From 8c863b1f0ee867785247dee3eb4da0c36d6c5181 Mon Sep 17 00:00:00 2001 From: Nicolas Meienberger Date: Wed, 8 Apr 2026 20:38:49 +0200 Subject: [PATCH] fix: pr feedbacks --- apps/agent/src/commands/backup-run.test.ts | 117 +++++++++++++ apps/agent/src/commands/backup-run.ts | 188 +++++++++++---------- 2 files changed, 212 insertions(+), 93 deletions(-) create mode 100644 apps/agent/src/commands/backup-run.test.ts diff --git a/apps/agent/src/commands/backup-run.test.ts b/apps/agent/src/commands/backup-run.test.ts new file mode 100644 index 00000000..70fc44ce --- /dev/null +++ b/apps/agent/src/commands/backup-run.test.ts @@ -0,0 +1,117 @@ +import { afterEach, expect, mock, spyOn, test } from "bun:test"; +import { Effect } from "effect"; +import waitForExpect from "wait-for-expect"; +import { fromPartial } from "@total-typescript/shoehorn"; +import { parseAgentMessage, type BackupCancelPayload, type BackupRunPayload } from "@zerobyte/contracts/agent-protocol"; +import * as resticServer from "@zerobyte/core/restic/server"; +import { handleBackupCancelCommand } from "./backup-cancel"; +import { handleBackupRunCommand } from "./backup-run"; +import type { ControllerCommandContext, RunningJob } from "../context"; + +const createDeferred = () => { + let resolve!: (value: T) => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + + return { promise, resolve }; +}; + +afterEach(() => { + mock.restore(); +}); + +test("waits for running-job registration before returning to the processor loop", async () => { + const outboundMessages: string[] = []; + const runningJobs = new Map(); + const setRunningJobGate = createDeferred(); + const backupGate = createDeferred<{ exitCode: number; result: null; warningDetails: null }>(); + let registeredAbortController: AbortController | undefined; + + spyOn(resticServer, "createRestic").mockReturnValue( + fromPartial({ + backup: () => + Effect.async<{ exitCode: number; result: null; warningDetails: null }, never>((resume) => { + void backupGate.promise.then((result) => { + resume(Effect.succeed(result)); + }); + }), + }), + ); + + const context: ControllerCommandContext = { + getRunningJob: (jobId) => Effect.succeed(runningJobs.get(jobId)), + setRunningJob: (jobId, job) => + Effect.async((resume) => { + void setRunningJobGate.promise.then(() => { + runningJobs.set(jobId, job); + registeredAbortController = job.abortController; + resume(Effect.void); + }); + }), + deleteRunningJob: (jobId) => + Effect.sync(() => { + runningJobs.delete(jobId); + }), + offerOutbound: (message) => + Effect.sync(() => { + outboundMessages.push(message); + return true; + }), + }; + + const runPayload = fromPartial({ + jobId: "job-1", + scheduleId: "schedule-1", + organizationId: "org-1", + sourcePath: "/tmp/source", + repositoryConfig: { + backend: "local", + path: "/tmp/repository", + }, + options: {}, + runtime: { + password: "password", + cacheDir: "/tmp/restic-cache", + passFile: "/tmp/restic-pass", + defaultExcludes: [], + }, + }); + const cancelPayload = fromPartial({ + jobId: "job-1", + scheduleId: "schedule-1", + }); + + const runPromise = Effect.runPromise(handleBackupRunCommand(context, runPayload)); + + try { + const returnedBeforeRegistration = await Promise.race([ + runPromise.then(() => true), + new Promise((resolve) => { + setTimeout(() => resolve(false), 0); + }), + ]); + + expect(returnedBeforeRegistration).toBe(false); + + setRunningJobGate.resolve(undefined); + await runPromise; + + await Effect.runPromise(handleBackupCancelCommand(context, cancelPayload)); + expect(registeredAbortController?.signal.aborted).toBe(true); + + backupGate.resolve({ exitCode: 0, result: null, warningDetails: null }); + + await waitForExpect(() => { + const cancelledMessage = outboundMessages + .map((message) => parseAgentMessage(message)) + .find((message) => message?.success && message.data.type === "backup.cancelled"); + + expect(cancelledMessage?.success).toBe(true); + expect(runningJobs.has("job-1")).toBe(false); + }); + } finally { + setRunningJobGate.resolve(undefined); + backupGate.resolve({ exitCode: 0, result: null, warningDetails: null }); + } +}); diff --git a/apps/agent/src/commands/backup-run.ts b/apps/agent/src/commands/backup-run.ts index 7af0cd62..a99b8f9f 100644 --- a/apps/agent/src/commands/backup-run.ts +++ b/apps/agent/src/commands/backup-run.ts @@ -7,107 +7,109 @@ import { toErrorDetails, toMessage } from "@zerobyte/core/utils"; import type { ControllerCommandContext } from "../context"; export const handleBackupRunCommand = (context: ControllerCommandContext, payload: BackupRunPayload) => { - return Effect.fork( - Effect.gen(function* () { - const existing = yield* context.getRunningJob(payload.jobId); - if (existing) { - yield* context.offerOutbound( - createAgentMessage("backup.failed", { - jobId: payload.jobId, - scheduleId: payload.scheduleId, - error: "Backup job is already running", - }), - ); - return; - } - - logger.info(`Starting backup ${payload.jobId} for schedule ${payload.scheduleId}`); - const abortController = new AbortController(); - yield* context.setRunningJob(payload.jobId, { scheduleId: payload.scheduleId, abortController }); - - const sendCancelled = () => { - return context.offerOutbound( - createAgentMessage("backup.cancelled", { - jobId: payload.jobId, - scheduleId: payload.scheduleId, - message: "Backup was cancelled", - }), - ); - }; - + return Effect.gen(function* () { + const existing = yield* context.getRunningJob(payload.jobId); + if (existing) { yield* context.offerOutbound( - createAgentMessage("backup.started", { + createAgentMessage("backup.failed", { jobId: payload.jobId, scheduleId: payload.scheduleId, + error: "Backup job is already running", }), ); + return; + } - const deps: ResticDeps = { - resolveSecret: async (encrypted) => encrypted, - getOrganizationResticPassword: async () => payload.runtime.password, - resticCacheDir: payload.runtime.cacheDir, - resticPassFile: payload.runtime.passFile, - defaultExcludes: payload.runtime.defaultExcludes, - hostname: payload.runtime.hostname, - }; + logger.info(`Starting backup ${payload.jobId} for schedule ${payload.scheduleId}`); + const abortController = new AbortController(); + yield* context.setRunningJob(payload.jobId, { scheduleId: payload.scheduleId, abortController }); - const restic = createRestic(deps); - const runtime = yield* Effect.runtime(); + yield* Effect.fork( + Effect.gen(function* () { + const sendCancelled = () => { + return context.offerOutbound( + createAgentMessage("backup.cancelled", { + jobId: payload.jobId, + scheduleId: payload.scheduleId, + message: "Backup was cancelled", + }), + ); + }; - yield* restic - .backup(payload.repositoryConfig, payload.sourcePath, { - organizationId: payload.organizationId, - ...payload.options, - signal: abortController.signal, - onProgress: (progress) => { - void Runtime.runPromise( - runtime, - context.offerOutbound( - createAgentMessage("backup.progress", { - jobId: payload.jobId, - scheduleId: payload.scheduleId, - progress, - }), - ), - ).catch((error) => { - logger.error(`Failed to send backup progress update: ${toMessage(error)}`); - }); - }, - }) - .pipe( - Effect.matchEffect({ - onSuccess: (result) => { - if (abortController.signal.aborted) { - return sendCancelled(); - } - - return context.offerOutbound( - createAgentMessage("backup.completed", { - jobId: payload.jobId, - scheduleId: payload.scheduleId, - exitCode: result.exitCode, - result: result.result, - warningDetails: result.warningDetails ?? undefined, - }), - ); - }, - onFailure: (error) => { - if (abortController.signal.aborted) { - return sendCancelled(); - } - - return context.offerOutbound( - createAgentMessage("backup.failed", { - jobId: payload.jobId, - scheduleId: payload.scheduleId, - error: toMessage(error), - errorDetails: toErrorDetails(error), - }), - ); - }, + yield* context.offerOutbound( + createAgentMessage("backup.started", { + jobId: payload.jobId, + scheduleId: payload.scheduleId, }), - Effect.ensuring(context.deleteRunningJob(payload.jobId)), ); - }), - ).pipe(Effect.asVoid); + + const deps: ResticDeps = { + resolveSecret: async (encrypted) => encrypted, + getOrganizationResticPassword: async () => payload.runtime.password, + resticCacheDir: payload.runtime.cacheDir, + resticPassFile: payload.runtime.passFile, + defaultExcludes: payload.runtime.defaultExcludes, + hostname: payload.runtime.hostname, + }; + + const restic = createRestic(deps); + const runtime = yield* Effect.runtime(); + + yield* restic + .backup(payload.repositoryConfig, payload.sourcePath, { + organizationId: payload.organizationId, + ...payload.options, + signal: abortController.signal, + onProgress: (progress) => { + void Runtime.runPromise( + runtime, + context.offerOutbound( + createAgentMessage("backup.progress", { + jobId: payload.jobId, + scheduleId: payload.scheduleId, + progress, + }), + ), + ).catch((error) => { + logger.error(`Failed to send backup progress update: ${toMessage(error)}`); + }); + }, + }) + .pipe( + Effect.matchEffect({ + onSuccess: (result) => { + if (abortController.signal.aborted) { + return sendCancelled(); + } + + return context.offerOutbound( + createAgentMessage("backup.completed", { + jobId: payload.jobId, + scheduleId: payload.scheduleId, + exitCode: result.exitCode, + result: result.result, + warningDetails: result.warningDetails ?? undefined, + }), + ); + }, + onFailure: (error) => { + if (abortController.signal.aborted) { + return sendCancelled(); + } + + return context.offerOutbound( + createAgentMessage("backup.failed", { + jobId: payload.jobId, + scheduleId: payload.scheduleId, + error: toMessage(error), + errorDetails: toErrorDetails(error), + }), + ); + }, + }), + Effect.ensuring(context.deleteRunningJob(payload.jobId)), + ); + }), + ); + }).pipe(Effect.asVoid); };