AkurAI Build
Menu

popagent

public

Latest change 7cecf6a89b6f39dae8c5678f5b369014199aeb3c - Add self-hosted Mastra observability by AkurAI Build

import { Cron } from "croner";
import { hookAudit } from "./hook-audit";
import { longTermMemory } from "./long-term-memory";
import { appLogger } from "./observability";
import { storage } from "./storage";

type MaintenanceStores = {
  memory: Pick<typeof longTermMemory, "pruneEpisodes">;
  audit: Pick<typeof hookAudit, "cleanupAll">;
  observability?: { prune(): Promise<unknown> };
};

function retentionDays(name: string, fallback: number): number {
  const value = Number(process.env[name] ?? fallback);
  if (!Number.isInteger(value) || value < 1 || value > 365) throw new Error(`${name} must be an integer from 1 to 365`);
  return value;
}

export const observabilityRetention = {
  async prune() {
    const store = await storage.getStore("observability");
    if (!store) throw new Error("Observability storage is unavailable");
    const signalDays = retentionDays("POPAGENT_OBSERVABILITY_RETENTION_DAYS", 30);
    const feedbackDays = retentionDays("POPAGENT_OBSERVABILITY_FEEDBACK_RETENTION_DAYS", 90);
    return store.prune({
      spans: { maxAge: `${signalDays}d` },
      metrics: { maxAge: `${signalDays}d` },
      logs: { maxAge: `${signalDays}d` },
      scores: { maxAge: `${signalDays}d` },
      feedback: { maxAge: `${feedbackDays}d` },
    }, { maxBatches: 100 });
  },
};

export async function runMaintenanceOnce(stores: MaintenanceStores = {
  memory: longTermMemory,
  audit: hookAudit,
  observability: observabilityRetention,
}): Promise<void> {
  await Promise.all([
    stores.memory.pruneEpisodes(),
    stores.audit.cleanupAll(),
    stores.observability?.prune(),
  ]);
}

export function startMaintenance(): Cron {
  const job = new Cron("0 3 * * *", { timezone: Intl.DateTimeFormat().resolvedOptions().timeZone }, () => {
    void runMaintenanceOnce().catch((error) => appLogger().error("maintenance.failed", { error }));
  });
  void runMaintenanceOnce().catch((error) => appLogger().error("maintenance.failed", { error }));
  return job;
}