Fix terminal log stream authorization
This commit is contained in:
@@ -778,6 +778,30 @@ describe("PlatformApiClient AI providers", () => {
|
||||
expect(client.serverLogEventsUrl("server-1")).toBe("/api/v1/server-instances/server-1/logs/events");
|
||||
});
|
||||
|
||||
it("streams terminal log SSE with bearer authorization and dispatches live events", async () => {
|
||||
const payload = "event: stream\ndata: {\"id\":\"log-1\",\"serverInstanceId\":\"server-1\"}\n\nevent: log\ndata: {\"streamId\":\"log-1\",\"entry\":{\"seq\":3,\"line\":\"live line\"}}\n\nevent: ready\ndata: {\"serverInstanceId\":\"server-1\"}\n\n";
|
||||
const fetchMock = vi.fn(async (_input: RequestInfo | URL, init?: RequestInit) => {
|
||||
expect(new Headers(init?.headers).get("Authorization")).toBe("Bearer terminal-session");
|
||||
expect(new Headers(init?.headers).get("Accept")).toBe("text/event-stream");
|
||||
return new Response(new ReadableStream<Uint8Array>({ start(controller) { controller.enqueue(new TextEncoder().encode(payload)); controller.close(); } }), { status: 200, headers: { "Content-Type": "text/event-stream" } });
|
||||
});
|
||||
vi.stubGlobal("fetch", fetchMock);
|
||||
const client = new PlatformApiClient("/api/v1", () => "terminal-session");
|
||||
const events: string[] = [];
|
||||
const stream = client.openServerLogEvents("server-1", { historyLimit: 2 });
|
||||
stream.addEventListener("stream", (event) => events.push(`stream:${event.data}`));
|
||||
stream.addEventListener("log", (event) => events.push(`log:${event.data}`));
|
||||
stream.addEventListener("ready", (event) => events.push(`ready:${event.data}`));
|
||||
|
||||
await vi.waitFor(() => expect(events).toEqual([
|
||||
'stream:{"id":"log-1","serverInstanceId":"server-1"}',
|
||||
'log:{"streamId":"log-1","entry":{"seq":3,"line":"live line"}}',
|
||||
'ready:{"serverInstanceId":"server-1"}'
|
||||
]));
|
||||
expect(String(fetchMock.mock.calls[0]?.[0])).toBe("/api/v1/server-instances/server-1/logs/events?historyLimit=2");
|
||||
stream.close();
|
||||
});
|
||||
|
||||
it("surfaces password confirmation denials without exposing generic forbidden text", async () => {
|
||||
vi.stubGlobal("fetch", vi.fn(async () => new Response(JSON.stringify({
|
||||
code: "forbidden",
|
||||
|
||||
@@ -152,6 +152,13 @@ export class PlatformApiError extends Error {
|
||||
}
|
||||
}
|
||||
|
||||
export interface PlatformEventStream {
|
||||
onerror: ((event: Event) => void) | null;
|
||||
addEventListener(type: string, listener: (event: MessageEvent) => void): void;
|
||||
removeEventListener(type: string, listener: (event: MessageEvent) => void): void;
|
||||
close(): void;
|
||||
}
|
||||
|
||||
export class PlatformApiClient {
|
||||
constructor(private readonly baseUrl = "/api/v1", private readonly sessionTokenProvider: () => string | null = () => platformApiSessionToken) {}
|
||||
|
||||
@@ -638,8 +645,13 @@ export class PlatformApiClient {
|
||||
return this.request<LogStreamListResponse>("/log-streams");
|
||||
}
|
||||
|
||||
openServerLogEvents(id: string, options: LogStreamEventOptions = {}): EventSource {
|
||||
return new EventSource(this.serverLogEventsUrl(id, options), { withCredentials: true });
|
||||
openServerLogEvents(id: string, options: LogStreamEventOptions = {}): PlatformEventStream {
|
||||
const url = this.serverLogEventsUrl(id, options);
|
||||
const sessionToken = this.sessionTokenProvider();
|
||||
if (!sessionToken) {
|
||||
return new EventSource(url, { withCredentials: true });
|
||||
}
|
||||
return new FetchServerSentEventStream(url, sessionToken);
|
||||
}
|
||||
|
||||
serverLogEventsUrl(id: string, options: LogStreamEventOptions = {}): string {
|
||||
@@ -777,6 +789,89 @@ async function responseError(response: Response): Promise<PlatformApiError> {
|
||||
return error;
|
||||
}
|
||||
|
||||
class FetchServerSentEventStream implements PlatformEventStream {
|
||||
onerror: ((event: Event) => void) | null = null;
|
||||
private readonly controller = new AbortController();
|
||||
private readonly listeners = new Map<string, Set<(event: MessageEvent) => void>>();
|
||||
|
||||
constructor(private readonly url: string, private readonly sessionToken: string) {
|
||||
void this.connect();
|
||||
}
|
||||
|
||||
addEventListener(type: string, listener: (event: MessageEvent) => void): void {
|
||||
const listeners = this.listeners.get(type) ?? new Set<(event: MessageEvent) => void>();
|
||||
listeners.add(listener);
|
||||
this.listeners.set(type, listeners);
|
||||
}
|
||||
|
||||
removeEventListener(type: string, listener: (event: MessageEvent) => void): void {
|
||||
this.listeners.get(type)?.delete(listener);
|
||||
}
|
||||
|
||||
close(): void {
|
||||
this.controller.abort();
|
||||
this.listeners.clear();
|
||||
}
|
||||
|
||||
private async connect(): Promise<void> {
|
||||
try {
|
||||
const response = await fetch(this.url, { method: "GET", credentials: "include", headers: { Accept: "text/event-stream", Authorization: `Bearer ${this.sessionToken}` }, signal: this.controller.signal });
|
||||
if (!response.ok || !response.body) {
|
||||
throw new Error(`event stream failed: ${response.status}`);
|
||||
}
|
||||
await this.read(response.body);
|
||||
} catch {
|
||||
if (!this.controller.signal.aborted) {
|
||||
this.onerror?.(new Event("error"));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async read(body: ReadableStream<Uint8Array>): Promise<void> {
|
||||
const reader = body.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
let buffer = "";
|
||||
let eventName = "message";
|
||||
let eventId = "";
|
||||
let dataLines: string[] = [];
|
||||
const processLine = (rawLine: string) => {
|
||||
const line = rawLine.endsWith("\r") ? rawLine.slice(0, -1) : rawLine;
|
||||
if (line === "") {
|
||||
if (dataLines.length > 0) {
|
||||
this.dispatch(eventName, dataLines.join("\n"), eventId);
|
||||
}
|
||||
eventName = "message";
|
||||
eventId = "";
|
||||
dataLines = [];
|
||||
return;
|
||||
}
|
||||
if (line.startsWith(":")) return;
|
||||
const separator = line.indexOf(":");
|
||||
const field = separator === -1 ? line : line.slice(0, separator);
|
||||
const value = separator === -1 ? "" : line.slice(separator + 1).replace(/^ /, "");
|
||||
if (field === "event") eventName = value || "message";
|
||||
if (field === "id") eventId = value;
|
||||
if (field === "data") dataLines.push(value);
|
||||
};
|
||||
while (!this.controller.signal.aborted) {
|
||||
const { value, done } = await reader.read();
|
||||
if (done) break;
|
||||
buffer += decoder.decode(value, { stream: true });
|
||||
const lines = buffer.split("\n");
|
||||
buffer = lines.pop() ?? "";
|
||||
lines.forEach(processLine);
|
||||
}
|
||||
}
|
||||
|
||||
private dispatch(type: string, data: string, lastEventId: string): void {
|
||||
const event = new MessageEvent(type, { data, lastEventId });
|
||||
this.listeners.get(type)?.forEach((listener) => listener(event));
|
||||
if (type !== "message") {
|
||||
this.listeners.get("message")?.forEach((listener) => listener(event));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function safeValidationMessage(apiError?: ApiErrorResponse | null): string {
|
||||
const details = apiError?.details
|
||||
?.map((detail) => safeValidationDetail(detail))
|
||||
|
||||
Reference in New Issue
Block a user