Current section

Files

Jump to
smol src adapt.ffi.mjs
Raw

src/adapt.ffi.mjs

import { Result$Ok, Result$isOk, Result$Ok$0, toList } from "./gleam.mjs";
import { Option$isSome, Option$Some$0 } from "../gleam_stdlib/gleam/option.mjs";
import * as uri from "../gleam_stdlib/gleam/uri.mjs";
import {
Uri$Uri$scheme,
Uri$Uri$host,
Uri$Uri$port,
Uri$Uri$path,
Uri$Uri$query,
} from "../gleam_stdlib/gleam/uri.mjs";
import {
parse_method,
scheme_from_string,
Scheme$Http,
} from "../gleam_http/gleam/http.mjs";
import { Request$Request } from "../gleam_http/gleam/http/request.mjs";
import {
Response$Response$status,
Response$Response$headers,
Response$Response$body,
} from "../gleam_http/gleam/http/response.mjs";
import { status_text } from "./smol.mjs";
export function adapt(handler) {
return async function (request) {
return gleamToResponse(await handler(requestToGleam(request)));
};
}
export function adaptNode(handler) {
const fetchHandler = adapt(handler);
const nodeHandler = async function (req, res) {
try {
const { Readable } = await import("node:stream");
const request = nodeToRequest({ toWebStream: Readable.toWeb }, req);
const response = await fetchHandler(request);
await responseToNode(response, res);
} catch (error) {
if (!res.headersSent) {
res.statusCode = 500;
res.end("Internal Server Error");
} else {
res.destroy(error);
}
throw err;
}
};
return nodeHandler;
}
export function requestToGleam(request) {
const method = parse_method(request.method);
if (!Result$isOk(method)) {
throw new globalThis.Error("Invalid method");
}
const headers = toList([...request.headers.entries()]);
const body = request.body;
const url = uri.parse(request.url);
if (!Result$isOk(url)) {
throw new globalThis.Error("Invalid url");
}
const urlVal = Result$Ok$0(url);
const urlScheme = Uri$Uri$scheme(urlVal);
const scheme = Option$isSome(urlScheme)
? scheme_from_string(Option$Some$0(urlScheme))
: Result$Ok(Scheme$Http());
if (!Result$isOk(scheme)) {
throw new globalThis.Error("Invalid scheme");
}
const host = option_unwrap(Uri$Uri$host(urlVal), "localhost");
return Request$Request(
Result$Ok$0(method),
headers,
body,
Result$Ok$0(scheme),
host,
Uri$Uri$port(urlVal),
Uri$Uri$path(urlVal),
Uri$Uri$query(urlVal),
);
}
function option_unwrap(option, default_value) {
if (Option$isSome(option)) {
return Option$Some$0(option);
}
return default_value;
}
export function gleamToResponse(response) {
const headers = new globalThis.Headers();
for (const [name, value] of Response$Response$headers(response)) {
headers.append(name, value);
}
return new globalThis.Response(Response$Response$body(response), {
status: Response$Response$status(response),
statusText: status_text(Response$Response$status(response)),
headers: headers,
});
}
export function nodeToRequest({ toWebStream }, 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" ? toWebStream(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;
}
export 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();
// 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();
}