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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 15 additions & 12 deletions docs/contributing/architecture/data-storage.md
Original file line number Diff line number Diff line change
Expand Up @@ -108,8 +108,8 @@ Deletion must cover these user-owned surfaces:
in `BUNDLE_ARTIFACTS_KV` are deleted before D1 projection rows are removed.
OAuth token/grant KV is owned by the OAuth provider and is handled through
provider grant revocation rather than app-level key scans.
- **Cloudflare Artifacts:** source repos referenced by `entity_sources` and
`repo_sessions` are deleted through the REST client in
- **Cloudflare Artifacts:** source repos referenced by `entity_sources` and the
per-user `RepoSessionIndex` catalog are deleted through the REST client in
`packages/worker/src/repo/artifacts.ts`.

## Account export inventory
Expand Down Expand Up @@ -257,12 +257,12 @@ should be rebuilt by reindexing after import.

Cloudflare Artifacts repo contents are not inlined in the JSON export. D1 stores
metadata/projections, while canonical package, job, and app source lives in the
Artifacts repos referenced by `entity_sources.repo_id` and
`repo_sessions.source_repo_id`. For account migration to a new Cloudflare
account, first run `account_export_manifest`, page through export sections as
needed, then separately fetch or clone every repo listed in `artifactRepos`
using Artifacts access and recreate those repos in the destination account
before importing D1 projections or republishing packages.
Artifacts repos referenced by `entity_sources.repo_id` and the
`RepoSessionIndex` catalog `source_repo_id`. For account migration to a new
Cloudflare account, first run `account_export_manifest`, page through export
sections as needed, then separately fetch or clone every repo listed in
`artifactRepos` using Artifacts access and recreate those repos in the
destination account before importing D1 projections or republishing packages.

## D1 (`APP_DB`)

Expand Down Expand Up @@ -305,6 +305,10 @@ The schema is defined by migrations in `packages/worker/migrations/`:
[Run records](./run-records.md)).
- There is no `jobs` or `archived_job_artifacts` table in `APP_DB` — those live
in the jobs worker's `JOBS_DB` (see [D1 (`JOBS_DB`)](#d1-jobs_db)).
- There is no `repo_sessions` table in `APP_DB` (dropped in migration `0013`).
The catalog lives in the per-user `RepoSessionIndex` Durable Object. D1 keeps
only the thin `repo_session_due_owners` hint and the platform-owned
`repo_session_storage_bucket_cursor`.
- `package_service_states` (`0095-package-service-states.sql`): per-service
liveness projection (`running` / `idle` / `stopped` / `error`) for discovery,
account export/deletion inventory, and disaster recovery. Upserted and
Expand Down Expand Up @@ -921,8 +925,8 @@ via `durableObjectNameFromParts`); domain helpers such as
untrimmed `userId` (like RunLog / UserMeter / Mailbox). Authority for the
per-user session catalog: rows, active counts, conversation resume, export,
and deletion inventory. Each index self-alarms; D1 keeps only the thin
`repo_session_due_owners` hint (one row per user with any session) plus
platform-owned hydrate and storage-bucket inventory cursors.
`repo_session_due_owners` hint (one row per user with any session) plus the
platform-owned storage-bucket inventory cursor.
- `RepoSession` — `repoSessionDurableObjectName(sessionId)` keyed by session id
only (not user-prefixed). Every RPC validates the catalog row's `user_id`
before touching the workspace. Account deletion enumerates the user's session
Expand Down Expand Up @@ -1019,8 +1023,7 @@ repos plus D1 `entity_sources` rows and a per-user `RepoSessionIndex` catalog.
`(user_id, entity_kind, entity_id)` to the repo identity and last published
commit.
- `RepoSessionIndex` stores mutable editing-fork catalog rows for repo session
Durable Objects. Leftover D1 `repo_sessions` rows hydrate into the index until
the Phase 2 DROP.
Durable Objects. D1 does not hold catalog rows.
- Published source snapshots and bundle artifacts are stored in
`BUNDLE_ARTIFACTS_KV` and keyed by `source_id` plus `published_commit`.

Expand Down
12 changes: 6 additions & 6 deletions docs/contributing/architecture/entitlements.md
Original file line number Diff line number Diff line change
Expand Up @@ -513,12 +513,12 @@ Rules:
- **Row-count limits** (saved packages, scheduled jobs, repo sessions, secrets,
running package services) are counted via built-in counters in `service.ts`.
Most are counted directly from their source D1 tables. **Repo sessions** count
`status = 'active'` rows. Unused (never-checkpointed) leftovers are swept
after 30 minutes idle; edited sessions after 7 days idle
(`repo_session_cleanup` lane, 100 rows per 5-minute tick). **Running package
services** are counted from the **per-user UserMeter DO**:
`countRunningPackageServices` counts `status = 'running'` rows with the 24h
staleness window from DO `source_updated_at`. D1 `package_service_states`
`status = 'active'` rows in the per-user `RepoSessionIndex` catalog. Unused
(never-checkpointed) leftovers are swept after 30 minutes idle; edited
sessions after 7 days idle (`repo_session_cleanup` lane, 100 rows per 5-minute
tick). **Running package services** are counted from the **per-user UserMeter
DO**: `countRunningPackageServices` counts `status = 'running'` rows with the
24h staleness window from DO `source_updated_at`. D1 `package_service_states`
remains only the enumeration index (discovery, export, deletion) — see
[Package service liveness](#package-service-liveness) and
[Run records](./run-records.md) (`state-vs-history`). **Concurrent workflows**
Expand Down
4 changes: 2 additions & 2 deletions docs/contributing/decisions/0002-data-placement.md
Original file line number Diff line number Diff line change
Expand Up @@ -121,8 +121,8 @@ placements are the decision, independent of any track's merge state):
defense in depth.
- **Per-user mailbox DO** for email metadata: planned second wave.
- **Per-user repo session catalog DO** (`RepoSessionIndex`) plus a thin D1
`repo_session_due_owners` reverse index: the shared `repo_sessions` table is
not the catalog authority. Workspace DOs stay keyed by session id.
`repo_session_due_owners` reverse index. Workspace DOs stay keyed by session
id. Shared D1 `repo_sessions` is dropped.

Deliberately stays in D1: users/auth, secrets/values/integrations config,
publish pointers and package projections, community tables, jobs schedule
Expand Down
2 changes: 1 addition & 1 deletion packages/shared/src/jobs/scheduled-lanes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
export const scheduledLaneNames = [
'reconcile_artifacts_pushes',
'repo_session_cleanup',
// Retired after leftover D1 DROP; kept so in-flight queue messages parse.
'repo_session_index_backfill',
'reconcile_inbound_deliveries',
'system_email_retention',
Expand Down Expand Up @@ -166,7 +167,6 @@ export function getScheduledLaneCadence(
const lanes: Array<ScheduledLaneName> = [
'reconcile_artifacts_pushes',
'repo_session_cleanup',
'repo_session_index_backfill',
'reconcile_inbound_deliveries',
'system_email_retention',
'storage_bucket_estimate_backfill',
Expand Down
5 changes: 5 additions & 0 deletions packages/worker/migrations/0013-drop-repo-sessions.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
-- Phase 2: leftover D1 catalog is empty. Catalog authority is RepoSessionIndex.
-- Indexes on repo_sessions drop with the table. The hydrate cursor is retired
-- with the backfill lane.
DROP TABLE IF EXISTS repo_sessions;
DROP TABLE IF EXISTS repo_session_index_backfill_cursor;
1 change: 0 additions & 1 deletion packages/worker/src/account/data-targets.node.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -209,7 +209,6 @@ test('operator-owned tables are explicit deletion/export exclusions', () => {
applyMigrations(db)
const expectedTables = [
'platform_oauth_apps',
'repo_session_index_backfill_cursor',
'repo_session_storage_bucket_cursor',
'system_email_attachments',
'system_email_delivery_events',
Expand Down
14 changes: 0 additions & 14 deletions packages/worker/src/account/data-targets.ts
Original file line number Diff line number Diff line change
Expand Up @@ -107,12 +107,6 @@ export const accountOperatorOwnedD1Surfaces = [
reason:
'Operator-provisioned built-in OAuth app registrations (global config like feature flags; no user data). Per-user connections and token secrets remain user-scoped and covered by their own targets.',
},
{
table: 'repo_session_index_backfill_cursor',
surface: 'repo_session_index_backfill_cursor',
reason:
'Platform-owned leftover D1 catalog hydrate cursor (one singleton row). Not user data; account deletion does not touch it.',
},
{
table: 'repo_session_storage_bucket_cursor',
surface: 'repo_session_storage_bucket_cursor',
Expand Down Expand Up @@ -258,14 +252,6 @@ export const accountUserDataTargets: ReadonlyArray<UserScopedDataTarget> = [
// service binding's purgeUser and export through listArchivedJobArtifacts /
// listJobsForUser.
{ kind: 'user_id', table: 'published_bundle_artifacts' },
{
kind: 'user_id',
table: 'repo_sessions',
includeInExport: false,
surface: 'repo_sessions',
reason:
'Leftover D1 catalog omitted from portable export; RepoSessionIndex is the export authority and hydrates leftover rows on first RPC.',
},
{
kind: 'user_id',
table: 'repo_session_due_owners',
Expand Down
11 changes: 2 additions & 9 deletions packages/worker/src/account/user-owned-surfaces.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ export type UserOwnedR2Surface = {
}

export type UserOwnedArtifactSurface = {
id: 'entity_sources' | 'repo_sessions'
id: 'entity_sources'
sourceTable: string
repoColumn: string
notes: string
Expand Down Expand Up @@ -160,7 +160,7 @@ export const accountUserOwnedDurableObjectSurfaces: ReadonlyArray<UserOwnedDurab
deletionResultKey: 'repoSessionIndexes',
export: 'include',
notes:
'Per-user repo session catalog (REPO_SESSION_INDEX binding; idFromName(userId)). Authority for session rows, active counts, conversation resume, export, and deletion inventory. Workspace bytes stay in per-session RepoSession DOs. D1 keeps only the thin repo_session_due_owners hint plus leftover repo_sessions rows until the Phase 2 DROP.',
'Per-user repo session catalog (REPO_SESSION_INDEX binding; idFromName(userId)). Authority for session rows, active counts, conversation resume, export, and deletion inventory. Workspace bytes stay in per-session RepoSession DOs. D1 keeps only the thin repo_session_due_owners hint plus the storage-bucket inventory cursor.',
},
{
id: 'mcp',
Expand Down Expand Up @@ -290,13 +290,6 @@ export const accountUserOwnedArtifactSurfaces: ReadonlyArray<UserOwnedArtifactSu
notes:
'Cloudflare Artifacts repos cleaned by cleanupAllUserArtifactRepos.',
},
{
id: 'repo_sessions',
sourceTable: 'repo_sessions',
repoColumn: 'source_repo_id',
notes:
'Cloudflare Artifacts repos cleaned by cleanupAllUserArtifactRepos. Catalog rows live in RepoSessionIndex; leftover D1 repo_sessions remain until Phase 2 DROP.',
},
] as const

const accountExportExcludedDurableObjectDisplayNames: Readonly<
Expand Down
51 changes: 49 additions & 2 deletions packages/worker/src/app/account-deletion.node.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,10 @@ import {
consoleWarn,
} from '#worker/test-support/console-spies.ts'
import { createInMemoryUserMeterEnv } from '#worker/test-support/user-meter.ts'
import {
insertRepoSession,
listRepoSessionsByUser,
} from '#worker/repo/repo-sessions.ts'
import { createInMemoryRepoSessionIndexEnv } from '#worker/test-support/repo-session-index.ts'
import { userMeterRpc } from '#worker/entitlements/user-meter-client.ts'
import { applyAllMigrations } from '#worker/test-support/apply-all-migrations.ts'
Expand Down Expand Up @@ -798,7 +802,6 @@ test('deleteUserAccount cascades user-scoped rows for the requested user', async
{ id: 'src-1', user_id: userAaa, published_commit: 'abc123' },
{ id: 'src-2', user_id: userBbb, published_commit: 'def456' },
],
repo_sessions: [{ id: 'rs-1', user_id: userAaa }],
password_resets: [
{ id: 1, user_id: 1 },
{ id: 2, user_id: 1 },
Expand Down Expand Up @@ -1215,6 +1218,25 @@ test('deleteUserAccount cascades user-scoped rows for the requested user', async
get: () => ({ fetch: doFetchMock }),
},
})
await insertRepoSession(env, {
id: 'rs-1',
user_id: userAaa,
source_id: 'src-1',
source_repo_id: '',
session_branch: 'sessions/rs-1',
source_branch: 'main',
base_commit: 'abc123',
source_root: '/',
conversation_id: null,
status: 'active',
expires_at: null,
last_checkpoint_at: null,
last_checkpoint_commit: null,
last_check_run_id: null,
last_check_tree_hash: null,
created_at: '2026-07-05T00:00:00.000Z',
updated_at: '2026-07-05T00:00:00.000Z',
})

// password_resets.user_id is the database integer id; the deletion
// service must use the dbUserId (1) to clear the deleted user's reset
Expand Down Expand Up @@ -1262,7 +1284,7 @@ test('deleteUserAccount cascades user-scoped rows for the requested user', async
expect(rows.entity_sources).toEqual([
{ id: 'src-2', user_id: userBbb, published_commit: 'def456' },
])
expect(rows.repo_sessions).toEqual([])
await expect(listRepoSessionsByUser(env, userAaa)).resolves.toEqual([])
expect(deletedEmailBlobKeys.sort()).toEqual([
'email-raw:v1:user-aaa/em-1',
'email-raw:v1:user-aaa/em-2',
Expand Down Expand Up @@ -1923,6 +1945,31 @@ test('account deletion reports missing Durable Object / blob bindings and remain
])
})

test('deleteUserAccount fails closed when REPO_SESSION_INDEX is missing', async () => {
const { db, rows } = createTestDb({
users: [{ id: 1, email: 'a@example.com', stable_user_id: 'user-aaa' }],
mcp_memories: [{ id: 'memory-a', user_id: 'user-aaa' }],
})
const env = createSuccessfulDeletionEnv(db)
const envWithoutIndex = { ...env }
delete envWithoutIndex.REPO_SESSION_INDEX
await expect(
deleteUserAccount({
env: envWithoutIndex,
dbUserId: 1,
mcpUserId: 'user-aaa',
}),
).rejects.toBeInstanceOf(AccountDeletionInventoryError)
expect(rows.users).toEqual([
expect.objectContaining({
id: 1,
email: 'a@example.com',
deleting_at: expect.any(String),
}),
])
expect(rows.mcp_memories).toEqual([{ id: 'memory-a', user_id: 'user-aaa' }])
})

test('deleteUserAccount fails closed when preflight inventory cannot be read', async () => {
const { db, rows } = createTestDb(
{
Expand Down
29 changes: 7 additions & 22 deletions packages/worker/src/app/account-deletion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,6 @@ import {
import { mailboxRpc } from '#worker/email/mailbox-client.ts'
import { repoSessionIndexRpc } from '#worker/repo/repo-session-index-client.ts'
import { listRepoSessionsByUser } from '#worker/repo/repo-sessions.ts'
import {
deleteLeftoverD1RepoSessionsByUser,
listLeftoverD1RepoSessionIdsByUser,
} from '#worker/repo/repo-session-leftover-d1.ts'
import {
listAccountUserPackageServices,
listAccountUserStorageIds,
Expand Down Expand Up @@ -275,21 +271,14 @@ async function listUserSavedPackages(env: Env, userId: string) {
}

async function listUserRepoSessions(env: Env, userId: string) {
const fromIndex = env.REPO_SESSION_INDEX
? await listRepoSessionsByUser(env, userId)
: []
const leftoverIds = await listLeftoverD1RepoSessionIdsByUser({
db: env.APP_DB,
userId,
})
const byId = new Map<string, { id: string }>()
for (const row of fromIndex) {
byId.set(row.id, { id: row.id })
}
for (const sessionId of leftoverIds) {
byId.set(sessionId, { id: sessionId })
if (!env.REPO_SESSION_INDEX) {
throw new Error(
'REPO_SESSION_INDEX binding is required for account deletion.',
)
}
return [...byId.values()]
return (await listRepoSessionsByUser(env, userId)).map((row) => ({
id: row.id,
}))
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

async function listUserRemoteConnectors(env: Env, userId: string) {
Expand Down Expand Up @@ -731,10 +720,6 @@ async function purgeRepoSessionIndex(input: {
warnings: Array<string>
}): Promise<number> {
try {
await deleteLeftoverD1RepoSessionsByUser({
db: input.env.APP_DB,
userId: input.userId,
})
await repoSessionIndexRpc({
env: input.env,
userId: input.userId,
Expand Down
15 changes: 0 additions & 15 deletions packages/worker/src/app/retention.node.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -85,14 +85,6 @@ function createRetentionDb() {
const sqlite = new DatabaseSync(':memory:')
sqlite.exec('PRAGMA foreign_keys = ON')
sqlite.exec(`
CREATE TABLE repo_sessions (
id TEXT PRIMARY KEY,
user_id TEXT NOT NULL,
source_id TEXT NOT NULL,
status TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE TABLE mcp_memory_conversation_suppressions (
user_id TEXT NOT NULL,
conversation_id TEXT NOT NULL,
Expand Down Expand Up @@ -538,13 +530,6 @@ test('published bundle artifact retention deletes stale rows, KV blobs, and sour
daysAgo(60),
)
}
sqlite
.prepare(
`INSERT INTO repo_sessions (
id, user_id, source_id, status, created_at, updated_at
) VALUES ('session-1', 'user-1', 'source-session', 'active', ?, ?)`,
)
.run(daysAgo(1), daysAgo(1))
for (const [id, sourceId, commit, createdAt] of [
[
'artifact-delete',
Expand Down
19 changes: 0 additions & 19 deletions packages/worker/src/community/community-flow-test-schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -223,25 +223,6 @@ export async function ensureCommunityFlowSchema(db: D1Database) {
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
)`,
`CREATE TABLE IF NOT EXISTS repo_sessions (
id TEXT PRIMARY KEY,
user_id TEXT NOT NULL,
source_id TEXT NOT NULL,
source_repo_id TEXT NOT NULL,
session_branch TEXT NOT NULL,
source_branch TEXT NOT NULL,
base_commit TEXT NOT NULL,
source_root TEXT NOT NULL DEFAULT '/',
conversation_id TEXT,
status TEXT NOT NULL DEFAULT 'active',
expires_at TEXT,
last_checkpoint_at TEXT,
last_checkpoint_commit TEXT,
last_check_run_id TEXT,
last_check_tree_hash TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
)`,
]
for (const statement of statements) {
await db.prepare(statement).run()
Expand Down
Loading
Loading