{"version":3,"file":"unix.d.ts","sourceRoot":"","sources":["../src/unix.ts"],"names":[],"mappings":"AAGA,OAAO,EAIN,KAAK,QAAQ,EACb,MAAM,6BAA6B,CAAC;AAGrC,OAAO,KAAK,EAAiB,oBAAoB,EAAyB,MAAM,gBAAgB,CAAC;AAOjG,MAAM,WAAW,oBAAoB;IACpC,IAAI,EAAE,MAAM,CAAC;IACb,eAAe,CAAC,EAAE,MAAM,CAAC;CACzB;AAED,MAAM,WAAW,eAAe;IAC/B,QAAQ,EAAE,QAAQ,CAAC;IACnB,IAAI,EAAE,MAAM,CAAC;CACb;AAED,MAAM,WAAW,0BAA0B;IAC1C,0DAA0D;IAC1D,SAAS,EAAE,MAAM,CAAC;IAClB,4EAA4E;IAC5E,SAAS,CAAC,EAAE,MAAM,CAAC;CACnB;AAED,iFAAiF;AACjF,wBAAsB,mBAAmB,CAAC,OAAO,EAAE,0BAA0B,GAAG,OAAO,CAAC,eAAe,EAAE,CAAC,CAgDzG;AAED,8GAA8G;AAC9G,wBAAgB,0BAA0B,CAAC,OAAO,EAAE,oBAAoB,GAAG,oBAAoB,CAG9F","sourcesContent":["import { lstat, readdir } from \"node:fs/promises\";\nimport { createConnection, type Socket } from \"node:net\";\nimport { join } from \"node:path\";\nimport {\n\tDEFAULT_MAX_FRAME_LENGTH,\n\tisServerId,\n\tProtocolValidationError,\n\ttype ServerId,\n} from \"@earendil-works/pi-protocol\";\nimport { Client } from \"./client.ts\";\nimport { DisconnectedError, ServerError } from \"./errors.ts\";\nimport type { ByteTransport, ByteTransportFactory, ByteTransportHandlers } from \"./transport.ts\";\n\nconst DEFAULT_DISCOVERY_TIMEOUT_MS = 1_000;\nconst MAX_TIMER_DELAY_MS = 2_147_483_647;\nconst UNIX_SOCKET_SUFFIX = \".sock\";\nconst MAX_CONCURRENT_DISCOVERY_PROBES = 16;\n\nexport interface UnixTransportOptions {\n\tpath: string;\n\tmaxPendingBytes?: number;\n}\n\nexport interface UnixServerRoute {\n\tserverId: ServerId;\n\tpath: string;\n}\n\nexport interface DiscoverUnixServersOptions {\n\t/** Directory containing server-addressed Unix sockets. */\n\tdirectory: string;\n\t/** Maximum time for each connection and handshake. Defaults to 1,000 ms. */\n\ttimeoutMs?: number;\n}\n\n/** Discover reachable local servers by probing server-addressed Unix sockets. */\nexport async function discoverUnixServers(options: DiscoverUnixServersOptions): Promise<UnixServerRoute[]> {\n\tif (process.platform === \"win32\") throw new Error(\"Unix transport is not supported on Windows\");\n\tconst directory = options.directory;\n\tconst timeoutMs = options.timeoutMs ?? DEFAULT_DISCOVERY_TIMEOUT_MS;\n\tif (!Number.isSafeInteger(timeoutMs) || timeoutMs <= 0 || timeoutMs > MAX_TIMER_DELAY_MS) {\n\t\tthrow new TypeError(`Unix discovery timeoutMs must be an integer between 1 and ${MAX_TIMER_DELAY_MS}`);\n\t}\n\n\tlet names: string[];\n\ttry {\n\t\tnames = await readdir(directory);\n\t} catch (error) {\n\t\tif (isErrorCode(error, \"ENOENT\")) return [];\n\t\tthrow error;\n\t}\n\n\tconst candidates = names.flatMap((name): UnixServerRoute[] => {\n\t\tif (!name.endsWith(UNIX_SOCKET_SUFFIX)) return [];\n\t\tconst serverId = name.slice(0, -UNIX_SOCKET_SUFFIX.length);\n\t\treturn isServerId(serverId) ? [{ serverId, path: join(directory, name) }] : [];\n\t});\n\tconst routes: UnixServerRoute[] = [];\n\tlet nextIndex = 0;\n\tlet failure: { error: unknown } | undefined;\n\tconst workerCount = Math.min(MAX_CONCURRENT_DISCOVERY_PROBES, candidates.length);\n\tawait Promise.all(\n\t\tArray.from({ length: workerCount }, async () => {\n\t\t\twhile (!failure) {\n\t\t\t\tconst candidate = candidates[nextIndex++];\n\t\t\t\tif (!candidate) return;\n\t\t\t\ttry {\n\t\t\t\t\ttry {\n\t\t\t\t\t\tif (!(await lstat(candidate.path)).isSocket()) continue;\n\t\t\t\t\t} catch (error) {\n\t\t\t\t\t\t// A socket can disappear between readdir and lstat during normal server shutdown.\n\t\t\t\t\t\tif (isErrorCode(error, \"ENOENT\")) continue;\n\t\t\t\t\t\tthrow error;\n\t\t\t\t\t}\n\t\t\t\t\tconst route = await probeUnixServer(candidate, timeoutMs);\n\t\t\t\t\tif (route) routes.push(route);\n\t\t\t\t} catch (error) {\n\t\t\t\t\tfailure ??= { error };\n\t\t\t\t}\n\t\t\t}\n\t\t}),\n\t);\n\tif (failure) throw failure.error;\n\treturn routes.sort((left, right) => left.serverId.localeCompare(right.serverId));\n}\n\n/** Creates fresh Unix-domain socket transports for Client connection attempts in Node-compatible runtimes. */\nexport function createUnixTransportFactory(options: UnixTransportOptions): ByteTransportFactory {\n\tconst maxPendingBytes = validateUnixTransportOptions(options);\n\treturn (handlers) => connectUnixSocket(options.path, maxPendingBytes, handlers);\n}\n\nfunction validateUnixTransportOptions(options: UnixTransportOptions): number {\n\tif (options.path.length === 0) throw new TypeError(\"Unix transport path must not be empty\");\n\tconst maxPendingBytes = options.maxPendingBytes ?? DEFAULT_MAX_FRAME_LENGTH * 4;\n\tif (!Number.isSafeInteger(maxPendingBytes) || maxPendingBytes <= 0) {\n\t\tthrow new TypeError(\"Unix transport maxPendingBytes must be a positive safe integer\");\n\t}\n\tif (process.platform === \"win32\") throw new Error(\"Unix transport is not supported on Windows\");\n\treturn maxPendingBytes;\n}\n\nfunction connectUnixSocket(\n\tpath: string,\n\tmaxPendingBytes: number,\n\thandlers: ByteTransportHandlers,\n\tonSocket?: (socket: Socket) => void,\n): Promise<ByteTransport> {\n\treturn new Promise<ByteTransport>((resolve, reject) => {\n\t\tconst socket = createConnection(path);\n\t\tonSocket?.(socket);\n\t\tlet connected = false;\n\t\tlet terminal = false;\n\n\t\tconst close = (): void => {\n\t\t\tif (terminal) return;\n\t\t\tterminal = true;\n\t\t\tsocket.destroy();\n\t\t\tif (connected) handlers.onClose();\n\t\t\telse reject(new Error(\"Unix transport closed before connecting\"));\n\t\t};\n\n\t\tsocket.once(\"connect\", () => {\n\t\t\tif (terminal) return;\n\t\t\tconnected = true;\n\t\t\tresolve(\n\t\t\t\tnew UnixByteTransport(socket, maxPendingBytes, () => {\n\t\t\t\t\tterminal = true;\n\t\t\t\t}),\n\t\t\t);\n\t\t});\n\t\tsocket.on(\"data\", (chunk) => {\n\t\t\tif (!terminal) handlers.onData(new Uint8Array(chunk.buffer, chunk.byteOffset, chunk.byteLength));\n\t\t});\n\t\tsocket.once(\"end\", close);\n\t\tsocket.once(\"close\", close);\n\t\tsocket.once(\"error\", (error) => {\n\t\t\tif (terminal) return;\n\t\t\tterminal = true;\n\t\t\tsocket.destroy();\n\t\t\tif (connected) handlers.onError(error);\n\t\t\telse reject(error);\n\t\t});\n\t});\n}\n\nclass UnixByteTransport implements ByteTransport {\n\treadonly #socket: Socket;\n\treadonly #maxPendingBytes: number;\n\treadonly #markLocalClose: () => void;\n\t#closed = false;\n\t#pendingBytes = 0;\n\t#writeTail: Promise<void> = Promise.resolve();\n\n\tconstructor(socket: Socket, maxPendingBytes: number, markLocalClose: () => void) {\n\t\tthis.#socket = socket;\n\t\tthis.#maxPendingBytes = maxPendingBytes;\n\t\tthis.#markLocalClose = markLocalClose;\n\t}\n\n\tsend(chunk: Uint8Array): Promise<void> {\n\t\tif (!(chunk instanceof Uint8Array)) {\n\t\t\treturn Promise.reject(new TypeError(\"Unix transport chunks must be Uint8Array\"));\n\t\t}\n\t\tif (this.#closed) return Promise.reject(new Error(\"Unix transport is closed\"));\n\t\tif (this.#pendingBytes + chunk.byteLength > this.#maxPendingBytes) {\n\t\t\treturn Promise.reject(new Error(\"Unix transport exceeded its pending byte limit\"));\n\t\t}\n\t\tthis.#pendingBytes += chunk.byteLength;\n\t\tconst bytes = chunk.slice();\n\t\tconst write = this.#writeTail.then(() => this.#write(bytes));\n\t\tconst tracked = write.finally(() => {\n\t\t\tthis.#pendingBytes -= bytes.byteLength;\n\t\t});\n\t\tthis.#writeTail = tracked.catch(() => {});\n\t\treturn tracked;\n\t}\n\n\tclose(): void {\n\t\tif (this.#closed) return;\n\t\tthis.#closed = true;\n\t\tthis.#markLocalClose();\n\t\tthis.#socket.destroy();\n\t}\n\n\t#write(chunk: Uint8Array): Promise<void> {\n\t\tif (this.#closed || !this.#socket.writable) return Promise.reject(new Error(\"Unix transport is closed\"));\n\t\treturn new Promise<void>((resolve, reject) => {\n\t\t\tlet callbackComplete = false;\n\t\t\tlet drainComplete = false;\n\t\t\tlet requiresDrain: boolean | undefined;\n\t\t\tlet settled = false;\n\n\t\t\tconst onDrain = (): void => {\n\t\t\t\tdrainComplete = true;\n\t\t\t\tfinish();\n\t\t\t};\n\t\t\tconst cleanup = (): void => {\n\t\t\t\tthis.#socket.off(\"drain\", onDrain);\n\t\t\t\tthis.#socket.off(\"close\", onClose);\n\t\t\t};\n\t\t\tconst fail = (error: Error): void => {\n\t\t\t\tif (settled) return;\n\t\t\t\tsettled = true;\n\t\t\t\tcleanup();\n\t\t\t\treject(error);\n\t\t\t};\n\t\t\tconst finish = (): void => {\n\t\t\t\tif (settled || !callbackComplete || requiresDrain === undefined) return;\n\t\t\t\tif (requiresDrain && !drainComplete) return;\n\t\t\t\tsettled = true;\n\t\t\t\tcleanup();\n\t\t\t\tresolve();\n\t\t\t};\n\t\t\tconst onClose = (): void => fail(new Error(\"Unix transport closed during write\"));\n\n\t\t\ttry {\n\t\t\t\tthis.#socket.once(\"close\", onClose);\n\t\t\t\tconst accepted = this.#socket.write(chunk, (error) => {\n\t\t\t\t\tif (error) {\n\t\t\t\t\t\tfail(error);\n\t\t\t\t\t\treturn;\n\t\t\t\t\t}\n\t\t\t\t\tcallbackComplete = true;\n\t\t\t\t\tfinish();\n\t\t\t\t});\n\t\t\t\trequiresDrain = !accepted;\n\t\t\t\tif (requiresDrain) this.#socket.once(\"drain\", onDrain);\n\t\t\t\tfinish();\n\t\t\t} catch (error) {\n\t\t\t\tfail(error instanceof Error ? error : new Error(String(error)));\n\t\t\t}\n\t\t});\n\t}\n}\n\nasync function probeUnixServer(route: UnixServerRoute, timeoutMs: number): Promise<UnixServerRoute | undefined> {\n\tconst maxPendingBytes = validateUnixTransportOptions({ path: route.path });\n\tlet socket: Socket | undefined;\n\tconst client = new Client({\n\t\tserverId: route.serverId,\n\t\ttransportFactory: (handlers) =>\n\t\t\tconnectUnixSocket(route.path, maxPendingBytes, handlers, (created) => {\n\t\t\t\tsocket = created;\n\t\t\t}),\n\t});\n\tlet timeout: ReturnType<typeof setTimeout> | undefined;\n\ttry {\n\t\tawait Promise.race([\n\t\t\tclient.connect(),\n\t\t\tnew Promise<never>((_, reject) => {\n\t\t\t\ttimeout = setTimeout(() => {\n\t\t\t\t\tsocket?.destroy();\n\t\t\t\t\treject(new UnixDiscoveryTimeoutError());\n\t\t\t\t}, timeoutMs);\n\t\t\t\ttimeout.unref();\n\t\t\t}),\n\t\t]);\n\t\treturn route;\n\t} catch (error) {\n\t\t// Missing/refused sockets are stale or shutting down. Protocol failures mean\n\t\t// the endpoint is not the advertised server. Both are safe to omit.\n\t\tif (\n\t\t\terror instanceof UnixDiscoveryTimeoutError ||\n\t\t\terror instanceof ProtocolValidationError ||\n\t\t\t(error instanceof DisconnectedError && error.cause === undefined) ||\n\t\t\t(error instanceof ServerError && error.code === \"version\") ||\n\t\t\tisErrorCode(error, \"ENOENT\") ||\n\t\t\tisErrorCode(error, \"ECONNREFUSED\") ||\n\t\t\tisErrorCode(error, \"ECONNRESET\") ||\n\t\t\tisErrorCode(error, \"EPIPE\") ||\n\t\t\tisErrorCode(error, \"ETIMEDOUT\")\n\t\t) {\n\t\t\treturn undefined;\n\t\t}\n\t\tthrow error;\n\t} finally {\n\t\tif (timeout) clearTimeout(timeout);\n\t\tawait client.dispose();\n\t\tconst activeSocket = socket;\n\t\tif (activeSocket && !activeSocket.destroyed) activeSocket.destroy();\n\t\tif (activeSocket && !activeSocket.closed) {\n\t\t\tawait new Promise<void>((resolve) => activeSocket.once(\"close\", resolve));\n\t\t}\n\t}\n}\n\nclass UnixDiscoveryTimeoutError extends Error {}\n\nfunction isErrorCode(error: unknown, code: string): boolean {\n\tlet current = error;\n\tconst seen = new Set<unknown>();\n\twhile (current instanceof Error && !seen.has(current)) {\n\t\tseen.add(current);\n\t\tif (\"code\" in current && current.code === code) return true;\n\t\tcurrent = current.cause;\n\t}\n\treturn false;\n}\n"]}