diff --git a/README.md b/README.md index 7896584..bd6b71b 100644 --- a/README.md +++ b/README.md @@ -1,27 +1,93 @@ # linji -atproto study-community client. Headless core (`src/core/`) + thin CLI shell. +The CLI for 林集 (linji.at) — an atproto study-community client. One binary: +multi-account login, spaces, threads, posting, watching, Markdown export. +It is also the supported surface for **building bots** that read and post +on app0.linji.at. -M001: one `linji` binary, multiple logged-in accounts, one topic on local zds. +## Build & install -## quickstart +Requires Zig 0.16 (`zvm`/`zigup` recommended) and a C toolchain. ```sh -# terminal 1: sandbox PDS (fresh, wipes state) + alice/bob/tongxi accounts -tools/sandbox.sh - -# terminal 2 -zig build -zig build run -- login alice.test --password alice-password -zig build run -- login bob.test --password bob-password -zig build run -- post --topic "课文讨论" "大家好" -zig build run -- post --as bob.test --topic "课文讨论" "你好" -zig build run -- read -zig build run -- watch +zig build # -> zig-out/bin/linji +zig build test # unit tests +zig build install --prefix ~/.local # -> ~/.local/bin/linji ``` -Tests: `zig build test` (unit), `zig build scenario` (M001 end-to-end; -needs `../zds` built — run `tools/sandbox.sh` once first, or -`(cd ../zds && zig build)`). +Dev scenarios (need `../zds` built — `cd ../zds && zig build +-Dtarget=x86_64-linux-gnu.2.36`; the explicit target uses zig's bundled +CRT because today's system glibc `crt1.o` carries `.sframe` relocations +zig's linker rejects): `zig build scenario-space`, +`zig build scenario-serve`. + +## Accounts + +Any atproto handle works. To **post in spaces hosted on app0.linji.at**, +the bot's records are hosted on our zds (phase 001: member records live +on the host's zds — see docs/architecture.dj "External members"), and an +admin must add the account to the space's member list: + +```sh +linji space add-member --space at:///space/at.linji.space/ --did +``` + +Login is OAuth. Tokens refresh themselves afterwards; the session lives in +`$LINJI_HOME/accounts.json` (default `~/.local/share/linji`). + +```sh +linji login --pds # prints an authorize URL, opens consent +``` + +**Headless servers:** log in once on a machine with a browser, then copy +`accounts.json` to the server (`chmod 600` — it carries the DPoP key and +refresh token). Refresh is automatic from then on. + +## Commands (bot-relevant) + +```sh +linji thread list --space # thread AT-URIs + titles +linji thread create --space --title # -> thread AT-URI +linji post --space --thread # -> post AT-URI +linji read --space [--tag ] # merged thread view +linji watch --space [--interval ms] # follow new posts +linji edit|delete --space --post … +linji export --space --thread # thread as Markdown +``` + +`--as ` picks a non-default account. `--tag` filters to threads +containing an inline hashtag. + +## A minimal bot + +The loop: `watch` for new posts, reply with `post`. Here reacting to a +`#question` tag: + +```sh +#!/bin/sh +SPACE="at://did:plc:community/space/at.linji.space/general" +THREAD="at://did:plc:bot/at.linji.thread/…" # reply thread + +linji watch --space "$SPACE" --interval 3000 | while IFS= read -r line; do + case "$line" in + *"#question"*) linji post --space "$SPACE" --thread "$THREAD" \ + "收到 — 让我想想。" ;; + esac +done +``` + +For machine-readable reads, the BFF serves `GET /api/space?uri=` +(JSON: threads + posts) once you hold a session cookie; the CLI's +rendered output is the stable interface for shell bots. + +Run it under your service manager. app0/dinit/ has examples for the dev +stack. + +## Dev rig + +```sh +dinit -d app0/dinit dev # zds + seed + BFF + vite (see app0/dinit/dev) +# seeded accounts: community/alice/bob (.test, password -password) +``` Layout, architecture, roadmap: see `../docs/`. diff --git a/src/core/accounts.zig b/src/core/accounts.zig index b40d744..8062902 100644 --- a/src/core/accounts.zig +++ b/src/core/accounts.zig @@ -174,7 +174,7 @@ fn mkdirPath(allocator: std.mem.Allocator, path: []const u8) !void { try mkdirOne(path, &buf); } -fn writeFile600(allocator: std.mem.Allocator, path: []const u8, data: []const u8) !void { +pub fn writeFile600(allocator: std.mem.Allocator, path: []const u8, data: []const u8) !void { const path_z = try allocator.dupeZ(u8, path); defer allocator.free(path_z); const file = std.c.fopen(path_z.ptr, "wb") orelse return error.OpenFailed; diff --git a/src/core/push.zig b/src/core/push.zig new file mode 100644 index 0000000..3c4ada0 --- /dev/null +++ b/src/core/push.zig @@ -0,0 +1,370 @@ +//! Web Push: RFC 8030 (push resource), RFC 8291 (aes128gcm message +//! encryption), RFC 8292 (VAPID). One VAPID keypair per server; +//! subscriptions are stored per account DID and fired when a post lands +//! in a space the subscriber belongs to. +//! +//! State under {LINJI_HOME} (see accounts.zig): +//! vapid.key base64url 32-byte P-256 secret (created once) +//! push-subscriptions.json {"subscriptions":[{endpoint,p256dh,auth,did}]} +//! +//! VAPID subject is fixed mailto:admin@linji.at — contact for push +//! service operators; revisit if we ever run someone else's BFF. + +const std = @import("std"); +const zat = @import("zat"); +const accounts = @import("accounts.zig"); + +const P256 = std.crypto.ecc.P256; +const Ecdsa = std.crypto.sign.ecdsa.EcdsaP256Sha256; +const Hmac = std.crypto.auth.hmac.sha2.HmacSha256; +const Aes = std.crypto.aead.aes_gcm.Aes128Gcm; + +const vapid_subject = "mailto:admin@linji.at"; +const jwt_header_b64 = "eyJhbGciOiJFUzI1NiIsInR5cCI6IkpXVCJ9"; // {"alg":"ES256","typ":"JWT"} + +fn hmac(key: []const u8, msg: []const u8) [32]u8 { + var out: [32]u8 = undefined; + Hmac.create(&out, msg, key); + return out; +} + +/// One push subscription as delivered by the browser. +pub const Subscription = struct { + endpoint: []const u8, + /// base64url, 65-byte uncompressed user-agent P-256 public key + p256dh: []const u8, + /// base64url, 16-byte auth secret + auth: []const u8, + /// owning account DID + did: []const u8, +}; + +/// JSON-file subscription store. Same single-writer assumption as +/// accounts.Store. +pub const Store = struct { + allocator: std.mem.Allocator, + path: []const u8, + subs: std.ArrayList(Subscription) = .empty, + + pub fn load(allocator: std.mem.Allocator, io: std.Io) !Store { + const dir = try accounts.Store.homeDir(allocator); + const path = try std.fs.path.join(allocator, &.{ dir, "push-subscriptions.json" }); + var store = Store{ .allocator = allocator, .path = path }; + const data = std.Io.Dir.cwd().readFileAlloc(io, path, allocator, .limited(1 << 20)) catch |err| switch (err) { + error.FileNotFound => return store, + else => return err, + }; + const parsed = try std.json.parseFromSlice(std.json.Value, allocator, data, .{}); + defer parsed.deinit(); + if (zat.json.getArray(parsed.value, "subscriptions")) |subs| { + for (subs) |s| { + const endpoint = zat.json.getString(s, "endpoint") orelse continue; + const p256dh = zat.json.getString(s, "p256dh") orelse continue; + const auth = zat.json.getString(s, "auth") orelse continue; + const did = zat.json.getString(s, "did") orelse continue; + try store.subs.append(allocator, .{ + .endpoint = try allocator.dupe(u8, endpoint), + .p256dh = try allocator.dupe(u8, p256dh), + .auth = try allocator.dupe(u8, auth), + .did = try allocator.dupe(u8, did), + }); + } + } + return store; + } + + /// Upsert by endpoint (a browser re-subscribing replaces its row). + pub fn add(self: *Store, sub: Subscription) !void { + for (self.subs.items) |*s| { + if (std.mem.eql(u8, s.endpoint, sub.endpoint)) { + s.* = sub; + return self.save(); + } + } + try self.subs.append(self.allocator, sub); + return self.save(); + } + + /// Remove by endpoint. Returns whether a row was removed. + pub fn remove(self: *Store, endpoint: []const u8) !bool { + for (self.subs.items, 0..) |s, i| { + if (std.mem.eql(u8, s.endpoint, endpoint)) { + _ = self.subs.orderedRemove(i); + try self.save(); + return true; + } + } + return false; + } + + fn save(self: *Store) !void { + var aw: std.Io.Writer.Allocating = .init(self.allocator); + defer aw.deinit(); + const w = &aw.writer; + try w.writeAll("{\"subscriptions\":["); + for (self.subs.items, 0..) |s, i| { + if (i > 0) try w.writeAll(","); + try w.print("{{\"endpoint\":{f},\"p256dh\":{f},\"auth\":{f},\"did\":{f}}}", .{ + std.json.fmt(s.endpoint, .{}), + std.json.fmt(s.p256dh, .{}), + std.json.fmt(s.auth, .{}), + std.json.fmt(s.did, .{}), + }); + } + try w.writeAll("]}"); + try accounts.writeFile600(self.allocator, self.path, aw.written()); + } +}; + +/// The server's VAPID identity (RFC 8292). The public key goes to +/// browsers as `applicationServerKey`; the secret signs the JWT each +/// push request carries. +pub const Vapid = struct { + secret: [32]u8, + + pub fn loadOrCreate(allocator: std.mem.Allocator, io: std.Io) !Vapid { + const dir = try accounts.Store.homeDir(allocator); + const path = try std.fs.path.join(allocator, &.{ dir, "vapid.key" }); + if (std.Io.Dir.cwd().readFileAlloc(io, path, allocator, .limited(256))) |data| { + const raw = try zat.jwt.base64UrlDecode(allocator, std.mem.trim(u8, data, " \n")); + if (raw.len != 32) return error.BadVapidKey; + return .{ .secret = raw[0..32].* }; + } else |err| switch (err) { + error.FileNotFound => {}, + else => return err, + } + const kp = Ecdsa.KeyPair.generate(io); + const b64 = try zat.jwt.base64UrlEncode(allocator, &kp.secret_key.bytes); + try accounts.writeFile600(allocator, path, b64); + return .{ .secret = kp.secret_key.bytes }; + } + + fn keyPair(self: *const Vapid) !Ecdsa.KeyPair { + return Ecdsa.KeyPair.fromSecretKey(.{ .bytes = self.secret }); + } + + /// base64url uncompressed public key — the browser's applicationServerKey. + pub fn publicKeyB64(self: *const Vapid, allocator: std.mem.Allocator) ![]u8 { + const kp = try self.keyPair(); + return zat.jwt.base64UrlEncode(allocator, &kp.public_key.p.toUncompressedSec1()); + } + + /// Sign a VAPID JWT for the given push-service origin (RFC 8292 §2). + fn jwt(self: *const Vapid, allocator: std.mem.Allocator, io: std.Io, aud: []const u8) ![]u8 { + const now = @divFloor(std.Io.Timestamp.now(io, .real).nanoseconds, std.time.ns_per_s); + const payload = try std.fmt.allocPrint(allocator, "{{\"aud\":{f},\"exp\":{d},\"sub\":{f}}}", .{ + std.json.fmt(aud, .{}), + now + 12 * 3600, + std.json.fmt(vapid_subject, .{}), + }); + const payload_b64 = try zat.jwt.base64UrlEncode(allocator, payload); + const signing_string = try std.fmt.allocPrint(allocator, "{s}.{s}", .{ jwt_header_b64, payload_b64 }); + const sig = try zat.jwt.signP256(signing_string, &self.secret); + const sig_b64 = try zat.jwt.base64UrlEncode(allocator, &sig.bytes); + return std.fmt.allocPrint(allocator, "{s}.{s}", .{ signing_string, sig_b64 }); + } +}; + +/// RFC 8291 §3.3/3.4 + RFC 8188: encrypt one record (aes128gcm) for the +/// user agent identified by ua_pub (SEC1 uncompressed) and auth_secret. +/// Returns the full body: salt(16) || rs(4,BE) || idlen(1) || keyid(65) +/// || ciphertext || tag(16). +pub fn encrypt( + allocator: std.mem.Allocator, + io: std.Io, + ua_pub: []const u8, + auth_secret: []const u8, + plaintext: []const u8, +) ![]u8 { + const eph = Ecdsa.KeyPair.generate(io); + const ua_point = try P256.fromSec1(ua_pub); + const shared = try ua_point.mulPublic(eph.secret_key.bytes, .big); + const ecdh_secret = shared.affineCoordinates().x.toBytes(.big); + + // RFC 8291 §3.3: IKM from the auth secret + const prk_key = hmac(auth_secret, &ecdh_secret); + const ikm = hmac(&prk_key, "Content-Encoding: auth\x00\x01"); + + var salt: [16]u8 = undefined; + io.random(&salt); + const prk = hmac(&salt, &ikm); + const cek_full = hmac(&prk, "Content-Encoding: aes128gcm\x00\x01"); + const nonce_full = hmac(&prk, "Content-Encoding: nonce\x00\x01"); + const cek: [16]u8 = cek_full[0..16].*; + const nonce: [12]u8 = nonce_full[0..12].*; + + // Single final record: plaintext || 0x02 delimiter, no padding. + const padded = try allocator.alloc(u8, plaintext.len + 1); + @memcpy(padded[0..plaintext.len], plaintext); + padded[plaintext.len] = 0x02; + const ct = try allocator.alloc(u8, padded.len); + var tag: [16]u8 = undefined; + Aes.encrypt(ct, &tag, padded, "", nonce, cek); + + const eph_pub = eph.public_key.p.toUncompressedSec1(); + var body: std.Io.Writer.Allocating = .init(allocator); + const w = &body.writer; + try w.writeAll(&salt); + try w.writeInt(u32, 4096, .big); // rs, must exceed record size + try w.writeByte(65); + try w.writeAll(&eph_pub); + try w.writeAll(ct); + try w.writeAll(&tag); + return body.written(); +} + +/// Decrypt helper used by tests: the inverse of encrypt(), recovering +/// keys from the body header and the UA secret key. +fn decryptForTest( + allocator: std.mem.Allocator, + ua_secret: [32]u8, + auth_secret: []const u8, + body: []const u8, +) ![]u8 { + const salt = body[0..16]; + const idlen = body[20]; + const eph_pub = body[21 .. 21 + @as(usize, idlen)]; + const ct = body[21 + idlen ..]; + const shared = try (try P256.fromSec1(eph_pub)).mulPublic(ua_secret, .big); + const ecdh_secret = shared.affineCoordinates().x.toBytes(.big); + const prk_key = hmac(auth_secret, &ecdh_secret); + const ikm = hmac(&prk_key, "Content-Encoding: auth\x00\x01"); + const prk = hmac(salt, &ikm); + const cek_full = hmac(&prk, "Content-Encoding: aes128gcm\x00\x01"); + const nonce_full = hmac(&prk, "Content-Encoding: nonce\x00\x01"); + const cek: [16]u8 = cek_full[0..16].*; + const nonce: [12]u8 = nonce_full[0..12].*; + const out = try allocator.alloc(u8, ct.len - 16); + const tag: [16]u8 = ct[ct.len - 16 ..][0..16].*; + try Aes.decrypt(out, ct[0 .. ct.len - 16], tag, "", nonce, cek); + return out; +} + +/// Push-service origin (scheme://host[:port]) — the VAPID `aud`. +fn originOf(endpoint: []const u8) ![]const u8 { + const scheme_end = std.mem.indexOf(u8, endpoint, "://") orelse return error.BadEndpoint; + const rest = endpoint[scheme_end + 3 ..]; + const path_start = std.mem.indexOfScalar(u8, rest, '/') orelse rest.len; + return endpoint[0 .. scheme_end + 3 + path_start]; +} + +/// Result of one delivery attempt. +pub const SendResult = enum { delivered, gone, failed }; + +/// Encrypt `payload` for `sub` and POST it to its push endpoint. +/// `gone` (404/410) means the subscription is dead — caller removes it. +pub fn send( + allocator: std.mem.Allocator, + io: std.Io, + vapid: *const Vapid, + sub: Subscription, + payload: []const u8, +) !SendResult { + const ua_pub = try zat.jwt.base64UrlDecode(allocator, sub.p256dh); + const auth = try zat.jwt.base64UrlDecode(allocator, sub.auth); + const body = try encrypt(allocator, io, ua_pub, auth, payload); + const token = try vapid.jwt(allocator, io, try originOf(sub.endpoint)); + const pub_b64 = try vapid.publicKeyB64(allocator); + const authorization = try std.fmt.allocPrint(allocator, "vapid t={s}, k={s}", .{ token, pub_b64 }); + + var transport = zat.HttpTransport.init(io, allocator); + defer transport.deinit(); + const extra = [_]std.http.Header{ + .{ .name = "content-encoding", .value = "aes128gcm" }, + .{ .name = "ttl", .value = "86400" }, + }; + var result = try transport.fetch(.{ + .url = sub.endpoint, + .method = .POST, + .payload = body, + .content_type = "application/octet-stream", + .authorization = authorization, + .extra_headers = &extra, + }); + defer result.deinit(allocator); + const status = result.status; + if (status == .created or status == .ok) return .delivered; + if (status == .not_found or status == .gone) return .gone; + std.debug.print("push: {s} -> {d}\n", .{ originOf(sub.endpoint) catch "?", @intFromEnum(status) }); + return .failed; +} + +/// Notify every subscriber among `member_dids` except `exclude_did`. +/// Dead subscriptions are pruned. Synchronous — push services answer in +/// well under a second and member counts are small; revisit with a queue +/// if that ever stops being true. +pub fn notifyMembers( + allocator: std.mem.Allocator, + io: std.Io, + vapid: *const Vapid, + store: *Store, + member_dids: []const []const u8, + exclude_did: []const u8, + payload: []const u8, +) !void { + for (store.subs.items) |sub| { + if (std.mem.eql(u8, sub.did, exclude_did)) continue; + var is_member = false; + for (member_dids) |d| { + if (std.mem.eql(u8, d, sub.did)) { + is_member = true; + break; + } + } + if (!is_member) continue; + const result = send(allocator, io, vapid, sub, payload) catch |err| { + std.debug.print("push: send failed: {s}\n", .{@errorName(err)}); + continue; + }; + if (result == .gone) _ = try store.remove(sub.endpoint); + } +} + +test "encrypt/decrypt roundtrip" { + const allocator = std.testing.allocator; + const io = std.testing.io; + const ua = Ecdsa.KeyPair.generate(io); + const auth = "auth-secret-1234"; + const plaintext = "你好,共修"; + const body = try encrypt(allocator, io, &ua.public_key.p.toUncompressedSec1(), auth, plaintext); + defer allocator.free(body); + // header: salt(16) || rs(4) || idlen(1)=65 || eph(65) = 86 bytes + try std.testing.expect(body.len == 86 + plaintext.len + 1 + 16); + try std.testing.expectEqual(@as(u32, 4096), std.mem.readInt(u32, body[16..20], .big)); + const recovered = try decryptForTest(allocator, ua.secret_key.bytes, auth, body); + defer allocator.free(recovered); + try std.testing.expectEqual(plaintext.len + 1, recovered.len); + try std.testing.expectEqualStrings(plaintext, recovered[0..plaintext.len]); + try std.testing.expectEqual(@as(u8, 0x02), recovered[plaintext.len]); +} + +test "store roundtrip" { + const allocator = std.testing.allocator; + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + const dir_path = try tmp.dir.realpathAlloc(allocator, "."); + defer allocator.free(dir_path); + var store = Store{ + .allocator = allocator, + .path = try std.fs.path.join(allocator, &.{ dir_path, "push-subscriptions.json" }), + }; + const sub = Subscription{ + .endpoint = "https://push.example/abc", + .p256dh = "cDI2NGRo", + .auth = "YXV0aA", + .did = "did:plc:alice", + }; + try store.add(sub); + try std.testing.expectEqual(@as(usize, 1), store.subs.items.len); + try store.add(sub); // upsert, not duplicate + try std.testing.expectEqual(@as(usize, 1), store.subs.items.len); + try std.testing.expect(try store.remove("https://push.example/abc")); + try std.testing.expectEqual(@as(usize, 0), store.subs.items.len); + try std.testing.expect(!try store.remove("https://push.example/abc")); +} + +test "originOf" { + try std.testing.expectEqualStrings("https://fcm.googleapis.com", try originOf("https://fcm.googleapis.com/fcm/send/abc")); + try std.testing.expectEqualStrings("https://push.services.mozilla.com", try originOf("https://push.services.mozilla.com/wpush/v2/xyz")); + try std.testing.expectError(error.BadEndpoint, originOf("not-a-url")); +} diff --git a/src/core/root.zig b/src/core/root.zig index 44bb0c1..bec0018 100644 --- a/src/core/root.zig +++ b/src/core/root.zig @@ -8,6 +8,7 @@ pub const hashtag = @import("hashtag.zig"); pub const loopback = @import("loopback.zig"); pub const oauth = @import("oauth.zig"); pub const pds = @import("pds.zig"); +pub const push = @import("push.zig"); pub const spaces = @import("spaces.zig"); pub const post = @import("post.zig"); pub const thread = @import("thread.zig"); diff --git a/src/server.zig b/src/server.zig index 8daeb32..9a46863 100644 --- a/src/server.zig +++ b/src/server.zig @@ -12,6 +12,7 @@ const linji = @import("linji"); const accounts = linji.accounts; const oauth = linji.oauth; const spaces = linji.spaces; +const push = linji.push; const App = struct { allocator: std.mem.Allocator, @@ -23,6 +24,8 @@ const App = struct { handle_domain: []const u8, invite_code: ?[]const u8, metadata: []const u8, + vapid: push.Vapid, + push_store: push.Store, mutex: std.Io.Mutex = .init, sessions: std.StringHashMapUnmanaged([]const u8) = .empty, // token -> did pending: std.StringHashMapUnmanaged(oauth.PendingLogin) = .empty, // state -> login @@ -92,6 +95,8 @@ pub fn serve( .handle_domain = handle_domain, .invite_code = invite_code, .metadata = metadata, + .vapid = try push.Vapid.loadOrCreate(allocator, io), + .push_store = try push.Store.load(allocator, io), }; var server = try httpz.Server(*App).init(io, allocator, .{ .address = .{ .ip = .{ .ip4 = try std.Io.net.Ip4Address.parse("127.0.0.1", port) } }, @@ -114,6 +119,8 @@ pub fn serve( router.post("/api/delete", deleteHandler, .{}); router.post("/api/move", moveHandler, .{}); router.get("/api/export", exportHandler, .{}); + router.post("/api/push/subscribe", pushSubscribeHandler, .{}); + router.post("/api/push/unsubscribe", pushUnsubscribeHandler, .{}); std.debug.print("linji serve listening on {s}\n", .{base_url}); try server.listen(); } @@ -128,10 +135,12 @@ fn clientMetadataHandler(app: *App, _: *httpz.Request, res: *httpz.Response) !vo } /// What the login page needs before any session exists. -fn configHandler(app: *App, _: *httpz.Request, res: *httpz.Response) !void { - const body = try std.fmt.allocPrint(res.arena, "{{\"handleDomain\":{f},\"inviteRequired\":{s}}}", .{ +fn configHandler(app: *App, req: *httpz.Request, res: *httpz.Response) !void { + const vapid_key = try app.vapid.publicKeyB64(req.arena); + const body = try std.fmt.allocPrint(res.arena, "{{\"handleDomain\":{f},\"inviteRequired\":{s},\"vapidKey\":{f}}}", .{ std.json.fmt(app.handle_domain, .{}), if (app.invite_code != null) "true" else "false", + std.json.fmt(vapid_key, .{}), }); try jsonBody(res, body); } @@ -369,6 +378,7 @@ fn writeOpHandler( const uri = spaces.createPost(core, app.io, app.store, account, space, thread_uri, text) catch return denied(res); const body = try std.fmt.allocPrint(req.arena, "{{\"uri\":{f}}}", .{std.json.fmt(uri, .{})}); try jsonBody(res, body); + notifySpaceMembers(app, core, space, account, thread_uri, text); }, .thread => { const title = zat.json.getString(parsed.value, "title") orelse return missing(res, "title"); @@ -422,6 +432,69 @@ fn moveHandler(app: *App, req: *httpz.Request, res: *httpz.Response) !void { try writeOpHandler(app, req, res, .move); } +/// Browser registers its push subscription (PushManager.subscribe JSON). +fn pushSubscribeHandler(app: *App, req: *httpz.Request, res: *httpz.Response) !void { + const account = app.sessionAccount(req) orelse { + res.status = 401; + res.body = "not logged in"; + return; + }; + const parsed = parseJsonBody(req, req.arena) orelse { + res.status = 400; + res.body = "invalid json body"; + return; + }; + const endpoint = zat.json.getString(parsed.value, "endpoint") orelse return missing(res, "endpoint"); + const keys = parsed.value.object.get("keys") orelse return missing(res, "keys"); + const p256dh = zat.json.getString(keys, "p256dh") orelse return missing(res, "p256dh"); + const auth = zat.json.getString(keys, "auth") orelse return missing(res, "auth"); + // Store outlives the request: dupe into the app allocator. + try app.push_store.add(.{ + .endpoint = try app.allocator.dupe(u8, endpoint), + .p256dh = try app.allocator.dupe(u8, p256dh), + .auth = try app.allocator.dupe(u8, auth), + .did = try app.allocator.dupe(u8, account.did), + }); + res.body = "ok"; +} + +fn pushUnsubscribeHandler(app: *App, req: *httpz.Request, res: *httpz.Response) !void { + _ = app.sessionAccount(req) orelse { + res.status = 401; + res.body = "not logged in"; + return; + }; + const parsed = parseJsonBody(req, req.arena) orelse { + res.status = 400; + res.body = "invalid json body"; + return; + }; + const endpoint = zat.json.getString(parsed.value, "endpoint") orelse return missing(res, "endpoint"); + _ = try app.push_store.remove(endpoint); + res.body = "ok"; +} + +/// After a post lands: push every subscribed space member except the +/// author. Runs before the handler returns (push services answer fast, +/// member counts are small — see push.notifyMembers). +fn notifySpaceMembers(app: *App, core: std.mem.Allocator, space: []const u8, author: *accounts.Account, thread_uri: []const u8, text: []const u8) void { + const members = spaces.listMembers(core, app.io, app.store, author, space) catch return; + const preview = truncateUtf8(text, 80); + const label = std.fmt.allocPrint(core, "{s}: {s}", .{ author.handle, preview }) catch return; + const payload = std.fmt.allocPrint(core, "{{\"title\":\"林集\",\"body\":{f},\"tag\":{f}}}", .{ + std.json.fmt(label, .{}), + std.json.fmt(thread_uri, .{}), + }) catch return; + push.notifyMembers(core, app.io, &app.vapid, &app.push_store, members, author.did, payload) catch {}; +} + +/// First `max` bytes of `s`, shrunk to a UTF-8 codepoint boundary. +fn truncateUtf8(s: []const u8, max: usize) []const u8 { + var end = @min(s.len, max); + while (end > 0 and end < s.len and (s[end] & 0xC0) == 0x80) end -= 1; + return s[0..end]; +} + fn exportHandler(app: *App, req: *httpz.Request, res: *httpz.Response) !void { var core_impl = std.heap.ArenaAllocator.init(std.heap.page_allocator); defer core_impl.deinit(); diff --git a/tests/scenario_serve.zig b/tests/scenario_serve.zig index f7c9bba..8f97c14 100644 --- a/tests/scenario_serve.zig +++ b/tests/scenario_serve.zig @@ -113,6 +113,7 @@ pub fn main(init: std.process.Init) !void { try expectCode(res.code, 200, "/api/config"); try expectContains(res.body, "\"handleDomain\":\"test\"", "config domain"); try expectContains(res.body, "\"inviteRequired\":true", "config invite"); + try expectContains(res.body, "\"vapidKey\":\"", "config vapid key"); res = try api(allocator, io, null, "GET", "/api/login?handle=alice.test&password=alice-password", null); try expectCode(res.code, 403, "login without invite denied"); @@ -175,6 +176,42 @@ pub fn main(init: std.process.Init) !void { try expectCode(res.code, 200, "GET /"); try expectContains(res.body, "linji pwa", "PWA served"); + // --- 6b. web push --- + // Subscribe requires a session. + const push_endpoint = api_base ++ "/api/no-such-push-endpoint"; // our own API 404s it -> 'gone' + const p256dh = "BGsX0fLhLEJH-Lzm5WOkQPJ3A32BLeszoPShOUXYmMKWT-NC4v4af5uO5-tKfA-eFivOM1drMV7Oy7ZAaDe_UfU"; // P-256 base point, valid test key + const sub_body = try std.fmt.allocPrint(allocator, + "{{\"endpoint\":\"{s}\",\"keys\":{{\"p256dh\":\"{s}\",\"auth\":\"MDEyMzQ1Njc4OWFiY2RlZg\"}}}}", .{ push_endpoint, p256dh }); + res = try api(allocator, io, null, "POST", "/api/push/subscribe", sub_body); + try expectCode(res.code, 401, "subscribe without session"); + res = try api(allocator, io, jar_alice, "POST", "/api/push/subscribe", sub_body); + try expectCode(res.code, 200, "alice subscribe"); + res = try api(allocator, io, jar_alice, "POST", "/api/push/subscribe", sub_body); + try expectCode(res.code, 200, "alice re-subscribe (upsert)"); + // The store file records it (one row despite the upsert). + const subs_path = try std.fs.path.join(allocator, &.{ home, "push-subscriptions.json" }); + var subs_data = try std.Io.Dir.cwd().readFileAlloc(io, subs_path, allocator, .limited(1 << 20)); + try expectContains(subs_data, push_endpoint, "subscription stored"); + if (std.mem.count(u8, subs_data, "endpoint") != 2) { // "endpoint" key + value once each + std.debug.print("FAIL: expected one subscription row, got:\n{s}\n", .{subs_data}); + return error.ScenarioFailed; + } + // bob posts -> server pushes alice's subscription synchronously -> our + // own API 404s the endpoint -> 'gone' -> the row is pruned. + res = try api(allocator, io, jar_bob, "POST", "/api/post", try std.fmt.allocPrint(allocator, + "{{\"space\":\"{s}\",\"thread\":\"{s}\",\"text\":\"触发推送\"}}", .{ space, thread_uri })); + try expectCode(res.code, 200, "bob post triggers push"); + subs_data = try std.Io.Dir.cwd().readFileAlloc(io, subs_path, allocator, .limited(1 << 20)); + if (std.mem.indexOf(u8, subs_data, push_endpoint) != null) { + std.debug.print("FAIL: dead subscription not pruned:\n{s}\n", .{subs_data}); + return error.ScenarioFailed; + } + std.debug.print("ok dead subscription pruned after push 404\n", .{}); + // unsubscribe is idempotent. + res = try api(allocator, io, jar_alice, "POST", "/api/push/unsubscribe", try std.fmt.allocPrint(allocator, + "{{\"endpoint\":\"{s}\"}}", .{push_endpoint})); + try expectCode(res.code, 200, "unsubscribe"); + // --- 7. logout revokes --- res = try api(allocator, io, jar_alice, "POST", "/api/logout", null); try expectCode(res.code, 200, "logout"); diff --git a/web/public/sw.js b/web/public/sw.js index b39ce35..7085020 100644 --- a/web/public/sw.js +++ b/web/public/sw.js @@ -2,7 +2,29 @@ // else -> network-first with cache fallback (deploys must not strand the // shell on dead asset hashes). /api is always network. Bump VERSION to // roll clients. -const VERSION = "linji-v2"; +// Web Push: every push MUST show a visible notification — silent pushes +// burn Firefox's quota and get the subscription revoked. +const VERSION = "linji-v3"; + +self.addEventListener("push", (e) => { + const data = e.data ? e.data.json() : {}; + e.waitUntil( + self.registration.showNotification(data.title || "林集", { + body: data.body || "", + tag: data.tag, + icon: "/icon.svg", + }), + ); +}); + +self.addEventListener("notificationclick", (e) => { + e.notification.close(); + e.waitUntil( + clients + .matchAll({ type: "window", includeUncontrolled: true }) + .then((list) => (list.length > 0 ? list[0].focus() : clients.openWindow("/"))), + ); +}); self.addEventListener("install", () => self.skipWaiting()); diff --git a/web/src/App.tsx b/web/src/App.tsx index 601b3e8..87c01b0 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -31,6 +31,7 @@ export function App() { +
{spaces().map((uri) => ( @@ -79,6 +80,76 @@ export function App() { ); } +/** base64url (VAPID key) -> Uint8Array for applicationServerKey. */ +function urlB64ToBytes(b64: string): Uint8Array { + const pad = "=".repeat((4 - (b64.length % 4)) % 4); + const raw = atob(b64.replace(/-/g, "+").replace(/_/g, "/") + pad); + const out = new Uint8Array(raw.length); + for (let i = 0; i < raw.length; i++) out[i] = raw.charCodeAt(i); + return out; +} + +/** Notifications toggle. "on" = an active push subscription lives on the + * server; clicking unsubscribes. Silent pushes are never sent (Firefox + * quota) — the SW shows a notification for every push. + */ +function PushBell() { + const [state, setState] = createSignal<"off" | "on" | "denied" | "unsupported">("off"); + + onMount(async () => { + if (!("serviceWorker" in navigator) || !("PushManager" in window)) { + setState("unsupported"); + return; + } + if (Notification.permission === "denied") { + setState("denied"); + return; + } + const reg = await navigator.serviceWorker.ready; + const sub = await reg.pushManager.getSubscription(); + setState(sub ? "on" : "off"); + }); + + const toggle = async () => { + const reg = await navigator.serviceWorker.ready; + const existing = await reg.pushManager.getSubscription(); + if (existing) { + await existing.unsubscribe(); + await api.pushUnsubscribe(existing.endpoint); + setState("off"); + return; + } + if (Notification.permission !== "granted") { + const perm = await Notification.requestPermission(); + if (perm !== "granted") { + setState("denied"); + return; + } + } + const config = await api.config(); + if (!config.vapidKey) return; + const sub = await reg.pushManager.subscribe({ + userVisibleOnly: true, + applicationServerKey: urlB64ToBytes(config.vapidKey) as BufferSource, + }); + await api.pushSubscribe(sub.toJSON()); + setState("on"); + }; + + return ( + + + + ); +} + function Login(props: { onLogin: (m: Me) => void }) { const [name, setName] = createSignal(""); const [password, setPassword] = createSignal(""); diff --git a/web/src/api.ts b/web/src/api.ts index ed6fe2b..d966e97 100644 --- a/web/src/api.ts +++ b/web/src/api.ts @@ -4,6 +4,7 @@ export interface Config { handleDomain: string; inviteRequired: boolean; + vapidKey?: string; } export interface Me { @@ -61,6 +62,18 @@ export const api = { headers: { "content-type": "application/json" }, body: JSON.stringify({ space, thread, text }), }), + pushSubscribe: (sub: unknown) => + req("/api/push/subscribe", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify(sub), + }), + pushUnsubscribe: (endpoint: string) => + req("/api/push/unsubscribe", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ endpoint }), + }), }; /** Spaces the user has joined on this device. Membership lives in zds;