diff --git a/docs/contributing/architecture/data-storage.md b/docs/contributing/architecture/data-storage.md index ab081153f0..9792bb8e48 100644 --- a/docs/contributing/architecture/data-storage.md +++ b/docs/contributing/architecture/data-storage.md @@ -388,8 +388,10 @@ layout is `index1 = userId`, `blob1 = event type`, `blob2 = delivery outcome`, `blob3 = source timestamp`, and `double1 = 1`. Admin queries return only platform-wide day/outcome counts and weight sampled rows by `_sample_interval`. When Analytics Engine SQL is unreachable, these two charts zero-fill while the -rest of the page renders. Local development uses the existing D1 counters and -delivery-event table because Wrangler's emulated dataset has no SQL API. +rest of the page renders. Local development cannot query Wrangler's emulated +Analytics Engine SQL API: email quota aggregates degrade to empty (with an +explicit warning) rather than reading the retired D1 mirror, while +delivery-outcome aggregates still read D1 `email_delivery_events`. **Mailbox expand-phase parity events** reuse the same `EMAIL_EVENTS` dataset with a separate row shape defined in @@ -637,21 +639,17 @@ time-pruned. Deletion-fence legacy lease rows are bounded by the D1 snapshot replace on `markDeleting` rather than time retention; DO-authority rows clear on release/repair/purge. -**Expand-phase D1 mirrors (daily counters only):** enforcement and point reads -are authoritative in UserMeter for daily counters. D1 -`entitlement_daily_counters` is **not** dropped — it remains a best-effort -mirror for existing readers and reporting. After each DO consume/refund/inbound -claim, the entitlements service schedules a non-awaited absolute mirror write -keyed by `(user_id, resource, day)` with a revision-ordered `updated_at` token -(`r/` + zero-padded revision from `userMeterMirrorUpdatedAtToken`) so late -writes cannot overwrite newer state. See +**D1 daily mirror writes stopped:** enforcement, point reads, bootstrap, and +mirror paths no longer read or write `entitlement_daily_counters`. The admin +parity report (`admin_user_meter_parity`) temporarily reads the table while it +exists for migration verification. The physical table stays quiescent until a +follow-up migration-only deploy drops it. See [Entitlements](./entitlements.md#usermeter-expand-phase). **Daily cold bootstrap:** a missing `(resource, day)` row returns -`needs_bootstrap`. The service performs one legacy D1 point read on -`entitlement_daily_counters`, then `initialize()` seeds the DO row with -`INSERT OR IGNORE` (concurrent callers cannot double-apply the baseline). Warm -daily paths never read D1 for enforcement. +`needs_bootstrap`. The service calls `initialize({ count: 0 })` with +`INSERT OR IGNORE` (concurrent callers stay safe). Warm daily paths never read +D1 for enforcement. Account deletion calls `UserMeter.purge()` (one RPC per user, no D1 id scan; `deleteAll` clears counters, claims, storage-byte shadow, package-service @@ -1664,10 +1662,13 @@ Current retention policies: After Mailbox cut-over, the same 365-day window and blob-before-row ordering are self-enforced by the Mailbox DO alarm; `system:email` stays on the D1 system-email retention job. -- `entitlement_daily_counters`: expand-phase **mirror** of UserMeter daily - counters (authoritative state lives in the per-user `UserMeter` DO). Rows keep - 400 days by `day` key until mirror retirement is verified after - reporting-off-D1 merges; the table is not dropped in this phase. +- `entitlement_daily_counters`: **quiescent pending drop** — mirror writes, + bootstrap reads, and scheduled retention pruning stop in the code deploy. The + admin parity report (`admin_user_meter_parity`) temporarily reads the table + while it exists; account deletion keeps removing rows while the physical table + exists; a follow-up code deploy removes that inventory target before the later + drop migration. Daily counter retention lives in the per-user `UserMeter` DO + (`userMeterDailyCounterRetentionDays`). - `usage_rollups`: per user/metric/month rollups keep 24 months by `month` key; raw Analytics Engine usage events follow platform retention. - `feature_flag_exposure_rollups`: local-dev/test flag exposure rollups keep 90 diff --git a/docs/contributing/architecture/entitlements.md b/docs/contributing/architecture/entitlements.md index 30117e2fc6..80dda96eeb 100644 --- a/docs/contributing/architecture/entitlements.md +++ b/docs/contributing/architecture/entitlements.md @@ -147,25 +147,22 @@ reserve bytes in D1 or UserMeter. `consumeInboundDelivery` RPCs check the plan limit and increment inside the DO. The Durable Object request model serializes mutations per user; counter updates use optimistic concurrency on monotonic `revision` so concurrent consumes cannot -overshoot. Missing `(resource, day)` rows return `needs_bootstrap` rather than -silently starting at zero on a warm account. - -**Cold bootstrap:** on `needs_bootstrap`, the service performs one legacy D1 -point read on `entitlement_daily_counters`, then `UserMeter.initialize()` seeds -the row with `INSERT OR IGNORE` (concurrent callers cannot double-apply the -baseline). The warm enforcement path awaits only the DO RPC — never APP_DB. - -**Non-awaited D1 mirror:** after each successful consume, refund, or inbound -delivery claim, the service schedules a best-effort absolute mirror write to -`entitlement_daily_counters` via `waitUntil` when available (otherwise a caught -void promise). Mirror ordering uses the DO-minted `mirrorUpdatedAt` token -(`r/` + zero-padded revision) in the existing `updated_at` TEXT column so late -writes cannot overwrite newer state, including refunds that lower `count`. -Mirror failures are logged and never affect enforcement. The D1 table is **not** -dropped in this phase — it remains for existing readers and reporting. +overshoot. Missing `(resource, day)` rows return `needs_bootstrap`; the service +then initializes that key at zero via `UserMeter.initialize()` +(`INSERT OR IGNORE`, concurrent-safe) before retrying. Warm enforcement awaits +only the DO RPC and never touches D1 daily counter state. + +**D1 daily mirror writes stopped:** consume, refund, inbound charge/read, +point-read surfaces, and retention no longer read or write +`entitlement_daily_counters`. Generic account export and deletion keep reading +or removing user rows while the physical table exists. A follow-up code deploy +removes those inventory targets before the later migration-only drop (migrations +apply before Workers); existing rows are otherwise quiescent historical mirror +state. Analytics Engine remains the production reporting path for email +send/receive aggregates. **Point-read surfaces** call `readDailyEntitlementResourceUsage` (UserMeter with -the same cold-bootstrap path): +the same cold zero-init path): - Account usage UI — `packages/worker/src/app/account-usage-data.ts` - Account email usage panel — `packages/worker/src/app/account-email-data.ts` @@ -173,23 +170,15 @@ the same cold-bootstrap path): - Admin per-user usage drill-down — `packages/worker/src/admin/user-usage-data.ts` -Non-daily resources and contexts without `USER_METER` still use -`readEntitlementResourceUsage` against D1. - -During a rolling deployment, requests already running on the previous Worker -version may still increment D1 after a new-version request bootstraps its DO -row. Cloudflare activation bounds that overlap to in-flight requests, but -operators should treat mirror parity during the deploy window as approximate; -post-deploy requests have one authority in UserMeter. +Non-daily resources still use `readEntitlementResourceUsage` against D1. +`readEntitlementResourceUsage` for daily resources throws and directs callers to +the UserMeter helpers above. **Inbound retry idempotency:** inbound receive quota uses `UserMeter.consumeInboundDelivery`, which atomically claims `delivery_id` and consumes one `email_receives_per_day` unit inside a SQLite transaction. Retries return the accepted counter without incrementing (`replayed: true`). -Cross-UTC-day retries use the original claim's resource/day. The legacy D1 -mirror is scheduled from the email path via -`scheduleAbsoluteDailyEntitlementMirror` so the email subsystem does not -duplicate mirror SQL. +Cross-UTC-day retries use the original claim's resource/day. ### D1 payload storage bytes — UserMeter shadow (expand phase slice 3) @@ -371,29 +360,38 @@ same-token rollout mirrors; email keeps its D1 lease path. Package-service and storage authority flips remain 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 -production. +**Daily-counter mirror retirement (three-deploy):** this code deploy stops all +D1 mirror/bootstrap/retention use while leaving the physical +`entitlement_daily_counters` table and account-deletion target in place. A +follow-up code deploy removes that target after this Worker is healthy; the +third, migration-only deploy drops the table. Production +`admin_user_meter_parity` scans across 38 users showed zero daily mismatches +with Analytics Engine reporting active — the deploy rationale for stopping +mirror writes before the drop. While the table exists, parity still compares +`d1Count === meterCount` (`mirrorRetired: false`); after the drop migration, +`mirrorRetired: true` reports meter counts only. ### Admin UserMeter parity gates (`admin_user_meter_parity`) -Production verification for mirror retirement and authority flips uses the -admin-only read-only capability `admin_user_meter_parity` (input: +Production verification for mirror retirement and remaining authority flips uses +the admin-only read-only capability `admin_user_meter_parity` (input: `stable_user_id`). It compares production-shaped D1 rows for one account against -direct UserMeter RPCs and never bootstraps or writes parity state. Opening a -cold UserMeter stub may still run Durable Object constructor schema maintenance -and opportunistic stale daily-counter pruning. Cold meter rows surface as +direct UserMeter RPCs and never bootstraps or writes parity state. Daily +comparison retires automatically once the drop migration removes +`entitlement_daily_counters` (`daily.mirrorRetired: true`). Opening a cold +UserMeter stub may still run Durable Object constructor schema maintenance and +opportunistic stale daily-counter pruning. Cold meter rows surface as `needsBootstrap` with `meterCount`/`meterBytes` null. Interpret the structured report as independent gates: -| Gate | Pass condition | -| -------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| Daily counters (current UTC day) | Each of the four daily resources has `parity: true` (`d1Count === meterCount`); aggregate `daily.mismatchCount === 0`. | -| Storage bytes | `storage.parity` — D1 `users.d1_storage_bytes` equals UserMeter `readStorageBytes` (not `needsBootstrap`). | -| Package services | `packageServices.parity` — inventory mismatch category counts are all zero (`d1Only` / `meterOnly` / `statusMismatch` / `startedAtMismatch` / `sourceUpdatedAtMismatch`), fresh-running counts match under the shared 24h stale window, and the meter page walk is not `truncated`. | -| Deletion tombstone | `deletion.deletingAtParity` — D1 `users.deleting_at` matches the meter tombstone. | -| Temporary D1 lease mirror | `deletion.mirrorLeaseParity` — `doOnly === 0`, `legacyWithoutD1 === 0`, inventory not truncated, and `d1ActiveLeaseCount >= doAuthorityLeaseCount` (same-token mirror coverage). `tokenSetMismatches.d1Only` is reported but does **not** fail this gate. | +| Gate | Pass condition | +| -------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| Daily counters (current UTC day) | While `daily.mirrorRetired` is false (table present): each of the four daily resources has `parity: true` (`d1Count === meterCount`); aggregate `daily.mismatchCount === 0`. After the drop migration (`mirrorRetired: true`): report meter counts only; `d1Count`/`delta` are null, each resource has `parity: true`, and `mismatchCount === 0` (no D1 comparison). | +| Storage bytes | `storage.parity` — D1 `users.d1_storage_bytes` equals UserMeter `readStorageBytes` (not `needsBootstrap`). | +| Package services | `packageServices.parity` — inventory mismatch category counts are all zero (`d1Only` / `meterOnly` / `statusMismatch` / `startedAtMismatch` / `sourceUpdatedAtMismatch`), fresh-running counts match under the shared 24h stale window, and the meter page walk is not `truncated`. | +| Deletion tombstone | `deletion.deletingAtParity` — D1 `users.deleting_at` matches the meter tombstone. | +| Temporary D1 lease mirror | `deletion.mirrorLeaseParity` — `doOnly === 0`, `legacyWithoutD1 === 0`, inventory not truncated, and `d1ActiveLeaseCount >= doAuthorityLeaseCount` (same-token mirror coverage). `tokenSetMismatches.d1Only` is reported but does **not** fail this gate. | **D1-only leases:** email and other transition paths that omit `env` still take exact D1 leases, so `d1Only > 0` is expected until that handoff. Mirror-removal @@ -402,12 +400,12 @@ are known email/transition holders via `admin_account_write_lease_list` and holder classification before retiring the temporary D1 mirror inventory. **Threshold:** treat unexplained mismatches as blocking for the corresponding -cutover (daily mirror retirement, storage authority flip, package-service -authority flip, or temporary D1 lease-mirror removal). Expected cold accounts -may report `needsBootstrap` until live traffic or an intentional bootstrap path -seeds the DO; that is a bootstrap gap, not a silent pass. Truncated inventories -fail closed (`parity` / `mirrorLeaseParity` false) so operators re-run or raise -the bounded page cap rather than approve a partial compare. +cutover (daily mirror retirement before the drop migration, storage authority +flip, package-service authority flip, or temporary D1 lease-mirror removal). +Expected cold accounts may report `needsBootstrap` until live traffic seeds the +DO; that is a bootstrap gap, not a silent pass. Truncated inventories fail +closed (`parity` / `mirrorLeaseParity` false) so operators re-run or raise the +bounded page cap rather than approve a partial compare. Module wiring: `consumeDailyEntitlement`, `refundDailyEntitlement`, and `readDailyEntitlementResourceUsage` require `env.USER_METER` and fail closed @@ -568,14 +566,13 @@ Rules: the limit abuse-resistant for permanent rejects (parse failures, entitlement/quota rejects). - **Cold bootstrap:** missing `(resource, day)` rows trigger one legacy D1 point - read and a single `UserMeter.initialize()` before retrying the consume. + **Cold bootstrap:** missing `(resource, day)` rows trigger + `UserMeter.initialize({ count: 0 })` (`INSERT OR IGNORE`) before retrying the + consume. Concurrent cold callers cannot double-apply a non-zero baseline. - **D1 mirror (expand phase):** after each successful consume/refund/inbound - claim, a best-effort, non-awaited absolute mirror write updates - `entitlement_daily_counters` with revision-ordered `updated_at` tokens. The - table is not dropped — it remains for existing readers and reporting until - reporting-off-D1 retirement is verified. + **D1 mirror writes stopped:** consume/refund/inbound charge/read paths never + touch `entitlement_daily_counters`; the physical table remains quiescent until + a follow-up migration drops it. A delivery claim remains charged when later storage fails. Cloudflare Email Routing retries replay that same `delivery_id` through @@ -583,9 +580,6 @@ Rules: across a UTC-day boundary. The retained claim is the idempotency boundary; production inbound handling does not call `refundDailyEntitlement`. - `incrementDailyEntitlementCounter` remains for raw D1 counter writes (tests, - backfills, and legacy paths). - - **Boolean allowances** (persistent package services) are modeled as limit `0` (not allowed) vs `1` (allowed) so the numeric contract stays uniform. - **Per-unit size limits** (`email_message_bytes`) compare one candidate value @@ -795,9 +789,10 @@ in [`../environment-variables.md`](../environment-variables.md). `0066-stripe-billing.sql`; owned by `packages/worker/src/billing/`, read by `getUserPlan` via `resolveEffectivePlan`. `stripe_plan` stays nullable because it is Stripe-derived; `max` is manual-only. -- `entitlement_daily_counters` — expand-phase **mirror** of UserMeter daily - counters (authoritative state in the per-user `UserMeter` DO), created by - migration `0048-user-plans-and-entitlement-counters.sql`; included in the - account-deletion cascade (`packages/worker/src/app/account-deletion.ts`). - Table retirement waits until reporting-off-D1 merges and mirror parity is - verified. +- `entitlement_daily_counters` — **quiescent pending drop**. Created by + migration `0048-user-plans-and-entitlement-counters.sql`; mirror writes, + bootstrap reads, and retention pruning stop in the code deploy. Account + deletion keeps removing user rows while the table exists. Daily counters are + authoritative in the per-user `UserMeter` DO; account export uses UserMeter + RPCs. A follow-up code deploy removes the deletion target before a later + migration-only PR drops the table. diff --git a/packages/worker/src/admin/user-meter-parity.node.test.ts b/packages/worker/src/admin/user-meter-parity.node.test.ts index 49b6b8a0e5..47a1e7c792 100644 --- a/packages/worker/src/admin/user-meter-parity.node.test.ts +++ b/packages/worker/src/admin/user-meter-parity.node.test.ts @@ -330,6 +330,7 @@ test('loadAdminUserMeterParityReport reports full parity across daily/storage/se stableUserId, daily: { day, + mirrorRetired: false, mismatchCount: 0, }, storage: { @@ -764,3 +765,155 @@ test('loadAdminUserMeterParityReport fails closed on zero-item write-lease curso mirrorLeaseParity: false, }) }) + +test('loadAdminUserMeterParityReport reports mirrorRetired when entitlement_daily_counters is absent', async () => { + const sqlite = new DatabaseSync(':memory:') + sqlite.exec(` + CREATE TABLE users ( + id INTEGER PRIMARY KEY, + stable_user_id TEXT UNIQUE NOT NULL, + username TEXT NOT NULL, + email TEXT NOT NULL, + d1_storage_bytes INTEGER NOT NULL DEFAULT 0, + deleting_at TEXT, + created_at TEXT NOT NULL DEFAULT '2026-01-01T00:00:00.000Z', + updated_at TEXT NOT NULL DEFAULT '2026-01-01T00:00:00.000Z' + ); + CREATE TABLE package_service_states ( + user_id TEXT NOT NULL, + package_id TEXT NOT NULL, + service_name TEXT NOT NULL, + status TEXT NOT NULL CHECK (status IN ('running', 'idle', 'stopped', 'error')), + started_at TEXT, + updated_at TEXT NOT NULL, + PRIMARY KEY (user_id, package_id, service_name) + ); + CREATE TABLE account_write_leases ( + token TEXT PRIMARY KEY, + user_id TEXT NOT NULL, + holder TEXT NOT NULL, + acquired_at TEXT NOT NULL, + released_at TEXT + ); + `) + const db = createD1FromSqlite(sqlite) + const meter = createInMemoryUserMeterEnv() + insertUser(sqlite, { stableUserId, d1StorageBytes: 0 }) + const meterStub = userMeterRpc({ env: meter.env, userId: stableUserId }) + await meterStub.initialize({ + resource: 'email_sends_per_day', + day, + count: 4, + updatedAt: now.toISOString(), + }) + await meterStub.initializeStorageBytes({ + bytes: 0, + updatedAt: now.toISOString(), + }) + + const report = await loadAdminUserMeterParityReport({ + db, + env: meter.env, + stableUserId, + now, + }) + expect(report?.daily).toMatchObject({ + mirrorRetired: true, + mismatchCount: 0, + }) + expect( + report?.daily.resources.find( + (row) => row.resource === 'email_sends_per_day', + ), + ).toMatchObject({ + d1Count: null, + meterCount: 4, + delta: null, + parity: true, + }) +}) + +test('loadAdminUserMeterParityReport still compares D1 daily counts when mirror table exists', async () => { + const sqlite = new DatabaseSync(':memory:') + sqlite.exec(` + CREATE TABLE users ( + id INTEGER PRIMARY KEY, + stable_user_id TEXT UNIQUE NOT NULL, + username TEXT NOT NULL, + email TEXT NOT NULL, + password_hash TEXT, + d1_storage_bytes INTEGER NOT NULL DEFAULT 0, + deleting_at TEXT, + active_write_count INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL DEFAULT '2026-01-01T00:00:00.000Z', + updated_at TEXT NOT NULL DEFAULT '2026-01-01T00:00:00.000Z' + ); + CREATE TABLE entitlement_daily_counters ( + user_id TEXT NOT NULL, + resource TEXT NOT NULL, + day TEXT NOT NULL, + count INTEGER NOT NULL DEFAULT 0, + updated_at TEXT NOT NULL, + PRIMARY KEY (user_id, resource, day) + ); + CREATE TABLE package_service_states ( + user_id TEXT NOT NULL, + package_id TEXT NOT NULL, + service_name TEXT NOT NULL, + status TEXT NOT NULL CHECK (status IN ('running', 'idle', 'stopped', 'error')), + started_at TEXT, + updated_at TEXT NOT NULL, + PRIMARY KEY (user_id, package_id, service_name) + ); + CREATE TABLE account_write_leases ( + token TEXT PRIMARY KEY, + user_id TEXT NOT NULL, + holder TEXT NOT NULL, + acquired_at TEXT NOT NULL, + released_at TEXT + ); + `) + const db = createD1FromSqlite(sqlite) + const meter = createInMemoryUserMeterEnv() + insertUser(sqlite, { stableUserId, d1StorageBytes: 0 }) + for (const resource of dailyEntitlementResources) { + sqlite + .prepare( + `INSERT INTO entitlement_daily_counters ( + user_id, resource, day, count, updated_at + ) VALUES (?, ?, ?, ?, ?)`, + ) + .run(stableUserId, resource, day, 1, now.toISOString()) + } + const meterStub = userMeterRpc({ env: meter.env, userId: stableUserId }) + for (const resource of dailyEntitlementResources) { + await meterStub.initialize({ + resource, + day, + count: resource === 'email_sends_per_day' ? 9 : 1, + updatedAt: now.toISOString(), + }) + } + await meterStub.initializeStorageBytes({ + bytes: 0, + updatedAt: now.toISOString(), + }) + const report = await loadAdminUserMeterParityReport({ + db, + env: meter.env, + stableUserId, + now, + }) + expect(report?.daily.mirrorRetired).toBe(false) + expect(report?.daily.mismatchCount).toBe(1) + expect( + report?.daily.resources.find( + (row) => row.resource === 'email_sends_per_day', + ), + ).toMatchObject({ + d1Count: 1, + meterCount: 9, + delta: -8, + parity: false, + }) +}) diff --git a/packages/worker/src/admin/user-meter-parity.ts b/packages/worker/src/admin/user-meter-parity.ts index f5dd00b59a..6cf77d1c6a 100644 --- a/packages/worker/src/admin/user-meter-parity.ts +++ b/packages/worker/src/admin/user-meter-parity.ts @@ -24,10 +24,22 @@ export const userMeterParityPageSize = 500 type DailyResourceParity = { resource: DailyEntitlementResource - d1Count: number + /** + * D1 mirror count for the current UTC day. `null` when + * `daily.mirrorRetired` is true (table dropped / comparison retired). + */ + d1Count: number | null meterCount: number | null needsBootstrap: boolean + /** + * `d1Count - meterCount` when both sides are present; `null` when the + * mirror is retired or the meter still needs bootstrap. + */ delta: number | null + /** + * When the mirror is active: true iff counts match. When + * `mirrorRetired`, always true (no D1 comparison is claimed). + */ parity: boolean } @@ -91,6 +103,12 @@ export type AdminUserMeterParityReport = { stableUserId: string daily: { day: string + /** + * True when D1 `entitlement_daily_counters` is absent (retired). The + * daily gate then reports meter counts only and does not claim a D1 + * comparison (`d1Count`/`delta` are null; `parity` stays true). + */ + mirrorRetired: boolean resources: Array mismatchCount: number } @@ -131,6 +149,23 @@ async function userExists(db: D1Database, stableUserId: string) { return row != null } +/** + * Detect retirement without querying the missing table. Uses sqlite_master + * only (safe when the table has already been dropped). + */ +async function isEntitlementDailyCountersMirrorRetired( + db: D1Database, +): Promise { + const row = await db + .prepare( + `SELECT 1 AS present + FROM sqlite_master + WHERE type = 'table' AND name = 'entitlement_daily_counters'`, + ) + .first<{ present: number }>() + return row == null +} + async function readD1DailyCounts(input: { db: D1Database stableUserId: string @@ -170,11 +205,11 @@ async function readDailyParity(input: { day: string generatedAt: string }): Promise { - const d1Counts = await readD1DailyCounts(input) + const mirrorRetired = await isEntitlementDailyCountersMirrorRetired(input.db) + const d1Counts = mirrorRetired ? null : await readD1DailyCounts(input) const meter = userMeterRpc({ env: input.env, userId: input.stableUserId }) const resources: Array = [] for (const resource of dailyEntitlementResources) { - const d1Count = d1Counts.get(resource) ?? 0 const meterRead = await meter.read({ resource, day: input.day, @@ -182,6 +217,19 @@ async function readDailyParity(input: { }) const needsBootstrap = meterRead.outcome === 'needs_bootstrap' const meterCount = needsBootstrap ? null : meterRead.count + if (mirrorRetired) { + resources.push({ + resource, + d1Count: null, + meterCount, + needsBootstrap, + delta: null, + // Mirror comparison is retired; do not claim a D1 mismatch. + parity: true, + }) + continue + } + const d1Count = d1Counts?.get(resource) ?? 0 const { delta, parity } = countDeltaParity({ d1Value: d1Count, meterValue: meterCount, @@ -198,8 +246,11 @@ async function readDailyParity(input: { } return { day: input.day, + mirrorRetired, resources, - mismatchCount: resources.filter((row) => !row.parity).length, + mismatchCount: mirrorRetired + ? 0 + : resources.filter((row) => !row.parity).length, } } diff --git a/packages/worker/src/admin/user-usage-data.node.test.ts b/packages/worker/src/admin/user-usage-data.node.test.ts index 4cf563112b..c9b6248c03 100644 --- a/packages/worker/src/admin/user-usage-data.node.test.ts +++ b/packages/worker/src/admin/user-usage-data.node.test.ts @@ -32,13 +32,6 @@ type UsageRollupRow = { total_bytes: number } -type CounterRow = { - user_id: string - resource: string - day: string - count: number -} - type ResourceCount = Partial< Record< | 'saved_packages' @@ -59,12 +52,10 @@ function normalizeQuery(query: string) { function createAdminUserUsageTestDb(input: { users: Array usageRollups?: Array - dailyCounters?: Array resourceCounts?: Record }) { const users = input.users.map((user) => ({ ...user })) const usageRollups = input.usageRollups?.map((row) => ({ ...row })) ?? [] - const dailyCounters = input.dailyCounters?.map((row) => ({ ...row })) ?? [] const resourceCounts = input.resourceCounts ?? {} function countForQuery(normalizedQuery: string, userId: string) { @@ -106,15 +97,6 @@ function createAdminUserUsageTestDb(input: { return (users.find((user) => user.stable_user_id === params[0]) ?? null) as T | null } - if (normalizedQuery.includes('from entitlement_daily_counters')) { - const row = dailyCounters.find( - (counter) => - counter.user_id === params[0] && - counter.resource === params[1] && - counter.day === params[2], - ) - return (row ? { count: row.count } : null) as T | null - } const count = countForQuery(normalizedQuery, String(params[0])) if (count !== null) return { count } as T throw new Error(`Unsupported first query: ${query}`) @@ -243,14 +225,6 @@ test('loadAdminUserUsageData warns above eighty percent of plan limits', async ( event_count: 40, }), ], - dailyCounters: [ - { - user_id: usageUserId, - resource: 'email_sends_per_day', - day: '2026-07-05', - count: 170, - }, - ], resourceCounts: { [usageUserId]: { saved_packages: 85, @@ -262,8 +236,15 @@ test('loadAdminUserUsageData warns above eighty percent of plan limits', async ( }, }) + const env = withUserMeter({ APP_DB: db }) + await env.meter.seed({ + userId: usageUserId, + resource: 'email_sends_per_day', + day: '2026-07-05', + count: 170, + }) const data = await loadAdminUserUsageData( - withUserMeter({ APP_DB: db }) as Env, + env as Env, usageUserId, new Date('2026-07-05T12:00:00.000Z'), ) @@ -293,14 +274,6 @@ test('loadAdminUserUsageData rejects an invalid stored plan', async () => { stable_user_id: usageUserId, }, ], - dailyCounters: [ - { - user_id: usageUserId, - resource: 'email_receives_per_day', - day: '2026-07-05', - count: 190, - }, - ], resourceCounts: { [usageUserId]: { stored_email_messages: 12 }, }, @@ -439,43 +412,42 @@ test('loadAdminUserUsageData keeps current-month and month-over-month rollups on expect(getEventCount(data?.monthUsage[1]?.usage, 'execute')).toBe(12) }) -test('loadAdminUserUsageData reads daily counts from UserMeter (bootstrap then warm)', async () => { +test('loadAdminUserUsageData reads daily counts from UserMeter (seeded then warm)', async () => { const now = new Date('2026-07-05T12:00:00.000Z') const day = utcDayKey(now) const bootstrapEmail = 'bootstrap-drilldown@example.com' const bootstrapUserId = await createStableUserIdFromEmail(bootstrapEmail) - const bootstrapped = await loadAdminUserUsageData( - withUserMeter({ - APP_DB: createAdminUserUsageTestDb({ - users: [ - { - id: 3, - username: 'bootstrap', - email: bootstrapEmail, - plan: 'pro', - stable_user_id: bootstrapUserId, - }, - ], - dailyCounters: [ - { - user_id: bootstrapUserId, - resource: 'email_receives_per_day', - day, - count: 55, - }, - { - user_id: bootstrapUserId, - resource: 'outbound_fetches_per_day', - day, - count: 66, - }, - ], - resourceCounts: { - [bootstrapUserId]: { secrets: 3 }, + const bootstrapEnv = withUserMeter({ + APP_DB: createAdminUserUsageTestDb({ + users: [ + { + id: 3, + username: 'bootstrap', + email: bootstrapEmail, + plan: 'pro', + stable_user_id: bootstrapUserId, }, - }), - }) as Env, + ], + resourceCounts: { + [bootstrapUserId]: { secrets: 3 }, + }, + }), + }) + await bootstrapEnv.meter.seed({ + userId: bootstrapUserId, + resource: 'email_receives_per_day', + day, + count: 55, + }) + await bootstrapEnv.meter.seed({ + userId: bootstrapUserId, + resource: 'outbound_fetches_per_day', + day, + count: 66, + }) + const bootstrapped = await loadAdminUserUsageData( + bootstrapEnv as Env, bootstrapUserId, now, ) @@ -508,32 +480,6 @@ test('loadAdminUserUsageData reads daily counts from UserMeter (bootstrap then w stable_user_id: meterUserId, }, ], - dailyCounters: [ - { - user_id: meterUserId, - resource: 'email_sends_per_day', - day, - count: 11, - }, - { - user_id: meterUserId, - resource: 'email_receives_per_day', - day, - count: 22, - }, - { - user_id: meterUserId, - resource: 'execute_calls_per_day', - day, - count: 33, - }, - { - user_id: meterUserId, - resource: 'outbound_fetches_per_day', - day, - count: 44, - }, - ], resourceCounts: { [meterUserId]: { saved_packages: 6, diff --git a/packages/worker/src/app/account-deletion.node.test.ts b/packages/worker/src/app/account-deletion.node.test.ts index 01eb593fe9..170f7cfd1f 100644 --- a/packages/worker/src/app/account-deletion.node.test.ts +++ b/packages/worker/src/app/account-deletion.node.test.ts @@ -897,10 +897,6 @@ test('deleteUserAccount cascades user-scoped rows for the requested user', async email_attachments: [{ id: 'ea-1', message_id: 'em-1' }], email_delivery_events: [{ id: 'ed-1', user_id: userAaa }], email_sender_identities: [{ id: 'ei-1', user_id: userAaa }], - entitlement_daily_counters: [ - { user_id: userAaa, resource: 'email_sends_per_day', day: '2026-07-05' }, - { user_id: userBbb, resource: 'email_sends_per_day', day: '2026-07-05' }, - ], platform_feedback: [ { id: 'feedback-submitted-by-a', @@ -1306,9 +1302,6 @@ test('deleteUserAccount cascades user-scoped rows for the requested user', async 'email-raw:v1:user-aaa/em-1', 'email-raw:v1:user-aaa/em-2', ]) - expect(rows.entitlement_daily_counters).toEqual([ - { user_id: userBbb, resource: 'email_sends_per_day', day: '2026-07-05' }, - ]) expect(rows.platform_feedback).toEqual([ { id: 'feedback-reviewed-by-a', diff --git a/packages/worker/src/app/account-retention-dispositions.ts b/packages/worker/src/app/account-retention-dispositions.ts index 396c01cdcf..d4f23694d0 100644 --- a/packages/worker/src/app/account-retention-dispositions.ts +++ b/packages/worker/src/app/account-retention-dispositions.ts @@ -16,7 +16,6 @@ export const accountRetentionDispositions: ReadonlyArray) { return { ...env, ...meter.env, ...runLog.env, meter, runLog } } -type DailyCounterRow = { - user_id: string - resource: string - day: string - count: number -} - function createUsageTestDb(input: { userId: number email: string plan: string stripePlan?: string | null packageCount?: number - dailyCounters?: Array d1StorageBytes?: number }) { const stableUserId = testStableUserIdFromEmail(input.email) - const dailyCounters = input.dailyCounters?.map((row) => ({ ...row })) ?? [] const d1StorageBytes = input.d1StorageBytes ?? 0 return { stableUserId, @@ -53,15 +44,6 @@ function createUsageTestDb(input: { if (normalized.includes('from saved_packages')) { return { count: input.packageCount ?? 0 } as T } - if (normalized.includes('from entitlement_daily_counters')) { - const row = dailyCounters.find( - (counter) => - counter.user_id === params[0] && - counter.resource === params[1] && - counter.day === params[2], - ) - return (row ? { count: row.count } : { count: 0 }) as T - } if (normalized.includes('d1_storage_bytes')) { return { bytes: d1StorageBytes } as T } @@ -123,23 +105,22 @@ test('loadAccountUsageData returns plan rows and authoritative UserMeter daily c email: bootstrapEmail, plan: 'pro', packageCount: 1, - dailyCounters: [ - { - user_id: bootstrapUserId, - resource: 'email_sends_per_day', - day, - count: 17, - }, - { - user_id: bootstrapUserId, - resource: 'execute_calls_per_day', - day, - count: 91, - }, - ], + }) + const bootstrapEnv = withUsageEnv({ APP_DB: bootstrapDb }) + await bootstrapEnv.meter.seed({ + userId: bootstrapUserId, + resource: 'email_sends_per_day', + day, + count: 17, + }) + await bootstrapEnv.meter.seed({ + userId: bootstrapUserId, + resource: 'execute_calls_per_day', + day, + count: 91, }) const bootstrapped = await loadAccountUsageData({ - env: withUsageEnv({ APP_DB: bootstrapDb }) as Env, + env: bootstrapEnv as Env, userId: 8, now, }) @@ -154,32 +135,6 @@ test('loadAccountUsageData returns plan rows and authoritative UserMeter daily c email: meterEmail, plan: 'pro', packageCount: 4, - dailyCounters: [ - { - user_id: meterUserId, - resource: 'email_sends_per_day', - day, - count: 11, - }, - { - user_id: meterUserId, - resource: 'email_receives_per_day', - day, - count: 22, - }, - { - user_id: meterUserId, - resource: 'execute_calls_per_day', - day, - count: 33, - }, - { - user_id: meterUserId, - resource: 'outbound_fetches_per_day', - day, - count: 44, - }, - ], }) const warmEnv = withUsageEnv({ APP_DB: warmDb }) warmEnv.runLog.setActiveWorkflowCount(meterUserId, 3) diff --git a/packages/worker/src/app/admin-insights-data.node.test.ts b/packages/worker/src/app/admin-insights-data.node.test.ts index 6b3b0e849c..256c879265 100644 --- a/packages/worker/src/app/admin-insights-data.node.test.ts +++ b/packages/worker/src/app/admin-insights-data.node.test.ts @@ -248,17 +248,6 @@ function createInsightsTestDb( ] as Array, } } - if (normalizedQuery.includes('from entitlement_daily_counters')) { - return { - results: [ - { - day: '2026-07-08', - resource: 'email_sends_per_day', - n: 3, - }, - ] as Array, - } - } if (normalizedQuery.includes('from email_delivery_events')) { return { results: [ @@ -322,12 +311,26 @@ function createInsightsTestDb( } test('loadAdminInsightsData assembles the dashboard payload', async () => { + consoleWarn.mockImplementation(() => {}) + consoleWarn.mockClear() const data = await loadAdminInsightsData( - { APP_DB: createInsightsTestDb() } as Env, + { + APP_DB: createInsightsTestDb(), + EMAIL_EVENTS: {} as AnalyticsEngineDataset, + WRANGLER_IS_LOCAL_DEV: 'true', + } as Env, now, ) expect(data.ok).toBe(true) + expect(consoleWarn).toHaveBeenCalledWith( + 'admin-insights-email-quota-aggregate-unavailable', + { reason: 'entitlement-daily-counters-retired-local-dev' }, + ) + expect(consoleWarn).not.toHaveBeenCalledWith( + 'admin-insights-email-quota-aggregate-unavailable', + { reason: 'missing-email-events-binding' }, + ) expect(data.totals).toEqual({ users: 8, verifiedUsers: 5, @@ -353,7 +356,7 @@ test('loadAdminInsightsData assembles the dashboard payload', async () => { expect(data.emailByDay).toHaveLength(28) expect(data.emailByDay.at(-1)).toEqual({ day: '2026-07-08', - sends: 3, + sends: 0, receives: 0, }) expect(data.emailDeliveryByDay).toHaveLength(28) @@ -395,9 +398,32 @@ test('loadAdminInsightsData assembles the dashboard payload', async () => { expect(data.activation.medianHoursToActivation).toBe(36) }) +test('loadAdminInsightsData warns when EMAIL_EVENTS binding is missing', async () => { + consoleWarn.mockImplementation(() => {}) + consoleWarn.mockClear() + const data = await loadAdminInsightsData( + { APP_DB: createInsightsTestDb() } as Env, + now, + ) + + expect(data.ok).toBe(true) + expect(consoleWarn).toHaveBeenCalledWith( + 'admin-insights-email-quota-aggregate-unavailable', + { reason: 'missing-email-events-binding' }, + ) + expect(consoleWarn).not.toHaveBeenCalledWith( + 'admin-insights-email-quota-aggregate-unavailable', + { reason: 'entitlement-daily-counters-retired-local-dev' }, + ) +}) + test('activation latency excludes users who have no usable verification date', async () => { + consoleWarn.mockImplementation(() => {}) const data = await loadAdminInsightsData( - { APP_DB: createInsightsTestDb({ activationLatencyRows: [] }) } as Env, + { + APP_DB: createInsightsTestDb({ activationLatencyRows: [] }), + WRANGLER_IS_LOCAL_DEV: 'true', + } as Env, now, ) // Two users reached the milestone, but neither has timing we can measure. diff --git a/packages/worker/src/app/admin-insights-data.ts b/packages/worker/src/app/admin-insights-data.ts index 9a2c3a7430..c979787e64 100644 --- a/packages/worker/src/app/admin-insights-data.ts +++ b/packages/worker/src/app/admin-insights-data.ts @@ -281,29 +281,31 @@ async function loadEmailInsightsRows(input: { deliveryRows: Array }> { // Wrangler exposes a local Analytics Engine binding, but its SQL API cannot - // query the emulated dataset. Local development therefore keeps using the - // existing D1 tables, matching the feature-flag exposure fallback. - if (!input.env.EMAIL_EVENTS || input.env.WRANGLER_IS_LOCAL_DEV === 'true') { - const [emailRows, deliveryRows] = await Promise.all([ - input.env.APP_DB.prepare( - `SELECT day, resource, SUM(count) AS n - FROM entitlement_daily_counters - WHERE day >= ? AND resource IN ('email_sends_per_day', 'email_receives_per_day') - GROUP BY day, resource`, - ) - .bind(input.dayCutoff) - .all(), - input.env.APP_DB.prepare( - `SELECT substr(created_at, 1, 10) AS day, event_type, COUNT(*) AS n - FROM email_delivery_events - WHERE provider = 'cloudflare-email' AND created_at >= ? - GROUP BY day, event_type`, - ) - .bind(input.dayCutoff) - .all(), - ]) + // query the emulated dataset. The D1 entitlement_daily_counters mirror is + // retired, so local/dev email quota aggregates degrade explicitly to empty + // (with a warning) rather than querying a removed table. Delivery-event + // aggregates still come from D1 email_delivery_events. + if (input.env.WRANGLER_IS_LOCAL_DEV === 'true') { + console.warn('admin-insights-email-quota-aggregate-unavailable', { + reason: 'entitlement-daily-counters-retired-local-dev', + }) + } + if (!input.env.EMAIL_EVENTS) { + console.warn('admin-insights-email-quota-aggregate-unavailable', { + reason: 'missing-email-events-binding', + }) + } + if (input.env.WRANGLER_IS_LOCAL_DEV === 'true' || !input.env.EMAIL_EVENTS) { + const deliveryRows = await input.env.APP_DB.prepare( + `SELECT substr(created_at, 1, 10) AS day, event_type, COUNT(*) AS n + FROM email_delivery_events + WHERE provider = 'cloudflare-email' AND created_at >= ? + GROUP BY day, event_type`, + ) + .bind(input.dayCutoff) + .all() return { - emailRows: emailRows.results ?? [], + emailRows: [], deliveryRows: deliveryRows.results ?? [], } } diff --git a/packages/worker/src/app/handlers/account-email.node.test.ts b/packages/worker/src/app/handlers/account-email.node.test.ts index 69293c4201..7110a388ff 100644 --- a/packages/worker/src/app/handlers/account-email.node.test.ts +++ b/packages/worker/src/app/handlers/account-email.node.test.ts @@ -4,6 +4,7 @@ import type * as EmailPlatformAddress from '#worker/email/platform-address.ts' import type * as EntitlementPlans from '#worker/entitlements/plans.ts' import type * as EntitlementService from '#worker/entitlements/service.ts' import { createInMemoryUserMeterEnv } from '#worker/test-support/user-meter.ts' +import { userMeterRpc } from '#worker/entitlements/user-meter-client.ts' const messageRow = { id: 'msg-1', @@ -251,41 +252,14 @@ function createListResult(rows: Array>) { function createEnv(input?: { meter?: ReturnType - dailyCounters?: Array<{ - user_id: string - resource: string - day: string - count: number - }> messageRows?: Array> messageTotal?: number }) { const meter = input?.meter ?? createInMemoryUserMeterEnv() - const dailyCounters = input?.dailyCounters?.map((row) => ({ ...row })) ?? [] const messageRows = input?.messageRows ?? [messageRow] const messageTotal = input?.messageTotal ?? messageRows.length mockModule.prepare.mockImplementation((query: string) => { const normalized = String(query).replace(/\s+/g, ' ').trim().toLowerCase() - if (normalized.includes('from entitlement_daily_counters')) { - return { - bind(...params: Array) { - return { - async first() { - const row = dailyCounters.find( - (counter) => - counter.user_id === params[0] && - counter.resource === params[1] && - counter.day === params[2], - ) - return (row ? { count: row.count } : { count: 0 }) as T - }, - async all() { - return { results: [] } - }, - } - }, - } - } if (normalized.includes('count(*)')) { return createCountResult(messageTotal) } @@ -321,20 +295,19 @@ test('email API lists messages with pagination, usage, and selected detail', asy const bootstrapMeter = createInMemoryUserMeterEnv() const env = createEnv({ meter: bootstrapMeter, - dailyCounters: [ - { - user_id: userId, - resource: 'email_sends_per_day', - day, - count: 2, - }, - { - user_id: userId, - resource: 'email_receives_per_day', - day, - count: 2, - }, - ], + }) + const meterStub = userMeterRpc({ env: bootstrapMeter.env, userId }) + await meterStub.initialize({ + resource: 'email_sends_per_day', + day, + count: 2, + updatedAt: new Date().toISOString(), + }) + await meterStub.initialize({ + resource: 'email_receives_per_day', + day, + count: 2, + updatedAt: new Date().toISOString(), }) const handler = createAccountEmailApiHandler(env) @@ -389,21 +362,25 @@ test('email API lists messages with pagination, usage, and selected detail', asy mockModule.prepare.mockClear() mockModule.getEmailMessageById.mockClear() + const selectionMeter = createInMemoryUserMeterEnv() const envWithSelection = createEnv({ - dailyCounters: [ - { - user_id: userId, - resource: 'email_sends_per_day', - day, - count: 2, - }, - { - user_id: userId, - resource: 'email_receives_per_day', - day, - count: 2, - }, - ], + meter: selectionMeter, + }) + const selectionStub = userMeterRpc({ + env: selectionMeter.env, + userId, + }) + await selectionStub.initialize({ + resource: 'email_sends_per_day', + day, + count: 2, + updatedAt: new Date().toISOString(), + }) + await selectionStub.initialize({ + resource: 'email_receives_per_day', + day, + count: 2, + updatedAt: new Date().toISOString(), }) const selectedResponse = await createAccountEmailApiHandler( envWithSelection, @@ -465,7 +442,7 @@ test('email API lists messages with pagination, usage, and selected detail', asy expect(unauthorizedResponse.status).toBe(401) }) -test('email API usage prefers initialized UserMeter daily counts over D1', async () => { +test('email API usage reads initialized UserMeter daily counts', async () => { using _frozenTime = useFrozenUtcTime('2026-07-31T15:00:00.000Z') const day = utcDayKey() const userId = 'stable-user-1' @@ -484,20 +461,6 @@ test('email API usage prefers initialized UserMeter daily counts over D1', async }) const env = createEnv({ meter, - dailyCounters: [ - { - user_id: userId, - resource: 'email_sends_per_day', - day, - count: 3, - }, - { - user_id: userId, - resource: 'email_receives_per_day', - day, - count: 5, - }, - ], }) const response = await createAccountEmailApiHandler(env).handler({ request: new Request('https://example.com/account/email.json'), @@ -513,7 +476,7 @@ test('email API usage prefers initialized UserMeter daily counts over D1', async receives_today: { count: 15, limit: 50 }, }), }) - expect(body.usage.sends_today.count).not.toBe(3) + expect(body.usage.sends_today.count).toBe(13) expect(body.usage.receives_today.count).not.toBe(5) }) diff --git a/packages/worker/src/app/retention.node.test.ts b/packages/worker/src/app/retention.node.test.ts index 95832df4cc..670863f3aa 100644 --- a/packages/worker/src/app/retention.node.test.ts +++ b/packages/worker/src/app/retention.node.test.ts @@ -6,14 +6,12 @@ import { auditEventRetentionDays, emailDeliveryEventRetentionDays, emailMessageRetentionDays, - entitlementDailyCounterRetentionDays, featureFlagExposureRetentionDays, getRetentionPolicyCoverage, memorySuppressionRetentionDays, platformFeedbackRetentionDays, pruneAuditEventsForRetention, pruneEmailDeliveryEventsForRetention, - pruneEntitlementDailyCountersForRetention, pruneFeatureFlagExposuresForRetention, pruneMemorySuppressionsForRetention, prunePlatformFeedbackForRetention, @@ -201,14 +199,6 @@ function createRetentionDb() { storage_key TEXT, created_at TEXT NOT NULL ); - CREATE TABLE entitlement_daily_counters ( - user_id TEXT NOT NULL, - resource TEXT NOT NULL, - day TEXT NOT NULL, - count INTEGER NOT NULL DEFAULT 0, - updated_at TEXT NOT NULL, - PRIMARY KEY (user_id, resource, day) - ); CREATE TABLE usage_rollups ( user_id TEXT NOT NULL, metric TEXT NOT NULL, @@ -865,22 +855,8 @@ test('retention prune reports selected separately from deleted when rows vanish ).toEqual({ selected: 2, deleted: 1 }) }) -test('entitlement counter and usage rollup retention respect boundaries', async () => { +test('usage rollup retention respects month boundaries', async () => { const { sqlite, db } = createRetentionDb() - const cutoffDay = daysAgo(entitlementDailyCounterRetentionDays).slice(0, 10) - for (const day of [ - daysAgo(entitlementDailyCounterRetentionDays + 1).slice(0, 10), - cutoffDay, - daysAgo(1).slice(0, 10), - ]) { - sqlite - .prepare( - `INSERT INTO entitlement_daily_counters ( - user_id, resource, day, count, updated_at - ) VALUES ('user-1', 'email_sends_per_day', ?, 1, ?)`, - ) - .run(day, now.toISOString()) - } // 24 months before 2026-07 keeps 2024-07 and later. for (const month of ['2024-06', '2024-07', '2026-06']) { sqlite @@ -892,22 +868,11 @@ test('entitlement counter and usage rollup retention respect boundaries', async .run(month, now.toISOString()) } - expect(await pruneEntitlementDailyCountersForRetention({ db, now })).toEqual({ - selected: 1, - deleted: 1, - }) expect(await pruneUsageRollupsForRetention({ db, now })).toEqual({ selected: 1, deleted: 1, }) - const days = sqlite - .prepare(`SELECT day FROM entitlement_daily_counters ORDER BY day`) - .all() as Array<{ day: string }> - expect(days.map((row) => row.day)).toEqual([ - cutoffDay, - daysAgo(1).slice(0, 10), - ]) const months = sqlite .prepare(`SELECT month FROM usage_rollups ORDER BY month`) .all() as Array<{ month: string }> diff --git a/packages/worker/src/app/retention.ts b/packages/worker/src/app/retention.ts index 423effa740..7db0133a1d 100644 --- a/packages/worker/src/app/retention.ts +++ b/packages/worker/src/app/retention.ts @@ -49,7 +49,6 @@ export const platformFeedbackRetentionDays = 365 export const publishedBundleArtifactRetentionDays = 30 export const emailDeliveryEventRetentionDays = 90 export const emailMessageRetentionDays = 365 -export const entitlementDailyCounterRetentionDays = 400 export const usageRollupRetentionMonths = 24 export const featureFlagExposureRetentionDays = 90 export const auditEventRetentionDays = 180 @@ -133,14 +132,6 @@ export const retentionPolicies: ReadonlyArray = [ description: 'Threads orphaned by email message retention are pruned for the affected users in the same run.', }, - { - table: 'entitlement_daily_counters', - scope: 'per-user', - retentionDays: entitlementDailyCounterRetentionDays, - batchSize: retentionDefaultBatchSize, - description: - 'Daily entitlement rate counters keep 400 days so year-over-year usage stays inspectable before rows age out.', - }, { table: 'usage_rollups', scope: 'per-user', @@ -210,7 +201,6 @@ export type RetentionPruneResult = { deletedAttachmentBlobs: number blobDeleteErrors: number } - entitlementDailyCounters: number usageRollups: number featureFlagExposureRollups: number auditEvents: number @@ -774,29 +764,6 @@ export async function pruneUserEmailMessagesForRetention(input: { return result } -export async function pruneEntitlementDailyCountersForRetention(input: { - db: D1Database - now?: Date - batchSize?: number -}) { - const cutoffDay = cutoffIso( - input.now ?? new Date(), - entitlementDailyCounterRetentionDays, - ).slice(0, 'YYYY-MM-DD'.length) - return selectAndDeleteByIds({ - db: input.db, - column: 'rowid', - bindings: [cutoffDay, input.batchSize ?? retentionDefaultBatchSize], - sql: `SELECT rowid - FROM entitlement_daily_counters - WHERE day < ? - ORDER BY day ASC, rowid ASC - LIMIT ?`, - table: 'entitlement_daily_counters', - idColumn: 'rowid', - }) -} - export async function pruneUsageRollupsForRetention(input: { db: D1Database now?: Date @@ -918,7 +885,6 @@ export async function pruneRetention(input: { deletedAttachmentBlobs: 0, blobDeleteErrors: 0, }, - entitlementDailyCounters: 0, usageRollups: 0, featureFlagExposureRollups: 0, auditEvents: 0, @@ -1011,13 +977,6 @@ export async function pruneRetention(input: { return batch.hasMore }, }, - countTask( - 'entitlement_daily_counters', - () => pruneEntitlementDailyCountersForRetention({ db, now }), - (count) => { - result.entitlementDailyCounters += count - }, - ), countTask( 'usage_rollups', () => pruneUsageRollupsForRetention({ db, now }), diff --git a/packages/worker/src/email/inbound-delivery.ts b/packages/worker/src/email/inbound-delivery.ts index f5499bf1fb..44726389c7 100644 --- a/packages/worker/src/email/inbound-delivery.ts +++ b/packages/worker/src/email/inbound-delivery.ts @@ -5,7 +5,6 @@ import { EntitlementLimitError, } from '#worker/entitlements/errors.ts' import { type PlanName } from '#worker/entitlements/plans.ts' -import { scheduleAbsoluteDailyEntitlementMirror } from '#worker/entitlements/service.ts' import { userMeterRpc, type UserMeterEnv, @@ -423,25 +422,11 @@ export async function pruneExpiredInboundDedupePointers(input: { ) } -async function readCounter(input: { +async function readSystemEmailDailyCounter(input: { db: D1Database - table: 'entitlement_daily_counters' | 'system_email_daily_counters' - userId?: string - localPart?: SystemEmailLocal + localPart: SystemEmailLocal day: string }) { - if (input.table === 'entitlement_daily_counters') { - const row = await input.db - .prepare( - `SELECT count FROM entitlement_daily_counters - WHERE user_id = ? - AND resource = 'email_receives_per_day' - AND day = ?`, - ) - .bind(input.userId, input.day) - .first<{ count: number }>() - return Number(row?.count ?? 0) - } const row = await input.db .prepare( `SELECT count FROM system_email_daily_counters @@ -452,7 +437,12 @@ async function readCounter(input: { return Number(row?.count ?? 0) } -/** Point-read today's user inbound receive count from UserMeter. */ +/** + * Point-read today's user inbound receive count from UserMeter. Cold meters + * initialize at zero; never touches the retired D1 daily counter table. + * + * `db` remains for call-site stability. + */ export async function readUserInboundReceiveCount(input: { db: D1Database env: UserMeterEnv @@ -460,6 +450,7 @@ export async function readUserInboundReceiveCount(input: { day: string now?: Date }) { + void input.db const now = input.now ?? new Date() const updatedAt = now.toISOString() const meter = userMeterRpc({ env: input.env, userId: input.userId }) @@ -469,16 +460,10 @@ export async function readUserInboundReceiveCount(input: { now: updatedAt, }) if (result.outcome === 'needs_bootstrap') { - const baseline = await readCounter({ - db: input.db, - table: 'entitlement_daily_counters', - userId: input.userId, - day: input.day, - }) await meter.initialize({ resource: 'email_receives_per_day', day: input.day, - count: baseline, + count: 0, updatedAt, }) result = await meter.read({ @@ -500,9 +485,8 @@ export async function readSystemInboundReceiveCount(input: { localPart: SystemEmailLocal day: string }) { - return await readCounter({ + return await readSystemEmailDailyCounter({ db: input.db, - table: 'system_email_daily_counters', localPart: input.localPart, day: input.day, }) @@ -521,7 +505,8 @@ async function resolveConcurrentDelivery(input: { /** * Charge one inbound receive via UserMeter, then persist the D1 delivery - * event. Delivery-id idempotency lives in the DO; D1 mirror is best-effort. + * event. Delivery-id idempotency lives in the DO. Never touches the retired + * D1 daily counter table. */ export async function chargeUserInboundDeliveryOnce(input: { db: D1Database @@ -530,8 +515,13 @@ export async function chargeUserInboundDeliveryOnce(input: { plan: PlanName limit: number now: Date + /** + * Retained for call-site stability. Daily counters no longer schedule D1 + * mirror work after mirror retirement. + */ waitUntil?: (promise: Promise) => void }): Promise { + void input.waitUntil const existing = await resolveConcurrentDelivery(input) if (existing) return existing @@ -548,16 +538,10 @@ export async function chargeUserInboundDeliveryOnce(input: { updatedAt, }) if (meterResult.outcome === 'needs_bootstrap') { - const baseline = await readCounter({ - db: input.db, - table: 'entitlement_daily_counters', - userId: input.delivery.userId, - day: input.delivery.quotaDay, - }) await meter.initialize({ resource: 'email_receives_per_day', day: input.delivery.quotaDay, - count: baseline, + count: 0, updatedAt, }) meterResult = await meter.consumeInboundDelivery({ @@ -585,16 +569,6 @@ export async function chargeUserInboundDeliveryOnce(input: { }) } - scheduleAbsoluteDailyEntitlementMirror({ - db: input.db, - userId: input.delivery.userId, - resource: meterResult.resource, - day: meterResult.day, - count: meterResult.count, - mirrorUpdatedAt: meterResult.mirrorUpdatedAt, - waitUntil: input.waitUntil, - }) - try { const insert = await input.db .prepare( @@ -635,9 +609,8 @@ export async function chargeSystemInboundDeliveryOnce(input: { }) { const existing = await resolveConcurrentDelivery(input) if (existing) return { delivery: existing, overLimit: false as const } - const current = await readCounter({ + const current = await readSystemEmailDailyCounter({ db: input.db, - table: 'system_email_daily_counters', localPart: input.localPart, day: input.delivery.quotaDay, }) diff --git a/packages/worker/src/email/inbound-entitlements.workers.test.ts b/packages/worker/src/email/inbound-entitlements.workers.test.ts index 53a9de2f46..fc9ca4746e 100644 --- a/packages/worker/src/email/inbound-entitlements.workers.test.ts +++ b/packages/worker/src/email/inbound-entitlements.workers.test.ts @@ -93,15 +93,6 @@ async function seedDailyReceiveCounter(userId: string, count: number) { updatedAt, ) }) - await env.APP_DB.prepare( - `INSERT INTO entitlement_daily_counters (user_id, resource, day, count, updated_at) - VALUES (?, 'email_receives_per_day', ?, ?, ?) - ON CONFLICT(user_id, resource, day) DO UPDATE SET - count = excluded.count, - updated_at = excluded.updated_at`, - ) - .bind(userId, day, count, updatedAt) - .run() } async function readDailyReceiveCounter(userId: string) { @@ -888,19 +879,4 @@ test('inbound UserMeter replay keeps original claim day across UTC midnight', as expect( await meter.read({ resource: 'email_receives_per_day', day: dayN1 }), ).toMatchObject({ outcome: 'needs_bootstrap' }) - - const mirroredDayN = await env.APP_DB.prepare( - `SELECT count FROM entitlement_daily_counters - WHERE user_id = ? AND resource = 'email_receives_per_day' AND day = ?`, - ) - .bind(userId, dayN) - .first<{ count: number }>() - const mirroredDayN1 = await env.APP_DB.prepare( - `SELECT count FROM entitlement_daily_counters - WHERE user_id = ? AND resource = 'email_receives_per_day' AND day = ?`, - ) - .bind(userId, dayN1) - .first<{ count: number }>() - expect(Number(mirroredDayN?.count ?? 0)).toBe(1) - expect(mirroredDayN1).toBeNull() }, 30_000) diff --git a/packages/worker/src/email/inbound-spam-controls.workers.test.ts b/packages/worker/src/email/inbound-spam-controls.workers.test.ts index 09aade2a83..9d36ecb790 100644 --- a/packages/worker/src/email/inbound-spam-controls.workers.test.ts +++ b/packages/worker/src/email/inbound-spam-controls.workers.test.ts @@ -1,5 +1,6 @@ import { env } from 'cloudflare:workers' import { expect, test, vi } from 'vitest' +import { userMeterRpc } from '#worker/entitlements/user-meter-client.ts' import type * as PackageSubscriptionsModule from './package-subscriptions.ts' import { handleInboundEmail } from './inbound.ts' import { processInboundDeliveryEffects } from './inbound-effects.ts' @@ -68,13 +69,11 @@ async function seedVerifiedAccount(input: { } async function readUserDailyReceiveCount(userId: string) { - const row = await env.APP_DB.prepare( - `SELECT count FROM entitlement_daily_counters - WHERE user_id = ? AND resource = 'email_receives_per_day' AND day = ?`, - ) - .bind(userId, new Date().toISOString().slice(0, 10)) - .first<{ count: number }>() - return Number(row?.count ?? 0) + const result = await userMeterRpc({ env, userId }).read({ + resource: 'email_receives_per_day', + day: new Date().toISOString().slice(0, 10), + }) + return result.outcome === 'ready' ? result.count : 0 } function dmarcFailAuthResults() { diff --git a/packages/worker/src/email/inbound.workers.test.ts b/packages/worker/src/email/inbound.workers.test.ts index 2505321ec9..e00a53a583 100644 --- a/packages/worker/src/email/inbound.workers.test.ts +++ b/packages/worker/src/email/inbound.workers.test.ts @@ -1,5 +1,6 @@ import { env } from 'cloudflare:workers' import { expect, test, vi } from 'vitest' +import { utcDayKey } from '@kody-internal/shared/date-keys.ts' import { userMeterRpc } from '#worker/entitlements/user-meter-client.ts' import { handleInboundEmail } from './inbound.ts' import { processInboundDeliveryEffects } from './inbound-effects.ts' @@ -576,13 +577,12 @@ test('inbound email handler rejects mail for unverified accounts', async () => { limit: 10, }) expect(messages).toEqual([]) - const counterRow = await env.APP_DB.prepare( - `SELECT count FROM entitlement_daily_counters - WHERE user_id = ? AND resource = 'email_receives_per_day'`, - ) - .bind(userId) - .first<{ count: number }>() - expect(counterRow).toBeNull() + expect( + await userMeterRpc({ env, userId }).read({ + resource: 'email_receives_per_day', + day: utcDayKey(), + }), + ).toEqual({ outcome: 'needs_bootstrap' }) const events = await env.APP_DB.prepare( `SELECT event_type, detail_json FROM email_delivery_events WHERE user_id = ?`, ) @@ -1339,13 +1339,12 @@ test('pointer-only retry after midnight enforces the current quota day', async ( }) const receiveLimit = planLimits.free.maxEmailReceivesPerDay if (receiveLimit == null) throw new Error('Expected finite receive limit.') - await env.APP_DB.prepare( - `INSERT INTO entitlement_daily_counters ( - user_id, resource, day, count, updated_at - ) VALUES (?, 'email_receives_per_day', '2026-07-23', ?, ?)`, - ) - .bind(userId, receiveLimit, retryNow.toISOString()) - .run() + await userMeterRpc({ env, userId }).initialize({ + resource: 'email_receives_per_day', + day: '2026-07-23', + count: receiveLimit, + updatedAt: retryNow.toISOString(), + }) vi.setSystemTime(retryNow) const retry = createForwardableEmailMessage({ diff --git a/packages/worker/src/email/outbound.ts b/packages/worker/src/email/outbound.ts index ea135691fc..5b3c2d0684 100644 --- a/packages/worker/src/email/outbound.ts +++ b/packages/worker/src/email/outbound.ts @@ -588,8 +588,7 @@ export async function sendOutboundEmail( }) + attachmentBytesTotal, }) - // Atomic check-and-increment via UserMeter. No ExecutionContext here; - // the legacy D1 mirror uses the service's caught best-effort fallback. + // Atomic check-and-increment via UserMeter (sole daily authority). await consumeDailyEntitlement({ db: input.env.APP_DB, env: input.env, diff --git a/packages/worker/src/entitlements/entitlements.node.test.ts b/packages/worker/src/entitlements/entitlements.node.test.ts index f95a7c9a1a..b5c454dded 100644 --- a/packages/worker/src/entitlements/entitlements.node.test.ts +++ b/packages/worker/src/entitlements/entitlements.node.test.ts @@ -33,18 +33,8 @@ import { refundDailyEntitlement, } from './service.ts' import { utcDayKey } from '@kody-internal/shared/date-keys.ts' -import { - createInMemoryUserMeterEnv, - createWaitUntilDrain, -} from '#worker/test-support/user-meter.ts' - -type CounterRow = { - user_id: string - resource: string - day: string - count: number - updated_at?: string -} +import { createInMemoryUserMeterEnv } from '#worker/test-support/user-meter.ts' +import { userMeterRpc } from './user-meter-client.ts' function createEntitlementsTestDb( input: { @@ -71,12 +61,10 @@ function createEntitlementsTestDb( number > > - counters?: Array } = {}, ) { const users = input.users ?? [] const counts = input.counts ?? {} - const counters = input.counters ?? [] const queries: Array<{ sql: string; params: Array }> = [] const initialStorageBytes = Object.values(counts).reduce( (total, count) => total + (count ?? 0), @@ -180,15 +168,6 @@ function createEntitlementsTestDb( const bytes = storageBytesByUser.get(String(params[0])) return (bytes === undefined ? null : { bytes }) as T | null } - if (query.includes('FROM entitlement_daily_counters')) { - const row = counters.find( - (counter) => - counter.user_id === params[0] && - counter.resource === params[1] && - counter.day === params[2], - ) - return (row ? { count: row.count } : null) as T | null - } const count = countFor(query) if (count !== null) { return { count } as T @@ -210,66 +189,6 @@ function createEntitlementsTestDb( storageBytesByUser.set(userId, existing + Number(params[0])) return { meta: { changes: 1 } } } - if (query.includes('INSERT INTO entitlement_daily_counters')) { - const isAbsoluteMirror = - query.includes('count = excluded.count') && - query.includes('updated_at < excluded.updated_at') - const isConditionalConsume = query.includes('count + 1 <= ?') - const amount = isConditionalConsume ? 1 : Number(params[3]) - const mirrorUpdatedAt = String(params[4] ?? '') - const existing = counters.find( - (counter) => - counter.user_id === params[0] && - counter.resource === params[1] && - counter.day === params[2], - ) - if (existing) { - if ( - isConditionalConsume && - existing.count + 1 > Number(params[4]) - ) { - return { meta: { changes: 0 } } - } - if (isAbsoluteMirror) { - const existingUpdatedAt = existing.updated_at ?? '' - if ( - existingUpdatedAt !== '' && - !(existingUpdatedAt < mirrorUpdatedAt) - ) { - return { meta: { changes: 0 } } - } - existing.count = amount - existing.updated_at = mirrorUpdatedAt - } else { - existing.count += amount - } - } else { - counters.push({ - user_id: String(params[0]), - resource: String(params[1]), - day: String(params[2]), - count: amount, - ...(isAbsoluteMirror - ? { updated_at: mirrorUpdatedAt } - : {}), - }) - } - return { meta: { changes: 1 } } - } - if (query.includes('UPDATE entitlement_daily_counters')) { - // bind(updated_at, user_id, resource, day) - const existing = counters.find( - (counter) => - counter.user_id === params[1] && - counter.resource === params[2] && - counter.day === params[3], - ) - if (existing) { - existing.count = Math.max(0, existing.count - 1) - return { meta: { changes: 1 } } - } - return { meta: { changes: 0 } } - } throw new Error(`Unsupported run query: ${query}`) }, } @@ -278,7 +197,28 @@ function createEntitlementsTestDb( }, } as unknown as D1Database - return { db, counters, queries } + return { db, queries } +} + +async function readMeterDailyCount(input: { + env: ReturnType['env'] + userId: string + resource: + | 'email_sends_per_day' + | 'email_receives_per_day' + | 'execute_calls_per_day' + | 'outbound_fetches_per_day' + now: Date +}) { + const result = await userMeterRpc({ + env: input.env, + userId: input.userId, + }).read({ + resource: input.resource, + day: utcDayKey(input.now), + now: input.now.toISOString(), + }) + return result.outcome === 'ready' ? result.count : 0 } const plannedEmail = 'planned@example.com' @@ -770,11 +710,10 @@ test('persistent package services are gated as a 0/1 limit', async () => { test('plan user daily entitlements increment, enforce at limit, and reset on a new UTC day', async () => { const userId = await createStableUserIdFromEmail(plannedEmail) const now = new Date('2026-07-05T15:00:00.000Z') - const { db, counters } = createEntitlementsTestDb({ + const { db } = createEntitlementsTestDb({ users: [{ email: plannedEmail, plan: 'free', stable_user_id: userId }], }) const { env } = createInMemoryUserMeterEnv() - const mirror = createWaitUntilDrain() expect(utcDayKey(now)).toBe('2026-07-05') const limit = planLimits.free.maxEmailSendsPerDay @@ -787,19 +726,16 @@ test('plan user daily entitlements increment, enforce at limit, and reset on a n email: plannedEmail, resource: 'email_sends_per_day', now, - waitUntil: mirror.waitUntil, }) } - await mirror.drain() - expect(counters).toEqual([ - { - user_id: userId, + expect( + await readMeterDailyCount({ + env, + userId, resource: 'email_sends_per_day', - day: '2026-07-05', - count: limit, - updated_at: expect.stringMatching(/^r\/0+[0-9]+$/), - }, - ]) + now, + }), + ).toBe(limit) await expect( assertWithinEntitlement({ db, @@ -808,7 +744,7 @@ test('plan user daily entitlements increment, enforce at limit, and reset on a n resource: 'email_sends_per_day', now, }), - ).rejects.toBeInstanceOf(EntitlementLimitError) + ).rejects.toThrow(/must be read from UserMeter/) const denied = await consumeDailyEntitlement({ db, @@ -817,7 +753,6 @@ test('plan user daily entitlements increment, enforce at limit, and reset on a n email: plannedEmail, resource: 'email_sends_per_day', now, - waitUntil: mirror.waitUntil, }).then( () => null, (thrown: unknown) => thrown, @@ -832,7 +767,6 @@ test('plan user daily entitlements increment, enforce at limit, and reset on a n limit, current: limit, }) - expect(counters[0]?.count).toBe(limit) const nextDay = new Date('2026-07-06T00:00:01.000Z') await consumeDailyEntitlement({ @@ -842,17 +776,28 @@ test('plan user daily entitlements increment, enforce at limit, and reset on a n email: plannedEmail, resource: 'email_sends_per_day', now: nextDay, - waitUntil: mirror.waitUntil, }) - await mirror.drain() - expect(counters).toHaveLength(2) - expect(counters[1]?.count).toBe(1) + expect( + await readMeterDailyCount({ + env, + userId, + resource: 'email_sends_per_day', + now: nextDay, + }), + ).toBe(1) + expect( + await readMeterDailyCount({ + env, + userId, + resource: 'email_sends_per_day', + now, + }), + ).toBe(limit) }) test('refundDailyEntitlement decrements the user/day counter and floors at zero', async () => { - const { db, counters } = createEntitlementsTestDb() + const { db } = createEntitlementsTestDb() const { env } = createInMemoryUserMeterEnv() - const mirror = createWaitUntilDrain() const now = new Date('2026-07-05T15:00:00.000Z') for (let index = 0; index < 2; index += 1) { await consumeDailyEntitlement({ @@ -862,7 +807,6 @@ test('refundDailyEntitlement decrements the user/day counter and floors at zero' email: null, resource: 'email_receives_per_day', now, - waitUntil: mirror.waitUntil, }) } for (let index = 0; index < 3; index += 1) { @@ -873,30 +817,30 @@ test('refundDailyEntitlement decrements the user/day counter and floors at zero' email: null, resource: 'email_receives_per_day', now, - waitUntil: mirror.waitUntil, }) } - await mirror.drain() await refundDailyEntitlement({ db, env, userId: 'user-1', resource: 'email_receives_per_day', now, - waitUntil: mirror.waitUntil, }) - await mirror.drain() expect( - counters.find( - (row) => - row.user_id === 'user-1' && row.resource === 'email_receives_per_day', - )?.count, + await readMeterDailyCount({ + env, + userId: 'user-1', + resource: 'email_receives_per_day', + now, + }), ).toBe(1) expect( - counters.find( - (row) => - row.user_id === 'user-2' && row.resource === 'email_receives_per_day', - )?.count, + await readMeterDailyCount({ + env, + userId: 'user-2', + resource: 'email_receives_per_day', + now, + }), ).toBe(3) await refundDailyEntitlement({ @@ -905,7 +849,6 @@ test('refundDailyEntitlement decrements the user/day counter and floors at zero' userId: 'user-1', resource: 'email_receives_per_day', now, - waitUntil: mirror.waitUntil, }) await refundDailyEntitlement({ db, @@ -913,21 +856,20 @@ test('refundDailyEntitlement decrements the user/day counter and floors at zero' userId: 'user-1', resource: 'email_receives_per_day', now, - waitUntil: mirror.waitUntil, }) - await mirror.drain() expect( - counters.find( - (row) => - row.user_id === 'user-1' && row.resource === 'email_receives_per_day', - )?.count, + await readMeterDailyCount({ + env, + userId: 'user-1', + resource: 'email_receives_per_day', + now, + }), ).toBe(0) }) test('missing-email lookups fail closed and honor free email caps', async () => { - const { db, counters } = createEntitlementsTestDb() + const { db } = createEntitlementsTestDb() const { env } = createInMemoryUserMeterEnv() - const mirror = createWaitUntilDrain() const sendLimit = planLimits.free.maxEmailSendsPerDay const now = new Date('2026-07-05T15:00:00.000Z') for (let index = 0; index < sendLimit; index += 1) { @@ -938,11 +880,16 @@ test('missing-email lookups fail closed and honor free email caps', async () => email: null, resource: 'email_sends_per_day', now, - waitUntil: mirror.waitUntil, }) } - await mirror.drain() - expect(counters[0]?.count).toBe(sendLimit) + expect( + await readMeterDailyCount({ + env, + userId: 'user-1', + resource: 'email_sends_per_day', + now, + }), + ).toBe(sendLimit) await expect( consumeDailyEntitlement({ db, @@ -951,7 +898,6 @@ test('missing-email lookups fail closed and honor free email caps', async () => email: null, resource: 'email_sends_per_day', now, - waitUntil: mirror.waitUntil, }), ).rejects.toBeInstanceOf(EntitlementLimitError) @@ -964,12 +910,15 @@ test('missing-email lookups fail closed and honor free email caps', async () => email: null, resource: 'email_receives_per_day', now, - waitUntil: mirror.waitUntil, }) } - await mirror.drain() expect( - counters.find((row) => row.resource === 'email_receives_per_day')?.count, + await readMeterDailyCount({ + env, + userId: 'user-1', + resource: 'email_receives_per_day', + now, + }), ).toBe(receiveLimit) const denied = await consumeDailyEntitlement({ @@ -979,7 +928,6 @@ test('missing-email lookups fail closed and honor free email caps', async () => email: null, resource: 'email_receives_per_day', now, - waitUntil: mirror.waitUntil, }).then( () => null, (thrown: unknown) => thrown, diff --git a/packages/worker/src/entitlements/service.ts b/packages/worker/src/entitlements/service.ts index 5d3a5e6ae5..27e9591435 100644 --- a/packages/worker/src/entitlements/service.ts +++ b/packages/worker/src/entitlements/service.ts @@ -204,37 +204,6 @@ export async function findUserAccountByStableUserId( } } -/** - * Legacy D1 mirror write for a daily entitlement counter. Expand-phase - * enforcement is authoritative in UserMeter; this helper remains for tests, - * backfills, and the best-effort mirror scheduled after DO consume/refund. - */ -export async function incrementDailyEntitlementCounter(input: { - db: D1Database - userId: string - resource: EntitlementResource - amount?: number - now?: Date -}) { - const now = input.now ?? new Date() - await input.db - .prepare( - `INSERT INTO entitlement_daily_counters (user_id, resource, day, count, updated_at) - VALUES (?, ?, ?, ?, ?) - ON CONFLICT(user_id, resource, day) DO UPDATE SET - count = entitlement_daily_counters.count + excluded.count, - updated_at = excluded.updated_at`, - ) - .bind( - input.userId, - input.resource, - utcDayKey(now), - input.amount ?? 1, - now.toISOString(), - ) - .run() -} - function assertDailyEntitlementResource( resource: EntitlementResource, ): DailyEntitlementResource { @@ -246,119 +215,33 @@ function assertDailyEntitlementResource( return resource } -// Best-effort mirror scheduling: prefer `waitUntil` when available; otherwise -// catch so the promise cannot surface as unhandled. Never awaited by callers. -function scheduleDailyEntitlementMirror( - work: Promise, - waitUntil?: (promise: Promise) => void, -) { - const tracked = work.catch((error: unknown) => { - console.warn('entitlement-daily-mirror-failed', error) - }) - if (waitUntil) { - waitUntil(tracked) - return - } - void tracked -} - -// Absolute D1 mirror ordered by DO-minted `mirrorUpdatedAt` (revision-primary) -// so a late older write cannot overwrite newer state, including refunds. -async function mirrorDailyEntitlementAbsoluteCount(input: { - db: D1Database - userId: string - resource: DailyEntitlementResource - day: string - count: number - mirrorUpdatedAt: string -}) { - await input.db - .prepare( - `INSERT INTO entitlement_daily_counters (user_id, resource, day, count, updated_at) - VALUES (?, ?, ?, ?, ?) - ON CONFLICT(user_id, resource, day) DO UPDATE SET - count = excluded.count, - updated_at = excluded.updated_at - WHERE entitlement_daily_counters.updated_at < excluded.updated_at`, - ) - .bind( - input.userId, - input.resource, - input.day, - input.count, - input.mirrorUpdatedAt, - ) - .run() -} - -/** Best-effort absolute D1 mirror of UserMeter daily counter state. */ -export function scheduleAbsoluteDailyEntitlementMirror(input: { - db: D1Database - userId: string - resource: EntitlementResource - day: string - count: number - mirrorUpdatedAt: string - waitUntil?: (promise: Promise) => void -}): void { - const resource = assertDailyEntitlementResource(input.resource) - scheduleDailyEntitlementMirror( - mirrorDailyEntitlementAbsoluteCount({ - db: input.db, - userId: input.userId, - resource, - day: input.day, - count: input.count, - mirrorUpdatedAt: input.mirrorUpdatedAt, - }), - input.waitUntil, - ) -} - -async function ensureUserMeterCounterInitialized(input: { - db: D1Database +/** + * Seed a missing UserMeter `(resource, day)` at zero. `INSERT OR IGNORE` + * inside the DO keeps concurrent cold callers safe. + */ +async function ensureUserMeterCounterInitializedAtZero(input: { env: UserMeterEnv userId: string resource: DailyEntitlementResource day: string updatedAt: string - now: Date }) { - const baseline = await readDailyEntitlementCounter({ - db: input.db, - userId: input.userId, - resource: input.resource, - now: input.now, - }) const meter = userMeterRpc({ env: input.env, userId: input.userId }) await meter.initialize({ resource: input.resource, day: input.day, - count: baseline, + count: 0, updatedAt: input.updatedAt, }) } -async function readDailyEntitlementCounter(input: { - db: D1Database - userId: string - resource: EntitlementResource - now: Date -}) { - const row = await input.db - .prepare( - `SELECT count FROM entitlement_daily_counters - WHERE user_id = ? AND resource = ? AND day = ?`, - ) - .bind(input.userId, input.resource, utcDayKey(input.now)) - .first<{ count: number }>() - return Number(row?.count ?? 0) -} - /** * Point-read one daily entitlement counter from UserMeter. Cold meters - * bootstrap once from the legacy D1 point row, then re-read; warm meters - * return the DO count without touching D1. + * initialize the `(resource, day)` at zero, then re-read; warm meters return + * the DO count. Never touches D1 daily counter state. + * + * `db` remains on the signature for call-site stability; plan/account lookups + * on other paths still need APP_DB. */ export async function readDailyEntitlementResourceUsage(input: { db: D1Database @@ -367,6 +250,7 @@ export async function readDailyEntitlementResourceUsage(input: { resource: EntitlementResource now?: Date }): Promise { + void input.db const resource = assertDailyEntitlementResource(input.resource) const now = input.now ?? new Date() const day = utcDayKey(now) @@ -378,14 +262,12 @@ export async function readDailyEntitlementResourceUsage(input: { now: updatedAt, }) if (result.outcome === 'needs_bootstrap') { - await ensureUserMeterCounterInitialized({ - db: input.db, + await ensureUserMeterCounterInitializedAtZero({ env: input.env, userId: input.userId, resource, day, updatedAt, - now, }) result = await meter.read({ resource, @@ -914,12 +796,12 @@ export async function readEntitlementResourceUsage(input: { case 'email_receives_per_day': case 'execute_calls_per_day': case 'outbound_fetches_per_day': - return await readDailyEntitlementCounter({ - db, - userId, - resource, - now, - }) + // Authoritative daily counters live in UserMeter. Callers must use + // consumeDailyEntitlement / readDailyEntitlementResourceUsage / + // readCurrentEntitlementResourceUsage — the retired D1 mirror is gone. + throw new Error( + `${resource} usage must be read from UserMeter (use readDailyEntitlementResourceUsage or readCurrentEntitlementResourceUsage).`, + ) case 'stored_email_messages': return await countRows( db, @@ -1161,25 +1043,29 @@ export async function assertWithinEntitlement( export type ConsumeDailyEntitlementInput = { db: D1Database - /** Must expose `USER_METER` (authoritative); D1 is a best-effort mirror. */ + /** Must expose `USER_METER` (sole daily counter authority). */ env: UserMeterEnv userId: string email: string | null | undefined resource: EntitlementResource now?: Date - /** Prefer `ctx.waitUntil` for the legacy D1 mirror; never awaited here. */ + /** + * Retained for call-site stability. Daily counters no longer schedule D1 + * mirror work; unused after mirror retirement. + */ waitUntil?: (promise: Promise) => void } /** * Atomically consume one daily entitlement unit via UserMeter, throwing * EntitlementLimitError when the plan limit would be exceeded. Cold keys - * bootstrap once from legacy D1; warm path awaits only the DO RPC. Mirror - * writes are best-effort and cannot affect enforcement. + * initialize at zero (concurrent-safe); warm path awaits only the DO RPC. + * Never touches the retired D1 daily counter table. */ export async function consumeDailyEntitlement( input: ConsumeDailyEntitlementInput, ): Promise { + void input.waitUntil const resource = assertDailyEntitlementResource(input.resource) const now = input.now ?? new Date() const day = utcDayKey(now) @@ -1200,14 +1086,12 @@ export async function consumeDailyEntitlement( updatedAt, }) if (result.outcome === 'needs_bootstrap') { - await ensureUserMeterCounterInitialized({ - db: input.db, + await ensureUserMeterCounterInitializedAtZero({ env: input.env, userId: input.userId, resource, day, updatedAt, - now, }) result = await meter.consume({ resource, @@ -1230,17 +1114,6 @@ export async function consumeDailyEntitlement( upgradeHint: buildEntitlementUpgradeHint(resource), }) } - scheduleDailyEntitlementMirror( - mirrorDailyEntitlementAbsoluteCount({ - db: input.db, - userId: input.userId, - resource, - day, - count: result.count, - mirrorUpdatedAt: result.mirrorUpdatedAt, - }), - input.waitUntil, - ) } export type RefundDailyEntitlementInput = { @@ -1249,38 +1122,31 @@ export type RefundDailyEntitlementInput = { userId: string resource: EntitlementResource now?: Date + /** + * Retained for call-site stability. Daily counters no longer schedule D1 + * mirror work; unused after mirror retirement. + */ waitUntil?: (promise: Promise) => void } /** * Atomically refund one previously consumed daily entitlement unit in * UserMeter (floors at zero). Pass the same `now` (day key) as the matching - * consume. The legacy D1 mirror is best-effort and never awaited. + * consume. Never touches the retired D1 daily counter table. */ export async function refundDailyEntitlement( input: RefundDailyEntitlementInput, ): Promise { + void input.db + void input.waitUntil const resource = assertDailyEntitlementResource(input.resource) const now = input.now ?? new Date() const day = utcDayKey(now) const updatedAt = now.toISOString() const meter = userMeterRpc({ env: input.env, userId: input.userId }) - const result = await meter.refund({ + await meter.refund({ resource, day, updatedAt, }) - // revision 0 means the key was never initialized; nothing to mirror. - if (result.revision < 1) return - scheduleDailyEntitlementMirror( - mirrorDailyEntitlementAbsoluteCount({ - db: input.db, - userId: input.userId, - resource, - day, - count: result.count, - mirrorUpdatedAt: result.mirrorUpdatedAt, - }), - input.waitUntil, - ) } diff --git a/packages/worker/src/entitlements/test-schema.ts b/packages/worker/src/entitlements/test-schema.ts index a032c1a79a..938f01a5f3 100644 --- a/packages/worker/src/entitlements/test-schema.ts +++ b/packages/worker/src/entitlements/test-schema.ts @@ -4,8 +4,9 @@ import { ensureUsersTestSchema } from '#worker/users-test-schema.ts' /** * Non-destructive schema for entitlement primitives in workers-unit tests, * where the D1 database starts empty and each suite provisions the tables it - * needs. Mirrors migrations 0048 (entitlement_daily_counters) and 0066 (Stripe - * billing columns) on top of the shared `users` schema. + * needs. Mirrors Stripe billing columns (0066) and package-service state on + * top of the shared `users` schema. Code no longer reads or writes the D1 + * `entitlement_daily_counters` mirror; daily counters live only in UserMeter. * * Also provisions `user_storage_buckets` because entitlement suites that touch * StorageRunner writes register ownership through that table. @@ -20,18 +21,6 @@ export async function ensureEntitlementTestSchema(db: D1Database) { 'stripe_plan_refreshed_at', ], }) - await db - .prepare( - `CREATE TABLE IF NOT EXISTS entitlement_daily_counters ( - user_id TEXT NOT NULL, - resource TEXT NOT NULL, - day TEXT NOT NULL, - count INTEGER NOT NULL DEFAULT 0, - updated_at TEXT NOT NULL, - PRIMARY KEY (user_id, resource, day) -)`, - ) - .run() await db .prepare( `CREATE TABLE IF NOT EXISTS package_service_states ( diff --git a/packages/worker/src/entitlements/user-meter-do.ts b/packages/worker/src/entitlements/user-meter-do.ts index 4a3eff0eac..21317b1a19 100644 --- a/packages/worker/src/entitlements/user-meter-do.ts +++ b/packages/worker/src/entitlements/user-meter-do.ts @@ -829,7 +829,7 @@ class UserMeterBase extends DurableObject { } } - /** Cold-key seed from legacy D1; INSERT OR IGNORE is concurrency-safe. */ + /** Cold-key seed (zero after D1 mirror retirement); INSERT OR IGNORE is concurrency-safe. */ async initialize(input: { resource: string day: string diff --git a/packages/worker/src/entitlements/user-meter.workers.test.ts b/packages/worker/src/entitlements/user-meter.workers.test.ts index 823adf7f6d..ac79978f38 100644 --- a/packages/worker/src/entitlements/user-meter.workers.test.ts +++ b/packages/worker/src/entitlements/user-meter.workers.test.ts @@ -11,6 +11,7 @@ import { planLimits } from './plans.ts' import { assertWithinStorageBytesEntitlement, consumeDailyEntitlement, + readDailyEntitlementResourceUsage, readUserD1StorageBytes, refundDailyEntitlement, } from './service.ts' @@ -36,22 +37,6 @@ async function seedFreeUser(emailPrefix: string) { return { email, userId } } -async function readD1DailyCount(input: { - userId: string - resource: string - day: string -}) { - const row = await env.APP_DB.prepare( - `SELECT count, updated_at FROM entitlement_daily_counters - WHERE user_id = ? AND resource = ? AND day = ?`, - ) - .bind(input.userId, input.resource, input.day) - .first<{ count: number; updated_at: string }>() - return row - ? { count: Number(row.count), updatedAt: String(row.updated_at) } - : null -} - async function waitFor( predicate: () => boolean, timeoutMs = 5_000, @@ -66,60 +51,24 @@ async function waitFor( } } -const mirrorSql = `INSERT INTO entitlement_daily_counters (user_id, resource, day, count, updated_at) - VALUES (?, ?, ?, ?, ?) - ON CONFLICT(user_id, resource, day) DO UPDATE SET - count = excluded.count, - updated_at = excluded.updated_at - WHERE entitlement_daily_counters.updated_at < excluded.updated_at` - -test('UserMeter bootstraps from legacy D1 once, then warms without APP_DB awaits', async () => { +test('cold daily consume initializes at zero without D1 prepare/run and first unit is 1', async () => { const now = new Date('2026-07-31T15:00:00.000Z') const day = utcDayKey(now) - const sendLimit = planLimits.free.maxEmailSendsPerDay - const user = await seedFreeUser('meter-bootstrap') - const drain = createWaitUntilDrain() + const user = await seedFreeUser('meter-cold-zero') const meter = userMeterRpc({ env, userId: user.userId }) - const legacyCount = sendLimit - 1 - await env.APP_DB.prepare( - `INSERT INTO entitlement_daily_counters (user_id, resource, day, count, updated_at) - VALUES (?, ?, ?, ?, ?)`, - ) - .bind( - user.userId, - 'email_sends_per_day', - day, - legacyCount, - '2026-07-31T12:00:00.000Z', - ) - .run() - - let d1PointReads = 0 + let dailyPrepareCalls = 0 using _patch = withPatchedDbPrepare(env.APP_DB, (originalPrepare) => { return ((query: string) => { - const statement = originalPrepare(query) - const isPointRead = - query.includes('SELECT count FROM entitlement_daily_counters') && - query.includes('WHERE user_id = ? AND resource = ? AND day = ?') - const originalBind = statement.bind.bind(statement) - return { - bind(...params: Array) { - const bound = originalBind(...params) - return { - first: async () => { - if (isPointRead) d1PointReads += 1 - return await bound.first() - }, - all: bound.all.bind(bound), - raw: bound.raw?.bind(bound), - run: bound.run.bind(bound), - } - }, - } + if (query.includes('entitlement_daily_counters')) dailyPrepareCalls += 1 + return originalPrepare(query) }) as D1Database['prepare'] }) + expect(await meter.read({ resource: 'email_sends_per_day', day })).toEqual({ + outcome: 'needs_bootstrap', + }) + await consumeDailyEntitlement({ db: env.APP_DB, env, @@ -127,285 +76,127 @@ test('UserMeter bootstraps from legacy D1 once, then warms without APP_DB awaits email: user.email, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }) - expect(d1PointReads).toBe(1) expect( await meter.read({ resource: 'email_sends_per_day', day }), ).toMatchObject({ outcome: 'ready', - count: sendLimit, + count: 1, }) + expect(dailyPrepareCalls).toBe(0) +}, 30_000) - const denied = await consumeDailyEntitlement({ - db: env.APP_DB, - env, - userId: user.userId, - email: user.email, +test('warm daily consume/read never prepares entitlement_daily_counters', async () => { + const now = new Date('2026-07-31T15:00:00.000Z') + const day = utcDayKey(now) + const user = await seedFreeUser('meter-warm-no-d1') + const meter = userMeterRpc({ env, userId: user.userId }) + await meter.initialize({ resource: 'email_sends_per_day', - now, - waitUntil: drain.waitUntil, - }).then( - () => null, - (thrown: unknown) => thrown, - ) - expect(denied).toBeInstanceOf(EntitlementLimitError) - expect(denied).toMatchObject({ - details: { - resource: 'email_sends_per_day', - plan: 'free', - limit: sendLimit, - current: sendLimit, - }, + day, + count: 0, + updatedAt: now.toISOString(), }) - expect(d1PointReads).toBe(1) - await refundDailyEntitlement({ + let dailyPrepareCalls = 0 + using _patch = withPatchedDbPrepare(env.APP_DB, (originalPrepare) => { + return ((query: string) => { + if (query.includes('entitlement_daily_counters')) dailyPrepareCalls += 1 + return originalPrepare(query) + }) as D1Database['prepare'] + }) + + await consumeDailyEntitlement({ db: env.APP_DB, env, userId: user.userId, + email: user.email, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, - }) - await drain.drain() - d1PointReads = 0 - let nonMirrorAppDbRuns = 0 - let releaseMirror!: () => void - const mirrorGate = new Promise((resolve) => { - releaseMirror = resolve - }) - let mirrorRunStarted = false - using _warmPatch = withPatchedDbPrepare(env.APP_DB, (originalPrepare) => { - return ((query: string) => { - const statement = originalPrepare(query) - const isMirrorWrite = - query.includes('INSERT INTO entitlement_daily_counters') && - query.includes('updated_at < excluded.updated_at') - const isPointRead = - query.includes('SELECT count FROM entitlement_daily_counters') && - query.includes('WHERE user_id = ? AND resource = ? AND day = ?') - const originalBind = statement.bind.bind(statement) - return { - bind(...params: Array) { - const bound = originalBind(...params) - return { - first: async () => { - if (isPointRead) d1PointReads += 1 - return await bound.first() - }, - all: bound.all.bind(bound), - raw: bound.raw?.bind(bound), - run: async () => { - if (isMirrorWrite) { - mirrorRunStarted = true - await mirrorGate - return await bound.run() - } - nonMirrorAppDbRuns += 1 - return await bound.run() - }, - } - }, - } - }) as D1Database['prepare'] }) - try { - await expect( - consumeDailyEntitlement({ - db: env.APP_DB, - env, - userId: user.userId, - email: user.email, - resource: 'email_sends_per_day', - now, - waitUntil: drain.waitUntil, - }), - ).resolves.toBeUndefined() - expect(mirrorRunStarted).toBe(true) - expect(d1PointReads).toBe(0) - expect(nonMirrorAppDbRuns).toBe(0) - } finally { - releaseMirror() - } - await drain.drain() + await expect( + readDailyEntitlementResourceUsage({ + db: env.APP_DB, + env, + userId: user.userId, + resource: 'email_sends_per_day', + now, + }), + ).resolves.toBe(1) + expect(dailyPrepareCalls).toBe(0) }, 30_000) -test('UserMeter mirror tokens reject out-of-order absolute D1 writes', async () => { - const now = new Date('2026-07-31T16:00:00.000Z') - const day = utcDayKey(now) - const user = await seedFreeUser('meter-mirror-order') +test('next UTC day cold consume starts at zero independently', async () => { + const dayOne = new Date('2026-07-31T15:00:00.000Z') + const dayTwo = new Date('2026-08-01T01:00:00.000Z') + const user = await seedFreeUser('meter-next-day') const meter = userMeterRpc({ env, userId: user.userId }) - await meter.initialize({ - resource: 'email_sends_per_day', - day, - count: 0, - updatedAt: now.toISOString(), - }) - const first = await meter.consume({ - resource: 'email_sends_per_day', - day, - limit: 100, - updatedAt: now.toISOString(), - }) - expect(first).toMatchObject({ - outcome: 'ready', - consumed: true, - count: 1, - revision: 2, - }) - const second = await meter.consume({ - resource: 'email_sends_per_day', - day, - limit: 100, - updatedAt: now.toISOString(), - }) - expect(second).toMatchObject({ - outcome: 'ready', - consumed: true, - count: 2, - revision: 3, - }) - const refunded = await meter.refund({ + await consumeDailyEntitlement({ + db: env.APP_DB, + env, + userId: user.userId, + email: user.email, resource: 'email_sends_per_day', - day, - updatedAt: now.toISOString(), + now: dayOne, }) - expect(refunded).toMatchObject({ - outcome: 'ready', - count: 1, - revision: 4, - }) - if (first.outcome !== 'ready') { - throw new Error('Expected ready consume results.') - } + expect( + await meter.read({ + resource: 'email_sends_per_day', + day: utcDayKey(dayOne), + }), + ).toMatchObject({ outcome: 'ready', count: 1 }) - type PendingMirror = { - count: number - mirrorUpdatedAt: string - resolve: () => void - promise: Promise - } - const pending: Array = [] - function deferMirror(count: number, mirrorUpdatedAt: string) { - let resolve!: () => void - const promise = new Promise((res) => { - resolve = res - }) - pending.push({ count, mirrorUpdatedAt, resolve, promise }) - return promise - } + expect( + await meter.read({ + resource: 'email_sends_per_day', + day: utcDayKey(dayTwo), + now: dayTwo.toISOString(), + }), + ).toEqual({ outcome: 'needs_bootstrap' }) + let dailyPrepareCalls = 0 using _patch = withPatchedDbPrepare(env.APP_DB, (originalPrepare) => { return ((query: string) => { - const statement = originalPrepare(query) - const isMirrorWrite = - query.includes('INSERT INTO entitlement_daily_counters') && - query.includes('updated_at < excluded.updated_at') - const originalBind = statement.bind.bind(statement) - return { - bind(...params: Array) { - const bound = originalBind(...params) - return { - first: bound.first.bind(bound), - all: bound.all.bind(bound), - raw: bound.raw?.bind(bound), - run: async () => { - if (!isMirrorWrite) return await bound.run() - await deferMirror(Number(params[3]), String(params[4])) - return await bound.run() - }, - } - }, - } + if (query.includes('entitlement_daily_counters')) dailyPrepareCalls += 1 + return originalPrepare(query) }) as D1Database['prepare'] }) - const older = env.APP_DB.prepare(mirrorSql) - .bind( - user.userId, - 'email_sends_per_day', - day, - first.count, - first.mirrorUpdatedAt, - ) - .run() - const newer = env.APP_DB.prepare(mirrorSql) - .bind( - user.userId, - 'email_sends_per_day', - day, - refunded.count, - refunded.mirrorUpdatedAt, - ) - .run() - - await waitFor(() => pending.length === 2, 5_000, 'mirror gates') - expect(pending.map((entry) => entry.mirrorUpdatedAt)).toEqual([ - userMeterMirrorUpdatedAtToken(2), - userMeterMirrorUpdatedAtToken(4), - ]) - const [olderGate, newerGate] = pending - newerGate!.resolve() - await newer - expect( - await readD1DailyCount({ - userId: user.userId, - resource: 'email_sends_per_day', - day, - }), - ).toEqual({ - count: 1, - updatedAt: userMeterMirrorUpdatedAtToken(4), + await consumeDailyEntitlement({ + db: env.APP_DB, + env, + userId: user.userId, + email: user.email, + resource: 'email_sends_per_day', + now: dayTwo, }) - - olderGate!.resolve() - await older expect( - await readD1DailyCount({ - userId: user.userId, + await meter.read({ resource: 'email_sends_per_day', - day, + day: utcDayKey(dayTwo), }), - ).toEqual({ - count: 1, - updatedAt: userMeterMirrorUpdatedAtToken(4), - }) - expect( - await meter.read({ resource: 'email_sends_per_day', day }), - ).toMatchObject({ outcome: 'ready', count: 1, revision: 4 }) + ).toMatchObject({ outcome: 'ready', count: 1 }) + expect(dailyPrepareCalls).toBe(0) }, 30_000) -test('UserMeter daily entitlement consume/refund/read/export/purge workflow is per-user and mirror-safe', async () => { +test('UserMeter daily entitlement consume/refund/read/export/purge workflow is per-user without D1 daily table', async () => { const now = new Date('2026-07-31T15:00:00.000Z') const day = utcDayKey(now) const sendLimit = planLimits.free.maxEmailSendsPerDay const userA = await seedFreeUser('meter-a') const userB = await seedFreeUser('meter-b') - const drain = createWaitUntilDrain() const meterA = userMeterRpc({ env, userId: userA.userId }) const meterB = userMeterRpc({ env, userId: userB.userId }) expect(userMeterDurableObjectName(userA.userId)).toBe(userA.userId) - await meterA.initialize({ - resource: 'email_sends_per_day', - day, - count: 0, - updatedAt: now.toISOString(), - }) - await expect( - meterA.consume({ - resource: 'email_sends_per_day', - day, - limit: sendLimit, - updatedAt: now.toISOString(), - }), - ).resolves.toMatchObject({ - outcome: 'ready', - consumed: true, - count: 1, + let dailyPrepareCalls = 0 + using _patch = withPatchedDbPrepare(env.APP_DB, (originalPrepare) => { + return ((query: string) => { + if (query.includes('entitlement_daily_counters')) dailyPrepareCalls += 1 + return originalPrepare(query) + }) as D1Database['prepare'] }) await consumeDailyEntitlement({ @@ -415,14 +206,12 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p email: userA.email, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }) - await drain.drain() expect( await meterA.read({ resource: 'email_sends_per_day', day }), - ).toMatchObject({ outcome: 'ready', count: 2 }) + ).toMatchObject({ outcome: 'ready', count: 1 }) - for (let index = 2; index < sendLimit; index += 1) { + for (let index = 1; index < sendLimit; index += 1) { await consumeDailyEntitlement({ db: env.APP_DB, env, @@ -430,7 +219,6 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p email: userA.email, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }) } const concurrent = await Promise.all( @@ -442,7 +230,6 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p email: userA.email, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }).then( () => null, (thrown: unknown) => thrown, @@ -454,17 +241,6 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p (result) => result instanceof EntitlementLimitError, ) expect(denials).toHaveLength(8) - for (const denial of denials) { - expect(denial).toMatchObject({ - details: { - code: 'entitlement_limit_exceeded', - resource: 'email_sends_per_day', - plan: 'free', - limit: sendLimit, - current: sendLimit, - }, - }) - } await consumeDailyEntitlement({ db: env.APP_DB, @@ -473,7 +249,6 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p email: userB.email, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }) expect( await meterB.read({ resource: 'email_sends_per_day', day }), @@ -488,7 +263,6 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p userId: userA.userId, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }) expect( await meterA.read({ resource: 'email_sends_per_day', day }), @@ -500,7 +274,6 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p userId: userA.userId, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }) } expect( @@ -514,7 +287,6 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p email: userA.email, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }) await consumeDailyEntitlement({ db: env.APP_DB, @@ -523,7 +295,6 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p email: userA.email, resource: 'execute_calls_per_day', now, - waitUntil: drain.waitUntil, }) const exportedA = await meterA.exportCounters({}) expect(exportedA.counters).toEqual( @@ -574,53 +345,20 @@ test('UserMeter daily entitlement consume/refund/read/export/purge workflow is p await meterB.read({ resource: 'email_sends_per_day', day }), ).toMatchObject({ outcome: 'ready', count: 1 }) - await env.APP_DB.prepare( - `DELETE FROM entitlement_daily_counters - WHERE user_id = ? AND resource = ? AND day = ?`, - ) - .bind(userA.userId, 'email_sends_per_day', day) - .run() - - consoleWarn.mockImplementation(() => {}) - const failingMirrorDb = { - prepare(query: string) { - if ( - query.includes('INSERT INTO entitlement_daily_counters') && - query.includes('updated_at < excluded.updated_at') - ) { - return { - bind() { - return { - async run() { - throw new Error('mirror write failed') - }, - } - }, - } - } - return env.APP_DB.prepare(query) - }, - } as unknown as D1Database - await expect( consumeDailyEntitlement({ - db: failingMirrorDb, + db: env.APP_DB, env, userId: userA.userId, email: userA.email, resource: 'email_sends_per_day', now, - waitUntil: drain.waitUntil, }), ).resolves.toBeUndefined() - await drain.drain() expect( await meterA.read({ resource: 'email_sends_per_day', day }), ).toMatchObject({ outcome: 'ready', count: 1 }) - expect(consoleWarn).toHaveBeenCalledWith( - 'entitlement-daily-mirror-failed', - expect.any(Error), - ) + expect(dailyPrepareCalls).toBe(0) const stub = env.USER_METER.get( env.USER_METER.idFromName(userMeterDurableObjectName(userA.userId)), diff --git a/packages/worker/src/mcp/capabilities/admin/admin-capabilities.node.test.ts b/packages/worker/src/mcp/capabilities/admin/admin-capabilities.node.test.ts index 7fa3dabbc9..c06cf8f4d6 100644 --- a/packages/worker/src/mcp/capabilities/admin/admin-capabilities.node.test.ts +++ b/packages/worker/src/mcp/capabilities/admin/admin-capabilities.node.test.ts @@ -285,9 +285,6 @@ function createAdminCapabilityTestDb(input: { raw_mime_key: message.raw_mime_key ?? null, } as T } - if (normalizedQuery.includes('from entitlement_daily_counters')) { - return null - } const count = countForQuery(normalizedQuery) if (count !== null) return { count } as T return null diff --git a/packages/worker/src/mcp/capabilities/admin/admin-user-meter-parity.node.test.ts b/packages/worker/src/mcp/capabilities/admin/admin-user-meter-parity.node.test.ts index ceccd65157..f14c29009c 100644 --- a/packages/worker/src/mcp/capabilities/admin/admin-user-meter-parity.node.test.ts +++ b/packages/worker/src/mcp/capabilities/admin/admin-user-meter-parity.node.test.ts @@ -138,6 +138,7 @@ test('admin_user_meter_parity returns null for missing users and omits lease sec ctx, ) expect(present.report?.stableUserId).toBe(stableUserId) + expect(present.report?.daily.mirrorRetired).toBe(false) expect(present.report?.deletion.mirrorLeaseParity).toBe(true) expect(mockModule.logAuditEvent).toHaveBeenCalledWith( expect.objectContaining({ diff --git a/packages/worker/src/mcp/capabilities/admin/admin-user-meter-parity.ts b/packages/worker/src/mcp/capabilities/admin/admin-user-meter-parity.ts index 2d40113791..d20b563b1f 100644 --- a/packages/worker/src/mcp/capabilities/admin/admin-user-meter-parity.ts +++ b/packages/worker/src/mcp/capabilities/admin/admin-user-meter-parity.ts @@ -13,7 +13,7 @@ const dailyResourceSchema = z.enum(dailyEntitlementResources) const dailyResourceParitySchema = z.object({ resource: dailyResourceSchema, - d1Count: z.number().int().nonnegative(), + d1Count: z.number().int().nonnegative().nullable(), meterCount: z.number().int().nonnegative().nullable(), needsBootstrap: z.boolean(), delta: z.number().int().nullable(), @@ -76,6 +76,11 @@ const reportSchema = z.object({ stableUserId: stableUserIdSchema, daily: z.object({ day: z.string(), + mirrorRetired: z + .boolean() + .describe( + 'True when D1 entitlement_daily_counters is absent. Daily gate then reports meter counts only; d1Count/delta are null and parity stays true (no D1 comparison).', + ), resources: z.array(dailyResourceParitySchema), mismatchCount: z.number().int().nonnegative(), }), @@ -98,7 +103,7 @@ export const adminUserMeterParityCapability = defineDomainCapability( ...adminCapabilityAccess, name: 'admin_user_meter_parity', description: - 'Read-only production verification report comparing one user D1 entitlement mirrors against UserMeter for daily counters, storage bytes, package-service liveness, and deletion write leases. Never bootstraps or writes parity state (DO constructor schema maintenance and stale daily pruning may still run); returns counts and parity only (no lease tokens/holders or email content). Admin-only.', + 'Read-only production verification report for one user: UserMeter daily counters plus D1↔UserMeter parity for storage bytes, package-service liveness, and deletion write leases. While `entitlement_daily_counters` exists, daily rows compare D1 mirror counts (`mirrorRetired: false`); after the drop migration, daily rows report meter counts only (`mirrorRetired: true`). Never bootstraps or writes parity state (DO constructor schema maintenance and stale daily pruning may still run); returns counts and parity only (no lease tokens/holders or email content). Admin-only.', keywords: [ 'admin', 'user meter', @@ -111,6 +116,7 @@ export const adminUserMeterParityCapability = defineDomainCapability( 'deletion', 'mirror', 'bootstrap', + 'retired', ], inputSchema, outputSchema, diff --git a/packages/worker/src/mcp/capabilities/email/email-usage-get.workers.test.ts b/packages/worker/src/mcp/capabilities/email/email-usage-get.workers.test.ts index 4b4b6a211b..c8a9fbe636 100644 --- a/packages/worker/src/mcp/capabilities/email/email-usage-get.workers.test.ts +++ b/packages/worker/src/mcp/capabilities/email/email-usage-get.workers.test.ts @@ -29,18 +29,13 @@ async function seedDailyCounter(input: { count: number day: string }) { - await env.APP_DB.prepare( - `INSERT INTO entitlement_daily_counters (user_id, resource, day, count, updated_at) - VALUES (?, ?, ?, ?, ?)`, - ) - .bind( - input.userId, - input.resource, - input.day, - input.count, - new Date().toISOString(), - ) - .run() + const meter = userMeterRpc({ env, userId: input.userId }) + await meter.initialize({ + resource: input.resource, + day: input.day, + count: input.count, + updatedAt: new Date().toISOString(), + }) } async function seedStoredMessages(userId: string, count: number) { @@ -180,32 +175,18 @@ test( const meterEmail = `usage-meter-${crypto.randomUUID()}@example.com` const meterUserId = await createStableUserIdFromEmail(meterEmail) await seedUsageAccount({ email: meterEmail, plan: 'pro' }) + await seedStoredMessages(meterUserId, 2) await seedDailyCounter({ userId: meterUserId, resource: 'email_sends_per_day', - count: 3, + count: 13, day, }) await seedDailyCounter({ userId: meterUserId, resource: 'email_receives_per_day', - count: 5, - day, - }) - await seedStoredMessages(meterUserId, 2) - const updatedAt = new Date().toISOString() - const meter = userMeterRpc({ env, userId: meterUserId }) - await meter.initialize({ - resource: 'email_sends_per_day', - day, - count: 13, - updatedAt, - }) - await meter.initialize({ - resource: 'email_receives_per_day', - day, count: 15, - updatedAt, + day, }) const meterResult = await emailUsageGetCapability.handler( {}, diff --git a/packages/worker/src/mcp/fetch-gateway.node.test.ts b/packages/worker/src/mcp/fetch-gateway.node.test.ts index 189235e61e..237c06683c 100644 --- a/packages/worker/src/mcp/fetch-gateway.node.test.ts +++ b/packages/worker/src/mcp/fetch-gateway.node.test.ts @@ -24,16 +24,6 @@ const env = { bind(...params: Array) { return { async run() { - if ( - normalizedQuery.includes( - 'insert into entitlement_daily_counters', - ) && - params[0] !== 'user-123' - ) { - throw new Error( - 'Entitlement counter upsert must bind the acting userId.', - ) - } return { meta: { changes: 1 } } }, async first() { @@ -622,8 +612,8 @@ test('gateway fetch records outbound_fetch usage metering', async () => { }) expect(successResponse.status).toBe(200) expect(fetchStub).toHaveBeenCalledTimes(1) - // waitUntil: entitlement D1 mirror + usage metering - expect(waitUntil).toHaveBeenCalledTimes(2) + // waitUntil: usage metering only (daily D1 mirror retired) + expect(waitUntil).toHaveBeenCalledTimes(1) expect(recordUsageSpy).toHaveBeenCalledTimes(1) expect(recordUsageSpy).toHaveBeenCalledWith(env, { userId: 'user-123', @@ -747,8 +737,8 @@ test('gateway fetch records outbound_fetch usage metering', async () => { waitUntil, }), ).rejects.toThrow('network failed with waitUntil') - // waitUntil: entitlement D1 mirror (before fetch) + usage metering (error path) - expect(waitUntil).toHaveBeenCalledTimes(2) + // waitUntil: usage metering only (daily D1 mirror retired) + expect(waitUntil).toHaveBeenCalledTimes(1) expect(recordUsageSpy).toHaveBeenCalledTimes(1) expect(recordUsageSpy.mock.calls[0]?.[1]).toMatchObject({ entityId: 'api.example.com',