use super::usage::{SessionUsageLedger, SessionUsageRecord, USAGE_ACTIVITY_PREFIX};
use crate::output::NormalizedUsageSnapshot;
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,
})
}
}