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
8 changes: 3 additions & 5 deletions packages/appkit-ui/src/js/sse/connect-sse.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,9 @@ export async function connectSSE<Payload = unknown>(
lastEventId: initialLastEventId = null,
retryDelay = 2000,
maxRetries = 3,
// 1 MiB — matches the server's `streamDefaults.maxEventSize`. SSE
// carries only short JSON control messages; bulk Arrow payloads stream
// back as the raw HTTP response body of the query request, so this
// buffer never needs to hold multi-MiB attachments.
maxBufferSize = 1 * 1024 * 1024,
// Matches the server's `streamDefaults.maxEventSize`: a smaller buffer here
// turns a legal server event into "Buffer size exceeded" plus retries.
maxBufferSize = 5 * 1024 * 1024,
timeout = 300000, // 5 minutes
onError,
} = options;
Expand Down
9 changes: 4 additions & 5 deletions packages/appkit/src/stream/defaults.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,9 @@
export const streamDefaults = {
bufferSize: 100,
// 1 MiB. SSE carries only short JSON control messages — JSON_ARRAY result
// rows (already row-size-bounded by the warehouse) plus warehouse-readiness
// and error events. ARROW_STREAM never uses SSE: the raw Arrow IPC bytes
// stream back on the query response body (`_handleArrowStreamQuery`).
maxEventSize: 1 * 1024 * 1024,
// Bounds the JSON_ARRAY path only: ARROW_STREAM bytes stream back on the
// query response body (`_handleArrowStreamQuery`), never over SSE.
// Keep in sync with `connectSSE`'s `maxBufferSize` default.
maxEventSize: 5 * 1024 * 1024,
bufferTTL: 10 * 60 * 1000, // 10 minutes
heartbeatInterval: 10 * 1000, // 10 seconds
maxActiveStreams: 1000, // 1000 streams
Expand Down
10 changes: 7 additions & 3 deletions packages/appkit/src/stream/stream-manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,7 @@ export class StreamManager {
const eventBuffer = new EventRingBuffer(
options?.bufferSize ?? streamDefaults.bufferSize,
);
const maxEventSize = options?.maxEventSize ?? this.maxEventSize;

// setup signals and heartbeat
const combinedSignal = this._combineSignals(
Expand Down Expand Up @@ -225,6 +226,7 @@ export class StreamManager {
lastAccess: Date.now(),
abortController,
traceContext,
maxEventSize,
};
this.streamRegistry.add(streamEntry);

Expand Down Expand Up @@ -270,10 +272,12 @@ export class StreamManager {
if (streamEntry.abortController.signal.aborted) break;
const eventId = randomUUID();
const eventData = JSON.stringify(event);
const { maxEventSize } = streamEntry;

// validate event size
if (eventData.length > this.maxEventSize) {
const errorMsg = `Event exceeds max size of ${this.maxEventSize} bytes`;
// UTF-8 bytes, not `String.length`: non-ASCII payloads used to slip
// past a limit they exceeded on the wire.
if (Buffer.byteLength(eventData, "utf8") > maxEventSize) {
const errorMsg = `Event exceeds max size of ${maxEventSize} bytes`;
const errorCode = SSEErrorCode.INVALID_REQUEST;
// broadcast error to all connected clients
this._broadcastErrorToClients(
Expand Down
2 changes: 2 additions & 0 deletions packages/appkit/src/stream/tests/stream-registry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import type { Context } from "@opentelemetry/api";
import { beforeEach, describe, expect, test, vi } from "vitest";

import { EventRingBuffer } from "../buffers";
import { streamDefaults } from "../defaults";
import { StreamRegistry } from "../stream-registry";
import type { StreamEntry } from "../types";
import { SSEErrorCode } from "../types";
Expand All @@ -20,6 +21,7 @@ function createMockStreamEntry(
lastAccess: Date.now(),
abortController: new AbortController(),
traceContext: {} as Context,
maxEventSize: streamDefaults.maxEventSize,
...overrides,
};
}
Expand Down
90 changes: 90 additions & 0 deletions packages/appkit/src/stream/tests/stream.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { beforeEach, describe, expect, test, vi } from "vitest";

import { streamDefaults } from "../defaults";
import { StreamManager } from "../index";

function createMockResponse(headers: Record<string, string> = {}) {
Expand Down Expand Up @@ -1319,3 +1320,92 @@ describe("StreamManager", () => {
});
});
});

/**
* `maxEventSize` was read off the StreamManager instance (constructor-only), so
* per-call values passed via `executeStream({ stream })` type-checked and were
* then silently ignored. The constructor path had coverage; the per-call path
* did not, which is how it survived.
*/
describe("maxEventSize", () => {
/** Tracks the default, so raising it can't quietly neuter these tests. */
function oversizedEvent() {
return {
type: "result",
data: "x".repeat(streamDefaults.maxEventSize + 64),
};
}

function errorEvents(events: string[]) {
// SSEWriter emits one res.write per line, so rejoin before parsing
return events
.join("")
.split("\n\n")
.filter((frame) => frame.includes("event: error"))
.map((frame) => {
const line = frame.split("\n").find((l) => l.startsWith("data: "));
return line ? JSON.parse(line.slice("data: ".length)) : undefined;
});
}

test("drops an oversized event under the default limit", async () => {
const manager = new StreamManager();
const { mockRes, events } = createMockResponse();

async function* generator() {
yield oversizedEvent();
}

await manager.stream(mockRes as any, generator);

const errors = errorEvents(events);
expect(errors).toHaveLength(1);
expect(errors[0].code).toBe("INVALID_REQUEST");
expect(events.some((e) => e.includes("event: result"))).toBe(false);
});

test("honours a per-call override", async () => {
const manager = new StreamManager();
const { mockRes, events } = createMockResponse();

async function* generator() {
yield oversizedEvent();
}

await manager.stream(mockRes as any, generator, {
maxEventSize: streamDefaults.maxEventSize * 2,
});

expect(errorEvents(events)).toHaveLength(0);
expect(events.some((e) => e.includes("event: result"))).toBe(true);
});

test("a per-call override beats the constructor value", async () => {
const manager = new StreamManager({ maxEventSize: 10_000 });
const { mockRes, events } = createMockResponse();

async function* generator() {
yield { type: "result", data: "x".repeat(1000) };
}

await manager.stream(mockRes as any, generator, { maxEventSize: 50 });

expect(errorEvents(events)).toHaveLength(1);
expect(errorEvents(events)[0].error).toContain("50 bytes");
});

test("measures UTF-8 bytes, not UTF-16 code units", async () => {
// ~730 UTF-16 units (under the limit, so String.length accepted it) but
// ~1430 UTF-8 bytes, which is what crosses the wire.
const manager = new StreamManager({ maxEventSize: 1000 });
const { mockRes, events } = createMockResponse();

async function* generator() {
yield { type: "result", data: "\u00e9".repeat(700) };
}

await manager.stream(mockRes as any, generator);

expect(errorEvents(events)).toHaveLength(1);
});
});
6 changes: 6 additions & 0 deletions packages/appkit/src/stream/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,12 @@ export interface StreamEntry {
lastAccess: number;
abortController: AbortController;
traceContext: Context;
/**
* UTF-8 bytes. Lives on the entry, not on `StreamManager`: one manager is
* shared by every stream a plugin opens, so reading it off the instance made
* per-call overrides silently ineffective.
*/
maxEventSize: number;
// pending grace-window abort, set while the last client is disconnected
disconnectGraceTimer?: NodeJS.Timeout;
/**
Expand Down
Loading