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
122 .source()
123 .set_task_status_reading(&id.native, category)
124 .await
125 .map_err(|error| source_failed(source, error))?
126 .ok_or_else(no_such_task)?;
127 let status = task.status.clone();
128 let delivers = targets(&task.delivers, &id.source);
129 let delivered = self
130 .deliver(id, status.category, &delivers, &delivers)
131 .await;
132 Ok(TaskStatusSet {
133 id: id.clone(),
134 status,
135 delivered,
136 })
137 }
138
139 pub(crate) async fn deliver(
142 &self,
143 deliverer: &GlobalId,
144 category: StatusCategory,
145 now: &[GlobalId],
146 before: &[GlobalId],
147 ) -> Vec<Delivered> {
148 let mut tickets: Vec<(&GlobalId, bool)> = Vec::new();
149 for ticket in now {
150 if !tickets.iter().any(|(held, _)| *held == ticket) {
151 tickets.push((ticket, true));
152 }
153 }
154 for ticket in before {
155 if !now.contains(ticket) && !tickets.iter().any(|(held, _)| *held == ticket) {
156 tickets.push((ticket, false));
157 }
158 }
159 let mut delivered = Vec::with_capacity(tickets.len());
160 for (ticket, kept) in tickets {
161 delivered.push(self.evaluate(deliverer, category, ticket, kept).await);
162 }
163 delivered
164 }
165
166 async fn evaluate(
169 &self,
170 deliverer: &GlobalId,
171 category: StatusCategory,
172 ticket: &GlobalId,
173 kept: bool,
174 ) -> Delivered {
175 let entry = |outcome, pruned| Delivered {
176 ticket: ticket.clone(),
177 deliverer: deliverer.clone(),
178 outcome,
179 pruned,
180 };
181 let failed = |from, error: &EngineError| DeliveryOutcome::Failed {
182 from,
183 failure: Failure::from(error),
184 };
185 let no_such_task = || EngineError::NoSuchTask {
186 id: ticket.to_string(),
187 };
188 let source = match self.built(&ticket.source) {
189 Ok(source) => source,
190 Err(error) => return entry(failed(None, &error), Vec::new()),
191 };
192 let task = match source.source().get_task(&ticket.native).await {
193 Ok(Some(task)) => task,
194 Ok(None) => return entry(failed(None, &no_such_task()), Vec::new()),
195 Err(error) => {
196 return entry(failed(None, &source_failed(source, error)), Vec::new());
197 }
198 };
199 let from = task.status.category;
200 let held = targets(&task.delivered_by, &ticket.source);
201 let mut named = held.clone();
202 if kept && !named.contains(deliverer) {
203 named.push(deliverer.clone());
204 } else if !kept {
205 named.retain(|other| other != deliverer);
206 }
207
208 let active = matches!(
211 from,
212 StatusCategory::Todo | StatusCategory::Queued | StatusCategory::InProgress
213 );
214 let mut categories = Vec::new();
215 let mut pruned = Vec::new();
216 let mut unreadable = None;
217 if active {
218 for other in &named {
219 if other == deliverer {
220 categories.push(category);
221 continue;
222 }
223 match self.category_of(other).await {
224 Ok(Some(found)) => categories.push(found),
225 Ok(None) => pruned.push(other.clone()),
226 Err(error) => {
227 unreadable = Some(error);
228 break;
229 }
230 }
231 }
232 }
233 if unreadable.is_some() {
236 pruned.clear();
237 }
238 let kept_by: Vec<GlobalId> = named
239 .into_iter()
240 .filter(|other| !pruned.contains(other))
241 .collect();
242 if kept_by != held {
243 let list: Vec<TaskRef> = kept_by
244 .iter()
245 .map(|other| TaskRef::qualified(&other.source, &other.native))
246 .collect();
247 match source
248 .source()
249 .set_delivered_by(&ticket.native, &list)
250 .await
251 {
252 Ok(Some(())) => {}
253 Ok(None) => return entry(failed(Some(from), &no_such_task()), Vec::new()),
254 Err(error) => {
255 return entry(
256 failed(Some(from), &source_failed(source, error)),
257 Vec::new(),
258 );
259 }
260 }
261 }
262 if let Some(error) = unreadable {
263 return entry(failed(Some(from), &error), Vec::new());
264 }
265 if !active {
266 return entry(DeliveryOutcome::Left { from }, pruned);
267 }
268 let Some(to) = settled(&categories).filter(|to| *to != from) else {
269 return entry(DeliveryOutcome::Unchanged { from }, pruned);
270 };
271 match source.source().set_task_status(&ticket.native, to).await {
272 Ok(Some(status)) => entry(
273 DeliveryOutcome::Written {
274 from,
275 to: status.category,
276 },
277 pruned,
278 ),
279 Ok(None) => entry(failed(Some(from), &no_such_task()), pruned),
280 Err(error) => entry(failed(Some(from), &source_failed(source, error)), pruned),
281 }
282 }
283
284 async fn category_of(
286 &self,
287 deliverer: &GlobalId,
288 ) -> Result<Option<StatusCategory>, EngineError> {
289 let source = self.built(&deliverer.source)?;
290 source
291 .source()
292 .get_task(&deliverer.native)
293 .await
294 .map(|task| task.map(|task| task.status.category))
295 .map_err(|error| source_failed(source, error))
296 }
297
298 pub(super) fn built(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
300 let name = self.known(name)?;
301 match self.sources.iter().find(|source| source.name() == &name) {
302 Some(ConfiguredSource::Ready(source)) => Ok(source),
303 Some(ConfiguredSource::Unavailable(source)) => Err(EngineError::SourceUnavailable {
304 name: name.to_string(),
305 error: source.error().clone(),
306 }),
307 None => Err(EngineError::NoSources),
309 }
310 }
311
312 fn status_writable(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
314 let source = self.built(name)?;
315 if source.source().writes().is_supported() {
316 return Ok(source);
317 }
318 Err(EngineError::StatusNotWritable {
319 name: source.name().to_string(),
320 kind: source.kind().to_owned(),
321 })
322 }
323}
324
325#[must_use]
339pub fn settled(categories: &[StatusCategory]) -> Option<StatusCategory> {
340 use StatusCategory::{Backlog, Cancelled, Done, Draft, InProgress, Queued, Todo, Unknown};
341 if categories.contains(&Done)
342 && categories
343 .iter()
344 .all(|category| matches!(category, Done | Cancelled))
345 {
346 return Some(Done);
347 }
348 if categories.contains(&InProgress) {
349 return Some(InProgress);
350 }
351 if categories.contains(&Queued) {
352 return Some(Queued);
353 }
354 if categories
355 .iter()
356 .any(|category| matches!(category, Todo | Cancelled | Unknown | Done))
357 {
358 return Some(Todo);
359 }
360 if categories.is_empty() {
361 return Some(Todo);
362 }
363 debug_assert!(
364 categories
365 .iter()
366 .all(|category| matches!(category, Draft | Backlog))
367 );
368 None
369}
370
371pub(crate) fn targets(list: &[TaskRef], near: &SourceName) -> Vec<GlobalId> {
374 list.iter()
375 .filter_map(|entry| entry.in_source(near).as_str().parse().ok())
376 .collect()
377}
378
379pub(crate) fn qualified_task(id: GlobalId, task: Task) -> Qualified<Task> {
382 let delivers = task
383 .delivers
384 .iter()
385 .map(|entry| entry.in_source(&id.source))
386 .collect();
387 let delivered_by = task
388 .delivered_by
389 .iter()
390 .map(|entry| entry.in_source(&id.source))
391 .collect();
392 Qualified {
393 id,
394 item: Task {
395 delivers,
396 delivered_by,
397 ..task
398 },
399 }
400}
401
402pub(super) fn source_failed(source: &ResolvedSource, error: SourceError) -> EngineError {
403 EngineError::SourceFailed {
404 name: source.name().to_string(),
405 error,
406 }
407}