Merge pull request #32 from timothyjlaurent/timothyjlaurent/Nautilus
fix(FN-31): resolve per-project HeartbeatMonitor in multi-project setups
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user