From 19ce2362621bcd66057e841dce9ae9fb94ea1fa4 Mon Sep 17 00:00:00 2001 From: Owen Date: Fri, 21 Aug 2026 17:22:00 -0400 Subject: [PATCH] Support AI session log streaming --- messages/en-US.json | 2 ++ server/db/pg/schema/privateSchema.ts | 3 ++ server/db/sqlite/schema/privateSchema.ts | 3 ++ .../lib/logStreaming/LogStreamingManager.ts | 34 ++++++++++++++++++- server/private/lib/logStreaming/types.ts | 5 +-- .../createEventStreamingDestination.ts | 6 ++-- .../listEventStreamingDestinations.ts | 4 ++- .../updateEventStreamingDestination.ts | 6 ++-- src/components/HttpDestinationCredenza.tsx | 30 +++++++++++++++- src/components/S3DestinationCredenza.tsx | 29 +++++++++++++++- 10 files changed, 112 insertions(+), 10 deletions(-) diff --git a/messages/en-US.json b/messages/en-US.json index 0c59990a7..62ed7f30a 100644 --- a/messages/en-US.json +++ b/messages/en-US.json @@ -4087,6 +4087,8 @@ "httpDestConnectionLogsDescription": "Site and tunnel connection events, including connects and disconnects.", "httpDestRequestLogsTitle": "HTTP Request Logs", "httpDestRequestLogsDescription": "HTTP request logs for proxied resources, including method, path, and response code.", + "httpDestAISessionLogsTitle": "AI Session Logs", + "httpDestAISessionLogsDescription": "AI gateway request and response sessions, including prompts, model responses, and token usage.", "httpDestSaveChanges": "Save Changes", "httpDestCreateDestination": "Create Destination", "httpDestUpdatedSuccess": "Destination updated successfully", diff --git a/server/db/pg/schema/privateSchema.ts b/server/db/pg/schema/privateSchema.ts index e10b459e9..8e88527d6 100644 --- a/server/db/pg/schema/privateSchema.ts +++ b/server/db/pg/schema/privateSchema.ts @@ -468,6 +468,9 @@ export const eventStreamingDestinations = pgTable( sendRequestLogs: boolean("sendRequestLogs").notNull().default(false), sendActionLogs: boolean("sendActionLogs").notNull().default(false), sendAccessLogs: boolean("sendAccessLogs").notNull().default(false), + sendAISessionLogs: boolean("sendAISessionLogs") + .notNull() + .default(false), type: varchar("type", { length: 50 }).notNull(), // e.g. "http", "kafka", etc. config: text("config").notNull(), // JSON string with the configuration for the destination enabled: boolean("enabled").notNull().default(true), diff --git a/server/db/sqlite/schema/privateSchema.ts b/server/db/sqlite/schema/privateSchema.ts index da77bfed2..0003ae432 100644 --- a/server/db/sqlite/schema/privateSchema.ts +++ b/server/db/sqlite/schema/privateSchema.ts @@ -459,6 +459,9 @@ export const eventStreamingDestinations = sqliteTable( sendAccessLogs: integer("sendAccessLogs", { mode: "boolean" }) .notNull() .default(false), + sendAISessionLogs: integer("sendAISessionLogs", { mode: "boolean" }) + .notNull() + .default(false), type: text("type").notNull(), // e.g. "http", "kafka", etc. config: text("config").notNull(), // JSON string with the configuration for the destination enabled: integer("enabled", { mode: "boolean" }) diff --git a/server/private/lib/logStreaming/LogStreamingManager.ts b/server/private/lib/logStreaming/LogStreamingManager.ts index 03efc2809..14ada27b0 100644 --- a/server/private/lib/logStreaming/LogStreamingManager.ts +++ b/server/private/lib/logStreaming/LogStreamingManager.ts @@ -19,7 +19,8 @@ import { requestAuditLog, actionAuditLog, accessAuditLog, - connectionAuditLog + connectionAuditLog, + aiSessionLog } from "@server/db"; import logger from "@server/logger"; import { and, eq, gt, desc, max, sql } from "drizzle-orm"; @@ -309,6 +310,7 @@ export class LogStreamingManager { if (dest.sendActionLogs) enabledTypes.push("action"); if (dest.sendAccessLogs) enabledTypes.push("access"); if (dest.sendConnectionLogs) enabledTypes.push("connection"); + if (dest.sendAISessionLogs) enabledTypes.push("aiSession"); if (enabledTypes.length === 0) return; @@ -585,6 +587,13 @@ export class LogStreamingManager { .where(eq(connectionAuditLog.orgId, orgId)); return row?.maxId ?? 0; } + case "aiSession": { + const [row] = await logsDb + .select({ maxId: max(aiSessionLog.id) }) + .from(aiSessionLog) + .where(eq(aiSessionLog.orgId, orgId)); + return row?.maxId ?? 0; + } } } catch (err) { logger.warn( @@ -670,6 +679,21 @@ export class LogStreamingManager { .limit(limit)) as Array< Record & { id: number } >; + + case "aiSession": + return (await logsDb + .select() + .from(aiSessionLog) + .where( + and( + eq(aiSessionLog.orgId, orgId), + gt(aiSessionLog.id, afterId) + ) + ) + .orderBy(aiSessionLog.id) + .limit(limit)) as Array< + Record & { id: number } + >; } } @@ -694,6 +718,14 @@ export class LogStreamingManager { timestamp = typeof row.startedAt === "number" ? row.startedAt : 0; break; + case "aiSession": + // createdAt is stored as epoch milliseconds; normalise to + // epoch seconds to match the other log types. + timestamp = + typeof row.createdAt === "number" + ? Math.floor(row.createdAt / 1000) + : 0; + break; } const orgId = typeof row.orgId === "string" ? row.orgId : ""; diff --git a/server/private/lib/logStreaming/types.ts b/server/private/lib/logStreaming/types.ts index 193a5ff6b..06d5603e9 100644 --- a/server/private/lib/logStreaming/types.ts +++ b/server/private/lib/logStreaming/types.ts @@ -15,13 +15,14 @@ // Log type identifiers // --------------------------------------------------------------------------- -export type LogType = "request" | "action" | "access" | "connection"; +export type LogType = "request" | "action" | "access" | "connection" | "aiSession"; export const LOG_TYPES: LogType[] = [ "request", "action", "access", - "connection" + "connection", + "aiSession" ]; // --------------------------------------------------------------------------- diff --git a/server/private/routers/eventStreamingDestination/createEventStreamingDestination.ts b/server/private/routers/eventStreamingDestination/createEventStreamingDestination.ts index 7b000c5d8..c7f28f48a 100644 --- a/server/private/routers/eventStreamingDestination/createEventStreamingDestination.ts +++ b/server/private/routers/eventStreamingDestination/createEventStreamingDestination.ts @@ -37,7 +37,8 @@ const bodySchema = z.strictObject({ sendConnectionLogs: z.boolean().optional().default(false), sendRequestLogs: z.boolean().optional().default(false), sendActionLogs: z.boolean().optional().default(false), - sendAccessLogs: z.boolean().optional().default(false) + sendAccessLogs: z.boolean().optional().default(false), + sendAISessionLogs: z.boolean().optional().default(false) }); export type CreateEventStreamingDestinationResponse = { @@ -122,7 +123,8 @@ export async function createEventStreamingDestination( sendAccessLogs: parsedBody.data.sendAccessLogs, sendActionLogs: parsedBody.data.sendActionLogs, sendConnectionLogs: parsedBody.data.sendConnectionLogs, - sendRequestLogs: parsedBody.data.sendRequestLogs + sendRequestLogs: parsedBody.data.sendRequestLogs, + sendAISessionLogs: parsedBody.data.sendAISessionLogs }) .returning(); diff --git a/server/private/routers/eventStreamingDestination/listEventStreamingDestinations.ts b/server/private/routers/eventStreamingDestination/listEventStreamingDestinations.ts index dc741d482..d0f3d75c2 100644 --- a/server/private/routers/eventStreamingDestination/listEventStreamingDestinations.ts +++ b/server/private/routers/eventStreamingDestination/listEventStreamingDestinations.ts @@ -60,6 +60,7 @@ export type ListEventStreamingDestinationsResponse = { sendRequestLogs: boolean; sendActionLogs: boolean; sendAccessLogs: boolean; + sendAISessionLogs: boolean; }[]; pagination: { total: number; @@ -83,7 +84,8 @@ const ListEventStreamingDestinationsResponseDataSchema = z.object({ sendConnectionLogs: z.boolean(), sendRequestLogs: z.boolean(), sendActionLogs: z.boolean(), - sendAccessLogs: z.boolean() + sendAccessLogs: z.boolean(), + sendAISessionLogs: z.boolean() }) ), pagination: z.object({ diff --git a/server/private/routers/eventStreamingDestination/updateEventStreamingDestination.ts b/server/private/routers/eventStreamingDestination/updateEventStreamingDestination.ts index 84202bf8e..ddf22cfba 100644 --- a/server/private/routers/eventStreamingDestination/updateEventStreamingDestination.ts +++ b/server/private/routers/eventStreamingDestination/updateEventStreamingDestination.ts @@ -40,7 +40,8 @@ const bodySchema = z.strictObject({ sendConnectionLogs: z.boolean().optional(), sendRequestLogs: z.boolean().optional(), sendActionLogs: z.boolean().optional(), - sendAccessLogs: z.boolean().optional() + sendAccessLogs: z.boolean().optional(), + sendAISessionLogs: z.boolean().optional() }); export type UpdateEventStreamingDestinationResponse = { @@ -125,7 +126,7 @@ export async function updateEventStreamingDestination( ); } - const { type, config: configToUpdate, enabled, sendAccessLogs, sendActionLogs, sendConnectionLogs, sendRequestLogs } = parsedBody.data; + const { type, config: configToUpdate, enabled, sendAccessLogs, sendActionLogs, sendConnectionLogs, sendRequestLogs, sendAISessionLogs } = parsedBody.data; const updateData: Record = { updatedAt: Date.now() @@ -141,6 +142,7 @@ export async function updateEventStreamingDestination( if (sendActionLogs !== undefined) updateData.sendActionLogs = sendActionLogs; if (sendConnectionLogs !== undefined) updateData.sendConnectionLogs = sendConnectionLogs; if (sendRequestLogs !== undefined) updateData.sendRequestLogs = sendRequestLogs; + if (sendAISessionLogs !== undefined) updateData.sendAISessionLogs = sendAISessionLogs; await db .update(eventStreamingDestinations) diff --git a/src/components/HttpDestinationCredenza.tsx b/src/components/HttpDestinationCredenza.tsx index 85d32fd5c..92b684756 100644 --- a/src/components/HttpDestinationCredenza.tsx +++ b/src/components/HttpDestinationCredenza.tsx @@ -57,6 +57,7 @@ export interface Destination { sendActionLogs: boolean; sendConnectionLogs: boolean; sendRequestLogs: boolean; + sendAISessionLogs: boolean; lastError: string | null; lastErrorAt: number | null; createdAt: number; @@ -180,6 +181,7 @@ export function HttpDestinationCredenza({ const [sendActionLogs, setSendActionLogs] = useState(false); const [sendConnectionLogs, setSendConnectionLogs] = useState(false); const [sendRequestLogs, setSendRequestLogs] = useState(false); + const [sendAISessionLogs, setSendAISessionLogs] = useState(false); useEffect(() => { if (open) { @@ -190,6 +192,7 @@ export function HttpDestinationCredenza({ setSendActionLogs(editing?.sendActionLogs ?? false); setSendConnectionLogs(editing?.sendConnectionLogs ?? false); setSendRequestLogs(editing?.sendRequestLogs ?? false); + setSendAISessionLogs(editing?.sendAISessionLogs ?? false); } }, [open, editing]); @@ -226,7 +229,8 @@ export function HttpDestinationCredenza({ sendAccessLogs, sendActionLogs, sendConnectionLogs, - sendRequestLogs + sendRequestLogs, + sendAISessionLogs }; if (editing) { await api.post( @@ -778,6 +782,30 @@ export function HttpDestinationCredenza({

+ +
+ + setSendAISessionLogs(v === true) + } + className="mt-0.5" + /> +
+ +

+ {t( + "httpDestAISessionLogsDescription" + )} +

+
+
diff --git a/src/components/S3DestinationCredenza.tsx b/src/components/S3DestinationCredenza.tsx index e6c128805..a66406f00 100644 --- a/src/components/S3DestinationCredenza.tsx +++ b/src/components/S3DestinationCredenza.tsx @@ -90,6 +90,7 @@ export function S3DestinationCredenza({ const [sendActionLogs, setSendActionLogs] = useState(false); const [sendConnectionLogs, setSendConnectionLogs] = useState(false); const [sendRequestLogs, setSendRequestLogs] = useState(false); + const [sendAISessionLogs, setSendAISessionLogs] = useState(false); useEffect(() => { if (open) { @@ -98,6 +99,7 @@ export function S3DestinationCredenza({ setSendActionLogs(editing?.sendActionLogs ?? false); setSendConnectionLogs(editing?.sendConnectionLogs ?? false); setSendRequestLogs(editing?.sendRequestLogs ?? false); + setSendAISessionLogs(editing?.sendAISessionLogs ?? false); } }, [open, editing]); @@ -121,7 +123,8 @@ export function S3DestinationCredenza({ sendAccessLogs, sendActionLogs, sendConnectionLogs, - sendRequestLogs + sendRequestLogs, + sendAISessionLogs }; if (editing) { await api.post( @@ -510,6 +513,30 @@ export function S3DestinationCredenza({

+ +
+ + setSendAISessionLogs(v === true) + } + className="mt-0.5" + /> +
+ +

+ {t( + "httpDestAISessionLogsDescription" + )} +

+
+