Skip to main content

onetaskgraph_core/engine/
mod.rs

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