Skip to main content

onetaskgraph_core/engine/
mod.rs

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