Compare commits

..
Author SHA1 Message Date
Garry TanandClaude Fable 5 3487e4b255 fix(embed): thread write column into sumStaleChunkChars — sync cost gate followed the legacy predicate (#1262)
Review follow-up: the PR threaded the write-side column through
countStaleChunks/listStaleChunks but not sumStaleChunkChars, so the
sync cost gate counted an alt-column brain's fully-embedded corpus as
phantom backlog on every gate (inflating the --full estimate and the
deferred-mode backlog note). Both engines already share
buildStaleChunkWhere, so this is a type widening + one call-site
thread + a contrast assertion.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-22 10:59:36 -07:00
159ddc3249 fix(embed): dynamic embedding column write target for upsertChunks + stale scans (#1262)
PostgresEngine/PGLiteEngine.upsertChunks hardcoded the legacy `embedding`
column (INSERT list + `::vector` cast + ON CONFLICT clauses), so a brain
with a registered alternate embedding column (e.g. halfvec(2560)) failed
every put_page/import/sync/embed write with a dimension mismatch — the
embedding_columns registry had read-side consumers only (PR #1164).

Fix (write-side symmetric to the read-side resolver):
- `resolveWriteColumn(cfg)` + `resolveWriteColumnForEngine(engine)` in
  search/embedding-column.ts: user-declared registry entry whose provider
  matches the current embedding model wins; no match => undefined (legacy
  column). Builtins are never consulted so the multimodal builtin can't
  capture text writes.
- `upsertChunks` accepts a caller-resolved `embeddingColumn` descriptor in
  BOTH engines; the target column + cast (`::vector` / `::halfvec(N)`)
  and the v0.40.3.0 D24 ON CONFLICT race-fix CASE arms follow the column.
- Stale scans follow the write column: `countStaleChunks`,
  `listStaleChunks` (both cursor arms), and the shared
  buildStaleChunkWhere accept the descriptor — without this, an
  alt-column brain re-selects (and re-pays for) already-embedded chunks
  forever.
- Boundary threading: runEmbedCore (embedPage/embedAll/embedAllStale),
  embedStaleForSource + the embed-backfill handler, importFromContent /
  importCodeFile / withImportTransaction, and the contextual-retrieval
  re-embed transaction all resolve once and pass the descriptor.
- `preflightDimMismatch` skips the legacy-column dim comparison when a
  non-default write column resolved (it would otherwise hard-block embed
  runs on alt-column brains).

Tests: resolveWriteColumn unit coverage; PGLite e2e for the write target,
D24 preserve-on-reupsert, stale-scan contrast, and an `embed --stale
--dry-run` convergence integration; Postgres e2e twins (DATABASE_URL-gated).

Takeover of PR #1263 rebased onto current master (D24 ON CONFLICT
semantics, batchRetry wrapper, signature-stale + embed-backfill paths).

Fixes #1262

Co-authored-by: DmitryBMsk <DmitryBMsk@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 15:04:38 -07:00
25 changed files with 578 additions and 574 deletions
+1 -3
View File
@@ -1411,8 +1411,7 @@ This is the dispatcher. Skills are the implementation. **Read the skill file bef
| "get more out of gbrain", "is my brain set up right", "weekly brain checkup", "advise me on my brain", "gbrain advisor" | `skills/gbrain-advisor/SKILL.md` |
| Save or load reports | `skills/reports/SKILL.md` |
| "Create a skill", "improve this skill" | `skills/skill-creator/SKILL.md` |
| "save this learning to the vault", "capture this skill in Obsidian", "record this workflow in my notes", "put this setup change in the vault" | `skills/skill-vault-capture-policy/SKILL.md` |
| "Skillify this", "is this a skill?", "make this proper", "add tests and evals for this" | `skills/skillify/SKILL.md` |
| "Skillify this", "is this a skill?", "make this proper" | `skills/skillify/SKILL.md` |
| "Compress my resolver", "AGENTS.md too large", "RESOLVER.md too big", "functional area dispatcher", "shrink routing table" | `skills/functional-area-resolver/SKILL.md` |
| "Is gbrain healthy?", morning health check, skillpack-check | `skills/skillpack-check/SKILL.md` |
| "harvest this skill into gbrain", "publish this skill to gbrain", "lift this skill upstream", "share this skill with other gbrain clients", "promote my skill to gbrain" | `skills/skillpack-harvest/SKILL.md` |
@@ -1430,7 +1429,6 @@ This is the dispatcher. Skills are the implementation. **Read the skill file bef
| "Set up GBrain", first boot | `skills/setup/SKILL.md` |
| "Now what?", "fill my brain", "cold start", "bootstrap", "import my data", "what should I import first" | `skills/cold-start/SKILL.md` |
| "Migrate from Obsidian/Notion/Logseq" | `skills/migrate/SKILL.md` |
| "Connect Obsidian to gbrain", "import my vault to gbrain", "sync vault and gbrain", "embed gbrain after vault update", "is gbrain synced with my vault" | `skills/obsidian-gbrain-safe-index/SKILL.md` |
| Brain health check, maintenance run | `skills/maintain/SKILL.md` |
| "Extract links", "build link graph", "populate timeline" | `skills/maintain/SKILL.md` (extraction sections) |
| "Run dream", "process today's session", "synthesize my conversations", "consolidate yesterday's conversations", "what patterns did you see", "did the dream cycle run" | `skills/maintain/SKILL.md` (dream cycle section) |
+2 -3
View File
@@ -39,8 +39,8 @@
"skills/briefing",
"skills/citation-fixer",
"skills/concept-synthesis",
"skills/cron-scheduler",
"skills/cross-modal-review",
"skills/cron-scheduler",
"skills/daily-task-manager",
"skills/daily-task-prep",
"skills/data-research",
@@ -54,11 +54,10 @@
"skills/media-ingest",
"skills/meeting-ingestion",
"skills/minion-orchestrator",
"skills/obsidian-gbrain-safe-index",
"skills/perplexity-research",
"skills/query",
"skills/repo-architecture",
"skills/reports",
"skills/repo-architecture",
"skills/signal-detector",
"skills/skill-creator",
"skills/skillify",
+1 -3
View File
@@ -60,8 +60,7 @@ This is the dispatcher. Skills are the implementation. **Read the skill file bef
| "get more out of gbrain", "is my brain set up right", "weekly brain checkup", "advise me on my brain", "gbrain advisor" | `skills/gbrain-advisor/SKILL.md` |
| Save or load reports | `skills/reports/SKILL.md` |
| "Create a skill", "improve this skill" | `skills/skill-creator/SKILL.md` |
| "save this learning to the vault", "capture this skill in Obsidian", "record this workflow in my notes", "put this setup change in the vault" | `skills/skill-vault-capture-policy/SKILL.md` |
| "Skillify this", "is this a skill?", "make this proper", "add tests and evals for this" | `skills/skillify/SKILL.md` |
| "Skillify this", "is this a skill?", "make this proper" | `skills/skillify/SKILL.md` |
| "Compress my resolver", "AGENTS.md too large", "RESOLVER.md too big", "functional area dispatcher", "shrink routing table" | `skills/functional-area-resolver/SKILL.md` |
| "Is gbrain healthy?", morning health check, skillpack-check | `skills/skillpack-check/SKILL.md` |
| "harvest this skill into gbrain", "publish this skill to gbrain", "lift this skill upstream", "share this skill with other gbrain clients", "promote my skill to gbrain" | `skills/skillpack-harvest/SKILL.md` |
@@ -79,7 +78,6 @@ This is the dispatcher. Skills are the implementation. **Read the skill file bef
| "Set up GBrain", first boot | `skills/setup/SKILL.md` |
| "Now what?", "fill my brain", "cold start", "bootstrap", "import my data", "what should I import first" | `skills/cold-start/SKILL.md` |
| "Migrate from Obsidian/Notion/Logseq" | `skills/migrate/SKILL.md` |
| "Connect Obsidian to gbrain", "import my vault to gbrain", "sync vault and gbrain", "embed gbrain after vault update", "is gbrain synced with my vault" | `skills/obsidian-gbrain-safe-index/SKILL.md` |
| Brain health check, maintenance run | `skills/maintain/SKILL.md` |
| "Extract links", "build link graph", "populate timeline" | `skills/maintain/SKILL.md` (extraction sections) |
| "Run dream", "process today's session", "synthesize my conversations", "consolidate yesterday's conversations", "what patterns did you see", "did the dream cycle run" | `skills/maintain/SKILL.md` (dream cycle section) |
+1 -11
View File
@@ -2,7 +2,7 @@
"name": "gbrain",
"version": "0.32.3.0",
"conformance_version": "1.0.0",
"description": "Personal knowledge brain with hybrid RAG search GStack mod for agent platforms",
"description": "Personal knowledge brain with hybrid RAG search \u2014 GStack mod for agent platforms",
"skills": [
{
"name": "ingest",
@@ -34,11 +34,6 @@
"path": "migrate/SKILL.md",
"description": "Universal migration from Obsidian, Notion, Logseq, markdown, CSV, JSON, Roam"
},
{
"name": "obsidian-gbrain-safe-index",
"path": "obsidian-gbrain-safe-index/SKILL.md",
"description": "Connect and sync an Obsidian-style Markdown vault with gbrain using a cost-controlled import-first workflow, explicit embedding gates, and durable skill/workflow capture."
},
{
"name": "setup",
"path": "setup/SKILL.md",
@@ -268,11 +263,6 @@
"name": "skill-optimizer",
"path": "skill-optimizer/SKILL.md",
"description": "Self-evolving skill optimization via gbrain skillopt — SkillOpt-paper-grounded text-space optimizer with validation gating (median-of-3 + epsilon=0.05), bundled-skill safety, bootstrap review sentinel, per-skill DB lock, and atomic versioned writes."
},
{
"name": "skill-vault-capture-policy",
"path": "skill-vault-capture-policy/SKILL.md",
"description": "Capture durable operational learnings, new agent skills, and environment/setup changes into the Obsidian vault instead of transient chat memory."
}
],
"dependencies": {
-137
View File
@@ -1,137 +0,0 @@
---
name: obsidian-gbrain-safe-index
version: 1.0.0
description: |
Connect, maintain, and sync an Obsidian-style Markdown vault with gbrain while preserving a cost-controlled workflow: the vault remains the source of truth, gbrain is the searchable/embedded index, durable skills/workflows are captured into the vault, and paid embedding runs only after explicit approval.
triggers:
- "connect Obsidian to gbrain"
- "import my vault to gbrain"
- "sync vault and gbrain"
- "capture this skill in my vault"
- "embed gbrain after vault update"
- "is gbrain synced with my vault"
tools:
- terminal
- read_file
- search_files
- write_file
- patch
mutating: true
---
# Obsidian → gbrain Safe Index and Capture
## Contract
This skill guarantees:
- Treats the user's Obsidian-style Markdown vault as the source of truth before gbrain indexing.
- Keeps gbrain in conservative mode unless the user explicitly approves a more expensive mode.
- Imports vault changes with `--no-embed` first, then embeds only after explicit approval for the specific paid action.
- Captures durable new skills, workflows, and environment learnings into the vault instead of leaving them only in chat memory.
- Verifies every sync with concrete `gbrain stats`, search mode, and, when embeddings run, exact embedded chunk counts.
## Phases
1. **Resolve the vault path.**
- Prefer an existing environment variable such as `OBSIDIAN_VAULT_PATH` or `WIKI_PATH`.
- If no path is configured, search likely note directories and ask the user before writing.
- Verify the directory exists and contains markdown files or an `.obsidian` directory.
2. **Read vault operating rules before writing.**
- If the vault has `SCHEMA.md`, `index.md`, `log.md`, `AGENTS.md`, or similar operating files, read them before ingest/query/major edit.
- Respect immutable source folders such as `raw/` when the vault declares them.
- Use the vault's native link convention, usually Obsidian `[[wikilinks]]`, for durable relationships.
3. **MECE/capture decision.**
- If new knowledge belongs on an existing page, update that page.
- If it is a distinct recurring workflow or operational policy, create a small meta or concept page following the vault schema.
- Update the vault index/catalog for every new page when the vault maintains one.
- Append a log entry for meaningful vault updates when the vault maintains a log.
4. **Safe gbrain import path.**
- Pre-check source directory; do not import a nonexistent path.
- Run `gbrain config set search.mode conservative` before/after risky reinit steps.
- Run `gbrain import "$OBSIDIAN_VAULT_PATH" --no-embed`.
- Run `gbrain extract links --source fs --dir "$OBSIDIAN_VAULT_PATH"` when wikilinks changed materially.
5. **Paid embedding gate.**
- Do not run `gbrain embed --stale` unless the user explicitly asks or a prior instruction clearly approved this exact paid action.
- Before embedding, verify provider readiness with `gbrain providers test --model <provider:model>`.
- Confirm the configured embedding dimensions match the local schema.
- After embedding, verify `gbrain stats` and record exact `Pages`, `Chunks`, `Embedded`, and `Links` counts.
6. **Final verification and vault echo.**
- Run `gbrain stats` and `gbrain search modes`.
- If vault files changed, re-import with `--no-embed`; if embedding was approved, embed stale chunks afterward.
- Report what changed, what was free/local, what used API billing, and what remains pending.
## Output Format
Use a compact status table:
| Item | Status |
|---|---|
| Vault path | `/path` |
| Vault updated | yes/no + files |
| gbrain mode | conservative/balanced/tokenmax |
| Import | `--no-embed` completed / skipped / failed |
| Pages/chunks | exact counts from `gbrain stats` |
| Embeddings | exact count; note whether this run used API billing |
| Links | exact count |
| Background jobs | none / list exact jobs |
Then include:
- **Safe next step:** free/local action.
- **Paid next step:** embedding/LLM action, if any, with explicit approval requirement.
## Anti-Patterns
- Creating a duplicate skill/page when an existing Obsidian, gbrain, or vault-ingest skill already covers the workflow.
- Running `gbrain embed --stale`, `gbrain dream`, `gbrain autopilot --install`, `gbrain onboard --auto`, or `tokenmax` without explicit cost approval.
- Importing a nonexistent or wrong directory and treating a zero-page import as success.
- Forgetting to update the vault index/catalog and log after creating or materially updating vault pages.
- Recording API keys, tokens, or raw secrets in the vault or final response.
## Tools Used
- `read_file` — read vault schema/index/log and target notes.
- `search_files` — find existing vault pages and avoid duplicates.
- `write_file` / `patch` — create or update vault pages.
- `terminal` — run `gbrain`, `git`, and environment checks with secret values redacted.
## Safe Commands
```bash
export PATH="$HOME/.bun/bin:$PATH"
gbrain config set search.mode conservative
gbrain import "$OBSIDIAN_VAULT_PATH" --no-embed
gbrain extract links --source fs --dir "$OBSIDIAN_VAULT_PATH"
gbrain stats
gbrain search modes
```
## Paid / Approval-Gated Commands
```bash
gbrain providers test --model <provider:model>
gbrain embed --stale
gbrain dream
gbrain autopilot --install
gbrain onboard --auto --max-usd 5
gbrain config set search.mode tokenmax
```
## Verification Checklist
- [ ] Vault path exists and is the intended source.
- [ ] Vault operating files were read before edits when present.
- [ ] Existing pages/skills were searched to avoid duplicates.
- [ ] New/updated vault pages follow the vault schema and link convention.
- [ ] Vault index/catalog updated for new pages when present.
- [ ] Vault log appended for meaningful actions when present.
- [ ] `gbrain import ... --no-embed` completed.
- [ ] `gbrain stats` recorded pages/chunks/embeddings/links.
- [ ] `gbrain search modes` confirms conservative mode unless a different mode was explicitly approved.
- [ ] No paid/background commands ran without approval.
@@ -1,11 +0,0 @@
// Routing eval fixtures for skills/obsidian-gbrain-safe-index.
// Positive cases: intents embed a trigger phrase in natural surrounding context
// (never verbatim-identical to a trigger — the fixture linter rejects tautologies).
{"intent": "help me connect Obsidian to gbrain for my notes", "expected_skill": "obsidian-gbrain-safe-index"}
{"intent": "import my vault to gbrain but skip embeddings for now", "expected_skill": "obsidian-gbrain-safe-index"}
{"intent": "please sync vault and gbrain after I edit notes", "expected_skill": "obsidian-gbrain-safe-index"}
{"intent": "run embed gbrain after vault update tonight", "expected_skill": "obsidian-gbrain-safe-index"}
{"intent": "hey is gbrain synced with my vault right now", "expected_skill": "obsidian-gbrain-safe-index"}
// Negative cases: related but owned by other skills. Assert NO route to this skill.
{"intent": "migrate my notes from Notion to gbrain", "expected_skill": null}
{"intent": "what is on my calendar tomorrow", "expected_skill": null}
-100
View File
@@ -1,100 +0,0 @@
---
name: skill-vault-capture-policy
version: 1.0.0
description: |
Use when the user wants a durable operational learning, new agent skill, or
important environment/setup change to be captured into the Obsidian vault
instead of left only in transient chat memory. Covers the capture rule,
preferred page patterns, and index/log update obligations.
triggers:
- "save this learning to the vault"
- "capture this skill in Obsidian"
- "record this workflow in my notes"
- "put this setup change in the vault"
- "should we add this to the knowledge base"
tools:
- read_file
- search_files
- write_file
- patch
mutating: true
---
# Skill Vault Capture Policy
## Contract
This skill guarantees:
- Durable operational learnings, newly adopted skills, and important environment/setup changes are captured into the Obsidian vault rather than left only in chat memory.
- An existing page is updated when the knowledge clearly belongs there; a new page is created only when the topic is distinct and likely to recur.
- Every new vault page is added to `index.md`.
- Every meaningful create/update appends a dated entry to `log.md`.
- Small linked pages are preferred over one giant running note.
## Phases
1. **Classify the learning.**
- New gbrain operating rule, cost control, or embedding/provider change.
- New agent-fork / harness / coding-tool integration fact.
- New Obsidian vault workflow or structure decision.
- New recurring agent skill that changes how the agent should operate here.
2. **Avoid duplicates.**
- Search the vault for an existing page that already owns the topic.
- If found, update it with a new section or dated note rather than creating a near-duplicate.
3. **Create when distinct.**
- Place new pages under the vault schema: `_meta/` for operating notes, `concepts/` for workflows, `entities/` for tools/people.
- Use YAML frontmatter and at least two `[[wikilinks]]` unless it is a short seed page.
4. **Update navigation.**
- Add the page to `index.md` under the correct type heading.
- Append a `## [YYYY-MM-DD] create|update | subject` entry to `log.md`.
5. **Report the capture.**
- State which files changed and whether `index.md` / `log.md` were updated.
## Output Format
Use a short status block:
| Item | Status |
|---|---|
| Learning classified | type |
| Page created/updated | path |
| index.md updated | yes/no |
| log.md updated | yes/no |
## Anti-Patterns
- Leaving durable learnings only in chat memory.
- Creating a near-duplicate page instead of updating the existing one.
- Forgetting to update `index.md` and `log.md`.
- Writing one giant running note instead of small linked pages.
- Recording secrets, API keys, or raw credentials in the vault.
## Tools Used
- `read_file` — read `SCHEMA.md`, `index.md`, `log.md`, and target pages.
- `search_files` — find existing pages to avoid duplicates.
- `write_file` / `patch` — create or update vault pages and navigation.
## Safe Commands
```bash
# inspect vault navigation before writing
read SCHEMA.md index.md log.md
# create or update a page, then refresh catalog/log
# index.md: add [[page-slug]] under the matching type heading
# log.md: append ## [YYYY-MM-DD] create|update | subject
```
## Verification Checklist
- [ ] Vault schema/read files were checked before writing.
- [ ] Existing pages were searched to avoid duplicates.
- [ ] New/updated page has frontmatter and wikilinks.
- [ ] `index.md` updated for new pages.
- [ ] `log.md` appended for meaningful actions.
- [ ] No secrets or raw credentials were written.
@@ -1,10 +0,0 @@
// Routing eval fixtures for skills/skill-vault-capture-policy.
// Positive cases: intents embed a trigger phrase in natural surrounding context
// (never verbatim-identical to a trigger — the fixture linter rejects tautologies).
{"intent": "please save this learning to the vault so we keep it", "expected_skill": "skill-vault-capture-policy", "ambiguous_with": ["idea-ingest"]}
{"intent": "we should capture this skill in Obsidian for reuse", "expected_skill": "skill-vault-capture-policy", "ambiguous_with": ["capture"]}
{"intent": "can you record this workflow in my notes for next time", "expected_skill": "skill-vault-capture-policy"}
{"intent": "put this setup change in the vault before we forget", "expected_skill": "skill-vault-capture-policy"}
// Negative cases: related but owned by other skills or out of scope.
{"intent": "connect my Obsidian vault to gbrain", "expected_skill": null}
{"intent": "what is on my calendar tomorrow", "expected_skill": null}
+43 -13
View File
@@ -1,6 +1,7 @@
import type { BrainEngine } from '../core/engine.ts';
import { embedBatch, currentEmbeddingSignature } from '../core/embedding.ts';
import type { ChunkInput } from '../core/types.ts';
import type { ChunkInput, ResolvedColumn } from '../core/types.ts';
import { resolveWriteColumnForEngine } from '../core/search/embedding-column.ts';
import { chunkText } from '../core/chunkers/recursive.ts';
import { createProgress, type ProgressReporter } from '../core/progress.ts';
import { getCliOptions, cliOptsToProgressOptions } from '../core/cli-options.ts';
@@ -183,8 +184,13 @@ export class EmbeddingDimMismatchError extends Error {
* fresh-install bug class at the very first invocation instead of letting
* the worker pool hammer N pages with raw 22000 errors.
*/
async function preflightDimMismatch(engine: BrainEngine, dryRun: boolean): Promise<void> {
async function preflightDimMismatch(engine: BrainEngine, dryRun: boolean, embeddingColumn?: ResolvedColumn): Promise<void> {
if (dryRun) return; // dry-run never embeds, no risk
// #1262: an alt-column brain writes to `embeddingColumn`, not the legacy
// `embedding` column — the legacy column's dims are irrelevant, and the
// registry entry (validated at resolve time) pins the target's dims. Only
// the legacy default path needs the schema-vs-gateway dim comparison.
if (embeddingColumn && embeddingColumn.name !== 'embedding') return;
const { readContentChunksEmbeddingDim, embeddingMismatchMessage } = await import('../core/embedding-dim-check.ts');
const { getEmbeddingDimensions, getEmbeddingModel } = await import('../core/ai/gateway.ts');
let existing;
@@ -238,7 +244,12 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
// v0.37.11.0 (Lane D.2): pre-flight dim-mismatch check. Catches the headline
// fresh-install bug class before the worker pool spends 20 parallel calls
// hitting raw Postgres dimension errors.
await preflightDimMismatch(engine, !!opts.dryRun);
// #1262: resolve the write-side embedding column ONCE at the boundary
// (merged config + gateway model) and thread the descriptor through every
// upsertChunks / stale-scan below. undefined => legacy `embedding` column.
const embeddingColumn = await resolveWriteColumnForEngine(engine);
await preflightDimMismatch(engine, !!opts.dryRun, embeddingColumn);
const result: EmbedResult = {
embedded: 0,
@@ -253,7 +264,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
for (const s of opts.slugs) {
if (isAborted(opts.signal)) break; // #1737: stop the per-slug loop on abort
try {
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal);
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal, embeddingColumn);
} catch (e: unknown) {
serr(` Error embedding ${s}: ${e instanceof Error ? e.message : e}`);
}
@@ -347,7 +358,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
catchUp: opts.catchUp,
pacer,
paceMaxConcurrency,
}, opts.signal);
}, opts.signal, embeddingColumn);
} finally {
// E1: surface pacing telemetry (human + structured) when pacing was on.
const snap = pacer.snapshot();
@@ -376,7 +387,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
return result;
}
if (opts.slug) {
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal);
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal, embeddingColumn);
return result;
}
throw new Error('No embed target specified. Pass { slug }, { slugs }, { all }, or { stale }.');
@@ -521,8 +532,13 @@ async function embedPage(
result: EmbedResult,
sourceId?: string,
signal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
const opts = sourceId ? { sourceId } : undefined;
// #1262: write-side descriptor rides only on WRITE calls (upsertChunks).
const chunkOpts = (sourceId || embeddingColumn)
? { ...(sourceId && { sourceId }), ...(embeddingColumn && { embeddingColumn }) }
: undefined;
const page = await engine.getPage(slug, opts);
if (!page) {
throw new Error(`Page not found: ${slug}`);
@@ -554,7 +570,7 @@ async function embedPage(
}
if (inputs.length > 0) {
await engine.upsertChunks(slug, inputs, opts);
await engine.upsertChunks(slug, inputs, chunkOpts);
chunks = await engine.getChunks(slug, opts);
}
}
@@ -589,7 +605,7 @@ async function embedPage(
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await engine.upsertChunks(slug, updated, opts);
await engine.upsertChunks(slug, updated, chunkOpts);
// v0.41.31: stamp provenance so a later model/dims swap is detectable as
// stale. embedPage is the per-slug path used by `gbrain embed <slug>` AND
// by `gbrain sync`'s post-import embed step (runEmbedCore({slugs})).
@@ -622,6 +638,7 @@ async function embedAll(
paceMaxConcurrency?: number;
},
signal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
// v0.41.31: current embedding provenance signature. Stamped onto pages
// when their chunks are (re)embedded so a later model/dimension swap is
@@ -644,7 +661,7 @@ async function embedAll(
// D7: thread sourceId so `gbrain embed --stale --source X` actually scopes.
// v0.41.18.0 (A13): thread batchSize/priority/catchUp into the stale path.
// #1737: thread the external abort signal so the cycle embed phase bails.
return await embedAllStale(engine, sourceId, dryRun, result, onProgress, staleOpts, signature, signal);
return await embedAllStale(engine, sourceId, dryRun, result, onProgress, staleOpts, signature, signal, embeddingColumn);
}
// --all path: pacer (no-op when off). E-1: lower the worker count to the
@@ -725,7 +742,10 @@ async function embedAll(
embedding: embeddingMap.get(c.chunk_index) ?? undefined,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(page.slug, updated, pageOpts));
await observed(pacer, () => engine.upsertChunks(page.slug, updated, {
...(pageSourceId && { sourceId: pageSourceId }),
...(embeddingColumn && { embeddingColumn }),
}));
// v0.41.31: stamp embedding provenance so a later model swap is
// detectable as stale.
await observed(pacer, () =>
@@ -805,10 +825,16 @@ async function embedAllStale(
},
signature?: string,
externalSignal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
// D7: thread sourceId so source-scoped runs only count + visit
// that source's NULL embeddings.
const sourceOpt = sourceId ? { sourceId } : undefined;
// #1262: the stale predicate follows the write-side column — without it an
// alt-column brain would perpetually re-select (and re-pay for) chunks whose
// target column is already populated.
const sourceOpt = (sourceId || embeddingColumn)
? { ...(sourceId && { sourceId }), ...(embeddingColumn && { embeddingColumn }) }
: undefined;
// v0.41.31: re-embed pages whose embedding_signature drifted (model/dims
// swap). dry-run must NOT mutate, so it counts signature-stale via the
@@ -967,6 +993,7 @@ async function embedAllStale(
afterUpdatedAt,
}),
...(sourceId && { sourceId }),
...(embeddingColumn && { embeddingColumn }),
}),
);
if (batch.length === 0) {
@@ -1019,7 +1046,10 @@ async function embedAllStale(
embedding: staleIdxToEmbedding.get(c.chunk_index) ?? undefined,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, { sourceId: keySourceId }));
await observed(pacer, () => engine.upsertChunks(slug, merged, {
sourceId: keySourceId,
...(embeddingColumn && { embeddingColumn }),
}));
// v0.41.31: stamp provenance after the page's chunks are embedded —
// but only when EVERY chunk was stale (fully re-embedded this pass).
// A partially-stale page keeps preserved chunks of unknown/old
@@ -1090,7 +1120,7 @@ async function embedAllStale(
// as a clean run — re-running won't help until the underlying failure is fixed.
if (staleOpts?.catchUp && !effectiveSignal.aborted && embedFailures > 0) {
const remaining = await engine.countStaleChunks(
signature ? { signature, ...(sourceId ? { sourceId } : {}) } : (sourceId ? { sourceId } : undefined),
signature ? { signature, ...sourceOpt } : sourceOpt,
);
if (remaining > 0) {
serr(`\n [embed] catch-up finished but ${remaining} chunk(s) remain stale after ${embedFailures} embed failure(s). These are not embeddable as-is; re-running won't clear them until the underlying error is resolved.`);
+8 -1
View File
@@ -576,9 +576,16 @@ async function runInlineCostGate(
// Stale backlog: cheap single SQL; fail-open to 0 so a transient DB hiccup
// never blocks the sync. Signature-aware (model/dims swap surfaces here).
// #1262: follow the write-side embedding column — otherwise an alt-column
// brain's fully-embedded corpus counts as phantom backlog on every gate.
let staleChars = 0;
try {
staleChars = await engine.sumStaleChunkChars({ signature: currentEmbeddingSignature() });
const { resolveWriteColumnForEngine } = await import('../core/search/embedding-column.ts');
const embeddingColumn = await resolveWriteColumnForEngine(engine);
staleChars = await engine.sumStaleChunkChars({
signature: currentEmbeddingSignature(),
...(embeddingColumn && { embeddingColumn }),
});
} catch {
staleChars = 0;
}
+5
View File
@@ -61,6 +61,7 @@ import {
type SynopsisFailureKind,
} from './audit-synopsis.ts';
import type { BrainEngine } from './engine.ts';
import { resolveWriteColumnForEngine } from './search/embedding-column.ts';
import type { ChunkInput, CRMode, Page } from './types.ts';
import type { SourceRow } from './sources-ops.ts';
@@ -286,9 +287,13 @@ export async function reembedPageWithContextualRetrieval(
// ── PHASE 2: single DB transaction ───────────────────────────
try {
// #1262: contextual re-embeds write TEXT embeddings — thread the
// caller-resolved write column like every other embed path.
const embeddingColumn = await resolveWriteColumnForEngine(args.engine);
await args.engine.transaction(async (tx) => {
await tx.upsertChunks(args.pageSlug, phase1.embeddedChunks, {
sourceId: args.sourceId,
...(embeddingColumn && { embeddingColumn }),
});
await tx.updatePageContextualRetrievalState(
args.pageSlug,
+13 -2
View File
@@ -18,7 +18,7 @@
*/
import type { BrainEngine } from './engine.ts';
import type { ChunkInput } from './types.ts';
import type { ChunkInput, ResolvedColumn } from './types.ts';
import { embedBatchWithBackoff } from '../commands/embed.ts';
import { type DbPacer, createNoopPacer, observed } from './db-pacer.ts';
import { AbortError } from './abort-check.ts';
@@ -61,6 +61,13 @@ export interface EmbedStaleOpts {
* Omit to keep the legacy `embedding IS NULL`-only behavior.
*/
embeddingSignature?: string;
/**
* #1262: caller-resolved write-side embedding column. Threaded into BOTH
* listStaleChunks (staleness predicate) and upsertChunks (write target) so
* an alt-column brain converges instead of re-selecting embedded rows.
* Resolve at the boundary via `resolveWriteColumnForEngine()`.
*/
embeddingColumn?: ResolvedColumn;
/**
* DB-contention pacer (paced-backfill). When enabled it (a) supplies the
* worker count via the caller passing `concurrency = bundle.maxConcurrency`
@@ -156,6 +163,7 @@ export async function embedStaleForSource(
afterPageId,
afterChunkIndex,
sourceId,
...(opts.embeddingColumn && { embeddingColumn: opts.embeddingColumn }),
}),
);
if (batch.length === 0) {
@@ -223,7 +231,10 @@ export async function embedStaleForSource(
doc_comment: c.doc_comment ?? undefined,
symbol_name_qualified: c.symbol_name_qualified ?? undefined,
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, { sourceId: keySourceId }));
await observed(pacer, () => engine.upsertChunks(slug, merged, {
sourceId: keySourceId,
...(opts.embeddingColumn && { embeddingColumn: opts.embeddingColumn }),
}));
// v0.41.31: stamp provenance only when EVERY chunk was stale (fully
// re-embedded this pass) — a partially-stale page keeps preserved
// chunks of unknown provenance, so don't claim current. After the
+22 -3
View File
@@ -12,6 +12,7 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
CodeEdgeInput, CodeEdgeResult,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
@@ -987,8 +988,13 @@ export interface BrainEngine {
* — Postgres rolls back automatically on conn drop, so commit-ambiguous
* failure replays to the same end state. Callers MUST NOT wrap externally;
* see {@link BatchOpts} retry-contract block.
*
* `opts.embeddingColumn` (optional) selects the content_chunks column that
* receives TEXT embeddings (#1262). The caller resolves the descriptor at
* the import/embed boundary via `resolveWriteColumn()`; engines never read
* config or choose columns themselves. Omitted => legacy `embedding`.
*/
upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void>;
upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void>;
/**
* Read every chunk for a page. `opts.sourceId` source-scopes the page
* lookup; without it, multi-source brains return chunks from every
@@ -1005,8 +1011,13 @@ export interface BrainEngine {
* counts across every source in the brain. Operators running
* `gbrain embed --stale --source media-corpus` expect only that
* source's NULLs touched; the caller threads `sourceId` here.
*
* `opts.embeddingColumn` switches the staleness predicate from the legacy
* `embedding` column to the resolved write-side column, so alt-column
* brains do not perpetually re-select rows whose target column is already
* populated (#1262). Must match the eventual upsertChunks target.
*/
countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number>;
countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number>;
/**
* Sum of LENGTH(chunk_text) over stale chunks — the character-count
* backlog the embed phase / embed-backfill will process. Sibling of
@@ -1020,8 +1031,13 @@ export interface BrainEngine {
* model signature (a model/dims swap). NULL signature is GRANDFATHERED
* (never counted) so the post-migration corpus isn't flagged en masse.
* Omit `signature` for the legacy `embedding IS NULL`-only count.
*
* `opts.embeddingColumn` switches the staleness predicate to the resolved
* write-side column (#1262) — same contract as countStaleChunks — so the
* sync cost gate doesn't count an alt-column brain's fully-embedded corpus
* as phantom backlog.
*/
sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number>;
sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number>;
/**
* Stamp `pages.embedding_signature = signature` for one page. Called after
* a page's chunks are (re)embedded so a later model swap can detect it as
@@ -1069,6 +1085,9 @@ export interface BrainEngine {
// both round-trip TIMESTAMPTZ as Date | string; ISO string is the
// common denominator on the wire).
afterUpdatedAt?: string | null;
// #1262: staleness predicate targets this column when set (must match
// countStaleChunks and the eventual upsertChunks write target).
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]>;
/**
* Delete every chunk for a page. Internal page-id lookup is sourceId-scoped
+25 -4
View File
@@ -10,7 +10,8 @@ import { findChunkForOffset } from './chunkers/edge-extractor.ts';
import { extractCodeRefs, imageOfCandidates } from './link-extraction.ts';
import { embedBatch, embedMultimodal, currentEmbeddingSignature } from './embedding.ts';
import { slugifyPath, slugifyCodePath, isCodeFilePath } from './sync.ts';
import type { ChunkInput, PageInput, PageType } from './types.ts';
import type { ChunkInput, PageInput, PageType, ResolvedColumn } from './types.ts';
import { resolveWriteColumnForEngine } from './search/embedding-column.ts';
import { computeEffectiveDate } from './effective-date.ts';
import { MARKDOWN_CHUNKER_VERSION } from './chunkers/recursive.ts';
import { logSlugFallback } from './audit-slug-fallback.ts';
@@ -740,6 +741,14 @@ export async function importFromContent(
// schema DEFAULT — required for multi-source brains; harmless ('default')
// for single-source callers.
const txOpts = sourceId ? { sourceId } : undefined;
// #1262: resolve the write-side embedding column once (merged config +
// gateway model) BEFORE the transaction; the descriptor rides only on
// upsertChunks so text embeddings land in the registered column.
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(sourceId || chunkWriteColumn)
? { ...(sourceId && { sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
await engine.transaction(async (tx) => {
if (existing) await tx.createVersion(slug, txOpts);
@@ -824,7 +833,7 @@ export async function importFromContent(
}
if (chunks.length > 0) {
await tx.upsertChunks(slug, chunks, txOpts);
await tx.upsertChunks(slug, chunks, chunkOpts);
// v0.41.31: stamp embedding provenance when this import actually
// embedded (not --no-embed), so a later model/dims swap is detectable
// as stale via embed --stale. The deferred/backfill + per-slug embed
@@ -1064,6 +1073,12 @@ export async function importCodeFile(
const title = `${relativePath} (${lang})`;
const sourceId = opts.sourceId;
const txOpts = sourceId ? { sourceId } : undefined;
// #1262: write-side embedding column descriptor (rides only on upsertChunks).
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(sourceId || chunkWriteColumn)
? { ...(sourceId && { sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
const byteLength = Buffer.byteLength(content, 'utf-8');
if (byteLength > MAX_FILE_SIZE) {
@@ -1183,7 +1198,7 @@ export async function importCodeFile(
await tx.addTag(slug, lang, txOpts);
if (chunks.length > 0) {
await tx.upsertChunks(slug, chunks, txOpts);
await tx.upsertChunks(slug, chunks, chunkOpts);
// v0.41.31: stamp embedding provenance ONLY when every chunk was
// freshly embedded with the current model this call (no reuse-by-hash
// carrying old-model vectors). Mixed pages stay unstamped rather than
@@ -1332,6 +1347,12 @@ export async function withImportTransaction(
): Promise<void> {
const sourceId = spec.sourceId ?? 'default';
const txOpts = spec.sourceId ? { sourceId: spec.sourceId } : undefined;
// #1262: write-side embedding column descriptor (rides only on upsertChunks).
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(spec.sourceId || chunkWriteColumn)
? { ...(spec.sourceId && { sourceId: spec.sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
await engine.transaction(async (tx) => {
if (spec.hadExisting) await tx.createVersion(spec.slug, txOpts);
await tx.putPage(spec.slug, spec.page, txOpts);
@@ -1347,7 +1368,7 @@ export async function withImportTransaction(
}
if (spec.chunks !== undefined) {
if (spec.chunks.length > 0) {
await tx.upsertChunks(spec.slug, spec.chunks, txOpts);
await tx.upsertChunks(spec.slug, spec.chunks, chunkOpts);
} else {
await tx.deleteChunks(spec.slug, txOpts);
}
@@ -35,6 +35,7 @@ import { tryAcquireDbLock } from '../../db-lock.ts';
import { BudgetTracker, BudgetExhausted } from '../../budget/budget-tracker.ts';
import { withBudgetTracker } from '../../ai/gateway.ts';
import { embedStaleForSource } from '../../embed-stale.ts';
import { resolveWriteColumnForEngine } from '../../search/embedding-column.ts';
import { currentEmbeddingSignature } from '../../embedding.ts';
import { type DbPacer, createDbPacer, createNoopPacer } from '../../db-pacer.ts';
import { resolvePaceMode, loadPaceModeConfig, readPaceEnv } from '../../pace-mode.ts';
@@ -164,12 +165,16 @@ export function makeEmbedBackfillHandler(engine: BrainEngine) {
// the supervisor, so pacing it is the headline win.
const { pacer, concurrency } = await resolveBackfillPacer(engine, job.data);
// #1262: resolve the write-side embedding column once at the job boundary.
const embeddingColumn = await resolveWriteColumnForEngine(engine);
try {
const result = await withBudgetTracker(tracker, async () =>
embedStaleForSource(engine, sourceId, {
batchSize,
signal: job.signal,
pacer,
...(embeddingColumn && { embeddingColumn }),
...(concurrency !== undefined && { concurrency }),
// v0.41.31: re-embed pages whose model signature drifted + stamp
// provenance as chunks land.
+42 -22
View File
@@ -40,6 +40,7 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
SalienceOpts, SalienceResult, AnomaliesOpts, AnomalyResult,
@@ -2230,12 +2231,20 @@ export class PGLiteEngine implements BrainEngine {
}
// Chunks
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void> {
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void> {
return this.batchRetry(opts?.auditSite ?? 'upsertChunks', opts?.signal, () => this._upsertChunksOnce(slug, chunks, opts), chunks.length);
}
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn }): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
// #1262: caller-resolved write target for TEXT embeddings. Descriptor
// names are identifier-validated + quoted by buildVectorCastFragment;
// omitted => legacy `embedding vector`. Mirrors postgres-engine.ts.
const targetFragment = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn)
: undefined;
const targetCol = targetFragment?.col ?? 'embedding';
const embeddingCast = targetFragment?.castSql.replace('$1::', '') ?? 'vector';
// Source-scope the page-id lookup so duplicate slugs in different sources
// do not return multiple rows or target the wrong page.
@@ -2270,7 +2279,7 @@ export class PGLiteEngine implements BrainEngine {
// list. Image chunks pass embedding=null + embedding_image=Float32Array
// (1024-dim Voyage). Text/code chunks pass embedding=Float32Array +
// embedding_image=null. Default modality='text' when omitted.
const cols = '(page_id, chunk_index, chunk_text, chunk_source, embedding, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)';
const cols = `(page_id, chunk_index, chunk_text, chunk_source, ${targetCol}, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)`;
const rowParts: string[] = [];
const params: unknown[] = [];
let paramIdx = 1;
@@ -2288,7 +2297,7 @@ export class PGLiteEngine implements BrainEngine {
const modality = chunk.modality ?? 'text';
// Inline ::vector NULL literals to avoid a per-branch placeholder.
const embeddingPh = embeddingStr ? `$${paramIdx++}::vector` : 'NULL';
const embeddingPh = embeddingStr ? `$${paramIdx++}::${embeddingCast}` : 'NULL';
const embeddedAtPh = embeddingStr ? 'now()' : 'NULL';
const embeddingImagePh = embeddingImageStr ? `$${paramIdx++}::vector` : 'NULL';
@@ -2327,19 +2336,19 @@ export class PGLiteEngine implements BrainEngine {
ON CONFLICT (page_id, chunk_index) DO UPDATE SET
chunk_text = EXCLUDED.chunk_text,
chunk_source = EXCLUDED.chunk_source,
embedding = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.embedding
WHEN content_chunks.embedding IS NULL THEN EXCLUDED.embedding
${targetCol} = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.${targetCol}
WHEN content_chunks.${targetCol} IS NULL THEN EXCLUDED.${targetCol}
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.embedding
ELSE content_chunks.embedding
THEN EXCLUDED.${targetCol}
ELSE content_chunks.${targetCol}
END,
model = COALESCE(EXCLUDED.model, content_chunks.model),
token_count = EXCLUDED.token_count,
embedded_at = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.embedding IS NULL THEN NULL
WHEN content_chunks.embedding IS NULL AND EXCLUDED.embedding IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.${targetCol} IS NULL THEN NULL
WHEN content_chunks.${targetCol} IS NULL AND EXCLUDED.${targetCol} IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.embedded_at
@@ -2377,14 +2386,19 @@ export class PGLiteEngine implements BrainEngine {
* drift (NULL grandfathered never stale). Shared by countStaleChunks +
* sumStaleChunkChars so they can't drift.
*/
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string }): { where: string; params: unknown[] } {
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): { where: string; params: unknown[] } {
// #1262: staleness targets the caller-resolved write column when set
// (identifier-validated + quoted); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
const params: unknown[] = [];
const conds: string[] = [];
if (opts?.signature !== undefined) {
params.push(opts.signature);
conds.push(`(cc.embedding IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
conds.push(`(cc.${staleCol} IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
} else {
conds.push(`cc.embedding IS NULL`);
conds.push(`cc.${staleCol} IS NULL`);
}
conds.push(`NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')`);
if (opts?.sourceId !== undefined) {
@@ -2394,7 +2408,7 @@ export class PGLiteEngine implements BrainEngine {
return { where: conds.join(' AND '), params };
}
async countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number> {
async countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
// D7: source-scoped count for `gbrain embed --stale --source X`. Always
// JOIN pages so embed-skip + signature predicates apply. PGLite is
// PostgreSQL 17.5 in WASM and supports the full JSONB operator set.
@@ -2410,7 +2424,7 @@ export class PGLiteEngine implements BrainEngine {
return Number(count);
}
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number> {
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
// Sibling of countStaleChunks: same stale predicate, summing chunk_text
// length for the sync cost preview. ::bigint guards int4 overflow.
const { where, params } = this.buildStaleChunkWhere(opts);
@@ -2463,11 +2477,17 @@ export class PGLiteEngine implements BrainEngine {
sourceId?: string;
orderBy?: 'page_id' | 'updated_desc';
afterUpdatedAt?: string | null;
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]> {
const limit = opts?.batchSize ?? 2000;
const afterPid = opts?.afterPageId ?? 0;
const afterIdx = opts?.afterChunkIndex ?? -1;
const orderBy = opts?.orderBy ?? 'page_id';
// #1262: staleness follows the caller-resolved write column (validated +
// quoted identifier); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
// v0.41.18.0 (A13, codex #9): --priority recent path. See postgres-engine
// sibling for full rationale. Same composite cursor + ORDER BY.
@@ -2481,7 +2501,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE cc.${staleCol} IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
LIMIT $1`,
@@ -2492,7 +2512,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE cc.${staleCol} IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
p.updated_at < $1::timestamptz
@@ -2511,7 +2531,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE cc.${staleCol} IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
@@ -2523,7 +2543,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE cc.${staleCol} IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
@@ -2548,7 +2568,7 @@ export class PGLiteEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE cc.${staleCol} IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > ($1, $2)
ORDER BY cc.page_id, cc.chunk_index
@@ -2562,7 +2582,7 @@ export class PGLiteEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE cc.${staleCol} IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > ($2, $3)
+43 -22
View File
@@ -50,6 +50,7 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
SalienceOpts, SalienceResult, AnomaliesOpts, AnomalyResult,
@@ -2380,13 +2381,21 @@ export class PostgresEngine implements BrainEngine {
}
// Chunks
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void> {
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void> {
return this.batchRetry(opts?.auditSite ?? 'upsertChunks', opts?.signal, () => this._upsertChunksOnce(slug, chunks, opts), chunks.length);
}
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn }): Promise<void> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
// #1262: caller-resolved write target for TEXT embeddings. Descriptor
// names are identifier-validated + quoted by buildVectorCastFragment;
// omitted => legacy `embedding vector`.
const targetFragment = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn)
: undefined;
const targetCol = targetFragment?.col ?? 'embedding';
const embeddingCast = targetFragment?.castSql.replace('$1::', '') ?? 'vector';
// Source-scope the page-id lookup. Without this filter, multi-source
// brains where the slug exists in 2+ sources return >1 row and the
@@ -2413,7 +2422,7 @@ export class PostgresEngine implements BrainEngine {
// scope metadata through upserts.
// v0.27.1 (Phase 8): added `modality` + `embedding_image` to the column
// list. Image chunks pass embedding=null + embedding_image=Float32Array.
const cols = '(page_id, chunk_index, chunk_text, chunk_source, embedding, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)';
const cols = `(page_id, chunk_index, chunk_text, chunk_source, ${targetCol}, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)`;
const rows: string[] = [];
const params: unknown[] = [];
let paramIdx = 1;
@@ -2430,7 +2439,7 @@ export class PostgresEngine implements BrainEngine {
: null;
const modality = chunk.modality ?? 'text';
const embeddingPh = embeddingStr ? `$${paramIdx++}::vector` : 'NULL';
const embeddingPh = embeddingStr ? `$${paramIdx++}::${embeddingCast}` : 'NULL';
const embeddedAtPh = embeddingStr ? 'now()' : 'NULL';
const embeddingImagePh = embeddingImageStr ? `$${paramIdx++}::vector` : 'NULL';
@@ -2478,19 +2487,19 @@ export class PostgresEngine implements BrainEngine {
ON CONFLICT (page_id, chunk_index) DO UPDATE SET
chunk_text = EXCLUDED.chunk_text,
chunk_source = EXCLUDED.chunk_source,
embedding = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.embedding
WHEN content_chunks.embedding IS NULL THEN EXCLUDED.embedding
${targetCol} = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.${targetCol}
WHEN content_chunks.${targetCol} IS NULL THEN EXCLUDED.${targetCol}
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.embedding
ELSE content_chunks.embedding
THEN EXCLUDED.${targetCol}
ELSE content_chunks.${targetCol}
END,
model = COALESCE(EXCLUDED.model, content_chunks.model),
token_count = EXCLUDED.token_count,
embedded_at = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.embedding IS NULL THEN NULL
WHEN content_chunks.embedding IS NULL AND EXCLUDED.embedding IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.${targetCol} IS NULL THEN NULL
WHEN content_chunks.${targetCol} IS NULL AND EXCLUDED.${targetCol} IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.embedded_at
@@ -2530,14 +2539,19 @@ export class PostgresEngine implements BrainEngine {
* embedding_signature drift (NULL grandfathered). Shared by
* countStaleChunks + sumStaleChunkChars (parity with the PGLite sibling).
*/
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string }): { where: string; params: unknown[] } {
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): { where: string; params: unknown[] } {
// #1262: staleness targets the caller-resolved write column when set
// (identifier-validated + quoted); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
const params: unknown[] = [];
const conds: string[] = [];
if (opts?.signature !== undefined) {
params.push(opts.signature);
conds.push(`(cc.embedding IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
conds.push(`(cc.${staleCol} IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
} else {
conds.push(`cc.embedding IS NULL`);
conds.push(`cc.${staleCol} IS NULL`);
}
conds.push(`NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')`);
if (opts?.sourceId !== undefined) {
@@ -2547,7 +2561,7 @@ export class PostgresEngine implements BrainEngine {
return { where: conds.join(' AND '), params };
}
async countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number> {
async countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
// Always JOIN pages so the embed_skip + signature predicates apply.
// D7: source_id scoping. v0.41.31: optional signature widens staleness
// to embedding_signature drift (NULL grandfathered).
@@ -2565,7 +2579,7 @@ export class PostgresEngine implements BrainEngine {
});
}
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number> {
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
// Sibling of countStaleChunks: same stale predicate, summing chunk_text
// length for the sync cost preview. ::bigint guards int4 overflow.
const { where, params } = this.buildStaleChunkWhere(opts);
@@ -2618,11 +2632,18 @@ export class PostgresEngine implements BrainEngine {
sourceId?: string;
orderBy?: 'page_id' | 'updated_desc';
afterUpdatedAt?: string | null;
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]> {
const limit = opts?.batchSize ?? 2000;
const afterPid = opts?.afterPageId ?? 0;
const afterIdx = opts?.afterChunkIndex ?? -1;
const orderBy = opts?.orderBy ?? 'page_id';
// #1262: staleness follows the caller-resolved write column (validated +
// quoted identifier); legacy `embedding` otherwise. Interpolated below as
// an unsafe FRAGMENT (identifiers can't be bound parameters).
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
// RLS scope binding (opt-in via GBRAIN_RLS_SCOPE_BINDING).
return await this.withScopedReadTransaction(undefined, opts?.sourceId, async (tx) => {
@@ -2639,7 +2660,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
LIMIT ${limit}
@@ -2649,7 +2670,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
p.updated_at < ${afterUpdated}::timestamptz
@@ -2667,7 +2688,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
@@ -2678,7 +2699,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
@@ -2698,7 +2719,7 @@ export class PostgresEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > (${afterPid}, ${afterIdx})
ORDER BY cc.page_id, cc.chunk_index
@@ -2711,7 +2732,7 @@ export class PostgresEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.embedding IS NULL
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > (${afterPid}, ${afterIdx})
+74
View File
@@ -443,6 +443,80 @@ export function resolveEmbeddingColumn(
};
}
/**
* Resolves the WRITE-side embedding column for the currently configured
* embedding model (#1262). The read-side resolver above answers "which
* column does this query search?"; this one answers "which column should
* newly produced text embeddings land in?".
*
* Unlike read-side search, writes take no per-call column override. The
* import/embed boundary resolves once from merged config + gateway state
* and passes the descriptor into `engine.upsertChunks`; engines stay
* config-free (same contract as the read-side descriptor).
*
* Behavior:
* - no user-declared `embedding_columns` => undefined (legacy brain,
* writes keep targeting the default `embedding` column)
* - a user-declared entry whose `provider` matches the current
* embedding model => that entry's descriptor
* - no provider match => undefined (fall back to legacy `embedding`)
*
* Only USER-declared entries are consulted — never the cfg-derived
* builtins. The `embedding_image` builtin's provider is the multimodal
* model; matching it here would misroute text embeddings into the image
* column. The no-match fallback is intentional: switching models before
* registering a matching column must not silently write vectors into an
* arbitrary column.
*/
export function resolveWriteColumn(cfg: GBrainConfig): ResolvedColumn | undefined {
const userColumns = cfg.embedding_columns;
if (
!userColumns ||
typeof userColumns !== 'object' ||
Array.isArray(userColumns) ||
Object.keys(userColumns).length === 0
) {
return undefined;
}
// Same model-resolution chain as the registry builtin: cfg > gateway > default.
let gwModel: string | undefined;
try {
const gw = require('../ai/gateway.ts') as typeof import('../ai/gateway.ts');
gwModel = gw.getEmbeddingModel();
} catch {
// Gateway unconfigured — fall through to the canonical default.
}
const currentModel = cfg.embedding_model ?? gwModel ?? DEFAULT_EMBEDDING_MODEL;
for (const [name, entry] of Object.entries(userColumns)) {
if (!entry) continue;
validateColumnKey(name);
validateColumnConfig(name, entry);
if (entry.provider !== currentModel) continue;
return {
name,
type: entry.type,
dimensions: entry.dimensions,
embeddingModel: entry.provider,
};
}
return undefined;
}
/**
* Engine-boundary convenience: merged config (file/env + DB plane) →
* resolveWriteColumn. Dynamic import keeps config.ts out of this module's
* static graph (mirrors the gateway require above).
*/
export async function resolveWriteColumnForEngine(
engine: { getConfig(key: string): Promise<string | null | undefined> },
): Promise<ResolvedColumn | undefined> {
const { loadConfigWithEngine } = await import('../config.ts');
const cfg = await loadConfigWithEngine(engine);
return cfg ? resolveWriteColumn(cfg) : undefined;
}
/**
* True when the resolved column is the default `embedding` name.
* Name-based check; does not compare embedding space.
+133
View File
@@ -241,3 +241,136 @@ describe('buildVectorCastFragment — engine SQL composer (D3)', () => {
expect(castSql).toBe('$1::halfvec(2560)');
});
});
describe('PGLite engine: upsertChunks write-side ResolvedColumn descriptor (#1262)', () => {
test('halfvec descriptor writes the text embedding to the alternate column, not legacy embedding', async () => {
await engine.putPage('docs/write-alt-pglite', {
type: 'concept',
title: 'Write alt column PGLite',
compiled_truth: 'PGLite write-side alternate embedding column test.',
});
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
await engine.upsertChunks('docs/write-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'PGLite write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{
has_default: boolean;
has_ze: boolean;
has_embedded_at: boolean;
}>(
`SELECT embedding IS NOT NULL AS has_default,
embedding_ze IS NOT NULL AS has_ze,
embedded_at IS NOT NULL AS has_embedded_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-pglite'`,
);
expect(rows.length).toBe(1);
expect(rows[0].has_default).toBe(false);
expect(rows[0].has_ze).toBe(true);
expect(rows[0].has_embedded_at).toBe(true);
});
test('text-unchanged re-upsert without a vector preserves the alternate-column embedding', async () => {
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
// Same chunk_text, no embedding: the ON CONFLICT CASE must keep the
// existing alternate-column vector (D24 semantics follow the column).
await engine.upsertChunks('docs/write-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'PGLite write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{ has_ze: boolean }>(
`SELECT embedding_ze IS NOT NULL AS has_ze
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-pglite'`,
);
expect(rows).toEqual([{ has_ze: true }]);
});
});
describe('PGLite: embed --stale converges on an alt-column brain (#1262)', () => {
test('boundary resolves the write column; stale scan does not re-select embedded rows', async () => {
const { runEmbedCore } = await import('../../src/commands/embed.ts');
const local = new PGLiteEngine();
const previousHome = process.env.GBRAIN_HOME;
process.env.GBRAIN_HOME = `/tmp/gbrain-write-col-stale-${Date.now()}`;
try {
await local.connect({});
await local.initSchema();
await (local as any).db.exec(
`ALTER TABLE content_chunks ADD COLUMN IF NOT EXISTS embedding_ze halfvec(2560)`,
);
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
await local.setConfig('embedding_columns', JSON.stringify({
embedding_ze: { provider: 'zeroentropyai:zembed-1', dimensions: 2560, type: 'halfvec' },
}));
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
await local.putPage('docs/stale-alt-pglite', {
type: 'concept',
title: 'Dynamic stale column',
compiled_truth: 'A chunk that is embedded only in the dynamic column.',
});
await local.upsertChunks('docs/stale-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'A chunk that is embedded only in the dynamic column.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
// Engine-level contrast: legacy predicate still sees the row as stale;
// the alt-column predicate does not.
expect(await local.countStaleChunks()).toBe(1);
expect(await local.countStaleChunks({ embeddingColumn: descriptor })).toBe(0);
// sumStaleChunkChars feeds the sync cost gate — same predicate contract.
expect(await local.sumStaleChunkChars()).toBeGreaterThan(0);
expect(await local.sumStaleChunkChars({ embeddingColumn: descriptor })).toBe(0);
expect(await local.listStaleChunks({ embeddingColumn: descriptor, batchSize: 100 })).toHaveLength(0);
expect(await local.listStaleChunks({ batchSize: 100 })).toHaveLength(1);
// Boundary-level: `embed --stale --dry-run` resolves the write column
// from merged config + gateway and reports NOTHING to embed. Without
// the fix this reports 1 (perpetual re-embed loop).
const result = await runEmbedCore(local, { stale: true, dryRun: true });
expect(result.would_embed).toBe(0);
} finally {
await local.disconnect();
if (previousHome === undefined) delete process.env.GBRAIN_HOME;
else process.env.GBRAIN_HOME = previousHome;
resetGateway();
}
});
});
@@ -224,4 +224,54 @@ if (!dbUrl) {
await engine.executeRaw(`UPDATE content_chunks SET embedding_voyage = '${v}'::vector WHERE id = ${dogId}`);
});
});
describe('Postgres: upsertChunks write-side ResolvedColumn descriptor (#1262)', () => {
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
test('halfvec descriptor writes the text embedding to the alternate column, not legacy embedding', async () => {
await engine.putPage('docs/write-alt-postgres', {
type: 'concept',
title: 'Write alt column Postgres',
compiled_truth: 'Postgres write-side alternate embedding column test.',
});
await engine.upsertChunks('docs/write-alt-postgres', [
{
chunk_index: 0,
chunk_text: 'Postgres write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{
has_default: boolean;
has_ze: boolean;
}>(
`SELECT embedding IS NOT NULL AS has_default,
embedding_ze IS NOT NULL AS has_ze
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-postgres'`,
);
expect(rows.length).toBe(1);
expect(rows[0].has_default).toBe(false);
expect(rows[0].has_ze).toBe(true);
}, 30_000);
test('stale scan follows the write-side column (count + list parity with the write target)', async () => {
// Legacy predicate: cat/dog/write-alt rows all have embedding NULL.
expect(await engine.countStaleChunks()).toBeGreaterThan(0);
// Alt-column predicate: every chunk has embedding_ze populated.
expect(await engine.countStaleChunks({ embeddingColumn: descriptor })).toBe(0);
expect(await engine.listStaleChunks({ embeddingColumn: descriptor, batchSize: 100 })).toHaveLength(0);
expect((await engine.listStaleChunks({ batchSize: 100 })).length).toBeGreaterThan(0);
// updated_desc arm uses the same predicate.
expect(await engine.listStaleChunks({ embeddingColumn: descriptor, orderBy: 'updated_desc', batchSize: 100 })).toHaveLength(0);
}, 30_000);
});
}
@@ -1,61 +0,0 @@
/**
* E2E smoke for skills/obsidian-gbrain-safe-index.
*
* Verifies the from-trigger-to-side-effect path that skillify requires:
* a real user trigger phrase routes to the skill, the resolver/check
* pipeline treats it as reachable, and the skill file exposes the
* gbrain commands the workflow actually runs.
*
* This stays local-only (no paid embedding, no external API): it asserts
* the documented safe/import path is present and parseable, not that it
* mutates a live brain.
*/
import { describe, expect, it } from 'bun:test';
import { existsSync, readFileSync } from 'fs';
import { join } from 'path';
const SKILLS = join(import.meta.dir, '..', '..', 'skills');
const SKILL_MD = join(SKILLS, 'obsidian-gbrain-safe-index', 'SKILL.md');
const RESOLVER = join(SKILLS, 'RESOLVER.md');
const TRIGGER_PHRASES = [
'connect my Obsidian vault to gbrain',
'import my vault to gbrain',
'sync vault and gbrain',
'capture this skill in my vault',
'embed gbrain after vault update',
'is gbrain synced with my vault',
];
describe('obsidian-gbrain-safe-index E2E', () => {
it('resolver maps real trigger phrasings to the skill', () => {
const resolver = readFileSync(RESOLVER, 'utf-8');
expect(resolver).toContain('obsidian-gbrain-safe-index/SKILL.md');
// Each representative phrase shares a token substring with a resolver row.
const rows = resolver
.split('\n')
.filter((l) => l.includes('obsidian-gbrain-safe-index/SKILL.md'))
.join('\n');
for (const phrase of TRIGGER_PHRASES) {
const hit = phrase
.toLowerCase()
.split(/\s+/)
.some((tok) => tok.length > 3 && rows.toLowerCase().includes(tok));
expect(hit, `no resolver token for: ${phrase}`).toBe(true);
}
});
it('skill documents the gbrain import-first safe path', () => {
const body = readFileSync(SKILL_MD, 'utf-8');
expect(body).toContain('gbrain import');
expect(body).toContain('--no-embed');
expect(body).toContain('gbrain config set search.mode conservative');
// Paid gate must be explicit, not a silent default.
expect(body).toContain('gbrain embed --stale');
});
it('skill is reachable from the skill tree', () => {
expect(existsSync(SKILL_MD)).toBe(true);
});
});
@@ -1,52 +0,0 @@
/**
* E2E smoke for skills/skill-vault-capture-policy.
*
* Verifies the from-trigger-to-side-effect path: a real capture request
* routes to the skill, and the skill documents the vault navigation
* obligations (index.md / log.md) that make a capture durable.
*
* Local-only: it asserts documented behavior, not live vault writes.
*/
import { describe, expect, it } from 'bun:test';
import { existsSync, readFileSync } from 'fs';
import { join } from 'path';
const SKILLS = join(import.meta.dir, '..', '..', 'skills');
const SKILL_MD = join(SKILLS, 'skill-vault-capture-policy', 'SKILL.md');
const RESOLVER = join(SKILLS, 'RESOLVER.md');
const TRIGGER_PHRASES = [
'save this learning to the vault',
'capture this skill in Obsidian',
'record this workflow in my notes',
'put this setup change in the knowledge base',
];
describe('skill-vault-capture-policy E2E', () => {
it('resolver maps real capture phrasings to the skill', () => {
const resolver = readFileSync(RESOLVER, 'utf-8');
expect(resolver).toContain('skill-vault-capture-policy/SKILL.md');
const rows = resolver
.split('\n')
.filter((l) => l.includes('skill-vault-capture-policy/SKILL.md'))
.join('\n');
for (const phrase of TRIGGER_PHRASES) {
const hit = phrase
.toLowerCase()
.split(/\s+/)
.some((tok) => tok.length > 3 && rows.toLowerCase().includes(tok));
expect(hit, `no resolver token for: ${phrase}`).toBe(true);
}
});
it('skill documents index.md and log.md update obligations', () => {
const body = readFileSync(SKILL_MD, 'utf-8');
expect(body).toContain('index.md');
expect(body).toContain('log.md');
});
it('skill is reachable from the skill tree', () => {
expect(existsSync(SKILL_MD)).toBe(true);
});
});
-51
View File
@@ -1,51 +0,0 @@
import { describe, expect, it } from 'bun:test';
import { readFileSync, existsSync } from 'fs';
import { join } from 'path';
const SKILL_DIR = join(import.meta.dir, '..', 'skills', 'obsidian-gbrain-safe-index');
const SKILL_MD = join(SKILL_DIR, 'SKILL.md');
const RESOLVER = join(import.meta.dir, '..', 'skills', 'RESOLVER.md');
function parseFrontmatter(raw: string): Record<string, unknown> {
const m = raw.match(/^---\n([\s\S]*?)\n---/);
if (!m) throw new Error('no frontmatter');
const out: Record<string, unknown> = {};
for (const line of m[1].split('\n')) {
const mm = line.match(/^([a-zA-Z_]+):\s*(.*)$/);
if (mm) out[mm[1]] = mm[2].trim();
}
return out;
}
describe('obsidian-gbrain-safe-index skill', () => {
it('has a SKILL.md with required frontmatter', () => {
expect(existsSync(SKILL_MD)).toBe(true);
const fm = parseFrontmatter(readFileSync(SKILL_MD, 'utf-8'));
expect(fm['name']).toBe('obsidian-gbrain-safe-index');
expect(fm['description']).toBeTruthy();
});
it('has the required conformance sections', () => {
const body = readFileSync(SKILL_MD, 'utf-8');
for (const section of ['## Contract', '## Phases', '## Output Format', '## Anti-Patterns']) {
expect(body.includes(section), `missing ${section}`).toBe(true);
}
});
it('is registered in RESOLVER.md', () => {
expect(existsSync(RESOLVER)).toBe(true);
const resolver = readFileSync(RESOLVER, 'utf-8');
expect(resolver.includes('obsidian-gbrain-safe-index/SKILL.md')).toBe(true);
});
it('has routing-eval fixtures that exercise real trigger phrasings', () => {
const evalPath = join(SKILL_DIR, 'routing-eval.jsonl');
expect(existsSync(evalPath)).toBe(true);
const lines = readFileSync(evalPath, 'utf-8')
.split('\n')
.filter((l) => l.trim() && !l.trim().startsWith('//'))
.map((l) => JSON.parse(l));
const positives = lines.filter((l) => l.expected_skill === 'obsidian-gbrain-safe-index');
expect(positives.length).toBeGreaterThanOrEqual(5);
});
});
+110 -1
View File
@@ -13,9 +13,10 @@
* throw on unknown string.
*/
import { describe, test, expect } from 'bun:test';
import { describe, test, expect, afterAll, afterEach } from 'bun:test';
import {
resolveEmbeddingColumn,
resolveWriteColumn,
getEmbeddingColumnRegistry,
buildVectorCastFragment,
quoteIdentifier,
@@ -34,6 +35,28 @@ import {
} from '../../src/core/search/embedding-column.ts';
import type { GBrainConfig } from '../../src/core/config.ts';
import type { ResolvedColumn } from '../../src/core/types.ts';
import { configureGateway, resetGateway } from '../../src/core/ai/gateway.ts';
/**
* Teardown: reset AND re-apply the legacy preload config
* (test/helpers/legacy-embedding-preload.ts). A bare resetGateway() would
* leave the slot empty for the NEXT file's beforeAll (the preload's
* per-test beforeEach only fires before tests, not before beforeAll), which
* would make sibling PGLite fixtures initSchema at the 1280 default instead
* of the legacy 1536 their seed vectors assume.
*/
function restorePreloadGateway() {
resetGateway();
configureGateway({
embedding_model: 'openai:text-embedding-3-large',
embedding_dimensions: 1536,
env: { ...process.env },
});
}
afterAll(() => {
restorePreloadGateway();
});
function cfg(overrides: Partial<GBrainConfig> = {}): GBrainConfig {
return { engine: 'pglite', ...overrides };
@@ -522,3 +545,89 @@ describe('codex /ship #4 — isCacheSafe (embedding-space-based skip)', () => {
expect(isCacheSafe(r, cfg())).toBe(true);
});
});
describe('resolveWriteColumn — write-side boundary resolution (#1262)', () => {
afterEach(() => {
restorePreloadGateway();
});
test('no registry / empty registry returns undefined (legacy single-column brain)', () => {
expect(resolveWriteColumn(cfg())).toBeUndefined();
expect(resolveWriteColumn(cfg({ embedding_columns: {} }))).toBeUndefined();
});
test('provider match via cfg.embedding_model returns the descriptor', () => {
const r = resolveWriteColumn(cfg({
embedding_model: 'voyage:voyage-3-large',
embedding_dimensions: 1024,
embedding_columns: {
embedding_voyage: { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
},
}));
expect(r).toEqual({
name: 'embedding_voyage',
type: 'vector',
dimensions: 1024,
embeddingModel: 'voyage:voyage-3-large',
});
});
test('provider match via gateway state (cfg.embedding_model unset) returns descriptor', () => {
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_ze: { provider: 'zeroentropyai:zembed-1', dimensions: 2560, type: 'halfvec' },
},
}));
expect(r).toEqual({
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
});
});
test('no provider match returns undefined instead of guessing a column', () => {
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_voyage: { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
},
}));
expect(r).toBeUndefined();
});
test('only USER-declared columns are consulted — multimodal builtin never captures text writes', () => {
// Current model equals the embedding_image BUILTIN's provider; a registry
// walk that consulted builtins would misroute text writes into the image
// column. resolveWriteColumn must return undefined here.
configureGateway({
embedding_model: 'voyage:voyage-multimodal-3',
embedding_dimensions: 1024,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_other: { provider: 'openai:text-embedding-3-large', dimensions: 1536, type: 'vector' },
},
}));
expect(r).toBeUndefined();
});
test('malformed registry entry throws loud (same validation as the read side)', () => {
expect(() => resolveWriteColumn(cfg({
embedding_model: 'voyage:voyage-3-large',
embedding_columns: {
'bad"col': { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
} as never,
}))).toThrow(EmbeddingColumnConfigError);
});
});
-64
View File
@@ -1,64 +0,0 @@
import { describe, expect, it } from 'bun:test';
import { readFileSync, existsSync } from 'fs';
import { join } from 'path';
import { DEFAULT_PRIVATE_PATTERNS } from '../src/core/skillpack/harvest-lint.ts';
const SKILL_DIR = join(import.meta.dir, '..', 'skills', 'skill-vault-capture-policy');
const SKILL_MD = join(SKILL_DIR, 'SKILL.md');
const RESOLVER = join(import.meta.dir, '..', 'skills', 'RESOLVER.md');
function parseFrontmatter(raw: string): Record<string, unknown> {
const m = raw.match(/^---\n([\s\S]*?)\n---/);
if (!m) throw new Error('no frontmatter');
const out: Record<string, unknown> = {};
for (const line of m[1].split('\n')) {
const mm = line.match(/^([a-zA-Z_]+):\s*(.*)$/);
if (mm) out[mm[1]] = mm[2].trim();
}
return out;
}
describe('skill-vault-capture-policy skill', () => {
it('has a SKILL.md with required frontmatter', () => {
expect(existsSync(SKILL_MD)).toBe(true);
const fm = parseFrontmatter(readFileSync(SKILL_MD, 'utf-8'));
expect(fm['name']).toBe('skill-vault-capture-policy');
expect(fm['description']).toBeTruthy();
});
it('has the required conformance sections', () => {
const body = readFileSync(SKILL_MD, 'utf-8');
for (const section of ['## Contract', '## Phases', '## Output Format', '## Anti-Patterns']) {
expect(body.includes(section), `missing ${section}`).toBe(true);
}
});
it('is registered in RESOLVER.md', () => {
expect(existsSync(RESOLVER)).toBe(true);
expect(readFileSync(RESOLVER, 'utf-8').includes('skill-vault-capture-policy/SKILL.md')).toBe(true);
});
it('contains no private user or agent-fork names (privacy rule)', () => {
for (const file of [SKILL_MD, join(SKILL_DIR, 'routing-eval.jsonl')]) {
const body = readFileSync(file, 'utf-8');
// DEFAULT_PRIVATE_PATTERNS[0] is the banned fork-name pattern; sourced
// from harvest-lint so this file never contains the literal itself
// (scripts/check-privacy.sh would reject it).
const forkName = new RegExp(DEFAULT_PRIVATE_PATTERNS[0], 'i');
for (const name of [/\bAdam\b/, /\bHermes\b/, /\bHerdr\b/, /\bArk\b/, forkName]) {
expect(name.test(body), `private name ${name} in ${file}`).toBe(false);
}
}
});
it('has routing-eval fixtures', () => {
const evalPath = join(SKILL_DIR, 'routing-eval.jsonl');
expect(existsSync(evalPath)).toBe(true);
const positives = readFileSync(evalPath, 'utf-8')
.split('\n')
.filter((l) => l.trim() && !l.trim().startsWith('//'))
.map((l) => JSON.parse(l))
.filter((l) => l.expected_skill === 'skill-vault-capture-policy');
expect(positives.length).toBeGreaterThanOrEqual(4);
});
});