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

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 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. 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 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 implement the counting logic that aborts the stream if the byte cap is exceeded.

  3. Upload Validation – Before processing, 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 demonstrates how to enforce upload limits before consuming the body:

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:

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 bridges the capped input stream with the event output stream:

// 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 wraps request bodies to enforce strict byte limits and provides an exceeded() probe for early validation.
  • emitRenderEvent in render-service/src/events.ts publishes lifecycle events to a RenderEventSink, defaulting to NDJSON output on stdout.
  • 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.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →