onetaskgraph_core/engine/
delivery.rs1use onetaskgraph_plugin_api::{SourceError, SourceName, Status, StatusCategory, Task, TaskRef};
19use schemars::JsonSchema;
20use serde::{Deserialize, Serialize};
21
22use super::{ConfiguredSource, Engine, EngineError, Qualified};
23use crate::resolve::ResolvedSource;
24use crate::{Failure, GlobalId};
25
26#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
28pub struct TaskStatusSet {
29 pub id: GlobalId,
31 pub status: Status,
33 pub delivered: Vec<Delivered>,
36}
37
38#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
40pub struct Delivered {
41 pub ticket: GlobalId,
43 pub deliverer: GlobalId,
46 #[serde(flatten)]
48 pub outcome: DeliveryOutcome,
49 #[serde(default, skip_serializing_if = "Vec::is_empty")]
52 #[schemars(!skip_serializing_if)]
53 pub pruned: Vec<GlobalId>,
54}
55
56impl Delivered {
57 #[must_use]
59 pub fn failed(&self) -> bool {
60 matches!(self.outcome, DeliveryOutcome::Failed { .. })
61 }
62}
63
64#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
67#[serde(tag = "outcome", rename_all = "kebab-case")]
68pub enum DeliveryOutcome {
69 Written {
71 from: StatusCategory,
73 to: StatusCategory,
75 },
76 Unchanged {
78 from: StatusCategory,
80 },
81 Left {
84 from: StatusCategory,
86 },
87 Failed {
89 #[serde(default, skip_serializing_if = "Option::is_none")]
91 from: Option<StatusCategory>,
92 failure: Failure,
96 },
97}
98
99impl Engine {
100 pub async fn set_task_status(
112 &self,
113 id: &GlobalId,
114 category: StatusCategory,
115 ) -> Result<TaskStatusSet, EngineError> {
116 let source = self.status_writable(&id.source)?;
117 let no_such_task = || EngineError::NoSuchTask { id: id.to_string() };
118 let task = source
119 .source()
120 .get_task(&id.native)
121 .await
122 .map_err(|error| source_failed(source, error))?
123 .ok_or_else(no_such_task)?;
124 let status = source
125 .source()
126 .set_task_status(&id.native, category)
127 .await
128 .map_err(|error| source_failed(source, error))?
129 .ok_or_else(no_such_task)?;
130 let delivers = targets(&task.delivers, &id.source);
131 let delivered = self
132 .deliver(id, status.category, &delivers, &delivers)
133 .await;
134 Ok(TaskStatusSet {
135 id: id.clone(),
136 status,
137 delivered,
138 })
139 }
140
141 pub(crate) async fn deliver(
144 &self,
145 deliverer: &GlobalId,
146 category: StatusCategory,
147 now: &[GlobalId],
148 before: &[GlobalId],
149 ) -> Vec<Delivered> {
150 let mut tickets: Vec<(&GlobalId, bool)> = Vec::new();
151 for ticket in now {
152 if !tickets.iter().any(|(held, _)| *held == ticket) {
153 tickets.push((ticket, true));
154 }
155 }
156 for ticket in before {
157 if !now.contains(ticket) && !tickets.iter().any(|(held, _)| *held == ticket) {
158 tickets.push((ticket, false));
159 }
160 }
161 let mut delivered = Vec::with_capacity(tickets.len());
162 for (ticket, kept) in tickets {
163 delivered.push(self.evaluate(deliverer, category, ticket, kept).await);
164 }
165 delivered
166 }
167
168 async fn evaluate(
171 &self,
172 deliverer: &GlobalId,
173 category: StatusCategory,
174 ticket: &GlobalId,
175 kept: bool,
176 ) -> Delivered {
177 let entry = |outcome, pruned| Delivered {
178 ticket: ticket.clone(),
179 deliverer: deliverer.clone(),
180 outcome,
181 pruned,
182 };
183 let failed = |from, error: &EngineError| DeliveryOutcome::Failed {
184 from,
185 failure: Failure::from(error),
186 };
187 let no_such_task = || EngineError::NoSuchTask {
188 id: ticket.to_string(),
189 };
190 let source = match self.built(&ticket.source) {
191 Ok(source) => source,
192 Err(error) => return entry(failed(None, &error), Vec::new()),
193 };
194 let task = match source.source().get_task(&ticket.native).await {
195 Ok(Some(task)) => task,
196 Ok(None) => return entry(failed(None, &no_such_task()), Vec::new()),
197 Err(error) => {
198 return entry(failed(None, &source_failed(source, error)), Vec::new());
199 }
200 };
201 let from = task.status.category;
202 let held = targets(&task.delivered_by, &ticket.source);
203 let mut named = held.clone();
204 if kept && !named.contains(deliverer) {
205 named.push(deliverer.clone());
206 } else if !kept {
207 named.retain(|other| other != deliverer);
208 }
209
210 let active = matches!(
213 from,
214 StatusCategory::Todo | StatusCategory::Queued | StatusCategory::InProgress
215 );
216 let mut categories = Vec::new();
217 let mut pruned = Vec::new();
218 let mut unreadable = None;
219 if active {
220 for other in &named {
221 if other == deliverer {
222 categories.push(category);
223 continue;
224 }
225 match self.category_of(other).await {
226 Ok(Some(found)) => categories.push(found),
227 Ok(None) => pruned.push(other.clone()),
228 Err(error) => {
229 unreadable = Some(error);
230 break;
231 }
232 }
233 }
234 }
235 if unreadable.is_some() {
238 pruned.clear();
239 }
240 let kept_by: Vec<GlobalId> = named
241 .into_iter()
242 .filter(|other| !pruned.contains(other))
243 .collect();
244 if kept_by != held {
245 let list: Vec<TaskRef> = kept_by
246 .iter()
247 .map(|other| TaskRef::qualified(&other.source, &other.native))
248 .collect();
249 match source
250 .source()
251 .set_delivered_by(&ticket.native, &list)
252 .await
253 {
254 Ok(Some(())) => {}
255 Ok(None) => return entry(failed(Some(from), &no_such_task()), Vec::new()),
256 Err(error) => {
257 return entry(
258 failed(Some(from), &source_failed(source, error)),
259 Vec::new(),
260 );
261 }
262 }
263 }
264 if let Some(error) = unreadable {
265 return entry(failed(Some(from), &error), Vec::new());
266 }
267 if !active {
268 return entry(DeliveryOutcome::Left { from }, pruned);
269 }
270 let Some(to) = settled(&categories).filter(|to| *to != from) else {
271 return entry(DeliveryOutcome::Unchanged { from }, pruned);
272 };
273 match source.source().set_task_status(&ticket.native, to).await {
274 Ok(Some(status)) => entry(
275 DeliveryOutcome::Written {
276 from,
277 to: status.category,
278 },
279 pruned,
280 ),
281 Ok(None) => entry(failed(Some(from), &no_such_task()), pruned),
282 Err(error) => entry(failed(Some(from), &source_failed(source, error)), pruned),
283 }
284 }
285
286 async fn category_of(
288 &self,
289 deliverer: &GlobalId,
290 ) -> Result<Option<StatusCategory>, EngineError> {
291 let source = self.built(&deliverer.source)?;
292 source
293 .source()
294 .get_task(&deliverer.native)
295 .await
296 .map(|task| task.map(|task| task.status.category))
297 .map_err(|error| source_failed(source, error))
298 }
299
300 pub(super) fn built(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
302 let name = self.known(name)?;
303 match self.sources.iter().find(|source| source.name() == &name) {
304 Some(ConfiguredSource::Ready(source)) => Ok(source),
305 Some(ConfiguredSource::Unavailable(source)) => Err(EngineError::SourceUnavailable {
306 name: name.to_string(),
307 error: source.error().clone(),
308 }),
309 None => Err(EngineError::NoSources),
311 }
312 }
313
314 fn status_writable(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
316 let source = self.built(name)?;
317 if source.source().writes().is_supported() {
318 return Ok(source);
319 }
320 Err(EngineError::StatusNotWritable {
321 name: source.name().to_string(),
322 kind: source.kind().to_owned(),
323 })
324 }
325}
326
327#[must_use]
341pub fn settled(categories: &[StatusCategory]) -> Option<StatusCategory> {
342 use StatusCategory::{Backlog, Cancelled, Done, Draft, InProgress, Queued, Todo, Unknown};
343 if categories.contains(&Done)
344 && categories
345 .iter()
346 .all(|category| matches!(category, Done | Cancelled))
347 {
348 return Some(Done);
349 }
350 if categories.contains(&InProgress) {
351 return Some(InProgress);
352 }
353 if categories.contains(&Queued) {
354 return Some(Queued);
355 }
356 if categories
357 .iter()
358 .any(|category| matches!(category, Todo | Cancelled | Unknown | Done))
359 {
360 return Some(Todo);
361 }
362 if categories.is_empty() {
363 return Some(Todo);
364 }
365 debug_assert!(
366 categories
367 .iter()
368 .all(|category| matches!(category, Draft | Backlog))
369 );
370 None
371}
372
373pub(crate) fn targets(list: &[TaskRef], near: &SourceName) -> Vec<GlobalId> {
376 list.iter()
377 .filter_map(|entry| entry.in_source(near).as_str().parse().ok())
378 .collect()
379}
380
381pub(crate) fn qualified_task(id: GlobalId, task: Task) -> Qualified<Task> {
384 let delivers = task
385 .delivers
386 .iter()
387 .map(|entry| entry.in_source(&id.source))
388 .collect();
389 let delivered_by = task
390 .delivered_by
391 .iter()
392 .map(|entry| entry.in_source(&id.source))
393 .collect();
394 Qualified {
395 id,
396 item: Task {
397 delivers,
398 delivered_by,
399 ..task
400 },
401 }
402}
403
404pub(super) fn source_failed(source: &ResolvedSource, error: SourceError) -> EngineError {
405 EngineError::SourceFailed {
406 name: source.name().to_string(),
407 error,
408 }
409}