ea1c0571c1
- registry 新增 POST /analytics/record:ANALYTICS_KV 計數器(stats:{hash_id}:{version})
為唯一真相源,衍生值(success_rate/avg_duration_ms/call_count)回填 comp: 記錄
——查詢讀取端讀哪就寫哪,不開第二真相源。KV 無 CAS,誠實標註非原子。
- cypher-executor execution-evaluator 從 stub 改真實作:執行收尾(/cypher/execute
成功與 ExecutionError 路徑+webhook 路徑)對 trace 裡每顆 Component 節點
fire-and-forget 回寫,waitUntil 包、不增加執行同步延遲(仿 recordRecipeStats 慣例)。
成敗判定=trace error 或 output.success===false(makeHttpRunner 非 2xx 不 throw)。
- registry 位置沿 search-nodes 慣例:REGISTRY_BASE_URL 覆蓋,未設走 wasmWorkerUrl。
- 新增單測 13 個全綠;本地雙 wrangler dev 端到端實測:http_request 跑 5 次
(3 成功+2 失敗)→ success_rate 1→0.6、call_count 0→5,/cypher/search 同步可見。
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
161 lines
5.9 KiB
TypeScript
161 lines
5.9 KiB
TypeScript
import type { Bindings, ExecutionGraph, ExecutionContext } from '../types';
|
||
import { ExecutionError } from '../types';
|
||
import { GraphExecutor } from '../graph-executor';
|
||
import { graphSchema } from '../lib/schemas';
|
||
import { createComponentLoader } from '../lib/component-loader';
|
||
import { recordTelemetry } from '../lib/telemetry';
|
||
import { recordComponentStats } from './execution-evaluator';
|
||
import type { GraphNode, TraceStep } from '../types';
|
||
|
||
/**
|
||
* kbdb-base §7.1+§7.5.h:一條工作流執行結束後,把這次用到的 recipe 各記一次成功/失敗到 KBDB 市場星數。
|
||
* 判定單位是「工作流執行」(n8n execution):整體成功 → 用到的每個 recipe key +1 成功;整體失敗 → 各 +1 失敗。
|
||
* **key = recipe uuid**(per-uuid,能區分同 canonical 的不同作者版本 §7.5.5;舊資料 fallback canonical_id)。
|
||
*
|
||
* fire-and-forget(用 ctx.waitUntil,仿 recordTelemetry):記錄失敗不影響工作流結果。
|
||
* KBDB 端點 POST {KBDB_BASE_URL}/recipe-stats/record { canonical_id, ok, at }——
|
||
* 該欄位名為 canonical_id 但語意已是 recipe key(uuid),KBDB 端只當 stat 的主鍵字串用。
|
||
*/
|
||
function recordRecipeStats(
|
||
env: Bindings,
|
||
recipeKeys: Set<string>,
|
||
ok: boolean,
|
||
at: number,
|
||
ctx?: ExecutionContext,
|
||
): void {
|
||
if (recipeKeys.size === 0) return;
|
||
const base = (env.KBDB_BASE_URL ?? 'https://kbdb.finally.click').replace(/\/$/, '');
|
||
const headers: Record<string, string> = { 'Content-Type': 'application/json' };
|
||
if (env.KBDB_INTERNAL_TOKEN) headers['Authorization'] = `Bearer ${env.KBDB_INTERNAL_TOKEN}`;
|
||
|
||
const promise = Promise.all(
|
||
[...recipeKeys].map(key =>
|
||
fetch(`${base}/recipe-stats/record`, {
|
||
method: 'POST',
|
||
headers,
|
||
body: JSON.stringify({ canonical_id: key, ok, at }),
|
||
}).catch(() => undefined),
|
||
),
|
||
).then(() => undefined);
|
||
|
||
if (ctx?.waitUntil) ctx.waitUntil(promise);
|
||
else void promise;
|
||
}
|
||
|
||
type WebhookRecord = {
|
||
graph: Record<string, unknown>;
|
||
description: string;
|
||
created_at: string;
|
||
};
|
||
|
||
export function generateToken(): string {
|
||
const tokenBytes = crypto.getRandomValues(new Uint8Array(16));
|
||
return Array.from(tokenBytes).map(b => b.toString(16).padStart(2, '0')).join('');
|
||
}
|
||
|
||
export async function validateAndParseWebhook(raw: string): Promise<WebhookRecord | null> {
|
||
try {
|
||
return JSON.parse(raw) as WebhookRecord;
|
||
} catch {
|
||
return null;
|
||
}
|
||
}
|
||
|
||
export async function executeWebhookGraph(
|
||
env: Bindings,
|
||
graph: Record<string, unknown>,
|
||
triggerContext: Record<string, unknown>,
|
||
token: string,
|
||
apiKey?: string,
|
||
ctx?: ExecutionContext, // 可選 — 用 waitUntil 把 telemetry 推到背景
|
||
userAgent?: string, // MCP / SDK client 帶過來
|
||
): Promise<{ success: boolean; data?: unknown; error?: string; trace?: unknown; duration_ms: number }> {
|
||
const parsed = graphSchema.safeParse(graph);
|
||
if (!parsed.success) {
|
||
return { success: false, error: '圖定義已失效', duration_ms: 0 };
|
||
}
|
||
|
||
const loader = createComponentLoader(env);
|
||
const executor = new GraphExecutor(loader, undefined, env, apiKey);
|
||
const start = Date.now();
|
||
|
||
try {
|
||
const result = await executor.execute(
|
||
parsed.data as ExecutionGraph,
|
||
{ ...triggerContext, _webhook_token: token },
|
||
env.EXEC_CONTEXT,
|
||
);
|
||
const duration_ms = Date.now() - start;
|
||
|
||
// Implicit telemetry:成功 run(含 paused 也算「成功啟動」由 trigger_workflow 那層分類)
|
||
recordTelemetry(env, apiKey, {
|
||
event_type: 'run_success',
|
||
workflow_name: token,
|
||
duration_ms,
|
||
agent_user_agent: userAgent,
|
||
}, ctx);
|
||
|
||
// kbdb-base §7.1:整體成功 → 用到的 recipe 各記成功一次。
|
||
recordRecipeStats(env, executor.usedRecipeKeys, true, Date.now(), ctx);
|
||
|
||
// arcrun-core-mvp「執行統計設計」:對用到的每顆零件回寫執行結果(fire-and-forget)。
|
||
{
|
||
const statsPromise = recordComponentStats(
|
||
env,
|
||
(parsed.data as ExecutionGraph).nodes as GraphNode[],
|
||
result.trace as TraceStep[],
|
||
);
|
||
if (ctx?.waitUntil) ctx.waitUntil(statsPromise);
|
||
else void statsPromise;
|
||
}
|
||
|
||
return { success: true, data: result.data, duration_ms };
|
||
} catch (err) {
|
||
const duration_ms = Date.now() - start;
|
||
const errMsg = err instanceof Error ? err.message : String(err);
|
||
const isPaused = /workflow paused/i.test(errMsg);
|
||
|
||
// Implicit telemetry:paused 算 run_success;真錯才 run_fail
|
||
recordTelemetry(env, apiKey, {
|
||
event_type: isPaused ? 'run_success' : 'run_fail',
|
||
workflow_name: token,
|
||
error_code: isPaused ? 'paused_awaiting_resume' : 'execution_error',
|
||
duration_ms,
|
||
agent_user_agent: userAgent,
|
||
}, ctx);
|
||
|
||
// kbdb-base §7.1:真錯(非 paused)→ 用到的 recipe 各記失敗一次。
|
||
// paused 是「執行中暫停等 callback」非失敗,不記(resume 後成功才會在那條路徑記成功)。
|
||
if (!isPaused) {
|
||
recordRecipeStats(env, executor.usedRecipeKeys, false, Date.now(), ctx);
|
||
}
|
||
|
||
// 零件統計失敗路徑:ExecutionError 帶完整 trace(失敗節點有 error、先前成功節點照記成功);
|
||
// paused 非失敗不記;非 ExecutionError 無 trace 可歸因 → 不記。
|
||
if (!isPaused && err instanceof ExecutionError) {
|
||
const statsPromise = recordComponentStats(
|
||
env,
|
||
(parsed.data as ExecutionGraph).nodes as GraphNode[],
|
||
err.trace,
|
||
);
|
||
if (ctx?.waitUntil) ctx.waitUntil(statsPromise);
|
||
else void statsPromise;
|
||
}
|
||
|
||
if (err instanceof ExecutionError) {
|
||
const traceFormatted = err.trace.map(s => ({
|
||
node: s.nodeId,
|
||
status: s.error ? 'failed' : 'success',
|
||
...(s.error ? { error: s.error } : {}),
|
||
}));
|
||
return {
|
||
success: false,
|
||
error: errMsg,
|
||
trace: traceFormatted,
|
||
duration_ms,
|
||
};
|
||
}
|
||
return { success: false, error: errMsg, duration_ms };
|
||
}
|
||
}
|