feat(km-wiki-ingest): 重做成跑在 cypher 上的 workflow(走 cypher binding,汰換 service-binding drainer,D28)
把知識庫 ingest 從 standalone drainer(誤用 Service Bindings)重表達為 cypher workflow:
- workflow.yaml(Phase 0 cron drain):watch_cron→load_cursor→list_cards→pick_card(code)
→fetch_card→parse_card(code)→upsert_entry→save_cursor→post_envelopes→post_one_envelope。
線性 pipe,跨-worker 全走 cypher binding(零件節點),零 service binding。
- workflow.delta.yaml(Gitea webhook 穩態):collect_changed(code)→foreach card→fetch/parse/upsert/foreach envelope。
- code/kbdb/graph 接法:code=canonical `code` 零件(arcrun-code);kbdb/graph=http_request 零件打
/entries/ingest、/triplets/ingest(server 端冪等);Gitea=http_request。
- 平台端最小補丁(各需 gated 部署):
1) cypher-executor component-loader:WASM_HTTP_RUNNER_IDS 加 'code'(canonical→arcrun-code,cypher binding 正解)。
2) kbdb base:POST /entries/ingest(page_name+content_hash 冪等 upsert,對稱 graph /triplets/ingest;
因 flow DSL 無資料條件分支,把 create/patch/skip 冪等推到 server 端)。
- drainer/DEPRECATED.md:標舊 standalone worker 退役計畫(新版穩定後 wrangler delete)。
- DEPLOY.md:部署順序、cron/webhook 掛法、subdomain 對齊、驗收與退役。
本輪不部署(待總管/leo 審架構)。
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HJiLCRUU2o3aSpPEzVCt2o
This commit is contained in:
@@ -23,6 +23,68 @@ entryRoutes.post('/', async (c) => {
|
||||
return c.json({ success: true, entry });
|
||||
});
|
||||
|
||||
// POST /entries/ingest — server 端冪等 upsert(page_name 當鍵 + content_hash skip-if-unchanged)。
|
||||
// 語義對稱 graph 的 POST /triplets/ingest:讓「工作流只打一發、冪等由 server 保證」,
|
||||
// 使 km_wiki_ingest_drain cypher workflow 不必在 flow DSL 裡做 lookup→decide→create/patch 分支
|
||||
// (DSL 無資料條件分支)。owner_id 走 query。body 收 parse_card 產的 entry 形狀:
|
||||
// { page_name, entry_type, content, source, tags?(array|json), metadata?(object)|metadata_json?, content_hash? }
|
||||
// 行為:無此 (page_name,owner_id) → create;有且 content_hash 同 → skip(不重嵌);有且不同/無 hash → update+重嵌。
|
||||
// 回:{ success, action:'created'|'updated'|'skipped', entry }。
|
||||
entryRoutes.post('/ingest', async (c) => {
|
||||
const body = await c.req.json().catch(() => null) as Record<string, unknown> | null;
|
||||
if (!body || !body.entry_type || !body.page_name) {
|
||||
return c.json({ success: false, error: 'entry_type 與 page_name 必填' }, 400);
|
||||
}
|
||||
const owner_id = c.req.query('owner_id') || (body.owner_id as string | undefined) || undefined;
|
||||
const page_name = String(body.page_name);
|
||||
|
||||
// 正規化 tags / metadata → *_json 字串(同時容忍呼叫端直接給 *_json)。
|
||||
const tags_json =
|
||||
typeof body.tags_json === 'string' ? body.tags_json
|
||||
: Array.isArray(body.tags) ? JSON.stringify(body.tags)
|
||||
: undefined;
|
||||
const metaObj = (body.metadata && typeof body.metadata === 'object') ? body.metadata as Record<string, unknown> : undefined;
|
||||
const metadata_json =
|
||||
typeof body.metadata_json === 'string' ? body.metadata_json
|
||||
: metaObj ? JSON.stringify(metaObj)
|
||||
: undefined;
|
||||
const newHash =
|
||||
(typeof body.content_hash === 'string' && body.content_hash)
|
||||
|| (metaObj && typeof metaObj.content_hash === 'string' ? metaObj.content_hash : undefined);
|
||||
|
||||
// 查現有(page_name 是冪等鍵;owner_id 隔離租戶)。
|
||||
const { entries } = await listEntries(c.env.DB, { page_name, owner_id, limit: 1 });
|
||||
const existing = entries[0];
|
||||
|
||||
if (existing) {
|
||||
let storedHash: string | undefined;
|
||||
try { storedHash = existing.metadata_json ? (JSON.parse(existing.metadata_json).content_hash as string) : undefined; } catch { /* ignore */ }
|
||||
if (newHash && storedHash && storedHash === newHash) {
|
||||
return c.json({ success: true, action: 'skipped', entry: existing });
|
||||
}
|
||||
const entry = await updateEntry(c.env.DB, existing.id, {
|
||||
content: body.content as string | undefined,
|
||||
...(tags_json !== undefined ? { tags_json } : {}),
|
||||
...(metadata_json !== undefined ? { metadata_json } : {}),
|
||||
});
|
||||
if (embedEnabled(c.env) && body.content !== undefined && entry) {
|
||||
c.executionCtx.waitUntil(embedOnWrite(c.env, entry).catch(() => {}));
|
||||
}
|
||||
return c.json({ success: true, action: 'updated', entry });
|
||||
}
|
||||
|
||||
const entry = await createEntry(c.env.DB, {
|
||||
entry_type: String(body.entry_type),
|
||||
content: body.content as string | undefined,
|
||||
owner_id,
|
||||
page_name,
|
||||
tags_json,
|
||||
metadata_json,
|
||||
});
|
||||
if (embedEnabled(c.env)) c.executionCtx.waitUntil(embedOnWrite(c.env, entry).catch(() => {}));
|
||||
return c.json({ success: true, action: 'created', entry });
|
||||
});
|
||||
|
||||
// GET /entries — list with filters (entry_type, owner_id, parent_id, page_name, source, q/search)
|
||||
// e.g. list workflows under a project: ?parent_id=PROJECT&entry_type=workflow
|
||||
// e.g. get one by idempotency key: ?page_name=skill-rag_with_arcrun
|
||||
|
||||
Reference in New Issue
Block a user