From d4f1a196929eacd7a9b64b94b324d7653cbec46c Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sat, 15 Aug 2026 00:59:25 -0500 Subject: [PATCH] fix: announce new accounts on the firehose; add invite CTA to signup createAccount previously wrote only the accounts row, so a new account emitted no firehose events and stayed invisible to relays until its first write. It now matches the official PDS and Tranquil sequencing: #identity, #account, an empty genesis #commit, and #sync, then a crawler notify. smoke asserts the genesis commit is servable and that exactly four events exist after signup. ZDS_OPERATOR_HANDLE / --operator-handle names the PDS operator; when set, the signup page's invite section links their Bluesky profile as the call to action for getting an invite code. Co-Authored-By: Claude Fable 5 --- README.md | 1 + docs/operations.md | 2 ++ fly.toml | 1 + src/atproto/server.zig | 6 ++++++ src/core/config.zig | 11 +++++++++++ src/internal/cli.zig | 8 ++++++++ src/internal/signup.zig | 40 ++++++++++++++++++++++++++++++++++++++-- src/main.zig | 1 + src/storage/store.zig | 16 +++++++++++++++- tools/smoke.sh | 11 +++++++++++ 10 files changed, 94 insertions(+), 3 deletions(-) diff --git a/README.md b/README.md index 4f36ecd..8558995 100644 --- a/README.md +++ b/README.md @@ -67,6 +67,7 @@ ZDS_BLOB_UPLOAD_LIMIT=100000000 \ ZDS_BLOBSTORE_PATH=/var/lib/zds/blobs \ ZDS_HANDLE_DOMAINS='.pds.example.com' \ ZDS_CRAWLERS='https://bsky.network,https://vsky.network' \ +ZDS_OPERATOR_HANDLE='operator.example.com' \ ZDS_MAX_CONCURRENT_REPO_EXPORTS=4 \ ZDS_PLC_ROTATION_KEY='64-hex-secp256k1-secret-or-private-multikey' \ ZDS_RECOVERY_DID_KEY='did:key:optionalRecoveryKey' \ diff --git a/docs/operations.md b/docs/operations.md index 79c24ae..ef9e8cd 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -133,6 +133,8 @@ Common deployment settings: - `ZDS_BLOB_UPLOAD_LIMIT`: upload body limit. Default: `100000000`. - `ZDS_BLOBSTORE_PATH`: disk blobstore root. - `ZDS_CRAWLERS`: comma-separated relay crawl targets. +- `ZDS_OPERATOR_HANDLE`: atproto handle of the person operating this PDS, + shown on the public signup page as the contact for invite codes. - `ZDS_PLC_DIRECTORY`: PLC directory origin. Default: `https://plc.directory`. - `ZDS_PROXY_SERVICE_URL`, `ZDS_PROXY_SERVICE_DID`, `ZDS_PROXY_SERVICE_ID`: default appview target for `atproto-proxy` forwarding. Defaults: diff --git a/fly.toml b/fly.toml index 6c558e0..d8b4da5 100644 --- a/fly.toml +++ b/fly.toml @@ -12,6 +12,7 @@ primary_region = "ord" ZDS_SERVER_DID = "did:web:pds.zat.dev" ZDS_HANDLE_DOMAINS = ".pds.zat.dev" ZDS_CRAWLERS = "https://bsky.network,https://vsky.network" + ZDS_OPERATOR_HANDLE = "zzstoatzz.io" ZDS_BLOB_UPLOAD_LIMIT = "100000000" ZDS_BLOBSTORE_PATH = "/data/blobs" ZDS_MAIL_PROVIDER = "comail" diff --git a/src/atproto/server.zig b/src/atproto/server.zig index 906da65..3ed3196 100644 --- a/src/atproto/server.zig +++ b/src/atproto/server.zig @@ -239,6 +239,12 @@ pub fn createAccount(request: *http_api.Request) !void { }; }; log.debug("xrpc createAccount ok did={s} handle={s}\n", .{ account.did, account.handle }); + if (existing_did == null) { + store.sequenceNewAccount(allocator, account) catch |err| { + log.err("xrpc createAccount failed sequence_new_account did={s} err={s}\n", .{ account.did, @errorName(err) }); + }; + sync.notifyCrawlers(true); + } const issued = try issueSessionTokens(allocator, account, "account_create", null); const body_out = try sessionJson(allocator, account, issued.access.token, issued.refresh.token); diff --git a/src/core/config.zig b/src/core/config.zig index e3d1ce0..e60b144 100644 --- a/src/core/config.zig +++ b/src/core/config.zig @@ -14,6 +14,7 @@ var blob_upload_limit_value: usize = 100_000_000; var blobstore_path_value: []const u8 = "dev/blobs"; var handle_domains_value: []const u8 = ".test"; var crawlers_value: []const u8 = "https://bsky.network,https://vsky.network"; +var operator_handle_value: []const u8 = ""; var max_concurrent_repo_exports_value: usize = 4; var proxy_service_did_value: []const u8 = "did:web:api.bsky.app"; var proxy_service_id_value: []const u8 = "bsky_appview"; @@ -92,6 +93,12 @@ pub fn crawlers() []const u8 { return crawlers_value; } +/// atproto handle of the person operating this PDS, or empty when unset. +/// Shown on the public signup page as the invite contact. +pub fn operatorHandle() []const u8 { + return operator_handle_value; +} + pub fn maxConcurrentRepoExports() usize { return max_concurrent_repo_exports_value; } @@ -192,6 +199,10 @@ pub fn setCrawlers(value: []const u8) void { crawlers_value = value; } +pub fn setOperatorHandle(value: []const u8) void { + operator_handle_value = value; +} + pub fn setMaxConcurrentRepoExports(value: usize) void { max_concurrent_repo_exports_value = value; } diff --git a/src/internal/cli.zig b/src/internal/cli.zig index ca51837..b6f1d18 100644 --- a/src/internal/cli.zig +++ b/src/internal/cli.zig @@ -18,6 +18,7 @@ pub const Options = struct { blobstore_path: ?[]const u8 = null, handle_domains: ?[]const u8 = null, crawlers: ?[]const u8 = null, + operator_handle: ?[]const u8 = null, max_concurrent_repo_exports: ?usize = null, proxy_service_did: ?[]const u8 = null, proxy_service_id: ?[]const u8 = null, @@ -36,6 +37,7 @@ pub const ParseError = error{ MissingBlobUploadLimit, MissingBlobstorePath, MissingCrawlers, + MissingOperatorHandle, MissingMaxConcurrentRepoExports, MissingProxyServiceDid, MissingProxyServiceId, @@ -80,6 +82,7 @@ pub fn parse(init: std.process.Init.Minimal) ParseError!Options { .blobstore_path = env("ZDS_BLOBSTORE_PATH"), .handle_domains = env("ZDS_HANDLE_DOMAINS"), .crawlers = env("ZDS_CRAWLERS"), + .operator_handle = env("ZDS_OPERATOR_HANDLE"), .max_concurrent_repo_exports = try envUsize("ZDS_MAX_CONCURRENT_REPO_EXPORTS"), .proxy_service_did = env("ZDS_PROXY_SERVICE_DID"), .proxy_service_id = env("ZDS_PROXY_SERVICE_ID"), @@ -203,6 +206,10 @@ fn parseSplitArg(options: *Options, arg: []const u8, args: *std.process.Args.Ite options.crawlers = args.next() orelse return error.MissingCrawlers; return true; } + if (std.mem.eql(u8, arg, "--operator-handle")) { + options.operator_handle = args.next() orelse return error.MissingOperatorHandle; + return true; + } if (std.mem.eql(u8, arg, "--max-concurrent-repo-exports")) { options.max_concurrent_repo_exports = try std.fmt.parseInt(usize, args.next() orelse return error.MissingMaxConcurrentRepoExports, 10); return true; @@ -281,6 +288,7 @@ const joined_string_options = [_]JoinedStringOption{ .{ .flag = "--blobstore-path=", .field = "blobstore_path" }, .{ .flag = "--handle-domains=", .field = "handle_domains" }, .{ .flag = "--crawlers=", .field = "crawlers" }, + .{ .flag = "--operator-handle=", .field = "operator_handle" }, .{ .flag = "--proxy-service-did=", .field = "proxy_service_did" }, .{ .flag = "--proxy-service-id=", .field = "proxy_service_id" }, .{ .flag = "--proxy-service-url=", .field = "proxy_service_url" }, diff --git a/src/internal/signup.zig b/src/internal/signup.zig index 87684a5..f343aec 100644 --- a/src/internal/signup.zig +++ b/src/internal/signup.zig @@ -1,10 +1,43 @@ const std = @import("std"); +const config = @import("../core/config.zig"); const http_api = @import("../http/api.zig"); const http = std.http; pub fn page(request: *http_api.Request) !void { - try http_api.respond(request, .ok, html, &html_headers); + var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); + defer arena.deinit(); + const body = try std.mem.concat(arena.allocator(), u8, &.{ + html_before_invite_cta, + try inviteCta(arena.allocator()), + html_after_invite_cta, + }); + try http_api.respond(request, .ok, body, &html_headers); +} + +fn inviteCta(allocator: std.mem.Allocator) ![]const u8 { + const operator = config.operatorHandle(); + if (operator.len == 0) return ""; + const escaped = try escapeHtml(allocator, operator); + return std.fmt.allocPrint( + allocator, + "

Need an invite? DM @{s}, the operator of this PDS.

", + .{ escaped, escaped }, + ); +} + +fn escapeHtml(allocator: std.mem.Allocator, text: []const u8) ![]const u8 { + var out: std.Io.Writer.Allocating = .init(allocator); + defer out.deinit(); + for (text) |c| switch (c) { + '&' => try out.writer.writeAll("&"), + '<' => try out.writer.writeAll("<"), + '>' => try out.writer.writeAll(">"), + '"' => try out.writer.writeAll("""), + '\'' => try out.writer.writeAll("'"), + else => try out.writer.writeByte(c), + }; + return out.toOwnedSlice(); } const html_headers = [_]http.Header{ @@ -15,7 +48,7 @@ const html_headers = [_]http.Header{ .{ .name = "connection", .value = "close" }, }; -const html = +const html_before_invite_cta = \\ \\ \\ @@ -63,6 +96,9 @@ const html = \\ \\ \\

This server requires an invite. If someone sent you a signup link, the code is already filled in.

+; + +const html_after_invite_cta = \\ \\

This is an experimental personal data server. Your account and repo live here; you can migrate them elsewhere later.

\\ diff --git a/src/main.zig b/src/main.zig index 4bab668..0447837 100644 --- a/src/main.zig +++ b/src/main.zig @@ -43,6 +43,7 @@ pub fn main(init: std.process.Init) !void { zds.core.config.setHandleDomains(value); } if (options.crawlers) |value| zds.core.config.setCrawlers(value); + if (options.operator_handle) |value| zds.core.config.setOperatorHandle(value); if (options.max_concurrent_repo_exports) |value| zds.core.config.setMaxConcurrentRepoExports(value); if (options.proxy_service_did) |value| zds.core.config.setProxyServiceDid(value); if (options.proxy_service_id) |value| zds.core.config.setProxyServiceId(value); diff --git a/src/storage/store.zig b/src/storage/store.zig index 01133a3..1302743 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -436,6 +436,9 @@ pub const ValidationMode = enum { pub const WriteOptions = struct { swap_commit: ?[]const u8 = null, validate: ValidationMode = .known, + /// Permit a commit with no record operations. Only the account-creation + /// genesis commit uses this; XRPC applyWrites keeps rejecting empty writes. + allow_empty_commit: bool = false, }; const BlobRef = struct { @@ -1211,6 +1214,17 @@ pub fn sequenceIdentityEvent(allocator: std.mem.Allocator, did: []const u8, hand eventlog.publish(seq); } +/// Firehose announcement for a freshly created account, matching the official +/// PDS sequencer order: #identity, #account, genesis #commit, #sync. The +/// genesis commit is an empty-ops write so consumers see the repo exist +/// before its first record. +pub fn sequenceNewAccount(allocator: std.mem.Allocator, account: auth.Account) !void { + try sequenceIdentityEvent(allocator, account.did, account.handle); + try sequenceAccountEvent(allocator, account.did, .active); + _ = try applyWritesWithOptions(allocator, account, &.{}, .{ .allow_empty_commit = true }); + try sequenceSyncEvent(allocator, account.did); +} + pub fn sequenceSyncEvent(allocator: std.mem.Allocator, did: []const u8) !void { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); @@ -2006,7 +2020,7 @@ pub fn applyWritesProfiled(allocator: std.mem.Allocator, account: auth.Account, fn applyWritesMeasured(allocator: std.mem.Allocator, account: auth.Account, ops: []const WriteOp, options: WriteOptions, profile: ?*WriteProfile) !WriteResult { const total_start = monotonicNs(); - if (ops.len == 0) return Error.MissingRecord; + if (ops.len == 0 and !options.allow_empty_commit) return Error.MissingRecord; const validation_start = monotonicNs(); for (ops) |op| switch (op) { diff --git a/tools/smoke.sh b/tools/smoke.sh index e42f504..ff231ba 100755 --- a/tools/smoke.sh +++ b/tools/smoke.sh @@ -55,6 +55,7 @@ ZDS_PLC_ROTATION_KEY=11111111111111111111111111111111111111111111111111111111111 --handle-domains .test \ --invite-required \ --admin-token smoke-admin-token \ + --operator-handle smoke-operator.test \ --plc-directory "http://127.0.0.1:${plc_port}" \ >"$log" 2>&1 & server_pid=$! @@ -127,6 +128,7 @@ signup_page=$(curl -fsS "$base/signup") printf '%s' "$signup_page" | grep -q 'join zds' printf '%s' "$signup_page" | grep -q 'com.atproto.server.describeServer' printf '%s' "$signup_page" | grep -q 'com.atproto.server.createAccount' +printf '%s' "$signup_page" | grep -q 'DM