diff --git a/docs/api/tasks/phase-1-mvp.md b/docs/api/tasks/phase-1-mvp.md index 1dded17..03c7296 100644 --- a/docs/api/tasks/phase-1-mvp.md +++ b/docs/api/tasks/phase-1-mvp.md @@ -215,9 +215,9 @@ Expose a usable public search API backed by Turso's Tantivy-backed FTS. ### Tasks -- [ ] Set up HTTP server with net/http router -- [ ] Implement `/healthz` (always 200) and `/readyz` (SELECT 1 against DB) -- [ ] Implement search repository with FTS queries: +- [x] Set up HTTP server with net/http router +- [x] Implement `/healthz` (always 200) and `/readyz` (SELECT 1 against DB) +- [x] Implement search repository with FTS queries: ```sql SELECT id, title, summary, repo_name, author_handle, collection, record_type, @@ -231,21 +231,21 @@ Expose a usable public search API backed by Turso's Tantivy-backed FTS. LIMIT ? OFFSET ?; ``` -- [ ] Implement request validation: +- [x] Implement request validation: - `q` required, non-empty - `limit` 1–100, default 20 - `offset` >= 0, default 0 - Reject unknown parameters with 400 -- [ ] Implement filters (as WHERE clauses): +- [x] Implement filters (as WHERE clauses): - `collection` → `d.collection = ?` - `type` → `d.record_type = ?` - `author` → `d.author_handle = ?` or `d.did = ?` - `repo` → `d.repo_name = ?` -- [ ] Implement `/documents/{id}` — full document response -- [ ] Implement stable JSON response contract (see spec 05-search.md) -- [ ] Exclude tombstoned documents (`deleted_at IS NOT NULL`) by default -- [ ] Add request logging middleware (method, path, status, duration) -- [ ] Add CORS headers if needed +- [x] Implement `/documents/{id}` — full document response +- [x] Implement stable JSON response contract (see spec 05-search.md) +- [x] Exclude tombstoned documents (`deleted_at IS NOT NULL`) by default +- [x] Add request logging middleware (method, path, status, duration) +- [x] Add CORS headers if needed ### Verification diff --git a/packages/api/internal/api/api.go b/packages/api/internal/api/api.go index 778f64e..a3523b4 100644 --- a/packages/api/internal/api/api.go +++ b/packages/api/internal/api/api.go @@ -1 +1,322 @@ package api + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "net/http" + "strconv" + "time" + + "tangled.org/desertthunder.dev/twister/internal/config" + "tangled.org/desertthunder.dev/twister/internal/search" + "tangled.org/desertthunder.dev/twister/internal/store" +) + +// Server is the HTTP search API server. +type Server struct { + search *search.Repository + store store.Store + cfg *config.Config + log *slog.Logger +} + +// New creates a new API server. +func New(searchRepo *search.Repository, st store.Store, cfg *config.Config, log *slog.Logger) *Server { + return &Server{ + search: searchRepo, + store: st, + cfg: cfg, + log: log, + } +} + +// Handler returns the HTTP handler with all routes registered. +func (s *Server) Handler() http.Handler { + mux := http.NewServeMux() + + // Health + mux.HandleFunc("GET /healthz", s.handleHealthz) + mux.HandleFunc("GET /readyz", s.handleReadyz) + + // Search — M5 + mux.HandleFunc("GET /search", s.handleSearch) + mux.HandleFunc("GET /search/keyword", s.handleSearchKeyword) + + // Search — placeholders (Phase 2/3) + mux.HandleFunc("GET /search/semantic", s.handleNotImplemented) + mux.HandleFunc("GET /search/hybrid", s.handleNotImplemented) + + // Documents + mux.HandleFunc("GET /documents/{id}", s.handleGetDocument) + + // Admin — placeholders (M7) + if s.cfg.EnableAdminEndpoints { + mux.HandleFunc("POST /admin/reindex", s.handleNotImplemented) + mux.HandleFunc("POST /admin/reembed", s.handleNotImplemented) + } + + return s.withMiddleware(mux) +} + +// Run starts the HTTP server and blocks until ctx is cancelled. +func (s *Server) Run(ctx context.Context) error { + srv := &http.Server{ + Addr: s.cfg.HTTPBindAddr, + Handler: s.Handler(), + ReadHeaderTimeout: 10 * time.Second, + IdleTimeout: 60 * time.Second, + } + + errCh := make(chan error, 1) + go func() { + s.log.Info("listening", slog.String("addr", s.cfg.HTTPBindAddr)) + errCh <- srv.ListenAndServe() + }() + + select { + case err := <-errCh: + return err + case <-ctx.Done(): + shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + return srv.Shutdown(shutdownCtx) + } +} + +// --- Middleware --- + +func (s *Server) withMiddleware(next http.Handler) http.Handler { + return s.corsMiddleware(s.loggingMiddleware(next)) +} + +func (s *Server) loggingMiddleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + start := time.Now() + rw := &responseWriter{ResponseWriter: w, status: 200} + next.ServeHTTP(rw, r) + s.log.Info("request", + slog.String("method", r.Method), + slog.String("path", r.URL.Path), + slog.String("query", r.URL.RawQuery), + slog.Int("status", rw.status), + slog.Int64("duration_ms", time.Since(start).Milliseconds()), + ) + }) +} + +func (s *Server) corsMiddleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Access-Control-Allow-Origin", "*") + w.Header().Set("Access-Control-Allow-Methods", "GET, OPTIONS") + w.Header().Set("Access-Control-Allow-Headers", "Content-Type") + if r.Method == http.MethodOptions { + w.WriteHeader(http.StatusNoContent) + return + } + next.ServeHTTP(w, r) + }) +} + +type responseWriter struct { + http.ResponseWriter + status int +} + +func (rw *responseWriter) WriteHeader(code int) { + rw.status = code + rw.ResponseWriter.WriteHeader(code) +} + +// --- Health Handlers --- + +func (s *Server) handleHealthz(w http.ResponseWriter, _ *http.Request) { + writeJSON(w, http.StatusOK, map[string]string{"status": "ok"}) +} + +func (s *Server) handleReadyz(w http.ResponseWriter, r *http.Request) { + if err := s.search.Ping(r.Context()); err != nil { + s.log.Error("readyz: db unreachable", slog.String("error", err.Error())) + writeJSON(w, http.StatusServiceUnavailable, errorBody("db_unreachable", "database is not reachable")) + return + } + writeJSON(w, http.StatusOK, map[string]string{"status": "ready"}) +} + +// --- Search Handlers --- + +// knownSearchParams is the whitelist of accepted query parameters for search endpoints. +var knownSearchParams = map[string]bool{ + "q": true, "mode": true, "limit": true, "offset": true, + "collection": true, "type": true, "author": true, "repo": true, + "language": true, "from": true, "to": true, "state": true, +} + +func (s *Server) handleSearch(w http.ResponseWriter, r *http.Request) { + mode := r.URL.Query().Get("mode") + if mode == "" { + mode = s.cfg.SearchDefaultMode + } + switch mode { + case "keyword": + s.handleSearchKeyword(w, r) + case "semantic", "hybrid": + s.handleNotImplemented(w, r) + default: + writeJSON(w, http.StatusBadRequest, errorBody("invalid_parameter", "mode must be keyword, semantic, or hybrid")) + } +} + +func (s *Server) handleSearchKeyword(w http.ResponseWriter, r *http.Request) { + // Reject unknown parameters. + for key := range r.URL.Query() { + if !knownSearchParams[key] { + writeJSON(w, http.StatusBadRequest, errorBody("unknown_parameter", fmt.Sprintf("unknown parameter: %s", key))) + return + } + } + + q := r.URL.Query().Get("q") + if q == "" { + writeJSON(w, http.StatusBadRequest, errorBody("invalid_parameter", "q is required")) + return + } + + limit, err := intParam(r, "limit", s.cfg.SearchDefaultLimit) + if err != nil || limit < 1 || limit > s.cfg.SearchMaxLimit { + writeJSON(w, http.StatusBadRequest, errorBody("invalid_parameter", fmt.Sprintf("limit must be between 1 and %d", s.cfg.SearchMaxLimit))) + return + } + + offset, err := intParam(r, "offset", 0) + if err != nil || offset < 0 { + writeJSON(w, http.StatusBadRequest, errorBody("invalid_parameter", "offset must be >= 0")) + return + } + + params := search.Params{ + Query: q, + Limit: limit, + Offset: offset, + Collection: r.URL.Query().Get("collection"), + Type: r.URL.Query().Get("type"), + Author: r.URL.Query().Get("author"), + Repo: r.URL.Query().Get("repo"), + Language: r.URL.Query().Get("language"), + From: r.URL.Query().Get("from"), + To: r.URL.Query().Get("to"), + State: r.URL.Query().Get("state"), + } + + resp, err := s.search.Keyword(r.Context(), params) + if err != nil { + s.log.Error("search failed", slog.String("error", err.Error()), slog.String("query", q)) + writeJSON(w, http.StatusInternalServerError, errorBody("search_error", "search failed")) + return + } + + writeJSON(w, http.StatusOK, resp) +} + +// --- Document Handler --- + +func (s *Server) handleGetDocument(w http.ResponseWriter, r *http.Request) { + id := r.PathValue("id") + if id == "" { + writeJSON(w, http.StatusBadRequest, errorBody("invalid_parameter", "document id is required")) + return + } + + // Path value may be URL-encoded with | separators. The mux already decodes it, + // but callers may use pipe-encoded or slash-separated IDs; accept as-is. + doc, err := s.store.GetDocument(r.Context(), id) + if err != nil { + s.log.Error("get document failed", slog.String("error", err.Error()), slog.String("id", id)) + writeJSON(w, http.StatusInternalServerError, errorBody("db_error", "failed to fetch document")) + return + } + if doc == nil { + writeJSON(w, http.StatusNotFound, errorBody("not_found", "document not found")) + return + } + if doc.DeletedAt != "" { + writeJSON(w, http.StatusNotFound, errorBody("not_found", "document not found")) + return + } + + writeJSON(w, http.StatusOK, documentResponse(doc)) +} + +// --- Placeholder --- + +func (s *Server) handleNotImplemented(w http.ResponseWriter, _ *http.Request) { + writeJSON(w, http.StatusNotImplemented, errorBody("not_implemented", "this endpoint is not yet available")) +} + +// --- Helpers --- + +func writeJSON(w http.ResponseWriter, status int, v any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(v) +} + +func errorBody(code, message string) map[string]string { + return map[string]string{"error": code, "message": message} +} + +func intParam(r *http.Request, key string, def int) (int, error) { + v := r.URL.Query().Get(key) + if v == "" { + return def, nil + } + n, err := strconv.Atoi(v) + if err != nil { + return 0, err + } + return n, nil +} + +type documentJSON struct { + ID string `json:"id"` + DID string `json:"did"` + Collection string `json:"collection"` + RKey string `json:"rkey"` + ATURI string `json:"at_uri"` + CID string `json:"cid"` + RecordType string `json:"record_type"` + Title string `json:"title"` + Body string `json:"body"` + Summary string `json:"summary,omitempty"` + RepoName string `json:"repo_name,omitempty"` + AuthorHandle string `json:"author_handle,omitempty"` + TagsJSON string `json:"tags_json,omitempty"` + Language string `json:"language,omitempty"` + CreatedAt string `json:"created_at,omitempty"` + UpdatedAt string `json:"updated_at,omitempty"` + IndexedAt string `json:"indexed_at"` +} + +func documentResponse(doc *store.Document) documentJSON { + return documentJSON{ + ID: doc.ID, + DID: doc.DID, + Collection: doc.Collection, + RKey: doc.RKey, + ATURI: doc.ATURI, + CID: doc.CID, + RecordType: doc.RecordType, + Title: doc.Title, + Body: doc.Body, + Summary: doc.Summary, + RepoName: doc.RepoName, + AuthorHandle: doc.AuthorHandle, + TagsJSON: doc.TagsJSON, + Language: doc.Language, + CreatedAt: doc.CreatedAt, + UpdatedAt: doc.UpdatedAt, + IndexedAt: doc.IndexedAt, + } +} + diff --git a/packages/api/internal/search/search.go b/packages/api/internal/search/search.go index 5ed8515..73ad9c0 100644 --- a/packages/api/internal/search/search.go +++ b/packages/api/internal/search/search.go @@ -1 +1,184 @@ package search + +import ( + "context" + "database/sql" + "fmt" + "strings" +) + +// Params holds validated search query parameters. +type Params struct { + Query string + Limit int + Offset int + Collection string + Type string + Author string + Repo string + Language string + From string + To string + State string +} + +// Result is a single search hit. +type Result struct { + ID string `json:"id"` + Collection string `json:"collection"` + RecordType string `json:"record_type"` + Title string `json:"title"` + BodySnippet string `json:"body_snippet,omitempty"` + Summary string `json:"summary,omitempty"` + RepoName string `json:"repo_name,omitempty"` + AuthorHandle string `json:"author_handle,omitempty"` + DID string `json:"did"` + ATURI string `json:"at_uri"` + Score float64 `json:"score"` + MatchedBy []string `json:"matched_by"` + CreatedAt string `json:"created_at,omitempty"` + UpdatedAt string `json:"updated_at,omitempty"` +} + +// Response is the search API response envelope. +type Response struct { + Query string `json:"query"` + Mode string `json:"mode"` + Total int `json:"total"` + Limit int `json:"limit"` + Offset int `json:"offset"` + Results []Result `json:"results"` +} + +// Repository executes search queries against the database. +type Repository struct { + db *sql.DB +} + +// NewRepository creates a search repository backed by the given database. +func NewRepository(db *sql.DB) *Repository { + return &Repository{db: db} +} + +// Ping checks database connectivity. +func (r *Repository) Ping(ctx context.Context) error { + return r.db.PingContext(ctx) +} + +// Keyword runs a full-text keyword search. +func (r *Repository) Keyword(ctx context.Context, p Params) (*Response, error) { + // Build filter conditions beyond the base FTS match. + var filters []string + var filterArgs []any + + if p.Collection != "" { + filters = append(filters, "d.collection = ?") + filterArgs = append(filterArgs, p.Collection) + } + if p.Type != "" { + filters = append(filters, "d.record_type = ?") + filterArgs = append(filterArgs, p.Type) + } + if p.Author != "" { + filters = append(filters, "(d.author_handle = ? OR d.did = ?)") + filterArgs = append(filterArgs, p.Author, p.Author) + } + if p.Repo != "" { + filters = append(filters, "(d.repo_name = ? OR d.repo_did = ?)") + filterArgs = append(filterArgs, p.Repo, p.Repo) + } + if p.Language != "" { + filters = append(filters, "d.language = ?") + filterArgs = append(filterArgs, p.Language) + } + if p.From != "" { + filters = append(filters, "d.created_at >= ?") + filterArgs = append(filterArgs, p.From) + } + if p.To != "" { + filters = append(filters, "d.created_at <= ?") + filterArgs = append(filterArgs, p.To) + } + + // State filter requires a JOIN. + var join string + if p.State != "" { + join = "JOIN record_state rs ON rs.subject_uri = d.at_uri" + filters = append(filters, "rs.state = ?") + filterArgs = append(filterArgs, p.State) + } + + where := "fts_match(d.title, d.body, d.summary, d.repo_name, d.author_handle, d.tags_json, ?) AND d.deleted_at IS NULL" + if len(filters) > 0 { + where += " AND " + strings.Join(filters, " AND ") + } + + // Count total matching documents. + countSQL := fmt.Sprintf("SELECT COUNT(*) FROM documents d %s WHERE %s", join, where) + countArgs := append([]any{p.Query}, filterArgs...) + + var total int + if err := r.db.QueryRowContext(ctx, countSQL, countArgs...).Scan(&total); err != nil { + return nil, fmt.Errorf("count: %w", err) + } + + // Fetch results with score and snippet. + resultsSQL := fmt.Sprintf(` + SELECT d.id, d.title, d.summary, d.repo_name, d.author_handle, + d.did, d.at_uri, d.collection, d.record_type, d.created_at, d.updated_at, + fts_score(d.title, d.body, d.summary, d.repo_name, d.author_handle, d.tags_json, ?) AS score, + fts_highlight(d.body, '', '', ?) AS body_snippet + FROM documents d + %s + WHERE %s + ORDER BY score DESC + LIMIT ? OFFSET ?`, join, where) + + resultsArgs := make([]any, 0, 3+len(filterArgs)+2) + resultsArgs = append(resultsArgs, p.Query, p.Query, p.Query) // score, highlight, match + resultsArgs = append(resultsArgs, filterArgs...) + resultsArgs = append(resultsArgs, p.Limit, p.Offset) + + rows, err := r.db.QueryContext(ctx, resultsSQL, resultsArgs...) + if err != nil { + return nil, fmt.Errorf("search: %w", err) + } + defer rows.Close() + + results := make([]Result, 0) + for rows.Next() { + var res Result + var title, summary, repoName, authorHandle sql.NullString + var createdAt, updatedAt sql.NullString + var bodySnippet sql.NullString + + if err := rows.Scan( + &res.ID, &title, &summary, &repoName, &authorHandle, + &res.DID, &res.ATURI, &res.Collection, &res.RecordType, + &createdAt, &updatedAt, &res.Score, &bodySnippet, + ); err != nil { + return nil, fmt.Errorf("scan: %w", err) + } + res.Title = title.String + res.Summary = summary.String + res.RepoName = repoName.String + res.AuthorHandle = authorHandle.String + res.BodySnippet = bodySnippet.String + res.CreatedAt = createdAt.String + res.UpdatedAt = updatedAt.String + res.MatchedBy = []string{"keyword"} + results = append(results, res) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("rows: %w", err) + } + + return &Response{ + Query: p.Query, + Mode: "keyword", + Total: total, + Limit: p.Limit, + Offset: p.Offset, + Results: results, + }, nil +} diff --git a/packages/api/main.go b/packages/api/main.go index 8de8797..0c42122 100644 --- a/packages/api/main.go +++ b/packages/api/main.go @@ -10,11 +10,13 @@ import ( "time" "github.com/spf13/cobra" + "tangled.org/desertthunder.dev/twister/internal/api" "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" "tangled.org/desertthunder.dev/twister/internal/observability" + "tangled.org/desertthunder.dev/twister/internal/search" "tangled.org/desertthunder.dev/twister/internal/store" "tangled.org/desertthunder.dev/twister/internal/tapclient" ) @@ -62,8 +64,9 @@ func baseContext() (context.Context, context.CancelFunc) { func newAPICmd() *cobra.Command { return &cobra.Command{ - Use: "api", - Short: "Start the HTTP search API", + Use: "api", + Aliases: []string{"serve"}, + Short: "Start the HTTP search API", RunE: func(cmd *cobra.Command, args []string) error { cfg, err := config.Load() if err != nil { @@ -71,9 +74,28 @@ func newAPICmd() *cobra.Command { } log := observability.NewLogger(cfg) log.Info("starting api", slog.String("service", "api"), slog.String("version", version), slog.String("addr", cfg.HTTPBindAddr)) + + db, err := store.Open(cfg.TursoURL, cfg.TursoToken) + if err != nil { + return fmt.Errorf("open database: %w", err) + } + defer db.Close() + + if err := store.Migrate(db); err != nil { + return fmt.Errorf("migrate database: %w", err) + } + + st := store.New(db) + searchRepo := search.NewRepository(db) + srv := api.New(searchRepo, st, cfg, log) + ctx, cancel := baseContext() defer cancel() - <-ctx.Done() + + if err := srv.Run(ctx); err != nil { + return fmt.Errorf("run api: %w", err) + } + log.Info("shutting down api") return nil },