Files
vercel__workflow/packages/next/src/socket-server.ts
Luke Sandberg 4cde3b962b fix bad socket file location (#2021)
Co-authored-by: JJ Kasper <jj@jjsweb.site>
2026-05-19 12:11:00 -07:00

259 lines
7.2 KiB
TypeScript

import { randomBytes } from 'node:crypto';
import { mkdir, rm, writeFile } from 'node:fs/promises';
import { createServer, type Server, type Socket } from 'node:net';
import { dirname, join } from 'node:path';
/**
* Magic preamble that must prefix all messages to authenticate them as workflow messages.
* This prevents accidental processing of messages from port scanners or other local processes.
*/
const MESSAGE_PREAMBLE = 'WF:';
/**
* Generate a random authentication token for this server session.
* Clients must include this token in all messages.
*/
function generateAuthToken(): string {
return randomBytes(16).toString('hex');
}
/**
* Message types that can be sent between loader and builder
*/
export type SocketMessage =
| {
type: 'file-discovered';
filePath: string;
hasWorkflow: boolean;
hasStep: boolean;
hasSerde: boolean;
}
| { type: 'trigger-build' }
| { type: 'build-complete' };
/**
* Configuration for the socket server
*/
export interface SocketServerConfig {
isDevServer: boolean;
onFileDiscovered: (
filePath: string,
hasWorkflow: boolean,
hasStep: boolean,
hasSerde: boolean
) => void;
onTriggerBuild: () => void;
socketInfoFilePath?: string;
}
/**
* Interface for the socket IO instance returned by createSocketServer
*/
export interface SocketIO {
emit(event: 'build-complete'): void;
getAuthToken(): string;
}
/**
* Filename for the socket-info file.
*/
export const SOCKET_INFO_FILENAME = 'workflow-socket.json';
/**
* Previous filesystem location for the socket-info file. This file lives
* inside `.next/cache/`, which Vercel and Turborepo preserve across builds —
* a stale file from a prior build would cause the loader to attempt to
* connect to a dead port (ECONNREFUSED). The current location is a sibling
* of `cache/` so it isn't preserved.
*
* Exported so the builders can unlink the legacy path at boot, cleaning up
* any leftover file written by older versions of the SDK.
*/
export const LEGACY_SOCKET_INFO_RELATIVE_PATH = join(
'cache',
SOCKET_INFO_FILENAME
);
function getDefaultSocketInfoFilePath(): string {
return join(process.cwd(), '.next', SOCKET_INFO_FILENAME);
}
/**
* Remove any stale socket-info files at boot.
* @param distDir absolute path to the project's `.next` directory.
*/
export async function cleanupStaleSocketInfoFiles(
distDir: string
): Promise<void> {
await Promise.all([
rm(join(distDir, SOCKET_INFO_FILENAME), { force: true }),
rm(join(distDir, LEGACY_SOCKET_INFO_RELATIVE_PATH), { force: true }),
]);
}
/**
* Serialize a message with authentication preamble
*/
export function serializeMessage(
message: SocketMessage,
authToken: string
): string {
return `${MESSAGE_PREAMBLE}${authToken}:${JSON.stringify(message)}\n`;
}
/**
* Parse and authenticate a message from the socket
* Returns the parsed message if valid, null otherwise
*/
export function parseMessage(
line: string,
authToken: string
): SocketMessage | null {
const trimmed = line.trim();
if (!trimmed) {
return null;
}
// Check for preamble
if (!trimmed.startsWith(MESSAGE_PREAMBLE)) {
console.warn('Received message without valid preamble, ignoring');
return null;
}
// Extract auth token and payload
const withoutPreamble = trimmed.slice(MESSAGE_PREAMBLE.length);
const colonIndex = withoutPreamble.indexOf(':');
if (colonIndex === -1) {
console.warn('Received message without auth token separator, ignoring');
return null;
}
const messageToken = withoutPreamble.slice(0, colonIndex);
const payload = withoutPreamble.slice(colonIndex + 1);
// Verify auth token
if (messageToken !== authToken) {
console.warn('Received message with invalid auth token, ignoring');
return null;
}
// Parse JSON payload
try {
return JSON.parse(payload) as SocketMessage;
} catch (error) {
console.error('Failed to parse socket message JSON:', error);
return null;
}
}
/**
* Create a TCP socket server for loader<->builder communication.
* Returns a SocketIO interface for broadcasting messages and the auth token.
*
* SECURITY: Server listens on 127.0.0.1 (localhost only) and uses
* message authentication to prevent processing of unauthorized messages.
*/
export async function createSocketServer(
config: SocketServerConfig
): Promise<SocketIO> {
const authToken = generateAuthToken();
const clients = new Set<Socket>();
let buildTriggered = false;
const server: Server = createServer((socket: Socket) => {
socket.setNoDelay(true);
clients.add(socket);
// Send build-complete if build already finished (production mode)
if (buildTriggered && !config.isDevServer) {
socket.write(serializeMessage({ type: 'build-complete' }, authToken));
}
let buffer = '';
socket.on('data', (data: Buffer) => {
buffer += data.toString();
// Process complete messages (newline-delimited)
let newlineIndex = buffer.indexOf('\n');
while (newlineIndex !== -1) {
const line = buffer.slice(0, newlineIndex);
buffer = buffer.slice(newlineIndex + 1);
newlineIndex = buffer.indexOf('\n');
const message = parseMessage(line, authToken);
if (!message) {
continue;
}
if (message.type === 'file-discovered') {
config.onFileDiscovered(
message.filePath,
message.hasWorkflow,
message.hasStep,
message.hasSerde
);
} else if (message.type === 'trigger-build') {
config.onTriggerBuild();
}
}
});
socket.on('end', () => {
clients.delete(socket);
});
socket.on('error', (err: Error) => {
console.error('Socket error:', err);
clients.delete(socket);
});
});
// Listen on random available port (localhost only)
await new Promise<void>((resolve, reject) => {
server.once('error', reject);
server.listen(0, '127.0.0.1', () => {
const address = server.address();
if (address && typeof address === 'object') {
const socketInfoFilePath =
config.socketInfoFilePath || getDefaultSocketInfoFilePath();
void (async () => {
try {
await mkdir(dirname(socketInfoFilePath), { recursive: true });
await writeFile(
socketInfoFilePath,
JSON.stringify(
{
port: address.port,
authToken,
},
null,
2
)
);
process.env.WORKFLOW_SOCKET_INFO_PATH = socketInfoFilePath;
process.env.WORKFLOW_SOCKET_PORT = String(address.port);
process.env.WORKFLOW_SOCKET_AUTH = authToken;
resolve();
} catch (error) {
reject(error);
}
})();
return;
}
reject(new Error('Failed to obtain workflow socket server address'));
});
});
return {
emit: (_event: 'build-complete') => {
buildTriggered = true;
const message = serializeMessage({ type: 'build-complete' }, authToken);
for (const client of clients) {
client.write(message);
}
},
getAuthToken: () => authToken,
};
}