Menu
popagent
publicLatest change 7cecf6a89b6f39dae8c5678f5b369014199aeb3c - Add self-hosted Mastra observability by AkurAI Build
import type { Agent, IterationCompleteContext } from "@mastra/core/agent";
import { createAgentRequestContext } from "./agent-context";
import { createFinalResponseGuard, toolCallConcurrency } from "./agent-autonomy";
import type { AgentRuntimeSettingsStore } from "./agent-runtime-settings";
import {
type AgentActivity,
type AgentTask,
type AgentTaskProgress,
} from "./api-types";
import type { BrowserRuntime } from "./browser-settings";
import { createDelegationHooks, createToolHooks, HookLifecycleProcessor } from "./hook-lifecycle";
import { createHookEvent, type HookRuntime } from "./hooks";
import type { LongTermMemoryStore } from "./long-term-memory";
import type { MemorySettingsStore } from "./memory-settings";
import { resolveModel } from "./models";
import { RESOURCE_ID } from "./sessions";
type ExecutionInput = {
sessionId: string;
turnId: string;
traceId: string;
model: string;
workspaceId?: string;
executionSource?: "chat" | "task" | "schedule";
taskId?: string;
scheduleId?: string;
prompt: string;
onIterationComplete?: (context: IterationCompleteContext) => Promise<void>;
};
export function createTraceId(): string {
return crypto.randomUUID().replaceAll("-", "");
}
type MemoryRepository = Pick<LongTermMemoryStore, "formatRecall" | "retainEpisode">;
type MemoryConfiguration = Pick<MemorySettingsStore, "get">;
type BrowserController = Pick<BrowserRuntime, "apply">;
type RuntimeConfiguration = Pick<AgentRuntimeSettingsStore, "get">;
export class AgentExecutionRuntime {
constructor(
private readonly agent: Agent,
private readonly hooks: HookRuntime,
private readonly memories: MemoryRepository,
private readonly browser: BrowserController,
private readonly memorySettings: MemoryConfiguration,
private readonly runtimeSettings: RuntimeConfiguration,
) {}
async prepare(input: ExecutionInput) {
const promptHook = await this.hooks.dispatch(createHookEvent("UserPromptSubmit", {
sessionId: input.sessionId,
turnId: input.turnId,
model: input.model,
detail: { prompt: input.prompt },
}));
const promptWasReplaced = typeof promptHook.replacement === "string";
const prompt = promptWasReplaced ? promptHook.replacement as string : input.prompt;
await this.browser.apply(this.agent);
const recalled = await this.memories.formatRecall({
resourceId: RESOURCE_ID,
query: prompt,
});
const context = [...promptHook.additionalContext, ...(recalled ? [recalled] : [])]
.map((text) => ({
role: "user" as const,
content: `<memory-context>\n${text}\n</memory-context>`,
}));
const [compacting, runtimeSettings] = await Promise.all([
this.memorySettings.get(),
this.runtimeSettings.get(),
]);
const requestContext = createAgentRequestContext({
resourceId: RESOURCE_ID,
sessionId: input.sessionId,
turnId: input.turnId,
model: input.model,
workspaceId: input.workspaceId,
executionSource: input.executionSource,
taskId: input.taskId,
scheduleId: input.scheduleId,
runtimeSettings,
hookRuntime: this.hooks,
iterationObserver: input.onIterationComplete,
});
const execution = { sessionId: input.sessionId, turnId: input.turnId, model: input.model };
const observeIteration = input.onIterationComplete;
const reserveFinalResponse = createFinalResponseGuard(runtimeSettings);
const onIterationComplete = observeIteration
? async (iteration: IterationCompleteContext) => {
await observeIteration(iteration);
return reserveFinalResponse(iteration);
}
: reserveFinalResponse;
return {
prompt,
promptWasReplaced,
options: {
model: resolveModel(input.model),
maxSteps: runtimeSettings.supervisorMaxSteps,
toolCallConcurrency: toolCallConcurrency(runtimeSettings),
onIterationComplete,
requestContext,
tracingOptions: {
traceId: input.traceId,
tags: [input.executionSource ?? "chat"],
metadata: {
sessionId: input.sessionId,
turnId: input.turnId,
workspaceId: input.workspaceId,
model: input.model,
executionSource: input.executionSource ?? "chat",
...(input.taskId ? { taskId: input.taskId } : {}),
...(input.scheduleId ? { scheduleId: input.scheduleId } : {}),
},
},
context,
hooks: createToolHooks(this.hooks, execution),
delegation: createDelegationHooks(this.hooks, execution, runtimeSettings),
outputProcessors: [new HookLifecycleProcessor(this.hooks, execution)],
errorProcessors: [new HookLifecycleProcessor(this.hooks, execution)],
maxProcessorRetries: runtimeSettings.maxProcessorRetries,
memoryConfig: {
observationalMemory: compacting.autoCompact ? {
enabled: true as const,
model: resolveModel(input.model),
observation: {
messageTokens: compacting.observationTokens,
bufferActivation: 1 - compacting.recentMessagePercent / 100,
bufferOnIdle: compacting.bufferOnIdle,
},
reflection: { observationTokens: compacting.reflectionTokens },
retrieval: true,
} : { enabled: false as const },
},
onFinish: async (result: { text: string }) => {
if (!result.text.trim()) return;
await this.memories.retainEpisode({
resourceId: RESOURCE_ID,
sessionId: input.sessionId,
userText: prompt,
assistantText: result.text,
});
},
},
};
}
}
export function createAgentTaskExecutor(agent: Agent, execution: AgentExecutionRuntime) {
return async (
task: AgentTask,
signal: AbortSignal,
turnId: string,
reportProgress: (progress: AgentTaskProgress) => Promise<void>,
): Promise<string> => {
const prepared = await execution.prepare({
sessionId: task.sessionId ?? `task:${task.id}`,
turnId,
traceId: createTraceId(),
model: task.model,
workspaceId: (task as { workspaceId?: string }).workspaceId,
executionSource: task.scheduleId ? "schedule" : "task",
taskId: task.id,
scheduleId: task.scheduleId ?? undefined,
prompt: task.prompt,
onIterationComplete: async (iteration) => {
await reportProgress({
stepsCompleted: iteration.iteration,
progress: iterationActivity(iteration),
});
},
});
const result = await agent.stream(prepared.prompt, {
...prepared.options,
abortSignal: signal,
memory: task.sessionId ? { thread: task.sessionId, resource: RESOURCE_ID } : undefined,
});
return await result.text;
};
}
export function agentActivity(turnId: string, iteration: IterationCompleteContext): AgentActivity {
return {
turnId,
runId: iteration.runId,
agentId: iteration.agentId,
agentName: iteration.agentName ?? iteration.agentId,
iteration: iteration.iteration,
maxIterations: iteration.maxIterations ?? null,
isFinal: iteration.isFinal,
finishReason: iteration.finishReason,
text: iteration.text,
tools: [...new Set(iteration.toolCalls.map((call) => call.name))],
};
}
function iterationActivity(iteration: IterationCompleteContext): string {
const tools = [...new Set(iteration.toolCalls.map((call) => call.name))];
const activity = tools.length
? `Using ${tools.join(", ")}`
: iteration.text.trim()
? "Preparing response"
: "Reasoning";
return iteration.agentId === "popagent"
? activity
: `${iteration.agentName ?? iteration.agentId}: ${activity}`;
}