fix: wait for agent to be ready before startup
This commit is contained in:
parent
a5d0cc986c
commit
188888ff9a
2 changed files with 30 additions and 1 deletions
|
|
@ -3,6 +3,7 @@ import type { BackupRunPayload, VolumeCommand, VolumeCommandResult } from "@zero
|
||||||
import { Effect } from "effect";
|
import { Effect } from "effect";
|
||||||
import { config } from "../../core/config";
|
import { config } from "../../core/config";
|
||||||
import { createAgentManagerRuntime, type AgentManagerEvent } from "./controller/server";
|
import { createAgentManagerRuntime, type AgentManagerEvent } from "./controller/server";
|
||||||
|
import { LOCAL_AGENT_ID } from "./constants";
|
||||||
import { spawnLocalAgentProcess, stopLocalAgentProcess } from "./local/process";
|
import { spawnLocalAgentProcess, stopLocalAgentProcess } from "./local/process";
|
||||||
import type { BackupExecutionProgress, BackupExecutionResult } from "./helpers/runtime-state";
|
import type { BackupExecutionProgress, BackupExecutionResult } from "./helpers/runtime-state";
|
||||||
import { createAgentRuntimeState } from "./helpers/runtime-state";
|
import { createAgentRuntimeState } from "./helpers/runtime-state";
|
||||||
|
|
@ -286,7 +287,16 @@ export const agentManager = {
|
||||||
};
|
};
|
||||||
|
|
||||||
export const startLocalAgent = async () => {
|
export const startLocalAgent = async () => {
|
||||||
await spawnLocalAgentProcess(getAgentRuntimeState());
|
const runtime = getAgentRuntimeState();
|
||||||
|
await spawnLocalAgentProcess(runtime);
|
||||||
|
|
||||||
|
if (!runtime.agentManager) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!(await runtime.agentManager.waitForAgentReady(LOCAL_AGENT_ID))) {
|
||||||
|
throw new Error("Local agent did not become ready before startup");
|
||||||
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
// fallow-ignore-next-line unused-export
|
// fallow-ignore-next-line unused-export
|
||||||
|
|
|
||||||
|
|
@ -46,6 +46,7 @@ class StopAgentManagerServerError extends Data.TaggedError("StopAgentManagerServ
|
||||||
export function createAgentManagerRuntime(onEvent: (event: AgentManagerEvent) => void) {
|
export function createAgentManagerRuntime(onEvent: (event: AgentManagerEvent) => void) {
|
||||||
let sessions = new Map<string, ControllerAgentSessionHandle>();
|
let sessions = new Map<string, ControllerAgentSessionHandle>();
|
||||||
let runtimeScope: Scope.CloseableScope | null = null;
|
let runtimeScope: Scope.CloseableScope | null = null;
|
||||||
|
const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));
|
||||||
|
|
||||||
const closeSession = (sessionHandle: ControllerAgentSessionHandle) =>
|
const closeSession = (sessionHandle: ControllerAgentSessionHandle) =>
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
|
|
@ -79,6 +80,11 @@ export function createAgentManagerRuntime(onEvent: (event: AgentManagerEvent) =>
|
||||||
const getSessionHandle = (agentId: string) => sessions.get(agentId);
|
const getSessionHandle = (agentId: string) => sessions.get(agentId);
|
||||||
const getSession = (agentId: string) => getSessionHandle(agentId)?.session;
|
const getSession = (agentId: string) => getSessionHandle(agentId)?.session;
|
||||||
|
|
||||||
|
const isAgentReady = (agentId: string) => {
|
||||||
|
const session = getSession(agentId);
|
||||||
|
return !!session && Effect.runSync(session.isReady());
|
||||||
|
};
|
||||||
|
|
||||||
const handleSessionEvent = (params: { agentId: string; agentName: string; sessionId: string }) => {
|
const handleSessionEvent = (params: { agentId: string; agentName: string; sessionId: string }) => {
|
||||||
const { agentId, agentName } = params;
|
const { agentId, agentName } = params;
|
||||||
|
|
||||||
|
|
@ -290,6 +296,19 @@ export function createAgentManagerRuntime(onEvent: (event: AgentManagerEvent) =>
|
||||||
|
|
||||||
return {
|
return {
|
||||||
start,
|
start,
|
||||||
|
waitForAgentReady: async (agentId: string, timeoutMs = 10_000) => {
|
||||||
|
const deadline = Date.now() + timeoutMs;
|
||||||
|
|
||||||
|
while (Date.now() < deadline) {
|
||||||
|
if (isAgentReady(agentId)) {
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
await sleep(50);
|
||||||
|
}
|
||||||
|
|
||||||
|
return isAgentReady(agentId);
|
||||||
|
},
|
||||||
sendBackup: (agentId: string, payload: BackupRunPayload) =>
|
sendBackup: (agentId: string, payload: BackupRunPayload) =>
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const session = getSession(agentId);
|
const session = getSession(agentId);
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue