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}