refactor: split agent runtime state
This commit is contained in:
parent
faf5c7f26d
commit
38b4ac2c5c
3 changed files with 87 additions and 73 deletions
|
|
@ -1,16 +1,15 @@
|
||||||
import type { ChildProcess } from "node:child_process";
|
|
||||||
import { logger } from "@zerobyte/core/node";
|
import { logger } from "@zerobyte/core/node";
|
||||||
import type { ResticBackupOutputDto } from "@zerobyte/core/restic";
|
import type { BackupRunPayload } from "@zerobyte/contracts/agent-protocol";
|
||||||
import type { BackupProgressPayload, BackupRunPayload } from "@zerobyte/contracts/agent-protocol";
|
import { config } from "../../core/config";
|
||||||
import type { AgentBackupEventHandlers, AgentManagerRuntime } from "./controller/server";
|
import type { AgentBackupEventHandlers } from "./controller/server";
|
||||||
import { spawnLocalAgentProcess, stopLocalAgentProcess } from "./local/process";
|
import { spawnLocalAgentProcess, stopLocalAgentProcess } from "./local/process";
|
||||||
|
import type { BackupExecutionProgress, BackupExecutionResult } from "./helpers/runtime-state";
|
||||||
|
import { createAgentRuntimeState } from "./helpers/runtime-state";
|
||||||
|
import { getDevAgentRuntimeState } from "./helpers/runtime-state.dev";
|
||||||
|
export type { BackupExecutionProgress, BackupExecutionResult } from "./helpers/runtime-state";
|
||||||
|
export type { ProcessWithAgentRuntime } from "./helpers/runtime-state.dev";
|
||||||
|
|
||||||
export type BackupExecutionProgress = BackupProgressPayload["progress"];
|
const productionRuntimeState = createAgentRuntimeState();
|
||||||
export type BackupExecutionResult =
|
|
||||||
| { status: "unavailable"; error: Error }
|
|
||||||
| { status: "completed"; exitCode: number; result: ResticBackupOutputDto | null; warningDetails: string | null }
|
|
||||||
| { status: "failed"; error: string }
|
|
||||||
| { status: "cancelled"; message?: string };
|
|
||||||
|
|
||||||
export type AgentRunBackupRequest = {
|
export type AgentRunBackupRequest = {
|
||||||
scheduleId: number;
|
scheduleId: number;
|
||||||
|
|
@ -19,71 +18,11 @@ export type AgentRunBackupRequest = {
|
||||||
onProgress: (progress: BackupExecutionProgress) => void;
|
onProgress: (progress: BackupExecutionProgress) => void;
|
||||||
};
|
};
|
||||||
|
|
||||||
type ActiveBackupRun = {
|
const getAgentRuntimeState = () => (config.__prod__ ? productionRuntimeState : getDevAgentRuntimeState());
|
||||||
scheduleId: number;
|
|
||||||
jobId: string;
|
|
||||||
scheduleShortId: string;
|
|
||||||
onProgress: (progress: BackupExecutionProgress) => void;
|
|
||||||
resolve: (result: BackupExecutionResult) => void;
|
|
||||||
cancellationRequested: boolean;
|
|
||||||
};
|
|
||||||
|
|
||||||
type AgentRuntimeState = {
|
|
||||||
agentManager: AgentManagerRuntime | null;
|
|
||||||
localAgent: ChildProcess | null;
|
|
||||||
isStoppingLocalAgent: boolean;
|
|
||||||
localAgentRestartTimeout: ReturnType<typeof setTimeout> | null;
|
|
||||||
activeBackupsByScheduleId: Map<number, ActiveBackupRun>;
|
|
||||||
activeBackupScheduleIdsByJobId: Map<string, number>;
|
|
||||||
};
|
|
||||||
|
|
||||||
type LegacyAgentRuntimeState = Omit<AgentRuntimeState, "activeBackupsByScheduleId" | "activeBackupScheduleIdsByJobId"> &
|
|
||||||
Partial<Pick<AgentRuntimeState, "activeBackupsByScheduleId" | "activeBackupScheduleIdsByJobId">>;
|
|
||||||
|
|
||||||
export type ProcessWithAgentRuntime = NodeJS.Process & {
|
|
||||||
__zerobyteAgentRuntime?: LegacyAgentRuntimeState;
|
|
||||||
};
|
|
||||||
|
|
||||||
const getAgentRuntimeState = () => {
|
|
||||||
const runtimeProcess = process as ProcessWithAgentRuntime;
|
|
||||||
const existingRuntime = runtimeProcess.__zerobyteAgentRuntime;
|
|
||||||
|
|
||||||
if (existingRuntime) {
|
|
||||||
const runtime: AgentRuntimeState = {
|
|
||||||
...existingRuntime,
|
|
||||||
activeBackupsByScheduleId: existingRuntime.activeBackupsByScheduleId ?? new Map<number, ActiveBackupRun>(),
|
|
||||||
activeBackupScheduleIdsByJobId: existingRuntime.activeBackupScheduleIdsByJobId ?? new Map<string, number>(),
|
|
||||||
};
|
|
||||||
|
|
||||||
runtimeProcess.__zerobyteAgentRuntime = runtime;
|
|
||||||
return runtime;
|
|
||||||
}
|
|
||||||
|
|
||||||
const runtime: AgentRuntimeState = {
|
|
||||||
agentManager: null,
|
|
||||||
localAgent: null,
|
|
||||||
isStoppingLocalAgent: false,
|
|
||||||
localAgentRestartTimeout: null,
|
|
||||||
activeBackupsByScheduleId: new Map(),
|
|
||||||
activeBackupScheduleIdsByJobId: new Map(),
|
|
||||||
};
|
|
||||||
|
|
||||||
runtimeProcess.__zerobyteAgentRuntime = runtime;
|
|
||||||
return runtime;
|
|
||||||
};
|
|
||||||
|
|
||||||
const getAgentManagerRuntime = () => getAgentRuntimeState().agentManager;
|
const getAgentManagerRuntime = () => getAgentRuntimeState().agentManager;
|
||||||
const getActiveBackupsByScheduleId = () => getAgentRuntimeState().activeBackupsByScheduleId;
|
const getActiveBackupsByScheduleId = () => getAgentRuntimeState().activeBackupsByScheduleId;
|
||||||
const getActiveBackupScheduleIdsByJobId = () => getAgentRuntimeState().activeBackupScheduleIdsByJobId;
|
const getActiveBackupScheduleIdsByJobId = () => getAgentRuntimeState().activeBackupScheduleIdsByJobId;
|
||||||
|
|
||||||
const getUnavailableError = (agentId: string) => {
|
|
||||||
if (agentId === "local") {
|
|
||||||
return new Error("Local backup agent is not connected");
|
|
||||||
}
|
|
||||||
|
|
||||||
return new Error(`Backup agent ${agentId} is not connected`);
|
|
||||||
};
|
|
||||||
|
|
||||||
const clearActiveBackupRun = (scheduleId: number) => {
|
const clearActiveBackupRun = (scheduleId: number) => {
|
||||||
const activeBackupsByScheduleId = getActiveBackupsByScheduleId();
|
const activeBackupsByScheduleId = getActiveBackupsByScheduleId();
|
||||||
const activeBackupScheduleIdsByJobId = getActiveBackupScheduleIdsByJobId();
|
const activeBackupScheduleIdsByJobId = getActiveBackupScheduleIdsByJobId();
|
||||||
|
|
@ -229,7 +168,7 @@ export const agentManager = {
|
||||||
if (!runtime) {
|
if (!runtime) {
|
||||||
return {
|
return {
|
||||||
status: "unavailable",
|
status: "unavailable",
|
||||||
error: getUnavailableError(agentId),
|
error: new Error(`Backup agent ${agentId} is not connected`),
|
||||||
} satisfies BackupExecutionResult;
|
} satisfies BackupExecutionResult;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -254,7 +193,7 @@ export const agentManager = {
|
||||||
clearActiveBackupRun(request.scheduleId);
|
clearActiveBackupRun(request.scheduleId);
|
||||||
return {
|
return {
|
||||||
status: "unavailable",
|
status: "unavailable",
|
||||||
error: getUnavailableError(agentId),
|
error: new Error(`Failed to send backup command to agent ${agentId}`),
|
||||||
} satisfies BackupExecutionResult;
|
} satisfies BackupExecutionResult;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
37
app/server/modules/agents/helpers/runtime-state.dev.ts
Normal file
37
app/server/modules/agents/helpers/runtime-state.dev.ts
Normal file
|
|
@ -0,0 +1,37 @@
|
||||||
|
import { createAgentRuntimeState, type AgentRuntimeState } from "./runtime-state";
|
||||||
|
|
||||||
|
type LegacyAgentRuntimeState = Omit<AgentRuntimeState, "activeBackupsByScheduleId" | "activeBackupScheduleIdsByJobId"> &
|
||||||
|
Partial<Pick<AgentRuntimeState, "activeBackupsByScheduleId" | "activeBackupScheduleIdsByJobId">>;
|
||||||
|
|
||||||
|
export type ProcessWithAgentRuntime = NodeJS.Process & {
|
||||||
|
__zerobyteAgentRuntime?: LegacyAgentRuntimeState;
|
||||||
|
};
|
||||||
|
|
||||||
|
const hasActiveBackupMaps = (runtime: LegacyAgentRuntimeState): runtime is AgentRuntimeState => {
|
||||||
|
return runtime.activeBackupsByScheduleId instanceof Map && runtime.activeBackupScheduleIdsByJobId instanceof Map;
|
||||||
|
};
|
||||||
|
|
||||||
|
const hydrateAgentRuntimeState = (runtime: LegacyAgentRuntimeState): AgentRuntimeState => ({
|
||||||
|
...runtime,
|
||||||
|
activeBackupsByScheduleId: runtime.activeBackupsByScheduleId ?? new Map(),
|
||||||
|
activeBackupScheduleIdsByJobId: runtime.activeBackupScheduleIdsByJobId ?? new Map(),
|
||||||
|
});
|
||||||
|
|
||||||
|
export const getDevAgentRuntimeState = (): AgentRuntimeState => {
|
||||||
|
// Bun reloads modules in place during development, so keep the live runtime on process.
|
||||||
|
const runtimeProcess = process as ProcessWithAgentRuntime;
|
||||||
|
const existingRuntime = runtimeProcess.__zerobyteAgentRuntime;
|
||||||
|
if (!existingRuntime) {
|
||||||
|
const runtime = createAgentRuntimeState();
|
||||||
|
runtimeProcess.__zerobyteAgentRuntime = runtime;
|
||||||
|
return runtime;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (hasActiveBackupMaps(existingRuntime)) {
|
||||||
|
return existingRuntime;
|
||||||
|
}
|
||||||
|
|
||||||
|
const runtime = hydrateAgentRuntimeState(existingRuntime);
|
||||||
|
runtimeProcess.__zerobyteAgentRuntime = runtime;
|
||||||
|
return runtime;
|
||||||
|
};
|
||||||
38
app/server/modules/agents/helpers/runtime-state.ts
Normal file
38
app/server/modules/agents/helpers/runtime-state.ts
Normal file
|
|
@ -0,0 +1,38 @@
|
||||||
|
import type { ChildProcess } from "node:child_process";
|
||||||
|
import type { ResticBackupOutputDto } from "@zerobyte/core/restic";
|
||||||
|
import type { BackupProgressPayload } from "@zerobyte/contracts/agent-protocol";
|
||||||
|
import type { AgentManagerRuntime } from "../controller/server";
|
||||||
|
|
||||||
|
export type BackupExecutionProgress = BackupProgressPayload["progress"];
|
||||||
|
export type BackupExecutionResult =
|
||||||
|
| { status: "unavailable"; error: Error }
|
||||||
|
| { status: "completed"; exitCode: number; result: ResticBackupOutputDto | null; warningDetails: string | null }
|
||||||
|
| { status: "failed"; error: string }
|
||||||
|
| { status: "cancelled"; message?: string };
|
||||||
|
|
||||||
|
type ActiveBackupRun = {
|
||||||
|
scheduleId: number;
|
||||||
|
jobId: string;
|
||||||
|
scheduleShortId: string;
|
||||||
|
onProgress: (progress: BackupExecutionProgress) => void;
|
||||||
|
resolve: (result: BackupExecutionResult) => void;
|
||||||
|
cancellationRequested: boolean;
|
||||||
|
};
|
||||||
|
|
||||||
|
export type AgentRuntimeState = {
|
||||||
|
agentManager: AgentManagerRuntime | null;
|
||||||
|
localAgent: ChildProcess | null;
|
||||||
|
isStoppingLocalAgent: boolean;
|
||||||
|
localAgentRestartTimeout: ReturnType<typeof setTimeout> | null;
|
||||||
|
activeBackupsByScheduleId: Map<number, ActiveBackupRun>;
|
||||||
|
activeBackupScheduleIdsByJobId: Map<string, number>;
|
||||||
|
};
|
||||||
|
|
||||||
|
export const createAgentRuntimeState = (): AgentRuntimeState => ({
|
||||||
|
agentManager: null,
|
||||||
|
localAgent: null,
|
||||||
|
isStoppingLocalAgent: false,
|
||||||
|
localAgentRestartTimeout: null,
|
||||||
|
activeBackupsByScheduleId: new Map(),
|
||||||
|
activeBackupScheduleIdsByJobId: new Map(),
|
||||||
|
});
|
||||||
Loading…
Reference in a new issue