From 74ae2ce0c698ee506a0b4c835aecc1eeb40518a2 Mon Sep 17 00:00:00 2001 From: Anso Date: Tue, 12 May 2026 15:58:30 -0400 Subject: [PATCH] fix: harden atomic deployment rollback (#1029) * fix: harden atomic deployment rollback * fix: update Docker toolchain to Go 1.26.3 * fix: repair Dockerfile tr argument split across lines * fix: bump protobufjs to clear npm audit high-severity advisories * fix: sanitize error objects in console.error to prevent log injection --- backend/src/__tests__/compose-service.test.ts | 40 +- .../src/__tests__/filesystem-backup.test.ts | 37 +- .../src/__tests__/scheduler-service.test.ts | 9 +- .../stacks-failure-notifications.test.ts | 57 +- backend/src/routes/imageUpdates.ts | 5 +- backend/src/routes/stacks.ts | 27 +- backend/src/routes/templates.ts | 13 +- backend/src/services/ComposeService.ts | 1004 +++++++++-------- backend/src/services/FileSystemService.ts | 21 +- backend/src/services/SchedulerService.ts | 4 + .../EditorLayout/hooks/useStackActions.ts | 52 +- 11 files changed, 739 insertions(+), 530 deletions(-) diff --git a/backend/src/__tests__/compose-service.test.ts b/backend/src/__tests__/compose-service.test.ts index 35d144ee..bcd8f5bf 100644 --- a/backend/src/__tests__/compose-service.test.ts +++ b/backend/src/__tests__/compose-service.test.ts @@ -100,7 +100,7 @@ vi.mock('../services/LogFormatter', () => ({ LogFormatter: { formatLine: (line: string) => line }, })); -import { ComposeService } from '../services/ComposeService'; +import { ComposeService, getComposeRollbackInfo } from '../services/ComposeService'; /** Creates an EventEmitter that mimics a child_process spawn result */ function createMockProcess() { @@ -235,7 +235,19 @@ describe('ComposeService - deployStack', () => { expect(mockBackupStackFiles).toHaveBeenCalledWith('my-stack'); }); - it('throws CONTAINER_CRASHED when exited container has non-zero exit code', async () => { + it('aborts atomic deploy before docker side effects when backup fails', async () => { + mockBackupStackFiles.mockRejectedValueOnce(new Error('disk full')); + + const svc = ComposeService.getInstance(1); + + await expect(svc.deployStack('my-stack', undefined, true)).rejects.toThrow( + 'Atomic deployment backup failed', + ); + expect(mockSpawn).not.toHaveBeenCalled(); + expect(mockGetContainersByStack).not.toHaveBeenCalled(); + }); + + it('throws sanitized CONTAINER_CRASHED when exited container has non-zero exit code', async () => { setupAutoCloseSpawn(); mockListContainers.mockResolvedValue([{ Id: 'crashed-c1', @@ -243,7 +255,7 @@ describe('ComposeService - deployStack', () => { Labels: { 'com.docker.compose.project': 'my-stack' }, }]); mockContainerInspect.mockResolvedValue({ State: { ExitCode: 1 } }); - mockContainerLogs.mockResolvedValue(Buffer.from('Error: something failed')); + mockContainerLogs.mockResolvedValue(Buffer.from('SECRET_TOKEN=leaked')); const svc = ComposeService.getInstance(1); // Attach catch handler immediately so rejection is never "unhandled" @@ -253,6 +265,8 @@ describe('ComposeService - deployStack', () => { const error = await result; expect(error).not.toBeNull(); expect(error!.message).toContain('CONTAINER_CRASHED'); + expect(error!.message).not.toContain('SECRET_TOKEN'); + expect(mockContainerLogs).not.toHaveBeenCalled(); }); it('rolls back on failure when atomic=true', async () => { @@ -272,9 +286,29 @@ describe('ComposeService - deployStack', () => { const error = await result; expect(error).not.toBeNull(); expect(error!.message).toContain('CONTAINER_CRASHED'); + expect(getComposeRollbackInfo(error)).toEqual({ attempted: true, rolledBack: true }); expect(mockRestoreStackFiles).toHaveBeenCalledWith('my-stack'); }); + it('reports rollback failure when atomic restore fails', async () => { + setupAutoCloseSpawn(); + mockListContainers.mockResolvedValue([{ + Id: 'crashed-c1', + State: 'exited', + Labels: { 'com.docker.compose.project': 'my-stack' }, + }]); + mockContainerInspect.mockResolvedValue({ State: { ExitCode: 1 } }); + mockRestoreStackFiles.mockRejectedValueOnce(new Error('restore denied')); + + const svc = ComposeService.getInstance(1); + const result = svc.deployStack('my-stack', undefined, true).then(() => null, (e: Error) => e); + + await vi.runAllTimersAsync(); + const error = await result; + expect(error).not.toBeNull(); + expect(getComposeRollbackInfo(error)).toEqual({ attempted: true, rolledBack: false }); + }); + it('does not roll back when atomic=false', async () => { setupAutoCloseSpawn(); mockListContainers.mockResolvedValue([{ diff --git a/backend/src/__tests__/filesystem-backup.test.ts b/backend/src/__tests__/filesystem-backup.test.ts index 91065f94..64db2d07 100644 --- a/backend/src/__tests__/filesystem-backup.test.ts +++ b/backend/src/__tests__/filesystem-backup.test.ts @@ -1,6 +1,6 @@ /** * Verifies that FileSystemService stores stack backups under - * /backups// rather than inside the user's compose + * /backups/// rather than inside the user's compose * folder. The old in-stack-folder location failed with EACCES whenever a * container had chowned the bind mount, breaking the atomic rollback * feature for those stacks. @@ -12,12 +12,12 @@ import { promises as fsPromises } from 'fs'; // Mutable state the mocked NodeRegistry reads. Each test rewrites these // before instantiating FileSystemService. -const mockState = { composeDir: '' }; +const mockState = { composeDir: '', composeDirs: new Map() }; vi.mock('../services/NodeRegistry', () => ({ NodeRegistry: { getInstance: () => ({ - getComposeDir: () => mockState.composeDir, + getComposeDir: (nodeId?: number) => mockState.composeDirs.get(nodeId ?? 1) ?? mockState.composeDir, getDefaultNodeId: () => 1, }), }, @@ -34,6 +34,7 @@ describe('FileSystemService backup location', () => { composeDir = await fsPromises.mkdtemp(path.join(os.tmpdir(), 'sencho-compose-')); dataDir = await fsPromises.mkdtemp(path.join(os.tmpdir(), 'sencho-data-')); mockState.composeDir = composeDir; + mockState.composeDirs = new Map([[1, composeDir]]); originalDataDir = process.env.DATA_DIR; process.env.DATA_DIR = dataDir; }); @@ -45,7 +46,7 @@ describe('FileSystemService backup location', () => { await fsPromises.rm(dataDir, { recursive: true, force: true }); }); - it('writes backups under /backups//, not inside the stack folder', async () => { + it('writes backups under /backups///, not inside the stack folder', async () => { const stackName = 'web'; const stackDir = path.join(composeDir, stackName); await fsPromises.mkdir(stackDir, { recursive: true }); @@ -55,7 +56,7 @@ describe('FileSystemService backup location', () => { const service = FileSystemService.getInstance(); await service.backupStackFiles(stackName); - const newBackupDir = path.join(dataDir, 'backups', stackName); + const newBackupDir = path.join(dataDir, 'backups', '1', stackName); const oldBackupDir = path.join(stackDir, '.sencho-backup'); // New location has every backed-up file @@ -83,6 +84,32 @@ describe('FileSystemService backup location', () => { expect(typeof after.timestamp).toBe('number'); }); + it('scopes backups by node id when stack names overlap', async () => { + const stackName = 'web'; + const secondComposeDir = await fsPromises.mkdtemp(path.join(os.tmpdir(), 'sencho-compose-')); + mockState.composeDirs.set(2, secondComposeDir); + try { + const nodeOneStackDir = path.join(composeDir, stackName); + const nodeTwoStackDir = path.join(secondComposeDir, stackName); + await fsPromises.mkdir(nodeOneStackDir, { recursive: true }); + await fsPromises.mkdir(nodeTwoStackDir, { recursive: true }); + await fsPromises.writeFile(path.join(nodeOneStackDir, 'compose.yaml'), 'services:\n one: {}\n', 'utf-8'); + await fsPromises.writeFile(path.join(nodeTwoStackDir, 'compose.yaml'), 'services:\n two: {}\n', 'utf-8'); + + await FileSystemService.getInstance(1).backupStackFiles(stackName); + await FileSystemService.getInstance(2).backupStackFiles(stackName); + + await expect( + fsPromises.readFile(path.join(dataDir, 'backups', '1', stackName, 'compose.yaml'), 'utf-8'), + ).resolves.toContain('one'); + await expect( + fsPromises.readFile(path.join(dataDir, 'backups', '2', stackName, 'compose.yaml'), 'utf-8'), + ).resolves.toContain('two'); + } finally { + await fsPromises.rm(secondComposeDir, { recursive: true, force: true }); + } + }); + it('restoreStackFiles copies files from the new location back to the stack dir', async () => { const stackName = 'db'; const stackDir = path.join(composeDir, stackName); diff --git a/backend/src/__tests__/scheduler-service.test.ts b/backend/src/__tests__/scheduler-service.test.ts index cd8ffa8f..4e06f86b 100644 --- a/backend/src/__tests__/scheduler-service.test.ts +++ b/backend/src/__tests__/scheduler-service.test.ts @@ -12,7 +12,7 @@ const { mockUpdateScheduledTask, mockCleanupOldTaskRuns, mockGetScheduledTask, mockGetNodes, mockGetNode, mockCreateSnapshot, mockInsertSnapshotFiles, mockClearStackUpdateStatus, mockMarkStaleRunsAsFailed, mockDeleteOldScans, - mockGetTier, mockGetVariant, + mockGetTier, mockGetVariant, mockGetProxyHeaders, mockGetContainersByStack, mockRestartContainer, mockPruneSystem, mockUpdateStack, mockGetStacks, mockGetStackContent, mockGetEnvContent, @@ -42,6 +42,7 @@ const { mockDeleteOldScans: vi.fn().mockReturnValue(0), mockGetTier: vi.fn().mockReturnValue('paid'), mockGetVariant: vi.fn().mockReturnValue('admiral'), + mockGetProxyHeaders: vi.fn().mockReturnValue({ tier: 'paid', variant: 'admiral' }), mockGetContainersByStack: vi.fn().mockResolvedValue([]), mockRestartContainer: vi.fn().mockResolvedValue(undefined), mockPruneSystem: vi.fn().mockResolvedValue({ success: true, reclaimedBytes: 0 }), @@ -94,6 +95,7 @@ vi.mock('../services/LicenseService', () => ({ getInstance: () => ({ getTier: mockGetTier, getVariant: mockGetVariant, + getProxyHeaders: mockGetProxyHeaders, }), }, })); @@ -1349,6 +1351,11 @@ describe('SchedulerService - executeUpdateRemote', () => { 'http://remote:1852/api/auto-update/execute', expect.objectContaining({ method: 'POST', + headers: expect.objectContaining({ + 'Authorization': 'Bearer test-token', + 'x-sencho-tier': 'paid', + 'x-sencho-variant': 'admiral', + }), body: JSON.stringify({ target: 'web-app' }), }) ); diff --git a/backend/src/__tests__/stacks-failure-notifications.test.ts b/backend/src/__tests__/stacks-failure-notifications.test.ts index 0c15dbe0..0ce7f027 100644 --- a/backend/src/__tests__/stacks-failure-notifications.test.ts +++ b/backend/src/__tests__/stacks-failure-notifications.test.ts @@ -9,7 +9,9 @@ */ import { describe, it, expect, beforeAll, afterAll, vi, beforeEach } from 'vitest'; import request from 'supertest'; -import { setupTestDb, cleanupTestDb, loginAsTestAdmin } from './helpers/setupTestDb'; +import jwt from 'jsonwebtoken'; +import { setupTestDb, cleanupTestDb, loginAsTestAdmin, TEST_JWT_SECRET } from './helpers/setupTestDb'; +import { ComposeRollbackError } from '../services/ComposeService'; // ── Hoisted mocks (must come before importing the app) ────────────────────── @@ -143,6 +145,46 @@ describe('deploy_failure notification on /deploy error', () => { expect(call[2]).toContain('network timeout'); expect(call[3]).toEqual({ stackName: 'webapp' }); }); + + it('returns rolledBack=true only when compose rollback completed', async () => { + mockDeployStack.mockRejectedValue( + new ComposeRollbackError(new Error('image pull failed'), true, true), + ); + + const res = await request(app) + .post('/api/stacks/myapp/deploy') + .set('Cookie', authCookie); + + expect(res.status).toBe(500); + expect(res.body).toMatchObject({ rolledBack: true }); + }); + + it('returns rolledBack=false when compose rollback failed', async () => { + mockDeployStack.mockRejectedValue( + new ComposeRollbackError(new Error('image pull failed'), true, false), + ); + + const res = await request(app) + .post('/api/stacks/myapp/deploy') + .set('Cookie', authCookie); + + expect(res.status).toBe(500); + expect(res.body).toMatchObject({ rolledBack: false }); + }); + + it('uses trusted proxy tier headers for remote atomic deploys', async () => { + mockDeployStack.mockResolvedValue(undefined); + const token = jwt.sign({ scope: 'node_proxy' }, TEST_JWT_SECRET, { expiresIn: '1m' }); + + const res = await request(app) + .post('/api/stacks/myapp/deploy') + .set('Authorization', `Bearer ${token}`) + .set('x-sencho-tier', 'paid') + .set('x-sencho-variant', 'skipper'); + + expect(res.status).toBe(200); + expect(mockDeployStack.mock.calls[0][2]).toBe(true); + }); }); describe('deploy_failure notification on /down error', () => { @@ -229,4 +271,17 @@ describe('deploy_failure notification on /update error', () => { { stackName: 'myapp' }, ); }); + + it('returns rollback completion status when updateStack throws rollback metadata', async () => { + mockUpdateStack.mockRejectedValue( + new ComposeRollbackError(new Error('image not found'), true, false), + ); + + const res = await request(app) + .post('/api/stacks/myapp/update') + .set('Cookie', authCookie); + + expect(res.status).toBe(500); + expect(res.body).toMatchObject({ rolledBack: false }); + }); }); diff --git a/backend/src/routes/imageUpdates.ts b/backend/src/routes/imageUpdates.ts index de6ebf5f..bd92ff8c 100644 --- a/backend/src/routes/imageUpdates.ts +++ b/backend/src/routes/imageUpdates.ts @@ -6,11 +6,10 @@ import { CacheService } from '../services/CacheService'; import { ImageUpdateService } from '../services/ImageUpdateService'; import { FileSystemService } from '../services/FileSystemService'; import { ComposeService } from '../services/ComposeService'; -import { LicenseService } from '../services/LicenseService'; import { NotificationService } from '../services/NotificationService'; import { enforcePolicyPreDeploy } from '../services/PolicyEnforcement'; import { authMiddleware } from '../middleware/auth'; -import { requireAdmin, requirePaid } from '../middleware/tierGates'; +import { effectiveTier, requireAdmin, requirePaid } from '../middleware/tierGates'; import { buildPolicyGateOptions } from '../helpers/policyGate'; import { isValidStackName } from '../utils/validation'; import { sanitizeForLog } from '../utils/safeLog'; @@ -215,7 +214,7 @@ autoUpdateRouter.post('/execute', authMiddleware, async (req: Request, res: Resp const imageUpdateService = ImageUpdateService.getInstance(); const compose = ComposeService.getInstance(req.nodeId); const db = DatabaseService.getInstance(); - const atomic = LicenseService.getInstance().getTier() === 'paid'; + const atomic = effectiveTier(req) === 'paid'; const results: string[] = []; for (const stackName of stackNames) { diff --git a/backend/src/routes/stacks.ts b/backend/src/routes/stacks.ts index ae33bfed..fed8c4db 100644 --- a/backend/src/routes/stacks.ts +++ b/backend/src/routes/stacks.ts @@ -3,16 +3,15 @@ import path from 'path'; import YAML from 'yaml'; import multer from 'multer'; import { FileSystemService } from '../services/FileSystemService'; -import { ComposeService } from '../services/ComposeService'; +import { ComposeService, getComposeRollbackInfo } from '../services/ComposeService'; import DockerController from '../services/DockerController'; import { DatabaseService } from '../services/DatabaseService'; import { CacheService } from '../services/CacheService'; -import { LicenseService } from '../services/LicenseService'; import { UpdatePreviewService } from '../services/UpdatePreviewService'; import { GitSourceService, GitSourceError, repoHost as gitRepoHost } from '../services/GitSourceService'; import { enforcePolicyPreDeploy } from '../services/PolicyEnforcement'; import { requirePermission } from '../middleware/permissions'; -import { requirePaid, requireAdmin } from '../middleware/tierGates'; +import { requirePaid, requireAdmin, effectiveTier } from '../middleware/tierGates'; import { NotificationService, type NotificationCategory } from '../services/NotificationService'; import { isValidStackName, isValidServiceName, isPathWithinBase, isValidRelativeStackPath } from '../utils/validation'; import { getErrorMessage } from '../utils/errors'; @@ -588,7 +587,7 @@ stacksRouter.post('/:stackName/deploy', async (req: Request, res: Response) => { try { if (!(await runPolicyGate(req, res, stackName, req.nodeId))) return; const debug = isDebugEnabled(); - const atomic = LicenseService.getInstance().getTier() === 'paid'; + const atomic = effectiveTier(req) === 'paid'; if (debug) console.debug('[Stacks:debug] Deploy starting', { stackName, atomic, nodeId: req.nodeId }); const t0 = Date.now(); await ComposeService.getInstance(req.nodeId).deployStack(stackName, getTerminalWs(), atomic); @@ -602,8 +601,13 @@ stacksRouter.post('/:stackName/deploy', async (req: Request, res: Response) => { ); } catch (error: unknown) { console.error('[Stacks] Deploy failed: %s', sanitizeForLog(stackName), error); - const rolledBack = LicenseService.getInstance().getTier() === 'paid'; - if (rolledBack) console.warn('[Stacks] Deploy failed, rolled back: %s', sanitizeForLog(stackName)); + const rollbackInfo = getComposeRollbackInfo(error); + const rolledBack = rollbackInfo?.rolledBack ?? false; + if (rolledBack) { + console.warn('[Stacks] Deploy failed, rolled back: %s', sanitizeForLog(stackName)); + } else if (rollbackInfo?.attempted) { + console.warn('[Stacks] Deploy failed, rollback did not complete: %s', sanitizeForLog(stackName)); + } const message = getErrorMessage(error, 'Failed to deploy stack'); notifyActionFailure('deploy', stackName, error); res.status(500).json({ error: message, rolledBack }); @@ -762,7 +766,7 @@ stacksRouter.post('/:stackName/update', async (req: Request, res: Response) => { try { if (!(await runPolicyGate(req, res, stackName, req.nodeId))) return; const debug = isDebugEnabled(); - const atomic = LicenseService.getInstance().getTier() === 'paid'; + const atomic = effectiveTier(req) === 'paid'; if (debug) console.debug('[Stacks:debug] Update starting', { stackName, atomic, nodeId: req.nodeId }); const t0 = Date.now(); await ComposeService.getInstance(req.nodeId).updateStack(stackName, getTerminalWs(), atomic); @@ -777,8 +781,13 @@ stacksRouter.post('/:stackName/update', async (req: Request, res: Response) => { ); } catch (error: unknown) { console.error('[Stacks] Update failed: %s', sanitizeForLog(stackName), error); - const rolledBack = LicenseService.getInstance().getTier() === 'paid'; - if (rolledBack) console.warn(`[Stacks] Update failed, rolled back: ${sanitizeForLog(stackName)}`); + const rollbackInfo = getComposeRollbackInfo(error); + const rolledBack = rollbackInfo?.rolledBack ?? false; + if (rolledBack) { + console.warn(`[Stacks] Update failed, rolled back: ${sanitizeForLog(stackName)}`); + } else if (rollbackInfo?.attempted) { + console.warn(`[Stacks] Update failed, rollback did not complete: ${sanitizeForLog(stackName)}`); + } notifyActionFailure('update', stackName, error); res.status(500).json({ error: getErrorMessage(error, 'Failed to update'), rolledBack }); } diff --git a/backend/src/routes/templates.ts b/backend/src/routes/templates.ts index f7891857..cbaf1333 100644 --- a/backend/src/routes/templates.ts +++ b/backend/src/routes/templates.ts @@ -2,13 +2,12 @@ import { Router, type Request, type Response } from 'express'; import path from 'path'; import { promises as fsPromises } from 'fs'; import { authMiddleware } from '../middleware/auth'; -import { requireAdmin } from '../middleware/tierGates'; +import { effectiveTier, requireAdmin } from '../middleware/tierGates'; import { requirePermission } from '../middleware/permissions'; import { templateService } from '../services/TemplateService'; import { FileSystemService } from '../services/FileSystemService'; import { ComposeService } from '../services/ComposeService'; import { DatabaseService } from '../services/DatabaseService'; -import { LicenseService } from '../services/LicenseService'; import { ErrorParser } from '../utils/ErrorParser'; import { isValidStackName, isPathWithinBase } from '../utils/validation'; import { isDebugEnabled } from '../utils/debug'; @@ -125,7 +124,7 @@ templatesRouter.post('/deploy', authMiddleware, async (req: Request, res: Respon } return; } - const atomic = LicenseService.getInstance().getTier() === 'paid'; + const atomic = effectiveTier(req) === 'paid'; await ComposeService.getInstance(req.nodeId).deployStack(stackName, getTerminalWs(), atomic); invalidateNodeCaches(req.nodeId); console.log(`[Templates] Deploy completed: ${stackName}`); @@ -141,25 +140,31 @@ templatesRouter.post('/deploy', authMiddleware, async (req: Request, res: Respon const parsed = ErrorParser.parse(rawError); const shouldRollback = parsed.rule ? parsed.rule.canSilentlyRollback : true; + let rolledBack = false; if (shouldRollback) { + let dockerDownCompleted = true; + let fileDeleteCompleted = true; try { await ComposeService.getInstance(req.nodeId).downStack(stackName); } catch (downErr) { + dockerDownCompleted = false; console.error("[Templates] Rollback Stage 1 (Docker down) failed:", downErr); } try { await fsService.deleteStack(stackName); } catch (fsErr) { + fileDeleteCompleted = false; console.error("[Templates] Rollback Stage 2 (File deletion) failed:", fsErr); } + rolledBack = dockerDownCompleted && fileDeleteCompleted; } invalidateNodeCaches(req.nodeId); res.status(500).json({ error: parsed.message, - rolledBack: shouldRollback, + rolledBack, ruleId: parsed.rule?.id || 'UNKNOWN' }); } diff --git a/backend/src/services/ComposeService.ts b/backend/src/services/ComposeService.ts index 492e4d50..0c14f0f5 100644 --- a/backend/src/services/ComposeService.ts +++ b/backend/src/services/ComposeService.ts @@ -1,488 +1,516 @@ -import { spawn } from 'child_process'; -import fs from 'fs'; -import os from 'os'; -import path from 'path'; -import WebSocket from 'ws'; -import DockerController from './DockerController'; -import { DatabaseService } from './DatabaseService'; -import { FileSystemService } from './FileSystemService'; -import { MeshService } from './MeshService'; -import { LogFormatter } from './LogFormatter'; -import { NodeRegistry } from './NodeRegistry'; -import { RegistryService } from './RegistryService'; - -import { isDebugEnabled } from '../utils/debug'; -import { isPathWithinBase, isValidStackName } from '../utils/validation'; -import { sanitizeForLog } from '../utils/safeLog'; - -/** - * ComposeService - local docker compose CLI execution. - * - * In the Distributed API model, remote node compose operations are handled - * by the remote Sencho instance. This service only executes commands locally. - */ -export class ComposeService { - private baseDir: string; - private nodeId: number; - - constructor(nodeId?: number) { - this.nodeId = nodeId ?? NodeRegistry.getInstance().getDefaultNodeId(); - this.baseDir = NodeRegistry.getInstance().getComposeDir(this.nodeId); - } - - public static getInstance(nodeId?: number): ComposeService { - return new ComposeService(nodeId); - } - - /** - * Build the `docker compose` argument prefix for a stack, splicing in the - * Sencho Mesh override file if the stack is opted into the mesh. When no - * override applies, returns args without `-f` so docker compose's built-in - * file discovery resolves the stack's actual compose filename. The user's - * source compose file is never mutated. - */ - private async composeArgs(stackName: string, action: string[]): Promise { - const args: string[] = ['compose']; - let overridePath: string | null = null; - try { - overridePath = await MeshService.getInstance().ensureStackOverride(this.nodeId, stackName); - } catch (err) { - console.warn('[ComposeService] mesh override skipped:', sanitizeForLog((err as Error).message)); - } - if (overridePath) { - const baseFilename = await FileSystemService.getInstance(this.nodeId).getComposeFilename(stackName); - args.push('-f', baseFilename, '-f', overridePath); - } - args.push(...action); - return args; - } - - private execute( - command: string, - args: string[], - cwd: string, - ws?: WebSocket, - throwOnError = true, - env?: Record - ): Promise { - return new Promise((resolve, reject) => { - const child = spawn(command, args, { - cwd, - env: env ?? { - ...process.env, - PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin' - } - }); - - let errorLog = ''; - - const onData = (data: Buffer) => { - const text = data.toString(); - errorLog += text; - if (ws && ws.readyState === WebSocket.OPEN) { - ws.send(text); - } - }; - - child.stdout.on('data', onData); - child.stderr.on('data', onData); - - child.on('close', (code: number | null) => { - if (ws && ws.readyState === WebSocket.OPEN) { - ws.send(`Command exited with code ${code}\n`); - } - if (code === 0) resolve(); - else if (throwOnError) reject(new Error(errorLog.trim() || `Command failed with code ${code}`)); - else resolve(); - }); - - child.on('error', (error: Error) => { - if (ws && ws.readyState === WebSocket.OPEN) { - ws.send(`Error: ${error.message}\n`); - } - if (throwOnError) reject(error); - else resolve(); - }); - }); - } - - private async withRegistryAuth( - fn: (env: Record) => Promise, - sendOutput?: (data: string) => void, - ): Promise { - const registries = DatabaseService.getInstance().getRegistries(); - if (registries.length === 0) { - return fn({ - ...process.env, - PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin', - }); - } - - const { config, warnings } = await RegistryService.getInstance().resolveDockerConfig(); - if (warnings.length > 0 && sendOutput) { - for (const warning of warnings) { - sendOutput(`[Sencho] Warning: ${warning}\n`); - } - } - - const tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'sencho-docker-')); - const configPath = path.join(tmpDir, 'config.json'); - - try { - fs.writeFileSync(configPath, JSON.stringify(config), { mode: 0o600 }); - return await fn({ - ...process.env, - DOCKER_CONFIG: tmpDir, - PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin', - }); - } finally { - // Best-effort cleanup; each step runs independently so a file that was never - // written (e.g., writeFileSync threw) does not prevent the directory removal. - try { fs.unlinkSync(configPath); } catch { /* file may not exist */ } - try { fs.rmdirSync(tmpDir); } catch (e) { - console.warn('[ComposeService] Could not remove temp Docker config dir:', (e as Error).message); - } - } - } - - async runCommand(stackName: string, action: 'down' | 'start' | 'stop' | 'restart', ws?: WebSocket): Promise { - const stackDir = path.join(this.baseDir, stackName); - await this.execute('docker', ['compose', action], stackDir, ws); - } - - async deployStack(stackName: string, ws?: WebSocket, atomic?: boolean): Promise { - const stackDir = path.join(this.baseDir, stackName); - const debug = isDebugEnabled(); - const t0 = Date.now(); - if (debug) console.debug('[ComposeService:debug] deployStack', { stackName, stackDir, atomic }); - const sendOutput = (data: string) => { - if (ws && ws.readyState === WebSocket.OPEN) ws.send(data); - }; - - // Atomic: backup files before deploying - if (atomic) { - try { - const fsSvc = FileSystemService.getInstance(this.nodeId); - await fsSvc.backupStackFiles(stackName); - sendOutput('=== Backup created for atomic deployment ===\n'); - } catch (e) { - console.warn('Failed to backup stack files for %s:', sanitizeForLog(stackName), e); - } - } - - try { - try { - const dockerController = DockerController.getInstance(this.nodeId); - const legacyContainers = await dockerController.getContainersByStack(stackName); - if (legacyContainers && legacyContainers.length > 0) { - sendOutput(`=== Cleaning up existing containers for clean deployment ===\n`); - await dockerController.removeContainers(legacyContainers.map((c: any) => c.Id)); - } - } catch (e) { - console.warn('Failed to clean up legacy containers for %s:', sanitizeForLog(stackName), e); - } - - await this.withRegistryAuth(async (env) => { - await this.execute('docker', await this.composeArgs(stackName, ['up', '-d', '--remove-orphans']), stackDir, ws, true, env); - }, sendOutput); - - // Post-Deploy Health Probe - await new Promise(resolve => setTimeout(resolve, 3000)); - - const dockerController = DockerController.getInstance(this.nodeId); - const containers = await dockerController.getDocker().listContainers({ - all: true, - filters: { label: [`com.docker.compose.project=${stackName}`] } - }); - - for (const containerInfo of containers) { - if (containerInfo.State === 'exited') { - const container = dockerController.getDocker().getContainer(containerInfo.Id); - const inspectData = await container.inspect(); - const exitCode = inspectData.State.ExitCode; - - if (exitCode !== 0) { - const logs = await container.logs({ stdout: true, stderr: true, tail: 50 }); - const logStr = logs.toString('utf-8'); - throw new Error(`CONTAINER_CRASHED\nExit Code: ${exitCode}\n${logStr}`); - } - } - } - if (debug) console.debug(`[ComposeService:debug] deployStack completed in ${Date.now() - t0}ms`, { stackName }); - } catch (deployError) { - // Atomic: auto-rollback on failure - if (atomic) { - sendOutput('\n=== Deployment failed - rolling back to previous version ===\n'); - try { - const fsSvc = FileSystemService.getInstance(this.nodeId); - await fsSvc.restoreStackFiles(stackName); - await this.withRegistryAuth(async (env) => { - await this.execute('docker', await this.composeArgs(stackName, ['up', '-d', '--remove-orphans']), stackDir, ws, true, env); - }, sendOutput); - sendOutput('=== Rolled back successfully ===\n'); - } catch (rollbackError) { - console.error('Rollback failed for %s:', sanitizeForLog(stackName), rollbackError); - sendOutput('=== Rollback failed - manual intervention may be required ===\n'); - } - } - throw deployError; - } - } - - streamLogs(stackName: string, ws: WebSocket) { - let isClosed = false; - let isFirstRun = true; - let isWaitingForActivity = false; - - ws.on('close', () => { isClosed = true; }); - - const startStream = async () => { - if (isClosed || ws.readyState !== WebSocket.OPEN) return; - - try { - const dockerController = DockerController.getInstance(this.nodeId); - const containers = await dockerController.getContainersByStack(stackName); - - if (!containers || containers.length === 0) { - if (!isWaitingForActivity) { - ws.send(`\r\n\x1b[33m[Sencho] No containers found. Waiting for activity...\x1b[0m\r\n`); - isWaitingForActivity = true; - } - setTimeout(startStream, 2000); - return; - } - - const runningContainers = containers.filter((c: any) => c.State === 'running'); - - if (!isFirstRun && runningContainers.length === 0) { - if (!isWaitingForActivity) { - ws.send(`\r\n\x1b[33m[Sencho] Log stream ended. Waiting for container activity...\x1b[0m\r\n`); - isWaitingForActivity = true; - } - setTimeout(startStream, 2000); - return; - } - - const containersToLog = isFirstRun ? containers : runningContainers; - isFirstRun = false; - isWaitingForActivity = false; - - let activeProcesses = 0; - let streamEndedHandled = false; - const localProcesses: ReturnType[] = []; - - const onWsClose = () => { - localProcesses.forEach(cp => { try { cp.kill(); } catch { } }); - }; - - ws.on('close', onWsClose); - - const handleProcessEnd = () => { - activeProcesses--; - if (activeProcesses <= 0 && !streamEndedHandled) { - streamEndedHandled = true; - ws.removeListener('close', onWsClose); - if (!isClosed && ws.readyState === WebSocket.OPEN) { - setTimeout(startStream, 1000); - } - } - }; - - for (const container of containersToLog) { - const containerName = container.Names?.[0]?.replace(/^\//, '') || container.Id; - activeProcesses++; - let lineBuffer = ''; - - const sendOutput = (data: Buffer) => { - if (ws.readyState === WebSocket.OPEN) { - lineBuffer += data.toString(); - const lines = lineBuffer.split(/\r?\n/); - lineBuffer = lines.pop() || ''; - for (const line of lines) { - ws.send(LogFormatter.process(line) + '\r\n'); - } - } - }; - - const flushBuffer = () => { - if (lineBuffer && ws.readyState === WebSocket.OPEN) { - ws.send(LogFormatter.process(lineBuffer) + '\r\n'); - lineBuffer = ''; - } - }; - - const child = spawn('docker', ['logs', '-f', '-t', '--tail', '100', containerName], { - env: { - ...process.env, - PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin' - } - }); - localProcesses.push(child); - child.stdout.on('data', sendOutput); - child.stderr.on('data', sendOutput); - child.on('error', handleProcessEnd); - child.on('close', () => { - flushBuffer(); - handleProcessEnd(); - }); - } - } catch (err) { - if (!isClosed && ws.readyState === WebSocket.OPEN) { - if (!isWaitingForActivity) { - ws.send(`\r\n\x1b[31m[Sencho] Error tracking containers. Retrying...\x1b[0m\r\n`); - isWaitingForActivity = true; - } - setTimeout(startStream, 2000); - } - } - }; - - startStream(); - } - - async updateStack(stackName: string, ws?: WebSocket, atomic?: boolean): Promise { - const stackDir = path.join(this.baseDir, stackName); - const debug = isDebugEnabled(); - const t0 = Date.now(); - if (debug) console.debug('[ComposeService:debug] updateStack', { stackName, stackDir, atomic }); - const sendOutput = (data: string) => { - if (ws && ws.readyState === WebSocket.OPEN) ws.send(data); - }; - - // Atomic: backup files before updating - if (atomic) { - try { - const fsSvc = FileSystemService.getInstance(this.nodeId); - await fsSvc.backupStackFiles(stackName); - sendOutput('=== Backup created for atomic update ===\n'); - } catch (e) { - console.warn('Failed to backup stack files for %s:', sanitizeForLog(stackName), e); - } - } - - try { - try { - const dockerController = DockerController.getInstance(this.nodeId); - const legacyContainers = await dockerController.getContainersByStack(stackName); - if (legacyContainers && legacyContainers.length > 0) { - sendOutput(`=== Cleaning up existing containers for clean update ===\n`); - await dockerController.removeContainers(legacyContainers.map((c: any) => c.Id)); - } - } catch (e) { - console.warn('Failed to clean up legacy containers for %s:', sanitizeForLog(stackName), e); - } - - await this.withRegistryAuth(async (env) => { - sendOutput('=== Pulling latest images ===\n'); - await this.execute('docker', ['compose', 'pull'], stackDir, ws, true, env); - - sendOutput('=== Recreating containers ===\n'); - await this.execute('docker', await this.composeArgs(stackName, ['up', '-d', '--remove-orphans']), stackDir, ws, true, env); - }, sendOutput); - - // Post-Update Health Probe - await new Promise(resolve => setTimeout(resolve, 3000)); - - const dockerController = DockerController.getInstance(this.nodeId); - const containers = await dockerController.getDocker().listContainers({ - all: true, - filters: { label: [`com.docker.compose.project=${stackName}`] } - }); - - for (const containerInfo of containers) { - if (containerInfo.State === 'exited') { - const container = dockerController.getDocker().getContainer(containerInfo.Id); - const inspectData = await container.inspect(); - const exitCode = inspectData.State.ExitCode; - - if (exitCode !== 0) { - const logs = await container.logs({ stdout: true, stderr: true, tail: 50 }); - const logStr = logs.toString('utf-8'); - throw new Error(`CONTAINER_CRASHED\nExit Code: ${exitCode}\n${logStr}`); - } - } - } - - sendOutput('=== Stack updated successfully ===\n'); - if (debug) console.debug(`[ComposeService:debug] updateStack completed in ${Date.now() - t0}ms`, { stackName }); - } catch (updateError) { - // Atomic: auto-rollback on failure - if (atomic) { - sendOutput('\n=== Update failed - rolling back to previous version ===\n'); - try { - const fsSvc = FileSystemService.getInstance(this.nodeId); - await fsSvc.restoreStackFiles(stackName); - await this.withRegistryAuth(async (env) => { - await this.execute('docker', await this.composeArgs(stackName, ['up', '-d', '--remove-orphans']), stackDir, ws, true, env); - }, sendOutput); - sendOutput('=== Rolled back successfully ===\n'); - } catch (rollbackError) { - console.error('Rollback failed for %s:', sanitizeForLog(stackName), rollbackError); - sendOutput('=== Rollback failed - manual intervention may be required ===\n'); - } - } - throw updateError; - } - } - - public async downStack(stackName: string): Promise { - const stackPath = path.join(this.baseDir, stackName); - try { - await this.execute('docker', ['compose', 'down', '--volumes', '--remove-orphans'], stackPath, undefined, false); - } catch (error) { - console.warn(`[Teardown] Docker down failed or nothing to clean up for ${sanitizeForLog(stackName)}`); - } - } - - /** - * Enumerate image references declared in a stack's compose file. - * - * Used by the pre-deploy policy gate to decide which images to scan before - * `docker compose up` runs. Path traversal is guarded against the node's - * compose base directory; missing / unreadable compose files or `.env` - * interpolation failures surface as a rejected Promise so the gate can - * block the deploy rather than silently allow it. - */ - public async listStackImages(stackName: string): Promise { - if (!isValidStackName(stackName)) { - throw new Error('Invalid stack path'); - } - const stackDir = path.resolve(this.baseDir, stackName); - if (!isPathWithinBase(stackDir, this.baseDir) || path.resolve(this.baseDir) === stackDir) { - throw new Error('Invalid stack path'); - } - const stdout = await this.captureCompose(['config', '--images'], stackDir); - const seen = new Set(); - const images: string[] = []; - for (const raw of stdout.split(/\r?\n/)) { - const line = raw.trim(); - if (!line) continue; - if (line.startsWith('sha256:')) continue; - if (seen.has(line)) continue; - seen.add(line); - images.push(line); - } - return images; - } - - private captureCompose(args: string[], cwd: string): Promise { - return new Promise((resolve, reject) => { - const child = spawn('docker', ['compose', ...args], { - cwd, - env: { - ...process.env, - PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin', - }, - }); - let stdout = ''; - let stderr = ''; - child.stdout.on('data', (data: Buffer) => { stdout += data.toString(); }); - child.stderr.on('data', (data: Buffer) => { stderr += data.toString(); }); - child.on('error', (err) => reject(err)); - child.on('close', (code) => { - if (code === 0) resolve(stdout); - else reject(new Error(stderr.trim() || `docker compose ${args.join(' ')} failed with code ${code}`)); - }); - }); - } -} +import { spawn } from 'child_process'; +import fs from 'fs'; +import os from 'os'; +import path from 'path'; +import WebSocket from 'ws'; +import DockerController from './DockerController'; +import { DatabaseService } from './DatabaseService'; +import { FileSystemService } from './FileSystemService'; +import { MeshService } from './MeshService'; +import { LogFormatter } from './LogFormatter'; +import { NodeRegistry } from './NodeRegistry'; +import { RegistryService } from './RegistryService'; + +import { isDebugEnabled } from '../utils/debug'; +import { getErrorMessage } from '../utils/errors'; +import { isPathWithinBase, isValidStackName } from '../utils/validation'; +import { sanitizeForLog } from '../utils/safeLog'; + +export class ComposeRollbackError extends Error { + public readonly rollbackAttempted: boolean; + public readonly rolledBack: boolean; + public readonly originalError: unknown; + + constructor(originalError: unknown, rollbackAttempted: boolean, rolledBack: boolean) { + super(getErrorMessage(originalError, 'Compose operation failed')); + this.name = 'ComposeRollbackError'; + this.rollbackAttempted = rollbackAttempted; + this.rolledBack = rolledBack; + this.originalError = originalError; + Object.setPrototypeOf(this, ComposeRollbackError.prototype); + } +} + +export function getComposeRollbackInfo(error: unknown): { attempted: boolean; rolledBack: boolean } | null { + if (!(error instanceof ComposeRollbackError)) { + return null; + } + return { attempted: error.rollbackAttempted, rolledBack: error.rolledBack }; +} + +/** + * ComposeService - local docker compose CLI execution. + * + * In the Distributed API model, remote node compose operations are handled + * by the remote Sencho instance. This service only executes commands locally. + */ +export class ComposeService { + private baseDir: string; + private nodeId: number; + + constructor(nodeId?: number) { + this.nodeId = nodeId ?? NodeRegistry.getInstance().getDefaultNodeId(); + this.baseDir = NodeRegistry.getInstance().getComposeDir(this.nodeId); + } + + public static getInstance(nodeId?: number): ComposeService { + return new ComposeService(nodeId); + } + + /** + * Build the `docker compose` argument prefix for a stack, splicing in the + * Sencho Mesh override file if the stack is opted into the mesh. When no + * override applies, returns args without `-f` so docker compose's built-in + * file discovery resolves the stack's actual compose filename. The user's + * source compose file is never mutated. + */ + private async composeArgs(stackName: string, action: string[]): Promise { + const args: string[] = ['compose']; + let overridePath: string | null = null; + try { + overridePath = await MeshService.getInstance().ensureStackOverride(this.nodeId, stackName); + } catch (err) { + console.warn('[ComposeService] mesh override skipped:', sanitizeForLog((err as Error).message)); + } + if (overridePath) { + const baseFilename = await FileSystemService.getInstance(this.nodeId).getComposeFilename(stackName); + args.push('-f', baseFilename, '-f', overridePath); + } + args.push(...action); + return args; + } + + private execute( + command: string, + args: string[], + cwd: string, + ws?: WebSocket, + throwOnError = true, + env?: Record + ): Promise { + return new Promise((resolve, reject) => { + const child = spawn(command, args, { + cwd, + env: env ?? { + ...process.env, + PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin' + } + }); + + let errorLog = ''; + + const onData = (data: Buffer) => { + const text = data.toString(); + errorLog += text; + if (ws && ws.readyState === WebSocket.OPEN) { + ws.send(text); + } + }; + + child.stdout.on('data', onData); + child.stderr.on('data', onData); + + child.on('close', (code: number | null) => { + if (ws && ws.readyState === WebSocket.OPEN) { + ws.send(`Command exited with code ${code}\n`); + } + if (code === 0) resolve(); + else if (throwOnError) reject(new Error(errorLog.trim() || `Command failed with code ${code}`)); + else resolve(); + }); + + child.on('error', (error: Error) => { + if (ws && ws.readyState === WebSocket.OPEN) { + ws.send(`Error: ${error.message}\n`); + } + if (throwOnError) reject(error); + else resolve(); + }); + }); + } + + private async withRegistryAuth( + fn: (env: Record) => Promise, + sendOutput?: (data: string) => void, + ): Promise { + const registries = DatabaseService.getInstance().getRegistries(); + if (registries.length === 0) { + return fn({ + ...process.env, + PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin', + }); + } + + const { config, warnings } = await RegistryService.getInstance().resolveDockerConfig(); + if (warnings.length > 0 && sendOutput) { + for (const warning of warnings) { + sendOutput(`[Sencho] Warning: ${warning}\n`); + } + } + + const tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'sencho-docker-')); + const configPath = path.join(tmpDir, 'config.json'); + + try { + fs.writeFileSync(configPath, JSON.stringify(config), { mode: 0o600 }); + return await fn({ + ...process.env, + DOCKER_CONFIG: tmpDir, + PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin', + }); + } finally { + // Best-effort cleanup; each step runs independently so a file that was never + // written (e.g., writeFileSync threw) does not prevent the directory removal. + try { fs.unlinkSync(configPath); } catch { /* file may not exist */ } + try { fs.rmdirSync(tmpDir); } catch (e) { + console.warn('[ComposeService] Could not remove temp Docker config dir:', (e as Error).message); + } + } + } + + private async createAtomicBackup( + stackName: string, + operation: 'deployment' | 'update', + sendOutput: (data: string) => void, + ): Promise { + try { + const fsSvc = FileSystemService.getInstance(this.nodeId); + await fsSvc.backupStackFiles(stackName); + sendOutput(`=== Backup created for atomic ${operation} ===\n`); + } catch (error) { + console.error('Atomic backup failed for %s:', sanitizeForLog(stackName), getErrorMessage(error, 'unknown error')); + sendOutput(`=== Atomic ${operation} backup failed. Operation aborted ===\n`); + throw new Error(`Atomic ${operation} backup failed: ${getErrorMessage(error, 'unknown error')}`); + } + } + + private async restoreAtomicBackup( + stackName: string, + stackDir: string, + ws: WebSocket | undefined, + sendOutput: (data: string) => void, + ): Promise { + try { + const fsSvc = FileSystemService.getInstance(this.nodeId); + await fsSvc.restoreStackFiles(stackName); + await this.withRegistryAuth(async (env) => { + await this.execute('docker', await this.composeArgs(stackName, ['up', '-d', '--remove-orphans']), stackDir, ws, true, env); + }, sendOutput); + sendOutput('=== Rolled back successfully ===\n'); + return true; + } catch (rollbackError) { + console.error('Rollback failed for %s:', sanitizeForLog(stackName), getErrorMessage(rollbackError, 'unknown error')); + sendOutput('=== Rollback failed. Manual intervention may be required ===\n'); + return false; + } + } + + private createContainerCrashError(exitCode: number): Error { + return new Error( + `CONTAINER_CRASHED\nExit Code: ${exitCode}\nContainer exited after deployment. Check container logs for details.` + ); + } + + async runCommand(stackName: string, action: 'down' | 'start' | 'stop' | 'restart', ws?: WebSocket): Promise { + const stackDir = path.join(this.baseDir, stackName); + await this.execute('docker', ['compose', action], stackDir, ws); + } + + async deployStack(stackName: string, ws?: WebSocket, atomic?: boolean): Promise { + const stackDir = path.join(this.baseDir, stackName); + const debug = isDebugEnabled(); + const t0 = Date.now(); + if (debug) console.debug('[ComposeService:debug] deployStack', { stackName, stackDir, atomic }); + const sendOutput = (data: string) => { + if (ws && ws.readyState === WebSocket.OPEN) ws.send(data); + }; + + if (atomic) { + await this.createAtomicBackup(stackName, 'deployment', sendOutput); + } + + try { + try { + const dockerController = DockerController.getInstance(this.nodeId); + const legacyContainers = await dockerController.getContainersByStack(stackName); + if (legacyContainers && legacyContainers.length > 0) { + sendOutput(`=== Cleaning up existing containers for clean deployment ===\n`); + await dockerController.removeContainers(legacyContainers.map((c: any) => c.Id)); + } + } catch (e) { + console.warn('Failed to clean up legacy containers for %s:', sanitizeForLog(stackName), e); + } + + await this.withRegistryAuth(async (env) => { + await this.execute('docker', await this.composeArgs(stackName, ['up', '-d', '--remove-orphans']), stackDir, ws, true, env); + }, sendOutput); + + // Post-Deploy Health Probe + await new Promise(resolve => setTimeout(resolve, 3000)); + + const dockerController = DockerController.getInstance(this.nodeId); + const containers = await dockerController.getDocker().listContainers({ + all: true, + filters: { label: [`com.docker.compose.project=${stackName}`] } + }); + + for (const containerInfo of containers) { + if (containerInfo.State === 'exited') { + const container = dockerController.getDocker().getContainer(containerInfo.Id); + const inspectData = await container.inspect(); + const exitCode = inspectData.State.ExitCode; + + if (exitCode !== 0) { + throw this.createContainerCrashError(exitCode); + } + } + } + if (debug) console.debug(`[ComposeService:debug] deployStack completed in ${Date.now() - t0}ms`, { stackName }); + } catch (deployError) { + if (atomic) { + sendOutput('\n=== Deployment failed - rolling back to previous version ===\n'); + const rolledBack = await this.restoreAtomicBackup(stackName, stackDir, ws, sendOutput); + throw new ComposeRollbackError(deployError, true, rolledBack); + } + throw deployError; + } + } + + streamLogs(stackName: string, ws: WebSocket) { + let isClosed = false; + let isFirstRun = true; + let isWaitingForActivity = false; + + ws.on('close', () => { isClosed = true; }); + + const startStream = async () => { + if (isClosed || ws.readyState !== WebSocket.OPEN) return; + + try { + const dockerController = DockerController.getInstance(this.nodeId); + const containers = await dockerController.getContainersByStack(stackName); + + if (!containers || containers.length === 0) { + if (!isWaitingForActivity) { + ws.send(`\r\n\x1b[33m[Sencho] No containers found. Waiting for activity...\x1b[0m\r\n`); + isWaitingForActivity = true; + } + setTimeout(startStream, 2000); + return; + } + + const runningContainers = containers.filter((c: any) => c.State === 'running'); + + if (!isFirstRun && runningContainers.length === 0) { + if (!isWaitingForActivity) { + ws.send(`\r\n\x1b[33m[Sencho] Log stream ended. Waiting for container activity...\x1b[0m\r\n`); + isWaitingForActivity = true; + } + setTimeout(startStream, 2000); + return; + } + + const containersToLog = isFirstRun ? containers : runningContainers; + isFirstRun = false; + isWaitingForActivity = false; + + let activeProcesses = 0; + let streamEndedHandled = false; + const localProcesses: ReturnType[] = []; + + const onWsClose = () => { + localProcesses.forEach(cp => { try { cp.kill(); } catch { } }); + }; + + ws.on('close', onWsClose); + + const handleProcessEnd = () => { + activeProcesses--; + if (activeProcesses <= 0 && !streamEndedHandled) { + streamEndedHandled = true; + ws.removeListener('close', onWsClose); + if (!isClosed && ws.readyState === WebSocket.OPEN) { + setTimeout(startStream, 1000); + } + } + }; + + for (const container of containersToLog) { + const containerName = container.Names?.[0]?.replace(/^\//, '') || container.Id; + activeProcesses++; + let lineBuffer = ''; + + const sendOutput = (data: Buffer) => { + if (ws.readyState === WebSocket.OPEN) { + lineBuffer += data.toString(); + const lines = lineBuffer.split(/\r?\n/); + lineBuffer = lines.pop() || ''; + for (const line of lines) { + ws.send(LogFormatter.process(line) + '\r\n'); + } + } + }; + + const flushBuffer = () => { + if (lineBuffer && ws.readyState === WebSocket.OPEN) { + ws.send(LogFormatter.process(lineBuffer) + '\r\n'); + lineBuffer = ''; + } + }; + + const child = spawn('docker', ['logs', '-f', '-t', '--tail', '100', containerName], { + env: { + ...process.env, + PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin' + } + }); + localProcesses.push(child); + child.stdout.on('data', sendOutput); + child.stderr.on('data', sendOutput); + child.on('error', handleProcessEnd); + child.on('close', () => { + flushBuffer(); + handleProcessEnd(); + }); + } + } catch (err) { + if (!isClosed && ws.readyState === WebSocket.OPEN) { + if (!isWaitingForActivity) { + ws.send(`\r\n\x1b[31m[Sencho] Error tracking containers. Retrying...\x1b[0m\r\n`); + isWaitingForActivity = true; + } + setTimeout(startStream, 2000); + } + } + }; + + startStream(); + } + + async updateStack(stackName: string, ws?: WebSocket, atomic?: boolean): Promise { + const stackDir = path.join(this.baseDir, stackName); + const debug = isDebugEnabled(); + const t0 = Date.now(); + if (debug) console.debug('[ComposeService:debug] updateStack', { stackName, stackDir, atomic }); + const sendOutput = (data: string) => { + if (ws && ws.readyState === WebSocket.OPEN) ws.send(data); + }; + + if (atomic) { + await this.createAtomicBackup(stackName, 'update', sendOutput); + } + + try { + try { + const dockerController = DockerController.getInstance(this.nodeId); + const legacyContainers = await dockerController.getContainersByStack(stackName); + if (legacyContainers && legacyContainers.length > 0) { + sendOutput(`=== Cleaning up existing containers for clean update ===\n`); + await dockerController.removeContainers(legacyContainers.map((c: any) => c.Id)); + } + } catch (e) { + console.warn('Failed to clean up legacy containers for %s:', sanitizeForLog(stackName), e); + } + + await this.withRegistryAuth(async (env) => { + sendOutput('=== Pulling latest images ===\n'); + await this.execute('docker', ['compose', 'pull'], stackDir, ws, true, env); + + sendOutput('=== Recreating containers ===\n'); + await this.execute('docker', await this.composeArgs(stackName, ['up', '-d', '--remove-orphans']), stackDir, ws, true, env); + }, sendOutput); + + // Post-Update Health Probe + await new Promise(resolve => setTimeout(resolve, 3000)); + + const dockerController = DockerController.getInstance(this.nodeId); + const containers = await dockerController.getDocker().listContainers({ + all: true, + filters: { label: [`com.docker.compose.project=${stackName}`] } + }); + + for (const containerInfo of containers) { + if (containerInfo.State === 'exited') { + const container = dockerController.getDocker().getContainer(containerInfo.Id); + const inspectData = await container.inspect(); + const exitCode = inspectData.State.ExitCode; + + if (exitCode !== 0) { + throw this.createContainerCrashError(exitCode); + } + } + } + + sendOutput('=== Stack updated successfully ===\n'); + if (debug) console.debug(`[ComposeService:debug] updateStack completed in ${Date.now() - t0}ms`, { stackName }); + } catch (updateError) { + if (atomic) { + sendOutput('\n=== Update failed - rolling back to previous version ===\n'); + const rolledBack = await this.restoreAtomicBackup(stackName, stackDir, ws, sendOutput); + throw new ComposeRollbackError(updateError, true, rolledBack); + } + throw updateError; + } + } + + public async downStack(stackName: string): Promise { + const stackPath = path.join(this.baseDir, stackName); + try { + await this.execute('docker', ['compose', 'down', '--volumes', '--remove-orphans'], stackPath, undefined, false); + } catch (error) { + console.warn(`[Teardown] Docker down failed or nothing to clean up for ${sanitizeForLog(stackName)}`); + } + } + + /** + * Enumerate image references declared in a stack's compose file. + * + * Used by the pre-deploy policy gate to decide which images to scan before + * `docker compose up` runs. Path traversal is guarded against the node's + * compose base directory; missing / unreadable compose files or `.env` + * interpolation failures surface as a rejected Promise so the gate can + * block the deploy rather than silently allow it. + */ + public async listStackImages(stackName: string): Promise { + if (!isValidStackName(stackName)) { + throw new Error('Invalid stack path'); + } + const stackDir = path.resolve(this.baseDir, stackName); + if (!isPathWithinBase(stackDir, this.baseDir) || path.resolve(this.baseDir) === stackDir) { + throw new Error('Invalid stack path'); + } + const stdout = await this.captureCompose(['config', '--images'], stackDir); + const seen = new Set(); + const images: string[] = []; + for (const raw of stdout.split(/\r?\n/)) { + const line = raw.trim(); + if (!line) continue; + if (line.startsWith('sha256:')) continue; + if (seen.has(line)) continue; + seen.add(line); + images.push(line); + } + return images; + } + + private captureCompose(args: string[], cwd: string): Promise { + return new Promise((resolve, reject) => { + const child = spawn('docker', ['compose', ...args], { + cwd, + env: { + ...process.env, + PATH: process.env.PATH || '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin', + }, + }); + let stdout = ''; + let stderr = ''; + child.stdout.on('data', (data: Buffer) => { stdout += data.toString(); }); + child.stderr.on('data', (data: Buffer) => { stderr += data.toString(); }); + child.on('error', (err) => reject(err)); + child.on('close', (code) => { + if (code === 0) resolve(stdout); + else reject(new Error(stderr.trim() || `docker compose ${args.join(' ')} failed with code ${code}`)); + }); + }); + } +} diff --git a/backend/src/services/FileSystemService.ts b/backend/src/services/FileSystemService.ts index 6d0b5eba..79ea1df8 100644 --- a/backend/src/services/FileSystemService.ts +++ b/backend/src/services/FileSystemService.ts @@ -51,11 +51,11 @@ const MIME_MAP: Record = { */ export class FileSystemService { private baseDir: string; + private nodeId: number; constructor(nodeId?: number) { - this.baseDir = NodeRegistry.getInstance().getComposeDir( - nodeId ?? NodeRegistry.getInstance().getDefaultNodeId() - ); + this.nodeId = nodeId ?? NodeRegistry.getInstance().getDefaultNodeId(); + this.baseDir = NodeRegistry.getInstance().getComposeDir(this.nodeId); } public static getInstance(nodeId?: number): FileSystemService { @@ -77,6 +77,13 @@ export class FileSystemService { return stackDir; } + private getBackupDir(stackName: string): string { + if (!isValidStackName(stackName)) { + throw Object.assign(new Error('Invalid stack name'), { code: 'INVALID_STACK_NAME' }); + } + return path.join(getBackupBaseDir(), String(this.nodeId), stackName); + } + async hasComposeFile(dir: string): Promise { this.assertWithinBase(dir); const composeFiles = ['compose.yaml', 'compose.yml', 'docker-compose.yaml', 'docker-compose.yml']; @@ -298,7 +305,7 @@ export class FileSystemService { /** * Backup stack files (compose.yaml + .env) into Sencho's data dir. * - * Backups live at /backups// (NOT inside the user's + * Backups live at /backups/// (NOT inside the user's * compose folder) so the operation always succeeds even when the stack * folder is owned by another UID (e.g., a container running as root has * chowned its bind mount). DATA_DIR is the same writable location that @@ -308,7 +315,7 @@ export class FileSystemService { const debug = isDebugEnabled(); const t0 = Date.now(); const stackDir = this.resolveStackDir(stackName); - const backupDir = path.join(getBackupBaseDir(), stackName); + const backupDir = this.getBackupDir(stackName); await fsPromises.mkdir(backupDir, { recursive: true }); // Copy compose file @@ -347,7 +354,7 @@ export class FileSystemService { const debug = isDebugEnabled(); const t0 = Date.now(); const stackDir = this.resolveStackDir(stackName); - const backupDir = path.join(getBackupBaseDir(), stackName); + const backupDir = this.getBackupDir(stackName); const items = await fsPromises.readdir(backupDir); for (const item of items) { @@ -358,7 +365,7 @@ export class FileSystemService { } async getBackupInfo(stackName: string): Promise<{ exists: boolean; timestamp: number | null }> { - const backupDir = path.join(getBackupBaseDir(), stackName); + const backupDir = this.getBackupDir(stackName); try { await fsPromises.access(backupDir); const tsFile = path.join(backupDir, '.timestamp'); diff --git a/backend/src/services/SchedulerService.ts b/backend/src/services/SchedulerService.ts index 2b153255..2d190887 100644 --- a/backend/src/services/SchedulerService.ts +++ b/backend/src/services/SchedulerService.ts @@ -2,6 +2,7 @@ import { CronExpressionParser } from 'cron-parser'; import { DatabaseService } from './DatabaseService'; import type { ScheduledTask } from './DatabaseService'; import { LicenseService } from './LicenseService'; +import { PROXY_TIER_HEADER, PROXY_VARIANT_HEADER } from './license-headers'; import DockerController from './DockerController'; import { ComposeService } from './ComposeService'; import { FileSystemService } from './FileSystemService'; @@ -625,6 +626,7 @@ export class SchedulerService { } const baseUrl = proxyTarget.apiUrl.replace(/\/$/, ''); + const proxyHeaders = LicenseService.getInstance().getProxyHeaders(); if (isDebugEnabled()) { console.log(`[SchedulerService] executeUpdateRemote: node=${nodeId} target=${target}`); } @@ -634,6 +636,8 @@ export class SchedulerService { headers: { 'Content-Type': 'application/json', 'Authorization': `Bearer ${proxyTarget.apiToken}`, + [PROXY_TIER_HEADER]: proxyHeaders.tier, + [PROXY_VARIANT_HEADER]: proxyHeaders.variant ?? '', }, body: JSON.stringify({ target }), signal: AbortSignal.timeout(300_000), // 5 minute timeout for long updates diff --git a/frontend/src/components/EditorLayout/hooks/useStackActions.ts b/frontend/src/components/EditorLayout/hooks/useStackActions.ts index ca250198..ff9d7b7a 100644 --- a/frontend/src/components/EditorLayout/hooks/useStackActions.ts +++ b/frontend/src/components/EditorLayout/hooks/useStackActions.ts @@ -14,8 +14,11 @@ import type { PolicyBlockPayload } from '../../stack/PolicyBlockDialog'; interface RunResult { ok: boolean; errorMessage?: string; + rolledBack?: boolean; } +type StackActionError = Error & { rolledBack?: boolean }; + type EditorState = ReturnType; type StackListState = ReturnType; type NavState = ReturnType; @@ -36,6 +39,30 @@ interface UseStackActionsOptions { diffPreviewEnabled: boolean; } +const isRecord = (value: unknown): value is Record => + typeof value === 'object' && value !== null; + +const parseStackActionError = (rawBody: string, fallback: string): StackActionError => { + let message = rawBody || fallback; + let rolledBack = false; + + try { + const parsed: unknown = JSON.parse(rawBody); + if (isRecord(parsed)) { + if (typeof parsed.error === 'string' && parsed.error.trim()) { + message = parsed.error; + } + rolledBack = parsed.rolledBack === true; + } + } catch { + /* not JSON */ + } + + const error = new Error(message) as StackActionError; + error.rolledBack = rolledBack; + return error; +}; + export function useStackActions(options: UseStackActionsOptions) { const { editorState, @@ -349,7 +376,7 @@ export function useStackActions(options: UseStackActionsOptions) { stackFile: string, ignorePolicy: boolean, started?: Promise, - ): Promise<{ ok: boolean; errorMessage?: string }> => { + ): Promise => { const previousStatus = stackListState.stackStatuses[stackFile]; stackListState.setOptimisticStatus(stackFile, 'running'); try { @@ -381,7 +408,7 @@ export function useStackActions(options: UseStackActionsOptions) { }; } } - throw new Error(rawBody || 'Deploy failed'); + throw parseStackActionError(rawBody, 'Deploy failed'); } overlayState.setPolicyBlock(null); toast.success( @@ -405,13 +432,14 @@ export function useStackActions(options: UseStackActionsOptions) { console.error('Failed to deploy:', error); if (previousStatus !== undefined) stackListState.setOptimisticStatus(stackFile, previousStatus as 'running' | 'exited'); - const errorMessage = (error as Error).message || 'Failed to deploy stack'; + const deployError = error as StackActionError; + const errorMessage = deployError.message || 'Failed to deploy stack'; toast.error( - isPaid + isPaid && deployError.rolledBack === true ? `${errorMessage} - automatically rolled back to previous version.` : errorMessage, ); - return { ok: false, errorMessage }; + return { ok: false, errorMessage, rolledBack: deployError.rolledBack }; } }; @@ -550,7 +578,12 @@ export function useStackActions(options: UseStackActionsOptions) { const response = await apiFetch(`/stacks/${stackName}/${endpoint}`, { method: 'POST' }); if (!response.ok) { const errText = await response.text(); - return { ok: false as const, errorMessage: errText || `${action} failed` }; + const actionError = parseStackActionError(errText, `${action} failed`); + return { + ok: false as const, + errorMessage: actionError.message, + rolledBack: actionError.rolledBack, + }; } toast.success(successMessage); if (action === 'update') stackListState.fetchImageUpdates(); @@ -695,7 +728,7 @@ export function useStackActions(options: UseStackActionsOptions) { const response = await apiFetch(`/stacks/${stackName}/${endpoint}`, { method: 'POST' }); if (!response.ok) { const errText = await response.text(); - throw new Error(errText || `${action} failed`); + throw parseStackActionError(errText, `${action} failed`); } toast.success(`Stack ${action}ed successfully!`); if (stackListState.selectedFile === stackFile) { @@ -714,9 +747,10 @@ export function useStackActions(options: UseStackActionsOptions) { } } catch (error) { console.error(`Failed to ${action}:`, error); - const msg = (error as Error).message || `Failed to ${action} stack`; + const actionError = error as StackActionError; + const msg = actionError.message || `Failed to ${action} stack`; toast.error( - action === 'deploy' && isPaid + action === 'deploy' && isPaid && actionError.rolledBack === true ? `${msg} - automatically rolled back to previous version.` : msg, );