class Croupier::Task

Overview

A Task is an object that may generate output

It has a Proc which is executed when the task is run It can have zero or more inputs It has zero or more outputs Tasks are connected by dependencies, where one task's output is another's input

Included Modules

Defined in:

task.cr

Constructors

Instance Method Summary

Constructor Detail

def self.new(ctx : YAML::ParseContext, node : YAML::Nodes::Node) #

[View source]
def self.new(outputs : Array(String) | String | Nil = nil, inputs : Array(String) = [] of String, proc : TaskProc | Nil = nil, no_save : Bool = false, id : String | Nil = nil, always_run : Bool = false, mergeable : Bool = true, mutex : String | Nil = nil, output : String | Nil = nil) #

[View source]
def self.new(*, __context_for_yaml_serializable ctx : YAML::ParseContext, __node_for_yaml_serializable node : YAML::Nodes::Node) #

[View source]
def self.new(outputs : Array(String) | String | Nil = nil, inputs : Array(String) = [] of String, no_save : Bool = false, id : String | Nil = nil, always_run : Bool = false, mergeable : Bool = true, mutex : String | Nil = nil, output : String | Nil = nil, &block : TaskProc) #

Create a task with zero or more outputs.

#outputs is an array of files or k/v store keys that the task generates. A single output can be passed as a string, or by name as output:. #inputs is an array of filesystem paths, task ids or k/v store keys that the task depends on The block (or proc:) is executed when the task is run no_save tells croupier that the task saves its outputs itself #id is a unique identifier for the task. If the task has no outputs, it must have an id. If not given, it's a hash of the outputs. always_run makes the task stale regardless of its inputs mergeable: if true, the task can be merged with others that share an output. Tasks with different mergeable values can NOT be merged together. #mutex names a lock held while the task's procs run, so tasks sharing it never run at the same time

k/v store keys are of the form kv://key, and are used to store intermediate data in a key/value store (in memory, or a file via TaskManager.use_persistent_store).

To access k/v data in your proc, use TaskManager.get(key).

Important: tasks are registered in TaskManager on creation, and creating one while a run is in progress raises UsageError. If the new task conflicts in id/outputs with others, it is merged into the existing one and the new object is NOT registered, so keeping references to Task objects you create is probably pointless.


[View source]

Instance Method Detail

def always_run=(always_run : Bool) #

[View source]
def always_run? : Bool #

[View source]
def id : String #

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

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

A task producing one of this task's inputs changed its outputs in the current run. compute_staleness can't see that: the producer is fresh once it ran, and the input scan predates the change. Set and cleared by the runner.


[View source]
def input_changed? : Bool #

A task producing one of this task's inputs changed its outputs in the current run. compute_staleness can't see that: the producer is fresh once it ran, and the input scan predates the change. Set and cleared by the runner.


[View source]
def inputs : Inputs #

The task's inputs: files, task ids or kv:// keys it depends on. Adding to it calls TaskManager.add_input.


[View source]
def keys #

Under what keys should this task be registered with TaskManager


[View source]
def merge(other : Task) #

Merge two tasks: inputs and outputs are joined, and the second task's procs are appended to the first's.


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

[View source]
def mergeable? : Bool #

[View source]
def mutex : String | Nil #

[View source]
def mutex=(name : String | Nil) #

Setting a mutex also registers it with the manager, which is where Task#run looks it up


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

[View source]
def no_save? : Bool #

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

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

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

[View source]
def outputs_changed? : Bool #

[View source]
def procs : Array(TaskProc) #

[View source]
def procs=(procs : Array(TaskProc)) #

[View source]
def ready?(run_all = false) #

A task is ready if it needs to run (stale, or run_all) and is not waiting for any input. always_run tasks are stale until they run.


[View source]
def recompute_staleness : Nil #

Early cutoff: one of the task's inputs turned out unchanged, so recompute staleness from all of them (another may still be stale, or may already have changed in this run).


[View source]
def register_with_manager(explicit_id : String | Nil) : Nil #

Register this task in the TaskManager: merge every task it has an output/id collision with into one, and register the survivor on every output/id of the merged set. Called by TaskManager.register_task. Not part of the public API. :nodoc:


[View source]
def run #

Executes the proc for the task


[View source]
def stale : Bool | Nil #

Tri-state staleness property: nil=unknown, true=stale, false=fresh.


[View source]
def stale=(value : Bool | Nil) #

[View source]
def stale? : Bool #

A task is stale if:

  • it is always_run, or has no inputs
  • one of its outputs is missing
  • one of its inputs was modified
  • one of its inputs is produced by a stale task

Staleness is tri-state: unknown, stale, fresh. TaskManager.propagate_staleness sets it for every task before a run, and running a task sets it to fresh. This method trusts an assigned value (dependents rely on a finished task reporting fresh even when it is always_run) and computes only while it is unknown.


[View source]
def stale_on_own? : Bool #

Whether the task is stale on its own account, regardless of the tasks producing its inputs: it is always_run or has no inputs, an output is missing (as a file or a k/v key), or an input was modified. TaskManager.propagate_staleness starts from these.


[View source]
def to_s(io) #

[View source]
def waiting? : Bool #

Early-exit version of waiting_for.empty?, used by ready?


[View source]
def waiting_for #

All inputs that are not satisfied yet


[View source]