Compare commits

..
Author SHA1 Message Date
a11ec9c468 fix(dream): stamp incremental extraction watermark (#2636)
The Dream cycle disables sync's inline extraction and routes changed
slugs through extractForSlugs, which flushed link/timeline batches but
never stamped links_extracted_at — so incrementally extracted pages
stayed permanently visible to `extract --stale` / doctor.

Collect processedRefs per successfully processed page and stamp them
via stampExtracted (best-effort) after both batch flushes, non-dry-run
mode 'all' only. Source-id threading from the original PR #2637 already
landed on master via #1503/#1747, so this rebase carries only the
missing watermark stamp plus regression tests.

Takeover of #2637.

Co-authored-by: JavanC <JavanC@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:34:22 -07:00
4 changed files with 60 additions and 219 deletions
+12
View File
@@ -1025,6 +1025,10 @@ async function extractForSlugs(
let linksCreated = 0;
let timelineCreated = 0;
let pagesProcessed = 0;
// #2636: successfully processed pages get their extraction watermark
// stamped after the final flush (mode 'all' only — a partial-mode run
// hasn't done the full extraction the watermark asserts).
const processedRefs: Array<{ slug: string; source_id: string }> = [];
// Issue #972: read the basename flag once per extract run.
const globalBasename = await isGlobalBasenameEnabled(engine);
@@ -1113,6 +1117,7 @@ async function extractForSlugs(
}
pagesProcessed++;
if (!dryRun) processedRefs.push({ slug, source_id: sourceId ?? 'default' });
} catch { /* skip unreadable */ }
progress.tick(1);
},
@@ -1120,6 +1125,13 @@ async function extractForSlugs(
await flushLinks();
await flushTimeline();
// #2636: the Dream cycle disables sync's inline extraction and routes
// changed slugs through this incremental path — without a stamp here,
// those pages never get links_extracted_at and stay permanently visible
// to `extract --stale` / doctor. Stamp only after BOTH batches flushed.
if (!dryRun && mode === 'all') {
await stampExtracted(engine, processedRefs);
}
progress.finish();
if (!jsonMode) {
+2 -112
View File
@@ -1,100 +1,8 @@
import type { Recipe } from '../types.ts';
/**
* MiniMax transport shim (#1977). MiniMax's `/v1/embeddings` endpoint is NOT
* OpenAI-compatible at the wire level despite the recipe's
* `implementation: 'openai-compatible'`:
* - Request: requires `texts` (the AI SDK sends `input`) plus an optional
* `type: 'db' | 'query'` asymmetric-retrieval field, and rejects OpenAI's
* `encoding_format`.
* - Response: returns `{vectors: number[][], total_tokens}` where the AI
* SDK's Zod schema expects `{data: [{embedding, index}], usage}`.
*
* Chat (`/chat/completions`) IS OpenAI-compatible, and this same fetch is
* applied to every openai-compatible touchpoint by `applyOpenAICompatConfig`,
* so everything outside the embeddings path passes through untouched — and
* the response rewrite parses via `resp.clone()` only (never consume the
* body of a response we return as-is; the DeepSeek shim rule). Fail-open:
* any rewrite error returns the original request/response.
*
* @internal exported for tests.
*/
// Cast through `unknown` because Bun's `typeof fetch` carries a `preconnect`
// member the arrow function does not implement (matches deepseek.ts).
export const minimaxCompatFetch = (async (
input: RequestInfo | URL,
init?: RequestInit,
): Promise<Response> => {
const url =
typeof input === 'string' ? input : input instanceof URL ? input.href : input.url;
const isEmbeddings = url.includes('/embeddings');
// OUTBOUND (embeddings only): `input` → `texts`, default `type: 'db'`
// (the recipe's documented symmetric default — the AI SDK adapter strips
// the `type` threaded via providerOptions before it reaches the wire,
// same class as #1400), and drop `encoding_format` (not a MiniMax param).
if (isEmbeddings && init?.body && typeof init.body === 'string') {
try {
const parsed = JSON.parse(init.body);
if (
parsed && typeof parsed === 'object' &&
parsed.input !== undefined && parsed.texts === undefined
) {
parsed.texts = Array.isArray(parsed.input) ? parsed.input : [parsed.input];
delete parsed.input;
delete parsed.encoding_format;
if (parsed.type === undefined) parsed.type = 'db';
// Drop Content-Length so fetch recomputes from the new body.
const headers = new Headers(init.headers ?? {});
headers.delete('content-length');
init = { ...init, body: JSON.stringify(parsed), headers };
}
} catch {
// Body wasn't JSON — pass through untouched.
}
}
const res = await fetch(input as any, init as any);
// INBOUND (embeddings only): `{vectors: [[...]]}` → `{data: [{embedding}]}`.
// Anything else (chat completions, MiniMax base_resp errors, non-JSON)
// returns the ORIGINAL response with its body unread.
if (!isEmbeddings || !res.ok) return res;
const ctype = res.headers.get('content-type') ?? '';
if (!ctype.toLowerCase().includes('application/json')) return res;
try {
const json = await res.clone().json();
if (!json || typeof json !== 'object' || !Array.isArray(json.vectors)) return res;
const totalTokens = typeof json.total_tokens === 'number' ? json.total_tokens : 0;
const rewritten = {
object: 'list',
data: (json.vectors as number[][]).map((embedding, index) => ({
object: 'embedding',
embedding,
index,
})),
model: typeof json.model === 'string' ? json.model : 'embo-01',
usage: { prompt_tokens: totalTokens, total_tokens: totalTokens },
};
// Fresh header set: the body changed, so upstream content-length /
// content-encoding would now be wrong.
const headers = new Headers(res.headers);
headers.delete('content-length');
headers.delete('content-encoding');
return new Response(JSON.stringify(rewritten), {
status: res.status,
statusText: res.statusText,
headers,
});
} catch {
return res;
}
}) as unknown as typeof fetch;
/**
* MiniMax (海螺AI). `/embeddings` endpoint at api.minimaxi.com (wire shape
* normalized by `minimaxCompatFetch` above); OpenAI-compatible
* `/chat/completions`. The flagship embedding model is `embo-01` (1536 dims).
* MiniMax (海螺AI). OpenAI-compatible /embeddings endpoint at
* api.minimax.chat. The flagship embedding model is `embo-01` (1536 dims).
*
* MiniMax's API takes an extra `type: 'db' | 'query'` field for asymmetric
* retrieval. gbrain currently has no notion of "this is a document vs a
@@ -130,25 +38,7 @@ export const minimax: Recipe = {
// halving in the gateway catches token-limit errors at runtime.
max_batch_tokens: 4096,
},
chat: {
// Model list from MiniMax's /v1/models (#1977). Chat is genuinely
// OpenAI-compatible — no wire rewrite needed (minimaxCompatFetch
// passes non-embedding requests through untouched).
models: [
'MiniMax-M3',
'MiniMax-M2.7',
'MiniMax-M2.7-highspeed',
'MiniMax-M2.5',
'MiniMax-M2.5-highspeed',
'MiniMax-M2.1',
'MiniMax-M2.1-highspeed',
'MiniMax-M2',
],
supports_tools: false,
supports_subagent_loop: false,
},
},
setup_hint:
'Get an API key at https://www.minimaxi.com, then `export MINIMAX_API_KEY=...`',
compat: { fetch: minimaxCompatFetch },
};
+1 -107
View File
@@ -6,14 +6,10 @@
* - default auth: MINIMAX_API_KEY → "Bearer <key>"; missing → AIConfigError
* - dimsProviderOptions threads `type: 'db'` for embo-01 (the asymmetric
* retrieval field default) — pins the v1 indexing-only behavior
* - #1977: chat touchpoint declared; minimaxCompatFetch rewrites the
* embedding wire shape both directions, passes chat through with the
* response body UNREAD (the consumed-body regression), fail-open.
*/
import { afterEach, describe, expect, test } from 'bun:test';
import { describe, expect, test } from 'bun:test';
import { getRecipe } from '../../src/core/ai/recipes/index.ts';
import { minimaxCompatFetch } from '../../src/core/ai/recipes/minimax.ts';
import { defaultResolveAuth } from '../../src/core/ai/gateway.ts';
import { dimsProviderOptions } from '../../src/core/ai/dims.ts';
import { AIConfigError } from '../../src/core/ai/errors.ts';
@@ -60,106 +56,4 @@ describe('recipe: minimax', () => {
expect(dimsProviderOptions('openai-compatible', 'voyage-3-lite', 512)).toBeUndefined();
expect(dimsProviderOptions('openai-compatible', 'nomic-embed-text', 768)).toBeUndefined();
});
test('chat touchpoint declared (#1977) so assertTouchpoint permits gbrain think', () => {
const r = getRecipe('minimax')!;
expect(r.touchpoints.chat).toBeDefined();
expect(r.touchpoints.chat!.models).toContain('MiniMax-M3');
expect(r.touchpoints.chat!.supports_tools).toBe(false);
expect(r.touchpoints.chat!.supports_subagent_loop).toBe(false);
});
test('recipe ships minimaxCompatFetch via compat.fetch (no env-templated base URL)', () => {
const r = getRecipe('minimax')!;
expect(r.compat?.fetch).toBe(minimaxCompatFetch);
// base_urls config override must keep working: no resolveOpenAICompatConfig.
expect(r.resolveOpenAICompatConfig).toBeUndefined();
});
});
describe('minimaxCompatFetch (#1977)', () => {
const realFetch = globalThis.fetch;
afterEach(() => { globalThis.fetch = realFetch; });
function stubFetch(body: unknown, init?: { status?: number; contentType?: string }) {
const calls: { url: string; init?: RequestInit }[] = [];
globalThis.fetch = (async (input: any, i?: RequestInit) => {
calls.push({ url: String(input), init: i });
return new Response(typeof body === 'string' ? body : JSON.stringify(body), {
status: init?.status ?? 200,
headers: { 'content-type': init?.contentType ?? 'application/json' },
});
}) as unknown as typeof fetch;
return calls;
}
const EMBED_URL = 'https://api.minimaxi.com/v1/embeddings';
const CHAT_URL = 'https://api.minimaxi.com/v1/chat/completions';
test('embedding request: input → texts, type:db injected, encoding_format dropped', async () => {
const calls = stubFetch({ vectors: [[0.1, 0.2]] });
await minimaxCompatFetch(EMBED_URL, {
method: 'POST',
headers: { 'content-type': 'application/json', 'content-length': '99' },
body: JSON.stringify({ model: 'embo-01', input: ['hello', 'world'], encoding_format: 'float' }),
});
const wire = JSON.parse(calls[0]!.init!.body as string);
expect(wire.texts).toEqual(['hello', 'world']);
expect(wire.input).toBeUndefined();
expect(wire.encoding_format).toBeUndefined();
expect(wire.type).toBe('db');
expect(new Headers(calls[0]!.init!.headers).get('content-length')).toBeNull();
});
test('embedding response: {vectors} rewritten to OpenAI {data:[{embedding}]}', async () => {
stubFetch({ vectors: [[0.1, 0.2], [0.3, 0.4]], total_tokens: 7 });
const res = await minimaxCompatFetch(EMBED_URL, {
method: 'POST',
body: JSON.stringify({ model: 'embo-01', input: ['a', 'b'] }),
});
const json = await res.json();
expect(json.data).toEqual([
{ object: 'embedding', embedding: [0.1, 0.2], index: 0 },
{ object: 'embedding', embedding: [0.3, 0.4], index: 1 },
]);
expect(json.usage).toEqual({ prompt_tokens: 7, total_tokens: 7 });
});
test('chat completion passes through with body UNREAD (consumed-body regression)', async () => {
stubFetch({ choices: [{ message: { role: 'assistant', content: 'hi' } }] });
const res = await minimaxCompatFetch(CHAT_URL, {
method: 'POST',
body: JSON.stringify({ model: 'MiniMax-M3', messages: [{ role: 'user', content: 'say hi' }] }),
});
expect(res.bodyUsed).toBe(false); // the broken PR #2882 wrapper consumed this
const json = await res.json(); // must NOT throw "Body already used"
expect(json.choices[0].message.content).toBe('hi');
});
test('chat request body is never rewritten (messages untouched, no type injected)', async () => {
const calls = stubFetch({ choices: [] });
const body = JSON.stringify({ model: 'MiniMax-M3', messages: [{ role: 'user', content: 'x' }] });
await minimaxCompatFetch(CHAT_URL, { method: 'POST', body });
expect(calls[0]!.init!.body).toBe(body);
});
test('embedding error response ({vectors:null, base_resp}) passes through re-readable', async () => {
stubFetch({ vectors: null, base_resp: { status_code: 2013, status_msg: 'invalid params' } });
const res = await minimaxCompatFetch(EMBED_URL, {
method: 'POST',
body: JSON.stringify({ model: 'embo-01', input: ['a'] }),
});
expect(res.bodyUsed).toBe(false);
const json = await res.json();
expect(json.base_resp.status_code).toBe(2013);
});
test('fail-open: non-JSON response body passes through untouched', async () => {
stubFetch('not json', { contentType: 'application/json' });
const res = await minimaxCompatFetch(EMBED_URL, {
method: 'POST',
body: JSON.stringify({ model: 'embo-01', input: ['a'] }),
});
expect(await res.text()).toBe('not json');
});
});
+45
View File
@@ -60,6 +60,51 @@ async function seedPage(slug: string, body: string): Promise<void> {
}
describe('runExtractCore — incremental cycle path (#417)', () => {
test('Dream incremental all-mode stamps the source-scoped extraction watermark (#2636)', async () => {
await engine.executeRaw(
`INSERT INTO sources (id, name, local_path) VALUES ($1, $2, $3)`,
['repo-a', 'repo-a', tempDir],
);
await engine.putPage('people/alice-example', {
type: 'person',
title: 'alice-example',
compiled_truth: '# alice',
timeline: '',
frontmatter: {},
content_hash: 'h',
}, { sourceId: 'repo-a' });
writeFileSync(join(tempDir, 'people/alice-example.md'), '# alice');
await runExtractCore(engine as unknown as BrainEngine, {
mode: 'all',
dir: tempDir,
slugs: ['people/alice-example'],
sourceId: 'repo-a',
});
const rows = await engine.executeRaw<{ links_extracted_at: string | null }>(
`SELECT links_extracted_at FROM pages WHERE slug = $1 AND source_id = $2`,
['people/alice-example', 'repo-a'],
);
expect(rows[0]?.links_extracted_at).not.toBeNull();
expect(await engine.countStalePagesForExtraction({ sourceId: 'repo-a' })).toBe(0);
});
test('Dream incremental dry-run does NOT stamp the watermark', async () => {
await seedPage('people/alice-example', '# alice');
await runExtractCore(engine as unknown as BrainEngine, {
mode: 'all',
dir: tempDir,
slugs: ['people/alice-example'],
dryRun: true,
});
const rows = await engine.executeRaw<{ links_extracted_at: string | null }>(
`SELECT links_extracted_at FROM pages WHERE slug = $1`,
['people/alice-example'],
);
expect(rows[0]?.links_extracted_at ?? null).toBeNull();
});
test('1. slugs: [] returns immediately with zero counts (early-return path)', async () => {
await seedPage('people/alice-example', '# alice');
const result = await runExtractCore(engine as unknown as BrainEngine, {