Skip to main content

onetaskgraph_core/engine/
update.rs

1//! One targeted update of one existing task: every field the caller names, and nothing else.
2//!
3//! One call to one source. The engine reads nothing of its own around it: the source's answer
4//! carries the task as it reads back, the fields it wrote and the `delivers` it held before,
5//! and every figure this reports is read off that answer. On a hosted backend that is the
6//! difference between one read of the item and two, on every update — which is the whole of
7//! why this exists beside a copy and the narrow verbs.
8//!
9//! `delivers` is kept exactly as `task status set` keeps it: whenever `status` or `delivers` is
10//! named, every task the list names now, and every task it dropped, is re-evaluated.
11
12use std::collections::BTreeSet;
13
14use onetaskgraph_plugin_api::{
15    DependencyEdge, DependencyEndpoint, ItemKind, NativeId, SourceName, Task, TaskUpdate,
16    UpdatedField,
17};
18use schemars::JsonSchema;
19use serde::{Deserialize, Serialize};
20
21use super::copy::{Spent, readings, spent_between};
22use super::delivery::{Delivered, qualified_task, source_failed, targets};
23use super::narrow::holds_priority;
24use super::{Engine, EngineError};
25use crate::GlobalId;
26
27/// What `task update` answers with.
28#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
29pub struct TaskUpdated {
30    /// The task updated.
31    pub id: GlobalId,
32    /// The task as its source reads it back, with every entry of its `delivers` and
33    /// `delivered_by` qualified.
34    pub task: Task,
35    /// The fields the source actually wrote — empty when nothing differed.
36    pub written: BTreeSet<UpdatedField>,
37    /// As `set_task_status` reports it; re-evaluated whenever `status` or `delivers` was
38    /// named, empty otherwise.
39    pub delivered: Vec<Delivered>,
40    /// What the source's meter says this call spent, as `CopyReport::spent` — and absent,
41    /// never zero, when the source does not meter its own requests.
42    #[serde(default, skip_serializing_if = "Option::is_none")]
43    pub spent: Option<Spent>,
44}
45
46impl Engine {
47    /// Apply a targeted update to one task — every field `update` names, and nothing else —
48    /// then keep every task it delivers in step with it when `status` or `delivers` was named.
49    ///
50    /// The source is asked once, and nothing else of it is read: a field already holding the
51    /// requested value is sent no write, and an update in which nothing differs writes nothing
52    /// at all — an update naming no field included, which answers the task as its source reads
53    /// it and re-evaluates nothing; the CLI refuses that one as a usage error before it gets
54    /// here, a library caller is answered. `update.depends_on` may name its far ends qualified; one in this task's own
55    /// source reaches the source as its native id, and each edge's `from` is this task.
56    ///
57    /// # Errors
58    ///
59    /// Returns [`EngineError::UnknownSource`] for a source nothing configures,
60    /// [`EngineError::UpdateNotWritable`] for one with no write side,
61    /// [`EngineError::NoPriority`] for a priority other than `none` to a source that holds
62    /// none — neither of which is asked — [`EngineError::NoSuchTask`] when the task is not
63    /// there, and [`EngineError::SourceFailed`] carrying the source's own [`SourceError`] when
64    /// it refuses or fails, so its class and its `retry_after_seconds` are what a caller
65    /// reads. An update both setting and removing one metadata key is refused that way too,
66    /// before anything is sent, in the words the source itself refuses it with. A delivered
67    /// task that cannot be kept in step is not an error: it is reported `failed` beside the
68    /// update that landed.
69    ///
70    /// [`SourceError`]: onetaskgraph_plugin_api::SourceError
71    pub async fn update_task(
72        &self,
73        id: &GlobalId,
74        update: &TaskUpdate,
75    ) -> Result<TaskUpdated, EngineError> {
76        let source = self.built(&id.source)?;
77        if !source.source().writes().is_supported() {
78            return Err(EngineError::UpdateNotWritable {
79                name: source.name().to_string(),
80                kind: source.kind().to_owned(),
81            });
82        }
83        update
84            .consistent()
85            .map_err(|error| source_failed(source, error))?;
86        if let Some(priority) = update.priority {
87            holds_priority(source, &id.to_string(), priority)?;
88        }
89        let update = TaskUpdate {
90            depends_on: update
91                .depends_on
92                .as_ref()
93                .map(|edges| near_edges(edges, &id.native, &id.source)),
94            ..update.clone()
95        };
96        let before = readings(&[source]).await;
97        let outcome = source
98            .source()
99            .update_task(&id.native, &update)
100            .await
101            .map_err(|error| source_failed(source, error))?
102            .ok_or_else(|| EngineError::NoSuchTask { id: id.to_string() })?;
103        let delivered = if update.status.is_some() || update.delivers.is_some() {
104            let now = targets(&outcome.task.delivers, &id.source);
105            let dropped = targets(&outcome.delivers_before, &id.source);
106            self.deliver(id, outcome.task.status.category, &now, &dropped)
107                .await
108        } else {
109            Vec::new()
110        };
111        let spent = spent_between(&before, &readings(&[source]).await);
112        Ok(TaskUpdated {
113            id: id.clone(),
114            task: qualified_task(id.clone(), outcome.task).item,
115            written: outcome.written,
116            delivered,
117            spent,
118        })
119    }
120}
121
122/// `edges` as the task `near` of `source` holds them: each starting at `near`, and each far
123/// end in `source` itself named by its native id, which is how a source names one of its own.
124fn near_edges(
125    edges: &[DependencyEdge],
126    near: &NativeId,
127    source: &SourceName,
128) -> Vec<DependencyEdge> {
129    edges
130        .iter()
131        .map(|edge| {
132            let to = match edge.to.id().parse::<GlobalId>() {
133                Ok(far) if &far.source == source => {
134                    DependencyEndpoint::from_native(far.native, edge.to.kind)
135                }
136                _ => edge.to.clone(),
137            };
138            DependencyEdge {
139                from: DependencyEndpoint::from_native(near.clone(), ItemKind::Task),
140                to,
141                kind: edge.kind,
142            }
143        })
144        .collect()
145}