Skip to content
123 changes: 120 additions & 3 deletions packages/cli/src/cli/commands/project/logs.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand All @@ -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 {
Expand Down Expand Up @@ -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<StreamEvent>,
levelFilter: string | undefined,
jsonMode: boolean,
startTime: string,
): Promise<StreamEnding> {
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<StreamAttemptResult> {
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<never> {
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<never> {
let state = initialState;
let first = state.lastTime === "";

while (true) {
const pollOptions = first ? options : { ...options, since: state.lastTime };
Expand All @@ -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);
}
}

Expand Down Expand Up @@ -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).",
Expand All @@ -295,6 +411,7 @@ async function logsAction(
options,
availableFunctionNames,
ctx.jsonMode,
ctx.log,
);
}

Expand Down
1 change: 1 addition & 0 deletions packages/cli/src/core/resources/function/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
185 changes: 185 additions & 0 deletions packages/cli/src/core/resources/function/stream-api.ts
Original file line number Diff line number Diff line change
@@ -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<typeof StreamLogEventSchema>;

const StreamEndEventSchema = z.object({
reason: z.string(),
retriable: z.boolean(),
});

export type StreamEndEvent = z.infer<typeof StreamEndEventSchema>;

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<Record<string, string>> {
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<unknown>;
releaseLock(): void;
}

async function readOrSilence(
reader: StreamReader,
): Promise<{ done: boolean; value?: Uint8Array } | "silence"> {
let timer: ReturnType<typeof setTimeout> | 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<Uint8Array>,
): AsyncGenerator<string> {
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<Uint8Array>,
): AsyncGenerator<StreamEvent> {
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<AsyncGenerator<StreamEvent> | 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);
}
Loading
Loading