diff --git a/knotserver/events.go b/knotserver/events.go --- a/knotserver/events.go +++ b/knotserver/events.go @@ -2,146 +2,36 @@ import ( "context" - "encoding/json" + "errors" "net/http" - "strconv" "time" "github.com/bluesky-social/indigo/xrpc" - "github.com/gorilla/websocket" "tangled.org/core/api/tangled" + "tangled.org/core/eventstream" "tangled.org/core/log" ) - -var upgrader = websocket.Upgrader{ - ReadBufferSize: 1024, - WriteBufferSize: 1024, -} func (h *Knot) Events(w http.ResponseWriter, r *http.Request) { l := log.SubLogger(h.l, "eventstream") l.Debug("received new connection") - conn, err := upgrader.Upgrade(w, r, nil) - if err != nil { - l.Error("websocket upgrade failed", "err", err) - w.WriteHeader(http.StatusInternalServerError) - return + err := eventstream.Stream(w, r, eventstream.StreamConfig{ + Backend: h.db, + Notifier: h.n, + Logger: l, + }) + if err != nil && !errors.Is(err, eventstream.ErrDrainCap) { + l.Error("event stream ended with error", "err", err) } - defer conn.Close() - l.Debug("upgraded http to wss") - ch := h.n.Subscribe() - defer h.n.Unsubscribe(ch) - - ctx, cancel := context.WithCancel(r.Context()) - defer cancel() go func() { - for { - if _, _, err := conn.NextReader(); err != nil { - l.Error("failed to read", "err", err) - cancel() - return - } + retryCtx, retryCancel := context.WithTimeout(context.Background(), 10*time.Second) + defer retryCancel() + if err := h.requestCrawl(retryCtx); err != nil { + l.Error("error requesting crawls", "err", err) } }() - - var cursor int64 - cursorStr := r.URL.Query().Get("cursor") - if cursorStr != "" { - cursor, err = strconv.ParseInt(cursorStr, 10, 64) - if err != nil { - l.Error("invalid cursor, starting from beginning", "invalidCursor", cursorStr) - cursor = 0 - } - } - - l.Debug("going through backfill", "cursor", cursor) - if err := h.drainBackfill(conn, &cursor, 10_000); err != nil { - l.Error("failed to backfill", "err", err) - return - } - - // try request crawl when connection closed - defer func() { - go func() { - retryCtx, retryCancel := context.WithTimeout(context.Background(), 10*time.Second) - defer retryCancel() - if err := h.requestCrawl(retryCtx); err != nil { - l.Error("error requesting crawls", "err", err) - } - }() - }() - - for { - // wait for new data or timeout - select { - case <-ctx.Done(): - l.Debug("stopping stream: client closed connection") - return - case <-ch: - l.Debug("going through live data", "cursor", cursor) - if _, err := h.streamOps(conn, &cursor); err != nil { - l.Error("failed to stream", "err", err) - return - } - case <-time.After(30 * time.Second): - // send a keep-alive - if err = conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { - l.Error("failed to write control", "err", err) - } - } - } -} - -func (h *Knot) drainBackfill(conn *websocket.Conn, cursor *int64, maxBatches int) error { - for range maxBatches { - n, err := h.streamOps(conn, cursor) - if err != nil { - return err - } - if n < 100 { - return nil - } - } - h.l.Warn("backfill hit batch limit", "maxBatches", maxBatches, "cursor", *cursor) - return nil -} - -func (h *Knot) streamOps(conn *websocket.Conn, cursor *int64) (int, error) { - events, err := h.db.GetEvents(*cursor) - if err != nil { - h.l.Error("failed to fetch events from db", "err", err, "cursor", cursor) - return 0, err - } - - for _, event := range events { - var eventJson map[string]any - err := json.Unmarshal([]byte(event.EventJson), &eventJson) - if err != nil { - h.l.Error("failed to unmarshal event", "err", err) - return 0, err - } - - jsonMsg, err := json.Marshal(map[string]any{ - "rkey": event.Rkey, - "nsid": event.Nsid, - "event": eventJson, - "created": event.Created, - }) - if err != nil { - h.l.Error("failed to marshal record", "err", err) - return 0, err - } - - if err := conn.WriteMessage(websocket.TextMessage, jsonMsg); err != nil { - h.l.Debug("err", "err", err) - return 0, err - } - *cursor = event.Created - } - - return len(events), nil } func (h *Knot) requestCrawl(ctx context.Context) error { diff --git a/knotserver/ingester.go b/knotserver/ingester.go --- a/knotserver/ingester.go +++ b/knotserver/ingester.go @@ -19,6 +19,7 @@ jmodels "github.com/bluesky-social/jetstream/pkg/models" "tangled.org/core/api/tangled" "tangled.org/core/appview/models" + "tangled.org/core/eventstream" "tangled.org/core/knotserver/db" "tangled.org/core/knotserver/git" knotxrpc "tangled.org/core/knotserver/xrpc" @@ -502,10 +503,10 @@ return fmt.Errorf("failed to marshal pipeline event: %w", err) } - ev := db.Event{ + ev := eventstream.Event{ Rkey: tid.TID(), Nsid: tangled.PipelineNSID, - EventJson: string(eventJson), + EventJson: eventJson, } l.Info("inserting pipeline event") diff --git a/knotserver/internal.go b/knotserver/internal.go --- a/knotserver/internal.go +++ b/knotserver/internal.go @@ -18,6 +18,7 @@ "github.com/go-chi/chi/v5/middleware" "github.com/go-git/go-git/v5/plumbing" "tangled.org/core/api/tangled" + "tangled.org/core/eventstream" "tangled.org/core/hook" "tangled.org/core/idresolver" "tangled.org/core/knotserver/config" @@ -313,10 +314,10 @@ return err } - event := db.Event{ + event := eventstream.Event{ Rkey: tid.TID(), Nsid: tangled.GitRefUpdateNSID, - EventJson: string(eventJson), + EventJson: eventJson, } return h.db.InsertEvent(event, h.n) @@ -417,10 +418,10 @@ return nil } - event := db.Event{ + event := eventstream.Event{ Rkey: tid.TID(), Nsid: tangled.PipelineNSID, - EventJson: string(eventJson), + EventJson: eventJson, } if h.c.LogsAddr != "" { diff --git a/spindle/server.go b/spindle/server.go --- a/spindle/server.go +++ b/spindle/server.go @@ -15,6 +15,7 @@ "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" "tangled.org/core/eventconsumer/cursor" + "tangled.org/core/eventstream" "tangled.org/core/idresolver" "tangled.org/core/jetstream" "tangled.org/core/log" @@ -183,7 +184,7 @@ // job in the above registered queue. ccfg := eventconsumer.NewConsumerConfig() ccfg.Logger = log.SubLogger(logger, "eventconsumer") - ccfg.Dev = cfg.Server.Dev + ccfg.URLFunc = eventconsumer.DefaultURL(cfg.Server.Dev) ccfg.ProcessFunc = spindle.processPipeline ccfg.CursorStore = cursorStore knownKnots, err := d.Knots() @@ -374,7 +375,7 @@ return x.Router() } -func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, msg eventconsumer.Message) error { +func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) error { if msg.Nsid == tangled.PipelineNSID { tpl := tangled.Pipeline{} err := json.Unmarshal(msg.EventJson, &tpl) @@ -391,8 +392,8 @@ return fmt.Errorf("no repo data found") } - if src.Key() != tpl.TriggerMetadata.Repo.Knot { - return fmt.Errorf("repo knot does not match event source: %s != %s", src.Key(), tpl.TriggerMetadata.Repo.Knot) + if src.Host != tpl.TriggerMetadata.Repo.Knot { + return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, tpl.TriggerMetadata.Repo.Knot) } repoDid, err := s.resolvePipelineRepoDid(tpl.TriggerMetadata.Repo) @@ -401,7 +402,7 @@ } pipelineId := models.PipelineId{ - Knot: src.Key(), + Knot: src.Host, Rkey: msg.Rkey, } diff --git a/spindle/stream.go b/spindle/stream.go --- a/spindle/stream.go +++ b/spindle/stream.go @@ -3,13 +3,14 @@ import ( "context" "encoding/json" + "errors" "fmt" "io" "net/http" "os" - "strconv" "time" + "tangled.org/core/eventstream" "tangled.org/core/log" "tangled.org/core/spindle/models" @@ -25,69 +26,15 @@ func (s *Spindle) Events(w http.ResponseWriter, r *http.Request) { l := log.SubLogger(s.l, "eventstream") - l.Debug("received new connection") - conn, err := upgrader.Upgrade(w, r, nil) - if err != nil { - l.Error("websocket upgrade failed", "err", err) - w.WriteHeader(http.StatusInternalServerError) - return - } - defer conn.Close() - l.Debug("upgraded http to wss") - - ch := s.n.Subscribe() - defer s.n.Unsubscribe(ch) - - ctx, cancel := context.WithCancel(r.Context()) - defer cancel() - go func() { - for { - if _, _, err := conn.NextReader(); err != nil { - l.Error("failed to read", "err", err) - cancel() - return - } - } - }() - - var cursor int64 - cursorStr := r.URL.Query().Get("cursor") - if cursorStr != "" { - cursor, err = strconv.ParseInt(cursorStr, 10, 64) - if err != nil { - l.Error("invalid cursor, starting from beginning", "invalidCursor", cursorStr) - cursor = 0 - } - } - - // complete backfill first before going to live data - l.Debug("going through backfill", "cursor", cursor) - if err := s.streamPipelines(conn, &cursor); err != nil { - l.Error("failed to backfill", "err", err) - return - } - - for { - // wait for new data or timeout - select { - case <-ctx.Done(): - l.Debug("stopping stream: client closed connection") - return - case <-ch: - // we have been notified of new data - l.Debug("going through live data", "cursor", cursor) - if err := s.streamPipelines(conn, &cursor); err != nil { - l.Error("failed to stream", "err", err) - return - } - case <-time.After(30 * time.Second): - // send a keep-alive - if err = conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { - l.Error("failed to write control", "err", err) - } - } + err := eventstream.Stream(w, r, eventstream.StreamConfig{ + Backend: s.db, + Notifier: s.n, + Logger: l, + }) + if err != nil && !errors.Is(err, eventstream.ErrDrainCap) { + l.Error("event stream ended with error", "err", err) } } @@ -220,43 +167,6 @@ } } } -} - -func (s *Spindle) streamPipelines(conn *websocket.Conn, cursor *int64) error { - events, err := s.db.GetEvents(*cursor) - if err != nil { - s.l.Debug("err", "err", err) - return err - } - - for _, event := range events { - // first extract the inner json into a map - var eventJson map[string]any - err := json.Unmarshal([]byte(event.EventJson), &eventJson) - if err != nil { - s.l.Error("failed to unmarshal event", "err", err) - return err - } - - jsonMsg, err := json.Marshal(map[string]any{ - "rkey": event.Rkey, - "nsid": event.Nsid, - "event": eventJson, - "created": event.Created, - }) - if err != nil { - s.l.Error("failed to marshal record", "err", err) - return err - } - - if err := conn.WriteMessage(websocket.TextMessage, jsonMsg); err != nil { - s.l.Debug("err", "err", err) - return err - } - *cursor = event.Created - } - - return nil } func getWorkflowID(r *http.Request) (models.WorkflowId, error) { diff --git a/knotserver/db/events.go b/knotserver/db/events.go --- a/knotserver/db/events.go +++ b/knotserver/db/events.go @@ -3,32 +3,14 @@ import ( "encoding/json" "fmt" - "time" + "tangled.org/core/eventstream" "tangled.org/core/notifier" "tangled.org/core/tid" ) -type Event struct { - Rkey string `json:"rkey"` - Nsid string `json:"nsid"` - EventJson string `json:"event"` - Created int64 `json:"created"` -} - -func (d *DB) InsertEvent(event Event, notifier *notifier.Notifier) error { - - _, err := d.db.Exec( - `insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, - event.Rkey, - event.Nsid, - event.EventJson, - time.Now().UnixNano(), - ) - - notifier.NotifyAll() - - return err +func (d *DB) InsertEvent(event eventstream.Event, n *notifier.Notifier) error { + return eventstream.Insert(d.db, event, n) } func (d *DB) EmitDIDAssign(n *notifier.Notifier, ownerDid, repoName, repoDid string) error { @@ -43,47 +25,13 @@ return fmt.Errorf("marshal didAssign event: %w", err) } - return d.InsertEvent(Event{ + return d.InsertEvent(eventstream.Event{ Rkey: tid.TID(), Nsid: RepoDIDAssignNSID, - EventJson: string(eventJson), + EventJson: eventJson, }, n) } -func (d *DB) GetEvents(cursor int64) ([]Event, error) { - whereClause := "" - args := []any{} - if cursor > 0 { - whereClause = "where created > ?" - args = append(args, cursor) - } - - query := fmt.Sprintf(` - select rkey, nsid, event, created - from events - %s - order by created asc - limit 100 - `, whereClause) - - rows, err := d.db.Query(query, args...) - if err != nil { - return nil, err - } - defer rows.Close() - - var evts []Event - for rows.Next() { - var ev Event - if err := rows.Scan(&ev.Rkey, &ev.Nsid, &ev.EventJson, &ev.Created); err != nil { - return nil, err - } - evts = append(evts, ev) - } - - if err := rows.Err(); err != nil { - return nil, err - } - - return evts, nil +func (d *DB) GetEvents(cursor int64, limit int) ([]eventstream.Event, error) { + return eventstream.List(d.db, cursor, limit) } diff --git a/knotserver/xrpc/set_default_branch.go b/knotserver/xrpc/set_default_branch.go --- a/knotserver/xrpc/set_default_branch.go +++ b/knotserver/xrpc/set_default_branch.go @@ -9,7 +9,7 @@ "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" - "tangled.org/core/knotserver/db" + "tangled.org/core/eventstream" "tangled.org/core/knotserver/git" "tangled.org/core/rbac" "tangled.org/core/tid" @@ -101,14 +101,19 @@ } eventJson, err := json.Marshal(refUpdate) if err != nil { + fail(xrpcerr.GenericError(err)) return } - x.Db.InsertEvent(db.Event{ + if err := x.Db.InsertEvent(eventstream.Event{ Rkey: tid.TID(), Nsid: tangled.GitRefUpdateNSID, - EventJson: string(eventJson), - }, x.Notifier) + EventJson: eventJson, + }, x.Notifier); err != nil { + l.Error("failed to insert event", "error", err) + writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) + return + } w.WriteHeader(http.StatusOK) } diff --git a/spindle/db/events.go b/spindle/db/events.go --- a/spindle/db/events.go +++ b/spindle/db/events.go @@ -2,72 +2,21 @@ import ( "encoding/json" - "fmt" "time" "tangled.org/core/api/tangled" + "tangled.org/core/eventstream" "tangled.org/core/notifier" "tangled.org/core/spindle/models" "tangled.org/core/tid" ) -type Event struct { - Rkey string `json:"rkey"` - Nsid string `json:"nsid"` - Created int64 `json:"created"` - EventJson string `json:"event"` +func (d *DB) insertEvent(event eventstream.Event, n *notifier.Notifier) error { + return eventstream.Insert(d, event, n) } -func (d *DB) insertEvent(event Event, notifier *notifier.Notifier) error { - _, err := d.Exec( - `insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, - event.Rkey, - event.Nsid, - event.EventJson, - time.Now().UnixNano(), - ) - - notifier.NotifyAll() - - return err -} - -func (d *DB) GetEvents(cursor int64) ([]Event, error) { - whereClause := "" - args := []any{} - if cursor > 0 { - whereClause = "where created > ?" - args = append(args, cursor) - } - - query := fmt.Sprintf(` - select rkey, nsid, event, created - from events - %s - order by created asc - limit 100 - `, whereClause) - - rows, err := d.Query(query, args...) - if err != nil { - return nil, err - } - defer rows.Close() - - var evts []Event - for rows.Next() { - var ev Event - if err := rows.Scan(&ev.Rkey, &ev.Nsid, &ev.EventJson, &ev.Created); err != nil { - return nil, err - } - evts = append(evts, ev) - } - - if err := rows.Err(); err != nil { - return nil, err - } - - return evts, nil +func (d *DB) GetEvents(cursor int64, limit int) ([]eventstream.Event, error) { + return eventstream.List(d, cursor, limit) } func (d *DB) createStatusEvent( @@ -93,15 +42,13 @@ return err } - event := Event{ + event := eventstream.Event{ Rkey: tid.TID(), Nsid: tangled.PipelineStatusNSID, - Created: now.UnixNano(), - EventJson: string(eventJson), + EventJson: eventJson, } return d.insertEvent(event, n) - } func (d *DB) GetStatus(workflowId models.WorkflowId) (*tangled.PipelineStatus, error) {