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