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}