feat(mesh): collapse sidecar into Sencho process via in-process forwarder (#1000)

The separate saelix/sencho-mesh sidecar container is gone. The
forwarder logic that previously lived in mesh-sidecar/src/forwarder.ts
moves into the Sencho process as backend/src/services/MeshForwarder.ts,
a thin per-port net.Server lifecycle wrapper. MeshService implements
the host interface and owns resolve plus splice; MeshForwarder owns
listener boilerplate. One container per node, no separate image to
publish, no control WebSocket.

Operator-facing change: the Sencho container now runs in
network_mode: host so the forwarder can bind alias ports on the host
network where meshed containers' extra_hosts host-gateway entries
point. Without host network mode the listeners would land in the
container's namespace and inbound traffic from peers would never
reach them. The 1852:1852 port publish becomes a no-op under host
mode and is commented out in the operator template.

Same-node forward path now dials the target container's bridge IP
via Dockerode (preferring the compose default network for
deterministic selection across daemon versions) instead of 127.0.0.1.
The legacy 127.0.0.1 path only worked when the target service
published its port to the host; the IP path works regardless.

Cross-node mesh routing in this phase is central -> pilot direction
only via PilotTunnelManager.openTcpStream. Pilot -> central and
pilot <-> pilot via central relay land in Phase B with the
tcp_open_reverse frame.

Deletions:
- mesh-sidecar/ package entirely (Dockerfile, package, sources, tests)
- backend/src/websocket/meshControl.ts
- MeshService sidecar lifecycle: spawnSidecar, stopSidecar,
  isSidecarRunning, mintSidecarToken, verifySidecarToken,
  attachSidecarSocket, handleSidecarResolve, sendSidecar
- POST /api/mesh/nodes/:id/sidecar/restart route
- /api/mesh/control WS dispatch in upgradeHandler

Type cleanup: 'sidecar' literal removed from MeshActivitySource and
MeshProbeResult.where (also the frontend mirror). MeshNodeStatus
sidecarRunning becomes localForwarderListening (boolean | null) so
non-local nodes get a null instead of an unconditional false; the
honest semantic is "this view only knows the local forwarder state;
remote forwarder status lands in Phase B." MeshNodeDiagnostic
sidecar object becomes forwarder { listening, listenerCount }.

Frontend MeshDiagnosticsSheet drops the restart-sidecar action and
sidecar liveness card; surfaces forwarder state plus a "runs
in-process; no separate container" caption.

Resolves audit findings C-1 (data plane non-functional), C-2 (sidecar
control WS not loopback-enforced), C-4 (sidecar lifecycle Dockerode-
on-remote, PR #999 closed), and C-5 (saelix/sencho-mesh:latest
unreachable). C-3 (PR #992) is unchanged. M-12 (PR #994) is
unchanged.
This commit is contained in:
Anso
2026-05-08 15:51:39 -04:00
committed by GitHub
parent b5463e1771
commit f599110386
23 changed files with 496 additions and 2558 deletions
@@ -0,0 +1,166 @@
/**
* Unit tests for MeshForwarder. Exercises the per-port listener lifecycle
* (listen/unlisten/shutdown) and the accept dispatch into the host
* (MeshService surrogate). MeshForwarder itself does not splice bytes —
* the host's `handleAccept` does — so the test injects a recording
* surrogate and asserts the call shape.
*/
import net from 'net';
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
import { MeshForwarder, type MeshForwarderHost } from '../services/MeshForwarder';
interface Accept {
port: number;
socket: net.Socket;
}
function makeRecordingHost(): { host: MeshForwarderHost; accepts: Accept[]; resolveAfter: (cb: (a: Accept) => void) => void } {
const accepts: Accept[] = [];
const subscribers: Array<(a: Accept) => void> = [];
const host: MeshForwarderHost = {
async handleAccept(port, socket) {
const a = { port, socket };
accepts.push(a);
for (const s of subscribers.splice(0)) s(a);
},
};
return {
host,
accepts,
resolveAfter: (cb) => { subscribers.push(cb); },
};
}
async function getEphemeralPort(): Promise<number> {
return new Promise((resolve, reject) => {
const s = net.createServer();
s.unref();
s.listen(0, '127.0.0.1', () => {
const addr = s.address();
if (!addr || typeof addr === 'string') {
reject(new Error('no address'));
return;
}
const port = addr.port;
s.close(() => resolve(port));
});
});
}
async function dial(port: number, host = '127.0.0.1'): Promise<net.Socket> {
return new Promise((resolve, reject) => {
const s = net.createConnection({ host, port });
s.once('connect', () => resolve(s));
s.once('error', reject);
});
}
async function waitFor<T>(check: () => T | undefined, timeoutMs = 1000): Promise<T> {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
const v = check();
if (v !== undefined && v !== null && (Array.isArray(v) ? (v as unknown[]).length > 0 : true)) {
return v as T;
}
await new Promise((r) => setTimeout(r, 10));
}
throw new Error('timeout waiting for condition');
}
describe('MeshForwarder', () => {
let forwarder: MeshForwarder | null = null;
afterEach(async () => {
if (forwarder) await forwarder.shutdown();
forwarder = null;
});
beforeEach(() => { /* fresh per test */ });
it('binds a listener on the requested port and reports it via getListenerPorts', async () => {
const { host } = makeRecordingHost();
forwarder = new MeshForwarder(host);
const port = await getEphemeralPort();
await forwarder.listen(port);
expect(forwarder.getListenerPorts()).toEqual([port]);
expect(forwarder.isListening(port)).toBe(true);
});
it('hands accepted sockets to the host with the original destination port', async () => {
const { host, accepts } = makeRecordingHost();
forwarder = new MeshForwarder(host);
const port = await getEphemeralPort();
await forwarder.listen(port);
const client = await dial(port);
const seen = await waitFor(() => accepts.length ? accepts : undefined);
expect(seen[0].port).toBe(port);
expect(seen[0].socket).toBeInstanceOf(net.Socket);
client.destroy();
seen[0].socket.destroy();
});
it('listen is idempotent (repeated calls on the same port no-op)', async () => {
const { host } = makeRecordingHost();
forwarder = new MeshForwarder(host);
const port = await getEphemeralPort();
await forwarder.listen(port);
await forwarder.listen(port);
expect(forwarder.getListenerPorts()).toEqual([port]);
});
it('two concurrent listen() calls on the same port produce a single bind, not EADDRINUSE', async () => {
const { host } = makeRecordingHost();
forwarder = new MeshForwarder(host);
const port = await getEphemeralPort();
// Fire both calls before either resolves. Without the in-flight
// dedup map, the second call would race past the listeners.has()
// check and fail with EADDRINUSE on its own bind attempt.
const [a, b] = await Promise.all([forwarder.listen(port), forwarder.listen(port)]);
expect(a).toBeUndefined();
expect(b).toBeUndefined();
expect(forwarder.getListenerPorts()).toEqual([port]);
});
it('unlisten releases the port so a new listener can bind it', async () => {
const { host } = makeRecordingHost();
forwarder = new MeshForwarder(host);
const port = await getEphemeralPort();
await forwarder.listen(port);
await forwarder.unlisten(port);
expect(forwarder.isListening(port)).toBe(false);
// Verify the port is genuinely free by binding a fresh net.Server.
await new Promise<void>((resolve, reject) => {
const probe = net.createServer();
probe.once('error', reject);
probe.listen(port, '127.0.0.1', () => probe.close(() => resolve()));
});
});
it('rejects new connections after shutdown', async () => {
const { host } = makeRecordingHost();
forwarder = new MeshForwarder(host);
const port = await getEphemeralPort();
await forwarder.listen(port);
await forwarder.shutdown();
await expect(dial(port)).rejects.toThrow();
forwarder = null;
});
it('destroys the source socket when the host handler throws', async () => {
const host: MeshForwarderHost = {
async handleAccept() { throw new Error('host blew up'); },
};
forwarder = new MeshForwarder(host);
const port = await getEphemeralPort();
await forwarder.listen(port);
const client = await dial(port);
const closed = new Promise<void>((resolve) => client.once('close', () => resolve()));
await closed;
});
});
+2 -2
View File
@@ -97,7 +97,7 @@ describe('MeshService activity log', () => {
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: 'sidecar', level: 'info', type: 'route.resolve.ok', alias: 'x.y.z.sencho', message: 'c' });
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);
@@ -172,7 +172,7 @@ describe('MeshService.testUpstream tunnel-down path', () => {
const result = await svc.testUpstream('nonexistent.sencho', localNodeId);
expect(result.ok).toBe(false);
expect(result.where).toBe('sidecar');
expect(result.where).toBe('no_route');
expect(result.code).toBe('no_route');
});
});
+1 -15
View File
@@ -173,24 +173,10 @@ meshRouter.get('/nodes/:nodeId/diagnostic', async (req: Request, res: Response):
}
});
meshRouter.post('/nodes/:nodeId/sidecar/restart', async (req: Request, res: Response): Promise<void> => {
if (!requireAdmiral(req, res)) return;
if (!requireAdmin(req, res)) return;
const nodeId = Number.parseInt(req.params.nodeId as string, 10);
if (!Number.isFinite(nodeId)) { res.status(400).json({ error: 'Invalid node id' }); return; }
try {
await MeshService.getInstance().stopSidecar(nodeId);
await MeshService.getInstance().spawnSidecar(nodeId);
res.json({ ok: true });
} catch (err) {
res.status(500).json({ error: (err as Error).message });
}
});
meshRouter.get('/activity', (req: Request, res: Response): void => {
if (!requireAdmiral(req, res)) return;
const alias = typeof req.query.alias === 'string' ? req.query.alias : undefined;
const source = typeof req.query.source === 'string' ? (req.query.source as 'sidecar' | 'pilot' | 'mesh') : undefined;
const source = typeof req.query.source === 'string' ? (req.query.source as 'pilot' | 'mesh') : undefined;
const level = typeof req.query.level === 'string' ? (req.query.level as 'info' | 'warn' | 'error') : undefined;
const limit = typeof req.query.limit === 'string' ? Number.parseInt(req.query.limit, 10) : 200;
const events = MeshService.getInstance().getActivity({ alias, source, level, limit });
+101
View File
@@ -0,0 +1,101 @@
import net from 'net';
import { sanitizeForLog } from '../utils/safeLog';
/**
* In-process mesh TCP forwarder. Owns per-port `net.Server` listeners on the
* host network and delegates accepted sockets to the host (MeshService) for
* resolve + splice. Replaces the separate `saelix/sencho-mesh` sidecar
* container that previously did this job over a control WebSocket. The
* resolve step is now a sync map lookup rather than a round-trip, so
* MeshForwarder is just a thin lifecycle layer; all routing + splicing
* lives on MeshService.
*
* Sencho's container must run in `network_mode: host` (Linux) for the
* listeners to bind on the host's network where meshed containers'
* `extra_hosts: <alias>:host-gateway` entries point. Without host network
* mode, `net.createServer().listen(port)` lands inside the container's
* namespace and inbound traffic from peers never reaches it.
*/
export interface MeshForwarderHost {
/** Called on each accepted inbound socket. The host owns the splice
* lifecycle; MeshForwarder only manages listener boilerplate. */
handleAccept(port: number, source: net.Socket): Promise<void>;
}
export class MeshForwarder {
private readonly listeners = new Map<number, net.Server>();
/**
* In-flight `listen(port)` promises so concurrent callers race-safely
* deduplicate. Without this guard, two concurrent calls to listen on
* the same port would both pass the `listeners.has(port)` check (which
* is only populated after the listening event resolves) and the second
* would fail with EADDRINUSE.
*/
private readonly pending = new Map<number, Promise<void>>();
private shuttingDown = false;
constructor(private readonly host: MeshForwarderHost) {}
public async listen(port: number): Promise<void> {
if (this.shuttingDown) return;
if (this.listeners.has(port)) return;
const inflight = this.pending.get(port);
if (inflight) return inflight;
const promise = (async () => {
const server = net.createServer((socket) => this.acceptConnection(port, socket));
try {
await new Promise<void>((resolve, reject) => {
const onError = (err: Error) => { server.removeListener('listening', onListening); reject(err); };
const onListening = () => { server.removeListener('error', onError); resolve(); };
server.once('error', onError);
server.once('listening', onListening);
// Bind on all interfaces. Under host network mode this is
// the host's own network; under bridge mode (mesh disabled
// at boot) this would be the container's namespace.
server.listen(port, '0.0.0.0');
});
this.listeners.set(port, server);
} finally {
this.pending.delete(port);
}
})();
this.pending.set(port, promise);
return promise;
}
public async unlisten(port: number): Promise<void> {
const server = this.listeners.get(port);
if (!server) return;
this.listeners.delete(port);
await new Promise<void>((resolve) => server.close(() => resolve()));
}
public async shutdown(): Promise<void> {
this.shuttingDown = true;
const ports = Array.from(this.listeners.keys());
await Promise.all(ports.map((p) => this.unlisten(p)));
}
public getListenerPorts(): number[] {
return Array.from(this.listeners.keys());
}
public isListening(port: number): boolean {
return this.listeners.has(port);
}
private acceptConnection(port: number, source: net.Socket): void {
if (this.shuttingDown) {
try { source.destroy(); } catch { /* ignore */ }
return;
}
// Defer to the host for resolve + splice. MeshForwarder itself does
// not look at the source bytes; routing lives on MeshService where
// the alias map and the cross-node bridge dispatch are.
this.host.handleAccept(port, source).catch((err) => {
console.warn('[MeshForwarder] accept handler failed:', sanitizeForLog((err as Error).message));
try { source.destroy(); } catch { /* ignore */ }
});
}
}
+200 -197
View File
@@ -2,11 +2,11 @@ import net from 'net';
import path from 'path';
import fs from 'fs/promises';
import { EventEmitter } from 'events';
import jwt from 'jsonwebtoken';
import { DatabaseService } from './DatabaseService';
import DockerController from './DockerController';
import { LicenseService } from './LicenseService';
import { PROXY_TIER_HEADER, PROXY_VARIANT_HEADER } from './license-headers';
import { MeshForwarder, type MeshForwarderHost } from './MeshForwarder';
import { NodeRegistry } from './NodeRegistry';
import { PilotTunnelManager } from './PilotTunnelManager';
import { generateOverrideYaml, MeshAlias } from './MeshComposeOverride';
@@ -15,13 +15,10 @@ import { isPathWithinBase, isValidStackName } from '../utils/validation';
const ACTIVITY_BUFFER_SIZE = 1000;
const ALIAS_REFRESH_INTERVAL_MS = 60_000;
const SIDECAR_CONTAINER_PREFIX = 'sencho-mesh-';
const DEFAULT_SIDECAR_IMAGE = process.env.SENCHO_MESH_IMAGE || 'saelix/sencho-mesh:latest';
const SIDECAR_TOKEN_TTL = '7d';
const PROBE_TIMEOUT_MS = 5_000;
const SLOW_PROBE_THRESHOLD_MS = 500;
export type MeshActivitySource = 'sidecar' | 'pilot' | 'mesh';
export type MeshActivitySource = 'pilot' | 'mesh';
export type MeshActivityLevel = 'info' | 'warn' | 'error';
export type MeshActivityType =
| 'route.resolve.ok' | 'route.resolve.denied'
@@ -29,7 +26,7 @@ export type MeshActivityType =
| 'opt_in' | 'opt_out'
| 'mesh.enable' | 'mesh.disable'
| 'probe.ok' | 'probe.fail'
| 'sidecar.start' | 'sidecar.stop' | 'sidecar.crash';
| 'forwarder.listen' | 'forwarder.unlisten' | 'forwarder.error';
export interface MeshActivityEvent {
ts: number;
@@ -65,7 +62,8 @@ export interface MeshNodeStatus {
nodeId: number;
nodeName: string;
enabled: boolean;
sidecarRunning: boolean;
/** Forwarder state for the LOCAL node (the Sencho instance answering this request). Always `null` for any non-local node — fetching the remote forwarder state requires a cross-node call which lands in Phase B. */
localForwarderListening: boolean | null;
pilotConnected: boolean;
optedInStacks: string[];
activeStreamCount: number;
@@ -73,7 +71,7 @@ export interface MeshNodeStatus {
export interface MeshNodeDiagnostic {
nodeId: number;
sidecar: { running: boolean; restartCount: number };
forwarder: { listening: boolean; listenerCount: number };
pilot: { connected: boolean; bufferedAmount: number; lastSeen: number | null };
activeStreams: Array<{ streamId: number; alias?: string; bytesIn: number; bytesOut: number; ageMs: number }>;
aliasCache: Array<{ host: string; targetNodeId: number; port: number }>;
@@ -91,7 +89,7 @@ export interface MeshRouteDiagnostic {
export interface MeshProbeResult {
ok: boolean;
latencyMs?: number;
where?: 'sidecar' | 'pilot_tunnel' | 'agent_resolve' | 'agent_dial' | 'target_port';
where?: 'no_route' | 'pilot_tunnel' | 'agent_resolve' | 'agent_dial' | 'target_port';
code?: string;
message?: string;
}
@@ -104,50 +102,44 @@ interface ActiveStreamRecord {
openedAt: number;
}
interface PendingResolve {
sidecarSocket: WebSocketLike;
connId: number;
port: number;
remoteAddr: string;
}
interface WebSocketLike {
send(data: string | Buffer, opts?: unknown, cb?: (err?: Error) => void): void;
readyState: number;
on(event: string, listener: (...args: unknown[]) => void): unknown;
}
/**
* Sencho Mesh orchestrator. Owns:
* - sidecar lifecycle (Dockerode-spawned per-instance)
* - in-process TCP forwarder (`MeshForwarder`) that binds host-network
* listeners on alias ports. Replaces the prior separate sidecar
* container; one container per node now.
* - opt-in / opt-out persistence and cascading override regeneration
* - global alias aggregation (across the fleet via the existing API)
* - request-based resolution from sidecar control WS
* - cross-node TCP forwarding via PilotTunnelManager
* - global alias aggregation (across the fleet via the existing HTTP
* proxy chain — see `inspectStackServices`)
* - cross-node TCP forwarding via `PilotTunnelManager` (central-side)
* - probe + diagnostics + activity ring buffer
*
* V1 limitations:
* - one cross-node alias per TCP port across the fleet (port-collision check at opt-in)
* - sidecar runs in host network mode; aliases resolve via `host-gateway` extra_hosts
* - pilot-to-pilot mesh routing is not supported (only central <-> pilot)
* - one cross-node alias per TCP port across the fleet (port-collision
* check at opt-in)
* - aliases resolve via `host-gateway` extra_hosts; Sencho's container
* must run with `network_mode: host` for the forwarder's listeners to
* bind on the host's network where meshed containers' `host-gateway`
* entries point
* - cross-node mesh routing is central → pilot in this phase. Pilot →
* central and pilot ↔ pilot via central relay land in Phase B.
*/
export class MeshService extends EventEmitter {
export class MeshService extends EventEmitter implements MeshForwarderHost {
private static instance: MeshService;
private started = false;
private aliasCache = new Map<string, MeshGlobalAlias>();
private aliasByPort = new Map<number, MeshGlobalAlias>();
private activity: MeshActivityEvent[] = [];
private activeStreams = new Map<number, ActiveStreamRecord>();
private pendingResolves = new Map<string, PendingResolve>();
private sidecarSockets = new Set<WebSocketLike>();
private aliasRefreshTimer?: NodeJS.Timeout;
private routeErrorMap = new Map<string, { ts: number; message: string }>();
private routeLatencyMap = new Map<string, number>();
private activityListeners = new Set<(e: MeshActivityEvent) => void>();
private readonly forwarder: MeshForwarder;
private constructor() {
super();
this.setMaxListeners(50);
this.forwarder = new MeshForwarder(this);
}
public static getInstance(): MeshService {
@@ -167,10 +159,16 @@ export class MeshService extends EventEmitter {
}));
await this.refreshAliasCache();
await this.syncForwarderListeners();
this.aliasRefreshTimer = setInterval(() => {
void this.refreshAliasCache().catch((err) => {
console.warn('[MeshService] alias refresh failed:', sanitizeForLog((err as Error).message));
});
void (async () => {
try {
await this.refreshAliasCache();
await this.syncForwarderListeners();
} catch (err) {
console.warn('[MeshService] alias refresh failed:', sanitizeForLog((err as Error).message));
}
})();
}, ALIAS_REFRESH_INTERVAL_MS);
this.logActivity({
@@ -186,10 +184,49 @@ export class MeshService extends EventEmitter {
clearInterval(this.aliasRefreshTimer);
this.aliasRefreshTimer = undefined;
}
for (const ws of this.sidecarSockets) {
try { (ws as { close?: (code: number) => void }).close?.(1000); } catch { /* ignore */ }
await this.forwarder.shutdown();
}
/**
* Bind the forwarder's listeners to the local-owned alias ports and
* release any listeners no longer in the alias set. Called from
* `start`, after each `refreshAliasCache` tick, and after every
* opt-in / opt-out / disable on the local node so the bound port set
* follows the DB state.
*/
private async syncForwarderListeners(): Promise<void> {
const localNodeId = NodeRegistry.getInstance().getDefaultNodeId();
const wantPorts = new Set<number>();
for (const alias of this.aliasByPort.values()) {
if (alias.nodeId === localNodeId) wantPorts.add(alias.port);
}
const havePorts = new Set(this.forwarder.getListenerPorts());
for (const port of havePorts) {
if (!wantPorts.has(port)) {
await this.forwarder.unlisten(port);
this.logActivity({
source: 'mesh', level: 'info', type: 'forwarder.unlisten',
nodeId: localNodeId, message: `forwarder released port ${port}`,
});
}
}
for (const port of wantPorts) {
if (havePorts.has(port)) continue;
try {
await this.forwarder.listen(port);
this.logActivity({
source: 'mesh', level: 'info', type: 'forwarder.listen',
nodeId: localNodeId, message: `forwarder listening on port ${port}`,
});
} catch (err) {
this.logActivity({
source: 'mesh', level: 'error', type: 'forwarder.error',
nodeId: localNodeId,
message: `forwarder bind failed on port ${port}: ${sanitizeForLog((err as Error).message)}`,
details: { port },
});
}
}
this.sidecarSockets.clear();
}
// --- Activity log ---
@@ -248,6 +285,7 @@ export class MeshService extends EventEmitter {
db.insertMeshStack(nodeId, stackName, actor);
await this.refreshAliasCache();
await this.syncForwarderListeners();
await this.regenerateOverridesForNode(nodeId);
this.logActivity({
@@ -271,6 +309,7 @@ export class MeshService extends EventEmitter {
db.deleteMeshStack(nodeId, stackName);
await this.removeStackOverride(nodeId, stackName);
await this.refreshAliasCache();
await this.syncForwarderListeners();
await this.regenerateOverridesForNode(nodeId);
this.logActivity({
@@ -301,6 +340,7 @@ export class MeshService extends EventEmitter {
await this.removeStackOverride(nodeId, s.stack_name);
}
await this.refreshAliasCache();
await this.syncForwarderListeners();
this.logActivity({
source: 'mesh', level: 'info', type: 'mesh.disable',
nodeId, message: `mesh disabled on node ${nodeId}`,
@@ -483,25 +523,50 @@ export class MeshService extends EventEmitter {
}
/**
* Forward bytes from a sidecar-accepted local socket to the target.
* Same-node target: open a direct TCP socket.
* Cross-node target: open a pilot-tunnel TcpStream to that node's agent.
* MeshForwarder calls this on every accepted inbound socket. Resolves
* the alias by destination port, then dispatches to the same-node fast
* path or the cross-node bridge path.
*/
public openTcp(target: MeshTarget, src: net.Socket, sourceNodeId: number): void {
if (target.nodeId === sourceNodeId) {
this.openSameNode(target, src);
public async handleAccept(port: number, src: net.Socket): Promise<void> {
const target = this.resolveByLocalPort(port);
if (!target) {
this.logActivity({
source: 'mesh', level: 'warn', type: 'route.resolve.denied',
message: `inbound on port ${port} has no registered alias`,
details: { port, remoteAddr: src.remoteAddress ?? '' },
});
try { src.destroy(); } catch { /* ignore */ }
return;
}
this.openCrossNode(target, src);
const localNodeId = NodeRegistry.getInstance().getDefaultNodeId();
if (target.nodeId === localNodeId) {
await this.openSameNode(target, src);
} else {
this.openCrossNode(target, src);
}
}
private openSameNode(target: MeshTarget, src: net.Socket): void {
// For same-node fast path, the agent's resolution logic isn't needed:
// we can dial via Dockerode's container IP. For V1 simplicity, we dial
// the host-gateway's published port if mapped, falling back to
// 127.0.0.1 (the sidecar runs in host network mode so localhost reaches
// the container's published port).
const upstream = net.createConnection({ host: '127.0.0.1', port: target.port });
/**
* Same-node forward: dial the target container's bridge IP directly.
* Sencho runs in `network_mode: host` so it sees the docker bridge
* networks and can reach container IPs without going through any
* host-port publish. Looks up the container by Compose's
* `<project>-<service>-<index>` naming convention; falls back to a
* label-filtered listContainers if the conventional name is absent
* (e.g. when the operator overrode the project name).
*/
private async openSameNode(target: MeshTarget, src: net.Socket): Promise<void> {
const ip = await this.resolveContainerIp(target);
if (!ip) {
this.logActivity({
source: 'mesh', level: 'error', type: 'route.resolve.denied',
alias: target.alias,
message: `cannot resolve container IP for ${target.alias}`,
});
try { src.destroy(); } catch { /* ignore */ }
return;
}
const upstream = net.createConnection({ host: ip, port: target.port });
upstream.setTimeout(PROBE_TIMEOUT_MS);
const stream = this.registerActiveStream(target.alias);
upstream.once('connect', () => {
@@ -509,7 +574,7 @@ export class MeshService extends EventEmitter {
this.logActivity({
source: 'mesh', level: 'info', type: 'route.resolve.ok',
alias: target.alias, streamId: stream.streamId,
message: `same-node connect to ${target.alias}`,
message: `same-node connect to ${target.alias} (${ip}:${target.port})`,
});
src.pipe(upstream);
upstream.pipe(src);
@@ -525,6 +590,60 @@ export class MeshService extends EventEmitter {
src.on('close', () => teardown());
}
/** Find the bridge-network IP of the first container of `<stack>/<service>`. */
private async resolveContainerIp(target: MeshTarget): Promise<string | null> {
try {
const docker = DockerController.getInstance().getDocker();
// Compose default container name pattern; -1 is the first replica.
const conventionalName = `${target.stack}-${target.service}-1`;
const info = await docker.getContainer(conventionalName).inspect().catch(() => null);
const fromInspect = info ? this.extractContainerIp(target.stack, info) : null;
if (fromInspect) return fromInspect;
// Fallback: filter by compose labels in case of a non-conventional
// container name (operator overrode `container_name` or compose
// project).
const containers = await docker.listContainers({
all: true,
filters: {
label: [
`com.docker.compose.project=${target.stack}`,
`com.docker.compose.service=${target.service}`,
],
},
});
if (containers.length === 0) return null;
const fallbackInfo = await docker.getContainer(containers[0].Id).inspect().catch(() => null);
return fallbackInfo ? this.extractContainerIp(target.stack, fallbackInfo) : null;
} catch (err) {
console.warn('[MeshService] container IP lookup failed:', sanitizeForLog((err as Error).message));
return null;
}
}
/**
* Pick a deterministic IP. Prefer the compose default network
* (`<stack>_default` or any network whose name starts with `<stack>_`),
* then any other declared network, then the legacy bridge `IPAddress`.
* Without this preference order, `Object.values(Networks)` ordering on
* containers attached to multiple networks varies across daemon
* versions and can make same-node forwarding flaky on a redeploy.
*/
private extractContainerIp(
stackName: string,
info: { NetworkSettings?: { Networks?: Record<string, { IPAddress?: string }>; IPAddress?: string } },
): string | null {
const networks = info.NetworkSettings?.Networks ?? {};
const composeDefault = networks[`${stackName}_default`];
if (composeDefault?.IPAddress) return composeDefault.IPAddress;
for (const [name, net] of Object.entries(networks)) {
if (name.startsWith(`${stackName}_`) && net?.IPAddress) return net.IPAddress;
}
for (const net of Object.values(networks)) {
if (net?.IPAddress) return net.IPAddress;
}
return info.NetworkSettings?.IPAddress || null;
}
private openCrossNode(target: MeshTarget, src: net.Socket): void {
const ptm = PilotTunnelManager.getInstance();
if (!ptm.hasActiveTunnel(target.nodeId)) {
@@ -598,7 +717,7 @@ export class MeshService extends EventEmitter {
public async testUpstream(alias: string, sourceNodeId: number): Promise<MeshProbeResult> {
const target = this.lookupAliasGlobal(alias);
if (!target) {
return { ok: false, where: 'sidecar', code: 'no_route', message: 'alias not found' };
return { ok: false, where: 'no_route', code: 'no_route', message: 'alias not found' };
}
if (!DatabaseService.getInstance().isMeshStackEnabled(target.nodeId, target.stackName)) {
return { ok: false, where: 'agent_resolve', code: 'denied', message: 'target stack not opted in' };
@@ -726,7 +845,8 @@ export class MeshService extends EventEmitter {
const ptm = PilotTunnelManager.getInstance();
const bridge = ptm.getBridge(nodeId);
const node = DatabaseService.getInstance().getNode(nodeId);
const sidecarRunning = await this.isSidecarRunning(nodeId);
const localNodeId = NodeRegistry.getInstance().getDefaultNodeId();
const isLocal = nodeId === localNodeId;
const aliasCacheRows = Array.from(this.aliasCache.values())
.filter((a) => a.nodeId === nodeId)
@@ -739,9 +859,13 @@ export class MeshService extends EventEmitter {
ageMs: now - s.openedAt,
}));
const listenerCount = isLocal ? this.forwarder.getListenerPorts().length : 0;
return {
nodeId,
sidecar: { running: sidecarRunning, restartCount: 0 },
forwarder: {
listening: isLocal && this.isLocalForwarderActive(),
listenerCount,
},
pilot: {
connected: !!bridge,
bufferedAmount: bridge?.getBufferedAmount() ?? 0,
@@ -755,151 +879,30 @@ export class MeshService extends EventEmitter {
public async getStatus(): Promise<MeshNodeStatus[]> {
const db = DatabaseService.getInstance();
const nodes = db.getNodes();
const out: MeshNodeStatus[] = [];
for (const node of nodes) {
const optedInStacks = db.listMeshStacks(node.id).map((s) => s.stack_name);
out.push({
nodeId: node.id,
nodeName: node.name,
enabled: db.getNodeMeshEnabled(node.id),
sidecarRunning: await this.isSidecarRunning(node.id),
pilotConnected: this.isMeshReachable(node.id),
optedInStacks,
activeStreamCount: Array.from(this.activeStreams.values()).length,
});
}
return out;
const localNodeId = NodeRegistry.getInstance().getDefaultNodeId();
const localListening = this.isLocalForwarderActive();
return nodes.map((node) => ({
nodeId: node.id,
nodeName: node.name,
enabled: db.getNodeMeshEnabled(node.id),
localForwarderListening: node.id === localNodeId ? localListening : null,
pilotConnected: this.isMeshReachable(node.id),
optedInStacks: db.listMeshStacks(node.id).map((s) => s.stack_name),
activeStreamCount: this.activeStreams.size,
}));
}
// --- Sidecar lifecycle (best-effort; real spawn happens on local node) ---
public async spawnSidecar(nodeId: number): Promise<void> {
const docker = DockerController.getInstance(nodeId).getDocker();
const name = `${SIDECAR_CONTAINER_PREFIX}${nodeId}`;
try {
const existing = docker.getContainer(name);
const info = await existing.inspect().catch(() => null);
if (info?.State?.Running) return;
if (info) await existing.remove({ force: true }).catch(() => undefined);
} catch { /* ignore */ }
const token = this.mintSidecarToken(nodeId);
const controlUrl = process.env.SENCHO_INTERNAL_URL || 'ws://127.0.0.1:1852/api/mesh/control';
try {
const container = await docker.createContainer({
name,
Image: DEFAULT_SIDECAR_IMAGE,
Env: [
`SENCHO_CONTROL_URL=${controlUrl}`,
`SENCHO_MESH_TOKEN=${token}`,
`MESH_NODE_ID=${nodeId}`,
],
HostConfig: {
NetworkMode: 'host',
RestartPolicy: { Name: 'unless-stopped' },
},
Labels: {
'sencho.mesh.role': 'sidecar',
'sencho.mesh.node_id': String(nodeId),
},
});
await container.start();
this.logActivity({
source: 'mesh', level: 'info', type: 'sidecar.start',
nodeId, message: `sidecar started for node ${nodeId}`,
});
} catch (err) {
this.logActivity({
source: 'mesh', level: 'error', type: 'sidecar.crash',
nodeId, message: `sidecar spawn failed: ${(err as Error).message}`,
});
throw err;
}
/** True when the local Sencho's forwarder is started and bound to at least one alias port. */
private isLocalForwarderActive(): boolean {
return this.started && this.forwarder.getListenerPorts().length > 0;
}
public async stopSidecar(nodeId: number): Promise<void> {
const docker = DockerController.getInstance(nodeId).getDocker();
const name = `${SIDECAR_CONTAINER_PREFIX}${nodeId}`;
try {
const c = docker.getContainer(name);
await c.stop({ t: 5 }).catch(() => undefined);
await c.remove({ force: true }).catch(() => undefined);
this.logActivity({
source: 'mesh', level: 'info', type: 'sidecar.stop',
nodeId, message: `sidecar stopped for node ${nodeId}`,
});
} catch { /* ignore */ }
}
private async isSidecarRunning(nodeId: number): Promise<boolean> {
try {
const docker = DockerController.getInstance(nodeId).getDocker();
const info = await docker.getContainer(`${SIDECAR_CONTAINER_PREFIX}${nodeId}`).inspect();
return !!info.State?.Running;
} catch {
return false;
}
}
public mintSidecarToken(nodeId: number): string {
const settings = DatabaseService.getInstance().getGlobalSettings();
const secret = settings.auth_jwt_secret;
if (!secret) throw new Error('JWT secret not configured');
return jwt.sign({ scope: 'mesh_sidecar', nodeId }, secret, { expiresIn: SIDECAR_TOKEN_TTL });
}
public verifySidecarToken(token: string): { nodeId: number } | null {
try {
const settings = DatabaseService.getInstance().getGlobalSettings();
const secret = settings.auth_jwt_secret;
if (!secret) return null;
const decoded = jwt.verify(token, secret) as { scope?: string; nodeId?: number };
if (decoded.scope !== 'mesh_sidecar' || typeof decoded.nodeId !== 'number') return null;
return { nodeId: decoded.nodeId };
} catch {
return null;
}
}
// --- Sidecar control WS attachment (called from websocket/meshControl.ts) ---
public attachSidecarSocket(ws: WebSocketLike, _nodeId: number): void {
this.sidecarSockets.add(ws);
ws.on('close', () => { this.sidecarSockets.delete(ws); });
}
/** Resolve an inbound sidecar request: "I have a connection on this port; who's it for?" */
public handleSidecarResolve(ws: WebSocketLike, nodeId: number, connId: number, port: number, remoteAddr: string): void {
const target = this.resolveByLocalPort(port);
if (!target) {
this.sendSidecar(ws, { t: 'resolve_err', connId, code: 'no_route', message: 'port not registered' });
this.logActivity({
source: 'sidecar', level: 'warn', type: 'route.resolve.denied',
nodeId, message: `unknown port ${port}`, details: { connId, remoteAddr },
});
return;
}
// The sidecar is the SOURCE; it asks for routing on its local node.
// Open a TCP path on this node's MeshService (same-node fast path or
// pilot tunnel) and acknowledge the resolve with a freshly allocated
// streamId. For V1 we do NOT bridge real bytes through the control WS
// until the sidecar package gains binary frame plumbing (Phase B).
// Instead we ack the resolve with the target metadata so the sidecar
// can dial directly on the host gateway.
this.sendSidecar(ws, { t: 'resolve_ok', connId, streamId: connId, alias: target.alias });
this.logActivity({
source: 'sidecar', level: 'info', type: 'route.resolve.ok',
nodeId, alias: target.alias,
message: `resolved port ${port} to ${target.alias}`,
details: { connId, remoteAddr },
});
}
private sendSidecar(ws: WebSocketLike, frame: Record<string, unknown>): void {
if (ws.readyState !== 1 /* OPEN */) return;
try { ws.send(JSON.stringify(frame)); } catch { /* ignore */ }
}
// mintSidecarToken / verifySidecarToken / spawnSidecar / stopSidecar /
// isSidecarRunning / handleSidecarResolve / sendSidecar /
// attachSidecarSocket are gone: the in-process MeshForwarder replaces
// the entire sidecar layer. Routing decisions happen via direct
// MeshService calls — no JWT minting, no separate container, no control
// WebSocket. See `docs/internal/architecture/mesh.md` for the new flow.
}
export class MeshError extends Error {
-57
View File
@@ -1,57 +0,0 @@
import type { IncomingMessage } from 'http';
import type { Duplex } from 'stream';
import type { WebSocketServer, WebSocket } from 'ws';
import { MeshService } from '../services/MeshService';
import { sanitizeForLog } from '../utils/safeLog';
import { rejectUpgrade as rejectSocket } from './reject';
/**
* Handle the local Sencho Mesh sidecar's control WebSocket. Authenticated
* with a `mesh_sidecar`-scoped JWT minted by MeshService when it spawned the
* sidecar; the JWT carries the node id the sidecar serves.
*
* The control WS is intentionally local-only: the sidecar runs in host
* network mode on the same Docker host as Sencho and reaches us via the
* loopback interface.
*/
export async function handleMeshControl(
req: IncomingMessage,
socket: Duplex,
head: Buffer,
wss: WebSocketServer,
): Promise<void> {
const authHeader = req.headers['authorization'];
const header = Array.isArray(authHeader) ? authHeader[0] : authHeader;
const token = header?.startsWith('Bearer ') ? header.slice(7) : null;
if (!token) return rejectSocket(socket, 401, 'Unauthorized');
const verified = MeshService.getInstance().verifySidecarToken(token);
if (!verified) return rejectSocket(socket, 401, 'Unauthorized');
wss.handleUpgrade(req, socket as never, head, (ws: WebSocket) => {
MeshService.getInstance().attachSidecarSocket(ws as unknown as never, verified.nodeId);
ws.on('message', (data, isBinary) => {
if (isBinary) return; // V1: control plane is JSON-only.
try {
const text = data.toString('utf8');
const frame = JSON.parse(text) as { t?: string; connId?: number; port?: number; remoteAddr?: string };
if (frame.t === 'resolve' && typeof frame.connId === 'number' && typeof frame.port === 'number') {
MeshService.getInstance().handleSidecarResolve(
ws as unknown as never,
verified.nodeId,
frame.connId,
frame.port,
frame.remoteAddr ?? '',
);
}
// hello / log / stream.stats / close are advisory; we accept
// them silently in V1. Future revisions can expand handling.
} catch (err) {
console.warn('[meshControl] bad frame:', sanitizeForLog((err as Error).message));
}
});
ws.on('error', () => { try { ws.close(); } catch { /* ignore */ } });
});
}
+1 -7
View File
@@ -7,7 +7,6 @@ import { DatabaseService, type UserRole } from '../services/DatabaseService';
import { NodeRegistry } from '../services/NodeRegistry';
import { COOKIE_NAME } from '../helpers/constants';
import { handlePilotTunnel } from './pilotTunnel';
import { handleMeshControl } from './meshControl';
import { handleNotificationsWs } from './notifications';
import { handleRemoteForwarder } from './remoteForwarder';
import { handleLogsWs } from './logs';
@@ -31,8 +30,7 @@ function parseCookies(req: IncomingMessage): Record<string, string> {
*
* Dispatch order (first match wins):
* 1. `/api/pilot/tunnel` -> handlePilotTunnel (own auth, own wss)
* 2. `/api/mesh/control` -> handleMeshControl (sidecar JWT, local-only)
* 3. shared cookie/Bearer auth + JWT verify (rejects unauthenticated)
* 2. shared cookie/Bearer auth + JWT verify (rejects unauthenticated)
* 3. API token scope gate (read-only / deploy-only restricted to logs + notifications)
* 4. `/ws/notifications` local -> handleNotificationsWs
* 5. remote nodeId path -> handleRemoteForwarder
@@ -56,10 +54,6 @@ export function attachUpgrade(
await handlePilotTunnel(req, socket, head, pilotTunnelWss);
return;
}
if (reqUrl.pathname === '/api/mesh/control') {
await handleMeshControl(req, socket, head, wss);
return;
}
} catch {
// URL parse error falls through and will be rejected below.
}
+6 -2
View File
@@ -4,8 +4,12 @@ services:
build: .
container_name: sencho
restart: unless-stopped
ports:
- "1852:1852"
# Sencho Mesh listens on host ports for cross-stack alias traffic. Host
# network mode is required for mesh to function. If you do not use mesh,
# comment `network_mode: host` out and uncomment the `ports` block below.
network_mode: host
# ports:
# - "1852:1852"
volumes:
# Required: Docker Socket for container orchestration
- /var/run/docker.sock:/var/run/docker.sock
@@ -1,8 +1,7 @@
import { useEffect, useState } from 'react';
import { apiFetch } from '@/lib/api';
import { toast } from '@/components/ui/toast-store';
import { SystemSheet, SheetSection } from '@/components/ui/system-sheet';
import { RefreshCw, ServerCog } from 'lucide-react';
import { RefreshCw } from 'lucide-react';
import { formatTimeAgo } from '@/lib/relativeTime';
import type { MeshNodeDiagnostic } from '@/types/mesh';
@@ -28,7 +27,6 @@ function ageFmt(ms: number): string {
export function MeshDiagnosticsSheet({ open, onOpenChange, nodeId, nodeName }: Props) {
const [diag, setDiag] = useState<MeshNodeDiagnostic | null>(null);
const [loading, setLoading] = useState(false);
const [restarting, setRestarting] = useState(false);
const [updatedAt, setUpdatedAt] = useState<number | null>(null);
const refresh = async () => {
@@ -50,29 +48,13 @@ export function MeshDiagnosticsSheet({ open, onOpenChange, nodeId, nodeName }: P
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [open, nodeId]);
const restart = async () => {
if (nodeId == null) return;
setRestarting(true);
try {
const res = await apiFetch(`/mesh/nodes/${nodeId}/sidecar/restart`, {
method: 'POST', localOnly: true,
});
if (res.ok) {
toast.success('Sidecar restart requested');
await refresh();
} else {
toast.error('Sidecar restart failed');
}
} finally {
setRestarting(false);
}
};
const sidecarLabel = diag ? (diag.sidecar.running ? 'sidecar running' : 'sidecar off') : 'sidecar ?';
const forwarderLabel = diag
? (diag.forwarder.listening ? `forwarder listening (${diag.forwarder.listenerCount})` : 'forwarder idle')
: 'forwarder ?';
const pilotLabel = diag ? (diag.pilot.connected ? 'pilot connected' : 'pilot disconnected') : 'pilot ?';
const streamsLabel = `${diag?.activeStreams.length ?? 0} streams`;
const aliasesLabel = `${diag?.aliasCache.length ?? 0} aliases`;
const meta = `${sidecarLabel} · ${pilotLabel} · ${streamsLabel} · ${aliasesLabel}`;
const meta = `${forwarderLabel} · ${pilotLabel} · ${streamsLabel} · ${aliasesLabel}`;
const footerContext = updatedAt ? `Updated ${formatTimeAgo(updatedAt)}` : (loading ? 'Loading…' : 'Never updated');
@@ -89,21 +71,15 @@ export function MeshDiagnosticsSheet({ open, onOpenChange, nodeId, nodeName }: P
onClick: () => { void refresh(); },
disabled: loading,
}}
secondaryActions={[
{
label: 'Restart sidecar',
icon: ServerCog,
onClick: () => { void restart(); },
disabled: restarting,
},
]}
footerContext={footerContext}
size="md"
>
<SheetSection title="Pilot · sidecar · transport">
<SheetSection title="Forwarder · pilot · transport">
<div className="grid grid-cols-2 gap-x-4 gap-y-1.5 text-xs">
<div className="text-stat-subtitle">Sidecar</div>
<div className="font-mono text-stat-value">{diag?.sidecar.running ? 'running' : 'off'}</div>
<div className="text-stat-subtitle">Forwarder</div>
<div className="font-mono text-stat-value">
{diag?.forwarder.listening ? `listening on ${diag.forwarder.listenerCount} port${diag.forwarder.listenerCount === 1 ? '' : 's'}` : 'idle'}
</div>
<div className="text-stat-subtitle">Pilot tunnel</div>
<div className="font-mono text-stat-value">{diag?.pilot.connected ? 'connected' : 'disconnected'}</div>
<div className="text-stat-subtitle">Buffered</div>
@@ -111,6 +87,9 @@ export function MeshDiagnosticsSheet({ open, onOpenChange, nodeId, nodeName }: P
<div className="text-stat-subtitle">Last seen</div>
<div className="font-mono text-stat-value">{diag?.pilot.lastSeen ? new Date(diag.pilot.lastSeen).toLocaleTimeString() : '-'}</div>
</div>
<p className="text-[11px] text-stat-subtitle leading-snug mt-2">
Mesh runs in-process on each node; no separate container.
</p>
</SheetSection>
<SheetSection title="Active streams">
@@ -2,7 +2,7 @@ import { useEffect, useState } from 'react';
import { apiFetch } from '@/lib/api';
import { SystemSheet, SheetSection } from '@/components/ui/system-sheet';
import { Badge } from '@/components/ui/badge';
import { Loader2, Activity, ServerCog, Hash } from 'lucide-react';
import { Loader2, Activity, Hash } from 'lucide-react';
import type { MeshRouteDiagnostic, MeshActivityEvent, MeshProbeResult } from '@/types/mesh';
import { meshRouteStateFromBackend, meshRouteStateTokens } from './meshRouteState';
@@ -146,7 +146,6 @@ export function MeshRouteDetailSheet({ open, onOpenChange, alias }: Props) {
)}
{events.map((e, i) => (
<div key={i} className="flex items-start gap-2 text-[11px] font-mono">
{e.source === 'sidecar' && <ServerCog className="w-3 h-3 mt-0.5 text-stat-subtitle" />}
{e.source === 'pilot' && <Hash className="w-3 h-3 mt-0.5 text-stat-subtitle" />}
{e.source === 'mesh' && <Activity className="w-3 h-3 mt-0.5 text-stat-subtitle" />}
<span className={`tabular-nums ${e.level === 'error' ? 'text-destructive' : e.level === 'warn' ? 'text-warning' : 'text-stat-value'}`}>
+5 -4
View File
@@ -13,7 +13,8 @@ export interface MeshNodeStatus {
nodeId: number;
nodeName: string;
enabled: boolean;
sidecarRunning: boolean;
/** Forwarder state for the LOCAL Sencho instance only. `null` for non-local nodes — cross-node forwarder state lands in Phase B. */
localForwarderListening: boolean | null;
pilotConnected: boolean;
optedInStacks: string[];
activeStreamCount: number;
@@ -36,7 +37,7 @@ export interface MeshRouteDiagnostic {
export interface MeshNodeDiagnostic {
nodeId: number;
sidecar: { running: boolean; restartCount: number };
forwarder: { listening: boolean; listenerCount: number };
pilot: { connected: boolean; bufferedAmount: number; lastSeen: number | null };
activeStreams: Array<{ streamId: number; alias?: string; bytesIn: number; bytesOut: number; ageMs: number }>;
aliasCache: Array<{ host: string; targetNodeId: number; port: number }>;
@@ -45,14 +46,14 @@ export interface MeshNodeDiagnostic {
export interface MeshProbeResult {
ok: boolean;
latencyMs?: number;
where?: 'sidecar' | 'pilot_tunnel' | 'agent_resolve' | 'agent_dial' | 'target_port';
where?: 'no_route' | 'pilot_tunnel' | 'agent_resolve' | 'agent_dial' | 'target_port';
code?: string;
message?: string;
}
export interface MeshActivityEvent {
ts: number;
source: 'sidecar' | 'pilot' | 'mesh';
source: 'pilot' | 'mesh';
level: 'info' | 'warn' | 'error';
type: string;
nodeId?: number;
-3
View File
@@ -1,3 +0,0 @@
node_modules
dist
*.tsbuildinfo
-30
View File
@@ -1,30 +0,0 @@
# Multi-stage build for the Sencho Mesh sidecar.
#
# Base image is pinned by digest, not tag, so a future republish of the
# `22-alpine` tag cannot silently shift the runtime under a release. Update
# the digest in the same PR that rolls the Node minor or upgrades from
# alpine. Resolve a fresh digest with:
# docker buildx imagetools inspect node:22-alpine
ARG NODE_BASE=node:22-alpine@sha256:8ea2348b068a9544dae7317b4f3aafcdc032df1647bb7d768a05a5cad1a7683f
FROM ${NODE_BASE} AS build
WORKDIR /app
COPY package.json tsconfig.json ./
RUN npm install --no-audit --no-fund
COPY src ./src
RUN npm run build
FROM ${NODE_BASE} AS runtime
WORKDIR /app
ENV NODE_ENV=production
COPY package.json ./
RUN npm install --omit=dev --no-audit --no-fund
COPY --from=build /app/dist ./dist
# No HEALTHCHECK: the sidecar has no inbound HTTP listener. It maintains
# an outbound websocket to the control plane and exits on connection loss
# so Docker's restart policy is the natural recovery loop. A liveness
# probe based on the node process being PID 1 would be tautological.
USER node
CMD ["node", "dist/index.js"]
-1333
View File
File diff suppressed because it is too large Load Diff
-23
View File
@@ -1,23 +0,0 @@
{
"name": "@sencho/mesh-sidecar",
"version": "0.0.1",
"private": true,
"description": "Sencho Mesh sidecar: per-node TCP forwarder + control WS bridge.",
"main": "dist/index.js",
"type": "commonjs",
"scripts": {
"build": "tsc",
"start": "node dist/index.js",
"test": "vitest run",
"lint": "eslint src"
},
"dependencies": {
"ws": "^8.19.0"
},
"devDependencies": {
"@types/node": "^25.3.0",
"@types/ws": "^8.18.1",
"typescript": "^6.0.2",
"vitest": "^4.1.0"
}
}
@@ -1,190 +0,0 @@
import net from 'net';
import { afterEach, describe, expect, it } from 'vitest';
import { Forwarder, MeshController } from '../forwarder';
interface RecordedCalls {
resolve: Array<{ connId: number; port: number; remoteAddr: string }>;
sendData: Array<{ streamId: number; payload: Buffer }>;
sendClose: number[];
sendStats: Array<{ streamId: number; bytesIn: number; bytesOut: number }>;
}
function makeRecordingController(): { controller: MeshController; calls: RecordedCalls } {
const calls: RecordedCalls = { resolve: [], sendData: [], sendClose: [], sendStats: [] };
const controller: MeshController = {
resolve: (connId, port, remoteAddr) => calls.resolve.push({ connId, port, remoteAddr }),
sendData: (streamId, payload) => calls.sendData.push({ streamId, payload: Buffer.from(payload) }),
sendClose: (streamId) => calls.sendClose.push(streamId),
sendStats: (streamId, bytesIn, bytesOut) => calls.sendStats.push({ streamId, bytesIn, bytesOut }),
};
return { controller, calls };
}
async function getEphemeralPort(): Promise<number> {
return new Promise((resolve, reject) => {
const server = net.createServer();
server.unref();
server.listen(0, '127.0.0.1', () => {
const addr = server.address();
if (!addr || typeof addr === 'string') {
reject(new Error('no address'));
return;
}
const port = addr.port;
server.close(() => resolve(port));
});
});
}
async function dial(port: number): Promise<net.Socket> {
return new Promise((resolve, reject) => {
const socket = net.createConnection({ host: '127.0.0.1', port });
socket.once('connect', () => resolve(socket));
socket.once('error', reject);
});
}
async function waitFor<T>(check: () => T | undefined, timeoutMs = 1000): Promise<T> {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
const value = check();
if (value !== undefined && value !== null && (Array.isArray(value) ? value.length > 0 : true)) {
return value as T;
}
await new Promise((r) => setTimeout(r, 10));
}
throw new Error('timeout waiting for condition');
}
describe('Forwarder', () => {
let forwarder: Forwarder | null = null;
afterEach(async () => {
if (forwarder) await forwarder.shutdown();
forwarder = null;
});
it('asks the controller to resolve when a client connects', async () => {
const { controller, calls } = makeRecordingController();
forwarder = new Forwarder(controller);
forwarder.start();
const port = await getEphemeralPort();
await forwarder.listen(port);
const client = await dial(port);
const resolved = await waitFor(() => calls.resolve.length ? calls.resolve : undefined);
expect(resolved[0].port).toBe(port);
expect(resolved[0].connId).toBeGreaterThan(0);
client.destroy();
});
it('splices bytes both directions after resolve_ok', async () => {
const { controller, calls } = makeRecordingController();
forwarder = new Forwarder(controller);
forwarder.start();
const port = await getEphemeralPort();
await forwarder.listen(port);
const received: Buffer[] = [];
const client = await dial(port);
client.on('data', (chunk: Buffer) => received.push(chunk));
const resolved = await waitFor(() => calls.resolve.length ? calls.resolve : undefined);
const connId = resolved[0].connId;
const streamId = 42;
forwarder.handleResolveOk(connId, streamId, 'test.alias.sencho');
client.write('hello');
const sent = await waitFor(() => calls.sendData.length ? calls.sendData : undefined);
expect(sent[0].streamId).toBe(streamId);
expect(sent[0].payload.toString()).toBe('hello');
forwarder.handleData(streamId, Buffer.from('hi back'));
await waitFor(() => received.length ? received : undefined);
expect(Buffer.concat(received).toString()).toBe('hi back');
client.destroy();
});
it('sends a close frame when the client socket disconnects', async () => {
const { controller, calls } = makeRecordingController();
forwarder = new Forwarder(controller);
forwarder.start();
const port = await getEphemeralPort();
await forwarder.listen(port);
const client = await dial(port);
const resolved = await waitFor(() => calls.resolve.length ? calls.resolve : undefined);
forwarder.handleResolveOk(resolved[0].connId, 7);
client.destroy();
const closed = await waitFor(() => calls.sendClose.length ? calls.sendClose : undefined);
expect(closed[0]).toBe(7);
});
it('drops the local socket on resolve_err', async () => {
const { controller, calls } = makeRecordingController();
forwarder = new Forwarder(controller);
forwarder.start();
const port = await getEphemeralPort();
await forwarder.listen(port);
const client = await dial(port);
const closedPromise = new Promise<void>((resolve) => client.once('close', () => resolve()));
const resolved = await waitFor(() => calls.resolve.length ? calls.resolve : undefined);
forwarder.handleResolveErr(resolved[0].connId, 'tunnel_down');
await closedPromise;
// No data was ever piped, so no close frame sent for this connId.
expect(calls.sendClose.length).toBe(0);
});
it('destroys the local socket on inbound close', async () => {
const { controller, calls } = makeRecordingController();
forwarder = new Forwarder(controller);
forwarder.start();
const port = await getEphemeralPort();
await forwarder.listen(port);
const client = await dial(port);
const closedPromise = new Promise<void>((resolve) => client.once('close', () => resolve()));
const resolved = await waitFor(() => calls.resolve.length ? calls.resolve : undefined);
forwarder.handleResolveOk(resolved[0].connId, 11);
forwarder.handleClose(11);
await closedPromise;
});
it('refuses new connections during shutdown', async () => {
const { controller } = makeRecordingController();
forwarder = new Forwarder(controller);
forwarder.start();
const port = await getEphemeralPort();
await forwarder.listen(port);
await forwarder.shutdown();
await expect(dial(port)).rejects.toThrow();
forwarder = null;
});
it('emits stats for active streams when the timer fires', async () => {
const { controller, calls } = makeRecordingController();
forwarder = new Forwarder(controller, { statsIntervalMs: 50 });
forwarder.start();
const port = await getEphemeralPort();
await forwarder.listen(port);
const client = await dial(port);
const resolved = await waitFor(() => calls.resolve.length ? calls.resolve : undefined);
forwarder.handleResolveOk(resolved[0].connId, 99);
client.write('x');
const stats = await waitFor(() => calls.sendStats.length ? calls.sendStats : undefined, 1500);
expect(stats[0].streamId).toBe(99);
expect(stats[0].bytesOut).toBeGreaterThan(0);
client.destroy();
});
});
@@ -1,73 +0,0 @@
import { describe, it, expect } from 'vitest';
import {
BinaryFrameType,
decodeControl,
decodeData,
encodeControl,
encodeData,
} from '../protocol';
describe('Mesh sidecar control frames', () => {
it('roundtrips a hello frame', () => {
const raw = encodeControl({ t: 'hello', version: 1, nodeId: 7, sidecarVersion: '0.0.1' });
const decoded = decodeControl(raw);
expect(decoded.t).toBe('hello');
if (decoded.t !== 'hello') throw new Error('narrowing');
expect(decoded.nodeId).toBe(7);
});
it('roundtrips a listen / unlisten pair', () => {
const listen = decodeControl(encodeControl({ t: 'listen', port: 5432 }));
const unlisten = decodeControl(encodeControl({ t: 'unlisten', port: 5432 }));
expect(listen.t).toBe('listen');
expect(unlisten.t).toBe('unlisten');
});
it('roundtrips resolve / resolve_ok / resolve_err', () => {
const resolve = decodeControl(encodeControl({ t: 'resolve', connId: 1, port: 5432, remoteAddr: '10.0.0.5' }));
const ok = decodeControl(encodeControl({ t: 'resolve_ok', connId: 1, streamId: 9, alias: 'db.api.opsix.sencho' }));
const err = decodeControl(encodeControl({ t: 'resolve_err', connId: 1, code: 'tunnel_down' }));
if (resolve.t !== 'resolve') throw new Error('narrowing');
if (ok.t !== 'resolve_ok') throw new Error('narrowing');
if (err.t !== 'resolve_err') throw new Error('narrowing');
expect(resolve.port).toBe(5432);
expect(ok.streamId).toBe(9);
expect(err.code).toBe('tunnel_down');
});
it('rejects malformed control frames', () => {
expect(() => decodeControl('not json')).toThrow();
expect(() => decodeControl('{}')).toThrow();
});
});
describe('Mesh sidecar data frames', () => {
it('encodes the 0x01 type discriminator', () => {
const buf = encodeData(1, Buffer.from('hello'));
expect(buf[0]).toBe(BinaryFrameType.Data);
});
it('roundtrips streamId + payload', () => {
const payload = Buffer.from('SELECT 1;');
const buf = encodeData(0xdeadbeef, payload);
const decoded = decodeData(buf);
expect(decoded.streamId).toBe(0xdeadbeef);
expect(decoded.payload.toString()).toBe('SELECT 1;');
});
it('preserves binary payloads byte-for-byte', () => {
const payload = Buffer.from([0x00, 0xff, 0x01, 0x80, 0x7f]);
const decoded = decodeData(encodeData(1, payload));
expect(decoded.payload.equals(payload)).toBe(true);
});
it('rejects an unknown binary frame type', () => {
const buf = Buffer.alloc(5);
buf.writeUInt8(0x99, 0);
expect(() => decodeData(buf)).toThrow(/unknown binary frame type/);
});
it('rejects too-short frames', () => {
expect(() => decodeData(Buffer.alloc(3))).toThrow(/too short/);
});
});
-175
View File
@@ -1,175 +0,0 @@
import WebSocket from 'ws';
import { Forwarder, MeshController } from './forwarder';
import {
ControlFrame,
PROTOCOL_VERSION,
decodeControl,
decodeData,
encodeControl,
encodeData,
wsDataToBuffer,
wsDataToString,
} from './protocol';
const RECONNECT_MIN_MS = 1_000;
const RECONNECT_MAX_MS = 30_000;
const PING_INTERVAL_MS = 30_000;
export interface ControlClientOptions {
controlUrl: string;
token: string;
nodeId: number;
sidecarVersion: string;
forwarder: Forwarder;
}
/**
* WS-backed implementation of MeshController. Keeps a single long-lived
* connection to the local Sencho instance, surfaces inbound frames to the
* Forwarder, and queues outbound frames when reconnecting.
*/
export class ControlClient implements MeshController {
private readonly options: ControlClientOptions;
private ws: WebSocket | null = null;
private backoff = RECONNECT_MIN_MS;
private pingTimer?: NodeJS.Timeout;
private reconnectTimer?: NodeJS.Timeout;
private shuttingDown = false;
constructor(options: ControlClientOptions) {
this.options = options;
}
public start(): void {
this.connect();
}
public async shutdown(): Promise<void> {
this.shuttingDown = true;
if (this.pingTimer) clearInterval(this.pingTimer);
if (this.reconnectTimer) { clearTimeout(this.reconnectTimer); this.reconnectTimer = undefined; }
try { this.ws?.close(1000, 'sidecar shutdown'); } catch { /* ignore */ }
}
// --- MeshController surface ---
public resolve(connId: number, port: number, remoteAddr: string): void {
this.send({ t: 'resolve', connId, port, remoteAddr });
}
public sendData(streamId: number, payload: Buffer): void {
const ws = this.ws;
if (!ws || ws.readyState !== WebSocket.OPEN) return;
try { ws.send(encodeData(streamId, payload), { binary: true }); } catch { /* ignore */ }
}
public sendClose(streamId: number): void {
this.send({ t: 'close', streamId });
}
public sendStats(streamId: number, bytesIn: number, bytesOut: number, lastActivity: number): void {
this.send({ t: 'stream.stats', streamId, bytesIn, bytesOut, lastActivity });
}
public sendLog(level: 'info' | 'warn' | 'error', message: string, details?: Record<string, unknown>): void {
this.send({ t: 'log', level, message, details });
}
// --- Connection lifecycle ---
private connect(): void {
if (this.shuttingDown) return;
if (this.pingTimer) { clearInterval(this.pingTimer); this.pingTimer = undefined; }
if (this.reconnectTimer) { clearTimeout(this.reconnectTimer); this.reconnectTimer = undefined; }
const ws = new WebSocket(this.options.controlUrl, {
headers: { Authorization: `Bearer ${this.options.token}` },
handshakeTimeout: 15_000,
});
this.ws = ws;
ws.on('open', () => {
this.backoff = RECONNECT_MIN_MS;
this.send({
t: 'hello',
version: PROTOCOL_VERSION,
nodeId: this.options.nodeId,
sidecarVersion: this.options.sidecarVersion,
});
this.pingTimer = setInterval(() => {
if (ws.readyState === WebSocket.OPEN) {
try { ws.ping(); } catch { /* error event will handle */ }
}
}, PING_INTERVAL_MS);
});
ws.on('message', (data, isBinary) => this.onMessage(data, isBinary));
ws.on('close', () => {
if (this.pingTimer) { clearInterval(this.pingTimer); this.pingTimer = undefined; }
this.scheduleReconnect();
});
ws.on('error', () => {
// 'close' will follow.
});
}
private scheduleReconnect(): void {
if (this.shuttingDown) return;
const jitter = Math.floor(Math.random() * 500);
const delay = this.backoff + jitter;
this.reconnectTimer = setTimeout(() => {
this.reconnectTimer = undefined;
this.connect();
}, delay);
this.backoff = Math.min(this.backoff * 2, RECONNECT_MAX_MS);
}
private onMessage(data: WebSocket.RawData, isBinary: boolean): void {
try {
if (isBinary) {
const buf = wsDataToBuffer(data);
if (!buf) return;
const decoded = decodeData(buf);
this.options.forwarder.handleData(decoded.streamId, decoded.payload);
} else {
const text = wsDataToString(data);
if (text == null) return;
const frame = decodeControl(text);
this.dispatchControl(frame);
}
} catch {
// Malformed frames are dropped; tunnel stays up.
}
}
private dispatchControl(frame: ControlFrame): void {
switch (frame.t) {
case 'listen':
void this.options.forwarder.listen(frame.port);
break;
case 'unlisten':
void this.options.forwarder.unlisten(frame.port);
break;
case 'resolve_ok':
this.options.forwarder.handleResolveOk(frame.connId, frame.streamId, frame.alias);
break;
case 'resolve_err':
this.options.forwarder.handleResolveErr(frame.connId, frame.code);
break;
case 'close':
this.options.forwarder.handleClose(frame.streamId);
break;
default:
// hello / resolve / stream.stats / log are sidecar-originated;
// ignore on the inbound path.
break;
}
}
private send(frame: ControlFrame): void {
const ws = this.ws;
if (!ws || ws.readyState !== WebSocket.OPEN) return;
try { ws.send(encodeControl(frame)); } catch { /* ignore */ }
}
}
-180
View File
@@ -1,180 +0,0 @@
import net from 'net';
const DEFAULT_STATS_INTERVAL_MS = 5_000;
export interface ForwarderOptions {
statsIntervalMs?: number;
}
/**
* Outbound channel to the Sencho instance hosting this sidecar. The
* production implementation is the control WS client; tests inject a fake.
*/
export interface MeshController {
resolve(connId: number, port: number, remoteAddr: string): void;
sendData(streamId: number, payload: Buffer): void;
sendClose(streamId: number): void;
sendStats(streamId: number, bytesIn: number, bytesOut: number, lastActivity: number): void;
}
interface PendingConn {
connId: number;
socket: net.Socket;
port: number;
}
interface ActiveStream {
streamId: number;
socket: net.Socket;
bytesIn: number;
bytesOut: number;
lastActivity: number;
alias?: string;
}
/**
* The Forwarder owns per-port TCP listeners and the live socket <-> stream
* mapping. It is wire-agnostic: all outbound frames flow through MeshController
* methods. Inbound frames are delivered via the handle* methods.
*/
export class Forwarder {
private readonly controller: MeshController;
private readonly statsIntervalMs: number;
private readonly listeners = new Map<number, net.Server>();
private readonly pendingConns = new Map<number, PendingConn>();
private readonly activeStreams = new Map<number, ActiveStream>();
private nextConnId = 1;
private statsTimer?: NodeJS.Timeout;
private shuttingDown = false;
constructor(controller: MeshController, options: ForwarderOptions = {}) {
this.controller = controller;
this.statsIntervalMs = options.statsIntervalMs ?? DEFAULT_STATS_INTERVAL_MS;
}
public start(): void {
this.statsTimer = setInterval(() => this.flushStats(), this.statsIntervalMs);
}
public async listen(port: number): Promise<void> {
if (this.listeners.has(port) || this.shuttingDown) return;
const server = net.createServer((socket) => this.acceptConnection(port, socket));
await new Promise<void>((resolve, reject) => {
const onError = (err: Error) => { server.removeListener('listening', onListening); reject(err); };
const onListening = () => { server.removeListener('error', onError); resolve(); };
server.once('error', onError);
server.once('listening', onListening);
server.listen(port);
});
this.listeners.set(port, server);
}
public async unlisten(port: number): Promise<void> {
const server = this.listeners.get(port);
if (!server) return;
this.listeners.delete(port);
await new Promise<void>((resolve) => {
server.close(() => resolve());
});
}
public handleResolveOk(connId: number, streamId: number, alias?: string): void {
const pending = this.pendingConns.get(connId);
if (!pending) {
// Sencho thinks we have a conn but we don't (race with socket close).
this.controller.sendClose(streamId);
return;
}
this.pendingConns.delete(connId);
const stream: ActiveStream = {
streamId,
socket: pending.socket,
bytesIn: 0,
bytesOut: 0,
lastActivity: Date.now(),
alias,
};
this.activeStreams.set(streamId, stream);
pending.socket.on('data', (chunk: Buffer) => {
stream.bytesOut += chunk.length;
stream.lastActivity = Date.now();
this.controller.sendData(streamId, chunk);
});
pending.socket.on('close', () => {
if (this.activeStreams.delete(streamId)) {
this.controller.sendClose(streamId);
}
});
pending.socket.on('error', () => {
if (this.activeStreams.delete(streamId)) {
this.controller.sendClose(streamId);
}
});
}
public handleResolveErr(connId: number, _code: string): void {
const pending = this.pendingConns.get(connId);
if (!pending) return;
this.pendingConns.delete(connId);
try { pending.socket.destroy(); } catch { /* ignore */ }
}
public handleData(streamId: number, payload: Buffer): void {
const stream = this.activeStreams.get(streamId);
if (!stream) return;
stream.bytesIn += payload.length;
stream.lastActivity = Date.now();
try { stream.socket.write(payload); } catch { /* ignore */ }
}
public handleClose(streamId: number): void {
const stream = this.activeStreams.get(streamId);
if (!stream) return;
this.activeStreams.delete(streamId);
try { stream.socket.destroy(); } catch { /* ignore */ }
}
public async shutdown(): Promise<void> {
this.shuttingDown = true;
if (this.statsTimer) { clearInterval(this.statsTimer); this.statsTimer = undefined; }
for (const [, stream] of this.activeStreams) {
try { stream.socket.destroy(); } catch { /* ignore */ }
}
this.activeStreams.clear();
for (const [, pending] of this.pendingConns) {
try { pending.socket.destroy(); } catch { /* ignore */ }
}
this.pendingConns.clear();
const ports = Array.from(this.listeners.keys());
await Promise.all(ports.map((p) => this.unlisten(p)));
}
public getActiveStreamCount(): number { return this.activeStreams.size; }
public getListenerPorts(): number[] { return Array.from(this.listeners.keys()); }
private acceptConnection(port: number, socket: net.Socket): void {
if (this.shuttingDown) {
try { socket.destroy(); } catch { /* ignore */ }
return;
}
const connId = this.nextConnId++;
this.pendingConns.set(connId, { connId, socket, port });
const remoteAddr = socket.remoteAddress ?? '';
socket.once('error', () => { this.pendingConns.delete(connId); });
socket.once('close', () => { this.pendingConns.delete(connId); });
this.controller.resolve(connId, port, remoteAddr);
}
private flushStats(): void {
const now = Date.now();
for (const [streamId, stream] of this.activeStreams) {
// Only emit for streams that saw activity in the last interval to
// keep the activity log readable.
if (now - stream.lastActivity <= this.statsIntervalMs * 2) {
this.controller.sendStats(streamId, stream.bytesIn, stream.bytesOut, stream.lastActivity);
}
}
}
}
-58
View File
@@ -1,58 +0,0 @@
import { ControlClient } from './control';
import { Forwarder } from './forwarder';
const SIDECAR_VERSION = '0.0.1';
function requireEnv(name: string): string {
const value = process.env[name];
if (!value) {
console.error(`[mesh-sidecar] missing required env: ${name}`);
process.exit(1);
}
return value;
}
function main(): void {
const controlUrl = requireEnv('SENCHO_CONTROL_URL');
const token = requireEnv('SENCHO_MESH_TOKEN');
const nodeIdRaw = requireEnv('MESH_NODE_ID');
const nodeId = Number.parseInt(nodeIdRaw, 10);
if (!Number.isFinite(nodeId)) {
console.error('[mesh-sidecar] MESH_NODE_ID must be an integer');
process.exit(1);
}
let client: ControlClient | null = null;
// Bridge with explicit named params so types stay tight; client is wired
// immediately after construction so the null check is only racing against
// the 5s stats timer first tick.
const forwarder = new Forwarder({
resolve: (connId, port, remoteAddr) => client?.resolve(connId, port, remoteAddr),
sendData: (streamId, payload) => client?.sendData(streamId, payload),
sendClose: (streamId) => client?.sendClose(streamId),
sendStats: (streamId, bytesIn, bytesOut, lastActivity) =>
client?.sendStats(streamId, bytesIn, bytesOut, lastActivity),
});
forwarder.start();
client = new ControlClient({
controlUrl,
token,
nodeId,
sidecarVersion: SIDECAR_VERSION,
forwarder,
});
client.start();
const shutdown = async () => {
await forwarder.shutdown();
await client?.shutdown();
process.exit(0);
};
process.on('SIGTERM', () => { void shutdown(); });
process.on('SIGINT', () => { void shutdown(); });
console.log(`[mesh-sidecar] started for node=${nodeId} version=${SIDECAR_VERSION}`);
}
main();
-146
View File
@@ -1,146 +0,0 @@
/**
* Mesh sidecar control protocol. Wire-compatible shape with the pilot tunnel:
* JSON text frames for control + a single binary frame type carrying the
* tunneled bytes.
*
* [ 1 byte: BinaryFrameType ][ 4 bytes: streamId (BE) ][ payload bytes... ]
*/
export const PROTOCOL_VERSION = 1;
export enum BinaryFrameType {
Data = 0x01,
}
export type ControlFrame =
| ListenFrame
| UnlistenFrame
| HelloFrame
| ResolveFrame
| ResolveOkFrame
| ResolveErrFrame
| StreamStatsFrame
| LogFrame
| CloseFrame;
export interface HelloFrame {
t: 'hello';
version: number;
nodeId: number;
sidecarVersion: string;
}
/** Sencho -> sidecar: start listening on a TCP port for inbound app traffic. */
export interface ListenFrame {
t: 'listen';
port: number;
}
/** Sencho -> sidecar: stop listening on the given port. */
export interface UnlistenFrame {
t: 'unlisten';
port: number;
}
/**
* Sidecar -> Sencho: a new local TCP connection arrived on this port; what
* stream should I attach it to?
*/
export interface ResolveFrame {
t: 'resolve';
connId: number;
port: number;
remoteAddr?: string;
}
/** Sencho -> sidecar: forward bytes for connId via streamId. */
export interface ResolveOkFrame {
t: 'resolve_ok';
connId: number;
streamId: number;
alias?: string;
}
/** Sencho -> sidecar: drop the connection with a reason. */
export interface ResolveErrFrame {
t: 'resolve_err';
connId: number;
code: 'no_route' | 'tunnel_down' | 'denied' | 'unreachable' | 'agent_error';
message?: string;
}
/** Sidecar -> Sencho: per-stream byte counters every ~5 seconds while open. */
export interface StreamStatsFrame {
t: 'stream.stats';
streamId: number;
bytesIn: number;
bytesOut: number;
lastActivity: number;
}
/** Sidecar -> Sencho: structured log line surfaced into the activity log. */
export interface LogFrame {
t: 'log';
level: 'info' | 'warn' | 'error';
message: string;
details?: Record<string, unknown>;
}
/** Either side -> close a stream. */
export interface CloseFrame {
t: 'close';
streamId: number;
}
export function encodeControl(frame: ControlFrame): string {
return JSON.stringify(frame);
}
export function decodeControl(raw: string): ControlFrame {
const parsed = JSON.parse(raw);
if (!parsed || typeof parsed !== 'object' || typeof parsed.t !== 'string') {
throw new Error('invalid control frame: missing type discriminator');
}
return parsed as ControlFrame;
}
export function encodeData(streamId: number, payload: Buffer): Buffer {
if (!Number.isInteger(streamId) || streamId < 0 || streamId > 0xffffffff) {
throw new Error(`invalid streamId: ${streamId}`);
}
const out = Buffer.allocUnsafe(5 + payload.length);
out.writeUInt8(BinaryFrameType.Data, 0);
out.writeUInt32BE(streamId, 1);
payload.copy(out, 5);
return out;
}
export interface DecodedData {
streamId: number;
payload: Buffer;
}
export function decodeData(buf: Buffer): DecodedData {
if (buf.length < 5) throw new Error(`data frame too short: ${buf.length} bytes`);
const type = buf.readUInt8(0);
if (type !== BinaryFrameType.Data) throw new Error(`unknown binary frame type: ${type}`);
return {
streamId: buf.readUInt32BE(1),
payload: buf.subarray(5),
};
}
export type WsRawData = Buffer | ArrayBuffer | Buffer[] | string;
export function wsDataToBuffer(data: WsRawData): Buffer | null {
if (Buffer.isBuffer(data)) return data;
if (data instanceof ArrayBuffer) return Buffer.from(data);
if (Array.isArray(data)) return Buffer.concat(data.map((d) => Buffer.isBuffer(d) ? d : Buffer.from(d as ArrayBuffer)));
return null;
}
export function wsDataToString(data: WsRawData): string | null {
if (typeof data === 'string') return data;
const buf = wsDataToBuffer(data);
return buf ? buf.toString('utf8') : null;
}
-16
View File
@@ -1,16 +0,0 @@
{
"compilerOptions": {
"target": "ES2022",
"module": "commonjs",
"lib": ["ES2022"],
"outDir": "./dist",
"rootDir": "./src",
"strict": true,
"esModuleInterop": true,
"skipLibCheck": true,
"forceConsistentCasingInFileNames": true,
"resolveJsonModule": true
},
"include": ["src/**/*"],
"exclude": ["node_modules", "dist"]
}
-11
View File
@@ -1,11 +0,0 @@
import { defineConfig } from 'vitest/config';
export default defineConfig({
test: {
environment: 'node',
include: ['src/__tests__/**/*.test.ts'],
exclude: ['dist/**', 'node_modules/**'],
pool: 'forks',
testTimeout: 15_000,
},
});