From e91800a0b59bcbf7bdb1b8569035778d60129148 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Fri, 14 Aug 2026 13:51:28 -0500 Subject: [PATCH] feat(space): align credential and sync semantics --- bench/main.zig | 14 ++ build.zig | 16 +++ docs/benchmarking.md | 6 + docs/permissioned-data-proposal-94.md | 23 +-- docs/permissioned-data.md | 31 ++-- src/atproto/space.zig | 197 ++++++++++++++++++++----- src/http/router.zig | 24 ++-- src/internal/dpop.zig | 81 +++++++++-- src/internal/permissioned_data.zig | 109 +++++++++++--- src/storage/store.zig | 200 +++++++++++++++++++------- tools/smoke-permissioned.sh | 51 ++++++- tools/space_dpop.zig | 28 ++++ 12 files changed, 629 insertions(+), 151 deletions(-) create mode 100644 tools/space_dpop.zig diff --git a/bench/main.zig b/bench/main.zig index 08eeb36..6ff6f66 100644 --- a/bench/main.zig +++ b/bench/main.zig @@ -840,6 +840,7 @@ fn benchSpace(allocator: std.mem.Allocator, options: Options) !void { (try benchSpaceListRecords(allocator, state.account, space.uri, records)).print(); (try benchSpaceListRepos(allocator, space.uri)).print(); (try benchSpaceOplog(allocator, state.account, space.uri)).print(); + (try benchSpaceListBlobs(allocator, state.account, space.uri)).print(); (try benchSpaceBlob(allocator, state.account, space.uri, cid)).print(); } @@ -1015,6 +1016,19 @@ fn benchSpaceBlob(allocator: std.mem.Allocator, account: zds.auth.tokens.Account return .{ .name = "space getBlob", .ops = iterations, .bytes = iterations * "hello permissioned audio".len, .elapsed_ns = nowNs() - start }; } +fn benchSpaceListBlobs(allocator: std.mem.Allocator, account: zds.auth.tokens.Account, space: []const u8) !BenchResult { + var arena = std.heap.ArenaAllocator.init(allocator); + defer arena.deinit(); + const iterations: usize = 1000; + const start = nowNs(); + for (0..iterations) |_| { + _ = arena.reset(.retain_capacity); + const cids = try zds.storage.store.listPermissionedBlobs(arena.allocator(), space, account.did, null, null, 500); + if (cids.len == 0) return error.MissingBlob; + } + return .{ .name = "space listBlobs", .ops = iterations, .elapsed_ns = nowNs() - start }; +} + fn benchDecodeRecord(allocator: std.mem.Allocator, options: Options) !void { try benchRecordBlockCpu(allocator, options, .decode); } diff --git a/build.zig b/build.zig index d50d469..c17aa84 100644 --- a/build.zig +++ b/build.zig @@ -120,6 +120,22 @@ pub fn build(b: *std.Build) void { const plc_repair_step = b.step("plc-repair", "Repair an existing account PLC rotation key set"); plc_repair_step.dependOn(&run_plc_repair.step); + const space_dpop = b.addExecutable(.{ + .name = "zds-space-dpop", + .root_module = b.createModule(.{ + .root_source_file = b.path("tools/space_dpop.zig"), + .target = target, + .optimize = optimize, + .imports = &.{.{ .name = "zds", .module = mod }}, + }), + }); + b.installArtifact(space_dpop); + + const run_space_dpop = b.addRunArtifact(space_dpop); + if (b.args) |args| run_space_dpop.addArgs(args); + const space_dpop_step = b.step("space-dpop", "Create a permissioned-data DPoP proof for local testing"); + space_dpop_step.dependOn(&run_space_dpop.step); + const tests = b.addTest(.{ .root_module = mod }); const run_tests = b.addRunArtifact(tests); diff --git a/docs/benchmarking.md b/docs/benchmarking.md index 16a598c..43b0abe 100644 --- a/docs/benchmarking.md +++ b/docs/benchmarking.md @@ -176,6 +176,12 @@ larger blob/CAR routes: | `sync.getRepo` | 37.9k req/s, p95 below 1 ms | | `sync.getRepo` HEAD | 36.6k req/s, p95 below 1 ms | +The separate permissioned-data scenario (`just bench space`) covers space +discovery, management-policy reads, record writes and reads, writer/oplog +enumeration, full repo CAR construction, referenced-blob enumeration, and blob +serving. Keep it separate from public-repo benchmarks because the storage and +commit formats are intentionally different. + The only pre-httpz route numbers currently available are the original narrow route set, run against commit `4841aec` with the same four probes: diff --git a/docs/permissioned-data-proposal-94.md b/docs/permissioned-data-proposal-94.md index 11ea06f..088d666 100644 --- a/docs/permissioned-data-proposal-94.md +++ b/docs/permissioned-data-proposal-94.md @@ -81,9 +81,10 @@ The current ZDS implementation broadly matches these proposal instincts: - 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 +## proposal alignment points -These are the main known gaps between ZDS's current prototype and PR #94: +These are the main design points ZDS tracks against proposal 0016 and its +follow-up pull requests: - Space declarations are lexicon definitions, not records. That appears to make a space type behave like a collection/schema authority: one NSID-shaped path @@ -103,25 +104,29 @@ These are the main known gaps between ZDS's current prototype and PR #94: of the authority's writer registry. `notifyWrite` advances that registry with the writer's revision and SHA-256 digest of its LtHash state; raw LtHash state remains private to the writer's repo host. -- `registerNotify` stores expiring whole-space or repo-scoped syncer - registrations. ZDS follows the branch's 24-hour registration window. +- `registerNotify` stores an expiring whole-space registration for a service + DID, and `unregisterNotify` removes it idempotently. ZDS follows the branch's + 24-hour registration window. - Delegation tokens identify a requesting user and contain no app identity. - Space credentials identify the authority and space and likewise contain no - app identity. Optional app identity arrives separately as a verified client - attestation. + Optional app identity arrives separately as a verified client attestation. + Credential issuance and use follow PR #99's ES256 DPoP binding: the + credential carries `cnf.jkt` and an attested `client_id` when present. - ZDS currently uses the proposal's permitted `#atproto` signing-key fallback. For remote authorities it resolves `#atproto_space_host` when present and falls back to `#atproto_pds`; ZDS does not yet publish dedicated space key or host entries for its own resident authorities. - Commit state uses versioned deniable commits: random `ikm`, a signature over - `(space, author, rev, ikm)`, and a context-derived MAC over the record-set - hash. Wire bytes use AT JSON `$bytes` values. + `(space, author, rev, ikm)`, and PR #100's direct HKDF-Expand-derived HMAC + over the record-set hash. The CAR record index uses canonical DAG-CBOR map + ordering. Wire bytes use AT JSON `$bytes` values. - OAuth `space:` scopes distinguish `read` from collection-constrained `read_self`, use `authority=`, and separate lifecycle/admin capability with repeated `manage=` parameters. Permission sets use `spaceType`, may name cross-namespace collections, and resolve `authority=self` at grant issuance. - Delegation tokens and verified client attestations are short-lived, single-use inputs consumed atomically when minting a space credential. +- `listBlobs` enumerates the permissioned repo's referenced blob CIDs for full + or revision-filtered sync; `getBlob` remains the authenticated byte path. - 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. diff --git a/docs/permissioned-data.md b/docs/permissioned-data.md index 4e6a311..ab395fc 100644 --- a/docs/permissioned-data.md +++ b/docs/permissioned-data.md @@ -58,8 +58,8 @@ ZDS splits protocol data routes from baseline PDS management routes: - `com.atproto.space.*`: `getSpace`, `listSpaces`, `listRepos`, `getDelegationToken`, `getSpaceCredential`, records, blobs, signed commits, - full and incremental repo sync, notification registration, write - notifications, and deletion notifications + full and incremental repo sync, blob enumeration, notification registration + and removal, write notifications, and deletion notifications - `com.atproto.simplespace.*`: `createSpace`, `updateSpace`, `deleteSpace`, `addMember`, `removeMember`, `listMembers`, and the managing-app `checkUserAccess` hook @@ -110,7 +110,10 @@ policy services replace it with their own application-layer decision. `{"$type":"com.atproto.simplespace.defs#allowList","allowed":[...]}`. Delegation tokens identify the user, not the app. For allow-list spaces, ZDS verifies the separately supplied `clientAttestation` against the client's -published metadata and JWKS before using its `client_id`. +published metadata and JWKS before using its `client_id`. Credential issuance +also requires an ES256 DPoP proof. The resulting credential records the proof +key thumbprint in `cnf.jkt` and, when attested, the `client_id`; every use of the +credential requires a fresh DPoP proof bound to the request and credential. Writer notifications are keyed by writer repo. `notifyWrite` verifies service auth from the writer repo to the space DID, records the writer's latest revision @@ -119,10 +122,11 @@ summary out to registered syncers. `listRepos` is space-credential-only and returns this authority-maintained writer set; it never exposes the raw LtHash state or a reader/member list. -`registerNotify` stores a 24-hour subscription. Omitting `repo` registers a -syncer with the space host for the whole space; naming `repo` registers it with -that repo's host only. Notifications remain best effort, so syncers compare -revisions and commit hashes and recover through direct pulls. +`registerNotify` stores a 24-hour whole-space subscription for a service DID, +resolving either the named service fragment or the DID's PDS endpoint. +`unregisterNotify` removes that subscription idempotently. Notifications remain +best effort, so syncers compare revisions and commit hashes and recover through +direct pulls. When notifying a remote authority, ZDS resolves its dedicated `#atproto_space_host` service when present and falls back to `#atproto_pds`. @@ -149,8 +153,8 @@ explicit space-scoped tables instead of the public repo tables: each known writer's latest revision and 32-byte commit digest - `permissioned_space_record_oplog`: incremental record changes by `(space, repo_did, rev, idx)` -- `permissioned_space_notify_registrations`: expiring whole-space and - repo-scoped write-notification subscriptions +- `permissioned_space_notify_registrations`: expiring service-identified + write-notification subscriptions - `permissioned_space_used_delegations`: consumed one-use delegation-token IDs - `permissioned_space_used_client_attestations`: consumed one-use client attestation IDs @@ -176,7 +180,8 @@ resolved to the granting account's DID when the token is issued. OAuth reads are limited to the authenticated account's own permissioned repo. Whole-space writer enumeration and cross-repo reads require a space credential. Delegation tokens and client attestations are short-lived and consumed together -in one transaction when the authority mints that credential. +in one transaction when the authority mints that credential. A missing, +invalid, or replayed issuance DPoP proof does not consume either token. Permissioned records and blobs are not public repo records. They must not be squeezed into `records`, `repo_blocks`, `commits`, or `seq_events`. @@ -239,8 +244,10 @@ record CID; deletes and superseded operations remain metadata-only. `com.atproto.space.getRepo` is the full-state recovery path. Its CAR declares the freshly generated deniable signed commit and a DAG-CBOR record index as its -two roots, then carries record blocks in lexicographic `collection/rkey` order. -Blobs are deliberately excluded and remain available through `getBlob`. +two roots. The record index uses canonical DAG-CBOR map ordering (key length, +then bytewise key order), and the commit MAC uses the proposal's direct +HKDF-Expand construction. Blobs are deliberately excluded; `listBlobs` +enumerates referenced CIDs and `getBlob` serves the authorized bytes. ## Zat first diff --git a/src/atproto/space.zig b/src/atproto/space.zig index f48015f..e480636 100644 --- a/src/atproto/space.zig +++ b/src/atproto/space.zig @@ -11,6 +11,7 @@ const config = @import("../core/config.zig"); const http_api = @import("../http/api.zig"); const permissioned = @import("../internal/permissioned_data.zig"); const client_attestation = @import("../internal/client_attestation.zig"); +const dpop = @import("../internal/dpop.zig"); const scopes = @import("../internal/scopes.zig"); const space_uris = @import("../internal/space_uri.zig"); const store = @import("../storage/store.zig"); @@ -43,6 +44,7 @@ pub fn dispatch(request: *http_api.Request) !void { if (std.mem.eql(u8, method, "com.atproto.space.getRecord")) return getRecord(request); if (std.mem.eql(u8, method, "com.atproto.space.listRecords")) return listRecords(request); if (std.mem.eql(u8, method, "com.atproto.space.getBlob")) return getBlob(request); + if (std.mem.eql(u8, method, "com.atproto.space.listBlobs")) return listBlobs(request); if (std.mem.eql(u8, method, "com.atproto.space.getLatestCommit")) return getLatestCommit(request); // newer name for the same operation in the atpint/bulleted stack; both // stay routable while the draft settles on one. @@ -50,6 +52,7 @@ pub fn dispatch(request: *http_api.Request) !void { if (std.mem.eql(u8, method, "com.atproto.space.getRepo")) return getRepo(request); if (std.mem.eql(u8, method, "com.atproto.space.listRepoOps")) return listRepoOps(request); if (std.mem.eql(u8, method, "com.atproto.space.registerNotify")) return registerNotify(request); + if (std.mem.eql(u8, method, "com.atproto.space.unregisterNotify")) return unregisterNotify(request); if (std.mem.eql(u8, method, "com.atproto.space.getSpaceCredential")) return getSpaceCredential(request); if (std.mem.eql(u8, method, "com.atproto.space.notifyWrite")) return notifyWrite(request); if (std.mem.eql(u8, method, "com.atproto.space.notifySpaceDeleted")) return notifySpaceDeleted(request); @@ -373,7 +376,7 @@ fn deleteSimpleSpace(request: *http_api.Request) !void { if (existing == null) return http_api.xrpcError(request, .not_found, "SpaceNotFound", "Space not found"); if (!existing.?.is_authority) return http_api.xrpcError(request, .forbidden, "NotSpaceOwner", "Not the space owner"); if (existing.?.deleted_at != null) return http_api.json(request, .ok, "{}"); - const recipients = try store.listNotificationRecipients(allocator, space, null, true); + const recipients = try store.listNotificationRecipients(allocator, space); try store.markSpaceDeleted(auth_ctx.account.did, space); try store.purgeAuthoritySpaceData(space); fireNotifySpaceDeleted(allocator, auth_ctx.account, space, recipients) catch {}; @@ -675,6 +678,46 @@ fn getBlob(request: *http_api.Request) !void { return writeBlobResponse(request, allocator, cid, blob); } +fn listBlobs(request: *http_api.Request) !void { + var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); + defer arena.deinit(); + const allocator = arena.allocator(); + + var space_buf: [1024]u8 = undefined; + const space = http_api.queryParam(request.url.raw, "space", &space_buf) orelse { + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing space"); + }; + var repo_buf: [256]u8 = undefined; + const repo = http_api.queryParam(request.url.raw, "repo", &repo_buf) orelse { + return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing repo"); + }; + _ = requireReadAccess(request, allocator, space, repo, null) catch return; + + var since_buf: [128]u8 = undefined; + const since = http_api.queryParam(request.url.raw, "since", &since_buf); + if (since) |rev| { + if (zat.Tid.parse(rev) == null) return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid since revision"); + } + var cursor_buf: [256]u8 = undefined; + const cursor = http_api.queryParam(request.url.raw, "cursor", &cursor_buf); + const limit = listBlobsLimit(request) catch return; + const cids = try store.listPermissionedBlobs(allocator, space, repo, since, cursor, limit); + + var out: std.Io.Writer.Allocating = .init(allocator); + defer out.deinit(); + try out.writer.writeAll("{\"cids\":["); + for (cids, 0..) |cid, idx| { + if (idx != 0) try out.writer.writeByte(','); + try out.writer.print("{f}", .{std.json.fmt(cid, .{})}); + } + try out.writer.writeByte(']'); + if (cids.len == limit and cids.len > 0) { + try out.writer.print(",\"cursor\":{f}", .{std.json.fmt(cids[cids.len - 1], .{})}); + } + try out.writer.writeByte('}'); + return http_api.json(request, .ok, out.written()); +} + fn getLatestCommit(request: *http_api.Request) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); @@ -782,32 +825,19 @@ fn registerNotify(request: *http_api.Request) !void { }; const input = parsed_body.value; const space = zat.json.getString(input, "space") orelse return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing space"); - const endpoint = zat.json.getString(input, "endpoint") orelse return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing endpoint"); - const repo = zat.json.getString(input, "repo"); - _ = std.Uri.parse(endpoint) catch return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid notification endpoint"); - if (!std.mem.startsWith(u8, endpoint, "https://") and !std.mem.startsWith(u8, endpoint, "http://")) { - return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Notification endpoint must be an HTTP URL"); - } + const service = zat.json.getString(input, "service") orelse return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing service"); _ = requireSpaceCredential(request, allocator, space) catch return; const parsed = parseSpaceUri(space) orelse return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid space URI"); - - if (repo) |repo_did| { - if (zat.Did.parse(repo_did) == null or (try store.findAccount(allocator, repo_did)) == null) { - return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"); - } - const state = try store.getSpaceRepoState(allocator, space, repo_did); - if (state.rev == null) return http_api.xrpcError(request, .not_found, "RepoNotFound", "Permissioned repo not found"); - } else { - if ((try store.findAccount(allocator, parsed.authority_did)) == null) { - return http_api.xrpcError(request, .not_found, "SpaceNotFound", "Space not found"); - } - const existing = try store.getSpace(allocator, parsed.authority_did, space); - if (existing == null or !existing.?.is_authority or existing.?.deleted_at != null) { - return http_api.xrpcError(request, .not_found, "SpaceNotFound", "Space not found"); - } + if ((try store.findAccount(allocator, parsed.authority_did)) == null) { + return http_api.xrpcError(request, .not_found, "SpaceNotFound", "Space not found"); } + const existing = try store.getSpace(allocator, parsed.authority_did, space); + if (existing == null or !existing.?.is_authority or existing.?.deleted_at != null) { + return http_api.xrpcError(request, .not_found, "SpaceNotFound", "Space not found"); + } + const endpoint = serviceEndpointForId(allocator, service) catch return http_api.xrpcError(request, .bad_request, "ServiceNotResolvable", "Could not resolve notification service"); const expires_at = unixNow() + 24 * 60 * 60; - try store.registerSpaceNotification(space, repo, endpoint, expires_at); + try store.registerSpaceNotification(space, service, endpoint, expires_at); return http_api.json(request, .ok, try std.fmt.allocPrint( allocator, "{{\"expiresAt\":{f}}}", @@ -815,11 +845,35 @@ fn registerNotify(request: *http_api.Request) !void { )); } +fn unregisterNotify(request: *http_api.Request) !void { + var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); + defer arena.deinit(); + const allocator = arena.allocator(); + const parsed_body = parseBody(request, allocator, 16 * 1024) catch |err| switch (err) { + error.HandledResponse => return, + else => return err, + }; + const input = parsed_body.value; + const space = zat.json.getString(input, "space") orelse return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing space"); + const service = zat.json.getString(input, "service") orelse return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Missing service"); + _ = requireSpaceCredential(request, allocator, space) catch return; + const parsed = parseSpaceUri(space) orelse return http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid space URI"); + const existing = try store.getSpace(allocator, parsed.authority_did, space); + if (existing == null or !existing.?.is_authority or existing.?.deleted_at != null) { + return http_api.xrpcError(request, .not_found, "SpaceNotFound", "Space not found"); + } + try store.unregisterSpaceNotification(space, service); + return http_api.json(request, .ok, "{}"); +} + fn getSpaceCredential(request: *http_api.Request) !void { var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer arena.deinit(); const allocator = arena.allocator(); const token = bearerToken(request) orelse return http_api.xrpcError(request, .unauthorized, "AuthenticationRequired", "Authentication required"); + const proof = dpop.verifySpaceRequest(allocator, request, null, null) catch { + return http_api.xrpcError(request, .unauthorized, "InvalidDpopProof", "A valid DPoP proof is required to obtain a space credential"); + }; const requester_did = permissioned.unverifiedStringClaim(allocator, token, "iss") catch { return http_api.xrpcError(request, .unauthorized, "InvalidToken", "Invalid delegation token"); }; @@ -863,7 +917,7 @@ fn getSpaceCredential(request: *http_api.Request) !void { return http_api.xrpcError(request, .unauthorized, "InvalidDelegationToken", "Credential exchange token has already been used"); } var keypair = store.signingKeypair(parsed.authority_did) catch return http_api.xrpcError(request, .not_found, "RepoNotFound", "Authority signing key not found"); - const credential = try permissioned.createSpaceCredential(allocator, store.currentIo(), parsed.authority_did, space, &keypair); + const credential = try permissioned.createSpaceCredential(allocator, store.currentIo(), parsed.authority_did, space, proof.jkt, attested_client_id, &keypair); return http_api.json(request, .ok, try std.fmt.allocPrint(allocator, "{{\"credential\":{f}}}", .{std.json.fmt(credential, .{})})); } @@ -1005,8 +1059,19 @@ fn bearerToken(request: *http_api.Request) ?[]const u8 { return std.mem.trim(u8, raw_header["bearer ".len..], " \t"); } +fn dpopToken(request: *http_api.Request) ?[]const u8 { + const raw_header = http_api.headerValue(request, "authorization") orelse return null; + if (!std.ascii.startsWithIgnoreCase(raw_header, "dpop ")) return null; + return std.mem.trim(u8, raw_header["dpop ".len..], " \t"); +} + fn signingKeyForDid(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { + return signingKeyForDidKid(allocator, did, "#atproto"); +} + +fn signingKeyForDidKid(allocator: std.mem.Allocator, did: []const u8, kid: []const u8) ![]const u8 { if (try store.findAccount(allocator, did) != null) { + if (!std.mem.eql(u8, kid, "#atproto")) return error.MissingSigningKey; var keypair = try store.signingKeypair(did); const did_key = try keypair.did(allocator); defer allocator.free(did_key); @@ -1018,8 +1083,18 @@ fn signingKeyForDid(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { defer resolver.deinit(); var doc = try resolver.resolve(parsed); defer doc.deinit(); - const signing_key = doc.signingKey() orelse return error.MissingSigningKey; - return allocator.dupe(u8, signing_key.public_key_multibase); + const absolute_kid = if (std.mem.startsWith(u8, kid, "#")) + try std.fmt.allocPrint(allocator, "{s}{s}", .{ did, kid }) + else + try allocator.dupe(u8, kid); + defer allocator.free(absolute_kid); + for (doc.verification_methods) |method| { + if (std.mem.eql(u8, method.id, kid) or std.mem.eql(u8, method.id, absolute_kid)) { + if (!std.mem.eql(u8, method.controller, did)) return error.InvalidSigningKey; + return allocator.dupe(u8, method.public_key_multibase); + } + } + return error.MissingSigningKey; } const ServiceAuth = struct { @@ -1089,20 +1164,23 @@ fn spaceHostEndpointForDid(allocator: std.mem.Allocator, did: []const u8) ![]con } fn serviceEndpointForId(allocator: std.mem.Allocator, service_id: []const u8) ![]const u8 { - const fragment_index = std.mem.indexOfScalar(u8, service_id, '#') orelse return error.InvalidServiceId; - const did_text = service_id[0..fragment_index]; - const fragment = service_id[fragment_index..]; + const fragment_index = std.mem.indexOfScalar(u8, service_id, '#'); + const did_text = if (fragment_index) |idx| service_id[0..idx] else service_id; + const fragment = if (fragment_index) |idx| service_id[idx..] else null; const did = zat.Did.parse(did_text) orelse return error.InvalidServiceId; var resolver = zat.DidResolver.init(store.currentIo(), allocator); defer resolver.deinit(); var doc = try resolver.resolve(did); defer doc.deinit(); - for (doc.services) |service| { - if (std.mem.eql(u8, service.id, service_id) or std.mem.eql(u8, service.id, fragment)) { - return allocator.dupe(u8, service.service_endpoint); + if (fragment) |service_fragment| { + for (doc.services) |service| { + if (std.mem.eql(u8, service.id, service_id) or std.mem.eql(u8, service.id, service_fragment)) { + return allocator.dupe(u8, service.service_endpoint); + } } + return error.MissingServiceEndpoint; } - return error.MissingServiceEndpoint; + return allocator.dupe(u8, doc.pdsEndpoint() orelse return error.MissingServiceEndpoint); } fn managingAppAllowsRequester( @@ -1207,21 +1285,21 @@ fn fireNotifyWriteForState(allocator: std.mem.Allocator, account: auth.Account, } fn fanoutNotifyWriteToRecipients(allocator: std.mem.Allocator, account: auth.Account, space: []const u8, repo: []const u8, rev: []const u8, hash: []const u8) !void { - const recipients = try store.listNotificationRecipients(allocator, space, repo, true); + const recipients = try store.listNotificationRecipients(allocator, space); for (recipients) |recipient| { const body = try std.fmt.allocPrint( allocator, "{{\"space\":{f},\"repo\":{f},\"rev\":{f},\"hash\":{s}}}", .{ std.json.fmt(space, .{}), std.json.fmt(repo, .{}), std.json.fmt(rev, .{}), try atJsonBytes(allocator, hash) }, ); - postSignedXrpc(allocator, account, recipient.service_endpoint, recipient.service_endpoint, "com.atproto.space.notifyWrite", body) catch {}; + postSignedXrpc(allocator, account, recipient.service_endpoint, recipient.service_id, "com.atproto.space.notifyWrite", body) catch {}; } } fn fireNotifySpaceDeleted(allocator: std.mem.Allocator, account: auth.Account, space: []const u8, recipients: []const store.CredentialRecipient) !void { const body = try std.fmt.allocPrint(allocator, "{{\"space\":{f}}}", .{std.json.fmt(space, .{})}); for (recipients) |recipient| { - postSignedXrpc(allocator, account, recipient.service_endpoint, recipient.service_endpoint, "com.atproto.space.notifySpaceDeleted", body) catch {}; + postSignedXrpc(allocator, account, recipient.service_endpoint, recipient.service_id, "com.atproto.space.notifySpaceDeleted", body) catch {}; } } @@ -1321,6 +1399,23 @@ fn requireReadAccess(request: *http_api.Request, allocator: std.mem.Allocator, s try http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid repo DID"); return error.HandledResponse; } + const repo_status = store.accountStatus(repo) catch |err| switch (err) { + error.RepoNotFound => { + try http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"); + return error.HandledResponse; + }, + else => return err, + }; + if (!repo_status.isActive()) { + switch (repo_status) { + .active => unreachable, + .takendown => try http_api.xrpcError(request, .bad_request, "RepoTakendown", "Repo has been taken down"), + .suspended => try http_api.xrpcError(request, .bad_request, "RepoSuspended", "Repo is suspended"), + .deactivated => try http_api.xrpcError(request, .bad_request, "RepoDeactivated", "Repo is deactivated"), + .deleted => try http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"), + } + return error.HandledResponse; + } const parsed = parseSpaceUri(space) orelse { try http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid space URI"); return error.HandledResponse; @@ -1337,12 +1432,26 @@ fn requireReadAccess(request: *http_api.Request, allocator: std.mem.Allocator, s return .{ .space_credential = try requireSpaceCredential(request, allocator, space) }; } +fn listBlobsLimit(request: *http_api.Request) !usize { + var buf: [16]u8 = undefined; + const raw = http_api.queryParam(request.url.raw, "limit", &buf) orelse return 500; + const limit = std.fmt.parseInt(usize, raw, 10) catch { + try http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid limit"); + return error.HandledResponse; + }; + if (limit < 1 or limit > 1000) { + try http_api.xrpcError(request, .bad_request, "InvalidRequest", "limit must be between 1 and 1000"); + return error.HandledResponse; + } + return limit; +} + fn requireSpaceCredential(request: *http_api.Request, allocator: std.mem.Allocator, space: []const u8) !permissioned.SpaceCredential { const parsed = parseSpaceUri(space) orelse { try http_api.xrpcError(request, .bad_request, "InvalidRequest", "Invalid space URI"); return error.HandledResponse; }; - const token = bearerToken(request) orelse { + const token = dpopToken(request) orelse { try http_api.xrpcError(request, .unauthorized, "AuthenticationRequired", "Authentication required"); return error.HandledResponse; }; @@ -1354,7 +1463,15 @@ fn requireSpaceCredential(request: *http_api.Request, allocator: std.mem.Allocat try http_api.xrpcError(request, .unauthorized, "InvalidToken", "Space credential issuer must be the space authority"); return error.HandledResponse; } - const public_key = signingKeyForDid(allocator, issuer) catch { + const kid = permissioned.unverifiedHeaderString(allocator, token, "kid") catch { + try http_api.xrpcError(request, .unauthorized, "InvalidToken", "Invalid space credential"); + return error.HandledResponse; + }; + if (!std.mem.eql(u8, kid, "#atproto") and !std.mem.eql(u8, kid, "#atproto_space")) { + try http_api.xrpcError(request, .unauthorized, "InvalidToken", "Invalid space credential signing key"); + return error.HandledResponse; + } + const public_key = signingKeyForDidKid(allocator, issuer, kid) catch { try http_api.xrpcError(request, .bad_gateway, "DidResolutionFailed", "could not resolve credential issuer did"); return error.HandledResponse; }; @@ -1366,6 +1483,10 @@ fn requireSpaceCredential(request: *http_api.Request, allocator: std.mem.Allocat try http_api.xrpcError(request, .bad_request, "InvalidRequest", "Credential space mismatch"); return error.HandledResponse; } + _ = dpop.verifySpaceRequest(allocator, request, token, credential.dpop_jkt) catch { + try http_api.xrpcError(request, .unauthorized, "InvalidDpopProof", "Invalid DPoP proof for space credential"); + return error.HandledResponse; + }; return credential; } diff --git a/src/http/router.zig b/src/http/router.zig index 67f8316..f2fabc2 100644 --- a/src/http/router.zig +++ b/src/http/router.zig @@ -217,9 +217,9 @@ pub const endpoints = [_]Endpoint{ .{ .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 = .identity_update_handle, .method = "POST", .path = "/xrpc/com.atproto.identity.updateHandle", .group = "identity", .auth = "bearer", .summary = "Update the signed-in account's handle.", .body = &.{"handle"} }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getSpace", .group = "space", .auth = "experimental space credential or resident bearer", .summary = "Read space configuration from its authority host.", .params = &.{"space"}, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getSpace", .group = "space", .auth = "experimental OAuth bearer or DPoP space credential", .summary = "Read space configuration from its authority host.", .params = &.{"space"}, .notes = permissioned_data_note }, .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listSpaces", .group = "space", .auth = "experimental bearer", .summary = "List permissioned repos held by the authenticated user, grouped by space.", .params = &.{ "did", "type", "limit", "cursor" }, .notes = permissioned_data_note }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listRepos", .group = "space", .auth = "experimental space credential", .summary = "List the authority's known writer repos and commit digests.", .params = &.{ "space", "limit", "cursor" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listRepos", .group = "space", .auth = "experimental DPoP space credential", .summary = "List the authority's known writer repos and commit digests.", .params = &.{ "space", "limit", "cursor" }, .notes = permissioned_data_note }, .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getDelegationToken", .group = "space", .auth = "experimental OAuth", .summary = "Create a delegation token for exchange with a space authority.", .params = &.{"space"}, .notes = permissioned_data_note }, @@ -227,17 +227,19 @@ pub const endpoints = [_]Endpoint{ .{ .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 records in a permissioned data space. Values are included unless excludeValues=true.", .params = &.{ "space", "repo", "collection", "limit", "cursor", "reverse", "excludeValues" }, .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.getLatestCommit", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read the current signed commit for a writer repo in a space.", .params = &.{ "space", "repo" }, .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 the current signed state for a writer repo in a space (newer name for getLatestCommit).", .params = &.{ "space", "repo" }, .notes = permissioned_data_note }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRepo", .group = "space", .auth = "experimental bearer or space credential", .summary = "Download a full permissioned repo CAR for recovery.", .params = &.{ "space", "repo" }, .notes = permissioned_data_note }, - .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listRepoOps", .group = "space", .auth = "experimental bearer or space credential", .summary = "Read incremental record operations. Values are included unless excludeValues=true.", .params = &.{ "space", "repo", "since", "limit", "excludeValues" }, .notes = permissioned_data_note }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.registerNotify", .group = "space", .auth = "experimental space credential", .summary = "Register an expiring write-notification endpoint for a space or repo.", .body = &.{ "space", "repo", "endpoint" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRecord", .group = "space", .auth = "experimental OAuth bearer or DPoP 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 OAuth bearer or DPoP space credential", .summary = "List records in a permissioned data space. Values are included unless excludeValues=true.", .params = &.{ "space", "repo", "collection", "limit", "cursor", "reverse", "excludeValues" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getBlob", .group = "space", .auth = "experimental OAuth bearer or DPoP 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.listBlobs", .group = "space", .auth = "experimental OAuth bearer or DPoP space credential", .summary = "List blobs referenced by a permissioned repo, optionally since a revision.", .params = &.{ "space", "repo", "since", "limit", "cursor" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getLatestCommit", .group = "space", .auth = "experimental OAuth bearer or DPoP space credential", .summary = "Read the current signed commit for a writer repo in a space.", .params = &.{ "space", "repo" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRepoState", .group = "space", .auth = "experimental OAuth bearer or DPoP space credential", .summary = "Read the current signed state for a writer repo in a space (newer name for getLatestCommit).", .params = &.{ "space", "repo" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.getRepo", .group = "space", .auth = "experimental OAuth bearer or DPoP space credential", .summary = "Download a full permissioned repo CAR for recovery.", .params = &.{ "space", "repo" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "GET", .path = "/xrpc/com.atproto.space.listRepoOps", .group = "space", .auth = "experimental OAuth bearer or DPoP space credential", .summary = "Read incremental record operations. Values are included unless excludeValues=true.", .params = &.{ "space", "repo", "since", "limit", "excludeValues" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.registerNotify", .group = "space", .auth = "experimental DPoP-bound space credential", .summary = "Register an expiring write-notification service for a space.", .body = &.{ "space", "service" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.unregisterNotify", .group = "space", .auth = "experimental DPoP-bound space credential", .summary = "Withdraw a write-notification registration.", .body = &.{ "space", "service" }, .notes = permissioned_data_note }, .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.notifyWrite", .group = "space", .auth = "experimental service", .summary = "Notify a space authority or syncing service of a permissioned data write.", .body = &.{ "space", "repo", "rev", "hash" }, .notes = permissioned_data_note }, - .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.getSpaceCredential", .group = "space", .auth = "experimental delegation token", .summary = "Exchange a delegation token for a space credential.", .body = &.{ "space", "clientAttestation" }, .notes = permissioned_data_note }, + .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.getSpaceCredential", .group = "space", .auth = "experimental delegation token and DPoP proof", .summary = "Exchange a delegation token for a DPoP-bound space credential.", .body = &.{ "space", "clientAttestation" }, .notes = permissioned_data_note }, .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.space.notifySpaceDeleted", .group = "space", .auth = "experimental service", .summary = "Notify a repo host or syncing service that a space was deleted.", .body = &.{"space"}, .notes = permissioned_data_note }, .{ .route = .permissioned_data, .method = "POST", .path = "/xrpc/com.atproto.simplespace.createSpace", .group = "simplespace", .auth = "experimental bearer", .summary = "Create or materialize a baseline PDS-managed permissioned data space.", .body = &.{ "did", "type", "skey", "config" }, .notes = permissioned_data_note }, diff --git a/src/internal/dpop.zig b/src/internal/dpop.zig index daad69b..f172456 100644 --- a/src/internal/dpop.zig +++ b/src/internal/dpop.zig @@ -51,11 +51,33 @@ pub fn verifyRequest( request: *const http_api.Request, access_token: ?[]const u8, expected_jkt: ?[]const u8, +) Error!Proof { + return verifyRequestWithPolicy(allocator, request, access_token, expected_jkt, true, ProofMaxAgeSecs, ProofMaxAgeSecs, false); +} + +pub fn verifySpaceRequest( + allocator: std.mem.Allocator, + request: *const http_api.Request, + credential: ?[]const u8, + expected_jkt: ?[]const u8, +) Error!Proof { + return verifyRequestWithPolicy(allocator, request, credential, expected_jkt, false, 60, 5, true); +} + +fn verifyRequestWithPolicy( + allocator: std.mem.Allocator, + request: *const http_api.Request, + access_token: ?[]const u8, + expected_jkt: ?[]const u8, + require_nonce: bool, + max_age_secs: i64, + future_skew_secs: i64, + require_es256: bool, ) Error!Proof { const dpop_header = http_api.headerValue(request, "dpop") orelse return error.MissingProof; const htu = try htuForRequest(allocator, request); defer allocator.free(htu); - return verifyProof(allocator, dpop_header, methodName(request.method), htu, access_token, expected_jkt); + return verifyProof(allocator, dpop_header, methodName(request.method), htu, access_token, expected_jkt, require_nonce, max_age_secs, future_skew_secs, require_es256); } pub fn maybeVerifyRequest( @@ -75,6 +97,10 @@ fn verifyProof( expected_htu: []const u8, access_token: ?[]const u8, expected_jkt: ?[]const u8, + require_nonce: bool, + max_age_secs: i64, + future_skew_secs: i64, + require_es256: bool, ) Error!Proof { var parts = std.mem.splitScalar(u8, proof, '.'); const header_part = parts.next() orelse return error.InvalidProof; @@ -97,6 +123,7 @@ fn verifyProof( if (!std.mem.eql(u8, zat.json.getString(header.value, "typ") orelse return error.InvalidProof, "dpop+jwt")) return error.InvalidProof; const alg_text = zat.json.getString(header.value, "alg") orelse return error.InvalidProof; const alg = zat.jwt.Algorithm.fromString(alg_text) orelse return error.InvalidProof; + if (require_es256 and alg != .ES256) return error.InvalidProof; const jwk = zat.json.getPath(header.value, "jwk") orelse return error.InvalidProof; const public_key = try publicKeyFromJwk(allocator, alg, jwk); @@ -115,10 +142,12 @@ fn verifyProof( const iat = zat.json.getInt(payload.value, "iat") orelse return error.InvalidProof; const ts = now(); - if (iat < ts - ProofMaxAgeSecs or iat > ts + ProofMaxAgeSecs) return error.InvalidProof; + if (iat < ts - max_age_secs or iat > ts + future_skew_secs) return error.InvalidProof; - const nonce = zat.json.getString(payload.value, "nonce") orelse return error.UseDpopNonce; - if (!validNonce(nonce, ts)) return error.UseDpopNonce; + if (require_nonce) { + const nonce = zat.json.getString(payload.value, "nonce") orelse return error.UseDpopNonce; + if (!validNonce(nonce, ts)) return error.UseDpopNonce; + } if (access_token) |token| { const expected_ath = try zat.oauth.accessTokenHash(allocator, token); @@ -133,7 +162,7 @@ fn verifyProof( if (expected_jkt) |expected| { if (!std.mem.eql(u8, jkt, expected)) return error.KeyBindingMismatch; } - const expires_at = @max(iat + ProofMaxAgeSecs, ts + ProofMaxAgeSecs); + const expires_at = @max(iat + max_age_secs, ts + max_age_secs); if (!(store.recordDpopJti(jti, expires_at) catch return error.InvalidProof)) return error.Replay; return .{ @@ -284,11 +313,11 @@ test "validates DPoP proof and rejects replay" { const expected_jkt = try keypair.jwkThumbprint(allocator); defer allocator.free(expected_jkt); - const verified = try verifyProof(allocator, proof, "GET", "https://pds.example/xrpc/com.atproto.repo.getRecord", access_token, expected_jkt); + const verified = try verifyProof(allocator, proof, "GET", "https://pds.example/xrpc/com.atproto.repo.getRecord", access_token, expected_jkt, true, ProofMaxAgeSecs, ProofMaxAgeSecs, false); defer allocator.free(verified.jkt); defer allocator.free(verified.jti); try std.testing.expectEqualStrings(expected_jkt, verified.jkt); - try std.testing.expectError(error.Replay, verifyProof(allocator, proof, "GET", "https://pds.example/xrpc/com.atproto.repo.getRecord", access_token, expected_jkt)); + try std.testing.expectError(error.Replay, verifyProof(allocator, proof, "GET", "https://pds.example/xrpc/com.atproto.repo.getRecord", access_token, expected_jkt, true, ProofMaxAgeSecs, ProofMaxAgeSecs, false)); } test "rejects DPoP proof with wrong token binding" { @@ -319,7 +348,7 @@ test "rejects DPoP proof with wrong token binding" { const expected_jkt = try keypair.jwkThumbprint(allocator); defer allocator.free(expected_jkt); - try std.testing.expectError(error.InvalidProof, verifyProof(allocator, proof, "POST", "https://pds.example/xrpc/com.atproto.repo.createRecord", "access-token-b", expected_jkt)); + try std.testing.expectError(error.InvalidProof, verifyProof(allocator, proof, "POST", "https://pds.example/xrpc/com.atproto.repo.createRecord", "access-token-b", expected_jkt, true, ProofMaxAgeSecs, ProofMaxAgeSecs, false)); } test "rejects DPoP proof with mismatched JKT" { @@ -346,5 +375,39 @@ test "rejects DPoP proof with mismatched JKT" { ); defer allocator.free(proof); - try std.testing.expectError(error.KeyBindingMismatch, verifyProof(allocator, proof, "POST", "https://pds.example/oauth/token", null, "not-the-key")); + try std.testing.expectError(error.KeyBindingMismatch, verifyProof(allocator, proof, "POST", "https://pds.example/oauth/token", null, "not-the-key", true, ProofMaxAgeSecs, ProofMaxAgeSecs, false)); +} + +test "space DPoP proof does not require an OAuth nonce" { + const allocator = std.testing.allocator; + try store.init(std.Options.debug_io, ":memory:"); + defer store.close(); + const keypair = try zat.Keypair.fromSecretKey(.p256, .{0x31} ** 32); + const proof = try zat.oauth.createDpopProof( + allocator, + std.Options.debug_io, + &keypair, + "POST", + "https://space.example/xrpc/com.atproto.space.getSpaceCredential", + "", + null, + ); + defer allocator.free(proof); + const expected_jkt = try keypair.jwkThumbprint(allocator); + defer allocator.free(expected_jkt); + const verified = try verifyProof( + allocator, + proof, + "POST", + "https://space.example/xrpc/com.atproto.space.getSpaceCredential", + null, + null, + false, + 60, + 5, + true, + ); + defer allocator.free(verified.jkt); + defer allocator.free(verified.jti); + try std.testing.expectEqualStrings(expected_jkt, verified.jkt); } diff --git a/src/internal/permissioned_data.zig b/src/internal/permissioned_data.zig index 0a3f62a..da80f0b 100644 --- a/src/internal/permissioned_data.zig +++ b/src/internal/permissioned_data.zig @@ -94,6 +94,8 @@ pub const DelegationToken = struct { pub const SpaceCredential = struct { authority_did: []const u8, space: []const u8, + dpop_jkt: []const u8, + client_id: ?[]const u8, exp: i64, }; @@ -158,13 +160,14 @@ pub fn serializeRepoCar( const index_entries = try allocator.alloc(zat.cbor.Value.MapEntry, records.len); for (records, 0..) |record, idx| { - if (idx > 0 and std.mem.order(u8, records[idx - 1].path, record.path) != .lt) { - return error.RecordsNotSorted; - } const computed = try zat.cbor.Cid.forDagCbor(allocator, record.data); if (!std.mem.eql(u8, computed.raw, record.cid.raw)) return error.RecordCidMismatch; index_entries[idx] = .{ .key = record.path, .value = .{ .cid = record.cid } }; } + std.mem.sort(zat.cbor.Value.MapEntry, index_entries, {}, canonicalMapKeyLessThan); + for (index_entries[1..], index_entries[0..index_entries.len -| 1]) |current, previous| { + if (std.mem.eql(u8, previous.key, current.key)) return error.DuplicateRecordPath; + } const index_bytes = try zat.cbor.encodeAlloc(allocator, .{ .map = index_entries }); const index_cid = try zat.cbor.Cid.forDagCbor(allocator, index_bytes); @@ -216,6 +219,8 @@ pub fn createSpaceCredential( io: std.Io, authority_did: []const u8, space: []const u8, + dpop_jkt: []const u8, + client_id: ?[]const u8, keypair: *const zat.Keypair, ) ![]const u8 { const iat = unixNow(); @@ -224,17 +229,24 @@ pub fn createSpaceCredential( defer allocator.free(jti); const header = try std.fmt.allocPrint(allocator, "{{\"typ\":\"atproto-space-credential+jwt\",\"alg\":\"{s}\",\"kid\":\"#atproto\"}}", .{@tagName(keypair.algorithm())}); defer allocator.free(header); - const payload = try std.fmt.allocPrint( - allocator, - "{{\"iss\":{f},\"sub\":{f},\"iat\":{d},\"exp\":{d},\"jti\":{f}}}", - .{ std.json.fmt(authority_did, .{}), std.json.fmt(space, .{}), iat, exp, std.json.fmt(jti, .{}) }, - ); + const payload = if (client_id) |client| + try std.fmt.allocPrint( + allocator, + "{{\"iss\":{f},\"sub\":{f},\"cnf\":{{\"jkt\":{f}}},\"client_id\":{f},\"iat\":{d},\"exp\":{d},\"jti\":{f}}}", + .{ std.json.fmt(authority_did, .{}), std.json.fmt(space, .{}), std.json.fmt(dpop_jkt, .{}), std.json.fmt(client, .{}), iat, exp, std.json.fmt(jti, .{}) }, + ) + else + try std.fmt.allocPrint( + allocator, + "{{\"iss\":{f},\"sub\":{f},\"cnf\":{{\"jkt\":{f}}},\"iat\":{d},\"exp\":{d},\"jti\":{f}}}", + .{ std.json.fmt(authority_did, .{}), std.json.fmt(space, .{}), std.json.fmt(dpop_jkt, .{}), iat, exp, std.json.fmt(jti, .{}) }, + ); defer allocator.free(payload); return zat.oauth.createJwt(allocator, header, payload, keypair); } pub fn verifyDelegationToken(allocator: std.mem.Allocator, token: []const u8, public_key_multibase: []const u8) !DelegationToken { - const parsed = try parseAndVerifyJwt(allocator, token, public_key_multibase, "atproto-space-delegation+jwt"); + const parsed = try parseAndVerifyJwt(allocator, token, public_key_multibase, "atproto-space-delegation+jwt", &.{"#atproto"}); defer parsed.deinit(); const exp = zat.json.getInt(parsed.payload.value, "exp") orelse return error.InvalidJwt; if (exp < unixNow()) return error.ExpiredJwt; @@ -250,13 +262,15 @@ pub fn verifyDelegationToken(allocator: std.mem.Allocator, token: []const u8, pu } pub fn verifySpaceCredential(allocator: std.mem.Allocator, token: []const u8, public_key_multibase: []const u8) !SpaceCredential { - const parsed = try parseAndVerifyJwt(allocator, token, public_key_multibase, "atproto-space-credential+jwt"); + const parsed = try parseAndVerifyJwt(allocator, token, public_key_multibase, "atproto-space-credential+jwt", &.{ "#atproto", "#atproto_space" }); defer parsed.deinit(); const exp = zat.json.getInt(parsed.payload.value, "exp") orelse return error.InvalidJwt; if (exp < unixNow()) return error.ExpiredJwt; return .{ .authority_did = try allocator.dupe(u8, zat.json.getString(parsed.payload.value, "iss") orelse return error.InvalidJwt), .space = try allocator.dupe(u8, zat.json.getString(parsed.payload.value, "sub") orelse return error.InvalidJwt), + .dpop_jkt = try allocator.dupe(u8, zat.json.getString(parsed.payload.value, "cnf.jkt") orelse return error.InvalidJwt), + .client_id = if (zat.json.getString(parsed.payload.value, "client_id")) |client_id| try allocator.dupe(u8, client_id) else null, .exp = exp, }; } @@ -270,6 +284,14 @@ fn authorityDidFromAudience(audience: []const u8) ?[]const u8 { } pub fn unverifiedStringClaim(allocator: std.mem.Allocator, token: []const u8, claim: []const u8) ![]const u8 { + return unverifiedStringPart(allocator, token, 1, claim); +} + +pub fn unverifiedHeaderString(allocator: std.mem.Allocator, token: []const u8, name: []const u8) ![]const u8 { + return unverifiedStringPart(allocator, token, 0, name); +} + +fn unverifiedStringPart(allocator: std.mem.Allocator, token: []const u8, part_index: usize, name: []const u8) ![]const u8 { var parts: [3][]const u8 = undefined; var part_count: usize = 0; var it = std.mem.splitScalar(u8, token, '.'); @@ -279,11 +301,11 @@ pub fn unverifiedStringClaim(allocator: std.mem.Allocator, token: []const u8, cl part_count += 1; } if (part_count != 3) return error.InvalidJwt; - const payload_json = try zat.jwt.base64UrlDecode(allocator, parts[1]); - defer allocator.free(payload_json); - const parsed = try std.json.parseFromSlice(std.json.Value, allocator, payload_json, .{}); + const json = try zat.jwt.base64UrlDecode(allocator, parts[part_index]); + defer allocator.free(json); + const parsed = try std.json.parseFromSlice(std.json.Value, allocator, json, .{}); defer parsed.deinit(); - return allocator.dupe(u8, zat.json.getString(parsed.value, claim) orelse return error.InvalidJwt); + return allocator.dupe(u8, zat.json.getString(parsed.value, name) orelse return error.InvalidJwt); } fn expandElement(element: []const u8) [lthash_state_bytes]u8 { @@ -308,7 +330,7 @@ const ParsedJwt = struct { } }; -fn parseAndVerifyJwt(allocator: std.mem.Allocator, token: []const u8, public_key_multibase: []const u8, expected_typ: []const u8) !ParsedJwt { +fn parseAndVerifyJwt(allocator: std.mem.Allocator, token: []const u8, public_key_multibase: []const u8, expected_typ: []const u8, allowed_kids: []const []const u8) !ParsedJwt { var parts: [3][]const u8 = undefined; var part_count: usize = 0; var it = std.mem.splitScalar(u8, token, '.'); @@ -334,6 +356,15 @@ fn parseAndVerifyJwt(allocator: std.mem.Allocator, token: []const u8, public_key const typ = zat.json.getString(header.value, "typ") orelse return error.InvalidJwt; if (!std.mem.eql(u8, typ, expected_typ)) return error.InvalidJwt; + const kid = zat.json.getString(header.value, "kid") orelse return error.InvalidJwt; + var kid_allowed = false; + for (allowed_kids) |allowed| { + if (std.mem.eql(u8, kid, allowed)) { + kid_allowed = true; + break; + } + } + if (!kid_allowed) return error.InvalidJwt; const alg_text = zat.json.getString(header.value, "alg") orelse return error.InvalidJwt; const alg = zat.jwt.Algorithm.fromString(alg_text) orelse return error.InvalidJwt; const key_bytes = try zat.multibase.decode(allocator, public_key_multibase); @@ -362,14 +393,18 @@ fn unixNow() i64 { } fn commitMac(ikm: *const [32]u8, context: []const u8, hash: *const [32]u8) [32]u8 { - const prk = HkdfSha256.extract("", ikm); var derived: [32]u8 = undefined; - HkdfSha256.expand(&derived, context, prk); + HkdfSha256.expand(&derived, context, ikm.*); var out: [32]u8 = undefined; HmacSha256.create(&out, hash, &derived); return out; } +fn canonicalMapKeyLessThan(_: void, left: zat.cbor.Value.MapEntry, right: zat.cbor.Value.MapEntry) bool { + if (left.key.len != right.key.len) return left.key.len < right.key.len; + return std.mem.order(u8, left.key, right.key) == .lt; +} + fn commitContext(allocator: std.mem.Allocator, ctx: SpaceContext, ikm: *const [32]u8) ![]const u8 { var out: std.Io.Writer.Allocating = .init(allocator); try out.writer.writeAll("atproto-space-v1"); @@ -458,6 +493,42 @@ test "permissioned repo CAR has commit and index roots followed by sorted record try std.testing.expectEqualSlices(u8, second_cid.raw, index.get("fm.example.note/two").?.cid.raw); } +test "permissioned commit MAC uses HKDF expand without extract" { + var ikm: [32]u8 = undefined; + var hash: [32]u8 = undefined; + for (&ikm, 0..) |*byte, idx| byte.* = @intCast(idx); + for (&hash, 0..) |*byte, idx| byte.* = @intCast(idx + 32); + const mac = commitMac(&ikm, "atproto-space-test-context", &hash); + const expected = "5f162f35ca2a0ad567359c92d907c5c30e15be7db0b8c786fd5bc79074f02369"; + try std.testing.expectEqualStrings(expected, &std.fmt.bytesToHex(mac, .lower)); +} + +test "permissioned repo CAR canonicalizes index keys by encoded length" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const allocator = arena.allocator(); + const data = try zat.cbor.encodeAlloc(allocator, .{ .map = &.{ + .{ .key = "$type", .value = .{ .text = "fm.example.note" } }, + } }); + const cid = try zat.cbor.Cid.forDagCbor(allocator, data); + const commit: SignedCommit = .{ + .hash = .{0} ** 32, + .mac = .{0} ** 32, + .ikm = .{0} ** 32, + .sig = .{0} ** 64, + .rev = "3mcanonicalmap", + }; + const records = [_]RepoRecordBlock{ + .{ .path = "longer/path", .cid = cid, .data = data }, + .{ .path = "z/x", .cid = cid, .data = data }, + }; + const bytes = try serializeRepoCar(allocator, commit, &records); + const car = try zat.car.read(allocator, bytes); + const index = try zat.cbor.decodeAll(allocator, car.blocks[1].data); + try std.testing.expectEqualStrings("z/x", index.map[0].key); + try std.testing.expectEqualStrings("longer/path", index.map[1].key); +} + test "LtHash snapshot vector" { var hash: LtHash = .{}; hash.add("atproto"); @@ -499,11 +570,15 @@ test "delegation token and space credential round trip" { std.Options.debug_io, "did:plc:spaceauthority", delegation.space, + "test-dpop-thumbprint", + "https://client.example", &keypair, ); const credential = try verifySpaceCredential(allocator, credential_token, public_key_multibase); try std.testing.expectEqualStrings("did:plc:spaceauthority", credential.authority_did); try std.testing.expectEqualStrings(delegation.space, credential.space); + try std.testing.expectEqualStrings("test-dpop-thumbprint", credential.dpop_jkt); + try std.testing.expectEqualStrings("https://client.example", credential.client_id.?); try std.testing.expectError(error.InvalidJwt, verifySpaceCredential(allocator, delegation_token, public_key_multibase)); try std.testing.expectError(error.InvalidJwt, verifyDelegationToken(allocator, credential_token, public_key_multibase)); diff --git a/src/storage/store.zig b/src/storage/store.zig index 84a26b2..83dbfe2 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -366,8 +366,8 @@ pub const SpaceWriterState = struct { }; pub const CredentialRecipient = struct { + service_id: []const u8, service_endpoint: []const u8, - repo_did: ?[]const u8, expires_at: i64, }; @@ -2817,6 +2817,60 @@ pub fn permissionedBlobReferenced(space: []const u8, repo_did: []const u8, cid: return false; } +pub fn listPermissionedBlobs( + allocator: std.mem.Allocator, + space: []const u8, + repo_did: []const u8, + since: ?[]const u8, + cursor: ?[]const u8, + limit: usize, +) ![][]const u8 { + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + try requireInitialized(); + const capped_limit: i64 = @intCast(@min(if (limit == 0) 500 else limit, 1000)); + var rows = if (since) |rev| + if (cursor) |after| + try conn.rows( + \\SELECT DISTINCT b.blob_cid + \\FROM permissioned_space_record_blobs b + \\JOIN permissioned_space_records r + \\ ON r.space = b.space AND r.repo_did = b.repo_did + \\ AND r.collection = b.collection AND r.rkey = b.rkey + \\WHERE b.space = ? AND b.repo_did = ? AND r.repo_rev > ? AND b.blob_cid > ? + \\ORDER BY b.blob_cid ASC LIMIT ? + , .{ space, repo_did, rev, after, capped_limit }) + else + try conn.rows( + \\SELECT DISTINCT b.blob_cid + \\FROM permissioned_space_record_blobs b + \\JOIN permissioned_space_records r + \\ ON r.space = b.space AND r.repo_did = b.repo_did + \\ AND r.collection = b.collection AND r.rkey = b.rkey + \\WHERE b.space = ? AND b.repo_did = ? AND r.repo_rev > ? + \\ORDER BY b.blob_cid ASC LIMIT ? + , .{ space, repo_did, rev, capped_limit }) + else if (cursor) |after| + try conn.rows( + \\SELECT DISTINCT blob_cid + \\FROM permissioned_space_record_blobs + \\WHERE space = ? AND repo_did = ? AND blob_cid > ? + \\ORDER BY blob_cid ASC LIMIT ? + , .{ space, repo_did, after, capped_limit }) + else + try conn.rows( + \\SELECT DISTINCT blob_cid + \\FROM permissioned_space_record_blobs + \\WHERE space = ? AND repo_did = ? + \\ORDER BY blob_cid ASC LIMIT ? + , .{ space, repo_did, capped_limit }); + defer rows.deinit(); + var cids: std.ArrayList([]const u8) = .empty; + while (rows.next()) |row| try cids.append(allocator, try allocator.dupe(u8, row.text(0))); + if (rows.err) |err| return err; + return cids.toOwnedSlice(allocator); +} + pub fn getPublicBlob( allocator: std.mem.Allocator, did: []const u8, @@ -3863,59 +3917,46 @@ pub fn purgeAuthoritySpaceData(space: []const u8) !void { try conn.commit(); } -pub fn registerSpaceNotification(space: []const u8, repo_did: ?[]const u8, service_endpoint: []const u8, expires_at: i64) !void { +pub fn registerSpaceNotification(space: []const u8, service_id: []const u8, service_endpoint: []const u8, expires_at: i64) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.exec( - \\INSERT INTO permissioned_space_notify_registrations (space, repo_did, service_endpoint, expires_at) + \\INSERT INTO permissioned_space_notify_registrations (space, service_id, service_endpoint, expires_at) \\VALUES (?, ?, ?, ?) - \\ON CONFLICT(space, repo_did, service_endpoint) DO UPDATE SET + \\ON CONFLICT(space, service_id) DO UPDATE SET + \\ service_endpoint = excluded.service_endpoint, \\ expires_at = excluded.expires_at - , .{ space, repo_did orelse "", service_endpoint, expires_at }); + , .{ space, service_id, service_endpoint, expires_at }); +} + +pub fn unregisterSpaceNotification(space: []const u8, service_id: []const u8) !void { + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + try requireInitialized(); + try conn.exec( + "DELETE FROM permissioned_space_notify_registrations WHERE space = ? AND service_id = ?", + .{ space, service_id }, + ); } -pub fn listNotificationRecipients(allocator: std.mem.Allocator, space: []const u8, repo_did: ?[]const u8, include_space_wide: bool) ![]CredentialRecipient { +pub fn listNotificationRecipients(allocator: std.mem.Allocator, space: []const u8) ![]CredentialRecipient { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); try conn.execNoArgs("DELETE FROM permissioned_space_notify_registrations WHERE expires_at <= unixepoch()"); - var rows = if (repo_did) |repo| - if (include_space_wide) - try conn.rows( - \\SELECT repo_did, service_endpoint, expires_at - \\FROM permissioned_space_notify_registrations - \\WHERE space = ? AND (repo_did = '' OR repo_did = ?) - \\ORDER BY repo_did ASC, service_endpoint ASC - , .{ space, repo }) - else - try conn.rows( - \\SELECT repo_did, service_endpoint, expires_at - \\FROM permissioned_space_notify_registrations - \\WHERE space = ? AND repo_did = ? - \\ORDER BY service_endpoint ASC - , .{ space, repo }) - else if (include_space_wide) - try conn.rows( - \\SELECT '', service_endpoint, max(expires_at) - \\FROM permissioned_space_notify_registrations - \\WHERE space = ? - \\GROUP BY service_endpoint - \\ORDER BY service_endpoint ASC - , .{space}) - else - try conn.rows( - \\SELECT repo_did, service_endpoint, expires_at - \\FROM permissioned_space_notify_registrations - \\WHERE space = ? AND repo_did = '' - \\ORDER BY service_endpoint ASC - , .{space}); + var rows = try conn.rows( + \\SELECT service_id, service_endpoint, expires_at + \\FROM permissioned_space_notify_registrations + \\WHERE space = ? + \\ORDER BY service_id ASC + , .{space}); defer rows.deinit(); var out: std.ArrayList(CredentialRecipient) = .empty; while (rows.next()) |row| { try out.append(allocator, .{ + .service_id = try allocator.dupe(u8, row.text(0)), .service_endpoint = try allocator.dupe(u8, row.text(1)), - .repo_did = if (row.text(0).len == 0) null else try allocator.dupe(u8, row.text(0)), .expires_at = row.int(2), }); } @@ -4536,6 +4577,7 @@ fn migratePermissionedDataTables() !void { try migratePermissionedSpaceRepoHashes(); try migratePermissionedSpaceWriters(); try migratePermissionedSpaceNotifyRegistrations(); + try migratePermissionedSpaceNotificationServices(); } fn migratePermissionedSpaceBlobRefs() !void { @@ -4592,6 +4634,7 @@ fn migratePermissionedSpaceNotifyRegistrations() !void { \\CREATE TABLE IF NOT EXISTS permissioned_space_notify_registrations ( \\ space TEXT NOT NULL, \\ repo_did TEXT NOT NULL DEFAULT '', + \\ service_id TEXT NOT NULL DEFAULT '', \\ service_endpoint TEXT NOT NULL, \\ expires_at INTEGER NOT NULL, \\ PRIMARY KEY (space, repo_did, service_endpoint) @@ -4600,6 +4643,22 @@ fn migratePermissionedSpaceNotifyRegistrations() !void { try markMigrationApplied(name); } +fn migratePermissionedSpaceNotificationServices() !void { + const name = "permissioned-space-notify-service-id"; + if (try migrationApplied(name)) return; + try conn.execNoArgs("DROP TABLE IF EXISTS permissioned_space_notify_registrations"); + try conn.execNoArgs( + \\CREATE TABLE permissioned_space_notify_registrations ( + \\ space TEXT NOT NULL, + \\ service_id TEXT NOT NULL, + \\ service_endpoint TEXT NOT NULL, + \\ expires_at INTEGER NOT NULL, + \\ PRIMARY KEY (space, service_id) + \\) + ); + try markMigrationApplied(name); +} + fn migratePermissionedSpaceUris() !void { const name = "permissioned-space-at-uri"; if (try migrationApplied(name)) return; @@ -6642,10 +6701,10 @@ const schema_statements = [_][*:0]const u8{ , \\CREATE TABLE IF NOT EXISTS permissioned_space_notify_registrations ( \\ space TEXT NOT NULL, - \\ repo_did TEXT NOT NULL DEFAULT '', + \\ service_id TEXT NOT NULL, \\ service_endpoint TEXT NOT NULL, \\ expires_at INTEGER NOT NULL, - \\ PRIMARY KEY (space, repo_did, service_endpoint) + \\ PRIMARY KEY (space, service_id) \\) , }; @@ -7350,9 +7409,9 @@ test "permissioned space foundation migration preserves rows and rebuilds hashes \\VALUES (?, ?, ?, ?, ?) , .{ canonical, did, "fm.example.note", "one", "bafkreiaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" }); try conn.exec( - \\INSERT INTO permissioned_space_notify_registrations (space, repo_did, service_endpoint, expires_at) - \\VALUES (?, '', ?, 4102444800) - , .{ canonical, "https://sync.example" }); + \\INSERT INTO permissioned_space_notify_registrations (space, service_id, service_endpoint, expires_at) + \\VALUES (?, ?, ?, 4102444800) + , .{ canonical, "did:web:sync.example", "https://sync.example" }); try conn.execNoArgs("PRAGMA foreign_keys = OFF"); inline for (.{ @@ -7399,7 +7458,7 @@ test "permissioned space foundation migration preserves rows and rebuilds hashes try std.testing.expectEqualSlices(u8, &expected.digest(), writers[0].hash); } -test "permissioned notification registrations are scoped and expire" { +test "permissioned notification registrations are service keyed and expire" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const allocator = arena.allocator(); @@ -7407,17 +7466,50 @@ test "permissioned notification registrations are scoped and expire" { defer close(); const space = "at://did:plc:authority/space/fm.example.private/self"; - try registerSpaceNotification(space, null, "https://whole.example", 4102444800); - try registerSpaceNotification(space, "did:plc:writer", "https://repo.example", 4102444800); - try registerSpaceNotification(space, null, "https://expired.example", 1); - - const whole = try listNotificationRecipients(allocator, space, null, false); - try std.testing.expectEqual(@as(usize, 1), whole.len); - try std.testing.expectEqualStrings("https://whole.example", whole[0].service_endpoint); - const repo = try listNotificationRecipients(allocator, space, "did:plc:writer", true); - try std.testing.expectEqual(@as(usize, 2), repo.len); - try std.testing.expectEqualStrings("https://whole.example", repo[0].service_endpoint); - try std.testing.expectEqualStrings("https://repo.example", repo[1].service_endpoint); + try registerSpaceNotification(space, "did:web:whole.example#syncer", "https://whole.example", 4102444800); + try registerSpaceNotification(space, "did:web:repo.example#syncer", "https://repo.example", 4102444800); + try registerSpaceNotification(space, "did:web:expired.example#syncer", "https://expired.example", 1); + + const whole = try listNotificationRecipients(allocator, space); + try std.testing.expectEqual(@as(usize, 2), whole.len); + try std.testing.expectEqualStrings("did:web:repo.example#syncer", whole[0].service_id); + try std.testing.expectEqualStrings("https://repo.example", whole[0].service_endpoint); + try std.testing.expectEqualStrings("did:web:whole.example#syncer", whole[1].service_id); + try std.testing.expectEqualStrings("https://whole.example", whole[1].service_endpoint); + + try unregisterSpaceNotification(space, "did:web:whole.example#syncer"); + try unregisterSpaceNotification(space, "did:web:whole.example#syncer"); + const after_unregister = try listNotificationRecipients(allocator, space); + try std.testing.expectEqual(@as(usize, 1), after_unregister.len); +} + +test "permissioned notification migration replaces endpoint registrations" { + 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(); + + try conn.execNoArgs("DROP TABLE permissioned_space_notify_registrations"); + try conn.execNoArgs( + \\CREATE TABLE permissioned_space_notify_registrations ( + \\ space TEXT NOT NULL, + \\ repo_did TEXT NOT NULL DEFAULT '', + \\ service_endpoint TEXT NOT NULL, + \\ expires_at INTEGER NOT NULL, + \\ PRIMARY KEY (space, repo_did, service_endpoint) + \\) + ); + try conn.execNoArgs("INSERT INTO permissioned_space_notify_registrations VALUES ('at://did:plc:old/space/fm.old/x', '', 'https://old.example', 4102444800)"); + try conn.exec("DELETE FROM zds_migrations WHERE name = ?", .{"permissioned-space-notify-service-id"}); + try migratePermissionedSpaceNotificationServices(); + + try registerSpaceNotification("at://did:plc:new/space/fm.new/x", "did:web:sync.example", "https://sync.example", 4102444800); + const registrations = try listNotificationRecipients(allocator, "at://did:plc:new/space/fm.new/x"); + try std.testing.expectEqual(@as(usize, 1), registrations.len); + const old = (try conn.row("SELECT count(*) FROM permissioned_space_notify_registrations WHERE service_endpoint = 'https://old.example'", .{})).?; + defer old.deinit(); + try std.testing.expectEqual(@as(i64, 0), old.int(0)); } test "permissioned credential exchange tokens are one use" { diff --git a/tools/smoke-permissioned.sh b/tools/smoke-permissioned.sh index 56afa50..0838985 100755 --- a/tools/smoke-permissioned.sh +++ b/tools/smoke-permissioned.sh @@ -28,6 +28,9 @@ rm -rf "$blob_root" mkdir -p "$blob_root" zig build +dpop_proof() { + ./zig-out/bin/zds-space-dpop "$@" 2>&1 +} ZDS_PERMISSIONED_DATA=true zig build run -- \ --host 127.0.0.1 \ --port "$port" \ @@ -109,16 +112,42 @@ printf '%s' "$space_get" | grep -q '"policy":"member-list"' delegation=$(curl -fsS -H "authorization: Bearer $token" "$base/xrpc/com.atproto.space.getDelegationToken?space=$encoded_space" | jq -r '.token') test -n "$delegation" + +# A missing issuance proof must not consume the one-use delegation grant. +space_credential_missing_proof=$(curl -sS -o /tmp/zds-space-credential-missing-proof.json -w '%{http_code}' -X POST "$base/xrpc/com.atproto.space.getSpaceCredential" \ + -H "authorization: Bearer $delegation" \ + -H 'content-type: application/json' \ + --data "$(jq -nc --arg space "$space_uri" '{space:$space}')") +test "$space_credential_missing_proof" = "401" +grep -q '"error":"InvalidDpopProof"' /tmp/zds-space-credential-missing-proof.json + +credential_exchange_url="$base/xrpc/com.atproto.space.getSpaceCredential" +issuance_proof=$(dpop_proof POST "$credential_exchange_url") space_credential_status=$(curl -sS -o /tmp/zds-space-credential.json -w '%{http_code}' -X POST "$base/xrpc/com.atproto.space.getSpaceCredential" \ -H "authorization: Bearer $delegation" \ + -H "dpop: $issuance_proof" \ -H 'content-type: application/json' \ --data "$(jq -nc --arg space "$space_uri" '{space:$space}')") test "$space_credential_status" = "200" space_credential=$(jq -r '.credential' /tmp/zds-space-credential.json) test -n "$space_credential" -credential_space_get=$(curl -fsS -H "authorization: Bearer $space_credential" "$base/xrpc/com.atproto.space.getSpace?space=$encoded_space") + +credential_bearer_status=$(curl -sS -o /tmp/zds-space-credential-bearer.json -w '%{http_code}' \ + -H "authorization: Bearer $space_credential" \ + "$base/xrpc/com.atproto.space.getSpace?space=$encoded_space") +test "$credential_bearer_status" = "401" + +space_get_url="$base/xrpc/com.atproto.space.getSpace" +space_get_proof=$(dpop_proof GET "$space_get_url" "$space_credential") +credential_space_get=$(curl -fsS -H "authorization: DPoP $space_credential" -H "dpop: $space_get_proof" "$base/xrpc/com.atproto.space.getSpace?space=$encoded_space") printf '%s' "$credential_space_get" | grep -q '"uri":"at://did:plc:permissionsmoke/space/fm.plyr.privateMedia/self"' +credential_replay_status=$(curl -sS -o /tmp/zds-space-credential-replay.json -w '%{http_code}' \ + -H "authorization: DPoP $space_credential" -H "dpop: $space_get_proof" \ + "$base/xrpc/com.atproto.space.getSpace?space=$encoded_space") +test "$credential_replay_status" = "401" +grep -q '"error":"InvalidDpopProof"' /tmp/zds-space-credential-replay.json + space_list=$(curl -fsS -H "authorization: Bearer $token" "$base/xrpc/com.atproto.space.listSpaces?did=did:plc:permissionsmoke&type=fm.plyr.privateMedia") printf '%s' "$space_list" | grep -q '"uri":"at://did:plc:permissionsmoke/space/fm.plyr.privateMedia/self"' printf '%s' "$space_list" | grep -q '"isOwner":true' @@ -142,6 +171,26 @@ public_blob_status=$(curl -sS -o /tmp/zds-space-public-blob.bin -w '%{http_code} "$base/xrpc/com.atproto.sync.getBlob?did=did:plc:permissionsmoke&cid=$blob_cid") test "$public_blob_status" = "404" +space_blobs=$(curl -fsS -H "authorization: Bearer $token" "$base/xrpc/com.atproto.space.listBlobs?space=$encoded_space&repo=did:plc:permissionsmoke") +printf '%s' "$space_blobs" | grep -q "\"$blob_cid\"" +space_rev=$(curl -fsS -H "authorization: Bearer $token" "$base/xrpc/com.atproto.space.getRepoState?space=$encoded_space&repo=did:plc:permissionsmoke" | jq -r '.commit.rev') +space_blobs_since=$(curl -fsS -H "authorization: Bearer $token" "$base/xrpc/com.atproto.space.listBlobs?space=$encoded_space&repo=did:plc:permissionsmoke&since=$space_rev") +test "$space_blobs_since" = '{"cids":[]}' +space_blobs_bad_limit=$(curl -sS -o /tmp/zds-space-blobs-bad-limit.json -w '%{http_code}' \ + -H "authorization: Bearer $token" \ + "$base/xrpc/com.atproto.space.listBlobs?space=$encoded_space&repo=did:plc:permissionsmoke&limit=0") +test "$space_blobs_bad_limit" = "400" +grep -q '"error":"InvalidRequest"' /tmp/zds-space-blobs-bad-limit.json + +unregister_url="$base/xrpc/com.atproto.space.unregisterNotify" +unregister_proof=$(dpop_proof POST "$unregister_url" "$space_credential") +curl -fsS -X POST "$unregister_url" \ + -H "authorization: DPoP $space_credential" \ + -H "dpop: $unregister_proof" \ + -H 'content-type: application/json' \ + --data "$(jq -nc --arg space "$space_uri" '{space:$space,service:"did:web:syncer.example"}')" \ + | grep -q '{}' + space_records=$(curl -fsS -H "authorization: Bearer $token" "$base/xrpc/com.atproto.space.listRecords?space=$encoded_space&repo=did:plc:permissionsmoke&collection=fm.plyr.track&limit=1") printf '%s' "$space_records" | grep -q '"collection":"fm.plyr.track"' printf '%s' "$space_records" | grep -q '"value":' diff --git a/tools/space_dpop.zig b/tools/space_dpop.zig new file mode 100644 index 0000000..c169e64 --- /dev/null +++ b/tools/space_dpop.zig @@ -0,0 +1,28 @@ +const std = @import("std"); +const zds = @import("zds"); + +const Io = std.Io; + +var app_threaded_io: Io.Threaded = undefined; +pub const std_options_debug_threaded_io: ?*Io.Threaded = &app_threaded_io; + +pub fn main(init: std.process.Init) !void { + const allocator = std.heap.smp_allocator; + app_threaded_io = Io.Threaded.init(allocator, .{}); + const io = app_threaded_io.io(); + var args = std.process.Args.Iterator.init(init.minimal.args); + _ = args.next(); + const method = args.next() orelse return usage(); + const htu = args.next() orelse return usage(); + const credential = args.next(); + + var keypair = try zds.zat.Keypair.fromSecretKey(.p256, .{0x42} ** 32); + const ath = if (credential) |token| try zds.zat.oauth.accessTokenHash(allocator, token) else null; + const proof = try zds.zat.oauth.createDpopProof(allocator, io, &keypair, method, htu, null, ath); + std.debug.print("{s}\n", .{proof}); +} + +fn usage() error{InvalidArguments} { + std.debug.print("usage: zds-space-dpop METHOD HTU [CREDENTIAL]\n", .{}); + return error.InvalidArguments; +} -- 2.51.2