Menu
popagent
publicLatest change ac87cc5fc89df21b21c6d2b8938cb47d98ed176c - Expose Popagent process memory metrics by AkurAI Build
import { TraceStatus, type ObservabilityStorage } from "@mastra/core/storage";
import type {
ObservabilityHealthResponse,
ObservabilityLogListResponse,
ObservabilityOverviewResponse,
ObservabilityTraceListResponse,
ObservabilityTraceResponse,
ObservabilityTraceSummary,
} from "./api-types";
import { storage } from "./storage";
const DAY_MS = 86_400_000;
const MAX_RANGE_MS = 30 * DAY_MS;
const MAX_PAGE_SIZE = 100;
export type ObservabilityQuery = {
from: Date;
to: Date;
page: number;
perPage: number;
};
export function observabilityFilters<T extends Record<string, unknown>>(filters: T): T {
return filters;
}
export function publicSpanPayload(_payload: unknown): null {
return null;
}
export function currentProcessMemory() {
const { rss, heapUsed } = process.memoryUsage();
return { processRssBytes: rss, processHeapUsedBytes: heapUsed };
}
export function parseObservabilityQuery(url: string, now = new Date()): ObservabilityQuery {
const params = new URL(url).searchParams;
const to = params.get("to") ? new Date(params.get("to")!) : now;
const from = params.get("from") ? new Date(params.get("from")!) : new Date(to.getTime() - DAY_MS);
const page = Number(params.get("page") ?? 0);
const perPage = Number(params.get("perPage") ?? 25);
if (!Number.isFinite(from.getTime()) || !Number.isFinite(to.getTime()) || from >= to) {
throw new Error("Invalid observability time range");
}
if (to.getTime() - from.getTime() > MAX_RANGE_MS) {
throw new Error("Observability time range cannot exceed 30 days");
}
if (!Number.isInteger(page) || page < 0 || !Number.isInteger(perPage) || perPage < 1 || perPage > MAX_PAGE_SIZE) {
throw new Error("Invalid observability pagination");
}
return { from, to, page, perPage };
}
async function getStore(): Promise<ObservabilityStorage> {
const store = await storage.getStore("observability");
if (!store) throw new Error("Observability storage is unavailable");
return store;
}
function summary(span: {
traceId: string;
name: string;
entityName?: string | null;
status?: "success" | "error" | "running";
startedAt: Date;
endedAt?: Date | null;
error?: unknown;
tags?: string[] | null;
metadata?: Record<string, unknown> | null;
}): ObservabilityTraceSummary {
const endedAt = span.endedAt ?? null;
return {
traceId: span.traceId,
name: span.name,
entityName: span.entityName ?? null,
status: span.status ?? (span.error ? "error" : endedAt ? "success" : "running"),
startedAt: span.startedAt.toISOString(),
endedAt: endedAt?.toISOString() ?? null,
durationMs: endedAt ? endedAt.getTime() - span.startedAt.getTime() : null,
tags: span.tags ?? [],
metadata: span.metadata ?? {},
};
}
export async function listObservabilityTraces(query: ObservabilityQuery): Promise<ObservabilityTraceListResponse> {
const result = await (await getStore()).listTraces({
filters: observabilityFilters({ startedAt: { start: query.from, end: query.to } }),
pagination: { page: query.page, perPage: query.perPage },
orderBy: { field: "startedAt", direction: "DESC" },
});
return {
traces: result.spans.map(summary),
page: result.pagination?.page ?? query.page,
perPage: Number(result.pagination?.perPage ?? query.perPage),
total: result.pagination?.total ?? result.spans.length,
hasMore: result.pagination?.hasMore ?? false,
};
}
export async function getObservabilityTrace(traceId: string): Promise<ObservabilityTraceResponse | null> {
const result = await (await getStore()).getTrace({ traceId });
if (!result) return null;
return {
traceId,
spans: result.spans.map((span) => ({
...summary(span),
spanId: span.spanId,
parentSpanId: span.parentSpanId ?? null,
spanType: span.spanType,
input: publicSpanPayload(span.input),
output: publicSpanPayload(span.output),
error: span.error ?? null,
})),
};
}
export async function listObservabilityLogs(query: ObservabilityQuery): Promise<ObservabilityLogListResponse> {
const result = await (await getStore()).listLogs({
filters: observabilityFilters({
timestamp: { start: query.from, end: query.to },
level: ["warn", "error", "fatal"] as ("warn" | "error" | "fatal")[],
}),
pagination: { page: query.page, perPage: query.perPage },
orderBy: { field: "timestamp", direction: "DESC" },
});
return {
logs: result.logs.map((log) => ({
id: log.logId ?? null,
timestamp: log.timestamp.toISOString(),
level: log.level,
message: log.message,
traceId: log.traceId ?? null,
spanId: log.spanId ?? null,
entityName: log.entityName ?? null,
data: log.data ?? {},
})),
page: result.pagination?.page ?? query.page,
perPage: Number(result.pagination?.perPage ?? query.perPage),
total: result.pagination?.total ?? result.logs.length,
hasMore: result.pagination?.hasMore ?? false,
};
}
export async function getObservabilityOverview(query: ObservabilityQuery): Promise<ObservabilityOverviewResponse> {
const store = await getStore();
const filters = { timestamp: { start: query.from, end: query.to } };
const [traces, errors, input, output] = await Promise.all([
store.listTraces({
filters: observabilityFilters({ startedAt: filters.timestamp }),
pagination: { page: 0, perPage: 100 },
orderBy: { field: "startedAt", direction: "DESC" },
}),
store.listTraces({
filters: observabilityFilters({ startedAt: filters.timestamp, status: TraceStatus.ERROR }),
pagination: { page: 0, perPage: 1 },
}),
store.getMetricAggregate({ name: ["mastra_model_total_input_tokens"], aggregation: "sum", filters: observabilityFilters(filters) }),
store.getMetricAggregate({ name: ["mastra_model_total_output_tokens"], aggregation: "sum", filters: observabilityFilters(filters) }),
]);
const durations = traces.spans
.map((span) => span.endedAt ? span.endedAt.getTime() - span.startedAt.getTime() : null)
.filter((value): value is number => value !== null)
.sort((a, b) => a - b);
const runs = traces.pagination?.total ?? traces.spans.length;
const errorCount = errors.pagination?.total ?? errors.spans.length;
return {
from: query.from.toISOString(),
to: query.to.toISOString(),
runs,
errors: errorCount,
errorRate: runs ? errorCount / runs : 0,
p95LatencyMs: durations.length ? durations[Math.ceil(durations.length * 0.95) - 1]! : null,
inputTokens: input.value ?? 0,
outputTokens: output.value ?? 0,
...currentProcessMemory(),
};
}
export async function observabilityHealth(): Promise<ObservabilityHealthResponse> {
const checkedAt = new Date().toISOString();
try {
await (await getStore()).getMetricNames({ prefix: "mastra_", limit: 1 });
return { status: "ok", checkedAt };
} catch (error) {
return { status: "error", checkedAt, error: error instanceof Error ? error.message : String(error) };
}
}