/** * Canonical, producer-facing description of what a play is doing. * * Producers report facts. They do not choose log copy, UI placement, polling * intervals, or whether a run is "healthy". The activity router and projector * own those decisions centrally. */ export type PlayActivityTarget = | { kind: 'provider'; provider: string; operation: string; label?: string; } | { kind: 'dataset'; tableNamespace: string; operation: 'materializing' | 'persisting' | 'indexing' | 'reading'; label?: string; } | { kind: 'lease'; resourceKind: 'run' | 'dataset' | 'row' | 'runner'; resourceLabel: string; } | { kind: 'runner'; backend: string; label?: string; } | { kind: 'step'; stepId: string; label?: string; } | { kind: 'capacity'; pool: 'worker' | 'provider' | 'dataset' | 'runner'; label?: string; }; export type PlayActivityProgress = { completed?: number; total?: number; failed?: number; message?: string; }; export type PlayActivityState = | { kind: 'active'; progress?: PlayActivityProgress; } | { kind: 'queued'; reason: 'capacity'; } | { kind: 'waiting'; reason: | { kind: 'external_event'; provider?: string; eventKey?: string; remaining?: number; deadlineAt?: number; } | { kind: 'dataset'; tableNamespace: string; requiredPhase: 'available' | 'persisted' | 'indexed'; } | { kind: 'lease'; resourceKind: 'run' | 'dataset' | 'row' | 'runner'; resourceLabel: string; expiresAt?: number; } | { /** Mixed-version fallback only; new producers must choose a typed reason. */ kind: 'unknown'; }; } | { kind: 'retrying'; reason: 'rate_limit' | 'provider_error' | 'transport_error'; retryAt: number; attempt?: number; } | { kind: 'scheduled'; reason: 'sleep'; resumeAt: number; } | { kind: 'completed'; } | { kind: 'failed'; code?: string; }; export type PlayActivityObservation = { schemaVersion: 1; /** Stable within one run, for example `tool:discover_ctos` or `dataset:leads`. */ activityId: string; /** Optional owning play step. */ stepId?: string; target: PlayActivityTarget; state: PlayActivityState; /** Producer clock. The runtime envelope supplies run identity and ordering. */ observedAt: number; }; export type PlayActivityEvent = { type: 'activity.observed'; observation: PlayActivityObservation; }; function isRecord(value: unknown): value is Record { return Boolean(value && typeof value === 'object' && !Array.isArray(value)); } export function isPlayActivityObservation( value: unknown, ): value is PlayActivityObservation { if ( !isRecord(value) || value.schemaVersion !== 1 || typeof value.activityId !== 'string' || !value.activityId.trim() || typeof value.observedAt !== 'number' || !Number.isFinite(value.observedAt) || !isRecord(value.target) || !isRecord(value.state) ) { return false; } const targetKinds = new Set([ 'provider', 'dataset', 'lease', 'runner', 'step', 'capacity', ]); const stateKinds = new Set([ 'active', 'queued', 'waiting', 'retrying', 'scheduled', 'completed', 'failed', ]); return ( typeof value.target.kind === 'string' && targetKinds.has(value.target.kind) && typeof value.state.kind === 'string' && stateKinds.has(value.state.kind) ); } export type PlayRunActivityProjection = { observation: PlayActivityObservation; summary: string; visibility: PlayActivityVisibility; /** False only when old payloads lacked enough typed facts to classify. */ classified: boolean; }; /** * Persistence lanes are a routing concern, not a producer concern: * * - durable: append a meaningful transition to the run ledger; * - snapshot: replace current progress for this activity; * - liveness: update last-seen evidence without growing the event ledger; * - drop: stale or duplicate observation. */ export type PlayActivityRoutingLane = | 'durable' | 'snapshot' | 'liveness' | 'drop'; function stableJson(value: unknown): string { return JSON.stringify(value); } function targetIdentity(target: PlayActivityTarget): string { return stableJson(target); } export function routePlayActivityObservation(input: { previous?: PlayActivityObservation | null; next: PlayActivityObservation; }): PlayActivityRoutingLane { const { previous, next } = input; if (!previous) return 'durable'; if (next.observedAt <= previous.observedAt) return 'drop'; if (targetIdentity(previous.target) !== targetIdentity(next.target)) { return 'durable'; } if (previous.state.kind !== next.state.kind) return 'durable'; if (next.state.kind !== 'active') { return stableJson(previous.state) === stableJson(next.state) ? 'liveness' : 'durable'; } return stableJson(previous.state) === stableJson(next.state) ? 'liveness' : 'snapshot'; } export type PlayActivityVisibility = | 'hidden' | 'timeline' | 'headline' | 'action'; /** * Central presentation policy. Short, self-healing retries do not become * alarming UI/log messages; blockers and failures do. */ export function derivePlayActivityVisibility(input: { observation: PlayActivityObservation; now: number; }): PlayActivityVisibility { const { state } = input.observation; if (state.kind === 'failed') return 'action'; if (state.kind === 'waiting') { return state.reason.kind === 'external_event' ? 'action' : 'headline'; } if (state.kind === 'queued' || state.kind === 'scheduled') return 'headline'; if (state.kind === 'retrying') { return state.retryAt - input.now <= 30_000 ? 'hidden' : 'timeline'; } if (state.kind === 'active') return 'headline'; return 'timeline'; } function targetLabel(target: PlayActivityTarget): string { switch (target.kind) { case 'provider': return target.label ?? `${target.provider} ${target.operation}`; case 'dataset': return target.label ?? `dataset ${target.tableNamespace}`; case 'lease': return `${target.resourceKind} ${target.resourceLabel}`; case 'runner': return target.label ?? `${target.backend} runner`; case 'step': return target.label ?? target.stepId; case 'capacity': return target.label ?? `${target.pool} capacity`; } } export function summarizePlayActivity( observation: PlayActivityObservation, ): string { const label = targetLabel(observation.target); const { state } = observation; switch (state.kind) { case 'active': { const progress = state.progress; const counts = typeof progress?.completed === 'number' && typeof progress.total === 'number' ? ` (${progress.completed}/${progress.total})` : ''; return `Active: ${label}${counts}.`; } case 'queued': return `Queued: waiting for ${label}.`; case 'waiting': if (state.reason.kind === 'external_event') { const provider = state.reason.provider ? ` from ${state.reason.provider}` : ''; return `Blocked: waiting for an external event${provider}.`; } if (state.reason.kind === 'dataset') { return `Blocked: waiting for dataset ${state.reason.tableNamespace} to become ${state.reason.requiredPhase}.`; } if (state.reason.kind === 'lease') { return `Contended: waiting for ${state.reason.resourceKind} lease ${state.reason.resourceLabel}.`; } return `Paused: ${label}; the legacy producer did not report a typed reason.`; case 'retrying': return `Delayed: ${label} will retry automatically.`; case 'scheduled': return `Scheduled: ${label} will resume automatically.`; case 'completed': return `Completed: ${label}.`; case 'failed': return `Failed: ${label}.`; } } export type PlayActivityReporter = { observe( input: Omit & { observedAt?: number; }, ): void; }; export function createPlayActivityReporter(input: { emit: (event: PlayActivityEvent) => void; now?: () => number; }): PlayActivityReporter { return { observe(observation) { input.emit({ type: 'activity.observed', observation: { ...observation, schemaVersion: 1, observedAt: observation.observedAt ?? input.now?.() ?? Date.now(), }, }); }, }; } type ProjectionStep = { nodeId: string; status?: string; label?: string; artifactTableNamespace?: string | null; updatedAt?: number | null; progress?: PlayActivityProgress | null; }; type ProjectionDataset = { datasetId?: string; tableNamespace: string; phase?: string; persistedRows?: number; succeededRows?: number; failedRows?: number; complete?: boolean; updatedAt?: number; }; function isTerminalActivityState(state: PlayActivityState): boolean { return state.kind === 'completed' || state.kind === 'failed'; } function projection( observation: PlayActivityObservation, now: number, classified = true, ): PlayRunActivityProjection { return { observation, summary: summarizePlayActivity(observation), visibility: derivePlayActivityVisibility({ observation, now }), classified, }; } /** * Read-side compatibility projector. Explicit typed observations win; existing * suspension, dataset, and step facts are adapted without requiring every * legacy producer to migrate in one deploy. */ export function projectPlayRunActivity(input: { runId: string; playName?: string | null; status: string; updatedAt?: number | null; waitKind?: string | null; waitUntil?: number | null; eventKey?: string | null; runtimeBackend?: string | null; activeNodeId?: string | null; nodeStates?: ProjectionStep[] | null; datasets?: ProjectionDataset[] | null; explicit?: PlayActivityObservation[] | null; now?: number; }): PlayRunActivityProjection | null { const now = input.now ?? Date.now(); const observedAt = input.updatedAt ?? now; const normalizedStatus = input.status.trim().toLowerCase(); if ( normalizedStatus === 'completed' || normalizedStatus === 'failed' || normalizedStatus === 'cancelled' || normalizedStatus === 'terminated' || normalizedStatus === 'timed_out' ) { return null; } const explicit = [...(input.explicit ?? [])] .filter((candidate) => !isTerminalActivityState(candidate.state)) .sort((left, right) => right.observedAt - left.observedAt)[0]; if (explicit) return projection(explicit, now); if (input.waitKind === 'detached_runner') { return projection( { schemaVersion: 1, activityId: 'runner:managed', target: { kind: 'runner', backend: input.runtimeBackend ?? 'managed', label: 'managed play runner', }, state: { kind: 'active' }, observedAt, }, now, ); } if (input.waitKind === 'sleep') { return projection( { schemaVersion: 1, activityId: 'run:sleep', target: { kind: 'step', stepId: input.activeNodeId ?? 'run', label: input.playName ?? 'play', }, state: { kind: 'scheduled', reason: 'sleep', resumeAt: input.waitUntil ?? observedAt, }, observedAt, }, now, ); } if ( input.waitKind === 'integration_event' || input.waitKind === 'integration_event_batch' ) { return projection( { schemaVersion: 1, activityId: `event:${input.activeNodeId ?? 'run'}`, stepId: input.activeNodeId ?? undefined, target: { kind: 'step', stepId: input.activeNodeId ?? 'run', label: input.playName ?? 'play', }, state: { kind: 'waiting', reason: { kind: 'external_event', eventKey: input.eventKey ?? undefined, deadlineAt: input.waitUntil ?? undefined, }, }, observedAt, }, now, ); } // Mixed-version readers may see the durable WAITING status before the // scheduler's typed wait projection is available. Do not let a stale // running step overwrite that stronger lifecycle fact. if (normalizedStatus === 'waiting') { return projection( { schemaVersion: 1, activityId: 'run:legacy-wait', target: { kind: 'step', stepId: input.activeNodeId ?? 'run', label: input.playName ?? 'play', }, state: { kind: 'waiting', reason: { kind: 'unknown' } }, observedAt, }, now, false, ); } const pendingDataset = [...(input.datasets ?? [])] .filter( (dataset) => dataset.phase === 'registered' || (dataset.complete !== true && dataset.phase !== 'failed'), ) .sort((left, right) => (right.updatedAt ?? 0) - (left.updatedAt ?? 0))[0]; if (pendingDataset) { const activeDatasetStep = input.nodeStates?.find( (candidate) => candidate.nodeId === input.activeNodeId && candidate.status === 'running' && candidate.artifactTableNamespace === pendingDataset.tableNamespace, ); const datasetStep = activeDatasetStep ?? input.nodeStates?.find( (candidate) => candidate.artifactTableNamespace === pendingDataset.tableNamespace, ); const datasetProgress = datasetStep?.progress ?? null; return projection( { schemaVersion: 1, activityId: `dataset:${pendingDataset.datasetId ?? pendingDataset.tableNamespace}`, target: { kind: 'dataset', tableNamespace: pendingDataset.tableNamespace, operation: 'materializing', }, state: { kind: 'active', progress: { completed: datasetProgress?.completed ?? pendingDataset.persistedRows, total: datasetProgress?.total ?? (typeof pendingDataset.succeededRows === 'number' || typeof pendingDataset.failedRows === 'number' ? (pendingDataset.succeededRows ?? 0) + (pendingDataset.failedRows ?? 0) : undefined), failed: datasetProgress?.failed ?? pendingDataset.failedRows, }, }, observedAt: datasetStep?.updatedAt ?? pendingDataset.updatedAt ?? observedAt, }, now, ); } const activeStep = input.nodeStates?.find( (candidate) => candidate.nodeId === input.activeNodeId, ) ?? input.nodeStates?.find((candidate) => candidate.status === 'running') ?? null; if (activeStep) { const toolId = activeStep.nodeId.startsWith('tool:') ? activeStep.nodeId.split(':').at(-1)?.trim() || null : null; const provider = toolId?.split(/[._]/)[0]?.trim() || null; return projection( { schemaVersion: 1, activityId: `step:${activeStep.nodeId}`, stepId: activeStep.nodeId, target: toolId && provider ? { kind: 'provider', provider, operation: toolId, label: activeStep.label ?? toolId, } : { kind: 'step', stepId: activeStep.nodeId, label: activeStep.label, }, state: { kind: 'active', progress: activeStep.progress ?? undefined, }, observedAt: activeStep.updatedAt ?? observedAt, }, now, ); } if (normalizedStatus === 'queued') { return projection( { schemaVersion: 1, activityId: 'capacity:worker', target: { kind: 'capacity', pool: 'worker' }, state: { kind: 'queued', reason: 'capacity' }, observedAt, }, now, ); } const fallback: PlayActivityObservation = { schemaVersion: 1, activityId: 'run:legacy', target: { kind: 'step', stepId: input.activeNodeId ?? 'run', label: input.playName ?? 'play', }, state: normalizedStatus === 'waiting' ? { kind: 'waiting', reason: { kind: 'unknown' } } : { kind: 'active' }, observedAt, }; return projection(fallback, now, normalizedStatus !== 'waiting'); }