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