export type TaskStatus = 'pending' | 'assigned' | 'running' | 'paused' | 'blocked' | 'completed' | 'failed' | 'timeout' | 'cancelled'; export interface TaskError { code?: string; message: string; details?: Record; } export interface TaskAuthConfig { rules: Array<{ match: { scope: PermissionScope[]; }; require: { claims?: Record; sub?: string[]; }; }>; } export interface WebhookConfig { url: string; filter?: SubscribeFilter; secret?: string; wrap?: boolean; retry?: RetryConfig; } export interface RetryConfig { retries: number; backoff: 'fixed' | 'exponential' | 'linear'; initialDelayMs: number; maxDelayMs: number; timeoutMs: number; } export type SeriesMode = 'keep-all' | 'accumulate' | 'latest'; export type Level = 'debug' | 'info' | 'warn' | 'error'; export type PermissionScope = 'task:create' | 'task:manage' | 'event:publish' | 'event:subscribe' | 'event:history' | 'webhook:create' | 'worker:connect' | 'worker:manage' | 'task:resolve' | 'task:signal' | '*'; export interface CleanupRule { name?: string; match?: { taskTypes?: string[]; status?: TaskStatus[]; }; trigger: { afterMs?: number; }; target: 'all' | 'events' | 'task'; eventFilter?: { types?: string[]; levels?: Level[]; olderThanMs?: number; seriesMode?: SeriesMode[]; }; } export interface BlockedRequest { type: string; data: unknown; } export type AssignMode = 'external' | 'pull' | 'ws-offer' | 'ws-race'; export type DisconnectPolicy = 'reassign' | 'mark' | 'fail'; export type WorkerStatus = 'idle' | 'busy' | 'draining' | 'offline'; export interface TagMatcher { all?: string[]; any?: string[]; none?: string[]; } export interface WorkerMatchRule { taskTypes?: string[]; tags?: TagMatcher; } export interface Worker { id: string; status: WorkerStatus; matchRule: WorkerMatchRule; capacity: number; usedSlots: number; weight: number; connectionMode: 'pull' | 'websocket'; connectedAt: number; lastHeartbeatAt: number; metadata?: Record; } export type WorkerAssignmentStatus = 'offered' | 'assigned' | 'running'; export interface WorkerAssignment { taskId: string; workerId: string; cost: number; assignedAt: number; status: WorkerAssignmentStatus; } export interface WorkerAuditEvent { id: string; workerId: string; timestamp: number; action: 'connected' | 'disconnected' | 'updated' | 'task_assigned' | 'task_declined' | 'task_reclaimed' | 'draining' | 'heartbeat_timeout' | 'pull_request'; data?: Record; } export interface Task { id: string; type?: string; status: TaskStatus; params?: Record; result?: Record; error?: TaskError; metadata?: Record; createdAt: number; updatedAt: number; completedAt?: number; ttl?: number; authConfig?: TaskAuthConfig; webhooks?: WebhookConfig[]; cleanup?: { rules: CleanupRule[]; }; tags?: string[]; assignMode?: AssignMode; cost?: number; assignedWorker?: string; reason?: string; resumeAt?: number; blockedRequest?: BlockedRequest; disconnectPolicy?: DisconnectPolicy; } export interface TaskEvent { id: string; taskId: string; index: number; timestamp: number; type: string; level: Level; data: unknown; seriesId?: string; seriesMode?: SeriesMode; seriesAccField?: string; seriesSnapshot?: boolean; /** Transient: accumulated data attached during broadcast, not persisted in ShortTermStore */ _accumulatedData?: unknown; } /** * Archive-persistable event shape. * * TaskArchive v1 stores a compacted, replayable event stream for one task: * indexes must be contiguous from 0, latest-mode histories are latest-only, * and accumulate-mode histories may be stored as accumulated snapshots. * Presentation/transient event fields such as collapsed `seriesSnapshot` events * and broadcast `_accumulatedData` are not valid archive data. */ export type TaskArchiveEvent = Omit; export interface TaskArchive { schema: 'taskcast.taskArchive'; version: 1; exportedAt: number; task: Task; /** Compacted, replayable event stream for the task, ordered by contiguous indexes from 0. */ events: TaskArchiveEvent[]; } export interface TaskArchiveImportOptions { overwrite?: boolean; } export interface TaskArchiveImportResult { taskId: string; eventCount: number; overwritten: boolean; } export interface SeriesLatestEntry { taskId: string; seriesId: string; event: TaskArchiveEvent; } export interface TaskArchiveRestoreData { task: Task; events: TaskArchiveEvent[]; nextIndex: number; seriesLatest: SeriesLatestEntry[]; } export interface SSEEnvelope { filteredIndex: number; rawIndex: number; eventId: string; taskId: string; type: string; timestamp: number; level: Level; data: unknown; seriesId?: string; seriesMode?: SeriesMode; seriesAccField?: string; seriesSnapshot?: boolean; } export interface SinceCursor { id?: string; index?: number; timestamp?: number; } export type SeriesFormat = 'delta' | 'accumulated'; export interface SubscribeFilter { since?: SinceCursor; types?: string[]; levels?: Level[]; includeStatus?: boolean; wrap?: boolean; seriesFormat?: SeriesFormat; } export interface EventQueryOptions { since?: SinceCursor; limit?: number; } export interface SeriesResult { /** The original delta event (stored in ShortTermStore) */ event: TaskEvent; /** The event with accumulated data (for LongTermStore + broadcast). Undefined for non-accumulate modes. */ accumulatedEvent?: TaskEvent; /** Whether processSeries already stored the event (e.g. latest mode uses replaceLastSeriesEvent). */ stored?: boolean; } export type StorageState = 'hot' | 'releasing' | 'cold'; export interface TaskStorageMetadata { taskId: string; storageState: StorageState; storageEpoch: number; activeReleaseGeneration: string | null; archiveWatermark: number; lastEventAt: number | null; coldAt: number | null; executionDeadlineAt: number | null; taskVersion: number; } export interface HotWriteToken { taskId: string; storageEpoch: number; } /** * A task plus an opaque compare-and-set token for its exact hot-store record. * Redis adapters use the raw serialized task so writes from older instances * that do not maintain a separate revision key are still detected. */ export interface TaskMutationSnapshot { task: Task; revision: string; } export interface StorageLease { taskId: string; lockToken: string; generation: string; storageEpoch: number; } export interface TaskWriteFence { taskId: string; acceptingWrites: boolean; storageEpoch: number; activeReleaseGeneration: string | null; } export interface ClosedWriteFence extends TaskWriteFence { acceptingWrites: false; highWatermark: number; } export interface ReleasePreconditions { expectedLastEventIndex: number; inactiveSince: number; } export interface ReleaseResult { taskId: string; storageState: StorageState; archiveWatermark: number; released: boolean; } export interface ArchiveSourceManifest { priorWatermark: number; targetWatermark: number; sourceEntryCount: number; sourceDigest: string; seriesStateDigest: string; expectedBatchOrdinals: number[]; } export type ArchiveGenerationStatus = 'open' | 'finalized' | 'aborted'; export interface ArchiveGeneration { taskId: string; generation: string; storageEpoch: number; targetWatermark: number; manifest: ArchiveSourceManifest; status: ArchiveGenerationStatus; createdAt: number; updatedAt: number; } export interface ArchiveBatchReceipt { taskId: string; generation: string; ordinal: number; previousBatchDigest: string | null; batchDigest: string; entryCount: number; firstIndex: number | null; lastIndex: number | null; } export interface ArchiveSourcePage { taskId: string; watermark: number; cursor: string | null; nextCursor: string | null; events: TaskEvent[]; done: boolean; } export interface DurableSeriesState { taskId: string; seriesId: string; mode: 'latest' | 'accumulate'; event: TaskEvent; throughIndex: number; } export interface RehydrateSnapshot { task: Task; archiveWatermark: number; maxEventIndex: number; replayEvents: TaskEvent[]; seriesLatest: DurableSeriesState[]; storageEpoch: number; } export interface CanonicalHistoryEntry { event: TaskEvent; seriesThroughIndex?: number; } export interface TtlClaim { taskId: string; claimToken: string; claimUntil: number; taskVersion: number; executionDeadlineAt: number; } export interface TerminalProjection { projectionId: string; task: Task; event: TaskEvent; assignment: WorkerAssignment | null; claimToken: string | null; claimUntil: number | null; } export interface TerminalProjectionResult { token: HotWriteToken; projected: boolean; } export interface TaskStoragePresence { task: boolean; eventCount: number; nextIndex: boolean; seriesStateCount: number; writeFence: boolean; } export interface StorageWriterRegistration { instanceId: string; storageProtocolVersion: number; build: string; expiresAt: number; } export interface StorageReleaseRequest { taskId: string; requestedAt: number; expectedLastEventIndex: number; inactiveSince: number; } export interface TaskStorageMetadataCas { taskId: string; expectedStorageState: StorageState; expectedStorageEpoch: number; expectedReleaseGeneration: string | null; next: TaskStorageMetadata; } export interface ArchiveBatch { receipt: ArchiveBatchReceipt; events: TaskEvent[]; seriesLatest: DurableSeriesState[]; } export declare abstract class TaskStorageError extends Error { abstract readonly code: string; abstract readonly retryable: boolean; protected constructor(message: string); } export declare class StorageFenceConflictError extends TaskStorageError { readonly code = "storage_fence_conflict"; readonly retryable = true; constructor(message?: string); } export declare class StorageBusyError extends TaskStorageError { readonly code = "storage_busy"; readonly retryable = true; constructor(message?: string); } export declare class StoragePreconditionError extends TaskStorageError { readonly code = "storage_precondition_failed"; readonly retryable = false; constructor(message?: string); } export declare class StorageIntegrityError extends TaskStorageError { readonly code = "storage_integrity_error"; readonly retryable = false; constructor(message?: string); } export declare class StorageUnavailableError extends TaskStorageError { readonly code = "storage_unavailable"; readonly retryable = true; constructor(message?: string); } export declare class StorageReleaseUnsupportedError extends TaskStorageError { readonly code = "storage_release_unsupported"; readonly retryable = false; constructor(message?: string); } export interface BroadcastProvider { publish(channel: string, event: TaskEvent): Promise; subscribe(channel: string, handler: (event: TaskEvent) => void): () => void; } export interface ShortTermStore { /** True only when every lifecycle operation below is implemented atomically. */ readonly supportsHotColdRelease?: boolean; saveTask(task: Task): Promise; getTask(taskId: string): Promise; /** Atomically allocates the next event index for a task. */ nextIndex(taskId: string): Promise; appendEvent(taskId: string, event: TaskEvent): Promise; getEvents(taskId: string, opts?: EventQueryOptions): Promise; setTTL(taskId: string, ttlSeconds: number): Promise; getSeriesLatest(taskId: string, seriesId: string): Promise; setSeriesLatest(taskId: string, seriesId: string, event: TaskEvent): Promise; /** Atomically read previous accumulated value, concatenate with new delta, write back. Returns the accumulated event. */ accumulateSeries(taskId: string, seriesId: string, event: TaskEvent, field: string): Promise; replaceLastSeriesEvent(taskId: string, seriesId: string, event: TaskEvent): Promise; /** Validates deterministic archive restore conflicts before mutation; engine calls this before multi-store restore. */ validateTaskArchiveRestore?(data: TaskArchiveRestoreData, options?: TaskArchiveImportOptions): Promise; /** Stores with native archive restore should implement this; engine import checks availability before use. */ restoreTaskArchive?(data: TaskArchiveRestoreData, options?: TaskArchiveImportOptions): Promise<{ overwritten: boolean; }>; acquireStorageLock?(taskId: string, lockToken: string, generation: string, ttlMs: number): Promise; renewStorageLock?(lease: StorageLease, ttlMs: number): Promise; releaseStorageLock?(lease: StorageLease): Promise; getWriteFence?(taskId: string): Promise; closeWriteFence?(lease: StorageLease, expectedEpoch: number): Promise; reopenWriteFence?(lease: StorageLease, expectedEpoch: number): Promise; commitEventFenced?(taskId: string, event: Omit, token: HotWriteToken): Promise; /** Atomically reads a task and the opaque token for its current mutation. */ getTaskMutationSnapshot?(taskId: string): Promise; /** Atomically stores a task mutation and its derived non-series events. */ commitTaskEventsFenced?(task: Task, expectedRevision: string, events: Omit[], token: HotWriteToken): Promise; saveTaskFenced?(task: Task, token: HotWriteToken): Promise; readArchiveSourcePage?(taskId: string, watermark: number, cursor: string | null, limit: number): Promise; deleteTaskStorageFenced?(lease: StorageLease, expectedEpoch: number): Promise; restoreHotTaskFenced?(snapshot: RehydrateSnapshot, lease: StorageLease, nextEpoch: number): Promise; /** Atomically projects a durable terminal task/event and settles its assignment. */ projectTerminalFenced?(projection: TerminalProjection, lease: StorageLease, expectedEpoch: number, nextEpoch: number): Promise; getTaskStoragePresence?(taskId: string): Promise; registerStorageWriter?(registration: StorageWriterRegistration, ttlMs: number): Promise; listStorageWriters?(): Promise; listTasks(filter: TaskFilter): Promise; saveWorker(worker: Worker): Promise; getWorker(workerId: string): Promise; listWorkers(filter?: WorkerFilter): Promise; deleteWorker(workerId: string): Promise; claimTask(taskId: string, workerId: string, cost: number): Promise; addAssignment(assignment: WorkerAssignment): Promise; removeAssignment(taskId: string): Promise; getWorkerAssignments(workerId: string): Promise; getTaskAssignment(taskId: string): Promise; clearTTL(taskId: string): Promise; listByStatus(statuses: TaskStatus[]): Promise; } export interface LongTermStore { /** True only for split-tier stores with a verifiable archive barrier. */ readonly supportsHotColdRelease?: boolean; /** True only when deadline claims and terminal projection are durable. */ readonly supportsDurableTtl?: boolean; saveTask(task: Task): Promise; /** Atomically claims a durable task identity. Returns false when it already exists. */ createTaskIfAbsent?(task: Task): Promise; /** Claims an explicit task identity until its hot copy is created. */ claimTaskCreation?(task: Task, creationToken: string, claimTtlMs: number): Promise; /** Marks a claimed identity as fully created. */ completeTaskCreation?(taskId: string, creationToken: string): Promise; /** Removes only the pristine identity owned by this creation token. */ abortTaskCreation?(taskId: string, creationToken: string): Promise; getTask(taskId: string): Promise; saveEvent(event: TaskEvent): Promise; /** Optional series-aware durable write for latest-mode series. */ replaceLastSeriesEvent?(taskId: string, seriesId: string, event: TaskEvent): Promise; /** Optional series-aware durable write for accumulate-mode series. Returns the accumulated event. */ accumulateSeries?(taskId: string, seriesId: string, event: TaskEvent, field: string): Promise; getEvents(taskId: string, opts?: EventQueryOptions): Promise; /** * True when short-term archive restore writes the same durable storage this * long-term store reads from. The engine still runs long-term preflight, but * skips a duplicate long-term final restore. */ sharesTaskArchiveRestoreStorage?: boolean; /** Validates deterministic archive restore conflicts before mutation; engine calls this before multi-store restore. */ validateTaskArchiveRestore?(data: TaskArchiveRestoreData, options?: TaskArchiveImportOptions): Promise; /** Stores with native archive restore should implement this; engine import checks availability before use. */ restoreTaskArchive?(data: TaskArchiveRestoreData, options?: TaskArchiveImportOptions): Promise<{ overwritten: boolean; }>; persistStorageReleaseRequest?(request: StorageReleaseRequest): Promise; clearStorageReleaseRequest?(request: StorageReleaseRequest): Promise; listStorageReleaseRequests?(limit: number): Promise; getTaskStorageMetadata?(taskId: string): Promise; compareAndSetTaskStorageMetadata?(update: TaskStorageMetadataCas): Promise; beginArchive?(generation: ArchiveGeneration): Promise; archiveBatch?(taskId: string, generation: string, batch: ArchiveBatch): Promise; finalizeArchive?(taskId: string, generation: string, task: Task, seriesLatest: DurableSeriesState[]): Promise; getArchiveWatermark?(taskId: string): Promise; getLastEventIndex?(taskId: string): Promise; getRecentEvents?(taskId: string, limit: number): Promise; getDurableSeriesState?(taskId: string): Promise; claimOverdueTasks?(limit: number, claimTtlMs: number): Promise; terminalizeTtlClaim?(claim: TtlClaim, task: Task, event: TaskEvent, assignment: WorkerAssignment | null): Promise; claimTerminalProjections?(limit: number, claimToken: string, claimTtlMs: number): Promise; completeTerminalProjection?(projection: TerminalProjection): Promise; saveDurableAssignment?(assignment: WorkerAssignment): Promise; deleteDurableAssignment?(taskId: string, assignmentId?: string): Promise; saveWorkerEvent(event: WorkerAuditEvent): Promise; getWorkerEvents(workerId: string, opts?: EventQueryOptions): Promise; } export interface TaskFilter { status?: TaskStatus[]; types?: string[]; tags?: TagMatcher; assignMode?: AssignMode[]; excludeTaskIds?: string[]; limit?: number; } export interface WorkerFilter { status?: WorkerStatus[]; connectionMode?: ('pull' | 'websocket')[]; } export interface ErrorContext { operation: string; taskId?: string; } export interface TaskcastHooks { onTaskFailed?(task: Task, error: TaskError): void; onTaskTimeout?(task: Task): void; onUnhandledError?(err: unknown, context: ErrorContext): void; onEventDropped?(event: TaskEvent, reason: string): void; onWebhookFailed?(config: WebhookConfig, err: unknown): void; onSSEConnect?(taskId: string, clientId: string): void; onSSEDisconnect?(taskId: string, clientId: string, duration: number): void; onTaskCreated?(task: Task): void; onTaskTransitioned?(task: Task, from: TaskStatus, to: TaskStatus): void; onWorkerConnected?(worker: Worker): void; onWorkerDisconnected?(worker: Worker, reason: string): void; onTaskAssigned?(task: Task, worker: Worker): void; onTaskDeclined?(task: Task, worker: Worker, blacklisted: boolean): void; } //# sourceMappingURL=types.d.ts.map