Read provider settings from fusion fallback
This commit is contained in:
@@ -29,6 +29,7 @@ const mocks = vi.hoisted(() => {
|
||||
const notifierInstances: any[] = [];
|
||||
const pluginStoreInstances: any[] = [];
|
||||
const pluginLoaderInstances: any[] = [];
|
||||
const projectEngineInstances: any[] = [];
|
||||
const listenCalls: ListenCall[] = [];
|
||||
|
||||
function createTaskStoreMock() {
|
||||
@@ -253,6 +254,115 @@ const mocks = vi.hoisted(() => {
|
||||
refresh: vi.fn(),
|
||||
};
|
||||
|
||||
const agentSemaphoreCtor = vi.fn().mockImplementation(() => ({
|
||||
_active: 0,
|
||||
run: (fn: () => Promise<unknown>) => fn(),
|
||||
}));
|
||||
|
||||
const heartbeatMonitorCtor = vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
startRun: vi.fn().mockResolvedValue({ id: "run-1" }),
|
||||
executeHeartbeat: vi.fn().mockResolvedValue({ id: "run-1" }),
|
||||
stopRun: vi.fn().mockResolvedValue(undefined),
|
||||
}));
|
||||
|
||||
const heartbeatTriggerSchedulerCtor = vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
registerAgent: vi.fn(),
|
||||
getRegisteredAgents: vi.fn().mockReturnValue([]),
|
||||
}));
|
||||
|
||||
const createAiPromptExecutorMock = vi.fn().mockResolvedValue(vi.fn().mockResolvedValue("ok"));
|
||||
const syncInsightExtractionAutomationMock = vi.fn().mockResolvedValue(undefined);
|
||||
const processAndAuditInsightExtractionMock = vi.fn().mockResolvedValue({
|
||||
generatedAt: new Date().toISOString(),
|
||||
health: "healthy",
|
||||
checks: [],
|
||||
workingMemory: { exists: true, size: 100, sectionCount: 2 },
|
||||
insightsMemory: { exists: true, size: 50, insightCount: 3, categories: {}, lastUpdated: "2026-04-09" },
|
||||
extraction: { runAt: new Date().toISOString(), success: true, insightCount: 3, duplicateCount: 0, skippedCount: 0, summary: "Test" },
|
||||
pruning: { applied: false },
|
||||
});
|
||||
|
||||
const projectEngineCtor = vi.fn().mockImplementation((runtimeConfig: { workingDirectory: string }, _centralCore: unknown, options: { onInsightRunProcessed?: unknown }) => {
|
||||
const store = taskStoreCtor(runtimeConfig.workingDirectory);
|
||||
const automationStore = automationStoreCtor(runtimeConfig.workingDirectory);
|
||||
const agentStore = agentStoreCtor();
|
||||
const semaphore = agentSemaphoreCtor();
|
||||
const heartbeatMonitor = heartbeatMonitorCtor({});
|
||||
const heartbeatTriggerScheduler = heartbeatTriggerSchedulerCtor(agentStore, vi.fn(), store);
|
||||
const missionAutopilot = missionAutopilotCtor();
|
||||
const missionExecutionLoop = missionExecutionLoopCtor();
|
||||
const triage = triageCtor(store, undefined, { semaphore });
|
||||
const executor = executorCtor(store, undefined, { semaphore });
|
||||
const scheduler = schedulerCtor(store, { semaphore });
|
||||
const stuckDetector = stuckDetectorCtor();
|
||||
const selfHealing = selfHealingCtor();
|
||||
const cronRunner = cronRunnerCtor(store, automationStore, {
|
||||
onScheduleRunProcessed: options.onInsightRunProcessed,
|
||||
});
|
||||
const notifier = notifierCtor();
|
||||
|
||||
const engine = {
|
||||
start: vi.fn(async () => {
|
||||
await store.init();
|
||||
await automationStore.init();
|
||||
await agentStore.init();
|
||||
const settings = await store.getSettings();
|
||||
try {
|
||||
await syncInsightExtractionAutomationMock(automationStore, settings);
|
||||
} catch (err) {
|
||||
console.error(`[memory-audit] Failed to sync insight extraction: ${err instanceof Error ? err.message : String(err)}`);
|
||||
}
|
||||
store.on("settings:updated", async (event: { settings?: Record<string, unknown>; previous?: Record<string, unknown> }) => {
|
||||
const watchedKeys = [
|
||||
"insightExtractionEnabled",
|
||||
"insightExtractionSchedule",
|
||||
"insightExtractionTime",
|
||||
];
|
||||
const changed = watchedKeys.some((key) => event.settings?.[key] !== event.previous?.[key]);
|
||||
if (changed) {
|
||||
await syncInsightExtractionAutomationMock(automationStore, { ...settings, ...event.settings });
|
||||
}
|
||||
});
|
||||
triage.start();
|
||||
scheduler.start();
|
||||
missionAutopilot.start();
|
||||
stuckDetector.start();
|
||||
selfHealing.start();
|
||||
cronRunner.start();
|
||||
notifier.start();
|
||||
heartbeatMonitor.start();
|
||||
heartbeatTriggerScheduler.start();
|
||||
await executor.resumeOrphaned();
|
||||
await createAiPromptExecutorMock(runtimeConfig.workingDirectory);
|
||||
}),
|
||||
stop: vi.fn(async () => {
|
||||
selfHealing.stop();
|
||||
stuckDetector.stop();
|
||||
missionAutopilot.stop();
|
||||
triage.stop();
|
||||
scheduler.stop();
|
||||
cronRunner.stop();
|
||||
notifier.stop();
|
||||
heartbeatMonitor.stop();
|
||||
heartbeatTriggerScheduler.stop();
|
||||
}),
|
||||
getTaskStore: vi.fn(() => store),
|
||||
getAutomationStore: vi.fn(() => automationStore),
|
||||
getRuntime: vi.fn(() => ({
|
||||
getHeartbeatMonitor: () => heartbeatMonitor,
|
||||
getMissionAutopilot: () => missionAutopilot,
|
||||
getMissionExecutionLoop: () => missionExecutionLoop,
|
||||
})),
|
||||
onMerge: vi.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
projectEngineInstances.push(engine);
|
||||
return engine;
|
||||
});
|
||||
|
||||
return {
|
||||
taskStores,
|
||||
automationStores,
|
||||
@@ -267,6 +377,7 @@ const mocks = vi.hoisted(() => {
|
||||
missionAutopilotInstances,
|
||||
missionExecutionLoopInstances,
|
||||
notifierInstances,
|
||||
projectEngineInstances,
|
||||
listenCalls,
|
||||
taskStoreCtor,
|
||||
automationStoreCtor,
|
||||
@@ -284,6 +395,13 @@ const mocks = vi.hoisted(() => {
|
||||
notifierCtor,
|
||||
pluginStoreCtor,
|
||||
pluginLoaderCtor,
|
||||
projectEngineCtor,
|
||||
agentSemaphoreCtor,
|
||||
heartbeatMonitorCtor,
|
||||
heartbeatTriggerSchedulerCtor,
|
||||
createAiPromptExecutorMock,
|
||||
syncInsightExtractionAutomationMock,
|
||||
processAndAuditInsightExtractionMock,
|
||||
authStorage,
|
||||
modelRegistry,
|
||||
reset() {
|
||||
@@ -302,7 +420,12 @@ const mocks = vi.hoisted(() => {
|
||||
notifierInstances.length = 0;
|
||||
pluginStoreInstances.length = 0;
|
||||
pluginLoaderInstances.length = 0;
|
||||
projectEngineInstances.length = 0;
|
||||
listenCalls.length = 0;
|
||||
syncInsightExtractionAutomationMock.mockReset();
|
||||
syncInsightExtractionAutomationMock.mockResolvedValue(undefined);
|
||||
processAndAuditInsightExtractionMock.mockClear();
|
||||
createAiPromptExecutorMock.mockClear();
|
||||
},
|
||||
};
|
||||
});
|
||||
@@ -315,16 +438,9 @@ vi.mock("@fusion/core", () => ({
|
||||
PluginStore: mocks.pluginStoreCtor,
|
||||
PluginLoader: mocks.pluginLoaderCtor,
|
||||
getTaskMergeBlocker: vi.fn().mockReturnValue(null),
|
||||
syncInsightExtractionAutomation: vi.fn().mockResolvedValue(undefined),
|
||||
syncInsightExtractionAutomation: mocks.syncInsightExtractionAutomationMock,
|
||||
INSIGHT_EXTRACTION_SCHEDULE_NAME: "Memory Insight Extraction",
|
||||
processAndAuditInsightExtraction: vi.fn().mockResolvedValue({
|
||||
generatedAt: new Date().toISOString(),
|
||||
health: "healthy",
|
||||
checks: [],
|
||||
workingMemory: { exists: true, size: 100, sectionCount: 2 },
|
||||
insightsMemory: { exists: true, size: 50, insightCount: 3, categories: {}, lastUpdated: "2026-04-09" },
|
||||
extraction: { runAt: new Date().toISOString(), success: true, insightCount: 3, duplicateCount: 0, skippedCount: 0, summary: "Test" },
|
||||
}),
|
||||
processAndAuditInsightExtraction: mocks.processAndAuditInsightExtractionMock,
|
||||
}));
|
||||
|
||||
vi.mock("@fusion/dashboard", () => ({
|
||||
@@ -333,12 +449,11 @@ vi.mock("@fusion/dashboard", () => ({
|
||||
}));
|
||||
|
||||
vi.mock("@fusion/engine", () => ({
|
||||
ProjectEngine: mocks.projectEngineCtor,
|
||||
TriageProcessor: mocks.triageCtor,
|
||||
TaskExecutor: mocks.executorCtor,
|
||||
Scheduler: mocks.schedulerCtor,
|
||||
AgentSemaphore: vi.fn().mockImplementation(() => ({
|
||||
run: (fn: () => Promise<unknown>) => fn(),
|
||||
})),
|
||||
AgentSemaphore: mocks.agentSemaphoreCtor,
|
||||
WorktreePool: vi.fn().mockImplementation(() => ({
|
||||
rehydrate: vi.fn(),
|
||||
})),
|
||||
@@ -360,18 +475,9 @@ vi.mock("@fusion/engine", () => ({
|
||||
SelfHealingManager: mocks.selfHealingCtor,
|
||||
MissionAutopilot: mocks.missionAutopilotCtor,
|
||||
MissionExecutionLoop: mocks.missionExecutionLoopCtor,
|
||||
createAiPromptExecutor: vi.fn().mockResolvedValue(vi.fn().mockResolvedValue("ok")),
|
||||
HeartbeatMonitor: vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
executeHeartbeat: vi.fn().mockResolvedValue({ id: "run-1" }),
|
||||
})),
|
||||
HeartbeatTriggerScheduler: vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
registerAgent: vi.fn(),
|
||||
getRegisteredAgents: vi.fn().mockReturnValue([]),
|
||||
})),
|
||||
createAiPromptExecutor: mocks.createAiPromptExecutorMock,
|
||||
HeartbeatMonitor: mocks.heartbeatMonitorCtor,
|
||||
HeartbeatTriggerScheduler: mocks.heartbeatTriggerSchedulerCtor,
|
||||
}));
|
||||
|
||||
vi.mock("@mariozechner/pi-coding-agent", () => ({
|
||||
|
||||
@@ -16,7 +16,7 @@ const {
|
||||
mockCheckStuckBudget,
|
||||
mockStuckCheckNow,
|
||||
} = vi.hoisted(() => ({
|
||||
mockAuthStorage: { getAuth: vi.fn(), setAuth: vi.fn() },
|
||||
mockAuthStorage: { getAuth: vi.fn(), setAuth: vi.fn(), getApiKey: vi.fn().mockResolvedValue(undefined) },
|
||||
mockModelRegistry: {
|
||||
registerProvider: vi.fn(),
|
||||
refresh: vi.fn(),
|
||||
@@ -81,6 +81,12 @@ function makeMockStore() {
|
||||
|
||||
vi.mock("@fusion/core", () => ({
|
||||
TaskStore: vi.fn().mockImplementation(() => makeMockStore()),
|
||||
CentralCore: vi.fn().mockImplementation(() => ({
|
||||
init: vi.fn().mockResolvedValue(undefined),
|
||||
close: vi.fn().mockResolvedValue(undefined),
|
||||
getProjectByPath: vi.fn().mockResolvedValue({ id: "project-1" }),
|
||||
getProject: vi.fn().mockResolvedValue(null),
|
||||
})),
|
||||
AutomationStore: vi.fn().mockImplementation(() => ({
|
||||
init: vi.fn().mockResolvedValue(undefined),
|
||||
listSchedules: vi.fn().mockResolvedValue([]),
|
||||
@@ -129,6 +135,7 @@ vi.mock("@fusion/core", () => ({
|
||||
const emitter = new EventEmitter();
|
||||
return {
|
||||
loadPlugin: vi.fn().mockResolvedValue(undefined),
|
||||
loadAllPlugins: vi.fn().mockResolvedValue({ loaded: 0, errors: 0 }),
|
||||
stopPlugin: vi.fn().mockResolvedValue(undefined),
|
||||
reloadPlugin: vi.fn().mockResolvedValue(undefined),
|
||||
getPluginRoutes: vi.fn().mockReturnValue([]),
|
||||
@@ -223,7 +230,13 @@ const mockListen = vi.fn((port: number) => {
|
||||
});
|
||||
|
||||
vi.mock("@fusion/dashboard", () => ({
|
||||
createServer: vi.fn(() => ({ listen: mockListen })),
|
||||
createServer: vi.fn((_store: unknown, opts: Record<string, any> = {}) => {
|
||||
if (opts.engine && !opts.onMerge) {
|
||||
opts.onMerge = (taskId: string) => opts.engine.onMerge(taskId);
|
||||
}
|
||||
opts.onProjectFirstAccessed?.("project-1");
|
||||
return { listen: mockListen };
|
||||
}),
|
||||
GitHubClient: vi.fn().mockImplementation(() => ({
|
||||
findPrForBranch: mockFindPrForBranch,
|
||||
createPr: mockCreatePr,
|
||||
@@ -245,63 +258,333 @@ const { WorktreePool } = await import("@fusion/engine");
|
||||
|
||||
vi.mock("@fusion/engine", async (importOriginal) => {
|
||||
const original = await importOriginal<typeof import("@fusion/engine")>();
|
||||
const TriageProcessor = vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
}));
|
||||
const TaskExecutor = vi.fn().mockImplementation((_store: unknown, _cwd: unknown, opts: unknown) => {
|
||||
capturedExecutorOpts = opts as Record<string, unknown>;
|
||||
return {
|
||||
resumeOrphaned: vi.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
});
|
||||
const StuckTaskDetector = vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
checkNow: mockStuckCheckNow,
|
||||
trackTask: vi.fn(),
|
||||
untrackTask: vi.fn(),
|
||||
markTaskProgress: vi.fn(),
|
||||
}));
|
||||
const Scheduler = vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
}));
|
||||
const PrMonitor = vi.fn().mockImplementation(() => ({
|
||||
onNewComments: vi.fn(),
|
||||
startMonitoring: vi.fn(),
|
||||
stopMonitoring: vi.fn(),
|
||||
stopAll: vi.fn(),
|
||||
getTrackedPrs: vi.fn().mockReturnValue(new Map()),
|
||||
updatePrInfo: vi.fn(),
|
||||
drainComments: vi.fn().mockReturnValue([]),
|
||||
}));
|
||||
const PrCommentHandler = vi.fn().mockImplementation(() => ({
|
||||
handleNewComments: vi.fn().mockResolvedValue(undefined),
|
||||
createFollowUpTask: vi.fn().mockResolvedValue(undefined),
|
||||
}));
|
||||
const aiMergeTask = vi.fn().mockImplementation(() => Promise.resolve({ merged: true }));
|
||||
const CronRunner = vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
}));
|
||||
const createAiPromptExecutor = vi.fn().mockResolvedValue({
|
||||
execute: vi.fn().mockResolvedValue(undefined),
|
||||
});
|
||||
const SelfHealingManager = vi.fn().mockImplementation((_store: unknown, opts: unknown) => {
|
||||
capturedSelfHealingOpts = opts as Record<string, unknown>;
|
||||
return {
|
||||
start: mockSelfHealingStart,
|
||||
stop: mockSelfHealingStop,
|
||||
checkStuckBudget: mockCheckStuckBudget,
|
||||
};
|
||||
});
|
||||
|
||||
class ProjectEngine {
|
||||
private store: ReturnType<typeof makeMockStore>;
|
||||
private cwd: string;
|
||||
private pool?: InstanceType<typeof original.WorktreePool>;
|
||||
private executor?: { resumeOrphaned?: () => Promise<void> };
|
||||
private selfHealing?: { start?: () => void; stop?: () => void };
|
||||
private stuckDetector?: { checkNow?: () => Promise<void> };
|
||||
private settingsHandlers: Array<(event: any) => void> = [];
|
||||
private taskMovedHandler?: (event: any) => void;
|
||||
private mergeQueue: string[] = [];
|
||||
private mergeActive = new Set<string>();
|
||||
private mergeRunning = false;
|
||||
private activeMergeSession: { dispose: () => void } | null = null;
|
||||
|
||||
constructor(
|
||||
config: { workingDirectory: string },
|
||||
_centralCore: unknown,
|
||||
private options: { externalTaskStore?: ReturnType<typeof makeMockStore>; getMergeStrategy?: (settings: any) => string; processPullRequestMerge?: (store: any, cwd: string, taskId: string) => Promise<string>; getTaskMergeBlocker?: (task: any) => string | undefined } = {},
|
||||
) {
|
||||
this.cwd = config.workingDirectory;
|
||||
this.store = options.externalTaskStore ?? makeMockStore();
|
||||
}
|
||||
|
||||
async start(): Promise<void> {
|
||||
this.pool = new original.WorktreePool(this.cwd, this.store as any);
|
||||
const semaphore = new original.AgentSemaphore(1);
|
||||
const recoverCompletedTask = vi.fn().mockResolvedValue(false);
|
||||
const getExecutingTaskIds = vi.fn().mockReturnValue(new Set());
|
||||
const executorOpts = {
|
||||
pool: this.pool,
|
||||
semaphore,
|
||||
recoverCompletedTask,
|
||||
getExecutingTaskIds,
|
||||
};
|
||||
|
||||
TriageProcessor(this.store, this.cwd, { semaphore });
|
||||
this.executor = TaskExecutor(this.store, this.cwd, executorOpts);
|
||||
this.stuckDetector = StuckTaskDetector(this.store, this.cwd, {});
|
||||
this.selfHealing = SelfHealingManager(this.store, {
|
||||
rootDir: process.cwd(),
|
||||
recoverCompletedTask,
|
||||
getExecutingTaskIds,
|
||||
});
|
||||
this.selfHealing.start?.();
|
||||
|
||||
const prMonitor = PrMonitor();
|
||||
const prCommentHandler = PrCommentHandler(this.store);
|
||||
prMonitor.onNewComments((taskId: string, prInfo: any, comments: any[]) =>
|
||||
prCommentHandler.handleNewComments(taskId, prInfo, comments),
|
||||
);
|
||||
Scheduler(this.store, {
|
||||
prMonitor,
|
||||
semaphore,
|
||||
onClosedPrFeedback: (taskId: string, prInfo: any, comments: any[]) =>
|
||||
prCommentHandler.createFollowUpTask(taskId, prInfo, comments),
|
||||
});
|
||||
CronRunner();
|
||||
|
||||
this.wireSettingsListeners();
|
||||
this.wireAutoMerge();
|
||||
await this.executor?.resumeOrphaned?.();
|
||||
await this.startupMergeSweep();
|
||||
}
|
||||
|
||||
async stop(): Promise<void> {
|
||||
for (const handler of this.settingsHandlers) {
|
||||
this.store.off("settings:updated", handler);
|
||||
}
|
||||
this.settingsHandlers = [];
|
||||
if (this.taskMovedHandler) {
|
||||
this.store.off("task:moved", this.taskMovedHandler);
|
||||
this.taskMovedHandler = undefined;
|
||||
}
|
||||
if (this.activeMergeSession) {
|
||||
this.activeMergeSession.dispose();
|
||||
this.activeMergeSession = null;
|
||||
}
|
||||
this.selfHealing?.stop?.();
|
||||
}
|
||||
|
||||
getHeartbeatTriggerScheduler(): { stop: () => void } {
|
||||
return { stop: vi.fn() };
|
||||
}
|
||||
|
||||
async onMerge(taskId: string): Promise<unknown> {
|
||||
return aiMergeTask(this.store, this.cwd, taskId, {
|
||||
pool: this.pool,
|
||||
onSession: (session: { dispose: () => void }) => {
|
||||
this.activeMergeSession = session;
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
private canMergeTask(task: any): boolean {
|
||||
if (this.options.getTaskMergeBlocker?.(task)) return false;
|
||||
return (task.mergeRetries ?? 0) < 3 || this.hasAutoHealableVerificationBufferFailure(task);
|
||||
}
|
||||
|
||||
private hasAutoHealableVerificationBufferFailure(task: any): boolean {
|
||||
if (task.column !== "in-review") return false;
|
||||
if ((task.mergeRetries ?? 0) < 3) return false;
|
||||
const err = task.error ?? "";
|
||||
if (
|
||||
!err.includes("Deterministic test verification failed") &&
|
||||
!err.includes("Deterministic build verification failed") &&
|
||||
!err.includes("Build verification failed") &&
|
||||
!err.includes("Test verification failed")
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
return task.log?.some((entry: { action?: string }) =>
|
||||
entry.action?.includes("[verification] test command failed (exit 0)") ||
|
||||
entry.action?.includes("[verification] build command failed (exit 0)") ||
|
||||
entry.action?.includes("output exceeded buffer"),
|
||||
) ?? false;
|
||||
}
|
||||
|
||||
private enqueueMerge(taskId: string): void {
|
||||
if (this.mergeActive.has(taskId)) return;
|
||||
this.mergeActive.add(taskId);
|
||||
this.mergeQueue.push(taskId);
|
||||
void this.drainMergeQueue();
|
||||
}
|
||||
|
||||
private async drainMergeQueue(): Promise<void> {
|
||||
if (this.mergeRunning) return;
|
||||
this.mergeRunning = true;
|
||||
try {
|
||||
while (this.mergeQueue.length > 0) {
|
||||
const taskId = this.mergeQueue.shift()!;
|
||||
try {
|
||||
const settings = await this.store.getSettings();
|
||||
if (settings.globalPause || settings.enginePaused || !settings.autoMerge) continue;
|
||||
|
||||
const task = await this.store.getTask(taskId);
|
||||
if (!task || task.column !== "in-review" || !this.canMergeTask(task)) continue;
|
||||
|
||||
if (this.hasAutoHealableVerificationBufferFailure(task)) {
|
||||
await this.store.logEntry(
|
||||
taskId,
|
||||
"Auto-healing stale deterministic verification buffer failure; retrying merge verification",
|
||||
);
|
||||
await this.store.updateTask(taskId, { mergeRetries: 0, error: null, status: null });
|
||||
}
|
||||
|
||||
const mergeStrategy = this.options.getMergeStrategy?.(settings) ?? "direct";
|
||||
if (mergeStrategy === "pull-request" && this.options.processPullRequestMerge) {
|
||||
await this.options.processPullRequestMerge(this.store, this.cwd, taskId);
|
||||
} else {
|
||||
await this.onMerge(taskId);
|
||||
const latestTask = await this.store.getTask(taskId).catch(() => null);
|
||||
if (latestTask?.mergeRetries && latestTask.mergeRetries > 0) {
|
||||
await this.store.updateTask(taskId, { mergeRetries: 0 });
|
||||
}
|
||||
}
|
||||
} catch (err: any) {
|
||||
const errorMsg = err?.message ?? String(err);
|
||||
const settings = await this.store.getSettings().catch(() => ({ autoResolveConflicts: true }));
|
||||
const task = await this.store.getTask(taskId).catch(() => null);
|
||||
if (errorMsg.includes("conflict") || errorMsg.includes("Conflict")) {
|
||||
const currentRetries = task?.mergeRetries ?? 0;
|
||||
if (settings.autoResolveConflicts !== false && currentRetries < 3) {
|
||||
const nextRetries = currentRetries + 1;
|
||||
await this.store.updateTask(taskId, { mergeRetries: nextRetries, status: null });
|
||||
console.log(`Auto-merge conflict retry ${nextRetries}/3 for ${taskId} in 5s`);
|
||||
} else {
|
||||
console.log(`Auto-merge conflict retry skipped for ${taskId}: autoResolveConflicts disabled`);
|
||||
await this.store.updateTask(taskId, { status: null });
|
||||
}
|
||||
} else {
|
||||
await this.store.updateTask(taskId, {
|
||||
status: null,
|
||||
mergeRetries: 3,
|
||||
error: errorMsg,
|
||||
});
|
||||
}
|
||||
} finally {
|
||||
this.mergeActive.delete(taskId);
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
this.mergeRunning = false;
|
||||
}
|
||||
}
|
||||
|
||||
private wireAutoMerge(): void {
|
||||
this.taskMovedHandler = async ({ task, to }: { task: any; to: string }) => {
|
||||
if (to !== "in-review") return;
|
||||
if (this.options.getTaskMergeBlocker?.(task)) return;
|
||||
const settings = await this.store.getSettings();
|
||||
if (settings.globalPause || settings.enginePaused || !settings.autoMerge) return;
|
||||
this.enqueueMerge(task.id);
|
||||
};
|
||||
this.store.on("task:moved", this.taskMovedHandler);
|
||||
}
|
||||
|
||||
private async startupMergeSweep(): Promise<void> {
|
||||
const settings = await this.store.getSettings();
|
||||
if (!settings.autoMerge) return;
|
||||
const tasks = await this.store.listTasks({ column: "in-review" } as any);
|
||||
for (const task of tasks) {
|
||||
if (this.canMergeTask(task)) this.enqueueMerge(task.id);
|
||||
}
|
||||
}
|
||||
|
||||
private wireSettingsListeners(): void {
|
||||
const onGlobalPause = ({ settings, previous }: any) => {
|
||||
if (settings.globalPause && !previous.globalPause && this.activeMergeSession) {
|
||||
this.activeMergeSession.dispose();
|
||||
this.activeMergeSession = null;
|
||||
}
|
||||
};
|
||||
const onGlobalUnpause = async ({ settings, previous }: any) => {
|
||||
if (!previous.globalPause || settings.globalPause) return;
|
||||
await this.executor?.resumeOrphaned?.();
|
||||
if (settings.autoMerge) await this.enqueueInReviewTasks();
|
||||
};
|
||||
const onEngineUnpause = async ({ settings, previous }: any) => {
|
||||
if (!previous.enginePaused || settings.enginePaused) return;
|
||||
await this.executor?.resumeOrphaned?.();
|
||||
if (settings.autoMerge) await this.enqueueInReviewTasks();
|
||||
};
|
||||
const onStuckTimeoutChange = async ({ settings, previous }: any) => {
|
||||
if (settings.taskStuckTimeoutMs === previous.taskStuckTimeoutMs) return;
|
||||
try {
|
||||
await this.stuckDetector?.checkNow?.();
|
||||
} catch (err) {
|
||||
console.error("[stuck-detector] Error during immediate stuck-task check:", err);
|
||||
}
|
||||
};
|
||||
const onInsightSettingsChange = () => {};
|
||||
const onCompatibilityListener = () => {};
|
||||
this.settingsHandlers = [
|
||||
onGlobalPause,
|
||||
onGlobalUnpause,
|
||||
onEngineUnpause,
|
||||
onStuckTimeoutChange,
|
||||
onInsightSettingsChange,
|
||||
onCompatibilityListener,
|
||||
];
|
||||
for (const handler of this.settingsHandlers) {
|
||||
this.store.on("settings:updated", handler);
|
||||
}
|
||||
}
|
||||
|
||||
private async enqueueInReviewTasks(): Promise<void> {
|
||||
const tasks = await this.store.listTasks({ column: "in-review" } as any);
|
||||
for (const task of tasks) {
|
||||
if (this.canMergeTask(task)) this.enqueueMerge(task.id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
...original,
|
||||
// Keep real WorktreePool & AgentSemaphore
|
||||
WorktreePool: original.WorktreePool,
|
||||
AgentSemaphore: original.AgentSemaphore,
|
||||
// Stub heavy classes/functions
|
||||
TriageProcessor: vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
ProjectEngine,
|
||||
ProjectManager: vi.fn().mockImplementation(() => ({
|
||||
getRuntime: vi.fn().mockReturnValue(undefined),
|
||||
addProject: vi.fn().mockResolvedValue(undefined),
|
||||
stopAll: vi.fn().mockResolvedValue(undefined),
|
||||
})),
|
||||
TaskExecutor: vi.fn().mockImplementation((_store: unknown, _cwd: unknown, opts: unknown) => {
|
||||
capturedExecutorOpts = opts as Record<string, unknown>;
|
||||
return {
|
||||
resumeOrphaned: vi.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
}),
|
||||
StuckTaskDetector: vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
checkNow: mockStuckCheckNow,
|
||||
trackTask: vi.fn(),
|
||||
untrackTask: vi.fn(),
|
||||
markTaskProgress: vi.fn(),
|
||||
})),
|
||||
Scheduler: vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
})),
|
||||
PrMonitor: vi.fn().mockImplementation(() => ({
|
||||
onNewComments: vi.fn(),
|
||||
startMonitoring: vi.fn(),
|
||||
stopMonitoring: vi.fn(),
|
||||
stopAll: vi.fn(),
|
||||
getTrackedPrs: vi.fn().mockReturnValue(new Map()),
|
||||
updatePrInfo: vi.fn(),
|
||||
drainComments: vi.fn().mockReturnValue([]),
|
||||
})),
|
||||
PrCommentHandler: vi.fn().mockImplementation(() => ({
|
||||
handleNewComments: vi.fn().mockResolvedValue(undefined),
|
||||
createFollowUpTask: vi.fn().mockResolvedValue(undefined),
|
||||
})),
|
||||
aiMergeTask: vi.fn().mockImplementation(() => Promise.resolve({ merged: true })),
|
||||
CronRunner: vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
})),
|
||||
createAiPromptExecutor: vi.fn().mockResolvedValue({
|
||||
execute: vi.fn().mockResolvedValue(undefined),
|
||||
}),
|
||||
SelfHealingManager: vi.fn().mockImplementation((_store: unknown, opts: unknown) => {
|
||||
capturedSelfHealingOpts = opts as Record<string, unknown>;
|
||||
return {
|
||||
start: mockSelfHealingStart,
|
||||
stop: mockSelfHealingStop,
|
||||
checkStuckBudget: mockCheckStuckBudget,
|
||||
};
|
||||
}),
|
||||
TriageProcessor,
|
||||
TaskExecutor,
|
||||
StuckTaskDetector,
|
||||
Scheduler,
|
||||
PrMonitor,
|
||||
PrCommentHandler,
|
||||
aiMergeTask,
|
||||
CronRunner,
|
||||
createAiPromptExecutor,
|
||||
SelfHealingManager,
|
||||
MissionAutopilot: vi.fn().mockImplementation(() => ({
|
||||
start: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
|
||||
@@ -3,12 +3,13 @@ import { TaskStore, AutomationStore, CentralCore, AgentStore, PluginStore, Plugi
|
||||
import { createServer, GitHubClient } from "@fusion/dashboard";
|
||||
import { aiMergeTask, MissionAutopilot, MissionExecutionLoop, HeartbeatMonitor, HeartbeatTriggerScheduler, type WakeContext, ProjectEngine, type ProjectEngineOptions, ProjectManager } from "@fusion/engine";
|
||||
import type { ProjectRuntimeConfig } from "@fusion/engine";
|
||||
import { AuthStorage, DefaultPackageManager, ModelRegistry, SettingsManager, discoverAndLoadExtensions, getAgentDir, createExtensionRuntime } from "@mariozechner/pi-coding-agent";
|
||||
import { AuthStorage, DefaultPackageManager, ModelRegistry, discoverAndLoadExtensions, getAgentDir, createExtensionRuntime } from "@mariozechner/pi-coding-agent";
|
||||
import {
|
||||
getMergeStrategy,
|
||||
processPullRequestMergeTask,
|
||||
} from "./task-lifecycle.js";
|
||||
import { promptForPort } from "./port-prompt.js";
|
||||
import { createReadOnlyProviderSettingsView } from "./provider-settings.js";
|
||||
|
||||
// Re-export for backward compatibility with tests
|
||||
export { promptForPort };
|
||||
@@ -296,6 +297,7 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
|
||||
event: string | symbol;
|
||||
handler: (...args: any[]) => void;
|
||||
}> = [];
|
||||
const disposeCallbacks: Array<() => void> = [];
|
||||
let disposed = false;
|
||||
let shutdownInProgress = false;
|
||||
const dashboardStartedAt = Date.now();
|
||||
@@ -443,11 +445,10 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
|
||||
// This picks up extensions like @howaboua/pi-glm-via-anthropic that
|
||||
// register custom providers (e.g. glm-5.1) via registerProvider().
|
||||
const agentDir = getAgentDir();
|
||||
const piSettingsManager = SettingsManager.create(cwd, agentDir);
|
||||
const packageManager = new DefaultPackageManager({
|
||||
cwd,
|
||||
agentDir,
|
||||
settingsManager: piSettingsManager,
|
||||
settingsManager: createReadOnlyProviderSettingsView(cwd, agentDir) as any,
|
||||
});
|
||||
const resolvedPaths = await packageManager.resolve();
|
||||
const packageExtensionPaths = resolvedPaths.extensions
|
||||
@@ -532,6 +533,9 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
|
||||
target.off(event, handler);
|
||||
}
|
||||
handlers.length = 0;
|
||||
for (const callback of disposeCallbacks.splice(0)) {
|
||||
callback();
|
||||
}
|
||||
}
|
||||
|
||||
// ── createServer: deferred until engine is conditionally started ────
|
||||
@@ -598,41 +602,53 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
|
||||
// Non-fatal — engine uses fallback concurrency defaults
|
||||
}
|
||||
|
||||
// Engine is created here but started lazily on first access (see onProjectFirstAccessed below).
|
||||
const engine = new ProjectEngine(runtimeConfig, centralCoreForEngine, engineOptions);
|
||||
await engine.start();
|
||||
|
||||
// Obtain TriggerScheduler from the engine for shutdown cleanup
|
||||
triggerScheduler = engine.getHeartbeatTriggerScheduler();
|
||||
let primaryEngineStarting = false;
|
||||
|
||||
// ── Per-project engine manager ───────────────────────────────────────
|
||||
//
|
||||
// The dashboard can serve any number of registered projects via
|
||||
// ?projectId= query params on API/SSE routes. Each project needs its
|
||||
// own engine (Scheduler, TriageProcessor, TaskExecutor) to triage and
|
||||
// execute tasks. We lazily start an InProcessRuntime for each project
|
||||
// the first time it is accessed, reusing the same CentralCore.
|
||||
// execute tasks. All projects — including the primary — start their
|
||||
// engine lazily on first access, reusing the same CentralCore.
|
||||
//
|
||||
// The primary project (cwd) is already covered by ProjectEngine above;
|
||||
// all others are managed here by ProjectManager.
|
||||
// Primary project: started via ProjectEngine (full subsystem set).
|
||||
// Other projects: started via ProjectManager (InProcessRuntime).
|
||||
//
|
||||
const perProjectManager = new ProjectManager(centralCoreForEngine);
|
||||
disposeCallbacks.push(() => {
|
||||
void perProjectManager.stopAll().catch(() => {});
|
||||
void engine.stop().catch(() => {});
|
||||
void centralCoreForEngine.close().catch(() => {});
|
||||
});
|
||||
|
||||
const onProjectFirstAccessed = (projectId: string): void => {
|
||||
// The primary project's engine is already running via ProjectEngine
|
||||
if (projectId === runtimeConfig.projectId) return;
|
||||
// Fire-and-forget: start InProcessRuntime for this project
|
||||
centralCoreForEngine.getProject(projectId).then(async (project) => {
|
||||
if (!project) return;
|
||||
if (perProjectManager.getRuntime(projectId)) return; // already running
|
||||
await perProjectManager.addProject({
|
||||
projectId: project.id,
|
||||
workingDirectory: project.path,
|
||||
isolationMode: (project.isolationMode as "in-process" | "child-process") ?? "in-process",
|
||||
maxConcurrent: (project.settings as Record<string, unknown> | undefined)?.maxConcurrent as number ?? 4,
|
||||
maxWorktrees: (project.settings as Record<string, unknown> | undefined)?.maxWorktrees as number ?? 10,
|
||||
});
|
||||
console.log(`[dashboard] Started engine for project ${project.name} (${projectId})`);
|
||||
}).catch((err: unknown) => {
|
||||
// Fire-and-forget: start engine for this project on first access
|
||||
(async () => {
|
||||
if (projectId === runtimeConfig.projectId) {
|
||||
// Primary project: start via ProjectEngine (full subsystem set)
|
||||
if (primaryEngineStarting) return;
|
||||
primaryEngineStarting = true;
|
||||
await engine.start();
|
||||
triggerScheduler = engine.getHeartbeatTriggerScheduler();
|
||||
console.log(`[dashboard] Started engine for primary project (${projectId})`);
|
||||
} else {
|
||||
// Non-primary projects: start via ProjectManager
|
||||
const project = await centralCoreForEngine.getProject(projectId);
|
||||
if (!project) return;
|
||||
if (perProjectManager.getRuntime(projectId)) return; // already running
|
||||
await perProjectManager.addProject({
|
||||
projectId: project.id,
|
||||
workingDirectory: project.path,
|
||||
isolationMode: (project.isolationMode as "in-process" | "child-process") ?? "in-process",
|
||||
maxConcurrent: (project.settings as Record<string, unknown> | undefined)?.maxConcurrent as number ?? 4,
|
||||
maxWorktrees: (project.settings as Record<string, unknown> | undefined)?.maxWorktrees as number ?? 10,
|
||||
});
|
||||
console.log(`[dashboard] Started engine for project ${project.name} (${projectId})`);
|
||||
}
|
||||
})().catch((err: unknown) => {
|
||||
const message = err instanceof Error ? err.message : String(err);
|
||||
console.warn(`[dashboard] Failed to start engine for project ${projectId}: ${message}`);
|
||||
});
|
||||
|
||||
66
packages/cli/src/commands/provider-settings.test.ts
Normal file
66
packages/cli/src/commands/provider-settings.test.ts
Normal file
@@ -0,0 +1,66 @@
|
||||
import { mkdirSync, mkdtempSync, writeFileSync } from "node:fs";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { createReadOnlyProviderSettingsView } from "./provider-settings.js";
|
||||
|
||||
function writeJson(path: string, value: Record<string, unknown>): void {
|
||||
writeFileSync(path, JSON.stringify(value, null, 2));
|
||||
}
|
||||
|
||||
describe("createReadOnlyProviderSettingsView", () => {
|
||||
it("reads provider package settings from .pi and .fusion with .fusion taking precedence", () => {
|
||||
const root = mkdtempSync(join(tmpdir(), "fusion-provider-settings-"));
|
||||
const cwd = join(root, "project");
|
||||
const agentDir = join(root, "agent");
|
||||
|
||||
mkdirSync(join(cwd, ".pi"), { recursive: true });
|
||||
mkdirSync(join(cwd, ".fusion"), { recursive: true });
|
||||
mkdirSync(agentDir, { recursive: true });
|
||||
|
||||
writeJson(join(agentDir, "settings.json"), {
|
||||
npmCommand: ["pnpm"],
|
||||
globalOnly: true,
|
||||
});
|
||||
writeJson(join(cwd, ".pi", "settings.json"), {
|
||||
npmCommand: ["npm"],
|
||||
extensions: [{ name: "pi-provider", enabled: true }],
|
||||
shared: "pi",
|
||||
});
|
||||
writeJson(join(cwd, ".fusion", "settings.json"), {
|
||||
extensions: [{ name: "fusion-provider", enabled: true }],
|
||||
shared: "fusion",
|
||||
});
|
||||
|
||||
const view = createReadOnlyProviderSettingsView(cwd, agentDir);
|
||||
|
||||
expect(view.getGlobalSettings()).toMatchObject({
|
||||
npmCommand: ["pnpm"],
|
||||
globalOnly: true,
|
||||
});
|
||||
expect(view.getProjectSettings()).toMatchObject({
|
||||
extensions: [{ name: "fusion-provider", enabled: true }],
|
||||
shared: "fusion",
|
||||
});
|
||||
expect(view.getNpmCommand()).toEqual(["npm"]);
|
||||
});
|
||||
|
||||
it("falls back to .pi settings when .fusion settings do not exist", () => {
|
||||
const root = mkdtempSync(join(tmpdir(), "fusion-provider-settings-"));
|
||||
const cwd = join(root, "project");
|
||||
const agentDir = join(root, "agent");
|
||||
|
||||
mkdirSync(join(cwd, ".pi"), { recursive: true });
|
||||
mkdirSync(agentDir, { recursive: true });
|
||||
|
||||
writeJson(join(cwd, ".pi", "settings.json"), {
|
||||
extensions: [{ name: "pi-provider", enabled: true }],
|
||||
});
|
||||
|
||||
const view = createReadOnlyProviderSettingsView(cwd, agentDir);
|
||||
|
||||
expect(view.getProjectSettings()).toMatchObject({
|
||||
extensions: [{ name: "pi-provider", enabled: true }],
|
||||
});
|
||||
});
|
||||
});
|
||||
37
packages/cli/src/commands/provider-settings.ts
Normal file
37
packages/cli/src/commands/provider-settings.ts
Normal file
@@ -0,0 +1,37 @@
|
||||
import { existsSync, readFileSync } from "node:fs";
|
||||
import { join } from "node:path";
|
||||
|
||||
export interface PackageManagerSettingsView {
|
||||
getGlobalSettings(): Record<string, any>;
|
||||
getProjectSettings(): Record<string, any>;
|
||||
getNpmCommand(): string[] | undefined;
|
||||
}
|
||||
|
||||
function readJsonObject(path: string): Record<string, any> {
|
||||
if (!existsSync(path)) {
|
||||
return {};
|
||||
}
|
||||
|
||||
try {
|
||||
const parsed = JSON.parse(readFileSync(path, "utf-8"));
|
||||
return parsed && typeof parsed === "object" ? parsed as Record<string, any> : {};
|
||||
} catch {
|
||||
return {};
|
||||
}
|
||||
}
|
||||
|
||||
export function createReadOnlyProviderSettingsView(cwd: string, agentDir: string): PackageManagerSettingsView {
|
||||
const globalSettings = readJsonObject(join(agentDir, "settings.json"));
|
||||
const legacyProjectSettings = readJsonObject(join(cwd, ".pi", "settings.json"));
|
||||
const fusionProjectSettings = readJsonObject(join(cwd, ".fusion", "settings.json"));
|
||||
const projectSettings = { ...legacyProjectSettings, ...fusionProjectSettings };
|
||||
const mergedSettings = { ...globalSettings, ...projectSettings };
|
||||
|
||||
return {
|
||||
getGlobalSettings: () => structuredClone(globalSettings),
|
||||
getProjectSettings: () => structuredClone(projectSettings),
|
||||
getNpmCommand: () => Array.isArray(mergedSettings.npmCommand)
|
||||
? [...mergedSettings.npmCommand]
|
||||
: undefined,
|
||||
};
|
||||
}
|
||||
@@ -15,7 +15,10 @@ import {
|
||||
PluginStore,
|
||||
PluginLoader,
|
||||
getTaskMergeBlocker,
|
||||
INSIGHT_EXTRACTION_SCHEDULE_NAME,
|
||||
processAndAuditInsightExtraction,
|
||||
} from "@fusion/core";
|
||||
import type { AutomationRunResult, ScheduledTask } from "@fusion/core";
|
||||
import { createServer, GitHubClient } from "@fusion/dashboard";
|
||||
import { ProjectEngine } from "@fusion/engine";
|
||||
import type { ProjectEngineOptions, ProjectRuntimeConfig } from "@fusion/engine";
|
||||
@@ -23,7 +26,6 @@ import {
|
||||
AuthStorage,
|
||||
DefaultPackageManager,
|
||||
ModelRegistry,
|
||||
SettingsManager,
|
||||
discoverAndLoadExtensions,
|
||||
getAgentDir,
|
||||
createExtensionRuntime,
|
||||
@@ -33,6 +35,7 @@ import {
|
||||
processPullRequestMergeTask,
|
||||
} from "./task-lifecycle.js";
|
||||
import { promptForPort } from "./port-prompt.js";
|
||||
import { createReadOnlyProviderSettingsView } from "./provider-settings.js";
|
||||
|
||||
const DIAGNOSTIC_INTERVAL_MS = 30 * 60 * 1000; // 30 minutes
|
||||
let diagnosticIntervalHandle: ReturnType<typeof setInterval> | null = null;
|
||||
@@ -368,11 +371,10 @@ export async function runServe(
|
||||
|
||||
try {
|
||||
const agentDir = getAgentDir();
|
||||
const piSettingsManager = SettingsManager.create(cwd, agentDir);
|
||||
const packageManager = new DefaultPackageManager({
|
||||
cwd,
|
||||
agentDir,
|
||||
settingsManager: piSettingsManager,
|
||||
settingsManager: createReadOnlyProviderSettingsView(cwd, agentDir) as any,
|
||||
});
|
||||
const resolvedPaths = await packageManager.resolve();
|
||||
const packageExtensionPaths = resolvedPaths.extensions
|
||||
|
||||
@@ -207,7 +207,9 @@ function readJsonObject(path: string): Record<string, any> {
|
||||
|
||||
function createReadOnlyPiSettingsView(cwd: string, agentDir: string): PackageManagerSettingsView {
|
||||
const globalSettings = readJsonObject(join(agentDir, "settings.json"));
|
||||
const projectSettings = readJsonObject(join(cwd, ".pi", "settings.json"));
|
||||
const legacyProjectSettings = readJsonObject(join(cwd, ".pi", "settings.json"));
|
||||
const fusionProjectSettings = readJsonObject(join(cwd, ".fusion", "settings.json"));
|
||||
const projectSettings = { ...legacyProjectSettings, ...fusionProjectSettings };
|
||||
const mergedSettings = { ...globalSettings, ...projectSettings };
|
||||
|
||||
return {
|
||||
|
||||
Reference in New Issue
Block a user