Add event stream via pubsub
- Add /event file for pub/sub event notifications - Events: n (new), u (update), r (rename), d (delete) - Format: '<type> <identifier>' per line - Blocking read with 30s timeout
This commit is contained in:
parent
1eb8c71e8a
commit
8a6c90f160
2
go.mod
2
go.mod
|
|
@ -3,3 +3,5 @@ module denotesrv
|
||||||
go 1.23.0
|
go 1.23.0
|
||||||
|
|
||||||
require 9fans.net/go v0.0.5
|
require 9fans.net/go v0.0.5
|
||||||
|
|
||||||
|
require github.com/simonfxr/pubsub v0.0.5 // indirect
|
||||||
|
|
|
||||||
8
go.sum
8
go.sum
|
|
@ -2,7 +2,13 @@
|
||||||
9fans.net/go v0.0.5/go.mod h1:Rxvbbc1e+1TyGMjAvLthGTyO97t+6JMQ6ly+Lcs9Uf0=
|
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=
|
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/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/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-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-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
|
||||||
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/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/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-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||||
golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/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=
|
||||||
|
|
|
||||||
|
|
@ -51,6 +51,7 @@ import (
|
||||||
|
|
||||||
"9fans.net/go/plan9"
|
"9fans.net/go/plan9"
|
||||||
"9fans.net/go/plan9/client"
|
"9fans.net/go/plan9/client"
|
||||||
|
ps "github.com/simonfxr/pubsub"
|
||||||
)
|
)
|
||||||
|
|
||||||
// File types in our filesystem
|
// File types in our filesystem
|
||||||
|
|
@ -92,6 +93,7 @@ type server struct {
|
||||||
callbacks Callbacks
|
callbacks Callbacks
|
||||||
filterQuery string
|
filterQuery string
|
||||||
filteredNotes []*metadata.Metadata
|
filteredNotes []*metadata.Metadata
|
||||||
|
events *ps.Bus
|
||||||
}
|
}
|
||||||
|
|
||||||
type connState struct {
|
type connState struct {
|
||||||
|
|
@ -100,11 +102,13 @@ type connState struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
type fid struct {
|
type fid struct {
|
||||||
qid plan9.Qid
|
qid plan9.Qid
|
||||||
path string
|
path string
|
||||||
offset int64
|
offset int64
|
||||||
mode uint8
|
mode uint8
|
||||||
writeBuf []byte // accumulates Twrite chunks, dispatched on Tclunk
|
writeBuf []byte // accumulates Twrite chunks, dispatched on Tclunk
|
||||||
|
eventCh chan string // for event subscribers
|
||||||
|
eventCancel func() // unsubscribe callback
|
||||||
}
|
}
|
||||||
|
|
||||||
var srv *server
|
var srv *server
|
||||||
|
|
@ -181,6 +185,7 @@ func NewServer(denoteDir string, callbacks Callbacks) (*Server, error) {
|
||||||
notes: notes,
|
notes: notes,
|
||||||
denoteDir: denoteDir,
|
denoteDir: denoteDir,
|
||||||
callbacks: callbacks,
|
callbacks: callbacks,
|
||||||
|
events: ps.NewBus(),
|
||||||
}
|
}
|
||||||
return &Server{s}, nil
|
return &Server{s}, nil
|
||||||
}
|
}
|
||||||
|
|
@ -199,6 +204,7 @@ func StartServer(initialData metadata.Results, denoteDir string, callbacks Callb
|
||||||
notes: initialData,
|
notes: initialData,
|
||||||
denoteDir: denoteDir,
|
denoteDir: denoteDir,
|
||||||
callbacks: callbacks,
|
callbacks: callbacks,
|
||||||
|
events: ps.NewBus(),
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get namespace and create Unix socket path
|
// 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)
|
qids = append(qids, qid)
|
||||||
path = "/dir"
|
path = "/dir"
|
||||||
found = true
|
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" {
|
} else if name == "n" {
|
||||||
qid := plan9.Qid{
|
qid := plan9.Qid{
|
||||||
Type: QTDir,
|
Type: QTDir,
|
||||||
|
|
@ -495,6 +509,24 @@ func (s *server) read(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
||||||
return errorFcall(fc, "fid not found")
|
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()
|
defer cs.mu.Unlock()
|
||||||
|
|
||||||
var data []byte
|
var data []byte
|
||||||
|
|
@ -622,6 +654,7 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall {
|
||||||
if s.callbacks.OnNew != nil {
|
if s.callbacks.OnNew != nil {
|
||||||
go s.callbacks.OnNew(identifier)
|
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))}
|
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 {
|
if s.callbacks.OnRename != nil {
|
||||||
s.callbacks.OnRename(noteID)
|
s.callbacks.OnRename(noteID)
|
||||||
}
|
}
|
||||||
|
s.events.Publish("events", "u "+noteID)
|
||||||
|
s.events.Publish("events", "r "+noteID)
|
||||||
case "keywords":
|
case "keywords":
|
||||||
if input == "" {
|
if input == "" {
|
||||||
note.Tags = []string{}
|
note.Tags = []string{}
|
||||||
|
|
@ -667,6 +702,8 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall {
|
||||||
if s.callbacks.OnRename != nil {
|
if s.callbacks.OnRename != nil {
|
||||||
s.callbacks.OnRename(noteID)
|
s.callbacks.OnRename(noteID)
|
||||||
}
|
}
|
||||||
|
s.events.Publish("events", "u "+noteID)
|
||||||
|
s.events.Publish("events", "r "+noteID)
|
||||||
case "signature":
|
case "signature":
|
||||||
note.Signature = input
|
note.Signature = input
|
||||||
if s.callbacks.OnUpdate != nil {
|
if s.callbacks.OnUpdate != nil {
|
||||||
|
|
@ -675,12 +712,15 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall {
|
||||||
if s.callbacks.OnRename != nil {
|
if s.callbacks.OnRename != nil {
|
||||||
s.callbacks.OnRename(noteID)
|
s.callbacks.OnRename(noteID)
|
||||||
}
|
}
|
||||||
|
s.events.Publish("events", "u "+noteID)
|
||||||
|
s.events.Publish("events", "r "+noteID)
|
||||||
case "ctl":
|
case "ctl":
|
||||||
switch input {
|
switch input {
|
||||||
case "d":
|
case "d":
|
||||||
if s.callbacks.OnDelete != nil {
|
if s.callbacks.OnDelete != nil {
|
||||||
s.callbacks.OnDelete(noteID)
|
s.callbacks.OnDelete(noteID)
|
||||||
}
|
}
|
||||||
|
s.events.Publish("events", "d "+noteID)
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
for i, n := range s.notes {
|
for i, n := range s.notes {
|
||||||
if n.Identifier == noteID {
|
if n.Identifier == noteID {
|
||||||
|
|
@ -699,6 +739,7 @@ func (s *server) dispatchWrite(f *fid, tag uint16) *plan9.Fcall {
|
||||||
if s.callbacks.OnRename != nil {
|
if s.callbacks.OnRename != nil {
|
||||||
s.callbacks.OnRename(noteID)
|
s.callbacks.OnRename(noteID)
|
||||||
}
|
}
|
||||||
|
s.events.Publish("events", "r "+noteID)
|
||||||
}
|
}
|
||||||
case "body":
|
case "body":
|
||||||
if err := s.writeBody(note.Path, input, note); err != nil {
|
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()
|
cs.mu.Lock()
|
||||||
}
|
}
|
||||||
|
// Clean up event subscription if any
|
||||||
|
if ok && f.eventCancel != nil {
|
||||||
|
f.eventCancel()
|
||||||
|
}
|
||||||
delete(cs.fids, fc.Fid)
|
delete(cs.fids, fc.Fid)
|
||||||
cs.mu.Unlock()
|
cs.mu.Unlock()
|
||||||
|
|
||||||
|
|
@ -809,6 +854,19 @@ func (s *server) readDir(path string, offset int64, count uint32) []byte {
|
||||||
Muid: "denote",
|
Muid: "denote",
|
||||||
Length: uint64(len(s.denoteDir)),
|
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
|
// add n directory
|
||||||
dirs = append(dirs, plan9.Dir{
|
dirs = append(dirs, plan9.Dir{
|
||||||
Qid: plan9.Qid{
|
Qid: plan9.Qid{
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue