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