mirror of
https://github.com/fosrl/pangolin.git
synced 2026-08-13 16:00:02 +02:00
send, process, store, display virtual api key ai information in usage and sessions
This commit is contained in:
@@ -1819,18 +1819,22 @@ export const aiUsageRecords = pgTable(
|
||||
.references(() => orgs.orgId, { onDelete: "cascade" }),
|
||||
providerId: integer("providerId")
|
||||
.notNull()
|
||||
.references(() => aiProviders.providerId, { onDelete: "cascade" }),
|
||||
.references(() => aiProviders.providerId, { onDelete: "set null" }),
|
||||
resourceId: integer("resourceId").references(
|
||||
() => resources.resourceId,
|
||||
{ onDelete: "cascade" }
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
siteResourceId: integer("siteResourceId").references(
|
||||
() => siteResources.siteResourceId,
|
||||
{ onDelete: "cascade" }
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
userId: varchar("userId").references(() => users.userId, {
|
||||
onDelete: "set null"
|
||||
}),
|
||||
virtualApiKeyId: varchar("virtualApiKeyId").references(
|
||||
() => virtualApiKeys.virtualApiKeyId,
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
// Links this usage record back to the aiSessionLog row for the same
|
||||
// request (aiSessionLog.sessionId), so token/cost usage can be shown
|
||||
// alongside the session transcript. Not a DB-level FK - aiSessionLog
|
||||
@@ -1870,6 +1874,11 @@ export const aiUsageRecords = pgTable(
|
||||
t.userId,
|
||||
t.createdAt
|
||||
),
|
||||
index("idx_ai_usage_records_org_virtual_api_key_created").on(
|
||||
t.orgId,
|
||||
t.virtualApiKeyId,
|
||||
t.createdAt
|
||||
),
|
||||
index("idx_ai_usage_records_session").on(t.sessionId)
|
||||
]
|
||||
);
|
||||
@@ -1927,19 +1936,23 @@ export const aiSessionLog = pgTable(
|
||||
}),
|
||||
providerId: integer("providerId")
|
||||
.notNull()
|
||||
.references(() => aiProviders.providerId, { onDelete: "cascade" }),
|
||||
.references(() => aiProviders.providerId, { onDelete: "set null" }),
|
||||
capability: varchar("capability").notNull(),
|
||||
resourceId: integer("resourceId").references(
|
||||
() => resources.resourceId,
|
||||
{ onDelete: "cascade" }
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
siteResourceId: integer("siteResourceId").references(
|
||||
() => siteResources.siteResourceId,
|
||||
{ onDelete: "cascade" }
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
userId: varchar("userId").references(() => users.userId, {
|
||||
onDelete: "set null"
|
||||
}),
|
||||
virtualApiKeyId: varchar("virtualApiKeyId").references(
|
||||
() => virtualApiKeys.virtualApiKeyId,
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
requestedModel: varchar("requestedModel"),
|
||||
isStream: boolean("isStream").notNull().default(false),
|
||||
requestBody: text("requestBody"),
|
||||
@@ -1979,6 +1992,11 @@ export const aiSessionLog = pgTable(
|
||||
t.userId,
|
||||
t.createdAt
|
||||
),
|
||||
index("idx_ai_session_log_org_virtual_api_key_created").on(
|
||||
t.orgId,
|
||||
t.virtualApiKeyId,
|
||||
t.createdAt
|
||||
),
|
||||
index("idx_ai_session_log_session").on(t.sessionId)
|
||||
]
|
||||
);
|
||||
|
||||
@@ -1807,18 +1807,22 @@ export const aiUsageRecords = sqliteTable(
|
||||
.references(() => orgs.orgId, { onDelete: "cascade" }),
|
||||
providerId: integer("providerId")
|
||||
.notNull()
|
||||
.references(() => aiProviders.providerId, { onDelete: "cascade" }),
|
||||
.references(() => aiProviders.providerId, { onDelete: "set null" }),
|
||||
resourceId: integer("resourceId").references(
|
||||
() => resources.resourceId,
|
||||
{ onDelete: "cascade" }
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
siteResourceId: integer("siteResourceId").references(
|
||||
() => siteResources.siteResourceId,
|
||||
{ onDelete: "cascade" }
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
userId: text("userId").references(() => users.userId, {
|
||||
onDelete: "set null"
|
||||
}),
|
||||
virtualApiKeyId: text("virtualApiKeyId").references(
|
||||
() => virtualApiKeys.virtualApiKeyId,
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
// Links this usage record back to the aiSessionLog row for the same
|
||||
// request (aiSessionLog.sessionId), so token/cost usage can be shown
|
||||
// alongside the session transcript. Not a DB-level FK - aiSessionLog
|
||||
@@ -1860,6 +1864,11 @@ export const aiUsageRecords = sqliteTable(
|
||||
t.userId,
|
||||
t.createdAt
|
||||
),
|
||||
index("idx_ai_usage_records_org_virtual_api_key_created").on(
|
||||
t.orgId,
|
||||
t.virtualApiKeyId,
|
||||
t.createdAt
|
||||
),
|
||||
index("idx_ai_usage_records_session").on(t.sessionId)
|
||||
]
|
||||
);
|
||||
@@ -1917,19 +1926,23 @@ export const aiSessionLog = sqliteTable(
|
||||
}),
|
||||
providerId: integer("providerId")
|
||||
.notNull()
|
||||
.references(() => aiProviders.providerId, { onDelete: "cascade" }),
|
||||
.references(() => aiProviders.providerId, { onDelete: "set null" }),
|
||||
capability: text("capability").notNull(),
|
||||
resourceId: integer("resourceId").references(
|
||||
() => resources.resourceId,
|
||||
{ onDelete: "cascade" }
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
siteResourceId: integer("siteResourceId").references(
|
||||
() => siteResources.siteResourceId,
|
||||
{ onDelete: "cascade" }
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
userId: text("userId").references(() => users.userId, {
|
||||
onDelete: "set null"
|
||||
}),
|
||||
virtualApiKeyId: text("virtualApiKeyId").references(
|
||||
() => virtualApiKeys.virtualApiKeyId,
|
||||
{ onDelete: "set null" }
|
||||
),
|
||||
requestedModel: text("requestedModel"),
|
||||
isStream: integer("isStream", { mode: "boolean" })
|
||||
.notNull()
|
||||
@@ -1973,6 +1986,11 @@ export const aiSessionLog = sqliteTable(
|
||||
t.userId,
|
||||
t.createdAt
|
||||
),
|
||||
index("idx_ai_session_log_org_virtual_api_key_created").on(
|
||||
t.orgId,
|
||||
t.virtualApiKeyId,
|
||||
t.createdAt
|
||||
),
|
||||
index("idx_ai_session_log_session").on(t.sessionId)
|
||||
]
|
||||
);
|
||||
|
||||
@@ -1,4 +1,14 @@
|
||||
import { and, eq, gte, inArray, isNull, or, sql, SQL, type InferInsertModel } from "drizzle-orm";
|
||||
import {
|
||||
and,
|
||||
eq,
|
||||
gte,
|
||||
inArray,
|
||||
isNull,
|
||||
or,
|
||||
sql,
|
||||
SQL,
|
||||
type InferInsertModel
|
||||
} from "drizzle-orm";
|
||||
import {
|
||||
AiBudget,
|
||||
aiBudgetBreachEvents,
|
||||
@@ -424,6 +434,7 @@ export type UsageRecordInput = {
|
||||
resourceId: number | null;
|
||||
siteResourceId: number | null;
|
||||
userId: string | null;
|
||||
virtualApiKeyId: string | null;
|
||||
requestedModel: string;
|
||||
usage: AiUsage;
|
||||
costUsd: number | null;
|
||||
@@ -460,7 +471,10 @@ async function flushUsageRecords() {
|
||||
|
||||
isUsageFlushInProgress = true;
|
||||
|
||||
const recordsToWrite = usageRecordBuffer.splice(0, usageRecordBuffer.length);
|
||||
const recordsToWrite = usageRecordBuffer.splice(
|
||||
0,
|
||||
usageRecordBuffer.length
|
||||
);
|
||||
|
||||
try {
|
||||
// Use a transaction to ensure all inserts succeed or fail together
|
||||
@@ -472,16 +486,25 @@ async function flushUsageRecords() {
|
||||
await tx.insert(aiUsageRecords).values(batch);
|
||||
}
|
||||
});
|
||||
logger.debug(`Flushed ${recordsToWrite.length} AI usage records to database`);
|
||||
logger.debug(
|
||||
`Flushed ${recordsToWrite.length} AI usage records to database`
|
||||
);
|
||||
} catch (error) {
|
||||
logger.error("Error flushing AI usage records:", error);
|
||||
// On transaction error, put records back at the front of the buffer
|
||||
// to retry, but only if the buffer isn't too large
|
||||
if (usageRecordBuffer.length < USAGE_MAX_BUFFER_SIZE - recordsToWrite.length) {
|
||||
if (
|
||||
usageRecordBuffer.length <
|
||||
USAGE_MAX_BUFFER_SIZE - recordsToWrite.length
|
||||
) {
|
||||
usageRecordBuffer.unshift(...recordsToWrite);
|
||||
logger.info(`Re-queued ${recordsToWrite.length} AI usage records for retry`);
|
||||
logger.info(
|
||||
`Re-queued ${recordsToWrite.length} AI usage records for retry`
|
||||
);
|
||||
} else {
|
||||
logger.error(`Buffer full, dropped ${recordsToWrite.length} AI usage records`);
|
||||
logger.error(
|
||||
`Buffer full, dropped ${recordsToWrite.length} AI usage records`
|
||||
);
|
||||
}
|
||||
} finally {
|
||||
isUsageFlushInProgress = false;
|
||||
@@ -544,6 +567,7 @@ export async function recordUsage(input: UsageRecordInput): Promise<void> {
|
||||
resourceId: input.resourceId,
|
||||
siteResourceId: input.siteResourceId,
|
||||
userId: input.userId,
|
||||
virtualApiKeyId: input.virtualApiKeyId,
|
||||
sessionId: input.sessionId,
|
||||
requestedModel: input.requestedModel,
|
||||
promptTokens: usage.promptTokens,
|
||||
|
||||
@@ -180,6 +180,7 @@ export function logAiSession(data: {
|
||||
resourceId: number | null;
|
||||
siteResourceId: number | null;
|
||||
requestUserId: string | null;
|
||||
virtualApiKeyId: string | null;
|
||||
}): void {
|
||||
(async () => {
|
||||
try {
|
||||
@@ -237,6 +238,9 @@ export function logAiSession(data: {
|
||||
resourceId: data.resourceId ?? undefined,
|
||||
siteResourceId: data.siteResourceId ?? undefined,
|
||||
userId: sanitizeString(data.requestUserId ?? undefined),
|
||||
virtualApiKeyId: sanitizeString(
|
||||
data.virtualApiKeyId ?? undefined
|
||||
),
|
||||
requestedModel: sanitizeString(data.requestedModel),
|
||||
isStream: data.isStream,
|
||||
requestBody: sanitizeString(requestBodyText.value),
|
||||
|
||||
@@ -161,6 +161,15 @@ export type RequestUser = {
|
||||
roleIds: number[];
|
||||
};
|
||||
|
||||
// Identity resolved for a gateway request: the app/session or virtual-API-key
|
||||
// user (if any) plus the virtual API key that authenticated the request (if
|
||||
// any) - a manual virtual API key with no associated user has a
|
||||
// virtualApiKeyId but no user.
|
||||
export type RequestIdentity = {
|
||||
user: RequestUser | null;
|
||||
virtualApiKeyId: string | null;
|
||||
};
|
||||
|
||||
// Identity headers forwarded to the upstream inference endpoint when the
|
||||
// requesting user is known. Omitted entirely (not sent empty) when we
|
||||
// couldn't resolve a user for the request.
|
||||
@@ -223,10 +232,12 @@ async function resolveRequestUser(
|
||||
req: Request,
|
||||
resourceId: number | null,
|
||||
orgId: string | null
|
||||
): Promise<RequestUser | null> {
|
||||
): Promise<RequestIdentity> {
|
||||
// Public inference: identity comes from Badger via Remote-* only when the
|
||||
// Traefik trust header proves the request passed verify-session (VAK).
|
||||
if (isAiGatewayTrustHeaderValid(req.headers as Record<string, string>)) {
|
||||
const virtualApiKeyId =
|
||||
getRequestHeader(req, "remote-virtual-api-key-id") || null;
|
||||
const userId = getRequestHeader(req, "remote-user-id");
|
||||
if (userId) {
|
||||
const username = getRequestHeader(req, "remote-user") || userId;
|
||||
@@ -236,32 +247,41 @@ async function resolveRequestUser(
|
||||
const orgRoles = orgId ? await getUserOrgRoles(userId, orgId) : [];
|
||||
|
||||
return {
|
||||
userId,
|
||||
username,
|
||||
email: email || null,
|
||||
name: name || null,
|
||||
role:
|
||||
role || orgRoles.map((r) => r.roleName).join(", ") || null,
|
||||
roleIds: orgRoles.map((r) => r.roleId)
|
||||
user: {
|
||||
userId,
|
||||
username,
|
||||
email: email || null,
|
||||
name: name || null,
|
||||
role:
|
||||
role ||
|
||||
orgRoles.map((r) => r.roleName).join(", ") ||
|
||||
null,
|
||||
roleIds: orgRoles.map((r) => r.roleId)
|
||||
},
|
||||
virtualApiKeyId
|
||||
};
|
||||
}
|
||||
|
||||
// Trusted request with no associated user (manual key without userId).
|
||||
// Trusted request with no associated user (manual key without
|
||||
// userId) - still attribute usage to the virtual API key itself.
|
||||
if (resourceId != null) {
|
||||
return null;
|
||||
return { user: null, virtualApiKeyId };
|
||||
}
|
||||
}
|
||||
|
||||
// Public inference must come through Badger; do not authorize via app session.
|
||||
if (resourceId != null) {
|
||||
return null;
|
||||
return { user: null, virtualApiKeyId: null };
|
||||
}
|
||||
|
||||
const sessionToken = req.cookies?.[SESSION_COOKIE_NAME];
|
||||
if (sessionToken) {
|
||||
const { session, user } = await validateSessionToken(sessionToken);
|
||||
if (session && user) {
|
||||
return buildRequestUser(user.userId, orgId);
|
||||
return {
|
||||
user: await buildRequestUser(user.userId, orgId),
|
||||
virtualApiKeyId: null
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -269,7 +289,7 @@ async function resolveRequestUser(
|
||||
|
||||
const ip = req.ip;
|
||||
if (!ip) {
|
||||
return null;
|
||||
return { user: null, virtualApiKeyId: null };
|
||||
}
|
||||
|
||||
const exitNodeRanges = await getExitNodeRanges();
|
||||
@@ -277,15 +297,18 @@ async function resolveRequestUser(
|
||||
isIpInCidr(ip, range)
|
||||
);
|
||||
if (!inExitNodeRange) {
|
||||
return null;
|
||||
return { user: null, virtualApiKeyId: null };
|
||||
}
|
||||
|
||||
const client = await findClientByIp(ip);
|
||||
if (!client || !client.userId) {
|
||||
return null;
|
||||
return { user: null, virtualApiKeyId: null };
|
||||
}
|
||||
|
||||
return buildRequestUser(client.userId, orgId);
|
||||
return {
|
||||
user: await buildRequestUser(client.userId, orgId),
|
||||
virtualApiKeyId: null
|
||||
};
|
||||
}
|
||||
|
||||
function getRequestHeader(req: Request, name: string): string | undefined {
|
||||
@@ -608,6 +631,7 @@ export function recordAiGatewayCompletion(args: {
|
||||
resourceId: number | null;
|
||||
siteResourceId: number | null;
|
||||
requestUserId: string | null;
|
||||
virtualApiKeyId: string | null;
|
||||
budgets: AiBudget[];
|
||||
}): void {
|
||||
const {
|
||||
@@ -623,6 +647,7 @@ export function recordAiGatewayCompletion(args: {
|
||||
resourceId,
|
||||
siteResourceId,
|
||||
requestUserId,
|
||||
virtualApiKeyId,
|
||||
budgets
|
||||
} = args;
|
||||
|
||||
@@ -668,6 +693,7 @@ export function recordAiGatewayCompletion(args: {
|
||||
resourceId,
|
||||
siteResourceId,
|
||||
userId: requestUserId,
|
||||
virtualApiKeyId,
|
||||
requestedModel: model ?? "unknown",
|
||||
usage,
|
||||
costUsd: cost?.totalCost ?? null,
|
||||
@@ -699,7 +725,8 @@ export function recordAiGatewayCompletion(args: {
|
||||
orgId,
|
||||
resourceId,
|
||||
siteResourceId,
|
||||
requestUserId
|
||||
requestUserId,
|
||||
virtualApiKeyId
|
||||
});
|
||||
}
|
||||
|
||||
@@ -770,7 +797,7 @@ export async function handleAiGatewayProxy(
|
||||
|
||||
const requestedModel = def.extractModel(req);
|
||||
|
||||
const [requestUser, selection] = await Promise.all([
|
||||
const [identity, selection] = await Promise.all([
|
||||
resolveRequestUser(req, resourceId, orgId),
|
||||
selectProvider(
|
||||
capableAttachments,
|
||||
@@ -778,6 +805,7 @@ export async function handleAiGatewayProxy(
|
||||
requestedModel
|
||||
)
|
||||
]);
|
||||
const requestUser = identity.user;
|
||||
|
||||
if (requestUser) {
|
||||
logger.debug(
|
||||
@@ -836,7 +864,8 @@ export async function handleAiGatewayProxy(
|
||||
resourceId,
|
||||
siteResourceId,
|
||||
requestedModel,
|
||||
budgets: appliedBudgets
|
||||
budgets: appliedBudgets,
|
||||
virtualApiKeyId: identity.virtualApiKeyId
|
||||
}
|
||||
);
|
||||
}
|
||||
@@ -987,6 +1016,7 @@ export async function handleAiGatewayProxy(
|
||||
resourceId,
|
||||
siteResourceId,
|
||||
requestUserId: requestUser?.userId ?? null,
|
||||
virtualApiKeyId: identity.virtualApiKeyId,
|
||||
budgets: appliedBudgets
|
||||
});
|
||||
}
|
||||
|
||||
@@ -179,6 +179,7 @@ export async function proxyAiGatewayToSiteTarget(
|
||||
siteResourceId: number | null;
|
||||
requestedModel: string | undefined;
|
||||
budgets: AiBudget[];
|
||||
virtualApiKeyId: string | null;
|
||||
}
|
||||
): Promise<void> {
|
||||
const providerTargets = await getProviderTargets(provider.providerId);
|
||||
@@ -319,6 +320,7 @@ export async function proxyAiGatewayToSiteTarget(
|
||||
resourceId: ctx.resourceId,
|
||||
siteResourceId: ctx.siteResourceId,
|
||||
requestUserId: requestUser?.userId ?? null,
|
||||
virtualApiKeyId: ctx.virtualApiKeyId,
|
||||
budgets: ctx.budgets
|
||||
});
|
||||
}
|
||||
|
||||
@@ -58,7 +58,8 @@ export const aiUsageAnalyticsFiltersQuery = z.object({
|
||||
.transform(Number)
|
||||
.pipe(z.int().positive())
|
||||
.optional(),
|
||||
userId: z.string().optional()
|
||||
userId: z.string().optional(),
|
||||
virtualApiKeyId: z.string().optional()
|
||||
});
|
||||
|
||||
export const aiUsageAnalyticsParams = z.object({
|
||||
@@ -114,6 +115,9 @@ export function buildAiUsageWhere(
|
||||
)
|
||||
: undefined,
|
||||
data.userId ? eq(aiUsageRecords.userId, data.userId) : undefined,
|
||||
data.virtualApiKeyId
|
||||
? eq(aiUsageRecords.virtualApiKeyId, data.virtualApiKeyId)
|
||||
: undefined,
|
||||
roleUserIds ? inArray(aiUsageRecords.userId, roleUserIds) : undefined
|
||||
);
|
||||
}
|
||||
|
||||
@@ -8,3 +8,4 @@ export * from "./queryAiUsageOverview";
|
||||
export * from "./queryAiUsageProviders";
|
||||
export * from "./queryAiUsageResources";
|
||||
export * from "./queryAiUsageUsersRoles";
|
||||
export * from "./queryAiUsageVirtualApiKeys";
|
||||
|
||||
@@ -6,6 +6,7 @@ import {
|
||||
resources,
|
||||
siteResources,
|
||||
users,
|
||||
virtualApiKeys,
|
||||
db,
|
||||
primaryDb
|
||||
} from "@server/db";
|
||||
@@ -67,6 +68,7 @@ export const queryAiSessionLogsQuery = z.strictObject({
|
||||
.pipe(z.int().positive())
|
||||
.optional(),
|
||||
actor: z.string().optional(),
|
||||
virtualApiKeyId: z.string().optional(),
|
||||
model: z.string().optional(),
|
||||
isStream: z
|
||||
.union([z.boolean(), z.string()])
|
||||
@@ -117,7 +119,9 @@ function getWhere(data: Q) {
|
||||
data.providerId
|
||||
? eq(aiSessionLog.providerId, data.providerId)
|
||||
: undefined,
|
||||
data.capability ? eq(aiSessionLog.capability, data.capability) : undefined,
|
||||
data.capability
|
||||
? eq(aiSessionLog.capability, data.capability)
|
||||
: undefined,
|
||||
data.resourceId
|
||||
? or(
|
||||
eq(aiSessionLog.resourceId, data.resourceId),
|
||||
@@ -125,9 +129,10 @@ function getWhere(data: Q) {
|
||||
)
|
||||
: undefined,
|
||||
data.actor ? eq(aiSessionLog.userId, data.actor) : undefined,
|
||||
data.model
|
||||
? eq(aiSessionLog.requestedModel, data.model)
|
||||
data.virtualApiKeyId
|
||||
? eq(aiSessionLog.virtualApiKeyId, data.virtualApiKeyId)
|
||||
: undefined,
|
||||
data.model ? eq(aiSessionLog.requestedModel, data.model) : undefined,
|
||||
data.isStream !== undefined
|
||||
? eq(aiSessionLog.isStream, data.isStream)
|
||||
: undefined
|
||||
@@ -145,6 +150,7 @@ export function queryAiSession(data: Q) {
|
||||
resourceId: aiSessionLog.resourceId,
|
||||
siteResourceId: aiSessionLog.siteResourceId,
|
||||
userId: aiSessionLog.userId,
|
||||
virtualApiKeyId: aiSessionLog.virtualApiKeyId,
|
||||
requestedModel: aiSessionLog.requestedModel,
|
||||
isStream: aiSessionLog.isStream,
|
||||
requestBody: aiSessionLog.requestBody,
|
||||
@@ -182,6 +188,14 @@ async function enrichWithDetails(
|
||||
)
|
||||
];
|
||||
|
||||
const virtualApiKeyIds = [
|
||||
...new Set(
|
||||
logs
|
||||
.map((log) => log.virtualApiKeyId)
|
||||
.filter((id): id is string => id !== null && id !== undefined)
|
||||
)
|
||||
];
|
||||
|
||||
const providerMap = new Map<
|
||||
number,
|
||||
{ name: string | null; type: string | null }
|
||||
@@ -254,6 +268,28 @@ async function enrichWithDetails(
|
||||
}
|
||||
}
|
||||
|
||||
const virtualApiKeyMap = new Map<
|
||||
string,
|
||||
{ name: string | null; lastChars: string }
|
||||
>();
|
||||
if (virtualApiKeyIds.length > 0) {
|
||||
const virtualApiKeyDetails = await db
|
||||
.select({
|
||||
virtualApiKeyId: virtualApiKeys.virtualApiKeyId,
|
||||
name: virtualApiKeys.name,
|
||||
lastChars: virtualApiKeys.lastChars
|
||||
})
|
||||
.from(virtualApiKeys)
|
||||
.where(inArray(virtualApiKeys.virtualApiKeyId, virtualApiKeyIds));
|
||||
|
||||
for (const k of virtualApiKeyDetails) {
|
||||
virtualApiKeyMap.set(k.virtualApiKeyId, {
|
||||
name: k.name,
|
||||
lastChars: k.lastChars
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
const usageMap = new Map<
|
||||
string,
|
||||
{
|
||||
@@ -328,6 +364,12 @@ async function enrichWithDetails(
|
||||
resourceName,
|
||||
resourceNiceId,
|
||||
userEmail: log.userId ? (userMap.get(log.userId) ?? null) : null,
|
||||
virtualApiKeyName: log.virtualApiKeyId
|
||||
? (virtualApiKeyMap.get(log.virtualApiKeyId)?.name ?? null)
|
||||
: null,
|
||||
virtualApiKeyLastChars: log.virtualApiKeyId
|
||||
? (virtualApiKeyMap.get(log.virtualApiKeyId)?.lastChars ?? null)
|
||||
: null,
|
||||
usage: usageMap.get(log.sessionId) ?? null
|
||||
};
|
||||
});
|
||||
@@ -385,7 +427,8 @@ async function queryUniqueFilterAttributes(
|
||||
uniqueUsers,
|
||||
uniqueResources,
|
||||
uniqueSiteResources,
|
||||
uniqueModels
|
||||
uniqueModels,
|
||||
uniqueVirtualApiKeys
|
||||
] = await Promise.all([
|
||||
logsDb
|
||||
.selectDistinct({ id: aiSessionLog.providerId })
|
||||
@@ -411,6 +454,11 @@ async function queryUniqueFilterAttributes(
|
||||
.selectDistinct({ model: aiSessionLog.requestedModel })
|
||||
.from(aiSessionLog)
|
||||
.where(baseConditions)
|
||||
.limit(DISTINCT_LIMIT + 1),
|
||||
logsDb
|
||||
.selectDistinct({ id: aiSessionLog.virtualApiKeyId })
|
||||
.from(aiSessionLog)
|
||||
.where(baseConditions)
|
||||
.limit(DISTINCT_LIMIT + 1)
|
||||
]);
|
||||
|
||||
@@ -498,10 +546,37 @@ async function queryUniqueFilterAttributes(
|
||||
];
|
||||
}
|
||||
|
||||
const virtualApiKeyIds = uniqueVirtualApiKeys
|
||||
.map((row) => row.id)
|
||||
.filter((id): id is string => id !== null);
|
||||
|
||||
let virtualApiKeyList: Array<{
|
||||
id: string;
|
||||
name: string | null;
|
||||
lastChars: string | null;
|
||||
}> = [];
|
||||
if (virtualApiKeyIds.length > 0) {
|
||||
const virtualApiKeyDetails = await primaryDb
|
||||
.select({
|
||||
virtualApiKeyId: virtualApiKeys.virtualApiKeyId,
|
||||
name: virtualApiKeys.name,
|
||||
lastChars: virtualApiKeys.lastChars
|
||||
})
|
||||
.from(virtualApiKeys)
|
||||
.where(inArray(virtualApiKeys.virtualApiKeyId, virtualApiKeyIds));
|
||||
|
||||
virtualApiKeyList = virtualApiKeyDetails.map((k) => ({
|
||||
id: k.virtualApiKeyId,
|
||||
name: k.name,
|
||||
lastChars: k.lastChars
|
||||
}));
|
||||
}
|
||||
|
||||
return {
|
||||
providers: sortNamedFilterOptions(providers),
|
||||
resources: sortNamedFilterOptions(resourcesWithNames),
|
||||
users: userList,
|
||||
virtualApiKeys: virtualApiKeyList,
|
||||
models: models.sort()
|
||||
};
|
||||
}
|
||||
|
||||
@@ -6,7 +6,8 @@ import {
|
||||
siteResources,
|
||||
users,
|
||||
roles,
|
||||
userOrgRoles
|
||||
userOrgRoles,
|
||||
virtualApiKeys
|
||||
} from "@server/db";
|
||||
import { registry } from "@server/openApi";
|
||||
import { NextFunction } from "express";
|
||||
@@ -43,8 +44,9 @@ const queryAiUsageFilterOptionsParams = z.object({
|
||||
orgId: z.string()
|
||||
});
|
||||
|
||||
const queryAiUsageFilterOptionsCombined =
|
||||
queryAiUsageFilterOptionsQuery.merge(queryAiUsageFilterOptionsParams);
|
||||
const queryAiUsageFilterOptionsCombined = queryAiUsageFilterOptionsQuery.merge(
|
||||
queryAiUsageFilterOptionsParams
|
||||
);
|
||||
type Q = z.infer<typeof queryAiUsageFilterOptionsCombined>;
|
||||
|
||||
function sortNamedFilterOptions<T extends { id: number; name: string | null }>(
|
||||
@@ -73,7 +75,8 @@ async function query(data: Q) {
|
||||
uniqueModels,
|
||||
uniqueResources,
|
||||
uniqueSiteResources,
|
||||
uniqueUsers
|
||||
uniqueUsers,
|
||||
uniqueVirtualApiKeys
|
||||
] = await Promise.all([
|
||||
db
|
||||
.selectDistinct({ id: aiUsageRecords.providerId })
|
||||
@@ -105,6 +108,13 @@ async function query(data: Q) {
|
||||
.selectDistinct({ userId: aiUsageRecords.userId })
|
||||
.from(aiUsageRecords)
|
||||
.where(and(baseConditions, not(isNull(aiUsageRecords.userId))))
|
||||
.limit(DISTINCT_LIMIT + 1),
|
||||
db
|
||||
.selectDistinct({ id: aiUsageRecords.virtualApiKeyId })
|
||||
.from(aiUsageRecords)
|
||||
.where(
|
||||
and(baseConditions, not(isNull(aiUsageRecords.virtualApiKeyId)))
|
||||
)
|
||||
.limit(DISTINCT_LIMIT + 1)
|
||||
]);
|
||||
|
||||
@@ -120,11 +130,17 @@ async function query(data: Q) {
|
||||
let providers: Array<{ id: number; name: string | null }> = [];
|
||||
if (providerIds.length > 0) {
|
||||
const providerDetails = await db
|
||||
.select({ providerId: aiProviders.providerId, name: aiProviders.name })
|
||||
.select({
|
||||
providerId: aiProviders.providerId,
|
||||
name: aiProviders.name
|
||||
})
|
||||
.from(aiProviders)
|
||||
.where(inArray(aiProviders.providerId, providerIds));
|
||||
|
||||
providers = providerDetails.map((p) => ({ id: p.providerId, name: p.name }));
|
||||
providers = providerDetails.map((p) => ({
|
||||
id: p.providerId,
|
||||
name: p.name
|
||||
}));
|
||||
}
|
||||
|
||||
const resourceIds = uniqueResources
|
||||
@@ -155,7 +171,10 @@ async function query(data: Q) {
|
||||
.where(inArray(siteResources.siteResourceId, siteResourceIds));
|
||||
|
||||
resourcesWithNames = resourcesWithNames.concat(
|
||||
siteResourceDetails.map((r) => ({ id: r.siteResourceId, name: r.name }))
|
||||
siteResourceDetails.map((r) => ({
|
||||
id: r.siteResourceId,
|
||||
name: r.name
|
||||
}))
|
||||
);
|
||||
}
|
||||
|
||||
@@ -190,11 +209,37 @@ async function query(data: Q) {
|
||||
roleList = [...roleMap.entries()].map(([id, name]) => ({ id, name }));
|
||||
}
|
||||
|
||||
const virtualApiKeyIds = uniqueVirtualApiKeys
|
||||
.map((row) => row.id)
|
||||
.filter((id): id is string => id !== null);
|
||||
|
||||
let virtualApiKeyList: Array<{
|
||||
id: string;
|
||||
name: string | null;
|
||||
lastChars: string;
|
||||
}> = [];
|
||||
if (virtualApiKeyIds.length > 0) {
|
||||
const virtualApiKeyDetails = await db
|
||||
.select({
|
||||
virtualApiKeyId: virtualApiKeys.virtualApiKeyId,
|
||||
name: virtualApiKeys.name,
|
||||
lastChars: virtualApiKeys.lastChars
|
||||
})
|
||||
.from(virtualApiKeys)
|
||||
.where(inArray(virtualApiKeys.virtualApiKeyId, virtualApiKeyIds));
|
||||
virtualApiKeyList = virtualApiKeyDetails.map((k) => ({
|
||||
id: k.virtualApiKeyId,
|
||||
name: k.name,
|
||||
lastChars: k.lastChars
|
||||
}));
|
||||
}
|
||||
|
||||
return {
|
||||
providers: sortNamedFilterOptions(providers),
|
||||
resources: sortNamedFilterOptions(resourcesWithNames),
|
||||
roles: sortNamedFilterOptions(roleList),
|
||||
users: userList,
|
||||
virtualApiKeys: virtualApiKeyList,
|
||||
models
|
||||
};
|
||||
}
|
||||
@@ -227,7 +272,9 @@ registry.registerPath({
|
||||
}
|
||||
});
|
||||
|
||||
export type QueryAiUsageFilterOptionsResponse = Awaited<ReturnType<typeof query>>;
|
||||
export type QueryAiUsageFilterOptionsResponse = Awaited<
|
||||
ReturnType<typeof query>
|
||||
>;
|
||||
|
||||
export async function queryAiUsageFilterOptions(
|
||||
req: Request,
|
||||
@@ -238,7 +285,10 @@ export async function queryAiUsageFilterOptions(
|
||||
const parsedQuery = queryAiUsageFilterOptionsQuery.safeParse(req.query);
|
||||
if (!parsedQuery.success) {
|
||||
return next(
|
||||
createHttpError(HttpCode.BAD_REQUEST, fromError(parsedQuery.error))
|
||||
createHttpError(
|
||||
HttpCode.BAD_REQUEST,
|
||||
fromError(parsedQuery.error)
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
@@ -247,7 +297,10 @@ export async function queryAiUsageFilterOptions(
|
||||
);
|
||||
if (!parsedParams.success) {
|
||||
return next(
|
||||
createHttpError(HttpCode.BAD_REQUEST, fromError(parsedParams.error))
|
||||
createHttpError(
|
||||
HttpCode.BAD_REQUEST,
|
||||
fromError(parsedParams.error)
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
import { db, aiUsageRecords, virtualApiKeys } from "@server/db";
|
||||
import { registry } from "@server/openApi";
|
||||
import { NextFunction } from "express";
|
||||
import { Request, Response } from "express";
|
||||
import { count, desc, inArray, sql } from "drizzle-orm";
|
||||
import { OpenAPITags } from "@server/openApi";
|
||||
import { z } from "zod";
|
||||
import createHttpError from "http-errors";
|
||||
import HttpCode from "@server/types/HttpCode";
|
||||
import { fromError } from "zod-validation-error";
|
||||
import response from "@server/lib/response";
|
||||
import logger from "@server/logger";
|
||||
import {
|
||||
aiUsageAnalyticsFiltersQuery,
|
||||
aiUsageAnalyticsParams,
|
||||
buildAiUsageWhere,
|
||||
resolveRoleUserIds,
|
||||
dayBucketExpr,
|
||||
pickTopNKeys,
|
||||
bucketTopNPerDay,
|
||||
DISTINCT_LIMIT,
|
||||
type AiUsageAnalyticsQuery
|
||||
} from "./aiUsageAnalyticsShared";
|
||||
|
||||
type Q = AiUsageAnalyticsQuery;
|
||||
|
||||
const UNKNOWN_VIRTUAL_API_KEY_KEY = "unknown";
|
||||
|
||||
async function query(data: Q) {
|
||||
const roleUserIds = await resolveRoleUserIds(data.orgId, data.roleId);
|
||||
const baseConditions = buildAiUsageWhere(data, roleUserIds);
|
||||
const dayExpr = dayBucketExpr();
|
||||
|
||||
const virtualApiKeyByDay = await db
|
||||
.select({
|
||||
day: dayExpr.as("day"),
|
||||
virtualApiKeyId: aiUsageRecords.virtualApiKeyId,
|
||||
cost: sql<number>`COALESCE(SUM(${aiUsageRecords.costUsd}), 0)`,
|
||||
tokens: sql<number>`COALESCE(SUM(${aiUsageRecords.totalTokens}), 0)`
|
||||
})
|
||||
.from(aiUsageRecords)
|
||||
.where(baseConditions)
|
||||
.groupBy(dayExpr, aiUsageRecords.virtualApiKeyId)
|
||||
.orderBy(dayExpr);
|
||||
|
||||
const virtualApiKeyTotalsRaw = await db
|
||||
.select({
|
||||
virtualApiKeyId: aiUsageRecords.virtualApiKeyId,
|
||||
requests: count(),
|
||||
totalTokens: sql<number>`COALESCE(SUM(${aiUsageRecords.totalTokens}), 0)`,
|
||||
costUsd: sql<number>`COALESCE(SUM(${aiUsageRecords.costUsd}), 0)`
|
||||
})
|
||||
.from(aiUsageRecords)
|
||||
.where(baseConditions)
|
||||
.groupBy(aiUsageRecords.virtualApiKeyId)
|
||||
.orderBy(desc(sql`COALESCE(SUM(${aiUsageRecords.costUsd}), 0)`))
|
||||
.limit(DISTINCT_LIMIT + 1);
|
||||
|
||||
if (virtualApiKeyTotalsRaw.length > DISTINCT_LIMIT) {
|
||||
throw createHttpError(
|
||||
HttpCode.BAD_REQUEST,
|
||||
"Too many distinct virtual API keys. Please narrow your query."
|
||||
);
|
||||
}
|
||||
|
||||
const virtualApiKeyIds = virtualApiKeyTotalsRaw
|
||||
.map((r) => r.virtualApiKeyId)
|
||||
.filter((id): id is string => id !== null);
|
||||
|
||||
const detailsMap = new Map<
|
||||
string,
|
||||
{ name: string | null; lastChars: string; kind: "user" | "manual" }
|
||||
>();
|
||||
if (virtualApiKeyIds.length > 0) {
|
||||
const details = await db
|
||||
.select({
|
||||
virtualApiKeyId: virtualApiKeys.virtualApiKeyId,
|
||||
name: virtualApiKeys.name,
|
||||
lastChars: virtualApiKeys.lastChars,
|
||||
kind: virtualApiKeys.kind
|
||||
})
|
||||
.from(virtualApiKeys)
|
||||
.where(inArray(virtualApiKeys.virtualApiKeyId, virtualApiKeyIds));
|
||||
for (const k of details) {
|
||||
detailsMap.set(k.virtualApiKeyId, {
|
||||
name: k.name,
|
||||
lastChars: k.lastChars,
|
||||
kind: k.kind
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
const topVirtualApiKeys = virtualApiKeyTotalsRaw.map((r) => {
|
||||
const details = r.virtualApiKeyId
|
||||
? detailsMap.get(r.virtualApiKeyId)
|
||||
: undefined;
|
||||
return {
|
||||
virtualApiKeyId: r.virtualApiKeyId,
|
||||
name: details?.name ?? null,
|
||||
lastChars: details?.lastChars ?? null,
|
||||
kind: details?.kind ?? null,
|
||||
requests: r.requests,
|
||||
totalTokens: r.totalTokens,
|
||||
costUsd: r.costUsd
|
||||
};
|
||||
});
|
||||
|
||||
const virtualApiKeyCostTotals = new Map<string, number>();
|
||||
const virtualApiKeyTokenTotals = new Map<string, number>();
|
||||
for (const row of virtualApiKeyByDay) {
|
||||
const key = row.virtualApiKeyId ?? UNKNOWN_VIRTUAL_API_KEY_KEY;
|
||||
virtualApiKeyCostTotals.set(
|
||||
key,
|
||||
(virtualApiKeyCostTotals.get(key) ?? 0) + row.cost
|
||||
);
|
||||
virtualApiKeyTokenTotals.set(
|
||||
key,
|
||||
(virtualApiKeyTokenTotals.get(key) ?? 0) + row.tokens
|
||||
);
|
||||
}
|
||||
const topVirtualApiKeysByCost = pickTopNKeys(virtualApiKeyCostTotals);
|
||||
const topVirtualApiKeysByTokens = pickTopNKeys(virtualApiKeyTokenTotals);
|
||||
|
||||
const virtualApiKeyCostPerDay = bucketTopNPerDay(
|
||||
virtualApiKeyByDay.map((r) => ({
|
||||
day: r.day,
|
||||
key: r.virtualApiKeyId ?? UNKNOWN_VIRTUAL_API_KEY_KEY,
|
||||
value: r.cost
|
||||
})),
|
||||
topVirtualApiKeysByCost
|
||||
);
|
||||
const virtualApiKeyTokensPerDay = bucketTopNPerDay(
|
||||
virtualApiKeyByDay.map((r) => ({
|
||||
day: r.day,
|
||||
key: r.virtualApiKeyId ?? UNKNOWN_VIRTUAL_API_KEY_KEY,
|
||||
value: r.tokens
|
||||
})),
|
||||
topVirtualApiKeysByTokens
|
||||
);
|
||||
|
||||
return {
|
||||
topVirtualApiKeys,
|
||||
virtualApiKeyCostPerDay,
|
||||
virtualApiKeyTokensPerDay
|
||||
};
|
||||
}
|
||||
|
||||
registry.registerPath({
|
||||
method: "get",
|
||||
path: "/org/{orgId}/logs/ai/usage/virtual-api-keys",
|
||||
description:
|
||||
"Query the AI usage analytics virtual API key breakdown for an organization",
|
||||
tags: [OpenAPITags.Logs],
|
||||
request: {
|
||||
query: aiUsageAnalyticsFiltersQuery,
|
||||
params: aiUsageAnalyticsParams
|
||||
},
|
||||
responses: {
|
||||
200: {
|
||||
description: "Successful response",
|
||||
content: {
|
||||
"application/json": {
|
||||
schema: z.object({
|
||||
data: z.record(z.string(), z.any()).nullable(),
|
||||
success: z.boolean(),
|
||||
error: z.boolean(),
|
||||
message: z.string(),
|
||||
status: z.number()
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
export type QueryAiUsageVirtualApiKeysResponse = Awaited<
|
||||
ReturnType<typeof query>
|
||||
>;
|
||||
|
||||
export async function queryAiUsageVirtualApiKeys(
|
||||
req: Request,
|
||||
res: Response,
|
||||
next: NextFunction
|
||||
): Promise<any> {
|
||||
try {
|
||||
const parsedQuery = aiUsageAnalyticsFiltersQuery.safeParse(req.query);
|
||||
if (!parsedQuery.success) {
|
||||
return next(
|
||||
createHttpError(
|
||||
HttpCode.BAD_REQUEST,
|
||||
fromError(parsedQuery.error)
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
const parsedParams = aiUsageAnalyticsParams.safeParse(req.params);
|
||||
if (!parsedParams.success) {
|
||||
return next(
|
||||
createHttpError(
|
||||
HttpCode.BAD_REQUEST,
|
||||
fromError(parsedParams.error)
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
const data = await query({ ...parsedQuery.data, ...parsedParams.data });
|
||||
|
||||
return response<QueryAiUsageVirtualApiKeysResponse>(res, {
|
||||
data,
|
||||
success: true,
|
||||
error: false,
|
||||
message:
|
||||
"AI usage virtual API key breakdown retrieved successfully",
|
||||
status: HttpCode.OK
|
||||
});
|
||||
} catch (error) {
|
||||
logger.error(error);
|
||||
return next(
|
||||
createHttpError(HttpCode.INTERNAL_SERVER_ERROR, "An error occurred")
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -110,6 +110,9 @@ export type QueryAiSessionLogResponse = {
|
||||
resourceType: "public" | "site" | null;
|
||||
userId: string | null;
|
||||
userEmail: string | null;
|
||||
virtualApiKeyId: string | null;
|
||||
virtualApiKeyName: string | null;
|
||||
virtualApiKeyLastChars: string | null;
|
||||
requestedModel: string | null;
|
||||
isStream: boolean;
|
||||
requestBody: string | null;
|
||||
@@ -148,6 +151,11 @@ export type QueryAiSessionLogResponse = {
|
||||
id: string;
|
||||
email: string | null;
|
||||
}[];
|
||||
virtualApiKeys: {
|
||||
id: string;
|
||||
name: string | null;
|
||||
lastChars: string | null;
|
||||
}[];
|
||||
models: string[];
|
||||
};
|
||||
};
|
||||
|
||||
@@ -96,6 +96,10 @@ export type VerifyUserResponse = {
|
||||
userData?: BasicUserData;
|
||||
pangolinVersion?: string;
|
||||
dontStripSession?: boolean;
|
||||
// Set independently of userData so a manual virtual API key with no
|
||||
// associated user still gets attributed to the key that authenticated
|
||||
// the request (see the mode === "inference" branch below).
|
||||
virtualApiKeyId?: string;
|
||||
};
|
||||
|
||||
export async function verifyResourceSession(
|
||||
@@ -401,7 +405,12 @@ export async function verifyResourceSession(
|
||||
parsedBody.data
|
||||
);
|
||||
|
||||
return allowed(res, vakUserData, dontStripSession);
|
||||
return allowed(
|
||||
res,
|
||||
vakUserData,
|
||||
dontStripSession,
|
||||
key.virtualApiKeyId
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -987,16 +996,20 @@ async function notAllowed(
|
||||
function allowed(
|
||||
res: Response,
|
||||
userData?: BasicUserData,
|
||||
dontStripSession?: boolean
|
||||
dontStripSession?: boolean,
|
||||
virtualApiKeyId?: string
|
||||
) {
|
||||
const baseData =
|
||||
userData !== undefined && userData !== null
|
||||
? { valid: true, ...userData, pangolinVersion: APP_VERSION }
|
||||
: { valid: true, pangolinVersion: APP_VERSION };
|
||||
const withVirtualApiKey = virtualApiKeyId
|
||||
? { ...baseData, virtualApiKeyId }
|
||||
: baseData;
|
||||
const data = {
|
||||
data: dontStripSession
|
||||
? { ...baseData, dontStripSession: true }
|
||||
: baseData,
|
||||
? { ...withVirtualApiKey, dontStripSession: true }
|
||||
: withVirtualApiKey,
|
||||
success: true,
|
||||
error: false,
|
||||
message: "Access allowed",
|
||||
|
||||
@@ -1537,6 +1537,13 @@ authenticated.get(
|
||||
logs.queryAiUsageUsersRoles
|
||||
);
|
||||
|
||||
authenticated.get(
|
||||
"/org/:orgId/logs/ai/usage/virtual-api-keys",
|
||||
verifyOrgAccess,
|
||||
verifyUserHasAction(ActionsEnum.viewLogs),
|
||||
logs.queryAiUsageVirtualApiKeys
|
||||
);
|
||||
|
||||
authenticated.get(
|
||||
"/org/:orgId/blueprints",
|
||||
verifyOrgAccess,
|
||||
|
||||
@@ -1580,6 +1580,13 @@ authenticated.get(
|
||||
logs.queryAiUsageUsersRoles
|
||||
);
|
||||
|
||||
authenticated.get(
|
||||
"/org/:orgId/logs/ai/usage/virtual-api-keys",
|
||||
verifyApiKeyOrgAccess,
|
||||
verifyApiKeyHasAction(ActionsEnum.viewLogs),
|
||||
logs.queryAiUsageVirtualApiKeys
|
||||
);
|
||||
|
||||
authenticated.get(
|
||||
"/org/:orgId/logs/analytics",
|
||||
verifyApiKeyOrgAccess,
|
||||
|
||||
Reference in New Issue
Block a user