{"_id":"@electric-ax/durable-streams-client-beta","_rev":"3-8fc9897e8f2ace362b9bffe2779be195","name":"@electric-ax/durable-streams-client-beta","dist-tags":{"latest":"0.3.1"},"versions":{"0.3.0":{"name":"@electric-ax/durable-streams-client-beta","version":"0.3.0","keywords":["durable-streams","streaming","client","typescript"],"author":{"name":"Durable Stream contributors"},"license":"Apache-2.0","_id":"@electric-ax/durable-streams-client-beta@0.3.0","maintainers":[{"name":"electricsql-bot","email":"infra@electric-sql.com"}],"homepage":"https://github.com/durable-streams/durable-streams#readme","bugs":{"url":"https://github.com/durable-streams/durable-streams/issues"},"bin":{"intent":"bin/intent.js"},"dist":{"shasum":"67b2e6ff6a8c13cd23e1e76f7d8b4351a777644e","tarball":"https://registry.npmjs.org/@electric-ax/durable-streams-client-beta/-/durable-streams-client-beta-0.3.0.tgz","fileCount":27,"integrity":"sha512-ECs0Q2pi6jxDfKpFaG2ydhRpmVXUIcjgcRlTIBtNTtZNvIA+EMCuCoosxP5KvqqYP9yx4XYNhw5DhEEfEM8hLQ==","signatures":[{"sig":"MEQCIC/kHxq2694LpXmZRpGvBSUW2Z0+su4OmYruxUmOL3G0AiBBj3mNbglU7Whap/ocntOXkzBnr2QHmU/ZTG4IOdErUA==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":547380},"main":"./dist/index.cjs","type":"module","_from":"file:electric-ax-durable-streams-client-beta-0.3.0.tgz","types":"./dist/index.d.ts","module":"./dist/index.js","engines":{"node":">=18.0.0"},"exports":{".":{"import":{"types":"./dist/index.d.ts","default":"./dist/index.js"},"require":{"types":"./dist/index.d.cts","default":"./dist/index.cjs"}},"./package.json":"./package.json"},"scripts":{"dev":"tsdown --watch","build":"tsdown","typecheck":"tsc --noEmit"},"_npmUser":{"name":"electricsql-bot","email":"infra@electric-sql.com"},"_resolved":"/private/var/folders/8h/10dynkk10g75zmhg0gy2sd5c0000gn/T/257113391b650b4045b2480872b0ee47/electric-ax-durable-streams-client-beta-0.3.0.tgz","_integrity":"sha512-ECs0Q2pi6jxDfKpFaG2ydhRpmVXUIcjgcRlTIBtNTtZNvIA+EMCuCoosxP5KvqqYP9yx4XYNhw5DhEEfEM8hLQ==","repository":{"url":"git+https://github.com/durable-streams/durable-streams.git","type":"git","directory":"packages/client"},"_npmVersion":"11.6.2","description":"TypeScript client for the Durable Streams protocol","directories":{},"sideEffects":false,"_nodeVersion":"24.12.0","dependencies":{"fastq":"^1.19.1","@microsoft/fetch-event-source":"^2.0.1"},"_hasShrinkwrap":false,"devDependencies":{"tsdown":"^0.9.0","fast-check":"^4.6.0","@tanstack/intent":"latest","@durable-streams/server":"0.2.2"},"_npmOperationalInternal":{"tmp":"tmp/durable-streams-client-beta_0.3.0_1776854233095_0.9137638662048564","host":"s3://npm-registry-packages-npm-production"}},"0.3.1":{"name":"@electric-ax/durable-streams-client-beta","description":"TypeScript client for the Durable Streams protocol","version":"0.3.1","author":{"name":"Durable Stream contributors"},"license":"Apache-2.0","repository":{"type":"git","url":"git+https://github.com/durable-streams/durable-streams.git","directory":"packages/client"},"bugs":{"url":"https://github.com/durable-streams/durable-streams/issues"},"keywords":["durable-streams","streaming","client","typescript"],"type":"module","main":"./dist/index.cjs","module":"./dist/index.js","types":"./dist/index.d.ts","exports":{".":{"import":{"types":"./dist/index.d.ts","default":"./dist/index.js"},"require":{"types":"./dist/index.d.cts","default":"./dist/index.cjs"}},"./package.json":"./package.json"},"sideEffects":false,"bin":{"intent":"bin/intent.js"},"dependencies":{"@microsoft/fetch-event-source":"^2.0.1","fastq":"^1.19.1"},"devDependencies":{"@tanstack/intent":"latest","fast-check":"^4.6.0","tsdown":"^0.9.0"},"engines":{"node":">=18.0.0"},"scripts":{"build":"tsdown","dev":"tsdown --watch","typecheck":"tsc --noEmit"},"_id":"@electric-ax/durable-streams-client-beta@0.3.1","homepage":"https://github.com/durable-streams/durable-streams#readme","_integrity":"sha512-smWzyfwrkA5TyRTnyO/HEbElJm54lyifRe6hgQpcfYaW0M3l3jedI5voOEvu2GTHfBKuOsFrFWsbMB0ltGm6qg==","_resolved":"/private/var/folders/8h/10dynkk10g75zmhg0gy2sd5c0000gn/T/cee84cba14a0732ac4a48a428fefb358/electric-ax-durable-streams-client-beta-0.3.1.tgz","_from":"file:electric-ax-durable-streams-client-beta-0.3.1.tgz","_nodeVersion":"24.12.0","_npmVersion":"11.6.2","dist":{"integrity":"sha512-smWzyfwrkA5TyRTnyO/HEbElJm54lyifRe6hgQpcfYaW0M3l3jedI5voOEvu2GTHfBKuOsFrFWsbMB0ltGm6qg==","shasum":"94f1c1518a4602bba90d9d4a93f8183307aefa7d","tarball":"https://registry.npmjs.org/@electric-ax/durable-streams-client-beta/-/durable-streams-client-beta-0.3.1.tgz","fileCount":28,"unpackedSize":559204,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIQCCHlG820qM57jlLkbqhNWPA3G0e/YodmesaL346NzHUwIgLsQZge9cujWCUAk6WAkhzDUxBcKZ7TCP5oFKe9y3R0E="}]},"_npmUser":{"name":"electricsql-bot","email":"infra@electric-sql.com"},"directories":{},"maintainers":[{"name":"electricsql-bot","email":"infra@electric-sql.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/durable-streams-client-beta_0.3.1_1777385640181_0.2862748482569799"},"_hasShrinkwrap":false}},"time":{"created":"2026-04-22T10:37:13.094Z","modified":"2026-04-28T14:14:00.469Z","0.3.0":"2026-04-22T10:37:13.249Z","0.3.1":"2026-04-28T14:14:00.361Z"},"bugs":{"url":"https://github.com/durable-streams/durable-streams/issues"},"author":{"name":"Durable Stream contributors"},"license":"Apache-2.0","homepage":"https://github.com/durable-streams/durable-streams#readme","keywords":["durable-streams","streaming","client","typescript"],"repository":{"type":"git","url":"git+https://github.com/durable-streams/durable-streams.git","directory":"packages/client"},"description":"TypeScript client for the Durable Streams protocol","maintainers":[{"name":"electricsql-bot","email":"infra@electric-sql.com"}],"readme":"# @durable-streams/client\n\nTypeScript client for the Durable Streams protocol.\n\n## Installation\n\n```bash\nnpm install @durable-streams/client\n```\n\n## Overview\n\nThe Durable Streams client provides three main APIs:\n\n1. **`stream()` function** - A fetch-like read-only API for consuming streams\n2. **`DurableStream` class** - A handle for read/write operations on a stream\n3. **`IdempotentProducer` class** - High-throughput producer with exactly-once write semantics (recommended for writes)\n\n## Key Features\n\n- **Exactly-Once Writes**: `IdempotentProducer` provides Kafka-style exactly-once semantics with automatic deduplication\n- **Automatic Batching**: Multiple writes are automatically batched together for high throughput\n- **Pipelining**: Up to 5 concurrent batches in flight by default for maximum throughput\n- **Streaming Reads**: `stream()` and `DurableStream.stream()` provide rich consumption options (promises, ReadableStreams, subscribers)\n- **Resumable**: Offset-based reads let you resume from any point\n- **Real-time**: Long-poll and SSE modes for live tailing with catch-up from any offset\n\n## Usage\n\n### Read-only: Using `stream()` (fetch-like API)\n\nThe `stream()` function provides a simple, fetch-like interface for reading from streams:\n\n```typescript\nimport { stream } from \"@durable-streams/client\"\n\n// Connect and get a StreamResponse\nconst res = await stream<{ message: string }>({\n  url: \"https://streams.example.com/my-account/chat/room-1\",\n  headers: {\n    Authorization: `Bearer ${process.env.DS_TOKEN!}`,\n  },\n  offset: savedOffset, // optional: resume from offset\n  live: true, // default: auto-select best live mode\n})\n\n// Accumulate all JSON items until up-to-date\nconst items = await res.json()\nconsole.log(\"All items:\", items)\n\n// Or stream live with a subscriber\nres.subscribeJson(async (batch) => {\n  for (const item of batch.items) {\n    console.log(\"item:\", item)\n    saveOffset(batch.offset) // persist for resumption\n  }\n})\n```\n\n### StreamResponse consumption methods\n\nThe `StreamResponse` object returned by `stream()` offers multiple ways to consume data:\n\n```typescript\n// Promise helpers (accumulate until first upToDate)\nconst bytes = await res.body() // Uint8Array\nconst items = await res.json() // Array<TJson>\nconst text = await res.text() // string\n\n// ReadableStreams\nconst byteStream = res.bodyStream() // ReadableStream<Uint8Array>\nconst jsonStream = res.jsonStream() // ReadableStream<TJson>\nconst textStream = res.textStream() // ReadableStream<string>\n\n// Subscribers (with backpressure)\nconst unsubscribe = res.subscribeJson(async (batch) => {\n  await processBatch(batch.items)\n})\nconst unsubscribe2 = res.subscribeBytes(async (chunk) => {\n  await processBytes(chunk.data)\n})\nconst unsubscribe3 = res.subscribeText(async (chunk) => {\n  await processText(chunk.text)\n})\n```\n\n### High-Throughput Writes: Using `IdempotentProducer` (Recommended)\n\nFor reliable, high-throughput writes with exactly-once semantics, use `IdempotentProducer`:\n\n```typescript\nimport { DurableStream, IdempotentProducer } from \"@durable-streams/client\"\n\nconst stream = await DurableStream.create({\n  url: \"https://streams.example.com/events\",\n  contentType: \"application/json\",\n})\n\nconst producer = new IdempotentProducer(stream, \"event-processor-1\", {\n  autoClaim: true,\n  onError: (err) => console.error(\"Batch failed:\", err), // Errors reported here\n})\n\n// Fire-and-forget - don't await, errors go to onError callback\nfor (const event of events) {\n  producer.append(event) // Objects serialized automatically for JSON streams\n}\n\n// IMPORTANT: Always flush before shutdown to ensure delivery\nawait producer.flush()\nawait producer.close()\n```\n\nFor high-throughput scenarios, `append()` is fire-and-forget (returns immediately):\n\n```typescript\n// Fire-and-forget - errors reported via onError callback\nfor (const event of events) {\n  producer.append(event) // Returns void, adds to batch\n}\n\n// Always flush before shutdown to ensure delivery\nawait producer.flush()\n```\n\n**Why use IdempotentProducer?**\n\n- **Exactly-once delivery**: Server deduplicates using `(producerId, epoch, seq)` tuple\n- **Automatic batching**: Multiple writes batched into single HTTP requests\n- **Pipelining**: Multiple batches in flight concurrently\n- **Zombie fencing**: Stale producers are rejected, preventing split-brain scenarios\n- **Network resilience**: Safe to retry on network errors (server deduplicates)\n\n### Read/Write: Using `DurableStream`\n\nFor simple write operations or when you need a persistent handle:\n\n```typescript\nimport { DurableStream } from \"@durable-streams/client\"\n\n// Create a new stream\nconst handle = await DurableStream.create({\n  url: \"https://streams.example.com/my-account/chat/room-1\",\n  headers: {\n    Authorization: `Bearer ${process.env.DS_TOKEN!}`,\n  },\n  contentType: \"application/json\",\n  ttlSeconds: 3600,\n})\n\n// Append data (simple API without exactly-once guarantees)\nawait handle.append(JSON.stringify({ type: \"message\", text: \"Hello\" }), {\n  seq: \"writer-1-000001\",\n})\n\n// Read using the new stream() API\nconst res = await handle.stream<{ type: string; text: string }>()\nres.subscribeJson(async (batch) => {\n  for (const item of batch.items) {\n    console.log(\"message:\", item.text)\n  }\n})\n```\n\n### Read from \"now\" (skip existing data)\n\n```typescript\n// HEAD gives you the current tail offset if the server exposes it\nconst handle = await DurableStream.connect({\n  url,\n  headers: { Authorization: `Bearer ${token}` },\n})\nconst { offset } = await handle.head()\n\n// Read only new data from that point on\nconst res = await handle.stream({ offset })\nres.subscribeBytes(async (chunk) => {\n  console.log(\"new data:\", new TextDecoder().decode(chunk.data))\n})\n```\n\n### Read catch-up only (no live updates)\n\n```typescript\n// Read existing data only, stop when up-to-date\nconst res = await stream({\n  url: \"https://streams.example.com/my-stream\",\n  live: false,\n})\n\nconst text = await res.text()\nconsole.log(\"All existing data:\", text)\n```\n\n## API\n\n### `stream(options): Promise<StreamResponse>`\n\nCreates a fetch-like streaming session:\n\n```typescript\nconst res = await stream<TJson>({\n  url: string | URL,              // Stream URL\n  headers?: HeadersRecord,        // Headers (static or function-based)\n  params?: ParamsRecord,          // Query params (static or function-based)\n  signal?: AbortSignal,           // Cancellation\n  fetch?: typeof fetch,           // Custom fetch implementation\n  backoffOptions?: BackoffOptions,// Retry backoff configuration\n  offset?: Offset,                // Starting offset (default: start of stream)\n  live?: LiveMode,                // Live mode (default: true)\n  json?: boolean,                 // Force JSON mode\n  onError?: StreamErrorHandler,   // Error handler\n})\n```\n\n### `DurableStream`\n\n```typescript\nclass DurableStream {\n  readonly url: string\n  readonly contentType?: string\n\n  constructor(opts: DurableStreamConstructorOptions)\n\n  // Static methods\n  static create(opts: CreateOptions): Promise<DurableStream>\n  static connect(opts: DurableStreamOptions): Promise<DurableStream>\n  static head(opts: DurableStreamOptions): Promise<HeadResult>\n  static delete(opts: DurableStreamOptions): Promise<void>\n\n  // Instance methods\n  head(opts?: { signal?: AbortSignal }): Promise<HeadResult>\n  create(opts?: CreateOptions): Promise<this>\n  delete(opts?: { signal?: AbortSignal }): Promise<void>\n  close(opts?: CloseOptions): Promise<CloseResult> // Close stream (EOF)\n  append(\n    body: BodyInit | Uint8Array | string,\n    opts?: AppendOptions\n  ): Promise<void>\n  appendStream(\n    source: AsyncIterable<Uint8Array | string>,\n    opts?: AppendOptions\n  ): Promise<void>\n\n  // Fetch-like read API\n  stream<TJson>(opts?: StreamOptions): Promise<StreamResponse<TJson>>\n}\n```\n\n### Live Modes\n\n```typescript\n// true (default): auto-select best live mode\n// - SSE for JSON streams, long-poll for binary\n// - Promise helpers (body/json/text): stop after upToDate\n// - Streams/subscribers: continue with live updates\n\n// false: catch-up only, stop at first upToDate\nconst res = await stream({ url, live: false })\n\n// \"long-poll\": explicit long-poll mode for live updates\nconst res = await stream({ url, live: \"long-poll\" })\n\n// \"sse\": explicit SSE mode for live updates\nconst res = await stream({ url, live: \"sse\" })\n```\n\n### Binary Streams with SSE\n\nFor binary content types (e.g., `application/octet-stream`), SSE mode requires the `encoding` option:\n\n```typescript\nconst stream = await DurableStream.create({\n  url: \"https://streams.example.com/my-binary-stream\",\n  contentType: \"application/octet-stream\",\n})\n\nconst response = await stream.read({\n  live: \"sse\",\n  encoding: \"base64\",\n})\n\nresponse.subscribe((chunk) => {\n  console.log(chunk.data) // Uint8Array - automatically decoded from base64\n})\n```\n\nThe client automatically decodes base64 data events before returning them. This is required for any content type other than `text/*` or `application/json` when using SSE mode.\n\n### Headers and Params\n\nHeaders and params support both static values and functions (sync or async) for dynamic values like authentication tokens.\n\n```typescript\n// Static headers\n{\n  headers: {\n    Authorization: \"Bearer my-token\",\n    \"X-Custom-Header\": \"value\",\n  }\n}\n\n// Function-based headers (sync)\n{\n  headers: {\n    Authorization: () => `Bearer ${getCurrentToken()}`,\n    \"X-Tenant-Id\": () => getCurrentTenant(),\n  }\n}\n\n// Async function headers (for refreshing tokens)\n{\n  headers: {\n    Authorization: async () => {\n      const token = await refreshToken()\n      return `Bearer ${token}`\n    }\n  }\n}\n\n// Mix static and function headers\n{\n  headers: {\n    \"X-Static\": \"always-the-same\",\n    Authorization: async () => `Bearer ${await getToken()}`,\n  }\n}\n\n// Query params work the same way\n{\n  params: {\n    tenant: \"static-tenant\",\n    region: () => getCurrentRegion(),\n    token: async () => await getSessionToken(),\n  }\n}\n```\n\n### Error Handling\n\n```typescript\nimport { stream, FetchError, DurableStreamError } from \"@durable-streams/client\"\n\nconst res = await stream({\n  url: \"https://streams.example.com/my-stream\",\n  headers: {\n    Authorization: \"Bearer my-token\",\n  },\n  onError: async (error) => {\n    if (error instanceof FetchError) {\n      if (error.status === 401) {\n        const newToken = await refreshAuthToken()\n        return { headers: { Authorization: `Bearer ${newToken}` } }\n      }\n    }\n    if (error instanceof DurableStreamError) {\n      console.error(`Stream error: ${error.code}`)\n    }\n    return {} // Retry with same params\n  },\n})\n```\n\n## StreamResponse Methods\n\nThe `StreamResponse` object provides multiple ways to consume stream data. All methods respect the `live` mode setting.\n\n### Promise Helpers\n\nThese methods accumulate data until the stream is up-to-date, then resolve.\n\n#### `body(): Promise<Uint8Array>`\n\nAccumulates all bytes until up-to-date.\n\n```typescript\nconst res = await stream({ url, live: false })\nconst bytes = await res.body()\nconsole.log(\"Total bytes:\", bytes.length)\n\n// Process as needed\nconst text = new TextDecoder().decode(bytes)\n```\n\n#### `json(): Promise<Array<TJson>>`\n\nAccumulates all JSON items until up-to-date. Only works with JSON content.\n\n```typescript\nconst res = await stream<{ id: number; name: string }>({\n  url,\n  live: false,\n})\nconst items = await res.json()\n\nfor (const item of items) {\n  console.log(`User ${item.id}: ${item.name}`)\n}\n```\n\n#### `text(): Promise<string>`\n\nAccumulates all text until up-to-date.\n\n```typescript\nconst res = await stream({ url, live: false })\nconst text = await res.text()\nconsole.log(\"Full content:\", text)\n```\n\n### ReadableStreams\n\nWeb Streams API for piping to other streams or using with streaming APIs. ReadableStreams can be consumed using either `getReader()` or `for await...of` syntax.\n\n> **Safari/iOS Compatibility**: The client ensures all returned streams are async-iterable by defining `[Symbol.asyncIterator]` on stream instances when missing. This allows `for await...of` consumption without requiring a global polyfill, while preserving `instanceof ReadableStream` behavior.\n>\n> **Derived streams**: Streams created via `.pipeThrough()` or similar transformations are NOT automatically patched. Use the exported `asAsyncIterableReadableStream()` helper:\n>\n> ```typescript\n> import { asAsyncIterableReadableStream } from \"@durable-streams/client\"\n>\n> const derived = res.bodyStream().pipeThrough(myTransform)\n> const iterable = asAsyncIterableReadableStream(derived)\n> for await (const chunk of iterable) { ... }\n> ```\n\n#### `bodyStream(): ReadableStream<Uint8Array> & AsyncIterable<Uint8Array>`\n\nRaw bytes as a ReadableStream.\n\n**Using `getReader()`:**\n\n```typescript\nconst res = await stream({ url, live: false })\nconst readable = res.bodyStream()\n\nconst reader = readable.getReader()\nwhile (true) {\n  const { done, value } = await reader.read()\n  if (done) break\n  console.log(\"Received:\", value.length, \"bytes\")\n}\n```\n\n**Using `for await...of`:**\n\n```typescript\nconst res = await stream({ url, live: false })\n\nfor await (const chunk of res.bodyStream()) {\n  console.log(\"Received:\", chunk.length, \"bytes\")\n}\n```\n\n**Piping to a file (Node.js):**\n\n```typescript\nimport { Readable } from \"node:stream\"\nimport { pipeline } from \"node:stream/promises\"\n\nconst res = await stream({ url, live: false })\nawait pipeline(\n  Readable.fromWeb(res.bodyStream()),\n  fs.createWriteStream(\"output.bin\")\n)\n```\n\n#### `jsonStream(): ReadableStream<TJson> & AsyncIterable<TJson>`\n\nIndividual JSON items as a ReadableStream.\n\n**Using `getReader()`:**\n\n```typescript\nconst res = await stream<{ id: number }>({ url, live: false })\nconst readable = res.jsonStream()\n\nconst reader = readable.getReader()\nwhile (true) {\n  const { done, value } = await reader.read()\n  if (done) break\n  console.log(\"Item:\", value)\n}\n```\n\n**Using `for await...of`:**\n\n```typescript\nconst res = await stream<{ id: number; name: string }>({ url, live: false })\n\nfor await (const item of res.jsonStream()) {\n  console.log(`User ${item.id}: ${item.name}`)\n}\n```\n\n#### `textStream(): ReadableStream<string> & AsyncIterable<string>`\n\nText chunks as a ReadableStream.\n\n**Using `getReader()`:**\n\n```typescript\nconst res = await stream({ url, live: false })\nconst readable = res.textStream()\n\nconst reader = readable.getReader()\nwhile (true) {\n  const { done, value } = await reader.read()\n  if (done) break\n  console.log(\"Text chunk:\", value)\n}\n```\n\n**Using `for await...of`:**\n\n```typescript\nconst res = await stream({ url, live: false })\n\nfor await (const text of res.textStream()) {\n  console.log(\"Text chunk:\", text)\n}\n```\n\n**Using with Response API:**\n\n```typescript\nconst res = await stream({ url, live: false })\nconst textResponse = new Response(res.textStream())\nconst fullText = await textResponse.text()\n```\n\n### Subscribers\n\nSubscribers provide callback-based consumption with backpressure. The next chunk isn't fetched until your callback's promise resolves. Returns an unsubscribe function.\n\n#### `subscribeJson(callback): () => void`\n\nSubscribe to JSON batches with metadata. Provides backpressure-aware consumption.\n\n```typescript\nconst res = await stream<{ event: string }>({ url, live: true })\n\nconst unsubscribe = res.subscribeJson(async (batch) => {\n  // Process items - next batch waits until this resolves\n  for (const item of batch.items) {\n    await processEvent(item)\n  }\n  await saveCheckpoint(batch.offset)\n})\n\n// Later: stop receiving updates\nsetTimeout(() => {\n  unsubscribe()\n}, 60000)\n```\n\n#### `subscribeBytes(callback): () => void`\n\nSubscribe to byte chunks with metadata.\n\n```typescript\nconst res = await stream({ url, live: true })\n\nconst unsubscribe = res.subscribeBytes(async (chunk) => {\n  console.log(\"Received bytes:\", chunk.data.length)\n  console.log(\"Offset:\", chunk.offset)\n  console.log(\"Up to date:\", chunk.upToDate)\n\n  await writeToFile(chunk.data)\n  await saveCheckpoint(chunk.offset)\n})\n```\n\n#### `subscribeText(callback): () => void`\n\nSubscribe to text chunks with metadata.\n\n```typescript\nconst res = await stream({ url, live: true })\n\nconst unsubscribe = res.subscribeText(async (chunk) => {\n  console.log(\"Text:\", chunk.text)\n  console.log(\"Offset:\", chunk.offset)\n\n  await appendToLog(chunk.text)\n})\n```\n\n### Lifecycle\n\n#### `cancel(reason?: unknown): void`\n\nCancel the stream session. Aborts any pending requests.\n\n```typescript\nconst res = await stream({ url, live: true })\n\n// Start consuming\nres.subscribeBytes(async (chunk) => {\n  console.log(\"Chunk:\", chunk)\n})\n\n// Cancel after 10 seconds\nsetTimeout(() => {\n  res.cancel(\"Timeout\")\n}, 10000)\n```\n\n#### `closed: Promise<void>`\n\nPromise that resolves when the session is complete or cancelled.\n\n```typescript\nconst res = await stream({ url, live: false })\n\n// Start consuming in background\nconst consumer = res.text()\n\n// Wait for completion\nawait res.closed\nconsole.log(\"Stream fully consumed\")\n```\n\n### State Properties\n\n```typescript\nconst res = await stream({ url })\n\nres.url // The stream URL\nres.contentType // Content-Type from response headers\nres.live // The live mode (true, \"long-poll\", \"sse\", or false)\nres.startOffset // The starting offset passed to stream()\nres.offset // Current offset (updates as data is consumed)\nres.cursor // Cursor for collapsing (if provided by server)\nres.upToDate // Whether we've caught up to the stream head\nres.streamClosed // Whether the stream is permanently closed (EOF)\n```\n\n---\n\n## DurableStream Methods\n\n### Static Methods\n\n#### `DurableStream.create(opts): Promise<DurableStream>`\n\nCreate a new stream on the server.\n\n```typescript\nconst handle = await DurableStream.create({\n  url: \"https://streams.example.com/my-stream\",\n  headers: {\n    Authorization: \"Bearer my-token\",\n  },\n  contentType: \"application/json\",\n  ttlSeconds: 3600, // Optional: auto-delete after 1 hour\n})\n\nawait handle.append('{\"hello\": \"world\"}')\n```\n\n#### `DurableStream.connect(opts): Promise<DurableStream>`\n\nConnect to an existing stream (validates it exists via HEAD).\n\n```typescript\nconst handle = await DurableStream.connect({\n  url: \"https://streams.example.com/my-stream\",\n  headers: {\n    Authorization: \"Bearer my-token\",\n  },\n})\n\nconsole.log(\"Content-Type:\", handle.contentType)\n```\n\n#### `DurableStream.head(opts): Promise<HeadResult>`\n\nGet stream metadata without creating a handle.\n\n```typescript\nconst metadata = await DurableStream.head({\n  url: \"https://streams.example.com/my-stream\",\n  headers: {\n    Authorization: \"Bearer my-token\",\n  },\n})\n\nconsole.log(\"Offset:\", metadata.offset)\nconsole.log(\"Content-Type:\", metadata.contentType)\n```\n\n#### `DurableStream.delete(opts): Promise<void>`\n\nDelete a stream without creating a handle.\n\n```typescript\nawait DurableStream.delete({\n  url: \"https://streams.example.com/my-stream\",\n  headers: {\n    Authorization: \"Bearer my-token\",\n  },\n})\n```\n\n### Instance Methods\n\n#### `head(opts?): Promise<HeadResult>`\n\nGet metadata for this stream.\n\n```typescript\nconst handle = new DurableStream({\n  url,\n  headers: { Authorization: `Bearer ${token}` },\n})\nconst metadata = await handle.head()\n\nconsole.log(\"Current offset:\", metadata.offset)\n```\n\n#### `create(opts?): Promise<this>`\n\nCreate this stream on the server.\n\n```typescript\nconst handle = new DurableStream({\n  url,\n  headers: { Authorization: `Bearer ${token}` },\n})\nawait handle.create({\n  contentType: \"text/plain\",\n  ttlSeconds: 7200,\n})\n```\n\n#### `delete(opts?): Promise<void>`\n\nDelete this stream.\n\n```typescript\nconst handle = new DurableStream({\n  url,\n  headers: { Authorization: `Bearer ${token}` },\n})\nawait handle.delete()\n```\n\n#### `append(body, opts?): Promise<void>`\n\nAppend data to the stream. By default, **automatic batching is enabled**: multiple `append()` calls made while a POST is in-flight will be batched together into a single request. This significantly improves throughput for high-frequency writes.\n\n```typescript\nconst handle = await DurableStream.connect({\n  url,\n  headers: { Authorization: `Bearer ${token}` },\n})\n\n// Append string\nawait handle.append(\"Hello, world!\")\n\n// Append with sequence number for ordering\nawait handle.append(\"Message 1\", { seq: \"writer-1-001\" })\nawait handle.append(\"Message 2\", { seq: \"writer-1-002\" })\n\n// For JSON streams, append objects directly (serialized automatically)\nawait handle.append({ event: \"click\", x: 100, y: 200 })\n\n// Batching happens automatically - these may be sent in a single request\nawait Promise.all([\n  handle.append({ event: \"msg1\" }),\n  handle.append({ event: \"msg2\" }),\n  handle.append({ event: \"msg3\" }),\n])\n```\n\n**Batching behavior:**\n\n- **JSON mode** (`contentType: \"application/json\"`): Multiple values are sent as a JSON array `[val1, val2, ...]`\n- **Byte mode**: Binary data is concatenated\n\n**Disabling batching:**\n\nIf you need to ensure each append is sent immediately (e.g., for precise timing or debugging):\n\n```typescript\nconst handle = new DurableStream({\n  url,\n  batching: false, // Disable automatic batching\n})\n```\n\n#### `appendStream(source, opts?): Promise<void>`\n\nAppend streaming data from an async iterable or ReadableStream. This method supports piping from any source.\n\n```typescript\nconst handle = await DurableStream.connect({\n  url,\n  headers: { Authorization: `Bearer ${token}` },\n})\n\n// From async generator\nasync function* generateData() {\n  for (let i = 0; i < 100; i++) {\n    yield `Line ${i}\\n`\n  }\n}\nawait handle.appendStream(generateData())\n\n// From ReadableStream\nconst readable = new ReadableStream({\n  start(controller) {\n    controller.enqueue(\"chunk 1\")\n    controller.enqueue(\"chunk 2\")\n    controller.close()\n  },\n})\nawait handle.appendStream(readable)\n\n// Pipe from a fetch response body\nconst response = await fetch(\"https://example.com/data\")\nawait handle.appendStream(response.body!)\n```\n\n#### `writable(opts?): WritableStream<Uint8Array | string>`\n\nCreate a WritableStream that can receive piped data. Useful for stream composition:\n\n```typescript\nconst handle = await DurableStream.connect({ url, auth })\n\n// Pipe from any ReadableStream\nawait someReadableStream.pipeTo(handle.writable())\n\n// Pipe through a transform\nconst readable = inputStream.pipeThrough(new TextEncoderStream())\nawait readable.pipeTo(handle.writable())\n```\n\n#### `stream(opts?): Promise<StreamResponse>`\n\nStart a read session (same as standalone `stream()` function).\n\n```typescript\nconst handle = await DurableStream.connect({\n  url,\n  headers: { Authorization: `Bearer ${token}` },\n})\n\nconst res = await handle.stream<{ message: string }>({\n  offset: savedOffset,\n  live: true,\n})\n\nres.subscribeJson(async (batch) => {\n  for (const item of batch.items) {\n    console.log(item.message)\n  }\n})\n```\n\n---\n\n## IdempotentProducer\n\nThe `IdempotentProducer` class provides Kafka-style exactly-once write semantics with automatic batching and pipelining.\n\n### Constructor\n\n```typescript\nnew IdempotentProducer(stream: DurableStream, producerId: string, opts?: IdempotentProducerOptions)\n```\n\n**Parameters:**\n\n- `stream` - The DurableStream to write to\n- `producerId` - Stable identifier for this producer (e.g., \"order-service-1\")\n- `opts` - Optional configuration\n\n**Options:**\n\n```typescript\ninterface IdempotentProducerOptions {\n  epoch?: number // Starting epoch (default: 0)\n  autoClaim?: boolean // On 403, retry with epoch+1 (default: false)\n  maxBatchBytes?: number // Max bytes before sending batch (default: 1MB)\n  lingerMs?: number // Max time to wait for more messages (default: 5ms)\n  maxInFlight?: number // Concurrent batches in flight (default: 5)\n  signal?: AbortSignal // Cancellation signal\n  fetch?: typeof fetch // Custom fetch implementation\n  onError?: (error: Error) => void // Error callback for batch failures\n}\n```\n\n### Methods\n\n#### `append(body): void`\n\nAppend data to the stream (fire-and-forget). For JSON streams, you can pass objects directly.\nReturns immediately after adding to the internal batch. Errors are reported via `onError` callback.\n\n```typescript\n// For JSON streams - pass objects directly\nproducer.append({ event: \"click\", x: 100 })\n\n// Or strings/bytes\nproducer.append(\"message data\")\nproducer.append(new Uint8Array([1, 2, 3]))\n\n// All appends are fire-and-forget - use flush() to wait for delivery\nawait producer.flush()\n```\n\n#### `flush(): Promise<void>`\n\nSend any pending batch immediately and wait for all in-flight batches to complete.\n\n```typescript\n// Always call before shutdown\nawait producer.flush()\n```\n\n#### `close(finalMessage?): Promise<CloseResult>`\n\nFlush pending messages and close the underlying **stream** (EOF). This is the typical way to end a producer session:\n\n1. Flushes all pending messages\n2. Optionally appends a final message atomically with close\n3. Closes the stream (no further appends permitted by any producer)\n\n**Idempotent**: Safe to retry on network failures - uses producer headers for deduplication.\n\n```typescript\n// Close stream (EOF)\nconst result = await producer.close()\nconsole.log(\"Final offset:\", result.finalOffset)\n\n// Close with final message (atomic append + close)\nconst result = await producer.close('{\"done\": true}')\n```\n\n#### `detach(): Promise<void>`\n\nStop the producer without closing the underlying stream. Use this when:\n\n- Handing off writing to another producer\n- Keeping the stream open for future writes\n- Stopping this producer but not signaling EOF to readers\n\n```typescript\nawait producer.detach() // Stream remains open\n```\n\n#### `restart(): Promise<void>`\n\nIncrement epoch and reset sequence. Call this when restarting the producer to establish a new session.\n\n```typescript\nawait producer.restart()\n```\n\n### Properties\n\n- `epoch: number` - Current epoch for this producer\n- `nextSeq: number` - Next sequence number to be assigned\n- `pendingCount: number` - Messages in the current pending batch\n- `inFlightCount: number` - Batches currently in flight\n\n### Error Handling\n\nErrors are delivered via the `onError` callback since `append()` is fire-and-forget:\n\n```typescript\nimport {\n  IdempotentProducer,\n  StaleEpochError,\n  SequenceGapError,\n} from \"@durable-streams/client\"\n\nconst producer = new IdempotentProducer(stream, \"my-producer\", {\n  onError: (error) => {\n    if (error instanceof StaleEpochError) {\n      // Another producer has a higher epoch - this producer is \"fenced\"\n      console.log(`Fenced by epoch ${error.currentEpoch}`)\n    } else if (error instanceof SequenceGapError) {\n      // Sequence gap detected (should not happen with proper usage)\n      console.log(`Expected seq ${error.expectedSeq}, got ${error.receivedSeq}`)\n    }\n  },\n})\n\nproducer.append(\"data\") // Fire-and-forget, errors go to onError\nawait producer.flush() // Wait for all batches to complete\n```\n\n---\n\n## Stream Closure (EOF)\n\nDurable Streams supports permanently closing streams to signal EOF (End of File). Once closed, no further appends are permitted, but data remains fully readable.\n\n### Writer Side\n\n#### Using DurableStream.close()\n\n```typescript\nconst stream = await DurableStream.connect({ url })\n\n// Simple close (no final message)\nconst result = await stream.close()\nconsole.log(\"Final offset:\", result.finalOffset)\n\n// Atomic append-and-close with final message\nconst result = await stream.close({\n  body: '{\"status\": \"complete\"}',\n})\n```\n\n**Options:**\n\n```typescript\ninterface CloseOptions {\n  body?: Uint8Array | string // Optional final message\n  contentType?: string // Content type (must match stream)\n  signal?: AbortSignal // Cancellation\n}\n\ninterface CloseResult {\n  finalOffset: Offset // The offset after the last byte\n}\n```\n\n**Idempotency:**\n\n- `close()` without body: Idempotent — safe to call multiple times\n- `close({ body })` with body: NOT idempotent — throws `StreamClosedError` if already closed. Use `IdempotentProducer.close(finalMessage)` for idempotent close-with-body.\n\n#### Using IdempotentProducer.close()\n\nFor reliable close with final message (safe to retry):\n\n```typescript\nconst producer = new IdempotentProducer(stream, \"producer-1\", {\n  autoClaim: true,\n})\n\n// Write some messages\nproducer.append('{\"event\": \"start\"}')\nproducer.append('{\"event\": \"data\"}')\n\n// Close with final message (idempotent, safe to retry)\nconst result = await producer.close('{\"event\": \"end\"}')\n```\n\n**Important:** `IdempotentProducer.close()` closes the **stream**, not just the producer. Use `detach()` to stop the producer without closing the stream.\n\n#### Creating Closed Streams\n\nCreate a stream that's immediately closed (useful for cached responses, errors, single-shot data):\n\n```typescript\n// Empty closed stream\nconst stream = await DurableStream.create({\n  url: \"https://streams.example.com/cached-response\",\n  contentType: \"application/json\",\n  closed: true,\n})\n\n// Closed stream with initial content\nconst stream = await DurableStream.create({\n  url: \"https://streams.example.com/error-response\",\n  contentType: \"application/json\",\n  body: '{\"error\": \"Service unavailable\"}',\n  closed: true,\n})\n```\n\n### Reader Side\n\n#### Detecting Closure\n\nThe `streamClosed` property indicates when a stream is permanently closed:\n\n```typescript\n// StreamResponse properties\nconst res = await stream({ url, live: true })\nconsole.log(res.streamClosed) // false initially\n\n// In subscribers - batch/chunk metadata includes streamClosed\nres.subscribeJson((batch) => {\n  console.log(\"Items:\", batch.items)\n  console.log(\"Stream closed:\", batch.streamClosed) // true when EOF reached\n})\n\n// In HEAD requests\nconst metadata = await stream.head()\nconsole.log(\"Stream closed:\", metadata.streamClosed)\n```\n\n#### Live Mode Behavior\n\nWhen a stream is closed:\n\n- **Long-poll**: Returns immediately with `streamClosed: true` (no waiting)\n- **SSE**: Sends `streamClosed: true` in final control event, then closes connection\n- **Subscribers**: Receive final batch with `streamClosed: true`, then stop\n\n```typescript\nconst res = await stream({ url, live: true })\n\nres.subscribeJson((batch) => {\n  for (const item of batch.items) {\n    process(item)\n  }\n\n  if (batch.streamClosed) {\n    console.log(\"Stream complete, no more data will arrive\")\n    // Connection will close automatically\n  }\n})\n```\n\n### Error Handling\n\nAttempting to append to a closed stream throws `StreamClosedError`:\n\n```typescript\nimport { StreamClosedError } from \"@durable-streams/client\"\n\ntry {\n  await stream.append(\"data\")\n} catch (error) {\n  if (error instanceof StreamClosedError) {\n    console.log(\"Stream is closed at offset:\", error.finalOffset)\n  }\n}\n```\n\n---\n\n## Types\n\nKey types exported from the package:\n\n- `Offset` - Opaque string for stream position\n- `StreamResponse` - Response object from stream() (includes `streamClosed` property)\n- `ByteChunk` - `{ data: Uint8Array, offset: Offset, upToDate: boolean, streamClosed: boolean, cursor?: string }`\n- `JsonBatch<T>` - `{ items: T[], offset: Offset, upToDate: boolean, streamClosed: boolean, cursor?: string }`\n- `TextChunk` - `{ text: string, offset: Offset, upToDate: boolean, streamClosed: boolean, cursor?: string }`\n- `HeadResult` - Metadata from HEAD requests (includes `streamClosed` property)\n- `CloseOptions` - Options for closing a stream\n- `CloseResult` - Result from closing a stream (includes `finalOffset`)\n- `IdempotentProducer` - Exactly-once producer class\n- `StaleEpochError` - Thrown when producer epoch is stale (zombie fencing)\n- `SequenceGapError` - Thrown when sequence numbers are out of order\n- `StreamClosedError` - Thrown when attempting to append to a closed stream (includes `finalOffset`)\n- `DurableStreamError` - Protocol-level errors with codes\n- `FetchError` - Transport/network errors\n\n## License\n\nApache-2.0\n","readmeFilename":"README.md"}