Files
sencho/backend/src/__tests__/mesh-service.test.ts
T

1945 lines
91 KiB
TypeScript

import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest';
import { EventEmitter } from 'events';
import fsSync from 'fs';
import fs from 'fs/promises';
import path from 'path';
import { setupTestDb, cleanupTestDb } from './helpers/setupTestDb';
import { getSenchoIpFromSubnet, MeshError, type MeshTarget, type MeshTcpStreamLike } from '../services/MeshService';
let tmpDir: string;
let MeshService: typeof import('../services/MeshService').MeshService;
let DatabaseService: typeof import('../services/DatabaseService').DatabaseService;
beforeAll(async () => {
tmpDir = await setupTestDb();
({ MeshService } = await import('../services/MeshService'));
({ DatabaseService } = await import('../services/DatabaseService'));
});
afterAll(() => {
cleanupTestDb(tmpDir);
});
beforeEach(() => {
process.env.SENCHO_MODE = 'server';
const db = DatabaseService.getInstance().getDb();
db.prepare('DELETE FROM mesh_stacks').run();
db.prepare('DELETE FROM nodes WHERE is_default = 0').run();
const svc = MeshService.getInstance() as unknown as {
aliasCache: Map<string, unknown>;
aliasByPort: Map<number, unknown>;
activity: unknown[];
activeStreams: Map<number, unknown>;
routeErrorMap: Map<string, unknown>;
routeLatencyMap: Map<string, unknown>;
senchoIp: string | null;
meshSubnet: string;
networkSetupError: string | null;
selfCentralNodeId: number | null;
proxyTunnelSelfCentralNodeId: number | null;
};
svc.aliasCache = new Map();
svc.aliasByPort = new Map();
svc.activity = [];
svc.activeStreams = new Map();
svc.routeErrorMap = new Map();
svc.routeLatencyMap = new Map();
svc.senchoIp = '172.30.0.2';
svc.meshSubnet = '172.30.0.0/24';
svc.networkSetupError = null;
svc.selfCentralNodeId = null;
svc.proxyTunnelSelfCentralNodeId = null;
delete process.env.SENCHO_ENROLL_TOKEN;
vi.restoreAllMocks();
});
describe('MeshService.optInStack', () => {
it('writes a mesh_stacks row and rejects duplicate ports', async () => {
const svc = MeshService.getInstance();
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'db', ports: [5432] }]);
vi.spyOn(svc as unknown as { regenerateOverridesAcrossFleet: (n?: number, s?: string) => Promise<void> }, 'regenerateOverridesAcrossFleet')
.mockResolvedValue(undefined);
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
await svc.optInStack(localNodeId, 'api', 'tester');
expect(db.isMeshStackEnabled(localNodeId, 'api')).toBe(true);
await expect(svc.optInStack(localNodeId, 'shadow', 'tester'))
.rejects.toThrow(/port 5432 is already claimed/);
});
it('rejects opt-in when every service declared empty ports (no routable target)', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([
{ service: 'web', ports: [] },
{ service: 'db', ports: [] },
]);
vi.spyOn(svc as unknown as { regenerateOverridesAcrossFleet: (n?: number, s?: string) => Promise<void> }, 'regenerateOverridesAcrossFleet')
.mockResolvedValue(undefined);
await expect(svc.optInStack(localNodeId, 'silent', 'tester'))
.rejects.toThrow(/no service ports to mesh/);
// Pre-fix this wrote a row with no aliases and the stack appeared
// "online but doesn't route". Verify the row was not written.
expect(db.isMeshStackEnabled(localNodeId, 'silent')).toBe(false);
});
it('opt-out removes the row and the override', async () => {
const svc = MeshService.getInstance();
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'db', ports: [5432] }]);
vi.spyOn(svc as unknown as { regenerateOverridesAcrossFleet: (n?: number, s?: string) => Promise<void> }, 'regenerateOverridesAcrossFleet')
.mockResolvedValue(undefined);
vi.spyOn(svc as unknown as { removeStackOverride: (n: number, s: string) => Promise<void> }, 'removeStackOverride')
.mockResolvedValue(undefined);
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
await svc.optInStack(localNodeId, 'api', 'tester');
await svc.optOutStack(localNodeId, 'api', 'tester');
expect(db.isMeshStackEnabled(localNodeId, 'api')).toBe(false);
});
it('restores opt-in authority when target override removal is rejected', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
db.insertMeshStack(localNodeId, 'busy-stack', 'setup');
vi.spyOn(svc, 'removeOverrideFromNode').mockRejectedValue(
new Error('HTTP 500: another operation is already in progress'),
);
await expect(svc.optOutStack(localNodeId, 'busy-stack', 'tester'))
.rejects.toThrow('another operation is already in progress');
expect(db.isMeshStackEnabled(localNodeId, 'busy-stack')).toBe(true);
});
it('serializes concurrent opt-out requests for the same node', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
db.insertMeshStack(localNodeId, 'queued-stack', 'setup');
let releaseRemoval!: () => void;
const removalPending = new Promise<void>((resolve) => {
releaseRemoval = resolve;
});
const removeSpy = vi.spyOn(svc, 'removeOverrideFromNode').mockReturnValue(removalPending);
vi.spyOn(svc as unknown as { regenerateOverridesAcrossFleet: () => Promise<void> }, 'regenerateOverridesAcrossFleet')
.mockResolvedValue(undefined);
vi.spyOn(svc as unknown as { cascadeRecomposeAcrossFleet: () => void }, 'cascadeRecomposeAcrossFleet')
.mockImplementation(() => { /* noop */ });
vi.spyOn(svc, 'triggerRedeploy').mockImplementation(() => { /* noop */ });
const first = svc.optOutStack(localNodeId, 'queued-stack', 'tester');
await vi.waitFor(() => expect(removeSpy).toHaveBeenCalledTimes(1));
const second = svc.optOutStack(localNodeId, 'queued-stack', 'tester');
await Promise.resolve();
expect(removeSpy).toHaveBeenCalledTimes(1);
releaseRemoval();
await Promise.all([first, second]);
expect(removeSpy).toHaveBeenCalledTimes(1);
expect(db.isMeshStackEnabled(localNodeId, 'queued-stack')).toBe(false);
});
it('serializes different Mesh mutations per node without blocking another node', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'independent-node', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.setNodeMeshEnabled(localNodeId, true);
db.insertMeshStack(localNodeId, 'held-stack', 'setup');
let releaseRemoval!: () => void;
const removalPending = new Promise<void>((resolve) => {
releaseRemoval = resolve;
});
vi.spyOn(svc, 'removeOverrideFromNode').mockReturnValue(removalPending);
vi.spyOn(svc as unknown as { regenerateOverridesAcrossFleet: () => Promise<void> }, 'regenerateOverridesAcrossFleet')
.mockResolvedValue(undefined);
vi.spyOn(svc as unknown as { cascadeRecomposeAcrossFleet: () => void }, 'cascadeRecomposeAcrossFleet')
.mockImplementation(() => { /* noop */ });
vi.spyOn(svc, 'triggerRedeploy').mockImplementation(() => { /* noop */ });
const optOut = svc.optOutStack(localNodeId, 'held-stack', 'tester');
await vi.waitFor(() => expect(svc.removeOverrideFromNode).toHaveBeenCalledTimes(1));
const disable = svc.disableForNode(localNodeId, 'tester');
await svc.enableForNode(remoteNodeId);
expect(db.getNodeMeshEnabled(localNodeId)).toBe(true);
expect(db.getNodeMeshEnabled(remoteNodeId)).toBe(true);
releaseRemoval();
await Promise.all([optOut, disable]);
expect(db.getNodeMeshEnabled(localNodeId)).toBe(false);
db.deleteNode(remoteNodeId);
});
it('rejects an invalid stack name (path traversal attempt)', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
await expect(svc.optInStack(localNodeId, '../../etc/passwd', 'tester'))
.rejects.toThrow(/invalid stack name/);
expect(db.isMeshStackEnabled(localNodeId, '../../etc/passwd')).toBe(false);
});
it('forwarder binds every alias port across the fleet, not just local-owned ports', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
// Seed local- and remote-owned aliases directly into
// aliasByPort so we exercise syncForwarderListeners without a
// live remote Sencho. The model: every meshed node binds every
// alias port because meshed containers' extra_hosts:host-gateway
// entries land on the SOURCE node's gateway, so the source node
// is where the inbound TCP connection is intercepted.
const aliasByPort = (svc as unknown as { aliasByPort: Map<number, unknown> }).aliasByPort;
aliasByPort.set(9000, {
host: 'echo.local-stack.Local.sencho',
nodeId: localNodeId, nodeName: 'Local',
stackName: 'local-stack', serviceName: 'echo', port: 9000,
});
aliasByPort.set(9001, {
host: 'echo.remote-stack.remote-pilot.sencho',
nodeId: remoteNodeId, nodeName: 'remote-pilot',
stackName: 'remote-stack', serviceName: 'echo', port: 9001,
});
// Stub the forwarder so the test never touches a real net.Server.
const listened: number[] = [];
const fwd = (svc as unknown as {
forwarder: { listen: (p: number) => Promise<void>; unlisten: (p: number) => Promise<void>; getListenerPorts: () => number[] };
}).forwarder;
const realListen = fwd.listen.bind(fwd);
const realUnlisten = fwd.unlisten.bind(fwd);
const realGetListenerPorts = fwd.getListenerPorts.bind(fwd);
fwd.listen = async (p: number) => { listened.push(p); };
fwd.unlisten = async () => { /* no-op */ };
fwd.getListenerPorts = () => [...listened];
try {
await (svc as unknown as { syncForwarderListeners: () => Promise<void> }).syncForwarderListeners();
expect(listened.sort()).toEqual([9000, 9001]);
} finally {
fwd.listen = realListen;
fwd.unlisten = realUnlisten;
fwd.getListenerPorts = realGetListenerPorts;
aliasByPort.clear();
db.deleteNode(remoteNodeId);
}
});
it('optInStack cascades override regen to every meshed node, not just the source node', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.insertMeshStack(remoteNodeId, 'remote-existing', 'setup');
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [8080] }]);
const pushed: Array<{ nodeId: number; stackName: string }> = [];
vi.spyOn(svc, 'pushOverrideToNode').mockImplementation(async (nodeId, stackName) => {
pushed.push({ nodeId, stackName });
});
vi.spyOn(svc as unknown as { triggerRedeploy: (n: number, s: string, a: string) => void }, 'triggerRedeploy')
.mockImplementation(() => { /* noop */ });
await svc.optInStack(localNodeId, 'new-stack', 'tester');
expect(pushed).toContainEqual({ nodeId: localNodeId, stackName: 'new-stack' });
// Other meshed nodes' existing stacks must regenerate so they pick
// up the new alias entry.
expect(pushed).toContainEqual({ nodeId: remoteNodeId, stackName: 'remote-existing' });
});
it('optOutStack cascades override regen to every other meshed node', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.insertMeshStack(localNodeId, 'to-remove', 'setup');
db.insertMeshStack(remoteNodeId, 'remote-existing', 'setup');
const pushed: Array<{ nodeId: number; stackName: string }> = [];
vi.spyOn(svc, 'pushOverrideToNode').mockImplementation(async (nodeId, stackName) => {
pushed.push({ nodeId, stackName });
});
vi.spyOn(svc as unknown as { removeOverrideFromNode: (n: number, s: string) => Promise<void> }, 'removeOverrideFromNode')
.mockResolvedValue(undefined);
vi.spyOn(svc as unknown as { triggerRedeploy: (n: number, s: string, a: string) => void }, 'triggerRedeploy')
.mockImplementation(() => { /* noop */ });
await svc.optOutStack(localNodeId, 'to-remove', 'tester');
// Other meshed nodes' existing stacks must regenerate so they drop
// the removed alias entry.
expect(pushed).toContainEqual({ nodeId: remoteNodeId, stackName: 'remote-existing' });
// The removed row was deleted before the cascade ran, and its file
// was unlinked separately, so it must not appear in the cascade.
expect(pushed.find((p) => p.stackName === 'to-remove')).toBeUndefined();
});
it('optInStack cascades recompose to every previously meshed stack across the fleet', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.insertMeshStack(localNodeId, 'local-existing', 'setup');
db.insertMeshStack(remoteNodeId, 'remote-existing', 'setup');
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [8080] }]);
vi.spyOn(svc, 'pushOverrideToNode').mockResolvedValue(undefined);
const redeployed: Array<{ nodeId: number; stackName: string }> = [];
vi.spyOn(svc, 'triggerRedeploy').mockImplementation((nodeId, stackName) => {
redeployed.push({ nodeId, stackName });
});
await svc.optInStack(localNodeId, 'new-stack', 'tester');
// The cascade must redeploy every previously meshed stack so its
// container's /etc/hosts picks up the newly-added alias.
expect(redeployed).toContainEqual({ nodeId: localNodeId, stackName: 'local-existing' });
expect(redeployed).toContainEqual({ nodeId: remoteNodeId, stackName: 'remote-existing' });
// The just-opted-in stack is still redeployed separately by the
// explicit triggerRedeploy at the end of optInStack.
expect(redeployed).toContainEqual({ nodeId: localNodeId, stackName: 'new-stack' });
});
it('optInStack does not double-redeploy the just-opted-in tuple via the cascade path', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.insertMeshStack(remoteNodeId, 'remote-existing', 'setup');
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [8080] }]);
vi.spyOn(svc, 'pushOverrideToNode').mockResolvedValue(undefined);
const redeployed: Array<{ nodeId: number; stackName: string }> = [];
vi.spyOn(svc, 'triggerRedeploy').mockImplementation((nodeId, stackName) => {
redeployed.push({ nodeId, stackName });
});
await svc.optInStack(localNodeId, 'new-stack', 'tester');
// new-stack appears exactly once: the explicit redeploy at the end
// of optInStack. The cascade walks db.listMeshStacks() (which now
// includes new-stack) but its skip tuple drops it.
const newStackHits = redeployed.filter(
(r) => r.nodeId === localNodeId && r.stackName === 'new-stack',
);
expect(newStackHits).toHaveLength(1);
});
it('optInStack cascade is a no-op when no other meshed stacks exist', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [8080] }]);
vi.spyOn(svc, 'pushOverrideToNode').mockResolvedValue(undefined);
const redeployed: Array<{ nodeId: number; stackName: string }> = [];
vi.spyOn(svc, 'triggerRedeploy').mockImplementation((nodeId, stackName) => {
redeployed.push({ nodeId, stackName });
});
await svc.optInStack(localNodeId, 'first-stack', 'tester');
// Only the just-opted-in stack is redeployed; the cascade has no
// peers to recompose so the early return short-circuits before
// the summary activity entry fires.
expect(redeployed).toEqual([{ nodeId: localNodeId, stackName: 'first-stack' }]);
const cascadeEntries = svc.getActivity({ limit: 1000 })
.filter((e) => e.message.startsWith('mesh cascade recompose'));
expect(cascadeEntries).toHaveLength(0);
});
it('optOutStack cascades recompose to every other meshed stack', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.insertMeshStack(localNodeId, 'to-remove', 'setup');
db.insertMeshStack(localNodeId, 'local-survivor', 'setup');
db.insertMeshStack(remoteNodeId, 'remote-survivor', 'setup');
vi.spyOn(svc, 'pushOverrideToNode').mockResolvedValue(undefined);
vi.spyOn(svc as unknown as { removeOverrideFromNode: (n: number, s: string) => Promise<void> }, 'removeOverrideFromNode')
.mockResolvedValue(undefined);
const redeployed: Array<{ nodeId: number; stackName: string }> = [];
vi.spyOn(svc, 'triggerRedeploy').mockImplementation((nodeId, stackName) => {
redeployed.push({ nodeId, stackName });
});
await svc.optOutStack(localNodeId, 'to-remove', 'tester');
// Surviving meshed stacks must recompose so their /etc/hosts drops
// the now-removed alias.
expect(redeployed).toContainEqual({ nodeId: localNodeId, stackName: 'local-survivor' });
expect(redeployed).toContainEqual({ nodeId: remoteNodeId, stackName: 'remote-survivor' });
// The opted-out stack is still redeployed separately by the
// explicit triggerRedeploy at the end of optOutStack so its own
// container's /etc/hosts is cleared.
expect(redeployed).toContainEqual({ nodeId: localNodeId, stackName: 'to-remove' });
});
});
describe('MeshService.disableForNode', () => {
it('redeploys affected stacks and cascades across the fleet on node disable', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.setNodeMeshEnabled(localNodeId, true);
db.insertMeshStack(localNodeId, 'alpha', 'setup');
db.insertMeshStack(localNodeId, 'beta', 'setup');
db.insertMeshStack(remoteNodeId, 'gamma', 'setup');
const regenSpy = vi.spyOn(
svc as unknown as { regenerateOverridesAcrossFleet: (n?: number, s?: string) => Promise<void> },
'regenerateOverridesAcrossFleet',
).mockResolvedValue(undefined);
const cascadeSpy = vi.spyOn(
svc as unknown as { cascadeRecomposeAcrossFleet: (n: number | undefined, s: string | undefined, a: string) => void },
'cascadeRecomposeAcrossFleet',
).mockImplementation(() => { /* noop */ });
const triggerSpy = vi.spyOn(svc, 'triggerRedeploy').mockImplementation(() => { /* noop */ });
const removeSpy = vi.spyOn(svc, 'removeOverrideFromNode').mockResolvedValue(undefined);
await svc.disableForNode(localNodeId, 'tester');
expect(db.getNodeMeshEnabled(localNodeId)).toBe(false);
// Disabled-node rows are deleted before the cascade so listMeshStacks
// does not return them. The other node's row stays.
expect(db.listMeshStacks(localNodeId)).toEqual([]);
expect(db.listMeshStacks(remoteNodeId).map((s) => s.stack_name)).toEqual(['gamma']);
expect(regenSpy).toHaveBeenCalledTimes(1);
// Cascade walks every remaining mesh_stacks row, no skip tuple.
expect(cascadeSpy).toHaveBeenCalledWith(undefined, undefined, 'tester');
// Each disabled-node stack triggers a direct redeploy so its
// container detaches from sencho_mesh. Cascade is mocked, so the
// only triggerRedeploy calls come from the disable loop itself
// (twice, not three: gamma stays put on the remote node).
expect(triggerSpy).toHaveBeenCalledTimes(2);
expect(triggerSpy).toHaveBeenCalledWith(localNodeId, 'alpha', 'tester');
expect(triggerSpy).toHaveBeenCalledWith(localNodeId, 'beta', 'tester');
// disableForNode routes through removeOverrideFromNode (not the
// local-only removeStackOverride) so remote nodes also have their
// pushed override files deleted via DELETE /api/mesh/local-override.
expect(removeSpy).toHaveBeenCalledTimes(2);
expect(removeSpy).toHaveBeenCalledWith(localNodeId, 'alpha');
expect(removeSpy).toHaveBeenCalledWith(localNodeId, 'beta');
db.deleteNode(remoteNodeId);
});
it('dispatches removeOverrideFromNode for each stack when the disabled node is remote', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const remoteNodeId = db.addNode({
name: 'remote-pilot-disable', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.setNodeMeshEnabled(remoteNodeId, true);
db.insertMeshStack(remoteNodeId, 'delta', 'setup');
db.insertMeshStack(remoteNodeId, 'epsilon', 'setup');
vi.spyOn(
svc as unknown as { regenerateOverridesAcrossFleet: (n?: number, s?: string) => Promise<void> },
'regenerateOverridesAcrossFleet',
).mockResolvedValue(undefined);
vi.spyOn(
svc as unknown as { cascadeRecomposeAcrossFleet: (n: number | undefined, s: string | undefined, a: string) => void },
'cascadeRecomposeAcrossFleet',
).mockImplementation(() => { /* noop */ });
vi.spyOn(svc, 'triggerRedeploy').mockImplementation(() => { /* noop */ });
const removeSpy = vi.spyOn(svc, 'removeOverrideFromNode').mockResolvedValue(undefined);
await svc.disableForNode(remoteNodeId, 'tester');
expect(db.getNodeMeshEnabled(remoteNodeId)).toBe(false);
// Remote-node disable must dispatch the remote-aware helper so the
// override file pushed earlier via applyLocalOverride is deleted
// on the peer via DELETE /api/mesh/local-override/:stack.
expect(removeSpy).toHaveBeenCalledTimes(2);
expect(removeSpy).toHaveBeenCalledWith(remoteNodeId, 'delta');
expect(removeSpy).toHaveBeenCalledWith(remoteNodeId, 'epsilon');
db.deleteNode(remoteNodeId);
});
it('keeps failed remote stacks authoritative when node disable is incomplete', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const remoteNodeId = db.addNode({
name: 'remote-partial-disable', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.setNodeMeshEnabled(remoteNodeId, true);
db.insertMeshStack(remoteNodeId, 'removed-stack', 'setup');
db.insertMeshStack(remoteNodeId, 'busy-stack', 'setup');
vi.spyOn(svc, 'removeOverrideFromNode').mockImplementation(async (_nodeId, stackName) => {
if (stackName === 'busy-stack') throw new Error('target busy');
});
vi.spyOn(
svc as unknown as { regenerateOverridesAcrossFleet: () => Promise<void> },
'regenerateOverridesAcrossFleet',
).mockResolvedValue(undefined);
vi.spyOn(
svc as unknown as { cascadeRecomposeAcrossFleet: () => void },
'cascadeRecomposeAcrossFleet',
).mockImplementation(() => { /* noop */ });
const redeployed: string[] = [];
vi.spyOn(svc, 'triggerRedeploy').mockImplementation((_nodeId, stackName) => {
redeployed.push(stackName);
});
vi.spyOn(svc as unknown as { refreshAliasCache: () => Promise<void> }, 'refreshAliasCache')
.mockRejectedValue(new Error('refresh failed'));
vi.spyOn(console, 'warn').mockImplementation(() => { /* silence */ });
await expect(svc.disableForNode(remoteNodeId, 'tester')).rejects.toThrow('busy-stack');
expect(db.getNodeMeshEnabled(remoteNodeId)).toBe(true);
expect(db.isMeshStackEnabled(remoteNodeId, 'removed-stack')).toBe(false);
expect(db.isMeshStackEnabled(remoteNodeId, 'busy-stack')).toBe(true);
expect(redeployed).toContain('removed-stack');
db.deleteNode(remoteNodeId);
});
it('defaults the actor when none is supplied so legacy callers still log a non-empty actor', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
db.setNodeMeshEnabled(localNodeId, true);
db.insertMeshStack(localNodeId, 'solo', 'setup');
vi.spyOn(svc as unknown as { regenerateOverridesAcrossFleet: () => Promise<void> }, 'regenerateOverridesAcrossFleet')
.mockResolvedValue(undefined);
const cascadeSpy = vi.spyOn(
svc as unknown as { cascadeRecomposeAcrossFleet: (n: number | undefined, s: string | undefined, a: string) => void },
'cascadeRecomposeAcrossFleet',
).mockImplementation(() => { /* noop */ });
const triggerSpy = vi.spyOn(svc, 'triggerRedeploy').mockImplementation(() => { /* noop */ });
vi.spyOn(svc, 'removeOverrideFromNode').mockResolvedValue(undefined);
await svc.disableForNode(localNodeId);
expect(db.getNodeMeshEnabled(localNodeId)).toBe(false);
expect(cascadeSpy.mock.calls[0][2]).toBe('system:mesh.disable');
expect(triggerSpy.mock.calls[0][2]).toBe('system:mesh.disable');
});
});
describe('MeshService activity log', () => {
it('keeps the most recent events under the 1000-cap', () => {
const svc = MeshService.getInstance();
for (let i = 0; i < 1100; i++) {
svc.logActivity({ source: 'mesh', level: 'info', type: 'opt_in', message: `evt-${i}` });
}
const all = svc.getActivity({ limit: 2000 });
expect(all.length).toBe(1000);
expect(all[0].message).toBe('evt-100');
expect(all[all.length - 1].message).toBe('evt-1099');
});
it('filters by alias / source / level', () => {
const svc = MeshService.getInstance();
svc.logActivity({ source: 'mesh', level: 'info', type: 'opt_in', alias: 'a.b.c.sencho', message: 'a' });
svc.logActivity({ source: 'pilot', level: 'error', type: 'tunnel.fail', alias: 'a.b.c.sencho', message: 'b' });
svc.logActivity({ source: 'mesh', level: 'info', type: 'route.resolve.ok', alias: 'x.y.z.sencho', message: 'c' });
expect(svc.getActivity({ alias: 'a.b.c.sencho' }).length).toBe(2);
expect(svc.getActivity({ source: 'pilot' }).length).toBe(1);
expect(svc.getActivity({ level: 'error' }).length).toBe(1);
});
it('subscribeActivity fires for new events and unsubscribes cleanly', () => {
const svc = MeshService.getInstance();
const seen: string[] = [];
const unsubscribe = svc.subscribeActivity((e) => seen.push(e.message));
svc.logActivity({ source: 'mesh', level: 'info', type: 'mesh.enable', message: 'one' });
unsubscribe();
svc.logActivity({ source: 'mesh', level: 'info', type: 'mesh.disable', message: 'two' });
expect(seen).toEqual(['one']);
});
});
describe('MeshService.testUpstream tunnel-down path', () => {
it('returns ok:false where=pilot_tunnel when no tunnel is registered', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'opsix', type: 'remote', is_default: false,
compose_dir: '/tmp', api_url: 'https://opsix.example',
api_token: 'tok', mode: 'pilot_agent',
});
(svc as unknown as { aliasCache: Map<string, unknown> }).aliasCache = new Map([
['db.api.opsix.sencho', {
host: 'db.api.opsix.sencho',
nodeId: remoteNodeId,
nodeName: 'opsix',
stackName: 'api',
serviceName: 'db',
port: 5432,
}],
]);
db.insertMeshStack(remoteNodeId, 'api', 'tester');
const result = await svc.testUpstream('db.api.opsix.sencho', localNodeId);
expect(result.ok).toBe(false);
expect(result.where).toBe('pilot_tunnel');
expect(result.code).toBe('tunnel_down');
});
it('returns ok:false where=agent_resolve when target stack is not opted in', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
(svc as unknown as { aliasCache: Map<string, unknown> }).aliasCache = new Map([
['db.api.opsix.sencho', {
host: 'db.api.opsix.sencho',
nodeId: localNodeId,
nodeName: 'opsix',
stackName: 'api',
serviceName: 'db',
port: 5432,
}],
]);
const result = await svc.testUpstream('db.api.opsix.sencho', localNodeId);
expect(result.ok).toBe(false);
expect(result.where).toBe('agent_resolve');
expect(result.code).toBe('denied');
});
it('returns ok:false where=sidecar when alias is unknown', async () => {
const svc = MeshService.getInstance();
const localNodeId = DatabaseService.getInstance().getNodes()[0].id;
const result = await svc.testUpstream('nonexistent.sencho', localNodeId);
expect(result.ok).toBe(false);
expect(result.where).toBe('no_route');
expect(result.code).toBe('no_route');
});
it('probes a proxy-mode remote via PilotTunnelManager.ensureBridge (bridge dialed on demand)', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'edge', type: 'remote', is_default: false,
compose_dir: '/tmp', api_url: 'https://edge.example',
api_token: 'tok', mode: 'proxy',
});
(svc as unknown as { aliasCache: Map<string, unknown> }).aliasCache = new Map([
['db.api.edge.sencho', {
host: 'db.api.edge.sencho',
nodeId: remoteNodeId,
nodeName: 'edge',
stackName: 'api',
serviceName: 'db',
port: 5432,
}],
]);
db.insertMeshStack(remoteNodeId, 'api', 'tester');
const { PilotTunnelManager } = await import('../services/PilotTunnelManager');
const fakeStream = new EventEmitter() as EventEmitter & { destroy: () => void };
fakeStream.destroy = vi.fn();
const fakeBridge = {
openTcpStream: vi.fn().mockReturnValue(fakeStream),
getActiveStreamCount: () => 0,
close: vi.fn(),
};
vi.spyOn(PilotTunnelManager.getInstance(), 'ensureBridge')
.mockResolvedValue(fakeBridge as unknown as Awaited<ReturnType<typeof PilotTunnelManager.prototype.ensureBridge>>);
const probe = svc.testUpstream('db.api.edge.sencho', localNodeId);
// Emit `open` so the probe resolves cleanly.
setImmediate(() => fakeStream.emit('open'));
const result = await probe;
expect(fakeBridge.openTcpStream).toHaveBeenCalledWith({ stack: 'api', service: 'db', port: 5432 });
expect(result.ok).toBe(true);
db.deleteNode(remoteNodeId);
});
it('returns ok:false where=pilot_tunnel when ensureBridge yields null (no reachable bridge)', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'edge2', type: 'remote', is_default: false,
compose_dir: '/tmp', api_url: 'https://edge2.example',
api_token: 'tok', mode: 'proxy',
});
(svc as unknown as { aliasCache: Map<string, unknown> }).aliasCache = new Map([
['db.api.edge2.sencho', {
host: 'db.api.edge2.sencho',
nodeId: remoteNodeId,
nodeName: 'edge2',
stackName: 'api',
serviceName: 'db',
port: 5432,
}],
]);
db.insertMeshStack(remoteNodeId, 'api', 'tester');
const { PilotTunnelManager } = await import('../services/PilotTunnelManager');
vi.spyOn(PilotTunnelManager.getInstance(), 'ensureBridge').mockResolvedValue(null);
const result = await svc.testUpstream('db.api.edge2.sencho', localNodeId);
expect(result.ok).toBe(false);
expect(result.where).toBe('pilot_tunnel');
expect(result.code).toBe('tunnel_down');
db.deleteNode(remoteNodeId);
});
});
describe('getSenchoIpFromSubnet', () => {
it('returns network+2 for the default /24', () => {
expect(getSenchoIpFromSubnet('172.30.0.0/24')).toBe('172.30.0.2');
});
it('handles a custom /24 in a different range', () => {
expect(getSenchoIpFromSubnet('10.42.7.0/24')).toBe('10.42.7.2');
});
it('handles a /16', () => {
expect(getSenchoIpFromSubnet('172.30.0.0/16')).toBe('172.30.0.2');
});
it('masks the input IP to the network address before adding 2', () => {
// 172.30.0.50/24 → network 172.30.0.0 → +2 = 172.30.0.2
expect(getSenchoIpFromSubnet('172.30.0.50/24')).toBe('172.30.0.2');
});
it('rejects a malformed CIDR', () => {
expect(() => getSenchoIpFromSubnet('not-a-cidr')).toThrow(/Invalid mesh subnet/);
expect(() => getSenchoIpFromSubnet('172.30.0.0')).toThrow(/Invalid mesh subnet/);
});
it('rejects prefixes too narrow to host two addresses', () => {
expect(() => getSenchoIpFromSubnet('172.30.0.0/31')).toThrow(/Invalid mesh subnet/);
});
it('rejects out-of-range octets', () => {
expect(() => getSenchoIpFromSubnet('172.30.0.999/24')).toThrow(/Invalid mesh subnet/);
});
});
describe('MeshService.optInStack rollback', () => {
it('rolls back the DB row when the just-inserted stack fails to push its override', async () => {
const svc = MeshService.getInstance();
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'db', ports: [5432] }]);
vi.spyOn(svc, 'pushOverrideToNode')
.mockRejectedValue(new Error('simulated remote pilot offline'));
vi.spyOn(svc as unknown as { triggerRedeploy: (n: number, s: string, a: string) => void }, 'triggerRedeploy')
.mockImplementation(() => { /* noop */ });
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
await expect(svc.optInStack(localNodeId, 'api', 'tester'))
.rejects.toThrow(/simulated remote pilot offline/);
expect(db.isMeshStackEnabled(localNodeId, 'api')).toBe(false);
});
it('retains remote authority when the push outcome is unknown', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const remoteNodeId = db.addNode({
name: 'ambiguous-push', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [8080] }]);
vi.spyOn(svc, 'pushOverrideToNode').mockRejectedValue(new Error('connection reset'));
await expect(svc.optInStack(remoteNodeId, 'ambiguous-stack', 'tester'))
.rejects.toThrow('connection reset');
expect(db.isMeshStackEnabled(remoteNodeId, 'ambiguous-stack')).toBe(true);
db.deleteNode(remoteNodeId);
});
it('treats a remote gateway error as an ambiguous push outcome', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const remoteNodeId = db.addNode({
name: 'gateway-error-push', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [8080] }]);
vi.spyOn(svc as unknown as { proxyFetch: () => Promise<Response> }, 'proxyFetch')
.mockResolvedValue(new Response('gateway timeout', { status: 502 }));
await expect(svc.optInStack(remoteNodeId, 'gateway-error-stack', 'tester'))
.rejects.toThrow('HTTP 502');
expect(db.isMeshStackEnabled(remoteNodeId, 'gateway-error-stack')).toBe(true);
db.deleteNode(remoteNodeId);
});
it('rolls back remote authority when the target explicitly rejects the push', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const remoteNodeId = db.addNode({
name: 'rejected-push', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [8080] }]);
vi.spyOn(svc, 'pushOverrideToNode')
.mockRejectedValue(new MeshError('push_failed', 'target rejected override'));
await expect(svc.optInStack(remoteNodeId, 'rejected-stack', 'tester'))
.rejects.toThrow('target rejected override');
expect(db.isMeshStackEnabled(remoteNodeId, 'rejected-stack')).toBe(false);
db.deleteNode(remoteNodeId);
});
});
describe('MeshService.optInStack guard rails (network setup)', () => {
it('rejects opt-in when senchoIp is null (mesh data plane unavailable)', async () => {
const svc = MeshService.getInstance();
(svc as unknown as { senchoIp: string | null }).senchoIp = null;
(svc as unknown as { networkSetupError: string | null }).networkSetupError = 'sencho_mesh subnet mismatch';
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
await expect(svc.optInStack(localNodeId, 'api', 'tester'))
.rejects.toThrow(/sencho_mesh subnet mismatch/);
expect(db.isMeshStackEnabled(localNodeId, 'api')).toBe(false);
});
it('rejects opt-in when a service exposes the reserved Sencho API port', async () => {
const svc = MeshService.getInstance();
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [1852] }]);
vi.spyOn(svc as unknown as { regenerateOverridesAcrossFleet: (n?: number, s?: string) => Promise<void> }, 'regenerateOverridesAcrossFleet')
.mockResolvedValue(undefined);
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
await expect(svc.optInStack(localNodeId, 'api', 'tester'))
.rejects.toThrow(/port 1852 is reserved/);
expect(db.isMeshStackEnabled(localNodeId, 'api')).toBe(false);
});
});
describe('MeshService.regenerateAllOverrides (F6: boot-time regen)', () => {
it('pushes every mesh_stacks row across the fleet and returns a summary', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.insertMeshStack(localNodeId, 'audit-mesh-prod', 'tester');
db.insertMeshStack(remoteNodeId, 'audit-mesh-pilot', 'tester');
const pushSpy = vi.spyOn(svc, 'pushOverrideToNode').mockResolvedValue(undefined);
try {
const summary = await svc.regenerateAllOverrides();
expect(summary.skipped).toBe(false);
expect(summary.regenerated).toBe(2);
expect(summary.failures).toEqual([]);
expect(pushSpy).toHaveBeenCalledTimes(2);
expect(pushSpy).toHaveBeenCalledWith(localNodeId, 'audit-mesh-prod');
expect(pushSpy).toHaveBeenCalledWith(remoteNodeId, 'audit-mesh-pilot');
const activity = svc.getActivity({ limit: 100 });
expect(activity.some((e) =>
e.message === 'mesh override regen complete: 2 succeeded, 0 failed across 0 node(s)',
)).toBe(true);
} finally {
db.deleteNode(remoteNodeId);
}
});
it('skips entirely with a reason when senchoIp is null (network setup failed)', async () => {
const svc = MeshService.getInstance();
(svc as unknown as { senchoIp: string | null }).senchoIp = null;
(svc as unknown as { networkSetupError: string | null }).networkSetupError = 'sencho_mesh subnet mismatch';
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
db.insertMeshStack(localNodeId, 'audit-mesh-prod', 'tester');
const pushSpy = vi.spyOn(svc, 'pushOverrideToNode').mockResolvedValue(undefined);
const summary = await svc.regenerateAllOverrides();
expect(summary.skipped).toBe(true);
expect(summary.reason).toBe('sencho_mesh subnet mismatch');
expect(summary.regenerated).toBe(0);
expect(pushSpy).not.toHaveBeenCalled();
const activity = svc.getActivity({ limit: 100 });
expect(activity.some((e) =>
e.level === 'warn' && /data plane unavailable/.test(e.message),
)).toBe(true);
});
it('records per-stack failures in the summary and aggregates by node', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
db.insertMeshStack(localNodeId, 'audit-mesh-prod', 'tester');
vi.spyOn(svc, 'pushOverrideToNode').mockRejectedValue(new Error('remote node offline'));
const summary = await svc.regenerateAllOverrides();
expect(summary.skipped).toBe(false);
expect(summary.regenerated).toBe(0);
expect(summary.failures).toEqual([{
nodeId: localNodeId,
stackName: 'audit-mesh-prod',
message: 'remote node offline',
}]);
const activity = svc.getActivity({ limit: 100 });
expect(activity.some((e) =>
e.level === 'warn' && /mesh override regen failed for audit-mesh-prod/.test(e.message),
)).toBe(true);
expect(activity.some((e) =>
/mesh override regen complete: 0 succeeded, 1 failed across 1 node\(s\)/.test(e.message),
)).toBe(true);
});
it('continues past version-skew 404 push failures and reports them in the summary', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const remoteNodeId = db.addNode({
name: 'remote-old-pilot', type: 'remote', mode: 'pilot_agent',
compose_dir: '/tmp', is_default: false, api_url: '', api_token: '',
});
db.insertMeshStack(localNodeId, 'audit-mesh-prod', 'tester');
db.insertMeshStack(remoteNodeId, 'audit-mesh-pilot', 'tester');
vi.spyOn(svc, 'pushOverrideToNode').mockImplementation(async (nodeId: number) => {
if (nodeId === remoteNodeId) {
throw new MeshError('push_failed', 'node remote-old-pilot does not support mesh override push (upgrade required)');
}
});
try {
const summary = await svc.regenerateAllOverrides();
expect(summary.regenerated).toBe(1);
expect(summary.failures).toHaveLength(1);
expect(summary.failures[0].nodeId).toBe(remoteNodeId);
expect(summary.failures[0].stackName).toBe('audit-mesh-pilot');
expect(summary.failures[0].message).toMatch(/upgrade required/);
} finally {
db.deleteNode(remoteNodeId);
}
});
it('completes both flows without throwing when regen and opt-in are scheduled concurrently', async () => {
// The race is benign by design: both writes target the same alias list
// with the same `senchoIp`, so even if regen reads `mesh_stacks` mid-
// opt-in and pushes a concurrent override for the new stack, the
// resulting on-disk file is identical. This test locks the no-throw
// contract; it does not attempt to force a specific interleave because
// `regenerateAllOverrides` snapshots the table synchronously before
// its first await, so under the JS event loop the two flows always
// resolve cleanly without a real read-after-write window.
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
db.insertMeshStack(localNodeId, 'existing-stack', 'tester');
vi.spyOn(svc as unknown as { inspectStackServices: (n: number, s: string) => Promise<unknown> }, 'inspectStackServices')
.mockResolvedValue([{ service: 'web', ports: [8080] }]);
const pushSpy = vi.spyOn(svc, 'pushOverrideToNode').mockResolvedValue(undefined);
vi.spyOn(svc as unknown as { triggerRedeploy: (n: number, s: string, a: string) => void }, 'triggerRedeploy')
.mockImplementation(() => { /* noop */ });
vi.spyOn(svc as unknown as { regenerateOverridesAcrossFleet: (n?: number, skip?: string) => Promise<void> }, 'regenerateOverridesAcrossFleet')
.mockResolvedValue(undefined);
const [summary] = await Promise.all([
svc.regenerateAllOverrides(),
svc.optInStack(localNodeId, 'concurrent-stack', 'tester'),
]);
expect(summary.regenerated).toBeGreaterThanOrEqual(1);
expect(db.isMeshStackEnabled(localNodeId, 'concurrent-stack')).toBe(true);
const existingCalls = pushSpy.mock.calls.filter((c) => c[1] === 'existing-stack');
expect(existingCalls.length).toBeGreaterThanOrEqual(1);
const concurrentCalls = pushSpy.mock.calls.filter((c) => c[1] === 'concurrent-stack');
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);
});
it('returns a pushed override file on pilot nodes where isMeshStackEnabled is always false', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
process.env.SENCHO_MODE = 'pilot';
// Simulate the pilot scenario: no mesh_stacks row (isMeshStackEnabled → false),
// but the override file already exists on disk, pushed by central via D-1.
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, 'pilot-stack.override.yml');
fsSync.writeFileSync(overrideFile, [
'services:',
' echo:',
' networks:',
' - sencho_mesh',
' extra_hosts:',
' - echo.pilot-stack.pilot.sencho:172.30.0.2',
'networks:',
' sencho_mesh:',
' external: true',
].join('\n'), 'utf8');
// No mesh_stacks row → isMeshStackEnabled returns false.
const result = await svc.ensureStackOverride(localNodeId, 'pilot-stack');
expect(result).toBe(overrideFile);
// Cleanup.
fsSync.unlinkSync(overrideFile);
process.env.SENCHO_MODE = 'server';
});
it('persists and removes proxy-target opt-in state with a pushed local override', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
vi.spyOn(svc, 'getDeclaredStackServiceNames').mockResolvedValue(['web']);
const file = await svc.applyLocalOverride('proxy-stack', []);
expect(file).not.toBeNull();
const yaml = fsSync.readFileSync(file as string, 'utf8');
expect(yaml).toContain('web:');
expect(yaml).toContain('sencho_mesh');
expect(fsSync.readdirSync(path.dirname(file as string)).some((name) => name.includes('proxy-stack') && name.endsWith('.tmp'))).toBe(false);
expect(db.isMeshStackEnabled(localNodeId, 'proxy-stack')).toBe(true);
await svc.removeLocalOverride('proxy-stack');
expect(db.isMeshStackEnabled(localNodeId, 'proxy-stack')).toBe(false);
});
it('does not write a proxy override when DB authority cannot be recorded', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
vi.spyOn(svc, 'getDeclaredStackServiceNames').mockResolvedValue(['web']);
vi.spyOn(db, 'insertMeshStack').mockImplementation(() => {
throw new Error('database unavailable');
});
await expect(svc.applyLocalOverride('db-failure', [])).rejects.toThrow('database unavailable');
const overrideFile = path.join(
process.env.DATA_DIR as string,
'mesh',
'overrides',
String(localNodeId),
'db-failure.override.yml',
);
expect(fsSync.existsSync(overrideFile)).toBe(false);
});
it('restores an existing override when DB authority cannot be recorded', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const overrideDir = path.join(process.env.DATA_DIR as string, 'mesh', 'overrides', String(localNodeId));
fsSync.mkdirSync(overrideDir, { recursive: true });
const overrideFile = path.join(overrideDir, 'db-replacement-failure.override.yml');
const originalYaml = 'services:\n prior:\n networks:\n - sencho_mesh\n';
fsSync.writeFileSync(overrideFile, originalYaml, 'utf8');
vi.spyOn(svc, 'getDeclaredStackServiceNames').mockResolvedValue(['web']);
vi.spyOn(db, 'insertMeshStack').mockImplementation(() => {
throw new Error('database unavailable');
});
await expect(svc.applyLocalOverride('db-replacement-failure', [])).rejects.toThrow('database unavailable');
expect(fsSync.readFileSync(overrideFile, 'utf8')).toBe(originalYaml);
expect(db.isMeshStackEnabled(localNodeId, 'db-replacement-failure')).toBe(false);
expect(fsSync.readdirSync(overrideDir).some((name) => name.includes('db-replacement-failure') && name.endsWith('.tmp'))).toBe(false);
});
it('does not publish DB authority or a final file when atomic override publication fails', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
vi.spyOn(svc, 'getDeclaredStackServiceNames').mockResolvedValue(['web']);
vi.spyOn(fs, 'rename').mockRejectedValue(new Error('rename failed'));
await expect(svc.applyLocalOverride('write-failure', [])).rejects.toThrow('rename failed');
const overrideDir = path.join(process.env.DATA_DIR as string, 'mesh', 'overrides', String(localNodeId));
expect(db.isMeshStackEnabled(localNodeId, 'write-failure')).toBe(false);
expect(fsSync.existsSync(path.join(overrideDir, 'write-failure.override.yml'))).toBe(false);
expect(fsSync.readdirSync(overrideDir).some((name) => name.includes('write-failure'))).toBe(false);
});
it('preserves the prior override and authority when atomic replacement fails', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const overrideDir = path.join(process.env.DATA_DIR as string, 'mesh', 'overrides', String(localNodeId));
fsSync.mkdirSync(overrideDir, { recursive: true });
const overrideFile = path.join(overrideDir, 'replacement-failure.override.yml');
const originalYaml = 'services:\n prior:\n networks:\n - sencho_mesh\n';
fsSync.writeFileSync(overrideFile, originalYaml, 'utf8');
db.insertMeshStack(localNodeId, 'replacement-failure', 'tester');
vi.spyOn(svc, 'getDeclaredStackServiceNames').mockResolvedValue(['web']);
vi.spyOn(fs, 'rename').mockRejectedValue(new Error('rename failed'));
await expect(svc.applyLocalOverride('replacement-failure', [])).rejects.toThrow('rename failed');
expect(fsSync.readFileSync(overrideFile, 'utf8')).toBe(originalYaml);
expect(db.isMeshStackEnabled(localNodeId, 'replacement-failure')).toBe(true);
});
it('prevents overlapping override mutations for the same stack', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
let releaseServices!: (services: string[]) => void;
const servicesPending = new Promise<string[]>((resolve) => {
releaseServices = resolve;
});
vi.spyOn(svc, 'getDeclaredStackServiceNames').mockReturnValue(servicesPending);
const first = svc.applyLocalOverride('concurrent-stack', []);
await vi.waitFor(() => expect(svc.getDeclaredStackServiceNames).toHaveBeenCalledTimes(1));
await expect(svc.applyLocalOverride('concurrent-stack', [])).rejects.toThrow('another operation');
releaseServices(['web']);
await first;
expect(db.isMeshStackEnabled(localNodeId, 'concurrent-stack')).toBe(true);
});
it('prevents removal from overlapping an in-flight override apply', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
let releaseServices!: (services: string[]) => void;
const servicesPending = new Promise<string[]>((resolve) => {
releaseServices = resolve;
});
vi.spyOn(svc, 'getDeclaredStackServiceNames').mockReturnValue(servicesPending);
const apply = svc.applyLocalOverride('apply-remove-stack', []);
await vi.waitFor(() => expect(svc.getDeclaredStackServiceNames).toHaveBeenCalledTimes(1));
await expect(svc.removeLocalOverride('apply-remove-stack')).rejects.toThrow('another operation');
releaseServices(['web']);
const file = await apply;
expect(file).not.toBeNull();
expect(fsSync.existsSync(file as string)).toBe(true);
expect(db.isMeshStackEnabled(localNodeId, 'apply-remove-stack')).toBe(true);
});
it('does not report committed removal as failed when alias refresh fails', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const overrideDir = path.join(process.env.DATA_DIR as string, 'mesh', 'overrides', String(localNodeId));
fsSync.mkdirSync(overrideDir, { recursive: true });
const overrideFile = path.join(overrideDir, 'refresh-failure.override.yml');
fsSync.writeFileSync(overrideFile, 'services: {}\n', 'utf8');
db.insertMeshStack(localNodeId, 'refresh-failure', 'setup');
(svc as unknown as { pilotAliasOverlay: Map<string, unknown> }).pilotAliasOverlay.set('refresh-failure', []);
vi.spyOn(svc as unknown as { refreshAliasCache: () => Promise<void> }, 'refreshAliasCache')
.mockRejectedValue(new Error('refresh failed'));
vi.spyOn(console, 'warn').mockImplementation(() => { /* silence */ });
await expect(svc.removeLocalOverride('refresh-failure')).resolves.toBeUndefined();
expect(fsSync.existsSync(overrideFile)).toBe(false);
expect(db.isMeshStackEnabled(localNodeId, 'refresh-failure')).toBe(false);
});
it('does not report committed apply as failed when alias refresh fails', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
vi.spyOn(svc, 'getDeclaredStackServiceNames').mockResolvedValue(['web']);
vi.spyOn(svc as unknown as { refreshAliasCache: () => Promise<void> }, 'refreshAliasCache')
.mockRejectedValue(new Error('refresh failed'));
vi.spyOn(console, 'warn').mockImplementation(() => { /* silence */ });
const file = await svc.applyLocalOverride('apply-refresh-failure', [], [{
host: 'web.example', nodeId: localNodeId, nodeName: 'local', stackName: 'apply-refresh-failure',
serviceName: 'web', port: 8080,
}]);
expect(file).not.toBeNull();
expect(fsSync.existsSync(file as string)).toBe(true);
expect(db.isMeshStackEnabled(localNodeId, 'apply-refresh-failure')).toBe(true);
});
it('returns null for pilot nodes when no pushed override file exists', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
process.env.SENCHO_MODE = 'pilot';
// No mesh_stacks row, no file on disk.
const result = await svc.ensureStackOverride(localNodeId, 'no-such-stack');
expect(result).toBeNull();
process.env.SENCHO_MODE = 'server';
});
it('does not use stale override presence as authority on a server', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
const overrideDir = path.join(process.env.DATA_DIR as string, 'mesh', 'overrides', String(localNodeId));
fsSync.mkdirSync(overrideDir, { recursive: true });
const overrideFile = path.join(overrideDir, 'stale-server.override.yml');
fsSync.writeFileSync(overrideFile, 'services: {}\n', 'utf8');
const result = await svc.ensureStackOverride(localNodeId, 'stale-server');
expect(result).toBeNull();
fsSync.unlinkSync(overrideFile);
});
it('restores proxy-target DB authority when override removal fails', async () => {
const svc = MeshService.getInstance();
const db = DatabaseService.getInstance();
const localNodeId = db.getNodes()[0].id;
db.insertMeshStack(localNodeId, 'unlink-failure', 'tester');
vi.spyOn(fs, 'unlink').mockRejectedValue(Object.assign(new Error('permission denied'), { code: 'EACCES' }));
await expect(svc.removeLocalOverride('unlink-failure')).rejects.toThrow('permission denied');
expect(db.isMeshStackEnabled(localNodeId, 'unlink-failure')).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<void> },
'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<void>;
refreshAliasCache: () => Promise<void>;
syncForwarderListeners: () => Promise<void>;
regenerateAllOverrides: () => Promise<unknown>;
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<typeof vi.fn>; end: ReturnType<typeof vi.fn>; on: ReturnType<typeof vi.fn>; write: ReturnType<typeof vi.fn> } {
return {
destroy: vi.fn(),
end: vi.fn(),
on: vi.fn(),
write: vi.fn(),
};
}
// Install a stub reverseDialer so openCrossNode routes through the
// multiplex path rather than depending on a live MeshProxyTunnelDialer
// bridge being open against a real peer.
const stubDialer = { openMeshTcpStream: vi.fn() };
beforeEach(() => { MeshService.getInstance().setReverseDialer(stubDialer); });
afterEach(() => { MeshService.getInstance().setReverseDialer(null); });
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();
}
});
// F-10 regression: src.on('close')/src.on('error') must delete the
// activeStreams entry, not just call tcpStream.destroy(). MeshTcpStream-
// Like.destroy() only sends a tcp_close frame; the handle's own 'close'
// event waits for the remote ack. If that ack never lands (peer gone,
// network drop), the record sits in activeStreams until tunnel idle-
// close. Use an EventEmitter-backed fake src so emit('close') actually
// delivers to the listener (the plain makeFakeSocket uses vi.fn for .on
// which records calls but never fires them).
function makeEmittingSocket(): EventEmitter & {
destroy: ReturnType<typeof vi.fn>;
end: ReturnType<typeof vi.fn>;
write: ReturnType<typeof vi.fn>;
} {
const ee = new EventEmitter() as EventEmitter & {
destroy: ReturnType<typeof vi.fn>;
end: ReturnType<typeof vi.fn>;
write: ReturnType<typeof vi.fn>;
};
ee.destroy = vi.fn();
ee.end = vi.fn();
ee.write = vi.fn();
return ee;
}
it('src.on(close) deletes activeStreams entry even when tcpStream never emits close', 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(50);
vi.spyOn(
svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null },
'dialMeshTcpStream',
).mockReturnValue(fakeStream);
const fakeSrc = makeEmittingSocket();
await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise<void> })
.openCrossNode(target, fakeSrc);
const activeStreams = (svc as unknown as { activeStreams: Map<number, unknown> }).activeStreams;
expect(activeStreams.size).toBe(1);
expect(activeStreams.has(50)).toBe(true);
fakeSrc.emit('close');
expect(activeStreams.size).toBe(0);
expect(fakeStream.destroy).toHaveBeenCalled();
});
it('src.on(error) deletes activeStreams entry symmetrically', 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(51);
vi.spyOn(
svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null },
'dialMeshTcpStream',
).mockReturnValue(fakeStream);
const fakeSrc = makeEmittingSocket();
await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise<void> })
.openCrossNode(target, fakeSrc);
const activeStreams = (svc as unknown as { activeStreams: Map<number, unknown> }).activeStreams;
expect(activeStreams.size).toBe(1);
fakeSrc.emit('error', new Error('connection reset by peer'));
expect(activeStreams.size).toBe(0);
expect(fakeStream.destroy).toHaveBeenCalled();
});
it('cleanup is idempotent when both src and tcpStream emit close', 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(52);
vi.spyOn(
svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null },
'dialMeshTcpStream',
).mockReturnValue(fakeStream);
const fakeSrc = makeEmittingSocket();
await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise<void> })
.openCrossNode(target, fakeSrc);
const activeStreams = (svc as unknown as { activeStreams: Map<number, unknown> }).activeStreams;
expect(activeStreams.size).toBe(1);
expect(() => {
fakeSrc.emit('close');
fakeStream.emit('close');
}).not.toThrow();
expect(activeStreams.size).toBe(0);
});
// Pins the ordering contract for protocols that send immediately after
// connect (HTTP, TLS, Redis, Postgres): src bytes arriving between
// dial and tcp_open_ack must reach the upstream in order, ahead of any
// post-ack writes. Without local buffering the first packet would race
// the ack on the wire and land on a stream the agent has not yet dialed.
it('buffers src bytes until tcpStream opens, then flushes them in order before any post-open writes', 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(60);
vi.spyOn(
svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null },
'dialMeshTcpStream',
).mockReturnValue(fakeStream);
const fakeSrc = makeEmittingSocket();
await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise<void> })
.openCrossNode(target, fakeSrc);
// Push data before the stream is open. Nothing should be written
// through to the tunnel.
fakeSrc.emit('data', Buffer.from('GET /healthz HTTP/1.1\r\n'));
fakeSrc.emit('data', Buffer.from('Host: echo.local\r\n\r\n'));
expect(fakeStream.write).not.toHaveBeenCalled();
// Open the stream. The pre-open buffer must flush in order, ahead
// of any post-open writes.
fakeStream.emit('open');
const writeCalls = (fakeStream.write as unknown as { mock: { calls: Buffer[][] } }).mock.calls;
expect(writeCalls.length).toBe(2);
expect(writeCalls[0][0].toString()).toBe('GET /healthz HTTP/1.1\r\n');
expect(writeCalls[1][0].toString()).toBe('Host: echo.local\r\n\r\n');
// Post-open data writes through directly (still in order).
fakeSrc.emit('data', Buffer.from('post-open'));
expect(writeCalls.length).toBe(3);
expect(writeCalls[2][0].toString()).toBe('post-open');
// Cleanup so the open-timer does not keep the worker alive.
fakeStream.emit('close');
});
it('tears down both sockets when buffered bytes exceed STREAM_PENDING_DATA_MAX_BYTES before open', async () => {
const { STREAM_PENDING_DATA_MAX_BYTES } = await import('../pilot/protocol');
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(61);
vi.spyOn(
svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null },
'dialMeshTcpStream',
).mockReturnValue(fakeStream);
const fakeSrc = makeEmittingSocket();
await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise<void> })
.openCrossNode(target, fakeSrc);
// Overflow the per-stream cap in one shot. The handler must destroy
// both sockets and not call write on the stream.
fakeSrc.emit('data', Buffer.alloc(STREAM_PENDING_DATA_MAX_BYTES + 1));
expect(fakeStream.write).not.toHaveBeenCalled();
expect(fakeSrc.destroy).toHaveBeenCalled();
expect(fakeStream.destroy).toHaveBeenCalled();
});
});
describe('MeshService.openCrossNode without reverseDialer', () => {
// Forward + reverse mesh traffic now share the same central-initiated
// bridge. When openCrossNode runs without a reverseDialer installed it
// falls straight through to dialMeshTcpStream, which on central uses
// MeshProxyTunnelDialer.ensureBridge to reach a proxy peer. The
// previous peer-recovery branch (which depended on the now-removed
// mesh_centrals table) is gone.
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<typeof vi.fn>; end: ReturnType<typeof vi.fn>; on: ReturnType<typeof vi.fn>; write: ReturnType<typeof vi.fn> } {
return { destroy: vi.fn(), end: vi.fn(), on: vi.fn(), write: vi.fn() };
}
beforeEach(() => {
MeshService.getInstance().setReverseDialer(null);
});
it('central (no reverseDialer) falls through to dialMeshTcpStream without route.resolve.fail forward-from-peer', async () => {
const svc = MeshService.getInstance();
const target: MeshTarget = {
nodeId: 7, stack: 'audit-mesh-proxy', service: 'echo',
port: 9002, alias: 'echo.audit-mesh-proxy.sencho-test-03.sencho',
};
const fakeStream = makeFakeStream(101);
const dialSpy = vi.spyOn(
svc as unknown as { dialMeshTcpStream: (t: MeshTarget) => MeshTcpStreamLike | null },
'dialMeshTcpStream',
).mockReturnValue(fakeStream);
const fakeSrc = makeFakeSocket();
await (svc as unknown as { openCrossNode: (t: MeshTarget, s: unknown) => Promise<void> })
.openCrossNode(target, fakeSrc);
const events = svc.getActivity({ limit: 50 });
expect(events.some((e) => e.type === 'route.dispatch')).toBe(true);
const wrongFail = events.find((e) =>
e.type === 'route.resolve.fail'
&& (e.details as { direction?: string } | undefined)?.direction === 'forward-from-peer',
);
expect(wrongFail).toBeUndefined();
expect(dialSpy).toHaveBeenCalledTimes(1);
expect(fakeSrc.destroy).not.toHaveBeenCalled();
fakeStream.emit('close');
});
});
describe('MeshService pilot handleAccept dispatch', () => {
function makeEnrollToken(nodeId: number): string {
const header = Buffer.from(JSON.stringify({ alg: 'HS256', typ: 'JWT' })).toString('base64url');
const payload = Buffer.from(JSON.stringify({ scope: 'pilot_enroll', nodeId })).toString('base64url');
return `${header}.${payload}.fakesig`;
}
it('resolveSelfCentralNodeId extracts nodeId from SENCHO_ENROLL_TOKEN', () => {
process.env.SENCHO_ENROLL_TOKEN = makeEnrollToken(14);
const svc = MeshService.getInstance() as unknown as {
resolveSelfCentralNodeId: () => number;
};
expect(svc.resolveSelfCentralNodeId()).toBe(14);
});
it('resolveSelfCentralNodeId falls back to local default when token is absent', () => {
delete process.env.SENCHO_ENROLL_TOKEN;
const svc = MeshService.getInstance() as unknown as {
resolveSelfCentralNodeId: () => number;
};
const db = DatabaseService.getInstance();
const defaultId = db.getDefaultNode()?.id ?? 1;
expect(svc.resolveSelfCentralNodeId()).toBe(defaultId);
});
it('resolveSelfCentralNodeId falls back to local default for a malformed token', () => {
process.env.SENCHO_ENROLL_TOKEN = 'not.a.jwt';
const svc = MeshService.getInstance() as unknown as {
resolveSelfCentralNodeId: () => number;
};
const db = DatabaseService.getInstance();
const defaultId = db.getDefaultNode()?.id ?? 1;
expect(svc.resolveSelfCentralNodeId()).toBe(defaultId);
});
it('handleAccept routes same-node alias to openSameNode on a pilot', async () => {
const svc = MeshService.getInstance();
const internals = svc as unknown as {
selfCentralNodeId: number | null;
aliasByPort: Map<number, unknown>;
openSameNode: (t: MeshTarget, s: unknown) => Promise<void>;
openCrossNode: (t: MeshTarget, s: unknown) => void;
};
internals.selfCentralNodeId = 14;
internals.aliasByPort.set(9001, {
host: 'echo.audit-mesh-pilot.sencho-pilot-test.sencho',
nodeId: 14,
nodeName: 'sencho-pilot-test',
stackName: 'audit-mesh-pilot',
serviceName: 'echo',
port: 9001,
});
const openSame = vi.spyOn(internals, 'openSameNode').mockResolvedValue(undefined);
const openCross = vi.spyOn(internals, 'openCrossNode').mockImplementation(() => undefined);
const fakeSrc = { remoteAddress: '127.0.0.1', destroy: vi.fn() } as unknown as import('net').Socket;
await svc.handleAccept(9001, fakeSrc);
expect(openSame).toHaveBeenCalledOnce();
expect(openCross).not.toHaveBeenCalled();
});
it('handleAccept routes cross-node alias to openCrossNode on a pilot', async () => {
const svc = MeshService.getInstance();
const internals = svc as unknown as {
selfCentralNodeId: number | null;
aliasByPort: Map<number, unknown>;
openSameNode: (t: MeshTarget, s: unknown) => Promise<void>;
openCrossNode: (t: MeshTarget, s: unknown) => void;
};
internals.selfCentralNodeId = 14;
internals.aliasByPort.set(9000, {
host: 'echo.audit-mesh-prod.Local.sencho',
nodeId: 1,
nodeName: 'Local',
stackName: 'audit-mesh-prod',
serviceName: 'echo',
port: 9000,
});
const openSame = vi.spyOn(internals, 'openSameNode').mockResolvedValue(undefined);
const openCross = vi.spyOn(internals, 'openCrossNode').mockImplementation(() => undefined);
const fakeSrc = { remoteAddress: '127.0.0.1', destroy: vi.fn() } as unknown as import('net').Socket;
await svc.handleAccept(9000, fakeSrc);
expect(openCross).toHaveBeenCalledOnce();
expect(openSame).not.toHaveBeenCalled();
});
it('handleAccept on a proxy peer uses proxyTunnelSelfCentralNodeId to route cross-node aliases correctly (R1)', async () => {
// Repro for the R1 bug: a proxy peer receives an overlay carrying
// central-namespace nodeIds (e.g., Local = 1, this peer = 14). Pre-R1
// the peer had no selfCentralNodeId source, fell back to its local DB
// default (always 1), and falsely matched alias.nodeId=1 to its own
// selfNodeId=1 — dispatching cross-node aliases as same-node.
const svc = MeshService.getInstance();
const internals = svc as unknown as {
proxyTunnelSelfCentralNodeId: number | null;
selfCentralNodeId: number | null;
aliasByPort: Map<number, unknown>;
openSameNode: (t: MeshTarget, s: unknown) => Promise<void>;
openCrossNode: (t: MeshTarget, s: unknown) => void;
};
// Proxy peer: selfCentralNodeId is null (no SENCHO_ENROLL_TOKEN),
// the proxy-tunnel handler installed central's view of this peer.
internals.selfCentralNodeId = null;
svc.setProxyTunnelSelfCentralNodeId(14);
// Overlay alias for central's own stack (Local = nodeId 1 in
// central's namespace).
internals.aliasByPort.set(9000, {
host: 'echo.audit-mesh-central.Local.sencho',
nodeId: 1,
nodeName: 'Local',
stackName: 'audit-mesh-central',
serviceName: 'echo',
port: 9000,
});
const openSame = vi.spyOn(internals, 'openSameNode').mockResolvedValue(undefined);
const openCross = vi.spyOn(internals, 'openCrossNode').mockImplementation(() => undefined);
const fakeSrc = { remoteAddress: '127.0.0.1', destroy: vi.fn() } as unknown as import('net').Socket;
await svc.handleAccept(9000, fakeSrc);
expect(openCross).toHaveBeenCalledOnce();
expect(openSame).not.toHaveBeenCalled();
});
it('handleAccept on a proxy peer routes same-node aliases (matching the proxy-tunnel nodeId) to openSameNode', async () => {
const svc = MeshService.getInstance();
const internals = svc as unknown as {
proxyTunnelSelfCentralNodeId: number | null;
selfCentralNodeId: number | null;
aliasByPort: Map<number, unknown>;
openSameNode: (t: MeshTarget, s: unknown) => Promise<void>;
openCrossNode: (t: MeshTarget, s: unknown) => void;
};
internals.selfCentralNodeId = null;
svc.setProxyTunnelSelfCentralNodeId(14);
// Alias for a stack on this peer (nodeId 14 in central's namespace).
internals.aliasByPort.set(9002, {
host: 'echo.audit-mesh-proxy.sencho-test-03.sencho',
nodeId: 14,
nodeName: 'sencho-test-03',
stackName: 'audit-mesh-proxy',
serviceName: 'echo',
port: 9002,
});
const openSame = vi.spyOn(internals, 'openSameNode').mockResolvedValue(undefined);
const openCross = vi.spyOn(internals, 'openCrossNode').mockImplementation(() => undefined);
const fakeSrc = { remoteAddress: '127.0.0.1', destroy: vi.fn() } as unknown as import('net').Socket;
await svc.handleAccept(9002, fakeSrc);
expect(openSame).toHaveBeenCalledOnce();
expect(openCross).not.toHaveBeenCalled();
});
});