FN-8911: allow direct database URLs for planning locks

Allow planning lifecycle locks to use non-pooled runtime PostgreSQL connections without a duplicate migration URL.

- Select runtime-direct endpoints for direct DATABASE_URL connections and retain migration overrides for pooled endpoints.
- Preserve lifecycle transport failures across triage retries and expose lock helpers through the core barrel.
- Document the connection behavior and add resolver, advisory-lock, and triage regression coverage.

Files changed:
 .../fn-8911-direct-database-url-planning-lock.md   |   7 +
 docs/architecture.md                               |   2 +-
 docs/multi-project.md                              |   4 +-
 .../__tests__/postgres/backend-resolver.test.ts    |  43 +++++-
 .../planning-lifecycle-advisory-lock.pg.test.ts    |  28 +++-
 packages/core/src/index.ts                         |   2 +
 packages/core/src/postgres/advisory-locks.ts       |   6 +-
 packages/core/src/postgres/backend-resolver.ts     |  29 +++-
 packages/core/src/postgres/index.ts                |   5 +
 ...ge-planning-lifecycle-transport-failure.test.ts | 159 +++++++++++++++++++++
 packages/engine/src/triage.ts                      | 100 +++++++++++--
 11 files changed, 360 insertions(+), 25 deletions(-)

Fusion-Task-Id: FN-8911
Fusion-Task-Lineage: ac2ec0d6-f409-40ba-a678-6dd7da7006dd
Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-09 15:24:21 -07:00
parent 3aa32a846c
commit 8a7ab1dedc
11 changed files with 360 additions and 25 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Allow direct DATABASE_URL connections to finalize planning lifecycle locks.
category: fix
dev: Adds runtime-direct provenance, core lock exports, and a durable planning lock transport marker.

View File

@@ -22,7 +22,7 @@ Each task also stores `lastNotifiedAtByReason`, an independent timestamp map key
Dependency changes and planning finalization share one outer lifecycle lock keyed by the canonical project ID and task ID. In PostgreSQL mode this is a dedicated, single-connection session advisory lock: it is acquired before the normal task lock and released before the mutation/finalization Promise settles. The operational runtime pool is never borrowed for this purpose. Dependency changes and planning finalization share one outer lifecycle lock keyed by the canonical project ID and task ID. In PostgreSQL mode this is a dedicated, single-connection session advisory lock: it is acquired before the normal task lock and released before the mutation/finalization Promise settles. The operational runtime pool is never borrowed for this purpose.
`ResolvedBackend.directSessionUrl` carries the executable lock endpoint. Embedded lifecycle startup supplies its lifecycle-created local URL for runtime, migration, and direct-session use with `embedded-lifecycle` provenance. External deployments must supply an explicit `DATABASE_MIGRATION_URL`; it is the only external direct-session endpoint and has `migration-override` provenance. Missing, mismatched, unavailable, or pooler-like endpoints fail closed rather than falling back to `DATABASE_URL`. `ResolvedBackend.directSessionUrl` carries the executable lock endpoint. Embedded lifecycle startup supplies its lifecycle-created local URL with `embedded-lifecycle` provenance. For external deployments, a non-pooler `DATABASE_URL` is itself the direct-session endpoint with `runtime-direct` provenance. `DATABASE_MIGRATION_URL` is required only when the runtime URL is pooler-like or a distinct schema-work endpoint is wanted; a direct override has `migration-override` provenance. Pooled, mismatched, or unavailable endpoints fail closed.
A newly added dependency atomically records `needs-replan` and clears superseded approval fingerprint/approval-park evidence before a new graph continuation can be eligible. Pending pre-execution continuation work is cancelled with that superseded handoff; running and historical graph work remain durable evidence. Finalizers take the same lock and therefore cannot restore stale handoff state after a re-seed. Manual approval is restored before graph continuation; Plan Review verdicts and results remain graph-owned. A newly added dependency atomically records `needs-replan` and clears superseded approval fingerprint/approval-park evidence before a new graph continuation can be eligible. Pending pre-execution continuation work is cancelled with that superseded handoff; running and historical graph work remain durable evidence. Finalizers take the same lock and therefore cannot restore stale handoff state after a re-seed. Manual approval is restored before graph continuation; Plan Review verdicts and results remain graph-owned.

View File

@@ -33,7 +33,7 @@ Fusion stores multi-project and multi-node coordination state in **PostgreSQL**:
**Default (single machine):** unset `DATABASE_URL` → embedded Postgres under `~/.fusion/embedded-postgres/`. That data directory is **local to the host**. Two laptops each running embedded Postgres do **not** share a board. **Default (single machine):** unset `DATABASE_URL` → embedded Postgres under `~/.fusion/embedded-postgres/`. That data directory is **local to the host**. Two laptops each running embedded Postgres do **not** share a board.
**Multi-node (shared board):** every Fusion node sets the **same external** `DATABASE_URL` (and `DATABASE_MIGRATION_URL` when the runtime URL is a transaction pooler). All nodes share one database; execution (worktrees, agent processes) stays per node. **Multi-node (shared board):** every Fusion node sets the **same external** `DATABASE_URL`. A direct URL needs no duplicate `DATABASE_MIGRATION_URL`; set that override only when the runtime URL is a transaction pooler or schema work needs a separate direct endpoint. All nodes share one database; execution (worktrees, agent processes) stays per node.
Core `central` tables (names as exposed by the data layer; SQL uses snake_case): Core `central` tables (names as exposed by the data layer; SQL uses snake_case):
@@ -59,7 +59,7 @@ Legacy SQLite paths (`~/.fusion/fusion-central.db`, `<repo>/.fusion/fusion.db`)
### Shared Postgres multi-node runbook ### Shared Postgres multi-node runbook
1. Provision one Postgres (local Docker, RDS, Supabase, etc.). 1. Provision one Postgres (local Docker, RDS, Supabase, etc.).
2. On **every** Fusion node: `export DATABASE_URL=...` (same URL). If you use PgBouncer/Supavisor in transaction mode, also set `DATABASE_MIGRATION_URL` to a direct (non-pooled) connection for schema work. 2. On **every** Fusion node: `export DATABASE_URL=...` (same URL). Do not duplicate a direct URL into `DATABASE_MIGRATION_URL`; if you use PgBouncer/Supavisor in transaction mode, set that override to a direct (non-pooled) connection for schema work and planning lifecycle locks.
3. Register projects and nodes so they appear in shared `central.projects` / `central.nodes`. 3. Register projects and nodes so they appear in shared `central.projects` / `central.nodes`.
4. For each host, set `project_node_path_mappings` so that host’s absolute checkout path is recorded for each project. 4. For each host, set `project_node_path_mappings` so that host’s absolute checkout path is recorded for each project.
5. Run `fn serve` / the engine on each node. Task IDs and settings are shared via Postgres; checkout exclusivity uses `task_claims`; abandoned-owner recovery uses `MeshLeaseManager`. 5. Run `fn serve` / the engine on each node. Task IDs and settings are shared via Postgres; checkout exclusivity uses `task_claims`; abandoned-owner recovery uses `MeshLeaseManager`.

View File

@@ -9,6 +9,7 @@ import {
DATABASE_URL_ENV, DATABASE_URL_ENV,
DATABASE_MIGRATION_URL_ENV, DATABASE_MIGRATION_URL_ENV,
} from "../../postgres/backend-resolver.js"; } from "../../postgres/backend-resolver.js";
import { PlanningLifecycleLockTransportError } from "../../index.js";
describe("backend-resolver: resolveBackend (env-based)", () => { describe("backend-resolver: resolveBackend (env-based)", () => {
it("resolves to embedded mode when DATABASE_URL is unset", () => { it("resolves to embedded mode when DATABASE_URL is unset", () => {
@@ -56,7 +57,7 @@ describe("backend-resolver: resolveBackendWithOptions", () => {
}); });
it("DATABASE_MIGRATION_URL routes schema work to it while runtime uses DATABASE_URL (VAL-CONN-003)", () => { it("DATABASE_MIGRATION_URL routes schema work to it while runtime uses DATABASE_URL (VAL-CONN-003)", () => {
const runtimeUrl = "postgresql://user:pass@pooler.supabase.com:6543/fusion"; const runtimeUrl = "postgresql://user:pass@xyz.pooler.supabase.com:6543/fusion";
const migrationUrl = "postgresql://user:pass@db.supabase.co:5432/fusion"; const migrationUrl = "postgresql://user:pass@db.supabase.co:5432/fusion";
const backend = resolveBackendWithOptions({ const backend = resolveBackendWithOptions({
databaseUrl: runtimeUrl, databaseUrl: runtimeUrl,
@@ -78,7 +79,38 @@ describe("backend-resolver: resolveBackendWithOptions", () => {
}); });
expect(backend.migrationUrl).toBe(url); expect(backend.migrationUrl).toBe(url);
expect(backend.migrationUrlOverridden).toBe(false); expect(backend.migrationUrlOverridden).toBe(false);
expect(backend.directSessionUrl).toBe(url);
expect(backend.directSessionProvenance).toBe("runtime-direct");
});
it.each([
["direct runtime with identical direct migration", "postgresql://u:p@localhost:5432/fusion", "postgresql://u:p@localhost:5432/fusion", "migration-override"],
["direct runtime with distinct direct migration", "postgresql://u:p@localhost:5432/fusion", "postgresql://u:p@db.example.com:5432/fusion", "migration-override"],
["pooled runtime with direct migration", "postgresql://u:p@xyz.pooler.supabase.com:6543/fusion", "postgresql://u:p@localhost:5432/fusion", "migration-override"],
["direct runtime with pooled migration", "postgresql://u:p@localhost:5432/fusion", "postgresql://u:p@xyz.pooler.supabase.com:6543/fusion", "runtime-direct"],
])("selects the correct endpoint for %s", (_name, databaseUrl, databaseMigrationUrl, provenance) => {
const backend = resolveBackendWithOptions({ databaseUrl, databaseMigrationUrl });
expect(backend.directSessionProvenance).toBe(provenance);
expect(backend.directSessionUrl).toBe(provenance === "migration-override" ? databaseMigrationUrl : databaseUrl);
});
it.each([
"postgresql://u:p@xyz.pooler.supabase.com:6543/fusion",
"postgresql://u:p@localhost:5432/fusion?pgbouncer=true",
"postgresql://u:p@localhost:5432/fusion?pool_mode=transaction",
])("fails closed without a direct endpoint for pooled runtime URL %s", (databaseUrl) => {
const backend = resolveBackendWithOptions({ databaseUrl });
expect(backend.directSessionUrl).toBeNull(); expect(backend.directSessionUrl).toBeNull();
expect(backend.directSessionProvenance).toBeNull();
});
it("fails closed when both runtime and migration URLs are pooled", () => {
const backend = resolveBackendWithOptions({
databaseUrl: "postgresql://u:p@xyz.pooler.supabase.com:6543/fusion",
databaseMigrationUrl: "postgresql://u:p@localhost:5432/fusion?pgbouncer=true",
});
expect(backend.directSessionUrl).toBeNull();
expect(backend.directSessionProvenance).toBeNull();
}); });
it("DATABASE_MIGRATION_URL without DATABASE_URL still resolves to embedded mode", () => { it("DATABASE_MIGRATION_URL without DATABASE_URL still resolves to embedded mode", () => {
@@ -175,6 +207,7 @@ describe("backend-resolver: describeBackendForLog (VAL-CONN-005)", () => {
expect(desc).toContain("localhost:5432"); expect(desc).toContain("localhost:5432");
expect(desc).not.toContain("hunter2"); expect(desc).not.toContain("hunter2");
expect(desc).toContain("********"); expect(desc).toContain("********");
expect(desc).toContain("planning lifecycle direct session: runtime URL");
}); });
it("migration URL override is logged with redacted URL", () => { it("migration URL override is logged with redacted URL", () => {
@@ -190,3 +223,11 @@ describe("backend-resolver: describeBackendForLog (VAL-CONN-005)", () => {
expect(desc).toContain("host2"); expect(desc).toContain("host2");
}); });
}); });
describe("backend-resolver: package barrel", () => {
it("exports a instanceof-classifiable planning lifecycle transport error", () => {
const error = new PlanningLifecycleLockTransportError("direct endpoint unavailable");
expect(error).toBeInstanceOf(PlanningLifecycleLockTransportError);
});
});

View File

@@ -9,6 +9,7 @@ import {
PlanningLifecycleLockTransportError, PlanningLifecycleLockTransportError,
withPlanningLifecycleAdvisoryLock, withPlanningLifecycleAdvisoryLock,
} from "../../postgres/advisory-locks.js"; } from "../../postgres/advisory-locks.js";
import { resolveBackendWithOptions } from "../../postgres/backend-resolver.js";
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({
prefix: "fusion_planning_lock", prefix: "fusion_planning_lock",
@@ -111,6 +112,28 @@ pgDescribe("planning lifecycle advisory lock", () => {
}, },
); );
it("uses a direct-only runtime URL to acquire and release the lock", async () => {
const backend = resolveBackendWithOptions({ databaseUrl: h.testUrl() });
let callbackRan = false;
await withPlanningLifecycleAdvisoryLock({
projectId: "project-a",
taskId: "FN-runtime-direct",
directSessionUrl: backend.directSessionUrl ?? null,
provenance: backend.directSessionProvenance ?? null,
runtimeUrl: backend.runtimeUrl,
migrationUrl: backend.migrationUrl,
}, async () => { callbackRan = true; });
expect(callbackRan).toBe(true);
await expect(withPlanningLifecycleAdvisoryLock({
projectId: "project-a",
taskId: "FN-runtime-direct",
directSessionUrl: backend.directSessionUrl ?? null,
provenance: backend.directSessionProvenance ?? null,
runtimeUrl: backend.runtimeUrl,
migrationUrl: backend.migrationUrl,
}, async () => {})).resolves.toBeUndefined();
});
it("bounds lock contention with a typed transport error", async () => { it("bounds lock contention with a typed transport error", async () => {
let releaseFirst!: () => void; let releaseFirst!: () => void;
const firstCanFinish = new Promise<void>((resolve) => { releaseFirst = resolve; }); const firstCanFinish = new Promise<void>((resolve) => { releaseFirst = resolve; });
@@ -151,9 +174,10 @@ pgDescribe("planning lifecycle advisory lock", () => {
}, callback)).rejects.toBeInstanceOf(PlanningLifecycleLockTransportError); }, callback)).rejects.toBeInstanceOf(PlanningLifecycleLockTransportError);
await expect(withPlanningLifecycleAdvisoryLock({ await expect(withPlanningLifecycleAdvisoryLock({
...base, ...base,
directSessionUrl: h.testUrl(), directSessionUrl: `${h.testUrl()}_other`,
provenance: "runtime-direct",
runtimeUrl: h.testUrl(), runtimeUrl: h.testUrl(),
migrationUrl: `${h.testUrl()}_other`, migrationUrl: h.testUrl(),
}, callback)).rejects.toBeInstanceOf(PlanningLifecycleLockTransportError); }, callback)).rejects.toBeInstanceOf(PlanningLifecycleLockTransportError);
await expect(withPlanningLifecycleAdvisoryLock({ await expect(withPlanningLifecycleAdvisoryLock({
...base, ...base,

View File

@@ -2563,6 +2563,8 @@ export {
TEST_MODE_ENV, TEST_MODE_ENV,
TEST_DATABASE_URL_ENV, TEST_DATABASE_URL_ENV,
TEST_DATABASE_MIGRATION_URL_ENV, TEST_DATABASE_MIGRATION_URL_ENV,
PlanningLifecycleLockTransportError,
withPlanningLifecycleAdvisoryLock,
} from "./postgres/index.js"; } from "./postgres/index.js";
export type { export type {
BackendMode, BackendMode,

View File

@@ -89,7 +89,7 @@ export async function withPlanningLifecycleAdvisoryLock<T>(
projectId: string; projectId: string;
taskId: string; taskId: string;
directSessionUrl: string | null; directSessionUrl: string | null;
provenance: "embedded-lifecycle" | "migration-override" | null; provenance: "embedded-lifecycle" | "migration-override" | "runtime-direct" | null;
runtimeUrl?: string | null; runtimeUrl?: string | null;
migrationUrl?: string | null; migrationUrl?: string | null;
/** Bounds dedicated-session setup and lock acquisition, not callback work. */ /** Bounds dedicated-session setup and lock acquisition, not callback work. */
@@ -102,7 +102,7 @@ export async function withPlanningLifecycleAdvisoryLock<T>(
const directUrl = input.directSessionUrl; const directUrl = input.directSessionUrl;
const timeoutMs = Math.max(1, input.timeoutMs ?? DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS); const timeoutMs = Math.max(1, input.timeoutMs ?? DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS);
if (!directUrl || !input.provenance || looksLikePoolerUrl(directUrl)) { if (!directUrl || !input.provenance || looksLikePoolerUrl(directUrl)) {
throw new PlanningLifecycleLockTransportError("Planning lifecycle lock requires a direct PostgreSQL session endpoint"); throw new PlanningLifecycleLockTransportError("Planning lifecycle lock requires a direct PostgreSQL session endpoint; set DATABASE_MIGRATION_URL to a direct, non-pooled connection");
} }
let directDatabase: string; let directDatabase: string;
@@ -114,7 +114,7 @@ export async function withPlanningLifecycleAdvisoryLock<T>(
if (!directDatabase) { if (!directDatabase) {
throw new PlanningLifecycleLockTransportError("Planning lifecycle lock direct endpoint must select a database"); throw new PlanningLifecycleLockTransportError("Planning lifecycle lock direct endpoint must select a database");
} }
const descriptorEndpoint = input.provenance === "embedded-lifecycle" const descriptorEndpoint = input.provenance === "embedded-lifecycle" || input.provenance === "runtime-direct"
? input.runtimeUrl ? input.runtimeUrl
: input.migrationUrl; : input.migrationUrl;
if (!descriptorEndpoint || descriptorEndpoint !== directUrl) { if (!descriptorEndpoint || descriptorEndpoint !== directUrl) {

View File

@@ -46,7 +46,7 @@ export interface ResolvedBackend {
readonly migrationUrlOverridden: boolean; readonly migrationUrlOverridden: boolean;
/** A proven direct endpoint for session-scoped lifecycle advisory locks. */ /** A proven direct endpoint for session-scoped lifecycle advisory locks. */
readonly directSessionUrl?: string | null; readonly directSessionUrl?: string | null;
readonly directSessionProvenance?: "embedded-lifecycle" | "migration-override" | null; readonly directSessionProvenance?: "embedded-lifecycle" | "migration-override" | "runtime-direct" | null;
} }
/** /**
@@ -102,14 +102,25 @@ export function resolveBackendWithOptions(
? databaseMigrationUrl ? databaseMigrationUrl
: runtimeUrl; : runtimeUrl;
// FNXC:PlanningDependencyReseed 2026-08-04-00:43: /*
// Advisory locks require one backend session. External runtime URLs can be FNXC:PostgresConnection 2026-08-09-21:53:
// transaction poolers, so only an explicit migration endpoint is eligible. Planning lifecycle locks require one stable PostgreSQL session. A direct runtime URL is
const directSessionUrl = migrationUrlOverridden ? databaseMigrationUrl : null; itself eligible; only a pooled runtime needs DATABASE_MIGRATION_URL because a transaction
pooler can swap sessions between pg_advisory_lock and pg_advisory_unlock.
*/
const directSessionUrl = migrationUrlOverridden && !looksLikePoolerUrl(databaseMigrationUrl)
? databaseMigrationUrl
: runtimeUrl && !looksLikePoolerUrl(runtimeUrl)
? runtimeUrl
: null;
return { return {
mode, runtimeUrl, migrationUrl, migrationUrlOverridden, mode, runtimeUrl, migrationUrl, migrationUrlOverridden,
directSessionUrl, directSessionUrl,
directSessionProvenance: directSessionUrl ? "migration-override" : null, directSessionProvenance: directSessionUrl
? migrationUrlOverridden && directSessionUrl === databaseMigrationUrl
? "migration-override"
: "runtime-direct"
: null,
}; };
} }
@@ -199,5 +210,11 @@ export function describeBackendForLog(backend: ResolvedBackend): string {
`DATABASE_MIGRATION_URL overrides schema-work target: ${redactConnectionString(backend.migrationUrl)}`, `DATABASE_MIGRATION_URL overrides schema-work target: ${redactConnectionString(backend.migrationUrl)}`,
); );
} }
const directSession = backend.directSessionProvenance === "runtime-direct"
? "runtime URL"
: backend.directSessionProvenance === "migration-override"
? "migration URL"
: "none (pooled endpoint — set DATABASE_MIGRATION_URL to a direct connection)";
parts.push(`planning lifecycle direct session: ${directSession}`);
return parts.join(" | "); return parts.join(" | ");
} }

View File

@@ -27,6 +27,11 @@ export {
type ResolveBackendOptions, type ResolveBackendOptions,
} from "./backend-resolver.js"; } from "./backend-resolver.js";
export {
PlanningLifecycleLockTransportError,
withPlanningLifecycleAdvisoryLock,
} from "./advisory-locks.js";
export { export {
createConnectionSet, createConnectionSet,
createConnectionSetFromUrl, createConnectionSetFromUrl,

View File

@@ -0,0 +1,159 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import type { Settings, Task, TaskStore } from "@fusion/core";
import { PlanningLifecycleLockTransportError } from "@fusion/core";
import { mkdtemp, mkdir, rm, writeFile } from "node:fs/promises";
import { join } from "node:path";
import { tmpdir } from "node:os";
const { mockCreateFnAgent, mockPromptWithFallback } = vi.hoisted(() => ({
mockCreateFnAgent: vi.fn(),
mockPromptWithFallback: vi.fn(),
}));
vi.mock("../pi.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("../pi.js")>();
return {
...actual,
createFnAgent: mockCreateFnAgent,
promptWithFallback: mockPromptWithFallback,
};
});
import { TriageProcessor } from "../triage.js";
const TRANSPORT_MARKER = "planning.lifecycleLockTransportFailure";
function createTask(overrides: Partial<Task> = {}): Task {
return {
id: "FN-8911",
description: "Plan a direct PostgreSQL deployment",
column: "triage",
status: null,
dependencies: [],
steps: [],
currentStep: 0,
log: [],
createdAt: "2026-08-09T00:00:00.000Z",
updatedAt: "2026-08-09T00:00:00.000Z",
customFields: { unrelated: "preserve-me" },
...overrides,
};
}
function createPersistedStore(initialTask: Task): { store: TaskStore; task: () => Task; logs: string[] } {
let persisted = initialTask;
const logs: string[] = [];
const update = async (_id: string, patch: Partial<Task>) => {
persisted = { ...persisted, ...patch };
return persisted;
};
const store = {
getTask: vi.fn(async () => ({ ...persisted, attachments: [], comments: [] })),
getSettings: vi.fn(async () => ({
maxConcurrent: 2,
maxWorktrees: 4,
pollIntervalMs: 10_000,
groupOverlappingFiles: false,
autoMerge: true,
} as Settings)),
getTaskWorkflowSelection: vi.fn(() => undefined),
getTaskWorkflowSelectionAsync: vi.fn(async () => undefined),
updateTask: vi.fn(update),
updateTaskAtomic: vi.fn(async (_id: string, patcher: (live: Task) => Partial<Task> | null) => {
const patch = patcher(persisted);
if (patch) await update(_id, patch);
return persisted;
}),
withPlanningLifecycleLock: vi.fn(async () => {
throw new PlanningLifecycleLockTransportError("direct PostgreSQL session endpoint is unavailable");
}),
logEntry: vi.fn(async (_id: string, message: string) => { logs.push(message); }),
appendAgentLog: vi.fn(async () => undefined),
getAgentLogs: vi.fn(async () => []),
listTasks: vi.fn(async () => []),
findRecentTasksBySourceParentTaskId: vi.fn(async () => []),
recordActivity: vi.fn(async () => undefined),
parseDependenciesFromPrompt: vi.fn(async () => []),
parseStepsFromPrompt: vi.fn(async () => []),
parseFileScopeFromPrompt: vi.fn(async () => []),
on: vi.fn(),
emit: vi.fn(),
} as unknown as TaskStore;
return { store, task: () => persisted, logs };
}
describe("triage planning lifecycle lock transport failures (FN-8911)", () => {
let root: string;
beforeEach(async () => {
root = await mkdtemp(join(tmpdir(), "fusion-triage-lock-transport-"));
mockCreateFnAgent.mockResolvedValue({
session: {
prompt: vi.fn(),
dispose: vi.fn(),
sessionManager: { getLeafId: vi.fn(() => null), navigateTree: vi.fn() },
},
});
});
afterEach(async () => {
mockCreateFnAgent.mockReset();
mockPromptWithFallback.mockReset();
await rm(root, { recursive: true, force: true });
});
it("persists transport failures across a fresh triage owner instead of laundering them into unchanged-PROMPT failures", async () => {
const fixture = createPersistedStore(createTask());
const promptPath = join(root, ".fusion", "tasks", "FN-8911", "PROMPT.md");
await mkdir(join(root, ".fusion", "tasks", "FN-8911"), { recursive: true });
mockPromptWithFallback.mockImplementationOnce(async () => {
await writeFile(promptPath, "# Direct PostgreSQL plan\n", "utf8");
});
await new TriageProcessor(fixture.store, root).specifyTask(fixture.task());
expect(fixture.task().error).toBeNull();
expect(fixture.task().recoveryRetryCount).toBe(1);
expect(fixture.task().customFields).toMatchObject({
unrelated: "preserve-me",
[TRANSPORT_MARKER]: {
message: "direct PostgreSQL session endpoint is unavailable",
attempt: 1,
at: expect.any(String),
},
});
expect(fixture.logs.join("\n")).toContain("Planning lifecycle lock transport failure");
// A new processor models a retry claimed after an engine restart or ownership change.
mockPromptWithFallback.mockResolvedValueOnce(undefined);
await new TriageProcessor(fixture.store, root).specifyTask(fixture.task());
const retryLog = fixture.logs.at(-1) ?? "";
expect(retryLog).toContain("Planning lifecycle lock transport failure recorded at");
expect(retryLog).not.toContain("did not update the authoritative PROMPT.md");
expect(fixture.task().recoveryRetryCount).toBe(2);
mockPromptWithFallback.mockResolvedValueOnce(undefined);
await new TriageProcessor(fixture.store, root).specifyTask(fixture.task());
mockPromptWithFallback.mockResolvedValueOnce(undefined);
await new TriageProcessor(fixture.store, root).specifyTask(fixture.task());
expect(fixture.task().status).toBe("failed");
expect(fixture.task().error).toContain("Planning lifecycle lock transport failure recorded at");
expect(fixture.task().error).not.toContain("did not update the authoritative PROMPT.md");
expect(fixture.task().customFields?.[TRANSPORT_MARKER]).toBeUndefined();
});
it("keeps the ordinary unchanged-PROMPT verdict when no persisted transport marker exists", async () => {
const fixture = createPersistedStore(createTask({ customFields: { unrelated: "preserve-me" } }));
const promptPath = join(root, ".fusion", "tasks", "FN-8911", "PROMPT.md");
await mkdir(join(root, ".fusion", "tasks", "FN-8911"), { recursive: true });
await writeFile(promptPath, "# Existing plan\n", "utf8");
mockPromptWithFallback.mockResolvedValue(undefined);
await new TriageProcessor(fixture.store, root).specifyTask(fixture.task());
expect(fixture.logs.join("\n")).toContain("Planner did not update the authoritative PROMPT.md");
expect(fixture.logs.join("\n")).not.toContain("Planning lifecycle lock transport failure recorded at");
});
});

View File

@@ -85,6 +85,24 @@ workflows; metadata is the only authoritative renamed-lane answer in this event
const LEGACY_PLANNER_WAKE_COLUMNS = new Set(["todo", "triage"]); const LEGACY_PLANNER_WAKE_COLUMNS = new Set(["todo", "triage"]);
const LEGACY_PLANNER_COLUMNS = new Set([...LEGACY_PLANNER_WAKE_COLUMNS, "in-progress"]); const LEGACY_PLANNER_COLUMNS = new Set([...LEGACY_PLANNER_WAKE_COLUMNS, "in-progress"]);
const PLANNING_LIFECYCLE_LOCK_TRANSPORT_FAILURE_KEY = "planning.lifecycleLockTransportFailure";
type PlanningLifecycleLockTransportFailure = { message: string; at: string; attempt: number | null };
function getPlanningLifecycleLockTransportFailure(task: Task): PlanningLifecycleLockTransportFailure | null {
const candidate = task.customFields?.[PLANNING_LIFECYCLE_LOCK_TRANSPORT_FAILURE_KEY];
if (!candidate || typeof candidate !== "object" || Array.isArray(candidate)) return null;
const marker = candidate as Partial<PlanningLifecycleLockTransportFailure>;
return typeof marker.message === "string" && typeof marker.at === "string"
? { message: marker.message, at: marker.at, attempt: typeof marker.attempt === "number" ? marker.attempt : null }
: null;
}
function isPlanningLifecycleLockTransportError(error: unknown): error is Error {
return error instanceof fusionCore.PlanningLifecycleLockTransportError
|| (error instanceof Error && error.name === "PlanningLifecycleLockTransportError");
}
/* /*
FNXC:PlanReviewReplan 2026-07-13-00:00: FNXC:PlanReviewReplan 2026-07-13-00:00:
The triage pre-execution Plan Review gate (runPlanReviewBeforeExecution) routes a REVISE The triage pre-execution Plan Review gate (runPlanReviewBeforeExecution) routes a REVISE
@@ -3493,11 +3511,15 @@ export class TriageProcessor {
} }
const artifactChangedByAttempt = planningAttempt.baseline !== written; const artifactChangedByAttempt = planningAttempt.baseline !== written;
if (fallbackDispatchBoundaryMissing || planningAttempt.fallbackEngaged || !artifactChangedByAttempt) { if (fallbackDispatchBoundaryMissing || planningAttempt.fallbackEngaged || !artifactChangedByAttempt) {
const failure = fallbackDispatchBoundaryMissing const liveTask = await Promise.resolve(this.store.getTask(task.id)).catch(() => task) ?? task;
? `Planner runtime ${runtimeId} did not provide a fallback-dispatch settlement boundary for attempt ${planningAttempt.id}` const transportFailure = getPlanningLifecycleLockTransportFailure(liveTask);
: planningAttempt.fallbackEngaged const failure = transportFailure
? `Planner fallback engaged during attempt ${planningAttempt.id}` ? `Planning lifecycle lock transport failure recorded at ${transportFailure.at}: ${transportFailure.message}`
: `Planner did not update the authoritative PROMPT.md during attempt ${planningAttempt.id}`; : fallbackDispatchBoundaryMissing
? `Planner runtime ${runtimeId} did not provide a fallback-dispatch settlement boundary for attempt ${planningAttempt.id}`
: planningAttempt.fallbackEngaged
? `Planner fallback engaged during attempt ${planningAttempt.id}`
: `Planner did not update the authoritative PROMPT.md during attempt ${planningAttempt.id}`;
const decision = computeRecoveryDecision({ const decision = computeRecoveryDecision({
recoveryRetryCount: task.recoveryRetryCount, recoveryRetryCount: task.recoveryRetryCount,
nextRecoveryAt: task.nextRecoveryAt, nextRecoveryAt: task.nextRecoveryAt,
@@ -3519,11 +3541,16 @@ export class TriageProcessor {
const failureMessage = `${failure} after ${MAX_RECOVERY_RETRIES} retries. Retry after adjusting the task prompt or model.`; const failureMessage = `${failure} after ${MAX_RECOVERY_RETRIES} retries. Retry after adjusting the task prompt or model.`;
planLog.error(`${task.id} clean planning attempt retry budget exhausted`); planLog.error(`${task.id} clean planning attempt retry budget exhausted`);
await this.store.logEntry(task.id, failureMessage); await this.store.logEntry(task.id, failureMessage);
if (await this.updatePlanningStateIfStillCurrent(task, { if (await this.updatePlanningStateIfStillCurrent(task, (live) => {
status: "failed", const customFields = { ...(live.customFields ?? {}) };
error: failureMessage, delete customFields[PLANNING_LIFECYCLE_LOCK_TRANSPORT_FAILURE_KEY];
recoveryRetryCount: null, return {
nextRecoveryAt: null, status: "failed",
error: failureMessage,
recoveryRetryCount: null,
nextRecoveryAt: null,
customFields,
};
})) { })) {
await this.backfillBlankTitleAfterTerminalTriageFailure(task); await this.backfillBlankTitleAfterTerminalTriageFailure(task);
} }
@@ -3732,6 +3759,49 @@ export class TriageProcessor {
await this.backfillBlankTitleAfterTerminalTriageFailure(task); await this.backfillBlankTitleAfterTerminalTriageFailure(task);
this.options.onSpecifyError?.(task, err instanceof Error ? err : new Error(errorMessage)); this.options.onSpecifyError?.(task, err instanceof Error ? err : new Error(errorMessage));
return; return;
} else if (isPlanningLifecycleLockTransportError(err)) {
/*
FNXC:PlanningDependencyReseed 2026-08-09-21:53:
A fail-closed lifecycle-lock transport rejection is infrastructure, not planner authoring.
Issue #3394 showed retries could turn a complete spec into a false unchanged-PROMPT verdict;
persist its marker because retry ownership can move across processes and restarts.
*/
const failureMessage = `Planning lifecycle lock transport failure: ${errorMessage}`;
const decision = computeRecoveryDecision({
recoveryRetryCount: task.recoveryRetryCount,
nextRecoveryAt: task.nextRecoveryAt,
});
const persistMarker = (live: Task) => ({
customFields: {
...(live.customFields ?? {}),
[PLANNING_LIFECYCLE_LOCK_TRANSPORT_FAILURE_KEY]: {
message: errorMessage,
at: new Date().toISOString(),
attempt: decision.nextState.recoveryRetryCount,
},
},
});
if (decision.shouldRetry) {
const retryMessage = `${failureMessage} — retry ${decision.nextState.recoveryRetryCount}/${MAX_RECOVERY_RETRIES} in ${formatDelay(decision.delayMs)}.`;
planLog.warn(`${task.id} ${retryMessage}`);
await this.store.logEntry(task.id, retryMessage).catch(() => undefined);
await this.updatePlanningStateIfStillCurrent(task, (live) => ({
...persistMarker(live),
status: this.restoreStatusAfterInterruptedTriageWork(task),
error: null,
recoveryRetryCount: decision.nextState.recoveryRetryCount,
nextRecoveryAt: decision.nextState.nextRecoveryAt,
}));
return;
}
await this.store.logEntry(task.id, failureMessage).catch(() => undefined);
await this.updatePlanningStateIfStillCurrent(task, (live) => {
const customFields = { ...(live.customFields ?? {}) };
delete customFields[PLANNING_LIFECYCLE_LOCK_TRANSPORT_FAILURE_KEY];
return { status: "failed", error: failureMessage, recoveryRetryCount: null, nextRecoveryAt: null, customFields };
});
await this.backfillBlankTitleAfterTerminalTriageFailure(task);
return;
} else if (isTransientError(errorMessage)) { } else if (isTransientError(errorMessage)) {
// Transient network/infrastructure error — use bounded recovery policy // Transient network/infrastructure error — use bounded recovery policy
const decision = computeRecoveryDecision({ const decision = computeRecoveryDecision({
@@ -4593,6 +4663,11 @@ export class TriageProcessor {
* near-duplicate decision; keep removes the marker before the next real plan. * near-duplicate decision; keep removes the marker before the next real plan.
*/ */
if (explicitDuplicateMarker) { if (explicitDuplicateMarker) {
await this.updatePlanningStateIfStillCurrent(task, (live) => {
const customFields = { ...(live.customFields ?? {}) };
delete customFields[PLANNING_LIFECYCLE_LOCK_TRANSPORT_FAILURE_KEY];
return { customFields };
});
const canonicalId = explicitDuplicateMarker.canonicalId; const canonicalId = explicitDuplicateMarker.canonicalId;
const duplicateSource = duplicateResolution.source ?? "prompt"; const duplicateSource = duplicateResolution.source ?? "prompt";
const canonicalTask = await this.store.getTask(canonicalId).catch(() => null); const canonicalTask = await this.store.getTask(canonicalId).catch(() => null);
@@ -5192,6 +5267,11 @@ export class TriageProcessor {
in the bookkeeping that follows does not un-hand-off a card that has already moved. in the bookkeeping that follows does not un-hand-off a card that has already moved.
*/ */
report.outcome = "released"; report.outcome = "released";
await this.updatePlanningStateIfStillCurrent(task, (live) => {
const customFields = { ...(live.customFields ?? {}) };
delete customFields[PLANNING_LIFECYCLE_LOCK_TRANSPORT_FAILURE_KEY];
return { customFields };
});
/* /*
FNXC:TriageStuckKill 2026-07-18-21:05: FNXC:TriageStuckKill 2026-07-18-21:05: