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
BlockedReason (type alias)
Section titled “BlockedReason (type alias)”Why a pending task cannot be reserved. Derived from the registry, never stored.
Signature
type BlockedReason = "missing_task" | "task_too_old" | "migration_failed"Since v0.0.0
Inspection (interface)
Section titled “Inspection (interface)”Signature
export interface Inspection { readonly scheduling: "paused" | "running" | "closing" /** Every live task: pending, running, waiting and completing. */ readonly tasks: ReadonlyArray<TaskInspection>}Since v0.0.0
Options (interface)
Section titled “Options (interface)”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>}Since v0.0.0
Scheduler (class)
Section titled “Scheduler (class)”Signature
declare class SchedulerSince v0.0.0
SchedulerOutcome (type alias)
Section titled “SchedulerOutcome (type alias)”Outcomes only the scheduler writes, without running task code.
Signature
type SchedulerOutcome = Extract<TaskOutcome, { readonly status: "faulted" | "orphaned" }>Since v0.0.0
SchedulerShape (interface)
Section titled “SchedulerShape (interface)”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>}Since v0.0.0
TaskInspection (interface)
Section titled “TaskInspection (interface)”Signature
export interface TaskInspection { readonly record: TaskRecord readonly state: TaskInspectionState}Since v0.0.0
TaskInspectionState (type alias)
Section titled “TaskInspectionState (type alias)”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 }Since v0.0.0
abortConversation
Section titled “abortConversation”Signature
declare const abortConversation: (id: ConversationId, options?: { readonly background?: boolean }) => anySince v0.0.0
abortTask
Section titled “abortTask”Signature
declare const abortTask: (id: TaskId) => anySince v0.0.0
awaitTask
Section titled “awaitTask”Signature
declare const awaitTask: (id: TaskId) => anySince v0.0.0
inspect
Section titled “inspect”Signature
declare const inspect: Effect.Effect<Inspection, StorageFailure | SessionPoisoned, Scheduler>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>Since v0.0.0
resume
Section titled “resume”Signature
declare const resume: Effect.Effect<void, never, Scheduler>Since v0.0.0
waitForIdle
Section titled “waitForIdle”Signature
declare const waitForIdle: (conversationId?: ConversationId) => anySince v0.0.0