diff --git a/docs/api/seeds.txt b/docs/api/seeds.txt new file mode 100644 index 0000000..a8f8488 --- /dev/null +++ b/docs/api/seeds.txt @@ -0,0 +1,9 @@ +# Example seed handles for Twister graph backfill +# One DID or handle per line. Comments and blank lines are ignored. + +anirudh.fi +atprotocol.dev +zzstoatzz.io +oppi.li +desertthunder.dev +tangled.org diff --git a/docs/api/specs/07-graph-backfill.md b/docs/api/specs/07-graph-backfill.md index b88b7c3..f491e3b 100644 --- a/docs/api/specs/07-graph-backfill.md +++ b/docs/api/specs/07-graph-backfill.md @@ -96,7 +96,7 @@ twister backfill --seeds seeds.txt --concurrency 5 | Flag | Default | Description | | --------------- | -------- | ----------------------------------------------- | -| `--seeds` | required | Path to seed file | +| `--seeds` | required | Seed source: file path or comma-separated list | | `--max-hops` | `2` | Max fan-out depth from seed users | | `--dry-run` | `false` | List discovered users without submitting to Tap | | `--concurrency` | `5` | Parallel discovery workers | diff --git a/docs/api/tasks/phase-1-mvp.md b/docs/api/tasks/phase-1-mvp.md index cf91990..aa9b782 100644 --- a/docs/api/tasks/phase-1-mvp.md +++ b/docs/api/tasks/phase-1-mvp.md @@ -134,48 +134,50 @@ Bootstrap the index with historical Tangled content by discovering and backfilli ### Tasks -- [ ] Implement `backfill` subcommand with flags: +- [x] Implement `backfill` subcommand with flags: - `--seeds ` — required seed file path - `--max-hops ` — depth limit for fan-out (default: 2) - `--dry-run` — print the discovery plan without mutating Tap - `--concurrency ` — parallel discovery workers (default: 5) - `--batch-size ` — DIDs per `/repos/add` request - `--batch-delay ` — delay between Tap registration batches -- [ ] Implement seed file parsing: +- [x] Implement seed file parsing: - One DID or handle per line - `#` comments allowed - Blank lines ignored - Handles resolved to DIDs before graph expansion -- [ ] Decide and document the initial seed file location for operators: +- [x] Decide and document the initial seed file location for operators: - Repository-managed example file for format/reference - Deployment-specific runtime file or mounted secret for real runs -- [ ] Implement graph discovery: + - Implemented: `docs/api/seeds.txt` and `packages/api/internal/backfill/doc.go` +- [x] Implement graph discovery: 1. Start from hop-0 seed users 2. Fetch `sh.tangled.graph.follow` records and collect subject DIDs 3. Fetch repo collaborators by inspecting repos, issues, PRs, and comments 4. Enqueue newly discovered DIDs with hop metadata 5. Stop expanding beyond `max-hops` -- [ ] Track discovery metadata for logs: +- [x] Track discovery metadata for logs: - source DID - hop depth - discovery reason (`seed`, `follow`, `collaborator`) -- [ ] Integrate with Tap admin endpoints: +- [x] Integrate with Tap admin endpoints: - `GET /info/:did` to skip already-tracked repos when practical - `POST /repos/add` to register new DIDs for backfill -- [ ] Make the command safe to re-run: +- [x] Make the command safe to re-run: - in-memory visited DID set during crawl - tolerate duplicate `/repos/add` - rely on index upsert idempotency for re-delivered records -- [ ] Add operator-friendly logging: +- [x] Add operator-friendly logging: - seed count - users discovered per hop - already-tracked vs newly-submitted DIDs - batch progress - final totals -- [ ] Add a short runbook covering: +- [x] Add a short runbook covering: - first bootstrap against an empty database - repeat run after expanding the seed list - dry-run before production mutation + - Implemented: `packages/api/internal/backfill/doc.go` ### Verification diff --git a/packages/api/go.mod b/packages/api/go.mod index 245d13e..335d36a 100644 --- a/packages/api/go.mod +++ b/packages/api/go.mod @@ -3,6 +3,8 @@ module tangled.org/desertthunder.dev/twister go 1.25.0 require ( + github.com/coder/websocket v1.8.12 + github.com/joho/godotenv v1.5.1 github.com/spf13/cobra v1.10.2 github.com/tursodatabase/libsql-client-go v0.0.0-20251219100830-236aa1ff8acc modernc.org/sqlite v1.47.0 @@ -10,7 +12,6 @@ require ( require ( github.com/antlr4-go/antlr/v4 v4.13.0 // indirect - github.com/coder/websocket v1.8.12 // indirect github.com/dustin/go-humanize v1.0.1 // indirect github.com/google/uuid v1.6.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect diff --git a/packages/api/go.sum b/packages/api/go.sum index 7ad79c2..166b5e2 100644 --- a/packages/api/go.sum +++ b/packages/api/go.sum @@ -13,6 +13,8 @@ github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= +github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0= +github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4= github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w= diff --git a/packages/api/internal/backfill/backfill.go b/packages/api/internal/backfill/backfill.go new file mode 100644 index 0000000..797f520 --- /dev/null +++ b/packages/api/internal/backfill/backfill.go @@ -0,0 +1,267 @@ +package backfill + +import ( + "context" + "fmt" + "log/slog" + "sort" + "sync" + "time" +) + +type discoveryStore interface { + GetRepoCollaborators(ctx context.Context, repoOwnerDID string) ([]string, error) +} + +// Runner executes seed resolution, graph discovery, and Tap registration. +type Runner struct { + store discoveryStore + tap tapAdmin + resolver handleResolver + follows followFetcher + log *slog.Logger +} + +func NewRunner(store discoveryStore, tap tapAdmin, resolver handleResolver, log *slog.Logger) *Runner { + return NewRunnerWithDeps(store, tap, resolver, NewHTTPFollowFetcher(), log) +} + +func NewRunnerWithDeps(store discoveryStore, tap tapAdmin, resolver handleResolver, follows followFetcher, log *slog.Logger) *Runner { + if log == nil { + log = slog.Default() + } + if follows == nil { + follows = NewHTTPFollowFetcher() + } + return &Runner{store: store, tap: tap, resolver: resolver, follows: follows, log: log} +} + +func (r *Runner) Run(ctx context.Context, opts Options) error { + if opts.SeedsPath == "" { + return fmt.Errorf("--seeds is required") + } + if opts.MaxHops < 0 { + return fmt.Errorf("--max-hops must be >= 0") + } + if opts.Concurrency <= 0 { + opts.Concurrency = 5 + } + if opts.BatchSize <= 0 { + opts.BatchSize = 10 + } + if opts.BatchDelay < 0 { + return fmt.Errorf("--batch-delay must be >= 0") + } + + seedEntries, err := parseSeedInput(opts.SeedsPath) + if err != nil { + return err + } + seeds, err := r.resolveSeeds(ctx, seedEntries) + if err != nil { + return err + } + if len(seeds) == 0 { + return fmt.Errorf("no valid seed DIDs resolved") + } + + r.log.Info("starting backfill discovery", + slog.Int("seed_count", len(seeds)), + slog.Int("max_hops", opts.MaxHops), + slog.Int("concurrency", opts.Concurrency), + ) + + discovered, err := r.discover(ctx, seeds, opts.MaxHops, opts.Concurrency) + if err != nil { + return err + } + + r.log.Info("discovery complete", slog.Int("discovered_total", len(discovered))) + if opts.DryRun { + r.log.Info("dry-run mode enabled; skipping Tap mutations") + return nil + } + + alreadyTracked := 0 + toSubmit := make([]string, 0, len(discovered)) + for _, user := range discovered { + tracked, err := r.tap.IsTracked(ctx, user.DID) + if err != nil { + return fmt.Errorf("tap info for %s: %w", user.DID, err) + } + if tracked { + alreadyTracked++ + continue + } + toSubmit = append(toSubmit, user.DID) + } + + r.log.Info("tap classification complete", + slog.Int("already_tracked", alreadyTracked), + slog.Int("to_submit", len(toSubmit)), + ) + + submitted := 0 + for i := 0; i < len(toSubmit); i += opts.BatchSize { + end := i + opts.BatchSize + if end > len(toSubmit) { + end = len(toSubmit) + } + batch := toSubmit[i:end] + if err := r.tap.AddRepos(ctx, batch); err != nil { + return fmt.Errorf("submit batch %d-%d: %w", i, end, err) + } + submitted += len(batch) + r.log.Info("submitted Tap batch", + slog.Int("batch_start", i), + slog.Int("batch_end", end), + slog.Int("batch_size", len(batch)), + slog.Int("submitted_total", submitted), + ) + if end < len(toSubmit) && opts.BatchDelay > 0 { + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(opts.BatchDelay): + } + } + } + + r.log.Info("backfill complete", + slog.Int("discovered_total", len(discovered)), + slog.Int("already_tracked", alreadyTracked), + slog.Int("submitted", submitted), + ) + return nil +} + +func (r *Runner) resolveSeeds(ctx context.Context, entries []seedEntry) ([]string, error) { + seen := map[string]bool{} + seeds := make([]string, 0, len(entries)) + for _, entry := range entries { + if entry.isDID { + seen[entry.raw] = true + seeds = append(seeds, entry.raw) + continue + } + did, err := r.resolver.Resolve(ctx, entry.raw) + if err != nil { + return nil, fmt.Errorf("resolve handle at line %d (%s): %w", entry.lineNo, entry.raw, err) + } + if seen[did] { + continue + } + seen[did] = true + seeds = append(seeds, did) + } + return seeds, nil +} + +func (r *Runner) discover(ctx context.Context, seeds []string, maxHops int, concurrency int) ([]DiscoveredUser, error) { + visited := map[string]DiscoveredUser{} + ordered := make([]DiscoveredUser, 0) + frontier := make([]DiscoveredUser, 0, len(seeds)) + for _, did := range seeds { + user := DiscoveredUser{DID: did, Hop: 0, Source: did, Reason: "seed"} + visited[did] = user + ordered = append(ordered, user) + frontier = append(frontier, user) + } + + for hop := 0; hop <= maxHops && len(frontier) > 0; hop++ { + r.log.Info("processing discovery hop", slog.Int("hop", hop), slog.Int("users", len(frontier))) + if hop == maxHops { + break + } + + type expansion struct { + node DiscoveredUser + follows []string + collaborators []string + err error + } + + jobs := make(chan DiscoveredUser) + results := make(chan expansion, len(frontier)) + var wg sync.WaitGroup + for i := 0; i < concurrency; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for node := range jobs { + follows, err := r.follows.ListFollowSubjects(ctx, node.DID) + if err != nil { + results <- expansion{node: node, err: fmt.Errorf("follows: %w", err)} + continue + } + collaborators, err := r.store.GetRepoCollaborators(ctx, node.DID) + if err != nil { + results <- expansion{node: node, err: fmt.Errorf("collaborators: %w", err)} + continue + } + results <- expansion{node: node, follows: follows, collaborators: collaborators} + } + }() + } + + go func() { + for _, node := range frontier { + jobs <- node + } + close(jobs) + wg.Wait() + close(results) + }() + + nextByDID := map[string]DiscoveredUser{} + for res := range results { + if res.err != nil { + r.log.Warn("discovery expansion failed", slog.String("did", res.node.DID), slog.String("error", res.err.Error())) + continue + } + + r.log.Info("discovery expansion", + slog.String("did", res.node.DID), + slog.Int("hop", hop), + slog.Int("follows", len(res.follows)), + slog.Int("collaborators", len(res.collaborators)), + ) + + for _, did := range res.follows { + if !isDID(did) { + continue + } + if _, exists := visited[did]; exists { + continue + } + if _, exists := nextByDID[did]; exists { + continue + } + nextByDID[did] = DiscoveredUser{DID: did, Hop: hop + 1, Source: res.node.DID, Reason: "follow"} + } + for _, did := range res.collaborators { + if !isDID(did) { + continue + } + if _, exists := visited[did]; exists { + continue + } + if _, exists := nextByDID[did]; exists { + continue + } + nextByDID[did] = DiscoveredUser{DID: did, Hop: hop + 1, Source: res.node.DID, Reason: "collaborator"} + } + } + + next := make([]DiscoveredUser, 0, len(nextByDID)) + for _, user := range nextByDID { + visited[user.DID] = user + next = append(next, user) + ordered = append(ordered, user) + } + sort.Slice(next, func(i, j int) bool { return next[i].DID < next[j].DID }) + frontier = next + } + + return ordered, nil +} diff --git a/packages/api/internal/backfill/backfill_test.go b/packages/api/internal/backfill/backfill_test.go new file mode 100644 index 0000000..24048dc --- /dev/null +++ b/packages/api/internal/backfill/backfill_test.go @@ -0,0 +1,118 @@ +package backfill + +import ( + "context" + "io" + "log/slog" + "os" + "path/filepath" + "testing" +) + +type fakeStore struct { + collaborators map[string][]string +} + +func (f *fakeStore) GetRepoCollaborators(_ context.Context, did string) ([]string, error) { + return f.collaborators[did], nil +} + +type fakeFollowFetcher struct { + follows map[string][]string +} + +func (f *fakeFollowFetcher) ListFollowSubjects(_ context.Context, did string) ([]string, error) { + return f.follows[did], nil +} + +type fakeTapAdmin struct { + tracked map[string]bool + added [][]string +} + +func (f *fakeTapAdmin) IsTracked(_ context.Context, did string) (bool, error) { + return f.tracked[did], nil +} + +func (f *fakeTapAdmin) AddRepos(_ context.Context, dids []string) error { + batch := make([]string, len(dids)) + copy(batch, dids) + f.added = append(f.added, batch) + return nil +} + +type fakeResolver struct { + mapping map[string]string +} + +func (r *fakeResolver) Resolve(_ context.Context, handle string) (string, error) { + if did, ok := r.mapping[handle]; ok { + return did, nil + } + return "", io.EOF +} + +func TestRunner_DiscoveryAndSubmit(t *testing.T) { + st := &fakeStore{ + collaborators: map[string][]string{ + "did:plc:seed": {"did:plc:c1"}, + }, + } + follows := &fakeFollowFetcher{follows: map[string][]string{"did:plc:seed": {"did:plc:f1"}}} + tap := &fakeTapAdmin{tracked: map[string]bool{"did:plc:f1": 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: 1, + Concurrency: 2, + BatchSize: 2, + }) + if err != nil { + t.Fatalf("run backfill: %v", err) + } + + if len(tap.added) != 1 { + t.Fatalf("expected one batch, got %d", len(tap.added)) + } + if len(tap.added[0]) != 2 { + t.Fatalf("expected 2 dids submitted, got %#v", tap.added[0]) + } +} + +func TestRunner_DryRunSkipsMutations(t *testing.T) { + st := &fakeStore{collaborators: map[string][]string{}} + follows := &fakeFollowFetcher{follows: map[string][]string{}} + tap := &fakeTapAdmin{tracked: map[string]bool{}} + 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, + DryRun: true, + Concurrency: 1, + BatchSize: 10, + }) + if err != nil { + t.Fatalf("run dry-run backfill: %v", err) + } + if len(tap.added) != 0 { + t.Fatalf("expected no tap submissions in dry-run, got %#v", tap.added) + } +} diff --git a/packages/api/internal/backfill/doc.go b/packages/api/internal/backfill/doc.go new file mode 100644 index 0000000..2f925cf --- /dev/null +++ b/packages/api/internal/backfill/doc.go @@ -0,0 +1,69 @@ +// Package backfill provides graph bootstrap tooling for Twister. +// +// # Backfill Runbook +// +// This runbook covers initial graph bootstrap and repeat runs using: +// +// twister backfill +// +// # Seeds Input +// +// The `--seeds` flag supports either of these forms: +// +// 1. File path: +// +// twister backfill --seeds /etc/twister/seeds.txt +// +// 2. Comma-separated inline list: +// +// twister backfill --seeds anirudh.fi,atprotocol.dev,oppi.li +// +// Supported seed entries are DIDs and handles. +// +// Repository-managed example seed file: +// +// docs/api/seeds.txt +// +// Runtime seed file is typically mounted outside the repo, for example: +// +// /etc/twister/seeds.txt +// +// # Prerequisites +// +// Required environment variables: +// +// - TURSO_DATABASE_URL +// - TURSO_AUTH_TOKEN (for non-file Turso URLs) +// - TAP_URL +// - TAP_AUTH_PASSWORD +// +// # First Bootstrap +// +// 1. Copy and customize seeds: +// +// cp docs/api/seeds.txt /tmp/twister-seeds.txt +// +// 2. Run dry-run first: +// +// twister backfill --seeds /tmp/twister-seeds.txt --max-hops 2 --dry-run +// +// 3. Run real backfill: +// +// twister backfill --seeds /tmp/twister-seeds.txt --max-hops 2 --concurrency 5 --batch-size 10 --batch-delay 1s +// +// Watch logs for seed count, hop-level discoveries, already-tracked vs submitted +// users, and batch progress totals. +// +// # Repeat Run +// +// Append new candidate users to the seed source, run dry-run, then run the real +// command again. Reruns are safe because discovery deduplicates in-memory and +// Tap /repos/add is treated as idempotent. +// +// # Dry-Run Safety +// +// Before production mutation: +// - include --dry-run +// - confirm only discovery output appears +// - confirm no Tap mutation side effects +package backfill diff --git a/packages/api/internal/backfill/follows.go b/packages/api/internal/backfill/follows.go new file mode 100644 index 0000000..2812753 --- /dev/null +++ b/packages/api/internal/backfill/follows.go @@ -0,0 +1,145 @@ +package backfill + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/url" + "strings" + "time" +) + +const ( + plcDirectoryBase = "https://plc.directory" + followCollection = "sh.tangled.graph.follow" +) + +type followFetcher interface { + ListFollowSubjects(ctx context.Context, did string) ([]string, error) +} + +// HTTPFollowFetcher resolves a DID's PDS endpoint and reads follow records +// directly from com.atproto.repo.listRecords. +type HTTPFollowFetcher struct { + client *http.Client +} + +func NewHTTPFollowFetcher() *HTTPFollowFetcher { + return &HTTPFollowFetcher{ + client: &http.Client{Timeout: 15 * time.Second}, + } +} + +func (f *HTTPFollowFetcher) ListFollowSubjects(ctx context.Context, did string) ([]string, error) { + pdsEndpoint, err := f.resolvePDSEndpoint(ctx, did) + if err != nil { + return nil, err + } + + seen := map[string]bool{} + var subjects []string + cursor := "" + + for { + u, err := url.Parse(strings.TrimSuffix(pdsEndpoint, "/") + "/xrpc/com.atproto.repo.listRecords") + if err != nil { + return nil, fmt.Errorf("build listRecords url: %w", err) + } + q := u.Query() + q.Set("repo", did) + q.Set("collection", followCollection) + q.Set("limit", "100") + if cursor != "" { + q.Set("cursor", cursor) + } + u.RawQuery = q.Encode() + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, u.String(), nil) + if err != nil { + return nil, fmt.Errorf("build listRecords request: %w", err) + } + + resp, err := f.client.Do(req) + if err != nil { + return nil, fmt.Errorf("listRecords request: %w", err) + } + + var payload struct { + Cursor string `json:"cursor"` + Records []struct { + Value map[string]any `json:"value"` + } `json:"records"` + } + if resp.StatusCode != http.StatusOK { + _ = resp.Body.Close() + return nil, fmt.Errorf("listRecords failed: status %d", resp.StatusCode) + } + if err := json.NewDecoder(resp.Body).Decode(&payload); err != nil { + _ = resp.Body.Close() + return nil, fmt.Errorf("decode listRecords response: %w", err) + } + _ = resp.Body.Close() + + for _, rec := range payload.Records { + subject, _ := rec.Value["subject"].(string) + if !isDID(subject) || seen[subject] { + continue + } + seen[subject] = true + subjects = append(subjects, subject) + } + + if payload.Cursor == "" || payload.Cursor == cursor { + break + } + cursor = payload.Cursor + } + + return subjects, nil +} + +func (f *HTTPFollowFetcher) resolvePDSEndpoint(ctx context.Context, did string) (string, error) { + var didDocURL string + switch { + case strings.HasPrefix(did, "did:plc:"): + didDocURL = plcDirectoryBase + "/" + url.PathEscape(did) + case strings.HasPrefix(did, "did:web:"): + hostAndPath := strings.TrimPrefix(did, "did:web:") + hostAndPath = strings.ReplaceAll(hostAndPath, ":", "/") + didDocURL = "https://" + hostAndPath + "/.well-known/did.json" + default: + return "", fmt.Errorf("unsupported did type for pds resolution: %s", did) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, didDocURL, nil) + if err != nil { + return "", fmt.Errorf("build did doc request: %w", err) + } + resp, err := f.client.Do(req) + if err != nil { + return "", fmt.Errorf("did doc request: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("did doc lookup failed: status %d", resp.StatusCode) + } + + var didDoc struct { + Service []struct { + Type string `json:"type"` + ServiceEndpoint string `json:"serviceEndpoint"` + } `json:"service"` + } + if err := json.NewDecoder(resp.Body).Decode(&didDoc); err != nil { + return "", fmt.Errorf("decode did doc: %w", err) + } + + for _, service := range didDoc.Service { + if service.Type == "AtprotoPersonalDataServer" && strings.TrimSpace(service.ServiceEndpoint) != "" { + return strings.TrimSpace(service.ServiceEndpoint), nil + } + } + + return "", fmt.Errorf("no atproto pds endpoint in did document") +} diff --git a/packages/api/internal/backfill/resolve.go b/packages/api/internal/backfill/resolve.go new file mode 100644 index 0000000..15b9d33 --- /dev/null +++ b/packages/api/internal/backfill/resolve.go @@ -0,0 +1,63 @@ +package backfill + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/url" + "time" +) + +const defaultIdentityService = "https://public.api.bsky.app" + +type handleResolver interface { + Resolve(ctx context.Context, handle string) (string, error) +} + +// HTTPHandleResolver resolves handles through com.atproto.identity.resolveHandle. +type HTTPHandleResolver struct { + baseURL string + client *http.Client +} + +func NewHTTPHandleResolver(baseURL string) *HTTPHandleResolver { + if baseURL == "" { + baseURL = defaultIdentityService + } + return &HTTPHandleResolver{ + baseURL: baseURL, + client: &http.Client{ + Timeout: 10 * time.Second, + }, + } +} + +func (r *HTTPHandleResolver) Resolve(ctx context.Context, handle string) (string, error) { + u := r.baseURL + "/xrpc/com.atproto.identity.resolveHandle?handle=" + url.QueryEscape(handle) + req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) + if err != nil { + return "", fmt.Errorf("build resolve handle request: %w", err) + } + + resp, err := r.client.Do(req) + if err != nil { + return "", fmt.Errorf("resolve handle request: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("resolve handle failed: status %d", resp.StatusCode) + } + + var payload struct { + DID string `json:"did"` + } + if err := json.NewDecoder(resp.Body).Decode(&payload); err != nil { + return "", fmt.Errorf("decode resolve handle response: %w", err) + } + if !isDID(payload.DID) { + return "", fmt.Errorf("resolve handle returned invalid did %q", payload.DID) + } + return payload.DID, nil +} diff --git a/packages/api/internal/backfill/seed.go b/packages/api/internal/backfill/seed.go new file mode 100644 index 0000000..16357e3 --- /dev/null +++ b/packages/api/internal/backfill/seed.go @@ -0,0 +1,98 @@ +package backfill + +import ( + "bufio" + "fmt" + "os" + "strings" +) + +type seedEntry struct { + raw string + isDID bool + lineNo int +} + +func parseSeedFile(path string) ([]seedEntry, error) { + f, err := os.Open(path) + if err != nil { + return nil, fmt.Errorf("open seed file: %w", err) + } + defer f.Close() + + seen := map[string]bool{} + entries := make([]seedEntry, 0) + s := bufio.NewScanner(f) + lineNo := 0 + for s.Scan() { + lineNo++ + line := strings.TrimSpace(s.Text()) + if line == "" || strings.HasPrefix(line, "#") { + continue + } + if seen[line] { + continue + } + seen[line] = true + entries = append(entries, seedEntry{ + raw: line, + isDID: isDID(line), + lineNo: lineNo, + }) + } + if err := s.Err(); err != nil { + return nil, fmt.Errorf("scan seed file: %w", err) + } + if len(entries) == 0 { + return nil, fmt.Errorf("seed file has no valid entries") + } + return entries, nil +} + +func parseSeedInput(input string) ([]seedEntry, error) { + input = strings.TrimSpace(input) + if input == "" { + return nil, fmt.Errorf("seeds input is required") + } + + if strings.Contains(input, ",") { + return parseSeedList(input) + } + + info, err := os.Stat(input) + if err == nil && !info.IsDir() { + return parseSeedFile(input) + } + + // Single inline DID/handle is supported for convenience. + return parseSeedList(input) +} + +func parseSeedList(list string) ([]seedEntry, error) { + parts := strings.Split(list, ",") + seen := map[string]bool{} + entries := make([]seedEntry, 0, len(parts)) + for i, part := range parts { + value := strings.TrimSpace(part) + if value == "" { + continue + } + if seen[value] { + continue + } + seen[value] = true + entries = append(entries, seedEntry{ + raw: value, + isDID: isDID(value), + lineNo: i + 1, + }) + } + if len(entries) == 0 { + return nil, fmt.Errorf("seed list has no valid entries") + } + return entries, nil +} + +func isDID(v string) bool { + return strings.HasPrefix(v, "did:") +} diff --git a/packages/api/internal/backfill/seed_test.go b/packages/api/internal/backfill/seed_test.go new file mode 100644 index 0000000..7b8d00c --- /dev/null +++ b/packages/api/internal/backfill/seed_test.go @@ -0,0 +1,73 @@ +package backfill + +import ( + "os" + "path/filepath" + "testing" +) + +func TestParseSeedFile(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "seeds.txt") + content := "\n# comment\ndid:plc:one\nalice.tangled.sh\ndid:plc:one\n\n" + if err := os.WriteFile(path, []byte(content), 0o644); err != nil { + t.Fatalf("write seed file: %v", err) + } + + entries, err := parseSeedFile(path) + if err != nil { + t.Fatalf("parse seed file: %v", err) + } + if len(entries) != 2 { + t.Fatalf("entries length: got %d want 2", len(entries)) + } + if !entries[0].isDID { + t.Fatalf("first entry should be did") + } + if entries[1].isDID { + t.Fatalf("second entry should be handle") + } +} + +func TestParseSeedInput_FilePath(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "seeds.txt") + content := "did:plc:one\nhandle.example\n" + if err := os.WriteFile(path, []byte(content), 0o644); err != nil { + t.Fatalf("write seed file: %v", err) + } + + entries, err := parseSeedInput(path) + if err != nil { + t.Fatalf("parse seed input from file: %v", err) + } + if len(entries) != 2 { + t.Fatalf("entries length: got %d want 2", len(entries)) + } +} + +func TestParseSeedInput_CommaSeparated(t *testing.T) { + entries, err := parseSeedInput("anirudh.fi, atprotocol.dev, did:plc:abc") + if err != nil { + t.Fatalf("parse comma-separated seeds: %v", err) + } + if len(entries) != 3 { + t.Fatalf("entries length: got %d want 3", len(entries)) + } + if !entries[2].isDID { + t.Fatalf("expected third entry to be did") + } +} + +func TestParseSeedInput_SingleInline(t *testing.T) { + entries, err := parseSeedInput("tangled.org") + if err != nil { + t.Fatalf("parse single inline seed: %v", err) + } + if len(entries) != 1 { + t.Fatalf("entries length: got %d want 1", len(entries)) + } + if entries[0].raw != "tangled.org" { + t.Fatalf("entry: got %q", entries[0].raw) + } +} diff --git a/packages/api/internal/backfill/tap_admin.go b/packages/api/internal/backfill/tap_admin.go new file mode 100644 index 0000000..0e3fec2 --- /dev/null +++ b/packages/api/internal/backfill/tap_admin.go @@ -0,0 +1,128 @@ +package backfill + +import ( + "bytes" + "context" + "encoding/base64" + "encoding/json" + "fmt" + "net/http" + "net/url" + "strings" + "time" +) + +type tapAdmin interface { + IsTracked(ctx context.Context, did string) (bool, error) + AddRepos(ctx context.Context, dids []string) error +} + +// HTTPTapAdmin calls Tap admin endpoints for backfill orchestration. +type HTTPTapAdmin struct { + baseURL string + password string + client *http.Client +} + +func NewHTTPTapAdmin(tapURL, password string) (*HTTPTapAdmin, error) { + baseURL, err := normalizeTapBaseURL(tapURL) + if err != nil { + return nil, err + } + return &HTTPTapAdmin{ + baseURL: baseURL, + password: password, + client: &http.Client{ + Timeout: 15 * time.Second, + }, + }, nil +} + +func (t *HTTPTapAdmin) IsTracked(ctx context.Context, did string) (bool, 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) + } + t.addAuth(req) + + resp, err := t.client.Do(req) + if err != nil { + return false, fmt.Errorf("tap info request: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode == http.StatusNotFound { + return false, nil + } + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return false, fmt.Errorf("tap info request failed: status %d", resp.StatusCode) + } + return true, nil +} + +func (t *HTTPTapAdmin) AddRepos(ctx context.Context, dids []string) error { + if len(dids) == 0 { + return nil + } + + payload, err := json.Marshal(map[string][]string{"dids": dids}) + if err != nil { + return fmt.Errorf("marshal repos add payload: %w", err) + } + + endpoint := t.baseURL + "/repos/add" + req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(payload)) + if err != nil { + return fmt.Errorf("build repos add request: %w", err) + } + req.Header.Set("Content-Type", "application/json") + t.addAuth(req) + + resp, err := t.client.Do(req) + if err != nil { + return fmt.Errorf("repos add request: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return fmt.Errorf("repos add failed: status %d", resp.StatusCode) + } + return nil +} + +func (t *HTTPTapAdmin) addAuth(req *http.Request) { + if t.password == "" { + return + } + token := base64.StdEncoding.EncodeToString([]byte("admin:" + t.password)) + req.Header.Set("Authorization", "Basic "+token) +} + +func normalizeTapBaseURL(raw string) (string, error) { + raw = strings.TrimSpace(raw) + if raw == "" { + return "", fmt.Errorf("tap url is required") + } + u, err := url.Parse(raw) + if err != nil { + return "", fmt.Errorf("parse tap url: %w", err) + } + switch u.Scheme { + case "ws": + u.Scheme = "http" + case "wss": + u.Scheme = "https" + case "http", "https": + default: + return "", fmt.Errorf("unsupported tap url scheme %q", u.Scheme) + } + + u.RawQuery = "" + u.Fragment = "" + u.Path = strings.TrimSuffix(u.Path, "/") + u.Path = strings.TrimSuffix(u.Path, "/channel") + if u.Path == "/" { + u.Path = "" + } + return strings.TrimSuffix(u.String(), "/"), nil +} diff --git a/packages/api/internal/backfill/types.go b/packages/api/internal/backfill/types.go new file mode 100644 index 0000000..4d56350 --- /dev/null +++ b/packages/api/internal/backfill/types.go @@ -0,0 +1,21 @@ +package backfill + +import "time" + +// Options configures a backfill run. +type Options struct { + SeedsPath string + MaxHops int + DryRun bool + Concurrency int + BatchSize int + BatchDelay time.Duration +} + +// DiscoveredUser contains crawl metadata for an included DID. +type DiscoveredUser struct { + DID string + Hop int + Source string + Reason string +} diff --git a/packages/api/internal/config/config.go b/packages/api/internal/config/config.go index 9e81fe8..cb9155c 100644 --- a/packages/api/internal/config/config.go +++ b/packages/api/internal/config/config.go @@ -3,8 +3,11 @@ package config import ( "errors" "os" + "path/filepath" "strconv" "strings" + + "github.com/joho/godotenv" ) type Config struct { @@ -32,6 +35,8 @@ type Config struct { } func Load() (*Config, error) { + loadDotEnv() + cfg := &Config{ TursoURL: os.Getenv("TURSO_DATABASE_URL"), TursoToken: os.Getenv("TURSO_AUTH_TOKEN"), @@ -69,6 +74,33 @@ func Load() (*Config, error) { return cfg, nil } +func loadDotEnv() { + seen := map[string]bool{} + candidates := make([]string, 0, 8) + + if explicit := strings.TrimSpace(os.Getenv("TWISTER_ENV_FILE")); explicit != "" { + candidates = append(candidates, explicit) + } + + if cwd, err := os.Getwd(); err == nil { + for _, rel := range []string{".env", "../.env", "../../.env"} { + candidates = append(candidates, filepath.Join(cwd, rel)) + } + } + + for _, candidate := range candidates { + if candidate == "" || seen[candidate] { + continue + } + seen[candidate] = true + if _, err := os.Stat(candidate); err != nil { + continue + } + // Load does not override existing process env vars. + _ = godotenv.Load(candidate) + } +} + func envOrDefault(key, def string) string { if v := os.Getenv(key); v != "" { return v diff --git a/packages/api/internal/ingest/ingest.go b/packages/api/internal/ingest/ingest.go index 7191976..4273b2a 100644 --- a/packages/api/internal/ingest/ingest.go +++ b/packages/api/internal/ingest/ingest.go @@ -6,6 +6,7 @@ import ( "log/slog" "math" "strings" + "sync" "time" "tangled.org/desertthunder.dev/twister/internal/normalize" @@ -15,6 +16,7 @@ import ( const ( defaultConsumerName = "indexer-tap-v1" maxDBRetryBackoff = 5 * time.Second + statusLogInterval = 30 * time.Second ) type client interface { @@ -31,6 +33,10 @@ type Runner struct { allowlist allowlist consumerName string log *slog.Logger + + statusMu sync.Mutex + lastCursor string + processedTick int64 } func NewRunner(st store.Store, registry *normalize.Registry, tap client, indexedCollections string, log *slog.Logger) *Runner { @@ -50,6 +56,8 @@ func NewRunner(st store.Store, registry *normalize.Registry, tap client, indexed func (r *Runner) Run(ctx context.Context) error { defer r.tap.Close() + go r.runStatusLogger(ctx) + for { if ctx.Err() != nil { return nil @@ -225,9 +233,46 @@ func (r *Runner) advanceCursorAndAck(ctx context.Context, eventID int64) error { if err := r.tap.AckEvent(ctx, eventID); err != nil { return err } + r.markProcessed(cursor) return nil } +func (r *Runner) runStatusLogger(ctx context.Context) { + ticker := time.NewTicker(statusLogInterval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + r.statusMu.Lock() + cursor := r.lastCursor + processed := r.processedTick + r.processedTick = 0 + r.statusMu.Unlock() + + docs, err := r.store.CountDocuments(ctx) + if err != nil { + r.log.Warn("indexer status failed", slog.String("error", err.Error())) + continue + } + r.log.Info("indexer status", + slog.String("cursor", cursor), + slog.Int64("events_processed", processed), + slog.Int64("documents", docs), + ) + } + } +} + +func (r *Runner) markProcessed(cursor string) { + r.statusMu.Lock() + r.lastCursor = cursor + r.processedTick++ + r.statusMu.Unlock() +} + type allowlist struct { entries []string } diff --git a/packages/api/internal/ingest/ingest_test.go b/packages/api/internal/ingest/ingest_test.go index 4f83942..cb43765 100644 --- a/packages/api/internal/ingest/ingest_test.go +++ b/packages/api/internal/ingest/ingest_test.go @@ -87,6 +87,18 @@ func (f *fakeStore) EnqueueEmbeddingJob(_ context.Context, documentID string) er return nil } +func (f *fakeStore) GetFollowSubjects(_ context.Context, _ string) ([]string, error) { + return nil, nil +} + +func (f *fakeStore) GetRepoCollaborators(_ context.Context, _ string) ([]string, error) { + return nil, nil +} + +func (f *fakeStore) CountDocuments(_ context.Context) (int64, error) { + return int64(len(f.docs)), 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) diff --git a/packages/api/internal/normalize/follow.go b/packages/api/internal/normalize/follow.go new file mode 100644 index 0000000..ad68469 --- /dev/null +++ b/packages/api/internal/normalize/follow.go @@ -0,0 +1,33 @@ +package normalize + +import "tangled.org/desertthunder.dev/twister/internal/store" + +const collectionFollow = "sh.tangled.graph.follow" + +// FollowAdapter normalizes follow edges for graph-backfill discovery. +type FollowAdapter struct{} + +func (a *FollowAdapter) Collection() string { return collectionFollow } +func (a *FollowAdapter) RecordType() string { return "follow" } + +func (a *FollowAdapter) Searchable(_ map[string]any) bool { return false } + +func (a *FollowAdapter) Normalize(event TapRecordEvent) (*store.Document, error) { + r := event.Record + rec := r.Record + subject := str(rec, "subject") + + return &store.Document{ + ID: StableID(r.DID, r.Collection, r.RKey), + DID: r.DID, + Collection: r.Collection, + RKey: r.RKey, + ATURI: BuildATURI(r.DID, r.Collection, r.RKey), + CID: r.CID, + RecordType: a.RecordType(), + Title: subject, + RepoDID: subject, + TagsJSON: "[]", + CreatedAt: str(rec, "createdAt"), + }, nil +} diff --git a/packages/api/internal/normalize/issue_comment.go b/packages/api/internal/normalize/issue_comment.go new file mode 100644 index 0000000..a1d1d32 --- /dev/null +++ b/packages/api/internal/normalize/issue_comment.go @@ -0,0 +1,50 @@ +package normalize + +import ( + "fmt" + + "tangled.org/desertthunder.dev/twister/internal/store" +) + +const collectionIssueComment = "sh.tangled.repo.issue.comment" + +// IssueCommentAdapter normalizes issue comments for collaborator discovery. +type IssueCommentAdapter struct{} + +func (a *IssueCommentAdapter) Collection() string { return collectionIssueComment } +func (a *IssueCommentAdapter) RecordType() string { return "issue_comment" } + +func (a *IssueCommentAdapter) Searchable(record map[string]any) bool { + return str(record, "body") != "" +} + +func (a *IssueCommentAdapter) Normalize(event TapRecordEvent) (*store.Document, error) { + r := event.Record + rec := r.Record + body := str(rec, "body") + issueURI := str(rec, "issue") + + repoDID := "" + if issueURI != "" { + did, _, _, err := ParseATURI(issueURI) + if err != nil { + return nil, fmt.Errorf("issue comment issue AT-URI: %w", err) + } + repoDID = did + } + + return &store.Document{ + ID: StableID(r.DID, r.Collection, r.RKey), + DID: r.DID, + Collection: r.Collection, + RKey: r.RKey, + ATURI: BuildATURI(r.DID, r.Collection, r.RKey), + CID: r.CID, + RecordType: a.RecordType(), + Body: body, + Summary: truncate(body, 200), + RepoDID: repoDID, + TagsJSON: "[]", + CreatedAt: str(rec, "createdAt"), + }, nil +} diff --git a/packages/api/internal/normalize/normalize_test.go b/packages/api/internal/normalize/normalize_test.go index 2ec2dce..ef91f74 100644 --- a/packages/api/internal/normalize/normalize_test.go +++ b/packages/api/internal/normalize/normalize_test.go @@ -276,6 +276,66 @@ func TestProfileAdapter(t *testing.T) { } } +func TestFollowAdapter(t *testing.T) { + event := loadFixture(t, "follow.json") + adapter := &normalize.FollowAdapter{} + + doc, err := adapter.Normalize(event) + if err != nil { + t.Fatalf("Normalize: %v", err) + } + + if doc.RecordType != "follow" { + t.Errorf("RecordType = %q", doc.RecordType) + } + if doc.RepoDID != "did:plc:bob" { + t.Errorf("RepoDID = %q, want did:plc:bob", doc.RepoDID) + } + if adapter.Searchable(event.Record.Record) { + t.Error("Searchable = true, want false") + } +} + +func TestIssueCommentAdapter(t *testing.T) { + event := loadFixture(t, "issue_comment.json") + adapter := &normalize.IssueCommentAdapter{} + + doc, err := adapter.Normalize(event) + if err != nil { + t.Fatalf("Normalize: %v", err) + } + + if doc.RecordType != "issue_comment" { + t.Errorf("RecordType = %q", doc.RecordType) + } + if doc.RepoDID != "did:plc:repoowner" { + t.Errorf("RepoDID = %q, want did:plc:repoowner", doc.RepoDID) + } + if !adapter.Searchable(event.Record.Record) { + t.Error("Searchable = false for non-empty comment body") + } +} + +func TestPullCommentAdapter(t *testing.T) { + event := loadFixture(t, "pull_comment.json") + adapter := &normalize.PullCommentAdapter{} + + doc, err := adapter.Normalize(event) + if err != nil { + t.Fatalf("Normalize: %v", err) + } + + if doc.RecordType != "pull_comment" { + t.Errorf("RecordType = %q", doc.RecordType) + } + if doc.RepoDID != "did:plc:repoowner" { + t.Errorf("RepoDID = %q, want did:plc:repoowner", doc.RepoDID) + } + if !adapter.Searchable(event.Record.Record) { + t.Error("Searchable = false for non-empty comment body") + } +} + // TestIssueStateHandler verifies record_state extraction. func TestIssueStateHandler(t *testing.T) { event := loadFixture(t, "issue_state.json") @@ -352,6 +412,9 @@ func TestRegistry(t *testing.T) { "sh.tangled.repo", "sh.tangled.repo.issue", "sh.tangled.repo.pull", + "sh.tangled.repo.issue.comment", + "sh.tangled.repo.pull.comment", + "sh.tangled.graph.follow", "sh.tangled.string", "sh.tangled.actor.profile", } diff --git a/packages/api/internal/normalize/pull_comment.go b/packages/api/internal/normalize/pull_comment.go new file mode 100644 index 0000000..1363abf --- /dev/null +++ b/packages/api/internal/normalize/pull_comment.go @@ -0,0 +1,50 @@ +package normalize + +import ( + "fmt" + + "tangled.org/desertthunder.dev/twister/internal/store" +) + +const collectionPullComment = "sh.tangled.repo.pull.comment" + +// PullCommentAdapter normalizes pull comments for collaborator discovery. +type PullCommentAdapter struct{} + +func (a *PullCommentAdapter) Collection() string { return collectionPullComment } +func (a *PullCommentAdapter) RecordType() string { return "pull_comment" } + +func (a *PullCommentAdapter) Searchable(record map[string]any) bool { + return str(record, "body") != "" +} + +func (a *PullCommentAdapter) Normalize(event TapRecordEvent) (*store.Document, error) { + r := event.Record + rec := r.Record + body := str(rec, "body") + pullURI := str(rec, "pull") + + repoDID := "" + if pullURI != "" { + did, _, _, err := ParseATURI(pullURI) + if err != nil { + return nil, fmt.Errorf("pull comment pull AT-URI: %w", err) + } + repoDID = did + } + + return &store.Document{ + ID: StableID(r.DID, r.Collection, r.RKey), + DID: r.DID, + Collection: r.Collection, + RKey: r.RKey, + ATURI: BuildATURI(r.DID, r.Collection, r.RKey), + CID: r.CID, + RecordType: a.RecordType(), + Body: body, + Summary: truncate(body, 200), + RepoDID: repoDID, + TagsJSON: "[]", + CreatedAt: str(rec, "createdAt"), + }, nil +} diff --git a/packages/api/internal/normalize/registry.go b/packages/api/internal/normalize/registry.go index a626178..21ffe99 100644 --- a/packages/api/internal/normalize/registry.go +++ b/packages/api/internal/normalize/registry.go @@ -16,7 +16,10 @@ func NewRegistry() *Registry { &RepoAdapter{}, &IssueAdapter{}, &PullAdapter{}, + &IssueCommentAdapter{}, + &PullCommentAdapter{}, &StringAdapter{}, + &FollowAdapter{}, &ProfileAdapter{}, } { r.adapters[a.Collection()] = a diff --git a/packages/api/internal/normalize/testdata/follow.json b/packages/api/internal/normalize/testdata/follow.json new file mode 100644 index 0000000..b373451 --- /dev/null +++ b/packages/api/internal/normalize/testdata/follow.json @@ -0,0 +1,14 @@ +{ + "id": 2001, + "type": "record", + "record": { + "live": true, + "rev": "3kb3follow1", + "did": "did:plc:alice", + "collection": "sh.tangled.graph.follow", + "rkey": "f1", + "action": "create", + "cid": "bafyreifollow", + "record": { "$type": "sh.tangled.graph.follow", "subject": "did:plc:bob", "createdAt": "2026-03-20T12:00:00.000Z" } + } +} diff --git a/packages/api/internal/normalize/testdata/issue_comment.json b/packages/api/internal/normalize/testdata/issue_comment.json new file mode 100644 index 0000000..0a18d8f --- /dev/null +++ b/packages/api/internal/normalize/testdata/issue_comment.json @@ -0,0 +1,19 @@ +{ + "id": 2002, + "type": "record", + "record": { + "live": true, + "rev": "3kb3ic1", + "did": "did:plc:commenter", + "collection": "sh.tangled.repo.issue.comment", + "rkey": "ic1", + "action": "create", + "cid": "bafyreic1", + "record": { + "$type": "sh.tangled.repo.issue.comment", + "issue": "at://did:plc:repoowner/sh.tangled.repo.issue/issue1", + "body": "I can help with this fix.", + "createdAt": "2026-03-20T13:00:00.000Z" + } + } +} diff --git a/packages/api/internal/normalize/testdata/pull_comment.json b/packages/api/internal/normalize/testdata/pull_comment.json new file mode 100644 index 0000000..8dedd82 --- /dev/null +++ b/packages/api/internal/normalize/testdata/pull_comment.json @@ -0,0 +1,19 @@ +{ + "id": 2003, + "type": "record", + "record": { + "live": true, + "rev": "3kb3pc1", + "did": "did:plc:reviewer", + "collection": "sh.tangled.repo.pull.comment", + "rkey": "pc1", + "action": "create", + "cid": "bafyreipc1", + "record": { + "$type": "sh.tangled.repo.pull.comment", + "pull": "at://did:plc:repoowner/sh.tangled.repo.pull/pull1", + "body": "Looks good to me.", + "createdAt": "2026-03-20T14:00:00.000Z" + } + } +} diff --git a/packages/api/internal/store/db.go b/packages/api/internal/store/db.go index c93ad4d..e0b4024 100644 --- a/packages/api/internal/store/db.go +++ b/packages/api/internal/store/db.go @@ -15,6 +15,8 @@ import ( //go:embed migrations/*.sql var migrationsFS embed.FS +var extensionMigrationNoticeLogged bool + // Open establishes a connection to the database. // For remote Turso URLs (libsql:// or https://) it uses the libsql-client-go driver. // For local file: URLs it uses the pure-Go SQLite driver (no CGo required). @@ -73,8 +75,13 @@ func execMigration(db *sql.DB, name, content string) error { if _, err := db.Exec(stmt); err != nil { upper := strings.ToUpper(stmt) if strings.Contains(upper, "USING FTS") || strings.Contains(upper, "LIBSQL_VECTOR_IDX") { - slog.Warn("migration: skipping extension index (not supported in this environment)", - "migration", name, "err", err) + if !extensionMigrationNoticeLogged { + extensionMigrationNoticeLogged = true + slog.Info("migration: skipping Turso extension indexes in this environment", + "migration", name, + "reason", "database engine does not support Turso-specific FTS/vector DDL", + ) + } continue } return fmt.Errorf("migration %s: exec failed: %w\nstatement: %s", name, err, stmt) diff --git a/packages/api/internal/store/sql_store.go b/packages/api/internal/store/sql_store.go index 298a6ea..b686c3a 100644 --- a/packages/api/internal/store/sql_store.go +++ b/packages/api/internal/store/sql_store.go @@ -178,6 +178,79 @@ func (s *SQLStore) EnqueueEmbeddingJob(ctx context.Context, documentID string) e return nil } +func (s *SQLStore) GetFollowSubjects(ctx context.Context, did string) ([]string, error) { + rows, err := s.db.QueryContext(ctx, ` + SELECT DISTINCT repo_did + FROM documents + WHERE did = ? + AND collection = 'sh.tangled.graph.follow' + AND deleted_at IS NULL + AND repo_did IS NOT NULL + AND repo_did != ''`, + did, + ) + if err != nil { + return nil, fmt.Errorf("get follow subjects: %w", err) + } + defer rows.Close() + + var subjects []string + for rows.Next() { + var subject string + if err := rows.Scan(&subject); err != nil { + return nil, fmt.Errorf("scan follow subject: %w", err) + } + subjects = append(subjects, subject) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate follow subjects: %w", err) + } + return subjects, nil +} + +func (s *SQLStore) GetRepoCollaborators(ctx context.Context, repoOwnerDID string) ([]string, error) { + rows, err := s.db.QueryContext(ctx, ` + SELECT DISTINCT did + FROM documents + WHERE repo_did = ? + AND did != ? + AND deleted_at IS NULL + AND collection IN ( + 'sh.tangled.repo.issue', + 'sh.tangled.repo.pull', + 'sh.tangled.repo.issue.comment', + 'sh.tangled.repo.pull.comment' + )`, + repoOwnerDID, repoOwnerDID, + ) + if err != nil { + return nil, fmt.Errorf("get repo collaborators: %w", err) + } + defer rows.Close() + + var collaborators []string + for rows.Next() { + var collaborator string + if err := rows.Scan(&collaborator); err != nil { + return nil, fmt.Errorf("scan collaborator: %w", err) + } + collaborators = append(collaborators, collaborator) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate collaborators: %w", err) + } + return collaborators, nil +} + +func (s *SQLStore) CountDocuments(ctx context.Context) (int64, error) { + var n int64 + err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM documents WHERE deleted_at IS NULL`).Scan(&n) + if err != nil { + return 0, fmt.Errorf("count documents: %w", err) + } + return n, 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 index 1673b0d..a1a5efc 100644 --- a/packages/api/internal/store/store.go +++ b/packages/api/internal/store/store.go @@ -51,4 +51,7 @@ type Store interface { 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 + GetFollowSubjects(ctx context.Context, did string) ([]string, error) + GetRepoCollaborators(ctx context.Context, repoOwnerDID string) ([]string, error) + CountDocuments(ctx context.Context) (int64, error) } diff --git a/packages/api/internal/store/store_test.go b/packages/api/internal/store/store_test.go index b8e9cfc..f1ee2ff 100644 --- a/packages/api/internal/store/store_test.go +++ b/packages/api/internal/store/store_test.go @@ -236,4 +236,83 @@ func TestIntegration(t *testing.T) { t.Fatalf("last_error: got %q, want NULL", lastError.String) } }) + + t.Run("follow subject discovery query", func(t *testing.T) { + followDoc := &store.Document{ + ID: "did:plc:owner|sh.tangled.graph.follow|f1", + DID: "did:plc:owner", + Collection: "sh.tangled.graph.follow", + RKey: "f1", + ATURI: "at://did:plc:owner/sh.tangled.graph.follow/f1", + CID: "cid-follow", + RecordType: "follow", + RepoDID: "did:plc:target", + } + if err := st.UpsertDocument(ctx, followDoc); err != nil { + t.Fatalf("upsert follow doc: %v", err) + } + + subjects, err := st.GetFollowSubjects(ctx, "did:plc:owner") + if err != nil { + t.Fatalf("get follow subjects: %v", err) + } + if len(subjects) != 1 || subjects[0] != "did:plc:target" { + t.Fatalf("subjects: got %#v", subjects) + } + }) + + t.Run("repo collaborator discovery query", func(t *testing.T) { + docs := []*store.Document{ + { + ID: "did:plc:collab1|sh.tangled.repo.issue|i1", + DID: "did:plc:collab1", + Collection: "sh.tangled.repo.issue", + RKey: "i1", + ATURI: "at://did:plc:collab1/sh.tangled.repo.issue/i1", + CID: "cid-c1", + RecordType: "issue", + RepoDID: "did:plc:owner", + }, + { + ID: "did:plc:collab2|sh.tangled.repo.pull.comment|pc1", + DID: "did:plc:collab2", + Collection: "sh.tangled.repo.pull.comment", + RKey: "pc1", + ATURI: "at://did:plc:collab2/sh.tangled.repo.pull.comment/pc1", + CID: "cid-c2", + RecordType: "pull_comment", + RepoDID: "did:plc:owner", + }, + { + ID: "did:plc:owner|sh.tangled.repo.issue|i-owner", + DID: "did:plc:owner", + Collection: "sh.tangled.repo.issue", + RKey: "i-owner", + ATURI: "at://did:plc:owner/sh.tangled.repo.issue/i-owner", + CID: "cid-owner", + RecordType: "issue", + RepoDID: "did:plc:owner", + }, + } + for _, doc := range docs { + if err := st.UpsertDocument(ctx, doc); err != nil { + t.Fatalf("upsert collaborator doc %s: %v", doc.ID, err) + } + } + + collaborators, err := st.GetRepoCollaborators(ctx, "did:plc:owner") + if err != nil { + t.Fatalf("get collaborators: %v", err) + } + if len(collaborators) != 2 { + t.Fatalf("collaborators length: got %d want 2 (%#v)", len(collaborators), collaborators) + } + got := map[string]bool{} + for _, did := range collaborators { + got[did] = true + } + if !got["did:plc:collab1"] || !got["did:plc:collab2"] { + t.Fatalf("collaborators: got %#v", collaborators) + } + }) } diff --git a/packages/api/internal/tapclient/tapclient.go b/packages/api/internal/tapclient/tapclient.go index 8ce042b..aea12ba 100644 --- a/packages/api/internal/tapclient/tapclient.go +++ b/packages/api/internal/tapclient/tapclient.go @@ -9,6 +9,7 @@ import ( "math/rand/v2" "net/http" "strconv" + "strings" "sync" "time" @@ -19,6 +20,8 @@ import ( const ( minReconnectBackoff = 500 * time.Millisecond maxReconnectBackoff = 10 * time.Second + keepAliveInterval = 20 * time.Second + keepAliveTimeout = 5 * time.Second ) // Client receives Tap events over WebSocket and sends acks after processing. @@ -85,6 +88,9 @@ func (c *Client) AckEvent(ctx context.Context, id int64) error { payload, _ := json.Marshal(map[string]int64{"id": id}) if err := conn.Write(ctx, websocket.MessageText, payload); err == nil { return nil + } else if isConnectionWriteError(err) { + c.resetConn(websocket.StatusInternalError, "ack json write failed") + return fmt.Errorf("ack event %d: %w", id, err) } c.log.Warn("tap ack json failed; trying plain id", slog.Int64("event_id", id)) @@ -144,6 +150,7 @@ func (c *Client) ensureConnected(ctx context.Context) (*websocket.Conn, error) { c.mu.Lock() if c.conn == nil { c.conn = conn + c.startKeepAlive(conn) } else { _ = conn.Close(websocket.StatusNormalClosure, "duplicate") } @@ -179,3 +186,39 @@ func (c *Client) resetConn(status websocket.StatusCode, reason string) { _ = c.conn.Close(status, reason) c.conn = nil } + +func (c *Client) startKeepAlive(conn *websocket.Conn) { + go func() { + ticker := time.NewTicker(keepAliveInterval) + defer ticker.Stop() + + for range ticker.C { + c.mu.Lock() + if c.conn != conn { + c.mu.Unlock() + return + } + c.mu.Unlock() + + ctx, cancel := context.WithTimeout(context.Background(), keepAliveTimeout) + err := conn.Ping(ctx) + cancel() + if err != nil { + c.log.Warn("tap keepalive ping failed", slog.String("error", err.Error())) + c.resetConn(websocket.StatusInternalError, "keepalive failed") + return + } + } + }() +} + +func isConnectionWriteError(err error) bool { + if err == nil { + return false + } + msg := strings.ToLower(err.Error()) + return strings.Contains(msg, "broken pipe") || + strings.Contains(msg, "connection reset") || + strings.Contains(msg, "closed network connection") || + strings.Contains(msg, "i/o timeout") +} diff --git a/packages/api/main.go b/packages/api/main.go index c555bb2..8de8797 100644 --- a/packages/api/main.go +++ b/packages/api/main.go @@ -7,8 +7,10 @@ import ( "os" "os/signal" "syscall" + "time" "github.com/spf13/cobra" + "tangled.org/desertthunder.dev/twister/internal/backfill" "tangled.org/desertthunder.dev/twister/internal/config" "tangled.org/desertthunder.dev/twister/internal/ingest" "tangled.org/desertthunder.dev/twister/internal/normalize" @@ -24,14 +26,17 @@ var ( func main() { root := &cobra.Command{ - Use: "twister", - Short: "Tangled search service", - Version: fmt.Sprintf("%s (%s)", version, commit), + Use: "twister", + Short: "Tangled search service", + Version: fmt.Sprintf("%s (%s)", version, commit), + SilenceUsage: true, + SilenceErrors: true, } root.AddCommand( newAPICmd(), newIndexerCmd(), + newBackfillCmd(), newEmbedWorkerCmd(), newReindexCmd(), newReembedCmd(), @@ -39,6 +44,7 @@ func main() { ) if err := root.Execute(); err != nil { + _, _ = fmt.Fprintln(os.Stderr, "Error:", err) os.Exit(1) } } @@ -138,6 +144,69 @@ func newEmbedWorkerCmd() *cobra.Command { } } +func newBackfillCmd() *cobra.Command { + var opts backfill.Options + + cmd := &cobra.Command{ + Use: "backfill", + Short: "Discover users from seeds and register repos for Tap backfill", + RunE: func(cmd *cobra.Command, args []string) error { + cfg, err := config.Load() + if err != nil { + return fmt.Errorf("config: %w", err) + } + log := observability.NewLogger(cfg) + log.Info("starting backfill", slog.String("service", "backfill"), slog.String("version", version)) + + if cfg.TapURL == "" { + return fmt.Errorf("TAP_URL is required for backfill") + } + + 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) + } + + tapAdmin, err := backfill.NewHTTPTapAdmin(cfg.TapURL, cfg.TapAuthPassword) + if err != nil { + return fmt.Errorf("tap admin client: %w", err) + } + + runner := backfill.NewRunner( + store.New(db), + tapAdmin, + backfill.NewHTTPHandleResolver(""), + log, + ) + + ctx, cancel := baseContext() + defer cancel() + + if err := runner.Run(ctx, opts); err != nil { + return fmt.Errorf("run backfill: %w", err) + } + + log.Info("shutting down backfill") + return nil + }, + } + + cmd.Flags().StringVar(&opts.SeedsPath, "seeds", "", "Seed source: file path or comma-separated DIDs/handles (required)") + cmd.Flags().IntVar(&opts.MaxHops, "max-hops", 2, "Max fan-out depth from seeds") + cmd.Flags().BoolVar(&opts.DryRun, "dry-run", false, "Print discovery plan without mutating Tap") + cmd.Flags().IntVar(&opts.Concurrency, "concurrency", 5, "Parallel discovery workers") + cmd.Flags().IntVar(&opts.BatchSize, "batch-size", 10, "DIDs per /repos/add request") + cmd.Flags().DurationVar(&opts.BatchDelay, "batch-delay", time.Second, "Delay between Tap /repos/add batches") + _ = cmd.MarkFlagRequired("seeds") + + return cmd +} + func newReindexCmd() *cobra.Command { return &cobra.Command{ Use: "reindex",