Migrate the keyed package-invocation idempotency ledger from D1 into the per-user RunLog DO - #1053
Conversation
📝 WalkthroughWalkthroughKeyed package invocation idempotency now uses a per-user RunLog Durable Object ledger. The ledger supports claims, replay responses, retention, export, deletion, and legacy D1 read fallback while module execution and subscription tests use the new flow. ChangesRunLog invocation ledger
Estimated code review effort: 5 (Critical) | ~120 minutes Sequence Diagram(s)sequenceDiagram
participant Caller as packages.invoke
participant Service as run-records service
participant RunLog as RunLog Durable Object
participant Module as saved package module
Caller->>Service: claim keyed invocation
Service->>RunLog: create or resolve ledger claim
Caller->>Module: execute when claim is owned
Caller->>Service: finish or release claim
Service->>RunLog: persist replay response and run record
Possibly related PRs
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches 💡 1📝 Generate docstrings 💡
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
|
🔎 Preview deployed: https://kody-pr-1053.kody-a99.workers.dev Worker: Mocks:
|
There was a problem hiding this comment.
Actionable comments posted: 2
🧹 Nitpick comments (2)
packages/worker/src/package-invocations/service.node.test.ts (1)
244-250: 🔒 Security & Privacy | 🔵 Trivial | ⚡ Quick winConsider making the fake namespace user-aware.
idFromNamediscards the name andget()returns one sharedrpcfor every id, so all users collapse into a single in-memory ledger. A regression that stopped namespacing the RunLog DO id byuserId(cross-user replay/claim) would still pass this suite. Capturing the requested names — or keyingrpcstate per id — would let these tests pin the scoping contract.♻️ Sketch
+ const requestedIds: Array<string> = [] return { namespace: { - idFromName: (name: string) => name as unknown as DurableObjectId, - get: () => rpc, + idFromName: (name: string) => { + requestedIds.push(name) + return name as unknown as DurableObjectId + }, + get: () => rpc, }, + requestedIds, ledgerRows, runRows,Then assert
db.runLog.requestedIds.every((id) => id.includes('user-123'))in the keyed-path tests.As per coding guidelines, "every Durable Object ID backing user-owned state must be namespaced by
userId".🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/package-invocations/service.node.test.ts` around lines 244 - 250, Update the fake namespace setup in the test helper so idFromName records each requested name and get returns RPC state keyed by that Durable Object ID instead of one shared rpc. Preserve existing behavior for each individual ID, then add assertions in the keyed-path tests that requested RunLog IDs include the expected userId, such as user-123.Source: Coding guidelines
packages/worker/src/run-records/service.ts (1)
680-711: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winConsider extracting the terminal run-row construction shared with
finishRunRecord.Lines 683–710 duplicate
finishRunRecord's finishedAt/durationMs/error-fields/result-metadata/buildRunRowblock verbatim. A small helper (e.g.buildTerminalRunRow({ handle, status, error, result })) keeps the two terminal paths from drifting as snapshot/truncation rules evolve.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/run-records/service.ts` around lines 680 - 711, Extract the duplicated terminal run-row assembly into a shared helper, such as buildTerminalRunRow, and use it from both this path and finishRunRecord. The helper should accept handle, status, error, and result, and centralize finishedAt, durationMs, error fields, result metadata snapshotting, and buildRunRow invocation so both terminal paths remain consistent.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@docs/contributing/architecture/run-records.md`:
- Around line 204-207: Update docs/contributing/architecture/run-records.md
lines 204-207 to describe the RunLog ledger as retained, bounded
idempotency/replay state rather than exactly-once execution. Update
docs/contributing/disaster-recovery.md lines 194-204 to explain that later
traffic recreates claims but cannot restore lost ledger rows, and qualify “not
data loss” so it excludes loss of correctness state.
In `@packages/worker/src/email/inbound.workers.test.ts`:
- Around line 3389-3397: Update the assertions around responses and invocations
so stored-message and approved-message expectations are matched using each
ledger row’s identifying fields rather than responses[0] or responses[1].
Preserve the existing sorting, or introduce a creation-order-based tiebreaker,
ensuring assertions remain correct when createdAt values are equal and UUID
ordering varies.
---
Nitpick comments:
In `@packages/worker/src/package-invocations/service.node.test.ts`:
- Around line 244-250: Update the fake namespace setup in the test helper so
idFromName records each requested name and get returns RPC state keyed by that
Durable Object ID instead of one shared rpc. Preserve existing behavior for each
individual ID, then add assertions in the keyed-path tests that requested RunLog
IDs include the expected userId, such as user-123.
In `@packages/worker/src/run-records/service.ts`:
- Around line 680-711: Extract the duplicated terminal run-row assembly into a
shared helper, such as buildTerminalRunRow, and use it from both this path and
finishRunRecord. The helper should accept handle, status, error, and result, and
centralize finishedAt, durationMs, error fields, result metadata snapshotting,
and buildRunRow invocation so both terminal paths remain consistent.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro Plus
Run ID: f21d005e-1b01-4c68-b446-7c0c9e5cb8f0
📒 Files selected for processing (20)
docs/contributing/architecture/data-storage.mddocs/contributing/architecture/invocation-overhead-guardrails.mddocs/contributing/architecture/run-records.mddocs/contributing/disaster-recovery.mdpackages/worker/src/account/export.node.test.tspackages/worker/src/account/export.tspackages/worker/src/account/user-owned-surfaces.tspackages/worker/src/app/retention.tspackages/worker/src/email/inbound.workers.test.tspackages/worker/src/email/system-email-subscriptions.workers.test.tspackages/worker/src/package-invocations/idempotency.tspackages/worker/src/package-invocations/idempotent-module-invocation.tspackages/worker/src/package-invocations/module-execution.tspackages/worker/src/package-invocations/repo.tspackages/worker/src/package-invocations/service.node.test.tspackages/worker/src/platform-feedback/platform-feedback-subscriptions.workers.test.tspackages/worker/src/run-records/invocation-ledger.workers.test.tspackages/worker/src/run-records/run-log-do.tspackages/worker/src/run-records/service.tspackages/worker/src/run-records/types.ts
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In @.github/workflows/backfill-mcp-agent-sessions.yml:
- Around line 7-9: Update the wrapped description values in the workflow,
including the descriptions near the dry-run text and the additional occurrence,
to use folded YAML scalar syntax with `description: >`; preserve their existing
wording while ensuring the workflow parses correctly.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro Plus
Run ID: 54bd99df-36b3-4637-9b36-55bae1a0d4d1
📒 Files selected for processing (1)
.github/workflows/backfill-mcp-agent-sessions.yml
| description: | ||
| Dry-run audits only; execute writes index rows and the completion | ||
| marker. |
There was a problem hiding this comment.
🎯 Functional Correctness | 🔴 Critical | ⚡ Quick win
Use folded YAML scalars for the wrapped descriptions.
These description values are invalid YAML because multiline plain scalars require > or |; the workflow will fail to parse. Use description: > before each wrapped value.
Proposed fix
- description:
+ description: >
Dry-run audits only; execute writes index rows and the completion
marker.
...
- description:
+ description: >
Cloudflare Workers script name that owns the MCP Durable Object
namespaceAlso applies to: 20-22
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In @.github/workflows/backfill-mcp-agent-sessions.yml around lines 7 - 9, Update
the wrapped description values in the workflow, including the descriptions near
the dry-run text and the additional occurrence, to use folded YAML scalar syntax
with `description: >`; preserve their existing wording while ensuring the
workflow parses correctly.
There was a problem hiding this comment.
Not applicable: multi-line plain scalars are valid YAML — a value may continue on following lines as long as they are indented more deeply than the key, which is exactly what oxfmt produced here. Verified locally with yaml.parse(...) (parses cleanly) and the workflow loads fine in Actions. The other two findings (assertion ordering, ledger-loss wording) were valid and are fixed in 9e69397.
There was a problem hiding this comment.
Skipped: comment is from another GitHub bot.
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes using default effort and found 3 potential issues.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 9e69397. Configure here.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@packages/worker/src/email/inbound.workers.test.ts`:
- Around line 3391-3399: Update the responseBodies mapping to validate each
row’s responseJson before JSON.parse: explicitly fail or handle null replay
data, and only parse non-null serialized responses before accessing body.
Preserve the existing response body extraction for valid ledger records.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro Plus
Run ID: 9639b5ae-67e1-4570-9c0c-74c9a3e96521
📒 Files selected for processing (3)
docs/contributing/architecture/run-records.mddocs/contributing/disaster-recovery.mdpackages/worker/src/email/inbound.workers.test.ts
🚧 Files skipped from review as they are similar to previous changes (2)
- docs/contributing/disaster-recovery.md
- docs/contributing/architecture/run-records.md
…er in the RunLog DO New package_invocation_ledger table (schema v5) inside the per-user RunLog Durable Object with combined claim+begin and finish+terminal-response RPCs (one on-path DO call each), atomic stale reclaim, fenced terminal writes, claim release, DO-local 90-day terminal-row retention wired into the existing retention passes/alarm, and export paging that continues from runs into ledger rows through one cursor. Service wrappers build the run rows and share the activation/run-error side effects with finishRunRecord.
The keyed packages.invoke path now claims (ledger row + eager run record) and finishes (terminal replay response + run record) in one awaited DO call each, preserving request-hash mismatch 409s, 100ms/1s in-progress polling, 15-minute stale reclaim, and the bounded replay cache. Dual-read window: the DO is read first with a read-only D1 fallback for pre-migration keys; the D1 write helpers are removed and no awaited D1 write remains on the keyed hot path (the node suite's fake D1 throws on any write). The withAccountWriteLease boundary wrapper is dropped to match the key-less lean path: it guarded D1 writes, none remain here, and the DO rows are purged by account deletion's clearRunRecords. The D1 retention sweep stays (marked legacy) to drain pre-migration rows until a follow-up drops the table; account export includes ledger rows via the run_records section.
…cate-execution risk
…orkers suites The system-email, platform-feedback, and inbound-email dispatch tests asserted idempotency rows in the D1 package_invocations table; they now read each owner's RunLog DO ledger through exportRunRecords.
…edger-loss wording
…t ledger rows in the export manifest Bugbot caught a dual-read gap: a stale DO claim reclaimed for a pre-migration key skipped the legacy lookup, so a terminal D1 row finished by a zombie pre-migration isolate would re-execute instead of replaying. The fallback now runs after every claim (fresh or reclaimed); DO replays still never touch D1. Also counts ledger rows in the run_records manifest section for full exports and hardens the inbound test against null replay caches.
c97cbb6 to
8d03cf8
Compare
There was a problem hiding this comment.
Actionable comments posted: 3
🧹 Nitpick comments (6)
packages/worker/src/run-records/invocation-ledger.workers.test.ts (1)
25-46: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winType the override bags with the real domain types instead of
Partial<Record<string, …>>.
Partial<Record<string, string>>accepts any key, so a typo likeidempotencyKeyssilently falls through to the default and the test passes for the wrong reason. Using the exported types also removes theString(...)coercions.♻️ Proposed typing
-function ledgerKey(overrides?: Partial<Record<string, string>>) { +function ledgerKey( + overrides?: Partial<PackageInvocationLedgerKey>, +): PackageInvocationLedgerKey { return { tokenId: overrides?.tokenId ?? 'token-1', packageId: overrides?.packageId ?? 'pkg-1', exportName: overrides?.exportName ?? './send-message', idempotencyKey: overrides?.idempotencyKey ?? 'evt-1', } } -function claimInput(overrides?: Partial<Record<string, string | null>>) { +function claimInput( + overrides?: Partial<PackageInvocationClaimInput>, +): PackageInvocationClaimInput { return { - id: String(overrides?.id ?? crypto.randomUUID()), - tokenId: String(overrides?.tokenId ?? 'token-1'), - packageId: String(overrides?.packageId ?? 'pkg-1'), - packageKodyId: String(overrides?.packageKodyId ?? 'pkg-one'), - exportName: String(overrides?.exportName ?? './send-message'), - idempotencyKey: String(overrides?.idempotencyKey ?? 'evt-1'), - requestHash: String(overrides?.requestHash ?? 'hash-1'), + id: overrides?.id ?? crypto.randomUUID(), + tokenId: overrides?.tokenId ?? 'token-1', + packageId: overrides?.packageId ?? 'pkg-1', + packageKodyId: overrides?.packageKodyId ?? 'pkg-one', + exportName: overrides?.exportName ?? './send-message', + idempotencyKey: overrides?.idempotencyKey ?? 'evt-1', + requestHash: overrides?.requestHash ?? 'hash-1', source: overrides?.source === undefined ? 'webhook' : overrides.source, topic: overrides?.topic === undefined ? null : overrides.topic, } }Requires adding
type PackageInvocationClaimInput, type PackageInvocationLedgerKeyto the./service.tsimport.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/run-records/invocation-ledger.workers.test.ts` around lines 25 - 46, Update the override parameters of ledgerKey and claimInput to use the exported PackageInvocationLedgerKey and PackageInvocationClaimInput types imported from ./service.ts, so unknown property names are rejected. Remove the String(...) coercions and construct these helpers using the correctly typed domain fields while preserving their existing defaults.packages/worker/src/run-records/run-log-do.ts (1)
1139-1203: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueReclaim path leaves the superseded attempt's
runningrun row behind.On reclaim the previous attempt's eager run row keeps
status = 'running'and is only healed by the surface stale TTL. Given the reclaim window is 15 minutes and export surface TTL is minutes-scale, this self-heals — just noting it's a deliberate gap rather than an oversight; a comment here would save the next reader a trip throughreconcileStaleRunning.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/run-records/run-log-do.ts` around lines 1139 - 1203, Add a concise comment in claimPackageInvocation’s reclaimable existing-row branch, near the reclaimed claim and insertRunningRun logic, documenting that the superseded attempt’s running run row is intentionally left for reconcileStaleRunning/stale-TTL healing rather than immediately marked complete. Do not alter the reclaim behavior.packages/worker/src/run-records/service.ts (1)
765-807: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winThe extracted helper's doc claims
finishRunRecordshares it, butfinishRunRecordstill inlines its own copy.
finishRunRecord(Lines 411–447) keeps the duplicated activation-milestone +run.errordispatch logic, so the two copies can now drift while the comment asserts they are shared. RoutingfinishRunRecordthroughdispatchTerminalRunRecordSideEffectsmakes the doc true and removes the duplication — the behaviour is equivalent (persistedRunnon-null is the same precondition as the helper'spersistedRunargument).♻️ Suggested consolidation in `finishRunRecord`
- if (input.status === 'success') { - const packageId = normalizeOptionalString(handle.context.packageId) - if (packageId) { - try { - await recordSuccessfulPackageRun(input.env, { - userId: handle.userId, - packageId, - surface: handle.context.surface, - }) - } catch (error) { - console.warn('run-record-activation-failed', error) - } - } - } - - if ( - persistedRun && - input.status === 'error' && - handle.context.surface !== 'subscription' - ) { - try { - const { dispatchRunErrorSubscriptionEvents } = - await import('./package-subscriptions.ts') - await dispatchRunErrorSubscriptionEvents({ - env: input.env, - userId: handle.userId, - run: persistedRun, - waitUntil: input.waitUntil, - }) - } catch (error) { - console.warn('run-error-subscription-dispatch-failed', error) - } - } + if (persistedRun) { + await dispatchTerminalRunRecordSideEffects({ + env: input.env, + handle, + persistedRun, + status: input.status, + waitUntil: input.waitUntil, + }) + }Note the success-path activation currently runs even when the run row failed to persist; if that is intentional, keep it outside the guard rather than dropping it.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/run-records/service.ts` around lines 765 - 807, Update finishRunRecord to call dispatchTerminalRunRecordSideEffects instead of maintaining its duplicated activation and run.error dispatch logic. Preserve the existing behavior by keeping success-path activation outside any persistedRun guard if it currently runs when persistence fails, and pass the same handle, status, environment, persisted row, and waitUntil values used by finishRunRecord.packages/worker/src/package-invocations/repo.ts (1)
23-27: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueStale doc reference to the legacy D1 write path.
This same change removes the D1 write helpers, so "Shared by the legacy D1 write path and the RunLog DO ledger" no longer has a second consumer.
📝 Suggested wording
- * Serialize a terminal response for the replay cache, dropping it when - * oversized. Shared by the legacy D1 write path and the RunLog DO ledger so - * both stores enforce the same restore-safe bound. + * Serialize a terminal response for the RunLog DO ledger replay cache, + * dropping it when oversized so the stored row stays restore-safe.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/package-invocations/repo.ts` around lines 23 - 27, Update the documentation comment above the terminal-response serialization helper to remove the stale reference to the removed legacy D1 write path, and describe only the RunLog DO ledger as its consumer while preserving the restore-safe size-bound description.packages/worker/src/package-invocations/service.node.test.ts (1)
244-250: 🔒 Security & Privacy | 🔵 Trivial | ⚡ Quick winFake
RUN_LOGnamespace ignores the DO id, so per-user scoping is untestable here.
get: () => rpcreturns the same ledger regardless of the name passed toidFromName, so a regression that stops namespacing the RunLog DO byuserId(or leaks one user's ledger into another's) would still pass these tests. Capturing the requested names and asserting they derive fromuser-123is cheap and guards the multi-user invariant.🧪 Suggested tweak
+ const requestedIds: Array<string> = [] return { namespace: { - idFromName: (name: string) => name as unknown as DurableObjectId, - get: () => rpc, + idFromName: (name: string) => { + requestedIds.push(name) + return name as unknown as DurableObjectId + }, + get: (id: DurableObjectId) => { + expect(String(id)).toContain('user-123') + return rpc + }, }, + requestedIds,As per coding guidelines, "every Durable Object ID backing user-owned state must be namespaced by
userId".🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/package-invocations/service.node.test.ts` around lines 244 - 250, Update the fake RUN_LOG namespace in the test fixture so idFromName records each requested name and get returns the corresponding user-scoped RPC/ledger instead of always sharing one instance. Add assertions in the relevant tests that the requested Durable Object names are derived from user-123, preserving separate ledgers for different users and guarding the user-owned state namespacing invariant.Source: Coding guidelines
packages/worker/src/package-invocations/idempotent-module-invocation.ts (1)
322-380: 🩺 Stability & Availability | 🔵 Trivial | 💤 Low valueOuter reclaim loop has no global deadline.
Each
existingiteration is bounded bypackageInvocationPollBudgetMs, but the outer loop can repeat indefinitely when competitors keep winning the stale reclaim, so total wall time is unbounded. Consider a single deadline computed before the loop and bailing out withinvocation_in_progressonce exceeded.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@packages/worker/src/package-invocations/idempotent-module-invocation.ts` around lines 322 - 380, Add one overall deadline before the outer existing-claim loop in the idempotency flow, and check it on each iteration before polling or reclaiming. When the deadline is exceeded, return the existing invocation-in-progress response path (using the established response builder and idempotency key) instead of retrying indefinitely; retain the per-iteration poll budget for individual lookups.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@packages/worker/src/email/system-email-subscriptions.workers.test.ts`:
- Line 390: Replace String(...) coercion with an explicit null/undefined guard
before JSON.parse in both migrated subscription tests:
packages/worker/src/email/system-email-subscriptions.workers.test.ts:390-390 and
packages/worker/src/platform-feedback/platform-feedback-subscriptions.workers.test.ts:367-367.
Fail or throw when responseJson is absent, then parse the narrowed responseJson
string within each test and per-package loop.
In `@packages/worker/src/package-invocations/idempotent-module-invocation.ts`:
- Around line 435-503: Update the `completed` and `failed` branches after
`finishPackageInvocationRecord` so a non-updated ledger with another
`in_progress` record returns the already-available `outcome.response`; only pass
terminal records (mismatch, completed, or failed) to `resolveLedgerRecord`.
Preserve the existing fallback error response when no record is available, and
keep the catch behavior unchanged.
In `@packages/worker/src/run-records/run-log-do.ts`:
- Around line 1231-1241: Replace the rowsWritten-based ledgerUpdated check in
finishPackageInvocation with an explicit verification that the ledger row was
updated while matching the expected invocationId, in_progress status, and
claimUpdatedAt fence. Apply the same owned-row/affected-row validation in
releasePackageInvocation before clearing the run row or returning released:
true; do not use rowsWritten as the ownership decision.
---
Nitpick comments:
In `@packages/worker/src/package-invocations/idempotent-module-invocation.ts`:
- Around line 322-380: Add one overall deadline before the outer existing-claim
loop in the idempotency flow, and check it on each iteration before polling or
reclaiming. When the deadline is exceeded, return the existing
invocation-in-progress response path (using the established response builder and
idempotency key) instead of retrying indefinitely; retain the per-iteration poll
budget for individual lookups.
In `@packages/worker/src/package-invocations/repo.ts`:
- Around line 23-27: Update the documentation comment above the
terminal-response serialization helper to remove the stale reference to the
removed legacy D1 write path, and describe only the RunLog DO ledger as its
consumer while preserving the restore-safe size-bound description.
In `@packages/worker/src/package-invocations/service.node.test.ts`:
- Around line 244-250: Update the fake RUN_LOG namespace in the test fixture so
idFromName records each requested name and get returns the corresponding
user-scoped RPC/ledger instead of always sharing one instance. Add assertions in
the relevant tests that the requested Durable Object names are derived from
user-123, preserving separate ledgers for different users and guarding the
user-owned state namespacing invariant.
In `@packages/worker/src/run-records/invocation-ledger.workers.test.ts`:
- Around line 25-46: Update the override parameters of ledgerKey and claimInput
to use the exported PackageInvocationLedgerKey and PackageInvocationClaimInput
types imported from ./service.ts, so unknown property names are rejected. Remove
the String(...) coercions and construct these helpers using the correctly typed
domain fields while preserving their existing defaults.
In `@packages/worker/src/run-records/run-log-do.ts`:
- Around line 1139-1203: Add a concise comment in claimPackageInvocation’s
reclaimable existing-row branch, near the reclaimed claim and insertRunningRun
logic, documenting that the superseded attempt’s running run row is
intentionally left for reconcileStaleRunning/stale-TTL healing rather than
immediately marked complete. Do not alter the reclaim behavior.
In `@packages/worker/src/run-records/service.ts`:
- Around line 765-807: Update finishRunRecord to call
dispatchTerminalRunRecordSideEffects instead of maintaining its duplicated
activation and run.error dispatch logic. Preserve the existing behavior by
keeping success-path activation outside any persistedRun guard if it currently
runs when persistence fails, and pass the same handle, status, environment,
persisted row, and waitUntil values used by finishRunRecord.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro Plus
Run ID: 101c33f0-62db-4021-a90d-f9bb496732ca
📒 Files selected for processing (20)
docs/contributing/architecture/data-storage.mddocs/contributing/architecture/invocation-overhead-guardrails.mddocs/contributing/architecture/run-records.mddocs/contributing/disaster-recovery.mdpackages/worker/src/account/export.node.test.tspackages/worker/src/account/export.tspackages/worker/src/account/user-owned-surfaces.tspackages/worker/src/app/retention.tspackages/worker/src/email/inbound.workers.test.tspackages/worker/src/email/system-email-subscriptions.workers.test.tspackages/worker/src/package-invocations/idempotency.tspackages/worker/src/package-invocations/idempotent-module-invocation.tspackages/worker/src/package-invocations/module-execution.tspackages/worker/src/package-invocations/repo.tspackages/worker/src/package-invocations/service.node.test.tspackages/worker/src/platform-feedback/platform-feedback-subscriptions.workers.test.tspackages/worker/src/run-records/invocation-ledger.workers.test.tspackages/worker/src/run-records/run-log-do.tspackages/worker/src/run-records/service.tspackages/worker/src/run-records/types.ts
🚧 Files skipped from review as they are similar to previous changes (6)
- packages/worker/src/app/retention.ts
- packages/worker/src/account/user-owned-surfaces.ts
- docs/contributing/architecture/invocation-overhead-guardrails.md
- docs/contributing/architecture/data-storage.md
- docs/contributing/architecture/run-records.md
- docs/contributing/disaster-recovery.md
| idempotencyKey: `email:${stored.id}:${adminPackage.packageId}:${systemTopic}`, | ||
| }) | ||
| const response = JSON.parse(String(invocation?.['response_json'])) as { | ||
| const response = JSON.parse(String(invocation?.responseJson)) as { |
There was a problem hiding this comment.
🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win
Nullable responseJson coerced with String(...) before JSON.parse in both migrated subscription tests. PackageInvocationLedgerRecord.responseJson is string | null (oversized responses drop the replay cache), so String(null) parses to null and the following property reads throw instead of producing a meaningful assertion failure. The equivalent pattern was already replaced with an explicit guard in packages/worker/src/email/inbound.workers.test.ts.
packages/worker/src/email/system-email-subscriptions.workers.test.ts#L390-L390: throw or fail explicitly wheninvocation?.responseJson == null, thenJSON.parse(responseJson)on the narrowed string.packages/worker/src/platform-feedback/platform-feedback-subscriptions.workers.test.ts#L367-L367: apply the same guard inside the per-package loop before parsing the found invocation'sresponseJson.
📍 Affects 2 files
packages/worker/src/email/system-email-subscriptions.workers.test.ts#L390-L390(this comment)packages/worker/src/platform-feedback/platform-feedback-subscriptions.workers.test.ts#L367-L367
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@packages/worker/src/email/system-email-subscriptions.workers.test.ts` at line
390, Replace String(...) coercion with an explicit null/undefined guard before
JSON.parse in both migrated subscription tests:
packages/worker/src/email/system-email-subscriptions.workers.test.ts:390-390 and
packages/worker/src/platform-feedback/platform-feedback-subscriptions.workers.test.ts:367-367.
Fail or throw when responseJson is absent, then parse the narrowed responseJson
string within each test and per-package loop.
| case 'completed': { | ||
| try { | ||
| const finished = await finishPackageInvocationRecord({ | ||
| env: input.env, | ||
| userId: input.actor.userId, | ||
| expectedUpdatedAt: existing.updated_at, | ||
| staleBefore: new Date( | ||
| now.getTime() - packageInvocationStaleAfterMs, | ||
| ).toISOString(), | ||
| now: reclaimedAt, | ||
| handle: claimed.handle, | ||
| invocationId: claimed.invocationId, | ||
| claimUpdatedAt: claimed.claimUpdatedAt, | ||
| ledgerStatus: 'completed', | ||
| responseJson: boundedResponseJson(outcome.response), | ||
| status: 'success', | ||
| logs: outcome.logs, | ||
| result: outcome.result, | ||
| waitUntil: input.waitUntil, | ||
| }) | ||
| if (!reclaimed) { | ||
| const current = await lookupInvocation() | ||
| if (!current) { | ||
| return buildJsonErrorResponse({ | ||
| if (finished.ledgerUpdated) return outcome.response | ||
| return finished.record | ||
| ? resolveLedgerRecord(finished.record) | ||
| : buildJsonErrorResponse({ | ||
| status: 500, | ||
| code: 'idempotency_conflict_unresolved', | ||
| message: 'Stale package invocation reclaim conflicted.', | ||
| code: 'idempotency_response_unavailable', | ||
| message: 'Package invocation result lost its recovery claim.', | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
| return resolveExistingInvocation({ | ||
| record: current, | ||
| requestHash, | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
| invocationId = existing.id | ||
| claimUpdatedAt = reclaimedAt | ||
| } catch (error) { | ||
| // The module already succeeded; do not poison the key with a | ||
| // stored permanent failure. Return the result to the caller and | ||
| // leave the claim in progress so the stale-reclaim path decides | ||
| // what happens on a later retry. | ||
| console.warn( | ||
| 'package invocation completed-result persistence failed', | ||
| getErrorMessage(error), | ||
| ) | ||
| return outcome.response | ||
| } | ||
|
|
||
| if (!claimUpdatedAt) { | ||
| let insertResult: Awaited<ReturnType<typeof insertPackageInvocationRow>> | ||
| try { | ||
| insertResult = await insertPackageInvocationRow({ | ||
| db: input.env.APP_DB, | ||
| row: { | ||
| id: invocationId, | ||
| userId: input.actor.userId, | ||
| tokenId: input.actor.tokenId, | ||
| packageId: input.savedPackage.id, | ||
| packageKodyId: input.savedPackage.kodyId, | ||
| exportName: input.invocationName, | ||
| idempotencyKey: input.idempotencyKey, | ||
| requestHash, | ||
| source: input.source, | ||
| topic: input.topic, | ||
| status: 'in_progress', | ||
| }, | ||
| }) | ||
| } catch (error) { | ||
| console.error( | ||
| 'package invocation idempotency persistence failed', | ||
| error, | ||
| ) | ||
| return buildJsonErrorResponse({ | ||
| status: 500, | ||
| code: 'idempotency_persistence_failed', | ||
| message: | ||
| 'Unable to persist the package invocation idempotency record. Please retry.', | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
| if (insertResult.inserted) { | ||
| claimUpdatedAt = insertResult.claimUpdatedAt | ||
| } else { | ||
| let current: Awaited<ReturnType<typeof getPackageInvocationByKey>> | ||
| try { | ||
| current = await lookupInvocation() | ||
| } catch (error) { | ||
| console.error('package invocation idempotency lookup failed', error) | ||
| return buildJsonErrorResponse({ | ||
| status: 500, | ||
| code: 'idempotency_lookup_failed', | ||
| message: | ||
| 'Unable to look up the package invocation idempotency record. Please retry.', | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
| if (!current) { | ||
| return buildJsonErrorResponse({ | ||
| } | ||
| case 'failed': { | ||
| try { | ||
| const finished = await finishPackageInvocationRecord({ | ||
| env: input.env, | ||
| userId: input.actor.userId, | ||
| handle: claimed.handle, | ||
| invocationId: claimed.invocationId, | ||
| claimUpdatedAt: claimed.claimUpdatedAt, | ||
| ledgerStatus: 'failed', | ||
| responseJson: boundedResponseJson(outcome.response), | ||
| status: 'error', | ||
| logs: outcome.logs, | ||
| error: outcome.error, | ||
| waitUntil: input.waitUntil, | ||
| }) | ||
| if (finished.ledgerUpdated) return outcome.response | ||
| return finished.record | ||
| ? resolveLedgerRecord(finished.record) | ||
| : buildJsonErrorResponse({ | ||
| status: 500, | ||
| code: 'idempotency_conflict_unresolved', | ||
| message: | ||
| 'Package invocation idempotency insert conflicted but no existing row was found.', | ||
| code: 'idempotency_response_unavailable', | ||
| message: 'Package invocation result lost its recovery claim.', | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
| return resolveExistingInvocation({ | ||
| record: current, | ||
| requestHash, | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
| } | ||
| if (!claimUpdatedAt) { | ||
| throw new Error( | ||
| 'Package invocation claim timestamp was not established.', | ||
| } catch (error) { | ||
| // Best effort; preserve the original invocation error. | ||
| console.warn( | ||
| 'package invocation terminal persistence failed', | ||
| getErrorMessage(error), | ||
| ) | ||
| return outcome.response | ||
| } | ||
| const persistClaimedResult = async ( | ||
| status: 'completed' | 'failed', | ||
| response: PackageInvocationStoredResponse, | ||
| ) => { | ||
| const updated = await updatePackageInvocationResult({ | ||
| db: input.env.APP_DB, | ||
| id: invocationId, | ||
| userId: input.actor.userId, | ||
| status, | ||
| response, | ||
| claimUpdatedAt, | ||
| }) | ||
| if (updated) return response | ||
| const current = await lookupInvocation() | ||
| if (current) { | ||
| return resolveExistingInvocation({ | ||
| record: current, | ||
| requestHash, | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
| return buildJsonErrorResponse({ | ||
| status: 500, | ||
| code: 'idempotency_response_unavailable', | ||
| message: 'Package invocation result lost its recovery claim.', | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
|
|
||
| const outcome = await runSavedPackageModuleOnce({ | ||
| env: input.env, | ||
| baseUrl: input.baseUrl, | ||
| actor: input.actor, | ||
| savedPackage: input.savedPackage, | ||
| invocationName: input.invocationName, | ||
| moduleSelector: input.moduleSelector, | ||
| params: input.params, | ||
| idempotencyKey: input.idempotencyKey, | ||
| invocationId, | ||
| source: input.source, | ||
| topic: input.topic, | ||
| notFoundCode: input.notFoundCode, | ||
| runtimeInvokeDepth: input.runtimeInvokeDepth, | ||
| toolFactories: input.toolFactories, | ||
| waitUntil: input.waitUntil, | ||
| preloadedModuleArtifact: input.preloadedModuleArtifact, | ||
| executorTimeoutMs: input.executorTimeoutMs, | ||
| }) | ||
| switch (outcome.kind) { | ||
| case 'artifact-unavailable': { | ||
| const released = await releasePackageInvocationClaim({ | ||
| db: input.env.APP_DB, | ||
| id: invocationId, | ||
| userId: input.actor.userId, | ||
| claimUpdatedAt, | ||
| }) | ||
| if (!released) { | ||
| const current = await lookupInvocation() | ||
| if (current?.status !== 'in_progress') { | ||
| return current | ||
| ? resolveExistingInvocation({ | ||
| record: current, | ||
| requestHash, | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| : buildJsonErrorResponse({ | ||
| status: 500, | ||
| code: 'idempotency_conflict_unresolved', | ||
| message: | ||
| 'Transient artifact preparation lost its invocation claim.', | ||
| idempotencyKey: input.idempotencyKey, | ||
| }) | ||
| } | ||
| } | ||
| return outcome.response | ||
| } | ||
| case 'completed': { | ||
| try { | ||
| return await persistClaimedResult('completed', outcome.response) | ||
| } catch (error) { | ||
| // The module already succeeded; do not poison the key with a | ||
| // stored permanent failure. Return the result to the caller and | ||
| // leave the claim in progress so the stale-reclaim path decides | ||
| // what happens on a later retry. | ||
| console.warn( | ||
| 'package invocation completed-result persistence failed', | ||
| getErrorMessage(error), | ||
| ) | ||
| return outcome.response | ||
| } | ||
| } | ||
| case 'failed': { | ||
| return await persistClaimedResult('failed', outcome.response).catch( | ||
| (error: unknown) => { | ||
| // Best effort; preserve the original invocation error. | ||
| console.warn( | ||
| 'package invocation terminal persistence failed', | ||
| getErrorMessage(error), | ||
| ) | ||
| return outcome.response | ||
| }, | ||
| ) | ||
| } | ||
| default: { | ||
| const exhaustive: never = outcome | ||
| void exhaustive | ||
| throw new Error('Unhandled package module run outcome.') | ||
| } | ||
| } | ||
| }, | ||
| }) | ||
| } |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
A lost recovery claim after a successful run returns 409 instead of the result.
When finished.ledgerUpdated is false and finished.record is another attempt's in_progress row, resolveLedgerRecord yields invocation_in_progress (409) even though this attempt executed successfully and its side effects already happened; the caller loses a result that is otherwise available in hand. Prefer serving outcome.response in that case and only resolving the record when it is terminal (mismatch/completed/failed) — same reasoning as the catch branch below, which deliberately returns outcome.response rather than poisoning the key.
🐛 Proposed fix (apply to both `completed` and `failed`)
if (finished.ledgerUpdated) return outcome.response
- return finished.record
- ? resolveLedgerRecord(finished.record)
- : buildJsonErrorResponse({
+ if (finished.record && finished.record.status !== 'in_progress') {
+ return resolveLedgerRecord(finished.record)
+ }
+ return finished.record
+ ? outcome.response
+ : buildJsonErrorResponse({
status: 500,
code: 'idempotency_response_unavailable',
message: 'Package invocation result lost its recovery claim.',
idempotencyKey: input.idempotencyKey,
})🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@packages/worker/src/package-invocations/idempotent-module-invocation.ts`
around lines 435 - 503, Update the `completed` and `failed` branches after
`finishPackageInvocationRecord` so a non-updated ledger with another
`in_progress` record returns the already-available `outcome.response`; only pass
terminal records (mismatch, completed, or failed) to `resolveLedgerRecord`.
Preserve the existing fallback error response when no record is available, and
keep the catch behavior unchanged.
| const updateCursor = this.ctx.storage.sql.exec( | ||
| `UPDATE package_invocation_ledger | ||
| SET status = ?, response_json = ?, updated_at = ? | ||
| WHERE id = ? AND status = 'in_progress' AND updated_at = ?`, | ||
| input.status, | ||
| input.responseJson, | ||
| now, | ||
| input.invocationId, | ||
| input.claimUpdatedAt, | ||
| ) | ||
| const ledgerUpdated = updateCursor.rowsWritten > 0 |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift
🧩 Analysis chain
🌐 Web query:
Cloudflare Durable Objects SqlStorageCursor rowsWritten semantics UPDATE DELETE affected rows
💡 Result:
In Cloudflare Durable Objects, the SqlStorageCursor property rowsWritten is a billing-relevant metric that tracks the total number of row-level write operations performed by a SQL query [1][2][3]. Semantics for UPDATE and DELETE: - UPDATE: Every row updated in the table counts as a write [1][2]. Furthermore, if the updated columns are part of any indexes, each index update is counted as an additional row written [1][2][4]. - DELETE: Deletions consistently report rowsWritten as 1 per deleted row [4]. Index cleanup does not appear to be counted toward the rowsWritten metric [4]. - General Write Counting: A row written corresponds to one table row modification plus one additional row for each index that requires an update [1][4]. This includes writes to virtual tables [1][2]. Important considerations for usage: - Billing: The final value of rowsWritten after a query finishes is used for SQL billing [1][2][3]. - Cursor Behavior: rowsWritten may increase incrementally as you iterate through the cursor [1][2][3]. If a transaction rolls back, these writes are not committed, but the cursor may have already reported them based on the temporary state during execution [2][3]. - Tracking Affected Rows: Because the built-in SqlStorage API does not directly return the number of affected rows (like a standard SQL changes function would), developers often use sqlite3 internal functions like changes or last_insert_rowid via supplementary queries to track affected rows accurately [5]. Sources: [1], [2], [3], [4], [5]
Citations:
- 1: https://developers.cloudflare.com/durable-objects/api/sqlite-storage-api/
- 2: https://developers.cloudflare.com/durable-objects/api/sqlite-storage-api/index.md
- 3: https://github.com/cloudflare/cloudflare-docs/blob/production/src/content/docs/durable-objects/api/sqlite-storage-api.mdx
- 4: https://lumenize.com/blog/do-sqlite-write-costs/
- 5: https://github.com/jeffwilde/kysely-durable-objects
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Locate and inspect the relevant function contexts and cursor usage.
fd -a 'run-log-do\.ts$' . | sed 's#^\./##'
file="$(fd 'run-log-do\.ts$' . | head -n1)"
echo "== line count =="
wc -l "$file"
echo "== outline =="
ast-grep outline "$file" --match finishPackageInvocation --view expanded || true
ast-grep outline "$file" --match releasePackageInvocation --view expanded || true
echo "== relevant sections =="
sed -n '1200,1305p' "$file" | nl -ba -v1200
echo "== other rowsWritten/SqlStorageCursor usages =="
rg -n "rowsWritten|SqlStorageCursor|changes\\(|last_insert_rowid" "$file"Repository: kentcdodds/kody
Length of output: 491
🌐 Web query:
Cloudflare Workers Durable Objects SQLite storage changes() SELECT affected rows rowWrite fence
💡 Result:
In Cloudflare Durable Objects with SQLite, tracking affected rows and ensuring data integrity is handled through specific API and architectural patterns rather than reliance on standard SQLite connection-level functions like changes() [1]. 1. Affected Rows: The standard SQLite changes() function is not natively returned by the exec() method in the Cloudflare Durable Objects SQL API [1][2]. Developers often implement tracking of affected rows by manually executing a separate query, such as SELECT changes(), immediately after a mutation within the same transactional context [1]. The exec() method itself provides rowsRead and rowsWritten in its returned object, which are used by Cloudflare for SQL billing and reflect the impact of the query, including auxiliary effects like index updates [2]. 2. RowWrite/Fence: The concept of a "write fence" in Durable Objects is primarily managed through the platform's execution model and consistency guarantees [3]. Durable Objects automatically coalesce synchronous operations into a single transaction if no await statements intervene [3]. To guarantee atomicity for complex operations that involve reading data and writing based on those results, you should use explicit transaction methods (such as ctx.storage.transaction()) [4][3]. The platform ensures consistency via an "output gate" that blocks outgoing responses until pending writes have successfully flushed to storage, preventing race conditions [3]. Because Cloudflare’s exec() method does not automatically return the value of changes(), the common pattern for libraries interacting with Durable Objects is to explicitly issue the additional SELECT changes() call to retrieve the count of rows modified by the preceding statement [1]. [4][1][2][3]
Citations:
- 1: https://github.com/jeffwilde/kysely-durable-objects
- 2: https://github.com/cloudflare/cloudflare-docs/blob/production/src/content/docs/durable-objects/api/sqlite-storage-api.mdx
- 3: https://mintlify.wiki/cloudflare/workerd/advanced/durable-objects-storage
- 4: https://developers.cloudflare.com/durable-objects/api/sqlite-storage-api/
🏁 Script executed:
#!/bin/bash
set -euo pipefail
file="./packages/worker/src/run-records/run-log-do.ts"
echo "== relevant sections =="
sed -n '1188,1302p' "$file" | while IFS= read -r line; do printf '%d\t%s\n' "$((START+NR))" "$line" 2>/dev/null || true; done
echo "== cursor usage context =="
python3 - <<'PY'
from pathlib import Path
p = Path("./packages/worker/src/run-records/run-log-do.ts")
text = p.read_text()
for needle in ["finishPackageInvocation", "releasePackageInvocation"]:
i = text.find(needle)
print(f"\n--- {needle} around {i} ---")
start = max(0, i-1200)
end = min(len(text), i+1800)
for j, line in enumerate(text[start:end].splitlines(), 1):
print(f"{start+j:4}: {line}")
PY
echo "== targeted searches =="
grep -nE "rowsWritten|changes|last_insert_rowid|SELECT id|invocationId|updated_at" "$file" | sed -n '1,220p'Repository: kentcdodds/kody
Length of output: 215
🏁 Script executed:
#!/bin/bash
set -euo pipefail
file="./packages/worker/src/run-records/run-log-do.ts"
start_line=1188
end_line=1302
count=$((end_line - start_line + 1))
sed -n "${start_line},${end_line}p" "$file" | awk '{printf "%d\t%s\n",NR+1187-1,$0}'
echo "== targeted searches =="
grep -nE "rowsWritten|changes|last_insert_rowid|SELECT id|invocationId|updated_at|package_invocation_ledger" "$file" | sed -n '1,220p'Repository: kentcdodds/kody
Length of output: 6547
Avoid treating rowsWritten as the ledger fence.
rowsWritten is a billing-side write-count metric, not an affected-row count. For finishPackageInvocation, even a stale updated_at match can still report writes if Cloudflare’s UPDATE write counting evaluates the SET updated_at = ? assignment before rejecting the row, or counts associated index writes as writes. For releasePackageInvocation, use an explicit affected-row/owned-row check instead of this billing counter before clearing the run row and returning released: true.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@packages/worker/src/run-records/run-log-do.ts` around lines 1231 - 1241,
Replace the rowsWritten-based ledgerUpdated check in finishPackageInvocation
with an explicit verification that the ledger row was updated while matching the
expected invocationId, in_progress status, and claimUpdatedAt fence. Apply the
same owned-row/affected-row validation in releasePackageInvocation before
clearing the run row or returning released: true; do not use rowsWritten as the
ownership decision.
|
Conductor production verification ✅ (deployed at |

Why
package_invocationsis per-user data on the global single-writer D1 database: every keyed invoke paid a claimINSERT, a terminalUPDATE, an account write lease (~5 more D1 statements), plus separate run-record DO round trips. The per-userRunLogDurable Object already indexes runs by idempotency key and executes serially, which gives strictly better claim atomicity thanINSERT OR IGNOREraces. This is the ledger slice of the static-first invocation program (follows #1046, complements #1045/#1047).What changed
run-log-do.ts, schema v5): newpackage_invocation_ledgertable (same shape as the D1 table minususer_id— the DO identity is the user; unique index ontoken_id, package_id, export_name, idempotency_key) with four RPCs:claimPackageInvocation— ledger claim plus eager run-record begin in one awaited DO call; atomic lookup-then-insert; in-place stale reclaim (15 min, matching request hash only).finishPackageInvocation— terminal replay response plus run-record finish in one awaited DO call; ledger update fenced on the claim timestamp; the attempt's own run row is always finished.releasePackageInvocation— deletes an unexecuted claim and itsrunningrun row (transient artifact failure / dual-read fallback hit).getPackageInvocation— conflict-poll lookups.idempotent-module-invocation.ts): rewired onto the DO. All keyed semantics preserved exactly: request-hash mismatch → 409idempotency_mismatch; in-progress polling at 100 ms with a 1 s budget; 15-minute stale reclaim; bounded response replay (maxStoredInvocationResponseJsonBytes) with thereplayedflag; identical error codes. The registry no longer opens a second run record for keyed runs — the claim handle owns it, andmodule-execution.tsenriches it with the published commit after artifact load.packageInvocationLedgerRetentionDays) wired into the existing retention passes and alarm scheduling;in_progressrows are never pruned. The D1package_invocationspolicy entry inapp/retention.tsstays to drain pre-migration rows and is marked legacy / remove in follow-up.exportRunspages runs first, then ledger rows, through one cursor (invocation-ledger:prefix; pre-existing run cursors stay valid), surfaced in therun_recordsexport section (manifest count includes ledger rows); account deletion's existingclearAllpurges ledger rows with run history. Disposition notes extended inaccount/user-owned-surfaces.ts.Decision:
withAccountWriteLeaseis dropped from the keyed pathIts purpose is guarding D1 writes against concurrent account deletion, and no D1 writes remain on this path. This is not a silent drop:
clearRunRecordsafter the lease barrier drains. A DO write racing past that purge is the same residual every run-record write already has (the key-less lean path in feat(invoke): key-less packages.invoke lean path + widen-phase deprecation shims #1046 shipped with exactly this profile, deliberately), and any leaked row self-deletes within the DO's own retention windows.jobs/service.ts,mcp/memory/service.ts,package-registry/service.ts, email, …) that take their own leases — again identical to the merged key-less lean path.Follow-up (deliberately out of this PR)
Once production traffic confirms the DO ledger (a deploy-window's worth of pre-migration keys aged out): remove the dual-read D1 fallback (
getPackageInvocationByKey+ the legacy poll/reclaim consults inidempotent-module-invocation.ts), thepackage_invocationsretention policy entry, the table's account data-target/disposition rows, and migration follow-ups to drop the table (including its never-pruned pre-migrationin_progressrows).Testing
packages/worker/src/run-records/invocation-ledger.workers.test.ts(new, real DO viacloudflare:test): claim+begin/finish+terminal journeys, duplicate claim, atomic stale reclaim, fence behavior for superseded finishes and releases, DO-local 90-day sweep (terminal pruned, in-progress and recent kept), one-cursor export paging across runs and ledger rows,clearAllpurge.packages/worker/src/package-invocations/service.node.test.ts: all pre-existing keyed-semantics tests preserved (replay, mismatch, corruption, persistence failure, stale reclaim, poll, oversized-response drop) against an in-memory RunLog fake whose D1 stand-in throws on any write; new dual-read tests (pre-migration terminal replay + mismatch, fresh legacy in-progress poll-to-replay, stale legacy takeover, stale-DO-claim-defers-to-terminal-legacy, DO-first precedence over a conflicting legacy row).Program report
MERGE-READY
Status: green and mergeable at head
c97cbb6b(mergeStateStatus CLEAN againstmain, includes merged #1045/#1047/#1052).npm run validateexit 0 on the head commit (format, lint, typecheck, node + workers unit, Playwright E2E, MCP E2E, backup:build, primitives/migrations checks).c97cbb6b, legacy fallback now runs on stale reclaims too, with a regression test); Bugbot medium (observer side effects vs ledger fence) answered — deliberately preserves pre-migration semantics; Bugbot low (manifest count) fixed. CodeRabbit: docs wording + test-ordering findings fixed; the "invalid YAML" finding dismissed with a parse proof (also note: main was format-broken by Add GitHub Actions workflow for MCP agent session backfill #1052's workflow file; this branch formats it,22711d7f).System recap — extends run-records and the keyed invocation path (medium risk)
Mode: recap · Base:
main@06aa3d7a· Head:c97cbb6bClassification: extends — the run-records primitive gains the keyed package-invocation idempotency ledger (new DO table + RPCs); the keyed
packages.invokepath is rewired onto it. No new primitive: the ledger is absorbed intorun-records.Primitives touched
run-recordspackage_invocation_ledgertable, combined claim/finish/release RPCs, 90-day DO-local sweep, export cursoraccount-exportrun_recordssection pages ledger rows after runs through the same cursor; manifest counts bothscheduled-cronpackage_invocationssweep kept but marked legacy for the dual-read windowemail,platform-feedbackUnmapped path:
packages/worker/src/package-invocations/(keyedpackages.invokesurface) — extends: claim/finish moved off D1; dual-read fallback added; account write lease dropped.System map
Keyed
packages.invokeclaims and finishes in the per-user RunLog DO; D1 is read-only fallback during the dual-read window.Legend: green = composes (wiring only) · amber = extended by this PR · red = new primitive · gray = context (unchanged, included only when an edge crosses it).
Before / after
INSERT OR IGNORE(+ stale-reclaim UPDATE)claimPackageInvocationRPC (atomic)startRunDO callUPDATE+ separatefinishRunDO callfinishPackageInvocationRPCInvariants
no-per-event-shared-writesstrengthened: per-invoke ledger writes leave the shared D1 writer for per-user DO SQLite.user_idcolumn needed — the DO identity is the user).Summary by CodeRabbit
New Features
Bug Fixes
Documentation