Skip to content

Commit 0a3dc7c

Browse files
author
moc
committed
fix: strip ledger contamination, bind upstream output into lineage, add MCP contract test
- store.mjs rebuilt from main: findCrossRunReuse only, no integrity-ledger code (the ledger belongs to MiniMax-AI#48; the previous round accidentally carried it) - lineageHash now includes each succeeded dependency's output hash, so an upstream that re-executes with different output (mcode node without an explicit model, tracked-file change during execution) invalidates downstream adoption — regression covers the maintainer's divergence shape - checks/cross-reuse-mcp.check.mjs: packaged MCP advertises reuseAcrossRuns on workflow_start/workflow_update and accepts/rejects it through the public tool surface (additionalProperties:false contract)
1 parent 1b1b1a6 commit 0a3dc7c

5 files changed

Lines changed: 50 additions & 90 deletions

File tree

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
import test from 'node:test';import assert from 'node:assert/strict';import {Client} from '@modelcontextprotocol/sdk/client/index.js';import {StdioClientTransport} from '@modelcontextprotocol/sdk/client/stdio.js';import {mkdtemp,rm,writeFile} from 'node:fs/promises';import {tmpdir} from 'node:os';import {join,resolve} from 'node:path';
2+
test('packaged MCP advertises reuseAcrossRuns and accepts it through the public tool surface',async()=>{
3+
const dir=await mkdtemp(join(tmpdir(),'wf-cross-mcp-'));await writeFile(join(dir,'settings.json'),JSON.stringify({workspace:dir,dataDir:dir}));
4+
const client=new Client({name:'cross-reuse-mcp-test',version:'1'});const transport=new StdioClientTransport({command:process.execPath,args:[resolve('dist/main.mjs'),'--stdio','--settings',join(dir,'settings.json')],stderr:'pipe'});
5+
try{
6+
await client.connect(transport);const {tools}=await client.listTools();
7+
for(const name of ['workflow_start','workflow_update']){const tool=tools.find(t=>t.name===name);assert.ok(tool,`${name} listed`);assert.equal(tool.inputSchema.properties.reuseAcrossRuns?.type,'boolean',`${name} schema must advertise reuseAcrossRuns (additionalProperties:false)`);}
8+
const started=await client.callTool({name:'workflow_start',arguments:{requestId:'cross-mcp',name:'Cross-run via public surface',executor:'demo',reuseAcrossRuns:true,script:'return await ctx.agent({id:"a",prompt:"p"});'}});
9+
assert.ok(!started.isError,started.content?.[0]?.text);const run=JSON.parse(started.content[0].text);
10+
assert.equal(run.reuseAcrossRuns,true,'the flag must survive the public tool surface');
11+
const rejected=await client.callTool({name:'workflow_start',arguments:{requestId:'cross-mcp-bad',name:'Bad flag type',executor:'demo',reuseAcrossRuns:'yes',script:'return 1;'}});
12+
assert.ok(rejected.isError,'a non-boolean flag must be rejected by the public surface');
13+
}finally{await client.close();await transport.close();await rm(dir,{recursive:true,force:true});}
14+
});

‎plugins/hetaoBackend/mcode-dynamic-workflows/checks/cross-reuse.check.mjs‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -153,3 +153,17 @@ test('chained adoption keeps reusedFrom on the immediate source and originalProd
153153
assert.equal(b3.usage,null);assert.deepEqual(b3.usageHistory,[]);
154154
}finally{await f.cleanup();}
155155
});
156+
test('an upstream that re-executes with a different output invalidates downstream adoption',async()=>{
157+
// Maintainer's divergence shape: the upstream is not a cross-run candidate itself
158+
// (mcode node without an explicit model re-executes every run) and its output
159+
// differs between runs; the eligible downstream must not adopt the stale result.
160+
let aOutput='old';const calls=[],f=await fixture(async s=>{calls.push(s.id);return {output:s.id==='a'?aOutput:s.id};});try{
161+
const script=`const a=await ctx.agent({id:'a',prompt:'write'});const b=await ctx.agent({id:'b',prompt:'read',model:'m2',dependsOn:['a'],input:{content:a.output}});return {a:a.output,b:b.output};`;
162+
const run1=await run(f.engine,script,{},{executor:'mcode'});
163+
aOutput='new';
164+
const run2=await run(f.engine,script,{},{executor:'mcode',reuseAcrossRuns:true});
165+
assert.deepEqual(calls,['a','b','a','b'],'a re-executes (no model => never a candidate); b must not adopt the stale run-1 result');
166+
assert.deepEqual(run2.result,{a:'new',b:'b'});
167+
assert.ok(!run2.steps.find(s=>s.id==='b').reusedFrom,'divergent upstream output must break lineage');
168+
}finally{await f.cleanup();}
169+
});

‎plugins/hetaoBackend/mcode-dynamic-workflows/dist/main.mjs‎

Lines changed: 12 additions & 63 deletions
Original file line numberDiff line numberDiff line change
@@ -7714,7 +7714,7 @@ import { spawn as spawn3 } from "node:child_process";
77147714
import { DatabaseSync } from "node:sqlite";
77157715
import { mkdirSync, openSync, writeFileSync, closeSync, readFileSync, unlinkSync } from "node:fs";
77167716
import { join } from "node:path";
7717-
import { createHash, randomUUID } from "node:crypto";
7717+
import { randomUUID } from "node:crypto";
77187718
var Store = class {
77197719
constructor(dir) {
77207720
mkdirSync(dir, { recursive: true, mode: 448 });
@@ -7740,7 +7740,6 @@ var Store = class {
77407740
this.fd = openSync(this.lock, "wx", 384);
77417741
}
77427742
this.owner = randomUUID();
7743-
this.txDepth = 0;
77447743
try {
77457744
writeFileSync(this.fd, JSON.stringify({ pid: process.pid, owner: this.owner }));
77467745
this.db = new DatabaseSync(join(dir, "workflows.sqlite"));
@@ -7752,8 +7751,7 @@ var Store = class {
77527751
CREATE TABLE IF NOT EXISTS steps(runId TEXT,id TEXT,body TEXT NOT NULL,PRIMARY KEY(runId,id));
77537752
CREATE TABLE IF NOT EXISTS repair_cache(runId TEXT,id TEXT,body TEXT NOT NULL,PRIMARY KEY(runId,id));
77547753
CREATE TABLE IF NOT EXISTS events(seq INTEGER PRIMARY KEY AUTOINCREMENT,runId TEXT,body TEXT NOT NULL);
7755-
CREATE INDEX IF NOT EXISTS run_events ON events(runId,seq);
7756-
CREATE TABLE IF NOT EXISTS integrity_rows(surface TEXT NOT NULL,pos INTEGER NOT NULL,key TEXT NOT NULL,hash TEXT NOT NULL,PRIMARY KEY(surface,pos));`);
7754+
CREATE INDEX IF NOT EXISTS run_events ON events(runId,seq);`);
77577755
const unfinished = this.db.prepare("SELECT body FROM runs WHERE json_extract(body,'$.status') IN ('running','queued','stopping','pausing')").all();
77587756
for (const row of unfinished) {
77597757
const run = JSON.parse(row.body);
@@ -7768,8 +7766,6 @@ var Store = class {
77687766
}
77697767
}
77707768
transaction(fn) {
7771-
if (this.txDepth) return fn();
7772-
this.txDepth = 1;
77737769
this.db.exec("BEGIN IMMEDIATE");
77747770
try {
77757771
const r = fn();
@@ -7778,8 +7774,6 @@ var Store = class {
77787774
} catch (e) {
77797775
this.db.exec("ROLLBACK");
77807776
throw e;
7781-
} finally {
7782-
this.txDepth = 0;
77837777
}
77847778
}
77857779
templates() {
@@ -7841,61 +7835,16 @@ var Store = class {
78417835
return r ? JSON.parse(r.body) : null;
78427836
}
78437837
saveRepairCandidate(runId, step) {
7844-
this.transaction(() => {
7845-
const rowid = Number(this.db.prepare("INSERT INTO repair_cache VALUES(?,?,?)").run(runId, step.id, JSON.stringify(step)).lastInsertRowid);
7846-
this.chainAdvance("repair", "repair", "SELECT rowid AS pos,runId,id,body FROM repair_cache WHERE rowid>? AND rowid<=? ORDER BY rowid", rowid, (r) => `${r.runId}/${r.id}`);
7847-
});
7838+
this.db.prepare("INSERT INTO repair_cache VALUES(?,?,?)").run(runId, step.id, JSON.stringify(step));
78487839
}
78497840
event(runId, type, data2 = {}) {
78507841
const event = { ...data2, type, time: Date.now() };
7851-
return this.transaction(() => {
7852-
const seq = Number(this.db.prepare("INSERT INTO events(runId,body) VALUES(?,?)").run(runId, JSON.stringify(event)).lastInsertRowid);
7853-
this.chainAdvance("event", "events", "SELECT seq AS pos,body FROM events WHERE seq>? AND seq<=? ORDER BY seq", seq, (r) => String(r.pos));
7854-
return { seq, ...event };
7855-
});
7842+
const seq = Number(this.db.prepare("INSERT INTO events(runId,body) VALUES(?,?)").run(runId, JSON.stringify(event)).lastInsertRowid);
7843+
return { seq, ...event };
78567844
}
78577845
events(runId, after = 0, limit = 150) {
78587846
return this.db.prepare("SELECT seq,body FROM events WHERE runId=? AND seq>? ORDER BY seq LIMIT ?").all(runId, after, limit).map((e) => ({ seq: e.seq, ...JSON.parse(e.body) }));
78597847
}
7860-
rowHash(prev, kind, key, body) {
7861-
return createHash("sha256").update(`${prev}:${kind}:${key}:${body}`).digest("hex");
7862-
}
7863-
chainAdvance(kind, surface, sql, newUpto, keyOf) {
7864-
const tail = this.setting(`integrity_${surface}`);
7865-
let prev = tail?.head ?? "0".repeat(64);
7866-
for (const r of this.db.prepare(sql).all(tail?.upto ?? 0, newUpto)) {
7867-
const k = keyOf(r);
7868-
prev = this.rowHash(prev, kind, k, r.body);
7869-
this.db.prepare("INSERT OR REPLACE INTO integrity_rows VALUES(?,?,?,?)").run(surface, r.pos, k, prev);
7870-
}
7871-
this.saveSetting(`integrity_${surface}`, { head: prev, upto: newUpto });
7872-
}
7873-
integrityHeads() {
7874-
return { events: this.setting("integrity_events") ?? null, repair: this.setting("integrity_repair") ?? null };
7875-
}
7876-
verifyIntegrity() {
7877-
const genesis = "0".repeat(64);
7878-
const face = (kind, surface, table, posCol) => {
7879-
const skey = `integrity_${surface}`;
7880-
const rec = this.setting(skey);
7881-
const total = Number(this.db.prepare(`SELECT COUNT(*) AS n FROM ${table}`).get().n);
7882-
if (!rec) return { head: null, upto: 0, verified: null, checked: 0, unchained: total, firstDivergence: null };
7883-
const rows = this.db.prepare("SELECT pos,key,hash FROM integrity_rows WHERE surface=? ORDER BY pos").all(surface);
7884-
let prev = genesis, firstDivergence = null;
7885-
for (const r of rows) {
7886-
const row = this.db.prepare(`SELECT body FROM ${table} WHERE ${posCol}=?`).get(r.pos);
7887-
const actual = row ? this.rowHash(prev, kind, r.key, row.body) : null;
7888-
if (!firstDivergence && (!row || actual !== r.hash)) firstDivergence = { key: r.key, expectedHead: r.hash, actualHead: actual };
7889-
prev = r.hash;
7890-
}
7891-
const verified = !firstDivergence && prev === rec.head;
7892-
return { head: rec.head, upto: rec.upto, verified, checked: rows.length, unchained: Number(this.db.prepare(`SELECT COUNT(*) AS n FROM ${table} WHERE ${posCol}>?`).get(rec.upto).n), firstDivergence };
7893-
};
7894-
return {
7895-
events: face("event", "events", "events", "seq"),
7896-
repair: face("repair", "repair", "repair_cache", "rowid")
7897-
};
7898-
}
78997848
releaseLock() {
79007849
closeSync(this.fd);
79017850
try {
@@ -7910,7 +7859,7 @@ var Store = class {
79107859
};
79117860

79127861
// src/common.mjs
7913-
import { createHash as createHash2 } from "node:crypto";
7862+
import { createHash } from "node:crypto";
79147863

79157864
// node_modules/acorn/dist/acorn.mjs
79167865
var astralIdentifierCodes = [509, 0, 227, 0, 150, 4, 294, 9, 1368, 2, 2, 1, 6, 3, 41, 2, 5, 0, 166, 1, 574, 3, 9, 9, 7, 9, 32, 4, 318, 1, 78, 5, 71, 10, 50, 3, 123, 2, 54, 14, 32, 10, 3, 1, 11, 3, 46, 10, 8, 0, 46, 9, 7, 2, 37, 13, 2, 9, 6, 1, 45, 0, 13, 2, 49, 13, 9, 3, 2, 11, 83, 11, 7, 0, 3, 0, 158, 11, 6, 9, 7, 3, 56, 1, 2, 6, 3, 1, 3, 2, 10, 0, 11, 1, 3, 6, 4, 4, 68, 8, 2, 0, 3, 0, 2, 3, 2, 4, 2, 0, 15, 1, 83, 17, 10, 9, 5, 0, 82, 19, 13, 9, 214, 6, 3, 8, 28, 1, 83, 16, 16, 9, 82, 12, 9, 9, 7, 19, 58, 14, 5, 9, 243, 14, 166, 9, 71, 5, 2, 1, 3, 3, 2, 0, 2, 1, 13, 9, 120, 6, 3, 6, 4, 0, 29, 9, 41, 6, 2, 3, 9, 0, 10, 10, 47, 15, 199, 7, 137, 9, 54, 7, 2, 7, 17, 9, 57, 21, 2, 13, 123, 5, 4, 0, 2, 1, 2, 6, 2, 0, 9, 9, 49, 4, 2, 1, 2, 4, 9, 9, 55, 9, 266, 3, 10, 1, 2, 0, 49, 6, 4, 4, 14, 10, 5350, 0, 7, 14, 11465, 27, 2343, 9, 87, 9, 39, 4, 60, 6, 26, 9, 535, 9, 470, 0, 2, 54, 8, 3, 82, 0, 12, 1, 19628, 1, 4178, 9, 519, 45, 3, 22, 543, 4, 4, 5, 9, 7, 3, 6, 31, 3, 149, 2, 1418, 49, 513, 54, 5, 49, 9, 0, 15, 0, 23, 4, 2, 14, 1361, 6, 2, 16, 3, 6, 2, 1, 2, 4, 101, 0, 161, 6, 10, 9, 357, 0, 62, 13, 499, 13, 245, 1, 2, 9, 233, 0, 3, 0, 8, 1, 6, 0, 475, 6, 110, 6, 6, 9, 4759, 9, 787719, 239];
@@ -13610,7 +13559,7 @@ function parse3(input, options) {
1361013559
}
1361113560

1361213561
// src/common.mjs
13613-
var hash = (value) => createHash2("sha256").update(typeof value === "string" || Buffer.isBuffer(value) ? value : stable(value)).digest("hex");
13562+
var hash = (value) => createHash("sha256").update(typeof value === "string" || Buffer.isBuffer(value) ? value : stable(value)).digest("hex");
1361413563
function stable(value) {
1361513564
return JSON.stringify(canonical(value));
1361613565
}
@@ -14641,7 +14590,7 @@ var Engine = class extends EventEmitter {
1464114590
const ctxHash = ctx.contextHash ?? (ctx.contextHash = hash({ workspace: ctx.run.workspace, input: ctx.run.input, executor: ctx.run.executor, fingerprints: ctx.run.fingerprints }));
1464214591
const lineageHash = hash({ requestHash, deps: deps.map((id2) => {
1464314592
const dep = this.store.step(ctx.run.id, id2);
14644-
return { id: id2, lineageHash: dep?.lineageHash ?? null };
14593+
return { id: id2, lineageHash: dep?.lineageHash ?? null, outputHash: dep?.status === "succeeded" ? hash(dep.output ?? null) : null };
1464514594
}) });
1464614595
const repair = ctx.run.repair, candidate = !previous && repair?.reuseStepIds.includes(spec.id) ? this.store.repairCandidate(ctx.run.id, spec.id) : null;
1464714596
if (candidate && candidate.requestHash === requestHash && repair.contextHash === hash({ workspace: ctx.run.workspace, input: ctx.run.input, executor: ctx.run.executor, fingerprints: ctx.run.fingerprints }) && deps.every((id2) => {
@@ -14821,7 +14770,7 @@ var Engine = class extends EventEmitter {
1482114770
};
1482214771

1482314772
// src/http.mjs
14824-
import { createHash as createHash3 } from "node:crypto";
14773+
import { createHash as createHash2 } from "node:crypto";
1482514774

1482614775
// web/graph-model.mjs
1482714776
var finished = /* @__PURE__ */ new Set(["succeeded", "completed_with_gaps", "failed", "cancelled"]);
@@ -26615,7 +26564,7 @@ async function startStdio(handler, tools = TOOLS) {
2661526564

2661626565
// src/http.mjs
2661726566
async function startHTTP(engine, { port = 0, webRoot = new URL("../web/", import.meta.url), exampleRoot = new URL("../examples/", import.meta.url) } = {}) {
26618-
const reportStyleHash = createHash3("sha256").update(REPORT_STYLES).digest("base64");
26567+
const reportStyleHash = createHash2("sha256").update(REPORT_STYLES).digest("base64");
2661926568
let origin;
2662026569
const sockets = /* @__PURE__ */ new Set();
2662126570
const server = http.createServer(async (req, res) => {
@@ -26708,7 +26657,7 @@ async function startHTTP(engine, { port = 0, webRoot = new URL("../web/", import
2670826657
}
2670926658

2671026659
// src/workspace-router.mjs
26711-
import { createHash as createHash4 } from "node:crypto";
26660+
import { createHash as createHash3 } from "node:crypto";
2671226661
import { realpath as realpath2, stat as stat2 } from "node:fs/promises";
2671326662
import { isAbsolute as isAbsolute2, relative as relative2, join as join3, sep as sep2 } from "node:path";
2671426663

@@ -27571,7 +27520,7 @@ async function canonicalWorkspace(value, pluginRoot) {
2757127520
return workspace;
2757227521
}
2757327522
function projectDataDir(base, workspace) {
27574-
return join3(base, "projects", createHash4("sha256").update(workspace).digest("hex"));
27523+
return join3(base, "projects", createHash3("sha256").update(workspace).digest("hex"));
2757527524
}
2757627525
function createWorkspaceRouter({ binary, pluginRoot, dataRoot, extraArgs = [] }) {
2757727526
const connections = /* @__PURE__ */ new Map();

‎plugins/hetaoBackend/mcode-dynamic-workflows/src/engine.mjs‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -187,11 +187,12 @@ const step={id:key,kind:'checkpoint',status:'succeeded',output:payload.value,req
187187
if(cached){check(cached.hash===requestHash,'重复 step id 参数冲突');return cached.promise;}
188188
// Computed once per dispatch attempt, before any candidate lookup. contextHash is
189189
// stamped on every new step so the store can filter cross-run candidates in SQL
190-
// before LIMIT; lineageHash pins the node to its own spec plus the lineage hashes
191-
// of its succeeded dependencies (in dependsOn order), so a changed upstream
190+
// before LIMIT; lineageHash pins the node to its own spec plus, for each
191+
// succeeded dependency (in dependsOn order), that dependency's lineage hash
192+
// AND output hash — so a changed, rerun, or differently-outcomed upstream
192193
// invalidates downstream candidates even when the downstream spec is unchanged.
193194
const ctxHash=ctx.contextHash??(ctx.contextHash=hash({workspace:ctx.run.workspace,input:ctx.run.input,executor:ctx.run.executor,fingerprints:ctx.run.fingerprints}));
194-
const lineageHash=hash({requestHash,deps:deps.map(id=>{const dep=this.store.step(ctx.run.id,id);return {id,lineageHash:dep?.lineageHash??null};})});
195+
const lineageHash=hash({requestHash,deps:deps.map(id=>{const dep=this.store.step(ctx.run.id,id);return {id,lineageHash:dep?.lineageHash??null,outputHash:dep?.status==='succeeded'?hash(dep.output??null):null};})});
195196
const repair=ctx.run.repair,candidate=!previous&&repair?.reuseStepIds.includes(spec.id)?this.store.repairCandidate(ctx.run.id,spec.id):null;
196197
if(candidate&&candidate.requestHash===requestHash
197198
&&repair.contextHash===hash({workspace:ctx.run.workspace,input:ctx.run.input,executor:ctx.run.executor,fingerprints:ctx.run.fingerprints})

0 commit comments

Comments
 (0)