From b6e554047ffb806c77a7c7cc39be5760c485aece Mon Sep 17 00:00:00 2001 From: Mitchell Hashimoto Date: Sat, 2 May 2026 08:59:38 -0700 Subject: [PATCH] jetstream: extract reusable consumer into internal/jetstream The consumer code in jetstream.go mixed Tack-specific concerns (collections, applyCommit, store cursor) with generic firehose mechanics: configuring the upstream client, looping reconnects, rewinding time-based cursors, persisting cursor progress, and distinguishing permanent bad-record failures from transient handler failures. None of those are specific to Tack and they are the bits most likely to be reused for future jetstream consumers. Move the generic mechanics into a new internal/jetstream package. --- internal/jetstream/consumer.go | 195 +++++++++++++++++++++++ internal/jetstream/consumer_test.go | 22 +++ internal/jetstream/cursor.go | 61 +++++++ internal/jetstream/doc.go | 16 ++ internal/jetstream/processor.go | 146 +++++++++++++++++ internal/jetstream/processor_test.go | 149 +++++++++++++++++ jetstream.go | 229 ++++----------------------- jetstream_test.go | 22 +++ 8 files changed, 639 insertions(+), 201 deletions(-) create mode 100644 internal/jetstream/consumer.go create mode 100644 internal/jetstream/consumer_test.go create mode 100644 internal/jetstream/cursor.go create mode 100644 internal/jetstream/doc.go create mode 100644 internal/jetstream/processor.go create mode 100644 internal/jetstream/processor_test.go diff --git a/internal/jetstream/consumer.go b/internal/jetstream/consumer.go new file mode 100644 index 0000000..f65ed4c --- /dev/null +++ b/internal/jetstream/consumer.go @@ -0,0 +1,195 @@ +package jetstream + +import ( + "context" + "errors" + "fmt" + "log/slog" + "time" + + "github.com/bluesky-social/jetstream/pkg/client" + "github.com/bluesky-social/jetstream/pkg/client/schedulers/sequential" +) + +const ( + // DefaultRewind is the cursor safety buffer used on reconnect. Jetstream + // cursors are time-based, so reconnecting a few seconds before the saved + // cursor avoids exact-boundary gaps at the cost of harmless duplicate + // deliveries for idempotent handlers. + DefaultRewind = 5 * time.Second + + // DefaultReconnectDelay is the pause between failed websocket reads. + DefaultReconnectDelay = 2 * time.Second +) + +// Config configures a reusable jetstream consumer. Use Processor directly when +// you already have events from another source and do not need websocket setup. +type Config struct { + // WebsocketURL is the jetstream endpoint used by Consumer. + WebsocketURL string + + // Collections is the set of record collection NSIDs to request from + // jetstream and process locally. An empty slice means "all collections", + // matching the upstream client behavior. + Collections []string + + CursorStore CursorStore + Handler Handler + Logger *slog.Logger + + // SchedulerIdent names the upstream sequential scheduler. It defaults to + // "jetstream". + SchedulerIdent string + + // Rewind controls how far before a saved cursor the Consumer reconnects. + // Zero uses DefaultRewind. + Rewind time.Duration + + // ReconnectDelay controls how long the Consumer waits before reconnecting + // after a failed read loop. Zero uses DefaultReconnectDelay. + ReconnectDelay time.Duration +} + +// Consumer owns the upstream jetstream client and reconnect loop. +type Consumer struct { + client *client.Client + processor *Processor + cursorStore CursorStore + logger *slog.Logger + rewind time.Duration + reconnectDelay time.Duration +} + +// NewConsumer builds a Consumer. Call Run to enter the blocking reconnect loop, +// or use Start to run it in a background goroutine. +func NewConsumer(cfg Config) (*Consumer, error) { + if cfg.WebsocketURL == "" { + return nil, errors.New("websocket URL is required") + } + + processor := &Processor{ + Collections: cfg.Collections, + CursorStore: cfg.CursorStore, + Handler: cfg.Handler, + Logger: cfg.Logger, + } + if err := processor.validate(); err != nil { + return nil, err + } + + cfg = withDefaults(cfg) + clientCfg := client.DefaultClientConfig() + clientCfg.WebsocketURL = cfg.WebsocketURL + clientCfg.WantedCollections = append([]string(nil), cfg.Collections...) + + c, err := client.NewClient( + clientCfg, + cfg.Logger, + sequential.NewScheduler(cfg.SchedulerIdent, cfg.Logger, processor.HandleEvent), + ) + if err != nil { + return nil, fmt.Errorf("new jetstream client: %w", err) + } + + return &Consumer{ + client: c, + processor: processor, + cursorStore: cfg.CursorStore, + logger: cfg.Logger, + rewind: cfg.Rewind, + reconnectDelay: cfg.ReconnectDelay, + }, nil +} + +// Start creates a Consumer and runs it in a background goroutine for the +// lifetime of ctx. +func Start(ctx context.Context, cfg Config) (*Consumer, error) { + c, err := NewConsumer(cfg) + if err != nil { + return nil, err + } + go c.Run(ctx) + return c, nil +} + +// Run consumes jetstream until ctx is cancelled. Connection failures are logged +// and retried after ReconnectDelay because the websocket is expected to be +// long-lived but not permanent. +func (c *Consumer) Run(ctx context.Context) { + for { + if ctx.Err() != nil { + return + } + + cur, err := c.cursorStore.LoadCursor(ctx) + if err != nil { + c.logger.Warn("ignoring unreadable cursor; resuming from now", "err", err) + cur = nil + } + + cursorForLog := cur + cursorForRead := c.rewindCursor(cur) + if cursorForRead != nil { + c.logger.Info("connecting to jetstream", + "cursor_us", *cursorForLog, + "rewound_us", *cursorForRead, + "rewind", c.rewind, + ) + } else { + c.logger.Info("connecting to jetstream from now (no cursor)") + } + + if err := c.client.ConnectAndRead(ctx, cursorForRead); err != nil { + if ctx.Err() != nil { + return + } + c.logger.Error("jetstream read loop", "err", err) + select { + case <-ctx.Done(): + return + case <-time.After(c.reconnectDelay): + } + continue + } + if ctx.Err() != nil { + return + } + } +} + +// rewindCursor returns the cursor value to hand to jetstream on reconnect. +// Jetstream cursors are microsecond timestamps and exact-boundary replay is not +// guaranteed, so we intentionally resume slightly before the last saved cursor. +// The duplicate window is expected to be safe for idempotent handlers; clamping +// at zero keeps brand-new or very-early cursors valid. +func (c *Consumer) rewindCursor(cur *int64) *int64 { + if cur == nil { + return nil + } + rewound := *cur - int64(c.rewind/time.Microsecond) + if rewound < 0 { + rewound = 0 + } + return &rewound +} + +func withDefaults(cfg Config) Config { + cfg.Logger = loggerOrDefault(cfg.Logger) + if cfg.SchedulerIdent == "" { + cfg.SchedulerIdent = "jetstream" + } + if cfg.Rewind == 0 { + cfg.Rewind = DefaultRewind + } + if cfg.ReconnectDelay == 0 { + cfg.ReconnectDelay = DefaultReconnectDelay + } + return cfg +} + +func loggerOrDefault(logger *slog.Logger) *slog.Logger { + if logger != nil { + return logger + } + return slog.Default() +} diff --git a/internal/jetstream/consumer_test.go b/internal/jetstream/consumer_test.go new file mode 100644 index 0000000..c095d31 --- /dev/null +++ b/internal/jetstream/consumer_test.go @@ -0,0 +1,22 @@ +package jetstream + +import ( + "testing" + "time" +) + +func TestConsumerRewindCursor(t *testing.T) { + c := &Consumer{rewind: 5 * time.Second} + + cursor := int64(10_000_000) + rewound := c.rewindCursor(&cursor) + if rewound == nil || *rewound != 5_000_000 { + t.Fatalf("rewound = %v, want 5000000", rewound) + } + + cursor = 2_000_000 + rewound = c.rewindCursor(&cursor) + if rewound == nil || *rewound != 0 { + t.Fatalf("rewound = %v, want 0", rewound) + } +} diff --git a/internal/jetstream/cursor.go b/internal/jetstream/cursor.go new file mode 100644 index 0000000..65b5ac8 --- /dev/null +++ b/internal/jetstream/cursor.go @@ -0,0 +1,61 @@ +package jetstream + +import ( + "context" + "sync" +) + +// CursorStore persists the jetstream cursor. Implementations commonly back this +// with a database row so process restarts can resume from the last applied +// event. +type CursorStore interface { + LoadCursor(context.Context) (*int64, error) + SaveCursor(context.Context, int64) error +} + +// MemoryCursorStore is an in-memory CursorStore implementation. It is useful +// for tests and for consumers that intentionally do not need cursor persistence +// across process restarts. +// +// The zero value is ready to use and starts with no cursor, which makes the +// consumer start from "now" until SaveCursor is called. +type MemoryCursorStore struct { + mu sync.Mutex + cursor *int64 +} + +var _ CursorStore = (*MemoryCursorStore)(nil) + +// NewMemoryCursorStore returns a MemoryCursorStore initialized with cursor. A +// nil cursor means no cursor has been saved yet. +func NewMemoryCursorStore(cursor *int64) *MemoryCursorStore { + s := &MemoryCursorStore{} + if cursor != nil { + v := *cursor + s.cursor = &v + } + return s +} + +// LoadCursor returns the last cursor saved in memory, or nil if none has been +// saved. The returned pointer is a copy so callers cannot mutate store state +// without going through SaveCursor. +func (s *MemoryCursorStore) LoadCursor(context.Context) (*int64, error) { + s.mu.Lock() + defer s.mu.Unlock() + + if s.cursor == nil { + return nil, nil + } + cursor := *s.cursor + return &cursor, nil +} + +// SaveCursor stores cursor in memory. +func (s *MemoryCursorStore) SaveCursor(_ context.Context, cursor int64) error { + s.mu.Lock() + defer s.mu.Unlock() + + s.cursor = &cursor + return nil +} diff --git a/internal/jetstream/doc.go b/internal/jetstream/doc.go new file mode 100644 index 0000000..cc1005c --- /dev/null +++ b/internal/jetstream/doc.go @@ -0,0 +1,16 @@ +// Package jetstream provides a reusable Bluesky Jetstream consumer. +// +// The package owns the mechanics that are common to firehose consumers: +// configuring the upstream websocket client, filtering by record collection, +// reconnecting after dropped reads, rewinding time-based cursors on reconnect, +// persisting cursor progress, and distinguishing permanent bad-record failures +// from transient handler failures. +// +// Callers provide a CursorStore and Handler, either through Config for a full +// Consumer or directly on Processor when they already have events from another +// source. The handler receives commit events for the configured collections and +// should apply domain-specific mutations. If a record is permanently unusable, +// the handler should return BadRecord(err) so the processor can advance the +// cursor and avoid replaying the same broken event forever. Any other error +// leaves the cursor unchanged so a later delivery can retry the event. +package jetstream diff --git a/internal/jetstream/processor.go b/internal/jetstream/processor.go new file mode 100644 index 0000000..546fff8 --- /dev/null +++ b/internal/jetstream/processor.go @@ -0,0 +1,146 @@ +package jetstream + +import ( + "context" + "errors" + "fmt" + "log/slog" + + jsmodels "github.com/bluesky-social/jetstream/pkg/models" +) + +// Handler applies one commit event. Returning BadRecord(err) tells the +// Processor that the event is permanently unusable and that the cursor should +// still advance; any other error is treated as transient and leaves the cursor +// unchanged for a later retry. +type Handler interface { + HandleJetstreamEvent(context.Context, *jsmodels.Event) error +} + +// HandlerFunc adapts a function to Handler. +type HandlerFunc func(context.Context, *jsmodels.Event) error + +// HandleJetstreamEvent calls f(ctx, event). +func (f HandlerFunc) HandleJetstreamEvent(ctx context.Context, event *jsmodels.Event) error { + return f(ctx, event) +} + +// Processor handles per-event policy that is independent of websocket IO. +type Processor struct { + // Collections is the set of record collection NSIDs to process locally. An + // empty slice means "all collections". + Collections []string + + // CursorStore persists progress after events are handled. + CursorStore CursorStore + + // Handler applies commit events that pass the collection filter. + Handler Handler + + // Logger receives apply errors and ignored-collection diagnostics. A nil + // Logger uses slog.Default(). + Logger *slog.Logger +} + +// HandleEvent applies a single jetstream event and advances the cursor when it +// is safe to do so. +func (p *Processor) HandleEvent(ctx context.Context, event *jsmodels.Event) error { + if event == nil || event.Kind != jsmodels.EventKindCommit || event.Commit == nil { + return nil + } + + if err := p.validate(); err != nil { + return err + } + + logger := loggerOrDefault(p.Logger) + wanted, err := p.wantsCollection(event.Commit.Collection) + if err != nil { + return err + } + if !wanted { + logger.Debug("ignoring unexpected collection", + "collection", event.Commit.Collection) + return p.saveCursor(ctx, event.TimeUS) + } + + applyErr := p.Handler.HandleJetstreamEvent(ctx, event) + if applyErr != nil { + logger.Error("apply commit", + "err", applyErr, + "did", event.Did, + "collection", event.Commit.Collection, + "op", event.Commit.Operation, + "rkey", event.Commit.RKey, + "transient", !IsBadRecord(applyErr), + ) + + if !IsBadRecord(applyErr) { + return applyErr + } + } + + return p.saveCursor(ctx, event.TimeUS) +} + +func (p *Processor) saveCursor(ctx context.Context, cursor int64) error { + if err := p.CursorStore.SaveCursor(ctx, cursor); err != nil { + return fmt.Errorf("save cursor: %w", err) + } + return nil +} + +func (p *Processor) validate() error { + if p.CursorStore == nil { + return errors.New("cursor store is required") + } + if p.Handler == nil { + return errors.New("handler is required") + } + for _, collection := range p.Collections { + if collection == "" { + return errors.New("collection must not be empty") + } + } + return nil +} + +func (p *Processor) wantsCollection(collection string) (bool, error) { + if len(p.Collections) == 0 { + return true, nil + } + for _, wanted := range p.Collections { + if wanted == "" { + return false, errors.New("collection must not be empty") + } + if wanted == collection { + return true, nil + } + } + return false, nil +} + +// badRecordError marks a handler failure as caused by the record itself being +// permanently unusable, e.g. malformed JSON or an unrecoverable schema +// violation. Processors advance the cursor past these errors so one bad event +// cannot stall every later event on restart. +type badRecordError struct{ err error } + +func (e *badRecordError) Error() string { return e.err.Error() } +func (e *badRecordError) Unwrap() error { return e.err } + +// BadRecord wraps err so Processor recognizes it as a permanent, +// cursor-advancing failure. Do not use this for storage, network, or other +// transient infrastructure failures. +func BadRecord(err error) error { + if err == nil { + return nil + } + return &badRecordError{err: err} +} + +// IsBadRecord reports whether err, or anything it wraps, came from BadRecord. +func IsBadRecord(err error) bool { + var b *badRecordError + return errors.As(err, &b) +} diff --git a/internal/jetstream/processor_test.go b/internal/jetstream/processor_test.go new file mode 100644 index 0000000..4d5cd9c --- /dev/null +++ b/internal/jetstream/processor_test.go @@ -0,0 +1,149 @@ +package jetstream + +import ( + "context" + "errors" + "io" + "log/slog" + "testing" + + jsmodels "github.com/bluesky-social/jetstream/pkg/models" +) + +func testLogger() *slog.Logger { + return slog.New(slog.NewTextHandler(io.Discard, nil)) +} + +func requireMemoryCursor(t *testing.T, store CursorStore, want *int64) { + t.Helper() + + got, err := store.LoadCursor(context.Background()) + if err != nil { + t.Fatalf("load cursor: %v", err) + } + if want == nil { + if got != nil { + t.Fatalf("cursor = %v, want nil", *got) + } + return + } + if got == nil || *got != *want { + t.Fatalf("cursor = %v, want %d", got, *want) + } +} + +func testCommit(timeUS int64, collection string) *jsmodels.Event { + return &jsmodels.Event{ + Did: "did:plc:test", + TimeUS: timeUS, + Kind: jsmodels.EventKindCommit, + Commit: &jsmodels.Commit{ + Operation: "create", + Collection: collection, + RKey: "rk", + }, + } +} + +func TestProcessorCommitAdvancesCursor(t *testing.T) { + store := NewMemoryCursorStore(nil) + called := false + + processor := &Processor{ + Collections: []string{"example.collection"}, + CursorStore: store, + Handler: HandlerFunc(func(ctx context.Context, event *jsmodels.Event) error { + called = true + return nil + }), + Logger: testLogger(), + } + + if err := processor.HandleEvent(context.Background(), testCommit(123, "example.collection")); err != nil { + t.Fatalf("handle: %v", err) + } + if !called { + t.Fatalf("handler was not called") + } + want := int64(123) + requireMemoryCursor(t, store, &want) +} + +func TestProcessorIgnoresNonCommit(t *testing.T) { + store := NewMemoryCursorStore(nil) + processor := &Processor{ + CursorStore: store, + Handler: HandlerFunc(func(ctx context.Context, event *jsmodels.Event) error { + t.Fatalf("handler should not be called") + return nil + }), + Logger: testLogger(), + } + + event := &jsmodels.Event{ + Did: "did:plc:test", + TimeUS: 123, + Kind: jsmodels.EventKindAccount, + } + if err := processor.HandleEvent(context.Background(), event); err != nil { + t.Fatalf("handle: %v", err) + } + requireMemoryCursor(t, store, nil) +} + +func TestProcessorUnexpectedCollectionAdvancesCursor(t *testing.T) { + store := NewMemoryCursorStore(nil) + processor := &Processor{ + Collections: []string{"wanted.collection"}, + CursorStore: store, + Handler: HandlerFunc(func(ctx context.Context, event *jsmodels.Event) error { + t.Fatalf("handler should not be called") + return nil + }), + Logger: testLogger(), + } + + if err := processor.HandleEvent(context.Background(), testCommit(456, "other.collection")); err != nil { + t.Fatalf("handle: %v", err) + } + want := int64(456) + requireMemoryCursor(t, store, &want) +} + +func TestProcessorBadRecordAdvancesCursor(t *testing.T) { + store := NewMemoryCursorStore(nil) + processor := &Processor{ + Collections: []string{"example.collection"}, + CursorStore: store, + Handler: HandlerFunc(func(ctx context.Context, event *jsmodels.Event) error { + return BadRecord(errors.New("decode failed")) + }), + Logger: testLogger(), + } + + if err := processor.HandleEvent(context.Background(), testCommit(789, "example.collection")); err != nil { + t.Fatalf("handle: %v", err) + } + want := int64(789) + requireMemoryCursor(t, store, &want) +} + +func TestProcessorTransientErrorDoesNotAdvanceCursor(t *testing.T) { + cursor := int64(100) + store := NewMemoryCursorStore(&cursor) + transientErr := errors.New("database busy") + processor := &Processor{ + Collections: []string{"example.collection"}, + CursorStore: store, + Handler: HandlerFunc(func(ctx context.Context, event *jsmodels.Event) error { + return transientErr + }), + Logger: testLogger(), + } + + err := processor.HandleEvent(context.Background(), testCommit(200, "example.collection")) + if !errors.Is(err, transientErr) { + t.Fatalf("handle error = %v, want %v", err, transientErr) + } + requireMemoryCursor(t, store, &cursor) +} diff --git a/jetstream.go b/jetstream.go index 54f996e..73411f8 100644 --- a/jetstream.go +++ b/jetstream.go @@ -20,32 +20,31 @@ package main import ( "context" "encoding/json" - "errors" "fmt" - "time" - "github.com/bluesky-social/jetstream/pkg/client" - "github.com/bluesky-social/jetstream/pkg/client/schedulers/sequential" jsmodels "github.com/bluesky-social/jetstream/pkg/models" + js "go.mitchellh.com/tack/internal/jetstream" "tangled.org/core/api/tangled" ) // jetstream operation strings. The jetstream protocol publishes these as // the Commit.Operation field; pulling them out as constants keeps the -// switch in handleJetstreamEvent honest about typos. +// switch in applyCommit honest about typos. const ( jsOpCreate = "create" jsOpUpdate = "update" jsOpDelete = "delete" ) -// jetstreamRewind is how far before the persisted cursor we resume on -// reconnect. Jetstream's docs recommend rewinding a few seconds because -// the cursor is a time-based filter and exact-boundary replay is not -// guaranteed gapless across disconnects. The handler's mutations are -// idempotent (UPSERTs and DELETEs keyed on (did, rkey)) so the small -// amount of duplicate replay this introduces is harmless. -const jetstreamRewind = 5 * time.Second +var _ js.CursorStore = (*store)(nil) + +// jetstreamCollections is the server-side and local filter for the Tangled +// records tack mirrors out of jetstream. +var jetstreamCollections = []string{ + tangled.SpindleMemberNSID, + tangled.RepoNSID, + tangled.RepoCollaboratorNSID, +} // startJetstream dials the configured jetstream endpoint and spawns a // background goroutine that consumes events for the lifetime of ctx. It @@ -63,201 +62,29 @@ const jetstreamRewind = 5 * time.Second func startJetstream(ctx context.Context, cfg config, st *store, knots KnotConsumer) error { logger := loggerFrom(ctx).With("component", "jetstream") - // `wantedCollections` is a server-side filter: jetstream will only send - // us commit events whose record collection (NSID) is in this list. The - // NSIDs come from tangled-core's generated lexicon types so they stay - // in sync with whatever the appview/knots are publishing. - collections := []string{ - tangled.SpindleMemberNSID, - tangled.RepoNSID, - tangled.RepoCollaboratorNSID, - } - - // Configure our JetStream client. - clientCfg := client.DefaultClientConfig() - clientCfg.WebsocketURL = cfg.JetstreamURL - clientCfg.WantedCollections = collections - // The handler closes over `st`, `knots`, the spindle hostname and // the owner DID so the scheduler signature stays plain // `func(ctx, *Event) error` and applyCommit can hand the knot // consumer new sources as soon as matching repo records arrive // (gated on the publisher being an authorized actor). - handler := func(ctx context.Context, evt *jsmodels.Event) error { - return handleJetstreamEvent(ctx, st, knots, cfg.Hostname, cfg.OwnerDID, evt) - } + handler := js.HandlerFunc(func(ctx context.Context, evt *jsmodels.Event) error { + return applyCommit(ctx, st, knots, cfg.Hostname, cfg.OwnerDID, evt) + }) // Re-attach the component-scoped logger so handler — which the - // scheduler invokes with the ctx we pass to ConnectAndRead — can - // pull it back out via loggerFrom. + // consumer invokes with the ctx we pass to ConnectAndRead — can pull + // it back out via loggerFrom. ctx = loggerInto(ctx, logger) - // The sequential scheduler processes events one-at-a-time in arrival - // order. That's the right default for a spindle: ordering matters - // (e.g. a member-added event must apply before any record from that - // member is processed), and our event volume is tiny. - c, err := client.NewClient( - clientCfg, - logger, - sequential.NewScheduler("tack", logger, handler), - ) - if err != nil { - return fmt.Errorf("new jetstream client: %w", err) - } - - go func() { - for { - // We re-read the cursor from the store at every (re)connect so we - // pick up any progress the previous connection persisted before - // dying. nil means "start from now", which is the right default on - // a brand-new install or after a corrupt cursor read. - cur, err := st.LoadCursor(ctx) - if err != nil { - logger.Warn("ignoring unreadable cursor; resuming from now", "err", err) - cur = nil - } - // Rewind a few seconds before the saved cursor on reconnect. - // Jetstream cursors are time-based and the docs explicitly note - // that exact-boundary replay is not guaranteed gapless, so a - // small negative buffer protects against missing events that - // straddle the disconnect. Duplicates the rewind produces are - // safe: every applyCommit path is an idempotent upsert/delete - // keyed on (did, rkey), and SaveCursor only moves forward in - // practice because TimeUS is monotonic. - if cur != nil { - rewound := *cur - int64(jetstreamRewind/time.Microsecond) - if rewound < 0 { - rewound = 0 - } - logger.Info("connecting to jetstream", - "cursor_us", *cur, - "rewound_us", rewound, - "rewind", jetstreamRewind, - ) - cur = &rewound - } else { - logger.Info("connecting to jetstream from now (no cursor)") - } - - // Reconnect loop. ConnectAndRead blocks on the websocket and returns - // either when the connection drops (transient network error, server - // restart, etc.) or when ctx is cancelled. On error we sleep briefly - // and reconnect; on ctx cancellation we exit cleanly. - if err := c.ConnectAndRead(ctx, cur); err != nil { - if ctx.Err() != nil { - return - } - logger.Error("jetstream read loop", "err", err) - time.Sleep(2 * time.Second) - continue - } - if ctx.Err() != nil { - return - } - } - }() - - return nil -} - -// handleJetstreamEvent is the per-event callback for the JetStream. It -// applies the event to the store and, when appropriate, advances the -// persisted cursor. Any returned error is logged by the scheduler but -// does not tear down the connection. -// -// Cursor-advancement policy: -// -// - Apply succeeded: advance the cursor. -// - Apply failed with a badRecordError (malformed record / unusable -// input): advance anyway, since replaying a permanently-broken -// event on every reconnect would stall the firehose forever. -// - Apply failed with anything else (treated as transient: store -// hiccup, SQLite busy, disk full, shutdown race, ...): do NOT -// advance. Saving the cursor here would permanently skip the -// record, which for membership/repo state means we'd silently -// lose a row that the rest of tack relies on. Leaving the cursor -// in place gives the next reconnect (which rewinds by -// jetstreamRewind) a chance to re-deliver and re-apply it. -func handleJetstreamEvent(ctx context.Context, st *store, knots KnotConsumer, hostname, ownerDID string, evt *jsmodels.Event) error { - // We only care about commits, which are the actual record CRUD - // operations on a user's PDS. Account/identity events are ignored - // for now; if we ever care about handle changes we can add them. - if evt.Kind != jsmodels.EventKindCommit || evt.Commit == nil { - return nil - } - logger := loggerFrom(ctx) - - // Dispatch on collection. Unknown collections shouldn't happen given - // our wantedCollections filter, but be defensive — jetstream may - // send schema changes ahead of us updating the filter. - applyErr := applyCommit(ctx, st, knots, hostname, ownerDID, evt) - if applyErr != nil { - logger.Error("apply commit", - "err", applyErr, - "did", evt.Did, - "collection", evt.Commit.Collection, - "op", evt.Commit.Operation, - "rkey", evt.Commit.RKey, - "transient", !isBadRecord(applyErr), - ) - - // Transient failure: bail without advancing the cursor so the - // next delivery can retry this event. Returning the error - // surfaces it in the scheduler's logs too. - if !isBadRecord(applyErr) { - return applyErr - } - // Otherwise (badRecordError) fall through to cursor save: a - // single bad record shouldn't stall the cursor forever and - // force us to re-process every subsequent event after a - // restart. - } - - // Advance the cursor. TimeUS is the jetstream-assigned microsecond - // timestamp; saving it after-apply means a crash mid-batch will at - // worst replay the failing event, never skip past it. - if err := st.SaveCursor(ctx, evt.TimeUS); err != nil { - // Returning the error logs it; it doesn't kill the scheduler. - return fmt.Errorf("save cursor: %w", err) - } - - return nil -} - -// badRecordError marks an applyCommit failure as caused by the *record -// itself* being permanently unusable, e.g. a malformed JSON body, or a -// field that violates an invariant we can't recover from on retry. -// -// We need this distinction because the jetstream cursor is time-based -// and once advanced past an event we will never see it again. So: -// -// - Permanent bad-input failures should still advance the cursor: -// replaying the same broken record on every restart accomplishes -// nothing and would stall progress on every later event behind it. -// - Transient infrastructure failures (SQLite busy, disk full, store -// closed mid-shutdown, etc.) must NOT advance the cursor: the -// record is fine and the next attempt, on reconnect or after the -// transient condition clears, should reapply it. Skipping past -// such a failure can permanently lose membership/repo state. -// -// Anything returned from applyCommit that isn't a badRecordError is -// treated as transient by handleJetstreamEvent. -type badRecordError struct{ err error } - -func (e *badRecordError) Error() string { return e.err.Error() } -func (e *badRecordError) Unwrap() error { return e.err } - -// badRecord wraps err so handleJetstreamEvent recognizes it as a -// permanent, cursor-advancing failure. Use this for any error caused by -// the contents of the record (decode errors, schema violations); never -// for store/IO errors. -func badRecord(err error) error { return &badRecordError{err: err} } - -// isBadRecord reports whether err (or anything it wraps) is a -// badRecordError. -func isBadRecord(err error) bool { - var b *badRecordError - return errors.As(err, &b) + _, err := js.Start(ctx, js.Config{ + WebsocketURL: cfg.JetstreamURL, + Collections: jetstreamCollections, + CursorStore: st, + Handler: handler, + Logger: logger, + SchedulerIdent: "tack", + }) + return err } // applyCommit routes a commit to the right store mutation based on its @@ -296,7 +123,7 @@ func applySpindleMember(ctx context.Context, st *store, knots KnotConsumer, host if err := json.Unmarshal(c.Record, &rec); err != nil { // Decode failures are a permanent property of the record's // bytes; mark as bad so the cursor can advance past it. - return badRecord(fmt.Errorf("decode spindle.member: %w", err)) + return js.BadRecord(fmt.Errorf("decode spindle.member: %w", err)) } if err := st.UpsertSpindleMember(ctx, did, c.RKey, rec.Instance, rec.Subject, rec.CreatedAt); err != nil { return err @@ -340,7 +167,7 @@ func applyRepo(ctx context.Context, st *store, knots KnotConsumer, hostname, own var rec tangled.Repo if err := json.Unmarshal(c.Record, &rec); err != nil { // See applySpindleMember: decode errors are permanent. - return badRecord(fmt.Errorf("decode repo: %w", err)) + return js.BadRecord(fmt.Errorf("decode repo: %w", err)) } // Capture the prior (knot, spindle) before the upsert so the @@ -507,7 +334,7 @@ func applyRepoCollaborator(ctx context.Context, st *store, did string, c *jsmode var rec tangled.RepoCollaborator if err := json.Unmarshal(c.Record, &rec); err != nil { // See applySpindleMember: decode errors are permanent. - return badRecord(fmt.Errorf("decode repo.collaborator: %w", err)) + return js.BadRecord(fmt.Errorf("decode repo.collaborator: %w", err)) } return st.UpsertRepoCollaborator(ctx, did, c.RKey, deref(rec.Repo), deref(rec.RepoDid), diff --git a/jetstream_test.go b/jetstream_test.go index 7f35a9b..5e8e187 100644 --- a/jetstream_test.go +++ b/jetstream_test.go @@ -14,6 +14,7 @@ import ( "testing" jsmodels "github.com/bluesky-social/jetstream/pkg/models" + js "go.mitchellh.com/tack/internal/jetstream" "tangled.org/core/api/tangled" ) @@ -55,6 +56,27 @@ func requireCursor(t *testing.T, s *store, want int64) { } } +// handleJetstreamEvent is the testable per-event entry point. The generic +// jetstream processor owns cursor advancement; this wrapper only supplies +// Tack's collection list and commit application callback. +func handleJetstreamEvent( + ctx context.Context, + st *store, + knots KnotConsumer, + hostname, ownerDID string, + evt *jsmodels.Event, +) error { + processor := &js.Processor{ + Collections: jetstreamCollections, + CursorStore: st, + Handler: js.HandlerFunc(func(ctx context.Context, evt *jsmodels.Event) error { + return applyCommit(ctx, st, knots, hostname, ownerDID, evt) + }), + Logger: loggerFrom(ctx), + } + return processor.HandleEvent(ctx, evt) +} + // TestHandleNonCommitEvent confirms account/identity events are ignored // without error and without advancing the cursor — they don't have a // TimeUS we want to commit to. -- 2.51.2