Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
82 changes: 79 additions & 3 deletions src/ts/parser.ts
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,8 @@ export interface CSVParserOptions<T = Record<string, string>> {
chunk?: (results: ChunkResult, parser: ParserHandle) => void;
/** Rows per chunk when using chunk callback (default: 1000) */
chunkSize?: number;
/** Explicitly treat string source as a URL to download (default: auto-detect http/https) */
download?: boolean;
}

/** Parse metadata (PapaParse-compatible) */
Expand Down Expand Up @@ -149,8 +151,9 @@ export class CSVParser<T = Record<string, string>>
private headerRow: string[] | null = null;
private startTime: number = 0;
private closed: boolean = false;
private aborted: boolean = false;
private truncated: boolean = false;
private needsAsyncInit: boolean = false;
private loaded: boolean = false;
private sourcePath: string | null = null;

// Error tracking
Expand Down Expand Up @@ -189,14 +192,23 @@ export class CSVParser<T = Record<string, string>>
this.options.delimiter = this.autoDetectDelimiter(source);
}

// Initialize immediately for file paths
if (typeof source === "string" && !source.startsWith("http")) {
// Determine input type and initialize
const isUrl = typeof source === "string" &&
(source.startsWith("http://") || source.startsWith("https://") || options.download === true);

if (isUrl || source instanceof ReadableStream) {
// Async sources: defer initialization until load() is called
this.needsAsyncInit = true;
} else if (typeof source === "string") {
this.sourcePath = source;
this.initFromFile(source);
this.loaded = true;
} else if (source instanceof Uint8Array) {
this.initFromBuffer(source);
this.loaded = true;
} else if (source instanceof ArrayBuffer) {
this.initFromBuffer(new Uint8Array(source));
this.loaded = true;
}
}

Expand Down Expand Up @@ -272,6 +284,59 @@ export class CSVParser<T = Record<string, string>>
}
}

/**
* Load data from async sources (URL or ReadableStream).
* Must be called before iteration when using URL or stream input.
* The async iterator calls this automatically.
*/
async load(): Promise<void> {
if (this.loaded) return;
if (!this.needsAsyncInit) {
throw new Error("load() is only needed for URL or ReadableStream sources");
}

const source = this.source;

if (typeof source === "string") {
// URL source: fetch and buffer
const response = await fetch(source);
if (!response.ok) {
throw new Error(
`Failed to fetch CSV from ${source}: ${response.status} ${response.statusText}`
);
}
const buffer = new Uint8Array(await response.arrayBuffer());
this.initFromBuffer(buffer);
} else if (source instanceof ReadableStream) {
// ReadableStream source: consume all chunks
const chunks: Uint8Array[] = [];
const reader = source.getReader();

while (true) {
const { done, value } = await reader.read();
if (done) break;
if (value instanceof Uint8Array) {
chunks.push(value);
} else {
chunks.push(new TextEncoder().encode(String(value)));
}
}

// Concatenate chunks into a single buffer
const totalLen = chunks.reduce((sum, c) => sum + c.length, 0);
const buffer = new Uint8Array(totalLen);
let offset = 0;
for (const chunk of chunks) {
buffer.set(chunk, offset);
offset += chunk.length;
}

this.initFromBuffer(buffer);
}

this.loaded = true;
}

/**
* Parse header row and build column name map.
*/
Expand Down Expand Up @@ -991,6 +1056,12 @@ export class CSVParser<T = Record<string, string>>
* Synchronous iterator for for-of loops.
*/
*[Symbol.iterator](): Iterator<CSVRow<T>> {
if (this.needsAsyncInit && !this.loaded) {
throw new Error(
"Parser requires async loading for URL or ReadableStream sources. " +
"Call await parser.load() first, or use for-await-of."
);
}
if (!this.handle) {
throw new Error("Parser not initialized");
}
Expand Down Expand Up @@ -1036,8 +1107,13 @@ export class CSVParser<T = Record<string, string>>

/**
* Async iterator for non-blocking iteration.
* Automatically loads URL/ReadableStream sources if not yet loaded.
*/
async *[Symbol.asyncIterator](): AsyncIterator<CSVRow<T>> {
// Auto-load for async sources
if (this.needsAsyncInit && !this.loaded) {
await this.load();
}
if (!this.handle) {
throw new Error("Parser not initialized");
}
Expand Down
232 changes: 232 additions & 0 deletions test/unit/stream-url.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,232 @@
/**
* Tests for Issue #10: ReadableStream and URL/HTTP input support
*/

import { describe, test, expect } from "bun:test";
import { CSVParser } from "../../src/ts/parser";

describe("ReadableStream input", () => {
test("parses from ReadableStream", async () => {
const csvData = "name,age\nAlice,30\nBob,25\n";
const stream = new ReadableStream({
start(controller) {
controller.enqueue(new TextEncoder().encode(csvData));
controller.close();
},
});

const parser = new CSVParser(stream);
await parser.load();

const rows: any[] = [];
for (const row of parser) {
rows.push(row.toObject());
}
parser.close();

expect(rows.length).toBe(2);
expect(rows[0].name).toBe("Alice");
expect(rows[0].age).toBe("30");
expect(rows[1].name).toBe("Bob");
});

test("parses from multi-chunk ReadableStream", async () => {
const chunks = [
"name,age\nAli",
"ce,30\nBob,25\n",
];
const stream = new ReadableStream({
start(controller) {
for (const chunk of chunks) {
controller.enqueue(new TextEncoder().encode(chunk));
}
controller.close();
},
});

const parser = new CSVParser(stream);
await parser.load();

const rows: any[] = [];
for (const row of parser) {
rows.push(row.toObject());
}
parser.close();

expect(rows.length).toBe(2);
expect(rows[0].name).toBe("Alice");
expect(rows[1].name).toBe("Bob");
});

test("async iterator auto-loads from ReadableStream", async () => {
const csvData = "name,age\nAlice,30\n";
const stream = new ReadableStream({
start(controller) {
controller.enqueue(new TextEncoder().encode(csvData));
controller.close();
},
});

const parser = new CSVParser(stream);
const rows: any[] = [];

// for-await-of should auto-load
for await (const row of parser) {
rows.push(row.toObject());
}
parser.close();

expect(rows.length).toBe(1);
expect(rows[0].name).toBe("Alice");
});

test("sync iterator throws without load() for stream", () => {
const stream = new ReadableStream({
start(controller) {
controller.enqueue(new TextEncoder().encode("name\nAlice\n"));
controller.close();
},
});

const parser = new CSVParser(stream);

expect(() => {
for (const _row of parser) { /* should not reach here */ }
}).toThrow("async loading");

parser.close();
});
});

describe("URL input", () => {
test("fetches and parses from URL", async () => {
const csvData = "name,age\nAlice,30\nBob,25\n";
const server = Bun.serve({
port: 0,
fetch() {
return new Response(csvData, {
headers: { "Content-Type": "text/csv" },
});
},
});

try {
const parser = new CSVParser(`http://localhost:${server.port}/data.csv`);
await parser.load();

const rows: any[] = [];
for (const row of parser) {
rows.push(row.toObject());
}
parser.close();

expect(rows.length).toBe(2);
expect(rows[0].name).toBe("Alice");
expect(rows[1].name).toBe("Bob");
} finally {
server.stop();
}
});

test("async iterator auto-loads from URL", async () => {
const csvData = "name,city\nAlice,NYC\n";
const server = Bun.serve({
port: 0,
fetch() {
return new Response(csvData, {
headers: { "Content-Type": "text/csv" },
});
},
});

try {
const parser = new CSVParser(`http://localhost:${server.port}/data.csv`);
const rows: any[] = [];

for await (const row of parser) {
rows.push(row.toObject());
}
parser.close();

expect(rows.length).toBe(1);
expect(rows[0].name).toBe("Alice");
expect(rows[0].city).toBe("NYC");
} finally {
server.stop();
}
});

test("download flag treats string as URL", async () => {
const csvData = "x\n1\n";
const server = Bun.serve({
port: 0,
fetch() {
return new Response(csvData);
},
});

try {
// Use download: true with full URL
const parser = new CSVParser(
`http://localhost:${server.port}/data`,
{ download: true }
);
await parser.load();

const rows: any[] = [];
for (const row of parser) {
rows.push(row.toObject());
}
parser.close();

expect(rows.length).toBe(1);
expect(rows[0].x).toBe("1");
} finally {
server.stop();
}
});

test("throws on HTTP error", async () => {
const server = Bun.serve({
port: 0,
fetch() {
return new Response("Not Found", { status: 404 });
},
});

try {
const parser = new CSVParser(`http://localhost:${server.port}/missing.csv`);

await expect(parser.load()).rejects.toThrow("404");
parser.close();
} finally {
server.stop();
}
});

test("load() is idempotent", async () => {
const csvData = "name\nAlice\n";
const server = Bun.serve({
port: 0,
fetch() {
return new Response(csvData);
},
});

try {
const parser = new CSVParser(`http://localhost:${server.port}/data.csv`);
await parser.load();
await parser.load(); // Second call should be a no-op

const rows: any[] = [];
for (const row of parser) {
rows.push(row.toObject());
}
parser.close();

expect(rows.length).toBe(1);
} finally {
server.stop();
}
});
});
Loading