git-pr

created pr with 25.1 on 2024-09-09T04:26:26Z · by c8ef7d19
cmds
checkout latest patchset:
ssh pr.pico.sh print 25 | git am -3
checkout any patchset in a patch request:
ssh pr.pico.sh print 25.[rev] | git am -3
add changes to patch request:
git format-patch main --stdout | ssh pr.pico.sh pr add 25
+18 -1 cmd/authorized_keys/main.go #
......@@ -32,6 +32,7 @@ func PubSubMiddleware(cfg *pubsub.Cfg) wish.Middleware {
3232 return
3333 }
3434
35+ ctx := sesh.Context()
3536 cmd := strings.TrimSpace(args[0])
3637 channel := args[1]
3738 logger := cfg.Logger.With(
......@@ -49,6 +50,13 @@ func PubSubMiddleware(cfg *pubsub.Cfg) wish.Middleware {
4950 Writer: sesh,
5051 Chan: make(chan error),
5152 }
53+ go func() {
54+ <-ctx.Done()
55+ err := cfg.PubSub.UnSub(sub)
56+ if err != nil {
57+ wish.Errorln(sesh, err)
58+ }
59+ }()
5260 err := cfg.PubSub.Sub(sub)
5361 if err != nil {
5462 wish.Errorln(sesh, err)
......@@ -58,6 +66,13 @@ func PubSubMiddleware(cfg *pubsub.Cfg) wish.Middleware {
5866 Name: channel,
5967 Reader: sesh,
6068 }
69+ go func() {
70+ <-ctx.Done()
71+ err := cfg.PubSub.UnPub(msg)
72+ if err != nil {
73+ wish.Errorln(sesh, err)
74+ }
75+ }()
6176 err := cfg.PubSub.Pub(msg)
6277 if err != nil {
6378 wish.Errorln(sesh, err)
......@@ -88,7 +103,9 @@ func main() {
88103 wish.WithAddress(fmt.Sprintf("%s:%s", host, port)),
89104 wish.WithHostKeyPath("ssh_data/term_info_ed25519"),
90105 wish.WithAuthorizedKeys(keyPath),
91- wish.WithMiddleware(PubSubMiddleware(cfg)),
106+ wish.WithMiddleware(
107+ PubSubMiddleware(cfg),
108+ ),
92109 )
93110 if err != nil {
94111 logger.Error(err.Error())
+20 -0 multicast.go #
......@@ -1,8 +1,10 @@
11 package pubsub
22
33 import (
4+ "fmt"
45 "io"
56 "log/slog"
7+ "time"
68
79 "github.com/google/uuid"
810 )
......@@ -67,6 +69,7 @@ func (b *PubSubMulticast) Pub(msg *Msg) error {
6769 writers := []io.Writer{}
6870 for _, sub := range b.subs {
6971 if b.PubMatcher(msg, sub) {
72+ log.Info("found match", "sub", sub.ID)
7073 matches = append(matches, sub)
7174 writers = append(writers, sub.Writer)
7275 }
......@@ -78,6 +81,11 @@ func (b *PubSubMulticast) Pub(msg *Msg) error {
7881 log.Info("no subs found, waiting for sub")
7982 sub = <-b.Chan
8083 if b.PubMatcher(msg, sub) {
84+ // empty subscriber is a signal to force a pub to stop
85+ // waiting for a sub
86+ if sub.Writer == nil {
87+ return fmt.Errorf("pub closed")
88+ }
8189 return b.Pub(msg)
8290 }
8391 }
......@@ -97,6 +105,18 @@ func (b *PubSubMulticast) Pub(msg *Msg) error {
97105 log.Error("unsub err", "err", err)
98106 }
99107 }
108+ del := time.Now()
109+ msg.Delivered = &del
100110
101111 return err
102112 }
113+
114+func (b *PubSubMulticast) UnPub(msg *Msg) error {
115+ b.Logger.Info("unpub", "channel", msg.Name)
116+ // if the message hasn't been delivered then send a cancel sub to
117+ // the multicast channel
118+ if msg.Delivered == nil {
119+ b.Chan <- &Subscriber{Name: msg.Name}
120+ }
121+ return nil
122+}
+7 -4 pubsub.go #
......@@ -3,6 +3,7 @@ package pubsub
33 import (
44 "io"
55 "log/slog"
6+ "time"
67 )
78
89 type Subscriber struct {
......@@ -18,15 +19,17 @@ func (s *Subscriber) Wait() error {
1819 }
1920
2021 type Msg struct {
21- Name string
22- Reader io.Reader
22+ Name string
23+ Reader io.Reader
24+ Delivered *time.Time
2325 }
2426
2527 type PubSub interface {
2628 GetSubs() []*Subscriber
27- Sub(l *Subscriber) error
28- UnSub(l *Subscriber) error
29+ Sub(sub *Subscriber) error
30+ UnSub(sub *Subscriber) error
2931 Pub(msg *Msg) error
32+ UnPub(msg *Msg) error
3033 // return true if message should be sent to this subscriber
3134 PubMatcher(msg *Msg, sub *Subscriber) bool
3235 }
Back to top