fix(fleet): lease-broker activation, symlink-safe unit placement, named launch refusal — Wall 6 (#1292) (#1297)
ci/woodpecker/push/publish Pipeline was successful

Co-authored-by: fargo <[email protected]>
This commit was merged in pull request #1297.
This commit is contained in:
2026-08-21 00:40:29 +00:00
committed by fred
parent 6306914965
commit 1d84bc3f3d
11 changed files with 1229 additions and 33 deletions
@@ -169,6 +169,16 @@ function reconcileDeps(host: FakeLifecycleHost): FleetReconcileDeps {
applyProjection: async () => undefined,
readRoster: async () => host.roster,
acquireMutationLock: async () => async () => undefined,
// Hermetic broker observation (#1297 F3): without this, the plan probes
// the REAL host filesystem, so the "stable JSON" fixtures answered true
// on any machine with a live lease broker and false elsewhere. Pointing
// both paths at fixtures that do not exist pins socketPresent:false and
// unitInstalled:false on every host, which is what these fixtures assert.
brokerSocketEnv: {
MOSAIC_LEASE_BROKER_SOCKET: '/nonexistent/mosaic-lease/broker.sock',
XDG_CONFIG_HOME: '/nonexistent/mosaic-config',
XDG_RUNTIME_DIR: '/nonexistent/run',
},
};
}
@@ -459,6 +469,7 @@ describe('FCM-M3-002 reconciler lifecycle acceptance', (): void => {
plan: {
generation: 7,
holder: 'owned',
broker: { unitInstalled: false, socketPresent: false },
agents: [
{
name: 'coder0',
@@ -1,6 +1,7 @@
import { chmod, mkdir, mkdtemp, readFile, rm, symlink, writeFile } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { createServer } from 'node:net';
import { afterEach, describe, expect, it } from 'vitest';
import {
acquirePrivateReconcileLock,
@@ -92,6 +93,179 @@ 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('probes the REAL filesystem when no seam is injected — live socket and unit report healthy, absent paths report absent (#1297 F3)', async (): Promise<void> => {
const dir = await mkdtemp(join(tmpdir(), 'mosaic-broker-probe-'));
cleanup = dir;
const configHome = join(dir, 'config');
const unitDir = join(configHome, 'systemd', 'user');
await mkdir(unitDir, { recursive: true });
await writeFile(join(unitDir, 'mosaic-lease-broker.service'), '[Unit]\n');
const sockPath = join(dir, 'broker.sock');
const server = createServer();
await new Promise<void>((resolve) => {
server.listen(sockPath, resolve);
});
try {
const result = await run('status', {
brokerSocketEnv: {
MOSAIC_LEASE_BROKER_SOCKET: sockPath,
XDG_CONFIG_HOME: configHome,
XDG_RUNTIME_DIR: dir,
},
});
expect(result.plan.broker).toEqual({ unitInstalled: true, socketPresent: true });
// Absent paths through the SAME seam-less path answer false — this is
// the half the old default got right; healthy is the half it got wrong.
const absent = await run('status', {
brokerSocketEnv: {
MOSAIC_LEASE_BROKER_SOCKET: join(dir, 'gone.sock'),
XDG_CONFIG_HOME: join(dir, 'gone-config'),
},
});
expect(absent.plan.broker).toEqual({ unitInstalled: false, socketPresent: false });
} finally {
await new Promise<void>((resolve) => {
server.close(() => resolve());
});
}
});
it('command start refuses with a named error when the broker socket does not appear after enable+start (#1297 F3)', async (): Promise<void> => {
const calls: string[][] = [];
await expect(
run('start', {
checkBrokerSocket: async () => false,
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 };
},
}),
).rejects.toThrow(/broker-absent/);
// Refused: broker enable+start attempted, no holder/agent unit touched.
const agentStarts = calls.filter(
(c) => c.join(' ') === 'systemctl --user start [email protected]',
);
expect(agentStarts).toHaveLength(0);
});
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', {
// Deterministic broker presence: without the seam this test answers the
// HOST's broker state (passes on a machine with a live broker, refuses
// on CI), not the ordering property it exists for (#1297 follow-up).
checkBrokerSocket: async () => true,
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,
// Deterministic broker presence (see start-ordering test note).
checkBrokerSocket: async () => true,
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');
@@ -375,6 +549,8 @@ describe('fleet roster-owned reconciler', (): void => {
expectedGeneration: 7,
deps: deps({
readRoster: async () => runningRoster,
// Deterministic broker presence (see start-ordering test note).
checkBrokerSocket: async () => true,
runner: async (command, args) => {
calls.push([command, ...args]);
if (command === 'tmux' && args.includes('list-sessions')) {
+149 -13
View File
@@ -1,5 +1,5 @@
import { constants } from 'node:fs';
import { lstat, open, readFile, unlink, type FileHandle } from 'node:fs/promises';
import { lstat, open, readFile, stat, unlink, type FileHandle } from 'node:fs/promises';
import { randomUUID } from 'node:crypto';
import { homedir } from 'node:os';
import { join } from 'node:path';
@@ -44,6 +44,10 @@ 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>;
@@ -75,6 +79,17 @@ 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[];
}
@@ -246,7 +261,17 @@ export async function executeFleetReconcile(
lifecycle: 'complete',
plan,
};
} catch {
} catch (error: unknown) {
// A named lifecycle precondition (broker-absent after enable+start,
// #1297 F3) must surface as itself — converting it to the generic
// recoverable result would hide the diagnosis and report a clean
// refusal where a loud one is the point.
if (
error instanceof FleetReconcileError &&
error.code === 'lifecycle-precondition-failed'
) {
throw error;
}
result = {
applied: false,
authoritativeRoster: 'unchanged',
@@ -315,6 +340,63 @@ 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.
*/
function resolveBrokerSocketPath(env: NodeJS.ProcessEnv): string {
const uid = typeof process.getuid === 'function' ? process.getuid() : 0;
const runtimeDir = env['XDG_RUNTIME_DIR'] ?? `/run/user/${uid}`;
return env['MOSAIC_LEASE_BROKER_SOCKET'] ?? join(runtimeDir, 'mosaic-lease', 'broker.sock');
}
/**
* Probe the broker socket. Seams take precedence, but with no seam injected
* the REAL stat().isSocket() runs (#1297 review F3): production passes no
* seams, and defaulting to false made plan/status/doctor report a healthy
* broker as absent — a dead broker was indistinguishable from noise.
*/
async function brokerSocketPresent(
deps: FleetReconcileDeps,
env: NodeJS.ProcessEnv,
): Promise<boolean> {
const socketPath = resolveBrokerSocketPath(env);
const check = deps.checkBrokerSocket;
if (check) return check(socketPath);
try {
return (await stat(socketPath)).isSocket();
} catch {
return false;
}
}
async function observeBroker(deps: FleetReconcileDeps): Promise<FleetReconcilePlan['broker']> {
const homeDirectory = deps.homeDirectory ?? homedir();
const env = (deps.brokerSocketEnv ?? process.env) as NodeJS.ProcessEnv;
const configHome = env['XDG_CONFIG_HOME'] ?? join(homeDirectory, '.config');
const unitPath = join(configHome, 'systemd', 'user', 'mosaic-lease-broker.service');
const statPath = deps.statPath;
let unitInstalled = false;
let socketPresent = false;
try {
// Same principle as the socket probe: no seam → look at the real
// filesystem. A unit file placed by installFleet (or a by-path residue
// symlink resolving to it) satisfies stat().isFile().
unitInstalled = statPath ? await statPath(unitPath) : (await stat(unitPath)).isFile();
} catch {
unitInstalled = false;
}
try {
socketPresent = await brokerSocketPresent(deps, env);
} catch {
socketPresent = false;
}
return { unitInstalled, socketPresent };
}
async function observeFleet(
roster: FleetRosterV2,
deps: FleetReconcileDeps,
@@ -325,10 +407,12 @@ 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: [],
};
@@ -351,6 +435,7 @@ async function observeFleet(
return {
generation: roster.generation,
holder,
broker,
agents: await observeAgents(roster, deps, sessions),
unmanagedSessions: Object.freeze(unmanagedSessions.sort()),
};
@@ -507,6 +592,17 @@ async function executeExplicitLifecycle(
plan: FleetReconcilePlan,
agents: readonly FleetRosterV2Agent[],
): Promise<FleetReconcileResult> {
const lifecycleApplyFailed = (): FleetReconcileResult => ({
applied: false,
authoritativeRoster: 'unchanged',
projections: 'not-applied',
lifecycle: 'incomplete',
plan,
recovery: {
code: 'lifecycle-apply-failed',
action: 'rerun-after-inspecting-owned-resources',
},
});
if (request.command === 'start') {
for (const agent of agents) {
if (!agent.lifecycle.enabled) {
@@ -517,6 +613,40 @@ async function executeExplicitLifecycle(
}
}
}
// 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.
try {
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',
]);
}
} catch {
return lifecycleApplyFailed();
}
if (request.command === 'start') {
// Socket re-check after start, as a NAMED precondition (#1297 review
// F3) — the same protection the v1 path in commands/fleet.ts has had all
// along: the unit reporting active is not the signal; the socket is.
// Deliberately outside the try/catch above: a swallowed FleetReconcileError
// here read as a generic recoverable failure, hiding the named refusal.
// Runs BEFORE any holder/agent unit is touched so nothing doomed starts.
const env = (request.deps.brokerSocketEnv ?? process.env) as NodeJS.ProcessEnv;
if (!(await brokerSocketPresent(request.deps, env))) {
throw new FleetReconcileError(
'lifecycle-precondition-failed',
'broker-absent: lease broker socket did not appear after enable+start (#1292; #1297 F3). Remedy: mosaic fleet install.',
);
}
}
try {
if (request.command === 'start' && plan.holder === 'missing') {
await runChecked(request.deps, 'systemctl', [
@@ -533,17 +663,7 @@ async function executeExplicitLifecycle(
]);
}
} catch {
return {
applied: false,
authoritativeRoster: 'unchanged',
projections: 'not-applied',
lifecycle: 'incomplete',
plan,
recovery: {
code: 'lifecycle-apply-failed',
action: 'rerun-after-inspecting-owned-resources',
},
};
return lifecycleApplyFailed();
}
return {
applied: true,
@@ -563,6 +683,22 @@ 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']);
// Same socket re-check as the explicit start path (#1297 F3): apply with
// running desired agents starts gated runtimes too, and a broker that
// starts but never binds dooms them the same way.
const env = (deps.brokerSocketEnv ?? process.env) as NodeJS.ProcessEnv;
if (!(await brokerSocketPresent(deps, env))) {
throw new FleetReconcileError(
'lifecycle-precondition-failed',
'broker-absent: lease broker socket did not appear after enable+start (#1292; #1297 F3). Remedy: mosaic fleet install.',
);
}
}
if (needsRunningAgent && plan.holder === 'missing') {
await runChecked(deps, 'systemctl', ['--user', 'start', 'mosaic-tmux-holder.service']);
}