diff --git a/docs/contributing/architecture/data-storage.md b/docs/contributing/architecture/data-storage.md index 868f8da81..6bbeda7bb 100644 --- a/docs/contributing/architecture/data-storage.md +++ b/docs/contributing/architecture/data-storage.md @@ -172,16 +172,16 @@ migration-safe chunked interface: with `section: "storage_runner"` and a `storage_id`, using the same StorageRunner `exportStorage({ pageSize, startAfter })` RPC as the dedicated storage export capability. User meter counters use `section: "user_meter"` and - the `UserMeter.exportCounters` RPC (daily counters plus additive - `storageBytesShadow` on the first page only when present; explicitly - non-authoritative). Mailbox metadata uses `section: "mailbox"` and the - `Mailbox.exportMailbox` RPC. R2 raw MIME, attachment, avatar, and icon objects - use `section: "r2_object"`; each response contains at most one 256 KiB base64 - chunk and an opaque cursor. Each request uses bounded `LIMIT 1` ownership - queries rather than reconstructing inventory. Continuation cursors bind the - source row, object key, size, and ETag; ownership/key mutations and object - overwrites are reported instead of mixing generations. Missing objects are - represented explicitly. + the `UserMeter.exportCounters` RPC (daily counters plus additive shadow fields + on the first page only when present: `storageBytesShadow` and + `packageServiceStatesShadow`; explicitly non-authoritative). Mailbox metadata + uses `section: "mailbox"` and the `Mailbox.exportMailbox` RPC. R2 raw MIME, + attachment, avatar, and icon objects use `section: "r2_object"`; each response + contains at most one 256 KiB base64 chunk and an opaque cursor. Each request + uses bounded `LIMIT 1` ownership queries rather than reconstructing inventory. + Continuation cursors bind the source row, object key, size, and ETag; + ownership/key mutations and object overwrites are reported instead of mixing + generations. Missing objects are represented explicitly. D1 manifest counts use bounded SQL `COUNT(*)` queries. D1 section rows are read with SQL-level keyset pagination: every query orders by the table's `rowid`, @@ -205,13 +205,15 @@ Durable Object export behavior: never pruned by retention. See [Run records](./run-records.md). - `UserMeter` exports daily entitlement counter rows through the `user_meter` section (`exportCounters` RPC; keyset pagination by UTC `day` and `resource`). - The same RPC may return additive `storageBytesShadow` on the first page only - (`startAfter` absent) when the schema-v4 shadow row exists (`null` on later - pages and when never shadowed); section totals count it as one row when - present, but it is explicitly **non-authoritative** — usage and enforcement - read D1 `users.d1_storage_bytes`. Retention is self-enforced inside the DO - (seven UTC days of counter and inbound-delivery-claim rows); shadow - storage-byte state is not time-pruned. See + The same RPC may return additive shadow fields on the first page only + (`startAfter` absent): `storageBytesShadow` when the schema-v4 row exists, and + `packageServiceStatesShadow` when schema-v5 service rows exist (`null` on + later pages and when never shadowed). Section totals count each shadow + inventory once when present, but both are explicitly **non-authoritative** — + usage and enforcement read D1 (`users.d1_storage_bytes` and + `package_service_states` respectively). Retention is self-enforced inside the + DO (seven UTC days of counter and inbound-delivery-claim rows); shadow + storage-byte and package-service liveness rows are not time-pruned. See [Entitlements](./entitlements.md#usermeter-expand-phase). - `Mailbox` exports per-user email metadata (threads, messages, attachments, delivery events) through the account-export `mailbox` section (`exportMailbox` @@ -262,8 +264,9 @@ The schema is defined by migrations in `packages/worker/migrations/`: [Platform accounts](./platform-accounts.md)). `d1_storage_bytes` and `d1_storage_bytes_updated_at` (migration 0122) are the **sole authority** for D1 payload storage-byte read, enforcement, and reconciliation. UserMeter - `storage_bytes_state` (schema v4) is an optional expand-phase shadow only — - see [Entitlements](./entitlements.md#usermeter-expand-phase). Inbound email + `storage_bytes_state` (schema v4) and `package_service_states` (schema v5) are + optional expand-phase shadows only — see + [Entitlements](./entitlements.md#usermeter-expand-phase). Inbound email routing does not reverse-resolve stable ids at all — it uses the indexed username lookup (`findPublicUserIdentityByUsername`). Contextless paths resolve stable ids with one indexed point read on `users.stable_user_id` (for @@ -293,8 +296,12 @@ The schema is defined by migrations in `packages/worker/migrations/`: [Run records](./run-records.md)). Execution history rows live in the same DO. - `package_service_states` (`0095-package-service-states.sql`): authoritative per-service liveness projection (`running` / `idle` / `stopped` / `error`) for - entitlement concurrency. Upserted and heartbeaten by the package-service - Durable Object; not derived from run history. + entitlement concurrency, discovery, and export/deletion inventory. Upserted + and heartbeaten (1h) by the `PackageServiceInstance` Durable Object; running + counts treat rows stale after 24h without a fresh heartbeat. Not derived from + run history. Expand-phase slice 4 Phase A also best-effort shadows each row + into the per-user `UserMeter` DO (schema v5); D1 remains sole authority in + that slice — see [Entitlements](./entitlements.md#usermeter-expand-phase). - `entity_sources`: durable mapping from user-facing entities to Artifacts repos and their latest published commit - `saved_packages`: package metadata/search projection derived from published @@ -519,10 +526,12 @@ Daily rate-style entitlement counters and inbound email delivery-id idempotency live in a per-user `UserMeter` Durable Object with SQLite (`packages/worker/src/entitlements/user-meter-do.ts`). Schema v4 adds an optional `storage_bytes_state` singleton as a **best-effort shadow** of D1 -`users.d1_storage_bytes` for future cutover — D1 remains sole authority for -reads, reserves, and reconciliation in the current additive slice. The Worker -binding is `USER_METER` (class `UserMeter`; Wrangler SQLite migration tag `v21` -via `new_sqlite_classes` in `packages/worker/wrangler.jsonc`). +`users.d1_storage_bytes`. Schema v5 adds an optional `package_service_states` +table as a **best-effort shadow** of D1 `package_service_states` for future +cutover — D1 remains sole authority for reads, running counts, discovery, and +`service_start` enforcement in expand-phase slice 4 Phase A. The Worker binding +is `USER_METER` (class `UserMeter`; Wrangler SQLite migration tag `v21` via +`new_sqlite_classes` in `packages/worker/wrangler.jsonc`). Naming matches `RunLog` and `JobManager`: one object per untrimmed stable MCP `userId` via `userMeterDurableObjectName(userId)` → `idFromName(userId)` in @@ -530,7 +539,7 @@ Naming matches `RunLog` and `JobManager`: one object per untrimmed stable MCP column inside the DO because the object identity is the user. SQLite ownership (schema version tracked in `user_meter_meta`; current version -**4**): +**5**): - `daily_counters` — authoritative UTC-day counters for `email_sends_per_day`, `email_receives_per_day`, `execute_calls_per_day`, and @@ -547,12 +556,23 @@ SQLite ownership (schema version tracked in `user_meter_meta`; current version and optional reconcile shadows; never read for enforcement or usage display. StorageRunner bucket estimates stay outside this row (see [Entitlements](./entitlements.md#usermeter-expand-phase)). +- `package_service_states` — per-service **shadow** of D1 liveness rows + (`package_id`, `service_name`, `status`, `started_at`, `source_updated_at`, + monotonic `revision`, `updated_at`; primary key `(package_id, service_name)`). + Added in schema v5. Populated by best-effort dual-writes from + `PackageServiceInstance` on every D1 projection/delete; monotonic on + `source_updated_at`. Never read for enforcement, running counts, discovery, or + usage display in expand-phase slice 4 Phase A. Cutover-support RPCs + (`listPackageServiceStates`, `countRunningPackageServices`, + `bootstrapPackageServiceStates`) mirror D1 semantics for future parity review + only. Retention is self-enforced inside the DO: every read/write path opportunistically deletes counter and claim rows older than seven UTC days (`userMeterDailyCounterRetentionDays`). Enforcement only needs the current day; the window covers timezone edge cases, recent account exports, and inbound -retries. Shadow storage-byte state is not time-pruned. +retries. Shadow storage-byte and package-service liveness rows are not +time-pruned. **Expand-phase D1 mirrors (daily counters only):** enforcement and point reads are authoritative in UserMeter for daily counters. D1 @@ -571,11 +591,12 @@ writes cannot overwrite newer state. See daily paths never read D1 for enforcement. Account deletion calls `UserMeter.purge()` (one RPC per user, no D1 id scan; -`deleteAll` clears counters, claims, and any shadow storage state). Account -export pages `UserMeter.exportCounters` through the `user_meter` manifest -section / `account_export_section` (daily counters plus additive -`storageBytesShadow` on the first page only when present; shadow field is -non-authoritative). +`deleteAll` clears counters, claims, and all shadow state including storage +bytes and package-service liveness). Account export pages +`UserMeter.exportCounters` through the `user_meter` manifest section / +`account_export_section` (daily counters plus additive `storageBytesShadow` and +`packageServiceStatesShadow` on the first page only when present; shadow fields +are non-authoritative). ## Durable Objects (`Mailbox`) @@ -718,7 +739,9 @@ storage homes as follows: `package:{encodeURIComponent(packageId)}` via `buildPackageStorageId` / `packageStorage()`. Shared durable data for every package surface. - **Package coordination** — `PackageServiceInstance` DO holds lifecycle and - alarms only; durable data stays in package storage. App facets and + alarms only; durable data stays in package storage. Each lifecycle projection + dual-writes D1 `package_service_states` (authority) and a best-effort + UserMeter shadow (expand-phase slice 4 Phase A). App facets and package-internal DO namespaces are extra StorageRunner buckets under the package id, not a general actor model. - **Package jobs** — schedule metadata in D1 `jobs`; run-local scratch in @@ -741,8 +764,8 @@ via `durableObjectNameFromParts`); domain helpers such as package activation counters/milestones). See [Run records](./run-records.md). - `UserMeter` — `userMeterDurableObjectName(userId)` → `idFromName(userId)`. One daily-entitlement meter DO per user (untrimmed stable id, same as `RunLog`), - plus optional schema-v4 D1 storage-byte shadow. See - [Entitlements](./entitlements.md#usermeter-expand-phase). + plus optional schema-v4 D1 storage-byte shadow and schema-v5 package-service + liveness shadow. See [Entitlements](./entitlements.md#usermeter-expand-phase). - `StripePlanRefresh` — `stripePlanRefreshDurableObjectName(userId)` → `idFromName(userId)`. One ephemeral, one-shot reconciliation alarm per user; checkout and subscription webhook activity arm it as a backstop to the @@ -1079,7 +1102,8 @@ to `durableObjectNameFromParts`). - `JobManager`: `idFromName(userId)` (no trim). - `RunLog`: `idFromName(userId)` (no trim); one execution-history DO per user. - `UserMeter`: `idFromName(userId)` (no trim); one daily-entitlement meter DO - per user, plus optional schema-v4 D1 storage-byte shadow. + per user, plus optional schema-v4 D1 storage-byte shadow and schema-v5 + package-service liveness shadow. - `StripePlanRefresh`: `idFromName(userId)` (no trim); one ephemeral billing reconciliation alarm DO per user. - `Mailbox`: `idFromName(userId)` (no trim); one email-metadata DO per user. diff --git a/docs/contributing/architecture/entitlements.md b/docs/contributing/architecture/entitlements.md index 4fee1293c..f72c68190 100644 --- a/docs/contributing/architecture/entitlements.md +++ b/docs/contributing/architecture/entitlements.md @@ -124,6 +124,12 @@ storage layout and naming are documented in [Data storage](./data-storage.md). `storage_bytes_state` singleton as a **best-effort shadow** for future cutover — it does not drive reads, reserves, or reconciliation in this additive slice. +**D1 package service liveness** (`package_service_states`) stays authoritative +for running counts, discovery, and `service_start` enforcement. UserMeter schema +v5 adds an optional per-service shadow table as **best-effort future-cutover +support only** — see +[Package service liveness — UserMeter shadow](#package-service-liveness--usermeter-shadow-expand-phase-slice-4-phase-a). + StorageRunner bucket `estimatedBytes` and the per-bucket inventory in `user_storage_buckets` stay a **separate** quota component. StorageRunner write chokepoints pass `getCurrent` as a check-only composed total (D1 payload bytes @@ -215,11 +221,78 @@ successful absolute reconciliation also advances the expand-phase shadow. `readUserD1StorageBytes` only. **Account export and purge:** `UserMeter.exportCounters` may return additive -non-authoritative `storageBytesShadow` on the first page only (`startAfter` -absent) when the shadow row exists; subsequent pages return `null` so paged -consumers never double-count it (still counted once in the `user_meter` section -total). `UserMeter.purge()` clears counters, inbound delivery claims, and any -shadow storage state via `deleteAll`. +non-authoritative shadow fields on the first page only (`startAfter` absent): +`storageBytesShadow` when the schema-v4 row exists, and +`packageServiceStatesShadow` when schema-v5 service rows exist. Subsequent pages +return `null` for each shadow so paged consumers never double-count them +(section totals still count each shadow inventory once when present). +`UserMeter.purge()` clears counters, inbound delivery claims, and all shadow +state (storage bytes and package-service liveness) via `deleteAll`. + +### Package service liveness — UserMeter shadow (expand phase slice 4, Phase A) + +D1 `package_service_states` remains the **sole authority** for running-service +**count**, **discovery**, and **`service_start` enforcement** in this PR. +Nothing in Phase A switches those reads or the `assertWithinEntitlement` path +for `package_services` / `persistent_package_services`. + +UserMeter schema **v5** adds a per-service `package_service_states` table inside +the DO as **best-effort shadow / future-cutover support only** (`status`, +`started_at`, monotonic `source_updated_at` from the D1 projection timestamp, +`revision`, `updated_at`). User scope is the DO identity — there is no `user_id` +column. Shadow rows are never read for usage display, entitlement enforcement, +or account-deletion inventory in this slice. + +**Dual-write from `PackageServiceInstance`:** every D1 projection also attempts +a best-effort UserMeter shadow on the same lifecycle surface: + +- lifecycle transitions and warm-start restore after upgrades + (`projectServiceStateToD1`) +- running-service heartbeat alarms (1h `packageServiceStateHeartbeatMs`, + unchanged) +- stop, error, and idle projections that clear `running` +- purge (`deleteProjectedServiceState` deletes D1 then shadow before + `deleteAll`) + +D1 upsert/delete runs first; shadow RPCs are optional when `USER_METER` is +unbound and failures log `package-service-user-meter-shadow-failed` without +affecting the service path. Shadow upserts reject stale/out-of-order writes when +`sourceUpdatedAt` is older than the existing shadow row so cold bootstrap cannot +clobber fresher state. + +**Timing unchanged:** live services heartbeat D1 `updated_at` every **1 hour** +(`packageServiceStateHeartbeatMs`). Running counts still treat rows as stale +after **24 hours** without a fresh heartbeat (`packageServiceStateStaleMs` in +`entitlements/service.ts`). The UserMeter cutover-support RPC +`countRunningPackageServices` uses the same 24h window on shadow +`source_updated_at` but is **not** wired to enforcement in Phase A. + +**Account export:** `UserMeter.exportCounters` returns additive +`packageServiceStatesShadow` on the first page only (`startAfter` absent); later +pages return `null`. Section totals count the shadow inventory once when +present; the field is explicitly non-authoritative — authoritative liveness +remains on D1. + +**Account purge:** `UserMeter.purge()` clears package-service shadow rows with +the rest of DO state via `deleteAll`. + +### Future package-service authority flip (contract follow-up) + +A separate **high-risk contract PR** — not a merge blocker for Phase A — will +flip running-service count/discovery/enforcement into UserMeter only after: + +1. at least one full **24h stale-window soak** with shadow/D1 parity review, and +2. a **cold-bootstrap design** for accounts whose DO shadow is empty while D1 + still holds rows (`bootstrapPackageServiceStates` / equivalent). + +Until that flip, shadow divergence is acceptable; D1 remains the contract. **D1 +likely stays the enumeration index** for account export, deletion, and admin +discovery until an alternate inventory exists — UserMeter would become the +running-count authority first, not a wholesale replacement for every D1 reader. + +**Remaining expand roadmap:** slice 4 Phase A (this shadow slice) is additive +only; slice 5 — account-deletion write fencing — follows independently. Storage +authority flip remains the separate contract follow-up above. ### Future storage authority flip (contract follow-up) @@ -231,9 +304,10 @@ Only then do reads, reserves, and reconciliation switch to UserMeter-first with D1 as mirror/cursor. Until that flip, shadow divergence is acceptable; D1 remains the contract. -**Remaining UserMeter expand roadmap:** slice 4 — `package_service_states` -running counts (service liveness); slice 5 — account-deletion write fencing. -Storage authority flip is tracked separately as the contract follow-up above. +**Remaining UserMeter expand roadmap:** slice 4 Phase A — +`package_service_states` UserMeter shadow (this PR; D1 authority unchanged); +slice 5 — account-deletion write fencing. Package-service and storage authority +flips are separate high-risk contract follow-ups after soak/parity review. **Daily-counter mirror retirement:** dropping D1 `entitlement_daily_counters` waits until reporting-off-D1 work merges and mirror parity is verified in @@ -379,12 +453,14 @@ Rules: running package services) are counted **directly from their source D1 tables at the enforcement point** via built-in counters in `service.ts`. They do not depend on any metering or rollup tables. Running package services are counted - from `package_service_states` (status `running` and freshly heartbeaten), not - from run-history rows — see [Run records](./run-records.md) - (`state-vs-history`). **Concurrent workflows** are authoritative in per-user - RunLog `workflow_projections`: create reserves atomically via - `reserveWorkflowProjectionSlot`, and usage readers call - `countActiveWorkflowProjections` through + from D1 `package_service_states` (status `running` and freshly heartbeaten; 1h + heartbeat, 24h staleness), not from run-history rows — see + [Run records](./run-records.md) (`state-vs-history`). Expand-phase slice 4 + Phase A dual-writes the same projection into UserMeter as a non-authoritative + shadow; enforcement and `service_start` still read D1 only. **Concurrent + workflows** are authoritative in per-user RunLog `workflow_projections`: + create reserves atomically via `reserveWorkflowProjectionSlot`, and usage + readers call `countActiveWorkflowProjections` through `readCurrentEntitlementResourceUsage`. Expand-phase D1 `workflow_runs` is a compatibility mirror only. - **Rate-style limits** (email sends/receives per day, execute calls per day, diff --git a/docs/contributing/architecture/primitives.yaml b/docs/contributing/architecture/primitives.yaml index c224fd876..776bcfe2f 100644 --- a/docs/contributing/architecture/primitives.yaml +++ b/docs/contributing/architecture/primitives.yaml @@ -516,7 +516,9 @@ primitives: name: User meter summary: Per-user UserMeter DO SQLite for daily entitlement counters, inbound - delivery idempotency, and expand-phase D1 storage-byte shadow. + delivery idempotency, expand-phase D1 storage-byte shadow, and + expand-phase package-service liveness shadow (D1 authority unchanged in + slice 4 Phase A). code: - packages/worker/src/entitlements/user-meter-do.ts - packages/worker/src/entitlements/user-meter-client.ts diff --git a/packages/worker/src/account/export.node.test.ts b/packages/worker/src/account/export.node.test.ts index 657b5a1f4..22fb0971d 100644 --- a/packages/worker/src/account/export.node.test.ts +++ b/packages/worker/src/account/export.node.test.ts @@ -1578,6 +1578,18 @@ test('account export includes user_meter counters, pages them, and warns on trun updatedAt: '2026-07-31T03:00:00.000Z', mirrorUpdatedAt: 'r/00000000000000000003', } + const packageServiceStatesShadow = [ + { + packageId: 'pkg-1', + serviceName: 'worker', + status: 'running' as const, + startedAt: '2026-07-31T03:00:00.000Z', + sourceUpdatedAt: '2026-07-31T03:05:00.000Z', + revision: 2, + updatedAt: '2026-07-31T03:05:00.000Z', + mirrorUpdatedAt: 'r/00000000000000000002', + }, + ] const exportCounters = vi.fn( async (input: { pageSize?: number; startAfter?: string | null }) => { const pageSize = input.pageSize ?? 100 @@ -1593,6 +1605,9 @@ test('account export includes user_meter counters, pages them, and warns on trun return { counters: page, storageBytesShadow: isFirstPage ? storageBytesShadow : null, + packageServiceStatesShadow: isFirstPage + ? packageServiceStatesShadow + : null, nextStartAfter: truncated ? `${page.at(-1)!.day}:${page.at(-1)!.resource}` : null, @@ -1615,10 +1630,11 @@ test('account export includes user_meter counters, pages them, and warns on trun mcpUserId: 'user-aaa', }) expect(idFromName).toHaveBeenCalledWith('user-aaa') - expect(accountExport.manifest.sections.user_meter?.count).toBe(4) + expect(accountExport.manifest.sections.user_meter?.count).toBe(5) expect(accountExport.durableObjects.userMeter).toEqual({ counters, storageBytesShadow, + packageServiceStatesShadow, nextStartAfter: null, truncated: false, }) @@ -1637,6 +1653,7 @@ test('account export includes user_meter counters, pages them, and warns on trun }) expect(first.items).toEqual(counters.slice(0, 2)) expect(first.storageBytesShadow).toEqual(storageBytesShadow) + expect(first.packageServiceStatesShadow).toEqual(packageServiceStatesShadow) expect(first.truncated).toBe(true) expect(first.nextStartAfter).toBe('2026-07-30:execute_calls_per_day') @@ -1650,6 +1667,7 @@ test('account export includes user_meter counters, pages them, and warns on trun }) expect(second.items).toEqual(counters.slice(2)) expect(second.storageBytesShadow).toBeNull() + expect(second.packageServiceStatesShadow).toBeNull() expect(second.truncated).toBe(false) expect(second.nextStartAfter).toBeNull() expect(exportCounters).toHaveBeenCalledWith( @@ -1662,6 +1680,7 @@ test('account export includes user_meter counters, pages them, and warns on trun exportCounters.mockImplementation(async () => ({ counters: [counters[0]!], storageBytesShadow: null, + packageServiceStatesShadow: null, nextStartAfter: 'cursor-more', truncated: true, })) diff --git a/packages/worker/src/account/export.ts b/packages/worker/src/account/export.ts index a5796c637..b3a31eb93 100644 --- a/packages/worker/src/account/export.ts +++ b/packages/worker/src/account/export.ts @@ -209,7 +209,13 @@ export type AccountExportArtifactRepo = { type RunRecordsExportPayload = Awaited> function countUserMeterExportEntries(result: UserMeterExportResult): number { - return result.counters.length + (result.storageBytesShadow == null ? 0 : 1) + return ( + result.counters.length + + (result.storageBytesShadow == null ? 0 : 1) + + (result.packageServiceStatesShadow == null + ? 0 + : result.packageServiceStatesShadow.length) + ) } function countRunRecordsExportEntries( @@ -293,6 +299,12 @@ export type AccountExportSectionResult = { * first `user_meter` page (`startAfter` absent); later pages omit it. */ storageBytesShadow?: UserMeterExportResult['storageBytesShadow'] + /** + * Non-authoritative UserMeter package-service shadow rows. Present only on + * the first `user_meter` page (`startAfter` absent); later pages set it to + * `null`. + */ + packageServiceStatesShadow?: UserMeterExportResult['packageServiceStatesShadow'] } function normalizePageSize(pageSize: number | undefined) { @@ -1989,6 +2001,7 @@ export async function readAccountExportSection(input: { section: input.section, items: page.counters, storageBytesShadow: page.storageBytesShadow, + packageServiceStatesShadow: page.packageServiceStatesShadow, truncated: page.truncated, nextStartAfter: page.nextStartAfter, pageSize, diff --git a/packages/worker/src/account/user-owned-surfaces.ts b/packages/worker/src/account/user-owned-surfaces.ts index 3e50742f6..065e34438 100644 --- a/packages/worker/src/account/user-owned-surfaces.ts +++ b/packages/worker/src/account/user-owned-surfaces.ts @@ -132,7 +132,7 @@ export const accountUserOwnedDurableObjectSurfaces: ReadonlyArray).includes( + status, + ) +} /** * Legacy D1 mirror `updated_at` token. Lexicographic order matches revision @@ -110,6 +133,42 @@ export type UserMeterStorageBytesSetResult = UserMeterStorageBytesReadyState & { created: boolean } +export type UserMeterPackageServiceState = { + packageId: string + serviceName: string + status: UserMeterPackageServiceStatus + startedAt: string | null + sourceUpdatedAt: string + revision: number + updatedAt: string + mirrorUpdatedAt: string +} + +export type UserMeterPackageServiceUpsertResult = { + applied: boolean + created: boolean + state: UserMeterPackageServiceState +} + +export type UserMeterPackageServiceDeleteResult = { + deleted: boolean +} + +export type UserMeterPackageServiceListResult = { + states: Array + nextStartAfter: string | null + truncated: boolean +} + +export type UserMeterPackageServiceCountResult = { + count: number +} + +export type UserMeterPackageServiceBootstrapResult = { + applied: number + skipped: number +} + export type UserMeterExportResult = { counters: Array /** @@ -118,6 +177,12 @@ export type UserMeterExportResult = { * paged consumers never double-count the singleton shadow. */ storageBytesShadow: UserMeterStorageBytesState | null + /** + * Non-authoritative D1 `package_service_states` shadow rows. Emitted only + * on the first export page (`startAfter` absent); subsequent pages return + * `null` so paged consumers never double-count the inventory. + */ + packageServiceStatesShadow: Array | null nextStartAfter: string | null truncated: boolean } @@ -141,6 +206,21 @@ type ExportCursor = { resource: string } +type PackageServiceExportCursor = { + packageId: string + serviceName: string +} + +type PackageServiceSqlRow = { + package_id: string + service_name: string + status: string + started_at: string | null + source_updated_at: string + revision: number + updated_at: string +} + function assertDailyResource(resource: string): DailyEntitlementResource { if (!isDailyEntitlementResource(resource)) { throw new Error( @@ -184,6 +264,52 @@ function assertInboundDeliveryId(deliveryId: string): string { return deliveryId } +function assertPackageServiceId(label: string, value: string): string { + if ( + typeof value !== 'string' || + value.length === 0 || + value.length > maxPackageServiceIdLength + ) { + throw new Error( + `UserMeter ${label} must be a non-empty string up to ${maxPackageServiceIdLength} characters.`, + ) + } + return value +} + +function assertPackageServiceStatus( + status: string, +): UserMeterPackageServiceStatus { + if (!isUserMeterPackageServiceStatus(status)) { + throw new Error( + `UserMeter package service status must be one of ${packageServiceShadowStatuses.join(', ')}; got ${JSON.stringify(status)}.`, + ) + } + return status +} + +function assertSourceUpdatedAt(sourceUpdatedAt: string): string { + if ( + typeof sourceUpdatedAt !== 'string' || + !packageServiceSourceUpdatedAtPattern.test(sourceUpdatedAt) + ) { + throw new Error( + 'UserMeter package service sourceUpdatedAt must be an ISO-8601 UTC timestamp.', + ) + } + // Pattern-first: avoid Invalid Date → RangeError from toISOString(). + const parsedMs = Date.parse(sourceUpdatedAt) + if ( + !Number.isFinite(parsedMs) || + new Date(parsedMs).toISOString() !== sourceUpdatedAt + ) { + throw new Error( + 'UserMeter package service sourceUpdatedAt must be an ISO-8601 UTC timestamp.', + ) + } + return sourceUpdatedAt +} + function retentionCutoffDay(now: Date): string { const cutoff = new Date(now) cutoff.setUTCDate( @@ -208,6 +334,26 @@ function decodeExportCursor(startAfter: string): ExportCursor | null { } } +function encodePackageServiceCursor(cursor: PackageServiceExportCursor) { + return JSON.stringify([cursor.packageId, cursor.serviceName]) +} + +function decodePackageServiceCursor( + startAfter: string, +): PackageServiceExportCursor | null { + try { + const parsed = JSON.parse(startAfter) as unknown + if (!Array.isArray(parsed) || parsed.length !== 2) return null + const [packageId, serviceName] = parsed + if (typeof packageId !== 'string' || typeof serviceName !== 'string') { + return null + } + return { packageId, serviceName } + } catch { + return null + } +} + function normalizePageSize(pageSize: number | undefined) { const requested = typeof pageSize === 'number' && Number.isFinite(pageSize) @@ -301,6 +447,23 @@ class UserMeterBase extends DurableObject { updated_at TEXT NOT NULL ) `) + // Expand-phase shadow of D1 `package_service_states` (schema v5). + this.ctx.storage.sql.exec(` + CREATE TABLE IF NOT EXISTS package_service_states ( + package_id TEXT NOT NULL, + service_name TEXT NOT NULL, + status TEXT NOT NULL, + started_at TEXT, + source_updated_at TEXT NOT NULL, + revision INTEGER NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (package_id, service_name) + ) + `) + this.ctx.storage.sql.exec( + `CREATE INDEX IF NOT EXISTS idx_package_service_states_status_source + ON package_service_states (status, source_updated_at)`, + ) this.ctx.storage.sql.exec( `INSERT INTO user_meter_meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value`, @@ -772,6 +935,276 @@ class UserMeterBase extends DurableObject { } } + private packageServiceStateFromRow( + row: PackageServiceSqlRow, + ): UserMeterPackageServiceState { + const revision = Math.max(0, Number(row.revision ?? 0)) + const status = assertPackageServiceStatus(String(row.status)) + return { + packageId: String(row.package_id), + serviceName: String(row.service_name), + status, + startedAt: + status === 'running' && row.started_at != null + ? String(row.started_at) + : null, + sourceUpdatedAt: String(row.source_updated_at), + revision, + updatedAt: String(row.updated_at), + mirrorUpdatedAt: userMeterMirrorUpdatedAtToken(revision), + } + } + + private readPackageServiceRow( + packageId: string, + serviceName: string, + ): UserMeterPackageServiceState | null { + const row = this.ctx.storage.sql + .exec( + `SELECT package_id, service_name, status, started_at, + source_updated_at, revision, updated_at + FROM package_service_states + WHERE package_id = ? AND service_name = ?`, + packageId, + serviceName, + ) + .toArray()[0] + if (!row) return null + return this.packageServiceStateFromRow(row) + } + + private listAllPackageServiceRows(): Array { + const rows = this.ctx.storage.sql + .exec( + `SELECT package_id, service_name, status, started_at, + source_updated_at, revision, updated_at + FROM package_service_states + ORDER BY package_id ASC, service_name ASC`, + ) + .toArray() + return rows.map((row) => this.packageServiceStateFromRow(row)) + } + + /** + * Expand-phase shadow upsert (monotonic on `sourceUpdatedAt`). Stale writes + * are rejected; expand-phase enforcement still reads D1. + */ + async upsertPackageServiceState(input: { + packageId: string + serviceName: string + status: string + startedAt?: string | null + sourceUpdatedAt: string + updatedAt?: string + }): Promise { + const packageId = assertPackageServiceId('packageId', input.packageId) + const serviceName = assertPackageServiceId('serviceName', input.serviceName) + const status = assertPackageServiceStatus(input.status) + const sourceUpdatedAt = assertSourceUpdatedAt(input.sourceUpdatedAt) + const updatedAt = + typeof input.updatedAt === 'string' && input.updatedAt.length > 0 + ? input.updatedAt + : sourceUpdatedAt + const startedAt = + status === 'running' && + typeof input.startedAt === 'string' && + input.startedAt.length > 0 + ? input.startedAt + : null + + const existing = this.readPackageServiceRow(packageId, serviceName) + if (existing && existing.sourceUpdatedAt > sourceUpdatedAt) { + return { + applied: false, + created: false, + state: existing, + } + } + if (!existing) { + this.ctx.storage.sql.exec( + `INSERT INTO package_service_states ( + package_id, service_name, status, started_at, + source_updated_at, revision, updated_at + ) VALUES (?, ?, ?, ?, ?, 1, ?)`, + packageId, + serviceName, + status, + startedAt, + sourceUpdatedAt, + updatedAt, + ) + const row = this.readPackageServiceRow(packageId, serviceName) + if (!row) { + throw new Error( + 'UserMeter upsertPackageServiceState failed to materialize shadow row.', + ) + } + return { applied: true, created: true, state: row } + } + + const nextRevision = existing.revision + 1 + this.ctx.storage.sql.exec( + `UPDATE package_service_states + SET status = ?, + started_at = ?, + source_updated_at = ?, + revision = ?, + updated_at = ? + WHERE package_id = ? AND service_name = ? AND revision = ?`, + status, + startedAt, + sourceUpdatedAt, + nextRevision, + updatedAt, + packageId, + serviceName, + existing.revision, + ) + const row = this.readPackageServiceRow(packageId, serviceName) + if (!row) { + throw new Error( + 'UserMeter upsertPackageServiceState lost the shadow row.', + ) + } + return { applied: true, created: false, state: row } + } + + /** Expand-phase shadow delete; discovery still uses D1. */ + async deletePackageServiceState(input: { + packageId: string + serviceName: string + }): Promise { + const packageId = assertPackageServiceId('packageId', input.packageId) + const serviceName = assertPackageServiceId('serviceName', input.serviceName) + const cursor = this.ctx.storage.sql.exec( + `DELETE FROM package_service_states + WHERE package_id = ? AND service_name = ?`, + packageId, + serviceName, + ) + return { deleted: cursor.rowsWritten > 0 } + } + + /** Cutover-support paged shadow list; expand-phase discovery uses D1. */ + async listPackageServiceStates(input: { + pageSize?: number + startAfter?: string | null + }): Promise { + const pageSize = normalizePageSize(input.pageSize) + const cursor = + typeof input.startAfter === 'string' && input.startAfter.length > 0 + ? decodePackageServiceCursor(input.startAfter) + : null + const rows = ( + cursor + ? this.ctx.storage.sql.exec( + `SELECT package_id, service_name, status, started_at, + source_updated_at, revision, updated_at + FROM package_service_states + WHERE package_id > ? + OR (package_id = ? AND service_name > ?) + ORDER BY package_id ASC, service_name ASC + LIMIT ?`, + cursor.packageId, + cursor.packageId, + cursor.serviceName, + pageSize + 1, + ) + : this.ctx.storage.sql.exec( + `SELECT package_id, service_name, status, started_at, + source_updated_at, revision, updated_at + FROM package_service_states + ORDER BY package_id ASC, service_name ASC + LIMIT ?`, + pageSize + 1, + ) + ).toArray() + const truncated = rows.length > pageSize + const pageRows = truncated ? rows.slice(0, pageSize) : rows + const states = pageRows.map((row) => this.packageServiceStateFromRow(row)) + const last = pageRows[pageRows.length - 1] + return { + states, + nextStartAfter: + truncated && last + ? encodePackageServiceCursor({ + packageId: String(last.package_id), + serviceName: String(last.service_name), + }) + : null, + truncated, + } + } + + /** + * Cutover-support running count from the shadow table. Expand-phase + * enforcement still uses D1 `countRunningPackageServices`. + */ + async countRunningPackageServices(input: { + staleAfterMs?: number + excludeService?: { packageId: string; serviceName: string } + now?: string + }): Promise { + const now = input.now ? new Date(input.now) : new Date() + const safeNow = Number.isNaN(now.valueOf()) ? new Date() : now + const staleAfterMs = + typeof input.staleAfterMs === 'number' && + Number.isFinite(input.staleAfterMs) && + input.staleAfterMs >= 0 + ? Math.trunc(input.staleAfterMs) + : userMeterPackageServiceStateStaleMs + const freshAfter = new Date(safeNow.valueOf() - staleAfterMs).toISOString() + const exclusion = input.excludeService + const packageId = exclusion + ? assertPackageServiceId('packageId', exclusion.packageId) + : null + const serviceName = exclusion + ? assertPackageServiceId('serviceName', exclusion.serviceName) + : null + const row = ( + packageId && serviceName + ? this.ctx.storage.sql.exec<{ count: number }>( + `SELECT COUNT(*) AS count + FROM package_service_states + WHERE status = 'running' + AND source_updated_at >= ? + AND NOT (package_id = ? AND service_name = ?)`, + freshAfter, + packageId, + serviceName, + ) + : this.ctx.storage.sql.exec<{ count: number }>( + `SELECT COUNT(*) AS count + FROM package_service_states + WHERE status = 'running' + AND source_updated_at >= ?`, + freshAfter, + ) + ).toArray()[0] + return { count: Math.max(0, Number(row?.count ?? 0)) } + } + + /** Cutover-support bulk seed; same monotonic guard as upsert. */ + async bootstrapPackageServiceStates(input: { + states: ReadonlyArray<{ + packageId: string + serviceName: string + status: string + startedAt?: string | null + sourceUpdatedAt: string + updatedAt?: string + }> + }): Promise { + let applied = 0 + let skipped = 0 + for (const state of input.states) { + const result = await this.upsertPackageServiceState(state) + if (result.applied) applied += 1 + else skipped += 1 + } + return { applied, skipped } + } + async purge(): Promise<{ ok: true }> { await this.ctx.blockConcurrencyWhile(async () => { await this.ctx.storage.deleteAll() @@ -842,11 +1275,14 @@ class UserMeterBase extends DurableObject { mirrorUpdatedAt: userMeterMirrorUpdatedAtToken(revision), }) } - // Shadow is a singleton outside keyset paging — emit it once on the - // first page only so section totals and multi-page consumers do not - // double-count. - const includeStorageShadow = cursor == null - const storageRow = includeStorageShadow ? this.readStorageRow() : null + // Shadow inventories sit outside counter keyset paging — emit them once + // on the first page only so section totals and multi-page consumers do + // not double-count. + const includeShadows = cursor == null + const storageRow = includeShadows ? this.readStorageRow() : null + const packageServiceStatesShadow = includeShadows + ? this.listAllPackageServiceRows() + : null const last = pageRows[pageRows.length - 1] return { counters, @@ -858,6 +1294,7 @@ class UserMeterBase extends DurableObject { mirrorUpdatedAt: userMeterMirrorUpdatedAtToken(storageRow.revision), } : null, + packageServiceStatesShadow, nextStartAfter: truncated && last ? encodeExportCursor({ @@ -923,6 +1360,42 @@ export type UserMeterRpc = { bytes: number updatedAt: string }) => Promise + /** Expand-phase shadow upsert; not used for enforcement. */ + upsertPackageServiceState: (input: { + packageId: string + serviceName: string + status: string + startedAt?: string | null + sourceUpdatedAt: string + updatedAt?: string + }) => Promise + /** Expand-phase shadow delete. */ + deletePackageServiceState: (input: { + packageId: string + serviceName: string + }) => Promise + /** Cutover-support paged shadow list; discovery uses D1. */ + listPackageServiceStates: (input: { + pageSize?: number + startAfter?: string | null + }) => Promise + /** Cutover-support running count; service-start still uses D1. */ + countRunningPackageServices: (input: { + staleAfterMs?: number + excludeService?: { packageId: string; serviceName: string } + now?: string + }) => Promise + /** Cutover-support bulk seed with sourceUpdatedAt monotonic guards. */ + bootstrapPackageServiceStates: (input: { + states: ReadonlyArray<{ + packageId: string + serviceName: string + status: string + startedAt?: string | null + sourceUpdatedAt: string + updatedAt?: string + }> + }) => Promise purge: () => Promise<{ ok: true }> exportCounters: (input: { pageSize?: number diff --git a/packages/worker/src/entitlements/user-meter.workers.test.ts b/packages/worker/src/entitlements/user-meter.workers.test.ts index 5c1440355..a330f75f5 100644 --- a/packages/worker/src/entitlements/user-meter.workers.test.ts +++ b/packages/worker/src/entitlements/user-meter.workers.test.ts @@ -561,6 +561,7 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p expect(await meterA.exportCounters({})).toEqual({ counters: [], storageBytesShadow: null, + packageServiceStatesShadow: [], nextStartAfter: null, truncated: false, }) @@ -690,6 +691,7 @@ test('UserMeter purge blocks concurrent RPCs across deleteAll and schema restore expect(exportDuringPurge).toEqual({ counters: [], storageBytesShadow: null, + packageServiceStatesShadow: [], nextStartAfter: null, truncated: false, }) @@ -963,6 +965,7 @@ test('UserMeter storage RPCs, export shadow field, and purge work additively', a updatedAt: '2026-07-31T17:02:00.000Z', mirrorUpdatedAt: userMeterMirrorUpdatedAtToken(3), }) + expect(firstPage.packageServiceStatesShadow).toEqual([]) const secondPage = await meter.exportCounters({ pageSize: 2, @@ -972,6 +975,7 @@ test('UserMeter storage RPCs, export shadow field, and purge work additively', a expect(secondPage.truncated).toBe(false) expect(secondPage.nextStartAfter).toBeNull() expect(secondPage.storageBytesShadow).toBeNull() + expect(secondPage.packageServiceStatesShadow).toBeNull() await expect( readUserD1StorageBytes({ db: env.APP_DB, userId: user.userId }), @@ -984,6 +988,7 @@ test('UserMeter storage RPCs, export shadow field, and purge work additively', a expect(await meter.exportCounters({})).toEqual({ counters: [], storageBytesShadow: null, + packageServiceStatesShadow: [], nextStartAfter: null, truncated: false, }) @@ -1024,3 +1029,246 @@ test('missing-user storage reserve preserves EntitlementLimitError semantics on }, }) }, 30_000) + +test('UserMeter package-service shadow upserts are monotonic, isolated, exportable, and purgeable', async () => { + const userA = await seedFreeUser('meter-pkg-svc-a') + const userB = await seedFreeUser('meter-pkg-svc-b') + const meterA = userMeterRpc({ env, userId: userA.userId }) + const meterB = userMeterRpc({ env, userId: userB.userId }) + + const created = await meterA.upsertPackageServiceState({ + packageId: 'pkg-1', + serviceName: 'worker', + status: 'running', + startedAt: '2026-08-01T10:00:00.000Z', + sourceUpdatedAt: '2026-08-01T10:00:00.000Z', + }) + expect(created).toMatchObject({ + applied: true, + created: true, + state: { + packageId: 'pkg-1', + serviceName: 'worker', + status: 'running', + startedAt: '2026-08-01T10:00:00.000Z', + sourceUpdatedAt: '2026-08-01T10:00:00.000Z', + revision: 1, + }, + }) + + const heartbeat = await meterA.upsertPackageServiceState({ + packageId: 'pkg-1', + serviceName: 'worker', + status: 'running', + startedAt: '2026-08-01T10:00:00.000Z', + sourceUpdatedAt: '2026-08-01T11:00:00.000Z', + }) + expect(heartbeat).toMatchObject({ + applied: true, + created: false, + state: { + status: 'running', + sourceUpdatedAt: '2026-08-01T11:00:00.000Z', + revision: 2, + }, + }) + + const stale = await meterA.upsertPackageServiceState({ + packageId: 'pkg-1', + serviceName: 'worker', + status: 'stopped', + startedAt: null, + sourceUpdatedAt: '2026-08-01T10:30:00.000Z', + }) + expect(stale).toMatchObject({ + applied: false, + created: false, + state: { + status: 'running', + sourceUpdatedAt: '2026-08-01T11:00:00.000Z', + revision: 2, + }, + }) + + await meterA.upsertPackageServiceState({ + packageId: 'pkg-2', + serviceName: 'idle-svc', + status: 'idle', + sourceUpdatedAt: '2026-08-01T11:05:00.000Z', + }) + await meterB.upsertPackageServiceState({ + packageId: 'pkg-1', + serviceName: 'worker', + status: 'running', + startedAt: '2026-08-01T11:00:00.000Z', + sourceUpdatedAt: '2026-08-01T11:00:00.000Z', + }) + + expect( + await meterA.countRunningPackageServices({ + now: '2026-08-01T11:30:00.000Z', + }), + ).toEqual({ count: 1 }) + expect( + await meterA.countRunningPackageServices({ + now: '2026-08-01T11:30:00.000Z', + excludeService: { packageId: 'pkg-1', serviceName: 'worker' }, + }), + ).toEqual({ count: 0 }) + + const listed = await meterA.listPackageServiceStates({ pageSize: 1 }) + expect(listed.states).toHaveLength(1) + expect(listed.truncated).toBe(true) + expect(listed.nextStartAfter).toEqual(expect.any(String)) + const listedRest = await meterA.listPackageServiceStates({ + pageSize: 10, + startAfter: listed.nextStartAfter, + }) + expect(listedRest.states).toHaveLength(1) + expect(listedRest.truncated).toBe(false) + + const bootstrap = await meterA.bootstrapPackageServiceStates({ + states: [ + { + packageId: 'pkg-1', + serviceName: 'worker', + status: 'error', + sourceUpdatedAt: '2026-08-01T10:00:00.000Z', + }, + { + packageId: 'pkg-3', + serviceName: 'new', + status: 'stopped', + sourceUpdatedAt: '2026-08-01T12:00:00.000Z', + }, + ], + }) + expect(bootstrap).toEqual({ applied: 1, skipped: 1 }) + + const day = '2026-08-01' + for (const resource of [ + 'email_receives_per_day', + 'email_sends_per_day', + ] as const) { + await meterA.initialize({ + resource, + day, + count: 1, + updatedAt: '2026-08-01T12:00:00.000Z', + }) + } + + const firstPage = await meterA.exportCounters({ pageSize: 1 }) + expect(firstPage.truncated).toBe(true) + expect(firstPage.nextStartAfter).toEqual(expect.any(String)) + expect(firstPage.packageServiceStatesShadow).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + packageId: 'pkg-1', + serviceName: 'worker', + status: 'running', + revision: 2, + }), + expect.objectContaining({ + packageId: 'pkg-2', + serviceName: 'idle-svc', + status: 'idle', + }), + expect.objectContaining({ + packageId: 'pkg-3', + serviceName: 'new', + status: 'stopped', + }), + ]), + ) + const secondPage = await meterA.exportCounters({ + pageSize: 1, + startAfter: firstPage.nextStartAfter, + }) + expect(secondPage.packageServiceStatesShadow).toBeNull() + + await expect( + meterA.deletePackageServiceState({ + packageId: 'pkg-1', + serviceName: 'worker', + }), + ).resolves.toEqual({ deleted: true }) + expect( + await meterA.countRunningPackageServices({ + now: '2026-08-01T11:30:00.000Z', + }), + ).toEqual({ count: 0 }) + + await expect(meterA.purge()).resolves.toEqual({ ok: true }) + expect(await meterA.listPackageServiceStates({})).toEqual({ + states: [], + nextStartAfter: null, + truncated: false, + }) + expect(await meterB.listPackageServiceStates({})).toMatchObject({ + states: [ + expect.objectContaining({ + packageId: 'pkg-1', + serviceName: 'worker', + status: 'running', + }), + ], + }) +}, 30_000) + +test('UserMeter package-service sourceUpdatedAt accepts only canonical UTC ISO timestamps', async () => { + const user = await seedFreeUser('meter-pkg-svc-iso') + const meter = userMeterRpc({ env, userId: user.userId }) + const stub = env.USER_METER.get( + env.USER_METER.idFromName(userMeterDurableObjectName(user.userId)), + ) + const valid = '2026-08-01T10:00:00.000Z' + + await expect( + meter.upsertPackageServiceState({ + packageId: 'pkg-iso', + serviceName: 'worker', + status: 'running', + startedAt: valid, + sourceUpdatedAt: valid, + }), + ).resolves.toMatchObject({ + applied: true, + state: { sourceUpdatedAt: valid }, + }) + + const rejected = [ + '', + 'not-a-timestamp', + '2026-08-01T10:00:00Z', + '2026-08-01 10:00:00.000Z', + '2026-08-01T10:00:00.000+02:00', + '2026-08-01T08:00:00.000+00:00', + '2026-13-01T00:00:00.000Z', + '2026-02-30T00:00:00.000Z', + ] as const + await runInDurableObject(stub, async (instance: UserMeter) => { + for (const sourceUpdatedAt of rejected) { + await expect( + instance.upsertPackageServiceState({ + packageId: 'pkg-iso', + serviceName: 'worker', + status: 'stopped', + startedAt: null, + sourceUpdatedAt, + }), + ).rejects.toThrow(/ISO-8601 UTC timestamp/) + } + }) + + // Prior valid row must remain; rejected writes must not throw RangeError past RPC. + expect(await meter.listPackageServiceStates({})).toMatchObject({ + states: [ + expect.objectContaining({ + packageId: 'pkg-iso', + status: 'running', + sourceUpdatedAt: valid, + }), + ], + }) +}, 30_000) diff --git a/packages/worker/src/package-runtime/package-service.node.test.ts b/packages/worker/src/package-runtime/package-service.node.test.ts index 5ca9d226d..99093306b 100644 --- a/packages/worker/src/package-runtime/package-service.node.test.ts +++ b/packages/worker/src/package-runtime/package-service.node.test.ts @@ -1,4 +1,9 @@ import { expect, test, vi } from 'vitest' +import { + consoleWarn, + silenceExpectedConsoleWarns, +} from '#worker/test-support/console-spies.ts' +import { createInMemoryUserMeterEnv } from '#worker/test-support/user-meter.ts' const mockModule = vi.hoisted(() => ({ getSavedPackageById: vi.fn(), @@ -174,6 +179,37 @@ async function flushWaitUntilTasks(waitUntilTasks: Array>) { await Promise.all(waitUntilTasks) } +/** + * Drain UserMeter shadow (and other promptly settling) waitUntil tasks without + * hanging on long-lived background service runs also scheduled via waitUntil. + */ +async function drainSettledWaitUntilTasks( + waitUntilTasks: Array>, + fromIndex = 0, +) { + const pending = waitUntilTasks.slice(fromIndex) + await Promise.all( + pending.map(async (task) => { + const state = { settled: false } + void Promise.resolve(task).then( + () => { + state.settled = true + }, + () => { + state.settled = true + }, + ) + for (let i = 0; i < 40; i++) { + await Promise.resolve() + if (state.settled) { + await task + return + } + } + }), + ) +} + function resetMocks() { mockModule.getSavedPackageById.mockReset() mockModule.loadPackageSourceBySourceId.mockReset() @@ -799,3 +835,619 @@ test('package service restore projects current state into package_service_states vi.useRealTimers() } }) + +function createRecordingPackageServiceDb() { + const upserts: Array> = [] + const deletes: Array> = [] + const appDb = { + prepare(query: string) { + return { + bind(...params: Array) { + return { + async run() { + if (query.includes('INSERT INTO package_service_states')) { + upserts.push(params) + } + if (query.includes('DELETE FROM package_service_states')) { + deletes.push(params) + } + return { meta: { changes: 1 } } + }, + } + }, + } + }, + } as unknown as D1Database + return { appDb, upserts, deletes } +} + +function createNoopPackageServiceDb() { + return { + prepare() { + return { + bind() { + return { + async run() { + return { meta: { changes: 1 } } + }, + } + }, + } + }, + } as unknown as D1Database +} + +function packageServiceShadowRow( + meter: ReturnType, + userId: string, + packageId: string, + serviceName: string, +) { + return meter.packageServicesByUser + .get(userId) + ?.get(`${packageId}\0${serviceName}`) +} + +async function postPackageService( + instance: InstanceType, + action: 'start' | 'stop' | 'purge', + binding: typeof serviceBinding = serviceBinding, +) { + return instance.fetch( + new Request(`https://package-service.invalid/service/${action}`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ binding }), + }), + ) +} + +test('package service start/heartbeat/stop/purge shadow UserMeter transitions', async () => { + resetMocks() + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-07-05T12:00:00.000Z')) + const { appDb, upserts, deletes } = createRecordingPackageServiceDb() + const meter = createInMemoryUserMeterEnv() + + try { + setupSavedPackage('persistent') + mockModule.runBundledModuleWithRegistry.mockImplementation( + () => new Promise(() => {}), + ) + + const created = await createPackageServiceInstance({ + APP_DB: appDb, + USER_METER: meter.env.USER_METER, + } as Env) + const start = await postPackageService(created.instance, 'start') + expect(start.status).toBe(200) + // D1 projection stays on the awaited lifecycle path. + expect(upserts[0]?.[3]).toBe('running') + await drainSettledWaitUntilTasks(created.waitUntilTasks) + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toMatchObject({ + status: 'running', + startedAt: '2026-07-05T12:00:00.000Z', + sourceUpdatedAt: '2026-07-05T12:00:00.000Z', + revision: 1, + }) + expect(created.getAlarmAt()).toBe(Date.parse('2026-07-05T13:00:00.000Z')) + + vi.setSystemTime(new Date('2026-07-05T13:00:00.000Z')) + const afterStartWaitUntilCount = created.waitUntilTasks.length + await created.instance.alarm() + await drainSettledWaitUntilTasks( + created.waitUntilTasks, + afterStartWaitUntilCount, + ) + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toMatchObject({ + status: 'running', + sourceUpdatedAt: '2026-07-05T13:00:00.000Z', + revision: 2, + }) + expect(created.getAlarmAt()).toBe(Date.parse('2026-07-05T14:00:00.000Z')) + + const afterHeartbeatWaitUntilCount = created.waitUntilTasks.length + await postPackageService(created.instance, 'stop') + await drainSettledWaitUntilTasks( + created.waitUntilTasks, + afterHeartbeatWaitUntilCount, + ) + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toMatchObject({ + status: 'running', + sourceUpdatedAt: '2026-07-05T13:00:00.000Z', + revision: 3, + }) + + const afterStopWaitUntilCount = created.waitUntilTasks.length + await postPackageService(created.instance, 'purge') + await drainSettledWaitUntilTasks( + created.waitUntilTasks, + afterStopWaitUntilCount, + ) + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toBeUndefined() + expect(deletes).toContainEqual([ + 'user-123', + 'package-1', + 'realtime-supervisor', + ]) + expect(upserts.length).toBeGreaterThanOrEqual(3) + } finally { + vi.useRealTimers() + } +}) + +test('package service serializes gated UserMeter shadows so purge stays deleted and running→stopped cannot regress', async () => { + resetMocks() + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-07-05T12:00:00.000Z')) + const { appDb } = createRecordingPackageServiceDb() + const meter = createInMemoryUserMeterEnv() + + type Gate = { + release: () => void + promise: Promise + label: string + } + const gates: Array = [] + const callOrder: Array = [] + const baseStub = meter.env.USER_METER.get( + meter.env.USER_METER.idFromName('user-123'), + ) + + function pushGate(label: string): Gate { + const gate: Gate = { + label, + release: () => {}, + promise: Promise.resolve(), + } + gate.promise = new Promise((resolve) => { + gate.release = resolve + }) + gates.push(gate) + return gate + } + + const gatedMeter = { + USER_METER: { + idFromName: (name: string) => ({ name, toString: () => name }), + get: () => ({ + async upsertPackageServiceState( + input: Parameters[0], + ) { + const gate = pushGate(`upsert:${input.status}`) + await gate.promise + callOrder.push(`upsert:${input.status}`) + return baseStub.upsertPackageServiceState(input) + }, + async deletePackageServiceState( + input: Parameters[0], + ) { + const gate = pushGate('delete') + await gate.promise + callOrder.push('delete') + return baseStub.deletePackageServiceState(input) + }, + }), + }, + } + + async function settleMicrotasks() { + await Promise.resolve() + await Promise.resolve() + await Promise.resolve() + } + + async function waitForGate(label: string) { + for (let i = 0; i < 100; i++) { + if (gates.some((gate) => gate.label === label)) return + await Promise.resolve() + } + throw new Error(`Timed out waiting for gated shadow op ${label}.`) + } + + try { + setupSavedPackage('bounded') + mockModule.runBundledModuleWithRegistry.mockImplementation(async () => { + vi.setSystemTime(new Date('2026-07-05T12:00:00.000Z')) + return { result: { ok: true }, error: null } + }) + + const created = await createPackageServiceInstance({ + APP_DB: appDb, + USER_METER: gatedMeter.USER_METER, + } as Env) + + const start = await postPackageService(created.instance, 'start') + expect(start.status).toBe(200) + // Lifecycle returned while the first shadow hop is still gated. + expect(created.waitUntilTasks.length).toBeGreaterThanOrEqual(2) + await waitForGate('upsert:running') + expect(gates.map((gate) => gate.label)).toEqual(['upsert:running']) + expect(callOrder).toEqual([]) + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toBeUndefined() + + // Finish the bounded background run while the running shadow stays gated so + // the same-ms stopped upsert is chained behind it (not started yet). + const backgroundRunTask = created.waitUntilTasks[1] + expect(backgroundRunTask).toBeDefined() + await backgroundRunTask + await settleMicrotasks() + expect(gates.map((gate) => gate.label)).toEqual(['upsert:running']) + expect(callOrder).toEqual([]) + + // Same-ms running→stopped: stopped cannot begin until running is released, + // so a reordered drain cannot let stopped win first and then regress. + gates[0]?.release() + await waitForGate('upsert:stopped') + expect(callOrder).toEqual(['upsert:running']) + expect(gates.map((gate) => gate.label)).toEqual([ + 'upsert:running', + 'upsert:stopped', + ]) + gates[1]?.release() + await drainSettledWaitUntilTasks(created.waitUntilTasks) + expect(callOrder).toEqual(['upsert:running', 'upsert:stopped']) + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toMatchObject({ + status: 'stopped', + sourceUpdatedAt: '2026-07-05T12:00:00.000Z', + }) + + callOrder.length = 0 + gates.length = 0 + const beforePurgeWaitUntilCount = created.waitUntilTasks.length + const purge = await postPackageService(created.instance, 'purge') + expect(purge.status).toBe(200) + await waitForGate('upsert:stopped') + // Purge schedules stopped upsert then delete; delete cannot start first. + expect(gates.map((gate) => gate.label)).toEqual(['upsert:stopped']) + expect(callOrder).toEqual([]) + + // Prefer settling the trailing waitUntil chain entry first while the head + // upsert is still gated. Serialization keeps delete behind the upsert, so + // the purge final state stays deleted instead of resurrecting. + const trailingPurgeChain = created.waitUntilTasks.at(-1) + expect(trailingPurgeChain).toBeDefined() + let trailingSettled = false + void Promise.resolve(trailingPurgeChain).then(() => { + trailingSettled = true + }) + await settleMicrotasks() + expect(trailingSettled).toBe(false) + expect(callOrder).toEqual([]) + + gates[0]?.release() + await waitForGate('delete') + expect(callOrder).toEqual(['upsert:stopped']) + expect(gates.map((gate) => gate.label)).toEqual([ + 'upsert:stopped', + 'delete', + ]) + expect(trailingSettled).toBe(false) + gates[1]?.release() + await drainSettledWaitUntilTasks( + created.waitUntilTasks, + beforePurgeWaitUntilCount, + ) + expect(trailingSettled).toBe(true) + expect(callOrder).toEqual(['upsert:stopped', 'delete']) + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toBeUndefined() + } finally { + vi.useRealTimers() + } +}) + +test('package service lifecycle does not await a gated UserMeter shadow', async () => { + resetMocks() + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-07-05T12:00:00.000Z')) + const { appDb, upserts } = createRecordingPackageServiceDb() + let releaseShadow!: () => void + const shadowGate = new Promise((resolve) => { + releaseShadow = resolve + }) + let shadowFinished = false + const gatedMeter = { + USER_METER: { + idFromName: (name: string) => ({ name, toString: () => name }), + get: () => ({ + async upsertPackageServiceState() { + await shadowGate + shadowFinished = true + return { + applied: true, + created: true, + state: { + packageId: 'package-1', + serviceName: 'realtime-supervisor', + status: 'running' as const, + startedAt: '2026-07-05T12:00:00.000Z', + sourceUpdatedAt: '2026-07-05T12:00:00.000Z', + revision: 1, + updatedAt: '2026-07-05T12:00:00.000Z', + mirrorUpdatedAt: 'r/00000000000000000001', + }, + } + }, + async deletePackageServiceState() { + return { deleted: true } + }, + }), + }, + } + + try { + setupSavedPackage('persistent') + mockModule.runBundledModuleWithRegistry.mockImplementation( + () => new Promise(() => {}), + ) + + const created = await createPackageServiceInstance({ + APP_DB: appDb, + USER_METER: gatedMeter.USER_METER, + } as Env) + const start = await postPackageService(created.instance, 'start') + expect(start.status).toBe(200) + expect(upserts[0]).toEqual([ + 'user-123', + 'package-1', + 'realtime-supervisor', + 'running', + '2026-07-05T12:00:00.000Z', + '2026-07-05T12:00:00.000Z', + ]) + // Lifecycle returned while the shadow DO hop is still gated. + expect(shadowFinished).toBe(false) + expect(created.waitUntilTasks.length).toBeGreaterThan(0) + + releaseShadow() + await drainSettledWaitUntilTasks(created.waitUntilTasks) + expect(shadowFinished).toBe(true) + } finally { + vi.useRealTimers() + } +}) + +test('package service UserMeter shadow failure cannot break D1 projection or lifecycle', async () => { + resetMocks() + silenceExpectedConsoleWarns(['package-service-user-meter-shadow-failed']) + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-07-05T12:00:00.000Z')) + const { appDb, upserts } = createRecordingPackageServiceDb() + const failingMeter = { + USER_METER: { + idFromName: (name: string) => ({ name, toString: () => name }), + get: () => ({ + async upsertPackageServiceState() { + throw new Error('shadow write failed') + }, + async deletePackageServiceState() { + throw new Error('shadow delete failed') + }, + }), + }, + } + + try { + setupSavedPackage('bounded') + mockModule.runBundledModuleWithRegistry.mockImplementation(async () => { + vi.setSystemTime(new Date('2026-07-05T12:00:05.000Z')) + return { result: { ok: true }, error: null } + }) + + const created = await createPackageServiceInstance({ + APP_DB: appDb, + USER_METER: failingMeter.USER_METER, + } as Env) + const start = await postPackageService(created.instance, 'start') + expect(start.status).toBe(200) + expect(upserts[0]).toEqual([ + 'user-123', + 'package-1', + 'realtime-supervisor', + 'running', + '2026-07-05T12:00:00.000Z', + '2026-07-05T12:00:00.000Z', + ]) + await flushWaitUntilTasks(created.waitUntilTasks) + expect(upserts.at(-1)).toEqual([ + 'user-123', + 'package-1', + 'realtime-supervisor', + 'stopped', + null, + '2026-07-05T12:00:05.000Z', + ]) + expect(consoleWarn).toHaveBeenCalledWith( + 'package-service-user-meter-shadow-failed', + expect.any(Error), + ) + + const beforePurgeWaitUntilCount = created.waitUntilTasks.length + const purge = await postPackageService(created.instance, 'purge') + expect(purge.status).toBe(200) + await drainSettledWaitUntilTasks( + created.waitUntilTasks, + beforePurgeWaitUntilCount, + ) + expect(consoleWarn).toHaveBeenCalledWith( + 'package-service-user-meter-shadow-failed', + expect.any(Error), + ) + } finally { + vi.useRealTimers() + } +}) + +test('package service UserMeter shadows are isolated per owning user', async () => { + resetMocks() + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-07-05T12:00:00.000Z')) + const appDb = createNoopPackageServiceDb() + const meter = createInMemoryUserMeterEnv() + const otherBinding = { + ...serviceBinding, + userId: 'user-other', + packageId: 'package-2', + kodyId: '@scope/package-2', + } + + try { + setupSavedPackage('bounded') + mockModule.runBundledModuleWithRegistry.mockImplementation( + () => new Promise(() => {}), + ) + mockModule.getSavedPackageById.mockImplementation( + async (_db: unknown, input: { packageId: string }) => + input.packageId === 'package-2' + ? { + id: 'package-2', + sourceId: 'source-1', + kodyId: '@scope/package-2', + } + : savedPackage, + ) + mockModule.loadPackageSourceBySourceId.mockImplementation(async () => + createPackageSource('bounded'), + ) + + const first = await createPackageServiceInstance({ + APP_DB: appDb, + USER_METER: meter.env.USER_METER, + } as Env) + const second = await createPackageServiceInstance({ + APP_DB: appDb, + USER_METER: meter.env.USER_METER, + } as Env) + + await postPackageService(first.instance, 'start') + await postPackageService(second.instance, 'start', otherBinding) + await drainSettledWaitUntilTasks(first.waitUntilTasks) + await drainSettledWaitUntilTasks(second.waitUntilTasks) + + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toMatchObject({ status: 'running' }) + expect( + packageServiceShadowRow( + meter, + 'user-other', + 'package-2', + 'realtime-supervisor', + ), + ).toMatchObject({ status: 'running' }) + expect(meter.packageServicesByUser.get('user-123')?.size).toBe(1) + expect(meter.packageServicesByUser.get('user-other')?.size).toBe(1) + expect( + meter.packageServicesByUser + .get('user-123') + ?.has('package-2\0realtime-supervisor'), + ).toBe(false) + } finally { + vi.useRealTimers() + } +}) + +test('package service restore shadows current D1 projection into UserMeter', async () => { + resetMocks() + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-07-05T14:00:00.000Z')) + const appDb = createNoopPackageServiceDb() + const meter = createInMemoryUserMeterEnv() + + try { + const restored = createPackageServiceState() + restored.persistedEntries.set('package-service-state', { + binding: serviceBinding, + autoStart: false, + mode: 'bounded', + timeoutMs: 30_000, + stopRequested: false, + currentRunId: null, + nextAlarmAt: null, + nextAlarmSource: null, + lastStartedAt: '2026-07-05T13:00:00.000Z', + lastStoppedAt: '2026-07-05T13:05:00.000Z', + status: 'error', + lastError: 'boom', + lastResult: null, + lastRunFinishedAt: '2026-07-05T13:05:00.000Z', + consecutiveFailureCount: 1, + }) + new PackageServiceInstance(restored.state, { + APP_DB: appDb, + USER_METER: meter.env.USER_METER, + } as Env) + await waitForRestoreState(restored.state) + await drainSettledWaitUntilTasks(restored.waitUntilTasks) + + expect( + packageServiceShadowRow( + meter, + 'user-123', + 'package-1', + 'realtime-supervisor', + ), + ).toMatchObject({ + status: 'error', + startedAt: null, + sourceUpdatedAt: '2026-07-05T14:00:00.000Z', + revision: 1, + }) + } finally { + vi.useRealTimers() + } +}) diff --git a/packages/worker/src/package-runtime/package-service.ts b/packages/worker/src/package-runtime/package-service.ts index 56e8bc9f6..af9846861 100644 --- a/packages/worker/src/package-runtime/package-service.ts +++ b/packages/worker/src/package-runtime/package-service.ts @@ -3,6 +3,10 @@ import * as Sentry from '@sentry/cloudflare' import { DurableObject } from 'cloudflare:workers' import { z } from 'zod' import { createMcpCallerContext } from '#mcp/context.ts' +import { + userMeterNamespace, + userMeterRpc, +} from '#worker/entitlements/user-meter-client.ts' import { getPackageServiceEntryPath, listPackageServices, @@ -340,6 +344,8 @@ class PackageServiceInstanceBase extends DurableObject { private stateSnapshot: PackageServiceState = createInitialPackageServiceState() private activeRunPromise: Promise | null = null + /** Per-instance shadow write chain so waitUntil hops cannot reorder. */ + private shadowQueue: Promise = Promise.resolve() constructor(state: DurableObjectState, env: Env) { super(state, env) @@ -415,26 +421,119 @@ class PackageServiceInstanceBase extends DurableObject { } /** - * Best-effort D1 projection of service liveness for entitlement counting. - * Failures must never break the service path; stop/error/idle writes still - * attempt to clear `running` so quota is released when D1 is healthy. + * Best-effort D1 projection of service liveness for entitlement counting, + * plus a non-awaited expand-phase UserMeter shadow via `waitUntil`. D1 + * remains sole count/discovery authority and stays on the awaited path; + * shadow DO hops must not gate lifecycle/heartbeat responses. + * Stop/error/idle writes still attempt to clear `running` so quota is + * released when D1 is healthy. */ private async projectServiceStateToD1() { const binding = this.stateSnapshot.binding if (!binding) return + const status = projectPackageServiceStatus(this.stateSnapshot.status) + const startedAt = this.stateSnapshot.lastStartedAt + const updatedAt = new Date().toISOString() try { await upsertPackageServiceState({ db: this.env.APP_DB, userId: binding.userId, packageId: binding.packageId, serviceName: binding.serviceName, - status: projectPackageServiceStatus(this.stateSnapshot.status), - startedAt: this.stateSnapshot.lastStartedAt, - updatedAt: new Date().toISOString(), + status, + startedAt, + updatedAt, }) } catch { // Best-effort: D1 outages must not take down package services. } + this.schedulePackageServiceStateShadow({ + binding, + status, + startedAt, + sourceUpdatedAt: updatedAt, + }) + } + + /** + * Schedule a best-effort UserMeter shadow upsert. Never awaited by + * lifecycle/heartbeat paths; the helper catches so `waitUntil` cannot + * reject. Appended to the per-instance shadow queue so rapid transitions + * and purge (stopped upsert then delete) cannot reorder. + */ + private schedulePackageServiceStateShadow(input: { + binding: PackageServiceBindingState + status: PackageServiceProjectedStatus + startedAt: string | null + sourceUpdatedAt: string + }) { + if (!userMeterNamespace(this.env)) return + this.enqueueShadowTask(() => + this.shadowPackageServiceStateToUserMeter(input), + ) + } + + /** Best-effort UserMeter shadow upsert; catches/logs so waitUntil settles. */ + private async shadowPackageServiceStateToUserMeter(input: { + binding: PackageServiceBindingState + status: PackageServiceProjectedStatus + startedAt: string | null + sourceUpdatedAt: string + }) { + try { + const meter = userMeterRpc({ + env: this.env, + userId: input.binding.userId, + }) + await meter.upsertPackageServiceState({ + packageId: input.binding.packageId, + serviceName: input.binding.serviceName, + status: input.status, + startedAt: input.startedAt, + sourceUpdatedAt: input.sourceUpdatedAt, + }) + } catch (error) { + console.warn('package-service-user-meter-shadow-failed', error) + } + } + + /** + * Schedule a best-effort UserMeter shadow delete. Never awaited by purge; + * the helper catches so `waitUntil` cannot reject. Queued after any prior + * shadow upserts for this instance. + */ + private schedulePackageServiceShadowDelete( + binding: PackageServiceBindingState, + ) { + if (!userMeterNamespace(this.env)) return + this.enqueueShadowTask(() => this.deletePackageServiceShadow(binding)) + } + + /** + * Run shadow writes in scheduling order and register the chain with + * `waitUntil`. Helpers already catch, so the queue itself does not reject. + */ + private enqueueShadowTask(task: () => Promise) { + this.shadowQueue = this.shadowQueue.then(task, task) + this.ctx.waitUntil(this.shadowQueue) + } + + /** Best-effort UserMeter shadow delete; catches/logs so waitUntil settles. */ + private async deletePackageServiceShadow( + binding: PackageServiceBindingState, + ) { + try { + const meter = userMeterRpc({ + env: this.env, + userId: binding.userId, + }) + await meter.deletePackageServiceState({ + packageId: binding.packageId, + serviceName: binding.serviceName, + }) + } catch (error) { + console.warn('package-service-user-meter-shadow-failed', error) + } } private async deleteProjectedServiceState( @@ -450,6 +549,7 @@ class PackageServiceInstanceBase extends DurableObject { } catch { // Best-effort cleanup on purge. } + this.schedulePackageServiceShadowDelete(binding) } private async ensureRunningHeartbeat() { @@ -956,8 +1056,9 @@ class PackageServiceInstanceBase extends DurableObject { await this.ctx.storage.deleteAlarm().catch(() => { // Best effort cleanup before deleteAll. }) - await this.ctx.storage.deleteAll() + // Clear D1 + UserMeter shadow before deleteAll. await this.deleteProjectedServiceState(binding) + await this.ctx.storage.deleteAll() return Response.json({ ok: true, }) diff --git a/packages/worker/src/test-support/user-meter.ts b/packages/worker/src/test-support/user-meter.ts index 8452d2d28..2eac2d4ec 100644 --- a/packages/worker/src/test-support/user-meter.ts +++ b/packages/worker/src/test-support/user-meter.ts @@ -1,13 +1,24 @@ import { isDailyEntitlementResource, + isUserMeterPackageServiceStatus, userMeterMirrorUpdatedAtToken, + userMeterPackageServiceStateStaleMs, type DailyEntitlementResource, + type UserMeterPackageServiceState, + type UserMeterPackageServiceStatus, type UserMeterStorageBytesState, } from '#worker/entitlements/user-meter-do.ts' import { type UserMeterEnv } from '#worker/entitlements/user-meter-client.ts' type MeterRow = { count: number; revision: number } type StorageRow = { bytes: number; revision: number; updatedAt: string } +type PackageServiceRow = { + status: UserMeterPackageServiceStatus + startedAt: string | null + sourceUpdatedAt: string + revision: number + updatedAt: string +} /** * In-memory UserMeter stub keyed by `idFromName` (stable userId) for node-unit @@ -17,16 +28,31 @@ type StorageRow = { bytes: number; revision: number; updatedAt: string } export function createInMemoryUserMeterEnv() { const metersByUser = new Map>() const storageByUser = new Map() + const packageServicesByUser = new Map< + string, + Map + >() function counterKey(resource: string, day: string) { return `${resource}\0${day}` } + function packageServiceKey(packageId: string, serviceName: string) { + return `${packageId}\0${serviceName}` + } + function meterFor(userId: string) { const existingRows = metersByUser.get(userId) const rows = existingRows ?? new Map() if (!existingRows) metersByUser.set(userId, rows) + const existingPackageServices = packageServicesByUser.get(userId) + const packageServices = + existingPackageServices ?? new Map() + if (!existingPackageServices) { + packageServicesByUser.set(userId, packageServices) + } + function readRow(resource: string, day: string) { return rows.get(counterKey(resource, day)) ?? null } @@ -49,6 +75,36 @@ export function createInMemoryUserMeterEnv() { } } + function packageServiceState( + packageId: string, + serviceName: string, + row: PackageServiceRow, + ): UserMeterPackageServiceState { + return { + packageId, + serviceName, + status: row.status, + startedAt: row.status === 'running' ? row.startedAt : null, + sourceUpdatedAt: row.sourceUpdatedAt, + revision: row.revision, + updatedAt: row.updatedAt, + mirrorUpdatedAt: userMeterMirrorUpdatedAtToken(row.revision), + } + } + + function listPackageServiceStatesSorted() { + return [...packageServices.entries()] + .map(([entryKey, row]) => { + const [packageId, serviceName] = entryKey.split('\0') + return packageServiceState(packageId!, serviceName!, row) + }) + .sort((left, right) => { + const byPackage = left.packageId.localeCompare(right.packageId) + if (byPackage !== 0) return byPackage + return left.serviceName.localeCompare(right.serviceName) + }) + } + return { async initialize(input: { resource: string @@ -177,9 +233,171 @@ export function createInMemoryUserMeterEnv() { storageByUser.set(userId, next) return { ...storageReady(next), created: false } }, + async upsertPackageServiceState(input: { + packageId: string + serviceName: string + status: string + startedAt?: string | null + sourceUpdatedAt: string + updatedAt?: string + }) { + if (!isUserMeterPackageServiceStatus(input.status)) { + throw new Error( + `UserMeter package service status must be running, idle, stopped, or error; got ${JSON.stringify(input.status)}.`, + ) + } + const status = input.status + const key = packageServiceKey(input.packageId, input.serviceName) + const existing = packageServices.get(key) + if (existing && existing.sourceUpdatedAt > input.sourceUpdatedAt) { + return { + applied: false, + created: false, + state: packageServiceState( + input.packageId, + input.serviceName, + existing, + ), + } + } + const startedAt = + status === 'running' && + typeof input.startedAt === 'string' && + input.startedAt.length > 0 + ? input.startedAt + : null + const updatedAt = + typeof input.updatedAt === 'string' && input.updatedAt.length > 0 + ? input.updatedAt + : input.sourceUpdatedAt + const row: PackageServiceRow = { + status, + startedAt, + sourceUpdatedAt: input.sourceUpdatedAt, + revision: existing ? existing.revision + 1 : 1, + updatedAt, + } + packageServices.set(key, row) + return { + applied: true, + created: !existing, + state: packageServiceState(input.packageId, input.serviceName, row), + } + }, + async deletePackageServiceState(input: { + packageId: string + serviceName: string + }) { + const deleted = packageServices.delete( + packageServiceKey(input.packageId, input.serviceName), + ) + return { deleted } + }, + async listPackageServiceStates( + input: { + pageSize?: number + startAfter?: string | null + } = {}, + ) { + const pageSize = + typeof input.pageSize === 'number' && Number.isFinite(input.pageSize) + ? Math.min(Math.max(Math.trunc(input.pageSize), 1), 500) + : 100 + const all = listPackageServiceStatesSorted() + let startIndex = 0 + if ( + typeof input.startAfter === 'string' && + input.startAfter.length > 0 + ) { + try { + const parsed = JSON.parse(input.startAfter) as unknown + if ( + Array.isArray(parsed) && + parsed.length === 2 && + typeof parsed[0] === 'string' && + typeof parsed[1] === 'string' + ) { + const packageId = parsed[0] + const serviceName = parsed[1] + startIndex = + all.findIndex( + (row) => + row.packageId === packageId && + row.serviceName === serviceName, + ) + 1 + } + } catch { + startIndex = 0 + } + } + const page = all.slice(startIndex, startIndex + pageSize) + const truncated = startIndex + pageSize < all.length + const last = page[page.length - 1] + return { + states: page, + nextStartAfter: + truncated && last + ? JSON.stringify([last.packageId, last.serviceName]) + : null, + truncated, + } + }, + async countRunningPackageServices( + input: { + staleAfterMs?: number + excludeService?: { packageId: string; serviceName: string } + now?: string + } = {}, + ) { + const now = input.now ? new Date(input.now) : new Date() + const staleAfterMs = + typeof input.staleAfterMs === 'number' && + Number.isFinite(input.staleAfterMs) && + input.staleAfterMs >= 0 + ? Math.trunc(input.staleAfterMs) + : userMeterPackageServiceStateStaleMs + const freshAfter = new Date(now.valueOf() - staleAfterMs).toISOString() + let count = 0 + for (const [entryKey, row] of packageServices) { + if (row.status !== 'running') continue + if (row.sourceUpdatedAt < freshAfter) continue + if (input.excludeService) { + const [packageId, serviceName] = entryKey.split('\0') + if ( + packageId === input.excludeService.packageId && + serviceName === input.excludeService.serviceName + ) { + continue + } + } + count += 1 + } + return { count } + }, + async bootstrapPackageServiceStates(input: { + states: ReadonlyArray<{ + packageId: string + serviceName: string + status: string + startedAt?: string | null + sourceUpdatedAt: string + updatedAt?: string + }> + }) { + let applied = 0 + let skipped = 0 + for (const state of input.states) { + const result = await this.upsertPackageServiceState(state) + if (result.applied) applied += 1 + else skipped += 1 + } + return { applied, skipped } + }, async purge() { rows.clear() storageByUser.delete(userId) + packageServices.clear() + packageServicesByUser.delete(userId) return { ok: true as const } }, async exportCounters( @@ -199,11 +417,9 @@ export function createInMemoryUserMeterEnv() { mirrorUpdatedAt: userMeterMirrorUpdatedAtToken(row.revision), } }) - const includeStorageShadow = + const includeShadows = typeof input.startAfter !== 'string' || input.startAfter.length === 0 - const storage = includeStorageShadow - ? storageByUser.get(userId) - : undefined + const storage = includeShadows ? storageByUser.get(userId) : undefined const storageBytesShadow: UserMeterStorageBytesState | null = storage ? { bytes: storage.bytes, @@ -212,9 +428,13 @@ export function createInMemoryUserMeterEnv() { mirrorUpdatedAt: userMeterMirrorUpdatedAtToken(storage.revision), } : null + const packageServiceStatesShadow = includeShadows + ? listPackageServiceStatesSorted() + : null return { counters, storageBytesShadow, + packageServiceStatesShadow, nextStartAfter: null, truncated: false, } @@ -233,6 +453,7 @@ export function createInMemoryUserMeterEnv() { env, metersByUser, storageByUser, + packageServicesByUser, async seed(input: { userId: string resource: DailyEntitlementResource