# How Event Streaming Works with Capped-Stream for Render Progress in OpenMAIC

> Discover how event streaming with capped-stream in OpenMAIC tracks render progress in real-time. Safely limit uploads and prevent memory overflow with this efficient technique.

- Repository: [MAIC/OpenMAIC](https://github.com/THU-MAIC/OpenMAIC)
- Tags: internals
- Published: 2026-09-13

---

**OpenMAIC uses a byte-capped stream wrapper to safely limit upload sizes while simultaneously emitting render lifecycle events as NDJSON lines, enabling real-time progress tracking without risking memory overflow.**

The OpenMAIC *render service* processes video export jobs that require both strict upload size limits and live progress feedback. According to the `THU-MAIC/OpenMAIC` source code, the system combines a protective byte-cap mechanism with an event streaming architecture to handle large file submissions while keeping clients updated on render status through newline-delimited JSON (NDJSON).

## The Two-Pipe Architecture: Capped-Stream and Event Streaming

The solution relies on two orthogonal mechanisms that operate concurrently during a render request.

### Byte-Limited Upload Protection

The **`capBodyStream`** utility in [`render-service/src/capped-stream.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/render-service/src/capped-stream.ts) creates a protective wrapper around the incoming request body. It instantiates a new `ReadableStream<Uint8Array>` that counts accumulated bytes (`total`) as chunks arrive. If the running total exceeds the configured `capBytes`, the wrapper immediately sets a `tripped = true` flag and aborts the underlying stream.

The function returns an object containing the capped `stream` and an `exceeded()` accessor, allowing the caller to check whether the limit was breached without buffering the entire payload.

### NDJSON Progress Events

Render lifecycle events are defined in [`render-service/src/events.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/render-service/src/events.ts). The **`emitRenderEvent`** function serializes `RenderEvent` objects—such as `render_job_started`, `render_job_progress`, and `render_job_finished`—and writes them to a configurable `RenderEventSink`. The default sink outputs one JSON line per event to **stdout**, which the HTTP handler forwards to the client as a streaming response.

## Request Flow from Upload to Progress Stream

The integration logic in [`render-service/src/main.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/render-service/src/main.ts) orchestrates these components through a strict sequence:

1. **Request Arrival** – The HTTP endpoint receives a `POST` to `/api/export-video/render` with a `ReadableStream<Uint8Array>` body.

2. **Capped-Stream Creation** – The handler invokes `capBodyStream(body, MAX_UPLOAD_BYTES, signal)` to wrap the raw body. Lines 19–83 in [`capped-stream.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/capped-stream.ts) implement the counting logic that aborts the stream if the byte cap is exceeded.

3. **Upload Validation** – Before processing, [`main.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/main.ts) checks `capped?.exceeded()`. If the probe returns `true`, the handler throws `UploadTooLargeError`, terminating the request before any buffering occurs.

4. **Render Execution** – The validated `capped.stream` is passed to the render executor. As the job processes, the executor calls `emitRenderEvent()` periodically to publish progress updates.

5. **Event Streaming to Client** – The HTTP handler wires the response stream to the same `RenderEventSink`. Each event is written as an NDJSON line, allowing the frontend to parse progress updates in real time as they arrive.

## Implementation Deep Dive

### Creating the Capped Stream

The following pattern from [`render-service/src/capped-stream.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/render-service/src/capped-stream.ts) demonstrates how to enforce upload limits before consuming the body:

```typescript
import { capBodyStream } from '@/lib/server/capped-stream';
import { UploadTooLargeError } from '@/lib/server/errors';

const capped = capBodyStream(req.body, MAX_UPLOAD_BYTES, req.signal);
if (capped?.exceeded()) {
  throw new UploadTooLargeError('Upload too large');
}
const body = capped.stream; // Pass this to the renderer

```

The `exceeded()` method provides an early-exit path that prevents oversized payloads from reaching the renderer or causing memory issues.

### Emitting Render Events

During execution, the render job reports progress via the event API defined in [`render-service/src/events.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/render-service/src/events.ts):

```typescript
import { emitRenderEvent, type RenderEventSink } from '@/render-service/src/events';

function reportProgress(sink: RenderEventSink, jobId: string, percent: number) {
  emitRenderEvent({
    event: 'render_job_progress',
    jobId,
    progress: percent,
    timestamp: Date.now(),
  });
}

```

Each call serializes the event and writes it to the configured sink, which defaults to stdout in the reference implementation.

### Wiring the HTTP Handler

The request handler in [`render-service/src/main.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/render-service/src/main.ts) bridges the capped input stream with the event output stream:

```typescript
// Inside the HTTP handler
const sink: RenderEventSink = (ev) => {
  // NDJSON line sent to the response stream
  res.write(JSON.stringify(ev) + '\n');
};

await renderExecutor.run(job, {
  eventSink: sink,
  body: capped?.stream,
});
res.end();

```

This configuration ensures that progress events flow to the client as NDJSON lines while the upload remains constrained by the byte cap.

## Design Benefits

**Back-pressure safety** – The capped stream respects the underlying socket back-pressure, reading only as fast as the downstream consumer can handle. This prevents memory blow-ups when clients send data faster than the server can process.

**Early abort** – Because `capBodyStream` aborts the moment the byte limit is exceeded, the service avoids buffering maliciously large uploads. The `exceeded()` check occurs before any `formData()` or `arrayBuffer()` calls, minimizing the attack surface.

**Stateless progress** – Render progress events are decoupled from the request body consumption. Once the capped stream is validated and passed to the renderer, it can be discarded, while progress events continue to stream independently via the `RenderEventSink`.

## Summary

- **`capBodyStream`** in [`render-service/src/capped-stream.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/render-service/src/capped-stream.ts) wraps request bodies to enforce strict byte limits and provides an `exceeded()` probe for early validation.
- **`emitRenderEvent`** in [`render-service/src/events.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/render-service/src/events.ts) publishes lifecycle events to a `RenderEventSink`, defaulting to NDJSON output on stdout.
- **[`main.ts`](https://github.com/THU-MAIC/OpenMAIC/blob/main/main.ts)** integrates both mechanisms, validating upload size before streaming progress events back to the client.
- The architecture uses raw streams rather than buffering, protecting against memory exhaustion while enabling real-time progress bars via NDJSON.

## Frequently Asked Questions

### What happens if the upload exceeds the byte cap?

The `capBodyStream` wrapper immediately aborts the underlying `ReadableStream` and sets an internal `tripped` flag. When the handler checks `capped.exceeded()`, it throws `UploadTooLargeError`, which terminates the HTTP request with an appropriate error response before any render processing begins.

### Why does OpenMAIC use NDJSON instead of WebSockets for progress updates?

NDJSON allows the service to stream events over a standard HTTP response without maintaining a persistent WebSocket connection. Each JSON line represents a discrete event that the client can parse incrementally, simplifying infrastructure requirements while still providing real-time updates suitable for progress bars.

### How does the capped-stream prevent memory attacks?

By counting bytes as chunks flow through the stream and aborting before buffering, the mechanism ensures that oversized payloads never occupy heap memory. The `capBodyStream` implementation reads the raw `Uint8Array` chunks directly without converting them to strings or buffers until after the size check passes, eliminating the risk of exhausting server memory with large uploads.

### Can the render event sink be customized?

Yes. While the default `emitRenderEvent` implementation writes to stdout, the `RenderEventSink` type is a simple function interface: `(event: RenderEvent) => void`. You can inject a custom sink—such as one that writes to a logging service, database, or alternative stream—when calling `renderExecutor.run()` to redirect events to any destination.