Compare commits

..
2 Commits
Author SHA1 Message Date
eSliderandGitHub 3f30052ea8 feat: in-process HTTP search; /get /stats /audit /ingest. (#13)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
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.
2026-08-13 17:52:15 +01:00
eSliderandGitHub c1ee920b0a feat: brain/index.go shebang; mail import is not a brain write (D14). (#12)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Commands live at bin/brain/{index,get,stats,eval,watch}.go and
bin/mail/import.go, bin/markdown/import.go, bin/postgres/query.go.
Python remains the Ladybug write worker. index_mail is a deprecation
shim that rebuilds via --with-mail.
2026-08-13 17:46:25 +01:00
39 changed files with 848 additions and 233 deletions
+11 -6
View File
@@ -36,11 +36,13 @@ PLAN.md decisions + execution + open questions
docs/ published docs docs/ published docs
skills/ in-project agent skills (vendored, no external links) skills/ in-project agent skills (vendored, no external links)
bin/ self-describing tools bin/{subject}/{method}.go (shebang) bin/ self-describing tools bin/{subject}/{method}.go (shebang)
bin/brain/ search.go, serve.go; libs in internal/brain and internal/httpapi bin/brain/ search.go serve.go index.go get.go stats.go eval.go watch.go
bin/chats/ sync.go import.go facts.go apply.go; libs in internal/chats bin/chats/ sync.go import.go facts.go apply.go; libs in internal/chats
bin/mail/ sync.go import.go (index_mail → brain/index.go)
bin/markdown/ import.go (mistune leafs)
bin/postgres/ query.go (read-only YAML)
internal/ shared Go (brain/rank is cgo-free; chats parsers too) internal/ shared Go (brain/rank is cgo-free; chats parsers too)
bin/watch/ corpus watcher (internal via bin/brain/watch later) bin/watch/ corpus watcher (used by bin/brain/watch.go)
bin/mail/ mail pipeline: sync (Go), import (md), index_mail (rebuild)
bin/tools/ vendored python libs behind bin/* (kblib, yamlout, websearch) bin/tools/ vendored python libs behind bin/* (kblib, yamlout, websearch)
bin/docker-entrypoint container entrypoint (brain index|search|serve|watch) bin/docker-entrypoint container entrypoint (brain index|search|serve|watch)
compose.yaml docker composition (root level, not docker/) compose.yaml docker composition (root level, not docker/)
@@ -54,8 +56,8 @@ var/ kb.lbug, var/mail/*, caches (gitignored)
```bash ```bash
bin/mail/sync.go --source onlyoffice,gmail --workers 8 --out var/mail # raw message.json + attachments bin/mail/sync.go --source onlyoffice,gmail --workers 8 --out var/mail # raw message.json + attachments
bin/mail/sync.go --source gmail --query 'from:example.com' --out var/mail # Gmail search (default in:inbox) bin/mail/sync.go --source gmail --query 'from:example.com' --out var/mail # Gmail search (default in:inbox)
bin/mail/import --from-raw var/mail # message.json → message.md (convert only) bin/mail/import.go --from-raw var/mail # message.json → message.md (convert only)
bin/mail/index_mail # rebuild brain incl. all mail (fresh DB) bin/brain/index.go --rebuild # rebuild brain incl. all mail (fresh DB)
``` ```
- `sync` (Go) downloads messages + attachments; Gmail uses paginated list + - `sync` (Go) downloads messages + attachments; Gmail uses paginated list +
@@ -64,7 +66,7 @@ bin/mail/index_mail # rebuil
`pdftotext -layout` fast path (~15ms); textless/scanned PDFs fall back to `pdftotext -layout` fast path (~15ms); textless/scanned PDFs fall back to
docling (isolated subprocess — its native onnx can segfault the parent). docling (isolated subprocess — its native onnx can segfault the parent).
Conversion never touches the brain DB (crash safety). Conversion never touches the brain DB (crash safety).
- `index_mail` always rebuilds from scratch (repo corpus + mail). Ladybug - `index_mail` is a deprecation shim for `bin/brain/index.go --rebuild`. Ladybug
corrupts its WAL when brand-new leafs are bulk-inserted while FTS/vector corrupts its WAL when brand-new leafs are bulk-inserted while FTS/vector
indexes exist; a fresh DB with indexes created last is the only safe path. indexes exist; a fresh DB with indexes created last is the only safe path.
Keep conversion + indexing separate so a conversion crash can't leave the Keep conversion + indexing separate so a conversion crash can't leave the
@@ -77,6 +79,9 @@ bin/facts/audit ["self"|"facts"|"info"|"stale"] # 2-source + staleness gate
bin/facts/crm [--dry-run] # proof person↔company/company↔project (ooCRM × corpus SoT) bin/facts/crm [--dry-run] # proof person↔company/company↔project (ooCRM × corpus SoT)
bin/kb/search "query" [--repo X] # deprecated wrapper → bin/brain/search.go bin/kb/search "query" [--repo X] # deprecated wrapper → bin/brain/search.go
bin/brain/search.go "query" [--root facts|info] # deduction search → YAML bin/brain/search.go "query" [--root facts|info] # deduction search → YAML
bin/brain/get.go <id> [--body]
bin/markdown/import.go [dir] # mistune leaves → YAML
bin/postgres/query.go --profile onlyoffice -c 'SELECT 1'
bin/md/tables # what the graph holds → YAML bin/md/tables # what the graph holds → YAML
bin/brain/deduce "question" # thinking wrapper bin/brain/deduce "question" # thinking wrapper
``` ```
+14 -9
View File
@@ -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”. | | 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. | | 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). | | 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`). | | 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. | | 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. | | D9 | facts/goal model | Who / What / How / Where / When + evidence + confidence on every edge. |
@@ -53,13 +53,17 @@ detective method: **a fact needs ≥2 independent sources or it is
bin/ bin/
facts/extract auto-pair 2 sources → lexicon yaml + graph facts/extract auto-pair 2 sources → lexicon yaml + graph
facts/audit ["self"|"facts"|"info"|"stale"] 2-source + staleness gate facts/audit ["self"|"facts"|"info"|"stale"] 2-source + staleness gate
kb/index build FTS + HNSW from corpus (Python, for now) kb/index Python write path (called by bin/brain/index.go)
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/search.go deduction: facts → info → web-search
kb/get kb/stats kb/eval brain/serve.go HTTP API in-process (internal/httpapi + internal/brain)
brain/serve.go HTTP API (internal/httpapi) mail/import.go JSON → markdown (no brain write)
markdown/import.go mistune leaves
postgres/query.go read-only YAML (wraps bin/db/psql-yq)
chats/sync.go import.go facts.go apply.go chats/sync.go import.go facts.go apply.go
(libs in internal/chats; no chats index) (libs in internal/chats; no chats index)
md/import md/select md/tables md/gaps (mistune) md/import (deprecated; bin/markdown/import.go)
brain/extract brain/audit brain/deduce (thinking wrapper) brain/extract brain/audit brain/deduce (thinking wrapper)
web/search (vendored) web/search (vendored)
db/psql-yq (vendored) db/psql-yq (vendored)
@@ -108,12 +112,13 @@ Common props on every node/edge: `root`, `confidence`, `evidence[]`, `how`,
1. `bin/mail/sync.go` (Go, 8 workers) — paginated Gmail/OnlyOffice download. 1. `bin/mail/sync.go` (Go, 8 workers) — paginated Gmail/OnlyOffice download.
Gmail attachments key off `body.attachmentId`, not MIME `partId`. Gmail attachments key off `body.attachmentId`, not MIME `partId`.
2. `bin/mail/import --from-raw` — message.json → message.md; PDFs via 2. `bin/mail/import.go --from-raw` — message.json → message.md; PDFs via
`pdftotext -layout` (~15ms) with docling subprocess fallback; ICS sidecars `pdftotext -layout` (~15ms) with docling subprocess fallback; ICS sidecars
Latin-1→UTF-8 normalized. Latin-1→UTF-8 normalized.
3. `bin/mail/index_mail` — fresh rebuild (repo corpus + mail) because ladybug 3. `bin/brain/index.go --rebuild` — fresh rebuild (repo corpus + mail) because ladybug
corrupts its WAL on bulk-insert into an already-indexed DB. Conversion and corrupts its WAL on bulk-insert into an already-indexed DB. Conversion and
indexing stay separate for crash safety. indexing stay separate for crash safety. `bin/mail/index_mail` is a
deprecation shim.
4. Result: 17,835 messages → 28,918 info leafs, FTS + HNSW healthy, searchable 4. Result: 17,835 messages → 28,918 info leafs, FTS + HNSW healthy, searchable
via `bin/brain/search.go`. via `bin/brain/search.go`.
@@ -125,7 +130,7 @@ Common props on every node/edge: `root`, `confidence`, `evidence[]`, `how`,
2. `go test ./internal/brain/rank` (cgo-free ranking + flag parser) 2. `go test ./internal/brain/rank` (cgo-free ranking + flag parser)
3. python -m unittest discover -s bin/tools (includes published-docs SoT) 3. python -m unittest discover -s bin/tools (includes published-docs SoT)
4. bin/facts/audit self (lexicon internal consistency) 4. bin/facts/audit self (lexicon internal consistency)
5. bin/kb/eval (recall@5 ≥ 0.95, gates index regressions) 5. bin/brain/eval.go (recall@5 ≥ 0.95, gates index regressions)
6. md-docs build/lint if docs tooling arrives. 6. md-docs build/lint if docs tooling arrives.
Feedback loop: every commit → PR → CI → green/gate → merge. Same discipline as Feedback loop: every commit → PR → CI → green/gate → merge. Same discipline as
+10 -10
View File
@@ -30,8 +30,8 @@ graph TB
subgraph dph["2dph tools"] subgraph dph["2dph tools"]
EX["bin/facts/extract<br/>2-source pairing"] EX["bin/facts/extract<br/>2-source pairing"]
AU["bin/facts/audit<br/>confidence + staleness"] AU["bin/facts/audit<br/>confidence + staleness"]
IDX["bin/kb/index<br/>chunk + embed"] IDX["bin/brain/index.go<br/>chunk + embed"]
MD["bin/md/import<br/>mistune leaves"] MD["bin/markdown/import.go<br/>mistune leaves"]
SR["bin/brain/search.go<br/>deduction"] SR["bin/brain/search.go<br/>deduction"]
end end
@@ -88,9 +88,9 @@ fact; conflicting sources or a single source → `hypothesis` → `(not confirme
bin/brain/search.go "Matrix federation over HTTPS" # facts → info → web-search bin/brain/search.go "Matrix federation over HTTPS" # facts → info → web-search
bin/brain/search.go "onlyoffice postgres" --root facts bin/brain/search.go "onlyoffice postgres" --root facts
bin/brain/search.go "where is cs-lexicon" --json | yq '.' bin/brain/search.go "where is cs-lexicon" --json | yq '.'
bin/kb/get <id> --body # full chunk on demand bin/brain/get.go <id> --body # full chunk on demand
bin/kb/stats # index health bin/brain/stats.go # index health
bin/kb/eval # recall@5 gate bin/brain/eval.go # recall@5 gate
``` ```
`--hop` is not implemented (needs File/FROM_FILE edges); the flag errors instead of walking. `bin/kb/search` is a deprecated wrapper around `bin/brain/search.go`. `--hop` is not implemented (needs File/FROM_FILE edges); the flag errors instead of walking. `bin/kb/search` is a deprecated wrapper around `bin/brain/search.go`.
@@ -99,8 +99,8 @@ Mail is a first-class corpus (retrievable through the same search):
```bash ```bash
bin/mail/sync.go --source onlyoffice,gmail --workers 8 --out var/mail # raw sync (Go) bin/mail/sync.go --source onlyoffice,gmail --workers 8 --out var/mail # raw sync (Go)
bin/mail/import --from-raw var/mail # JSON → markdown bin/mail/import.go --from-raw var/mail # JSON → markdown
bin/mail/index_mail # rebuild brain incl. mail bin/brain/index.go --rebuild # rebuild brain (incl. mail)
bin/brain/search.go "invoice from last week" # same search over mail leafs bin/brain/search.go "invoice from last week" # same search over mail leafs
``` ```
@@ -111,7 +111,7 @@ bin/brain/search.go "invoice from last week" # same s
readers. **Never `DROP INDEX` FTS/VECTOR** on Ladybug 0.19: DROP leaves readers. **Never `DROP INDEX` FTS/VECTOR** on Ladybug 0.19: DROP leaves
ghost catalog tables (`_0_Leaf_vec_UPPER`) so recreate fails while ghost catalog tables (`_0_Leaf_vec_UPPER`) so recreate fails while
`SHOW_INDEXES` omits HNSW. Fresh indexes = delete `var/kb.lbug` + `SHOW_INDEXES` omits HNSW. Fresh indexes = delete `var/kb.lbug` +
`bin/kb/index --rebuild`. Use `ensure_indexes()` after upserts. `bin/brain/index.go --rebuild`. Use `ensure_indexes()` after upserts.
- **model2vec** — `potion-multilingual-128M` static embeddings (256-dim), - **model2vec** — `potion-multilingual-128M` static embeddings (256-dim),
CPU-fast, deterministic, no Ollama runtime dependency. CPU-fast, deterministic, no Ollama runtime dependency.
- facts and info split semantically by `root` column but written inside the - facts and info split semantically by `root` column but written inside the
@@ -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 `bin/{subject}/{method}.go` — self-describing: shebang on line 1, usage comment
from line 2. Shared code in `internal/`. YAML default output, `--json` for from line 2. Shared code in `internal/`. YAML default output, `--json` for
machines. Tests gate every commit. HTTP: `bin/brain/serve.go` (default search machines. Tests gate every commit. HTTP: `bin/brain/serve.go` calls
binary `var/bin/brain-search`, not Python). `internal/brain` in-process (`/health` `/search` `/get` `/stats` `/audit` `/ingest`).
## Development ## Development
+2 -2
View File
@@ -1,3 +1,3 @@
// Commands in this directory are shebang mains (search.go). // Commands in this directory are shebang mains (search.go, serve.go, index.go,
// search.go is behind the system_ladybug build tag (cgo). // get.go, stats.go, eval.go, watch.go), each behind an exclusive build tag.
package main package main
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=brain_eval "$0" "$@"; exit
//go:build brain_eval
//
// bin/brain/eval.go - recall@5 gate.
//
// ./bin/brain/eval.go
// ./bin/brain/eval.go --json
//
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/cmdbin"
)
func main() {
os.Exit(cmdbin.ExecFile("bin/kb/eval", os.Args[1:]))
}
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=brain_get "$0" "$@"; exit
//go:build brain_get
//
// bin/brain/get.go - read one leaf by id.
//
// ./bin/brain/get.go <id>
// ./bin/brain/get.go <id> --body
//
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/cmdbin"
)
func main() {
os.Exit(cmdbin.ExecFile("bin/kb/get", os.Args[1:]))
}
+24
View File
@@ -0,0 +1,24 @@
//usr/bin/env go run -tags=brain_index "$0" "$@"; exit
//go:build brain_index
//
// bin/brain/index.go - rebuild the Ladybug graph (Python write path).
//
// ./bin/brain/index.go --rebuild
// ./bin/brain/index.go --rebuild --with-mail
// ./bin/brain/index.go --dry-run --with-mail
//
// v1 write is always a rebuild when mail is included (live FTS/HNSW + bulk
// insert corrupts Ladybug 0.19 WAL). `add` is v2.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/cmdbin"
)
func main() {
args := append([]string{"--with-mail"}, os.Args[1:]...)
os.Exit(cmdbin.ExecFile("bin/kb/index", args))
}
+11 -6
View File
@@ -1,18 +1,20 @@
//usr/bin/env go run -tags=brain_serve "$0" "$@"; exit //usr/bin/env go run -tags=brain_serve,system_ladybug "$0" "$@"; exit
//go:build brain_serve //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_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. // NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main package main
import ( import (
"log"
"os" "os"
"github.com/eSlider/2dph/internal/brain"
"github.com/eSlider/2dph/internal/httpapi" "github.com/eSlider/2dph/internal/httpapi"
) )
@@ -22,5 +24,8 @@ func main() {
os.Setenv("KB_ROOT", wd) os.Setenv("KB_ROOT", wd)
} }
} }
httpapi.Run() if err := brain.Ready(); err != nil {
log.Fatal(err)
}
httpapi.Run(brain.HTTP{})
} }
+20
View File
@@ -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)
}
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=brain_stats "$0" "$@"; exit
//go:build brain_stats
//
// bin/brain/stats.go - index health.
//
// ./bin/brain/stats.go
// ./bin/brain/stats.go --json
//
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/cmdbin"
)
func main() {
os.Exit(cmdbin.ExecFile("bin/kb/stats", os.Args[1:]))
}
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=brain_watch "$0" "$@"; exit
//go:build brain_watch
//
// bin/brain/watch.go - re-index when corpus files change.
//
// ./bin/brain/watch.go [dir...]
// KB_WATCH_INTERVAL=15 ./bin/brain/watch.go
//
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/bin/watch"
)
func main() {
watch.Run(os.Args[1:])
}
+5 -5
View File
@@ -2,10 +2,10 @@
# bin/docker-entrypoint - run 2dph tools inside the container. # bin/docker-entrypoint - run 2dph tools inside the container.
# #
# brain shell (default) # brain shell (default)
# brain search <q> bin/kb/search # brain search <q> bin/brain/search.go
# brain index bin/kb/index # brain index bin/kb/index --with-mail
# brain watch <dir> watchdog re-indexer (bin/kb/watch) # brain watch <dir> compiled /app/bin/watch (bin/brain/watch.go)
# brain serve async Go HTTP server (bin/serve) # brain serve compiled /app/bin/serve (bin/brain/serve.go)
# brain extract bin/facts/extract (docker×compose pairing) # brain extract bin/facts/extract (docker×compose pairing)
# brain audit bin/facts/audit # brain audit bin/facts/audit
# #
@@ -18,7 +18,7 @@ shift || true
case "$CMD" in case "$CMD" in
shell) exec bash ;; shell) exec bash ;;
search) exec "$KB_PY" /app/bin/kb/search "$@" ;; search) exec "$KB_PY" /app/bin/kb/search "$@" ;;
index) exec "$KB_PY" /app/bin/kb/index "$@" ;; index) exec "$KB_PY" /app/bin/kb/index --with-mail "$@" ;;
watch) exec /app/bin/watch "$@" ;; watch) exec /app/bin/watch "$@" ;;
serve) exec /app/bin/serve "$@" ;; serve) exec /app/bin/serve "$@" ;;
extract) exec "$KB_PY" /app/bin/facts/extract "$@" ;; extract) exec "$KB_PY" /app/bin/facts/extract "$@" ;;
+19 -3
View File
@@ -26,6 +26,7 @@ from kblib import ( # noqa: E402
open_readonly, stats, open_readonly, stats,
) )
from mdleaves import read_markdown, to_all, walk_markdown # noqa: E402 from mdleaves import read_markdown, to_all, walk_markdown # noqa: E402
from mailleafs import from_mail_root # noqa: E402
CORPUS_DEFAULTS = ["README.md", "PLAN.md", "AGENTS.md", "docs", "skills"] CORPUS_DEFAULTS = ["README.md", "PLAN.md", "AGENTS.md", "docs", "skills"]
@@ -99,6 +100,9 @@ def main(argv: list[str]) -> int:
p = argparse.ArgumentParser(description="build the 2dph brain index") p = argparse.ArgumentParser(description="build the 2dph brain index")
p.add_argument("--corpus", action="append", help="extra markdown dir/file to index (may repeat)") p.add_argument("--corpus", action="append", help="extra markdown dir/file to index (may repeat)")
p.add_argument("--rebuild", action="store_true", help="fresh db + indexes") p.add_argument("--rebuild", action="store_true", help="fresh db + indexes")
p.add_argument("--with-mail", action="store_true", help="include var/mail message.md leafs")
p.add_argument("--since", default="", help="with --with-mail, only messages dated >= YYYY-MM-DD")
p.add_argument("--dry-run", action="store_true", help="count leafs, write nothing")
p.add_argument( p.add_argument(
"--skip-indexes", "--skip-indexes",
action="store_true", action="store_true",
@@ -109,14 +113,26 @@ def main(argv: list[str]) -> int:
a = p.parse_args(argv) a = p.parse_args(argv)
from kblib import DB_PATH, VAR from kblib import DB_PATH, VAR
VAR.mkdir(exist_ok=True)
if a.rebuild and DB_PATH.exists():
DB_PATH.unlink()
leafs = load_corpus(ROOT) leafs = load_corpus(ROOT)
if a.corpus: if a.corpus:
for source in a.corpus: for source in a.corpus:
leafs.extend(load_corpus_glob(source)) leafs.extend(load_corpus_glob(source))
mail_n = 0
if a.with_mail:
mail = from_mail_root(ROOT / "var" / "mail", since=a.since)
mail_n = len(mail)
leafs.extend(mail)
if a.dry_run:
msg = {"indexed": 0, "corpus_total": len(leafs), "mail_leafs": mail_n, "dry_run": True}
print(json.dumps(msg, indent=2) if a.json else
f"brain/index: {len(leafs)} leafs would be indexed (mail={mail_n})")
return 0
VAR.mkdir(exist_ok=True)
if a.rebuild and DB_PATH.exists():
DB_PATH.unlink()
db, conn = connect(DB_PATH, read_only=False) db, conn = connect(DB_PATH, read_only=False)
init_schema(conn) init_schema(conn)
+1 -1
View File
@@ -1,5 +1,5 @@
//usr/bin/env go run "$0" "$@"; exit //usr/bin/env go run "$0" "$@"; exit
// bin/kb/watch.go - re-index the 2dph brain when corpus files change. // bin/kb/watch.go — deprecated. Use bin/brain/watch.go.
// //
// Usage: // Usage:
// //
+2 -2
View File
@@ -15,8 +15,8 @@ Writes one directory per message: var/mail/{folder}/{message_id}/
attachments/ raw attachment files (zips unpacked to _unpacked/) attachments/ raw attachment files (zips unpacked to _unpacked/)
attachments/*.md converted attachment content attachments/*.md converted attachment content
Indexing is a separate step (bin/mail/index_mail): conversion can crash in Indexing is a separate step (`bin/brain/index.go --rebuild`): conversion can
native docling and must not leave the brain DB mid-transaction. crash in native docling and must not leave the brain DB mid-transaction.
Requires ONLYOFFICE_URL/USER/PASS in .env (or env). Idempotent: a message Requires ONLYOFFICE_URL/USER/PASS in .env (or env). Idempotent: a message
already present (message.md exists) is skipped unless --force. already present (message.md exists) is skipped unless --force.
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=mail_import "$0" "$@"; exit
//go:build mail_import
//
// bin/mail/import.go - message.json → markdown (no brain write).
//
// ./bin/mail/import.go --from-raw var/mail
//
// Indexing is bin/brain/index.go --rebuild, not this command.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/cmdbin"
)
func main() {
os.Exit(cmdbin.ExecFile("bin/mail/import", os.Args[1:]))
}
+11 -120
View File
@@ -1,135 +1,26 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
"""mail/index_mail - rebuild the brain with every markdown under var/mail. """mail/index_mail — deprecated. Use bin/brain/index.go --rebuild --with-mail.
Ladybug corrupts its WAL when brand-new leafs are bulk-inserted while the Ladybug corrupts its WAL on bulk-insert into an already-indexed DB, so this
FTS/VECTOR indexes already exist, so indexing ALWAYS runs as a fresh rebuild shim always rebuilds (repo corpus + var/mail). Conversion stays in mail/import.
(repo corpus + var/mail), matching the proven-safe `kb/index --rebuild` path.
Conversion and indexing stay separate: conversion can crash in native docling
and must not leave the brain DB mid-transaction.
bin/mail/index_mail rebuild the index incl. all mail
bin/mail/index_mail --dry-run count without writing
bin/mail/index_mail --limit N cap messages included
bin/mail/index_mail --since D only messages dated >= D (YYYY-MM-DD)
""" """
from __future__ import annotations from __future__ import annotations
import argparse import os
import json
import sys import sys
from pathlib import Path from pathlib import Path
ROOT = Path(__file__).resolve().parents[2] ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(ROOT / "bin" / "tools"))
from kblib import DB_PATH, VAR, connect, ensure_indexes, init_schema, stats, upsert_leaf # noqa: E402
from mdleaves import read_markdown, to_all, walk_markdown # noqa: E402
def msg_date(md: Path) -> str:
j = md.parent / "message.json"
try:
d = json.loads(j.read_text(encoding="utf-8"))
return (d.get("receivedDate") or d.get("receivedAt") or "")[:10]
except Exception:
return ""
def mail_leafs(limit: int, since: str, repo: str = "ooMail") -> list[dict]:
root = ROOT / "var" / "mail"
mds = sorted(root.rglob("message.md"))
if since:
mds = [m for m in mds if msg_date(m) >= since]
if limit:
mds = mds[:limit]
leafs: list[dict] = []
for md in mds:
files = [md] + sorted((md.parent / "attachments").glob("*.md"))
for f in files:
if not f.exists():
continue
for lf in to_all(read_markdown(f), f, repo=repo):
lf["source"] = f"ooMail:{md.parent.name}:{f.name}"
lf["how"] = "mail/import"
leafs.append(lf)
return leafs
def main(argv: list[str]) -> int: def main(argv: list[str]) -> int:
p = argparse.ArgumentParser(description="rebuild the brain incl. all mail") print(
p.add_argument("--dry-run", action="store_true", help="count only, write nothing") "bin/mail/index_mail is deprecated; use bin/brain/index.go --rebuild --with-mail",
p.add_argument("--limit", type=int, default=0, help="cap messages included") file=sys.stderr,
p.add_argument("--since", default="", help="only messages dated >= YYYY-MM-DD") )
p.add_argument("--json", action="store_true") index = ROOT / "bin" / "kb" / "index"
a = p.parse_args(argv) os.execv(sys.executable, [sys.executable, str(index), "--rebuild", "--with-mail", *argv])
return 1
mail = mail_leafs(a.limit, a.since)
if a.dry_run:
print(f"mail/index_mail: {len(mail)} mail leafs would be indexed")
return 0
# Fresh rebuild: delete DB, index repo corpus + mail, create indexes once
# at the end. Never insert into an already-indexed DB (WAL corruption).
VAR.mkdir(exist_ok=True)
if DB_PATH.exists():
DB_PATH.unlink()
corpus = _load_corpus()
leafs = corpus + mail
db, conn = connect(DB_PATH, read_only=False)
init_schema(conn)
embed = _embedder()
done, total = _index_leafs(conn, leafs, embed)
ensure_indexes(conn)
s = stats(conn)
conn.close()
db.close()
result = {"indexed": done, "corpus_total": total, "mail_leafs": len(mail),
**{k: v for k, v in s.items() if k in ("total", "by_root")}}
print(json.dumps(result, indent=2) if a.json else
f"mail/index_mail: indexed {done}/{total} leafs (mail={len(mail)}); db total {s['total']}")
return 0
CORPUS_DEFAULTS = ["README.md", "PLAN.md", "AGENTS.md", "docs", "skills"]
def _load_corpus() -> list[dict]:
files: list[Path] = []
for entry in CORPUS_DEFAULTS:
p = ROOT / entry
if p.is_file():
files.append(p)
elif p.is_dir():
files.extend(walk_markdown(p))
leafs: list[dict] = []
for path in files:
try:
leafs.extend(to_all(read_markdown(path), path, repo="eSlider/2dph"))
except OSError as e:
print(f"mail/index_mail: skip {path}: {e}", file=sys.stderr)
return leafs
def _index_leafs(conn, leafs: list[dict], embed_fn) -> tuple[int, int]:
count = 0
for lf in leafs:
query = f"{lf['heading']}\n\n{lf['text']}"
emb = embed_fn(lf["text"]) if lf["text"] else None
upsert_leaf(conn, text=query, root="info", confidence="confirmed",
source=lf["source"], source_rev="mail" if lf.get("how") == "mail/import" else "working-tree",
how=lf.get("how", "kb/index"), loc=lf["source"], type_=lf.get("type", "reference"),
embedding=emb)
count += 1
return count, len(leafs)
def _embedder():
from model2vec import StaticModel
model = StaticModel.from_pretrained("minishlab/potion-multilingual-128M")
return lambda text: model.encode([text])[0].astype(float).tolist()
if __name__ == "__main__": if __name__ == "__main__":
+1 -1
View File
@@ -6,7 +6,7 @@
// ./bin/mail/sync.go --dry-run // ./bin/mail/sync.go --dry-run
// //
// Writes raw message.json + attachments under var/mail/<folder>/<id>/; run // Writes raw message.json + attachments under var/mail/<folder>/<id>/; run
// bin/mail/import --from-raw afterwards to convert everything to markdown. // bin/mail/import.go --from-raw afterwards to convert everything to markdown.
// //
// Shebang trick: first line is a Go `//` comment; the real code lives in the // Shebang trick: first line is a Go `//` comment; the real code lives in the
// importable package (module path, never a relative import). // importable package (module path, never a relative import).
+3
View File
@@ -0,0 +1,3 @@
// Commands in this directory are shebang mains (import.go), tagged so
// `go build ./bin/markdown` does not see two mains.
package main
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=markdown_import "$0" "$@"; exit
//go:build markdown_import
//
// bin/markdown/import.go - split markdown into leafs (mistune).
//
// ./bin/markdown/import.go [dir]
// ./bin/markdown/import.go --files a.md,b.md --json
//
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/cmdbin"
)
func main() {
os.Exit(cmdbin.ExecFile("bin/md/import", os.Args[1:]))
}
+2
View File
@@ -0,0 +1,2 @@
// Commands in this directory are shebang mains (query.go).
package main
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=postgres_query "$0" "$@"; exit
//go:build postgres_query
//
// bin/postgres/query.go - read-only Postgres as YAML.
//
// ./bin/postgres/query.go --profile onlyoffice -c 'SELECT 1'
//
// Profiles: $HOME/.config/brain/db-profiles.yml (credentials stay out of git).
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/cmdbin"
)
func main() {
os.Exit(cmdbin.ExecFile("bin/db/psql-yq", os.Args[1:]))
}
+1 -1
View File
@@ -18,5 +18,5 @@ func main() {
os.Setenv("KB_ROOT", wd) os.Setenv("KB_ROOT", wd)
} }
} }
httpapi.Run() httpapi.Run(nil)
} }
+4 -4
View File
@@ -135,7 +135,7 @@ def create_fts_and_vector(conn: ladybug.Connection, force: bool = False) -> None
`force=True` is accepted for API compatibility but does **not** drop. `force=True` is accepted for API compatibility but does **not** drop.
Fresh indexes require deleting `var/kb.lbug` and rebuilding Fresh indexes require deleting `var/kb.lbug` and rebuilding
(`bin/kb/index --rebuild`). (`bin/brain/index.go --rebuild`).
""" """
del force # API compat; DROP is unsafe — see docstring del force # API compat; DROP is unsafe — see docstring
names = leaf_index_names(conn) names = leaf_index_names(conn)
@@ -145,7 +145,7 @@ def create_fts_and_vector(conn: ladybug.Connection, force: bool = False) -> None
except Exception as e: except Exception as e:
raise RuntimeError( raise RuntimeError(
"CREATE_FTS_INDEX failed (often ghost catalog after DROP INDEX). " "CREATE_FTS_INDEX failed (often ghost catalog after DROP INDEX). "
"Delete var/kb.lbug and run bin/kb/index --rebuild. " "Delete var/kb.lbug and run bin/brain/index.go --rebuild. "
f"Cause: {e}" f"Cause: {e}"
) from e ) from e
if "Leaf_vec" not in names: if "Leaf_vec" not in names:
@@ -158,7 +158,7 @@ def create_fts_and_vector(conn: ladybug.Connection, force: bool = False) -> None
raise RuntimeError( raise RuntimeError(
"CREATE_VECTOR_INDEX failed (often ghost catalog after DROP INDEX " "CREATE_VECTOR_INDEX failed (often ghost catalog after DROP INDEX "
"Leaf.Leaf_vec → `_0_Leaf_vec_UPPER already exists in catalog`). " "Leaf.Leaf_vec → `_0_Leaf_vec_UPPER already exists in catalog`). "
"Delete var/kb.lbug and run bin/kb/index --rebuild. " "Delete var/kb.lbug and run bin/brain/index.go --rebuild. "
f"Cause: {e}" f"Cause: {e}"
) from e ) from e
names = leaf_index_names(conn) names = leaf_index_names(conn)
@@ -237,6 +237,6 @@ def stats(conn: ladybug.Connection) -> dict:
def open_readonly() -> tuple[ladybug.Database, ladybug.Connection]: def open_readonly() -> tuple[ladybug.Database, ladybug.Connection]:
if not DB_PATH.exists(): if not DB_PATH.exists():
raise FileNotFoundError(f"{DB_PATH} missing - run bin/kb/index first") raise FileNotFoundError(f"{DB_PATH} missing - run bin/brain/index.go --rebuild first")
db, conn = connect(read_only=True) db, conn = connect(read_only=True)
return db, conn return db, conn
+37
View File
@@ -0,0 +1,37 @@
"""Mail markdown under var/mail → info leafs. Conversion stays off the brain DB."""
from __future__ import annotations
import json
from pathlib import Path
from mdleaves import read_markdown, to_all
def msg_date(md: Path) -> str:
j = md.parent / "message.json"
try:
d = json.loads(j.read_text(encoding="utf-8"))
return (d.get("receivedDate") or d.get("receivedAt") or "")[:10]
except (OSError, json.JSONDecodeError, TypeError):
return ""
def from_mail_root(root: Path, limit: int = 0, since: str = "", repo: str = "ooMail") -> list[dict]:
if not root.is_dir():
return []
mds = sorted(root.rglob("message.md"))
if since:
mds = [m for m in mds if msg_date(m) >= since]
if limit:
mds = mds[:limit]
leafs: list[dict] = []
for md in mds:
files = [md] + sorted((md.parent / "attachments").glob("*.md"))
for f in files:
if not f.exists():
continue
for lf in to_all(read_markdown(f), f, repo=repo):
lf["source"] = f"ooMail:{md.parent.name}:{f.name}"
lf["how"] = "mail/import"
leafs.append(lf)
return leafs
+28
View File
@@ -63,3 +63,31 @@ class BinLayoutTest(unittest.TestCase):
(ROOT / "bin" / "chats" / "linkedin.go").exists(), (ROOT / "bin" / "chats" / "linkedin.go").exists(),
"parser must not stay under bin/chats as a second main", "parser must not stay under bin/chats as a second main",
) )
def _assert_shebang(self, rel: str) -> None:
p = ROOT / rel
self.assertTrue(p.is_file(), f"missing {rel}")
first = p.read_text().splitlines()[0]
self.assertTrue(
first.startswith("//usr/bin/env go run"),
f"{rel} shebang, got {first!r}",
)
def test_brain_methods_are_shebangs(self) -> None:
for method in ("index.go", "get.go", "stats.go", "eval.go", "watch.go"):
self._assert_shebang(f"bin/brain/{method}")
def test_mail_import_is_shebang_not_brain_write(self) -> None:
self._assert_shebang("bin/mail/import.go")
index_mail = (ROOT / "bin" / "mail" / "index_mail").read_text()
self.assertIn(
"bin/brain/index.go",
index_mail,
"index_mail must point at bin/brain/index.go",
)
def test_markdown_import_is_shebang(self) -> None:
self._assert_shebang("bin/markdown/import.go")
def test_postgres_query_is_shebang(self) -> None:
self._assert_shebang("bin/postgres/query.go")
+46
View File
@@ -0,0 +1,46 @@
"""Mail markdown → leafs (no Ladybug). Brain index --with-mail uses this."""
from __future__ import annotations
import json
import sys
import tempfile
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent))
import mailleafs # noqa: E402
class MailLeafsTest(unittest.TestCase):
def test_message_md_becomes_info_leaf(self) -> None:
root = Path(tempfile.mkdtemp())
msg = root / "inbox" / "alice-1"
msg.mkdir(parents=True)
(msg / "message.json").write_text(
json.dumps({"receivedDate": "2026-01-15T10:00:00Z", "subject": "Hello"}),
encoding="utf-8",
)
(msg / "message.md").write_text(
"---\nroot: info\n---\n\n# Hello\n\nFrom Alice to Bob.\n",
encoding="utf-8",
)
leafs = mailleafs.from_mail_root(root)
self.assertEqual(len(leafs), 1)
self.assertIn("Alice", leafs[0]["text"])
self.assertTrue(leafs[0]["source"].startswith("ooMail:"))
self.assertEqual(leafs[0]["how"], "mail/import")
def test_since_filters_by_message_json_date(self) -> None:
root = Path(tempfile.mkdtemp())
for name, day in (("old", "2025-01-01"), ("new", "2026-06-01")):
d = root / "inbox" / name
d.mkdir(parents=True)
(d / "message.json").write_text(
json.dumps({"receivedDate": f"{day}T00:00:00Z"}),
encoding="utf-8",
)
(d / "message.md").write_text(f"# {name}\n\nbody\n", encoding="utf-8")
leafs = mailleafs.from_mail_root(root, since="2026-01-01")
self.assertEqual(len(leafs), 1)
self.assertIn("new", leafs[0]["text"])
+9
View File
@@ -30,6 +30,15 @@ class PublishedDocsTest(unittest.TestCase):
"README deduction search must name bin/brain/search.go", "README deduction search must name bin/brain/search.go",
) )
def test_readme_index_is_brain_not_index_mail(self) -> None:
text = (ROOT / "README.md").read_text()
self.assertIn("bin/brain/index.go", text)
self.assertNotIn(
"bin/mail/index_mail",
text,
"mail index is a brain write; README must name bin/brain/index.go",
)
def test_docs_do_not_claim_hop_walks(self) -> None: def test_docs_do_not_claim_hop_walks(self) -> None:
paths = [ paths = [
ROOT / "README.md", ROOT / "README.md",
+4 -4
View File
@@ -1,4 +1,4 @@
// Package watch polls corpus directories for changes and re-runs bin/kb/index. // Package watch polls corpus directories for changes and re-runs brain/index.
// //
// Port of the former bin/kb-watch bash script to an importable, testable Go // Port of the former bin/kb-watch bash script to an importable, testable Go
// package. Polls file mtimes (no inotify deps); cheap and reliable. // package. Polls file mtimes (no inotify deps); cheap and reliable.
@@ -18,8 +18,8 @@ import (
type Options struct { type Options struct {
Dirs []string Dirs []string
Interval time.Duration Interval time.Duration
// IndexCmd is the kb/index command template. %s is replaced by the repo // IndexCmd is the index command template. %s is replaced by the repo
// root (from KB_ROOT). Defaults to `python3 <root>/bin/kb/index`. // root (from KB_ROOT). Defaults to `python3 <root>/bin/kb/index --with-mail`.
IndexCmd string IndexCmd string
} }
@@ -67,7 +67,7 @@ func fromEnv(args []string) Options {
if pys == "" { if pys == "" {
pys = "python3" pys = "python3"
} }
opts.IndexCmd = pys + " <root>/bin/kb/index" opts.IndexCmd = pys + " <root>/bin/kb/index --with-mail"
return opts return opts
} }
+6 -2
View File
@@ -3,6 +3,7 @@ package watch
import ( import (
"os" "os"
"path/filepath" "path/filepath"
"strings"
"testing" "testing"
"time" "time"
) )
@@ -43,7 +44,10 @@ func TestFromEnvDefaults(t *testing.T) {
if opts.Interval != 30*time.Second { if opts.Interval != 30*time.Second {
t.Fatalf("default interval = %s, want 30s", opts.Interval) t.Fatalf("default interval = %s, want 30s", opts.Interval)
} }
if opts.IndexCmd == "" { if !strings.Contains(opts.IndexCmd, "kb/index") {
t.Fatal("default index cmd is empty") t.Fatalf("default index cmd = %q, want kb/index", opts.IndexCmd)
}
if !strings.Contains(opts.IndexCmd, "--with-mail") {
t.Fatalf("default index cmd must include --with-mail, got %q", opts.IndexCmd)
} }
} }
+2 -1
View File
@@ -8,7 +8,8 @@ Brain/ops/eSlider stack. Facts need proof or they are
- [design](design.md) — schema, deduction model, sources - [design](design.md) — schema, deduction model, sources
- [Gitea issues](https://git.produktor.io/eSlider/2dph/issues) — work board (origin) - [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. not a walk; the flag errors until File/FROM_FILE edges exist.
Published docs live here and mirror the project state. Published docs live here and mirror the project state.
+154
View File
@@ -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
}
}
+19 -15
View File
@@ -51,25 +51,13 @@ func runSearch(args []string) int {
} }
defer closeBrain() defer closeBrain()
emb, err := embedQuery(query) hits, err := searchHits(query, root, repo, limit)
if err != nil { if err != nil {
fmt.Fprintf(os.Stderr, "embed: %v\n", err) fmt.Fprintf(os.Stderr, "search: %v\n", err)
return 1 return 1
} }
fts, err := queryFTS(query, limit*3) results := hits
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)
for i := range results { for i := range results {
if results[i].Text != "" { if results[i].Text != "" {
runes := []rune(results[i].Text) runes := []rune(results[i].Text)
@@ -97,6 +85,22 @@ func runSearch(args []string) int {
return 0 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 { func b2i(err error) int {
if err != nil { if err != nil {
return 1 return 1
+54
View File
@@ -0,0 +1,54 @@
package cmdbin
import (
"errors"
"os"
"os/exec"
"path/filepath"
"strings"
)
// Root is the 2dph checkout (KB_ROOT, or walk up for .git / var).
func Root() string {
if v := os.Getenv("KB_ROOT"); v != "" {
return v
}
wd, err := os.Getwd()
if err != nil {
return "."
}
for i := 0; i < 10; i++ {
if _, err := os.Stat(filepath.Join(wd, ".git")); err == nil {
return wd
}
if _, err := os.Stat(filepath.Join(wd, "var")); err == nil {
return wd
}
parent := filepath.Dir(wd)
if parent == wd {
break
}
wd = parent
}
return "."
}
// ExecFile runs repo-relative path (python/bash shebang scripts) with stdio.
func ExecFile(rel string, args []string) int {
path := filepath.Join(Root(), filepath.FromSlash(rel))
cmd := exec.Command(path, args...)
cmd.Stdin = os.Stdin
cmd.Stdout = os.Stdout
cmd.Stderr = os.Stderr
cmd.Dir = Root()
if err := cmd.Run(); err != nil {
if ee, ok := err.(*exec.ExitError); ok {
return ee.ExitCode()
}
if errors.Is(err, os.ErrNotExist) || strings.Contains(err.Error(), "no such file") {
return 127
}
return 1
}
return 0
}
+34
View File
@@ -0,0 +1,34 @@
package cmdbin
import (
"os"
"path/filepath"
"testing"
)
func TestRootHonorsKBROOT(t *testing.T) {
dir := t.TempDir()
t.Setenv("KB_ROOT", dir)
if got := Root(); got != dir {
t.Fatalf("Root() = %q, want %q", got, dir)
}
}
func TestExecFileMissingIs127(t *testing.T) {
t.Setenv("KB_ROOT", t.TempDir())
if code := ExecFile("no/such-tool", nil); code != 127 {
t.Fatalf("exit = %d, want 127", code)
}
}
func TestExecFileRuns(t *testing.T) {
root := t.TempDir()
script := filepath.Join(root, "echo.sh")
if err := os.WriteFile(script, []byte("#!/bin/sh\nexit 3\n"), 0o755); err != nil {
t.Fatal(err)
}
t.Setenv("KB_ROOT", root)
if code := ExecFile("echo.sh", nil); code != 3 {
t.Fatalf("exit = %d, want 3", code)
}
}
+107 -36
View File
@@ -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 // Async by design: every request runs on its own goroutine, and CPU-heavy
// searches are serialized through a bounded worker pool (a counting // searches are serialized through a bounded worker pool so N requests can't
// semaphore) so N requests can't spawn N search processes at once. // 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 package httpapi
import ( import (
@@ -21,30 +21,45 @@ import (
"time" "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) 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 { type Server struct {
searcher Searcher api API
semaphore chan struct{} semaphore chan struct{}
} }
const defaultPort = 8630 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{ return &Server{
searcher: searcher, api: api,
semaphore: make(chan struct{}, workers), semaphore: make(chan struct{}, workers),
} }
} }
func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
switch { switch r.URL.Path {
case r.URL.Path == "/health": case "/health":
writeJSON(w, http.StatusOK, map[string]any{"status": "ok"}) writeJSON(w, http.StatusOK, map[string]any{"status": "ok"})
case r.URL.Path == "/search": case "/search":
s.handleSearch(w, r) 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: default:
writeJSON(w, http.StatusNotFound, map[string]any{"error": "not found"}) 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 limit = n
} }
if !s.acquire(w, r) {
// 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():
return 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 { 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 return
} }
writeRaw(w, http.StatusOK, body) writeRaw(w, http.StatusOK, body)
@@ -95,17 +147,20 @@ func writeRaw(w http.ResponseWriter, code int, body []byte) {
w.Write(body) w.Write(body)
} }
// brainSearcher shells out to the Go brain-search binary (not Python). // ExecSearcher shells out to var/bin/brain-search. Fallback when the serve
// A single search is bounded and short-lived; the worker pool keeps at most N live. // binary is built without ladybug cgo (CI / tags=brain_serve only).
type brainSearcher struct { type ExecSearcher struct {
cmdPath string CmdPath string
timeout time.Duration Timeout time.Duration
} }
func (b *brainSearcher) Search(ctx context.Context, query string, limit int) ([]byte, error) { func (b ExecSearcher) Search(ctx context.Context, query string, limit int) ([]byte, error) {
ctx, cancel := context.WithTimeout(ctx, b.timeout) if b.Timeout == 0 {
b.Timeout = 60 * time.Second
}
ctx, cancel := context.WithTimeout(ctx, b.Timeout)
defer cancel() 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() out, err := cmd.Output()
if err != nil { if err != nil {
var exitErr *exec.ExitError var exitErr *exec.ExitError
@@ -117,6 +172,18 @@ func (b *brainSearcher) Search(ctx context.Context, query string, limit int) ([]
return out, nil 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 { func defaultSearchCmd(root string) string {
if env := os.Getenv("KB_SEARCH_CMD"); env != "" { if env := os.Getenv("KB_SEARCH_CMD"); env != "" {
return env return env
@@ -124,11 +191,7 @@ func defaultSearchCmd(root string) string {
return filepath.Join(root, "var", "bin", "brain-search") return filepath.Join(root, "var", "bin", "brain-search")
} }
// Run starts the HTTP server. Reads env: KB_SEARCH_CMD (default func workersAndPort() (int, int) {
// $KB_ROOT/var/bin/brain-search), KB_WORKERS (default 4), KB_PORT (default 8630).
func Run() {
root := os.Getenv("KB_ROOT")
searchPath := defaultSearchCmd(root)
workers := 4 workers := 4
if raw := os.Getenv("KB_WORKERS"); raw != "" { if raw := os.Getenv("KB_WORKERS"); raw != "" {
if n, err := strconv.Atoi(raw); err == nil && n > 0 { if n, err := strconv.Atoi(raw); err == nil && n > 0 {
@@ -141,11 +204,19 @@ func Run() {
port = n port = n
} }
} }
return workers, port
}
searcher := &brainSearcher{cmdPath: searchPath, timeout: 60 * time.Second} // Run starts the HTTP server with an injected API (in-process brain, or ExecSearcher).
handler := NewServer(searcher, workers) 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) 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 { if err := http.ListenAndServe(addr, handler); err != nil {
log.Fatal(err) log.Fatal(err)
} }
+62
View File
@@ -5,6 +5,7 @@ import (
"encoding/json" "encoding/json"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"os"
"strings" "strings"
"sync" "sync"
"sync/atomic" "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 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 { func (f *fakeSearcher) count() int {
f.mu.Lock() f.mu.Lock()
defer f.mu.Unlock() 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) { func TestDefaultSearchCmdIsBrainNotPython(t *testing.T) {
t.Setenv("KB_SEARCH_CMD", "") t.Setenv("KB_SEARCH_CMD", "")
cmd := defaultSearchCmd("/repo") cmd := defaultSearchCmd("/repo")
+1 -1
View File
@@ -30,7 +30,7 @@ related:
--- ---
``` ```
`bin/kb/index` reads this. `type` becomes a searchable column. `related:` is `bin/brain/index.go` reads this. `type` becomes a searchable column. `related:` is
frontmatter for humans; graph hops from it are not implemented yet. frontmatter for humans; graph hops from it are not implemented yet.
```bash ```bash
+4 -4
View File
@@ -26,9 +26,9 @@ second independent source when local roots cannot confirm. An answer is
bin/brain/search.go "Matrix federation" # pointers + snippets, YAML bin/brain/search.go "Matrix federation" # pointers + snippets, YAML
bin/brain/search.go "onlyoffice postgres" --root facts # restrict to confirmed bin/brain/search.go "onlyoffice postgres" --root facts # restrict to confirmed
bin/brain/search.go "where is cs-lexicon" --json | yq '.[].ref' bin/brain/search.go "where is cs-lexicon" --json | yq '.[].ref'
bin/kb/get <id> --body # full chunk only when needed bin/brain/get.go <id> --body # full chunk only when needed
bin/kb/stats # index health bin/brain/stats.go # index health
bin/kb/eval # recall@5 >= 0.95 gate bin/brain/eval.go # recall@5 >= 0.95 gate
``` ```
`bin/kb/search` is a deprecated wrapper. `--hop` errors (File/FROM_FILE edges `bin/kb/search` is a deprecated wrapper. `--hop` errors (File/FROM_FILE edges
@@ -39,7 +39,7 @@ are not wired yet); do not treat it as a graph walk.
- Search before you read. Never grep a repo for a concept the graph covers. - Search before you read. Never grep a repo for a concept the graph covers.
- `--root facts` returns only confirmed evidence-linked answers. Default shows - `--root facts` returns only confirmed evidence-linked answers. Default shows
facts first, then info leafs clearly marked `(not confirmed)`. facts first, then info leafs clearly marked `(not confirmed)`.
- If recall looks wrong, run `bin/kb/eval`; it gates control questions and - If recall looks wrong, run `bin/brain/eval.go`; it gates control questions and
should stay at or above 95% recall@5. should stay at or above 95% recall@5.
- Escalate to `web-search` (the `web-search` skill) as the independent second - Escalate to `web-search` (the `web-search` skill) as the independent second
source when both local roots cannot confirm; never report an unconfirmed source when both local roots cannot confirm; never report an unconfirmed