Skip to main content

onetaskgraph_core/engine/
mod.rs

1//! The query engine: fan-out, capability compensation, and the plan it reports.
2//!
3//! Given the sources a configuration resolved to and a query, the engine addresses one,
4//! several or every one of them **at once**, and decides per source and per predicate
5//! whether to push the predicate down or to apply it itself. A predicate a source
6//! declares [`Native`](Support::Native) is passed in the query and never re-applied. A
7//! predicate it declares [`Unsupported`](Support::Unsupported) is **removed** from what
8//! that source sees — the source would ignore it anyway, and leaving it in would invite
9//! a source to half-apply it — and the engine narrows the wider result set in memory.
10//!
11//! What makes that worth the machinery is that the two are not the same plan. A source
12//! with real server-side search keeps it, and a folder of Markdown beside it is
13//! compensated for, and the caller can see which of the two it got: every response
14//! carries a [`QueryPlan`], `--explain` renders it, and `--json` publishes it.
15//!
16//! Nothing here writes anything down. See [`fetch`] for the walk that makes that true.
17
18mod comment;
19mod copy;
20mod delivery;
21mod fetch;
22mod join;
23mod local;
24mod metadata;
25mod narrow;
26mod rendered;
27mod resume;
28mod update;
29
30use std::collections::BTreeMap;
31use std::num::NonZeroU32;
32use std::sync::atomic::{AtomicU32, Ordering};
33
34use chrono::{DateTime, Utc};
35use onetaskgraph_plugin_api::{
36    Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
37    MetadataRecord, NativeId, Page, PageRequest, Priority, Project, ProjectFilter, ProjectQuery,
38    SecretResolver, SourceError, SourceName, StatusCategory, Task, TaskQuery, TextFields,
39    TextQuery,
40};
41use schemars::JsonSchema;
42use serde::{Deserialize, Serialize};
43
44use crate::GlobalId;
45use crate::config::Config;
46use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
47use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
48
49use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
50use join::join_all;
51use local::{LocalDocuments, LocalProjects, LocalTasks};
52pub(crate) use resume::{Owed, Resumption, StreamState};
53use resume::{Resume, StreamKind};
54
55pub use comment::{CommentList, DeletedComment, TaskDetail};
56pub use copy::{
57    BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
58    Spent,
59};
60pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
61pub use local::ProjectSelector;
62pub use metadata::MetadataSet;
63pub use narrow::{TaskContentSet, TaskPrioritySet};
64pub use rendered::{
65    Body, DocumentCreate, Regenerated, Regeneration, RenderRequest, RenderTemplate, RenderedRecord,
66    TaskCreate, TaskCreated, TemplateAnswers, UnusedAnswers,
67};
68pub use update::TaskUpdated;
69
70/// One item, under the qualified id the engine addresses it by.
71///
72/// A plugin only ever deals in its own [`NativeId`]; qualifying one is the engine's job,
73/// so this type is the engine's and a plugin never constructs one.
74#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
75pub struct Qualified<T> {
76    /// `<source>:<native>`, the form a user types back at the command line.
77    pub id: GlobalId,
78    /// The item as its source reported it, unchanged.
79    pub item: T,
80}
81
82/// One dependency edge with both ends qualified.
83///
84/// An end may belong to another source. The near plugin owns that qualified far id and the
85/// engine reports it without resolving or fetching the far source, so it holds no index.
86#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
87pub struct QualifiedEdge {
88    /// The item the edge starts at, and the one that **depends on** the other.
89    ///
90    /// The direction a caller asked in says which end they named, not which end the edge
91    /// starts at: a forward read and the matching reverse read report the same edge.
92    pub from: QualifiedEndpoint,
93    /// The item the edge points at, and the one that must finish first.
94    pub to: QualifiedEndpoint,
95    /// What the edge means.
96    pub kind: onetaskgraph_plugin_api::DependencyKind,
97}
98
99/// One typed, qualified endpoint in an engine dependency response.
100#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
101pub struct QualifiedEndpoint {
102    /// `<source>:<native>`, preserved when a plugin reports another source.
103    pub id: GlobalId,
104    /// Whether this endpoint names a task or project.
105    pub kind: onetaskgraph_plugin_api::ItemKind,
106}
107
108impl std::fmt::Display for QualifiedEndpoint {
109    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
110        self.id.fmt(formatter)
111    }
112}
113
114/// One hit of a search that may cross entities.
115#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
116#[serde(tag = "kind", rename_all = "kebab-case")]
117pub enum SearchHit {
118    /// A task matched.
119    Task(Qualified<Task>),
120    /// A project matched.
121    Project(Qualified<Project>),
122}
123
124/// Which entities a search covers.
125#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
126#[serde(rename_all = "kebab-case")]
127pub enum SearchKind {
128    /// Tasks only.
129    Tasks,
130    /// Projects only.
131    Projects,
132    /// Both, interleaved.
133    #[default]
134    Both,
135}
136
137/// What one configured source is, as `sources list` reports it.
138#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
139pub struct SourceListing {
140    /// The name the configuration gave it.
141    pub source: SourceName,
142    /// The plugin kind behind it.
143    ///
144    /// A `String` because the vocabulary is open, not because it was not thought about:
145    /// this is the kind a source reports, and a subprocess-hosted plugin reports one
146    /// arriving over the wire from a binary this workspace never compiled. No
147    /// compile-time enumeration can hold that, and a newtype over the same string would
148    /// only move where an unrelated value is accepted.
149    // llmlint: ignore[invalid_states_unrepresentable] the reason above, and the one
150    // recorded for the same field of `SourcePlan` in plan.rs: `kind` is an open
151    // vocabulary a subprocess plugin extends at run time, and `SourcePlan.kind: String`
152    // is approved contract text this field is rendered beside.
153    pub kind: String,
154    /// Whether it built, and what it can do if it did.
155    #[serde(flatten)]
156    pub state: SourceState,
157}
158
159/// Whether a configured source is answering.
160#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
161#[serde(tag = "state", rename_all = "kebab-case")]
162pub enum SourceState {
163    /// The source built, and declares this.
164    Available {
165        /// What it applies itself.
166        capabilities: Capabilities,
167    },
168    /// The source could not be built at all.
169    Unavailable {
170        /// Why not.
171        error: SourceError,
172    },
173}
174
175/// Which page of a result set the caller wants.
176#[derive(Debug, Clone, PartialEq)]
177pub struct Paging {
178    /// The most items to return.
179    pub limit: NonZeroU32,
180    /// Where to resume, or `None` to start at the beginning.
181    pub token: Option<PageToken>,
182}
183
184/// The filters every list verb shares.
185#[derive(Debug, Clone, Default, PartialEq)]
186pub struct Filters {
187    /// Free-text search, when the caller asked for one.
188    pub text: Option<TextQuery>,
189    /// Label membership, by name.
190    pub labels: LabelFilter,
191    /// Status categories to keep. Empty means unfiltered.
192    pub statuses: Vec<StatusCategory>,
193}
194
195/// A request for a page of tasks.
196#[derive(Debug, Clone)]
197pub struct TaskRequest {
198    /// Which sources to address. Empty means the configuration's own selection.
199    pub sources: Vec<SourceName>,
200    /// What to keep.
201    pub filters: Filters,
202    /// Which project the tasks belong to.
203    pub project: ProjectSelector,
204    /// Priorities to keep: a task matches when its priority is any one of these. Empty means
205    /// unfiltered.
206    ///
207    /// Here rather than in [`Filters`], because a project has no priority: the filters a
208    /// project list shares with a task list are the ones both kinds of item carry.
209    pub priorities: Vec<Priority>,
210    /// Comment activity to keep: a task matches when one of its comments was created, or last
211    /// edited, at or after this instant. `None` means unfiltered.
212    ///
213    /// Here rather than in [`Filters`] for the reason [`priorities`](Self::priorities) is: a
214    /// project has no comments. A source that does not apply it natively is narrowed by a read
215    /// of its comments for every task its other predicates kept.
216    pub commented_since: Option<DateTime<Utc>>,
217    /// Which page.
218    pub paging: Paging,
219}
220
221/// A request for a page of projects.
222#[derive(Debug, Clone)]
223pub struct ProjectRequest {
224    /// Which sources to address. Empty means the configuration's own selection.
225    pub sources: Vec<SourceName>,
226    /// What to keep.
227    pub filters: Filters,
228    /// Which page.
229    pub paging: Paging,
230}
231
232/// The filters a document list carries.
233///
234/// Its own type rather than [`Filters`], and the difference is the whole of it: there are
235/// no statuses here. A document is not work and carries no status, so a status filter has
236/// nothing to compare against — and a shared type with a `statuses` field would let a
237/// caller write one down and leave every reader to decide what it meant.
238#[derive(Debug, Clone, Default, PartialEq)]
239pub struct DocumentFilters {
240    /// Free-text search, when the caller asked for one.
241    pub text: Option<TextQuery>,
242    /// Label membership, by name.
243    pub labels: LabelFilter,
244}
245
246/// A request for a page of documents.
247#[derive(Debug, Clone)]
248pub struct DocumentRequest {
249    /// Which sources to address. Empty means the configuration's own selection.
250    pub sources: Vec<SourceName>,
251    /// What to keep.
252    pub filters: DocumentFilters,
253    /// Which project the documents live in.
254    pub project: ProjectSelector,
255    /// Which page.
256    pub paging: Paging,
257}
258
259/// A request for a page of labels.
260#[derive(Debug, Clone)]
261pub struct LabelRequest {
262    /// Which sources to address. Empty means the configuration's own selection.
263    pub sources: Vec<SourceName>,
264    /// Which page.
265    pub paging: Paging,
266}
267
268/// A request for a page of search hits.
269#[derive(Debug, Clone)]
270pub struct SearchRequest {
271    /// Which sources to address. Empty means the configuration's own selection.
272    pub sources: Vec<SourceName>,
273    /// What to look for, and where.
274    pub text: TextQuery,
275    /// Which entities to cover.
276    pub kind: SearchKind,
277    /// Which page.
278    pub paging: Paging,
279}
280
281/// A request for a page of one item's dependency edges.
282#[derive(Debug, Clone)]
283pub struct DependencyRequest {
284    /// The qualified item to walk from.
285    pub id: GlobalId,
286    /// Which way to walk.
287    pub direction: Direction,
288    /// Which page.
289    pub paging: Paging,
290}
291
292/// What the engine refuses before it asks any source.
293///
294/// Distinct from a [`SourceFailure`], which is one source failing while the others
295/// answer: everything here means the request itself cannot be run at all.
296///
297/// `Eq` is deliberately absent: three variants carry the [`SourceError`] the source
298/// itself gave, which is what a user has to act on, and that type is `PartialEq` alone.
299#[derive(Debug, Clone, PartialEq, thiserror::Error)]
300pub enum EngineError {
301    /// A `--source` named something no configuration configures.
302    #[error(
303        "no source named {name:?} is configured\n\
304         next: name one of the configured sources ({configured}), or add {name:?} under \
305         `sources` — `onetaskgraph sources list` shows what this configuration has."
306    )]
307    UnknownSource {
308        /// The name that was asked for.
309        name: String,
310        /// The names that exist, for the message.
311        configured: String,
312    },
313
314    /// A `--page` token that decodes but does not belong to this query.
315    #[error(
316        "{message}\n\
317         next: page with a token exactly as the previous page reported it, and against \
318         the same configuration — or drop `--page` to start the walk again."
319    )]
320    Token {
321        /// What the token claims that this configuration cannot honour.
322        message: String,
323    },
324
325    /// Nothing at all is configured, so there is nothing to ask.
326    #[error(
327        "no sources are configured\n\
328         next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
329         prints what each plugin accepts."
330    )]
331    NoSources,
332
333    /// A copy named a destination whose plugin has no write side.
334    #[error(
335        "source {name} cannot be written: its plugin is {kind}, which has no write \
336         side\n\
337         next: copy into a source whose plugin can be written — `onetaskgraph sources \
338         list` reports each one's plugin."
339    )]
340    NotWritable {
341        /// The configured name of the destination.
342        name: String,
343        /// The plugin behind it.
344        kind: String,
345    },
346
347    /// A document read or write named a source whose plugin has no documents.
348    ///
349    /// The same shape [`NotWritable`](Self::NotWritable) has, and for the same reason: the
350    /// declaration is read once at the handshake, so a source that says it has no
351    /// documents is refused before anything is read rather than asked and then apologised
352    /// for.
353    #[error(
354        "source {name} has no documents: its plugin is {kind}, which holds none\n\
355         next: name a source whose plugin has documents — `onetaskgraph sources list` \
356         reports each one's plugin and what it declares."
357    )]
358    NoDocuments {
359        /// The configured name of the source.
360        name: String,
361        /// The plugin behind it.
362        kind: String,
363    },
364
365    /// A comment verb named a task of a source whose plugin has no comments.
366    ///
367    /// The shape [`NoDocuments`](Self::NoDocuments) has, for its reason: the declaration is
368    /// read once at the handshake, so the source is refused before anything is read.
369    #[error(
370        "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
371         next: name a task of a source whose plugin has comments — `onetaskgraph sources \
372         list` reports each one's plugin."
373    )]
374    NoComments {
375        /// The configured name of the source.
376        name: String,
377        /// The plugin behind it.
378        kind: String,
379    },
380
381    /// A comment verb that writes named a source whose comments can be read but not written.
382    #[error(
383        "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
384         but not added to, edited or removed\n\
385         next: list them with `onetaskgraph task comment list`, or comment on a task of a \
386         source whose plugin can be written — `onetaskgraph sources list` reports each \
387         one's plugin."
388    )]
389    CommentsNotWritable {
390        /// The configured name of the source.
391        name: String,
392        /// The plugin behind it.
393        kind: String,
394    },
395
396    /// `task status set` named a source whose plugin has no write side.
397    #[error(
398        "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
399         next: set the status in that source itself, or name a task of a source whose plugin \
400         can be written — `onetaskgraph sources list` reports each one's plugin."
401    )]
402    StatusNotWritable {
403        /// The configured name of the source.
404        name: String,
405        /// The plugin behind it.
406        kind: String,
407    },
408
409    /// A metadata verb named a record of a source whose plugin has no write side.
410    #[error(
411        "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
412         write side\n\
413         next: set the key in that source itself, or name a {record} of a source whose plugin \
414         can be written — `onetaskgraph sources list` reports each one's plugin."
415    )]
416    MetadataNotWritable {
417        /// The configured name of the source.
418        name: String,
419        /// The plugin behind it.
420        kind: String,
421        /// Which kind of record was named.
422        record: MetadataRecord,
423    },
424
425    /// `task priority set` named a source whose plugin has no write side.
426    #[error(
427        "source {name} cannot write a priority: its plugin is {kind}, which has no write side\n\
428         next: set the priority in that source itself, or name a task of a source whose plugin \
429         can be written — `onetaskgraph sources list` reports each one's plugin."
430    )]
431    PriorityNotWritable {
432        /// The configured name of the source.
433        name: String,
434        /// The plugin behind it.
435        kind: String,
436    },
437
438    /// `task content set` named a source whose plugin has no write side.
439    #[error(
440        "source {name} cannot write a task's content: its plugin is {kind}, which has no write \
441         side\n\
442         next: edit the content in that source itself, or name a task of a source whose plugin \
443         can be written — `onetaskgraph sources list` reports each one's plugin."
444    )]
445    ContentNotWritable {
446        /// The configured name of the source.
447        name: String,
448        /// The plugin behind it.
449        kind: String,
450    },
451
452    /// `task update` named a source whose plugin has no write side.
453    #[error(
454        "source {name} cannot update a task: its plugin is {kind}, which has no write side\n\
455         next: change the task in that source itself, or name a task of a source whose plugin \
456         can be written — `onetaskgraph sources list` reports each one's plugin."
457    )]
458    UpdateNotWritable {
459        /// The configured name of the source.
460        name: String,
461        /// The plugin behind it.
462        kind: String,
463    },
464
465    /// `task create` or `document create` named a source whose plugin has no write side.
466    #[error(
467        "source {name} cannot create a {record}: its plugin is {kind}, which has no write side\n\
468         next: create it in a source whose plugin can be written — `onetaskgraph sources list` \
469         reports each one's plugin."
470    )]
471    NotCreatable {
472        /// The configured name of the source.
473        name: String,
474        /// The plugin behind it.
475        kind: String,
476        /// Which kind of record was to be created.
477        record: MetadataRecord,
478    },
479
480    /// `task render` or `document render` named a source whose plugin has no write side.
481    #[error(
482        "source {name} cannot write a {record}'s rendering: its plugin is {kind}, which has no \
483         write side\n\
484         next: render it with --dry-run to read the result, or regenerate a {record} of a source \
485         whose plugin can be written."
486    )]
487    RenderingNotWritable {
488        /// The configured name of the source.
489        name: String,
490        /// The plugin behind it.
491        kind: String,
492        /// A task or a document.
493        record: RenderedRecord,
494    },
495
496    /// A regenerate left required variables unanswered, because the answers stored beside the
497    /// item could not be used or did not cover them.
498    #[error(
499        "supply every required answer to regenerate {id}: {} unanswered, and {reason}\n\
500         next: answer {} with --var NAME=VALUE or an answers file (--answers FILE), or run \
501         interactively to be asked.",
502        names.join(", "),
503        if names.len() == 1 { "it" } else { "each" }
504    )]
505    MissingAnswers {
506        /// The item regenerated.
507        id: String,
508        /// Every required variable left unanswered, in declaration order.
509        names: Vec<String>,
510        /// What became of the stored answers, as a clause: why they were not the base, or
511        /// that they were and did not answer these.
512        reason: String,
513    },
514
515    /// A template could not be loaded or rendered, or refused the answers it was given.
516    #[error("{error}")]
517    Template {
518        /// What the template refused.
519        error: crate::template::TemplateError,
520    },
521
522    /// An answers read named an item with no answers stored beside it.
523    #[error(
524        "{record} {id} has no stored template answers: {reason}\n\
525         next: regenerate it with every required answer (`onetaskgraph {record} render {id} \
526         --var NAME=VALUE`), or read its provenance with `onetaskgraph {record} show {id}`."
527    )]
528    NoStoredAnswers {
529        /// A task or a document.
530        record: RenderedRecord,
531        /// The item.
532        id: String,
533        /// Why none are there.
534        reason: String,
535    },
536
537    /// A regenerate named no template, and the item records none.
538    #[error(
539        "{record} {id} records no template it was rendered from, and none was given\n\
540         next: name one with --template FILE or --template-loader FILE."
541    )]
542    NoTemplate {
543        /// A task or a document.
544        record: RenderedRecord,
545        /// The item.
546        id: String,
547    },
548
549    /// A regenerate named no template, and the item's `onetaskgraph.template` entry is not one
550    /// this product writes, so it names nothing to render.
551    #[error(
552        "{record} {id} records a template entry this product did not write — {problem}\n\
553         next: regenerate it with --template FILE or --template-loader FILE and every required \
554         answer, which records a fresh entry."
555    )]
556    MalformedProvenance {
557        /// A task or a document.
558        record: RenderedRecord,
559        /// The item.
560        id: String,
561        /// What is wrong with the entry.
562        problem: String,
563    },
564
565    /// A regenerate named no template, and the one the item records is not a readable file.
566    #[error(
567        "{record} {id} was rendered from {reference:?}, which is not a readable file, so it \
568         cannot be re-read; a recorded reference is never turned into a location\n\
569         next: supply the template with --template-loader FILE (a loader document naming what \
570         to render), or name a template file with --template FILE."
571    )]
572    TemplateNotAFile {
573        /// A task or a document.
574        record: RenderedRecord,
575        /// The item.
576        id: String,
577        /// The provenance `template` it records.
578        reference: String,
579    },
580
581    /// A priority other than `none` was to be written — by `task priority set` or by a copy —
582    /// to a source whose plugin declares it holds no priority.
583    ///
584    /// Refused before that source is asked anything, so a copy it stops has written nothing.
585    #[error(
586        "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
587         be written to it: its plugin is {kind}, which declares priority unsupported\n\
588         next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
589         reports what each declares — or set the task's priority to none first; a \
590         github-projects source holds one once its configuration sets priority_mapping."
591    )]
592    NoPriority {
593        /// The configured name of the source that cannot hold it.
594        name: String,
595        /// The plugin behind it.
596        kind: String,
597        /// The qualified id of the task whose priority it is.
598        task: String,
599        /// The priority that could not be written.
600        priority: onetaskgraph_plugin_api::Priority,
601    },
602
603    /// A metadata verb named a project its source does not hold.
604    #[error(
605        "no project with the id {id}\n\
606         next: check the id, or list what is there — `onetaskgraph project list` reports every \
607         project the configured sources hold."
608    )]
609    NoSuchProject {
610        /// The qualified id that named nothing.
611        id: String,
612    },
613
614    /// A metadata verb named a document its source does not hold.
615    #[error(
616        "no document with the id {id}\n\
617         next: check the id, or list what is there — `onetaskgraph document list` reports every \
618         document the configured sources hold."
619    )]
620    NoSuchDocument {
621        /// The qualified id that named nothing.
622        id: String,
623    },
624
625    /// A comment verb named a task its source does not hold.
626    #[error(
627        "no task with the id {id}\n\
628         next: check the id, or list what is there — `onetaskgraph task list` reports every \
629         task the configured sources hold."
630    )]
631    NoSuchTask {
632        /// The qualified id that named nothing.
633        id: String,
634    },
635
636    /// A comment verb named a comment the task does not have.
637    #[error(
638        "task {task} has no comment with the id {comment}\n\
639         next: list its comments — `onetaskgraph task comment list {task}` reports each \
640         one's id."
641    )]
642    NoSuchComment {
643        /// The qualified id of the task.
644        task: String,
645        /// The comment id that named nothing on it.
646        comment: String,
647    },
648
649    /// A source the verb reads or writes one item of is configured but could not be built.
650    ///
651    /// Distinct from [`DestinationUnavailable`](Self::DestinationUnavailable), whose next
652    /// action is about a copy.
653    #[error(
654        "source {name} could not be built: {error}\n\
655         next: fix that source — `onetaskgraph sources list` reports its state — then run the \
656         command again."
657    )]
658    SourceUnavailable {
659        /// The configured name of the source.
660        name: String,
661        /// Why it did not build.
662        error: SourceError,
663    },
664
665    /// A source refused, or failed, the one call a verb made of it.
666    ///
667    /// Distinct from [`SourceRefused`](Self::SourceRefused), whose next action is about a
668    /// copy, and from a [`SourceFailure`], which leaves other sources' results standing: a
669    /// comment is one call to one source, so there is nothing else to report beside it.
670    #[error(
671        "source {name} could not do it: {error}\n\
672         next: fix what the source named above, then run the command again."
673    )]
674    SourceFailed {
675        /// The source that failed.
676        name: String,
677        /// What it said.
678        error: SourceError,
679    },
680
681    /// A copy named a destination that is configured but could not be built.
682    #[error(
683        "the destination source {name} could not be built: {error}\n\
684         next: fix that source — `onetaskgraph sources list` reports its state — then \
685         copy again."
686    )]
687    DestinationUnavailable {
688        /// The configured name of the destination.
689        name: String,
690        /// Why it did not build.
691        error: SourceError,
692    },
693
694    /// A copy named an id its own source does not hold.
695    #[error(
696        "no item with the id {id}\n\
697         next: check the id, or list what is there — `onetaskgraph task list` and \
698         `onetaskgraph project list` report what the configured sources hold."
699    )]
700    NoSuchItem {
701        /// The qualified id that named nothing.
702        id: String,
703    },
704
705    /// An item's recorded origin names an item the destination no longer holds.
706    ///
707    /// Refused rather than created: the destination item was deleted or moved on purpose,
708    /// and creating a second one there would duplicate work somebody removed.
709    #[error(
710        "{item} was copied from {origin}, which that destination no longer holds\n\
711         next: re-run with --recreate to create a new item there instead, or restore \
712         {origin}."
713    )]
714    StaleOrigin {
715        /// The item being copied.
716        item: String,
717        /// The origin it records.
718        origin: String,
719    },
720
721    /// A member copy named a task that is not a member of the project being copied.
722    ///
723    /// Refused before anything is written, because a copy that names members names the
724    /// part of one project it carries: a task of some other project, or of no project, is
725    /// not a narrower version of that copy but a different one.
726    #[error(
727        "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
728         next: name a task `onetaskgraph task list --project {project}` reports, or copy \
729         {id} on its own with `onetaskgraph task copy`.",
730        project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
731    )]
732    NotAMember {
733        /// The id that was named as a member.
734        id: GlobalId,
735        /// The projects the copy carries, none of which it is a member of.
736        projects: Vec<GlobalId>,
737    },
738
739    /// A member copy's item depends on a member it was not told to carry, and that member
740    /// records no origin naming the destination.
741    ///
742    /// Refused before anything is written. The destination id of a member the copy does
743    /// not carry comes from that member's own recorded origin and from nowhere else — not
744    /// from a walk of the destination, which is the read a member copy exists to avoid —
745    /// and writing the edge without it would either drop it silently or point it at the
746    /// id the member has at its source, which the destination has never heard of.
747    #[error(
748        "{item} depends on {member}, which this copy was not told to carry and which records \
749         no origin in {destination}\n\
750         next: name {member} with --member as well, record its {destination} id at \
751         onetaskgraph.origin, or copy the whole project without --member."
752    )]
753    UnrecordedMember {
754        /// The item whose edge could not be resolved.
755        item: GlobalId,
756        /// The member that edge points at.
757        member: GlobalId,
758        /// The destination it records no origin in.
759        destination: SourceName,
760    },
761
762    /// A source refused something the copy asked of it.
763    ///
764    /// Distinct from a [`SourceFailure`], which leaves the other sources' results
765    /// standing: a copy is one write into one destination, and half of one is not an
766    /// answer.
767    #[error(
768        "source {name} could not do it: {error}\n\
769         next: fix what the source named above, then copy again."
770    )]
771    SourceRefused {
772        /// The source that refused.
773        name: String,
774        /// What it said.
775        error: SourceError,
776    },
777
778    /// A copy failed, and the destination could not be put back the way it was found.
779    ///
780    /// A copy is either complete or it never happened, and the engine undoes its own
781    /// writes to make that true. When the destination will not take one of them back, the
782    /// failure has to say so and name what is still there: a user told only that the copy
783    /// failed would copy again over a destination nobody has described to them, which is
784    /// the retry that trips a hosted destination's rate limiter.
785    #[error(
786        "the copy failed and could not be undone.\n\
787         it failed because: {error}\n\
788         it could not be undone because: {refusal}\n\
789         so the destination still holds: {left_behind}\n\
790         next: remove those items at the destination, then copy again."
791    )]
792    CopyNotUndone {
793        /// Why the copy failed in the first place.
794        error: Box<EngineError>,
795        /// The items the copy created or overwrote and could not take back.
796        ///
797        /// The ids themselves rather than a sentence about them: a caller acting on this —
798        /// removing them, or reporting them — needs the ids, and the joined form only the
799        /// message needs is [`LeftBehind`]'s own `Display`.
800        left_behind: LeftBehind,
801        /// What the destination said when the copy tried to take them back.
802        refusal: SourceError,
803    },
804}
805
806/// The items a failed copy left at the destination — one of them at least, always.
807///
808/// A plain `Vec` here would let [`EngineError::CopyNotUndone`] be built naming nothing
809/// still there, and naming what is still there is the whole reason that failure is
810/// distinct from the one it wraps: a user told only that the copy could not be undone,
811/// and then handed an empty list, has been told about a destination nobody described to
812/// them, which is exactly the blind retry this mechanism exists to remove. The first id
813/// is a field of its own, so the empty case cannot be written down.
814#[derive(Debug, Clone, PartialEq)]
815pub struct LeftBehind {
816    /// The first item the destination would not take back.
817    first: GlobalId,
818    /// The ones it would not take back after that, in the order it refused them.
819    rest: Vec<GlobalId>,
820}
821
822impl LeftBehind {
823    /// The list holding the one item every such refusal has to name.
824    #[must_use]
825    pub fn new(first: GlobalId) -> Self {
826        Self {
827            first,
828            rest: Vec::new(),
829        }
830    }
831
832    /// Record another item the destination would not take back.
833    pub fn push(&mut self, id: GlobalId) {
834        self.rest.push(id);
835    }
836
837    /// Every item still at the destination, in the order the destination refused them.
838    pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
839        std::iter::once(&self.first).chain(self.rest.iter())
840    }
841}
842
843impl std::fmt::Display for LeftBehind {
844    /// The qualified ids, comma-separated — the form the failure message reads in.
845    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
846        write!(formatter, "{}", self.first)?;
847        for id in &self.rest {
848            write!(formatter, ", {id}")?;
849        }
850        Ok(())
851    }
852}
853
854/// One configured source, in exactly one of the two states a configured source has.
855///
856/// A sum rather than two lists side by side: with a `ready` vector and an `unavailable`
857/// one, a source appearing in both is a shape the type permits and every reader has to
858/// decide about — and they would not all decide the same way, since one of them fans a
859/// query out and another names failures.
860pub enum ConfiguredSource {
861    /// It built, and answers queries.
862    Ready(ResolvedSource),
863    /// It did not, and every response says so instead.
864    Unavailable(UnavailableSource),
865}
866
867impl ConfiguredSource {
868    /// The name the configuration gave it, whichever state it is in.
869    #[must_use]
870    pub fn name(&self) -> &SourceName {
871        match self {
872            Self::Ready(source) => source.name(),
873            Self::Unavailable(source) => source.name(),
874        }
875    }
876}
877
878/// The sources a configuration resolved to, and the queries they answer.
879pub struct Engine {
880    /// Every configured source, in configured-name order, each in one state.
881    sources: Vec<ConfiguredSource>,
882    /// Which sources answer when a request names none.
883    selection: Vec<SourceName>,
884}
885
886impl Engine {
887    /// Build every source a configuration names.
888    ///
889    /// A source whose plugin refuses to build — a credential that is not there, a
890    /// plugin whose implementation has not landed — is **not** fatal: it becomes an
891    /// entry in every response's `errors`, exactly as a source that fails mid-query
892    /// does, and the other sources still answer. A user with three sources and one
893    /// expired token gets the other two rather than nothing.
894    #[must_use]
895    pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
896        let (ready, unavailable) = resolve_available(config, secrets);
897        Self::new(
898            ready
899                .into_iter()
900                .map(ConfiguredSource::Ready)
901                .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
902                .collect(),
903            config.selected_sources(),
904        )
905    }
906
907    /// Drive sources built elsewhere — the engine's own tests, and any caller holding a
908    /// source it did not resolve from a configuration document.
909    #[must_use]
910    pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
911        Self { sources, selection }
912    }
913
914    /// Every source that built, in configured-name order.
915    fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
916        self.sources.iter().filter_map(|source| match source {
917            ConfiguredSource::Ready(ready) => Some(ready),
918            ConfiguredSource::Unavailable(_) => None,
919        })
920    }
921
922    /// Every source that did not, in the same order.
923    fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
924        self.sources.iter().filter_map(|source| match source {
925            ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
926            ConfiguredSource::Ready(_) => None,
927        })
928    }
929
930    /// Every configured source, whether or not it built, in name order.
931    #[must_use]
932    pub fn listing(&self) -> Vec<SourceListing> {
933        let mut listings: Vec<SourceListing> = self
934            .ready()
935            .map(|source| SourceListing {
936                source: source.name().clone(),
937                kind: source.kind().to_owned(),
938                state: SourceState::Available {
939                    capabilities: source.source().capabilities(),
940                },
941            })
942            .chain(self.unavailable().map(|source| SourceListing {
943                source: source.name().clone(),
944                kind: source.kind().to_owned(),
945                state: SourceState::Unavailable {
946                    error: source.error().clone(),
947                },
948            }))
949            .collect();
950        listings.sort_by(|left, right| left.source.cmp(&right.source));
951        listings
952    }
953
954    /// Whether this configuration has a source called `name`, built or not.
955    ///
956    /// A caller reading a `--project` argument needs this: `urn:project:1` is a qualified
957    /// id only if `urn` is a source here, and a native id full of colons otherwise. That
958    /// rule cannot be applied without knowing what is configured.
959    #[must_use]
960    pub fn has(&self, name: &SourceName) -> bool {
961        self.sources.iter().any(|source| source.name() == name)
962    }
963
964    /// One page of tasks.
965    ///
966    /// # Errors
967    ///
968    /// Returns [`EngineError`] when the request names a source nothing configures, or
969    /// carries a page token this engine did not issue. One source failing is not an
970    /// error: it lands in the response's `errors`.
971    pub async fn tasks(
972        &self,
973        request: &TaskRequest,
974    ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
975        let mut names = self.resolve_selection(&request.sources)?;
976        // A qualified project id names one project of one source, so no other source can
977        // hold a task in it. Narrowing here means the plan reports the source that was
978        // actually asked rather than a row of empty entries for sources that could not
979        // have answered.
980        if let ProjectSelector::Qualified(id) = &request.project {
981            self.known(&id.source)?;
982            names.retain(|name| name == &id.source);
983        }
984        let query = shape(
985            "task-list",
986            &names,
987            &(
988                &request.filters,
989                &request.project,
990                &request.priorities,
991                &request.commented_since,
992            ),
993        );
994        let states = resumption(
995            self,
996            request.paging.token.as_ref(),
997            &[StreamKind::Items],
998            &query,
999        )?;
1000        let budget = request.paging.limit.get();
1001
1002        let mut answer = Answer::new();
1003        let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1004
1005        let shapes: Vec<TaskShape> = ready
1006            .iter()
1007            .map(|source| {
1008                shape_tasks(
1009                    &source.source().capabilities(),
1010                    &request.filters,
1011                    &project_filter(&request.project),
1012                    &request.priorities,
1013                    request.commented_since,
1014                )
1015            })
1016            .collect();
1017        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1018        let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1019
1020        let walks = ready
1021            .iter()
1022            .enumerate()
1023            .map(|(index, source)| {
1024                fetch_tasks(
1025                    source,
1026                    &shapes[index],
1027                    &starts[index],
1028                    budget,
1029                    &counters[index],
1030                )
1031            })
1032            .collect();
1033
1034        let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1035        answer.finish(
1036            streams,
1037            budget,
1038            owed(&states),
1039            &query,
1040            |name, task: Task| {
1041                delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1042            },
1043        )
1044    }
1045
1046    /// One page of projects.
1047    ///
1048    /// # Errors
1049    ///
1050    /// As [`tasks`](Self::tasks).
1051    pub async fn projects(
1052        &self,
1053        request: &ProjectRequest,
1054    ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1055        let names = self.resolve_selection(&request.sources)?;
1056        let query = shape("project-list", &names, &request.filters);
1057        let states = resumption(
1058            self,
1059            request.paging.token.as_ref(),
1060            &[StreamKind::Items],
1061            &query,
1062        )?;
1063        let budget = request.paging.limit.get();
1064
1065        let mut answer = Answer::new();
1066        // A source declaring `projects: unsupported` has no project table at all, so
1067        // there is nothing to compensate for and nothing to ask: the predicate is
1068        // reported unavailable and that source contributes no rows. This is the one
1069        // outcome the engine cannot narrow its way out of, which is what `unavailable`
1070        // in the plan is for.
1071        let mut with_projects = Vec::new();
1072        for source in answer.split(self, &names) {
1073            if source.source().capabilities().projects.is_native() {
1074                with_projects.push(source);
1075            } else {
1076                answer.unreachable_predicate(source, Predicate::Project);
1077            }
1078        }
1079        let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1080
1081        let shapes: Vec<ProjectShape> = ready
1082            .iter()
1083            .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1084            .collect();
1085        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1086        let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1087
1088        let walks = ready
1089            .iter()
1090            .enumerate()
1091            .map(|(index, source)| {
1092                fetch_projects(
1093                    source,
1094                    &shapes[index],
1095                    &starts[index],
1096                    budget,
1097                    &counters[index],
1098                )
1099            })
1100            .collect();
1101
1102        let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1103        answer.finish(
1104            streams,
1105            budget,
1106            owed(&states),
1107            &query,
1108            |name, project: Project| Qualified {
1109                id: GlobalId::new(name.clone(), project.id.clone()),
1110                item: project,
1111            },
1112        )
1113    }
1114
1115    /// One page of documents.
1116    ///
1117    /// A source declaring it has no documents is **not asked**: the declaration is read
1118    /// once here, that source contributes no rows, and the plan reports
1119    /// [`Predicate::Document`] unavailable for it. So a document list spanning a mixed set
1120    /// of sources answers with what the document-bearing ones hold and says nothing
1121    /// alarming about the others — the same shape a source with no project table takes.
1122    ///
1123    /// # Errors
1124    ///
1125    /// As [`tasks`](Self::tasks).
1126    pub async fn documents(
1127        &self,
1128        request: &DocumentRequest,
1129    ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1130        let mut names = self.resolve_selection(&request.sources)?;
1131        // A qualified project id names one project of one source, exactly as it does for
1132        // tasks, so no other source can hold a document in it.
1133        if let ProjectSelector::Qualified(id) = &request.project {
1134            self.known(&id.source)?;
1135            names.retain(|name| name == &id.source);
1136        }
1137        let query = shape(
1138            "document-list",
1139            &names,
1140            &(&request.filters, &request.project),
1141        );
1142        let states = resumption(
1143            self,
1144            request.paging.token.as_ref(),
1145            &[StreamKind::Items],
1146            &query,
1147        )?;
1148        let budget = request.paging.limit.get();
1149
1150        let mut answer = Answer::new();
1151        let mut with_documents = Vec::new();
1152        for source in answer.split(self, &names) {
1153            if source.source().capabilities().documents.is_native() {
1154                with_documents.push(source);
1155            } else {
1156                answer.unreachable_predicate(source, Predicate::Document);
1157            }
1158        }
1159        let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1160
1161        let shapes: Vec<DocumentShape> = ready
1162            .iter()
1163            .map(|source| {
1164                shape_documents(
1165                    &source.source().capabilities(),
1166                    &request.filters,
1167                    &project_filter(&request.project),
1168                )
1169            })
1170            .collect();
1171        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1172        let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1173
1174        let walks = ready
1175            .iter()
1176            .enumerate()
1177            .map(|(index, source)| {
1178                fetch_documents(
1179                    source,
1180                    &shapes[index],
1181                    &starts[index],
1182                    budget,
1183                    &counters[index],
1184                )
1185            })
1186            .collect();
1187
1188        let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1189        answer.finish(
1190            streams,
1191            budget,
1192            owed(&states),
1193            &query,
1194            |name, document: Document| Qualified {
1195                id: GlobalId::new(name.clone(), document.id.clone()),
1196                item: document,
1197            },
1198        )
1199    }
1200
1201    /// One page of labels.
1202    ///
1203    /// # Errors
1204    ///
1205    /// As [`tasks`](Self::tasks).
1206    pub async fn labels(
1207        &self,
1208        request: &LabelRequest,
1209    ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1210        let names = self.resolve_selection(&request.sources)?;
1211        let query = shape("label-list", &names, &());
1212        let states = resumption(
1213            self,
1214            request.paging.token.as_ref(),
1215            &[StreamKind::Items],
1216            &query,
1217        )?;
1218        let budget = request.paging.limit.get();
1219
1220        let mut answer = Answer::new();
1221        let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1222        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1223        let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1224
1225        let walks = ready
1226            .iter()
1227            .enumerate()
1228            .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1229            .collect();
1230
1231        let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1232        answer.finish(
1233            streams,
1234            budget,
1235            owed(&states),
1236            &query,
1237            |name, label: Label| Qualified {
1238                id: GlobalId::new(name.clone(), label.id.clone()),
1239                item: label,
1240            },
1241        )
1242    }
1243
1244    /// One page of search hits, over tasks, projects, or both.
1245    ///
1246    /// # Errors
1247    ///
1248    /// As [`tasks`](Self::tasks).
1249    pub async fn search(
1250        &self,
1251        request: &SearchRequest,
1252    ) -> Result<QueryResponse<SearchHit>, EngineError> {
1253        let names = self.resolve_selection(&request.sources)?;
1254        // The streams this search reads, which is what a token resuming it may name. A
1255        // `--kind both` walk that has exhausted one half carries only the other, so this
1256        // is what a token may name rather than what it must.
1257        let reads: &[StreamKind] = match request.kind {
1258            SearchKind::Tasks => &[StreamKind::Tasks],
1259            SearchKind::Projects => &[StreamKind::Projects],
1260            SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1261        };
1262        // The scope is deliberately not in the fingerprint: which streams a search covers
1263        // is exactly what `reads` checks below, name by name and with a message that says
1264        // which half a token names. Folding it in here would refuse the same mistake one
1265        // layer earlier and less clearly, and leave that check unreachable.
1266        let query = shape("search", &names, &request.text);
1267        let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1268        let budget = request.paging.limit.get();
1269        let filters = Filters {
1270            text: Some(request.text.clone()),
1271            ..Filters::default()
1272        };
1273
1274        let mut answer = Answer::new();
1275
1276        // One stream per (source, entity), because a search over both entities reads two
1277        // result sets from each source and each has its own place to resume.
1278        let mut ready = Vec::new();
1279        let mut kinds = Vec::new();
1280        let mut starts = Vec::new();
1281        for source in answer.split(self, &names) {
1282            let mut streams = Vec::new();
1283            if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1284                streams.push(StreamKind::Tasks);
1285            }
1286            if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1287                if source.source().capabilities().projects.is_native() {
1288                    streams.push(StreamKind::Projects);
1289                } else {
1290                    answer.unreachable_predicate(source, Predicate::Project);
1291                }
1292            }
1293            for stream in streams {
1294                if let Some(resume) = resume_at(&states, source.name(), stream) {
1295                    ready.push(source);
1296                    kinds.push(stream);
1297                    starts.push(resume);
1298                }
1299            }
1300        }
1301
1302        let shapes: Vec<HitShape> = ready
1303            .iter()
1304            .zip(kinds.iter())
1305            .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1306            .collect();
1307        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1308        let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1309
1310        let walks = ready
1311            .iter()
1312            .enumerate()
1313            .map(|(index, source)| {
1314                fetch_hits(
1315                    source,
1316                    &shapes[index],
1317                    &starts[index],
1318                    budget,
1319                    &counters[index],
1320                )
1321            })
1322            .collect();
1323
1324        let streams =
1325            answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1326        answer.finish(
1327            streams,
1328            budget,
1329            owed(&states),
1330            &query,
1331            |name, found: Found| match found {
1332                Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1333                    GlobalId::new(name.clone(), task.id.clone()),
1334                    task,
1335                )),
1336                Found::Project(project) => SearchHit::Project(Qualified {
1337                    id: GlobalId::new(name.clone(), project.id.clone()),
1338                    item: project,
1339                }),
1340            },
1341        )
1342    }
1343
1344    /// One task by its qualified id, or an empty page when there is no such task.
1345    ///
1346    /// # Errors
1347    ///
1348    /// Returns [`EngineError::UnknownSource`] when the id names a source nothing
1349    /// configures.
1350    pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1351        let name = self.known(&id.source)?;
1352        let mut answer = Answer::new();
1353        let selected = answer.split(self, std::slice::from_ref(&name));
1354        let Some(source) = selected.first() else {
1355            return answer.nothing();
1356        };
1357        let found = source.source().get_task(&id.native).await;
1358        let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1359        answer.one(source, found, |task| {
1360            delivery::qualified_task(qualified, task)
1361        })
1362    }
1363
1364    /// One project by its qualified id, or an empty page when there is no such project.
1365    ///
1366    /// # Errors
1367    ///
1368    /// As [`task`](Self::task).
1369    pub async fn project(
1370        &self,
1371        id: &GlobalId,
1372    ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1373        let name = self.known(&id.source)?;
1374        let mut answer = Answer::new();
1375        let selected = answer.split(self, std::slice::from_ref(&name));
1376        let Some(source) = selected.first() else {
1377            return answer.nothing();
1378        };
1379        let found = source.source().get_project(&id.native).await;
1380        let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1381        answer.one(source, found, |project| Qualified {
1382            id: qualified,
1383            item: project,
1384        })
1385    }
1386
1387    /// One document by its qualified id, or an empty page when there is no such document.
1388    ///
1389    /// A source declaring it has no documents holds none, so it is not asked and the
1390    /// answer is the empty page with [`Predicate::Document`] reported unavailable — the
1391    /// same answer the list verb gives, rather than a failure.
1392    ///
1393    /// # Errors
1394    ///
1395    /// As [`task`](Self::task).
1396    pub async fn document(
1397        &self,
1398        id: &GlobalId,
1399    ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1400        let name = self.known(&id.source)?;
1401        let mut answer = Answer::new();
1402        let selected = answer.split(self, std::slice::from_ref(&name));
1403        let Some(source) = selected.first() else {
1404            return answer.nothing();
1405        };
1406        if !source.source().capabilities().documents.is_native() {
1407            answer.unreachable_predicate(source, Predicate::Document);
1408            return answer.nothing();
1409        }
1410        let found = source.source().get_document(&id.native).await;
1411        let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1412        answer.one(source, found, |document| Qualified {
1413            id: qualified,
1414            item: document,
1415        })
1416    }
1417
1418    /// One page of a task's dependency edges.
1419    ///
1420    /// # Errors
1421    ///
1422    /// As [`task`](Self::task), plus [`EngineError::Token`] for a page token this engine
1423    /// did not issue.
1424    pub async fn task_dependencies(
1425        &self,
1426        request: &DependencyRequest,
1427    ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1428        self.dependencies(request, Entity::Task).await
1429    }
1430
1431    /// One page of a project's dependency edges.
1432    ///
1433    /// # Errors
1434    ///
1435    /// As [`task_dependencies`](Self::task_dependencies).
1436    pub async fn project_dependencies(
1437        &self,
1438        request: &DependencyRequest,
1439    ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1440        self.dependencies(request, Entity::Project).await
1441    }
1442
1443    /// Both dependency verbs, which differ only in which of a source's two edge sets
1444    /// they read and which of its two declarations governs the reverse direction.
1445    async fn dependencies(
1446        &self,
1447        request: &DependencyRequest,
1448        entity: Entity,
1449    ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1450        let name = self.known(&request.id.source)?;
1451        let query = shape(
1452            "dependencies",
1453            std::slice::from_ref(&name),
1454            &(entity, &request.id.native, request.direction),
1455        );
1456        let states = resumption(
1457            self,
1458            request.paging.token.as_ref(),
1459            &[StreamKind::Items],
1460            &query,
1461        )?;
1462        let budget = request.paging.limit.get();
1463
1464        let mut answer = Answer::new();
1465        let (ready, starts) = walking(
1466            answer.split(self, std::slice::from_ref(&name)),
1467            &states,
1468            StreamKind::Items,
1469        );
1470        let Some(source) = ready.first() else {
1471            return answer.nothing();
1472        };
1473
1474        let capabilities = source.source().capabilities();
1475        let support = match entity {
1476            Entity::Task => capabilities.task_dependencies,
1477            Entity::Project => capabilities.project_dependencies,
1478        };
1479        // `DependencySupport` has no unsupported variant on purpose: a dependency read is
1480        // answered natively or emulated by the scan below, never abandoned and never
1481        // silently empty.
1482        let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1483        let mut outcomes = Outcomes::default();
1484        if request.direction == Direction::DependedOnBy {
1485            if emulating {
1486                outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1487            } else {
1488                outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1489            }
1490        }
1491
1492        let counters = vec![AtomicU32::new(0)];
1493        let walked = fetch_edges(
1494            source,
1495            &request.id.native,
1496            request.direction,
1497            entity,
1498            emulating,
1499            &starts[0],
1500            budget,
1501            &counters[0],
1502        )
1503        .await;
1504
1505        let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1506        answer.finish(
1507            streams,
1508            budget,
1509            owed(&states),
1510            &query,
1511            |name, edge: DependencyEdge| QualifiedEdge {
1512                from: qualify_endpoint(name, edge.from),
1513                to: qualify_endpoint(name, edge.to),
1514                kind: edge.kind,
1515            },
1516        )
1517    }
1518
1519    /// The names a request addresses: the ones it gave, or the configuration's own.
1520    fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1521        if asked.is_empty() {
1522            if self.selection.is_empty() {
1523                return Err(EngineError::NoSources);
1524            }
1525            return Ok(self.selection.clone());
1526        }
1527        asked.iter().map(|name| self.known(name)).collect()
1528    }
1529
1530    /// `name` when this configuration has a source called that.
1531    fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1532        if self.has(name) {
1533            return Ok(name.clone());
1534        }
1535        if self.sources.is_empty() {
1536            return Err(EngineError::NoSources);
1537        }
1538        Err(EngineError::UnknownSource {
1539            name: name.to_string(),
1540            configured: self
1541                .listing()
1542                .iter()
1543                .map(|listing| listing.source.to_string())
1544                .collect::<Vec<_>>()
1545                .join(", "),
1546        })
1547    }
1548}
1549
1550fn qualify_endpoint(
1551    source: &SourceName,
1552    endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1553) -> QualifiedEndpoint {
1554    let kind = endpoint.kind;
1555    let is_qualified = endpoint.is_qualified();
1556    let endpoint_id = endpoint.into_id();
1557    QualifiedEndpoint {
1558        id: if is_qualified {
1559            endpoint_id
1560                .parse()
1561                .expect("plugin-api validates qualified dependency endpoints")
1562        } else {
1563            GlobalId::new(source.clone(), NativeId(endpoint_id))
1564        },
1565        kind,
1566    }
1567}
1568
1569/// Which of a source's two dependency graphs a request walks.
1570#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1571enum Entity {
1572    /// Task dependencies.
1573    Task,
1574    /// Project dependencies.
1575    Project,
1576}
1577
1578/// A search hit before it is qualified.
1579enum Found {
1580    /// A task matched.
1581    Task(Task),
1582    /// A project matched.
1583    Project(Project),
1584}
1585
1586/// What happened to one predicate against one source.
1587#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1588enum Outcome {
1589    /// Applied by the source itself.
1590    PushedDown,
1591    /// Applied by the engine over a wider result set.
1592    AppliedLocally,
1593    /// Answered by a bounded scan of the source.
1594    Emulated,
1595    /// Neither side could answer it, so this source contributed nothing for it.
1596    Unavailable,
1597}
1598
1599/// What happened to each predicate against one source.
1600///
1601/// Keyed by predicate, one outcome each, because those are the only states there are: a
1602/// predicate the source applied was not also applied here, and one nobody could answer
1603/// was not also pushed down. The four lists [`SourcePlan`] carries are this map fanned
1604/// out at the boundary — held *as* four lists a predicate could sit in all four at once,
1605/// and four contradictory claims about one predicate is the one thing the part of the
1606/// answer whose whole job is to say which of them is true must not be able to say.
1607///
1608/// A `BTreeMap` rather than a `HashMap` so the lists come out in one order and two runs
1609/// of a query render the same plan.
1610#[derive(Debug, Clone, Default, PartialEq)]
1611struct Outcomes(BTreeMap<Predicate, Outcome>);
1612
1613impl Outcomes {
1614    /// Record what happened to one predicate, replacing whatever was recorded before.
1615    ///
1616    /// Replacing rather than refusing: shaping a query decides each predicate once, and a
1617    /// second decision about the same one is the later one — there is no case here where
1618    /// both were meant to stand.
1619    fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1620        self.0.insert(predicate, outcome);
1621    }
1622
1623    /// Record the same outcome for several predicates, as a text search does for the two
1624    /// fields it covers.
1625    fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1626        for predicate in predicates {
1627            self.record(predicate, outcome);
1628        }
1629    }
1630
1631    /// The predicates this outcome befell, in the map's own stable order.
1632    fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1633        self.0
1634            .iter()
1635            .filter(|(_, recorded)| **recorded == outcome)
1636            .map(|(predicate, _)| *predicate)
1637            .collect()
1638    }
1639}
1640
1641/// The query one source sees, and the predicates left to the engine.
1642struct TaskShape {
1643    /// What the source is asked.
1644    pushed: TaskQuery,
1645    /// What the engine narrows afterwards.
1646    local: LocalTasks,
1647    /// What to report.
1648    outcomes: Outcomes,
1649}
1650
1651/// As [`TaskShape`], for projects.
1652struct ProjectShape {
1653    /// What the source is asked.
1654    pushed: ProjectQuery,
1655    /// What the engine narrows afterwards.
1656    local: LocalProjects,
1657    /// What to report.
1658    outcomes: Outcomes,
1659}
1660
1661/// As [`TaskShape`], for documents.
1662struct DocumentShape {
1663    /// What the source is asked.
1664    pushed: DocumentQuery,
1665    /// What the engine narrows afterwards.
1666    local: LocalDocuments,
1667    /// What to report.
1668    outcomes: Outcomes,
1669}
1670
1671/// As [`TaskShape`], for one entity's half of a search.
1672struct HitShape {
1673    /// Which entity this stream reads.
1674    stream: StreamKind,
1675    /// What a task stream asks.
1676    tasks: TaskQuery,
1677    /// What a project stream asks.
1678    projects: ProjectQuery,
1679    /// What the engine narrows afterwards, for tasks.
1680    local_tasks: LocalTasks,
1681    /// What the engine narrows afterwards, for projects.
1682    local_projects: LocalProjects,
1683    /// What to report.
1684    outcomes: Outcomes,
1685}
1686
1687/// The plan and the failures a response carries, accumulated as the verb runs.
1688///
1689/// One type rather than three parallel vectors threaded through every verb, because the
1690/// rule they enforce together is one rule: a source that fails contributes an error and
1691/// still leaves every other source's results standing.
1692struct Answer {
1693    /// One entry per source the engine addressed, merged by source at the end.
1694    plans: Vec<SourcePlan>,
1695    /// Every source that could not answer.
1696    errors: Vec<SourceFailure>,
1697}
1698
1699impl Answer {
1700    fn new() -> Self {
1701        Self {
1702            plans: Vec::new(),
1703            errors: Vec::new(),
1704        }
1705    }
1706
1707    /// The selected sources that built, recording the ones that did not as failures.
1708    ///
1709    /// A source that never built is reported and skipped rather than fatal: that is the
1710    /// same rule as a source failing mid-query, applied one step earlier.
1711    fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1712        let mut selected = Vec::new();
1713        for name in names {
1714            match engine.sources.iter().find(|source| source.name() == name) {
1715                Some(ConfiguredSource::Ready(source)) => selected.push(source),
1716                Some(ConfiguredSource::Unavailable(source)) => {
1717                    self.errors.push(source.failure());
1718                }
1719                None => {}
1720            }
1721        }
1722        selected
1723    }
1724
1725    /// Record that a source could answer nothing for `predicate`.
1726    fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1727        let mut outcomes = Outcomes::default();
1728        outcomes.record(predicate, Outcome::Unavailable);
1729        self.plans.push(plan_for(source, outcomes, 0));
1730    }
1731
1732    /// Turn each source's walk into a stream, keeping a failed source's failure.
1733    fn collect<T>(
1734        &mut self,
1735        ready: &[&ResolvedSource],
1736        walked: Vec<Result<Fetched<T>, SourceError>>,
1737        counters: &[AtomicU32],
1738        outcomes: Vec<Outcomes>,
1739    ) -> Vec<Stream<T>> {
1740        let kinds = vec![StreamKind::Items; ready.len()];
1741        self.collect_streams(ready, &kinds, walked, counters, outcomes)
1742    }
1743
1744    /// As [`collect`](Self::collect), where a source may contribute more than one stream.
1745    fn collect_streams<T>(
1746        &mut self,
1747        ready: &[&ResolvedSource],
1748        kinds: &[StreamKind],
1749        walked: Vec<Result<Fetched<T>, SourceError>>,
1750        counters: &[AtomicU32],
1751        outcomes: Vec<Outcomes>,
1752    ) -> Vec<Stream<T>> {
1753        let mut streams = Vec::new();
1754        for (index, result) in walked.into_iter().enumerate() {
1755            let source = ready[index];
1756            let pages = counters[index].load(Ordering::Relaxed);
1757            self.plans
1758                .push(plan_for(source, outcomes[index].clone(), pages));
1759            match result {
1760                Ok(fetched) => streams.push(Stream {
1761                    source: source.name().clone(),
1762                    kind: kinds[index],
1763                    fetched,
1764                }),
1765                // A stream that failed leaves the token, so a walk always terminates: a
1766                // source failing on every page would otherwise page forever.
1767                Err(error) => self.errors.push(SourceFailure {
1768                    source: source.name().clone(),
1769                    error,
1770                }),
1771            }
1772        }
1773        streams
1774    }
1775
1776    /// The response for a verb that reads exactly one item from exactly one source.
1777    fn one<T, U>(
1778        mut self,
1779        source: &ResolvedSource,
1780        found: Result<Option<T>, SourceError>,
1781        qualify: impl FnOnce(T) -> U,
1782    ) -> Result<QueryResponse<U>, EngineError> {
1783        self.plans.push(plan_for(source, Outcomes::default(), 1));
1784        let items = match found {
1785            Ok(Some(item)) => vec![qualify(item)],
1786            Ok(None) => Vec::new(),
1787            Err(error) => {
1788                self.errors.push(SourceFailure {
1789                    source: source.name().clone(),
1790                    error,
1791                });
1792                Vec::new()
1793            }
1794        };
1795        Ok(QueryResponse {
1796            items,
1797            next: None,
1798            plan: QueryPlan {
1799                per_source: merge_plans(self.plans),
1800            },
1801            errors: self.errors,
1802        })
1803    }
1804
1805    /// The response for a verb with nothing left to ask.
1806    fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1807        Ok(QueryResponse {
1808            items: Vec::new(),
1809            next: None,
1810            plan: QueryPlan {
1811                per_source: merge_plans(self.plans),
1812            },
1813            errors: self.errors,
1814        })
1815    }
1816
1817    /// Merge the streams into the caller's page and mint the token that resumes it.
1818    ///
1819    /// `first` is the stream the token being resumed says is owed the next row, so the
1820    /// round-robin picks up where the previous page stopped rather than restarting.
1821    fn finish<T, U>(
1822        self,
1823        streams: Vec<Stream<T>>,
1824        budget: u32,
1825        first: Option<&Owed>,
1826        query: &str,
1827        qualify: impl Fn(&SourceName, T) -> U,
1828    ) -> Result<QueryResponse<U>, EngineError> {
1829        let (rows, states, owed) = merge(streams, budget, first);
1830        let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1831        Ok(QueryResponse {
1832            items: rows
1833                .into_iter()
1834                .map(|(name, item)| qualify(&name, item))
1835                .collect(),
1836            next,
1837            plan: QueryPlan {
1838                per_source: merge_plans(self.plans),
1839            },
1840            errors: self.errors,
1841        })
1842    }
1843}
1844
1845/// One source's plan entry: the outcomes fanned out into the four lists the contract's
1846/// [`SourcePlan`] carries, each in one order so two runs read the same.
1847fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1848    SourcePlan {
1849        source: source.name().clone(),
1850        kind: source.kind().to_owned(),
1851        pushed_down: outcomes.with(Outcome::PushedDown),
1852        applied_locally: outcomes.with(Outcome::AppliedLocally),
1853        emulated: outcomes.with(Outcome::Emulated),
1854        unavailable: outcomes.with(Outcome::Unavailable),
1855        pages_fetched: pages,
1856    }
1857}
1858
1859/// One entry per source, however many streams that source contributed.
1860///
1861/// `search --kind both` reads two streams from each source, and a plan is per source:
1862/// two entries for one name would say the engine addressed it twice.
1863fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1864    let mut merged: Vec<SourcePlan> = Vec::new();
1865    for plan in plans {
1866        if let Some(existing) = merged
1867            .iter_mut()
1868            .find(|existing| existing.source == plan.source)
1869        {
1870            existing.pushed_down.extend(plan.pushed_down);
1871            existing.applied_locally.extend(plan.applied_locally);
1872            existing.emulated.extend(plan.emulated);
1873            existing.unavailable.extend(plan.unavailable);
1874            existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1875            for list in [
1876                &mut existing.pushed_down,
1877                &mut existing.applied_locally,
1878                &mut existing.emulated,
1879                &mut existing.unavailable,
1880            ] {
1881                list.sort_unstable();
1882                list.dedup();
1883            }
1884        } else {
1885            merged.push(plan);
1886        }
1887    }
1888    merged
1889}
1890
1891/// The sources still walking, with where each picks up.
1892///
1893/// A token names every stream that has more to give, so a source **absent** from one has
1894/// been exhausted and is not asked again. Without that, the second page of a walk would
1895/// restart every finished source from its first row.
1896fn walking<'a>(
1897    selected: Vec<&'a ResolvedSource>,
1898    states: &Option<Resumption>,
1899    kind: StreamKind,
1900) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1901    let mut ready = Vec::new();
1902    let mut starts = Vec::new();
1903    for source in selected {
1904        if let Some(resume) = resume_at(states, source.name(), kind) {
1905            ready.push(source);
1906            starts.push(resume);
1907        }
1908    }
1909    (ready, starts)
1910}
1911
1912/// A fingerprint of everything about a query that decides which rows it returns, and in
1913/// what order — the verb, the sources it addresses, and every filter it carries.
1914///
1915/// Written from the request's own [`Debug`] rather than field by field, and that is the
1916/// point: a filter added to `Filters` next year joins the fingerprint by existing. A
1917/// hand-written canonical form would keep compiling with the new field missing, and the
1918/// tokens it minted would silently stop distinguishing the queries that differ by it —
1919/// which is the whole failure this exists to prevent, reintroduced quietly.
1920///
1921/// Hashed rather than carried whole so a token stays a thing a person can paste. This is
1922/// not a signature and there is nothing secret in a token — see [`PageToken`]. It detects
1923/// a caller resuming the wrong walk, which is a mistake rather than an attack, so FNV-1a
1924/// is enough and needs no dependency the supply-chain gate would then have to weigh.
1925///
1926/// [`Debug`] output is not promised to be stable across compiler releases, and that is
1927/// survivable here: a token outstanding across a rebuild is refused with the message
1928/// above rather than honoured wrongly, which is the safe direction to fail in.
1929fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1930    let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1931    fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1932}
1933
1934/// FNV-1a over `text`, as sixteen hex digits.
1935fn fingerprint(text: &str) -> String {
1936    let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1937    for byte in text.as_bytes() {
1938        hash ^= u64::from(*byte);
1939        hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1940    }
1941    format!("{hash:016x}")
1942}
1943
1944/// The stream a token says is owed the next row, or `None` when it says none is.
1945///
1946/// A fresh query has no token and every token whose last round ended evenly carries no
1947/// such stream, so `None` is the common case and means "begin at the first stream".
1948fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1949    document.as_ref()?.owed.as_ref()
1950}
1951
1952/// Where one stream picks up, or `None` when the token says it is finished.
1953fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1954    match states {
1955        None => Some(Resume::default()),
1956        Some(document) => document
1957            .streams
1958            .iter()
1959            .find(|state| &state.source == source && state.stream == kind)
1960            .map(|state| state.resume.clone()),
1961    }
1962}
1963
1964/// Read a caller's page token against the sources this configuration has.
1965///
1966/// [`PageToken::parse`] has already established that the string is this engine's own
1967/// resume document — that is structural and happens where the caller's string enters.
1968/// What it cannot establish is that the document belongs *here*, because only the engine
1969/// knows which sources are configured and what page ceiling each one declares. So the
1970/// three things a token this engine wrote is always true of are checked here:
1971///
1972/// 1. every stream it names belongs to a configured source, so a token carried over from
1973///    another configuration is refused rather than quietly resuming half a walk;
1974/// 2. every stream it names is one **this verb reads**. A token minted by
1975///    `search --kind both` carries task and project streams, and `task list` reads
1976///    neither; without this it would find nothing to resume, drop every source, and
1977///    answer with an empty page and a zero exit — a wrong answer that looks like an
1978///    exhausted walk, which is the worst shape this failure could take;
1979/// 3. no stream appears twice, because a walk has one place to pick up per stream;
1980/// 4. no `skip` reaches a source's declared page ceiling, because the engine's `skip` is
1981///    an index among the surviving rows of one source page and can never reach it.
1982///
1983/// 5. the stream owed the next row, when the document names one, is a stream the
1984///    document also resumes.
1985///
1986/// A sixth thing needs no check here: at most one stream is owed the next row, because
1987/// [`Resumption`] holds that as one optional stream rather than as a flag on each of
1988/// them, so a document naming two has no spelling.
1989///
1990/// None of this is a security boundary and a page token is not a credential: nothing in
1991/// one is secret, and a forged cursor is handed straight back to the source that would
1992/// have issued it and refused there. What it buys is that a stale or hand-edited token
1993/// fails saying so, instead of silently returning a page from somewhere else in the walk.
1994fn resumption(
1995    engine: &Engine,
1996    token: Option<&PageToken>,
1997    reads: &[StreamKind],
1998    query: &str,
1999) -> Result<Option<Resumption>, EngineError> {
2000    let Some(document) = token.map(PageToken::decode) else {
2001        return Ok(None);
2002    };
2003
2004    // Every cursor below is an offset into the result set *one* query produced. Handed to
2005    // a different one — the same verb with another `--label`, another `--search`, another
2006    // `--direction` — each source picks up at a position that meant something in a walk
2007    // the caller is no longer doing, and the rows that come back are real rows at exit
2008    // zero. Nothing about that answer says it is arbitrary, which is what makes it worth
2009    // refusing rather than serving.
2010    if document.query != query {
2011        return Err(EngineError::Token {
2012            message: "this page token was written by a different query — resume the walk it \
2013                      came from, or drop --page to start this one from the beginning"
2014                .to_owned(),
2015        });
2016    }
2017    let states = &document.streams;
2018
2019    let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2020    for state in states {
2021        if !reads.contains(&state.stream) {
2022            return Err(EngineError::Token {
2023                message: format!(
2024                    "this page token resumes {}, which this command does not read — it \
2025                     was written by a different query",
2026                    state.stream.describe()
2027                ),
2028            });
2029        }
2030        let ceiling = engine
2031            .ready()
2032            .find(|source| source.name() == &state.source)
2033            .map(ceiling);
2034        if ceiling.is_none() && !engine.has(&state.source) {
2035            return Err(EngineError::Token {
2036                message: format!(
2037                    "this page token resumes a source called {:?}, which this \
2038                     configuration does not have",
2039                    state.source.as_str()
2040                ),
2041            });
2042        }
2043        if let Some(ceiling) = ceiling
2044            && state.resume.skip >= ceiling
2045        {
2046            return Err(EngineError::Token {
2047                message: format!(
2048                    "this page token resumes {} rows into a page of source {:?}, which \
2049                     serves at most {ceiling}",
2050                    state.resume.skip,
2051                    state.source.as_str()
2052                ),
2053            });
2054        }
2055        if seen.contains(&(&state.source, state.stream)) {
2056            return Err(EngineError::Token {
2057                message: format!(
2058                    "this page token gives source {:?} two places to resume from",
2059                    state.source.as_str()
2060                ),
2061            });
2062        }
2063        seen.push((&state.source, state.stream));
2064    }
2065
2066    // The stream owed the next row has to be one of the streams this document resumes.
2067    // Ignoring a stray one would be harmless in its effect — the merge would start at the
2068    // first stream instead — but it would be a value from outside accepted without a
2069    // reading, and the next thing to depend on it would inherit that.
2070    if let Some(owed) = &document.owed
2071        && !document
2072            .streams
2073            .iter()
2074            .any(|state| state.source == owed.source && state.stream == owed.stream)
2075    {
2076        return Err(EngineError::Token {
2077            message: format!(
2078                "this page token owes the next row to a stream it does not resume, \
2079                 {:?}'s {}",
2080                owed.source.as_str(),
2081                owed.stream.describe()
2082            ),
2083        });
2084    }
2085
2086    Ok(Some(document))
2087}
2088
2089/// Which project a task must belong to, as a source sees it.
2090///
2091/// A qualified id becomes a plain native one because by the time this runs the selection
2092/// holds only that id's own source — so there is no "some other source" case to get
2093/// wrong, and none to leave untested.
2094fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2095    match selector {
2096        ProjectSelector::Any => ProjectFilter::Any,
2097        ProjectSelector::Orphans => ProjectFilter::Orphans,
2098        ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2099        ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2100    }
2101}
2102
2103/// The predicates one text query is made of.
2104fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2105    match fields {
2106        TextFields::Title => vec![Predicate::SearchTitle],
2107        TextFields::Content => vec![Predicate::SearchContent],
2108        TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2109    }
2110}
2111
2112/// Whether a source searches **every** field this query names.
2113///
2114/// Every, not any: a `title-or-content` search pushed to a source that searches only
2115/// titles would come back missing every row that matches in the body alone — a narrower
2116/// result than the truth, which is the one thing compensation cannot repair. So a
2117/// half-capable source is not asked at all and the engine searches both fields itself.
2118fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2119    match fields {
2120        TextFields::Title => capabilities.search_title.is_native(),
2121        TextFields::Content => capabilities.search_content.is_native(),
2122        TextFields::TitleOrContent => {
2123            capabilities.search_title.is_native() && capabilities.search_content.is_native()
2124        }
2125    }
2126}
2127
2128/// Split a task query between the source and the engine.
2129fn shape_tasks(
2130    capabilities: &Capabilities,
2131    filters: &Filters,
2132    project: &ProjectFilter,
2133    priorities: &[Priority],
2134    commented_since: Option<DateTime<Utc>>,
2135) -> TaskShape {
2136    let mut pushed = TaskQuery::default();
2137    let mut local = LocalTasks::default();
2138    let mut outcomes = Outcomes::default();
2139
2140    if !filters.labels.is_empty() {
2141        if capabilities.filter_by_label.is_native() {
2142            pushed.labels = filters.labels.clone();
2143            outcomes.record(Predicate::Label, Outcome::PushedDown);
2144        } else {
2145            local.labels = Some(filters.labels.clone());
2146            outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2147        }
2148    }
2149    if !filters.statuses.is_empty() {
2150        if capabilities.filter_by_status.is_native() {
2151            pushed.statuses.clone_from(&filters.statuses);
2152            outcomes.record(Predicate::Status, Outcome::PushedDown);
2153        } else {
2154            local.statuses.clone_from(&filters.statuses);
2155            outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2156        }
2157    }
2158    if !priorities.is_empty() {
2159        if capabilities.filter_by_priority.is_native() {
2160            pushed.priorities = priorities.to_vec();
2161            outcomes.record(Predicate::Priority, Outcome::PushedDown);
2162        } else {
2163            local.priorities = priorities.to_vec();
2164            outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2165        }
2166    }
2167    if let Some(since) = commented_since {
2168        if capabilities.filter_by_comment_activity.is_native() {
2169            pushed.commented_since = Some(since);
2170            outcomes.record(Predicate::CommentedSince, Outcome::PushedDown);
2171        } else {
2172            local.commented_since = Some(since);
2173            outcomes.record(Predicate::CommentedSince, Outcome::AppliedLocally);
2174        }
2175    }
2176    if let Some(text) = &filters.text {
2177        let predicates = text_predicates(text.fields);
2178        if searches_natively(capabilities, text.fields) {
2179            pushed.text = Some(text.clone());
2180            outcomes.record_all(predicates, Outcome::PushedDown);
2181        } else {
2182            local.text = Some(text.clone());
2183            outcomes.record_all(predicates, Outcome::AppliedLocally);
2184        }
2185    }
2186    match project {
2187        ProjectFilter::Any => {}
2188        ProjectFilter::Orphans => {
2189            if capabilities.orphan_tasks.is_native() {
2190                pushed.project = ProjectFilter::Orphans;
2191                outcomes.record(Predicate::Project, Outcome::PushedDown);
2192            } else {
2193                local.project = Some(ProjectFilter::Orphans);
2194                outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2195            }
2196        }
2197        ProjectFilter::Is(id) => {
2198            if capabilities.projects.is_native() {
2199                pushed.project = ProjectFilter::Is(id.clone());
2200                outcomes.record(Predicate::Project, Outcome::PushedDown);
2201            } else {
2202                local.project = Some(ProjectFilter::Is(id.clone()));
2203                outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2204            }
2205        }
2206    }
2207
2208    TaskShape {
2209        pushed,
2210        local,
2211        outcomes,
2212    }
2213}
2214
2215/// Split a project query between the source and the engine.
2216fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2217    let mut pushed = ProjectQuery::default();
2218    let mut local = LocalProjects::default();
2219    let mut outcomes = Outcomes::default();
2220
2221    if !filters.labels.is_empty() {
2222        if capabilities.filter_by_label.is_native() {
2223            pushed.labels = filters.labels.clone();
2224            outcomes.record(Predicate::Label, Outcome::PushedDown);
2225        } else {
2226            local.labels = Some(filters.labels.clone());
2227            outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2228        }
2229    }
2230    if !filters.statuses.is_empty() {
2231        if capabilities.filter_by_status.is_native() {
2232            pushed.statuses.clone_from(&filters.statuses);
2233            outcomes.record(Predicate::Status, Outcome::PushedDown);
2234        } else {
2235            local.statuses.clone_from(&filters.statuses);
2236            outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2237        }
2238    }
2239    if let Some(text) = &filters.text {
2240        let predicates = text_predicates(text.fields);
2241        if searches_natively(capabilities, text.fields) {
2242            pushed.text = Some(text.clone());
2243            outcomes.record_all(predicates, Outcome::PushedDown);
2244        } else {
2245            local.text = Some(text.clone());
2246            outcomes.record_all(predicates, Outcome::AppliedLocally);
2247        }
2248    }
2249
2250    ProjectShape {
2251        pushed,
2252        local,
2253        outcomes,
2254    }
2255}
2256
2257/// Split a document query between the source and the engine.
2258///
2259/// The same three predicates a task query is split by, minus the status filter a document
2260/// has nothing to compare against.
2261fn shape_documents(
2262    capabilities: &Capabilities,
2263    filters: &DocumentFilters,
2264    project: &ProjectFilter,
2265) -> DocumentShape {
2266    let mut pushed = DocumentQuery::default();
2267    let mut local = LocalDocuments::default();
2268    let mut outcomes = Outcomes::default();
2269
2270    if !filters.labels.is_empty() {
2271        if capabilities.filter_by_label.is_native() {
2272            pushed.labels = filters.labels.clone();
2273            outcomes.record(Predicate::Label, Outcome::PushedDown);
2274        } else {
2275            local.labels = Some(filters.labels.clone());
2276            outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2277        }
2278    }
2279    if let Some(text) = &filters.text {
2280        let predicates = text_predicates(text.fields);
2281        if searches_natively(capabilities, text.fields) {
2282            pushed.text = Some(text.clone());
2283            outcomes.record_all(predicates, Outcome::PushedDown);
2284        } else {
2285            local.text = Some(text.clone());
2286            outcomes.record_all(predicates, Outcome::AppliedLocally);
2287        }
2288    }
2289    match project {
2290        ProjectFilter::Any => {}
2291        ProjectFilter::Orphans => {
2292            if capabilities.orphan_tasks.is_native() {
2293                pushed.project = ProjectFilter::Orphans;
2294                outcomes.record(Predicate::Project, Outcome::PushedDown);
2295            } else {
2296                local.project = Some(ProjectFilter::Orphans);
2297                outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2298            }
2299        }
2300        ProjectFilter::Is(id) => {
2301            if capabilities.projects.is_native() {
2302                pushed.project = ProjectFilter::Is(id.clone());
2303                outcomes.record(Predicate::Project, Outcome::PushedDown);
2304            } else {
2305                local.project = Some(ProjectFilter::Is(id.clone()));
2306                outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2307            }
2308        }
2309    }
2310
2311    DocumentShape {
2312        pushed,
2313        local,
2314        outcomes,
2315    }
2316}
2317
2318/// Split one entity's half of a search between the source and the engine.
2319fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2320    match stream {
2321        StreamKind::Projects => {
2322            let shaped = shape_projects(capabilities, filters);
2323            HitShape {
2324                stream,
2325                tasks: TaskQuery::default(),
2326                projects: shaped.pushed,
2327                local_tasks: LocalTasks::default(),
2328                local_projects: shaped.local,
2329                outcomes: shaped.outcomes,
2330            }
2331        }
2332        StreamKind::Items | StreamKind::Tasks => {
2333            let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any, &[], None);
2334            HitShape {
2335                stream,
2336                tasks: shaped.pushed,
2337                projects: ProjectQuery::default(),
2338                local_tasks: shaped.local,
2339                local_projects: LocalProjects::default(),
2340                outcomes: shaped.outcomes,
2341            }
2342        }
2343    }
2344}
2345
2346/// How large a page to ask a source for.
2347///
2348/// Exactly what is needed when every predicate went down, and the source's own ceiling
2349/// when the engine is narrowing — because a compensating walk cannot know how many rows
2350/// of a page will survive, and asking for the caller's limit would turn one filtered
2351/// page into a page per surviving row.
2352fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2353    if compensating {
2354        ceiling
2355    } else {
2356        budget.min(ceiling)
2357    }
2358}
2359
2360/// The largest page this source will serve, never zero.
2361fn ceiling(source: &ResolvedSource) -> u32 {
2362    source.source().capabilities().max_page_size.max(1)
2363}
2364
2365/// Walk one source's tasks, narrowing whatever it did not apply itself.
2366async fn fetch_tasks(
2367    source: &ResolvedSource,
2368    shape: &TaskShape,
2369    start: &Resume,
2370    budget: u32,
2371    calls: &AtomicU32,
2372) -> Result<Fetched<Task>, SourceError> {
2373    let compensating = shape.local != LocalTasks::default();
2374    walk(
2375        start,
2376        budget,
2377        page_size(compensating, budget, ceiling(source)),
2378        |task| shape.local.keeps(task),
2379        |cursor, limit| async move {
2380            calls.fetch_add(1, Ordering::Relaxed);
2381            let request = PageRequest { cursor, limit };
2382            let page = source.source().query_tasks(&shape.pushed, &request).await?;
2383            match shape.local.commented_since {
2384                None => Ok(page),
2385                Some(since) => comment::commented_since(source, &shape.local, page, since).await,
2386            }
2387        },
2388    )
2389    .await
2390}
2391
2392/// Walk one source's projects, narrowing whatever it did not apply itself.
2393async fn fetch_projects(
2394    source: &ResolvedSource,
2395    shape: &ProjectShape,
2396    start: &Resume,
2397    budget: u32,
2398    calls: &AtomicU32,
2399) -> Result<Fetched<Project>, SourceError> {
2400    let compensating = shape.local != LocalProjects::default();
2401    walk(
2402        start,
2403        budget,
2404        page_size(compensating, budget, ceiling(source)),
2405        |project| shape.local.keeps(project),
2406        |cursor, limit| async move {
2407            calls.fetch_add(1, Ordering::Relaxed);
2408            let request = PageRequest { cursor, limit };
2409            source
2410                .source()
2411                .query_projects(&shape.pushed, &request)
2412                .await
2413        },
2414    )
2415    .await
2416}
2417
2418/// Walk one source's documents, narrowing whatever it did not apply itself.
2419async fn fetch_documents(
2420    source: &ResolvedSource,
2421    shape: &DocumentShape,
2422    start: &Resume,
2423    budget: u32,
2424    calls: &AtomicU32,
2425) -> Result<Fetched<Document>, SourceError> {
2426    let compensating = shape.local != LocalDocuments::default();
2427    walk(
2428        start,
2429        budget,
2430        page_size(compensating, budget, ceiling(source)),
2431        |document| shape.local.keeps(document),
2432        |cursor, limit| async move {
2433            calls.fetch_add(1, Ordering::Relaxed);
2434            let request = PageRequest { cursor, limit };
2435            source
2436                .source()
2437                .query_documents(&shape.pushed, &request)
2438                .await
2439        },
2440    )
2441    .await
2442}
2443
2444/// Walk one source's labels. There is no predicate to compensate for.
2445async fn fetch_labels(
2446    source: &ResolvedSource,
2447    start: &Resume,
2448    budget: u32,
2449    calls: &AtomicU32,
2450) -> Result<Fetched<Label>, SourceError> {
2451    walk(
2452        start,
2453        budget,
2454        page_size(false, budget, ceiling(source)),
2455        |_| true,
2456        |cursor, limit| async move {
2457            calls.fetch_add(1, Ordering::Relaxed);
2458            let request = PageRequest { cursor, limit };
2459            source.source().labels(&request).await
2460        },
2461    )
2462    .await
2463}
2464
2465/// Walk one entity's half of a search.
2466async fn fetch_hits(
2467    source: &ResolvedSource,
2468    shape: &HitShape,
2469    start: &Resume,
2470    budget: u32,
2471    calls: &AtomicU32,
2472) -> Result<Fetched<Found>, SourceError> {
2473    let ceiling = ceiling(source);
2474    match shape.stream {
2475        StreamKind::Projects => {
2476            let compensating = shape.local_projects != LocalProjects::default();
2477            walk(
2478                start,
2479                budget,
2480                page_size(compensating, budget, ceiling),
2481                |found| match found {
2482                    Found::Project(project) => shape.local_projects.keeps(project),
2483                    Found::Task(_) => true,
2484                },
2485                |cursor, limit| async move {
2486                    calls.fetch_add(1, Ordering::Relaxed);
2487                    let request = PageRequest { cursor, limit };
2488                    let page = source
2489                        .source()
2490                        .query_projects(&shape.projects, &request)
2491                        .await?;
2492                    Ok(Page {
2493                        items: page.items.into_iter().map(Found::Project).collect(),
2494                        next: page.next,
2495                    })
2496                },
2497            )
2498            .await
2499        }
2500        StreamKind::Items | StreamKind::Tasks => {
2501            let compensating = shape.local_tasks != LocalTasks::default();
2502            walk(
2503                start,
2504                budget,
2505                page_size(compensating, budget, ceiling),
2506                |found| match found {
2507                    Found::Task(task) => shape.local_tasks.keeps(task),
2508                    Found::Project(_) => true,
2509                },
2510                |cursor, limit| async move {
2511                    calls.fetch_add(1, Ordering::Relaxed);
2512                    let request = PageRequest { cursor, limit };
2513                    let page = source.source().query_tasks(&shape.tasks, &request).await?;
2514                    Ok(Page {
2515                        items: page.items.into_iter().map(Found::Task).collect(),
2516                        next: page.next,
2517                    })
2518                },
2519            )
2520            .await
2521        }
2522    }
2523}
2524
2525/// One page of an item's forward edges.
2526async fn forward_edges(
2527    source: &ResolvedSource,
2528    entity: Entity,
2529    id: &NativeId,
2530    request: &PageRequest,
2531) -> Result<Page<DependencyEdge>, SourceError> {
2532    match entity {
2533        Entity::Task => {
2534            source
2535                .source()
2536                .task_dependencies(id, Direction::DependsOn, request)
2537                .await
2538        }
2539        Entity::Project => {
2540            source
2541                .source()
2542                .project_dependencies(id, Direction::DependsOn, request)
2543                .await
2544        }
2545    }
2546}
2547
2548/// Walk one item's dependency edges, emulating the reverse direction when the source
2549/// only reports forward ones.
2550///
2551/// The emulation is the bounded page-by-page scan the contract describes: a page of the
2552/// source's items, each asked for its own forward edges, keeping the ones that point at
2553/// `native`. It is indexless by construction — nothing is retained between pages beyond
2554/// the caller's own page — which is why a source that cannot walk backwards costs a scan
2555/// rather than a stored reverse index.
2556#[expect(
2557    clippy::too_many_arguments,
2558    reason = "every argument is one axis of one walk — the source, the item, the \
2559              direction, which of its two graphs, whether the reverse is emulated, where \
2560              to resume, how many rows to return and where to count calls. Grouping them \
2561              into a struct would name the same eight values one indirection further from \
2562              the loop that reads them."
2563)]
2564async fn fetch_edges(
2565    source: &ResolvedSource,
2566    native: &NativeId,
2567    direction: Direction,
2568    entity: Entity,
2569    emulating: bool,
2570    start: &Resume,
2571    budget: u32,
2572    calls: &AtomicU32,
2573) -> Result<Fetched<DependencyEdge>, SourceError> {
2574    let ceiling = ceiling(source);
2575    if !emulating {
2576        return walk(
2577            start,
2578            budget,
2579            page_size(false, budget, ceiling),
2580            |_| true,
2581            |cursor, limit| async move {
2582                calls.fetch_add(1, Ordering::Relaxed);
2583                let request = PageRequest { cursor, limit };
2584                match entity {
2585                    Entity::Task => {
2586                        source
2587                            .source()
2588                            .task_dependencies(native, direction, &request)
2589                            .await
2590                    }
2591                    Entity::Project => {
2592                        source
2593                            .source()
2594                            .project_dependencies(native, direction, &request)
2595                            .await
2596                    }
2597                }
2598            },
2599        )
2600        .await;
2601    }
2602
2603    walk(
2604        start,
2605        budget,
2606        ceiling,
2607        |_| true,
2608        |cursor, limit| async move {
2609            calls.fetch_add(1, Ordering::Relaxed);
2610            let request = PageRequest { cursor, limit };
2611            let (ids, next) = match entity {
2612                Entity::Task => {
2613                    let page = source
2614                        .source()
2615                        .query_tasks(&TaskQuery::default(), &request)
2616                        .await?;
2617                    let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2618                    (ids, page.next)
2619                }
2620                Entity::Project => {
2621                    let page = source
2622                        .source()
2623                        .query_projects(&ProjectQuery::default(), &request)
2624                        .await?;
2625                    let ids: Vec<NativeId> =
2626                        page.items.into_iter().map(|project| project.id).collect();
2627                    (ids, page.next)
2628                }
2629            };
2630
2631            let mut edges = Vec::new();
2632            for id in ids {
2633                let mut inner: Option<Cursor> = None;
2634                loop {
2635                    calls.fetch_add(1, Ordering::Relaxed);
2636                    let request = PageRequest {
2637                        cursor: inner.clone(),
2638                        limit,
2639                    };
2640                    let page = forward_edges(source, entity, &id, &request).await?;
2641                    // The inner half of the same bound the walk holds on its own pages:
2642                    // this scan keeps every matching edge of one source page, so a source
2643                    // that overruns here overruns the engine's memory just as surely.
2644                    fits(page.items.len(), limit)?;
2645                    edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2646                    unrepeated(
2647                        page.next.as_ref(),
2648                        inner.as_ref(),
2649                        "its forward edges were being scanned",
2650                    )?;
2651                    match page.next {
2652                        Some(cursor) => inner = Some(cursor),
2653                        None => break,
2654                    }
2655                }
2656            }
2657
2658            Ok(Page { items: edges, next })
2659        },
2660    )
2661    .await
2662}