use onetaskgraph_plugin_api::{SourceError, SourceName, Status, StatusCategory, Task, TaskRef};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use super::{ConfiguredSource, Engine, EngineError, Qualified};
use crate::resolve::ResolvedSource;
use crate::{Failure, GlobalId};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct TaskStatusSet {
pub id: GlobalId,
pub status: Status,
pub delivered: Vec<Delivered>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct Delivered {
pub ticket: GlobalId,
pub deliverer: GlobalId,
#[serde(flatten)]
pub outcome: DeliveryOutcome,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
#[schemars(!skip_serializing_if)]
pub pruned: Vec<GlobalId>,
}
impl Delivered {
#[must_use]
pub fn failed(&self) -> bool {
matches!(self.outcome, DeliveryOutcome::Failed { .. })
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(tag = "outcome", rename_all = "kebab-case")]
pub enum DeliveryOutcome {
Written {
from: StatusCategory,
to: StatusCategory,
},
Unchanged {
from: StatusCategory,
},
Left {
from: StatusCategory,
},
Failed {
#[serde(default, skip_serializing_if = "Option::is_none")]
from: Option<StatusCategory>,
failure: Failure,
},
}
impl Engine {
pub async fn set_task_status(
&self,
id: &GlobalId,
category: StatusCategory,
) -> Result<TaskStatusSet, EngineError> {
let source = self.status_writable(&id.source)?;
let no_such_task = || EngineError::NoSuchTask { id: id.to_string() };
let task = source
.source()
.get_task(&id.native)
.await
.map_err(|error| source_failed(source, error))?
.ok_or_else(no_such_task)?;
let status = source
.source()
.set_task_status(&id.native, category)
.await
.map_err(|error| source_failed(source, error))?
.ok_or_else(no_such_task)?;
let delivers = targets(&task.delivers, &id.source);
let delivered = self
.deliver(id, status.category, &delivers, &delivers)
.await;
Ok(TaskStatusSet {
id: id.clone(),
status,
delivered,
})
}
pub(crate) async fn deliver(
&self,
deliverer: &GlobalId,
category: StatusCategory,
now: &[GlobalId],
before: &[GlobalId],
) -> Vec<Delivered> {
let mut tickets: Vec<(&GlobalId, bool)> = Vec::new();
for ticket in now {
if !tickets.iter().any(|(held, _)| *held == ticket) {
tickets.push((ticket, true));
}
}
for ticket in before {
if !now.contains(ticket) && !tickets.iter().any(|(held, _)| *held == ticket) {
tickets.push((ticket, false));
}
}
let mut delivered = Vec::with_capacity(tickets.len());
for (ticket, kept) in tickets {
delivered.push(self.evaluate(deliverer, category, ticket, kept).await);
}
delivered
}
async fn evaluate(
&self,
deliverer: &GlobalId,
category: StatusCategory,
ticket: &GlobalId,
kept: bool,
) -> Delivered {
let entry = |outcome, pruned| Delivered {
ticket: ticket.clone(),
deliverer: deliverer.clone(),
outcome,
pruned,
};
let failed = |from, error: &EngineError| DeliveryOutcome::Failed {
from,
failure: Failure::from(error),
};
let no_such_task = || EngineError::NoSuchTask {
id: ticket.to_string(),
};
let source = match self.built(&ticket.source) {
Ok(source) => source,
Err(error) => return entry(failed(None, &error), Vec::new()),
};
let task = match source.source().get_task(&ticket.native).await {
Ok(Some(task)) => task,
Ok(None) => return entry(failed(None, &no_such_task()), Vec::new()),
Err(error) => {
return entry(failed(None, &source_failed(source, error)), Vec::new());
}
};
let from = task.status.category;
let held = targets(&task.delivered_by, &ticket.source);
let mut named = held.clone();
if kept && !named.contains(deliverer) {
named.push(deliverer.clone());
} else if !kept {
named.retain(|other| other != deliverer);
}
let active = matches!(
from,
StatusCategory::Todo | StatusCategory::Queued | StatusCategory::InProgress
);
let mut categories = Vec::new();
let mut pruned = Vec::new();
let mut unreadable = None;
if active {
for other in &named {
if other == deliverer {
categories.push(category);
continue;
}
match self.category_of(other).await {
Ok(Some(found)) => categories.push(found),
Ok(None) => pruned.push(other.clone()),
Err(error) => {
unreadable = Some(error);
break;
}
}
}
}
if unreadable.is_some() {
pruned.clear();
}
let kept_by: Vec<GlobalId> = named
.into_iter()
.filter(|other| !pruned.contains(other))
.collect();
if kept_by != held {
let list: Vec<TaskRef> = kept_by
.iter()
.map(|other| TaskRef::qualified(&other.source, &other.native))
.collect();
match source
.source()
.set_delivered_by(&ticket.native, &list)
.await
{
Ok(Some(())) => {}
Ok(None) => return entry(failed(Some(from), &no_such_task()), Vec::new()),
Err(error) => {
return entry(
failed(Some(from), &source_failed(source, error)),
Vec::new(),
);
}
}
}
if let Some(error) = unreadable {
return entry(failed(Some(from), &error), Vec::new());
}
if !active {
return entry(DeliveryOutcome::Left { from }, pruned);
}
let Some(to) = settled(&categories).filter(|to| *to != from) else {
return entry(DeliveryOutcome::Unchanged { from }, pruned);
};
match source.source().set_task_status(&ticket.native, to).await {
Ok(Some(status)) => entry(
DeliveryOutcome::Written {
from,
to: status.category,
},
pruned,
),
Ok(None) => entry(failed(Some(from), &no_such_task()), pruned),
Err(error) => entry(failed(Some(from), &source_failed(source, error)), pruned),
}
}
async fn category_of(
&self,
deliverer: &GlobalId,
) -> Result<Option<StatusCategory>, EngineError> {
let source = self.built(&deliverer.source)?;
source
.source()
.get_task(&deliverer.native)
.await
.map(|task| task.map(|task| task.status.category))
.map_err(|error| source_failed(source, error))
}
pub(super) fn built(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
let name = self.known(name)?;
match self.sources.iter().find(|source| source.name() == &name) {
Some(ConfiguredSource::Ready(source)) => Ok(source),
Some(ConfiguredSource::Unavailable(source)) => Err(EngineError::SourceUnavailable {
name: name.to_string(),
error: source.error().clone(),
}),
None => Err(EngineError::NoSources),
}
}
fn status_writable(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
let source = self.built(name)?;
if source.source().writes().is_supported() {
return Ok(source);
}
Err(EngineError::StatusNotWritable {
name: source.name().to_string(),
kind: source.kind().to_owned(),
})
}
}
#[must_use]
pub fn settled(categories: &[StatusCategory]) -> Option<StatusCategory> {
use StatusCategory::{Backlog, Cancelled, Done, Draft, InProgress, Queued, Todo, Unknown};
if categories.contains(&Done)
&& categories
.iter()
.all(|category| matches!(category, Done | Cancelled))
{
return Some(Done);
}
if categories.contains(&InProgress) {
return Some(InProgress);
}
if categories.contains(&Queued) {
return Some(Queued);
}
if categories
.iter()
.any(|category| matches!(category, Todo | Cancelled | Unknown | Done))
{
return Some(Todo);
}
if categories.is_empty() {
return Some(Todo);
}
debug_assert!(
categories
.iter()
.all(|category| matches!(category, Draft | Backlog))
);
None
}
pub(crate) fn targets(list: &[TaskRef], near: &SourceName) -> Vec<GlobalId> {
list.iter()
.filter_map(|entry| entry.in_source(near).as_str().parse().ok())
.collect()
}
pub(crate) fn qualified_task(id: GlobalId, task: Task) -> Qualified<Task> {
let delivers = task
.delivers
.iter()
.map(|entry| entry.in_source(&id.source))
.collect();
let delivered_by = task
.delivered_by
.iter()
.map(|entry| entry.in_source(&id.source))
.collect();
Qualified {
id,
item: Task {
delivers,
delivered_by,
..task
},
}
}
pub(super) fn source_failed(source: &ResolvedSource, error: SourceError) -> EngineError {
EngineError::SourceFailed {
name: source.name().to_string(),
error,
}
}