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