i love deno-fmt

This commit is contained in:
laura 2025-11-01 23:05:48 -03:00
parent 5197e3316d
commit 995161f491
Signed by: w
GPG key ID: BCD2117C99E69817
9 changed files with 435 additions and 430 deletions

View file

@ -1,26 +1,30 @@
{ {
"tasks": { "tasks": {
"dev": "deno run --allow-net src/main.ts" "dev": "deno run --allow-net src/main.ts"
}, },
"imports": { "imports": {
"interest/jsx-runtime": "./src/jsx.ts" "interest/jsx-runtime": "./src/jsx.ts"
}, },
"lint": { "fmt": {
"rules": { "useTabs": true,
"tags": [ "semiColons": true
"recommended" },
], "lint": {
"exclude": ["no-explicit-any", "require-await"] "rules": {
} "tags": [
}, "recommended"
"compilerOptions": { ],
"jsx": "react-jsx", "exclude": ["no-explicit-any", "require-await"]
"jsxImportSource": "interest", }
"lib": [ },
"deno.ns", "compilerOptions": {
"esnext", "jsx": "react-jsx",
"dom", "jsxImportSource": "interest",
"dom.iterable" "lib": [
] "deno.ns",
} "esnext",
"dom",
"dom.iterable"
]
}
} }

View file

@ -15,17 +15,17 @@ const copyrightHeader = `/**
const dir = "./"; const dir = "./";
for await ( for await (
const entry of walk(dir, { const entry of walk(dir, {
exts: [".ts", ".tsx"], exts: [".ts", ".tsx"],
includeDirs: false, includeDirs: false,
skip: [/node_modules/, /copyright\.ts$/], skip: [/node_modules/, /copyright\.ts$/],
}) })
) { ) {
const filePath = entry.path; const filePath = entry.path;
const content = await Deno.readTextFile(filePath); const content = await Deno.readTextFile(filePath);
if (!content.startsWith(copyrightHeader)) { if (!content.startsWith(copyrightHeader)) {
await Deno.writeTextFile(filePath, copyrightHeader + "\n" + content); await Deno.writeTextFile(filePath, copyrightHeader + "\n" + content);
console.log(`Added header to ${filePath}`); console.log(`Added header to ${filePath}`);
} }
} }

View file

@ -6,21 +6,22 @@
import { close, open } from "interest/jsx-runtime"; import { close, open } from "interest/jsx-runtime";
async function* Fruits() { async function* Fruits() {
const fruits = ["TSX", "Apple", "Banana", "Cherry"]; const fruits = ["TSX", "Apple", "Banana", "Cherry"];
yield open("ol");
for (const fruit of fruits) { yield open("ol");
await new Promise((r) => setTimeout(r, 500)); for (const fruit of fruits) {
yield <li>{fruit}</li>; await new Promise((r) => setTimeout(r, 500));
} yield <li>{fruit}</li>;
yield close("ol"); }
yield close("ol");
} }
export default function App() { export default function App() {
return ( return (
<> <>
<h1>JSX Page</h1> <h1>JSX Page</h1>
<p class="oh hey">meowing chunk by chunk</p> <p class="oh hey">meowing chunk by chunk</p>
<Fruits /> <Fruits />
</> </>
); );
} }

50
src/global.d.ts vendored
View file

@ -8,34 +8,34 @@
import type { JsxElement } from "interest/jsx-runtime"; import type { JsxElement } from "interest/jsx-runtime";
type HTMLAttributeMap<T = HTMLElement> = Partial< type HTMLAttributeMap<T = HTMLElement> = Partial<
Omit<T, keyof Element | "children" | "style"> & { Omit<T, keyof Element | "children" | "style"> & {
style?: string; style?: string;
class?: string; class?: string;
children?: any; children?: any;
[key: `data-${string}`]: string | number | boolean | null | undefined; [key: `data-${string}`]: string | number | boolean | null | undefined;
[key: `aria-${string}`]: string | number | boolean | null | undefined; [key: `aria-${string}`]: string | number | boolean | null | undefined;
} }
>; >;
declare global { declare global {
namespace JSX { namespace JSX {
type Element = JsxElement; type Element = JsxElement;
export interface ElementChildrenAttribute { export interface ElementChildrenAttribute {
// deno-lint-ignore ban-types // deno-lint-ignore ban-types
children: {}; children: {};
} }
export type IntrinsicElements = export type IntrinsicElements =
& { & {
[K in keyof HTMLElementTagNameMap]: HTMLAttributeMap< [K in keyof HTMLElementTagNameMap]: HTMLAttributeMap<
HTMLElementTagNameMap[K] HTMLElementTagNameMap[K]
>; >;
} }
& { & {
[K in keyof SVGElementTagNameMap]: HTMLAttributeMap< [K in keyof SVGElementTagNameMap]: HTMLAttributeMap<
SVGElementTagNameMap[K] SVGElementTagNameMap[K]
>; >;
}; };
} }
} }

View file

@ -7,131 +7,131 @@ import { type Chunk } from "./http.ts";
import { ChunkedStream } from "./stream.ts"; import { ChunkedStream } from "./stream.ts";
export type Attrs = Record< export type Attrs = Record<
string, string,
string | number | boolean | null | undefined string | number | boolean | null | undefined
>; >;
export const VOID_TAGS = new Set([ export const VOID_TAGS = new Set([
"area", "area",
"base", "base",
"br", "br",
"col", "col",
"embed", "embed",
"hr", "hr",
"img", "img",
"input", "input",
"link", "link",
"meta", "meta",
"param", "param",
"source", "source",
"track", "track",
"wbr", "wbr",
]); ]);
const ESCAPE_MAP: Record<string, string> = { const ESCAPE_MAP: Record<string, string> = {
"&": "&amp;", "&": "&amp;",
"<": "&lt;", "<": "&lt;",
">": "&gt;", ">": "&gt;",
'"': "&quot;", '"': "&quot;",
"'": "&#039;", "'": "&#039;",
}; };
export function escape(input: string): string { export function escape(input: string): string {
let result = ""; let result = "";
let lastIndex = 0; let lastIndex = 0;
for (let i = 0; i < input.length; i++) { for (let i = 0; i < input.length; i++) {
const replacement = ESCAPE_MAP[input[i]]; const replacement = ESCAPE_MAP[input[i]];
if (replacement) { if (replacement) {
result += input.slice(lastIndex, i) + replacement; result += input.slice(lastIndex, i) + replacement;
lastIndex = i + 1; lastIndex = i + 1;
} }
} }
return lastIndex ? result + input.slice(lastIndex) : input; return lastIndex ? result + input.slice(lastIndex) : input;
} }
function serialize(attrs: Attrs | undefined): string { function serialize(attrs: Attrs | undefined): string {
if (!attrs) return ""; if (!attrs) return "";
let output = ""; let output = "";
for (const key in attrs) { for (const key in attrs) {
const val = attrs[key]; const val = attrs[key];
if (val == null || val === false) continue; if (val == null || val === false) continue;
output += " "; output += " ";
output += val === true ? key : `${key}="${escape(String(val))}"`; output += val === true ? key : `${key}="${escape(String(val))}"`;
} }
return output; return output;
} }
type TagRes = void | Promise<void>; type TagRes = void | Promise<void>;
type TagFn = { type TagFn = {
(attrs: Attrs, ...children: Chunk[]): TagRes; (attrs: Attrs, ...children: Chunk[]): TagRes;
(attrs: Attrs, fn: () => any): TagRes; (attrs: Attrs, fn: () => any): TagRes;
(...children: Chunk[]): TagRes; (...children: Chunk[]): TagRes;
(template: TemplateStringsArray, ...values: Chunk[]): TagRes; (template: TemplateStringsArray, ...values: Chunk[]): TagRes;
(fn: () => any): TagRes; (fn: () => any): TagRes;
}; };
export type HtmlProxy = { [K in keyof HTMLElementTagNameMap]: TagFn } & { export type HtmlProxy = { [K in keyof HTMLElementTagNameMap]: TagFn } & {
[key: string]: TagFn; [key: string]: TagFn;
}; };
const isTemplateLiteral = (arg: any): arg is TemplateStringsArray => const isTemplateLiteral = (arg: any): arg is TemplateStringsArray =>
Array.isArray(arg) && "raw" in arg; Array.isArray(arg) && "raw" in arg;
const isAttributes = (arg: any): arg is Record<string, any> => const isAttributes = (arg: any): arg is Record<string, any> =>
arg && typeof arg === "object" && !isTemplateLiteral(arg) && arg && typeof arg === "object" && !isTemplateLiteral(arg) &&
!Array.isArray(arg) && !(arg instanceof Promise); !Array.isArray(arg) && !(arg instanceof Promise);
async function render(child: unknown): Promise<string> { async function render(child: unknown): Promise<string> {
if (child == null) return ""; if (child == null) return "";
if (typeof child === "string") return escape(child); if (typeof child === "string") return escape(child);
if (typeof child === "number") return String(child); if (typeof child === "number") return String(child);
if (typeof child === "boolean") return String(Number(child)); if (typeof child === "boolean") return String(Number(child));
if (child instanceof Promise) return render(await child); if (child instanceof Promise) return render(await child);
if (Array.isArray(child)) { if (Array.isArray(child)) {
return (await Promise.all(child.map(render))).join(""); return (await Promise.all(child.map(render))).join("");
} }
if (typeof child === "function") return render(await child()); if (typeof child === "function") return render(await child());
return escape(String(child)); return escape(String(child));
} }
export function html(chunks: ChunkedStream<string>): HtmlProxy { export function html(chunks: ChunkedStream<string>): HtmlProxy {
const cache = new Map<string, TagFn>(); const cache = new Map<string, TagFn>();
const write = (buf: string) => !chunks.closed && chunks.write(buf); const write = (buf: string) => !chunks.closed && chunks.write(buf);
const handler: ProxyHandler<Record<string, TagFn>> = { const handler: ProxyHandler<Record<string, TagFn>> = {
get(_, tag: string) { get(_, tag: string) {
let fn = cache.get(tag); let fn = cache.get(tag);
if (fn) return fn; if (fn) return fn;
fn = async (...args: unknown[]) => { fn = async (...args: unknown[]) => {
const attrs = isAttributes(args[0]) ? args.shift() : undefined; const attrs = isAttributes(args[0]) ? args.shift() : undefined;
const isVoid = VOID_TAGS.has(tag.toLowerCase()); const isVoid = VOID_TAGS.has(tag.toLowerCase());
const attributes = serialize(attrs as Attrs); const attributes = serialize(attrs as Attrs);
write(`<${tag}${attributes}${isVoid ? "/" : ""}>`); write(`<${tag}${attributes}${isVoid ? "/" : ""}>`);
if (isVoid) return; if (isVoid) return;
for (const child of args) { for (const child of args) {
write(await render(child)); write(await render(child));
} }
write(`</${tag}>`); write(`</${tag}>`);
}; };
return cache.set(tag, fn), fn; return cache.set(tag, fn), fn;
}, },
}; };
const proxy = new Proxy({}, handler) as HtmlProxy; const proxy = new Proxy({}, handler) as HtmlProxy;
return proxy; return proxy;
} }

View file

@ -6,127 +6,127 @@
import { ChunkedStream } from "./stream.ts"; import { ChunkedStream } from "./stream.ts";
export interface StreamOptions { export interface StreamOptions {
headContent?: string; headContent?: string;
bodyAttributes?: string; bodyAttributes?: string;
lang?: string; lang?: string;
} }
export type Chunk = export type Chunk =
| string | string
| AsyncIterable<string> | AsyncIterable<string>
| Promise<string> | Promise<string>
| Iterable<string> | Iterable<string>
| null | null
| undefined; | undefined;
async function* normalize( async function* normalize(
value: Chunk | undefined | null, value: Chunk | undefined | null,
): AsyncIterable<string> { ): AsyncIterable<string> {
if (value == null) return; if (value == null) return;
if (typeof value === "string") { if (typeof value === "string") {
yield value; yield value;
} else if (value instanceof Promise) { } else if (value instanceof Promise) {
const resolved = await value; const resolved = await value;
if (resolved != null) yield String(resolved); if (resolved != null) yield String(resolved);
} else if (Symbol.asyncIterator in value || Symbol.iterator in value) { } else if (Symbol.asyncIterator in value || Symbol.iterator in value) {
for await (const chunk of value as AsyncIterable<string>) { for await (const chunk of value as AsyncIterable<string>) {
if (chunk != null) yield String(chunk); if (chunk != null) yield String(chunk);
} }
} else { } else {
yield String(value); yield String(value);
} }
} }
export type ChunkedWriter = ( export type ChunkedWriter = (
strings: TemplateStringsArray, strings: TemplateStringsArray,
...values: Chunk[] ...values: Chunk[]
) => Promise<void>; ) => Promise<void>;
export const makeChunkWriter = export const makeChunkWriter =
(stream: ChunkedStream<string>): ChunkedWriter => (stream: ChunkedStream<string>): ChunkedWriter =>
async (strings, ...values) => { async (strings, ...values) => {
const emit = (chunk: string) => const emit = (chunk: string) =>
!stream.closed && !stream.closed &&
(chunk === "EOF" ? stream.close() : stream.write(chunk)); (chunk === "EOF" ? stream.close() : stream.write(chunk));
for (let i = 0; i < strings.length; i++) { for (let i = 0; i < strings.length; i++) {
strings[i] && emit(strings[i]); strings[i] && emit(strings[i]);
for await (const chunk of normalize(values[i])) { for await (const chunk of normalize(values[i])) {
emit(chunk); emit(chunk);
} }
} }
}; };
export function chunkedHtml() { export function chunkedHtml() {
const chunks = new ChunkedStream<string>(); const chunks = new ChunkedStream<string>();
const stream = new ReadableStream<Uint8Array>({ const stream = new ReadableStream<Uint8Array>({
async start(controller) { async start(controller) {
const encoder = new TextEncoder(); const encoder = new TextEncoder();
try { try {
for await (const chunk of chunks) { for await (const chunk of chunks) {
controller.enqueue(encoder.encode(chunk)); controller.enqueue(encoder.encode(chunk));
} }
controller.close(); controller.close();
} catch (error) { } catch (error) {
controller.error(error); controller.error(error);
} }
}, },
cancel: chunks.close, cancel: chunks.close,
}); });
return { chunks, stream }; return { chunks, stream };
} }
const DOCUMENT_TYPE = "<!DOCTYPE html>"; const DOCUMENT_TYPE = "<!DOCTYPE html>";
const HTML_BEGIN = (lang: string) => const HTML_BEGIN = (lang: string) =>
`<html lang="${lang}"><head><meta charset="utf-8"><meta name="viewport" content="width=device-width, initial-scale=1.0">`; `<html lang="${lang}"><head><meta charset="utf-8"><meta name="viewport" content="width=device-width, initial-scale=1.0">`;
const HEAD_END = "</head><body"; const HEAD_END = "</head><body";
const BODY_END = ">"; const BODY_END = ">";
const HTML_END = "</body></html>"; const HTML_END = "</body></html>";
export interface HtmlStream { export interface HtmlStream {
write: ChunkedWriter; write: ChunkedWriter;
blob: ReadableStream<Uint8Array>; blob: ReadableStream<Uint8Array>;
chunks: ChunkedStream<string>; chunks: ChunkedStream<string>;
close(): void; close(): void;
error(err: Error): void; error(err: Error): void;
readonly response: Response; readonly response: Response;
} }
export async function createHtmlStream( export async function createHtmlStream(
options: StreamOptions = {}, options: StreamOptions = {},
): Promise<HtmlStream> { ): Promise<HtmlStream> {
const { chunks, stream } = chunkedHtml(); const { chunks, stream } = chunkedHtml();
const writer = makeChunkWriter(chunks); const writer = makeChunkWriter(chunks);
chunks.write(DOCUMENT_TYPE); chunks.write(DOCUMENT_TYPE);
chunks.write(HTML_BEGIN(options.lang || "en")); chunks.write(HTML_BEGIN(options.lang || "en"));
options.headContent && chunks.write(options.headContent); options.headContent && chunks.write(options.headContent);
chunks.write(HEAD_END); chunks.write(HEAD_END);
options.bodyAttributes && chunks.write(" " + options.bodyAttributes); options.bodyAttributes && chunks.write(" " + options.bodyAttributes);
chunks.write(BODY_END); chunks.write(BODY_END);
return { return {
write: writer, write: writer,
blob: stream, blob: stream,
chunks, chunks,
close() { close() {
if (!chunks.closed) { if (!chunks.closed) {
chunks.write(HTML_END); chunks.write(HTML_END);
chunks.close(); chunks.close();
} }
}, },
error: chunks.error, error: chunks.error,
response: new Response(stream, { response: new Response(stream, {
headers: { headers: {
"Content-Type": "text/html; charset=utf-8", "Content-Type": "text/html; charset=utf-8",
"Transfer-Encoding": "chunked", "Transfer-Encoding": "chunked",
"Cache-Control": "no-cache", "Cache-Control": "no-cache",
"Connection": "keep-alive", "Connection": "keep-alive",
}, },
}), }),
}; };
} }

View file

@ -12,111 +12,111 @@ type Props = Attrs & { children?: any };
type Component = (props: Props) => JsxElement | AsyncGenerator<JsxElement>; type Component = (props: Props) => JsxElement | AsyncGenerator<JsxElement>;
export type JsxElement = export type JsxElement =
| ((chunks: ChunkedStream<string>) => Promise<void>) | ((chunks: ChunkedStream<string>) => Promise<void>)
| AsyncGenerator<JsxElement, void, unknown>; | AsyncGenerator<JsxElement, void, unknown>;
async function render( async function render(
child: any, child: any,
chunks: ChunkedStream<string>, chunks: ChunkedStream<string>,
context: ReturnType<typeof html>, context: ReturnType<typeof html>,
): Promise<void> { ): Promise<void> {
if (child == null || child === false || child === true) return; if (child == null || child === false || child === true) return;
if (typeof child === "string") return chunks.write(escape(child)); if (typeof child === "string") return chunks.write(escape(child));
if (typeof child === "number") return chunks.write(String(child)); if (typeof child === "number") return chunks.write(String(child));
if (typeof child === "function") return await child(chunks); if (typeof child === "function") return await child(chunks);
if (child instanceof Promise) { if (child instanceof Promise) {
return await render(await child, chunks, context); return await render(await child, chunks, context);
} }
if (typeof child === "object" && Symbol.asyncIterator in child) { if (typeof child === "object" && Symbol.asyncIterator in child) {
(async () => { (async () => {
for await (const item of child) { for await (const item of child) {
await render(item, chunks, context); await render(item, chunks, context);
} }
})(); })();
return; return;
} }
if (Array.isArray(child)) { if (Array.isArray(child)) {
for (const item of child) await render(item, chunks, context); for (const item of child) await render(item, chunks, context);
return; return;
} }
chunks.write(escape(String(child))); chunks.write(escape(String(child)));
} }
export function jsx( export function jsx(
tag: string | Component | typeof Fragment, tag: string | Component | typeof Fragment,
props: Props | null = {}, props: Props | null = {},
): JsxElement { ): JsxElement {
props ||= {}; props ||= {};
return async (chunks: ChunkedStream<string>) => { return async (chunks: ChunkedStream<string>) => {
const context = html(chunks); const context = html(chunks);
const { children, ...attrs } = props; const { children, ...attrs } = props;
if (tag === Fragment) { if (tag === Fragment) {
if (!Array.isArray(children)) { if (!Array.isArray(children)) {
return await render([children], chunks, context); return await render([children], chunks, context);
} }
for (const child of children) { for (const child of children) {
await render(child, chunks, context); await render(child, chunks, context);
} }
return; return;
} }
if (typeof tag === "function") { if (typeof tag === "function") {
return await render(tag(props), chunks, context); return await render(tag(props), chunks, context);
} }
const childr = children == null ? [] : [].concat(children); const childr = children == null ? [] : [].concat(children);
const attributes = Object.keys(attrs).length ? attrs : {}; const attributes = Object.keys(attrs).length ? attrs : {};
if (!childr.length || VOID_TAGS.has(tag)) { if (!childr.length || VOID_TAGS.has(tag)) {
await context[tag](childr); await context[tag](childr);
} else { } else {
await context[tag](attributes, async () => { await context[tag](attributes, async () => {
for (const child of childr) { for (const child of childr) {
await render(child, chunks, context); await render(child, chunks, context);
} }
}); });
} }
}; };
} }
export const jsxs = jsx; export const jsxs = jsx;
async function renderJsx( async function renderJsx(
element: JsxElement | JsxElement[], element: JsxElement | JsxElement[],
chunks: ChunkedStream<string>, chunks: ChunkedStream<string>,
): Promise<void> { ): Promise<void> {
if (Array.isArray(element)) { if (Array.isArray(element)) {
for (const el of element) { for (const el of element) {
await renderJsx(el, chunks); await renderJsx(el, chunks);
} }
return; return;
} }
if (typeof element === "object" && Symbol.asyncIterator in element) { if (typeof element === "object" && Symbol.asyncIterator in element) {
for await (const item of element) { for await (const item of element) {
await renderJsx(item, chunks); await renderJsx(item, chunks);
} }
return; return;
} }
if (typeof element === "function") { if (typeof element === "function") {
await element(chunks); await element(chunks);
} }
} }
export const raw = export const raw =
(html: string): JsxElement => async (chunks: ChunkedStream<string>) => (html: string): JsxElement => async (chunks: ChunkedStream<string>) =>
void (!chunks.closed && chunks.write(html)); void (!chunks.closed && chunks.write(html));
export const open = <K extends keyof HTMLElementTagNameMap>(tag: K) => export const open = <K extends keyof HTMLElementTagNameMap>(tag: K) =>
raw(`<${tag}>`); raw(`<${tag}>`);
export const close = <K extends keyof HTMLElementTagNameMap>(tag: K) => export const close = <K extends keyof HTMLElementTagNameMap>(tag: K) =>
raw(`</${tag}>`); raw(`</${tag}>`);
export { renderJsx as render }; export { renderJsx as render };

View file

@ -8,10 +8,10 @@ import App from "./app.tsx";
import { render } from "interest/jsx-runtime"; import { render } from "interest/jsx-runtime";
Deno.serve({ Deno.serve({
port: 8080, port: 8080,
async handler() { async handler() {
const stream = await createHtmlStream({ lang: "en" }); const stream = await createHtmlStream({ lang: "en" });
await render(App(), stream.chunks); await render(App(), stream.chunks);
return stream.response; return stream.response;
}, },
}); });

View file

@ -4,141 +4,141 @@
*/ */
export class ChunkedStream<T> implements AsyncIterable<T> { export class ChunkedStream<T> implements AsyncIterable<T> {
private readonly chunks: T[] = []; private readonly chunks: T[] = [];
private readonly resolvers: ((result: IteratorResult<T>) => void)[] = []; private readonly resolvers: ((result: IteratorResult<T>) => void)[] = [];
private readonly rejectors: ((error: Error) => void)[] = []; private readonly rejectors: ((error: Error) => void)[] = [];
private _error: Error | null = null; private _error: Error | null = null;
private _closed = false; private _closed = false;
get closed(): boolean { get closed(): boolean {
return this._closed; return this._closed;
} }
write(chunk: T) { write(chunk: T) {
if (this._closed) throw new Error("Cannot write to closed stream"); if (this._closed) throw new Error("Cannot write to closed stream");
const resolver = this.resolvers.shift(); const resolver = this.resolvers.shift();
if (resolver) { if (resolver) {
this.rejectors.shift(); this.rejectors.shift();
resolver({ value: chunk, done: false }); resolver({ value: chunk, done: false });
} else { } else {
this.chunks.push(chunk); this.chunks.push(chunk);
} }
} }
close(): void { close(): void {
this._closed = true; this._closed = true;
while (this.resolvers.length) { while (this.resolvers.length) {
this.rejectors.shift(); this.rejectors.shift();
this.resolvers.shift()!({ value: undefined! as any, done: true }); this.resolvers.shift()!({ value: undefined! as any, done: true });
} }
} }
error(err: Error): void { error(err: Error): void {
if (this._closed) return; if (this._closed) return;
this._error = err; this._error = err;
this._closed = true; this._closed = true;
while (this.rejectors.length) { while (this.rejectors.length) {
this.rejectors.shift()!(err); this.rejectors.shift()!(err);
this.resolvers.shift(); this.resolvers.shift();
} }
} }
async next(): Promise<IteratorResult<T>> { async next(): Promise<IteratorResult<T>> {
if (this._error) { if (this._error) {
throw this._error; throw this._error;
} }
if (this.chunks.length) { if (this.chunks.length) {
return { value: this.chunks.shift()!, done: false }; return { value: this.chunks.shift()!, done: false };
} }
if (this._closed) return { value: undefined as any, done: true }; if (this._closed) return { value: undefined as any, done: true };
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
this.resolvers.push(resolve); this.resolvers.push(resolve);
this.rejectors.push(reject); this.rejectors.push(reject);
}); });
} }
[Symbol.asyncIterator](): AsyncIterableIterator<T> { [Symbol.asyncIterator](): AsyncIterableIterator<T> {
return this; return this;
} }
} }
export const mapStream = <T, U>( export const mapStream = <T, U>(
fn: (chunk: T, index: number) => U | Promise<U>, fn: (chunk: T, index: number) => U | Promise<U>,
) => ) =>
async function* (source: AsyncIterable<T>): AsyncIterable<U> { async function* (source: AsyncIterable<T>): AsyncIterable<U> {
let index = 0; let index = 0;
for await (const chunk of source) yield await fn(chunk, index++); for await (const chunk of source) yield await fn(chunk, index++);
}; };
export const filterStream = <T>( export const filterStream = <T>(
pred: (chunk: T, index: number) => boolean | Promise<boolean>, pred: (chunk: T, index: number) => boolean | Promise<boolean>,
) => ) =>
async function* (source: AsyncIterable<T>): AsyncIterable<T> { async function* (source: AsyncIterable<T>): AsyncIterable<T> {
let index = 0; let index = 0;
for await (const chunk of source) { for await (const chunk of source) {
if (await pred(chunk, index++)) yield chunk; if (await pred(chunk, index++)) yield chunk;
} }
}; };
export const takeStream = <T>(count: number) => export const takeStream = <T>(count: number) =>
async function* (source: AsyncIterable<T>): AsyncIterable<T> { async function* (source: AsyncIterable<T>): AsyncIterable<T> {
let taken = 0; let taken = 0;
for await (const chunk of source) { for await (const chunk of source) {
if (taken++ >= count) return; if (taken++ >= count) return;
yield chunk; yield chunk;
} }
}; };
export const skipStream = <T>(count: number) => export const skipStream = <T>(count: number) =>
async function* (source: AsyncIterable<T>): AsyncIterable<T> { async function* (source: AsyncIterable<T>): AsyncIterable<T> {
let index = 0; let index = 0;
for await (const chunk of source) { for await (const chunk of source) {
if (index++ >= count) yield chunk; if (index++ >= count) yield chunk;
} }
}; };
export const batchStream = <T>(size: number) => export const batchStream = <T>(size: number) =>
async function* (source: AsyncIterable<T>): AsyncIterable<T[]> { async function* (source: AsyncIterable<T>): AsyncIterable<T[]> {
let batch: T[] = []; let batch: T[] = [];
for await (const chunk of source) { for await (const chunk of source) {
batch.push(chunk); batch.push(chunk);
if (batch.length >= size) { if (batch.length >= size) {
yield batch; yield batch;
batch = []; batch = [];
} }
} }
if (batch.length) yield batch; if (batch.length) yield batch;
}; };
export const tapStream = <T>( export const tapStream = <T>(
fn: (chunk: T, index: number) => void | Promise<void>, fn: (chunk: T, index: number) => void | Promise<void>,
) => ) =>
async function* (source: AsyncIterable<T>): AsyncIterable<T> { async function* (source: AsyncIterable<T>): AsyncIterable<T> {
let index = 0; let index = 0;
for await (const chunk of source) { for await (const chunk of source) {
yield chunk; yield chunk;
await fn(chunk, index++); await fn(chunk, index++);
} }
}; };
export const catchStream = <T>( export const catchStream = <T>(
handler: (error: Error) => void | Promise<void>, handler: (error: Error) => void | Promise<void>,
) => ) =>
async function* (source: AsyncIterable<T>): AsyncIterable<T> { async function* (source: AsyncIterable<T>): AsyncIterable<T> {
try { try {
for await (const chunk of source) yield chunk; for await (const chunk of source) yield chunk;
} catch (err) { } catch (err) {
await handler(err instanceof Error ? err : new Error(String(err))); await handler(err instanceof Error ? err : new Error(String(err)));
} }
}; };
export const pipe = export const pipe =
<T>(...fns: Array<(src: AsyncIterable<T>) => AsyncIterable<any>>) => <T>(...fns: Array<(src: AsyncIterable<T>) => AsyncIterable<any>>) =>
(source: AsyncIterable<T>) => fns.reduce((acc, fn) => fn(acc), source); (source: AsyncIterable<T>) => fns.reduce((acc, fn) => fn(acc), source);