refactor: simplify chat subscription API
Room.Subscribe() now returns a channel of parsed Message structs instead of raw NATS messages. The room handles NATS subscription and message parsing internally, so callers no longer need to call Receive() separately.
This commit is contained in:
42
chat/chat.go
42
chat/chat.go
@@ -76,12 +76,11 @@ func (r *Room) Send(msg Message) {
|
||||
}
|
||||
}
|
||||
|
||||
// Receive processes an incoming NATS message, appending it to the buffer.
|
||||
// Returns the new message and a snapshot of all messages.
|
||||
func (r *Room) Receive(data []byte) (Message, []Message) {
|
||||
// receive processes an incoming NATS message, appending it to the buffer.
|
||||
func (r *Room) receive(data []byte) (Message, bool) {
|
||||
var msg Message
|
||||
if err := json.Unmarshal(data, &msg); err != nil {
|
||||
return msg, nil
|
||||
return msg, false
|
||||
}
|
||||
|
||||
r.mu.Lock()
|
||||
@@ -89,11 +88,9 @@ func (r *Room) Receive(data []byte) (Message, []Message) {
|
||||
if len(r.messages) > maxMessages {
|
||||
r.messages = r.messages[len(r.messages)-maxMessages:]
|
||||
}
|
||||
snapshot := make([]Message, len(r.messages))
|
||||
copy(snapshot, r.messages)
|
||||
r.mu.Unlock()
|
||||
|
||||
return msg, snapshot
|
||||
return msg, true
|
||||
}
|
||||
|
||||
// Messages returns a snapshot of the current message buffer.
|
||||
@@ -105,15 +102,32 @@ func (r *Room) Messages() []Message {
|
||||
return snapshot
|
||||
}
|
||||
|
||||
// Subscribe creates a NATS channel subscription for the room's subject.
|
||||
// Caller is responsible for unsubscribing.
|
||||
func (r *Room) Subscribe() (chan *nats.Msg, *nats.Subscription, error) {
|
||||
ch := make(chan *nats.Msg, 64)
|
||||
sub, err := r.nc.ChanSubscribe(r.subject, ch)
|
||||
// Subscribe returns a channel of parsed messages and a cleanup function.
|
||||
// The room handles NATS subscription internally and buffers messages.
|
||||
func (r *Room) Subscribe() (<-chan Message, func()) {
|
||||
natsCh := make(chan *nats.Msg, 64)
|
||||
msgCh := make(chan Message, 64)
|
||||
|
||||
sub, err := r.nc.ChanSubscribe(r.subject, natsCh)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
close(msgCh)
|
||||
return msgCh, func() {}
|
||||
}
|
||||
return ch, sub, nil
|
||||
|
||||
go func() {
|
||||
for natsMsg := range natsCh {
|
||||
if msg, ok := r.receive(natsMsg.Data); ok {
|
||||
msgCh <- msg
|
||||
}
|
||||
}
|
||||
close(msgCh)
|
||||
}()
|
||||
|
||||
cleanup := func() {
|
||||
_ = sub.Unsubscribe()
|
||||
}
|
||||
|
||||
return msgCh, cleanup
|
||||
}
|
||||
|
||||
func (r *Room) saveMessage(msg Message) {
|
||||
|
||||
Reference in New Issue
Block a user