mirror of
https://github.com/vercel/workflow.git
synced 2026-09-14 19:59:43 +08:00
4cde3b962b
Co-authored-by: JJ Kasper <jj@jjsweb.site>
259 lines
7.2 KiB
TypeScript
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,
|
|
};
|
|
}
|