Skip to content

feat(bps): lite pubsub - #5626

Draft
acud wants to merge 21 commits into
masterfrom
bps-simplified
Draft

acud wants to merge 21 commits into
masterfrom
bps-simplified

Conversation

@acud

@acud acud commented Sep 23, 2026

Copy link
Copy Markdown
Contributor

Checklist

  • I have read the coding guide.
  • My change requires a documentation update, and I have done it.
  • I have added tests to cover my changes.
  • I have filled out the description and linked the related issues.

Description

Open API Spec Version Changes (if applicable)

Motivation and Context (Optional)

Related Issue (Optional)

Screenshots (if appropriate):

AI Disclosure

  • This PR contains code that has been generated by an LLM.
  • I have reviewed the AI generated code thoroughly.
  • I possess the technical expertise to responsibly review the code generated in this PR.

@zelig

zelig commented Sep 24, 2026 •

Copy link
Copy Markdown
Member

Review against SWIP-74 rev 3 (ethersphere/SWIPs #111), which now specifies the challenge–response claim this PR started from. The findings were checked adversarially against the code, the spec and bee's p2p wrapper on master. (Edited: cursor semantics, test flow and counter list aligned with SWIP-74 rev 3 as pushed.)

State of the PR. Seven commits 2026-09-12..24; 424 lines of non-generated Go and
proto plus generated bps.pb.go; the description is the empty template; not wired into
the node; CI red — TestJoin fails with write join msg: stream closed (fullNode is
never set, so the handler resets every stream), and lint fails on the commit message
bps init. It is a scaffold — most
functions are stubs with "replace later" comments — so this review is about the shape
it commits to, not the stubs. It is also a restart: acud's previous bee PR #5597 (closed
2026-09-10, ~9,900 lines) implemented SWIP-60 rev 3 in full — bindings, broker, cohort,
session, WS bridge, metrics, hostile tests — and none of it is carried over yet; the
validation, queue and metrics machinery there is the obvious quarry for what is missing.

Headline

The scaffold follows SWIP-74 in three places and diverges from it in one that matters;
and it has one lifecycle error that no amount of filling in the stubs will fix.

Follows: one create-or-attach join (Jopen in the registry: "joins an existing cohort
or creates one by joining a previously unknown topic"); a per-cohort lastSeen index
"used to prevent replay and circumvent dedup logic" — that is the cursor; one frame type
both directions after the handshake (SWIP-74's Message, though the frame must carry the
whole chunk — see the wire table).

Diverges — and the divergence is the right one. The publisher role is claimed by a
challenge–response (JoinAck{challenge} → Claim{sig}), which is what SWIP-74 rev 3
now says: a joining peer is a receiver at once, and a peer that signs the challenge
upgrades to a publisher stream. SWIP-74 rev 2's static-preimage Auth-in-Join was
replayable by design, and a replayed signature would have bought an identity — a stream
carried as the admin's, exempt from the fan-out bound, and in SWIP-60 admitted to a
closed cohort; the challenge makes the identity worth exactly the key. What the scaffold
gets wrong is the shape of the challenge, in four ways (SWIP-74 rev 3, Handshake):

  • the challenge is a per-cohort value (Cohort.challenge, rotated after each claim)
    handed to every joiner — it must be derived per address from a boot secret:
    S = H(S_C ‖ S_c ‖ addr) with S_C drawn once at boot and never persisted,
    S_c = H(S_C ‖ H(Marshal(spec))) (the spec's canonical serialisation, which also keys
    the registry), addr the address the joiner declares in its Join; the broker stores
    nothing and recomputes S at claim time, and an address that moves to another node
    keeps its claim;
  • Claim opens a second stream — it must be sent on the joined stream, which
    upgrades in place; the second stream is why the admin's Join stream would otherwise
    stay a fan-out target receiving the admin's own messages;
  • the claim takes over the cohort and rotates the challenge "so that the publisher
    can reclaim later" — no takeover: a reconnecting admin joins and claims on its new
    stream, concurrent admin streams are all publisher streams, the cursor arbitrates;
  • Claim{sig} carries a bare signature over a bare challenge, on a stream with no
    topic — a claim is Claim{addr, index, auth} with auth = {r, s, v} over
    H("bps-claim:v1" ‖ S ‖ O_B ‖ index): domain-separated from SOC signatures, bound
    through O_B (the broker's overlay) to the verifier — otherwise a relay can forward
    the challenge and land the admin's signature on its own stream at the honest broker —
    and carrying the publisher's cursor, signed: the claim that its next message has an
    index of at least index, which a reconnecting admin uses to move the broker's cursor
    forward.

Lifecycle: the broker handler is fire-and-forget — it spawns the per-stream goroutine
on the handler's context and returns nil at once. In bee's libp2p wrapper the handler
context is cancelled the moment the handler returns (libp2p.go removes the stream and
calls its cancel), so on a real node every subscriber and publisher stream is torn down
right after the ack. streamtest hands the handler context.Background(), which is why
the test does not show it — and why the goroutine will instead leak and trip goleak
once fullNode is set. Every long-lived bee handler (pushsync, pullsync, retrieval)
blocks for the stream's lifetime for exactly this reason, and it is also the only way to
return p2p.NewBlockPeerError for a violation, which the spec needs (below).

Wire, message by message

PR #5626 SWIP-74 note
SystemMessage{oneof Join ǀ Claim} no envelope Join is the only first frame; after it, what a frame is follows from the stream's role — a subscriber stream sends at most one Claim, a publisher stream sends Message — so nothing needs a discriminator
Join{topic} Join{CohortSpec{topic, binding, admin}, addr, claim?} the broker needs topic to rebuild keccak256(topic ‖ index) and admin to check the owner; a hash of both (the registry comment: "feed topic + owner hash") gives it neither. The spec is the cohort's identity and the registry key — carry it, hash its canonical serialisation for the key. addr declares the address the stream will publish as; a returning publisher claims right here
JoinAck{challenge} Ack{status, challenge} — OK / FULL / REJECTED the challenge is right, and issued only to a Join that declared an address; the status is missing: the peer cannot tell FULL (capacity — back off and rejoin) from REJECTED (the spec itself is refused — rejoining unchanged is pointless); SWIP-74 mandates backoff after either
Claim{sig} on a fresh stream Claim{addr, index, auth} on the joined stream same stream, upgrade in place; signature over H("bps-claim:v1" ‖ S ‖ O_B ‖ index), not over the bare challenge — as sent, Claim names neither cohort, verifier nor cursor
ClaimAck{} nothing the client waits for it and the handler never sends it: Claim blocks until its context ends. No reply at all: a publisher pipelines its first Message behind the Claim, and a bad claim resets the stream
Broadcast{soc} Message{address, data} the chunk data alone has no address, so the ordinary SOC validation is vacuous (ecrecover always yields an address): carry the address, and the frame is a whole chunk validated by the ordinary path
bps/1.0.0, stream bps pubsub/1.0.0 bee registers /swarm/<name>/<version>/<stream>; decide the name once, in the spec, and match it

The takeover semantics in the Claim comment ("the current stream becomes the publisher
stream and the next challenge changes randomly, so that ... the publisher can reclaim
later") is the supersede rule SWIP-74 considered and dropped: a reconnecting admin needs
no reclaim, its new stream is simply another publisher stream, and a stale one closing
later is a no-op.

Streams and state

  1. Broker streams die after the ack (above). Make the handler resident: the read loop
    inline on a publisher stream, the writer loop draining a bounded per-stream queue
    inline on a subscriber stream; exits are the peer's EOF, the queue overflow, the
    inactivity reclaim. A violation returns p2p.NewBlockPeerError(...) so the wrapper
    resets and blocklists; a non-OK ack returns an error so the wrapper closes.
  2. Subscriber streams are never read: the join branch is write-only. SWIP-74 item 5
    (a Message on a subscriber stream: drop, reset, blocklist, count wrong_stream) is
    unsatisfiable in this shape, and a subscriber that closes is not noticed until the
    next fan-out write fails, so idle streams keep counting against the per-cohort bound.
  3. Client-side lifetimes: every error path after NewStream returns without
    stream.Reset() (bee's idiom is Reset on error, FullClose on a clean end — the
    CI failure is such a path); the subscription's whole life is the caller's context, so
    a request-scoped one kills it on return and Background makes it uncancellable, and
    Join returns no handle; rxCh is never closed, so a consumer ranging over it hangs
    after the stream dies; in Claim, once the writer goroutine exits on a broker reset
    nobody reads ch and the publisher's next send blocks forever. Join/Claim should
    return a handle (close, done, error), the Service should own a long-lived context.
  4. Subscriber drops deliveries: select { case rxCh <- msg.Soc: default: } on an
    unbuffered channel discards any frame the consumer is not ready for at that instant.
    Loss policy is the broker's (bounded queue, reset on overflow); the subscriber blocks
    or buffers, and re-verifies (SWIP-74 item 10).
  5. Identity is modelled as a node: Cohort.publisher swarm.Address, Jopen(overlay, topic), members map[string]string. SWIP-74 binds the identity to the stream, never
    the node — the same key works from any node, and several streams may be the admin's
    at once — and a cohort is a set of streams each with its own queue, which a map of
    strings cannot hold. And lastSeen models the cursor the wrong way round: SWIP-74's
    cursor is the lowest index accepted next, a plain uint64 initially 0, moved to
    n + 1 on acceptance and to max(cursor, index) on a claim, so index 0 is accepted
    on a fresh cohort without an "absent" value.

Not started, and expected not to be

Validation on the publish path (id substitution, SOC check with the address, owner ==
admin, index ≥ cursor — lastSeen is declared and never read), fan-out from the registry,
the bounds and the inactivity deadline, the three statuses with backoff, and the eight
counters. All of it exists in some form in #5597.

Tests

TestJoin asserts a challenge round trip and never sends a Message, so it would pass
with no validation, no cursor and no fan-out; streamtest.Records waits for the handler
to return, so it can only pass with the fire-and-forget handler that has to go, and once
fullNode is set the leaked goroutines fail goleak. Suggested flow: join a subscriber,
join the admin declaring addr, claim on the challenge from its Ack (and once more with
the claim in the Join), publish index 0, receive it on the subscriber, close both, then
Records. Cases worth their own test: index 0 accepted on a fresh cohort; n < cursor
counted as a retransmit; owner ≠ admin dropped; a subscriber Message → reset and
blocklist; two concurrent admin streams arbitrated by the cursor; inactivity reclaim.

Small things

Claim sends Sig: topic; two log lines say "read join ack" in the write loop and the
broadcast loop; the package doc still says "hello-world wire protocol used to greet other
peers"; TestJoin prints the challenge and its failure message says "want 2" for a check
of 1; registry_test.go is a bare package line; a SystemMessage with neither member
set falls through to return nil with the stream neither closed nor reset.

What to ask for

  • Join{CohortSpec, addr, claim?}, Ack{status, challenge}, Claim{addr, index, auth},
    Message{address, data} — field numbers as SWIP-74 rev 3 fixes them; keep Jopen, key
    the registry by the canonical serialisation of the spec; rename Broadcast to Message;
    drop SystemMessage and ClaimAck.
  • A resident handler; a read loop on subscriber streams that accepts at most one Claim
    and treats everything else as a violation; BlockPeerError on violation.
  • No challenge state: one boot secret, S recomputed at claim time; the declared addr
    and the role per stream; several admin streams; no takeover; the cursor as "lowest index
    accepted next", set forward by the claim.
  • The p2p handshake should verify that a peer's signed address record names the
    connection's authenticated peer ID — the claim's binding to the broker's overlay rests on
    it, and bee's handshake checks the record, not the binding, today. Separate PR.
  • Client handles with close/done; Reset on error paths; close rxCh; no drop at the
    subscriber.
  • Then the validation steps and the cursor, fan-out to subscriber streams only, bounded
    per-stream queues, the bounds and the inactivity deadline, statuses with backoff, the
    counters — from feat: bps protocol draft #5597 where it fits.
  • The protocol name decided in one place.

Open

  • The protocol name, bps in code versus pubsub in both SWIPs.
  • The subscribers-per-cohort bound admits one extra stream while the admin is absent — a
    Join declaring the admin's address — and disconnects it if it has not claimed within
    the claim deadline (SWIP-74 rev 3, Resource bounds).

@zelig

zelig commented Sep 28, 2026 •

Copy link
Copy Markdown
Member

Update 2026-10-01 (SWIP-74 rev 8): SWIP-74 is at rev 8 (ethersphere/SWIPs #111). Against this branch: Broadcast becomes {soc, kind, challenge, index} with the chunk an ordinary SOC whose id is keccak256(keccak256(prefix ‖ topic ‖ challenge) ‖ index); Claim goes — the first valid Broadcast on a stream declaring the admin's address is the claim, an empty AUTH chunk if there is nothing to publish yet; Join is {cohort, addr}; JoinAck{challenge} stays, random per stream, 32 bytes; the challenge is not signed into a payload together with the broker's overlay, it salts the id. The WS bridge added in 2d7adcf signs a claim over challenge ‖ broker — that goes with Claim. The points below about validation, the cursor, bounds, statuses and counters stand.

Round 2, against the six commits since the review (d28ced3..14b89a5) and the seven "design decisions" from our chat, checked against SWIP-74 rev 4 (#111), which since yesterday carries the claim as a SOC — your idea, with the spec's fields inside it. Where a point was already in the review above and the code has not changed, it is only listed as still open.

What improved

  • One stream; the claim rides the joined stream (the review's first ask).
  • The handler is resident (the review's central finding): it waits for its goroutines.
  • The broker's overlay is in the signed material; each joiner gets its own challenge.
  • A real test: join, claim, publish, receive, and the publisher does not get its own
    message back.

The scheme, decision by decision

  1. Topic T = keccak(id ‖ owner), "the SOC address of the claim". True for the claim,
    false for everything after it. A publication is a feed update with id
    keccak(topic ‖ n) and address keccak(id ‖ admin); none of these ever hashes to T,
    so T can validate exactly one chunk shape — the claim. The code accordingly validates
    nothing on the publish path (ch1 <- m.Soc, bps.go:246) and the test enshrines it: the
    accepted broadcast is the bytes "hello cohort", not a SOC. A subscriber holding only T
    can never verify a delivery — it has neither the feed topic to rebuild the id nor the
    admin to compare the owner — so "the broker withholds, never forges" is gone and the
    broker becomes trusted. The only publications T could validate are same-id SOCs, one
    address for every message: the ANCHOR binding the review and SWIP-74 rejected.
  2. The challenge in the SOC payload. Fresh and tamper-proof, yes — and a genuine,
    storable, stampable single-owner chunk signed by the admin at the channel's own
    address T, with a payload the broker chose. The broker, the one party that receives it,
    can upload it with its own stamp and overwrite whatever the admin keeps at T; every new
    claim replaces it again. SWIP-74's "bps-claim:v1" ‖ S ‖ O_B ‖ index is an 84-byte
    preimage that is never a chunk digest; that separation is what the scheme deletes. The
    claim also carries no index, and the code has no cursor at all (lastSeen is now
    commented out).
  3. Owner compared to T, not merely recovered. Correct, and what the spec does with
    addr.
  4. A random nonce per member. Argued against acud's own previous code (one rotating
    per-cohort value), not against the spec: SWIP-74's S is already per address and never
    handed to two parties. Against the spec it is a trade-off he states one side of: it
    closes replay by the admin's own node (row 2 of the table) and costs a wallet
    signature per reconnect, state per stream, and the identity-from-any-node property.
    And it is not per stream: it is stored per (topic, overlay), so it works on whichever
    stream of that node is registered last.
  5. Broker overlay in the payload. Correct; SWIP-74's O_B.
  6. Same stream. Correct; what the review asked for.
  7. "The broker already knows T, so sending the address would be redundant." Right
    for the claim; and, it turns out, right for publications too — see "A point for the
    spec" below.

The arguments in chat

  • "The admin address leaks when explicit." Every delivery reveals the owner to every
    subscriber (soc.FromChunk recovers it), and a subscriber must know the admin to verify
    anything. T hides the admin from someone who holds the invite and never sees the stream.
  • "Subscribers may open the cohort with the wrong settings and block the publisher."
    A different spec is a different cohort that can accept no message; the admin's cohort is
    created by the first Join with the right spec, whoever sends it. Viktor's "it's a
    different stream channel" is the spec's answer.
  • "Topic must be unique per broker." It is not; cohorts are keyed by the canonical
    serialisation of the spec, and a broker holds several on one topic.
  • "Wall-of-text invite." SWIP-74's invite is topic + admin + broker, 84 bytes; T + broker
    is 64. Twenty bytes.
  • "More cohort state makes revocation a PITA." The spec is the cohort's key, not state;
    there is no revocation in lite, and in SWIP-60 it is a roster update on the admin's
    feed, never a spec change.
  • "We can do without the spec for the first iteration and change it later." The first
    iteration is the base. Going from Join{bytes topic = 1} to Join{CohortSpec cohort = 1; addr = 2; claim = 3} changes field 1's type and meaning; JoinAck{bytes challenge = 1} to Ack{Status status = 1} changes the wire type; Broadcast{bytes soc = 1} to
    Message{bytes address = 1; data = 2} makes every old subscriber read a 32-byte address
    as the chunk body. That is a bps/2.0.0, and the second implementation the dev line
    calls for cannot interoperate on a single frame.

The code

  • Members keyed by overlay, not stream (bps.go:269, 288, 303): a node's second stream
    to the same cohort overwrites the first's entry, the first stream's exit deletes the
    second's, and one ordering deadlocks a publisher handler for good. SWIP-74 needs one
    stream per (peer, cohort, identity) — the admin's node that also watches the stream is
    the ordinary case.
  • Publications unvalidated (above); no cursor.
  • Silent loss at both ends: a 1-slot buffer and select-default drop at the broker's
    fan-out and at the client's receive channel. The spec is lossy at the stream level, not
    the message level: a 64-frame queue whose overflow resets the stream and is counted, so
    the subscriber knows and rejoins. With a buffer of one, loss is the normal case under
    any burst — a live video stream would shed frames continuously and silently.
  • No status, no bounds, no deadlines, no backoff, no counters; JoinAck cannot say
    FULL or REJECTED.
  • Violations reset inline and return nil, so the libp2p wrapper never sees a
    BlockPeerError and nobody is blocklisted.
  • A publisher stream stays in the fan-out set (excluded only from its own frames by
    pointer equality); two admin streams feed each other.
  • Client side, still as reviewed: no Reset on error paths, rxCh never closed, no
    handle to close a subscription.

Where the two designs meet — a construction for the SWIP

acud's best idea is to carry the claim as a SOC body so that one existing call verifies
it. It works with the spec's fields intact — and it is now the spec: the claim is a Broadcast
whose chunk payload is a service message of kind CLAIM (SWIP-74 rev 6), and Broadcast is
the data frame's name again:

id      = keccak256("bps-claim:v1" ‖ topic)      44-byte preimage; a feed id's is 40 bytes
owner   = addr                                    so the address is keccak256(id ‖ addr)
payload = S ‖ O_B ‖ index                         72 bytes

The broker (or any receiver of a claim) computes the expected address from the id it
derives and the declared addr, and runs soc.Valid(NewChunk(expected, body)) — steps
"signer == addr" and "the digest is right" in one call. The domain separator lives in the
id, so the claim can never be a feed update; it is storable, but only at a namespaced
address nobody reads as content. And it settles the signing convention that SWIP-74 rev 3
marks (?): a claim is signed exactly as a SOC. Cost: the claim frame is a SOC body (~185
bytes) instead of ~120.

A point for the spec

The same reasoning shows the address field on Broadcast is computable by every
receiver in every configuration: the expected owner is always known — the admin, the
stream's claimed address, the declared address under ALL — so a receiver forms
keccak256(keccak256(topic ‖ n) ‖ owner) itself and runs soc.Valid against it, which is
steps 3 and 4 in one call and forces the owner rather than checking it afterwards.
Your decision 7 is the same observation, and it is now the spec's: SWIP-74 rev 5 carries
the chunk data alone in Broadcast{soc}, and every receiver validates it at the address it
forms from the feed id and the owner it already knows.

What to ask acud

  • Reply on the PR; the review is the record.
  • Join{CohortSpec, addr, auth?} (auth = the claim chunk's bytes) now — not for access control, but because without the
    feed topic and the admin neither the broker nor a subscriber can validate one
    publication.
  • The claim as a domain-separated SOC (construction above, now SWIP-74 rev 6's CLAIM
    service message on Broadcast), with index, over the derived S — keeping your
    one-call verification and dropping the per-stream nonce.
  • Validate publications: rewrite the bare index, soc.Valid against the expected
    address, n ≥ cursor, cursor := n + 1.
  • Key members by stream; a bounded queue with reset; BlockPeerError on violation;
    statuses; the bounds and the two deadlines; the counters.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants