Compare commits

...
16 Commits
Author SHA1 Message Date
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
89 changed files with 9669 additions and 208 deletions
-1
View File
@@ -10,5 +10,4 @@ __pycache__
.cache .cache
.secrets .secrets
.skills-tmp .skills-tmp
serve/serve
docs/.build docs/.build
+2 -4
View File
@@ -31,18 +31,16 @@ 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
- 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 (server + watch packages)
run: | run: |
go vet ./... go vet ./...
go test ./... -count=1 go test ./... -count=1
working-directory: serve
- name: facts/audit self (lexicon consistency, no network) - name: facts/audit self (lexicon consistency, no network)
run: | run: |
+1
View File
@@ -9,3 +9,4 @@ __pycache__/
*.env *.env
.env .env
.secrets/ .secrets/
lib-ladybug/
+42 -5
View File
@@ -36,20 +36,44 @@ 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} (shebang)
bin/kb-watch corpus watcher (mtimes, no inotify deps) bin/serve.go async Go HTTP server entry (self-executing go run shebang)
bin/watch/ corpus watcher Go package (mtimes, no inotify deps)
bin/server/ async Go HTTP server (goroutines, bounded worker pool)
bin/mail/ mail pipeline: sync (Go), import (md), index_mail (rebuild)
bin/tools/ vendored python libs behind bin/* (kblib, yamlout, websearch)
bin/docker-entrypoint container entrypoint (brain index|search|serve|watch) bin/docker-entrypoint container entrypoint (brain index|search|serve|watch)
serve/ async Go HTTP server (goroutines, bounded worker pool)
tools/ vendored python libs behind bin/* (yamlout, websearch)
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 multi-stage: python deps + static Go binaries
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 --from-raw var/mail # message.json → message.md (convert only)
bin/mail/index_mail # rebuild brain incl. all mail (fresh DB)
```
- `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 fall back to
docling (isolated subprocess — its native onnx can segfault the parent).
Conversion never touches the brain DB (crash safety).
- `index_mail` always rebuilds from scratch (repo corpus + mail). Ladybug
corrupts its WAL when brand-new leafs are bulk-inserted while FTS/vector
indexes exist; a fresh DB with indexes created last is the only safe path.
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 ["self"|"facts"|"info"|"stale"] # 2-source + staleness gate
bin/facts/crm [--dry-run] # proof person↔company/company↔project (ooCRM × corpus SoT)
bin/kb/search "query" [--hop N] [--repo X] # deduction search → YAML bin/kb/search "query" [--hop N] [--repo X] # deduction search → YAML
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
@@ -58,6 +82,19 @@ 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.
## 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
Same tone as the corpus: plain, lists, no hype. Sign-off `Andriy Oblivantsev`. Same tone as the corpus: plain, lists, no hype. Sign-off `Andriy Oblivantsev`.
+14 -10
View File
@@ -14,23 +14,27 @@ 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 # Go services: static binaries, no interpreter at runtime
FROM golang:1.25 AS serve-build FROM golang:1.25 AS go-build
WORKDIR /src/serve WORKDIR /src
COPY serve/go.mod serve/go.sum* ./ COPY go.mod ./
COPY serve . COPY bin/server ./bin/server
RUN CGO_ENABLED=0 go build -o /serve -ldflags="-s -w" . COPY bin/watch ./bin/watch
RUN CGO_ENABLED=0 go build -o /serve ./bin/server \
&& CGO_ENABLED=0 go build -o /watch ./bin/watch
# runtime: python toolchain + Go server # runtime: python toolchain + Go services
FROM base FROM base
COPY . . COPY . .
COPY --from=serve-build /serve /app/serve/serve COPY --from=go-build /serve /app/bin/serve
RUN chmod +x /app/bin/kb-watch /app/bin/docker-entrypoint \ COPY --from=go-build /watch /app/bin/watch
RUN chmod +x /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
+18 -2
View File
@@ -25,7 +25,7 @@ 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 | import `web-search`, retire local `searxng-ops`. Vendored here, no remote link. |
| 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). |
@@ -93,8 +93,24 @@ Common props on every node/edge: `root`, `confidence`, `evidence[]`, `how`,
- OQ1: mutually-contradicting evidence — how to resolve (authority weighting, - OQ1: mutually-contradicting evidence — how to resolve (authority weighting,
temporal freshness, audit adjudication). temporal freshness, audit adjudication).
- OQ2: OCR pipeline for pdfs/images/docs (late phase). - OQ2: OCR pipeline for pdfs/images/docs — mostly solved: poppler pdftotext
fast-path for born-digital PDFs, docling fallback for the ~5% textless ones.
- OQ3: optional duckdb-md layer for `SELECT … FORMAT MARKDOWN` export/write-back. - OQ3: optional duckdb-md layer for `SELECT … FORMAT MARKDOWN` export/write-back.
- 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 --from-raw` — message.json → message.md; PDFs via
`pdftotext -layout` (~15ms) with docling subprocess fallback; ICS sidecars
Latin-1→UTF-8 normalized.
3. `bin/mail/index_mail` — 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.
4. Result: 17,835 messages → 28,918 info leafs, FTS + HNSW healthy, searchable
via `bin/kb/search`.
## CI/CD pipeline (D15) ## CI/CD pipeline (D15)
+15 -3
View File
@@ -93,11 +93,23 @@ bin/kb/stats # index health
bin/kb/eval # recall@5 gate bin/kb/eval # recall@5 gate
``` ```
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 --from-raw var/mail # JSON → markdown
bin/mail/index_mail # rebuild brain incl. mail
bin/kb/search "Mietwagen Nürnberg invoice" # now answers from mail
```
## Storage ## Storage
- **LadybugDB** — single `var/kb.lbug`, Cypher property graph, HNSW + BM25 - **LadybugDB** — single `var/kb.lbug`, Cypher property graph, HNSW + BM25
in one engine, embedded (no server), ACID, read-only-safe for concurrent in one engine, embedded (no server), ACID, read-only-safe for concurrent
readers. readers. **Never `DROP INDEX` FTS/VECTOR** on Ladybug 0.19: DROP leaves
ghost catalog tables (`_0_Leaf_vec_UPPER`) so recreate fails while
`SHOW_INDEXES` omits HNSW. Fresh indexes = delete `var/kb.lbug` +
`bin/kb/index --rebuild`. Use `ensure_indexes()` after upserts.
- **model2vec** — `potion-multilingual-128M` static embeddings (256-dim), - **model2vec** — `potion-multilingual-128M` static embeddings (256-dim),
CPU-fast, deterministic, no Ollama runtime dependency. CPU-fast, deterministic, no Ollama runtime dependency.
- facts and info split semantically by `root` column but written inside the - facts and info split semantically by `root` column but written inside the
@@ -116,7 +128,7 @@ touches network/db is read-only, throttled, cached. Tests gate every commit.
uv venv .venv # Python 3.12, uv-managed uv venv .venv # Python 3.12, uv-managed
uv pip install -r requirements.lock.txt # pinned toolchain uv pip install -r requirements.lock.txt # pinned toolchain
bin/facts/audit self # lexicon consistency gate bin/facts/audit self # lexicon consistency gate
go test ./... && python -m unittest discover -s tools -t . go test ./... && python -m unittest discover -s bin/tools -t .
``` ```
Docker (optional, cached model + var volumes): Docker (optional, cached model + var volumes):
@@ -134,6 +146,6 @@ docker compose up brain-watch # auto re-index on change
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`, `db-yaml`, …) 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. See [PLAN.md](PLAN.md) for decisions, execution status, and v2 open questions.
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" "$@"
+317
View File
@@ -0,0 +1,317 @@
package main
import (
"bytes"
"encoding/json"
"flag"
"fmt"
"os"
"os/exec"
"path/filepath"
"strings"
)
type ooContact struct {
ID int `json:"id"`
DisplayName string `json:"displayName"`
FirstName string `json:"firstName"`
LastName string `json:"lastName"`
About string `json:"about"`
CommonData []struct {
InfoType int `json:"infoType"`
Data string `json:"data"`
Category string `json:"categoryName"`
} `json:"commonData"`
}
func runApply(args []string) int {
fs := flag.NewFlagSet("chats apply", flag.ContinueOnError)
dryRun := fs.Bool("dry-run", false, "show what would be done without writing")
help := fs.Bool("help", false, "")
fs.SetOutput(os.Stderr)
if err := fs.Parse(args); err != nil {
return 2
}
if *help {
fmt.Fprintln(os.Stderr, "usage: chats apply [--dry-run]")
return 0
}
ooCLI := findOO()
if ooCLI == "" {
fmt.Fprintln(os.Stderr, "chats apply: oo CLI not found; set OO_CLI or install go-onlyoffice")
return 1
}
facts, err := loadFacts()
if err != nil {
fmt.Fprintf(os.Stderr, "chats apply: %v\n", err)
return 1
}
if len(facts) == 0 {
fmt.Println("chats apply: no facts to process")
return 0
}
phoneFacts := filterFacts(facts, "phone")
emailFacts := filterFacts(facts, "email")
phoneFacts = dedupeFacts(phoneFacts)
emailFacts = dedupeFacts(emailFacts)
type resolvedFact struct {
Fact ExtractedFact
OoID int
OoName string
Action string // "info-add" or "persons-create"
}
var resolved []resolvedFact
for _, f := range phoneFacts {
contact, err := searchContact(ooCLI, f.ChatName)
if err != nil || contact == nil {
fmt.Printf(" ✗ %s: phone %s — not found in CRM\n", f.ChatName, f.Value)
resolved = append(resolved, resolvedFact{Fact: f, Action: "persons-create"})
continue
}
hasPhone := false
for _, d := range contact.CommonData {
if d.InfoType == 2 {
hasPhone = true
break
}
}
if hasPhone {
fmt.Printf(" ✓ %s (ID %d): phone %s — already has phone, skip\n", contact.DisplayName, contact.ID, f.Value)
continue
}
fmt.Printf(" → %s (ID %d): add phone %s\n", contact.DisplayName, contact.ID, f.Value)
resolved = append(resolved, resolvedFact{
Fact: f, OoID: contact.ID, OoName: contact.DisplayName, Action: "info-add",
})
}
for _, f := range emailFacts {
if strings.EqualFold(f.Value, envVar("ONLYOFFICE_USER", "")) ||
strings.EqualFold(f.Value, envVar("OO_USER", "")) ||
strings.EqualFold(f.Value, os.Getenv("EMAIL")) {
continue
}
contact, err := searchContact(ooCLI, f.ChatName)
if err != nil || contact == nil {
fmt.Printf(" ✗ %s: email %s — not found in CRM\n", f.ChatName, f.Value)
resolved = append(resolved, resolvedFact{Fact: f, Action: "persons-create"})
continue
}
hasEmail := false
for _, d := range contact.CommonData {
if d.InfoType == 1 && d.Data == f.Value {
hasEmail = true
break
}
}
if hasEmail {
fmt.Printf(" ✓ %s (ID %d): email %s — already exists\n", contact.DisplayName, contact.ID, f.Value)
continue
}
fmt.Printf(" → %s (ID %d): add email %s\n", contact.DisplayName, contact.ID, f.Value)
resolved = append(resolved, resolvedFact{
Fact: f, OoID: contact.ID, OoName: contact.DisplayName, Action: "info-add",
})
}
if len(resolved) == 0 {
fmt.Println("chats apply: nothing to apply")
return 0
}
fmt.Printf("\nchats apply: %d actions to apply\n", len(resolved))
if *dryRun {
for _, r := range resolved {
switch r.Action {
case "info-add":
infoType := "Phone"
if r.Fact.FactType == "email" {
infoType = "Email"
}
fmt.Printf(" [dry-run] oo contacts info-add %d --type %s --value %s\n",
r.OoID, infoType, r.Fact.Value)
case "persons-create":
fmt.Printf(" [dry-run] oo persons create --first %q --about %q\n",
r.Fact.ChatName, "Contact from Telegram chat")
}
}
return 0
}
success := 0
failed := 0
for _, r := range resolved {
switch r.Action {
case "info-add":
infoType := "Phone"
if r.Fact.FactType == "email" {
infoType = "Email"
}
if err := ooInfoAdd(ooCLI, r.OoID, infoType, r.Fact.Value); err != nil {
fmt.Fprintf(os.Stderr, " ✗ info-add %s: %v\n", r.Fact.Value, err)
failed++
} else {
fmt.Printf(" ✓ %s → %s (ID %d)\n", r.Fact.Value, r.OoName, r.OoID)
success++
}
case "persons-create":
fmt.Printf(" - create %s (skipped — needs review)\n", r.Fact.ChatName)
success++
}
}
fmt.Printf("\nchats apply: %d succeeded, %d failed\n", success, failed)
if failed > 0 {
return 1
}
return 0
}
func loadFacts() ([]ExtractedFact, error) {
factsPath := filepath.Join(chatsDir(), "facts", "chat-facts.json")
data, err := os.ReadFile(factsPath)
if err != nil {
if os.IsNotExist(err) {
return nil, fmt.Errorf("no facts at %s; run 'chats facts' first", factsPath)
}
return nil, fmt.Errorf("read facts: %w", err)
}
var facts []ExtractedFact
if err := json.Unmarshal(data, &facts); err != nil {
return nil, fmt.Errorf("parse facts: %w", err)
}
return facts, nil
}
func dedupeFacts(facts []ExtractedFact) []ExtractedFact {
seen := make(map[string]bool)
var result []ExtractedFact
for _, f := range facts {
norm := normalizePhone(f.Value)
key := f.ChatName + ":" + factTypeKey(f.FactType) + ":" + norm
if seen[key] {
continue
}
seen[key] = true
f.Value = norm
result = append(result, f)
}
return result
}
func normalizePhone(s string) string {
var digits []rune
for _, r := range s {
if r >= '0' && r <= '9' {
digits = append(digits, r)
}
}
if len(digits) > 0 {
return string(digits)
}
return s
}
func factTypeKey(t string) string {
switch t {
case "phone":
return "p"
case "email":
return "e"
default:
return t
}
}
func findOO() string {
if v := os.Getenv("OO_CLI"); v != "" {
if _, err := os.Stat(v); err == nil {
return v
}
}
candidates := []string{
filepath.Join(os.Getenv("HOME"), "go", "bin", "oo"),
}
for _, c := range candidates {
if _, err := os.Stat(c); err == nil {
return c
}
}
return ""
}
func searchContact(ooCLI, name string) (*ooContact, error) {
query := name
// Try full name first
if c, _ := searchByQuery(ooCLI, query); c != nil {
return c, nil
}
// Try first word
firstWord := strings.Fields(name)[0]
if firstWord != name {
if c, _ := searchByQuery(ooCLI, firstWord); c != nil {
return c, nil
}
}
return nil, nil
}
func searchByQuery(ooCLI, query string) (*ooContact, error) {
cmd := exec.Command(ooCLI, "persons", "list", "--search", query, "-o", "json")
var outBuf, errBuf bytes.Buffer
cmd.Stdout = &outBuf
cmd.Stderr = &errBuf
cmd.Env = os.Environ()
if err := cmd.Run(); err != nil {
return nil, fmt.Errorf("oo persons list: %w\n%s", err, errBuf.String())
}
var contacts []ooContact
if err := json.Unmarshal(outBuf.Bytes(), &contacts); err != nil {
return nil, nil
}
for _, c := range contacts {
lower := strings.ToLower(c.DisplayName)
lowerQuery := strings.ToLower(query)
if strings.EqualFold(c.DisplayName, query) ||
strings.Contains(lower, lowerQuery) ||
strings.Contains(lowerQuery, strings.ToLower(c.FirstName)) {
return &c, nil
}
for _, d := range c.CommonData {
if d.InfoType == 1 && strings.Contains(strings.ToLower(d.Data), lowerQuery) {
return &c, nil
}
}
}
if len(contacts) > 0 {
return &contacts[0], nil
}
return nil, nil
}
func ooInfoAdd(ooCLI string, contactID int, infoType, value string) error {
cmd := exec.Command(ooCLI, "contacts", "info-add",
fmt.Sprintf("%d", contactID),
"--type", infoType,
"--value", value,
"--category", "Work",
"-o", "json",
)
var outBuf, errBuf bytes.Buffer
cmd.Stdout = &outBuf
cmd.Stderr = &errBuf
cmd.Env = os.Environ()
if err := cmd.Run(); err != nil {
return fmt.Errorf("info-add: %w\n%s", err, errBuf.String())
}
return nil
}
+202
View File
@@ -0,0 +1,202 @@
// System tests for bin/chats.
//
// These are integration tests using real data and real Telegram API (when
// credentials are available). They follow the TDD workflow pattern:
// sync → import → facts → verify.
package main
import (
"encoding/json"
"os"
"path/filepath"
"strings"
"testing"
)
// TestChatsImport validates JSONL → MD conversion with a synthetic fixture.
func TestChatsImport(t *testing.T) {
dir := t.TempDir()
root := filepath.Join(dir, "var", "chats")
chatDir := filepath.Join(root, "telegram", "test_user_123")
if err := os.MkdirAll(chatDir, 0755); err != nil {
t.Fatal(err)
}
jsonlPath := filepath.Join(chatDir, "messages.jsonl")
f, err := os.Create(jsonlPath)
if err != nil {
t.Fatal(err)
}
defer f.Close()
enc := json.NewEncoder(f)
messages := []Message{
{ID: "tg_1", Timestamp: "2026-01-15T10:30:00Z", From: "Alice", Text: "Hello!", Platform: "telegram"},
{ID: "tg_2", Timestamp: "2026-01-15T10:31:00Z", From: "Bob", Text: "Hi Alice, my phone is +34 612 345 678", Platform: "telegram"},
{ID: "tg_3", Timestamp: "2026-01-15T10:32:00Z", From: "Alice", Text: "Check my LinkedIn: https://linkedin.com/in/alice-test", Platform: "telegram"},
{ID: "tg_4", Timestamp: "2026-01-15T10:33:00Z", From: "Bob", Text: "My email is bob@example.com, working on Project X", Platform: "telegram"},
}
for _, m := range messages {
if err := enc.Encode(m); err != nil {
t.Fatal(err)
}
}
f.Close()
cwd, _ := os.Getwd()
os.Chdir(dir)
t.Cleanup(func() { os.Chdir(cwd) })
t.Setenv("KB_ROOT", dir)
exitCode := runImport([]string{})
if exitCode != 0 {
t.Fatalf("import exit code %d", exitCode)
}
mdGlob := filepath.Join(root, "md", "telegram", "*", "messages.md")
matches, err := filepath.Glob(mdGlob)
if err != nil {
t.Fatal(err)
}
if len(matches) == 0 {
t.Fatal("no markdown files created by import")
}
mdData, err := os.ReadFile(matches[0])
if err != nil {
t.Fatal(err)
}
content := string(mdData)
if !strings.Contains(content, "Alice") {
t.Error("markdown missing sender name 'Alice'")
}
if !strings.Contains(content, "2026-01-15") {
t.Error("markdown missing date")
}
if !strings.Contains(content, "---") {
t.Error("markdown missing YAML frontmatter")
}
}
// TestChatsFacts validates fact extraction from JSONL fixture.
func TestChatsFacts(t *testing.T) {
dir := t.TempDir()
root := filepath.Join(dir, "var", "chats")
chatDir := filepath.Join(root, "telegram", "test_user_facts")
if err := os.MkdirAll(chatDir, 0755); err != nil {
t.Fatal(err)
}
jsonlPath := filepath.Join(chatDir, "messages.jsonl")
f, err := os.Create(jsonlPath)
if err != nil {
t.Fatal(err)
}
defer f.Close()
enc := json.NewEncoder(f)
messages := []Message{
{ID: "tg_10", Timestamp: "2026-06-01T12:00:00Z", From: "Charlie", Text: "Call me at +1 555 123 4567", Platform: "telegram"},
{ID: "tg_11", Timestamp: "2026-06-01T12:01:00Z", From: "Charlie", Text: "My LinkedIn is linkedin.com/in/charlie-dev", Platform: "telegram"},
{ID: "tg_12", Timestamp: "2026-06-01T12:02:00Z", From: "Charlie", Text: "Email: charlie@dev.com", Platform: "telegram"},
{ID: "tg_13", Timestamp: "2026-06-01T12:03:00Z", From: "Charlie", Text: "I work at Acme Corp on Project Mercury", Platform: "telegram"},
}
for _, m := range messages {
if err := enc.Encode(m); err != nil {
t.Fatal(err)
}
}
f.Close()
facts, _ := extractFacts(jsonlPath, "test_user_facts")
if len(facts) == 0 {
t.Fatal("expected facts, got none")
}
types := make(map[string]int)
for _, f := range facts {
types[f.FactType]++
}
if types["phone"] < 1 {
t.Errorf("expected >=1 phone fact, got %d", types["phone"])
}
if types["email"] < 1 {
t.Errorf("expected >=1 email fact, got %d", types["email"])
}
if types["linkedin"] < 1 {
t.Errorf("expected >=1 linkedin fact, got %d", types["linkedin"])
}
if types["skill"] < 1 {
t.Errorf("expected >=1 skill fact, got %d", types["skill"])
}
}
// TestChatsImportEmptyDir tests that import handles no JSONL gracefully.
func TestChatsImportEmpty(t *testing.T) {
dir := t.TempDir()
cwd, _ := os.Getwd()
os.Chdir(dir)
t.Cleanup(func() { os.Chdir(cwd) })
t.Setenv("KB_ROOT", dir)
exitCode := runImport([]string{})
if exitCode == 0 {
t.Fatal("expected non-zero exit for empty data dir")
}
}
// TestChatsRoundTrip creates a synthetic JSONL, imports it, then verifies
// the markdown structure is parseable and contains YAML frontmatter.
func TestChatsRoundTrip(t *testing.T) {
dir := t.TempDir()
root := filepath.Join(dir, "var", "chats")
chatDir := filepath.Join(root, "telegram", "rt_user")
if err := os.MkdirAll(chatDir, 0755); err != nil {
t.Fatal(err)
}
jsonlPath := filepath.Join(chatDir, "messages.jsonl")
f, err := os.Create(jsonlPath)
if err != nil {
t.Fatal(err)
}
enc := json.NewEncoder(f)
enc.Encode(Message{ID: "tg_100", Timestamp: "2026-07-01T08:00:00Z", From: "Diana", Text: "Hey", Platform: "telegram"})
enc.Encode(Message{ID: "tg_101", Timestamp: "2026-07-01T08:01:00Z", From: "Diana", Text: "How are you?", Platform: "telegram"})
f.Close()
cwd, _ := os.Getwd()
os.Chdir(dir)
t.Cleanup(func() { os.Chdir(cwd) })
t.Setenv("KB_ROOT", dir)
if code := runImport([]string{}); code != 0 {
t.Fatalf("import exit %d", code)
}
mdGlob := filepath.Join(root, "md", "telegram", "*", "messages.md")
matches, _ := filepath.Glob(mdGlob)
if len(matches) == 0 {
t.Fatal("no markdown produced")
}
data, err := os.ReadFile(matches[0])
if err != nil {
t.Fatal(err)
}
content := string(data)
if !strings.HasPrefix(content, "---") {
t.Error("markdown should start with YAML frontmatter delimiter")
}
if !strings.Contains(content, "platform: telegram") {
t.Error("markdown should contain platform field")
}
if !strings.Contains(content, "message_count: 2") {
t.Error("markdown should contain correct message count")
}
if !strings.Contains(content, "Diana") {
t.Error("markdown should contain participants")
}
}
+316
View File
@@ -0,0 +1,316 @@
package main
import (
"bufio"
"bytes"
"encoding/json"
"flag"
"fmt"
"os"
"os/exec"
"path/filepath"
"regexp"
"strings"
)
var (
phoneRegex = regexp.MustCompile(`[+\d][\d\s\-()]{6,25}\d`)
dateRegex = regexp.MustCompile(`^\d{2,4}[-/]\d{1,2}[-/]\d{2,4}$`)
rangeRegex = regexp.MustCompile(`^\d+\s*[-]\s*\d+$`)
linkedinRegex = regexp.MustCompile(`linkedin\.com/in/[\w-]+`)
emailRegex = regexp.MustCompile(`[\w.+-]+@[\w-]+\.[\w.-]+`)
projectRegex = regexp.MustCompile(`(?i)project\s*[:/]\s*(.+)`)
dealRegex = regexp.MustCompile(`(?i)(deal|opportunity)\s*[:/]\s*(.+)`)
skillRegex = regexp.MustCompile(`(?i)(works?|worked|working)\s+(at|on|with)\s+([A-Z][\w\s]+)`)
)
func isValidPhone(s string) bool {
s = strings.TrimSpace(s)
s = strings.Trim(s, "+()-\t ")
if len(s) < 6 || len(s) > 25 {
return false
}
if dateRegex.MatchString(s) || rangeRegex.MatchString(s) {
return false
}
if strings.ContainsAny(s, "/abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ") {
return false
}
if strings.Contains(s, "000") || strings.Contains(s, "500 ") || strings.Contains(s, "000 ") {
return false
}
digits := 0
for _, r := range s {
if r >= '0' && r <= '9' {
digits++
}
}
if digits < 7 || digits > 15 {
return false
}
// Card number pattern: 16 digits with possible spaces
if digits == 16 {
return false
}
// Date-like: 8 digits starting with 20xx or 19xx
if len(s) <= 8 && digits == 8 && (strings.HasPrefix(s, "20") || strings.HasPrefix(s, "19")) {
return false
}
// 11+ digits starting with 2 - unlikely phone
if digits >= 11 && strings.HasPrefix(s, "2") && !strings.HasPrefix(s, "+") {
return false
}
// Must start with + or be at least 7 digits
if !strings.HasPrefix(s, "+") && digits < 7 {
return false
}
return true
}
type ExtractedFact struct {
ChatID string `json:"chat_id"`
ChatName string `json:"chat_name"`
Platform string `json:"platform"`
FactType string `json:"fact_type"`
Value string `json:"value"`
Source string `json:"source"`
MessageID string `json:"message_id"`
}
func runFacts(args []string) int {
fs := flag.NewFlagSet("chats facts", flag.ContinueOnError)
help := fs.Bool("help", false, "")
fs.SetOutput(os.Stderr)
if err := fs.Parse(args); err != nil {
return 2
}
if *help {
fmt.Fprintln(os.Stderr, "usage: chats facts")
return 0
}
root := chatsDir()
telegramDir := filepath.Join(root, "telegram")
entries, err := os.ReadDir(telegramDir)
if err != nil {
fmt.Fprintf(os.Stderr, "chats facts: read %s: %v\n", telegramDir, err)
return 1
}
var allFacts []ExtractedFact
for _, entry := range entries {
if !entry.IsDir() {
continue
}
chatID := entry.Name()
jsonlPath := filepath.Join(telegramDir, chatID, "messages.jsonl")
info, err := os.Stat(jsonlPath)
if err != nil {
continue
}
if info.Size() == 0 {
continue
}
facts, chatName := extractFacts(jsonlPath, chatID)
allFacts = append(allFacts, facts...)
_ = chatName
}
if len(allFacts) == 0 {
fmt.Println("chats facts: no facts extracted")
return 0
}
phoneFacts := filterFacts(allFacts, "phone")
emailFacts := filterFacts(allFacts, "email")
linkedinFacts := filterFacts(allFacts, "linkedin")
projectFacts := filterFacts(allFacts, "project")
skillFacts := filterFacts(allFacts, "skill")
fmt.Printf("chats facts: extracted %d facts (%d phone, %d email, %d linkedin, %d project, %d skill)\n",
len(allFacts), len(phoneFacts), len(emailFacts), len(linkedinFacts), len(projectFacts), len(skillFacts))
factsDir := filepath.Join(root, "facts")
if err := os.MkdirAll(factsDir, 0755); err != nil {
fmt.Fprintf(os.Stderr, "chats facts: mkdir %s: %v\n", factsDir, err)
return 1
}
factsPath := filepath.Join(factsDir, "chat-facts.json")
data, err := json.MarshalIndent(allFacts, "", " ")
if err != nil {
fmt.Fprintf(os.Stderr, "chats facts: marshal: %v\n", err)
return 1
}
if err := os.WriteFile(factsPath, data, 0644); err != nil {
fmt.Fprintf(os.Stderr, "chats facts: write %s: %v\n", factsPath, err)
return 1
}
fmt.Printf("chats facts: saved to %s\n", factsPath)
writeFactsToBrain(root, allFacts)
return 0
}
func extractFacts(jsonlPath, chatID string) ([]ExtractedFact, string) {
f, err := os.Open(jsonlPath)
if err != nil {
return nil, ""
}
defer f.Close()
var facts []ExtractedFact
chatName := ""
scanner := bufio.NewScanner(f)
scanner.Buffer(make([]byte, 1<<20), 1<<20)
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
if line == "" {
continue
}
var msg Message
if err := json.Unmarshal([]byte(line), &msg); err != nil {
continue
}
if chatName == "" && msg.From != "" {
chatName = msg.From
}
text := msg.Text
phones := phoneRegex.FindAllString(text, -1)
for _, p := range phones {
p = strings.TrimSpace(p)
p = strings.Trim(p, "()- \t")
if isValidPhone(p) {
facts = append(facts, ExtractedFact{
ChatID: chatID,
ChatName: chatName,
Platform: "telegram",
FactType: "phone",
Value: p,
Source: "chat:" + msg.ID,
MessageID: msg.ID,
})
}
}
emails := emailRegex.FindAllString(text, -1)
for _, e := range emails {
facts = append(facts, ExtractedFact{
ChatID: chatID,
ChatName: chatName,
Platform: "telegram",
FactType: "email",
Value: strings.ToLower(e),
Source: "chat:" + msg.ID,
MessageID: msg.ID,
})
}
linkedins := linkedinRegex.FindAllString(text, -1)
for _, l := range linkedins {
facts = append(facts, ExtractedFact{
ChatID: chatID,
ChatName: chatName,
Platform: "telegram",
FactType: "linkedin",
Value: "https://" + l,
Source: "chat:" + msg.ID,
MessageID: msg.ID,
})
}
if matches := projectRegex.FindStringSubmatch(text); len(matches) > 1 {
facts = append(facts, ExtractedFact{
ChatID: chatID,
ChatName: chatName,
Platform: "telegram",
FactType: "project",
Value: strings.TrimSpace(matches[1]),
Source: "chat:" + msg.ID,
MessageID: msg.ID,
})
}
if matches := dealRegex.FindStringSubmatch(text); len(matches) > 2 {
facts = append(facts, ExtractedFact{
ChatID: chatID,
ChatName: chatName,
Platform: "telegram",
FactType: "deal",
Value: strings.TrimSpace(matches[2]),
Source: "chat:" + msg.ID,
MessageID: msg.ID,
})
}
if matches := skillRegex.FindStringSubmatch(text); len(matches) > 3 {
facts = append(facts, ExtractedFact{
ChatID: chatID,
ChatName: chatName,
Platform: "telegram",
FactType: "skill",
Value: strings.TrimSpace(matches[0]),
Source: "chat:" + msg.ID,
MessageID: msg.ID,
})
}
}
return facts, chatName
}
func filterFacts(facts []ExtractedFact, factType string) []ExtractedFact {
var result []ExtractedFact
for _, f := range facts {
if f.FactType == factType {
result = append(result, f)
}
}
return result
}
func writeFactsToBrain(root string, facts []ExtractedFact) {
indexScript := filepath.Join(root, "bin", "kb", "index")
if _, err := os.Stat(indexScript); os.IsNotExist(err) {
fmt.Fprintf(os.Stderr, "chats facts: kb/index not found, skipping brain write\n")
return
}
mdDir := filepath.Join(chatsDir(), "facts")
if err := os.MkdirAll(mdDir, 0755); err != nil {
fmt.Fprintf(os.Stderr, "chats facts: mkdir %s: %v\n", mdDir, err)
return
}
var sb strings.Builder
sb.WriteString("---\n")
sb.WriteString("root: facts\n")
sb.WriteString("---\n\n")
sb.WriteString("# Chat-Derived Facts\n\n")
for _, f := range facts {
sb.WriteString(fmt.Sprintf("- **%s**: %s (source: %s, chat: %s)\n",
f.FactType, f.Value, f.Source, f.ChatName))
}
sb.WriteString("\n")
factsMD := filepath.Join(mdDir, "chat-facts.md")
if err := os.WriteFile(factsMD, []byte(sb.String()), 0644); err != nil {
fmt.Fprintf(os.Stderr, "chats facts: write %s: %v\n", factsMD, err)
return
}
cmd := exec.Command(indexScript, "--corpus", mdDir, "--skip-indexes")
var outBuf, errBuf bytes.Buffer
cmd.Stdout = &outBuf
cmd.Stderr = &errBuf
cmd.Dir = root
if err := cmd.Run(); err != nil {
fmt.Fprintf(os.Stderr, "chats facts: brain index: %v\n%s", err, errBuf.String())
return
}
fmt.Printf("chats facts: written to brain (%s)\n", strings.TrimSpace(outBuf.String()))
}
+3
View File
@@ -0,0 +1,3 @@
module github.com/eSlider/2dph/bin/chats
go 1.25.0
+204
View File
@@ -0,0 +1,204 @@
package main
import (
"bufio"
"bytes"
"encoding/json"
"flag"
"fmt"
"html"
"os"
"path/filepath"
"sort"
"strings"
)
func runImport(args []string) int {
fs := flag.NewFlagSet("chats import", flag.ContinueOnError)
help := fs.Bool("help", false, "")
fs.SetOutput(os.Stderr)
if err := fs.Parse(args); err != nil {
return 2
}
if *help {
fmt.Fprintln(os.Stderr, "usage: chats import")
return 0
}
root := chatsDir()
mdRoot := filepath.Join(root, "md")
glob := filepath.Join(root, "telegram", "*", "messages.jsonl")
matches, err := filepath.Glob(glob)
if err != nil {
fmt.Fprintf(os.Stderr, "chats import: glob %s: %v\n", glob, err)
return 1
}
if len(matches) == 0 {
fmt.Fprintf(os.Stderr, "chats import: no messages.jsonl found under %s\n", root)
return 1
}
written := 0
failed := 0
for _, jsonlPath := range matches {
chatID := filepath.Base(filepath.Dir(jsonlPath))
messages, chatName, err := readJSONL(jsonlPath)
if err != nil {
fmt.Fprintf(os.Stderr, "chats import: read %s: %v\n", jsonlPath, err)
failed++
continue
}
if len(messages) == 0 {
continue
}
if chatName == "" {
chatName = chatID
}
participants := collectParticipants(messages)
chatType := "personal"
if len(participants) > 3 {
chatType = "group"
}
firstID := ""
if len(messages) > 0 {
firstID = messages[0].ID
}
var b bytes.Buffer
b.WriteString("---\n")
fmt.Fprintf(&b, "id: %s\n", firstID)
fmt.Fprintf(&b, "platform: telegram\n")
fmt.Fprintf(&b, "chat_id: %s\n", chatID)
fmt.Fprintf(&b, "chat_name: %s\n", escapeYAML(chatName))
fmt.Fprintf(&b, "participants: [")
for i, p := range participants {
if i > 0 {
b.WriteString(", ")
}
b.WriteString(escapeYAML(p))
}
b.WriteString("]\n")
fmt.Fprintf(&b, "message_count: %d\n", len(messages))
fmt.Fprintf(&b, "type: %s\n", chatType)
b.WriteString("---\n\n")
fmt.Fprintf(&b, "# Чат с %s\n\n", chatName)
for _, msg := range messages {
ts := msg.Timestamp
if len(ts) > 10 {
ts = ts[:10]
}
text := msg.Text
text = html.UnescapeString(text)
text = strings.ReplaceAll(text, "\n", "\n ")
line := fmt.Sprintf("**%s** — %s: %s", ts, msg.From, text)
if msg.Media != nil {
line += " *(" + *msg.Media + ")*"
}
b.WriteString(line + "\n\n")
}
mdFile := filepath.Join(mdRoot, "telegram", sanitizeDir(chatName), "messages.md")
if err := os.MkdirAll(filepath.Dir(mdFile), 0755); err != nil {
fmt.Fprintf(os.Stderr, "chats import: mkdir %s: %v\n", filepath.Dir(mdFile), err)
failed++
continue
}
if err := os.WriteFile(mdFile, b.Bytes(), 0644); err != nil {
fmt.Fprintf(os.Stderr, "chats import: write %s: %v\n", mdFile, err)
failed++
continue
}
written++
}
fmt.Printf("chats import: %d chats written", written)
if failed > 0 {
fmt.Printf(", %d failed", failed)
}
fmt.Println()
if failed > 0 {
return 1
}
return 0
}
func readJSONL(path string) ([]Message, string, error) {
f, err := os.Open(path)
if err != nil {
return nil, "", err
}
defer f.Close()
var messages []Message
scanner := bufio.NewScanner(f)
scanner.Buffer(make([]byte, 1<<20), 1<<20)
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
if line == "" {
continue
}
var msg Message
if err := json.Unmarshal([]byte(line), &msg); err != nil {
continue
}
messages = append(messages, msg)
}
if err := scanner.Err(); err != nil {
return messages, "", err
}
chatName := ""
if len(messages) > 0 {
nameCounts := make(map[string]int)
for _, msg := range messages {
nameCounts[msg.From]++
}
best := ""
bestN := 0
for name, n := range nameCounts {
if name != "" && name != "unknown" && n > bestN {
best = name
bestN = n
}
}
if best != "" {
chatName = best
}
}
return messages, chatName, nil
}
func collectParticipants(messages []Message) []string {
seen := make(map[string]bool)
var result []string
for _, msg := range messages {
if msg.From == "" || seen[msg.From] {
continue
}
seen[msg.From] = true
result = append(result, msg.From)
}
sort.Strings(result)
return result
}
func escapeYAML(s string) string {
if strings.ContainsAny(s, ":#,[]{}'\"") || strings.HasPrefix(s, "-") {
return `"` + strings.ReplaceAll(s, `"`, `\"`) + `"`
}
return s
}
func sanitizeDir(name string) string {
r := strings.NewReplacer(
"/", "_", "\\", "_", ":", "_", "*", "_",
"?", "_", "\"", "_", "<", "_", ">", "_", "|", "_",
" ", "_",
)
return strings.TrimSpace(r.Replace(name))
}
+56
View File
@@ -0,0 +1,56 @@
package main
import (
"bytes"
"flag"
"fmt"
"os"
"os/exec"
"path/filepath"
"strings"
)
func runIndex(args []string) int {
fs := flag.NewFlagSet("chats index", flag.ContinueOnError)
help := fs.Bool("help", false, "")
fs.SetOutput(os.Stderr)
if err := fs.Parse(args); err != nil {
return 2
}
if *help {
fmt.Fprintln(os.Stderr, "usage: chats index")
return 0
}
root := repoRoot()
mdDir := filepath.Join(chatsDir(), "md")
_, err := os.Stat(mdDir)
if os.IsNotExist(err) {
fmt.Fprintf(os.Stderr, "chats index: no chat markdown at %s; run 'chats import' first\n", mdDir)
return 1
}
indexScript := filepath.Join(root, "bin", "kb", "index")
if _, err := os.Stat(indexScript); os.IsNotExist(err) {
fmt.Fprintf(os.Stderr, "chats index: %s not found\n", indexScript)
return 1
}
cmd := exec.Command(indexScript, "--corpus", mdDir)
var outBuf, errBuf bytes.Buffer
cmd.Stdout = &outBuf
cmd.Stderr = &errBuf
cmd.Dir = root
if err := cmd.Run(); err != nil {
fmt.Fprintf(os.Stderr, "chats index: %v\n%s", err, errBuf.String())
return 1
}
result := strings.TrimSpace(outBuf.String())
if result == "" {
result = strings.TrimSpace(errBuf.String())
}
fmt.Printf("chats index: %s\n", result)
return 0
}
+356
View File
@@ -0,0 +1,356 @@
package main
import (
"bufio"
"context"
"encoding/json"
"fmt"
"os"
"os/exec"
"path/filepath"
"strings"
"time"
)
type LinkedInMCPSource struct {
userDataDir string
limit int
}
type lnInboxItem struct {
ThreadID string `json:"thread_id"`
Participants string `json:"participants"`
LastMessage string `json:"last_message"`
LastMessageDate string `json:"last_message_date"`
Unread bool `json:"unread"`
}
type lnInboxEnvelope struct {
Results []lnInboxItem `json:"results"`
HasMore bool `json:"hasMore"`
}
type lnMessage struct {
From string `json:"from"`
Date string `json:"date"`
Text string `json:"text"`
}
type lnConvEnvelope struct {
Results []lnMessage `json:"results"`
HasMore bool `json:"hasMore"`
TotalCount int `json:"total_count"`
}
func NewLinkedInMCPSource(userDataDir string) *LinkedInMCPSource {
return &LinkedInMCPSource{userDataDir: userDataDir}
}
func (s *LinkedInMCPSource) Name() string { return "linkedin" }
func (s *LinkedInMCPSource) Sync(ctx context.Context, outDir string, limit int) error {
if limit > 0 {
s.limit = limit
}
client, err := newLinkedInMCP(ctx, s.userDataDir)
if err != nil {
return fmt.Errorf("linkedin mcp: %w", err)
}
defer client.Close()
inbox, err := client.GetInbox(ctx, 50)
if err != nil {
return fmt.Errorf("get_inbox: %w", err)
}
if len(inbox) == 0 {
fmt.Println("chats: no LinkedIn conversations found")
return nil
}
fmt.Printf("chats: found %d LinkedIn conversations\n", len(inbox))
for _, conv := range inbox {
convID := sanitizeDir(conv.ThreadID)
if convID == "" {
convID = fmt.Sprintf("conv_%d", time.Now().UnixNano())
}
parts := strings.SplitN(conv.Participants, ",", 2)
chatName := strings.TrimSpace(parts[0])
if chatName == "" {
chatName = convID
}
chatDir := filepath.Join(outDir, "linkedin", convID)
if err := os.MkdirAll(chatDir, 0755); err != nil {
fmt.Fprintf(os.Stderr, "chats: mkdir %s: %v\n", chatDir, err)
continue
}
msgLimit := 100
if s.limit > 0 {
msgLimit = s.limit
}
msgs, err := client.GetConversation(ctx, "", conv.ThreadID, msgLimit)
if err != nil {
fmt.Fprintf(os.Stderr, "chats: get_conversation %s: %v\n", convID, err)
continue
}
jsonlPath := filepath.Join(chatDir, "messages.jsonl")
f, err := os.Create(jsonlPath)
if err != nil {
fmt.Fprintf(os.Stderr, "chats: create %s: %v\n", jsonlPath, err)
continue
}
enc := json.NewEncoder(f)
written := 0
for i, m := range msgs {
text := m.Text
if text == "" {
continue
}
ts := m.Date
if t, err := time.Parse("2006-01-02T15:04:05Z07:00", m.Date); err == nil {
ts = t.UTC().Format(time.RFC3339)
} else if t, err := time.Parse(time.RFC3339, m.Date); err == nil {
ts = t.UTC().Format(time.RFC3339)
}
chatMsg := Message{
ID: fmt.Sprintf("li_%s_%d", convID, i),
Timestamp: ts,
From: m.From,
Text: text,
Platform: "linkedin",
}
if err := enc.Encode(chatMsg); err != nil {
fmt.Fprintf(os.Stderr, "chats: encode: %v\n", err)
continue
}
written++
}
f.Close()
fmt.Printf("chats: synced %s (%s) — %d messages\n", chatName, convID, written)
}
return nil
}
type linkedInMCPClient struct {
cmd *exec.Cmd
stdin *bufio.Writer
stdout *bufio.Scanner
msgID int
}
func newLinkedInMCP(ctx context.Context, userDataDir string) (*linkedInMCPClient, error) {
args := []string{
"mcp-server-linkedin@latest",
"--user-data-dir", userDataDir,
"--no-auto-import",
"--transport", "stdio",
"--login-timeout", "10",
"--browser-wait", "1",
"--browser-idle-timeout", "10",
"--log-level", "ERROR",
}
cmd := exec.CommandContext(ctx, "uvx", args...)
cmd.Env = os.Environ()
stdin, err := cmd.StdinPipe()
if err != nil {
return nil, fmt.Errorf("stdin pipe: %w", err)
}
stdout, err := cmd.StdoutPipe()
if err != nil {
return nil, fmt.Errorf("stdout pipe: %w", err)
}
cmd.Stderr = os.Stderr
if err := cmd.Start(); err != nil {
return nil, fmt.Errorf("start: %w", err)
}
c := &linkedInMCPClient{
cmd: cmd,
stdin: bufio.NewWriter(stdin),
stdout: bufio.NewScanner(stdout),
msgID: 0,
}
c.stdout.Buffer(make([]byte, 1<<20), 1<<20)
if err := c.initialize(ctx); err != nil {
c.Close()
return nil, fmt.Errorf("init: %w", err)
}
return c, nil
}
func (c *linkedInMCPClient) nextID() int {
c.msgID++
return c.msgID
}
func (c *linkedInMCPClient) initialize(ctx context.Context) error {
params := map[string]interface{}{
"protocolVersion": "2024-11-05",
"capabilities": map[string]interface{}{},
"clientInfo": map[string]string{
"name": "chats-sync",
"version": "0.1.0",
},
}
_, err := c.send(ctx, "initialize", params)
return err
}
func (c *linkedInMCPClient) send(ctx context.Context, method string, params interface{}) (json.RawMessage, error) {
id := c.nextID()
req := map[string]interface{}{
"jsonrpc": "2.0",
"id": id,
"method": method,
}
if params != nil {
req["params"] = params
}
body, err := json.Marshal(req)
if err != nil {
return nil, fmt.Errorf("marshal: %w", err)
}
if _, err := c.stdin.Write(body); err != nil {
return nil, fmt.Errorf("write: %w", err)
}
if err := c.stdin.WriteByte('\n'); err != nil {
return nil, fmt.Errorf("newline: %w", err)
}
if err := c.stdin.Flush(); err != nil {
return nil, fmt.Errorf("flush: %w", err)
}
for c.stdout.Scan() {
line := c.stdout.Text()
if line == "" {
continue
}
var resp struct {
JSONRPC string `json:"jsonrpc"`
ID int `json:"id"`
Result json.RawMessage `json:"result,omitempty"`
Error *struct {
Code int `json:"code"`
Message string `json:"message"`
} `json:"error,omitempty"`
}
if err := json.Unmarshal([]byte(line), &resp); err != nil {
return nil, fmt.Errorf("unmarshal: %w\nline: %s", err, line[:min(len(line), 500)])
}
if resp.Error != nil {
return nil, fmt.Errorf("rpc error %d: %s", resp.Error.Code, resp.Error.Message)
}
return resp.Result, nil
}
return nil, fmt.Errorf("no response: %w", c.stdout.Err())
}
func (c *linkedInMCPClient) GetInbox(ctx context.Context, limit int) ([]lnInboxItem, error) {
params := map[string]interface{}{
"limit": limit,
}
result, err := c.send(ctx, "tools/call", map[string]interface{}{
"name": "get_inbox",
"arguments": params,
})
if err != nil {
return nil, err
}
var toolRes struct {
Content []struct {
Type string `json:"type"`
Text string `json:"text"`
} `json:"content"`
IsError bool `json:"isError"`
}
if err := json.Unmarshal(result, &toolRes); err != nil {
return nil, fmt.Errorf("unmarshal tool: %w", err)
}
if toolRes.IsError {
return nil, fmt.Errorf("get_inbox error")
}
if len(toolRes.Content) == 0 {
return nil, nil
}
text := toolRes.Content[0].Text
var env lnInboxEnvelope
if err := json.Unmarshal([]byte(text), &env); err != nil {
var arr []lnInboxItem
if err2 := json.Unmarshal([]byte(text), &arr); err2 == nil {
return arr, nil
}
return nil, fmt.Errorf("parse inbox: %w", err)
}
return env.Results, nil
}
func (c *linkedInMCPClient) GetConversation(ctx context.Context, username, threadID string, limit int) ([]lnMessage, error) {
params := map[string]interface{}{
"linkedin_username": username,
"thread_id": threadID,
"index": limit,
}
result, err := c.send(ctx, "tools/call", map[string]interface{}{
"name": "get_conversation",
"arguments": params,
})
if err != nil {
return nil, err
}
var toolRes struct {
Content []struct {
Type string `json:"type"`
Text string `json:"text"`
} `json:"content"`
IsError bool `json:"isError"`
}
if err := json.Unmarshal(result, &toolRes); err != nil {
return nil, fmt.Errorf("unmarshal tool: %w", err)
}
if toolRes.IsError {
return nil, nil
}
if len(toolRes.Content) == 0 {
return nil, nil
}
text := toolRes.Content[0].Text
var env lnConvEnvelope
if err := json.Unmarshal([]byte(text), &env); err != nil {
var arr []lnMessage
if err2 := json.Unmarshal([]byte(text), &arr); err2 == nil {
return arr, nil
}
return nil, fmt.Errorf("parse conv: %w", err)
}
return env.Results, nil
}
func (c *linkedInMCPClient) Close() error {
if c.stdin != nil {
c.stdin.Flush()
}
if c.cmd != nil && c.cmd.Process != nil {
c.cmd.Process.Kill()
}
return nil
}
+115
View File
@@ -0,0 +1,115 @@
// bin/chats - sync, import, index, extract facts, and apply chat data
// from Telegram, WhatsApp, LinkedIn into the brain and OnlyOffice CRM.
//
// Usage:
//
// chats sync telegram [--limit N] [--since DATE] [--phone PHONE]
// chats sync whatsapp [--qr] [--limit N]
// chats sync linkedin [--limit N]
// chats import # JSONL → MD (all sources)
// chats index # rebuild var/kb.lbug with chats
// chats facts # extract + cross-check
// chats apply [--dry-run] # push to OnlyOffice CRM
package main
import (
"fmt"
"os"
"strings"
)
func main() {
if len(os.Args) < 2 {
usage()
os.Exit(2)
}
cmd := os.Args[1]
args := os.Args[2:]
switch cmd {
case "sync":
if len(args) < 1 {
usage()
os.Exit(2)
}
platform := args[0]
platformArgs := args[1:]
switch platform {
case "telegram":
os.Exit(runSyncTelegram(platformArgs))
case "whatsapp":
fmt.Fprintf(os.Stderr, "chats: WhatsApp not implemented yet\n")
os.Exit(1)
case "linkedin":
os.Exit(runSyncLinkedIn(platformArgs))
default:
fmt.Fprintf(os.Stderr, "chats: unknown platform %q\n", platform)
os.Exit(2)
}
case "import":
os.Exit(runImport(args))
case "index":
os.Exit(runIndex(args))
case "facts":
os.Exit(runFacts(args))
case "apply":
os.Exit(runApply(args))
case "help", "-h", "--help":
usage()
return
default:
fmt.Fprintf(os.Stderr, "chats: unknown command %q\n", cmd)
usage()
os.Exit(2)
}
}
func usage() {
w := os.Stderr
fmt.Fprintln(w, `Usage: chats <command> [args]
Commands:
sync telegram [--limit N] [--since DATE] [--phone PHONE]
sync whatsapp [--qr] [--limit N]
sync linkedin [--limit N]
import JSONL → MD (all sources)
index rebuild var/kb.lbug with chats
facts extract + cross-check facts
apply [--dry-run] push to OnlyOffice CRM
Output layout:
var/chats/<platform>/<chat_id>/messages.jsonl
var/chats/md/<platform>/<chat_name>/messages.md`)
}
// repoRoot locates the 2dph project root by walking up from the binary.
func repoRoot() string {
if v := os.Getenv("KB_ROOT"); v != "" {
return v
}
wd, err := os.Getwd()
if err != nil {
return "."
}
for i := 0; i < 10; i++ {
if _, err := os.Stat(wd + "/var"); err == nil {
return wd
}
if _, err := os.Stat(wd + "/.git"); err == nil {
return wd
}
parent := wd
if idx := strings.LastIndex(wd, "/"); idx >= 0 {
parent = wd[:idx]
}
if parent == wd {
break
}
wd = parent
}
return "."
}
// chatsDir returns var/chats under the repo root.
func chatsDir() string {
return repoRoot() + "/var/chats"
}
+424
View File
@@ -0,0 +1,424 @@
package main
import (
"bufio"
"context"
"encoding/json"
"fmt"
"os"
"os/exec"
"path/filepath"
"strings"
"time"
)
type MCPClient struct {
cmd *exec.Cmd
stdin *bufio.Writer
stdout *bufio.Scanner
msgID int
}
type mcpRequest struct {
JSONRPC string `json:"jsonrpc"`
ID int `json:"id"`
Method string `json:"method"`
Params interface{} `json:"params,omitempty"`
}
type mcpResponse struct {
JSONRPC string `json:"jsonrpc"`
ID int `json:"id"`
Result json.RawMessage `json:"result,omitempty"`
Error *struct {
Code int `json:"code"`
Message string `json:"message"`
} `json:"error,omitempty"`
}
type mcpToolResult struct {
Content []struct {
Type string `json:"type"`
Text string `json:"text"`
} `json:"content"`
IsError bool `json:"isError,omitempty"`
}
type ListChatsResult struct {
ChatID int64 `json:"chat_id"`
Title string `json:"name"`
Type string `json:"type"`
Username string `json:"username,omitempty"`
}
type listChatsEnvelope struct {
Results []ListChatsResult `json:"results"`
}
type historyEnvelope struct {
Results []GetHistoryResult `json:"results"`
}
type GetHistoryResult struct {
ID int `json:"id"`
Sender string `json:"sender"`
Date string `json:"date"`
Text string `json:"text"`
Media string `json:"media,omitempty"`
Out bool `json:"out,omitempty"`
}
func NewMCPClient(ctx context.Context, apiID int, apiHash, phone, sessionString, mcpDir string) (*MCPClient, error) {
env := os.Environ()
env = append(env,
fmt.Sprintf("TELEGRAM_API_ID=%d", apiID),
fmt.Sprintf("TELEGRAM_API_HASH=%s", apiHash),
fmt.Sprintf("TELEGRAM_PHONE=%s", phone),
fmt.Sprintf("TELEGRAM_SESSION_STRING=%s", sessionString),
"MCP_TRANSPORT=stdio",
)
serverPath := filepath.Join(mcpDir, ".venv", "bin", "python3")
mainPath := filepath.Join(mcpDir, "main.py")
cmd := exec.CommandContext(ctx, serverPath, mainPath)
cmd.Env = env
cmd.Dir = mcpDir
stdin, err := cmd.StdinPipe()
if err != nil {
return nil, fmt.Errorf("stdin pipe: %w", err)
}
stdout, err := cmd.StdoutPipe()
if err != nil {
return nil, fmt.Errorf("stdout pipe: %w", err)
}
cmd.Stderr = os.Stderr
if err := cmd.Start(); err != nil {
return nil, fmt.Errorf("start mcp: %w", err)
}
c := &MCPClient{
cmd: cmd,
stdin: bufio.NewWriter(stdin),
stdout: bufio.NewScanner(stdout),
msgID: 0,
}
c.stdout.Buffer(make([]byte, 1<<20), 1<<20)
if err := c.initialize(ctx); err != nil {
c.Close()
return nil, fmt.Errorf("initialize: %w", err)
}
return c, nil
}
func (c *MCPClient) nextID() int {
c.msgID++
return c.msgID
}
func (c *MCPClient) sendRequest(ctx context.Context, method string, params interface{}) (json.RawMessage, error) {
id := c.nextID()
req := mcpRequest{
JSONRPC: "2.0",
ID: id,
Method: method,
Params: params,
}
body, err := json.Marshal(req)
if err != nil {
return nil, fmt.Errorf("marshal: %w", err)
}
if _, err := c.stdin.Write(body); err != nil {
return nil, fmt.Errorf("write: %w", err)
}
if err := c.stdin.WriteByte('\n'); err != nil {
return nil, fmt.Errorf("write newline: %w", err)
}
if err := c.stdin.Flush(); err != nil {
return nil, fmt.Errorf("flush: %w", err)
}
for c.stdout.Scan() {
line := c.stdout.Text()
if line == "" {
continue
}
var resp mcpResponse
if err := json.Unmarshal([]byte(line), &resp); err != nil {
return nil, fmt.Errorf("unmarshal response: %w\nline: %s", err, line[:min(len(line), 500)])
}
if resp.Error != nil {
return nil, fmt.Errorf("rpc error %d: %s", resp.Error.Code, resp.Error.Message)
}
return resp.Result, nil
}
return nil, fmt.Errorf("no response: %w", c.stdout.Err())
}
func (c *MCPClient) initialize(ctx context.Context) error {
params := map[string]interface{}{
"protocolVersion": "2024-11-05",
"capabilities": map[string]interface{}{},
"clientInfo": map[string]string{
"name": "chats-sync",
"version": "0.1.0",
},
}
_, err := c.sendRequest(ctx, "initialize", params)
return err
}
func (c *MCPClient) ListChats(ctx context.Context, chatType string, limit int) ([]ListChatsResult, error) {
args := map[string]interface{}{
"chat_type": chatType,
"limit": limit,
}
result, err := c.sendRequest(ctx, "tools/call", map[string]interface{}{
"name": "list_chats",
"arguments": args,
})
if err != nil {
return nil, err
}
var toolRes mcpToolResult
if err := json.Unmarshal(result, &toolRes); err != nil {
return nil, fmt.Errorf("unmarshal tool result: %w", err)
}
if toolRes.IsError {
msg := "unknown"
if len(toolRes.Content) > 0 {
msg = toolRes.Content[0].Text
}
return nil, fmt.Errorf("list_chats error: %s", msg)
}
if len(toolRes.Content) == 0 {
return nil, nil
}
text := toolRes.Content[0].Text
if text == "" || text == "No chats found matching the criteria." {
return nil, nil
}
var env listChatsEnvelope
if err := json.Unmarshal([]byte(text), &env); err != nil {
var arr []ListChatsResult
if err2 := json.Unmarshal([]byte(text), &arr); err2 != nil {
return nil, fmt.Errorf("parse chats: %w (also tried array: %v)\nbody: %s", err, err2, text[:min(len(text), 500)])
}
return arr, nil
}
return env.Results, nil
}
func (c *MCPClient) GetHistory(ctx context.Context, chatID int64, limit int) ([]GetHistoryResult, error) {
args := map[string]interface{}{
"chat_id": chatID,
"limit": limit,
}
result, err := c.sendRequest(ctx, "tools/call", map[string]interface{}{
"name": "get_history",
"arguments": args,
})
if err != nil {
return nil, err
}
var toolRes mcpToolResult
if err := json.Unmarshal(result, &toolRes); err != nil {
return nil, fmt.Errorf("unmarshal tool result: %w", err)
}
if toolRes.IsError {
msg := "unknown"
if len(toolRes.Content) > 0 {
msg = toolRes.Content[0].Text
}
return nil, fmt.Errorf("get_history error: %s", msg)
}
if len(toolRes.Content) == 0 {
return nil, nil
}
text := toolRes.Content[0].Text
if text == "" || text == "No messages found for this page." || text == "No messages found matching the criteria." {
return nil, nil
}
var env historyEnvelope
if err := json.Unmarshal([]byte(text), &env); err != nil {
var arr []GetHistoryResult
if err2 := json.Unmarshal([]byte(text), &arr); err2 != nil {
return nil, fmt.Errorf("parse history: %w (also tried array: %v)\nbody: %s", err, err2, text[:min(len(text), 500)])
}
return arr, nil
}
return env.Results, nil
}
func (c *MCPClient) Close() error {
if c.stdin != nil {
c.stdin.Flush()
}
if c.cmd != nil && c.cmd.Process != nil {
c.cmd.Process.Kill()
}
return nil
}
type TelegramMCPSource struct {
mcpDir string
apiID int
apiHash string
phone string
sessionStr string
limit int
}
func NewTelegramMCPSource(apiID int, apiHash, phone, sessionString, mcpDir string) *TelegramMCPSource {
return &TelegramMCPSource{
mcpDir: mcpDir,
apiID: apiID,
apiHash: apiHash,
phone: phone,
sessionStr: sessionString,
}
}
func (s *TelegramMCPSource) Name() string { return "telegram" }
func (s *TelegramMCPSource) Sync(ctx context.Context, outDir string, limit int) error {
if limit > 0 {
s.limit = limit
}
client, err := NewMCPClient(ctx, s.apiID, s.apiHash, s.phone, s.sessionStr, s.mcpDir)
if err != nil {
return fmt.Errorf("mcp client: %w", err)
}
defer client.Close()
chats, err := client.ListChats(ctx, "user", 100)
if err != nil {
return fmt.Errorf("list chats: %w", err)
}
if len(chats) == 0 {
fmt.Println("chats: no personal chats found")
return nil
}
fmt.Printf("chats: found %d personal chats\n", len(chats))
var filtered []ListChatsResult
for _, c := range chats {
if strings.Contains(strings.ToLower(c.Username), "bot") {
continue
}
if c.ChatID == 777000 { // Telegram service
continue
}
filtered = append(filtered, c)
}
fmt.Printf("chats: %d after filter (bots excluded)\n", len(filtered))
for _, chat := range filtered {
chatID := fmt.Sprintf("user_%d", chat.ChatID)
chatName := chat.Title
if chatName == "" {
chatName = chatID
}
chatDir := filepath.Join(outDir, "telegram", chatID)
if err := os.MkdirAll(chatDir, 0755); err != nil {
fmt.Fprintf(os.Stderr, "chats: mkdir %s: %v\n", chatDir, err)
continue
}
jsonlPath := filepath.Join(chatDir, "messages.jsonl")
f, err := os.Create(jsonlPath)
if err != nil {
fmt.Fprintf(os.Stderr, "chats: create %s: %v\n", jsonlPath, err)
continue
}
msgLimit := 100
if s.limit > 0 {
msgLimit = s.limit
}
msgs, err := client.GetHistory(ctx, chat.ChatID, msgLimit)
if err != nil {
fmt.Fprintf(os.Stderr, "chats: get_history for %s: %v\n", chatName, err)
f.Close()
continue
}
enc := json.NewEncoder(f)
written := 0
for _, m := range msgs {
if m.Out {
continue
}
text := m.Text
if text == "" && m.Media != "" {
text = fmt.Sprintf("[%s]", m.Media)
}
if text == "" {
continue
}
sender := cleanSender(m.Sender)
ts := m.Date
if t, err := time.Parse(time.RFC3339, m.Date); err == nil {
ts = t.UTC().Format(time.RFC3339)
}
chatMsg := Message{
ID: fmt.Sprintf("tg_%d_%d", chat.ChatID, m.ID),
Timestamp: ts,
From: sender,
Text: text,
Platform: "telegram",
}
if m.Media != "" {
desc := fmt.Sprintf("[%s]", m.Media)
chatMsg.Media = &desc
}
if err := enc.Encode(chatMsg); err != nil {
fmt.Fprintf(os.Stderr, "chats: encode msg: %v\n", err)
continue
}
written++
}
f.Close()
if written > 0 {
fmt.Printf("chats: synced %s (%s) — %d messages\n", chatName, chatID, written)
}
}
return nil
}
func cleanSender(sender string) string {
if idx := strings.Index(sender, " ("); idx > 0 {
sender = sender[:idx]
} else if idx := strings.Index(sender, " @"); idx > 0 {
sender = sender[:idx]
}
if idx := strings.Index(sender, " ["); idx > 0 {
sender = sender[:idx]
}
return sender
}
func min(a, b int) int {
if a < b {
return a
}
return b
}
+52
View File
@@ -0,0 +1,52 @@
package main
import (
"context"
"fmt"
"os"
"time"
)
type Source interface {
Name() string
Sync(ctx context.Context, outDir string, limit int) error
}
type Message struct {
ID string `json:"id"`
Timestamp string `json:"ts"`
From string `json:"from"`
Text string `json:"text"`
Media *string `json:"media,omitempty"`
Platform string `json:"platform"`
}
type ChatInfo struct {
ID string `json:"id"`
Platform string `json:"platform"`
Name string `json:"name"`
Participants []string `json:"participants"`
Type string `json:"type"`
MessageCount int `json:"messageCount"`
LastTS string `json:"lastTs,omitempty"`
}
func envVar(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
}
return fallback
}
func parseSince(s string) (time.Time, error) {
for _, layout := range []string{
time.RFC3339,
"2006-01-02T15:04:05",
"2006-01-02",
} {
if t, err := time.Parse(layout, s); err == nil {
return t, nil
}
}
return time.Time{}, fmt.Errorf("cannot parse --since %q; use YYYY-MM-DD or RFC3339", s)
}
+87
View File
@@ -0,0 +1,87 @@
package main
import (
"context"
"flag"
"fmt"
"os"
"path/filepath"
"strconv"
"strings"
"time"
)
func runSyncTelegram(args []string) int {
fs := flag.NewFlagSet("chats sync telegram", flag.ContinueOnError)
limit := fs.Int("limit", 0, "max messages per chat (0 = all)")
phone := fs.String("phone", "", "phone number (default env TELEGRAM_PHONE)")
help := fs.Bool("help", false, "")
fs.SetOutput(os.Stderr)
if err := fs.Parse(args); err != nil {
return 2
}
if *help {
fmt.Fprintln(os.Stderr, "usage: chats sync telegram [--limit N] [--phone PHONE]")
return 0
}
apiIDStr := envVar("TELEGRAM_API_ID", "")
apiHash := envVar("TELEGRAM_API_HASH", "")
sessionStr := envVar("TELEGRAM_SESSION_STRING", "")
phoneNum := *phone
if phoneNum == "" {
phoneNum = envVar("TELEGRAM_PHONE", "")
}
if apiIDStr == "" || apiHash == "" || phoneNum == "" {
fmt.Fprintln(os.Stderr, "chats: need TELEGRAM_API_ID, TELEGRAM_API_HASH, TELEGRAM_PHONE in env")
return 2
}
apiID, err := strconv.Atoi(apiIDStr)
if err != nil {
fmt.Fprintf(os.Stderr, "chats: invalid TELEGRAM_API_ID %q\n", apiIDStr)
return 2
}
mcpDir := envVar("TELEGRAM_MCP_DIR", "")
if mcpDir == "" {
fmt.Fprintln(os.Stderr, "chats: set TELEGRAM_MCP_DIR to telegram-mcp directory")
return 1
}
if _, err := os.Stat(filepath.Join(mcpDir, "main.py")); err != nil {
fmt.Fprintf(os.Stderr, "chats: TELEGRAM_MCP_DIR=%s: main.py not found\n", mcpDir)
return 1
}
if sessionStr == "" {
envPath := filepath.Join(mcpDir, ".env")
if data, err := os.ReadFile(envPath); err == nil {
for _, line := range strings.Split(string(data), "\n") {
if strings.HasPrefix(line, "TELEGRAM_SESSION_STRING=") {
sessionStr = strings.TrimPrefix(line, "TELEGRAM_SESSION_STRING=")
sessionStr = strings.Trim(sessionStr, "\"'")
break
}
}
}
}
if sessionStr == "" {
sessionStr = envVar("TELEGRAM_SESSION_STRING", "")
}
if sessionStr == "" {
fmt.Fprintln(os.Stderr, "chats: TELEGRAM_SESSION_STRING not found; set env or in TELEGRAM_MCP_DIR/.env")
return 1
}
src := NewTelegramMCPSource(apiID, apiHash, phoneNum, sessionStr, mcpDir)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
defer cancel()
start := time.Now()
if err := src.Sync(ctx, chatsDir(), *limit); err != nil {
fmt.Fprintf(os.Stderr, "chats sync telegram: %v\n", err)
return 1
}
fmt.Printf("chats sync telegram: completed in %s\n", time.Since(start).Round(time.Millisecond))
return 0
}
+69
View File
@@ -0,0 +1,69 @@
package main
import (
"context"
"flag"
"fmt"
"os"
"os/exec"
"strings"
"time"
)
func checkLinkedInSession(userDataDir string) (bool, error) {
cmd := exec.Command("uvx", "mcp-server-linkedin@latest",
"--user-data-dir", userDataDir,
"--no-auto-import",
"--status",
)
out, err := cmd.CombinedOutput()
if err != nil {
return true, fmt.Errorf("status check: %w\n%s", err, string(out))
}
return !strings.Contains(string(out), "✅"), nil
}
func runSyncLinkedIn(args []string) int {
fs := flag.NewFlagSet("chats sync linkedin", flag.ContinueOnError)
limit := fs.Int("limit", 0, "max messages per conversation (0 = all)")
help := fs.Bool("help", false, "")
fs.SetOutput(os.Stderr)
if err := fs.Parse(args); err != nil {
return 2
}
if *help {
fmt.Fprintln(os.Stderr, "usage: chats sync linkedin [--limit N]")
return 0
}
userDataDir := envVar("LINKEDIN_USER_DATA_DIR", "")
if userDataDir == "" {
home, _ := os.UserHomeDir()
userDataDir = home + "/.linkedin-mcp/profile"
}
// Check session first
loginNeeded, err := checkLinkedInSession(userDataDir)
if err != nil {
fmt.Fprintf(os.Stderr, "chats: linkedin status check: %v\n", err)
}
if loginNeeded {
fmt.Fprintf(os.Stderr, "chats: LinkedIn session expired. Run:\n")
fmt.Fprintf(os.Stderr, " uvx mcp-server-linkedin@latest --user-data-dir %s --login\n", userDataDir)
fmt.Fprintf(os.Stderr, "Then retry 'chats sync linkedin'\n")
return 1
}
src := NewLinkedInMCPSource(userDataDir)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
defer cancel()
start := time.Now()
if err := src.Sync(ctx, chatsDir(), *limit); err != nil {
fmt.Fprintf(os.Stderr, "chats sync linkedin: %v\n", err)
return 1
}
fmt.Printf("chats sync linkedin: completed in %s\n", time.Since(start).Round(time.Millisecond))
return 0
}
+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
Regular → Executable
+8 -4
View File
@@ -4,8 +4,10 @@
# brain shell (default) # brain shell (default)
# brain search <q> bin/kb/search # brain search <q> bin/kb/search
# brain index bin/kb/index # brain index bin/kb/index
# brain watch <dir> watchdog re-indexer # brain watch <dir> watchdog re-indexer (bin/kb/watch)
# brain serve async Go HTTP server (serve/) # brain serve async Go HTTP server (bin/serve)
# brain extract bin/facts/extract (docker×compose pairing)
# brain audit bin/facts/audit
# #
# Usage comment starts at line 2 (self-describing convention). # Usage comment starts at line 2 (self-describing convention).
set -euo pipefail set -euo pipefail
@@ -17,7 +19,9 @@ 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 "$@" ;;
watch) exec bash /app/bin/kb-watch "$@" ;; watch) exec /app/bin/watch "$@" ;;
serve) exec /app/serve/serve "$@" ;; serve) exec /app/bin/serve "$@" ;;
extract) exec "$KB_PY" /app/bin/facts/extract "$@" ;;
audit) exec "$KB_PY" /app/bin/facts/audit "$@" ;;
*) echo "unknown command: $CMD" >&2; exit 2 ;; *) echo "unknown command: $CMD" >&2; exit 2 ;;
esac esac
+1 -1
View File
@@ -20,7 +20,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"))
def audit_db() -> list[str]: def audit_db() -> list[str]:
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())
+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()
Executable
+154
View File
@@ -0,0 +1,154 @@
#!/usr/bin/env python3
"""git/import - import git history (commits, authors, files) into the brain.
bin/git/import [REPO] import all commits -> leafs + graph
bin/git/import --json emit import leafs as JSON, no write
bin/git/import --limit 100 cap commits processed
bin/git/import --since 2026-01-01 only recent commits
bin/git/import --root DIR run per repo dir under DIR
bin/git/import --no-env never read .env anywhere (default: true)
Reads `git log --no-merges --name-only` from the repo, maps commits to
`info` leafs (root=info, type=commit) and writes the version graph
`File -[:HAS_VERSION]-> Commit -[:AUTHORED]-> Person` into var/kb.lbug.
Idempotent: leaf MERGE by (source,text via leaf_id), graph MERGE by sha.
"""
from __future__ import annotations
import json
import subprocess
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
connect, ensure_indexes, init_schema, upsert_leaf,
)
from gitimport import commits_to_leafs, ensure_git_schema, index_commits, parse_log # noqa: E402
LOG_FMT = "--format=%x1e%H%x1f%an%x1f%ae%x1f%aI%x1f%s"
def git_log(repo: Path, limit: int = 0, since: str = "") -> str:
cmd = ["git", "-C", str(repo), "log", "--no-merges", "--name-only", LOG_FMT]
if since:
cmd += ["--since", since]
if limit:
cmd += ["-n", str(limit)]
try:
out = subprocess.run(cmd, capture_output=True, text=True, timeout=120)
except (FileNotFoundError, subprocess.TimeoutExpired):
return ""
if out.returncode != 0:
print(f"git/import: {repo}: {out.stderr.strip()}", file=sys.stderr)
return ""
return out.stdout
def repo_name(repo: Path) -> str:
try:
out = subprocess.run(
["git", "-C", str(repo), "remote", "get-url", "origin"],
capture_output=True, text=True, timeout=20)
url = out.stdout.strip()
return url.rstrip("/").split("/")[-1].removesuffix(".git") if url else repo.name
except (FileNotFoundError, subprocess.TimeoutExpired):
return repo.name
def embedder():
from model2vec import StaticModel
model = StaticModel.from_pretrained("minishlab/potion-multilingual-128M")
return lambda text: model.encode([text])[0].astype(float).tolist()
def import_repo(conn, repo: Path, embed, limit: int, since: str,
no_write: bool = False) -> tuple[int, int]:
raw = git_log(repo, limit, since)
commits = parse_log(raw)
leafs = commits_to_leafs(commits, repo_name(repo))
if no_write:
return len(commits), 0
written = 0
for lf in leafs:
query = f"{lf['heading']}\n\n{lf['text']}"
emb = embed(lf["text"]) if lf["text"] else None
upsert_leaf(conn, text=query, root="info", confidence="confirmed",
source=lf["source"], source_rev="git", how="git/import",
loc=lf["source"], type_=lf.get("type", "commit"),
embedding=emb)
written += 1
index_commits(conn, commits, repo_name(repo))
return len(commits), written
def main(argv: list[str]) -> int:
import argparse
p = argparse.ArgumentParser(description="import git history into the brain")
p.add_argument("repo", nargs="?", default=None)
p.add_argument("--root", default=None, help="directory of repos to import (each git dir separately)")
p.add_argument("--limit", type=int, default=0)
p.add_argument("--since", default="")
p.add_argument("--json", action="store_true")
p.add_argument("--dry-run", action="store_true", help="parse + report, no db write")
a = p.parse_args(argv)
repos: list[Path] = []
if a.repo:
repos = [Path(a.repo)]
elif a.root:
root = Path(a.root)
if root.is_file():
repos = [root]
else:
repos = [dp for dp in sorted(root.iterdir()) if (dp / ".git").exists() or dp.is_file()]
else:
repos = [ROOT]
total_commits = 0
results: list[dict] = []
if a.dry_run:
for repo in repos:
if not repo.exists():
continue
commits = parse_log(git_log(repo, a.limit, a.since))
name = repo_name(repo)
total_commits += len(commits)
results.append({"repo": name, "commits": len(commits),
"leafs": len(commits_to_leafs(commits, name)), "path": str(repo)})
if a.json:
print(json.dumps(results, indent=2))
else:
for r in results:
print(f"{r['repo']:<24} {r['commits']:>5} commits -> {r['leafs']} leafs {r['path']}")
return 0
# Never DROP FTS/VECTOR (ghost catalog). Upsert while indexes exist is OK;
# ensure_indexes only CREATEs when missing.
db, conn = connect(ROOT / "var" / "kb.lbug", read_only=False)
init_schema(conn)
embed = embedder()
rows: list[dict] = []
for repo in repos:
if not repo.exists():
continue
reached, written = import_repo(conn, repo, embed, a.limit, a.since)
total_commits += reached
rows.append({"repo": repo_name(repo), "commits": reached, "written": written})
ensure_indexes(conn)
conn.close()
db.close()
if a.json:
print(json.dumps(rows, indent=2))
else:
for r in rows:
print(f"imported {r['commits']:>5} commits -> {r['written']} leafs {r['repo']}")
print(f"total: {total_commits} commits")
return 0
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))
-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
+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
+14 -12
View File
@@ -19,10 +19,10 @@ 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, connect, ensure_indexes, init_schema, upsert_leaf,
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
@@ -99,6 +99,11 @@ def main(argv: list[str]) -> int:
p = argparse.ArgumentParser(description="build the 2dph brain index") p = argparse.ArgumentParser(description="build the 2dph brain index")
p.add_argument("--corpus", action="append", help="extra markdown dir/file to index (may repeat)") p.add_argument("--corpus", action="append", help="extra markdown dir/file to index (may repeat)")
p.add_argument("--rebuild", action="store_true", help="fresh db + indexes") p.add_argument("--rebuild", action="store_true", help="fresh db + indexes")
p.add_argument(
"--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)
@@ -116,26 +121,23 @@ def main(argv: list[str]) -> int:
db, conn = connect(DB_PATH, read_only=False) db, conn = connect(DB_PATH, read_only=False)
init_schema(conn) init_schema(conn)
if not (a.rebuild or _already_indexed(conn)): # Never DROP FTS/VECTOR (ghost catalog). Write leafs, then ensure indexes
create_fts_and_vector(conn, force=True) # unless --skip-indexes (seed facts first — MERGE under live FTS corrupts it).
# --rebuild already deleted kb.lbug above, so CREATE runs on a clean DB.
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)) 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 = {"indexed": done, "corpus_total": total, **{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} leafs; db total {s['total']}") print(json.dumps(result, indent=2) if a.json else f"indexed {done}/{total} leafs; 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:]))
+29 -67
View File
@@ -1,72 +1,34 @@
#!/usr/bin/env python3 #!/usr/bin/env bash
"""kb/search - deduction search over the 2dph brain. # bin/kb/search - Go deduction search over the brain (model served by daemon).
# Builds the kbsearch binary on first run / when source changes, then execs it.
set -euo pipefail
bin/kb/search "query" # hybrid facts+info, YAML out KB="$(cd "$(dirname "$0")/../.." && pwd)"
bin/kb/search "query" --root facts # confirmed facts only BIN="$KB/var/bin/kbsearch"
bin/kb/search "query" --hop 1 # follow graph edges after hitting SRC="$KB/bin/kbsearch"
bin/kb/search "query" --json | yq '.'
bin/kb/search "query" -n 5 # more results
Deduction order: facts root first (confirmed answers with evidence links), mkdir -p "$KB/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 # Rebuild if binary missing or any .go source newer
import sys need_build=0
from pathlib import Path if [ ! -x "$BIN" ]; then
need_build=1
else
# Check if any .go in kbsearch is newer than binary
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
ROOT = Path(__file__).resolve().parents[2] if [ "$need_build" -eq 1 ]; then
sys.path.insert(0, str(ROOT / "tools")) echo "Building kbsearch..." >&2
(cd "$SRC" && \
CGO_CFLAGS="-I$KB/lib-ladybug" \
CGO_LDFLAGS="-L$KB/lib-ladybug -Wl,-rpath,$KB/lib-ladybug" \
go build -tags system_ladybug -o "$BIN" .) || exit 1
fi
from kblib import connect, hybrid_search, init_schema, open_readonly, query_fts # noqa: E402 exec "$BIN" "$@"
from yamlout import to_yaml # noqa: E402
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 - re-index the 2dph brain when corpus files change.
//
// 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:])
}
+88
View File
@@ -0,0 +1,88 @@
// Brain connection management using go-ladybug.
package main
import (
"fmt"
"os"
"path/filepath"
lbug "github.com/LadybugDB/go-ladybug"
)
var (
db *lbug.Database
conn *lbug.Connection
)
func repoRoot() string {
// Try KB_ROOT env, then walk up from binary
if v := os.Getenv("KB_ROOT"); v != "" {
return v
}
self, err := os.Executable()
if err == nil {
dir := filepath.Dir(self)
for i := 0; i < 5; i++ {
if _, err := os.Stat(filepath.Join(dir, "var")); err == nil {
return dir
}
if _, err := os.Stat(filepath.Join(dir, ".git")); err == nil {
return dir
}
parent := filepath.Dir(dir)
if parent == dir {
break
}
dir = parent
}
}
return "."
}
func dbPath() string {
return filepath.Join(repoRoot(), "var", "kb.lbug")
}
func openBrain() error {
return openWithOpts(2, eps())
}
func openWithOpts(allow int, epsv string) error {
cfg := lbug.DefaultSystemConfig()
cfg.MaxNumThreads = 8
cfg.BufferPoolSize = 1 << 30 // 1GB
var err error
db, err = lbug.OpenDatabase(dbPath(), cfg)
if err != nil {
return fmt.Errorf("OpenDatabase: %w", err)
}
if epsv != "" {
if _, err := conn.Query("SET STREAM_SANDBOX = '" + epsv + "'"); err != nil {
return err
}
}
conn, err = lbug.OpenConnection(db)
if err != nil {
return fmt.Errorf("OpenConnection: %w", err)
}
if _, err := conn.Query("LOAD EXTENSION FTS"); err != nil {
return fmt.Errorf("LOAD EXTENSION FTS: %w", err)
}
if _, err := conn.Query("LOAD EXTENSION VECTOR"); err != nil {
return fmt.Errorf("LOAD EXTENSION VECTOR: %w", err)
}
return nil
}
func closeBrain() {
if conn != nil {
conn.Close()
conn = nil
}
if db != nil {
db.Close()
db = nil
}
}
+23
View File
@@ -0,0 +1,23 @@
module github.com/eSlider/2dph/bin/kbsearch
go 1.26.0
require (
github.com/LadybugDB/go-ladybug v0.17.0
github.com/chewxy/math32 v1.11.2
github.com/daulet/tokenizers v1.27.0
)
require (
github.com/apache/arrow-go/v18 v18.6.0 // indirect
github.com/goccy/go-json v0.10.6 // indirect
github.com/google/flatbuffers v25.12.19+incompatible // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/klauspost/compress v1.18.5 // indirect
github.com/klauspost/cpuid/v2 v2.3.0 // indirect
github.com/pierrec/lz4/v4 v4.1.26 // indirect
github.com/shopspring/decimal v1.4.0 // indirect
github.com/zeebo/xxh3 v1.1.0 // indirect
golang.org/x/exp v0.0.0-20260112195511-716be5621a96 // indirect
golang.org/x/sys v0.43.0 // indirect
)
+44
View File
@@ -0,0 +1,44 @@
github.com/LadybugDB/go-ladybug v0.17.0 h1:RXDbkBjrbRmLdEbhGl4CLOIEzSt09gbP0n9UbKDEfwI=
github.com/LadybugDB/go-ladybug v0.17.0/go.mod h1:GeIXmE8XyF5TFS94NAuTag7vgCC+no/HTBMRA6Rd5Cs=
github.com/andybalholm/brotli v1.2.1 h1:R+f5xP285VArJDRgowrfb9DqL18yVK0gKAW/F+eTWro=
github.com/andybalholm/brotli v1.2.1/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY=
github.com/apache/arrow-go/v18 v18.6.0 h1:GX/Jyd3R7mCLiECAwY9FWbbaYblie2WXBSz4Sw8fNpM=
github.com/apache/arrow-go/v18 v18.6.0/go.mod h1:gm3MiPpY82fLYK5VKPB3WoJbsiLVDfT7flD5/vHReKw=
github.com/apache/thrift v0.22.0 h1:r7mTJdj51TMDe6RtcmNdQxgn9XcyfGDOzegMDRg47uc=
github.com/apache/thrift v0.22.0/go.mod h1:1e7J/O1Ae6ZQMTYdy9xa3w9k+XHWPfRvdPyJeynQ+/g=
github.com/chewxy/math32 v1.11.2 h1:IufN08Zwr1NKuWfY+4Tz55BcwKmyKKNdOP7KtumehnM=
github.com/chewxy/math32 v1.11.2/go.mod h1:dOB2rcuFrCn6UHrze36WSLVPKtzPMRAQvBvUwkSsLqs=
github.com/daulet/tokenizers v1.27.0 h1:MmFYAEDFz69s/nNQfHg59DWqHz3v94m99kEZ/JbL+s4=
github.com/daulet/tokenizers v1.27.0/go.mod h1:YjFY1o1HGMyWkQgbXJDghhvke/yFDp2vGdIO2hYs4MQ=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/goccy/go-json v0.10.6 h1:p8HrPJzOakx/mn/bQtjgNjdTcN+/S6FcG2CTtQOrHVU=
github.com/goccy/go-json v0.10.6/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M=
github.com/google/flatbuffers v25.12.19+incompatible h1:haMV2JRRJCe1998HeW/p0X9UaMTK6SDo0ffLn2+DbLs=
github.com/google/flatbuffers v25.12.19+incompatible/go.mod h1:1AeVuKshWv4vARoZatz6mlQ0JxURH0Kv5+zNeJKJCa8=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE=
github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y=
github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0=
github.com/pierrec/lz4/v4 v4.1.26 h1:GrpZw1gZttORinvzBdXPUXATeqlJjqUG/D87TKMnhjY=
github.com/pierrec/lz4/v4 v4.1.26/go.mod h1:EoQMVJgeeEOMsCqCzqFm2O0cJvljX2nGZjcRIPL34O4=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/shopspring/decimal v1.4.0 h1:bxl37RwXBklmTi0C79JfXCEBD1cqqHt0bbgBAGFp81k=
github.com/shopspring/decimal v1.4.0/go.mod h1:gawqmDU56v4yIKSwfBSFip1HdCCXN8/+DMd9qYNcwME=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ=
github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0=
github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs=
github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s=
golang.org/x/exp v0.0.0-20260112195511-716be5621a96 h1:Z/6YuSHTLOHfNFdb8zVZomZr7cqNgTJvA8+Qz75D8gU=
golang.org/x/exp v0.0.0-20260112195511-716be5621a96/go.mod h1:nzimsREAkjBCIEFtHiYkrJyT+2uy9YZJB7H1k68CXZU=
golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI=
golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+44
View File
@@ -0,0 +1,44 @@
// bin/kbsearch - the Go implementation of bin/kb/search (nested module so the
// root `go test ./...` and CI never compile it against native ladyships).
//
// Usage (built/run by ./bin/kb/search):
//
// kbsearch "query" [--root facts|info] [--repo P] [-n N] [--json]
// kbsearch serve [port] start the embedding daemon
// kbsearch --list-model print the resolved model dir
//
// The potion-multilingual model is loaded only in `serve`; a CLI reuses the
// daemon over localhost HTTP (falling back to in-process embedding).
package main
import (
"fmt"
"log"
"os"
"strconv"
)
func main() {
if len(os.Args) > 1 && os.Args[1] == "serve" {
port := 17830
if len(os.Args) > 2 {
if p, err := strconv.Atoi(os.Args[2]); err == nil {
port = p
}
}
if err := serve(port); err != nil {
log.Fatalf("kbsearch serve: %v", err)
}
return
}
if len(os.Args) > 1 && os.Args[1] == "--list-model" {
dir, err := modelDir()
if err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
fmt.Println(dir)
return
}
os.Exit(runSearch(os.Args[1:]))
}
+236
View File
@@ -0,0 +1,236 @@
// StaticModel wraps the potion-multilingual-128m embedding model.
//
// Mirrors model2vec.StaticModel: tokenizer (daulet Unigram) + safetensors matrix.
// Embed(text) applies the same preprocessing: median_token_length pre-truncation,
// add_special_tokens=false, drop unk (id=1), truncate to 512, mean pool, L2 normalize +1e-32.
package main
import (
"encoding/json"
"fmt"
"io"
"math"
"os"
"path/filepath"
"sort"
"github.com/chewxy/math32"
"github.com/daulet/tokenizers"
)
type StaticModel struct {
tok *tokenizers.Tokenizer
mat []float32 // row-major: vocab_size x 128
medianLen int
vocabSize int
dim int
}
func loadModel() (*StaticModel, error) {
dir, err := modelDir()
if err != nil {
return nil, err
}
tok, err := tokenizers.FromFile(filepath.Join(dir, "tokenizer.json"))
if err != nil {
return nil, fmt.Errorf("tokenizer: %w", err)
}
mat, vocabSize, dim, err := loadMatrix(filepath.Join(dir, "model.safetensors"))
if err != nil {
return nil, fmt.Errorf("safetensors: %w", err)
}
median := medianTokenLength(filepath.Join(dir, "tokenizer.json"))
return &StaticModel{
tok: tok,
mat: mat,
vocabSize: vocabSize,
dim: dim,
medianLen: median,
}, nil
}
func (m *StaticModel) Close() error {
if m.tok != nil {
m.tok.Close()
m.tok = nil
}
return nil
}
func (m *StaticModel) Embed(text string) ([]float64, error) {
const maxLen = 512
if m.medianLen > 0 {
maxChars := maxLen * m.medianLen
runes := []rune(text)
if len(runes) > maxChars {
text = string(runes[:maxChars])
}
}
ids, _, err := m.tok.EncodeErr(text, false)
if err != nil {
return nil, fmt.Errorf("encode: %w", err)
}
filtered := make([]uint32, 0, len(ids))
for _, id := range ids {
if id != 1 {
filtered = append(filtered, id)
}
if len(filtered) >= maxLen {
break
}
}
if len(filtered) == 0 {
return make([]float64, m.dim), nil
}
acc := make([]float32, m.dim)
for _, id := range filtered {
if int(id) >= m.vocabSize {
continue
}
off := int(id) * m.dim
for d := 0; d < m.dim; d++ {
acc[d] += m.mat[off+d]
}
}
inv := 1.0 / float32(len(filtered))
for d := 0; d < m.dim; d++ {
acc[d] *= inv
}
var norm float32
for d := 0; d < m.dim; d++ {
norm += acc[d] * acc[d]
}
norm = math32.Sqrt(norm) + 1e-32
for d := 0; d < m.dim; d++ {
acc[d] /= norm
}
out := make([]float64, m.dim)
for d := 0; d < m.dim; d++ {
out[d] = float64(acc[d])
}
return out, nil
}
func medianTokenLength(tokenizerPath string) int {
data, err := os.ReadFile(tokenizerPath)
if err != nil {
return 0
}
var parsed struct {
Model struct {
Vocab [][]json.RawMessage `json:"vocab"`
} `json:"model"`
}
if err := json.Unmarshal(data, &parsed); err != nil {
return 0
}
vocab := parsed.Model.Vocab
if len(vocab) == 0 {
return 0
}
lengths := make([]int, 0, len(vocab))
for _, pair := range vocab {
if len(pair) < 1 {
continue
}
var tok string
if err := json.Unmarshal(pair[0], &tok); err != nil {
continue
}
lengths = append(lengths, len([]rune(tok)))
}
if len(lengths) == 0 {
return 0
}
sort.Ints(lengths)
return lengths[len(lengths)/2]
}
func loadMatrix(path string) ([]float32, int, int, error) {
f, err := os.Open(path)
if err != nil {
return nil, 0, 0, err
}
defer f.Close()
var hdrLen uint64
if err := binaryRead(f, &hdrLen); err != nil {
return nil, 0, 0, err
}
hdrBytes := make([]byte, hdrLen)
if _, err := io.ReadFull(f, hdrBytes); err != nil {
return nil, 0, 0, err
}
var hdr struct {
Embeddings struct {
Dtype string `json:"dtype"`
Shape []int `json:"shape"`
Offset []uint64 `json:"data_offsets"`
} `json:"embeddings"`
}
if err := json.Unmarshal(hdrBytes, &hdr); err != nil {
return nil, 0, 0, err
}
if hdr.Embeddings.Dtype != "F32" {
return nil, 0, 0, fmt.Errorf("unsupported dtype %s", hdr.Embeddings.Dtype)
}
if len(hdr.Embeddings.Shape) != 2 {
return nil, 0, 0, fmt.Errorf("expected 2D shape, got %v", hdr.Embeddings.Shape)
}
vocabSize := hdr.Embeddings.Shape[0]
dim := hdr.Embeddings.Shape[1]
if len(hdr.Embeddings.Offset) != 2 {
return nil, 0, 0, fmt.Errorf("bad offsets")
}
start := hdr.Embeddings.Offset[0]
end := hdr.Embeddings.Offset[1]
size := end - start
if size != uint64(vocabSize*dim*4) {
return nil, 0, 0, fmt.Errorf("size mismatch")
}
if _, err := f.Seek(int64(8+hdrLen+start), io.SeekStart); err != nil {
return nil, 0, 0, err
}
buf := make([]byte, size)
if _, err := io.ReadFull(f, buf); err != nil {
return nil, 0, 0, err
}
mat := make([]float32, vocabSize*dim)
for i := 0; i < len(mat); i++ {
off := i * 4
mat[i] = math.Float32frombits(
uint32(buf[off]) |
uint32(buf[off+1])<<8 |
uint32(buf[off+2])<<16 |
uint32(buf[off+3])<<24,
)
}
return mat, vocabSize, dim, nil
}
func binaryRead(r io.Reader, v any) error {
switch p := v.(type) {
case *uint64:
var b [8]byte
if _, err := io.ReadFull(r, b[:]); err != nil {
return err
}
*p = uint64(b[0]) | uint64(b[1])<<8 | uint64(b[2])<<16 | uint64(b[3])<<24 |
uint64(b[4])<<32 | uint64(b[5])<<40 | uint64(b[6])<<48 | uint64(b[7])<<56
}
return nil
}
+60
View File
@@ -0,0 +1,60 @@
// modelDir returns the resolved potion-multilingual-128m model directory.
package main
import (
"fmt"
"os"
"path/filepath"
"strings"
)
func modelDir() (string, error) {
// 1. Explicit env
if v := os.Getenv("KBSEARCH_MODEL"); v != "" {
return v, nil
}
// 2. Next to the binary (dev or installed)
self, err := os.Executable()
if err == nil {
if dir, err := filepath.EvalSymlinks(filepath.Dir(self)); err == nil {
cand := filepath.Join(dir, "potion-multilingual-128m")
if st, err := os.Stat(cand); err == nil && st.IsDir() {
return cand, nil
}
}
}
// 3. Repo root lib/ (where other scripts expect it)
if v := os.Getenv("KB_ROOT"); v != "" {
cand := filepath.Join(v, "lib", "potion-multilingual-128m")
if st, err := os.Stat(cand); err == nil && st.IsDir() {
return cand, nil
}
}
// 4. HF cache (new layout: models--*/snapshots/*)
if v := os.Getenv("HF_HOME"); v != "" {
base := filepath.Join(v, "hub")
if entries, err := os.ReadDir(base); err == nil {
for _, e := range entries {
if strings.HasPrefix(e.Name(), "models--") {
snapDir := filepath.Join(base, e.Name(), "snapshots")
if snaps, err := os.ReadDir(snapDir); err == nil {
for _, s := range snaps {
cand := filepath.Join(snapDir, s.Name())
if st, _ := os.Stat(cand); st != nil && st.IsDir() {
return cand, nil
}
}
}
}
}
}
}
// 5. Legacy HF cache (symlinked model dir)
if v := os.Getenv("HF_HOME"); v != "" {
cand := filepath.Join(v, "potion-multilingual-128m")
if st, err := os.Stat(cand); err == nil && st.IsDir() {
return cand, nil
}
}
return "", fmt.Errorf("model not found (set KBSEARCH_MODEL or KB_ROOT, or download to HF cache)")
}
+465
View File
@@ -0,0 +1,465 @@
// Hybrid FTS + vector search implementation, plus daemon client/server.
package main
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"log"
"net"
"net/http"
"os"
"os/exec"
"path/filepath"
"sort"
"strconv"
"strings"
"time"
lbug "github.com/LadybugDB/go-ladybug"
)
const defaultPort = 17830
const daemonPath = "/embed"
const healthPath = "/health"
func runSearch(args []string) int {
// Manual flag parsing to allow flags after query (like Python argparse)
root := ""
repo := ""
limit := 20
jsonOut := false
listModel := false
var queryArgs []string
for i := 0; i < len(args); i++ {
switch args[i] {
case "--root":
if i+1 < len(args) {
root = args[i+1]
i++
}
case "--repo":
if i+1 < len(args) {
repo = args[i+1]
i++
}
case "-n":
if i+1 < len(args) {
if n, err := strconv.Atoi(args[i+1]); err == nil {
limit = n
}
i++
}
case "--json":
jsonOut = true
case "--list-model":
listModel = true
default:
if !strings.HasPrefix(args[i], "-") {
queryArgs = append(queryArgs, args[i])
}
}
}
if listModel {
dir, err := modelDir()
if err != nil {
fmt.Fprintln(os.Stderr, err)
return 1
}
fmt.Println(dir)
return 0
}
query := strings.TrimSpace(strings.Join(queryArgs, " "))
if query == "" {
fmt.Fprintln(os.Stderr, "usage: kbsearch \"query\" [--root facts|info] [--repo REPO] [-n N] [--json]")
return 1
}
if err := openBrain(); err != nil {
fmt.Fprintf(os.Stderr, "open brain: %v\n", err)
return 1
}
defer closeBrain()
emb, err := embedQuery(query)
if err != nil {
fmt.Fprintf(os.Stderr, "embed: %v\n", err)
return 1
}
fts, err := queryFTS(query, limit*3)
if err != nil {
fmt.Fprintf(os.Stderr, "fts: %v\n", err)
return 1
}
var vec []Hit
if vec, err = queryVector(emb, limit*3); err != nil {
fmt.Fprintf(os.Stderr, "vec: %v\n", err)
}
results := hybrid(fts, vec, limit)
if root != "" {
results = filterRoot(results, root)
}
if repo != "" {
results = filterRepo(results, repo)
}
if len(results) > limit {
results = results[:limit]
}
for i := range results {
if results[i].Text != "" {
runes := []rune(results[i].Text)
if len(runes) > 280 {
runes = runes[:280]
}
results[i].Snippet = string(runes)
}
}
out := Dict{
{"query", query},
{"root_filter", root},
{"count", len(results)},
{"results", resultsToDicts(results)},
}
if jsonOut {
enc := json.NewEncoder(os.Stdout)
enc.SetIndent("", " ")
enc.SetEscapeHTML(false)
return b2i(enc.Encode(toJSONOut(results, query, root)))
}
fmt.Print(toYAML(out, 0))
return 0
}
func b2i(err error) int {
if err != nil {
return 1
}
return 0
}
func queryFTS(text string, limit int) ([]Hit, error) {
stmt, err := conn.Prepare(
"CALL QUERY_FTS_INDEX('Leaf', 'id', $q) " +
"RETURN node.id, node.text, node.root, node.source, score ORDER BY score LIMIT $n",
)
if err != nil {
return nil, err
}
defer stmt.Close()
res, err := conn.Execute(stmt, map[string]any{"q": text, "n": limit})
if err != nil {
return nil, err
}
return rowsToHits(res)
}
func queryVector(emb []float64, limit int) ([]Hit, error) {
embList := make([]any, len(emb))
for i, v := range emb {
embList[i] = v
}
stmt, err := conn.Prepare(
"CALL QUERY_VECTOR_INDEX('Leaf', 'Leaf_vec', $q, $n) " +
"RETURN node.id, node.text, node.root, node.source, distance ORDER BY distance LIMIT $n",
)
if err != nil {
return nil, err
}
defer stmt.Close()
res, err := conn.Execute(stmt, map[string]any{"q": embList, "n": limit})
if err != nil {
return nil, err
}
hits, err := rowsToHits(res)
if err != nil {
return nil, err
}
for i := range hits {
hits[i].Score = 1.0 - hits[i].Score
}
return hits, nil
}
func rowsToHits(res *lbug.QueryResult) ([]Hit, error) {
var hits []Hit
for res.HasNext() {
row, err := res.Next()
if err != nil {
return nil, err
}
vals, err := row.GetAsSlice()
if err != nil || len(vals) < 5 {
continue
}
id := fmt.Sprint(vals[0])
text := fmt.Sprint(vals[1])
root := fmt.Sprint(vals[2])
source := fmt.Sprint(vals[3])
score := float64(vals[4].(float64))
hits = append(hits, Hit{ID: id, Text: text, Root: root, Source: source, Score: score})
}
return hits, nil
}
// JSON output types
type jsonOut struct {
Query string `json:"query"`
RootFilter string `json:"root_filter"`
Count int `json:"count"`
Results []jsonHit `json:"results"`
}
type jsonHit struct {
ID string `json:"id"`
Text string `json:"text"`
Root string `json:"root"`
Score float64 `json:"score"`
Snippet string `json:"snippet,omitempty"`
}
func toJSONOut(hits []Hit, query, rootFilter string) *jsonOut {
out := make([]jsonHit, len(hits))
for i, h := range hits {
out[i] = jsonHit{
ID: h.ID,
Text: h.Text,
Root: h.Root,
Score: h.Score,
Snippet: h.Snippet,
}
}
return &jsonOut{
Query: query,
RootFilter: rootFilter,
Count: len(hits),
Results: out,
}
}
func hybrid(fts, vec []Hit, limit int) []Hit {
byID := make(map[string]Hit)
rrf := make(map[string]float64)
for rank, h := range fts {
byID[h.ID] = h
rrf[h.ID] += 1.0 / (60 + float64(rank+1))
}
for rank, h := range vec {
if _, ok := byID[h.ID]; !ok {
byID[h.ID] = h
} else {
existing := byID[h.ID]
if existing.Score == 0 {
existing.Score = h.Score
byID[h.ID] = existing
}
}
rrf[h.ID] += 1.0 / (60 + float64(rank+1))
}
type scored struct {
id string
rrf float64
}
var scoredList []scored
for id, v := range rrf {
scoredList = append(scoredList, scored{id, v})
}
sort.Slice(scoredList, func(i, j int) bool {
return scoredList[i].rrf > scoredList[j].rrf
})
var out []Hit
for i, s := range scoredList {
if i >= limit {
break
}
h := byID[s.id]
out = append(out, h)
}
return out
}
func filterRoot(hits []Hit, root string) []Hit {
var out []Hit
for _, h := range hits {
if h.Root == root {
out = append(out, h)
}
}
return out
}
func filterRepo(hits []Hit, repo string) []Hit {
var out []Hit
for _, h := range hits {
if strings.Contains(h.Source, repo) {
out = append(out, h)
}
}
return out
}
func resultsToDicts(hits []Hit) []any {
out := make([]any, len(hits))
for i, h := range hits {
d := Dict{
{"id", h.ID},
{"text", h.Text},
{"root", h.Root},
{"score", h.Score},
}
if h.Snippet != "" {
d = append(d, KV{"snippet", h.Snippet})
}
out[i] = d
}
return out
}
// --- Daemon server ---
func serve(port int) error {
model, err := loadModel()
if err != nil {
return fmt.Errorf("load model: %w", err)
}
defer model.Close()
mux := http.NewServeMux()
mux.HandleFunc(healthPath, func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
})
mux.HandleFunc(daemonPath, func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
w.WriteHeader(http.StatusMethodNotAllowed)
return
}
var req struct {
Text string `json:"text"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
w.WriteHeader(http.StatusBadRequest)
return
}
vec, err := model.Embed(req.Text)
if err != nil {
w.WriteHeader(http.StatusInternalServerError)
json.NewEncoder(w).Encode(map[string]string{"error": err.Error()})
return
}
json.NewEncoder(w).Encode(map[string]any{"vector": vec})
})
addr := fmt.Sprintf("127.0.0.1:%d", port)
log.Printf("kbsearch daemon listening on %s", addr)
return http.ListenAndServe(addr, mux)
}
// --- Daemon client ---
var daemonClient = &http.Client{
Timeout: 5 * time.Second,
Transport: &http.Transport{
DialContext: (&net.Dialer{Timeout: 2 * time.Second}).DialContext,
},
}
func embedQuery(text string) ([]float64, error) {
port := defaultPort
if envPort := os.Getenv("KBSEARCH_PORT"); envPort != "" {
if p, err := strconv.Atoi(envPort); err == nil {
port = p
}
}
emb, err := tryDaemon(text, port)
if err == nil {
return emb, nil
}
model, err := loadModel()
if err != nil {
return nil, fmt.Errorf("fallback load model: %w", err)
}
defer model.Close()
return model.Embed(text)
}
func tryDaemon(text string, port int) ([]float64, error) {
url := fmt.Sprintf("http://127.0.0.1:%d%s", port, daemonPath)
payload := map[string]string{"text": text}
body, _ := json.Marshal(payload)
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, "POST", url, bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
resp, err := daemonClient.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("daemon HTTP %d", resp.StatusCode)
}
var r struct {
Vector []float64 `json:"vector"`
Error string `json:"error"`
}
if err := json.NewDecoder(resp.Body).Decode(&r); err != nil {
return nil, err
}
if r.Error != "" {
return nil, errors.New(r.Error)
}
return r.Vector, nil
}
func ensureDaemon(port int) error {
url := fmt.Sprintf("http://127.0.0.1:%d%s", port, healthPath)
ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, "GET", url, nil)
if resp, err := daemonClient.Do(req); err == nil {
resp.Body.Close()
if resp.StatusCode == http.StatusOK {
return nil
}
}
self, err := os.Executable()
if err != nil {
return err
}
cmd := exec.Command(self, "serve", strconv.Itoa(port))
cmd.Dir, _ = filepath.Split(self)
cmd.Stdout = nil
cmd.Stderr = nil
if err := cmd.Start(); err != nil {
return err
}
for i := 0; i < 40; i++ {
time.Sleep(250 * time.Millisecond)
req, _ := http.NewRequestWithContext(context.Background(), "GET", url, nil)
if resp, err := daemonClient.Do(req); err == nil {
resp.Body.Close()
if resp.StatusCode == http.StatusOK {
return nil
}
}
}
return fmt.Errorf("daemon failed to start on port %d", port)
}
+16
View File
@@ -0,0 +1,16 @@
// Common types and helpers for kbsearch.
package main
import "os"
func eps() string { return os.Getenv("KBTEST_EPS") }
// Hit is one search result, mirroring the python script's dict shape.
type Hit struct {
ID string `json:"id"`
Text string `json:"text"`
Root string `json:"root"`
Source string `json:"-"` // for repo filtering, not in output
Score float64 `json:"score"`
Snippet string `json:"snippet,omitempty"`
}
+107
View File
@@ -0,0 +1,107 @@
// YAML emitter ported from bin/kb/yamlout.py — preserves insertion order.
package main
import (
"fmt"
"strconv"
"strings"
)
// KV is an ordered key-value pair for maps.
type KV struct {
K string
V any
}
// Dict is an ordered map (slice of KV).
type Dict []KV
func toYAML(node any, indent int) string {
pad := strings.Repeat(" ", indent)
switch n := node.(type) {
case Dict:
if len(n) == 0 {
return pad + "{}\n"
}
var b strings.Builder
for _, kv := range n {
switch nv := kv.V.(type) {
case Dict:
if len(nv) == 0 {
b.WriteString(pad + kv.K + ": {}\n")
} else {
b.WriteString(pad + kv.K + ":\n" + toYAML(nv, indent+1))
}
case []any:
if len(nv) == 0 {
b.WriteString(pad + kv.K + ": []\n")
} else {
b.WriteString(pad + kv.K + ":\n" + toYAML(nv, indent+1))
}
default:
b.WriteString(pad + kv.K + ": " + scalar(nv) + "\n")
}
}
return b.String()
case []any:
if len(n) == 0 {
return pad + "[]\n"
}
var b strings.Builder
for _, item := range n {
if d, ok := item.(Dict); ok {
b.WriteString(pad + "-\n" + toYAML(d, indent+1))
} else {
b.WriteString(pad + "- " + scalar(item) + "\n")
}
}
return b.String()
default:
return pad + scalar(node) + "\n"
}
}
func scalar(v any) string {
switch t := v.(type) {
case nil:
return "null"
case bool:
if t {
return "true"
}
return "false"
case int:
return strconv.Itoa(t)
case int64:
return strconv.FormatInt(t, 10)
case float64:
return fmtFloat(t)
case float32:
return fmtFloat(float64(t))
case string:
return quoteIfNeeded(t)
default:
// fallback
return fmt.Sprintf("%v", v)
}
}
func fmtFloat(f float64) string {
s := strconv.FormatFloat(f, 'g', -1, 64)
if !strings.ContainsAny(s, ".eE") {
s += ".0"
}
return s
}
func quoteIfNeeded(s string) string {
if strings.Contains(s, "\n") {
return strconv.Quote(s)
}
if s == "" || strings.ContainsAny(s, ":#'\"[]{}&*!|>%@`") || s != strings.TrimSpace(s) {
return strconv.Quote(s)
}
return s
}
+450
View File
@@ -0,0 +1,450 @@
#!/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 scanned PDFs/images via docling
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/mail/index_mail): conversion can crash in
native docling and must not leave the brain DB mid-transaction.
Requires ONLYOFFICE_URL/USER/PASS in .env (or env). Idempotent: a message
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,
html_to_markdown,
is_convertible,
normalize_markdown,
subject_to_filename,
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 _convert_pdf(path, ocr)
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_pdf(path: Path, ocr: bool) -> str:
"""Convert one PDF to markdown.
Fast path: poppler's pdftotext (-layout) extracts exact text from
born-digital PDFs in ~15ms vs docling's 1-3s. Only textless PDFs (scanned
pages, layout-heavy) fall back to docling, which runs isolated in a
subprocess because its native onnx/RT-DETR has segfaulted the main process.
"""
text = _pdf_fast_text(path)
if ocr or text is None or not text.strip():
return _convert_pdf_docling(path, ocr)
return normalize_markdown(text)
def _pdf_fast_text(path: Path) -> str | None:
"""pdftotext -layout; None when poppler is unavailable (or the PDF has no text layer)."""
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 _convert_pdf_docling(path: Path, ocr: bool) -> str:
try:
proc = subprocess.run(
[sys.executable, os.path.abspath(__file__), "--pdf-worker", str(path),
"--ocr" if ocr else "--no-ocr"],
capture_output=True, text=True, timeout=600)
except subprocess.TimeoutExpired:
return "\n<!-- pdf conversion timed out -->\n"
if proc.returncode != 0:
tail = proc.stderr.strip().splitlines()[-3:]
return f"\n<!-- pdf conversion failed: {proc.returncode}: {' | '.join(tail)} -->\n"
return proc.stdout
def _pdf_worker(path: Path, ocr: bool) -> None:
"""docling worker entry: prints converted markdown on stdout, exits non-zero on error."""
try:
from docling.document_converter import DocumentConverter, PdfFormatOption
from docling.datamodel.pipeline_options import PdfPipelineOptions
opts = PdfPipelineOptions()
opts.do_ocr = bool(ocr)
opts.do_table_structure = True
conv = DocumentConverter(format_options={"pdf": PdfFormatOption(pipeline_options=opts)})
res = conv.convert(str(path))
sys.stdout.write(normalize_markdown(res.document.export_to_markdown()))
sys.exit(0)
except Exception as e:
# errors/stacktraces to stderr; the caller only reports a one-liner
print(f"pdf-worker: {e}", file=sys.stderr)
import traceback
traceback.print_exc(file=sys.stderr)
sys.exit(1)
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 scanned PDFs/images via docling")
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")
p.add_argument("--pdf-worker", default="", help=argparse.SUPPRESS)
p.add_argument("--no-ocr", action="store_true", help=argparse.SUPPRESS)
a = p.parse_args(argv)
if a.pdf_worker:
_pdf_worker(Path(a.pdf_worker), ocr=not a.no_ocr)
return 0
conf = load_env()
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:
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:]))
+136
View File
@@ -0,0 +1,136 @@
#!/usr/bin/env python3
"""mail/index_mail - rebuild the brain with every markdown under var/mail.
Ladybug corrupts its WAL when brand-new leafs are bulk-inserted while the
FTS/VECTOR indexes already exist, so indexing ALWAYS runs as a fresh rebuild
(repo corpus + var/mail), matching the proven-safe `kb/index --rebuild` path.
Conversion and indexing stay separate: conversion can crash in native docling
and must not leave the brain DB mid-transaction.
bin/mail/index_mail rebuild the index incl. all mail
bin/mail/index_mail --dry-run count without writing
bin/mail/index_mail --limit N cap messages included
bin/mail/index_mail --since D only messages dated >= D (YYYY-MM-DD)
"""
from __future__ import annotations
import argparse
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 DB_PATH, VAR, connect, ensure_indexes, init_schema, stats, upsert_leaf # noqa: E402
from mdleaves import read_markdown, to_all, walk_markdown # noqa: E402
def msg_date(md: Path) -> str:
j = md.parent / "message.json"
try:
d = json.loads(j.read_text(encoding="utf-8"))
return (d.get("receivedDate") or d.get("receivedAt") or "")[:10]
except Exception:
return ""
def mail_leafs(limit: int, since: str, repo: str = "ooMail") -> list[dict]:
root = ROOT / "var" / "mail"
mds = sorted(root.rglob("message.md"))
if since:
mds = [m for m in mds if msg_date(m) >= since]
if limit:
mds = mds[:limit]
leafs: list[dict] = []
for md in mds:
files = [md] + sorted((md.parent / "attachments").glob("*.md"))
for f in files:
if not f.exists():
continue
for lf in to_all(read_markdown(f), f, repo=repo):
lf["source"] = f"ooMail:{md.parent.name}:{f.name}"
lf["how"] = "mail/import"
leafs.append(lf)
return leafs
def main(argv: list[str]) -> int:
p = argparse.ArgumentParser(description="rebuild the brain incl. all mail")
p.add_argument("--dry-run", action="store_true", help="count only, write nothing")
p.add_argument("--limit", type=int, default=0, help="cap messages included")
p.add_argument("--since", default="", help="only messages dated >= YYYY-MM-DD")
p.add_argument("--json", action="store_true")
a = p.parse_args(argv)
mail = mail_leafs(a.limit, a.since)
if a.dry_run:
print(f"mail/index_mail: {len(mail)} mail leafs would be indexed")
return 0
# Fresh rebuild: delete DB, index repo corpus + mail, create indexes once
# at the end. Never insert into an already-indexed DB (WAL corruption).
VAR.mkdir(exist_ok=True)
if DB_PATH.exists():
DB_PATH.unlink()
corpus = _load_corpus()
leafs = corpus + mail
db, conn = connect(DB_PATH, read_only=False)
init_schema(conn)
embed = _embedder()
done, total = _index_leafs(conn, leafs, embed)
ensure_indexes(conn)
s = stats(conn)
conn.close()
db.close()
result = {"indexed": done, "corpus_total": total, "mail_leafs": len(mail),
**{k: v for k, v in s.items() if k in ("total", "by_root")}}
print(json.dumps(result, indent=2) if a.json else
f"mail/index_mail: indexed {done}/{total} leafs (mail={len(mail)}); db total {s['total']}")
return 0
CORPUS_DEFAULTS = ["README.md", "PLAN.md", "AGENTS.md", "docs", "skills"]
def _load_corpus() -> list[dict]:
files: list[Path] = []
for entry in CORPUS_DEFAULTS:
p = ROOT / entry
if p.is_file():
files.append(p)
elif p.is_dir():
files.extend(walk_markdown(p))
leafs: list[dict] = []
for path in files:
try:
leafs.extend(to_all(read_markdown(path), path, repo="eSlider/2dph"))
except OSError as e:
print(f"mail/index_mail: skip {path}: {e}", file=sys.stderr)
return leafs
def _index_leafs(conn, leafs: list[dict], embed_fn) -> tuple[int, int]:
count = 0
for lf in leafs:
query = f"{lf['heading']}\n\n{lf['text']}"
emb = embed_fn(lf["text"]) if lf["text"] else None
upsert_leaf(conn, text=query, root="info", confidence="confirmed",
source=lf["source"], source_rev="mail" if lf.get("how") == "mail/import" else "working-tree",
how=lf.get("how", "kb/index"), loc=lf["source"], type_=lf.get("type", "reference"),
embedding=emb)
count += 1
return count, len(leafs)
def _embedder():
from model2vec import StaticModel
model = StaticModel.from_pretrained("minishlab/potion-multilingual-128M")
return lambda text: model.encode([text])[0].astype(float).tolist()
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))
+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 --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 -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
Executable
+27
View File
@@ -0,0 +1,27 @@
//usr/bin/env go run "$0" "$@"; exit
// bin/serve.go - async Go HTTP server for the 2dph brain (see bin/server).
//
// KB_ROOT=/path/to/2dph ./bin/serve.go # serve the brain
// KB_SEARCH_CMD=... KB_WORKERS=4 KB_PORT=8630 ./bin/serve.go
//
// Shebang trick: the first line is a Go `//` comment; when executed, env runs
// `go run "$0"` so this file doubles as an executable script. 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/server"
)
func main() {
if env := os.Getenv("KB_ROOT"); env == "" {
if wd, err := os.Getwd(); err == nil {
os.Setenv("KB_ROOT", wd)
}
}
server.Run()
}
+15 -4
View File
@@ -1,9 +1,16 @@
// Package main serves the 2dph brain over HTTP. // Package server serves the 2dph brain over HTTP.
// //
// Async by design: every request runs on its own goroutine, and CPU-heavy // Async by design: every request runs on its own goroutine, and CPU-heavy
// searches are serialized through a bounded worker pool (a counting // searches are serialized through a bounded worker pool (a counting
// semaphore) so N requests can't spawn N Python interpreters at once. // semaphore) so N requests can't spawn N Python interpreters at once.
package main //
// Used by bin/serve.go which is a self-executing shebang script:
//
// ///usr/bin/env go run "$0" "$@"; exit
// package main
// import "github.com/eSlider/2dph/bin/server"
// func main() { server.Run() }
package server
import ( import (
"context" "context"
@@ -115,10 +122,14 @@ func (b *brainSearcher) Search(ctx context.Context, query string, limit int) ([]
return out, nil return out, nil
} }
func main() { // Run starts the HTTP server. Reads env: KB_SEARCH_CMD (default bin/kb/search,
// relative to the repo root given by KB_ROOT), KB_WORKERS (default 4), KB_PORT
// (default 8630).
func Run() {
root := os.Getenv("KB_ROOT")
searchPath := os.Getenv("KB_SEARCH_CMD") searchPath := os.Getenv("KB_SEARCH_CMD")
if searchPath == "" { if searchPath == "" {
searchPath = filepath.Join("bin", "kb", "search") searchPath = filepath.Join(root, "bin", "kb", "search")
} }
workers := 4 workers := 4
if raw := os.Getenv("KB_WORKERS"); raw != "" { if raw := os.Getenv("KB_WORKERS"); raw != "" {
@@ -1,11 +1,8 @@
package main package server
import ( import (
"bytes"
"context" "context"
"encoding/json" "encoding/json"
"fmt"
"io"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"sync" "sync"
@@ -55,10 +52,6 @@ func (f *fakeSearcher) count() int {
return f.calls return f.calls
} }
func newTestServer(s Searcher, workers int) http.Handler {
return NewServer(s, workers)
}
func get(t *testing.T, h http.Handler, path string) (int, []byte) { func get(t *testing.T, h http.Handler, path string) (int, []byte) {
t.Helper() t.Helper()
req := httptest.NewRequest(http.MethodGet, path, nil) req := httptest.NewRequest(http.MethodGet, path, nil)
@@ -68,7 +61,7 @@ func get(t *testing.T, h http.Handler, path string) (int, []byte) {
} }
func TestHealth(t *testing.T) { func TestHealth(t *testing.T) {
h := newTestServer(&fakeSearcher{}, 1) h := NewServer(&fakeSearcher{}, 1)
code, body := get(t, h, "/health") code, body := get(t, h, "/health")
if code != http.StatusOK { if code != http.StatusOK {
t.Fatalf("health code = %d, want 200", code) t.Fatalf("health code = %d, want 200", code)
@@ -83,7 +76,7 @@ func TestHealth(t *testing.T) {
} }
func TestSearchMissingQuery(t *testing.T) { func TestSearchMissingQuery(t *testing.T) {
h := newTestServer(&fakeSearcher{}, 1) h := NewServer(&fakeSearcher{}, 1)
if code, _ := get(t, h, "/search"); code != http.StatusBadRequest { if code, _ := get(t, h, "/search"); code != http.StatusBadRequest {
t.Fatalf("code = %d, want 400", code) t.Fatalf("code = %d, want 400", code)
} }
@@ -93,7 +86,7 @@ func TestSearchReturnsSearcherResult(t *testing.T) {
fs := &fakeSearcher{callback: func(q string, limit int) ([]byte, error) { fs := &fakeSearcher{callback: func(q string, limit int) ([]byte, error) {
return []byte(`{"query":"` + q + `","count":1,"results":[{"id":"x"}]}`), nil return []byte(`{"query":"` + q + `","count":1,"results":[{"id":"x"}]}`), nil
}} }}
h := newTestServer(fs, 1) h := NewServer(fs, 1)
code, body := get(t, h, "/search?q=matrix") code, body := get(t, h, "/search?q=matrix")
if code != http.StatusOK { if code != http.StatusOK {
t.Fatalf("code = %d, want 200", code) t.Fatalf("code = %d, want 200", code)
@@ -113,7 +106,7 @@ func TestSearchReturnsSearcherResult(t *testing.T) {
func TestSearchConcurrencyBounded(t *testing.T) { func TestSearchConcurrencyBounded(t *testing.T) {
// 8 parallel requests on a 3-worker pool: at most 3 concurrent searches. // 8 parallel requests on a 3-worker pool: at most 3 concurrent searches.
fs := &fakeSearcher{delay: 20 * time.Millisecond} fs := &fakeSearcher{delay: 20 * time.Millisecond}
h := newTestServer(fs, 3) h := NewServer(fs, 3)
var wg sync.WaitGroup var wg sync.WaitGroup
for i := 0; i < 8; i++ { for i := 0; i < 8; i++ {
@@ -142,7 +135,7 @@ func TestSearchConcurrencyBounded(t *testing.T) {
} }
func TestSearchRejectsBadLimit(t *testing.T) { func TestSearchRejectsBadLimit(t *testing.T) {
h := newTestServer(&fakeSearcher{}, 1) h := NewServer(&fakeSearcher{}, 1)
if code, _ := get(t, h, "/search?q=x&n=hundred"); code != http.StatusBadRequest { if code, _ := get(t, h, "/search?q=x&n=hundred"); code != http.StatusBadRequest {
t.Fatalf("code = %d, want 400", code) t.Fatalf("code = %d, want 400", code)
} }
@@ -171,7 +164,4 @@ func TestSearchTimeout(t *testing.T) {
case <-time.After(500 * time.Millisecond): case <-time.After(500 * time.Millisecond):
t.Fatal("request hung after context cancellation") t.Fatal("request hung after context cancellation")
} }
_ = io.Discard
_ = bytes.MinRead
_ = fmt.Sprintf
} }
View File
+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
+117
View File
@@ -0,0 +1,117 @@
"""gitimport - parse `git log` output and turn commits into brain leafs.
Pure, testable functions. Field grammar (see bin/git/import):
git log --no-merges --name-only \
--format='%x1e%H%x1f%an%x1f%ae%x1f%aI%x1f%s'
0x1e = record separator, 0x1f = field separator.
Files: newline-separated lines following each record's subject.
"""
from __future__ import annotations
from dataclasses import dataclass, field
REC_SEP = "\x1e"
FIELD_SEP = "\x1f"
@dataclass
class Commit:
sha: str
author: str
email: str
date: str
subject: str
files: list[str] = field(default_factory=list)
def leaf_text(self, repo: str) -> str:
head = f"commit {self.sha[:12]} in {repo}{self.subject}"
body = [head, f"Author: {self.author} <{self.email}>", f"Date: {self.date}"]
if self.files:
body.append("Changing: " + ", ".join(self.files))
return "\n".join(body)
def parse_log(text: str) -> list[Commit]:
"""Parse `git log` output into Commit records.
Records are separated by 0x1e. A record is fields joined by 0x1f,
followed by optional newline-separated file paths inside the next
segment (git emits blank line + files after each record).
"""
commits: list[Commit] = []
# field records and file lists alternate; simpler: split on REC_SEP,
# each chunk = header line, possibly followed by newline + files.
for chunk in text.split(REC_SEP):
chunk = chunk.strip("\n")
if not chunk:
continue
lines = chunk.split("\n", 1)
header = lines[0].split(FIELD_SEP)
if len(header) < 5:
continue
sha, author, email, date, subject = header[:5]
files = [ln.strip() for ln in lines[1].splitlines() if ln.strip()] if len(lines) > 1 else []
commits.append(Commit(sha=sha, author=author, email=email,
date=date, subject=subject, files=files))
return commits
def commits_to_leafs(commits: list[Commit], repo: str) -> list[dict]:
"""Map commits to the leaf shape bin/kb/index expects (source/repo/...)."""
out: list[dict] = []
for c in commits:
out.append({
"source": f"{repo}@{c.sha}",
"repo": repo,
"heading": f"commit {c.sha[:12]}{c.subject}",
"text": c.leaf_text(repo),
"type": "commit",
"status": "current",
"related": ",".join(c.files),
})
return out
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)
+94 -12
View File
@@ -22,7 +22,17 @@ ROOT_FACTS = "facts"
ROOT_INFO = "info" ROOT_INFO = "info"
CONF_CONFIRMED = "confirmed" CONF_CONFIRMED = "confirmed"
VAR = Path(__file__).resolve().parents[1] / "var" 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" DB_PATH = VAR / "kb.lbug"
@@ -68,6 +78,19 @@ def init_schema(conn: ladybug.Connection) -> None:
conn.execute( conn.execute(
"CREATE REL TABLE IF NOT EXISTS RUNS_ON (FROM Leaf TO Host)" "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: def leaf_id(text: str, source: str) -> str:
@@ -95,18 +118,77 @@ def upsert_leaf(conn: ladybug.Connection, *, text: str, root: str, confidence: s
return lid return lid
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: def create_fts_and_vector(conn: ladybug.Connection, force: bool = False) -> None:
if force: """Create FTS (BM25) + HNSW vector indexes if missing.
conn.execute("DROP INDEX IF EXISTS Leaf.Leaf_fts")
conn.execute("DROP INDEX IF EXISTS Leaf.Leaf_vec") Never DROP INDEX for FTS/VECTOR. Ladybug 0.19 leaves ghost catalog
try: entries after DROP (`_0_Leaf_vec_UPPER`, `0_id_docs`), so a later
conn.execute("CALL CREATE_FTS_INDEX('Leaf', 'id', ['text'])") CREATE fails with "already exists in catalog" while SHOW_INDEXES
except Exception: still omits the index. Swallowing that error made HNSW look "OK"
pass until the first QUERY_VECTOR_INDEX.
try:
conn.execute("CALL CREATE_VECTOR_INDEX('Leaf', 'Leaf_vec', 'embedding', metric := 'cosine')") `force=True` is accepted for API compatibility but does **not** drop.
except Exception: Fresh indexes require deleting `var/kb.lbug` and rebuilding
pass (`bin/kb/index --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/kb/index --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/kb/index --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]: def query_fts(conn: ladybug.Connection, text: str, limit: int = 10) -> list[dict]:
+148
View File
@@ -0,0 +1,148 @@
"""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 re
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/docling skip them; we try
# pandoc first, else leave a stub.
LEGACY_OFFICE_SUFFIXES = {".doc", ".xls", ".ppt"}
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
+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()
+71
View File
@@ -0,0 +1,71 @@
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
SAMPLE = (
"\x1e" + "a1b2c3d" + "\x1f" + "Ada Lovelace" + "\x1f" + "ada@example.com"
+ "\x1f" + "2026-08-10T12:00:00+01:00" + "\x1f" + "feat: first commit"
+ "\n\nREADME.md\nsrc/main.c\n"
)
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)"
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):
cs = gitimport.parse_log(SAMPLE)
gitimport.index_commits(self.conn, cs, "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")
# File -[:HAS_VERSION]-> Commit -[:AUTHORED]-> Person
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 = gitimport.parse_log(SAMPLE)
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()
+57
View File
@@ -0,0 +1,57 @@
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent))
import gitimport # noqa: E402
SAMPLE = (
"\x1e" + "a1b2c3d" + "\x1f" + "Ada Lovelace" + "\x1f" + "ada@example.com"
+ "\x1f" + "2026-08-10T12:00:00+01:00" + "\x1f" + "feat: first commit"
+ "\n\nREADME.md\nsrc/main.c\n"
+ "\x1e" + "e4f5a6b" + "\x1f" + "Bob Babbage" + "\x1f" + "bob@example.com"
+ "\x1f" + "2026-08-11T09:30:00+01:00" + "\x1f" + "fix: typo"
+ "\n\ndocs/notes.md"
)
class GitparseTest(unittest.TestCase):
def test_parses_records(self):
cs = gitimport.parse_log(SAMPLE)
self.assertEqual(len(cs), 2)
def test_parses_commit_fields(self):
cs = gitimport.parse_log(SAMPLE)
c = cs[0]
self.assertEqual(c.sha, "a1b2c3d")
self.assertEqual(c.author, "Ada Lovelace")
self.assertEqual(c.email, "ada@example.com")
self.assertEqual(c.date, "2026-08-10T12:00:00+01:00")
self.assertEqual(c.subject, "feat: first commit")
def test_parses_changed_files(self):
cs = gitimport.parse_log(SAMPLE)
self.assertEqual(cs[0].files, ["README.md", "src/main.c"])
self.assertEqual(cs[1].files, ["docs/notes.md"])
def test_ignores_empty(self):
self.assertEqual(gitimport.parse_log(""), [])
def test_skip_malformed_record(self):
self.assertEqual(gitimport.parse_log("\x1eweird\x1e"), [])
def test_commit_leaf_shape(self):
leafs = gitimport.commits_to_leafs(gitimport.parse_log(SAMPLE), "sample-repo")
self.assertEqual(len(leafs), 2)
lf = leafs[0]
self.assertEqual(lf["type"], "commit")
self.assertEqual(lf["repo"], "sample-repo")
self.assertEqual(lf["source"], "sample-repo@a1b2c3d")
self.assertIn("Ada Lovelace", lf["text"])
self.assertIn("README.md", lf["related"])
self.assertIn("feat: first commit", lf["heading"])
if __name__ == "__main__":
unittest.main()
@@ -35,7 +35,7 @@ class KblibTest(unittest.TestCase):
confidence="confirmed", source="s", source_rev="r1", confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference", how="test", loc="/tmp", type_="reference",
embedding=make_emb(1.0)) embedding=make_emb(1.0))
kblib.create_fts_and_vector(self.conn, force=True) kblib.ensure_indexes(self.conn)
hits = kblib.query_fts(self.conn, "fox", 5) hits = kblib.query_fts(self.conn, "fox", 5)
self.assertEqual(len(hits), 1) self.assertEqual(len(hits), 1)
self.assertEqual(hits[0]["root"], "info") self.assertEqual(hits[0]["root"], "info")
@@ -49,12 +49,42 @@ class KblibTest(unittest.TestCase):
confidence="confirmed", source="s", source_rev="r1", confidence="confirmed", source="s", source_rev="r1",
how="test", loc="/tmp", type_="reference", how="test", loc="/tmp", type_="reference",
embedding=make_emb(0.0)) embedding=make_emb(0.0))
kblib.create_fts_and_vector(self.conn, force=True) kblib.ensure_indexes(self.conn)
result = kblib.hybrid_search(self.conn, make_emb(1.0), [], 5) result = kblib.hybrid_search(self.conn, make_emb(1.0), [], 5)
self.assertTrue(result) self.assertTrue(result)
self.assertIn("rrf", result[0]) self.assertIn("rrf", result[0])
self.assertEqual(result[0]["text"], "the quick brown fox") 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_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): def test_stats_counts_roots(self):
kblib.upsert_leaf(self.conn, text="a fact leaf", root="facts", kblib.upsert_leaf(self.conn, text="a fact leaf", root="facts",
confidence="confirmed", source="s", source_rev="r1", confidence="confirmed", source="s", source_rev="r1",
+123
View File
@@ -0,0 +1,123 @@
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
clean_email_address,
html_to_markdown,
is_convertible,
normalize_markdown,
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 _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()
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 bin/kb/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 kb/index command template. %s is replaced by the repo
// root (from KB_ROOT). Defaults to `python3 <root>/bin/kb/index`.
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"
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")
}
}
+49
View File
@@ -0,0 +1,49 @@
package watch
import (
"os"
"path/filepath"
"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 opts.IndexCmd == "" {
t.Fatal("default index cmd is empty")
}
}
+1 -1
View File
@@ -25,7 +25,7 @@ import urllib.parse
import urllib.request import urllib.request
from pathlib import Path from pathlib import Path
TOOLS = Path(__file__).resolve().parents[1].parent / "tools" TOOLS = Path(__file__).resolve().parents[1] / "tools"
sys.path.insert(0, str(TOOLS)) sys.path.insert(0, str(TOOLS))
sys.path.insert(0, str(TOOLS / "web-search")) sys.path.insert(0, str(TOOLS / "web-search"))
+23
View File
@@ -0,0 +1,23 @@
# Chat Import Pipeline
Plan: https://git.produktor.io/eSlider/brain-chats-import/issues/1
## Env vars (set in shell, never committed)
```
TELEGRAM_MCP_DIR
TELEGRAM_API_ID / TELEGRAM_API_HASH / TELEGRAM_PHONE
TELEGRAM_SESSION_STRING
ONLYOFFICE_URL / ONLYOFFICE_USER / ONLYOFFICE_PASS
OO_CLI (default: $HOME/go/bin/oo)
```
## Quick reference
```
./bin/chat sync telegram --limit 100
./bin/chat import
./bin/chat index
./bin/chat facts
./bin/chat apply --dry-run
```
+8
View File
@@ -0,0 +1,8 @@
module github.com/eSlider/2dph
go 1.25.0
require (
github.com/arran4/golang-ical v0.3.5
golang.org/x/text v0.40.0
)
+14
View File
@@ -0,0 +1,14 @@
github.com/arran4/golang-ical v0.3.5 h1:bbz6ld4dC+MmCKiFfOd6SkmIGnhNMBACZ485ULh7p9A=
github.com/arran4/golang-ical v0.3.5/go.mod h1:OnguFgjN0Hmx8jzpmWcC+AkHio94ujmLHKoaef7xQh8=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI=
github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/testify v1.7.0 h1:nwc3DEeHmmLAfoZucVR881uASk0Mfjw8xYJ99tb5CcY=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs=
golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+2
View File
@@ -6,7 +6,9 @@ readme = "README.md"
requires-python = ">=3.12" requires-python = ">=3.12"
license = { text = "MIT" } license = { text = "MIT" }
dependencies = [ dependencies = [
"docling>=2.119.0",
"ladybug==0.19.1", "ladybug==0.19.1",
"markitdown[docx,epub,html,image-exif,pdf,pptx,xlsx,zip]>=0.1.7",
"mistune==3.3.4", "mistune==3.3.4",
"model2vec==0.8.2", "model2vec==0.8.2",
"numpy>=2.5.2", "numpy>=2.5.2",
+76
View File
@@ -0,0 +1,76 @@
#!/usr/bin/env python3
"""Load test: bulk insert performance (writing facts to the brain)."""
from __future__ import annotations
import json
import time
import sys
from pathlib import Path
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT / "bin" / "tools"))
from kblib import connect, init_schema, upsert_leaf, ensure_indexes
from model2vec import StaticModel
DB_PATH = ROOT / "var" / "kb.lbug"
def run_bulk_tests(count: int = 100, _drop_indexes: bool = True) -> dict:
results: dict = {}
print(f"Bulk insert test: {count} leafs")
# Fresh DB for load test — never DROP INDEX on a live catalog.
if DB_PATH.exists():
DB_PATH.unlink()
db, conn = connect(str(DB_PATH), read_only=False)
init_schema(conn)
model = StaticModel.from_pretrained("minishlab/potion-multilingual-128M")
t0 = time.time()
for i in range(count):
text = f"bulk load test fact {i:03d} running container on host"
emb = model.encode([text])[0].astype(float).tolist()
upsert_leaf(
conn,
text=text,
root="facts",
confidence="confirmed",
source="load-test",
source_rev="2dph",
how="load_test_bulk",
loc=f"test:{i}",
type_="fact",
embedding=emb,
)
t1 = time.time()
# Create indexes once after bulk write (DROP+recreate is unsafe on Ladybug 0.19).
ensure_indexes(conn)
conn.close()
db.close()
elapsed = t1 - t0
results["total_seconds"] = elapsed
results["throughput"] = count / elapsed # leafs/sec
results["count"] = count
return results
def main():
count_str = sys.argv[1] if len(sys.argv) > 1 else "200"
count = int(count_str)
drop_str = sys.argv[2] if len(sys.argv) > 2 else "true"
drop_indexes = drop_str.lower() in ("1", "true", "yes")
r = run_bulk_tests(count=count, _drop_indexes=drop_indexes)
print(json.dumps(r, indent=2))
if __name__ == "__main__":
main()
+118
View File
@@ -0,0 +1,118 @@
# Brain Load Test Summary
**Date**: 2026-08-11
**Project**: 2dph (deductionphile)
**Target**: LadybugDB-embedded knowledge graph brain
## Test Suite
Four independent load tests were written and executed in `qa/`:
| Test | Purpose | Key Finding |
|------|---------|-------------|
| `load_test_search.py` | FTS, vector, hybrid search latency | FTS: 2.8ms, Vector: 1.9ms, Hybrid: 3.3ms |
| `load_test_graph.py` | Cypher hop traversal (1-hop, 2-hop, 3-hop) | 1-hop: 2.8ms, 2-hop: 4.7ms, 3-hop: 6.6ms |
| `load_test_queries.py` | Query pattern diversity (9 patterns) | All patterns under 25ms |
| `load_test_bulk.py` | Bulk insert throughput (leafs/sec) | 251 leafs/sec (with index drop/recreate) |
## Results
### 1. Search Performance (`load_test_search.py`, 10 iterations)
| Mode | Avg Latency (ms) | Description |
|------|-----------------|-------------|
| FTS (BM25) | **2.8 ms** | Pure keyword search |
| Vector (HNSW cosine) | **1.9 ms** | Embedding similarity search |
| Hybrid (RRF merge) | **3.3 ms** | FTS + vector fusion |
**Observation**: All modes under 5ms. Hybrid is ~1.8x slower than individual modes due to RRF overhead, but still well under 25ms per query.
### 2. Graph Traversal (`load_test_graph.py`, 10 iterations)
| Pattern | Avg Latency (ms) | Description |
|---------|-----------------|-------------|
| 1-hop (Leaf -FROM_FILE-> File) | **2.8 ms** | Simple edge traversal |
| 2-hop (Leaf -> File -> Commit) | **4.7 ms** | Two-hop path with mix node types |
| 3-hop (facts -from_file-> File -> HAS_VERSION-> Commit -AUTHORED-> Person) | **6.6 ms** | Three-hop path with root filter |
| Degree centrality (avg children per file) | **4.2 ms** | Aggregation query |
**Observation**: Graph queries are very fast (<10ms even for 3 hops) on the knowledge graph.
### 3. Query Pattern Diversity (`load_test_queries.py`, 10 iterations)
| Query Pattern | Avg Latency (ms) |
|---------------|-----------------|
| fact_source (docker, root=facts) | 5.3 |
| info_docker (docker, root=info) | 3.4 |
| info_k8s (kubernetes, root=info) | 3.2 |
| repo_2dph (search term, repo=eSlider/2dph) | 3.6 |
| facts_no_root (search, root=facts) | 4.1 |
| hybrid_container (container, hybrid search) | 3.9 |
| hybrid_service (service, hybrid search) | 3.7 |
| multi_obs (observability, multi-word) | 4.3 |
| multi_container (container orchestration, multi-word) | 4.7 |
**Observation**: All 9 query patterns complete in under 25ms. The system correctly handles root-filtered and repo-filtered searches.
### 4. Bulk Insert (`load_test_bulk.py`, 30 leafs, indexes dropped before insert)
| Metric | Value |
|--------|-------|
| Total time for 30 leafs | 0.12s |
| Throughput | **251 leafs/sec** |
**Critical observation (corrected 2026-08-12)**: LadybugDB **0.19** must **not**
`DROP INDEX` for FTS/VECTOR and recreate. DROP leaves ghost catalog tables
(`_0_Leaf_vec_UPPER`, `0_id_docs`); CREATE then fails with "already exists in
catalog" while `SHOW_INDEXES` omits the index — HNSW looks dead until
`var/kb.lbug` is deleted. Upsert while indexes exist keeps HNSW queryable.
Fresh indexes: delete the DB file and `bin/kb/index --rebuild`. See
`kblib.create_fts_and_vector` / `ensure_indexes`.
## Critical Assessment - Evidence Rule Working
The most important finding: **the evidence-based audit correctly enforces the two-source rule for facts**.
- `bin/facts/audit db` runs against `var/kb.lbug` and asserts each `root=facts` leaf has:
- A `source` field containing " x " (indicating two independent sources, e.g., "docker ps x compose:docker-compose.yml")
- A non-empty `loc` (evidence pointer)
- `confidence='confirmed'`
- **Before cleanup**: Database had 50 test facts with `source="load-test"` (single source) → audit correctly flagged all as failing the 2-source rule
- **After cleanup (12 facts from extract)**: Audit passes (`ok: true, problems: []`) because the 12 facts have proper 2-source evidence:
- 11 facts: `source="docker ps x compose:..."` or `source="docker ps x compose:..."`
- 1 fact: `source="ssh config x docs(README.md, PLAN.md, AGENTS.md)"`
This validates the core design principle from PLAN.md (D8/D11): **a fact needs ≥2 independent sources or it is `(not confirmed)`**.
## Database State (After Cleanup)
| Metric | Value |
|--------|-------|
| Total leaves | 89 (47 info + 12 facts) |
| Facts (root=facts) | 12, all with 2-source evidence |
| Info (root=info) | 47 (from markdown corpus) |
| Audit result | `ok: true, problems: []` |
## Files in `qa/`
- `load_test_search.py` - Search latency test (FT/Vector/Hybrid)
- `load_test_graph.py` - Graph traversal test (1-hop, 2-hop, 3-hop)
- `load_test_queries.py` - Query pattern diversity test (9 patterns)
- `load_test_bulk.py` - Bulk insert throughput test
- `load_test_summary.md` - This summary
## Verdict
The brain performs well within design parameters:
- **Search/retrieval latency**: sub-25ms across all modes
- **Graph traversal**: under 10ms even for 3-hop paths
- **Bulk insertion**: ~250 leafs/sec (with proper index management)
- **Evidence enforcement**: The two-source audit correctly validates facts, confirming the detective method works as designed (`facts` root = strong assertions, `info` root = weak claims)
The system is ready for production use with the understanding that:
1. Bulk inserts must drop/recreate indexes to avoid corruption
2. Facts are only stored when backed by >=2 independent sources (enforced by audit)
3. The info root holds the narrative corpus (28K+ markdown-derived leafs)
4. Facts root holds confirmed assertions with evidence links
-3
View File
@@ -1,3 +0,0 @@
module github.com/eSlider/2dph/serve
go 1.25
Generated
+2379 -18
View File
File diff suppressed because it is too large Load Diff