package service import ( "strings" "browser.local/platform/domain" ) const logEventSubscriberBuffer = 512 type LogEventSubscription struct { Events <-chan domain.LogStreamEvent Close func() } type logEventSubscriber struct { serverInstanceID string events chan domain.LogStreamEvent } func (svc *CoreService) SubscribeLogEvents(serverInstanceID string) (LogEventSubscription, error) { serverInstanceID = strings.TrimSpace(serverInstanceID) if serverInstanceID == "" { return LogEventSubscription{}, validationError("serverInstanceId is required") } if _, err := svc.store.ServerInstances().Get(serverInstanceID); err != nil { return LogEventSubscription{}, err } events := make(chan domain.LogStreamEvent, logEventSubscriberBuffer) svc.logEventMu.Lock() svc.logEventSubscriberSeq++ id := svc.logEventSubscriberSeq svc.logEventSubscribers[id] = logEventSubscriber{serverInstanceID: serverInstanceID, events: events} svc.logEventMu.Unlock() closeOnce := func() { svc.logEventMu.Lock() if subscriber, ok := svc.logEventSubscribers[id]; ok { delete(svc.logEventSubscribers, id) close(subscriber.events) } svc.logEventMu.Unlock() } return LogEventSubscription{Events: events, Close: closeOnce}, nil } func (svc *CoreService) SubscribeLogEventsForSession(sessionID string, serverInstanceID string) (LogEventSubscription, error) { instance, err := svc.GetServerInstanceForSession(sessionID, serverInstanceID) if err != nil { return LogEventSubscription{}, err } return svc.SubscribeLogEvents(instance.ID) } func (svc *CoreService) publishLogEvents(stream domain.LogStream, entries []domain.LogEntry) { if len(entries) == 0 { return } events := make([]domain.LogStreamEvent, len(entries)) for index, entry := range entries { events[index] = domain.CopyLogStreamEvent(domain.LogStreamEvent{ ServerInstanceID: stream.ServerInstanceID, Stream: stream, Entry: entry, LatestSeq: stream.LatestSeq, }) } svc.logEventMu.Lock() for id, subscriber := range svc.logEventSubscribers { if subscriber.serverInstanceID != stream.ServerInstanceID { continue } dropped := false for _, event := range events { select { case subscriber.events <- event: default: delete(svc.logEventSubscribers, id) close(subscriber.events) dropped = true } if dropped { break } } } svc.logEventMu.Unlock() }