vv-agent 0.7.1

VectorVein agent runtime, SDK, CLI, tools, and workspace backends
Documentation
use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::sync::Arc;

use serde_json::Value;

use crate::config::load_memory_summary_defaults_from_file;
use crate::events::RunEvent;
use crate::memory::token_utils::count_messages_tokens;
use crate::memory::{
    provider::block_on_memory_future, MemoryManager, MemoryManagerConfig, MemoryProvider,
};
use crate::model::ModelProvider;
use crate::runtime::context::ExecutionContext;
use crate::types::{AgentTask, Message};

mod callbacks;
mod metadata;
mod session;
mod token_limits;

use callbacks::build_memory_summary_callback;
use metadata::{
    metadata_path, read_bool_metadata, read_f64_metadata, read_optional_string_metadata,
    read_string_metadata, read_string_set_metadata, read_u64_metadata, read_usize_metadata,
};
use session::build_session_memory;
use token_limits::resolve_runtime_model_token_limits;

pub(super) fn build_memory_manager(
    task: &AgentTask,
    workspace_path: PathBuf,
    memory_model_provider: Option<Arc<dyn ModelProvider>>,
    settings_file: Option<&Path>,
    default_backend: Option<&str>,
) -> MemoryManager {
    let workspace = task.use_workspace.then_some(workspace_path);
    let local_summary_defaults = settings_file
        .map(load_memory_summary_defaults_from_file)
        .unwrap_or_default();
    let summary_backend = read_optional_string_metadata(
        &task.metadata,
        &[
            "memory_summary_backend",
            "compress_memory_summary_backend",
            "memory_compress_backend",
        ],
    )
    .or(local_summary_defaults.backend)
    .or_else(|| default_backend.map(str::to_string));
    let summary_model = read_optional_string_metadata(
        &task.metadata,
        &[
            "memory_summary_model",
            "compress_memory_summary_model",
            "memory_compress_model",
        ],
    )
    .or(local_summary_defaults.model)
    .unwrap_or_else(|| task.model.clone());
    let has_memory_route = summary_backend
        .as_deref()
        .is_some_and(|value| !value.trim().is_empty())
        || read_optional_string_metadata(&task.metadata, &["session_memory_extraction_backend"])
            .is_some();
    let summary_callback = if has_memory_route {
        memory_model_provider.map(|provider| {
            build_memory_summary_callback(provider, summary_backend.clone(), summary_model.clone())
        })
    } else {
        None
    };
    let (resolved_context_window, resolved_max_output_tokens) =
        resolve_runtime_model_token_limits(settings_file, default_backend, &task.model);

    MemoryManager::new(MemoryManagerConfig {
        compact_threshold: task.memory_compact_threshold,
        keep_recent_messages: read_usize_metadata(
            &task.metadata,
            "memory_keep_recent_messages",
            10,
        ),
        model: task.model.clone(),
        model_context_window: read_u64_metadata(
            &task.metadata,
            "model_context_window",
            resolved_context_window.unwrap_or(200_000),
        ),
        reserved_output_tokens: read_u64_metadata(
            &task.metadata,
            "reserved_output_tokens",
            resolved_max_output_tokens.unwrap_or(16_000),
        ),
        autocompact_buffer_tokens: read_u64_metadata(
            &task.metadata,
            "autocompact_buffer_tokens",
            13_000,
        ),
        language: read_string_metadata(&task.metadata, "language", "zh-CN"),
        warning_threshold_percentage: task.memory_threshold_percentage.clamp(1, 100),
        include_memory_warning: read_bool_metadata(&task.metadata, "include_memory_warning", false),
        summary_event_limit: read_usize_metadata(&task.metadata, "summary_event_limit", 40),
        summary_backend: summary_backend.clone(),
        summary_model: Some(summary_model.clone()),
        summary_callback: summary_callback.clone(),
        tool_result_compact_threshold: read_usize_metadata(
            &task.metadata,
            "tool_result_compact_threshold",
            2_000,
        ),
        tool_result_keep_last: read_usize_metadata(&task.metadata, "tool_result_keep_last", 3),
        tool_result_excerpt_head: read_usize_metadata(
            &task.metadata,
            "tool_result_excerpt_head",
            200,
        ),
        tool_result_excerpt_tail: read_usize_metadata(
            &task.metadata,
            "tool_result_excerpt_tail",
            200,
        ),
        tool_calls_keep_last: read_usize_metadata(&task.metadata, "tool_calls_keep_last", 3),
        assistant_no_tool_keep_last: read_usize_metadata(
            &task.metadata,
            "assistant_no_tool_keep_last",
            1,
        ),
        tool_result_artifact_dir: metadata_path(
            &task.metadata,
            "tool_result_artifact_dir",
            ".memory/tool_results",
        ),
        microcompact_trigger_ratio: read_f64_metadata(
            &task.metadata,
            "microcompact_trigger_ratio",
            0.75,
            0.0,
            Some(1.0),
        ),
        microcompact_keep_recent_cycles: read_usize_metadata(
            &task.metadata,
            "microcompact_keep_recent_cycles",
            3,
        ),
        microcompact_min_result_length: read_usize_metadata(
            &task.metadata,
            "microcompact_min_result_length",
            500,
        ),
        microcompact_compactable_tools: read_string_set_metadata(
            &task.metadata,
            "microcompact_compactable_tools",
        ),
        workspace: workspace.clone(),
        session_memory: build_session_memory(
            task,
            workspace,
            summary_callback.clone(),
            summary_backend,
            summary_model,
        ),
    })
}

#[allow(clippy::too_many_arguments)]
pub(super) fn memory_compact_started_event(
    execution_context: Option<&ExecutionContext>,
    memory_manager: &MemoryManager,
    task: &AgentTask,
    cycle_index: u32,
    messages: &[Message],
    previous_prompt_tokens: Option<u64>,
    recent_tool_call_ids: Option<&BTreeSet<String>>,
    force: bool,
) -> Option<RunEvent> {
    if !force
        && !memory_manager.should_attempt_compaction(
            messages,
            previous_prompt_tokens,
            recent_tool_call_ids,
        )
    {
        return None;
    }
    let identity = execution_context.map(|context| &context.metadata);
    let run_id = identity
        .and_then(|metadata| metadata.get("_vv_agent_run_id"))
        .and_then(Value::as_str)
        .unwrap_or(&task.task_id)
        .to_string();
    let trace_id = identity
        .and_then(|metadata| metadata.get("_vv_agent_trace_id"))
        .and_then(Value::as_str)
        .or_else(|| task.metadata.get("trace_id").and_then(Value::as_str))
        .unwrap_or(&run_id)
        .to_string();
    let agent_name = identity
        .and_then(|metadata| metadata.get("_vv_agent_agent_name"))
        .or_else(|| task.metadata.get("agent_name"))
        .and_then(Value::as_str)
        .unwrap_or(&task.task_id)
        .to_string();
    let event = RunEvent::memory_compact_started(
        run_id,
        trace_id,
        agent_name,
        cycle_index,
        messages.len(),
        previous_prompt_tokens.or_else(|| {
            Some(count_messages_tokens(
                messages,
                &memory_manager.config.model,
            ))
        }),
    );
    Some(
        match identity
            .and_then(|metadata| metadata.get("_vv_agent_session_id"))
            .and_then(Value::as_str)
        {
            Some(session_id) => event.with_session_id(session_id),
            None => event,
        },
    )
}

pub(super) fn notify_memory_before_compact(
    execution_context: Option<&ExecutionContext>,
    mut event: RunEvent,
    messages: &[Message],
) -> RunEvent {
    let provider_event = event.clone().with_metadata(
        "messages",
        serde_json::to_value(messages).unwrap_or(Value::Null),
    );
    let mut results = BTreeMap::new();
    let mut errors = Vec::new();
    let mut seen_names = BTreeMap::new();
    for (index, provider) in memory_providers(execution_context).into_iter().enumerate() {
        let provider_name = memory_provider_name(provider, index, &mut seen_names);
        match block_on_memory_future(provider.before_compact(&provider_event)) {
            Ok(result) if !result.metadata.is_empty() => {
                results.insert(
                    provider_name,
                    Value::Object(result.metadata.into_iter().collect()),
                );
            }
            Ok(_) => {}
            Err(error) => errors.push(memory_provider_error(
                provider_name,
                "before_compact",
                error,
            )),
        }
    }
    if !results.is_empty() {
        event = event.with_metadata(
            "memory_provider_results",
            Value::Object(results.into_iter().collect()),
        );
    }
    if !errors.is_empty() {
        event = event.with_metadata("memory_provider_errors", Value::Array(errors));
    }
    event
}

pub(super) fn notify_memory_after_compact(
    execution_context: Option<&ExecutionContext>,
    mut event: RunEvent,
) -> RunEvent {
    let mut errors = Vec::new();
    let mut seen_names = BTreeMap::new();
    for (index, provider) in memory_providers(execution_context).into_iter().enumerate() {
        let provider_name = memory_provider_name(provider, index, &mut seen_names);
        if let Err(error) = block_on_memory_future(provider.after_compact(&event)) {
            errors.push(memory_provider_error(provider_name, "after_compact", error));
        }
    }
    if !errors.is_empty() {
        event = event.with_metadata("memory_provider_errors", Value::Array(errors));
    }
    event
}

fn memory_providers(execution_context: Option<&ExecutionContext>) -> Vec<&Arc<dyn MemoryProvider>> {
    execution_context
        .map(|context| context.memory_providers.iter().collect())
        .unwrap_or_default()
}

pub(super) fn memory_compact_completed_event(
    started_event: &RunEvent,
    cycle_index: u32,
    before_messages: &[Message],
    after_messages: &[Message],
    model: &str,
) -> RunEvent {
    let event = RunEvent::memory_compact_completed(
        started_event.run_id(),
        started_event.trace_id(),
        started_event
            .agent_name()
            .expect("memory compact event has agent identity"),
        cycle_index,
        before_messages.len(),
        after_messages.len(),
        Some(count_messages_tokens(after_messages, model)),
    );
    match started_event.session_id() {
        Some(session_id) => event.with_session_id(session_id),
        None => event,
    }
}

pub(super) fn memory_compact_event_payload(event: &RunEvent) -> BTreeMap<String, Value> {
    let mut payload = event.metadata().clone();
    if let Some(cycle_index) = event.cycle_index() {
        payload.insert("cycle".to_string(), Value::from(cycle_index));
    }
    match event.payload() {
        crate::events::RunEventPayload::MemoryCompactStarted {
            message_count,
            estimated_tokens,
        } => {
            payload.insert("message_count".to_string(), Value::from(*message_count));
            if let Some(estimated_tokens) = estimated_tokens {
                payload.insert(
                    "estimated_tokens".to_string(),
                    Value::from(*estimated_tokens),
                );
            }
        }
        crate::events::RunEventPayload::MemoryCompactCompleted {
            before_count,
            after_count,
            summary_tokens,
        } => {
            payload.insert("before_count".to_string(), Value::from(*before_count));
            payload.insert("after_count".to_string(), Value::from(*after_count));
            if let Some(summary_tokens) = summary_tokens {
                payload.insert("summary_tokens".to_string(), Value::from(*summary_tokens));
            }
        }
        _ => {}
    }
    payload
}

fn memory_provider_name(
    provider: &Arc<dyn MemoryProvider>,
    index: usize,
    seen_names: &mut BTreeMap<String, usize>,
) -> String {
    let base_name = provider
        .provider_name()
        .rsplit("::")
        .next()
        .unwrap_or("MemoryProvider")
        .to_string();
    let seen = seen_names.entry(base_name.clone()).or_insert(0);
    let name = if *seen == 0 {
        base_name
    } else {
        format!("{base_name}#{}", index + 1)
    };
    *seen += 1;
    name
}

fn memory_provider_error(
    provider_name: String,
    stage: &str,
    error: crate::memory::MemoryError,
) -> Value {
    eprintln!("warning: Memory provider {provider_name} {stage} failed: {error}");
    serde_json::json!({
        "provider": provider_name,
        "stage": stage,
        "error": error.to_string(),
        "error_type": "MemoryError",
    })
}