Skip to content

feat(bps): align with SWIP-74 rev 8 (BPS-lite) - #5644

Merged
acud merged 4 commits into
bps-simplifiedfrom
bps-swip74
Oct 9, 2026
Merged

acud merged 4 commits into
bps-simplifiedfrom
bps-swip74

Conversation

@acud

@acud acud commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

Brings pkg/bps in line with SWIP-74 rev 8 per @zelig's review: Claim goes, Join carries the spec and addr, and the challenge salts the chunk id instead of being signed into a payload.

  • Wire: Join{cohort, addr}, Ack{status, challenge}, Broadcast{soc, kind, challenge, index}; pubsub/1.0.0.
  • Broker: cohorts keyed by canonical spec; random per-stream challenge; pending admin streams outside the fan-out bound with a 30s claim deadline; frame validation in spec order (wrong stream / invalid soc → reset + blocklist; unknown kind, wrong challenge, retransmit → dropped and counted); one cursor per cohort, no dedup window; bounds per cohort/broker/peer; 64-frame subscriber queues; 10 min inactivity reclaim; per-cohort counters + prometheus.
  • Subscriber: re-verifies every delivery against spec and its own cursor.
  • API/OpenAPI: subscribe frames are challenge ‖ index ‖ soc (optional cursor); publish takes kind ‖ index ‖ soc frames, no JSON claim.

Not done (see pkg/bps/TODO): client rejoin backoff, spec open points taken as drafted, SWIP-60 interop.

Tests: go test -race ./pkg/bps/... ./pkg/api/ ./pkg/node/ pass. make lint not run (golangci-lint crashes on the local go1.27 toolchain).

🤖 Generated with Claude Code

- wire: Join{cohort, addr}, Ack{status, challenge}, Broadcast{soc, kind,
  challenge, index}; protocol id pubsub/1.0.0
- per-stream random 32-byte challenge salts the session feed topic; the
  first valid frame (DATA or empty AUTH) claims the publisher role
- pending admin streams outside the fan-out bound with a claim deadline
- per-cohort cursor replaces dedup; ordered validation with violations
  reset and blocklisted
- resource bounds, inactivity reclaim, per-cohort counters and metrics
- subscribers re-verify deliveries against spec and cursor
- API websocket bridge and OpenAPI updated to the new frames

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

@zelig zelig left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

preliminary look

Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/pb/bps.proto Outdated
Comment thread pkg/bps/pb/bps.proto Outdated
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/TODO Outdated
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go
InactivityDeadline: 10 * time.Minute,
ClaimDeadline: 30 * time.Second,
QueueSize: 64,
ViolationBlocklist: 10 * time.Minute,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It is a a Timeout, shall we end on that word?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Still open, cosmetic: ViolationBlocklist and BlocklistDuration are an asymmetric pair — ViolationBlocklistDuration / AuthTimeoutBlocklistDuration would say what each is.

Comment thread pkg/bps/bps.go Outdated

@zelig zelig left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

some extra comments

Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go Outdated
// stream was reset must back off before rejoining.
func (s *Service) Join(ctx context.Context, req JoinRequest) (Session, error) {
spec := req.spec()
if err := validateJoin(&pb.Join{Cohort: spec, Addr: req.Addr}); err != nil {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

you should not just send any Addr, it should be added in the API by SwarmID or something

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Still open: /bps/subscribe takes any 20 bytes as identity. Default it to the node's own address, and reject identity == owner there — as it stands that makes a pending stream that receives nothing and gets this node blocklisted at the broker after 30 s.

Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go Outdated
Comment thread pkg/bps/bps.go
cursor uint64 // the lowest DATA index accepted next
lastActivity time.Time
members map[*member]struct{}
subscribers int

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is this needed ? members - publishers

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Still open, low: the counter is the O(1) way to enforce the fan-out bound with one map; splitting members by role would drop it and make "outside the fan-out bound" structural. Nice to have.

@acud
acud merged commit 0adf704 into bps-simplified Oct 9, 2026
10 of 11 checks passed
@acud
acud deleted the bps-swip74 branch October 9, 2026 12:43
@zelig zelig mentioned this pull request Oct 9, 2026
7 tasks
zelig pushed a commit that referenced this pull request Oct 10, 2026
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