diff --git a/apps/twisted/package.json b/apps/twisted/package.json index 60dde17..5d89f3c 100644 --- a/apps/twisted/package.json +++ b/apps/twisted/package.json @@ -15,6 +15,7 @@ "dependencies": { "@atcute/bluesky": "^3.3.0", "@atcute/client": "^4.2.1", + "@atcute/oauth-browser-client": "^3.0.0", "@atcute/tangled": "^1.0.17", "@capacitor/android": "^8.2.0", "@capacitor/app": "8.0.1", diff --git a/docs/api/specs/04-data-pipeline.md b/docs/api/specs/04-data-pipeline.md index edb5268..08bdec0 100644 --- a/docs/api/specs/04-data-pipeline.md +++ b/docs/api/specs/04-data-pipeline.md @@ -114,6 +114,9 @@ For each event, the indexer: ### Cursor Persistence Rules - If DB commit fails → cursor does not advance → event will be retried +- After successful DB writes, ack Tap first, then persist cursor for operator-visible resume +- If ack fails → cursor does not advance +- If ack succeeds but cursor persistence fails → retry cursor persistence until successful or process exit - If normalization fails → log error, optionally dead-letter, skip → cursor advances - If embedding scheduling fails → document remains keyword-searchable → cursor advances diff --git a/docs/api/specs/07-graph-backfill.md b/docs/api/specs/07-graph-backfill.md index f491e3b..0ad2285 100644 --- a/docs/api/specs/07-graph-backfill.md +++ b/docs/api/specs/07-graph-backfill.md @@ -59,7 +59,10 @@ Discovered DIDs are added to a queue, deduplicated by DID. Each entry tracks: For each discovered user: -1. **Check if already tracked**: Query Tap's `/info/:did` endpoint — if the repo is already tracked and backfilled, skip +1. **Check Tap status**: Query Tap's `/info/:did` endpoint and classify by status: + - tracked + backfilled: skip + - tracked + backfilling/in-progress: skip and let current backfill finish + - untracked or tracked-without-backfill-state: submit to `/repos/add` 2. **Register with Tap**: POST to `/repos/add` with the DID — Tap handles the actual repo export and event delivery 3. **Tap backfill flow**: Tap fetches full repo history from PDS via `com.atproto.sync.getRepo`, then delivers historical events (`live: false`) through the normal WebSocket channel 4. **Indexer processes normally**: The indexer's existing ingestion loop handles backfill events the same as live events — no special backfill code path needed diff --git a/docs/api/tasks/phase-1-mvp.md b/docs/api/tasks/phase-1-mvp.md index aa9b782..ff61666 100644 --- a/docs/api/tasks/phase-1-mvp.md +++ b/docs/api/tasks/phase-1-mvp.md @@ -203,7 +203,7 @@ Expose a usable public search API backed by Turso's Tantivy-backed FTS. ### Deliverables -- HTTP server (chi or net/http) +- HTTP server (net/http) - `GET /healthz` — liveness - `GET /readyz` — readiness (DB connectivity) - `GET /search` — keyword search with configurable mode @@ -214,7 +214,7 @@ Expose a usable public search API backed by Turso's Tantivy-backed FTS. ### Tasks -- [ ] Set up HTTP server with chi router +- [ ] Set up HTTP server with net/http router - [ ] Implement `/healthz` (always 200) and `/readyz` (SELECT 1 against DB) - [ ] Implement search repository with FTS queries: diff --git a/docs/app/tasks/phase-4.md b/docs/app/tasks/phase-4.md index 000328a..d8fa6df 100644 --- a/docs/app/tasks/phase-4.md +++ b/docs/app/tasks/phase-4.md @@ -3,7 +3,7 @@ ## OAuth Setup - [ ] Install `@atcute/oauth-browser-client` -- [ ] Host OAuth client metadata JSON at a public URL (or configure for local dev) +- [ ] Host OAuth client metadata JSON at a public URL & configure for local dev - [ ] Create `core/auth/oauth.ts` — call `configureOAuth()` with client metadata URL and redirect URI - [ ] Create `core/auth/session.ts` — session management: get, list, delete stored sessions - [ ] Create `core/auth/store.ts` — Pinia auth store with state machine (idle → authenticating → authenticated → error) @@ -62,12 +62,6 @@ - [ ] Add reaction button/picker to PR and issue detail views - [ ] Show reaction counts grouped by type -## Personalized Feed - -- [ ] When signed in, filter activity feed to show activity from followed users and starred repos -- [ ] Add "For You" / "Global" toggle on Activity tab -- [ ] If appview provides a personalized endpoint, use it; otherwise filter client-side - ## Profile Tab (Authenticated) - [ ] Wire Profile tab to show current user's profile data @@ -75,6 +69,12 @@ - [ ] Add logout button - [ ] Add account switcher UI +## Personalized Feed + +- [ ] When signed in, filter activity feed to show activity from followed users and starred repos +- [ ] Add "For You" / "Global" toggle on Activity tab +- [ ] If appview provides a personalized endpoint, use it; otherwise filter client-side + ## Quality - [ ] Test full OAuth flow: login → browse → star → follow → logout diff --git a/packages/api/internal/backfill/backfill.go b/packages/api/internal/backfill/backfill.go index 797f520..3a02102 100644 --- a/packages/api/internal/backfill/backfill.go +++ b/packages/api/internal/backfill/backfill.go @@ -83,21 +83,27 @@ func (r *Runner) Run(ctx context.Context, opts Options) error { } alreadyTracked := 0 + inProgress := 0 toSubmit := make([]string, 0, len(discovered)) for _, user := range discovered { - tracked, err := r.tap.IsTracked(ctx, user.DID) + status, err := r.tap.RepoStatus(ctx, user.DID) if err != nil { return fmt.Errorf("tap info for %s: %w", user.DID, err) } - if tracked { + if status.Tracked && status.Backfilled { alreadyTracked++ continue } + if status.Tracked && status.Backfilling { + inProgress++ + continue + } toSubmit = append(toSubmit, user.DID) } r.log.Info("tap classification complete", slog.Int("already_tracked", alreadyTracked), + slog.Int("backfill_in_progress", inProgress), slog.Int("to_submit", len(toSubmit)), ) @@ -130,6 +136,7 @@ func (r *Runner) Run(ctx context.Context, opts Options) error { r.log.Info("backfill complete", slog.Int("discovered_total", len(discovered)), slog.Int("already_tracked", alreadyTracked), + slog.Int("backfill_in_progress", inProgress), slog.Int("submitted", submitted), ) return nil diff --git a/packages/api/internal/backfill/backfill_test.go b/packages/api/internal/backfill/backfill_test.go index 24048dc..62acbc2 100644 --- a/packages/api/internal/backfill/backfill_test.go +++ b/packages/api/internal/backfill/backfill_test.go @@ -26,12 +26,15 @@ func (f *fakeFollowFetcher) ListFollowSubjects(_ context.Context, did string) ([ } type fakeTapAdmin struct { - tracked map[string]bool - added [][]string + statuses map[string]RepoStatus + added [][]string } -func (f *fakeTapAdmin) IsTracked(_ context.Context, did string) (bool, error) { - return f.tracked[did], nil +func (f *fakeTapAdmin) RepoStatus(_ context.Context, did string) (RepoStatus, error) { + if status, ok := f.statuses[did]; ok { + return status, nil + } + return RepoStatus{Found: false, Tracked: false}, nil } func (f *fakeTapAdmin) AddRepos(_ context.Context, dids []string) error { @@ -59,7 +62,7 @@ func TestRunner_DiscoveryAndSubmit(t *testing.T) { }, } follows := &fakeFollowFetcher{follows: map[string][]string{"did:plc:seed": {"did:plc:f1"}}} - tap := &fakeTapAdmin{tracked: map[string]bool{"did:plc:f1": true}} + tap := &fakeTapAdmin{statuses: map[string]RepoStatus{"did:plc:f1": {Found: true, Tracked: true, Backfilled: true}}} resolver := &fakeResolver{mapping: map[string]string{"alice.tangled.sh": "did:plc:seed"}} log := slog.New(slog.NewTextHandler(io.Discard, nil)) r := NewRunnerWithDeps(st, tap, resolver, follows, log) @@ -91,7 +94,7 @@ func TestRunner_DiscoveryAndSubmit(t *testing.T) { func TestRunner_DryRunSkipsMutations(t *testing.T) { st := &fakeStore{collaborators: map[string][]string{}} follows := &fakeFollowFetcher{follows: map[string][]string{}} - tap := &fakeTapAdmin{tracked: map[string]bool{}} + tap := &fakeTapAdmin{statuses: map[string]RepoStatus{}} resolver := &fakeResolver{mapping: map[string]string{"alice.tangled.sh": "did:plc:seed"}} log := slog.New(slog.NewTextHandler(io.Discard, nil)) r := NewRunnerWithDeps(st, tap, resolver, follows, log) @@ -116,3 +119,28 @@ func TestRunner_DryRunSkipsMutations(t *testing.T) { t.Fatalf("expected no tap submissions in dry-run, got %#v", tap.added) } } + +func TestRunner_SkipsInProgressBackfills(t *testing.T) { + st := &fakeStore{collaborators: map[string][]string{}} + follows := &fakeFollowFetcher{follows: map[string][]string{}} + tap := &fakeTapAdmin{statuses: map[string]RepoStatus{ + "did:plc:seed": {Found: true, Tracked: true, Backfilling: true}, + }} + resolver := &fakeResolver{mapping: map[string]string{"alice.tangled.sh": "did:plc:seed"}} + log := slog.New(slog.NewTextHandler(io.Discard, nil)) + r := NewRunnerWithDeps(st, tap, resolver, follows, log) + + dir := t.TempDir() + seedsPath := filepath.Join(dir, "seeds.txt") + if err := os.WriteFile(seedsPath, []byte("alice.tangled.sh\n"), 0o644); err != nil { + t.Fatalf("write seeds: %v", err) + } + + err := r.Run(context.Background(), Options{SeedsPath: seedsPath, MaxHops: 0}) + if err != nil { + t.Fatalf("run backfill: %v", err) + } + if len(tap.added) != 0 { + t.Fatalf("expected no submission for in-progress did, got %#v", tap.added) + } +} diff --git a/packages/api/internal/backfill/tap_admin.go b/packages/api/internal/backfill/tap_admin.go index 0e3fec2..b883010 100644 --- a/packages/api/internal/backfill/tap_admin.go +++ b/packages/api/internal/backfill/tap_admin.go @@ -6,6 +6,7 @@ import ( "encoding/base64" "encoding/json" "fmt" + "io" "net/http" "net/url" "strings" @@ -13,10 +14,18 @@ import ( ) type tapAdmin interface { - IsTracked(ctx context.Context, did string) (bool, error) + RepoStatus(ctx context.Context, did string) (RepoStatus, error) AddRepos(ctx context.Context, dids []string) error } +type RepoStatus struct { + Found bool + Tracked bool + Backfilled bool + Backfilling bool + State string +} + // HTTPTapAdmin calls Tap admin endpoints for backfill orchestration. type HTTPTapAdmin struct { baseURL string @@ -38,27 +47,55 @@ func NewHTTPTapAdmin(tapURL, password string) (*HTTPTapAdmin, error) { }, nil } -func (t *HTTPTapAdmin) IsTracked(ctx context.Context, did string) (bool, error) { +func (t *HTTPTapAdmin) RepoStatus(ctx context.Context, did string) (RepoStatus, error) { endpoint := fmt.Sprintf("%s/info/%s", t.baseURL, url.PathEscape(did)) req, err := http.NewRequestWithContext(ctx, http.MethodGet, endpoint, nil) if err != nil { - return false, fmt.Errorf("build tap info request: %w", err) + return RepoStatus{}, fmt.Errorf("build tap info request: %w", err) } t.addAuth(req) resp, err := t.client.Do(req) if err != nil { - return false, fmt.Errorf("tap info request: %w", err) + return RepoStatus{}, fmt.Errorf("tap info request: %w", err) } defer resp.Body.Close() if resp.StatusCode == http.StatusNotFound { - return false, nil + return RepoStatus{Found: false}, nil } if resp.StatusCode < 200 || resp.StatusCode >= 300 { - return false, fmt.Errorf("tap info request failed: status %d", resp.StatusCode) + return RepoStatus{}, fmt.Errorf("tap info request failed: status %d", resp.StatusCode) + } + + body, err := io.ReadAll(resp.Body) + if err != nil { + return RepoStatus{}, fmt.Errorf("read tap info response: %w", err) + } + if len(bytes.TrimSpace(body)) == 0 { + return RepoStatus{Found: true, Tracked: true}, nil + } + + var payload map[string]any + if err := json.Unmarshal(body, &payload); err != nil { + return RepoStatus{}, fmt.Errorf("decode tap info response: %w", err) } - return true, nil + + status := RepoStatus{Found: true, Tracked: true} + if tracked, ok := boolFromAnyWithPresence(payload, "tracked", "isTracked", "enabled", "registered"); ok { + status.Tracked = tracked + } + status.Backfilled = boolFromAny(payload, "backfilled", "isBackfilled", "complete", "done") + status.Backfilling = boolFromAny(payload, "backfilling", "inProgress", "in_progress", "pendingBackfill") + status.State = stringFromAny(payload, "status", "state") + + if stateImpliesBackfilled(status.State) { + status.Backfilled = true + } + if stateImpliesBackfilling(status.State) { + status.Backfilling = true + } + return status, nil } func (t *HTTPTapAdmin) AddRepos(ctx context.Context, dids []string) error { @@ -126,3 +163,62 @@ func normalizeTapBaseURL(raw string) (string, error) { } return strings.TrimSuffix(u.String(), "/"), nil } + +func boolFromAny(payload map[string]any, keys ...string) bool { + v, _ := boolFromAnyWithPresence(payload, keys...) + return v +} + +func boolFromAnyWithPresence(payload map[string]any, keys ...string) (bool, bool) { + for _, key := range keys { + raw, ok := payload[key] + if !ok { + continue + } + switch v := raw.(type) { + case bool: + return v, true + case float64: + return v != 0, true + case string: + switch strings.ToLower(strings.TrimSpace(v)) { + case "true", "1", "yes", "y", "active", "complete", "done", "backfilled", "backfilling", "in_progress", "in-progress": + return true, true + case "false", "0", "no", "n", "inactive": + return false, true + } + } + } + return false, false +} + +func stringFromAny(payload map[string]any, keys ...string) string { + for _, key := range keys { + raw, ok := payload[key] + if !ok { + continue + } + if value, ok := raw.(string); ok { + return strings.TrimSpace(strings.ToLower(value)) + } + } + return "" +} + +func stateImpliesBackfilled(state string) bool { + switch strings.TrimSpace(strings.ToLower(state)) { + case "backfilled", "complete", "completed", "done", "ready", "synced": + return true + default: + return false + } +} + +func stateImpliesBackfilling(state string) bool { + switch strings.TrimSpace(strings.ToLower(state)) { + case "backfilling", "in-progress", "in_progress", "pending", "queued", "running": + return true + default: + return false + } +} diff --git a/packages/api/internal/backfill/tap_admin_test.go b/packages/api/internal/backfill/tap_admin_test.go new file mode 100644 index 0000000..ee58466 --- /dev/null +++ b/packages/api/internal/backfill/tap_admin_test.go @@ -0,0 +1,86 @@ +package backfill + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "testing" +) + +func TestRepoStatusNotFound(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/info/did:plc:missing" { + t.Fatalf("unexpected path: %s", r.URL.Path) + } + http.NotFound(w, r) + })) + defer ts.Close() + + admin, err := NewHTTPTapAdmin(ts.URL+"/channel", "") + if err != nil { + t.Fatalf("new admin: %v", err) + } + + status, err := admin.RepoStatus(context.Background(), "did:plc:missing") + if err != nil { + t.Fatalf("repo status: %v", err) + } + if status.Found { + t.Fatalf("expected found=false, got %#v", status) + } + if status.Tracked { + t.Fatalf("expected tracked=false for 404") + } +} + +func TestRepoStatusParsesBackfillState(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/info/did:plc:abc" { + t.Fatalf("unexpected path: %s", r.URL.Path) + } + w.Header().Set("Content-Type", "application/json") + _, _ = fmt.Fprint(w, `{"tracked":true,"status":"backfilling"}`) + })) + defer ts.Close() + + admin, err := NewHTTPTapAdmin(ts.URL+"/channel", "") + if err != nil { + t.Fatalf("new admin: %v", err) + } + + status, err := admin.RepoStatus(context.Background(), "did:plc:abc") + if err != nil { + t.Fatalf("repo status: %v", err) + } + if !status.Found || !status.Tracked { + t.Fatalf("expected tracked found status, got %#v", status) + } + if !status.Backfilling { + t.Fatalf("expected backfilling=true, got %#v", status) + } + if status.Backfilled { + t.Fatalf("expected backfilled=false, got %#v", status) + } +} + +func TestRepoStatusParsesExplicitTrackedFalse(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = fmt.Fprint(w, `{"tracked":false}`) + })) + defer ts.Close() + + admin, err := NewHTTPTapAdmin(ts.URL+"/channel", "") + if err != nil { + t.Fatalf("new admin: %v", err) + } + + status, err := admin.RepoStatus(context.Background(), "did:plc:abc") + if err != nil { + t.Fatalf("repo status: %v", err) + } + if status.Tracked { + t.Fatalf("expected tracked=false, got %#v", status) + } +} diff --git a/packages/api/internal/ingest/ingest.go b/packages/api/internal/ingest/ingest.go index 4273b2a..a1b19c1 100644 --- a/packages/api/internal/ingest/ingest.go +++ b/packages/api/internal/ingest/ingest.go @@ -5,6 +5,7 @@ import ( "fmt" "log/slog" "math" + "strconv" "strings" "sync" "time" @@ -33,6 +34,7 @@ type Runner struct { allowlist allowlist consumerName string log *slog.Logger + resumeCursor int64 statusMu sync.Mutex lastCursor string @@ -55,6 +57,9 @@ func NewRunner(st store.Store, registry *normalize.Registry, tap client, indexed func (r *Runner) Run(ctx context.Context) error { defer r.tap.Close() + if err := r.initializeCursor(ctx); err != nil { + return err + } go r.runStatusLogger(ctx) @@ -72,6 +77,18 @@ func (r *Runner) Run(ctx context.Context) error { continue } + if r.shouldSkipEvent(event.ID) { + if err := r.tap.AckEvent(ctx, event.ID); err != nil { + r.log.Warn("tap ack skipped event failed", + slog.Int64("event_id", event.ID), + slog.String("error", err.Error()), + ) + continue + } + r.log.Info("skipped previously-processed event", slog.Int64("event_id", event.ID), slog.Int64("resume_cursor", r.resumeCursor)) + continue + } + if err := r.processWithRetry(ctx, event); err != nil { if ctx.Err() != nil { return nil @@ -81,6 +98,37 @@ func (r *Runner) Run(ctx context.Context) error { } } +func (r *Runner) initializeCursor(ctx context.Context) error { + state, err := r.store.GetSyncState(ctx, r.consumerName) + if err != nil { + return fmt.Errorf("load sync cursor: %w", err) + } + if state == nil || strings.TrimSpace(state.Cursor) == "" { + r.log.Info("indexer cursor resume disabled", slog.String("reason", "no prior sync_state")) + return nil + } + + cursor, err := strconv.ParseInt(strings.TrimSpace(state.Cursor), 10, 64) + if err != nil { + r.log.Warn("indexer cursor parse failed; resume disabled", + slog.String("cursor", state.Cursor), + slog.String("error", err.Error()), + ) + return nil + } + + r.resumeCursor = cursor + r.statusMu.Lock() + r.lastCursor = state.Cursor + r.statusMu.Unlock() + r.log.Info("indexer cursor resume enabled", slog.Int64("resume_cursor", cursor)) + return nil +} + +func (r *Runner) shouldSkipEvent(eventID int64) bool { + return r.resumeCursor > 0 && eventID <= r.resumeCursor +} + func (r *Runner) processWithRetry(ctx context.Context, event normalize.TapRecordEvent) error { attempt := 0 for { @@ -227,16 +275,42 @@ func (r *Runner) processRecordEvent(ctx context.Context, event normalize.TapReco func (r *Runner) advanceCursorAndAck(ctx context.Context, eventID int64) error { cursor := fmt.Sprintf("%d", eventID) - if err := r.store.SetSyncState(ctx, r.consumerName, cursor); err != nil { + if err := r.tap.AckEvent(ctx, eventID); err != nil { return err } - if err := r.tap.AckEvent(ctx, eventID); err != nil { + if err := r.persistCursorWithRetry(ctx, cursor, eventID); err != nil { return err } r.markProcessed(cursor) return nil } +func (r *Runner) persistCursorWithRetry(ctx context.Context, cursor string, eventID int64) error { + attempt := 0 + for { + if ctx.Err() != nil { + return ctx.Err() + } + if err := r.store.SetSyncState(ctx, r.consumerName, cursor); err == nil { + return nil + } else { + attempt++ + backoff := retryBackoff(attempt) + r.log.Error("cursor persist failed after ack", + slog.Int64("event_id", eventID), + slog.Int("attempt", attempt), + slog.Duration("retry_in", backoff), + slog.String("error", err.Error()), + ) + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(backoff): + } + } + } +} + func (r *Runner) runStatusLogger(ctx context.Context) { ticker := time.NewTicker(statusLogInterval) defer ticker.Stop() diff --git a/packages/api/internal/ingest/ingest_test.go b/packages/api/internal/ingest/ingest_test.go index cb43765..dff7c60 100644 --- a/packages/api/internal/ingest/ingest_test.go +++ b/packages/api/internal/ingest/ingest_test.go @@ -12,6 +12,7 @@ import ( type fakeTapClient struct { acked []int64 + onAck func(id int64) } func (f *fakeTapClient) ReadEvent(_ context.Context) (normalize.TapRecordEvent, error) { @@ -19,6 +20,9 @@ func (f *fakeTapClient) ReadEvent(_ context.Context) (normalize.TapRecordEvent, } func (f *fakeTapClient) AckEvent(_ context.Context, id int64) error { + if f.onAck != nil { + f.onAck(id) + } f.acked = append(f.acked, id) return nil } @@ -29,9 +33,11 @@ type fakeStore struct { docs map[string]*store.Document deleted map[string]bool syncCursor string + initialSync *store.SyncState recordStates map[string]string handles map[string]string enqueued map[string]bool + onSetSync func() } func newFakeStore() *fakeStore { @@ -60,10 +66,20 @@ func (f *fakeStore) MarkDeleted(_ context.Context, id string) error { } func (f *fakeStore) GetSyncState(_ context.Context, _ string) (*store.SyncState, error) { - return nil, nil + if f.initialSync != nil { + state := *f.initialSync + return &state, nil + } + if f.syncCursor == "" { + return nil, nil + } + return &store.SyncState{ConsumerName: "indexer-tap-v1", Cursor: f.syncCursor}, nil } func (f *fakeStore) SetSyncState(_ context.Context, _ string, cursor string) error { + if f.onSetSync != nil { + f.onSetSync() + } f.syncCursor = cursor return nil } @@ -264,3 +280,60 @@ func TestAllowlistMatching(t *testing.T) { t.Fatal("unexpected match") } } + +func TestRunner_InitializeCursorResume(t *testing.T) { + st := newFakeStore() + st.initialSync = &store.SyncState{ConsumerName: "indexer-tap-v1", Cursor: "150"} + tap := &fakeTapClient{} + r := newRunnerForTest(st, tap, "sh.tangled.*") + + if err := r.initializeCursor(context.Background()); err != nil { + t.Fatalf("initialize cursor: %v", err) + } + if r.resumeCursor != 150 { + t.Fatalf("resume cursor: got %d want 150", r.resumeCursor) + } + if !r.shouldSkipEvent(149) { + t.Fatalf("expected event 149 to be skipped") + } + if !r.shouldSkipEvent(150) { + t.Fatalf("expected event 150 to be skipped") + } + if r.shouldSkipEvent(151) { + t.Fatalf("expected event 151 to be processed") + } +} + +func TestRunner_AckBeforeCursorPersist(t *testing.T) { + st := newFakeStore() + tap := &fakeTapClient{} + r := newRunnerForTest(st, tap, "sh.tangled.*") + + acked := false + tap.onAck = func(_ int64) { acked = true } + st.onSetSync = func() { + if !acked { + t.Fatalf("cursor persisted before ack") + } + } + + event := normalize.TapRecordEvent{ + ID: 901, + Type: "record", + Record: &normalize.TapRecord{ + DID: "did:plc:author", + Collection: "sh.tangled.repo", + RKey: "repo1", + Action: "create", + CID: "cid-1", + Record: map[string]any{ + "name": "repo-one", + "description": "test repo", + }, + }, + } + + if err := r.processEvent(context.Background(), event); err != nil { + t.Fatalf("process event: %v", err) + } +} diff --git a/packages/api/internal/tapclient/tapclient.go b/packages/api/internal/tapclient/tapclient.go index aea12ba..8b20840 100644 --- a/packages/api/internal/tapclient/tapclient.go +++ b/packages/api/internal/tapclient/tapclient.go @@ -8,6 +8,7 @@ import ( "log/slog" "math/rand/v2" "net/http" + "os" "strconv" "strings" "sync" @@ -30,20 +31,23 @@ type Client struct { password string log *slog.Logger - mu sync.Mutex - conn *websocket.Conn - ackAsJSON bool + mu sync.Mutex + conn *websocket.Conn + ackAsJSON bool + disableAcks bool } func New(url, password string, log *slog.Logger) *Client { if log == nil { log = slog.Default() } + disableAcks, _ := strconv.ParseBool(strings.TrimSpace(os.Getenv("TAP_DISABLE_ACKS"))) return &Client{ - url: url, - password: password, - log: log, - ackAsJSON: true, + url: url, + password: password, + log: log, + ackAsJSON: true, + disableAcks: disableAcks, } } @@ -75,6 +79,10 @@ func (c *Client) ReadEvent(ctx context.Context) (normalize.TapRecordEvent, error } func (c *Client) AckEvent(ctx context.Context, id int64) error { + if c.disableAcks { + return nil + } + conn, err := c.ensureConnected(ctx) if err != nil { return err @@ -91,6 +99,8 @@ func (c *Client) AckEvent(ctx context.Context, id int64) error { } else if isConnectionWriteError(err) { c.resetConn(websocket.StatusInternalError, "ack json write failed") return fmt.Errorf("ack event %d: %w", id, err) + } else if !isAckFormatError(err) { + return fmt.Errorf("ack event %d: %w", id, err) } c.log.Warn("tap ack json failed; trying plain id", slog.Int64("event_id", id)) @@ -222,3 +232,14 @@ func isConnectionWriteError(err error) bool { strings.Contains(msg, "closed network connection") || strings.Contains(msg, "i/o timeout") } + +func isAckFormatError(err error) bool { + if err == nil { + return false + } + msg := strings.ToLower(err.Error()) + return strings.Contains(msg, "invalid") || + strings.Contains(msg, "unsupported") || + strings.Contains(msg, "bad payload") || + strings.Contains(msg, "unexpected message") +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index c01e68f..2b1dddc 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -16,6 +16,9 @@ importers: '@atcute/client': specifier: ^4.2.1 version: 4.2.1 + '@atcute/oauth-browser-client': + specifier: ^3.0.0 + version: 3.0.0(@atcute/identity@1.1.4) '@atcute/tangled': specifier: ^1.0.17 version: 1.0.17 @@ -152,18 +155,41 @@ packages: '@atcute/client@4.2.1': resolution: {integrity: sha512-ZBFM2pW075JtgGFu5g7HHZBecrClhlcNH8GVP9Zz1aViWR+cjjBsTpeE63rJs+FCOHFYlirUyo5L8SGZ4kMINw==} + '@atcute/identity-resolver@1.2.2': + resolution: {integrity: sha512-eUh/UH4bFvuXS0X7epYCeJC/kj4rbBXfSRumLEH4smMVwNOgTo7cL/0Srty+P/qVPoZEyXdfEbS0PHJyzoXmHw==} + peerDependencies: + '@atcute/identity': ^1.0.0 + '@atcute/identity@1.1.4': resolution: {integrity: sha512-RCw1IqflfuSYCxK5m0lZCm0UnvIzcUnuhngiBhJEJb9a9Mc2SEf1xP3H8N5r8pvEH1LoAYd6/zrvCNU+uy9esw==} '@atcute/lexicons@1.2.9': resolution: {integrity: sha512-/RRHm2Cw9o8Mcsrq0eo8fjS9okKYLGfuFwrQ0YoP/6sdSDsXshaTLJsvLlcUcaDaSJ1YFOuHIo3zr2Om2F/16g==} + '@atcute/multibase@1.2.0': + resolution: {integrity: sha512-ZK2GRra+qIYq9nNuQB52m2ul0hOmCQEtPobGfTSUxm7pF0OGEkWGkWHugFhNEDVzHzTwPxHp6VGotdZFue4lYQ==} + + '@atcute/oauth-browser-client@3.0.0': + resolution: {integrity: sha512-7AbKV8tTe7aRJNJV7gCcWHSVEADb2nr58O1p7dQsf73HSe9pvlBkj/Vk1yjjtH691uAVYkwhHSh0bC7D8XdwJw==} + + '@atcute/oauth-crypto@0.1.0': + resolution: {integrity: sha512-qZYDCNLF/4B6AndYT1rsQelN8621AC5u/sL5PHvlr/qqAbmmUwCBGjEgRSyZtHE1AqD60VNiSMlOgAuEQTSl3w==} + + '@atcute/oauth-keyset@0.1.0': + resolution: {integrity: sha512-+wqT/+I5Lg9VzKnKY3g88+N45xbq+wsdT6bHDGqCVa2u57gRvolFF4dY+weMfc/OX641BIZO6/o+zFtKBsMQnQ==} + + '@atcute/oauth-types@0.1.1': + resolution: {integrity: sha512-u+3KMjse3Uc/9hDyilu1QVN7IpcnjVXgRzhddzBB8Uh6wePHNVBDdi9wQvFTVVA3zmxtMJVptXRyLLg6Ou9bqg==} + '@atcute/tangled@1.0.17': resolution: {integrity: sha512-YE2KXYSawWbASt/XB6bi0gPGkF089+wsbxQcqtksrSy7j02LGOKUFgX1gfG1ZImv40eQHjFnhLHr1OZVXeflZA==} '@atcute/uint8array@1.1.1': resolution: {integrity: sha512-3LsC8XB8TKe9q/5hOA5sFuzGaIFdJZJNewC5OKa3o/eU6+K7JR6see9Zy2JbQERNVnRl11EzbNov1efgLMAs4g==} + '@atcute/util-fetch@1.0.5': + resolution: {integrity: sha512-qjHj01BGxjSjIFdPiAjSARnodJIIyKxnCMMEcXMESo9TAyND6XZQqrie5fia+LlYWVXdpsTds8uFQwc9jdKTig==} + '@atcute/util-text@1.2.0': resolution: {integrity: sha512-b8WSh+Z7K601eUFFmTFj8QPKDO8Ic0VDDj63sdKzpkm+ySQKsYT5nXekViGqFVKbyKj1V5FyvZvgXad6/aI4QQ==} @@ -2550,6 +2576,11 @@ packages: engines: {node: ^10 || ^12 || ^13.7 || ^14 || >=15.0.1} hasBin: true + nanoid@5.1.7: + resolution: {integrity: sha512-ua3NDgISf6jdwezAheMOk4mbE1LXjm1DfMUDMuJf4AqxLFK3ccGpgWizwa5YV7Yz9EpXwEaWoRXSb/BnV0t5dQ==} + engines: {node: ^18 || >=20} + hasBin: true + native-run@2.0.3: resolution: {integrity: sha512-U1PllBuzW5d1gfan+88L+Hky2eZx+9gv3Pf6rNBxKbORxi7boHzqiA6QFGSnqMem4j0A9tZ08NMIs5+0m/VS1Q==} engines: {node: '>=16.0.0'} @@ -3402,6 +3433,13 @@ snapshots: '@atcute/identity': 1.1.4 '@atcute/lexicons': 1.2.9 + '@atcute/identity-resolver@1.2.2(@atcute/identity@1.1.4)': + dependencies: + '@atcute/identity': 1.1.4 + '@atcute/lexicons': 1.2.9 + '@atcute/util-fetch': 1.0.5 + '@badrap/valita': 0.4.6 + '@atcute/identity@1.1.4': dependencies: '@atcute/lexicons': 1.2.9 @@ -3414,6 +3452,40 @@ snapshots: '@standard-schema/spec': 1.1.0 esm-env: 1.2.2 + '@atcute/multibase@1.2.0': + dependencies: + '@atcute/uint8array': 1.1.1 + + '@atcute/oauth-browser-client@3.0.0(@atcute/identity@1.1.4)': + dependencies: + '@atcute/client': 4.2.1 + '@atcute/identity-resolver': 1.2.2(@atcute/identity@1.1.4) + '@atcute/lexicons': 1.2.9 + '@atcute/multibase': 1.2.0 + '@atcute/oauth-crypto': 0.1.0 + '@atcute/oauth-types': 0.1.1 + nanoid: 5.1.7 + transitivePeerDependencies: + - '@atcute/identity' + + '@atcute/oauth-crypto@0.1.0': + dependencies: + '@atcute/multibase': 1.2.0 + '@atcute/uint8array': 1.1.1 + '@badrap/valita': 0.4.6 + nanoid: 5.1.7 + + '@atcute/oauth-keyset@0.1.0': + dependencies: + '@atcute/oauth-crypto': 0.1.0 + + '@atcute/oauth-types@0.1.1': + dependencies: + '@atcute/identity': 1.1.4 + '@atcute/lexicons': 1.2.9 + '@atcute/oauth-keyset': 0.1.0 + '@badrap/valita': 0.4.6 + '@atcute/tangled@1.0.17': dependencies: '@atcute/atproto': 3.1.10 @@ -3421,6 +3493,10 @@ snapshots: '@atcute/uint8array@1.1.1': {} + '@atcute/util-fetch@1.0.5': + dependencies: + '@badrap/valita': 0.4.6 + '@atcute/util-text@1.2.0': dependencies: unicode-segmenter: 0.14.5 @@ -6029,6 +6105,8 @@ snapshots: nanoid@3.3.11: {} + nanoid@5.1.7: {} + native-run@2.0.3: dependencies: '@ionic/utils-fs': 3.1.7