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.