diff --git a/appview/firehose.go b/appview/firehose.go index 10e33bd..1160c4f 100644 --- a/appview/firehose.go +++ b/appview/firehose.go @@ -150,6 +150,47 @@ func (srv *Server) RunFirehoseConsumer(ctx context.Context) error { srv.logger.Debug("identity event", "did", evt.Did) return nil }, + // #sync replaces the deprecated tooBig commit path: the relay emits it + // to re-baseline a repo after data loss or a broken stream. We don't + // walk the CAR here because effem-namespaced records are low-volume + // enough that the next commit covering them will re-index; fuller + // reconciliation is future work if we see drift in practice. + RepoSync: func(evt *comatproto.SyncSubscribeRepos_Sync) error { + atomic.StoreInt64(&srv.lastSeq, evt.Seq) + srv.lastSeqTime.Store(time.Now()) + metrics.IncFirehoseEvent("sync", "") + srv.logger.Info("sync event", "did", evt.Did, "rev", evt.Rev, "seq", evt.Seq) + return nil + }, + // #account surfaces PDS/Relay status changes (takedowns, deactivations). + // We log + count for now; enforcement (hiding records from inactive + // accounts) will land when we add an account-status table. + RepoAccount: func(evt *comatproto.SyncSubscribeRepos_Account) error { + atomic.StoreInt64(&srv.lastSeq, evt.Seq) + srv.lastSeqTime.Store(time.Now()) + action := "inactive" + if evt.Active { + action = "active" + } + metrics.IncFirehoseEvent("account", action) + status := "" + if evt.Status != nil { + status = *evt.Status + } + srv.logger.Info("account event", "did", evt.Did, "active", evt.Active, "status", status, "seq", evt.Seq) + return nil + }, + // #info carries out-of-band notices like OutdatedCursor. No Seq field, + // so cursor/watchdog state is untouched. + RepoInfo: func(evt *comatproto.SyncSubscribeRepos_Info) error { + metrics.IncFirehoseEvent("info", evt.Name) + msg := "" + if evt.Message != nil { + msg = *evt.Message + } + srv.logger.Warn("firehose info event", "name", evt.Name, "message", msg) + return nil + }, } scheduler := newBoundedScheduler( diff --git a/appview/metrics/metrics.go b/appview/metrics/metrics.go index f526388..8387e5e 100644 --- a/appview/metrics/metrics.go +++ b/appview/metrics/metrics.go @@ -50,7 +50,7 @@ var ( prometheus.CounterOpts{ Namespace: namespace, Name: "firehose_events_total", - Help: "AT Protocol firehose events observed, partitioned by kind (commit/identity) and action (create/update/delete).", + Help: "AT Protocol firehose events observed, partitioned by kind (commit/identity/sync/account/info) and action (create/update/delete for commits, active/inactive for accounts, info name otherwise).", }, []string{"kind", "action"}, ) @@ -182,9 +182,11 @@ func HTTPMiddleware() echo.MiddlewareFunc { } } -// IncFirehoseEvent increments the event counter. Kind is "commit" or -// "identity"; action is "create"/"update"/"delete" for commits, or empty -// for identity events. +// IncFirehoseEvent increments the event counter. Kind is one of +// "commit"/"identity"/"sync"/"account"/"info". Action is +// "create"/"update"/"delete" for commits, "active"/"inactive" for +// account status changes, the info event Name for info events, and +// empty for identity/sync. func IncFirehoseEvent(kind, action string) { firehoseEventsTotal.WithLabelValues(kind, action).Inc() }