pico

created pr with 105.1 on 2026-01-25T17:03:08Z · by c8ef7d19
added 105.2 on 2026-01-27T04:11:48Z · by c8ef7d19
1: 603ca6e = 1: 603ca6e chore(pubsub): add more tests
2: 501c042 ! 2: 17e00b2 feat(pubsub): round robin
added 105.3 on 2026-01-29T01:03:33Z · by c8ef7d19
1: 603ca6e = 1: 603ca6e chore(pubsub): add more tests
2: 17e00b2 = 2: 17e00b2 feat(pubsub): round robin
-: ------- > 3: e3136bd fix(pubsub): check for eof before processing and skip empty byte reads
-: ------- > 4: 9a6d19e fix: rr
-: ------- > 5: 5b3f3a1 fix: sending 0 byte read
added 105.4 on 2026-02-01T17:16:43Z · by c8ef7d19
1: 603ca6e = 1: 603ca6e chore(pubsub): add more tests
2: 17e00b2 = 2: 17e00b2 feat(pubsub): round robin
3: e3136bd = 3: e3136bd fix(pubsub): check for eof before processing and skip empty byte reads
4: 9a6d19e = 4: 9a6d19e fix: rr
5: 5b3f3a1 = 5: 5b3f3a1 fix: sending 0 byte read
-: ------- > 6: 4fff471 refactor: fixes
-: ------- > 7: d0dfc85 chore: SetDispatch on Broker
changed status to accepted on 2026-02-23T02:02:07Z · by c8ef7d19
cmds
checkout latest patchset:
ssh pr.pico.sh print 105 | git am -3
checkout any patchset in a patch request:
ssh pr.pico.sh print 105.[rev] | git am -3
add changes to patch request:
git format-patch main --stdout | ssh pr.pico.sh pr add 105
set PR to open (enables RSS notifications):
ssh pr.pico.sh pr open 105
set PR to draft (stops RSS notifications):
ssh pr.pico.sh pr draft 105
+0 -14 pkg/pubsub/broker.go #
......@@ -105,20 +105,6 @@ func (b *BaseBroker) Connect(client *Client, channels []*Channel) (error, error)
105105 data := make([]byte, 32*1024)
106106 n, err := client.ReadWriter.Read(data)
107107
108- // TODO: Skip empty reads
109- /*
110- if err != nil {
111- if errors.Is(err, io.EOF) {
112- return
113- }
114- inputErr = err
115- return
116- }
117-
118- if n == 0 {
119- continue
120- }
121- */
122108 data = data[:n]
123109
124110 channelMessage := ChannelMessage{
+12 -8 pkg/pubsub/channel.go #
......@@ -71,11 +71,6 @@ func (c *Channel) Cleanup() {
7171 }
7272
7373 func (c *Channel) Handle() {
74- // If no dispatcher is set, use multicast as default
75- if c.Dispatcher == nil {
76- c.Dispatcher = &MulticastDispatcher{}
77- }
78-
7974 c.handleOnce.Do(func() {
8075 go func() {
8176 defer func() {
......@@ -100,10 +95,19 @@ func (c *Channel) Handle() {
10095 }
10196
10297 // Collect eligible subscribers
103- subscribers := dispatcherForGetClients(c.GetClients(), data)
98+ subscribers := make([]*Client, 0)
99+ for _, client := range c.GetClients() {
100+ // Skip input-only clients and senders (unless replay is enabled)
101+ if client.Direction == ChannelDirectionInput || (client.ID == data.ClientID && !client.Replay) {
102+ continue
103+ }
104+ subscribers = append(subscribers, client)
105+ }
104106
105- // Dispatch message using the configured dispatcher
106- _ = c.Dispatcher.Dispatch(data, subscribers, c.Done)
107+ if len(data.Data) > 0 {
108+ // Dispatch message using the configured dispatcher
109+ _ = c.Dispatcher.Dispatch(data, subscribers, c.Done)
110+ }
107111 }
108112 }
109113 }()
+0 -18 pkg/pubsub/dispatcher.go #
......@@ -1,26 +1,8 @@
11 package pubsub
22
3-import (
4- "iter"
5-)
6-
73 // MessageDispatcher defines how messages are dispatched to subscribers.
84 type MessageDispatcher interface {
95 // Dispatch sends a message to the appropriate subscriber(s).
106 // It receives the message, all subscribers, and the channel's sync primitives.
117 Dispatch(msg ChannelMessage, subscribers []*Client, channelDone chan struct{}) error
128 }
13-
14-// dispatcherForGetClients collects eligible clients for dispatching.
15-// Returns clients that should receive messages (output direction, not the sender unless replay).
16-func dispatcherForGetClients(getClients iter.Seq2[string, *Client], msg ChannelMessage) []*Client {
17- subscribers := make([]*Client, 0)
18- for _, client := range getClients {
19- // Skip input-only clients and senders (unless replay is enabled)
20- if client.Direction == ChannelDirectionInput || (client.ID == msg.ClientID && !client.Replay) {
21- continue
22- }
23- subscribers = append(subscribers, client)
24- }
25- return subscribers
26-}
+9 -4 pkg/pubsub/multicast.go #
......@@ -59,9 +59,14 @@ func (p *Multicast) GetSubs() iter.Seq2[string, *Client] {
5959 return p.getClients(ChannelDirectionOutput)
6060 }
6161
62-func (p *Multicast) connect(ctx context.Context, ID string, rw io.ReadWriter, channels []*Channel, direction ChannelDirection, blockWrite bool, replay, keepAlive bool) (error, error) {
62+func (p *Multicast) connect(ctx context.Context, ID string, rw io.ReadWriter, channels []*Channel, direction ChannelDirection, blockWrite bool, replay, keepAlive bool, dispatcher MessageDispatcher) (error, error) {
6363 client := NewClient(ID, rw, direction, blockWrite, replay, keepAlive)
6464
65+ // Set dispatcher on all channels
66+ for _, ch := range channels {
67+ ch.Dispatcher = dispatcher
68+ }
69+
6570 go func() {
6671 <-ctx.Done()
6772 client.Cleanup()
......@@ -71,15 +76,15 @@ func (p *Multicast) connect(ctx context.Context, ID string, rw io.ReadWriter, ch
7176 }
7277
7378 func (p *Multicast) Pipe(ctx context.Context, ID string, rw io.ReadWriter, channels []*Channel, replay bool) (error, error) {
74- return p.connect(ctx, ID, rw, channels, ChannelDirectionInputOutput, false, replay, false)
79+ return p.connect(ctx, ID, rw, channels, ChannelDirectionInputOutput, false, replay, false, &MulticastDispatcher{})
7580 }
7681
7782 func (p *Multicast) Pub(ctx context.Context, ID string, rw io.ReadWriter, channels []*Channel, blockWrite bool) error {
78- return errors.Join(p.connect(ctx, ID, rw, channels, ChannelDirectionInput, blockWrite, false, false))
83+ return errors.Join(p.connect(ctx, ID, rw, channels, ChannelDirectionInput, blockWrite, false, false, &MulticastDispatcher{}))
7984 }
8085
8186 func (p *Multicast) Sub(ctx context.Context, ID string, rw io.ReadWriter, channels []*Channel, keepAlive bool) error {
82- return errors.Join(p.connect(ctx, ID, rw, channels, ChannelDirectionOutput, false, false, keepAlive))
87+ return errors.Join(p.connect(ctx, ID, rw, channels, ChannelDirectionOutput, false, false, keepAlive, &MulticastDispatcher{}))
8388 }
8489
8590 // MulticastDispatcher sends each message to all eligible subscribers.
Back to top