diff --git a/docs/TODO.md b/docs/TODO.md new file mode 100644 --- /dev/null +++ b/docs/TODO.md @@ -0,0 +1,13 @@ +--- +title: To-Dos +updated: 2026-03-23 +--- + +A catch-all for ideas, issues/bugs, and future work that doesn't fit into the current specs or tasks. This is a "parking lot." + +## App + +- Repo stars, forks, etc. are not properly parsed from JSON. +- ATOM/RSS feed link for repos: (`tangled.org/{did}/{repo}/feed.atom`) + +## API diff --git a/packages/api/main.go b/packages/api/main.go --- a/packages/api/main.go +++ b/packages/api/main.go @@ -10,7 +10,11 @@ "github.com/spf13/cobra" "tangled.org/desertthunder.dev/twister/internal/config" + "tangled.org/desertthunder.dev/twister/internal/ingest" + "tangled.org/desertthunder.dev/twister/internal/normalize" "tangled.org/desertthunder.dev/twister/internal/observability" + "tangled.org/desertthunder.dev/twister/internal/store" + "tangled.org/desertthunder.dev/twister/internal/tapclient" ) var ( @@ -81,9 +85,33 @@ } log := observability.NewLogger(cfg) log.Info("starting indexer", slog.String("service", "indexer"), slog.String("version", version)) + + if cfg.TapURL == "" { + return fmt.Errorf("TAP_URL is required for indexer") + } + + db, err := store.Open(cfg.TursoURL, cfg.TursoToken) + if err != nil { + return fmt.Errorf("open database: %w", err) + } + defer db.Close() + + if err := store.Migrate(db); err != nil { + return fmt.Errorf("migrate database: %w", err) + } + + st := store.New(db) + registry := normalize.NewRegistry() + tap := tapclient.New(cfg.TapURL, cfg.TapAuthPassword, log) + runner := ingest.NewRunner(st, registry, tap, cfg.IndexedCollections, log) + ctx, cancel := baseContext() defer cancel() - <-ctx.Done() + + if err := runner.Run(ctx); err != nil { + return fmt.Errorf("run indexer: %w", err) + } + log.Info("shutting down indexer") return nil }, diff --git a/docs/api/tasks/phase-1-mvp.md b/docs/api/tasks/phase-1-mvp.md --- a/docs/api/tasks/phase-1-mvp.md +++ b/docs/api/tasks/phase-1-mvp.md @@ -51,7 +51,7 @@ ### Tasks -- [ ] Define Tap event DTOs matching the documented event shape: +- [x] Define Tap event DTOs matching the documented event shape: ```go type TapEvent struct { @@ -78,12 +78,12 @@ } ``` -- [ ] Implement WebSocket client: +- [x] Implement WebSocket client: - Connect to `TAP_URL` (e.g., `wss://tap.railway.internal/channel`) - HTTP Basic auth with `admin:TAP_AUTH_PASSWORD` - Auto-reconnect with exponential backoff - Ack protocol: send event `id` back after successful processing -- [ ] Implement ingestion loop: +- [x] Implement ingestion loop: 1. Receive event from WebSocket 2. If `type == "identity"` → update handle cache, ack, continue 3. If `type == "record"` → check collection allowlist @@ -91,13 +91,13 @@ 5. Decode `record.record` via adapter registry 6. Normalize to `Document` 7. Upsert to store - 8. Schedule embedding job if eligible (Phase 2) + 8. Schedule embedding job if eligible ([Phase 2](phase-2-semantic.md)) 9. Persist cursor (event ID) after successful DB commit 10. Ack the event -- [ ] Implement collection allowlist from `INDEXED_COLLECTIONS` config -- [ ] Handle state events (`sh.tangled.repo.issue.state`, `sh.tangled.repo.pull.status`) → update `record_state` -- [ ] Handle normalization failures: log, skip, advance cursor -- [ ] Handle DB failures: retry with backoff, do not advance cursor +- [x] Implement collection allowlist from `INDEXED_COLLECTIONS` config +- [x] Handle state events (`sh.tangled.repo.issue.state`, `sh.tangled.repo.pull.status`) → update `record_state` +- [x] Handle normalization failures: log, skip, advance cursor +- [x] Handle DB failures: retry with backoff, do not advance cursor ### Verification diff --git a/packages/api/internal/ingest/ingest.go b/packages/api/internal/ingest/ingest.go --- a/packages/api/internal/ingest/ingest.go +++ b/packages/api/internal/ingest/ingest.go @@ -1,1 +1,281 @@ package ingest + +import ( + "context" + "fmt" + "log/slog" + "math" + "strings" + "time" + + "tangled.org/desertthunder.dev/twister/internal/normalize" + "tangled.org/desertthunder.dev/twister/internal/store" +) + +const ( + defaultConsumerName = "indexer-tap-v1" + maxDBRetryBackoff = 5 * time.Second +) + +type client interface { + ReadEvent(ctx context.Context) (normalize.TapRecordEvent, error) + AckEvent(ctx context.Context, id int64) error + Close() error +} + +// Runner ingests Tap events into the store. +type Runner struct { + store store.Store + registry *normalize.Registry + tap client + allowlist allowlist + consumerName string + log *slog.Logger +} + +func NewRunner(st store.Store, registry *normalize.Registry, tap client, indexedCollections string, log *slog.Logger) *Runner { + if log == nil { + log = slog.Default() + } + return &Runner{ + store: st, + registry: registry, + tap: tap, + allowlist: parseAllowlist(indexedCollections), + consumerName: defaultConsumerName, + log: log, + } +} + +func (r *Runner) Run(ctx context.Context) error { + defer r.tap.Close() + + for { + if ctx.Err() != nil { + return nil + } + + event, err := r.tap.ReadEvent(ctx) + if err != nil { + if ctx.Err() != nil { + return nil + } + r.log.Warn("tap read error", slog.String("error", err.Error())) + continue + } + + if err := r.processWithRetry(ctx, event); err != nil { + if ctx.Err() != nil { + return nil + } + return err + } + } +} + +func (r *Runner) processWithRetry(ctx context.Context, event normalize.TapRecordEvent) error { + attempt := 0 + for { + if ctx.Err() != nil { + return ctx.Err() + } + + err := r.processEvent(ctx, event) + if err == nil { + return nil + } + + attempt++ + backoff := retryBackoff(attempt) + r.log.Warn("ingest retry", + slog.Int64("event_id", event.ID), + 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) processEvent(ctx context.Context, event normalize.TapRecordEvent) error { + switch event.Type { + case "identity": + if event.Identity == nil { + return r.advanceCursorAndAck(ctx, event.ID) + } + id := event.Identity + if err := r.store.UpsertIdentityHandle(ctx, id.DID, id.Handle, id.IsActive, id.Status); err != nil { + return err + } + return r.advanceCursorAndAck(ctx, event.ID) + case "record": + return r.processRecordEvent(ctx, event) + default: + return r.advanceCursorAndAck(ctx, event.ID) + } +} + +func (r *Runner) processRecordEvent(ctx context.Context, event normalize.TapRecordEvent) error { + if event.Record == nil { + return r.advanceCursorAndAck(ctx, event.ID) + } + + record := event.Record + if !r.allowlist.match(record.Collection) { + return r.advanceCursorAndAck(ctx, event.ID) + } + + if handler, ok := r.registry.StateHandler(record.Collection); ok { + if record.Action == "delete" { + return r.advanceCursorAndAck(ctx, event.ID) + } + update, err := handler.HandleState(event) + if err != nil { + r.log.Warn("state normalization failed", + slog.Int64("event_id", event.ID), + slog.String("collection", record.Collection), + slog.String("did", record.DID), + slog.String("rkey", record.RKey), + slog.String("error", err.Error()), + ) + return r.advanceCursorAndAck(ctx, event.ID) + } + if err := r.store.UpdateRecordState(ctx, update.SubjectURI, update.State); err != nil { + return err + } + return r.advanceCursorAndAck(ctx, event.ID) + } + + adapter, ok := r.registry.Adapter(record.Collection) + if !ok { + return r.advanceCursorAndAck(ctx, event.ID) + } + + switch record.Action { + case "delete": + docID := normalize.StableID(record.DID, record.Collection, record.RKey) + if err := r.store.MarkDeleted(ctx, docID); err != nil { + return err + } + return r.advanceCursorAndAck(ctx, event.ID) + case "create", "update": + if record.Record == nil { + r.log.Warn("record payload missing", + slog.Int64("event_id", event.ID), + slog.String("collection", record.Collection), + slog.String("did", record.DID), + slog.String("rkey", record.RKey), + ) + return r.advanceCursorAndAck(ctx, event.ID) + } + default: + return r.advanceCursorAndAck(ctx, event.ID) + } + + doc, err := adapter.Normalize(event) + if err != nil { + r.log.Warn("normalization failed", + slog.Int64("event_id", event.ID), + slog.String("collection", record.Collection), + slog.String("did", record.DID), + slog.String("rkey", record.RKey), + slog.String("error", err.Error()), + ) + return r.advanceCursorAndAck(ctx, event.ID) + } + + handle, err := r.store.GetIdentityHandle(ctx, record.DID) + if err != nil { + return err + } + if handle != "" { + doc.AuthorHandle = handle + if doc.RecordType == "profile" { + doc.Title = handle + } + } + + if err := r.store.UpsertDocument(ctx, doc); err != nil { + return err + } + + if adapter.Searchable(record.Record) { + if err := r.store.EnqueueEmbeddingJob(ctx, doc.ID); err != nil { + r.log.Warn("embedding enqueue failed", + slog.Int64("event_id", event.ID), + slog.String("document_id", doc.ID), + slog.String("error", err.Error()), + ) + } + } + + return r.advanceCursorAndAck(ctx, event.ID) +} + +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 { + return err + } + if err := r.tap.AckEvent(ctx, eventID); err != nil { + return err + } + return nil +} + +type allowlist struct { + entries []string +} + +func parseAllowlist(raw string) allowlist { + if strings.TrimSpace(raw) == "" { + return allowlist{entries: nil} + } + parts := strings.FieldsFunc(raw, func(r rune) bool { + return r == ',' || r == ' ' || r == '\n' || r == '\t' + }) + entries := make([]string, 0, len(parts)) + for _, part := range parts { + entry := strings.TrimSpace(part) + if entry == "" { + continue + } + entries = append(entries, entry) + } + return allowlist{entries: entries} +} + +func (a allowlist) match(collection string) bool { + if len(a.entries) == 0 { + return true + } + for _, entry := range a.entries { + if entry == collection { + return true + } + if strings.HasSuffix(entry, "*") { + prefix := strings.TrimSuffix(entry, "*") + if strings.HasPrefix(collection, prefix) { + return true + } + } + } + return false +} + +func retryBackoff(attempt int) time.Duration { + if attempt < 1 { + attempt = 1 + } + exponent := math.Pow(2, float64(attempt-1)) + d := time.Duration(float64(time.Second) * exponent) + if d > maxDBRetryBackoff { + return maxDBRetryBackoff + } + return d +} diff --git a/packages/api/internal/ingest/ingest_test.go b/packages/api/internal/ingest/ingest_test.go new file mode 100644 --- /dev/null +++ b/packages/api/internal/ingest/ingest_test.go @@ -0,0 +1,254 @@ +package ingest + +import ( + "context" + "io" + "log/slog" + "testing" + + "tangled.org/desertthunder.dev/twister/internal/normalize" + "tangled.org/desertthunder.dev/twister/internal/store" +) + +type fakeTapClient struct { + acked []int64 +} + +func (f *fakeTapClient) ReadEvent(_ context.Context) (normalize.TapRecordEvent, error) { + return normalize.TapRecordEvent{}, io.EOF +} + +func (f *fakeTapClient) AckEvent(_ context.Context, id int64) error { + f.acked = append(f.acked, id) + return nil +} + +func (f *fakeTapClient) Close() error { return nil } + +type fakeStore struct { + docs map[string]*store.Document + deleted map[string]bool + syncCursor string + recordStates map[string]string + handles map[string]string + enqueued map[string]bool +} + +func newFakeStore() *fakeStore { + return &fakeStore{ + docs: make(map[string]*store.Document), + deleted: make(map[string]bool), + recordStates: make(map[string]string), + handles: make(map[string]string), + enqueued: make(map[string]bool), + } +} + +func (f *fakeStore) UpsertDocument(_ context.Context, doc *store.Document) error { + clone := *doc + f.docs[doc.ID] = &clone + return nil +} + +func (f *fakeStore) GetDocument(_ context.Context, id string) (*store.Document, error) { + return f.docs[id], nil +} + +func (f *fakeStore) MarkDeleted(_ context.Context, id string) error { + f.deleted[id] = true + return nil +} + +func (f *fakeStore) GetSyncState(_ context.Context, _ string) (*store.SyncState, error) { + return nil, nil +} + +func (f *fakeStore) SetSyncState(_ context.Context, _ string, cursor string) error { + f.syncCursor = cursor + return nil +} + +func (f *fakeStore) UpdateRecordState(_ context.Context, subjectURI string, state string) error { + f.recordStates[subjectURI] = state + return nil +} + +func (f *fakeStore) UpsertIdentityHandle(_ context.Context, did, handle string, _ bool, _ string) error { + f.handles[did] = handle + return nil +} + +func (f *fakeStore) GetIdentityHandle(_ context.Context, did string) (string, error) { + return f.handles[did], nil +} + +func (f *fakeStore) EnqueueEmbeddingJob(_ context.Context, documentID string) error { + f.enqueued[documentID] = true + return nil +} + +func newRunnerForTest(st *fakeStore, tap *fakeTapClient, indexedCollections string) *Runner { + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + return NewRunner(st, normalize.NewRegistry(), tap, indexedCollections, logger) +} + +func TestRunner_ProcessIdentityEvent(t *testing.T) { + st := newFakeStore() + tap := &fakeTapClient{} + r := newRunnerForTest(st, tap, "sh.tangled.*") + + event := normalize.TapRecordEvent{ + ID: 101, + Type: "identity", + Identity: &normalize.TapIdentity{ + DID: "did:plc:abc", + Handle: "alice.tangled.org", + IsActive: true, + Status: "active", + }, + } + + if err := r.processEvent(context.Background(), event); err != nil { + t.Fatalf("process identity event: %v", err) + } + if got := st.handles["did:plc:abc"]; got != "alice.tangled.org" { + t.Fatalf("handle: got %q", got) + } + if st.syncCursor != "101" { + t.Fatalf("cursor: got %q, want 101", st.syncCursor) + } + if len(tap.acked) != 1 || tap.acked[0] != 101 { + t.Fatalf("acks: got %#v", tap.acked) + } +} + +func TestRunner_ProcessCreateAndDelete(t *testing.T) { + st := newFakeStore() + st.handles["did:plc:author"] = "author.tangled.org" + tap := &fakeTapClient{} + r := newRunnerForTest(st, tap, "sh.tangled.*") + + createEvent := normalize.TapRecordEvent{ + ID: 201, + 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(), createEvent); err != nil { + t.Fatalf("process create event: %v", err) + } + + docID := normalize.StableID("did:plc:author", "sh.tangled.repo", "repo1") + doc := st.docs[docID] + if doc == nil { + t.Fatalf("document %q not found", docID) + } + if doc.AuthorHandle != "author.tangled.org" { + t.Fatalf("author handle: got %q", doc.AuthorHandle) + } + if !st.enqueued[docID] { + t.Fatalf("embedding job not enqueued for %q", docID) + } + + deleteEvent := normalize.TapRecordEvent{ + ID: 202, + Type: "record", + Record: &normalize.TapRecord{ + DID: "did:plc:author", + Collection: "sh.tangled.repo", + RKey: "repo1", + Action: "delete", + }, + } + if err := r.processEvent(context.Background(), deleteEvent); err != nil { + t.Fatalf("process delete event: %v", err) + } + if !st.deleted[docID] { + t.Fatalf("expected tombstone for %q", docID) + } + if st.syncCursor != "202" { + t.Fatalf("cursor: got %q, want 202", st.syncCursor) + } +} + +func TestRunner_ProcessStateEvent(t *testing.T) { + st := newFakeStore() + tap := &fakeTapClient{} + r := newRunnerForTest(st, tap, "sh.tangled.*") + + event := normalize.TapRecordEvent{ + ID: 301, + Type: "record", + Record: &normalize.TapRecord{ + DID: "did:plc:abc", + Collection: "sh.tangled.repo.issue.state", + RKey: "state1", + Action: "create", + Record: map[string]any{ + "subject": "at://did:plc:abc/sh.tangled.repo.issue/1", + "status": "closed", + }, + }, + } + + if err := r.processEvent(context.Background(), event); err != nil { + t.Fatalf("process state event: %v", err) + } + if got := st.recordStates["at://did:plc:abc/sh.tangled.repo.issue/1"]; got != "closed" { + t.Fatalf("record state: got %q", got) + } +} + +func TestRunner_NormalizationFailureAdvancesCursor(t *testing.T) { + st := newFakeStore() + tap := &fakeTapClient{} + r := newRunnerForTest(st, tap, "sh.tangled.*") + + event := normalize.TapRecordEvent{ + ID: 401, + Type: "record", + Record: &normalize.TapRecord{ + DID: "did:plc:abc", + Collection: "sh.tangled.repo.issue", + RKey: "bad-issue", + Action: "create", + CID: "cid-bad", + Record: map[string]any{ + "title": "bad issue", + "repo": "not-an-at-uri", + }, + }, + } + + if err := r.processEvent(context.Background(), event); err != nil { + t.Fatalf("process malformed issue event: %v", err) + } + if st.syncCursor != "401" { + t.Fatalf("cursor: got %q, want 401", st.syncCursor) + } + if len(st.docs) != 0 { + t.Fatalf("expected no documents, got %d", len(st.docs)) + } +} + +func TestAllowlistMatching(t *testing.T) { + a := parseAllowlist("sh.tangled.repo, sh.tangled.string sh.tangled.actor.*") + if !a.match("sh.tangled.repo") { + t.Fatal("expected exact match") + } + if !a.match("sh.tangled.actor.profile") { + t.Fatal("expected wildcard prefix match") + } + if a.match("app.bsky.feed.post") { + t.Fatal("unexpected match") + } +} diff --git a/packages/api/internal/store/sql_store.go b/packages/api/internal/store/sql_store.go --- a/packages/api/internal/store/sql_store.go +++ b/packages/api/internal/store/sql_store.go @@ -130,6 +130,54 @@ return nil } +func (s *SQLStore) UpsertIdentityHandle(ctx context.Context, did, handle string, isActive bool, status string) error { + now := time.Now().UTC().Format(time.RFC3339) + _, err := s.db.ExecContext(ctx, ` + INSERT INTO identity_handles (did, handle, is_active, status, updated_at) + VALUES (?, ?, ?, ?, ?) + ON CONFLICT(did) DO UPDATE SET + handle = excluded.handle, + is_active = excluded.is_active, + status = excluded.status, + updated_at = excluded.updated_at`, + did, handle, isActive, status, now, + ) + if err != nil { + return fmt.Errorf("upsert identity handle: %w", err) + } + return nil +} + +func (s *SQLStore) GetIdentityHandle(ctx context.Context, did string) (string, error) { + var handle sql.NullString + err := s.db.QueryRowContext(ctx, `SELECT handle FROM identity_handles WHERE did = ?`, did).Scan(&handle) + if errors.Is(err, sql.ErrNoRows) { + return "", nil + } + if err != nil { + return "", fmt.Errorf("get identity handle: %w", err) + } + return handle.String, nil +} + +func (s *SQLStore) EnqueueEmbeddingJob(ctx context.Context, documentID string) error { + now := time.Now().UTC().Format(time.RFC3339) + _, err := s.db.ExecContext(ctx, ` + INSERT INTO embedding_jobs (document_id, status, attempts, last_error, scheduled_at, updated_at) + VALUES (?, 'pending', 0, NULL, ?, ?) + ON CONFLICT(document_id) DO UPDATE SET + status = 'pending', + last_error = NULL, + scheduled_at = excluded.scheduled_at, + updated_at = excluded.updated_at`, + documentID, now, now, + ) + if err != nil { + return fmt.Errorf("enqueue embedding job: %w", err) + } + return nil +} + func scanDocument(row *sql.Row) (*Document, error) { doc := &Document{} var ( diff --git a/packages/api/internal/store/store.go b/packages/api/internal/store/store.go --- a/packages/api/internal/store/store.go +++ b/packages/api/internal/store/store.go @@ -48,4 +48,7 @@ GetSyncState(ctx context.Context, consumer string) (*SyncState, error) SetSyncState(ctx context.Context, consumer string, cursor string) error UpdateRecordState(ctx context.Context, subjectURI string, state string) error + UpsertIdentityHandle(ctx context.Context, did, handle string, isActive bool, status string) error + GetIdentityHandle(ctx context.Context, did string) (string, error) + EnqueueEmbeddingJob(ctx context.Context, documentID string) error } diff --git a/packages/api/internal/store/store_test.go b/packages/api/internal/store/store_test.go --- a/packages/api/internal/store/store_test.go +++ b/packages/api/internal/store/store_test.go @@ -2,6 +2,7 @@ import ( "context" + "database/sql" "os" "path/filepath" "testing" @@ -154,6 +155,85 @@ } if err := st.UpdateRecordState(ctx, uri, "closed"); err != nil { t.Fatalf("update record state to closed: %v", err) + } + }) + + t.Run("identity handle upsert and get", func(t *testing.T) { + const did = "did:plc:identity1" + + handle, err := st.GetIdentityHandle(ctx, did) + if err != nil { + t.Fatalf("get missing identity handle: %v", err) + } + if handle != "" { + t.Fatalf("expected empty missing handle, got %q", handle) + } + + if err := st.UpsertIdentityHandle(ctx, did, "alice.tangled.org", true, "active"); err != nil { + t.Fatalf("upsert identity handle: %v", err) + } + + handle, err = st.GetIdentityHandle(ctx, did) + if err != nil { + t.Fatalf("get identity handle: %v", err) + } + if handle != "alice.tangled.org" { + t.Fatalf("handle: got %q, want %q", handle, "alice.tangled.org") + } + + if err := st.UpsertIdentityHandle(ctx, did, "alice2.tangled.org", false, "inactive"); err != nil { + t.Fatalf("update identity handle: %v", err) + } + + handle, err = st.GetIdentityHandle(ctx, did) + if err != nil { + t.Fatalf("get identity handle after update: %v", err) + } + if handle != "alice2.tangled.org" { + t.Fatalf("handle after update: got %q, want %q", handle, "alice2.tangled.org") + } + }) + + t.Run("enqueue embedding job is idempotent", func(t *testing.T) { + doc := &store.Document{ + ID: "did:plc:embed|sh.tangled.string|abc", + DID: "did:plc:embed", + Collection: "sh.tangled.string", + RKey: "abc", + ATURI: "at://did:plc:embed/sh.tangled.string/abc", + CID: "bafyreienqueue", + RecordType: "string", + Title: "foo.go", + Body: "package main", + } + if err := st.UpsertDocument(ctx, doc); err != nil { + t.Fatalf("upsert doc for embedding queue: %v", err) + } + + if err := st.EnqueueEmbeddingJob(ctx, doc.ID); err != nil { + t.Fatalf("enqueue embedding job: %v", err) + } + if err := st.EnqueueEmbeddingJob(ctx, doc.ID); err != nil { + t.Fatalf("enqueue embedding job second call: %v", err) + } + + row := db.QueryRowContext(ctx, `SELECT status, attempts, last_error FROM embedding_jobs WHERE document_id = ?`, doc.ID) + var ( + status string + attempts int + lastError sql.NullString + ) + if err := row.Scan(&status, &attempts, &lastError); err != nil { + t.Fatalf("query embedding job: %v", err) + } + if status != "pending" { + t.Fatalf("status: got %q, want pending", status) + } + if attempts != 0 { + t.Fatalf("attempts: got %d, want 0", attempts) + } + if lastError.Valid { + t.Fatalf("last_error: got %q, want NULL", lastError.String) } }) } diff --git a/packages/api/internal/tapclient/tapclient.go b/packages/api/internal/tapclient/tapclient.go --- a/packages/api/internal/tapclient/tapclient.go +++ b/packages/api/internal/tapclient/tapclient.go @@ -1,1 +1,181 @@ package tapclient + +import ( + "context" + "encoding/base64" + "encoding/json" + "fmt" + "log/slog" + "math/rand/v2" + "net/http" + "strconv" + "sync" + "time" + + "github.com/coder/websocket" + "tangled.org/desertthunder.dev/twister/internal/normalize" +) + +const ( + minReconnectBackoff = 500 * time.Millisecond + maxReconnectBackoff = 10 * time.Second +) + +// Client receives Tap events over WebSocket and sends acks after processing. +type Client struct { + url string + password string + log *slog.Logger + + mu sync.Mutex + conn *websocket.Conn + ackAsJSON bool +} + +func New(url, password string, log *slog.Logger) *Client { + if log == nil { + log = slog.Default() + } + return &Client{ + url: url, + password: password, + log: log, + ackAsJSON: true, + } +} + +func (c *Client) ReadEvent(ctx context.Context) (normalize.TapRecordEvent, error) { + for { + conn, err := c.ensureConnected(ctx) + if err != nil { + return normalize.TapRecordEvent{}, err + } + + _, data, err := conn.Read(ctx) + if err != nil { + c.log.Warn("tap read failed", slog.String("error", err.Error())) + c.resetConn(websocket.StatusInternalError, "read failed") + if ctx.Err() != nil { + return normalize.TapRecordEvent{}, ctx.Err() + } + continue + } + + var event normalize.TapRecordEvent + if err := json.Unmarshal(data, &event); err != nil { + c.log.Warn("tap decode failed", slog.String("error", err.Error())) + continue + } + + return event, nil + } +} + +func (c *Client) AckEvent(ctx context.Context, id int64) error { + conn, err := c.ensureConnected(ctx) + if err != nil { + return err + } + + c.mu.Lock() + ackAsJSON := c.ackAsJSON + c.mu.Unlock() + + if ackAsJSON { + payload, _ := json.Marshal(map[string]int64{"id": id}) + if err := conn.Write(ctx, websocket.MessageText, payload); err == nil { + return nil + } + + c.log.Warn("tap ack json failed; trying plain id", slog.Int64("event_id", id)) + plain := []byte(strconv.FormatInt(id, 10)) + if err := conn.Write(ctx, websocket.MessageText, plain); err != nil { + c.resetConn(websocket.StatusInternalError, "ack failed") + return fmt.Errorf("ack event %d: %w", id, err) + } + + c.mu.Lock() + c.ackAsJSON = false + c.mu.Unlock() + return nil + } + + if err := conn.Write(ctx, websocket.MessageText, []byte(strconv.FormatInt(id, 10))); err != nil { + c.resetConn(websocket.StatusInternalError, "ack failed") + return fmt.Errorf("ack event %d: %w", id, err) + } + return nil +} + +func (c *Client) Close() error { + c.mu.Lock() + defer c.mu.Unlock() + if c.conn == nil { + return nil + } + err := c.conn.Close(websocket.StatusNormalClosure, "shutdown") + c.conn = nil + return err +} + +func (c *Client) ensureConnected(ctx context.Context) (*websocket.Conn, error) { + c.mu.Lock() + if c.conn != nil { + conn := c.conn + c.mu.Unlock() + return conn, nil + } + c.mu.Unlock() + + backoff := minReconnectBackoff + for { + if ctx.Err() != nil { + return nil, ctx.Err() + } + + h := http.Header{} + if c.password != "" { + token := base64.StdEncoding.EncodeToString([]byte("admin:" + c.password)) + h.Set("Authorization", "Basic "+token) + } + + conn, _, err := websocket.Dial(ctx, c.url, &websocket.DialOptions{HTTPHeader: h}) + if err == nil { + c.mu.Lock() + if c.conn == nil { + c.conn = conn + } else { + _ = conn.Close(websocket.StatusNormalClosure, "duplicate") + } + existing := c.conn + c.mu.Unlock() + c.log.Info("tap connected") + return existing, nil + } + + c.log.Warn("tap connect failed", slog.String("error", err.Error()), slog.Duration("retry_in", backoff)) + + jitter := time.Duration(rand.Int64N(int64(backoff / 2))) + wait := backoff + jitter + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-time.After(wait): + } + + backoff *= 2 + if backoff > maxReconnectBackoff { + backoff = maxReconnectBackoff + } + } +} + +func (c *Client) resetConn(status websocket.StatusCode, reason string) { + c.mu.Lock() + defer c.mu.Unlock() + if c.conn == nil { + return + } + _ = c.conn.Close(status, reason) + c.conn = nil +} diff --git a/packages/api/internal/store/migrations/002_identity_handles.sql b/packages/api/internal/store/migrations/002_identity_handles.sql new file mode 100644 --- /dev/null +++ b/packages/api/internal/store/migrations/002_identity_handles.sql @@ -0,0 +1,9 @@ +CREATE TABLE IF NOT EXISTS identity_handles ( + did TEXT PRIMARY KEY, + handle TEXT NOT NULL, + is_active INTEGER NOT NULL DEFAULT 1, + status TEXT, + updated_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_identity_handles_handle ON identity_handles(handle);