diff --git a/pkg/atproto/teleport_endstream.go b/pkg/atproto/teleport_endstream.go index 840430c5..858c6d51 100644 --- a/pkg/atproto/teleport_endstream.go +++ b/pkg/atproto/teleport_endstream.go @@ -2,6 +2,7 @@ package atproto import ( "context" + "errors" "fmt" "time" @@ -9,6 +10,7 @@ import ( "github.com/bluesky-social/indigo/util" "github.com/bluesky-social/indigo/xrpc" "github.com/streamplace/oatproxy/pkg/oatproxy" + "gorm.io/gorm" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/placestream" @@ -40,8 +42,10 @@ import ( // // Best-effort: any failure only logs and returns. It never blocks the teleport // arrival notification (the caller publishes that before invoking this). It is a -// no-op when there is no stored session, no origin livestream strongRef, the -// referenced record is gone, or the livestream is already ended. +// no-op when there is no stored session, no OAuth proxy to act through (the +// `streamplace sync` subcommand runs the synchronizer without one), no origin +// livestream strongRef, the referenced record is gone, or the livestream is +// already ended. func (atsync *ATProtoSynchronizer) endLivestreamForTeleport(ctx context.Context, repoDID string, livestreamRef *comatproto.RepoStrongRef) { // A teleport without an origin livestream strongRef (the field is // optional, and records from before it existed won't have one) cannot be @@ -57,6 +61,10 @@ func (atsync *ATProtoSynchronizer) endLivestreamForTeleport(ctx context.Context, // this node (e.g. multi-node setups). Without a stored session we have no // credentials to write the record update, so there is nothing to do. session, err := atsync.StatefulDB.GetSessionByDID(repoDID) + if errors.Is(err, gorm.ErrRecordNotFound) { + log.Debug(ctx, "no stored session for streamer, cannot end livestream for teleport") + return + } if err != nil { log.Error(ctx, "failed to get streamer session for teleport stream-end", "err", err) return @@ -66,6 +74,14 @@ func (atsync *ATProtoSynchronizer) endLivestreamForTeleport(ctx context.Context, return } + // Some synchronizers run without an OAuth proxy (`streamplace sync`), and + // this is a best-effort callback on a timer goroutine — a panic here takes + // down the whole process, so a missing proxy is a logged no-op. + if atsync.OATProxy == nil { + log.Error(ctx, "no OAuth proxy configured, cannot end livestream for teleport") + return + } + // Refresh the session if its tokens are stale, then build a client that // signs requests as the streamer. session, err = atsync.OATProxy.RefreshIfNeeded(session) diff --git a/pkg/atproto/teleport_endstream_test.go b/pkg/atproto/teleport_endstream_test.go new file mode 100644 index 00000000..60582338 --- /dev/null +++ b/pkg/atproto/teleport_endstream_test.go @@ -0,0 +1,42 @@ +package atproto + +import ( + "context" + "testing" + + "github.com/streamplace/oatproxy/pkg/oatproxy" + "github.com/stretchr/testify/require" + comatproto "stream.place/streamplace/pkg/comatproto" +) + +// TestEndLivestreamForTeleportWithoutOATProxy: ending a livestream for a +// teleport needs an OAuth proxy to act as the streamer, but not every +// synchronizer has one — `streamplace sync` runs without it, and the serve path +// once forgot to wire it. With a stored session present the code used to march +// straight into OATProxy.RefreshIfNeeded and nil-deref; since this runs on a +// time.AfterFunc goroutine, that panic took down the whole node in production. +// Best-effort means a missing proxy is a logged no-op, never a crash. +func TestEndLivestreamForTeleportWithoutOATProxy(t *testing.T) { + ctx := context.Background() + atsync, _, _ := offlineSynchronizer(t) + require.Nil(t, atsync.OATProxy) + + streamer := "did:plc:cccccccccccccccccccccccc" + ref := &comatproto.RepoStrongRef{ + Uri: "at://" + streamer + "/place.stream.livestream/3lteleport0001", + Cid: "bafyreihdwdcefgh4dqkjv67uzcmw7ojee6xedzdetojuzjevtenxquvyku", + } + + // No stored session: returns before ever touching the proxy. + atsync.endLivestreamForTeleport(ctx, streamer, ref) + + // Stored session: this is the production crash — the session lookup + // succeeds and the next step is the nil proxy. + require.NoError(t, atsync.StatefulDB.CreateOAuthSession("test-jkt", &oatproxy.OAuthSession{ + DID: streamer, + Handle: "streamer.test", + PDSUrl: "http://127.0.0.1:1", + DownstreamDPoPJKT: "test-jkt", + })) + atsync.endLivestreamForTeleport(ctx, streamer, ref) +} diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index d1ee2c2e..ad0cec23 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -268,15 +268,6 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu Noter: noter, Bus: b, } - // Sync every repo we know about, at boot and every --sweep-interval after: - // a repair pass for repos left half-indexed by a previous run, a head check - // that finds the ones that drifted while we were not listening, and history - // deepening, which on a fresh node runs for as long as the network is big. - // Nothing below depends on it, so it runs in the background off the serve - // context -- shutdown cancels it -- and the node is up and serving in the - // meantime. - go atsync.SweepForever(ctx) - mm, err := media.MakeMediaManager(ctx, cli, signer, mod, b, atsync, ldb) if err != nil { return err @@ -339,6 +330,20 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu Public: cli.PublicOAuth, }) state.OATProxy = op + // The synchronizer uses stored sessions to act as a streamer (e.g. ending + // the source livestream when a teleport fires); without this it would + // nil-deref the first time a teleport arrives for a logged-in streamer. + atsync.OATProxy = op + + // Sync every repo we know about, at boot and every --sweep-interval after: + // a repair pass for repos left half-indexed by a previous run, a head check + // that finds the ones that drifted while we were not listening, and history + // deepening, which on a fresh node runs for as long as the network is big. + // Nothing below depends on it, so it runs in the background off the serve + // context -- shutdown cancels it -- and the node is up and serving in the + // meantime. Started only now that atsync is fully wired (OATProxy above): + // the sweep can index records and schedule callbacks that use every field. + go atsync.SweepForever(ctx) var replicator replication.Replicator = nil if slices.Contains(cli.Replicators, config.ReplicatorIroh) {