diff --git a/packages/cli/src/cli/commands/project/logs.ts b/packages/cli/src/cli/commands/project/logs.ts index 503551a69..4a3586204 100644 --- a/packages/cli/src/cli/commands/project/logs.ts +++ b/packages/cli/src/cli/commands/project/logs.ts @@ -1,3 +1,4 @@ +import type { Logger } from "@base44-cli/logger"; import type { Command } from "commander"; import { Option } from "commander"; import type { CLIContext, RunCommandResult } from "@/cli/types.js"; @@ -9,12 +10,17 @@ import type { FunctionLogsResponse, LogEnv, LogLevel, + LogStreamFilters, + StreamEndEvent, + StreamEvent, + StreamLogEvent, } from "@/core/resources/function/index.js"; import { fetchFunctionLogs, LogEnvSchema, LogLevelSchema, listDeployedFunctions, + openLogStream, } from "@/core/resources/function/index.js"; interface LogsOptions { @@ -123,14 +129,119 @@ function writeFollowLine(entry: LogEntry, jsonMode: boolean): void { process.stdout.write(`${line}\n`); } +const delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); + +function streamEventToLogEntry(event: StreamLogEvent): LogEntry { + return { + time: event.time, + level: event.level, + message: event.function + ? `[${event.function}] ${event.message}` + : event.message, + source: event.function ?? "", + }; +} + +interface StreamEnding { + lastTime: string; + producedEvents: boolean; + end: StreamEndEvent | null; +} + +async function printStreamUntilEnd( + stream: AsyncGenerator, + levelFilter: string | undefined, + jsonMode: boolean, + startTime: string, +): Promise { + let lastTime = startTime; + let producedEvents = false; + try { + for await (const event of stream) { + producedEvents = true; + if (event.kind === "end") + return { lastTime, producedEvents, end: event.end }; + if (levelFilter && event.log.level !== levelFilter) continue; + writeFollowLine(streamEventToLogEntry(event.log), jsonMode); + if (event.log.time > lastTime) lastTime = event.log.time; + } + } catch {} + return { lastTime, producedEvents, end: null }; +} + +const STREAM_RECONNECT_DELAY_MS = 1_000; +const MAX_DROPS_SINCE_LAST_EVENT = 2; + +interface StreamAttemptResult { + everConnected: boolean; + lastTime: string; +} + +async function streamUntilExhausted( + options: LogsOptions, + jsonMode: boolean, +): Promise { + const filters: LogStreamFilters = { + functions: parseFunctionNames(options.function), + env: options.env, + }; + let everConnected = false; + let lastTime = ""; + let dropsSinceLastEvent = 0; + while (true) { + const stream = await openLogStream(filters); + if (!stream) break; + everConnected = true; + const ending = await printStreamUntilEnd( + stream, + options.level, + jsonMode, + lastTime, + ); + lastTime = ending.lastTime; + if (ending.end) { + if (!ending.end.retriable) break; + dropsSinceLastEvent = 0; + } else { + dropsSinceLastEvent = ending.producedEvents ? 1 : dropsSinceLastEvent + 1; + if (dropsSinceLastEvent >= MAX_DROPS_SINCE_LAST_EVENT) break; + } + await delay(STREAM_RECONNECT_DELAY_MS); + } + return { everConnected, lastTime }; +} + async function followLogs( functionNames: string[], options: LogsOptions, availableFunctionNames: string[], jsonMode: boolean, + logger: Logger, ): Promise { - let state: FollowState = { lastTime: "", boundaryKeys: new Set() }; - let first = true; + const { everConnected, lastTime } = await streamUntilExhausted( + options, + jsonMode, + ); + logger.warn( + everConnected + ? "Realtime stream disconnected — falling back to polling (lines may lag ~20-30s)." + : "Realtime stream unavailable — falling back to polling (lines may lag ~20-30s).", + ); + return pollLogs(functionNames, options, availableFunctionNames, jsonMode, { + lastTime, + boundaryKeys: new Set(), + }); +} + +async function pollLogs( + functionNames: string[], + options: LogsOptions, + availableFunctionNames: string[], + jsonMode: boolean, + initialState: FollowState, +): Promise { + let state = initialState; + let first = state.lastTime === ""; while (true) { const pollOptions = first ? options : { ...options, since: state.lastTime }; @@ -144,7 +255,7 @@ async function followLogs( fresh.sort((a, b) => a.time.localeCompare(b.time)); for (const entry of fresh) writeFollowLine(entry, jsonMode); first = false; - await new Promise((resolve) => setTimeout(resolve, 2000)); + await delay(2000); } } @@ -279,6 +390,11 @@ async function logsAction( } if (options.follow) { + if (options.since) { + throw new InvalidInputError( + "--since cannot be combined with --follow yet (the realtime stream starts from now).", + ); + } if (options.until) { throw new InvalidInputError( "--until cannot be combined with --follow (a stream has no end).", @@ -295,6 +411,7 @@ async function logsAction( options, availableFunctionNames, ctx.jsonMode, + ctx.log, ); } diff --git a/packages/cli/src/core/resources/function/index.ts b/packages/cli/src/core/resources/function/index.ts index 08b63274b..514e1513a 100644 --- a/packages/cli/src/core/resources/function/index.ts +++ b/packages/cli/src/core/resources/function/index.ts @@ -4,3 +4,4 @@ export * from "./deploy.js"; export * from "./pull.js"; export * from "./resource.js"; export * from "./schema.js"; +export * from "./stream-api.js"; diff --git a/packages/cli/src/core/resources/function/stream-api.ts b/packages/cli/src/core/resources/function/stream-api.ts new file mode 100644 index 000000000..1ae784dcc --- /dev/null +++ b/packages/cli/src/core/resources/function/stream-api.ts @@ -0,0 +1,185 @@ +import { z } from "zod"; +import { + getWorkspaceApiKeyFromEnv, + isTokenExpired, + isWorkspaceApiKey, + readAuth, + refreshAndSaveTokens, +} from "@/core/auth/config.js"; +import { getBase44ApiUrl } from "@/core/config.js"; +import { getAppContext } from "@/core/project/index.js"; +import { + type LogEnv, + LogLevelSchema, +} from "@/core/resources/function/schema.js"; + +export const StreamLogEventSchema = z.object({ + time: z.string(), + level: z.preprocess( + (value) => (value === "warn" ? "warning" : value), + LogLevelSchema, + ), + function: z.string().nullable(), + message: z.string(), +}); + +export type StreamLogEvent = z.infer; + +const StreamEndEventSchema = z.object({ + reason: z.string(), + retriable: z.boolean(), +}); + +export type StreamEndEvent = z.infer; + +export type StreamEvent = + | { kind: "log"; log: StreamLogEvent } + | { kind: "end"; end: StreamEndEvent }; + +export interface LogStreamFilters { + functions?: string[]; + env?: LogEnv; +} + +function buildStreamUrl(filters: LogStreamFilters): string { + const { id } = getAppContext(); + const url = new URL( + `/api/apps/${id}/functions-mgmt/logs/stream`, + getBase44ApiUrl(), + ); + if (filters.functions?.length) { + url.searchParams.set("function", filters.functions.join(",")); + } + if (filters.env) { + url.searchParams.set("env", filters.env); + } + return url.href; +} + +async function buildStreamAuthHeaders(): Promise> { + const workspaceApiKey = getWorkspaceApiKeyFromEnv(); + if (workspaceApiKey && isWorkspaceApiKey(workspaceApiKey)) { + return { api_key: workspaceApiKey }; + } + const auth = await readAuth(); + if (isTokenExpired(auth)) { + const refreshedToken = await refreshAndSaveTokens(); + if (refreshedToken) { + return { Authorization: `Bearer ${refreshedToken}` }; + } + } + return { Authorization: `Bearer ${auth.accessToken}` }; +} + +export function parseStreamEvent( + eventName: string, + data: string, +): StreamEvent | null { + try { + const payload = JSON.parse(data); + if (eventName === "end") { + const result = StreamEndEventSchema.safeParse(payload); + return result.success ? { kind: "end", end: result.data } : null; + } + if (eventName === "") { + const result = StreamLogEventSchema.safeParse(payload); + return result.success ? { kind: "log", log: result.data } : null; + } + return null; + } catch { + return null; + } +} + +const STREAM_SILENCE_TIMEOUT_MS = 60_000; + +interface StreamReader { + read(): Promise<{ done: boolean; value?: Uint8Array }>; + cancel(): Promise; + releaseLock(): void; +} + +async function readOrSilence( + reader: StreamReader, +): Promise<{ done: boolean; value?: Uint8Array } | "silence"> { + let timer: ReturnType | undefined; + const silence = new Promise<"silence">((resolve) => { + timer = setTimeout(() => resolve("silence"), STREAM_SILENCE_TIMEOUT_MS); + }); + try { + return await Promise.race([reader.read(), silence]); + } finally { + clearTimeout(timer); + } +} + +async function* readLines( + body: ReadableStream, +): AsyncGenerator { + const reader: StreamReader = body.getReader(); + const decoder = new TextDecoder(); + let buffered = ""; + try { + while (true) { + const result = await readOrSilence(reader); + if (result === "silence") { + await reader.cancel(); + return; + } + if (result.done || !result.value) return; + buffered += decoder.decode(result.value, { stream: true }); + const lines = buffered.split("\n"); + buffered = lines.pop() ?? ""; + yield* lines; + } + } finally { + reader.releaseLock(); + } +} + +async function* readStreamEvents( + body: ReadableStream, +): AsyncGenerator { + let eventName = ""; + for await (const line of readLines(body)) { + if (line.startsWith("event:")) { + eventName = line.slice(6).trim(); + continue; + } + if (line.startsWith("data:")) { + const event = parseStreamEvent(eventName, line.slice(5)); + eventName = ""; + if (event) yield event; + continue; + } + if (line.trim() === "") eventName = ""; + } +} + +const STREAM_CONNECT_TIMEOUT_MS = 10_000; + +export async function openLogStream( + filters: LogStreamFilters, +): Promise | null> { + const connectPhase = new AbortController(); + const connectTimer = setTimeout( + () => connectPhase.abort(), + STREAM_CONNECT_TIMEOUT_MS, + ); + let response: Response; + try { + response = await fetch(buildStreamUrl(filters), { + headers: { + Accept: "text/event-stream", + ...(await buildStreamAuthHeaders()), + }, + signal: connectPhase.signal, + }); + } catch { + return null; + } finally { + clearTimeout(connectTimer); + } + if (!response.ok || !response.body) return null; + return readStreamEvents(response.body); +} diff --git a/packages/cli/tests/cli/logs.spec.ts b/packages/cli/tests/cli/logs.spec.ts index 0da92dd34..1aa5c3db4 100644 --- a/packages/cli/tests/cli/logs.spec.ts +++ b/packages/cli/tests/cli/logs.spec.ts @@ -4,6 +4,7 @@ import { type LogEntry, selectNewEntries, } from "@/cli/commands/project/logs.js"; +import { parseStreamEvent } from "@/core/resources/function/index.js"; import { fixture, setupCLITests } from "./testkit/index.js"; function entry(time: string, message: string): LogEntry { @@ -72,6 +73,67 @@ describe("selectNewEntries (follow dedup)", () => { }); }); +describe("parseStreamEvent (SSE log stream)", () => { + it("parses an unnamed data payload into a log event", () => { + const event = parseStreamEvent( + "", + '{"time":"2024-01-15T10:00:00Z","level":"info","function":"my-fn","message":"hello"}', + ); + + expect(event).toEqual({ + kind: "log", + log: { + time: "2024-01-15T10:00:00Z", + level: "info", + function: "my-fn", + message: "hello", + }, + }); + }); + + it("normalizes level warn to warning", () => { + const event = parseStreamEvent( + "", + '{"time":"2024-01-15T10:00:00Z","level":"warn","function":"my-fn","message":"careful"}', + ); + + expect(event?.kind === "log" && event.log.level).toBe("warning"); + }); + + it("keeps unattributed lines (null function)", () => { + const event = parseStreamEvent( + "", + '{"time":"2024-01-15T10:00:00Z","level":"error","function":null,"message":"boom"}', + ); + + expect(event?.kind === "log" && event.log.function).toBeNull(); + }); + + it("parses the typed end event with reason and retriable", () => { + const event = parseStreamEvent( + "end", + '{"reason":"tail_unavailable","retriable":false}', + ); + + expect(event).toEqual({ + kind: "end", + end: { reason: "tail_unavailable", retriable: false }, + }); + }); + + it("ignores unknown event names", () => { + expect( + parseStreamEvent("progress", '{"reason":"x","retriable":true}'), + ).toBeNull(); + }); + + it("ignores malformed payloads", () => { + expect(parseStreamEvent("", "not-json")).toBeNull(); + expect(parseStreamEvent("", '{"level":"info"}')).toBeNull(); + expect(parseStreamEvent("end", '{"reason":"x"}')).toBeNull(); + }); +}); + describe("logs command", () => { const t = setupCLITests(); @@ -265,6 +327,22 @@ describe("logs command", () => { t.expectResult(result).toContain("No production logs found"); }); + it("rejects --follow combined with --since", async () => { + await t.givenLoggedInWithProject(fixture("basic")); + + const result = await t.run( + "logs", + "--function", + "my-function", + "--follow", + "--since", + "1h", + ); + + t.expectResult(result).toFail(); + t.expectResult(result).toContain("--since cannot be combined"); + }); + it("rejects --follow combined with --order", async () => { await t.givenLoggedInWithProject(fixture("basic"));