class
Croupier::TaskManagerType
- Croupier::TaskManagerType
- Reference
- Object
Overview
TaskManager is a singleton that keeps track of all tasks. Its methods live in focused files under src/croupier/: kv_store.cr, hash_state.cr, graph.cr, runner.cr and watcher.cr.
Defined in:
croupier.crcroupier/graph.cr
croupier/hash_state.cr
croupier/kv_store.cr
croupier/runner.cr
croupier/watcher.cr
Constant Summary
-
AUTORUN_RETRY_MAX_DELAY =
1.0 -
AUTORUN_RETRY_MIN_DELAY =
0.01 -
Autorun retry backoff bounds, in seconds: consecutive failures slow the loop from the change-poll minimum to at most one attempt per second; any success resets to the minimum.
-
FAST_MODE_GRACE =
1.0 -
Seconds subtracted from the fast-mode baseline, absorbing timestamp granularity and clock skew. The cost is sometimes re-detecting an input modified just before the previous scan.
-
NO_CONSUMERS =
[] of Task -
Shared empty for consumers.fetch misses
-
SCAN_CHUNK_SIZE =
64 -
How many files each scan worker hashes at a time: small enough that one huge file doesn't hold up much other work, big enough that channel overhead stays negligible
-
STATE_VERSION =
"1" -
Version of the state-file schema, stored as __version. A mismatch (or a file with no version) discards all recorded hashes: one full rebuild instead of comparing hashes computed by a different scheme.
Instance Method Summary
-
#add_input(task_key : String, input : String) : Bool
Add
inputto the inputs of the task registered astask_key. -
#add_mutex(name : String)
Register the mutex
name, keeping the existing lock if there is one: replacing it could swap out a lock a running task holds, and tasks sharing the name would stop excluding each other. - #all_inputs : Set(String)
-
#auto_mode=(auto_mode : Bool)
If true, it's running in auto mode
-
#auto_mode? : Bool
If true, it's running in auto mode
- #auto_run(targets : Array(String) = [] of String)
- #auto_stop
-
#before_run_hook : Proc(Set(String), Nil)
If set, it's called in auto mode after changes are detected but before tasks run Receives the set of modified paths (files and kv:// keys)
-
#before_run_hook=(before_run_hook : Proc(Set(String), Nil))
If set, it's called in auto mode after changes are detected but before tasks run Receives the set of modified paths (files and kv:// keys)
-
#check_dependencies(targets : Array(String) | Nil = nil)
Raise UnknownInputsError unless every input is a kv:// key, a task, or an existing file.
-
#cleanup
Remove all tasks and everything else (good for tests)
-
#dependencies(outputs : Array(String))
The tasks needed to produce
outputs(including themselves), in execution order. -
#dependencies(output : String)
Single-output convenience overload.
-
#depends_on(input : String)
Outputs of every task that depends, directly or not, on
input. - #depends_on(inputs : Array(String))
-
#early_cutoff=(early_cutoff : Bool)
If true, enable early cutoff optimization (skip tasks when upstream outputs unchanged)
-
#early_cutoff? : Bool
If true, enable early cutoff optimization (skip tasks when upstream outputs unchanged)
-
#fast_dirs=(fast_dirs : Bool)
If true, directories depend on a list of files, not its contents
-
#fast_dirs? : Bool
If true, directories depend on a list of files, not its contents
-
#fast_mode=(fast_mode : Bool)
If true, only compare file dates
-
#fast_mode? : Bool
If true, only compare file dates
- #get(key)
-
#hash_directory(path : String) : String
Hash a single directory input.
-
#inputs(targets : Array(String))
Every input of the given targets and their dependencies.
-
#invalidate_graph_cache
Mark the cached graph and input set stale.
-
#last_run : Hash(String, String)
Hashes recorded by the previous run (from the state file, or folded in by the previous auto mode cycle).
-
#last_run=(last_run : Hash(String, String))
Hashes recorded by the previous run (from the state file, or folded in by the previous auto mode cycle).
- #lock_mutex(name : String)
-
#mark_stale_inputs(run_all : Bool = false, targets : Array(String) | Nil = nil)
Compare inputs against the last run and leave the changed ones in @modified for propagate_staleness.
-
#modified : Set(String)
Inputs (files and kv:// keys) modified since the last run; they make the tasks that consume them stale.
-
#modified=(modified : Set(String))
Inputs (files and kv:// keys) modified since the last run; they make the tasks that consume them stale.
-
#modified?(key : String) : Bool
Whether
key(a file or kv:// key) was modified since the last run. -
#mutexes : Hash(String, Sync::Mutex)
A hash of mutexes required by tasks
-
#mutexes=(mutexes : Hash(String, Sync::Mutex))
A hash of mutexes required by tasks
-
#next_run : Hash(String, String)
Output hashes recorded by tasks during this run.
-
#next_run=(next_run : Hash(String, String))
Output hashes recorded by tasks during this run.
-
#progress_callback : Proc(String, Nil)
If set, it's called after every task finishes
-
#progress_callback=(progress_callback : Proc(String, Nil))
If set, it's called after every task finishes
-
#propagate_staleness(run_all : Bool = false)
Mark every task stale or fresh in one O(V+E) pass: find the tasks stale on their own, then everything downstream of them.
-
#register_task(task : Task, explicit_id : String | Nil) : Nil
Register a newly constructed task (called by Task.new).
-
#remove_task(task_key : String) : Nil
Remove the task registered as
task_key, under every key it is registered under (one per output) and its id. -
#run_tasks(targets : Array(String) | Nil = nil, run_all : Bool = false, dry_run : Bool = false, parallel : Bool = false, keep_going : Bool = false, early_cutoff : Bool | Nil = nil)
Run the stale tasks needed to create or update
targets(every task when nil), in dependency order -
#save_run
Save the current state.
-
#scan_inputs(scope : Set(String) | Nil = nil)
Scan the given inputs (all of them by default) and return a hash with their sha1.
-
#set(key, value) : Bool
Store a value and return whether it changed.
-
#state_file : String
Path to the state file that stores hashes between runs
-
#state_file=(state_file : String)
Path to the state file that stores hashes between runs
-
#swap_output_hash(output : String, new_hash : String) : String | Nil
Record the hash of a task output for the next run's state file, and return the hash the last run recorded for it, in one locked step.
-
#tasks : Croupier::TaskRegistry
Registry of all tasks, keyed by each output (or by id for tasks without outputs).
-
#this_run : Hash(String, String)
Input hashes scanned at the start of this run.
-
#this_run=(this_run : Hash(String, String))
Input hashes scanned at the start of this run.
- #unlock_mutex(name : String)
-
#use_persistent_store(path : String)
Use a persistent k/v store in this path instead of the default memory store
-
#watch(targets : Array(String) = [] of String) : Nil
Watch the inputs of
targets(all tasks by default) and queue changed paths in @queued_changes. -
#with_state_lock(dry_run : Bool, &)
Serialize runs on the same state file across processes.
Instance Method Detail
Add input to the inputs of the task registered as task_key.
The supported way to grow a task's dependencies, including from task procs on parallel workers. The graph and staleness are computed before tasks run, so the new dependency takes effect on the next run: during a run the addition is queued, and applied when the last overlapping run ends.
Returns false if the task already had the input. Raises
UnknownTaskError for an unregistered task_key, and CycleError
if input is one of the task's own keys.
Register the mutex name, keeping the existing lock if there is
one: replacing it could swap out a lock a running task holds,
and tasks sharing the name would stop excluding each other.
If set, it's called in auto mode after changes are detected but before tasks run Receives the set of modified paths (files and kv:// keys)
If set, it's called in auto mode after changes are detected but before tasks run Receives the set of modified paths (files and kv:// keys)
Raise UnknownInputsError unless every input is a kv:// key, a
task, or an existing file. With targets, only the inputs of
their dependency closure are checked.
The tasks needed to produce outputs (including themselves),
in execution order.
If true, enable early cutoff optimization (skip tasks when upstream outputs unchanged)
If true, enable early cutoff optimization (skip tasks when upstream outputs unchanged)
Hash a single directory input.
The digest is a hash of hashes, like a git tree: every file in the tree is hashed (in parallel), and the per-file hashes are folded into one SHA1 together with the sorted entry list. The entry list and the file hashes are separate, delimited fields, so two different trees can't produce the same input bytes.
Public because Task#run hashes no_save directory outputs with it. The digest must match what scan_inputs computes for the same directory, or a task consuming it would re-run every time. Safe to call from parallel task workers (it only uses local channels).
Every input of the given targets and their dependencies. Raises UnknownTaskError for an unknown target.
Mark the cached graph and input set stale. The autorun loop checks @graph_invalidated to decide whether to re-run.
Hashes recorded by the previous run (from the state file, or folded in by the previous auto mode cycle).
Concurrency contract shared by last_run / this_run / next_run: task workers only touch them through the @hashes_lock accessors in hash_state.cr (swap_output_hash). Every other access happens on the coordinating fiber (the run_tasks caller, or the autorun fiber in auto mode), so those sites take no lock.
Hashes recorded by the previous run (from the state file, or folded in by the previous auto mode cycle).
Concurrency contract shared by last_run / this_run / next_run: task workers only touch them through the @hashes_lock accessors in hash_state.cr (swap_output_hash). Every other access happens on the coordinating fiber (the run_tasks caller, or the autorun fiber in auto mode), so those sites take no lock.
Compare inputs against the last run and leave the changed ones in @modified for propagate_staleness. Three modes:
auto the watcher reported what changed; hashing confirms it, so unchanged rewrites don't retrigger fast mtime against the last run's scan start, no hashing content content hashes against the last run's hashes
run_all skips fast mode's mtime scan (every task re-runs
anyway). Content mode still scans, because @this_run feeds the
state file. targets limits the scan to those tasks' inputs;
other inputs keep their recorded hashes.
Inputs (files and kv:// keys) modified since the last run; they make the tasks that consume them stale.
Task procs touch this set from parallel workers (#set marks
kv:// keys, #modified? reads it), so every internal access goes
through @modified_lock. Mutating it directly from a running proc
races.
Inputs (files and kv:// keys) modified since the last run; they make the tasks that consume them stale.
Task procs touch this set from parallel workers (#set marks
kv:// keys, #modified? reads it), so every internal access goes
through @modified_lock. Mutating it directly from a running proc
races.
Whether key (a file or kv:// key) was modified since the last run.
Output hashes recorded by tasks during this run.
See last_run for the concurrency contract.
Output hashes recorded by tasks during this run.
See last_run for the concurrency contract.
If set, it's called after every task finishes
Mark every task stale or fresh in one O(V+E) pass: find the tasks stale on their own, then everything downstream of them.
With run_all every task stays stale. Staleness then only orders the run (dependents wait for stale dependencies), so the root scan is skipped.
Register a newly constructed task (called by Task.new).
The task set is fixed while a run executes: runs read the registries without a lock. Raises UsageError during a run. The check and the registration happen under one lock acquisition, so a run can't start in between.
Remove the task registered as task_key, under every key it is
registered under (one per output) and its id.
Raises UnknownTaskError for an unknown key, and UsageError during a run.
Run the stale tasks needed to create or update targets (every
task when nil), in dependency order
If run_all is true, run non-stale tasks too
If dry_run is true, only log what would be done, but don't do it
If parallel is true, run tasks in parallel
If keep_going is true, keep going even if a task fails
If early_cutoff is true, skip tasks when upstream outputs are
unchanged (defaults to TaskManager.early_cutoff?)
Save the current state. It is written to a temporary file and renamed into place, so a crash can't leave a truncated state file. The temp name includes the PID so two processes sharing a directory never write the same temp file.
Scan the given inputs (all of them by default) and return a hash with their sha1. Files, including those inside directory inputs, are hashed in parallel by a worker pool bounded by CPU count.
Store a value and return whether it changed. A same-value set is a no-op, so identical kv outputs don't re-stale their dependents.
Record the hash of a task output for the next run's state file, and return the hash the last run recorded for it, in one locked step. Thread-safe for parallel task workers.
Registry of all tasks, keyed by each output (or by id for tasks
without outputs). Read-only: change it with Task.new and
#remove_task, which raise UsageError during a run (runs read
the registry without locks).
Input hashes scanned at the start of this run.
See last_run for the concurrency contract.
Input hashes scanned at the start of this run.
See last_run for the concurrency contract.
Use a persistent k/v store in this path instead of the default memory store
Watch the inputs of targets (all tasks by default) and queue
changed paths in @queued_changes. Changes made before this call
are not detected.
Serialize runs on the same state file across processes. Without this, two processes in one directory would race their read-scan-run-save cycles and the last writer would erase the other's results.
flock is released by the kernel when the holder dies, so there is no stale lock to clean up. The lock file is never deleted: that would let a third process lock a new inode while the old one is still held. The lock is polled without blocking so other fibers keep running while waiting. Not reentrant: nothing inside a run may call run_tasks again.