From 66c87842e2726555384b028d97e5343ac7b5b462 Mon Sep 17 00:00:00 2001 From: Andrey Oblivantsev Date: Thu, 13 Aug 2026 17:52:15 +0100 Subject: [PATCH] feat: in-process HTTP search; /get /stats /audit /ingest. (#13) bin/brain/serve.go (ladybug tags) calls internal/brain instead of exec. HTTP tests inject a fake API so CI stays cgo-free. ExecSearcher remains the fallback when the binary is built without system_ladybug. --- PLAN.md | 4 +- README.md | 4 +- bin/brain/serve.go | 17 ++-- bin/brain/serve_exec.go | 20 +++++ bin/serve.go | 2 +- docs/README.md | 3 +- internal/brain/http.go | 154 ++++++++++++++++++++++++++++++++ internal/brain/search.go | 34 +++---- internal/httpapi/server.go | 143 +++++++++++++++++++++-------- internal/httpapi/server_test.go | 62 +++++++++++++ 10 files changed, 380 insertions(+), 63 deletions(-) create mode 100644 bin/brain/serve_exec.go create mode 100644 internal/brain/http.go diff --git a/PLAN.md b/PLAN.md index 7a107ef..1c534fc 100644 --- a/PLAN.md +++ b/PLAN.md @@ -29,7 +29,7 @@ detective method: **a fact needs ≥2 independent sources or it is | D3 | web search | Vendored client; SearXNG URL is config. Optional Compose instance (sanitized settings). Do not run a second copy on a host that already has one. Empty/`throttled` ≠ “nothing exists”. | | D4 | embeddings | **model2vec** `minishlab/potion-multilingual-128M` instead of embeddinggemma. | | D5 | parser | **mistune** for MD → leaf extraction (duckdb-md documented as future optional SQL/export layer, not v1). | -| D6 | graph engine | **LadybugDB**. Go is the service (`bin/brain/search.go`, `internal/brain`); Python remains for index/write until the Go write path is safe. | +| D6 | graph engine | **LadybugDB**. Go is the service (`bin/brain/search.go`, `bin/brain/serve.go` in-process, `internal/brain`); Python remains for index/write until the Go write path is safe. | | D7 | db access | `db-yaml`/`psql-yq`-style, read-only, YAML out. OnlyOffice Postgres via SSH tunnel (`127.0.0.1:5433`). | | D8 | evidence | detective method: ≥2 independent sources or `(not confirmed)`. Auto-pair docker ps × compose × ssh-config × docs. | | D9 | facts/goal model | Who / What / How / Where / When + evidence + confidence on every edge. | @@ -57,7 +57,7 @@ detective method: **a fact needs ≥2 independent sources or it is brain/index.go rebuild FTS + HNSW (incl. --with-mail) brain/get.go stats.go eval.go watch.go brain/search.go deduction: facts → info → web-search - brain/serve.go HTTP API (internal/httpapi) + brain/serve.go HTTP API in-process (internal/httpapi + internal/brain) mail/import.go JSON → markdown (no brain write) markdown/import.go mistune leaves postgres/query.go read-only YAML (wraps bin/db/psql-yq) diff --git a/README.md b/README.md index c668225..ce46bad 100644 --- a/README.md +++ b/README.md @@ -121,8 +121,8 @@ bin/brain/search.go "invoice from last week" # same s `bin/{subject}/{method}.go` — self-describing: shebang on line 1, usage comment from line 2. Shared code in `internal/`. YAML default output, `--json` for -machines. Tests gate every commit. HTTP: `bin/brain/serve.go` (default search -binary `var/bin/brain-search`, not Python). +machines. Tests gate every commit. HTTP: `bin/brain/serve.go` calls +`internal/brain` in-process (`/health` `/search` `/get` `/stats` `/audit` `/ingest`). ## Development diff --git a/bin/brain/serve.go b/bin/brain/serve.go index 03be4c6..1f24dac 100755 --- a/bin/brain/serve.go +++ b/bin/brain/serve.go @@ -1,18 +1,20 @@ -//usr/bin/env go run -tags=brain_serve "$0" "$@"; exit -//go:build brain_serve +//usr/bin/env go run -tags=brain_serve,system_ladybug "$0" "$@"; exit +//go:build brain_serve && cgo && system_ladybug // -// bin/brain/serve.go - HTTP API for the 2dph brain. +// bin/brain/serve.go - HTTP API (in-process ladybug search). // // KB_ROOT=/path/to/2dph ./bin/brain/serve.go -// KB_SEARCH_CMD=... KB_WORKERS=4 KB_PORT=8630 ./bin/brain/serve.go +// KB_WORKERS=4 KB_PORT=8630 ./bin/brain/serve.go // -// Default search backend is var/bin/brain-search (Go), not Python. +// Needs CGO + libladybug (same as bin/brain/search.go). // NOTE: never run `gofmt -w` on this file — it breaks the shebang. package main import ( + "log" "os" + "github.com/eSlider/2dph/internal/brain" "github.com/eSlider/2dph/internal/httpapi" ) @@ -22,5 +24,8 @@ func main() { os.Setenv("KB_ROOT", wd) } } - httpapi.Run() + if err := brain.Ready(); err != nil { + log.Fatal(err) + } + httpapi.Run(brain.HTTP{}) } diff --git a/bin/brain/serve_exec.go b/bin/brain/serve_exec.go new file mode 100644 index 0000000..fb36c07 --- /dev/null +++ b/bin/brain/serve_exec.go @@ -0,0 +1,20 @@ +//go:build brain_serve && !system_ladybug +// +// Fallback serve when ladybug cgo is not in the build (CI / tags=brain_serve). +// Production shebang is serve.go (in-process). +package main + +import ( + "os" + + "github.com/eSlider/2dph/internal/httpapi" +) + +func main() { + if os.Getenv("KB_ROOT") == "" { + if wd, err := os.Getwd(); err == nil { + os.Setenv("KB_ROOT", wd) + } + } + httpapi.Run(nil) +} diff --git a/bin/serve.go b/bin/serve.go index 3feaa3d..8fbddad 100755 --- a/bin/serve.go +++ b/bin/serve.go @@ -18,5 +18,5 @@ func main() { os.Setenv("KB_ROOT", wd) } } - httpapi.Run() + httpapi.Run(nil) } diff --git a/docs/README.md b/docs/README.md index bc6923e..485e487 100644 --- a/docs/README.md +++ b/docs/README.md @@ -8,7 +8,8 @@ Brain/ops/eSlider stack. Facts need proof or they are - [design](design.md) — schema, deduction model, sources - [Gitea issues](https://git.produktor.io/eSlider/2dph/issues) — work board (origin) -Search: `bin/brain/search.go "query"` (HTTP: `bin/brain/serve.go`). `--hop` is +Search: `bin/brain/search.go "query"` (HTTP: `bin/brain/serve.go` — +`/health` `/search` `/get` `/stats` `/audit` `/ingest`). `--hop` is not a walk; the flag errors until File/FROM_FILE edges exist. Published docs live here and mirror the project state. diff --git a/internal/brain/http.go b/internal/brain/http.go new file mode 100644 index 0000000..fa153bc --- /dev/null +++ b/internal/brain/http.go @@ -0,0 +1,154 @@ +//go:build cgo && system_ladybug + +package brain + +import ( + "bytes" + "context" + "encoding/json" + "fmt" +) + +// Ready opens the Ladybug file for the life of the serve process. +func Ready() error { + return openBrain() +} + +// HTTP is the in-process API used by bin/brain/serve.go. +type HTTP struct{} + +func (HTTP) Search(_ context.Context, query string, limit int) ([]byte, error) { + hits, err := searchHits(query, "", "", limit) + if err != nil { + return nil, err + } + for i := range hits { + if hits[i].Text != "" { + runes := []rune(hits[i].Text) + if len(runes) > 280 { + runes = runes[:280] + } + hits[i].Snippet = string(runes) + } + } + var buf bytes.Buffer + enc := json.NewEncoder(&buf) + enc.SetEscapeHTML(false) + if err := enc.Encode(toJSONOut(hits, query, "")); err != nil { + return nil, err + } + return buf.Bytes(), nil +} + +func (HTTP) Get(_ context.Context, id string, body bool) ([]byte, error) { + if conn == nil { + return nil, fmt.Errorf("brain not open") + } + stmt, err := conn.Prepare( + "MATCH (l:Leaf {id:$id}) RETURN l.id, l.text, l.root, l.confidence, l.source, l.type", + ) + if err != nil { + return nil, err + } + defer stmt.Close() + res, err := conn.Execute(stmt, map[string]any{"id": id}) + if err != nil { + return nil, err + } + if !res.HasNext() { + return nil, fmt.Errorf("no leaf %s", id) + } + row, err := res.Next() + if err != nil { + return nil, err + } + vals, err := row.GetAsSlice() + if err != nil || len(vals) < 6 { + return nil, fmt.Errorf("leaf row") + } + out := map[string]any{ + "id": fmt.Sprint(vals[0]), + "root": fmt.Sprint(vals[2]), + "confidence": fmt.Sprint(vals[3]), + "source": fmt.Sprint(vals[4]), + "type": fmt.Sprint(vals[5]), + } + if body { + out["text"] = fmt.Sprint(vals[1]) + } + return json.Marshal(out) +} + +func (HTTP) Stats(context.Context) ([]byte, error) { + if conn == nil { + return nil, fmt.Errorf("brain not open") + } + res, err := conn.Query("MATCH (l:Leaf) RETURN l.root, count(*)") + if err != nil { + return nil, err + } + byRoot := map[string]int{} + total := 0 + for res.HasNext() { + row, err := res.Next() + if err != nil { + return nil, err + } + vals, err := row.GetAsSlice() + if err != nil || len(vals) < 2 { + continue + } + n := int(asInt(vals[1])) + byRoot[fmt.Sprint(vals[0])] = n + total += n + } + return json.Marshal(map[string]any{"total": total, "by_root": byRoot, "db": dbPath()}) +} + +func (HTTP) Audit(context.Context) ([]byte, error) { + if conn == nil { + return nil, fmt.Errorf("brain not open") + } + res, err := conn.Query("MATCH (l:Leaf) RETURN l.root, l.confidence, count(*)") + if err != nil { + return nil, err + } + var rows []map[string]any + for res.HasNext() { + row, err := res.Next() + if err != nil { + return nil, err + } + vals, err := row.GetAsSlice() + if err != nil || len(vals) < 3 { + continue + } + rows = append(rows, map[string]any{ + "root": fmt.Sprint(vals[0]), + "confidence": fmt.Sprint(vals[1]), + "count": asInt(vals[2]), + }) + } + return json.Marshal(map[string]any{"status": "ok", "by_confidence": rows}) +} + +func (HTTP) Ingest(context.Context) ([]byte, error) { + return json.Marshal(map[string]any{ + "mode": "rebuild", + "command": "bin/brain/index.go --rebuild", + "add": "v2", + }) +} + +func asInt(v any) int64 { + switch n := v.(type) { + case int64: + return n + case int: + return int64(n) + case float64: + return int64(n) + default: + return 0 + } +} diff --git a/internal/brain/search.go b/internal/brain/search.go index d150186..c0d9b38 100644 --- a/internal/brain/search.go +++ b/internal/brain/search.go @@ -51,25 +51,13 @@ func runSearch(args []string) int { } defer closeBrain() - emb, err := embedQuery(query) + hits, err := searchHits(query, root, repo, limit) if err != nil { - fmt.Fprintf(os.Stderr, "embed: %v\n", err) + fmt.Fprintf(os.Stderr, "search: %v\n", err) return 1 } - fts, err := queryFTS(query, limit*3) - if err != nil { - fmt.Fprintf(os.Stderr, "fts: %v\n", err) - return 1 - } - - var vec []Hit - if vec, err = queryVector(emb, limit*3); err != nil { - fmt.Fprintf(os.Stderr, "vec: %v\n", err) - } - - results := rank.RankAndFilter(fts, vec, root, repo, limit) - + results := hits for i := range results { if results[i].Text != "" { runes := []rune(results[i].Text) @@ -97,6 +85,22 @@ func runSearch(args []string) int { return 0 } +func searchHits(query, root, repo string, limit int) ([]Hit, error) { + emb, err := embedQuery(query) + if err != nil { + return nil, fmt.Errorf("embed: %w", err) + } + fts, err := queryFTS(query, limit*3) + if err != nil { + return nil, fmt.Errorf("fts: %w", err) + } + var vec []Hit + if vec, err = queryVector(emb, limit*3); err != nil { + fmt.Fprintf(os.Stderr, "vec: %v\n", err) + } + return rank.RankAndFilter(fts, vec, root, repo, limit), nil +} + func b2i(err error) int { if err != nil { return 1 diff --git a/internal/httpapi/server.go b/internal/httpapi/server.go index e190014..4aaee9b 100644 --- a/internal/httpapi/server.go +++ b/internal/httpapi/server.go @@ -1,10 +1,10 @@ -// Package server serves the 2dph brain over HTTP. +// Package httpapi serves the 2dph brain over HTTP. // // Async by design: every request runs on its own goroutine, and CPU-heavy -// searches are serialized through a bounded worker pool (a counting -// semaphore) so N requests can't spawn N search processes at once. +// searches are serialized through a bounded worker pool so N requests can't +// spawn N backends at once. // -// Used by bin/brain/serve.go. +// Used by bin/brain/serve.go. Tests inject a fake API (no exec, no ladybug). package httpapi import ( @@ -21,30 +21,45 @@ import ( "time" ) -type Searcher interface { +// API is the in-process brain surface. Production serve.go wires internal/brain. +type API interface { Search(ctx context.Context, query string, limit int) ([]byte, error) + Get(ctx context.Context, id string, body bool) ([]byte, error) + Stats(ctx context.Context) ([]byte, error) + Audit(ctx context.Context) ([]byte, error) + Ingest(ctx context.Context) ([]byte, error) } type Server struct { - searcher Searcher + api API semaphore chan struct{} } const defaultPort = 8630 -func NewServer(searcher Searcher, workers int) http.Handler { +var errUnimplemented = errors.New("not implemented") + +func NewServer(api API, workers int) http.Handler { return &Server{ - searcher: searcher, + api: api, semaphore: make(chan struct{}, workers), } } func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { - switch { - case r.URL.Path == "/health": + switch r.URL.Path { + case "/health": writeJSON(w, http.StatusOK, map[string]any{"status": "ok"}) - case r.URL.Path == "/search": + case "/search": s.handleSearch(w, r) + case "/get": + s.handleGet(w, r) + case "/stats": + s.handleJSON(w, r, s.api.Stats) + case "/audit": + s.handleJSON(w, r, s.api.Audit) + case "/ingest": + s.handleJSON(w, r, s.api.Ingest) default: writeJSON(w, http.StatusNotFound, map[string]any{"error": "not found"}) } @@ -65,19 +80,56 @@ func (s *Server) handleSearch(w http.ResponseWriter, r *http.Request) { } limit = n } - - // Worker pool: block until a slot frees, so burst concurrency still - // bounds memory (no unbounded python processes). - select { - case s.semaphore <- struct{}{}: - defer func() { <-s.semaphore }() - case <-r.Context().Done(): + if !s.acquire(w, r) { return } + defer s.release() + body, err := s.api.Search(r.Context(), q, limit) + writeAPI(w, body, err) +} - body, err := s.searcher.Search(r.Context(), q, limit) +func (s *Server) handleGet(w http.ResponseWriter, r *http.Request) { + id := strings.TrimSpace(r.URL.Query().Get("id")) + if id == "" { + writeJSON(w, http.StatusBadRequest, map[string]any{"error": "id required"}) + return + } + body := r.URL.Query().Get("body") == "1" || r.URL.Query().Get("body") == "true" + if !s.acquire(w, r) { + return + } + defer s.release() + out, err := s.api.Get(r.Context(), id, body) + writeAPI(w, out, err) +} + +func (s *Server) handleJSON(w http.ResponseWriter, r *http.Request, fn func(context.Context) ([]byte, error)) { + if !s.acquire(w, r) { + return + } + defer s.release() + body, err := fn(r.Context()) + writeAPI(w, body, err) +} + +func (s *Server) acquire(w http.ResponseWriter, r *http.Request) bool { + select { + case s.semaphore <- struct{}{}: + return true + case <-r.Context().Done(): + return false + } +} + +func (s *Server) release() { <-s.semaphore } + +func writeAPI(w http.ResponseWriter, body []byte, err error) { if err != nil { - writeJSON(w, http.StatusGatewayTimeout, map[string]any{"error": err.Error()}) + code := http.StatusBadGateway + if errors.Is(err, errUnimplemented) { + code = http.StatusNotImplemented + } + writeJSON(w, code, map[string]any{"error": err.Error()}) return } writeRaw(w, http.StatusOK, body) @@ -95,17 +147,20 @@ func writeRaw(w http.ResponseWriter, code int, body []byte) { w.Write(body) } -// brainSearcher shells out to the Go brain-search binary (not Python). -// A single search is bounded and short-lived; the worker pool keeps at most N live. -type brainSearcher struct { - cmdPath string - timeout time.Duration +// ExecSearcher shells out to var/bin/brain-search. Fallback when the serve +// binary is built without ladybug cgo (CI / tags=brain_serve only). +type ExecSearcher struct { + CmdPath string + Timeout time.Duration } -func (b *brainSearcher) Search(ctx context.Context, query string, limit int) ([]byte, error) { - ctx, cancel := context.WithTimeout(ctx, b.timeout) +func (b ExecSearcher) Search(ctx context.Context, query string, limit int) ([]byte, error) { + if b.Timeout == 0 { + b.Timeout = 60 * time.Second + } + ctx, cancel := context.WithTimeout(ctx, b.Timeout) defer cancel() - cmd := exec.CommandContext(ctx, b.cmdPath, "--json", "-n", strconv.Itoa(limit), query) + cmd := exec.CommandContext(ctx, b.CmdPath, "--json", "-n", strconv.Itoa(limit), query) out, err := cmd.Output() if err != nil { var exitErr *exec.ExitError @@ -117,6 +172,18 @@ func (b *brainSearcher) Search(ctx context.Context, query string, limit int) ([] return out, nil } +func (ExecSearcher) Get(context.Context, string, bool) ([]byte, error) { + return nil, errUnimplemented +} +func (ExecSearcher) Stats(context.Context) ([]byte, error) { return nil, errUnimplemented } +func (ExecSearcher) Audit(context.Context) ([]byte, error) { return nil, errUnimplemented } +func (ExecSearcher) Ingest(context.Context) ([]byte, error) { + return json.Marshal(map[string]any{ + "mode": "rebuild", + "command": "bin/brain/index.go --rebuild", + }) +} + func defaultSearchCmd(root string) string { if env := os.Getenv("KB_SEARCH_CMD"); env != "" { return env @@ -124,11 +191,7 @@ func defaultSearchCmd(root string) string { return filepath.Join(root, "var", "bin", "brain-search") } -// Run starts the HTTP server. Reads env: KB_SEARCH_CMD (default -// $KB_ROOT/var/bin/brain-search), KB_WORKERS (default 4), KB_PORT (default 8630). -func Run() { - root := os.Getenv("KB_ROOT") - searchPath := defaultSearchCmd(root) +func workersAndPort() (int, int) { workers := 4 if raw := os.Getenv("KB_WORKERS"); raw != "" { if n, err := strconv.Atoi(raw); err == nil && n > 0 { @@ -141,11 +204,19 @@ func Run() { port = n } } + return workers, port +} - searcher := &brainSearcher{cmdPath: searchPath, timeout: 60 * time.Second} - handler := NewServer(searcher, workers) +// Run starts the HTTP server with an injected API (in-process brain, or ExecSearcher). +func Run(api API) { + if api == nil { + root := os.Getenv("KB_ROOT") + api = ExecSearcher{CmdPath: defaultSearchCmd(root), Timeout: 60 * time.Second} + } + workers, port := workersAndPort() + handler := NewServer(api, workers) addr := "127.0.0.1:" + strconv.Itoa(port) - log.Printf("serve: %s (workers=%d cmd=%s)", addr, workers, searchPath) + log.Printf("serve: %s (workers=%d)", addr, workers) if err := http.ListenAndServe(addr, handler); err != nil { log.Fatal(err) } diff --git a/internal/httpapi/server_test.go b/internal/httpapi/server_test.go index 2e226a4..2ea3462 100644 --- a/internal/httpapi/server_test.go +++ b/internal/httpapi/server_test.go @@ -5,6 +5,7 @@ import ( "encoding/json" "net/http" "net/http/httptest" + "os" "strings" "sync" "sync/atomic" @@ -47,6 +48,26 @@ func (f *fakeSearcher) Search(ctx context.Context, query string, limit int) ([]b return []byte(`{"query":"` + query + `","count":0,"results":[]}`), nil } +func (f *fakeSearcher) Get(_ context.Context, id string, body bool) ([]byte, error) { + out := map[string]any{"id": id, "root": "info"} + if body { + out["text"] = "fake body" + } + return json.Marshal(out) +} + +func (f *fakeSearcher) Stats(context.Context) ([]byte, error) { + return []byte(`{"total":0,"by_root":{}}`), nil +} + +func (f *fakeSearcher) Audit(context.Context) ([]byte, error) { + return []byte(`{"status":"ok"}`), nil +} + +func (f *fakeSearcher) Ingest(context.Context) ([]byte, error) { + return []byte(`{"mode":"rebuild","command":"bin/brain/index.go --rebuild"}`), nil +} + func (f *fakeSearcher) count() int { f.mu.Lock() defer f.mu.Unlock() @@ -142,6 +163,47 @@ func TestSearchRejectsBadLimit(t *testing.T) { } } +func TestGetLeaf(t *testing.T) { + fs := &fakeSearcher{callback: func(q string, limit int) ([]byte, error) { + return []byte(`{}`), nil + }} + h := NewServer(fs, 1) + if code, _ := get(t, h, "/get"); code != http.StatusBadRequest { + t.Fatalf("missing id code = %d, want 400", code) + } + code, body := get(t, h, "/get?id=leaf-1&body=1") + if code != http.StatusOK { + t.Fatalf("get code = %d, want 200 body=%s", code, body) + } + if !strings.Contains(string(body), "leaf-1") { + t.Fatalf("get body %s missing id", body) + } +} + +func TestStatsAuditIngest(t *testing.T) { + h := NewServer(&fakeSearcher{}, 1) + for _, path := range []string{"/stats", "/audit", "/ingest"} { + code, body := get(t, h, path) + if code != http.StatusOK { + t.Fatalf("%s code = %d, want 200 (%s)", path, code, body) + } + if !json.Valid(body) { + t.Fatalf("%s body not json: %s", path, body) + } + } +} + +func TestHTTPPackageDoesNotExecPython(t *testing.T) { + raw, err := os.ReadFile("server.go") + if err != nil { + t.Fatal(err) + } + lower := strings.ToLower(string(raw)) + if strings.Contains(lower, "python3") || strings.Contains(lower, "bin/kb/search") { + t.Fatal("httpapi must not exec Python or bin/kb/search") + } +} + func TestDefaultSearchCmdIsBrainNotPython(t *testing.T) { t.Setenv("KB_SEARCH_CMD", "") cmd := defaultSearchCmd("/repo")