From cf210e59993c33d566dbde44449318d9ccabe3f6 Mon Sep 17 00:00:00 2001 From: rachel-mp4 Date: Mon, 6 Oct 2025 13:35:19 -0400 Subject: [PATCH] add mediainit/mediapub --- go.mod | 2 +- go.sum | 8 +-- options.go | 27 ++++++--- server.go | 160 ++++++++++++++++++++++++++++++++++++++++------------- 4 files changed, 145 insertions(+), 52 deletions(-) diff --git a/go.mod b/go.mod index d55e758..6e3bb83 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,6 @@ go 1.24.2 require ( github.com/gorilla/websocket v1.5.3 - github.com/rachel-mp4/lrcproto v0.0.0-20250905151943-8e3a1989ea5a + github.com/rachel-mp4/lrcproto v1.2.0 google.golang.org/protobuf v1.36.6 ) diff --git a/go.sum b/go.sum index 27ec954..e74be1a 100644 --- a/go.sum +++ b/go.sum @@ -2,12 +2,8 @@ github.com/google/go-cmp v0.5.5 h1:Khx7svrCpmxxtHBq5j2mp/xVjsi8hQMfNLvJFAlrGgU= github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= -github.com/rachel-mp4/lrcproto v0.0.0-20250720164211-c6162669b709 h1:P//gJE0zFv9Qvfn8dvp9ZrnG0FZh2MVcAX+uOP2flRw= -github.com/rachel-mp4/lrcproto v0.0.0-20250720164211-c6162669b709/go.mod h1:hQzO36tQELGbkmRnUtKeM6NMU34t79ZcTlhM+MO7pHw= -github.com/rachel-mp4/lrcproto v0.0.0-20250905145450-74a49183ee1f h1:CZVsuwuS/5fDwB4X9AgobE1bjnhVFaIuedgZOuVtqeA= -github.com/rachel-mp4/lrcproto v0.0.0-20250905145450-74a49183ee1f/go.mod h1:hQzO36tQELGbkmRnUtKeM6NMU34t79ZcTlhM+MO7pHw= -github.com/rachel-mp4/lrcproto v0.0.0-20250905151943-8e3a1989ea5a h1:W3zPeGz/jHYyWj8ZfsuA9yigPMNlA/h7Fu+SQKgiBmQ= -github.com/rachel-mp4/lrcproto v0.0.0-20250905151943-8e3a1989ea5a/go.mod h1:hQzO36tQELGbkmRnUtKeM6NMU34t79ZcTlhM+MO7pHw= +github.com/rachel-mp4/lrcproto v1.2.0 h1:nZI80WQKO6yKgX0O5H6OO1cM/tiJqugs0p52KCoIDOw= +github.com/rachel-mp4/lrcproto v1.2.0/go.mod h1:hQzO36tQELGbkmRnUtKeM6NMU34t79ZcTlhM+MO7pHw= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543 h1:E7g+9GITq07hpfrRu66IVDexMakfv52eLZ2CXBWiKr4= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= google.golang.org/protobuf v1.36.6 h1:z1NpPI8ku2WgiWnf+t9wTPsn6eP1L7ksHUlkfLvd9xY= diff --git a/options.go b/options.go index b255f13..f66a7ad 100644 --- a/options.go +++ b/options.go @@ -8,14 +8,15 @@ import ( ) type options struct { - uri string - secret string - welcome *string - writer *io.Writer - verbose bool - pubChan chan PubEvent - initChan chan lrcpb.Event_Init - initialID *uint32 + uri string + secret string + welcome *string + writer *io.Writer + verbose bool + pubChan chan PubEvent + initChan chan lrcpb.Event_Init + mediainitChan chan lrcpb.Event_Mediainit + initialID *uint32 } type Option func(option *options) error @@ -55,6 +56,16 @@ func WithInitChannel(initChan chan lrcpb.Event_Init) Option { } } +func WithMediainitChannel(mediainitChan chan lrcpb.Event_Mediainit) Option { + return func(options *options) error { + if mediainitChan == nil { + return errors.New("must provide a channel") + } + options.mediainitChan = mediainitChan + return nil + } +} + func WithPubChannel(pubChan chan PubEvent) Option { return func(options *options) error { if pubChan == nil { diff --git a/server.go b/server.go index a7d493c..69fe901 100644 --- a/server.go +++ b/server.go @@ -3,6 +3,7 @@ package lrcd import ( "context" "errors" + "fmt" "github.com/gorilla/websocket" "github.com/rachel-mp4/lrcproto/gen/go" "google.golang.org/protobuf/proto" @@ -15,23 +16,23 @@ import ( ) type Server struct { - secret string - uri string - eventBus chan clientEvent - ctx context.Context - cancel context.CancelFunc - clients map[*client]bool - clientsMu sync.Mutex - idmapsMu sync.Mutex - clientToID map[*client]*uint32 - idToClient map[uint32]*client - lastID uint32 - logger *log.Logger - debugLogger *log.Logger - welcomeEvt []byte - pongEvt []byte - initChan chan lrcpb.Event_Init - pubChan chan PubEvent + secret string + uri string + eventBus chan clientEvent + ctx context.Context + cancel context.CancelFunc + clients map[*client]bool + clientsMu sync.Mutex + idmapsMu sync.Mutex + idToClient map[uint32]*client + lastID uint32 + logger *log.Logger + debugLogger *log.Logger + welcomeEvt []byte + pongEvt []byte + initChan chan lrcpb.Event_Init + mediainitChan chan lrcpb.Event_Mediainit + pubChan chan PubEvent } type PubEvent struct { @@ -47,6 +48,8 @@ type client struct { muteMap map[*client]bool mutedBy map[*client]bool myIDs []uint32 + textID *uint32 + mediaID *uint32 post *string nick *string externID *string @@ -85,6 +88,9 @@ func NewServer(opts ...Option) (*Server, error) { if options.initChan != nil { s.initChan = options.initChan } + if options.mediainitChan != nil { + s.mediainitChan = options.mediainitChan + } if options.pubChan != nil { s.pubChan = options.pubChan } @@ -97,7 +103,6 @@ func NewServer(opts ...Option) (*Server, error) { s.clients = make(map[*client]bool) s.clientsMu = sync.Mutex{} s.idmapsMu = sync.Mutex{} - s.clientToID = make(map[*client]*uint32) s.idToClient = make(map[uint32]*client) s.eventBus = make(chan clientEvent, 100) return &s, nil @@ -214,7 +219,6 @@ func (s *Server) WSHandler() http.HandlerFunc { s.handlePub(client) s.idmapsMu.Lock() - delete(s.clientToID, client) for _, id := range client.myIDs { // remove myself from the idToClient map delete(s.idToClient, id) } @@ -307,8 +311,12 @@ func (s *Server) broadcaster() { continue case *lrcpb.Event_Init: s.handleInit(msg, client) + case *lrcpb.Event_Mediainit: + s.handleMediainit(msg, client) case *lrcpb.Event_Pub: s.handlePub(client) + case *lrcpb.Event_Mediapub: + s.handleMediapub(msg, client) case *lrcpb.Event_Insert: s.handleInsert(msg, client) case *lrcpb.Event_Delete: @@ -330,17 +338,16 @@ func (s *Server) broadcaster() { } func (s *Server) handleInit(msg *lrcpb.Event_Init, client *client) { - s.idmapsMu.Lock() - curID := s.clientToID[client] + curID := client.textID if curID != nil { - s.idmapsMu.Unlock() return } + s.idmapsMu.Lock() newID := s.lastID + 1 s.lastID = newID - s.clientToID[client] = &newID s.idToClient[newID] = client s.idmapsMu.Unlock() + client.textID = &newID client.myIDs = append(client.myIDs, newID) newpost := "" client.post = &newpost @@ -393,16 +400,74 @@ func (s *Server) broadcastInit(msg *lrcpb.Event_Init, client *client) { } } } +func (s *Server) handleMediainit(msg *lrcpb.Event_Mediainit, client *client) { + curId := client.mediaID + if curId != nil { + return + } + s.idmapsMu.Lock() + newID := s.lastID + 1 + s.lastID = newID + s.idToClient[newID] = client + s.idmapsMu.Unlock() + client.mediaID = &newID + client.myIDs = append(client.myIDs, newID) + msg.Mediainit.Id = &newID + msg.Mediainit.Nick = client.nick + msg.Mediainit.ExternalID = client.externID + msg.Mediainit.Color = client.color + echoed := false + msg.Mediainit.Echoed = &echoed + msg.Mediainit.Nonce = nil + if s.mediainitChan != nil { + select { + case s.mediainitChan <- *msg: + default: + s.log("initchan blocked, closing channel") + close(s.mediainitChan) + s.mediainitChan = nil + } + } + s.broadcastMediainit(msg, client) +} + +func (s *Server) broadcastMediainit(msg *lrcpb.Event_Mediainit, client *client) { + stdEvent := &lrcpb.Event{Msg: msg} + stdData, _ := proto.Marshal(stdEvent) + echoed := true + msg.Mediainit.Echoed = &echoed + msg.Mediainit.Nonce = GenerateNonce(*msg.Mediainit.Id, s.uri, s.secret) + echoEvent := &lrcpb.Event{Msg: msg} + echoData, _ := proto.Marshal(echoEvent) + muteEvent := &lrcpb.Event{Msg: &lrcpb.Event_Mute{Mute: &lrcpb.Mute{Id: msg.Mediainit.GetId()}}} + muteData, _ := proto.Marshal(muteEvent) + s.clientsMu.Lock() + defer s.clientsMu.Unlock() + for c := range s.clients { + var dts []byte + if c == client { + dts = echoData + } else if client.mutedBy[c] { + dts = muteData + } else { + dts = stdData + } + select { + case c.dataChan <- dts: + s.logDebug("b mediainit") + default: + s.log("kicked client") + client.cancel() + } + } +} func (s *Server) handlePub(client *client) { - s.idmapsMu.Lock() - curID := s.clientToID[client] + curID := client.textID if curID == nil { - s.idmapsMu.Unlock() return } - s.clientToID[client] = nil - s.idmapsMu.Unlock() + client.textID = nil event := &lrcpb.Event{Msg: &lrcpb.Event_Pub{Pub: &lrcpb.Pub{Id: curID}}} if s.pubChan != nil { select { @@ -417,10 +482,35 @@ func (s *Server) handlePub(client *client) { s.broadcast(event, client) } +func (s *Server) handleMediapub(msg *lrcpb.Event_Mediapub, client *client) { + curID := client.mediaID + if curID == nil { + return + } + client.mediaID = nil + msg.Mediapub.Id = curID + body := "external media." + if msg.Mediapub.Alt != nil { + body += fmt.Sprintf(" alt=%s.", *msg.Mediapub.Alt) + } + if msg.Mediapub.ContentAddress != nil { + body += fmt.Sprintf(" cid=%s.", *msg.Mediapub.ContentAddress) + } + if s.pubChan != nil { + select { + case s.pubChan <- PubEvent{ID: *curID, Body: body}: + default: + s.log("pubchan blocked, closing channel") + close(s.pubChan) + s.pubChan = nil + } + } + event := &lrcpb.Event{Msg: msg} + s.broadcast(event, client) +} + func (s *Server) handleInsert(msg *lrcpb.Event_Insert, client *client) { - s.idmapsMu.Lock() - curID := s.clientToID[client] - s.idmapsMu.Unlock() + curID := client.textID if curID == nil { return } @@ -452,9 +542,7 @@ func insertAtUTF16Index(base string, index uint32, insert string) (string, error } func (s *Server) handleDelete(msg *lrcpb.Event_Delete, client *client) { - s.idmapsMu.Lock() - curID := s.clientToID[client] - s.idmapsMu.Unlock() + curID := client.textID if curID == nil { return } @@ -502,9 +590,7 @@ func (s *Server) broadcast(event *lrcpb.Event, client *client) { } func (s *Server) handleEditBatch(msg *lrcpb.Event_Editbatch, client *client) { - s.idmapsMu.Lock() - curID := s.clientToID[client] - s.idmapsMu.Unlock() + curID := client.textID if curID == nil { return } -- 2.51.2