From 3320c48f78bbd3125fa8721696ab755fd5aa8152 Mon Sep 17 00:00:00 2001 From: Owais Jamil Date: Wed, 25 Mar 2026 10:30:03 -0500 Subject: [PATCH] feat: read-through indexing job processing - job retry with exponential backoff. - update local SQLite configuration with PRAGMA settings for better concurrency. - smoke checks/tests for API endpoints using uv --- .gitignore | 5 + README.md | 33 +- docs/reference/api.md | 32 +- docs/roadmap.md | 6 +- justfile | 10 +- package.json | 4 +- packages/api/README.md | 18 + packages/api/internal/api/actors.go | 16 + packages/api/internal/api/api.go | 21 + packages/api/internal/api/readthrough.go | 202 ++++++++++ packages/api/internal/ingest/ingest_test.go | 16 + packages/api/internal/store/db.go | 27 ++ packages/api/internal/store/db_test.go | 30 ++ .../store/migrations/005_indexing_jobs.sql | 16 + packages/api/internal/store/sql_store.go | 120 +++++- packages/api/internal/store/store.go | 38 +- packages/api/internal/store/store_test.go | 64 ++++ packages/api/justfile | 24 +- scripts/api/README.md | 24 ++ scripts/api/pyproject.toml | 20 + scripts/api/src/twister_api_smoke/__init__.py | 1 + scripts/api/src/twister_api_smoke/cli.py | 362 ++++++++++++++++++ scripts/api/uv.lock | 8 + 23 files changed, 1049 insertions(+), 48 deletions(-) create mode 100644 packages/api/internal/api/readthrough.go create mode 100644 packages/api/internal/store/migrations/005_indexing_jobs.sql create mode 100644 scripts/api/README.md create mode 100644 scripts/api/pyproject.toml create mode 100644 scripts/api/src/twister_api_smoke/__init__.py create mode 100644 scripts/api/src/twister_api_smoke/cli.py create mode 100644 scripts/api/uv.lock diff --git a/.gitignore b/.gitignore index 242e550..889f125 100644 --- a/.gitignore +++ b/.gitignore @@ -33,5 +33,10 @@ platforms plugins www +# SQLite *.db +*.db-shm +*.db-wal + .env +__pycache__/ diff --git a/README.md b/README.md index 76e5da0..423b2ca 100644 --- a/README.md +++ b/README.md @@ -41,37 +41,42 @@ pnpm install Start the Ionic/Vite app: ```bash -pnpm dev -# or: just dev +pnpm dev # or: just dev ``` That serves the client from `apps/twisted` with Vite. -To run the Go API locally, make sure `packages/api/.env` has at least: +To run the Go API locally for routine experimentation, no Turso credentials are required. -- `TURSO_DATABASE_URL` -- `TURSO_AUTH_TOKEN` +Start the API in local file mode: + +```bash +pnpm api:run:api # or: just api-dev +``` -Then start the API: +This serves the API and search site on `http://localhost:8080` using +`packages/api/twister-dev.db`. + +To run the API against remote Turso instead: ```bash -pnpm api:run:api -# or: just api-dev +just api-dev remote ``` -This serves the API and search site on `http://localhost:8080`. +To run the indexer in local file mode as well: -To run the indexer as well, `packages/api/.env` also needs: +```bash +pnpm api:run:indexer # or: just api-run-indexer +``` + +To run the indexer against remote Turso, `packages/api/.env` needs: - `TAP_URL` - `TAP_AUTH_PASSWORD` - `INDEXED_COLLECTIONS` -Then start the indexer in a separate terminal: - ```bash -pnpm api:run:indexer -# or: just api-run-indexer +just api-run-indexer remote ``` Typical local setup is three terminals: diff --git a/docs/reference/api.md b/docs/reference/api.md index 101b7f1..7eb3a71 100644 --- a/docs/reference/api.md +++ b/docs/reference/api.md @@ -133,22 +133,22 @@ The backfill command discovers users from a seed file and registers them with Ta All configuration is via environment variables (with `.env` file support): -| Variable | Default | Purpose | -| -------------------------- | ----------------------- | ---------------------------------------------- | -| `TURSO_DATABASE_URL` | — | Database connection (required) | -| `TURSO_AUTH_TOKEN` | — | Auth token (required for remote) | -| `TAP_URL` | — | Tap WebSocket URL | -| `TAP_AUTH_PASSWORD` | — | Tap admin password | -| `INDEXED_COLLECTIONS` | all | Collection allowlist (CSV, supports wildcards) | -| `HTTP_BIND_ADDR` | `:8080` | API server bind address | -| `INDEXER_HEALTH_ADDR` | `:9090` | Indexer health probe address | -| `LOG_LEVEL` | info | debug/info/warn/error | -| `LOG_FORMAT` | json | json or text | -| `ENABLE_ADMIN_ENDPOINTS` | false | Enable admin routes | -| `ADMIN_AUTH_TOKEN` | — | Bearer token for admin | -| `ENABLE_INGEST_ENRICHMENT` | true | XRPC enrichment at ingest time | -| `PLC_DIRECTORY_URL` | `https://plc.directory` | PLC Directory | -| `XRPC_TIMEOUT` | 15s | XRPC HTTP timeout | +| Variable | Default | Purpose | +| -------------------------- | ----------------------- | ----------------------------------------------- | +| `TURSO_DATABASE_URL` | — | Database connection (required unless `--local`) | +| `TURSO_AUTH_TOKEN` | — | Auth token (required for remote) | +| `TAP_URL` | — | Tap WebSocket URL | +| `TAP_AUTH_PASSWORD` | — | Tap admin password | +| `INDEXED_COLLECTIONS` | all | Collection allowlist (CSV, supports wildcards) | +| `HTTP_BIND_ADDR` | `:8080` | API server bind address | +| `INDEXER_HEALTH_ADDR` | `:9090` | Indexer health probe address | +| `LOG_LEVEL` | info | debug/info/warn/error | +| `LOG_FORMAT` | json | json or text | +| `ENABLE_ADMIN_ENDPOINTS` | false | Enable admin routes | +| `ADMIN_AUTH_TOKEN` | — | Bearer token for admin | +| `ENABLE_INGEST_ENRICHMENT` | true | XRPC enrichment at ingest time | +| `PLC_DIRECTORY_URL` | `https://plc.directory` | PLC Directory | +| `XRPC_TIMEOUT` | 15s | XRPC HTTP timeout | ## Deployment diff --git a/docs/roadmap.md b/docs/roadmap.md index ae6b2d8..805f3df 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -7,19 +7,19 @@ updated: 2026-03-25 Highest priority. This work blocks further investment in semantic search, hybrid ranking, and broader discovery features. -- [ ] Stabilize local development and experimentation around a local `file:` database +- [x] Stabilize local development and experimentation around a local `file:` database - [x] Document backup, restore, and disk-growth procedures for the experimental local DB - [x] Research production backend options: PostgreSQL, Turso remote/libSQL, and Turso embedded replicas - [x] Write a production storage decision record with workload and operational tradeoffs, using `docs/adr/pg.md` and `docs/adr/turso.md` - [x] Define the migration path from the experimental local setup to the chosen production backend -- [ ] Add cURL smoke tests for `healthz`, `readyz`, `search`, `documents`, indexing, and activity in `scripts/api/` +- [x] Add a durable read-through indexing job queue for records fetched through the API +- [x] Add API smoke tests for `healthz`, `readyz`, `search`, `documents`, indexing, and activity in `scripts/api/` - desertthunder.dev DID: `did:plc:xg2vq45muivyy3xwatcehspu` - Twisted AT URI: `at://did:plc:xg2vq45muivyy3xwatcehspu/sh.tangled.repo/3mho6hukiei22` - Profile AT URI: `at://did:plc:xg2vq45muivyy3xwatcehspu/sh.tangled.actor.profile/self` - Follow AT URI (desertthunder.dev follows npmx): `at://did:plc:xg2vq45muivyy3xwatcehspu/sh.tangled.graph.follow/3mhofstanru22` - Star AT URI (desertthunder.dev stars microcosm-rs): `at://did:plc:lulmyldiq4sb2ikags5sfb25/sh.tangled.repo/3lvsxzinfz222` - ~~Add `just` targets for smoke-test runs locally and against a remote base URL~~ directly invoking the scripts is fine. -- [ ] Add a durable read-through indexing job queue for records fetched through the API - [ ] Reuse the existing normalization and upsert path for on-demand indexing jobs - [ ] Trigger indexing jobs from repo, issue, PR, profile, and similar fetch handlers - [ ] Add dedupe, retries, and observability for indexing jobs diff --git a/justfile b/justfile index 2c7929a..ab4585c 100644 --- a/justfile +++ b/justfile @@ -41,11 +41,13 @@ app-cap-android: api-build: just --justfile packages/api/justfile build -api-dev: - just --justfile packages/api/justfile run-api +# Run API. Usage: just api-dev [mode], mode: local|remote (default local) +api-dev mode="local": + just --justfile packages/api/justfile run-api {{mode}} -api-run-indexer: - just --justfile packages/api/justfile run-indexer +# Run indexer. Usage: just api-run-indexer [mode], mode: local|remote (default local) +api-run-indexer mode="local": + just --justfile packages/api/justfile run-indexer {{mode}} api-test: just --justfile packages/api/justfile test diff --git a/package.json b/package.json index 76a4d1f..0629cfc 100644 --- a/package.json +++ b/package.json @@ -13,7 +13,7 @@ "app:cap:run:android": "pnpm --dir apps/twisted exec cap run android", "api:build": "just --justfile packages/api/justfile build", "api:test": "just --justfile packages/api/justfile test", - "api:run:api": "just --justfile packages/api/justfile run-api", - "api:run:indexer": "just --justfile packages/api/justfile run-indexer" + "api:run:api": "just api-dev", + "api:run:indexer": "just api-run-indexer" } } diff --git a/packages/api/README.md b/packages/api/README.md index 32e313b..73e7618 100644 --- a/packages/api/README.md +++ b/packages/api/README.md @@ -18,6 +18,24 @@ go run . api --local The server listens on `:8080` by default. Logs are printed as text when `--local` is set. +## API Smoke Tests + +Smoke checks for the API surface live in a uv-managed Python project at +`scripts/api/`. + +From the repo root: + +```sh +uv run --project scripts/api twister-api-smoke +``` + +Optional base URL override: + +```sh +TWISTER_API_BASE_URL=http://localhost:8080 \ + uv run --project scripts/api twister-api-smoke +``` + ## Experimental Local DB Operations The experimental local database lives at `packages/api/twister-dev.db` when you run Twister from `packages/api` with `--local`. diff --git a/packages/api/internal/api/actors.go b/packages/api/internal/api/actors.go index f44146b..4d3a897 100644 --- a/packages/api/internal/api/actors.go +++ b/packages/api/internal/api/actors.go @@ -84,6 +84,7 @@ func (s *Server) resolveRepo(r *http.Request, handleOrDID, repoName string) (*re if err != nil { return nil, fmt.Errorf("list repos for %s: %w", actor.DID, err) } + s.enqueueXRPCList(r.Context(), entries) for _, entry := range entries { name, _ := entry.Value["name"].(string) @@ -208,6 +209,7 @@ func (s *Server) handleGetActor(w http.ResponseWriter, r *http.Request) { s.actorError(w, err) return } + s.enqueueXRPCRecord(r.Context(), rec.URI, rec.CID, rec.Value) var bsky *bskyProfileResponse if linked, _ := rec.Value["bluesky"].(bool); linked { @@ -243,6 +245,7 @@ func (s *Server) handleListActorRepos(w http.ResponseWriter, r *http.Request) { s.actorError(w, err) return } + s.enqueueXRPCList(r.Context(), entries) records := make([]recordEntry, len(entries)) for i, e := range entries { @@ -279,6 +282,7 @@ func (s *Server) handleGetActorRepo(w http.ResponseWriter, r *http.Request) { s.actorError(w, err) return } + s.enqueueXRPCRecord(r.Context(), rec.URI, rec.CID, rec.Value) writeJSON(w, http.StatusOK, map[string]any{ "did": repo.DID, @@ -443,6 +447,7 @@ func (s *Server) handleRepoIssues(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusBadGateway, errorBody("upstream_error", "failed to fetch issues")) return } + s.enqueueXRPCList(r.Context(), issues) var records []issueEntry for _, e := range issues { @@ -480,6 +485,7 @@ func (s *Server) handleRepoPulls(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusBadGateway, errorBody("upstream_error", "failed to fetch pulls")) return } + s.enqueueXRPCList(r.Context(), pulls) var records []pullEntry for _, e := range pulls { @@ -518,6 +524,7 @@ func (s *Server) handleActorIssues(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusBadGateway, errorBody("upstream_error", "failed to fetch issues")) return } + s.enqueueXRPCList(r.Context(), issues) records := make([]issueEntry, len(issues)) for i, e := range issues { @@ -548,6 +555,7 @@ func (s *Server) handleActorPulls(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusBadGateway, errorBody("upstream_error", "failed to fetch pulls")) return } + s.enqueueXRPCList(r.Context(), pulls) records := make([]pullEntry, len(pulls)) for i, e := range pulls { @@ -578,6 +586,7 @@ func (s *Server) handleActorFollowing(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusBadGateway, errorBody("upstream_error", "failed to fetch follows")) return } + s.enqueueXRPCList(r.Context(), entries) records := make([]recordEntry, len(entries)) for i, e := range entries { @@ -605,6 +614,7 @@ func (s *Server) handleActorStrings(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusBadGateway, errorBody("upstream_error", "failed to fetch strings")) return } + s.enqueueXRPCList(r.Context(), entries) records := make([]recordEntry, len(entries)) for i, e := range entries { @@ -635,6 +645,7 @@ func (s *Server) handleIssueDetail(w http.ResponseWriter, r *http.Request) { s.actorError(w, err) return } + s.enqueueXRPCRecord(r.Context(), rec.URI, rec.CID, rec.Value) _, stateMap, err := s.fetchIssuesAndStates(r, actor.PDS, actor.DID) if err != nil { @@ -667,6 +678,7 @@ func (s *Server) handleIssueComments(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusBadGateway, errorBody("upstream_error", "failed to fetch comments")) return } + s.enqueueXRPCList(r.Context(), entries) var records []recordEntry for _, e := range entries { @@ -704,6 +716,7 @@ func (s *Server) handlePullDetail(w http.ResponseWriter, r *http.Request) { s.actorError(w, err) return } + s.enqueueXRPCRecord(r.Context(), rec.URI, rec.CID, rec.Value) _, statusMap, err := s.fetchPullsAndStatuses(r, actor.PDS, actor.DID) if err != nil { @@ -736,6 +749,7 @@ func (s *Server) handlePullComments(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusBadGateway, errorBody("upstream_error", "failed to fetch comments")) return } + s.enqueueXRPCList(r.Context(), entries) var records []recordEntry for _, e := range entries { @@ -792,6 +806,7 @@ func (s *Server) fetchIssuesAndStates(r *http.Request, pds, did string) ([]xrpc. } stateMap := make(map[string]string, len(states)) + s.enqueueXRPCList(r.Context(), states) for _, e := range states { issueURI, _ := e.Value["issue"].(string) state, _ := e.Value["state"].(string) @@ -839,6 +854,7 @@ func (s *Server) fetchPullsAndStatuses(r *http.Request, pds, did string) ([]xrpc } statusMap := make(map[string]string, len(statuses)) + s.enqueueXRPCList(r.Context(), statuses) for _, e := range statuses { pullURI, _ := e.Value["pull"].(string) status, _ := e.Value["status"].(string) diff --git a/packages/api/internal/api/api.go b/packages/api/internal/api/api.go index 1b22706..bdab831 100644 --- a/packages/api/internal/api/api.go +++ b/packages/api/internal/api/api.go @@ -1,10 +1,12 @@ package api import ( + "bufio" "context" "encoding/json" "fmt" "log/slog" + "net" "net/http" "strconv" "strings" @@ -13,6 +15,7 @@ import ( "tangled.org/desertthunder.dev/twister/internal/config" "tangled.org/desertthunder.dev/twister/internal/constellation" + "tangled.org/desertthunder.dev/twister/internal/normalize" "tangled.org/desertthunder.dev/twister/internal/reindex" "tangled.org/desertthunder.dev/twister/internal/search" "tangled.org/desertthunder.dev/twister/internal/store" @@ -28,6 +31,7 @@ type Server struct { log *slog.Logger constellation *constellation.Client xrpc *xrpc.Client + registry *normalize.Registry } // New creates a new API server. @@ -39,6 +43,7 @@ func New(searchRepo *search.Repository, st store.Store, cfg *config.Config, log log: log, constellation: constellation, xrpc: xrpcClient, + registry: normalize.NewRegistry(), } } @@ -114,6 +119,8 @@ func (s *Server) Run(ctx context.Context) error { } errCh := make(chan error, 1) + go s.runReadThroughIndexer(ctx) + go func() { s.log.Info("listening", slog.String("addr", s.cfg.HTTPBindAddr)) errCh <- srv.ListenAndServe() @@ -171,6 +178,20 @@ func (rw *responseWriter) WriteHeader(code int) { rw.ResponseWriter.WriteHeader(code) } +func (rw *responseWriter) Hijack() (net.Conn, *bufio.ReadWriter, error) { + h, ok := rw.ResponseWriter.(http.Hijacker) + if !ok { + return nil, nil, fmt.Errorf("response writer does not support hijacking") + } + return h.Hijack() +} + +func (rw *responseWriter) Flush() { + if f, ok := rw.ResponseWriter.(http.Flusher); ok { + f.Flush() + } +} + func (s *Server) handleHealthz(w http.ResponseWriter, _ *http.Request) { writeJSON(w, http.StatusOK, map[string]string{"status": "ok"}) } diff --git a/packages/api/internal/api/readthrough.go b/packages/api/internal/api/readthrough.go new file mode 100644 index 0000000..99e1005 --- /dev/null +++ b/packages/api/internal/api/readthrough.go @@ -0,0 +1,202 @@ +package api + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "time" + + "tangled.org/desertthunder.dev/twister/internal/normalize" + "tangled.org/desertthunder.dev/twister/internal/store" + "tangled.org/desertthunder.dev/twister/internal/xrpc" +) + +const readThroughIdlePoll = 1 * time.Second + +func (s *Server) runReadThroughIndexer(ctx context.Context) { + ticker := time.NewTicker(readThroughIdlePoll) + defer ticker.Stop() + + s.log.Info("read-through indexer worker started") + for { + if ctx.Err() != nil { + s.log.Info("read-through indexer worker stopped") + return + } + + job, err := s.store.ClaimIndexingJob(ctx) + if err != nil { + s.log.Warn("read-through claim failed", slog.String("error", err.Error())) + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + continue + } + if job == nil { + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + continue + } + + if err := s.processReadThroughJob(ctx, job); err != nil { + nextDelay := retryDelay(job.Attempts + 1) + nextAt := time.Now().UTC().Add(nextDelay).Format(time.RFC3339) + retryErr := s.store.RetryIndexingJob(ctx, job.DocumentID, nextAt, truncateErr(err)) + if retryErr != nil { + s.log.Error("read-through retry update failed", + slog.String("document_id", job.DocumentID), + slog.String("error", retryErr.Error()), + ) + continue + } + s.log.Warn("read-through job failed; scheduled retry", + slog.String("document_id", job.DocumentID), + slog.Int("attempt", job.Attempts+1), + slog.Duration("retry_in", nextDelay), + slog.String("error", err.Error()), + ) + continue + } + + if err := s.store.CompleteIndexingJob(ctx, job.DocumentID); err != nil { + s.log.Error("read-through complete failed", + slog.String("document_id", job.DocumentID), + slog.String("error", err.Error()), + ) + continue + } + } +} + +func (s *Server) processReadThroughJob(ctx context.Context, job *store.IndexingJob) error { + record := map[string]any{} + if err := json.Unmarshal([]byte(job.RecordJSON), &record); err != nil { + return fmt.Errorf("decode record json: %w", err) + } + + event := normalize.TapRecordEvent{ + Type: "record", + Record: &normalize.TapRecord{ + DID: job.DID, + Collection: job.Collection, + RKey: job.RKey, + Action: "create", + CID: job.CID, + Record: record, + }, + } + + if handler, ok := s.registry.StateHandler(job.Collection); ok { + update, err := handler.HandleState(event) + if err != nil { + return fmt.Errorf("state normalize: %w", err) + } + if err := s.store.UpdateRecordState(ctx, update.SubjectURI, update.State); err != nil { + return fmt.Errorf("update state: %w", err) + } + return nil + } + + adapter, ok := s.registry.Adapter(job.Collection) + if !ok { + return nil + } + + doc, err := adapter.Normalize(event) + if err != nil { + return fmt.Errorf("normalize record: %w", err) + } + + handle, err := s.store.GetIdentityHandle(ctx, job.DID) + if err != nil { + return fmt.Errorf("lookup identity handle: %w", err) + } + if handle != "" { + doc.AuthorHandle = handle + if doc.RecordType == "profile" { + doc.Title = handle + } + } + + if err := s.store.UpsertDocument(ctx, doc); err != nil { + return fmt.Errorf("upsert document: %w", err) + } + + if adapter.Searchable(record) { + if err := s.store.EnqueueEmbeddingJob(ctx, doc.ID); err != nil { + s.log.Warn("read-through enqueue embedding failed", + slog.String("document_id", doc.ID), + slog.String("error", err.Error()), + ) + } + } + return nil +} + +func (s *Server) enqueueXRPCRecord(ctx context.Context, uri, cid string, value map[string]any) { + did, collection, rkey, err := normalize.ParseATURI(uri) + if err != nil { + s.log.Debug("read-through skip invalid at-uri", slog.String("uri", uri), slog.String("error", err.Error())) + return + } + payload, err := json.Marshal(value) + if err != nil { + s.log.Debug("read-through skip unmarshalable record", slog.String("uri", uri), slog.String("error", err.Error())) + return + } + input := store.IndexingJobInput{ + DocumentID: normalize.StableID(did, collection, rkey), + DID: did, + Collection: collection, + RKey: rkey, + CID: cid, + RecordJSON: string(payload), + } + if err := s.store.EnqueueIndexingJob(ctx, input); err != nil { + s.log.Warn("enqueue read-through indexing job failed", + slog.String("document_id", input.DocumentID), + slog.String("error", err.Error()), + ) + } +} + +func (s *Server) enqueueXRPCList(ctx context.Context, entries []xrpc.ListRecordEntry) { + for _, e := range entries { + s.enqueueXRPCRecord(ctx, e.URI, e.CID, e.Value) + } +} + +func retryDelay(attempt int) time.Duration { + if attempt < 1 { + attempt = 1 + } + base := time.Second * time.Duration(1< 5*time.Minute { + return 5 * time.Minute + } + return base +} + +func truncateErr(err error) string { + if err == nil { + return "" + } + msg := err.Error() + if len(msg) > 500 { + return msg[:500] + } + return msg +} + +func minInt(a, b int) int { + if a < b { + return a + } + return b +} diff --git a/packages/api/internal/ingest/ingest_test.go b/packages/api/internal/ingest/ingest_test.go index be52bad..73e556b 100644 --- a/packages/api/internal/ingest/ingest_test.go +++ b/packages/api/internal/ingest/ingest_test.go @@ -103,6 +103,22 @@ func (f *fakeStore) EnqueueEmbeddingJob(_ context.Context, documentID string) er return nil } +func (f *fakeStore) EnqueueIndexingJob(_ context.Context, _ store.IndexingJobInput) error { + return nil +} + +func (f *fakeStore) ClaimIndexingJob(_ context.Context) (*store.IndexingJob, error) { + return nil, nil +} + +func (f *fakeStore) CompleteIndexingJob(_ context.Context, _ string) error { + return nil +} + +func (f *fakeStore) RetryIndexingJob(_ context.Context, _, _, _ string) error { + return nil +} + func (f *fakeStore) GetFollowSubjects(_ context.Context, _ string) ([]string, error) { return nil, nil } diff --git a/packages/api/internal/store/db.go b/packages/api/internal/store/db.go index b59c414..96ee8bb 100644 --- a/packages/api/internal/store/db.go +++ b/packages/api/internal/store/db.go @@ -7,6 +7,7 @@ import ( "log/slog" "sort" "strings" + "time" _ "github.com/tursodatabase/libsql-client-go/libsql" _ "modernc.org/sqlite" @@ -31,6 +32,12 @@ func Open(url, token string) (*sql.DB, error) { if err != nil { return nil, fmt.Errorf("open db: %w", err) } + if strings.HasPrefix(url, "file:") { + if err := configureLocalSQLite(db); err != nil { + db.Close() + return nil, err + } + } if err := db.Ping(); err != nil { db.Close() return nil, fmt.Errorf("ping db: %w", err) @@ -38,6 +45,26 @@ func Open(url, token string) (*sql.DB, error) { return db, nil } +func configureLocalSQLite(db *sql.DB) error { + // Busy timeout gives the writer a window to wait instead of failing fast with "database is locked". + if _, err := db.Exec(`PRAGMA busy_timeout = 5000`); err != nil { + return fmt.Errorf("configure sqlite busy_timeout: %w", err) + } + // WAL mode allows concurrent readers with a writer and is the default for multi-process local dev. + if _, err := db.Exec(`PRAGMA journal_mode = WAL`); err != nil { + return fmt.Errorf("configure sqlite wal mode: %w", err) + } + if _, err := db.Exec(`PRAGMA synchronous = NORMAL`); err != nil { + return fmt.Errorf("configure sqlite synchronous mode: %w", err) + } + + db.SetMaxOpenConns(1) + db.SetMaxIdleConns(1) + db.SetConnMaxLifetime(0) + db.SetConnMaxIdleTime(5 * time.Minute) + return nil +} + // driverAndDSN returns the sql driver name and DSN for the given URL. // file: URLs use the pure-Go "sqlite" driver; all others use "libsql". func driverAndDSN(url, token string) (driver, dsn string) { diff --git a/packages/api/internal/store/db_test.go b/packages/api/internal/store/db_test.go index cf3bf43..f86d68e 100644 --- a/packages/api/internal/store/db_test.go +++ b/packages/api/internal/store/db_test.go @@ -2,6 +2,7 @@ package store import ( "database/sql" + "path/filepath" "strings" "testing" @@ -42,3 +43,32 @@ func TestExecMigrationFailsForRemoteWhenNativeFTSUnavailable(t *testing.T) { t.Fatalf("unexpected error: %v", err) } } + +func TestOpenLocalSQLiteAppliesPragmasAndPoolLimits(t *testing.T) { + path := filepath.Join(t.TempDir(), "local.db") + db, err := Open("file:"+path, "") + if err != nil { + t.Fatalf("open local sqlite: %v", err) + } + t.Cleanup(func() { _ = db.Close() }) + + if got := db.Stats().MaxOpenConnections; got != 1 { + t.Fatalf("max open conns: got %d, want 1", got) + } + + var mode string + if err := db.QueryRow(`PRAGMA journal_mode`).Scan(&mode); err != nil { + t.Fatalf("pragma journal_mode: %v", err) + } + if !strings.EqualFold(mode, "wal") { + t.Fatalf("journal_mode: got %q, want wal", mode) + } + + var timeout int + if err := db.QueryRow(`PRAGMA busy_timeout`).Scan(&timeout); err != nil { + t.Fatalf("pragma busy_timeout: %v", err) + } + if timeout != 5000 { + t.Fatalf("busy_timeout: got %d, want 5000", timeout) + } +} diff --git a/packages/api/internal/store/migrations/005_indexing_jobs.sql b/packages/api/internal/store/migrations/005_indexing_jobs.sql new file mode 100644 index 0000000..25cf968 --- /dev/null +++ b/packages/api/internal/store/migrations/005_indexing_jobs.sql @@ -0,0 +1,16 @@ +CREATE TABLE IF NOT EXISTS indexing_jobs ( + document_id TEXT PRIMARY KEY, + did TEXT NOT NULL, + collection TEXT NOT NULL, + rkey TEXT NOT NULL, + cid TEXT NOT NULL, + record_json TEXT NOT NULL, + status TEXT NOT NULL, + attempts INTEGER NOT NULL DEFAULT 0, + last_error TEXT, + scheduled_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_indexing_jobs_status_scheduled + ON indexing_jobs(status, scheduled_at, updated_at); diff --git a/packages/api/internal/store/sql_store.go b/packages/api/internal/store/sql_store.go index 2b3e569..41c9742 100644 --- a/packages/api/internal/store/sql_store.go +++ b/packages/api/internal/store/sql_store.go @@ -98,7 +98,7 @@ func (s *SQLStore) ListDocuments(ctx context.Context, filter DocumentFilter) ([] for rows.Next() { doc := &Document{} var ( - title, body, summary, repoDID, repoName, authorHandle sql.NullString + title, body, summary, repoDID, repoName, authorHandle sql.NullString tagsJSON, language, createdAt, updatedAt, webURL, deletedAt sql.NullString ) if err := rows.Scan( @@ -271,6 +271,122 @@ func (s *SQLStore) EnqueueEmbeddingJob(ctx context.Context, documentID string) e return nil } +func (s *SQLStore) EnqueueIndexingJob(ctx context.Context, input IndexingJobInput) error { + now := time.Now().UTC().Format(time.RFC3339) + _, err := s.db.ExecContext(ctx, ` + INSERT INTO indexing_jobs ( + document_id, did, collection, rkey, cid, record_json, + status, attempts, last_error, scheduled_at, updated_at + ) VALUES (?, ?, ?, ?, ?, ?, 'pending', 0, NULL, ?, ?) + ON CONFLICT(document_id) DO UPDATE SET + did = excluded.did, + collection = excluded.collection, + rkey = excluded.rkey, + cid = excluded.cid, + record_json = excluded.record_json, + status = 'pending', + last_error = NULL, + scheduled_at = excluded.scheduled_at, + updated_at = excluded.updated_at`, + input.DocumentID, input.DID, input.Collection, input.RKey, input.CID, input.RecordJSON, now, now, + ) + if err != nil { + return fmt.Errorf("enqueue indexing job: %w", err) + } + return nil +} + +func (s *SQLStore) ClaimIndexingJob(ctx context.Context) (*IndexingJob, error) { + now := time.Now().UTC().Format(time.RFC3339) + staleCutoff := time.Now().UTC().Add(-5 * time.Minute).Format(time.RFC3339) + + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return nil, fmt.Errorf("begin claim indexing job tx: %w", err) + } + defer tx.Rollback() + + var job IndexingJob + row := tx.QueryRowContext(ctx, ` + SELECT document_id, did, collection, rkey, cid, record_json, + attempts, status, COALESCE(last_error, ''), scheduled_at, updated_at + FROM indexing_jobs + WHERE (status = 'pending' AND scheduled_at <= ?) + OR (status = 'processing' AND updated_at <= ?) + ORDER BY scheduled_at ASC, updated_at ASC + LIMIT 1`, now, staleCutoff) + + err = row.Scan( + &job.DocumentID, + &job.DID, + &job.Collection, + &job.RKey, + &job.CID, + &job.RecordJSON, + &job.Attempts, + &job.Status, + &job.LastError, + &job.ScheduledAt, + &job.UpdatedAt, + ) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("select indexing job: %w", err) + } + + res, err := tx.ExecContext(ctx, ` + UPDATE indexing_jobs + SET status = 'processing', updated_at = ? + WHERE document_id = ?`, + now, job.DocumentID, + ) + if err != nil { + return nil, fmt.Errorf("mark indexing job processing: %w", err) + } + affected, err := res.RowsAffected() + if err != nil { + return nil, fmt.Errorf("rows affected claim indexing job: %w", err) + } + if affected == 0 { + return nil, nil + } + + if err := tx.Commit(); err != nil { + return nil, fmt.Errorf("commit claim indexing job tx: %w", err) + } + job.Status = "processing" + job.UpdatedAt = now + return &job, nil +} + +func (s *SQLStore) CompleteIndexingJob(ctx context.Context, documentID string) error { + _, err := s.db.ExecContext(ctx, `DELETE FROM indexing_jobs WHERE document_id = ?`, documentID) + if err != nil { + return fmt.Errorf("complete indexing job: %w", err) + } + return nil +} + +func (s *SQLStore) RetryIndexingJob(ctx context.Context, documentID string, nextScheduledAt string, lastError string) error { + now := time.Now().UTC().Format(time.RFC3339) + _, err := s.db.ExecContext(ctx, ` + UPDATE indexing_jobs + SET status = 'pending', + attempts = attempts + 1, + last_error = ?, + scheduled_at = ?, + updated_at = ? + WHERE document_id = ?`, + lastError, nextScheduledAt, now, documentID, + ) + if err != nil { + return fmt.Errorf("retry indexing job: %w", err) + } + return nil +} + func (s *SQLStore) GetFollowSubjects(ctx context.Context, did string) ([]string, error) { rows, err := s.db.QueryContext(ctx, ` SELECT DISTINCT repo_did @@ -351,7 +467,7 @@ func (s *SQLStore) Ping(ctx context.Context) error { func scanDocument(row *sql.Row) (*Document, error) { doc := &Document{} var ( - title, body, summary, repoDID, repoName, authorHandle sql.NullString + title, body, summary, repoDID, repoName, authorHandle sql.NullString tagsJSON, language, createdAt, updatedAt, webURL, deletedAt sql.NullString ) err := row.Scan( diff --git a/packages/api/internal/store/store.go b/packages/api/internal/store/store.go index bc15c9c..a967b6a 100644 --- a/packages/api/internal/store/store.go +++ b/packages/api/internal/store/store.go @@ -41,11 +41,39 @@ type RecordState struct { UpdatedAt string } +// IndexingJob stores queued read-through indexing work fetched through the API. +type IndexingJob struct { + DocumentID string + DID string + Collection string + RKey string + CID string + RecordJSON string + Attempts int + Status string + LastError string + ScheduledAt string + UpdatedAt string +} + +// IndexingJobInput is the payload used to enqueue or refresh an indexing job. +type IndexingJobInput struct { + DocumentID string + DID string + Collection string + RKey string + CID string + RecordJSON string +} + // DocumentFilter scopes a ListDocuments query to a subset of documents. type DocumentFilter struct { - Collection string // filter by collection NSID - DID string // filter by author DID - DocumentID string // filter to a single document by stable ID + // filter by collection NSID + Collection string + // filter by author DID + DID string + // filter to a single document by stable ID + DocumentID string } // Store is the persistence interface for Twister. @@ -61,6 +89,10 @@ 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 + EnqueueIndexingJob(ctx context.Context, input IndexingJobInput) error + ClaimIndexingJob(ctx context.Context) (*IndexingJob, error) + CompleteIndexingJob(ctx context.Context, documentID string) error + RetryIndexingJob(ctx context.Context, documentID string, nextScheduledAt string, lastError 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 8303dea..54aeef7 100644 --- a/packages/api/internal/store/store_test.go +++ b/packages/api/internal/store/store_test.go @@ -315,4 +315,68 @@ func TestIntegration(t *testing.T) { t.Fatalf("collaborators: got %#v", collaborators) } }) + + t.Run("indexing jobs enqueue claim retry complete", func(t *testing.T) { + job := store.IndexingJobInput{ + DocumentID: "did:plc:owner|sh.tangled.repo|repo1", + DID: "did:plc:owner", + Collection: "sh.tangled.repo", + RKey: "repo1", + CID: "cid-repo1", + RecordJSON: `{"name":"repo1"}`, + } + if err := st.EnqueueIndexingJob(ctx, job); err != nil { + t.Fatalf("enqueue indexing job: %v", err) + } + if err := st.EnqueueIndexingJob(ctx, job); err != nil { + t.Fatalf("enqueue indexing job second call: %v", err) + } + + claimed, err := st.ClaimIndexingJob(ctx) + if err != nil { + t.Fatalf("claim indexing job: %v", err) + } + if claimed == nil { + t.Fatal("expected claimed indexing job") + } + if claimed.DocumentID != job.DocumentID { + t.Fatalf("claimed document id: got %q want %q", claimed.DocumentID, job.DocumentID) + } + + if err := st.RetryIndexingJob(ctx, job.DocumentID, "9999-12-31T23:59:59Z", "boom"); err != nil { + t.Fatalf("retry indexing job: %v", err) + } + + none, err := st.ClaimIndexingJob(ctx) + if err != nil { + t.Fatalf("claim delayed indexing job: %v", err) + } + if none != nil { + t.Fatalf("expected no claim before schedule time, got %#v", none) + } + + if err := st.RetryIndexingJob(ctx, job.DocumentID, "1970-01-01T00:00:00Z", "retry-now"); err != nil { + t.Fatalf("retry indexing job now: %v", err) + } + + claimed, err = st.ClaimIndexingJob(ctx) + if err != nil { + t.Fatalf("claim retried indexing job: %v", err) + } + if claimed == nil { + t.Fatal("expected claimed retried indexing job") + } + + if err := st.CompleteIndexingJob(ctx, job.DocumentID); err != nil { + t.Fatalf("complete indexing job: %v", err) + } + + claimed, err = st.ClaimIndexingJob(ctx) + if err != nil { + t.Fatalf("claim after complete: %v", err) + } + if claimed != nil { + t.Fatalf("expected no job after complete, got %#v", claimed) + } + }) } diff --git a/packages/api/justfile b/packages/api/justfile index 3c50d3f..ae20541 100644 --- a/packages/api/justfile +++ b/packages/api/justfile @@ -5,11 +5,27 @@ ldflags := "-s -w -X main.version=" + version + " -X main.commit=" + commit build: CGO_ENABLED=0 go build -ldflags "{{ldflags}}" -o twister ./main.go -run-api: - go run -ldflags "{{ldflags}}" ./main.go api +# Run the API server. Usage: just run-api [mode], mode: local|remote (default local) +run-api mode="local": + if [ "{{mode}}" = "local" ]; then \ + go run -ldflags "{{ldflags}}" ./main.go api --local; \ + elif [ "{{mode}}" = "remote" ]; then \ + go run -ldflags "{{ldflags}}" ./main.go api; \ + else \ + echo "invalid mode '{{mode}}' (expected local or remote)" >&2; \ + exit 1; \ + fi -run-indexer: - go run -ldflags "{{ldflags}}" ./main.go indexer +# Run the indexer. Usage: just run-indexer [mode], mode: local|remote (default local) +run-indexer mode="local": + if [ "{{mode}}" = "local" ]; then \ + go run -ldflags "{{ldflags}}" ./main.go indexer --local; \ + elif [ "{{mode}}" = "remote" ]; then \ + go run -ldflags "{{ldflags}}" ./main.go indexer; \ + else \ + echo "invalid mode '{{mode}}' (expected local or remote)" >&2; \ + exit 1; \ + fi test: go test ./... diff --git a/scripts/api/README.md b/scripts/api/README.md new file mode 100644 index 0000000..8d74062 --- /dev/null +++ b/scripts/api/README.md @@ -0,0 +1,24 @@ +# Twister API Smoke Checks + +Python smoke checks for Twister API endpoints, managed with uv. + +## Usage + +From the repo root: + +```sh +# Run all +uv run --project scripts/api twister-api-smoke +# Run specific checks (healthz | readyz | search | documents | indexing | activity) +uv run --project scripts/api twister-api-smoke --check healthz +``` + +## Options + +- `--verbose` for detailed output of API responses (JSON) +- `--base-url` (or env `TWISTER_API_BASE_URL`, default `http://localhost:8080`) +- `--query` for search check (default `twisted`) +- `--document-id` for documents check +- `--actor-handle` for indexing check (default `desertthunder.dev`) +- `--repo-at-uri` for repo fixture indexing/search checks +- `--profile-at-uri` for profile fixture indexing/search checks diff --git a/scripts/api/pyproject.toml b/scripts/api/pyproject.toml new file mode 100644 index 0000000..4f8e45e --- /dev/null +++ b/scripts/api/pyproject.toml @@ -0,0 +1,20 @@ +[project] +name = "twister-api-smoke" +version = "0.1.0" +description = "Smoke checks for Twister API endpoints" +readme = "README.md" +requires-python = ">=3.11" +dependencies = [] + +[project.scripts] +twister-api-smoke = "twister_api_smoke.cli:main" + +[build-system] +requires = ["hatchling>=1.24.0"] +build-backend = "hatchling.build" + +[tool.uv] +package = true + +[tool.ruff] +line-length = 100 diff --git a/scripts/api/src/twister_api_smoke/__init__.py b/scripts/api/src/twister_api_smoke/__init__.py new file mode 100644 index 0000000..b3a4bb1 --- /dev/null +++ b/scripts/api/src/twister_api_smoke/__init__.py @@ -0,0 +1 @@ +"""Twister API smoke checks.""" diff --git a/scripts/api/src/twister_api_smoke/cli.py b/scripts/api/src/twister_api_smoke/cli.py new file mode 100644 index 0000000..601c3bb --- /dev/null +++ b/scripts/api/src/twister_api_smoke/cli.py @@ -0,0 +1,362 @@ +import argparse +import enum +import json +import os +import sys +import time +from dataclasses import dataclass +from http.client import HTTPConnection +from typing import Any, NoReturn +from urllib.error import HTTPError, URLError +from urllib.parse import quote, urlencode, urljoin, urlparse +from urllib.request import Request, urlopen + + +class ANSI(enum.StrEnum): + RED = "\033[31m" + GREEN = "\033[32m" + YELLOW = "\033[33m" + CYAN = "\033[36m" + MAGENTA = "\033[35m" + RESET = "\033[0m" + + def colorize(self, msg: str) -> str: + return f"{self.value}{msg}{ANSI.RESET.value}" + + +def echo(msg: str) -> None: + tag = ANSI.GREEN.colorize("[smoke]") + print(f"{tag} {msg}") + + +def fail(msg: str) -> NoReturn: + tag = ANSI.RED.colorize("[smoke] FAIL") + print(f"{tag}: {msg}", file=sys.stderr) + raise SystemExit(1) + + +@dataclass(frozen=True) +class Options: + base_url: str + check: str + query: str + document_id: str + actor_handle: str + repo_at_uri: str + profile_at_uri: str + verbose: bool + + +def http_get_status(url: str) -> int: + req = Request(url, method="GET") + try: + with urlopen(req, timeout=10) as resp: + return resp.status + except HTTPError as err: + return err.code + except URLError as err: + fail(f"request failed for {url}: {err}") + + +def http_get_json(url: str, params: dict[str, str] | None = None) -> Any: + if params: + query = urlencode(params) + sep = "&" if "?" in url else "?" + url = f"{url}{sep}{query}" + + req = Request(url, method="GET") + try: + with urlopen(req, timeout=15) as resp: + payload = resp.read().decode("utf-8") + return json.loads(payload) + except HTTPError as err: + body = err.read().decode("utf-8", errors="replace") + fail(f"{url} returned {err.code}: {body}") + except URLError as err: + fail(f"request failed for {url}: {err}") + except json.JSONDecodeError as err: + fail(f"invalid JSON from {url}: {err}") + + +def assert_status(url: str, expected: int) -> None: + actual = http_get_status(url) + if actual != expected: + fail(f"{url} returned {actual} (expected {expected})") + + +def at_uri_to_document_id(at_uri: str) -> str: + if not at_uri.startswith("at://"): + fail(f"invalid at uri: {at_uri}") + trimmed = at_uri[len("at://") :] + parts = trimmed.split("/", 2) + if len(parts) != 3 or not all(parts): + fail(f"invalid at uri: {at_uri}") + did, collection, rkey = parts + return f"{did}|{collection}|{rkey}" + + +def encode_document_id(document_id: str) -> str: + return quote(document_id, safe="") + + +def colorize_json(json_str: str, indent=2) -> str: + """Recursively colorize JSON string for terminal output, retaining indentation. + + CYAN: keys, GREEN: string values, YELLOW: numbers, MAGENTA: booleans/null + """ + + def colorize_and_indent(value: Any, level: int = 0) -> str: + indent_str = " " * (indent * level) + if isinstance(value, dict): + items = [] + for k, v in value.items(): + colored_key = ANSI.CYAN.colorize(json.dumps(k)) + colored_value = colorize_and_indent(v, level + 1) + items.append(f"{indent_str} {colored_key}: {colored_value}") + return "{\n" + ",\n".join(items) + f"\n{indent_str}}}" + elif isinstance(value, list): + items = [colorize_and_indent(v, level + 1) for v in value] + return "[\n" + ",\n".join(f"{indent_str} {i}" for i in items) + f"\n{indent_str}]" + elif isinstance(value, str): + return ANSI.GREEN.colorize(json.dumps(value)) + elif isinstance(value, (int, float)): + return ANSI.YELLOW.colorize(str(value)) + elif isinstance(value, bool) or value is None: + return ANSI.MAGENTA.colorize(str(value).lower()) + else: + return json.dumps(value) + + try: + parsed = json.loads(json_str) + return colorize_and_indent(parsed) + except json.JSONDecodeError: + return json_str + + +def maybe_log_json(opts: Options, label: str, payload: Any) -> None: + if not opts.verbose: + return + pretty = json.dumps(payload, indent=2, sort_keys=True) + echo(f"{label} JSON:\n{colorize_json(pretty)}") + + +def check_healthz(opts: Options) -> None: + echo("checking GET /healthz") + payload = http_get_json(urljoin(opts.base_url, "/healthz")) + maybe_log_json(opts, "healthz", payload) + echo("healthz ok") + + +def check_readyz(opts: Options) -> None: + echo("checking GET /readyz") + payload = http_get_json(urljoin(opts.base_url, "/readyz")) + maybe_log_json(opts, "readyz", payload) + echo("readyz ok") + + +def check_search(opts: Options) -> None: + repo_id = at_uri_to_document_id(opts.repo_at_uri) + profile_id = at_uri_to_document_id(opts.profile_at_uri) + + echo("checking search with indexed repo fixture") + repo_payload = http_get_json(urljoin(opts.base_url, "/search"), {"q": "twisted"}) + if not isinstance(repo_payload, dict): + fail("search response must be a JSON object") + repo_results = repo_payload.get("results") + if not isinstance(repo_results, list): + fail("search response is missing expected fields") + repo_ids = { + r.get("id") for r in repo_results if isinstance(r, dict) and isinstance(r.get("id"), str) + } + if repo_id not in repo_ids: + fail(f"repo fixture id not found in search results: {repo_id}") + maybe_log_json(opts, "search-repo", repo_payload) + + echo("checking search with indexed profile fixture") + profile_payload = http_get_json(urljoin(opts.base_url, "/search"), {"q": opts.actor_handle}) + if not isinstance(profile_payload, dict): + fail("profile search response must be a JSON object") + profile_results = profile_payload.get("results") + if not isinstance(profile_results, list): + fail("profile search response is missing expected fields") + profile_ids = { + r.get("id") for r in profile_results if isinstance(r, dict) and isinstance(r.get("id"), str) + } + if profile_id not in profile_ids: + fail(f"profile fixture id not found in search results: {profile_id}") + maybe_log_json(opts, "search-profile", profile_payload) + echo("search ok") + + +def check_documents(opts: Options) -> None: + document_id = opts.document_id + if not document_id: + echo("no document id provided, deriving first result from search") + payload = http_get_json(urljoin(opts.base_url, "/search"), {"q": opts.query}) + results = payload.get("results") if isinstance(payload, dict) else None + if not isinstance(results, list) or not results: + fail("document id required or search must return at least one result") + first = results[0] + if not isinstance(first, dict) or not isinstance(first.get("id"), str): + fail("search result missing document id") + document_id = first["id"] + + encoded = encode_document_id(document_id) + echo(f"checking GET /documents/{{id}} for {document_id}") + payload = http_get_json(urljoin(opts.base_url, f"/documents/{encoded}")) + maybe_log_json(opts, "documents", payload) + echo("documents ok") + + +def check_indexing(opts: Options) -> None: + repo_id = at_uri_to_document_id(opts.repo_at_uri) + profile_id = at_uri_to_document_id(opts.profile_at_uri) + repo_encoded = encode_document_id(repo_id) + profile_encoded = encode_document_id(profile_id) + + profile_url = urljoin(opts.base_url, f"/documents/{profile_encoded}") + repo_url = urljoin(opts.base_url, f"/documents/{repo_encoded}") + + echo(f"triggering read-through fetch via /actors/{opts.actor_handle}") + assert_status(urljoin(opts.base_url, f"/actors/{opts.actor_handle}"), 200) + + echo(f"triggering read-through fetch via /actors/{opts.actor_handle}/repos") + assert_status(urljoin(opts.base_url, f"/actors/{opts.actor_handle}/repos"), 200) + + echo(f"waiting for queued indexing of profile fixture {profile_id}") + for _ in range(30): + if http_get_status(profile_url) == 200: + if opts.verbose: + payload = http_get_json(profile_url) + maybe_log_json(opts, "indexing-profile", payload) + break + time.sleep(1) + else: + fail(f"profile fixture did not become available at /documents/{profile_encoded} within 30s") + + echo(f"waiting for queued indexing of repo fixture {repo_id}") + for _ in range(30): + if http_get_status(repo_url) == 200: + if opts.verbose: + payload = http_get_json(repo_url) + maybe_log_json(opts, "indexing-repo", payload) + echo("indexing ok") + return + time.sleep(1) + + fail(f"repo fixture did not become available at /documents/{repo_encoded} within 30s") + + +def check_activity(opts: Options) -> None: + echo("checking websocket handshake on /activity/stream") + parsed = urlparse(opts.base_url) + if parsed.scheme not in {"http", "https"}: + fail(f"unsupported base url scheme: {parsed.scheme}") + + host = parsed.hostname + if host is None: + fail(f"invalid base url: {opts.base_url}") + + port = parsed.port + if port is None: + port = 443 if parsed.scheme == "https" else 80 + + path = parsed.path.rstrip("/") + "/activity/stream?wantedCollections=sh.tangled.repo" + + conn = HTTPConnection(host, port, timeout=10) + try: + conn.putrequest("GET", path) + conn.putheader("Connection", "Upgrade") + conn.putheader("Upgrade", "websocket") + conn.putheader("Sec-WebSocket-Version", "13") + conn.putheader("Sec-WebSocket-Key", "dGhlIHNhbXBsZSBub25jZQ==") + conn.endheaders() + resp = conn.getresponse() + if resp.status != 101: + fail(f"/activity/stream websocket handshake returned {resp.status} (expected 101)") + echo("activity handshake ok") + finally: + conn.close() + + +CHECKS: dict[str, Any] = { + "healthz": check_healthz, + "readyz": check_readyz, + "search": check_search, + "documents": check_documents, + "indexing": check_indexing, + "activity": check_activity, +} + + +def parse_args(argv: list[str]) -> Options: + parser = argparse.ArgumentParser(description="Twister API smoke checks") + parser.add_argument( + "--base-url", + default=None, + help="Twister API base URL (default env TWISTER_API_BASE_URL or http://localhost:8080)", + ) + parser.add_argument( + "--check", + choices=["all", *CHECKS.keys()], + default="all", + help="Which check to run", + ) + parser.add_argument( + "--query", default="twisted", help="Search query for search/documents checks" + ) + parser.add_argument("--document-id", default="", help="Document ID for documents check") + parser.add_argument( + "--actor-handle", + default="desertthunder.dev", + help="Actor handle used to trigger read-through indexing", + ) + parser.add_argument( + "--repo-at-uri", + default="at://did:plc:xg2vq45muivyy3xwatcehspu/sh.tangled.repo/3mho6hukiei22", + help="Repo AT URI expected to be fetched and indexed by smoke checks", + ) + parser.add_argument( + "--profile-at-uri", + default="at://did:plc:xg2vq45muivyy3xwatcehspu/sh.tangled.actor.profile/self", + help="Profile AT URI expected to be fetched and indexed by smoke checks", + ) + parser.add_argument( + "--verbose", action="store_true", help="Print JSON payloads returned by smoke endpoints" + ) + + ns = parser.parse_args(argv) + base_url = ns.base_url or os.environ.get("TWISTER_API_BASE_URL", "http://localhost:8080") + return Options( + base_url=base_url, + check=ns.check, + query=ns.query, + document_id=ns.document_id, + actor_handle=ns.actor_handle, + repo_at_uri=ns.repo_at_uri, + profile_at_uri=ns.profile_at_uri, + verbose=ns.verbose, + ) + + +def main(argv: list[str] | None = None) -> int: + opts = parse_args(sys.argv[1:] if argv is None else argv) + if opts.check == "all": + for name in ( + "healthz", + "readyz", + "indexing", + "search", + "documents", + "activity", + ): + CHECKS[name](opts) + echo("all API smoke checks passed") + return 0 + + CHECKS[opts.check](opts) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/api/uv.lock b/scripts/api/uv.lock new file mode 100644 index 0000000..12c5d3c --- /dev/null +++ b/scripts/api/uv.lock @@ -0,0 +1,8 @@ +version = 1 +revision = 3 +requires-python = ">=3.11" + +[[package]] +name = "twister-api-smoke" +version = "0.1.0" +source = { editable = "." } -- 2.51.2