Skip to content

effect

WatchdogQueuedJob = Schema.Schema.Type<typeof WatchdogQueuedJobSchema>>

Generic durable work waiting for execution.


WatchdogSettledJob = Schema.Schema.Type<typeof WatchdogSettledJobSchema>>

Generic durable work with a terminal outcome.


WatchdogSettledRecord = Readonly<{ sequence: number; job: WatchdogSettledJob; }>

One bounded terminal-job recovery record with its durable ordering cursor.


WatchdogRunningJob = Schema.Schema.Type<typeof WatchdogRunningJobSchema>>

Generic durable work held by Watchdog’s private claim.


WatchdogStoredJob = Schema.Schema.Type<typeof WatchdogStoredJobSchema>>

Any publicly readable lifecycle state for a generic job.


WatchdogExecutionDecision = Schema.Schema.Type<typeof WatchdogExecutionDecisionSchema>>

Terminal decision returned by the host executor.


WatchdogCancellationIntent = Schema.Schema.Type<typeof WatchdogCancellationIntentSchema>>

Owner-issued durable cancellation intent.


WatchdogCancellationDecision = Readonly<{ status: "recorded" | "replayed"; intent: WatchdogCancellationIntent; }> | Readonly<{ status: "conflict" | "missing"; }>

Result of recording an owner cancellation in its surrounding transaction.


WatchdogTransactionalProjection = Readonly<{ readQueueHead: (queue) => WatchdogQueueHead; readLatestSettledJob: (queue) => WatchdogSettledJob | undefined; enqueueAcceptedJob: (job) => void; recordCancellationIntent: (intent) => WatchdogCancellationDecision; }>

Narrow synchronous seam for owner-atomic queue inspection, ingest, and cancellation.


WatchdogReadResult = Schema.Schema.Type<typeof WatchdogReadResultSchema>>

Result of reading one generic job and its cancellation state.


WatchdogQueueHead = Schema.Schema.Type<typeof WatchdogQueueHeadSchema>>

Current queued or running job for one queue, excluding terminal history.


WatchdogQueueTransition = Readonly<{ queue: string; head: WatchdogQueueHead; settled?: WatchdogSettledJob; }>

One queue’s state after a durable transition.


WatchdogQueueProjector = (transition) => void

Owner projection written inside the same SQLite transaction as every queue transition: ingest, claim, cancellation, settlement, recovery and purge. It must be a deterministic owner-local write. A throw fails the transition exactly like any other storage failure.

WatchdogQueueTransition

void


WatchdogIngestResult = Schema.Schema.Type<typeof WatchdogIngestResultSchema>>

Result of idempotent durable ingest.


WatchdogCancellationResult = Schema.Schema.Type<typeof WatchdogCancellationResultSchema>>

Result of recording a cancellation through the asynchronous runtime.


EffectWatchdogExecutor = Readonly<{ execute: (job) => Effect.Effect<WatchdogExecutionDecision, WatchdogError>>; }>

Effect-native host executor for one generic running job.


WatchdogSqliteValue = Schema.Schema.Type<typeof watchdogSqliteValueSchema>>

Parsed native SQLite scalar shared by the durable storage adapters.


WatchdogSqliteRow = Schema.Schema.Type<typeof watchdogSqliteRowSchema>>

Parsed scalar record; query-specific columns require their owner schema.


WatchdogSqliteCursor = Readonly<{ toArray: () => readonly WatchdogSqliteRow[]; }>

Structural result cursor required by Watchdog’s SQLite implementation.


WatchdogSqliteOwner = Readonly<{ sql: Readonly<{ exec: (statement, …bindings) => WatchdogSqliteCursor; }>; transactionSync: <A>>(operation) => A; }>

Native owner-local SQLite capability consumed by Watchdog.


EffectWatchdogOptions = Readonly<{ owner: WatchdogSqliteOwner; isAlive: (job) => Effect.Effect<boolean, WatchdogError>>; executor: EffectWatchdogExecutor; wake: Readonly<{ recompute: () => Effect.Effect<void, WatchdogError>>; }>; project?: WatchdogQueueProjector; }>

Effect-native construction inputs for the deep SQLite Watchdog.


WatchdogLivenessUnknown = Readonly<{ status: "liveness_unknown"; jobId: string; reason: WatchdogError["cause"]; }>

A running job’s liveness probe failed. The job is left exactly as it was: Watchdog neither treats it as alive nor fabricates a disappearance.


WatchdogTickResult = Readonly<{ status: "busy"; }> | Readonly<{ status: "no_ready"; }> | WatchdogLivenessUnknown | Readonly<{ status: "settled"; job: WatchdogSettledJob; }>

Result of one bounded deep tick.


WatchdogRuntimeEffect = Readonly<{ hasPendingWork: Effect.Effect<boolean, WatchdogError>>; ingest: (job) => Effect.Effect<WatchdogIngestResult, WatchdogError>>; read: (jobId) => Effect.Effect<WatchdogReadResult, WatchdogError>>; readQueueHead: (queue) => Effect.Effect<WatchdogQueueHead, WatchdogError>>; readSettled: (limit, afterSequence?) => Effect.Effect<readonly WatchdogSettledRecord[], WatchdogError>>; tick: () => Effect.Effect<WatchdogTickResult, WatchdogError>>; cancel: (input) => Effect.Effect<WatchdogCancellationResult, WatchdogError>>; purge: (queue) => Effect.Effect<void, WatchdogError>>; transactional: WatchdogTransactionalProjection; }>

Effect-native Watchdog runtime surface.

const WatchdogQueuedJobSchema: Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: Int; }>>; state: Literal<"queued">>; }>

The public queued-job grammar. Payload is deliberately uninterpreted.


const WatchdogSettledJobSchema: Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: Int; }>>; state: Literal<"settled">>; outcome: Literals<readonly ["completed", "failed", "cancelled", "outcome_unknown", "interrupted"]>; failureReason: optionalKey<String>>; }>

The public terminal-job grammar.


const WatchdogRunningJobSchema: Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: Int; }>>; state: Literal<"running">>; }>

The public running-job grammar. Private claim data is deliberately absent.


const WatchdogStoredJobSchema: Union<readonly [Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: Int; }>>; state: Literal<"queued">>; }>, Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: Int; }>>; state: Literal<"running">>; }>, Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: Int; }>>; state: Literal<"settled">>; outcome: Literals<readonly ["completed", "failed", "cancelled", "outcome_unknown", "interrupted"]>; failureReason: optionalKey<String>>; }>]>

The complete public job lifecycle grammar.


const WatchdogExecutionDecisionSchema: Union<readonly [Struct<{ kind: Literal<"completed">>; }>, Struct<{ kind: Literal<"failed">>; reason: optionalKey<String>>; }>, Struct<{ kind: Literal<"cancelled">>; }>, Struct<{ kind: Literal<"outcome_unknown">>; }>]>

Terminal decisions returned by the execution adapter.


const WatchdogCancellationIntentSchema: Struct<{ jobId: String; cancellationId: String; onlyIfQueued: optional<Boolean>>; }>

Durable cancellation requested by the enclosing owner.


const WatchdogReadResultSchema: Union<readonly [Struct<{ status: Literal<"found">>; job: Union<readonly [Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: …; }>>; state: Literal<"queued">>; }>, Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: …; }>>; state: Literal<"running">>; }>, Struct<{ jobId: String; queue: String; lane: String; priority: Int; payload: Unknown; recovery: optional<Struct<{ maxRecoveries: …; }>>; state: Literal<"settled">>; outcome: Literals<readonly ["completed", "failed", "cancelled", "outcome_unknown", "interrupted"]>; failureReason: optionalKey<String>>; }>]>; cancellation: Union<readonly [Struct<{ status: Literal<"absent">>; }>, Struct<{ status: Literal<"present">>; cancellation: Struct<{ jobId: String; cancellationId: String; onlyIfQueued: optional<…>; }>; }>]>; }>, Struct<{ status: Literal<"missing">>; }>]>

Results returned by public job inspection.


const WatchdogQueueHeadSchema: Union<readonly [Struct<{ state: Literal<"idle">>; }>, Struct<{ state: Literals<readonly ["queued", "running"]>; jobId: String; cancellationRequested: Boolean; }>]>

Bounded current-work projection for one generic queue.


const WatchdogIngestResultSchema: Union<readonly [Struct<{ status: Literal<"queued">>; }>, Struct<{ status: Literal<"duplicate">>; }>]>

Results returned by idempotent public ingest.


const WatchdogCancellationResultSchema: Union<readonly [Struct<{ status: Literals<readonly ["recorded", "replayed"]>; }>, Struct<{ status: Literals<readonly ["conflict", "missing"]>; }>]>

Results returned by durable cancellation recording.


const watchdogSqliteValueSchema: Union<readonly [instanceOf<ArrayBuffer, unknown>>, String, Number, Null]>

Scalar grammar shared by native SQLite adapters and their storage owners.


const watchdogSqliteRowSchema: $Record<String, Union<readonly [instanceOf<ArrayBuffer, unknown>>, String, Number, Null]>>

A native row is a scalar record; each query owner decodes its own projection.

makeRuntime(options): Effect<Readonly<{ hasPendingWork: Effect<boolean, WatchdogError>>; ingest: (job) => Effect<{ status: "queued"; } | { status: "duplicate"; }, WatchdogError>>; read: (jobId) => Effect<{ status: "found"; job: { jobId: string; queue: string; lane: string; priority: number; payload: unknown; recovery?: { maxRecoveries: …; }; state: "queued"; } | { jobId: string; queue: string; lane: string; priority: number; payload: unknown; recovery?: { maxRecoveries: …; }; state: "running"; } | { jobId: string; queue: string; lane: string; priority: number; payload: unknown; recovery?: { maxRecoveries: …; }; state: "settled"; outcome: "failed" | "completed" | "interrupted" | "cancelled" | "outcome_unknown"; failureReason?: string; }; cancellation: { status: "absent"; } | { status: "present"; cancellation: { jobId: string; cancellationId: string; onlyIfQueued?: … | … | …; }; }; } | { status: "missing"; }, WatchdogError>>; readQueueHead: (queue) => Effect<{ state: "idle"; } | { state: "running" | "queued"; jobId: string; cancellationRequested: boolean; }, WatchdogError>>; readSettled: (limit, afterSequence?) => Effect<readonly Readonly<{ sequence: number; job: { jobId: string; queue: string; lane: string; priority: number; payload: unknown; recovery?: { maxRecoveries: …; }; state: "settled"; outcome: "failed" | "completed" | "interrupted" | "cancelled" | "outcome_unknown"; failureReason?: string; }; }>[], WatchdogError>>; tick: () => Effect<WatchdogTickResult, WatchdogError>>; cancel: (input) => Effect<{ status: "replayed" | "recorded"; } | { status: "missing" | "conflict"; }, WatchdogError>>; purge: (queue) => Effect<void, WatchdogError>>; transactional: WatchdogTransactionalProjection; }>, WatchdogError>>

One deep tick over Watchdog-owned owner-local SQLite state.

EffectWatchdogOptions

Effect<Readonly<{ hasPendingWork: Effect<boolean, WatchdogError>; ingest: (job) => Effect<{ status: "queued"; } | { status: "duplicate"; }, WatchdogError>; read: (jobId) => Effect<{ status: "found"; job: { jobId: string; queue: string; lane: string; priority: number; payload: unknown; recovery?: { maxRecoveries: …; }; state: "queued"; } | { jobId: string; queue: string; lane: string; priority: number; payload: unknown; recovery?: { maxRecoveries: …; }; state: "running"; } | { jobId: string; queue: string; lane: string; priority: number; payload: unknown; recovery?: { maxRecoveries: …; }; state: "settled"; outcome: "failed" | "completed" | "interrupted" | "cancelled" | "outcome_unknown"; failureReason?: string; }; cancellation: { status: "absent"; } | { status: "present"; cancellation: { jobId: string; cancellationId: string; onlyIfQueued?: … | … | …; }; }; } | { status: "missing"; }, WatchdogError>; readQueueHead: (queue) => Effect<{ state: "idle"; } | { state: "running" | "queued"; jobId: string; cancellationRequested: boolean; }, WatchdogError>; readSettled: (limit, afterSequence?) => Effect<readonly Readonly<{ sequence: number; job: { jobId: string; queue: string; lane: string; priority: number; payload: unknown; recovery?: { maxRecoveries: …; }; state: "settled"; outcome: "failed" | "completed" | "interrupted" | "cancelled" | "outcome_unknown"; failureReason?: string; }; }>[], WatchdogError>; tick: () => Effect<WatchdogTickResult, WatchdogError>; cancel: (input) => Effect<{ status: "replayed" | "recorded"; } | { status: "missing" | "conflict"; }, WatchdogError>; purge: (queue) => Effect<void, WatchdogError>; transactional: WatchdogTransactionalProjection; }>, WatchdogError>