Merge pull request #32 from timothyjlaurent/timothyjlaurent/Nautilus

fix(FN-31): resolve per-project HeartbeatMonitor in multi-project setups
This commit is contained in:
gsxdsm
2026-05-04 07:15:35 -07:00
committed by GitHub
2 changed files with 90 additions and 20 deletions

View File

@@ -998,6 +998,30 @@ export function createApiRoutes(store: TaskStore, options?: ServerOptions): Rout
}
}
/**
* Resolve the HeartbeatMonitor for the engine that owns the given scopedStore.
*
* In multi-project setups each ProjectEngine has its own HeartbeatMonitor.
* This function walks all engines in the engineManager and returns the one
* whose working directory matches the scopedStore's root.
* Returns undefined when no matching engine is found.
*/
function resolveHeartbeatMonitor(scopedStore: TaskStore): ServerOptions["heartbeatMonitor"] {
const engineManager = options?.engineManager;
if (!engineManager) return undefined;
try {
const storeRoot = resolve(scopedStore.getRootDir());
for (const engine of engineManager.getAllEngines().values()) {
if (resolve(engine.getWorkingDirectory()) === storeRoot) {
return engine.getHeartbeatMonitor() ?? undefined;
}
}
} catch {
// path resolution failure — fall through
}
return undefined;
}
/**
* Trigger a heartbeat wake for an assigned agent based on a comment event.
*
@@ -1026,9 +1050,14 @@ export function createApiRoutes(store: TaskStore, options?: ServerOptions): Rout
return;
}
// Guard: heartbeatMonitor is bound to a specific project root directory.
// Skip the wake when the scoped store belongs to a different project.
if (!isHeartbeatMonitorForProject(scopedStore)) {
// Resolve the correct HeartbeatMonitor for this project.
const resolvedMonitor =
isHeartbeatMonitorForProject(scopedStore)
? heartbeatMonitor
: resolveHeartbeatMonitor(scopedStore);
// Skip: no heartbeat executor available for this project
if (!resolvedMonitor) {
return;
}
@@ -1063,7 +1092,7 @@ export function createApiRoutes(store: TaskStore, options?: ServerOptions): Rout
triggeringCommentType: wake.triggeringCommentType,
};
await heartbeatMonitor.executeHeartbeat({
await resolvedMonitor.executeHeartbeat({
agentId: assignedAgent.id,
source: "on_demand",
triggerDetail: wake.triggerDetail,
@@ -2987,6 +3016,7 @@ export function createApiRoutes(store: TaskStore, options?: ServerOptions): Rout
hasHeartbeatExecutor,
heartbeatMonitor,
isHeartbeatMonitorForProject,
resolveHeartbeatMonitor,
runExcerptToAgentLogs,
parseRunAuditFilters,
normalizeRunAuditEvent,

View File

@@ -16,6 +16,10 @@ interface AgentRuntimeRouteDeps {
hasHeartbeatExecutor: boolean;
heartbeatMonitor: import("../server.js").ServerOptions["heartbeatMonitor"];
isHeartbeatMonitorForProject: (scopedStore: import("@fusion/core").TaskStore) => boolean;
/** Resolve the HeartbeatMonitor for the engine backing a scoped store.
* Used for multi-project setups where each engine has its own monitor.
* Returns undefined when no matching engine is found. */
resolveHeartbeatMonitor: (scopedStore: import("@fusion/core").TaskStore) => import("../server.js").ServerOptions["heartbeatMonitor"];
runExcerptToAgentLogs: (run: import("@fusion/core").AgentHeartbeatRun) => import("@fusion/core").AgentLogEntry[];
parseRunAuditFilters: (query: Record<string, unknown>) => {
taskId?: string;
@@ -42,6 +46,7 @@ export function registerAgentRuntimeRoutes(ctx: ApiRoutesContext, deps: AgentRun
hasHeartbeatExecutor,
heartbeatMonitor,
isHeartbeatMonitorForProject,
resolveHeartbeatMonitor,
runExcerptToAgentLogs,
parseRunAuditFilters,
normalizeRunAuditEvent,
@@ -823,16 +828,22 @@ export function registerAgentRuntimeRoutes(ctx: ApiRoutesContext, deps: AgentRun
// Optionally trigger execution
let run: import("@fusion/core").AgentHeartbeatRun | undefined;
if (triggerExecution && hasHeartbeatExecutor && heartbeatMonitor && isHeartbeatMonitorForProject(scopedStore)) {
run = await heartbeatMonitor.executeHeartbeat({
agentId: req.params.id,
source: "on_demand",
triggerDetail: "Triggered from heartbeat",
contextSnapshot: {
wakeReason: "on_demand",
if (triggerExecution && hasHeartbeatExecutor && heartbeatMonitor) {
const resolvedMonitor =
isHeartbeatMonitorForProject(scopedStore)
? heartbeatMonitor
: resolveHeartbeatMonitor(scopedStore);
if (resolvedMonitor) {
run = await resolvedMonitor.executeHeartbeat({
agentId: req.params.id,
source: "on_demand",
triggerDetail: "Triggered from heartbeat",
},
});
contextSnapshot: {
wakeReason: "on_demand",
triggerDetail: "Triggered from heartbeat",
},
});
}
}
res.json(run ? { event, run } : event);
@@ -968,10 +979,15 @@ export function registerAgentRuntimeRoutes(ctx: ApiRoutesContext, deps: AgentRun
// Check for existing active run
const { store: scopedStore } = await getProjectContext(req);
// Guard: heartbeatMonitor is bound to a specific project root directory.
// Reject when the scoped store belongs to a different project.
if (!isHeartbeatMonitorForProject(scopedStore)) {
throw new ApiError(400, "Agent execution is only available for the server's primary project. The heartbeat monitor is not bound to this project.");
// Resolve the correct HeartbeatMonitor for this project.
// In multi-project setups, each engine has its own monitor.
const resolvedMonitor =
isHeartbeatMonitorForProject(scopedStore)
? heartbeatMonitor
: resolveHeartbeatMonitor(scopedStore);
if (!resolvedMonitor) {
throw new ApiError(400, "No heartbeat executor available for this project.");
}
const { AgentStore: AgentStoreClass } = await import("@fusion/core");
@@ -989,7 +1005,7 @@ export function registerAgentRuntimeRoutes(ctx: ApiRoutesContext, deps: AgentRun
}
// Execute heartbeat end-to-end (single run record, no duplicate startRun call)
const run = await heartbeatMonitor.executeHeartbeat({
const run = await resolvedMonitor.executeHeartbeat({
agentId: req.params.id,
source: invocationSource,
triggerDetail: trigger,
@@ -1062,8 +1078,32 @@ export function registerAgentRuntimeRoutes(ctx: ApiRoutesContext, deps: AgentRun
return;
}
if (hasHeartbeatExecutor && heartbeatMonitor && isHeartbeatMonitorForProject(scopedStore)) {
await heartbeatMonitor.stopRun(req.params.id);
if (hasHeartbeatExecutor && heartbeatMonitor) {
const resolvedMonitor =
isHeartbeatMonitorForProject(scopedStore)
? heartbeatMonitor
: resolveHeartbeatMonitor(scopedStore);
if (resolvedMonitor) {
await resolvedMonitor.stopRun(req.params.id);
} else {
const existingRun = await agentStore.getRunDetail(req.params.id, activeRun.id);
if (existingRun) {
await agentStore.saveRun({
...existingRun,
endedAt: new Date().toISOString(),
status: "terminated",
stderrExcerpt: existingRun.stderrExcerpt ?? "Run stopped by user",
});
}
await agentStore.endHeartbeatRun(activeRun.id, "terminated");
try {
await agentStore.updateAgentState(req.params.id, "active");
} catch {
// Best effort to restore an idle/active state for follow-up runs.
}
}
} else {
const existingRun = await agentStore.getRunDetail(req.params.id, activeRun.id);
if (existingRun) {