Skip to content

Scheduler

The durable task scheduler: picks up pending tasks, runs each invocation in its own fiber, and applies the step rules, holds, aborts, waits, versioning and recovery described in site/src/content/docs/design/analysis-and-plan.md §4.6.

Building the layer recovers: tasks that were running go back to pending, abort cascades are re-derived, held outcomes are finalized, and (unless resume: false) work resumes. Closing the layer’s scope plays the role of pi’s close(): every invocation is interrupted, nothing is written, and the work resumes when the storage is reopened.

Since v0.0.0


Why a pending task cannot be reserved. Derived from the registry, never stored.

Signature

type BlockedReason = "missing_task" | "task_too_old" | "migration_failed"

Source

Since v0.0.0

Signature

export interface Inspection {
readonly scheduling: "paused" | "running" | "closing"
/** Every live task: pending, running, waiting and completing. */
readonly tasks: ReadonlyArray<TaskInspection>
}

Source

Since v0.0.0

Signature

export interface Options {
/** Start scheduling when the layer is built. Default `true`; with `false`, call `Scheduler.resume`. */
readonly resume?: boolean
/**
* Runs in the commit that makes an outcome the scheduler wrote itself (`faulted`, `orphaned`) terminal, which is the
* final commit after a hold. The Harness uses it to settle runs; the scheduler knows nothing about task kinds.
*/
readonly settleOutcome?: (record: TaskRecord, outcome: SchedulerOutcome) => Effect.Effect<void, WriteError, Tx>
/**
* Runs in the commit that admits a conversation abort, once for the aborted conversation and, with `background`, once
* for every other conversation of a reached task. The Harness uses it to withdraw queued inputs.
*/
readonly withdrawInputs?: (conversationId: ConversationId) => Effect.Effect<void, WriteError, Tx>
/**
* How the conversation handles of `TaskRuntime.conversation` admit and follow submissions. The Harness supplies it;
* without it, a handle's `submit` dies.
*/
readonly submitter?: Submitter
/**
* Called with the cause of a failure the scheduler cannot surface otherwise: a rejected commit it will retry, a
* failed migration, an invocation that died. It runs on the commit line in some cases, so it must not use the
* Session. Default: log a warning.
*/
readonly report?: (cause: Cause.Cause<unknown>) => Effect.Effect<void>
}

Source

Since v0.0.0

Signature

declare class Scheduler

Source

Since v0.0.0

Outcomes only the scheduler writes, without running task code.

Signature

type SchedulerOutcome = Extract<TaskOutcome, { readonly status: "faulted" | "orphaned" }>

Source

Since v0.0.0

Signature

export interface SchedulerShape {
/** Start scheduling if the layer was built with `resume: false`. Idempotent. */
readonly resume: Effect.Effect<void>
/**
* Commit an abort mark, interrupt and join the task's run invocation, and return. The abort handler then runs once
* the task's ordinary owned work is no longer live; wait for the outcome with `Task.await`. A task that no
* registered definition can take is settled `orphaned` instead. A `completing` task is only marked. Interrupting the
* caller only stops the wait; the mark stays.
*/
readonly abortTask: (id: TaskId) => Effect.Effect<"marked" | "terminal", TaskNotFound | SchedulerClosed | CommitError>
/**
* Mark every live non-background task that ordinary traversal from the conversation reaches (with `background`,
* every task it reaches), then wait until they are terminal and the conversation is idle.
*/
readonly abortConversation: (
id: ConversationId,
options?: { readonly background?: boolean }
) => Effect.Effect<void, SchedulerClosed | CommitError>
/** The terminal record of the task, waiting for it if it is live. */
readonly awaitTask: (
id: TaskId
) => Effect.Effect<TaskRecord, TaskNotFound | SchedulerClosed | StorageFailure | SessionPoisoned>
/**
* Wait until ordinary traversal from the conversation (or from every ownerless conversation) reaches no live
* non-background task. `completing` and blocked tasks count as live.
*/
readonly waitForIdle: (conversationId?: ConversationId) => Effect.Effect<void, SchedulerClosed>
/** Scheduling state and every live task with its derived state. Runs no task code. */
readonly inspect: Effect.Effect<Inspection, StorageFailure | SessionPoisoned>
}

Source

Since v0.0.0

Signature

export interface TaskInspection {
readonly record: TaskRecord
readonly state: TaskInspectionState
}

Source

Since v0.0.0

The derived scheduling state of one live task.

Signature

type TaskInspectionState =
| { readonly kind: "running" }
| { readonly kind: "completing" }
| { readonly kind: "waiting"; readonly on: ReadonlyArray<TaskId> }
| { readonly kind: "blocked"; readonly reason: BlockedReason; readonly error?: unknown }
| { readonly kind: "ready"; readonly migrates: boolean }

Source

Since v0.0.0

Signature

declare const abortConversation: (id: ConversationId, options?: { readonly background?: boolean }) => any

Source

Since v0.0.0

Signature

declare const abortTask: (id: TaskId) => any

Source

Since v0.0.0

Signature

declare const awaitTask: (id: TaskId) => any

Source

Since v0.0.0

Signature

declare const inspect: Effect.Effect<Inspection, StorageFailure | SessionPoisoned, Scheduler>

Source

Since v0.0.0

The scheduler over the Session, registry and dispatcher in context. The layer’s scope owns every invocation; the layer fails if recovery cannot commit.

Signature

declare const layer: (options?: Options) => Layer.Layer<Scheduler, CommitError, Session | TaskRegistry | TaskDispatcher>

Source

Since v0.0.0

Signature

declare const resume: Effect.Effect<void, never, Scheduler>

Source

Since v0.0.0

Signature

declare const waitForIdle: (conversationId?: ConversationId) => any

Source

Since v0.0.0