diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go index 4e5dd95a..e2d1ee14 100644 --- a/pkg/statedb/queue_processor.go +++ b/pkg/statedb/queue_processor.go @@ -354,11 +354,29 @@ func (state *StatefulDB) processFinalizeLivestreamTask(ctx context.Context, task // truly abandoned the heartbeat stays frozen and we end it on the next // pass; if it's a heartbeat-lag artifact, the heartbeat catches up and the // rescheduled task hits the "active" early-return above. + // + // BUT the heartbeat-lag guard only applies to repos this node is actively + // ingesting — the heartbeat runs inside StreamSession.doUpdateLivestream, + // which requires the streamer to have a local OAuth session on this node. + // If no session exists, there is no heartbeat to wait for and no write + // access to end the record anyway: the record arrived via firehose sync + // from an account that never connected here. Rescheduling in that case just + // respawns the task every idle window forever (the rescheduled key embeds a + // fresh timestamp, so dedup never fires), flooding the logs and growing the + // task table without bound. So drop the task instead. latest, err := state.model.GetLatestLivestreamForRepo(livestream.RepoDID) if err != nil { return fmt.Errorf("failed to get latest livestream for repo: %w", err) } if latest != nil && latest.URI == livestream.URI { + // Check for a local session before rescheduling. GetSessionByDID + // returns gorm.ErrRecordNotFound (or nil session via callers that + // swallow it) when the repo has never logged in here. + session, err := state.GetSessionByDID(livestream.RepoDID) + if err != nil || session == nil { + log.Debug(ctx, "stale latest livestream has no local session; dropping finalize task (firehose-observed, no heartbeat to wait for)", "uri", livestream.URI, "lastSeenAt", lastSeenTime) + return state.CompleteTask(ctx, task.ID) + } rescheduledAt := time.Now().Add(time.Duration(*rec.IdleTimeoutSeconds) * time.Second).UTC() rescheduledKey := fmt.Sprintf("finalize-livestream::%s::%s", livestream.URI, rescheduledAt.Format(util.ISO8601)) _, err = state.EnqueueTask(ctx, TaskFinalizeLivestream, finalizeLivestreamTask, WithTaskKey(rescheduledKey), WithScheduledAt(rescheduledAt)) diff --git a/pkg/statedb/queue_processor_finalize_test.go b/pkg/statedb/queue_processor_finalize_test.go index 9edb1858..ec4426ad 100644 --- a/pkg/statedb/queue_processor_finalize_test.go +++ b/pkg/statedb/queue_processor_finalize_test.go @@ -6,7 +6,9 @@ import ( "testing" "time" + "github.com/streamplace/oatproxy/pkg/oatproxy" "github.com/stretchr/testify/require" + "gorm.io/gorm" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/streamplace" @@ -46,9 +48,23 @@ func seedLivestream(t *testing.T, mod model.Model, did, rkey string, createdAgo func ptr[T any](v T) *T { return &v } +// seedSession inserts a non-revoked OAuth session row for did, so the repo +// counts as "has a local session on this node." The finalize guard reschedules +// only when a local session exists (the heartbeat-lag scenario is real only +// for repos this node is actively ingesting). Without this, the stale-but-latest +// guard drops the task instead (firehose-observed, never connected here). +func seedSession(t *testing.T, state *StatefulDB, did string) { + t.Helper() + require.NoError(t, state.CreateOAuthSession("jkt-"+did, &oatproxy.OAuthSession{ + DID: did, + DownstreamDPoPJKT: "jkt-" + did, + })) +} + // TestFinalizeLivestreamReschedulesWhenLatestButStale proves the guard added to // processFinalizeLivestreamTask: when a livestream's lastSeenAt is older than -// its idleTimeoutSeconds but the record is still the streamer's latest, the +// its idleTimeoutSeconds but the record is still the streamer's latest, AND +// the repo has a local OAuth session (so heartbeat lag is possible), the // task must NOT set endedAt (which would take the active stream pre-live). It // must instead reschedule itself for one more idle window and return nil, so a // heartbeat that's lagging behind actual ingestion gets a chance to catch up. @@ -70,14 +86,19 @@ func TestFinalizeLivestreamReschedulesWhenLatestButStale(t *testing.T) { IdleTimeoutSeconds: ptr(int64(300)), }) + // A local session exists, so the heartbeat-lag guard applies and + // the task reschedules instead of ending or dropping. + seedSession(t, state, did) + task := &AppTask{ID: 1, Type: TaskFinalizeLivestream, Payload: mustMarshal(t, FinalizeLivestreamTask{ LivestreamURI: uri, })} - // Must not reach the PDS client (no session is configured), and must - // not error — it reschedules and returns nil. + // Must not reach the PDS client (no real PDS is configured, only + // the session row), and must not error — it reschedules and + // returns nil. err := state.processFinalizeLivestreamTask(ctx, task) - require.NoError(t, err, "stale-but-latest must reschedule, not error or end") + require.NoError(t, err, "stale-but-latest with a session must reschedule, not error or end") // A rescheduled task must have been enqueued, keyed to the same URI. tasks, err := state.ListTasks(ctx, TaskFilters{Type: TaskFinalizeLivestream, Limit: 10}) @@ -142,6 +163,103 @@ func TestFinalizeLivestreamEndsSupersededRecord(t *testing.T) { }) } +// TestFinalizeLivestreamDropsStaleLatestWithoutSession proves the fix for the +// dev-environment log flood: when a stale-but-latest livestream has NO local +// OAuth session (it arrived via firehose sync from an account that never +// connected to this node), the task must be dropped — not rescheduled. +// +// Without a local session there is no heartbeat to wait for (doUpdateLivestream +// only runs inside a local StreamSession) and no write access to end the record. +// Rescheduling would respawn the task every idle window forever — the +// rescheduled key embeds a fresh timestamp, so dedup never fires — flooding the +// logs and growing the task table without bound. +func TestFinalizeLivestreamDropsStaleLatestWithoutSession(t *testing.T) { + WithAllDatabases(t, func(state *StatefulDB) { + ctx := context.Background() + did := "did:plc:firehose-only" + + // The record is "latest" (only one for this repo) and its lastSeenAt + // is well past the 300s idle timeout. + uri := seedLivestream(t, state.model, did, "latest", 1*time.Hour, &streamplace.Livestream{ + LexiconTypeID: "place.stream.livestream", + CreatedAt: time.Now().Add(-1 * time.Hour).Format(time.RFC3339), + LastSeenAt: ptr(time.Now().Add(-10 * time.Minute).Format(time.RFC3339)), + IdleTimeoutSeconds: ptr(int64(300)), + }) + // Deliberately NO seedSession: this repo was only seen on the firehose. + + // Create the task row so CompleteTask has something to mark done. + createdTask, err := state.EnqueueTask(ctx, TaskFinalizeLivestream, FinalizeLivestreamTask{ + LivestreamURI: uri, + }, WithTaskKey("finalize-livestream::"+uri+"::initial")) + require.NoError(t, err) + + // Must not error, must not reschedule, must not reach the PDS client. + // The task is dropped (completed) because there's no session and thus + // no heartbeat-lag scenario to guard. + err = state.processFinalizeLivestreamTask(ctx, createdTask) + require.NoError(t, err, "stale-but-latest with no session must drop, not error or reschedule") + + // No rescheduled (pending) task should have been enqueued. The + // original task row is COMPLETED; filter to PENDING to assert no + // reschedule was spawned. + tasks, err := state.ListTasks(ctx, TaskFilters{Type: TaskFinalizeLivestream, Status: TaskStatusPending, Limit: 10}) + require.NoError(t, err) + require.Empty(t, tasks, "no reschedule task should be enqueued for a session-less stale-but-latest record") + + // The record must NOT have been ended (no session = no write access). + ls, err := state.model.GetLivestream(uri) + require.NoError(t, err) + view, err := ls.ToLivestreamView() + require.NoError(t, err) + rec, ok := view.Record.Val.(*streamplace.Livestream) + require.True(t, ok) + require.Nil(t, rec.EndedAt, "endedAt must not be set without a local session to write it") + }) +} + +// TestFinalizeLivestreamSupersededWithoutSessionDropsNotErrors is a regression +// guard: a superseded record with no session should still be droppable. In +// practice superseded records proceed past the latest-record guard to the +// session-lookup path, where GetSessionByDID returns gorm.ErrRecordNotFound. +// This test documents that the superseded path still errors at session lookup +// (it doesn't silently drop), confirming the no-session drop only applies to +// the stale-but-latest branch. +func TestFinalizeLivestreamSupersededWithoutSessionErrorsAtSessionLookup(t *testing.T) { + WithAllDatabases(t, func(state *StatefulDB) { + ctx := context.Background() + did := "did:plc:superseded-nosession" + + // Older record: stale, will be superseded. + _ = seedLivestream(t, state.model, did, "old", 2*time.Hour, &streamplace.Livestream{ + LexiconTypeID: "place.stream.livestream", + CreatedAt: time.Now().Add(-2 * time.Hour).Format(time.RFC3339), + LastSeenAt: ptr(time.Now().Add(-10 * time.Minute).Format(time.RFC3339)), + IdleTimeoutSeconds: ptr(int64(300)), + }) + // Newer record: makes "old" no longer latest. + _ = seedLivestream(t, state.model, did, "new", 1*time.Minute, &streamplace.Livestream{ + LexiconTypeID: "place.stream.livestream", + CreatedAt: time.Now().Add(-1 * time.Minute).Format(time.RFC3339), + LastSeenAt: ptr(time.Now().Format(time.RFC3339)), + IdleTimeoutSeconds: ptr(int64(300)), + }) + // No session for this repo. + + oldURI := "at://" + did + "/place.stream.livestream/old" + task := &AppTask{ID: 2, Type: TaskFinalizeLivestream, Payload: mustMarshal(t, FinalizeLivestreamTask{ + LivestreamURI: oldURI, + })} + + err := state.processFinalizeLivestreamTask(ctx, task) + // Superseded records proceed past the guard to session lookup, which + // fails with gorm.ErrRecordNotFound — the no-session drop only + // applies to the stale-but-latest branch. + require.Error(t, err, "superseded record without session should error at session lookup, not drop") + require.ErrorIs(t, err, gorm.ErrRecordNotFound) + }) +} + // smoke: the memory-mode StatefulDB used above needs config.DBURL only for the // non-draft test bootstrap path; reference it so an unused import can't bite // if these tests grow. -- 2.51.2 From 54d901b24e9cf75d933b9e8b57c982113d0e807f Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 15 Jul 2026 14:30:03 -0700 Subject: [PATCH 2/2] Update pkg/statedb/queue_processor.go Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> --- pkg/statedb/queue_processor.go | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/pkg/statedb/queue_processor.go b/pkg/statedb/queue_processor.go index e2d1ee14..fc688b3b 100644 --- a/pkg/statedb/queue_processor.go +++ b/pkg/statedb/queue_processor.go @@ -373,10 +373,13 @@ func (state *StatefulDB) processFinalizeLivestreamTask(ctx context.Context, task // returns gorm.ErrRecordNotFound (or nil session via callers that // swallow it) when the repo has never logged in here. session, err := state.GetSessionByDID(livestream.RepoDID) - if err != nil || session == nil { + if errors.Is(err, gorm.ErrRecordNotFound) || (err == nil && session == nil) { log.Debug(ctx, "stale latest livestream has no local session; dropping finalize task (firehose-observed, no heartbeat to wait for)", "uri", livestream.URI, "lastSeenAt", lastSeenTime) return state.CompleteTask(ctx, task.ID) } + if err != nil { + return fmt.Errorf("failed to get session for finalize-livestream guard: %w", err) + } rescheduledAt := time.Now().Add(time.Duration(*rec.IdleTimeoutSeconds) * time.Second).UTC() rescheduledKey := fmt.Sprintf("finalize-livestream::%s::%s", livestream.URI, rescheduledAt.Format(util.ISO8601)) _, err = state.EnqueueTask(ctx, TaskFinalizeLivestream, finalizeLivestreamTask, WithTaskKey(rescheduledKey), WithScheduledAt(rescheduledAt))