git-pr
created pr with
25.1
cmds
checkout latest patchset:
ssh pr.pico.sh print 25 | git am -3checkout any patchset in a patch request:
ssh pr.pico.sh print 25.[rev] | git am -3add changes to patch request:
git format-patch main --stdout | ssh pr.pico.sh pr add 25
Patchset
25.1
refactor: unsub method to interface
Eric Bower
2024-09-09T04:24:45Zfix(cmd): proper cleanup with ssh session is closed (e.g. ctr+c)
Semantic diff summary
1 added,
7 modified,
0 signature changed,
0 removed
across 3 analyzed files
+18
-1
cmd/authorized_keys/main.go
#
| ... | ... | @@ -49,6 +50,13 @@ func PubSubMiddleware(cfg *pubsub.Cfg) wish.Middleware { | |
| 49 | 50 | Writer: sesh, | |
| 50 | 51 | Chan: make(chan error), | |
| 51 | 52 | } | |
| 53 | + | go func() { | |
| 54 | + | <-ctx.Done() | |
| 55 | + | err := cfg.PubSub.UnSub(sub) | |
| 56 | + | if err != nil { | |
| 57 | + | wish.Errorln(sesh, err) | |
| 58 | + | } | |
| 59 | + | }() | |
| 52 | 60 | err := cfg.PubSub.Sub(sub) | |
| 53 | 61 | if err != nil { | |
| 54 | 62 | wish.Errorln(sesh, err) |
| ... | ... | @@ -58,6 +66,13 @@ func PubSubMiddleware(cfg *pubsub.Cfg) wish.Middleware { | |
| 58 | 66 | Name: channel, | |
| 59 | 67 | Reader: sesh, | |
| 60 | 68 | } | |
| 69 | + | go func() { | |
| 70 | + | <-ctx.Done() | |
| 71 | + | err := cfg.PubSub.UnPub(msg) | |
| 72 | + | if err != nil { | |
| 73 | + | wish.Errorln(sesh, err) | |
| 74 | + | } | |
| 75 | + | }() | |
| 61 | 76 | err := cfg.PubSub.Pub(msg) | |
| 62 | 77 | if err != nil { | |
| 63 | 78 | wish.Errorln(sesh, err) |
| ... | ... | @@ -88,7 +103,9 @@ func main() { | |
| 88 | 103 | wish.WithAddress(fmt.Sprintf("%s:%s", host, port)), | |
| 89 | 104 | wish.WithHostKeyPath("ssh_data/term_info_ed25519"), | |
| 90 | 105 | wish.WithAuthorizedKeys(keyPath), | |
| 91 | - | wish.WithMiddleware(PubSubMiddleware(cfg)), | |
| 106 | + | wish.WithMiddleware( | |
| 107 | + | PubSubMiddleware(cfg), | |
| 108 | + | ), | |
| 92 | 109 | ) | |
| 93 | 110 | if err != nil { | |
| 94 | 111 | logger.Error(err.Error()) |
+20
-0
multicast.go
#
| ... | ... | @@ -78,6 +81,11 @@ func (b *PubSubMulticast) Pub(msg *Msg) error { | |
| 78 | 81 | log.Info("no subs found, waiting for sub") | |
| 79 | 82 | sub = <-b.Chan | |
| 80 | 83 | 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 | + | } | |
| 81 | 89 | return b.Pub(msg) | |
| 82 | 90 | } | |
| 83 | 91 | } |
| ... | ... | @@ -97,6 +105,18 @@ func (b *PubSubMulticast) Pub(msg *Msg) error { | |
| 97 | 105 | log.Error("unsub err", "err", err) | |
| 98 | 106 | } | |
| 99 | 107 | } | |
| 108 | + | del := time.Now() | |
| 109 | + | msg.Delivered = &del | |
| 100 | 110 | ||
| 101 | 111 | return err | |
| 102 | 112 | } | |
| 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
#
| ... | ... | @@ -18,15 +19,17 @@ func (s *Subscriber) Wait() error { | |
| 18 | 19 | } | |
| 19 | 20 | ||
| 20 | 21 | type Msg struct { | |
| 21 | - | Name string | |
| 22 | - | Reader io.Reader | |
| 22 | + | Name string | |
| 23 | + | Reader io.Reader | |
| 24 | + | Delivered *time.Time | |
| 23 | 25 | } | |
| 24 | 26 | ||
| 25 | 27 | type PubSub interface { | |
| 26 | 28 | GetSubs() []*Subscriber | |
| 27 | - | Sub(l *Subscriber) error | |
| 28 | - | UnSub(l *Subscriber) error | |
| 29 | + | Sub(sub *Subscriber) error | |
| 30 | + | UnSub(sub *Subscriber) error | |
| 29 | 31 | Pub(msg *Msg) error | |
| 32 | + | UnPub(msg *Msg) error | |
| 30 | 33 | // return true if message should be sent to this subscriber | |
| 31 | 34 | PubMatcher(msg *Msg, sub *Subscriber) bool | |
| 32 | 35 | } |