mirror of
https://github.com/vercel/eve.git
synced 2026-09-20 05:35:39 +08:00
fix(eve): lease session stream responses (#3407)
Signed-off-by: owenkephart <owen.kephart@vercel.com>
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"eve": patch
|
||||
---
|
||||
|
||||
Bound negotiated session stream responses with renewable leases so abandoned serverless invocations release their durable stream readers. The eve client renews these responses from its cursor without exposing transport heartbeats or lease records as session events; tail-relative reads and streams with reconnection disabled remain unleased.
|
||||
@@ -251,6 +251,8 @@ Reset terminally retires the exact session ID. A reset ID never becomes a new se
|
||||
|
||||
The stream is durable. Every event is recorded before a step completes, so consumers can reconnect from their cursor when an HTTP connection ends. A nonnegative `startIndex` is an absolute event count: use it to pick up where you dropped off or pass `0` to rewind to the start.
|
||||
|
||||
When automatic reconnection is enabled and `startIndex` is nonnegative, the TypeScript client requests renewable leases over that durable stream. A server that supports leases sends transport heartbeats during quiet periods, then ends the response so the client reconnects from its current cursor. This bounds server-side stream readers even when a host does not report that the client disconnected. The lease renewal and heartbeats are transport details: they do not stop the run or appear as session events. Tail-relative reads and streams with reconnection disabled do not request leases.
|
||||
|
||||
If a reconnect overlaps events you already handled, [`meta.id`](#the-event-envelope) identifies the duplicates: it is unchanged across reconnects and rewinds, so a consumer keyed on it can replay safely.
|
||||
|
||||
```bash
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { type MessageStreamEvent } from "#protocol/message.js";
|
||||
import { EVE_STREAM_LEASE_ENDED_CONTROL, type MessageStreamEvent } from "#protocol/message.js";
|
||||
import {
|
||||
normalizeMessageStreamEvent,
|
||||
type MessageStreamEventForVersion,
|
||||
@@ -41,7 +41,9 @@ export function isStreamDisconnectError(error: unknown): boolean {
|
||||
export async function* readNdjsonStream(
|
||||
body: ReadableStream<Uint8Array>,
|
||||
options: {
|
||||
readonly controlVersion?: "1";
|
||||
readonly idleTimeoutMs?: number;
|
||||
readonly onLeaseEnded?: () => void;
|
||||
readonly streamVersion: MessageStreamVersion;
|
||||
},
|
||||
): AsyncGenerator<MessageStreamEvent> {
|
||||
@@ -72,7 +74,12 @@ export async function* readNdjsonStream(
|
||||
buffer = buffer.slice(newlineIndex + 1);
|
||||
|
||||
if (line.length > 0) {
|
||||
yield parseMessageStreamEvent(line, options.streamVersion);
|
||||
const value = JSON.parse(line) as unknown;
|
||||
if (options.controlVersion === "1" && isLeaseEndedControl(value)) {
|
||||
options.onLeaseEnded?.();
|
||||
} else {
|
||||
yield parseMessageStreamEvent(value, options.streamVersion);
|
||||
}
|
||||
}
|
||||
|
||||
newlineIndex = buffer.indexOf("\n");
|
||||
@@ -82,7 +89,12 @@ export async function* readNdjsonStream(
|
||||
// Yield any trailing content without a final newline.
|
||||
const trailing = buffer.trim();
|
||||
if (trailing.length > 0) {
|
||||
yield parseMessageStreamEvent(trailing, options.streamVersion);
|
||||
const value = JSON.parse(trailing) as unknown;
|
||||
if (options.controlVersion === "1" && isLeaseEndedControl(value)) {
|
||||
options.onLeaseEnded?.();
|
||||
} else {
|
||||
yield parseMessageStreamEvent(value, options.streamVersion);
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
if (!reachedEof) {
|
||||
@@ -94,12 +106,20 @@ export async function* readNdjsonStream(
|
||||
}
|
||||
}
|
||||
|
||||
function isLeaseEndedControl(value: unknown): boolean {
|
||||
if (value === null || typeof value !== "object") return false;
|
||||
const record = value as Record<string, unknown>;
|
||||
return (
|
||||
record.$eve === EVE_STREAM_LEASE_ENDED_CONTROL.$eve &&
|
||||
record.version === EVE_STREAM_LEASE_ENDED_CONTROL.version
|
||||
);
|
||||
}
|
||||
|
||||
function parseMessageStreamEvent<Version extends MessageStreamVersion>(
|
||||
line: string,
|
||||
value: unknown,
|
||||
version: Version,
|
||||
): MessageStreamEvent {
|
||||
const event = JSON.parse(line) as MessageStreamEventForVersion<Version>;
|
||||
return normalizeMessageStreamEvent(version, event);
|
||||
return normalizeMessageStreamEvent(version, value as MessageStreamEventForVersion<Version>);
|
||||
}
|
||||
|
||||
async function readWithIdleTimeout(
|
||||
|
||||
@@ -1,10 +1,19 @@
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
import { openStreamBody } from "./open-stream.js";
|
||||
import { EVE_MESSAGE_STREAM_VERSION, EVE_STREAM_VERSION_HEADER } from "#protocol/message.js";
|
||||
import { followStreamIterable, openStreamBody } from "./open-stream.js";
|
||||
import type { StreamReconnectPolicy } from "#client/types.js";
|
||||
import {
|
||||
EVE_MESSAGE_STREAM_VERSION,
|
||||
EVE_STREAM_CONTROL_VERSION,
|
||||
EVE_STREAM_CONTROL_VERSION_QUERY,
|
||||
EVE_STREAM_LEASE_ENDED_CONTROL,
|
||||
EVE_STREAM_TAIL_INDEX_HEADER,
|
||||
EVE_STREAM_VERSION_HEADER,
|
||||
} from "#protocol/message.js";
|
||||
|
||||
afterEach(() => {
|
||||
vi.unstubAllGlobals();
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
describe("openStreamBody", () => {
|
||||
@@ -36,4 +45,158 @@ describe("openStreamBody", () => {
|
||||
expect(cancel).toHaveBeenCalledOnce();
|
||||
expect(signal?.aborted).toBe(true);
|
||||
});
|
||||
|
||||
it("advertises support for versioned stream controls", async () => {
|
||||
let requestUrl: URL | undefined;
|
||||
vi.stubGlobal(
|
||||
"fetch",
|
||||
vi.fn(async (input: Parameters<typeof fetch>[0]) => {
|
||||
requestUrl = new URL(String(input));
|
||||
return new Response("\n", {
|
||||
headers: {
|
||||
[EVE_STREAM_VERSION_HEADER]: EVE_MESSAGE_STREAM_VERSION,
|
||||
},
|
||||
status: 200,
|
||||
});
|
||||
}),
|
||||
);
|
||||
|
||||
const connection = await openStreamBody({
|
||||
host: "https://agent.example",
|
||||
resolveHeaders: () => Promise.resolve(new Headers()),
|
||||
sessionId: "session_1",
|
||||
startIndex: 0,
|
||||
});
|
||||
connection.close();
|
||||
|
||||
expect(requestUrl?.searchParams.get(EVE_STREAM_CONTROL_VERSION_QUERY)).toBe(
|
||||
EVE_STREAM_CONTROL_VERSION,
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
describe("followStreamIterable", () => {
|
||||
it.each([
|
||||
{ startIndex: -1, streamReconnectPolicy: undefined },
|
||||
{ startIndex: 0, streamReconnectPolicy: { reconnect: false } },
|
||||
{
|
||||
startIndex: 0,
|
||||
streamReconnectPolicy: { streamIdleReconnectPolicy: { maxAttempts: 0 } },
|
||||
},
|
||||
] satisfies Array<{
|
||||
startIndex: number;
|
||||
streamReconnectPolicy: StreamReconnectPolicy | undefined;
|
||||
}>)("does not request a lease when it cannot renew: %j", async (options) => {
|
||||
const requestUrls: URL[] = [];
|
||||
vi.stubGlobal(
|
||||
"fetch",
|
||||
vi.fn(async (input: Parameters<typeof fetch>[0]) => {
|
||||
requestUrls.push(new URL(String(input)));
|
||||
return new Response("\n", {
|
||||
headers: { [EVE_STREAM_VERSION_HEADER]: EVE_MESSAGE_STREAM_VERSION },
|
||||
});
|
||||
}),
|
||||
);
|
||||
|
||||
for await (const event of followStreamIterable({
|
||||
host: "https://agent.example",
|
||||
resolveHeaders: () => Promise.resolve(new Headers()),
|
||||
sessionId: "session_1",
|
||||
...options,
|
||||
})) {
|
||||
expect.unreachable(`Unexpected event: ${event.type}`);
|
||||
}
|
||||
|
||||
expect(requestUrls).toHaveLength(1);
|
||||
expect(requestUrls[0]?.searchParams.has(EVE_STREAM_CONTROL_VERSION_QUERY)).toBe(false);
|
||||
});
|
||||
|
||||
it.each([true, false])(
|
||||
"resumes from the event cursor across leases without gaps or duplicates (follow: %s)",
|
||||
async (follow) => {
|
||||
const events = Array.from({ length: 7 }, (_, index) => ({
|
||||
type: "message.appended",
|
||||
data: {
|
||||
messageDelta: String(index),
|
||||
sequence: 0,
|
||||
stepIndex: 0,
|
||||
turnId: "turn_1",
|
||||
},
|
||||
meta: { id: `evt_${index}`, at: "2026-09-15T00:00:00.000Z" },
|
||||
}));
|
||||
const cursors: number[] = [];
|
||||
const tailRequests: Array<string | null> = [];
|
||||
vi.stubGlobal(
|
||||
"fetch",
|
||||
vi.fn(async (input: Parameters<typeof fetch>[0]) => {
|
||||
const url = new URL(String(input));
|
||||
expect(url.searchParams.get(EVE_STREAM_CONTROL_VERSION_QUERY)).toBe(
|
||||
EVE_STREAM_CONTROL_VERSION,
|
||||
);
|
||||
const cursor = Number(url.searchParams.get("startIndex") ?? "0");
|
||||
cursors.push(cursor);
|
||||
tailRequests.push(url.searchParams.get("includeTailIndex"));
|
||||
expect(cursors.length).toBeLessThanOrEqual(3);
|
||||
const records = events.slice(cursor, cursor + 2).map((event) => JSON.stringify(event));
|
||||
return new Response(
|
||||
`\n${records.join("\n\n")}\n\n${JSON.stringify(EVE_STREAM_LEASE_ENDED_CONTROL)}\n`,
|
||||
{
|
||||
headers: {
|
||||
[EVE_STREAM_VERSION_HEADER]: EVE_MESSAGE_STREAM_VERSION,
|
||||
...(url.searchParams.has("includeTailIndex")
|
||||
? { [EVE_STREAM_TAIL_INDEX_HEADER]: String(events.length - 1) }
|
||||
: {}),
|
||||
},
|
||||
},
|
||||
);
|
||||
}),
|
||||
);
|
||||
|
||||
const received = [];
|
||||
for await (const event of followStreamIterable({
|
||||
host: "https://agent.example",
|
||||
resolveHeaders: () => Promise.resolve(new Headers()),
|
||||
sessionId: "session_1",
|
||||
startIndex: 1,
|
||||
follow,
|
||||
})) {
|
||||
received.push(event);
|
||||
if (follow && received.length === events.length - 1) break;
|
||||
}
|
||||
|
||||
expect(received).toEqual(events.slice(1));
|
||||
expect(cursors).toEqual([1, 3, 5]);
|
||||
expect(tailRequests).toEqual(follow ? [null, null, null] : ["1", null, null]);
|
||||
},
|
||||
);
|
||||
|
||||
it("renews an explicitly ended lease without charging the idle budget", async () => {
|
||||
let connections = 0;
|
||||
const abort = new AbortController();
|
||||
vi.stubGlobal(
|
||||
"fetch",
|
||||
vi.fn(async () => {
|
||||
connections += 1;
|
||||
if (connections === 8) abort.abort();
|
||||
return new Response(`${JSON.stringify(EVE_STREAM_LEASE_ENDED_CONTROL)}\n`, {
|
||||
headers: {
|
||||
[EVE_STREAM_VERSION_HEADER]: EVE_MESSAGE_STREAM_VERSION,
|
||||
},
|
||||
status: 200,
|
||||
});
|
||||
}),
|
||||
);
|
||||
|
||||
for await (const _event of followStreamIterable({
|
||||
host: "https://agent.example",
|
||||
resolveHeaders: () => Promise.resolve(new Headers()),
|
||||
sessionId: "session_1",
|
||||
signal: abort.signal,
|
||||
startIndex: 0,
|
||||
})) {
|
||||
// No durable events are expected.
|
||||
}
|
||||
|
||||
expect(connections).toBe(8);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
import type { MessageStreamEvent } from "#protocol/message.js";
|
||||
import { EVE_STREAM_TAIL_INDEX_HEADER } from "#protocol/message.js";
|
||||
import {
|
||||
EVE_STREAM_CONTROL_VERSION,
|
||||
EVE_STREAM_CONTROL_VERSION_QUERY,
|
||||
EVE_STREAM_TAIL_INDEX_HEADER,
|
||||
} from "#protocol/message.js";
|
||||
import type { MessageStreamVersion } from "#protocol/message-version.js";
|
||||
import { createEveSessionStreamRoutePath } from "#protocol/routes.js";
|
||||
import { ClientError } from "#client/client-error.js";
|
||||
@@ -162,9 +166,14 @@ export async function* followStreamIterable(
|
||||
}
|
||||
|
||||
let deliveredEvent = false;
|
||||
let leaseEnded = false;
|
||||
try {
|
||||
for await (const event of readNdjsonStream(connection.body, {
|
||||
controlVersion: connection.controlVersion,
|
||||
idleTimeoutMs: input.streamReadIdleTimeoutMs ?? DEFAULT_STREAM_READ_IDLE_TIMEOUT_MS,
|
||||
onLeaseEnded: () => {
|
||||
leaseEnded = true;
|
||||
},
|
||||
streamVersion: connection.streamVersion,
|
||||
})) {
|
||||
startIndex += 1;
|
||||
@@ -187,6 +196,10 @@ export async function* followStreamIterable(
|
||||
return;
|
||||
}
|
||||
|
||||
if (leaseEnded) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (
|
||||
input.keepAlive !== true &&
|
||||
!deliveredEvent &&
|
||||
@@ -209,6 +222,7 @@ export async function* followStreamIterable(
|
||||
interface OpenedStream {
|
||||
readonly body: ReadableStream<Uint8Array>;
|
||||
close(): void;
|
||||
readonly controlVersion: "1" | undefined;
|
||||
readonly streamVersion: MessageStreamVersion;
|
||||
readonly tailIndex: number | undefined;
|
||||
}
|
||||
@@ -222,14 +236,22 @@ interface OpenedStream {
|
||||
export async function openStreamBody(
|
||||
input: OpenStreamInput & { readonly retryPolicy?: ResolvedStreamReconnectPolicy },
|
||||
): Promise<OpenedStream> {
|
||||
const retryPolicy = input.retryPolicy ?? DEFAULT_STREAM_RECONNECT_POLICY;
|
||||
const retryPolicy =
|
||||
input.retryPolicy ?? resolveStreamReconnectPolicy(input.streamReconnectPolicy);
|
||||
const openRetryPolicy = retryPolicy.streamOpenReconnectPolicy;
|
||||
let lastStatus: number | undefined;
|
||||
let lastBody: string | undefined;
|
||||
let lastHeaders: Headers | undefined;
|
||||
let retryDelayMs = openRetryPolicy.baseDelayMs;
|
||||
|
||||
const controlVersion =
|
||||
input.startIndex >= 0 && retryPolicy.streamIdleReconnectPolicy.maxAttempts > 0
|
||||
? EVE_STREAM_CONTROL_VERSION
|
||||
: undefined;
|
||||
const searchParams: Record<string, string> = {};
|
||||
if (controlVersion !== undefined) {
|
||||
searchParams[EVE_STREAM_CONTROL_VERSION_QUERY] = controlVersion;
|
||||
}
|
||||
if (input.startIndex !== 0) {
|
||||
searchParams.startIndex = String(input.startIndex);
|
||||
}
|
||||
@@ -287,6 +309,7 @@ export async function openStreamBody(
|
||||
response.body?.cancel().catch(() => {});
|
||||
connectionController.abort();
|
||||
},
|
||||
controlVersion,
|
||||
streamVersion: readMessageStreamVersion(response.headers),
|
||||
tailIndex: parseTailIndexHeader(response.headers),
|
||||
};
|
||||
|
||||
@@ -19,6 +19,7 @@ import {
|
||||
} from "#internal/nitro/routes/channel-route-context.js";
|
||||
import {
|
||||
EVE_SESSION_ID_HEADER,
|
||||
EVE_STREAM_CONTROL_VERSION_QUERY,
|
||||
EVE_STREAM_FORMAT_HEADER,
|
||||
EVE_STREAM_TAIL_INDEX_HEADER,
|
||||
EVE_STREAM_VERSION_HEADER,
|
||||
@@ -643,6 +644,10 @@ export function eveChannel(input: EveChannelInput): EveChannel {
|
||||
if (startIndex !== undefined) {
|
||||
upstreamUrl.searchParams.set("startIndex", String(startIndex));
|
||||
}
|
||||
const controlVersion = new URL(req.url).searchParams.get(EVE_STREAM_CONTROL_VERSION_QUERY);
|
||||
if (controlVersion !== null) {
|
||||
upstreamUrl.searchParams.set(EVE_STREAM_CONTROL_VERSION_QUERY, controlVersion);
|
||||
}
|
||||
if (includeTailIndex) {
|
||||
upstreamUrl.searchParams.set("includeTailIndex", "1");
|
||||
}
|
||||
|
||||
@@ -0,0 +1,123 @@
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
import type { Session } from "#channel/session.js";
|
||||
import { createSessionStreamResponse } from "#eve-channel/request.js";
|
||||
import {
|
||||
EVE_STREAM_CONTROL_VERSION,
|
||||
EVE_STREAM_CONTROL_VERSION_QUERY,
|
||||
EVE_STREAM_LEASE_ENDED_CONTROL,
|
||||
} from "#protocol/message.js";
|
||||
|
||||
const encoder = new TextEncoder();
|
||||
|
||||
function stubSession(events: ReadableStream<unknown>): Session {
|
||||
return {
|
||||
id: "session_1",
|
||||
async getEventStream() {
|
||||
return events;
|
||||
},
|
||||
} as Session;
|
||||
}
|
||||
|
||||
function leasedRequest(): Request {
|
||||
const url = new URL("https://eve.test/eve/v1/session/session_1/stream");
|
||||
url.searchParams.set(EVE_STREAM_CONTROL_VERSION_QUERY, EVE_STREAM_CONTROL_VERSION);
|
||||
return new Request(url);
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
describe("createSessionStreamResponse", () => {
|
||||
it("ends a negotiated response after its lease and cancels the source", async () => {
|
||||
vi.useFakeTimers();
|
||||
let cancelled = false;
|
||||
const events = new ReadableStream({
|
||||
cancel() {
|
||||
cancelled = true;
|
||||
},
|
||||
});
|
||||
|
||||
const response = await createSessionStreamResponse(leasedRequest(), stubSession(events));
|
||||
const reader = response.body!.getReader();
|
||||
|
||||
await expect(reader.read()).resolves.toEqual({ done: false, value: encoder.encode("\n") });
|
||||
for (let elapsed = 10_000; elapsed < 60_000; elapsed += 10_000) {
|
||||
const heartbeat = reader.read();
|
||||
await vi.advanceTimersByTimeAsync(10_000);
|
||||
await expect(heartbeat).resolves.toEqual({ done: false, value: encoder.encode("\n") });
|
||||
}
|
||||
|
||||
const control = reader.read();
|
||||
await vi.advanceTimersByTimeAsync(10_000);
|
||||
await expect(control).resolves.toEqual({
|
||||
done: false,
|
||||
value: encoder.encode(`${JSON.stringify(EVE_STREAM_LEASE_ENDED_CONTROL)}\n`),
|
||||
});
|
||||
await expect(reader.read()).resolves.toEqual({ done: true, value: undefined });
|
||||
await vi.waitFor(() => expect(cancelled).toBe(true));
|
||||
});
|
||||
|
||||
it("resets heartbeats after events without extending the lease", async () => {
|
||||
vi.useFakeTimers();
|
||||
let eventController: ReadableStreamDefaultController<unknown> | undefined;
|
||||
const events = new ReadableStream({
|
||||
start(controller) {
|
||||
eventController = controller;
|
||||
},
|
||||
});
|
||||
|
||||
const response = await createSessionStreamResponse(leasedRequest(), stubSession(events));
|
||||
const reader = response.body!.getReader();
|
||||
await reader.read();
|
||||
|
||||
await vi.advanceTimersByTimeAsync(9_000);
|
||||
eventController!.enqueue({ type: "test" });
|
||||
await expect(reader.read()).resolves.toEqual({
|
||||
done: false,
|
||||
value: encoder.encode(`${JSON.stringify({ type: "test" })}\n`),
|
||||
});
|
||||
|
||||
const heartbeat = reader.read();
|
||||
await vi.advanceTimersByTimeAsync(9_999);
|
||||
let heartbeatSettled = false;
|
||||
void heartbeat.then(() => {
|
||||
heartbeatSettled = true;
|
||||
});
|
||||
await vi.runAllTicks();
|
||||
expect(heartbeatSettled).toBe(false);
|
||||
await vi.advanceTimersByTimeAsync(1);
|
||||
await expect(heartbeat).resolves.toEqual({ done: false, value: encoder.encode("\n") });
|
||||
|
||||
for (let elapsed = 20_000; elapsed < 60_000; elapsed += 10_000) {
|
||||
const nextHeartbeat = reader.read();
|
||||
await vi.advanceTimersByTimeAsync(10_000);
|
||||
await expect(nextHeartbeat).resolves.toEqual({ done: false, value: encoder.encode("\n") });
|
||||
}
|
||||
const control = reader.read();
|
||||
await vi.advanceTimersByTimeAsync(1_000);
|
||||
await expect(control).resolves.toMatchObject({ done: false });
|
||||
await expect(reader.read()).resolves.toEqual({ done: true, value: undefined });
|
||||
});
|
||||
|
||||
it("does not lease a response when the client does not negotiate controls", async () => {
|
||||
vi.useFakeTimers();
|
||||
const response = await createSessionStreamResponse(
|
||||
new Request("https://eve.test/eve/v1/session/session_1/stream"),
|
||||
stubSession(new ReadableStream()),
|
||||
);
|
||||
const reader = response.body!.getReader();
|
||||
await reader.read();
|
||||
|
||||
let settled = false;
|
||||
const pending = reader.read().then((result) => {
|
||||
settled = true;
|
||||
return result;
|
||||
});
|
||||
await vi.advanceTimersByTimeAsync(120_000);
|
||||
expect(settled).toBe(false);
|
||||
await reader.cancel();
|
||||
await pending;
|
||||
});
|
||||
});
|
||||
@@ -19,7 +19,10 @@ import {
|
||||
EVE_MESSAGE_STREAM_FORMAT,
|
||||
EVE_MESSAGE_STREAM_VERSION,
|
||||
EVE_SESSION_ID_HEADER,
|
||||
EVE_STREAM_CONTROL_VERSION,
|
||||
EVE_STREAM_CONTROL_VERSION_QUERY,
|
||||
EVE_STREAM_FORMAT_HEADER,
|
||||
EVE_STREAM_LEASE_ENDED_CONTROL,
|
||||
EVE_STREAM_TAIL_INDEX_HEADER,
|
||||
EVE_STREAM_VERSION_HEADER,
|
||||
} from "#protocol/message.js";
|
||||
@@ -32,6 +35,9 @@ import { isInputResponse, type ValidatedInputResponse } from "#shared/input.js";
|
||||
import { parseJsonObject, type JsonObject } from "#shared/json.js";
|
||||
import type { RunMode } from "#shared/run-mode.js";
|
||||
|
||||
const SESSION_STREAM_HEARTBEAT_MS = 10_000;
|
||||
const SESSION_STREAM_LEASE_MS = 60_000;
|
||||
|
||||
interface ParsedCreateBody {
|
||||
activityObserver?: ActivityObserverConfig;
|
||||
callback?: SessionCallback;
|
||||
@@ -310,6 +316,11 @@ export async function createSessionStreamResponse(
|
||||
try {
|
||||
const tailIndex = includeTailIndex ? await session.getStreamTailIndex() : undefined;
|
||||
const events = await session.getEventStream({ startIndex });
|
||||
const controlVersion =
|
||||
new URL(request.url).searchParams.get(EVE_STREAM_CONTROL_VERSION_QUERY) ===
|
||||
EVE_STREAM_CONTROL_VERSION
|
||||
? EVE_STREAM_CONTROL_VERSION
|
||||
: undefined;
|
||||
const headers = new Headers({
|
||||
"cache-control": "no-store, no-transform",
|
||||
"content-type": EVE_MESSAGE_STREAM_CONTENT_TYPE,
|
||||
@@ -322,7 +333,12 @@ export async function createSessionStreamResponse(
|
||||
headers.set(EVE_STREAM_TAIL_INDEX_HEADER, String(tailIndex));
|
||||
}
|
||||
return new Response(
|
||||
serializeAsNdjson(events, request.signal, streamEventLimit(startIndex, tailIndex)),
|
||||
serializeAsNdjson(
|
||||
events,
|
||||
request.signal,
|
||||
streamEventLimit(startIndex, tailIndex),
|
||||
controlVersion !== undefined,
|
||||
),
|
||||
{ headers },
|
||||
);
|
||||
} catch {
|
||||
@@ -615,20 +631,69 @@ function serializeAsNdjson(
|
||||
events: ReadableStream<unknown>,
|
||||
signal: AbortSignal,
|
||||
eventLimit?: number,
|
||||
leased = false,
|
||||
): ReadableStream<Uint8Array> {
|
||||
const encoder = new TextEncoder();
|
||||
let eventCount = 0;
|
||||
let heartbeat: ReturnType<typeof setTimeout> | undefined;
|
||||
let lease: ReturnType<typeof setTimeout> | undefined;
|
||||
|
||||
const clearTimers = () => {
|
||||
clearTimeout(heartbeat);
|
||||
clearTimeout(lease);
|
||||
heartbeat = undefined;
|
||||
lease = undefined;
|
||||
};
|
||||
const scheduleHeartbeat = (controller: TransformStreamDefaultController<Uint8Array>) => {
|
||||
clearTimeout(heartbeat);
|
||||
heartbeat = setTimeout(() => {
|
||||
try {
|
||||
controller.enqueue(encoder.encode("\n"));
|
||||
scheduleHeartbeat(controller);
|
||||
} catch {
|
||||
clearTimers();
|
||||
}
|
||||
}, SESSION_STREAM_HEARTBEAT_MS);
|
||||
};
|
||||
const startLease = (controller: TransformStreamDefaultController<Uint8Array>) => {
|
||||
scheduleHeartbeat(controller);
|
||||
lease = setTimeout(() => {
|
||||
clearTimers();
|
||||
try {
|
||||
controller.enqueue(encoder.encode(`${JSON.stringify(EVE_STREAM_LEASE_ENDED_CONTROL)}\n`));
|
||||
controller.terminate();
|
||||
} catch {
|
||||
// The response was cancelled while the lease callback was already queued.
|
||||
}
|
||||
}, SESSION_STREAM_LEASE_MS);
|
||||
};
|
||||
|
||||
const transform = new TransformStream<unknown, Uint8Array>({
|
||||
start(controller) {
|
||||
controller.enqueue(encoder.encode("\n"));
|
||||
if (eventLimit === 0) controller.terminate();
|
||||
if (eventLimit === 0) {
|
||||
controller.terminate();
|
||||
} else if (leased) {
|
||||
startLease(controller);
|
||||
}
|
||||
},
|
||||
transform(event, controller) {
|
||||
controller.enqueue(encoder.encode(`${JSON.stringify(event)}\n`));
|
||||
eventCount += 1;
|
||||
if (eventCount === eventLimit) controller.terminate();
|
||||
if (eventCount === eventLimit) {
|
||||
clearTimers();
|
||||
controller.terminate();
|
||||
} else if (leased) {
|
||||
scheduleHeartbeat(controller);
|
||||
}
|
||||
},
|
||||
flush() {
|
||||
clearTimers();
|
||||
},
|
||||
});
|
||||
void events.pipeTo(transform.writable, { signal }).catch(() => {});
|
||||
void events
|
||||
.pipeTo(transform.writable, { signal })
|
||||
.catch(() => {})
|
||||
.finally(clearTimers);
|
||||
return transform.readable;
|
||||
}
|
||||
|
||||
@@ -29,6 +29,15 @@ export const EVE_MESSAGE_STREAM_CONTENT_TYPE = "application/x-ndjson; charset=ut
|
||||
export const EVE_MESSAGE_STREAM_FORMAT = "ndjson";
|
||||
export const EVE_MESSAGE_STREAM_VERSION = "25";
|
||||
|
||||
/** Version of transport control records understood by this eve release. */
|
||||
export const EVE_STREAM_CONTROL_VERSION = "1";
|
||||
export const EVE_STREAM_CONTROL_VERSION_QUERY = "streamControlVersion";
|
||||
/** Internal record emitted when a leased HTTP response should be renewed. */
|
||||
export const EVE_STREAM_LEASE_ENDED_CONTROL = {
|
||||
$eve: "stream.lease-ended",
|
||||
version: 1,
|
||||
} as const;
|
||||
|
||||
/**
|
||||
* eve-owned finish reason for one completed assistant step.
|
||||
*
|
||||
|
||||
@@ -437,7 +437,7 @@ describe("useEveAgent", () => {
|
||||
|
||||
expect(fetchMock.mock.calls[0]?.[0]).toBe("/eve/agents/support/eve/v1/session");
|
||||
expect(fetchMock.mock.calls[1]?.[0]).toBe(
|
||||
"/eve/agents/support/eve/v1/session/session_1/stream",
|
||||
"/eve/agents/support/eve/v1/session/session_1/stream?streamControlVersion=1",
|
||||
);
|
||||
});
|
||||
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
---
|
||||
issue: https://github.com/vercel/eve/issues/1159
|
||||
status: in-progress
|
||||
last_updated: "2026-09-15"
|
||||
---
|
||||
|
||||
# Leased session stream responses
|
||||
|
||||
## Purpose
|
||||
|
||||
A client that abandons a session stream expects `request.signal` to release the
|
||||
server-side Workflow reader. Some serverless transports do not propagate that
|
||||
cancellation, so each client reconnect can leave the previous invocation alive
|
||||
until the host timeout.
|
||||
|
||||
Treat an HTTP response as a renewable lease over the durable event stream. A
|
||||
capable client opts into a bounded response, receives transport heartbeats, and
|
||||
reconnects from its cursor after an explicit lease-end control record. Session
|
||||
events remain durable and the public `session.stream()` API does not change.
|
||||
|
||||
## Wire protocol
|
||||
|
||||
The client advertises the control records it can decode:
|
||||
|
||||
```http
|
||||
GET /eve/v1/session/:id/stream?streamControlVersion=1
|
||||
```
|
||||
|
||||
For version 1, the server:
|
||||
|
||||
- sends an ignored blank line every ten seconds without an event, keeping the
|
||||
existing 15-second client read timer from treating a healthy response as
|
||||
stalled;
|
||||
- ends the response after a fixed 60-second lease;
|
||||
- writes `{"$eve":"stream.lease-ended","version":1}` immediately before the
|
||||
intentional close; and
|
||||
- cancels the response's Workflow reader when the lease ends.
|
||||
|
||||
The client requests leases only for nonnegative cursors with stream reconnection
|
||||
enabled. Tail-relative reads and streams with reconnection disabled remain
|
||||
unleased because they cannot renew. The client consumes the control record
|
||||
internally, reconnects immediately from its absolute event cursor, and does not
|
||||
charge the lease renewal against its empty-stream retry budget. EOF without the
|
||||
control record retains the existing bounded, backed-off reconnect behavior.
|
||||
|
||||
The lease bounds cleanup even if heartbeats or the final control record are
|
||||
buffered by an intermediary. In that case the client's ordinary read timeout
|
||||
reconnects, while the abandoned server invocation ends no later than its lease.
|
||||
|
||||
## Compatibility
|
||||
|
||||
| Client | Server | Behavior |
|
||||
| ------- | ------- | ------------------------------------------------------------------------------- |
|
||||
| Current | Current | Leased response with heartbeats and explicit renewal |
|
||||
| Current | Older | Query parameter is ignored; existing read timeout and reconnect behavior remain |
|
||||
| Older | Current | No capability query parameter, so existing unleased response behavior remains |
|
||||
|
||||
The capability query parameter is forwarded through remote subagent stream
|
||||
proxies. An unknown control version is not negotiated. This avoids emitting records that an
|
||||
older client could mistake for session events.
|
||||
|
||||
## Non-goals
|
||||
|
||||
- Changing event persistence, cursors, or the public client API.
|
||||
- Solving intermediary compression independently of reconnect recovery.
|
||||
- Negotiating lease duration; duration remains an eve implementation detail.
|
||||
- Retrofitting bounded responses onto clients that do not advertise support.
|
||||
Reference in New Issue
Block a user