Skip to content

Commit

Permalink
feat: pooling broadcast
Browse files Browse the repository at this point in the history
  • Loading branch information
nick-bisonai committed Aug 28, 2024
1 parent e714bf2 commit 7ef45c8
Showing 1 changed file with 9 additions and 8 deletions.
17 changes: 9 additions & 8 deletions node/pkg/dal/hub/hub.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"bisonai.com/miko/node/pkg/dal/collector"
dalcommon "bisonai.com/miko/node/pkg/dal/common"
"bisonai.com/miko/node/pkg/dal/utils/stats"
"bisonai.com/miko/node/pkg/utils/pool"
"github.com/rs/zerolog/log"
"nhooyr.io/websocket"
"nhooyr.io/websocket/wsjson"
Expand All @@ -29,10 +30,12 @@ type Hub struct {
Unregister chan *websocket.Conn
broadcast map[string]chan *dalcommon.OutgoingSubmissionData
mu sync.RWMutex
pool *pool.Pool
}

const (
MAX_CONNECTIONS = 10
WORKER_COUNT = 20
CleanupInterval = time.Hour
)

Expand All @@ -53,10 +56,13 @@ func NewHub(configs map[string]Config) *Hub {
Register: make(chan *websocket.Conn),
Unregister: make(chan *websocket.Conn),
broadcast: make(map[string]chan *dalcommon.OutgoingSubmissionData),
pool: pool.NewPool(WORKER_COUNT),
}
}

func (h *Hub) Start(ctx context.Context, collector *collector.Collector) {
h.pool.Run(ctx)

go h.handleClientRegistration()

h.initializeBroadcastChannels(collector)
Expand Down Expand Up @@ -151,24 +157,19 @@ func (h *Hub) broadcastDataForSymbol(ctx context.Context, symbol string) {
}

func (h *Hub) castSubmissionData(ctx context.Context, data *dalcommon.OutgoingSubmissionData, symbol *string) {
var wg sync.WaitGroup

h.mu.RLock()
defer h.mu.RUnlock()

for client, subscriptions := range h.Clients {
if _, ok := subscriptions[*symbol]; ok {
wg.Add(1)
go func(entry *websocket.Conn) {
defer wg.Done()
err := wsjson.Write(ctx, entry, data)
h.pool.AddJob(func() {
err := wsjson.Write(ctx, client, data)
if err != nil {
log.Warn().Err(err).Msg("failed to write message to client")
}
}(client)
})
}
}
wg.Wait()
}

func (h *Hub) cleanupJob(ctx context.Context) {
Expand Down

0 comments on commit 7ef45c8

Please sign in to comment.