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