magi-code 0.96.2

Repository-aware CLI coding agent for terminal work
Documentation
//! Worker-owned usage persistence shared by terminal and local clients.
use super::usage::{SessionUsageLedger, SessionUsageRecord, USAGE_ACTIVITY_PREFIX};
use crate::output::NormalizedUsageSnapshot;

/// Called only by execution workers, before attempting bounded event delivery.
pub(crate) struct SessionUsageRecorder {
    run_id: String,
    session: Option<crate::sessions::Session>,
    cwd: std::path::PathBuf,
    pub(crate) ledger: std::sync::Arc<std::sync::Mutex<SessionUsageLedger>>,
}

impl SessionUsageRecorder {
    pub(crate) fn new(
        session: Option<crate::sessions::Session>,
        cwd: std::path::PathBuf,
        ledger: std::sync::Arc<std::sync::Mutex<SessionUsageLedger>>,
    ) -> Self {
        Self {
            run_id: uuid::Uuid::new_v4().to_string(),
            session,
            cwd,
            ledger,
        }
    }

    pub(crate) fn request_started(
        &self,
        source: &str,
        request_sequence: u64,
    ) -> anyhow::Result<Option<crate::output::ActivityEvent>> {
        let record = self
            .ledger
            .lock()
            .unwrap_or_else(|error| error.into_inner())
            .pending_request(
                format!("{USAGE_ACTIVITY_PREFIX}{}/{source}", self.run_id),
                request_sequence,
            );
        record
            .map(|record| self.record(source, record.usage, request_sequence, false))
            .transpose()
    }

    pub(crate) fn record_activity(
        &self,
        event: &crate::output::ActivityEvent,
    ) -> anyhow::Result<Option<crate::output::ActivityEvent>> {
        use crate::output::ActivityEvent;
        match event {
            ActivityEvent::UsageSnapshot {
                id,
                usage,
                request_sequence,
                final_usage,
            } => self
                .record(
                    &format!("activity/{}", id.0),
                    usage.whole_run,
                    *request_sequence,
                    *final_usage,
                )
                .map(Some),
            ActivityEvent::UsageUpdate {
                id,
                request_sequence,
                ..
            } => self.request_started(&format!("activity/{}", id.0), *request_sequence),
            _ => Ok(None),
        }
    }

    pub(crate) fn record(
        &self,
        source: &str,
        usage: NormalizedUsageSnapshot,
        request_sequence: u64,
        final_usage: bool,
    ) -> anyhow::Result<crate::output::ActivityEvent> {
        let record = SessionUsageRecord {
            id: format!("{USAGE_ACTIVITY_PREFIX}{}/{source}", self.run_id),
            usage,
            request_sequence,
            final_usage,
        };
        self.ledger
            .lock()
            .unwrap_or_else(|error| error.into_inner())
            .observe(record.clone());
        crate::sessions::record_session_event(
            self.session.as_ref(),
            &self.cwd,
            crate::sessions::SessionEventKind::SessionUsage,
            serde_json::to_value(&record)?,
        )?;
        Ok(crate::output::ActivityEvent::UsageSnapshot {
            id: crate::output::ActivityId::new(record.id),
            usage: crate::output::NormalizedUsageAggregate {
                whole_run: usage,
                ..Default::default()
            },
            request_sequence,
            final_usage,
        })
    }
}