From 48921727ab7ce92698f22d57da083a863215323f Mon Sep 17 00:00:00 2001 From: Val Alexander Date: Sun, 27 Sep 2026 05:29:50 -0500 Subject: [PATCH] feat(coven): page one automation's occurrence history Add automations.occurrenceHistory(automationId, { limit, cursor }) over coven.automations.occurrence.history.v1, returning an sdk-core Page, and iterateOccurrenceHistory() on iteratePages. Cursors are validated before transport I/O; pages that cross automations, break newest-first order, repeat rows, misreport hasMore/next or fail to echo the requested cursor are refused. History responses get their own 256 KiB cap. Co-Authored-By: Claude Opus 5.5 (1M context) --- .changeset/coven-occurrence-history.md | 11 ++ api-baselines/coven.d.ts | 36 ++++- packages/coven/README.md | 54 ++++++- packages/coven/src/automations-definitions.ts | 33 ++++- packages/coven/src/automations-occurrences.ts | 71 ++++++++++ packages/coven/src/automations-socket.ts | 4 +- packages/coven/src/automations.ts | 46 +++++- packages/coven/src/index.ts | 2 + tests/coven-automations-platforms.spec.ts | 27 ++++ tests/coven-automations.spec.ts | 134 +++++++++++++++++- 10 files changed, 401 insertions(+), 17 deletions(-) create mode 100644 .changeset/coven-occurrence-history.md diff --git a/.changeset/coven-occurrence-history.md b/.changeset/coven-occurrence-history.md new file mode 100644 index 00000000..b0e9358f --- /dev/null +++ b/.changeset/coven-occurrence-history.md @@ -0,0 +1,11 @@ +--- +"@opencoven/coven-client": minor +--- + +Add `automations.occurrenceHistory(automationId, { limit, cursor })` over the +capability-gated `coven.automations.occurrence.history.v1` producer action: one +automation's occurrences in every state, newest first, as an sdk-core `Page` +with an opaque keyset cursor. Add `iterateOccurrenceHistory()` on sdk-core +`iteratePages`. The SDK validates the cursor before transport I/O, refuses pages +that cross automations, break ordering or do not echo the requested cursor, and +caps history responses at 256 KiB. Requires OpenCoven/coven#1157. diff --git a/api-baselines/coven.d.ts b/api-baselines/coven.d.ts index 5fe1efe8..25a201d2 100644 --- a/api-baselines/coven.d.ts +++ b/api-baselines/coven.d.ts @@ -1,6 +1,6 @@ // Entrypoint: . // Declaration: dist/index.d.ts -import { OperationContext, OperationDefaults, OperationOptions, DiscoveryEndpoint, NormalizedError } from '@opencoven/sdk-core'; +import { OperationContext, OperationDefaults, OperationOptions, BoundedPageOptions, DiscoveryEndpoint, NormalizedError } from '@opencoven/sdk-core'; /** Hand-authored read projection of Coven 4e35dd4c99013159fcee4c1ab2f183accdf7a5f8, not generated types. */ interface CovenAutomationReceiptDigest { @@ -358,6 +358,26 @@ interface CovenAutomationOccurrenceDetail extends CovenAutomationOccurrence { interface CovenAutomationOccurrencesResult { readonly occurrences: readonly CovenAutomationOccurrence[]; } +interface CovenAutomationOccurrenceHistoryOptions { + /** 1–100, default 20. */ + readonly limit?: number; + /** An opaque `cursor.next` from an earlier page of the same automation. */ + readonly cursor?: string; +} +/** + * One page of an automation's occurrences in every state, newest first by + * `scheduledFor` then `id`. Shaped as an sdk-core `Page`. + */ +interface CovenAutomationOccurrenceHistoryPage { + readonly automationId: string; + readonly data: readonly CovenAutomationOccurrence[]; + readonly cursor: { + readonly hasMore: boolean; + /** The cursor this page was read from; absent on the first page. */ + readonly current?: string; + readonly next?: string; + }; +} interface CovenAutomationOccurrenceResult { readonly occurrence: CovenAutomationOccurrenceDetail | null; } @@ -446,6 +466,11 @@ type CovenAutomationDefinitionReadRequest = { readonly view: CovenAutomationOccurrenceView; readonly limit: number; readonly automationId?: string; +} | { + readonly action: 'coven.automations.occurrence.history.v1'; + readonly automationId: string; + readonly limit: number; + readonly cursor?: string; } | { readonly action: 'coven.automations.occurrence.get.v1'; readonly id: string; @@ -513,6 +538,13 @@ declare class CovenAutomationsClient { health(id: string, options?: OperationOptions): Promise; runs(id: string, query?: CovenAutomationRunsOptions, options?: OperationOptions): Promise; occurrences(query: CovenAutomationOccurrencesOptions, options?: OperationOptions): Promise; + /** One page of one automation's occurrences in every state, newest first. */ + occurrenceHistory(automationId: string, query?: CovenAutomationOccurrenceHistoryOptions, options?: OperationOptions): Promise; + /** + * Every occurrence of one automation, newest first, across pages. Requires + * `maxPages` or a caller-owned `signal`; `timeoutMs` bounds the whole walk. + */ + iterateOccurrenceHistory(automationId: string, options: BoundedPageOptions): AsyncGenerator; /** One run and its attempts by run id, read from one producer snapshot. */ getRun(runId: string, options?: OperationOptions): Promise; getOccurrence(id: string, options?: OperationOptions): Promise; @@ -1025,4 +1057,4 @@ type CovenAutomationEventReduction = { /** Bounded supplied-batch reference projection; no persistence, transport or execution authority. */ declare function reduceAutomationEvents(events: unknown): CovenAutomationEventReduction; -export { COVEN_DAEMON_PROTOCOL, COVEN_SESSION_POLICY_CONTRACT, COVEN_SESSION_POLICY_PROFILE, type CovenAutomationAttempt, type CovenAutomationCapabilities, type CovenAutomationCapabilityProfile, type CovenAutomationDefinition, type CovenAutomationDefinitionDigestResult, type CovenAutomationDefinitionDocument, type CovenAutomationDefinitionList, type CovenAutomationDefinitionReadRequest, type CovenAutomationEvent, type CovenAutomationEventIntegrity, type CovenAutomationEventPage, type CovenAutomationEventReduction, type CovenAutomationEventStream, type CovenAutomationEventsOptions, type CovenAutomationEventsRequest, type CovenAutomationHealth, type CovenAutomationHealthResult, type CovenAutomationListOptions, type CovenAutomationOccurrence, type CovenAutomationOccurrenceDetail, type CovenAutomationOccurrenceResult, type CovenAutomationOccurrenceRun, type CovenAutomationOccurrenceView, type CovenAutomationOccurrencesOptions, type CovenAutomationOccurrencesResult, type CovenAutomationProjectionJson, type CovenAutomationReceipt, type CovenAutomationReceiptDigest, type CovenAutomationReceiptReadVerification, type CovenAutomationReceiptResult, type CovenAutomationReceiptTrustContext, type CovenAutomationReceiptVerification, type CovenAutomationReceiptVerificationCheck, type CovenAutomationReceiptVerificationReason, type CovenAutomationRoutine, type CovenAutomationRun, type CovenAutomationRunCancellation, type CovenAutomationRunResult, type CovenAutomationRunsOptions, type CovenAutomationRunsResult, type CovenAutomationVariant, CovenAutomationsClient, type CovenAutomationsClientOptions, type CovenAutomationsTransport, type CovenAutomationsUnixTransportOptions, type CovenAutomationsWindowsTransportOptions, CovenClient, CovenClientError, type CovenClientOptions, type CovenConnectedSocket, type CovenDaemonFailure, CovenDaemonResponseError, type CovenDiscoveredClientOptions, type CovenDiscoveredEndpoint, type CovenDiscoveredUnixClientOptions, type CovenDiscoveredUnixTransportOptions, type CovenDiscoveredWindowsClientOptions, type CovenDiscoveredWindowsTransportOptions, type CovenDiscoveryDependencies, type CovenDiscoveryFileIdentity, type CovenDiscoverySource, type CovenEndpointFreshness, type CovenEndpointOwner, type CovenExecFile, type CovenExecFileError, type CovenExecFileOptions, type CovenExecutableResolver, type CovenHealth, type CovenHealthResponse, type CovenHealthTransportLimits, type CovenIpcDiagnostics, CovenIpcError, type CovenIpcErrorCode, type CovenMetadataFileHandle, type CovenRestrictedLaunchRequest, CovenSessionPolicyClient, type CovenSessionPolicyClientOptions, type CovenSessionPolicyDelivery, type CovenSessionPolicyDiscovery, CovenSessionPolicyError, type CovenSessionPolicyErrorCode, type CovenSessionPolicyRefusal, type CovenSessionPolicyTransport, type CovenSessionPolicyTransportRequest, type CovenSessionPolicyTransportResponse, type CovenSessionPolicyUnixTransportOptions, type CovenSocket, type CovenSocketConnector, type CovenTransport, type CovenTransportSecurityProvider, type CovenUnixFileIdentity, type CovenUnixPeerIdentity, type CovenUnixPeerIdentityAdapter, type CovenUnixTransportDependencies, type CovenUnixTransportOptions, type CovenUnixTransportSecurityProvider, type CovenWindowsFileTrustValidator, type CovenWindowsPipeIdentity, type CovenWindowsPipeOwnershipAdapter, type CovenWindowsTransportDependencies, type CovenWindowsTransportOptions, type CovenWindowsTransportSecurityProvider, type DiscoverCovenEndpointOptions, computeDefinitionDigest, createCovenAutomationsClient, createCovenAutomationsUnixTransport, createCovenAutomationsWindowsTransport, createCovenClient, createCovenSessionPolicyClient, createCovenSessionPolicyUnixTransport, createCovenUnixTransport, createCovenWindowsTransport, createDiscoveredCovenClient, discoverCovenEndpoint, isCovenClientError, isCovenDaemonResponseError, isCovenIpcError, isCovenSessionPolicyError, normalizeCovenError, reduceAutomationEvents, verifyEventIntegrity, verifyReceipt }; +export { COVEN_DAEMON_PROTOCOL, COVEN_SESSION_POLICY_CONTRACT, COVEN_SESSION_POLICY_PROFILE, type CovenAutomationAttempt, type CovenAutomationCapabilities, type CovenAutomationCapabilityProfile, type CovenAutomationDefinition, type CovenAutomationDefinitionDigestResult, type CovenAutomationDefinitionDocument, type CovenAutomationDefinitionList, type CovenAutomationDefinitionReadRequest, type CovenAutomationEvent, type CovenAutomationEventIntegrity, type CovenAutomationEventPage, type CovenAutomationEventReduction, type CovenAutomationEventStream, type CovenAutomationEventsOptions, type CovenAutomationEventsRequest, type CovenAutomationHealth, type CovenAutomationHealthResult, type CovenAutomationListOptions, type CovenAutomationOccurrence, type CovenAutomationOccurrenceDetail, type CovenAutomationOccurrenceHistoryOptions, type CovenAutomationOccurrenceHistoryPage, type CovenAutomationOccurrenceResult, type CovenAutomationOccurrenceRun, type CovenAutomationOccurrenceView, type CovenAutomationOccurrencesOptions, type CovenAutomationOccurrencesResult, type CovenAutomationProjectionJson, type CovenAutomationReceipt, type CovenAutomationReceiptDigest, type CovenAutomationReceiptReadVerification, type CovenAutomationReceiptResult, type CovenAutomationReceiptTrustContext, type CovenAutomationReceiptVerification, type CovenAutomationReceiptVerificationCheck, type CovenAutomationReceiptVerificationReason, type CovenAutomationRoutine, type CovenAutomationRun, type CovenAutomationRunCancellation, type CovenAutomationRunResult, type CovenAutomationRunsOptions, type CovenAutomationRunsResult, type CovenAutomationVariant, CovenAutomationsClient, type CovenAutomationsClientOptions, type CovenAutomationsTransport, type CovenAutomationsUnixTransportOptions, type CovenAutomationsWindowsTransportOptions, CovenClient, CovenClientError, type CovenClientOptions, type CovenConnectedSocket, type CovenDaemonFailure, CovenDaemonResponseError, type CovenDiscoveredClientOptions, type CovenDiscoveredEndpoint, type CovenDiscoveredUnixClientOptions, type CovenDiscoveredUnixTransportOptions, type CovenDiscoveredWindowsClientOptions, type CovenDiscoveredWindowsTransportOptions, type CovenDiscoveryDependencies, type CovenDiscoveryFileIdentity, type CovenDiscoverySource, type CovenEndpointFreshness, type CovenEndpointOwner, type CovenExecFile, type CovenExecFileError, type CovenExecFileOptions, type CovenExecutableResolver, type CovenHealth, type CovenHealthResponse, type CovenHealthTransportLimits, type CovenIpcDiagnostics, CovenIpcError, type CovenIpcErrorCode, type CovenMetadataFileHandle, type CovenRestrictedLaunchRequest, CovenSessionPolicyClient, type CovenSessionPolicyClientOptions, type CovenSessionPolicyDelivery, type CovenSessionPolicyDiscovery, CovenSessionPolicyError, type CovenSessionPolicyErrorCode, type CovenSessionPolicyRefusal, type CovenSessionPolicyTransport, type CovenSessionPolicyTransportRequest, type CovenSessionPolicyTransportResponse, type CovenSessionPolicyUnixTransportOptions, type CovenSocket, type CovenSocketConnector, type CovenTransport, type CovenTransportSecurityProvider, type CovenUnixFileIdentity, type CovenUnixPeerIdentity, type CovenUnixPeerIdentityAdapter, type CovenUnixTransportDependencies, type CovenUnixTransportOptions, type CovenUnixTransportSecurityProvider, type CovenWindowsFileTrustValidator, type CovenWindowsPipeIdentity, type CovenWindowsPipeOwnershipAdapter, type CovenWindowsTransportDependencies, type CovenWindowsTransportOptions, type CovenWindowsTransportSecurityProvider, type DiscoverCovenEndpointOptions, computeDefinitionDigest, createCovenAutomationsClient, createCovenAutomationsUnixTransport, createCovenAutomationsWindowsTransport, createCovenClient, createCovenSessionPolicyClient, createCovenSessionPolicyUnixTransport, createCovenUnixTransport, createCovenWindowsTransport, createDiscoveredCovenClient, discoverCovenEndpoint, isCovenClientError, isCovenDaemonResponseError, isCovenIpcError, isCovenSessionPolicyError, normalizeCovenError, reduceAutomationEvents, verifyEventIntegrity, verifyReceipt }; diff --git a/packages/coven/README.md b/packages/coven/README.md index 89148036..a33de7a2 100644 --- a/packages/coven/README.md +++ b/packages/coven/README.md @@ -85,7 +85,8 @@ Discovery runs once. Health and Automations reuse that exact endpoint, security provider, and platform dependencies. Construction sends no health, capability, or action request. Each later request authenticates its own connection. The `unix`/`windows` health response-limit options still apply only to health; -Automations keeps its fixed 16 KiB body limit. +Automations keeps its fixed 16 KiB body limit, except for the event (1 MiB) and +occurrence-history (256 KiB) reads described below. `client.automations` is `CovenAutomationsClient | undefined`. `client.requireAutomations()` returns that same namespace or throws @@ -133,6 +134,7 @@ does not poll or activate Automations. | `capabilities()` | Authenticated GET | Authenticated GET | | `list()`, `get()`, `health()` | Allowlisted reads | Same actions and decoders | | `runs()`, `occurrences()`, `getOccurrence()`, `getRun()` | Bounded diagnostic reads | Same actions and decoders | +| `occurrenceHistory()`, `iterateOccurrenceHistory()` | Keyset-paged diagnostic reads | Same actions and decoders | | `getReceipt()` | Public/operational receipt result | Same result and privacy checks | | `events()`, `subscribe()` | Bounded domain event pages | Same pages and cancellation | | Normal/discovered client and `sdk.coven` | Explicit opt-in | Explicit opt-in | @@ -207,7 +209,7 @@ missing action names fail with `capability_unsupported` without posting an actio Custom capability-only transports remain compatible; reads without the optional `readDefinitions` hook fail with `unsupported_operation`. -The built-in Unix and Windows transports send only these eight allowlisted JSON actions to +The built-in Unix and Windows transports send only these nine allowlisted JSON actions to `POST /api/v1/actions`. It authenticates each connection, including the separate capability request, under one client deadline/cancellation scope. It cannot send mutations through this hook. IDs are trimmed as the producer does; the SDK @@ -233,7 +235,7 @@ with `event.kind: "automations.changed"` even for reads, and defines owns the compatibility routine fields. No GET definition routes, normative command-envelope adaptation, pagination, changefeed emission, or certification are inferred from the schemas in `spec/coven-automations/v1`. Existing artifact -pins are unchanged. Individual run reads, per-automation occurrence history, receipt authentication, +pins are unchanged. Receipt authentication, global feed subscriptions and authority-bearing phases remain separate #80 work. ### Routine health @@ -252,7 +254,7 @@ Failure/exhaustion counters are nonnegative safe integers and `maxAttempts` is Missing routines produce sanitized `action_rejected`, not an invented null result. Health is store-derived diagnostic data, not execution or receipt authority. Custom transports use the existing optional `readDefinitions` hook, whose -historical name now covers all eight explicitly allowlisted read actions. +historical name now covers all nine explicitly allowlisted read actions. Health source authority was independently read from Coven [`b3b2d043a4ee586ccbf25ef6aad21db8a1171a54`](https://github.com/OpenCoven/coven/tree/b3b2d043a4ee586ccbf25ef6aad21db8a1171a54): @@ -575,7 +577,8 @@ exact advertised `coven.automations.occurrence.list.v1` action. Without producer restricts the view to that automation inside its query, so `limit` bounds that automation's rows; the SDK refuses a filtered page that names any other automation. An empty or non-string `automationId` is rejected before any -transport I/O. There is still no cursor in this producer contract. +transport I/O. This action has no cursor; page one automation's full history +with `occurrenceHistory()` below. The required view is `due`, `eligible`, `claimed`, `running`, or `recovery_required`; limits are integers 1–100 (default 20). Unsupported query fields are rejected rather than silently suggesting filtering or pagination. @@ -592,6 +595,45 @@ cancellation projection. Both reads need a Coven producer that advertises them ([coven#1155](https://github.com/OpenCoven/coven/issues/1155)); an older producer yields `capability_unsupported`, never a fallback. +### Occurrence history + +```ts +const page = await automations.occurrenceHistory('morning', { limit: 50 }); +for await (const occurrence of automations.iterateOccurrenceHistory('morning', { + limit: 50, + maxPages: 10, +})) { + console.log(occurrence.scheduledFor, occurrence.state); +} +``` + +`occurrenceHistory(automationId, { limit?, cursor? }, operationOptions?)` calls +`coven.automations.occurrence.history.v1` and returns one sdk-core `Page`: +`{ automationId, data, cursor: { hasMore, current?, next? } }`. It covers every +state, terminal ones included, newest first by `scheduledFor` then `id`, and +reflects whatever the producer's history retention has kept. `limit` is 1–100 +(default 20). `cursor` is the opaque `next` of an earlier page for the same +automation; the SDK accepts only canonical unpadded base64url of at most 512 +characters and rejects anything else, and any other query field, before +transport I/O. The producer pages by keyset, so occurrences planned while a +caller pages never shift a later page. + +The SDK refuses a page that names another automation, holds more than `limit` +rows, is not strictly newest first, repeats a row, claims `hasMore` without a +fresh `next` or on a short page, carries `next` without `hasMore`, or does not +echo the requested cursor as `current`. History responses are capped at +**256 KiB**, enough for 100 records; a larger page fails closed, and a smaller +`limit` reads it. + +`iterateOccurrenceHistory(automationId, { limit?, cursor?, maxPages?, signal?, +timeoutMs?, observer? })` walks pages with sdk-core `iteratePages`: it needs +`maxPages` or a caller-owned `signal`, stops at `hasMore: false`, and refuses a +cursor that does not advance. Its `limit` defaults to sdk-core's 50, and +`timeoutMs` bounds the whole walk, while each page read keeps the client's +per-call timeout. Both need a Coven producer that advertises the action +([coven#1157](https://github.com/OpenCoven/coven/issues/1157)); an older +producer yields `capability_unsupported`. + `getOccurrence(id, operationOptions?)` calls `coven.automations.occurrence.get.v1` and returns `{ occurrence: null } for absence, or a diagnostic detail record. The detail contains up to **20 oldest @@ -690,7 +732,7 @@ events on reconnect. Event responses are capped at 1 MiB and 100 wire events, with the existing 16-level JSON depth limit. Oversized or malformed pages fail closed without yielding partial events. All other Automations and policy response limits -remain 16 KiB. A producer page larger than the event byte cap cannot be +remain 16 KiB, except occurrence history (256 KiB). A producer page larger than the event byte cap cannot be retried with a smaller subscription limit because the producer forbids that field. The SDK does not silently truncate it. diff --git a/packages/coven/src/automations-definitions.ts b/packages/coven/src/automations-definitions.ts index 3da9b43a..b9931bf3 100644 --- a/packages/coven/src/automations-definitions.ts +++ b/packages/coven/src/automations-definitions.ts @@ -7,8 +7,9 @@ import { import { decodeReceiptRead, receiptId, type CovenAutomationReceiptResult } from './automations-receipts.js'; import { runsPayload, type CovenAutomationRunsResult } from './automations-runs.js'; import { - occurrencePayload, occurrencesPayload, occurrenceView, runPayload, - type CovenAutomationOccurrenceResult, type CovenAutomationOccurrencesResult, type CovenAutomationOccurrenceView, + AUTOMATION_HISTORY_MAX_BYTES, historyCursor, occurrenceHistoryPayload, occurrencePayload, occurrencesPayload, + occurrenceView, runPayload, + type CovenAutomationOccurrenceHistoryPage, type CovenAutomationOccurrenceResult, type CovenAutomationOccurrencesResult, type CovenAutomationOccurrenceView, type CovenAutomationRunResult, } from './automations-occurrences.js'; @@ -91,6 +92,12 @@ export type CovenAutomationDefinitionReadRequest = readonly limit: number; readonly automationId?: string; } + | { + readonly action: 'coven.automations.occurrence.history.v1'; + readonly automationId: string; + readonly limit: number; + readonly cursor?: string; + } | { readonly action: 'coven.automations.occurrence.get.v1'; readonly id: string } | { readonly action: 'coven.automations.run.get.v1'; readonly id: string } | { readonly action: 'coven.automations.receipt.get.v1'; readonly id: string } @@ -148,6 +155,18 @@ export function definitionReadBytes(request: CovenAutomationDefinitionReadReques Buffer.byteLength(automationId) > 4_096 || !automationId.isWellFormed()) return invalid(); return Buffer.from(JSON.stringify({ action, view, limit, automationId: automationId.trim() })); } + if (action === 'coven.automations.occurrence.history.v1') { + const automationId = own('automationId'); + const limit = own('limit'); + const paged = Object.hasOwn(descriptors, 'cursor'); + const cursor = own('cursor'); + if (keys.length !== (paged ? 4 : 3) || !integer(limit, 1, 100) || (paged && !historyCursor(cursor)) || + typeof automationId !== 'string' || automationId.trim().length === 0 || + Buffer.byteLength(automationId) > 4_096 || !automationId.isWellFormed()) return invalid(); + return Buffer.from(JSON.stringify({ + action, automationId: automationId.trim(), limit, ...(paged ? { cursor } : {}), + })); + } if (keys.length !== (action === 'coven.automations.runs' ? 3 : 2)) return invalid(); if (action === 'coven.automations.definition.list.v1' && typeof own('includeTombstoned') === 'boolean') { return Buffer.from(JSON.stringify({ action, includeTombstoned: own('includeTombstoned') })); @@ -172,12 +191,13 @@ export function decodeDefinitionRead( request: CovenAutomationDefinitionReadRequest, operation: string, ): CovenAutomationDefinitionList | CovenAutomationDefinition | CovenAutomationHealthResult | CovenAutomationRunsResult | - CovenAutomationOccurrencesResult | CovenAutomationOccurrenceResult | CovenAutomationRunResult | CovenAutomationReceiptResult | - CovenAutomationEventPage { + CovenAutomationOccurrencesResult | CovenAutomationOccurrenceHistoryPage | CovenAutomationOccurrenceResult | + CovenAutomationRunResult | CovenAutomationReceiptResult | CovenAutomationEventPage { const invalid = (): never => definitionReadFailure('invalid_response', operation); let value: unknown; try { - value = parsePolicyJson(bytes, request.action === 'coven.automations.events.subscribe.v1' ? AUTOMATION_EVENTS_MAX_BYTES : 16_384); + value = parsePolicyJson(bytes, request.action === 'coven.automations.events.subscribe.v1' ? AUTOMATION_EVENTS_MAX_BYTES + : request.action === 'coven.automations.occurrence.history.v1' ? AUTOMATION_HISTORY_MAX_BYTES : 16_384); } catch { return invalid(); } @@ -199,6 +219,9 @@ export function decodeDefinitionRead( if (request.action === 'coven.automations.occurrence.list.v1') { return occurrencesPayload(payload, request.limit, request.automationId?.trim()) ?? invalid(); } + if (request.action === 'coven.automations.occurrence.history.v1') { + return occurrenceHistoryPayload(payload, request.automationId.trim(), request.limit, request.cursor) ?? invalid(); + } if (request.action === 'coven.automations.occurrence.get.v1') { return occurrencePayload(payload, request.id.trim()) ?? invalid(); } diff --git a/packages/coven/src/automations-occurrences.ts b/packages/coven/src/automations-occurrences.ts index c11a8a41..69cc4c21 100644 --- a/packages/coven/src/automations-occurrences.ts +++ b/packages/coven/src/automations-occurrences.ts @@ -1,3 +1,5 @@ +import { normalizePageOptions } from '@opencoven/sdk-core'; + import { integer, object } from './automations-read-validation.js'; import { runsPayload, type CovenAutomationRun } from './automations-runs.js'; @@ -51,6 +53,31 @@ export interface CovenAutomationOccurrencesResult { readonly occurrences: readonly CovenAutomationOccurrence[]; } +/** Largest history response accepted: room for 100 records, far above the 16 KiB read cap. */ +export const AUTOMATION_HISTORY_MAX_BYTES = 262_144; + +export interface CovenAutomationOccurrenceHistoryOptions { + /** 1–100, default 20. */ + readonly limit?: number; + /** An opaque `cursor.next` from an earlier page of the same automation. */ + readonly cursor?: string; +} + +/** + * One page of an automation's occurrences in every state, newest first by + * `scheduledFor` then `id`. Shaped as an sdk-core `Page`. + */ +export interface CovenAutomationOccurrenceHistoryPage { + readonly automationId: string; + readonly data: readonly CovenAutomationOccurrence[]; + readonly cursor: { + readonly hasMore: boolean; + /** The cursor this page was read from; absent on the first page. */ + readonly current?: string; + readonly next?: string; + }; +} + export interface CovenAutomationOccurrenceResult { readonly occurrence: CovenAutomationOccurrenceDetail | null; } @@ -86,6 +113,50 @@ export function occurrencesPayload( return { occurrences: value.occurrences }; } +/** A cursor sdk-core would accept: canonical unpadded base64url, at most 512 characters. */ +export function historyCursor(value: unknown): value is string { + if (typeof value !== 'string') return false; + try { + normalizePageOptions({ cursor: value }); + return true; + } catch { + return false; + } +} + +function newerThan(left: CovenAutomationOccurrence, right: CovenAutomationOccurrence): boolean { + // The producer orders by SQLite's byte comparison, not UTF-16 code units. + const order = Buffer.compare(Buffer.from(left.scheduledFor), Buffer.from(right.scheduledFor)); + return order > 0 || (order === 0 && Buffer.compare(Buffer.from(left.id), Buffer.from(right.id)) > 0); +} + +export function occurrenceHistoryPayload( + value: Record, + automationId: string, + limit: number, + cursor: string | undefined, +): CovenAutomationOccurrenceHistoryPage | undefined { + const page = value.automationId === automationId ? occurrencesPayload(value, limit, automationId) : undefined; + const position = value.cursor; + if (page === undefined || !object(position) || typeof position.hasMore !== 'boolean' || + Reflect.ownKeys(position).some((key) => key !== 'hasMore' && key !== 'current' && key !== 'next') || + (cursor === undefined ? Object.hasOwn(position, 'current') : position.current !== cursor)) return undefined; + const rows = page.occurrences; + if (position.hasMore + ? !historyCursor(position.next) || position.next === cursor || rows.length !== limit + : Object.hasOwn(position, 'next')) return undefined; + if (rows.some((row, index) => index > 0 && !newerThan(rows[index - 1]!, row))) return undefined; + return { + automationId, + data: rows, + cursor: { + hasMore: position.hasMore, + ...(cursor === undefined ? {} : { current: cursor }), + ...(position.hasMore ? { next: position.next as string } : {}), + }, + }; +} + function occurrenceRunFields(record: Record): boolean { return typeof record.occurrenceId === 'string' && record.occurrenceId.length > 0 && integer(record.automationRevision, 1) && diff --git a/packages/coven/src/automations-socket.ts b/packages/coven/src/automations-socket.ts index 787806e2..24a2f78d 100644 --- a/packages/coven/src/automations-socket.ts +++ b/packages/coven/src/automations-socket.ts @@ -3,6 +3,7 @@ import { runOperation, type OperationContext } from '@opencoven/sdk-core'; import type { CovenAutomationsTransport } from './automations.js'; import { definitionReadBytes } from './automations-definitions.js'; import { AUTOMATION_EVENTS_MAX_BYTES } from './automations-events.js'; +import { AUTOMATION_HISTORY_MAX_BYTES } from './automations-occurrences.js'; import { CovenClientError, normalizeCovenError } from './client-errors.js'; import { requestCovenPolicyOverSocket, @@ -49,7 +50,8 @@ export function createCovenAutomationsSocketTransport( `Content-Length: ${body.byteLength}\r\n\r\n`, ), body, - ]), context, 'automations.read', action === 'coven.automations.events.subscribe.v1' ? AUTOMATION_EVENTS_MAX_BYTES : 16_384); + ]), context, 'automations.read', action === 'coven.automations.events.subscribe.v1' ? AUTOMATION_EVENTS_MAX_BYTES + : action === 'coven.automations.occurrence.history.v1' ? AUTOMATION_HISTORY_MAX_BYTES : 16_384); }, }; } diff --git a/packages/coven/src/automations.ts b/packages/coven/src/automations.ts index f2c155a6..0be8eed9 100644 --- a/packages/coven/src/automations.ts +++ b/packages/coven/src/automations.ts @@ -1,5 +1,7 @@ import { + iteratePages, runOperation, + type BoundedPageOptions, type OperationContext, type OperationDefaults, type OperationOptions, @@ -18,7 +20,10 @@ import { eventsOptions, eventsRequest, subscribeEvents, type CovenAutomationEventPage, type CovenAutomationEventsOptions, } from './automations-events.js'; import { + historyCursor, occurrenceView, + type CovenAutomationOccurrence, + type CovenAutomationOccurrenceHistoryOptions, type CovenAutomationOccurrenceHistoryPage, type CovenAutomationOccurrencesOptions, type CovenAutomationOccurrencesResult, type CovenAutomationOccurrenceResult, type CovenAutomationRunResult, } from './automations-occurrences.js'; @@ -204,6 +209,42 @@ export class CovenAutomationsClient { }, options) as CovenAutomationOccurrencesResult; } + /** One page of one automation's occurrences in every state, newest first. */ + async occurrenceHistory( + automationId: string, + query: CovenAutomationOccurrenceHistoryOptions = {}, + options: OperationOptions = {}, + ): Promise { + if (!object(query) || (query.limit !== undefined && !integer(query.limit, 1, 100)) || + Reflect.ownKeys(query).some((key) => key !== 'limit' && key !== 'cursor') || + (Object.hasOwn(query, 'cursor') && !historyCursor(query.cursor))) { + return definitionReadFailure('invalid_options', 'automations.occurrenceHistory'); + } + return await this.#read({ + action: 'coven.automations.occurrence.history.v1', automationId, limit: query.limit ?? 20, + ...(Object.hasOwn(query, 'cursor') ? { cursor: query.cursor as string } : {}), + }, options) as CovenAutomationOccurrenceHistoryPage; + } + + /** + * Every occurrence of one automation, newest first, across pages. Requires + * `maxPages` or a caller-owned `signal`; `timeoutMs` bounds the whole walk. + */ + iterateOccurrenceHistory(automationId: string, options: BoundedPageOptions): AsyncGenerator { + const defaults = this.#options.operation; + const timeoutMs = object(options) ? options.timeoutMs ?? defaults?.timeoutMs : undefined; + const observer = object(options) ? options.observer ?? defaults?.observer : undefined; + return iteratePages( + ({ limit, cursor, signal }) => + this.occurrenceHistory(automationId, cursor === undefined ? { limit } : { limit, cursor }, { signal }), + object(options) ? { + ...options, + ...(timeoutMs === undefined ? {} : { timeoutMs }), + ...(observer === undefined ? {} : { observer }), + } : options, + ); + } + /** One run and its attempts by run id, read from one producer snapshot. */ async getRun(runId: string, options: OperationOptions = {}): Promise { return await this.#read({ action: 'coven.automations.run.get.v1', id: runId }, options) as CovenAutomationRunResult; @@ -234,11 +275,12 @@ export class CovenAutomationsClient { request: CovenAutomationDefinitionReadRequest, options: OperationOptions, ): Promise { + CovenAutomationOccurrencesResult | CovenAutomationOccurrenceHistoryPage | CovenAutomationOccurrenceResult | + CovenAutomationRunResult | CovenAutomationReceiptResult | CovenAutomationEventPage> { const operation = request.action === 'coven.automations.definition.list.v1' ? 'automations.list' : request.action === 'coven.automations.health' ? 'automations.health' : request.action === 'coven.automations.occurrence.list.v1' ? 'automations.occurrences' + : request.action === 'coven.automations.occurrence.history.v1' ? 'automations.occurrenceHistory' : request.action === 'coven.automations.occurrence.get.v1' ? 'automations.getOccurrence' : request.action === 'coven.automations.run.get.v1' ? 'automations.getRun' : request.action === 'coven.automations.receipt.get.v1' ? 'automations.getReceipt' diff --git a/packages/coven/src/index.ts b/packages/coven/src/index.ts index 2fd3de0b..3bde5236 100644 --- a/packages/coven/src/index.ts +++ b/packages/coven/src/index.ts @@ -33,6 +33,8 @@ export type { export type { CovenAutomationOccurrence, CovenAutomationOccurrenceDetail, + CovenAutomationOccurrenceHistoryOptions, + CovenAutomationOccurrenceHistoryPage, CovenAutomationOccurrenceResult, CovenAutomationRunResult, CovenAutomationOccurrenceRun, diff --git a/tests/coven-automations-platforms.spec.ts b/tests/coven-automations-platforms.spec.ts index 983e292f..3c446fbd 100644 --- a/tests/coven-automations-platforms.spec.ts +++ b/tests/coven-automations-platforms.spec.ts @@ -211,6 +211,33 @@ describe.each(['unix', 'windows'] as const)('%s Automations parity', (platform) expect(sockets.every((socket) => socket.destroyed)).toBe(true); }); + test.skipIf(platform === 'unix' && process.platform === 'win32')('reads a 100-row history page with a history-only 256 KiB bound', async () => { + const { client, configure, sockets } = setup(platform); + const action = 'coven.automations.occurrence.history.v1'; + const occurrences = Array.from({ length: 100 }, (_, index) => ({ + id: `occurrence-${String(999 - index).padStart(3, '0')}`, automationId: 'morning', automationRevision: 1, + definitionDigest: 'a'.repeat(64), scheduledFor: '2026-09-14T00:00:00.000Z', kind: 'schedule', state: 'succeeded', + leaseOwner: null, leaseExpiresAt: null, schedulerGeneration: 3, fenceGeneration: 1, failureReason: null, + createdAt: '2026-09-14T00:00:00.000Z', updatedAt: '2026-09-14T00:01:00.000Z', + })); + const payload = { automationId: 'morning', occurrences, cursor: { hasMore: false } }; + const body = Buffer.from(JSON.stringify({ + ok: true, accepted: true, action, status: 'completed', event: { kind: 'automations.changed', action, payload }, + })); + expect(body.length).toBeGreaterThan(16_384); + configure.mockImplementation((socket, index) => { + socket.response = index === 0 ? Buffer.from(JSON.stringify(advertisement([action]))) : body; + }); + expect((await client.occurrenceHistory('morning', { limit: 100 })).data).toEqual(occurrences); + expect(sockets[1]?.writes[0]?.split('\r\n\r\n')[1]).toBe(JSON.stringify({ action, automationId: 'morning', limit: 100 })); + configure.mockImplementation((socket, index) => { + if (index === 2) socket.response = Buffer.from(JSON.stringify(advertisement([action]))); + else socket.rawResponse = Buffer.from('HTTP/1.1 200 OK\r\nContent-Length: 262145\r\n\r\n'); + }); + await expect(client.occurrenceHistory('morning')).rejects.toMatchObject({ code: 'body_limit' }); + expect(sockets.every((socket) => socket.destroyed)).toBe(true); + }); + test.skipIf(platform === 'unix' && process.platform === 'win32')('return closes a pending subscription socket before any late response', async () => { const { client, configure, sockets } = setup(platform); configure.mockImplementation((socket, index) => { diff --git a/tests/coven-automations.spec.ts b/tests/coven-automations.spec.ts index 1736dbdb..3e0abc75 100644 --- a/tests/coven-automations.spec.ts +++ b/tests/coven-automations.spec.ts @@ -13,6 +13,7 @@ import { type CovenAutomationListOptions, type CovenAutomationRunsOptions, type CovenAutomationRunsResult, + type CovenAutomationOccurrenceHistoryPage, type CovenAutomationOccurrencesOptions, type CovenAutomationOccurrencesResult, type CovenAutomationOccurrenceResult, @@ -31,7 +32,7 @@ function advertisement() { adapter: 'coven-daemon', status: 'available', policy: 'allow', actions: ['coven.automations.definition.get.v1', 'coven.automations.definition.list.v1', 'coven.automations.health', 'coven.automations.runs', 'coven.automations.occurrence.list.v1', 'coven.automations.occurrence.get.v1', 'coven.automations.run.get.v1', - 'coven.automations.receipt.get.v1'], + 'coven.automations.receipt.get.v1', 'coven.automations.occurrence.history.v1'], variantNegotiation: { version: 1, contractProfile: 'coven.automations.v1', description: 'Variant negotiation', supported: { @@ -1078,3 +1079,134 @@ test('does not read a run when the producer does not advertise the action', asyn await expect(client.getRun('run-1')).rejects.toMatchObject({ code: 'capability_unsupported' }); expect(transport.readDefinitions).not.toHaveBeenCalled(); }); + +const historyAction = 'coven.automations.occurrence.history.v1'; + +function historyRow(day: number, state = 'planned') { + return { + ...occurrenceSnapshot(), id: `morning-${day}`, state, + scheduledFor: `2026-09-${String(day).padStart(2, '0')}T09:00:00.000Z`, + }; +} + +function historyCursorFor(row: { scheduledFor: string; id: string }): string { + return Buffer.from(JSON.stringify([row.scheduledFor, row.id])).toString('base64url'); +} + +function historyPage(rows: readonly unknown[], cursor: Record) { + return { automationId: 'morning', occurrences: rows, cursor }; +} + +test('reads one page of an automation\'s history and sends the trimmed id', async () => { + const rows = [historyRow(5), historyRow(4, 'succeeded')]; + const next = historyCursorFor(rows[1]!); + const { client, transport } = readSetup(historyPage(rows, { hasMore: true, next }), historyAction); + const result = await client.occurrenceHistory(' morning ', { limit: 2 }); + expectTypeOf(result).toEqualTypeOf(); + expect(result).toEqual({ automationId: 'morning', data: rows, cursor: { hasMore: true, next } }); + expect(transport.readDefinitions.mock.calls[0]?.[0]).toEqual({ action: historyAction, automationId: ' morning ', limit: 2 }); + expect(Object.isFrozen(transport.readDefinitions.mock.calls[0]?.[0])).toBe(true); +}); + +test('reads a later page by cursor and requires the producer to echo it', async () => { + const cursor = historyCursorFor(historyRow(4)); + const rows = [historyRow(3, 'failed')]; + const { client, transport } = readSetup(historyPage(rows, { hasMore: false, current: cursor }), historyAction); + expect(await client.occurrenceHistory('morning', { cursor })).toEqual({ + automationId: 'morning', data: rows, cursor: { hasMore: false, current: cursor }, + }); + expect(transport.readDefinitions.mock.calls[0]?.[0]).toEqual({ action: historyAction, automationId: 'morning', limit: 20, cursor }); +}); + +test('iterates every page newest first within maxPages', async () => { + const rows = [5, 4, 3, 2, 1].map((day) => historyRow(day)); + const { client, transport } = readSetup(); + transport.readDefinitions.mockImplementation((request) => { + const read = request as { cursor?: string; limit: number }; + const start = read.cursor === undefined ? 0 : rows.findIndex((row) => historyCursorFor(row) === read.cursor) + 1; + const slice = rows.slice(start, start + read.limit); + const more = start + read.limit < rows.length; + const cursor = { + hasMore: more, + ...(read.cursor === undefined ? {} : { current: read.cursor }), + ...(more ? { next: historyCursorFor(slice.at(-1)!) } : {}), + }; + return Promise.resolve({ + status: 200, body: Buffer.from(JSON.stringify(readEnvelope(historyAction, historyPage(slice, cursor)))), + }); + }); + const seen: string[] = []; + for await (const occurrence of client.iterateOccurrenceHistory('morning', { limit: 2, maxPages: 5 })) { + seen.push(occurrence.id); + } + expect(seen).toEqual(['morning-5', 'morning-4', 'morning-3', 'morning-2', 'morning-1']); + expect(transport.readDefinitions).toHaveBeenCalledTimes(3); + + const bounded: string[] = []; + for await (const occurrence of client.iterateOccurrenceHistory('morning', { limit: 2, maxPages: 1 })) { + bounded.push(occurrence.id); + } + expect(bounded).toEqual(['morning-5', 'morning-4']); + expect(() => client.iterateOccurrenceHistory('morning', { limit: 2 } as never)).toThrow(/maxPages or a caller-owned signal/u); +}); + +test('accepts a full history page above the 16 KiB read cap', async () => { + const rows = Array.from({ length: 100 }, (_, index) => ({ + ...historyRow(1), id: `morning-${String(999 - index).padStart(3, '0')}`, failureReason: 'x'.repeat(200), + })); + const page = historyPage(rows, { hasMore: false }); + expect(Buffer.byteLength(JSON.stringify(readEnvelope(historyAction, page)))).toBeGreaterThan(16_384); + const { client } = readSetup(page, historyAction); + expect((await client.occurrenceHistory('morning', { limit: 100 })).data).toHaveLength(100); + const { client: oversized, transport } = readSetup(); + transport.readDefinitions.mockResolvedValue({ status: 200, body: Buffer.alloc(262_145, 32) }); + await expect(oversized.occurrenceHistory('morning')).rejects.toMatchObject({ code: 'invalid_response' }); +}); + +const historyCursor = historyCursorFor(historyRow(9)); +test.each([ + ['another automation id', { automationId: 'evening' }, undefined], + ['a row from another automation', { occurrences: [{ ...historyRow(5), automationId: 'evening' }] }, undefined], + ['rows oldest first', { occurrences: [historyRow(4), historyRow(5)] }, undefined], + ['a repeated row', { occurrences: [historyRow(5), historyRow(5)] }, undefined], + ['more rows than the limit', { occurrences: [historyRow(5), historyRow(4), historyRow(3)] }, undefined], + ['no cursor', { cursor: undefined }, undefined], + ['hasMore without next', { cursor: { hasMore: true } }, undefined], + ['next without hasMore', { cursor: { hasMore: false, next: historyCursor } }, undefined], + ['a malformed next', { cursor: { hasMore: true, next: 'not a cursor' } }, undefined], + ['a short page claiming more', { occurrences: [historyRow(5)], cursor: { hasMore: true, next: historyCursor } }, undefined], + ['an unknown cursor field', { cursor: { hasMore: false, previous: historyCursor } }, undefined], + ['a current on the first page', { cursor: { hasMore: false, current: historyCursor } }, undefined], + ['no current echo', { cursor: { hasMore: false } }, historyCursor], + ['a different current', { cursor: { hasMore: false, current: historyCursorFor(historyRow(8)) } }, historyCursor], + ['a next equal to the cursor', { cursor: { hasMore: true, current: historyCursor, next: historyCursor } }, historyCursor], +])('rejects a history page with %s', async (_label, change, cursor) => { + const page = { ...historyPage([historyRow(5), historyRow(4)], { hasMore: false }), ...change }; + const { client } = readSetup(page, historyAction); + await expect(client.occurrenceHistory('morning', { limit: 2, ...(cursor === undefined ? {} : { cursor }) })) + .rejects.toMatchObject({ code: 'invalid_response' }); +}); + +test.each([ + ['an empty id', '', {}], + ['a blank id', ' ', {}], + ['a zero limit', 'morning', { limit: 0 }], + ['an extra field', 'morning', { view: 'due' }], + ['a padded cursor', 'morning', { cursor: `${historyCursor}=` }], + ['a non-string cursor', 'morning', { cursor: 7 }], + ['a cursor over 512 characters', 'morning', { cursor: 'A'.repeat(516) }], + ['a non-canonical cursor', 'morning', { cursor: 'AB' }], +])('rejects history with %s before transport', async (_label, automationId, query) => { + const { client, transport } = readSetup(); + await expect(client.occurrenceHistory(automationId, query as never)).rejects.toMatchObject({ code: 'invalid_options' }); + expect(transport.readDefinitions).not.toHaveBeenCalled(); +}); + +test('does not read history when the producer does not advertise the action', async () => { + const { client, transport } = readSetup(historyPage([], { hasMore: false }), historyAction); + const advertised = advertisement(); + advertised.capabilities[0]!.actions = advertised.capabilities[0]!.actions.filter((action) => action !== historyAction); + transport.capabilities.mockResolvedValue({ status: 200, body: Buffer.from(JSON.stringify(advertised)) }); + await expect(client.occurrenceHistory('morning')).rejects.toMatchObject({ code: 'capability_unsupported' }); + expect(transport.readDefinitions).not.toHaveBeenCalled(); +});