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