Current section

Files

Jump to
smol src node.ffi.mjs
Raw

src/node.ffi.mjs

import { createReadStream } from 'node:fs';
import { stat } from 'node:fs/promises'
import { Readable } from 'node:stream';
import http, { ServerResponse } from 'node:http';
import WebSocket from './vendor/ws.js';
import { Ok, Error } from './gleam.mjs';
import { None } from '../gleam_stdlib/gleam/option.mjs';
import { adapt as adaptFetch } from './adapt.ffi.mjs';
import { upgrade } from './smol.ffi.mjs';
import {
// Error
IsDir, NoEntry, NoAccess, UnknownFileError,
//
started, File
} from './smol.mjs';
export async function start(builder) {
try {
const listener = adapt(builder.handler);
const options = {};
const server = http.createServer(options, listener);
server.addListener('upgrade', adaptUpgrade(builder.handler));
let disposed = false;
const stop = () => {
if (disposed) {
return;
}
disposed = true;
console.log('smol: Gracefully shutting down. Please wait...');
process.off('SIGINT', onExit)
process.off('SIGTERM', onExit)
process.off('exit', onExit)
return new Promise((resolve, reject) => {
server.close(err => err ? reject(err) : resolve())
});
}
const onExit = async () => {
try {
await stop()
} finally {
process.exit()
}
}
await new Promise((resolve, reject) => {
server.once('error', reject)
server.listen(builder.port, builder.interface, 511, resolve)
})
process.on('SIGINT', onExit);
process.on('SIGTERM', onExit);
process.on('exit', onExit);
server.on('error', err => {
console.log('smol:', err.stack);
// FIXME: can we not just die here? maybe restart the server?
process.exit(1);
});
const addressInfo = server.address();
const port = addressInfo.port;
const address = addressInfo.address;
builder.after_start(port, address);
return started(port, address, stop);
} catch(e) {
return new Error();
}
}
export function adapt(handler) {
const fetchHandler = adaptFetch(handler);
const nodeHandler = async function (req, res) {
try {
const request = nodeToRequest(req);
const response = await fetchHandler(request);
await responseToNode(response, res);
} catch (error) {
console.error('smol:', error.stack)
if (!res.headersSent) {
res.statusCode = 500;
res.end('Internal Server Error');
} else {
res.destroy(error);
}
}
}
return nodeHandler;
}
function adaptUpgrade(handler) {
const fetchHandler = adaptFetch(handler);
const wss = new WebSocket.Server({ noServer: true, perMessageDeflate: true });
const upgradeHandler = async (req, socket, head) => {
const request = nodeToRequest(req);
const response = await fetchHandler(request);
if (response.body[upgrade]) {
const runtime = response.body[upgrade]
wss.handleUpgrade(req, socket, head, (ws) => {
runtime.send = message => ws.send(message);
runtime.close = () => ws.close();
runtime.onOpen();
ws.on('message', (data, isBinary) => runtime.onMessage(isBinary ? data : data.toString()))
ws.on('close', runtime.onClose)
})
} else {
const res = new ServerResponse(req)
res.assignSocket(socket)
responseToNode(response, res)
}
}
return upgradeHandler
}
function nodeToRequest(req) {
const protocol = req.socket?.encrypted ? 'https:' : 'http:';
const host = req.headers.host || 'localhost';
const url = new URL(req.url, `${protocol}//${host}`);
const headers = new Headers();
for (const [key, value] of Object.entries(req.headers)) {
if (Array.isArray(value)) {
for (const v of value) {
headers.append(key, v);
}
} else {
headers.append(key, value);
}
}
// Only include a body if this is not a GET or HEAD request
const body = (req.method !== 'GET' && req.method !== 'HEAD')
? Readable.toWeb(req)
: null;
const request = new globalThis.Request(url, {
method: req.method,
headers: headers,
body: body,
duplex: 'half' // Required for compatibility with Node's http streams
});
return request;
}
async function responseToNode(response, res) {
// Set headers
res.statusCode = response.status;
for (const [key, value] of response.headers.entries()) {
res.setHeader(key, value);
}
if (response.body) {
const reader = response.body.getReader();
try {
// Process the stream chunk by chunk
while (true) {
const { done, value } = await reader.read();
if (done) {
break;
}
const canContinue = res.write(Buffer.from(value));
if (!canContinue) {
// Wait for the drain event if backpressure is detected
await new Promise(resolve => res.once('drain', resolve));
}
}
res.end();
} catch (error) {
console.error('smol: sending response failed:', error.stack);
if (!res.headersSent) {
res.statusCode = 500;
res.end('Internal Server Error: Stream processing failed');
} else {
res.destroy(error);
}
}
} else {
res.end();
}
}
export async function from_file(path, offset, limit) {
try {
const fileInfo = await stat(path)
if (fileInfo.isDirectory()) {
return new Error(new IsDir());
}
const body = Readable.toWeb(createReadStream(path, {
start: offset,
end: limit > 0 ? offset + limit : Infinity
}));
return new Ok(new File(body, fileInfo.size, new None()));
} catch(e) {
if (e.code === 'ENOENT') {
return new Error(new NoEntry());
} else if (e.code === 'EACCES') {
return new Error(new NoAccess());
} else {
return new Error(new UnknownFileError());
}
}
}