Compare commits

...
Author SHA1 Message Date
eSlider 2b8f946edb feat: D16 adjudication — 2v2 stays hypothesis until a rule fires.
Tests / Test (push) Skipped
Tests / OCR (tesseract fixture) (push) Skipped
Tests / Release (semver) (push) Skipped
Tests / Test (pull_request) Failing after 6s
Tests / OCR (tesseract fixture) (pull_request) Failing after 4s
Tests / Release (semver) (pull_request) Skipped
temporal_freshness then authority_pairing; unresolved keeps (not confirmed).
bin/facts/audit contradict. Gitea #29.
2026-08-14 11:40:16 +01:00
eSliderandGitHub a331042488 feat: DuckDB quantiles in-process (gcc CGO), not Ladybug. (#34)
Tests / Test (push) Failing after 6s
Tests / OCR (tesseract fixture) (push) Failing after 4s
Tests / Release (semver) (push) Skipped
OQ3: internal/duckstats + bin/qa/stats.go. Zig stays Ladybug-only.
mikefarah/yq for small structured slices. Gitea #30.
2026-08-14 11:32:23 +01:00
eSliderandGitHub f99dfea104 feat: OCR scans with tesseract, drop docling from the default path. (#33)
Tests / Test (push) Failing after 5s
Tests / OCR (tesseract fixture) (push) Failing after 4s
Tests / Release (semver) (push) Skipped
pdftotext still wins on born-digital PDFs. Empty text layers go through
pdftoppm + tesseract eng+deu. Optional OCR_ENGINE=paddle. Gitea #6.
2026-08-14 11:26:04 +01:00
eSliderandGitHub bae1494258 ci: recall@5 SoT is Zig bin/brain/eval.go (#19) (#32)
Tests / Test (push) Failing after 6s
Tests / Release (semver) (push) Skipped
* ci: recall@5 SoT is Zig bin/brain/eval.go.

GitHub Actions rebuilds the repo corpus and runs the Go eval gate
instead of Python bin/kb/eval. Audit self is a hard fail (Gitea #19).

* ci: point eval at the checkout and keep DevOps in the corpus.

/tmp/brain-eval cannot walk up to var/kb.lbug; KB_ROOT is required.
Recall fragments must exist in README/docs so a repo-only rebuild gates.
2026-08-14 11:03:39 +01:00
eSliderandGitHub c5be3f19be feat: index facts extract and chats markdown on rebuild. (#31)
Tests / Test (push) Failing after 5s
Tests / Release (semver) (push) Skipped
--with-facts / --facts-json write root=facts leafs; --with-chats
picks up var/chats/md as info. WhatsApp sync stays out of v1 (Gitea #18).
2026-08-14 10:57:40 +01:00
eSliderandGitHub a88dbb490c feat: search --hop walks FROM_FILE to Person. (#30)
Tests / Test (push) Failing after 4s
Tests / Release (semver) (push) Skipped
Parser no longer errors; hop 1 returns File, hop 3 reaches Person.
Rebuild writes Leaf-[:FROM_FILE]->File so the walk is not empty
on a fresh index (Gitea #17).
2026-08-14 10:51:35 +01:00
eSliderandGitHub 9d1a3f3c70 feat: write leafs incrementally without rebuilding the graph. (#29)
Tests / Test (push) Failing after 4s
Tests / Release (semver) (push) Skipped
Ladybug 0.19 stays FTS/HNSW queryable on MERGE of new ids; DROP INDEX
was the fatal path. bin/brain/add.go and POST /ingest land facts+info
in one transaction so watch/mail/git can become leafs now (Gitea #14).
2026-08-14 10:46:59 +01:00
eSliderandGitHub 8b5be9b659 docs: record v1 detective epic and remaining gaps. (#28)
Tests / Test (push) Failing after 4s
Tests / Release (semver) (push) Skipped
Read path and MCP are in; write, hops, corpus, and CI eval are still open. Board is Gitea epic #16.
2026-08-14 10:32:19 +01:00
eSliderandGitHub 1edb158f35 docs: portable runbook and Diataxis index (#27)
Tests / Test (push) Failing after 5s
Tests / Release (semver) (push) Skipped
Public face is a product, not a laptop tool. Run steps live in
docs/runbook.md; decisions D3/D6/D14/D15/D17/D18 stay in the docs index.
2026-08-13 23:31:41 +01:00
eSliderandGitHub 85caff90b9 feat: markdown import splits H2 leafs in Go (#26)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Drop the Python ExecFile wrapper. Conversion still does not write
Ladybug; add a dry-run fixture check for the index adapter split.
2026-08-13 23:27:42 +01:00
eSliderandGitHub b317968a4c feat: CPU reasoner bake-off for Qwen3.5-9B vs Bonsai (#25)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Measure OpenAI tool_calls (search/get/audit) on a CPU Ollama sidecar
instead of gating D18 on GPU or PicoClaw. Weights stay out of the image.
2026-08-13 23:19:51 +01:00
eSliderandGitHub 1f9bdb0bf6 feat: compile ladybug CGO with Zig, not gcc (#24)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
API image is Go-only (compose target api). Python rebuild is profile
index. bin/cgo/zig pins Zig 0.14.1, liblbug, and tokenizers.
2026-08-13 21:54:31 +01:00
eSliderandGitHub 0a05803f4b feat: PicoClaw compose profile exposes brain MCP on localhost (#23)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Agent is not shipped. docker compose --profile picoclaw up brain-mcp
and point the client at 127.0.0.1:8630/mcp.
2026-08-13 21:40:40 +01:00
eSliderandGitHub ad83e2a12f docs: PicoClaw fact-check tool order before a factual reply (#22)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
search then get then audit. throttled is not absence. 2dph is the
gate, not the agent loop.
2026-08-13 21:36:31 +01:00
eSliderandGitHub d894c6609f feat: generate brain tools skill from OpenAPI; rename db-yaml to postgres (#21)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Cursor skills stay in lockstep with serve handlers. CI checks every bin/
path named in SKILL.md exists.
2026-08-13 21:32:57 +01:00
eSliderandGitHub ff80359684 feat: OpenAPI and MCP from the same serve handlers (#20)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Agents get GET /openapi.json and POST /mcp. Tool names match
search/get/stats/audit paths so PicoClaw does not need shebangs.
2026-08-13 21:26:06 +01:00
eSliderandGitHub 36976d9b53 feat: facts audit/extract/crm shebang wrappers (D14) (#19)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Python stays the implementation. Commands are bin/facts/{audit,extract,crm}.go
like postgres/query.go.
2026-08-13 21:20:28 +01:00
eSliderandGitHub 8e6f67cc97 feat: brain get/stats/eval call internal/brain, not Python (#18)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Read path is cgo like search. Control questions live in rank so CI can
test them without ladybug. Python bin/kb/{get,stats,eval} stays the
runner fallback.
2026-08-13 21:16:48 +01:00
eSliderandGitHub aca05626bd feat: escalate brain search to web when facts cannot confirm (#17)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
2026-08-13 20:24:30 +01:00
eSliderandGitHub 39ae2abe8d feat: Go SearXNG client; throttled is not absence (#16)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
2026-08-13 19:53:31 +01:00
eSliderandGitHub ba5cc3a6e2 feat: read git history with go-git, not the git binary (#15)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
2026-08-13 18:07:56 +01:00
eSliderandGitHub 15d59054ff docs: delete agent-cost; rename kb-search skill to brain (#14)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
* docs: delete agent-cost; rename kb-search skill to brain.

bin/agents/cost does not exist. CI unittest now fails if a SKILL.md names a
missing bin/ path.

* test: gate SKILL.md bin/ paths; name the brain skill brain.

Follow-up to the agent-cost delete: unittest fails if a skill names a missing
tool. Frontmatter name is brain, not kb-search.
2026-08-13 17:55:41 +01:00
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
eSliderandGitHub cec0161ff6 refactor: chats method shebangs; drop chats index (D14). (#11)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Parsers and commands live in internal/chats. bin/chats/{sync,import,facts,apply}.go
are tagged shebang mains. Brain ingest is not a chats command.
2026-08-13 17:31:03 +01:00
eSliderandGitHub 46310f8773 docs: name bin/brain/search.go; --hop is not a graph walk. (#10)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Published docs and skills still taught bin/kb/search --hop 1. Search lives
at bin/brain/search.go; --hop errors until File edges exist. A unittest
gates the SoT so the lie cannot return.
2026-08-13 17:23:40 +01:00
eSliderandGitHub 7e0f3c9e06 feat: bin/brain/serve.go; search backend is Go not Python (#9)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
* feat(brain): HTTP serve from bin/brain/serve.go, default Go search binary.

Move the HTTP package to internal/httpapi. Default backend is
var/bin/brain-search, not Python. bin/serve.go stays as a deprecation shim.

* feat(httpapi): default search backend is var/bin/brain-search.

bin/brain/serve.go is the command; bin/serve.go stays as a tagged
deprecation shim. Tests fail if the default path still names Python.
2026-08-13 14:32:21 +01:00
eSliderandGitHub f14025304e refactor: one Go module; brain search in bin/brain + internal/brain. (#8)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
Collapse nested kbsearch/chats go.mod into the root module. Ranking stays
cgo-free under internal/brain/rank so CI does not need ladybug. bin/kb/search
is a deprecation wrapper that still sets CGO and builds the binary.
2026-08-13 14:26:54 +01:00
eSliderandGitHub dd6d7e9395 docs: point issues at Gitea origin (D15). (#7)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
GitHub stays the public clone for PRs and Actions. Work board is
https://git.produktor.io/eSlider/2dph/issues.
2026-08-13 14:11:00 +01:00
eSliderandGitHub 68d478224f feat(chats): parse LinkedIn MCP v4.22 inbox/conversation blobs. (#6)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
get_inbox/get_conversation return a sections+references envelope, not a
message list. Parser is covered by synthetic Alice/Bob fixtures; CI now
runs the nested bin/chats tests. Session check no longer launches Chromium.
2026-08-13 12:25:51 +01:00
eSliderandGitHub 117f3c2cfd fix(kbsearch): rank FTS correctly, filter before -n, start the daemon. (#5)
Go search took worst BM25 hits (ORDER BY score), cut to -n before --root,
and never called ensureDaemon. Ranking and flag parsing move to a cgo-free
package so CI can fail those regressions without ladybug. --hop errors
instead of being swallowed into the query.
2026-08-13 12:19:26 +01:00
eSliderandGitHub 140d86a4b9 Add Gmail --query to mail/sync (default in:inbox) (#4)
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
* Add --query to Gmail mail/sync instead of always listing in:inbox.

Callers keep the search string; default remains in:inbox.

* Document Gmail --query on the mail/sync pipeline.

* test(mail): assert Gmail --query reaches ListIDs, not only the CLI flag.

ParseCLI coverage left a hole: an empty query still has to become in:inbox
and a custom q has to be the string the client lists with.
2026-08-13 12:19:22 +01:00
eSlider c96c393a4a feat(chats): LinkedIn source — MCP client via get_inbox + get_conversation
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
- LinkedInMCPSource: MCP JSON-RPC, как TelegramMCPSource
- sync linkedin --limit N: выгрузка сообщений из LinkedIn
- Проверка сессии: uvx mcp-server-linkedin --status
- Вывод инструкции если сессия истекла
- JSONL в var/chats/linkedin/<thread_id>/messages.jsonl
2026-08-13 00:10:42 +01:00
eSlider 3d0d95cf00 docs: add edelweiss to GitHub safety rules
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
2026-08-13 00:07:01 +01:00
eSlider ed28fdbd2a chore: remove edelweiss references from public repo
Tests / Test (push) Skipped
Tests / Release (semver) (push) Skipped
2026-08-13 00:06:49 +01:00
eSlider 7e511d5b78 docs: GitHub safety rules — no absolute paths, PII, secrets, curasoft 2026-08-13 00:02:09 +01:00
eSlider 0d26519fab fix: resolve plan.md conflict, remove remaining /mnt/ paths 2026-08-13 00:00:48 +01:00
eSlider a429b823e5 chore: clean absolute paths, curasoft refs, secrets from history
- bin/chats/: env-based paths, no /mnt/ /home/ hardcodes
- bin/edelweiss-pilot: remove curasoft, use DOCS_BASE env var
- bin/facts/crm: use KNOWLEDGE_MESH_SEED env var
- compose.edelweiss.yml: remove curasoft volumes, use DOCS_BASE
- docs/chat-import-plan.md: link to Gitea issue, no secrets
- bin/seed-edelweiss-facts.py: removed (curasoft-only)
2026-08-13 00:00:15 +01:00
eSlider 6847233183 bin/chats: Phase 1 MVP — Telegram sync/import/index/facts/apply
- bin/chats/ — nested Go module (как bin/kbsearch/)
  - sync telegram — MCP JSON-RPC клиент, 31 личный чат, 922 сообщения
  - import — конвертация JSONL → MD с YAML frontmatter
  - index — делегирует bin/kb/index --corpus (132 leafs в brain)
  - facts — regex extraction phone/email/linkedin с валидацией
    (исключены: даты, суммы, номера карт, инвойсы)
  - apply — oo CLI cross-check + dry-run
- Source interface для будущих WhatsApp/LinkedIn
- 4 system tests (import, facts, empty, roundtrip) — синтетические данные
- bin/chat — build+exec wrapper
- docs/chat-import-plan.md — прогресс, пути к env (без секретов)

Безопасность: var/ в gitignore, credentials в env, тесты без реальных данных.
2026-08-12 23:59:29 +01:00
eSlider 6d7638ab73 docs: chat import pipeline plan — link to Gitea issue #1 2026-08-12 23:59:20 +01:00
eSliderandCursor bd1a91dab7 fix(kb): seed facts before CREATE indexes (FTS MERGE corruption)
Upsert under live FTS raises "document for node offset N is missing".
Add --skip-indexes; edelweiss-pilot index = write → seed → ensure_indexes.
Ship seed-edelweiss-facts.py (paired lexicon/OO/interview/QEMU facts).

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-12 16:09:02 +01:00
eSliderandCursor a5a1f91d95 fix(kb): stop DROP INDEX killing HNSW via Ladybug ghost catalog
Ladybug 0.19 DROP INDEX leaves `_0_Leaf_vec_UPPER` / `0_id_docs` in catalog so
CREATE fails while SHOW_INDEXES omits the index; create_fts_and_vector used to
swallow that. Never drop FTS/VECTOR; ensure_indexes after upserts; rebuild =
delete kb.lbug. Add compose.edelweiss.yml + regression tests.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-12 16:05:14 +01:00
eSlider 80e3b7a1cf Remove curasoft references, rename to detective method
- PLAN.md: replace 'curasoft-detective' with 'detective method'
- README.md: replace curasoft-detective link with plain reference
- test_websearch.py: fix test domain from ticket.curasoft.de to example.com
- Rewrote git history with git-filter-repo to remove all traces
2026-08-12 13:46:54 +01:00
eSlider ef4189c72d kbsearch: Go implementation with daemon model serving
- New nested module bin/kbsearch with Go implementation of bin/kb/search
- Embedding model (potion-multilingual-128M) served by localhost daemon
  so repeated CLI calls reuse the loaded model
- Bash launcher bin/kb/search builds binary on first run, caches to var/bin/
- Hybrid FTS + vector search (RRF k=60) matching Python kblib behavior
- YAML output via port of yamlout.py (ordered keys, same format)
- JSON output with proper field order
- All flags: --root, --repo, -n, --json, --list-model
- Root go.mod reverted to 1.25.0 (kbsearch is isolated nested module)
- CI passes: go test ./... and go vet ./... unaffected by kbsearch
2026-08-11 23:57:39 +01:00
eSlider 60c20ed98d feat(mail): full Gmail+OnlyOffice sync, import, and brain indexing
- bin/mail/sync.go: async Go sync engine (8 workers, paginated Gmail via
  API + OnlyOffice IMAP); Gmail attachments key off body.attachmentId, not
  MIME partId; ICS sidecars Latin-1->UTF-8 normalized (TestICSToMarkdownNormalizesLatin1)
- bin/mail/import: message.json -> markdown; PDFs via pdftotext -layout
  fast path with docling subprocess fallback for the ~5% textless files
- bin/mail/index_mail: fresh-rebuild indexer (repo corpus + mail) avoiding
  ladybug WAL corruption on bulk-insert into indexed DBs; split from import
- bin/kb/index: keep FTS/VECTOR indexes across incremental runs (drop+recreate
  leaves stale backing tables killing the vector index)
- docs: README/PLAN/AGENTS cover the mail pipeline

Result: 17,835 messages -> 28,918 info leafs, FTS+HNSW healthy.
2026-08-11 21:57:38 +01:00
eSlider 9f22380e82 refactor(tools): bin/{subject}/{method} layout; Go serve+watch modules
Move serve/ (module) -> bin/server, tools/ -> bin/tools, replace bin/kb-watch
bash with bin/watch Go package; self-executing Go shebangs bin/serve.go and
bin/kb/watch.go; Docker + CI + git/import + docs repointed. Multi-stage image
builds static serve+watch binaries (no Go runtime in container).
2026-08-11 09:52:20 +01:00
eSlider e2eff3b9c7 feat(kb): CRM association proof via oo, fix ssh-tunnel self-ref + oo creds
- bin/facts/crm: prove person<->company/company<->project against ooCRM
  x corpus SoT (knowledge-mesh-seed.yaml), write 78 facts (root=facts)
- tools/crmfacts.py + test_crm_facts.py: parser under unit tests (26 pass)
- docs/crm-associations-proof.md: provable graph, mistakes, fixes
- oo merge 759->763 resolves duplicate GoldenRatio.Exchange legal entity
- bin/db/ssh-tunnel: "$0" self-check + accept-new/BatchMode ssh flags
- AGENTS.md: document bin/facts/crm
2026-08-10 23:22:34 +01:00
198 changed files with 17529 additions and 1025 deletions
+1 -1
View File
@@ -3,6 +3,7 @@
var var
.git .git
.github .github
lib-ladybug
__pycache__ __pycache__
*.pyc *.pyc
*.lbug *.lbug
@@ -10,5 +11,4 @@ __pycache__
.cache .cache
.secrets .secrets
.skills-tmp .skills-tmp
serve/serve
docs/.build docs/.build
+47 -11
View File
@@ -19,6 +19,10 @@ jobs:
with: with:
fetch-depth: 0 fetch-depth: 0
- uses: actions/setup-go@v5
with:
go-version-file: go.mod
- name: Install uv - name: Install uv
uses: astral-sh/setup-uv@v6 uses: astral-sh/setup-uv@v6
with: with:
@@ -31,31 +35,63 @@ jobs:
run: | run: |
bash -n bin/db/psql-yq bash -n bin/db/psql-yq
bash -n bin/db/ssh-tunnel bash -n bin/db/ssh-tunnel
bash -n bin/kb-watch
bash -n bin/docker-entrypoint bash -n bin/docker-entrypoint
bash -n bin/kb/search
bash -n bin/cgo/zig
sh -n bin/cgo/zcc
sh -n bin/cgo/zc++
- name: Python unit tests (offline, vendored tools) - name: Python unit tests (offline, vendored tools)
run: | run: |
uv run python -m unittest discover -s tools -t . uv run python -m unittest discover -s bin/tools -t .
- name: Go serve tests (async, goroutine-bounded) - name: Go tests (root module; duckdb-go CGO via gcc, no ladybug)
run: | run: |
go vet ./... CC=gcc CXX=g++ CGO_CFLAGS= CGO_LDFLAGS= go vet ./...
go test ./... -count=1 CC=gcc CXX=g++ CGO_CFLAGS= CGO_LDFLAGS= go test ./... -count=1
working-directory: serve
- name: brain ranking tests (no cgo / no ladybug)
run: go test ./internal/brain/rank -count=1
- name: facts/audit self (lexicon consistency, no network) - name: facts/audit self (lexicon consistency, no network)
run: | run: ./bin/facts/audit self
./bin/facts/audit self 2>/dev/null || echo "audit: not yet implemented; gate skipped"
- name: kb/eval recall gate - name: CGO via Zig (compile brain/search + eval)
run: | run: |
./bin/kb/eval 2>/dev/null || echo "eval: not yet implemented; gate skipped" chmod +x bin/cgo/zig bin/cgo/zcc bin/cgo/zc++
bin/cgo/zig go build -tags system_ladybug -o /tmp/brain-search ./bin/brain/search.go
bin/cgo/zig go build -tags 'system_ladybug,brain_eval' -o /tmp/brain-eval ./bin/brain/eval.go
- uses: actions/cache@v4
with:
path: ~/.cache/huggingface
key: ${{ runner.os }}-hf-potion-multilingual-128M
- name: recall@5 SoT (Zig bin/brain/eval.go)
run: |
uv run python bin/kb/index --rebuild --json
KB_ROOT="$PWD" /tmp/brain-eval --json
ocr:
name: OCR (tesseract fixture)
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version-file: go.mod
- name: Install tesseract + poppler
run: |
sudo apt-get update
sudo apt-get install -y --no-install-recommends \
tesseract-ocr tesseract-ocr-eng tesseract-ocr-deu poppler-utils
- name: Go OCR tests (synthetic HELLO PNG)
run: go test ./internal/ocr -count=1
release: release:
name: Release (semver) name: Release (semver)
if: github.event_name == 'push' && github.ref == 'refs/heads/main' if: github.event_name == 'push' && github.ref == 'refs/heads/main'
needs: test needs: [test, ocr]
runs-on: ubuntu-latest runs-on: ubuntu-latest
permissions: permissions:
contents: write contents: write
+5
View File
@@ -9,3 +9,8 @@ __pycache__/
*.env *.env
.env .env
.secrets/ .secrets/
lib-ladybug/
go.work.local
models/
# Purged from git history. Do not re-add.
docs/crm-associations-proof.md
+80 -12
View File
@@ -3,7 +3,8 @@
Evidence-first brain over the ops/eSlider stack. Facts need proof or they are Evidence-first brain over the ops/eSlider stack. Facts need proof or they are
`(not confirmed)`. `(not confirmed)`.
Read first: [PLAN](PLAN.md) → [docs](docs/). Read first: [PLAN](PLAN.md) → [docs](docs/) → [roadmap](docs/roadmap.md)
(epic [#16](https://git.produktor.io/eSlider/2dph/issues/16)).
## Method (detective, no fork) ## Method (detective, no fork)
@@ -15,6 +16,9 @@ Read first: [PLAN](PLAN.md) → [docs](docs/).
- `info` root = descriptive/narrative leafs, searchable, never asserted as fact. - `info` root = descriptive/narrative leafs, searchable, never asserted as fact.
- Search is deduction: `facts``info``web-search` (second independent - Search is deduction: `facts``info``web-search` (second independent
source). An answer is `confirmed` only if it comes off the facts root. source). An answer is `confirmed` only if it comes off the facts root.
- Fact-check every *claim* (facts → info → live → web), not every edit or
syntax tweak. PicoClaw: `search` then `get` then `audit` before a factual
reply (`skills/picoclaw/SKILL.md`). `throttled` is not a negative finding.
## Hard rules ## Hard rules
@@ -24,7 +28,7 @@ Read first: [PLAN](PLAN.md) → [docs](docs/).
2. **Read-only data sources.** Ladybug `var/kb.lbug` and Postgres are opened 2. **Read-only data sources.** Ladybug `var/kb.lbug` and Postgres are opened
read-only for queries. Index rebuilds write to `var/` (gitignored). read-only for queries. Index rebuilds write to `var/` (gitignored).
3. **PII.** `brain-test`, `cs_brain` client data is never read or quoted. 3. **PII.** `brain-test`, `cs_brain` client data is never read or quoted.
4. **No main pushes.** Feature branches + PR via `gh`; CI must be green. 4. **No main pushes.** Feature branches + GitHub PR (`gh`); CI (Actions) must be green. Work board: [Gitea issues](https://git.produktor.io/eSlider/2dph/issues).
5. **TDD.** Failing test before tool code. Unit tests run offline against 5. **TDD.** Failing test before tool code. Unit tests run offline against
fixtures; network/db calls are wrapped. fixtures; network/db calls are wrapped.
6. **docs reflect behaviour.** Any change updates `docs/` + `PLAN.md` status. 6. **docs reflect behaviour.** Any change updates `docs/` + `PLAN.md` status.
@@ -35,28 +39,92 @@ Read first: [PLAN](PLAN.md) → [docs](docs/).
PLAN.md decisions + execution + open questions 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} (shebang) bin/ self-describing tools bin/{subject}/{method}.go (shebang)
bin/kb-watch corpus watcher (mtimes, no inotify deps) bin/brain/ search.go serve.go index.go add.go get.go stats.go eval.go watch.go
bin/docker-entrypoint container entrypoint (brain index|search|serve|watch) bin/chats/ sync.go import.go facts.go apply.go; libs in internal/chats
serve/ async Go HTTP server (goroutines, bounded worker pool) bin/mail/ sync.go import.go ocr.go (index_mail → brain/index.go)
tools/ vendored python libs behind bin/* (yamlout, websearch) bin/markdown/ import.go (H2 leaf split; Python bin/md/import fallback)
bin/postgres/ query.go (read-only YAML)
bin/git/ import.go (go-git history; Python shim execs it)
bin/web/ search.go (SearXNG; Python shim execs it)
bin/reasoner/ bakeoff.go (D18 CPU OpenAI tool-call bake-off)
internal/ shared Go (brain/rank is cgo-free; facts D16; chats; gitlog; websearch; reasoner; duckstats)
bin/qa/ stats.go (DuckDB quantiles / JSONL count; gcc CGO, not Zig)
bin/watch/ corpus watcher (used by bin/brain/watch.go)
bin/tools/ vendored python libs behind bin/* (kblib, yamlout, websearch)
bin/cgo/ zig zcc zc++ (CGO via zig cc, not gcc)
bin/docker-entrypoint container entrypoint (api: serve|search|watch; index: python)
compose.yaml docker composition (root level, not docker/) compose.yaml docker composition (root level, not docker/)
Dockerfile multi-stage: python deps + static Go serve Dockerfile api (Zig CGO, no Python) + index (Python write)
var/ kb.lbug, caches (gitignored) var/ kb.lbug, var/mail/*, caches (gitignored)
.venv/ ladybug + model2vec + mistune .venv/ ladybug + model2vec + mistune
``` ```
## Mail pipeline
```bash
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/import.go --from-raw var/mail # message.json → message.md (convert only)
bin/brain/index.go --rebuild --with-facts --with-chats
```
- `sync` (Go) downloads messages + attachments; Gmail uses paginated list +
`body.attachmentId` (not partId) for attachments.
- `import` converts body + attachments to markdown. PDFs use poppler
`pdftotext -layout` fast path (~15ms); textless/scanned PDFs use
`pdftoppm` + tesseract `eng+deu` (`bin/mail/ocr.go`). Optional
`OCR_ENGINE=paddle`. Conversion never touches the brain DB (crash safety).
- `index_mail` is a deprecation shim for `bin/brain/index.go --rebuild`. Bulk
rebuild still deletes `var/kb.lbug` and creates FTS/HNSW last. Single-leaf
write is `bin/brain/add.go` (safe while indexes exist; do not DROP INDEX).
Keep conversion + indexing separate so a conversion crash can't leave the
DB mid-transaction.
## Tools ## Tools
```bash ```bash
bin/facts/audit ["self"|"facts"|"info"|"stale"] # 2-source + staleness gate bin/facts/audit.go ["self"|"db"|"contradict"] # 2-source + D16 adjudication
bin/kb/search "query" [--hop N] [--repo X] # deduction search → YAML bin/facts/crm.go [--dry-run] # proof person↔company/company↔project (ooCRM × corpus SoT)
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" --no-web # local graph only
eval "$(bin/cgo/zig env)" # Zig cc + liblbug (not gcc)
bin/brain/index.go --rebuild [--with-mail] [--with-facts] [--with-chats]
bin/brain/add.go --text T --root facts --source "a.md x b.md" # incremental write
bin/brain/add.go --json # stdin leaf or {leafs:[...]}
bin/brain/get.go <id> [--body] [--json] # Go read; Python bin/kb/get CI fallback
bin/brain/stats.go [--json]
bin/brain/eval.go [--json] # recall@5; questions in internal/brain/rank
bin/brain/serve.go # HTTP :8630; GET /openapi.json POST /mcp
bin/markdown/import.go [dir] # H2 leafs → YAML; Python bin/md/import fallback
bin/git/import.go [REPO] [--json] [--limit N] # go-git history → commit leafs
bin/web/search.go "query" [--json] # SearXNG; throttled ≠ absence
bin/reasoner/bakeoff.go [--model ID] [--json] # D18 CPU tool-call bake-off
bin/postgres/query.go --profile onlyoffice -c 'SELECT 1'
bin/qa/stats.go # D22 DuckDB quantiles / JSONL (gcc CGO)
bin/mail/ocr.go <image|pdf> # tesseract eng+deu (scans)
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
``` ```
Never start a shell command with `cd` — use the tool working-directory Never start a shell command with `cd` — use the tool working-directory
parameter. Search before reading whole files. parameter. Search before reading whole files. For YAML/JSON/XML/CSV/TOML/HCL
prefer mikefarah/yq (`skills/yq/SKILL.md`). For bulk rows and quantiles use
duckdb-go (`internal/duckstats`, `skills/duckdb/SKILL.md`), not Ladybug.
## GitHub safety rules (ABSOLUTE — never violate)
1. **No absolute paths in committed files.** Replace `/mnt/`, `/home/<user>/`,
`/Users/<user>/` with env vars (`$HOME`, `$PROJECTS_ROOT`, `$DOCS_BASE`).
2. **No PII in commits.** No real names, phones, emails of third parties.
Test data must be synthetic (Alice, Bob, Charlie, Diana, example.com).
3. **No credentials/secrets in commits.** API keys, tokens, passwords, session
strings, phone numbers only in gitignored `.env` files, referenced by path.
4. **Curasoft, edelweiss — no files, no mentions.** Remove all traces if found.
5. **Check git history before push.** If any commit contains leaks, rewrite
history (rebase + force push) AND delete affected GitHub releases/tags.
6. **`docs/chat-import-plan.md`** — reference Gitea issue, never embed secrets.
## Communication ## Communication
+63 -15
View File
@@ -1,5 +1,13 @@
# syntax=docker/dockerfile:1 # syntax=docker/dockerfile:1
FROM python:3.12-slim AS base #
# docker build --target api -t 2dph:api .
# docker build --target index -t 2dph:index .
#
# API: Go + ladybug via Zig CGO (no CPython).
# Index: Python write path (profile `index` until brain/add is v2).
# --- Python sidecar (Ladybug write / rebuild) ---
FROM python:3.12-slim AS index
ENV PYTHONUNBUFFERED=1 \ ENV PYTHONUNBUFFERED=1 \
PYTHONDONTWRITEBYTECODE=1 \ PYTHONDONTWRITEBYTECODE=1 \
@@ -8,30 +16,70 @@ ENV PYTHONUNBUFFERED=1 \
WORKDIR /app WORKDIR /app
RUN id -u 2dph 2>/dev/null || useradd --create-home --uid 1001 2dph RUN id -u 2dph 2>/dev/null || useradd --create-home --uid 1001 2dph
RUN apt-get update \
&& apt-get install -y --no-install-recommends \
poppler-utils tesseract-ocr tesseract-ocr-eng tesseract-ocr-deu \
&& rm -rf /var/lib/apt/lists/*
# deps layer-first: rebuild only on dependency change
COPY requirements.lock.txt /tmp/requirements.lock.txt COPY requirements.lock.txt /tmp/requirements.lock.txt
RUN python -m pip install --no-cache-dir -r /tmp/requirements.lock.txt \ RUN python -m pip install --no-cache-dir -r /tmp/requirements.lock.txt \
&& rm /tmp/requirements.lock.txt && rm /tmp/requirements.lock.txt
# Go serve: static binary, no interpreter at runtime
FROM golang:1.25 AS serve-build
WORKDIR /src/serve
COPY serve/go.mod serve/go.sum* ./
COPY serve .
RUN CGO_ENABLED=0 go build -o /serve -ldflags="-s -w" .
# runtime: python toolchain + Go server
FROM base
COPY . . COPY . .
COPY --from=serve-build /serve /app/serve/serve RUN chmod +x /app/bin/docker-entrypoint \
RUN chmod +x /app/bin/kb-watch /app/bin/docker-entrypoint \
&& chown -R 2dph:2dph /app && chown -R 2dph:2dph /app
USER 2dph USER 2dph
ENV PATH="/app/bin:${PATH}" \ ENV PATH="/app/bin:${PATH}" \
KB_PY=python3 KB_PY=python3 \
KB_ROOT=/app
HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \ HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \
CMD python -c "import model2vec, ladybug, mistune; print('ok')" || exit 1 CMD python -c "import model2vec, ladybug, mistune; print('ok')" || exit 1
ENTRYPOINT ["/app/bin/docker-entrypoint"] ENTRYPOINT ["/app/bin/docker-entrypoint"]
# --- Go API: CGO with Zig, not gcc ---
FROM golang:1.26-bookworm AS api-build
WORKDIR /src
RUN apt-get update \
&& apt-get install -y --no-install-recommends curl xz-utils ca-certificates \
&& rm -rf /var/lib/apt/lists/*
COPY bin/cgo ./bin/cgo
RUN chmod +x bin/cgo/zig bin/cgo/zcc bin/cgo/zc++ \
&& ./bin/cgo/zig env >/dev/null
COPY go.mod go.sum ./
RUN go mod download
COPY . .
ENV CGO_RPATH=/usr/local/lib
RUN eval "$(./bin/cgo/zig env)" \
&& go build -tags brain_serve,system_ladybug -o /out/brain-serve ./bin/brain/serve.go \
&& go build -tags system_ladybug -o /out/brain-search ./bin/brain/search.go \
&& CGO_ENABLED=0 go build -tags brain_watch -o /out/brain-watch ./bin/brain/watch.go
FROM debian:bookworm-slim AS api
RUN apt-get update \
&& apt-get install -y --no-install-recommends libssl3 ca-certificates wget \
&& rm -rf /var/lib/apt/lists/* \
&& useradd --create-home --uid 1001 2dph
COPY --from=api-build /out/brain-serve /usr/local/bin/brain-serve
COPY --from=api-build /out/brain-search /usr/local/bin/brain-search
COPY --from=api-build /out/brain-watch /usr/local/bin/brain-watch
COPY --from=api-build /src/lib-ladybug/liblbug.so.0.19.1 /usr/local/lib/liblbug.so.0.19.1
COPY bin/docker-entrypoint /usr/local/bin/docker-entrypoint
RUN chmod +x /usr/local/bin/docker-entrypoint \
&& ln -s liblbug.so.0.19.1 /usr/local/lib/liblbug.so.0 \
&& ln -s liblbug.so.0 /usr/local/lib/liblbug.so \
&& ldconfig
USER 2dph
ENV KB_ROOT=/data \
KB_PORT=8630 \
LD_LIBRARY_PATH=/usr/local/lib \
HF_HOME=/data/hf
WORKDIR /data
EXPOSE 8630
HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \
CMD wget -qO- http://127.0.0.1:8630/health || exit 1
ENTRYPOINT ["/usr/local/bin/docker-entrypoint"]
CMD ["serve"]
+97 -29
View File
@@ -4,7 +4,12 @@ A brain that loves facts and deduction. Evidence-first knowledge graph + hybrid
RAG over the operational Brain/ops/eSlider stack. Built like Sherlock RAG over the operational Brain/ops/eSlider stack. Built like Sherlock
Holmes: nothing is asserted unless it has proof. Holmes: nothing is asserted unless it has proof.
Status: **in progress** — this file is the plan and the record of decisions. Status: **v1 in** (epic [#16](https://git.produktor.io/eSlider/2dph/issues/16) closed).
v2 board: milestone [v2](https://git.produktor.io/eSlider/2dph/milestone/13) —
OCR [#6](https://git.produktor.io/eSlider/2dph/issues/6) in,
[#29](https://git.produktor.io/eSlider/2dph/issues/29) OQ1 in,
[#30](https://git.produktor.io/eSlider/2dph/issues/30) OQ3 in.
Gap: [docs/roadmap.md](docs/roadmap.md).
## What ## What
@@ -25,21 +30,27 @@ detective method: **a fact needs ≥2 independent sources or it is
| # | Question | Answer | | # | Question | Answer |
|---|----------|--------| |---|----------|--------|
| D1 | RAG corpus | ops stack (chat, onlyoffice, gitea/NPM, searchxng, observability, ai-bot, mcp-servers, `~/.ssh/config`) + portfolio. Exclude `office.dev` + jobs/applications. | | D1 | RAG corpus | ops stack (chat, onlyoffice, gitea/NPM, searchxng, observability, ai-bot, mcp-servers, `~/.ssh/config`) + portfolio. Exclude `office.dev` + jobs/applications. |
| D2 | skill merging | integrate skills **in this project** `skills/`; skip gitea / brain-detective-depe ndent skills. | | D2 | skill merging | integrate skills **in this project** `skills/`; skip gitea / brain-dependent skills. |
| D3 | web search | import `web-search`, retire local `searxng-ops`. Vendored here, no remote link. | | D3 | web search | Go client `bin/web/search.go` (`internal/websearch`). SearXNG URL is config (`BRAIN_SEARCH_URL`). Optional Compose profile `searxng` (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** (Kuzu successor, MIT, embedded, native FTS+vector+Cypher). Python binding for `bin/*`; Go shebang for golang tools. | | D6 | graph engine | **LadybugDB**. Go is the service (`bin/brain/search.go`, `bin/brain/serve.go` in-process, `internal/brain`). Read path is Go + Zig CGO (D21). Python `bin/kb/{get,stats,eval}` is the CI fallback when Zig/libs are not fetched. Incremental write is Python `bin/kb/add` (`bin/brain/add.go`). Bulk rebuild stays `compose --profile index` 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)`. 2-source 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. |
| D10 | versioning | everything is a leaf with `sha256 + observed_at + source_rev`; `File-[:HAS_VERSION]->Commit-[:AUTHORED]->Person`. Stale = `source_rev` < git HEAD. | | D10 | versioning | everything is a leaf with `sha256 + observed_at + source_rev`; `File-[:HAS_VERSION]->Commit-[:AUTHORED]->Person`. Stale = `source_rev` < git HEAD. |
| D11 | strong/weak | `root` column: `facts` (strong) vs `info` (weak). Answer is `confirmed` only from facts root. | | D11 | strong/weak | `root` column: `facts` (strong) vs `info` (weak). Answer is `confirmed` only from facts root. |
| D12 | transactional | facts and info split by root but **written in the same Ladybug transaction (ACID)** on every write. | | D12 | transactional | facts and info split by root but **written in the same Ladybug transaction (ACID)** on every write. |
| D13 | portfolio | start graph `(Person:eslider)-[:HAS]->(Portfolio)`, associate other natural/juristic persons later. | | D13 | portfolio | start graph `(Person:eslider)-[:HAS]->(Portfolio)`, associate other natural/juristic persons later. |
| D14 | tooling style | `bin/{subject}/{method}` self-describing: shebang line 1, usage comment from line 2. Go shebang: `///usr/bin/env go run "$0" "$@"; exit`. | | D14 | tooling style | `bin/{subject}/{method}.go` shebang (e.g. `bin/brain/search.go`). Shared code in `internal/`. One root `go.mod` + `go.work`. No `bin/*/main.go`, no nested modules. |
| D15 | repo | GitHub `eSlider/2dph`, public (like sibling repos), push/commit via `gh`, TDD + commit every change, CI/CD. | | D15 | repo | Gitea [`eSlider/2dph`](https://git.produktor.io/eSlider/2dph) is origin + [issues](https://git.produktor.io/eSlider/2dph/issues). GitHub `eSlider/2dph` is the public clone (PRs + Actions CI). No direct `main` pushes. TDD → PR → CI green → merge. |
| D16 | contradictions | ≥2 yes vs ≥2 no → unrelated sources conflict → hypothesis → `(not confirmed)`. Resolution (authority, staleness adjudication) = **v2**, tracked as open question. | | D16 | contradictions | ≥2 yes vs ≥2 no → hypothesis → `(not confirmed)` until a rule fires. Order: **temporal_freshness** (fresh ≥2 vs stale minority), then **authority_pairing** (runtime/config A×B beats narrative C). Store as `a x b vs c x d` on hypothesis leafs. `bin/facts/audit contradict`. [#29](https://git.produktor.io/eSlider/2dph/issues/29). |
| D17 | assertion gate | Fact-check every *claim* (facts → info → live → web), not every edit. `bin/brain/search.go` adds a `web` block when there is no facts hit (`throttled`/`skipped`/`refused` ≠ absence). `--root` and `--no-web` stay local. Missing graph ≠ “does not exist”. |
| D18 | reasoner | Pluggable OpenAI-compatible URL (`REASONER_BASE_URL`). RAM: `Qwen/Qwen3.5-9B`. Quality: `prism-ml/Bonsai-27B-gguf` or `Qwen/Qwen3.6-27B`. No official Qwen3.6-9B. CPU bake-off: `bin/reasoner/bakeoff.go` + compose profile `reasoner` (`OLLAMA_NUM_GPU=0`, `:11435`). PicoClaw is compose profile `picoclaw`; tools are `search`/`get`/`audit`. Weights are not copied into the 2dph image. Agent lever/loop: [#15](https://git.produktor.io/eSlider/2dph/issues/15). |
| D19 | git history | [go-git](https://github.com/go-git/go-git) via `bin/git/import.go`. No subprocess of the git binary. Conversion prints commit leafs; brain write is `bin/brain/index.go`. |
| D20 | agent API | OpenAPI + MCP are generated from the same `internal/httpapi.Ops` table as `bin/brain/serve.go` handlers. `GET /openapi.json`, `POST /mcp` (JSON-RPC tools/list + tools/call). Tool names match OpenAPI paths (`search`/`get`/`stats`/`audit`/`ingest`). |
| D21 | CGO | Ladybug/tokenizers CGO is compiled with **Zig** (`bin/cgo/zcc``zig cc -target …-linux-gnu`), not gcc. `bin/cgo/zig` pins Zig 0.14.1 + liblbug 0.19.1 + libtokenizers 1.27.0. Compose `target: api` has no CPython; write/rebuild is profile `index`. |
| D22 | analytics | **duckdb-go** in-process (`internal/duckstats`, `bin/qa/stats.go`) for quantiles/JSONL. Links with **gcc/g++**, not Zig. Ladybug stays the graph; web-search cache stays modernc sqlite. Slice small structured docs with **mikefarah/yq**, not kislyuk/jq. [#30](https://git.produktor.io/eSlider/2dph/issues/30). |
## Architecture ## Architecture
@@ -47,16 +58,30 @@ detective method: **a fact needs ≥2 independent sources or it is
2dph/ 2dph/
PLAN.md / AGENTS.md PLAN.md / AGENTS.md
docs/ published docs (this conversation → docs/ as md) docs/ published docs (this conversation → docs/ as md)
skills/ in-project skills (web-search, db-yaml, kb-search, agent-cost, diataxis-docs, …) skills/ in-project skills (web-search, postgres, brain, picoclaw, diataxis-docs)
bin/ bin/
facts/extract auto-pair 2 sources → lexicon yaml + graph facts/extract.go audit.go crm.go # D14 shebang; Python implementation
facts/audit ["self"|"facts"|"info"|"stale"] 2-source + staleness gate kb/index Python bulk write (called by bin/brain/index.go)
kb/index build FTS + HNSW from corpus kb/add Python incremental write (called by bin/brain/add.go)
kb/search deduction: facts → info → web-search; --hop N brain/index.go rebuild FTS + HNSW (incl. --with-mail)
kb/get kb/stats kb/eval brain/add.go incremental leaf write (no rebuild)
md/import md/select md/tables md/gaps (mistune) brain/get.go stats.go eval.go # Go read (cgo); Python bin/kb/* CI fallback
brain/watch.go
brain/search.go deduction: facts → info → web-search
brain/serve.go HTTP API in-process + OpenAPI/MCP (D20); Zig CGO (D21)
cgo/zig zcc zc++ CGO toolchain (zig cc, not gcc)
mail/import.go JSON → markdown (no brain write)
markdown/import.go H2 leaf split (Go); Python bin/md/import fallback
postgres/query.go read-only YAML (wraps bin/db/psql-yq)
git/import.go go-git history (no git binary; conversion only)
web/search.go SearXNG client (throttled ≠ absence)
reasoner/bakeoff.go CPU tool-call bake-off (D18; OpenAI tools)
chats/sync.go import.go facts.go apply.go
(libs in internal/chats; no chats index)
mail/ocr.go tesseract eng+deu (pdftoppm scans)
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 (deprecated shim → web/search.go)
db/psql-yq (vendored) db/psql-yq (vendored)
ssh-tunnel onlyoffice pg tunnel 5433 ssh-tunnel onlyoffice pg tunnel 5433
var/kb.lbug single embedded store (gitignored) var/kb.lbug single embedded store (gitignored)
@@ -67,7 +92,10 @@ detective method: **a fact needs ≥2 independent sources or it is
Node tables: `Person, Service, Host, Container, Repo, File, Commit, Leaf`. Node tables: `Person, Service, Host, Container, Repo, File, Commit, Leaf`.
`Leaf(embedding FLOAT[N])` — FTS on `text`, HNSW vector index on `embedding`. `Leaf(embedding FLOAT[N])` — FTS on `text`, HNSW vector index on `embedding`.
Edges: `RUNS / USES / HAS_VERSION / AUTHORED / ABOUT / ASSOCIATED / SIMILAR_0.85`. Edges: `RUNS / USES / FROM_FILE / HAS_VERSION / AUTHORED / ABOUT / ASSOCIATED / SIMILAR_0.85`.
`FROM_FILE` / `HAS_VERSION` / `AUTHORED`: `bin/brain/search.go --hop N` walks
them from each hit (1=File, 2=Commit, 3=Person). Rebuild writes
`Leaf-[:FROM_FILE]->File`; git import writes the rest.
Common props on every node/edge: `root`, `confidence`, `evidence[]`, `how`, Common props on every node/edge: `root`, `confidence`, `evidence[]`, `how`,
`where`, `when`, `source_rev`. `where`, `when`, `source_rev`.
@@ -85,26 +113,49 @@ Common props on every node/edge: `root`, `confidence`, `evidence[]`, `how`,
- `bin/{subject}/{method}` — line 2 is a usage comment (mirrors `psql-yq`). - `bin/{subject}/{method}` — line 2 is a usage comment (mirrors `psql-yq`).
- bash + python primary; golang via Go shebang when a compiled helper is right. - bash + python primary; golang via Go shebang when a compiled helper is right.
- YAML default output, `--json` for machines. Slice with `yq`. - YAML default output, `--json` for machines. Slice with mikefarah/yq.
- Everything that touches the network / DB is read-only, throttled, cached. - Everything that touches the network / DB is read-only, throttled, cached.
- Tests (TDD) gate every commit; `gh` + CI/CD on every push. - Tests (TDD) gate every commit; `gh` + CI/CD on every push.
## Open questions (v2) ## Open questions (v2)
- OQ1: mutually-contradicting evidence — how to resolve (authority weighting, - OQ1: **in** — D16 adjudication: `temporal_freshness` then `authority_pairing`.
temporal freshness, audit adjudication). Unresolved 2v2 stays hypothesis. [#29](https://git.produktor.io/eSlider/2dph/issues/29).
- OQ2: OCR pipeline for pdfs/images/docs (late phase). - OQ2: OCR **in**. `pdftotext -layout` first; scans `pdftoppm` + tesseract
- OQ3: optional duckdb-md layer for `SELECT … FORMAT MARKDOWN` export/write-back. `eng+deu` (`bin/mail/ocr.go`, `internal/ocr`). No gocv, no gosseract CGO
(D21 Zig owns Ladybug CGO). Optional `OCR_ENGINE=paddle` / compose profile
`ocr-paddle`. Docling left the default path. [#6](https://git.produktor.io/eSlider/2dph/issues/6).
- OQ3: **in** — duckdb-go (`internal/duckstats`, `bin/qa/stats.go`) for
quantiles / JSONL count. Not a second graph. [#30](https://git.produktor.io/eSlider/2dph/issues/30).
- OQ4: YAML-first storage for leafs — deferred: JSON is ~10x faster to
serialize and unambiguous; YAML only where humans edit files.
## Mail pipeline (done)
1. `bin/mail/sync.go` (Go, 8 workers) — paginated Gmail/OnlyOffice download.
Gmail attachments key off `body.attachmentId`, not MIME `partId`.
2. `bin/mail/import.go --from-raw` — message.json → message.md; PDFs via
`pdftotext -layout` (~15ms); textless/scanned PDFs `pdftoppm` + tesseract
`eng+deu`. ICS sidecars
Latin-1→UTF-8 normalized.
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
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
via `bin/brain/search.go`.
## CI/CD pipeline (D15) ## CI/CD pipeline (D15)
`.github/workflows/ci.yml`: `.github/workflows/ci.yml`:
1. go vet + go test ./... (Go tools) 1. go vet + go test ./... (root module; packages without ladybug cgo)
2. python -m unittest discover + pytest (Py tools) 2. `go test ./internal/brain/rank` (cgo-free ranking + flag parser)
3. bin/facts/audit self (lexicon internal consistency) 3. python -m unittest discover -s bin/tools (includes published-docs SoT)
4. bin/kb/eval (recall@5 ≥ 0.95, gates index regressions) 4. `bin/facts/audit self` (lexicon internal consistency; `bin/facts/audit.go` is the D14 wrapper)
5. md-docs build/lint if docs tooling arrives. 5. `bin/brain/eval.go` via Zig (recall@5 ≥ 0.95). Python `bin/kb/eval` is an
explicit fallback, not the CI SoT.
6. `bin/cgo/zig go build -tags system_ladybug` (compile search with zig cc; fetches pinned zig+libs).
Feedback loop: every commit → PR → CI → green/gate → merge. Same discipline as Feedback loop: every commit → PR → CI → green/gate → merge. Same discipline as
`db/tech-poc`: contract first where there is an OpenAPI/message shape. `db/tech-poc`: contract first where there is an OpenAPI/message shape.
@@ -113,9 +164,26 @@ Feedback loop: every commit → PR → CI → green/gate → merge. Same discipl
1. scaffold repo (:done after this file + AGENTS.md + .gitignore + ci) 1. scaffold repo (:done after this file + AGENTS.md + .gitignore + ci)
2. gh repo create eSlider/2dph --private + initial commit + CI 2. gh repo create eSlider/2dph --private + initial commit + CI
3. vendored skill integration (web-search, db-yaml, kb-search, agent-cost, diataxis-docs) — no remote links 3. vendored skill integration (web-search, postgres, brain, diataxis-docs) — no remote links
4. .venv: ladybug + model2vec + mistune 4. .venv: ladybug + model2vec + mistune
5. schema + tools with TDD (kb + md + facts + brain) 5. schema + tools with TDD (kb + md + facts + brain)
6. ~/.config/brain config 6. ~/.config/brain config
7. corpus extraction (facts/info) 7. corpus extraction (facts/info)**in**: [#18](https://git.produktor.io/eSlider/2dph/issues/18)
8. verify: web-search smoke, onlyoffice pg, md-db round-trip, eval, audit 8. verify: web-search smoke, onlyoffice pg, md-db round-trip, eval, audit
## Gap to v1 (epic #16)
Remaining: none for epic #16 (v1). Board:
[epic #16](https://git.produktor.io/eSlider/2dph/issues/16),
milestone [v1 detective brain](https://git.produktor.io/eSlider/2dph/milestone/12).
Narrative: [docs/roadmap.md](docs/roadmap.md).
| Order | Issue | Gap |
|-------|-------|-----|
| 1 | [#14](https://git.produktor.io/eSlider/2dph/issues/14) | **in**`bin/brain/add.go` / `POST /ingest` write facts+info without deleting `kb.lbug`. Bulk corpus still `--rebuild`. Leftover Python (mail/facts) is not the living-graph blocker. |
| 2 | [#17](https://git.produktor.io/eSlider/2dph/issues/17) | **in**`--hop N` walks `FROM_FILE``HAS_VERSION``AUTHORED` (max 3). |
| 3 | [#18](https://git.produktor.io/eSlider/2dph/issues/18) | **in**`--with-facts` / `--facts-json` land `root=facts`; `--with-chats` indexes `var/chats/md`. WhatsApp sync is out of v1. |
| 4 | [#15](https://git.produktor.io/eSlider/2dph/issues/15) | **in** — lever/loop documented (`search``get``audit`). |
| 5 | [#19](https://git.produktor.io/eSlider/2dph/issues/19) | **in** — CI recall SoT is `bin/brain/eval.go` via Zig. Python `bin/kb/eval` stays as an explicit fallback. |
Does **not** block epic close: OQ4. OCR [#6](https://git.produktor.io/eSlider/2dph/issues/6), OQ1 [#29](https://git.produktor.io/eSlider/2dph/issues/29), OQ3 [#30](https://git.produktor.io/eSlider/2dph/issues/30) are **in**.
+87 -40
View File
@@ -7,14 +7,16 @@
[![Latest Release](https://img.shields.io/github/v/tag/eSlider/2dph?sort=semver&label=release)](https://github.com/eSlider/2dph/releases) [![Latest Release](https://img.shields.io/github/v/tag/eSlider/2dph?sort=semver&label=release)](https://github.com/eSlider/2dph/releases)
[![GitHub Stars](https://img.shields.io/github/stars/eSlider/2dph?style=social)](https://github.com/eSlider/2dph/stargazers) [![GitHub Stars](https://img.shields.io/github/stars/eSlider/2dph?style=social)](https://github.com/eSlider/2dph/stargazers)
An evidence-first brain over the operational eSlider stack. **Facts need two An evidence-first brain. **Facts need two independent sources, or they are
independent sources, or they are `(not confirmed)`.** `(not confirmed)`.** Cursor is not the runtime.
`2dph` is a single embedded knowledge graph (LadybugDB = Kuzu successor) with `2dph` is a single embedded knowledge graph (LadybugDB) with native **HNSW
native **HNSW vector** + **BM25 full-text** indexes, built from markdown, vector** + **BM25 full-text** indexes. Search is *deduction*: confirmed facts
compose files, ssh config, docker state, and git history. Search is first, supporting info second, `web-search` as the independent second source
*deduction*: confirmed facts first, supporting info second, `web-search` as when the local graph cannot confirm.
the independent second source when the local graph cannot confirm.
Run it: [docs/runbook.md](docs/runbook.md). Design: [docs/design.md](docs/design.md).
Docs index: [docs/README.md](docs/README.md).
## Architecture ## Architecture
@@ -28,11 +30,11 @@ graph TB
end end
subgraph dph["2dph tools"] subgraph dph["2dph tools"]
EX["bin/facts/extract<br/>2-source pairing"] EX["bin/facts/extract.go<br/>2-source pairing"]
AU["bin/facts/audit<br/>confidence + staleness"] AU["bin/facts/audit.go<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/>H2 leaf split"]
SR["bin/kb/search<br/>deduction + --hop"] SR["bin/brain/search.go<br/>deduction"]
end end
subgraph store["Ladybug var/kb.lbug"] subgraph store["Ladybug var/kb.lbug"]
@@ -74,7 +76,7 @@ graph TB
## The method ## The method
Every assertion is `Who / What / How / Where / When + evidence + confidence`, Every assertion is `Who / What / How / Where / When + evidence + confidence`,
mirroring the detective detective skill: **≥2 independent sources confirm a mirroring the detective method: **≥2 independent sources confirm a
fact; conflicting sources or a single source → `hypothesis``(not confirmed)`.** fact; conflicting sources or a single source → `hypothesis``(not confirmed)`.**
| root | meaning | used for answers | | root | meaning | used for answers |
@@ -85,55 +87,100 @@ fact; conflicting sources or a single source → `hypothesis` → `(not confirme
## Deduction search ## Deduction search
```bash ```bash
bin/kb/search "Matrix federation over HTTPS" # facts → info → web-search bin/brain/search.go "Matrix federation over HTTPS" # facts → info → web
bin/kb/search "what runs on arc-2" --hop 1 # walk graph edges bin/brain/search.go "onlyoffice postgres" --root facts
bin/kb/search "where is cs-lexicon" --json | yq '.' # YAML by default bin/brain/search.go "where is cs-lexicon" --json | yq '.'
bin/kb/get <id> --body # full chunk on demand bin/brain/search.go "upstream flag" --no-web # local graph only
bin/kb/stats # index health bin/brain/get.go <id> --body # full chunk on demand
bin/kb/eval # recall@5 gate bin/brain/stats.go # index health
bin/brain/eval.go # recall@5 gate
```
`--hop N` walks File/Commit/Person from each hit (max 3). `bin/kb/search` is a deprecated wrapper around `bin/brain/search.go`.
Git history is read with [go-git](https://github.com/go-git/go-git) (no git binary):
```bash
bin/git/import.go --json --limit 100 # commit leafs for this repo
bin/git/import.go --root "$PROJECTS_ROOT" --json # one pass per .git under root
```
Conversion only. Graph write (`File-[:HAS_VERSION]->Commit-[:AUTHORED]->Person`) stays with `bin/brain/index.go`.
Web search (second independent source) goes through SearXNG. Empty results mean **throttled**, not “nothing exists”:
```bash
bin/web/search.go "LadybugDB vector index" --json
# Optional local instance (skip if BRAIN_SEARCH_URL already points at one):
# SEARXNG_SECRET=$(openssl rand -hex 32) docker compose --profile searxng up -d
```
Mail is a first-class corpus (retrievable through the same search):
```bash
bin/mail/sync.go --source onlyoffice,gmail --workers 8 --out var/mail # raw sync (Go)
bin/mail/import.go --from-raw var/mail # JSON → markdown
bin/brain/add.go --text T --root facts --source "a.md x b.md"
bin/brain/index.go --rebuild --with-facts --with-chats # facts extract + chats md
bin/brain/index.go --rebuild # rebuild brain (incl. mail)
bin/brain/search.go "invoice from last week" # same search over mail leafs
``` ```
## Storage ## Storage
- **LadybugDB** — single `var/kb.lbug`, Cypher property graph, HNSW + BM25 - **LadybugDB** — single `var/kb.lbug`, Cypher + HNSW + BM25, embedded.
in one engine, embedded (no server), ACID, read-only-safe for concurrent Read tools (`get` / `stats` / `eval`) are Go + Zig CGO (`bin/cgo/zcc`).
readers. Python fallbacks stay for CI until the runner fetches Zig. Incremental
- **model2vec** — `potion-multilingual-128M` static embeddings (256-dim), write is `bin/brain/add.go` (Python `kblib.add_leafs`). Bulk rebuild is
CPU-fast, deterministic, no Ollama runtime dependency. Compose profile `index` (`bin/brain/index.go --rebuild`).
- facts and info split semantically by `root` column but written inside the - **model2vec** — `potion-multilingual-128M` (256-dim), CPU, no Ollama
same transaction. runtime dependency.
- facts and info split by `root` but written in the same transaction.
Ladybug 0.19 DROP INDEX warning: [docs/runbook.md](docs/runbook.md).
## Tooling conventions ## Tooling conventions
`bin/{subject}/{method}` — self-describing: shebang on line 1, usage comment `bin/{subject}/{method}.go` — self-describing: shebang on line 1, usage comment
from line 2. bash + python primary; golang via the Go shebang when a compiled from line 2. Shared code in `internal/`. YAML default output, `--json` for
helper is right. YAML default output, `--json` for machines. Everything that machines. Tests gate every commit. HTTP: `bin/brain/serve.go` calls
touches network/db is read-only, throttled, cached. Tests gate every commit. `internal/brain` in-process (`/health` `/search` `/get` `/stats` `/audit` `/ingest` `/openapi.json` `/mcp`).
## Development ## Development
See the portable runbook: [docs/runbook.md](docs/runbook.md).
```bash ```bash
uv venv .venv # Python 3.12, uv-managed uv venv .venv
uv pip install -r requirements.lock.txt # pinned toolchain uv pip install -r requirements.lock.txt
bin/facts/audit self # lexicon consistency gate bin/facts/audit.go self
go test ./... && python -m unittest discover -s tools -t . go test ./... && uv run python -m unittest discover -s bin/tools -t .
``` ```
Docker (optional, cached model + var volumes): Docker (optional, cached model + var volumes):
```bash ```bash
docker compose run --rm brain index # (re)index corpus docker compose up -d brain # API (Zig CGO serve :8630)
docker compose run --rm brain search "query" # one-shot query docker compose --profile index run --rm index # Python Ladybug rebuild
docker compose run --rm brain serve # async Go HTTP server docker compose --profile picoclaw up brain-mcp # MCP on 127.0.0.1:8630
docker compose --profile reasoner up -d reasoner # CPU Ollama 127.0.0.1:11435
docker compose up brain-watch # auto re-index on change docker compose up brain-watch # auto re-index on change
``` ```
## Related ## Related
eSlider DevOps engineer practice: ops, OnlyOffice, and mail feed the facts
root through `bin/facts/extract` (two-source pairing).
- [go-second-brain](https://github.com/eSlider/go-second-brain) — the earlier - [go-second-brain](https://github.com/eSlider/go-second-brain) — the earlier
Neo4j + Qdrant + Matrix RAG brain Neo4j + Qdrant + Matrix RAG brain
- [agent-skills](https://github.com/eSlider/agent-skills) — upstream - [agent-skills](https://github.com/eSlider/agent-skills) — upstream
skills (`web-search`, `db-yaml`, …) that 2dph integrates skills (`web-search`, `postgres`, …) that 2dph integrates
- [detective](https://github.com/detective) — the two-source method - detective method — the two-source method
See [PLAN.md](PLAN.md) for decisions, execution status, and v2 open questions. Work board (issues): [epic #16](https://git.produktor.io/eSlider/2dph/issues/16)
on [git.produktor.io/eSlider/2dph/issues](https://git.produktor.io/eSlider/2dph/issues).
PRs and CI: GitHub [`eSlider/2dph`](https://github.com/eSlider/2dph).
See [PLAN.md](PLAN.md) for decisions, [docs/roadmap.md](docs/roadmap.md) for
the gap to v1, and v2 open questions.
+21
View File
@@ -0,0 +1,21 @@
//usr/bin/env go run -tags=brain_add "$0" "$@"; exit
//go:build brain_add
//
// bin/brain/add.go - incremental leaf write (Python kblib, no rebuild).
//
// ./bin/brain/add.go --text T --root facts --source "a.md x b.md"
// ./bin/brain/add.go --json
//
// D6: write stays Python. Does not delete var/kb.lbug.
// 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/add", os.Args[1:]))
}
+3
View File
@@ -0,0 +1,3 @@
// Commands in this directory are shebang mains (search.go, serve.go, index.go,
// get.go, stats.go, eval.go, watch.go), each behind an exclusive build tag.
package main
+22
View File
@@ -0,0 +1,22 @@
//usr/bin/env go run -tags=system_ladybug,brain_eval "$0" "$@"; exit
//go:build cgo && system_ladybug && brain_eval
//
// bin/brain/eval.go - recall@5 gate.
//
// ./bin/brain/eval.go
// ./bin/brain/eval.go --json
//
// Needs CGO + libladybug. Python bin/kb/eval is the CI fallback (no cgo).
// Control questions live in internal/brain/rank (cgo-free).
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/brain"
)
func main() {
os.Exit(brain.MainEval(os.Args[1:]))
}
+23
View File
@@ -0,0 +1,23 @@
//usr/bin/env go run -tags=system_ladybug,brain_get "$0" "$@"; exit
//go:build cgo && system_ladybug && brain_get
//
// bin/brain/get.go - read one leaf by id.
//
// ./bin/brain/get.go <id>
// ./bin/brain/get.go <id> --body
// ./bin/brain/get.go <id> --json
//
// Needs CGO + libladybug. Python bin/kb/get is the CI fallback (no cgo).
// CGO compiler is Zig (`eval "$(bin/cgo/zig env)"`), not gcc.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/brain"
)
func main() {
os.Exit(brain.MainGet(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 --with-facts --with-chats
// ./bin/brain/index.go --rebuild --with-mail
// ./bin/brain/index.go --dry-run --with-mail
//
// v1 write: bin/brain/add.go for one/few leafs (indexes may already exist).
// Bulk mail/corpus still --rebuild (fresh file, indexes last).
// 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))
}
+23
View File
@@ -0,0 +1,23 @@
//usr/bin/env go run -tags=system_ladybug "$0" "$@"; exit
//go:build cgo && system_ladybug
//
// bin/brain/search.go - deduction search over the 2dph brain.
//
// ./bin/brain/search.go "query" [--root facts|info] [--repo P] [-n N] [--hop N] [--json] [--no-web]
// ./bin/brain/search.go serve [port]
// ./bin/brain/search.go --list-model
//
// Needs CGO + libladybug via Zig (`eval "$(bin/cgo/zig env)"`), not gcc.
// bin/kb/search which sets those and builds a binary for the embed daemon.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/brain"
)
func main() {
os.Exit(brain.Main(os.Args[1:]))
}
+34
View File
@@ -0,0 +1,34 @@
//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 (in-process ladybug search).
//
// KB_ROOT=/path/to/2dph ./bin/brain/serve.go
// KB_WORKERS=4 KB_PORT=8630 ./bin/brain/serve.go
//
// GET /openapi.json same Ops table as the handlers
// POST /mcp JSON-RPC tools/list + tools/call
//
// 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"
)
func main() {
if os.Getenv("KB_ROOT") == "" {
if wd, err := os.Getwd(); err == nil {
os.Setenv("KB_ROOT", wd)
}
}
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)
}
+21
View File
@@ -0,0 +1,21 @@
//usr/bin/env go run -tags=system_ladybug,brain_stats "$0" "$@"; exit
//go:build cgo && system_ladybug && brain_stats
//
// bin/brain/stats.go - index health.
//
// ./bin/brain/stats.go
// ./bin/brain/stats.go --json
//
// Needs CGO + libladybug. Python bin/kb/stats is the CI fallback (no cgo).
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/brain"
)
func main() {
os.Exit(brain.MainStats(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:])
}
Executable
+23
View File
@@ -0,0 +1,23 @@
#!/bin/sh
# bin/cgo/zc++ — CGO CXX. Zig, not g++.
set -eu
ROOT="$(CDPATH= cd -- "$(dirname "$0")/../.." && pwd)"
case "$(uname -m)" in
x86_64|amd64) TARGET=x86_64-linux-gnu ;;
aarch64|arm64) TARGET=aarch64-linux-gnu ;;
*)
echo "zc++: unsupported arch $(uname -m)" >&2
exit 2
;;
esac
if [ -n "${ZIG:-}" ] && [ -x "$ZIG" ]; then
:
elif [ -x "$ROOT/var/zig/zig" ]; then
ZIG="$ROOT/var/zig/zig"
elif command -v zig >/dev/null 2>&1; then
ZIG="$(command -v zig)"
else
echo "zc++: zig missing; run bin/cgo/zig first" >&2
exit 127
fi
exec "$ZIG" c++ -target "$TARGET" "$@"
Executable
+24
View File
@@ -0,0 +1,24 @@
#!/bin/sh
# bin/cgo/zcc — CGO CC. Zig, not gcc.
# Go invokes CC with many args; a wrapper avoids spaces in $CC.
set -eu
ROOT="$(CDPATH= cd -- "$(dirname "$0")/../.." && pwd)"
case "$(uname -m)" in
x86_64|amd64) TARGET=x86_64-linux-gnu ;;
aarch64|arm64) TARGET=aarch64-linux-gnu ;;
*)
echo "zcc: unsupported arch $(uname -m)" >&2
exit 2
;;
esac
if [ -n "${ZIG:-}" ] && [ -x "$ZIG" ]; then
:
elif [ -x "$ROOT/var/zig/zig" ]; then
ZIG="$ROOT/var/zig/zig"
elif command -v zig >/dev/null 2>&1; then
ZIG="$(command -v zig)"
else
echo "zcc: zig missing; run bin/cgo/zig first" >&2
exit 127
fi
exec "$ZIG" cc -target "$TARGET" "$@"
Executable
+130
View File
@@ -0,0 +1,130 @@
#!/usr/bin/env bash
# bin/cgo/zig — CGO toolchain: zig cc (not gcc) + pinned liblbug + libtokenizers.
#
# eval "$(bin/cgo/zig env)" # export CC/CXX/CGO_*
# bin/cgo/zig go build ... # ensure, then exec with env
# bin/cgo/zig ./bin/brain/search.go "query"
#
# Pins live in this file. Downloads land in var/ (gitignored).
set -euo pipefail
ROOT="$(CDPATH= cd -- "$(dirname "$0")/../.." && pwd)"
ZIG_VERSION=0.14.1
LBUG_VERSION=0.19.1
TOKENIZERS_VERSION=1.27.0
arch="$(uname -m)"
case "$arch" in
x86_64|amd64)
ZIG_ARCH=x86_64
LBUG_ARCH=x86_64
TOK_ARCH=x86_64
ZIG_SHA=24aeeec8af16c381934a6cd7d95c807a8cb2cf7df9fa40d359aa884195c4716c
LBUG_SHA=ed263ae913f68cb0ddba0b98548b58edaac49929766d03bdaaa83be46c68847d
TOK_SHA=72556cdca798dd4ea7cdaba308e5f0d68a8cb93b67c96edf485b7a0edd7b07f4
;;
aarch64|arm64)
ZIG_ARCH=aarch64
LBUG_ARCH=aarch64
TOK_ARCH=aarch64
ZIG_SHA=f7a654acc967864f7a050ddacfaa778c7504a0eca8d2b678839c21eea47c992b
LBUG_SHA=b07df2cd533c3976a2a3025866d6420a5f35514d0a822ecc4b2902d55b4725b7
TOK_SHA=e96545ad05930c26f51f63d932ee6d3bbd32bbed149e102c5290d587a2293067
;;
*)
echo "bin/cgo/zig: unsupported arch $arch" >&2
exit 2
;;
esac
CACHE="$ROOT/var/cache"
LIB="$ROOT/lib-ladybug"
ZIG_DIR="$ROOT/var/zig-dist"
ZIG_BIN="$ROOT/var/zig/zig"
sha256of() {
if command -v sha256sum >/dev/null 2>&1; then
sha256sum "$1" | awk '{print $1}'
else
shasum -a 256 "$1" | awk '{print $1}'
fi
}
fetch() {
local url="$1" dest="$2" expect="$3"
if [ -f "$dest" ] && [ "$(sha256of "$dest")" = "$expect" ]; then
return 0
fi
mkdir -p "$(dirname "$dest")"
echo "fetch $url" >&2
curl -fsSL "$url" -o "$dest"
local got
got="$(sha256of "$dest")"
if [ "$got" != "$expect" ]; then
echo "checksum mismatch $dest: got $got want $expect" >&2
rm -f "$dest"
exit 1
fi
}
ensure_zig() {
if [ -n "${ZIG:-}" ] && [ -x "$ZIG" ]; then
return 0
fi
if [ -x "$ZIG_BIN" ]; then
export ZIG="$ZIG_BIN"
return 0
fi
if command -v zig >/dev/null 2>&1; then
export ZIG
ZIG="$(command -v zig)"
return 0
fi
local tar="$CACHE/zig-${ZIG_ARCH}-linux-${ZIG_VERSION}.tar.xz"
fetch "https://ziglang.org/download/${ZIG_VERSION}/zig-${ZIG_ARCH}-linux-${ZIG_VERSION}.tar.xz" \
"$tar" "$ZIG_SHA"
mkdir -p "$CACHE"
rm -rf "$ZIG_DIR"
tar -xJf "$tar" -C "$CACHE"
mv "$CACHE/zig-${ZIG_ARCH}-linux-${ZIG_VERSION}" "$ZIG_DIR"
mkdir -p "$ROOT/var/zig"
ln -sfn "$ZIG_DIR/zig" "$ZIG_BIN"
export ZIG="$ZIG_BIN"
}
ensure_libs() {
mkdir -p "$LIB"
if [ ! -f "$LIB/liblbug.so" ]; then
local tar="$CACHE/liblbug-linux-${LBUG_ARCH}.tar.gz"
fetch "https://github.com/LadybugDB/ladybug/releases/download/v${LBUG_VERSION}/liblbug-linux-${LBUG_ARCH}.tar.gz" \
"$tar" "$LBUG_SHA"
tar -xzf "$tar" -C "$LIB"
fi
if [ ! -f "$LIB/libtokenizers.a" ]; then
local tar="$CACHE/libtokenizers.linux-${TOK_ARCH}.tar.gz"
fetch "https://github.com/daulet/tokenizers/releases/download/v${TOKENIZERS_VERSION}/libtokenizers.linux-${TOK_ARCH}.tar.gz" \
"$tar" "$TOK_SHA"
tar -xzf "$tar" -C "$LIB"
fi
}
print_env() {
printf 'export ZIG=%q\n' "$ZIG"
printf 'export CC=%q\n' "$ROOT/bin/cgo/zcc"
printf 'export CXX=%q\n' "$ROOT/bin/cgo/zc++"
printf 'export CGO_ENABLED=1\n'
printf 'export CGO_CFLAGS=%q\n' "-I$LIB"
printf 'export CGO_LDFLAGS=%q\n' "-L$LIB -Wl,-rpath,${CGO_RPATH:-$LIB}"
}
ensure_zig
ensure_libs
cmd="${1:-env}"
if [ "$cmd" = "env" ]; then
print_env
exit 0
fi
eval "$(print_env)"
exec "$@"
Executable
+29
View File
@@ -0,0 +1,29 @@
#!/usr/bin/env bash
# bin/chats - sync, import, index, facts, apply for Telegram/WhatsApp/LinkedIn.
# Builds the chats binary on first run / when source changes, then execs it.
set -euo pipefail
ROOT="$(cd "$(dirname "$0")/.." && pwd)"
BIN="$ROOT/var/bin/chats"
SRC="$ROOT/bin/chats"
mkdir -p "$ROOT/var/bin"
need_build=0
if [ ! -x "$BIN" ]; then
need_build=1
else
while IFS= read -r -d '' f; do
if [ "$f" -nt "$BIN" ]; then
need_build=1
break
fi
done < <(find "$SRC" -name '*.go' -print0 2>/dev/null)
fi
if [ "$need_build" -eq 1 ]; then
echo "Building chats..." >&2
(cd "$SRC" && go build -o "$BIN" .) || exit 1
fi
exec "$BIN" "$@"
+19
View File
@@ -0,0 +1,19 @@
//usr/bin/env go run -tags=chats_apply "$0" "$@"; exit
//go:build chats_apply
//
// bin/chats/apply.go - push extracted chat facts to OnlyOffice CRM.
//
// ./bin/chats/apply.go [--dry-run]
//
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/chats"
)
func main() {
os.Exit(chats.RunApply(os.Args[1:]))
}
+4
View File
@@ -0,0 +1,4 @@
// Commands in this directory are shebang mains (sync.go, import.go, facts.go,
// apply.go), each behind an exclusive build tag so `go build ./bin/chats`
// does not see two mains. Shared code lives in internal/chats.
package main
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=chats_facts "$0" "$@"; exit
//go:build chats_facts
//
// bin/chats/facts.go - extract phone/email/linkedin facts from JSONL.
//
// ./bin/chats/facts.go
//
// Writes var/chats/facts/. Does not index the brain.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/chats"
)
func main() {
os.Exit(chats.RunFacts(os.Args[1:]))
}
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=chats_import "$0" "$@"; exit
//go:build chats_import
//
// bin/chats/import.go - JSONL → markdown under var/chats/md/.
//
// ./bin/chats/import.go
//
// Conversion only. Brain ingest is bin/brain/index.go, not this command.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/internal/chats"
)
func main() {
os.Exit(chats.RunImport(os.Args[1:]))
}
+155
View File
@@ -0,0 +1,155 @@
#!/usr/bin/env python3
"""chats/refresh-linkedin-session - refresh LinkedIn MCP session from webtop CDP.
bin/chats/refresh-linkedin-session [--cdp URL] [--root DIR]
"""
Reads the current LinkedIn cookies out of the running Thorium browser in the
work-webtop container via CDP (Network.getAllCookies), copies the live browser
profile onto the source profile directory, and rewrites the portable
cookies.json + source-state.json that mcp-server-linkedin requires.
Usage:
refresh-linkedin-session [--cdp http://127.0.0.1:9222] [--root /var/tmp/liprofile]
[--container work-webtop] [--profile thorium-profile]
After the headless driver uses a copied profile, LinkedIn rotates the session
in that copy, so this must run before every sync.
"""
import asyncio
import json
import os
import shutil
import subprocess
import sys
import tempfile
import urllib.request
import websockets
def cdp_tab(ws_json):
for t in ws_json:
if t.get("webSocketDebuggerUrl"):
return t["webSocketDebuggerUrl"]
return None
async def get_cookies(ws_url):
async with websockets.connect(ws_url, max_size=50_000_000) as ws:
await ws.send(json.dumps({"id": 1, "method": "Network.getAllCookies", "params": {}}))
resp = await ws.recv()
return json.loads(resp).get("result", {}).get("cookies", [])
def write_source_state(root, profile_dir):
# Reuse the linkedin-mcp-server session_state module to write a valid
# source-state.json (same schema the daemon reads).
try:
from linkedin_mcp_server.session_state import canonical, write_source_state
write_source_state(canonical(__import__("pathlib").Path(profile_dir)))
return
except Exception:
pass
# Fallback: minimal schema-compatible state.
import uuid
state = {
"version": 1,
"source_runtime_id": "linux-amd64-host",
"login_generation": str(uuid.uuid4()),
"created_at": None,
"profile_path": profile_dir,
"cookies_path": os.path.join(root, "cookies.json"),
}
from datetime import datetime, timezone
state["created_at"] = datetime.now(timezone.utc).isoformat()
with open(os.path.join(root, "source-state.json"), "w") as f:
json.dump(state, f, indent=2)
def main():
args = sys.argv[1:]
cdp = "http://127.0.0.1:9222"
root = "/var/tmp/liprofile"
container = "work-webtop"
cprofile = "thorium-profile"
for i in range(0, len(args), 2):
k = args[i]
v = args[i + 1] if i + 1 < len(args) else ""
if k == "--cdp":
cdp = v
elif k == "--root":
root = v
elif k == "--container":
container = v
elif k == "--profile":
cprofile = v
profile_dir = os.path.join(root, "profile")
os.makedirs(profile_dir, exist_ok=True)
# 1. Clear stale daemon/browser locks so the server can claim the profile.
for lock in ("profile-claim.lock", "profile.lock", "daemon.lock", "lease.lock"):
p = os.path.join(root, lock)
if os.path.exists(p):
os.remove(p)
for name in os.listdir(profile_dir):
if name.startswith("Singleton"):
os.remove(os.path.join(profile_dir, name))
for name in os.listdir(root):
if name.startswith("invalid-state-"):
shutil.rmtree(os.path.join(root, name), ignore_errors=True)
# 1. Copy the live browser profile (cookies DB + Local State) so the
# session the driver launches carries the current login.
subprocess.run(
["docker", "cp", f"{container}:/config/{cprofile}/Default", os.path.join(profile_dir, "Default")],
check=True, capture_output=True,
)
subprocess.run(
["docker", "cp", f"{container}:/config/{cprofile}/Local State", os.path.join(profile_dir, "Local State")],
check=True, capture_output=True,
)
for lock in ("SingletonLock", "SingletonCookie", "SingletonSocket"):
p = os.path.join(profile_dir, lock)
if os.path.exists(p):
os.remove(p)
# 2. Pull the live cookies out of the running browser.
with urllib.request.urlopen(f"{cdp}/json", timeout=5) as r:
tabs = json.loads(r.read())
ws_url = cdp_tab(tabs)
if not ws_url:
sys.stderr.write("refresh-linkedin-session: no CDP tab\n")
sys.exit(1)
cookies = asyncio.run(get_cookies(ws_url))
li = [c for c in cookies if "linkedin" in c.get("domain", "")]
out = []
for c in li:
domain = c.get("domain", "")
if domain in (".www.linkedin.com", "www.linkedin.com"):
domain = ".linkedin.com"
out.append({
"name": c["name"],
"value": c["value"].strip('"'),
"domain": domain,
"path": c.get("path", "/"),
"expires": c.get("expires", -1),
"httpOnly": c.get("httpOnly", False),
"secure": c.get("secure", False),
"sameSite": c.get("sameSite", "None"),
})
with open(os.path.join(root, "cookies.json"), "w") as f:
json.dump(out, f, indent=2)
write_source_state(root, profile_dir)
sys.stderr.write(f"refresh-linkedin-session: {len(out)} cookies, profile refreshed\n")
if __name__ == "__main__":
main()
+42
View File
@@ -0,0 +1,42 @@
//usr/bin/env go run -tags=chats_sync "$0" "$@"; exit
//go:build chats_sync
//
// bin/chats/sync.go - download chat messages to var/chats/<platform>/.
//
// ./bin/chats/sync.go telegram [--limit N] [--phone PHONE]
// ./bin/chats/sync.go linkedin [--limit N] [--refresh]
//
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"fmt"
"os"
"github.com/eSlider/2dph/internal/chats"
)
func main() {
if len(os.Args) < 2 {
fmt.Fprintln(os.Stderr, `usage: bin/chats/sync.go telegram|linkedin [flags]`)
os.Exit(2)
}
platform := os.Args[1]
args := os.Args[2:]
switch platform {
case "telegram":
os.Exit(chats.RunSyncTelegram(args))
case "linkedin":
os.Exit(chats.RunSyncLinkedIn(args))
case "whatsapp":
fmt.Fprintln(os.Stderr, "chats: WhatsApp sync is out of v1")
os.Exit(1)
case "help", "-h", "--help":
fmt.Fprintln(os.Stderr, `usage: bin/chats/sync.go telegram|linkedin [flags]
WhatsApp sync is out of v1.`)
return
default:
fmt.Fprintf(os.Stderr, "chats: unknown platform %q\n", platform)
os.Exit(2)
}
}
+1 -1
View File
@@ -17,7 +17,7 @@ import subprocess
import sys import sys
from pathlib import Path from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "tools")) sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "tools"))
from semver import bump_type, bump_version # noqa: E402 from semver import bump_type, bump_version # noqa: E402
+3 -1
View File
@@ -34,11 +34,13 @@ case "${1:-}" in
;; ;;
"") "")
[ -f "$HOME/.ssh/config" ] || { echo "db/ssh-tunnel: ~/.ssh/config missing" >&2; exit 1; } [ -f "$HOME/.ssh/config" ] || { echo "db/ssh-tunnel: ~/.ssh/config missing" >&2; exit 1; }
if db/ssh-tunnel --check; then if "$0" --check; then
echo "tunnel already up on ${SRC}" echo "tunnel already up on ${SRC}"
exit 0 exit 0
fi fi
ssh -f -N -M -S "$HOME/.ssh/2dph-tunnel.sock" \ ssh -f -N -M -S "$HOME/.ssh/2dph-tunnel.sock" \
-o StrictHostKeyChecking=accept-new \
-o BatchMode=yes \
-L "${SRC}:${DST}" -p "$SSH_PORT" "${SSH_USER}@${SSH_HOST}" \ -L "${SRC}:${DST}" -p "$SSH_PORT" "${SSH_USER}@${SSH_HOST}" \
&& echo "tunnel up on ${SRC} (-> vm:${DST})" && echo "tunnel up on ${SRC} (-> vm:${DST})"
exit 0 exit 0
+2
View File
@@ -0,0 +1,2 @@
// Deprecated shebang mains at bin root (serve.go is tagged brain_serve).
package main
Regular → Executable
+26 -11
View File
@@ -1,11 +1,10 @@
#!/usr/bin/env bash #!/usr/bin/env bash
# bin/docker-entrypoint - run 2dph tools inside the container. # bin/docker-entrypoint - run 2dph tools inside the container.
# #
# brain shell (default) # API image (Zig CGO binaries):
# brain search <q> bin/kb/search # serve | search | watch
# brain index bin/kb/index # Index image (Python write path, compose profile `index`):
# brain watch <dir> watchdog re-indexer # index | extract | audit | search (deprecated python wrapper)
# brain serve async Go HTTP server (serve/)
# #
# Usage comment starts at line 2 (self-describing convention). # Usage comment starts at line 2 (self-describing convention).
set -euo pipefail set -euo pipefail
@@ -13,11 +12,27 @@ set -euo pipefail
CMD="${1:-shell}" CMD="${1:-shell}"
shift || true shift || true
if [ -x /usr/local/bin/brain-serve ]; then
case "$CMD" in
shell) exec bash ;;
serve) exec /usr/local/bin/brain-serve "$@" ;;
search) exec /usr/local/bin/brain-search "$@" ;;
watch) exec /usr/local/bin/brain-watch "$@" ;;
index)
echo "index is the Python sidecar: docker compose --profile index run --rm index" >&2
exit 2
;;
*) echo "unknown command: $CMD (api: serve|search|watch)" >&2; exit 2 ;;
esac
fi
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 bash /app/bin/kb-watch "$@" ;; watch) exec /app/bin/watch "$@" ;;
serve) exec /app/serve/serve "$@" ;; serve) exec /app/bin/serve "$@" ;;
*) echo "unknown command: $CMD" >&2; exit 2 ;; extract) exec "$KB_PY" /app/bin/facts/extract "$@" ;;
audit) exec "$KB_PY" /app/bin/facts/audit "$@" ;;
*) echo "unknown command: $CMD" >&2; exit 2 ;;
esac esac
+46 -19
View File
@@ -1,14 +1,15 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
"""facts/audit - evidence & lexicon checks for the 2dph brain. """facts/audit - evidence & lexicon checks for the 2dph brain.
bin/facts/audit self # lexicon: every fact in db has >=2 sources bin/facts/audit self # lexicon: docs + two-source rule
bin/facts/audit db # evidence gate: run against var/kb.lbug bin/facts/audit db # evidence gate against var/kb.lbug
bin/facts/audit contradict # D16 adjudication (JSON claim(s) on stdin)
`self` mode checks the repo itself (no network, no runtime deps). It greps `self` mode checks the repo itself (no network, no runtime deps).
for known-good two-source pairings and confirms the docs are consistent. `db` mode loads every Leaf with root=facts. Confirmed facts need ` x `;
`db` mode loads every Leaf with root=facts and asserts each has source_rev hypothesis contradictions need `a x b vs c x d` (both sides ≥2).
and a non-empty `loc` (the "where did you see it" evidence pointer) and that `contradict` applies temporal_freshness then authority_pairing; ≥2 vs ≥2
'confirmed' facts carry a two-source `source` field. with no rule stays hypothesis / `(not confirmed)`.
Exit 0 = all checks pass, 1 = audit failures, 2 = could not evaluate. Exit 0 = all checks pass, 1 = audit failures, 2 = could not evaluate.
""" """
@@ -20,7 +21,9 @@ 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 / "tools")) sys.path.insert(0, str(ROOT / "bin" / "tools"))
from contradict import adjudicate, check_fact_row # noqa: E402
def audit_db() -> list[str]: def audit_db() -> list[str]:
@@ -33,14 +36,8 @@ def audit_db() -> list[str]:
r = conn.execute("MATCH (l:Leaf {root:'facts'}) RETURN l.id, l.source, l.loc, l.how, l.confidence") r = conn.execute("MATCH (l:Leaf {root:'facts'}) RETURN l.id, l.source, l.loc, l.how, l.confidence")
problems: list[str] = [] problems: list[str] = []
for lid, source, loc, how, conf in r.get_all(): for lid, source, loc, how, conf in r.get_all():
if conf != "confirmed": problems.extend(check_fact_row(str(lid), str(source or ""), str(loc or ""),
problems.append(f"{lid}: facts require confidence='confirmed', got '{conf}'") str(how or ""), str(conf or "")))
if not source or " x " not in source:
problems.append(f"{lid}: needs 2-source evidence in source, got '{source}'")
if not loc:
problems.append(f"{lid}: missing loc (evidence pointer)")
if not how:
problems.append(f"{lid}: missing how")
conn.close() conn.close()
db.close() db.close()
return problems return problems
@@ -54,20 +51,50 @@ def audit_self() -> list[str]:
problems.append("PLAN.md missing recall@5 gate") problems.append("PLAN.md missing recall@5 gate")
if re.search(r"(?i)facts must have.*2 sources|2.source", plan) is None: if re.search(r"(?i)facts must have.*2 sources|2.source", plan) is None:
problems.append("PLAN.md missing the two-source evidence rule for facts") problems.append("PLAN.md missing the two-source evidence rule for facts")
if "temporal_freshness" not in plan or "authority_pairing" not in plan:
problems.append("PLAN.md missing D16 adjudication rules")
if re.search(r"(?i)HNSW|BM25|deduction", (ROOT / "README.md").read_text()) is None: if re.search(r"(?i)HNSW|BM25|deduction", (ROOT / "README.md").read_text()) is None:
problems.append("README.md missing search/retrieval description") problems.append("README.md missing search/retrieval description")
return problems return problems
def audit_contradict(raw: str) -> tuple[list[str], list[dict]]:
raw = raw.strip()
if not raw:
return ["contradict: empty stdin (JSON claim or {claims:[...]})"], []
try:
payload = json.loads(raw)
except json.JSONDecodeError as e:
return [f"contradict: invalid JSON: {e}"], []
if isinstance(payload, dict) and "claims" in payload:
claims = list(payload.get("claims") or [])
elif isinstance(payload, dict):
claims = [payload]
elif isinstance(payload, list):
claims = payload
else:
return ["contradict: expected object or list"], []
details = [adjudicate(c) for c in claims]
return [], details
def main(argv: list[str]) -> int: def main(argv: list[str]) -> int:
import argparse import argparse
p = argparse.ArgumentParser(description="evidence & lexicon audit") p = argparse.ArgumentParser(description="evidence & lexicon audit")
p.add_argument("mode", choices=("self", "db")) p.add_argument("mode", choices=("self", "db", "contradict"))
p.add_argument("--json", action="store_true") p.add_argument("--json", action="store_true")
a = p.parse_args(argv) a = p.parse_args(argv)
problems = audit_self() if a.mode == "self" else audit_db() details: list[dict] = []
out = {"mode": a.mode, "ok": not problems, "problems": problems} if a.mode == "self":
problems = audit_self()
elif a.mode == "db":
problems = audit_db()
else:
problems, details = audit_contradict(sys.stdin.read())
out: dict = {"mode": a.mode, "ok": not problems, "problems": problems}
if details:
out["contradictions"] = details
if a.json: if a.json:
print(json.dumps(out, indent=2)) print(json.dumps(out, indent=2))
else: else:
+22
View File
@@ -0,0 +1,22 @@
//usr/bin/env go run -tags=facts_audit "$0" "$@"; exit
//go:build facts_audit
//
// bin/facts/audit.go - 2-source + lexicon checks.
//
// ./bin/facts/audit.go self
// ./bin/facts/audit.go db
// ./bin/facts/audit.go contradict --json < claim.json
//
// Python bin/facts/audit is the implementation (CI runs it directly).
// 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/facts/audit", os.Args[1:]))
}
Executable
+124
View File
@@ -0,0 +1,124 @@
#!/usr/bin/env python3
"""facts/crm - prove person->company and company->project associations.
Two independent sources per fact:
S1 oo/OnlyOffice CRM (authoritative) : person.company_id -> company,
project.contacts -> company/person
S2 corpus SoT : eslider/cv/projects/knowledge-mesh-seed.yaml
(orgs: employer/client/... + projects)
Only associations supported by BOTH sources are written as root=facts.
Mismatches are reported (or, with --fix-crm, printed as oo CLI commands).
Usage:
bin/facts/crm write proven facts (needs var/kb.lbug)
bin/facts/crm --dry-run show proposed facts + mismatches only
bin/facts/crm --mismatches show associations found in only one side
"""
from __future__ import annotations
import json
import os
import sys
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(ROOT / "bin" / "tools"))
from kblib import upsert_leaf, connect, leaf_id # noqa: E402
MESH_ENV = os.environ.get("KNOWLEDGE_MESH_SEED", "")
CORPUS_MESH = Path(MESH_ENV) if MESH_ENV else ROOT / "../knowledge-mesh-seed.yaml"
def corpus_orgs(raw: str) -> dict[str, dict]:
"""Delegate to tools.crmfacts.corpus_orgs (tested in tools/)."""
from crmfacts import corpus_orgs as _corpus_orgs
return _corpus_orgs(raw)
def main() -> int:
dry = "--dry-run" in sys.argv
mism = "--mismatches" in sys.argv
mesh = CORPUS_MESH.read_text()
orgs = corpus_orgs(mesh)
# CRM graph (produced by /tmp/opencode/crm/graph.py -> /tmp/opencode/crm/graph.json)
graph = json.load(open("/tmp/opencode/crm/graph.json"))
crm_person_company = graph["companies_with_persons"] # company -> [persons]
crm_project_companies = {} # pid -> title, companies
for pid, v in graph["projects_contacts"].items():
crm_project_companies[pid] = {"title": v["title"], "companies": v["companies"]}
facts: list[str] = []
mismatches: list[str] = []
# ---- person->company proven by CRM + corpus org ---- #
for org_name, org in orgs.items():
token = org.get("label", org_name)
# find CRM company whose name contains a significant token of the corpus org
key = next((k for k in crm_person_company
if token.split()[0].lower() in k.lower() or any(
t.lower() in k.lower() for t in org.get("label", "").split(" / "))),
None)
persons = crm_person_company.get(key, []) if key else []
if persons and org:
for p in persons:
facts.append(f"{p} is associated with {org.get('label')} "
f"(role: {org.get('kind', '?')}, {org.get('period', '')})")
elif org and key and not persons:
mismatches.append(f"corpus org '{org_name}' ({org.get('label')}) has no CRM persons")
elif org and not key:
mismatches.append(f"corpus org '{org_name}' ({org.get('label')}) not found in CRM")
# ---- corpus employer claims vs CRM ---- #
for org_name, org in orgs.items():
if not org or not org.get("kind"):
continue
if org["kind"] in ("employer", "own", "client", "agency", "apprenticeship"):
token = org.get("label", org_name).split()[0]
if not any(token.lower() in k.lower() for k in crm_person_company):
mismatches.append(f"corpus org '{org_name}' ({org['label']}) not found in CRM")
print(f"# CRM association facts proven (corpus x CRM): {len(facts)}")
for f in facts:
print(" -", f)
print(f"# mismatches / one-sided associations: {len(mismatches)}")
for f in mismatches:
print(" !", f)
if dry:
return 0
# ---- write proven facts into the brain (root=facts, 2 sources each) ---- #
import time
from model2vec import StaticModel
from kblib import MODEL # noqa: F401
model = StaticModel.from_pretrained(MODEL)
db, conn = connect(read_only=False)
try:
r = conn.execute("MATCH (l:Leaf) WHERE l.root='facts' RETURN count(*) AS n")
stats_before = r.get_all()[0][0]
except Exception:
stats_before = 0
rev = time.strftime("%Y%m%d-%H%M%S")
written = 0
for f in facts:
src = f"ooCRM x {CORPUS_MESH.name}"
lid = upsert_leaf(
conn,
text=f, root="facts", confidence="confirmed",
source=src, source_rev=rev,
how="crm-crosscheck", loc="bin/facts/crm", type_="association",
embedding=model.encode(f).tolist(),
)
written += 1
conn.close()
print(f"# wrote {written} facts into var/kb.lbug (facts was {stats_before})")
return 0
if __name__ == "__main__":
sys.exit(main())
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=facts_crm "$0" "$@"; exit
//go:build facts_crm
//
// bin/facts/crm.go - prove person↔company / company↔project (ooCRM × corpus).
//
// ./bin/facts/crm.go [--dry-run] [--mismatches]
//
// Python bin/facts/crm is the implementation. Graph write stays Python.
// 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/facts/crm", os.Args[1:]))
}
+4 -3
View File
@@ -21,7 +21,7 @@ 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 / "tools")) sys.path.insert(0, str(ROOT / "bin" / "tools"))
COMPOSE_FILES = [ROOT / "docker" / "compose.yaml", ROOT / "compose.yaml"] COMPOSE_FILES = [ROOT / "docker" / "compose.yaml", ROOT / "compose.yaml"]
DOC_MARKERS = ["README.md", "PLAN.md", "AGENTS.md"] DOC_MARKERS = ["README.md", "PLAN.md", "AGENTS.md"]
@@ -176,8 +176,7 @@ def dedupe(facts: list[dict]) -> list[dict]:
def write_facts(facts: list[dict]) -> None: def write_facts(facts: list[dict]) -> None:
from kblib import connect, init_schema, upsert_leaf from kblib import VAR, connect, ensure_indexes, init_schema, upsert_leaf
from kblib import VAR
VAR.mkdir(exist_ok=True) VAR.mkdir(exist_ok=True)
db, conn = connect(VAR / "kb.lbug", read_only=False) db, conn = connect(VAR / "kb.lbug", read_only=False)
init_schema(conn) init_schema(conn)
@@ -188,6 +187,8 @@ def write_facts(facts: list[dict]) -> None:
upsert_leaf(conn, text=f["text"], root="facts", confidence="confirmed", upsert_leaf(conn, text=f["text"], root="facts", confidence="confirmed",
source=f["source"], source_rev=REPO, how=f["how"], source=f["source"], source_rev=REPO, how=f["how"],
loc=f["loc"], type_="fact", embedding=emb) loc=f["loc"], type_="fact", embedding=emb)
# Upsert-with-index is safe; never DROP+recreate (ghost catalog kills HNSW).
ensure_indexes(conn)
conn.close() conn.close()
db.close() db.close()
+20
View File
@@ -0,0 +1,20 @@
//usr/bin/env go run -tags=facts_extract "$0" "$@"; exit
//go:build facts_extract
//
// bin/facts/extract.go - acquire confirmed facts (2-source each).
//
// ./bin/facts/extract.go [--json] [--dry-run]
//
// Python bin/facts/extract is the implementation. Graph write stays Python.
// 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/facts/extract", os.Args[1:]))
}
Executable
+26
View File
@@ -0,0 +1,26 @@
#!/usr/bin/env python3
"""git/import — deprecated. Use bin/git/import.go (go-git, no git binary).
bin/git/import.go [REPO] [--json] [--limit N] [--since DATE]
"""
from __future__ import annotations
import os
import sys
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
def main(argv: list[str]) -> int:
print(
"bin/git/import is deprecated; use bin/git/import.go (go-git)",
file=sys.stderr,
)
target = ROOT / "bin" / "git" / "import.go"
os.execvp("go", ["go", "run", str(target), *argv])
return 1
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))
+143
View File
@@ -0,0 +1,143 @@
//usr/bin/env go run "$0" "$@"; exit
//
// bin/git/import.go - read git history with go-git (no git binary).
//
// ./bin/git/import.go [REPO]
// ./bin/git/import.go --json
// ./bin/git/import.go --limit 100 --since 2026-01-01
// ./bin/git/import.go --root DIR
//
// Conversion only: prints commit leafs. Brain write is bin/brain/index.go.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"encoding/json"
"fmt"
"os"
"path/filepath"
"strconv"
"time"
"github.com/eSlider/2dph/internal/cmdbin"
"github.com/eSlider/2dph/internal/gitlog"
)
func main() {
os.Exit(run(os.Args[1:]))
}
func run(args []string) int {
var repo, root, since string
limit := 0
jsonOut := false
i := 0
for i < len(args) {
a := args[i]
switch {
case a == "--json":
jsonOut = true
case a == "--limit" && i+1 < len(args):
i++
n, err := strconv.Atoi(args[i])
if err != nil || n < 0 {
fmt.Fprintf(os.Stderr, "git/import: --limit must be a non-negative integer\n")
return 2
}
limit = n
case a == "--since" && i+1 < len(args):
i++
since = args[i]
case a == "--root" && i+1 < len(args):
i++
root = args[i]
case a == "-h" || a == "--help":
fmt.Fprintln(os.Stderr, `usage: bin/git/import.go [REPO] [--json] [--limit N] [--since DATE] [--root DIR]`)
return 0
case len(a) > 0 && a[0] != '-':
repo = a
default:
fmt.Fprintf(os.Stderr, "git/import: unknown flag %s\n", a)
return 2
}
i++
}
var sinceT time.Time
if since != "" {
var err error
sinceT, err = parseSince(since)
if err != nil {
fmt.Fprintf(os.Stderr, "git/import: %v\n", err)
return 2
}
}
repos := []string{}
if repo != "" {
repos = []string{repo}
} else if root != "" {
entries, err := os.ReadDir(root)
if err != nil {
fmt.Fprintf(os.Stderr, "git/import: %v\n", err)
return 1
}
for _, e := range entries {
p := filepath.Join(root, e.Name())
if _, err := os.Stat(filepath.Join(p, ".git")); err == nil {
repos = append(repos, p)
}
}
} else {
repos = []string{cmdbin.Root()}
}
opt := gitlog.Options{Limit: limit, Since: sinceT}
type row struct {
Repo string `json:"repo"`
Path string `json:"path"`
Commits int `json:"commits"`
Leafs []gitlog.Leaf `json:"leafs,omitempty"`
}
var rows []row
for _, p := range repos {
name, err := gitlog.RepoName(p)
if err != nil && name == "" {
fmt.Fprintf(os.Stderr, "git/import: %s: %v\n", p, err)
continue
}
cs, err := gitlog.Log(p, opt)
if err != nil {
fmt.Fprintf(os.Stderr, "git/import: %s: %v\n", p, err)
return 1
}
leafs := make([]gitlog.Leaf, 0, len(cs))
for _, c := range cs {
leafs = append(leafs, gitlog.ToLeaf(c, name))
}
rows = append(rows, row{Repo: name, Path: p, Commits: len(cs), Leafs: leafs})
}
if jsonOut {
enc := json.NewEncoder(os.Stdout)
enc.SetIndent("", " ")
enc.SetEscapeHTML(false)
if err := enc.Encode(rows); err != nil {
return 1
}
return 0
}
for _, r := range rows {
fmt.Printf("%-24s %5d commits %s\n", r.Repo, r.Commits, r.Path)
}
return 0
}
func parseSince(s string) (time.Time, error) {
for _, layout := range []string{time.RFC3339, "2006-01-02"} {
if t, err := time.Parse(layout, s); err == nil {
return t, nil
}
}
return time.Time{}, fmt.Errorf("cannot parse --since %q", s)
}
-27
View File
@@ -1,27 +0,0 @@
#!/usr/bin/env bash
# kb-watch - re-index 2dph when corpus files change.
#
# kb-watch [dir...] [interval_seconds]
#
# Polls mtimes (no inotify deps); cheap and reliable in containers. Defaults:
# dirs = /corpus (compose) or . ; interval = 30s.
set -euo pipefail
DEFAULT_DIRS="${KB_WATCH_DIRS:-/corpus}"
DIRS=("$@")
[[ ${#DIRS[@]} -eq 0 ]] && DIRS=(${DEFAULT_DIRS})
INTERVAL="${KB_WATCH_INTERVAL:-30}"
index() { "${KB_PY:-python3}" /app/bin/kb/index; }
LAST_STAMP=""
while true; do
STAMP=$(find "${DIRS[@]}" -type f -newermt "-${INTERVAL} seconds" 2>/dev/null \
| head -1 | md5sum)
if [[ -n "$STAMP" && "$STAMP" != "$LAST_STAMP" ]]; then
echo "kb-watch: changes detected, re-indexing" >&2
index || echo "kb-watch: index failed; will retry" >&2
LAST_STAMP="$STAMP"
fi
sleep "$INTERVAL"
done
Executable
+114
View File
@@ -0,0 +1,114 @@
#!/usr/bin/env python3
"""kb/add - incremental leaf write (no rebuild).
bin/kb/add --text T --root facts|info --source S
bin/kb/add --json # stdin: one object or {"leafs":[...]}
bin/kb/add --db PATH --json
Writes facts+info in one Ladybug transaction. Does not delete kb.lbug.
Embedding is used when provided; otherwise model2vec encodes the text.
"""
from __future__ import annotations
import json
import sys
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(ROOT / "bin" / "tools"))
from kblib import ( # noqa: E402
EMBED_DIM,
add_leafs,
connect,
ensure_indexes,
init_schema,
)
def _as_leafs(payload: object) -> list[dict]:
if isinstance(payload, list):
return [dict(x) for x in payload]
if isinstance(payload, dict):
if "leafs" in payload:
return [dict(x) for x in payload["leafs"]]
return [dict(payload)]
raise ValueError("json must be an object, a list, or {leafs:[...]}")
def _embed_missing(leafs: list[dict]) -> None:
missing = [lf for lf in leafs if not lf.get("embedding")]
if not missing:
return
from model2vec import StaticModel
model = StaticModel.from_pretrained("minishlab/potion-multilingual-128M")
for lf in missing:
text = str(lf.get("text") or "")
vec = model.encode([text])[0].astype(float).tolist()
if len(vec) != EMBED_DIM:
vec = (vec + [0.0] * EMBED_DIM)[:EMBED_DIM]
lf["embedding"] = vec
def main(argv: list[str]) -> int:
import argparse
p = argparse.ArgumentParser(description="add leafs without rebuilding the brain")
p.add_argument("--db", default="", help="path to kb.lbug (default var/kb.lbug)")
p.add_argument("--json", action="store_true", help="read leaf JSON from stdin")
p.add_argument("--text", default="", help="leaf text")
p.add_argument("--root", default="info", choices=("facts", "info"))
p.add_argument("--source", default="")
p.add_argument("--confidence", default="confirmed")
p.add_argument("--source-rev", default="working-tree")
p.add_argument("--how", default="brain/add")
p.add_argument("--loc", default="")
p.add_argument("--type", default="reference", dest="type_")
args = p.parse_args(argv)
if args.json:
raw = sys.stdin.read()
if not raw.strip():
print("kb/add: empty stdin", file=sys.stderr)
return 2
leafs = _as_leafs(json.loads(raw))
else:
if not args.text or not args.source:
print("kb/add: --text and --source are required (or --json)", file=sys.stderr)
return 2
leafs = [{
"text": args.text,
"root": args.root,
"source": args.source,
"confidence": args.confidence,
"source_rev": args.source_rev,
"how": args.how,
"loc": args.loc or args.source,
"type": args.type_,
}]
for lf in leafs:
if not lf.get("text") or not lf.get("source"):
print("kb/add: each leaf needs text and source", file=sys.stderr)
return 2
_embed_missing(leafs)
from kblib import DB_PATH, VAR
dbpath = Path(args.db) if args.db else DB_PATH
dbpath.parent.mkdir(parents=True, exist_ok=True)
VAR.mkdir(exist_ok=True)
db, conn = connect(dbpath, read_only=False)
init_schema(conn)
ids = add_leafs(conn, leafs)
ensure_indexes(conn)
conn.close()
db.close()
print(json.dumps({"mode": "add", "ids": ids, "db": str(dbpath)}))
return 0
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))
+1 -1
View File
@@ -13,7 +13,7 @@ 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 / "tools")) sys.path.insert(0, str(ROOT / "bin" / "tools"))
from kblib import open_readonly, query_fts # noqa: E402 from kblib import open_readonly, query_fts # noqa: E402
from yamlout import to_yaml # noqa: E402 from yamlout import to_yaml # noqa: E402
+1 -1
View File
@@ -10,7 +10,7 @@ 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 / "tools")) sys.path.insert(0, str(ROOT / "bin" / "tools"))
from kblib import open_readonly # noqa: E402 from kblib import open_readonly # noqa: E402
from yamlout import to_yaml # noqa: E402 from yamlout import to_yaml # noqa: E402
+123 -21
View File
@@ -2,12 +2,15 @@
"""kb/index - build the 2dph brain from markdown + factual leafs. """kb/index - build the 2dph brain from markdown + factual leafs.
bin/kb/index [--corpus DIR] [--rebuild] [--limit N] bin/kb/index [--corpus DIR] [--rebuild] [--limit N]
bin/kb/index --rebuild --with-facts --with-chats
bin/kb/index --json # emit stats as JSON bin/kb/index --json # emit stats as JSON
Reads every .md under the corpus (default: repo root docs, skills, READMEs) Reads every .md under the corpus (default: repo root docs, skills, READMEs)
as `info` leafs, embeds them with model2vec (potion-multilingual-128M), and as `info` leafs, embeds them with model2vec (potion-multilingual-128M), and
writes them into var/kb.lbug with FTS + HNSW indexes. `facts` leafs come writes them into var/kb.lbug with FTS + HNSW indexes. `facts` leafs come
from bin/facts/extract (docker x compose x ssh-config pairing). from bin/facts/extract (docker × compose × ssh-config pairing) when
`--with-facts` is set. `--with-chats` indexes markdown under var/chats/md
(or a given dir) as info. WhatsApp sync stays out of v1.
--rebuild drops the database file and indexes from scratch. Without it a run --rebuild drops the database file and indexes from scratch. Without it a run
is idempotent (MERGE by (source,text) id). is idempotent (MERGE by (source,text) id).
@@ -19,13 +22,14 @@ 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 / "tools")) sys.path.insert(0, str(ROOT / "bin" / "tools"))
from kblib import ( # noqa: E402 from kblib import ( # noqa: E402
connect, create_fts_and_vector, init_schema, upsert_leaf, add_leafs, connect, ensure_indexes, init_schema, upsert_leaf, link_from_file,
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"]
@@ -80,10 +84,11 @@ def index_leafs(conn, leafs: list[dict], embed_fn, limit: int) -> tuple[int, int
for lf in leafs[:limit] if limit else leafs: for lf in leafs[:limit] if limit else leafs:
query = f"{lf['heading']}\n\n{lf['text']}" query = f"{lf['heading']}\n\n{lf['text']}"
emb = embed_fn(lf["text"]) if lf["text"] else None emb = embed_fn(lf["text"]) if lf["text"] else None
upsert_leaf(conn, text=query, root="info", confidence="confirmed", lid = upsert_leaf(conn, text=query, root="info", confidence="confirmed",
source=lf["source"], source_rev="working-tree", source=lf["source"], source_rev="working-tree",
how="kb/index", loc=lf["source"], type_=lf.get("type", "reference"), how="kb/index", loc=lf["source"], type_=lf.get("type", "reference"),
embedding=emb) embedding=emb)
link_from_file(conn, lid, lf["source"], repo=str(lf.get("repo") or ""))
count += 1 count += 1
return count, len(leafs) return count, len(leafs)
@@ -94,48 +99,145 @@ def embedder():
return lambda text: model.encode([text])[0].astype(float).tolist() return lambda text: model.encode([text])[0].astype(float).tolist()
def index_fact_dicts(conn, facts: list[dict], embed_fn) -> int:
"""Write extract-shaped dicts as root=facts leafs (2-source source field)."""
leafs = []
for f in facts:
text = str(f.get("text") or "")
source = str(f.get("source") or "")
if not text or not source:
continue
leafs.append({
"text": text,
"root": "facts",
"confidence": "confirmed",
"source": source,
"source_rev": f.get("source_rev") or "working-tree",
"how": f.get("how") or "facts/extract",
"loc": f.get("loc") or source,
"type": "fact",
"embedding": embed_fn(text) if text else None,
})
return len(add_leafs(conn, leafs))
def facts_from_extract() -> list[dict]:
import subprocess
proc = subprocess.run(
[sys.executable, str(ROOT / "bin" / "facts" / "extract"), "--json", "--dry-run"],
cwd=ROOT,
capture_output=True,
text=True,
check=False,
)
if proc.returncode != 0:
print(f"kb/index: facts/extract failed: {proc.stderr}", file=sys.stderr)
return []
try:
payload = json.loads(proc.stdout)
except json.JSONDecodeError:
print("kb/index: facts/extract produced non-JSON", file=sys.stderr)
return []
return list(payload.get("facts") or [])
def main(argv: list[str]) -> int: def main(argv: list[str]) -> int:
import argparse import argparse
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("--db", default="", help="path to kb.lbug (default var/kb.lbug)")
p.add_argument("--no-defaults", action="store_true", help="do not index repo README/docs/skills")
p.add_argument("--with-mail", action="store_true", help="include var/mail message.md leafs")
p.add_argument("--with-facts", action="store_true", help="run facts/extract into root=facts")
p.add_argument("--facts-json", default="", help="JSON list (or {facts:[...]}) of fact dicts")
p.add_argument(
"--with-chats",
nargs="?",
const=str(ROOT / "var" / "chats" / "md"),
default="",
help="index chat markdown as info (default var/chats/md)",
)
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(
"--skip-indexes",
action="store_true",
help="write leafs only; caller runs ensure_indexes after seeding facts",
)
p.add_argument("--limit", type=int, default=0, help="max leafs to embed") p.add_argument("--limit", type=int, default=0, help="max leafs to embed")
p.add_argument("--json", action="store_true") p.add_argument("--json", action="store_true")
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) dbpath = Path(a.db) if a.db else DB_PATH
leafs: list[dict] = [] if a.no_defaults else 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))
chat_n = 0
if a.with_chats:
chats = load_corpus_glob(a.with_chats)
chat_n = len(chats)
leafs.extend(chats)
mail_n = 0
if a.with_mail:
mail = from_mail_root(ROOT / "var" / "mail", since=a.since)
mail_n = len(mail)
leafs.extend(mail)
db, conn = connect(DB_PATH, read_only=False) facts: list[dict] = []
if a.facts_json:
raw = Path(a.facts_json).read_text(encoding="utf-8")
payload = json.loads(raw)
facts = list(payload.get("facts") if isinstance(payload, dict) else payload)
if a.with_facts:
facts.extend(facts_from_extract())
if a.dry_run:
msg = {
"indexed": 0,
"corpus_total": len(leafs),
"mail_leafs": mail_n,
"chat_leafs": chat_n,
"facts_leafs": len(facts),
"dry_run": True,
}
print(json.dumps(msg, indent=2) if a.json else
f"brain/index: {len(leafs)} info + {len(facts)} facts would be indexed")
return 0
VAR.mkdir(exist_ok=True)
dbpath.parent.mkdir(parents=True, exist_ok=True)
if a.rebuild and dbpath.exists():
dbpath.unlink()
db, conn = connect(dbpath, read_only=False)
init_schema(conn) init_schema(conn)
if not (a.rebuild or _already_indexed(conn)):
create_fts_and_vector(conn, force=True)
embed = embedder() embed = embedder()
done, total = index_leafs(conn, leafs, embed, a.limit) done, total = index_leafs(conn, leafs, embed, a.limit)
create_fts_and_vector(conn, force=(done > 0 or a.rebuild)) fact_n = index_fact_dicts(conn, facts, embed) if facts else 0
if not a.skip_indexes:
ensure_indexes(conn)
s = stats(conn) s = stats(conn)
conn.close() conn.close()
db.close() db.close()
result = {"indexed": done, "corpus_total": total, **{k: v for k, v in s.items() if k in ("total", "by_root")}} result = {
print(json.dumps(result, indent=2) if a.json else f"indexed {done}/{total} leafs; db total {s['total']}") "indexed": done,
"corpus_total": total,
"facts_leafs": fact_n,
"chat_leafs": chat_n,
**{k: v for k, v in s.items() if k in ("total", "by_root")},
}
if a.skip_indexes:
result["indexes"] = "skipped"
print(json.dumps(result, indent=2) if a.json else
f"indexed {done}/{total} info + {fact_n} facts; db total {s['total']}")
return 0 return 0
def _already_indexed(conn) -> bool:
try:
return conn.execute("MATCH (l:Leaf) RETURN count(*)").get_all()[0][0] > 0
except Exception:
return False
if __name__ == "__main__": if __name__ == "__main__":
sys.exit(main(sys.argv[1:])) sys.exit(main(sys.argv[1:]))
+30 -67
View File
@@ -1,72 +1,35 @@
#!/usr/bin/env python3 #!/usr/bin/env bash
"""kb/search - deduction search over the 2dph brain. # bin/kb/search deprecated wrapper. Use bin/brain/search.go.
# CGO via Zig (bin/cgo/zig), not gcc. Builds a binary then execs it.
set -euo pipefail
bin/kb/search "query" # hybrid facts+info, YAML out ROOT="$(CDPATH= cd -- "$(dirname "$0")/../.." && pwd)"
bin/kb/search "query" --root facts # confirmed facts only BIN="$ROOT/var/bin/brain-search"
bin/kb/search "query" --hop 1 # follow graph edges after hitting SRC="$ROOT/internal/brain"
bin/kb/search "query" --json | yq '.' CMD="$ROOT/bin/brain"
bin/kb/search "query" -n 5 # more results
Deduction order: facts root first (confirmed answers with evidence links), mkdir -p "$ROOT/var/bin"
then info root (marked `(not confirmed)`). --root restricts to one root.
--hop N walks FROM_FILE edges (sibling leafs in the same source file).
"""
from __future__ import annotations
import json need_build=0
import sys if [ ! -x "$BIN" ]; then
from pathlib import Path need_build=1
else
while IFS= read -r -d '' f; do
if [ "$f" -nt "$BIN" ]; then
need_build=1
break
fi
done < <(find "$SRC" "$CMD" -name '*.go' -print0 2>/dev/null)
fi
ROOT = Path(__file__).resolve().parents[2] if [ "$need_build" -eq 1 ]; then
sys.path.insert(0, str(ROOT / "tools")) echo "Building brain/search (zig cc)..." >&2
(
cd "$ROOT" &&
eval "$("$ROOT/bin/cgo/zig" env)" &&
go build -tags system_ladybug -o "$BIN" ./bin/brain
) || exit 1
fi
from kblib import connect, hybrid_search, init_schema, open_readonly, query_fts # noqa: E402 echo "bin/kb/search is deprecated; use bin/brain/search.go" >&2
from yamlout import to_yaml # noqa: E402 exec "$BIN" "$@"
import ladybug # noqa: E402
def main(argv: list[str]) -> int:
import argparse
p = argparse.ArgumentParser(description="deduction search over the brain")
p.add_argument("query")
p.add_argument("--root", choices=("facts", "info", None), default=None)
p.add_argument("--hop", type=int, default=0)
p.add_argument("-n", "--limit", type=int, default=10)
p.add_argument("--json", action="store_true")
a = p.parse_args(argv)
try:
db, conn = open_readonly()
except FileNotFoundError as e:
print(e, file=sys.stderr)
return 1
from model2vec import StaticModel
model = StaticModel.from_pretrained("minishlab/potion-multilingual-128M")
emb = model.encode([a.query])[0].astype(float).tolist()
rhs: list[dict] = []
try:
rhs = query_fts(conn, a.query, a.limit * 2)
except Exception:
rhs = []
results = hybrid_search(conn, emb, rhs, a.limit)
if a.root:
results = [h for h in results if h["root"] == a.root]
for hit in results:
hit.pop("rrf", None)
if hit.get("text"):
hit["snippet"] = hit["text"][:280]
out = {"query": a.query, "root_filter": a.root or "facts+info",
"count": len(results), "results": results}
print(json.dumps(out, indent=2, ensure_ascii=False) if a.json else to_yaml(out))
conn.close()
db.close()
return 0
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))
+1 -1
View File
@@ -11,7 +11,7 @@ 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 / "tools")) sys.path.insert(0, str(ROOT / "bin" / "tools"))
from kblib import open_readonly, stats # noqa: E402 from kblib import open_readonly, stats # noqa: E402
from yamlout import to_yaml # noqa: E402 from yamlout import to_yaml # noqa: E402
+23
View File
@@ -0,0 +1,23 @@
//usr/bin/env go run "$0" "$@"; exit
// bin/kb/watch.go — deprecated. Use bin/brain/watch.go.
//
// Usage:
//
// ./bin/kb/watch.go [dir...] # dirs default /corpus
// KB_WATCH_INTERVAL=15 ./bin/kb/watch.go
//
// Shebang trick: first line is a Go `//` comment; the real code lives in the
// importable package (module path, never a relative import).
// NOTE: never run `gofmt -w` on this file - it rewrites `//usr/bin/env` to
// `// usr/...` and breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/bin/watch"
)
func main() {
watch.Run(os.Args[1:])
}
+384
View File
@@ -0,0 +1,384 @@
#!/usr/bin/env python3
"""mail/import - pull OnlyOffice mails into var/mail/ as markdown.
bin/mail/import --from-raw var/mail convert Go-synced message.json to md
bin/mail/import import newest inbox messages
bin/mail/import --folder sent import sent folder
bin/mail/import --since 2026-01-01 only messages after a date
bin/mail/import --limit 50 cap messages per run
bin/mail/import --no-attachments body only, skip attachment conversion
bin/mail/import --ocr OCR images (PDFs OCR when textless)
bin/mail/import --dry-run list messages without writing anything
Writes one directory per message: var/mail/{folder}/{message_id}/
message.md frontmatter + markdown body
attachments/ raw attachment files (zips unpacked to _unpacked/)
attachments/*.md converted attachment content
Indexing is a separate step (`bin/brain/index.go --rebuild`): conversion can
crash and must not leave the brain DB mid-transaction.
Requires ONLYOFFICE_URL/USER/PASS in .env (or env) except `--from-raw`.
Idempotent: a message already present (message.md exists) is skipped unless
--force.
"""
from __future__ import annotations
import argparse
import json
import os
import re
import subprocess
import sys
import time
import urllib.parse
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(ROOT / "bin" / "tools"))
from mailconv import ( # noqa: E402
ARCHIVE_SUFFIXES,
IMAGE_SUFFIXES,
LEGACY_OFFICE_SUFFIXES,
TEXT_SUFFIXES,
convert_pdf,
html_to_markdown,
normalize_markdown,
ocr_image,
zip_extract_safe,
)
import requests # noqa: E402
FOLDER_IDS = {"inbox": 1, "sent": 2, "drafts": 3, "trash": 4, "spam": 5}
DEFAULT_LIMIT = 25
def load_env() -> dict:
env = {k: v for k, v in os.environ.items()}
envfile = ROOT / ".env"
if envfile.exists():
for line in envfile.read_text().splitlines():
line = line.strip()
if not line or line.startswith("#") or "=" not in line:
continue
k, _, v = line.partition("=")
env.setdefault(k.strip(), v.strip().strip("\"'"))
url = env.get("ONLYOFFICE_URL") or env.get("OO_URL")
user = env.get("ONLYOFFICE_USER") or env.get("OO_USER")
password = env.get("ONLYOFFICE_PASS") or env.get("OO_PASSWORD")
missing = [n for n, v in (("ONLYOFFICE_URL", url), ("ONLYOFFICE_USER", user),
("ONLYOFFICE_PASS", password)) if not v]
if missing:
sys.exit(f"mail/import: missing {', '.join(missing)} (need .env or env)")
return {"url": url.rstrip("/"), "user": user, "password": password}
class OOClient:
def __init__(self, conf: dict):
self.base = conf["url"]
self.session = requests.Session()
self.token = None
self._login(conf)
def _login(self, conf: dict) -> None:
r = self.session.post(f"{self.base}/api/2.0/authentication.json",
json={"userName": conf["user"], "password": conf["password"], "type": 0},
timeout=30)
r.raise_for_status()
body = r.json()
self.token = (body.get("response") or {}).get("token", "")
if not self.token:
sys.exit("mail/import: authentication failed (empty token)")
def _headers(self) -> dict:
return {"Authorization": f"Bearer {self.token}", "Accept": "application/json"}
def get(self, path: str, params: dict | None = None):
r = self.session.get(f"{self.base}{path}", params=params, headers=self._headers(), timeout=30)
r.raise_for_status()
return r.json()
def list_messages(self, folder: int, page: int = 1, count: int = DEFAULT_LIMIT) -> list[dict]:
data = self.get("/api/2.0/mail/messages",
params={"folder": folder, "page": page, "count": count})
return data.get("response", [])
def get_message(self, message_id: str) -> dict:
data = self.get(f"/api/2.0/mail/messages/{message_id}")
return data.get("response", {})
def download_attachment(self, attach_id, dest: Path) -> bool:
"""Download one attachment via the portal session cookie (.ashx handler)."""
url = f"{self.base}/addons/mail/httphandlers/download.ashx?attachid={attach_id}"
r = self.session.get(url, timeout=60)
if r.status_code != 200:
return False
dest.parent.mkdir(parents=True, exist_ok=True)
dest.write_bytes(r.content)
return True
def folder_id(name: str) -> int:
if name in FOLDER_IDS:
return FOLDER_IDS[name]
if name.isdigit():
return int(name)
sys.exit(f"mail/import: unknown folder '{name}' (use {', '.join(FOLDER_IDS)})")
def safe_attachment_name(att: dict) -> str:
name = att.get("fileName") or att.get("storedName") or "attachment"
name = re.sub(r"[^\w.\- ]+", "_", name)
return name
def convert_file_to_md(path: Path, ocr: bool) -> str | None:
"""Convert one attachment file to markdown text; None when not convertible."""
suffix = path.suffix.lower()
if suffix in TEXT_SUFFIXES:
return normalize_markdown(path.read_text(encoding="utf-8", errors="replace"))
if suffix in (".docx", ".pptx", ".xlsx", ".html", ".htm", ".epub", ".eml", ".msg"):
try:
from markitdown import MarkItDown
md = MarkItDown()
result = md.convert(str(path))
return normalize_markdown(result.text_content)
except Exception as e:
return f"\n<!-- conversion failed: {e} -->\n"
if suffix == ".pdf":
return convert_pdf(path, ocr)
if suffix in IMAGE_SUFFIXES and ocr:
return ocr_image(path) or "\n<!-- ocr unavailable -->\n"
if suffix in LEGACY_OFFICE_SUFFIXES:
return _convert_legacy(path)
if suffix in ARCHIVE_SUFFIXES:
return None # handled by caller (unpack + recurse)
return None
def _convert_legacy(path: Path) -> str:
"""Legacy .doc/.xls/.ppt -> md via pandoc (installed) or a stub."""
try:
out = subprocess.run(["pandoc", str(path), "-t", "markdown"],
capture_output=True, text=True, timeout=120)
if out.returncode == 0 and out.stdout.strip():
return normalize_markdown(out.stdout)
except (FileNotFoundError, subprocess.TimeoutExpired):
pass
return f"\n<!-- legacy {path.suffix} not convertible (pandoc unavailable) -->\n"
def write_message_md(msg: dict, folder: str, out_dir: Path, target_dir: Path | None = None) -> Path:
import yaml
body_html = msg.get("htmlBody") or ""
body_text = msg.get("textBody") or ""
body_md = ""
if body_html.strip():
body_md = html_to_markdown(body_html)
elif body_text.strip():
body_md = normalize_markdown(body_text)
# accept both OnlyOffice (receivedDate) and Go-sync (receivedAt) date keys
date = msg.get("receivedDate") or msg.get("receivedAt") or ""
if date and not isinstance(date, str):
date = str(date)
meta = {
"id": msg.get("id"),
"source": msg.get("source"),
"folder": folder,
"subject": msg.get("subject", ""),
"from": msg.get("from", ""),
"to": msg.get("to", ""),
"cc": msg.get("cc", ""),
"date": date,
"has_attachments": bool(msg.get("hasAttachments")),
"mime_message_id": msg.get("mimeMessageId", ""),
"calendar_uid": msg.get("calendarUid", ""),
"type": "mail",
}
meta = {k: v for k, v in meta.items() if v not in (None, "")}
frontmatter = "---\n" + yaml.safe_dump(meta, sort_keys=False, allow_unicode=True).strip() + "\n---\n"
content = f"{frontmatter}\n# {meta.get('subject','')}\n\n{body_md}".strip() + "\n"
if target_dir is not None:
msg_dir = target_dir
else:
msg_dir = out_dir / folder / str(meta.get("id"))
msg_dir.mkdir(parents=True, exist_ok=True)
md_path = msg_dir / "message.md"
md_path.write_text(content, encoding="utf-8")
return md_path
def convert_attachments(msg: dict, msg_dir: Path, ocr: bool) -> list[dict]:
"""Download + convert each attachment; returns [{name, md, raw}] summaries.
raw file keeps the API storedName (unique hash, avoids collisions); the
markdown is named after the friendly fileName when available.
In --from-raw mode attachments are already on disk (Go sync wrote them;
.ics already has a structured .md sidecar). Files with an existing .md
sidecar are left as-is, only unconverted raws are converted here.
"""
out: list[dict] = []
atts = msg.get("attachments") or []
att_dir = msg_dir / "attachments"
for att in atts:
aid = att.get("fileId")
display = safe_attachment_name(att)
stored = att.get("storedName")
raw_name = safe_attachment_name({"storedName": stored}) if stored else display
raw = att_dir / raw_name
if aid and not raw.exists() and OOCLIENT is not None:
if not OOCLIENT.download_attachment(aid, raw):
out.append({"name": display, "md": "\n<!-- download failed -->\n", "raw": str(raw)})
continue
if not raw.exists():
out.append({"name": display, "md": "\n<!-- raw missing -->\n", "raw": str(raw)})
continue
md_stem = Path(display).stem or raw.stem
md_path = att_dir / f"{md_stem}.md"
# Go sync pre-wrote structured .md for .ics; keep it.
if not md_path.exists():
md_text = _convert_att_recursive(raw, ocr)
md_path.write_text(f"# Attachment: {display}\n\n{md_text}\n", encoding="utf-8")
else:
md_text = md_path.read_text(encoding="utf-8", errors="replace")
out.append({"name": display, "md": md_text, "raw": str(raw), "md_file": str(md_path)})
return out
def _convert_att_recursive(path: Path, ocr: bool) -> str:
if path.suffix.lower() in ARCHIVE_SUFFIXES:
parts: list[str] = []
unpack = path.parent / "_unpacked" / path.stem
files = zip_extract_safe(path, unpack)
for f in files:
sub = _convert_att_recursive(f, ocr)
if sub and sub.strip():
parts.append(f"## {f.name}\n\n{sub}")
return "\n\n".join(parts) if parts else "\n<!-- empty zip -->\n"
text = convert_file_to_md(path, ocr)
return text or "\n<!-- not convertible -->\n"
# module-level client for attachment downloads in convert_attachments
OOCLIENT: OOClient | None = None
def convert_one(msg_dir: Path, full: dict, folder: str, out_root: Path,
ocr: bool, no_attachments: bool, target_dir: Path | None = None) -> dict:
"""Write message.md + convert attachments for one message dict.
Works for both live API messages and the Go-sync message.json shape
(source field optional; attachments read from attachments/ dir).
target_dir overrides the derived path (used by --from-raw where the
directory layout is authoritative, not the message folder field).
"""
mid = str(full.get("id"))
write_message_md(full, folder, out_root, target_dir=target_dir)
converted: list[dict] = []
if not no_attachments:
converted = convert_attachments(full, msg_dir, ocr)
return {"id": mid, "subject": full.get("subject", ""),
"date": full.get("receivedDate", "") or full.get("receivedAt", ""),
"attachments": len(converted)}
def main(argv: list[str]) -> int:
global OOCLIENT
p = argparse.ArgumentParser(description="pull OnlyOffice mails to var/mail as markdown")
p.add_argument("--folder", default="inbox", help="inbox|sent|drafts|trash|spam or numeric id")
p.add_argument("--limit", type=int, default=DEFAULT_LIMIT, help="max messages per run")
p.add_argument("--offset", type=int, default=0, help="skip N messages")
p.add_argument("--since", default="", help="only messages received after YYYY-MM-DD")
p.add_argument("--id", action="append", default=[], help="import specific message id (repeatable)")
p.add_argument("--from-raw", default="",
help="convert Go-synced dirs (var/mail/<folder>/<id>/message.json) to markdown")
p.add_argument("--no-attachments", action="store_true", help="skip attachment download+convert")
p.add_argument("--ocr", action="store_true", help="OCR images (PDFs OCR when textless)")
p.add_argument("--force", action="store_true", help="re-import even if message.md exists")
p.add_argument("--dry-run", action="store_true", help="list messages, write nothing")
p.add_argument("--json", action="store_true")
a = p.parse_args(argv)
fid = folder_id(a.folder)
out_root = ROOT / "var" / "mail"
summary: list[dict] = []
if a.from_raw:
OOCLIENT = None
raw_root = Path(a.from_raw)
for msg_dir in sorted(raw_root.rglob("message.json")):
mid = msg_dir.parent.name
entry = {"id": mid, "subject": "", "date": "",
"attachments": 0, "skipped": False}
md_path = msg_dir.parent / "message.md"
if md_path.exists() and not a.force:
entry["skipped"] = True
summary.append(entry)
continue
if a.dry_run:
entry["skipped"] = "dry-run"
summary.append(entry)
continue
full = json.loads(msg_dir.read_text(encoding="utf-8"))
entry.update(convert_one(msg_dir.parent, full, full.get("folder") or a.folder,
raw_root, a.ocr, a.no_attachments,
target_dir=msg_dir.parent))
summary.append(entry)
else:
conf = load_env()
OOCLIENT = OOClient(conf)
if a.id:
messages = [{"id": i} for i in a.id]
else:
page = 1
messages = []
want = a.offset + a.limit
while len(messages) < want:
count = min(DEFAULT_LIMIT, want - len(messages))
chunk = OOCLIENT.list_messages(fid, page=page, count=count)
if not chunk:
break
messages.extend(chunk)
if len(chunk) < count:
break
page += 1
messages = messages[a.offset:a.offset + a.limit]
if a.since:
messages = [m for m in messages
if (m.get("receivedDate") or "") >= a.since]
for m in messages:
mid = str(m.get("id"))
entry = {"id": mid, "subject": m.get("subject", ""),
"date": m.get("receivedDate", ""), "attachments": 0, "skipped": False}
msg_dir = out_root / a.folder / mid
md_path = msg_dir / "message.md"
if md_path.exists() and not a.force:
entry["skipped"] = True
summary.append(entry)
continue
if a.dry_run:
entry["skipped"] = "dry-run"
summary.append(entry)
continue
full = OOCLIENT.get_message(mid)
entry.update(convert_one(msg_dir, full, a.folder, out_root, a.ocr, a.no_attachments,
target_dir=msg_dir))
summary.append(entry)
if a.json:
print(json.dumps(summary, ensure_ascii=False, indent=2))
else:
imported = [e for e in summary if not e["skipped"]]
print(f"mail/import: folder={a.folder} checked={len(summary)} "
f"imported={len(imported)} (skipped={sum(e['skipped'] is True for e in summary)})")
for e in summary:
flag = "skip" if e["skipped"] is True else ("dry" if e["skipped"] == "dry-run" else "ok ")
print(f" [{flag}] {e['id']} {e['date'][:10]} {e['subject'][:60]}"
f" (atts={e['attachments']})")
return 0
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))
+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:]))
}
+27
View File
@@ -0,0 +1,27 @@
#!/usr/bin/env python3
"""mail/index_mail — deprecated. Use bin/brain/index.go --rebuild --with-mail.
Ladybug corrupts its WAL on bulk-insert into an already-indexed DB, so this
shim always rebuilds (repo corpus + var/mail). Conversion stays in mail/import.
"""
from __future__ import annotations
import os
import sys
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
def main(argv: list[str]) -> int:
print(
"bin/mail/index_mail is deprecated; use bin/brain/index.go --rebuild --with-mail",
file=sys.stderr,
)
index = ROOT / "bin" / "kb" / "index"
os.execv(sys.executable, [sys.executable, str(index), "--rebuild", "--with-mail", *argv])
return 1
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))
+48
View File
@@ -0,0 +1,48 @@
//usr/bin/env go run -tags=mail_ocr "$0" "$@"; exit
//go:build mail_ocr
//
// bin/mail/ocr.go - OCR an image or scanned PDF (tesseract eng+deu).
//
// ./bin/mail/ocr.go scan.png
// ./bin/mail/ocr.go scan.pdf
// OCR_ENGINE=paddle ./bin/mail/ocr.go scan.png
//
// PDFs try pdftotext -layout first; empty text layer uses pdftoppm + tesseract.
// No gocv. Tesseract CGO bindings are not used (D21 Zig owns Ladybug CGO).
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"fmt"
"os"
"strings"
"github.com/eSlider/2dph/internal/ocr"
)
func main() {
os.Exit(run(os.Args[1:]))
}
func run(args []string) int {
if len(args) != 1 || strings.HasPrefix(args[0], "-") {
fmt.Fprintln(os.Stderr, `usage: bin/mail/ocr.go <image|pdf>`)
return 2
}
path := args[0]
var (
text string
err error
)
if strings.HasSuffix(strings.ToLower(path), ".pdf") {
text, err = ocr.PDFFile(path)
} else {
text, err = ocr.ImageFile(path)
}
if err != nil {
fmt.Fprintf(os.Stderr, "mail/ocr: %v\n", err)
return 1
}
fmt.Println(text)
return 0
}
+25
View File
@@ -0,0 +1,25 @@
//usr/bin/env go run "$0" "$@"; exit
// bin/mail/sync.go - async download of OnlyOffice and Gmail mail to var/mail/.
//
// ./bin/mail/sync.go --source onlyoffice,gmail --limit 50 --workers 8
// ./bin/mail/sync.go --source gmail --force
// ./bin/mail/sync.go --dry-run
//
// Writes raw message.json + attachments under var/mail/<folder>/<id>/; run
// 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
// importable package (module path, never a relative import).
// NOTE: never run `gofmt -w` on this file - it rewrites `//usr/bin/env` to
// `// usr/...` and breaks the shebang.
package main
import (
"os"
"github.com/eSlider/2dph/bin/mail/sync"
)
func main() {
os.Exit(sync.Main(os.Args[1:]))
}
+156
View File
@@ -0,0 +1,156 @@
// Package synccmd wires the sync library to a CLI: reads .env, parses flags,
// picks sources, prints stats. Kept separate from the library so unit tests
// don't depend on os.Args/env.
package sync
import (
"context"
"flag"
"fmt"
"os"
"path/filepath"
"strings"
"time"
)
// CLIConfig is a superset of SyncConfig plus flag parsing results.
type CLIConfig struct {
Sync SyncConfig
Env string // .env path; default <cwd>/.env
Sources string
Help bool
}
// ParseCLI reads os.Args into a CLIConfig. Exit codes: 0 ok, 2 usage.
func ParseCLI(args []string) (CLIConfig, int, error) {
fs := flag.NewFlagSet("mail/sync", flag.ContinueOnError)
var (
env = fs.String("env", "", ".env file (default: <cwd>/.env)")
out = fs.String("out", "", "var/mail root (default: <cwd>/var/mail)")
workers = fs.Int("workers", 4, "concurrent downloads")
limit = fs.Int("limit", 0, "max messages per source (0 = all)")
offset = fs.Int("offset", 0, "skip first N messages per source")
force = fs.Bool("force", false, "overwrite existing message.json + attachments")
dryRun = fs.Bool("dry-run", false, "list message counts without writing")
query = fs.String("query", "in:inbox", "Gmail search query (gmail source only)")
srcs = fs.String("source", "onlyoffice", "comma list: onlyoffice,gmail (default onlyoffice)")
help = fs.Bool("help", false, "usage")
)
fs.SetOutput(os.Stderr)
if err := fs.Parse(args); err != nil {
return CLIConfig{}, 2, err
}
if *help || fs.NArg() > 0 {
return CLIConfig{Help: true}, 0, nil
}
wd, err := os.Getwd()
if err != nil {
return CLIConfig{}, 2, err
}
if *env == "" {
*env = filepath.Join(wd, ".env")
}
if *out == "" {
*out = filepath.Join(wd, "var", "mail")
}
envVars := readEnv(*env)
cfg := SyncConfig{
Out: *out,
Workers: *workers,
Limit: *limit,
Offset: *offset,
Force: *force,
DryRun: *dryRun,
Query: *query,
Policy: RetryPolicy{},
}
cli := CLIConfig{Sync: cfg, Env: *env, Sources: *srcs}
for _, s := range strings.Split(*srcs, ",") {
switch strings.TrimSpace(s) {
case "onlyoffice":
u := pick(envVars["ONLYOFFICE_URL"], envVars["OO_URL"])
user := pick(envVars["ONLYOFFICE_USER"], envVars["OO_USER"])
pass := pick(envVars["ONLYOFFICE_PASS"], envVars["OO_PASSWORD"])
if u == "" || user == "" || pass == "" {
return CLIConfig{}, 2, fmt.Errorf("onlyoffice source needs ONLYOFFICE_URL/USER/PASS in %s", *env)
}
cfg.OO = &OOConfig{URL: u, User: user, Password: pass}
case "gmail":
home, _ := os.UserHomeDir()
cfg.Gmail = &GmailCredentials{
CredentialsPath: filepath.Join(home, ".gmail-mcp", "credentials.json"),
KeysPath: filepath.Join(home, ".gmail-mcp", "gcp-oauth.keys.json"),
}
default:
return CLIConfig{}, 2, fmt.Errorf("unknown source %q", s)
}
}
cli.Sync = cfg
return cli, 0, nil
}
// Main is the CLI entry: returns process exit code.
func Main(args []string) int {
cli, code, err := ParseCLI(args)
if err != nil {
fmt.Fprintln(os.Stderr, "mail/sync:", err)
return code
}
if cli.Help {
fmt.Fprintln(os.Stderr, "usage: bin/mail/sync.go [--source onlyoffice,gmail] [--query GMAIL_Q] [--limit N] [--offset N] [--workers N] [--force] [--dry-run]")
return 0
}
ctx, cancel := context.WithTimeout(context.Background(), 6*time.Hour)
defer cancel()
start := time.Now()
stats, err := Run(ctx, cli.Sync)
if err != nil {
fmt.Fprintln(os.Stderr, "mail/sync:", err)
return 1
}
if cli.Sync.DryRun {
fmt.Printf("mail/sync: dry-run checked=%d (no writes)\n", stats.Checked)
return 0
}
fmt.Printf("mail/sync: checked=%d new=%d skipped=%d failed=%d in %s\n",
stats.Checked, stats.New, stats.Skipped, stats.Failed, time.Since(start).Round(time.Millisecond))
if stats.Failed > 0 {
return 1
}
return 0
}
// readEnv parses KEY=VALUE lines (ignoring comments) with KEY=PATH override.
func readEnv(path string) map[string]string {
out := map[string]string{}
b, err := os.ReadFile(path)
if err != nil {
return out
}
for _, line := range strings.Split(string(b), "\n") {
line = strings.TrimSpace(line)
if line == "" || strings.HasPrefix(line, "#") || !strings.Contains(line, "=") {
continue
}
k, v, _ := strings.Cut(line, "=")
out[strings.TrimSpace(k)] = strings.Trim(strings.TrimSpace(v), "\"'")
}
// env overrides file
for _, kv := range os.Environ() {
k, v, ok := strings.Cut(kv, "=")
if !ok {
continue
}
if strings.HasPrefix(k, "ONLYOFFICE_") || strings.HasPrefix(k, "OO_") {
out[k] = v
}
}
return out
}
func pick(a, b string) string {
if a != "" {
return a
}
return b
}
+356
View File
@@ -0,0 +1,356 @@
package sync
import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"os"
"path/filepath"
"strings"
"time"
)
// GmailCredentials holds the OAuth files produced by the gmail MCP
// (@gongrzhe/server-gmail-autoauth-mcp) auto-auth flow.
type GmailCredentials struct {
CredentialsPath string // ~/.gmail-mcp/credentials.json
KeysPath string // ~/.gmail-mcp/gcp-oauth.keys.json
User string // fixed: the authed account
}
// gmailToken is the JSON shape of credentials.json + refresh response.
type gmailToken struct {
AccessToken string `json:"access_token"`
RefreshToken string `json:"refresh_token"`
Expiry int64 `json:"expiry_date"` // ms epoch
}
type gmailKeys struct {
Installed *gmailKeyBlock `json:"installed"`
Web *gmailKeyBlock `json:"web"`
}
type gmailKeyBlock struct {
ClientID string `json:"client_id"`
ClientSecret string `json:"client_secret"`
}
// GmailClient talks to the Gmail REST API using the OAuth refresh token from
// ~/.gmail-mcp/. Token is refreshed lazily with a mutex-guarded cache.
type GmailClient struct {
creds GmailCredentials
client *http.Client
mu chan struct{}
token *gmailToken
user string
}
func NewGmailClient(creds GmailCredentials) (*GmailClient, error) {
if creds.CredentialsPath == "" {
home, _ := os.UserHomeDir()
creds.CredentialsPath = filepath.Join(home, ".gmail-mcp", "credentials.json")
creds.KeysPath = filepath.Join(home, ".gmail-mcp", "gcp-oauth.keys.json")
}
g := &GmailClient{
creds: creds,
client: &http.Client{Timeout: 60 * time.Second},
mu: make(chan struct{}, 1),
}
g.mu <- struct{}{}
return g, nil
}
// accessToken returns a fresh bearer token, refreshing via the Google token
// endpoint when the cached one is missing or about to expire.
func (g *GmailClient) accessToken(ctx context.Context) (string, error) {
select {
case <-g.mu:
case <-ctx.Done():
return "", ctx.Err()
}
defer func() { g.mu <- struct{}{} }()
if g.token != nil && g.token.AccessToken != "" && g.token.Expiry > time.Now().UnixMilli()+300_000 {
return g.token.AccessToken, nil
}
return g.refreshLocked(ctx)
}
func (g *GmailClient) refreshLocked(ctx context.Context) (string, error) {
cred, err := os.ReadFile(g.creds.CredentialsPath)
if err != nil {
return "", fmt.Errorf("read gmail credentials %s: %w", g.creds.CredentialsPath, err)
}
var t gmailToken
if err := json.Unmarshal(cred, &t); err != nil {
return "", fmt.Errorf("parse gmail credentials: %w", err)
}
if t.RefreshToken == "" {
return "", errors.New("gmail credentials.json has no refresh_token (run the gmail MCP auth flow)")
}
keys, err := os.ReadFile(g.creds.KeysPath)
if err != nil {
return "", fmt.Errorf("read gmail keys %s: %w", g.creds.KeysPath, err)
}
var k gmailKeys
if err := json.Unmarshal(keys, &k); err != nil {
return "", fmt.Errorf("parse gmail keys: %w", err)
}
block := k.Installed
if block == nil {
block = k.Web
}
if block == nil {
return "", errors.New("gmail gcp-oauth.keys.json has no installed/web block")
}
form := url.Values{}
form.Set("client_id", block.ClientID)
form.Set("client_secret", block.ClientSecret)
form.Set("refresh_token", t.RefreshToken)
form.Set("grant_type", "refresh_token")
req, err := http.NewRequestWithContext(ctx, http.MethodPost, "https://oauth2.googleapis.com/token",
strings.NewReader(form.Encode()))
if err != nil {
return "", err
}
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
resp, err := g.client.Do(req)
if err != nil {
return "", fmt.Errorf("gmail token refresh: %w", err)
}
defer resp.Body.Close()
body, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
if resp.StatusCode != http.StatusOK {
var e struct {
Error string `json:"error"`
Desc string `json:"error_description"`
}
_ = json.Unmarshal(body, &e)
if e.Error == "invalid_grant" {
return "", fmt.Errorf("gmail OAuth token invalid/expired - re-auth via: npx -y @gongrzhe/server-gmail-autoauth-mcp auth (uses ~/.gmail-mcp)")
}
return "", fmt.Errorf("gmail token refresh status %d: %s", resp.StatusCode, truncate(string(body), 300))
}
var out struct {
AccessToken string `json:"access_token"`
ExpiresIn int64 `json:"expires_in"`
}
if err := json.Unmarshal(body, &out); err != nil {
return "", fmt.Errorf("gmail token refresh parse: %w", err)
}
g.token = &gmailToken{
AccessToken: out.AccessToken,
RefreshToken: t.RefreshToken,
Expiry: time.Now().UnixMilli() + out.ExpiresIn*1000,
}
return out.AccessToken, nil
}
// ListIDs returns message ids matching q, walking nextPageToken up to maxIDs
// (0 = unlimited). Thread-level pagination via the messages.list endpoint.
func (g *GmailClient) ListIDs(ctx context.Context, q string, maxIDs int, pageToken string) (ids []string, next string, err error) {
for {
params := url.Values{}
params.Set("q", q)
params.Set("maxResults", "100")
if pageToken != "" {
params.Set("pageToken", pageToken)
}
var out struct {
Messages []struct {
ID string `json:"id"`
} `json:"messages"`
NextPageToken string `json:"nextPageToken"`
}
if err := g.getJSON(ctx, "/gmail/v1/users/me/messages?"+params.Encode(), &out); err != nil {
return nil, "", err
}
for _, m := range out.Messages {
ids = append(ids, m.ID)
if maxIDs > 0 && len(ids) >= maxIDs {
return ids, out.NextPageToken, nil
}
}
if out.NextPageToken == "" {
break
}
pageToken = out.NextPageToken
}
return ids, "", nil
}
// GetMessage fetches a message in format=full and normalizes it.
func (g *GmailClient) GetMessage(ctx context.Context, id string) (*Message, error) {
var raw struct {
ID string `json:"id"`
ThreadID string `json:"threadId"`
InternalDate string `json:"internalDate"` // ms epoch string
Payload gmailPart
}
path := "/gmail/v1/users/me/messages/" + url.PathEscape(id) + "?format=full"
if err := g.getJSON(ctx, path, &raw); err != nil {
return nil, err
}
m := &Message{
Source: "gmail",
ID: raw.ID,
Folder: "gmail",
}
for _, h := range raw.Payload.Headers {
switch strings.ToLower(h.Name) {
case "subject":
m.Subject = h.Value
case "from":
m.From = h.Value
case "to":
m.To = h.Value
case "cc":
m.CC = h.Value
case "bcc":
m.BCC = h.Value
case "message-id":
m.MimeMessageID = h.Value
case "date":
if t, err := time.Parse(time.RFC1123Z, h.Value); err == nil {
m.ReceivedAt = t
}
}
}
if ms, err := parseMS(raw.InternalDate); err == nil {
m.ReceivedAt = ms
}
m.TextBody, m.HTMLBody, m.Attachments = collectParts(raw.Payload, "root", m.ID, 0)
m.HasAttachments = len(m.Attachments) > 0
return m, nil
}
type gmailPart struct {
PartID string `json:"partId"`
MimeType string `json:"mimeType"`
Filename string `json:"filename"`
Body gmailBody `json:"body"`
Headers []gmailHeader `json:"headers"`
Parts []gmailPart `json:"parts"`
}
type gmailHeader struct {
Name string `json:"name"`
Value string `json:"value"`
}
type gmailBody struct {
Size int64 `json:"size"`
Data string `json:"data"`
AttachmentID string `json:"attachmentId"`
}
// collectParts walks the MIME tree: text bodies into plain/html, anything with
// a filename into attachments (returned with base64 ids for later download).
func collectParts(p gmailPart, mime string, msgID string, depth int) (text, html string, atts []Attachment) {
if depth > 16 {
return
}
mt := strings.ToLower(p.MimeType)
if p.Filename != "" && mt != "text/plain" && mt != "text/html" {
pid := p.PartID
if pid == "" {
pid = fmt.Sprintf("%d", depth)
}
// Gmail's attachments API keys off body.attachmentId, not partId.
attID := p.Body.AttachmentID
if attID == "" {
attID = pid
}
atts = append(atts, Attachment{
FileID: msgID + ":" + attID,
FileName: p.Filename,
StoredName: p.Filename,
Size: p.Body.Size,
ContentType: p.MimeType,
})
} else if data, err := base64.URLEncoding.DecodeString(p.Body.Data); err == nil && len(p.Body.Data) > 0 {
s := string(data)
if mt == "text/html" && html == "" {
html = s
} else if (mt == "text/plain" || mt == "") && text == "" {
text = s
}
}
for _, child := range p.Parts {
t, h, a := collectParts(child, mt, msgID, depth+1)
if text == "" {
text = t
}
if html == "" {
html = h
}
atts = append(atts, a...)
}
return
}
// DownloadAttachment fetches an attachment's bytes from the Gmail API.
func (g *GmailClient) DownloadAttachment(ctx context.Context, msgID, attID string) ([]byte, error) {
// attID format is "<msgId>:<partId>"; the API needs the bare attachment id.
partID := attID
if i := strings.Index(attID, ":"); i >= 0 {
partID = attID[i+1:]
}
var out struct {
Data string `json:"data"`
}
path := "/gmail/v1/users/me/messages/" + url.PathEscape(msgID) + "/attachments/" + url.PathEscape(partID)
if err := g.getJSON(ctx, path, &out); err != nil {
return nil, err
}
return base64.URLEncoding.DecodeString(out.Data)
}
func (g *GmailClient) getJSON(ctx context.Context, path string, out any) error {
tok, err := g.accessToken(ctx)
if err != nil {
return err
}
u := "https://gmail.googleapis.com" + path
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+tok)
resp, err := g.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
body, _ := io.ReadAll(io.LimitReader(resp.Body, 16<<20))
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("gmail %s: status %d: %s", path, resp.StatusCode, truncate(string(body), 300))
}
if out != nil {
return json.Unmarshal(body, out)
}
return nil
}
func parseMS(s string) (time.Time, error) {
if s == "" {
return time.Time{}, errors.New("empty")
}
var ms int64
if _, err := fmt.Sscanf(s, "%d", &ms); err != nil {
return time.Time{}, err
}
return time.UnixMilli(ms), nil
}
func truncate(s string, n int) string {
if len(s) <= n {
return s
}
return s[:n] + "…"
}
var _ = bytes.MinRead
+258
View File
@@ -0,0 +1,258 @@
package sync
import (
"fmt"
"strings"
"time"
"unicode/utf8"
ics "github.com/arran4/golang-ical"
"golang.org/x/text/encoding/charmap"
)
// ICSToMarkdown parses a VCALENDAR/VEVENT payload and renders a compact
// structured markdown block: what / when / where / organizer / attendees.
// Returns the raw text when the payload is not a calendar.
func ICSToMarkdown(data []byte) string {
data = normalizeEncoding(data)
cal, err := ics.ParseCalendar(strings.NewReader(string(data)))
if err != nil {
return normalizeMarkdown(string(data))
}
method := ""
for _, p := range cal.CalendarProperties {
if p.IANAToken == string(ics.ComponentPropertyMethod) {
method = p.Value
break
}
}
method = strings.TrimSpace(method)
var out []string
for _, ev := range cal.Events() {
summary := strings.TrimSpace(propValue(ev, ics.ComponentPropertySummary))
if summary != "" {
out = append(out, "# "+summary)
}
if when := eventWhen(ev); when != "" {
out = append(out, "- **When:** "+when)
}
if loc := strings.TrimSpace(propValue(ev, ics.ComponentPropertyLocation)); loc != "" {
out = append(out, "- **Where:** "+loc)
}
if desc := strings.TrimSpace(stripHTML(propValue(ev, ics.ComponentPropertyDescription))); desc != "" {
out = append(out, "- **What:** "+desc)
}
if org := propValue(ev, ics.ComponentPropertyOrganizer); org != "" {
out = append(out, "- **Organizer:** "+attendeeFmt(org))
}
for _, a := range ev.Attendees() {
cn := strings.TrimSpace(firstParam(a.ICalParameters, "CN"))
partstat := string(a.ParticipationStatus())
name := cn
if name == "" {
name = a.Email()
}
line := name
if email := a.Email(); email != "" && email != name {
line = name + " <" + email + ">"
}
if partstat != "" && !strings.EqualFold(partstat, "NEEDS-ACTION") {
line += " (" + strings.Title(strings.ToLower(strings.ReplaceAll(partstat, "_", " "))) + ")"
}
out = append(out, "- **Attendee:** "+line)
}
}
if len(out) == 0 {
return normalizeMarkdown(string(data))
}
if method != "" {
out = append([]string{"*Calendar method: " + method + "*"}, out...)
}
return normalizeMarkdown(strings.Join(out, "\n\n"))
}
func eventWhen(ev *ics.VEvent) string {
start, errStart := ev.GetStartAt()
end, errEnd := ev.GetEndAt()
// All-day events: golang-ical has dedicated getters.
if errStart != nil {
if allDay, err := ev.GetAllDayStartAt(); err == nil {
start = allDay
errStart = nil
}
}
if errEnd != nil {
if allDay, err := ev.GetAllDayEndAt(); err == nil {
end = allDay
errEnd = nil
}
}
if errStart != nil {
// Non-IANA TZID (e.g. "W. Europe Standard Time"): parse the raw
// property text instead of failing.
return rawWhen(ev)
}
if errEnd != nil || end.Equal(start) {
return dtFmt(start)
}
return dtFmt(start) + " → " + dtFmt(end)
}
// rawWhen parses DTSTART/DTEND property values that golang-ical cannot resolve
// because the TZID is not an IANA zone. Formats: 20260812T120000 or 20260812.
func rawWhen(ev *ics.VEvent) string {
start := rawPropValue(ev, ics.ComponentPropertyDtStart)
end := rawPropValue(ev, ics.ComponentPropertyDtEnd)
if start == "" {
return ""
}
if end == "" || end == start {
return rawDTFmt(start)
}
return rawDTFmt(start) + " → " + rawDTFmt(end)
}
func rawPropValue(ev *ics.VEvent, prop ics.ComponentProperty) string {
p := ev.GetProperty(prop)
if p == nil {
return ""
}
return p.Value
}
// rawDTFmt turns 20260812T120000 into 2026-08-12 12:00; 20260812 into 2026-08-12.
func rawDTFmt(s string) string {
s = strings.TrimSpace(s)
if len(s) >= 8 && isDigits(s[:8]) {
y, m, d := s[:4], s[4:6], s[6:8]
if len(s) > 8 && (s[8] == 'T' || s[8] == 't') && len(s) >= 15 && isDigits(s[9:15]) {
h, mi := s[9:11], s[11:13]
return fmt.Sprintf("%s-%s-%s %s:%s", y, m, d, h, mi)
}
return fmt.Sprintf("%s-%s-%s", y, m, d)
}
return s
}
func isDigits(s string) bool {
for _, c := range s {
if c < '0' || c > '9' {
return false
}
}
return s != ""
}
// dtFmt renders a time as local "2006-01-02 15:04" (tz label when meaningful).
func dtFmt(t time.Time) string {
loc := t.Local()
label := ""
if loc.Location() != time.Local {
label = " " + loc.Location().String()
}
return loc.Format("2006-01-02 15:04") + label
}
// propertyGetter is satisfied by both *ics.Calendar and *ics.VEvent.
type propertyGetter interface {
GetProperty(ics.ComponentProperty) *ics.IANAProperty
}
func propValue(ev propertyGetter, prop ics.ComponentProperty) string {
p := ev.GetProperty(prop)
if p == nil {
return ""
}
return p.Value
}
func firstParam(params map[string][]string, key string) string {
if vs, ok := params[key]; ok && len(vs) > 0 {
return vs[0]
}
return ""
}
func attendeeFmt(raw string) string {
raw = strings.TrimSpace(raw)
if i := strings.Index(raw, ":"); i >= 0 {
raw = raw[i+1:]
}
return raw
}
// stripHTML removes tags and decodes entities from an ics DESCRIPTION that may
// carry HTML (Outlook/Exchange style), keeping text lines readable.
func stripHTML(s string) string {
if !strings.Contains(s, "<") {
return s
}
lines := strings.Split(s, "\n")
for i, l := range lines {
var b strings.Builder
depth := 0
for j := 0; j < len(l); j++ {
c := l[j]
if c == '<' {
if j+1 < len(l) && l[j+1] == '/' {
depth--
} else {
depth++
}
for j < len(l) && l[j] != '>' {
j++
}
continue
}
if c == '>' {
continue
}
if depth == 0 {
b.WriteByte(c)
}
}
lines[i] = strings.TrimSpace(b.String())
}
return strings.Join(lines, "\n")
}
// normalizeMarkdown collapses blank-line runs and strips control chars.
func normalizeMarkdown(s string) string {
s = strings.ReplaceAll(s, "\x00", "")
for _, ch := range []string{"\ufeff", "\u200b", "\u034f", "\u00ad", "\u2007", "\u2008", "\u200a", "\u2002"} {
s = strings.ReplaceAll(s, ch, "")
}
lines := strings.Split(s, "\n")
var out []string
blank := 0
for _, l := range lines {
if strings.TrimSpace(l) == "" {
blank++
if blank > 1 {
continue
}
} else {
blank = 0
}
out = append(out, l)
}
return strings.Join(out, "\n")
}
// normalizeEncoding re-encodes legacy single-byte text as UTF-8. ICS files
// exported by some portals are Latin-1 (e.g. "N\xfcrnberg"); golang-ical
// passes the bytes through, producing invalid UTF-8 in the markdown output.
// Valid UTF-8 is returned untouched.
func normalizeEncoding(data []byte) []byte {
if utf8.Valid(data) {
return data
}
dec := charmap.ISO8859_1.NewDecoder()
out, err := dec.Bytes(data)
if err != nil {
return data
}
return out
}
var _ = fmt.Sprintf // keep fmt import if helpers change
+222
View File
@@ -0,0 +1,222 @@
package sync
import (
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"net/http/cookiejar"
"net/url"
"strings"
"time"
)
// OOConfig mirrors the .env / environment used by bin/mail/import.
type OOConfig struct {
URL string
User string
Password string
}
// OOClient is a minimal OnlyOffice API client: authentication.json for the
// bearer token plus the session cookie jar required by the .ashx download
// handler. It mirrors the endpoint contract bin/mail/import already uses.
type OOClient struct {
cfg OOConfig
client *http.Client
mu chan struct{}
token string
folderID int
}
func NewOOClient(cfg OOConfig, folderID int) (*OOClient, error) {
jar, err := cookiejar.New(nil)
if err != nil {
return nil, err
}
c := &OOClient{
cfg: cfg,
client: &http.Client{Jar: jar, Timeout: 60 * time.Second},
mu: make(chan struct{}, 1),
folderID: folderID,
}
c.mu <- struct{}{}
if err := c.authenticate(context.Background()); err != nil {
return nil, err
}
return c, nil
}
func (o *OOClient) authenticate(ctx context.Context) error {
select {
case <-o.mu:
case <-ctx.Done():
return ctx.Err()
}
defer func() { o.mu <- struct{}{} }()
body, _ := json.Marshal(map[string]any{
"userName": o.cfg.User, "password": o.cfg.Password, "type": 0,
})
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
strings.TrimRight(o.cfg.URL, "/")+"/api/2.0/authentication.json",
strings.NewReader(string(body)))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
resp, err := o.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
data, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<20))
if resp.StatusCode < 200 || resp.StatusCode > 299 {
return fmt.Errorf("oo authenticate status %d: %s", resp.StatusCode, truncate(string(data), 200))
}
var out struct {
Response struct {
Token string `json:"token"`
} `json:"response"`
}
if err := json.Unmarshal(data, &out); err != nil {
return err
}
if out.Response.Token == "" {
return fmt.Errorf("oo authenticate: empty token")
}
o.token = out.Response.Token
return nil
}
// get performs an authenticated GET and decodes the JSON body into out.
func (o *OOClient) get(ctx context.Context, path string, out any) error {
u := strings.TrimRight(o.cfg.URL, "/") + path
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+o.token)
req.Header.Set("Accept", "application/json")
resp, err := o.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
data, _ := io.ReadAll(io.LimitReader(resp.Body, 16<<20))
if resp.StatusCode < 200 || resp.StatusCode > 299 {
return fmt.Errorf("oo %s: status %d: %s", path, resp.StatusCode, truncate(string(data), 300))
}
if out != nil {
return json.Unmarshal(data, out)
}
return nil
}
// ooMessage mirrors the OnlyOffice mail message JSON (subset we need).
type ooMessage struct {
ID int `json:"id"`
Subject string `json:"subject"`
From string `json:"from"`
To string `json:"to"`
CC string `json:"cc"`
BCC string `json:"bcc"`
ReceivedDate string `json:"receivedDate"`
HTMLBody string `json:"htmlBody"`
TextBody string `json:"textBody"`
HasAttachments bool `json:"hasAttachments"`
MimeMessageID string `json:"mimeMessageId"`
Attachments []struct {
FileID int `json:"fileId"`
FileName string `json:"fileName"`
StoredName string `json:"storedName"`
Size int64 `json:"size"`
ContentType string `json:"contentType"`
} `json:"attachments"`
}
// ListIDs returns message ids in the configured folder, paginating pages until
// maxIDs is reached (0 = all).
func (o *OOClient) ListIDs(ctx context.Context, maxIDs int, page int) (ids []int, next int, err error) {
var out struct {
Response []ooMessage `json:"response"`
}
count := 100
if maxIDs > 0 && maxIDs < count {
count = maxIDs
}
path := fmt.Sprintf("/api/2.0/mail/messages?folder=%d&page=%d&count=%d", o.folderID, page, count)
if err := o.get(ctx, path, &out); err != nil {
return nil, 0, err
}
for _, m := range out.Response {
ids = append(ids, m.ID)
if maxIDs > 0 && len(ids) >= maxIDs {
break
}
}
next = page + 1
return ids, next, nil
}
// GetMessage fetches the full message by id and normalizes into Message.
func (o *OOClient) GetMessage(ctx context.Context, id int) (*Message, error) {
var out struct {
Response ooMessage `json:"response"`
}
path := fmt.Sprintf("/api/2.0/mail/messages/%d", id)
if err := o.get(ctx, path, &out); err != nil {
return nil, err
}
m := out.Response
msg := &Message{
Source: "onlyoffice",
ID: fmt.Sprintf("%d", m.ID),
Folder: "oo",
Subject: m.Subject,
From: m.From,
To: m.To,
CC: m.CC,
BCC: m.BCC,
HTMLBody: m.HTMLBody,
TextBody: m.TextBody,
HasAttachments: m.HasAttachments,
MimeMessageID: m.MimeMessageID,
}
if t, err := time.Parse(time.RFC3339Nano, m.ReceivedDate); err == nil {
msg.ReceivedAt = t
}
for _, a := range m.Attachments {
msg.Attachments = append(msg.Attachments, Attachment{
FileID: fmt.Sprintf("%d", a.FileID),
FileName: a.FileName,
StoredName: a.StoredName,
Size: a.Size,
ContentType: a.ContentType,
})
}
return msg, nil
}
// DownloadAttachment fetches attachment bytes via the .ashx handler, which
// requires the session cookie (client.Jar) captured during authenticate().
func (o *OOClient) DownloadAttachment(ctx context.Context, fileID string) ([]byte, error) {
u := strings.TrimRight(o.cfg.URL, "/") + "/addons/mail/httphandlers/download.ashx?attachid=" + url.QueryEscape(fileID)
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
if err != nil {
return nil, err
}
resp, err := o.client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
data, err := io.ReadAll(io.LimitReader(resp.Body, 64<<20))
if err != nil {
return nil, err
}
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("oo download attach %s: status %d", fileID, resp.StatusCode)
}
return data, nil
}
+397
View File
@@ -0,0 +1,397 @@
package sync
import (
"context"
"encoding/json"
"errors"
"fmt"
"math"
"math/rand"
"os"
"path/filepath"
"strings"
"sync"
"sync/atomic"
"time"
)
// RetryPolicy is the exponential-backoff strategy applied to transient HTTP
// failures (5xx, timeouts, network errors). Callers wrap transient errors with
// retryWrap; everything else aborts immediately.
type RetryPolicy struct {
MaxAttempts int // total attempts (>=1); 0 => 5
BaseDelay time.Duration // first backoff; 0 => 250ms
MaxDelay time.Duration // cap; 0 => 15s
Jitter float64 // 0..1 multiplier; 0 => 0.2
}
func (p RetryPolicy) withDefaults() RetryPolicy {
if p.MaxAttempts <= 0 {
p.MaxAttempts = 5
}
if p.BaseDelay <= 0 {
p.BaseDelay = 250 * time.Millisecond
}
if p.MaxDelay <= 0 {
p.MaxDelay = 15 * time.Second
}
if p.Jitter <= 0 {
p.Jitter = 0.2
}
return p
}
// delay returns the wait before attempt n (1-based): base * 2^(n-2) + jitter,
// capped at MaxDelay. Attempt 1 waits 0, attempt 2 waits base, then doubles.
func (p RetryPolicy) delay(attempt int) time.Duration {
if attempt <= 1 {
return 0
}
exp := math.Min(float64(attempt-2), 10)
d := float64(p.BaseDelay) * math.Pow(2, exp)
if p.Jitter > 0 {
d *= 1 - p.Jitter + 2*p.Jitter*rand.Float64()
}
if d > float64(p.MaxDelay) {
d = float64(p.MaxDelay)
}
return time.Duration(d)
}
type errRetry struct{ err error }
func (e *errRetry) Error() string { return e.err.Error() }
func (e *errRetry) Unwrap() error { return e.err }
func isRetriable(err error) bool {
var r *errRetry
return errors.As(err, &r)
}
func retryWrap(err error) error {
if err == nil {
return nil
}
if isRetriable(err) {
return err
}
return &errRetry{err: err}
}
// Retry runs fn up to MaxAttempts times with exponential backoff between
// attempts. Non-retriable errors abort immediately. Returns the last error.
func Retry(ctx context.Context, policy RetryPolicy, fn func() error) error {
policy = policy.withDefaults()
var err error
for attempt := 1; attempt <= policy.MaxAttempts; attempt++ {
if err = fn(); err == nil {
return nil
}
if !isRetriable(err) {
return err
}
if attempt == policy.MaxAttempts {
return fmt.Errorf("after %d attempts: %w", policy.MaxAttempts, err)
}
select {
case <-time.After(policy.delay(attempt)):
case <-ctx.Done():
return ctx.Err()
}
}
return err
}
// SyncConfig wires up a sync run.
type SyncConfig struct {
OO *OOConfig // OnlyOffice source (optional)
Gmail *GmailCredentials // Gmail source (optional)
Out string // var/mail root; default <repo>/var/mail
Workers int // concurrency; default 4
Limit int // max messages per source (0 = all)
Offset int // skip first N messages per source
Force bool // overwrite existing message.json + attachments
DryRun bool // list without writing
Query string // Gmail search query; default in:inbox
Policy RetryPolicy
}
// SyncStats is returned by Run.
type SyncStats struct {
Checked int
New int32
Failed int32
Skipped int32
}
// Source abstracts the two backends for the worker pool.
type Source interface {
// ListIDs yields ids (string form) to fetch. cursor resumes pagination.
ListIDs(ctx context.Context, limit int, cursor string) (ids []string, next string, err error)
Get(ctx context.Context, id string) (*Message, error)
DownloadAttachment(ctx context.Context, msg *Message, att Attachment) ([]byte, error)
Folder() string
}
type ooSource struct {
c *OOClient
page int
}
// gmailAPI is the Gmail client surface gmailSource needs. *GmailClient implements it.
type gmailAPI interface {
ListIDs(ctx context.Context, q string, maxIDs int, pageToken string) ([]string, string, error)
GetMessage(ctx context.Context, id string) (*Message, error)
DownloadAttachment(ctx context.Context, msgID, attID string) ([]byte, error)
}
type gmailSource struct {
c gmailAPI
cur string
query string
}
func (s *ooSource) Folder() string { return "inbox" }
func (s *gmailSource) Folder() string { return "gmail" }
func (s *ooSource) ListIDs(ctx context.Context, limit int, cursor string) ([]string, string, error) {
page := s.page
if page == 0 {
page = 1
}
ids, next, err := s.c.ListIDs(ctx, limit, page)
s.page = next
strs := make([]string, len(ids))
for i, id := range ids {
strs[i] = fmt.Sprintf("%d", id)
}
return strs, "", err
}
func (s *ooSource) Get(ctx context.Context, id string) (*Message, error) {
var mid int
if _, err := fmt.Sscanf(id, "%d", &mid); err != nil {
return nil, fmt.Errorf("oo id %q: %w", id, err)
}
return s.c.GetMessage(ctx, mid)
}
func (s *ooSource) DownloadAttachment(ctx context.Context, msg *Message, att Attachment) ([]byte, error) {
return s.c.DownloadAttachment(ctx, att.FileID)
}
func (s *gmailSource) ListIDs(ctx context.Context, limit int, cursor string) ([]string, string, error) {
q := s.query
if q == "" {
q = "in:inbox"
}
ids, next, err := s.c.ListIDs(ctx, q, limit, cursor)
return ids, next, err
}
func (s *gmailSource) Get(ctx context.Context, id string) (*Message, error) {
return s.c.GetMessage(ctx, id)
}
func (s *gmailSource) DownloadAttachment(ctx context.Context, msg *Message, att Attachment) ([]byte, error) {
return s.c.DownloadAttachment(ctx, msg.ID, att.FileID)
}
// Run executes the sync across the configured sources with a worker pool.
func Run(ctx context.Context, cfg SyncConfig) (*SyncStats, error) {
if cfg.Out == "" {
cfg.Out = "var/mail"
}
if cfg.Workers <= 0 {
cfg.Workers = 4
}
if err := os.MkdirAll(cfg.Out, 0o755); err != nil {
return nil, err
}
var sources []Source
if cfg.OO != nil {
oo, err := NewOOClient(*cfg.OO, 1) // folder inbox
if err != nil {
return nil, fmt.Errorf("onlyoffice auth: %w", err)
}
sources = append(sources, &ooSource{c: oo})
}
if cfg.Gmail != nil {
gm, err := NewGmailClient(*cfg.Gmail)
if err != nil {
return nil, fmt.Errorf("gmail init: %w", err)
}
sources = append(sources, &gmailSource{c: gm, query: cfg.Query})
}
if len(sources) == 0 {
return nil, errors.New("sync: no source configured (need OO, Gmail, or both)")
}
stats := &SyncStats{}
var jobs []struct {
src Source
id string
}
for _, src := range sources {
ids, _, err := src.ListIDs(ctx, cfg.Offset+cfg.Limit, "")
if err != nil {
return nil, fmt.Errorf("list %s: %w", src.Folder(), err)
}
if cfg.Offset > 0 {
if cfg.Offset >= len(ids) {
ids = nil
} else {
ids = ids[cfg.Offset:]
}
}
if cfg.Limit > 0 && len(ids) > cfg.Limit {
ids = ids[:cfg.Limit]
}
stats.Checked += len(ids)
for _, id := range ids {
jobs = append(jobs, struct {
src Source
id string
}{src: src, id: id})
}
}
var (
wg sync.WaitGroup
mu sync.Mutex
failures []string
)
jobsCh := make(chan struct {
src Source
id string
})
for i := 0; i < cfg.Workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for j := range jobsCh {
status, err := processOne(ctx, j.src, j.id, cfg)
switch status {
case statusFailed:
mu.Lock()
failures = append(failures, j.src.Folder()+"/"+j.id+": "+err.Error())
mu.Unlock()
atomic.AddInt32(&stats.Failed, 1)
case statusNew:
atomic.AddInt32(&stats.New, 1)
case statusSkipped:
atomic.AddInt32(&stats.Skipped, 1)
}
}
}()
}
for _, j := range jobs {
select {
case jobsCh <- j:
case <-ctx.Done():
close(jobsCh)
wg.Wait()
return stats, ctx.Err()
}
}
close(jobsCh)
wg.Wait()
if len(failures) > 0 {
fmt.Fprintf(os.Stderr, "sync: %d failures:\n %s\n", len(failures), strings.Join(failures, "\n "))
}
return stats, nil
}
type status int
const (
statusNew status = iota
statusSkipped
statusFailed
)
func processOne(ctx context.Context, src Source, id string, cfg SyncConfig) (status, error) {
if cfg.DryRun {
return statusNew, nil
}
dir := filepath.Join(cfg.Out, src.Folder(), id)
jsonPath := filepath.Join(dir, "message.json")
if !cfg.Force {
if _, err := os.Stat(jsonPath); err == nil {
return statusSkipped, nil
}
}
var msg *Message
err := Retry(ctx, cfg.Policy, func() error {
m, err := src.Get(ctx, id)
if err != nil {
return retryWrap(err)
}
m.Folder = src.Folder() // directory layout is authoritative
if err := writeMessage(jsonPath, m); err != nil {
return err
}
msg = m
return nil
})
if err != nil {
return statusFailed, err
}
for _, att := range msg.Attachments {
attDir := filepath.Join(dir, "attachments")
if err := os.MkdirAll(attDir, 0o755); err != nil {
return statusFailed, err
}
attPath := filepath.Join(attDir, sanitize(att.StoredName))
if _, err := os.Stat(attPath); err == nil && !cfg.Force {
continue
}
var data []byte
err := Retry(ctx, cfg.Policy, func() error {
b, err := src.DownloadAttachment(ctx, msg, att)
if err != nil {
return retryWrap(err)
}
data = b
return os.WriteFile(attPath, b, 0o644)
})
if err != nil {
return statusFailed, fmt.Errorf("attachment %s: %w", att.FileName, err)
}
// ICS attachments get structured markdown immediately (same name the
// Python converter would use: <display stem>.md).
if isICS(att.FileName) {
stem := att.FileName
if i := strings.LastIndex(stem, "."); i >= 0 {
stem = stem[:i]
}
mdPath := filepath.Join(attDir, sanitize(stem)+".md")
if err := os.WriteFile(mdPath, []byte(ICSToMarkdown(data)), 0o644); err != nil {
return statusFailed, err
}
}
}
return statusNew, nil
}
func writeMessage(path string, m *Message) error {
b, err := json.MarshalIndent(m, "", " ")
if err != nil {
return err
}
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
return err
}
return os.WriteFile(path, b, 0o644)
}
func sanitize(name string) string {
r := strings.NewReplacer("/", "_", "\\", "_", ":", "_", "*", "_", "?", "_", "\"", "_",
"<", "_", ">", "_", "|", "_", " ", "_")
return r.Replace(name)
}
func isICS(name string) bool {
n := strings.ToLower(name)
return strings.HasSuffix(n, ".ics") || strings.HasSuffix(n, ".ical")
}
+305
View File
@@ -0,0 +1,305 @@
package sync
import (
"context"
"encoding/base64"
"errors"
"testing"
"time"
"unicode/utf8"
)
func TestRetrySucceedsOnSecondTry(t *testing.T) {
attempts := 0
err := Retry(context.Background(), RetryPolicy{BaseDelay: time.Millisecond, MaxDelay: 5 * time.Millisecond}, func() error {
attempts++
if attempts == 1 {
return retryWrap(errors.New("boom"))
}
return nil
})
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if attempts != 2 {
t.Fatalf("expected 2 attempts, got %d", attempts)
}
}
func TestRetryExhaustsAttempts(t *testing.T) {
attempts := 0
err := Retry(context.Background(), RetryPolicy{MaxAttempts: 3, BaseDelay: time.Millisecond, MaxDelay: 5 * time.Millisecond}, func() error {
attempts++
return retryWrap(errors.New("nope"))
})
if err == nil {
t.Fatal("expected error after exhaustion")
}
if attempts != 3 {
t.Fatalf("expected 3 attempts, got %d", attempts)
}
}
func TestRetryNonRetriableAbortsImmediately(t *testing.T) {
attempts := 0
err := Retry(context.Background(), RetryPolicy{MaxAttempts: 5, BaseDelay: time.Millisecond}, func() error {
attempts++
return errors.New("permanent")
})
if err == nil {
t.Fatal("expected error")
}
if attempts != 1 {
t.Fatalf("expected 1 attempt for non-retriable, got %d", attempts)
}
}
func TestRetryRespectsContext(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
err := Retry(ctx, RetryPolicy{MaxAttempts: 5, BaseDelay: time.Millisecond}, func() error {
return retryWrap(errors.New("x"))
})
if !errors.Is(err, context.Canceled) {
t.Fatalf("expected context.Canceled, got %v", err)
}
}
func TestDelayGrows(t *testing.T) {
p := RetryPolicy{BaseDelay: time.Second, MaxDelay: 30 * time.Second, Jitter: 0}
d1 := p.delay(1) // attempt 1 => 0
d2 := p.delay(2)
d3 := p.delay(3)
if d1 != 0 {
t.Fatalf("attempt 1 delay should be 0, got %v", d1)
}
if d2 != time.Second {
t.Fatalf("attempt 2 delay should be 1s, got %v", d2)
}
if d3 != 2*time.Second {
t.Fatalf("attempt 3 delay should be 2s, got %v", d3)
}
}
func TestSanitize(t *testing.T) {
cases := map[string]string{
"a/b\\c:d*e": "a_b_c_d_e",
"normal.txt": "normal.txt",
"../evil": ".._evil",
"a b c.pdf": "a_b_c.pdf",
}
for in, want := range cases {
if got := sanitize(in); got != want {
t.Errorf("sanitize(%q) = %q, want %q", in, got, want)
}
}
}
func TestIsICS(t *testing.T) {
if !isICS("reply.ics") || !isICS("x.ICAL") {
t.Fatal("ics extensions not detected")
}
if isICS("invoice.pdf") {
t.Fatal("pdf misdetected as ics")
}
}
const fixtureReplyICS = `BEGIN:VCALENDAR
METHOD:REPLY
PRODID:Microsoft Exchange Server 2010
VERSION:2.0
BEGIN:VTIMEZONE
TZID:W. Europe Standard Time
BEGIN:STANDARD
DTSTART:16010101T030000
TZOFFSETFROM:+0200
TZOFFSETTO:+0100
RRULE:FREQ=YEARLY;INTERVAL=1;BYDAY=-1SU;BYMONTH=10
END:STANDARD
BEGIN:DAYLIGHT
DTSTART:16010101T020000
TZOFFSETFROM:+0100
TZOFFSETTO:+0200
RRULE:FREQ=YEARLY;INTERVAL=1;BYDAY=-1SU;BYMONTH=3
END:DAYLIGHT
END:VTIMEZONE
BEGIN:VEVENT
ATTENDEE;PARTSTAT=ACCEPTED;CN="Baker, Ben":mailto:bbaker1@teksystems.com
UID:bvlnr1i35ug30kn6rvu9dop00g@google.com
SUMMARY;LANGUAGE=en-US:Accepted: Appointment (Ben Baker)
DTSTART;TZID=W. Europe Standard Time:20260812T120000
DTEND;TZID=W. Europe Standard Time:20260812T123000
CLASS:PUBLIC
STATUS:CONFIRMED
LOCATION;LANGUAGE=en-US:https://meet.google.com/sxh-ubud-jrd
END:VEVENT
END:VCALENDAR`
func TestICSToMarkdown(t *testing.T) {
out := ICSToMarkdown([]byte(fixtureReplyICS))
for _, want := range []string{
"Accepted: Appointment",
"When:",
"Where:",
"meet.google.com",
"Attendee:",
"Baker, Ben",
"Accepted",
"Calendar method: REPLY",
} {
if !contains(out, want) {
t.Errorf("output missing %q:\n%s", want, out)
}
}
if contains(out, "BEGIN:VCALENDAR") {
t.Errorf("raw ICS leaked into markdown:\n%s", out)
}
}
func TestICSToMarkdownFallback(t *testing.T) {
out := ICSToMarkdown([]byte("not a calendar"))
if !contains(out, "not a calendar") {
t.Fatalf("expected raw fallback, got %q", out)
}
}
func TestICSToMarkdownNormalizesLatin1(t *testing.T) {
// Real-world ICS from a rental portal: summary in UTF-8, location Latin-1
// ("N\xfcrnberg"). The markdown output must be valid UTF-8 everywhere.
raw := "BEGIN:VCALENDAR\r\nVERSION:2.0\r\nBEGIN:VEVENT\r\n" +
"SUMMARY:Mietwagen-Buchung: N\xc3\xbcrnberg\r\n" +
"LOCATION:N\xfcrnberg\r\nDTSTART:20200101T090000Z\r\nDTEND:20200101T180000Z\r\n" +
"END:VEVENT\r\nEND:VCALENDAR\r\n"
out := ICSToMarkdown([]byte(raw))
if !utf8.ValidString(out) {
t.Fatalf("output is not valid UTF-8:\n%q", out)
}
if !contains(out, "Nürnberg") {
t.Errorf("expected Nürnberg in output:\n%s", out)
}
if contains(out, "N\xfcrnberg") {
t.Errorf("Latin-1 bytes leaked into output:\n%q", out)
}
}
func TestICSToMarkdownAllDay(t *testing.T) {
ics := `BEGIN:VCALENDAR
VERSION:2.0
BEGIN:VEVENT
UID:y@google.com
SUMMARY:All day thing
DTSTART;VALUE=DATE:20260815
DTEND;VALUE=DATE:20260816
END:VEVENT
END:VCALENDAR`
out := ICSToMarkdown([]byte(ics))
if !contains(out, "All day thing") || !contains(out, "2026-08-15") {
t.Errorf("all-day event not parsed:\n%s", out)
}
}
func TestCollectParts(t *testing.T) {
p := gmailPart{
MimeType: "multipart/mixed",
Parts: []gmailPart{
{MimeType: "multipart/alternative", Parts: []gmailPart{
{MimeType: "text/plain", Body: gmailBody{Data: b64("plain text")}},
{MimeType: "text/html", Body: gmailBody{Data: b64("<p>html</p>")}},
}},
{PartID: "2", MimeType: "application/pdf", Filename: "invoice.pdf", Body: gmailBody{Size: 100}},
},
}
text, html, atts := collectParts(p, "root", "abc123", 0)
if text != "plain text" {
t.Errorf("text = %q", text)
}
if html != "<p>html</p>" {
t.Errorf("html = %q", html)
}
if len(atts) != 1 || atts[0].FileName != "invoice.pdf" || atts[0].FileID != "abc123:2" {
t.Errorf("atts = %+v", atts)
}
}
type fakeGmailAPI struct {
lastQ string
lastLimit int
ids []string
}
func (f *fakeGmailAPI) ListIDs(_ context.Context, q string, maxIDs int, _ string) ([]string, string, error) {
f.lastQ = q
f.lastLimit = maxIDs
return f.ids, "", nil
}
func (f *fakeGmailAPI) GetMessage(context.Context, string) (*Message, error) {
return nil, errors.New("unused")
}
func (f *fakeGmailAPI) DownloadAttachment(context.Context, string, string) ([]byte, error) {
return nil, errors.New("unused")
}
func TestGmailSourcePassesQueryToListIDs(t *testing.T) {
fake := &fakeGmailAPI{ids: []string{"m1"}}
src := &gmailSource{c: fake, query: "from:alice@example.com"}
ids, _, err := src.ListIDs(context.Background(), 10, "")
if err != nil {
t.Fatal(err)
}
if fake.lastQ != "from:alice@example.com" {
t.Fatalf("ListIDs q=%q, want from:alice@example.com", fake.lastQ)
}
if fake.lastLimit != 10 {
t.Fatalf("ListIDs limit=%d, want 10", fake.lastLimit)
}
if len(ids) != 1 || ids[0] != "m1" {
t.Fatalf("ids=%v", ids)
}
}
func TestGmailSourceEmptyQueryDefaultsToInbox(t *testing.T) {
fake := &fakeGmailAPI{}
src := &gmailSource{c: fake, query: ""}
if _, _, err := src.ListIDs(context.Background(), 5, ""); err != nil {
t.Fatal(err)
}
if fake.lastQ != "in:inbox" {
t.Fatalf("empty query q=%q, want in:inbox", fake.lastQ)
}
}
func TestParseCLIGmailQuery(t *testing.T) {
cli, code, err := ParseCLI([]string{
"--source", "gmail",
"--query", "from:letrado@example.com",
"--out", t.TempDir(),
"--dry-run",
})
if err != nil || code != 0 {
t.Fatalf("ParseCLI: code=%d err=%v", code, err)
}
if cli.Sync.Query != "from:letrado@example.com" {
t.Fatalf("query=%q", cli.Sync.Query)
}
if cli.Sync.Gmail == nil {
t.Fatal("gmail source not configured")
}
}
func b64(s string) string {
return base64.URLEncoding.EncodeToString([]byte(s))
}
func contains(s, sub string) bool {
return len(s) >= len(sub) && (s == sub || len(sub) == 0 ||
indexOf(s, sub) >= 0)
}
func indexOf(s, sub string) int {
for i := 0; i+len(sub) <= len(s); i++ {
if s[i:i+len(sub)] == sub {
return i
}
}
return -1
}
+46
View File
@@ -0,0 +1,46 @@
// Package sync downloads OnlyOffice and Gmail messages to var/mail/ as raw
// JSON + attachment files, then hands off to bin/mail/import --from-raw for
// markdown conversion.
//
// On-disk schema (per message):
//
// var/mail/<folder>/<id>/message.json # Message (this package)
// var/mail/<folder>/<id>/attachments/ # raw attachment bytes (storedName)
//
// The Message JSON is the contract shared with the Python converter. Fields
// deliberately mirror what bin/mail/import already reads from the OnlyOffice
// API, so conversion is source-agnostic.
package sync
import (
"time"
)
// Attachment describes one attachment of a Message. FileID/FileName/StoredName
// mirror OnlyOffice; Gmail fills them from its own ids. StoredName is always
// unique (hash/attachment id) so raw files never collide.
type Attachment struct {
FileID string `json:"fileId,omitempty"`
FileName string `json:"fileName"`
StoredName string `json:"storedName"`
Size int64 `json:"size,omitempty"`
ContentType string `json:"contentType,omitempty"`
}
// Message is the normalized record written to var/mail/<folder>/<id>/message.json.
type Message struct {
Source string `json:"source"` // "onlyoffice" | "gmail"
ID string `json:"id"`
Folder string `json:"folder"`
Subject string `json:"subject,omitempty"`
From string `json:"from,omitempty"`
To string `json:"to,omitempty"`
CC string `json:"cc,omitempty"`
BCC string `json:"bcc,omitempty"`
ReceivedAt time.Time `json:"receivedAt,omitempty"`
HTMLBody string `json:"htmlBody,omitempty"`
TextBody string `json:"textBody,omitempty"`
HasAttachments bool `json:"hasAttachments,omitempty"`
Attachments []Attachment `json:"attachments,omitempty"`
MimeMessageID string `json:"mimeMessageId,omitempty"`
}
+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
+100
View File
@@ -0,0 +1,100 @@
//usr/bin/env go run "$0" "$@"; exit
//
// bin/markdown/import.go - split markdown into leafs (H2 boundaries).
//
// ./bin/markdown/import.go [dir]
// ./bin/markdown/import.go --files a.md,b.md --json
//
// Conversion only. Brain write is bin/brain/index.go.
// Python bin/md/import remains as a fallback.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"fmt"
"os"
"strings"
"github.com/eSlider/2dph/internal/mdleaves"
)
func main() {
os.Exit(run(os.Args[1:]))
}
func run(args []string) int {
jsonOut := false
files := ""
root := "."
for i := 0; i < len(args); i++ {
a := args[i]
switch {
case a == "--json":
jsonOut = true
case a == "--files" && i+1 < len(args):
i++
files = args[i]
case strings.HasPrefix(a, "--files="):
files = strings.TrimPrefix(a, "--files=")
case a == "-h" || a == "--help":
fmt.Fprintln(os.Stderr, "bin/markdown/import.go [dir] [--files a.md,b.md] [--json]")
return 0
case strings.HasPrefix(a, "-"):
fmt.Fprintln(os.Stderr, "unknown arg:", a)
return 2
default:
root = a
}
}
var paths []string
if files != "" {
for _, f := range strings.Split(files, ",") {
f = strings.TrimSpace(f)
if f != "" {
paths = append(paths, f)
}
}
} else {
st, err := os.Stat(root)
if err != nil {
fmt.Fprintf(os.Stderr, "md/import: no such path %s\n", root)
return 2
}
if !st.IsDir() {
paths = []string{root}
} else {
var err error
paths, err = mdleaves.WalkMarkdown(root)
if err != nil {
fmt.Fprintf(os.Stderr, "md/import: %v\n", err)
return 1
}
}
}
if len(paths) == 0 {
fmt.Fprintln(os.Stderr, "md/import: no markdown files")
return 1
}
var all []mdleaves.Leaf
for _, p := range paths {
raw, err := os.ReadFile(p)
if err != nil {
fmt.Fprintf(os.Stderr, "md/import: %s: %v\n", p, err)
continue
}
all = append(all, mdleaves.ToAll(string(raw), p, "")...)
}
if jsonOut {
s, err := mdleaves.EncodeJSON(all)
if err != nil {
fmt.Fprintln(os.Stderr, err)
return 1
}
fmt.Print(s)
return 0
}
fmt.Print(mdleaves.EncodeYAML(all))
return 0
}
+3 -2
View File
@@ -1,9 +1,10 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
import lib from __future__ import annotations
import sys import sys
from pathlib import Path from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent)) sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "tools"))
from mdleaves import leaves_to_json, read_markdown, to_all, walk_markdown # noqa: E402 from mdleaves import leaves_to_json, read_markdown, to_all, walk_markdown # noqa: E402
from yamlout import to_yaml # noqa: E402 from yamlout import to_yaml # noqa: E402
+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:]))
}
+73
View File
@@ -0,0 +1,73 @@
//usr/bin/env go run -tags=qa_stats "$0" "$@"; exit
//go:build qa_stats
//
// bin/qa/stats.go - DuckDB quantiles over a JSON number array or JSONL count.
//
// ./bin/qa/stats.go <<< '[1,2,3,4,5]'
// ./bin/qa/stats.go --jsonl rows.jsonl
//
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
// DuckDB CGO needs gcc/g++ (not Zig). After eval "$(bin/cgo/zig env)":
// CC=gcc CXX=g++ CGO_CFLAGS= CGO_LDFLAGS= ./bin/qa/stats.go
package main
import (
"encoding/json"
"fmt"
"io"
"os"
"strings"
"github.com/eSlider/2dph/internal/duckstats"
)
func main() {
os.Exit(run(os.Args[1:]))
}
func run(args []string) int {
jsonl := ""
for i := 0; i < len(args); i++ {
a := args[i]
switch {
case a == "--jsonl" && i+1 < len(args):
i++
jsonl = args[i]
case strings.HasPrefix(a, "--jsonl="):
jsonl = strings.TrimPrefix(a, "--jsonl=")
case a == "-h" || a == "--help":
fmt.Fprintln(os.Stderr, "bin/qa/stats.go [--jsonl FILE] # stdin = JSON [float,…]")
return 0
default:
fmt.Fprintln(os.Stderr, "unknown arg:", a)
return 2
}
}
if jsonl != "" {
n, err := duckstats.CountJSONL(jsonl)
if err != nil {
fmt.Fprintln(os.Stderr, err)
return 1
}
fmt.Printf("n: %d\n", n)
return 0
}
raw, err := io.ReadAll(os.Stdin)
if err != nil {
fmt.Fprintln(os.Stderr, err)
return 1
}
var samples []float64
if err := json.Unmarshal(raw, &samples); err != nil {
fmt.Fprintln(os.Stderr, err)
return 1
}
s, err := duckstats.Quantiles(samples)
if err != nil {
fmt.Fprintln(os.Stderr, err)
return 1
}
fmt.Printf("n: %d\nmin: %g\np50: %g\np95: %g\nmax: %g\navg: %g\n",
s.N, s.Min, s.P50, s.P95, s.Max, s.Avg)
return 0
}
+102
View File
@@ -0,0 +1,102 @@
//usr/bin/env go run -tags=reasoner_bakeoff "$0" "$@"; exit
//go:build reasoner_bakeoff
//
// bin/reasoner/bakeoff.go - CPU tool-call bake-off against an OpenAI-compatible URL (D18).
//
// REASONER_BASE_URL=http://127.0.0.1:11435/v1 REASONER_MODEL=qwen3.5:9b ./bin/reasoner/bakeoff.go
// ./bin/reasoner/bakeoff.go --model MichelRosselli/bonsai-27b:Q1_0 --json
//
// Measures OpenAI tool_calls (search/get/audit) and RSS from Ollama /api/ps, not VRAM.
// PicoClaw is compose profile picoclaw; tool names match internal/httpapi MCP ops.
// NOTE: never run `gofmt -w` on this file — it breaks the shebang.
package main
import (
"encoding/json"
"fmt"
"os"
"strings"
"github.com/eSlider/2dph/internal/duckstats"
"github.com/eSlider/2dph/internal/reasoner"
)
func main() {
os.Exit(run(os.Args[1:]))
}
func run(args []string) int {
base := os.Getenv("REASONER_BASE_URL")
if base == "" {
base = "http://127.0.0.1:11435/v1"
}
model := os.Getenv("REASONER_MODEL")
if model == "" {
model = reasoner.OllamaRAM
}
jsonOut := false
device := "cpu"
for i := 0; i < len(args); i++ {
a := args[i]
switch {
case a == "--json":
jsonOut = true
case a == "--model" && i+1 < len(args):
i++
model = args[i]
case strings.HasPrefix(a, "--model="):
model = strings.TrimPrefix(a, "--model=")
case a == "--base-url" && i+1 < len(args):
i++
base = args[i]
case a == "--device" && i+1 < len(args):
i++
device = args[i]
case a == "-h" || a == "--help":
fmt.Fprintln(os.Stderr, "bin/reasoner/bakeoff.go [--model ID] [--base-url URL] [--device cpu] [--json]")
return 0
default:
fmt.Fprintln(os.Stderr, "unknown arg:", a)
return 2
}
}
c := reasoner.Client{BaseURL: base, Model: model, Device: device}
rep := reasoner.Run(c)
lat := make([]float64, 0, len(rep.Prompts))
for _, p := range rep.Prompts {
lat = append(lat, float64(p.LatencyMS))
}
if st, err := duckstats.Quantiles(lat); err == nil {
rep.LatencyP50MS = st.P50
rep.LatencyP95MS = st.P95
}
raw, err := json.MarshalIndent(rep, "", " ")
if err != nil {
fmt.Fprintln(os.Stderr, err)
return 1
}
if jsonOut {
fmt.Println(string(raw))
} else {
fmt.Printf("model: %s\n", rep.Model)
fmt.Printf("hf_id: %s\n", rep.HF)
fmt.Printf("device: %s\n", rep.Device)
fmt.Printf("tool_call: %d/%d\n", rep.ToolCallOK, rep.ToolCallN)
fmt.Printf("xml_leak: %d\n", rep.XMLLeak)
fmt.Printf("rss_mb: %d\n", rep.RSSMB)
fmt.Printf("vram_mb: %d\n", rep.VRAMMB)
fmt.Printf("latency_p50_ms: %g\n", rep.LatencyP50MS)
fmt.Printf("latency_p95_ms: %g\n", rep.LatencyP95MS)
for _, p := range rep.Prompts {
status := "fail"
if p.OK {
status = "ok"
}
fmt.Printf(" %s: %s wanted=%s got=%s xml=%v %dms %s\n", p.WantedTool, status, p.WantedTool, p.ToolName, p.XMLLeak, p.LatencyMS, p.Err)
}
}
if rep.ToolCallN == 0 {
return 1
}
return 0
}
+2
View File
@@ -0,0 +1,2 @@
// Commands in this directory are shebang mains (bakeoff.go).
package main
Executable
+22
View File
@@ -0,0 +1,22 @@
//usr/bin/env go run -tags=brain_serve "$0" "$@"; exit
//go:build brain_serve
//
// bin/serve.go — deprecated; use bin/brain/serve.go.
package main
import (
"fmt"
"os"
"github.com/eSlider/2dph/internal/httpapi"
)
func main() {
fmt.Fprintln(os.Stderr, "bin/serve.go is deprecated; use bin/brain/serve.go")
if os.Getenv("KB_ROOT") == "" {
if wd, err := os.Getwd(); err == nil {
os.Setenv("KB_ROOT", wd)
}
}
httpapi.Run(nil)
}
+103
View File
@@ -0,0 +1,103 @@
"""D16 contradiction adjudication (same rules as internal/facts)."""
from __future__ import annotations
from typing import Any
CONF_CONFIRMED = "confirmed"
CONF_HYPOTHESIS = "hypothesis"
RULE_UNRESOLVED = "unresolved"
RULE_TEMPORAL = "temporal_freshness"
RULE_AUTHORITY = "authority_pairing"
RULE_TWO_SOURCE = "two_source"
RULE_SINGLE = "single_source"
KIND_RUNTIME = "runtime"
KIND_CONFIG = "config"
KIND_NARRATIVE = "narrative"
def _independent(sources: list[dict]) -> int:
seen: set[str] = set()
for i, s in enumerate(sources):
sid = str(s.get("id") or "") or f"{s.get('kind', '')}#{i}"
seen.add(sid)
return len(seen)
def _fresh_n(sources: list[dict]) -> int:
return sum(1 for s in sources if not s.get("stale"))
def _strong_n(sources: list[dict]) -> int:
return sum(1 for s in sources if s.get("kind") in (KIND_RUNTIME, KIND_CONFIG))
def adjudicate(claim: dict[str, Any]) -> dict[str, Any]:
yes = list(claim.get("yes") or [])
no = list(claim.get("no") or [])
yes_n, no_n = _independent(yes), _independent(no)
text = str(claim.get("text") or "")
def out(conf: str, rule: str, winner: str = "") -> dict[str, Any]:
return {
"text": text,
"confidence": conf,
"confirmed": conf == CONF_CONFIRMED,
"rule": rule,
"winner": winner,
"yes": yes_n,
"no": no_n,
}
if yes_n < 2 or no_n < 2:
if yes_n >= 2:
return out(CONF_CONFIRMED, RULE_TWO_SOURCE, "yes")
if no_n >= 2:
return out(CONF_CONFIRMED, RULE_TWO_SOURCE, "no")
return out(CONF_HYPOTHESIS, RULE_SINGLE)
yf, nf = _fresh_n(yes), _fresh_n(no)
if yf >= 2 and nf < 2:
return out(CONF_CONFIRMED, RULE_TEMPORAL, "yes")
if nf >= 2 and yf < 2:
return out(CONF_CONFIRMED, RULE_TEMPORAL, "no")
ys, ns = _strong_n(yes), _strong_n(no)
if ys >= 2 and ns < 2:
return out(CONF_CONFIRMED, RULE_AUTHORITY, "yes")
if ns >= 2 and ys < 2:
return out(CONF_CONFIRMED, RULE_AUTHORITY, "no")
return out(CONF_HYPOTHESIS, RULE_UNRESOLVED)
def parse_source_field(source: str) -> tuple[str, str]:
"""Split `a x b vs c x d` into (yes, no). Empty no if no ` vs `."""
if " vs " not in source:
return source, ""
yes, _, no = source.partition(" vs ")
return yes.strip(), no.strip()
def check_fact_row(lid: str, source: str, loc: str, how: str, conf: str) -> list[str]:
"""Lexicon checks for one facts leaf (no Ladybug)."""
problems: list[str] = []
src = source or ""
if conf == CONF_CONFIRMED:
if " vs " in src:
problems.append(f"{lid}: confirmed fact cannot keep a vs-contradiction")
if " x " not in src:
problems.append(f"{lid}: needs 2-source evidence in source, got '{source}'")
elif conf == CONF_HYPOTHESIS:
yes, no = parse_source_field(src)
if not no or " x " not in yes or " x " not in no:
problems.append(
f"{lid}: hypothesis contradiction needs 'a x b vs c x d', got '{source}'"
)
elif conf == "partial":
pass
else:
problems.append(f"{lid}: unknown confidence '{conf}'")
if not loc:
problems.append(f"{lid}: missing loc (evidence pointer)")
if not how:
problems.append(f"{lid}: missing how")
return problems
+29
View File
@@ -0,0 +1,29 @@
"""crmfacts - pure helpers for bin/facts/crm (association proofing).
Shared with tools/ unit tests so the corpus-org parser is covered in CI.
"""
import re
def corpus_orgs(raw: str) -> dict[str, dict]:
"""Parse the orgs block of the CV knowledge-mesh YAML into id -> fields.
Fields kept: label, kind, period, website. Stops at the first sibling
top-level key (clients, timeline, ...).
"""
m = re.search(r"^orgs:\n(.*?)\n^(?:clients|timeline|tech_weights|nodes|edges):", raw, re.S | re.M)
if not m:
return {}
orgs: dict[str, dict] = {}
cur = None
for line in m.group(1).splitlines():
lm = re.match(r"^\s*- id:\s*(\S+)", line)
if lm:
cur = lm.group(1)
orgs[cur] = {}
continue
fm = re.match(r"^\s+(\w+):\s*(.*)$", line)
if fm and cur and fm.group(1) in ("label", "kind", "period", "website"):
orgs[cur][fm.group(1)] = fm.group(2).strip()
return orgs
+60
View File
@@ -0,0 +1,60 @@
"""gitimport - Ladybug graph writes for Commit/File/Person (no git binary).
Commit records come from bin/git/import.go (go-git). This module only MERGEs
the version graph File-[:HAS_VERSION]->Commit-[:AUTHORED]->Person.
"""
from __future__ import annotations
from dataclasses import dataclass, field
@dataclass
class Commit:
sha: str
author: str
email: str
date: str
subject: str
files: list[str] = field(default_factory=list)
GIT_SCHEMA = (
"CREATE NODE TABLE IF NOT EXISTS Commit (id STRING, repo STRING, subject STRING, "
"author STRING, email STRING, date STRING, PRIMARY KEY(id))",
"CREATE NODE TABLE IF NOT EXISTS Person (id STRING, name STRING, email STRING, PRIMARY KEY(id))",
"CREATE REL TABLE IF NOT EXISTS HAS_VERSION (FROM File TO Commit)",
"CREATE REL TABLE IF NOT EXISTS AUTHORED (FROM Commit TO Person)",
)
def ensure_git_schema(conn) -> None:
for stmt in GIT_SCHEMA:
conn.execute(stmt)
def index_commits(conn, commits: list[Commit], repo: str) -> int:
"""Write Commit/File/Person nodes + edges, one per commit (idempotent by sha)."""
ensure_git_schema(conn)
for c in commits:
conn.execute(
"MERGE (c:Commit {id:$sha}) SET c.repo=$repo, c.subject=$subject, "
"c.author=$author, c.email=$email, c.date=$date",
parameters={"sha": c.sha, "repo": repo, "subject": c.subject,
"author": c.author, "email": c.email, "date": c.date},
)
conn.execute(
"MERGE (p:Person {id:$email}) SET p.name=$name, p.email=$email",
parameters={"email": c.email, "name": c.author},
)
conn.execute("MATCH (c:Commit {id:$sha}), (p:Person {id:$email}) "
"MERGE (c)-[:AUTHORED]->(p)",
parameters={"sha": c.sha, "email": c.email})
for path in c.files:
conn.execute(
"MERGE (f:File {id:$fid}) SET f.path=$path, f.repo=$repo",
parameters={"fid": f"{repo}:{path}", "path": path, "repo": repo},
)
conn.execute("MATCH (f:File {id:$fid}), (c:Commit {id:$sha}) "
"MERGE (f)-[:HAS_VERSION]->(c)",
parameters={"fid": f"{repo}:{path}", "sha": c.sha})
return len(commits)
+333
View File
@@ -0,0 +1,333 @@
"""kblib - the 2dph brain core over LadybugDB.
Single embedded graph `var/kb.lbug`. Two roots: facts (assertions backed by
>=2 independent sources) and info (narrative leafs). Hybrid retrieval: BM25
(FTS extension) + HNSW cosine (VECTOR extension) + Cypher graph hops.
All access is read-only unless `--rebuild` (kb/index) or `kb/add`.
"""
from __future__ import annotations
import hashlib
import json
import time
import zlib
from pathlib import Path
import ladybug
MODEL = "minishlab/potion-multilingual-128M"
EMBED_DIM = 256
ROOT_FACTS = "facts"
ROOT_INFO = "info"
CONF_CONFIRMED = "confirmed"
def _repo_root() -> Path:
p = Path(__file__).resolve().parent
while True:
if (p / "var").is_dir() or (p / ".git").is_dir() or (p / "pyproject.toml").is_file():
return p
if p.parent == p:
return Path(__file__).resolve().parents[2]
p = p.parent
VAR = _repo_root() / "var"
DB_PATH = VAR / "kb.lbug"
def sha256_b64(text: str) -> str:
return hashlib.sha256(text.encode()).hexdigest()
def _load_extension(conn: ladybug.Connection, name: str) -> None:
"""Install (download once) and load a ladybug extension."""
try:
conn.execute(f"INSTALL {name}")
except Exception:
pass # already installed / offline-ok when present
conn.execute(f"LOAD EXTENSION {name}")
def connect(path: Path | str | None = None, read_only: bool = True) -> tuple[ladybug.Database, ladybug.Connection]:
db = ladybug.Database(str(path or DB_PATH), read_only=read_only)
conn = ladybug.Connection(db)
_load_extension(conn, "FTS")
_load_extension(conn, "VECTOR")
return db, conn
def init_schema(conn: ladybug.Connection) -> None:
conn.execute(
"CREATE NODE TABLE IF NOT EXISTS Leaf ("
" id STRING, text STRING, root STRING, confidence STRING, "
" sha256 STRING, source STRING, source_rev STRING, observed_at STRING, "
" how STRING, loc STRING, type STRING, embedding FLOAT[256], "
" PRIMARY KEY(id))"
)
conn.execute(
"CREATE NODE TABLE IF NOT EXISTS File ("
" id STRING, path STRING, repo STRING, mtime STRING, PRIMARY KEY(id))"
)
conn.execute(
"CREATE REL TABLE IF NOT EXISTS FROM_FILE (FROM Leaf TO File)"
)
conn.execute(
"CREATE NODE TABLE IF NOT EXISTS Host (id STRING, hostname STRING, user STRING, PRIMARY KEY(id))"
)
conn.execute(
"CREATE REL TABLE IF NOT EXISTS RUNS_ON (FROM Leaf TO Host)"
)
conn.execute(
"CREATE NODE TABLE IF NOT EXISTS Commit (id STRING, repo STRING, subject STRING, "
"author STRING, email STRING, date STRING, PRIMARY KEY(id))"
)
conn.execute(
"CREATE NODE TABLE IF NOT EXISTS Person (id STRING, name STRING, email STRING, PRIMARY KEY(id))"
)
conn.execute(
"CREATE REL TABLE IF NOT EXISTS HAS_VERSION (FROM File TO Commit)"
)
conn.execute(
"CREATE REL TABLE IF NOT EXISTS AUTHORED (FROM Commit TO Person)"
)
def leaf_id(text: str, source: str) -> str:
return sha256_b64(f"{source}\0{text}")[:24]
def upsert_leaf(conn: ladybug.Connection, *, text: str, root: str, confidence: str,
source: str, source_rev: str, how: str, loc: str, type_: str,
embedding: list[float] | None) -> str:
lid = leaf_id(text, source)
obs = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())
conn.execute(
"MERGE (l:Leaf {id:$id}) "
"SET l.text=$text, l.root=$root, l.confidence=$confidence, "
" l.sha256=$sha, l.source=$source, l.source_rev=$rev, l.observed_at=$obs, "
" l.how=$how, l.loc=$location, l.type=$type"
+ (", l.embedding=$emb" if embedding else ""),
parameters={
"id": lid, "text": text, "root": root, "confidence": confidence,
"sha": sha256_b64(text), "source": source, "rev": source_rev,
"obs": obs, "how": how, "location": loc, "type": type_,
"emb": (embedding if embedding else None),
},
)
return lid
def add_leafs(conn: ladybug.Connection, leafs: list[dict]) -> list[str]:
"""Write facts+info leafs in one transaction. Safe while FTS/HNSW exist.
Each leaf dict: text, source, optional root/confidence/source_rev/how/loc/type/embedding.
Does not delete the database file. Measured on Ladybug 0.19: MERGE of new
ids (and updates) stays FTS+HNSW queryable; DROP INDEX is the fatal path.
"""
if not leafs:
return []
started = False
try:
conn.execute("BEGIN TRANSACTION")
started = True
except Exception:
started = False
ids: list[str] = []
try:
for lf in leafs:
ids.append(
upsert_leaf(
conn,
text=str(lf["text"]),
root=str(lf.get("root") or ROOT_INFO),
confidence=str(lf.get("confidence") or CONF_CONFIRMED),
source=str(lf["source"]),
source_rev=str(lf.get("source_rev") or "working-tree"),
how=str(lf.get("how") or "brain/add"),
loc=str(lf.get("loc") or lf.get("source") or ""),
type_=str(lf.get("type") or lf.get("type_") or "reference"),
embedding=lf.get("embedding"),
)
)
if started:
conn.execute("COMMIT")
except Exception:
if started:
try:
conn.execute("ROLLBACK")
except Exception:
pass
raise
return ids
def file_id(repo: str, path: str) -> str:
"""Stable File.id matching gitimport (`repo:path`)."""
return f"{repo}:{path}" if repo else path
def link_from_file(conn: ladybug.Connection, leaf_id: str, path: str,
repo: str = "", mtime: str = "") -> str:
"""MERGE File and Leaf-[:FROM_FILE]->File so --hop 1 can walk."""
fid = file_id(repo, path)
conn.execute(
"MERGE (f:File {id:$id}) SET f.path=$path, f.repo=$repo, f.mtime=$mtime",
parameters={"id": fid, "path": path, "repo": repo, "mtime": mtime},
)
conn.execute(
"MATCH (l:Leaf {id:$lid}), (f:File {id:$fid}) "
"MERGE (l)-[:FROM_FILE]->(f)",
parameters={"lid": leaf_id, "fid": fid},
)
return fid
HOP_STMTS = {
1: "MATCH (l:Leaf {id:$id})-[:FROM_FILE]->(f:File) RETURN f.id, f.path, 1",
2: ("MATCH (l:Leaf {id:$id})-[:FROM_FILE]->(f:File)-[:HAS_VERSION]->(c:Commit) "
"RETURN c.id, c.subject, 2"),
3: ("MATCH (l:Leaf {id:$id})-[:FROM_FILE]->(f:File)-[:HAS_VERSION]->(c:Commit)"
"-[:AUTHORED]->(p:Person) RETURN p.id, p.name, 3"),
}
HOP_LABELS = {1: "File", 2: "Commit", 3: "Person"}
def hop_walk(conn: ladybug.Connection, leaf_id: str, n: int) -> list[dict]:
"""Walk Leaf → File → Commit → Person up to n hops (max 3)."""
depth = min(max(int(n), 0), 3)
out: list[dict] = []
for d in range(1, depth + 1):
rows = conn.execute(HOP_STMTS[d], parameters={"id": leaf_id}).get_all()
for row in rows:
out.append({
"id": row[0],
"label": HOP_LABELS[d],
"name": row[1],
"depth": int(row[2]),
})
return out
def leaf_index_names(conn: ladybug.Connection) -> set[str]:
"""Return index names on the Leaf table (e.g. {'id', 'Leaf_vec', '_PK'})."""
rows = conn.execute("CALL SHOW_INDEXES() RETURN *").get_all()
return {row[1] for row in rows if row[0] == "Leaf"}
def create_fts_and_vector(conn: ladybug.Connection, force: bool = False) -> None:
"""Create FTS (BM25) + HNSW vector indexes if missing.
Never DROP INDEX for FTS/VECTOR. Ladybug 0.19 leaves ghost catalog
entries after DROP (`_0_Leaf_vec_UPPER`, `0_id_docs`), so a later
CREATE fails with "already exists in catalog" while SHOW_INDEXES
still omits the index. Swallowing that error made HNSW look "OK"
until the first QUERY_VECTOR_INDEX.
`force=True` is accepted for API compatibility but does **not** drop.
Fresh indexes require deleting `var/kb.lbug` and rebuilding
(`bin/brain/index.go --rebuild`).
"""
del force # API compat; DROP is unsafe — see docstring
names = leaf_index_names(conn)
if "id" not in names:
try:
conn.execute("CALL CREATE_FTS_INDEX('Leaf', 'id', ['text'])")
except Exception as e:
raise RuntimeError(
"CREATE_FTS_INDEX failed (often ghost catalog after DROP INDEX). "
"Delete var/kb.lbug and run bin/brain/index.go --rebuild. "
f"Cause: {e}"
) from e
if "Leaf_vec" not in names:
try:
conn.execute(
"CALL CREATE_VECTOR_INDEX('Leaf', 'Leaf_vec', 'embedding', "
"metric := 'cosine')"
)
except Exception as e:
raise RuntimeError(
"CREATE_VECTOR_INDEX failed (often ghost catalog after DROP INDEX "
"Leaf.Leaf_vec → `_0_Leaf_vec_UPPER already exists in catalog`). "
"Delete var/kb.lbug and run bin/brain/index.go --rebuild. "
f"Cause: {e}"
) from e
names = leaf_index_names(conn)
missing = {"id", "Leaf_vec"} - names
if missing:
raise RuntimeError(
f"Leaf indexes incomplete after create: missing {sorted(missing)}; "
f"have {sorted(names)}. Delete var/kb.lbug and --rebuild."
)
def ensure_indexes(conn: ladybug.Connection) -> None:
"""Idempotent: create FTS + HNSW only when missing. Safe after upserts."""
create_fts_and_vector(conn, force=False)
def drop_indexes(conn: ladybug.Connection) -> None:
"""No-op. Kept for callers; DROP INDEX is fatal on Ladybug 0.19.
Historical note claimed "drop before bulk MERGE". Measured on 0.19:
- DROP FTS/VECTOR leaves ghost catalog → CREATE fails permanently until
`var/kb.lbug` is deleted.
- MERGE/upsert while **FTS** exists can corrupt FTS
("document for node offset N is missing during delete").
- Upsert while **HNSW** exists stays queryable.
Bulk rebuilders must delete `var/kb.lbug`, write all leafs (info+facts)
with no indexes, then `ensure_indexes()` once.
"""
return
def query_fts(conn: ladybug.Connection, text: str, limit: int = 10) -> list[dict]:
r = conn.execute(
"CALL QUERY_FTS_INDEX('Leaf', 'id', $q) "
"RETURN node.id, node.text, node.root, score ORDER BY score DESC LIMIT $n",
parameters={"q": text, "n": limit},
)
return [{"id": row[0], "text": row[1], "root": row[2], "score": row[3]} for row in r.get_all()]
def query_vector(conn: ladybug.Connection, embedding: list[float], limit: int = 10) -> list[dict]:
r = conn.execute(
"CALL QUERY_VECTOR_INDEX('Leaf', 'Leaf_vec', $q, $n) "
"RETURN node.id, node.text, node.root, distance ORDER BY distance LIMIT $n",
parameters={"q": embedding, "n": limit},
)
out = []
for row in r.get_all():
# distance -> similarity reasonable for cosine
score = 1.0 - row[3] if row[3] is not None else 0.0
out.append({"id": row[0], "text": row[1], "root": row[2], "score": score})
return out
def hybrid_search(conn: ladybug.Connection, embedding: list[float], fts_hits: list[dict],
limit: int = 10) -> list[dict]:
"""Merge FTS + vector by reciprocal rank fusion."""
fused: dict[str, dict] = {}
for rank, hit in enumerate(fts_hits):
fused.setdefault(hit["id"], {**hit, "rrf": 0.0})["rrf"] = 1.0 / (60 + rank + 1)
for rank, hit in enumerate(query_vector(conn, embedding, limit * 3)):
entry = fused.setdefault(hit["id"], {**hit, "rrf": 0.0})
entry["rrf"] += 1.0 / (60 + rank + 1)
entry.setdefault("score", hit.get("score", 0.0))
ranked = sorted(fused.values(), key=lambda h: h.get("rrf", 0.0), reverse=True)
return ranked[:limit]
def stats(conn: ladybug.Connection) -> dict:
r = conn.execute("MATCH (l:Leaf) RETURN l.root, count(*)")
rows = {row[0]: row[1] for row in r.get_all()}
total = conn.execute("MATCH (l:Leaf) RETURN count(*)").get_all()[0][0]
return {"total": total, "by_root": rows, "db": str(DB_PATH), "model": MODEL}
def open_readonly() -> tuple[ladybug.Database, ladybug.Connection]:
if not DB_PATH.exists():
raise FileNotFoundError(f"{DB_PATH} missing - run bin/brain/index.go --rebuild first")
db, conn = connect(read_only=True)
return db, conn
+231
View File
@@ -0,0 +1,231 @@
"""mailconv - pure helpers for bin/mail/import (mail -> markdown + attachments).
Shared with unit tests in bin/tools/test_mailconv.py. No network, no OnlyOffice
dependencies here: everything is `str -> str` or `Path -> str` so the tests run
offline against fixtures.
"""
from __future__ import annotations
import html
import os
import re
import subprocess
import tempfile
import zipfile
from pathlib import Path
# Body part / attachment file suffixes we know how to turn into markdown text.
TEXT_SUFFIXES = {".md", ".markdown", ".txt", ".csv", ".json", ".xml", ".yaml", ".yml", ".log", ".tsv",
".ics", ".ical", ".vcf", ".eml"}
OFFICE_SUFFIXES = {".docx", ".pptx", ".xlsx", ".html", ".htm", ".epub", ".eml", ".msg"}
PDF_SUFFIXES = {".pdf"}
IMAGE_SUFFIXES = {".png", ".jpg", ".jpeg", ".gif", ".bmp", ".tiff", ".tif", ".webp"}
ARCHIVE_SUFFIXES = {".zip"}
# Legacy binary Office (doc/xls/ppt) — markitdown skip them; we try
# pandoc first, else leave a stub.
LEGACY_OFFICE_SUFFIXES = {".doc", ".xls", ".ppt"}
TESS_LANG = "eng+deu"
CONVERTIBLE_SUFFIXES = (
TEXT_SUFFIXES | OFFICE_SUFFIXES | PDF_SUFFIXES | IMAGE_SUFFIXES | ARCHIVE_SUFFIXES | LEGACY_OFFICE_SUFFIXES
)
def clean_email_address(raw: str) -> str:
"""Extract the bare email from '"Name" <a@b.c>' and strip control chars."""
m = re.search(r"<([^<>@\s]+@[^<>@\s]+)>", raw)
return (m.group(1) if m else raw).strip()
def subject_to_filename(subject: str, max_len: int = 80) -> str:
"""Turn a mail subject into a filesystem-safe slug (keep first token readable)."""
s = re.sub(r"[^\w\-. ]+", "", subject).strip()
s = re.sub(r"\s+", "_", s)
s = s.strip("._")
if not s:
s = "untitled"
return s[:max_len] or "untitled"
def strip_html(html_text: str) -> str:
"""Naive HTML -> plain text fallback (used only if markitdown is missing)."""
import re as _re
text = _re.sub(r"(?is)<(script|style)[^>]*>.*?</\1>", "", html_text)
text = _re.sub(r"(?s)<br\s*/?>", "\n", text)
text = _re.sub(r"(?s)</p>", "\n\n", text)
text = _re.sub(r"(?s)<[^>]+>", "", text)
return html.unescape(text).strip()
def _unwrap_tables(html_text: str) -> str:
"""Unwrap mail HTML tables into pipe-joined text lines.
Outlook/Stripe-style emails wrap content in nested spacer/frame tables that
markitdown renders as hundreds of `--- |` cells and duplicated blocks.
Every <table> becomes plain "cell1 | cell2" lines (key-value pairs survive),
so only headings/paragraphs/links reach markitdown and no table noise is left.
"""
try:
from bs4 import BeautifulSoup
except Exception:
return html_text
soup = BeautifulSoup(html_text, "html.parser")
for table in reversed(soup.find_all("table")):
lines: list[str] = []
for row in table.find_all("tr"):
cells = [c.get_text(" ", strip=True) for c in row.find_all(["td", "th"])]
line = " | ".join(x for x in cells if x)
if line:
lines.append(line)
if lines:
table.replace_with(BeautifulSoup("\n".join(lines), "html.parser"))
else:
table.decompose()
return str(soup)
def html_to_markdown(html_text: str) -> str:
"""Convert a mail HTML body to markdown using markitdown when available."""
html_text = _unwrap_tables(html_text)
try:
from markitdown import MarkItDown
import io
md = MarkItDown()
result = md.convert_stream(io.BytesIO(html_text.encode("utf-8", errors="replace")),
file_extension=".html")
text = result.text_content.strip()
if text:
return normalize_markdown(text)
except Exception:
pass
return normalize_markdown(strip_html(html_text))
def normalize_markdown(text: str) -> str:
"""Collapse the pdfminer/markitdown NUL artifacts and stray control chars."""
# NUL bytes that pdfminer inserts between digits/letters.
text = text.replace("\x00", "")
# Email spacer noise: zero-width chars, soft hyphens, figure spaces,
# combining grapheme joiner, BOM.
for ch in ("\ufeff", "\u200b", "\u034f", "\u00ad", "\u2007", "\u2008", "\u200a", "\u2002"):
text = text.replace(ch, "")
text = re.sub(r"[ \t]{2,}", " ", text)
# Trim trailing whitespace per line so space-only spacer rows collapse.
text = "\n".join(l.rstrip() for l in text.split("\n"))
# Collapse 3+ blank lines to two.
text = re.sub(r"\n{3,}", "\n\n", text)
# Remove weird trailing control chars.
text = "".join(ch for ch in text if ch >= " " or ch in "\n\t")
return text.strip()
def split_zip_members(zip_path: Path) -> list[str]:
"""Return safe member names of a zip archive (skips dir entries)."""
try:
with zipfile.ZipFile(zip_path) as zf:
return [m for m in zf.namelist() if not m.endswith("/")]
except zipfile.BadZipFile:
return []
def zip_extract_safe(zip_path: Path, dest: Path) -> list[Path]:
"""Extract a zip into dest guarding against path traversal; returns files."""
out: list[Path] = []
try:
with zipfile.ZipFile(zip_path) as zf:
for member in zf.infolist():
if member.is_dir():
continue
target = (dest / member.filename).resolve()
if not target.is_relative_to(dest.resolve()):
continue
target.parent.mkdir(parents=True, exist_ok=True)
with zf.open(member) as src, open(target, "wb") as dst:
dst.write(src.read())
out.append(target)
except zipfile.BadZipFile:
return []
return out
def is_convertible(suffix: str) -> bool:
return suffix.lower() in CONVERTIBLE_SUFFIXES
def convert_pdf(path: Path, ocr: bool = False) -> str:
"""pdftotext -layout first; empty text layer → pdftoppm + tesseract.
`ocr` is unused for born-digital PDFs (text layer wins). Scans OCR
automatically. This path never execs an ONNX document converter.
"""
del ocr # scans OCR when the text layer is empty; flag is for images
text = pdf_fast_text(path)
if text and text.strip():
return normalize_markdown(text)
scanned = ocr_pdf(path)
if scanned and scanned.strip():
return normalize_markdown(scanned)
if text:
return normalize_markdown(text)
return "\n<!-- pdf has no text layer (ocr unavailable) -->\n"
def pdf_fast_text(path: Path) -> str | None:
"""pdftotext -layout; None when poppler is missing or the command fails."""
try:
proc = subprocess.run(
["pdftotext", "-layout", str(path), "-"],
capture_output=True, timeout=60)
except (OSError, subprocess.TimeoutExpired):
return None
if proc.returncode != 0:
return None
return proc.stdout.decode("utf-8", errors="replace")
def ocr_pdf(path: Path) -> str:
"""Rasterize with pdftoppm and OCR each page (tesseract or paddle)."""
try:
with tempfile.TemporaryDirectory(prefix="2dph-ocr-") as tmp:
prefix = str(Path(tmp) / "page")
proc = subprocess.run(
["pdftoppm", "-png", "-r", "200", str(path), prefix],
capture_output=True, timeout=120)
if proc.returncode != 0:
return ""
pages = sorted(Path(tmp).glob("page*.png"))
parts = [ocr_image(p) for p in pages]
return "\n\n".join(p for p in parts if p and p.strip())
except (OSError, subprocess.TimeoutExpired):
return ""
def ocr_image(path: Path) -> str:
engine = os.environ.get("OCR_ENGINE", "tesseract")
if engine == "paddle":
return _ocr_paddle(path)
return _ocr_tesseract(path)
def _ocr_tesseract(path: Path) -> str:
try:
proc = subprocess.run(
["tesseract", str(path), "stdout", "-l", TESS_LANG, "--psm", "6"],
capture_output=True, timeout=120)
except (OSError, subprocess.TimeoutExpired):
return ""
if proc.returncode != 0:
return ""
return proc.stdout.decode("utf-8", errors="replace").strip()
def _ocr_paddle(path: Path) -> str:
try:
proc = subprocess.run(
["paddleocr", "ocr", "-i", str(path)],
capture_output=True, timeout=180)
except (OSError, subprocess.TimeoutExpired):
return ""
if proc.returncode != 0:
return ""
return proc.stdout.decode("utf-8", errors="replace").strip()
+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
+300
View File
@@ -0,0 +1,300 @@
"""D14 layout: bin/{subject}/{method}.go, libs in internal/, one go.mod."""
from __future__ import annotations
import os
import unittest
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
class BinLayoutTest(unittest.TestCase):
def test_brain_search_shebang_exists(self) -> None:
p = ROOT / "bin" / "brain" / "search.go"
self.assertTrue(p.is_file(), "missing bin/brain/search.go")
first = p.read_text().splitlines()[0]
self.assertTrue(
first.startswith("//usr/bin/env go run"),
f"shebang first line, got {first!r}",
)
def test_no_nested_go_mod_under_bin(self) -> None:
nested = list((ROOT / "bin").rglob("go.mod"))
self.assertEqual(nested, [], f"nested go.mod files: {nested}")
def test_rank_lives_in_internal_brain(self) -> None:
self.assertTrue(
(ROOT / "internal" / "brain" / "rank" / "rank.go").is_file(),
"ranking must live in internal/brain/rank (cgo-free)",
)
self.assertFalse(
(ROOT / "bin" / "kbsearch").exists(),
"bin/kbsearch nested module must be gone",
)
def test_no_main_go_under_bin_brain(self) -> None:
main = ROOT / "bin" / "brain" / "main.go"
self.assertFalse(main.exists(), "bin/brain/main.go is not a method")
def test_chats_methods_are_shebangs_not_main(self) -> None:
chats = ROOT / "bin" / "chats"
self.assertFalse(
(chats / "main.go").exists(),
"bin/chats/main.go is a dispatcher, not a method",
)
self.assertFalse(
(chats / "index_cmd.go").exists(),
"chats index is a brain write hiding under the wrong subject",
)
for method in ("sync.go", "import.go", "facts.go", "apply.go"):
p = chats / method
self.assertTrue(p.is_file(), f"missing bin/chats/{method}")
first = p.read_text().splitlines()[0]
self.assertTrue(
first.startswith("//usr/bin/env go run"),
f"{method} shebang, got {first!r}",
)
def test_chats_lib_lives_in_internal(self) -> None:
self.assertTrue(
(ROOT / "internal" / "chats" / "linkedin.go").is_file(),
"LinkedIn parser must live in internal/chats",
)
self.assertFalse(
(ROOT / "bin" / "chats" / "linkedin.go").exists(),
"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", "add.go", "get.go", "stats.go", "eval.go", "watch.go"):
self._assert_shebang(f"bin/brain/{method}")
def test_brain_add_is_python_write_not_rebuild(self) -> None:
self._assert_shebang("bin/brain/add.go")
text = (ROOT / "bin" / "brain" / "add.go").read_text()
self.assertIn("cmdbin.ExecFile", text)
self.assertIn("bin/kb/add", text)
self.assertNotIn("--rebuild", text)
py = (ROOT / "bin" / "kb" / "add").read_text()
self.assertIn("add_leafs", py)
self.assertIn("--json", py)
self.assertNotIn("unlink", py.lower())
def test_brain_get_stats_eval_are_not_python_exec(self) -> None:
for method in ("get.go", "stats.go", "eval.go"):
text = (ROOT / "bin" / "brain" / method).read_text()
self.assertNotIn(
"ExecFile",
text,
f"bin/brain/{method} must call internal/brain, not ExecFile Python",
)
self.assertNotIn(
"cmdbin",
text,
f"bin/brain/{method} must not import internal/cmdbin",
)
self.assertIn(
"system_ladybug",
text.splitlines()[0],
f"bin/brain/{method} shebang must pass -tags=system_ladybug",
)
self.assertIn(
"github.com/eSlider/2dph/internal/brain",
text,
)
def test_eval_control_questions_live_in_rank(self) -> None:
rank = (ROOT / "internal" / "brain" / "rank" / "evalq.go").read_text()
py = (ROOT / "bin" / "kb" / "eval").read_text()
for frag in ("BM25", "DevOps", "LadybugDB"):
self.assertIn(frag, rank)
self.assertIn(frag, py)
self.assertIn("0.95", rank)
def test_facts_methods_are_shebangs(self) -> None:
for method in ("audit.go", "extract.go", "crm.go"):
self._assert_shebang(f"bin/facts/{method}")
text = (ROOT / "bin" / "facts" / method).read_text()
self.assertIn("cmdbin.ExecFile", text)
self.assertIn(f"bin/facts/{method.removesuffix('.go')}", text)
def test_d16_adjudication_is_cgo_free(self) -> None:
self.assertTrue((ROOT / "internal" / "facts" / "contradict.go").is_file())
go = (ROOT / "internal" / "facts" / "contradict.go").read_text()
py = (ROOT / "bin" / "tools" / "contradict.py").read_text()
audit = (ROOT / "bin" / "facts" / "audit").read_text()
for token in ("temporal_freshness", "authority_pairing", "unresolved"):
self.assertIn(token, go)
self.assertIn(token, py)
self.assertIn("contradict", audit)
self.assertIn(" vs ", py)
plan = (ROOT / "PLAN.md").read_text()
self.assertIn("temporal_freshness", plan)
self.assertIn("authority_pairing", plan)
shebang = (ROOT / "bin" / "facts" / "audit.go").read_text()
self.assertIn("contradict", shebang)
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_mail_ocr_is_tesseract_not_docling(self) -> None:
self._assert_shebang("bin/mail/ocr.go")
ocr = (ROOT / "bin" / "mail" / "ocr.go").read_text()
self.assertIn("internal/ocr", ocr)
self.assertIn("mail_ocr", ocr)
self.assertNotIn("github.com/otiai10/gosseract", ocr)
py = (ROOT / "bin" / "mail" / "import").read_text()
self.assertNotIn("from docling", py)
self.assertNotIn("import docling", py)
self.assertIn("convert_pdf", py)
conv = (ROOT / "bin" / "tools" / "mailconv.py").read_text()
self.assertIn("pdftotext", conv)
self.assertIn("pdftoppm", conv)
self.assertIn("tesseract", conv)
self.assertIn("eng+deu", conv)
self.assertNotIn("from docling", conv)
self.assertNotIn("import docling", conv)
self.assertNotIn("gocv", conv.lower())
proj = (ROOT / "pyproject.toml").read_text()
self.assertNotIn("docling", proj)
ci = (ROOT / ".github" / "workflows" / "ci.yml").read_text()
self.assertIn("tesseract-ocr", ci)
self.assertIn("./internal/ocr", ci)
compose = (ROOT / "compose.yaml").read_text()
self.assertIn("ocr-paddle", compose)
self.assertIn("OCR_ENGINE", compose)
def test_markdown_import_is_go_not_python_exec(self) -> None:
self._assert_shebang("bin/markdown/import.go")
text = (ROOT / "bin" / "markdown" / "import.go").read_text()
self.assertNotIn("ExecFile", text)
self.assertNotIn("cmdbin", text)
self.assertIn("internal/mdleaves", text)
self.assertNotIn("kb.lbug", text)
def test_import_adapters_do_not_write_ladybug(self) -> None:
for rel in (
"bin/mail/import.go",
"bin/mail/import",
"bin/markdown/import.go",
"bin/chats/import.go",
"bin/git/import.go",
):
text = (ROOT / rel).read_text()
self.assertNotIn("upsert_leaf", text, rel)
self.assertNotIn("kb.lbug", text, rel)
self.assertNotIn("var/brain.lbug", text, rel)
index = (ROOT / "bin" / "brain" / "index.go").read_text()
self.assertIn("bin/kb/index", index)
def test_postgres_query_is_shebang(self) -> None:
self._assert_shebang("bin/postgres/query.go")
def test_git_import_is_gogit_shebang(self) -> None:
self._assert_shebang("bin/git/import.go")
py = (ROOT / "bin" / "git" / "import").read_text()
self.assertNotIn(
'["git"',
py,
"Python git/import must not subprocess the git binary",
)
self.assertIn("bin/git/import.go", py)
def test_web_search_is_shebang(self) -> None:
self._assert_shebang("bin/web/search.go")
py = (ROOT / "bin" / "web" / "search").read_text()
self.assertIn("bin/web/search.go", py)
def test_gitimport_py_has_no_git_binary(self) -> None:
py = (ROOT / "bin" / "tools" / "gitimport.py").read_text()
self.assertNotIn("subprocess", py)
self.assertNotIn("git log", py)
def test_gogit_is_direct_go_mod_require(self) -> None:
text = (ROOT / "go.mod").read_text()
first = text.split("require (")[1].split(")")[0]
self.assertRegex(first, r"github.com/go-git/go-git/v5\s+v")
for line in first.splitlines():
if "go-git/go-git" in line:
self.assertNotIn("indirect", line)
def test_duckdb_go_is_direct_require(self) -> None:
text = (ROOT / "go.mod").read_text()
first = text.split("require (")[1].split(")")[0]
self.assertRegex(first, r"github.com/duckdb/duckdb-go/v2\s+v")
for line in first.splitlines():
if "duckdb/duckdb-go" in line:
self.assertNotIn("indirect", line)
skill = (ROOT / "skills" / "duckdb" / "SKILL.md").read_text()
self.assertIn("github.com/duckdb/duckdb-go", skill)
self.assertIn("Ladybug", skill)
self.assertIn("sqlite", skill.lower())
self.assertIn("gcc", skill.lower())
self.assertIn("Zig", skill)
plan = (ROOT / "PLAN.md").read_text()
self.assertIn("D22", plan)
self.assertIn("duckdb-go", plan)
self._assert_shebang("bin/qa/stats.go")
reasoner = (ROOT / "internal" / "reasoner" / "client.go").read_text()
self.assertNotIn("duckdb", reasoner)
self.assertNotIn("duckstats", reasoner)
bakeoff = (ROOT / "bin" / "reasoner" / "bakeoff.go").read_text()
self.assertIn("internal/duckstats", bakeoff)
webcache = (ROOT / "internal" / "websearch" / "cache.go").read_text()
self.assertNotIn("duckdb", webcache)
self.assertIn("modernc.org/sqlite", webcache)
def test_cgo_uses_zig_not_gcc(self) -> None:
for rel in ("bin/cgo/zig", "bin/cgo/zcc", "bin/cgo/zc++"):
p = ROOT / rel
self.assertTrue(p.is_file(), f"missing {rel}")
self.assertTrue(
os.access(p, os.X_OK),
f"{rel} must be executable",
)
zig = (ROOT / "bin" / "cgo" / "zig").read_text()
self.assertIn("zig cc", zig)
self.assertIn("0.14.1", zig)
zcc = (ROOT / "bin" / "cgo" / "zcc").read_text()
self.assertIn('exec "$ZIG" cc', zcc)
self.assertNotIn("command -v gcc", zcc)
search = (ROOT / "bin" / "kb" / "search").read_text()
self.assertIn("bin/cgo/zig", search)
self.assertNotIn("command -v gcc", search)
def test_ci_recall_sot_is_zig_brain_eval(self) -> None:
ci = (ROOT / ".github" / "workflows" / "ci.yml").read_text()
self.assertIn("bin/brain/eval.go", ci)
self.assertIn("system_ladybug,brain_eval", ci)
self.assertIn("/tmp/brain-eval", ci)
self.assertIn("KB_ROOT", ci)
self.assertNotIn("bin/kb/eval", ci)
self.assertNotIn("gate skipped", ci)
self.assertIn("./bin/facts/audit self", ci)
def test_eval_fragments_live_in_default_corpus(self) -> None:
"""CI --rebuild indexes README/PLAN/docs/skills; fragments must be there."""
corpus = []
for rel in ("README.md", "PLAN.md", "AGENTS.md"):
corpus.append((ROOT / rel).read_text())
for d in ("docs", "skills"):
for p in (ROOT / d).rglob("*.md"):
corpus.append(p.read_text())
blob = "\n".join(corpus)
for frag in ("BM25", "DevOps", "LadybugDB"):
self.assertIn(frag, blob, f"{frag} must appear in default index corpus")
+104
View File
@@ -0,0 +1,104 @@
import os
import sys
import unittest
sys.path.insert(0, os.path.dirname(__file__))
from contradict import ( # noqa: E402
RULE_AUTHORITY,
RULE_SINGLE,
RULE_TEMPORAL,
RULE_TWO_SOURCE,
RULE_UNRESOLVED,
adjudicate,
check_fact_row,
parse_source_field,
)
def src(i, kind, stale=False):
return {"id": i, "kind": kind, "stale": stale}
class TestContradict(unittest.TestCase):
def test_two_vs_two_stays_hypothesis(self):
r = adjudicate({
"text": "svc listens on 443",
"yes": [src("docker-ps", "runtime"), src("compose", "config")],
"no": [src("docker-old", "runtime"), src("compose-old", "config")],
})
self.assertFalse(r["confirmed"])
self.assertEqual(r["rule"], RULE_UNRESOLVED)
self.assertEqual(r["winner"], "")
def test_temporal_freshness(self):
r = adjudicate({
"text": "svc listens on 443",
"yes": [src("docker-ps", "runtime"), src("compose", "config")],
"no": [src("old-readme", "narrative", True), src("old-wiki", "narrative", True)],
})
self.assertTrue(r["confirmed"])
self.assertEqual(r["rule"], RULE_TEMPORAL)
self.assertEqual(r["winner"], "yes")
def test_authority_pairing(self):
r = adjudicate({
"text": "svc listens on 443",
"yes": [src("docker-ps", "runtime"), src("compose", "config")],
"no": [src("readme", "narrative"), src("wiki", "narrative")],
})
self.assertTrue(r["confirmed"])
self.assertEqual(r["rule"], RULE_AUTHORITY)
self.assertEqual(r["winner"], "yes")
def test_two_source_and_single(self):
two = adjudicate({
"text": "arc-1 runs Matrix",
"yes": [src("compose", "config"), src("docker-ps", "runtime")],
})
self.assertTrue(two["confirmed"])
self.assertEqual(two["rule"], RULE_TWO_SOURCE)
one = adjudicate({"text": "maybe", "yes": [src("readme", "narrative")]})
self.assertFalse(one["confirmed"])
self.assertEqual(one["rule"], RULE_SINGLE)
def test_parse_source_field(self):
yes, no = parse_source_field("docker ps x compose.yml vs old.md x wiki.md")
self.assertIn(" x ", yes)
self.assertIn(" x ", no)
def test_check_fact_row_allows_hypothesis_vs(self):
p = check_fact_row(
"L1", "a.md x b.md vs c.md x d.md", "var/", "audit", "hypothesis",
)
self.assertEqual(p, [])
p = check_fact_row("L2", "a.md x b.md", "var/", "audit", "confirmed")
self.assertEqual(p, [])
p = check_fact_row("L3", "a.md x b.md vs c.md x d.md", "var/", "audit", "confirmed")
self.assertTrue(any("vs-contradiction" in x for x in p))
p = check_fact_row("L4", "only-one.md", "var/", "audit", "hypothesis")
self.assertTrue(any("a x b vs" in x for x in p))
def test_audit_contradict_cli_unresolved(self):
import json
import subprocess
from pathlib import Path
root = Path(__file__).resolve().parents[2]
payload = json.dumps({
"text": "svc 443",
"yes": [src("a", "runtime"), src("b", "config")],
"no": [src("c", "runtime"), src("d", "config")],
})
proc = subprocess.run(
[sys.executable, str(root / "bin" / "facts" / "audit"), "contradict", "--json"],
input=payload, capture_output=True, text=True, check=False,
)
self.assertEqual(proc.returncode, 0, proc.stderr)
out = json.loads(proc.stdout)
self.assertTrue(out["ok"])
self.assertEqual(out["contradictions"][0]["rule"], RULE_UNRESOLVED)
self.assertFalse(out["contradictions"][0]["confirmed"])
if __name__ == "__main__":
unittest.main()
+46
View File
@@ -0,0 +1,46 @@
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent))
import crmfacts # noqa: E402
FM = """\
schema: 2
meta:
title: x
orgs:
- id: produktor
label: ProProdukt SL / produktor.io
kind: own
period: 2006present
website: https://produktor.io
- id: dyvenia
label: Dyvenia
kind: employer
period: 20232025
clients:
- name: One
- name: Two
timeline:
- start: 2001
"""
class CorpusOrgsTest(unittest.TestCase):
def test_parses_label_kind_period(self):
orgs = crmfacts.corpus_orgs(FM)
self.assertEqual(orgs["produktor"]["label"], "ProProdukt SL / produktor.io")
self.assertEqual(orgs["produktor"]["kind"], "own")
self.assertEqual(orgs["dyvenia"]["kind"], "employer")
def test_does_not_leak_clients_into_orgs(self):
orgs = crmfacts.corpus_orgs(FM)
self.assertNotIn("One", orgs)
self.assertNotIn("Two", orgs)
self.assertNotIn("timeline", orgs)
if __name__ == "__main__":
unittest.main()
+74
View File
@@ -0,0 +1,74 @@
import os
import sys
import tempfile
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent))
import kblib # noqa: E402
import gitimport # noqa: E402
COMMIT_PERSON_SCHEMA = (
"CREATE NODE TABLE IF NOT EXISTS Commit (id STRING, repo STRING, subject STRING, "
"author STRING, email STRING, date STRING, PRIMARY KEY(id))"
)
PERSON_SCHEMA = (
"CREATE NODE TABLE IF NOT EXISTS Person (id STRING, name STRING, email STRING, PRIMARY KEY(id))"
)
HAS_VERSION_SCHEMA = "CREATE REL TABLE IF NOT EXISTS HAS_VERSION (FROM File TO Commit)"
AUTHORED_SCHEMA = "CREATE REL TABLE IF NOT EXISTS AUTHORED (FROM Commit TO Person)"
def sample_commit() -> gitimport.Commit:
return gitimport.Commit(
sha="a1b2c3d",
author="Ada Lovelace",
email="ada@example.com",
date="2026-08-10T12:00:00+01:00",
subject="feat: first commit",
files=["README.md", "src/main.c"],
)
class GitGraphTest(unittest.TestCase):
def setUp(self):
self.dir = tempfile.mkdtemp()
self.dbpath = os.path.join(self.dir, "kb.lbug")
self.db, self.conn = kblib.connect(self.dbpath, read_only=False)
kblib.init_schema(self.conn)
self.conn.execute(COMMIT_PERSON_SCHEMA)
self.conn.execute(PERSON_SCHEMA)
self.conn.execute(HAS_VERSION_SCHEMA)
self.conn.execute(AUTHORED_SCHEMA)
def tearDown(self):
self.conn.close()
self.db.close()
def test_index_commits_creates_nodes_and_edges(self):
gitimport.index_commits(self.conn, [sample_commit()], "sample-repo")
rp = self.conn.execute("MATCH (p:Person) RETURN p.name, p.email").get_all()
self.assertEqual([tuple(r) for r in rp], [("Ada Lovelace", "ada@example.com")])
rc = self.conn.execute("MATCH (c:Commit) RETURN c.id, c.repo").get_all()
self.assertEqual(len(rc), 1)
self.assertEqual(rc[0][1], "sample-repo")
rf = self.conn.execute(
"MATCH (f:File)-[:HAS_VERSION]->(c:Commit)-[:AUTHORED]->(p:Person) "
"RETURN f.path, c.id, p.email").get_all()
paths = sorted(r[0] for r in rf)
self.assertEqual(paths, ["README.md", "src/main.c"])
self.assertTrue(all(r[2] == "ada@example.com" for r in rf))
def test_index_commits_idempotent(self):
cs = [sample_commit()]
gitimport.index_commits(self.conn, cs, "sample-repo")
gitimport.index_commits(self.conn, cs, "sample-repo")
n = self.conn.execute("MATCH (c:Commit) RETURN count(*)").get_all()[0][0]
self.assertEqual(n, 1)
p = self.conn.execute("MATCH (p:Person) RETURN count(*)").get_all()[0][0]
self.assertEqual(p, 1)
if __name__ == "__main__":
unittest.main()
+97
View File
@@ -0,0 +1,97 @@
"""Import adapters write files only. Index rebuild is brain/index (D14 / Gitea #7)."""
from __future__ import annotations
import json
import os
import subprocess
import sys
import tempfile
import unittest
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
class IndexAdapterTest(unittest.TestCase):
def test_dry_run_fixture_corpus_does_not_write_lbug(self) -> None:
tmp = Path(tempfile.mkdtemp())
(tmp / "note.md").write_text("# Fixture\n\n## Leaf\n\nhello corpus\n", encoding="utf-8")
lbug = tmp / "kb.lbug"
try:
import ladybug # noqa: F401
except ImportError:
venv_py = ROOT / ".venv" / "bin" / "python"
if not venv_py.is_file():
self.skipTest("ladybug missing")
py = str(venv_py)
else:
py = sys.executable
proc = subprocess.run(
[py, str(ROOT / "bin" / "kb" / "index"), "--dry-run", "--json", "--corpus", str(tmp)],
cwd=ROOT,
capture_output=True,
text=True,
env=os.environ.copy(),
check=False,
)
self.assertEqual(proc.returncode, 0, proc.stderr)
msg = json.loads(proc.stdout)
self.assertTrue(msg.get("dry_run"))
self.assertGreaterEqual(msg.get("corpus_total", 0), 1)
self.assertFalse(lbug.exists(), "dry-run must not create a Ladybug file")
def test_facts_json_and_chats_land_on_rebuild(self) -> None:
"""Gitea #18: facts (2-source) + chats markdown become leafs on rebuild."""
tmp = Path(tempfile.mkdtemp())
dbpath = tmp / "kb.lbug"
chats = tmp / "chats"
chats.mkdir()
(chats / "alice.md").write_text(
"# Chat\n\n## Alice and Bob\n\nhello from chats fixture unique-chat-token\n",
encoding="utf-8",
)
facts_path = tmp / "facts.json"
facts_path.write_text(json.dumps([{
"text": "container 'brain' unique-fact-token is running and declared in compose.yaml",
"source": "docker ps x compose.yaml",
"loc": "compose.yaml:brain",
"how": "facts/extract",
}]), encoding="utf-8")
venv_py = ROOT / ".venv" / "bin" / "python"
py = str(venv_py) if venv_py.is_file() else sys.executable
proc = subprocess.run(
[
py, str(ROOT / "bin" / "kb" / "index"),
"--rebuild", "--db", str(dbpath), "--no-defaults",
"--with-chats", str(chats),
"--facts-json", str(facts_path),
"--json",
],
cwd=ROOT,
capture_output=True,
text=True,
env=os.environ.copy(),
check=False,
)
self.assertEqual(proc.returncode, 0, proc.stderr)
msg = json.loads(proc.stdout)
self.assertGreaterEqual(msg.get("facts_leafs", 0), 1)
self.assertGreaterEqual(msg.get("chat_leafs", 0), 1)
self.assertTrue(dbpath.exists())
sys.path.insert(0, str(ROOT / "bin" / "tools"))
import kblib
db, conn = kblib.connect(dbpath, read_only=True)
try:
stats = kblib.stats(conn)
self.assertGreaterEqual(stats["by_root"].get("facts", 0), 1)
fts = kblib.query_fts(conn, "unique-chat-token", 5)
self.assertTrue(fts, "chats markdown must be FTS-searchable")
fact_hits = kblib.query_fts(conn, "unique-fact-token", 5)
self.assertTrue(any(h.get("root") == "facts" for h in fact_hits))
src = conn.execute(
"MATCH (l:Leaf {root:'facts'}) RETURN l.source"
).get_all()
self.assertTrue(any(" x " in str(r[0]) for r in src))
finally:
conn.close()
db.close()
+69
View File
@@ -0,0 +1,69 @@
"""Incremental add writes leafs without deleting kb.lbug."""
from __future__ import annotations
import json
import os
import subprocess
import sys
import tempfile
import unittest
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
class KbAddCLITest(unittest.TestCase):
def test_json_add_does_not_delete_db(self) -> None:
tmp = Path(tempfile.mkdtemp())
dbpath = tmp / "kb.lbug"
py = sys.executable
venv_py = ROOT / ".venv" / "bin" / "python"
if venv_py.is_file():
py = str(venv_py)
payload = {
"text": "cli zebra leaf",
"root": "info",
"source": "cli-test",
"confidence": "confirmed",
"how": "test",
"loc": str(tmp),
"type": "reference",
"embedding": [0.0] * 256,
}
payload["embedding"][0] = 0.3
proc = subprocess.run(
[py, str(ROOT / "bin" / "kb" / "add"), "--db", str(dbpath), "--json"],
cwd=ROOT,
input=json.dumps(payload),
capture_output=True,
text=True,
env=os.environ.copy(),
check=False,
)
self.assertEqual(proc.returncode, 0, proc.stderr)
self.assertTrue(dbpath.exists(), "add must create the db, not skip write")
out = json.loads(proc.stdout)
self.assertEqual(out.get("mode"), "add")
self.assertEqual(len(out.get("ids") or []), 1)
again = subprocess.run(
[py, str(ROOT / "bin" / "kb" / "add"), "--db", str(dbpath), "--json"],
cwd=ROOT,
input=json.dumps({
**payload,
"text": "second moose leaf",
"source": "cli-test-2",
}),
capture_output=True,
text=True,
env=os.environ.copy(),
check=False,
)
self.assertEqual(again.returncode, 0, again.stderr)
self.assertTrue(dbpath.exists())
second = json.loads(again.stdout)
self.assertEqual(len(second.get("ids") or []), 1)
self.assertNotEqual(out["ids"][0], second["ids"][0])
if __name__ == "__main__":
unittest.main()
+199
View File
@@ -0,0 +1,199 @@
import os
import sys
import tempfile
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent))
import kblib # noqa: E402
def make_emb(value: float) -> list[float]:
vec = [0.0] * kblib.EMBED_DIM
vec[0] = value
return vec
class KblibTest(unittest.TestCase):
def setUp(self):
self.dir = tempfile.mkdtemp()
self.dbpath = os.path.join(self.dir, "kb.lbug")
self.db, self.conn = kblib.connect(self.dbpath, read_only=False)
kblib.init_schema(self.conn)
def tearDown(self):
self.conn.close()
self.db.close()
def test_leaf_id_is_stable(self):
self.assertEqual(kblib.leaf_id("abc", "src"), kblib.leaf_id("abc", "src"))
self.assertNotEqual(kblib.leaf_id("abc", "src"), kblib.leaf_id("abd", "src"))
def test_upsert_roundtrip(self):
kblib.upsert_leaf(self.conn, text="the quick brown fox", root="info",
confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference",
embedding=make_emb(1.0))
kblib.ensure_indexes(self.conn)
hits = kblib.query_fts(self.conn, "fox", 5)
self.assertEqual(len(hits), 1)
self.assertEqual(hits[0]["root"], "info")
def test_hybrid_ranks_vector_match(self):
kblib.upsert_leaf(self.conn, text="the quick brown fox", root="info",
confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference",
embedding=make_emb(1.0))
kblib.upsert_leaf(self.conn, text="a lazy dog sleeps", root="info",
confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference",
embedding=make_emb(0.0))
kblib.ensure_indexes(self.conn)
result = kblib.hybrid_search(self.conn, make_emb(1.0), [], 5)
self.assertTrue(result)
self.assertIn("rrf", result[0])
self.assertEqual(result[0]["text"], "the quick brown fox")
def test_upsert_keeps_hnsw_queryable(self):
"""Upsert while HNSW exists must not kill vector search."""
kblib.upsert_leaf(self.conn, text="seed leaf", root="info",
confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference",
embedding=make_emb(0.2))
kblib.ensure_indexes(self.conn)
self.assertIn("Leaf_vec", kblib.leaf_index_names(self.conn))
kblib.upsert_leaf(self.conn, text="added after index", root="facts",
confidence="confirmed", source="a.md x b.md",
source_rev="r1", how="test", loc="/tmp", type_="fact",
embedding=make_emb(0.9))
hits = kblib.query_vector(self.conn, make_emb(0.9), 5)
self.assertTrue(hits)
self.assertIn("Leaf_vec", kblib.leaf_index_names(self.conn))
def test_add_after_indexes_keeps_fts_queryable(self):
"""Incremental add after FTS+HNSW must find the new leaf on both indexes."""
kblib.upsert_leaf(self.conn, text="seed fox leaf", root="info",
confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference",
embedding=make_emb(0.1))
kblib.ensure_indexes(self.conn)
ids = kblib.add_leafs(self.conn, [{
"text": "added zebra after index",
"root": "facts",
"confidence": "confirmed",
"source": "a.md x b.md",
"source_rev": "r1",
"how": "test",
"loc": "/tmp",
"type": "fact",
"embedding": make_emb(0.9),
}])
self.assertEqual(len(ids), 1)
fts = kblib.query_fts(self.conn, "zebra", 5)
self.assertTrue(fts)
self.assertIn("zebra", fts[0]["text"])
self.assertEqual(fts[0]["root"], "facts")
vec = kblib.query_vector(self.conn, make_emb(0.9), 5)
self.assertTrue(any("zebra" in h["text"] for h in vec))
fox = kblib.query_fts(self.conn, "fox", 5)
self.assertTrue(fox)
self.assertIn("fox", fox[0]["text"])
def test_add_facts_and_info_one_transaction(self):
"""D12: facts and info land in the same transaction."""
kblib.ensure_indexes(self.conn)
ids = kblib.add_leafs(self.conn, [
{
"text": "tx fact leaf two-source",
"root": "facts",
"confidence": "confirmed",
"source": "compose.yml x docker ps",
"source_rev": "r1",
"how": "test",
"loc": "/tmp",
"type": "fact",
"embedding": make_emb(0.4),
},
{
"text": "tx info narrative",
"root": "info",
"confidence": "confirmed",
"source": "note.md",
"source_rev": "r1",
"how": "test",
"loc": "/tmp",
"type": "reference",
"embedding": make_emb(0.5),
},
])
self.assertEqual(len(ids), 2)
stats = kblib.stats(self.conn)
self.assertEqual(stats["by_root"].get("facts"), 1)
self.assertEqual(stats["by_root"].get("info"), 1)
self.assertTrue(kblib.query_fts(self.conn, "two-source", 5))
self.assertTrue(kblib.query_fts(self.conn, "narrative", 5))
def test_drop_vector_then_create_raises_clear_error(self):
"""DROP INDEX leaves ghost catalog; create_fts_and_vector must raise."""
kblib.upsert_leaf(self.conn, text="seed", root="info",
confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference",
embedding=make_emb(0.1))
kblib.ensure_indexes(self.conn)
self.conn.execute("DROP INDEX IF EXISTS Leaf.Leaf_vec")
with self.assertRaises(RuntimeError) as ctx:
kblib.create_fts_and_vector(self.conn, force=True)
msg = str(ctx.exception)
self.assertIn("CREATE_VECTOR_INDEX failed", msg)
self.assertIn("--rebuild", msg)
def test_stats_counts_roots(self):
kblib.upsert_leaf(self.conn, text="a fact leaf", root="facts",
confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference",
embedding=make_emb(0.5))
kblib.upsert_leaf(self.conn, text="an info leaf", root="info",
confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference",
embedding=make_emb(0.5))
stats = kblib.stats(self.conn)
self.assertEqual(stats["total"], 2)
self.assertEqual(stats["by_root"], {"facts": 1, "info": 1})
def test_hop_1_returns_file_hop_3_reaches_person(self):
"""--hop walks FROM_FILE / HAS_VERSION / AUTHORED (Gitea #17)."""
import gitimport
lid = kblib.upsert_leaf(
self.conn, text="readme hop fixture", root="info",
confidence="confirmed", source="README.md", source_rev="r1",
how="test", loc="README.md", type_="reference",
embedding=make_emb(0.3),
)
kblib.link_from_file(self.conn, lid, "README.md", repo="sample-repo")
gitimport.index_commits(self.conn, [gitimport.Commit(
sha="a1b2c3d",
author="Ada Lovelace",
email="ada@example.com",
date="2026-08-10T12:00:00Z",
subject="feat: first commit",
files=["README.md"],
)], "sample-repo")
hop1 = kblib.hop_walk(self.conn, lid, 1)
self.assertEqual(len(hop1), 1)
self.assertEqual(hop1[0]["label"], "File")
self.assertEqual(hop1[0]["name"], "README.md")
self.assertEqual(hop1[0]["depth"], 1)
hop3 = kblib.hop_walk(self.conn, lid, 3)
labels = {n["label"] for n in hop3}
self.assertIn("File", labels)
self.assertIn("Commit", labels)
self.assertIn("Person", labels)
person = [n for n in hop3 if n["label"] == "Person"][0]
self.assertEqual(person["name"], "Ada Lovelace")
self.assertEqual(person["depth"], 3)
if __name__ == "__main__":
unittest.main()
+212
View File
@@ -0,0 +1,212 @@
import io
import os
import sys
import unittest
import zipfile
from pathlib import Path
sys.path.insert(0, os.path.dirname(__file__))
from mailconv import ( # noqa: E402
TESS_LANG,
clean_email_address,
convert_pdf,
html_to_markdown,
is_convertible,
normalize_markdown,
ocr_image,
split_zip_members,
subject_to_filename,
zip_extract_safe,
)
from mailconv import _unwrap_tables # noqa: E402
class TestMailConv(unittest.TestCase):
def test_clean_email_address(self):
self.assertEqual(clean_email_address('"Ben Baker" <bb@teks.com>'), "bb@teks.com")
self.assertEqual(clean_email_address("eslider@gmail.com"), "eslider@gmail.com")
self.assertEqual(clean_email_address("<a@b.c>"), "a@b.c")
def test_subject_to_filename(self):
self.assertEqual(subject_to_filename("Your receipt #2422"), "Your_receipt_2422")
self.assertEqual(subject_to_filename("a/b\\c:d*e"), "abcde")
self.assertEqual(subject_to_filename(" "), "untitled")
def test_html_to_markdown(self):
out = html_to_markdown("<html><body><h1>Hi</h1><p>Some <b>bold</b> text.</p></body></html>")
self.assertIn("Hi", out)
self.assertIn("**bold**", out)
def test_html_strip_fallback(self):
from mailconv import strip_html
self.assertEqual(strip_html("<p>a</p><p>b</p>"), "a\n\nb")
def test_flatten_layout_tables(self):
html = ("<table><tr>"
+ "".join(f"<td>spacer{i}</td>" for i in range(12))
+ "</tr></table>"
+ "<p>real</p>"
+ "<table><tr><td>a</td><td>b</td></tr></table>")
out = _unwrap_tables(html)
# tables unwrapped into pipe text; no <td> left; content preserved
self.assertNotIn("<td>spacer0</td>", out)
self.assertIn("spacer0 | spacer1", out)
self.assertIn("a | b", out)
self.assertIn("real", out)
def test_html_to_markdown_layout_clean(self):
html = "<table><tr>" + "".join(f"<td>x{i}</td>" for i in range(12)) + "</tr></table><h1>Hi</h1>"
out = html_to_markdown(html)
self.assertIn("Hi", out)
self.assertNotIn("| ---", out)
def test_normalize_markdown_removes_nul(self):
self.assertEqual(normalize_markdown("Z0\x00A\x00Y\x00B"), "Z0AYB")
self.assertEqual(normalize_markdown("a\n\n\n\nb"), "a\n\nb")
def test_normalize_strips_email_noise(self):
noisy = "\ufeffa\u200b\u034f\u00ad\u2007\u2002 b\u200a c\u2008"
out = normalize_markdown(noisy)
self.assertNotIn("\u200b", out)
self.assertNotIn("\ufeff", out)
self.assertNotIn("\u034f", out)
self.assertIn("a b c", out)
def test_split_zip_members(self):
p = Path(self._mk_zip(["a.txt", "sub/b.txt"]))
self.assertEqual(split_zip_members(p), ["a.txt", "sub/b.txt"])
def test_zip_extract_safe(self):
zip_path = self._mk_zip(["a.txt", "dir/b.txt"])
dest = Path(self._tmp("x"))
files = zip_extract_safe(zip_path, dest)
self.assertEqual(len(files), 2)
self.assertTrue((dest / "a.txt").exists())
self.assertTrue((dest / "dir" / "b.txt").exists())
def test_zip_extract_safe_blocks_traversal(self):
# member "../evil.txt" must not escape dest
zip_path = Path(self._tmp("evil.zip"))
with zipfile.ZipFile(zip_path, "w") as zf:
zf.writestr("../evil.txt", "boom")
dest = Path(self._tmp("out"))
files = zip_extract_safe(zip_path, dest)
self.assertEqual(files, [])
self.assertFalse((dest.parent / "evil.txt").exists())
def test_is_convertible(self):
self.assertTrue(is_convertible(".pdf"))
self.assertTrue(is_convertible(".zip"))
self.assertTrue(is_convertible(".docx"))
self.assertTrue(is_convertible(".TXT"))
self.assertFalse(is_convertible(".exe"))
self.assertFalse(is_convertible(".unknown"))
def test_convert_pdf_prefers_pdftotext(self):
import mailconv as mc
calls: list[list[str]] = []
def fake_run(cmd, **kwargs):
calls.append(list(cmd))
class P:
returncode = 0
stdout = b"Invoice BM25 layout"
stderr = b""
return P()
self._patch_run(mc, fake_run)
out = convert_pdf(Path(self._tmp("born.pdf")))
self.assertIn("BM25", out)
self.assertEqual(calls[0][:2], ["pdftotext", "-layout"])
self.assertFalse(any(c[0] == "tesseract" for c in calls))
self.assertFalse(any(c[0] == "pdftoppm" for c in calls))
def test_convert_pdf_empty_layer_uses_pdftoppm_tesseract(self):
import mailconv as mc
calls: list[list[str]] = []
def fake_run(cmd, **kwargs):
calls.append(list(cmd))
class P:
returncode = 0
stdout = b""
stderr = b""
if cmd[0] == "pdftotext":
P.stdout = b" \n"
return P()
if cmd[0] == "pdftoppm":
prefix = Path(cmd[-1])
(prefix.parent / "page-1.png").write_bytes(b"fake")
return P()
if cmd[0] == "tesseract":
P.stdout = b"scanned HELLO"
return P()
return P()
self._patch_run(mc, fake_run)
out = convert_pdf(Path(self._tmp("scan.pdf")))
self.assertIn("HELLO", out)
bins = [c[0] for c in calls]
self.assertIn("pdftotext", bins)
self.assertIn("pdftoppm", bins)
self.assertIn("tesseract", bins)
tess = next(c for c in calls if c[0] == "tesseract")
self.assertIn(TESS_LANG, tess)
self.assertNotIn("docling", " ".join(bins))
def test_ocr_image_paddle_engine(self):
import mailconv as mc
calls: list[list[str]] = []
def fake_run(cmd, **kwargs):
calls.append(list(cmd))
class P:
returncode = 0
stdout = b"paddle text"
stderr = b""
return P()
self._patch_run(mc, fake_run)
os.environ["OCR_ENGINE"] = "paddle"
try:
out = ocr_image(Path(self._tmp("x.png")))
finally:
os.environ.pop("OCR_ENGINE", None)
self.assertEqual(out, "paddle text")
self.assertEqual(calls[0][:2], ["paddleocr", "ocr"])
def _patch_run(self, mod, fn) -> None:
self.addCleanup(setattr, mod.subprocess, "run", mod.subprocess.run)
mod.subprocess.run = fn
def _mk_zip(self, members):
zpath = Path(self._tmp("arc.zip"))
with zipfile.ZipFile(zpath, "w") as zf:
for m in members:
zf.writestr(m, "content")
return str(zpath)
def _tmp(self, name):
d = self.__class__._td
p = Path(d) / name
p.parent.mkdir(parents=True, exist_ok=True)
return str(p)
@classmethod
def setUpClass(cls):
import tempfile
cls._td = tempfile.mkdtemp(prefix="mailconv_test_")
if __name__ == "__main__":
unittest.main()
+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"])
+188
View File
@@ -0,0 +1,188 @@
"""Published docs must match live commands (Gitea SoT, brain/search)."""
from __future__ import annotations
import unittest
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
class PublishedDocsTest(unittest.TestCase):
def test_readme_points_issues_at_gitea(self) -> None:
text = (ROOT / "README.md").read_text()
self.assertIn(
"https://git.produktor.io/eSlider/2dph/issues",
text,
"README must point issues at Gitea",
)
def test_plan_d15_names_gitea_origin(self) -> None:
text = (ROOT / "PLAN.md").read_text()
self.assertIn("D15", text)
self.assertIn("git.produktor.io/eSlider/2dph", text)
def test_readme_primary_search_is_brain(self) -> None:
text = (ROOT / "README.md").read_text()
self.assertIn(
"bin/brain/search.go",
text,
"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_readme_git_import_is_gogit(self) -> None:
text = (ROOT / "README.md").read_text()
self.assertIn("bin/git/import.go", text)
self.assertIn("go-git", text)
self.assertIn("D19", (ROOT / "PLAN.md").read_text())
def test_web_search_is_go_not_ops_host(self) -> None:
readme = (ROOT / "README.md").read_text()
self.assertIn("bin/web/search.go", readme)
skill = (ROOT / "skills" / "web-search" / "SKILL.md").read_text()
self.assertIn("bin/web/search.go", skill)
self.assertNotIn("search.ops.io", skill)
self.assertNotIn("search.ops.io", readme)
compose = (ROOT / "compose.yaml").read_text()
self.assertIn("searxng", compose)
self.assertNotIn("search.ops.io", compose)
settings = (ROOT / "deploy" / "searxng" / "settings.yml").read_text()
self.assertNotIn("password", settings.lower())
self.assertIn("json", settings)
def test_picoclaw_compose_profile_has_mcp_example(self) -> None:
compose = (ROOT / "compose.yaml").read_text()
self.assertIn('profiles: ["picoclaw"]', compose)
self.assertIn("127.0.0.1:8630", compose)
example = (ROOT / "deploy" / "picoclaw" / "mcp.json.example").read_text()
self.assertIn("127.0.0.1:8630/mcp", example)
self.assertNotIn("password", example.lower())
self.assertNotIn("token", example.lower())
docs = (ROOT / "docs" / "picoclaw.md").read_text()
self.assertIn("search", docs)
self.assertIn("throttled", docs)
def test_readme_read_path_is_go(self) -> None:
plan = (ROOT / "PLAN.md").read_text()
self.assertIn("get.go", plan)
self.assertIn("CI fallback", plan)
design = (ROOT / "docs" / "design.md").read_text()
self.assertIn("internal/brain/rank", design)
self.assertIn("They do not exec Python", design)
def test_openapi_mcp_from_same_handlers(self) -> None:
plan = (ROOT / "PLAN.md").read_text()
self.assertIn("D20", plan)
self.assertIn("/openapi.json", (ROOT / "README.md").read_text())
self.assertIn("/mcp", (ROOT / "README.md").read_text())
skill = (ROOT / "skills" / "brain" / "SKILL.md").read_text()
self.assertIn("/mcp", skill)
self.assertFalse((ROOT / "skills" / "db-yaml").exists())
self.assertTrue((ROOT / "skills" / "postgres" / "SKILL.md").is_file())
def test_cgo_zig_and_index_profile(self) -> None:
plan = (ROOT / "PLAN.md").read_text()
self.assertIn("D21", plan)
self.assertIn("zig cc", plan)
dockerfile = (ROOT / "Dockerfile").read_text()
self.assertIn("bin/cgo/zcc", dockerfile)
self.assertIn("FROM debian:bookworm-slim AS api", dockerfile)
self.assertIn("FROM python:3.12-slim AS index", dockerfile)
api = dockerfile[dockerfile.index("FROM debian:bookworm-slim AS api") :]
self.assertNotIn("pip install", api)
compose = (ROOT / "compose.yaml").read_text()
self.assertIn('profiles: ["index"]', compose)
self.assertIn("target: api", compose)
def test_reasoner_docs_name_real_hf_ids_cpu_sidecar(self) -> None:
docs = (ROOT / "docs" / "reasoner.md").read_text()
for hf in (
"Qwen/Qwen3.5-9B",
"Qwen/Qwen3.6-27B",
"prism-ml/Bonsai-27B-gguf",
):
self.assertIn(hf, docs)
self.assertIn("no official qwen3.6-9b", docs.lower())
self.assertIn("OLLAMA_NUM_GPU", docs)
self.assertIn("rss_mb", docs)
self.assertIn("vram_mb", docs)
self.assertIn("3/3", docs)
self.assertIn("Do not claim 9B is better at tools", docs)
self.assertNotIn("Qwen/Qwen3.6-9B", docs)
plan = (ROOT / "PLAN.md").read_text()
self.assertIn("D18", plan)
self.assertIn("Qwen/Qwen3.5-9B", plan)
compose = (ROOT / "compose.yaml").read_text()
self.assertIn('"reasoner"', compose)
self.assertIn("OLLAMA_NUM_GPU", compose)
self.assertIn("127.0.0.1:11435", compose)
dockerfile = (ROOT / "Dockerfile").read_text()
self.assertNotIn(".gguf", dockerfile.lower())
self.assertNotIn(".safetensors", dockerfile.lower())
api = dockerfile[dockerfile.index("FROM debian:bookworm-slim AS api") :]
self.assertNotIn("COPY models", api)
self.assertNotIn("qwen", api.lower())
def test_readme_search_escalates_web(self) -> None:
text = (ROOT / "README.md").read_text()
self.assertIn("--no-web", text)
self.assertIn("D17", (ROOT / "PLAN.md").read_text())
skill = (ROOT / "skills" / "brain" / "SKILL.md").read_text()
self.assertIn("`web` block", skill)
def test_docs_say_hop_walks_from_file(self) -> None:
paths = [
ROOT / "README.md",
ROOT / "docs" / "design.md",
ROOT / "skills" / "brain" / "SKILL.md",
ROOT / "docs" / "runbook.md",
ROOT / "docs" / "README.md",
]
for path in paths:
text = path.read_text()
self.assertIn("--hop", text, f"{path.relative_to(ROOT)} must document --hop")
self.assertNotIn(
"not implemented",
text.lower(),
f"{path.relative_to(ROOT)} still says hop is not implemented",
)
def test_docs_are_portable_diataxis(self) -> None:
index = (ROOT / "docs" / "README.md").read_text()
self.assertIn("type: reference", index)
for d in ("D3", "D6", "D14", "D15", "D17", "D18"):
self.assertIn(d, index)
runbook = (ROOT / "docs" / "runbook.md").read_text()
self.assertIn("type: howto", runbook)
self.assertIn("bin/brain/search.go", runbook)
self.assertIn("bin/brain/index.go", runbook)
self.assertNotIn("search.ops.io", runbook)
self.assertNotIn("/mnt/", runbook)
self.assertNotIn("/home/", runbook)
readme = (ROOT / "README.md").read_text()
self.assertIn("docs/runbook.md", readme)
self.assertNotIn("search.ops.io", readme)
def test_v1_epic_is_named_in_docs(self) -> None:
plan = (ROOT / "PLAN.md").read_text()
self.assertIn("Gap to v1", plan)
self.assertIn("eSlider/2dph/issues/16", plan)
self.assertIn("eSlider/2dph/issues/17", plan)
self.assertIn("eSlider/2dph/milestone/12", plan)
road = (ROOT / "docs" / "roadmap.md").read_text()
self.assertIn("type: explanation", road)
self.assertIn("issues/16", road)
self.assertIn("issues/14", road)
index = (ROOT / "docs" / "README.md").read_text()
self.assertIn("roadmap.md", index)
self.assertIn("epic #16", index)
agents = (ROOT / "AGENTS.md").read_text()
self.assertIn("roadmap.md", agents)
+63
View File
@@ -0,0 +1,63 @@
"""Skills must name live commands; every bin/ path in SKILL.md must exist."""
from __future__ import annotations
import re
import unittest
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
BIN_PATH = re.compile(r"(bin/[A-Za-z0-9_./-]+)")
class SkillsTest(unittest.TestCase):
def test_db_yaml_renamed_to_postgres(self) -> None:
self.assertFalse(
(ROOT / "skills" / "db-yaml").exists(),
"skills/db-yaml must be skills/postgres",
)
self.assertTrue((ROOT / "skills" / "postgres" / "SKILL.md").is_file())
text = (ROOT / "skills" / "postgres" / "SKILL.md").read_text()
self.assertIn("bin/postgres/query.go", text)
self.assertNotIn("search.ops.io", text)
def test_every_bin_path_in_skills_exists(self) -> None:
missing: list[str] = []
for path in (ROOT / "skills").rglob("SKILL.md"):
text = path.read_text()
for m in BIN_PATH.finditer(text):
rel = m.group(1).rstrip(")`.,;")
candidate = ROOT / rel
if not candidate.exists():
missing.append(f"{path.relative_to(ROOT)}: {rel}")
self.assertEqual(missing, [], "skill bin paths must exist")
def test_brain_skill_lists_generated_tools(self) -> None:
tools = (ROOT / "skills" / "brain" / "tools.md").read_text()
skill = (ROOT / "skills" / "brain" / "SKILL.md").read_text()
self.assertIn("tools.md", skill)
for name in ("search", "get", "stats", "audit"):
self.assertIn(f"`{name}`", tools)
def test_picoclaw_lists_tool_order(self) -> None:
skill = (ROOT / "skills" / "picoclaw" / "SKILL.md").read_text()
agents = (ROOT / "AGENTS.md").read_text()
self.assertIn("**`search`**", skill)
self.assertIn("**`get`**", skill)
self.assertIn("**`audit`**", skill)
self.assertIn("throttled", skill.lower())
self.assertIn("not a negative finding", agents)
self.assertIn("Fact-check every", agents)
def test_yq_is_mikefarah_for_structured_data(self) -> None:
skill = (ROOT / "skills" / "yq" / "SKILL.md").read_text()
self.assertIn("https://github.com/mikefarah/yq", skill)
for fmt in ("YAML", "JSON", "XML", "CSV", "TOML", "HCL"):
self.assertIn(fmt, skill)
self.assertIn("not kislyuk", skill.lower())
plan = (ROOT / "PLAN.md").read_text()
self.assertIn("mikefarah/yq", plan)
agents = (ROOT / "AGENTS.md").read_text()
self.assertIn("mikefarah/yq", agents)
web = (ROOT / "skills" / "web-search" / "SKILL.md").read_text()
self.assertIn("| yq ", web)
self.assertNotIn("| jq ", web)
+33
View File
@@ -0,0 +1,33 @@
"""Every bin/ path named in skills/ must exist on disk."""
from __future__ import annotations
import re
import unittest
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
BIN_PATH = re.compile(r"\b(bin/[A-Za-z0-9_./-]+)")
class SkillsBinPathsTest(unittest.TestCase):
def test_agent_cost_skill_is_gone(self) -> None:
self.assertFalse(
(ROOT / "skills" / "agent-cost").exists(),
"skills/agent-cost documents bin/agents/cost which does not exist",
)
def test_brain_skill_replaces_kb_search(self) -> None:
self.assertTrue((ROOT / "skills" / "brain" / "SKILL.md").is_file())
self.assertFalse((ROOT / "skills" / "kb-search").exists())
def test_skill_bin_paths_exist(self) -> None:
missing: list[str] = []
for skill in sorted((ROOT / "skills").rglob("SKILL.md")):
text = skill.read_text()
for match in BIN_PATH.findall(text):
rel = match.rstrip("`'.,")
if rel.endswith(".go") or Path(rel).suffix == "" or Path(rel).suffix in {".go", ".py"}:
p = ROOT / rel
if not p.exists():
missing.append(f"{skill.relative_to(ROOT)}: {rel}")
self.assertEqual(missing, [], "SKILL.md names bin/ paths that do not exist")
+36
View File
@@ -0,0 +1,36 @@
"""qa/system_perf.py is an offline-gated system test (no live brain in CI)."""
from __future__ import annotations
import ast
import unittest
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
class SystemPerfScriptTest(unittest.TestCase):
def test_script_compiles_and_is_read_only(self) -> None:
path = ROOT / "qa" / "system_perf.py"
src = path.read_text()
compile(src, str(path), "exec")
self.assertIn("--json", src)
self.assertIn("qwen3.5:9b", src)
self.assertIn("--picoclaw", src)
self.assertIn("BRAIN_URL", src)
self.assertIn("tools/list", src)
self.assertIn("tools/call", src)
self.assertIn("GATE_HEALTH_MS", src)
self.assertIn("GATE_GET_P50_MS", src)
self.assertNotIn("kb.lbug", src)
self.assertNotIn("password", src.lower())
self.assertNotIn("token", src.lower())
def test_script_does_not_write_ladybug(self) -> None:
tree = ast.parse((ROOT / "qa" / "system_perf.py").read_text())
writes = [
n.func.attr
for n in ast.walk(tree)
if isinstance(n, ast.Call) and isinstance(n.func, ast.Attribute)
and n.func.attr in {"write_text", "write_bytes", "dump"}
]
self.assertEqual(writes, [], f"system_perf must not write files: {writes}")
View File
@@ -76,7 +76,7 @@ class PhiGuard(unittest.TestCase):
def test_plain_technical_query_passes(self): def test_plain_technical_query_passes(self):
self.assertIsNone(ws.phi_reason("Pflegegrad SGB XI Einstufung")) self.assertIsNone(ws.phi_reason("Pflegegrad SGB XI Einstufung"))
self.assertIsNone(ws.phi_reason("site:ticket.detective.de Toureffizienz")) self.assertIsNone(ws.phi_reason("site:example.com technical query"))
def test_long_digit_run_is_refused(self): def test_long_digit_run_is_refused(self):
self.assertIsNotNone(ws.phi_reason("Kunde 4711220385 Adresse")) self.assertIsNotNone(ws.phi_reason("Kunde 4711220385 Adresse"))
+105
View File
@@ -0,0 +1,105 @@
// 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
// package. Polls file mtimes (no inotify deps); cheap and reliable.
package watch
import (
"log"
"os"
"os/exec"
"path/filepath"
"strconv"
"strings"
"time"
)
// Options controls the polling loop. Zero value uses defaults.
type Options struct {
Dirs []string
Interval time.Duration
// IndexCmd is the index command template. %s is replaced by the repo
// root (from KB_ROOT). Defaults to `python3 <root>/bin/kb/index --with-mail`.
IndexCmd string
}
// Run blocks forever polling Dirs (defaults: KB_WATCH_DIRS or /corpus) every
// Interval (default 30s) and re-indexing when files change. KB_ROOT names the
// repo root used to locate bin/kb/index.
func Run(args []string) {
opts := fromEnv(args)
root, _ := os.Getwd()
if r := os.Getenv("KB_ROOT"); r != "" {
root = r
}
log.Printf("watch: dirs=%v interval=%s root=%s", opts.Dirs, opts.Interval, root)
var last string
for {
if flag := Stamp(opts.Dirs); flag != "" && flag != last {
last = flag
reindex(opts.IndexCmd, root)
}
time.Sleep(opts.Interval)
}
}
func fromEnv(args []string) Options {
opts := Options{Interval: 30 * time.Second}
if raw := os.Getenv("KB_WATCH_INTERVAL"); raw != "" {
if n, err := strconv.Atoi(raw); err == nil && n > 0 {
opts.Interval = time.Duration(n) * time.Second
}
}
defDirs := "/corpus"
if raw := os.Getenv("KB_WATCH_DIRS"); raw != "" {
defDirs = raw
}
if len(args) > 0 {
opts.Dirs = args
} else {
for _, d := range strings.Split(defDirs, " ") {
if d != "" {
opts.Dirs = append(opts.Dirs, d)
}
}
}
pys := os.Getenv("KB_PY")
if pys == "" {
pys = "python3"
}
opts.IndexCmd = pys + " <root>/bin/kb/index --with-mail"
return opts
}
// Stamp returns a rolling fingerprint (newest mtime under dirs) that changes
// whenever any corpus file is touched. Empty when no files found.
func Stamp(dirs []string) string {
var newest time.Time
for _, dir := range dirs {
_ = filepath.WalkDir(dir, func(path string, _ os.DirEntry, err error) error {
if err != nil {
return nil
}
if info, e := os.Stat(path); e == nil && info.ModTime().After(newest) {
newest = info.ModTime()
}
return nil
})
}
if newest.IsZero() {
return ""
}
return strconv.FormatInt(newest.UnixNano(), 10)
}
func reindex(template, root string) {
cmd := strings.ReplaceAll(template, "<root>", root)
parts := strings.Fields(cmd)
c := exec.Command(parts[0], parts[1:]...)
out, err := c.CombinedOutput()
if err != nil {
log.Printf("watch: index failed: %v\n%s", err, out)
} else {
log.Printf("watch: re-indexed")
}
}
+53
View File
@@ -0,0 +1,53 @@
package watch
import (
"os"
"path/filepath"
"strings"
"testing"
"time"
)
func TestStampChangesWhenFileTouched(t *testing.T) {
dir := t.TempDir()
a := filepath.Join(dir, "a.md")
if err := os.WriteFile(a, []byte("x"), 0o644); err != nil {
t.Fatal(err)
}
s1 := Stamp([]string{dir})
if s1 == "" {
t.Fatal("stamp empty for a dir with a file")
}
time.Sleep(10 * time.Millisecond)
if err := os.WriteFile(a, []byte("y"), 0o644); err != nil {
t.Fatal(err)
}
if s2 := Stamp([]string{dir}); s2 == s1 {
t.Fatal("stamp did not change after the file was modified")
}
}
func TestStampEmptyForMissingDir(t *testing.T) {
if s := Stamp([]string{filepath.Join(t.TempDir(), "nope")}); s != "" {
t.Fatalf("stamp = %q, want empty for missing dir", s)
}
}
func TestFromEnvDefaults(t *testing.T) {
t.Setenv("KB_WATCH_INTERVAL", "")
t.Setenv("KB_WATCH_DIRS", "")
t.Setenv("KB_PY", "")
opts := fromEnv(nil)
if len(opts.Dirs) == 0 || opts.Dirs[0] != "/corpus" {
t.Fatalf("default dirs = %v, want [/corpus]", opts.Dirs)
}
if opts.Interval != 30*time.Second {
t.Fatalf("default interval = %s, want 30s", opts.Interval)
}
if !strings.Contains(opts.IndexCmd, "kb/index") {
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)
}
}

Some files were not shown because too many files have changed in this diff Show More