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, TaskDetails};
56pub(crate) use copy::malformed_links;
57pub use copy::{
58    BudgetSpent, CopyAction, CopyItems, CopyLink, CopyLookup, CopyOutcome, CopyReport, CopyRequest,
59    CopyScope, 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    /// `--create` was given beside a flag that says how to look for a counterpart, when
749    /// `--create` says there is none to look for.
750    #[error(
751        "--create cannot be given with {flag}: --create asserts the destination holds no \
752         counterpart, so there is nothing for {flag} to look for\n\
753         next: drop {flag} to create each item without looking, or drop --create to look."
754    )]
755    CreateWith {
756        /// The flag given beside `--create`.
757        flag: CopyLookup,
758    },
759
760    /// `--create` named an item that itself records a counterpart at the destination.
761    ///
762    /// Refused rather than created: the item's own `onetaskgraph.copies` link — or its
763    /// origin — names an item there, so the caller's assertion that the destination holds
764    /// none is wrong for it, and creating another would duplicate it.
765    #[error(
766        "{item} already records a counterpart at the destination, {carrier}, and --create \
767         asserts it has none\n\
768         next: copy it without --create, which updates {carrier}."
769    )]
770    CreateCarried {
771        /// The item being copied.
772        item: GlobalId,
773        /// The destination item its link or its origin names.
774        carrier: GlobalId,
775    },
776
777    /// A member copy named a task that is not a member of the project being copied.
778    ///
779    /// Refused before anything is written, because a copy that names members names the
780    /// part of one project it carries: a task of some other project, or of no project, is
781    /// not a narrower version of that copy but a different one.
782    #[error(
783        "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
784         next: name a task `onetaskgraph task list --project {project}` reports, or copy \
785         {id} on its own with `onetaskgraph task copy`.",
786        project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
787    )]
788    NotAMember {
789        /// The id that was named as a member.
790        id: GlobalId,
791        /// The projects the copy carries, none of which it is a member of.
792        projects: Vec<GlobalId>,
793    },
794
795    /// A member copy's item depends on a member it was not told to carry, and that member
796    /// records no origin naming the destination.
797    ///
798    /// Refused before anything is written. The destination id of a member the copy does
799    /// not carry comes from that member's own recorded origin and from nowhere else — not
800    /// from a walk of the destination, which is the read a member copy exists to avoid —
801    /// and writing the edge without it would either drop it silently or point it at the
802    /// id the member has at its source, which the destination has never heard of.
803    #[error(
804        "{item} depends on {member}, which this copy was not told to carry and which records \
805         no origin in {destination}\n\
806         next: name {member} with --member as well, record its {destination} id at \
807         onetaskgraph.origin, or copy the whole project without --member."
808    )]
809    UnrecordedMember {
810        /// The item whose edge could not be resolved.
811        item: GlobalId,
812        /// The member that edge points at.
813        member: GlobalId,
814        /// The destination it records no origin in.
815        destination: SourceName,
816    },
817
818    /// A source refused something the copy asked of it.
819    ///
820    /// Distinct from a [`SourceFailure`], which leaves the other sources' results
821    /// standing: a copy is one write into one destination, and half of one is not an
822    /// answer.
823    #[error(
824        "source {name} could not do it: {error}\n\
825         next: fix what the source named above, then copy again."
826    )]
827    SourceRefused {
828        /// The source that refused.
829        name: String,
830        /// What it said.
831        error: SourceError,
832    },
833
834    /// A copy failed, and the destination could not be put back the way it was found.
835    ///
836    /// A copy is either complete or it never happened, and the engine undoes its own
837    /// writes to make that true. When the destination will not take one of them back, the
838    /// failure has to say so and name what is still there: a user told only that the copy
839    /// failed would copy again over a destination nobody has described to them, which is
840    /// the retry that trips a hosted destination's rate limiter.
841    #[error(
842        "the copy failed and could not be undone.\n\
843         it failed because: {error}\n\
844         it could not be undone because: {refusal}\n\
845         so these still hold what it wrote: {left_behind}\n\
846         next: remove or put back those items, then copy again."
847    )]
848    CopyNotUndone {
849        /// Why the copy failed in the first place.
850        error: Box<EngineError>,
851        /// The items the copy created or overwrote and could not take back.
852        ///
853        /// The ids themselves rather than a sentence about them: a caller acting on this —
854        /// removing them, or reporting them — needs the ids, and the joined form only the
855        /// message needs is [`LeftBehind`]'s own `Display`.
856        left_behind: LeftBehind,
857        /// What the destination said when the copy tried to take them back.
858        refusal: SourceError,
859    },
860}
861
862/// The items a failed copy left at the destination — one of them at least, always.
863///
864/// A plain `Vec` here would let [`EngineError::CopyNotUndone`] be built naming nothing
865/// still there, and naming what is still there is the whole reason that failure is
866/// distinct from the one it wraps: a user told only that the copy could not be undone,
867/// and then handed an empty list, has been told about a destination nobody described to
868/// them, which is exactly the blind retry this mechanism exists to remove. The first id
869/// is a field of its own, so the empty case cannot be written down.
870#[derive(Debug, Clone, PartialEq)]
871pub struct LeftBehind {
872    /// The first item the destination would not take back.
873    first: GlobalId,
874    /// The ones it would not take back after that, in the order it refused them.
875    rest: Vec<GlobalId>,
876}
877
878impl LeftBehind {
879    /// The list holding the one item every such refusal has to name.
880    #[must_use]
881    pub fn new(first: GlobalId) -> Self {
882        Self {
883            first,
884            rest: Vec::new(),
885        }
886    }
887
888    /// Record another item the destination would not take back.
889    pub fn push(&mut self, id: GlobalId) {
890        self.rest.push(id);
891    }
892
893    /// Every item still at the destination, in the order the destination refused them.
894    pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
895        std::iter::once(&self.first).chain(self.rest.iter())
896    }
897}
898
899impl std::fmt::Display for LeftBehind {
900    /// The qualified ids, comma-separated — the form the failure message reads in.
901    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
902        write!(formatter, "{}", self.first)?;
903        for id in &self.rest {
904            write!(formatter, ", {id}")?;
905        }
906        Ok(())
907    }
908}
909
910/// One configured source, in exactly one of the two states a configured source has.
911///
912/// A sum rather than two lists side by side: with a `ready` vector and an `unavailable`
913/// one, a source appearing in both is a shape the type permits and every reader has to
914/// decide about — and they would not all decide the same way, since one of them fans a
915/// query out and another names failures.
916pub enum ConfiguredSource {
917    /// It built, and answers queries.
918    Ready(ResolvedSource),
919    /// It did not, and every response says so instead.
920    Unavailable(UnavailableSource),
921}
922
923impl ConfiguredSource {
924    /// The name the configuration gave it, whichever state it is in.
925    #[must_use]
926    pub fn name(&self) -> &SourceName {
927        match self {
928            Self::Ready(source) => source.name(),
929            Self::Unavailable(source) => source.name(),
930        }
931    }
932}
933
934/// The sources a configuration resolved to, and the queries they answer.
935pub struct Engine {
936    /// Every configured source, in configured-name order, each in one state.
937    sources: Vec<ConfiguredSource>,
938    /// Which sources answer when a request names none.
939    selection: Vec<SourceName>,
940}
941
942impl Engine {
943    /// Build every source a configuration names.
944    ///
945    /// A source whose plugin refuses to build — a credential that is not there, a
946    /// plugin whose implementation has not landed — is **not** fatal: it becomes an
947    /// entry in every response's `errors`, exactly as a source that fails mid-query
948    /// does, and the other sources still answer. A user with three sources and one
949    /// expired token gets the other two rather than nothing.
950    #[must_use]
951    pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
952        let (ready, unavailable) = resolve_available(config, secrets);
953        Self::new(
954            ready
955                .into_iter()
956                .map(ConfiguredSource::Ready)
957                .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
958                .collect(),
959            config.selected_sources(),
960        )
961    }
962
963    /// Drive sources built elsewhere — the engine's own tests, and any caller holding a
964    /// source it did not resolve from a configuration document.
965    #[must_use]
966    pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
967        Self { sources, selection }
968    }
969
970    /// Every source that built, in configured-name order.
971    fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
972        self.sources.iter().filter_map(|source| match source {
973            ConfiguredSource::Ready(ready) => Some(ready),
974            ConfiguredSource::Unavailable(_) => None,
975        })
976    }
977
978    /// Every source that did not, in the same order.
979    fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
980        self.sources.iter().filter_map(|source| match source {
981            ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
982            ConfiguredSource::Ready(_) => None,
983        })
984    }
985
986    /// Every configured source, whether or not it built, in name order.
987    #[must_use]
988    pub fn listing(&self) -> Vec<SourceListing> {
989        let mut listings: Vec<SourceListing> = self
990            .ready()
991            .map(|source| SourceListing {
992                source: source.name().clone(),
993                kind: source.kind().to_owned(),
994                state: SourceState::Available {
995                    capabilities: source.source().capabilities(),
996                },
997            })
998            .chain(self.unavailable().map(|source| SourceListing {
999                source: source.name().clone(),
1000                kind: source.kind().to_owned(),
1001                state: SourceState::Unavailable {
1002                    error: source.error().clone(),
1003                },
1004            }))
1005            .collect();
1006        listings.sort_by(|left, right| left.source.cmp(&right.source));
1007        listings
1008    }
1009
1010    /// Whether this configuration has a source called `name`, built or not.
1011    ///
1012    /// A caller reading a `--project` argument needs this: `urn:project:1` is a qualified
1013    /// id only if `urn` is a source here, and a native id full of colons otherwise. That
1014    /// rule cannot be applied without knowing what is configured.
1015    #[must_use]
1016    pub fn has(&self, name: &SourceName) -> bool {
1017        self.sources.iter().any(|source| source.name() == name)
1018    }
1019
1020    /// One page of tasks.
1021    ///
1022    /// # Errors
1023    ///
1024    /// Returns [`EngineError`] when the request names a source nothing configures, or
1025    /// carries a page token this engine did not issue. One source failing is not an
1026    /// error: it lands in the response's `errors`.
1027    pub async fn tasks(
1028        &self,
1029        request: &TaskRequest,
1030    ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1031        let mut names = self.resolve_selection(&request.sources)?;
1032        // A qualified project id names one project of one source, so no other source can
1033        // hold a task in it. Narrowing here means the plan reports the source that was
1034        // actually asked rather than a row of empty entries for sources that could not
1035        // have answered.
1036        if let ProjectSelector::Qualified(id) = &request.project {
1037            self.known(&id.source)?;
1038            names.retain(|name| name == &id.source);
1039        }
1040        let query = shape(
1041            "task-list",
1042            &names,
1043            &(
1044                &request.filters,
1045                &request.project,
1046                &request.priorities,
1047                &request.commented_since,
1048                &request.metadata,
1049                &request.origin,
1050            ),
1051        );
1052        let states = resumption(
1053            self,
1054            request.paging.token.as_ref(),
1055            &[StreamKind::Items],
1056            &query,
1057        )?;
1058        let budget = request.paging.limit.get();
1059
1060        let mut answer = Answer::new();
1061        let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1062
1063        let shapes: Vec<TaskShape> = ready
1064            .iter()
1065            .map(|source| {
1066                shape_tasks(
1067                    &source.source().capabilities(),
1068                    &request.filters,
1069                    &project_filter(&request.project),
1070                    &request.priorities,
1071                    request.commented_since,
1072                    &request.metadata,
1073                    request.origin.as_ref(),
1074                )
1075            })
1076            .collect();
1077        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1078        let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1079
1080        let walks = ready
1081            .iter()
1082            .enumerate()
1083            .map(|(index, source)| {
1084                fetch_tasks(
1085                    source,
1086                    &shapes[index],
1087                    &starts[index],
1088                    budget,
1089                    &counters[index],
1090                )
1091            })
1092            .collect();
1093
1094        let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1095        answer.finish(
1096            streams,
1097            budget,
1098            owed(&states),
1099            &query,
1100            |name, task: Task| {
1101                delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1102            },
1103        )
1104    }
1105
1106    /// One page of projects.
1107    ///
1108    /// # Errors
1109    ///
1110    /// As [`tasks`](Self::tasks).
1111    pub async fn projects(
1112        &self,
1113        request: &ProjectRequest,
1114    ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1115        let names = self.resolve_selection(&request.sources)?;
1116        let query = shape("project-list", &names, &request.filters);
1117        let states = resumption(
1118            self,
1119            request.paging.token.as_ref(),
1120            &[StreamKind::Items],
1121            &query,
1122        )?;
1123        let budget = request.paging.limit.get();
1124
1125        let mut answer = Answer::new();
1126        // A source declaring `projects: unsupported` has no project table at all, so
1127        // there is nothing to compensate for and nothing to ask: the predicate is
1128        // reported unavailable and that source contributes no rows. This is the one
1129        // outcome the engine cannot narrow its way out of, which is what `unavailable`
1130        // in the plan is for.
1131        let mut with_projects = Vec::new();
1132        for source in answer.split(self, &names) {
1133            if source.source().capabilities().projects.is_native() {
1134                with_projects.push(source);
1135            } else {
1136                answer.unreachable_predicate(source, Predicate::Project);
1137            }
1138        }
1139        let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1140
1141        let shapes: Vec<ProjectShape> = ready
1142            .iter()
1143            .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1144            .collect();
1145        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1146        let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1147
1148        let walks = ready
1149            .iter()
1150            .enumerate()
1151            .map(|(index, source)| {
1152                fetch_projects(
1153                    source,
1154                    &shapes[index],
1155                    &starts[index],
1156                    budget,
1157                    &counters[index],
1158                )
1159            })
1160            .collect();
1161
1162        let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1163        answer.finish(
1164            streams,
1165            budget,
1166            owed(&states),
1167            &query,
1168            |name, project: Project| Qualified {
1169                id: GlobalId::new(name.clone(), project.id.clone()),
1170                item: project,
1171            },
1172        )
1173    }
1174
1175    /// One page of documents.
1176    ///
1177    /// A source declaring it has no documents is **not asked**: the declaration is read
1178    /// once here, that source contributes no rows, and the plan reports
1179    /// [`Predicate::Document`] unavailable for it. So a document list spanning a mixed set
1180    /// of sources answers with what the document-bearing ones hold and says nothing
1181    /// alarming about the others — the same shape a source with no project table takes.
1182    ///
1183    /// # Errors
1184    ///
1185    /// As [`tasks`](Self::tasks).
1186    pub async fn documents(
1187        &self,
1188        request: &DocumentRequest,
1189    ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1190        let mut names = self.resolve_selection(&request.sources)?;
1191        // A qualified project id names one project of one source, exactly as it does for
1192        // tasks, so no other source can hold a document in it.
1193        if let ProjectSelector::Qualified(id) = &request.project {
1194            self.known(&id.source)?;
1195            names.retain(|name| name == &id.source);
1196        }
1197        let query = shape(
1198            "document-list",
1199            &names,
1200            &(&request.filters, &request.project),
1201        );
1202        let states = resumption(
1203            self,
1204            request.paging.token.as_ref(),
1205            &[StreamKind::Items],
1206            &query,
1207        )?;
1208        let budget = request.paging.limit.get();
1209
1210        let mut answer = Answer::new();
1211        let mut with_documents = Vec::new();
1212        for source in answer.split(self, &names) {
1213            if source.source().capabilities().documents.is_native() {
1214                with_documents.push(source);
1215            } else {
1216                answer.unreachable_predicate(source, Predicate::Document);
1217            }
1218        }
1219        let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1220
1221        let shapes: Vec<DocumentShape> = ready
1222            .iter()
1223            .map(|source| {
1224                shape_documents(
1225                    &source.source().capabilities(),
1226                    &request.filters,
1227                    &project_filter(&request.project),
1228                )
1229            })
1230            .collect();
1231        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1232        let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1233
1234        let walks = ready
1235            .iter()
1236            .enumerate()
1237            .map(|(index, source)| {
1238                fetch_documents(
1239                    source,
1240                    &shapes[index],
1241                    &starts[index],
1242                    budget,
1243                    &counters[index],
1244                )
1245            })
1246            .collect();
1247
1248        let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1249        answer.finish(
1250            streams,
1251            budget,
1252            owed(&states),
1253            &query,
1254            |name, document: Document| Qualified {
1255                id: GlobalId::new(name.clone(), document.id.clone()),
1256                item: document,
1257            },
1258        )
1259    }
1260
1261    /// One page of labels.
1262    ///
1263    /// # Errors
1264    ///
1265    /// As [`tasks`](Self::tasks).
1266    pub async fn labels(
1267        &self,
1268        request: &LabelRequest,
1269    ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1270        let names = self.resolve_selection(&request.sources)?;
1271        let query = shape("label-list", &names, &());
1272        let states = resumption(
1273            self,
1274            request.paging.token.as_ref(),
1275            &[StreamKind::Items],
1276            &query,
1277        )?;
1278        let budget = request.paging.limit.get();
1279
1280        let mut answer = Answer::new();
1281        let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1282        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1283        let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1284
1285        let walks = ready
1286            .iter()
1287            .enumerate()
1288            .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1289            .collect();
1290
1291        let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1292        answer.finish(
1293            streams,
1294            budget,
1295            owed(&states),
1296            &query,
1297            |name, label: Label| Qualified {
1298                id: GlobalId::new(name.clone(), label.id.clone()),
1299                item: label,
1300            },
1301        )
1302    }
1303
1304    /// One page of search hits, over tasks, projects, or both.
1305    ///
1306    /// # Errors
1307    ///
1308    /// As [`tasks`](Self::tasks).
1309    pub async fn search(
1310        &self,
1311        request: &SearchRequest,
1312    ) -> Result<QueryResponse<SearchHit>, EngineError> {
1313        let names = self.resolve_selection(&request.sources)?;
1314        // The streams this search reads, which is what a token resuming it may name. A
1315        // `--kind both` walk that has exhausted one half carries only the other, so this
1316        // is what a token may name rather than what it must.
1317        let reads: &[StreamKind] = match request.kind {
1318            SearchKind::Tasks => &[StreamKind::Tasks],
1319            SearchKind::Projects => &[StreamKind::Projects],
1320            SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1321        };
1322        // The scope is deliberately not in the fingerprint: which streams a search covers
1323        // is exactly what `reads` checks below, name by name and with a message that says
1324        // which half a token names. Folding it in here would refuse the same mistake one
1325        // layer earlier and less clearly, and leave that check unreachable.
1326        let query = shape("search", &names, &request.text);
1327        let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1328        let budget = request.paging.limit.get();
1329        let filters = Filters {
1330            text: Some(request.text.clone()),
1331            ..Filters::default()
1332        };
1333
1334        let mut answer = Answer::new();
1335
1336        // One stream per (source, entity), because a search over both entities reads two
1337        // result sets from each source and each has its own place to resume.
1338        let mut ready = Vec::new();
1339        let mut kinds = Vec::new();
1340        let mut starts = Vec::new();
1341        for source in answer.split(self, &names) {
1342            let mut streams = Vec::new();
1343            if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1344                streams.push(StreamKind::Tasks);
1345            }
1346            if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1347                if source.source().capabilities().projects.is_native() {
1348                    streams.push(StreamKind::Projects);
1349                } else {
1350                    answer.unreachable_predicate(source, Predicate::Project);
1351                }
1352            }
1353            for stream in streams {
1354                if let Some(resume) = resume_at(&states, source.name(), stream) {
1355                    ready.push(source);
1356                    kinds.push(stream);
1357                    starts.push(resume);
1358                }
1359            }
1360        }
1361
1362        let shapes: Vec<HitShape> = ready
1363            .iter()
1364            .zip(kinds.iter())
1365            .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1366            .collect();
1367        let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1368        let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1369
1370        let walks = ready
1371            .iter()
1372            .enumerate()
1373            .map(|(index, source)| {
1374                fetch_hits(
1375                    source,
1376                    &shapes[index],
1377                    &starts[index],
1378                    budget,
1379                    &counters[index],
1380                )
1381            })
1382            .collect();
1383
1384        let streams =
1385            answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1386        answer.finish(
1387            streams,
1388            budget,
1389            owed(&states),
1390            &query,
1391            |name, found: Found| match found {
1392                Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1393                    GlobalId::new(name.clone(), task.id.clone()),
1394                    task,
1395                )),
1396                Found::Project(project) => SearchHit::Project(Qualified {
1397                    id: GlobalId::new(name.clone(), project.id.clone()),
1398                    item: project,
1399                }),
1400            },
1401        )
1402    }
1403
1404    /// One task by its qualified id, or an empty page when there is no such task.
1405    ///
1406    /// # Errors
1407    ///
1408    /// Returns [`EngineError::UnknownSource`] when the id names a source nothing
1409    /// configures.
1410    pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1411        let name = self.known(&id.source)?;
1412        let mut answer = Answer::new();
1413        let selected = answer.split(self, std::slice::from_ref(&name));
1414        let Some(source) = selected.first() else {
1415            return answer.nothing();
1416        };
1417        let found = source.source().get_task(&id.native).await;
1418        let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1419        answer.one(source, found, |task| {
1420            delivery::qualified_task(qualified, task)
1421        })
1422    }
1423
1424    /// One project by its qualified id, or an empty page when there is no such project.
1425    ///
1426    /// # Errors
1427    ///
1428    /// As [`task`](Self::task).
1429    pub async fn project(
1430        &self,
1431        id: &GlobalId,
1432    ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1433        let name = self.known(&id.source)?;
1434        let mut answer = Answer::new();
1435        let selected = answer.split(self, std::slice::from_ref(&name));
1436        let Some(source) = selected.first() else {
1437            return answer.nothing();
1438        };
1439        let found = source.source().get_project(&id.native).await;
1440        let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1441        answer.one(source, found, |project| Qualified {
1442            id: qualified,
1443            item: project,
1444        })
1445    }
1446
1447    /// One document by its qualified id, or an empty page when there is no such document.
1448    ///
1449    /// A source declaring it has no documents holds none, so it is not asked and the
1450    /// answer is the empty page with [`Predicate::Document`] reported unavailable — the
1451    /// same answer the list verb gives, rather than a failure.
1452    ///
1453    /// # Errors
1454    ///
1455    /// As [`task`](Self::task).
1456    pub async fn document(
1457        &self,
1458        id: &GlobalId,
1459    ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1460        let name = self.known(&id.source)?;
1461        let mut answer = Answer::new();
1462        let selected = answer.split(self, std::slice::from_ref(&name));
1463        let Some(source) = selected.first() else {
1464            return answer.nothing();
1465        };
1466        if !source.source().capabilities().documents.is_native() {
1467            answer.unreachable_predicate(source, Predicate::Document);
1468            return answer.nothing();
1469        }
1470        let found = source.source().get_document(&id.native).await;
1471        let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1472        answer.one(source, found, |document| Qualified {
1473            id: qualified,
1474            item: document,
1475        })
1476    }
1477
1478    /// One page of a task's dependency edges.
1479    ///
1480    /// # Errors
1481    ///
1482    /// As [`task`](Self::task), plus [`EngineError::Token`] for a page token this engine
1483    /// did not issue.
1484    pub async fn task_dependencies(
1485        &self,
1486        request: &DependencyRequest,
1487    ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1488        self.dependencies(request, Entity::Task).await
1489    }
1490
1491    /// One page of a project's dependency edges.
1492    ///
1493    /// # Errors
1494    ///
1495    /// As [`task_dependencies`](Self::task_dependencies).
1496    pub async fn project_dependencies(
1497        &self,
1498        request: &DependencyRequest,
1499    ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1500        self.dependencies(request, Entity::Project).await
1501    }
1502
1503    /// Both dependency verbs, which differ only in which of a source's two edge sets
1504    /// they read and which of its two declarations governs the reverse direction.
1505    async fn dependencies(
1506        &self,
1507        request: &DependencyRequest,
1508        entity: Entity,
1509    ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1510        let name = self.known(&request.id.source)?;
1511        let query = shape(
1512            "dependencies",
1513            std::slice::from_ref(&name),
1514            &(entity, &request.id.native, request.direction),
1515        );
1516        let states = resumption(
1517            self,
1518            request.paging.token.as_ref(),
1519            &[StreamKind::Items],
1520            &query,
1521        )?;
1522        let budget = request.paging.limit.get();
1523
1524        let mut answer = Answer::new();
1525        let (ready, starts) = walking(
1526            answer.split(self, std::slice::from_ref(&name)),
1527            &states,
1528            StreamKind::Items,
1529        );
1530        let Some(source) = ready.first() else {
1531            return answer.nothing();
1532        };
1533
1534        let capabilities = source.source().capabilities();
1535        let support = match entity {
1536            Entity::Task => capabilities.task_dependencies,
1537            Entity::Project => capabilities.project_dependencies,
1538        };
1539        // `DependencySupport` has no unsupported variant on purpose: a dependency read is
1540        // answered natively or emulated by the scan below, never abandoned and never
1541        // silently empty.
1542        let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1543        let mut outcomes = Outcomes::default();
1544        if request.direction == Direction::DependedOnBy {
1545            if emulating {
1546                outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1547            } else {
1548                outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1549            }
1550        }
1551
1552        let counters = vec![AtomicU32::new(0)];
1553        let walked = fetch_edges(
1554            source,
1555            &request.id.native,
1556            request.direction,
1557            entity,
1558            emulating,
1559            &starts[0],
1560            budget,
1561            &counters[0],
1562        )
1563        .await;
1564
1565        let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1566        answer.finish(
1567            streams,
1568            budget,
1569            owed(&states),
1570            &query,
1571            |name, edge: DependencyEdge| QualifiedEdge {
1572                from: qualify_endpoint(name, edge.from),
1573                to: qualify_endpoint(name, edge.to),
1574                kind: edge.kind,
1575            },
1576        )
1577    }
1578
1579    /// The names a request addresses: the ones it gave, or the configuration's own.
1580    fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1581        if asked.is_empty() {
1582            if self.selection.is_empty() {
1583                return Err(EngineError::NoSources);
1584            }
1585            return Ok(self.selection.clone());
1586        }
1587        asked.iter().map(|name| self.known(name)).collect()
1588    }
1589
1590    /// `name` when this configuration has a source called that.
1591    fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1592        if self.has(name) {
1593            return Ok(name.clone());
1594        }
1595        if self.sources.is_empty() {
1596            return Err(EngineError::NoSources);
1597        }
1598        Err(EngineError::UnknownSource {
1599            name: name.to_string(),
1600            configured: self
1601                .listing()
1602                .iter()
1603                .map(|listing| listing.source.to_string())
1604                .collect::<Vec<_>>()
1605                .join(", "),
1606        })
1607    }
1608}
1609
1610fn qualify_endpoint(
1611    source: &SourceName,
1612    endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1613) -> QualifiedEndpoint {
1614    let kind = endpoint.kind;
1615    let is_qualified = endpoint.is_qualified();
1616    let endpoint_id = endpoint.into_id();
1617    QualifiedEndpoint {
1618        id: if is_qualified {
1619            endpoint_id
1620                .parse()
1621                .expect("plugin-api validates qualified dependency endpoints")
1622        } else {
1623            GlobalId::new(source.clone(), NativeId(endpoint_id))
1624        },
1625        kind,
1626    }
1627}
1628
1629/// Which of a source's two dependency graphs a request walks.
1630#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1631enum Entity {
1632    /// Task dependencies.
1633    Task,
1634    /// Project dependencies.
1635    Project,
1636}
1637
1638/// A search hit before it is qualified.
1639enum Found {
1640    /// A task matched.
1641    Task(Task),
1642    /// A project matched.
1643    Project(Project),
1644}
1645
1646/// What happened to one predicate against one source.
1647#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1648enum Outcome {
1649    /// Applied by the source itself.
1650    PushedDown,
1651    /// Applied by the engine over a wider result set.
1652    AppliedLocally,
1653    /// Answered by a bounded scan of the source.
1654    Emulated,
1655    /// Neither side could answer it, so this source contributed nothing for it.
1656    Unavailable,
1657}
1658
1659/// What happened to each predicate against one source.
1660///
1661/// Keyed by predicate, one outcome each, because those are the only states there are: a
1662/// predicate the source applied was not also applied here, and one nobody could answer
1663/// was not also pushed down. The four lists [`SourcePlan`] carries are this map fanned
1664/// out at the boundary — held *as* four lists a predicate could sit in all four at once,
1665/// and four contradictory claims about one predicate is the one thing the part of the
1666/// answer whose whole job is to say which of them is true must not be able to say.
1667///
1668/// A `BTreeMap` rather than a `HashMap` so the lists come out in one order and two runs
1669/// of a query render the same plan.
1670#[derive(Debug, Clone, Default, PartialEq)]
1671struct Outcomes(BTreeMap<Predicate, Outcome>);
1672
1673impl Outcomes {
1674    /// Record what happened to one predicate, replacing whatever was recorded before.
1675    ///
1676    /// Replacing rather than refusing: shaping a query decides each predicate once, and a
1677    /// second decision about the same one is the later one — there is no case here where
1678    /// both were meant to stand.
1679    fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1680        self.0.insert(predicate, outcome);
1681    }
1682
1683    /// Record the same outcome for several predicates, as a text search does for the two
1684    /// fields it covers.
1685    fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1686        for predicate in predicates {
1687            self.record(predicate, outcome);
1688        }
1689    }
1690
1691    /// The predicates this outcome befell, in the map's own stable order.
1692    fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1693        self.0
1694            .iter()
1695            .filter(|(_, recorded)| **recorded == outcome)
1696            .map(|(predicate, _)| *predicate)
1697            .collect()
1698    }
1699}
1700
1701/// The query one source sees, and the predicates left to the engine.
1702struct TaskShape {
1703    /// What the source is asked.
1704    pushed: TaskQuery,
1705    /// What the engine narrows afterwards.
1706    local: LocalTasks,
1707    /// What to report.
1708    outcomes: Outcomes,
1709}
1710
1711/// As [`TaskShape`], for projects.
1712struct ProjectShape {
1713    /// What the source is asked.
1714    pushed: ProjectQuery,
1715    /// What the engine narrows afterwards.
1716    local: LocalProjects,
1717    /// What to report.
1718    outcomes: Outcomes,
1719}
1720
1721/// As [`TaskShape`], for documents.
1722struct DocumentShape {
1723    /// What the source is asked.
1724    pushed: DocumentQuery,
1725    /// What the engine narrows afterwards.
1726    local: LocalDocuments,
1727    /// What to report.
1728    outcomes: Outcomes,
1729}
1730
1731/// As [`TaskShape`], for one entity's half of a search.
1732struct HitShape {
1733    /// Which entity this stream reads.
1734    stream: StreamKind,
1735    /// What a task stream asks.
1736    tasks: TaskQuery,
1737    /// What a project stream asks.
1738    projects: ProjectQuery,
1739    /// What the engine narrows afterwards, for tasks.
1740    local_tasks: LocalTasks,
1741    /// What the engine narrows afterwards, for projects.
1742    local_projects: LocalProjects,
1743    /// What to report.
1744    outcomes: Outcomes,
1745}
1746
1747/// The plan and the failures a response carries, accumulated as the verb runs.
1748///
1749/// One type rather than three parallel vectors threaded through every verb, because the
1750/// rule they enforce together is one rule: a source that fails contributes an error and
1751/// still leaves every other source's results standing.
1752struct Answer {
1753    /// One entry per source the engine addressed, merged by source at the end.
1754    plans: Vec<SourcePlan>,
1755    /// Every source that could not answer.
1756    errors: Vec<SourceFailure>,
1757}
1758
1759impl Answer {
1760    fn new() -> Self {
1761        Self {
1762            plans: Vec::new(),
1763            errors: Vec::new(),
1764        }
1765    }
1766
1767    /// The selected sources that built, recording the ones that did not as failures.
1768    ///
1769    /// A source that never built is reported and skipped rather than fatal: that is the
1770    /// same rule as a source failing mid-query, applied one step earlier.
1771    fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1772        let mut selected = Vec::new();
1773        for name in names {
1774            match engine.sources.iter().find(|source| source.name() == name) {
1775                Some(ConfiguredSource::Ready(source)) => selected.push(source),
1776                Some(ConfiguredSource::Unavailable(source)) => {
1777                    self.errors.push(source.failure());
1778                }
1779                None => {}
1780            }
1781        }
1782        selected
1783    }
1784
1785    /// Record that a source could answer nothing for `predicate`.
1786    fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1787        let mut outcomes = Outcomes::default();
1788        outcomes.record(predicate, Outcome::Unavailable);
1789        self.plans.push(plan_for(source, outcomes, 0));
1790    }
1791
1792    /// Turn each source's walk into a stream, keeping a failed source's failure.
1793    fn collect<T>(
1794        &mut self,
1795        ready: &[&ResolvedSource],
1796        walked: Vec<Result<Fetched<T>, SourceError>>,
1797        counters: &[AtomicU32],
1798        outcomes: Vec<Outcomes>,
1799    ) -> Vec<Stream<T>> {
1800        let kinds = vec![StreamKind::Items; ready.len()];
1801        self.collect_streams(ready, &kinds, walked, counters, outcomes)
1802    }
1803
1804    /// As [`collect`](Self::collect), where a source may contribute more than one stream.
1805    fn collect_streams<T>(
1806        &mut self,
1807        ready: &[&ResolvedSource],
1808        kinds: &[StreamKind],
1809        walked: Vec<Result<Fetched<T>, SourceError>>,
1810        counters: &[AtomicU32],
1811        outcomes: Vec<Outcomes>,
1812    ) -> Vec<Stream<T>> {
1813        let mut streams = Vec::new();
1814        for (index, result) in walked.into_iter().enumerate() {
1815            let source = ready[index];
1816            let pages = counters[index].load(Ordering::Relaxed);
1817            self.plans
1818                .push(plan_for(source, outcomes[index].clone(), pages));
1819            match result {
1820                Ok(fetched) => streams.push(Stream {
1821                    source: source.name().clone(),
1822                    kind: kinds[index],
1823                    fetched,
1824                }),
1825                // A stream that failed leaves the token, so a walk always terminates: a
1826                // source failing on every page would otherwise page forever.
1827                Err(error) => self.errors.push(SourceFailure {
1828                    source: source.name().clone(),
1829                    error,
1830                }),
1831            }
1832        }
1833        streams
1834    }
1835
1836    /// The response for a verb that reads exactly one item from exactly one source.
1837    fn one<T, U>(
1838        self,
1839        source: &ResolvedSource,
1840        found: Result<Option<T>, SourceError>,
1841        qualify: impl FnOnce(T) -> U,
1842    ) -> Result<QueryResponse<U>, EngineError> {
1843        Ok(self.one_response(source, found, qualify))
1844    }
1845
1846    /// [`one`](Self::one), for a caller with no refusal to thread through.
1847    fn one_response<T, U>(
1848        mut self,
1849        source: &ResolvedSource,
1850        found: Result<Option<T>, SourceError>,
1851        qualify: impl FnOnce(T) -> U,
1852    ) -> QueryResponse<U> {
1853        self.plans.push(plan_for(source, Outcomes::default(), 1));
1854        let items = match found {
1855            Ok(Some(item)) => vec![qualify(item)],
1856            Ok(None) => Vec::new(),
1857            Err(error) => {
1858                self.errors.push(SourceFailure {
1859                    source: source.name().clone(),
1860                    error,
1861                });
1862                Vec::new()
1863            }
1864        };
1865        QueryResponse {
1866            items,
1867            next: None,
1868            plan: QueryPlan {
1869                per_source: merge_plans(self.plans),
1870            },
1871            errors: self.errors,
1872        }
1873    }
1874
1875    /// The response for a verb with nothing left to ask.
1876    fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1877        Ok(QueryResponse {
1878            items: Vec::new(),
1879            next: None,
1880            plan: QueryPlan {
1881                per_source: merge_plans(self.plans),
1882            },
1883            errors: self.errors,
1884        })
1885    }
1886
1887    /// Merge the streams into the caller's page and mint the token that resumes it.
1888    ///
1889    /// `first` is the stream the token being resumed says is owed the next row, so the
1890    /// round-robin picks up where the previous page stopped rather than restarting.
1891    fn finish<T, U>(
1892        self,
1893        streams: Vec<Stream<T>>,
1894        budget: u32,
1895        first: Option<&Owed>,
1896        query: &str,
1897        qualify: impl Fn(&SourceName, T) -> U,
1898    ) -> Result<QueryResponse<U>, EngineError> {
1899        let (rows, states, owed) = merge(streams, budget, first);
1900        let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1901        Ok(QueryResponse {
1902            items: rows
1903                .into_iter()
1904                .map(|(name, item)| qualify(&name, item))
1905                .collect(),
1906            next,
1907            plan: QueryPlan {
1908                per_source: merge_plans(self.plans),
1909            },
1910            errors: self.errors,
1911        })
1912    }
1913}
1914
1915/// One source's plan entry: the outcomes fanned out into the four lists the contract's
1916/// [`SourcePlan`] carries, each in one order so two runs read the same.
1917fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1918    SourcePlan {
1919        source: source.name().clone(),
1920        kind: source.kind().to_owned(),
1921        pushed_down: outcomes.with(Outcome::PushedDown),
1922        applied_locally: outcomes.with(Outcome::AppliedLocally),
1923        emulated: outcomes.with(Outcome::Emulated),
1924        unavailable: outcomes.with(Outcome::Unavailable),
1925        pages_fetched: pages,
1926    }
1927}
1928
1929/// One entry per source, however many streams that source contributed.
1930///
1931/// `search --kind both` reads two streams from each source, and a plan is per source:
1932/// two entries for one name would say the engine addressed it twice.
1933fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1934    let mut merged: Vec<SourcePlan> = Vec::new();
1935    for plan in plans {
1936        if let Some(existing) = merged
1937            .iter_mut()
1938            .find(|existing| existing.source == plan.source)
1939        {
1940            existing.pushed_down.extend(plan.pushed_down);
1941            existing.applied_locally.extend(plan.applied_locally);
1942            existing.emulated.extend(plan.emulated);
1943            existing.unavailable.extend(plan.unavailable);
1944            existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1945            for list in [
1946                &mut existing.pushed_down,
1947                &mut existing.applied_locally,
1948                &mut existing.emulated,
1949                &mut existing.unavailable,
1950            ] {
1951                list.sort_unstable();
1952                list.dedup();
1953            }
1954        } else {
1955            merged.push(plan);
1956        }
1957    }
1958    merged
1959}
1960
1961/// The sources still walking, with where each picks up.
1962///
1963/// A token names every stream that has more to give, so a source **absent** from one has
1964/// been exhausted and is not asked again. Without that, the second page of a walk would
1965/// restart every finished source from its first row.
1966fn walking<'a>(
1967    selected: Vec<&'a ResolvedSource>,
1968    states: &Option<Resumption>,
1969    kind: StreamKind,
1970) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1971    let mut ready = Vec::new();
1972    let mut starts = Vec::new();
1973    for source in selected {
1974        if let Some(resume) = resume_at(states, source.name(), kind) {
1975            ready.push(source);
1976            starts.push(resume);
1977        }
1978    }
1979    (ready, starts)
1980}
1981
1982/// A fingerprint of everything about a query that decides which rows it returns, and in
1983/// what order — the verb, the sources it addresses, and every filter it carries.
1984///
1985/// Written from the request's own [`Debug`] rather than field by field, and that is the
1986/// point: a filter added to `Filters` next year joins the fingerprint by existing. A
1987/// hand-written canonical form would keep compiling with the new field missing, and the
1988/// tokens it minted would silently stop distinguishing the queries that differ by it —
1989/// which is the whole failure this exists to prevent, reintroduced quietly.
1990///
1991/// Hashed rather than carried whole so a token stays a thing a person can paste. This is
1992/// not a signature and there is nothing secret in a token — see [`PageToken`]. It detects
1993/// a caller resuming the wrong walk, which is a mistake rather than an attack, so FNV-1a
1994/// is enough and needs no dependency the supply-chain gate would then have to weigh.
1995///
1996/// [`Debug`] output is not promised to be stable across compiler releases, and that is
1997/// survivable here: a token outstanding across a rebuild is refused with the message
1998/// above rather than honoured wrongly, which is the safe direction to fail in.
1999fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
2000    let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
2001    fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
2002}
2003
2004/// FNV-1a over `text`, as sixteen hex digits.
2005fn fingerprint(text: &str) -> String {
2006    let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
2007    for byte in text.as_bytes() {
2008        hash ^= u64::from(*byte);
2009        hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
2010    }
2011    format!("{hash:016x}")
2012}
2013
2014/// The stream a token says is owed the next row, or `None` when it says none is.
2015///
2016/// A fresh query has no token and every token whose last round ended evenly carries no
2017/// such stream, so `None` is the common case and means "begin at the first stream".
2018fn owed(document: &Option<Resumption>) -> Option<&Owed> {
2019    document.as_ref()?.owed.as_ref()
2020}
2021
2022/// Where one stream picks up, or `None` when the token says it is finished.
2023fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
2024    match states {
2025        None => Some(Resume::default()),
2026        Some(document) => document
2027            .streams
2028            .iter()
2029            .find(|state| &state.source == source && state.stream == kind)
2030            .map(|state| state.resume.clone()),
2031    }
2032}
2033
2034/// Read a caller's page token against the sources this configuration has.
2035///
2036/// [`PageToken::parse`] has already established that the string is this engine's own
2037/// resume document — that is structural and happens where the caller's string enters.
2038/// What it cannot establish is that the document belongs *here*, because only the engine
2039/// knows which sources are configured and what page ceiling each one declares. So the
2040/// three things a token this engine wrote is always true of are checked here:
2041///
2042/// 1. every stream it names belongs to a configured source, so a token carried over from
2043///    another configuration is refused rather than quietly resuming half a walk;
2044/// 2. every stream it names is one **this verb reads**. A token minted by
2045///    `search --kind both` carries task and project streams, and `task list` reads
2046///    neither; without this it would find nothing to resume, drop every source, and
2047///    answer with an empty page and a zero exit — a wrong answer that looks like an
2048///    exhausted walk, which is the worst shape this failure could take;
2049/// 3. no stream appears twice, because a walk has one place to pick up per stream;
2050/// 4. no `skip` reaches a source's declared page ceiling, because the engine's `skip` is
2051///    an index among the surviving rows of one source page and can never reach it.
2052///
2053/// 5. the stream owed the next row, when the document names one, is a stream the
2054///    document also resumes.
2055///
2056/// A sixth thing needs no check here: at most one stream is owed the next row, because
2057/// [`Resumption`] holds that as one optional stream rather than as a flag on each of
2058/// them, so a document naming two has no spelling.
2059///
2060/// None of this is a security boundary and a page token is not a credential: nothing in
2061/// one is secret, and a forged cursor is handed straight back to the source that would
2062/// have issued it and refused there. What it buys is that a stale or hand-edited token
2063/// fails saying so, instead of silently returning a page from somewhere else in the walk.
2064fn resumption(
2065    engine: &Engine,
2066    token: Option<&PageToken>,
2067    reads: &[StreamKind],
2068    query: &str,
2069) -> Result<Option<Resumption>, EngineError> {
2070    let Some(document) = token.map(PageToken::decode) else {
2071        return Ok(None);
2072    };
2073
2074    // Every cursor below is an offset into the result set *one* query produced. Handed to
2075    // a different one — the same verb with another `--label`, another `--search`, another
2076    // `--direction` — each source picks up at a position that meant something in a walk
2077    // the caller is no longer doing, and the rows that come back are real rows at exit
2078    // zero. Nothing about that answer says it is arbitrary, which is what makes it worth
2079    // refusing rather than serving.
2080    if document.query != query {
2081        return Err(EngineError::Token {
2082            message: "this page token was written by a different query — resume the walk it \
2083                      came from, or drop --page to start this one from the beginning"
2084                .to_owned(),
2085        });
2086    }
2087    let states = &document.streams;
2088
2089    let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2090    for state in states {
2091        if !reads.contains(&state.stream) {
2092            return Err(EngineError::Token {
2093                message: format!(
2094                    "this page token resumes {}, which this command does not read — it \
2095                     was written by a different query",
2096                    state.stream.describe()
2097                ),
2098            });
2099        }
2100        let ceiling = engine
2101            .ready()
2102            .find(|source| source.name() == &state.source)
2103            .map(ceiling);
2104        if ceiling.is_none() && !engine.has(&state.source) {
2105            return Err(EngineError::Token {
2106                message: format!(
2107                    "this page token resumes a source called {:?}, which this \
2108                     configuration does not have",
2109                    state.source.as_str()
2110                ),
2111            });
2112        }
2113        if let Some(ceiling) = ceiling
2114            && state.resume.skip >= ceiling
2115        {
2116            return Err(EngineError::Token {
2117                message: format!(
2118                    "this page token resumes {} rows into a page of source {:?}, which \
2119                     serves at most {ceiling}",
2120                    state.resume.skip,
2121                    state.source.as_str()
2122                ),
2123            });
2124        }
2125        if seen.contains(&(&state.source, state.stream)) {
2126            return Err(EngineError::Token {
2127                message: format!(
2128                    "this page token gives source {:?} two places to resume from",
2129                    state.source.as_str()
2130                ),
2131            });
2132        }
2133        seen.push((&state.source, state.stream));
2134    }
2135
2136    // The stream owed the next row has to be one of the streams this document resumes.
2137    // Ignoring a stray one would be harmless in its effect — the merge would start at the
2138    // first stream instead — but it would be a value from outside accepted without a
2139    // reading, and the next thing to depend on it would inherit that.
2140    if let Some(owed) = &document.owed
2141        && !document
2142            .streams
2143            .iter()
2144            .any(|state| state.source == owed.source && state.stream == owed.stream)
2145    {
2146        return Err(EngineError::Token {
2147            message: format!(
2148                "this page token owes the next row to a stream it does not resume, \
2149                 {:?}'s {}",
2150                owed.source.as_str(),
2151                owed.stream.describe()
2152            ),
2153        });
2154    }
2155
2156    Ok(Some(document))
2157}
2158
2159/// Which project a task must belong to, as a source sees it.
2160///
2161/// A qualified id becomes a plain native one because by the time this runs the selection
2162/// holds only that id's own source — so there is no "some other source" case to get
2163/// wrong, and none to leave untested.
2164fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2165    match selector {
2166        ProjectSelector::Any => ProjectFilter::Any,
2167        ProjectSelector::Orphans => ProjectFilter::Orphans,
2168        ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2169        ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2170    }
2171}
2172
2173/// The predicates one text query is made of.
2174fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2175    match fields {
2176        TextFields::Title => vec![Predicate::SearchTitle],
2177        TextFields::Content => vec![Predicate::SearchContent],
2178        TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2179    }
2180}
2181
2182/// Whether a source searches **every** field this query names.
2183///
2184/// Every, not any: a `title-or-content` search pushed to a source that searches only
2185/// titles would come back missing every row that matches in the body alone — a narrower
2186/// result than the truth, which is the one thing compensation cannot repair. So a
2187/// half-capable source is not asked at all and the engine searches both fields itself.
2188fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2189    match fields {
2190        TextFields::Title => capabilities.search_title.is_native(),
2191        TextFields::Content => capabilities.search_content.is_native(),
2192        TextFields::TitleOrContent => {
2193            capabilities.search_title.is_native() && capabilities.search_content.is_native()
2194        }
2195    }
2196}
2197
2198/// Split a task query between the source and the engine.
2199fn shape_tasks(
2200    capabilities: &Capabilities,
2201    filters: &Filters,
2202    project: &ProjectFilter,
2203    priorities: &[Priority],
2204    commented_since: Option<DateTime<Utc>>,
2205    metadata: &[MetadataMatch],
2206    origin: Option<&GlobalId>,
2207) -> TaskShape {
2208    let mut pushed = TaskQuery::default();
2209    let mut local = LocalTasks::default();
2210    let mut outcomes = Outcomes::default();
2211
2212    if !filters.labels.is_empty() {
2213        if capabilities.filter_by_label.is_native() {
2214            pushed.labels = filters.labels.clone();
2215            outcomes.record(Predicate::Label, Outcome::PushedDown);
2216        } else {
2217            local.labels = Some(filters.labels.clone());
2218            outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2219        }
2220    }
2221    if !filters.statuses.is_empty() {
2222        if capabilities.filter_by_status.is_native() {
2223            pushed.statuses.clone_from(&filters.statuses);
2224            outcomes.record(Predicate::Status, Outcome::PushedDown);
2225        } else {
2226            local.statuses.clone_from(&filters.statuses);
2227            outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2228        }
2229    }
2230    if !priorities.is_empty() {
2231        if capabilities.filter_by_priority.is_native() {
2232            pushed.priorities = priorities.to_vec();
2233            outcomes.record(Predicate::Priority, Outcome::PushedDown);
2234        } else {
2235            local.priorities = priorities.to_vec();
2236            outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2237        }
2238    }
2239    if let Some(since) = commented_since {
2240        if capabilities.filter_by_comment_activity.is_native() {
2241            pushed.commented_since = Some(since);
2242            outcomes.record(Predicate::CommentedSince, Outcome::PushedDown);
2243        } else {
2244            local.commented_since = Some(since);
2245            outcomes.record(Predicate::CommentedSince, Outcome::AppliedLocally);
2246        }
2247    }
2248    if !metadata.is_empty() {
2249        if capabilities.filter_by_metadata.is_native() {
2250            pushed.metadata = metadata.to_vec();
2251            outcomes.record(Predicate::Metadata, Outcome::PushedDown);
2252        } else {
2253            local.metadata = metadata.to_vec();
2254            outcomes.record(Predicate::Metadata, Outcome::AppliedLocally);
2255        }
2256    }
2257    if let Some(origin) = origin {
2258        if capabilities.filter_by_origin.is_native() {
2259            // The qualified id as a copy stores it: a source compares the string and never
2260            // parses it.
2261            pushed.origin = Some(origin.to_string());
2262            outcomes.record(Predicate::Origin, Outcome::PushedDown);
2263        } else {
2264            local.origin = Some(origin.clone());
2265            outcomes.record(Predicate::Origin, Outcome::AppliedLocally);
2266        }
2267    }
2268    if let Some(text) = &filters.text {
2269        let predicates = text_predicates(text.fields);
2270        if searches_natively(capabilities, text.fields) {
2271            pushed.text = Some(text.clone());
2272            outcomes.record_all(predicates, Outcome::PushedDown);
2273        } else {
2274            local.text = Some(text.clone());
2275            outcomes.record_all(predicates, Outcome::AppliedLocally);
2276        }
2277    }
2278    match project {
2279        ProjectFilter::Any => {}
2280        ProjectFilter::Orphans => {
2281            if capabilities.orphan_tasks.is_native() {
2282                pushed.project = ProjectFilter::Orphans;
2283                outcomes.record(Predicate::Project, Outcome::PushedDown);
2284            } else {
2285                local.project = Some(ProjectFilter::Orphans);
2286                outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2287            }
2288        }
2289        ProjectFilter::Is(id) => {
2290            if capabilities.projects.is_native() {
2291                pushed.project = ProjectFilter::Is(id.clone());
2292                outcomes.record(Predicate::Project, Outcome::PushedDown);
2293            } else {
2294                local.project = Some(ProjectFilter::Is(id.clone()));
2295                outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2296            }
2297        }
2298    }
2299
2300    TaskShape {
2301        pushed,
2302        local,
2303        outcomes,
2304    }
2305}
2306
2307/// Split a project query between the source and the engine.
2308fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2309    let mut pushed = ProjectQuery::default();
2310    let mut local = LocalProjects::default();
2311    let mut outcomes = Outcomes::default();
2312
2313    if !filters.labels.is_empty() {
2314        if capabilities.filter_by_label.is_native() {
2315            pushed.labels = filters.labels.clone();
2316            outcomes.record(Predicate::Label, Outcome::PushedDown);
2317        } else {
2318            local.labels = Some(filters.labels.clone());
2319            outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2320        }
2321    }
2322    if !filters.statuses.is_empty() {
2323        if capabilities.filter_by_status.is_native() {
2324            pushed.statuses.clone_from(&filters.statuses);
2325            outcomes.record(Predicate::Status, Outcome::PushedDown);
2326        } else {
2327            local.statuses.clone_from(&filters.statuses);
2328            outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2329        }
2330    }
2331    if let Some(text) = &filters.text {
2332        let predicates = text_predicates(text.fields);
2333        if searches_natively(capabilities, text.fields) {
2334            pushed.text = Some(text.clone());
2335            outcomes.record_all(predicates, Outcome::PushedDown);
2336        } else {
2337            local.text = Some(text.clone());
2338            outcomes.record_all(predicates, Outcome::AppliedLocally);
2339        }
2340    }
2341
2342    ProjectShape {
2343        pushed,
2344        local,
2345        outcomes,
2346    }
2347}
2348
2349/// Split a document query between the source and the engine.
2350///
2351/// The same three predicates a task query is split by, minus the status filter a document
2352/// has nothing to compare against.
2353fn shape_documents(
2354    capabilities: &Capabilities,
2355    filters: &DocumentFilters,
2356    project: &ProjectFilter,
2357) -> DocumentShape {
2358    let mut pushed = DocumentQuery::default();
2359    let mut local = LocalDocuments::default();
2360    let mut outcomes = Outcomes::default();
2361
2362    if !filters.labels.is_empty() {
2363        if capabilities.filter_by_label.is_native() {
2364            pushed.labels = filters.labels.clone();
2365            outcomes.record(Predicate::Label, Outcome::PushedDown);
2366        } else {
2367            local.labels = Some(filters.labels.clone());
2368            outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2369        }
2370    }
2371    if let Some(text) = &filters.text {
2372        let predicates = text_predicates(text.fields);
2373        if searches_natively(capabilities, text.fields) {
2374            pushed.text = Some(text.clone());
2375            outcomes.record_all(predicates, Outcome::PushedDown);
2376        } else {
2377            local.text = Some(text.clone());
2378            outcomes.record_all(predicates, Outcome::AppliedLocally);
2379        }
2380    }
2381    match project {
2382        ProjectFilter::Any => {}
2383        ProjectFilter::Orphans => {
2384            if capabilities.orphan_tasks.is_native() {
2385                pushed.project = ProjectFilter::Orphans;
2386                outcomes.record(Predicate::Project, Outcome::PushedDown);
2387            } else {
2388                local.project = Some(ProjectFilter::Orphans);
2389                outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2390            }
2391        }
2392        ProjectFilter::Is(id) => {
2393            if capabilities.projects.is_native() {
2394                pushed.project = ProjectFilter::Is(id.clone());
2395                outcomes.record(Predicate::Project, Outcome::PushedDown);
2396            } else {
2397                local.project = Some(ProjectFilter::Is(id.clone()));
2398                outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2399            }
2400        }
2401    }
2402
2403    DocumentShape {
2404        pushed,
2405        local,
2406        outcomes,
2407    }
2408}
2409
2410/// Split one entity's half of a search between the source and the engine.
2411fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2412    match stream {
2413        StreamKind::Projects => {
2414            let shaped = shape_projects(capabilities, filters);
2415            HitShape {
2416                stream,
2417                tasks: TaskQuery::default(),
2418                projects: shaped.pushed,
2419                local_tasks: LocalTasks::default(),
2420                local_projects: shaped.local,
2421                outcomes: shaped.outcomes,
2422            }
2423        }
2424        StreamKind::Items | StreamKind::Tasks => {
2425            let shaped = shape_tasks(
2426                capabilities,
2427                filters,
2428                &ProjectFilter::Any,
2429                &[],
2430                None,
2431                &[],
2432                None,
2433            );
2434            HitShape {
2435                stream,
2436                tasks: shaped.pushed,
2437                projects: ProjectQuery::default(),
2438                local_tasks: shaped.local,
2439                local_projects: LocalProjects::default(),
2440                outcomes: shaped.outcomes,
2441            }
2442        }
2443    }
2444}
2445
2446/// How large a page to ask a source for.
2447///
2448/// Exactly what is needed when every predicate went down, and the source's own ceiling
2449/// when the engine is narrowing — because a compensating walk cannot know how many rows
2450/// of a page will survive, and asking for the caller's limit would turn one filtered
2451/// page into a page per surviving row.
2452fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2453    if compensating {
2454        ceiling
2455    } else {
2456        budget.min(ceiling)
2457    }
2458}
2459
2460/// The largest page this source will serve, never zero.
2461fn ceiling(source: &ResolvedSource) -> u32 {
2462    source.source().capabilities().max_page_size.max(1)
2463}
2464
2465/// Walk one source's tasks, narrowing whatever it did not apply itself.
2466async fn fetch_tasks(
2467    source: &ResolvedSource,
2468    shape: &TaskShape,
2469    start: &Resume,
2470    budget: u32,
2471    calls: &AtomicU32,
2472) -> Result<Fetched<Task>, SourceError> {
2473    let compensating = shape.local != LocalTasks::default();
2474    walk(
2475        start,
2476        budget,
2477        page_size(compensating, budget, ceiling(source)),
2478        |task| shape.local.keeps(task),
2479        |cursor, limit| async move {
2480            calls.fetch_add(1, Ordering::Relaxed);
2481            let request = PageRequest { cursor, limit };
2482            let page = source.source().query_tasks(&shape.pushed, &request).await?;
2483            match shape.local.commented_since {
2484                None => Ok(page),
2485                Some(since) => comment::commented_since(source, &shape.local, page, since).await,
2486            }
2487        },
2488    )
2489    .await
2490}
2491
2492/// Walk one source's projects, narrowing whatever it did not apply itself.
2493async fn fetch_projects(
2494    source: &ResolvedSource,
2495    shape: &ProjectShape,
2496    start: &Resume,
2497    budget: u32,
2498    calls: &AtomicU32,
2499) -> Result<Fetched<Project>, SourceError> {
2500    let compensating = shape.local != LocalProjects::default();
2501    walk(
2502        start,
2503        budget,
2504        page_size(compensating, budget, ceiling(source)),
2505        |project| shape.local.keeps(project),
2506        |cursor, limit| async move {
2507            calls.fetch_add(1, Ordering::Relaxed);
2508            let request = PageRequest { cursor, limit };
2509            source
2510                .source()
2511                .query_projects(&shape.pushed, &request)
2512                .await
2513        },
2514    )
2515    .await
2516}
2517
2518/// Walk one source's documents, narrowing whatever it did not apply itself.
2519async fn fetch_documents(
2520    source: &ResolvedSource,
2521    shape: &DocumentShape,
2522    start: &Resume,
2523    budget: u32,
2524    calls: &AtomicU32,
2525) -> Result<Fetched<Document>, SourceError> {
2526    let compensating = shape.local != LocalDocuments::default();
2527    walk(
2528        start,
2529        budget,
2530        page_size(compensating, budget, ceiling(source)),
2531        |document| shape.local.keeps(document),
2532        |cursor, limit| async move {
2533            calls.fetch_add(1, Ordering::Relaxed);
2534            let request = PageRequest { cursor, limit };
2535            source
2536                .source()
2537                .query_documents(&shape.pushed, &request)
2538                .await
2539        },
2540    )
2541    .await
2542}
2543
2544/// Walk one source's labels. There is no predicate to compensate for.
2545async fn fetch_labels(
2546    source: &ResolvedSource,
2547    start: &Resume,
2548    budget: u32,
2549    calls: &AtomicU32,
2550) -> Result<Fetched<Label>, SourceError> {
2551    walk(
2552        start,
2553        budget,
2554        page_size(false, budget, ceiling(source)),
2555        |_| true,
2556        |cursor, limit| async move {
2557            calls.fetch_add(1, Ordering::Relaxed);
2558            let request = PageRequest { cursor, limit };
2559            source.source().labels(&request).await
2560        },
2561    )
2562    .await
2563}
2564
2565/// Walk one entity's half of a search.
2566async fn fetch_hits(
2567    source: &ResolvedSource,
2568    shape: &HitShape,
2569    start: &Resume,
2570    budget: u32,
2571    calls: &AtomicU32,
2572) -> Result<Fetched<Found>, SourceError> {
2573    let ceiling = ceiling(source);
2574    match shape.stream {
2575        StreamKind::Projects => {
2576            let compensating = shape.local_projects != LocalProjects::default();
2577            walk(
2578                start,
2579                budget,
2580                page_size(compensating, budget, ceiling),
2581                |found| match found {
2582                    Found::Project(project) => shape.local_projects.keeps(project),
2583                    Found::Task(_) => true,
2584                },
2585                |cursor, limit| async move {
2586                    calls.fetch_add(1, Ordering::Relaxed);
2587                    let request = PageRequest { cursor, limit };
2588                    let page = source
2589                        .source()
2590                        .query_projects(&shape.projects, &request)
2591                        .await?;
2592                    Ok(Page {
2593                        items: page.items.into_iter().map(Found::Project).collect(),
2594                        next: page.next,
2595                    })
2596                },
2597            )
2598            .await
2599        }
2600        StreamKind::Items | StreamKind::Tasks => {
2601            let compensating = shape.local_tasks != LocalTasks::default();
2602            walk(
2603                start,
2604                budget,
2605                page_size(compensating, budget, ceiling),
2606                |found| match found {
2607                    Found::Task(task) => shape.local_tasks.keeps(task),
2608                    Found::Project(_) => true,
2609                },
2610                |cursor, limit| async move {
2611                    calls.fetch_add(1, Ordering::Relaxed);
2612                    let request = PageRequest { cursor, limit };
2613                    let page = source.source().query_tasks(&shape.tasks, &request).await?;
2614                    Ok(Page {
2615                        items: page.items.into_iter().map(Found::Task).collect(),
2616                        next: page.next,
2617                    })
2618                },
2619            )
2620            .await
2621        }
2622    }
2623}
2624
2625/// One page of an item's forward edges.
2626async fn forward_edges(
2627    source: &ResolvedSource,
2628    entity: Entity,
2629    id: &NativeId,
2630    request: &PageRequest,
2631) -> Result<Page<DependencyEdge>, SourceError> {
2632    match entity {
2633        Entity::Task => {
2634            source
2635                .source()
2636                .task_dependencies(id, Direction::DependsOn, request)
2637                .await
2638        }
2639        Entity::Project => {
2640            source
2641                .source()
2642                .project_dependencies(id, Direction::DependsOn, request)
2643                .await
2644        }
2645    }
2646}
2647
2648/// Walk one item's dependency edges, emulating the reverse direction when the source
2649/// only reports forward ones.
2650///
2651/// The emulation is the bounded page-by-page scan the contract describes: a page of the
2652/// source's items, each asked for its own forward edges, keeping the ones that point at
2653/// `native`. It is indexless by construction — nothing is retained between pages beyond
2654/// the caller's own page — which is why a source that cannot walk backwards costs a scan
2655/// rather than a stored reverse index.
2656#[expect(
2657    clippy::too_many_arguments,
2658    reason = "every argument is one axis of one walk — the source, the item, the \
2659              direction, which of its two graphs, whether the reverse is emulated, where \
2660              to resume, how many rows to return and where to count calls. Grouping them \
2661              into a struct would name the same eight values one indirection further from \
2662              the loop that reads them."
2663)]
2664async fn fetch_edges(
2665    source: &ResolvedSource,
2666    native: &NativeId,
2667    direction: Direction,
2668    entity: Entity,
2669    emulating: bool,
2670    start: &Resume,
2671    budget: u32,
2672    calls: &AtomicU32,
2673) -> Result<Fetched<DependencyEdge>, SourceError> {
2674    let ceiling = ceiling(source);
2675    if !emulating {
2676        return walk(
2677            start,
2678            budget,
2679            page_size(false, budget, ceiling),
2680            |_| true,
2681            |cursor, limit| async move {
2682                calls.fetch_add(1, Ordering::Relaxed);
2683                let request = PageRequest { cursor, limit };
2684                match entity {
2685                    Entity::Task => {
2686                        source
2687                            .source()
2688                            .task_dependencies(native, direction, &request)
2689                            .await
2690                    }
2691                    Entity::Project => {
2692                        source
2693                            .source()
2694                            .project_dependencies(native, direction, &request)
2695                            .await
2696                    }
2697                }
2698            },
2699        )
2700        .await;
2701    }
2702
2703    walk(
2704        start,
2705        budget,
2706        ceiling,
2707        |_| true,
2708        |cursor, limit| async move {
2709            calls.fetch_add(1, Ordering::Relaxed);
2710            let request = PageRequest { cursor, limit };
2711            let (ids, next) = match entity {
2712                Entity::Task => {
2713                    let page = source
2714                        .source()
2715                        .query_tasks(&TaskQuery::default(), &request)
2716                        .await?;
2717                    let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2718                    (ids, page.next)
2719                }
2720                Entity::Project => {
2721                    let page = source
2722                        .source()
2723                        .query_projects(&ProjectQuery::default(), &request)
2724                        .await?;
2725                    let ids: Vec<NativeId> =
2726                        page.items.into_iter().map(|project| project.id).collect();
2727                    (ids, page.next)
2728                }
2729            };
2730
2731            let mut edges = Vec::new();
2732            for id in ids {
2733                let mut inner: Option<Cursor> = None;
2734                loop {
2735                    calls.fetch_add(1, Ordering::Relaxed);
2736                    let request = PageRequest {
2737                        cursor: inner.clone(),
2738                        limit,
2739                    };
2740                    let page = forward_edges(source, entity, &id, &request).await?;
2741                    // The inner half of the same bound the walk holds on its own pages:
2742                    // this scan keeps every matching edge of one source page, so a source
2743                    // that overruns here overruns the engine's memory just as surely.
2744                    fits(page.items.len(), limit)?;
2745                    edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2746                    unrepeated(
2747                        page.next.as_ref(),
2748                        inner.as_ref(),
2749                        "its forward edges were being scanned",
2750                    )?;
2751                    match page.next {
2752                        Some(cursor) => inner = Some(cursor),
2753                        None => break,
2754                    }
2755                }
2756            }
2757
2758            Ok(Page { items: edges, next })
2759        },
2760    )
2761    .await
2762}