package appview import ( "bytes" "context" "fmt" "net/http" "net/url" "strings" "sync/atomic" "time" "tangled.org/sparrowtek.com/effem-AppView/appview/database" comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/events" "github.com/bluesky-social/indigo/events/schedulers/parallel" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/repomgr" "github.com/gorilla/websocket" ) const effemNSPrefix = "xyz.effem." func (srv *Server) RunFirehoseConsumer(ctx context.Context) error { cursor, err := srv.loadFirehoseCursor() if err != nil { srv.logger.Warn("no saved cursor, starting from live", "err", err) } u, err := url.Parse(srv.config.RelayHost) if err != nil { return fmt.Errorf("parsing relay host: %w", err) } u.Path = "xrpc/com.atproto.sync.subscribeRepos" if cursor > 0 { u.RawQuery = fmt.Sprintf("cursor=%d", cursor) } srv.logger.Info("connecting to firehose", "url", u.String(), "cursor", cursor) dialer := websocket.DefaultDialer con, _, err := dialer.DialContext(ctx, u.String(), http.Header{ "User-Agent": []string{"effem-appview/1.0"}, }) if err != nil { return fmt.Errorf("dialing firehose: %w", err) } defer con.Close() rsc := &events.RepoStreamCallbacks{ RepoCommit: func(evt *comatproto.SyncSubscribeRepos_Commit) error { atomic.StoreInt64(&srv.lastSeq, evt.Seq) return srv.handleCommit(ctx, evt) }, RepoIdentity: func(evt *comatproto.SyncSubscribeRepos_Identity) error { atomic.StoreInt64(&srv.lastSeq, evt.Seq) srv.logger.Debug("identity event", "did", evt.Did) return nil }, } scheduler := parallel.NewScheduler( srv.config.FirehoseParallel, 1000, srv.config.RelayHost, rsc.EventHandler, ) go srv.persistCursorLoop(ctx) srv.logger.Info("firehose consumer running") return events.HandleRepoStream(ctx, con, scheduler, srv.logger) } func (srv *Server) handleCommit(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Commit) error { if evt.TooBig { srv.logger.Warn("skipping tooBig commit", "repo", evt.Repo, "seq", evt.Seq) return nil } var rr *repo.Repo if len(evt.Blocks) > 0 { var err error rr, err = repo.ReadRepoFromCar(ctx, bytes.NewReader(evt.Blocks)) if err != nil { srv.logger.Warn("failed to read repo from car", "err", err, "did", evt.Repo) return nil } } for _, op := range evt.Ops { collection, rkey, err := syntax.ParseRepoPath(op.Path) if err != nil { srv.logger.Warn("invalid repo path", "path", op.Path, "err", err) continue } collectionName := collection.String() if !strings.HasPrefix(collectionName, effemNSPrefix) { continue } srv.logger.Info("effem record event", "action", op.Action, "collection", collectionName, "rkey", rkey.String(), "repo", evt.Repo) switch repomgr.EventKind(op.Action) { case repomgr.EvtKindCreateRecord, repomgr.EvtKindUpdateRecord: if rr == nil { srv.logger.Warn("missing CAR blocks for create/update", "path", op.Path) continue } recCID, recCBOR, err := rr.GetRecordBytes(ctx, op.Path) if err != nil { srv.logger.Warn("failed to get record", "err", err, "path", op.Path) continue } if op.Cid != nil && lexutil.LexLink(recCID) != *op.Cid { srv.logger.Warn("record CID mismatch", "path", op.Path, "carCID", recCID, "opCID", op.Cid) continue } if recCBOR == nil { srv.logger.Warn("nil record payload", "path", op.Path) continue } if err := srv.indexer.IndexRecord(ctx, evt.Repo, collectionName, rkey.String(), *recCBOR); err != nil { srv.logger.Warn("failed to index record", "err", err, "path", op.Path) } case repomgr.EvtKindDeleteRecord: if err := srv.indexer.DeleteRecord(ctx, evt.Repo, collectionName, rkey.String()); err != nil { srv.logger.Warn("failed to delete record", "err", err, "path", op.Path) } default: continue } } return nil } func (srv *Server) persistCursorLoop(ctx context.Context) { ticker := time.NewTicker(5 * time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): seq := atomic.LoadInt64(&srv.lastSeq) if seq > 0 { if err := srv.saveFirehoseCursor(seq); err != nil { srv.logger.Warn("failed to persist final cursor", "err", err, "seq", seq) } } return case <-ticker.C: seq := atomic.LoadInt64(&srv.lastSeq) if seq > 0 { if err := srv.saveFirehoseCursor(seq); err != nil { srv.logger.Warn("failed to persist cursor", "err", err, "seq", seq) } } } } } func (srv *Server) loadFirehoseCursor() (int64, error) { var cur database.FirehoseCursor err := srv.db.First(&cur).Error if err != nil { return 0, err } return cur.Seq, nil } func (srv *Server) saveFirehoseCursor(seq int64) error { var cur database.FirehoseCursor err := srv.db.First(&cur).Error if err != nil { cur = database.FirehoseCursor{Seq: seq} return srv.db.Create(&cur).Error } cur.Seq = seq return srv.db.Save(&cur).Error }