Skip to main content

onetaskgraph_core/engine/
mod.rs

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