From cd9fb88c1d353846ae48a10d5c67e19d29ec663c Mon Sep 17 00:00:00 2001 From: dholms Date: Tue, 10 Mar 2026 17:02:23 -0500 Subject: [PATCH] swithc from limit to on/off --- cmd/tap/README.md | 8 ++++---- cmd/tap/firehose.go | 24 +++++++++------------- cmd/tap/firehose_test.go | 44 ++++++++++++---------------------------- cmd/tap/main.go | 15 ++++++++------ cmd/tap/models/models.go | 5 ++--- cmd/tap/tap.go | 2 +- 6 files changed, 39 insertions(+), 59 deletions(-) diff --git a/cmd/tap/README.md b/cmd/tap/README.md index 61d1de8b..c3909c25 100644 --- a/cmd/tap/README.md +++ b/cmd/tap/README.md @@ -43,13 +43,13 @@ When running Tap locally, restarts are frequent (e.g. sleep/wake, code changes). Recommended local dev config: ```bash -go run ./cmd/tap run --disable-acks=true --firehose-replay-limit=30s +go run ./cmd/tap run --no-replay --disable-acks ``` Tips: -- **Set `--firehose-replay-limit=0`**: Always connect to the firehose head on restart instead of replaying from a stale cursor. Repos that fell behind will resync when their next firehose event arrives. +- **Set `--no-replay`**: Always connect to the firehose head on restart instead of replaying from a stale cursor. Repos that fell behind will resync when their next firehose event arrives. This flag is incompatible with `--full-network` and is not recommended for production. - **Don't use `--full-network`**: Full network mode tracks every repo on the network and takes days to backfill. Instead, add specific DIDs with `/repos/add`. -- **Use `--disable-acks=true` until you setup webhooks or event acks**: Use a simple WebSocket client like `websocat` to inspect events. +- **Use `--disable-acks` until you setup webhooks or event acks**: Use a simple WebSocket client like `websocat` to inspect events. ## HTTP API @@ -80,7 +80,7 @@ Environment variables or CLI flags: - `TAP_RESYNC_PARALLELISM`: concurrent resync workers (default: `5`) - `TAP_OUTBOX_PARALLELISM`: concurrent outbox workers (default: `1`) - `TAP_CURSOR_SAVE_INTERVAL`: how often to persist upstream firehose cursor (default: `1s`) -- `TAP_FIREHOSE_REPLAY_LIMIT`: max age of saved firehose cursor before skipping to live head (default: `24h`, set to `0` to always start from live) +- `TAP_NO_REPLAY`: skip saved cursor and connect to firehose head on startup; incompatible with `TAP_FULL_NETWORK`, not recommended for production (default: `false`) - `TAP_REPO_FETCH_TIMEOUT`: timeout for fetching repo CARs from PDS (default: `300s`) - `TAP_IDENT_CACHE_SIZE`: size of in-process identity cache (default: `2000000`) - `TAP_OUTBOX_CAPACITY`: rough size of outbox before back pressure is applied (default: `100000`) diff --git a/cmd/tap/firehose.go b/cmd/tap/firehose.go index 542920d1..0465d30a 100644 --- a/cmd/tap/firehose.go +++ b/cmd/tap/firehose.go @@ -9,6 +9,7 @@ import ( "sync/atomic" "time" + "github.com/bluesky-social/go-util/pkg/bus/cursor" comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/atdata" "github.com/bluesky-social/indigo/atproto/repo" @@ -36,7 +37,7 @@ type FirehoseProcessor struct { collectionFilters []string parallelism int cursorSaveInterval time.Duration - replayLimit time.Duration + noReplay bool lastSeq atomic.Int64 } @@ -53,7 +54,7 @@ func NewFirehoseProcessor(logger *slog.Logger, db *gorm.DB, events *EventManager collectionFilters: config.CollectionFilters, parallelism: config.FirehoseParallelism, cursorSaveInterval: config.FirehoseCursorSaveInterval, - replayLimit: config.FirehoseReplayLimit, + noReplay: config.NoReplay, } } @@ -363,9 +364,8 @@ func (fp *FirehoseProcessor) saveCursor(ctx context.Context) error { } return fp.db.WithContext(ctx).Save(&models.FirehoseCursor{ - Url: fp.relayUrl, - Cursor: seq, - SavedAt: time.Now().Unix(), + Url: fp.relayUrl, + Cursor: seq, }).Error } @@ -385,6 +385,11 @@ func (fp *FirehoseProcessor) RunCursorSaver(ctx context.Context) { } func (fp *FirehoseProcessor) GetCursor(ctx context.Context) (int64, error) { + if fp.noReplay { + fp.logger.Info("firehose replay disabled, skipping to live", "skippedCursor", cursor.Cursor) + return 0, nil + } + var cursor models.FirehoseCursor if err := fp.db.WithContext(ctx).Where("url = ?", fp.relayUrl).First(&cursor).Error; err != nil { if err == gorm.ErrRecordNotFound { @@ -394,15 +399,6 @@ func (fp *FirehoseProcessor) GetCursor(ctx context.Context) (int64, error) { return 0, err } - if cursor.SavedAt > 0 && time.Since(time.Unix(cursor.SavedAt, 0)) > fp.replayLimit { - fp.logger.Warn("saved cursor is too old, skipping to live", - "cursorAge", time.Since(time.Unix(cursor.SavedAt, 0)), - "replayLimit", fp.replayLimit, - "skippedCursor", cursor.Cursor, - ) - return 0, nil - } - return cursor.Cursor, nil } diff --git a/cmd/tap/firehose_test.go b/cmd/tap/firehose_test.go index df680399..e7808a28 100644 --- a/cmd/tap/firehose_test.go +++ b/cmd/tap/firehose_test.go @@ -16,7 +16,6 @@ func newTestFirehoseProcessor(te *testEnv, fullNetwork bool) *FirehoseProcessor FullNetworkMode: fullNetwork, FirehoseParallelism: 1, FirehoseCursorSaveInterval: 5 * time.Second, - FirehoseReplayLimit: 24 * time.Hour, EventCacheSize: 1000, } return NewFirehoseProcessor(te.server.logger, te.db, te.events, te.repos, config) @@ -535,39 +534,36 @@ func TestCursorSaveAndLoad(t *testing.T) { } } -func TestGetCursor_ReplayLimit(t *testing.T) { +func TestGetCursor_NoReplay(t *testing.T) { relayUrl := "wss://relay.test.example" - // helper to insert a cursor row with a specific SavedAt - saveCursor := func(te *testEnv, seq int64, savedAt int64) { + saveCursor := func(te *testEnv, seq int64) { t.Helper() err := te.db.Save(&models.FirehoseCursor{ - Url: relayUrl, - Cursor: seq, - SavedAt: savedAt, + Url: relayUrl, + Cursor: seq, }).Error if err != nil { t.Fatalf("failed to insert cursor: %v", err) } } - // helper to create a processor with a specific replay limit - makeProcessor := func(te *testEnv, replayLimit time.Duration) *FirehoseProcessor { + makeProcessor := func(te *testEnv, disableReplay bool) *FirehoseProcessor { t.Helper() config := &TapConfig{ RelayUrl: relayUrl, FirehoseParallelism: 1, FirehoseCursorSaveInterval: 5 * time.Second, - FirehoseReplayLimit: replayLimit, + NoReplay: disableReplay, EventCacheSize: 1000, } return NewFirehoseProcessor(te.server.logger, te.db, te.events, te.repos, config) } - t.Run("fresh cursor is returned", func(t *testing.T) { + t.Run("replay enabled returns saved cursor", func(t *testing.T) { te := newTestEnv(t, testEnvOpts{}) - saveCursor(te, 100, time.Now().Unix()) - fp := makeProcessor(te, 24*time.Hour) + saveCursor(te, 100) + fp := makeProcessor(te, false) cursor, err := fp.GetCursor(te.ctx) if err != nil { @@ -578,31 +574,17 @@ func TestGetCursor_ReplayLimit(t *testing.T) { } }) - t.Run("stale cursor is skipped", func(t *testing.T) { + t.Run("replay disabled skips saved cursor", func(t *testing.T) { te := newTestEnv(t, testEnvOpts{}) - saveCursor(te, 100, time.Now().Add(-48*time.Hour).Unix()) - fp := makeProcessor(te, 24*time.Hour) + saveCursor(te, 100) + fp := makeProcessor(te, true) cursor, err := fp.GetCursor(te.ctx) if err != nil { t.Fatalf("unexpected error: %v", err) } if cursor != 0 { - t.Fatalf("expected cursor=0 for stale cursor, got %d", cursor) - } - }) - - t.Run("replay limit 0 always skips", func(t *testing.T) { - te := newTestEnv(t, testEnvOpts{}) - saveCursor(te, 100, time.Now().Unix()) - fp := makeProcessor(te, 0) - - cursor, err := fp.GetCursor(te.ctx) - if err != nil { - t.Fatalf("unexpected error: %v", err) - } - if cursor != 0 { - t.Fatalf("expected cursor=0 with replay limit 0, got %d", cursor) + t.Fatalf("expected cursor=0 with replay disabled, got %d", cursor) } }) } diff --git a/cmd/tap/main.go b/cmd/tap/main.go index a686830a..d6b179c1 100644 --- a/cmd/tap/main.go +++ b/cmd/tap/main.go @@ -100,11 +100,10 @@ func run(args []string) error { Value: 1 * time.Second, Sources: cli.EnvVars("TAP_CURSOR_SAVE_INTERVAL"), }, - &cli.DurationFlag{ - Name: "firehose-replay-limit", - Usage: "max age of saved cursor before skipping to live (0 = always skip to live)", - Value: 24 * time.Hour, - Sources: cli.EnvVars("TAP_FIREHOSE_REPLAY_LIMIT"), + &cli.BoolFlag{ + Name: "no-replay", + Usage: "skip saved cursor and connect to firehose head on startup (incompatible with --full-network, not recommended for production)", + Sources: cli.EnvVars("TAP_NO_REPLAY"), }, &cli.DurationFlag{ Name: "repo-fetch-timeout", @@ -201,6 +200,10 @@ func runTap(ctx context.Context, cmd *cli.Command) error { return fmt.Errorf("plc-url must start with http:// or https://") } + if cmd.Bool("no-replay") && cmd.Bool("full-network") { + return fmt.Errorf("--no-replay cannot be used with --full-network") + } + config := TapConfig{ DatabaseURL: cmd.String("db-url"), DBMaxConns: int(cmd.Int("max-db-conn")), @@ -210,7 +213,7 @@ func runTap(ctx context.Context, cmd *cli.Command) error { ResyncParallelism: int(cmd.Int("resync-parallelism")), OutboxParallelism: int(cmd.Int("outbox-parallelism")), FirehoseCursorSaveInterval: cmd.Duration("cursor-save-interval"), - FirehoseReplayLimit: cmd.Duration("firehose-replay-limit"), + NoReplay: cmd.Bool("no-replay"), RepoFetchTimeout: cmd.Duration("repo-fetch-timeout"), IdentityCacheSize: int(cmd.Int("ident-cache-size")), EventCacheSize: int(cmd.Int("outbox-capacity")), diff --git a/cmd/tap/models/models.go b/cmd/tap/models/models.go index 23383e3e..7a4ce5bb 100644 --- a/cmd/tap/models/models.go +++ b/cmd/tap/models/models.go @@ -56,9 +56,8 @@ type RepoRecord struct { } type FirehoseCursor struct { - Url string `gorm:"primaryKey"` - Cursor int64 `gorm:"not null"` - SavedAt int64 `gorm:"not null;default:0"` + Url string `gorm:"primaryKey"` + Cursor int64 `gorm:"not null"` } type ListReposCursor struct { diff --git a/cmd/tap/tap.go b/cmd/tap/tap.go index 7eecc3d2..67046772 100644 --- a/cmd/tap/tap.go +++ b/cmd/tap/tap.go @@ -40,7 +40,7 @@ type TapConfig struct { ResyncParallelism int OutboxParallelism int FirehoseCursorSaveInterval time.Duration - FirehoseReplayLimit time.Duration + NoReplay bool RepoFetchTimeout time.Duration IdentityCacheSize int EventCacheSize int -- 2.51.2