index.mjs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640
  1. // src/server.ts
  2. import { createServer as createServerHTTP } from "http";
  3. // src/listener.ts
  4. import { Http2ServerRequest as Http2ServerRequest2, constants as h2constants } from "http2";
  5. // src/request.ts
  6. import { Http2ServerRequest } from "http2";
  7. import { Readable } from "stream";
  8. var RequestError = class extends Error {
  9. constructor(message, options) {
  10. super(message, options);
  11. this.name = "RequestError";
  12. }
  13. };
  14. var toRequestError = (e) => {
  15. if (e instanceof RequestError) {
  16. return e;
  17. }
  18. return new RequestError(e.message, { cause: e });
  19. };
  20. var GlobalRequest = global.Request;
  21. var Request = class extends GlobalRequest {
  22. constructor(input, options) {
  23. if (typeof input === "object" && getRequestCache in input) {
  24. input = input[getRequestCache]();
  25. }
  26. if (typeof options?.body?.getReader !== "undefined") {
  27. ;
  28. options.duplex ??= "half";
  29. }
  30. super(input, options);
  31. }
  32. };
  33. var newHeadersFromIncoming = (incoming) => {
  34. const headerRecord = [];
  35. const rawHeaders = incoming.rawHeaders;
  36. for (let i = 0; i < rawHeaders.length; i += 2) {
  37. const { [i]: key, [i + 1]: value } = rawHeaders;
  38. if (key.charCodeAt(0) !== /*:*/
  39. 58) {
  40. headerRecord.push([key, value]);
  41. }
  42. }
  43. return new Headers(headerRecord);
  44. };
  45. var wrapBodyStream = Symbol("wrapBodyStream");
  46. var newRequestFromIncoming = (method, url, headers, incoming, abortController) => {
  47. const init = {
  48. method,
  49. headers,
  50. signal: abortController.signal
  51. };
  52. if (method === "TRACE") {
  53. init.method = "GET";
  54. const req = new Request(url, init);
  55. Object.defineProperty(req, "method", {
  56. get() {
  57. return "TRACE";
  58. }
  59. });
  60. return req;
  61. }
  62. if (!(method === "GET" || method === "HEAD")) {
  63. if ("rawBody" in incoming && incoming.rawBody instanceof Buffer) {
  64. init.body = new ReadableStream({
  65. start(controller) {
  66. controller.enqueue(incoming.rawBody);
  67. controller.close();
  68. }
  69. });
  70. } else if (incoming[wrapBodyStream]) {
  71. let reader;
  72. init.body = new ReadableStream({
  73. async pull(controller) {
  74. try {
  75. reader ||= Readable.toWeb(incoming).getReader();
  76. const { done, value } = await reader.read();
  77. if (done) {
  78. controller.close();
  79. } else {
  80. controller.enqueue(value);
  81. }
  82. } catch (error) {
  83. controller.error(error);
  84. }
  85. }
  86. });
  87. } else {
  88. init.body = Readable.toWeb(incoming);
  89. }
  90. }
  91. return new Request(url, init);
  92. };
  93. var getRequestCache = Symbol("getRequestCache");
  94. var requestCache = Symbol("requestCache");
  95. var incomingKey = Symbol("incomingKey");
  96. var urlKey = Symbol("urlKey");
  97. var headersKey = Symbol("headersKey");
  98. var abortControllerKey = Symbol("abortControllerKey");
  99. var getAbortController = Symbol("getAbortController");
  100. var requestPrototype = {
  101. get method() {
  102. return this[incomingKey].method || "GET";
  103. },
  104. get url() {
  105. return this[urlKey];
  106. },
  107. get headers() {
  108. return this[headersKey] ||= newHeadersFromIncoming(this[incomingKey]);
  109. },
  110. [getAbortController]() {
  111. this[getRequestCache]();
  112. return this[abortControllerKey];
  113. },
  114. [getRequestCache]() {
  115. this[abortControllerKey] ||= new AbortController();
  116. return this[requestCache] ||= newRequestFromIncoming(
  117. this.method,
  118. this[urlKey],
  119. this.headers,
  120. this[incomingKey],
  121. this[abortControllerKey]
  122. );
  123. }
  124. };
  125. [
  126. "body",
  127. "bodyUsed",
  128. "cache",
  129. "credentials",
  130. "destination",
  131. "integrity",
  132. "mode",
  133. "redirect",
  134. "referrer",
  135. "referrerPolicy",
  136. "signal",
  137. "keepalive"
  138. ].forEach((k) => {
  139. Object.defineProperty(requestPrototype, k, {
  140. get() {
  141. return this[getRequestCache]()[k];
  142. }
  143. });
  144. });
  145. ["arrayBuffer", "blob", "clone", "formData", "json", "text"].forEach((k) => {
  146. Object.defineProperty(requestPrototype, k, {
  147. value: function() {
  148. return this[getRequestCache]()[k]();
  149. }
  150. });
  151. });
  152. Object.setPrototypeOf(requestPrototype, Request.prototype);
  153. var newRequest = (incoming, defaultHostname) => {
  154. const req = Object.create(requestPrototype);
  155. req[incomingKey] = incoming;
  156. const incomingUrl = incoming.url || "";
  157. if (incomingUrl[0] !== "/" && // short-circuit for performance. most requests are relative URL.
  158. (incomingUrl.startsWith("http://") || incomingUrl.startsWith("https://"))) {
  159. if (incoming instanceof Http2ServerRequest) {
  160. throw new RequestError("Absolute URL for :path is not allowed in HTTP/2");
  161. }
  162. try {
  163. const url2 = new URL(incomingUrl);
  164. req[urlKey] = url2.href;
  165. } catch (e) {
  166. throw new RequestError("Invalid absolute URL", { cause: e });
  167. }
  168. return req;
  169. }
  170. const host = (incoming instanceof Http2ServerRequest ? incoming.authority : incoming.headers.host) || defaultHostname;
  171. if (!host) {
  172. throw new RequestError("Missing host header");
  173. }
  174. let scheme;
  175. if (incoming instanceof Http2ServerRequest) {
  176. scheme = incoming.scheme;
  177. if (!(scheme === "http" || scheme === "https")) {
  178. throw new RequestError("Unsupported scheme");
  179. }
  180. } else {
  181. scheme = incoming.socket && incoming.socket.encrypted ? "https" : "http";
  182. }
  183. const url = new URL(`${scheme}://${host}${incomingUrl}`);
  184. if (url.hostname.length !== host.length && url.hostname !== host.replace(/:\d+$/, "")) {
  185. throw new RequestError("Invalid host header");
  186. }
  187. req[urlKey] = url.href;
  188. return req;
  189. };
  190. // src/response.ts
  191. var responseCache = Symbol("responseCache");
  192. var getResponseCache = Symbol("getResponseCache");
  193. var cacheKey = Symbol("cache");
  194. var GlobalResponse = global.Response;
  195. var Response2 = class _Response {
  196. #body;
  197. #init;
  198. [getResponseCache]() {
  199. delete this[cacheKey];
  200. return this[responseCache] ||= new GlobalResponse(this.#body, this.#init);
  201. }
  202. constructor(body, init) {
  203. let headers;
  204. this.#body = body;
  205. if (init instanceof _Response) {
  206. const cachedGlobalResponse = init[responseCache];
  207. if (cachedGlobalResponse) {
  208. this.#init = cachedGlobalResponse;
  209. this[getResponseCache]();
  210. return;
  211. } else {
  212. this.#init = init.#init;
  213. headers = new Headers(init.#init.headers);
  214. }
  215. } else {
  216. this.#init = init;
  217. }
  218. if (typeof body === "string" || typeof body?.getReader !== "undefined" || body instanceof Blob || body instanceof Uint8Array) {
  219. ;
  220. this[cacheKey] = [init?.status || 200, body, headers || init?.headers];
  221. }
  222. }
  223. get headers() {
  224. const cache = this[cacheKey];
  225. if (cache) {
  226. if (!(cache[2] instanceof Headers)) {
  227. cache[2] = new Headers(
  228. cache[2] || { "content-type": "text/plain; charset=UTF-8" }
  229. );
  230. }
  231. return cache[2];
  232. }
  233. return this[getResponseCache]().headers;
  234. }
  235. get status() {
  236. return this[cacheKey]?.[0] ?? this[getResponseCache]().status;
  237. }
  238. get ok() {
  239. const status = this.status;
  240. return status >= 200 && status < 300;
  241. }
  242. };
  243. ["body", "bodyUsed", "redirected", "statusText", "trailers", "type", "url"].forEach((k) => {
  244. Object.defineProperty(Response2.prototype, k, {
  245. get() {
  246. return this[getResponseCache]()[k];
  247. }
  248. });
  249. });
  250. ["arrayBuffer", "blob", "clone", "formData", "json", "text"].forEach((k) => {
  251. Object.defineProperty(Response2.prototype, k, {
  252. value: function() {
  253. return this[getResponseCache]()[k]();
  254. }
  255. });
  256. });
  257. Object.setPrototypeOf(Response2, GlobalResponse);
  258. Object.setPrototypeOf(Response2.prototype, GlobalResponse.prototype);
  259. // src/utils.ts
  260. async function readWithoutBlocking(readPromise) {
  261. return Promise.race([readPromise, Promise.resolve().then(() => Promise.resolve(void 0))]);
  262. }
  263. function writeFromReadableStreamDefaultReader(reader, writable, currentReadPromise) {
  264. const cancel = (error) => {
  265. reader.cancel(error).catch(() => {
  266. });
  267. };
  268. writable.on("close", cancel);
  269. writable.on("error", cancel);
  270. (currentReadPromise ?? reader.read()).then(flow, handleStreamError);
  271. return reader.closed.finally(() => {
  272. writable.off("close", cancel);
  273. writable.off("error", cancel);
  274. });
  275. function handleStreamError(error) {
  276. if (error) {
  277. writable.destroy(error);
  278. }
  279. }
  280. function onDrain() {
  281. reader.read().then(flow, handleStreamError);
  282. }
  283. function flow({ done, value }) {
  284. try {
  285. if (done) {
  286. writable.end();
  287. } else if (!writable.write(value)) {
  288. writable.once("drain", onDrain);
  289. } else {
  290. return reader.read().then(flow, handleStreamError);
  291. }
  292. } catch (e) {
  293. handleStreamError(e);
  294. }
  295. }
  296. }
  297. function writeFromReadableStream(stream, writable) {
  298. if (stream.locked) {
  299. throw new TypeError("ReadableStream is locked.");
  300. } else if (writable.destroyed) {
  301. return;
  302. }
  303. return writeFromReadableStreamDefaultReader(stream.getReader(), writable);
  304. }
  305. var buildOutgoingHttpHeaders = (headers) => {
  306. const res = {};
  307. if (!(headers instanceof Headers)) {
  308. headers = new Headers(headers ?? void 0);
  309. }
  310. const cookies = [];
  311. for (const [k, v] of headers) {
  312. if (k === "set-cookie") {
  313. cookies.push(v);
  314. } else {
  315. res[k] = v;
  316. }
  317. }
  318. if (cookies.length > 0) {
  319. res["set-cookie"] = cookies;
  320. }
  321. res["content-type"] ??= "text/plain; charset=UTF-8";
  322. return res;
  323. };
  324. // src/utils/response/constants.ts
  325. var X_ALREADY_SENT = "x-hono-already-sent";
  326. // src/globals.ts
  327. import crypto from "crypto";
  328. if (typeof global.crypto === "undefined") {
  329. global.crypto = crypto;
  330. }
  331. // src/listener.ts
  332. var outgoingEnded = Symbol("outgoingEnded");
  333. var incomingDraining = Symbol("incomingDraining");
  334. var DRAIN_TIMEOUT_MS = 500;
  335. var MAX_DRAIN_BYTES = 64 * 1024 * 1024;
  336. var drainIncoming = (incoming) => {
  337. const incomingWithDrainState = incoming;
  338. if (incoming.destroyed || incomingWithDrainState[incomingDraining]) {
  339. return;
  340. }
  341. incomingWithDrainState[incomingDraining] = true;
  342. if (incoming instanceof Http2ServerRequest2) {
  343. try {
  344. ;
  345. incoming.stream?.close?.(h2constants.NGHTTP2_NO_ERROR);
  346. } catch {
  347. }
  348. return;
  349. }
  350. let bytesRead = 0;
  351. const cleanup = () => {
  352. clearTimeout(timer);
  353. incoming.off("data", onData);
  354. incoming.off("end", cleanup);
  355. incoming.off("error", cleanup);
  356. };
  357. const forceClose = () => {
  358. cleanup();
  359. const socket = incoming.socket;
  360. if (socket && !socket.destroyed) {
  361. socket.destroySoon();
  362. }
  363. };
  364. const timer = setTimeout(forceClose, DRAIN_TIMEOUT_MS);
  365. timer.unref?.();
  366. const onData = (chunk) => {
  367. bytesRead += chunk.length;
  368. if (bytesRead > MAX_DRAIN_BYTES) {
  369. forceClose();
  370. }
  371. };
  372. incoming.on("data", onData);
  373. incoming.on("end", cleanup);
  374. incoming.on("error", cleanup);
  375. incoming.resume();
  376. };
  377. var handleRequestError = () => new Response(null, {
  378. status: 400
  379. });
  380. var handleFetchError = (e) => new Response(null, {
  381. status: e instanceof Error && (e.name === "TimeoutError" || e.constructor.name === "TimeoutError") ? 504 : 500
  382. });
  383. var handleResponseError = (e, outgoing) => {
  384. const err = e instanceof Error ? e : new Error("unknown error", { cause: e });
  385. if (err.code === "ERR_STREAM_PREMATURE_CLOSE") {
  386. console.info("The user aborted a request.");
  387. } else {
  388. console.error(e);
  389. if (!outgoing.headersSent) {
  390. outgoing.writeHead(500, { "Content-Type": "text/plain" });
  391. }
  392. outgoing.end(`Error: ${err.message}`);
  393. outgoing.destroy(err);
  394. }
  395. };
  396. var flushHeaders = (outgoing) => {
  397. if ("flushHeaders" in outgoing && outgoing.writable) {
  398. outgoing.flushHeaders();
  399. }
  400. };
  401. var responseViaCache = async (res, outgoing) => {
  402. let [status, body, header] = res[cacheKey];
  403. let hasContentLength = false;
  404. if (!header) {
  405. header = { "content-type": "text/plain; charset=UTF-8" };
  406. } else if (header instanceof Headers) {
  407. hasContentLength = header.has("content-length");
  408. header = buildOutgoingHttpHeaders(header);
  409. } else if (Array.isArray(header)) {
  410. const headerObj = new Headers(header);
  411. hasContentLength = headerObj.has("content-length");
  412. header = buildOutgoingHttpHeaders(headerObj);
  413. } else {
  414. for (const key in header) {
  415. if (key.length === 14 && key.toLowerCase() === "content-length") {
  416. hasContentLength = true;
  417. break;
  418. }
  419. }
  420. }
  421. if (!hasContentLength) {
  422. if (typeof body === "string") {
  423. header["Content-Length"] = Buffer.byteLength(body);
  424. } else if (body instanceof Uint8Array) {
  425. header["Content-Length"] = body.byteLength;
  426. } else if (body instanceof Blob) {
  427. header["Content-Length"] = body.size;
  428. }
  429. }
  430. outgoing.writeHead(status, header);
  431. if (typeof body === "string" || body instanceof Uint8Array) {
  432. outgoing.end(body);
  433. } else if (body instanceof Blob) {
  434. outgoing.end(new Uint8Array(await body.arrayBuffer()));
  435. } else {
  436. flushHeaders(outgoing);
  437. await writeFromReadableStream(body, outgoing)?.catch(
  438. (e) => handleResponseError(e, outgoing)
  439. );
  440. }
  441. ;
  442. outgoing[outgoingEnded]?.();
  443. };
  444. var isPromise = (res) => typeof res.then === "function";
  445. var responseViaResponseObject = async (res, outgoing, options = {}) => {
  446. if (isPromise(res)) {
  447. if (options.errorHandler) {
  448. try {
  449. res = await res;
  450. } catch (err) {
  451. const errRes = await options.errorHandler(err);
  452. if (!errRes) {
  453. return;
  454. }
  455. res = errRes;
  456. }
  457. } else {
  458. res = await res.catch(handleFetchError);
  459. }
  460. }
  461. if (cacheKey in res) {
  462. return responseViaCache(res, outgoing);
  463. }
  464. const resHeaderRecord = buildOutgoingHttpHeaders(res.headers);
  465. if (res.body) {
  466. const reader = res.body.getReader();
  467. const values = [];
  468. let done = false;
  469. let currentReadPromise = void 0;
  470. if (resHeaderRecord["transfer-encoding"] !== "chunked") {
  471. let maxReadCount = 2;
  472. for (let i = 0; i < maxReadCount; i++) {
  473. currentReadPromise ||= reader.read();
  474. const chunk = await readWithoutBlocking(currentReadPromise).catch((e) => {
  475. console.error(e);
  476. done = true;
  477. });
  478. if (!chunk) {
  479. if (i === 1) {
  480. await new Promise((resolve) => setTimeout(resolve));
  481. maxReadCount = 3;
  482. continue;
  483. }
  484. break;
  485. }
  486. currentReadPromise = void 0;
  487. if (chunk.value) {
  488. values.push(chunk.value);
  489. }
  490. if (chunk.done) {
  491. done = true;
  492. break;
  493. }
  494. }
  495. if (done && !("content-length" in resHeaderRecord)) {
  496. resHeaderRecord["content-length"] = values.reduce((acc, value) => acc + value.length, 0);
  497. }
  498. }
  499. outgoing.writeHead(res.status, resHeaderRecord);
  500. values.forEach((value) => {
  501. ;
  502. outgoing.write(value);
  503. });
  504. if (done) {
  505. outgoing.end();
  506. } else {
  507. if (values.length === 0) {
  508. flushHeaders(outgoing);
  509. }
  510. await writeFromReadableStreamDefaultReader(reader, outgoing, currentReadPromise);
  511. }
  512. } else if (resHeaderRecord[X_ALREADY_SENT]) {
  513. } else {
  514. outgoing.writeHead(res.status, resHeaderRecord);
  515. outgoing.end();
  516. }
  517. ;
  518. outgoing[outgoingEnded]?.();
  519. };
  520. var getRequestListener = (fetchCallback, options = {}) => {
  521. const autoCleanupIncoming = options.autoCleanupIncoming ?? true;
  522. if (options.overrideGlobalObjects !== false && global.Request !== Request) {
  523. Object.defineProperty(global, "Request", {
  524. value: Request
  525. });
  526. Object.defineProperty(global, "Response", {
  527. value: Response2
  528. });
  529. }
  530. return async (incoming, outgoing) => {
  531. let res, req;
  532. try {
  533. req = newRequest(incoming, options.hostname);
  534. let incomingEnded = !autoCleanupIncoming || incoming.method === "GET" || incoming.method === "HEAD";
  535. if (!incomingEnded) {
  536. ;
  537. incoming[wrapBodyStream] = true;
  538. incoming.on("end", () => {
  539. incomingEnded = true;
  540. });
  541. if (incoming instanceof Http2ServerRequest2) {
  542. ;
  543. outgoing[outgoingEnded] = () => {
  544. if (!incomingEnded) {
  545. setTimeout(() => {
  546. if (!incomingEnded) {
  547. setTimeout(() => {
  548. drainIncoming(incoming);
  549. });
  550. }
  551. });
  552. }
  553. };
  554. }
  555. outgoing.on("finish", () => {
  556. if (!incomingEnded) {
  557. drainIncoming(incoming);
  558. }
  559. });
  560. }
  561. outgoing.on("close", () => {
  562. const abortController = req[abortControllerKey];
  563. if (abortController) {
  564. if (incoming.errored) {
  565. req[abortControllerKey].abort(incoming.errored.toString());
  566. } else if (!outgoing.writableFinished) {
  567. req[abortControllerKey].abort("Client connection prematurely closed.");
  568. }
  569. }
  570. if (!incomingEnded) {
  571. setTimeout(() => {
  572. if (!incomingEnded) {
  573. setTimeout(() => {
  574. drainIncoming(incoming);
  575. });
  576. }
  577. });
  578. }
  579. });
  580. res = fetchCallback(req, { incoming, outgoing });
  581. if (cacheKey in res) {
  582. return responseViaCache(res, outgoing);
  583. }
  584. } catch (e) {
  585. if (!res) {
  586. if (options.errorHandler) {
  587. res = await options.errorHandler(req ? e : toRequestError(e));
  588. if (!res) {
  589. return;
  590. }
  591. } else if (!req) {
  592. res = handleRequestError();
  593. } else {
  594. res = handleFetchError(e);
  595. }
  596. } else {
  597. return handleResponseError(e, outgoing);
  598. }
  599. }
  600. try {
  601. return await responseViaResponseObject(res, outgoing, options);
  602. } catch (e) {
  603. return handleResponseError(e, outgoing);
  604. }
  605. };
  606. };
  607. // src/server.ts
  608. var createAdaptorServer = (options) => {
  609. const fetchCallback = options.fetch;
  610. const requestListener = getRequestListener(fetchCallback, {
  611. hostname: options.hostname,
  612. overrideGlobalObjects: options.overrideGlobalObjects,
  613. autoCleanupIncoming: options.autoCleanupIncoming
  614. });
  615. const createServer = options.createServer || createServerHTTP;
  616. const server = createServer(options.serverOptions || {}, requestListener);
  617. return server;
  618. };
  619. var serve = (options, listeningListener) => {
  620. const server = createAdaptorServer(options);
  621. server.listen(options?.port ?? 3e3, options.hostname, () => {
  622. const serverInfo = server.address();
  623. listeningListener && listeningListener(serverInfo);
  624. });
  625. return server;
  626. };
  627. export {
  628. RequestError,
  629. createAdaptorServer,
  630. getRequestListener,
  631. serve
  632. };