onetaskgraph-core 0.2.55

The onetaskgraph engine: the plugin registry, global-id qualification, and the plan every response carries.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
//! Comments on a task: the four verbs that read and write them, and the task detail
//! `task show` renders them in.
//!
//! Every verb here addresses exactly one task of exactly one source, so none of them fans
//! out, pages for a caller or compensates for anything. What they owe instead is the
//! refusals: a source declaring no comments is refused before anything is read, a source
//! whose comments cannot be written is refused before a write is attempted, and a task or a
//! comment that is not there is named rather than answered with an empty result.
//!
//! Nothing here writes anything down. A list walks the source's pages to the end and hands
//! the caller exactly what it read; a copy never reaches this module at all, which is what
//! keeps a copy from reading or writing a comment at either end.

use chrono::{DateTime, Utc};
use onetaskgraph_plugin_api::{
    Comment, CommentBody, Cursor, NativeId, NewComment, Page, PageRequest, SourceError, SourceName,
    Task, TaskDetailRead, TaskQuery,
};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};

use super::fetch::{fits, unrepeated};
use super::{Answer, ConfiguredSource, Engine, EngineError, Qualified, delivery};
use crate::GlobalId;
use crate::plan::{QueryPlan, QueryResponse, SourceFailure};
use crate::resolve::ResolvedSource;

/// Every comment on one task, oldest first: what `task comment list` answers with.
///
/// An object rather than a bare list, so a later member — a total, say — is an addition a
/// reader already written against this shape can ignore.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct CommentList {
    /// The task's comments in the order they were written. Empty when it has none.
    pub comments: Vec<Comment>,
}

/// What `task comment delete` answers with: the id of the comment it removed.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct DeletedComment {
    /// The id the comment was removed under, exactly as `list` reported it.
    pub deleted: NativeId,
}

/// One task as `task show` reports it: the response every show verb answers with, and the
/// task's comments beside it.
///
/// The response is flattened rather than nested, so a reader of `task show --json` written
/// before comments existed reads exactly the members it read before.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct TaskDetail {
    /// The task, its plan and any failure, exactly as [`Engine::task`] answers.
    #[serde(flatten)]
    pub response: QueryResponse<Qualified<Task>>,
    /// The task's comments, oldest first, for a source whose tasks have comments.
    ///
    /// **Absent** rather than empty for a source declaring none, for a task that was not
    /// found, and for a task whose comments could not be read — the last with the failure in
    /// the response's `errors`, a source refusing the read (a GitHub draft, which has none)
    /// included, so showing such a task is a partial answer that says why. An empty list says
    /// the source has comments and this task holds none, which is a different thing to tell a
    /// reader.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub comments: Option<Vec<Comment>>,
}

/// Several tasks as `task show-many` reports them: one [`TaskDetail`] per id asked for, in
/// the order they were asked for.
///
/// An object rather than a bare list, for the reason [`CommentList`] is one.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct TaskDetails {
    /// One detail per id, in request order — the document `task show <ID>` answers for that
    /// id, except that an id it would refuse outright, or answer with nothing, carries why
    /// in that detail's own `errors` instead, so it never refuses the others.
    pub details: Vec<TaskDetail>,
}

impl Engine {
    /// One task by its qualified id, with its comments when its source has them.
    ///
    /// The task is read exactly as [`task`](Self::task) reads it, and the comments are read
    /// only once the task was found — a source declaring no comments is never asked. Both
    /// halves go through the source's one [`get_task_details`] call, so a source that reads
    /// an item and its first page of comments in one request answers in one.
    ///
    /// [`get_task_details`]: onetaskgraph_plugin_api::TaskSource::get_task_details
    ///
    /// # Errors
    ///
    /// As [`task`](Self::task). A comment read that fails is not an error: it lands in the
    /// response's `errors` beside the task that was read.
    pub async fn task_detail(&self, id: &GlobalId) -> Result<TaskDetail, EngineError> {
        let name = self.known(&id.source)?;
        let mut details = self
            .source_details(&name, std::slice::from_ref(id), true)
            .await;
        Ok(details.remove(0))
    }

    /// Several tasks by their qualified ids, each with its comments when `comments` is set and
    /// its source has them — the ids may span sources.
    ///
    /// Each source is asked once, for every id naming it, through its
    /// [`get_task_details`](onetaskgraph_plugin_api::TaskSource::get_task_details); what each
    /// detail holds is what [`task_detail`](Self::task_detail), or [`task`](Self::task) when
    /// `comments` is unset, answers for that id. An id that answer would refuse — it names no
    /// configured source — or answer with nothing — its source holds no such task — carries
    /// that in its own detail's `errors`, so one unreadable id never refuses the others, and a
    /// caller reads every failure from the same place.
    pub async fn task_details(&self, ids: &[GlobalId], comments: bool) -> TaskDetails {
        let mut details: Vec<Option<TaskDetail>> = vec![None; ids.len()];
        let mut named: Vec<&SourceName> = Vec::new();
        for id in ids {
            if !named.contains(&&id.source) {
                named.push(&id.source);
            }
        }
        for name in named {
            let at: Vec<usize> = (0..ids.len())
                .filter(|index| ids[*index].source == *name)
                .collect();
            let asked: Vec<GlobalId> = at.iter().map(|index| ids[*index].clone()).collect();
            let answered = match self.known(name) {
                Ok(name) => self.source_details(&name, &asked, comments).await,
                Err(refusal) => asked
                    .iter()
                    .map(|_| TaskDetail {
                        response: failed_response(
                            name,
                            SourceError::Config {
                                message: refusal.to_string(),
                            },
                        ),
                        comments: None,
                    })
                    .collect(),
            };
            for ((index, id), mut detail) in at.into_iter().zip(&asked).zip(answered) {
                if detail.response.items.is_empty() && detail.response.errors.is_empty() {
                    detail.response.errors.push(SourceFailure {
                        source: id.source.clone(),
                        error: SourceError::Refused {
                            message: no_such_task(id).to_string(),
                        },
                    });
                }
                details[index] = Some(detail);
            }
        }
        TaskDetails {
            details: details.into_iter().flatten().collect(),
        }
    }

    /// One detail per id of `ids`, every one of which names the configured source `name`, read
    /// with one call of that source.
    async fn source_details(
        &self,
        name: &SourceName,
        ids: &[GlobalId],
        comments: bool,
    ) -> Vec<TaskDetail> {
        let mut answer = Answer::new();
        let selected = answer.split(self, std::slice::from_ref(name));
        let Some(source) = selected.first() else {
            // A source that never built: every id is answered with the failure `task` answers
            // with, and nothing is asked.
            return ids
                .iter()
                .map(|_| TaskDetail {
                    response: QueryResponse {
                        items: Vec::new(),
                        next: None,
                        plan: QueryPlan::default(),
                        errors: answer.errors.clone(),
                    },
                    comments: None,
                })
                .collect();
        };
        let commented = comments && source.source().capabilities().comments.is_native();
        let page = PageRequest {
            cursor: None,
            limit: source.source().capabilities().max_page_size.max(1),
        };
        let natives: Vec<NativeId> = ids.iter().map(|id| id.native.clone()).collect();
        let mut read = source
            .source()
            .get_task_details(&natives, commented.then_some(&page))
            .await;
        if read.len() != ids.len() {
            let message = format!(
                "the source answered {} task details for the {} ids it was asked for",
                read.len(),
                ids.len()
            );
            read = ids
                .iter()
                .map(|_| {
                    Err(SourceError::Malformed {
                        message: message.clone(),
                    })
                })
                .collect();
        }
        let mut details = Vec::with_capacity(ids.len());
        for (id, result) in ids.iter().zip(read) {
            let (found, first) = match result {
                Ok(Some(TaskDetailRead { task, comments })) => (Ok(Some(task)), comments),
                Ok(None) => (Ok(None), None),
                Err(error) => (Err(error), None),
            };
            let qualified = GlobalId::new(source.name().clone(), id.native.clone());
            let mut response = Answer::new().one_response(source, found, |task| {
                delivery::qualified_task(qualified, task)
            });
            let comments = if commented && !response.items.is_empty() {
                let walked = match first {
                    Some(Ok(Some(first))) => walk_from(source, &id.native, first, page.limit).await,
                    Some(Ok(None)) => Ok(None),
                    Some(Err(error)) => Err(error),
                    // A source that was asked for the first page and answered none of it is
                    // asked the way a source of one comment page at a time always is.
                    None => walk(source, &id.native).await,
                };
                match walked {
                    Ok(comments) => comments,
                    Err(error) => {
                        response.errors.push(SourceFailure {
                            source: source.name().clone(),
                            error,
                        });
                        None
                    }
                }
            } else {
                None
            };
            details.push(TaskDetail { response, comments });
        }
        details
    }

    /// Every comment on one task, oldest first.
    ///
    /// # Errors
    ///
    /// Returns [`EngineError::NoComments`] for a source whose tasks have none, before
    /// anything is read; [`EngineError::NoSuchTask`] when the source holds no such task; and
    /// [`EngineError::SourceFailed`] when the source could not answer.
    pub async fn comments(&self, task: &GlobalId) -> Result<CommentList, EngineError> {
        let source = self.commented(&task.source)?;
        match walk(source, &task.native).await {
            Ok(Some(comments)) => Ok(CommentList { comments }),
            Ok(None) => Err(no_such_task(task)),
            Err(error) => Err(failed(source, error)),
        }
    }

    /// Add one comment to a task, answering with the comment as its source now holds it.
    ///
    /// # Errors
    ///
    /// As [`comments`](Self::comments), plus [`EngineError::CommentsNotWritable`] for a
    /// source whose comments cannot be written, before anything is written.
    pub async fn add_comment(
        &self,
        task: &GlobalId,
        comment: &NewComment,
    ) -> Result<Comment, EngineError> {
        let source = self.writable_comments(&task.source)?;
        match source.source().add_comment(&task.native, comment).await {
            Ok(Some(added)) => Ok(added),
            Ok(None) => Err(no_such_task(task)),
            Err(error) => Err(failed(source, error)),
        }
    }

    /// Replace one comment's body, answering with the comment as its source now holds it.
    ///
    /// # Errors
    ///
    /// As [`add_comment`](Self::add_comment), plus [`EngineError::NoSuchComment`] when the
    /// task has no comment under `comment`.
    pub async fn edit_comment(
        &self,
        task: &GlobalId,
        comment: &NativeId,
        body: &CommentBody,
    ) -> Result<Comment, EngineError> {
        let source = self.writable_comments(&task.source)?;
        match source
            .source()
            .edit_comment(&task.native, comment, body)
            .await
        {
            Ok(Some(edited)) => Ok(edited),
            Ok(None) => Err(missing(source, task, comment).await),
            Err(error) => Err(failed(source, error)),
        }
    }

    /// Remove one comment from a task, answering with the id it removed.
    ///
    /// # Errors
    ///
    /// As [`edit_comment`](Self::edit_comment).
    pub async fn delete_comment(
        &self,
        task: &GlobalId,
        comment: &NativeId,
    ) -> Result<DeletedComment, EngineError> {
        let source = self.writable_comments(&task.source)?;
        match source.source().delete_comment(&task.native, comment).await {
            Ok(Some(deleted)) => Ok(DeletedComment { deleted }),
            Ok(None) => Err(missing(source, task, comment).await),
            Err(error) => Err(failed(source, error)),
        }
    }

    /// The configured source called `name`, in whichever state it is in.
    fn configured(&self, name: &SourceName) -> Option<&ConfiguredSource> {
        self.sources.iter().find(|source| source.name() == name)
    }

    /// The built source called `name`, when its tasks have comments.
    fn commented(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
        let name = self.known(name)?;
        match self.configured(&name) {
            Some(ConfiguredSource::Ready(source)) => {
                if source.source().capabilities().comments.is_native() {
                    Ok(source)
                } else {
                    Err(EngineError::NoComments {
                        name: name.to_string(),
                        kind: source.kind().to_owned(),
                    })
                }
            }
            Some(ConfiguredSource::Unavailable(source)) => Err(EngineError::SourceUnavailable {
                name: name.to_string(),
                error: source.error().clone(),
            }),
            // `known` has just said a source by this name is configured, and `configured`
            // reads the same list it read.
            None => Err(EngineError::NoSources),
        }
    }

    /// The built source called `name`, when its tasks have comments it can write.
    fn writable_comments(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
        let source = self.commented(name)?;
        if source.source().writes().is_supported() {
            return Ok(source);
        }
        Err(EngineError::CommentsNotWritable {
            name: source.name().to_string(),
            kind: source.kind().to_owned(),
        })
    }
}

/// Every comment on `task`, walked to the end of the source's pages, or `None` when the
/// source holds no such task.
///
/// Each page is asked at the source's own ceiling, and each is held to the two refusals every
/// pagination loop of this engine owes: a page longer than the one asked for, and a cursor
/// handed back unchanged. What the walk accumulates is the caller's answer and nothing else.
async fn walk(
    source: &ResolvedSource,
    task: &NativeId,
) -> Result<Option<Vec<Comment>>, SourceError> {
    let limit = source.source().capabilities().max_page_size.max(1);
    let request = PageRequest {
        cursor: None,
        limit,
    };
    let Some(first) = source.source().task_comments(task, &request).await? else {
        return Ok(None);
    };
    walk_from(source, task, first, limit).await
}

/// [`walk`], from a first page already in hand — the one a detail read carried beside the
/// task — asked at `limit`.
async fn walk_from(
    source: &ResolvedSource,
    task: &NativeId,
    first: Page<Comment>,
    limit: u32,
) -> Result<Option<Vec<Comment>>, SourceError> {
    let mut comments = Vec::new();
    let mut cursor: Option<Cursor> = None;
    let mut page = first;
    loop {
        fits(page.items.len(), limit)?;
        unrepeated(
            page.next.as_ref(),
            cursor.as_ref(),
            "walking a task's comments",
        )?;
        comments.extend(page.items);
        let Some(next) = page.next else {
            return Ok(Some(comments));
        };
        cursor = Some(next);
        let request = PageRequest {
            cursor: cursor.clone(),
            limit,
        };
        // A task that is gone part way through a walk is gone: reporting the comments read
        // before it went would describe a task nobody can address any more.
        let Some(read) = source.source().task_comments(task, &request).await? else {
            return Ok(None);
        };
        page = read;
    }
}

/// Narrow one page of a source's tasks to those with a comment created or last edited at or
/// after `since`, for a source that does not apply that predicate itself.
///
/// This is the one predicate the engine cannot answer from the row: a task does not carry its
/// comments. So each task the other predicates left to the engine already keep — `local` —
/// has its comments read, page by page, until one matches or they run out; a task every
/// other predicate drops is never asked about. What is held is one task's page of comments
/// at a time, and nothing of it outlives the call. A source whose tasks have no comments at
/// all holds no comment activity, so none of its tasks is kept and it is asked nothing, and a
/// task gone from the source by the time its comments are read is gone from the answer.
///
/// # Errors
///
/// Returns whatever the source returned for a comment read, and the two refusals every
/// pagination loop of this engine owes.
pub(super) async fn commented_since(
    source: &ResolvedSource,
    local: &super::local::LocalTasks,
    page: Page<Task>,
    since: DateTime<Utc>,
) -> Result<Page<Task>, SourceError> {
    if !source.source().capabilities().comments.is_native() {
        return Ok(Page {
            items: Vec::new(),
            next: page.next,
        });
    }
    let query = TaskQuery {
        commented_since: Some(since),
        ..TaskQuery::default()
    };
    let mut kept = Vec::new();
    for task in page.items {
        if local.keeps(&task) && any_comment_matches(source, &task.id, &query).await? {
            kept.push(task);
        }
    }
    Ok(Page {
        items: kept,
        next: page.next,
    })
}

/// Whether one of `task`'s comments satisfies `query`'s comment activity, walking the source's
/// comment pages only as far as the first that does.
async fn any_comment_matches(
    source: &ResolvedSource,
    task: &NativeId,
    query: &TaskQuery,
) -> Result<bool, SourceError> {
    let limit = source.source().capabilities().max_page_size.max(1);
    let mut cursor: Option<Cursor> = None;
    loop {
        let request = PageRequest {
            cursor: cursor.clone(),
            limit,
        };
        let Some(page) = source.source().task_comments(task, &request).await? else {
            return Ok(false);
        };
        fits(page.items.len(), limit)?;
        unrepeated(
            page.next.as_ref(),
            cursor.as_ref(),
            "walking a task's comments",
        )?;
        if query.comments_match(&page.items) {
            return Ok(true);
        }
        match page.next {
            Some(next) => cursor = Some(next),
            None => return Ok(false),
        }
    }
}

/// Which of the two things an edit or a delete named was not there.
///
/// Asked only once the source has already said one of them is missing, so a comment verb
/// that succeeds costs its source exactly one call.
async fn missing(source: &ResolvedSource, task: &GlobalId, comment: &NativeId) -> EngineError {
    match source.source().get_task(&task.native).await {
        Ok(Some(_)) => EngineError::NoSuchComment {
            task: task.to_string(),
            comment: comment.to_string(),
        },
        Ok(None) => no_such_task(task),
        Err(error) => failed(source, error),
    }
}

/// A response for an id the engine could not ask any source about, carrying why.
fn failed_response(source: &SourceName, error: SourceError) -> QueryResponse<Qualified<Task>> {
    QueryResponse {
        items: Vec::new(),
        next: None,
        plan: QueryPlan::default(),
        errors: vec![SourceFailure {
            source: source.clone(),
            error,
        }],
    }
}

fn no_such_task(task: &GlobalId) -> EngineError {
    EngineError::NoSuchTask {
        id: task.to_string(),
    }
}

fn failed(source: &ResolvedSource, error: SourceError) -> EngineError {
    EngineError::SourceFailed {
        name: source.name().to_string(),
        error,
    }
}