package sse import ( "context" "fmt" "ruben/inventory2/domains/raw_events" ) type ( DBEventPublisher struct { queue sender getAccountID func(context.Context, raw_events.Event) (int64, error) getEventType func(acctID int64, e raw_events.Event) (string, error) } ) func (q *Queue) NewDBEventPublisher( getAccountID func(context.Context, raw_events.Event) (int64, error), getEventType func(acctID int64, e raw_events.Event) (string, error), ) *DBEventPublisher { return &DBEventPublisher{ queue: q, getAccountID: getAccountID, getEventType: getEventType, } } // Notify publishes e to the SSE queue. The publish itself is quick and // synchronous, so it acks inline once it succeeds - there's no separate // async completion to wait for here. func (p *DBEventPublisher) Notify(ctx context.Context, e raw_events.Event, ack func(context.Context) error) error { acctID, err := p.getAccountID(ctx, e) if err != nil { return fmt.Errorf("failed to get account id: %w", err) } et, err := p.getEventType(acctID, e) if err != nil { return fmt.Errorf("failed to compute event type: %w", err) } if err := p.queue.Send(ctx, StandardEvent(acctID, et)); err != nil { return fmt.Errorf("failed to send sse event to listener: %w", err) } if err := ack(ctx); err != nil { return fmt.Errorf("failed to ack event: %w", err) } return nil }