package sse import ( "bytes" "context" "fmt" "net/http" "sync" ) type ( Queue struct { in chan Event out map[int]chan Event ctx context.Context cancel context.CancelFunc lock sync.Mutex prevID int } Event struct { Type string Data []byte } ) func NewQueue() *Queue { ctx, cancel := context.WithCancel(context.Background()) return &Queue{ in: make(chan Event), out: make(map[int]chan Event), ctx: ctx, cancel: cancel, } } func (q *Queue) Start(ctx context.Context) error { defer q.cancel() for { select { case <-ctx.Done(): // we're done piping events return nil case e := <-q.in: // share event will all subscribers q.lock.Lock() for _, out := range q.out { out <- e } q.lock.Unlock() } } } func (q *Queue) Listen(ctx context.Context, fn func(context.Context, *Event) error) error { // create new out pipe and append it to the queue out := make(chan Event, 1) q.lock.Lock() q.prevID += 1 id := q.prevID q.out[id] = out q.lock.Unlock() // delete the pipe when done listening defer func() { q.lock.Lock() delete(q.out, id) q.lock.Unlock() }() for { select { case <-q.ctx.Done(): // queue is shut down return nil case <-ctx.Done(): // done listening to events return nil case e := <-out: // pass event to the caller if err := fn(ctx, &e); err != nil { return err } } } } func (q *Queue) Send(ctx context.Context, eventType string, data []byte) error { select { case <-q.ctx.Done(): // the queue has closed return fmt.Errorf("event queue closed: %w", q.ctx.Err()) case <-ctx.Done(): // sender ran out of time return fmt.Errorf("provided context canceled: %w", ctx.Err()) // push the event onto the queue case q.in <- Event{ Type: eventType, Data: data, }: } return nil } func (e *Event) Write(w http.ResponseWriter) { fmt.Fprintf( w, "event: %s\ndata: %s\n\n", e.Type, bytes.ReplaceAll(e.Data, []byte("\n"), []byte(" ")), ) w.(http.Flusher).Flush() }