Skip to main content

onetaskgraph_core/engine/
delivery.rs

1//! A task's status set on its own, and the `delivers` relation the store keeps in step.
2//!
3//! A task names the tasks it **delivers**: finishing it finishes them. Whenever this engine
4//! writes a task whose `delivers` is not empty, or was not before the write, it keeps two
5//! things true of every task that list names now and every task it dropped:
6//!
7//! 1. **The back-reference.** A delivered task's `delivered_by` holds the deliverer's
8//!    qualified id while the deliverer names it, and does not once the deliverer drops it.
9//! 2. **The status.** The delivered task is re-evaluated over every deliverer it names, by
10//!    [`settled`], and its status is written through the status-only write when it is at
11//!    `todo`, `queued` or `in-progress` and the result differs from what it holds.
12//!
13//! Every write of such a task re-evaluates, a write that changes nothing about the deliverer
14//! included, so a retried write re-fires the rule. Nothing here is written down outside the
15//! plugins: `delivered_by` lives on the delivered task, in its own source, and every
16//! evaluation reads the deliverers afresh.
17
18use 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/// What `task status set` answers with.
27#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
28pub struct TaskStatusSet {
29    /// The task whose status was set.
30    pub id: GlobalId,
31    /// Its status as its source reads it back.
32    pub status: Status,
33    /// One entry per task it delivers, or dropped, that this write re-evaluated. Empty when
34    /// it delivers nothing.
35    pub delivered: Vec<Delivered>,
36}
37
38/// What keeping one delivered task in step with one deliverer came to.
39#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
40pub struct Delivered {
41    /// The delivered task.
42    pub ticket: GlobalId,
43    /// The deliverer whose write re-evaluated it — the one that just dropped it, when it was
44    /// re-evaluated for being dropped.
45    pub deliverer: GlobalId,
46    /// What happened to it.
47    #[serde(flatten)]
48    pub outcome: DeliveryOutcome,
49    /// Deliverers its source read as not found, removed from its `delivered_by` on this
50    /// write. Left out when there were none.
51    #[serde(default, skip_serializing_if = "Vec::is_empty")]
52    #[schemars(!skip_serializing_if)]
53    pub pruned: Vec<GlobalId>,
54}
55
56impl Delivered {
57    /// Whether the rule could not keep this task in step, which is what exits `4`.
58    #[must_use]
59    pub fn failed(&self) -> bool {
60        matches!(self.outcome, DeliveryOutcome::Failed { .. })
61    }
62}
63
64/// The four things re-evaluating a delivered task can come to, and the categories each is
65/// about.
66#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
67#[serde(tag = "outcome", rename_all = "kebab-case")]
68pub enum DeliveryOutcome {
69    /// Its status was written.
70    Written {
71        /// The category it read.
72        from: StatusCategory,
73        /// The category written.
74        to: StatusCategory,
75    },
76    /// The rule asked for what it already holds, or for nothing at all.
77    Unchanged {
78        /// The category it read.
79        from: StatusCategory,
80    },
81    /// It is at `draft`, `backlog`, `unknown`, `done` or `cancelled`, which a claim never
82    /// accepts, reopens or un-defers on a person's behalf.
83    Left {
84        /// The category it read.
85        from: StatusCategory,
86    },
87    /// It, or a deliverer it names, could not be read or written; the rule did not guess.
88    Failed {
89        /// The category it read, when it could be read.
90        #[serde(default, skip_serializing_if = "Option::is_none")]
91        from: Option<StatusCategory>,
92        /// Why, as the failure object itself: the same `class`, `kind`, `source`, `message`
93        /// and `retry_after_seconds` a failure document carries under its own `failure`
94        /// member — not that whole document nested again.
95        failure: Failure,
96    },
97}
98
99impl Engine {
100    /// Set one task's status, and nothing else about it, then keep every task it delivers in
101    /// step with it.
102    ///
103    /// # Errors
104    ///
105    /// Returns [`EngineError::UnknownSource`] for a source nothing configures,
106    /// [`EngineError::StatusNotWritable`] for one with no write side,
107    /// [`EngineError::NoSuchTask`] when the task is not there, and
108    /// [`EngineError::SourceFailed`] when the source refuses — a category it has disabled
109    /// included. A delivered task that cannot be kept in step is not an error: it is reported
110    /// `failed` beside the status that was set.
111    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        // One call answering the task as it now reads, rather than a read and then a write: a
119        // source whose status write answers the whole task saves the read, and every other
120        // source makes exactly those two calls itself.
121        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    /// Keep every task `deliverer` delivers `now`, and every one it delivered `before` and
140    /// dropped, in step with it — one entry each, in that order.
141    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    /// Re-evaluate one delivered task for one deliverer, which names it when `kept` and has
167    /// just dropped it otherwise.
168    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        // Only a task the rule could write is worth reading its deliverers for: one it leaves
209        // alone is left alone whatever they say.
210        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        // A rule that never finished reading prunes nothing: what it read before the failure
234        // is half an answer, and the deliverer it could not read stays named.
235        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    /// A deliverer's category, or `None` when its source reads it as not found.
285    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    /// The built source called `name`.
299    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            // `known` has just said a source by this name is configured.
308            None => Err(EngineError::NoSources),
309        }
310    }
311
312    /// The built source called `name`, when it can be written through.
313    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/// What a delivered task's status should be, over the categories of every deliverer that
326/// remains — or `None` when it should not be written at all.
327///
328/// The first branch that matches decides:
329///
330/// 1. `done` when every deliverer is `done` or `cancelled` and at least one is `done`. A
331///    `cancelled` deliverer releases its claim and never counts as completion; it only stops
332///    blocking a sibling's.
333/// 2. `in-progress` when any deliverer is `in-progress`.
334/// 3. `queued` when any deliverer is `queued`.
335/// 4. `todo` when any deliverer releases — `todo`, `cancelled`, `unknown` — or is `done`.
336/// 5. `todo` when no deliverer remains: the claim is released.
337/// 6. `None` when every remaining deliverer is `draft` or `backlog`, which count as nothing.
338#[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
371/// Every entry of one task's list as a qualified id, reading a bare one as naming a task of
372/// `near`, the source holding the list.
373pub(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
379/// A task as a verb reports it: under its qualified id, with every entry of its two lists
380/// qualified too, so a reader never has to know which source a bare entry meant.
381pub(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}