Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4917df1f07 |
@@ -34,6 +34,7 @@ export default tseslint.config(
|
||||
'packages/storage/vitest.config.ts',
|
||||
'packages/mosaic/vitest.config.ts',
|
||||
'packages/mosaic/__tests__/*.ts',
|
||||
'packages/forge/__tests__/*.ts',
|
||||
'tools/federation-harness/*.ts',
|
||||
],
|
||||
},
|
||||
|
||||
@@ -539,3 +539,43 @@ Not every brief needs full Board of Directors review. The classification system
|
||||
### Backward compatibility
|
||||
|
||||
Existing briefs without a `class` field are auto-classified. The default (no matching keywords) is `strategic`, so all existing runs get the full pipeline unless keywords trigger `technical`.
|
||||
|
||||
---
|
||||
|
||||
## Fail-Closed Execution & Explicit Simulation (SDLC-D-035)
|
||||
|
||||
**Added:** 2026-08-17
|
||||
|
||||
Forge fails closed when a required capability is missing. It never runs a
|
||||
pipeline with a stub executor and reports success.
|
||||
|
||||
### Normal mode (default)
|
||||
|
||||
- No task executor wired → the CLI exits nonzero with the typed capability
|
||||
error `FORGE_NO_EXECUTOR`. No run is created.
|
||||
- A stage whose gate is approval-based (board approval, planning approvals,
|
||||
remediation re-review, discovery/analysis attestations) records a typed
|
||||
`waiting-for-authority` stage result and raises `FORGE_AUTHORITY_REQUIRED`.
|
||||
It never passes vacuously.
|
||||
- A stage whose gate requires an unwired provider (AI reviewer, CI pipeline)
|
||||
records a typed `blocked` stage result and raises `FORGE_NO_REVIEWER` /
|
||||
`FORGE_NO_CI_PIPELINE`. The synthetic echo-review approval in `06-review`
|
||||
and all vacuous `true` gates were removed.
|
||||
|
||||
### Explicit simulation (`--simulate`)
|
||||
|
||||
Opts into stub/synthetic execution. Every stage result, every gate result, and
|
||||
the run manifest carry the distinct typed status `simulated` (manifest also
|
||||
records `mode: "simulated"`). `simulated` is a non-satisfying outcome:
|
||||
`isSatisfyingOutcome()` and all completion/gate consumers treat only `passed`
|
||||
as satisfying. The CLI exits 0 for a simulated run only because the caller
|
||||
explicitly passed `--simulate`, and prints a loud SIMULATED banner.
|
||||
|
||||
### Typed outcome model
|
||||
|
||||
Every gate/task outcome is one of the closed set
|
||||
`passed | failed | blocked | error | waiting-for-authority | simulated |
|
||||
not-applicable`, with the reason recorded on the stage status and each gate
|
||||
result in `manifest.json`. Missing implementations, missing gate evidence,
|
||||
unknown stages, process errors, and timeouts map to fail-closed members —
|
||||
never to `passed`.
|
||||
|
||||
@@ -0,0 +1,319 @@
|
||||
import fs from 'node:fs';
|
||||
import os from 'node:os';
|
||||
import path from 'node:path';
|
||||
import { describe, it, expect, beforeEach, afterEach } from 'vitest';
|
||||
|
||||
import { generateBoardTasks } from '../src/board-tasks.js';
|
||||
import { STAGE_SPECS } from '../src/constants.js';
|
||||
import { ForgeCapabilityError } from '../src/errors.js';
|
||||
import {
|
||||
evaluateStageGates,
|
||||
gateLabel,
|
||||
isCommandGate,
|
||||
isSatisfyingOutcome,
|
||||
} from '../src/outcomes.js';
|
||||
import { loadManifest, runPipeline } from '../src/pipeline-runner.js';
|
||||
import type { ForgeTask, ForgeTaskResult, TaskExecutor } from '../src/types.js';
|
||||
|
||||
/**
|
||||
* Mock real executor that returns typed results.
|
||||
*
|
||||
* Command gates are "verified" by the mock so normal-mode runs can pass
|
||||
* mechanically gated stages; authority/provider gates are never reported
|
||||
* because they have no mechanical implementation.
|
||||
*/
|
||||
function createTypedExecutor(options?: {
|
||||
failStage?: string;
|
||||
gateOutcomes?: Record<string, 'passed' | 'failed' | 'simulated' | 'error' | 'blocked'>;
|
||||
}): TaskExecutor & { submittedTasks: ForgeTask[] } {
|
||||
const submittedTasks: ForgeTask[] = [];
|
||||
return {
|
||||
submittedTasks,
|
||||
async submitTask(task: ForgeTask) {
|
||||
submittedTasks.push(task);
|
||||
},
|
||||
async waitForCompletion(taskId: string): Promise<ForgeTaskResult> {
|
||||
const task = submittedTasks.find((t) => t.id === taskId);
|
||||
const stageName = task?.metadata?.['stageName'] as string | undefined;
|
||||
|
||||
if (options?.failStage && stageName === options.failStage) {
|
||||
return {
|
||||
task_id: taskId,
|
||||
outcome: 'failed',
|
||||
reason: 'mock task failure',
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 1,
|
||||
gate_results: [],
|
||||
};
|
||||
}
|
||||
|
||||
const gateResults = (task?.qualityGates ?? [])
|
||||
.filter((gate) => isCommandGate(gate))
|
||||
.map((gate) => {
|
||||
const label = gateLabel(gate);
|
||||
const outcome = options?.gateOutcomes?.[label] ?? 'passed';
|
||||
return {
|
||||
gate: label,
|
||||
outcome,
|
||||
reason: outcome === 'passed' ? 'mock verified' : `mock gate outcome: ${outcome}`,
|
||||
};
|
||||
});
|
||||
|
||||
return {
|
||||
task_id: taskId,
|
||||
outcome: 'passed',
|
||||
reason: 'mock verified',
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 0,
|
||||
gate_results: gateResults,
|
||||
};
|
||||
},
|
||||
async getTaskStatus() {
|
||||
return 'completed' as const;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
describe('fail-closed: no executor wired', () => {
|
||||
let tmpDir: string;
|
||||
let briefPath: string;
|
||||
|
||||
beforeEach(() => {
|
||||
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'forge-failclosed-'));
|
||||
briefPath = path.join(tmpDir, 'brief.md');
|
||||
fs.writeFileSync(briefPath, '# Fix bug\n\nA bugfix for lint cleanup.');
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
fs.rmSync(tmpDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
it('throws a typed FORGE_NO_EXECUTOR capability error without --simulate', async () => {
|
||||
await expect(
|
||||
runPipeline(briefPath, tmpDir, {
|
||||
// no executor, no simulate — must fail closed, never run with a stub
|
||||
stages: ['00-intake'],
|
||||
}),
|
||||
).rejects.toMatchObject({
|
||||
name: 'ForgeCapabilityError',
|
||||
code: 'FORGE_NO_EXECUTOR',
|
||||
capability: 'task-executor',
|
||||
});
|
||||
});
|
||||
|
||||
it('does not create a run directory when failing closed on a missing executor', async () => {
|
||||
try {
|
||||
await runPipeline(briefPath, tmpDir, { stages: ['00-intake'] });
|
||||
} catch {
|
||||
// expected
|
||||
}
|
||||
expect(fs.existsSync(path.join(tmpDir, '.forge', 'runs'))).toBe(false);
|
||||
});
|
||||
|
||||
it('completes with every result typed simulated when simulate is set', async () => {
|
||||
const result = await runPipeline(briefPath, tmpDir, {
|
||||
simulate: true,
|
||||
stages: ['00-intake', '00b-discovery', '02-planning-1', '06-review'],
|
||||
});
|
||||
|
||||
expect(result.manifest.mode).toBe('simulated');
|
||||
expect(result.manifest.status).toBe('simulated');
|
||||
|
||||
for (const stage of result.stages) {
|
||||
const stageStatus = result.manifest.stages[stage];
|
||||
expect(stageStatus?.status, `stage ${stage}`).toBe('simulated');
|
||||
expect(stageStatus?.status, `stage ${stage}`).not.toBe('passed');
|
||||
expect(stageStatus?.reason, `stage ${stage}`).toBeTruthy();
|
||||
for (const gateResult of stageStatus?.gateResults ?? []) {
|
||||
expect(gateResult.outcome, `gate ${gateResult.gate} of ${stage}`).toBe('simulated');
|
||||
expect(gateResult.outcome, `gate ${gateResult.gate} of ${stage}`).not.toBe('passed');
|
||||
}
|
||||
}
|
||||
|
||||
// The persisted manifest agrees.
|
||||
const persisted = loadManifest(result.runDir);
|
||||
expect(persisted.mode).toBe('simulated');
|
||||
expect(persisted.status).toBe('simulated');
|
||||
expect(persisted.stages['02-planning-1']?.status).toBe('simulated');
|
||||
});
|
||||
});
|
||||
|
||||
describe('fail-closed: typed outcome model', () => {
|
||||
it('only passed satisfies the gate/dependency predicate', () => {
|
||||
expect(isSatisfyingOutcome('passed')).toBe(true);
|
||||
expect(isSatisfyingOutcome('failed')).toBe(false);
|
||||
expect(isSatisfyingOutcome('blocked')).toBe(false);
|
||||
expect(isSatisfyingOutcome('error')).toBe(false);
|
||||
expect(isSatisfyingOutcome('waiting-for-authority')).toBe(false);
|
||||
expect(isSatisfyingOutcome('simulated')).toBe(false);
|
||||
expect(isSatisfyingOutcome('not-applicable')).toBe(false);
|
||||
});
|
||||
|
||||
it('a simulated gate result cannot satisfy the stage gate evaluation', () => {
|
||||
const evaluation = evaluateStageGates('05-coding', STAGE_SPECS['05-coding']!.qualityGates, {
|
||||
task_id: 'FORGE-x-05',
|
||||
outcome: 'passed',
|
||||
reason: 'executor claims success',
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 0,
|
||||
gate_results: [{ gate: 'pnpm lint', outcome: 'simulated', reason: 'simulated gate' }],
|
||||
});
|
||||
expect(isSatisfyingOutcome(evaluation.outcome)).toBe(false);
|
||||
expect(evaluation.outcome).toBe('error');
|
||||
});
|
||||
|
||||
it('a simulated task outcome cannot satisfy evaluation in normal mode', () => {
|
||||
const evaluation = evaluateStageGates('00-intake', [], {
|
||||
task_id: 'FORGE-x-00',
|
||||
outcome: 'simulated',
|
||||
reason: 'executor reported simulated',
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 0,
|
||||
gate_results: [],
|
||||
});
|
||||
expect(isSatisfyingOutcome(evaluation.outcome)).toBe(false);
|
||||
});
|
||||
|
||||
it('a missing gate result blocks the stage instead of passing vacuously', () => {
|
||||
const evaluation = evaluateStageGates('05-coding', STAGE_SPECS['05-coding']!.qualityGates, {
|
||||
task_id: 'FORGE-x-05',
|
||||
outcome: 'passed',
|
||||
reason: 'executor claims success',
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 0,
|
||||
gate_results: [],
|
||||
});
|
||||
expect(evaluation.outcome).toBe('blocked');
|
||||
});
|
||||
});
|
||||
|
||||
describe('fail-closed: authority and provider gates', () => {
|
||||
let tmpDir: string;
|
||||
let briefPath: string;
|
||||
|
||||
beforeEach(() => {
|
||||
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'forge-authority-'));
|
||||
briefPath = path.join(tmpDir, 'brief.md');
|
||||
fs.writeFileSync(briefPath, '# Fix bug\n\nA bugfix for lint cleanup.');
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
fs.rmSync(tmpDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
it.each(['02-planning-1', '03-planning-2', '04-planning-3', '07-remediate'])(
|
||||
'planning/remediation stage %s yields waiting-for-authority (not passed) in normal mode',
|
||||
async (stage) => {
|
||||
const executor = createTypedExecutor();
|
||||
let runDir: string | undefined;
|
||||
|
||||
try {
|
||||
await runPipeline(briefPath, tmpDir, {
|
||||
executor,
|
||||
stages: [stage as string],
|
||||
});
|
||||
expect.unreachable('runPipeline should have failed closed');
|
||||
} catch (err) {
|
||||
expect(err).toBeInstanceOf(ForgeCapabilityError);
|
||||
expect((err as ForgeCapabilityError).code).toBe('FORGE_AUTHORITY_REQUIRED');
|
||||
runDir = path.join(tmpDir, '.forge', 'runs');
|
||||
}
|
||||
|
||||
const runIds = fs.readdirSync(runDir!);
|
||||
expect(runIds).toHaveLength(1);
|
||||
const manifest = loadManifest(path.join(runDir!, runIds[0]!));
|
||||
expect(manifest.stages[stage]?.status).toBe('waiting-for-authority');
|
||||
expect(manifest.stages[stage]?.status).not.toBe('passed');
|
||||
expect(manifest.status).toBe('waiting-for-authority');
|
||||
},
|
||||
);
|
||||
|
||||
it('review stage fails closed with a typed FORGE_NO_REVIEWER error in normal mode', async () => {
|
||||
const executor = createTypedExecutor();
|
||||
|
||||
try {
|
||||
await runPipeline(briefPath, tmpDir, {
|
||||
executor,
|
||||
stages: ['06-review'],
|
||||
});
|
||||
expect.unreachable('runPipeline should have failed closed');
|
||||
} catch (err) {
|
||||
expect(err).toBeInstanceOf(ForgeCapabilityError);
|
||||
expect((err as ForgeCapabilityError).code).toBe('FORGE_NO_REVIEWER');
|
||||
expect((err as ForgeCapabilityError).capability).toBe('reviewer');
|
||||
}
|
||||
|
||||
const runsDir = path.join(tmpDir, '.forge', 'runs');
|
||||
const runIds = fs.readdirSync(runsDir);
|
||||
const manifest = loadManifest(path.join(runsDir, runIds[0]!));
|
||||
expect(manifest.stages['06-review']?.status).toBe('blocked');
|
||||
expect(manifest.stages['06-review']?.status).not.toBe('passed');
|
||||
expect(manifest.status).toBe('failed');
|
||||
});
|
||||
|
||||
it('review stage produces simulated results under --simulate', async () => {
|
||||
const result = await runPipeline(briefPath, tmpDir, {
|
||||
simulate: true,
|
||||
stages: ['06-review'],
|
||||
});
|
||||
|
||||
expect(result.manifest.mode).toBe('simulated');
|
||||
expect(result.manifest.stages['06-review']?.status).toBe('simulated');
|
||||
for (const gateResult of result.manifest.stages['06-review']?.gateResults ?? []) {
|
||||
expect(gateResult.outcome).toBe('simulated');
|
||||
}
|
||||
});
|
||||
|
||||
it('deploy stage fails closed without a wired ci-pipeline provider in normal mode', async () => {
|
||||
const executor = createTypedExecutor();
|
||||
|
||||
await expect(
|
||||
runPipeline(briefPath, tmpDir, {
|
||||
executor,
|
||||
stages: ['09-deploy'],
|
||||
}),
|
||||
).rejects.toMatchObject({
|
||||
name: 'ForgeCapabilityError',
|
||||
code: 'FORGE_NO_CI_PIPELINE',
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
describe('fail-closed: no vacuous gate commands remain', () => {
|
||||
it('stage constants contain no echo/synthetic-approval, vacuous true, or empty gate commands', () => {
|
||||
for (const [stageName, spec] of Object.entries(STAGE_SPECS)) {
|
||||
for (const gate of spec.qualityGates) {
|
||||
const serialized = JSON.stringify(gate);
|
||||
// The echo-review synthetic approval must be gone.
|
||||
expect(serialized, `stage ${stageName} gate ${serialized}`).not.toContain('echo');
|
||||
expect(serialized, `stage ${stageName} gate ${serialized}`).not.toMatch(/"verdict"\s*:/);
|
||||
expect(serialized, `stage ${stageName} gate ${serialized}`).not.toMatch(
|
||||
/"summary"\s*:\s*"review-pass"/,
|
||||
);
|
||||
// No vacuous literal `true` gate.
|
||||
expect(gate, `stage ${stageName}`).not.toBe('true');
|
||||
// Command gates must carry a real, non-empty command.
|
||||
if (isCommandGate(gate)) {
|
||||
const command = typeof gate === 'string' ? gate : gate.command;
|
||||
expect(command.trim().length, `stage ${stageName} gate ${serialized}`).toBeGreaterThan(0);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
it('board tasks contain no vacuous true gates', () => {
|
||||
const tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'forge-board-gates-'));
|
||||
try {
|
||||
const tasks = generateBoardTasks('# Brief', [], tmpDir, 'BOARD-TEST');
|
||||
for (const task of tasks) {
|
||||
for (const gate of task.qualityGates) {
|
||||
expect(gate, `task ${task.id}`).not.toBe('true');
|
||||
const serialized = JSON.stringify(gate);
|
||||
expect(serialized, `task ${task.id} gate ${serialized}`).not.toContain('echo');
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
fs.rmSync(tmpDir, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -12,10 +12,10 @@ import {
|
||||
resumePipeline,
|
||||
getPipelineStatus,
|
||||
} from '../src/pipeline-runner.js';
|
||||
import type { ForgeTask, RunManifest, TaskExecutor } from '../src/types.js';
|
||||
import type { TaskResult } from '@mosaicstack/macp';
|
||||
import type { ForgeTask, ForgeTaskResult, RunManifest, TaskExecutor } from '../src/types.js';
|
||||
import { gateLabel, isCommandGate } from '../src/outcomes.js';
|
||||
|
||||
/** Mock TaskExecutor that records submitted tasks and returns success. */
|
||||
/** Mock TaskExecutor that records submitted tasks and returns typed results. */
|
||||
function createMockExecutor(options?: {
|
||||
failStage?: string;
|
||||
}): TaskExecutor & { submittedTasks: ForgeTask[] } {
|
||||
@@ -25,7 +25,7 @@ function createMockExecutor(options?: {
|
||||
async submitTask(task: ForgeTask) {
|
||||
submittedTasks.push(task);
|
||||
},
|
||||
async waitForCompletion(taskId: string): Promise<TaskResult> {
|
||||
async waitForCompletion(taskId: string): Promise<ForgeTaskResult> {
|
||||
const failStage = options?.failStage;
|
||||
const task = submittedTasks.find((t) => t.id === taskId);
|
||||
const stageName = task?.metadata?.['stageName'] as string | undefined;
|
||||
@@ -33,7 +33,8 @@ function createMockExecutor(options?: {
|
||||
if (failStage && stageName === failStage) {
|
||||
return {
|
||||
task_id: taskId,
|
||||
status: 'failed',
|
||||
outcome: 'failed',
|
||||
reason: 'mock task failure',
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 1,
|
||||
gate_results: [],
|
||||
@@ -41,10 +42,17 @@ function createMockExecutor(options?: {
|
||||
}
|
||||
return {
|
||||
task_id: taskId,
|
||||
status: 'completed',
|
||||
outcome: 'passed',
|
||||
reason: 'mock verified',
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 0,
|
||||
gate_results: [],
|
||||
gate_results: (task?.qualityGates ?? [])
|
||||
.filter((gate) => isCommandGate(gate))
|
||||
.map((gate) => ({
|
||||
gate: gateLabel(gate),
|
||||
outcome: 'passed' as const,
|
||||
reason: 'mock verified',
|
||||
})),
|
||||
};
|
||||
},
|
||||
async getTaskStatus() {
|
||||
@@ -156,12 +164,13 @@ describe('runPipeline', () => {
|
||||
const executor = createMockExecutor();
|
||||
const result = await runPipeline(briefPath, tmpDir, {
|
||||
executor,
|
||||
stages: ['00-intake', '00b-discovery'],
|
||||
stages: ['00-intake', '05-coding'],
|
||||
});
|
||||
|
||||
expect(result.runId).toMatch(/^\d{8}-\d{6}$/);
|
||||
expect(result.stages).toEqual(['00-intake', '00b-discovery']);
|
||||
expect(result.stages).toEqual(['00-intake', '05-coding']);
|
||||
expect(result.manifest.status).toBe('completed');
|
||||
expect(result.manifest.mode).toBe('normal');
|
||||
expect(executor.submittedTasks).toHaveLength(2);
|
||||
});
|
||||
|
||||
@@ -180,12 +189,17 @@ describe('runPipeline', () => {
|
||||
const executor = createMockExecutor();
|
||||
const result = await runPipeline(briefPath, tmpDir, {
|
||||
executor,
|
||||
stages: ['00-intake', '00b-discovery'],
|
||||
stages: ['00-intake', '05-coding'],
|
||||
});
|
||||
|
||||
const manifest = loadManifest(result.runDir);
|
||||
expect(manifest.stages['00-intake']?.status).toBe('passed');
|
||||
expect(manifest.stages['00b-discovery']?.status).toBe('passed');
|
||||
expect(manifest.stages['05-coding']?.status).toBe('passed');
|
||||
expect(manifest.stages['05-coding']?.gateResults?.map((g) => g.outcome)).toEqual([
|
||||
'passed',
|
||||
'passed',
|
||||
'passed',
|
||||
]);
|
||||
});
|
||||
|
||||
it('respects CLI class override', async () => {
|
||||
@@ -215,7 +229,7 @@ describe('runPipeline', () => {
|
||||
const executor = createMockExecutor();
|
||||
await runPipeline(briefPath, tmpDir, {
|
||||
executor,
|
||||
stages: ['00-intake', '00b-discovery', '02-planning-1'],
|
||||
stages: ['00-intake', '05-coding', '08-test'],
|
||||
});
|
||||
|
||||
expect(executor.submittedTasks[0]!.dependsOn).toBeUndefined();
|
||||
@@ -224,14 +238,14 @@ describe('runPipeline', () => {
|
||||
});
|
||||
|
||||
it('handles stage failure', async () => {
|
||||
const executor = createMockExecutor({ failStage: '00b-discovery' });
|
||||
const executor = createMockExecutor({ failStage: '05-coding' });
|
||||
|
||||
await expect(
|
||||
runPipeline(briefPath, tmpDir, {
|
||||
executor,
|
||||
stages: ['00-intake', '00b-discovery'],
|
||||
stages: ['00-intake', '05-coding'],
|
||||
}),
|
||||
).rejects.toThrow('Stage 00b-discovery failed');
|
||||
).rejects.toThrow('Stage 05-coding failed');
|
||||
});
|
||||
|
||||
it('marks manifest as failed on stage failure', async () => {
|
||||
@@ -270,30 +284,143 @@ describe('resumePipeline', () => {
|
||||
fs.rmSync(tmpDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
it('resumes from first incomplete stage', async () => {
|
||||
// First run fails on discovery
|
||||
const executor1 = createMockExecutor({ failStage: '00b-discovery' });
|
||||
let runDir: string;
|
||||
it('resumes from first incomplete stage and fails closed at the next provider gate', async () => {
|
||||
// Simulate a run whose authority stages were approved out-of-band
|
||||
// (recorded as passed) and whose coding stage failed mechanically.
|
||||
const runId = '20260101-000000';
|
||||
const runDir = path.join(tmpDir, '.forge', 'runs', runId);
|
||||
fs.mkdirSync(runDir, { recursive: true });
|
||||
const passed = { status: 'passed' as const, startedAt: '2026-01-01T00:00:00Z' };
|
||||
saveManifest(runDir, {
|
||||
runId,
|
||||
brief: briefPath,
|
||||
codebase: tmpDir,
|
||||
briefClass: 'hotfix',
|
||||
classSource: 'frontmatter',
|
||||
forceBoard: false,
|
||||
mode: 'normal',
|
||||
createdAt: '2026-01-01T00:00:00Z',
|
||||
updatedAt: '2026-01-01T00:00:00Z',
|
||||
currentStage: '05-coding',
|
||||
status: 'failed',
|
||||
stages: {
|
||||
'00-intake': passed,
|
||||
'00b-discovery': passed,
|
||||
'02-planning-1': passed,
|
||||
'03-planning-2': passed,
|
||||
'04-planning-3': passed,
|
||||
'05-coding': { status: 'failed', reason: 'gate failed' },
|
||||
},
|
||||
});
|
||||
|
||||
try {
|
||||
await runPipeline(briefPath, tmpDir, {
|
||||
executor: executor1,
|
||||
stages: ['00-intake', '00b-discovery', '02-planning-1'],
|
||||
});
|
||||
} catch {
|
||||
// expected
|
||||
// Resume re-runs 05-coding (the first non-passed stage), then fails
|
||||
// closed at 06-review because no reviewer provider is wired.
|
||||
const executor = createMockExecutor();
|
||||
await expect(resumePipeline(runDir, executor)).rejects.toMatchObject({
|
||||
name: 'ForgeCapabilityError',
|
||||
code: 'FORGE_NO_REVIEWER',
|
||||
});
|
||||
|
||||
const manifest = loadManifest(runDir);
|
||||
expect(manifest.stages['05-coding']?.status).toBe('passed');
|
||||
expect(manifest.stages['06-review']?.status).toBe('blocked');
|
||||
expect(manifest.status).toBe('failed');
|
||||
});
|
||||
|
||||
it('resumes to completion as simulated under explicit simulate', async () => {
|
||||
const runId = '20260101-000003';
|
||||
const runDir = path.join(tmpDir, '.forge', 'runs', runId);
|
||||
fs.mkdirSync(runDir, { recursive: true });
|
||||
const passed = { status: 'passed' as const, startedAt: '2026-01-01T00:00:00Z' };
|
||||
saveManifest(runDir, {
|
||||
runId,
|
||||
brief: briefPath,
|
||||
codebase: tmpDir,
|
||||
briefClass: 'hotfix',
|
||||
classSource: 'frontmatter',
|
||||
forceBoard: false,
|
||||
mode: 'normal',
|
||||
createdAt: '2026-01-01T00:00:00Z',
|
||||
updatedAt: '2026-01-01T00:00:00Z',
|
||||
currentStage: '05-coding',
|
||||
status: 'failed',
|
||||
stages: {
|
||||
'00-intake': passed,
|
||||
'00b-discovery': passed,
|
||||
'02-planning-1': passed,
|
||||
'03-planning-2': passed,
|
||||
'04-planning-3': passed,
|
||||
'05-coding': { status: 'failed', reason: 'gate failed' },
|
||||
},
|
||||
});
|
||||
|
||||
const result = await resumePipeline(runDir, undefined, { simulate: true });
|
||||
|
||||
expect(result.manifest.status).toBe('simulated');
|
||||
expect(result.manifest.mode).toBe('simulated');
|
||||
expect(result.stages[0]).toBe('05-coding');
|
||||
for (const stage of result.stages) {
|
||||
expect(result.manifest.stages[stage]?.status).toBe('simulated');
|
||||
}
|
||||
});
|
||||
|
||||
const runsDir = path.join(tmpDir, '.forge', 'runs');
|
||||
runDir = path.join(runsDir, fs.readdirSync(runsDir)[0]!);
|
||||
it('fails closed on resume when the next stage needs authority sign-off', async () => {
|
||||
const runId = '20260101-000001';
|
||||
const runDir = path.join(tmpDir, '.forge', 'runs', runId);
|
||||
fs.mkdirSync(runDir, { recursive: true });
|
||||
saveManifest(runDir, {
|
||||
runId,
|
||||
brief: briefPath,
|
||||
codebase: tmpDir,
|
||||
briefClass: 'hotfix',
|
||||
classSource: 'frontmatter',
|
||||
forceBoard: false,
|
||||
mode: 'normal',
|
||||
createdAt: '2026-01-01T00:00:00Z',
|
||||
updatedAt: '2026-01-01T00:00:00Z',
|
||||
currentStage: '00-intake',
|
||||
status: 'in_progress',
|
||||
stages: {
|
||||
'00-intake': { status: 'passed' },
|
||||
},
|
||||
});
|
||||
|
||||
// Resume should pick up from 00b-discovery
|
||||
const executor2 = createMockExecutor();
|
||||
const result = await resumePipeline(runDir, executor2);
|
||||
const executor = createMockExecutor();
|
||||
await expect(resumePipeline(runDir, executor)).rejects.toMatchObject({
|
||||
name: 'ForgeCapabilityError',
|
||||
code: 'FORGE_AUTHORITY_REQUIRED',
|
||||
});
|
||||
|
||||
expect(result.manifest.status).toBe('completed');
|
||||
// Should have re-run from 00b-discovery onward
|
||||
expect(result.stages[0]).toBe('00b-discovery');
|
||||
const manifest = loadManifest(runDir);
|
||||
expect(manifest.stages['00b-discovery']?.status).toBe('waiting-for-authority');
|
||||
expect(manifest.status).toBe('waiting-for-authority');
|
||||
});
|
||||
|
||||
it('fails closed on resume without an executor or --simulate', async () => {
|
||||
const runId = '20260101-000002';
|
||||
const runDir = path.join(tmpDir, '.forge', 'runs', runId);
|
||||
fs.mkdirSync(runDir, { recursive: true });
|
||||
saveManifest(runDir, {
|
||||
runId,
|
||||
brief: briefPath,
|
||||
codebase: tmpDir,
|
||||
briefClass: 'hotfix',
|
||||
classSource: 'frontmatter',
|
||||
forceBoard: false,
|
||||
mode: 'normal',
|
||||
createdAt: '2026-01-01T00:00:00Z',
|
||||
updatedAt: '2026-01-01T00:00:00Z',
|
||||
currentStage: '00-intake',
|
||||
status: 'in_progress',
|
||||
stages: {
|
||||
'00-intake': { status: 'passed' },
|
||||
},
|
||||
});
|
||||
|
||||
await expect(resumePipeline(runDir)).rejects.toMatchObject({
|
||||
name: 'ForgeCapabilityError',
|
||||
code: 'FORGE_NO_EXECUTOR',
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -95,7 +95,14 @@ export function generateBoardTasks(
|
||||
briefPath,
|
||||
resultPath: resultRelPath,
|
||||
timeoutSeconds: 120,
|
||||
qualityGates: ['true'],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 'board-approval',
|
||||
reason:
|
||||
'persona evaluation is judged by board synthesis (authority review); no mechanical gate exists',
|
||||
},
|
||||
],
|
||||
metadata: {
|
||||
personaName: persona.name,
|
||||
personaSlug: persona.slug,
|
||||
@@ -121,7 +128,13 @@ export function generateBoardTasks(
|
||||
timeoutSeconds: 120,
|
||||
dependsOn: personaTaskIds,
|
||||
dependsOnPolicy: 'all_terminal',
|
||||
qualityGates: ['true'],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 'board-approval',
|
||||
reason: 'board synthesis is an authority decision; no mechanical gate exists',
|
||||
},
|
||||
],
|
||||
metadata: {
|
||||
resultOutputPath: synthesisResult,
|
||||
inputResultPaths: personaResultPaths,
|
||||
|
||||
@@ -1,7 +1,11 @@
|
||||
import fs from 'node:fs';
|
||||
import os from 'node:os';
|
||||
import path from 'node:path';
|
||||
import { Command } from 'commander';
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { describe, expect, it, vi, beforeEach, afterEach } from 'vitest';
|
||||
|
||||
import { registerForgeCommand } from './cli.js';
|
||||
import { loadManifest } from './pipeline-runner.js';
|
||||
|
||||
describe('registerForgeCommand', () => {
|
||||
it('registers a "forge" command on the parent program', () => {
|
||||
@@ -55,3 +59,94 @@ describe('registerForgeCommand', () => {
|
||||
}).not.toThrow();
|
||||
});
|
||||
});
|
||||
|
||||
describe('forge run fail-closed behavior (SDLC-D-035)', () => {
|
||||
let tmpDir: string;
|
||||
let briefPath: string;
|
||||
let errSpy: ReturnType<typeof vi.spyOn>;
|
||||
let logSpy: ReturnType<typeof vi.spyOn>;
|
||||
let prevExitCode: string | number | null | undefined;
|
||||
|
||||
const parse = (args: string[]) => {
|
||||
const program = new Command();
|
||||
registerForgeCommand(program);
|
||||
return program.parseAsync(['forge', ...args], { from: 'user' });
|
||||
};
|
||||
|
||||
beforeEach(() => {
|
||||
tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'forge-cli-failclosed-'));
|
||||
briefPath = path.join(tmpDir, 'brief.md');
|
||||
fs.writeFileSync(briefPath, '# Fix bug\n\nA bugfix for lint cleanup.');
|
||||
errSpy = vi.spyOn(console, 'error').mockImplementation(() => {});
|
||||
logSpy = vi.spyOn(console, 'log').mockImplementation(() => {});
|
||||
prevExitCode = process.exitCode;
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
errSpy.mockRestore();
|
||||
logSpy.mockRestore();
|
||||
process.exitCode = prevExitCode;
|
||||
fs.rmSync(tmpDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
it('exits nonzero with a typed FORGE_NO_EXECUTOR error when no executor is wired and --simulate is absent', async () => {
|
||||
await parse(['run', '--brief', briefPath, '--codebase', tmpDir]);
|
||||
|
||||
expect(process.exitCode).toBe(1);
|
||||
const errText = errSpy.mock.calls.map((c) => c.join(' ')).join('\n');
|
||||
expect(errText).toContain('FORGE_NO_EXECUTOR');
|
||||
// It must never run the pipeline with a stub and report success.
|
||||
expect(fs.existsSync(path.join(tmpDir, '.forge', 'runs'))).toBe(false);
|
||||
});
|
||||
|
||||
it('completes with typed simulated results and exit 0 under explicit --simulate', async () => {
|
||||
await parse(['run', '--brief', briefPath, '--codebase', tmpDir, '--simulate']);
|
||||
|
||||
expect(process.exitCode).toBeUndefined();
|
||||
|
||||
// Loud simulated-mode summary.
|
||||
const logText = logSpy.mock.calls.map((c) => c.join(' ')).join('\n');
|
||||
expect(logText).toContain('SIMULATED');
|
||||
|
||||
// Manifest records the mode and simulated per-result statuses.
|
||||
const runsDir = path.join(tmpDir, '.forge', 'runs');
|
||||
const runIds = fs.readdirSync(runsDir);
|
||||
expect(runIds).toHaveLength(1);
|
||||
const manifest = loadManifest(path.join(runsDir, runIds[0]!));
|
||||
expect(manifest.mode).toBe('simulated');
|
||||
expect(manifest.status).toBe('simulated');
|
||||
for (const stageStatus of Object.values(manifest.stages)) {
|
||||
expect(stageStatus?.status).toBe('simulated');
|
||||
for (const gateResult of stageStatus?.gateResults ?? []) {
|
||||
expect(gateResult.outcome).toBe('simulated');
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
it('resume exits nonzero with a typed FORGE_NO_EXECUTOR error without --simulate', async () => {
|
||||
const runDir = path.join(tmpDir, '.forge', 'runs', '20260101-000000');
|
||||
fs.mkdirSync(runDir, { recursive: true });
|
||||
fs.writeFileSync(
|
||||
path.join(runDir, 'manifest.json'),
|
||||
JSON.stringify({
|
||||
runId: '20260101-000000',
|
||||
brief: briefPath,
|
||||
codebase: tmpDir,
|
||||
briefClass: 'hotfix',
|
||||
classSource: 'frontmatter',
|
||||
forceBoard: false,
|
||||
createdAt: '2026-01-01T00:00:00Z',
|
||||
updatedAt: '2026-01-01T00:00:00Z',
|
||||
currentStage: '00-intake',
|
||||
status: 'in_progress',
|
||||
stages: { '00-intake': { status: 'passed' } },
|
||||
}),
|
||||
);
|
||||
|
||||
await parse(['resume', '20260101-000000', '--project', tmpDir]);
|
||||
|
||||
expect(process.exitCode).toBe(1);
|
||||
const errText = errSpy.mock.calls.map((c) => c.join(' ')).join('\n');
|
||||
expect(errText).toContain('FORGE_NO_EXECUTOR');
|
||||
});
|
||||
});
|
||||
|
||||
+122
-48
@@ -5,37 +5,47 @@ import type { Command } from 'commander';
|
||||
|
||||
import { classifyBrief } from './brief-classifier.js';
|
||||
import { STAGE_LABELS, STAGE_SEQUENCE } from './constants.js';
|
||||
import { ForgeCapabilityError } from './errors.js';
|
||||
import { getEffectivePersonas, loadBoardPersonas } from './persona-loader.js';
|
||||
import { generateRunId, getPipelineStatus, loadManifest, runPipeline } from './pipeline-runner.js';
|
||||
import type { PipelineOptions, RunManifest, TaskExecutor } from './types.js';
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Stub executor — used when no real executor is wired at CLI invocation time.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
const stubExecutor: TaskExecutor = {
|
||||
async submitTask(task) {
|
||||
console.log(` [forge] stage submitted: ${task.id} (${task.title})`);
|
||||
},
|
||||
async waitForCompletion(taskId, _timeoutMs) {
|
||||
console.log(` [forge] stage complete: ${taskId}`);
|
||||
return {
|
||||
task_id: taskId,
|
||||
status: 'completed' as const,
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 0,
|
||||
gate_results: [],
|
||||
};
|
||||
},
|
||||
async getTaskStatus(_taskId) {
|
||||
return 'completed' as const;
|
||||
},
|
||||
};
|
||||
import { createSimulatedExecutor } from './simulated-executor.js';
|
||||
import type { PipelineOptions, RunManifest, RunMode } from './types.js';
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Helpers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/** Resolve a run's effective mode, defaulting legacy manifests to normal. */
|
||||
function runModeOf(manifest: RunManifest): RunMode {
|
||||
return manifest.mode ?? 'normal';
|
||||
}
|
||||
|
||||
/** Print a loud banner so a simulated run can never be misread as verified. */
|
||||
function printSimulatedBanner(): void {
|
||||
console.log('');
|
||||
console.log('[forge] ===============================================================');
|
||||
console.log('[forge] MODE: SIMULATED — no stage or gate was really executed.');
|
||||
console.log('[forge] All results are synthetic and MUST NOT be read as verified');
|
||||
console.log('[forge] success. Wire a real executor/providers and re-run to verify.');
|
||||
console.log('[forge] ===============================================================');
|
||||
}
|
||||
|
||||
/** Print a typed error line for fail-closed capability errors. */
|
||||
function printCapabilityError(err: ForgeCapabilityError): void {
|
||||
console.error(`[forge] error ${err.code}: ${err.message}`);
|
||||
console.error(`[forge] missing capability: ${err.capability}`);
|
||||
}
|
||||
|
||||
/** Handle a pipeline error uniformly: typed capability errors get their code. */
|
||||
function handlePipelineError(err: unknown): void {
|
||||
if (err instanceof ForgeCapabilityError) {
|
||||
printCapabilityError(err);
|
||||
} else {
|
||||
console.error(`[forge] pipeline failed: ${err instanceof Error ? err.message : String(err)}`);
|
||||
}
|
||||
process.exitCode = 1;
|
||||
}
|
||||
|
||||
function formatDuration(startedAt?: string, completedAt?: string): string {
|
||||
if (!startedAt || !completedAt) return '-';
|
||||
const ms = new Date(completedAt).getTime() - new Date(startedAt).getTime();
|
||||
@@ -44,19 +54,24 @@ function formatDuration(startedAt?: string, completedAt?: string): string {
|
||||
}
|
||||
|
||||
function printManifestTable(manifest: RunManifest): void {
|
||||
const mode = runModeOf(manifest);
|
||||
console.log(`\nRun ID : ${manifest.runId}`);
|
||||
console.log(`Status : ${manifest.status}`);
|
||||
console.log(`Mode : ${mode}`);
|
||||
if (mode === 'simulated') {
|
||||
console.log('WARNING: SIMULATED RUN — results are synthetic, not verified success.');
|
||||
}
|
||||
console.log(`Brief : ${manifest.brief}`);
|
||||
console.log(`Class : ${manifest.briefClass} (${manifest.classSource})`);
|
||||
console.log(`Updated: ${manifest.updatedAt}`);
|
||||
console.log('');
|
||||
console.log('Stage'.padEnd(22) + 'Status'.padEnd(14) + 'Duration');
|
||||
console.log('-'.repeat(50));
|
||||
console.log('Stage'.padEnd(22) + 'Status'.padEnd(24) + 'Duration');
|
||||
console.log('-'.repeat(60));
|
||||
for (const stage of STAGE_SEQUENCE) {
|
||||
const s = manifest.stages[stage];
|
||||
if (!s) continue;
|
||||
const label = (STAGE_LABELS[stage] ?? stage).padEnd(22);
|
||||
const status = s.status.padEnd(14);
|
||||
const status = s.status.padEnd(24);
|
||||
const dur = formatDuration(s.startedAt, s.completedAt);
|
||||
console.log(`${label}${status}${dur}`);
|
||||
}
|
||||
@@ -90,23 +105,58 @@ function listRecentRuns(projectRoot?: string): void {
|
||||
}
|
||||
|
||||
console.log('\nRecent runs:');
|
||||
console.log('Run ID'.padEnd(22) + 'Status'.padEnd(14) + 'Brief');
|
||||
console.log('-'.repeat(70));
|
||||
console.log('Run ID'.padEnd(22) + 'Status'.padEnd(24) + 'Mode'.padEnd(12) + 'Brief');
|
||||
console.log('-'.repeat(80));
|
||||
|
||||
for (const runId of entries) {
|
||||
const runDir = path.join(runsDir, runId);
|
||||
try {
|
||||
const manifest = loadManifest(runDir);
|
||||
const status = manifest.status.padEnd(14);
|
||||
const status = manifest.status.padEnd(24);
|
||||
const mode = runModeOf(manifest).padEnd(12);
|
||||
const brief = path.basename(manifest.brief);
|
||||
console.log(`${runId.padEnd(22)}${status}${brief}`);
|
||||
console.log(`${runId.padEnd(22)}${status}${mode}${brief}`);
|
||||
} catch {
|
||||
console.log(`${runId.padEnd(22)}${'(unreadable)'.padEnd(14)}`);
|
||||
console.log(`${runId.padEnd(22)}${'(unreadable)'.padEnd(24)}`);
|
||||
}
|
||||
}
|
||||
console.log('');
|
||||
}
|
||||
|
||||
/**
|
||||
* Apply the exit-code policy for a finished pipeline run (SDLC-D-035):
|
||||
*
|
||||
* - exit 0 only for a verified `completed` normal run, or for an overall
|
||||
* `simulated` run when the caller explicitly passed --simulate;
|
||||
* - anything else exits nonzero so it can never be read as success.
|
||||
*/
|
||||
function applyRunExitPolicy(result: { manifest: RunManifest; runDir: string }, simulate: boolean) {
|
||||
const { manifest } = result;
|
||||
|
||||
if (runModeOf(manifest) === 'simulated') {
|
||||
if (!simulate || manifest.status !== 'simulated') {
|
||||
console.error(
|
||||
'[forge] error FORGE_MODE_MISMATCH: run reports simulated results without an explicit, ' +
|
||||
'consistent --simulate request; refusing to report success.',
|
||||
);
|
||||
process.exitCode = 1;
|
||||
return;
|
||||
}
|
||||
printSimulatedBanner();
|
||||
console.log(`[forge] run directory: ${result.runDir}`);
|
||||
return; // exit 0 — the caller explicitly opted into simulation
|
||||
}
|
||||
|
||||
if (manifest.status !== 'completed') {
|
||||
console.error(`[forge] run did not complete: terminal status '${manifest.status}'`);
|
||||
process.exitCode = 1;
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(`[forge] pipeline complete (mode: normal): ${manifest.runId}`);
|
||||
console.log(`[forge] run directory: ${result.runDir}`);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Register function
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -129,6 +179,11 @@ export function registerForgeCommand(parent: Command): void {
|
||||
.option('--config <path>', 'Path to forge config file (.forge/config.yaml)')
|
||||
.option('--codebase <path>', 'Codebase root to pass to the pipeline', process.cwd())
|
||||
.option('--dry-run', 'Print planned stages without executing', false)
|
||||
.option(
|
||||
'--simulate',
|
||||
'Simulate execution without real providers (every result is typed simulated, never verified)',
|
||||
false,
|
||||
)
|
||||
.action(
|
||||
async (opts: {
|
||||
brief: string;
|
||||
@@ -137,6 +192,7 @@ export function registerForgeCommand(parent: Command): void {
|
||||
config?: string;
|
||||
codebase: string;
|
||||
dryRun: boolean;
|
||||
simulate: boolean;
|
||||
}) => {
|
||||
const briefPath = path.resolve(opts.brief);
|
||||
|
||||
@@ -149,14 +205,22 @@ export function registerForgeCommand(parent: Command): void {
|
||||
const briefContent = fs.readFileSync(briefPath, 'utf-8');
|
||||
const briefClass = classifyBrief(briefContent);
|
||||
const projectRoot = opts.codebase;
|
||||
// A real executor is never wired at CLI invocation time today, so the
|
||||
// only executor we may construct is the explicitly-requested simulated
|
||||
// one. Normal mode fails closed with FORGE_NO_EXECUTOR.
|
||||
const executor = opts.simulate ? createSimulatedExecutor() : undefined;
|
||||
|
||||
if (opts.resume) {
|
||||
const runId = opts.runId ?? generateRunId();
|
||||
const runDir = resolveRunDir(runId, projectRoot);
|
||||
console.log(`[forge] resuming run: ${runId}`);
|
||||
const { resumePipeline } = await import('./pipeline-runner.js');
|
||||
const result = await resumePipeline(runDir, stubExecutor);
|
||||
console.log(`[forge] pipeline complete: ${result.runId}`);
|
||||
try {
|
||||
const { resumePipeline } = await import('./pipeline-runner.js');
|
||||
const result = await resumePipeline(runDir, executor, { simulate: opts.simulate });
|
||||
applyRunExitPolicy(result, opts.simulate);
|
||||
} catch (err) {
|
||||
handlePipelineError(err);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -164,7 +228,8 @@ export function registerForgeCommand(parent: Command): void {
|
||||
briefClass,
|
||||
codebase: projectRoot,
|
||||
dryRun: opts.dryRun,
|
||||
executor: stubExecutor,
|
||||
executor,
|
||||
simulate: opts.simulate,
|
||||
};
|
||||
|
||||
if (opts.dryRun) {
|
||||
@@ -180,16 +245,15 @@ export function registerForgeCommand(parent: Command): void {
|
||||
|
||||
console.log(`[forge] starting pipeline for brief: ${briefPath}`);
|
||||
console.log(`[forge] classified as: ${briefClass}`);
|
||||
if (opts.simulate) {
|
||||
console.log('[forge] mode: SIMULATED (explicit --simulate)');
|
||||
}
|
||||
|
||||
try {
|
||||
const result = await runPipeline(briefPath, projectRoot, pipelineOptions);
|
||||
console.log(`[forge] pipeline complete: ${result.runId}`);
|
||||
console.log(`[forge] run directory: ${result.runDir}`);
|
||||
applyRunExitPolicy(result, opts.simulate);
|
||||
} catch (err) {
|
||||
console.error(
|
||||
`[forge] pipeline failed: ${err instanceof Error ? err.message : String(err)}`,
|
||||
);
|
||||
process.exitCode = 1;
|
||||
handlePipelineError(err);
|
||||
}
|
||||
},
|
||||
);
|
||||
@@ -224,7 +288,12 @@ export function registerForgeCommand(parent: Command): void {
|
||||
.command('resume <runId>')
|
||||
.description('Resume a stopped or failed pipeline run')
|
||||
.option('--project <path>', 'Project root (defaults to cwd)', process.cwd())
|
||||
.action(async (runId: string, opts: { project: string }) => {
|
||||
.option(
|
||||
'--simulate',
|
||||
'Simulate execution without real providers (every result is typed simulated, never verified)',
|
||||
false,
|
||||
)
|
||||
.action(async (runId: string, opts: { project: string; simulate: boolean }) => {
|
||||
const runDir = resolveRunDir(runId, opts.project);
|
||||
|
||||
if (!fs.existsSync(runDir)) {
|
||||
@@ -234,15 +303,20 @@ export function registerForgeCommand(parent: Command): void {
|
||||
}
|
||||
|
||||
console.log(`[forge] resuming run: ${runId}`);
|
||||
if (opts.simulate) {
|
||||
console.log('[forge] mode: SIMULATED (explicit --simulate)');
|
||||
}
|
||||
|
||||
// No real executor is wired at CLI invocation time; only the explicitly
|
||||
// requested simulated executor may be constructed (fail closed otherwise).
|
||||
const executor = opts.simulate ? createSimulatedExecutor() : undefined;
|
||||
|
||||
try {
|
||||
const { resumePipeline } = await import('./pipeline-runner.js');
|
||||
const result = await resumePipeline(runDir, stubExecutor);
|
||||
console.log(`[forge] pipeline complete: ${result.runId}`);
|
||||
console.log(`[forge] run directory: ${result.runDir}`);
|
||||
const result = await resumePipeline(runDir, executor, { simulate: opts.simulate });
|
||||
applyRunExitPolicy(result, opts.simulate);
|
||||
} catch (err) {
|
||||
console.error(`[forge] resume failed: ${err instanceof Error ? err.message : String(err)}`);
|
||||
process.exitCode = 1;
|
||||
handlePipelineError(err);
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
@@ -9,7 +9,16 @@ export const PACKAGE_ROOT = path.resolve(path.dirname(fileURLToPath(import.meta.
|
||||
/** Pipeline asset directory (stages, agents, rails, gates, templates). */
|
||||
export const PIPELINE_DIR = path.join(PACKAGE_ROOT, 'pipeline');
|
||||
|
||||
/** Stage specifications — defines every pipeline stage. */
|
||||
/** Stage specifications — defines every pipeline stage.
|
||||
*\n * Gate semantics (SDLC-D-035): every gate is one of
|
||||
* - a real command string / GateEntry a mechanical runner can execute,
|
||||
* - an `authority` gate (human/board sign-off; produces waiting-for-authority),
|
||||
* - a `provider` gate (requires a wired provider such as a reviewer or CI pipeline).
|
||||
*
|
||||
* Vacuous gates (`true`, echo'd synthetic approvals, placeholder ci-pipeline
|
||||
* commands) are forbidden: a stage whose gate has no real implementation
|
||||
* fails closed instead of passing.
|
||||
*/
|
||||
export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
'00-intake': {
|
||||
number: '00',
|
||||
@@ -27,7 +36,13 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
type: 'research',
|
||||
gate: 'discovery-complete',
|
||||
promptFile: '00b-discovery.md',
|
||||
qualityGates: ['true'],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 'discovery-complete',
|
||||
reason: 'discovery completion is attested by an authority; no mechanical check exists',
|
||||
},
|
||||
],
|
||||
},
|
||||
'01-board': {
|
||||
number: '01',
|
||||
@@ -36,7 +51,13 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
type: 'review',
|
||||
gate: 'board-approval',
|
||||
promptFile: '01-board.md',
|
||||
qualityGates: [{ type: 'ci-pipeline', command: 'board-approval (via board-tasks)' }],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 'board-approval',
|
||||
reason: 'board approval is a board/human decision; no mechanical gate exists',
|
||||
},
|
||||
],
|
||||
},
|
||||
'01b-brief-analyzer': {
|
||||
number: '01b',
|
||||
@@ -45,7 +66,13 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
type: 'research',
|
||||
gate: 'brief-analysis-complete',
|
||||
promptFile: '01-board.md',
|
||||
qualityGates: ['true'],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 'brief-analysis-complete',
|
||||
reason: 'brief analysis completion is attested by an authority; no mechanical check exists',
|
||||
},
|
||||
],
|
||||
},
|
||||
'02-planning-1': {
|
||||
number: '02',
|
||||
@@ -54,7 +81,13 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
type: 'research',
|
||||
gate: 'architecture-approval',
|
||||
promptFile: '02-planning-1-architecture.md',
|
||||
qualityGates: ['true'],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 'architecture-approval',
|
||||
reason: 'ADR approval requires authority sign-off; no mechanical check exists',
|
||||
},
|
||||
],
|
||||
},
|
||||
'03-planning-2': {
|
||||
number: '03',
|
||||
@@ -63,7 +96,14 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
type: 'research',
|
||||
gate: 'implementation-approval',
|
||||
promptFile: '03-planning-2-implementation.md',
|
||||
qualityGates: ['true'],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 'implementation-approval',
|
||||
reason:
|
||||
'implementation spec approval requires authority sign-off; no mechanical check exists',
|
||||
},
|
||||
],
|
||||
},
|
||||
'04-planning-3': {
|
||||
number: '04',
|
||||
@@ -72,7 +112,14 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
type: 'research',
|
||||
gate: 'decomposition-approval',
|
||||
promptFile: '04-planning-3-decomposition.md',
|
||||
qualityGates: ['true'],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 'decomposition-approval',
|
||||
reason:
|
||||
'task decomposition approval requires authority sign-off; no mechanical check exists',
|
||||
},
|
||||
],
|
||||
},
|
||||
'05-coding': {
|
||||
number: '05',
|
||||
@@ -92,9 +139,10 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
promptFile: '06-review.md',
|
||||
qualityGates: [
|
||||
{
|
||||
type: 'ai-review',
|
||||
command:
|
||||
'echo \'{"summary":"review-pass","verdict":"approve","findings":[],"stats":{"blockers":0,"should_fix":0,"suggestions":0}}\'',
|
||||
kind: 'provider',
|
||||
capability: 'reviewer',
|
||||
reason:
|
||||
'review verdicts require a wired reviewer provider; synthetic approvals are not permitted',
|
||||
},
|
||||
],
|
||||
},
|
||||
@@ -105,7 +153,13 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
type: 'coding',
|
||||
gate: 're-review',
|
||||
promptFile: '07-remediate.md',
|
||||
qualityGates: ['true'],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'authority',
|
||||
capability: 're-review',
|
||||
reason: 'remediation re-review is an approval-based gate; no mechanical check exists',
|
||||
},
|
||||
],
|
||||
},
|
||||
'08-test': {
|
||||
number: '08',
|
||||
@@ -123,7 +177,13 @@ export const STAGE_SPECS: Record<string, StageSpec> = {
|
||||
type: 'deploy',
|
||||
gate: 'deploy-verification',
|
||||
promptFile: '09-deploy.md',
|
||||
qualityGates: [{ type: 'ci-pipeline', command: 'deploy-verification' }],
|
||||
qualityGates: [
|
||||
{
|
||||
kind: 'provider',
|
||||
capability: 'ci-pipeline',
|
||||
reason: 'deploy verification requires a wired CI pipeline provider',
|
||||
},
|
||||
],
|
||||
},
|
||||
};
|
||||
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
/**
|
||||
* Typed fail-closed capability errors (SDLC-D-035).
|
||||
*
|
||||
* A Forge run must fail closed when a required capability (executor, reviewer
|
||||
* provider, CI pipeline, authority sign-off) is missing. These typed errors
|
||||
* name the missing capability so callers can distinguish "not wired" from
|
||||
* ordinary execution failures.
|
||||
*/
|
||||
|
||||
/** Closed set of typed Forge capability error codes. */
|
||||
export const FORGE_ERROR_CODES = [
|
||||
'FORGE_NO_EXECUTOR',
|
||||
'FORGE_NO_REVIEWER',
|
||||
'FORGE_NO_CI_PIPELINE',
|
||||
'FORGE_NO_PROVIDER',
|
||||
'FORGE_AUTHORITY_REQUIRED',
|
||||
] as const;
|
||||
|
||||
export type ForgeErrorCode = (typeof FORGE_ERROR_CODES)[number];
|
||||
|
||||
/** Raised when a required capability is missing and the pipeline must fail closed. */
|
||||
export class ForgeCapabilityError extends Error {
|
||||
/** Typed error code from the closed FORGE_ERROR_CODES set. */
|
||||
readonly code: ForgeErrorCode;
|
||||
/** The missing capability, e.g. `task-executor`, `reviewer`, `board-approval`. */
|
||||
readonly capability: string;
|
||||
|
||||
constructor(code: ForgeErrorCode, capability: string, message: string) {
|
||||
super(message);
|
||||
this.name = 'ForgeCapabilityError';
|
||||
this.code = code;
|
||||
this.capability = capability;
|
||||
}
|
||||
}
|
||||
|
||||
/** Map a provider gate capability to its typed error code. */
|
||||
export function providerErrorCode(capability: string): ForgeErrorCode {
|
||||
switch (capability) {
|
||||
case 'reviewer':
|
||||
return 'FORGE_NO_REVIEWER';
|
||||
case 'ci-pipeline':
|
||||
return 'FORGE_NO_CI_PIPELINE';
|
||||
default:
|
||||
return 'FORGE_NO_PROVIDER';
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,13 @@ export type {
|
||||
StageSpec,
|
||||
BriefClass,
|
||||
ClassSource,
|
||||
ForgeOutcome,
|
||||
AuthorityGate,
|
||||
ProviderGate,
|
||||
ForgeGate,
|
||||
ForgeGateResult,
|
||||
ForgeTaskResult,
|
||||
RunMode,
|
||||
StageStatus,
|
||||
RunManifest,
|
||||
ForgeTaskStatus,
|
||||
@@ -81,5 +88,24 @@ export {
|
||||
getPipelineStatus,
|
||||
} from './pipeline-runner.js';
|
||||
|
||||
// Fail-closed errors and typed outcome model (SDLC-D-035)
|
||||
export { FORGE_ERROR_CODES, ForgeCapabilityError, providerErrorCode } from './errors.js';
|
||||
export type { ForgeErrorCode } from './errors.js';
|
||||
export {
|
||||
isSatisfyingOutcome,
|
||||
isCapabilityGate,
|
||||
isCommandGate,
|
||||
gateLabel,
|
||||
uniformGateResults,
|
||||
simulatedGateResults,
|
||||
waitingGateResults,
|
||||
blockedGateResults,
|
||||
evaluateStageGates,
|
||||
} from './outcomes.js';
|
||||
export type { StageEvaluation } from './outcomes.js';
|
||||
|
||||
// Simulated executor (explicit --simulate only)
|
||||
export { createSimulatedExecutor } from './simulated-executor.js';
|
||||
|
||||
// CLI
|
||||
export { registerForgeCommand } from './cli.js';
|
||||
|
||||
@@ -0,0 +1,147 @@
|
||||
import type { GateEntry } from '@mosaicstack/macp';
|
||||
|
||||
import type {
|
||||
AuthorityGate,
|
||||
ForgeGate,
|
||||
ForgeGateResult,
|
||||
ForgeOutcome,
|
||||
ForgeTaskResult,
|
||||
ProviderGate,
|
||||
} from './types.js';
|
||||
|
||||
/**
|
||||
* Gate and dependency satisfaction predicate (SDLC-D-035).
|
||||
*
|
||||
* ONLY a verified `passed` outcome satisfies. Every other member of the closed
|
||||
* outcome set — including `simulated` — is non-satisfying, so a simulated or
|
||||
* authority-blocked result can never be read as success-by-verification.
|
||||
*/
|
||||
export function isSatisfyingOutcome(outcome: ForgeOutcome): boolean {
|
||||
return outcome === 'passed';
|
||||
}
|
||||
|
||||
/** Whether a gate is an authority or provider gate (capability-based, command-less). */
|
||||
export function isCapabilityGate(gate: ForgeGate): gate is AuthorityGate | ProviderGate {
|
||||
if (typeof gate !== 'object' || gate === null) return false;
|
||||
const kind = (gate as Record<string, unknown>)['kind'];
|
||||
return kind === 'authority' || kind === 'provider';
|
||||
}
|
||||
|
||||
/** Whether a gate definition carries a real command a mechanical runner can execute. */
|
||||
export function isCommandGate(gate: ForgeGate): gate is string | GateEntry {
|
||||
if (typeof gate === 'string') {
|
||||
return gate.trim().length > 0;
|
||||
}
|
||||
if (isCapabilityGate(gate)) {
|
||||
// Authority and provider gates are satisfied by a capability, not a command.
|
||||
return false;
|
||||
}
|
||||
return typeof gate.command === 'string' && gate.command.trim().length > 0;
|
||||
}
|
||||
|
||||
/** Typed label identifying a gate in results and logs. */
|
||||
export function gateLabel(gate: ForgeGate): string {
|
||||
if (typeof gate === 'string') return gate;
|
||||
if (isCapabilityGate(gate)) return `${gate.kind}:${gate.capability}`;
|
||||
return gate.command || gate.type || 'unnamed-gate';
|
||||
}
|
||||
|
||||
/** Reason string stamped on every simulated gate result. */
|
||||
export const SIMULATED_GATE_REASON =
|
||||
'simulated execution (--simulate): gate was not evaluated by a real implementation';
|
||||
|
||||
/** Build typed gate results with a uniform outcome for a stage's declared gates. */
|
||||
export function uniformGateResults(
|
||||
gates: ForgeGate[],
|
||||
outcome: ForgeOutcome,
|
||||
reason: string,
|
||||
): ForgeGateResult[] {
|
||||
return gates.map((gate) => ({ gate: gateLabel(gate), outcome, reason }));
|
||||
}
|
||||
|
||||
/** Typed simulated gate results — used exclusively in `--simulate` runs. */
|
||||
export function simulatedGateResults(gates: ForgeGate[]): ForgeGateResult[] {
|
||||
return uniformGateResults(gates, 'simulated', SIMULATED_GATE_REASON);
|
||||
}
|
||||
|
||||
/** Typed waiting-for-authority gate results for approval-based stages. */
|
||||
export function waitingGateResults(gates: ForgeGate[], reason: string): ForgeGateResult[] {
|
||||
return uniformGateResults(gates, 'waiting-for-authority', reason);
|
||||
}
|
||||
|
||||
/** Typed blocked gate results for stages whose provider capability is not wired. */
|
||||
export function blockedGateResults(gates: ForgeGate[], reason: string): ForgeGateResult[] {
|
||||
return uniformGateResults(gates, 'blocked', reason);
|
||||
}
|
||||
|
||||
/** Outcome of evaluating a completed stage in normal mode. */
|
||||
export interface StageEvaluation {
|
||||
outcome: ForgeOutcome;
|
||||
reason: string;
|
||||
gateResults: ForgeGateResult[];
|
||||
}
|
||||
|
||||
/**
|
||||
* Evaluate a stage's declared gates against the executor's typed result.
|
||||
*
|
||||
* Fail-closed mapping:
|
||||
* - a `simulated` task or gate outcome in normal mode maps to `error`
|
||||
* - a missing gate result for a required command gate maps to `blocked`
|
||||
* - a non-passing task outcome propagates as the stage outcome
|
||||
* - only verified `passed` task and gate outcomes yield a `passed` stage
|
||||
*/
|
||||
export function evaluateStageGates(
|
||||
stageName: string,
|
||||
gates: ForgeGate[],
|
||||
result: ForgeTaskResult,
|
||||
): StageEvaluation {
|
||||
const gateResults = result.gate_results ?? [];
|
||||
|
||||
if (result.outcome === 'simulated') {
|
||||
return {
|
||||
outcome: 'error',
|
||||
reason: `executor reported a simulated outcome for stage '${stageName}' in normal mode — refusing to treat simulated results as verified`,
|
||||
gateResults,
|
||||
};
|
||||
}
|
||||
|
||||
if (!isSatisfyingOutcome(result.outcome)) {
|
||||
return {
|
||||
outcome: result.outcome,
|
||||
reason: `task outcome is '${result.outcome}': ${result.reason}`,
|
||||
gateResults,
|
||||
};
|
||||
}
|
||||
|
||||
for (const gate of gates) {
|
||||
// Authority and provider gates are pre-flighted before execution; they have
|
||||
// no mechanical result to verify here.
|
||||
if (!isCommandGate(gate)) continue;
|
||||
|
||||
const label = gateLabel(gate);
|
||||
const gateResult = gateResults.find((r) => r.gate === label);
|
||||
if (!gateResult) {
|
||||
return {
|
||||
outcome: 'blocked',
|
||||
reason: `no gate result was reported for required gate '${label}' (stage '${stageName}')`,
|
||||
gateResults,
|
||||
};
|
||||
}
|
||||
if (!isSatisfyingOutcome(gateResult.outcome)) {
|
||||
return {
|
||||
outcome: gateResult.outcome === 'simulated' ? 'error' : gateResult.outcome,
|
||||
reason: `gate '${label}' outcome is '${gateResult.outcome}': ${gateResult.reason}`,
|
||||
gateResults,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
outcome: 'passed',
|
||||
reason:
|
||||
gates.length === 0
|
||||
? "stage declares no gates; task outcome 'passed' accepted"
|
||||
: 'all declared gates verified passed',
|
||||
gateResults,
|
||||
};
|
||||
}
|
||||
@@ -1,18 +1,33 @@
|
||||
import fs from 'node:fs';
|
||||
import path from 'node:path';
|
||||
|
||||
import { STAGE_SEQUENCE } from './constants.js';
|
||||
import { STAGE_SEQUENCE, STAGE_SPECS } from './constants.js';
|
||||
import { determineBriefClass, stagesForClass } from './brief-classifier.js';
|
||||
import { ForgeCapabilityError, providerErrorCode } from './errors.js';
|
||||
import {
|
||||
blockedGateResults,
|
||||
evaluateStageGates,
|
||||
isCapabilityGate,
|
||||
simulatedGateResults,
|
||||
waitingGateResults,
|
||||
} from './outcomes.js';
|
||||
import { mapStageToTask } from './stage-adapter.js';
|
||||
import { createSimulatedExecutor } from './simulated-executor.js';
|
||||
import type {
|
||||
ForgeTask,
|
||||
ForgeTaskResult,
|
||||
PipelineOptions,
|
||||
PipelineResult,
|
||||
RunManifest,
|
||||
RunMode,
|
||||
StageStatus,
|
||||
TaskExecutor,
|
||||
} from './types.js';
|
||||
|
||||
/** Reason stamped on stages that complete under explicit simulation. */
|
||||
const SIMULATED_STAGE_REASON =
|
||||
'simulated execution (--simulate): stage was not executed by a real executor';
|
||||
|
||||
/**
|
||||
* Generate a timestamp-based run ID.
|
||||
*/
|
||||
@@ -47,6 +62,7 @@ function createManifest(opts: {
|
||||
briefClass: RunManifest['briefClass'];
|
||||
classSource: RunManifest['classSource'];
|
||||
forceBoard: boolean;
|
||||
mode: RunMode;
|
||||
runDir: string;
|
||||
}): RunManifest {
|
||||
const ts = nowISO();
|
||||
@@ -57,6 +73,7 @@ function createManifest(opts: {
|
||||
briefClass: opts.briefClass,
|
||||
classSource: opts.classSource,
|
||||
forceBoard: opts.forceBoard,
|
||||
mode: opts.mode,
|
||||
createdAt: ts,
|
||||
updatedAt: ts,
|
||||
currentStage: '',
|
||||
@@ -108,20 +125,199 @@ export function selectStages(stages?: string[], skipTo?: string): string[] {
|
||||
return selected.slice(skipIndex);
|
||||
}
|
||||
|
||||
/**
|
||||
* Fail closed when the required executor capability is missing (SDLC-D-035).
|
||||
*/
|
||||
function requireExecutor(executor: TaskExecutor | undefined, simulate: boolean): TaskExecutor {
|
||||
if (executor) return executor;
|
||||
if (simulate) return createSimulatedExecutor({ log: false });
|
||||
throw new ForgeCapabilityError(
|
||||
'FORGE_NO_EXECUTOR',
|
||||
'task-executor',
|
||||
'no task executor is wired; refusing to run the pipeline with a stub executor (fail closed). ' +
|
||||
'Pass --simulate to opt into explicitly simulated execution.',
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Pre-flight a stage's gates in normal mode (fail closed, SDLC-D-035).
|
||||
*
|
||||
* - authority gates: record a typed `waiting-for-authority` stage result and
|
||||
* raise FORGE_AUTHORITY_REQUIRED — approval-based gates never pass vacuously.
|
||||
* - provider gates: record a typed `blocked` stage result and raise the typed
|
||||
* capability error for the missing provider.
|
||||
*
|
||||
* Returns the stage status to record when the pre-flight blocks, or undefined
|
||||
* when the stage may proceed.
|
||||
*/
|
||||
function preflightStageGates(
|
||||
stageName: string,
|
||||
manifest: RunManifest,
|
||||
): { status: StageStatus; error: ForgeCapabilityError } | undefined {
|
||||
const spec = STAGE_SPECS[stageName];
|
||||
if (!spec) throw new Error(`Unknown Forge stage: ${stageName}`);
|
||||
|
||||
for (const gate of spec.qualityGates) {
|
||||
if (!isCapabilityGate(gate)) continue;
|
||||
|
||||
const startedAt = manifest.stages[stageName]?.startedAt;
|
||||
const completedAt = nowISO();
|
||||
|
||||
if (gate.kind === 'authority') {
|
||||
const reason = `gate '${gate.capability}' requires authority sign-off; no mechanical implementation exists (${gate.reason})`;
|
||||
return {
|
||||
status: {
|
||||
status: 'waiting-for-authority',
|
||||
reason,
|
||||
startedAt,
|
||||
completedAt,
|
||||
gateResults: waitingGateResults(spec.qualityGates, reason),
|
||||
},
|
||||
error: new ForgeCapabilityError(
|
||||
'FORGE_AUTHORITY_REQUIRED',
|
||||
gate.capability,
|
||||
`stage '${stageName}' is blocked on authority gate '${gate.capability}': ${gate.reason}. ` +
|
||||
'The pipeline fails closed instead of passing vacuously. Record the approval out-of-band ' +
|
||||
'or run with --simulate for explicitly simulated execution.',
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
const reason = `gate '${gate.capability}' requires provider '${gate.capability}' and none is wired (${gate.reason})`;
|
||||
return {
|
||||
status: {
|
||||
status: 'blocked',
|
||||
reason,
|
||||
startedAt,
|
||||
completedAt,
|
||||
gateResults: blockedGateResults(spec.qualityGates, reason),
|
||||
},
|
||||
error: new ForgeCapabilityError(
|
||||
providerErrorCode(gate.capability),
|
||||
gate.capability,
|
||||
`stage '${stageName}' requires provider '${gate.capability}' which is not wired: ${gate.reason}. ` +
|
||||
'The pipeline fails closed instead of passing vacuously.',
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
return undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Execute the given stage tasks sequentially, updating the manifest.
|
||||
*
|
||||
* Normal mode requires a real executor and evaluates every declared command
|
||||
* gate through the typed outcome model; any non-verified result fails closed.
|
||||
* Simulate mode types every stage and gate result as `simulated`.
|
||||
*/
|
||||
async function executeStages(opts: {
|
||||
manifest: RunManifest;
|
||||
runDir: string;
|
||||
tasks: ForgeTask[];
|
||||
stageNames: string[];
|
||||
executor: TaskExecutor;
|
||||
simulate: boolean;
|
||||
}): Promise<void> {
|
||||
const { manifest, runDir, tasks, stageNames, executor, simulate } = opts;
|
||||
|
||||
for (let i = 0; i < tasks.length; i++) {
|
||||
const task = tasks[i]!;
|
||||
const stageName = stageNames[i]!;
|
||||
const spec = STAGE_SPECS[stageName];
|
||||
if (!spec) throw new Error(`Unknown Forge stage: ${stageName}`);
|
||||
|
||||
// Update manifest: stage in progress
|
||||
manifest.currentStage = stageName;
|
||||
manifest.stages[stageName] = {
|
||||
status: 'in_progress',
|
||||
startedAt: nowISO(),
|
||||
};
|
||||
saveManifest(runDir, manifest);
|
||||
|
||||
// Fail-closed pre-flight (normal mode only): authority/provider gates have
|
||||
// no mechanical implementation and must never pass vacuously.
|
||||
if (!simulate) {
|
||||
const blocked = preflightStageGates(stageName, manifest);
|
||||
if (blocked) {
|
||||
manifest.stages[stageName] = blocked.status;
|
||||
manifest.status =
|
||||
blocked.status.status === 'waiting-for-authority' ? 'waiting-for-authority' : 'failed';
|
||||
saveManifest(runDir, manifest);
|
||||
throw blocked.error;
|
||||
}
|
||||
}
|
||||
|
||||
let result: ForgeTaskResult;
|
||||
try {
|
||||
await executor.submitTask(task);
|
||||
result = await executor.waitForCompletion(task.id, task.timeoutSeconds * 1000);
|
||||
} catch (error) {
|
||||
// Process errors (including timeouts) map to the fail-closed `error` outcome.
|
||||
const reason = error instanceof Error ? error.message : String(error);
|
||||
manifest.stages[stageName] = {
|
||||
status: 'error',
|
||||
reason: `executor error: ${reason}`,
|
||||
startedAt: manifest.stages[stageName]?.startedAt,
|
||||
completedAt: nowISO(),
|
||||
gateResults: [],
|
||||
};
|
||||
manifest.status = 'failed';
|
||||
saveManifest(runDir, manifest);
|
||||
throw error instanceof Error ? error : new Error(reason);
|
||||
}
|
||||
|
||||
if (simulate) {
|
||||
manifest.stages[stageName] = {
|
||||
status: 'simulated',
|
||||
reason: SIMULATED_STAGE_REASON,
|
||||
startedAt: manifest.stages[stageName]?.startedAt,
|
||||
completedAt: nowISO(),
|
||||
gateResults: simulatedGateResults(spec.qualityGates),
|
||||
};
|
||||
saveManifest(runDir, manifest);
|
||||
continue;
|
||||
}
|
||||
|
||||
const evaluation = evaluateStageGates(stageName, spec.qualityGates, result);
|
||||
manifest.stages[stageName] = {
|
||||
status: evaluation.outcome,
|
||||
reason: evaluation.reason,
|
||||
startedAt: manifest.stages[stageName]?.startedAt,
|
||||
completedAt: nowISO(),
|
||||
gateResults: evaluation.gateResults,
|
||||
};
|
||||
|
||||
if (evaluation.outcome !== 'passed') {
|
||||
manifest.status =
|
||||
evaluation.outcome === 'waiting-for-authority' ? 'waiting-for-authority' : 'failed';
|
||||
saveManifest(runDir, manifest);
|
||||
throw new Error(`Stage ${stageName} ${evaluation.outcome}: ${evaluation.reason}`);
|
||||
}
|
||||
|
||||
saveManifest(runDir, manifest);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Run the Forge pipeline.
|
||||
*
|
||||
* 1. Classify the brief
|
||||
* 2. Generate a run ID and create run directory
|
||||
* 3. Map stages to tasks and submit to TaskExecutor
|
||||
* 4. Track manifest with stage statuses
|
||||
* 5. Return pipeline result
|
||||
* 1. Fail closed unless a real executor is wired or simulation is explicit
|
||||
* 2. Classify the brief
|
||||
* 3. Generate a run ID and create run directory
|
||||
* 4. Map stages to tasks and submit to TaskExecutor
|
||||
* 5. Track manifest with typed stage outcomes
|
||||
* 6. Return pipeline result
|
||||
*/
|
||||
export async function runPipeline(
|
||||
briefPath: string,
|
||||
projectRoot: string,
|
||||
options: PipelineOptions,
|
||||
): Promise<PipelineResult> {
|
||||
const simulate = options.simulate ?? false;
|
||||
const executor = requireExecutor(options.executor, simulate);
|
||||
const mode: RunMode = simulate ? 'simulated' : 'normal';
|
||||
|
||||
const resolvedRoot = path.resolve(projectRoot);
|
||||
const resolvedBrief = path.resolve(briefPath);
|
||||
const briefContent = fs.readFileSync(resolvedBrief, 'utf-8');
|
||||
@@ -146,6 +342,7 @@ export async function runPipeline(
|
||||
briefClass,
|
||||
classSource,
|
||||
forceBoard: options.forceBoard ?? false,
|
||||
mode,
|
||||
runDir,
|
||||
});
|
||||
|
||||
@@ -172,54 +369,10 @@ export async function runPipeline(
|
||||
}
|
||||
|
||||
// Execute stages
|
||||
const { executor } = options;
|
||||
for (let i = 0; i < tasks.length; i++) {
|
||||
const task = tasks[i]!;
|
||||
const stageName = selectedStages[i]!;
|
||||
await executeStages({ manifest, runDir, tasks, stageNames: selectedStages, executor, simulate });
|
||||
|
||||
// Update manifest: stage in progress
|
||||
manifest.currentStage = stageName;
|
||||
manifest.stages[stageName] = {
|
||||
status: 'in_progress',
|
||||
startedAt: nowISO(),
|
||||
};
|
||||
saveManifest(runDir, manifest);
|
||||
|
||||
try {
|
||||
await executor.submitTask(task);
|
||||
const result = await executor.waitForCompletion(task.id, task.timeoutSeconds * 1000);
|
||||
|
||||
// Update manifest: stage completed or failed
|
||||
const stageStatus: StageStatus = {
|
||||
status: result.status === 'completed' ? 'passed' : 'failed',
|
||||
startedAt: manifest.stages[stageName]!.startedAt,
|
||||
completedAt: nowISO(),
|
||||
};
|
||||
manifest.stages[stageName] = stageStatus;
|
||||
|
||||
if (result.status !== 'completed') {
|
||||
manifest.status = 'failed';
|
||||
saveManifest(runDir, manifest);
|
||||
throw new Error(`Stage ${stageName} failed with status: ${result.status}`);
|
||||
}
|
||||
|
||||
saveManifest(runDir, manifest);
|
||||
} catch (error) {
|
||||
if (!manifest.stages[stageName]?.completedAt) {
|
||||
manifest.stages[stageName] = {
|
||||
status: 'failed',
|
||||
startedAt: manifest.stages[stageName]?.startedAt,
|
||||
completedAt: nowISO(),
|
||||
};
|
||||
}
|
||||
manifest.status = 'failed';
|
||||
saveManifest(runDir, manifest);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
// All stages passed
|
||||
manifest.status = 'completed';
|
||||
// All stages reached a terminal state for this mode
|
||||
manifest.status = simulate ? 'simulated' : 'completed';
|
||||
saveManifest(runDir, manifest);
|
||||
|
||||
return {
|
||||
@@ -234,22 +387,30 @@ export async function runPipeline(
|
||||
}
|
||||
|
||||
/**
|
||||
* Resume a pipeline from the last incomplete stage.
|
||||
* Resume a pipeline from the last non-passed stage.
|
||||
*/
|
||||
export async function resumePipeline(
|
||||
runDir: string,
|
||||
executor: TaskExecutor,
|
||||
executor?: TaskExecutor,
|
||||
options?: { simulate?: boolean },
|
||||
): Promise<PipelineResult> {
|
||||
const simulate = options?.simulate ?? false;
|
||||
const wiredExecutor = requireExecutor(executor, simulate);
|
||||
const mode: RunMode = simulate ? 'simulated' : 'normal';
|
||||
|
||||
const manifest = loadManifest(runDir);
|
||||
const resolvedRoot = path.dirname(path.dirname(path.dirname(runDir))); // .forge/runs/{id} → project root
|
||||
|
||||
const briefContent = fs.readFileSync(manifest.brief, 'utf-8');
|
||||
const allStages = stagesForClass(manifest.briefClass, manifest.forceBoard);
|
||||
|
||||
// Find first non-passed stage
|
||||
manifest.mode = mode;
|
||||
|
||||
// Find first non-satisfying stage (only a verified `passed` counts as done;
|
||||
// simulated and waiting-for-authority stages are re-run).
|
||||
const resumeFrom = allStages.find((s) => manifest.stages[s]?.status !== 'passed');
|
||||
if (!resumeFrom) {
|
||||
manifest.status = 'completed';
|
||||
manifest.status = mode === 'simulated' ? 'simulated' : 'completed';
|
||||
saveManifest(runDir, manifest);
|
||||
return {
|
||||
runId: manifest.runId,
|
||||
@@ -284,49 +445,16 @@ export async function resumePipeline(
|
||||
tasks.push(task);
|
||||
}
|
||||
|
||||
for (let i = 0; i < tasks.length; i++) {
|
||||
const task = tasks[i]!;
|
||||
const stageName = remainingStages[i]!;
|
||||
await executeStages({
|
||||
manifest,
|
||||
runDir,
|
||||
tasks,
|
||||
stageNames: remainingStages,
|
||||
executor: wiredExecutor,
|
||||
simulate,
|
||||
});
|
||||
|
||||
manifest.currentStage = stageName;
|
||||
manifest.stages[stageName] = {
|
||||
status: 'in_progress',
|
||||
startedAt: nowISO(),
|
||||
};
|
||||
saveManifest(runDir, manifest);
|
||||
|
||||
try {
|
||||
await executor.submitTask(task);
|
||||
const result = await executor.waitForCompletion(task.id, task.timeoutSeconds * 1000);
|
||||
|
||||
manifest.stages[stageName] = {
|
||||
status: result.status === 'completed' ? 'passed' : 'failed',
|
||||
startedAt: manifest.stages[stageName]!.startedAt,
|
||||
completedAt: nowISO(),
|
||||
};
|
||||
|
||||
if (result.status !== 'completed') {
|
||||
manifest.status = 'failed';
|
||||
saveManifest(runDir, manifest);
|
||||
throw new Error(`Stage ${stageName} failed with status: ${result.status}`);
|
||||
}
|
||||
|
||||
saveManifest(runDir, manifest);
|
||||
} catch (error) {
|
||||
if (!manifest.stages[stageName]?.completedAt) {
|
||||
manifest.stages[stageName] = {
|
||||
status: 'failed',
|
||||
startedAt: manifest.stages[stageName]?.startedAt,
|
||||
completedAt: nowISO(),
|
||||
};
|
||||
}
|
||||
manifest.status = 'failed';
|
||||
saveManifest(runDir, manifest);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
manifest.status = 'completed';
|
||||
manifest.status = simulate ? 'simulated' : 'completed';
|
||||
saveManifest(runDir, manifest);
|
||||
|
||||
return {
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
import type { ForgeTask, ForgeTaskResult, TaskExecutor } from './types.js';
|
||||
|
||||
/**
|
||||
* Simulated executor — used ONLY when the caller explicitly passes --simulate.
|
||||
*
|
||||
* It submits no real work and returns typed `simulated` results so a simulated
|
||||
* run can never be confused with a verified one. In normal mode (no --simulate)
|
||||
* the CLI refuses to run at all with FORGE_NO_EXECUTOR instead of wiring this
|
||||
* stub (fail closed, SDLC-D-035).
|
||||
*/
|
||||
export function createSimulatedExecutor(options?: { log?: boolean }): TaskExecutor {
|
||||
const log = options?.log ?? true;
|
||||
return {
|
||||
async submitTask(task: ForgeTask) {
|
||||
if (log) console.log(` [forge:simulated] stage submitted: ${task.id} (${task.title})`);
|
||||
},
|
||||
async waitForCompletion(taskId: string): Promise<ForgeTaskResult> {
|
||||
if (log) console.log(` [forge:simulated] stage complete: ${taskId}`);
|
||||
return {
|
||||
task_id: taskId,
|
||||
outcome: 'simulated',
|
||||
reason: 'no executor wired; simulated execution requested via --simulate',
|
||||
completed_at: new Date().toISOString(),
|
||||
exit_code: 0,
|
||||
gate_results: [],
|
||||
};
|
||||
},
|
||||
async getTaskStatus() {
|
||||
return 'completed' as const;
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -1,4 +1,4 @@
|
||||
import type { GateEntry, TaskResult } from '@mosaicstack/macp';
|
||||
import type { GateEntry } from '@mosaicstack/macp';
|
||||
|
||||
/** Stage dispatch mode. */
|
||||
export type StageDispatch = 'exec' | 'yolo' | 'pi';
|
||||
@@ -6,6 +6,58 @@ export type StageDispatch = 'exec' | 'yolo' | 'pi';
|
||||
/** Stage type — determines agent selection and gate requirements. */
|
||||
export type StageType = 'research' | 'review' | 'coding' | 'deploy';
|
||||
|
||||
/**
|
||||
* Typed outcome for every gate and stage evaluation — closed set (SDLC-D-035).
|
||||
*
|
||||
* Only `passed` means "verified by a real implementation". `simulated` is
|
||||
* produced exclusively in explicit `--simulate` runs and is never satisfying.
|
||||
*/
|
||||
export type ForgeOutcome =
|
||||
| 'passed'
|
||||
| 'failed'
|
||||
| 'blocked'
|
||||
| 'error'
|
||||
| 'waiting-for-authority'
|
||||
| 'simulated'
|
||||
| 'not-applicable';
|
||||
|
||||
/** A gate that requires authority (human/board) sign-off; no mechanical command can satisfy it. */
|
||||
export interface AuthorityGate {
|
||||
kind: 'authority';
|
||||
capability: string;
|
||||
reason: string;
|
||||
}
|
||||
|
||||
/** A gate that requires a wired provider (e.g. an AI reviewer, CI pipeline) to evaluate. */
|
||||
export interface ProviderGate {
|
||||
kind: 'provider';
|
||||
capability: string;
|
||||
reason: string;
|
||||
}
|
||||
|
||||
/** Forge quality gate: a real command, an authority sign-off, or a provider-backed check. */
|
||||
export type ForgeGate = string | GateEntry | AuthorityGate | ProviderGate;
|
||||
|
||||
/** Typed result of evaluating a single quality gate. */
|
||||
export interface ForgeGateResult {
|
||||
gate: string;
|
||||
outcome: ForgeOutcome;
|
||||
reason: string;
|
||||
exitCode?: number;
|
||||
output?: string;
|
||||
timedOut?: boolean;
|
||||
}
|
||||
|
||||
/** Typed result of a task/stage execution returned by a TaskExecutor. */
|
||||
export interface ForgeTaskResult {
|
||||
task_id: string;
|
||||
outcome: ForgeOutcome;
|
||||
reason: string;
|
||||
completed_at: string;
|
||||
exit_code: number;
|
||||
gate_results: ForgeGateResult[];
|
||||
}
|
||||
|
||||
/** Stage specification — defines a single pipeline stage. */
|
||||
export interface StageSpec {
|
||||
number: string;
|
||||
@@ -14,7 +66,7 @@ export interface StageSpec {
|
||||
type: StageType;
|
||||
gate: string;
|
||||
promptFile: string;
|
||||
qualityGates: (string | GateEntry)[];
|
||||
qualityGates: ForgeGate[];
|
||||
}
|
||||
|
||||
/** Brief classification. */
|
||||
@@ -25,11 +77,18 @@ export type ClassSource = 'cli' | 'frontmatter' | 'auto';
|
||||
|
||||
/** Per-stage status within a run manifest. */
|
||||
export interface StageStatus {
|
||||
status: 'pending' | 'in_progress' | 'passed' | 'failed';
|
||||
status: 'pending' | 'in_progress' | ForgeOutcome;
|
||||
/** Why the stage reached its current (terminal) outcome, when applicable. */
|
||||
reason?: string;
|
||||
startedAt?: string;
|
||||
completedAt?: string;
|
||||
/** Typed per-gate results recorded alongside the stage outcome. */
|
||||
gateResults?: ForgeGateResult[];
|
||||
}
|
||||
|
||||
/** Execution mode of a run. */
|
||||
export type RunMode = 'normal' | 'simulated';
|
||||
|
||||
/** Run manifest — persisted to disk as manifest.json. */
|
||||
export interface RunManifest {
|
||||
runId: string;
|
||||
@@ -38,10 +97,23 @@ export interface RunManifest {
|
||||
briefClass: BriefClass;
|
||||
classSource: ClassSource;
|
||||
forceBoard: boolean;
|
||||
/**
|
||||
* Execution mode. `simulated` runs stub execution; their results are typed
|
||||
* `simulated` and must never be read as verified success. Optional because
|
||||
* manifests written before this field existed default to `normal`.
|
||||
*/
|
||||
mode?: RunMode;
|
||||
createdAt: string;
|
||||
updatedAt: string;
|
||||
currentStage: string;
|
||||
status: 'in_progress' | 'completed' | 'failed' | 'interrupted' | 'rejected';
|
||||
status:
|
||||
| 'in_progress'
|
||||
| 'completed'
|
||||
| 'failed'
|
||||
| 'interrupted'
|
||||
| 'rejected'
|
||||
| 'simulated'
|
||||
| 'waiting-for-authority';
|
||||
stages: Record<string, StageStatus>;
|
||||
}
|
||||
|
||||
@@ -65,7 +137,7 @@ export interface ForgeTask {
|
||||
briefPath: string;
|
||||
resultPath: string;
|
||||
timeoutSeconds: number;
|
||||
qualityGates: (string | GateEntry)[];
|
||||
qualityGates: ForgeGate[];
|
||||
worktree?: string;
|
||||
command?: string;
|
||||
dependsOn?: string[];
|
||||
@@ -76,7 +148,7 @@ export interface ForgeTask {
|
||||
/** Abstract task executor — decouples from packages/coord. */
|
||||
export interface TaskExecutor {
|
||||
submitTask(task: ForgeTask): Promise<void>;
|
||||
waitForCompletion(taskId: string, timeoutMs: number): Promise<TaskResult>;
|
||||
waitForCompletion(taskId: string, timeoutMs: number): Promise<ForgeTaskResult>;
|
||||
getTaskStatus(taskId: string): Promise<ForgeTaskStatus>;
|
||||
}
|
||||
|
||||
@@ -122,7 +194,16 @@ export interface PipelineOptions {
|
||||
stages?: string[];
|
||||
skipTo?: string;
|
||||
dryRun?: boolean;
|
||||
executor: TaskExecutor;
|
||||
/**
|
||||
* Real task executor. Required in normal mode: the pipeline fails closed
|
||||
* with FORGE_NO_EXECUTOR when it is absent.
|
||||
*/
|
||||
executor?: TaskExecutor;
|
||||
/**
|
||||
* Explicit opt-in to simulated execution. Every stage and gate result is
|
||||
* typed `simulated` and is never satisfying.
|
||||
*/
|
||||
simulate?: boolean;
|
||||
}
|
||||
|
||||
/** Pipeline run result. */
|
||||
|
||||
@@ -233,36 +233,8 @@ assert_owned_tmux_server() {
|
||||
fail "tmux server ownership or environment validation failed"
|
||||
}
|
||||
|
||||
# Lease-broker socket preflight (#1292). The gated runtime (`mosaic yolo …` →
|
||||
# launch-runtime.py) registers with the broker or dies ~4 seconds in, with the
|
||||
# diagnostic invisible because tmux destroys the dead pane. This check runs
|
||||
# BEFORE any tmux effect — including the ownership probe below — so a host
|
||||
# without a broker produces a named, surviving refusal instead of a doomed
|
||||
# pane. Exit 75 (EX_TEMPFAIL), distinct from 64 (bad projection) and 69 (host
|
||||
# not ready for other reasons); the agent@ unit is Type=oneshot with no
|
||||
# Restart=, so the failed unit keeps its message instead of looping. Socket
|
||||
# resolution matches launch.ts's defaultLeaseBrokerSocket precedence exactly.
|
||||
# This preflight DETECTS and REFUSES — it never starts the broker (activation
|
||||
# belongs to the fleet control plane; a component that both detects and fixes
|
||||
# cannot be used to measure whether the fix worked).
|
||||
broker_socket_path() {
|
||||
if [ -n "${MOSAIC_LEASE_BROKER_SOCKET:-}" ]; then
|
||||
printf '%s\n' "$MOSAIC_LEASE_BROKER_SOCKET"
|
||||
return 0
|
||||
fi
|
||||
local runtime_dir="${XDG_RUNTIME_DIR:-/run/user/$(id -u)}"
|
||||
printf '%s\n' "${runtime_dir}/mosaic-lease/broker.sock"
|
||||
}
|
||||
|
||||
if [ "$MODE" = "launch" ]; then
|
||||
_broker_socket=$(broker_socket_path)
|
||||
if [ ! -S "$_broker_socket" ]; then
|
||||
echo "[fleet] FAIL_LAUNCH broker-absent: lease broker socket ${_broker_socket} missing; runtime launch denied (#1292)." >&2
|
||||
echo "[fleet] remedy: systemctl --user enable --now mosaic-lease-broker.service (or reinstall via: mosaic fleet install)" >&2
|
||||
exit 75
|
||||
fi
|
||||
fi
|
||||
|
||||
# Validate exact server ownership before querying, cleaning, or creating any
|
||||
# managed session. An unmanaged or contaminated named socket is never repaired.
|
||||
assert_owned_tmux_server
|
||||
|
||||
if [ "$MODE" = interaction ]; then
|
||||
|
||||
@@ -1,216 +0,0 @@
|
||||
#!/usr/bin/env bash
|
||||
# CI-fit regression suite for the #1292 lease-broker socket preflight in
|
||||
# start-agent-session.sh.
|
||||
#
|
||||
# WHY THIS SUITE IS CI-FIT WHERE test-start-agent-session.sh IS NOT (#1017/#1270
|
||||
# context): that older suite's precondition is "the host does not have the pi
|
||||
# binary", which a CI image that ships pi violates — its guard correctly
|
||||
# refuses to report a pass there, so it is excluded from the chain. THIS suite
|
||||
# controls its own preconditions instead of inheriting them from the host: a
|
||||
# fake tmux on PATH, a fake mosaic on PATH, a real unix socket created in a
|
||||
# tmpdir, a hermetic env (env -i, fake HOME, GIT_CONFIG_GLOBAL severed). It
|
||||
# never depends on what the host has installed, so a green here means the same
|
||||
# thing on every host. Anyone adding cases: keep that property — no case may
|
||||
# depend on host state.
|
||||
#
|
||||
# The failure this suite is written down to catch (#1292): a seat launched on a
|
||||
# host with no lease broker dies ~4 seconds in at registration, with the
|
||||
# diagnostic invisible because tmux destroys the dead pane. The preflight runs
|
||||
# BEFORE any tmux effect and refuses with a NAMED code (exit 75, EX_TEMPFAIL)
|
||||
# so the message survives. The agent@ unit is Type=oneshot with no Restart=,
|
||||
# so a failed unit keeps its output instead of looping.
|
||||
#
|
||||
# Cases:
|
||||
# 1. absent socket -> exit 75, message names broker-absent + socket path +
|
||||
# remedy, and NO tmux session was ever created (the doomed-pane half).
|
||||
# 2. present socket (real unix socket in tmpdir) -> proceeds PAST the
|
||||
# preflight (the suite then stops at the next precondition, proving the
|
||||
# preflight was not the refusal).
|
||||
# 3. explicit MOSAIC_LEASE_BROKER_SOCKET wins over XDG_RUNTIME_DIR default.
|
||||
# 4. --stop mode does NOT require the broker (teardown must not be fenced on
|
||||
# a component whose absence is exactly what teardown may follow).
|
||||
#
|
||||
# Sabotage control, run by the developer (not in-suite): remove the preflight
|
||||
# block from start-agent-session.sh, re-run — case 1 fails (a tmux session is
|
||||
# created / exit is not 75), cases 2-4 still pass; restore byte-identically.
|
||||
|
||||
set -euo pipefail
|
||||
|
||||
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
|
||||
WORK_DIR="${MOSAIC_TEST_WORK_DIR:-$PWD/.mosaic-test-work/agent-session-broker-preflight}"
|
||||
FAKE_HOME="$WORK_DIR/home"
|
||||
BIN_DIR="$WORK_DIR/bin"
|
||||
ENV_DIR="$WORK_DIR/env"
|
||||
SOCK_DIR="$WORK_DIR/sockets"
|
||||
LOG_FILE="$WORK_DIR/tmux-calls.log"
|
||||
|
||||
rm -rf "$WORK_DIR"
|
||||
# The script asserts a managed directory tree under MOSAIC_HOME: mosaic/,
|
||||
# mosaic/fleet/, mosaic/fleet/agents/ — private (0700/0750-style) modes, no
|
||||
# symlinks — plus a per-agent env projection. Build the full tree the launcher
|
||||
# expects so the suite reaches the BROKER preflight rather than dying at
|
||||
# environment validation.
|
||||
mkdir -p "$FAKE_HOME/.config/mosaic/fleet/agents" "$BIN_DIR" "$SOCK_DIR"
|
||||
chmod 700 "$FAKE_HOME/.config/mosaic" "$FAKE_HOME/.config/mosaic/fleet/agents"
|
||||
chmod 750 "$FAKE_HOME/.config/mosaic/fleet"
|
||||
cat > "$FAKE_HOME/.config/mosaic/fleet/agents/preflight-test.env.generated" <<'ENVEOF'
|
||||
MOSAIC_AGENT_NAME=preflight-test
|
||||
MOSAIC_AGENT_CLASS=worker
|
||||
MOSAIC_AGENT_RUNTIME=pi
|
||||
MOSAIC_AGENT_MODEL=
|
||||
MOSAIC_AGENT_REASONING=
|
||||
MOSAIC_AGENT_TOOL_POLICY=code
|
||||
MOSAIC_AGENT_WORKDIR=/tmp
|
||||
MOSAIC_TMUX_SOCKET=mosaic-fleet
|
||||
ENVEOF
|
||||
chmod 600 "$FAKE_HOME/.config/mosaic/fleet/agents/preflight-test.env.generated"
|
||||
|
||||
# ─── Fake tmux: records every invocation; new-session marks the marker. ────
|
||||
: > "$LOG_FILE"
|
||||
cat > "$BIN_DIR/tmux" <<SH
|
||||
#!/usr/bin/env bash
|
||||
printf 'tmux %s\n' "\$*" >> "$LOG_FILE"
|
||||
if [[ "\$*" == *new-session* ]]; then
|
||||
echo "TMUX-NEW-SESSION-INVOKED" >> "$LOG_FILE"
|
||||
fi
|
||||
exit 0
|
||||
SH
|
||||
chmod +x "$BIN_DIR/tmux"
|
||||
|
||||
# ─── Fake mosaic/pi binaries so the script proceeds past its own lookups. ───
|
||||
for bin in mosaic pi claude; do
|
||||
printf '#!/usr/bin/env bash\nexit 0\n' > "$BIN_DIR/$bin"
|
||||
chmod +x "$BIN_DIR/$bin"
|
||||
done
|
||||
|
||||
# ─── Minimal launch environment the script expects. ────────────────────────
|
||||
# (Enough for the preflight to be reached; later stages will still fail in
|
||||
# case 2 — that is expected and asserted.)
|
||||
run_session_script() {
|
||||
local mode="$1"; shift
|
||||
(
|
||||
cd "$WORK_DIR"
|
||||
env -i HOME="$FAKE_HOME" PATH="$BIN_DIR:/usr/bin:/bin" \
|
||||
GIT_CONFIG_GLOBAL=/dev/null GIT_CONFIG_SYSTEM=/dev/null \
|
||||
MOSAIC_HOME="$FAKE_HOME/.config/mosaic" \
|
||||
AGENT_NAME=preflight-test \
|
||||
"$@" \
|
||||
bash "$SCRIPT_DIR/start-agent-session.sh" $mode preflight-test
|
||||
)
|
||||
}
|
||||
|
||||
fail=0
|
||||
assert() {
|
||||
local desc="$1" expected="$2" actual="$3"
|
||||
if [[ "$expected" != "$actual" ]]; then
|
||||
echo "FAIL: $desc — expected '$expected', got '$actual'" >&2
|
||||
fail=1
|
||||
fi
|
||||
}
|
||||
assert_contains() {
|
||||
local desc="$1" haystack="$2" needle="$3"
|
||||
[[ "$haystack" == *"$needle"* ]] || { echo "FAIL: $desc — missing '$needle' in: $haystack" >&2; fail=1; }
|
||||
}
|
||||
assert_not_contains() {
|
||||
local desc="$1" haystack="$2" needle="$3"
|
||||
if [[ "$haystack" == *"$needle"* ]]; then
|
||||
echo "FAIL: $desc — must not contain '$needle'" >&2
|
||||
fail=1
|
||||
fi
|
||||
return 0
|
||||
}
|
||||
|
||||
# ─── 1. Absent socket → named refusal, NO tmux session. ────────────────────
|
||||
: > "$LOG_FILE"
|
||||
stderr_file="$WORK_DIR/stderr-1.tmp"
|
||||
set +e
|
||||
out=$(run_session_script "" MOSAIC_LEASE_BROKER_SOCKET="$SOCK_DIR/absent.sock" 2>"$stderr_file")
|
||||
rc=$?
|
||||
set -e
|
||||
assert "absent socket exit code" "75" "$rc"
|
||||
err=$(cat "$stderr_file")
|
||||
assert_contains "absent socket names the failure" "$err" "FAIL_LAUNCH broker-absent"
|
||||
assert_contains "absent socket names the socket path" "$err" "$SOCK_DIR/absent.sock"
|
||||
assert_contains "absent socket names a remedy" "$err" "mosaic fleet install"
|
||||
log1=$(cat "$LOG_FILE")
|
||||
assert_not_contains "absent socket must not create a tmux session" "$log1" "TMUX-NEW-SESSION-INVOKED"
|
||||
|
||||
# ─── 2. Present socket → passes the preflight. ─────────────────────────────
|
||||
# Expected: ownership/env checks AFTER the preflight may refuse (fixture is
|
||||
# minimal by design); the assertion is only that the refusal is NOT
|
||||
# broker-absent and the exit is NOT 75.
|
||||
# Create a REAL unix socket: a detached python holder binds it and stays alive
|
||||
# for the duration (bash cannot create sockets; a foreground python would
|
||||
# close the socket on exit and -S on a closed-but-unlinked path fails). Written
|
||||
# as a script file + setsid nohup so no job-control/heredoc interaction with
|
||||
# set -e can silently kill the suite.
|
||||
# AF_UNIX binds cap at 108 path bytes; the suite's workdir exceeds that, so
|
||||
# the live socket lives at a SHORT path under /tmp (unique per run, cleaned
|
||||
# with the suite). The preflight takes its socket path explicitly, so this
|
||||
# stays fully controlled.
|
||||
LIVE_SOCK=$(mktemp -u /tmp/mosaic-preflight-XXXXXX.sock)
|
||||
trap 'rm -f "$LIVE_SOCK"' EXIT
|
||||
rm -f "$SOCK_DIR/live.sock" "$LIVE_SOCK"
|
||||
cat > "$SOCK_DIR/holder.py" <<'PY'
|
||||
import socket, sys, time
|
||||
path = sys.argv[1]
|
||||
s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
s.bind(path)
|
||||
s.listen(1)
|
||||
time.sleep(120)
|
||||
PY
|
||||
python3 "$SOCK_DIR/holder.py" "$LIVE_SOCK" >/dev/null 2>"$SOCK_DIR/holder.err" &
|
||||
HOLDER_PID=$!
|
||||
# Wait for the socket object to exist (bind is near-instant, but do not race it).
|
||||
for _ in $(seq 1 50); do
|
||||
[ -S "$LIVE_SOCK" ] && break
|
||||
sleep 0.1
|
||||
done
|
||||
if [ ! -S "$LIVE_SOCK" ]; then
|
||||
echo "FAIL: could not create live socket fixture (holder pid $HOLDER_PID)" >&2
|
||||
ps -p "$HOLDER_PID" -o pid,stat,cmd --no-headers >&2 || echo "(holder exited)" >&2
|
||||
cat "$SOCK_DIR/holder.err" >&2 || true
|
||||
exit 1
|
||||
fi
|
||||
: > "$LOG_FILE"
|
||||
set +e
|
||||
out=$(run_session_script "" MOSAIC_LEASE_BROKER_SOCKET="$LIVE_SOCK" 2>"$WORK_DIR/stderr-2.tmp")
|
||||
rc=$?
|
||||
set -e
|
||||
# The preflight PASSED if the failure (whatever later stage refused) is NOT
|
||||
# the broker refusal, and tmux was reached or a later precondition named
|
||||
# something else.
|
||||
err2=$(cat "$WORK_DIR/stderr-2.tmp")
|
||||
assert_not_contains "live socket must not refuse broker-absent" "$err2" "broker-absent"
|
||||
if [[ "$rc" == "75" ]]; then
|
||||
echo "FAIL: live socket — preflight still refused (exit 75) with a live socket" >&2
|
||||
fail=1
|
||||
fi
|
||||
|
||||
# ─── 3. Explicit socket env wins over XDG default. ─────────────────────────
|
||||
set +e
|
||||
out=$(run_session_script "" XDG_RUNTIME_DIR="$SOCK_DIR/no-runtime-here" MOSAIC_LEASE_BROKER_SOCKET="$SOCK_DIR/absent2.sock" 2>"$WORK_DIR/stderr-3.tmp")
|
||||
rc=$?
|
||||
set -e
|
||||
assert "explicit env wins (exit 75)" "75" "$rc"
|
||||
assert_contains "explicit env path named" "$(cat "$WORK_DIR/stderr-3.tmp")" "$SOCK_DIR/absent2.sock"
|
||||
|
||||
# ─── 4. --stop is not fenced on the broker. ────────────────────────────────
|
||||
: > "$LOG_FILE"
|
||||
set +e
|
||||
out=$(run_session_script "--stop" MOSAIC_LEASE_BROKER_SOCKET="$SOCK_DIR/absent3.sock" 2>"$WORK_DIR/stderr-4.tmp")
|
||||
rc=$?
|
||||
set -e
|
||||
err4=$(cat "$WORK_DIR/stderr-4.tmp")
|
||||
assert_not_contains "--stop must not refuse broker-absent" "$err4" "broker-absent"
|
||||
if [[ "$rc" == "75" ]]; then
|
||||
echo "FAIL: --stop — exit 75 means teardown was fenced on the broker" >&2
|
||||
fail=1
|
||||
fi
|
||||
|
||||
kill "$HOLDER_PID" 2>/dev/null || true
|
||||
|
||||
if [[ "$fail" -eq 0 ]]; then
|
||||
echo "start-agent-session lease-broker preflight regression passed"
|
||||
fi
|
||||
exit "$fail"
|
||||
@@ -1,177 +0,0 @@
|
||||
import { lstat, mkdir, mkdtemp, readFile, rm, symlink, writeFile } from 'node:fs/promises';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { afterEach, describe, expect, it } from 'vitest';
|
||||
|
||||
import { placeUnitFile, resolveLeaseBrokerSocketForPreflight } from './fleet.js';
|
||||
|
||||
/**
|
||||
* Unit-placement regression harness for #1292.
|
||||
*
|
||||
* The two measured defects this suite pins:
|
||||
* 1. `systemctl enable <name>` does NOT rewrite an existing by-path
|
||||
* wants-symlink — so placement must remove stale residue explicitly, and
|
||||
* acceptance asserts on the RESULTING SYMLINK TARGET, never on the enable
|
||||
* call's argument (asserting the call cannot see where the link ended up).
|
||||
* 2. Node's copyFile FOLLOWS a by-path symlink at the destination and
|
||||
* overwrites the SEED template. Acceptance asserts on the SEED's bytes
|
||||
* AND mtime — unchanged — which is the only check that can redden for
|
||||
* finding 2. The symlink-target assertion catches finding 1; these are
|
||||
* different defects with different failure modes.
|
||||
*
|
||||
* Fixtures are entirely inside tmpdirs (source template, active systemd dir,
|
||||
* wants dir) — no real host paths are touched by this suite.
|
||||
*/
|
||||
|
||||
describe('placeUnitFile (#1292 unit placement)', () => {
|
||||
const cleanup: string[] = [];
|
||||
afterEach(async () => {
|
||||
while (cleanup.length > 0) {
|
||||
await rm(cleanup.pop()!, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
|
||||
async function fixture() {
|
||||
const root = await mkdtemp(join(tmpdir(), 'place-unit-'));
|
||||
cleanup.push(root);
|
||||
const seedDir = join(root, 'seed');
|
||||
const activeDir = join(root, 'active');
|
||||
await mkdir(seedDir, { recursive: true });
|
||||
await mkdir(activeDir, { recursive: true });
|
||||
const seedTemplate = join(seedDir, 'unit-under-test.service');
|
||||
await writeFile(
|
||||
seedTemplate,
|
||||
'[Unit]\nDescription=seed template\n[Service]\nType=oneshot\nExecStart=/bin/true\n[Install]\nWantedBy=default.target\n',
|
||||
);
|
||||
const activeSource = join(root, 'active-source.service');
|
||||
await writeFile(
|
||||
activeSource,
|
||||
'[Unit]\nDescription=active copy v2\n[Service]\nType=oneshot\nExecStart=/bin/true\n[Install]\nWantedBy=default.target\n',
|
||||
);
|
||||
return { root, seedDir, activeDir, seedTemplate, activeSource };
|
||||
}
|
||||
|
||||
it('places a regular file on a clean host (negative control: no residue anywhere)', async () => {
|
||||
const f = await fixture();
|
||||
const result = await placeUnitFile(f.activeSource, f.activeDir, 'unit-under-test.service');
|
||||
expect(result.unlinkedDestinationSymlink).toBe(false);
|
||||
expect(result.removedStaleWantsSymlink).toBe(false);
|
||||
const info = await lstat(join(f.activeDir, 'unit-under-test.service'));
|
||||
expect(info.isSymbolicLink()).toBe(false);
|
||||
expect(await readFile(join(f.activeDir, 'unit-under-test.service'), 'utf8')).toContain(
|
||||
'active copy v2',
|
||||
);
|
||||
// Seed untouched by construction — but assert it, so the clean-host case
|
||||
// cannot silently regress into seed-mutation.
|
||||
expect(await readFile(f.seedTemplate, 'utf8')).toContain('seed template');
|
||||
});
|
||||
|
||||
it('by-path residue: unlinks destination symlink, places the file, seed bytes AND mtime unchanged (finding 2)', async () => {
|
||||
const f = await fixture();
|
||||
const seedBefore = await readFile(f.seedTemplate, 'utf8');
|
||||
const mtimeBefore = (await lstat(f.seedTemplate)).mtimeMs;
|
||||
// The fomo-lin convention: by-path enable left a symlink AT the unit name
|
||||
// pointing at the seed template, plus a wants-symlink doing the same.
|
||||
await symlink(f.seedTemplate, join(f.activeDir, 'unit-under-test.service'));
|
||||
const wantsDir = join(f.activeDir, 'default.target.wants');
|
||||
await mkdir(wantsDir, { recursive: true });
|
||||
await symlink(f.seedTemplate, join(wantsDir, 'unit-under-test.service'));
|
||||
|
||||
const result = await placeUnitFile(f.activeSource, f.activeDir, 'unit-under-test.service');
|
||||
expect(result.unlinkedDestinationSymlink).toBe(true);
|
||||
expect(result.removedStaleWantsSymlink).toBe(true);
|
||||
|
||||
// FINDING 2's check: the seed is byte-identical and its mtime did not move.
|
||||
expect(await readFile(f.seedTemplate, 'utf8')).toBe(seedBefore);
|
||||
expect((await lstat(f.seedTemplate)).mtimeMs).toBe(mtimeBefore);
|
||||
|
||||
// The destination is now a regular file carrying the ACTIVE content.
|
||||
const destInfo = await lstat(join(f.activeDir, 'unit-under-test.service'));
|
||||
expect(destInfo.isSymbolicLink()).toBe(false);
|
||||
expect(await readFile(join(f.activeDir, 'unit-under-test.service'), 'utf8')).toContain(
|
||||
'active copy v2',
|
||||
);
|
||||
});
|
||||
|
||||
it('by-path residue: no wants-symlink remains pointing at the seed (finding 1 residue cleared)', async () => {
|
||||
const f = await fixture();
|
||||
await symlink(f.seedTemplate, join(f.activeDir, 'unit-under-test.service'));
|
||||
const wantsDir = join(f.activeDir, 'default.target.wants');
|
||||
await mkdir(wantsDir, { recursive: true });
|
||||
await symlink(f.seedTemplate, join(wantsDir, 'unit-under-test.service'));
|
||||
|
||||
await placeUnitFile(f.activeSource, f.activeDir, 'unit-under-test.service');
|
||||
|
||||
// After placement the stale wants link is GONE (enable-by-name recreates
|
||||
// it correctly). A link still present must not point at the seed.
|
||||
try {
|
||||
const link = await lstat(join(wantsDir, 'unit-under-test.service'));
|
||||
if (link.isSymbolicLink()) {
|
||||
const target = await readFile(join(wantsDir, 'unit-under-test.service'), 'utf8').catch(
|
||||
async () => '',
|
||||
);
|
||||
expect(target).not.toContain('seed template');
|
||||
}
|
||||
} catch {
|
||||
// absent wants link — the expected post-placement state
|
||||
}
|
||||
});
|
||||
|
||||
it('idempotence: second placement on a reconciled host is a no-op producing the identical final state', async () => {
|
||||
const f = await fixture();
|
||||
// Reconciled starting state: regular file at the name, wants link to the active copy.
|
||||
await writeFile(
|
||||
join(f.activeDir, 'unit-under-test.service'),
|
||||
await readFile(f.activeSource, 'utf8'),
|
||||
);
|
||||
const wantsDir = join(f.activeDir, 'default.target.wants');
|
||||
await mkdir(wantsDir, { recursive: true });
|
||||
await symlink(
|
||||
join(f.activeDir, 'unit-under-test.service'),
|
||||
join(wantsDir, 'unit-under-test.service'),
|
||||
);
|
||||
const before = await readFile(join(f.activeDir, 'unit-under-test.service'), 'utf8');
|
||||
|
||||
const result = await placeUnitFile(f.activeSource, f.activeDir, 'unit-under-test.service');
|
||||
// No destructive step fired: no unlink, no wants removal.
|
||||
expect(result.unlinkedDestinationSymlink).toBe(false);
|
||||
expect(result.removedStaleWantsSymlink).toBe(false);
|
||||
// Identical final state.
|
||||
expect(await readFile(join(f.activeDir, 'unit-under-test.service'), 'utf8')).toBe(before);
|
||||
const link = await lstat(join(wantsDir, 'unit-under-test.service'));
|
||||
expect(link.isSymbolicLink()).toBe(true);
|
||||
});
|
||||
|
||||
it('double install on by-path residue converges to the identical reconciled state', async () => {
|
||||
const f = await fixture();
|
||||
await symlink(f.seedTemplate, join(f.activeDir, 'unit-under-test.service'));
|
||||
const wantsDir = join(f.activeDir, 'default.target.wants');
|
||||
await mkdir(wantsDir, { recursive: true });
|
||||
await symlink(f.seedTemplate, join(wantsDir, 'unit-under-test.service'));
|
||||
|
||||
await placeUnitFile(f.activeSource, f.activeDir, 'unit-under-test.service');
|
||||
const first = await readFile(join(f.activeDir, 'unit-under-test.service'), 'utf8');
|
||||
const secondRun = await placeUnitFile(f.activeSource, f.activeDir, 'unit-under-test.service');
|
||||
const second = await readFile(join(f.activeDir, 'unit-under-test.service'), 'utf8');
|
||||
expect(secondRun.unlinkedDestinationSymlink).toBe(false);
|
||||
expect(second).toBe(first);
|
||||
});
|
||||
});
|
||||
|
||||
describe('resolveLeaseBrokerSocketForPreflight (#1292 preflight resolution)', () => {
|
||||
it('explicit MOSAIC_LEASE_BROKER_SOCKET wins', () => {
|
||||
expect(
|
||||
resolveLeaseBrokerSocketForPreflight({ MOSAIC_LEASE_BROKER_SOCKET: '/custom/sock' }, 1000),
|
||||
).toBe('/custom/sock');
|
||||
});
|
||||
it('XDG_RUNTIME_DIR next', () => {
|
||||
expect(resolveLeaseBrokerSocketForPreflight({ XDG_RUNTIME_DIR: '/run/user/1001' }, 1000)).toBe(
|
||||
'/run/user/1001/mosaic-lease/broker.sock',
|
||||
);
|
||||
});
|
||||
it('falls back to /run/user/<uid>', () => {
|
||||
expect(resolveLeaseBrokerSocketForPreflight({}, 1002)).toBe(
|
||||
'/run/user/1002/mosaic-lease/broker.sock',
|
||||
);
|
||||
});
|
||||
});
|
||||
@@ -835,25 +835,13 @@ describe('fleet command construction', () => {
|
||||
};
|
||||
const program = new Command();
|
||||
program.exitOverride();
|
||||
// #1292: inject a present broker socket so the preflight passes and this
|
||||
// spec keeps testing its ORIGINAL property (holder-before-agent ordering).
|
||||
// The preflight's own refusal behavior has dedicated specs below.
|
||||
registerFleetCommand(program, {
|
||||
runner,
|
||||
mosaicHome: home,
|
||||
checkBrokerSocket: async () => true,
|
||||
});
|
||||
registerFleetCommand(program, { runner, mosaicHome: home });
|
||||
|
||||
try {
|
||||
await program.parseAsync(['node', 'mosaic', 'fleet', 'start']);
|
||||
await program.parseAsync(['node', 'mosaic', 'fleet', 'stop']);
|
||||
|
||||
expect(calls).toEqual([
|
||||
// #1292: fleet start enables + starts the broker FIRST (enable is
|
||||
// idempotent; the unit exists after install), re-checking the socket
|
||||
// before any holder/agent lifecycle effect.
|
||||
['systemctl', '--user', 'enable', 'mosaic-lease-broker.service'],
|
||||
['systemctl', '--user', 'start', 'mosaic-lease-broker.service'],
|
||||
['systemctl', '--user', 'start', 'mosaic-tmux-holder.service'],
|
||||
['systemctl', '--user', 'start', '[email protected]'],
|
||||
['systemctl', '--user', 'stop', '[email protected]'],
|
||||
@@ -864,92 +852,6 @@ describe('fleet command construction', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('fleet start refuses with a named error when the broker socket does not appear (#1292)', async () => {
|
||||
const home = await tempDir();
|
||||
const rosterPath = join(home, 'fleet', 'roster.yaml');
|
||||
await mkdir(join(home, 'fleet'), { recursive: true });
|
||||
await writeFile(
|
||||
rosterPath,
|
||||
['version: 1', 'transport: tmux', 'agents:', ' - name: coder0', ' runtime: codex'].join(
|
||||
'\n',
|
||||
),
|
||||
);
|
||||
const calls: string[][] = [];
|
||||
const runner: CommandRunner = async (command, args) => {
|
||||
calls.push([command, ...args]);
|
||||
return { stdout: '', stderr: '', exitCode: 0 };
|
||||
};
|
||||
const program = new Command();
|
||||
program.exitOverride();
|
||||
const errors: string[] = [];
|
||||
const origError = console.error;
|
||||
console.error = (...args: unknown[]) => {
|
||||
errors.push(args.join(' '));
|
||||
};
|
||||
registerFleetCommand(program, {
|
||||
runner,
|
||||
mosaicHome: home,
|
||||
checkBrokerSocket: async () => false,
|
||||
});
|
||||
try {
|
||||
await program.parseAsync(['node', 'mosaic', 'fleet', 'start']);
|
||||
// Refused: no holder/agent starts were issued after the broker attempt.
|
||||
expect(calls).toEqual([
|
||||
['systemctl', '--user', 'enable', 'mosaic-lease-broker.service'],
|
||||
['systemctl', '--user', 'start', 'mosaic-lease-broker.service'],
|
||||
]);
|
||||
expect(errors.join('\n')).toContain('broker-absent');
|
||||
expect(errors.join('\n')).toContain('mosaic fleet install');
|
||||
} finally {
|
||||
console.error = origError;
|
||||
await rm(home, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
|
||||
it('fleet start re-probes the broker on the SECOND invocation — no ActiveState trust (#1292 sticky half)', async () => {
|
||||
const home = await tempDir();
|
||||
const rosterPath = join(home, 'fleet', 'roster.yaml');
|
||||
await mkdir(join(home, 'fleet'), { recursive: true });
|
||||
await writeFile(
|
||||
rosterPath,
|
||||
['version: 1', 'transport: tmux', 'agents:', ' - name: coder0', ' runtime: codex'].join(
|
||||
'\n',
|
||||
),
|
||||
);
|
||||
const calls: string[][] = [];
|
||||
const runner: CommandRunner = async (command, args) => {
|
||||
calls.push([command, ...args]);
|
||||
return { stdout: '', stderr: '', exitCode: 0 };
|
||||
};
|
||||
const program = new Command();
|
||||
program.exitOverride();
|
||||
// Broker socket NEVER appears — the second start must refuse exactly like
|
||||
// the first; RemainAfterExit-style stale unit state changes nothing
|
||||
// because the check is the socket, not systemctl.
|
||||
registerFleetCommand(program, {
|
||||
runner,
|
||||
mosaicHome: home,
|
||||
checkBrokerSocket: async () => false,
|
||||
});
|
||||
const errors: string[] = [];
|
||||
const origError = console.error;
|
||||
console.error = (...args: unknown[]) => {
|
||||
errors.push(args.join(' '));
|
||||
};
|
||||
try {
|
||||
await program.parseAsync(['node', 'mosaic', 'fleet', 'start']);
|
||||
await program.parseAsync(['node', 'mosaic', 'fleet', 'start']);
|
||||
// Two invocations, each refusing after its own broker attempt:
|
||||
expect(
|
||||
calls.filter((c) => c.join(' ') === 'systemctl --user start [email protected]'),
|
||||
).toHaveLength(0);
|
||||
expect(errors.filter((e) => e.includes('broker-absent')).length).toBeGreaterThanOrEqual(2);
|
||||
} finally {
|
||||
console.error = origError;
|
||||
await rm(home, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
|
||||
it('waits for an in-flight restart to clear before relaunching (re-entry guard)', async () => {
|
||||
const home = await tempDir();
|
||||
const rosterPath = join(home, 'fleet', 'roster.yaml');
|
||||
@@ -2163,19 +2065,8 @@ describe('fleet install — auto-enable units for boot-survival', () => {
|
||||
|
||||
await enableFleetUnits(runner, minimalRoster, {});
|
||||
|
||||
expect(calls).toContainEqual(['systemctl', '--user', 'enable', 'mosaic-lease-broker.service']);
|
||||
expect(calls).toContainEqual(['systemctl', '--user', 'enable', 'mosaic-tmux-holder.service']);
|
||||
expect(calls).toContainEqual(['systemctl', '--user', 'enable', '[email protected]']);
|
||||
// The broker must be enabled BEFORE the holder and agents: a start of any
|
||||
// gated runtime without the broker is exactly the #1292 4-second death.
|
||||
const brokerIndex = calls.findIndex(
|
||||
(c) => c.join(' ') === 'systemctl --user enable mosaic-lease-broker.service',
|
||||
);
|
||||
const holderIndex = calls.findIndex(
|
||||
(c) => c.join(' ') === 'systemctl --user enable mosaic-tmux-holder.service',
|
||||
);
|
||||
expect(brokerIndex).toBeGreaterThanOrEqual(0);
|
||||
expect(brokerIndex).toBeLessThan(holderIndex);
|
||||
});
|
||||
|
||||
it('install still succeeds when systemctl enable returns non-zero (non-fatal)', async () => {
|
||||
|
||||
@@ -3,11 +3,9 @@ import {
|
||||
access,
|
||||
chmod,
|
||||
copyFile,
|
||||
lstat,
|
||||
mkdir,
|
||||
open,
|
||||
readFile,
|
||||
readlink,
|
||||
stat,
|
||||
unlink,
|
||||
writeFile,
|
||||
@@ -91,8 +89,6 @@ export type SleepFn = (ms: number) => Promise<void>;
|
||||
|
||||
export interface FleetCommandDeps {
|
||||
runner?: CommandRunner;
|
||||
/** Test seam for the #1292 fleet-start broker preflight (socket presence). */
|
||||
checkBrokerSocket?: (path: string) => Promise<boolean> | boolean;
|
||||
/** Injectable interactive runner for commands needing inherited TTY (e.g., `tmux attach`). */
|
||||
interactiveRunner?: InteractiveRunner;
|
||||
/**
|
||||
@@ -801,96 +797,6 @@ export function buildSystemdEnableCommand(unit: string): string[] {
|
||||
return ['systemctl', '--user', 'enable', unit];
|
||||
}
|
||||
|
||||
/**
|
||||
* Place a unit file into the ACTIVE systemd user directory, never through a
|
||||
* symlink (#1292, measured 2026-08-17).
|
||||
*
|
||||
* ⚠ SET-INDEPENDENCE (fomo-lin, 2026-08-17): the set of unit names carrying
|
||||
* by-path residue and the set of unit names this install copies are
|
||||
* INDEPENDENT. Until 0.0.50 they were disjoint only by accident of which
|
||||
* units the install happened to name — fomo-lin survived copy-through solely
|
||||
* because its one by-path symlink (the broker) was the one unit the install
|
||||
* did NOT copy. Adding the broker to the copy set made the intersection
|
||||
* non-empty on the first run. Whoever adds a fifth unit to the placement
|
||||
* list inherits this helper and its unlink step; do not place units with a
|
||||
* bare copyFile.
|
||||
*
|
||||
* A host provisioned by the enable-by-path convention carries a symlink AT
|
||||
* the unit-name path in ~/.config/systemd/user/ pointing at the shipped
|
||||
* template under ~/.config/mosaic/systemd/user/. Node's copyFile FOLLOWS
|
||||
* that link and overwrites the SEED template instead of placing the active
|
||||
* unit (verified with fs.copyFile on a throwaway systemd user instance) —
|
||||
* silent, rc=0, and it mutates the directory every later reseed reads from.
|
||||
* The same measurement showed `systemctl enable <name>` does NOT rewrite an
|
||||
* existing by-path wants-symlink, so reconciliation must be explicit.
|
||||
*
|
||||
* Placement therefore: if the destination is a symlink, unlink it first
|
||||
* (unlink → copy — copy-then-unlink would mutate the seed and then destroy
|
||||
* the evidence that it did); then copy. Also removes a stale
|
||||
* `default.target.wants/<name>` symlink that points outside the active
|
||||
* directory (readlink — NOT readFile, which follows the link and returns the
|
||||
* target's CONTENT), so the subsequent enable-by-name recreates it against
|
||||
* the active copy. Idempotent: on a clean or already-reconciled destination
|
||||
* every step is a no-op (the copy rewrites identical bytes).
|
||||
*
|
||||
* Returns what was done, for assertions and install reporting.
|
||||
*/
|
||||
export interface PlaceUnitResult {
|
||||
readonly unit: string;
|
||||
readonly destination: string;
|
||||
/** A symlink at the unit-name path was unlinked (by-path residue). */
|
||||
readonly unlinkedDestinationSymlink: boolean;
|
||||
/** A stale wants-symlink pointing outside the active dir was removed. */
|
||||
readonly removedStaleWantsSymlink: boolean;
|
||||
}
|
||||
|
||||
export async function placeUnitFile(
|
||||
source: string,
|
||||
systemdUserDir: string,
|
||||
unit: string,
|
||||
): Promise<PlaceUnitResult> {
|
||||
const destination = join(systemdUserDir, unit);
|
||||
let unlinkedDestinationSymlink = false;
|
||||
try {
|
||||
const destInfo = await lstat(destination);
|
||||
if (destInfo.isSymbolicLink()) {
|
||||
await unlink(destination);
|
||||
unlinkedDestinationSymlink = true;
|
||||
}
|
||||
} catch {
|
||||
// absent destination — nothing to unlink
|
||||
}
|
||||
await copyFile(source, destination);
|
||||
|
||||
let removedStaleWantsSymlink = false;
|
||||
const wantsLink = join(systemdUserDir, 'default.target.wants', unit);
|
||||
try {
|
||||
const wantsInfo = await lstat(wantsLink);
|
||||
if (wantsInfo.isSymbolicLink()) {
|
||||
// readlink — NOT readFile: readFile FOLLOWS the link and returns the
|
||||
// target file's CONTENT, which is not the question being asked.
|
||||
let target: string | undefined;
|
||||
try {
|
||||
target = await readlink(wantsLink);
|
||||
} catch {
|
||||
target = undefined;
|
||||
}
|
||||
// Normalize (systemctl writes absolute targets; a relative one resolves
|
||||
// against the wants dir). A wants-symlink pointing anywhere other than
|
||||
// the active copy (the by-path convention points at the seed template)
|
||||
// survives enable-by-name unchanged — remove it so enable recreates it.
|
||||
if (target !== undefined && resolve(dirname(wantsLink), target) !== destination) {
|
||||
await unlink(wantsLink);
|
||||
removedStaleWantsSymlink = true;
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
// absent wants link — nothing to reconcile
|
||||
}
|
||||
|
||||
return { unit, destination, unlinkedDestinationSymlink, removedStaleWantsSymlink };
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the systemctl --user disable command for a given unit.
|
||||
* Used by `fleet remove` so a removed agent's enabled unit cannot resurrect on
|
||||
@@ -925,22 +831,6 @@ export async function enableFleetUnits(
|
||||
let succeeded = 0;
|
||||
let failed = 0;
|
||||
|
||||
// The lease broker ships with the fleet and every gated runtime needs it
|
||||
// (#1292): seats die at lease registration without it, and no documented
|
||||
// path ever enabled it. Enabled first — alongside the holder — and the
|
||||
// unit must have been placed by installFleet's placeUnitFile step.
|
||||
const brokerResult = await runner(
|
||||
...splitCommand(buildSystemdEnableCommand('mosaic-lease-broker.service')),
|
||||
);
|
||||
if (brokerResult.exitCode === 0) {
|
||||
succeeded++;
|
||||
} else {
|
||||
failed++;
|
||||
process.stderr.write(
|
||||
`Warning: could not enable mosaic-lease-broker.service: ${brokerResult.stderr || brokerResult.stdout || 'non-zero exit'}\n`,
|
||||
);
|
||||
}
|
||||
|
||||
const holderResult = await runner(
|
||||
...splitCommand(buildSystemdEnableCommand('mosaic-tmux-holder.service')),
|
||||
);
|
||||
@@ -1637,7 +1527,7 @@ export function registerFleetCommand(program: Command, deps: FleetCommandDeps =
|
||||
.description('Install local fleet tools and user systemd units')
|
||||
.option('--no-enable', 'Skip enabling units for boot-survival')
|
||||
.action(async (opts: { enable?: boolean }) => {
|
||||
await installFleet(cmd, frameworkRoot, runner);
|
||||
await installFleet(cmd, frameworkRoot);
|
||||
// Unit enablement needs agent names only, so it reads either version.
|
||||
const roster = await loadRosterReadModel(cmd);
|
||||
await enableFleetUnits(runner, roster, opts);
|
||||
@@ -1648,7 +1538,7 @@ export function registerFleetCommand(program: Command, deps: FleetCommandDeps =
|
||||
.description('Install local fleet tools and user systemd units')
|
||||
.option('--no-enable', 'Skip enabling units for boot-survival')
|
||||
.action(async (opts: { enable?: boolean }) => {
|
||||
await installFleet(cmd, frameworkRoot, runner);
|
||||
await installFleet(cmd, frameworkRoot);
|
||||
// Unit enablement needs agent names only, so it reads either version.
|
||||
const roster = await loadRosterReadModel(cmd);
|
||||
await enableFleetUnits(runner, roster, opts);
|
||||
@@ -1701,37 +1591,6 @@ export function registerFleetCommand(program: Command, deps: FleetCommandDeps =
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (action === 'start') {
|
||||
// Broker preflight (#1292), re-probed on EVERY invocation: a
|
||||
// gated runtime started without a live lease broker dies ~4s in
|
||||
// while the unit reports active (RemainAfterExit) — enabling +
|
||||
// starting here and then RE-CHECKING the socket refuses loudly
|
||||
// instead of reporting rc0 over a doomed start. This is the
|
||||
// second-start check as much as the first: it never trusts unit
|
||||
// ActiveState.
|
||||
await runChecked(runner, [
|
||||
'systemctl',
|
||||
'--user',
|
||||
'enable',
|
||||
'mosaic-lease-broker.service',
|
||||
]);
|
||||
await runChecked(runner, [
|
||||
'systemctl',
|
||||
'--user',
|
||||
'start',
|
||||
'mosaic-lease-broker.service',
|
||||
]);
|
||||
if (!(await brokerSocketPresent(deps))) {
|
||||
console.error(
|
||||
'[fleet] broker-absent: lease broker socket did not appear after enable+start (#1292).',
|
||||
);
|
||||
console.error(
|
||||
'[fleet] remedy: mosaic fleet install (it reconciles either enable convention)',
|
||||
);
|
||||
process.exitCode = 1;
|
||||
return;
|
||||
}
|
||||
}
|
||||
if (action === 'restart') {
|
||||
// Serialize the holder+agents teardown/relaunch behind the restart lock
|
||||
// so a re-entrant restart waits for clean shutdown before relaunching,
|
||||
@@ -2490,11 +2349,7 @@ export function registerFleetAgentCommands(
|
||||
});
|
||||
}
|
||||
|
||||
async function installFleet(
|
||||
cmd: Command,
|
||||
frameworkRoot: string,
|
||||
runner: CommandRunner,
|
||||
): Promise<void> {
|
||||
async function installFleet(cmd: Command, frameworkRoot: string): Promise<void> {
|
||||
const activePaths = resolveFleetPaths(cmd.opts<{ mosaicHome: string }>().mosaicHome);
|
||||
assertDefaultMosaicHomeForSystemd(activePaths.mosaicHome);
|
||||
// Read model first: every file this function places is roster-independent, and
|
||||
@@ -2546,40 +2401,18 @@ async function installFleet(
|
||||
for (const toolPath of executableToolPaths) {
|
||||
await chmod(toolPath, 0o755);
|
||||
}
|
||||
// Unit placement (#1292): every unit goes through placeUnitFile — never a
|
||||
// bare copyFile — so a by-path-enable symlink at the destination is
|
||||
// unlinked rather than written through (copy-through would silently
|
||||
// overwrite the SEED template, measured 2026-08-17). The lease broker unit
|
||||
// is placed here too: previously the install named three units and omitted
|
||||
// the broker entirely, which is why no documented path ever enabled it.
|
||||
const placedUnits = await Promise.all(
|
||||
[
|
||||
'mosaic-tmux-holder.service',
|
||||
'[email protected]',
|
||||
'[email protected]',
|
||||
'mosaic-lease-broker.service',
|
||||
].map((unit) =>
|
||||
placeUnitFile(join(frameworkRoot, 'systemd', 'user', unit), activePaths.systemdUserDir, unit),
|
||||
),
|
||||
await copyFile(
|
||||
join(frameworkRoot, 'systemd', 'user', 'mosaic-tmux-holder.service'),
|
||||
join(activePaths.systemdUserDir, 'mosaic-tmux-holder.service'),
|
||||
);
|
||||
const reconciled = placedUnits.filter(
|
||||
(result) => result.unlinkedDestinationSymlink || result.removedStaleWantsSymlink,
|
||||
await copyFile(
|
||||
join(frameworkRoot, 'systemd', 'user', '[email protected]'),
|
||||
join(activePaths.systemdUserDir, '[email protected]'),
|
||||
);
|
||||
await copyFile(
|
||||
join(frameworkRoot, 'systemd', 'user', '[email protected]'),
|
||||
join(activePaths.systemdUserDir, '[email protected]'),
|
||||
);
|
||||
if (reconciled.length > 0) {
|
||||
console.log(
|
||||
`Reconciled ${reconciled.length} unit placement(s) from by-path enable residue: ${reconciled.map((r) => r.unit).join(', ')}`,
|
||||
);
|
||||
}
|
||||
// systemd will not see a replaced unit file without a reload; do it once
|
||||
// after all placements, before any enable call below. runCommand never
|
||||
// rejects (it resolves exitCode 127 on spawn error), so a plain await with
|
||||
// an exitCode check matches the rest of this file's systemctl handling.
|
||||
const reloadResult = await runner(...splitCommand(['systemctl', '--user', 'daemon-reload']));
|
||||
if (reloadResult.exitCode !== 0) {
|
||||
process.stderr.write(
|
||||
`Warning: systemctl --user daemon-reload after unit placement failed (non-systemd host?): ${reloadResult.stderr || reloadResult.stdout || 'non-zero exit'}\n`,
|
||||
);
|
||||
}
|
||||
|
||||
// On roster v2 the reconciler owns the generated env: `apply` writes it and
|
||||
// `regen` rebuilds it, both from projectRosterV2AgentGeneratedEnv. Writing it
|
||||
@@ -2776,40 +2609,6 @@ function splitCommand(command: string[]): [string, string[]] {
|
||||
return [bin, args];
|
||||
}
|
||||
|
||||
/**
|
||||
* Lease-broker socket presence for the fleet-start preflight (#1292).
|
||||
* Resolution precedence matches launch.ts's defaultLeaseBrokerSocket and
|
||||
* start-agent-session.sh's broker_socket_path: explicit
|
||||
* MOSAIC_LEASE_BROKER_SOCKET, else $XDG_RUNTIME_DIR/mosaic-lease/broker.sock,
|
||||
* else /run/user/<uid>/mosaic-lease/broker.sock. Pure filesystem check — this
|
||||
* deliberately does NOT consult systemd state: a unit can be active
|
||||
* (RemainAfterExit) with no live socket, and the socket is the thing the
|
||||
* gated runtime connects to. Injectable via deps for tests.
|
||||
*/
|
||||
export function resolveLeaseBrokerSocketForPreflight(
|
||||
env: NodeJS.ProcessEnv = process.env,
|
||||
uid: number = typeof process.getuid === 'function' ? process.getuid() : 0,
|
||||
): string {
|
||||
if (env['MOSAIC_LEASE_BROKER_SOCKET']) return env['MOSAIC_LEASE_BROKER_SOCKET'];
|
||||
const runtimeDir = env['XDG_RUNTIME_DIR'] ?? `/run/user/${uid}`;
|
||||
return join(runtimeDir, 'mosaic-lease', 'broker.sock');
|
||||
}
|
||||
|
||||
async function brokerSocketPresent(
|
||||
deps: FleetCommandDeps,
|
||||
env: NodeJS.ProcessEnv = process.env,
|
||||
): Promise<boolean> {
|
||||
const check = deps.checkBrokerSocket;
|
||||
if (check) return check(resolveLeaseBrokerSocketForPreflight(env));
|
||||
try {
|
||||
const socketPath = resolveLeaseBrokerSocketForPreflight(env);
|
||||
await access(socketPath, constants.S_IFSOCK);
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/** All supported fleet profile names. */
|
||||
export type FleetProfile =
|
||||
| 'general'
|
||||
|
||||
@@ -205,12 +205,6 @@ export async function runLeaseEnforcementDoctorCheck(
|
||||
message:
|
||||
`Lease-enforcement hooks (${matchedMarkers.join(', ')}) are wired in ~/.claude/settings.json, but ${reasons.join(' and ')}. ` +
|
||||
'Every gated tool call will fail closed and BRICK this agent (see #869). ' +
|
||||
// #1292: one remedy, correct under BOTH enable conventions (by-path on
|
||||
// the seed template, and copy-then-enable in the active dir). Written
|
||||
// from the 2026-08-17 symlink measurement: `systemctl enable` by name
|
||||
// does NOT rewrite an existing by-path wants-symlink, so teaching a
|
||||
// manual systemctl line here could leave a host with two competing
|
||||
// wants links. fleet install reconciles either shape.
|
||||
'Remedy: run `mosaic fleet install` (it reconciles either enable convention), or remove the enforcement hooks from ~/.claude/settings.json.',
|
||||
'Remediate by activating the lease-broker supervisor (systemd unit + socket) or by removing the enforcement hooks from ~/.claude/settings.json.',
|
||||
};
|
||||
}
|
||||
|
||||
@@ -459,7 +459,6 @@ describe('FCM-M3-002 reconciler lifecycle acceptance', (): void => {
|
||||
plan: {
|
||||
generation: 7,
|
||||
holder: 'owned',
|
||||
broker: { unitInstalled: false, socketPresent: false },
|
||||
agents: [
|
||||
{
|
||||
name: 'coder0',
|
||||
|
||||
@@ -92,112 +92,6 @@ async function run(command: FleetReconcileCommand, overrides: Partial<FleetRecon
|
||||
}
|
||||
|
||||
describe('fleet roster-owned reconciler', (): void => {
|
||||
// ── #1292: broker as first-class plan member + broker-first start ordering ──
|
||||
|
||||
it('reports broker unit and socket state in the plan (socket is the signal, not unit state)', async (): Promise<void> => {
|
||||
const result = await run('status', {
|
||||
statPath: async () => true,
|
||||
checkBrokerSocket: async () => true,
|
||||
});
|
||||
expect(result.plan.broker).toEqual({ unitInstalled: true, socketPresent: true });
|
||||
});
|
||||
|
||||
it('reports a dead broker as socketPresent=false even when the unit is installed (enabled-but-dead is the #1292 shape)', async (): Promise<void> => {
|
||||
const result = await run('status', {
|
||||
statPath: async () => true,
|
||||
checkBrokerSocket: async () => false,
|
||||
});
|
||||
expect(result.plan.broker).toEqual({ unitInstalled: true, socketPresent: false });
|
||||
});
|
||||
|
||||
it('reports broker-absent when neither seam is present (defaults false, never guesses healthy)', async (): Promise<void> => {
|
||||
const result = await run('status');
|
||||
expect(result.plan.broker).toEqual({ unitInstalled: false, socketPresent: false });
|
||||
});
|
||||
|
||||
it('command start enables and starts the broker BEFORE the holder and any agent unit', async (): Promise<void> => {
|
||||
const calls: string[][] = [];
|
||||
const result = await run('start', {
|
||||
runner: async (command, args) => {
|
||||
calls.push([command, ...args]);
|
||||
if (command === 'tmux' && args.includes('list-sessions')) {
|
||||
return { stdout: '_holder\ncoder0\n', stderr: '', exitCode: 0 };
|
||||
}
|
||||
if (command === 'tmux' && args.includes('show-environment')) {
|
||||
return {
|
||||
stdout:
|
||||
'HOME=/home/mosaic\nMOSAIC_FLEET_OWNER=11111111-1111-4111-8111-111111111111\nMOSAIC_TMUX_HOLDER=_holder\nMOSAIC_TMUX_SOCKET=mosaic-fleet\nPATH=/usr/bin:/bin\nPWD=/home/mosaic\n',
|
||||
stderr: '',
|
||||
exitCode: 0,
|
||||
};
|
||||
}
|
||||
return { stdout: '', stderr: '', exitCode: 0 };
|
||||
},
|
||||
});
|
||||
expect(result.lifecycle).toBe('complete');
|
||||
const brokerEnable = calls.findIndex(
|
||||
(c) => c.join(' ') === 'systemctl --user enable mosaic-lease-broker.service',
|
||||
);
|
||||
const brokerStart = calls.findIndex(
|
||||
(c) => c.join(' ') === 'systemctl --user start mosaic-lease-broker.service',
|
||||
);
|
||||
const holderStart = calls.findIndex(
|
||||
(c) => c.join(' ') === 'systemctl --user start mosaic-tmux-holder.service',
|
||||
);
|
||||
const agentStart = calls.findIndex(
|
||||
(c) => c.join(' ') === 'systemctl --user start [email protected]',
|
||||
);
|
||||
expect(brokerEnable).toBeGreaterThanOrEqual(0);
|
||||
expect(brokerStart).toBeGreaterThan(brokerEnable);
|
||||
// Holder start may be absent (holder 'owned' in this fixture); if present it must follow the broker.
|
||||
if (holderStart >= 0) expect(holderStart).toBeGreaterThan(brokerStart);
|
||||
expect(agentStart).toBeGreaterThan(brokerStart);
|
||||
});
|
||||
|
||||
it('apply with a running desired agent also enables and starts the broker first', async (): Promise<void> => {
|
||||
const calls: string[][] = [];
|
||||
const runningRoster: FleetRosterV2 = {
|
||||
...roster,
|
||||
agents: roster.agents.map((agent) =>
|
||||
agent.name === 'coder0'
|
||||
? { ...agent, lifecycle: { enabled: true, desiredState: 'running' as const } }
|
||||
: agent,
|
||||
),
|
||||
};
|
||||
const result = await executeFleetReconcile({
|
||||
roster: runningRoster,
|
||||
command: 'apply',
|
||||
expectedGeneration: 7,
|
||||
deps: deps({
|
||||
readRoster: async () => runningRoster,
|
||||
runner: async (command, args) => {
|
||||
calls.push([command, ...args]);
|
||||
if (command === 'tmux' && args.includes('list-sessions')) {
|
||||
return { stdout: '_holder\n', stderr: '', exitCode: 0 };
|
||||
}
|
||||
if (command === 'tmux' && args.includes('show-environment')) {
|
||||
return {
|
||||
stdout:
|
||||
'HOME=/home/mosaic\nMOSAIC_FLEET_OWNER=11111111-1111-4111-8111-111111111111\nMOSAIC_TMUX_HOLDER=_holder\nMOSAIC_TMUX_SOCKET=mosaic-fleet\nPATH=/usr/bin:/bin\nPWD=/home/mosaic\n',
|
||||
stderr: '',
|
||||
exitCode: 0,
|
||||
};
|
||||
}
|
||||
return { stdout: '', stderr: '', exitCode: 0 };
|
||||
},
|
||||
}),
|
||||
});
|
||||
expect(result.applied).toBe(true);
|
||||
const brokerStart = calls.findIndex(
|
||||
(c) => c.join(' ') === 'systemctl --user start mosaic-lease-broker.service',
|
||||
);
|
||||
const agentStart = calls.findIndex(
|
||||
(c) => c.join(' ') === 'systemctl --user start [email protected]',
|
||||
);
|
||||
expect(brokerStart).toBeGreaterThanOrEqual(0);
|
||||
expect(agentStart).toBeGreaterThan(brokerStart);
|
||||
});
|
||||
|
||||
it('fails closed on a symlinked fleet ancestor without touching its target', async (): Promise<void> => {
|
||||
const home = await lockHome();
|
||||
const fleet = join(home, 'fleet');
|
||||
|
||||
@@ -43,10 +43,6 @@ export interface FleetReconcileDeps {
|
||||
readonly overrideDir?: string;
|
||||
readonly homeDirectory?: string;
|
||||
readonly readHolderIdentity?: () => Promise<string>;
|
||||
/** Test/observation seams for the lease-broker plan member (#1292). */
|
||||
readonly statPath?: (path: string) => Promise<boolean> | boolean;
|
||||
readonly checkBrokerSocket?: (path: string) => Promise<boolean> | boolean;
|
||||
readonly brokerSocketEnv?: NodeJS.ProcessEnv;
|
||||
readonly validateRoster?: (roster: FleetRosterV2) => Promise<void>;
|
||||
readonly prepareProjections?: (roster: FleetRosterV2) => Promise<readonly unknown[]>;
|
||||
readonly applyProjection?: (prepared: unknown) => Promise<unknown>;
|
||||
@@ -78,17 +74,6 @@ export interface FleetReconcileObservedAgent {
|
||||
export interface FleetReconcilePlan {
|
||||
readonly generation: number;
|
||||
readonly holder: 'owned' | 'missing' | 'ownership-mismatch';
|
||||
/**
|
||||
* Lease broker observation (#1292): every gated runtime registers with the
|
||||
* broker or dies ~4s in — a broker not in the plan cannot be reported as
|
||||
* drifted, which made "broker died an hour ago" and "broker fine"
|
||||
* produce identical output. `unitInstalled` = unit file present in the
|
||||
* active dir; `socketPresent` = live broker at the resolved socket path.
|
||||
*/
|
||||
readonly broker: {
|
||||
readonly unitInstalled: boolean;
|
||||
readonly socketPresent: boolean;
|
||||
};
|
||||
readonly agents: readonly FleetReconcileObservedAgent[];
|
||||
readonly unmanagedSessions: readonly string[];
|
||||
}
|
||||
@@ -329,39 +314,6 @@ function isObservational(command: FleetReconcileCommand): boolean {
|
||||
return command === 'plan' || command === 'status' || command === 'verify' || command === 'doctor';
|
||||
}
|
||||
|
||||
/**
|
||||
* Observe the lease broker for the plan (#1292). Unit presence via systemctl
|
||||
* is-system-running is NOT the signal — a unit can be enabled-but-dead. The
|
||||
* authoritative signal is the socket the gated runtimes connect to, matching
|
||||
* broker-supervisor.ts's `checkBrokerSupervisorHealth` (healthy ===
|
||||
* socketPresent). Injectable so tests drive every branch without a broker.
|
||||
*/
|
||||
async function observeBroker(deps: FleetReconcileDeps): Promise<FleetReconcilePlan['broker']> {
|
||||
const homeDirectory = deps.homeDirectory ?? homedir();
|
||||
const env = (deps.brokerSocketEnv ?? process.env) as NodeJS.ProcessEnv;
|
||||
const uid = typeof process.getuid === 'function' ? process.getuid() : 0;
|
||||
const runtimeDir = env['XDG_RUNTIME_DIR'] ?? `/run/user/${uid}`;
|
||||
const socketPath =
|
||||
env['MOSAIC_LEASE_BROKER_SOCKET'] ?? join(runtimeDir, 'mosaic-lease', 'broker.sock');
|
||||
const configHome = env['XDG_CONFIG_HOME'] ?? join(homeDirectory, '.config');
|
||||
const unitPath = join(configHome, 'systemd', 'user', 'mosaic-lease-broker.service');
|
||||
const statPath = deps.statPath;
|
||||
const checkBrokerSocket = deps.checkBrokerSocket;
|
||||
let unitInstalled = false;
|
||||
let socketPresent = false;
|
||||
try {
|
||||
unitInstalled = statPath ? await statPath(unitPath) : false;
|
||||
} catch {
|
||||
unitInstalled = false;
|
||||
}
|
||||
try {
|
||||
socketPresent = checkBrokerSocket ? await checkBrokerSocket(socketPath) : false;
|
||||
} catch {
|
||||
socketPresent = false;
|
||||
}
|
||||
return { unitInstalled, socketPresent };
|
||||
}
|
||||
|
||||
async function observeFleet(
|
||||
roster: FleetRosterV2,
|
||||
deps: FleetReconcileDeps,
|
||||
@@ -372,12 +324,10 @@ async function observeFleet(
|
||||
'-F',
|
||||
'#{session_name}',
|
||||
]);
|
||||
const broker = await observeBroker(deps);
|
||||
if (sessionsResult.exitCode !== 0) {
|
||||
return {
|
||||
generation: roster.generation,
|
||||
holder: 'missing',
|
||||
broker,
|
||||
agents: await observeAgents(roster, deps, new Set<string>()),
|
||||
unmanagedSessions: [],
|
||||
};
|
||||
@@ -400,7 +350,6 @@ async function observeFleet(
|
||||
return {
|
||||
generation: roster.generation,
|
||||
holder,
|
||||
broker,
|
||||
agents: await observeAgents(roster, deps, sessions),
|
||||
unmanagedSessions: Object.freeze(unmanagedSessions.sort()),
|
||||
};
|
||||
@@ -568,24 +517,6 @@ async function executeExplicitLifecycle(
|
||||
}
|
||||
}
|
||||
try {
|
||||
// Broker FIRST (#1292): a gated runtime started without a running lease
|
||||
// broker dies ~4 seconds in at registration — enable the unit (install
|
||||
// places it) and start it before any holder/agent lifecycle effect. The
|
||||
// socket re-check after start is the same probe observeBroker uses, so a
|
||||
// unit that starts but never produces a socket is caught here, not four
|
||||
// seconds later inside a doomed seat.
|
||||
if (request.command === 'start') {
|
||||
await runChecked(request.deps, 'systemctl', [
|
||||
'--user',
|
||||
'enable',
|
||||
'mosaic-lease-broker.service',
|
||||
]);
|
||||
await runChecked(request.deps, 'systemctl', [
|
||||
'--user',
|
||||
'start',
|
||||
'mosaic-lease-broker.service',
|
||||
]);
|
||||
}
|
||||
if (request.command === 'start' && plan.holder === 'missing') {
|
||||
await runChecked(request.deps, 'systemctl', [
|
||||
'--user',
|
||||
@@ -631,12 +562,6 @@ async function applyDesiredLifecycle(
|
||||
(agent: FleetRosterV2Agent): boolean =>
|
||||
agent.lifecycle.enabled && agent.lifecycle.desiredState === 'running',
|
||||
);
|
||||
// Broker before any running agent, same ordering and reason as the
|
||||
// command-driven path above (#1292).
|
||||
if (needsRunningAgent) {
|
||||
await runChecked(deps, 'systemctl', ['--user', 'enable', 'mosaic-lease-broker.service']);
|
||||
await runChecked(deps, 'systemctl', ['--user', 'start', 'mosaic-lease-broker.service']);
|
||||
}
|
||||
if (needsRunningAgent && plan.holder === 'missing') {
|
||||
await runChecked(deps, 'systemctl', ['--user', 'start', 'mosaic-tmux-holder.service']);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user