feat(FN-4823): wire central claim store into mesh lease manager
Fusion-Task-Id: FN-4823 Fusion-Task-Lineage: 0034c04f-df82-4a1b-9f25-b93643a2c157
This commit is contained in:
committed by
gsxdsm
parent
6fa9aef630
commit
bbcc0a269c
@@ -11,7 +11,7 @@ import type {
|
||||
MessageStore,
|
||||
RoutineStore,
|
||||
} from "@fusion/core";
|
||||
import { ChatStore, isEphemeralAgent } from "@fusion/core";
|
||||
import { ChatStore, createCentralDatabase, isEphemeralAgent } from "@fusion/core";
|
||||
import { Scheduler } from "../scheduler.js";
|
||||
import type { PrMonitor, PrComment } from "../pr-monitor.js";
|
||||
import type { PrInfo } from "@fusion/core";
|
||||
@@ -96,6 +96,7 @@ export class InProcessRuntime
|
||||
private usageLimitPauser?: UsageLimitPauser;
|
||||
private selfHealingManager?: SelfHealingManager;
|
||||
private leaseManager?: MeshLeaseManager;
|
||||
private leaseCentralClaimStore?: ReturnType<typeof createCentralDatabase>;
|
||||
private agentStore?: AgentStore;
|
||||
private heartbeatMonitor?: HeartbeatMonitor;
|
||||
private triggerScheduler?: HeartbeatTriggerScheduler;
|
||||
@@ -308,11 +309,22 @@ export class InProcessRuntime
|
||||
})
|
||||
: undefined;
|
||||
|
||||
// FN-4823/FN-4819 §2.5: central-claim-aware recovery when central DB is reachable;
|
||||
// fallback to local-only recovery remains in MeshLeaseManager for single-node contexts.
|
||||
try {
|
||||
this.leaseCentralClaimStore = createCentralDatabase(this.centralCore.getGlobalDir());
|
||||
this.leaseCentralClaimStore.init();
|
||||
} catch (error) {
|
||||
runtimeLog.warn(`Failed to initialize central claim store for mesh lease recovery: ${error instanceof Error ? error.message : String(error)}`);
|
||||
this.leaseCentralClaimStore = undefined;
|
||||
}
|
||||
|
||||
this.leaseManager = new MeshLeaseManager({
|
||||
taskStore: this.taskStore,
|
||||
agentStore: this.agentStore,
|
||||
getHandoffPolicy: () => this.taskStore.getSettings().then((settings) => settings.owningNodeHandoffPolicy),
|
||||
getExecutingTaskIds: () => this.executor?.getExecutingTaskIds() ?? new Set<string>(),
|
||||
centralClaimStore: this.leaseCentralClaimStore,
|
||||
projectId: this.config.projectId,
|
||||
});
|
||||
|
||||
@@ -897,6 +909,11 @@ export class InProcessRuntime
|
||||
}
|
||||
}
|
||||
|
||||
if (this.leaseCentralClaimStore) {
|
||||
this.leaseCentralClaimStore.close();
|
||||
this.leaseCentralClaimStore = undefined;
|
||||
}
|
||||
|
||||
this.setStatus("stopped");
|
||||
runtimeLog.log(`InProcessRuntime stopped for project ${this.config.projectId}`);
|
||||
} catch (error) {
|
||||
|
||||
Reference in New Issue
Block a user