Skip to main content

onetaskgraph_core/engine/
mod.rs

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