AkurAI Build
Menu

popagent

public

Latest change 6edc7e928c57dab23c0973ed3f7ee70b3ebb72d5 - Initial popagent baseline 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;
  model: string;
  workspaceId?: string;
  prompt: string;
  onIterationComplete?: (context: IterationCompleteContext) => Promise<void>;
};

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,
      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,
        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,
      model: task.model,
      workspaceId: (task as { workspaceId?: string }).workspaceId,
      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}`;
}