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