fix(ci): stabilize team launch reconciliation

This commit is contained in:
777genius 2026-04-24 01:28:11 +03:00
parent e9f46c5bd4
commit ba4de3775b
4 changed files with 198 additions and 47 deletions

View file

@ -601,7 +601,7 @@ export function normalizePersistedLaunchSnapshot(
name, name,
launchState: failed ? 'failed_to_start' : confirmed ? 'confirmed_alive' : 'starting', launchState: failed ? 'failed_to_start' : confirmed ? 'confirmed_alive' : 'starting',
agentToolAccepted: true, agentToolAccepted: true,
runtimeAlive: confirmed, runtimeAlive: false,
bootstrapConfirmed: confirmed, bootstrapConfirmed: confirmed,
hardFailure: failed, hardFailure: failed,
hardFailureReason: failed hardFailureReason: failed
@ -620,7 +620,7 @@ export function normalizePersistedLaunchSnapshot(
typeof maybeLegacy.leadSessionId === 'string' && maybeLegacy.leadSessionId.trim().length > 0 typeof maybeLegacy.leadSessionId === 'string' && maybeLegacy.leadSessionId.trim().length > 0
? maybeLegacy.leadSessionId.trim() ? maybeLegacy.leadSessionId.trim()
: undefined, : undefined,
launchPhase: 'finished', launchPhase: 'reconciled',
members, members,
updatedAt, updatedAt,
}); });

View file

@ -24,6 +24,18 @@ export function getTeamLaunchSummaryPath(teamName: string): string {
return path.join(getTeamsBasePath(), teamName, TEAM_LAUNCH_SUMMARY_FILE); return path.join(getTeamsBasePath(), teamName, TEAM_LAUNCH_SUMMARY_FILE);
} }
async function isMissingTeamDirectoryWriteRace(teamName: string, error: unknown): Promise<boolean> {
if ((error as NodeJS.ErrnoException).code !== 'ENOENT') {
return false;
}
try {
await fs.promises.access(path.dirname(getTeamLaunchStatePath(teamName)));
return false;
} catch {
return true;
}
}
export class TeamLaunchStateStore { export class TeamLaunchStateStore {
async read(teamName: string): Promise<PersistedTeamLaunchSnapshot | null> { async read(teamName: string): Promise<PersistedTeamLaunchSnapshot | null> {
const targetPath = getTeamLaunchStatePath(teamName); const targetPath = getTeamLaunchStatePath(teamName);
@ -50,6 +62,9 @@ export class TeamLaunchStateStore {
`${JSON.stringify(createPersistedLaunchSummaryProjection(snapshot), null, 2)}\n` `${JSON.stringify(createPersistedLaunchSummaryProjection(snapshot), null, 2)}\n`
); );
} catch (error) { } catch (error) {
if (await isMissingTeamDirectoryWriteRace(teamName, error)) {
return;
}
logger.warn( logger.warn(
`[${teamName}] Failed to persist launch-state: ${ `[${teamName}] Failed to persist launch-state: ${
error instanceof Error ? error.message : String(error) error instanceof Error ? error.message : String(error)

View file

@ -1522,6 +1522,10 @@ function matchesTeamMemberIdentity(leftName: string, rightName: string): boolean
); );
} }
function matchesObservedMemberNameForExpected(observedName: string, expectedName: string): boolean {
return matchesMemberNameOrBase(observedName, expectedName);
}
function matchesExactTeamMemberName(candidateName: string, memberName: string): boolean { function matchesExactTeamMemberName(candidateName: string, memberName: string): boolean {
const left = candidateName.trim().toLowerCase(); const left = candidateName.trim().toLowerCase();
const right = memberName.trim().toLowerCase(); const right = memberName.trim().toLowerCase();
@ -1534,6 +1538,7 @@ interface MemberSpawnInboxCursor {
} }
type LeadInboxMemberSpawnMessage = InboxMessage & { messageId: string }; type LeadInboxMemberSpawnMessage = InboxMessage & { messageId: string };
type LeadInboxLaunchReconcileMessage = Pick<InboxMessage, 'from' | 'text' | 'timestamp'>;
function compareMemberSpawnInboxCursor( function compareMemberSpawnInboxCursor(
left: MemberSpawnInboxCursor, left: MemberSpawnInboxCursor,
@ -4148,7 +4153,9 @@ export class TeamProvisioningService {
private getMixedSecondaryLaunchPhase(run: ProvisioningRun): PersistedTeamLaunchPhase { private getMixedSecondaryLaunchPhase(run: ProvisioningRun): PersistedTeamLaunchPhase {
return (run.mixedSecondaryLanes ?? []).some( return (run.mixedSecondaryLanes ?? []).some(
(lane) => lane.state !== 'finished' || lane.result?.teamLaunchState === 'partial_pending' (lane) =>
(!lane.result && lane.state !== 'finished') ||
lane.result?.teamLaunchState === 'partial_pending'
) )
? 'active' ? 'active'
: 'finished'; : 'finished';
@ -5047,7 +5054,7 @@ export class TeamProvisioningService {
} }
const matches = expectedMembers.filter((memberName) => const matches = expectedMembers.filter((memberName) =>
matchesTeamMemberIdentity(memberName, trimmedCandidate) matchesObservedMemberNameForExpected(trimmedCandidate, memberName)
); );
return matches.length === 1 ? (matches[0] ?? null) : null; return matches.length === 1 ? (matches[0] ?? null) : null;
} }
@ -6510,7 +6517,9 @@ export class TeamProvisioningService {
launchPhase: run.provisioningComplete ? 'finished' : 'active', launchPhase: run.provisioningComplete ? 'finished' : 'active',
statuses: this.buildRuntimeSpawnStatusRecord(run), statuses: this.buildRuntimeSpawnStatusRecord(run),
}); });
const snapshot = liveSnapshot ?? persisted; const rawSnapshot = liveSnapshot ?? persisted;
const metaMembers = await this.membersMetaStore.getMembers(teamName).catch(() => []);
const snapshot = this.filterRemovedMembersFromLaunchSnapshot(rawSnapshot, metaMembers);
const statuses = await this.attachLiveRuntimeMetadataToStatuses( const statuses = await this.attachLiveRuntimeMetadataToStatuses(
teamName, teamName,
snapshotToMemberSpawnStatuses(snapshot) snapshotToMemberSpawnStatuses(snapshot)
@ -11643,7 +11652,7 @@ export class TeamProvisioningService {
? memberName ? memberName
: (() => { : (() => {
const matches = Object.keys(nextStatuses).filter((candidateName) => const matches = Object.keys(nextStatuses).filter((candidateName) =>
matchesTeamMemberIdentity(candidateName, memberName) matchesObservedMemberNameForExpected(memberName, candidateName)
); );
return matches.length === 1 ? matches[0] : null; return matches.length === 1 ? matches[0] : null;
})(); })();
@ -12664,6 +12673,70 @@ export class TeamProvisioningService {
return hasMixedPersistedLaunchMetadata(snapshot); return hasMixedPersistedLaunchMetadata(snapshot);
} }
private hasMixedSecondaryLaunchMetadata(snapshot: PersistedTeamLaunchSnapshot | null): boolean {
if (!snapshot) {
return false;
}
return Object.values(snapshot.members).some(
(member) =>
member?.laneKind === 'secondary' ||
(typeof member?.laneId === 'string' && member.laneId.startsWith('secondary:'))
);
}
private hasPrimaryOnlyLaneAwareLaunchMetadata(
snapshot: PersistedTeamLaunchSnapshot | null
): boolean {
if (!snapshot || this.hasMixedSecondaryLaunchMetadata(snapshot)) {
return false;
}
return Object.values(snapshot.members).some(
(member) =>
Boolean(member?.laneId) ||
Boolean(member?.laneKind) ||
Boolean(member?.laneOwnerProviderId) ||
Boolean(member?.launchIdentity)
);
}
private hasLeadInboxLaunchReconcileHeartbeat(
snapshot: PersistedTeamLaunchSnapshot,
messages: readonly LeadInboxLaunchReconcileMessage[]
): boolean {
const expectedMembers = this.getPersistedLaunchMemberNames(snapshot);
if (expectedMembers.length === 0 || messages.length === 0) {
return false;
}
return messages.some((message) => {
if (
typeof message.from !== 'string' ||
typeof message.text !== 'string' ||
typeof message.timestamp !== 'string' ||
!isMeaningfulBootstrapCheckInMessage(message.text)
) {
return false;
}
const expected = this.resolveExpectedLaunchMemberName(expectedMembers, message.from);
if (!expected) {
return false;
}
const current = snapshot.members[expected];
const firstAcceptedAt = current?.firstSpawnAcceptedAt
? Date.parse(current.firstSpawnAcceptedAt)
: NaN;
const messageTs = Date.parse(message.timestamp);
return (
!Number.isFinite(firstAcceptedAt) ||
!Number.isFinite(messageTs) ||
messageTs >= firstAcceptedAt
);
});
}
private shouldRecoverStalePersistedMixedLaunchSnapshot( private shouldRecoverStalePersistedMixedLaunchSnapshot(
snapshot: PersistedTeamLaunchSnapshot snapshot: PersistedTeamLaunchSnapshot
): boolean { ): boolean {
@ -12701,13 +12774,16 @@ export class TeamProvisioningService {
return null; return null;
} }
if (snapshot.teamLaunchState === 'clean_success' && launchPhase !== 'active') { const metaMembers = await this.membersMetaStore.getMembers(run.teamName).catch(() => []);
const filteredSnapshot = this.filterRemovedMembersFromLaunchSnapshot(snapshot, metaMembers);
if (filteredSnapshot.teamLaunchState === 'clean_success' && launchPhase !== 'active') {
await this.clearPersistedLaunchState(run.teamName); await this.clearPersistedLaunchState(run.teamName);
return null; return null;
} }
await this.launchStateStore.write(run.teamName, snapshot); await this.launchStateStore.write(run.teamName, filteredSnapshot);
return snapshot; return filteredSnapshot;
} }
private async launchSingleMixedSecondaryLane( private async launchSingleMixedSecondaryLane(
@ -12742,6 +12818,7 @@ export class TeamProvisioningService {
lane.warnings = []; lane.warnings = [];
lane.diagnostics = [message]; lane.diagnostics = [message];
await this.publishMixedSecondaryLaneStatusChange(run, lane); await this.publishMixedSecondaryLaneStatusChange(run, lane);
lane.state = 'finished';
return; return;
} }
@ -12802,7 +12879,6 @@ export class TeamProvisioningService {
this.deleteSecondaryRuntimeRun(run.teamName, lane.laneId); this.deleteSecondaryRuntimeRun(run.teamName, lane.laneId);
return; return;
} }
lane.state = 'finished';
lane.result = result; lane.result = result;
lane.warnings = [...result.warnings]; lane.warnings = [...result.warnings];
lane.diagnostics = [...migration.diagnostics, ...result.diagnostics]; lane.diagnostics = [...migration.diagnostics, ...result.diagnostics];
@ -12829,7 +12905,6 @@ export class TeamProvisioningService {
return; return;
} }
const message = error instanceof Error ? error.message : String(error); const message = error instanceof Error ? error.message : String(error);
lane.state = 'finished';
lane.result = { lane.result = {
runId: lane.runId, runId: lane.runId,
teamName: run.teamName, teamName: run.teamName,
@ -12864,6 +12939,7 @@ export class TeamProvisioningService {
} }
await this.publishMixedSecondaryLaneStatusChange(run, lane); await this.publishMixedSecondaryLaneStatusChange(run, lane);
lane.state = 'finished';
} }
private async stopSingleMixedSecondaryRuntimeLane( private async stopSingleMixedSecondaryRuntimeLane(
@ -12931,7 +13007,6 @@ export class TeamProvisioningService {
logger.warn( logger.warn(
`[${run.teamName}] OpenCode secondary lane ${lane.laneId} crashed during launch orchestration: ${message}` `[${run.teamName}] OpenCode secondary lane ${lane.laneId} crashed during launch orchestration: ${message}`
); );
lane.state = 'finished';
lane.result = createUnexpectedMixedSecondaryLaneFailureResult({ lane.result = createUnexpectedMixedSecondaryLaneFailureResult({
runId: lane.runId ?? randomUUID(), runId: lane.runId ?? randomUUID(),
teamName: run.teamName, teamName: run.teamName,
@ -12949,6 +13024,7 @@ export class TeamProvisioningService {
}).catch(() => undefined); }).catch(() => undefined);
this.deleteSecondaryRuntimeRun(run.teamName, lane.laneId); this.deleteSecondaryRuntimeRun(run.teamName, lane.laneId);
await this.publishMixedSecondaryLaneStatusChange(run, lane).catch(() => undefined); await this.publishMixedSecondaryLaneStatusChange(run, lane).catch(() => undefined);
lane.state = 'finished';
} }
})(); })();
} }
@ -13010,7 +13086,7 @@ export class TeamProvisioningService {
): Promise<PersistedTeamLaunchSnapshot | null> { ): Promise<PersistedTeamLaunchSnapshot | null> {
if ( if (
persistedSnapshot && persistedSnapshot &&
this.hasMixedLaunchMetadata(persistedSnapshot) && this.hasMixedSecondaryLaunchMetadata(persistedSnapshot) &&
!this.shouldRecoverStalePersistedMixedLaunchSnapshot(persistedSnapshot) !this.shouldRecoverStalePersistedMixedLaunchSnapshot(persistedSnapshot)
) { ) {
return persistedSnapshot; return persistedSnapshot;
@ -13245,6 +13321,39 @@ export class TeamProvisioningService {
} }
} }
private async readLeadInboxMessagesForLaunchReconcile(
teamName: string,
leadName: string
): Promise<LeadInboxLaunchReconcileMessage[]> {
const inboxPath = path.join(getTeamsBasePath(), teamName, 'inboxes', `${leadName}.json`);
try {
const raw = await tryReadRegularFileUtf8(inboxPath, {
timeoutMs: TEAM_JSON_READ_TIMEOUT_MS,
maxBytes: TEAM_INBOX_MAX_BYTES,
});
if (!raw) {
return [];
}
const parsed = JSON.parse(raw) as unknown;
if (!Array.isArray(parsed)) {
return [];
}
return parsed.flatMap((item): LeadInboxLaunchReconcileMessage[] => {
if (!item || typeof item !== 'object') {
return [];
}
const row = item as Partial<InboxMessage>;
return typeof row.from === 'string' &&
typeof row.text === 'string' &&
typeof row.timestamp === 'string'
? [{ from: row.from, text: row.text, timestamp: row.timestamp }]
: [];
});
} catch {
return [];
}
}
private async reconcilePersistedLaunchState(teamName: string): Promise<{ private async reconcilePersistedLaunchState(teamName: string): Promise<{
snapshot: ReturnType<typeof createPersistedLaunchSnapshot> | null; snapshot: ReturnType<typeof createPersistedLaunchSnapshot> | null;
statuses: Record<string, MemberSpawnStatusEntry>; statuses: Record<string, MemberSpawnStatusEntry>;
@ -13310,11 +13419,19 @@ export class TeamProvisioningService {
// best-effort // best-effort
} }
let leadInboxMessages: Awaited<ReturnType<TeamInboxReader['getMessagesFor']>> = []; const leadInboxMessages = await this.readLeadInboxMessagesForLaunchReconcile(
try { teamName,
leadInboxMessages = await this.inboxReader.getMessagesFor(teamName, leadName); leadName
} catch { );
// best-effort
if (
this.hasPrimaryOnlyLaneAwareLaunchMetadata(filteredPersisted) &&
!this.hasLeadInboxLaunchReconcileHeartbeat(filteredPersisted, leadInboxMessages)
) {
return {
snapshot: filteredPersisted,
statuses: snapshotToMemberSpawnStatuses(filteredPersisted),
};
} }
const liveAgentNames = await this.getLiveTeamAgentNames(teamName); const liveAgentNames = await this.getLiveTeamAgentNames(teamName);
@ -13347,11 +13464,12 @@ export class TeamProvisioningService {
current.lastHeartbeatAt = current.lastHeartbeatAt ?? bootstrapMember.lastHeartbeatAt; current.lastHeartbeatAt = current.lastHeartbeatAt ?? bootstrapMember.lastHeartbeatAt;
} }
const matchedConfigNames = [...configMembers].filter((name) => const matchedConfigNames = [...configMembers].filter((name) =>
matchesTeamMemberIdentity(name, expected) matchesObservedMemberNameForExpected(name, expected)
); );
const runtimeAlive = [...liveAgentNames].some((name) => const observedRuntimeAlive = [...liveAgentNames].some((name) =>
matchesTeamMemberIdentity(name, expected) matchesObservedMemberNameForExpected(name, expected)
); );
const runtimeAlive = current.runtimeAlive === true || observedRuntimeAlive;
const heartbeatMessage = leadInboxMessages.find((message) => { const heartbeatMessage = leadInboxMessages.find((message) => {
if ( if (
typeof message.from !== 'string' || typeof message.from !== 'string' ||
@ -13384,7 +13502,7 @@ export class TeamProvisioningService {
const acceptedAtMs = const acceptedAtMs =
current.firstSpawnAcceptedAt != null ? Date.parse(current.firstSpawnAcceptedAt) : NaN; current.firstSpawnAcceptedAt != null ? Date.parse(current.firstSpawnAcceptedAt) : NaN;
current.runtimeAlive = runtimeAlive; current.runtimeAlive = runtimeAlive;
current.lastRuntimeAliveAt = runtimeAlive ? now : current.lastRuntimeAliveAt; current.lastRuntimeAliveAt = observedRuntimeAlive ? now : current.lastRuntimeAliveAt;
current.sources = { current.sources = {
...(current.sources ?? {}), ...(current.sources ?? {}),
processAlive: runtimeAlive || undefined, processAlive: runtimeAlive || undefined,
@ -13463,7 +13581,7 @@ export class TeamProvisioningService {
teamName, teamName,
expectedMembers: persistedMemberNames, expectedMembers: persistedMemberNames,
leadSessionId: filteredPersisted.leadSessionId, leadSessionId: filteredPersisted.leadSessionId,
launchPhase: filteredPersisted.launchPhase === 'active' ? 'active' : 'reconciled', launchPhase: filteredPersisted.launchPhase,
members: nextMembers, members: nextMembers,
updatedAt: now, updatedAt: now,
}); });

View file

@ -2711,8 +2711,7 @@ describe('TeamProvisioningService', () => {
await (svc as any).launchMixedSecondaryLaneIfNeeded(run); await (svc as any).launchMixedSecondaryLaneIfNeeded(run);
await vi.waitFor(async () => { await vi.waitFor(async () => {
expect(adapterLaunch).toHaveBeenCalledTimes(1); expect(adapterLaunch).toHaveBeenCalledTimes(1);
await expect(readOpenCodeRuntimeLaneIndex(tempTeamsBase, teamName)).resolves.toMatchObject( await expect(readOpenCodeRuntimeLaneIndex(tempTeamsBase, teamName)).resolves.toMatchObject({
{
lanes: { lanes: {
'secondary:opencode:bob': { 'secondary:opencode:bob': {
state: 'degraded', state: 'degraded',
@ -2721,8 +2720,7 @@ describe('TeamProvisioningService', () => {
]), ]),
}, },
}, },
} });
);
}); });
}); });
@ -4429,6 +4427,7 @@ describe('TeamProvisioningService', () => {
}; };
const membersMetaStore = { const membersMetaStore = {
writeMembers: vi.fn(async () => {}), writeMembers: vi.fn(async () => {}),
getMembers: vi.fn(async () => []),
getMeta: vi.fn(async () => null), getMeta: vi.fn(async () => null),
}; };
const teamMetaStore = { const teamMetaStore = {
@ -4534,7 +4533,9 @@ describe('TeamProvisioningService', () => {
stdio: ['pipe', 'pipe', 'pipe'], stdio: ['pipe', 'pipe', 'pipe'],
}); });
const spawnArgs = spawnCall?.[1] as string[]; const spawnArgs = spawnCall?.[1] as string[];
expect(spawnArgs).toEqual(expect.arrayContaining(['--model', 'gpt-5.4', '--effort', 'medium'])); expect(spawnArgs).toEqual(
expect.arrayContaining(['--model', 'gpt-5.4', '--effort', 'medium'])
);
const bootstrapSpec = readBootstrapSpecFromSpawnArgs(spawnArgs); const bootstrapSpec = readBootstrapSpecFromSpawnArgs(spawnArgs);
expect(bootstrapSpec).toMatchObject({ expect(bootstrapSpec).toMatchObject({
@ -4682,11 +4683,18 @@ describe('TeamProvisioningService', () => {
); );
const config = JSON.parse( const config = JSON.parse(
fs.readFileSync(path.join(tempTeamsBase, 'safe-opencode-only-launch', 'config.json'), 'utf8') fs.readFileSync(
path.join(tempTeamsBase, 'safe-opencode-only-launch', 'config.json'),
'utf8'
)
) as { members: Array<{ name: string; providerId?: string; model?: string }> }; ) as { members: Array<{ name: string; providerId?: string; model?: string }> };
expect(config.members).toEqual( expect(config.members).toEqual(
expect.arrayContaining([ expect.arrayContaining([
expect.objectContaining({ name: 'team-lead', providerId: 'opencode', model: 'big-pickle' }), expect.objectContaining({
name: 'team-lead',
providerId: 'opencode',
model: 'big-pickle',
}),
expect.objectContaining({ expect.objectContaining({
name: 'bob', name: 'bob',
providerId: 'opencode', providerId: 'opencode',
@ -4895,7 +4903,9 @@ describe('TeamProvisioningService', () => {
status: 'online', status: 'online',
launchState: 'confirmed_alive', launchState: 'confirmed_alive',
}); });
expect(publicStatuses.expectedMembers).toEqual(expect.arrayContaining(['alice', 'bob', 'tom'])); expect(publicStatuses.expectedMembers).toEqual(
expect.arrayContaining(['alice', 'bob', 'tom'])
);
await svc.cancelProvisioning(runId); await svc.cancelProvisioning(runId);
}); });
@ -7018,7 +7028,11 @@ describe('TeamProvisioningService', () => {
write: vi.fn(async () => {}), write: vi.fn(async () => {}),
clear: vi.fn(async () => {}), clear: vi.fn(async () => {}),
}; };
vi.spyOn((svc as any).inboxReader, 'getMessagesFor').mockResolvedValue([ fs.mkdirSync(path.join(tempTeamsBase, teamName, 'inboxes'), { recursive: true });
fs.writeFileSync(
path.join(tempTeamsBase, teamName, 'inboxes', 'team-lead.json'),
JSON.stringify(
[
{ {
from: 'alice-2', from: 'alice-2',
text: 'heartbeat', text: 'heartbeat',
@ -7026,7 +7040,11 @@ describe('TeamProvisioningService', () => {
messageId: 'msg-suffixed-reconcile', messageId: 'msg-suffixed-reconcile',
read: false, read: false,
}, },
]); ],
null,
2
)
);
(svc as any).getLiveTeamAgentNames = vi.fn(async () => new Set<string>()); (svc as any).getLiveTeamAgentNames = vi.fn(async () => new Set<string>());
const result = await (svc as any).reconcilePersistedLaunchState(teamName); const result = await (svc as any).reconcilePersistedLaunchState(teamName);