AkurAI Build
Menu

popagent

public

Latest change 7cecf6a89b6f39dae8c5678f5b369014199aeb3c - Add self-hosted Mastra observability 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 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: { startedAt: { start: query.from, end: query.to }, serviceName: "popagent" },
    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: span.input ?? null,
      output: span.output ?? null,
      error: span.error ?? null,
    })),
  };
}

export async function listObservabilityLogs(query: ObservabilityQuery): Promise<ObservabilityLogListResponse> {
  const result = await (await getStore()).listLogs({
    filters: {
      timestamp: { start: query.from, end: query.to },
      level: ["warn", "error", "fatal"],
      serviceName: "popagent",
    },
    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 }, serviceName: "popagent" };
  const [traces, errors, input, output] = await Promise.all([
    store.listTraces({
      filters: { startedAt: filters.timestamp, serviceName: "popagent" },
      pagination: { page: 0, perPage: 100 },
      orderBy: { field: "startedAt", direction: "DESC" },
    }),
    store.listTraces({
      filters: { startedAt: filters.timestamp, serviceName: "popagent", status: TraceStatus.ERROR },
      pagination: { page: 0, perPage: 1 },
    }),
    store.getMetricAggregate({ name: ["mastra_model_total_input_tokens"], aggregation: "sum", filters }),
    store.getMetricAggregate({ name: ["mastra_model_total_output_tokens"], aggregation: "sum", 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,
  };
}

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) };
  }
}