From 9bb5d2c9169d512690d4723d295176de6b3477f1 Mon Sep 17 00:00:00 2001 From: Nicolas Meienberger Date: Fri, 1 May 2026 18:20:16 +0200 Subject: [PATCH] fix: mark agent offline only if the session was removed properly --- app/server/modules/agents/agents-manager.ts | 3 +- .../modules/agents/controller/server.ts | 18 ++++---- .../modules/agents/controller/session.ts | 46 ++++++++----------- apps/agent/src/commands/backup-run.ts | 2 +- packages/core/src/node/logger.ts | 7 +++ 5 files changed, 36 insertions(+), 40 deletions(-) diff --git a/app/server/modules/agents/agents-manager.ts b/app/server/modules/agents/agents-manager.ts index 9342b072..519ab9f9 100644 --- a/app/server/modules/agents/agents-manager.ts +++ b/app/server/modules/agents/agents-manager.ts @@ -1,7 +1,7 @@ import { logger } from "@zerobyte/core/node"; import type { BackupRunPayload } from "@zerobyte/contracts/agent-protocol"; import { config } from "../../core/config"; -import type { AgentBackupEventHandlers } from "./controller/server"; +import { createAgentManagerRuntime, type AgentBackupEventHandlers } from "./controller/server"; import { spawnLocalAgentProcess, stopLocalAgentProcess } from "./local/process"; import type { BackupExecutionProgress, BackupExecutionResult } from "./helpers/runtime-state"; import { createAgentRuntimeState } from "./helpers/runtime-state"; @@ -155,7 +155,6 @@ export const startAgentController = async () => { runtime.agentManager = null; } - const { createAgentManagerRuntime } = await import("./controller/server"); const nextAgentManager = createAgentManagerRuntime(); nextAgentManager.setBackupEventHandlers(backupEventHandlers); diff --git a/app/server/modules/agents/controller/server.ts b/app/server/modules/agents/controller/server.ts index a763b5ea..3aec9a7f 100644 --- a/app/server/modules/agents/controller/server.ts +++ b/app/server/modules/agents/controller/server.ts @@ -197,10 +197,11 @@ export function createAgentManagerRuntime() { }); }, close: (ws) => { - removeSession(ws.data.agentId, ws.data.id); - void agentsService.markAgentOffline(ws.data.agentId).catch((error) => { - logger.error(`Failed to mark agent ${ws.data.agentId} as offline: ${toMessage(error)}`); - }); + if (removeSession(ws.data.agentId, ws.data.id)) { + void agentsService.markAgentOffline(ws.data.agentId).catch((error) => { + logger.error(`Failed to mark agent ${ws.data.agentId} as offline: ${toMessage(error)}`); + }); + } logger.info(`Agent "${ws.data.agentName}" (${ws.data.agentId}) disconnected`); }, }, @@ -214,11 +215,9 @@ export function createAgentManagerRuntime() { catch: (error) => new StopAgentManagerServerError({ cause: error }), }), ), - Effect.catchAll((error) => - Effect.sync(() => { - logger.error(`Failed to stop Agent Manager server: ${toMessage(error.cause)}`); - }), - ), + Effect.catchAll((error) => { + return logger.effect.error(`Failed to stop Agent Manager server: ${toMessage(error.cause)}`); + }), ), ); @@ -233,6 +232,7 @@ export function createAgentManagerRuntime() { await Effect.runPromise(Scope.close(scope, Exit.succeed(undefined))); }; + // TODO: Move the effect boundary up const start = async () => { if (runtimeScope) { await stop(); diff --git a/app/server/modules/agents/controller/session.ts b/app/server/modules/agents/controller/session.ts index a057ce4d..b8f0bdc0 100644 --- a/app/server/modules/agents/controller/session.ts +++ b/app/server/modules/agents/controller/session.ts @@ -136,16 +136,13 @@ export const createControllerAgentSession = ( yield* Effect.addFinalizer(() => closeSession()); const handleSendFailure = (reason: string) => { - logger.error( - `Closing session for agent ${socket.data.agentId} on ${socket.data.id} after an outbound websocket send failed: ${reason}`, - ); - - socket.close(); - - void Effect.runPromise(closeSession()).catch((error) => { + return Effect.gen(function* () { logger.error( - `Failed to close session for agent ${socket.data.agentId} on ${socket.data.id}: ${toMessage(error)}`, + `Closing session for agent ${socket.data.agentId} on ${socket.data.id} after an outbound websocket send failed: ${reason}`, ); + + yield* Effect.sync(() => socket.close()); + yield* closeSession(); }); }; @@ -154,17 +151,16 @@ export const createControllerAgentSession = ( Effect.forever( Effect.gen(function* () { const message = yield* Queue.take(outboundQueue); - yield* Effect.sync(() => { - try { - const sendResult = socket.send(message); - if (sendResult === 0) { - handleSendFailure("connection issue"); - } - } catch (error) { - handleSendFailure(toMessage(error)); - } + + const sendResult = yield* Effect.try({ + try: () => socket.send(message), + catch: (error) => toMessage(error), }); - }), + + if (sendResult === 0) { + yield* handleSendFailure("connection issue"); + } + }).pipe(Effect.catchAll((reason) => handleSendFailure(reason))), ), ); @@ -200,9 +196,7 @@ export const createControllerAgentSession = ( at: readyAt, }); - yield* Effect.sync(() => { - logger.info(`Agent "${socket.data.agentName}" (${socket.data.agentId}) is ready`); - }); + yield* logger.effect.info(`Agent "${socket.data.agentName}" (${socket.data.agentId}) is ready`); break; } case "backup.started": { @@ -210,7 +204,7 @@ export const createControllerAgentSession = ( scheduleId: message.payload.scheduleId, state: "active", }); - logger.info( + yield* logger.effect.info( `Backup ${message.payload.jobId} started on agent ${socket.data.agentId} for schedule ${message.payload.scheduleId}`, ); yield* handlers.onBackupStarted(message.payload); @@ -258,16 +252,12 @@ export const createControllerAgentSession = ( const parsed = parseAgentMessage(data); if (parsed === null) { - yield* Effect.sync(() => { - logger.warn(`Invalid JSON from agent ${socket.data.agentId}`); - }); + yield* logger.effect.warn(`Invalid JSON from agent ${socket.data.agentId}`); return; } if (!parsed.success) { - yield* Effect.sync(() => { - logger.warn(`Invalid agent message from ${socket.data.agentId}: ${parsed.error.message}`); - }); + yield* logger.effect.warn(`Invalid agent message from ${socket.data.agentId}: ${parsed.error.message}`); return; } diff --git a/apps/agent/src/commands/backup-run.ts b/apps/agent/src/commands/backup-run.ts index 1112fe2b..500ef2a3 100644 --- a/apps/agent/src/commands/backup-run.ts +++ b/apps/agent/src/commands/backup-run.ts @@ -21,7 +21,7 @@ export const handleBackupRunCommand = (context: ControllerCommandContext, payloa return; } - logger.info(`Starting backup ${payload.jobId} for schedule ${payload.scheduleId}`); + yield* logger.effect.info(`Starting backup ${payload.jobId} for schedule ${payload.scheduleId}`); const abortController = new AbortController(); yield* context.setRunningJob(payload.jobId, { scheduleId: payload.scheduleId, abortController }); diff --git a/packages/core/src/node/logger.ts b/packages/core/src/node/logger.ts index f9d35719..f1f163d2 100644 --- a/packages/core/src/node/logger.ts +++ b/packages/core/src/node/logger.ts @@ -2,6 +2,7 @@ import { format } from "date-fns"; import { createConsola, type ConsolaReporter } from "consola"; import { formatWithOptions } from "node:util"; import { sanitizeSensitiveData } from "../utils/sanitize"; +import { Effect } from "effect"; type LogLevel = "debug" | "info" | "warn" | "error"; @@ -102,4 +103,10 @@ export const logger = { info: (...messages: unknown[]) => consola.info(formatMessages(messages).join(" ")), warn: (...messages: unknown[]) => consola.warn(formatMessages(messages).join(" ")), error: (...messages: unknown[]) => consola.error(formatMessages(messages).join(" ")), + effect: { + debug: (...messages: unknown[]) => Effect.sync(() => consola.debug(formatMessages(messages).join(" "))), + info: (...messages: unknown[]) => Effect.sync(() => consola.info(formatMessages(messages).join(" "))), + warn: (...messages: unknown[]) => Effect.sync(() => consola.warn(formatMessages(messages).join(" "))), + error: (...messages: unknown[]) => Effect.sync(() => consola.error(formatMessages(messages).join(" "))), + }, };