use std::collections::BTreeMap;
use std::sync::Arc;
use serde::Serialize;
use taquba::object_store::ObjectStore;
use taquba::{JobRecord, QueueReader, QueueStats, ReaderMode, ReaderOptions};
use crate::definition_store::{DefinitionError, DefinitionStore};
use crate::partition::Partition;
use crate::readiness::{NodeState, current_records, node_states};
use crate::records::{
self, Entry, GRAPH_RUNS_PREFIX, GRAPHS_PREFIX, GraphRecord, GraphRunRecord, GraphRunState,
NodeRecord, ReadError, RecordError, RequestRecord,
};
use crate::request::RequestId;
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error(transparent)]
Queue(#[from] taquba::Error),
#[error(transparent)]
Definition(#[from] DefinitionError),
#[error(transparent)]
Record(#[from] RecordError),
#[error("definition `{0}` is not in the definition store")]
UnknownDefinition(String),
}
impl From<ReadError> for Error {
fn from(e: ReadError) -> Self {
match e {
ReadError::Queue(e) => Error::Queue(e),
ReadError::Record(e) => Error::Record(e),
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize)]
pub struct RunCounts {
pub active: usize,
pub cancelled: usize,
pub complete: usize,
pub failed: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct RunSummary {
pub partition: Partition,
pub record: GraphRunRecord,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct GraphStatus {
pub name: String,
pub adopted: Option<GraphRecord>,
pub runs: RunCounts,
pub latest: Option<RunSummary>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct NodeStatus {
pub name: String,
pub pool: String,
pub state: NodeState,
pub record: Option<NodeRecord>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct GraphRunStatus {
pub graph: String,
pub partition: Partition,
pub record: GraphRunRecord,
pub nodes: Vec<NodeStatus>,
}
pub struct StatusReader {
reader: QueueReader,
definitions: Arc<DefinitionStore>,
}
impl StatusReader {
pub async fn open(
store: Arc<dyn ObjectStore>,
queue_path: &str,
definitions: Arc<DefinitionStore>,
) -> Result<Self, Error> {
let options = ReaderOptions::default().mode(ReaderMode::FollowLatest);
let reader = QueueReader::open_with_options(store, queue_path, options).await?;
Ok(StatusReader::new(reader, definitions))
}
pub fn new(reader: QueueReader, definitions: Arc<DefinitionStore>) -> Self {
StatusReader {
reader,
definitions,
}
}
pub async fn graphs(&self) -> Result<Vec<GraphStatus>, Error> {
let mut graphs: BTreeMap<String, GraphStatus> = BTreeMap::new();
let adopted: Vec<Entry<GraphRecord>> =
records::scan(self.reader.view(), GRAPHS_PREFIX.as_bytes()).await?;
for Entry { key, record, .. } in adopted {
let Some(name) = records::parse_graph_key(&key) else {
continue;
};
graph_entry(&mut graphs, &name).adopted = Some(record);
}
let runs: Vec<Entry<GraphRunRecord>> =
records::scan(self.reader.view(), GRAPH_RUNS_PREFIX.as_bytes()).await?;
for Entry { key, record, .. } in runs {
let Some((name, partition)) = records::parse_graph_run_key(&key) else {
continue;
};
let graph = graph_entry(&mut graphs, &name);
match record.state {
GraphRunState::Active => graph.runs.active += 1,
GraphRunState::Cancelled => graph.runs.cancelled += 1,
GraphRunState::Complete => graph.runs.complete += 1,
GraphRunState::Failed => graph.runs.failed += 1,
}
graph.latest = Some(RunSummary { partition, record });
}
Ok(graphs.into_values().collect())
}
pub async fn runs(&self, graph: &str) -> Result<Vec<RunSummary>, Error> {
let prefix = format!("{GRAPH_RUNS_PREFIX}{graph}/");
let mut runs = Vec::new();
let entries: Vec<Entry<GraphRunRecord>> =
records::scan(self.reader.view(), prefix.as_bytes()).await?;
for Entry { key, record, .. } in entries {
let Some((_, partition)) = records::parse_graph_run_key(&key) else {
continue;
};
runs.push(RunSummary { partition, record });
}
Ok(runs)
}
pub async fn run(
&self,
graph: &str,
partition: &Partition,
) -> Result<Option<GraphRunStatus>, Error> {
let key = records::graph_run_key(graph, partition);
let Some(record) = records::read::<GraphRunRecord>(self.reader.view(), &key).await? else {
return Ok(None);
};
let definition = self
.definitions
.get(&record.definition)
.await?
.ok_or_else(|| Error::UnknownDefinition(record.definition.clone()))?;
let mut node_records = BTreeMap::new();
for node in definition.nodes() {
let key = records::node_record_key(graph, partition, node);
if let Some(node_record) = records::read::<NodeRecord>(self.reader.view(), &key).await?
{
node_records.insert(node.name().to_string(), node_record);
}
}
let states = node_states(&definition, ¤t_records(&record, &node_records));
let nodes = definition
.nodes()
.iter()
.map(|node| NodeStatus {
name: node.name().to_string(),
pool: node.pool().to_string(),
state: states[node.name()],
record: node_records.remove(node.name()),
})
.collect();
Ok(Some(GraphRunStatus {
graph: graph.to_string(),
partition: partition.clone(),
record,
nodes,
}))
}
pub async fn request(&self, id: &RequestId) -> Result<Option<RequestRecord>, Error> {
Ok(records::read(self.reader.view(), &records::request_key(id)).await?)
}
pub async fn queues(&self) -> Result<Vec<QueueStats>, Error> {
let mut names = self.reader.view().list_queues().await?;
names.sort();
let mut stats = Vec::with_capacity(names.len());
for name in names {
stats.push(self.reader.view().stats(&name).await?);
}
Ok(stats)
}
pub async fn dead_jobs(
&self,
queue: &str,
limit: usize,
) -> Result<Option<Vec<JobRecord>>, Error> {
if !self
.reader
.view()
.list_queues()
.await?
.iter()
.any(|q| q == queue)
{
return Ok(None);
}
Ok(Some(
self.reader.view().dead_jobs(queue, None, limit).await?,
))
}
pub async fn close(self) -> Result<(), Error> {
Ok(self.reader.close().await?)
}
}
fn graph_entry<'a>(
graphs: &'a mut BTreeMap<String, GraphStatus>,
name: &str,
) -> &'a mut GraphStatus {
graphs
.entry(name.to_string())
.or_insert_with(|| GraphStatus {
name: name.to_string(),
adopted: None,
runs: RunCounts::default(),
latest: None,
})
}