diff --git a/backend/src/__tests__/mesh-service.test.ts b/backend/src/__tests__/mesh-service.test.ts index 78fb78dc..3cf012fa 100644 --- a/backend/src/__tests__/mesh-service.test.ts +++ b/backend/src/__tests__/mesh-service.test.ts @@ -1,6 +1,9 @@ import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'; +import { EventEmitter } from 'events'; +import fsSync from 'fs'; +import path from 'path'; import { setupTestDb, cleanupTestDb } from './helpers/setupTestDb'; -import { getSenchoIpFromSubnet, MeshError } from '../services/MeshService'; +import { getSenchoIpFromSubnet, MeshError, type MeshTarget, type MeshTcpStreamLike } from '../services/MeshService'; let tmpDir: string; let MeshService: typeof import('../services/MeshService').MeshService; @@ -501,3 +504,267 @@ describe('MeshService.regenerateAllOverrides (F6: boot-time regen)', () => { expect(concurrentCalls.length).toBeGreaterThanOrEqual(1); }); }); + +describe('MeshService.getDeclaredStackServiceNames (BUG-1)', () => { + function writeStackFile(stack: string, contents: string): void { + const composeDir = process.env.COMPOSE_DIR as string; + const dir = path.join(composeDir, stack); + fsSync.mkdirSync(dir, { recursive: true }); + fsSync.writeFileSync(path.join(dir, 'compose.yaml'), contents, 'utf8'); + } + + it('returns the keys of the compose services map', async () => { + const svc = MeshService.getInstance(); + writeStackFile('declared-stack', [ + 'services:', + ' echo:', + ' image: busybox:latest', + ' expose: ["9000"]', + ' prober:', + ' image: busybox:latest', + ].join('\n')); + + const names = await svc.getDeclaredStackServiceNames('declared-stack'); + expect(names.sort()).toEqual(['echo', 'prober']); + }); + + it('returns [] when the compose file is missing', async () => { + const svc = MeshService.getInstance(); + const names = await svc.getDeclaredStackServiceNames('does-not-exist'); + expect(names).toEqual([]); + }); + + it('returns [] for an invalid stack name (path traversal attempt)', async () => { + const svc = MeshService.getInstance(); + const names = await svc.getDeclaredStackServiceNames('../etc/passwd'); + expect(names).toEqual([]); + }); + + it('returns [] when YAML has no services key', async () => { + const svc = MeshService.getInstance(); + writeStackFile('no-services', 'version: "3.9"\n'); + const names = await svc.getDeclaredStackServiceNames('no-services'); + expect(names).toEqual([]); + }); +}); + +describe('MeshService.ensureStackOverride (BUG-1 fix)', () => { + function writeStackFile(stack: string, contents: string): void { + const composeDir = process.env.COMPOSE_DIR as string; + const dir = path.join(composeDir, stack); + fsSync.mkdirSync(dir, { recursive: true }); + fsSync.writeFileSync(path.join(dir, 'compose.yaml'), contents, 'utf8'); + } + + it('writes a non-empty services map sourced from the compose file (BUG-1 case)', async () => { + const svc = MeshService.getInstance(); + const db = DatabaseService.getInstance(); + const localNodeId = db.getNodes()[0].id; + + writeStackFile('audit-mesh-prod', [ + 'services:', + ' echo:', + ' image: busybox:latest', + ' prober:', + ' image: busybox:latest', + ].join('\n')); + + // The override generator must derive services from the compose + // file, not from runtime containers. Spy on + // inspectLocalStackServices to assert it is NOT consulted by the + // override-write path (regression guard). + const inspectSpy = vi.spyOn(svc, 'inspectLocalStackServices').mockResolvedValue([]); + db.insertMeshStack(localNodeId, 'audit-mesh-prod', 'tester'); + + const overridePath = await svc.ensureStackOverride(localNodeId, 'audit-mesh-prod'); + expect(overridePath).not.toBeNull(); + const yaml = fsSync.readFileSync(overridePath as string, 'utf8'); + expect(yaml).toContain('echo:'); + expect(yaml).toContain('prober:'); + expect(yaml).toContain('sencho_mesh'); + expect(yaml).not.toMatch(/^services:\s*\{\}/m); + expect(inspectSpy).not.toHaveBeenCalled(); + }); + + it('preserves an existing non-empty override when declared services come back empty', async () => { + const svc = MeshService.getInstance(); + const db = DatabaseService.getInstance(); + const localNodeId = db.getNodes()[0].id; + + // No compose file on disk; getDeclaredStackServiceNames returns []. + // Pre-seed an existing override file with a populated services map. + const dataDir = process.env.DATA_DIR as string; + const overrideDir = path.join(dataDir, 'mesh', 'overrides', String(localNodeId)); + fsSync.mkdirSync(overrideDir, { recursive: true }); + const overrideFile = path.join(overrideDir, 'orphaned-stack.override.yml'); + fsSync.writeFileSync(overrideFile, [ + 'services:', + ' webapp:', + ' networks:', + ' - sencho_mesh', + 'networks:', + ' sencho_mesh:', + ' external: true', + ].join('\n'), 'utf8'); + const originalContent = fsSync.readFileSync(overrideFile, 'utf8'); + + db.insertMeshStack(localNodeId, 'orphaned-stack', 'tester'); + + const result = await svc.ensureStackOverride(localNodeId, 'orphaned-stack'); + expect(result).toBe(overrideFile); + + const after = fsSync.readFileSync(overrideFile, 'utf8'); + expect(after).toBe(originalContent); + + const activity = svc.getActivity({ limit: 100 }); + expect(activity.some((e) => + e.type === 'mesh.override.preserved' + && /orphaned-stack/.test(e.message), + )).toBe(true); + }); +}); + +describe('MeshService tunnel-up regen (BUG-2)', () => { + it('triggers regenerateOverridesForNode for the firing nodeId once tunnel comes up', async () => { + const svc = MeshService.getInstance(); + const ptm = (await import('../services/PilotTunnelManager')).PilotTunnelManager.getInstance() as unknown as EventEmitter; + + const regenSpy = vi.spyOn( + svc as unknown as { regenerateOverridesForNode: (n: number) => Promise }, + 'regenerateOverridesForNode', + ).mockResolvedValue(undefined); + + // Drive start() while stubbing the heavy bits. setupMeshNetwork + // is what actually creates the bridge network; refreshAliasCache + // and syncForwarderListeners would touch real Dockerode and net + // state. regenerateAllOverrides we stub so the test only + // exercises the tunnel-up listener. + const internals = svc as unknown as { + started: boolean; + setupMeshNetwork: () => Promise; + refreshAliasCache: () => Promise; + syncForwarderListeners: () => Promise; + regenerateAllOverrides: () => Promise; + aliasRefreshTimer: NodeJS.Timeout | undefined; + }; + internals.started = false; + const origNetwork = internals.setupMeshNetwork.bind(svc); + const origRefresh = internals.refreshAliasCache.bind(svc); + const origSync = internals.syncForwarderListeners.bind(svc); + const origRegenAll = internals.regenerateAllOverrides.bind(svc); + internals.setupMeshNetwork = vi.fn().mockResolvedValue(undefined); + internals.refreshAliasCache = vi.fn().mockResolvedValue(undefined); + internals.syncForwarderListeners = vi.fn().mockResolvedValue(undefined); + internals.regenerateAllOverrides = vi.fn().mockResolvedValue({ regenerated: 0, failures: [], skipped: false }); + + try { + await svc.start(); + ptm.emit('tunnel-up', 14); + // Give the void-promise chain a tick to invoke the spy. + await new Promise((r) => setImmediate(r)); + + expect(regenSpy).toHaveBeenCalledWith(14); + + const activity = svc.getActivity({ limit: 50 }); + expect(activity.some((e) => + e.type === 'tunnel.open' && e.nodeId === 14, + )).toBe(true); + } finally { + internals.setupMeshNetwork = origNetwork; + internals.refreshAliasCache = origRefresh; + internals.syncForwarderListeners = origSync; + internals.regenerateAllOverrides = origRegenAll; + internals.started = false; + if (internals.aliasRefreshTimer) { + clearInterval(internals.aliasRefreshTimer); + internals.aliasRefreshTimer = undefined; + } + // Drain any tunnel-up listeners we registered during start(). + ptm.removeAllListeners('tunnel-up'); + ptm.removeAllListeners('tunnel-down'); + } + }); +}); + +describe('MeshService.openCrossNode (BUG-4)', () => { + function makeFakeStream(streamId: number): MeshTcpStreamLike & EventEmitter { + const ee = new EventEmitter() as MeshTcpStreamLike & EventEmitter & { destroyed: boolean }; + ee.destroyed = false; + Object.defineProperty(ee, 'streamId', { value: streamId, writable: false }); + ee.write = vi.fn().mockReturnValue(true); + ee.end = vi.fn(); + ee.destroy = vi.fn(() => { ee.destroyed = true; }); + return ee; + } + + function makeFakeSocket(): { destroy: ReturnType; end: ReturnType; on: ReturnType; write: ReturnType } { + return { + destroy: vi.fn(), + end: vi.fn(), + on: vi.fn(), + write: vi.fn(), + }; + } + + it('emits route.dispatch immediately on cross-node entry', async () => { + const svc = MeshService.getInstance(); + const target: MeshTarget = { + nodeId: 14, stack: 'audit-mesh-pilot', service: 'echo', + port: 9001, alias: 'echo.audit-mesh-pilot.sencho-pilot-test.sencho', + }; + const fakeStream = makeFakeStream(42); + vi.spyOn( + svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null }, + 'dialMeshTcpStream', + ).mockReturnValue(fakeStream); + + const fakeSrc = makeFakeSocket(); + // openCrossNode is private; cast to call it. + (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => void }) + .openCrossNode(target, fakeSrc); + + const dispatch = svc.getActivity({ limit: 50 }).find((e) => e.type === 'route.dispatch'); + expect(dispatch).toBeDefined(); + expect(dispatch?.alias).toBe(target.alias); + expect(dispatch?.nodeId).toBe(14); + // Clean up the open-timer so the test process exits cleanly. + fakeStream.emit('close'); + }); + + it('emits tunnel.fail after PROBE_TIMEOUT_MS when tcp_open_ack never arrives', async () => { + vi.useFakeTimers(); + try { + const svc = MeshService.getInstance(); + const target: MeshTarget = { + nodeId: 14, stack: 'audit-mesh-pilot', service: 'echo', + port: 9001, alias: 'echo.audit-mesh-pilot.sencho-pilot-test.sencho', + }; + const fakeStream = makeFakeStream(43); + vi.spyOn( + svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null }, + 'dialMeshTcpStream', + ).mockReturnValue(fakeStream); + + const fakeSrc = makeFakeSocket(); + (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => void }) + .openCrossNode(target, fakeSrc); + + // Before the timeout, only route.dispatch should be present. + expect(svc.getActivity({ limit: 50 }).some((e) => e.type === 'route.resolve.ok')).toBe(false); + + await vi.advanceTimersByTimeAsync(5_000); + + const events = svc.getActivity({ limit: 50 }); + const timeoutEvent = events.find((e) => + e.type === 'tunnel.fail' && /timed out waiting for tcp_open_ack/.test(e.message), + ); + expect(timeoutEvent).toBeDefined(); + expect(timeoutEvent?.nodeId).toBe(14); + expect(timeoutEvent?.alias).toBe(target.alias); + expect(fakeStream.destroy).toHaveBeenCalled(); + expect(fakeSrc.destroy).toHaveBeenCalled(); + } finally { + vi.useRealTimers(); + } + }); +}); diff --git a/backend/src/services/MeshService.ts b/backend/src/services/MeshService.ts index bd649612..e100b215 100644 --- a/backend/src/services/MeshService.ts +++ b/backend/src/services/MeshService.ts @@ -2,9 +2,11 @@ import net from 'net'; import path from 'path'; import fs from 'fs/promises'; import { EventEmitter } from 'events'; +import * as YAML from 'yaml'; import { ComposeService } from './ComposeService'; import { DatabaseService } from './DatabaseService'; import DockerController from './DockerController'; +import { FileSystemService } from './FileSystemService'; import { LicenseService } from './LicenseService'; import { PROXY_TIER_HEADER, PROXY_VARIANT_HEADER } from './license-headers'; import { MeshForwarder, type MeshForwarderHost } from './MeshForwarder'; @@ -51,10 +53,11 @@ export function getSenchoIpFromSubnet(subnet: string): string { export type MeshActivitySource = 'pilot' | 'mesh'; export type MeshActivityLevel = 'info' | 'warn' | 'error'; export type MeshActivityType = - | 'route.resolve.ok' | 'route.resolve.denied' + | 'route.dispatch' | 'route.resolve.ok' | 'route.resolve.denied' | 'tunnel.open' | 'tunnel.fail' | 'tunnel.backpressure' | 'opt_in' | 'opt_out' | 'mesh.enable' | 'mesh.disable' + | 'mesh.override.preserved' | 'probe.ok' | 'probe.fail' | 'forwarder.listen' | 'forwarder.unlisten' | 'forwarder.error'; @@ -199,10 +202,25 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { const ptm = PilotTunnelManager.getInstance(); ptm.on('tunnel-down', (nodeId: number) => this.onTunnelDown(nodeId)); - ptm.on('tunnel-up', (nodeId: number) => this.logActivity({ - source: 'pilot', level: 'info', type: 'tunnel.open', - nodeId, message: `pilot tunnel up for node ${nodeId}`, - })); + ptm.on('tunnel-up', (nodeId: number) => { + this.logActivity({ + source: 'pilot', level: 'info', type: 'tunnel.open', + nodeId, message: `pilot tunnel up for node ${nodeId}`, + }); + // Boot regen runs before any pilot tunnel comes up, so any + // pilot-mode node misses its initial override push. Now that + // the tunnel is live, retry the regen for this node so its + // overrides on disk match what central holds. Idempotent: + // pushOverrideToNode writes the same file every time, and a + // tunnel reconnect during runtime regenerates harmlessly. + void this.regenerateOverridesForNode(nodeId).catch((err) => { + this.logActivity({ + source: 'mesh', level: 'warn', type: 'forwarder.error', + nodeId, + message: `tunnel-up regen failed for node ${nodeId}: ${sanitizeForLog((err as Error).message)}`, + }); + }); + }); await this.setupMeshNetwork(); try { @@ -578,17 +596,40 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { if (!this.senchoIp) return null; const aliases: MeshAlias[] = Array.from(this.aliasCache.values()).map((a) => ({ host: a.host })); - const services = await this.inspectStackServices(nodeId, stackName); - const yaml = generateOverrideYaml({ - services: services.map((s) => s.service), - aliases, - senchoIp: this.senchoIp, - }); + const serviceNames = await this.getDeclaredStackServiceNames(stackName, nodeId); const dir = this.overrideDirFor(nodeId); await fs.mkdir(dir, { recursive: true }); - const file = path.resolve(dir, `${stackName}.override.yml`); + // path.basename mirrors the applyLocalOverride pattern (and is + // the form CodeQL's path-injection model recognizes). + const file = path.resolve(dir, `${path.basename(stackName)}.override.yml`); if (!isPathWithinBase(file, dir)) return null; + + // Defensive fallback: a deploy that runs `compose down` immediately + // before `compose up` removes the containers, but the compose file + // is still on disk so getDeclaredStackServiceNames returns the + // declared services. The fallback below covers the much narrower + // case where the compose file itself is unreadable (permission + // glitch, transient FS error, mid-write rename); in that case + // keep any existing override rather than overwrite with `services: {}`. + if (serviceNames.length === 0) { + const existing = await this.readExistingOverrideServiceNames(dir, stackName); + if (existing.length > 0) { + this.logActivity({ + source: 'mesh', level: 'warn', type: 'mesh.override.preserved', + nodeId, + message: `mesh override preserved for ${stackName}: declared services unreadable, keeping ${existing.length} existing entries`, + details: { stackName, preservedServices: existing }, + }); + return file; + } + } + + const yaml = generateOverrideYaml({ + services: serviceNames, + aliases, + senchoIp: this.senchoIp, + }); await fs.writeFile(file, yaml, 'utf8'); return file; } @@ -610,14 +651,9 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { this.networkSetupError || 'mesh data plane unavailable on this node', ); } - const services = await this.inspectLocalStackServices(stackName); - const yaml = generateOverrideYaml({ - services: services.map((s) => s.service), - aliases, - senchoIp: this.senchoIp, - }); - const localNodeId = NodeRegistry.getInstance().getDefaultNodeId(); + const serviceNames = await this.getDeclaredStackServiceNames(stackName, localNodeId); + const dir = this.overrideDirFor(localNodeId); await fs.mkdir(dir, { recursive: true }); // path.basename strips any directory component as defense-in-depth @@ -625,6 +661,30 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { // CodeQL's path-injection model. const file = path.resolve(dir, `${path.basename(stackName)}.override.yml`); if (!isPathWithinBase(file, dir)) return null; + + // Defensive fallback (mirror of ensureStackOverride): keep any + // existing override when the compose file is transiently + // unreadable. The remote that pushed this update will retry on + // its next regen tick, so a one-shot read failure should not + // wipe out a working override. + if (serviceNames.length === 0) { + const existing = await this.readExistingOverrideServiceNames(dir, stackName); + if (existing.length > 0) { + this.logActivity({ + source: 'mesh', level: 'warn', type: 'mesh.override.preserved', + nodeId: localNodeId, + message: `mesh override preserved for ${stackName}: declared services unreadable, keeping ${existing.length} existing entries`, + details: { stackName, preservedServices: existing }, + }); + return file; + } + } + + const yaml = generateOverrideYaml({ + services: serviceNames, + aliases, + senchoIp: this.senchoIp, + }); await fs.writeFile(file, yaml, 'utf8'); return file; } @@ -769,6 +829,73 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { return Array.from(this.aliasCache.values()); } + /** + * Read the local stack's compose file and return its declared service + * names. Used by the override-write paths so a deploy that has just + * torn containers down still emits a complete services map even + * though Dockerode briefly returns no containers. Independent of + * runtime container state, so the override stays correct across the + * deploy lifecycle. + * + * Returns [] when the compose file is missing, unreadable, or fails + * to parse. Combined with the defensive fallback in + * {@link ensureStackOverride} / {@link applyLocalOverride} a transient + * empty result does not clobber a known-good override. + * + * LIMITATION: stacks that pull services in via compose `extends:` or + * `include:` will not have those external services covered. The + * top-level YAML.parse is sufficient for the supported compose + * shapes; if extends/include usage emerges, swap to + * `docker compose config --services` (subprocess). + */ + public async getDeclaredStackServiceNames(stackName: string, nodeId?: number): Promise { + if (!isValidStackName(stackName)) return []; + const targetNodeId = nodeId ?? NodeRegistry.getInstance().getDefaultNodeId(); + try { + const fsSvc = FileSystemService.getInstance(targetNodeId); + const filename = await fsSvc.getComposeFilename(stackName); + const baseDir = fsSvc.getBaseDir(); + // path.basename strips any directory component as defense-in-depth + // on top of isValidStackName + isPathWithinBase. Recognized by + // CodeQL's path-injection model. + const composePath = path.join(baseDir, path.basename(stackName), filename); + if (!isPathWithinBase(composePath, baseDir)) return []; + const content = await fs.readFile(composePath, 'utf8'); + const parsed = YAML.parse(content) as { services?: Record } | null; + const services = parsed?.services && typeof parsed.services === 'object' ? parsed.services : null; + if (!services) return []; + return Object.keys(services).filter((name) => /^[A-Za-z0-9_][A-Za-z0-9_.-]*$/.test(name)); + } catch (err) { + console.warn( + '[MeshService] getDeclaredStackServiceNames failed:', + sanitizeForLog((err as Error).message), + ); + return []; + } + } + + /** + * Parse an existing mesh override file and extract the service names + * it already lists. Used by the defensive fallback so a transient + * empty compose-file read does not clobber a known-good override. + * Takes (dir, stackName) rather than a pre-built filePath so the + * path-injection sanitizer pattern (path.basename + isPathWithinBase) + * lives next to the read sink and stays recognizable to CodeQL. + */ + private async readExistingOverrideServiceNames(dir: string, stackName: string): Promise { + const file = path.resolve(dir, `${path.basename(stackName)}.override.yml`); + if (!isPathWithinBase(file, dir)) return []; + try { + const content = await fs.readFile(file, 'utf8'); + const parsed = YAML.parse(content) as { services?: Record } | null; + const services = parsed?.services && typeof parsed.services === 'object' ? parsed.services : null; + if (!services) return []; + return Object.keys(services); + } catch { + return []; + } + } + /** * Inspect a stack and return its running services with the ports they * listen on. For the LOCAL Docker daemon only; callers targeting a @@ -1181,6 +1308,17 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { } private openCrossNode(target: MeshTarget, src: net.Socket): void { + // Log every dispatch entry. Same-node logs route.resolve.ok on + // its TCP `connect` event; cross-node only logs route.resolve.ok + // once tcp_open_ack arrives from the agent. Without this entry + // log, a stuck cross-node dial leaves zero events in the activity + // buffer even though the prober's TCP handshake completed. + this.logActivity({ + source: 'mesh', level: 'info', type: 'route.dispatch', + nodeId: target.nodeId, alias: target.alias, + message: `cross-node dispatch to ${target.alias} on node ${target.nodeId}`, + }); + const tcpStream = this.dialMeshTcpStream(target); if (!tcpStream) { this.logActivity({ @@ -1195,7 +1333,27 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { } const record = this.registerActiveStream(target.alias, tcpStream.streamId); const t0 = Date.now(); + + // Timer guards against the agent never returning a tcp_open_ack + // (broken pilot, frame dropped, dial stuck after handshake). + // Without this the stream sits forever and the operator sees + // nothing in the activity log between dispatch and close. + let openTimer: NodeJS.Timeout | null = setTimeout(() => { + openTimer = null; + this.logActivity({ + source: 'pilot', level: 'warn', type: 'tunnel.fail', + nodeId: target.nodeId, alias: target.alias, streamId: record.streamId, + message: `cross-node dial to ${target.alias} timed out waiting for tcp_open_ack`, + }); + try { tcpStream.destroy(); } catch { /* ignore */ } + try { src.destroy(); } catch { /* ignore */ } + }, PROBE_TIMEOUT_MS); + const clearOpenTimer = () => { + if (openTimer) { clearTimeout(openTimer); openTimer = null; } + }; + tcpStream.on('open', () => { + clearOpenTimer(); this.logActivity({ source: 'mesh', level: 'info', type: 'route.resolve.ok', nodeId: target.nodeId, alias: target.alias, streamId: record.streamId, @@ -1208,6 +1366,7 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { try { src.write(chunk); } catch { /* ignore */ } }); tcpStream.on('error', (err: Error) => { + clearOpenTimer(); this.logActivity({ source: 'pilot', level: 'error', type: 'tunnel.fail', nodeId: target.nodeId, alias: target.alias, streamId: record.streamId, @@ -1217,6 +1376,7 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { try { src.destroy(); } catch { /* ignore */ } }); tcpStream.on('close', () => { + clearOpenTimer(); this.activeStreams.delete(record.streamId); try { src.end(); } catch { /* ignore */ } }); @@ -1225,8 +1385,8 @@ export class MeshService extends EventEmitter implements MeshForwarderHost { tcpStream.write(chunk); }); src.on('end', () => tcpStream.end()); - src.on('close', () => tcpStream.destroy()); - src.on('error', () => tcpStream.destroy()); + src.on('close', () => { clearOpenTimer(); tcpStream.destroy(); }); + src.on('error', () => { clearOpenTimer(); tcpStream.destroy(); }); } private registerActiveStream(alias: string, streamId?: number): ActiveStreamRecord {