From 59042f2a3960855199ae03ca29e749e3cb8c88ec Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 28 Aug 2026 05:50:27 -0700 Subject: [PATCH] Re-vendor @intx/hub-sessions at a8bc06ae; retire three deltas upstream absorbed Upstream now owns what three workbench deltas patched: fab86ca9 emits TurnUsage once per finalized turn (replacing the per-event inference.usage forward, which double-counted the cumulative usage ticks), a1d419c3 serializes collector dispatch at the registry (replacing the collector's own promise chain), and a203a057 + 9e11829f sequence prepare -> INSERT anchor -> emit with frameSent-tagged DeployFrameFailure (replacing the insert-before-frame ordering and DeployFrameNotSentError). The hub's usage sink now consumes TurnUsage, which carries tenant/session/run identity itself. Re-applied onto upstream's files: the adopted deploy front + sourceRef threading (CL-6324), the wire-projection writer (CL-6324), pack acceptance (ownsWorkflowRunRepo, anchorAddressForPackSource, decideTerminalRunFlip), malformed tool-call-name sanitization (CL-6478) and the sealed-run terminal-status backfill (CL-6595), the last now classifying through upstream's classifyTerminalEvent. Tests for the retired deltas go with them; the surviving ones stay. The still-pinned workflow-host gains a second bridging edit: its boot replay reads ownedMessageIds from scanRunsForBoot (upstream f89bb51b), since readOwnedMessageIds no longer exists. --- VENDORED.md | 54 +- apps/hub/src/index.ts | 16 +- packages/insights/src/on-usage-wire.test.ts | 32 +- scripts/checks/kill-dates.txt | 4 +- vendor/intx/hub-sessions/VENDORED-FROM | 8 +- vendor/intx/hub-sessions/package.json | 13 +- vendor/intx/hub-sessions/src/agent-repo.ts | 19 +- .../src/event-collector-registry.test.ts | 114 --- .../src/event-collector-registry.ts | 73 +- .../hub-sessions/src/event-collector.test.ts | 108 +-- .../intx/hub-sessions/src/event-collector.ts | 110 ++- vendor/intx/hub-sessions/src/index.ts | 9 +- .../hub-sessions/src/session-service.test.ts | 183 ----- .../intx/hub-sessions/src/session-service.ts | 256 +++++-- vendor/intx/hub-sessions/src/substrate.ts | 5 +- .../hub-sessions/src/workflow-run-kind.ts | 725 ++++++++++++++++-- .../hub-sessions/src/ws/sidecar-handler.ts | 91 ++- vendor/intx/hub-sessions/tsconfig.json | 8 +- vendor/intx/workflow-host/VENDORED-FROM | 2 +- .../src/supervisor/supervisor.ts | 6 +- 20 files changed, 1059 insertions(+), 777 deletions(-) delete mode 100644 vendor/intx/hub-sessions/src/event-collector-registry.test.ts diff --git a/VENDORED.md b/VENDORED.md index 5203d117f..ae0e9390e 100644 --- a/VENDORED.md +++ b/VENDORED.md @@ -28,7 +28,7 @@ never a convenience. | `vendor/intx/agent` | `@intx/agent` source (`src/`, manifest, tsconfig) | [faremeter/interchange](https://github.com/faremeter/interchange) @ `a8bc06ae` (origin/main, 2026-08-27) | npm 0.3.0 predates the operator-configurable doom-loop threshold (`afd0c82b`, `c421c092`) the re-vendored `workflow-host` configures; no local delta; retired by the next `@intx/agent` publish | sawyer | 2026-10-26 | `check:killdates` | | `vendor/intx/db` | `@intx/db` source (`src/`, `migrations/`, drizzle config, manifest, tsconfigs) | [faremeter/interchange](https://github.com/faremeter/interchange) @ `a8bc06ae` (origin/main, 2026-08-27) | npm 0.3.0 covers the base package but not the `wire_projection` column/loader delta (CL-6324) or the `workflow_definition.origin` column separating a definition from the per-run record of one folded run's deploy (CL-6452), shipped as migrations `0086`/`0087` behind upstream's `0085_add_approval_run_idx`, plus `0088` rewriting the retired `onBodyFailure: "continue"` literal to upstream's `"tolerate"` in stored wire projections; retired when upstream absorbs the deltas | sawyer | 2026-10-26 | `check:killdates` | | `vendor/intx/hub-api` | `@intx/hub-api` source (`src/`, manifest, tsconfig) | [faremeter/interchange](https://github.com/faremeter/interchange) @ `b5580a02` (v0.3.0) | npm 0.3.0 covers the base package but not the exported null-principal `resolveApproval` (CL-6345) or the bearer-authenticated workflow-deploy mirror (`middleware/workflow-run-deploy-auth.ts`, CL-workflow-deploy-bearer); retired when upstream absorbs the deltas | sawyer | 2026-09-19 | `check:killdates` | -| `vendor/intx/hub-sessions` | `@intx/hub-sessions` source (`src/`, manifest, tsconfig) | [faremeter/interchange](https://github.com/faremeter/interchange) @ `b5580a02` (v0.3.0) | npm 0.3.0 covers the base package but not the usage forward (CL-5879), pack-acceptance fixes, adopted deploy front, wire-projection writer, event-collector serialization, anchor ordering, or malformed tool-call-name sanitization (CL-6478) | sawyer | 2026-09-19 | `check:killdates` | +| `vendor/intx/hub-sessions` | `@intx/hub-sessions` source (`src/`, manifest, tsconfig) | [faremeter/interchange](https://github.com/faremeter/interchange) @ `a8bc06ae` (origin/main, 2026-08-27) | npm 0.3.0 covers the base package but not the pack-acceptance fixes (`ownsWorkflowRunRepo`, `anchorAddressForPackSource`, `decideTerminalRunFlip`), the adopted deploy front + `sourceRef` (CL-6324), the wire-projection writer (CL-6324), malformed tool-call-name sanitization (CL-6478) or the sealed-run terminal-status backfill (CL-6595); retired when upstream absorbs the deltas | sawyer | 2026-10-26 | `check:killdates` | | `vendor/intx/inference` | `@intx/inference` source (`src/`, manifest, tsconfig) | [faremeter/interchange](https://github.com/faremeter/interchange) @ `a8bc06ae` (origin/main, 2026-08-27) | npm 0.3.0 predates doom-loop detection (`8da4c827`, `afd0c82b`, `c421c092`); one local delta: `providers/google-genai-files.ts` builds its upload body as `new Uint8Array(bytes)` because TS 6's lib.dom `BodyInit` rejects `Uint8Array` (upstream compiles ESNext-only under TS 5.9); retired by the next publish | sawyer | 2026-10-26 | `check:killdates` | | `vendor/intx/mail-memory` | `@intx/mail-memory` source (`src/`, manifest, tsconfig) | [faremeter/interchange](https://github.com/faremeter/interchange) @ `a8bc06ae` (origin/main, 2026-08-27) | npm 0.3.0 predates the `@intx/mailbox` extraction (`af03bb90`), on-demand body reads (`54f7c239`) and `expunge` returning the swept uids (`bcabb1f8`) that the re-vendored `workflow-host` binds against; no local delta; retired by the next publish | sawyer | 2026-10-26 | `check:killdates` | | `vendor/intx/mailbox` | `@intx/mailbox` source (`src/`, manifest, tsconfig) | [faremeter/interchange](https://github.com/faremeter/interchange) @ `a8bc06ae` (origin/main, 2026-08-27) | Never published: a new package at the target pin (`af03bb90`) that `workflow-host`'s substrate mailbox store and supervisor-backed transport import; no local delta; retired by its first publish | sawyer | 2026-10-26 | `check:killdates` | @@ -47,10 +47,11 @@ published `@intx/harness`, `@intx/hub-agent`, `@intx/tool-packaging`, `@intx/authz`, … resolve their own `@intx/*` dependencies onto the vendored copies instead of a second npm copy) and keep the unchanged names on `0.3.0`. `vendor/intx/workflow-host` (still at `b5580a02` until its own re-pin) -carries one bridging edit against the re-pinned `@intx/types`: its -supervisor-backed transport's `expunge` stub returns -`Promise<{ expungedUids: number[] }>` (upstream `bcabb1f8`); it disappears -with that tree's re-pin. +carries two bridging edits against the re-pinned `@intx/types` and +`@intx/hub-sessions`: its supervisor-backed transport's `expunge` stub +returns `Promise<{ expungedUids: number[] }>` (upstream `bcabb1f8`), and its +boot replay reads `ownedMessageIds` from `scanRunsForBoot` (upstream +`f89bb51b`); both disappear with that tree's re-pin. The pinned commit `b5580a02` is upstream's `v0.3.0` release tag, 16 commits past the previous pin `4ed8baf4`: a workflow-host supervisor @@ -99,14 +100,7 @@ same asset tenant-scoping (a foreign-tenant asset already read as `not_found` before this change and still does), same install/probe/gate/freeze call into `sessionService.deployWorkflowFromSource`. A request with no bearer credential falls through unchanged to the session path. -`vendor/intx/hub-sessions` (CL-5879) forwards `inference.usage` events — -previously matched by `event-collector.ts`'s "not persisted" default and -dropped — to a new optional `onUsage` callback on `createEventCollector` -and `createEventCollectorRegistry`, carrying `{turnId, provider, model, usage}` plus -the registry's own `tenantId`/`sessionId`; no new persistence lands in the -vendored copy itself, only the forward, and `apps/hub/src/index.ts` wires -it to `@corbits/insights`' `createUsageSink` so `usage_turn` rows are -written for the first time. `vendor/intx/hub-sessions` also drops the +`vendor/intx/hub-sessions` also drops the live-status gate on `receiveWorkflowRunPack`'s anchor lookup: the gate is now the exported pure helper `ownsWorkflowRunRepo` (a self-anchored `workflow_run` row with a routable address), with the allocation fences unchanged. Upstream @@ -195,42 +189,12 @@ folded run's system prompt, tool pins, model, credential bindings) had nowhere to get it. Keyed to the approved wire hash and stored beside it, this is one store per concept, not a second copy: the projection and the hash that addresses it are written and read together. -`vendor/intx/hub-sessions` (CL-6379) serializes the event -collector's `onEvent`/`abandon` through an internal promise chain: the -registry's dispatch is deliberately fire-and-forget, and without the chain -two events interleave across their DB awaits — a `connector.reply` finalize -nulls the current turn while `inference.done` is still inserting parts -(dropped as "no active turn"), and a finalize processed during the next -`inference.start`'s begin-insert marks the NEW turn finalized, leaving its -row "running" forever. The same change classifies an accepted workflow-run +`vendor/intx/hub-sessions` (CL-6379) classifies an accepted workflow-run pack's newly-terminal runs through the new pure `decideTerminalRunFlip` before the DB flip: a section occurrence's repo-local child run (`turn__`) has no `workflow_run` row by design and is skipped quietly instead of being logged as a foreign-deployment violation on every turn. -`vendor/intx/hub-sessions` (CL-6388) reorders -`deployCodeSourcedWorkflow`, the SHARED code-sourced deploy front, to -persist its anchor `workflow_run` row BEFORE emitting the source-ref -deploy frame. Upstream inserted the row only after the sidecar's deploy -ack, but the frame spawns the deployment's child, whose first -`refs/heads/events` pack push races that ack back to the hub — -`receiveWorkflowRunPack` fails closed (`path_violation`) on a missing -anchor row, so every fresh deployment's first events pack was rejected -and the durable event log never bootstrapped. The row is born with a -null `publicKey` (the reconnect challenge keeps failing closed until the -ack), the acked supervisor key is stamped afterwards, and a failed frame -emit deletes the pre-inserted row. The prepared and adopted fronts -already had their anchor row pre-frame; the shared front now matches -them. `vendor/intx/hub-sessions` (CL-6395) narrows that failed-emit -deletion: CL-6388 deleted the pre-inserted anchor on ANY -`emitSourceRefDeployFrame` rejection, but an ack-timeout or socket-drop -rejection fires strictly after the `agent.deploy` frame already reached -the sidecar, so deleting the row could permanently orphan an -already-spawned child on the missing-anchor `path_violation` path. -`ws/sidecar-handler.ts` now exports `DeployFrameNotSentError`, thrown only -where a guard clause or the `conn.send()` call itself fails before the -frame could have reached the wire; `deployCodeSourcedWorkflow` deletes the -row only on that error and otherwise keeps the row and logs a -reconciliation line. `vendor/intx/hub-sessions` (CL-6478) adds a +`vendor/intx/hub-sessions` (CL-6478) adds a `sanitize-tool-name.ts` module and calls it from `event-collector.ts`'s `tool_call` handling: `@intx/inference`'s `decodeToolName` is deliberately total and returns a hallucinated or provider-mangled function name diff --git a/apps/hub/src/index.ts b/apps/hub/src/index.ts index a2eaa61be..49c79fd3d 100644 --- a/apps/hub/src/index.ts +++ b/apps/hub/src/index.ts @@ -768,16 +768,13 @@ export async function createHub(config: HubConfig) { db: withTurnPartPersistGuard(withTurnPartWriteDefaults(db)), onTurnFinalized: (agentAddress, turn) => artifactDeliveryHandlerRef.current?.(agentAddress, turn), - // The vendored `onUsage` forward (see VENDORED.md) — the platform - // collector accumulates turns per session but never persisted - // `inference.usage`; this is the first point tenantId + turnId + - // provider + model + tokens are all in scope at once. - onUsage: (_agentAddress, tenantId, sessionId, usage) => { + // Per-turn usage, emitted once when the collector finalizes a turn. + onUsage: (_agentAddress, usage) => { void usageSink .handle({ turnId: usage.turnId, - tenantId, - sessionId, + tenantId: usage.tenantId, + sessionId: usage.sessionId, provider: usage.provider, model: usage.model, tokens: usage.usage, @@ -1618,9 +1615,8 @@ export async function createHub(config: HubConfig) { // (see @corbits/insights' createDrizzleRunTraceReader) — no new storage, // same `db` handle every other platform-table reader in this file uses. // The sink itself is constructed earlier, alongside `eventCollectors` - // (see the vendored `onUsage` forward on `createEventCollectorRegistry` - // above), since that's the only place tenantId/turnId/model land - // together on an `inference.usage` event. + // (see the `onUsage` hook on `createEventCollectorRegistry` above), + // which reports each finalized turn's usage with its run identity. app.route( `${TENANT_PREFIX}/insights`, createInsightsRoutes({ diff --git a/packages/insights/src/on-usage-wire.test.ts b/packages/insights/src/on-usage-wire.test.ts index b2d80f405..c4f7c6404 100644 --- a/packages/insights/src/on-usage-wire.test.ts +++ b/packages/insights/src/on-usage-wire.test.ts @@ -4,12 +4,15 @@ import { createUsageSink } from "./collector"; import { createMemoryUsageStore } from "./store"; /** - * Exact `UsageForwarded` payload the vendored event-collector emits - * (`vendor/intx/hub-sessions/src/event-collector.ts` `UsageForwarded`). - * Hub `onUsage` (`apps/hub/src/index.ts` ~572) remaps it to `UsageEvent` as: + * Exact `TurnUsage` payload `@intx/hub-sessions`' event collector emits once + * per finalized turn. Hub `onUsage` (`apps/hub/src/index.ts`) remaps it to + * `UsageEvent` as * `{ turnId, tenantId, sessionId, provider, model, tokens: usage.usage }` */ -type HubUsageForwarded = { +type HubTurnUsage = { + tenantId: string; + sessionId: string; + runId: string; turnId: string; provider: string; model: string; @@ -22,15 +25,11 @@ type HubUsageForwarded = { }; }; -function hubOnUsageHandleArg( - tenantId: string, - sessionId: string, - usage: HubUsageForwarded, -) { +function hubOnUsageHandleArg(usage: HubTurnUsage) { return { turnId: usage.turnId, - tenantId, - sessionId, + tenantId: usage.tenantId, + sessionId: usage.sessionId, provider: usage.provider, model: usage.model, tokens: usage.usage, @@ -38,7 +37,7 @@ function hubOnUsageHandleArg( } describe("hub onUsage wire → createUsageSink", () => { - test("inserts a row from the exact hub remapping of UsageForwarded", async () => { + test("inserts a row from the exact hub remapping of TurnUsage", async () => { const store = createMemoryUsageStore(); let n = 0; const sink = createUsageSink({ @@ -46,7 +45,10 @@ describe("hub onUsage wire → createUsageSink", () => { generateId: () => `id-${++n}`, }); - const forwarded: HubUsageForwarded = { + const forwarded: HubTurnUsage = { + tenantId: "tenant-acme", + sessionId: "session-1", + runId: "run-1", turnId: "turn-wire-1", provider: "anthropic", model: "claude-sonnet", @@ -59,9 +61,7 @@ describe("hub onUsage wire → createUsageSink", () => { }, }; - const status = await sink.handle( - hubOnUsageHandleArg("tenant-acme", "session-1", forwarded), - ); + const status = await sink.handle(hubOnUsageHandleArg(forwarded)); expect(status).toBe("inserted"); const rows = await store.listUsageByTenants(["tenant-acme"]); diff --git a/scripts/checks/kill-dates.txt b/scripts/checks/kill-dates.txt index 0d7902fc6..c13a14849 100644 --- a/scripts/checks/kill-dates.txt +++ b/scripts/checks/kill-dates.txt @@ -17,7 +17,7 @@ apps/sidecar | sawyer | 2026-09-19 vendor/intx/agent | sawyer | 2026-10-26 | d0d56d9f452b78f4b541ad8f4e89f975e8069446bb98f2c8097b90de4b020243 vendor/intx/db | sawyer | 2026-09-19 | 0a4cdb9a8a6ff19d5d4713cbc4f5cc9257aad839b1fa393b2026e6d5afd828b9 vendor/intx/hub-api | sawyer | 2026-09-19 | f93a383cb5d6acdf50a461b43e4c8991dbdf7e13d34a4e5598e556eaee66308b -vendor/intx/hub-sessions | sawyer | 2026-09-19 | e3344253ed6f2b0998599f9f8712d4e9b0757895b31a3e7df6a245c2bb76af3c +vendor/intx/hub-sessions | sawyer | 2026-09-19 | 53addc3090ad9f54bc4bac8fb50ad8d567ccf46f30bb5403d447351cb16b4fb6 vendor/intx/inference | sawyer | 2026-10-26 | 77fec29b078e8d03e686747c70e6b62ac1fd1434db0fb2c1e12e84b6dc71465f vendor/intx/mail-memory | sawyer | 2026-10-26 | 9f3601a7fb22e2d1c63daa976f3afccbd79af2187c155a0080c0d60c82450b92 vendor/intx/mailbox | sawyer | 2026-10-26 | d36d7ffcc32018571276e4922a8c2714b7ee0bb5deb80b01e73859245975d4c6 @@ -25,7 +25,7 @@ vendor/intx/mime | sawyer | 2026-10-26 | d02e5f8f1429eac7c27d3a37eec31111f8a1053 vendor/intx/types | sawyer | 2026-10-26 | ec1de14b859007b4db137da1533d4ce79d11024ad69c8b36938017319d6d8e86 vendor/intx/workflow | sawyer | 2026-09-19 | 4b51b9bd6a124cfaa0c916e2b26c04ac9170618bb092f0c8e1c312263fc84fdf vendor/intx/workflow-deploy | sawyer | 2026-10-26 | 960a2ae408223649fe8be0e3b9d63f2b0cca25259bc0ae06761bca521bb738e5 -vendor/intx/workflow-host | sawyer | 2026-09-19 | 48ae3e34c6f14b99a3a940ede98b119e52dfe7410a86f644d3b78cf6e7d44f4d +vendor/intx/workflow-host | sawyer | 2026-09-19 | 001bea028bf1484f150fa00c9331b897a2db888f6aacd0a109f4e2739854cf74 packages/folded-runs | sawyer | 2026-11-01 diff --git a/vendor/intx/hub-sessions/VENDORED-FROM b/vendor/intx/hub-sessions/VENDORED-FROM index b882821e1..b4c90f050 100644 --- a/vendor/intx/hub-sessions/VENDORED-FROM +++ b/vendor/intx/hub-sessions/VENDORED-FROM @@ -1,8 +1,4 @@ Source: https://github.com/faremeter/interchange (packages/hub-sessions) -Commit: b5580a02fb918eebccc33ded7727ffee781ffbd1 (tag v0.3.0) +Commit: a8bc06ae38661c5e0ed91ded8559bf09f502213d (origin/main, 2026-08-27) License: LGPL-2.1-only (see vendor/intx/LICENSE) -Local modifications: exports map repointed from the upstream intx-src condition to direct TypeScript source resolution (types/default -> ./src/...); dist references removed. CL-5879: event-collector.ts's inference.usage case (previously falling into the "not persisted" default) now forwards {turnId, provider, model, usage} to an optional `onUsage` callback, threaded through event-collector-registry.ts's EventCollectorRegistryConfig as `onUsage(agentAddress, tenantId, sessionId, usage)` — the collector's own turn/tenant state is the only place these identifiers meet an inference.usage event. No persistence added upstream; the app wires the callback to @corbits/insights' usage sink. Terminal-anchor pack acceptance: hub-session-lookups.ts's receiveWorkflowRunPack no longer gates the anchor lookup on liveWorkflowRunStatuses — the ownership gate is the exported pure helper ownsWorkflowRunRepo (self-anchored row with a routable address), so a terminal run can still land the inbox-enqueue and markConsumed-rejection packs that retire mail which arrived in its teardown window. Upstream's live-status gate made that pair unresolvable: pack rejected as path_violation -> ack withheld -> hub redelivers, forever. CL-6324: a third code-sourced deploy front, `deployAdoptedCodeSourcedWorkflow` (plus the `deployAdoptedWorkflowFromSource` service method and its `AdoptingWorkflowDeployer` type), deploys onto shared capacity while ADOPTING an anchor `workflow_run` row the caller already owns. Upstream's two fronts cannot: `deployWorkflowFromSource` INSERTs its anchor (a primary-key collision against a folded run's existing row) and threads no `credentialCipher`, and `deployPreparedCodeSourcedWorkflow` does both correctly but only under the allocation-ownership lock. The new front composes the same private halves (`emitSourceRefDeployFrame`, `buildInertProjectionStepSources`) and follows the prepared front's semantics minus the allocation lock: ownership is the anchor row's own tenant plus self-anchoring, checked before the frame and re-asserted on the guarded UPDATE that stamps `definitionId`/`publicKey`. See VENDORED.md and docs/revendor-inventory.md. -CL-6388: `deployCodeSourcedWorkflow` now INSERTs its anchor `workflow_run` row BEFORE emitting the source-ref deploy frame (publicKey null until the ack stamps it; a failed emit deletes the row). Upstream's frame-then-insert ordering let the spawned child's first refs/heads/events pack push race the deploy ack, and receiveWorkflowRunPack fails closed (path_violation) on the missing anchor row, so every fresh deployment's first events pack was rejected and the durable event log never bootstrapped. -CL-6395: CL-6388's "a failed emit deletes the row" was too broad — any rejection from `emitSourceRefDeployFrame`, including an ack-timeout or socket-drop that fires strictly AFTER the `agent.deploy` frame already reached the sidecar, deleted the anchor row and permanently orphaned an already-spawned child on the missing-anchor `path_violation` path. `ws/sidecar-handler.ts` now exports `DeployFrameNotSentError`, thrown only by a guard clause that runs before `conn.send()` or by `conn.send()` itself throwing synchronously — the sole cases that provably never reached the wire; every other deploy rejection (timeout, disconnect, reconnect takeover, ack-processing failure) is raised through the pending-deploy's `reject()`, which by construction only fires after the send. `deployCodeSourcedWorkflow` deletes the pre-inserted row only on `DeployFrameNotSentError`; any other failure keeps the row and logs one reconciliation line. Also corrects an overclaiming comment in hub-session-lookups.ts: `markTerminal`'s null return means no row in a LIVE status (`deployed` or `running`) matched, not specifically "running". -CL-6478: `event-collector.ts`'s `tool_call` handling in `handleInferenceDone` now runs `block.name` through a new `sanitize-tool-name.ts` module before persisting it. `@intx/inference`'s `decodeToolName` is deliberately total — a hallucinated or provider-mangled function name is returned verbatim rather than throwing — but `encodeToolName` throws when that same name is later put back on the wire to build the next turn's outbound request, so persisting a decoded name unchecked wedged the room forever once the bad name was durable. `sanitizeToolNameForPersistence` round-trips the name through `encodeToolName` before it is written; a name that cannot be re-encoded collapses to a stable `malformed_tool_call` placeholder instead. `@intx/inference` is added to this package's own `package.json` dependencies for the check. -CL-6595: `workflow-run-kind.ts`'s newly-terminal detection in `validatePush` only ever scanned a run's per-event `runs//events/.json` blobs; `enumerateEventBlobs` explicitly skips a run whose events already live in a combined `events.jsonl` (`hasCombined` -> `continue`), so a run sealed from birth — its entire event log, including the terminal event, arriving pre-combined in a single push with no per-event blobs ever landing — was never surfaced as newly terminal and `markTerminal` never fired, leaving `workflow_run.status` stuck live forever despite the run having genuinely finished. `validatePush` now also walks `validateCombinedEventRuns`' `combinedRunIds` and reports a newly-sealed run (absent from the prior tree's combined form) as terminal by reading its combined log's last (terminal, by `checkCombinedStructure`'s own invariant) event. A new `readCommittedWorkflowRunTerminalStatus` export mirrors `readCommittedWorkflowRunLifecycle` but returns the mapped `workflow_run.status` value instead of just live/terminal/absent; `hub-session-lookups.ts`'s pack-receive path now calls it as a same-push defense-in-depth backfill (calling `markTerminal` directly) whenever the committed log proves a run terminal, independent of whether the primary per-push detection caught it. +Local modifications: exports map repointed from the upstream intx-src condition to direct TypeScript source resolution (types/default -> ./src/...); dist references removed. Terminal-anchor pack acceptance: hub-session-lookups.ts's receiveWorkflowRunPack gates the anchor lookup on the exported pure helper ownsWorkflowRunRepo (self-anchored row with a routable address, no liveness requirement) instead of upstream's `status in (deployed, running)`, so a terminal run can still land the inbox-enqueue and markConsumed-rejection packs that retire mail arriving in its teardown window; CL-6361 widens the same lookup to peel a per-step pack source address back to its base run's anchor via anchorAddressForPackSource; CL-6379 classifies an accepted pack's newly-terminal runs through decideTerminalRunFlip so a section occurrence's repo-local child run (turn__) is skipped quietly. CL-6324: a third code-sourced deploy front, deployAdoptedCodeSourcedWorkflow (plus the deployAdoptedWorkflowFromSource service method and its AdoptingWorkflowDeployer / DeployAdoptedWorkflowFromSourceParams types), deploys onto shared capacity while ADOPTING an anchor workflow_run row the caller already owns; it composes upstream's emitSourceRefDeployFrame and a guarded UPDATE. DeployWorkflowFromSourceParams / DeployPreparedCodeSourcedWorkflowParams gain an optional sourceRef threaded through bindAssetAttachmentResolver / bindSourceAttachmentResolver so a per-run source tree on refs/heads/runs/ packs the pinned commit. CL-6324: workflow-probe-gate.ts's PersistFrozenApprovalFn carries the inert projection and createDbFrozenApprovalWriter stamps workflow_definition_version.wire_projection in the same transaction as approved_wire_hash. CL-6478: sanitize-tool-name.ts + event-collector.ts's tool_call case persist only a tool-call name encodeToolName can re-invert (anything else collapses to MALFORMED_TOOL_NAME), so one bad name fails its turn instead of wedging the room; @intx/inference is a dependency for this. CL-6595: workflow-run-kind.ts's validatePush also reports a run sealed from birth (combined events.jsonl with no per-event blobs) as newly terminal, and the new readCommittedWorkflowRunTerminalStatus export backs a same-push markTerminal backfill in hub-session-lookups.ts; both classify through upstream's classifyTerminalEvent. Retired at this pin: the per-event inference.usage forward (upstream fab86ca9 emits TurnUsage once per turn), the collector-level event serialization (upstream a1d419c3 serializes at the registry) and the anchor-before-frame ordering with DeployFrameNotSentError (upstream a203a057 + 9e11829f, isDeployFrameFailure). diff --git a/vendor/intx/hub-sessions/package.json b/vendor/intx/hub-sessions/package.json index 56ecd1baf..7fa1ca8fe 100644 --- a/vendor/intx/hub-sessions/package.json +++ b/vendor/intx/hub-sessions/package.json @@ -15,23 +15,22 @@ } }, "scripts": { - "typecheck": "tsc --noEmit", - "test": "bun test" + "typecheck": "tsc --noEmit" }, "dependencies": { - "@intx/agent": "0.3.0", + "@intx/agent": "workspace:*", "@intx/crypto": "0.3.0", "@intx/db": "workspace:*", "@intx/hub-common": "0.3.0", - "@intx/inference": "0.3.0", + "@intx/inference": "workspace:*", "@intx/log": "0.3.0", - "@intx/mime": "0.3.0", + "@intx/mime": "workspace:*", "@intx/pack-transport": "0.3.0", "@intx/storage-isogit": "0.3.0", "@intx/tool-packaging": "0.3.0", - "@intx/types": "0.3.0", + "@intx/types": "workspace:*", "@intx/workflow": "workspace:*", - "@intx/workflow-deploy": "0.3.0", + "@intx/workflow-deploy": "workspace:*", "arktype": "catalog:", "drizzle-orm": "catalog:", "isomorphic-git": "catalog:", diff --git a/vendor/intx/hub-sessions/src/agent-repo.ts b/vendor/intx/hub-sessions/src/agent-repo.ts index f6c14ba46..84311b76c 100644 --- a/vendor/intx/hub-sessions/src/agent-repo.ts +++ b/vendor/intx/hub-sessions/src/agent-repo.ts @@ -122,6 +122,21 @@ export type AgentRepoStore = { readonly repoStore: RepoStore; }; +/** + * Repo kinds eligible for write-path object GC. Deliberately EXCLUDES + * "workflow-run". The warm-agent mailbox physically expunges `.eml` + * blobs from the live tree; the raw bytes then persist only through the + * parent commit in git history, and they survive a push ONLY because a + * `workflow-run` repo's objects are never GC'd. Adding "workflow-run" + * here -- above all with a retention other than "keep-history" -- would + * make an expunged message's bytes prunable and silently destroy the + * mailbox audit trail. The mailbox subtree contract in `workflow-run-kind` + * documents this dependency; `agent-repo.test.ts` pins it. Do not add + * "workflow-run" without replacing physical expunge with a tip-reachable + * retain-bytes scheme first. + */ +export const DEFAULT_GC_KINDS = ["agent-state"] as const; + export function createAgentRepoStore(config: { dataDir: string; signingKey: { privateKey: Uint8Array; publicKey: Uint8Array }; @@ -177,7 +192,9 @@ export function createAgentRepoStore(config: { }, authorize, signingCallback: () => signer, - ...(gc === undefined ? {} : { gc: { kinds: ["agent-state"], ...gc } }), + ...(gc === undefined + ? {} + : { gc: { kinds: [...DEFAULT_GC_KINDS], ...gc } }), }); const hub: AgentStateHubPrincipal = { kind: "hub" }; diff --git a/vendor/intx/hub-sessions/src/event-collector-registry.test.ts b/vendor/intx/hub-sessions/src/event-collector-registry.test.ts deleted file mode 100644 index 348ddf249..000000000 --- a/vendor/intx/hub-sessions/src/event-collector-registry.test.ts +++ /dev/null @@ -1,114 +0,0 @@ -import { describe, expect, test } from "bun:test"; - -import type { InferenceEvent } from "@intx/types/runtime"; -import type { DB } from "@intx/db"; - -import { createEventCollectorRegistry } from "./event-collector-registry"; - -// Captures every insert/update the collector issues, resolving each on a -// later tick so concurrent onEvent calls would interleave exactly the way -// they do against a real database. Rows are told apart by shape: an -// inference_turn insert carries `model`, a turn_part insert carries -// `ordinal`. -function createRecordingDb(): { - db: DB["db"]; - turnInserts: Record[]; - partInserts: Record[]; - turnUpdates: Record[]; -} { - const turnInserts: Record[] = []; - const partInserts: Record[] = []; - const turnUpdates: Record[] = []; - const later = () => new Promise((resolve) => setTimeout(resolve, 1)); - const db = { - insert: () => ({ - values: async (row: Record) => { - await later(); - if ("model" in row) turnInserts.push(row); - else partInserts.push(row); - }, - }), - update: () => ({ - set: (row: Record) => ({ - where: async () => { - await later(); - turnUpdates.push(row); - }, - }), - }), - } as unknown as DB["db"]; - return { db, turnInserts, partInserts, turnUpdates }; -} - -function turnEvents(seqBase: number, text: string): InferenceEvent[] { - return [ - { - type: "inference.start", - seq: seqBase, - data: { model: "test-model" }, - }, - { - type: "inference.done", - seq: seqBase + 1, - data: { turn: { content: [{ type: "text", text }] } }, - }, - { - type: "connector.reply", - seq: seqBase + 2, - data: { content: text }, - }, - ] as InferenceEvent[]; -} - -async function until(check: () => boolean, ms: number): Promise { - const deadline = Date.now() + ms; - while (!check() && Date.now() < deadline) { - await new Promise((resolve) => setTimeout(resolve, 5)); - } -} - -describe("event collector registry dispatch ordering (CL-6379)", () => { - test("back-to-back turn events dispatched without awaiting persist every part, in order, and finalize every turn", async () => { - const { db, turnInserts, partInserts, turnUpdates } = createRecordingDb(); - const registry = createEventCollectorRegistry({ db }); - registry.create("agent@test", "tenant-1", "ses_1", "run-1"); - - // Fire-and-forget, exactly as the session orchestrator's `agent.event` - // listener does: two full turns arrive faster than any DB roundtrip. - for (const event of [...turnEvents(1, "first"), ...turnEvents(4, "second")]) { - registry.dispatch("agent@test", event); - } - - await until(() => turnUpdates.length >= 2, 2000); - - // Two turns opened, two finalized — the second turn must not be left - // "running" because the first turn's finalize interleaved with it. - expect(turnInserts).toHaveLength(2); - expect(turnUpdates).toHaveLength(2); - expect(turnUpdates.map((u) => u["status"])).toEqual([ - "completed", - "completed", - ]); - - // Each turn persists step-start, text, step-finish — nothing dropped - // by "no active turn", and ordinals reflect event order. - const byTurn = new Map[]>(); - for (const part of partInserts) { - const parts = byTurn.get(part["turnId"]) ?? []; - parts.push(part); - byTurn.set(part["turnId"], parts); - } - expect(byTurn.size).toBe(2); - for (const parts of byTurn.values()) { - expect(parts.map((p) => p["type"])).toEqual([ - "step-start", - "text", - "step-finish", - ]); - expect(parts.map((p) => p["ordinal"])).toEqual([0, 1, 2]); - } - - // After both turns settle the collector is idle: no current turn. - expect(registry.getCurrentTurnId("agent@test")).toBeNull(); - }); -}); diff --git a/vendor/intx/hub-sessions/src/event-collector-registry.ts b/vendor/intx/hub-sessions/src/event-collector-registry.ts index a240e36e2..a81fe4f0d 100644 --- a/vendor/intx/hub-sessions/src/event-collector-registry.ts +++ b/vendor/intx/hub-sessions/src/event-collector-registry.ts @@ -14,7 +14,7 @@ import { createEventCollector, type EventCollector, type TurnFinalized, - type UsageForwarded, + type TurnUsage, } from "./event-collector"; const log = getLogger(["hub", "event-collector-registry"]); @@ -38,12 +38,7 @@ export type EventCollectorRegistry = { export type EventCollectorRegistryConfig = { db: DB["db"]; onTurnFinalized?: (agentAddress: string, turn: TurnFinalized) => void; - onUsage?: ( - agentAddress: string, - tenantId: string, - sessionId: string, - usage: UsageForwarded, - ) => void; + onUsage?: (agentAddress: string, usage: TurnUsage) => void; }; export function deriveStatus(event: InferenceEvent): SessionStatus | null { @@ -74,6 +69,9 @@ export function createEventCollectorRegistry( const { db, onTurnFinalized, onUsage } = config; const collectors = new Map(); const statuses = new Map(); + // Per-address tail promise: serializes onEvent/abandon work for one run + // address so their DB writes cannot interleave. Reaped when it drains. + const tails = new Map>(); function create( agentAddress: string, @@ -99,8 +97,7 @@ export function createEventCollectorRegistry( : {}), ...(onUsage ? { - onUsage: (usage: UsageForwarded) => - onUsage(agentAddress, tenantId, sessionId, usage), + onUsage: (usage: TurnUsage) => onUsage(agentAddress, usage), } : {}), }); @@ -113,6 +110,31 @@ export function createEventCollectorRegistry( statuses.delete(agentAddress); } + // Chain `work` onto the address's tail so per-address work runs in order and + // never interleaves, while the caller stays non-blocking. `onError` swallows + // a failure so one bad event cannot wedge the chain; `onSettled` runs after + // the work settles. The tail entry is reaped once no later work is queued. + function enqueue( + agentAddress: string, + work: () => Promise, + onError: (err: unknown) => void, + onSettled?: () => void, + ): void { + const prev = tails.get(agentAddress) ?? Promise.resolve(); + const next = prev + .then(work) + .catch(onError) + .finally(() => { + if (onSettled !== undefined) { + onSettled(); + } + if (tails.get(agentAddress) === next) { + tails.delete(agentAddress); + } + }); + tails.set(agentAddress, next); + } + function dispatch(agentAddress: string, event: InferenceEvent): void { const collector = collectors.get(agentAddress); if (collector === undefined) { @@ -128,27 +150,38 @@ export function createEventCollectorRegistry( event.type === "reactor.done" || (event.type === "reactor.error" && event.data.fatal); - collector - .onEvent(event) - .catch((err: unknown) => { + enqueue( + agentAddress, + () => collector.onEvent(event), + (err: unknown) => { log.warn`Failed to persist event ${event.type} seq=${String(event.seq)} for ${agentAddress}: ${err instanceof Error ? err.message : String(err)}`; - }) - .finally(() => { - if (isTerminal) { + }, + () => { + if (isTerminal && collectors.get(agentAddress) === collector) { removeCollector(agentAddress); } - }); + }, + ); } function abandon(agentAddress: string): void { const collector = collectors.get(agentAddress); if (collector === undefined) return; - collector.abandon().catch((err: unknown) => { - log.warn`Failed to abandon collector for ${agentAddress}: ${err instanceof Error ? err.message : String(err)}`; - }); - + // Stop NEW dispatches immediately; the queued closures keep their own + // reference so already-queued events still drain before the abandon runs. removeCollector(agentAddress); + + // Chain the abandon onto the tail so it runs AFTER any queued onEvents + // instead of racing them. Otherwise a queued beginTurn could create a + // fresh `running` turn row after the collector was finalized, orphaning it. + enqueue( + agentAddress, + () => collector.abandon(), + (err: unknown) => { + log.warn`Failed to abandon collector for ${agentAddress}: ${err instanceof Error ? err.message : String(err)}`; + }, + ); } function has(agentAddress: string): boolean { diff --git a/vendor/intx/hub-sessions/src/event-collector.test.ts b/vendor/intx/hub-sessions/src/event-collector.test.ts index 384a77e4c..f08d7f383 100644 --- a/vendor/intx/hub-sessions/src/event-collector.test.ts +++ b/vendor/intx/hub-sessions/src/event-collector.test.ts @@ -4,26 +4,12 @@ import { encodeToolName } from "@intx/inference"; import type { InferenceEvent } from "@intx/types/runtime"; import type { DB } from "@intx/db"; -import { createEventCollector, type UsageForwarded } from "./event-collector"; +import { createEventCollector } from "./event-collector"; import { MALFORMED_TOOL_NAME } from "./sanitize-tool-name"; -// Minimal stand-in for the two drizzle call chains event-collector.ts -// drives (insert().values() and update().set().where()) — enough to -// exercise the inference.usage forwarding without a real database. -function createFakeDb(): DB["db"] { - const chain = { - values: () => Promise.resolve(), - set: () => chain, - where: () => Promise.resolve(), - }; - return { - insert: () => chain, - update: () => chain, - } as unknown as DB["db"]; -} - -// Same shape as createFakeDb, but records every turnPart insert so a test -// can inspect what actually landed in "history". +// Stand-in for the drizzle call chains event-collector.ts drives, recording +// every turnPart insert so a test can inspect what actually landed in +// "history". function createRecordingDb(): { db: DB["db"]; insertedParts: { type: string; metadata?: Record }[]; @@ -55,89 +41,6 @@ function createRecordingDb(): { return { db, insertedParts }; } -describe("createEventCollector inference.usage forwarding", () => { - test("forwards turnId, provider, model, and usage to onUsage — CL-5879 kill-date 2026-09-05", async () => { - const forwarded: UsageForwarded[] = []; - const collector = createEventCollector({ - db: createFakeDb(), - sessionId: "session-1", - runId: "run-1", - tenantId: "tenant-acme", - onUsage: (usage) => forwarded.push(usage), - }); - - await collector.onEvent({ - type: "inference.start", - seq: 1, - data: { model: "claude-sonnet" }, - } as InferenceEvent); - - const turnId = collector.getCurrentTurnId(); - if (turnId === null) throw new Error("expected an active turn"); - - await collector.onEvent({ - type: "inference.usage", - seq: 2, - data: { - usage: { - input: 100, - output: 50, - cacheRead: 0, - cacheWrite: 0, - thinking: 0, - }, - source: { - sourceId: "source-1", - provider: "anthropic", - model: "claude-sonnet", - }, - }, - } as InferenceEvent); - - expect(forwarded).toHaveLength(1); - expect(forwarded[0]).toEqual({ - turnId, - provider: "anthropic", - model: "claude-sonnet", - usage: { - input: 100, - output: 50, - cacheRead: 0, - cacheWrite: 0, - thinking: 0, - }, - }); - }); - - test("drops inference.usage with no active turn", async () => { - const forwarded: UsageForwarded[] = []; - const collector = createEventCollector({ - db: createFakeDb(), - sessionId: "session-1", - runId: "run-1", - tenantId: "tenant-acme", - onUsage: (usage) => forwarded.push(usage), - }); - - await collector.onEvent({ - type: "inference.usage", - seq: 1, - data: { - usage: { - input: 1, - output: 1, - cacheRead: 0, - cacheWrite: 0, - thinking: 0, - }, - source: { sourceId: "s", provider: "anthropic", model: "m" }, - }, - } as InferenceEvent); - - expect(forwarded).toHaveLength(0); - }); -}); - describe("CL-6478: malformed tool-call names never wedge the room", () => { const OPENAI_LIMIT = { provider: "openai", maxLength: 64 } as const; @@ -179,7 +82,8 @@ describe("CL-6478: malformed tool-call names never wedge the room", () => { const toolCallPart = insertedParts.find( (part) => part.metadata?.kind === "call", ); - if (toolCallPart === undefined) throw new Error("expected a tool call part"); + if (toolCallPart === undefined) + throw new Error("expected a tool call part"); const persistedName = toolCallPart.metadata?.name; // The malformed name was never written to history as-is. diff --git a/vendor/intx/hub-sessions/src/event-collector.ts b/vendor/intx/hub-sessions/src/event-collector.ts index 81dbdbf20..e99aaceb0 100644 --- a/vendor/intx/hub-sessions/src/event-collector.ts +++ b/vendor/intx/hub-sessions/src/event-collector.ts @@ -12,6 +12,7 @@ import type { InferenceEvent, ContentBlock, TokenUsage, + LastCycleSource, } from "@intx/types/runtime"; import { type DB, parseTurnPartType } from "@intx/db"; @@ -39,11 +40,13 @@ export type TurnFinalized = { toolErrors: { name: string; content: string }[]; }; -// Forwarded on every `inference.usage` event so product-side consumers -// (see @corbits/insights' createUsageSink) can persist tokens without the -// collector's private turn-accumulation state. Deliberately as narrow as -// `TokenUsage` + the turn/source identity needed to attribute it. -export type UsageForwarded = { +// Per-turn token usage, emitted once when a turn finalizes. Carries full run +// identity so a host accounting sink never has to join against a separate +// address-to-tenant map (which would go stale at run teardown). +export type TurnUsage = { + tenantId: string; + sessionId: string; + runId: string; turnId: string; provider: string; model: string; @@ -64,7 +67,7 @@ export type EventCollectorConfig = { runId: string; tenantId: string; onTurnFinalized?: (turn: TurnFinalized) => void; - onUsage?: (usage: UsageForwarded) => void; + onUsage?: (usage: TurnUsage) => void; }; export function createEventCollector( @@ -81,11 +84,6 @@ export function createEventCollector( // after the turn commits but before the collector is removed. let lastTurnId: string | null = null; let ordinal = 0; - // Model of the currently open turn, stamped by beginTurn. Read by - // inference.usage forwarding to attribute tokens without re-deriving it - // from the event's own `source.model` (both agree; this is the turn's - // system of record). - let currentModel = "unknown"; // Prevents double-finalization when reactor.done and abandon() race. let finalized = false; // Set when inference.error fires so connector.reply knows to persist its @@ -113,31 +111,14 @@ export function createEventCollector( let accumulatedToolCalls: TurnToolCall[] = []; // Tool results that reported isError, accumulated for TurnFinalized. let accumulatedToolErrors: { name: string; content: string }[] = []; - - // Serializes event processing. Callers fire onEvent without awaiting - // (the registry's dispatch is deliberately fire-and-forget so it never - // blocks the websocket message loop), so without this chain two events - // interleave across their DB awaits: a connector.reply's finalize can - // null currentTurnId while inference.done is still inserting parts - // (dropping them as "no active turn"), and a finalize processed while - // the next inference.start's beginTurn is mid-insert marks the NEW - // turn finalized, leaving its row "running" forever (CL-6379). Each - // event fully settles before the next begins, restoring wire order. - let eventTail: Promise = Promise.resolve(); - - function enqueue(work: () => Promise): Promise { - const run = eventTail.then(work); - // A rejected event must not wedge every later event; the caller - // still observes the rejection through the returned promise. - eventTail = run.catch(() => undefined); - return run; - } - - function onEvent(event: InferenceEvent): Promise { - return enqueue(() => processEvent(event)); - } - - async function processEvent(event: InferenceEvent): Promise { + // Final cumulative token usage for the current turn. `inference.usage` events + // carry a running cumulative total and fire multiple times per step, so these + // are OVERWRITTEN (not summed); the last value before finalize is the + // authoritative per-turn total. Both reset on each new turn. + let turnUsage: TokenUsage | null = null; + let turnSource: LastCycleSource | null = null; + + async function onEvent(event: InferenceEvent): Promise { switch (event.type) { case "inference.start": await beginTurn(event.data.model); @@ -149,16 +130,8 @@ export function createEventCollector( case "inference.done": await handleInferenceDone(event.data.turn.content); streamingText = ""; - break; - case "inference.usage": - if (currentTurnId !== null && onUsage) { - onUsage({ - turnId: currentTurnId, - provider: event.data.source.provider, - model: currentModel, - usage: event.data.usage, - }); - } + turnUsage = event.data.usage; + turnSource = event.data.source; break; case "tool.done": { const callId = event.data.result.callId; @@ -243,10 +216,15 @@ export function createEventCollector( }); } break; + case "inference.usage": + // Cumulative running total that fires several times per step; overwrite + // so the last value before finalize is the authoritative per-turn + // usage, emitted once via onUsage from finalizeTurn. + turnUsage = event.data.usage; + turnSource = event.data.source; + break; default: - // reactor.start, streaming deltas, and other events are not - // persisted. inference.usage is handled above (forwarded, not - // persisted here — see UsageForwarded). + // reactor.start, streaming deltas, and other events are not persisted. break; } } @@ -260,7 +238,6 @@ export function createEventCollector( currentTurnId = generateId("inferenceTurn"); lastTurnId = currentTurnId; - currentModel = model; ordinal = 0; finalized = false; pendingError = false; @@ -272,6 +249,8 @@ export function createEventCollector( callArgs.clear(); accumulatedToolCalls = []; accumulatedToolErrors = []; + turnUsage = null; + turnSource = null; await db.insert(inferenceTurn).values({ id: currentTurnId, @@ -436,20 +415,33 @@ export function createEventCollector( }); } + // Report per-turn usage once, alongside the finalize notify. Gated on + // turnUsage being present so a turn that ran no inference emits nothing. + // NOTE: abandon() finalizes with notify=false, so a turn abandoned + // mid-step (e.g. a sidecar disconnect) reports no usage even if the + // provider already billed input tokens -- an accepted gap until durable + // usage persistence lands on the turn row. + if (notify && onUsage && turnUsage !== null && turnSource !== null) { + onUsage({ + tenantId, + sessionId, + runId, + turnId, + provider: turnSource.provider, + model: turnSource.model, + usage: turnUsage, + }); + } + currentTurnId = null; } - function abandon(): Promise { - // Chained behind any in-flight events so an abandon issued while a - // turn's events are still persisting closes the turn they produce, - // not a half-processed intermediate state. - return enqueue(async () => { - if (currentTurnId === null || finalized) return; + async function abandon(): Promise { + if (currentTurnId === null || finalized) return; - log.warn`Abandoning running turn ${currentTurnId} for session ${sessionId}`; + log.warn`Abandoning running turn ${currentTurnId} for session ${sessionId}`; - await finalizeTurn("failed", false, false); - }); + await finalizeTurn("failed", false, false); } async function insertPart( diff --git a/vendor/intx/hub-sessions/src/index.ts b/vendor/intx/hub-sessions/src/index.ts index 8c0173640..cc85482f8 100644 --- a/vendor/intx/hub-sessions/src/index.ts +++ b/vendor/intx/hub-sessions/src/index.ts @@ -8,16 +8,16 @@ export { SessionLaunchError, bridgeOrchestratorDeployContent, deployCodeSourcedWorkflow, - deployAdoptedCodeSourcedWorkflow, type SessionService, type DeployWorkflowDefinitionResult, type DeployWorkflowFromSourceParams, type DeployPreparedCodeSourcedWorkflowParams, type InstallAndApproveWorkflowSourceParams, type PreparedWorkflowDeployer, + type DeployCodeSourcedWorkflowArgs, + deployAdoptedCodeSourcedWorkflow, type AdoptingWorkflowDeployer, type DeployAdoptedWorkflowFromSourceParams, - type DeployCodeSourcedWorkflowArgs, } from "./session-service"; export { installAndApproveWorkflowDefinition, @@ -159,9 +159,9 @@ export { dequeueToProcessing, readProcessingEntry, markConsumed, - readOwnedMessageIds, + classifyTerminalEvent, + scanRunsForBoot, readCommittedWorkflowRunLifecycle, - readCommittedWorkflowRunTerminalStatus, readWorkflowRunLifecycle, replayProcessingToInbox, WORKFLOW_RUN_GITIGNORE_PATH, @@ -194,6 +194,7 @@ export { type WorkflowRunSidecarPrincipal, type WorkflowRunWorkflowProcessPrincipal, type WorkflowRunSupervisorPrincipal, + readCommittedWorkflowRunTerminalStatus, } from "./workflow-run-kind"; export { restoreWorkflowRunToAllocation, diff --git a/vendor/intx/hub-sessions/src/session-service.test.ts b/vendor/intx/hub-sessions/src/session-service.test.ts index 8aa38078b..ade9e6062 100644 --- a/vendor/intx/hub-sessions/src/session-service.test.ts +++ b/vendor/intx/hub-sessions/src/session-service.test.ts @@ -16,10 +16,8 @@ import { describe, expect, test } from "bun:test"; import { deployAdoptedCodeSourcedWorkflow, - deployCodeSourcedWorkflow, type DeployCodeSourcedWorkflowArgs, } from "./session-service"; -import { DeployFrameNotSentError } from "./ws/sidecar-handler"; const TENANT = "tnt_adopt"; const ANCHOR_RUN_ID = "run_adopted_anchor"; @@ -175,184 +173,3 @@ describe("deployAdoptedCodeSourcedWorkflow", () => { expect(db.inserts).toBe(0); }); }); - -// CL-6388: the SHARED code-sourced front must persist its anchor row BEFORE -// the deploy frame goes out. The frame spawns the deployment's child, whose -// first `refs/heads/events` pack push races the deploy ack back to the hub -- -// and `receiveWorkflowRunPack` fails closed (`path_violation`) on a missing -// anchor row, so an insert-after-ack ordering rejects every fresh -// deployment's first events pack. These tests pin the ordering contract: -// insert, then frame, then the acked-key stamp; a failed frame removes the -// pre-inserted row so a retried deploy does not collide on the primary key. -type OrderedFakeDb = { - handle: DeployCodeSourcedWorkflowArgs["db"]; - events: string[]; - insertedValues: Record[]; - updatedSets: Record[]; -}; - -function orderedFakeDb(): OrderedFakeDb { - const state: OrderedFakeDb = { - handle: undefined as unknown as DeployCodeSourcedWorkflowArgs["db"], - events: [], - insertedValues: [], - updatedSets: [], - }; - const handle = { - query: { - workflowDefinition: { - findFirst: () => Promise.resolve({ id: DEFINITION_ID }), - }, - workflowRun: { - findFirst: () => Promise.resolve(undefined), - }, - }, - insert: () => ({ - values: (values: Record) => { - state.events.push("insert"); - state.insertedValues.push(values); - return Promise.resolve(undefined); - }, - }), - update: () => ({ - set: (values: Record) => ({ - where: () => { - state.events.push("update"); - state.updatedSets.push(values); - return Promise.resolve(undefined); - }, - }), - }), - delete: () => ({ - where: () => { - state.events.push("delete"); - return Promise.resolve(undefined); - }, - }), - }; - state.handle = handle as unknown as DeployCodeSourcedWorkflowArgs["db"]; - return state; -} - -function sharedDeployArgs( - db: OrderedFakeDb, - sendAgentDeploy: (agentAddress: string) => Promise<{ publicKey: string }>, -): DeployCodeSourcedWorkflowArgs { - const projection = { - id: "wf_shared", - triggers: [{ type: "manual" }], - stepOrder: [], - steps: {}, - }; - const args = { - approved: { - approval: { - ok: true, - definitionId: DEFINITION_ID, - approvedWireHash: "sha256:frozen", - approvedGrants: new Set(), - projection, - }, - projection, - closure: { entries: [] }, - }, - sidecarRouter: { - sendAgentDeploy: (agentAddress: string) => sendAgentDeploy(agentAddress), - }, - agentAddress: `${ANCHOR_RUN_ID}@${DEPLOYMENT_DOMAIN}`, - config: { sources: [], defaultSource: "default", principalId: "prn_x" }, - sources: {}, - db: db.handle, - tenantId: TENANT, - anchorRunId: ANCHOR_RUN_ID, - deploymentDomain: DEPLOYMENT_DOMAIN, - source: { kind: "registry", registry: "npm" }, - }; - return args as unknown as DeployCodeSourcedWorkflowArgs; -} - -describe("deployCodeSourcedWorkflow (CL-6388)", () => { - test("persists the anchor run before the deploy frame is emitted", async () => { - const db = orderedFakeDb(); - - const result = await deployCodeSourcedWorkflow( - sharedDeployArgs(db, () => { - db.events.push("frame"); - return Promise.resolve({ publicKey: SUPERVISOR_KEY }); - }), - ); - - expect(result.publicKey).toBe(SUPERVISOR_KEY); - expect(db.events).toEqual(["insert", "frame", "update"]); - expect(db.insertedValues[0]).toMatchObject({ - id: ANCHOR_RUN_ID, - anchorRunId: ANCHOR_RUN_ID, - tenantId: TENANT, - definitionId: DEFINITION_ID, - address: `${ANCHOR_RUN_ID}@${DEPLOYMENT_DOMAIN}`, - status: "deployed", - }); - }); - - test("stamps the acked supervisor key onto the anchor after the frame", async () => { - const db = orderedFakeDb(); - - await deployCodeSourcedWorkflow( - sharedDeployArgs(db, () => { - db.events.push("frame"); - return Promise.resolve({ publicKey: SUPERVISOR_KEY }); - }), - ); - - // The key is only known from the ack, so the pre-frame insert cannot - // carry it; the post-frame stamp is where it lands. - expect(db.insertedValues[0]).toMatchObject({ publicKey: null }); - expect(db.updatedSets).toEqual([{ publicKey: SUPERVISOR_KEY }]); - }); - - // CL-6395: a `DeployFrameNotSentError` is proof the `agent.deploy` frame - // never left the hub (e.g. no sidecar available, or the send itself threw - // synchronously) -- no child could have spawned against this anchor, so - // the pre-inserted row is safe to remove. - test("removes the pre-inserted anchor when the frame provably was not sent", async () => { - const db = orderedFakeDb(); - - await expect( - deployCodeSourcedWorkflow( - sharedDeployArgs(db, () => { - db.events.push("frame"); - return Promise.reject( - new DeployFrameNotSentError("No sidecar available"), - ); - }), - ), - ).rejects.toThrow(/No sidecar available/); - - expect(db.events).toEqual(["insert", "frame", "delete"]); - expect(db.updatedSets).toHaveLength(0); - }); - - // CL-6395: an ack-timeout or socket-drop rejection is raised AFTER the - // frame already went out -- the sidecar may have already spawned the - // deployment's child against this anchor. Deleting the row here would - // permanently strand that child on the missing-anchor `path_violation` - // path with no grants, so the row must survive and the failure is logged - // as a reconciliation signal instead. - test("keeps the pre-inserted anchor when the frame failure is post-send", async () => { - const db = orderedFakeDb(); - - await expect( - deployCodeSourcedWorkflow( - sharedDeployArgs(db, () => { - db.events.push("frame"); - return Promise.reject( - new Error('Deploy of "run_adopted_anchor@runs.example.test" timed out after 30000ms'), - ); - }), - ), - ).rejects.toThrow(/timed out/); - - expect(db.events).toEqual(["insert", "frame"]); - expect(db.updatedSets).toHaveLength(0); - }); -}); diff --git a/vendor/intx/hub-sessions/src/session-service.ts b/vendor/intx/hub-sessions/src/session-service.ts index edd26e9c5..455bcf79a 100644 --- a/vendor/intx/hub-sessions/src/session-service.ts +++ b/vendor/intx/hub-sessions/src/session-service.ts @@ -1,5 +1,5 @@ import { type } from "arktype"; -import { and, eq } from "drizzle-orm"; +import { and, eq, isNull } from "drizzle-orm"; import { getLogger } from "@intx/log"; import { @@ -58,6 +58,7 @@ import { buildInertProjectionStepSources, deriveRunAddress, enumerateInertOnTriggerBodies, + inertLoopBody, pickStepInferenceSource, WorkflowDefinitionInvalidError, type DeployContent as OrchestratorDeployContent, @@ -69,12 +70,12 @@ import { type Asset, type AssetService, } from "./asset-service"; -import { - DeployFrameNotSentError, - type AllocatedSidecarTarget, - type SidecarAllocationRouter, - type SidecarRouter, +import type { + AllocatedSidecarTarget, + SidecarAllocationRouter, + SidecarRouter, } from "./ws/sidecar-handler"; +import { isDeployFrameFailure } from "./ws/sidecar-handler"; import type { Principal, RepoId, RepoKind } from "./repo-store"; import { buildSourceAssetMounts, @@ -646,8 +647,7 @@ export type DeployCodeSourcedAssetArgs = DeployCodeSourcedCommonArgs & { }; export type DeployCodeSourcedWorkflowArgs = - | DeployCodeSourcedRegistryArgs - | DeployCodeSourcedAssetArgs; + DeployCodeSourcedRegistryArgs | DeployCodeSourcedAssetArgs; function isAssetDeployArgs( args: DeployCodeSourcedWorkflowArgs, @@ -676,19 +676,22 @@ function isAssetDeployArgs( * A gate outcome that did not approve cannot deploy: an unapproved `approval` * fails closed here rather than shipping an unfrozen definition. * - * This emits the source-ref deploy frame ONLY -- it does NOT write the anchor - * `workflow_run` row. `deployCodeSourcedWorkflow` wraps it with the shared-path - * INSERT; the prepared exclusive path wraps it with an UPDATE-under-lock of the - * anchor row that already exists from prepare time. It returns the frozen - * definition id so each wrapper writes the same content-addressed identity the - * gate persisted. + * This does the READ-ONLY preparation ONLY: it runs the guards, resolves + * credential material, pins the body sources, and builds the asset mounts, then + * returns the frozen definition id and the assembled send args. It emits NO + * frame and writes NO row, so it has no side effect to unwind. The shared path + * (`deployCodeSourcedWorkflow`) sequences prepare -> INSERT anchor -> emit so + * the anchor is visible before the frame spawns the child; `emitSourceRefDeployFrame` + * composes prepare -> emit for the prepared exclusive path, whose anchor row + * already exists from prepare time. It returns the frozen definition id so each + * caller writes the same content-addressed identity the gate persisted. */ -async function emitSourceRefDeployFrame( +async function prepareSourceRefDeploy( args: DeployCodeSourcedWorkflowArgs & { allocationTarget?: AllocatedSidecarTarget; sidecarAllocationRouter?: SidecarAllocationRouter; }, -): Promise<{ publicKey: string; definitionId: string }> { +): Promise<{ definitionId: string; sendArgs: SendMultiStepDeployFrameArgs }> { const { approval, projection, closure } = args.approved; if (!approval.ok) { throw new Error( @@ -788,6 +791,19 @@ async function emitSourceRefDeployFrame( enumerateInertOnTriggerBodies(projection).map(async (body) => { const sources: Record = {}; for (const bodyStepId of body.definition.stepOrder) { + // A loop nested inside an onTrigger body is not yet supported: this + // per-body pin does not recurse into the loop's own body, so the + // loop-body steps' inference sources would be unpinned and the child + // would fail loud at the first iteration. Reject at deploy instead of + // shipping that latent crash. Top-level loops ARE pinned (the source + // pin recurses into their bodies); this gap is only the + // loop-in-onTrigger-body combination, tracked as a follow-on. + if (inertLoopBody(body.definition.steps[bodyStepId]) !== null) { + throw new WorkflowDefinitionInvalidError( + body.ref, + `loop step ${bodyStepId} is nested inside an onTrigger body, which is not yet supported: its body steps' inference sources are not pinned. Move the loop to the top level.`, + ); + } // Agent-bearing body steps run inference and need a source pinned // through the approval gate. A non-agent body step (sleep, // awaitSignal) declares no preference and runs no inference, so it @@ -837,7 +853,7 @@ async function emitSourceRefDeployFrame( ? await buildSourceAssetMounts(closure, args.resolveAttachment) : []; - const result = await sendMultiStepDeployFrame({ + const sendArgs: SendMultiStepDeployFrameArgs = { lineage: "source-ref", sidecarRouter: args.sidecarRouter, ...(args.sidecarAllocationRouter !== undefined @@ -854,90 +870,170 @@ async function emitSourceRefDeployFrame( ...(credentials !== undefined ? { credentials } : {}), ...(referencedDefinitions.length > 0 ? { referencedDefinitions } : {}), ...(assets.length > 0 ? { assets } : {}), - }); + }; - return { publicKey: result.publicKey, definitionId: approval.definitionId }; + return { definitionId: approval.definitionId, sendArgs }; +} + +/** + * Prepare then emit the source-ref deploy frame, for the prepared exclusive + * path whose anchor `workflow_run` row already exists (inserted at prepare + * time). It emits the frame but does NOT touch the anchor row: the caller + * (`deployPreparedCodeSourcedWorkflow`) stamps the acked key under the + * allocation-ownership lock. On emit failure it throws the raw + * `DeployFrameFailure` verbatim, which the caller's own error handling wraps. + * The shared path does NOT use this wrapper -- it must interleave the anchor + * INSERT between prepare and emit, so it drives `prepareSourceRefDeploy` and + * `sendMultiStepDeployFrame` directly. + */ +async function emitSourceRefDeployFrame( + args: DeployCodeSourcedWorkflowArgs & { + allocationTarget?: AllocatedSidecarTarget; + sidecarAllocationRouter?: SidecarAllocationRouter; + }, +): Promise<{ publicKey: string; definitionId: string }> { + const { definitionId, sendArgs } = await prepareSourceRefDeploy(args); + const result = await sendMultiStepDeployFrame(sendArgs); + return { publicKey: result.publicKey, definitionId }; } /** * The single public composition entrypoint for a SHARED code-sourced (npm) - * deploy: INSERT the deployment's anchor `workflow_run` row, then emit the - * source-ref frame, then stamp the acked supervisor key onto the row. The - * anchor row is the deployment's first-class record that owns its routing - * address and public key. Run-grant materialization keys off this row - * (address + live status), so WITHOUT it no per-run grants (tool, capability, OR - * credential) ever materialize for a source-ref deployment. Born "deployed" - * (live but pre-trigger): the first trigger's materialization flips it to - * "running" via `anchorWithPrincipal`'s guarded update, which a row born - * "running" would skip. Its `anchorRunId` equals its own id, so the anchor - * references itself. The deployer read grant is deferred to the production - * route, which carries the authenticated deployer principal; this stays a - * single insert with no grant row to pair atomically. + * deploy: prepare, INSERT the deployment's anchor `workflow_run` row, THEN emit + * the source-ref frame. The anchor row is the deployment's first-class record + * that owns its routing address and public key. Run-grant materialization keys + * off this row (address + live status), so WITHOUT it no per-run grants (tool, + * capability, OR credential) ever materialize for a source-ref deployment. Born + * "deployed" (live but pre-trigger) with a null public key: the first trigger's + * materialization flips it to "running" via `anchorWithPrincipal`'s guarded + * update, which a row born "running" would skip. Its `anchorRunId` equals its + * own id, so the anchor references itself. The deployer read grant is deferred + * to the production route, which carries the authenticated deployer principal. + * + * ORDERING IS LOAD-BEARING. The anchor row must be committed and visible to the + * pack-receipt connection BEFORE the frame reaches the wire: the frame spawns + * the child, whose first events pack races the ack back, and + * `receiveWorkflowRunPack` fails closed on a missing live anchor. Emitting first + * (the previous order) rejected that first pack and never bootstrapped the log. + * This works because `args.db` is the autocommit handle (`DB["db"]`, which the + * type forbids from being a transaction) and the INSERT is NOT wrapped in a + * transaction with the emit -- so the row is durably visible the instant the + * INSERT statement returns. Do NOT relax `db` to a transaction executor or wrap + * anchor+emit in one transaction to make them atomic: that reopens the race. + * + * On emit failure the anchor row is rolled back or fenced by the `frameSent` + * evidence from the transport. `leakedAgent: false` (safe to fully roll back) is + * the STRONG claim and is made only on positive proof the frame never reached + * the wire (`isDeployFrameFailure && frameSent === false`); every other failure + * -- a sent-but-unacked frame OR any untagged error -- is treated as + * possibly-live: the anchor is fenced `deployed` -> `failed` and the error is + * `leakedAgent: true`. * - * The prepared exclusive path does NOT use this wrapper: its anchor row already - * exists from prepare time, so it wraps `emitSourceRefDeployFrame` with an - * UPDATE-under-allocation-lock instead of this INSERT. + * The prepared exclusive path does NOT use this composition: its anchor row + * already exists from prepare time, so it drives `emitSourceRefDeployFrame` and + * an UPDATE-under-allocation-lock instead. */ export async function deployCodeSourcedWorkflow( args: DeployCodeSourcedWorkflowArgs, ): Promise<{ publicKey: string }> { - const { approval } = args.approved; - if (!approval.ok) { - throw new Error( - `deployCodeSourcedWorkflow: refusing to deploy an unapproved workflow (gate reason: ${approval.reason})`, - ); - } + const { definitionId, sendArgs } = await prepareSourceRefDeploy(args); - // CL-6388 (workbench-local; see VENDORED.md): the anchor row must exist - // BEFORE the deploy frame goes out. The frame spawns the deployment's - // child, whose first `refs/heads/events` pack push races the deploy ack - // back to the hub -- and `receiveWorkflowRunPack` fails closed - // (`path_violation`) on a missing anchor row, so an insert-after-ack - // ordering rejected every fresh deployment's first events pack. The - // prepared and adopted fronts already have their anchor row pre-frame; - // this front now matches them. The supervisor key is only known from the - // ack, so the row is born with a null `publicKey` (which keeps the - // reconnect challenge failing closed until the ack) and the key is - // stamped afterwards. - // - // CL-6395 (workbench-local; see VENDORED.md): the pre-inserted row is - // deleted ONLY when `emitSourceRefDeployFrame` fails with a - // `DeployFrameNotSentError` -- proof the `agent.deploy` frame never left - // the hub, so no child could have spawned against this anchor. Any OTHER - // failure (ack timeout, socket drop, reconnect takeover) is raised AFTER - // the frame was already sent: the spawned child may already be running - // against this anchor, and deleting the row would permanently strand it - // on the missing-anchor `path_violation` path with no grants. That row - // is kept and the failure logged as a reconciliation signal instead. - await args.db.insert(workflowRunTable).values({ - id: args.anchorRunId, - tenantId: args.tenantId, - anchorRunId: args.anchorRunId, - definitionId: approval.definitionId, - address: args.agentAddress, - publicKey: null, - status: "deployed", - createdAt: new Date(), - }); + // INSERT the anchor before the frame. A collision or DB error here spawned + // nothing (no frame went out), so it is a clean, non-leaking failure. + try { + await args.db.insert(workflowRunTable).values({ + id: args.anchorRunId, + tenantId: args.tenantId, + anchorRunId: args.anchorRunId, + definitionId, + address: args.agentAddress, + publicKey: null, + status: "deployed", + createdAt: new Date(), + }); + } catch (cause) { + throw new SessionLaunchError("start", cause, false); + } let publicKey: string; try { - ({ publicKey } = await emitSourceRefDeployFrame(args)); - } catch (error) { - if (error instanceof DeployFrameNotSentError) { - await args.db + const result = await sendMultiStepDeployFrame(sendArgs); + publicKey = result.publicKey; + } catch (cause) { + if (isDeployFrameFailure(cause) && cause.frameSent === false) { + // Positive proof the frame never reached the wire: nothing spawned, so + // fully roll the anchor back. The guard (`deployed`, null key) is a + // tripwire on the `frameSent: false` contract -- a 0-row delete means the + // row advanced or vanished, so the contract lied and a child may be live; + // surface that loudly and refuse to claim it is safe to roll back. + const deleted = await args.db .delete(workflowRunTable) - .where(eq(workflowRunTable.id, args.anchorRunId)); + .where( + and( + eq(workflowRunTable.id, args.anchorRunId), + eq(workflowRunTable.anchorRunId, args.anchorRunId), + eq(workflowRunTable.tenantId, args.tenantId), + eq(workflowRunTable.status, "deployed"), + isNull(workflowRunTable.publicKey), + ), + ) + .returning({ id: workflowRunTable.id }); + if (deleted.length === 0) { + logger.error`anchor-before-frame rollback found no deployed/null-key row for ${args.anchorRunId} after a frameSent:false failure; the never-sent contract was violated and a child may be live`; + throw new SessionLaunchError("start", cause, true); + } + throw new SessionLaunchError("start", cause, false); + } + // A sent-but-unacked frame, OR any untagged/unexpected error: no positive + // proof of a clean send, so treat the agent as possibly-live. Fence the + // anchor `deployed` -> `failed` (guarded so a self-flip to "running" by a + // trigger that already landed is left alone). Do NOT delete: a live child + // needs the anchor to bootstrap. + const flipped = await args.db + .update(workflowRunTable) + .set({ status: "failed" }) + .where( + and( + eq(workflowRunTable.id, args.anchorRunId), + eq(workflowRunTable.anchorRunId, args.anchorRunId), + eq(workflowRunTable.tenantId, args.tenantId), + eq(workflowRunTable.status, "deployed"), + isNull(workflowRunTable.publicKey), + ), + ) + .returning({ id: workflowRunTable.id }); + if (flipped.length === 0) { + // The anchor already advanced past deployed -- a trigger flipped it to + // "running", so the deploy actually succeeded and the run is progressing + // despite the ack failure. Leave it; the leaked-agent disposition still + // holds because the frame was (or may have been) sent. + logger.warn`anchor-before-frame: anchor ${args.anchorRunId} already advanced past deployed on an unacked/failed emit; the agent is live and the run is progressing despite the ack failure`; } else { - logger.error`deployCodeSourcedWorkflow: anchor run ${args.anchorRunId} kept after a post-send deploy-frame failure; the row needs manual reconciliation: ${error instanceof Error ? error.message : String(error)}`; + logger.warn`anchor-before-frame: fenced anchor ${args.anchorRunId} deployed->failed on an unacked/failed emit; the agent may be leaked but the run is dead`; } - throw error; + throw new SessionLaunchError("start", cause, true); } - await args.db + // Emit succeeded: stamp the acked key. No status guard -- the key is a fact + // regardless of whether the pack-ack race already flipped the row to + // "running", and skipping the stamp there would strand a live run with a null + // key. A 0-row update is an anomaly (nothing should remove a deployed anchor + // on the success path), but the deploy succeeded, so log it rather than + // failing a live run. + const stamped = await args.db .update(workflowRunTable) .set({ publicKey }) - .where(eq(workflowRunTable.id, args.anchorRunId)); + .where( + and( + eq(workflowRunTable.id, args.anchorRunId), + eq(workflowRunTable.anchorRunId, args.anchorRunId), + eq(workflowRunTable.tenantId, args.tenantId), + ), + ) + .returning({ id: workflowRunTable.id }); + if (stamped.length === 0) { + logger.error`anchor-before-frame: anchor ${args.anchorRunId} vanished before its public key could be stamped on a successful deploy`; + } return { publicKey }; } diff --git a/vendor/intx/hub-sessions/src/substrate.ts b/vendor/intx/hub-sessions/src/substrate.ts index aaa40cff5..7d29aa1d3 100644 --- a/vendor/intx/hub-sessions/src/substrate.ts +++ b/vendor/intx/hub-sessions/src/substrate.ts @@ -16,12 +16,13 @@ // module. The barrel still re-exports all of these for hub-side consumers. export { + classifyTerminalEvent, DEFAULT_CONSUMED_RETENTION_MS, dequeueToProcessing, enqueueInbox, markConsumed, parseEventSeq, - readOwnedMessageIds, + scanRunsForBoot, readCommittedWorkflowRunLifecycle, readWorkflowRunLifecycle, readProcessingEntry, @@ -29,8 +30,10 @@ export { requireEventSeq, StaleInboxEnqueueError, WORKFLOW_RUN_AGENT_STATE_PREFIX, + WORKFLOW_RUN_PARTS_DIR, WORKFLOW_RUN_EVENTS_DIR, WORKFLOW_RUN_RUNS_PREFIX, + MAX_MAIL_PART_PATH_COMPONENT_BYTES, } from "./workflow-run-kind"; export type { WorkflowRunSupervisorPrincipal, diff --git a/vendor/intx/hub-sessions/src/workflow-run-kind.ts b/vendor/intx/hub-sessions/src/workflow-run-kind.ts index 5548a1d30..34067c72e 100644 --- a/vendor/intx/hub-sessions/src/workflow-run-kind.ts +++ b/vendor/intx/hub-sessions/src/workflow-run-kind.ts @@ -222,6 +222,21 @@ export const WORKFLOW_RUN_INBOX_DIR = "inbox"; export const WORKFLOW_RUN_PROCESSING_DIR = "processing"; export const WORKFLOW_RUN_CONSUMED_DIR = "consumed"; +/** + * Per-run inbound mail-part subtree. Non-text inbound mail content + * (image/audio/video/document mail parts) is committed here as real + * files rather than inlined into the JSON event log, whose serialization + * boundary would corrupt binary bytes. The layout is + * `runs//parts//-`: one + * directory per inbound message (so a long-lived run's successive turns + * never collide), and one file per mail part carrying its verbatim + * bytes. The workflow-host ingest writes the bytes and records a + * lightweight `{ name, contentType, ref }` reference into the run's + * trigger / signal payload; the step invoker reads the bytes back at + * `agent.send` time. Files are immutable once written, like `blobs/`. + */ +export const WORKFLOW_RUN_PARTS_DIR = "parts"; + /** * Filename of the per-address retention watermark blob, a direct child * of `addresses//` (a file, not a directory). Carries the @@ -281,6 +296,39 @@ export const DEFAULT_CONSUMED_RETENTION_MS = 24 * 60 * 60 * 1000; */ export const WORKFLOW_RUN_AGENT_STATE_PREFIX = "agent-state"; +/** + * Conversational-mailbox subtree for the warm single-step agent. The + * substrate mailbox backing commits the agent's durable inbox under + * `mailbox/INBOX/` so the full message history replicates to the hub + * alongside the run state. The layout is: + * + * - `mailbox/INBOX/index.json` — the mailbox index. MUTABLE: the + * backing rewrites it on every flush, so it is exempt from the + * retained-blob byte-equality walk the `.eml` blobs are subject + * to. It must nonetheless PERSIST once it existed: dropping it drops + * the whole mailbox and resets uidValidity/uidNext on the next open. + * - `mailbox/INBOX/.eml` — one message per file, carrying the + * raw signed message bytes. `` is a decimal integer >= 1. A + * RETAINED `.eml` (present in both prior and prospective) is + * opaque and IMMUTABLE: it must reappear byte-identically, exactly + * like `runs//blobs/`. A prior `.eml` may be ABSENT from + * the prospective tree — that is the warm agent expunging a message. + * The raw bytes are not lost: they stay reachable through the parent + * commit, and a `workflow-run` repo's objects are never GC'd (its + * kind is off the GC allow-list), so an expunged message survives in + * history for the life of the run repo. The audit trail rests on + * "these objects are never pruned", not on the live tree being + * monotonic. + * + * The only entries permitted under `mailbox/` are the `INBOX/` + * directory; the only entries permitted under `mailbox/INBOX/` are + * `index.json` and `.eml` message files. Anything else fails the + * push. + */ +export const WORKFLOW_RUN_MAILBOX_PREFIX = "mailbox"; +export const WORKFLOW_RUN_MAILBOX_INBOX_DIR = "INBOX"; +export const WORKFLOW_RUN_MAILBOX_INDEX_FILE = "index.json"; + /** * Allowed top-level entries in the prospective tree. Anything else * fails the push. `control/` has no v1 use and stays absent. @@ -289,9 +337,18 @@ const ALLOWED_TOP_LEVEL = new Set([ WORKFLOW_RUN_RUNS_PREFIX, WORKFLOW_RUN_ADDRESSES_PREFIX, WORKFLOW_RUN_AGENT_STATE_PREFIX, + WORKFLOW_RUN_MAILBOX_PREFIX, WORKFLOW_RUN_GITIGNORE_PATH, ]); +/** + * Per-message filename shape for the `mailbox/INBOX/` subtree: + * `.eml`, where `` is a decimal integer >= 1 (no leading zero, + * never `0`). Pins the shape so a malformed message name fails the push + * at the boundary rather than landing silently. + */ +const MAILBOX_EML_FILENAME_RE = /^[1-9][0-9]*\.eml$/; + const CLAIM_CHECK_SUBDIRS = new Set([ WORKFLOW_RUN_INBOX_DIR, WORKFLOW_RUN_PROCESSING_DIR, @@ -347,6 +404,33 @@ export function requireEventSeq(filename: string, context: string): number { */ const BLOB_FILENAME_RE = /^[0-9a-f]{64}$/; +/** + * Per-mail-part filename shape for the + * `runs//parts//` subtree: + * `-`, where `index` is the decimal position of the + * mail part within its inbound message and `name` is a non-empty + * (possibly sanitized) filename. The workflow-host ingest owns the exact + * encoding and sanitizes untrusted names to satisfy this shape; this regex + * pins the shape so a malformed name fails the push rather than landing + * silently. The bytes themselves are opaque and immutable, exactly like + * `blobs/`. + */ +const PART_FILENAME_RE = /^(0|[1-9][0-9]*)-(.+)$/; + +/** + * Maximum byte length of a single mail part path component (the + * URL-encoded message segment, and each `-` filename). + * messageIds and mail part names arrive from untrusted inbound mail; + * an over-long RFC 5322 message-id URL-encodes past the filesystem's + * 255-byte component limit and would otherwise fail at disk-write time, + * downstream of validation. Reject it at the boundary instead. + */ +export const MAX_MAIL_PART_PATH_COMPONENT_BYTES = 255; + +function mailPartComponentByteLength(component: string): number { + return new TextEncoder().encode(component).length; +} + /** * Entries the kind handler accepts under `runs//`. The `events/` * subtree carries the append-only event log; the `blobs/` subtree carries @@ -365,6 +449,10 @@ const RUN_DIR_ALLOWED_CHILDREN = new Set([ // files into one combined file by a compaction commit. WORKFLOW_RUN_EVENTS_FILE, WORKFLOW_RUN_GRANTS_FILE, + // Inbound-mail mail part bytes committed as real files. See + // WORKFLOW_RUN_PARTS_DIR; validated by enumerateRunParts and + // held immutable by the same prior-tree byte-equality walk as blobs. + WORKFLOW_RUN_PARTS_DIR, ]); /** @@ -464,12 +552,11 @@ export type WatermarkEnvelope = typeof WatermarkEnvelope.infer; * stay in sync with the canonical runtime definition: * if the runtime adds or removes a terminal run phase, update this map too. * Drift silently reopens the restore-time double-driver collision that - * `readOwnedMessageIds` (below) exists to prevent. + * `scanRunsForBoot` (below) exists to prevent. */ -const TERMINAL_EVENT_STATUS: ReadonlyMap< - string, - "completed" | "failed" | "cancelled" -> = new Map([ +type TerminalRunStatus = "completed" | "failed" | "cancelled"; + +const TERMINAL_EVENT_STATUS: ReadonlyMap = new Map([ ["RunCompleted", "completed"], ["RunFailed", "failed"], ["RunCancelled", "cancelled"], @@ -482,6 +569,34 @@ const TERMINAL_EVENT_STATUS: ReadonlyMap< */ const TERMINAL_EVENT_TYPES = new Set(TERMINAL_EVENT_STATUS.keys()); +/** + * Classify a workflow-run event type against the terminal-status vocabulary. + * `TERMINAL_EVENT_STATUS` is the sole authority (see above), so a type absent + * from it is by definition not terminal: no separate membership set is + * consulted and no "unmapped terminal type" case can arise. + */ +export function classifyTerminalEvent( + eventType: string, +): { terminal: true; status: TerminalRunStatus } | { terminal: false } { + const status = TERMINAL_EVENT_STATUS.get(eventType); + return status === undefined + ? { terminal: false } + : { terminal: true, status }; +} + +/** + * True when a blob is absent from the prior tree -- i.e. this commit is the + * one that authored it. Consumers act only on a newly-added blob; a blob + * already present in the prior tree was carried forward unchanged by a + * compaction commit and must not be re-acted upon. + */ +async function blobIsNewlyAdded( + blobPath: string, + priorReadBlob: (path: string) => Promise, +): Promise { + return (await priorReadBlob(blobPath)) === null; +} + /** * Recognised CancelRequested origins. Mirrors the workflow package's * `CANCEL_ORIGINS` vocabulary; inlined here so the substrate does @@ -635,7 +750,7 @@ async function enumerateEventBlobs( if (offender !== undefined) { return { ok: false, - reason: `run directory ${runDirPath} contains unexpected entry ${JSON.stringify(offender)}; only "${WORKFLOW_RUN_EVENTS_DIR}", "${WORKFLOW_RUN_BLOBS_DIR}", "${WORKFLOW_RUN_EVENTS_FILE}", and "${WORKFLOW_RUN_GRANTS_FILE}" are allowed`, + reason: `run directory ${runDirPath} contains unexpected entry ${JSON.stringify(offender)}; only "${WORKFLOW_RUN_EVENTS_DIR}", "${WORKFLOW_RUN_BLOBS_DIR}", "${WORKFLOW_RUN_EVENTS_FILE}", "${WORKFLOW_RUN_GRANTS_FILE}", and "${WORKFLOW_RUN_PARTS_DIR}" are allowed`, }; } const hasCombined = runChildren.includes(WORKFLOW_RUN_EVENTS_FILE); @@ -650,16 +765,18 @@ async function enumerateEventBlobs( // by the combined-form path, not this per-event enumeration. if (hasCombined) continue; if (!hasPerEvent) { - // A run dir whose only child is `grants.json` is the legitimate - // pre-first-event window: the hub's `run.grants` frame writes the - // grants ahead of the trigger, so the grants file lands before the - // child emits its first event. Carry it forward untouched -- there is - // no event log to enumerate yet. Any other events-less shape (e.g. a - // bare `blobs/` with no events) remains rejected below. - if ( - runChildren.length === 1 && - runChildren[0] === WORKFLOW_RUN_GRANTS_FILE - ) { + // The pre-first-event window: `grants.json` (the hub's `run.grants` + // frame writes grants ahead of the trigger) and `parts/` (the + // supervisor commits inbound-mail mail part bytes before firing the + // trigger) may both land before the child emits its first event. Carry + // such a run forward untouched -- there is no event log to enumerate + // yet, and the mail parts subtree is validated by its own walk. Any + // other events-less shape (e.g. a bare `blobs/` with no events) remains + // rejected below. + const nonPreEvent = runChildren.filter( + (c) => c !== WORKFLOW_RUN_GRANTS_FILE && c !== WORKFLOW_RUN_PARTS_DIR, + ); + if (nonPreEvent.length === 0) { continue; } return { @@ -938,6 +1055,197 @@ async function enumerateRunBlobs( return { ok: true, blobs: out }; } +type RunPartEntry = { + runId: string; + messageSegment: string; + filename: string; + blobPath: string; +}; + +/** + * Walk every `runs//parts//` directory and + * validate each entry: the `` is a URL-encoded messageId + * that must round-trip cleanly (the same canonical-encoding discipline the + * claim-check `addresses/` subtree enforces) and stay within the path- + * component byte cap, must be a directory rather than a dangling blob, and + * each filename matches the `-` shape within the same cap. The + * `parts/` subdirectory is optional: a run that never received a + * non-text inbound message never produces one. Returns the flat entry list + * so the caller can apply immutability checks against the prior tree, + * exactly as it does for blobs. + */ +async function enumerateRunParts( + listDir: (path: string) => Promise, + scopeRunIds?: ReadonlySet, +): Promise< + { ok: true; parts: RunPartEntry[] } | { ok: false; reason: string } +> { + const out: RunPartEntry[] = []; + // See enumerateRunBlobs: a defined `scopeRunIds` walks only the commit's + // touched runs; an untouched run's mail parts are carried forward + // byte-identical and were validated when written. + const runIds = + scopeRunIds === undefined + ? await listDir(WORKFLOW_RUN_RUNS_PREFIX) + : Array.from(scopeRunIds); + for (const runId of runIds) { + const runDirPath = `${WORKFLOW_RUN_RUNS_PREFIX}/${runId}`; + const runChildren = await listDir(runDirPath); + if (!runChildren.includes(WORKFLOW_RUN_PARTS_DIR)) continue; + const partsDirPath = `${runDirPath}/${WORKFLOW_RUN_PARTS_DIR}`; + const messageSegments = await listDir(partsDirPath); + for (const messageSegment of messageSegments) { + const roundTrip = checkUrlSegmentRoundTrip(messageSegment); + if (!roundTrip.ok) { + return { + ok: false, + reason: `mail part message ${roundTrip.reason} under ${partsDirPath}`, + }; + } + if ( + mailPartComponentByteLength(messageSegment) > + MAX_MAIL_PART_PATH_COMPONENT_BYTES + ) { + return { + ok: false, + reason: `mail part message segment ${JSON.stringify(messageSegment)} under ${partsDirPath} exceeds the ${String(MAX_MAIL_PART_PATH_COMPONENT_BYTES)}-byte path-component limit`, + }; + } + const messageDirPath = `${partsDirPath}/${messageSegment}`; + const filenames = await listDir(messageDirPath); + // A message segment must be a directory carrying at least one + // mail part file. An empty listing means either a dangling blob + // committed directly at `parts/` (the substrate lists a + // blob path as empty, exactly as the agent-state walk detects) or an + // empty directory; both are rejected so untrusted inbound content has + // no silent-accept path. + if (filenames.length === 0) { + return { + ok: false, + reason: `mail part message segment ${JSON.stringify(messageSegment)} under ${partsDirPath} is not a directory carrying mail part files`, + }; + } + for (const filename of filenames) { + if (!PART_FILENAME_RE.test(filename)) { + return { + ok: false, + reason: `mail part filename ${messageDirPath}/${filename} does not match -`, + }; + } + if ( + mailPartComponentByteLength(filename) > + MAX_MAIL_PART_PATH_COMPONENT_BYTES + ) { + return { + ok: false, + reason: `mail part filename ${messageDirPath}/${filename} exceeds the ${String(MAX_MAIL_PART_PATH_COMPONENT_BYTES)}-byte path-component limit`, + }; + } + // Each mail part entry must be a leaf blob, not a nested directory. + // The `-` shape is permissive (`.+`), so a directory + // named e.g. `0-foo` would otherwise pass and admit an arbitrarily + // nested subtree, breaking the one-file-per-mail-part invariant and + // -- on a pack-receive validation with no `listDirOids` -- driving the + // immutability resolver to `readBlob` a tree path and throw. A blob + // lists as empty; a directory lists its children. + const filePath = `${messageDirPath}/${filename}`; + if ((await listDir(filePath)).length > 0) { + return { + ok: false, + reason: `mail part ${filePath} is a directory; each mail part must be a single file`, + }; + } + out.push({ + runId, + messageSegment, + filename, + blobPath: filePath, + }); + } + } + } + return { ok: true, parts: out }; +} + +/** + * Validate the per-run `parts/` subtree's immutability against the + * prior tree, both directions, mirroring the blobs walk. Mail part files + * are write-once. A path present in the prior tree must carry the same git + * blob OID in the prospective tree: unlike `blobs/`, a mail part filename + * is `-` (not a content hash), so the same path can carry + * different bytes -- the OID compare is the load-bearing immutability guard, + * and comparing the OID git already computed avoids re-reading tens of MB of + * mail part bytes on every commit that merely touches the run. A prior path + * must reappear (no deletion): run reclaim drops the whole `runs//` + * subtree outside `validatePush`, so a partial mail part deletion is always + * a violation. Structural shape is enforced by `enumerateRunParts`. + */ +async function validateRunPartsSubtree(args: { + listDir: (path: string) => Promise; + priorListDir: (path: string) => Promise; + readBlob: (path: string) => Promise; + priorReadBlob: (path: string) => Promise; + listDirOids: + ((path: string) => Promise<{ name: string; oid: string }[]>) | undefined; + priorListDirOids: + ((path: string) => Promise<{ name: string; oid: string }[]>) | undefined; + scopeRunIds: ReadonlySet | undefined; +}): Promise { + const prospective = await enumerateRunParts(args.listDir, args.scopeRunIds); + if (!prospective.ok) return prospective; + const prior = await enumerateRunParts(args.priorListDir, args.scopeRunIds); + if (!prior.ok) { + return { + ok: false, + reason: `prior tree's mail parts subtree is structurally invalid: ${prior.reason}`, + }; + } + + const prospectiveOid = makeListingOidResolver( + "prospective", + args.listDirOids, + async (p) => (await git.hashBlob({ object: await args.readBlob(p) })).oid, + ); + const priorOid = makeListingOidResolver( + "prior", + args.priorListDirOids, + async (p) => { + const bytes = await args.priorReadBlob(p); + if (bytes === null) { + throw new Error( + `mail parts: prior entry ${p} was enumerated but its bytes could not be read`, + ); + } + return (await git.hashBlob({ object: bytes })).oid; + }, + ); + + const prospectivePaths = new Set(prospective.parts.map((e) => e.blobPath)); + const priorPaths = new Set(prior.parts.map((e) => e.blobPath)); + + for (const entry of prospective.parts) { + if (!priorPaths.has(entry.blobPath)) continue; // newly added + const [next, before] = await Promise.all([ + prospectiveOid(entry.blobPath), + priorOid(entry.blobPath), + ]); + if (next !== before) { + return { + ok: false, + reason: `mail part ${entry.blobPath} bytes diverge from the prior tree (blob OID ${next} vs ${before}); mail part files are immutable once written`, + }; + } + } + for (const entry of prior.parts) { + if (prospectivePaths.has(entry.blobPath)) continue; + return { + ok: false, + reason: `mail part ${entry.blobPath} present in the prior tree is missing from the prospective tree; mail part files are immutable once written`, + }; + } + return { ok: true }; +} + /** * Enforce blob immutability via prior-tree byte equality. The blob * value itself is opaque bytes (no JSON envelope, no arktype @@ -1070,13 +1378,16 @@ async function checkPriorByteEquality( } /** - * Round-trip an `` segment through decode then - * encode. A divergence means the segment is not the canonical - * encoding of any address, which would leave consumers guessing - * which encoding to use when reading the subtree. Surface as a - * concrete rejection at push time. + * Round-trip a URL-encoded path segment through decode then encode. A + * divergence means the segment is not the canonical encoding of any + * value, which would leave consumers guessing which encoding to use + * when reading the subtree. Surface as a concrete rejection at push + * time. Shared by every subtree that keys a directory by an + * `encodeURIComponent`-encoded identity (`addresses/`, `agent-state/`, + * and the per-run `parts/` message segment); callers prepend + * their own subtree context to the returned reason. */ -function checkAddressSegmentRoundTrip(segment: string): +function checkUrlSegmentRoundTrip(segment: string): | { ok: true; decoded: string; @@ -1091,7 +1402,7 @@ function checkAddressSegmentRoundTrip(segment: string): } catch (cause) { return { ok: false, - reason: `address segment ${JSON.stringify(segment)} is not a valid URL-encoded string: ${ + reason: `segment ${JSON.stringify(segment)} is not a valid URL-encoded string: ${ cause instanceof Error ? cause.message : String(cause) }`, }; @@ -1100,7 +1411,7 @@ function checkAddressSegmentRoundTrip(segment: string): if (reencoded !== segment) { return { ok: false, - reason: `address segment ${JSON.stringify(segment)} does not round-trip URL-encoding (re-encoded as ${JSON.stringify(reencoded)})`, + reason: `segment ${JSON.stringify(segment)} does not round-trip URL-encoding (re-encoded as ${JSON.stringify(reencoded)})`, }; } return { ok: true, decoded }; @@ -1183,8 +1494,10 @@ async function enumerateClaimCheckBlobs( const perAddress = new Map(); const segments = await listDir(WORKFLOW_RUN_ADDRESSES_PREFIX); for (const segment of segments) { - const roundTrip = checkAddressSegmentRoundTrip(segment); - if (!roundTrip.ok) return roundTrip; + const roundTrip = checkUrlSegmentRoundTrip(segment); + if (!roundTrip.ok) { + return { ok: false, reason: `address ${roundTrip.reason}` }; + } const addrDir = `${WORKFLOW_RUN_ADDRESSES_PREFIX}/${segment}`; const children = await listDir(addrDir); for (const child of children) { @@ -1484,8 +1797,7 @@ async function hashConsumedBlobOid(bytes: Uint8Array): Promise { function makeListingOidResolver( sideLabel: string, listDirOids: - | ((path: string) => Promise<{ name: string; oid: string }[]>) - | undefined, + ((path: string) => Promise<{ name: string; oid: string }[]>) | undefined, hashFallback: (blobPath: string) => Promise, ): (blobPath: string) => Promise { const dirOidCache = new Map>(); @@ -1517,8 +1829,7 @@ function makeListingOidResolver( function makePriorConsumedOidResolver( priorReadBlob: (path: string) => Promise, priorListDirOids: - | ((path: string) => Promise<{ name: string; oid: string }[]>) - | undefined, + ((path: string) => Promise<{ name: string; oid: string }[]>) | undefined, ): (blobPath: string) => Promise { return makeListingOidResolver("prior", priorListDirOids, async (blobPath) => { const bytes = await priorReadBlob(blobPath); @@ -1996,11 +2307,11 @@ async function validateAgentStateSubtree( } const segments = await listDir(WORKFLOW_RUN_AGENT_STATE_PREFIX); for (const segment of segments) { - const roundTrip = checkAddressSegmentRoundTrip(segment); + const roundTrip = checkUrlSegmentRoundTrip(segment); if (!roundTrip.ok) { return { ok: false, - reason: `agent-state segment ${JSON.stringify(segment)} does not round-trip URL-encoding; ${roundTrip.reason}`, + reason: `agent-state ${roundTrip.reason}`, }; } // Reject a blob dangling directly at `agent-state/`: every @@ -2020,6 +2331,165 @@ async function validateAgentStateSubtree( return { ok: true }; } +/** + * Enforce mailbox `.eml` immutability via prior-tree byte equality. + * A message blob RETAINED from the prior tree (present in both) must carry + * byte-identical contents in the prospective tree. A prior blob absent + * from the prospective tree is a legal expunge and never reaches here + * (the caller only checks retained paths). Mirrors the blob-immutability + * discipline (`checkBlobPriorByteEquality`) with mailbox-specific wording. + */ +async function checkMailboxEmlPriorByteEquality( + emlPath: string, + readBlob: (path: string) => Promise, + priorReadBlob: (path: string) => Promise, +): Promise { + const prior = await priorReadBlob(emlPath); + if (prior === null) return { ok: true }; + const prospective = await readBlob(emlPath); + if (prior.byteLength !== prospective.byteLength) { + return { + ok: false, + reason: `mailbox message ${emlPath} bytes diverge from the prior tree (lengths ${String(prior.byteLength)} vs ${String(prospective.byteLength)}); a retained mailbox message is immutable`, + }; + } + for (let i = 0; i < prior.byteLength; i++) { + if (prior[i] !== prospective[i]) { + return { + ok: false, + reason: `mailbox message ${emlPath} bytes diverge from the prior tree at offset ${String(i)}; a retained mailbox message is immutable`, + }; + } + } + return { ok: true }; +} + +/** + * Walk the `mailbox/INBOX/` subtree and validate its shape. The only + * entry permitted directly under `mailbox/` is the `INBOX/` directory; + * the only entries permitted under `mailbox/INBOX/` are the mutable + * `index.json` file and `.eml` message files. Each entry must be a + * leaf blob rather than a nested directory (a blob lists as empty, a + * directory lists its children -- the same discrimination the agent-state + * and mail-parts walks use). Returns the flat set of `.eml` blob + * paths so the caller can hold retained messages immutable against the + * prior tree, plus `indexPresent` -- whether `index.json` exists under + * the INBOX -- so the caller can enforce index continuity. An absent + * `mailbox/` subtree lists as empty and contributes no message paths. + */ +async function enumerateMailboxInbox( + listDir: (path: string) => Promise, +): Promise< + | { ok: true; emlPaths: Set; indexPresent: boolean } + | { ok: false; reason: string } +> { + const emlPaths = new Set(); + let indexPresent = false; + const mailboxChildren = await listDir(WORKFLOW_RUN_MAILBOX_PREFIX); + if (mailboxChildren.length === 0) return { ok: true, emlPaths, indexPresent }; + for (const child of mailboxChildren) { + if (child !== WORKFLOW_RUN_MAILBOX_INBOX_DIR) { + return { + ok: false, + reason: `mailbox subtree contains unexpected entry ${JSON.stringify(child)} under ${WORKFLOW_RUN_MAILBOX_PREFIX}/; only "${WORKFLOW_RUN_MAILBOX_INBOX_DIR}" is allowed`, + }; + } + } + const inboxPath = `${WORKFLOW_RUN_MAILBOX_PREFIX}/${WORKFLOW_RUN_MAILBOX_INBOX_DIR}`; + const inboxEntries = await listDir(inboxPath); + // A `mailbox/` top-level whose `INBOX` child carries no entries is a + // dangling blob committed directly at `mailbox/INBOX` (a blob lists as + // empty), not the required directory. Git never records an empty + // directory, so a present-but-empty listing is always the blob case. + if (inboxEntries.length === 0) { + return { + ok: false, + reason: `mailbox ${inboxPath} is a blob, not a directory carrying "${WORKFLOW_RUN_MAILBOX_INDEX_FILE}" and .eml message files`, + }; + } + for (const entry of inboxEntries) { + const entryPath = `${inboxPath}/${entry}`; + if (entry === WORKFLOW_RUN_MAILBOX_INDEX_FILE) { + if ((await listDir(entryPath)).length > 0) { + return { + ok: false, + reason: `mailbox ${entryPath} is a directory; the mailbox index must be a single file`, + }; + } + indexPresent = true; + continue; + } + if (!MAILBOX_EML_FILENAME_RE.test(entry)) { + return { + ok: false, + reason: `mailbox entry ${entryPath} does not match "${WORKFLOW_RUN_MAILBOX_INDEX_FILE}" or .eml (uid a decimal integer >= 1)`, + }; + } + if ((await listDir(entryPath)).length > 0) { + return { + ok: false, + reason: `mailbox message ${entryPath} is a directory; each message must be a single .eml file`, + }; + } + emlPaths.add(entryPath); + } + return { ok: true, emlPaths, indexPresent }; +} + +/** + * Validate the `mailbox/INBOX/` subtree (design conversational-mailbox). + * Enforces the subtree shape via `enumerateMailboxInbox`, then holds a + * RETAINED `.eml` message blob byte-identical against the prior tree. + * A prior `.eml` absent from the prospective tree is a legal expunge: + * the warm agent physically removes a message from the live INBOX. The raw + * bytes are not lost -- they stay reachable through the parent commit, and + * a `workflow-run` repo's objects are never GC'd (its kind is excluded + * from the GC allow-list; see `DEFAULT_GC_KINDS` in `agent-repo`), so the + * expunged message survives in history for the life of the run repo. The + * audit trail therefore rests on "these objects are never pruned", not on + * the live tree being monotonic. + * + * The `index.json` entry is mutable, but must PERSIST once it existed: if + * the prior tree carried an index and the prospective tree drops it, the + * whole mailbox has vanished, which would reset `uidValidity` / `uidNext` + * on the next open and force uid reuse from 1. That is rejected. A + * well-behaved backing always rewrites `index.json` on flush, so the guard + * fails no legitimate push. Mirrors the blob-immutability discipline the + * `runs//blobs/` subtree uses, minus the deletion direction. + */ +async function validateMailboxSubtree( + listDir: (path: string) => Promise, + readBlob: (path: string) => Promise, + priorListDir: (path: string) => Promise, + priorReadBlob: (path: string) => Promise, +): Promise { + const prospective = await enumerateMailboxInbox(listDir); + if (!prospective.ok) return prospective; + const prior = await enumerateMailboxInbox(priorListDir); + if (!prior.ok) { + return { + ok: false, + reason: `prior tree's mailbox subtree is structurally invalid: ${prior.reason}`, + }; + } + for (const emlPath of prospective.emlPaths) { + if (!prior.emlPaths.has(emlPath)) continue; // newly added + const immutable = await checkMailboxEmlPriorByteEquality( + emlPath, + readBlob, + priorReadBlob, + ); + if (!immutable.ok) return immutable; + } + if (prior.indexPresent && !prospective.indexPresent) { + return { + ok: false, + reason: `mailbox ${WORKFLOW_RUN_MAILBOX_PREFIX}/${WORKFLOW_RUN_MAILBOX_INBOX_DIR}/${WORKFLOW_RUN_MAILBOX_INDEX_FILE} present in the prior tree is missing from the prospective tree; the mailbox index must persist so uidValidity/uidNext continuity holds across reopens`, + }; + } + return { ok: true }; +} + export const workflowRunKindHandler: KindHandler = { kind: "workflow-run", directoryPrefix: "workflow-runs", @@ -2063,7 +2533,7 @@ export const workflowRunKindHandler: KindHandler = { if (!ALLOWED_TOP_LEVEL.has(entry)) { return { ok: false, - reason: `unexpected top-level entry ${JSON.stringify(entry)}; allowed: "${WORKFLOW_RUN_RUNS_PREFIX}", "${WORKFLOW_RUN_ADDRESSES_PREFIX}", "${WORKFLOW_RUN_AGENT_STATE_PREFIX}", "${WORKFLOW_RUN_GITIGNORE_PATH}"`, + reason: `unexpected top-level entry ${JSON.stringify(entry)}; allowed: "${WORKFLOW_RUN_RUNS_PREFIX}", "${WORKFLOW_RUN_ADDRESSES_PREFIX}", "${WORKFLOW_RUN_AGENT_STATE_PREFIX}", "${WORKFLOW_RUN_MAILBOX_PREFIX}", "${WORKFLOW_RUN_GITIGNORE_PATH}"`, }; } } @@ -2111,6 +2581,28 @@ export const workflowRunKindHandler: KindHandler = { } } + const mailboxPresent = + topLevelTreePaths.includes(WORKFLOW_RUN_MAILBOX_PREFIX) || + priorTopLevels.includes(WORKFLOW_RUN_MAILBOX_PREFIX); + if (mailboxPresent) { + // Enter mailbox validation when the prospective OR prior tree + // carries a `mailbox/` subtree. A prospective tree that drops the + // subtree while the prior tree held one must still go through the + // walk so the index-continuity guard fires on the vanished + // `index.json` (a message-only expunge is legal; dropping the whole + // mailbox is not). + const mailboxCheck = await validateMailboxSubtree( + listDir, + readBlob, + priorListDir, + priorReadBlob, + ); + if (!mailboxCheck.ok) { + logger.debug`workflow-run validatePush rejected ${repoId.kind}/${repoId.id} on ${ref}: ${mailboxCheck.reason}`; + return mailboxCheck; + } + } + const runsPresent = topLevelTreePaths.includes(WORKFLOW_RUN_RUNS_PREFIX) || priorTopLevels.includes(WORKFLOW_RUN_RUNS_PREFIX); @@ -2200,7 +2692,7 @@ export const workflowRunKindHandler: KindHandler = { // checkPriorByteEquality rejects it first. Mirrors the newly-terminal // gate below, which likewise acts only on a blob absent from the // prior tree. - if ((await priorReadBlob(entry.blobPath)) === null) { + if (await blobIsNewlyAdded(entry.blobPath, priorReadBlob)) { const principalCheck = checkCancelOriginPrincipal( entry.blobPath, origin, @@ -2218,7 +2710,8 @@ export const workflowRunKindHandler: KindHandler = { reason: `run ${runId} has event at seq ${String(entry.filenameSeq)} after terminal ${terminalType} at seq ${String(terminalSeq)}`, }; } - if (TERMINAL_EVENT_TYPES.has(parsed.parsed.body.type)) { + const classified = classifyTerminalEvent(parsed.parsed.body.type); + if (classified.terminal) { terminalSeq = entry.filenameSeq; terminalType = parsed.parsed.body.type; // Surface the run as newly terminal only when this commit is @@ -2229,17 +2722,11 @@ export const workflowRunKindHandler: KindHandler = { // terminal blob already present in the prior tree and emits no // signal, so a downstream consumer keyed on the signal does // not double-fire. - if ((await priorReadBlob(entry.blobPath)) === null) { - const status = TERMINAL_EVENT_STATUS.get(parsed.parsed.body.type); - if (status === undefined) { - throw new Error( - `terminal event type ${parsed.parsed.body.type} has no workflow_run.status mapping`, - ); - } + if (await blobIsNewlyAdded(entry.blobPath, priorReadBlob)) { const terminalBytes = await readBlob(entry.blobPath); newlyTerminalRuns.push({ runId, - status, + status: classified.status, terminalEventJson: new TextDecoder().decode(terminalBytes), }); } @@ -2287,14 +2774,20 @@ export const workflowRunKindHandler: KindHandler = { typeof body === "object" && body !== null && "type" in body ? body.type : undefined; - const status = - typeof type === "string" ? TERMINAL_EVENT_STATUS.get(type) : undefined; - if (status === undefined) { + const classified = + typeof type === "string" + ? classifyTerminalEvent(type) + : ({ terminal: false } as const); + if (!classified.terminal) { throw new Error( `combined event log ${combinedPath} for run ${runId} sealed without a recognized terminal event type`, ); } - newlyTerminalRuns.push({ runId, status, terminalEventJson: lastLine }); + newlyTerminalRuns.push({ + runId, + status: classified.status, + terminalEventJson: lastLine, + }); } const blobsEnumerated = await enumerateRunBlobs(listDir, scopeRunIds); @@ -2369,6 +2862,20 @@ export const workflowRunKindHandler: KindHandler = { }; } + const partsCheck = await validateRunPartsSubtree({ + listDir, + priorListDir, + readBlob, + priorReadBlob, + listDirOids, + priorListDirOids, + scopeRunIds, + }); + if (!partsCheck.ok) { + logger.debug`workflow-run validatePush rejected ${repoId.kind}/${repoId.id} on ${ref}: ${partsCheck.reason}`; + return partsCheck; + } + return { ok: true, newlyTerminalRuns }; }, onRefUpdated() { @@ -2772,10 +3279,7 @@ export type EnqueueInboxResult = { * receipt on the enqueue may safely acknowledge on any of them. */ export type EnqueueAlreadyPresentReason = - | "duplicate" - | "already_inbox" - | "processing" - | "consumed"; + "duplicate" | "already_inbox" | "processing" | "consumed"; /** * Outcome of an `enqueueInbox` call. Modeled as a value (not an exception) @@ -3260,30 +3764,40 @@ export type ReplayProcessingToInboxResult = { replayedKeys: string[]; }; +export type ScanRunsForBootResult = { + ownedMessageIds: Set; + pendingSealRunIds: string[]; +}; + /** - * Read the run event logs under `runs/` and return the set of - * `consumedMessageId`s belonging to NON-terminal runs -- the messages a - * live run still owns. The caller (the supervisor's spawn-time replay) - * feeds this into `replayProcessingToInbox`'s `ownedMessageIds` so a - * parked run's message is not re-admitted to inbox and dispatched a - * second time while the run is recovered by re-driving its durable log. - * Without this, the re-drive AND the re-triggered fresh run both re-park - * the same awaitSignal gate on the same runId, and the two concurrent - * runtime bodies race to a corrupt terminal. + * Walk `runs/` once and return the two boot-recovery inputs the supervisor's + * spawn needs, from a single traversal of the working tree via `getRepoDir`: + * + * - `ownedMessageIds`: the `consumedMessageId`s of NON-terminal runs -- the + * messages a live run still owns. Spawn feeds this into + * `replayProcessingToInbox`'s `ownedMessageIds` so a parked run's message is + * not re-admitted to inbox and dispatched a second time while the run is + * recovered by re-driving its durable log. Without this, the re-drive AND the + * re-triggered fresh run both re-park the same awaitSignal gate on the same + * runId, and the two concurrent runtime bodies race to a corrupt terminal. + * - `pendingSealRunIds`: runs that are terminal but still in per-event form -- + * an interrupted fold left them unsealed. Spawn hands these to the recovery + * sweep, which re-runs the idempotent fold. A terminal event is a *proposal*: + * the authoritative decision is `compactRunEvents`, which independently + * re-checks the run's max-seq event and no-ops a run that is not actually + * terminal, so this scan may be loose. * - * Reads the substrate's working tree via `getRepoDir`, mirroring the - * child's `discoverInFlightRuns`. The working tree tracks the run-event - * ref (`refs/heads/main`); the claim-check ref (`refs/heads/events`) - * cannot see it, which is why this lives at the caller rather than inside - * `replayProcessingToInbox`'s single-ref delta. A run whose log is sealed - * (combined `events.json`, only permitted for a terminated run) or - * carries a terminal event is excluded; an absent `runs/` directory - * yields an empty set. + * The working tree tracks the run-event ref (`refs/heads/main`); the + * claim-check ref (`refs/heads/events`) cannot see it, which is why this lives + * at the caller rather than inside `replayProcessingToInbox`'s single-ref + * delta. A run whose log is sealed (combined `events.jsonl`, only permitted for + * a terminated run) contributes to neither set; an absent `runs/` directory + * yields empty results. */ -export async function readOwnedMessageIds( +export async function scanRunsForBoot( store: RepoStore, repoId: RepoId, -): Promise> { +): Promise { const fs = await import("node:fs/promises"); const path = await import("node:path"); const repoDir = store.getRepoDir(repoId); @@ -3293,21 +3807,33 @@ export async function readOwnedMessageIds( runIds = await fs.readdir(runsDir); } catch (cause) { if (cause instanceof Error && "code" in cause && cause.code === "ENOENT") { - return new Set(); + return { ownedMessageIds: new Set(), pendingSealRunIds: [] }; } throw cause; } const owned = new Set(); + const pendingSealRunIds: string[] = []; for (const runId of runIds) { const runDir = path.join(runsDir, runId); // A sealed run (combined events file) is terminal by the handler's // own invariant -- only a terminated run is sealed -- so it owns - // nothing. Its presence also means the per-event directory is absent. + // nothing and is already folded. Its presence also means the per-event + // directory is absent. let sealed = false; try { await fs.access(path.join(runDir, WORKFLOW_RUN_EVENTS_FILE)); sealed = true; - } catch { + } catch (cause) { + // ENOENT is the normal "not sealed" case. Any other stat error leaves + // the run to fall through to the events-dir read below, which resolves + // it, so this catch is benign; warn so the anomaly is still visible. + if ( + !(cause instanceof Error) || + !("code" in cause) || + cause.code !== "ENOENT" + ) { + logger.warn`scanRunsForBoot: stat of the sealed-log file for run ${runId} failed: ${cause instanceof Error ? cause.message : String(cause)}`; + } sealed = false; } if (sealed) continue; @@ -3315,7 +3841,21 @@ export async function readOwnedMessageIds( let files: string[]; try { files = await fs.readdir(eventsDir); - } catch { + } catch (cause) { + // ENOENT means the run has neither a sealed log nor a per-event + // directory (grants may be staged before the first event); skip it. + // A non-ENOENT error drops the run from BOTH result sets, and a live + // run dropped from ownedMessageIds gets its message re-admitted and + // dispatched a second time on the same runId -- the double-driver + // corruption this scan exists to prevent. Surface it, but still skip: + // aborting the whole boot scan over one run is worse. + if ( + !(cause instanceof Error) || + !("code" in cause) || + cause.code !== "ENOENT" + ) { + logger.error`scanRunsForBoot: reading events for run ${runId} failed; skipping it may re-admit its message and start a second run on the same runId: ${cause instanceof Error ? cause.message : String(cause)}`; + } continue; } let terminal = false; @@ -3327,7 +3867,19 @@ export async function readOwnedMessageIds( parsed = JSON.parse( await fs.readFile(path.join(eventsDir, file), "utf8"), ); - } catch { + } catch (cause) { + // A corrupt or unreadable event file drops this run's + // classification: a missed RunStarted re-admits its message (a + // second run on the same runId), a missed terminal event skips a + // needed seal. Surface it, but skip the file rather than abort the + // scan. ENOENT here is a benign race (the file vanished mid-scan). + if ( + !(cause instanceof Error) || + !("code" in cause) || + cause.code !== "ENOENT" + ) { + logger.error`scanRunsForBoot: reading event ${file} for run ${runId} failed; skipping it may re-admit its message and start a second run on the same runId: ${cause instanceof Error ? cause.message : String(cause)}`; + } continue; } if ( @@ -3348,10 +3900,13 @@ export async function readOwnedMessageIds( if (typeof mid === "string") consumedMessageId = mid; } } - if (terminal) continue; + if (terminal) { + pendingSealRunIds.push(runId); + continue; + } if (consumedMessageId !== undefined) owned.add(consumedMessageId); } - return owned; + return { ownedMessageIds: owned, pendingSealRunIds }; } export type WorkflowRunLifecycle = "absent" | "live" | "terminal"; @@ -3446,14 +4001,16 @@ export async function readCommittedWorkflowRunTerminalStatus( typeof body === "object" && body !== null && "type" in body ? body.type : undefined; - const status = - typeof type === "string" ? TERMINAL_EVENT_STATUS.get(type) : undefined; - if (status === undefined) { + const classified = + typeof type === "string" + ? classifyTerminalEvent(type) + : ({ terminal: false } as const); + if (!classified.terminal) { throw new Error( `sealed run ${runId} has no recognized terminal event type`, ); } - return status; + return classified.status; } const eventsPath = `${runPath}/${WORKFLOW_RUN_EVENTS_DIR}`; @@ -3487,9 +4044,9 @@ export async function readCommittedWorkflowRunTerminalStatus( typeof parsed === "object" && parsed !== null && "type" in parsed ? parsed.type : undefined; - return typeof type === "string" - ? (TERMINAL_EVENT_STATUS.get(type) ?? null) - : null; + if (typeof type !== "string") return null; + const classified = classifyTerminalEvent(type); + return classified.terminal ? classified.status : null; } /** diff --git a/vendor/intx/hub-sessions/src/ws/sidecar-handler.ts b/vendor/intx/hub-sessions/src/ws/sidecar-handler.ts index c213dbfa9..6e89ae956 100644 --- a/vendor/intx/hub-sessions/src/ws/sidecar-handler.ts +++ b/vendor/intx/hub-sessions/src/ws/sidecar-handler.ts @@ -52,22 +52,30 @@ import { const logger = getLogger(["hub", "ws", "sidecar"]); /** - * A deploy frame failure that PROVABLY never reached the wire: thrown only by - * a guard clause that runs before `conn.send()`, or by `conn.send()` itself - * throwing synchronously. A caller that inserted a row anticipating the frame - * (e.g. `deployCodeSourcedWorkflow`'s pre-inserted anchor `workflow_run`) may - * safely undo that insert on this error class specifically. - * - * Any OTHER deploy rejection (ack timeout, socket drop, reconnect takeover, - * ack-processing failure) is raised after `conn.send()` already ran -- the - * frame may have reached and been acted on by the sidecar -- and must NOT be - * treated as not-sent. + * A deploy-frame send failure, tagged with whether the `agent.deploy` frame + * reached the wire. `frameSent: false` means the send was refused before + * `conn.send` (a guard failed, or the send threw synchronously) -- the deploy + * provably never started, so a caller may safely roll back anything it staged. + * `frameSent: true` means the frame was sent and the failure came afterward (ack + * timeout, sidecar disconnect), so the sidecar may hold a live agent. */ -export class DeployFrameNotSentError extends Error { - constructor(message: string) { - super(message); - this.name = "DeployFrameNotSentError"; - } +export interface DeployFrameFailure extends Error { + readonly frameSent: boolean; +} + +function deployFrameFailure( + message: string, + frameSent: boolean, +): DeployFrameFailure { + return Object.assign(new Error(message), { frameSent }); +} + +export function isDeployFrameFailure(err: unknown): err is DeployFrameFailure { + return ( + err instanceof Error && + "frameSent" in err && + typeof err.frameSent === "boolean" + ); } export type SidecarConnection = { @@ -1934,6 +1942,15 @@ export function createSidecarRouter( isRunAddress(recipient) ) { const runId = deriveWorkflowRunId(recipient); + // This does NOT let mail mutate a run's authorization. First delivery + // reserves and commits the run's grants (the mail IS the trigger); + // every later delivery only RE-READS the current committed grants + // (`loadCommittedRunGrants`) and re-asserts them ahead of the dispatch. + // The committed rows already carry any standing-approval change (an + // approve/reject-with-`always` resolution mutates them through its own + // path), so this re-send is idempotent -- it re-establishes the run's + // current floor on the sidecar, self-healing a `grants.json` a sidecar + // may have lost, and never overwrites it with anything staler. const result = await lookups.materializeMailTriggeredRunGrants({ agentAddress: recipient, runId, @@ -3107,26 +3124,30 @@ export function createSidecarRouter( workflow?: AgentDeployFrame["workflow"], ): Promise<{ publicKey: string }> { if (hubPublicKeyHex === undefined) { - throw new DeployFrameNotSentError( + throw deployFrameFailure( "Hub signing key is required for agent deployment", + false, ); } const ws = addressIndex.get(agentAddress) ?? findSidecarForNewAgent(agentAddress); if (ws === undefined) { - throw new DeployFrameNotSentError( + throw deployFrameFailure( `No sidecar available for agent "${agentAddress}"`, + false, ); } const conn = connections.get(ws); if (conn === undefined) { - throw new DeployFrameNotSentError( + throw deployFrameFailure( `No sidecar connected for agent "${agentAddress}"`, + false, ); } if (conn.identity.kind !== "shared") { - throw new DeployFrameNotSentError( + throw deployFrameFailure( `Allocated sidecar ${conn.sidecarId} requires allocation-bound deploy routing`, + false, ); } return sendAgentDeployOnConnection( @@ -3146,14 +3167,16 @@ export function createSidecarRouter( workflow?: AgentDeployFrame["workflow"], ): Promise<{ publicKey: string }> { if (hubPublicKeyHex === undefined) { - throw new DeployFrameNotSentError( + throw deployFrameFailure( "Hub signing key is required for agent deployment", + false, ); } if (pendingDeploys.has(agentAddress)) { - throw new DeployFrameNotSentError( + throw deployFrameFailure( `Deploy already in progress for agent "${agentAddress}"`, + false, ); } @@ -3172,8 +3195,9 @@ export function createSidecarRouter( addressIndex.delete(agentAddress); } reject( - new Error( + deployFrameFailure( `Deploy of "${agentAddress}" timed out after ${requestTimeoutMs}ms`, + true, ), ); }, requestTimeoutMs); @@ -3189,18 +3213,11 @@ export function createSidecarRouter( addressSet.delete(agentAddress); addressIndex.delete(agentAddress); } - reject(new Error(error)); + reject(deployFrameFailure(error, true)); }, timer, }); - // Once `pendingDeploys` carries this entry, every OTHER rejection path - // (timeout, disconnect, reconnect takeover, ack-processing failure) is - // reached only through `req.reject()` above -- which fires after this - // send, by construction. A synchronous throw HERE is the sole - // post-registration case that provably never left the process, so it - // is the one case converted to `DeployFrameNotSentError` rather than - // rejecting through `req.reject()`. try { conn.send({ type: "agent.deploy", @@ -3211,6 +3228,9 @@ export function createSidecarRouter( ...(workflow !== undefined ? { workflow } : {}), }); } catch (err) { + // A synchronous send failure means the frame never reached the wire. + // Tear down the pending entry and timer we just registered, and reject + // as not-sent so a caller may safely roll back what it staged. clearTimeout(timer); pendingDeploys.delete(agentAddress); if (addressIndex.get(agentAddress) === ws) { @@ -3218,8 +3238,9 @@ export function createSidecarRouter( addressIndex.delete(agentAddress); } reject( - new DeployFrameNotSentError( - `Failed to send agent.deploy frame to "${agentAddress}": ${err instanceof Error ? err.message : String(err)}`, + deployFrameFailure( + `Deploy of "${agentAddress}" failed to send: ${err instanceof Error ? err.message : String(err)}`, + false, ), ); } @@ -3234,18 +3255,18 @@ export function createSidecarRouter( ): Promise<{ publicKey: string }> { const { ws, conn } = await getAllocatedConnection(target, "routing"); if (conn.identity.kind !== "allocated") { - throw new DeployFrameNotSentError( + throw new Error( `Allocation ${target.allocationId} resolved to a shared sidecar`, ); } if (agentAddress !== conn.identity.workflowRunAddress) { - throw new DeployFrameNotSentError( + throw new Error( `Allocation ${target.allocationId} cannot deploy unrelated address ${agentAddress}`, ); } const existing = addressIndex.get(agentAddress); if (existing !== undefined && existing !== ws) { - throw new DeployFrameNotSentError( + throw new Error( `Deployment ${agentAddress} is already routed to another sidecar`, ); } diff --git a/vendor/intx/hub-sessions/tsconfig.json b/vendor/intx/hub-sessions/tsconfig.json index d984862c9..dbbb0384b 100644 --- a/vendor/intx/hub-sessions/tsconfig.json +++ b/vendor/intx/hub-sessions/tsconfig.json @@ -1,11 +1,11 @@ { "extends": "../tsconfig.base.json", + "include": [ + "src/**/*.ts" + ], "compilerOptions": { "types": [ "bun" ] - }, - "include": [ - "src/**/*.ts" - ] + } } diff --git a/vendor/intx/workflow-host/VENDORED-FROM b/vendor/intx/workflow-host/VENDORED-FROM index 3fcca01bd..631711869 100644 --- a/vendor/intx/workflow-host/VENDORED-FROM +++ b/vendor/intx/workflow-host/VENDORED-FROM @@ -1,4 +1,4 @@ Source: https://github.com/faremeter/interchange (packages/workflow-host) Commit: b5580a02fb918eebccc33ded7727ffee781ffbd1 (tag v0.3.0) License: LGPL-2.1-only (see vendor/intx/LICENSE) -Local modifications: exports map repointed from the upstream intx-src condition to direct TypeScript source resolution (types/default -> ./src/...); dist references removed. CL-6164: the supervisor's signal.deliver branch drops mail whose extracted conversation body is empty (the new hasConversationText gate in conversation-text.ts), recording an empty_conversation_content rejection, instead of delivering "" -- which throws in agent.send and kills the run with StepFailed/retriesExhausted. Attachments-only conversation.message mail (e.g. @corbits/chat's workbench.agent-joined event send) is exactly that shape. CL-6325: adds the action-primitive adapters (adapters/action-invoker.ts, adapters/effect-ledger.ts, adapters/run-blobs.ts and their tests) -- copied from gtm-workbench's packages/workflow-host workspace fork, not from upstream, which has no action-primitive adapters at the pinned commit (see VENDORED.md) -- and the completed run-child bind in child/run-child.ts: RunWorkflowChildBindings gains resolveActionHandler (awaited once per child with the re-verified definition and the live CredentialWiring) and loopFns (both defaulting to the fail-closed empty registries), and buildRuntimeEnv wires effects, invokeAction, loopFns, and runLoopIteration (createLoopIteration) into every run's WorkflowRuntimeEnv. buildRuntimeEnv is exported (child/index.ts, index.ts) so a host's runtime-env-level probe can exercise the bind without the full control-channel harness. Bridging edit against the re-pinned @intx/types (a8bc06ae): child/supervisor-backed-transport.ts's expunge stub returns Promise<{ expungedUids: number[] }> (upstream bcabb1f8); gone with this tree's own re-pin. +Local modifications: exports map repointed from the upstream intx-src condition to direct TypeScript source resolution (types/default -> ./src/...); dist references removed. CL-6164: the supervisor's signal.deliver branch drops mail whose extracted conversation body is empty (the new hasConversationText gate in conversation-text.ts), recording an empty_conversation_content rejection, instead of delivering "" -- which throws in agent.send and kills the run with StepFailed/retriesExhausted. Attachments-only conversation.message mail (e.g. @corbits/chat's workbench.agent-joined event send) is exactly that shape. CL-6325: adds the action-primitive adapters (adapters/action-invoker.ts, adapters/effect-ledger.ts, adapters/run-blobs.ts and their tests) -- copied from gtm-workbench's packages/workflow-host workspace fork, not from upstream, which has no action-primitive adapters at the pinned commit (see VENDORED.md) -- and the completed run-child bind in child/run-child.ts: RunWorkflowChildBindings gains resolveActionHandler (awaited once per child with the re-verified definition and the live CredentialWiring) and loopFns (both defaulting to the fail-closed empty registries), and buildRuntimeEnv wires effects, invokeAction, loopFns, and runLoopIteration (createLoopIteration) into every run's WorkflowRuntimeEnv. buildRuntimeEnv is exported (child/index.ts, index.ts) so a host's runtime-env-level probe can exercise the bind without the full control-channel harness. Bridging edits against the re-pinned @intx/types and @intx/hub-sessions (a8bc06ae): child/supervisor-backed-transport.ts's expunge stub returns Promise<{ expungedUids: number[] }> (upstream bcabb1f8), and supervisor.ts's boot replay reads ownedMessageIds from scanRunsForBoot (upstream f89bb51b) since readOwnedMessageIds no longer exists; both gone with this tree's own re-pin. diff --git a/vendor/intx/workflow-host/src/supervisor/supervisor.ts b/vendor/intx/workflow-host/src/supervisor/supervisor.ts index 87e16b246..dfa2578d8 100644 --- a/vendor/intx/workflow-host/src/supervisor/supervisor.ts +++ b/vendor/intx/workflow-host/src/supervisor/supervisor.ts @@ -55,7 +55,7 @@ import { enqueueInbox as defaultEnqueueInbox, dequeueToProcessing as defaultDequeueToProcessing, markConsumed as defaultMarkConsumed, - readOwnedMessageIds, + scanRunsForBoot, readWorkflowRunLifecycle, replayProcessingToInbox as defaultReplayProcessingToInbox, StaleInboxEnqueueError, @@ -1973,11 +1973,11 @@ export function createWorkflowSupervisor( // first `dequeueToProcessing` so a fresh inbound mail that lands // during the replay window cannot ship ahead of the orphan once // the replay completes. - const replayDone = readOwnedMessageIds( + const replayDone = scanRunsForBoot( bindings.repoStore, bindings.workflowRunRepoId, ) - .then((ownedMessageIds) => + .then(({ ownedMessageIds }) => inboxPrimitives.replayProcessingToInbox( bindings.repoStore, inboxWritePrincipal,