diff --git a/packages/core/schema.json b/packages/core/schema.json index a330c6136c..7df86d7742 100644 --- a/packages/core/schema.json +++ b/packages/core/schema.json @@ -1,9 +1,9 @@ { "version": "7", "dialect": "sqlite", - "id": "4142b961-0712-4834-b475-16ea4a74c43c", + "id": "874d8e74-d354-4dcb-b98c-c893660c9371", "prevIds": [ - "7e8e00e9-7bbb-443e-996b-f646ec030c2b" + "4142b961-0712-4834-b475-16ea4a74c43c" ], "ddl": [ { @@ -682,6 +682,16 @@ "entityType": "columns", "table": "workflow_node" }, + { + "type": "integer", + "notNull": true, + "autoincrement": false, + "default": "false", + "generated": null, + "name": "superseded", + "entityType": "columns", + "table": "workflow_node" + }, { "type": "integer", "notNull": true, @@ -822,6 +832,16 @@ "entityType": "columns", "table": "workflow" }, + { + "type": "integer", + "notNull": true, + "autoincrement": false, + "default": "1", + "generated": null, + "name": "graph_rev", + "entityType": "columns", + "table": "workflow" + }, { "type": "integer", "notNull": false, diff --git a/packages/core/src/dag/projector.ts b/packages/core/src/dag/projector.ts index f255e43311..ced752825a 100644 --- a/packages/core/src/dag/projector.ts +++ b/packages/core/src/dag/projector.ts @@ -141,33 +141,62 @@ export const layer = Layer.effectDiscard( yield* events.project(DagEvent.WorkflowReplanned, (event) => Effect.gen(function* () { - // Atomic wake can reach the parent only after a leaf checkpoint has - // completed the current graph. An additive extend emits this event to - // reopen that completed workflow without changing completed nodes. + // Rev-view (v1.0.15 Train A): every replan opens a new graph + // revision. The two legs run in THIS order so each event bumps + // graph_rev exactly once: the seq-bump leg matches the active + // statuses first (status unchanged), and the reopen leg then matches + // completed rows still untouched by it — reversing the legs would let + // the reopened (now running) row match the seq-bump leg too and + // double-bump. yield* db .update(WorkflowTable) .set({ - status: "running", - wake_reported: false, - completed_at: null, + graph_rev: sql`${WorkflowTable.graph_rev} + 1`, seq: event.durable!.seq, time_updated: toMillis(event.data.timestamp), }) .where(and( eq(WorkflowTable.id, event.data.dagID), - inArray(WorkflowTable.status, [...WorkflowStatusProjection.replanReopen.from]), + inArray(WorkflowTable.status, ["pending", "running", "paused", "stepping"]), )) .run() .pipe(Effect.orDie) + // Atomic wake can reach the parent only after a leaf checkpoint has + // completed the current graph. An additive extend emits this event to + // reopen that completed workflow without changing completed nodes. yield* db .update(WorkflowTable) - .set({ seq: event.durable!.seq, time_updated: toMillis(event.data.timestamp) }) + .set({ + status: "running", + wake_reported: false, + completed_at: null, + graph_rev: sql`${WorkflowTable.graph_rev} + 1`, + seq: event.durable!.seq, + time_updated: toMillis(event.data.timestamp), + }) .where(and( eq(WorkflowTable.id, event.data.dagID), - inArray(WorkflowTable.status, ["pending", "running", "paused", "stepping"]), + inArray(WorkflowTable.status, [...WorkflowStatusProjection.replanReopen.from]), )) .run() .pipe(Effect.orDie) + // Rev-view: mark the nodes this replan pushed out of the current + // revision — terminal rows the fragment bypassed (a failed node the + // new path routes around). plan.cancel rows are marked via the + // NodeCancelled projection instead. Idempotent fold: the marker is + // monotonic, replaying the event never resurrects a marked row. + const superseded = event.data.superseded + if (superseded && superseded.length > 0) { + yield* db + .update(WorkflowNodeTable) + .set({ superseded: true, seq: event.durable!.seq, time_updated: toMillis(event.data.timestamp) }) + .where(and( + eq(WorkflowNodeTable.workflow_id, event.data.dagID), + inArray(WorkflowNodeTable.id, [...superseded]), + )) + .run() + .pipe(Effect.orDie) + } }), ) @@ -361,10 +390,16 @@ export const layer = Layer.effectDiscard( // therefore never hold status="cancelled"; see NodeStatusProjection.cancelled // above and the canonical proof in // packages/opencode/test/dag/dag-escalation-clear-flag.test.ts:130-148. + // + // Rev-view (v1.0.15 Train A): a cancelled node leaves the current graph + // revision — the projection also sets the superseded marker so view and + // aggregation reads (summaries, status, node lists, rebuild input, wake + // attribution) show only the current rev. This covers both plan.cancel + // rows and explicit dag.nodeCancelled publishes (U1: shared semantics). yield* events.project(DagEvent.NodeCancelled, (event) => db .update(WorkflowNodeTable) - .set({ status: "failed", error_reason: "cancelled via replan", escalation_pending: false, seq: event.durable!.seq, time_updated: toMillis(event.data.timestamp) }) + .set({ status: "failed", superseded: true, error_reason: "cancelled via replan", escalation_pending: false, seq: event.durable!.seq, time_updated: toMillis(event.data.timestamp) }) .where(and( eq(WorkflowNodeTable.workflow_id, event.data.dagID), eq(WorkflowNodeTable.id, event.data.nodeID), diff --git a/packages/core/src/dag/sql.ts b/packages/core/src/dag/sql.ts index cc946745b4..9b22077345 100644 --- a/packages/core/src/dag/sql.ts +++ b/packages/core/src/dag/sql.ts @@ -47,6 +47,11 @@ export const WorkflowTable = sqliteTable( config: text().notNull(), // YAML string seq: integer().notNull(), // latest durable event seq wake_reported: integer({ mode: "boolean" }).notNull().default(false), // D3: has workflow terminal been reported to parent? + // Rev-view (v1.0.15 Train A): the current graph-revision counter. Bumped + // by the WorkflowReplanned projection; audit/telemetry only — the view + // predicate is the per-node `superseded` marker below. Default 1: legacy + // rows predate the concept and render exactly as before. + graph_rev: integer().notNull().default(1), started_at: integer(), completed_at: integer(), ...Timestamps, @@ -77,13 +82,20 @@ export const WorkflowNodeTable = sqliteTable( output: text({ mode: "json" }).$type(), error_reason: text(), error_class: text(), // dag.node.failed trigger (timeout/exec_failed/verdict_fail/push_exhausted) for failure triage - captured_output: text({ mode: "json" }).$type(), // durable payload from submit_result; survives a process crash, reset to null on a replan-restart via NodeStarted + captured_output: text({ mode: "json" }).$type(), // durable payload from submit_result, or a Train B file-ref record ({content_ref, size, sha256, summary}); survives a process crash, reset to null on a replan-restart via NodeStarted deadline_ms: integer(), // absolute deadline (spawnedAt + timeout_ms) for D0 termination boundary wake_eligible: integer({ mode: "boolean" }).notNull().default(false), // D6: node has report_to_parent=true wake_reported: integer({ mode: "boolean" }).notNull().default(false), // D3: has this node's terminal event been injected into the parent session? replan_attempts: integer().notNull().default(0), // D4: per-node replan counter for circuit breaker timeout_extensions: integer().notNull().default(0), // timeout escalation count (node stays running; main agent adjudicates) escalation_pending: integer({ mode: "boolean" }).notNull().default(false), // set on escalate, cleared on adjudication (extend) or new attempt — "awaiting main-agent adjudication" + // Rev-view (v1.0.15 Train A): this node was pushed OUT of the current + // graph revision by a replan (cancelled via replan, or a terminal row the + // fragment bypassed). Durable data is untouched — the marker only filters + // VIEW/aggregation reads (summaries, status, node lists, rebuild input, + // wake attribution) to the current revision. Monotonic: once true, stays + // true. Default false: legacy rows render exactly as before. + superseded: integer({ mode: "boolean" }).notNull().default(false), seq: integer().notNull(), // latest durable event seq for this node started_at: integer(), completed_at: integer(), diff --git a/packages/core/src/dag/store.ts b/packages/core/src/dag/store.ts index 40bc65d45b..861bafa78d 100644 --- a/packages/core/src/dag/store.ts +++ b/packages/core/src/dag/store.ts @@ -24,6 +24,8 @@ export interface WorkflowRow { config: string seq: number wakeReported: boolean + /** Rev-view (v1.0.15 Train A): current graph-revision counter (audit/telemetry). */ + graphRev: number startedAt: number | null completedAt: number | null timeCreated: number @@ -51,6 +53,8 @@ export interface NodeRow { replanAttempts: number timeoutExtensions: number escalationPending: boolean + /** Rev-view (v1.0.15 Train A): pushed out of the current graph revision by a replan. */ + superseded: boolean seq: number startedAt: number | null completedAt: number | null @@ -91,6 +95,7 @@ const mapWorkflow = (r: typeof WorkflowTable.$inferSelect): WorkflowRow => ({ config: r.config, seq: r.seq, wakeReported: r.wake_reported, + graphRev: r.graph_rev, startedAt: r.started_at, completedAt: r.completed_at, timeCreated: r.time_created, @@ -118,6 +123,7 @@ const mapNode = (r: typeof WorkflowNodeTable.$inferSelect): NodeRow => ({ replanAttempts: r.replan_attempts, timeoutExtensions: r.timeout_extensions, escalationPending: r.escalation_pending, + superseded: r.superseded, seq: r.seq, startedAt: r.started_at, completedAt: r.completed_at, @@ -156,6 +162,15 @@ export interface Interface { readonly getWorkflowSummaries: (sessionId: string) => Effect.Effect readonly getNodes: (workflowId: string) => Effect.Effect + /** + * Rev-view (v1.0.15 Train A): the CURRENT graph revision only — rows the + * replan pushed out of the graph (superseded) are filtered out. This is the + * read for VIEW and terminal-aggregation consumers: summaries, status/node + * listings, the loop's rebuild/recovery/completion input, and wake failure + * attribution. Durable truth is untouched — getNodes still returns every + * row, and completed old-rev outputs stay resolvable for input mapping. + */ + readonly getCurrentNodes: (workflowId: string) => Effect.Effect readonly getNode: (workflowId: string, nodeId: string) => Effect.Effect readonly getRunningNodes: (workflowId: string) => Effect.Effect readonly setCapturedOutput: (childSessionID: string, payload: unknown) => Effect.Effect @@ -257,6 +272,9 @@ export const layer = Layer.effect( if (wfRows.length === 0) return [] // P1-4: aggregate in SQL — pulling every node row into JS made each // dag.* event burst scale with total node count across the session. + // Rev-view (v1.0.15 Train A): superseded rows are filtered out so the + // counts reflect ONLY the current graph revision — a replaced segment + // neither counts toward nodeCount nor inflates failedNodes. const countRows = yield* db .select({ workflowId: WorkflowNodeTable.workflow_id, @@ -265,7 +283,7 @@ export const layer = Layer.effect( }) .from(WorkflowNodeTable) .innerJoin(WorkflowTable, eq(WorkflowNodeTable.workflow_id, WorkflowTable.id)) - .where(eq(WorkflowTable.session_id, sessionId)) + .where(and(eq(WorkflowTable.session_id, sessionId), eq(WorkflowNodeTable.superseded, false))) .groupBy(WorkflowNodeTable.workflow_id, WorkflowNodeTable.status) .all() .pipe(Effect.orDie) @@ -284,6 +302,7 @@ export const layer = Layer.effect( eq(WorkflowTable.session_id, sessionId), eq(WorkflowNodeTable.status, "running"), eq(WorkflowNodeTable.escalation_pending, true), + eq(WorkflowNodeTable.superseded, false), )) .groupBy(WorkflowNodeTable.workflow_id) .all() @@ -320,6 +339,17 @@ export const layer = Layer.effect( return rows.map(mapNode) }), + getCurrentNodes: Effect.fn("DagStore.getCurrentNodes")(function* (workflowId) { + const rows = yield* db + .select() + .from(WorkflowNodeTable) + .where(and(eq(WorkflowNodeTable.workflow_id, workflowId), eq(WorkflowNodeTable.superseded, false))) + .orderBy(desc(WorkflowNodeTable.seq)) + .all() + .pipe(Effect.orDie) + return rows.map(mapNode) + }), + getNode: Effect.fn("DagStore.getNode")(function* (workflowId, nodeId) { const row = yield* db .select() diff --git a/packages/core/src/database/migration.gen.ts b/packages/core/src/database/migration.gen.ts index a69db114e0..3d44bcf97d 100644 --- a/packages/core/src/database/migration.gen.ts +++ b/packages/core/src/database/migration.gen.ts @@ -54,6 +54,7 @@ export const migrations = ( import("./migration/20260811060000_goal_outcome"), import("./migration/20260813020344_bored_skaar"), import("./migration/20260813040429_workflow_directory"), + import("./migration/20260815044858_dag_graph_rev_view"), import("./migration/20260815083000_workflow_directory_convergence"), ]) ).map((module) => module.default) satisfies DatabaseMigration.Migration[] diff --git a/packages/core/src/database/migration/20260815044858_dag_graph_rev_view.ts b/packages/core/src/database/migration/20260815044858_dag_graph_rev_view.ts new file mode 100644 index 0000000000..7391cf1eab --- /dev/null +++ b/packages/core/src/database/migration/20260815044858_dag_graph_rev_view.ts @@ -0,0 +1,23 @@ +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +import { Effect } from "effect" +import type { DatabaseMigration } from "../migration" + +// Rev-view (v1.0.15 Train A, workflows/dag-engine-optimization.md). Legacy +// policy: existing rows migrate in place with superseded=false / graph_rev=1, +// so every pre-feature workflow renders EXACTLY as before — including its +// cancelled-via-replan rows, which stay visible and counted (config +// membership cannot be the current-rev predicate: <=v1.0.14 merged configs +// already drop cancelled nodes, so it would hide rows that render today). +// Marking only ever happens via the WorkflowReplanned and NodeCancelled +// projections after this migration runs. +export default { + id: "20260815044858_dag_graph_rev_view", + up(tx) { + return Effect.gen(function* () { + yield* tx.run(`ALTER TABLE \`workflow_node\` ADD \`superseded\` integer DEFAULT false NOT NULL;`) + yield* tx.run(`ALTER TABLE \`workflow\` ADD \`graph_rev\` integer DEFAULT 1 NOT NULL;`) + }) + }, +} satisfies DatabaseMigration.Migration diff --git a/packages/core/src/database/schema.gen.ts b/packages/core/src/database/schema.gen.ts index 7504e0fbea..25c9b5657b 100644 --- a/packages/core/src/database/schema.gen.ts +++ b/packages/core/src/database/schema.gen.ts @@ -91,6 +91,7 @@ export default { \`replan_attempts\` integer DEFAULT 0 NOT NULL, \`timeout_extensions\` integer DEFAULT 0 NOT NULL, \`escalation_pending\` integer DEFAULT false NOT NULL, + \`superseded\` integer DEFAULT false NOT NULL, \`seq\` integer NOT NULL, \`started_at\` integer, \`completed_at\` integer, @@ -111,6 +112,7 @@ export default { \`config\` text NOT NULL, \`seq\` integer NOT NULL, \`wake_reported\` integer DEFAULT false NOT NULL, + \`graph_rev\` integer DEFAULT 1 NOT NULL, \`started_at\` integer, \`completed_at\` integer, \`time_created\` integer NOT NULL, diff --git a/packages/core/test/dag-rev-view-legacy.test.ts b/packages/core/test/dag-rev-view-legacy.test.ts new file mode 100644 index 0000000000..7a4520d4ae --- /dev/null +++ b/packages/core/test/dag-rev-view-legacy.test.ts @@ -0,0 +1,131 @@ +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +/** + * Train A probe A-p4 (PIN, green before and after) — legacy no-rev workflows + * render unchanged (workflows/dag-engine-optimization.md, v1.0.15 ledger A4; + * evidence §7 legacy policy). + * + * Rows written before the rev-view feature carry no rev marker. The migration + * defaults (workflow.graph_rev = 1, workflow_node.superseded = false) must + * keep EVERY legacy row visible and counted exactly as v1.0.14 rendered it — + * including legacy cancelled-via-replan rows (status failed, error_reason + * "cancelled via replan"), which today appear in getNodes and count toward + * failedNodes in getWorkflowSummaries. + * + * The seed deliberately omits the new columns so the DB defaults apply, the + * same state a migrated install carries. + */ +import { describe, expect, test } from "bun:test" +import { Effect, Layer } from "effect" +import { Database } from "@opencode-ai/core/database/database" +import { ProjectV2 } from "@opencode-ai/core/project" +import { ProjectTable } from "@opencode-ai/core/project/sql" +import { WorkflowNodeTable, WorkflowTable } from "@opencode-ai/core/dag/sql" +import { DagStore } from "@opencode-ai/core/dag/store" +import { AbsolutePath } from "@opencode-ai/core/schema" +import { SessionSchema } from "@opencode-ai/core/session/schema" +import { SessionTable } from "@opencode-ai/core/session/sql" + +function storeLayer() { + const database = Database.layerFromPath(":memory:") + const store = DagStore.layer.pipe(Layer.provide(database)) + return Layer.merge(database, store) +} + +const projectID = ProjectV2.ID.make("project-1") +const sessionID = SessionSchema.ID.create() + +function node(workflowId: string, id: string, status: string, seq: number, errorReason?: string, errorClass?: string) { + return { + id, + workflow_id: workflowId, + name: id, + worker_type: "build", + status, + required: true, + depends_on: [], + wake_eligible: false, + wake_reported: false, + seq, + ...(errorReason !== undefined ? { error_reason: errorReason } : {}), + ...(errorClass !== undefined ? { error_class: errorClass } : {}), + } +} + +function seed() { + return Effect.gen(function* () { + const database = yield* Database.Service + yield* database.db.insert(ProjectTable).values({ + id: projectID, + worktree: AbsolutePath.make(process.cwd()), + sandboxes: [], + }).run().pipe(Effect.orDie) + yield* database.db.insert(SessionTable).values({ + id: sessionID, + project_id: projectID, + slug: "parent", + directory: process.cwd(), + title: "Parent", + version: "test", + }).run().pipe(Effect.orDie) + yield* database.db.insert(WorkflowTable).values({ + id: "wf-legacy", + project_id: projectID, + session_id: sessionID, + title: "Legacy", + status: "running", + config: "{}", + seq: 1, + wake_reported: false, + time_created: 1, + }).run().pipe(Effect.orDie) + yield* database.db.insert(WorkflowNodeTable).values([ + node("wf-legacy", "a", "completed", 1), + node("wf-legacy", "b", "completed", 2), + // Legacy cancelled-via-replan row: v1.0.14 renders it visible and + // counts it in failedNodes. That rendering must survive the feature. + node("wf-legacy", "old_cancelled", "failed", 3, "cancelled via replan"), + node("wf-legacy", "f1", "failed", 4, "boom", "exec_failed"), + ]).run().pipe(Effect.orDie) + }) +} + +describe("Train A rev-view — legacy rows render unchanged (A-p4 PIN)", () => { + test("legacy workflow (no rev markers) keeps every row visible and counted as in v1.0.14", async () => { + await Effect.runPromise( + Effect.gen(function* () { + const store = yield* DagStore.Service + yield* seed() + + const summaries = yield* store.getWorkflowSummaries(sessionID) + expect(summaries).toHaveLength(1) + expect(summaries[0]).toEqual({ + id: "wf-legacy", + title: "Legacy", + status: "running", + nodeCount: 4, + completedNodes: 2, + runningNodes: 0, + failedNodes: 2, + skippedNodes: 0, + queuedNodes: 0, + escalatedNodes: 0, + }) + + const rows = yield* store.getNodes("wf-legacy") + expect(rows.map((r) => r.id).sort()).toEqual(["a", "b", "f1", "old_cancelled"]) + const cancelled = rows.find((r) => r.id === "old_cancelled")! + expect(cancelled.status).toBe("failed") + expect(cancelled.errorReason).toBe("cancelled via replan") + + // Migration defaults (the legacy policy itself): rows that predate the + // feature carry superseded=false, and the workflow carries graph_rev=1. + // Any other default would change the rendering pinned above. + expect(rows.every((r) => !r.superseded)).toBe(true) + const wf = yield* store.getWorkflow("wf-legacy") + expect(wf?.graphRev).toBe(1) + }).pipe(Effect.provide(storeLayer()), Effect.scoped), + ) + }) +}) diff --git a/packages/core/test/dag-rev-view-projection.test.ts b/packages/core/test/dag-rev-view-projection.test.ts new file mode 100644 index 0000000000..4b37468cc8 --- /dev/null +++ b/packages/core/test/dag-rev-view-projection.test.ts @@ -0,0 +1,193 @@ +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +/** + * Train A rev-view — projector marking (implementation-side proof for + * workflows/dag-engine-optimization.md, v1.0.15 ledger A4 decision M1). + * + * Pins the projection seams directly on the real projector SQL: + * - NodeCancelled marks the row superseded (plan.cancel + explicit cancel). + * - WorkflowReplanned marks its optional superseded list (terminal rows the + * fragment bypassed, which the engine never cancels) and bumps graph_rev. + * - Legacy WorkflowReplanned events WITHOUT the superseded field still decode + * and project (replay safety — the field is optional, dag-event.ts pattern). + * - The NodeRegistered upsert never resets the marker (replace/restart + * re-registration must not resurrect a superseded row). + * - The reopen leg (completed → running) bumps graph_rev as well. + */ +import { describe, expect, test } from "bun:test" +import { DateTime, Effect, Layer } from "effect" +import { Database } from "@opencode-ai/core/database/database" +import { EventV2 } from "@opencode-ai/core/event" +import { DagProjector } from "@opencode-ai/core/dag/projector" +import { DagStore } from "@opencode-ai/core/dag/store" +import { WorkflowNodeTable, WorkflowTable } from "@opencode-ai/core/dag/sql" +import { ProjectV2 } from "@opencode-ai/core/project" +import { ProjectTable } from "@opencode-ai/core/project/sql" +import { AbsolutePath } from "@opencode-ai/core/schema" +import { SessionSchema } from "@opencode-ai/core/session/schema" +import { SessionTable } from "@opencode-ai/core/session/sql" +import { DagEvent } from "@opencode-ai/schema/dag-event" + +function projectorLayer() { + const database = Database.layerFromPath(":memory:") + const eventLayer = EventV2.layer.pipe(Layer.provide(database)) + const projector = DagProjector.layer.pipe(Layer.provide(Layer.merge(database, eventLayer))) + const store = DagStore.layer.pipe(Layer.provide(database)) + return Layer.mergeAll(database, eventLayer, projector, store) +} + +const projectID = ProjectV2.ID.make("project-1") +const sessionID = SessionSchema.ID.create() + +function seed(workflowStatus: string) { + return Effect.gen(function* () { + const { db } = yield* Database.Service + yield* db.insert(ProjectTable).values({ + id: projectID, + worktree: AbsolutePath.make(process.cwd()), + sandboxes: [], + }).run().pipe(Effect.orDie) + yield* db.insert(SessionTable).values({ + id: sessionID, + project_id: projectID, + slug: "parent", + directory: process.cwd(), + title: "Parent", + version: "test", + }).run().pipe(Effect.orDie) + yield* db.insert(WorkflowTable).values({ + id: "dag_rev", + project_id: projectID, + session_id: sessionID, + title: "Rev projection", + status: workflowStatus, + config: "{}", + seq: 1, + wake_reported: false, + }).run().pipe(Effect.orDie) + yield* db.insert(WorkflowNodeTable).values([ + { + id: "n1", + workflow_id: "dag_rev", + name: "n1", + worker_type: "build", + status: "completed", + required: true, + depends_on: [], + wake_eligible: false, + wake_reported: false, + seq: 1, + }, + { + id: "n2", + workflow_id: "dag_rev", + name: "n2", + worker_type: "build", + status: "failed", + required: true, + depends_on: ["n1"], + error_reason: "simulated exec failure", + error_class: "exec_failed", + wake_eligible: false, + wake_reported: false, + seq: 2, + }, + { + id: "n3", + workflow_id: "dag_rev", + name: "n3", + worker_type: "build", + status: "pending", + required: true, + depends_on: ["n1"], + wake_eligible: false, + wake_reported: false, + seq: 3, + }, + ]).run().pipe(Effect.orDie) + }) +} + +function replan(input: { dagID: string; superseded?: string[]; added?: number }) { + return Effect.gen(function* () { + const events = yield* EventV2.Service + yield* events.publish(DagEvent.WorkflowReplanned, { + dagID: DagEvent.DagID.make(input.dagID), + added: input.added ?? 1, + removed: 0, + replaced: 0, + restarted: 0, + ...(input.superseded + ? { superseded: input.superseded.map((id) => DagEvent.NodeID.make(id)) } + : {}), + timestamp: yield* DateTime.now, + }) + }) +} + +describe("Train A rev-view — projector marking", () => { + test("NodeCancelled marks superseded; WorkflowReplanned marks its supersede list and bumps graph_rev", async () => { + await Effect.runPromise( + Effect.gen(function* () { + yield* seed("running") + const events = yield* EventV2.Service + const store = yield* DagStore.Service + + // plan.cancel seam: cancel pending n3 via NodeCancelled. + yield* events.publish(DagEvent.NodeCancelled, { + dagID: DagEvent.DagID.make("dag_rev"), + nodeID: DagEvent.NodeID.make("n3"), + timestamp: yield* DateTime.now, + }) + expect((yield* store.getNode("dag_rev", "n3"))?.superseded).toBe(true) + expect((yield* store.getNode("dag_rev", "n3"))?.status).toBe("failed") + + // Bypassed-failure seam: n2 is terminal-failed and never cancelled — + // the WorkflowReplanned supersede list marks it. Legacy n1 stays. + yield* replan({ dagID: "dag_rev", superseded: ["n2"] }) + expect((yield* store.getNode("dag_rev", "n2"))?.superseded).toBe(true) + expect((yield* store.getNode("dag_rev", "n1"))?.superseded).toBe(false) + expect((yield* store.getWorkflow("dag_rev"))?.graphRev).toBe(2) + + // Replay-safe legacy shape: NO superseded field still decodes, + // bumps the revision, and marks nothing new. + yield* replan({ dagID: "dag_rev" }) + expect((yield* store.getWorkflow("dag_rev"))?.graphRev).toBe(3) + expect((yield* store.getNode("dag_rev", "n1"))?.superseded).toBe(false) + + // The NodeRegistered upsert must NOT reset the marker — replace and + // restart re-publish definitions for the same id. + yield* events.publish(DagEvent.NodeRegistered, { + dagID: DagEvent.DagID.make("dag_rev"), + nodeID: DagEvent.NodeID.make("n2"), + name: "n2-renamed", + workerType: "build", + dependsOn: ["n1"].map((id) => DagEvent.NodeID.make(id)), + required: true, + timestamp: yield* DateTime.now, + }) + const n2 = yield* store.getNode("dag_rev", "n2") + expect(n2?.superseded).toBe(true) + expect(n2?.name).toBe("n2-renamed") + + // View reads see only the current revision; durable reads see all. + expect((yield* store.getCurrentNodes("dag_rev")).map((n) => n.id)).toEqual(["n1"]) + expect((yield* store.getNodes("dag_rev")).map((n) => n.id).sort()).toEqual(["n1", "n2", "n3"]) + }).pipe(Effect.provide(projectorLayer()), Effect.scoped), + ) + }) + + test("the reopen leg (completed workflow) bumps graph_rev too", async () => { + await Effect.runPromise( + Effect.gen(function* () { + yield* seed("completed") + const store = yield* DagStore.Service + yield* replan({ dagID: "dag_rev" }) + const wf = yield* store.getWorkflow("dag_rev") + expect(wf?.status).toBe("running") + expect(wf?.graphRev).toBe(2) + }).pipe(Effect.provide(projectorLayer()), Effect.scoped), + ) + }) +}) diff --git a/packages/core/test/dag-rev-view-wash.test.ts b/packages/core/test/dag-rev-view-wash.test.ts new file mode 100644 index 0000000000..c75d8f8801 --- /dev/null +++ b/packages/core/test/dag-rev-view-wash.test.ts @@ -0,0 +1,101 @@ +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +/** + * Train A probe A-p5(b) (PIN, green before and after) — reopen wash by + * workflow status (workflows/dag-engine-optimization.md, v1.0.15 ledger A4; + * evidence §7 U3). + * + * A completed workflow reopened by an additive extend (WorkflowReplanned + * projection's sanctioned completed→running exception) stops re-waking + * because the wake read filters on TERMINAL workflow status. This pin locks + * the observable part of A-p5(b): after the reopen projection the workflow + * is no longer returned as an unreported wake, regardless of rev semantics. + */ +import { describe, expect, test } from "bun:test" +import { DateTime, Effect, Layer } from "effect" +import { Database } from "@opencode-ai/core/database/database" +import { EventV2 } from "@opencode-ai/core/event" +import { DagProjector } from "@opencode-ai/core/dag/projector" +import { DagStore } from "@opencode-ai/core/dag/store" +import { WorkflowTable } from "@opencode-ai/core/dag/sql" +import { ProjectV2 } from "@opencode-ai/core/project" +import { ProjectTable } from "@opencode-ai/core/project/sql" +import { AbsolutePath } from "@opencode-ai/core/schema" +import { SessionSchema } from "@opencode-ai/core/session/schema" +import { SessionTable } from "@opencode-ai/core/session/sql" +import { DagEvent } from "@opencode-ai/schema/dag-event" + +function projectorLayer() { + const database = Database.layerFromPath(":memory:") + const eventLayer = EventV2.layer.pipe(Layer.provide(database)) + const projector = DagProjector.layer.pipe(Layer.provide(Layer.merge(database, eventLayer))) + const store = DagStore.layer.pipe(Layer.provide(database)) + return Layer.mergeAll(database, eventLayer, projector, store) +} + +const projectID = ProjectV2.ID.make("project-1") +const sessionID = SessionSchema.ID.create() + +function seed() { + return Effect.gen(function* () { + const { db } = yield* Database.Service + yield* db.insert(ProjectTable).values({ + id: projectID, + worktree: AbsolutePath.make(process.cwd()), + sandboxes: [], + }).run().pipe(Effect.orDie) + yield* db.insert(SessionTable).values({ + id: sessionID, + project_id: projectID, + slug: "parent", + directory: process.cwd(), + title: "Parent", + version: "test", + }).run().pipe(Effect.orDie) + yield* db.insert(WorkflowTable).values({ + id: "dag_reopen", + project_id: projectID, + session_id: sessionID, + title: "Reopen wash", + status: "completed", + config: "{}", + seq: 1, + wake_reported: true, + started_at: 1, + completed_at: 2, + }).run().pipe(Effect.orDie) + }) +} + +describe("Train A rev-view — reopen wash (A-p5(b) PIN)", () => { + test("WorkflowReplanned reopen puts the workflow back to running and out of the terminal wake read", async () => { + await Effect.runPromise( + Effect.gen(function* () { + yield* seed() + const events = yield* EventV2.Service + const store = yield* DagStore.Service + + yield* events.publish(DagEvent.WorkflowReplanned, { + dagID: DagEvent.DagID.make("dag_reopen"), + added: 1, + removed: 0, + replaced: 0, + restarted: 0, + timestamp: yield* DateTime.now, + }) + + const wf = yield* store.getWorkflow("dag_reopen") + expect(wf?.status).toBe("running") + expect(wf?.wakeReported).toBe(false) + expect(wf?.completedAt).toBeNull() + + // Washed by status: the terminal wake read only returns workflows in + // a terminal status — a reopened (running) workflow is not among the + // unreported wakes, so it stops re-waking the parent. + const unreported = yield* store.getUnreportedWakeWorkflows(sessionID) + expect(unreported.map((w) => w.id)).toEqual([]) + }).pipe(Effect.provide(projectorLayer()), Effect.scoped), + ) + }) +}) diff --git a/packages/opencode/src/dag/dag.ts b/packages/opencode/src/dag/dag.ts index e5f97c6e02..bb0b980538 100644 --- a/packages/opencode/src/dag/dag.ts +++ b/packages/opencode/src/dag/dag.ts @@ -449,7 +449,11 @@ export const layer = Layer.effect( // counts as in-flight: the node is durably admitted and will start once // a permit frees (P0-2) — stepping alongside it would put two nodes in // flight. - const nodes = yield* store.getNodes(dagID) + // Rev-view (v1.0.15 Train A): step computes readiness over the CURRENT + // graph revision — superseded rows must not re-seed their failures or + // edges into the transient runtime. Superseded rows are terminal, so + // the in-flight check is unaffected by the filter. + const nodes = yield* store.getCurrentNodes(dagID) const hasInFlight = nodes.some((n) => n.status === "running" || n.status === "queued") if (hasInFlight) return yield* Effect.fail(new Error(`Node still in-flight: cannot step ${dagID}`)) // Compute ready nodes using a transient WorkflowRuntime. @@ -681,6 +685,14 @@ export const layer = Layer.effect( removed: effectivePlan.cancel.length as never, replaced: effectivePlan.replace.length as never, restarted: effectivePlan.restart.length as never, + // Rev-view (v1.0.15 Train A): the terminal-FAILED rows at replan time + // are the segment the new revision replaces — left in the rebuild + // input they re-seed as required-unsatisfied and weld the workflow to + // failure (the wake-up bug this train breaks). plan.cancel rows are + // marked superseded by the NodeCancelled projection instead; this + // list carries the genuine failures the fragment bypasses, which the + // engine never cancels. Durable rows stay untouched. + superseded: nodes.filter((n) => n.status === "failed").map((n) => DagEvent.NodeID.make(n.id)), timestamp: yield* DateTime.now, }) return { cancel: effectivePlan.cancel, restart: effectivePlan.restart, replace: effectivePlan.replace, add: effectivePlan.add, ignore: effectivePlan.ignore } diff --git a/packages/opencode/src/dag/runtime/loop.ts b/packages/opencode/src/dag/runtime/loop.ts index 24a2ef8f50..9076807b7c 100644 --- a/packages/opencode/src/dag/runtime/loop.ts +++ b/packages/opencode/src/dag/runtime/loop.ts @@ -263,6 +263,7 @@ const serviceLayer = Layer.effect( nodeID, node, parentSessionID: entry.parentSessionID, + directory: ctx.directory, promptParts, outputSchema: nodeConfig?.output_schema as Record | undefined, timeoutMs: nodeConfig?.worker_config?.timeout_ms, @@ -301,7 +302,10 @@ const serviceLayer = Layer.effect( // cancellation handler can reach this point before WorkflowReplanned // rebuilds the in-memory graph. Durable active nodes prove that the // apparent completion belongs to an obsolete graph generation. - const nodes = yield* store.getNodes(dagID) + // Rev-view (v1.0.15 Train A): the guard read filters to the current + // revision — superseded rows are terminal and never trip the guard, + // and the review-outcome input below must judge only current nodes. + const nodes = yield* store.getCurrentNodes(dagID) const hasUnseenActiveNode = nodes.some( (node) => !isNodeTerminalStatus(node.status as never) && !entry.runtime.containsNode(node.id), ) @@ -419,7 +423,12 @@ const serviceLayer = Layer.effect( ownershipLost: recovery.ownershipLost, }) } - const nodes = yield* store.getNodes(dagID) + // Rev-view (v1.0.15 Train A): the recovery runtime rebuilds from + // the CURRENT graph revision — a crash must not revive the + // replaced segment's failures (same rebuild input as the + // WorkflowReplanned handler). Durable truth (the unfiltered read + // in recovery reconcile above) is untouched. + const nodes = yield* store.getCurrentNodes(dagID) const maxConcurrency = Math.max(1, config?.max_concurrency ?? Dag.DEFAULT_WORKFLOW_CONFIG.maxConcurrency) const runtime = new WorkflowRuntime(toSchedulingNodes(nodes), maxConcurrency) const semaphore = Semaphore.makeUnsafe(maxConcurrency) @@ -549,7 +558,10 @@ const serviceLayer = Layer.effect( // directory contexts. if (!(yield* DagLocation.ownsWorkflow(dagID, ctx.directory))) return const config = parseWorkflowConfig(wf.config) - const nodes = yield* store.getNodes(dagID) + // Rev-view (v1.0.15 Train A): same current-revision rebuild + // input as recoverWorkflow — adoption must not resurrect the + // replaced segment's failures either. + const nodes = yield* store.getCurrentNodes(dagID) const maxConcurrency = Math.max(1, config?.max_concurrency ?? Dag.DEFAULT_WORKFLOW_CONFIG.maxConcurrency) const runtime = new WorkflowRuntime(toSchedulingNodes(nodes), maxConcurrency) const semaphore = Semaphore.makeUnsafe(maxConcurrency) @@ -848,7 +860,15 @@ const serviceLayer = Layer.effect( const wf = yield* store.getWorkflow(dagID).pipe(Effect.orDie) const oldConfig = entry.config if (wf) entry.config = parseWorkflowConfig(wf.config) - const nodes = yield* store.getNodes(dagID) + // Rev-view (v1.0.15 Train A): THE aggregation filter point. + // The rebuild input is the CURRENT graph revision only — + // superseded rows (cancelled via replan, or terminal + // failures the fragment bypassed) must not re-seed as + // required-unsatisfied, or the wake-up bug fails a workflow + // on its replaced segment despite the new path succeeding. + // The re-time loop below only sees running rows (superseded + // rows are terminal), so the filter changes nothing there. + const nodes = yield* store.getCurrentNodes(dagID) entry.runtime.rebuildGraph(toSchedulingNodes(nodes)) // Timeout extension (Q6): a running node gets a recomputed // deadline (now + new timeout) ONLY when the replan carries a @@ -1231,7 +1251,11 @@ const serviceLayer = Layer.effect( const failuresByWorkflow = new Map() for (const workflow of batch.workflows) { if (workflow.status !== "failed") continue - const failedNodes = yield* store.getNodes(workflow.id).pipe( + // Rev-view (v1.0.15 Train A): attribution is terminal + // aggregation — only CURRENT-revision failures are attributed. + // Superseded replaced failures (cancelled rows already fall + // out via errorClass null) must not be re-attributed. + const failedNodes = yield* store.getCurrentNodes(workflow.id).pipe( Effect.map((nodes) => nodes.filter((node): node is DagStore.NodeRow & { errorClass: string } => node.status === "failed" && node.errorClass !== null)), Effect.catchCause((cause) => Effect.gen(function* () { diff --git a/packages/opencode/src/dag/runtime/output-ref.ts b/packages/opencode/src/dag/runtime/output-ref.ts new file mode 100644 index 0000000000..0ceceda422 --- /dev/null +++ b/packages/opencode/src/dag/runtime/output-ref.ts @@ -0,0 +1,136 @@ +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +/** + * Node-output file references (v1.0.15 Train B ledger B1–B4). + * + * A long-report node without an output_schema may submit an absolute file + * path instead of inlining its text. When the submitted string IS an + * absolute path, the runtime validates existence + size>0 at submit time and + * records `{content_ref, size, sha256}` — plus the summary the result seam + * serves — in captured_output: an integrity receipt for later verification + * (B2). output_schema nodes are untouched (inline JSON, B1) and existing + * inline outputs coexist with no migration (B3): any anomaly (missing file, + * empty file, directory, prose around the path) simply leaves the legacy + * inline behavior. + * + * Report-area convention: `/.opencode/workflow-reports/` (named + * after the existing `.opencode/workflow-drafts/` area). On the first + * capture into the area the runtime ensures the project `.gitignore` carries + * the area entry — append-only, idempotent, never overwrite (B4). + */ + +import { appendFile, readFile, stat } from "node:fs/promises" +import path from "node:path" +import { Effect } from "effect" +import { FSUtil } from "@opencode-ai/core/fs-util" +import { Hash } from "@opencode-ai/core/util/hash" + +export interface OutputFileRef { + /** Discriminator against output_schema captured payloads (inline JSON). */ + kind: "file_ref" + /** Durable reference for the result seam — the parent agent fetches content itself. */ + content_ref: string + /** Absolute path the read tool fetches. */ + path: string + /** Byte size at submit time. */ + size: number + /** SHA-256 of the file bytes at submit time (integrity for later verification). */ + sha256: string + /** First ~200 chars of content at submit time — stable even if the file drifts later. */ + summary: string +} + +const SUMMARY_CHARS = 200 +// The summary only needs the leading chars; decoding a bounded prefix keeps a +// giant report from being copied twice (once for the digest, once for text). +const SUMMARY_DECODE_BYTES = 4096 +const MAX_PATH_CHARS = 4096 + +export const REPORT_AREA = path.join(".opencode", "workflow-reports") +// Gitignore patterns are slash-separated on every platform. +const REPORT_GITIGNORE_ENTRY = ".opencode/workflow-reports/" + +export function isOutputFileRef(value: unknown): value is OutputFileRef { + if (!isRecord(value)) return false + return ( + value.kind === "file_ref" && + typeof value.content_ref === "string" && + typeof value.path === "string" && + typeof value.size === "number" && + typeof value.sha256 === "string" && + typeof value.summary === "string" + ) +} + +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value) +} + +/** + * Submit-time detection: the trimmed reply must BE one absolute path (no + * surrounding prose, no inner whitespace — the single-token contract keeps a + * sentence that mentions a path inline). Existence + regular file + size>0 + * are validated at submit time; every anomaly resolves to `undefined` so the + * caller falls back to the exact legacy inline behavior (never fails the + * node — the capture is audit metadata, not the settlement). + */ +export function captureOutputFileRef(rawText: string): Effect.Effect { + const candidate = rawText.trim() + if ( + candidate.length === 0 || + candidate.length > MAX_PATH_CHARS || + /\s/.test(candidate) || + !path.isAbsolute(candidate) + ) { + return Effect.succeed(undefined) + } + return Effect.gen(function* () { + const info = yield* Effect.promise(() => stat(candidate).catch(() => undefined)) + if (!info || !info.isFile() || info.size === 0) return undefined + const bytes = yield* Effect.promise(() => + Bun.file(candidate) + .arrayBuffer() + .catch(() => undefined), + ) + if (!bytes || bytes.byteLength === 0) return undefined + const text = new TextDecoder().decode(bytes.slice(0, SUMMARY_DECODE_BYTES)) + const summary = text.length > SUMMARY_CHARS ? `${text.slice(0, SUMMARY_CHARS)}\u2026` : text + return { + kind: "file_ref" as const, + content_ref: candidate, + path: candidate, + size: bytes.byteLength, + sha256: Hash.sha256(Buffer.from(bytes)), + summary, + } + }).pipe(Effect.orElseSucceed(() => undefined)) +} + +/** + * First-write gitignore guarantee for the project `.opencode/` report area + * (B4). Fires only when the captured ref lies inside + * `/.opencode/workflow-reports/` — cross-worktree refs + * (/private/tmp/...) have no project gitignore to touch. Append-only and + * idempotent: an existing entry (or an already-covering `.opencode/` rule) + * leaves the file untouched; pre-existing entries are preserved. Best-effort + * — a permission blip must never fail the node completion. + */ +export function ensureReportAreaGitignore(directory: string, refPath: string): Effect.Effect { + if (!FSUtil.contains(path.join(directory, REPORT_AREA), refPath)) return Effect.void + const gitignorePath = path.join(directory, ".gitignore") + return Effect.gen(function* () { + const existing = yield* Effect.promise(() => readFile(gitignorePath, "utf8").catch(() => undefined)) + const covered = existing + ?.split("\n") + .map((line) => line.trim()) + .some((line) => [REPORT_GITIGNORE_ENTRY, ".opencode/workflow-reports", ".opencode/", ".opencode"].includes(line)) + if (covered) return + const separator = existing === undefined || existing.length === 0 || existing.endsWith("\n") ? "" : "\n" + yield* Effect.promise(() => appendFile(gitignorePath, `${separator}${REPORT_GITIGNORE_ENTRY}\n`)) + } ).pipe( + Effect.catchCause((cause) => + Effect.logWarning("failed to ensure the workflow-reports gitignore entry", { directory, cause }), + ), + ) +} diff --git a/packages/opencode/src/dag/runtime/spawn.ts b/packages/opencode/src/dag/runtime/spawn.ts index d705c2890b..8ccad71b82 100644 --- a/packages/opencode/src/dag/runtime/spawn.ts +++ b/packages/opencode/src/dag/runtime/spawn.ts @@ -46,6 +46,7 @@ import type { DagStore } from "@opencode-ai/core/dag/store" import { ModelV2 } from "@opencode-ai/core/model" import { ProviderV2 } from "@opencode-ai/core/provider" import { registerCaptureSlot, clearCaptureSlot, settleCapturedOutput } from "./capture" +import { captureOutputFileRef, ensureReportAreaGitignore } from "./output-ref" type PromptParts = SessionPrompt.PromptInput["parts"] @@ -55,6 +56,8 @@ export interface NodeSpawnInput { node: DagStore.NodeRow parentSessionID: string promptParts: PromptParts + /** Workflow execution directory — keys the report-area gitignore guarantee (Train B, B4). */ + directory?: string outputSchema?: Record timeoutMs?: number reportToParent?: boolean @@ -493,6 +496,23 @@ export function spawnNode( ) return } + // Train B (v1.0.15 B2): submit-time file-ref detection — when the + // reply IS an existing non-empty absolute path, record + // {content_ref, size, sha256, summary} in captured_output (the + // same durable column submit_result uses; reset on NodeStarted). + // The settlement stays the raw string: input mapping, wake + // digests, and legacy readers keep the exact inline behavior, + // and the result seam serves the pointer (B3). Any anomaly keeps + // the inline path — the capture must never fail the node. + const fileRef = yield* captureOutputFileRef(rawText) + if (fileRef && childSessionID) { + yield* dag.store.setCapturedOutput(childSessionID, fileRef).pipe( + Effect.catchCause((cause) => + Effect.logWarning("output-ref capture persistence failed — inline output preserved", { cause }), + ), + ) + if (input.directory) yield* ensureReportAreaGitignore(input.directory, fileRef.path) + } yield* dag.nodeCompleted(input.dagID, input.nodeID, rawText).pipe( Effect.catchIf( isTransitionRejection, diff --git a/packages/opencode/src/server/routes/instance/httpapi/handlers/dag.ts b/packages/opencode/src/server/routes/instance/httpapi/handlers/dag.ts index 2a73c1afcc..bc069897ec 100644 --- a/packages/opencode/src/server/routes/instance/httpapi/handlers/dag.ts +++ b/packages/opencode/src/server/routes/instance/httpapi/handlers/dag.ts @@ -129,7 +129,11 @@ export const dagHandlers = HttpApiBuilder.group(InstanceHttpApi, "dag", (handler const nodes = Effect.fn("DagHttpApi.nodes")(function* (ctx: { params: { dagID: string } }) { yield* requireWorkflow(ctx.params.dagID) - const rows = yield* dag.store.getNodes(ctx.params.dagID).pipe(Effect.orDie) + // Rev-view (v1.0.15 Train A): the TUI node list is a view seam — it + // renders the CURRENT graph revision only (zero TUI changes: the + // server filters). nodeDetail below keeps the unfiltered read so a + // superseded node's durable state stays auditable by id. + const rows = yield* dag.store.getCurrentNodes(ctx.params.dagID).pipe(Effect.orDie) return rows.map(node) }) diff --git a/packages/opencode/src/tool/workflow.ts b/packages/opencode/src/tool/workflow.ts index 314164d29a..58c42b86de 100644 --- a/packages/opencode/src/tool/workflow.ts +++ b/packages/opencode/src/tool/workflow.ts @@ -16,6 +16,7 @@ import { Provider } from "@/provider/provider" import { Session } from "@/session/session" import { SessionID } from "@/session/schema" import { createAdmissionRecord } from "@/dag/admission" +import { isOutputFileRef } from "@/dag/runtime/output-ref" import { TerminalViolationError } from "@opencode-ai/core/dag/core/types" import { FSUtil } from "@opencode-ai/core/fs-util" import { stringify as yamlStringify } from "yaml" @@ -445,7 +446,10 @@ export const WorkflowTool = Tool.define< } case "status": { const workflow = yield* requireOwnedWorkflow(params.workflow_id, ctx.sessionID) - const nodes = yield* dag.store.getNodes(params.workflow_id).pipe(Effect.orDie) + // Rev-view (v1.0.15 Train A): status is a view seam — it shows + // the CURRENT graph revision only. Superseded replaced segments + // stay reachable via the result seam (getNode is unfiltered). + const nodes = yield* dag.store.getCurrentNodes(params.workflow_id).pipe(Effect.orDie) const config = Dag.parseWorkflowConfig(workflow.config) return { title: `Workflow status: ${workflow.title}`, @@ -496,6 +500,37 @@ export const WorkflowTool = Tool.define< if (!node) { return yield* Effect.die(new Error(`Workflow node not found: ${params.workflow_id}/${params.node_id}`)) } + // Train B (v1.0.15 B3): a submit-time file ref reads as a durable + // pointer — content_ref + summary + path — and the parent agent + // fetches the content itself (read tool). No paging: the pointer + // is bounded by construction and next_cursor is never issued, so + // inline outputs below keep the exact legacy paged read. + if (isOutputFileRef(node.capturedOutput)) { + return { + title: `Workflow result: ${node.name}`, + output: JSON.stringify( + { + workflow_id: params.workflow_id, + node_id: params.node_id, + status: node.status, + content_ref: node.capturedOutput.content_ref, + path: node.capturedOutput.path, + summary: node.capturedOutput.summary, + size: node.capturedOutput.size, + sha256: node.capturedOutput.sha256, + truncated: false, + next_cursor: null, + }, + null, + 2, + ), + metadata: { + workflowId: params.workflow_id, + nodeId: params.node_id, + truncated: false, + } as Metadata, + } + } const cursor = params.cursor ? decodeResultCursor(Buffer.from(params.cursor, "base64url").toString()) : Option.some( diff --git a/packages/opencode/test/dag/dag-output-ref-result.test.ts b/packages/opencode/test/dag/dag-output-ref-result.test.ts new file mode 100644 index 0000000000..f0577fc4c7 --- /dev/null +++ b/packages/opencode/test/dag/dag-output-ref-result.test.ts @@ -0,0 +1,234 @@ +// oxlint-disable typescript-eslint/no-unsafe-type-assertion -- tool-level +// probes deliberately mirror dag-rev-view-status.test.ts: mocked service +// layers and row fixtures use `as never`-style shims (mock objects implement +// only the interface slice the scenario exercises). The shims are type-only; +// converting them would fork the template's shape without changing behavior. +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +/** + * Train B probe B-p3 — workflow tool `result` action seam + * (workflows/dag-engine-optimization.md, v1.0.15 ledger, decision B3). + * + * The `result` action is how the parent agent reads a durable node output. + * For a node whose submit-time capture recorded a file_ref in captured_output, + * the action must return the durable pointer — content_ref + summary (first + * ~200 chars, captured at submit time) + path — instead of paging raw text; + * the parent agent fetches content itself with the read tool. Inline outputs + * (legacy strings and output_schema payloads) keep the existing paged read + * verbatim — no migration, no shape change for them. + * + * RED on the unmodified engine: the result action never inspects + * captured_output and pages the raw output string, so the content_ref field + * is absent and the summary is missing. + */ +import { describe, expect } from "bun:test" +import os from "node:os" +import path from "node:path" +import { Effect, Layer } from "effect" +import { DagStore } from "@opencode-ai/core/dag/store" +import { ProjectV2 } from "@opencode-ai/core/project" +import { ModelV2 } from "@opencode-ai/core/model" +import { ProviderV2 } from "@opencode-ai/core/provider" +import { Event } from "@opencode-ai/schema/event" +import { Agent } from "@/agent/agent" +import { Dag } from "@/dag/dag" +import { EventV2Bridge } from "@/event-v2-bridge" +import { Question } from "@/question" +import { Session } from "@/session/session" +import { MessageID, SessionID } from "@/session/schema" +import { Skill } from "@/skill" +import type { Tool } from "@/tool/tool" +import { Truncate } from "@/tool/truncate" +import { WorkflowTool } from "@/tool/workflow" +import { Provider } from "@/provider/provider" +import { testEffect } from "../lib/effect" + +const projectID = ProjectV2.ID.make("project_test") + +const refPath = path.join(os.tmpdir(), "dag-ref-result", "report.md") +const refRecord = { + kind: "file_ref" as const, + content_ref: refPath, + path: refPath, + size: 4096, + sha256: "0".repeat(64), + summary: "Report summary captured at submit time", +} +const refRow = { + id: "node_ref", + workflowId: "dag_output_ref", + name: "Report node", + workerType: "build", + status: "completed", + required: true, + dependsOn: [], + modelId: null, + modelProviderId: null, + childSessionId: null, + output: refPath, + capturedOutput: refRecord, + errorReason: null, + errorClass: null, + deadlineMs: null, + wakeEligible: false, + wakeReported: false, + replanAttempts: 0, + timeoutExtensions: 0, + escalationPending: false, + superseded: false, + seq: 1, + startedAt: 1, + completedAt: 2, + timeCreated: 1, + timeUpdated: 2, +} +const inlineRow = { + ...refRow, + id: "node_inline", + name: "Inline node", + output: "plain inline output", + capturedOutput: null, + seq: 2, +} +const rows = [refRow, inlineRow] + +const storeImpl = { + getWorkflow: (id: string) => + Effect.succeed( + id === "dag_output_ref" + ? { + id, + projectId: projectID, + sessionId: "ses_ref_parent", + directory: null, + title: "Output ref workflow", + status: "completed", + config: "{}", + seq: 2, + wakeReported: true, + graphRev: 1, + startedAt: 1, + completedAt: 2, + timeCreated: 1, + timeUpdated: 2, + } + : undefined, + ), + getNodes: () => Effect.succeed(rows), + getCurrentNodes: () => Effect.succeed(rows), + getNode: (workflowID: string, nodeID: string) => + Effect.succeed(rows.find((node) => node.workflowId === workflowID && node.id === nodeID)), +} + +const dag = Dag.layer.pipe( + Layer.provide(Layer.mock(DagStore.Service, storeImpl)), + Layer.provide(Layer.mock(EventV2Bridge.Service, { + publish: (definition, data) => + Effect.succeed({ id: Event.ID.create(), type: definition.type, data }), + })), +) + +const runtime = testEffect( + Layer.mergeAll( + Layer.mock(Agent.Service, { + get: () => + Effect.succeed({ + name: "build", + mode: "all", + permission: [], + options: {}, + description: "", + prompt: "", + model: { providerID: ProviderV2.ID.make("test"), modelID: ModelV2.ID.make("test-model") }, + tools: {}, + hooks: {}, + }), + list: () => Effect.succeed([]), + }), + Layer.mock(Skill.Service, { all: () => Effect.succeed([]) }), + Layer.mock(Truncate.Service, { + output: (content) => Effect.succeed({ content, truncated: false }), + }), + Layer.mock(Question.Service, { ask: () => Effect.succeed([]) }), + Layer.mock(Provider.Service, { + list: () => Effect.succeed({}), + getModel: (providerID, modelID) => Effect.fail(new Provider.ModelNotFoundError({ providerID, modelID })), + }), + dag, + Layer.mock(Session.Service, { + get: (id: Parameters[0]) => + Effect.succeed({ + id, + slug: "output-ref", + projectID, + directory: process.cwd(), + parentID: undefined, + title: "Output ref", + version: "test", + time: { created: 0, updated: 0 }, + model: { providerID: ProviderV2.ID.make("test"), id: ModelV2.ID.make("test-model") }, + } satisfies Session.Info), + }), + ), +) + +function toolContext() { + return { + sessionID: SessionID.make("ses_ref_parent"), + messageID: MessageID.ascending(), + agent: "build", + abort: new AbortController().signal, + messages: [], + metadata: () => Effect.void, + ask: () => Effect.void, + } satisfies Tool.Context +} + +describe("workflow tool result — file-ref view (Train B, B-p3)", () => { + runtime.effect("result returns content_ref + summary + path for a captured file ref", () => + Effect.gen(function* () { + const info = yield* WorkflowTool + const workflow = yield* info.init() + const result = yield* workflow.execute( + { params: { action: "result", workflow_id: Dag.ID.make("dag_output_ref"), node_id: Dag.NodeID.make("node_ref") } }, + toolContext(), + ) + const parsed = JSON.parse(result.output) + expect(parsed).toEqual(expect.objectContaining({ + workflow_id: "dag_output_ref", + node_id: "node_ref", + status: "completed", + content_ref: refPath, + path: refPath, + summary: "Report summary captured at submit time", + size: 4096, + sha256: "0".repeat(64), + truncated: false, + next_cursor: null, + })) + expect(parsed.content).toBeUndefined() + }), + ) + + runtime.effect("result keeps the paged inline read verbatim for legacy outputs", () => + Effect.gen(function* () { + const info = yield* WorkflowTool + const workflow = yield* info.init() + const result = yield* workflow.execute( + { params: { action: "result", workflow_id: Dag.ID.make("dag_output_ref"), node_id: Dag.NodeID.make("node_inline") } }, + toolContext(), + ) + const parsed = JSON.parse(result.output) + expect(parsed).toEqual(expect.objectContaining({ + workflow_id: "dag_output_ref", + node_id: "node_inline", + status: "completed", + content: "plain inline output", + truncated: false, + next_cursor: null, + })) + expect(parsed.content_ref).toBeUndefined() + }), + ) +}) diff --git a/packages/opencode/test/dag/dag-output-ref.test.ts b/packages/opencode/test/dag/dag-output-ref.test.ts new file mode 100644 index 0000000000..ece144f13b --- /dev/null +++ b/packages/opencode/test/dag/dag-output-ref.test.ts @@ -0,0 +1,391 @@ +// oxlint-disable typescript-eslint/no-unsafe-type-assertion -- spawn-level +// probes deliberately mirror dag-structured-output.test.ts: mocked service +// layers and row fixtures use `as never` type shims (mock objects implement +// only the interface slice the scenario exercises). The shims are type-only; +// converting them would fork the template's shape without changing behavior. +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +/** + * Train B probes — node-output file references (workflows/dag-engine-optimization.md, + * v1.0.15 ledger, decisions B1–B4). + * + * Dual track (B1): report nodes without an output_schema may submit an absolute + * file path instead of inlining long text; output_schema nodes keep inline JSON. + * At submit time the runtime validates existence + size>0 and records + * `{content_ref, size, sha256}` in captured_output (B2); on first capture into + * the project `.opencode/` report area the runtime ensures the `.gitignore` + * entry exists (B4, append, never overwrite). The report area convention is + * `.opencode/workflow-reports/` — no better-fit existing directory was found + * under `.opencode/` (checked the worktree; naming mirrors the existing + * `.opencode/workflow-drafts/` convention). + * + * NOTE: the run's evidence.md was not present in the worktree, the config + * workflow repo, or the opencode data dir when this train started — the probe + * contract below is derived directly from the settled ledger above. + * + * Probe map: + * - B-p1: submit-time absolute-path detection — a child reply that IS an + * absolute path to an existing non-empty regular file captures a file_ref + * record ({kind, content_ref, path, size, sha256, summary}) into + * captured_output while nodeCompleted keeps the path string as the inline + * output (backward compatible: input mapping, wake digests, and legacy + * readers all see the same string they always did). RED on the unmodified + * engine: non-schema nodes never write captured_output. + * - B-p2 (PIN): output_schema settlement is untouched — a captured payload + * containing an absolute-path string stays inline JSON, no file_ref rewrite. + * Green before AND after the feature. + * - B-p4: `.gitignore` auto-entry — a capture whose ref path lies inside + * /.opencode/workflow-reports/ appends the report-area entry to + * /.gitignore exactly once, preserving pre-existing entries. RED + * on the unmodified engine: nothing touches the project `.gitignore`. + */ +import { afterAll, describe, expect, it } from "bun:test" +import { createHash } from "node:crypto" +import fs from "node:fs/promises" +import os from "node:os" +import path from "node:path" +import { Effect, Fiber, Layer, Semaphore } from "effect" +import type { SessionV1 } from "@opencode-ai/core/v1/session" +import { SessionPrompt } from "@/session/prompt" +import { MessageID } from "@/session/schema" +import { Dag } from "@/dag/dag" +import { Agent } from "@/agent/agent" +import { Session } from "@/session/session" +import { spawnNode, type NodeSpawnInput } from "@/dag/runtime/spawn" +import { registerCaptureSlot, validatePayload } from "@/dag/runtime/capture" +import { captureOutputFileRef, ensureReportAreaGitignore, REPORT_AREA } from "@/dag/runtime/output-ref" +import { makeNodeRow } from "./fixtures" +import type { DagStore } from "@opencode-ai/core/dag/store" + +const tmpRoots: string[] = [] + +function tmpRoot(prefix: string) { + return fs.mkdtemp(path.join(os.tmpdir(), prefix)).then((dir) => { + tmpRoots.push(dir) + return dir + }) +} + +afterAll(async () => { + for (const dir of tmpRoots) await fs.rm(dir, { recursive: true, force: true }) +}) + +type TrackedEvent = { type: string; nodeID: string; output?: unknown; reason?: string; trigger?: string } + +let capturedStore: Map = new Map() +let capturedCalls: unknown[] = [] + +function makeEventTracker() { + const events: TrackedEvent[] = [] + capturedStore = new Map() + capturedCalls = [] + const storeStub: Partial = { + tryClaimAdoption: () => Effect.succeed(true), + getNode: Effect.fn("s")((_workflowID: string, nodeID: string) => + Effect.sync(() => ({ ...makeNodeRow({ id: nodeID }), capturedOutput: capturedStore.get(nodeID) }))), + setCapturedOutput: Effect.fn("s")((_childSessionID: string, payload: unknown) => + Effect.sync(() => { + capturedCalls.push(payload) + capturedStore.set("node-1", payload) + })), + } + const dagLayer = Layer.mock(Dag.Service, { + store: storeStub as DagStore.Interface, + nodeQueued: Effect.fn("s")((_dagID: string, _nodeID: string) => Effect.void), + nodeStarted: Effect.fn("s")((_dagID: string, _nodeID: string) => Effect.void), + nodeCompleted: Effect.fn("s")((_dagID: string, nodeID: string, output: unknown) => + Effect.sync(() => events.push({ type: "nodeCompleted", nodeID, output }))), + nodeFailed: Effect.fn("s")((_dagID: string, nodeID: string, reason: string, trigger?: string) => + Effect.sync(() => events.push({ type: "nodeFailed", nodeID, reason, trigger }))), + nodeSkipped: Effect.fn("s")((_dagID: string, nodeID: string) => + Effect.sync(() => events.push({ type: "nodeSkipped", nodeID }))), + }) + return { events, dagLayer } +} + +const agentLayer = Layer.mock(Agent.Service, { + get: () => Effect.succeed({ + name: "build", mode: "all", permission: [], options: {}, description: "", prompt: "", + model: { providerID: "test" as never, modelID: "test-model" as never }, + tools: {}, hooks: {}, + }), + list: () => Effect.succeed([]), + defaultAgent: () => Effect.succeed("build"), +}) + +const sessionLayer = Layer.mock(Session.Service, { + get: () => Effect.succeed({ id: "ses_parent" as never, permission: [], agent: "build" } as never), + create: () => Effect.succeed({ id: "ses_child" as never } as never), + list: () => Effect.succeed([]), + messages: () => Effect.succeed([]), +}) + +function reply(text: string): SessionV1.WithParts { + return { + info: { + id: MessageID.ascending(), role: "assistant", parentID: MessageID.ascending(), + sessionID: "ses_child" as never, mode: "build", agent: "build", cost: 0, + path: { cwd: "/tmp", root: "/tmp" }, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + modelID: "test-model" as never, providerID: "test" as never, + time: { created: Date.now() }, finish: "stop", + }, + parts: text ? [{ type: "text", text }] as never : [], + } +} + +function makePromptLayer(result: SessionV1.WithParts): Layer.Layer { + return Layer.mock(SessionPrompt.Service, { + prompt: () => Effect.succeed(result), + }) +} + +function makeSpawnInput( + outputSchema?: Record, + overrides: Partial = {}, +): NodeSpawnInput { + return { + dagID: "wf-1", nodeID: "node-1", node: makeNodeRow(), + parentSessionID: "ses_parent", + promptParts: [{ type: "text", text: "do the thing" }], + outputSchema, + ...overrides, + } +} + +async function runSpawn( + dagLayer: Layer.Layer, + promptLayer: Layer.Layer, + outputSchema?: Record, + overrides: Partial = {}, +) { + const semaphore = Semaphore.makeUnsafe(1) + const fullLayer = Layer.mergeAll(dagLayer, agentLayer, sessionLayer, promptLayer) + await Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const result = yield* spawnNode(semaphore, makeSpawnInput(outputSchema, overrides)) + yield* Fiber.await(result.fiber) + }), + ).pipe(Effect.provide(fullLayer)) as Effect.Effect, + ) +} + +// Train A cast idiom: the extra `directory` field is cast away pre-feature so +// the probe compiles against the baseline NodeSpawnInput while offering the +// post-feature seam (gitignore guarantee keyed on the workflow directory). +const directoryOverride = (directory: string) => + ({ directory }) as unknown as Partial + +describe("submit-time absolute-path capture (Train B, B-p1)", () => { + it("B-p1(a) captures {content_ref, size, sha256} when the reply IS an existing non-empty absolute path", async () => { + const dir = await tmpRoot("dag-ref-") + const content = `${"report line\n".repeat(40)}SENTINEL` + const reportPath = path.join(dir, "report.md") + await Bun.write(reportPath, content) + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(reportPath))) + const completed = events.find((event) => event.type === "nodeCompleted") + expect(completed).toBeDefined() + expect(completed!.output).toBe(reportPath) + expect(capturedCalls).toHaveLength(1) + expect(capturedCalls[0]).toEqual({ + kind: "file_ref", + content_ref: reportPath, + path: reportPath, + size: Buffer.byteLength(content), + sha256: createHash("sha256").update(content).digest("hex"), + summary: `${content.slice(0, 200)}\u2026`, + }) + }) + + it("B-p1(a2) keeps the full text as summary when the file is at most 200 chars", async () => { + const dir = await tmpRoot("dag-ref-") + const content = "short report" + const reportPath = path.join(dir, "short.md") + await Bun.write(reportPath, content) + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(reportPath))) + expect(events.find((event) => event.type === "nodeCompleted")?.output).toBe(reportPath) + expect(capturedCalls).toHaveLength(1) + expect(capturedCalls[0]).toEqual(expect.objectContaining({ summary: "short report", size: Buffer.byteLength(content) })) + }) + + it("B-p1(b) leaves output inline when the absolute path does not exist", async () => { + const ghostPath = path.join(os.tmpdir(), `dag-ref-ghost-${Date.now()}.md`) + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(ghostPath))) + expect(events.find((event) => event.type === "nodeCompleted")?.output).toBe(ghostPath) + expect(capturedCalls).toHaveLength(0) + }) + + it("B-p1(c) leaves output inline when the file exists but is empty", async () => { + const dir = await tmpRoot("dag-ref-") + const emptyPath = path.join(dir, "empty.md") + await Bun.write(emptyPath, "") + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(emptyPath))) + expect(events.find((event) => event.type === "nodeCompleted")?.output).toBe(emptyPath) + expect(capturedCalls).toHaveLength(0) + }) + + it("B-p1(d) leaves output inline when the reply mentions a path inside prose", async () => { + const dir = await tmpRoot("dag-ref-") + const reportPath = path.join(dir, "report.md") + await Bun.write(reportPath, "prose-embedded report") + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(`Report written to ${reportPath}`))) + expect(events.find((event) => event.type === "nodeCompleted")?.output).toBe(`Report written to ${reportPath}`) + expect(capturedCalls).toHaveLength(0) + }) + + it("B-p1(e) leaves output inline for a plain-text reply (baseline behavior)", async () => { + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply("Task completed"))) + expect(events.find((event) => event.type === "nodeCompleted")?.output).toBe("Task completed") + expect(capturedCalls).toHaveLength(0) + }) + + it("B-p1(f) leaves output inline when the reply is a directory path", async () => { + const dir = await tmpRoot("dag-ref-") + const subDir = path.join(dir, "subdir") + await fs.mkdir(subDir) + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(subDir))) + expect(events.find((event) => event.type === "nodeCompleted")?.output).toBe(subDir) + expect(capturedCalls).toHaveLength(0) + }) +}) + +describe("output_schema dual track (Train B, B-p2 pin)", () => { + it("B-p2 keeps a captured payload containing an absolute path as inline JSON (no file_ref rewrite)", async () => { + const dir = await tmpRoot("dag-ref-") + const reportPath = path.join(dir, "schema-report.md") + await Bun.write(reportPath, "schema node report") + const { events, dagLayer } = makeEventTracker() + const schema = { type: "object" as const, required: ["report"] } + const payload = { report: reportPath } + const promptLayer = Layer.mock(SessionPrompt.Service, { + prompt: () => Effect.gen(function* () { + registerCaptureSlot("ses_child", schema) + const result = validatePayload("ses_child", payload) + if (result.ok) capturedStore.set("node-1", payload) + return reply("ignored text") + }), + }) + await runSpawn(dagLayer, promptLayer, schema) + const completed = events.find((event) => event.type === "nodeCompleted") + expect(completed).toBeDefined() + expect(completed!.output).toEqual(payload) + }) +}) + +describe("report-area gitignore entry (Train B, B-p4)", () => { + it("B-p4(a) appends the report-area entry to the project .gitignore on first capture into the report area", async () => { + const projectDir = await tmpRoot("dag-ref-proj-") + await Bun.write(path.join(projectDir, "existing-code.ts"), "export const x = 1\n") + const reportArea = path.join(projectDir, ".opencode", "workflow-reports") + await fs.mkdir(reportArea, { recursive: true }) + const reportPath = path.join(reportArea, "run-1.md") + await Bun.write(reportPath, "# report\n") + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(reportPath)), undefined, directoryOverride(projectDir)) + expect(events.find((event) => event.type === "nodeCompleted")?.output).toBe(reportPath) + const gitignore = await fs.readFile(path.join(projectDir, ".gitignore"), "utf8") + expect(gitignore.split("\n").map((line) => line.trim())).toContain(".opencode/workflow-reports/") + }) + + it("B-p4(b) is append-only and idempotent: pre-existing entries survive, the report entry lands exactly once", async () => { + const projectDir = await tmpRoot("dag-ref-proj-") + await Bun.write(path.join(projectDir, ".gitignore"), "node_modules\n*.log\n") + const reportArea = path.join(projectDir, ".opencode", "workflow-reports") + await fs.mkdir(reportArea, { recursive: true }) + const reportPath = path.join(reportArea, "run-2.md") + await Bun.write(reportPath, "# report 2\n") + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(reportPath)), undefined, directoryOverride(projectDir)) + expect(events.find((event) => event.type === "nodeCompleted")).toBeDefined() + const { events: secondEvents, dagLayer: secondDagLayer } = makeEventTracker() + await runSpawn(secondDagLayer, makePromptLayer(reply(reportPath)), undefined, directoryOverride(projectDir)) + expect(secondEvents.find((event) => event.type === "nodeCompleted")).toBeDefined() + const gitignore = await fs.readFile(path.join(projectDir, ".gitignore"), "utf8") + const lines = gitignore.split("\n").map((line) => line.trim()).filter((line) => line.length > 0) + expect(lines).toContain("node_modules") + expect(lines).toContain("*.log") + expect(lines.filter((line) => line === ".opencode/workflow-reports/")).toHaveLength(1) + }) + + it("B-p4(c) does not touch the project .gitignore for refs outside the report area", async () => { + const projectDir = await tmpRoot("dag-ref-proj-") + await Bun.write(path.join(projectDir, ".gitignore"), "node_modules\n") + const outsideDir = await tmpRoot("dag-ref-outside-") + const reportPath = path.join(outsideDir, "elsewhere.md") + await Bun.write(reportPath, "# elsewhere\n") + const { events, dagLayer } = makeEventTracker() + await runSpawn(dagLayer, makePromptLayer(reply(reportPath)), undefined, directoryOverride(projectDir)) + expect(events.find((event) => event.type === "nodeCompleted")?.output).toBe(reportPath) + expect(await fs.readFile(path.join(projectDir, ".gitignore"), "utf8")).toBe("node_modules\n") + }) +}) + +describe("output-ref module rules (Train B, post-feature units)", () => { + it("rejects relative paths and paths containing whitespace even when the file exists", async () => { + const dir = await tmpRoot("dag-ref-") + const spacedPath = path.join(dir, "two words.md") + await Bun.write(spacedPath, "spaced") + expect(await Effect.runPromise(captureOutputFileRef(` ${spacedPath} `))).toBeUndefined() + const relative = path.relative(process.cwd(), spacedPath) + expect(await Effect.runPromise(captureOutputFileRef(relative))).toBeUndefined() + }) + + it("captures cross-worktree refs and normalizes trailing whitespace only", async () => { + const dir = await tmpRoot("dag-ref-") + const reportPath = path.join(dir, "cross.md") + const content = "cross-worktree report" + await Bun.write(reportPath, content) + const ref = await Effect.runPromise(captureOutputFileRef(`\n${reportPath}\n`)) + expect(ref).toEqual(expect.objectContaining({ + kind: "file_ref", + content_ref: reportPath, + path: reportPath, + size: Buffer.byteLength(content), + sha256: createHash("sha256").update(content).digest("hex"), + summary: content, + })) + }) + + it("creates a missing .gitignore with only the report-area entry", async () => { + const projectDir = await tmpRoot("dag-ref-proj-") + const refPath = path.join(projectDir, REPORT_AREA, "x.md") + await Effect.runPromise(ensureReportAreaGitignore(projectDir, refPath)) + expect(await fs.readFile(path.join(projectDir, ".gitignore"), "utf8")).toBe(".opencode/workflow-reports/\n") + }) + + it("separates the entry onto its own line when the existing file lacks a trailing newline", async () => { + const projectDir = await tmpRoot("dag-ref-proj-") + await Bun.write(path.join(projectDir, ".gitignore"), "node_modules") + const refPath = path.join(projectDir, REPORT_AREA, "x.md") + await Effect.runPromise(ensureReportAreaGitignore(projectDir, refPath)) + expect(await fs.readFile(path.join(projectDir, ".gitignore"), "utf8")).toBe("node_modules\n.opencode/workflow-reports/\n") + }) + + it("leaves the file byte-identical when an entry or a covering .opencode/ rule already exists", async () => { + for (const covering of [".opencode/workflow-reports/", ".opencode/workflow-reports", ".opencode/", ".opencode"]) { + const projectDir = await tmpRoot("dag-ref-proj-") + const before = `${covering}\nkeep-me\n` + await Bun.write(path.join(projectDir, ".gitignore"), before) + const refPath = path.join(projectDir, REPORT_AREA, "x.md") + await Effect.runPromise(ensureReportAreaGitignore(projectDir, refPath)) + expect(await fs.readFile(path.join(projectDir, ".gitignore"), "utf8")).toBe(before) + } + }) + + it("does not create a .gitignore for refs outside the report area", async () => { + const projectDir = await tmpRoot("dag-ref-proj-") + const outsideDir = await tmpRoot("dag-ref-outside-") + const refPath = path.join(outsideDir, "y.md") + await Effect.runPromise(ensureReportAreaGitignore(projectDir, refPath)) + expect(await Bun.file(path.join(projectDir, ".gitignore")).exists()).toBe(false) + }) +}) diff --git a/packages/opencode/test/dag/dag-rev-view-status.test.ts b/packages/opencode/test/dag/dag-rev-view-status.test.ts new file mode 100644 index 0000000000..5f5da0b19b --- /dev/null +++ b/packages/opencode/test/dag/dag-rev-view-status.test.ts @@ -0,0 +1,205 @@ +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +/** + * Train A probe A-p1 (status-action seam) — workflows/dag-engine-optimization.md, + * v1.0.15 ledger. + * + * The workflow tool's `status` action is one of the view seams the ledger + * filters to the current revision (evidence §8: workflow.ts status action). + * The store mock offers BOTH reads: the legacy all-rows read and the + * current-revision read. Before the feature, status consumes the all-rows + * read and surfaces superseded replaced segments; after, it consumes the + * current-revision read only. + * + * RED on the unmodified engine: the superseded rows leak into the status + * output. + */ +import { describe, expect } from "bun:test" +import { Effect, Layer } from "effect" +import { DagStore } from "@opencode-ai/core/dag/store" +import { ProjectV2 } from "@opencode-ai/core/project" +import { ModelV2 } from "@opencode-ai/core/model" +import { ProviderV2 } from "@opencode-ai/core/provider" +import { Event } from "@opencode-ai/schema/event" +import { Agent } from "@/agent/agent" +import { Dag } from "@/dag/dag" +import { EventV2Bridge } from "@/event-v2-bridge" +import { Question } from "@/question" +import { Session } from "@/session/session" +import { MessageID, SessionID } from "@/session/schema" +import { Skill } from "@/skill" +import type { Tool } from "@/tool/tool" +import { Truncate } from "@/tool/truncate" +import { WorkflowTool } from "@/tool/workflow" +import { Provider } from "@/provider/provider" +import { testEffect } from "../lib/effect" + +const projectID = ProjectV2.ID.make("project_test") + +const currentRow = { + id: "current_a", + workflowId: "dag_rev_status", + name: "Current node", + workerType: "build", + status: "completed", + required: true, + dependsOn: [], + modelId: null, + modelProviderId: null, + childSessionId: null, + output: "kept", + capturedOutput: null, + errorReason: null, + errorClass: null, + deadlineMs: null, + wakeEligible: false, + wakeReported: false, + replanAttempts: 0, + timeoutExtensions: 0, + escalationPending: false, + superseded: false, + seq: 1, + startedAt: 1, + completedAt: 2, + timeCreated: 1, + timeUpdated: 2, +} +const supersededFail = { + ...currentRow, + id: "superseded_c", + name: "Replaced failure", + status: "failed", + errorReason: "simulated exec failure", + errorClass: "exec_failed", + output: null, + startedAt: null, + completedAt: 3, + seq: 2, +} +const supersededCancelled = { + ...currentRow, + id: "superseded_d", + name: "Replaced cancelled", + status: "failed", + errorReason: "cancelled via replan", + output: null, + startedAt: null, + completedAt: 3, + seq: 3, +} +const allRows = [currentRow, supersededFail, supersededCancelled] +const currentRows = [currentRow] + +// The extra getCurrentNodes property is cast away pre-feature: the probe must +// compile against the baseline interface while offering the post-feature seam. +const storeImpl = { + getWorkflow: (id: string) => + Effect.succeed( + id === "dag_rev_status" + ? { + id, + projectId: projectID, + sessionId: "ses_rev_parent", + directory: null, + title: "Rev status", + status: "running", + config: "{}", + seq: 3, + wakeReported: false, + graphRev: 2, + startedAt: 1, + completedAt: null, + timeCreated: 1, + timeUpdated: 2, + } + : undefined, + ), + getNodes: () => Effect.succeed(allRows), + getCurrentNodes: () => Effect.succeed(currentRows), + getNode: (workflowID: string, nodeID: string) => + Effect.succeed(allRows.find((node) => node.workflowId === workflowID && node.id === nodeID)), +} + +const dag = Dag.layer.pipe( + Layer.provide(Layer.mock(DagStore.Service, storeImpl)), + Layer.provide(Layer.mock(EventV2Bridge.Service, { + publish: (definition, data) => + Effect.succeed({ id: Event.ID.create(), type: definition.type, data }), + })), +) + +const runtime = testEffect( + Layer.mergeAll( + Layer.mock(Agent.Service, { + get: () => + Effect.succeed({ + name: "build", + mode: "all", + permission: [], + options: {}, + description: "", + prompt: "", + model: { providerID: ProviderV2.ID.make("test"), modelID: ModelV2.ID.make("test-model") }, + tools: {}, + hooks: {}, + }), + list: () => Effect.succeed([]), + }), + Layer.mock(Skill.Service, { all: () => Effect.succeed([]) }), + Layer.mock(Truncate.Service, { + output: (content) => Effect.succeed({ content, truncated: false }), + }), + Layer.mock(Question.Service, { ask: () => Effect.succeed([]) }), + Layer.mock(Provider.Service, { + list: () => Effect.succeed({}), + getModel: (providerID, modelID) => Effect.fail(new Provider.ModelNotFoundError({ providerID, modelID })), + }), + dag, + Layer.mock(Session.Service, { + get: (id: Parameters[0]) => + Effect.succeed({ + id, + slug: "rev-status", + projectID, + directory: process.cwd(), + parentID: undefined, + title: "Rev status", + version: "test", + time: { created: 0, updated: 0 }, + model: { providerID: ProviderV2.ID.make("test"), id: ModelV2.ID.make("test-model") }, + } satisfies Session.Info), + }), + ), +) + +function toolContext() { + return { + sessionID: SessionID.make("ses_rev_parent"), + messageID: MessageID.ascending(), + agent: "build", + abort: new AbortController().signal, + messages: [], + metadata: () => Effect.void, + ask: () => Effect.void, + } satisfies Tool.Context +} + +describe("workflow tool status — current-revision filtering (Train A, A-p1 status seam)", () => { + runtime.effect("status exposes only the current revision; replaced segments stay hidden", () => + Effect.gen(function* () { + const info = yield* WorkflowTool + const workflow = yield* info.init() + const result = yield* workflow.execute( + { params: { action: "status", workflow_id: Dag.ID.make("dag_rev_status") } }, + toolContext(), + ) + expect(result.output).toContain("current_a") + expect(result.output).not.toContain("superseded_c") + expect(result.output).not.toContain("superseded_d") + const ids = [...result.output.matchAll(/"id": "(\w+)"/g)].map((match) => match[1]) + expect(ids).toContain("dag_rev_status") + expect(ids.filter((id) => id.startsWith("current_") || id.startsWith("superseded_"))).toEqual(["current_a"]) + }), + ) +}) diff --git a/packages/opencode/test/dag/dag-rev-view.test.ts b/packages/opencode/test/dag/dag-rev-view.test.ts new file mode 100644 index 0000000000..34867d29ff --- /dev/null +++ b/packages/opencode/test/dag/dag-rev-view.test.ts @@ -0,0 +1,495 @@ +// SPDX-FileCopyrightText: 2026 LeXwDeX +// SPDX-License-Identifier: AGPL-3.0-or-later + +/** + * Train A probes — graph-revision VIEW semantics (workflows/dag-engine-optimization.md, + * v1.0.15 ledger, decisions A1–A6). + * + * Durable data is untouched by this feature: replaced segments stay in the + * read model and in EventV2 history; only the CURRENT revision is exposed to + * view and terminal-aggregation seams (summaries, status, loop rebuilds, wake + * attribution). Replan/cancel engine semantics are unchanged (A5/A6). + * + * Probe map (evidence §9): + * - A-p1/A-p3: a replan that bypasses a failed required node hides the + * replaced segment from summaries/status/rebuild input, and the workflow + * COMPLETES via the new path instead of failing on the old failure + * (the wake-up bug: replaced fails re-seeded as required-unsatisfied on + * every rebuild). + * - A-p2 (PIN): a genuine current-rev required failure stays visible, is + * counted, and fails the workflow — before AND after the feature. + * - A-p5(a): the wake digest's failure attribution excludes superseded + * replaced failures (only current-rev failures are attributed). + */ +import { describe, expect, it } from "bun:test" +import { Deferred, Effect, Layer, Option, Queue } from "effect" +import type { SessionV1 } from "@opencode-ai/core/v1/session" +import { Database } from "@opencode-ai/core/database/database" +import { DagProjector } from "@opencode-ai/core/dag/projector" +import { DagStore } from "@opencode-ai/core/dag/store" +import { EventV2 } from "@opencode-ai/core/event" +import { ProjectV2 } from "@opencode-ai/core/project" +import { ProjectTable } from "@opencode-ai/core/project/sql" +import { AbsolutePath } from "@opencode-ai/core/schema" +import { SessionSchema } from "@opencode-ai/core/session/schema" +import { SessionTable } from "@opencode-ai/core/session/sql" +import { Model } from "@opencode-ai/schema/model" +import { Provider } from "@opencode-ai/schema/provider" +import { Agent } from "@/agent/agent" +import { Dag, type NodeConfig } from "@/dag/dag" +import { DagLoop } from "@/dag/runtime/loop" +import { InstanceRef } from "@/effect/instance-ref" +import { EventV2Bridge } from "@/event-v2-bridge" +import { SessionPrompt } from "@/session/prompt" +import { MessageID, PartID, SessionID } from "@/session/schema" +import { Session } from "@/session/session" +import { SessionStatus } from "@/session/status" +import { pollWithTimeout } from "../lib/effect" +import { withIdleAdmission } from "../lib/session-prompt" + +/** Child release: reply text (node completes) or a simulated exec failure. */ +type ChildOutcome = { text: string } | { fail: true } + +interface PromptGate { + readonly title: string + readonly release: Deferred.Deferred +} + +interface ParentPromptGate { + readonly text: string + readonly release: Deferred.Deferred<"success" | "failure"> +} + +function takeWithin(queue: Queue.Queue, message: string) { + return Queue.take(queue).pipe( + Effect.timeoutOption("2 seconds"), + Effect.flatMap(Option.match({ + onNone: () => Effect.fail(new Error(message)), + onSome: Effect.succeed, + })), + ) +} + +function reply(sessionID: string, text: string): SessionV1.WithParts { + const sid = SessionID.make(sessionID) + const messageID = MessageID.ascending() + const part: SessionV1.TextPart = { + id: PartID.ascending(), + sessionID: sid, + messageID, + type: "text", + text, + } + return { + info: { + id: messageID, + role: "assistant", + parentID: MessageID.ascending(), + sessionID: sid, + mode: "build", + agent: "build", + cost: 0, + path: { cwd: process.cwd(), root: process.cwd() }, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + modelID: Model.ID.make("test-model"), + providerID: Provider.ID.make("test"), + time: { created: Date.now() }, + finish: "stop", + }, + parts: text ? [part] : [], + } +} + +function node(id: string, dependsOn: string[] = []): NodeConfig { + return { + id, + name: id, + worker_type: "build", + depends_on: dependsOn, + required: true, + prompt_template: { inline: id }, + report_to_parent: true, + } +} + +function loopLayer(input: { + readonly childPrompts: Queue.Queue + readonly parentPrompts: Queue.Queue +}) { + const database = Database.layerFromPath(":memory:") + const events = EventV2.layer.pipe(Layer.provide(database)) + const bridge = EventV2Bridge.layer.pipe(Layer.provide(events)) + const store = DagStore.layer.pipe(Layer.provide(database)) + const status = SessionStatus.layer.pipe(Layer.provide(bridge)) + const projector = DagProjector.layer.pipe( + Layer.provide(events), + Layer.provide(database), + ) + const dag = Dag.layer.pipe( + Layer.provide(bridge), + Layer.provide(store), + ) + const base = Layer.mergeAll(database, events, bridge, store, projector, dag, status) + const childTitles = new Map() + const created: string[] = [] + const projectID = ProjectV2.ID.make("project-1") + const sessionInfo = (id: SessionSchema.ID, title: string): Session.Info => ({ + id, + slug: title, + projectID, + directory: process.cwd(), + title, + version: "test", + time: { created: 0, updated: 0 }, + }) + const session = Layer.mock(Session.Service, { + get: (id) => Effect.succeed(sessionInfo(id, "Parent")), + create: (value) => + Effect.sync(() => { + const id = SessionID.make(`ses_child_${created.length + 1}`) + created.push(id) + childTitles.set(id, (value?.title ?? id).replace(" (DAG node)", "")) + return sessionInfo(id, value?.title ?? id) + }), + messages: () => Effect.succeed([]), + }) + const deliver = Effect.fn("test.SessionPrompt.deliver")(function* (value: SessionPrompt.PromptInput) { + const sessionID = value.sessionID as string + const text = value.parts + .map((part) => (part.type === "text" ? (part as { text: string }).text : "")) + .join("\n") + if (sessionID === "ses_parent") { + const release = yield* Deferred.make<"success" | "failure">() + yield* Queue.offer(input.parentPrompts, { text, release }) + const outcome = yield* Deferred.await(release) + if (outcome === "failure") return yield* Effect.die(new Error("provider unavailable")) + return reply(sessionID, "parent handled wake") + } + const release = yield* Deferred.make() + yield* Queue.offer(input.childPrompts, { + title: childTitles.get(sessionID) ?? sessionID, + release, + }) + const outcome = yield* Deferred.await(release) + // Simulated execution failure: the spawn fiber's catchCause publishes a + // durable NodeFailed with trigger exec_failed — a GENUINE failure, the + // same class of failure the ledger's replaced-segment scenarios carry. + if ("fail" in outcome) return yield* Effect.die(new Error("simulated exec failure")) + return reply(sessionID, outcome.text) + }) + const prompt = Layer.mock(SessionPrompt.Service, withIdleAdmission({ + cancel: () => Effect.void, + prompt: deliver, + promptIfIdle: (value) => deliver(value).pipe(Effect.map(Option.some)), + })) + const agent = Layer.mock(Agent.Service, { + get: () => Effect.succeed({ + name: "build", + mode: "all", + permission: [], + options: {}, + description: "", + prompt: "", + model: { providerID: Provider.ID.make("test"), modelID: Model.ID.make("test-model") }, + tools: {}, + hooks: {}, + }), + }) + const loop = DagLoop.layer.pipe( + Layer.provide(base), + Layer.provide(session), + Layer.provide(prompt), + Layer.provide(agent), + ) + return Layer.merge(base, loop) +} + +function runLoopTest( + test: (services: { + readonly dag: Dag.Interface + readonly store: DagStore.Interface + readonly childPrompts: Queue.Queue + readonly parentPrompts: Queue.Queue + }) => Effect.Effect, +) { + return Effect.gen(function* () { + const childPrompts = yield* Queue.unbounded() + const parentPrompts = yield* Queue.unbounded() + return yield* Effect.gen(function* () { + const dag = yield* Dag.Service + const loop = yield* DagLoop.Service + const store = yield* DagStore.Service + const database = yield* Database.Service + const projectID = ProjectV2.ID.make("project-1") + const parentSessionID = SessionID.make("ses_parent") + yield* database.db.insert(ProjectTable).values({ + id: projectID, + worktree: AbsolutePath.make(process.cwd()), + sandboxes: [], + }).run().pipe(Effect.orDie) + yield* database.db.insert(SessionTable).values({ + id: parentSessionID, + project_id: projectID, + slug: "parent", + directory: process.cwd(), + title: "Parent", + version: "test", + }).run().pipe(Effect.orDie) + yield* loop.init() + return yield* test({ dag, store, childPrompts, parentPrompts }) + }).pipe( + Effect.provide(loopLayer({ childPrompts, parentPrompts })), + Effect.provideService(InstanceRef, { + directory: process.cwd(), + worktree: process.cwd(), + project: { + id: ProjectV2.ID.make("project-1"), + worktree: process.cwd(), + time: { created: 0, updated: 0 }, + sandboxes: [], + }, + }), + Effect.scoped, + ) + }) +} + +describe("Train A rev-view (durable data untouched, view = current revision only)", () => { + // A-p1 + A-p3: the ledger wake-up scenario. A→B→C→D all required plus an + // independent Z (its running state keeps the workflow from going terminal + // the moment C fails — the orchestrator gets its chance to replan, exactly + // like a real wake-driven replan). C fails genuinely; the agent bypasses it + // with a new E→G→H suffix and drops D. After the replan, the CURRENT + // revision is {A,B,E,G,H,Z}: the view seams expose only it, and completing + // the new path COMPLETES the workflow — the replaced C/D failures must not + // re-seed as required-unsatisfied and fail it. + // + // RED on the unmodified engine: getWorkflowSummaries counts C+D (nodeCount + // 8, failedNodes 2) and the WorkflowReplanned rebuild seeds them as + // required-unsatisfied, so the workflow FAILS as soon as the runtime is + // complete instead of completing. + it("A-p1/A-p3: replan bypassing a failed node hides the replaced segment and the workflow completes", async () => { + await Effect.runPromise( + runLoopTest(({ dag, store, childPrompts }) => + Effect.gen(function* () { + const dagID = yield* dag.create({ + projectID: "project-1", + sessionID: "ses_parent", + title: "Rev view", + config: { + name: "rev-view", + nodes: [node("a"), node("b", ["a"]), node("c", ["b"]), node("d", ["c"]), node("z")], + }, + }) + const first = yield* takeWithin(childPrompts, "a/z did not start") + const second = yield* takeWithin(childPrompts, "a/z did not start") + const gateA = first.title === "a" ? first : second + const gateZ = first.title === "a" ? second : first + expect([gateA.title, gateZ.title].sort()).toEqual(["a", "z"]) + + // Run A then B; hold Z for the whole scenario. + yield* Deferred.succeed(gateA.release, { text: "a done" }) + const gateB = yield* takeWithin(childPrompts, "b did not start") + yield* Deferred.succeed(gateB.release, { text: "b done" }) + + // C fails genuinely (exec_failed) while Z is still running, so the + // workflow stays RUNNING and the orchestrator can replan. + const gateC = yield* takeWithin(childPrompts, "c did not start") + yield* Deferred.succeed(gateC.release, { fail: true }) + yield* pollWithTimeout( + store.getNode(dagID, "c").pipe( + Effect.map((n) => (n?.status === "failed" && n.errorClass === "exec_failed" ? n : undefined)), + ), + "c did not fail with exec_failed", + ) + + // Bypass C: new suffix E→G→H off B; D (pending) is dropped by the + // fragment and cancels; C (terminal failed, absent from fragment) is + // the replaced segment the view must hide. + const plan = yield* dag.replan(dagID, { + nodes: [node("e", ["b"]), node("g", ["e"]), node("h", ["g"])], + }) + expect(plan.cancel).toEqual(["d"]) + expect(plan.add.sort()).toEqual(["e", "g", "h"]) + + // A-p1 VIEW seams, immediately after the replan: the summary counts + // ONLY the current revision. The replaced C+D rows must not count. + const summary = (yield* store.getWorkflowSummaries("ses_parent")).find((s) => s.id === dagID) + expect(summary).toBeDefined() + expect(summary!.nodeCount).toBe(6) + expect(summary!.completedNodes).toBe(2) + expect(summary!.failedNodes).toBe(0) + + // DURABLE TRUTH UNTOUCHED (A1/A2): the full read still carries all + // eight rows with their outcomes intact — replaced segments survive + // in the read model and remain resolvable via the result store. + const all = yield* store.getNodes(dagID) + expect(all.map((n) => n.id).sort()).toEqual(["a", "b", "c", "d", "e", "g", "h", "z"]) + const c = all.find((n) => n.id === "c")! + expect(c.status).toBe("failed") + expect(c.errorClass).toBe("exec_failed") + const d = all.find((n) => n.id === "d")! + expect(d.status).toBe("failed") + expect(d.errorReason).toBe("cancelled via replan") + + // MARKER level: the replaced segment carries the superseded flag and + // the workflow bumped its graph revision — while current-rev rows + // (completed A, B; added E) stay unmarked. + expect((yield* store.getNode(dagID, "c"))?.superseded).toBe(true) + expect((yield* store.getNode(dagID, "d"))?.superseded).toBe(true) + expect((yield* store.getNode(dagID, "a"))?.superseded).toBe(false) + expect((yield* store.getNode(dagID, "e"))?.superseded).toBe(false) + expect((yield* store.getWorkflow(dagID))?.graphRev).toBe(2) + + // Complete the new path (and Z). The workflow must COMPLETE — not + // fail on the replaced required failures. + const gateE = yield* takeWithin(childPrompts, "e did not start after replan") + yield* Deferred.succeed(gateE.release, { text: "e done" }) + const gateG = yield* takeWithin(childPrompts, "g did not start") + yield* Deferred.succeed(gateG.release, { text: "g done" }) + const gateH = yield* takeWithin(childPrompts, "h did not start") + yield* Deferred.succeed(gateH.release, { text: "h done" }) + yield* Deferred.succeed(gateZ.release, { text: "z done" }) + + const terminal = yield* pollWithTimeout( + store.getWorkflow(dagID).pipe( + Effect.map((wf) => (wf && (wf.status === "completed" || wf.status === "failed") ? wf : undefined)), + ), + "workflow did not reach terminal after the new path completed", + ) + expect(terminal.status).toBe("completed") + + const finalSummary = (yield* store.getWorkflowSummaries("ses_parent")).find((s) => s.id === dagID) + expect(finalSummary!.nodeCount).toBe(6) + expect(finalSummary!.completedNodes).toBe(6) + expect(finalSummary!.failedNodes).toBe(0) + }), + ), + ) + }) + + // A-p2 (PIN, green before and after): a genuine failure of the CURRENT + // revision stays visible, is counted in failedNodes, and fails the + // workflow. Regression guard for "current-rev true failures visible". + it("A-p2: a genuine current-rev required failure stays visible and fails the workflow", async () => { + await Effect.runPromise( + runLoopTest(({ dag, store, childPrompts }) => + Effect.gen(function* () { + const dagID = yield* dag.create({ + projectID: "project-1", + sessionID: "ses_parent", + title: "Current rev failure", + config: { name: "current-rev-failure", nodes: [node("g")] }, + }) + const gate = yield* takeWithin(childPrompts, "g did not start") + yield* Deferred.succeed(gate.release, { fail: true }) + + yield* pollWithTimeout( + store.getWorkflow(dagID).pipe( + Effect.map((wf) => (wf?.status === "failed" ? wf : undefined)), + ), + "workflow did not fail on the genuine required failure", + ) + + // Visible + counted: the failure belongs to the current revision. + const summary = (yield* store.getWorkflowSummaries("ses_parent")).find((s) => s.id === dagID) + expect(summary!.nodeCount).toBe(1) + expect(summary!.failedNodes).toBe(1) + const rows = yield* store.getNodes(dagID) + expect(rows.map((n) => n.id)).toEqual(["g"]) + expect(rows[0].status).toBe("failed") + expect(rows[0].errorClass).toBe("exec_failed") + }), + ), + ) + }) + + // A-p5(a): wake digest failure attribution is terminal aggregation — it + // must attribute only CURRENT-revision failures. The replaced failure C is + // superseded by the replan that adds G; when G fails genuinely and the + // workflow goes terminal, the digest's "Failed nodes:" block names G and + // NOT C. + // + // RED on the unmodified engine: the attribution reads every failed row + // with an error class, so the superseded C is attributed alongside G. + it("A-p5(a): wake attribution excludes superseded replaced failures", async () => { + await Effect.runPromise( + runLoopTest(({ dag, store, childPrompts, parentPrompts }) => + Effect.gen(function* () { + const dagID = yield* dag.create({ + projectID: "project-1", + sessionID: "ses_parent", + title: "Wake attribution", + config: { + name: "wake-attribution", + nodes: [node("a"), node("b", ["a"]), node("c", ["b"]), node("z")], + }, + }) + const first = yield* takeWithin(childPrompts, "a/z did not start") + const second = yield* takeWithin(childPrompts, "a/z did not start") + const gateA = first.title === "a" ? first : second + const gateZ = first.title === "a" ? second : first + + yield* Deferred.succeed(gateA.release, { text: "a done" }) + const gateB = yield* takeWithin(childPrompts, "b did not start") + yield* Deferred.succeed(gateB.release, { text: "b done" }) + const gateC = yield* takeWithin(childPrompts, "c did not start") + yield* Deferred.succeed(gateC.release, { fail: true }) + yield* pollWithTimeout( + store.getNode(dagID, "c").pipe( + Effect.map((n) => (n?.status === "failed" && n.errorClass === "exec_failed" ? n : undefined)), + ), + "c did not fail with exec_failed", + ) + + // Bypass the failure with G; C becomes the replaced segment. + const plan = yield* dag.replan(dagID, { nodes: [node("g", ["b"])] }) + expect(plan.add).toEqual(["g"]) + expect((yield* store.getNode(dagID, "c"))?.superseded).toBe(true) + expect((yield* store.getNode(dagID, "g"))?.superseded).toBe(false) + + // G fails genuinely — the current revision's true failure. + const gateG = yield* takeWithin(childPrompts, "g did not start after replan") + yield* Deferred.succeed(gateG.release, { fail: true }) + yield* pollWithTimeout( + store.getNode(dagID, "g").pipe( + Effect.map((n) => (n?.status === "failed" && n.errorClass === "exec_failed" ? n : undefined)), + ), + "g did not fail with exec_failed", + ) + + // Z completes; the required G failure then fails the workflow. + yield* Deferred.succeed(gateZ.release, { text: "z done" }) + yield* pollWithTimeout( + store.getWorkflow(dagID).pipe( + Effect.map((wf) => (wf?.status === "failed" ? wf : undefined)), + ), + "workflow did not fail on the current-rev required failure", + ) + + // Drain wakes until the terminal one arrives. Node-terminal and idle + // stimuli can deliver intermediate "actionable" wakes while the + // workflow is still running — those carry node results but no + // workflow-level failure attribution. The attribution block lives on + // the wake delivered once the workflow is terminally failed. + const wakes: string[] = [] + let terminalWake = "" + for (;;) { + const wake = yield* takeWithin(parentPrompts, "terminal wake did not reach the parent") + wakes.push(wake.text) + yield* Deferred.succeed(wake.release, "success") + if (wake.text.includes("[DAG Workflow failed]")) { + terminalWake = wake.text + break + } + } + // Attribution block format: `- "" (): `. + // The current-rev failure G is attributed; the superseded replaced + // failure C must not be attributed on ANY delivered wake. + expect(terminalWake).toContain("- \"g\" (exec_failed)") + for (const text of wakes) expect(text.includes("- \"c\"")).toBe(false) + }), + ), + ) + }) +}) diff --git a/packages/opencode/test/dag/dag-summary-publisher-behavior.test.ts b/packages/opencode/test/dag/dag-summary-publisher-behavior.test.ts index 83e70442b7..15da4955fa 100644 --- a/packages/opencode/test/dag/dag-summary-publisher-behavior.test.ts +++ b/packages/opencode/test/dag/dag-summary-publisher-behavior.test.ts @@ -57,6 +57,7 @@ function workflow(id: string, sessionId: string, projectId: string): WorkflowRow config: "", seq: 0, wakeReported: false, + graphRev: 1, startedAt: null, completedAt: null, timeCreated: 0, diff --git a/packages/opencode/test/dag/fixtures.ts b/packages/opencode/test/dag/fixtures.ts index e9ce2acd22..7fc681d304 100644 --- a/packages/opencode/test/dag/fixtures.ts +++ b/packages/opencode/test/dag/fixtures.ts @@ -22,6 +22,7 @@ export function makeNodeRow(overrides: Partial = {}): DagStore replanAttempts: 0, timeoutExtensions: 0, escalationPending: false, + superseded: false, seq: 0, startedAt: null, completedAt: null, diff --git a/packages/opencode/test/dag/workflow-tool.test.ts b/packages/opencode/test/dag/workflow-tool.test.ts index f9a7ce84d8..c2fef15d61 100644 --- a/packages/opencode/test/dag/workflow-tool.test.ts +++ b/packages/opencode/test/dag/workflow-tool.test.ts @@ -109,6 +109,99 @@ const resultNodes = [ output: "other", }), ] +const mockNodes = (id: string) => + id === "dag_status" + ? [ + { + id: "node_running", + workflowId: "dag_status", + name: "Running node", + workerType: "build", + status: "running", + required: true, + dependsOn: [], + modelId: null, + modelProviderId: null, + childSessionId: "ses_child", + output: null, + capturedOutput: null, + errorReason: null, + errorClass: null, + deadlineMs: null, + wakeEligible: true, + wakeReported: false, + replanAttempts: 0, + seq: 1, + timeoutExtensions: 0, + escalationPending: false, + superseded: false, + startedAt: 1, + completedAt: null, + timeCreated: 1, + timeUpdated: 2, + }, + { + id: "node_failed", + workflowId: "dag_status", + name: "Failed node", + workerType: "build", + status: "failed", + required: false, + dependsOn: ["node_running"], + modelId: null, + modelProviderId: null, + childSessionId: "ses_failed_child", + output: null, + capturedOutput: null, + errorReason: "node exceeded timeout of 600000ms", + errorClass: "timeout", + deadlineMs: null, + wakeEligible: false, + wakeReported: false, + replanAttempts: 0, + seq: 2, + timeoutExtensions: 0, + escalationPending: false, + superseded: false, + startedAt: 1, + completedAt: 2, + timeCreated: 1, + timeUpdated: 2, + }, + ] + : id === "dag_step" + ? [ + { + id: "node_ready", + workflowId: "dag_step", + name: "Ready node", + workerType: "build", + status: "pending", + required: true, + dependsOn: [], + modelId: null, + modelProviderId: null, + childSessionId: null, + output: null, + capturedOutput: null, + errorReason: null, + errorClass: null, + deadlineMs: null, + wakeEligible: false, + wakeReported: false, + replanAttempts: 0, + seq: 1, + timeoutExtensions: 0, + escalationPending: false, + superseded: false, + startedAt: null, + completedAt: null, + timeCreated: 1, + timeUpdated: 1, + }, + ] + : [] + const store = Layer.mock(DagStore.Service, { getWorkflow: (id: string) => Effect.succeed( @@ -123,6 +216,7 @@ const store = Layer.mock(DagStore.Service, { config: "{}", seq: 1, wakeReported: false, + graphRev: 1, startedAt: 1, completedAt: null, timeCreated: 1, @@ -139,6 +233,7 @@ const store = Layer.mock(DagStore.Service, { config: "{}", seq: 1, wakeReported: true, + graphRev: 1, startedAt: 1, completedAt: 2, timeCreated: 1, @@ -155,6 +250,7 @@ const store = Layer.mock(DagStore.Service, { config: "{}", seq: 1, wakeReported: false, + graphRev: 1, startedAt: 1, completedAt: null, timeCreated: 1, @@ -179,6 +275,7 @@ const store = Layer.mock(DagStore.Service, { }), seq: 1, wakeReported: false, + graphRev: 1, startedAt: 1, completedAt: null, timeCreated: 1, @@ -210,6 +307,7 @@ const store = Layer.mock(DagStore.Service, { }), seq: 1, wakeReported: false, + graphRev: 1, startedAt: 1, completedAt: null, timeCreated: 1, @@ -217,97 +315,10 @@ const store = Layer.mock(DagStore.Service, { } : undefined, ), - getNodes: (id: string) => - Effect.succeed( - id === "dag_status" - ? [ - { - id: "node_running", - workflowId: "dag_status", - name: "Running node", - workerType: "build", - status: "running", - required: true, - dependsOn: [], - modelId: null, - modelProviderId: null, - childSessionId: "ses_child", - output: null, - capturedOutput: null, - errorReason: null, - errorClass: null, - deadlineMs: null, - wakeEligible: true, - wakeReported: false, - replanAttempts: 0, - seq: 1, - timeoutExtensions: 0, - escalationPending: false, - startedAt: 1, - completedAt: null, - timeCreated: 1, - timeUpdated: 2, - }, - { - id: "node_failed", - workflowId: "dag_status", - name: "Failed node", - workerType: "build", - status: "failed", - required: false, - dependsOn: ["node_running"], - modelId: null, - modelProviderId: null, - childSessionId: "ses_failed_child", - output: null, - capturedOutput: null, - errorReason: "node exceeded timeout of 600000ms", - errorClass: "timeout", - deadlineMs: null, - wakeEligible: false, - wakeReported: false, - replanAttempts: 0, - seq: 2, - timeoutExtensions: 0, - escalationPending: false, - startedAt: 1, - completedAt: 2, - timeCreated: 1, - timeUpdated: 2, - }, - ] - : id === "dag_step" - ? [ - { - id: "node_ready", - workflowId: "dag_step", - name: "Ready node", - workerType: "build", - status: "pending", - required: true, - dependsOn: [], - modelId: null, - modelProviderId: null, - childSessionId: null, - output: null, - capturedOutput: null, - errorReason: null, - errorClass: null, - deadlineMs: null, - wakeEligible: false, - wakeReported: false, - replanAttempts: 0, - seq: 1, - timeoutExtensions: 0, - escalationPending: false, - startedAt: null, - completedAt: null, - timeCreated: 1, - timeUpdated: 1, - }, - ] - : [], - ), + getNodes: (id: string) => Effect.succeed(mockNodes(id)), + // Rev-view: status reads the current graph revision; the mock offers the + // same rows (none superseded) so legacy expectations stay intact. + getCurrentNodes: (id: string) => Effect.succeed(mockNodes(id)), getNode: (workflowID: string, nodeID: string) => Effect.succeed(resultNodes.find((node) => node.workflowId === workflowID && node.id === nodeID)), }) diff --git a/packages/schema/src/dag-event.ts b/packages/schema/src/dag-event.ts index dafd8d6c46..c084438ff2 100644 --- a/packages/schema/src/dag-event.ts +++ b/packages/schema/src/dag-event.ts @@ -179,6 +179,13 @@ export const WorkflowReplanned = Event.define({ removed: NonNegativeInt, replaced: NonNegativeInt, restarted: NonNegativeInt, + // Rev-view (v1.0.15 Train A): nodes the replan pushes OUT of the current + // graph revision — terminal rows the fragment bypasses (a failed node the + // new path routes around). NodeCancelled covers the plan.cancel bucket + // separately at projection; this list covers replacements the engine + // never cancels. Optional so legacy durable events still decode (same + // precedent as WorkflowCreated.directory); absent lists project nothing. + superseded: Schema.optional(Schema.Array(NodeID)), }, }) export type WorkflowReplanned = typeof WorkflowReplanned.Type