Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

## [0.14.7] - 2026-09-24

### Fixed

- A pong no longer holds `receiveTimeout` past its deadline. A ping was answered with a write that had no bound, so a relay that sent pings and had stopped reading filled the socket, and the pong waited for room that never came while holding the connection's write lock, which also held up any other write on the connection. With a deadline, the pong is now skipped when another write holds the lock or when the socket has no room before the deadline; a peer that is not reading would not read it, and RFC 6455 allows answering only the latest of several pings. With no deadline, a ping is answered as before.

### Added

- A `Connection` stream may provide `writableBy(deadline)`, which `IoStream` implements as a wait for room that writes nothing. A stream without it is taken to always have room.

## [0.14.6] - 2026-09-23

### Added
Expand Down
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ it like [Notary](https://github.com/zig-nostr/notary), a remote signer that keep
your key off every client. Full docs, benchmarks, and the ecosystem overview
live at [zignostr.com](https://zignostr.com).

> **Status: early (`v0.14.6`).** The library core, transport, local-first store
> **Status: early (`v0.14.7`).** The library core, transport, local-first store
> and signer protocol have shipped and are covered by tests. Two native apps run
> on it today. APIs may still change before 1.0.

Expand Down Expand Up @@ -77,7 +77,7 @@ Methodology and the full write-up are on the
Add the library to your `build.zig.zon`:

```sh
zig fetch --save https://github.com/zig-nostr/nostr/archive/refs/tags/v0.14.6.tar.gz
zig fetch --save https://github.com/zig-nostr/nostr/archive/refs/tags/v0.14.7.tar.gz
```

Wire the module in `build.zig`:
Expand Down
2 changes: 1 addition & 1 deletion build.zig.zon
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
.{
.name = .nostr,
.version = "0.14.6",
.version = "0.14.7",
// Generated at project creation; never regenerate for this repo.
.fingerprint = 0x208aa38fcd8fcc08,
.minimum_zig_version = "0.16.0",
Expand Down
131 changes: 126 additions & 5 deletions src/relay.zig
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,9 @@ pub const ConnectionError = error{
/// — 0 means the peer closed. A null deadline waits forever, which is what
/// every caller did before `receiveTimeout` existed.
/// * `fn writeAll(self, bytes: []const u8) !void`
/// * optionally `fn writableBy(self, deadline: std.Io.Clock.Timestamp) bool`,
/// whether a small write would go through before the deadline. A stream
/// without it is taken to always have room.
///
/// The connection owns two growable buffers (raw receive bytes and the
/// reassembled message) but not the stream; call `deinit` to free them.
Expand Down Expand Up @@ -329,7 +332,7 @@ pub fn Connection(comptime Stream: type) type {
const n = @min(frame.payload.len, echo.len);
@memcpy(echo[0..n], frame.payload[0..n]);
self.consume(frame.frame_len);
try self.sendFrame(self.io, .pong, echo[0..n]);
try self.answerPing(echo[0..n], deadline);
if (self.passed(deadline)) return error.Timeout;
},
.pong => {
Expand Down Expand Up @@ -387,12 +390,9 @@ pub fn Connection(comptime Stream: type) type {
}

fn sendFrame(self: *Self, io: std.Io, opcode: websocket.Opcode, payload: []const u8) !void {
var mask: [4]u8 = undefined;
io.randomSecure(&mask) catch return ConnectionError.RandomFailed;

var frame: std.ArrayList(u8) = .empty;
defer frame.deinit(self.allocator);
try websocket.appendClientFrame(&frame, self.allocator, opcode, payload, mask);
try self.buildFrame(io, &frame, opcode, payload);

// Uncancelable: a half-written frame is a broken stream, and the
// only thing to do after being cancelled here would be to finish
Expand All @@ -402,6 +402,39 @@ pub fn Connection(comptime Stream: type) type {
try self.stream.writeAll(frame.items);
}

fn buildFrame(self: *Self, io: std.Io, frame: *std.ArrayList(u8), opcode: websocket.Opcode, payload: []const u8) !void {
var mask: [4]u8 = undefined;
io.randomSecure(&mask) catch return ConnectionError.RandomFailed;
try websocket.appendClientFrame(frame, self.allocator, opcode, payload, mask);
}

/// Answers a ping with a pong.
///
/// With no deadline this is an ordinary write. With one, the pong must
/// not become the thing that holds the call past it. A relay that has
/// stopped reading fills the socket, and a pong written into a full
/// socket waits for room that never comes, holding the write lock the
/// whole time. So the pong is skipped when another write holds the
/// lock, or when the socket has no room before the deadline. Skipping
/// one is harmless: a peer that is not reading would not read it, and
/// RFC 6455 lets an endpoint answer only the latest of several pings.
fn answerPing(self: *Self, payload: []const u8, deadline: ?std.Io.Clock.Timestamp) !void {
const d = deadline orelse return self.sendFrame(self.io, .pong, payload);
var frame: std.ArrayList(u8) = .empty;
defer frame.deinit(self.allocator);
try self.buildFrame(self.io, &frame, .pong, payload);
if (!self.write_lock.tryLock()) return;
defer self.write_lock.unlock(self.io);
const StreamType = switch (@typeInfo(Stream)) {
.pointer => |ptr| ptr.child,
else => Stream,
};
if (comptime @hasDecl(StreamType, "writableBy")) {
if (!self.stream.writableBy(d)) return;
}
try self.stream.writeAll(frame.items);
}

/// Reads more bytes into `recv`. Returns false on EOF.
fn fill(self: *Self, deadline: ?std.Io.Clock.Timestamp) !bool {
var tmp: [4096]u8 = undefined;
Expand Down Expand Up @@ -532,6 +565,22 @@ pub const IoStream = struct {
return n;
}

/// Whether the socket has room for a small write before `deadline`.
///
/// A wait for room, not a write, so it changes nothing when it gives up.
/// A control frame is far smaller than the kernel's low-water mark, so a
/// socket that reports room takes the whole frame without blocking. With
/// no socket to ask (the hermetic tests) the answer is yes.
pub fn writableBy(self: IoStream, deadline: std.Io.Clock.Timestamp) bool {
const io = self.io orelse return true;
const sock = self.socket orelse return true;
const left = deadline.raw.nanoseconds - std.Io.Clock.Timestamp.now(io, deadline.clock).raw.nanoseconds;
const ms: i32 = if (left <= 0) 0 else @intCast(@min(@divFloor(left + std.time.ns_per_ms - 1, std.time.ns_per_ms), std.math.maxInt(i32)));
var fds = [_]std.posix.pollfd{.{ .fd = sock.handle, .events = std.posix.POLL.OUT, .revents = 0 }};
const ready = std.posix.poll(&fds, ms) catch return true;
return ready > 0;
}

pub fn writeAll(self: IoStream, bytes: []const u8) std.Io.Writer.Error!void {
try self.writer.writeAll(bytes);
try self.writer.flush();
Expand Down Expand Up @@ -832,6 +881,9 @@ const FakeStream = struct {
/// Captured client->server bytes.
written: *std.ArrayList(u8),
allocator: std.mem.Allocator,
/// Whether the peer is taking what is written. False stands in for a
/// relay that has stopped reading.
room: bool = true,

fn read(self: *FakeStream, buffer: []u8, deadline: ?std.Io.Clock.Timestamp) error{}!usize {
// Never blocks, so a deadline is meaningless here.
Expand All @@ -843,6 +895,11 @@ const FakeStream = struct {
return n;
}

fn writableBy(self: *FakeStream, deadline: std.Io.Clock.Timestamp) bool {
_ = deadline;
return self.room;
}

fn writeAll(self: *FakeStream, bytes: []const u8) !void {
try self.written.appendSlice(self.allocator, bytes);
}
Expand Down Expand Up @@ -1108,6 +1165,70 @@ test "pings and pongs arriving back to back still stop at the deadline" {
try std.testing.expectEqualStrings("after", m.value.notice.message);
}

test "with a deadline, a pong is skipped rather than waited on when the peer has no room" {
const allocator = std.testing.allocator;
var written: std.ArrayList(u8) = .empty;
defer written.deinit(allocator);
var script: std.ArrayList(u8) = .empty;
defer script.deinit(allocator);
try script.appendSlice(allocator, &[_]u8{ 0x89, 0x01, 'p' }); // ping
try appendServerText(&script, allocator, "[\"NOTICE\",\"after\"]");

// A peer that has stopped reading: a pong written to it would wait for
// room that never comes.
var server = FakeStream{ .to_read = script.items, .written = &written, .allocator = allocator, .room = false };
var conn = newConn(allocator, &server);
defer conn.deinit();

var m = (try conn.receiveTimeout(.{ .duration = .{ .raw = .fromMilliseconds(1000), .clock = .awake } })).?;
defer m.deinit();
try std.testing.expectEqualStrings("after", m.value.notice.message);
try std.testing.expectEqual(@as(usize, 0), written.items.len);
}

test "with a deadline, a pong is skipped rather than queued behind another write" {
const allocator = std.testing.allocator;
var written: std.ArrayList(u8) = .empty;
defer written.deinit(allocator);
var script: std.ArrayList(u8) = .empty;
defer script.deinit(allocator);
try script.appendSlice(allocator, &[_]u8{ 0x89, 0x01, 'p' }); // ping
try appendServerText(&script, allocator, "[\"NOTICE\",\"after\"]");

var server = FakeStream{ .to_read = script.items, .written = &written, .allocator = allocator };
var conn = newConn(allocator, &server);
defer conn.deinit();

// Another writer holds the lock, as a send stuck on a full socket would.
try std.testing.expect(conn.write_lock.tryLock());
defer conn.write_lock.unlock(std.testing.io);

var m = (try conn.receiveTimeout(.{ .duration = .{ .raw = .fromMilliseconds(1000), .clock = .awake } })).?;
defer m.deinit();
try std.testing.expectEqualStrings("after", m.value.notice.message);
try std.testing.expectEqual(@as(usize, 0), written.items.len);
}

test "with a deadline and room to write, a ping is still answered" {
const allocator = std.testing.allocator;
var written: std.ArrayList(u8) = .empty;
defer written.deinit(allocator);
var script: std.ArrayList(u8) = .empty;
defer script.deinit(allocator);
try script.appendSlice(allocator, &[_]u8{ 0x89, 0x02, 'p', 'q' }); // ping
try appendServerText(&script, allocator, "[\"NOTICE\",\"after\"]");

var server = FakeStream{ .to_read = script.items, .written = &written, .allocator = allocator };
var conn = newConn(allocator, &server);
defer conn.deinit();

var m = (try conn.receiveTimeout(.{ .duration = .{ .raw = .fromMilliseconds(1000), .clock = .awake } })).?;
defer m.deinit();
const pong = (try websocket.decodeFrame(written.items)).?;
try std.testing.expectEqual(websocket.Opcode.pong, pong.opcode);
try std.testing.expectEqualStrings("pq", pong.payload);
}

test "receive answers a ping with a pong and continues" {
const allocator = std.testing.allocator;
var written: std.ArrayList(u8) = .empty;
Expand Down
2 changes: 1 addition & 1 deletion src/root.zig
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ const std = @import("std");

/// Kept in step with `build.zig.zon` by hand, and it had drifted three
/// releases behind, so anything reading it was told the wrong number.
pub const version = "0.14.6";
pub const version = "0.14.7";

pub const bech32 = @import("bech32.zig");
pub const nip19 = @import("nip19.zig");
Expand Down
Loading