perf(team): reduce launch status IO churn
This commit is contained in:
parent
fcb9179990
commit
ad7c4e24ad
2 changed files with 589 additions and 15 deletions
|
|
@ -313,6 +313,11 @@ interface OpenCodeRuntimeControlAck {
|
||||||
observedAt: string;
|
observedAt: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
interface LaunchStateWriteResult {
|
||||||
|
snapshot: PersistedTeamLaunchSnapshot;
|
||||||
|
wrote: boolean;
|
||||||
|
}
|
||||||
|
|
||||||
type BootstrapTranscriptOutcome =
|
type BootstrapTranscriptOutcome =
|
||||||
| {
|
| {
|
||||||
kind: 'success';
|
kind: 'success';
|
||||||
|
|
@ -336,6 +341,7 @@ import type {
|
||||||
LeadContextUsage,
|
LeadContextUsage,
|
||||||
MemberLaunchState,
|
MemberLaunchState,
|
||||||
MemberSpawnLivenessSource,
|
MemberSpawnLivenessSource,
|
||||||
|
MemberSpawnStatusesSnapshot,
|
||||||
MemberSpawnStatus,
|
MemberSpawnStatus,
|
||||||
MemberSpawnStatusEntry,
|
MemberSpawnStatusEntry,
|
||||||
PersistedTeamLaunchMemberState,
|
PersistedTeamLaunchMemberState,
|
||||||
|
|
@ -4505,6 +4511,8 @@ export class TeamProvisioningService {
|
||||||
private static readonly SAME_TEAM_RUN_START_SKEW_MS = 1_000;
|
private static readonly SAME_TEAM_RUN_START_SKEW_MS = 1_000;
|
||||||
private static readonly SAME_TEAM_PERSIST_RETRY_MS = 2_000;
|
private static readonly SAME_TEAM_PERSIST_RETRY_MS = 2_000;
|
||||||
private static readonly AGENT_RUNTIME_SNAPSHOT_CACHE_TTL_MS = 2_000;
|
private static readonly AGENT_RUNTIME_SNAPSHOT_CACHE_TTL_MS = 2_000;
|
||||||
|
private static readonly MEMBER_SPAWN_STATUS_SNAPSHOT_CACHE_TTL_MS = 500;
|
||||||
|
private static readonly LAUNCH_STATE_NOOP_REFRESH_MS = 15_000;
|
||||||
|
|
||||||
private readonly runs = new Map<string, ProvisioningRun>();
|
private readonly runs = new Map<string, ProvisioningRun>();
|
||||||
private readonly provisioningRunByTeam = new Map<string, string>();
|
private readonly provisioningRunByTeam = new Map<string, string>();
|
||||||
|
|
@ -4590,8 +4598,27 @@ export class TeamProvisioningService {
|
||||||
}
|
}
|
||||||
>();
|
>();
|
||||||
private readonly runtimeSnapshotCacheGenerationByTeam = new Map<string, number>();
|
private readonly runtimeSnapshotCacheGenerationByTeam = new Map<string, number>();
|
||||||
|
private readonly memberSpawnStatusesSnapshotCache = new Map<
|
||||||
|
string,
|
||||||
|
{
|
||||||
|
expiresAtMs: number;
|
||||||
|
generation: number;
|
||||||
|
runId: string;
|
||||||
|
snapshot: MemberSpawnStatusesSnapshot;
|
||||||
|
}
|
||||||
|
>();
|
||||||
|
private readonly memberSpawnStatusesInFlightByTeam = new Map<
|
||||||
|
string,
|
||||||
|
{
|
||||||
|
generationAtStart: number;
|
||||||
|
runIdAtStart: string;
|
||||||
|
promise: Promise<MemberSpawnStatusesSnapshot>;
|
||||||
|
}
|
||||||
|
>();
|
||||||
|
private readonly memberSpawnStatusesCacheGenerationByTeam = new Map<string, number>();
|
||||||
private readonly launchStateStore = new TeamLaunchStateStore();
|
private readonly launchStateStore = new TeamLaunchStateStore();
|
||||||
private readonly launchStateStoreQueue = new Map<string, Promise<unknown>>();
|
private readonly launchStateStoreQueue = new Map<string, Promise<unknown>>();
|
||||||
|
private readonly launchStateWrittenRunIdByTeam = new Map<string, string>();
|
||||||
private readonly memberLogsFinder: TeamMemberLogsFinder;
|
private readonly memberLogsFinder: TeamMemberLogsFinder;
|
||||||
private readonly transcriptProjectResolver: TeamTranscriptProjectResolver;
|
private readonly transcriptProjectResolver: TeamTranscriptProjectResolver;
|
||||||
private teamChangeEmitter: ((event: TeamChangeEvent) => void) | null = null;
|
private teamChangeEmitter: ((event: TeamChangeEvent) => void) | null = null;
|
||||||
|
|
@ -4676,6 +4703,19 @@ export class TeamProvisioningService {
|
||||||
return this.runtimeSnapshotCacheGenerationByTeam.get(teamName) ?? 0;
|
return this.runtimeSnapshotCacheGenerationByTeam.get(teamName) ?? 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private getMemberSpawnStatusesCacheGeneration(teamName: string): number {
|
||||||
|
return this.memberSpawnStatusesCacheGenerationByTeam.get(teamName) ?? 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
private invalidateMemberSpawnStatusesCache(teamName: string): void {
|
||||||
|
this.memberSpawnStatusesCacheGenerationByTeam.set(
|
||||||
|
teamName,
|
||||||
|
this.getMemberSpawnStatusesCacheGeneration(teamName) + 1
|
||||||
|
);
|
||||||
|
this.memberSpawnStatusesSnapshotCache.delete(teamName);
|
||||||
|
this.memberSpawnStatusesInFlightByTeam.delete(teamName);
|
||||||
|
}
|
||||||
|
|
||||||
private invalidateRuntimeSnapshotCaches(teamName: string): void {
|
private invalidateRuntimeSnapshotCaches(teamName: string): void {
|
||||||
this.runtimeSnapshotCacheGenerationByTeam.set(
|
this.runtimeSnapshotCacheGenerationByTeam.set(
|
||||||
teamName,
|
teamName,
|
||||||
|
|
@ -4687,6 +4727,27 @@ export class TeamProvisioningService {
|
||||||
this.liveTeamAgentRuntimeMetadataInFlightByTeam.delete(teamName);
|
this.liveTeamAgentRuntimeMetadataInFlightByTeam.delete(teamName);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private cloneMemberSpawnStatusesSnapshot(
|
||||||
|
snapshot: MemberSpawnStatusesSnapshot
|
||||||
|
): MemberSpawnStatusesSnapshot {
|
||||||
|
return {
|
||||||
|
...snapshot,
|
||||||
|
statuses: Object.fromEntries(
|
||||||
|
Object.entries(snapshot.statuses).map(([memberName, entry]) => [
|
||||||
|
memberName,
|
||||||
|
{
|
||||||
|
...entry,
|
||||||
|
...(entry.pendingPermissionRequestIds
|
||||||
|
? { pendingPermissionRequestIds: [...entry.pendingPermissionRequestIds] }
|
||||||
|
: {}),
|
||||||
|
},
|
||||||
|
])
|
||||||
|
),
|
||||||
|
...(snapshot.expectedMembers ? { expectedMembers: [...snapshot.expectedMembers] } : {}),
|
||||||
|
...(snapshot.summary ? { summary: { ...snapshot.summary } } : {}),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
private cloneLiveTeamAgentRuntimeMetadata(
|
private cloneLiveTeamAgentRuntimeMetadata(
|
||||||
metadata: ReadonlyMap<string, LiveTeamAgentRuntimeMetadata>
|
metadata: ReadonlyMap<string, LiveTeamAgentRuntimeMetadata>
|
||||||
): Map<string, LiveTeamAgentRuntimeMetadata> {
|
): Map<string, LiveTeamAgentRuntimeMetadata> {
|
||||||
|
|
@ -10409,6 +10470,68 @@ export class TeamProvisioningService {
|
||||||
return readPersistedStatuses(runId);
|
return readPersistedStatuses(runId);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (!this.shouldCacheMemberSpawnStatusesSnapshot(run)) {
|
||||||
|
return this.buildMemberSpawnStatusesSnapshotForRun(run);
|
||||||
|
}
|
||||||
|
|
||||||
|
const generationAtStart = this.getMemberSpawnStatusesCacheGeneration(teamName);
|
||||||
|
const cached = this.memberSpawnStatusesSnapshotCache.get(teamName);
|
||||||
|
if (
|
||||||
|
cached &&
|
||||||
|
cached.expiresAtMs > Date.now() &&
|
||||||
|
cached.runId === run.runId &&
|
||||||
|
cached.generation === generationAtStart
|
||||||
|
) {
|
||||||
|
return this.cloneMemberSpawnStatusesSnapshot(cached.snapshot);
|
||||||
|
}
|
||||||
|
|
||||||
|
const existingRequest = this.memberSpawnStatusesInFlightByTeam.get(teamName);
|
||||||
|
if (
|
||||||
|
existingRequest &&
|
||||||
|
existingRequest.generationAtStart === generationAtStart &&
|
||||||
|
existingRequest.runIdAtStart === run.runId
|
||||||
|
) {
|
||||||
|
const snapshot = await existingRequest.promise;
|
||||||
|
if (
|
||||||
|
this.getMemberSpawnStatusesCacheGeneration(teamName) === generationAtStart &&
|
||||||
|
this.getTrackedRunId(teamName) === run.runId
|
||||||
|
) {
|
||||||
|
return this.cloneMemberSpawnStatusesSnapshot(snapshot);
|
||||||
|
}
|
||||||
|
return this.getMemberSpawnStatuses(teamName);
|
||||||
|
}
|
||||||
|
|
||||||
|
const request = this.buildMemberSpawnStatusesSnapshotForRun(run, generationAtStart).finally(
|
||||||
|
() => {
|
||||||
|
if (this.memberSpawnStatusesInFlightByTeam.get(teamName)?.promise === request) {
|
||||||
|
this.memberSpawnStatusesInFlightByTeam.delete(teamName);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
);
|
||||||
|
this.memberSpawnStatusesInFlightByTeam.set(teamName, {
|
||||||
|
generationAtStart,
|
||||||
|
runIdAtStart: run.runId,
|
||||||
|
promise: request,
|
||||||
|
});
|
||||||
|
const snapshot = await request;
|
||||||
|
if (
|
||||||
|
this.getMemberSpawnStatusesCacheGeneration(teamName) === generationAtStart &&
|
||||||
|
this.getTrackedRunId(teamName) === run.runId
|
||||||
|
) {
|
||||||
|
return this.cloneMemberSpawnStatusesSnapshot(snapshot);
|
||||||
|
}
|
||||||
|
return this.getMemberSpawnStatuses(teamName);
|
||||||
|
}
|
||||||
|
|
||||||
|
private shouldCacheMemberSpawnStatusesSnapshot(run: ProvisioningRun): boolean {
|
||||||
|
return run.isLaunch === true && run.provisioningComplete !== true;
|
||||||
|
}
|
||||||
|
|
||||||
|
private async buildMemberSpawnStatusesSnapshotForRun(
|
||||||
|
run: ProvisioningRun,
|
||||||
|
generationAtStart?: number
|
||||||
|
): Promise<MemberSpawnStatusesSnapshot> {
|
||||||
|
const teamName = run.teamName;
|
||||||
await this.refreshMemberSpawnStatusesFromLeadInbox(run);
|
await this.refreshMemberSpawnStatusesFromLeadInbox(run);
|
||||||
await this.maybeAuditMemberSpawnStatuses(run);
|
await this.maybeAuditMemberSpawnStatuses(run);
|
||||||
await this.persistLaunchStateSnapshot(run, run.provisioningComplete ? 'finished' : 'active');
|
await this.persistLaunchStateSnapshot(run, run.provisioningComplete ? 'finished' : 'active');
|
||||||
|
|
@ -10428,23 +10551,37 @@ export class TeamProvisioningService {
|
||||||
});
|
});
|
||||||
const rawSnapshot = liveSnapshot ?? persisted;
|
const rawSnapshot = liveSnapshot ?? persisted;
|
||||||
const metaMembers = await this.membersMetaStore.getMembers(teamName).catch(() => []);
|
const metaMembers = await this.membersMetaStore.getMembers(teamName).catch(() => []);
|
||||||
const snapshot = this.filterRemovedMembersFromLaunchSnapshot(rawSnapshot, metaMembers);
|
const launchSnapshot = this.filterRemovedMembersFromLaunchSnapshot(rawSnapshot, metaMembers);
|
||||||
const statuses = await this.attachLiveRuntimeMetadataToStatuses(
|
const statuses = await this.attachLiveRuntimeMetadataToStatuses(
|
||||||
teamName,
|
teamName,
|
||||||
snapshotToMemberSpawnStatuses(snapshot)
|
snapshotToMemberSpawnStatuses(launchSnapshot)
|
||||||
);
|
);
|
||||||
const expectedMembers = this.getPersistedLaunchMemberNames(snapshot);
|
const expectedMembers = this.getPersistedLaunchMemberNames(launchSnapshot);
|
||||||
const summary = summarizeMemberSpawnStatusRecord(expectedMembers, statuses);
|
const summary = summarizeMemberSpawnStatusRecord(expectedMembers, statuses);
|
||||||
return {
|
const spawnSnapshot: MemberSpawnStatusesSnapshot = {
|
||||||
statuses,
|
statuses,
|
||||||
runId,
|
runId: run.runId,
|
||||||
teamLaunchState: deriveTeamLaunchAggregateState(summary),
|
teamLaunchState: deriveTeamLaunchAggregateState(summary),
|
||||||
launchPhase: snapshot.launchPhase,
|
launchPhase: launchSnapshot.launchPhase,
|
||||||
expectedMembers,
|
expectedMembers,
|
||||||
updatedAt: snapshot.updatedAt,
|
updatedAt: launchSnapshot.updatedAt,
|
||||||
summary,
|
summary,
|
||||||
source: persisted ? 'merged' : 'live',
|
source: persisted ? 'merged' : 'live',
|
||||||
};
|
};
|
||||||
|
if (
|
||||||
|
generationAtStart != null &&
|
||||||
|
this.shouldCacheMemberSpawnStatusesSnapshot(run) &&
|
||||||
|
this.getMemberSpawnStatusesCacheGeneration(teamName) === generationAtStart &&
|
||||||
|
this.getTrackedRunId(teamName) === run.runId
|
||||||
|
) {
|
||||||
|
this.memberSpawnStatusesSnapshotCache.set(teamName, {
|
||||||
|
expiresAtMs: Date.now() + TeamProvisioningService.MEMBER_SPAWN_STATUS_SNAPSHOT_CACHE_TTL_MS,
|
||||||
|
generation: generationAtStart,
|
||||||
|
runId: run.runId,
|
||||||
|
snapshot: this.cloneMemberSpawnStatusesSnapshot(spawnSnapshot),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
return spawnSnapshot;
|
||||||
}
|
}
|
||||||
|
|
||||||
async getTeamAgentRuntimeSnapshot(teamName: string): Promise<TeamAgentRuntimeSnapshot> {
|
async getTeamAgentRuntimeSnapshot(teamName: string): Promise<TeamAgentRuntimeSnapshot> {
|
||||||
|
|
@ -18041,6 +18178,7 @@ export class TeamProvisioningService {
|
||||||
|
|
||||||
private async clearPersistedLaunchStateNow(teamName: string): Promise<void> {
|
private async clearPersistedLaunchStateNow(teamName: string): Promise<void> {
|
||||||
await this.launchStateStore.clear(teamName);
|
await this.launchStateStore.clear(teamName);
|
||||||
|
this.launchStateWrittenRunIdByTeam.delete(teamName);
|
||||||
await clearBootstrapState(teamName);
|
await clearBootstrapState(teamName);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -18302,15 +18440,17 @@ export class TeamProvisioningService {
|
||||||
teamName: string,
|
teamName: string,
|
||||||
snapshot: PersistedTeamLaunchSnapshot
|
snapshot: PersistedTeamLaunchSnapshot
|
||||||
): Promise<PersistedTeamLaunchSnapshot> {
|
): Promise<PersistedTeamLaunchSnapshot> {
|
||||||
return this.enqueueLaunchStateStoreOperation(teamName, () =>
|
const result = await this.enqueueLaunchStateStoreOperation(teamName, () =>
|
||||||
this.writeLaunchStateSnapshotNow(teamName, snapshot)
|
this.writeLaunchStateSnapshotNow(teamName, snapshot)
|
||||||
);
|
);
|
||||||
|
return result.snapshot;
|
||||||
}
|
}
|
||||||
|
|
||||||
private async writeLaunchStateSnapshotNow(
|
private async writeLaunchStateSnapshotNow(
|
||||||
teamName: string,
|
teamName: string,
|
||||||
snapshot: PersistedTeamLaunchSnapshot
|
snapshot: PersistedTeamLaunchSnapshot,
|
||||||
): Promise<PersistedTeamLaunchSnapshot> {
|
options?: { allowNoopSkip?: boolean; runId?: string }
|
||||||
|
): Promise<LaunchStateWriteResult> {
|
||||||
const previousSnapshot = await this.launchStateStore.read(teamName).catch(() => null);
|
const previousSnapshot = await this.launchStateStore.read(teamName).catch(() => null);
|
||||||
const metaMembers = await this.membersMetaStore.getMembers(teamName).catch(() => []);
|
const metaMembers = await this.membersMetaStore.getMembers(teamName).catch(() => []);
|
||||||
const overlaidSnapshot = await this.applyOpenCodeSecondaryEvidenceOverlay({
|
const overlaidSnapshot = await this.applyOpenCodeSecondaryEvidenceOverlay({
|
||||||
|
|
@ -18319,8 +18459,74 @@ export class TeamProvisioningService {
|
||||||
previousSnapshot,
|
previousSnapshot,
|
||||||
metaMembers,
|
metaMembers,
|
||||||
});
|
});
|
||||||
|
if (
|
||||||
|
options?.allowNoopSkip === true &&
|
||||||
|
typeof options.runId === 'string' &&
|
||||||
|
this.launchStateWrittenRunIdByTeam.get(teamName) === options.runId &&
|
||||||
|
previousSnapshot &&
|
||||||
|
this.areLaunchStateSnapshotsSemanticallyEqual(previousSnapshot, overlaidSnapshot) &&
|
||||||
|
!this.isLaunchStateNoopRefreshDue(previousSnapshot)
|
||||||
|
) {
|
||||||
|
return { snapshot: previousSnapshot, wrote: false };
|
||||||
|
}
|
||||||
await this.launchStateStore.write(teamName, overlaidSnapshot);
|
await this.launchStateStore.write(teamName, overlaidSnapshot);
|
||||||
return overlaidSnapshot;
|
if (typeof options?.runId === 'string') {
|
||||||
|
this.launchStateWrittenRunIdByTeam.set(teamName, options.runId);
|
||||||
|
}
|
||||||
|
return { snapshot: overlaidSnapshot, wrote: true };
|
||||||
|
}
|
||||||
|
|
||||||
|
private isLaunchStateNoopRefreshDue(snapshot: PersistedTeamLaunchSnapshot): boolean {
|
||||||
|
const updatedAtMs = Date.parse(snapshot.updatedAt);
|
||||||
|
return (
|
||||||
|
!Number.isFinite(updatedAtMs) ||
|
||||||
|
Date.now() - updatedAtMs >= TeamProvisioningService.LAUNCH_STATE_NOOP_REFRESH_MS
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
private areLaunchStateSnapshotsSemanticallyEqual(
|
||||||
|
left: PersistedTeamLaunchSnapshot,
|
||||||
|
right: PersistedTeamLaunchSnapshot
|
||||||
|
): boolean {
|
||||||
|
return (
|
||||||
|
JSON.stringify(this.toLaunchStateSemanticValue(left)) ===
|
||||||
|
JSON.stringify(this.toLaunchStateSemanticValue(right))
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
private toLaunchStateSemanticValue(snapshot: PersistedTeamLaunchSnapshot): unknown {
|
||||||
|
const { updatedAt: _updatedAt, members, ...rest } = snapshot;
|
||||||
|
const stableMembers = Object.fromEntries(
|
||||||
|
Object.entries(members)
|
||||||
|
.sort(([leftName], [rightName]) => leftName.localeCompare(rightName))
|
||||||
|
.map(([memberName, member]) => {
|
||||||
|
const {
|
||||||
|
lastEvaluatedAt: _lastEvaluatedAt,
|
||||||
|
lastRuntimeAliveAt: _lastRuntimeAliveAt,
|
||||||
|
...stableMember
|
||||||
|
} = member;
|
||||||
|
return [memberName, stableMember];
|
||||||
|
})
|
||||||
|
);
|
||||||
|
return this.toStableJsonValue({
|
||||||
|
...rest,
|
||||||
|
members: stableMembers,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
private toStableJsonValue(value: unknown): unknown {
|
||||||
|
if (Array.isArray(value)) {
|
||||||
|
return value.map((item) => this.toStableJsonValue(item));
|
||||||
|
}
|
||||||
|
if (!value || typeof value !== 'object') {
|
||||||
|
return value;
|
||||||
|
}
|
||||||
|
return Object.fromEntries(
|
||||||
|
Object.entries(value as Record<string, unknown>)
|
||||||
|
.filter(([, entryValue]) => entryValue !== undefined)
|
||||||
|
.sort(([leftKey], [rightKey]) => leftKey.localeCompare(rightKey))
|
||||||
|
.map(([key, entryValue]) => [key, this.toStableJsonValue(entryValue)])
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
private async enqueueLaunchStateStoreOperation<T>(
|
private async enqueueLaunchStateStoreOperation<T>(
|
||||||
|
|
@ -18657,6 +18863,7 @@ export class TeamProvisioningService {
|
||||||
run: Pick<ProvisioningRun, 'teamName' | 'runId'>,
|
run: Pick<ProvisioningRun, 'teamName' | 'runId'>,
|
||||||
memberName: string
|
memberName: string
|
||||||
) {
|
) {
|
||||||
|
this.invalidateMemberSpawnStatusesCache(run.teamName);
|
||||||
this.teamChangeEmitter?.({
|
this.teamChangeEmitter?.({
|
||||||
type: 'member-spawn',
|
type: 'member-spawn',
|
||||||
teamName: run.teamName,
|
teamName: run.teamName,
|
||||||
|
|
@ -19008,9 +19215,14 @@ export class TeamProvisioningService {
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
const writtenSnapshot = await this.writeLaunchStateSnapshotNow(run.teamName, filteredSnapshot);
|
const writeResult = await this.writeLaunchStateSnapshotNow(run.teamName, filteredSnapshot, {
|
||||||
this.invalidateRuntimeSnapshotCaches(run.teamName);
|
allowNoopSkip: true,
|
||||||
return writtenSnapshot;
|
runId: run.runId,
|
||||||
|
});
|
||||||
|
if (writeResult.wrote) {
|
||||||
|
this.invalidateRuntimeSnapshotCaches(run.teamName);
|
||||||
|
}
|
||||||
|
return writeResult.snapshot;
|
||||||
}
|
}
|
||||||
|
|
||||||
private async launchSingleMixedSecondaryLane(
|
private async launchSingleMixedSecondaryLane(
|
||||||
|
|
@ -23879,6 +24091,7 @@ export class TeamProvisioningService {
|
||||||
}
|
}
|
||||||
if (!hasNewerTrackedRun) {
|
if (!hasNewerTrackedRun) {
|
||||||
this.invalidateRuntimeSnapshotCaches(run.teamName);
|
this.invalidateRuntimeSnapshotCaches(run.teamName);
|
||||||
|
this.invalidateMemberSpawnStatusesCache(run.teamName);
|
||||||
this.leadInboxRelayInFlight.delete(run.teamName);
|
this.leadInboxRelayInFlight.delete(run.teamName);
|
||||||
this.relayedLeadInboxMessageIds.delete(run.teamName);
|
this.relayedLeadInboxMessageIds.delete(run.teamName);
|
||||||
this.pendingCrossTeamFirstReplies.delete(run.teamName);
|
this.pendingCrossTeamFirstReplies.delete(run.teamName);
|
||||||
|
|
|
||||||
|
|
@ -129,7 +129,10 @@ import {
|
||||||
initializeAutoResumeService,
|
initializeAutoResumeService,
|
||||||
} from '@main/services/team/AutoResumeService';
|
} from '@main/services/team/AutoResumeService';
|
||||||
import { getTeamBootstrapStatePath } from '@main/services/team/TeamBootstrapStateReader';
|
import { getTeamBootstrapStatePath } from '@main/services/team/TeamBootstrapStateReader';
|
||||||
import { createPersistedLaunchSnapshot } from '@main/services/team/TeamLaunchStateEvaluator';
|
import {
|
||||||
|
createPersistedLaunchSnapshot,
|
||||||
|
snapshotFromRuntimeMemberStatuses,
|
||||||
|
} from '@main/services/team/TeamLaunchStateEvaluator';
|
||||||
import { getTeamLaunchStatePath } from '@main/services/team/TeamLaunchStateStore';
|
import { getTeamLaunchStatePath } from '@main/services/team/TeamLaunchStateStore';
|
||||||
import { TeamConfigReader } from '@main/services/team/TeamConfigReader';
|
import { TeamConfigReader } from '@main/services/team/TeamConfigReader';
|
||||||
import {
|
import {
|
||||||
|
|
@ -628,6 +631,364 @@ describe('TeamProvisioningService', () => {
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe('member spawn status launch reads', () => {
|
||||||
|
it('coalesces concurrent active launch status reads and serves a short cached follow-up', async () => {
|
||||||
|
const svc = new TeamProvisioningService();
|
||||||
|
const teamName = 'spawn-cache-team';
|
||||||
|
const run = createMemberSpawnRun({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
});
|
||||||
|
run.isLaunch = true;
|
||||||
|
run.progress = {
|
||||||
|
teamName,
|
||||||
|
state: 'launching',
|
||||||
|
message: 'Launching',
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
};
|
||||||
|
run.onProgress = vi.fn();
|
||||||
|
(svc as any).aliveRunByTeam.set(teamName, run.runId);
|
||||||
|
(svc as any).runs.set(run.runId, run);
|
||||||
|
(svc as any).membersMetaStore = { getMembers: vi.fn(async () => []) };
|
||||||
|
(svc as any).launchStateStore = {
|
||||||
|
read: vi.fn(async () => null),
|
||||||
|
write: vi.fn(async () => {}),
|
||||||
|
clear: vi.fn(async () => {}),
|
||||||
|
};
|
||||||
|
(svc as any).getLiveTeamAgentRuntimeMetadata = vi.fn(async () => new Map());
|
||||||
|
const refreshDeferred = createDeferred<void>();
|
||||||
|
const refreshLeadInbox = vi
|
||||||
|
.spyOn(svc as any, 'refreshMemberSpawnStatusesFromLeadInbox')
|
||||||
|
.mockImplementation(async () => refreshDeferred.promise);
|
||||||
|
const auditStatuses = vi
|
||||||
|
.spyOn(svc as any, 'maybeAuditMemberSpawnStatuses')
|
||||||
|
.mockResolvedValue(undefined);
|
||||||
|
const persistSnapshot = vi
|
||||||
|
.spyOn(svc as any, 'persistLaunchStateSnapshot')
|
||||||
|
.mockResolvedValue(null);
|
||||||
|
|
||||||
|
const first = svc.getMemberSpawnStatuses(teamName);
|
||||||
|
const second = svc.getMemberSpawnStatuses(teamName);
|
||||||
|
await Promise.resolve();
|
||||||
|
|
||||||
|
expect(refreshLeadInbox).toHaveBeenCalledTimes(1);
|
||||||
|
refreshDeferred.resolve();
|
||||||
|
const [firstSnapshot, secondSnapshot] = await Promise.all([first, second]);
|
||||||
|
|
||||||
|
expect(firstSnapshot).toEqual(secondSnapshot);
|
||||||
|
expect(auditStatuses).toHaveBeenCalledTimes(1);
|
||||||
|
expect(persistSnapshot).toHaveBeenCalledTimes(1);
|
||||||
|
|
||||||
|
await svc.getMemberSpawnStatuses(teamName);
|
||||||
|
|
||||||
|
expect(refreshLeadInbox).toHaveBeenCalledTimes(1);
|
||||||
|
expect(auditStatuses).toHaveBeenCalledTimes(1);
|
||||||
|
expect(persistSnapshot).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('invalidates the short status cache when a real member-spawn change is emitted', async () => {
|
||||||
|
const svc = new TeamProvisioningService();
|
||||||
|
const teamName = 'spawn-cache-invalidated-team';
|
||||||
|
const run = createMemberSpawnRun({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
});
|
||||||
|
run.isLaunch = true;
|
||||||
|
run.progress = {
|
||||||
|
teamName,
|
||||||
|
state: 'launching',
|
||||||
|
message: 'Launching',
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
};
|
||||||
|
run.onProgress = vi.fn();
|
||||||
|
(svc as any).aliveRunByTeam.set(teamName, run.runId);
|
||||||
|
(svc as any).runs.set(run.runId, run);
|
||||||
|
(svc as any).membersMetaStore = { getMembers: vi.fn(async () => []) };
|
||||||
|
(svc as any).launchStateStore = {
|
||||||
|
read: vi.fn(async () => null),
|
||||||
|
write: vi.fn(async () => {}),
|
||||||
|
clear: vi.fn(async () => {}),
|
||||||
|
};
|
||||||
|
(svc as any).getLiveTeamAgentRuntimeMetadata = vi.fn(async () => new Map());
|
||||||
|
const refreshLeadInbox = vi
|
||||||
|
.spyOn(svc as any, 'refreshMemberSpawnStatusesFromLeadInbox')
|
||||||
|
.mockResolvedValue(undefined);
|
||||||
|
vi.spyOn(svc as any, 'maybeAuditMemberSpawnStatuses').mockResolvedValue(undefined);
|
||||||
|
vi.spyOn(svc as any, 'persistLaunchStateSnapshot').mockResolvedValue(null);
|
||||||
|
|
||||||
|
await svc.getMemberSpawnStatuses(teamName);
|
||||||
|
expect(refreshLeadInbox).toHaveBeenCalledTimes(1);
|
||||||
|
|
||||||
|
(svc as any).setMemberSpawnStatus(
|
||||||
|
run,
|
||||||
|
'alice',
|
||||||
|
'online',
|
||||||
|
undefined,
|
||||||
|
'heartbeat',
|
||||||
|
'2026-05-02T10:00:05.000Z'
|
||||||
|
);
|
||||||
|
await svc.getMemberSpawnStatuses(teamName);
|
||||||
|
|
||||||
|
expect(refreshLeadInbox).toHaveBeenCalledTimes(2);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('retries the owner status request when a member-spawn change lands while it is building', async () => {
|
||||||
|
const svc = new TeamProvisioningService();
|
||||||
|
const teamName = 'spawn-cache-owner-retry-team';
|
||||||
|
const run = createMemberSpawnRun({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
});
|
||||||
|
run.isLaunch = true;
|
||||||
|
run.progress = {
|
||||||
|
teamName,
|
||||||
|
state: 'launching',
|
||||||
|
message: 'Launching',
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
};
|
||||||
|
run.onProgress = vi.fn();
|
||||||
|
(svc as any).aliveRunByTeam.set(teamName, run.runId);
|
||||||
|
(svc as any).runs.set(run.runId, run);
|
||||||
|
(svc as any).membersMetaStore = { getMembers: vi.fn(async () => []) };
|
||||||
|
(svc as any).launchStateStore = {
|
||||||
|
read: vi.fn(async () => null),
|
||||||
|
write: vi.fn(async () => {}),
|
||||||
|
clear: vi.fn(async () => {}),
|
||||||
|
};
|
||||||
|
(svc as any).getLiveTeamAgentRuntimeMetadata = vi.fn(async () => new Map());
|
||||||
|
const firstRefresh = createDeferred<void>();
|
||||||
|
const refreshLeadInbox = vi
|
||||||
|
.spyOn(svc as any, 'refreshMemberSpawnStatusesFromLeadInbox')
|
||||||
|
.mockImplementationOnce(async () => firstRefresh.promise)
|
||||||
|
.mockResolvedValue(undefined);
|
||||||
|
vi.spyOn(svc as any, 'maybeAuditMemberSpawnStatuses').mockResolvedValue(undefined);
|
||||||
|
vi.spyOn(svc as any, 'persistLaunchStateSnapshot').mockResolvedValue(null);
|
||||||
|
|
||||||
|
const pending = svc.getMemberSpawnStatuses(teamName);
|
||||||
|
await Promise.resolve();
|
||||||
|
expect(refreshLeadInbox).toHaveBeenCalledTimes(1);
|
||||||
|
|
||||||
|
(svc as any).setMemberSpawnStatus(
|
||||||
|
run,
|
||||||
|
'alice',
|
||||||
|
'online',
|
||||||
|
undefined,
|
||||||
|
'heartbeat',
|
||||||
|
'2026-05-02T10:00:05.000Z'
|
||||||
|
);
|
||||||
|
firstRefresh.resolve();
|
||||||
|
const result = await pending;
|
||||||
|
|
||||||
|
expect(refreshLeadInbox).toHaveBeenCalledTimes(2);
|
||||||
|
expect(result.statuses.alice).toMatchObject({
|
||||||
|
status: 'online',
|
||||||
|
launchState: 'confirmed_alive',
|
||||||
|
bootstrapConfirmed: true,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe('launch-state no-op persistence guard', () => {
|
||||||
|
it('does not rewrite launch-state or invalidate runtime cache for a recent semantic no-op', async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
vi.setSystemTime(new Date('2026-05-02T10:00:05.000Z'));
|
||||||
|
const svc = new TeamProvisioningService();
|
||||||
|
const teamName = 'launch-state-noop-team';
|
||||||
|
const status = createMemberSpawnStatusEntry({
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
firstSpawnAcceptedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
});
|
||||||
|
const previousSnapshot = snapshotFromRuntimeMemberStatuses({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
launchPhase: 'active',
|
||||||
|
updatedAt: '2026-05-02T10:00:02.000Z',
|
||||||
|
statuses: { alice: status as any },
|
||||||
|
});
|
||||||
|
const run = createMemberSpawnRun({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
memberSpawnStatuses: new Map([['alice', status]]),
|
||||||
|
});
|
||||||
|
run.isLaunch = true;
|
||||||
|
(svc as any).launchStateStore = {
|
||||||
|
read: vi.fn(async () => previousSnapshot),
|
||||||
|
write: vi.fn(async () => {}),
|
||||||
|
clear: vi.fn(async () => {}),
|
||||||
|
};
|
||||||
|
(svc as any).launchStateWrittenRunIdByTeam.set(teamName, run.runId);
|
||||||
|
(svc as any).membersMetaStore = { getMembers: vi.fn(async () => []) };
|
||||||
|
const invalidateRuntime = vi.spyOn(svc as any, 'invalidateRuntimeSnapshotCaches');
|
||||||
|
|
||||||
|
const result = await (svc as any).persistLaunchStateSnapshotNow(run, 'active');
|
||||||
|
|
||||||
|
expect(result).toBe(previousSnapshot);
|
||||||
|
expect((svc as any).launchStateStore.write).not.toHaveBeenCalled();
|
||||||
|
expect(invalidateRuntime).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it('keeps a bounded launch-state heartbeat for unchanged active snapshots', async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
vi.setSystemTime(new Date('2026-05-02T10:00:20.000Z'));
|
||||||
|
const svc = new TeamProvisioningService();
|
||||||
|
const teamName = 'launch-state-heartbeat-team';
|
||||||
|
const status = createMemberSpawnStatusEntry({
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
firstSpawnAcceptedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
});
|
||||||
|
const previousSnapshot = snapshotFromRuntimeMemberStatuses({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
launchPhase: 'active',
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
statuses: { alice: status as any },
|
||||||
|
});
|
||||||
|
const run = createMemberSpawnRun({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
memberSpawnStatuses: new Map([['alice', status]]),
|
||||||
|
});
|
||||||
|
run.isLaunch = true;
|
||||||
|
(svc as any).launchStateStore = {
|
||||||
|
read: vi.fn(async () => previousSnapshot),
|
||||||
|
write: vi.fn(async () => {}),
|
||||||
|
clear: vi.fn(async () => {}),
|
||||||
|
};
|
||||||
|
(svc as any).membersMetaStore = { getMembers: vi.fn(async () => []) };
|
||||||
|
const invalidateRuntime = vi.spyOn(svc as any, 'invalidateRuntimeSnapshotCaches');
|
||||||
|
|
||||||
|
const result = await (svc as any).persistLaunchStateSnapshotNow(run, 'active');
|
||||||
|
|
||||||
|
expect(result.updatedAt).toBe('2026-05-02T10:00:20.000Z');
|
||||||
|
expect((svc as any).launchStateStore.write).toHaveBeenCalledTimes(1);
|
||||||
|
expect(invalidateRuntime).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('does not skip the first service-owned launch-state write for an existing snapshot', async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
vi.setSystemTime(new Date('2026-05-02T10:00:05.000Z'));
|
||||||
|
const svc = new TeamProvisioningService();
|
||||||
|
const teamName = 'launch-state-first-write-team';
|
||||||
|
const status = createMemberSpawnStatusEntry({
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
firstSpawnAcceptedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
});
|
||||||
|
const previousSnapshot = snapshotFromRuntimeMemberStatuses({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
launchPhase: 'active',
|
||||||
|
updatedAt: '2026-05-02T10:00:02.000Z',
|
||||||
|
statuses: { alice: status as any },
|
||||||
|
});
|
||||||
|
const run = createMemberSpawnRun({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
memberSpawnStatuses: new Map([['alice', status]]),
|
||||||
|
});
|
||||||
|
run.isLaunch = true;
|
||||||
|
(svc as any).launchStateStore = {
|
||||||
|
read: vi.fn(async () => previousSnapshot),
|
||||||
|
write: vi.fn(async () => {}),
|
||||||
|
clear: vi.fn(async () => {}),
|
||||||
|
};
|
||||||
|
(svc as any).membersMetaStore = { getMembers: vi.fn(async () => []) };
|
||||||
|
const invalidateRuntime = vi.spyOn(svc as any, 'invalidateRuntimeSnapshotCaches');
|
||||||
|
|
||||||
|
await (svc as any).persistLaunchStateSnapshotNow(run, 'active');
|
||||||
|
|
||||||
|
expect((svc as any).launchStateStore.write).toHaveBeenCalledTimes(1);
|
||||||
|
expect(invalidateRuntime).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('does not skip the first write for a new run even when the previous snapshot is semantically equal', async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
vi.setSystemTime(new Date('2026-05-02T10:00:05.000Z'));
|
||||||
|
const svc = new TeamProvisioningService();
|
||||||
|
const teamName = 'launch-state-new-run-team';
|
||||||
|
const status = createMemberSpawnStatusEntry({
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
firstSpawnAcceptedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
});
|
||||||
|
const previousSnapshot = snapshotFromRuntimeMemberStatuses({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
launchPhase: 'active',
|
||||||
|
updatedAt: '2026-05-02T10:00:02.000Z',
|
||||||
|
statuses: { alice: status as any },
|
||||||
|
});
|
||||||
|
const run = createMemberSpawnRun({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
memberSpawnStatuses: new Map([['alice', status]]),
|
||||||
|
});
|
||||||
|
run.isLaunch = true;
|
||||||
|
(svc as any).launchStateStore = {
|
||||||
|
read: vi.fn(async () => previousSnapshot),
|
||||||
|
write: vi.fn(async () => {}),
|
||||||
|
clear: vi.fn(async () => {}),
|
||||||
|
};
|
||||||
|
(svc as any).launchStateWrittenRunIdByTeam.set(teamName, 'previous-run-id');
|
||||||
|
(svc as any).membersMetaStore = { getMembers: vi.fn(async () => []) };
|
||||||
|
const invalidateRuntime = vi.spyOn(svc as any, 'invalidateRuntimeSnapshotCaches');
|
||||||
|
|
||||||
|
const result = await (svc as any).persistLaunchStateSnapshotNow(run, 'active');
|
||||||
|
|
||||||
|
expect(result.updatedAt).toBe('2026-05-02T10:00:05.000Z');
|
||||||
|
expect((svc as any).launchStateStore.write).toHaveBeenCalledTimes(1);
|
||||||
|
expect(invalidateRuntime).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('writes and invalidates runtime cache when launch-state semantics change', async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
vi.setSystemTime(new Date('2026-05-02T10:00:05.000Z'));
|
||||||
|
const svc = new TeamProvisioningService();
|
||||||
|
const teamName = 'launch-state-change-team';
|
||||||
|
const previousStatus = createMemberSpawnStatusEntry({
|
||||||
|
status: 'waiting',
|
||||||
|
launchState: 'runtime_pending_bootstrap',
|
||||||
|
updatedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
firstSpawnAcceptedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
});
|
||||||
|
const nextStatus = createMemberSpawnStatusEntry({
|
||||||
|
status: 'online',
|
||||||
|
launchState: 'confirmed_alive',
|
||||||
|
runtimeAlive: true,
|
||||||
|
bootstrapConfirmed: true,
|
||||||
|
hardFailure: false,
|
||||||
|
livenessSource: 'heartbeat',
|
||||||
|
updatedAt: '2026-05-02T10:00:05.000Z',
|
||||||
|
firstSpawnAcceptedAt: '2026-05-02T10:00:00.000Z',
|
||||||
|
lastHeartbeatAt: '2026-05-02T10:00:05.000Z',
|
||||||
|
});
|
||||||
|
const previousSnapshot = snapshotFromRuntimeMemberStatuses({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
launchPhase: 'active',
|
||||||
|
updatedAt: '2026-05-02T10:00:02.000Z',
|
||||||
|
statuses: { alice: previousStatus as any },
|
||||||
|
});
|
||||||
|
const run = createMemberSpawnRun({
|
||||||
|
teamName,
|
||||||
|
expectedMembers: ['alice'],
|
||||||
|
memberSpawnStatuses: new Map([['alice', nextStatus]]),
|
||||||
|
});
|
||||||
|
run.isLaunch = true;
|
||||||
|
(svc as any).launchStateStore = {
|
||||||
|
read: vi.fn(async () => previousSnapshot),
|
||||||
|
write: vi.fn(async () => {}),
|
||||||
|
clear: vi.fn(async () => {}),
|
||||||
|
};
|
||||||
|
(svc as any).membersMetaStore = { getMembers: vi.fn(async () => []) };
|
||||||
|
const invalidateRuntime = vi.spyOn(svc as any, 'invalidateRuntimeSnapshotCaches');
|
||||||
|
|
||||||
|
const result = await (svc as any).persistLaunchStateSnapshotNow(run, 'active');
|
||||||
|
|
||||||
|
expect(result.members.alice?.launchState).toBe('confirmed_alive');
|
||||||
|
expect((svc as any).launchStateStore.write).toHaveBeenCalledTimes(1);
|
||||||
|
expect(invalidateRuntime).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
describe('getTeamAgentRuntimeSnapshot', () => {
|
describe('getTeamAgentRuntimeSnapshot', () => {
|
||||||
it('dedupes concurrent runtime snapshot probes for the same team', async () => {
|
it('dedupes concurrent runtime snapshot probes for the same team', async () => {
|
||||||
const svc = new TeamProvisioningService();
|
const svc = new TeamProvisioningService();
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue