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