use std::pin::Pin; use dashmap::DashMap; use eyre::Result; use time::OffsetDateTime; use tokio_stream::Stream; use tonic::{Request, Response, Status}; use tracing::{Level, instrument}; use crate::{ aclient::{ history::{History, HistoryId, store::HistoryStore}, settings::Settings, }, api::{ DAEMON_PROTOCOL_VERSION, generated::history::{ EndHistoryReply, EndHistoryRequest, HistoryEntry, HistoryEventKind, StartHistoryReply, StartHistoryRequest, TailHistoryReply, TailHistoryRequest, history_server::{History as HistorySvc, HistoryServer}, }, }, daemon::DaemonHandle, events::DaemonEvent, }; /// The gRPC service implementation. /// /// This is a thin wrapper that delegates to the component's shared state. pub(crate) struct HistoryService { /// Commands currently running (not yet completed). running: DashMap, /// Handle to the daemon (set during start). pub(crate) handle: DaemonHandle, /// History store for pushing records (set during start). pub(crate) history_store: HistoryStore, } impl HistoryService { pub(crate) async fn new(handle: DaemonHandle) -> Result { let host_id = Settings::host_id().await?; let history_store = HistoryStore::new(handle.store().clone(), host_id, *handle.encryption_key()); Ok(Self { running: DashMap::new(), handle, history_store, }) } /// Get a tonic server for this service. pub(crate) fn into_server(self) -> HistoryServer { HistoryServer::new(self) } } fn history_to_tail_reply(kind: HistoryEventKind, history: History) -> TailHistoryReply { TailHistoryReply { kind: kind as i32, history: Some(HistoryEntry { timestamp: history.timestamp.unix_timestamp_nanos() as u64, id: history.id.0, command: history.command, cwd: history.cwd, session: history.session, hostname: history.hostname, author: history.author, intent: history.intent.unwrap_or_default(), exit: history.exit, duration: history.duration, }), } } #[tonic::async_trait] impl HistorySvc for HistoryService { type TailHistoryStream = Pin> + Send>>; #[instrument(skip_all, level = Level::INFO)] async fn start_history( &self, request: Request, ) -> Result, Status> { let req = request.into_inner(); let timestamp = OffsetDateTime::from_unix_timestamp_nanos(i128::from(req.timestamp)) .map_err(|_| { Status::invalid_argument( "failed to parse timestamp as unix time (expected nanos since epoch)", ) })?; let h: History = History::daemon() .timestamp(timestamp) .command(req.command) .cwd(req.cwd) .session(req.session) .hostname(req.hostname) .author(req.author) .intent(req.intent) .build() .into(); self.handle.emit(DaemonEvent::HistoryStarted(h.clone())); let id = h.id.clone(); tracing::info!(id = id.to_string(), "start history"); self.running.insert(id.clone(), h); let reply = StartHistoryReply { id: id.to_string(), version: env!("CARGO_PKG_VERSION").to_string(), protocol: DAEMON_PROTOCOL_VERSION, }; Ok(Response::new(reply)) } #[instrument(skip_all, level = Level::INFO)] #[expect(clippy::significant_drop_tightening, reason = "Would be a logic-bug")] async fn end_history( &self, request: Request, ) -> Result, Status> { let req = request.into_inner(); let id = HistoryId(req.id); if let Some((_, mut history)) = self.running.remove(&id) { history.exit = req.exit; history.duration = match req.duration { 0 => i64::try_from( (OffsetDateTime::now_utc() - history.timestamp).whole_nanoseconds(), ) .expect("failed to convert calculated duration to i64"), value => i64::try_from(value).expect("failed to get i64 duration"), }; self.handle .history_db() .save(&history) .await .map_err(|e| Status::internal(format!("failed to write to db: {e:?}")))?; tracing::info!(id = id.0, duration = history.duration, "end history"); let (record_id, idx) = self .history_store .push(history.clone()) .await .map_err(|e| Status::internal(format!("failed to push record to store: {e:?}")))?; self.handle.emit(DaemonEvent::HistoryEnded(history)); let reply = EndHistoryReply { id: record_id.0.to_string(), idx, version: env!("CARGO_PKG_VERSION").to_string(), protocol: DAEMON_PROTOCOL_VERSION, }; return Ok(Response::new(reply)); } Err(Status::not_found(format!( "could not find history with id: {id}" ))) } #[instrument(skip_all, level = Level::INFO)] #[expect(clippy::significant_drop_tightening, reason = "Would be a logic-bug")] async fn tail_history( &self, _request: Request, ) -> Result, Status> { let mut rx = self.handle.subscribe(); let (tx, out_rx) = tokio::sync::mpsc::channel::>(128); tokio::spawn(async move { loop { let event = match rx.recv().await { Ok(event) => event, Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => { drop( tx.send(Err(Status::resource_exhausted(format!( "tail stream lagged behind and dropped {skipped} events" )))) .await, ); break; } Err(tokio::sync::broadcast::error::RecvError::Closed) => break, }; let reply = match event { DaemonEvent::HistoryStarted(history) => { Some(history_to_tail_reply(HistoryEventKind::Started, history)) } DaemonEvent::HistoryEnded(history) => { Some(history_to_tail_reply(HistoryEventKind::Ended, history)) } _ => None, }; if let Some(reply) = reply && tx.send(Ok(reply)).await.is_err() { break; } } }); let stream = tokio_stream::wrappers::ReceiverStream::new(out_rx); Ok(Response::new(Box::pin(stream))) } }