class Croupier::TaskManagerType

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.cr
croupier/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

Instance Method Detail

def add_input(task_key : String, input : String) : Bool #

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.


[View source]
def 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.


[View source]
def all_inputs : Set(String) #

[View source]
def auto_mode=(auto_mode : Bool) #

If true, it's running in auto mode


[View source]
def auto_mode? : Bool #

If true, it's running in auto mode


[View source]
def auto_run(targets : Array(String) = [] of String) #

[View source]
def auto_stop #

[View source]
def 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)


[View source]
def 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)


[View source]
def check_dependencies(targets : Array(String) | Nil = nil) #

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.


[View source]
def cleanup #

Remove all tasks and everything else (good for tests)


[View source]
def dependencies(outputs : Array(String)) #

The tasks needed to produce outputs (including themselves), in execution order.


[View source]
def dependencies(output : String) #

Single-output convenience overload.


[View source]
def depends_on(input : String) #

Outputs of every task that depends, directly or not, on input.


[View source]
def depends_on(inputs : Array(String)) #

[View source]
def early_cutoff=(early_cutoff : Bool) #

If true, enable early cutoff optimization (skip tasks when upstream outputs unchanged)


[View source]
def early_cutoff? : Bool #

If true, enable early cutoff optimization (skip tasks when upstream outputs unchanged)


[View source]
def fast_dirs=(fast_dirs : Bool) #

If true, directories depend on a list of files, not its contents


[View source]
def fast_dirs? : Bool #

If true, directories depend on a list of files, not its contents


[View source]
def fast_mode=(fast_mode : Bool) #

If true, only compare file dates


[View source]
def fast_mode? : Bool #

If true, only compare file dates


[View source]
def get(key) #

[View source]
def hash_directory(path : String) : String #

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).


[View source]
def inputs(targets : Array(String)) #

Every input of the given targets and their dependencies. Raises UnknownTaskError for an unknown target.


[View source]
def invalidate_graph_cache #

Mark the cached graph and input set stale. The autorun loop checks @graph_invalidated to decide whether to re-run.


[View source]
def last_run : Hash(String, String) #

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.


[View source]
def 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).

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.


[View source]
def lock_mutex(name : String) #

[View source]
def 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. 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.


[View source]
def modified : Set(String) #

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.


[View source]
def modified=(modified : Set(String)) #

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.


[View source]
def modified?(key : String) : Bool #

Whether key (a file or kv:// key) was modified since the last run.


[View source]
def mutexes : Hash(String, Sync::Mutex) #

A hash of mutexes required by tasks


[View source]
def mutexes=(mutexes : Hash(String, Sync::Mutex)) #

A hash of mutexes required by tasks


[View source]
def next_run : Hash(String, String) #

Output hashes recorded by tasks during this run.

See last_run for the concurrency contract.


[View source]
def next_run=(next_run : Hash(String, String)) #

Output hashes recorded by tasks during this run.

See last_run for the concurrency contract.


[View source]
def progress_callback : Proc(String, Nil) #

If set, it's called after every task finishes


[View source]
def progress_callback=(progress_callback : Proc(String, Nil)) #

If set, it's called after every task finishes


[View source]
def 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.

With run_all every task stays stale. Staleness then only orders the run (dependents wait for stale dependencies), so the root scan is skipped.


[View source]
def register_task(task : Task, explicit_id : String | Nil) : Nil #

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.


[View source]
def 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.

Raises UnknownTaskError for an unknown key, and UsageError during a run.


[View source]
def 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

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?)


[View source]
def save_run #

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.


[View source]
def scan_inputs(scope : Set(String) | Nil = nil) #

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.


[View source]
def set(key, value) : Bool #

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.


[View source]
def state_file : String #

Path to the state file that stores hashes between runs


[View source]
def state_file=(state_file : String) #

Path to the state file that stores hashes between runs


[View source]
def 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. Thread-safe for parallel task workers.


[View source]

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).


[View source]
def this_run : Hash(String, String) #

Input hashes scanned at the start of this run.

See last_run for the concurrency contract.


[View source]
def this_run=(this_run : Hash(String, String)) #

Input hashes scanned at the start of this run.

See last_run for the concurrency contract.


[View source]
def unlock_mutex(name : String) #

[View source]
def use_persistent_store(path : String) #

Use a persistent k/v store in this path instead of the default memory store


[View source]
def watch(targets : Array(String) = [] of String) : Nil #

Watch the inputs of targets (all tasks by default) and queue changed paths in @queued_changes. Changes made before this call are not detected.


[View source]
def with_state_lock(dry_run : Bool, &) #

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.


[View source]