Skip to main content

onetaskgraph_core/engine/
comment.rs

1//! Comments on a task: the four verbs that read and write them, and the task detail
2//! `task show` renders them in.
3//!
4//! Every verb here addresses exactly one task of exactly one source, so none of them fans
5//! out, pages for a caller or compensates for anything. What they owe instead is the
6//! refusals: a source declaring no comments is refused before anything is read, a source
7//! whose comments cannot be written is refused before a write is attempted, and a task or a
8//! comment that is not there is named rather than answered with an empty result.
9//!
10//! Nothing here writes anything down. A list walks the source's pages to the end and hands
11//! the caller exactly what it read; a copy never reaches this module at all, which is what
12//! keeps a copy from reading or writing a comment at either end.
13
14use chrono::{DateTime, Utc};
15use onetaskgraph_plugin_api::{
16    Comment, CommentBody, Cursor, NativeId, NewComment, Page, PageRequest, SourceError, SourceName,
17    Task, TaskDetailRead, TaskQuery,
18};
19use schemars::JsonSchema;
20use serde::{Deserialize, Serialize};
21
22use super::fetch::{fits, unrepeated};
23use super::{Answer, ConfiguredSource, Engine, EngineError, Qualified, delivery};
24use crate::GlobalId;
25use crate::plan::{QueryPlan, QueryResponse, SourceFailure};
26use crate::resolve::ResolvedSource;
27
28/// Every comment on one task, oldest first: what `task comment list` answers with.
29///
30/// An object rather than a bare list, so a later member — a total, say — is an addition a
31/// reader already written against this shape can ignore.
32#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
33pub struct CommentList {
34    /// The task's comments in the order they were written. Empty when it has none.
35    pub comments: Vec<Comment>,
36}
37
38/// What `task comment delete` answers with: the id of the comment it removed.
39#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
40pub struct DeletedComment {
41    /// The id the comment was removed under, exactly as `list` reported it.
42    pub deleted: NativeId,
43}
44
45/// One task as `task show` reports it: the response every show verb answers with, and the
46/// task's comments beside it.
47///
48/// The response is flattened rather than nested, so a reader of `task show --json` written
49/// before comments existed reads exactly the members it read before.
50#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
51pub struct TaskDetail {
52    /// The task, its plan and any failure, exactly as [`Engine::task`] answers.
53    #[serde(flatten)]
54    pub response: QueryResponse<Qualified<Task>>,
55    /// The task's comments, oldest first, for a source whose tasks have comments.
56    ///
57    /// **Absent** rather than empty for a source declaring none, for a task that was not
58    /// found, and for a task whose comments could not be read — the last with the failure in
59    /// the response's `errors`, a source refusing the read (a GitHub draft, which has none)
60    /// included, so showing such a task is a partial answer that says why. An empty list says
61    /// the source has comments and this task holds none, which is a different thing to tell a
62    /// reader.
63    #[serde(default, skip_serializing_if = "Option::is_none")]
64    pub comments: Option<Vec<Comment>>,
65}
66
67/// Several tasks as `task show-many` reports them: one [`TaskDetail`] per id asked for, in
68/// the order they were asked for.
69///
70/// An object rather than a bare list, for the reason [`CommentList`] is one.
71#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
72pub struct TaskDetails {
73    /// One detail per id, in request order — the document `task show <ID>` answers for that
74    /// id, except that an id it would refuse outright, or answer with nothing, carries why
75    /// in that detail's own `errors` instead, so it never refuses the others.
76    pub details: Vec<TaskDetail>,
77}
78
79impl Engine {
80    /// One task by its qualified id, with its comments when its source has them.
81    ///
82    /// The task is read exactly as [`task`](Self::task) reads it, and the comments are read
83    /// only once the task was found — a source declaring no comments is never asked. Both
84    /// halves go through the source's one [`get_task_details`] call, so a source that reads
85    /// an item and its first page of comments in one request answers in one.
86    ///
87    /// [`get_task_details`]: onetaskgraph_plugin_api::TaskSource::get_task_details
88    ///
89    /// # Errors
90    ///
91    /// As [`task`](Self::task). A comment read that fails is not an error: it lands in the
92    /// response's `errors` beside the task that was read.
93    pub async fn task_detail(&self, id: &GlobalId) -> Result<TaskDetail, EngineError> {
94        let name = self.known(&id.source)?;
95        let mut details = self
96            .source_details(&name, std::slice::from_ref(id), true)
97            .await;
98        Ok(details.remove(0))
99    }
100
101    /// Several tasks by their qualified ids, each with its comments when `comments` is set and
102    /// its source has them — the ids may span sources.
103    ///
104    /// Each source is asked once, for every id naming it, through its
105    /// [`get_task_details`](onetaskgraph_plugin_api::TaskSource::get_task_details); what each
106    /// detail holds is what [`task_detail`](Self::task_detail), or [`task`](Self::task) when
107    /// `comments` is unset, answers for that id. An id that answer would refuse — it names no
108    /// configured source — or answer with nothing — its source holds no such task — carries
109    /// that in its own detail's `errors`, so one unreadable id never refuses the others, and a
110    /// caller reads every failure from the same place.
111    pub async fn task_details(&self, ids: &[GlobalId], comments: bool) -> TaskDetails {
112        let mut details: Vec<Option<TaskDetail>> = vec![None; ids.len()];
113        let mut named: Vec<&SourceName> = Vec::new();
114        for id in ids {
115            if !named.contains(&&id.source) {
116                named.push(&id.source);
117            }
118        }
119        for name in named {
120            let at: Vec<usize> = (0..ids.len())
121                .filter(|index| ids[*index].source == *name)
122                .collect();
123            let asked: Vec<GlobalId> = at.iter().map(|index| ids[*index].clone()).collect();
124            let answered = match self.known(name) {
125                Ok(name) => self.source_details(&name, &asked, comments).await,
126                Err(refusal) => asked
127                    .iter()
128                    .map(|_| TaskDetail {
129                        response: failed_response(
130                            name,
131                            SourceError::Config {
132                                message: refusal.to_string(),
133                            },
134                        ),
135                        comments: None,
136                    })
137                    .collect(),
138            };
139            for ((index, id), mut detail) in at.into_iter().zip(&asked).zip(answered) {
140                if detail.response.items.is_empty() && detail.response.errors.is_empty() {
141                    detail.response.errors.push(SourceFailure {
142                        source: id.source.clone(),
143                        error: SourceError::Refused {
144                            message: no_such_task(id).to_string(),
145                        },
146                    });
147                }
148                details[index] = Some(detail);
149            }
150        }
151        TaskDetails {
152            details: details.into_iter().flatten().collect(),
153        }
154    }
155
156    /// One detail per id of `ids`, every one of which names the configured source `name`, read
157    /// with one call of that source.
158    async fn source_details(
159        &self,
160        name: &SourceName,
161        ids: &[GlobalId],
162        comments: bool,
163    ) -> Vec<TaskDetail> {
164        let mut answer = Answer::new();
165        let selected = answer.split(self, std::slice::from_ref(name));
166        let Some(source) = selected.first() else {
167            // A source that never built: every id is answered with the failure `task` answers
168            // with, and nothing is asked.
169            return ids
170                .iter()
171                .map(|_| TaskDetail {
172                    response: QueryResponse {
173                        items: Vec::new(),
174                        next: None,
175                        plan: QueryPlan::default(),
176                        errors: answer.errors.clone(),
177                    },
178                    comments: None,
179                })
180                .collect();
181        };
182        let commented = comments && source.source().capabilities().comments.is_native();
183        let page = PageRequest {
184            cursor: None,
185            limit: source.source().capabilities().max_page_size.max(1),
186        };
187        let natives: Vec<NativeId> = ids.iter().map(|id| id.native.clone()).collect();
188        let mut read = source
189            .source()
190            .get_task_details(&natives, commented.then_some(&page))
191            .await;
192        if read.len() != ids.len() {
193            let message = format!(
194                "the source answered {} task details for the {} ids it was asked for",
195                read.len(),
196                ids.len()
197            );
198            read = ids
199                .iter()
200                .map(|_| {
201                    Err(SourceError::Malformed {
202                        message: message.clone(),
203                    })
204                })
205                .collect();
206        }
207        let mut details = Vec::with_capacity(ids.len());
208        for (id, result) in ids.iter().zip(read) {
209            let (found, first) = match result {
210                Ok(Some(TaskDetailRead { task, comments })) => (Ok(Some(task)), comments),
211                Ok(None) => (Ok(None), None),
212                Err(error) => (Err(error), None),
213            };
214            let qualified = GlobalId::new(source.name().clone(), id.native.clone());
215            let mut response = Answer::new().one_response(source, found, |task| {
216                delivery::qualified_task(qualified, task)
217            });
218            let comments = if commented && !response.items.is_empty() {
219                let walked = match first {
220                    Some(Ok(Some(first))) => walk_from(source, &id.native, first, page.limit).await,
221                    Some(Ok(None)) => Ok(None),
222                    Some(Err(error)) => Err(error),
223                    // A source that was asked for the first page and answered none of it is
224                    // asked the way a source of one comment page at a time always is.
225                    None => walk(source, &id.native).await,
226                };
227                match walked {
228                    Ok(comments) => comments,
229                    Err(error) => {
230                        response.errors.push(SourceFailure {
231                            source: source.name().clone(),
232                            error,
233                        });
234                        None
235                    }
236                }
237            } else {
238                None
239            };
240            details.push(TaskDetail { response, comments });
241        }
242        details
243    }
244
245    /// Every comment on one task, oldest first.
246    ///
247    /// # Errors
248    ///
249    /// Returns [`EngineError::NoComments`] for a source whose tasks have none, before
250    /// anything is read; [`EngineError::NoSuchTask`] when the source holds no such task; and
251    /// [`EngineError::SourceFailed`] when the source could not answer.
252    pub async fn comments(&self, task: &GlobalId) -> Result<CommentList, EngineError> {
253        let source = self.commented(&task.source)?;
254        match walk(source, &task.native).await {
255            Ok(Some(comments)) => Ok(CommentList { comments }),
256            Ok(None) => Err(no_such_task(task)),
257            Err(error) => Err(failed(source, error)),
258        }
259    }
260
261    /// Add one comment to a task, answering with the comment as its source now holds it.
262    ///
263    /// # Errors
264    ///
265    /// As [`comments`](Self::comments), plus [`EngineError::CommentsNotWritable`] for a
266    /// source whose comments cannot be written, before anything is written.
267    pub async fn add_comment(
268        &self,
269        task: &GlobalId,
270        comment: &NewComment,
271    ) -> Result<Comment, EngineError> {
272        let source = self.writable_comments(&task.source)?;
273        match source.source().add_comment(&task.native, comment).await {
274            Ok(Some(added)) => Ok(added),
275            Ok(None) => Err(no_such_task(task)),
276            Err(error) => Err(failed(source, error)),
277        }
278    }
279
280    /// Replace one comment's body, answering with the comment as its source now holds it.
281    ///
282    /// # Errors
283    ///
284    /// As [`add_comment`](Self::add_comment), plus [`EngineError::NoSuchComment`] when the
285    /// task has no comment under `comment`.
286    pub async fn edit_comment(
287        &self,
288        task: &GlobalId,
289        comment: &NativeId,
290        body: &CommentBody,
291    ) -> Result<Comment, EngineError> {
292        let source = self.writable_comments(&task.source)?;
293        match source
294            .source()
295            .edit_comment(&task.native, comment, body)
296            .await
297        {
298            Ok(Some(edited)) => Ok(edited),
299            Ok(None) => Err(missing(source, task, comment).await),
300            Err(error) => Err(failed(source, error)),
301        }
302    }
303
304    /// Remove one comment from a task, answering with the id it removed.
305    ///
306    /// # Errors
307    ///
308    /// As [`edit_comment`](Self::edit_comment).
309    pub async fn delete_comment(
310        &self,
311        task: &GlobalId,
312        comment: &NativeId,
313    ) -> Result<DeletedComment, EngineError> {
314        let source = self.writable_comments(&task.source)?;
315        match source.source().delete_comment(&task.native, comment).await {
316            Ok(Some(deleted)) => Ok(DeletedComment { deleted }),
317            Ok(None) => Err(missing(source, task, comment).await),
318            Err(error) => Err(failed(source, error)),
319        }
320    }
321
322    /// The configured source called `name`, in whichever state it is in.
323    fn configured(&self, name: &SourceName) -> Option<&ConfiguredSource> {
324        self.sources.iter().find(|source| source.name() == name)
325    }
326
327    /// The built source called `name`, when its tasks have comments.
328    fn commented(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
329        let name = self.known(name)?;
330        match self.configured(&name) {
331            Some(ConfiguredSource::Ready(source)) => {
332                if source.source().capabilities().comments.is_native() {
333                    Ok(source)
334                } else {
335                    Err(EngineError::NoComments {
336                        name: name.to_string(),
337                        kind: source.kind().to_owned(),
338                    })
339                }
340            }
341            Some(ConfiguredSource::Unavailable(source)) => Err(EngineError::SourceUnavailable {
342                name: name.to_string(),
343                error: source.error().clone(),
344            }),
345            // `known` has just said a source by this name is configured, and `configured`
346            // reads the same list it read.
347            None => Err(EngineError::NoSources),
348        }
349    }
350
351    /// The built source called `name`, when its tasks have comments it can write.
352    fn writable_comments(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
353        let source = self.commented(name)?;
354        if source.source().writes().is_supported() {
355            return Ok(source);
356        }
357        Err(EngineError::CommentsNotWritable {
358            name: source.name().to_string(),
359            kind: source.kind().to_owned(),
360        })
361    }
362}
363
364/// Every comment on `task`, walked to the end of the source's pages, or `None` when the
365/// source holds no such task.
366///
367/// Each page is asked at the source's own ceiling, and each is held to the two refusals every
368/// pagination loop of this engine owes: a page longer than the one asked for, and a cursor
369/// handed back unchanged. What the walk accumulates is the caller's answer and nothing else.
370async fn walk(
371    source: &ResolvedSource,
372    task: &NativeId,
373) -> Result<Option<Vec<Comment>>, SourceError> {
374    let limit = source.source().capabilities().max_page_size.max(1);
375    let request = PageRequest {
376        cursor: None,
377        limit,
378    };
379    let Some(first) = source.source().task_comments(task, &request).await? else {
380        return Ok(None);
381    };
382    walk_from(source, task, first, limit).await
383}
384
385/// [`walk`], from a first page already in hand — the one a detail read carried beside the
386/// task — asked at `limit`.
387async fn walk_from(
388    source: &ResolvedSource,
389    task: &NativeId,
390    first: Page<Comment>,
391    limit: u32,
392) -> Result<Option<Vec<Comment>>, SourceError> {
393    let mut comments = Vec::new();
394    let mut cursor: Option<Cursor> = None;
395    let mut page = first;
396    loop {
397        fits(page.items.len(), limit)?;
398        unrepeated(
399            page.next.as_ref(),
400            cursor.as_ref(),
401            "walking a task's comments",
402        )?;
403        comments.extend(page.items);
404        let Some(next) = page.next else {
405            return Ok(Some(comments));
406        };
407        cursor = Some(next);
408        let request = PageRequest {
409            cursor: cursor.clone(),
410            limit,
411        };
412        // A task that is gone part way through a walk is gone: reporting the comments read
413        // before it went would describe a task nobody can address any more.
414        let Some(read) = source.source().task_comments(task, &request).await? else {
415            return Ok(None);
416        };
417        page = read;
418    }
419}
420
421/// Narrow one page of a source's tasks to those with a comment created or last edited at or
422/// after `since`, for a source that does not apply that predicate itself.
423///
424/// This is the one predicate the engine cannot answer from the row: a task does not carry its
425/// comments. So each task the other predicates left to the engine already keep — `local` —
426/// has its comments read, page by page, until one matches or they run out; a task every
427/// other predicate drops is never asked about. What is held is one task's page of comments
428/// at a time, and nothing of it outlives the call. A source whose tasks have no comments at
429/// all holds no comment activity, so none of its tasks is kept and it is asked nothing, and a
430/// task gone from the source by the time its comments are read is gone from the answer.
431///
432/// # Errors
433///
434/// Returns whatever the source returned for a comment read, and the two refusals every
435/// pagination loop of this engine owes.
436pub(super) async fn commented_since(
437    source: &ResolvedSource,
438    local: &super::local::LocalTasks,
439    page: Page<Task>,
440    since: DateTime<Utc>,
441) -> Result<Page<Task>, SourceError> {
442    if !source.source().capabilities().comments.is_native() {
443        return Ok(Page {
444            items: Vec::new(),
445            next: page.next,
446        });
447    }
448    let query = TaskQuery {
449        commented_since: Some(since),
450        ..TaskQuery::default()
451    };
452    let mut kept = Vec::new();
453    for task in page.items {
454        if local.keeps(&task) && any_comment_matches(source, &task.id, &query).await? {
455            kept.push(task);
456        }
457    }
458    Ok(Page {
459        items: kept,
460        next: page.next,
461    })
462}
463
464/// Whether one of `task`'s comments satisfies `query`'s comment activity, walking the source's
465/// comment pages only as far as the first that does.
466async fn any_comment_matches(
467    source: &ResolvedSource,
468    task: &NativeId,
469    query: &TaskQuery,
470) -> Result<bool, SourceError> {
471    let limit = source.source().capabilities().max_page_size.max(1);
472    let mut cursor: Option<Cursor> = None;
473    loop {
474        let request = PageRequest {
475            cursor: cursor.clone(),
476            limit,
477        };
478        let Some(page) = source.source().task_comments(task, &request).await? else {
479            return Ok(false);
480        };
481        fits(page.items.len(), limit)?;
482        unrepeated(
483            page.next.as_ref(),
484            cursor.as_ref(),
485            "walking a task's comments",
486        )?;
487        if query.comments_match(&page.items) {
488            return Ok(true);
489        }
490        match page.next {
491            Some(next) => cursor = Some(next),
492            None => return Ok(false),
493        }
494    }
495}
496
497/// Which of the two things an edit or a delete named was not there.
498///
499/// Asked only once the source has already said one of them is missing, so a comment verb
500/// that succeeds costs its source exactly one call.
501async fn missing(source: &ResolvedSource, task: &GlobalId, comment: &NativeId) -> EngineError {
502    match source.source().get_task(&task.native).await {
503        Ok(Some(_)) => EngineError::NoSuchComment {
504            task: task.to_string(),
505            comment: comment.to_string(),
506        },
507        Ok(None) => no_such_task(task),
508        Err(error) => failed(source, error),
509    }
510}
511
512/// A response for an id the engine could not ask any source about, carrying why.
513fn failed_response(source: &SourceName, error: SourceError) -> QueryResponse<Qualified<Task>> {
514    QueryResponse {
515        items: Vec::new(),
516        next: None,
517        plan: QueryPlan::default(),
518        errors: vec![SourceFailure {
519            source: source.clone(),
520            error,
521        }],
522    }
523}
524
525fn no_such_task(task: &GlobalId) -> EngineError {
526    EngineError::NoSuchTask {
527        id: task.to_string(),
528    }
529}
530
531fn failed(source: &ResolvedSource, error: SourceError) -> EngineError {
532    EngineError::SourceFailed {
533        name: source.name().to_string(),
534        error,
535    }
536}