From c2de0765c3935e6b1b8c9f95838ef95119636e05 Mon Sep 17 00:00:00 2001 From: sepehr-safari Date: Wed, 23 Sep 2026 15:22:44 +0300 Subject: [PATCH 1/2] feat: publish says what it published Each event that at least one relay accepted is printed on stdout once the relays have answered or its deadline has passed, so what publish writes out is what was published. Every answer on stderr carries the event's id, and an acceptance shows the relay's message, so a relay that already had the event says so. The exit code covers every event: 0 only when every event was accepted by at least one relay, 1 when any was not, and the rest of the stream still goes out. Before, one acceptance anywhere in the run was enough to exit 0. Events are checked before they are sent, so one whose signature does not match its content is never offered. The relays are dialled at once under the same bound as req and fetch, sent each event at once, and given one absolute deadline per event to take it and answer, 10 seconds by default and --timeout to change it. A relay still being sent to at the deadline has stopped reading and is dropped; its reader is stopped first, because a send can be waiting for the write lock behind that reader's pong rather than for the socket. Each relay has its own reader for the whole run. It answers pings, hands answers and notices over through a queue, and reports the moment its connection goes away. So an answer is counted as soon as it arrives whichever relay sends it, a relay that stalls in the middle of a message holds up only its own reader, and a relay that closed while publish waited for input is dialled again before the next event. An event sent down a connection that then closes before answering is offered once more on a fresh one, when there is time left to do it. Control characters in text a relay sends are shown as escapes, C1 included, and notices are shown up to eight per relay across its connections. More than 32 relays is refused, and req and fetch now say when they leave relays out. The help for all three says that looking up a relay's name is the one step that cannot be cut short. The nostr pin moves to 0.14.5, where pings and pongs arriving back to back no longer hold a read past its deadline. Tests drive the command against real relays on loopback: accepting, refusing, already having the event, never answering the dial, never reading, pinging while never reading, closing after an answer, pinging and then closing, closing on the next event, answering about another event, answering behind a ping or seventy notices beside twenty silent relays, sending notices or pongs without answering, notices across reconnections, and trying to forge output. CI jobs now time out at 20 minutes, so a regression that brings back an unbounded wait fails the job rather than holding a runner. Closes #19. Closes #21. --- .github/workflows/ci.yml | 4 + README.md | 9 +- build.zig.zon | 4 +- src/cmd_fetch.zig | 1 + src/cmd_publish.zig | 1241 +++++++++++++++++++++++++++++++++++--- src/cmd_req.zig | 1 + src/relayset.zig | 28 + src/testrelay.zig | 140 ++++- 8 files changed, 1315 insertions(+), 113 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 08eb755..d37927c 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -13,6 +13,10 @@ jobs: matrix: os: [ubuntu-latest, macos-latest] runs-on: ${{ matrix.os }} + # The tests talk to relays over real sockets, so a regression that brings + # back an unbounded wait shows up as a hang. It should fail the job in + # minutes, not hold a runner for hours. + timeout-minutes: 20 steps: - uses: actions/checkout@v4 with: diff --git a/README.md b/README.md index 2432723..e03488f 100644 --- a/README.md +++ b/README.md @@ -19,7 +19,7 @@ macOS and Linux, Intel and ARM. It works out which build this machine wants, che The script is short and worth reading before you pipe anything into bash. If you would rather do it yourself: ```sh -VERSION=0.2.0 +VERSION=0.3.0 PLATFORM=macos-aarch64 # or macos-x86_64, linux-x86_64, linux-aarch64 BASE=https://github.com/zig-nostr/deed/releases/download/v$VERSION @@ -53,7 +53,7 @@ deed reaches relays and keeps what it finds. Every verb that does not need a soc | `verify` | check that events are correctly signed | | `req` | build a subscription, and run it | | `fetch` | get the events a code names | -| `publish` | offer signed events to relays | +| `publish` | offer signed events to relays, and print the ones they accepted | `deed help ` explains any of them. @@ -85,8 +85,11 @@ export NOSTR_SECRET_KEY=$(deed key generate) deed event -c "hello" | deed verify # builds one, signs it, checks it cat drafts.jsonl | deed event - | deed verify +cat drafts.jsonl | deed event - | deed publish wss://relay.example > sent.jsonl ``` +`publish` prints each event a relay accepted, once the relays have answered or the deadline has passed, so what it writes out is what was published. + A key can be passed with `--sec`, but a key on a command line lands in your shell history and in the process table, so `$NOSTR_SECRET_KEY` is the better habit. A bad record fails that record alone. The stream carries on, the reason goes to standard error, and the exit code reports that something in the run failed. @@ -98,7 +101,7 @@ Scripts branch on these, so they are part of the interface and not free to drift | | | | --- | --- | | `0` | it worked | -| `1` | the command ran and failed: a bad signature, an unreadable key, a malformed code | +| `1` | the command ran and failed: a bad signature, an unreadable key, a malformed code, an event no relay accepted | | `2` | the command was not understood: unknown verb, unknown flag, missing argument. Nothing was attempted | | `141` | the reader on the other end of the pipe went away, as in `deed decode … \| head -1` | diff --git a/build.zig.zon b/build.zig.zon index e5ae9c4..61ddc08 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -10,8 +10,8 @@ // themselves later, and "whatever main happened to be that afternoon" // is not that. .nostr = .{ - .url = "https://github.com/zig-nostr/nostr/archive/refs/tags/v0.14.4.tar.gz", - .hash = "nostr-0.14.4-CMyPze4BCQAUoi8JLNWfaR4ZDF1e4C5dhM2vMwGJ_qIT", + .url = "https://github.com/zig-nostr/nostr/archive/refs/tags/v0.14.5.tar.gz", + .hash = "nostr-0.14.5-CMyPzUgLCQBev0r22ht88wNHfIfc6vajpRq-hylCECfB", }, }, .paths = .{ diff --git a/src/cmd_fetch.zig b/src/cmd_fetch.zig index 04a929c..fbe3b6b 100644 --- a/src/cmd_fetch.zig +++ b/src/cmd_fetch.zig @@ -28,6 +28,7 @@ pub const usage = \\ \\A relay that has not accepted the connection within five seconds, or \\within --timeout if that is shorter, is named on stderr and left out. + \\Looking up a relay's name is the one step that cannot be cut short. \\ \\A bare npub fetches that person's profile rather than everything they have \\ever written, which is what `npub` on its own can sensibly mean. diff --git a/src/cmd_publish.zig b/src/cmd_publish.zig index 9d78fb2..1b42784 100644 --- a/src/cmd_publish.zig +++ b/src/cmd_publish.zig @@ -1,39 +1,65 @@ -//! `deed publish`: offer signed events to relays. +//! `deed publish`: offer signed events to relays, and say which were published. //! //! Reads events the way every other verb reads its records, so the thing that //! made them is somebody else's business: `deed event -c "hi" | deed publish //! wss://...` is the ordinary path, and a file of events works the same way. +//! +//! What comes out on stdout is what was published: each event that at least one +//! relay accepted, once the relays have answered or its deadline has passed. +//! Every answer goes to stderr with the event's id on it, and a notice with the +//! relay's URL. const std = @import("std"); const nostr = @import("nostr"); const cli = @import("cli.zig"); +const dial = @import("dial.zig"); +const hex = nostr.hex; const relay = nostr.relay; /// An event too big for this is one no relay would take either. const max_record_bytes = 1 << 20; -/// How long one relay gets to say yes or no. +/// How long the relays get to take one event and answer it, all of them +/// together. /// -/// A relay that accepts the frame and never answers is the case this exists -/// for. nak forces the same shape with a 7 second deadline on the OK and ten -/// seconds around the whole flow; this is the middle of those. -const ok_timeout_ms: i64 = 8_000; +/// An absolute deadline, not a wait that starts again with every message: a +/// relay that keeps sending notices and never answers would otherwise hold the +/// run for as long as it liked. +pub const default_timeout_ms: i64 = 10_000; pub const usage = \\deed publish: offer signed events to relays \\ \\Usage: - \\ deed publish ... [...] + \\ deed publish [--timeout ] ... [...] \\ \\Reads events as newline-delimited JSON on stdin when given no event - \\arguments, and offers each to every relay named. + \\arguments, checks each one, and offers it to every relay named. \\ \\ deed event -c "hello" | deed publish wss://relay.example \\ - \\Every relay's answer is reported on stderr. Exits 0 when at least one - \\relay accepted an event, because an event one relay holds is published; - \\exits 1 when none did, and 2 when the command made no sense. + \\Prints each event that at least one relay accepted, once the relays have + \\answered or the deadline has passed, so what comes out is what was + \\published: + \\ + \\ deed event - < drafts.jsonl | deed publish wss://a wss://b > sent.jsonl + \\ + \\Every relay's answer is reported on stderr with the event's id. An event + \\no relay accepted, or whose signature does not check out, is reported + \\there and not printed. + \\ + \\Options: + \\ --timeout how long the relays get to take each event and + \\ answer it (default 10000) + \\ + \\A relay that has not accepted the connection within five seconds, or + \\within --timeout if that is shorter, is named on stderr and left out. + \\Looking up a relay's name is the one step that cannot be cut short. + \\ + \\Exits 0 when every event was accepted by at least one relay, 1 when any + \\was not or there was nothing to publish, and 2 when the command made no + \\sense. \\ ; @@ -48,12 +74,31 @@ pub fn run( defer urls.deinit(gpa); var positionals: std.ArrayList([]const u8) = .empty; defer positionals.deinit(gpa); + var timeout_ms: i64 = default_timeout_ms; - for (args) |a| { + var i: usize = 0; + while (i < args.len) : (i += 1) { + const a = args[i]; if (cli.isOneOf(a, &.{ "help", "-h", "--help" })) { try out.writeAll(usage); return cli.exit_ok; } + if (std.mem.eql(u8, a, "--timeout")) { + i += 1; + if (i == args.len) { + try err.writeAll("deed publish: --timeout needs a number of milliseconds\n"); + return cli.exit_usage; + } + timeout_ms = std.fmt.parseInt(i64, args[i], 10) catch { + try err.print("deed publish: --timeout needs a number of milliseconds, not '{s}'\n", .{args[i]}); + return cli.exit_usage; + }; + if (timeout_ms <= 0) { + try err.writeAll("deed publish: --timeout must be more than zero\n"); + return cli.exit_usage; + } + continue; + } if (std.mem.startsWith(u8, a, "-")) { try err.print("deed publish: unknown option '{s}'\n", .{a}); return cli.exit_usage; @@ -70,32 +115,9 @@ pub fn run( try err.writeAll("deed publish: name at least one relay to publish to\n"); return cli.exit_usage; } - - // Dialled BEFORE anything is read, so a run that cannot reach a single - // relay says so without having consumed the events it was given. nak orders - // it the same way, for the sharper version of the reason: there, signing - // costs a round trip to a remote signer, and burning one to then discover - // no relay is reachable is a bad trade. - var conns: [32]?*relay.Relay = @splat(null); - const n = @min(urls.items.len, conns.len); - defer for (conns[0..n]) |maybe| { - if (maybe) |r| { - r.shutdown(io); - r.deinit(); - } - }; - - var live: usize = 0; - for (urls.items[0..n], 0..) |url, i| { - conns[i] = relay.dial(gpa, io, url) catch |e| { - try err.print("deed publish: {s}: {s}\n", .{ url, @errorName(e) }); - continue; - }; - live += 1; - } - if (live == 0) { - try err.writeAll("deed publish: no relay could be reached\n"); - return cli.exit_fail; + if (urls.items.len > dial.max_relays) { + try err.print("deed publish: name at most {d} relays\n", .{dial.max_relays}); + return cli.exit_usage; } const stdin_buf = try gpa.alloc(u8, max_record_bytes); @@ -103,93 +125,1124 @@ pub fn run( var stdin_reader = std.Io.File.stdin().readerStreaming(io, stdin_buf); var input = cli.Input.init(positionals.items, &stdin_reader.interface); - var accepted_any = false; - var offered: usize = 0; - var refusals: usize = 0; + var signer = nostr.keys.Signer.init(); + defer signer.deinit(); + + // Dialled when the first event that checks out arrives, not before. Empty + // or unreadable input has nothing to publish, and opens no connection. + var set: Set = undefined; + set.init(gpa, io, urls.items); + defer set.deinit(); + + // Every record is something the caller asked to have published, whether + // or not it turned out to be an event. + var records: usize = 0; + var published: usize = 0; while (try input.next()) |record| { + records += 1; const json = switch (record) { .line => |l| std.mem.trim(u8, l, " \t"), .too_long => { try err.print("deed publish: skipped an event longer than {d} bytes\n", .{max_record_bytes}); - refusals += 1; + try err.flush(); continue; }, }; var parsed = nostr.event.fromJson(gpa, json) catch |e| { try err.print("deed publish: not an event: {s}\n", .{@errorName(e)}); - refusals += 1; + try err.flush(); continue; }; defer parsed.deinit(); - offered += 1; - - for (conns[0..n], 0..) |maybe, i| { - const r = maybe orelse continue; - const url = urls.items[i]; - r.publish(parsed.value) catch |e| { - try err.print("deed publish: {s}: {s}\n", .{ url, @errorName(e) }); - refusals += 1; - continue; - }; - if (awaitOk(r, parsed.value.id, io, err, url)) accepted_any = true else refusals += 1; + const ev = parsed.value; + + const id = try hex.encode(gpa, &ev.id); + defer gpa.free(id); + + // Checked before anything is sent. A relay would refuse it, but the + // refusal would read as the relay's problem rather than the event's. + if (!(nostr.event.verify(gpa, signer, ev) catch false)) { + try err.print("deed publish: {s}: not sent, bad signature\n", .{id}); + try err.flush(); + continue; + } + + try set.refresh(err); + try set.ready(@min(dial.default_timeout_ms, timeout_ms), err); + if (set.live() == 0) { + try err.print("deed publish: {s}: not published, no relay could be reached\n", .{id}); + try err.flush(); + continue; } + + if (try set.offer(ev, id, timeout_ms, err)) { + const wire = try nostr.event.toJson(gpa, ev); + defer gpa.free(wire); + try out.print("{s}\n", .{wire}); + try out.flush(); + published += 1; + } else { + try err.print("deed publish: {s}: not published, no relay accepted it in time\n", .{id}); + } + try err.flush(); } - if (offered == 0) { + try set.reportAll(err); + + if (records == 0) { try err.writeAll("deed publish: nothing to publish\n"); return cli.exit_fail; } - // An event one relay holds is published. nak draws the line in the same - // place: it fails only when every relay failed, because a run that reported - // failure after reaching the network would have scripts retrying something - // that already happened. - if (!accepted_any) return cli.exit_fail; - if (refusals > 0) try err.print("deed publish: {d} relay answers were not an acceptance\n", .{refusals}); + if (published < records) { + try err.print("deed publish: {d} of {d} published\n", .{ published, records }); + return cli.exit_fail; + } return cli.exit_ok; } -/// Waits for this relay's answer about this event. +/// The relays one run publishes to. /// -/// Other messages keep arriving on the same socket while this waits, so -/// anything that is not the OK being waited for is passed over rather than -/// treated as an answer. -fn awaitOk( - r: *relay.Relay, - id: [32]u8, +/// Each connected relay has its own reader for the whole run. It answers pings, +/// hands answers and notices to this thread through `heard`, and says the +/// moment its connection goes away. So an answer is seen as soon as it arrives +/// whichever relay sends it, a relay that closes while the run waits for input +/// is known about before the next event is sent to it, and a relay that stalls +/// in the middle of a read holds up only its own reader. +const Set = struct { + gpa: std.mem.Allocator, io: std.Io, - err: *std.Io.Writer, - url: []const u8, -) bool { - _ = io; - const deadline = std.Io.Timeout{ .duration = .{ .raw = .fromMilliseconds(ok_timeout_ms), .clock = .awake } }; - while (true) { - var msg = (r.receiveTimeout(deadline) catch { - try_print(err, "deed publish: {s}: gave up waiting for an answer\n", .{url}); - return false; - }) orelse { - try_print(err, "deed publish: {s}: closed before answering\n", .{url}); - return false; + urls: []const []const u8, + conns: [dial.max_relays]?*relay.Relay, + state: [dial.max_relays]State, + readers: [dial.max_relays]?std.Io.Future(void), + /// Set before a reader is cancelled, so it hands nothing over on its way + /// out. A put after the cancel could otherwise wait on a full queue that + /// nobody is reading while this thread waits for the reader. + quiet: [dial.max_relays]std.atomic.Value(bool), + /// Which connection a message came from. Bumped whenever a relay is let + /// go, so what its old reader left in the queue is recognised and dropped. + gen: [dial.max_relays]u32, + /// Notices a relay has sent this run, across every connection to it. + /// Counted by the readers, so a relay sending them as fast as it can does + /// not fill the queue with lines nobody will see. + notices: [dial.max_relays]std.atomic.Value(u32), + heard: std.Io.Queue(Heard), + heard_buf: [256]Heard, + /// Which event a deadline marker in `heard` belongs to. + seq: u64, + + const State = enum { + /// Not dialled yet. + idle, + connected, + /// Was connected and went away. Dialled again for the next event: + /// a relay that drops an idle connection is still a relay. + closed, + /// Not tried again this run: it failed to dial, or it stopped + /// reading what was sent to it. + gone, + }; + + /// What a reader hands over. Text is copied out of the message, which the + /// reader frees, and is freed here once it has been dealt with. + const Heard = struct { + index: usize = 0, + gen: u32 = 0, + what: union(enum) { + ok: struct { id: [32]u8, accepted: bool, message: []u8 }, + notice: []u8, + more_notices, + /// The connection went away: null when the relay closed it. + lost: ?anyerror, + /// This event's time is up. Everything a reader handed over + /// before it is ahead of it in the queue. + deadline: u64, + }, + }; + + fn init(self: *Set, gpa: std.mem.Allocator, io: std.Io, urls: []const []const u8) void { + self.* = .{ + .gpa = gpa, + .io = io, + .urls = urls, + .conns = @splat(null), + .state = @splat(.idle), + .readers = @splat(null), + .quiet = @splat(.init(false)), + .gen = @splat(0), + .notices = @splat(.init(0)), + .heard = undefined, + .heard_buf = undefined, + .seq = 0, }; - defer msg.deinit(); - switch (msg.value) { - .ok => |o| { - if (!std.mem.eql(u8, &o.event_id, &id)) continue; - if (o.accepted) { - try_print(err, "deed publish: {s}: accepted\n", .{url}); - return true; + self.heard = .init(&self.heard_buf); + } + + fn deinit(self: *Set) void { + for (0..self.urls.len) |i| _ = self.release(i); + // What the readers left behind, text included. + var buf: [32]Heard = undefined; + while (true) { + const n = self.heard.get(self.io, &buf, 0) catch break; + if (n == 0) break; + for (buf[0..n]) |h| self.free(h); + } + } + + fn free(self: *Set, h: Heard) void { + switch (h.what) { + .ok => |o| self.gpa.free(o.message), + .notice => |t| self.gpa.free(t), + .more_notices, .lost, .deadline => {}, + } + } + + fn live(self: *const Set) usize { + var n: usize = 0; + for (self.conns[0..self.urls.len]) |c| { + if (c != null) n += 1; + } + return n; + } + + /// Stops relay `i`'s reader and waits for it. + fn stopReader(self: *Set, i: usize) void { + if (self.readers[i]) |*f| { + self.quiet[i].store(true, .release); + f.cancel(self.io); + self.readers[i] = null; + } + } + + /// Stops relay `i`'s reader, waits for it, and closes the connection. + /// Returns how many messages it could not read. + fn release(self: *Set, i: usize) u64 { + self.stopReader(i); + const r = self.conns[i] orelse return 0; + self.conns[i] = null; + self.gen[i] +%= 1; + const unreadable = r.unreadable(); + r.shutdown(self.io); + r.deinit(); + return unreadable; + } + + /// Lets go of relay `i`, saying first what it sent that could not be read. + fn drop(self: *Set, i: usize, next: State, err: *std.Io.Writer) !void { + const unreadable = self.release(i); + self.state[i] = next; + try reportUnreadable(err, self.urls[i], unreadable); + } + + fn startReader(self: *Set, i: usize) !void { + self.quiet[i].store(false, .release); + self.readers[i] = try self.io.concurrent(readLoop, .{ self, i, self.gen[i], self.conns[i].? }); + } + + /// Reads relay `i` until its connection goes away or this run lets it go. + /// + /// `receive` rather than `receiveTimeout`: with no deadline the read goes + /// straight to the socket, where a cancel reaches it, including in the + /// middle of a TLS record. + fn readLoop(self: *Set, i: usize, gen: u32, r: *relay.Relay) void { + while (true) { + var msg = (r.receive() catch |e| { + _ = self.tell(i, gen, .{ .lost = e }); + return; + }) orelse { + _ = self.tell(i, gen, .{ .lost = null }); + return; + }; + defer msg.deinit(); + switch (msg.value) { + .ok => |o| { + const text = self.gpa.dupe(u8, o.message) catch { + _ = self.tell(i, gen, .{ .lost = error.OutOfMemory }); + return; + }; + if (!self.tell(i, gen, .{ .ok = .{ .id = o.event_id, .accepted = o.accepted, .message = text } })) { + self.gpa.free(text); + return; + } + }, + .notice => |n| { + const count = self.notices[i].fetchAdd(1, .monotonic) + 1; + if (count <= max_notices) { + const text = self.gpa.dupe(u8, n.message) catch { + _ = self.tell(i, gen, .{ .lost = error.OutOfMemory }); + return; + }; + if (!self.tell(i, gen, .{ .notice = text })) { + self.gpa.free(text); + return; + } + } else if (count == max_notices + 1) { + if (!self.tell(i, gen, .more_notices)) return; + } + }, + // A challenge is not a refusal yet. If the relay needs it + // answered, the OK that follows says so. + .auth, .event, .eose, .closed => {}, + } + } + } + + fn tell(self: *Set, i: usize, gen: u32, what: @FieldType(Heard, "what")) bool { + if (self.quiet[i].load(.acquire)) return false; + self.heard.putOne(self.io, .{ .index = i, .gen = gen, .what = what }) catch return false; + return true; + } + + fn deadlineAfter(q: *std.Io.Queue(Heard), io: std.Io, ms: i64, seq: u64) void { + io.sleep(.fromMilliseconds(ms), .awake) catch return; + q.putOne(io, .{ .what = .{ .deadline = seq } }) catch return; + } + + /// Deals with what the readers handed over while nothing was waiting for + /// an answer: a relay that went away is let go, to be dialled again for + /// the next event, and notices are reported. + fn refresh(self: *Set, err: *std.Io.Writer) !void { + var buf: [32]Heard = undefined; + // At most a queue's worth. Relays still sending as this reads could + // otherwise keep it going for as long as they liked. + var taken: usize = 0; + while (taken < self.heard_buf.len) { + const n = self.heard.get(self.io, &buf, 0) catch return; + if (n == 0) return; + taken += n; + for (buf[0..n], 0..) |h, k| { + errdefer for (buf[k + 1 .. n]) |rest| self.free(rest); + defer self.free(h); + if (h.what == .deadline or h.gen != self.gen[h.index]) continue; + switch (h.what) { + .lost => try self.drop(h.index, .closed, err), + .notice => |t| try self.printNotice(err, h.index, t), + .more_notices => try err.print("deed publish: {s}: more notices, not shown\n", .{self.urls[h.index]}), + .ok, .deadline => {}, } - try_print(err, "deed publish: {s}: refused: {s}\n", .{ url, o.message }); - return false; + } + } + } + + /// Dials every relay not dialled yet, and every relay that went away + /// since the last event, all at once. + fn ready(self: *Set, bound_ms: i64, err: *std.Io.Writer) !void { + var which: [dial.max_relays]usize = undefined; + var n: usize = 0; + for (0..self.urls.len) |i| { + if (self.state[i] != .idle and self.state[i] != .closed) continue; + which[n] = i; + n += 1; + } + try self.dialSome(which[0..n], bound_ms, .gone, err); + } + + /// Dials the relays in `which` at once. A relay that cannot be reached is + /// left in `failed`: `.gone` when it could never be reached, `.closed` when + /// only this attempt, squeezed into what was left of a deadline, failed. + fn dialSome(self: *Set, which: []const usize, bound_ms: i64, failed: State, err: *std.Io.Writer) !void { + if (which.len == 0) return; + var urls: [dial.max_relays][]const u8 = undefined; + for (which, 0..) |i, k| urls[k] = self.urls[i]; + var outcomes: [dial.max_relays]dial.Outcome = undefined; + dial.all(self.gpa, self.io, urls[0..which.len], bound_ms, &outcomes); + // Every connection is stored before anything is written, so a failed + // write cannot strand one that `deinit` would never see. + for (outcomes[0..which.len], which) |o, i| switch (o) { + .connected => |r| { + self.conns[i] = r; + self.state[i] = .connected; + }, + .failed, .timed_out => self.state[i] = failed, + }; + for (outcomes[0..which.len], which) |o, i| switch (o) { + .connected => self.startReader(i) catch |e| { + try err.print("deed publish: {s}: {s}\n", .{ self.urls[i], @errorName(e) }); + try self.drop(i, .gone, err); }, - .notice => |notice| try_print(err, "deed publish: {s}: notice: {s}\n", .{ url, notice.message }), - else => {}, + .failed => |e| try err.print("deed publish: {s}: {s}\n", .{ self.urls[i], @errorName(e) }), + .timed_out => try err.print("deed publish: {s}: no answer within {d} ms\n", .{ self.urls[i], bound_ms }), + }; + } + + const Sent = struct { + index: usize, + result: anyerror!void, + }; + + const SendArrival = union(enum) { + sent: Sent, + timer: std.Io.Cancelable!void, + }; + + fn sendOne(r: *relay.Relay, ev: nostr.event.Event, index: usize) Sent { + return .{ .index = index, .result = r.publish(ev) }; + } + + /// Sends `ev` to the relays in `which` at once, and marks each one that + /// took it as waiting. A send can block: a relay that has stopped reading + /// fills the socket, and the write waits for room that never comes, so + /// whatever is still sending at the deadline is cancelled and let go. + /// + /// Every send is finished or cancelled before this returns, whatever + /// happens, because each one holds the event and its relay. + fn sendSome( + self: *Set, + which: []const usize, + ev: nostr.event.Event, + id: []const u8, + deadline: std.Io.Clock.Timestamp, + timeout_ms: i64, + stalled: State, + waiting: *[dial.max_relays]bool, + err: *std.Io.Writer, + ) !void { + var results: [dial.max_relays]?Sent = @splat(null); + var late: [dial.max_relays]bool = @splat(false); + var started: [dial.max_relays]bool = @splat(false); + { + var buf: [dial.max_relays + 1]SendArrival = undefined; + var sel = std.Io.Select(SendArrival).init(self.io, &buf); + defer while (sel.cancel()) |a| switch (a) { + .sent => |s| { + results[s.index] = s; + late[s.index] = true; + }, + .timer => {}, + }; + var pending: usize = 0; + for (which) |i| { + const r = self.conns[i] orelse continue; + sel.concurrent(.sent, sendOne, .{ r, ev, i }) catch { + // No concurrency to spare. Sent here instead, which is only + // as bounded as the relay is willing to read. + results[i] = .{ .index = i, .result = r.publish(ev) }; + continue; + }; + started[i] = true; + pending += 1; + } + if (pending > 0) sel.concurrent(.timer, std.Io.sleep, .{ self.io, .fromMilliseconds(msLeft(self.io, deadline)), .awake }) catch {}; + while (pending > 0) { + const a = sel.await() catch break; + switch (a) { + .sent => |s| { + results[s.index] = s; + pending -= 1; + }, + .timer => break, + } + } + // A send still going at the deadline may be waiting for the + // connection's write lock rather than for the socket, behind the + // relay's own reader writing a pong to a relay that has stopped + // reading. That wait cannot be cancelled, so the reader is stopped + // first, which lets go of the lock, before the sends are joined. + if (pending > 0) for (which) |i| { + if (started[i] and results[i] == null) self.stopReader(i); + }; + } + // Every send has been joined by now, so what follows can write, and + // fail to, without anything still using the event or a relay. + for (which) |i| { + const s = results[i] orelse continue; + if (s.result) { + waiting[i] = true; + } else |e| { + if (late[i]) { + // Still sending at the deadline. Whatever reached the + // socket is half a frame, so this connection is done. + try err.print("deed publish: {s}: {s}: could not send within {d} ms\n", .{ id, self.urls[i], timeout_ms }); + try self.drop(i, stalled, err); + } else { + try err.print("deed publish: {s}: {s}: {s}\n", .{ id, self.urls[i], @errorName(e) }); + try self.drop(i, .closed, err); + } + } + } + } + + /// Offers `ev` to every connected relay at once and collects their + /// answers, all against one deadline. Returns whether any relay accepted. + fn offer( + self: *Set, + ev: nostr.event.Event, + id: []const u8, + timeout_ms: i64, + err: *std.Io.Writer, + ) !bool { + const io = self.io; + const deadline = std.Io.Clock.Timestamp.fromNow(io, .{ .raw = .fromMilliseconds(timeout_ms), .clock = .awake }); + var waiting: [dial.max_relays]bool = @splat(false); + var resent: [dial.max_relays]bool = @splat(false); + + var which: [dial.max_relays]usize = undefined; + var n: usize = 0; + for (0..self.urls.len) |i| { + if (self.conns[i] == null) continue; + which[n] = i; + n += 1; + } + try self.sendSome(which[0..n], ev, id, deadline, timeout_ms, .gone, &waiting, err); + + self.seq += 1; + var timer = try io.concurrent(deadlineAfter, .{ &self.heard, io, msLeft(io, deadline), self.seq }); + defer timer.cancel(io); + + var accepted = false; + while (anyOf(&waiting)) { + const h = self.heard.getOne(io) catch break; + defer self.free(h); + if (h.what == .deadline) { + if (h.what.deadline != self.seq) continue; + for (0..self.urls.len) |i| { + if (waiting[i]) try err.print("deed publish: {s}: {s}: no answer within {d} ms\n", .{ id, self.urls[i], timeout_ms }); + } + break; + } + if (h.gen != self.gen[h.index]) continue; + const i = h.index; + const url = self.urls[i]; + switch (h.what) { + .ok => |o| { + // An answer about some other event, such as one an + // earlier wait gave up on, is not an answer to this. + if (!waiting[i] or !std.mem.eql(u8, &o.id, &ev.id)) continue; + waiting[i] = false; + if (o.accepted) accepted = true; + try err.print("deed publish: {s}: {s}: {s}", .{ id, url, if (o.accepted) "accepted" else "refused: " }); + if (o.accepted and o.message.len > 0) { + try err.writeAll(" ("); + try writeEscaped(err, o.message); + try err.writeAll(")"); + } else if (!o.accepted) { + try writeEscaped(err, o.message); + } + try err.writeAll("\n"); + }, + .notice => |t| try self.printNotice(err, i, t), + .more_notices => try err.print("deed publish: {s}: more notices, not shown\n", .{url}), + .lost => |why| { + const was_waiting = waiting[i]; + waiting[i] = false; + try self.drop(i, .closed, err); + if (!was_waiting) continue; + // A relay that dropped an idle connection only finds out + // when something is sent down it. One fresh connection and + // one more try, inside the same deadline. + // Only with time left to do it in; a relay that cannot be + // reached in the last few milliseconds is still dialled + // again for the next event. + if (!resent[i] and msLeft(io, deadline) >= min_resend_ms) { + resent[i] = true; + try self.dialSome(&.{i}, @min(dial.default_timeout_ms, msLeft(io, deadline)), .closed, err); + if (self.conns[i] != null) { + try self.sendSome(&.{i}, ev, id, deadline, timeout_ms, .closed, &waiting, err); + if (waiting[i]) continue; + } + } + if (why) |e| { + try err.print("deed publish: {s}: {s}: {s}\n", .{ id, url, @errorName(e) }); + } else { + try err.print("deed publish: {s}: {s}: closed before answering\n", .{ id, url }); + } + }, + .deadline => unreachable, + } + } + return accepted; + } + + fn printNotice(self: *Set, err: *std.Io.Writer, i: usize, text: []const u8) !void { + try err.print("deed publish: {s}: notice: ", .{self.urls[i]}); + try writeEscaped(err, text); + try err.writeAll("\n"); + } + + fn reportAll(self: *Set, err: *std.Io.Writer) !void { + for (0..self.urls.len) |i| { + if (self.conns[i] == null) continue; + try reportUnreadable(err, self.urls[i], self.release(i)); + } + } +}; + +fn anyOf(flags: *const [dial.max_relays]bool) bool { + for (flags) |f| { + if (f) return true; + } + return false; +} + +/// Notices shown per relay per run before the rest are counted and not shown. +const max_notices = 8; + +/// The least time left in which a relay that closed before answering is dialled +/// and sent the event again. +const min_resend_ms = 250; + +fn reportUnreadable(err: *std.Io.Writer, url: []const u8, n: u64) !void { + if (n > 0) try err.print("deed publish: {s}: skipped {d} unreadable {s}\n", .{ url, n, if (n == 1) "message" else "messages" }); +} + +/// Milliseconds until `deadline`, never below zero. +fn msLeft(io: std.Io, deadline: std.Io.Clock.Timestamp) i64 { + const left = deadline.raw.nanoseconds - std.Io.Clock.Timestamp.now(io, .awake).raw.nanoseconds; + if (left <= 0) return 0; + return @intCast(@min(@divFloor(left, std.time.ns_per_ms), std.math.maxInt(i64))); +} + +/// Writes text a relay sent, with every control character shown as an escape: +/// C0 and DEL as `\xNN`, and C1 (U+0080 to U+009F) as `\u00NN`. +/// +/// A relay chooses these bytes. Printed raw, a newline would let it write a +/// line of its own that looks like one of these, such as a relay claiming to +/// have accepted an event, and a control sequence would reach the terminal. +/// C1 matters because some terminals read U+009B as the start of one. +fn writeEscaped(w: *std.Io.Writer, text: []const u8) !void { + var i: usize = 0; + while (i < text.len) : (i += 1) { + const c = text[i]; + if (c < 0x20 or c == 0x7f) { + try w.print("\\x{x:0>2}", .{c}); + } else if (c == 0xc2 and i + 1 < text.len and text[i + 1] >= 0x80 and text[i + 1] <= 0x9f) { + try w.print("\\u00{x:0>2}", .{text[i + 1]}); + i += 1; + } else { + try w.writeByte(c); } } } -/// Writing a progress line must never be the thing that fails a publish. -fn try_print(w: *std.Io.Writer, comptime fmt: []const u8, args: anytype) void { - w.print(fmt, args) catch {}; +// --------------------------------------------------------------------------- +// Tests: the real command, against real relays on loopback. +// --------------------------------------------------------------------------- + +const testrelay = @import("testrelay.zig"); + +/// A signed kind 1 event as one line of JSON, from a fixed key. +fn signedEvent(arena: std.mem.Allocator, io: std.Io, content: []const u8) ![]const u8 { + var signer = try nostr.keys.Signer.initRandomized(io); + defer signer.deinit(); + const kp = try signer.keyPairFromSecretKey([_]u8{0x11} ** 32); + const ev = try nostr.event.create(arena, signer, kp, 1700000000, 1, &.{}, content, null); + return nostr.event.toJson(arena, ev); +} + +fn idOf(arena: std.mem.Allocator, json: []const u8) ![]const u8 { + var parsed = try nostr.event.fromJson(arena, json); + defer parsed.deinit(); + return hex.encode(arena, &parsed.value.id); +} + +const Ran = struct { + code: u8, + out: []const u8, + err: []const u8, + ms: i64, +}; + +/// Runs `deed publish` with `args`, dialling through an allocator that +/// captures no stack traces (see `testrelay.DialAllocator`). +fn runPublish(arena: std.mem.Allocator, args: []const []const u8) !Ran { + const io = std.testing.io; + var da: testrelay.DialAllocator = .init; + defer if (da.deinit() == .leak) @panic("publish leaked"); + var out: std.Io.Writer.Allocating = .init(arena); + var err: std.Io.Writer.Allocating = .init(arena); + const started = std.Io.Timestamp.now(io, .awake).toMilliseconds(); + const code = try run(da.allocator(), io, args, &out.writer, &err.writer); + return .{ + .code = code, + .out = out.written(), + .err = err.written(), + .ms = std.Io.Timestamp.now(io, .awake).toMilliseconds() - started, + }; +} + +fn contains(haystack: []const u8, needle: []const u8) bool { + return std.mem.indexOf(u8, haystack, needle) != null; +} + +test "an accepted event is printed, and every relay's answer is on stderr with its id" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var yes: testrelay.Relay = undefined; + try yes.start(io, .accept); + defer yes.stop(io); + var no: testrelay.Relay = undefined; + try no.start(io, .refuse); + defer no.stop(io); + + var bufs: [2][40]u8 = undefined; + const yes_url = try testrelay.url(&bufs[0], yes.port()); + const no_url = try testrelay.url(&bufs[1], no.port()); + const ev = try signedEvent(arena, io, "hello"); + const id = try idOf(arena, ev); + + const r = try runPublish(arena, &.{ yes_url, no_url, ev }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n", .{ev}), r.out); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "deed publish: {s}: {s}: accepted\n", .{ id, yes_url }))); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "deed publish: {s}: {s}: refused: invalid: test refusal\n", .{ id, no_url }))); +} + +test "an event no relay accepted is not printed, and the run fails" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var no: testrelay.Relay = undefined; + try no.start(io, .refuse); + defer no.stop(io); + var buf: [40]u8 = undefined; + const ev = try signedEvent(arena, io, "hello"); + const id = try idOf(arena, ev); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, no.port()), ev }); + try std.testing.expectEqual(cli.exit_fail, r.code); + try std.testing.expectEqualStrings("", r.out); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "deed publish: {s}: not published, no relay accepted it in time\n", .{id}))); + try std.testing.expect(contains(r.err, "deed publish: 0 of 1 published\n")); +} + +test "one event not published fails the run, and the others still go out" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var yes: testrelay.Relay = undefined; + try yes.start(io, .accept); + defer yes.stop(io); + var buf: [40]u8 = undefined; + const good = try signedEvent(arena, io, "hello"); + const later = try signedEvent(arena, io, "later"); + // Signed, then altered: its id no longer matches its content. + const forged = try std.mem.replaceOwned(u8, arena, try signedEvent(arena, io, "original"), "original", "tampered"); + const forged_id = try idOf(arena, forged); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, yes.port()), good, forged, "not json", later }); + try std.testing.expectEqual(cli.exit_fail, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n{s}\n", .{ good, later }), r.out); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "deed publish: {s}: not sent, bad signature\n", .{forged_id}))); + try std.testing.expect(contains(r.err, "deed publish: not an event: ")); + try std.testing.expect(contains(r.err, "deed publish: 2 of 4 published\n")); + // The forgery never reached the relay. + try std.testing.expectEqual(@as(usize, 2), yes.events.load(.monotonic)); +} + +test "a relay that already had the event counts as accepting it, and says so" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var dup: testrelay.Relay = undefined; + try dup.start(io, .duplicate); + defer dup.stop(io); + var buf: [40]u8 = undefined; + const ev = try signedEvent(arena, io, "hello"); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, dup.port()), ev }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n", .{ev}), r.out); + try std.testing.expect(contains(r.err, ": accepted (duplicate: already have it)\n")); +} + +test "the answers have one deadline, however much a relay talks" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + // A notice every 100 ms and never an answer. A wait that started again + // with every message would never end. + var chatty: testrelay.Relay = undefined; + try chatty.start(io, .chatty); + defer chatty.stop(io); + var buf: [40]u8 = undefined; + const url = try testrelay.url(&buf, chatty.port()); + const ev = try signedEvent(arena, io, "hello"); + + const r = try runPublish(arena, &.{ "--timeout", "400", url, ev }); + try std.testing.expectEqual(cli.exit_fail, r.code); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "{s}: no answer within 400 ms\n", .{url}))); + try std.testing.expect(contains(r.err, ": notice: still here\n")); + try std.testing.expect(r.ms < 3_000); +} + +test "a relay that never answers the dial does not hold the others" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var silent = try testrelay.listen(io); + defer silent.deinit(io); + var yes: testrelay.Relay = undefined; + try yes.start(io, .accept); + defer yes.stop(io); + var bufs: [2][40]u8 = undefined; + const silent_url = try testrelay.url(&bufs[0], silent.socket.address.ip4.port); + const ev = try signedEvent(arena, io, "hello"); + + const r = try runPublish(arena, &.{ "--timeout", "1000", silent_url, try testrelay.url(&bufs[1], yes.port()), ev }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n", .{ev}), r.out); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "deed publish: {s}: no answer within 1000 ms\n", .{silent_url}))); + try std.testing.expect(r.ms < 3_000); +} + +test "a message the relay sends that cannot be read does not cost the answer" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var odd: testrelay.Relay = undefined; + try odd.start(io, .unreadable_first); + defer odd.stop(io); + var buf: [40]u8 = undefined; + const ev = try signedEvent(arena, io, "hello"); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, odd.port()), ev }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expect(contains(r.err, ": accepted\n")); + try std.testing.expect(contains(r.err, ": skipped 1 unreadable message\n")); +} + +test "text from a relay cannot forge a line or reach the terminal" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var bad: testrelay.Relay = undefined; + try bad.start(io, .hostile); + defer bad.stop(io); + var buf: [40]u8 = undefined; + const ev = try signedEvent(arena, io, "hello"); + const id = try idOf(arena, ev); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, bad.port()), ev }); + try std.testing.expectEqual(cli.exit_fail, r.code); + try std.testing.expect(std.mem.indexOfScalar(u8, r.err, 0x1b) == null); + // The claim is there, but inside the refusal, on the refusal's own line. + try std.testing.expect(!contains(r.err, try std.fmt.allocPrint(arena, "\ndeed publish: {s}: ws://x: accepted", .{id}))); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "refused: invalid: \\x1b[31mred\\x0adeed publish: {s}: ws://x: accepted\n", .{id}))); +} + +test "an event carrying a key NIP-01 does not name is published as the event that was signed" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var yes: testrelay.Relay = undefined; + try yes.start(io, .accept); + defer yes.stop(io); + var buf: [40]u8 = undefined; + const ev = try signedEvent(arena, io, "hello"); + const extra = try std.mem.concat(arena, u8, &.{ "{\"seen_on\":[\"wss://a\"],", ev[1..] }); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, yes.port()), extra }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n", .{ev}), r.out); +} + +test "a command that makes no sense is refused before anything is dialled" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + try std.testing.expectEqual(cli.exit_usage, (try runPublish(arena, &.{"{}"})).code); + try std.testing.expectEqual(cli.exit_usage, (try runPublish(arena, &.{ "--nope", "ws://127.0.0.1:1" })).code); + try std.testing.expectEqual(cli.exit_usage, (try runPublish(arena, &.{ "--timeout", "soon", "ws://127.0.0.1:1" })).code); + try std.testing.expectEqual(cli.exit_usage, (try runPublish(arena, &.{ "--timeout", "0", "ws://127.0.0.1:1" })).code); + + var many: [dial.max_relays + 1][]const u8 = undefined; + for (&many) |*u| u.* = "ws://127.0.0.1:1"; + try std.testing.expectEqual(cli.exit_usage, (try runPublish(arena, &many)).code); +} + +test "a relay that stops reading is dropped at the deadline, and the others still get every event" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var deaf: testrelay.Relay = undefined; + try deaf.start(io, .deaf); + defer deaf.stop(io); + var yes: testrelay.Relay = undefined; + try yes.start(io, .unreadable_first); + defer yes.stop(io); + var bufs: [2][40]u8 = undefined; + const deaf_url = try testrelay.url(&bufs[0], deaf.port()); + + // Big enough that a relay reading none of them fills the socket within a + // few events, and the write to it waits for room that never comes. + // + // The other relay sends a message that cannot be read ahead of each + // answer, and its answer to the event stuck behind the deaf relay arrives + // while that write is still waiting. It is counted all the same. + const big = try arena.alloc(u8, 900 * 1024); + @memset(big, 'x'); + var args: [8][]const u8 = undefined; + args[0] = "--timeout"; + args[1] = "500"; + args[2] = deaf_url; + args[3] = try testrelay.url(&bufs[1], yes.port()); + for (args[4..], 0..) |*a, i| { + big[0] = @intCast('a' + i); + a.* = try signedEvent(arena, io, big); + } + + const r = try runPublish(arena, &args); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqual(@as(usize, 4), std.mem.count(u8, r.out, "\n")); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "{s}: could not send within 500 ms\n", .{deaf_url}))); + try std.testing.expectEqual(@as(usize, 4), yes.events.load(.monotonic)); + try std.testing.expect(r.ms < 10_000); +} + +test "a relay that closed while the run waited for input is dialled again" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + // Closes after each answer, and takes a second connection. + var once: testrelay.Relay = undefined; + try once.startFor(io, .close_after_ok, 2); + defer once.stop(io); + var buf: [40]u8 = undefined; + const first = try signedEvent(arena, io, "first"); + const second = try signedEvent(arena, io, "second"); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, once.port()), first, second }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n{s}\n", .{ first, second }), r.out); + try std.testing.expectEqual(@as(usize, 2), once.events.load(.monotonic)); +} + +test "an answer about some other event is not taken as this event's" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var stale: testrelay.Relay = undefined; + try stale.start(io, .stale_ok); + defer stale.stop(io); + var buf: [40]u8 = undefined; + const ev = try signedEvent(arena, io, "hello"); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, stale.port()), ev }); + try std.testing.expectEqual(cli.exit_fail, r.code); + try std.testing.expectEqualStrings("", r.out); + try std.testing.expect(contains(r.err, ": refused: invalid: not this one\n")); +} + +test "nothing is dialled for input that has nothing to publish" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + // A closed port: dialling it would be reported, and it is not. + var closed = try testrelay.listen(io); + const port = closed.socket.address.ip4.port; + closed.deinit(io); + var buf: [40]u8 = undefined; + const url = try testrelay.url(&buf, port); + + const r = try runPublish(arena, &.{ url, "not json" }); + try std.testing.expectEqual(cli.exit_fail, r.code); + try std.testing.expect(contains(r.err, "deed publish: not an event: ")); + try std.testing.expect(!contains(r.err, url)); +} + +test "pongs arriving back to back do not keep the wait going" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var pongs: testrelay.Relay = undefined; + try pongs.start(io, .pongs); + defer pongs.stop(io); + var buf: [40]u8 = undefined; + const url = try testrelay.url(&buf, pongs.port()); + const ev = try signedEvent(arena, io, "hello"); + + const r = try runPublish(arena, &.{ "--timeout", "300", url, ev }); + try std.testing.expectEqual(cli.exit_fail, r.code); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "{s}: no answer within 300 ms\n", .{url}))); + try std.testing.expect(r.ms < 3_000); +} + +test "a relay's notices are shown up to a point" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var burst: testrelay.Relay = undefined; + try burst.start(io, .notice_burst); + defer burst.stop(io); + var buf: [40]u8 = undefined; + const ev = try signedEvent(arena, io, "hello"); + + const r = try runPublish(arena, &.{ "--timeout", "200", try testrelay.url(&buf, burst.port()), ev }); + try std.testing.expectEqual(@as(usize, max_notices), std.mem.count(u8, r.err, ": notice: again\n")); + try std.testing.expectEqual(@as(usize, 1), std.mem.count(u8, r.err, ": more notices, not shown\n")); +} + +test "control characters in an acceptance or a notice are shown as escapes too" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var bad: testrelay.Relay = undefined; + try bad.start(io, .hostile_accept); + defer bad.stop(io); + var buf: [40]u8 = undefined; + const ev = try signedEvent(arena, io, "hello"); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, bad.port()), ev }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expect(std.mem.indexOfScalar(u8, r.err, 0x1b) == null); + try std.testing.expect(std.mem.indexOfScalar(u8, r.err, 0x7f) == null); + try std.testing.expect(std.mem.indexOf(u8, r.err, "\xc2\x9b") == null); + try std.testing.expect(contains(r.err, ": notice: \\x1bbad\\x7f\\u009b31m\n")); + try std.testing.expect(contains(r.err, ": accepted (note\\x1b[2J\\u009b2J\\x7f)\n")); +} + +test "an answer behind a ping or a run of notices is counted, however many relays are silent" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var pinger: testrelay.Relay = undefined; + try pinger.start(io, .ping_then_ok); + defer pinger.stop(io); + var talker: testrelay.Relay = undefined; + try talker.start(io, .notices_then_ok); + defer talker.stop(io); + var silent: [20]testrelay.Relay = undefined; + for (&silent) |*r| try r.start(io, .silent); + defer for (&silent) |*r| r.stop(io); + + var bufs: [22][40]u8 = undefined; + var args: [25][]const u8 = undefined; + args[0] = "--timeout"; + args[1] = "1000"; + for (&silent, 0..) |*r, i| args[2 + i] = try testrelay.url(&bufs[i], r.port()); + const pinger_url = try testrelay.url(&bufs[20], pinger.port()); + const talker_url = try testrelay.url(&bufs[21], talker.port()); + args[22] = pinger_url; + args[23] = talker_url; + const ev = try signedEvent(arena, io, "hello"); + args[24] = ev; + + const r = try runPublish(arena, &args); + // Accepted by two relays is published, whatever the silent ones do. + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n", .{ev}), r.out); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "{s}: accepted\n", .{pinger_url}))); + try std.testing.expect(contains(r.err, try std.fmt.allocPrint(arena, "{s}: accepted\n", .{talker_url}))); +} + +test "a relay that pings and then closes is dialled again for the next event" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var r1: testrelay.Relay = undefined; + try r1.startFor(io, .ping_then_close, 2); + defer r1.stop(io); + var buf: [40]u8 = undefined; + const first = try signedEvent(arena, io, "first"); + const second = try signedEvent(arena, io, "second"); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, r1.port()), first, second }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n{s}\n", .{ first, second }), r.out); +} + +test "an event sent down a connection that then closes is offered again on a fresh one" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + // The second event reaches the first connection, which closes without + // answering. It is offered once more on a second connection. + var flaky: testrelay.Relay = undefined; + try flaky.startFor(io, .close_on_second_event, 2); + defer flaky.stop(io); + var buf: [40]u8 = undefined; + const first = try signedEvent(arena, io, "first"); + const second = try signedEvent(arena, io, "second"); + + const r = try runPublish(arena, &.{ try testrelay.url(&buf, flaky.port()), first, second }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n{s}\n", .{ first, second }), r.out); + try std.testing.expectEqual(@as(usize, 3), flaky.events.load(.monotonic)); +} + +test "a relay that pings and never reads does not hold the run, although its reader is stuck writing a pong" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var pinger: testrelay.Relay = undefined; + try pinger.start(io, .pings_deaf); + defer pinger.stop(io); + var yes: testrelay.Relay = undefined; + try yes.start(io, .accept); + defer yes.stop(io); + var bufs: [2][40]u8 = undefined; + const pinger_url = try testrelay.url(&bufs[0], pinger.port()); + const first = try signedEvent(arena, io, "first"); + const second = try signedEvent(arena, io, "second"); + + // The second send to it waits behind the reader's pong for the write lock. + const r = try runPublish(arena, &.{ "--timeout", "500", pinger_url, try testrelay.url(&bufs[1], yes.port()), first, second }); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqualStrings(try std.fmt.allocPrint(arena, "{s}\n{s}\n", .{ first, second }), r.out); + try std.testing.expect(r.ms < 5_000); +} + +test "notices are counted per relay across its connections" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const io = std.testing.io; + + var talker: testrelay.Relay = undefined; + try talker.startFor(io, .notices_then_close, 3); + defer talker.stop(io); + var buf: [40]u8 = undefined; + const url = try testrelay.url(&buf, talker.port()); + var args: [4][]const u8 = undefined; + args[0] = url; + for (args[1..], 0..) |*a, i| a.* = try signedEvent(arena, io, try std.fmt.allocPrint(arena, "event {d}", .{i})); + + const r = try runPublish(arena, &args); + try std.testing.expectEqual(cli.exit_ok, r.code); + try std.testing.expectEqual(@as(usize, max_notices), std.mem.count(u8, r.err, ": notice: again\n")); + try std.testing.expectEqual(@as(usize, 1), std.mem.count(u8, r.err, ": more notices, not shown\n")); } diff --git a/src/cmd_req.zig b/src/cmd_req.zig index 3ed1299..3c16272 100644 --- a/src/cmd_req.zig +++ b/src/cmd_req.zig @@ -38,6 +38,7 @@ pub const usage = \\ \\A relay that has not accepted the connection within five seconds, or \\within --timeout if that is shorter, is named on stderr and left out. + \\Looking up a relay's name is the one step that cannot be cut short. \\ \\Given no relay, it prints what it would send and stops, so a filter can be \\read before it is asked of anybody: diff --git a/src/relayset.zig b/src/relayset.zig index ba7ba42..82ce076 100644 --- a/src/relayset.zig +++ b/src/relayset.zig @@ -87,6 +87,9 @@ pub fn query( var relays: [max_relays]?*relay.Relay = @splat(null); var done: [max_relays]bool = @splat(false); const n = @min(urls.len, max_relays); + if (urls.len > max_relays) { + try err.print("deed: {d} relays were named, and only the first {d} are asked\n", .{ urls.len, max_relays }); + } defer for (relays[0..n]) |maybe| { if (maybe) |r| { @@ -260,6 +263,31 @@ test "a relay that never answers the dial costs the run its deadline, not foreve try std.testing.expect(took < 3_000); } +test "relays past the most one run dials are named as left out, not dropped quietly" { + const io = std.testing.io; + const testrelay = @import("testrelay.zig"); + var da: testrelay.DialAllocator = .init; + defer if (da.deinit() == .leak) @panic("the query leaked"); + const gpa = da.allocator(); + + // Closed ports: each dial is refused at once. + var closed = try testrelay.listen(io); + const port = closed.socket.address.ip4.port; + closed.deinit(io); + var buf: [40]u8 = undefined; + const url = try testrelay.url(&buf, port); + var urls: [max_relays + 1][]const u8 = undefined; + for (&urls) |*u| u.* = url; + const filters = [_]filter.Filter{.{ .kinds = &.{1} }}; + + var out: std.Io.Writer.Allocating = .init(std.testing.allocator); + defer out.deinit(); + var err: std.Io.Writer.Allocating = .init(std.testing.allocator); + defer err.deinit(); + _ = try query(gpa, io, &urls, &filters, &out.writer, &err.writer, .{ .deadline_ms = 500, .poll_ms = 50 }); + try std.testing.expect(std.mem.indexOf(u8, err.written(), "deed: 33 relays were named, and only the first 32 are asked\n") != null); +} + test { // Forces this file to be analysed. `_ = @import("relayset.zig")` alone // imports it without ever compiling a function body nobody references, so diff --git a/src/testrelay.zig b/src/testrelay.zig index fc87691..fe015c9 100644 --- a/src/testrelay.zig +++ b/src/testrelay.zig @@ -19,11 +19,43 @@ pub const Mode = enum { duplicate, /// Answers the upgrade and then never sends another byte. silent, - /// Never answers an EVENT, and sends a NOTICE every 100 ms, so a wait - /// that restarts on every message would never end. + /// Never answers an EVENT, and sends a NOTICE every 100 ms for four + /// seconds, so a wait that restarts on every message outlasts a test that + /// expects it to end sooner, and fails it rather than hanging it. chatty, + /// Sends twenty notices as soon as it is connected, and never answers. + notice_burst, + /// Sends pong frames back to back, and never answers. + pongs, + /// Answers the upgrade and never reads again, so what is sent to it + /// fills the socket and the sender's write waits. + deaf, + /// OK true to every EVENT, then closes the connection. + close_after_ok, + /// An OK true about some other event, then OK false about this one. + stale_ok, + /// A notice, then OK true, both carrying C0, DEL and C1 controls. + hostile_accept, + /// A ping, then OK true. + ping_then_ok, + /// Seventy notices, then OK true. + notices_then_ok, + /// OK true, then a ping and a close frame, and the connection ends. + ping_then_close, + /// OK true to the first EVENT on a connection. On the second, it closes + /// the connection without answering. + close_on_second_event, + /// Never reads, and sends pings back to back. The answering pongs fill + /// the socket, so the reader writing them waits, holding the connection's + /// write lock. + pings_deaf, + /// Twenty notices, then OK true, then the connection ends. + notices_then_close, /// Sends a message the parser has no case for before each OK true. unreadable_first, + /// Refuses every EVENT with a reason carrying a terminal escape and a + /// newline followed by a line that claims the event was accepted. + hostile, }; /// The allocator a test hands to anything that dials. @@ -50,12 +82,30 @@ pub const Relay = struct { mode: Mode, /// EVENT messages this relay has received, across every connection. events: std.atomic.Value(usize) = .init(0), + /// Set by `close_after_ok` once it has answered, to end the connection. + closing: bool = false, + /// Connections served before the task returns. + connections: usize = 1, task: Io.Future(void), /// Starts serving in the background. `self` must not move until `stop`. pub fn start(self: *Relay, io: Io, mode: Mode) !void { - self.* = .{ .server = try listen(io), .mode = mode, .task = undefined }; + return self.startFor(io, mode, 1); + } + + /// Like `start`, serving `connections` connections one after another. + /// Only for a relay whose earlier connections end on their own: see + /// `serve` for why a cancel must never land in the middle of one. + pub fn startFor(self: *Relay, io: Io, mode: Mode, connections: usize) !void { + self.* = .{ .server = try listen(io), .mode = mode, .connections = connections, .task = undefined }; errdefer self.server.deinit(io); + if (mode == .deaf or mode == .pings_deaf) { + // A small window, inherited by the connection it accepts, so what + // is sent to it fills the sender's socket after a few hundred + // kilobytes however large the system's defaults are. + const size: c_int = 4096; + try std.posix.setsockopt(self.server.socket.handle, std.posix.SOL.SOCKET, std.posix.SO.RCVBUF, std.mem.asBytes(&size)); + } self.task = try io.concurrent(serve, .{ self, io }); } @@ -74,9 +124,11 @@ pub const Relay = struct { /// cancel, a blocking call it makes afterwards can no longer be woken, so /// `stop` would wait on that `accept` forever. fn serve(self: *Relay, io: Io) void { - const conn = self.server.accept(io) catch return; - defer conn.close(io); - self.handle(io, conn) catch {}; + for (0..self.connections) |_| { + const conn = self.server.accept(io) catch return; + defer conn.close(io); + self.handle(io, conn) catch {}; + } } fn handle(self: *Relay, io: Io, conn: Io.net.Stream) !void { @@ -84,6 +136,7 @@ pub const Relay = struct { // this task is cancelled when the test stops the relay, and a cancel // swallowed by a stack capture would leave `stop` waiting forever. const gpa = std.heap.page_allocator; + self.closing = false; var rbuf: [8192]u8 = undefined; var wbuf: [8192]u8 = undefined; var r = conn.reader(io, &rbuf); @@ -104,16 +157,34 @@ pub const Relay = struct { try w.interface.print("HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: {s}\r\n\r\n", .{&accept_key}); try w.interface.flush(); - if (self.mode == .chatty) { - while (true) { - try io.sleep(.fromMilliseconds(100), .awake); - try sendText(&w.interface, "[\"NOTICE\",\"still here\"]"); - } + switch (self.mode) { + .chatty => { + for (0..40) |_| { + try io.sleep(.fromMilliseconds(100), .awake); + try sendText(&w.interface, "[\"NOTICE\",\"still here\"]"); + } + return; + }, + .notice_burst => for (0..20) |_| try sendText(&w.interface, "[\"NOTICE\",\"again\"]"), + .pongs => while (true) { + try w.interface.writeAll(&.{ 0x8A, 0x00 }); + try w.interface.flush(); + }, + .deaf => while (true) try io.sleep(.fromMilliseconds(1000), .awake), + // The largest a ping may carry, so the pongs echoing it fill the + // socket within a fraction of a second rather than many. + .pings_deaf => while (true) { + try w.interface.writeAll(&(.{ 0x89, 125 } ++ .{'p'} ** 125)); + try w.interface.flush(); + }, + else => {}, } var frames: std.ArrayList(u8) = .empty; defer frames.deinit(gpa); + var seen: usize = 0; while (true) { + if (self.closing) return; // Not `readSliceShort`: it blocks until its destination is full, // so a frame smaller than the buffer would never surface. r.interface.fillMore() catch return; @@ -122,7 +193,8 @@ pub const Relay = struct { r.interface.toss(avail.len); while (try ws.decodeFrame(frames.items)) |f| { if (f.opcode == .close) return; - if (f.opcode == .text) try self.answer(gpa, &w.interface, f.payload); + if (f.opcode == .text) try self.answer(gpa, &w.interface, f.payload, &seen); + if (self.closing) return; const n = f.frame_len; std.mem.copyForwards(u8, frames.items, frames.items[n..]); frames.shrinkRetainingCapacity(frames.items.len - n); @@ -130,7 +202,7 @@ pub const Relay = struct { } } - fn answer(self: *Relay, gpa: std.mem.Allocator, w: *Io.Writer, payload: []const u8) !void { + fn answer(self: *Relay, gpa: std.mem.Allocator, w: *Io.Writer, payload: []const u8, seen: *usize) !void { if (self.mode == .silent) return; var parsed = try std.json.parseFromSlice(std.json.Value, gpa, payload, .{}); defer parsed.deinit(); @@ -141,16 +213,56 @@ pub const Relay = struct { try sendText(w, try std.fmt.bufPrint(&text, "[\"EOSE\",\"{s}\"]", .{items[1].string})); } else if (std.mem.eql(u8, kind, "EVENT")) { _ = self.events.fetchAdd(1, .monotonic); + seen.* += 1; const id = items[1].object.get("id").?.string; switch (self.mode) { .accept => try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"\"]", .{id})), .refuse => try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",false,\"invalid: test refusal\"]", .{id})), .duplicate => try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"duplicate: already have it\"]", .{id})), + .hostile => try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",false,\"invalid: \\u001b[31mred\\ndeed publish: {s}: ws://x: accepted\"]", .{ id, id })), + .close_after_ok => { + try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"\"]", .{id})); + self.closing = true; + }, + .stale_ok => { + try sendText(w, "[\"OK\",\"" ++ "00" ** 32 ++ "\",true,\"\"]"); + try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",false,\"invalid: not this one\"]", .{id})); + }, + .hostile_accept => { + try sendText(w, "[\"NOTICE\",\"\\u001bbad\\u007f\\u009b31m\"]"); + try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"note\\u001b[2J\\u009b2J\\u007f\"]", .{id})); + }, + .ping_then_ok => { + try w.writeAll(&.{ 0x89, 0x00 }); + try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"\"]", .{id})); + }, + .notices_then_ok => { + for (0..70) |_| try sendText(w, "[\"NOTICE\",\"busy\"]"); + try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"\"]", .{id})); + }, + .notices_then_close => { + for (0..20) |_| try sendText(w, "[\"NOTICE\",\"again\"]"); + try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"\"]", .{id})); + self.closing = true; + }, + .ping_then_close => { + try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"\"]", .{id})); + try w.writeAll(&.{ 0x89, 0x00, 0x88, 0x00 }); + try w.flush(); + self.closing = true; + }, + .close_on_second_event => { + if (seen.* >= 2) { + self.closing = true; + } else { + try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"\"]", .{id})); + } + }, .unreadable_first => { try sendText(w, "[\"COUNT\",\"x\",{\"count\":1}]"); try sendText(w, try std.fmt.bufPrint(&text, "[\"OK\",\"{s}\",true,\"\"]", .{id})); }, - .silent, .chatty => {}, + .silent, .chatty, .notice_burst, .pongs, .deaf, .pings_deaf => {}, } } } From aa444a5377dac877b7a28fa216ca61811a71a89a Mon Sep 17 00:00:00 2001 From: sepehr-safari Date: Wed, 23 Sep 2026 15:22:44 +0300 Subject: [PATCH 2/2] chore(release): 0.3.0 Publish says what it published. --- .github/RELEASE_NOTES.md | 20 +++++++++++++++++++- build.zig.zon | 2 +- src/main.zig | 2 +- 3 files changed, 21 insertions(+), 3 deletions(-) diff --git a/.github/RELEASE_NOTES.md b/.github/RELEASE_NOTES.md index 74a3df4..cb02aa6 100644 --- a/.github/RELEASE_NOTES.md +++ b/.github/RELEASE_NOTES.md @@ -2,6 +2,24 @@ Every artifact below is published with a `.sha256` beside it, so the download can be checked against a digest that was written by the same job that built it. +### What's new in v0.3.0 + +`publish` says what it published. + +**What comes out of `publish` is what was published.** Each event that at least one relay accepted is printed once the relays have answered or the deadline has passed, so `deed publish wss://a wss://b < events.jsonl > sent.jsonl` leaves a record of exactly which events a relay accepted in time. Every relay's answer goes to stderr with the event's id on it, including a relay saying it already had the event. + +**The exit code covers every event.** `publish` exits 0 only when every event was accepted by at least one relay. An event no relay accepted, a record that is not an event, or an event whose signature does not check out is reported on stderr and makes the run exit 1, and the rest of the stream still goes out. In v0.2.0 one acceptance anywhere in the run was enough to exit 0. + +**Events are checked before they are sent.** An event whose id or signature does not match its content is not offered to any relay. + +**One deadline per event.** Every relay is sent the event at once, and the relays have ten seconds from then to take it and answer, `--timeout` to change it. Each relay is read on its own for the whole run, so an answer is counted the moment it arrives, a relay that keeps sending notices, pings or messages that cannot be read does not keep the wait going, and a relay that stalls in the middle of a message holds up nobody else. A relay that stops reading what is sent to it is dropped when the deadline passes. + +**A relay that never answers the dial no longer holds a run open.** `req`, `fetch` and `publish` dial the relays at once, and a relay that has not accepted the connection within five seconds is named and left out. The name lookup is the one step that cannot be cut short. Before, a relay that accepted the TCP connection and never answered the websocket upgrade held the whole run, and every relay listed after it, forever. A relay that closes the connection while `publish` waits for more input is dialled again for the next event, and an event sent down a connection that closes before answering is offered once more on a fresh one. + +**Text from a relay is shown, not obeyed.** Control characters in a relay's refusal or notice are printed as escapes, so a relay cannot write a line of its own into the output or send control sequences to the terminal. Notices are shown up to eight per relay. + +**Two fixes from the nostr library, now at 0.14.5.** An event that carries a key NIP-01 does not name is accepted as the event that was signed, and a message from a relay that cannot be read costs that message only, rather than every message after it on that connection. + ### What's new in v0.2.0 deed reaches relays now, and keeps what it finds. @@ -15,7 +33,7 @@ deed req -k 1 -l 50 --store ~/.deed/db wss://relay.example deed req -k 1 -l 50 --store ~/.deed/db --local ``` -The second dials nothing. The events are the same events, and they still verify, because what is stored is what was signed. Other nostr command lines do not do this: the one most people use keeps its local database behind a build tag for Linux on x86_64 only, and even there its own query verbs never write to it. +The second dials nothing. The events are the same events, and they still verify, because what is stored is what was signed. **Events are checked before they are kept or printed.** A relay can send anything. A signature that does not verify is dropped, and so is an event that does not answer the question that was asked. Both are reported rather than silently skipped. diff --git a/build.zig.zon b/build.zig.zon index 61ddc08..9efc878 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -1,6 +1,6 @@ .{ .name = .deed, - .version = "0.2.0", + .version = "0.3.0", .fingerprint = 0x89498c2094a1a4a3, .minimum_zig_version = "0.16.0", .dependencies = .{ diff --git a/src/main.zig b/src/main.zig index ad2bd6b..33cbd9c 100644 --- a/src/main.zig +++ b/src/main.zig @@ -15,7 +15,7 @@ const cmd_publish = @import("cmd_publish.zig"); const cmd_req = @import("cmd_req.zig"); const cmd_verify = @import("cmd_verify.zig"); -pub const version = "0.2.0"; +pub const version = "0.3.0"; const usage = \\deed: the nostr command line