Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 9 additions & 45 deletions VENDORED.md

Large diffs are not rendered by default.

16 changes: 6 additions & 10 deletions apps/hub/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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({
Expand Down
32 changes: 16 additions & 16 deletions packages/insights/src/on-usage-wire.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -22,31 +25,30 @@ 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,
};
}

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({
store,
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",
Expand All @@ -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"]);
Expand Down
4 changes: 2 additions & 2 deletions scripts/checks/kill-dates.txt
Original file line number Diff line number Diff line change
Expand Up @@ -17,15 +17,15 @@ 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
vendor/intx/mime | sawyer | 2026-10-26 | d02e5f8f1429eac7c27d3a37eec31111f8a1053c91fbfae77ac58a0d63c823ed
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

8 changes: 2 additions & 6 deletions vendor/intx/hub-sessions/VENDORED-FROM
Original file line number Diff line number Diff line change
@@ -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/<runId>/events/<seq>.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__<n>) 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/<runId> 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).
13 changes: 6 additions & 7 deletions vendor/intx/hub-sessions/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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:",
Expand Down
19 changes: 18 additions & 1 deletion vendor/intx/hub-sessions/src/agent-repo.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<uid>.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 };
Expand Down Expand Up @@ -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" };
Expand Down
Loading
Loading