Files
NixMsg/internal/broker/queue.go

46 lines
893 B
Go

package broker
import (
"context"
"sync"
"git.asio.asia/nixevol/NixMsg/internal/app/port"
)
type uplinkItem struct {
conn port.ConnInfo
payload []byte
}
// uplinkQueue 每端串行队列,长度 256,满了堵住 OnPublish(背压)。
type uplinkQueue struct {
b *Broker
endpointID string
ch chan uplinkItem
once sync.Once
}
func newUplinkQueue(b *Broker, endpointID string) *uplinkQueue {
q := &uplinkQueue{
b: b,
endpointID: endpointID,
ch: make(chan uplinkItem, uplinkQueueSize),
}
go q.loop()
return q
}
func (q *uplinkQueue) push(item uplinkItem) {
q.ch <- item // 满则阻塞读循环,形成背压
}
func (q *uplinkQueue) close() {
q.once.Do(func() { close(q.ch) })
}
func (q *uplinkQueue) loop() {
for item := range q.ch {
_ = q.b.uplink.HandleUplink(context.Background(), item.conn, item.payload)
}
}