From 8a6c90f16091797c4863ae7080700037a41a57c4 Mon Sep 17 00:00:00 2001 From: Levi Neely Date: Fri, 20 Mar 2026 18:00:34 +0100 Subject: [PATCH] Add event stream via pubsub - Add /event file for pub/sub event notifications - Events: n (new), u (update), r (rename), d (delete) - Format: ' ' per line - Blocking read with 30s timeout --- go.mod | 2 ++ go.sum | 8 +++++ internal/p9/server/server.go | 68 +++++++++++++++++++++++++++++++++--- 3 files changed, 73 insertions(+), 5 deletions(-) diff --git a/go.mod b/go.mod index 75e083a..f724fa0 100644 --- a/go.mod +++ b/go.mod @@ -3,3 +3,5 @@ module denotesrv go 1.23.0 require 9fans.net/go v0.0.5 + +require github.com/simonfxr/pubsub v0.0.5 // indirect diff --git a/go.sum b/go.sum index 5089ea2..e8c7e00 100644 --- a/go.sum +++ b/go.sum @@ -2,7 +2,13 @@ 9fans.net/go v0.0.5/go.mod h1:Rxvbbc1e+1TyGMjAvLthGTyO97t+6JMQ6ly+Lcs9Uf0= dmitri.shuralyov.com/gpu/mtl v0.0.0-20201218220906-28db891af037/go.mod h1:H6x//7gZCb22OMCxBHrMx7a5I7Hp++hsVxbQ4BYO7hU= github.com/BurntSushi/xgb v0.0.0-20160522181843-27f122750802/go.mod h1:IVnqGOEym/WlBOVXweHU+Q+/VP0lqqI8lqeDx9IjBqo= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/go-gl/glfw/v3.3/glfw v0.0.0-20200222043503-6f7a984d4dc4/go.mod h1:tQ2UAYgL5IevRw8kRxooKSPJfGvJ9fJQFa0TUsXzTg8= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/simonfxr/pubsub v0.0.5 h1:DJfvFoglqGvwJriIOC5NI5um34n2YX8KAU5+7jv768w= +github.com/simonfxr/pubsub v0.0.5/go.mod h1:bQ+B2NEEHZ08VY/0xVFGaGL1KMJ7cl3yLaXe8xvyC24= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= @@ -32,3 +38,5 @@ golang.org/x/tools v0.0.0-20200117012304-6edc0a871e69/go.mod h1:TB2adYChydJhpapK golang.org/x/tools v0.0.0-20200207183749-b753a1ba74fa/go.mod h1:TB2adYChydJhpapKDTa4BR/hXlZSLoq2Wpct/0txZ28= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/internal/p9/server/server.go b/internal/p9/server/server.go index 1a028f7..41da298 100644 --- a/internal/p9/server/server.go +++ b/internal/p9/server/server.go @@ -51,6 +51,7 @@ import ( "9fans.net/go/plan9" "9fans.net/go/plan9/client" + ps "github.com/simonfxr/pubsub" ) // File types in our filesystem @@ -92,6 +93,7 @@ type server struct { callbacks Callbacks filterQuery string filteredNotes []*metadata.Metadata + events *ps.Bus } type connState struct { @@ -100,11 +102,13 @@ type connState struct { } type fid struct { - qid plan9.Qid - path string - offset int64 - mode uint8 - writeBuf []byte // accumulates Twrite chunks, dispatched on Tclunk + qid plan9.Qid + path string + offset int64 + mode uint8 + writeBuf []byte // accumulates Twrite chunks, dispatched on Tclunk + eventCh chan string // for event subscribers + eventCancel func() // unsubscribe callback } var srv *server @@ -181,6 +185,7 @@ func NewServer(denoteDir string, callbacks Callbacks) (*Server, error) { notes: notes, denoteDir: denoteDir, callbacks: callbacks, + events: ps.NewBus(), } return &Server{s}, nil } @@ -199,6 +204,7 @@ func StartServer(initialData metadata.Results, denoteDir string, callbacks Callb notes: initialData, denoteDir: denoteDir, callbacks: callbacks, + events: ps.NewBus(), } // Get namespace and create Unix socket path @@ -403,6 +409,14 @@ func (s *server) walk(cs *connState, fc *plan9.Fcall) *plan9.Fcall { qids = append(qids, qid) path = "/dir" found = true + } else if name == "event" { + qid := plan9.Qid{ + Type: QTFile, + Path: uint64(qidEvent), + } + qids = append(qids, qid) + path = "/event" + found = true } else if name == "n" { qid := plan9.Qid{ Type: QTDir, @@ -495,6 +509,24 @@ func (s *server) read(cs *connState, fc *plan9.Fcall) *plan9.Fcall { return errorFcall(fc, "fid not found") } + // Handle event stream (blocking read with pubsub) + if f.path == "/event" { + if f.eventCh == nil { + ch := make(chan string, 64) + sub := s.events.SubscribeChan("events", ch, ps.CloseOnUnsubscribe) + f.eventCh = ch + f.eventCancel = func() { s.events.Unsubscribe(sub) } + } + cs.mu.Unlock() + select { + case msg := <-f.eventCh: + data := []byte(msg + "\n") + return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(data)), Data: data} + case <-time.After(30 * time.Second): + return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0, Data: []byte{}} + } + } + defer cs.mu.Unlock() var data []byte @@ -622,6 +654,7 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall { if s.callbacks.OnNew != nil { go s.callbacks.OnNew(identifier) } + s.events.Publish("events", "n "+identifier) return &plan9.Fcall{Type: plan9.Rwrite, Tag: fc.Tag, Count: uint32(len(f.writeBuf))} } @@ -651,6 +684,8 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall { if s.callbacks.OnRename != nil { s.callbacks.OnRename(noteID) } + s.events.Publish("events", "u "+noteID) + s.events.Publish("events", "r "+noteID) case "keywords": if input == "" { note.Tags = []string{} @@ -667,6 +702,8 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall { if s.callbacks.OnRename != nil { s.callbacks.OnRename(noteID) } + s.events.Publish("events", "u "+noteID) + s.events.Publish("events", "r "+noteID) case "signature": note.Signature = input if s.callbacks.OnUpdate != nil { @@ -675,12 +712,15 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall { if s.callbacks.OnRename != nil { s.callbacks.OnRename(noteID) } + s.events.Publish("events", "u "+noteID) + s.events.Publish("events", "r "+noteID) case "ctl": switch input { case "d": if s.callbacks.OnDelete != nil { s.callbacks.OnDelete(noteID) } + s.events.Publish("events", "d "+noteID) s.mu.Lock() for i, n := range s.notes { if n.Identifier == noteID { @@ -699,6 +739,7 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall { if s.callbacks.OnRename != nil { s.callbacks.OnRename(noteID) } + s.events.Publish("events", "r "+noteID) } case "body": if err := s.writeBody(note.Path, input, note); err != nil { @@ -744,6 +785,10 @@ func (s *server) clunk(cs *connState, fc *plan9.Fcall) *plan9.Fcall { } cs.mu.Lock() } + // Clean up event subscription if any + if ok && f.eventCancel != nil { + f.eventCancel() + } delete(cs.fids, fc.Fid) cs.mu.Unlock() @@ -809,6 +854,19 @@ func (s *server) readDir(path string, offset int64, count uint32) []byte { Muid: "denote", Length: uint64(len(s.denoteDir)), }) + // add event node + dirs = append(dirs, plan9.Dir{ + Qid: plan9.Qid{ + Type: QTFile, + Path: uint64(qidEvent), + }, + Mode: 0444, + Name: "event", + Uid: "denote", + Gid: "denote", + Muid: "denote", + Length: 0, + }) // add n directory dirs = append(dirs, plan9.Dir{ Qid: plan9.Qid{