package api import ( "context" "database/sql" "encoding/json" "fmt" "time" "github.com/cosmic-clash/cosmic-clash/server/store" ) // RunProposalOutboxDispatcher delivers committed proposal changes to the // authenticated WebSocket subscribers. It only reads proposal_changed rows; // result and other outbox event types remain owned by their own consumers. // Delivery is at-least-once because the row is acknowledged only after every // participant publication succeeds. func RunProposalOutboxDispatcher(ctx context.Context, db *sql.DB, service *Service) { if db == nil || service == nil { return } ticker := time.NewTicker(100 * time.Millisecond) defer ticker.Stop() dispatcher := store.NewOutboxDispatcher(db, func(deliveryCtx context.Context, event store.OutboxEvent) error { return deliverProposalOutboxEvent(deliveryCtx, event, service) }) for { select { case <-ctx.Done(): return case <-ticker.C: events, err := store.ReadUnpublishedProposalOutbox(ctx, db, 100) if err != nil { continue } _ = dispatchOutboxEvents(ctx, dispatcher, events) } } } func dispatchOutboxEvents(ctx context.Context, dispatcher *store.OutboxDispatcher, events []store.OutboxEvent) error { if len(events) == 0 { return nil } // Use the same delivery-before-ack contract as the general dispatcher, // while keeping the already-filtered batch from being read a second time. for _, event := range events { if event.EventID == "" { return fmt.Errorf("outbox event has no ID") } if err := dispatcher.Deliver(ctx, event); err != nil { return err } if err := dispatcher.Ack(ctx, event.EventID, time.Now().UTC()); err != nil { return err } } return nil } func deliverProposalOutboxEvent(_ context.Context, event store.OutboxEvent, service *Service) error { var envelope struct { Event string `json:"event"` Revision uint64 `json:"revision"` ResourceID string `json:"resource_id"` OccurredAt time.Time `json:"occurred_at"` State string `json:"state"` PlayerIDs []string `json:"player_ids"` } if err := json.Unmarshal(event.Payload, &envelope); err != nil { return fmt.Errorf("decode proposal outbox event: %w", err) } if envelope.Event != "proposal_changed" || envelope.ResourceID == "" || len(envelope.PlayerIDs) == 0 { return fmt.Errorf("invalid proposal outbox event") } for _, playerID := range envelope.PlayerIDs { if err := service.PublishControlPlaneEvent(ControlPlaneEvent{ Event: envelope.Event, Revision: envelope.Revision, ResourceID: envelope.ResourceID, OccurredAt: envelope.OccurredAt, State: envelope.State, PlayerID: playerID, }); err != nil { return err } } return nil }