1mod assets;
19mod comment;
20mod copy;
21mod delivery;
22mod fetch;
23mod join;
24mod local;
25mod metadata;
26mod narrow;
27mod rendered;
28mod resume;
29mod update;
30
31use std::collections::BTreeMap;
32use std::num::NonZeroU32;
33use std::sync::atomic::{AtomicU32, Ordering};
34
35use chrono::{DateTime, Utc};
36use onetaskgraph_plugin_api::{
37 Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
38 MetadataMatch, MetadataRecord, NativeId, Page, PageRequest, Priority, Project, ProjectFilter,
39 ProjectQuery, Repository, SecretResolver, SharedClock, SourceError, SourceName, StatusCategory,
40 Task, TaskQuery, TextFields, TextQuery, system_clock,
41};
42use schemars::JsonSchema;
43use serde::{Deserialize, Serialize};
44
45use crate::GlobalId;
46use crate::config::{Config, Placement, Routes};
47use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
48use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available_with_clock};
49
50use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
51use join::join_all;
52use local::{LocalDocuments, LocalProjects, LocalTasks};
53pub(crate) use resume::{Owed, Resumption, StreamState};
54use resume::{Resume, StreamKind};
55
56pub use comment::{CommentList, DeletedComment, DocumentDetail, TaskDetail, TaskDetails};
57pub(crate) use copy::malformed_links;
58pub use copy::{
59 BudgetSpent, CopyAction, CopyItems, CopyLink, CopyLookup, CopyOutcome, CopyReport, CopyRequest,
60 CopyScope, CopyVia, MatchBy, NoCounterpart, Spent,
61};
62pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
63pub use local::ProjectSelector;
64pub use metadata::MetadataSet;
65pub use narrow::{TaskContentSet, TaskPrioritySet};
66pub use rendered::{
67 Body, DocumentCreate, ProjectCreate, Regenerated, Regeneration, RenderRequest, RenderTemplate,
68 RenderedRecord, TaskCreate, TaskCreated, TemplateAnswers, UnusedAnswers,
69};
70pub use update::TaskUpdated;
71
72#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
77pub struct Qualified<T> {
78 pub id: GlobalId,
80 pub item: T,
82}
83
84#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
89pub struct QualifiedEdge {
90 pub from: QualifiedEndpoint,
95 pub to: QualifiedEndpoint,
97 pub kind: onetaskgraph_plugin_api::DependencyKind,
99}
100
101#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
103pub struct QualifiedEndpoint {
104 pub id: GlobalId,
106 pub kind: onetaskgraph_plugin_api::ItemKind,
108}
109
110impl std::fmt::Display for QualifiedEndpoint {
111 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
112 self.id.fmt(formatter)
113 }
114}
115
116#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
118#[serde(tag = "kind", rename_all = "kebab-case")]
119pub enum SearchHit {
120 Task(Qualified<Task>),
122 Project(Qualified<Project>),
124}
125
126#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
128#[serde(rename_all = "kebab-case")]
129pub enum SearchKind {
130 Tasks,
132 Projects,
134 #[default]
136 Both,
137}
138
139#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
141pub struct SourceListing {
142 pub source: SourceName,
144 pub kind: String,
156 #[serde(flatten)]
158 pub state: SourceState,
159}
160
161#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
163#[serde(tag = "state", rename_all = "kebab-case")]
164pub enum SourceState {
165 Available {
167 capabilities: Capabilities,
169 },
170 Unavailable {
172 error: SourceError,
174 },
175}
176
177#[derive(Debug, Clone, PartialEq)]
179pub struct Paging {
180 pub limit: NonZeroU32,
182 pub token: Option<PageToken>,
184}
185
186#[derive(Debug, Clone, Default, PartialEq)]
188pub struct Filters {
189 pub text: Option<TextQuery>,
191 pub labels: LabelFilter,
193 pub statuses: Vec<StatusCategory>,
195}
196
197#[derive(Debug, Clone)]
199pub struct TaskRequest {
200 pub sources: Vec<SourceName>,
202 pub filters: Filters,
204 pub project: ProjectSelector,
206 pub priorities: Vec<Priority>,
212 pub commented_since: Option<DateTime<Utc>>,
219 pub metadata: Vec<MetadataMatch>,
225 pub origin: Option<GlobalId>,
228 pub include_members: bool,
236 pub paging: Paging,
238}
239
240#[derive(Debug, Clone)]
242pub struct ProjectRequest {
243 pub sources: Vec<SourceName>,
245 pub filters: Filters,
247 pub paging: Paging,
249}
250
251#[derive(Debug, Clone, Default, PartialEq)]
258pub struct DocumentFilters {
259 pub text: Option<TextQuery>,
261 pub labels: LabelFilter,
263}
264
265#[derive(Debug, Clone)]
267pub struct DocumentRequest {
268 pub sources: Vec<SourceName>,
270 pub filters: DocumentFilters,
272 pub project: ProjectSelector,
274 pub paging: Paging,
276}
277
278#[derive(Debug, Clone)]
280pub struct LabelRequest {
281 pub sources: Vec<SourceName>,
283 pub paging: Paging,
285}
286
287#[derive(Debug, Clone)]
289pub struct SearchRequest {
290 pub sources: Vec<SourceName>,
292 pub text: TextQuery,
294 pub kind: SearchKind,
296 pub paging: Paging,
298}
299
300#[derive(Debug, Clone)]
302pub struct DependencyRequest {
303 pub id: GlobalId,
305 pub direction: Direction,
307 pub paging: Paging,
309}
310
311#[derive(Debug, Clone, PartialEq, thiserror::Error)]
319pub enum EngineError {
320 #[error(
322 "no source named {name:?} is configured\n\
323 next: name one of the configured sources ({configured}), or add {name:?} under \
324 `sources` — `onetaskgraph sources list` shows what this configuration has."
325 )]
326 UnknownSource {
327 name: String,
329 configured: String,
331 },
332
333 #[error(
335 "{message}\n\
336 next: page with a token exactly as the previous page reported it, and against \
337 the same configuration — or drop `--page` to start the walk again."
338 )]
339 Token {
340 message: String,
342 },
343
344 #[error(
346 "no sources are configured\n\
347 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
348 prints what each plugin accepts."
349 )]
350 NoSources,
351
352 #[error(
354 "source {name} cannot be written: its plugin is {kind}, which has no write \
355 side\n\
356 next: copy into a source whose plugin can be written — `onetaskgraph sources \
357 list` reports each one's plugin."
358 )]
359 NotWritable {
360 name: String,
362 kind: String,
364 },
365
366 #[error(
373 "source {name} has no documents: its plugin is {kind}, which holds none\n\
374 next: name a source whose plugin has documents — `onetaskgraph sources list` \
375 reports each one's plugin and what it declares."
376 )]
377 NoDocuments {
378 name: String,
380 kind: String,
382 },
383
384 #[error(
389 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
390 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
391 list` reports each one's plugin."
392 )]
393 NoComments {
394 name: String,
396 kind: String,
398 },
399
400 #[error(
402 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
403 but not added to, edited or removed\n\
404 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
405 source whose plugin can be written — `onetaskgraph sources list` reports each \
406 one's plugin."
407 )]
408 CommentsNotWritable {
409 name: String,
411 kind: String,
413 },
414
415 #[error(
417 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
418 next: set the status in that source itself, or name a task of a source whose plugin \
419 can be written — `onetaskgraph sources list` reports each one's plugin."
420 )]
421 StatusNotWritable {
422 name: String,
424 kind: String,
426 },
427
428 #[error(
430 "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
431 write side\n\
432 next: set the key in that source itself, or name a {record} of a source whose plugin \
433 can be written — `onetaskgraph sources list` reports each one's plugin."
434 )]
435 MetadataNotWritable {
436 name: String,
438 kind: String,
440 record: MetadataRecord,
442 },
443
444 #[error(
446 "source {name} cannot write a priority: its plugin is {kind}, which has no write side\n\
447 next: set the priority in that source itself, or name a task of a source whose plugin \
448 can be written — `onetaskgraph sources list` reports each one's plugin."
449 )]
450 PriorityNotWritable {
451 name: String,
453 kind: String,
455 },
456
457 #[error(
459 "source {name} cannot write a task's content: its plugin is {kind}, which has no write \
460 side\n\
461 next: edit the content in that source itself, or name a task of a source whose plugin \
462 can be written — `onetaskgraph sources list` reports each one's plugin."
463 )]
464 ContentNotWritable {
465 name: String,
467 kind: String,
469 },
470
471 #[error(
473 "source {name} cannot update a task: its plugin is {kind}, which has no write side\n\
474 next: change the task in that source itself, or name a task of a source whose plugin \
475 can be written — `onetaskgraph sources list` reports each one's plugin."
476 )]
477 UpdateNotWritable {
478 name: String,
480 kind: String,
482 },
483
484 #[error(
486 "the {record} was given the asset {asset}, which its content does not reference\n\
487 next: reference it in the content as ``, or drop that --asset."
488 )]
489 AssetNotReferenced {
490 record: String,
492 asset: String,
494 },
495
496 #[error(
498 "the {record}'s content references the asset ./{asset}, which the {record} would not \
499 hold\n\
500 next: give it with --asset <PATH> naming a file called {asset}, or remove the \
501 reference."
502 )]
503 AssetNotGiven {
504 record: String,
506 asset: String,
508 },
509
510 #[error(
512 "the {record} was given two assets named {asset}\n\
513 next: give each asset once — an asset is stored under its file's base name, so \
514 rename one of the two files."
515 )]
516 AssetGivenTwice {
517 record: String,
519 asset: String,
521 },
522
523 #[error(
526 "source {name} cannot store the asset {asset} of {record}: its plugin is {kind}, which \
527 declares image assets unsupported\n\
528 next: write the record to a source whose plugin stores assets — such as local-md — \
529 or remove the reference to ./{asset} from its content."
530 )]
531 AssetsUnsupported {
532 name: String,
534 kind: String,
536 record: String,
538 asset: String,
540 },
541
542 #[error(
544 "{record}'s content references the asset ./{asset}, which {record} does not hold\n\
545 next: store the asset with the record — `onetaskgraph task render` or `document \
546 render` with --asset <PATH> — or remove the reference, then copy again."
547 )]
548 AssetNotHeld {
549 record: String,
551 asset: String,
553 },
554
555 #[error(
557 "source {name} cannot create a {record}: its plugin is {kind}, which has no write side\n\
558 next: create it in a source whose plugin can be written — `onetaskgraph sources list` \
559 reports each one's plugin."
560 )]
561 NotCreatable {
562 name: String,
564 kind: String,
566 record: MetadataRecord,
568 },
569
570 #[error(
573 "source {name} cannot write a {record}'s rendering: its plugin is {kind}, which has no \
574 write side\n\
575 next: render it with --dry-run to read the result, or regenerate a {record} of a source \
576 whose plugin can be written."
577 )]
578 RenderingNotWritable {
579 name: String,
581 kind: String,
583 record: RenderedRecord,
585 },
586
587 #[error(
590 "supply every required answer to regenerate {id}: {} unanswered, and {reason}\n\
591 next: answer {} with --var NAME=VALUE or an answers file (--answers FILE), or run \
592 interactively to be asked.",
593 names.join(", "),
594 if names.len() == 1 { "it" } else { "each" }
595 )]
596 MissingAnswers {
597 id: String,
599 names: Vec<String>,
601 reason: String,
604 },
605
606 #[error("{error}")]
608 Template {
609 error: crate::template::TemplateError,
611 },
612
613 #[error(
615 "{record} {id} has no stored template answers: {reason}\n\
616 next: regenerate it with every required answer (`onetaskgraph {record} render {id} \
617 --var NAME=VALUE`), or read its provenance with `onetaskgraph {record} show {id}`."
618 )]
619 NoStoredAnswers {
620 record: RenderedRecord,
622 id: String,
624 reason: String,
626 },
627
628 #[error(
630 "{record} {id} records no template it was rendered from, and none was given\n\
631 next: name one with --template FILE or --template-loader FILE."
632 )]
633 NoTemplate {
634 record: RenderedRecord,
636 id: String,
638 },
639
640 #[error(
643 "{record} {id} records a template entry this product did not write — {problem}\n\
644 next: regenerate it with --template FILE or --template-loader FILE and every required \
645 answer, which records a fresh entry."
646 )]
647 MalformedProvenance {
648 record: RenderedRecord,
650 id: String,
652 problem: String,
654 },
655
656 #[error(
658 "{record} {id} was rendered from {reference:?}, which is not a readable file, so it \
659 cannot be re-read; a recorded reference is never turned into a location\n\
660 next: supply the template with --template-loader FILE (a loader document naming what \
661 to render), or name a template file with --template FILE."
662 )]
663 TemplateNotAFile {
664 record: RenderedRecord,
666 id: String,
668 reference: String,
670 },
671
672 #[error(
677 "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
678 be written to it: its plugin is {kind}, which declares priority unsupported\n\
679 next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
680 reports what each declares — or set the task's priority to none first; a \
681 github-projects source holds one once its configuration sets priority_mapping."
682 )]
683 NoPriority {
684 name: String,
686 kind: String,
688 task: String,
690 priority: onetaskgraph_plugin_api::Priority,
692 },
693
694 #[error(
696 "no project with the id {id}\n\
697 next: check the id, or list what is there — `onetaskgraph project list` reports every \
698 project the configured sources hold."
699 )]
700 NoSuchProject {
701 id: String,
703 },
704
705 #[error(
707 "no document with the id {id}\n\
708 next: check the id, or list what is there — `onetaskgraph document list` reports every \
709 document the configured sources hold."
710 )]
711 NoSuchDocument {
712 id: String,
714 },
715
716 #[error(
718 "no task with the id {id}\n\
719 next: check the id, or list what is there — `onetaskgraph task list` reports every \
720 task the configured sources hold."
721 )]
722 NoSuchTask {
723 id: String,
725 },
726
727 #[error(
729 "task {task} has no comment with the id {comment}\n\
730 next: list its comments — `onetaskgraph task comment list {task}` reports each \
731 one's id."
732 )]
733 NoSuchComment {
734 task: String,
736 comment: String,
738 },
739
740 #[error(
745 "source {name} could not be built: {error}\n\
746 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
747 command again."
748 )]
749 SourceUnavailable {
750 name: String,
752 error: SourceError,
754 },
755
756 #[error(
762 "source {name} could not do it: {error}\n\
763 next: fix what the source named above, then run the command again."
764 )]
765 SourceFailed {
766 name: String,
768 error: SourceError,
770 },
771
772 #[error(
774 "the destination source {name} could not be built: {error}\n\
775 next: fix that source — `onetaskgraph sources list` reports its state — then \
776 copy again."
777 )]
778 DestinationUnavailable {
779 name: String,
781 error: SourceError,
783 },
784
785 #[error(
787 "no item with the id {id}\n\
788 next: check the id, or list what is there — `onetaskgraph task list` and \
789 `onetaskgraph project list` report what the configured sources hold."
790 )]
791 NoSuchItem {
792 id: String,
794 },
795
796 #[error(
801 "{item} was copied from {origin}, which that destination no longer holds\n\
802 next: re-run with --recreate to create a new item there instead, or restore \
803 {origin}."
804 )]
805 StaleOrigin {
806 item: String,
808 origin: String,
810 },
811
812 #[error(
818 "{item} was last copied to {link}, which that destination no longer holds\n\
819 next: re-run with --recreate to create a new item there instead, or restore \
820 {link}."
821 )]
822 StaleLink {
823 item: String,
825 link: String,
827 },
828
829 #[error(
832 "--create cannot be given with {flag}: --create asserts the destination holds no \
833 counterpart, so there is nothing for {flag} to look for\n\
834 next: drop {flag} to create each item without looking, or drop --create to look."
835 )]
836 CreateWith {
837 flag: CopyLookup,
839 },
840
841 #[error(
847 "{item} already records a counterpart at the destination, {carrier}, and --create \
848 asserts it has none\n\
849 next: copy it without --create, which updates {carrier}."
850 )]
851 CreateCarried {
852 item: GlobalId,
854 carrier: GlobalId,
856 },
857
858 #[error(
864 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
865 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
866 {id} on its own with `onetaskgraph task copy`.",
867 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
868 )]
869 NotAMember {
870 id: GlobalId,
872 projects: Vec<GlobalId>,
874 },
875
876 #[error(
885 "{item} depends on {member}, which this copy was not told to carry and which records \
886 no origin in {destination}\n\
887 next: name {member} with --member as well, record its {destination} id at \
888 onetaskgraph.origin, or copy the whole project without --member."
889 )]
890 UnrecordedMember {
891 item: GlobalId,
893 member: GlobalId,
895 destination: SourceName,
897 },
898
899 #[error(
906 "{item} was copied to {counterpart}, but its repositories now route it to {route}\n\
907 next: move or remove {counterpart} by hand and copy again, or restore {item}'s \
908 repositories so it routes to {source} again.",
909 source = .counterpart.source
910 )]
911 Misrouted {
912 item: GlobalId,
914 counterpart: GlobalId,
916 route: SourceName,
918 },
919
920 #[error(
926 "source {name} could not do it: {error}\n\
927 next: fix what the source named above, then copy again."
928 )]
929 SourceRefused {
930 name: String,
932 error: SourceError,
934 },
935
936 #[error(
944 "the copy failed and could not be undone.\n\
945 it failed because: {error}\n\
946 it could not be undone because: {refusal}\n\
947 so these still hold what it wrote: {left_behind}\n\
948 next: remove or put back those items, then copy again."
949 )]
950 CopyNotUndone {
951 error: Box<EngineError>,
953 left_behind: LeftBehind,
959 refusal: SourceError,
961 },
962}
963
964#[derive(Debug, Clone, PartialEq)]
973pub struct LeftBehind {
974 first: GlobalId,
976 rest: Vec<GlobalId>,
978}
979
980impl LeftBehind {
981 #[must_use]
983 pub fn new(first: GlobalId) -> Self {
984 Self {
985 first,
986 rest: Vec::new(),
987 }
988 }
989
990 pub fn push(&mut self, id: GlobalId) {
992 self.rest.push(id);
993 }
994
995 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
997 std::iter::once(&self.first).chain(self.rest.iter())
998 }
999}
1000
1001impl std::fmt::Display for LeftBehind {
1002 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1004 write!(formatter, "{}", self.first)?;
1005 for id in &self.rest {
1006 write!(formatter, ", {id}")?;
1007 }
1008 Ok(())
1009 }
1010}
1011
1012pub enum ConfiguredSource {
1019 Ready(ResolvedSource),
1021 Unavailable(UnavailableSource),
1023}
1024
1025impl ConfiguredSource {
1026 #[must_use]
1028 pub fn name(&self) -> &SourceName {
1029 match self {
1030 Self::Ready(source) => source.name(),
1031 Self::Unavailable(source) => source.name(),
1032 }
1033 }
1034}
1035
1036pub struct Engine {
1038 sources: Vec<ConfiguredSource>,
1040 selection: Vec<SourceName>,
1042 routes: Routes,
1045}
1046
1047impl Engine {
1048 #[must_use]
1056 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
1057 Self::build_with_clock(config, secrets, &system_clock())
1058 }
1059
1060 #[must_use]
1064 pub fn build_with_clock(
1065 config: &Config,
1066 secrets: &dyn SecretResolver,
1067 clock: &SharedClock,
1068 ) -> Self {
1069 let (ready, unavailable) = resolve_available_with_clock(config, secrets, clock);
1070 Self::new(
1071 ready
1072 .into_iter()
1073 .map(ConfiguredSource::Ready)
1074 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
1075 .collect(),
1076 config.selected_sources(),
1077 )
1078 .with_routes(config.routes())
1079 }
1080
1081 #[must_use]
1084 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
1085 Self {
1086 sources,
1087 selection,
1088 routes: Routes::default(),
1089 }
1090 }
1091
1092 #[must_use]
1097 pub fn with_routes(mut self, routes: Routes) -> Self {
1098 self.routes = routes;
1099 self
1100 }
1101
1102 #[must_use]
1104 pub fn place(&self, source: &SourceName, repositories: &[Repository]) -> Placement {
1105 self.routes.place(source, repositories)
1106 }
1107
1108 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
1110 self.sources.iter().filter_map(|source| match source {
1111 ConfiguredSource::Ready(ready) => Some(ready),
1112 ConfiguredSource::Unavailable(_) => None,
1113 })
1114 }
1115
1116 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
1118 self.sources.iter().filter_map(|source| match source {
1119 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
1120 ConfiguredSource::Ready(_) => None,
1121 })
1122 }
1123
1124 #[must_use]
1126 pub fn listing(&self) -> Vec<SourceListing> {
1127 let mut listings: Vec<SourceListing> = self
1128 .ready()
1129 .map(|source| SourceListing {
1130 source: source.name().clone(),
1131 kind: source.kind().to_owned(),
1132 state: SourceState::Available {
1133 capabilities: source.source().capabilities(),
1134 },
1135 })
1136 .chain(self.unavailable().map(|source| SourceListing {
1137 source: source.name().clone(),
1138 kind: source.kind().to_owned(),
1139 state: SourceState::Unavailable {
1140 error: source.error().clone(),
1141 },
1142 }))
1143 .collect();
1144 listings.sort_by(|left, right| left.source.cmp(&right.source));
1145 listings
1146 }
1147
1148 #[must_use]
1154 pub fn has(&self, name: &SourceName) -> bool {
1155 self.sources.iter().any(|source| source.name() == name)
1156 }
1157
1158 pub async fn end_command(&self) -> Result<(), EngineError> {
1178 let mut first = None;
1179 for source in self.ready() {
1180 if let Err(error) = source.source().end_command().await
1181 && first.is_none()
1182 {
1183 first = Some(EngineError::SourceFailed {
1184 name: source.name().to_string(),
1185 error,
1186 });
1187 }
1188 }
1189 first.map_or(Ok(()), Err)
1190 }
1191
1192 pub async fn tasks(
1200 &self,
1201 request: &TaskRequest,
1202 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1203 let mut names = self.resolve_selection(&request.sources)?;
1204 let mut members: Vec<GlobalId> = Vec::new();
1209 let mut unread: Vec<SourceFailure> = Vec::new();
1210 if let ProjectSelector::Qualified(id) = &request.project {
1211 self.known(&id.source)?;
1212 names.retain(|name| name == &id.source);
1213 if request.include_members {
1216 (members, unread) = self.members(id).await;
1217 names.extend(members.iter().map(|member| member.source.clone()));
1218 }
1219 }
1220 let selector = |name: &SourceName| -> ProjectSelector {
1221 members
1222 .iter()
1223 .find(|member| &member.source == name)
1224 .map_or_else(
1225 || request.project.clone(),
1226 |member| ProjectSelector::Qualified(member.clone()),
1227 )
1228 };
1229 let query = shape(
1230 "task-list",
1231 &names,
1232 &(
1233 &request.filters,
1234 &request.project,
1235 &request.priorities,
1236 &request.commented_since,
1237 &request.metadata,
1238 &request.origin,
1239 &members,
1240 ),
1241 );
1242 let states = resumption(
1243 self,
1244 request.paging.token.as_ref(),
1245 &[StreamKind::Items],
1246 &query,
1247 )?;
1248 let budget = request.paging.limit.get();
1249
1250 let mut answer = Answer::new();
1251 answer.errors.extend(unread);
1252 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1253
1254 let shapes: Vec<TaskShape> = ready
1255 .iter()
1256 .map(|source| {
1257 shape_tasks(
1258 &source.source().capabilities(),
1259 &request.filters,
1260 &project_filter(&selector(source.name())),
1261 &request.priorities,
1262 request.commented_since,
1263 &request.metadata,
1264 request.origin.as_ref(),
1265 )
1266 })
1267 .collect();
1268 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1269 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1270
1271 let walks = ready
1272 .iter()
1273 .enumerate()
1274 .map(|(index, source)| {
1275 fetch_tasks(
1276 source,
1277 &shapes[index],
1278 &starts[index],
1279 budget,
1280 &counters[index],
1281 )
1282 })
1283 .collect();
1284
1285 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1286 answer.finish(
1287 streams,
1288 budget,
1289 owed(&states),
1290 &query,
1291 |name, task: Task| {
1292 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1293 },
1294 )
1295 }
1296
1297 async fn members(&self, home: &GlobalId) -> (Vec<GlobalId>, Vec<SourceFailure>) {
1308 let mut members = Vec::new();
1309 let mut unread = Vec::new();
1310 let Some(source) = self.ready().find(|source| source.name() == &home.source) else {
1311 return (members, unread);
1312 };
1313 let held = match source.source().get_project(&home.native).await {
1314 Ok(Some(held)) => held,
1315 Ok(None) => {
1316 unread.push(SourceFailure {
1317 source: home.source.clone(),
1318 error: SourceError::Refused {
1319 message: format!(
1320 "{home} names no project {} holds, so its member projects cannot \
1321 be read",
1322 home.source
1323 ),
1324 },
1325 });
1326 return (members, unread);
1327 }
1328 Err(error) => {
1329 unread.push(SourceFailure {
1330 source: home.source.clone(),
1331 error,
1332 });
1333 return (members, unread);
1334 }
1335 };
1336 let named = match copy::members_of(home, &held.metadata) {
1337 Ok(named) => named,
1338 Err(message) => {
1341 unread.push(SourceFailure {
1342 source: home.source.clone(),
1343 error: SourceError::Malformed { message },
1344 });
1345 return (members, unread);
1346 }
1347 };
1348 for member in named {
1349 if !self.has(&member.source) {
1350 unread.push(SourceFailure {
1351 source: member.source.clone(),
1352 error: SourceError::Config {
1353 message: format!(
1354 "{home} names {member} as a member project, but no source named \
1355 {} is configured",
1356 member.source
1357 ),
1358 },
1359 });
1360 continue;
1361 }
1362 let Some(there) = self.ready().find(|source| source.name() == &member.source) else {
1365 members.push(member);
1366 continue;
1367 };
1368 match there.source().get_project(&member.native).await {
1369 Ok(Some(_)) => members.push(member),
1370 Ok(None) => unread.push(SourceFailure {
1371 source: member.source.clone(),
1372 error: SourceError::Refused {
1373 message: format!(
1374 "{home} names {member} as a member project, and {} holds no such \
1375 project",
1376 member.source
1377 ),
1378 },
1379 }),
1380 Err(error) => unread.push(SourceFailure {
1381 source: member.source.clone(),
1382 error,
1383 }),
1384 }
1385 }
1386 (members, unread)
1387 }
1388
1389 pub async fn projects(
1395 &self,
1396 request: &ProjectRequest,
1397 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1398 let names = self.resolve_selection(&request.sources)?;
1399 let query = shape("project-list", &names, &request.filters);
1400 let states = resumption(
1401 self,
1402 request.paging.token.as_ref(),
1403 &[StreamKind::Items],
1404 &query,
1405 )?;
1406 let budget = request.paging.limit.get();
1407
1408 let mut answer = Answer::new();
1409 let mut with_projects = Vec::new();
1415 for source in answer.split(self, &names) {
1416 if source.source().capabilities().projects.is_native() {
1417 with_projects.push(source);
1418 } else {
1419 answer.unreachable_predicate(source, Predicate::Project);
1420 }
1421 }
1422 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1423
1424 let shapes: Vec<ProjectShape> = ready
1425 .iter()
1426 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1427 .collect();
1428 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1429 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1430
1431 let walks = ready
1432 .iter()
1433 .enumerate()
1434 .map(|(index, source)| {
1435 fetch_projects(
1436 source,
1437 &shapes[index],
1438 &starts[index],
1439 budget,
1440 &counters[index],
1441 )
1442 })
1443 .collect();
1444
1445 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1446 answer.finish(
1447 streams,
1448 budget,
1449 owed(&states),
1450 &query,
1451 |name, project: Project| Qualified {
1452 id: GlobalId::new(name.clone(), project.id.clone()),
1453 item: project,
1454 },
1455 )
1456 }
1457
1458 pub async fn documents(
1470 &self,
1471 request: &DocumentRequest,
1472 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1473 let mut names = self.resolve_selection(&request.sources)?;
1474 if let ProjectSelector::Qualified(id) = &request.project {
1477 self.known(&id.source)?;
1478 names.retain(|name| name == &id.source);
1479 }
1480 let query = shape(
1481 "document-list",
1482 &names,
1483 &(&request.filters, &request.project),
1484 );
1485 let states = resumption(
1486 self,
1487 request.paging.token.as_ref(),
1488 &[StreamKind::Items],
1489 &query,
1490 )?;
1491 let budget = request.paging.limit.get();
1492
1493 let mut answer = Answer::new();
1494 let mut with_documents = Vec::new();
1495 for source in answer.split(self, &names) {
1496 if source.source().capabilities().documents.is_native() {
1497 with_documents.push(source);
1498 } else {
1499 answer.unreachable_predicate(source, Predicate::Document);
1500 }
1501 }
1502 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1503
1504 let shapes: Vec<DocumentShape> = ready
1505 .iter()
1506 .map(|source| {
1507 shape_documents(
1508 &source.source().capabilities(),
1509 &request.filters,
1510 &project_filter(&request.project),
1511 )
1512 })
1513 .collect();
1514 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1515 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1516
1517 let walks = ready
1518 .iter()
1519 .enumerate()
1520 .map(|(index, source)| {
1521 fetch_documents(
1522 source,
1523 &shapes[index],
1524 &starts[index],
1525 budget,
1526 &counters[index],
1527 )
1528 })
1529 .collect();
1530
1531 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1532 answer.finish(
1533 streams,
1534 budget,
1535 owed(&states),
1536 &query,
1537 |name, document: Document| Qualified {
1538 id: GlobalId::new(name.clone(), document.id.clone()),
1539 item: document,
1540 },
1541 )
1542 }
1543
1544 pub async fn labels(
1550 &self,
1551 request: &LabelRequest,
1552 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1553 let names = self.resolve_selection(&request.sources)?;
1554 let query = shape("label-list", &names, &());
1555 let states = resumption(
1556 self,
1557 request.paging.token.as_ref(),
1558 &[StreamKind::Items],
1559 &query,
1560 )?;
1561 let budget = request.paging.limit.get();
1562
1563 let mut answer = Answer::new();
1564 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1565 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1566 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1567
1568 let walks = ready
1569 .iter()
1570 .enumerate()
1571 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1572 .collect();
1573
1574 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1575 answer.finish(
1576 streams,
1577 budget,
1578 owed(&states),
1579 &query,
1580 |name, label: Label| Qualified {
1581 id: GlobalId::new(name.clone(), label.id.clone()),
1582 item: label,
1583 },
1584 )
1585 }
1586
1587 pub async fn search(
1593 &self,
1594 request: &SearchRequest,
1595 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1596 let names = self.resolve_selection(&request.sources)?;
1597 let reads: &[StreamKind] = match request.kind {
1601 SearchKind::Tasks => &[StreamKind::Tasks],
1602 SearchKind::Projects => &[StreamKind::Projects],
1603 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1604 };
1605 let query = shape("search", &names, &request.text);
1610 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1611 let budget = request.paging.limit.get();
1612 let filters = Filters {
1613 text: Some(request.text.clone()),
1614 ..Filters::default()
1615 };
1616
1617 let mut answer = Answer::new();
1618
1619 let mut ready = Vec::new();
1622 let mut kinds = Vec::new();
1623 let mut starts = Vec::new();
1624 for source in answer.split(self, &names) {
1625 let mut streams = Vec::new();
1626 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1627 streams.push(StreamKind::Tasks);
1628 }
1629 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1630 if source.source().capabilities().projects.is_native() {
1631 streams.push(StreamKind::Projects);
1632 } else {
1633 answer.unreachable_predicate(source, Predicate::Project);
1634 }
1635 }
1636 for stream in streams {
1637 if let Some(resume) = resume_at(&states, source.name(), stream) {
1638 ready.push(source);
1639 kinds.push(stream);
1640 starts.push(resume);
1641 }
1642 }
1643 }
1644
1645 let shapes: Vec<HitShape> = ready
1646 .iter()
1647 .zip(kinds.iter())
1648 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1649 .collect();
1650 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1651 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1652
1653 let walks = ready
1654 .iter()
1655 .enumerate()
1656 .map(|(index, source)| {
1657 fetch_hits(
1658 source,
1659 &shapes[index],
1660 &starts[index],
1661 budget,
1662 &counters[index],
1663 )
1664 })
1665 .collect();
1666
1667 let streams =
1668 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1669 answer.finish(
1670 streams,
1671 budget,
1672 owed(&states),
1673 &query,
1674 |name, found: Found| match found {
1675 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1676 GlobalId::new(name.clone(), task.id.clone()),
1677 task,
1678 )),
1679 Found::Project(project) => SearchHit::Project(Qualified {
1680 id: GlobalId::new(name.clone(), project.id.clone()),
1681 item: project,
1682 }),
1683 },
1684 )
1685 }
1686
1687 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1694 let name = self.known(&id.source)?;
1695 let mut answer = Answer::new();
1696 let selected = answer.split(self, std::slice::from_ref(&name));
1697 let Some(source) = selected.first() else {
1698 return answer.nothing();
1699 };
1700 let found = source.source().get_task(&id.native).await;
1701 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1702 answer.one(source, found, |task| {
1703 delivery::qualified_task(qualified, task)
1704 })
1705 }
1706
1707 pub async fn project(
1713 &self,
1714 id: &GlobalId,
1715 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1716 let name = self.known(&id.source)?;
1717 let mut answer = Answer::new();
1718 let selected = answer.split(self, std::slice::from_ref(&name));
1719 let Some(source) = selected.first() else {
1720 return answer.nothing();
1721 };
1722 let found = source.source().get_project(&id.native).await;
1723 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1724 answer.one(source, found, |project| Qualified {
1725 id: qualified,
1726 item: project,
1727 })
1728 }
1729
1730 pub async fn document(
1740 &self,
1741 id: &GlobalId,
1742 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1743 let name = self.known(&id.source)?;
1744 let mut answer = Answer::new();
1745 let selected = answer.split(self, std::slice::from_ref(&name));
1746 let Some(source) = selected.first() else {
1747 return answer.nothing();
1748 };
1749 if !source.source().capabilities().documents.is_native() {
1750 answer.unreachable_predicate(source, Predicate::Document);
1751 return answer.nothing();
1752 }
1753 let found = source.source().get_document(&id.native).await;
1754 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1755 answer.one(source, found, |document| Qualified {
1756 id: qualified,
1757 item: document,
1758 })
1759 }
1760
1761 pub async fn task_dependencies(
1768 &self,
1769 request: &DependencyRequest,
1770 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1771 self.dependencies(request, Entity::Task).await
1772 }
1773
1774 pub async fn project_dependencies(
1780 &self,
1781 request: &DependencyRequest,
1782 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1783 self.dependencies(request, Entity::Project).await
1784 }
1785
1786 async fn dependencies(
1789 &self,
1790 request: &DependencyRequest,
1791 entity: Entity,
1792 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1793 let name = self.known(&request.id.source)?;
1794 let query = shape(
1795 "dependencies",
1796 std::slice::from_ref(&name),
1797 &(entity, &request.id.native, request.direction),
1798 );
1799 let states = resumption(
1800 self,
1801 request.paging.token.as_ref(),
1802 &[StreamKind::Items],
1803 &query,
1804 )?;
1805 let budget = request.paging.limit.get();
1806
1807 let mut answer = Answer::new();
1808 let (ready, starts) = walking(
1809 answer.split(self, std::slice::from_ref(&name)),
1810 &states,
1811 StreamKind::Items,
1812 );
1813 let Some(source) = ready.first() else {
1814 return answer.nothing();
1815 };
1816
1817 let capabilities = source.source().capabilities();
1818 let support = match entity {
1819 Entity::Task => capabilities.task_dependencies,
1820 Entity::Project => capabilities.project_dependencies,
1821 };
1822 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1826 let mut outcomes = Outcomes::default();
1827 if request.direction == Direction::DependedOnBy {
1828 if emulating {
1829 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1830 } else {
1831 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1832 }
1833 }
1834
1835 let counters = vec![AtomicU32::new(0)];
1836 let walked = fetch_edges(
1837 source,
1838 &request.id.native,
1839 request.direction,
1840 entity,
1841 emulating,
1842 &starts[0],
1843 budget,
1844 &counters[0],
1845 )
1846 .await;
1847
1848 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1849 answer.finish(
1850 streams,
1851 budget,
1852 owed(&states),
1853 &query,
1854 |name, edge: DependencyEdge| QualifiedEdge {
1855 from: qualify_endpoint(name, edge.from),
1856 to: qualify_endpoint(name, edge.to),
1857 kind: edge.kind,
1858 },
1859 )
1860 }
1861
1862 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1864 if asked.is_empty() {
1865 if self.selection.is_empty() {
1866 return Err(EngineError::NoSources);
1867 }
1868 return Ok(self.selection.clone());
1869 }
1870 asked.iter().map(|name| self.known(name)).collect()
1871 }
1872
1873 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1875 if self.has(name) {
1876 return Ok(name.clone());
1877 }
1878 if self.sources.is_empty() {
1879 return Err(EngineError::NoSources);
1880 }
1881 Err(EngineError::UnknownSource {
1882 name: name.to_string(),
1883 configured: self
1884 .listing()
1885 .iter()
1886 .map(|listing| listing.source.to_string())
1887 .collect::<Vec<_>>()
1888 .join(", "),
1889 })
1890 }
1891}
1892
1893fn qualify_endpoint(
1894 source: &SourceName,
1895 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1896) -> QualifiedEndpoint {
1897 let kind = endpoint.kind;
1898 let is_qualified = endpoint.is_qualified();
1899 let endpoint_id = endpoint.into_id();
1900 QualifiedEndpoint {
1901 id: if is_qualified {
1902 endpoint_id
1903 .parse()
1904 .expect("plugin-api validates qualified dependency endpoints")
1905 } else {
1906 GlobalId::new(source.clone(), NativeId(endpoint_id))
1907 },
1908 kind,
1909 }
1910}
1911
1912#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1914enum Entity {
1915 Task,
1917 Project,
1919}
1920
1921enum Found {
1923 Task(Task),
1925 Project(Project),
1927}
1928
1929#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1931enum Outcome {
1932 PushedDown,
1934 AppliedLocally,
1936 Emulated,
1938 Unavailable,
1940}
1941
1942#[derive(Debug, Clone, Default, PartialEq)]
1954struct Outcomes(BTreeMap<Predicate, Outcome>);
1955
1956impl Outcomes {
1957 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1963 self.0.insert(predicate, outcome);
1964 }
1965
1966 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1969 for predicate in predicates {
1970 self.record(predicate, outcome);
1971 }
1972 }
1973
1974 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1976 self.0
1977 .iter()
1978 .filter(|(_, recorded)| **recorded == outcome)
1979 .map(|(predicate, _)| *predicate)
1980 .collect()
1981 }
1982}
1983
1984struct TaskShape {
1986 pushed: TaskQuery,
1988 local: LocalTasks,
1990 outcomes: Outcomes,
1992}
1993
1994struct ProjectShape {
1996 pushed: ProjectQuery,
1998 local: LocalProjects,
2000 outcomes: Outcomes,
2002}
2003
2004struct DocumentShape {
2006 pushed: DocumentQuery,
2008 local: LocalDocuments,
2010 outcomes: Outcomes,
2012}
2013
2014struct HitShape {
2016 stream: StreamKind,
2018 tasks: TaskQuery,
2020 projects: ProjectQuery,
2022 local_tasks: LocalTasks,
2024 local_projects: LocalProjects,
2026 outcomes: Outcomes,
2028}
2029
2030struct Answer {
2036 plans: Vec<SourcePlan>,
2038 errors: Vec<SourceFailure>,
2040}
2041
2042impl Answer {
2043 fn new() -> Self {
2044 Self {
2045 plans: Vec::new(),
2046 errors: Vec::new(),
2047 }
2048 }
2049
2050 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
2055 let mut selected = Vec::new();
2056 for name in names {
2057 match engine.sources.iter().find(|source| source.name() == name) {
2058 Some(ConfiguredSource::Ready(source)) => selected.push(source),
2059 Some(ConfiguredSource::Unavailable(source)) => {
2060 self.errors.push(source.failure());
2061 }
2062 None => {}
2063 }
2064 }
2065 selected
2066 }
2067
2068 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
2070 let mut outcomes = Outcomes::default();
2071 outcomes.record(predicate, Outcome::Unavailable);
2072 self.plans.push(plan_for(source, outcomes, 0));
2073 }
2074
2075 fn collect<T>(
2077 &mut self,
2078 ready: &[&ResolvedSource],
2079 walked: Vec<Result<Fetched<T>, SourceError>>,
2080 counters: &[AtomicU32],
2081 outcomes: Vec<Outcomes>,
2082 ) -> Vec<Stream<T>> {
2083 let kinds = vec![StreamKind::Items; ready.len()];
2084 self.collect_streams(ready, &kinds, walked, counters, outcomes)
2085 }
2086
2087 fn collect_streams<T>(
2089 &mut self,
2090 ready: &[&ResolvedSource],
2091 kinds: &[StreamKind],
2092 walked: Vec<Result<Fetched<T>, SourceError>>,
2093 counters: &[AtomicU32],
2094 outcomes: Vec<Outcomes>,
2095 ) -> Vec<Stream<T>> {
2096 let mut streams = Vec::new();
2097 for (index, result) in walked.into_iter().enumerate() {
2098 let source = ready[index];
2099 let pages = counters[index].load(Ordering::Relaxed);
2100 self.plans
2101 .push(plan_for(source, outcomes[index].clone(), pages));
2102 match result {
2103 Ok(fetched) => streams.push(Stream {
2104 source: source.name().clone(),
2105 kind: kinds[index],
2106 fetched,
2107 }),
2108 Err(error) => self.errors.push(SourceFailure {
2111 source: source.name().clone(),
2112 error,
2113 }),
2114 }
2115 }
2116 streams
2117 }
2118
2119 fn one<T, U>(
2121 self,
2122 source: &ResolvedSource,
2123 found: Result<Option<T>, SourceError>,
2124 qualify: impl FnOnce(T) -> U,
2125 ) -> Result<QueryResponse<U>, EngineError> {
2126 Ok(self.one_response(source, found, qualify))
2127 }
2128
2129 fn one_response<T, U>(
2131 mut self,
2132 source: &ResolvedSource,
2133 found: Result<Option<T>, SourceError>,
2134 qualify: impl FnOnce(T) -> U,
2135 ) -> QueryResponse<U> {
2136 self.plans.push(plan_for(source, Outcomes::default(), 1));
2137 let items = match found {
2138 Ok(Some(item)) => vec![qualify(item)],
2139 Ok(None) => Vec::new(),
2140 Err(error) => {
2141 self.errors.push(SourceFailure {
2142 source: source.name().clone(),
2143 error,
2144 });
2145 Vec::new()
2146 }
2147 };
2148 QueryResponse {
2149 items,
2150 next: None,
2151 plan: QueryPlan {
2152 per_source: merge_plans(self.plans),
2153 },
2154 errors: self.errors,
2155 }
2156 }
2157
2158 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
2160 Ok(QueryResponse {
2161 items: Vec::new(),
2162 next: None,
2163 plan: QueryPlan {
2164 per_source: merge_plans(self.plans),
2165 },
2166 errors: self.errors,
2167 })
2168 }
2169
2170 fn finish<T, U>(
2175 self,
2176 streams: Vec<Stream<T>>,
2177 budget: u32,
2178 first: Option<&Owed>,
2179 query: &str,
2180 qualify: impl Fn(&SourceName, T) -> U,
2181 ) -> Result<QueryResponse<U>, EngineError> {
2182 let (rows, states, owed) = merge(streams, budget, first);
2183 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
2184 Ok(QueryResponse {
2185 items: rows
2186 .into_iter()
2187 .map(|(name, item)| qualify(&name, item))
2188 .collect(),
2189 next,
2190 plan: QueryPlan {
2191 per_source: merge_plans(self.plans),
2192 },
2193 errors: self.errors,
2194 })
2195 }
2196}
2197
2198fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
2201 SourcePlan {
2202 source: source.name().clone(),
2203 kind: source.kind().to_owned(),
2204 pushed_down: outcomes.with(Outcome::PushedDown),
2205 applied_locally: outcomes.with(Outcome::AppliedLocally),
2206 emulated: outcomes.with(Outcome::Emulated),
2207 unavailable: outcomes.with(Outcome::Unavailable),
2208 pages_fetched: pages,
2209 }
2210}
2211
2212fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
2217 let mut merged: Vec<SourcePlan> = Vec::new();
2218 for plan in plans {
2219 if let Some(existing) = merged
2220 .iter_mut()
2221 .find(|existing| existing.source == plan.source)
2222 {
2223 existing.pushed_down.extend(plan.pushed_down);
2224 existing.applied_locally.extend(plan.applied_locally);
2225 existing.emulated.extend(plan.emulated);
2226 existing.unavailable.extend(plan.unavailable);
2227 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
2228 for list in [
2229 &mut existing.pushed_down,
2230 &mut existing.applied_locally,
2231 &mut existing.emulated,
2232 &mut existing.unavailable,
2233 ] {
2234 list.sort_unstable();
2235 list.dedup();
2236 }
2237 } else {
2238 merged.push(plan);
2239 }
2240 }
2241 merged
2242}
2243
2244fn walking<'a>(
2250 selected: Vec<&'a ResolvedSource>,
2251 states: &Option<Resumption>,
2252 kind: StreamKind,
2253) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
2254 let mut ready = Vec::new();
2255 let mut starts = Vec::new();
2256 for source in selected {
2257 if let Some(resume) = resume_at(states, source.name(), kind) {
2258 ready.push(source);
2259 starts.push(resume);
2260 }
2261 }
2262 (ready, starts)
2263}
2264
2265fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
2283 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
2284 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
2285}
2286
2287fn fingerprint(text: &str) -> String {
2289 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
2290 for byte in text.as_bytes() {
2291 hash ^= u64::from(*byte);
2292 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
2293 }
2294 format!("{hash:016x}")
2295}
2296
2297fn owed(document: &Option<Resumption>) -> Option<&Owed> {
2302 document.as_ref()?.owed.as_ref()
2303}
2304
2305fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
2307 match states {
2308 None => Some(Resume::default()),
2309 Some(document) => document
2310 .streams
2311 .iter()
2312 .find(|state| &state.source == source && state.stream == kind)
2313 .map(|state| state.resume.clone()),
2314 }
2315}
2316
2317fn resumption(
2348 engine: &Engine,
2349 token: Option<&PageToken>,
2350 reads: &[StreamKind],
2351 query: &str,
2352) -> Result<Option<Resumption>, EngineError> {
2353 let Some(document) = token.map(PageToken::decode) else {
2354 return Ok(None);
2355 };
2356
2357 if document.query != query {
2364 return Err(EngineError::Token {
2365 message: "this page token was written by a different query — resume the walk it \
2366 came from, or drop --page to start this one from the beginning"
2367 .to_owned(),
2368 });
2369 }
2370 let states = &document.streams;
2371
2372 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2373 for state in states {
2374 if !reads.contains(&state.stream) {
2375 return Err(EngineError::Token {
2376 message: format!(
2377 "this page token resumes {}, which this command does not read — it \
2378 was written by a different query",
2379 state.stream.describe()
2380 ),
2381 });
2382 }
2383 let ceiling = engine
2384 .ready()
2385 .find(|source| source.name() == &state.source)
2386 .map(ceiling);
2387 if ceiling.is_none() && !engine.has(&state.source) {
2388 return Err(EngineError::Token {
2389 message: format!(
2390 "this page token resumes a source called {:?}, which this \
2391 configuration does not have",
2392 state.source.as_str()
2393 ),
2394 });
2395 }
2396 if let Some(ceiling) = ceiling
2397 && state.resume.skip >= ceiling
2398 {
2399 return Err(EngineError::Token {
2400 message: format!(
2401 "this page token resumes {} rows into a page of source {:?}, which \
2402 serves at most {ceiling}",
2403 state.resume.skip,
2404 state.source.as_str()
2405 ),
2406 });
2407 }
2408 if seen.contains(&(&state.source, state.stream)) {
2409 return Err(EngineError::Token {
2410 message: format!(
2411 "this page token gives source {:?} two places to resume from",
2412 state.source.as_str()
2413 ),
2414 });
2415 }
2416 seen.push((&state.source, state.stream));
2417 }
2418
2419 if let Some(owed) = &document.owed
2424 && !document
2425 .streams
2426 .iter()
2427 .any(|state| state.source == owed.source && state.stream == owed.stream)
2428 {
2429 return Err(EngineError::Token {
2430 message: format!(
2431 "this page token owes the next row to a stream it does not resume, \
2432 {:?}'s {}",
2433 owed.source.as_str(),
2434 owed.stream.describe()
2435 ),
2436 });
2437 }
2438
2439 Ok(Some(document))
2440}
2441
2442fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2448 match selector {
2449 ProjectSelector::Any => ProjectFilter::Any,
2450 ProjectSelector::Orphans => ProjectFilter::Orphans,
2451 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2452 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2453 }
2454}
2455
2456fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2458 match fields {
2459 TextFields::Title => vec![Predicate::SearchTitle],
2460 TextFields::Content => vec![Predicate::SearchContent],
2461 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2462 }
2463}
2464
2465fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2472 match fields {
2473 TextFields::Title => capabilities.search_title.is_native(),
2474 TextFields::Content => capabilities.search_content.is_native(),
2475 TextFields::TitleOrContent => {
2476 capabilities.search_title.is_native() && capabilities.search_content.is_native()
2477 }
2478 }
2479}
2480
2481fn shape_tasks(
2483 capabilities: &Capabilities,
2484 filters: &Filters,
2485 project: &ProjectFilter,
2486 priorities: &[Priority],
2487 commented_since: Option<DateTime<Utc>>,
2488 metadata: &[MetadataMatch],
2489 origin: Option<&GlobalId>,
2490) -> TaskShape {
2491 let mut pushed = TaskQuery::default();
2492 let mut local = LocalTasks::default();
2493 let mut outcomes = Outcomes::default();
2494
2495 if !filters.labels.is_empty() {
2496 if capabilities.filter_by_label.is_native() {
2497 pushed.labels = filters.labels.clone();
2498 outcomes.record(Predicate::Label, Outcome::PushedDown);
2499 } else {
2500 local.labels = Some(filters.labels.clone());
2501 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2502 }
2503 }
2504 if !filters.statuses.is_empty() {
2505 if capabilities.filter_by_status.is_native() {
2506 pushed.statuses.clone_from(&filters.statuses);
2507 outcomes.record(Predicate::Status, Outcome::PushedDown);
2508 } else {
2509 local.statuses.clone_from(&filters.statuses);
2510 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2511 }
2512 }
2513 if !priorities.is_empty() {
2514 if capabilities.filter_by_priority.is_native() {
2515 pushed.priorities = priorities.to_vec();
2516 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2517 } else {
2518 local.priorities = priorities.to_vec();
2519 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2520 }
2521 }
2522 if let Some(since) = commented_since {
2523 if capabilities.filter_by_comment_activity.is_native() {
2524 pushed.commented_since = Some(since);
2525 outcomes.record(Predicate::CommentedSince, Outcome::PushedDown);
2526 } else {
2527 local.commented_since = Some(since);
2528 outcomes.record(Predicate::CommentedSince, Outcome::AppliedLocally);
2529 }
2530 }
2531 if !metadata.is_empty() {
2532 if capabilities.filter_by_metadata.is_native() {
2533 pushed.metadata = metadata.to_vec();
2534 outcomes.record(Predicate::Metadata, Outcome::PushedDown);
2535 } else {
2536 local.metadata = metadata.to_vec();
2537 outcomes.record(Predicate::Metadata, Outcome::AppliedLocally);
2538 }
2539 }
2540 if let Some(origin) = origin {
2541 if capabilities.filter_by_origin.is_native() {
2542 pushed.origin = Some(origin.to_string());
2545 outcomes.record(Predicate::Origin, Outcome::PushedDown);
2546 } else {
2547 local.origin = Some(origin.clone());
2548 outcomes.record(Predicate::Origin, Outcome::AppliedLocally);
2549 }
2550 }
2551 if let Some(text) = &filters.text {
2552 let predicates = text_predicates(text.fields);
2553 if searches_natively(capabilities, text.fields) {
2554 pushed.text = Some(text.clone());
2555 outcomes.record_all(predicates, Outcome::PushedDown);
2556 } else {
2557 local.text = Some(text.clone());
2558 outcomes.record_all(predicates, Outcome::AppliedLocally);
2559 }
2560 }
2561 match project {
2562 ProjectFilter::Any => {}
2563 ProjectFilter::Orphans => {
2564 if capabilities.orphan_tasks.is_native() {
2565 pushed.project = ProjectFilter::Orphans;
2566 outcomes.record(Predicate::Project, Outcome::PushedDown);
2567 } else {
2568 local.project = Some(ProjectFilter::Orphans);
2569 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2570 }
2571 }
2572 ProjectFilter::Is(id) => {
2573 if capabilities.projects.is_native() {
2574 pushed.project = ProjectFilter::Is(id.clone());
2575 outcomes.record(Predicate::Project, Outcome::PushedDown);
2576 } else {
2577 local.project = Some(ProjectFilter::Is(id.clone()));
2578 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2579 }
2580 }
2581 }
2582
2583 TaskShape {
2584 pushed,
2585 local,
2586 outcomes,
2587 }
2588}
2589
2590fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2592 let mut pushed = ProjectQuery::default();
2593 let mut local = LocalProjects::default();
2594 let mut outcomes = Outcomes::default();
2595
2596 if !filters.labels.is_empty() {
2597 if capabilities.filter_by_label.is_native() {
2598 pushed.labels = filters.labels.clone();
2599 outcomes.record(Predicate::Label, Outcome::PushedDown);
2600 } else {
2601 local.labels = Some(filters.labels.clone());
2602 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2603 }
2604 }
2605 if !filters.statuses.is_empty() {
2606 if capabilities.filter_by_status.is_native() {
2607 pushed.statuses.clone_from(&filters.statuses);
2608 outcomes.record(Predicate::Status, Outcome::PushedDown);
2609 } else {
2610 local.statuses.clone_from(&filters.statuses);
2611 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2612 }
2613 }
2614 if let Some(text) = &filters.text {
2615 let predicates = text_predicates(text.fields);
2616 if searches_natively(capabilities, text.fields) {
2617 pushed.text = Some(text.clone());
2618 outcomes.record_all(predicates, Outcome::PushedDown);
2619 } else {
2620 local.text = Some(text.clone());
2621 outcomes.record_all(predicates, Outcome::AppliedLocally);
2622 }
2623 }
2624
2625 ProjectShape {
2626 pushed,
2627 local,
2628 outcomes,
2629 }
2630}
2631
2632fn shape_documents(
2637 capabilities: &Capabilities,
2638 filters: &DocumentFilters,
2639 project: &ProjectFilter,
2640) -> DocumentShape {
2641 let mut pushed = DocumentQuery::default();
2642 let mut local = LocalDocuments::default();
2643 let mut outcomes = Outcomes::default();
2644
2645 if !filters.labels.is_empty() {
2646 if capabilities.filter_by_label.is_native() {
2647 pushed.labels = filters.labels.clone();
2648 outcomes.record(Predicate::Label, Outcome::PushedDown);
2649 } else {
2650 local.labels = Some(filters.labels.clone());
2651 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2652 }
2653 }
2654 if let Some(text) = &filters.text {
2655 let predicates = text_predicates(text.fields);
2656 if searches_natively(capabilities, text.fields) {
2657 pushed.text = Some(text.clone());
2658 outcomes.record_all(predicates, Outcome::PushedDown);
2659 } else {
2660 local.text = Some(text.clone());
2661 outcomes.record_all(predicates, Outcome::AppliedLocally);
2662 }
2663 }
2664 match project {
2665 ProjectFilter::Any => {}
2666 ProjectFilter::Orphans => {
2667 if capabilities.orphan_tasks.is_native() {
2668 pushed.project = ProjectFilter::Orphans;
2669 outcomes.record(Predicate::Project, Outcome::PushedDown);
2670 } else {
2671 local.project = Some(ProjectFilter::Orphans);
2672 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2673 }
2674 }
2675 ProjectFilter::Is(id) => {
2676 if capabilities.projects.is_native() {
2677 pushed.project = ProjectFilter::Is(id.clone());
2678 outcomes.record(Predicate::Project, Outcome::PushedDown);
2679 } else {
2680 local.project = Some(ProjectFilter::Is(id.clone()));
2681 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2682 }
2683 }
2684 }
2685
2686 DocumentShape {
2687 pushed,
2688 local,
2689 outcomes,
2690 }
2691}
2692
2693fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2695 match stream {
2696 StreamKind::Projects => {
2697 let shaped = shape_projects(capabilities, filters);
2698 HitShape {
2699 stream,
2700 tasks: TaskQuery::default(),
2701 projects: shaped.pushed,
2702 local_tasks: LocalTasks::default(),
2703 local_projects: shaped.local,
2704 outcomes: shaped.outcomes,
2705 }
2706 }
2707 StreamKind::Items | StreamKind::Tasks => {
2708 let shaped = shape_tasks(
2709 capabilities,
2710 filters,
2711 &ProjectFilter::Any,
2712 &[],
2713 None,
2714 &[],
2715 None,
2716 );
2717 HitShape {
2718 stream,
2719 tasks: shaped.pushed,
2720 projects: ProjectQuery::default(),
2721 local_tasks: shaped.local,
2722 local_projects: LocalProjects::default(),
2723 outcomes: shaped.outcomes,
2724 }
2725 }
2726 }
2727}
2728
2729fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2736 if compensating {
2737 ceiling
2738 } else {
2739 budget.min(ceiling)
2740 }
2741}
2742
2743fn ceiling(source: &ResolvedSource) -> u32 {
2745 source.source().capabilities().max_page_size.max(1)
2746}
2747
2748async fn fetch_tasks(
2750 source: &ResolvedSource,
2751 shape: &TaskShape,
2752 start: &Resume,
2753 budget: u32,
2754 calls: &AtomicU32,
2755) -> Result<Fetched<Task>, SourceError> {
2756 let compensating = shape.local != LocalTasks::default();
2757 walk(
2758 start,
2759 budget,
2760 page_size(compensating, budget, ceiling(source)),
2761 |task| shape.local.keeps(task),
2762 |cursor, limit| async move {
2763 calls.fetch_add(1, Ordering::Relaxed);
2764 let request = PageRequest { cursor, limit };
2765 let page = source.source().query_tasks(&shape.pushed, &request).await?;
2766 match shape.local.commented_since {
2767 None => Ok(page),
2768 Some(since) => comment::commented_since(source, &shape.local, page, since).await,
2769 }
2770 },
2771 )
2772 .await
2773}
2774
2775async fn fetch_projects(
2777 source: &ResolvedSource,
2778 shape: &ProjectShape,
2779 start: &Resume,
2780 budget: u32,
2781 calls: &AtomicU32,
2782) -> Result<Fetched<Project>, SourceError> {
2783 let compensating = shape.local != LocalProjects::default();
2784 walk(
2785 start,
2786 budget,
2787 page_size(compensating, budget, ceiling(source)),
2788 |project| shape.local.keeps(project),
2789 |cursor, limit| async move {
2790 calls.fetch_add(1, Ordering::Relaxed);
2791 let request = PageRequest { cursor, limit };
2792 source
2793 .source()
2794 .query_projects(&shape.pushed, &request)
2795 .await
2796 },
2797 )
2798 .await
2799}
2800
2801async fn fetch_documents(
2803 source: &ResolvedSource,
2804 shape: &DocumentShape,
2805 start: &Resume,
2806 budget: u32,
2807 calls: &AtomicU32,
2808) -> Result<Fetched<Document>, SourceError> {
2809 let compensating = shape.local != LocalDocuments::default();
2810 walk(
2811 start,
2812 budget,
2813 page_size(compensating, budget, ceiling(source)),
2814 |document| shape.local.keeps(document),
2815 |cursor, limit| async move {
2816 calls.fetch_add(1, Ordering::Relaxed);
2817 let request = PageRequest { cursor, limit };
2818 source
2819 .source()
2820 .query_documents(&shape.pushed, &request)
2821 .await
2822 },
2823 )
2824 .await
2825}
2826
2827async fn fetch_labels(
2829 source: &ResolvedSource,
2830 start: &Resume,
2831 budget: u32,
2832 calls: &AtomicU32,
2833) -> Result<Fetched<Label>, SourceError> {
2834 walk(
2835 start,
2836 budget,
2837 page_size(false, budget, ceiling(source)),
2838 |_| true,
2839 |cursor, limit| async move {
2840 calls.fetch_add(1, Ordering::Relaxed);
2841 let request = PageRequest { cursor, limit };
2842 source.source().labels(&request).await
2843 },
2844 )
2845 .await
2846}
2847
2848async fn fetch_hits(
2850 source: &ResolvedSource,
2851 shape: &HitShape,
2852 start: &Resume,
2853 budget: u32,
2854 calls: &AtomicU32,
2855) -> Result<Fetched<Found>, SourceError> {
2856 let ceiling = ceiling(source);
2857 match shape.stream {
2858 StreamKind::Projects => {
2859 let compensating = shape.local_projects != LocalProjects::default();
2860 walk(
2861 start,
2862 budget,
2863 page_size(compensating, budget, ceiling),
2864 |found| match found {
2865 Found::Project(project) => shape.local_projects.keeps(project),
2866 Found::Task(_) => true,
2867 },
2868 |cursor, limit| async move {
2869 calls.fetch_add(1, Ordering::Relaxed);
2870 let request = PageRequest { cursor, limit };
2871 let page = source
2872 .source()
2873 .query_projects(&shape.projects, &request)
2874 .await?;
2875 Ok(Page {
2876 items: page.items.into_iter().map(Found::Project).collect(),
2877 next: page.next,
2878 })
2879 },
2880 )
2881 .await
2882 }
2883 StreamKind::Items | StreamKind::Tasks => {
2884 let compensating = shape.local_tasks != LocalTasks::default();
2885 walk(
2886 start,
2887 budget,
2888 page_size(compensating, budget, ceiling),
2889 |found| match found {
2890 Found::Task(task) => shape.local_tasks.keeps(task),
2891 Found::Project(_) => true,
2892 },
2893 |cursor, limit| async move {
2894 calls.fetch_add(1, Ordering::Relaxed);
2895 let request = PageRequest { cursor, limit };
2896 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2897 Ok(Page {
2898 items: page.items.into_iter().map(Found::Task).collect(),
2899 next: page.next,
2900 })
2901 },
2902 )
2903 .await
2904 }
2905 }
2906}
2907
2908async fn forward_edges(
2910 source: &ResolvedSource,
2911 entity: Entity,
2912 id: &NativeId,
2913 request: &PageRequest,
2914) -> Result<Page<DependencyEdge>, SourceError> {
2915 match entity {
2916 Entity::Task => {
2917 source
2918 .source()
2919 .task_dependencies(id, Direction::DependsOn, request)
2920 .await
2921 }
2922 Entity::Project => {
2923 source
2924 .source()
2925 .project_dependencies(id, Direction::DependsOn, request)
2926 .await
2927 }
2928 }
2929}
2930
2931#[expect(
2940 clippy::too_many_arguments,
2941 reason = "every argument is one axis of one walk — the source, the item, the \
2942 direction, which of its two graphs, whether the reverse is emulated, where \
2943 to resume, how many rows to return and where to count calls. Grouping them \
2944 into a struct would name the same eight values one indirection further from \
2945 the loop that reads them."
2946)]
2947async fn fetch_edges(
2948 source: &ResolvedSource,
2949 native: &NativeId,
2950 direction: Direction,
2951 entity: Entity,
2952 emulating: bool,
2953 start: &Resume,
2954 budget: u32,
2955 calls: &AtomicU32,
2956) -> Result<Fetched<DependencyEdge>, SourceError> {
2957 let ceiling = ceiling(source);
2958 if !emulating {
2959 return walk(
2960 start,
2961 budget,
2962 page_size(false, budget, ceiling),
2963 |_| true,
2964 |cursor, limit| async move {
2965 calls.fetch_add(1, Ordering::Relaxed);
2966 let request = PageRequest { cursor, limit };
2967 match entity {
2968 Entity::Task => {
2969 source
2970 .source()
2971 .task_dependencies(native, direction, &request)
2972 .await
2973 }
2974 Entity::Project => {
2975 source
2976 .source()
2977 .project_dependencies(native, direction, &request)
2978 .await
2979 }
2980 }
2981 },
2982 )
2983 .await;
2984 }
2985
2986 walk(
2987 start,
2988 budget,
2989 ceiling,
2990 |_| true,
2991 |cursor, limit| async move {
2992 calls.fetch_add(1, Ordering::Relaxed);
2993 let request = PageRequest { cursor, limit };
2994 let (ids, next) = match entity {
2995 Entity::Task => {
2996 let page = source
2997 .source()
2998 .query_tasks(&TaskQuery::default(), &request)
2999 .await?;
3000 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
3001 (ids, page.next)
3002 }
3003 Entity::Project => {
3004 let page = source
3005 .source()
3006 .query_projects(&ProjectQuery::default(), &request)
3007 .await?;
3008 let ids: Vec<NativeId> =
3009 page.items.into_iter().map(|project| project.id).collect();
3010 (ids, page.next)
3011 }
3012 };
3013
3014 let mut edges = Vec::new();
3015 for id in ids {
3016 let mut inner: Option<Cursor> = None;
3017 loop {
3018 calls.fetch_add(1, Ordering::Relaxed);
3019 let request = PageRequest {
3020 cursor: inner.clone(),
3021 limit,
3022 };
3023 let page = forward_edges(source, entity, &id, &request).await?;
3024 fits(page.items.len(), limit)?;
3028 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
3029 unrepeated(
3030 page.next.as_ref(),
3031 inner.as_ref(),
3032 "its forward edges were being scanned",
3033 )?;
3034 match page.next {
3035 Some(cursor) => inner = Some(cursor),
3036 None => break,
3037 }
3038 }
3039 }
3040
3041 Ok(Page { items: edges, next })
3042 },
3043 )
3044 .await
3045}