> For the complete documentation index, see [llms.txt](https://docs.collieai.io/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://docs.collieai.io/sdks/node-sdk.md).

# Node SDK

The `@collieai/sdk` Node/TypeScript SDK is the recommended way to do customer-owned streaming from a Node backend. You run your own model; the SDK checks the prompt before you call it and lets you stream back **only CollieAi-released text** — never raw model output. It owns the chunk protocol (sequence numbers, retries, idempotency, finalization), so you don't touch the [raw chunk endpoints](/async-jobs/customer-owned-streaming.md). It mirrors the [Python SDK](/sdks/python-sdk.md) — same event model, same typed errors — and has **zero runtime dependencies** (it uses the global `fetch`; Node ≥ 18).

{% hint style="warning" %}
**Forward `event.text`, never the raw model delta.** The whole point of the SDK is that your users only ever see text CollieAi has released. The examples below keep the raw provider stream inside the `rawStreamFactory`; you only forward SDK events.
{% endhint %}

## Install

```bash
npm install @collieai/sdk
```

## Construct the client

```ts
import { CollieClient } from "@collieai/sdk";

const collie = new CollieClient({
  apiKey: "clai_...",
  baseUrl: "https://app.collieai.io",
  projectId: "project_123",
});
```

## Check an input before calling your LLM

```ts
const result = await collie.moderate.input({
  prompt: userPrompt,
  conversationId,          // optional: groups a conversation
  correlationId: turnId,   // optional: pins one turn
});

if (result.blocked) return result.blockMessage ?? "Input blocked by policy.";

// If your input policy MASKS, send the FILTERED prompt to your model —
// null-check, never truthiness: "" is a legitimate full wipe.
const promptForModel = result.filteredText ?? userPrompt;
```

A policy block is a normal result (`result.blocked === true`), not an error. No webhook is required — the SDK polls for you. `filteredText` is the post-mask prompt: the raw values a mask rule removed must never reach your model, so `promptForModel` — not `userPrompt` — goes into your LLM call. When the wrapper's own gate masks a prompt, `protectStream` / `protectBuffered` throw `MaskedInputError` before your factory runs; the recipe is to gate manually as above, build the factory over `result.filteredText`, and pass `inputResult: result` (leave `checkInput` unset/true — combining `inputResult` with `checkInput: false` is a contradiction and throws `TypeError`; in 2.0 it was silently ignored).

## Analyze context alongside the prompt

Pass `context` to analyze the data you're about to send the model — retrieved documents, tool output, a record — alongside the prompt. Structured `context` travels as JSON; `contextFormat` (`"auto"` | `"json"` | `"text"`) is an optional hint for a raw string.

```ts
const result = await collie.moderate.input({
  prompt: userPrompt,
  context: { transaction: { memo: retrievedMemo } },
});

const ctx = result.context;
if (ctx && ctx.status !== "clean" && ctx.status !== "not_provided") {
  console.log(ctx.status, ctx.triggeringPointer, ctx.triggeringRuleType); // path, never a value
}
if (result.blocked) return result.blockMessage ?? "Blocked by policy.";
// result.blockedBy is "prompt" | "context" | "none"
```

`result.context` carries the closed `status` enum, the triggering JSON Pointer + rule, and degraded markers; the pointer is a **path, never a value**. The same `context` works on `protectStream` / `protectBuffered` when input checking is enabled. See [Context analysis](/security-rules/context-analysis.md).

{% hint style="info" %}
Until context analysis is enabled (host flag **and** policy switch), `context` is inert — a safe no-op.
{% endhint %}

## Check an output produced outside a wrapper

`protectStream` / `protectBuffered` already moderate the streamed answer. Use `moderate.output` for assistant text that never went through a wrapper — proactive notifications, escalation messages, any side channel. Don't route such text through `moderate.input`: that evaluates it with **input** rules, so output-safety and masking rules silently never run, and injection detectors can false-block assistant-style imperatives ("You need to submit…").

```ts
const result = await collie.moderate.output({
  response: assistantText,
  conversationId,          // optional
  correlationId: noteId,   // optional
});

if (result.blocked) return; // don't send it
send(result.filteredText ?? assistantText); // masking applies here
```

`filteredText` carries the **masked** output — send it, not the original, or masking rules silently do nothing. `moderate.output` takes no `context`: context is an input surface; for context-aware output filtering use `protectBuffered` / `protectStream`.

## Stream safely (Express)

`protectStream` checks the input, calls your LLM **only if it passes**, batches the output, and yields only safe events. Pass a *factory* (a function returning a fresh async iterable; it may accept the SDK-provided `AbortSignal`), not an already-started stream — the SDK calls it once, after the check.

By default it streams optimistically — if the policy can't stream, the first push throws `ChunkStreamingUnsupported`. To choose the UX *before* generating, use [preflight](#choose-the-ux-up-front-preflight) or pass `requireStreaming: true`.

```ts
import express from "express";
import OpenAI from "openai";
import { CollieClient } from "@collieai/sdk";

const collie = new CollieClient({ apiKey: process.env.COLLIEAI_API_KEY! });
const openai = new OpenAI();
const app = express();
app.use(express.json());

async function* openaiDeltas(prompt: string, signal?: AbortSignal) {
  const stream = await openai.chat.completions.create(
    {
      model: "gpt-4o-mini",
      messages: [{ role: "user", content: prompt }],
      stream: true,
    },
    signal ? { signal } : undefined,
  );
  try {
    for await (const part of stream as AsyncIterable<any>) {
      const delta = part.choices?.[0]?.delta?.content;
      if (delta) yield delta;
    }
  } finally {
    stream.controller?.abort(); // stop generation if abandoned early
  }
}

app.post("/chat", async (req, res) => {
  res.setHeader("content-type", "text/plain");
  for await (const event of collie.streaming.protectStream({
    input: req.body.prompt,
    rawStreamFactory: (signal) => openaiDeltas(req.body.prompt, signal),
  })) {
    if (event.type === "delta") res.write(event.text);
    else if (event.type === "blocked" || event.type === "input_blocked") {
      res.write(event.blockMessage ?? "Blocked by policy.");
      break;
    }
  }
  res.end();
});
```

That's the whole integration. You never touch chunk sequence numbers, retries, or batching. The same loop works in a Next.js App Router route handler — return a `ReadableStream` that enqueues `event.text`.

{% hint style="info" %}
**Provider adapters.** Instead of writing `openaiDeltas` by hand, use the packaged factory (the provider SDK stays an optional peer dependency):

```ts
import { openaiFactory } from "@collieai/sdk/adapters/openai";
// or: import { anthropicFactory } from "@collieai/sdk/adapters/anthropic";

const factory = openaiFactory(openai, {
  model: "gpt-4o-mini",
  messages: [{ role: "user", content: prompt }],
});
for await (const event of collie.streaming.protectStream({ input: prompt, rawStreamFactory: factory })) {
  // ...
}
```

The factory opens the provider stream lazily (only after the input check passes) and closes it on early exit.
{% endhint %}

## Choose the UX up front (preflight)

Ask whether the policy can stream **before** calling the LLM, and branch into a token-stream UI or a "checking response…" UI:

```ts
const cap = await collie.streaming.preflight(); // cached until cap.validUntil

if (cap.recommendedClientBehavior === "stream") {
  // protectStream → token-stream UI
} else if (cap.recommendedClientBehavior === "buffer_then_show") {
  const result = await collie.streaming.protectBuffered({
    input: userPrompt,
    rawStreamFactory: (signal) => openaiDeltas(userPrompt, signal),
  });
  // "checking response…", then show result
} else {
  throw new Error(cap.reasonDetail ?? cap.reason ?? "policy unavailable"); // fail_fast
}
```

Prefer not to branch yourself? Pass `requireStreaming: true` to `protectStream`: it preflights first and throws `BufferedFallbackRequired` (policy must buffer) or a `PreflightError` (can't be served) **before** calling your LLM.

### Which rule is responsible — `cap.rules`

`cap.rules` explains the verdict rule by rule. Each entry carries `executionRole`, the role that rule plays in a streaming request:

| `executionRole`       | meaning                                                                                          |
| --------------------- | ------------------------------------------------------------------------------------------------ |
| `enforce_streaming`   | blocks or masks **as tokens arrive**                                                             |
| `enforce_postflight`  | blocks or masks, but needs the **whole** response first — a policy containing one always buffers |
| `stream_observed`     | monitor-only; observed **as tokens arrive** while the text streams through untouched             |
| `postflight_observed` | monitor-only; observed after the response completes                                              |
| `null`                | the rule could not be planned at all (e.g. an unrecognised type)                                 |

```ts
const blockers = cap.rules.filter((r) => r.executionRole === "enforce_postflight");
if (blockers.length) {
  console.log("buffered because:", blockers.map((r) => r.ruleName).join(", "));
}
```

{% hint style="info" %}
`streamingSupported` is `false` for **every** monitor rule, so it cannot tell a `stream_observed` rule from a `postflight_observed` one. Use `executionRole` when you need that distinction — typically while running a policy in monitor mode before turning on blocking.
{% endhint %}

Three things `executionRole` does not tell you. It is **per-rule** and ignores policy-level gates, so a rule can say `stream_observed` while the request still buffers — `cap.mode` remains the answer to "does this request stream at all". It reflects your project's current streaming setting, so a project set to `buffered` reports full-context roles throughout. And it reflects the **server that answered**: on current servers a monitor-only policy streams — monitor rules observe the text as it passes and their findings appear in your CollieAi audit log only, never in the chunk responses your users see. A server that has not been upgraded yet still buffers any policy containing a monitor rule and answers `cap.mode = "buffered"`, `cap.reason = "monitor_mode"`; the SDK handles both, and `cap.mode` is always the delivery verdict. The roles are what let you see, ahead of time, whether monitor→enforce will be a config flip or a re-integration.

The `ExecutionRole` type keeps an open `string` arm on purpose: compare against the names you know and fall through on anything else, so a role added later doesn't break your branch.

## Buffered fallback

When a policy can't stream, check the whole response at once — same input gate and factory contract, but it returns a single result instead of events:

```ts
const result = await collie.streaming.protectBuffered({
  input: userPrompt,
  rawStreamFactory: (signal) => openaiDeltas(userPrompt, signal),
});
return result.blocked ? result.blockMessage : result.filteredText;
```

One difference from `protectStream` matters if you reuse an input result: on the streaming path a passed `inputResult` carrying a `jobId` becomes, against a 2.1+ server, a **server-verified, expiring, single-use claim** (the `input_gate_*` 409s can surface; without a `jobId` or on an older server the session runs claimless and re-filters the input in its own async pass — the pre-2.1 shape: billed, and the verdict can land after streaming has started); on the buffered path it is a **trusted-client reuse** — the server neither verifies nor consumes it, and no `input_gate_*` error can occur. Because nothing server-side checks freshness there, obtain the result in the same turn, from the same client and project, immediately before the call — never cache or reuse it.

## Relay safe output to a browser (SSE)

When your backend submits chunks for a job, it can also **subscribe** to that job's CollieAi SSE stream and relay safe events to a browser — with automatic reconnect. `session.streamEvents()` yields the same typed events as `protectStream`, plus `interrupted` on idle-timeout / disconnect.

```ts
import { toSse } from "@collieai/sdk";

for await (const event of session.streamEvents()) {  // auto-resumes from Last-Event-ID
  if (event.type === "interrupted") continue;        // reconnecting; nothing to forward
  res.write(toSse(event));                           // SSE bytes for text/event-stream
  if (event.type === "blocked" || event.type === "finished") break;
}
```

On a dropped connection or server idle-timeout the SDK reconnects from the last seen event id and **deduplicates replayed frames**, so a delta is never shown twice. To handle resume yourself, pass `autoResume: false` — iteration then stops at the first `interrupted` event, whose `lastEventId` you pass back to `streamEvents({ lastEventId })` later.

### Let a browser subscribe directly (stream tokens)

To skip the backend relay, mint a short-lived, job-scoped **stream token** and hand it to the browser — your API key stays server-side:

```ts
// Backend: mint and return a token for the browser.
const st = await session.mintStreamToken();
return { url: st.url, expiresIn: st.expiresIn };
```

```javascript
// Browser: subscribe with the token in the query string.
const es = new EventSource(url); // .../v1/jobs/{id}/stream?stream_token=...
es.addEventListener("chunk", (e) => {
  const chunk = JSON.parse(e.data);
  if (chunk.content) append(chunk.content);
  if (chunk.blocked) {
    append(chunk.block_message ?? "Blocked by policy.");
    es.close();
  }
});
es.addEventListener("end", () => es.close());
```

{% hint style="info" %}
**Token scope & delegation.** A stream token is read-only and valid only for its one job's stream — it can't submit chunks or call any other endpoint, and it can't be used for a different job. *Minting* is fully authenticated (active user

* IP allowlist, via your API key); the token is then a **delegated capability** that the browser redeems from its own IP, so it is *not* re-checked against the IP allowlist. Keep the API key server-side; give the browser only the token, and re-mint before it expires. Operators can disable browser tokens entirely with `STREAM_TOKEN_ENABLED=false` (then use the backend relay above).
  {% endhint %}

{% hint style="warning" %}
**The token is a credential in the URL.** Treat the full `?stream_token=...` URL as secret: serve it over HTTPS, don't log it, and redact `stream_token` at your proxy/access-log layer. A direct browser `EventSource` to CollieAi is **cross-origin** — your site's origin must be in the API's CORS allowlist, or use the backend relay above instead.
{% endhint %}

## Errors

Catch typed errors instead of parsing strings. All extend `CollieError`. Policy blocks are **events** (`type: "blocked"` / `"input_blocked"`), not errors.

| Error                                                                                                                             | When                                                                                                                      | What to do                                                                                                                                                                         |
| --------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `ChunkRetryExhausted`                                                                                                             | transient failures exceeded the retry budget                                                                              | retry the turn, or fall back to a non-streaming path                                                                                                                               |
| `ChunkPolicyChanged`                                                                                                              | policy changed mid-stream                                                                                                 | start a new stream/session (it picks up the new policy)                                                                                                                            |
| `ChunkSessionUnrecoverable`                                                                                                       | the stream entered an unrepairable state                                                                                  | start a new stream/session                                                                                                                                                         |
| `ChunkQuotaExceeded`                                                                                                              | rate-limited with no usable `Retry-After`                                                                                 | back off and retry later                                                                                                                                                           |
| `ChunkStreamingUnsupported`                                                                                                       | the policy can't be served by the streaming engine                                                                        | use `protectBuffered` instead                                                                                                                                                      |
| `ChunkResolutionUnavailable`                                                                                                      | `503 chunk_resolution_unavailable` — the server couldn't **resolve** the policy (a dependency outage, not a policy shape) | nothing at first: the SDK retries it automatically; if the outage outlasts the retry budget it is the `cause` of `ChunkRetryExhausted`                                             |
| `BufferedFallbackRequired`                                                                                                        | `requireStreaming: true` but the policy must buffer                                                                       | switch to `protectBuffered`                                                                                                                                                        |
| `ProjectNotFound` / `StreamingFeatureDisabled` / `PlanNotEntitled` / `UnknownRuleType` / `PolicyNotStreamable` (`PreflightError`) | preflight says the policy can't be served                                                                                 | fix project/policy configuration; other reason codes (e.g. `rule_unplannable`, `resolution_error:<cause>`) surface as the base `PreflightError` — treat the code as an open string |
| `ProviderStreamFactoryRequired`                                                                                                   | passed a started stream (or non-async-iterable) instead of a factory                                                      | pass a factory returning a fresh async iterable: `(signal) => myStream(signal)`                                                                                                    |
| `ConcurrentSessionUseError`                                                                                                       | overlapping `push()` calls on one low-level session                                                                       | serialize submits per session                                                                                                                                                      |
| `ModerationError`                                                                                                                 | a `moderate.input` / `moderate.output` job failed/expired or timed out                                                    | retry the check                                                                                                                                                                    |
| `MaskedInputError`                                                                                                                | the wrapper's own input gate MASKED the prompt — streaming would send the unmasked original to your model                 | gate manually with `moderate.input`, build the factory over `result.filteredText` (`""` is a legitimate full wipe), pass `inputResult: result`                                     |
| `CollieConnectionError`                                                                                                           | transport failure (timeout, connection refused)                                                                           | retry; check connectivity to `baseUrl`                                                                                                                                             |
| `CollieApiError`                                                                                                                  | unexpected HTTP error or malformed response (`code === "invalid_response"`)                                               | inspect `statusCode`/`code`; retry or report                                                                                                                                       |

## Retry behavior

The SDK retries the **same** chunk sequence on transient failures (network timeouts, `503` — including the typed `chunk_resolution_unavailable` — `504 chunk_filter_timeout`, `429` with a usable `Retry-After`) with exponential backoff + jitter — default base 250 ms, max 4 s, 3 attempts per chunk, 10 s ceiling (`new CollieClient({ retryMaxPerChunkS: ... })`). A retried chunk never produces a duplicate visible delta.

## Failure policy: fail-closed by default

When retries are exhausted the SDK **fails closed** — it throws and releases no text. There is no fail-open switch: shipping unchecked model output is a risk decision only you can make, so it has to be your explicit code, not a default.

Two things decide what that code should do.

**1. Where you caught it.** `rawStreamFactory` is a deferred factory: the SDK calls it exactly once, *after* the input check passes and the session exists. So if CollieAi is unreachable, the failure lands before your LLM ever ran — nothing was generated, nothing reached the user, and passing through is a clean choice. A failure *mid-stream* is different: some text is already on screen and the engine is holding more, so re-running the model would duplicate output. Track it with one flag.

**2. Not every `CollieError` is an outage.** `ChunkStreamingUnsupported` means the policy demands buffered checking, and `ChunkRetryExhausted` means a mid-stream abort. Treating the base class as an outage would bypass a policy that deliberately asked to inspect the whole answer. Match the transport failures — `CollieConnectionError` and `CollieApiError` with a 5xx — and nothing else.

```ts
import {
  BufferedFallbackRequired, ChunkStreamingUnsupported,
  CollieApiError, CollieConnectionError, CollieError,
} from "@collieai/sdk";

let started = false;
const rawStreamFactory = () => customerLlm(prompt);

try {
  for await (const ev of collie.streaming.protectStream({ input: prompt, rawStreamFactory })) {
    if (ev.type === "delta") {
      started = true;
      await write(ev.text);
    } else if (ev.type === "input_blocked") {
      await write(ev.blockMessage ?? "Request rejected."); return;
    } else if (ev.type === "blocked") {
      await write(ev.blockMessage ?? "Response rejected."); return;
    }
  }
} catch (e) {
  // Fail-open ONLY on infrastructure failure AND only before the first delta.
  const infra =
    e instanceof CollieConnectionError ||
    (e instanceof CollieApiError && (e.statusCode ?? 0) >= 500);
  if (infra && FAIL_OPEN_ALLOWED && !started) {
    for await (const delta of customerLlm(prompt)) await write(delta); // unfiltered
    return;
  }
  // Policy demands buffered checking -- an instruction, not an outage. The
  // !started guard avoids re-emitting an answer the user has already partly seen.
  // NOTE: protectBuffered re-invokes your LLM — a second, billable generation.
  if ((e instanceof BufferedFallbackRequired || e instanceof ChunkStreamingUnsupported) && !started) {
    const r = await collie.streaming.protectBuffered({ input: prompt, rawStreamFactory });
    await write(r.blocked ? (r.blockMessage ?? "Response rejected.") : (r.filteredText ?? ""));
    return;
  }
  // Everything else, including a mid-stream abort.
  if (e instanceof CollieError) {
    await write("Could not verify the response. Please try again.");
    return;
  }
  throw e;
}
```

An `AbortError` from your own `AbortSignal` is not a `CollieError`, so a client disconnect never triggers fail-open.

If you need fail-open **mid-stream** too, you must tee the provider stream: wrap your LLM iterable so deltas also accumulate in a local buffer, then on failure emit `buffer.slice(alreadyWrittenChars)` and continue raw. Without the tee the text the engine was still holding is lost and the answer has a hole in the middle.

Whatever you choose, make fail-open **observable**: log every occurrence and put it behind a flag you can turn off. A silent `catch` means an outage can leave you serving unchecked model output for hours with everything looking healthy.

## Advanced: low-level session

If you need to drive batching yourself, use the session directly. It does **not** check the input — call `moderate.input(...)` first.

```ts
const check = await collie.moderate.input({ prompt: userPrompt });
if (check.blocked) return check.blockMessage ?? "Input blocked by policy.";

// The split that matters when your input policy masks:
// - the SESSION gets the ORIGINAL prompt (it must byte-match the gate);
// - your MODEL gets the FILTERED prompt ("" is a legitimate full wipe).
const promptForModel = check.filteredText ?? userPrompt;

// inputJobId: proves the prompt was just gated, so the server skips
// re-filtering it on the session job (one input pass per turn).
const session = collie.streaming.session({
  input: userPrompt,
  inputJobId: check.jobId ?? undefined,
});
await session.open();
try {
  for await (const rawDelta of yourLlmStream(promptForModel)) {
    const result = await session.push(rawDelta);
    for (const emit of result.emits) forward(emit.text); // forward ONLY safe emits
    if (result.finished) break;
  }
  await session.finish();
} finally {
  await session.close();
}
```

Manual sessions do not auto-retry the claim protocol: `input_gate_stale` means the policy changed since your gate ran — re-run `moderate.input` (with the same context, if any) and open a new session with the fresh job id; `input_gate_unverifiable` means retry without the gate reference; `input_gate_claimed` means that gate was already consumed. `protectStream` handles all of this for you — but ONLY when it runs its own gate: with an external `inputResult` (this very recipe) the typed errors surface to YOUR code by design, since the wrapper cannot re-check a context it never saw.
