use anyhow::Result;
use futures::{Stream, StreamExt};
use kube::Client;
use kube::api::{Api, ListParams};
use kube::runtime::{watcher, watcher::Config as WatcherConfig};
use polyc_controller::Conversation;
use serde::Deserialize;
use crate::components::fleet::ConversationRow;
pub(crate) async fn build_client() -> Result<Client> {
Ok(Client::try_default().await?)
}
pub(crate) async fn load_fleet(client: &Client, namespace: &str) -> Result<Vec<ConversationRow>> {
let api: Api<Conversation> = Api::namespaced(client.clone(), namespace);
let params = ListParams::default();
let list = api.list(¶ms).await?;
Ok(list.items.iter().map(row_from).collect())
}
#[allow(clippy::unused_async)]
pub(crate) async fn watch_fleet(
client: &Client,
namespace: &str,
) -> Result<impl Stream<Item = Result<ConversationRow>>> {
let api: Api<Conversation> = Api::namespaced(client.clone(), namespace);
let stream = watcher(api, WatcherConfig::default()).filter_map(|event| async move {
match event {
Ok(watcher::Event::Apply(obj) | watcher::Event::InitApply(obj)) => {
Some(Ok(row_from(&obj)))
}
Ok(watcher::Event::Init | watcher::Event::InitDone | watcher::Event::Delete(_)) => None,
Err(e) => Some(Err(anyhow::Error::new(e))),
}
});
Ok(stream)
}
#[must_use]
pub(crate) fn row_from(conv: &Conversation) -> ConversationRow {
let id = conv.metadata.name.clone().unwrap_or_default();
let status = conv.status.as_ref();
ConversationRow {
id,
model: conv.spec.model.clone(),
principal_ref: conv.spec.principal_ref.clone(),
phase: status.and_then(|s| s.phase.clone()),
pod_ip: status.and_then(|s| s.pod_ip.clone()),
harness_ready: status.is_some_and(|s| s.harness_ready),
closed: status.is_some_and(|s| s.closed),
committed_turns: None,
input_tokens: None,
output_tokens: None,
pending_approvals: None,
}
}
#[derive(Debug, Clone, Deserialize)]
pub(crate) struct OverviewResponse {
pub committed_turns: usize,
pub usage: UsageJson,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub(crate) struct UsageJson {
pub input_tokens: u64,
pub output_tokens: u64,
}
#[derive(Debug, Clone, Deserialize)]
pub(crate) struct ApprovalsCountResponse {
pub approvals: Vec<ApprovalCountEntry>,
}
#[derive(Debug, Clone, Deserialize)]
pub(crate) struct ApprovalCountEntry {
pub response: Option<serde_json::Value>,
}
pub(crate) async fn enrich_row(forensics_base_url: &str, row: &mut ConversationRow) -> Result<()> {
let metrics = fetch_metrics(forensics_base_url, &row.id).await?;
row.committed_turns = Some(metrics.committed_turns);
row.input_tokens = Some(metrics.input_tokens);
row.output_tokens = Some(metrics.output_tokens);
row.pending_approvals = Some(metrics.pending_approvals);
Ok(())
}
#[derive(Debug, Clone, Default)]
pub(crate) struct ConversationMetrics {
pub committed_turns: usize,
pub input_tokens: u64,
pub output_tokens: u64,
pub pending_approvals: usize,
}
pub(crate) async fn fetch_metrics(
forensics_base_url: &str,
conversation_id: &str,
) -> Result<ConversationMetrics> {
let client = crate::data::http_client();
let base = forensics_base_url.trim_end_matches('/');
let overview: OverviewResponse = client
.get(format!("{base}/api/conversations/{conversation_id}"))
.send()
.await?
.error_for_status()?
.json()
.await?;
let approvals: ApprovalsCountResponse = client
.get(format!(
"{base}/api/conversations/{conversation_id}/approvals"
))
.send()
.await?
.error_for_status()?
.json()
.await?;
let pending_approvals = approvals
.approvals
.iter()
.filter(|a| a.response.is_none())
.count();
Ok(ConversationMetrics {
committed_turns: overview.committed_turns,
input_tokens: overview.usage.input_tokens,
output_tokens: overview.usage.output_tokens,
pending_approvals,
})
}
#[derive(Debug, Clone, Deserialize)]
pub(crate) struct ConversationListResponse {
pub conversation_ids: Vec<String>,
}
pub(crate) async fn list_conversations(forensics_base_url: &str) -> Result<Vec<String>> {
let base = forensics_base_url.trim_end_matches('/');
let resp: ConversationListResponse = crate::data::http_client()
.get(format!("{base}/api/conversations"))
.send()
.await?
.error_for_status()?
.json()
.await?;
Ok(resp.conversation_ids)
}
#[must_use]
pub(crate) fn forensics_row(id: &str) -> ConversationRow {
ConversationRow {
id: id.to_owned(),
model: String::new(),
principal_ref: String::new(),
phase: Some("event-log".to_owned()),
pod_ip: None,
harness_ready: false,
closed: false,
committed_turns: None,
input_tokens: None,
output_tokens: None,
pending_approvals: None,
}
}