Skip to content

Commit af349e9

Browse files
authored
broadcast: session reconstruction signaling + disposal (SESS reset 0x01) (#71)
Adds the session lifecycle state machine, aligned with ethp2p `broadcast/session.go`: - `SessionStage` (consuming/decoding/reconstructed/origin) with ordering semantics (`isDecoded` == stage >= reconstructed), - `sessCodeReconstructed = 0x01` — the QUIC app error code a peer's SESS stream is reset with once reconstructed (`broadcast/peer_in.go`), - `signalReconstructed`, per-peer `markPeerCompleted`/`dropPeer`, and `maybeDispose` (not-decoded → all peers dropped; decoded → all attached peers completed; peerless sessions left for TTL cleanup), - `ChannelRs.disposeSession` — removes + frees the session and emits `OnSessionDisposed` (completing another Observer callback from #58). `sessionDecode` now calls `signalReconstructed` (no-op for origin). Origin sessions start at `.origin`, relay sessions at `.consuming`. The actual per-peer SESS-stream reset with 0x01 is emitted by the QUIC transport when it observes reconstruction — a transport wiring follow-up, since `SessionRs` is transport-agnostic. Tests: stage ordering + reset-code constant, origin disposal gated on peer completion (with `OnSessionDisposed` emit), and relay consuming→reconstructed + drop-based disposal — all leak-checked. README implementation-status table updated. Closes #62.
1 parent 68cb82b commit af349e9

4 files changed

Lines changed: 178 additions & 1 deletion

File tree

README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ Zig helpers for the wire formats of **[ethp2p](https://github.com/ethp2p/ethp2p)
1717
| `Verdict` / `ChunkHandle` / `SessionRole` | [`broadcast/types.go`](https://github.com/ethp2p/ethp2p/blob/main/broadcast/types.go), [`broadcast/observer.go`](https://github.com/ethp2p/ethp2p/blob/main/broadcast/observer.go) | `layer.broadcast_types` |
1818
| Broadcast event `Observer` interface + `NoOpObserver` (15 callbacks, type-erased vtable) — emitted for channel-attach, session start/decode, chunk receive, peer subscribe/unsubscribe (remaining callbacks land with their subsystems) | [`broadcast/observer.go`](https://github.com/ethp2p/ethp2p/blob/main/broadcast/observer.go) | `broadcast.observer` |
1919
| `FullMessage` + channel `Subscribe` / decoded-message delivery (bounded FIFO `Subscription`, drop-on-full; `AlreadySubscribed` / `UnbufferedSubscription`) | [`broadcast/channel.go`](https://github.com/ethp2p/ethp2p/blob/main/broadcast/channel.go) `Subscribe` / `deliverMessage`, [`broadcast/types.go`](https://github.com/ethp2p/ethp2p/blob/main/broadcast/types.go) `FullMessage` | `broadcast.channel_rs` (`Subscription`), `layer.broadcast_types` |
20+
| Session lifecycle: `SessionStage` machine, `signalReconstructed`, per-peer complete/drop tracking, `maybeDispose`, `disposeSession` (emits `OnSessionDisposed`); `sessCodeReconstructed = 0x01` SESS-reset code (per-peer stream reset itself is a transport follow-up) | [`broadcast/session.go`](https://github.com/ethp2p/ethp2p/blob/main/broadcast/session.go), [`broadcast/peer_in.go`](https://github.com/ethp2p/ethp2p/blob/main/broadcast/peer_in.go) `sessCodeReconstructed` | `broadcast.session_rs`, `broadcast.channel_rs` |
2021
| `DedupCancel` + `DedupGroup` token (session hook) | session → strategy `takeChunk` | `layer.broadcast_types`, `layer.dedup` |
2122
| Engine-wide dedup registry (`channel` + `message` + chunk index) | multi-peer ingest dedup | `layer.dedup_registry`, `Engine.enable_cross_session_dedup` |
2223
| Verify result FIFO (single-threaded `Verified()` shim) | async verify channels | `layer.verify_queue` |

src/broadcast/channel_rs.zig

Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -183,6 +183,12 @@ pub const ChannelRs = struct {
183183
}
184184
const stats = try self.allocator.alloc(broadcast_types.PeerSessionStats, n);
185185
errdefer self.allocator.free(stats);
186+
const completed = try self.allocator.alloc(bool, n);
187+
errdefer self.allocator.free(completed);
188+
@memset(completed, false);
189+
const dropped = try self.allocator.alloc(bool, n);
190+
errdefer self.allocator.free(dropped);
191+
@memset(dropped, false);
186192

187193
while (populated < n) : (populated += 1) {
188194
member_ids[populated] = try self.allocator.dupe(u8, self.members.items[populated]);
@@ -201,6 +207,10 @@ pub const ChannelRs = struct {
201207
.strategy = strat.?,
202208
.member_ids = member_ids,
203209
.stats = stats,
210+
.stage = .origin,
211+
.completed = completed,
212+
.dropped = dropped,
213+
.ever_had_peers = n > 0,
204214
};
205215
strat = null;
206216

@@ -227,6 +237,12 @@ pub const ChannelRs = struct {
227237
}
228238
const stats = try self.allocator.alloc(broadcast_types.PeerSessionStats, n);
229239
errdefer self.allocator.free(stats);
240+
const completed = try self.allocator.alloc(bool, n);
241+
errdefer self.allocator.free(completed);
242+
@memset(completed, false);
243+
const dropped = try self.allocator.alloc(bool, n);
244+
errdefer self.allocator.free(dropped);
245+
@memset(dropped, false);
230246

231247
while (populated < n) : (populated += 1) {
232248
member_ids[populated] = try self.allocator.dupe(u8, self.members.items[populated]);
@@ -245,6 +261,10 @@ pub const ChannelRs = struct {
245261
.strategy = strat,
246262
.member_ids = member_ids,
247263
.stats = stats,
264+
.stage = .consuming,
265+
.completed = completed,
266+
.dropped = dropped,
267+
.ever_had_peers = n > 0,
248268
};
249269

250270
try self.sessions.put(self.allocator, mid, sess);
@@ -271,13 +291,27 @@ pub const ChannelRs = struct {
271291
const slot = self.sessions.getPtr(message_id) orelse return error.InvalidMessage;
272292
const out = try slot.*.strategy.decode();
273293
errdefer self.allocator.free(out);
294+
// Relay session reconstructed (no-op for origin); mirrors Go
295+
// `signalReconstructed` after a successful decode.
296+
slot.*.signalReconstructed();
274297
// Latency is not tracked in the Zig port yet (no per-session start clock);
275298
// emit 0 for now — Go passes `time.Since(session.start)`.
276299
self.engine.config.observer.sessionDecoded(self.id, message_id, 0);
277300
try self.deliverMessage(message_id, out);
278301
return out;
279302
}
280303

304+
/// Remove `message_id`'s session and free it, emitting `OnSessionDisposed`.
305+
/// Mirrors ethp2p `Channel.disposeSession`; no-op if there is no session.
306+
pub fn disposeSession(self: *ChannelRs, message_id: []const u8, reason: []const u8) void {
307+
if (self.sessions.fetchRemove(message_id)) |kv| {
308+
self.engine.config.observer.sessionDisposed(self.id, kv.key, reason);
309+
kv.value.deinit();
310+
self.allocator.destroy(kv.value);
311+
self.allocator.free(kv.key);
312+
}
313+
}
314+
281315
/// After a successful decode, clears engine `DedupRegistry` keys for this `(channel_id, message_id)`
282316
/// when `EngineConfig.enable_cross_session_dedup` is set (no-op otherwise).
283317
pub fn sessionDecodeClearEngineDedup(self: *ChannelRs, message_id: []const u8) ![]u8 {
@@ -493,6 +527,63 @@ test "subscriber drops decoded messages when full" {
493527
try std.testing.expectEqual(@as(usize, 1), sub.dropped());
494528
}
495529

530+
const session_rs_mod = @import("session_rs.zig");
531+
532+
test "origin session disposes once its peer completes; disposeSession emits" {
533+
const gpa = std.testing.allocator;
534+
var rec: @import("observer.zig").Recording = .{};
535+
var eng = try Engine.init(gpa, "local", .{ .observer = rec.observer() });
536+
defer eng.deinit();
537+
const ch = try eng.attachChannelRs("topic", sub_test_cfg);
538+
try ch.addMember("p1");
539+
540+
const payload = [_]u8{ 'h', 'i', '!', '!', '!' };
541+
try ch.publish("m1", &payload);
542+
543+
const sess = ch.sessions.get("m1").?;
544+
try std.testing.expectEqual(session_rs_mod.SessionStage.origin, sess.stage);
545+
// Decoded (origin) but the peer has not completed → not disposable yet.
546+
try std.testing.expect(!sess.maybeDispose());
547+
sess.markPeerCompleted("p1");
548+
try std.testing.expect(sess.maybeDispose());
549+
550+
ch.disposeSession("m1", "reconstructed");
551+
try std.testing.expect(ch.sessions.get("m1") == null);
552+
try std.testing.expectEqual(@as(usize, 1), rec.session_disposed);
553+
}
554+
555+
test "relay session: consuming stage, reconstruct signal, drop-based disposal" {
556+
const gpa = std.testing.allocator;
557+
var eng = try Engine.init(gpa, "local", .{});
558+
defer eng.deinit();
559+
const ch = try eng.attachChannelRs("topic", sub_test_cfg);
560+
try ch.addMember("peerA");
561+
try ch.addMember("peerB");
562+
563+
const payload = [_]u8{ 9, 8, 7 };
564+
var origin = try RsStrategy.newOrigin(gpa, sub_test_cfg, &payload);
565+
defer origin.deinit();
566+
try ch.attachRelaySession("m1", &origin.preamble);
567+
568+
const sess = ch.sessions.get("m1").?;
569+
try std.testing.expectEqual(session_rs_mod.SessionStage.consuming, sess.stage);
570+
571+
// Not decoded: disposes only when every peer is dropped.
572+
try std.testing.expect(!sess.maybeDispose());
573+
sess.dropPeer("peerA");
574+
try std.testing.expect(!sess.maybeDispose());
575+
sess.dropPeer("peerB");
576+
try std.testing.expect(sess.maybeDispose());
577+
578+
// Reconstruction moves it to the decoded branch.
579+
sess.signalReconstructed();
580+
try std.testing.expectEqual(session_rs_mod.SessionStage.reconstructed, sess.stage);
581+
try std.testing.expect(sess.maybeDispose());
582+
583+
ch.disposeSession("m1", "ttl_expired");
584+
try std.testing.expect(ch.sessions.get("m1") == null);
585+
}
586+
496587
test "relay ingest registry blocks second claim same chunk index" {
497588
const gpa = std.testing.allocator;
498589
var eng = try Engine.init(gpa, "local", .{});

src/broadcast/observer.zig

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -144,6 +144,7 @@ const testing = std.testing;
144144
pub const Recording = struct {
145145
session_started: usize = 0,
146146
session_decoded: usize = 0,
147+
session_disposed: usize = 0,
147148
chunk_rcvd: usize = 0,
148149
peer_subscribed: usize = 0,
149150
peer_unsubscribed: usize = 0,
@@ -164,7 +165,7 @@ pub const Recording = struct {
164165
.onPeerGone = noop_impl.peerGone,
165166
.onSessionStarted = onSessionStarted,
166167
.onSessionDecoded = onSessionDecoded,
167-
.onSessionDisposed = noop_impl.sessionDisposed,
168+
.onSessionDisposed = onSessionDisposed,
168169
.onChunkSent = noop_impl.chunkSent,
169170
.onChunkRcvd = onChunkRcvd,
170171
.onChunkError = noop_impl.chunkError,
@@ -190,6 +191,9 @@ pub const Recording = struct {
190191
fn onSessionDecoded(ctx: *anyopaque, _: []const u8, _: []const u8, _: i64) void {
191192
cast(ctx).session_decoded += 1;
192193
}
194+
fn onSessionDisposed(ctx: *anyopaque, _: []const u8, _: []const u8, _: []const u8) void {
195+
cast(ctx).session_disposed += 1;
196+
}
193197
fn onChunkRcvd(ctx: *anyopaque, _: []const u8, _: []const u8, _: []const u8, verdict: Verdict) void {
194198
const self = cast(ctx);
195199
self.chunk_rcvd += 1;

src/broadcast/session_rs.zig

Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,19 +17,92 @@ pub const SendRsChunkFn = *const fn (
1717
payload: []const u8,
1818
) anyerror!void;
1919

20+
/// QUIC application error code a peer's SESS stream is reset with once the
21+
/// session is reconstructed, telling the remote to stop sending chunks.
22+
/// Mirrors ethp2p `sessCodeReconstructed` (`broadcast/peer_in.go`).
23+
pub const sess_code_reconstructed: u64 = 0x01;
24+
25+
/// Linear session lifecycle (ethp2p `sessionStage`). Ordering is meaningful:
26+
/// `stage >= .reconstructed` means "decoded" (a relay that reconstructed, or an
27+
/// origin that always had the payload — `.origin` sorts highest).
28+
pub const SessionStage = enum(u8) {
29+
consuming = 0, // relay, accepting chunks
30+
decoding = 1, // relay, decode running
31+
reconstructed = 2, // relay, decode succeeded
32+
origin = 3, // origin, has the payload from the start
33+
34+
/// Whether the session holds the full message (origin or reconstructed relay).
35+
pub fn isDecoded(self: SessionStage) bool {
36+
return @intFromEnum(self) >= @intFromEnum(SessionStage.reconstructed);
37+
}
38+
};
39+
2040
pub const SessionRs = struct {
2141
allocator: Allocator,
2242
/// Owned by the parent `ChannelRs` map key; not freed here.
2343
message_id: []const u8,
2444
strategy: RsStrategy,
2545
member_ids: [][]u8,
2646
stats: []broadcast_types.PeerSessionStats,
47+
stage: SessionStage = .consuming,
48+
/// Per-member flags, parallel to `member_ids` (Go tracks these in `peers`).
49+
completed: []bool = &.{},
50+
dropped: []bool = &.{},
51+
/// Whether any peer was ever attached (Go `everHadPeers`); disposal is
52+
/// gated on this so peerless sessions are left for TTL cleanup.
53+
ever_had_peers: bool = false,
2754

2855
pub fn deinit(self: *SessionRs) void {
2956
self.strategy.deinit();
3057
for (self.member_ids) |m| self.allocator.free(m);
3158
self.allocator.free(self.member_ids);
3259
self.allocator.free(self.stats);
60+
self.allocator.free(self.completed);
61+
self.allocator.free(self.dropped);
62+
}
63+
64+
/// Mark the session reconstructed after a successful decode
65+
/// (ethp2p `signalReconstructed`). No-op for origin sessions.
66+
pub fn signalReconstructed(self: *SessionRs) void {
67+
if (self.stage == .origin) return;
68+
self.stage = .reconstructed;
69+
}
70+
71+
fn peerIndex(self: *const SessionRs, peer_id: []const u8) ?usize {
72+
for (self.member_ids, 0..) |m, i| {
73+
if (std.mem.eql(u8, m, peer_id)) return i;
74+
}
75+
return null;
76+
}
77+
78+
/// Mark a peer as having completed this session (ethp2p `handlePeerCompleted`).
79+
pub fn markPeerCompleted(self: *SessionRs, peer_id: []const u8) void {
80+
if (self.peerIndex(peer_id)) |i| self.completed[i] = true;
81+
}
82+
83+
/// Mark a peer as dropped (disconnected). Dropped peers are excluded from
84+
/// disposal accounting (Go removes them from the `peers` map).
85+
pub fn dropPeer(self: *SessionRs, peer_id: []const u8) void {
86+
if (self.peerIndex(peer_id)) |i| self.dropped[i] = true;
87+
}
88+
89+
/// Whether the session has no remaining work and can be disposed
90+
/// (ethp2p `maybeDispose`):
91+
/// - not decoded: dispose once every peer has been dropped,
92+
/// - decoded: dispose once every still-attached peer has completed.
93+
/// Peerless sessions never auto-dispose (left for TTL cleanup).
94+
pub fn maybeDispose(self: *const SessionRs) bool {
95+
if (!self.ever_had_peers) return false;
96+
if (!self.stage.isDecoded()) {
97+
for (self.dropped) |d| {
98+
if (!d) return false; // a live peer still needs chunks from us
99+
}
100+
return true;
101+
}
102+
for (self.completed, self.dropped) |c, d| {
103+
if (!d and !c) return false; // an attached peer has not finished
104+
}
105+
return true;
33106
}
34107

35108
/// One scheduling step: emit dispatches, then mark sends successful (simulates transport ACK).
@@ -78,3 +151,11 @@ pub const SessionRs = struct {
78151
return total;
79152
}
80153
};
154+
155+
test "session stage ordering and reconstructed reset code" {
156+
try std.testing.expectEqual(@as(u64, 0x01), sess_code_reconstructed);
157+
try std.testing.expect(!SessionStage.consuming.isDecoded());
158+
try std.testing.expect(!SessionStage.decoding.isDecoded());
159+
try std.testing.expect(SessionStage.reconstructed.isDecoded());
160+
try std.testing.expect(SessionStage.origin.isDecoded());
161+
}

0 commit comments

Comments
 (0)