diff --git a/README.md b/README.md index 41b3ec8..73c7493 100644 --- a/README.md +++ b/README.md @@ -138,9 +138,10 @@ just docker-publish-release v0.1.1 credential IDs, public keys, counters, names, and last-use timestamps. - `/security` is the local account-security page for passkey and app-password management. -- `com.atproto.space.*` permissioned-data routes are experimental and operator - gated with `ZDS_PERMISSIONED_DATA`. ZDS implements the storage and credential - substrate, while reader/group semantics are left to applications. +- `com.atproto.space.*` permissioned-data routes are an experimental prototype + and operator gated with `ZDS_PERMISSIONED_DATA`. ZDS implements enough storage + and credential substrate for local experiments, but the upstream proposal is + still moving and this surface is not a stable compatibility contract. ## references diff --git a/docs/account-takedown-runbook.md b/docs/account-takedown-runbook.md new file mode 100644 index 0000000..71a3d95 --- /dev/null +++ b/docs/account-takedown-runbook.md @@ -0,0 +1,300 @@ +# account takedown runbook + +This runbook is for operator action against accounts hosted on this ZDS +instance. It covers test-account cleanup, emergency suspension, and long-term +host takedown. It is not a global identity takedown mechanism. + +## authority boundary + +A PDS operator controls local hosting for accounts on that PDS: + +- whether this PDS reports a hosted account as active +- whether this PDS serves the account's repo, records, and blobs +- whether this PDS accepts authenticated writes for the account +- whether this PDS emits account-status events on its local firehose +- whether this PDS retains or purges local account data according to policy + +A PDS operator does not, by this action alone, erase all downstream copies, +control every relay or appview, or tombstone a PLC identity. Downstream services +are expected to respect account-status changes, but propagation is hop-by-hop +and can lag or diverge. + +References: + +- Account lifecycle guide: + +- Account hosting status spec: + +- Sync and `#account` events: + + +## status vocabulary + +Atproto distinguishes several inactive hosting states: + +- `deactivated`: user or host temporarily pauses the account. Content should not + be displayed or redistributed, but infrastructure does not need to delete + local copies. +- `suspended`: host temporarily pauses the account. +- `takendown`: host or service removes the repository long-term for policy or + terms reasons. +- `deleted`: account deletion. Downstream services should stop serving content + and may delete data according to policy. + +ZDS stores account hosting state in `accounts.account_status`. The currently +implemented operator states are: + +- `active` +- `deactivated` +- `takendown` + +`suspended` and `deleted` are reserved vocabulary from the protocol, but ZDS +does not yet expose operator workflows for them. + +Current behavior: + +- `com.atproto.server.deactivateAccount` sets `account_status: + "deactivated"`. +- `com.atproto.server.activateAccount` returns the account to `active` only + from `deactivated`. +- `com.atproto.admin.updateSubjectStatus` can apply or clear `takendown` for a + repo subject. +- applying `takendown` revokes local session tokens and OAuth token families. +- ZDS emits a `#account` event with `active: false` and the precise inactive + status. +- `com.atproto.sync.getRepoStatus` reports `active: false` and the precise + inactive status, and omits `rev` while inactive. +- public repo/blob/read XRPCs and authenticated repo writes are denied while the + account is inactive. + +## response levels + +### test account cleanup + +Use this when the operator owns the test account and wants to remove clutter. + +Default action: + +1. Verify the account DID and handle. +2. Export or snapshot enough metadata to undo mistakes. +3. Deactivate first. +4. Leave physical purge as a separate, reviewed step. + +### temporary safety pause + +Use this when an account may be compromised, abusive, or causing operational +harm, but the final policy decision is not ready. + +Default action: + +1. Capture evidence and timestamps. +2. Deactivate the account. +3. Revoke active local sessions and OAuth grants if the compromise involves + credentials. +4. Verify `getRepoStatus` reports inactive. +5. Watch relay/appview behavior and logs. +6. Decide later whether to reactivate, keep suspended, or perform long-term + takedown. + +### long-term host takedown + +Use this when the PDS will no longer host the account. + +ZDS has a first-class `takendown` state for hosted repo subjects. It stops local +hosting without deleting rows or blob bytes. Treat takedown as the normal +long-term operator action; treat physical purge as a separate retention-policy +decision. + +Behavior: + +- durable status distinct from user deactivation +- `#account` event with `active: false` and `status: "takendown"` +- `getRepoStatus` returns `active: false` and `status: "takendown"` +- writes, repo exports, record reads, and blob reads are denied +- sessions and OAuth grants are revoked + +Remaining gap: + +- ZDS should add a richer admin audit row for operator status changes. Until + then, preserve external notes with actor, reason, timestamps, and affected + DID. + +## investigation checklist + +Before changing account state: + +```sh +fly ssh console -a zds-pds -C "sqlite3 /data/zds.sqlite3 \ + \"select handle,did,account_status,datetime(created_at,'unixepoch'),\ + datetime(activated_at,'unixepoch'),email \ + from accounts order by created_at;\"" +``` + +For a candidate DID, inspect footprint before deciding: + +```sh +did='did:plc:example' +fly ssh console -a zds-pds -C "sqlite3 /data/zds.sqlite3 \" +select 'records', count(*) from records where did='$did' +union all select 'repo_blocks', count(*) from repo_blocks where did='$did' +union all select 'commits', count(*) from commits where did='$did' +union all select 'blobs', count(*) from blobs where did='$did' +union all select 'spaces_owned', count(*) from permissioned_spaces where owner_did='$did' +union all select 'space_actor_state', count(*) from permissioned_space_actor_state where actor_did='$did' +union all select 'space_records', count(*) from permissioned_space_records where repo_did='$did' +union all select 'space_repos', count(*) from permissioned_space_repos where repo_did='$did' +union all select 'session_tokens', count(*) from session_tokens where did='$did' +union all select 'oauth_tokens', count(*) from oauth_tokens where did='$did'; +\"" +``` + +Record: + +- handle and DID +- who requested action +- why action is justified +- whether the account is operator-owned test data or another user's account +- current record/blob/space footprint +- whether a backup is needed before action +- planned status and retention policy + +## current ZDS actions + +For user-initiated deactivation, use the existing XRPC with the user's own +bearer token. + +For operator-driven takedown, use the admin XRPC with `ZDS_ADMIN_TOKEN`: + +```sh +curl -fsS -X POST "https://pds.zat.dev/xrpc/com.atproto.admin.updateSubjectStatus" \ + -H "authorization: Bearer $ZDS_ADMIN_TOKEN" \ + -H "content-type: application/json" \ + --data-binary @- < + +## working posture + +ZDS is a sandcastle for this surface. It is acceptable to make subtractive +changes while there are no external production users, but avoid deepening older +experimental shapes when the proposal has already moved away from them. + +Prefer small alignment work that preserves optionality: + +- keep the API reference and operator docs explicit that this is experimental +- isolate prototype code under the existing `space` and permissioned-data + modules +- add tests around durable semantics that are likely to survive +- avoid adding new member-list or grant behavior in the old vocabulary + +ZDS's immediate goal is modest: get ahead of the likely direction enough to +store private data on the user's PDS, while keeping the implementation easy to +reshape. Do not write code or docs as if the draft is settled. + +## recent design pressure + +On 2026-06-24, Richard Barnes pushed back publicly on the proposal's broad +scope. His useful objection is that "permissioned data" appears to combine +several product shapes that can have different needs: + +- personal data: bookmarks, mutes, drafts +- gated content: paid newsletters, subscriber-only posts +- socially shared data: private posts and stories +- groups: private forums, communities, group chats + +He also questioned the cost of creating a parallel data universe and called out +URI-scheme complexity around `ats://`. This is worth tracking even if the +thread did not fully engage Daniel's earlier long-form diaries. Richard has +done foundational protocol/security work, so treat the critique as meaningful +design pressure rather than as implementation guidance. + +Eli Mallon pointed back to the diary motivation: for permissioned data, +"rebroadcastability" is an antipattern. That is the core distinction that makes +this more than ordinary public repo records with app-level filtering. If a +reader can freely rebroadcast the same bytes into the public sync graph, the +privacy boundary has already failed. + +Daniel's reply clarified the current bet: the motivating use case is groups, +and the proposal tries not to be dogmatic about one protocol per product +category. If the group-shaped primitive is expressive enough for the simpler +cases without becoming overwrought, that is a reasonable design win. He also +said he is still wrestling with URI scheme choice and is leaning back toward +keeping `at`. + +ZDS implications: + +- keep our docs careful about `ats://`; do not invest in naming churn until the + proposal lands +- keep private personal data as the motivating local use case, not a claim that + ZDS has solved groups +- do not encode a universal group model into ZDS while the proposal is still + deliberately pushing group semantics up to applications/space hosts +- preserve the ability to store private records and blobs now, but keep the + protocol surface clearly experimental + +## broad matches + +The current ZDS implementation broadly matches these proposal instincts: + +- permissioned data is separate from public repo data +- permissioned records use an `ats://` address that names space authority, + space type, space key, writer repo, collection, and record key +- one writer has one permissioned repo per space +- record sets are represented with LtHash-style set commitment state +- general reader/group semantics should not be encoded as a universal protocol + member list +- apps that receive access to a space should expect whole-space read semantics, + not just access to the authorizing user's own writer repo + +## proposal deltas to track + +These are the main known gaps between ZDS's current prototype and PR #94: + +- Space declarations are lexicon definitions, not records. That appears to make + a space type behave like a collection/schema authority: one NSID-shaped path + declares a reusable space type that many space repos, accounts, spaces, and + apps can share. This trades away the direct repo-authority feel of records in + favor of namespace authority, fallback behavior, and a single path segment. + App designers or project maintainers may publish the initial declaration + using DNS and lexicon-publishing authority, but the declaration becomes a + reusable commons resource after that. Lexicon unpublishing also should not + break existing spaces that depend on the declaration, which makes lexicon + resolution override behavior appropriate in a way that would be surprising + for normal AT records. +- `com.atproto.space.getMemberGrant` is older vocabulary. The proposal uses + `com.atproto.space.getDelegationToken` as a PDS-issued OAuth query and then + exchanges that with the space host through `getSpaceCredential`. +- `getRepoOplog` is older vocabulary. The proposal uses `listRepoOps`. +- The proposal adds `com.atproto.space.listRepos` so a space host can list known + writer repos in a space without enumerating readers. +- The proposal adds `registerNotify` as an explicit syncer registration method. + ZDS currently records credential recipients and supports write/deletion + notifications, but does not expose this method shape. +- Space credentials and delegation tokens have more specific JWT `typ`, `sub`, + `aud`, `client_id`, lifetime, and verification expectations than ZDS's + prototype helper names imply. +- The proposal introduces optional client attestation for app-bound space + credentials. +- Space-authority DID documents are expected to publish `#atproto_space` and + `#atproto_space_host` material. ZDS currently does not model this separately + from ordinary account/PDS signing material. +- Commit state is more than an LtHash value. The draft includes random commit + `ikm`, a signature over commit context, and a MAC over the record-set hash. + This is a deliberate deniability design: do not accidentally make a + permissioned commit into a durable public proof of private content. +- OAuth `space:` scopes distinguish `read` from `read_self` and separate + lifecycle/admin capability with `manage=`. +- Permission-set records use proposal vocabulary such as `spaceType`; older ZDS + permission-set handling should be audited before relying on it. +- Space deletion semantics need review. The draft distinguishes space-authority + deletion from member repo ownership and does not imply arbitrary erasure of + every writer's local data. + +## implementation traps + +These are low-level details to keep visible before attempting parity work: + +- LtHash live-record elements must stay unique. The draft shape uses + `{collection}/{rkey}/{record_cid}` so the same live element is not added more + than once. +- Byte order matters. LtHash lanes are little-endian, while the + TLS-style length-prefixed context used around signed commit material is + big-endian. +- The operation log is a sync shortcut, not permanent history. A host may + compact or drop it; syncers must be able to compare set hashes and recover by + listing records. +- Permissioned sync is relay-less and pull-based. Real-time behavior comes from + write notifications through the space authority, but correctness cannot rely + on every notification arriving. +- Space types should be particular. Avoid designing one universal community + space that gives every app access to every modality; apps can group many + narrower spaces under a community authority at the application layer. + +## export and backup pressure + +The draft does not currently look like it wants a `getSpaceRepo` CAR equivalent +for permissioned repos. That may be intentional: a permissioned repo is not +necessarily the public MST repo shape, and backfill can be expressed through +space enumeration plus `listRecords`. + +For ZDS, avoid inventing a `getSpaceRepo` endpoint unless the proposal moves +that way. The safer design pressure is exportability: + +- residents should be able to enumerate their own permissioned spaces +- residents should be able to enumerate their own writer repos and records + inside those spaces +- backup/export can be a higher-level artifact, such as a manifest plus records + grouped by space +- blobs should probably be referenced or streamed separately from record export + instead of assuming one giant blob archive is always practical +- any future export path should be resumable or easy to retry for large media + sets + +This is separate from sync. Sync methods need protocol parity; backup/export can +be an operator or resident affordance built from stable lower-level primitives. + +## simplespace surprise + +PR #94 removes the protocol-level member list as a universal authorization +substrate, but it also sketches a required PDS baseline called +`com.atproto.simplespace.*`. + +That baseline currently includes: + +- `createSpace` +- `updateSpace` +- `deleteSpace` +- `addMember` +- `removeMember` +- `listMembers` + +It also introduces simple policies such as `public`, `member-list`, and +`managing-app`, plus a managing-app access check hook. + +This is not the same as reviving the older generic member-list experiment. If +the proposal keeps `simplespace`, ZDS should treat it as a small baseline space +management API while still leaving richer application/group semantics to apps. + +Some commentary describes the boring access model as a list of DIDs. In ZDS +terms, do not read that as permission to restore the old protocol member-list +sync surface. Treat it as space-management policy, especially for +`com.atproto.simplespace.*` and managing-app credential minting. + +## useful tests before more surface + +Before adding more endpoints, prefer tests that exercise durable semantics: + +- per-space and per-writer repo isolation +- `ats://` parsing and formatting +- OAuth scope gating for owner reads, self reads, writes, and management +- ranged blob reads through permissioned-data auth +- write oplog ordering and cursor behavior +- record-set commitment changes across create, update, and delete +- deletion behavior for owner spaces versus writer repos + +These tests should stay separate from public-repo conformance tests so the +experimental surface can move without muddying stable PDS behavior. diff --git a/docs/permissioned-data.md b/docs/permissioned-data.md index 84d153c..bf23764 100644 --- a/docs/permissioned-data.md +++ b/docs/permissioned-data.md @@ -1,8 +1,10 @@ # permissioned data -Permissioned data support is experimental and operator gated. ZDS documents the -`com.atproto.space.*` surface as experimental, and only enables the handlers -when explicitly configured: +Permissioned data support is experimental and operator gated. Treat the current +ZDS surface as a prototype for local experiments, not as a claim of parity with +the evolving upstream permissioned-data proposal. ZDS documents the +`com.atproto.space.*` routes as experimental, and only enables the handlers when +explicitly configured: ```sh ZDS_PERMISSIONED_DATA=true @@ -24,21 +26,31 @@ even when an operator has not enabled it. - Reference implementation path: `packages/pds/src/actor-store/space` on the permissioned-data branch +- Bluesky proposal PR: + +- Local alignment notes: + [permissioned-data proposal 94](permissioned-data-proposal-94.md) The upstream design is still moving. The May branch included protocol-level -member lists. The later discussion is moving toward making space credentials -the protocol substrate and leaving reader/group semantics to applications or -space-host policy. ZDS follows that newer direction because this deployment has -no external production users, so subtractive experimental changes are cheaper -than carrying compatibility for a shape we no longer want. +member lists. The later discussion moved toward making space credentials the +protocol substrate and leaving reader/group semantics to applications or +space-host policy. Proposal PR #94 keeps that broad direction, but also sketches +a baseline `com.atproto.simplespace.*` management surface for PDS-managed +spaces. ZDS has not implemented that newer proposal shape yet. Feedback should respect the research behind the sketch and avoid overfitting to ZDS or plyr.fm. Name concrete implementation pressure, but frame it as input for the broader protocol shape. -## current surface +The current local motivation is private data storage on a resident's PDS. That +is narrower than proving every group/community use case. Recent proposal +discussion reinforces the same posture: groups are the hard case that may +justify the primitive, but ZDS should not bake group semantics into the PDS just +to get private records and blobs working early. -ZDS currently keeps the permissioned-data substrate: +## current prototype surface + +ZDS currently keeps a pre-proposal-94 permissioned-data substrate: - space lifecycle: `createSpace`, `getSpace`, `listSpaces`, `updateSpaceConfig`, `deleteSpace` @@ -48,7 +60,12 @@ ZDS currently keeps the permissioned-data substrate: - credentials and deletion fanout: `getMemberGrant`, `getSpaceCredential`, `notifySpaceDeleted` -ZDS deliberately does not expose protocol member-list routes: +The proposal draft has since moved some names and responsibilities. In +particular, `getMemberGrant` maps only loosely to the newer +`getDelegationToken`, and `getRepoOplog` maps only loosely to `listRepoOps`. +Keep that vocabulary difference visible when changing this area. + +ZDS deliberately does not expose the older protocol member-list routes: - `addMember` - `removeMember` @@ -64,11 +81,16 @@ ZDS deliberately does not expose protocol member-list routes: ZDS treats permissioned data as a space credential and writer-repo substrate, not as a universal group-membership system. -Applications own reader/group semantics: label rosters, supporter access, +Applications own rich reader/group semantics: label rosters, supporter access, personal private libraries, community roles, follower-only access, and similar -policy all belong above the PDS substrate. If an application needs portable -access state, it should model that state as application records in a -permissioned space rather than relying on a PDS-wide member list. +policy all belong above the low-level PDS substrate. If an application needs +portable access state, it should model that state as application records in a +permissioned space rather than relying on a universal PDS-wide member list. + +The newer proposal draft may still require a small PDS-managed +`com.atproto.simplespace.*` baseline for simple spaces. That is a compatibility +target to evaluate, not permission to rebuild the older generic member-list +experiment. For private spaces, ZDS currently mints space credentials only to the space owner DID unless the space is public. That is the conservative default until @@ -93,7 +115,8 @@ explicit space-scoped tables instead of the public repo tables: - `permissioned_space_repos`: writer repo state and current record-set hash - `permissioned_space_record_oplog`: incremental record changes by `(space, repo_did, rev, idx)` -- `permissioned_space_credentials`: short-lived grants and credentials +- `permissioned_space_credentials`: short-lived prototype grants and + credentials - `permissioned_space_credential_recipients`: services to notify for writes and space deletion @@ -117,7 +140,7 @@ Likely future Zat candidates: - `SpaceUri` parsing/formatting - LtHash and set commitment primitives -- member-grant and space-credential JWT helpers +- proposal-shaped delegation-token and space-credential JWT helpers Do not edit or patch the sibling `zat` repo from ZDS without stopping and making the proposed Zat change explicit. diff --git a/docs/references.md b/docs/references.md index d552188..18646a4 100644 --- a/docs/references.md +++ b/docs/references.md @@ -185,6 +185,9 @@ the code or in the relevant docs page. state, not merely JWT formatting concerns. - Delegation-ready code should keep subject, actor, and controller identities distinct even before user-facing account delegation exists. +- Account-status actions are local hosting decisions. For operator action, + distinguish temporary deactivation from host takedown, deletion, and physical + purge; see [account takedown runbook](account-takedown-runbook.md). ## active comparison work diff --git a/src/atproto/repo.zig b/src/atproto/repo.zig index ae539d2..df4895b 100644 --- a/src/atproto/repo.zig +++ b/src/atproto/repo.zig @@ -18,6 +18,7 @@ pub fn createRecord(request: *http_api.Request) !void { const auth_ctx = requireAccount(request, allocator) catch return; const account = auth_ctx.account; + try requireActiveAccount(request, account.did); const body = http_api.readBodyAlloc(request, allocator, max_repo_write_body_len) catch |err| switch (err) { error.BodyTooLarge => return http_api.xrpcError(request, .payload_too_large, "InvalidRequest", "record write body is too large"), else => return err, @@ -58,6 +59,7 @@ pub fn putRecord(request: *http_api.Request) !void { const auth_ctx = requireAccount(request, allocator) catch return; const account = auth_ctx.account; + try requireActiveAccount(request, account.did); const body = http_api.readBodyAlloc(request, allocator, max_repo_write_body_len) catch |err| switch (err) { error.BodyTooLarge => return http_api.xrpcError(request, .payload_too_large, "InvalidRequest", "record write body is too large"), else => return err, @@ -107,6 +109,7 @@ pub fn describeRepo(request: *http_api.Request) !void { const account = store.resolveRepo(repo) orelse { return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"); }; + try requirePublicRepoAvailable(request, account.did); const collections = try store.listCollectionsJson(allocator, account.did); const body = try std.fmt.allocPrint( allocator, @@ -135,6 +138,7 @@ pub fn getRecord(request: *http_api.Request) !void { const account = store.resolveRepo(repo) orelse { return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"); }; + try requirePublicRepoAvailable(request, account.did); var collection_buf: [256]u8 = undefined; const collection = http_api.queryParam(request.url.raw, "collection", &collection_buf) orelse { return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing collection"); @@ -169,6 +173,7 @@ pub fn listRecords(request: *http_api.Request) !void { const account = store.resolveRepo(repo) orelse { return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"); }; + try requirePublicRepoAvailable(request, account.did); var collection_buf: [256]u8 = undefined; const collection = http_api.queryParam(request.url.raw, "collection", &collection_buf) orelse { return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing collection"); @@ -187,6 +192,7 @@ pub fn deleteRecord(request: *http_api.Request) !void { const auth_ctx = requireAccount(request, allocator) catch return; const account = auth_ctx.account; + try requireActiveAccount(request, account.did); const body = http_api.readBodyAlloc(request, allocator, max_repo_write_body_len) catch |err| switch (err) { error.BodyTooLarge => return http_api.xrpcError(request, .payload_too_large, "InvalidRequest", "record write body is too large"), else => return err, @@ -223,6 +229,7 @@ pub fn applyWrites(request: *http_api.Request) !void { const auth_ctx = requireAccount(request, allocator) catch return; const account = auth_ctx.account; + try requireActiveAccount(request, account.did); const body = http_api.readBodyAlloc(request, allocator, max_apply_writes_body_len) catch |err| switch (err) { error.BodyTooLarge => return http_api.xrpcError(request, .payload_too_large, "InvalidRequest", "applyWrites body is too large"), else => return err, @@ -340,6 +347,7 @@ pub fn importRepo(request: *http_api.Request) !void { const auth_ctx = requireAccount(request, allocator) catch return; const account = auth_ctx.account; + try requireActiveAccount(request, account.did); const body = http_api.readBodyAlloc(request, allocator, 128 * 1024 * 1024) catch |err| switch (err) { error.BodyTooLarge => return http_api.xrpcError(request, .payload_too_large, "InvalidRequest", "repo CAR is too large"), else => return err, @@ -414,6 +422,7 @@ pub fn uploadBlob(io: std.Io, request: *http_api.Request) !void { const auth_ctx = requireBlobUploadAuth(request, allocator) catch return; const account = auth_ctx.account; + try requireActiveAccount(request, account.did); const mime_type = try allocator.dupe(u8, http_api.headerValue(request, "content-type") orelse "application/octet-stream"); if (auth_ctx.oauth_scope) |scope| { try requireBlobScope(request, scope, mime_type); @@ -553,6 +562,7 @@ pub fn listMissingBlobs(request: *http_api.Request) !void { const auth_ctx = requireAccount(request, allocator) catch return; const account = auth_ctx.account; + try requireActiveAccount(request, account.did); var cursor_buf: [256]u8 = undefined; const body = try store.writeMissingBlobsJson( allocator, @@ -573,6 +583,34 @@ fn requireAccount(request: *http_api.Request, allocator: std.mem.Allocator) !htt }; } +fn requireActiveAccount(request: *http_api.Request, did: []const u8) !void { + const status = store.accountStatus(did) catch |err| switch (err) { + error.RepoNotFound => return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"), + else => return err, + }; + if (status.isActive()) return; + return accountStatusError(request, status); +} + +fn requirePublicRepoAvailable(request: *http_api.Request, did: []const u8) !void { + const status = store.accountStatus(did) catch |err| switch (err) { + error.RepoNotFound => return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"), + else => return err, + }; + if (status.isActive()) return; + return accountStatusError(request, status); +} + +fn accountStatusError(request: *http_api.Request, status: store.AccountStatus) !void { + return switch (status) { + .active => {}, + .takendown => http_api.xrpcError(request, .bad_request, "RepoTakendown", "Repo has been taken down"), + .suspended => http_api.xrpcError(request, .bad_request, "RepoSuspended", "Repo is suspended"), + .deactivated => http_api.xrpcError(request, .bad_request, "RepoDeactivated", "Repo is deactivated"), + .deleted => http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"), + }; +} + fn requireRepoScope(request: *http_api.Request, maybe_scope: ?[]const u8, action: scopes.RepoAction, collection: []const u8) !void { if (scopes.repoAllows(maybe_scope, action, collection)) return; return http_api.xrpcError(request, .forbidden, "InsufficientScope", "Insufficient scope"); diff --git a/src/atproto/server.zig b/src/atproto/server.zig index 2a6760c..eedf2bf 100644 --- a/src/atproto/server.zig +++ b/src/atproto/server.zig @@ -753,6 +753,13 @@ pub fn createSession(request: *http_api.Request) !void { log.debug("xrpc createSession rejected account_not_found identifier={s}\n", .{identifier}); return http_api.xrpcError(request, .unauthorized, "AuthenticationRequired", "Invalid identifier or password"); }; + const account_status = try store.accountStatus(account.did); + switch (account_status) { + .active, .deactivated => {}, + .takendown => return http_api.xrpcError(request, .unauthorized, "AccountTakedown", "Account has been taken down"), + .suspended => return http_api.xrpcError(request, .unauthorized, "AccountSuspended", "Account is suspended"), + .deleted => return http_api.xrpcError(request, .unauthorized, "AuthenticationRequired", "Invalid identifier or password"), + } const primary_password_matches = auth.passwordMatches(account, password); const app_password = if (primary_password_matches) null @@ -772,9 +779,8 @@ pub fn createSession(request: *http_api.Request) !void { const info = store.getEmailInfo(arena.allocator(), account.did) orelse { return http_api.xrpcError(request, .not_found, "AccountNotFound", "Account not found"); }; - const active = store.isAccountActive(account.did); - - const body_out = try sessionJsonWithInfo(arena.allocator(), account, info.email, info.email_confirmed, active, issued.access.token, issued.refresh.token); + const status = try store.accountStatus(account.did); + const body_out = try sessionJsonWithInfo(arena.allocator(), account, info.email, info.email_confirmed, status, issued.access.token, issued.refresh.token); try http_api.json(request, .ok, body_out); } @@ -804,8 +810,8 @@ pub fn refreshSession(request: *http_api.Request) !void { const info = store.getEmailInfo(allocator, account.did) orelse { return http_api.xrpcError(request, .not_found, "AccountNotFound", "Account not found"); }; - const active = store.isAccountActive(account.did); - const body_out = try sessionJsonWithInfo(allocator, account, info.email, info.email_confirmed, active, issued.access.token, issued.refresh.token); + const status = try store.accountStatus(account.did); + const body_out = try sessionJsonWithInfo(allocator, account, info.email, info.email_confirmed, status, issued.access.token, issued.refresh.token); try http_api.json(request, .ok, body_out); } @@ -822,8 +828,8 @@ pub fn getSession(request: *http_api.Request) !void { const info = store.getEmailInfo(allocator, account.did) orelse { return http_api.xrpcError(request, .not_found, "AccountNotFound", "Account not found"); }; - const active = store.isAccountActive(account.did); - const body_out = try sessionJsonWithInfo(allocator, account, info.email, info.email_confirmed, active, null, null); + const status = try store.accountStatus(account.did); + const body_out = try sessionJsonWithInfo(allocator, account, info.email, info.email_confirmed, status, null, null); return http_api.json(request, .ok, body_out); } @@ -834,7 +840,7 @@ fn sessionJson( refresh: []const u8, ) ![]const u8 { const info = store.getEmailInfo(allocator, account.did) orelse return error.AccountNotFound; - return sessionJsonWithInfo(allocator, account, info.email, info.email_confirmed, store.isAccountActive(account.did), access, refresh); + return sessionJsonWithInfo(allocator, account, info.email, info.email_confirmed, try store.accountStatus(account.did), access, refresh); } fn issueSessionTokens(allocator: std.mem.Allocator, account: auth.Account, auth_method: []const u8, app_password_name: ?[]const u8) !auth.SessionPair { @@ -872,7 +878,7 @@ fn sessionJsonWithInfo( account: auth.Account, email: []const u8, email_confirmed: bool, - active: bool, + status: store.AccountStatus, access: ?[]const u8, refresh: ?[]const u8, ) ![]const u8 { @@ -898,8 +904,8 @@ fn sessionJsonWithInfo( .handle = account.handle, .email = email, .emailConfirmed = email_confirmed, - .active = active, - .status = if (active) null else "deactivated", + .active = status.isActive(), + .status = if (status.isActive()) null else status.asString(), }, .{ .emit_null_optional_fields = false }); } @@ -925,8 +931,11 @@ pub fn activateAccount(request: *http_api.Request) !void { defer arena.deinit(); const allocator = arena.allocator(); const account = requireAccount(request, allocator) catch return; - try store.setAccountActive(account.did, true); - try store.sequenceAccountEvent(allocator, account.did, true); + store.setAccountActive(account.did, true) catch |err| switch (err) { + error.InvalidAccountStatus => return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Account cannot be activated while takendown, suspended, or deleted"), + else => return err, + }; + try store.sequenceAccountEvent(allocator, account.did, .active); try store.sequenceIdentityEvent(allocator, account.did, account.handle); try store.sequenceSyncEvent(allocator, account.did); sync.notifyCrawlers(true); @@ -939,11 +948,88 @@ pub fn deactivateAccount(request: *http_api.Request) !void { const allocator = arena.allocator(); const account = requireAccount(request, allocator) catch return; try store.setAccountActive(account.did, false); - try store.sequenceAccountEvent(allocator, account.did, false); + try store.sequenceAccountEvent(allocator, account.did, .deactivated); sync.notifyCrawlers(true); return http_api.json(request, .ok, "{}"); } +pub fn updateSubjectStatus(request: *http_api.Request) !void { + var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); + defer arena.deinit(); + const allocator = arena.allocator(); + + requireAdminToken(request) catch return; + var body_buf: [8192]u8 = undefined; + const body = try http_api.readBody(request, &body_buf); + const parsed = try http_api.parseJsonBody(request, allocator, body); + const subject = switch (parsed.value) { + .object => |object| object.get("subject") orelse { + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing subject"); + }, + else => return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Expected object"), + }; + const did = switch (subject) { + .object => |object| switch (object.get("did") orelse return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing subject did")) { + .string => |value| value, + else => return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid subject did"), + }, + else => return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid subject"), + }; + if (zat.Did.parse(did) == null) { + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid subject did"); + } + + const takedown = objectValue(parsed.value, "takedown"); + const deactivated = objectValue(parsed.value, "deactivated"); + const takedown_applied = if (takedown) |value| valueBool(value, "applied") orelse false else false; + const deactivated_applied = if (deactivated) |value| valueBool(value, "applied") else null; + if (takedown_applied and deactivated_applied != null and deactivated_applied.? == false) { + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Cannot activate and takedown an account at the same time"); + } + + if (takedown) |value| { + const applied = valueBool(value, "applied") orelse { + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing takedown.applied"); + }; + if (!applied) { + try store.setAccountTakendown(did, false, null); + } + } + if (!takedown_applied) { + if (deactivated) |value| { + const applied = valueBool(value, "applied") orelse { + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing deactivated.applied"); + }; + store.setAccountActive(did, !applied) catch |err| switch (err) { + error.InvalidAccountStatus => return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Account cannot be activated while takendown, suspended, or deleted"), + else => return err, + }; + } + } + if (takedown) |value| { + const applied = valueBool(value, "applied") orelse { + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing takedown.applied"); + }; + if (applied) { + const status_ref = zat.json.getString(value, "ref"); + try store.setAccountTakendown(did, true, status_ref); + } + } + if (takedown == null and deactivated == null) { + _ = try store.accountStatus(did); + } + + const status = try store.accountStatus(did); + try store.sequenceAccountEvent(allocator, did, status); + sync.notifyCrawlers(true); + const body_out = try std.fmt.allocPrint( + allocator, + "{{\"subject\":{{\"$type\":\"com.atproto.admin.defs#repoRef\",\"did\":{f}}},\"active\":{},\"status\":{f}}}", + .{ std.json.fmt(did, .{}), status.isActive(), std.json.fmt(status.asString(), .{}) }, + ); + return http_api.json(request, .ok, body_out); +} + pub fn getServiceAuth(request: *http_api.Request) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); @@ -1227,6 +1313,23 @@ fn requireAdminToken(request: *http_api.Request) !void { } } +fn objectValue(value: std.json.Value, key: []const u8) ?std.json.Value { + return switch (value) { + .object => |object| object.get(key), + else => null, + }; +} + +fn valueBool(value: std.json.Value, key: []const u8) ?bool { + return switch (value) { + .object => |object| switch (object.get(key) orelse return null) { + .bool => |boolean| boolean, + else => null, + }, + else => null, + }; +} + fn expiresInTenMinutes() i64 { return store.nowMs() + (10 * 60 * 1000); } diff --git a/src/atproto/sync.zig b/src/atproto/sync.zig index bccceae..62ed753 100644 --- a/src/atproto/sync.zig +++ b/src/atproto/sync.zig @@ -177,6 +177,7 @@ pub fn getBlob(request: *http_api.Request) !void { const cid = http_api.queryParam(request.url.raw, "cid", &cid_buf) orelse { return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing cid"); }; + try requirePublicRepoAvailable(request, did); const blob = store.getPublicBlob(allocator, did, cid) orelse { return http_api.xrpcError(request, .not_found, "BlobNotFound", "Blob not found"); }; @@ -205,6 +206,7 @@ pub fn getRepo(request: *http_api.Request) !void { }; var since_buf: [256]u8 = undefined; const since = http_api.queryParam(request.url.raw, "since", &since_buf); + try requirePublicRepoAvailable(request, did); const body = store.writeRepoCarSince(allocator, did, since) catch { return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"); }; @@ -228,6 +230,7 @@ pub fn listBlobs(request: *http_api.Request) !void { }; var since_buf: [256]u8 = undefined; var cursor_buf: [256]u8 = undefined; + try requirePublicRepoAvailable(request, did); const body = try store.writeBlobListJson( allocator, did, @@ -247,6 +250,7 @@ pub fn getLatestCommit(request: *http_api.Request) !void { const did = http_api.queryParam(request.url.raw, "did", &did_buf) orelse { return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing did"); }; + try requirePublicRepoAvailable(request, did); const body = store.writeLatestCommitJson(allocator, did) catch |err| switch (err) { error.RepoNotFound => return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"), else => return err, @@ -316,6 +320,20 @@ pub fn requestCrawl(request: *http_api.Request) !void { return http_api.json(request, .ok, "{}"); } +fn requirePublicRepoAvailable(request: *http_api.Request, did: []const u8) !void { + const status = store.accountStatus(did) catch |err| switch (err) { + error.RepoNotFound => return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"), + else => return err, + }; + switch (status) { + .active => return, + .takendown => return http_api.xrpcError(request, .bad_request, "RepoTakendown", "Repo has been taken down"), + .suspended => return http_api.xrpcError(request, .bad_request, "RepoSuspended", "Repo is suspended"), + .deactivated => return http_api.xrpcError(request, .bad_request, "RepoDeactivated", "Repo is deactivated"), + .deleted => return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"), + } +} + pub fn notifyCrawlers(force: bool) void { const now = nowMillis(); if (!force) { diff --git a/src/http/router.zig b/src/http/router.zig index 67c38c2..4cd707c 100644 --- a/src/http/router.zig +++ b/src/http/router.zig @@ -35,6 +35,7 @@ pub const Route = enum { create_invite_code, create_invite_codes, get_account_invite_codes, + admin_update_subject_status, list_app_passwords, create_app_password, revoke_app_password, @@ -104,6 +105,8 @@ pub const Endpoint = struct { } }; +const permissioned_data_note = "Experimental prototype gated by ZDS_PERMISSIONED_DATA; shape may change with upstream permissioned-data drafts."; + pub const endpoints = [_]Endpoint{ .{ .route = .api_docs, .method = "GET", .path = "/api", .group = "zds", .auth = "public", .summary = "Interactive endpoint inventory for this ZDS instance." }, .{ .route = .api_docs, .method = "HEAD", .path = "/api", .group = "zds", .auth = "public", .summary = "Endpoint inventory probe without a response body." }, @@ -146,6 +149,7 @@ pub const endpoints = [_]Endpoint{ .{ .route = .create_invite_code, .method = "POST", .path = "/xrpc/com.atproto.server.createInviteCode", .group = "server", .auth = "admin", .summary = "Create one invite code.", .body = &.{ "useCount", "forAccount" } }, .{ .route = .create_invite_codes, .method = "POST", .path = "/xrpc/com.atproto.server.createInviteCodes", .group = "server", .auth = "admin", .summary = "Create invite codes in bulk.", .body = &.{ "codeCount", "useCount", "forAccounts" } }, .{ .route = .get_account_invite_codes, .method = "GET", .path = "/xrpc/com.atproto.server.getAccountInviteCodes", .group = "server", .auth = "bearer", .summary = "List invite codes associated with the signed-in account." }, + .{ .route = .admin_update_subject_status, .method = "POST", .path = "/xrpc/com.atproto.admin.updateSubjectStatus", .group = "admin", .auth = "admin", .summary = "Update account takedown or deactivation status.", .body = &.{ "subject", "takedown", "deactivated" } }, .{ .route = .list_app_passwords, .method = "GET", .path = "/xrpc/com.atproto.server.listAppPasswords", .group = "server", .auth = "bearer", .summary = "List app passwords for the signed-in account." }, .{ .route = .create_app_password, .method = "POST", .path = "/xrpc/com.atproto.server.createAppPassword", .group = "server", .auth = "bearer", .summary = "Create an app password.", .body = &.{ "name", "privileged" } }, .{ .route = .revoke_app_password, .method = "POST", .path = "/xrpc/com.atproto.server.revokeAppPassword", .group = "server", .auth = "bearer", .summary = "Revoke an app password.", .body = &.{"name"} }, @@ -201,27 +205,27 @@ pub const endpoints = [_]Endpoint{ .{ .route = .identity_submit_plc_operation, .method = "POST", .path = "/xrpc/com.atproto.identity.submitPlcOperation", .group = "identity", .auth = "bearer", .summary = "Submit a PLC operation." }, .{ .route = .identity_resolve_handle, .method = "GET", .path = "/xrpc/com.atproto.identity.resolveHandle", .group = "identity", .auth = "public", .summary = "Resolve a handle to a DID.", .params = &.{"handle"} }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.createSpace", .group = "space", .auth = "experimental bearer", .summary = "Create a permissioned data space.", .body = &.{ "did", "type", "skey", "managingApp", "isPublic", "appAccessMode", "appExceptions" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getSpace", .group = "space", .auth = "experimental bearer", .summary = "Read permissioned data space configuration.", .params = &.{"space"}, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listSpaces", .group = "space", .auth = "experimental bearer", .summary = "List spaces the authenticated user participates in.", .params = &.{ "did", "type", "limit", "cursor" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.updateSpaceConfig", .group = "space", .auth = "experimental bearer", .summary = "Update permissioned data space configuration.", .body = &.{ "space", "managingApp", "isPublic", "appAccessMode", "appExceptions" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.deleteSpace", .group = "space", .auth = "experimental bearer", .summary = "Tombstone a permissioned data space.", .body = &.{"space"}, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.createSpace", .group = "space", .auth = "experimental bearer", .summary = "Create a permissioned data space.", .body = &.{ "did", "type", "skey", "managingApp", "isPublic", "appAccessMode", "appExceptions" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getSpace", .group = "space", .auth = "experimental bearer", .summary = "Read permissioned data space configuration.", .params = &.{"space"}, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listSpaces", .group = "space", .auth = "experimental bearer", .summary = "List spaces the authenticated user participates in.", .params = &.{ "did", "type", "limit", "cursor" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.updateSpaceConfig", .group = "space", .auth = "experimental bearer", .summary = "Update permissioned data space configuration.", .body = &.{ "space", "managingApp", "isPublic", "appAccessMode", "appExceptions" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.deleteSpace", .group = "space", .auth = "experimental bearer", .summary = "Tombstone a permissioned data space.", .body = &.{"space"}, .notes = permissioned_data_note }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getMemberGrant", .group = "space", .auth = "experimental OAuth", .summary = "Create a member grant for exchange with a space owner.", .params = &.{"space"}, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getMemberGrant", .group = "space", .auth = "experimental OAuth", .summary = "Create a member grant for exchange with a space owner.", .params = &.{"space"}, .notes = permissioned_data_note }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.createRecord", .group = "space", .auth = "experimental bearer", .summary = "Create a record inside a permissioned data space.", .body = &.{ "space", "repo", "collection", "rkey", "validate", "record" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.putRecord", .group = "space", .auth = "experimental bearer", .summary = "Create or update a record inside a permissioned data space.", .body = &.{ "space", "repo", "collection", "rkey", "validate", "record" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.deleteRecord", .group = "space", .auth = "experimental bearer", .summary = "Delete a record inside a permissioned data space.", .body = &.{ "space", "repo", "collection", "rkey" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.applyWrites", .group = "space", .auth = "experimental bearer", .summary = "Apply a batch of writes inside a permissioned data space.", .body = &.{ "space", "repo", "validate", "writes" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRecord", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read a record from a permissioned data space.", .params = &.{ "space", "repo", "collection", "rkey" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listRecords", .group = "space", .auth = "experimental bearer or space credential", .summary = "List record keys and CIDs in a permissioned data space.", .params = &.{ "space", "repo", "collection", "limit", "cursor", "reverse" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getBlob", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read a blob referenced from a permissioned data record.", .params = &.{ "space", "repo", "cid" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRepoState", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read current record-set commitment state for a writer repo in a space.", .params = &.{ "space", "repo" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRepoOplog", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read incremental record operations for a writer repo in a space.", .params = &.{ "space", "repo", "since", "limit" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.notifyWrite", .group = "space", .auth = "experimental service", .summary = "Notify a space owner or syncing service of a permissioned data write.", .body = &.{ "space", "repo", "rev" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.createRecord", .group = "space", .auth = "experimental bearer", .summary = "Create a record inside a permissioned data space.", .body = &.{ "space", "repo", "collection", "rkey", "validate", "record" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.putRecord", .group = "space", .auth = "experimental bearer", .summary = "Create or update a record inside a permissioned data space.", .body = &.{ "space", "repo", "collection", "rkey", "validate", "record" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.deleteRecord", .group = "space", .auth = "experimental bearer", .summary = "Delete a record inside a permissioned data space.", .body = &.{ "space", "repo", "collection", "rkey" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.applyWrites", .group = "space", .auth = "experimental bearer", .summary = "Apply a batch of writes inside a permissioned data space.", .body = &.{ "space", "repo", "validate", "writes" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRecord", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read a record from a permissioned data space.", .params = &.{ "space", "repo", "collection", "rkey" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listRecords", .group = "space", .auth = "experimental bearer or space credential", .summary = "List record keys and CIDs in a permissioned data space.", .params = &.{ "space", "repo", "collection", "limit", "cursor", "reverse" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getBlob", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read a blob referenced from a permissioned data record.", .params = &.{ "space", "repo", "cid" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRepoState", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read current record-set commitment state for a writer repo in a space.", .params = &.{ "space", "repo" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRepoOplog", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read incremental record operations for a writer repo in a space.", .params = &.{ "space", "repo", "since", "limit" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.notifyWrite", .group = "space", .auth = "experimental service", .summary = "Notify a space owner or syncing service of a permissioned data write.", .body = &.{ "space", "repo", "rev" }, .notes = permissioned_data_note }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.getSpaceCredential", .group = "space", .auth = "experimental member grant", .summary = "Exchange a member grant for a space credential.", .body = &.{ "space", "notifyEndpoint" }, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.notifySpaceDeleted", .group = "space", .auth = "experimental service", .summary = "Notify a member PDS or syncing service that a space was deleted.", .body = &.{"space"}, .notes = "Experimental and gated by ZDS_PERMISSIONED_DATA." }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.getSpaceCredential", .group = "space", .auth = "experimental member grant", .summary = "Exchange a member grant for a space credential.", .body = &.{ "space", "notifyEndpoint" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.notifySpaceDeleted", .group = "space", .auth = "experimental service", .summary = "Notify a member PDS or syncing service that a space was deleted.", .body = &.{"space"}, .notes = permissioned_data_note }, }; pub fn route(method: httpz.Method, target: []const u8) Route { @@ -301,6 +305,7 @@ test "routes pds probes" { try std.testing.expectEqual(Route.refresh_session, route(.POST, "/xrpc/com.atproto.server.refreshSession")); try std.testing.expectEqual(Route.get_session, route(.GET, "/xrpc/com.atproto.server.getSession")); try std.testing.expectEqual(Route.get_service_auth, route(.GET, "/xrpc/com.atproto.server.getServiceAuth?aud=did%3Aplc%3Aservice")); + try std.testing.expectEqual(Route.admin_update_subject_status, route(.POST, "/xrpc/com.atproto.admin.updateSubjectStatus")); try std.testing.expectEqual(Route.app_preferences_get, route(.GET, "/xrpc/app.bsky.actor.getPreferences")); try std.testing.expectEqual(Route.app_preferences_put, route(.POST, "/xrpc/app.bsky.actor.putPreferences")); try std.testing.expectEqual(Route.repo_list_records, route(.GET, "/xrpc/com.atproto.repo.listRecords?repo=alice.test&collection=app.bsky.feed.post")); diff --git a/src/http/server.zig b/src/http/server.zig index 98a78af..29f5ae1 100644 --- a/src/http/server.zig +++ b/src/http/server.zig @@ -104,6 +104,7 @@ const App = struct { .create_invite_code => try atproto_server.createInviteCode(request), .create_invite_codes => try atproto_server.createInviteCodes(request), .get_account_invite_codes => try atproto_server.getAccountInviteCodes(request), + .admin_update_subject_status => try atproto_server.updateSubjectStatus(request), .list_app_passwords => try atproto_server.listAppPasswords(request), .create_app_password => try atproto_server.createAppPassword(request), .revoke_app_password => try atproto_server.revokeAppPassword(request), diff --git a/src/storage/store.zig b/src/storage/store.zig index 0fdd048..e4cfc54 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -25,9 +25,40 @@ pub const Error = error{ InvalidRepoPath, InvalidDagCbor, InvalidRefreshSession, + InvalidAccountStatus, StoreNotInitialized, }; +pub const AccountStatus = enum { + active, + takendown, + suspended, + deactivated, + deleted, + + pub fn asString(self: AccountStatus) []const u8 { + return switch (self) { + .active => "active", + .takendown => "takendown", + .suspended => "suspended", + .deactivated => "deactivated", + .deleted => "deleted", + }; + } + + pub fn parse(raw: []const u8) AccountStatus { + if (std.mem.eql(u8, raw, "takendown")) return .takendown; + if (std.mem.eql(u8, raw, "suspended")) return .suspended; + if (std.mem.eql(u8, raw, "deactivated")) return .deactivated; + if (std.mem.eql(u8, raw, "deleted")) return .deleted; + return .active; + } + + pub fn isActive(self: AccountStatus) bool { + return self == .active; + } +}; + pub const Record = struct { did: []const u8, collection: []const u8, @@ -529,7 +560,7 @@ pub fn listResidents(allocator: std.mem.Allocator, limit: usize) ![]Resident { const actual_limit = if (limit == 0) 24 else @min(limit, 100); var rows = try conn.rows( - \\SELECT a.handle, a.did, a.activated_at, a.deactivated_at, COALESCE(c.rev, ''), COUNT(r.uri) + \\SELECT a.handle, a.did, a.account_status, COALESCE(c.rev, ''), COUNT(r.uri) \\FROM accounts a \\LEFT JOIN ( \\ SELECT c1.did, c1.rev @@ -553,9 +584,9 @@ pub fn listResidents(allocator: std.mem.Allocator, limit: usize) ![]Resident { try residents.append(allocator, .{ .handle = try allocator.dupe(u8, row.text(0)), .did = try allocator.dupe(u8, row.text(1)), - .active = row.nullableInt(2) != null and row.nullableInt(3) == null, - .rev = try allocator.dupe(u8, row.text(4)), - .record_count = @intCast(row.int(5)), + .active = AccountStatus.parse(row.text(2)).isActive(), + .rev = try allocator.dupe(u8, row.text(3)), + .record_count = @intCast(row.int(4)), }); } if (rows.err) |err| return err; @@ -774,9 +805,9 @@ pub fn createAccountWithSigningKeyAndInvite( try conn.execNoArgs("BEGIN IMMEDIATE"); errdefer conn.execNoArgs("ROLLBACK") catch {}; try conn.exec( - \\INSERT INTO accounts (did, handle, email, password_hash, activated_at, signing_key_type, signing_key) - \\VALUES (?, ?, ?, ?, CASE WHEN ? THEN unixepoch() ELSE NULL END, 'secp256k1', ?) - , .{ did, handle, email, password_hash, activated, zqlite.blob(&signing_key) }); + \\INSERT INTO accounts (did, handle, email, password_hash, activated_at, account_status, signing_key_type, signing_key) + \\VALUES (?, ?, ?, ?, CASE WHEN ? THEN unixepoch() ELSE NULL END, CASE WHEN ? THEN 'active' ELSE 'deactivated' END, 'secp256k1', ?) + , .{ did, handle, email, password_hash, activated, activated, zqlite.blob(&signing_key) }); if (invite_code) |code| { try recordInviteUseLocked(code, did); } @@ -934,27 +965,63 @@ pub fn setAccountActive(did: []const u8, active: bool) !void { defer db_mutex.unlock(store_io); try requireInitialized(); if (active) { + const status = try accountStatusLocked(did); + if (status != .active and status != .deactivated) return Error.InvalidAccountStatus; try conn.exec( \\UPDATE accounts - \\SET activated_at = COALESCE(activated_at, unixepoch()), - \\ deactivated_at = NULL + \\SET account_status = 'active', + \\ activated_at = COALESCE(activated_at, unixepoch()) \\WHERE did = ? , .{did}); } else { try conn.exec( \\UPDATE accounts - \\SET deactivated_at = unixepoch() + \\SET account_status = 'deactivated', + \\ activated_at = COALESCE(activated_at, unixepoch()) + \\WHERE did = ? + , .{did}); + } +} + +pub fn setAccountTakendown(did: []const u8, applied: bool, status_ref: ?[]const u8) !void { + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + try requireInitialized(); + if (applied) { + try conn.execNoArgs("BEGIN IMMEDIATE"); + errdefer conn.execNoArgs("ROLLBACK") catch {}; + try conn.exec( + \\UPDATE accounts + \\SET account_status = 'takendown', + \\ account_status_ref = ? \\WHERE did = ? + , .{ status_ref, did }); + try revokeAccountTokensLocked(did); + try conn.execNoArgs("COMMIT"); + } else { + try conn.exec( + \\UPDATE accounts + \\SET account_status = 'active', + \\ account_status_ref = NULL, + \\ activated_at = COALESCE(activated_at, unixepoch()) + \\WHERE did = ? AND account_status = 'takendown' , .{did}); } } -pub fn sequenceAccountEvent(allocator: std.mem.Allocator, did: []const u8, active: bool) !void { +pub fn accountStatus(did: []const u8) !AccountStatus { + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + try requireInitialized(); + return accountStatusLocked(did); +} + +pub fn sequenceAccountEvent(allocator: std.mem.Allocator, did: []const u8, status: AccountStatus) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); const seq = try nextSeqLocked(); - const frame = try accountEventFrame(allocator, seq, did, active); + const frame = try accountEventFrame(allocator, seq, did, status); try insertSeqEventLocked(seq, did, "", frame); eventlog.publish(seq); } @@ -2527,15 +2594,17 @@ pub fn writeAccountStatusJson(allocator: std.mem.Allocator, did: []const u8) ![] "SELECT COUNT(*) FROM repo_blocks WHERE did = ?", did, ); - const active = try accountActiveLocked(did); + const status = try accountStatusLocked(did); + const active = status.isActive(); const root = try latestRootLocked(allocator, did); return std.fmt.allocPrint( allocator, - "{{\"activated\":{},\"validDid\":{},\"repoCommit\":{f},\"repoRev\":{f},\"repoBlocks\":{d},\"indexedRecords\":{d},\"privateStateValues\":0,\"expectedBlobs\":{d},\"importedBlobs\":{d}}}", + "{{\"activated\":{},\"validDid\":{},\"status\":{f},\"repoCommit\":{f},\"repoRev\":{f},\"repoBlocks\":{d},\"indexedRecords\":{d},\"privateStateValues\":0,\"expectedBlobs\":{d},\"importedBlobs\":{d}}}", .{ active, active, + std.json.fmt(status.asString(), .{}), std.json.fmt(root.cid, .{}), std.json.fmt(root.rev, .{}), block_count, @@ -2666,7 +2735,7 @@ pub fn writeRepoListJson(allocator: std.mem.Allocator, cursor: ?[]const u8, limi const parsed_cursor = try parseRepoListCursor(cursor); var rows = if (parsed_cursor) |after| try conn.rows( - \\SELECT a.did, c.cid, c.rev, a.activated_at, a.deactivated_at, a.created_at + \\SELECT a.did, c.cid, c.rev, a.account_status, a.created_at \\FROM accounts a \\JOIN ( \\ SELECT did, MAX(seq) AS seq @@ -2680,7 +2749,7 @@ pub fn writeRepoListJson(allocator: std.mem.Allocator, cursor: ?[]const u8, limi , .{ after.created_at, after.created_at, after.did, @as(i64, @intCast(actual_limit)) }) else try conn.rows( - \\SELECT a.did, c.cid, c.rev, a.activated_at, a.deactivated_at, a.created_at + \\SELECT a.did, c.cid, c.rev, a.account_status, a.created_at \\FROM accounts a \\JOIN ( \\ SELECT did, MAX(seq) AS seq @@ -2704,8 +2773,9 @@ pub fn writeRepoListJson(allocator: std.mem.Allocator, cursor: ?[]const u8, limi while (rows.next()) |row| { const did = row.text(0); last_did = try allocator.dupe(u8, did); - last_created_at = row.int(5); - const active = row.nullableInt(3) != null and row.nullableInt(4) == null; + last_created_at = row.int(4); + const status = AccountStatus.parse(row.text(3)); + const active = status.isActive(); try json.beginObject(); try json.objectField("did"); try json.write(did); @@ -2717,7 +2787,7 @@ pub fn writeRepoListJson(allocator: std.mem.Allocator, cursor: ?[]const u8, limi try json.write(active); if (!active) { try json.objectField("status"); - try json.write("deactivated"); + try json.write(status.asString()); } try json.endObject(); } @@ -2799,8 +2869,9 @@ pub fn writeRepoStatusJson(allocator: std.mem.Allocator, did: []const u8) ![]con db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); - const root = try latestRootLocked(allocator, did); - const active = try accountActiveLocked(did); + const status = try accountStatusLocked(did); + const active = status.isActive(); + const root = if (active) try latestRootLocked(allocator, did) else null; var out: std.Io.Writer.Allocating = .init(allocator); defer out.deinit(); var json: std.json.Stringify = .{ .writer = &out.writer }; @@ -2811,10 +2882,11 @@ pub fn writeRepoStatusJson(allocator: std.mem.Allocator, did: []const u8) ![]con try json.write(active); if (!active) { try json.objectField("status"); - try json.write("deactivated"); + try json.write(status.asString()); + } else { + try json.objectField("rev"); + try json.write(root.?.rev); } - try json.objectField("rev"); - try json.write(root.rev); try json.endObject(); return out.toOwnedSlice(); } @@ -3618,6 +3690,7 @@ fn parseSpaceParts(uri: []const u8) ?SpaceParts { fn migrate() !void { inline for (schema_statements) |sql| try conn.execNoArgs(sql); try migratePermissionedDataTables(); + const had_deactivated_at = try accountsColumnExistsLocked("deactivated_at"); conn.execNoArgs("ALTER TABLE accounts ADD COLUMN email TEXT") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN email_confirmed_at INTEGER") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN auth_code TEXT") catch {}; @@ -3625,7 +3698,12 @@ fn migrate() !void { conn.execNoArgs("ALTER TABLE accounts ADD COLUMN pending_email TEXT") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN invites_disabled INTEGER NOT NULL DEFAULT 0") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN activated_at INTEGER") catch {}; - conn.execNoArgs("ALTER TABLE accounts ADD COLUMN deactivated_at INTEGER") catch {}; + conn.execNoArgs("ALTER TABLE accounts ADD COLUMN account_status TEXT NOT NULL DEFAULT 'active'") catch {}; + conn.execNoArgs("ALTER TABLE accounts ADD COLUMN account_status_ref TEXT") catch {}; + conn.execNoArgs("UPDATE accounts SET account_status = 'deactivated' WHERE activated_at IS NULL AND account_status = 'active'") catch {}; + if (had_deactivated_at) { + conn.execNoArgs("UPDATE accounts SET account_status = 'deactivated' WHERE deactivated_at IS NOT NULL AND account_status = 'active'") catch {}; + } conn.execNoArgs("ALTER TABLE accounts ADD COLUMN signing_key_type TEXT") catch {}; conn.execNoArgs("ALTER TABLE accounts ADD COLUMN signing_key BLOB") catch {}; conn.execNoArgs("ALTER TABLE oauth_requests ADD COLUMN login_hint TEXT") catch {}; @@ -3671,6 +3749,16 @@ fn ensureMigrationTable() !void { ); } +fn accountsColumnExistsLocked(column: []const u8) !bool { + var rows = try conn.rows("PRAGMA table_info(accounts)", .{}); + defer rows.deinit(); + while (rows.next()) |row| { + if (std.mem.eql(u8, row.text(1), column)) return true; + } + if (rows.err) |err| return err; + return false; +} + fn migrationApplied(name: []const u8) !bool { const row = try conn.row("SELECT 1 FROM zds_migrations WHERE name = ?", .{name}); if (row == null) return false; @@ -4918,9 +5006,9 @@ fn importedCommitEventFrame( return commitEventFrameFromCar(allocator, seq, did, commit_cid, rev, null, null, car_bytes, &.{}); } -fn accountEventFrame(allocator: std.mem.Allocator, seq: u64, did: []const u8, active: bool) ![]const u8 { +fn accountEventFrame(allocator: std.mem.Allocator, seq: u64, did: []const u8, status: AccountStatus) ![]const u8 { const header = try eventHeader(allocator, "#account"); - const body = if (active) + const body = if (status.isActive()) try zat.cbor.encodeAlloc(allocator, .{ .map = &.{ .{ .key = "seq", .value = .{ .unsigned = seq } }, .{ .key = "did", .value = .{ .text = did } }, @@ -4932,7 +5020,7 @@ fn accountEventFrame(allocator: std.mem.Allocator, seq: u64, did: []const u8, ac .{ .key = "seq", .value = .{ .unsigned = seq } }, .{ .key = "did", .value = .{ .text = did } }, .{ .key = "active", .value = .{ .boolean = false } }, - .{ .key = "status", .value = .{ .text = "deactivated" } }, + .{ .key = "status", .value = .{ .text = status.asString() } }, .{ .key = "time", .value = .{ .text = try nowIso(allocator) } }, } }); return joinFrame(allocator, header, body); @@ -5064,14 +5152,31 @@ fn scalarCountLocked(sql: [:0]const u8, did: []const u8) !u64 { } fn accountActiveLocked(did: []const u8) !bool { + return (try accountStatusLocked(did)).isActive(); +} + +fn accountStatusLocked(did: []const u8) !AccountStatus { const row = try conn.row( - \\SELECT activated_at, deactivated_at + \\SELECT account_status \\FROM accounts \\WHERE did = ? , .{did}); - if (row == null) return false; + if (row == null) return Error.RepoNotFound; defer row.?.deinit(); - return row.?.nullableInt(0) != null and row.?.nullableInt(1) == null; + return AccountStatus.parse(row.?.text(0)); +} + +fn revokeAccountTokensLocked(did: []const u8) !void { + try conn.exec( + \\UPDATE session_tokens + \\SET revoked_at = unixepoch(), updated_at = unixepoch() + \\WHERE did = ? AND revoked_at IS NULL + , .{did}); + try conn.exec( + \\UPDATE oauth_tokens + \\SET revoked_at = unixepoch() + \\WHERE did = ? AND revoked_at IS NULL + , .{did}); } fn cidForJson(allocator: std.mem.Allocator, json: []const u8) ![]const u8 { @@ -5185,7 +5290,8 @@ const schema_statements = [_][*:0]const u8{ \\ pending_email TEXT, \\ invites_disabled INTEGER NOT NULL DEFAULT 0, \\ activated_at INTEGER, - \\ deactivated_at INTEGER, + \\ account_status TEXT NOT NULL DEFAULT 'active', + \\ account_status_ref TEXT, \\ signing_key_type TEXT, \\ signing_key BLOB, \\ password_hash TEXT NOT NULL, @@ -6347,6 +6453,55 @@ test "session tokens are durable and refresh rotation invalidates old jti" { try std.testing.expect(!inactive_sessions[0].active); } +test "account takedown sets precise sync status and revokes tokens" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const allocator = arena.allocator(); + + try init(std.Options.debug_io, ":memory:"); + defer close(); + + const account = try createAccount( + allocator, + "status.test", + "status@test.com", + "password", + "did:plc:status", + true, + ); + var record_json = try std.json.parseFromSlice(std.json.Value, allocator, "{\"text\":\"hello\"}", .{}); + defer record_json.deinit(); + _ = try create(allocator, account, "app.bsky.feed.post", "3jtest", record_json.value); + + _ = try createSessionTokenRow( + allocator, + account.did, + "status-access", + "status-refresh", + 4102444800, + 4102444800, + "password", + null, + null, + ); + try std.testing.expect(try sessionTokenIsActive(account.did, "status-access", "com.atproto.access")); + + try setAccountTakendown(account.did, true, "mod-ref-1"); + try std.testing.expectEqual(AccountStatus.takendown, try accountStatus(account.did)); + try std.testing.expect(!try sessionTokenIsActive(account.did, "status-access", "com.atproto.access")); + + const status_json = try writeRepoStatusJson(allocator, account.did); + try std.testing.expect(std.mem.indexOf(u8, status_json, "\"active\":false") != null); + try std.testing.expect(std.mem.indexOf(u8, status_json, "\"status\":\"takendown\"") != null); + try std.testing.expect(std.mem.indexOf(u8, status_json, "\"rev\"") == null); + + const frame = try accountEventFrame(allocator, 1, account.did, .takendown); + try std.testing.expect(std.mem.indexOf(u8, frame, "takendown") != null); + + try setAccountTakendown(account.did, false, null); + try std.testing.expectEqual(AccountStatus.active, try accountStatus(account.did)); +} + test "app passwords are durable matchable and revoke their sessions" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit();