1mod copy;
19mod fetch;
20mod join;
21mod local;
22mod resume;
23
24use std::collections::BTreeMap;
25use std::num::NonZeroU32;
26use std::sync::atomic::{AtomicU32, Ordering};
27
28use onetaskgraph_plugin_api::{
29 Capabilities, Cursor, DependencyEdge, Direction, Label, LabelFilter, NativeId, Page,
30 PageRequest, Project, ProjectFilter, ProjectQuery, SecretResolver, SourceError, SourceName,
31 StatusCategory, Task, TaskQuery, TextFields, TextQuery,
32};
33use schemars::JsonSchema;
34use serde::{Deserialize, Serialize};
35
36use crate::GlobalId;
37use crate::config::Config;
38use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
39use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
40
41use fetch::{Fetched, Stream, fits, merge, walk};
42use join::join_all;
43use local::{LocalProjects, LocalTasks};
44pub(crate) use resume::{Owed, Resumption, StreamState};
45use resume::{Resume, StreamKind};
46
47pub use copy::{CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy};
48pub use local::ProjectSelector;
49
50#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
55pub struct Qualified<T> {
56 pub id: GlobalId,
58 pub item: T,
60}
61
62#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
67pub struct QualifiedEdge {
68 pub from: QualifiedEndpoint,
73 pub to: QualifiedEndpoint,
75 pub kind: onetaskgraph_plugin_api::DependencyKind,
77}
78
79#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
81pub struct QualifiedEndpoint {
82 pub id: GlobalId,
84 pub kind: onetaskgraph_plugin_api::ItemKind,
86}
87
88impl std::fmt::Display for QualifiedEndpoint {
89 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
90 self.id.fmt(formatter)
91 }
92}
93
94#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
96#[serde(tag = "kind", rename_all = "kebab-case")]
97pub enum SearchHit {
98 Task(Qualified<Task>),
100 Project(Qualified<Project>),
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
106#[serde(rename_all = "kebab-case")]
107pub enum SearchKind {
108 Tasks,
110 Projects,
112 #[default]
114 Both,
115}
116
117#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
119pub struct SourceListing {
120 pub source: SourceName,
122 pub kind: String,
134 #[serde(flatten)]
136 pub state: SourceState,
137}
138
139#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
141#[serde(tag = "state", rename_all = "kebab-case")]
142pub enum SourceState {
143 Available {
145 capabilities: Capabilities,
147 },
148 Unavailable {
150 error: SourceError,
152 },
153}
154
155#[derive(Debug, Clone, PartialEq)]
157pub struct Paging {
158 pub limit: NonZeroU32,
160 pub token: Option<PageToken>,
162}
163
164#[derive(Debug, Clone, Default, PartialEq)]
166pub struct Filters {
167 pub text: Option<TextQuery>,
169 pub labels: LabelFilter,
171 pub statuses: Vec<StatusCategory>,
173}
174
175#[derive(Debug, Clone)]
177pub struct TaskRequest {
178 pub sources: Vec<SourceName>,
180 pub filters: Filters,
182 pub project: ProjectSelector,
184 pub paging: Paging,
186}
187
188#[derive(Debug, Clone)]
190pub struct ProjectRequest {
191 pub sources: Vec<SourceName>,
193 pub filters: Filters,
195 pub paging: Paging,
197}
198
199#[derive(Debug, Clone)]
201pub struct LabelRequest {
202 pub sources: Vec<SourceName>,
204 pub paging: Paging,
206}
207
208#[derive(Debug, Clone)]
210pub struct SearchRequest {
211 pub sources: Vec<SourceName>,
213 pub text: TextQuery,
215 pub kind: SearchKind,
217 pub paging: Paging,
219}
220
221#[derive(Debug, Clone)]
223pub struct DependencyRequest {
224 pub id: GlobalId,
226 pub direction: Direction,
228 pub paging: Paging,
230}
231
232#[derive(Debug, Clone, PartialEq, thiserror::Error)]
240pub enum EngineError {
241 #[error(
243 "no source named {name:?} is configured\n\
244 next: name one of the configured sources ({configured}), or add {name:?} under \
245 `sources` — `onetaskgraph sources list` shows what this configuration has."
246 )]
247 UnknownSource {
248 name: String,
250 configured: String,
252 },
253
254 #[error(
256 "{message}\n\
257 next: page with a token exactly as the previous page reported it, and against \
258 the same configuration — or drop `--page` to start the walk again."
259 )]
260 Token {
261 message: String,
263 },
264
265 #[error(
267 "no sources are configured\n\
268 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
269 prints what each plugin accepts."
270 )]
271 NoSources,
272
273 #[error(
275 "source {name} cannot be written: its plugin is {kind}, which has no write \
276 side\n\
277 next: copy into a source whose plugin can be written — `onetaskgraph sources \
278 list` reports each one's plugin."
279 )]
280 NotWritable {
281 name: String,
283 kind: String,
285 },
286
287 #[error(
289 "the destination source {name} could not be built: {error}\n\
290 next: fix that source — `onetaskgraph sources list` reports its state — then \
291 copy again."
292 )]
293 DestinationUnavailable {
294 name: String,
296 error: SourceError,
298 },
299
300 #[error(
302 "no item with the id {id}\n\
303 next: check the id, or list what is there — `onetaskgraph task list` and \
304 `onetaskgraph project list` report what the configured sources hold."
305 )]
306 NoSuchItem {
307 id: String,
309 },
310
311 #[error(
316 "{item} was copied from {origin}, which that destination no longer holds\n\
317 next: re-run with --recreate to create a new item there instead, or restore \
318 {origin}."
319 )]
320 StaleOrigin {
321 item: String,
323 origin: String,
325 },
326
327 #[error(
333 "source {name} could not do it: {error}\n\
334 next: fix what the source named above, then copy again."
335 )]
336 SourceRefused {
337 name: String,
339 error: SourceError,
341 },
342
343 #[error(
351 "the copy failed and could not be undone.\n\
352 it failed because: {error}\n\
353 it could not be undone because: {refusal}\n\
354 so the destination still holds: {left_behind}\n\
355 next: remove those items at the destination, then copy again."
356 )]
357 CopyNotUndone {
358 error: Box<EngineError>,
360 left_behind: LeftBehind,
366 refusal: SourceError,
368 },
369}
370
371#[derive(Debug, Clone, PartialEq)]
380pub struct LeftBehind {
381 first: GlobalId,
383 rest: Vec<GlobalId>,
385}
386
387impl LeftBehind {
388 #[must_use]
390 pub fn new(first: GlobalId) -> Self {
391 Self {
392 first,
393 rest: Vec::new(),
394 }
395 }
396
397 pub fn push(&mut self, id: GlobalId) {
399 self.rest.push(id);
400 }
401
402 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
404 std::iter::once(&self.first).chain(self.rest.iter())
405 }
406}
407
408impl std::fmt::Display for LeftBehind {
409 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
411 write!(formatter, "{}", self.first)?;
412 for id in &self.rest {
413 write!(formatter, ", {id}")?;
414 }
415 Ok(())
416 }
417}
418
419pub enum ConfiguredSource {
426 Ready(ResolvedSource),
428 Unavailable(UnavailableSource),
430}
431
432impl ConfiguredSource {
433 #[must_use]
435 pub fn name(&self) -> &SourceName {
436 match self {
437 Self::Ready(source) => source.name(),
438 Self::Unavailable(source) => source.name(),
439 }
440 }
441}
442
443pub struct Engine {
445 sources: Vec<ConfiguredSource>,
447 selection: Vec<SourceName>,
449}
450
451impl Engine {
452 #[must_use]
460 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
461 let (ready, unavailable) = resolve_available(config, secrets);
462 Self::new(
463 ready
464 .into_iter()
465 .map(ConfiguredSource::Ready)
466 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
467 .collect(),
468 config.selected_sources(),
469 )
470 }
471
472 #[must_use]
475 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
476 Self { sources, selection }
477 }
478
479 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
481 self.sources.iter().filter_map(|source| match source {
482 ConfiguredSource::Ready(ready) => Some(ready),
483 ConfiguredSource::Unavailable(_) => None,
484 })
485 }
486
487 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
489 self.sources.iter().filter_map(|source| match source {
490 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
491 ConfiguredSource::Ready(_) => None,
492 })
493 }
494
495 #[must_use]
497 pub fn listing(&self) -> Vec<SourceListing> {
498 let mut listings: Vec<SourceListing> = self
499 .ready()
500 .map(|source| SourceListing {
501 source: source.name().clone(),
502 kind: source.kind().to_owned(),
503 state: SourceState::Available {
504 capabilities: source.source().capabilities(),
505 },
506 })
507 .chain(self.unavailable().map(|source| SourceListing {
508 source: source.name().clone(),
509 kind: source.kind().to_owned(),
510 state: SourceState::Unavailable {
511 error: source.error().clone(),
512 },
513 }))
514 .collect();
515 listings.sort_by(|left, right| left.source.cmp(&right.source));
516 listings
517 }
518
519 #[must_use]
525 pub fn has(&self, name: &SourceName) -> bool {
526 self.sources.iter().any(|source| source.name() == name)
527 }
528
529 pub async fn tasks(
537 &self,
538 request: &TaskRequest,
539 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
540 let mut names = self.resolve_selection(&request.sources)?;
541 if let ProjectSelector::Qualified(id) = &request.project {
546 self.known(&id.source)?;
547 names.retain(|name| name == &id.source);
548 }
549 let query = shape("task-list", &names, &(&request.filters, &request.project));
550 let states = resumption(
551 self,
552 request.paging.token.as_ref(),
553 &[StreamKind::Items],
554 &query,
555 )?;
556 let budget = request.paging.limit.get();
557
558 let mut answer = Answer::new();
559 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
560
561 let shapes: Vec<TaskShape> = ready
562 .iter()
563 .map(|source| {
564 shape_tasks(
565 &source.source().capabilities(),
566 &request.filters,
567 &project_filter(&request.project),
568 )
569 })
570 .collect();
571 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
572 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
573
574 let walks = ready
575 .iter()
576 .enumerate()
577 .map(|(index, source)| {
578 fetch_tasks(
579 source,
580 &shapes[index],
581 &starts[index],
582 budget,
583 &counters[index],
584 )
585 })
586 .collect();
587
588 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
589 answer.finish(
590 streams,
591 budget,
592 owed(&states),
593 &query,
594 |name, task: Task| Qualified {
595 id: GlobalId::new(name.clone(), task.id.clone()),
596 item: task,
597 },
598 )
599 }
600
601 pub async fn projects(
607 &self,
608 request: &ProjectRequest,
609 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
610 let names = self.resolve_selection(&request.sources)?;
611 let query = shape("project-list", &names, &request.filters);
612 let states = resumption(
613 self,
614 request.paging.token.as_ref(),
615 &[StreamKind::Items],
616 &query,
617 )?;
618 let budget = request.paging.limit.get();
619
620 let mut answer = Answer::new();
621 let mut with_projects = Vec::new();
627 for source in answer.split(self, &names) {
628 if source.source().capabilities().projects.is_native() {
629 with_projects.push(source);
630 } else {
631 answer.unreachable_predicate(source, Predicate::Project);
632 }
633 }
634 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
635
636 let shapes: Vec<ProjectShape> = ready
637 .iter()
638 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
639 .collect();
640 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
641 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
642
643 let walks = ready
644 .iter()
645 .enumerate()
646 .map(|(index, source)| {
647 fetch_projects(
648 source,
649 &shapes[index],
650 &starts[index],
651 budget,
652 &counters[index],
653 )
654 })
655 .collect();
656
657 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
658 answer.finish(
659 streams,
660 budget,
661 owed(&states),
662 &query,
663 |name, project: Project| Qualified {
664 id: GlobalId::new(name.clone(), project.id.clone()),
665 item: project,
666 },
667 )
668 }
669
670 pub async fn labels(
676 &self,
677 request: &LabelRequest,
678 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
679 let names = self.resolve_selection(&request.sources)?;
680 let query = shape("label-list", &names, &());
681 let states = resumption(
682 self,
683 request.paging.token.as_ref(),
684 &[StreamKind::Items],
685 &query,
686 )?;
687 let budget = request.paging.limit.get();
688
689 let mut answer = Answer::new();
690 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
691 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
692 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
693
694 let walks = ready
695 .iter()
696 .enumerate()
697 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
698 .collect();
699
700 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
701 answer.finish(
702 streams,
703 budget,
704 owed(&states),
705 &query,
706 |name, label: Label| Qualified {
707 id: GlobalId::new(name.clone(), label.id.clone()),
708 item: label,
709 },
710 )
711 }
712
713 pub async fn search(
719 &self,
720 request: &SearchRequest,
721 ) -> Result<QueryResponse<SearchHit>, EngineError> {
722 let names = self.resolve_selection(&request.sources)?;
723 let reads: &[StreamKind] = match request.kind {
727 SearchKind::Tasks => &[StreamKind::Tasks],
728 SearchKind::Projects => &[StreamKind::Projects],
729 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
730 };
731 let query = shape("search", &names, &request.text);
736 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
737 let budget = request.paging.limit.get();
738 let filters = Filters {
739 text: Some(request.text.clone()),
740 ..Filters::default()
741 };
742
743 let mut answer = Answer::new();
744
745 let mut ready = Vec::new();
748 let mut kinds = Vec::new();
749 let mut starts = Vec::new();
750 for source in answer.split(self, &names) {
751 let mut streams = Vec::new();
752 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
753 streams.push(StreamKind::Tasks);
754 }
755 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
756 if source.source().capabilities().projects.is_native() {
757 streams.push(StreamKind::Projects);
758 } else {
759 answer.unreachable_predicate(source, Predicate::Project);
760 }
761 }
762 for stream in streams {
763 if let Some(resume) = resume_at(&states, source.name(), stream) {
764 ready.push(source);
765 kinds.push(stream);
766 starts.push(resume);
767 }
768 }
769 }
770
771 let shapes: Vec<HitShape> = ready
772 .iter()
773 .zip(kinds.iter())
774 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
775 .collect();
776 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
777 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
778
779 let walks = ready
780 .iter()
781 .enumerate()
782 .map(|(index, source)| {
783 fetch_hits(
784 source,
785 &shapes[index],
786 &starts[index],
787 budget,
788 &counters[index],
789 )
790 })
791 .collect();
792
793 let streams =
794 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
795 answer.finish(
796 streams,
797 budget,
798 owed(&states),
799 &query,
800 |name, found: Found| match found {
801 Found::Task(task) => SearchHit::Task(Qualified {
802 id: GlobalId::new(name.clone(), task.id.clone()),
803 item: task,
804 }),
805 Found::Project(project) => SearchHit::Project(Qualified {
806 id: GlobalId::new(name.clone(), project.id.clone()),
807 item: project,
808 }),
809 },
810 )
811 }
812
813 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
820 let name = self.known(&id.source)?;
821 let mut answer = Answer::new();
822 let selected = answer.split(self, std::slice::from_ref(&name));
823 let Some(source) = selected.first() else {
824 return answer.nothing();
825 };
826 let found = source.source().get_task(&id.native).await;
827 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
828 answer.one(source, found, |task| Qualified {
829 id: qualified,
830 item: task,
831 })
832 }
833
834 pub async fn project(
840 &self,
841 id: &GlobalId,
842 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
843 let name = self.known(&id.source)?;
844 let mut answer = Answer::new();
845 let selected = answer.split(self, std::slice::from_ref(&name));
846 let Some(source) = selected.first() else {
847 return answer.nothing();
848 };
849 let found = source.source().get_project(&id.native).await;
850 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
851 answer.one(source, found, |project| Qualified {
852 id: qualified,
853 item: project,
854 })
855 }
856
857 pub async fn task_dependencies(
864 &self,
865 request: &DependencyRequest,
866 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
867 self.dependencies(request, Entity::Task).await
868 }
869
870 pub async fn project_dependencies(
876 &self,
877 request: &DependencyRequest,
878 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
879 self.dependencies(request, Entity::Project).await
880 }
881
882 async fn dependencies(
885 &self,
886 request: &DependencyRequest,
887 entity: Entity,
888 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
889 let name = self.known(&request.id.source)?;
890 let query = shape(
891 "dependencies",
892 std::slice::from_ref(&name),
893 &(entity, &request.id.native, request.direction),
894 );
895 let states = resumption(
896 self,
897 request.paging.token.as_ref(),
898 &[StreamKind::Items],
899 &query,
900 )?;
901 let budget = request.paging.limit.get();
902
903 let mut answer = Answer::new();
904 let (ready, starts) = walking(
905 answer.split(self, std::slice::from_ref(&name)),
906 &states,
907 StreamKind::Items,
908 );
909 let Some(source) = ready.first() else {
910 return answer.nothing();
911 };
912
913 let capabilities = source.source().capabilities();
914 let support = match entity {
915 Entity::Task => capabilities.task_dependencies,
916 Entity::Project => capabilities.project_dependencies,
917 };
918 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
922 let mut outcomes = Outcomes::default();
923 if request.direction == Direction::DependedOnBy {
924 if emulating {
925 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
926 } else {
927 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
928 }
929 }
930
931 let counters = vec![AtomicU32::new(0)];
932 let walked = fetch_edges(
933 source,
934 &request.id.native,
935 request.direction,
936 entity,
937 emulating,
938 &starts[0],
939 budget,
940 &counters[0],
941 )
942 .await;
943
944 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
945 answer.finish(
946 streams,
947 budget,
948 owed(&states),
949 &query,
950 |name, edge: DependencyEdge| QualifiedEdge {
951 from: qualify_endpoint(name, edge.from),
952 to: qualify_endpoint(name, edge.to),
953 kind: edge.kind,
954 },
955 )
956 }
957
958 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
960 if asked.is_empty() {
961 if self.selection.is_empty() {
962 return Err(EngineError::NoSources);
963 }
964 return Ok(self.selection.clone());
965 }
966 asked.iter().map(|name| self.known(name)).collect()
967 }
968
969 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
971 if self.has(name) {
972 return Ok(name.clone());
973 }
974 if self.sources.is_empty() {
975 return Err(EngineError::NoSources);
976 }
977 Err(EngineError::UnknownSource {
978 name: name.to_string(),
979 configured: self
980 .listing()
981 .iter()
982 .map(|listing| listing.source.to_string())
983 .collect::<Vec<_>>()
984 .join(", "),
985 })
986 }
987}
988
989fn qualify_endpoint(
990 source: &SourceName,
991 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
992) -> QualifiedEndpoint {
993 let kind = endpoint.kind;
994 let is_qualified = endpoint.is_qualified();
995 let endpoint_id = endpoint.into_id();
996 QualifiedEndpoint {
997 id: if is_qualified {
998 endpoint_id
999 .parse()
1000 .expect("plugin-api validates qualified dependency endpoints")
1001 } else {
1002 GlobalId::new(source.clone(), NativeId(endpoint_id))
1003 },
1004 kind,
1005 }
1006}
1007
1008#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1010enum Entity {
1011 Task,
1013 Project,
1015}
1016
1017enum Found {
1019 Task(Task),
1021 Project(Project),
1023}
1024
1025#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1027enum Outcome {
1028 PushedDown,
1030 AppliedLocally,
1032 Emulated,
1034 Unavailable,
1036}
1037
1038#[derive(Debug, Clone, Default, PartialEq)]
1050struct Outcomes(BTreeMap<Predicate, Outcome>);
1051
1052impl Outcomes {
1053 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1059 self.0.insert(predicate, outcome);
1060 }
1061
1062 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1065 for predicate in predicates {
1066 self.record(predicate, outcome);
1067 }
1068 }
1069
1070 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1072 self.0
1073 .iter()
1074 .filter(|(_, recorded)| **recorded == outcome)
1075 .map(|(predicate, _)| *predicate)
1076 .collect()
1077 }
1078}
1079
1080struct TaskShape {
1082 pushed: TaskQuery,
1084 local: LocalTasks,
1086 outcomes: Outcomes,
1088}
1089
1090struct ProjectShape {
1092 pushed: ProjectQuery,
1094 local: LocalProjects,
1096 outcomes: Outcomes,
1098}
1099
1100struct HitShape {
1102 stream: StreamKind,
1104 tasks: TaskQuery,
1106 projects: ProjectQuery,
1108 local_tasks: LocalTasks,
1110 local_projects: LocalProjects,
1112 outcomes: Outcomes,
1114}
1115
1116struct Answer {
1122 plans: Vec<SourcePlan>,
1124 errors: Vec<SourceFailure>,
1126}
1127
1128impl Answer {
1129 fn new() -> Self {
1130 Self {
1131 plans: Vec::new(),
1132 errors: Vec::new(),
1133 }
1134 }
1135
1136 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1141 let mut selected = Vec::new();
1142 for name in names {
1143 match engine.sources.iter().find(|source| source.name() == name) {
1144 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1145 Some(ConfiguredSource::Unavailable(source)) => {
1146 self.errors.push(source.failure());
1147 }
1148 None => {}
1149 }
1150 }
1151 selected
1152 }
1153
1154 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1156 let mut outcomes = Outcomes::default();
1157 outcomes.record(predicate, Outcome::Unavailable);
1158 self.plans.push(plan_for(source, outcomes, 0));
1159 }
1160
1161 fn collect<T>(
1163 &mut self,
1164 ready: &[&ResolvedSource],
1165 walked: Vec<Result<Fetched<T>, SourceError>>,
1166 counters: &[AtomicU32],
1167 outcomes: Vec<Outcomes>,
1168 ) -> Vec<Stream<T>> {
1169 let kinds = vec![StreamKind::Items; ready.len()];
1170 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1171 }
1172
1173 fn collect_streams<T>(
1175 &mut self,
1176 ready: &[&ResolvedSource],
1177 kinds: &[StreamKind],
1178 walked: Vec<Result<Fetched<T>, SourceError>>,
1179 counters: &[AtomicU32],
1180 outcomes: Vec<Outcomes>,
1181 ) -> Vec<Stream<T>> {
1182 let mut streams = Vec::new();
1183 for (index, result) in walked.into_iter().enumerate() {
1184 let source = ready[index];
1185 let pages = counters[index].load(Ordering::Relaxed);
1186 self.plans
1187 .push(plan_for(source, outcomes[index].clone(), pages));
1188 match result {
1189 Ok(fetched) => streams.push(Stream {
1190 source: source.name().clone(),
1191 kind: kinds[index],
1192 fetched,
1193 }),
1194 Err(error) => self.errors.push(SourceFailure {
1197 source: source.name().clone(),
1198 error,
1199 }),
1200 }
1201 }
1202 streams
1203 }
1204
1205 fn one<T, U>(
1207 mut self,
1208 source: &ResolvedSource,
1209 found: Result<Option<T>, SourceError>,
1210 qualify: impl FnOnce(T) -> U,
1211 ) -> Result<QueryResponse<U>, EngineError> {
1212 self.plans.push(plan_for(source, Outcomes::default(), 1));
1213 let items = match found {
1214 Ok(Some(item)) => vec![qualify(item)],
1215 Ok(None) => Vec::new(),
1216 Err(error) => {
1217 self.errors.push(SourceFailure {
1218 source: source.name().clone(),
1219 error,
1220 });
1221 Vec::new()
1222 }
1223 };
1224 Ok(QueryResponse {
1225 items,
1226 next: None,
1227 plan: QueryPlan {
1228 per_source: merge_plans(self.plans),
1229 },
1230 errors: self.errors,
1231 })
1232 }
1233
1234 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1236 Ok(QueryResponse {
1237 items: Vec::new(),
1238 next: None,
1239 plan: QueryPlan {
1240 per_source: merge_plans(self.plans),
1241 },
1242 errors: self.errors,
1243 })
1244 }
1245
1246 fn finish<T, U>(
1251 self,
1252 streams: Vec<Stream<T>>,
1253 budget: u32,
1254 first: Option<&Owed>,
1255 query: &str,
1256 qualify: impl Fn(&SourceName, T) -> U,
1257 ) -> Result<QueryResponse<U>, EngineError> {
1258 let (rows, states, owed) = merge(streams, budget, first);
1259 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1260 Ok(QueryResponse {
1261 items: rows
1262 .into_iter()
1263 .map(|(name, item)| qualify(&name, item))
1264 .collect(),
1265 next,
1266 plan: QueryPlan {
1267 per_source: merge_plans(self.plans),
1268 },
1269 errors: self.errors,
1270 })
1271 }
1272}
1273
1274fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1277 SourcePlan {
1278 source: source.name().clone(),
1279 kind: source.kind().to_owned(),
1280 pushed_down: outcomes.with(Outcome::PushedDown),
1281 applied_locally: outcomes.with(Outcome::AppliedLocally),
1282 emulated: outcomes.with(Outcome::Emulated),
1283 unavailable: outcomes.with(Outcome::Unavailable),
1284 pages_fetched: pages,
1285 }
1286}
1287
1288fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1293 let mut merged: Vec<SourcePlan> = Vec::new();
1294 for plan in plans {
1295 if let Some(existing) = merged
1296 .iter_mut()
1297 .find(|existing| existing.source == plan.source)
1298 {
1299 existing.pushed_down.extend(plan.pushed_down);
1300 existing.applied_locally.extend(plan.applied_locally);
1301 existing.emulated.extend(plan.emulated);
1302 existing.unavailable.extend(plan.unavailable);
1303 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1304 for list in [
1305 &mut existing.pushed_down,
1306 &mut existing.applied_locally,
1307 &mut existing.emulated,
1308 &mut existing.unavailable,
1309 ] {
1310 list.sort_unstable();
1311 list.dedup();
1312 }
1313 } else {
1314 merged.push(plan);
1315 }
1316 }
1317 merged
1318}
1319
1320fn walking<'a>(
1326 selected: Vec<&'a ResolvedSource>,
1327 states: &Option<Resumption>,
1328 kind: StreamKind,
1329) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1330 let mut ready = Vec::new();
1331 let mut starts = Vec::new();
1332 for source in selected {
1333 if let Some(resume) = resume_at(states, source.name(), kind) {
1334 ready.push(source);
1335 starts.push(resume);
1336 }
1337 }
1338 (ready, starts)
1339}
1340
1341fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1359 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1360 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1361}
1362
1363fn fingerprint(text: &str) -> String {
1365 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1366 for byte in text.as_bytes() {
1367 hash ^= u64::from(*byte);
1368 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1369 }
1370 format!("{hash:016x}")
1371}
1372
1373fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1378 document.as_ref()?.owed.as_ref()
1379}
1380
1381fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1383 match states {
1384 None => Some(Resume::default()),
1385 Some(document) => document
1386 .streams
1387 .iter()
1388 .find(|state| &state.source == source && state.stream == kind)
1389 .map(|state| state.resume.clone()),
1390 }
1391}
1392
1393fn resumption(
1424 engine: &Engine,
1425 token: Option<&PageToken>,
1426 reads: &[StreamKind],
1427 query: &str,
1428) -> Result<Option<Resumption>, EngineError> {
1429 let Some(document) = token.map(PageToken::decode) else {
1430 return Ok(None);
1431 };
1432
1433 if document.query != query {
1440 return Err(EngineError::Token {
1441 message: "this page token was written by a different query — resume the walk it \
1442 came from, or drop --page to start this one from the beginning"
1443 .to_owned(),
1444 });
1445 }
1446 let states = &document.streams;
1447
1448 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
1449 for state in states {
1450 if !reads.contains(&state.stream) {
1451 return Err(EngineError::Token {
1452 message: format!(
1453 "this page token resumes {}, which this command does not read — it \
1454 was written by a different query",
1455 state.stream.describe()
1456 ),
1457 });
1458 }
1459 let ceiling = engine
1460 .ready()
1461 .find(|source| source.name() == &state.source)
1462 .map(ceiling);
1463 if ceiling.is_none() && !engine.has(&state.source) {
1464 return Err(EngineError::Token {
1465 message: format!(
1466 "this page token resumes a source called {:?}, which this \
1467 configuration does not have",
1468 state.source.as_str()
1469 ),
1470 });
1471 }
1472 if let Some(ceiling) = ceiling
1473 && state.resume.skip >= ceiling
1474 {
1475 return Err(EngineError::Token {
1476 message: format!(
1477 "this page token resumes {} rows into a page of source {:?}, which \
1478 serves at most {ceiling}",
1479 state.resume.skip,
1480 state.source.as_str()
1481 ),
1482 });
1483 }
1484 if seen.contains(&(&state.source, state.stream)) {
1485 return Err(EngineError::Token {
1486 message: format!(
1487 "this page token gives source {:?} two places to resume from",
1488 state.source.as_str()
1489 ),
1490 });
1491 }
1492 seen.push((&state.source, state.stream));
1493 }
1494
1495 if let Some(owed) = &document.owed
1500 && !document
1501 .streams
1502 .iter()
1503 .any(|state| state.source == owed.source && state.stream == owed.stream)
1504 {
1505 return Err(EngineError::Token {
1506 message: format!(
1507 "this page token owes the next row to a stream it does not resume, \
1508 {:?}'s {}",
1509 owed.source.as_str(),
1510 owed.stream.describe()
1511 ),
1512 });
1513 }
1514
1515 Ok(Some(document))
1516}
1517
1518fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
1524 match selector {
1525 ProjectSelector::Any => ProjectFilter::Any,
1526 ProjectSelector::Orphans => ProjectFilter::Orphans,
1527 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
1528 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
1529 }
1530}
1531
1532fn text_predicates(fields: TextFields) -> Vec<Predicate> {
1534 match fields {
1535 TextFields::Title => vec![Predicate::SearchTitle],
1536 TextFields::Content => vec![Predicate::SearchContent],
1537 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
1538 }
1539}
1540
1541fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
1548 match fields {
1549 TextFields::Title => capabilities.search_title.is_native(),
1550 TextFields::Content => capabilities.search_content.is_native(),
1551 TextFields::TitleOrContent => {
1552 capabilities.search_title.is_native() && capabilities.search_content.is_native()
1553 }
1554 }
1555}
1556
1557fn shape_tasks(
1559 capabilities: &Capabilities,
1560 filters: &Filters,
1561 project: &ProjectFilter,
1562) -> TaskShape {
1563 let mut pushed = TaskQuery::default();
1564 let mut local = LocalTasks::default();
1565 let mut outcomes = Outcomes::default();
1566
1567 if !filters.labels.is_empty() {
1568 if capabilities.filter_by_label.is_native() {
1569 pushed.labels = filters.labels.clone();
1570 outcomes.record(Predicate::Label, Outcome::PushedDown);
1571 } else {
1572 local.labels = Some(filters.labels.clone());
1573 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1574 }
1575 }
1576 if !filters.statuses.is_empty() {
1577 if capabilities.filter_by_status.is_native() {
1578 pushed.statuses.clone_from(&filters.statuses);
1579 outcomes.record(Predicate::Status, Outcome::PushedDown);
1580 } else {
1581 local.statuses.clone_from(&filters.statuses);
1582 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1583 }
1584 }
1585 if let Some(text) = &filters.text {
1586 let predicates = text_predicates(text.fields);
1587 if searches_natively(capabilities, text.fields) {
1588 pushed.text = Some(text.clone());
1589 outcomes.record_all(predicates, Outcome::PushedDown);
1590 } else {
1591 local.text = Some(text.clone());
1592 outcomes.record_all(predicates, Outcome::AppliedLocally);
1593 }
1594 }
1595 match project {
1596 ProjectFilter::Any => {}
1597 ProjectFilter::Orphans => {
1598 if capabilities.orphan_tasks.is_native() {
1599 pushed.project = ProjectFilter::Orphans;
1600 outcomes.record(Predicate::Project, Outcome::PushedDown);
1601 } else {
1602 local.project = Some(ProjectFilter::Orphans);
1603 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1604 }
1605 }
1606 ProjectFilter::Is(id) => {
1607 if capabilities.projects.is_native() {
1608 pushed.project = ProjectFilter::Is(id.clone());
1609 outcomes.record(Predicate::Project, Outcome::PushedDown);
1610 } else {
1611 local.project = Some(ProjectFilter::Is(id.clone()));
1612 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1613 }
1614 }
1615 }
1616
1617 TaskShape {
1618 pushed,
1619 local,
1620 outcomes,
1621 }
1622}
1623
1624fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
1626 let mut pushed = ProjectQuery::default();
1627 let mut local = LocalProjects::default();
1628 let mut outcomes = Outcomes::default();
1629
1630 if !filters.labels.is_empty() {
1631 if capabilities.filter_by_label.is_native() {
1632 pushed.labels = filters.labels.clone();
1633 outcomes.record(Predicate::Label, Outcome::PushedDown);
1634 } else {
1635 local.labels = Some(filters.labels.clone());
1636 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1637 }
1638 }
1639 if !filters.statuses.is_empty() {
1640 if capabilities.filter_by_status.is_native() {
1641 pushed.statuses.clone_from(&filters.statuses);
1642 outcomes.record(Predicate::Status, Outcome::PushedDown);
1643 } else {
1644 local.statuses.clone_from(&filters.statuses);
1645 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1646 }
1647 }
1648 if let Some(text) = &filters.text {
1649 let predicates = text_predicates(text.fields);
1650 if searches_natively(capabilities, text.fields) {
1651 pushed.text = Some(text.clone());
1652 outcomes.record_all(predicates, Outcome::PushedDown);
1653 } else {
1654 local.text = Some(text.clone());
1655 outcomes.record_all(predicates, Outcome::AppliedLocally);
1656 }
1657 }
1658
1659 ProjectShape {
1660 pushed,
1661 local,
1662 outcomes,
1663 }
1664}
1665
1666fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
1668 match stream {
1669 StreamKind::Projects => {
1670 let shaped = shape_projects(capabilities, filters);
1671 HitShape {
1672 stream,
1673 tasks: TaskQuery::default(),
1674 projects: shaped.pushed,
1675 local_tasks: LocalTasks::default(),
1676 local_projects: shaped.local,
1677 outcomes: shaped.outcomes,
1678 }
1679 }
1680 StreamKind::Items | StreamKind::Tasks => {
1681 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any);
1682 HitShape {
1683 stream,
1684 tasks: shaped.pushed,
1685 projects: ProjectQuery::default(),
1686 local_tasks: shaped.local,
1687 local_projects: LocalProjects::default(),
1688 outcomes: shaped.outcomes,
1689 }
1690 }
1691 }
1692}
1693
1694fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
1701 if compensating {
1702 ceiling
1703 } else {
1704 budget.min(ceiling)
1705 }
1706}
1707
1708fn ceiling(source: &ResolvedSource) -> u32 {
1710 source.source().capabilities().max_page_size.max(1)
1711}
1712
1713async fn fetch_tasks(
1715 source: &ResolvedSource,
1716 shape: &TaskShape,
1717 start: &Resume,
1718 budget: u32,
1719 calls: &AtomicU32,
1720) -> Result<Fetched<Task>, SourceError> {
1721 let compensating = shape.local != LocalTasks::default();
1722 walk(
1723 start,
1724 budget,
1725 page_size(compensating, budget, ceiling(source)),
1726 |task| shape.local.keeps(task),
1727 |cursor, limit| async move {
1728 calls.fetch_add(1, Ordering::Relaxed);
1729 let request = PageRequest { cursor, limit };
1730 source.source().query_tasks(&shape.pushed, &request).await
1731 },
1732 )
1733 .await
1734}
1735
1736async fn fetch_projects(
1738 source: &ResolvedSource,
1739 shape: &ProjectShape,
1740 start: &Resume,
1741 budget: u32,
1742 calls: &AtomicU32,
1743) -> Result<Fetched<Project>, SourceError> {
1744 let compensating = shape.local != LocalProjects::default();
1745 walk(
1746 start,
1747 budget,
1748 page_size(compensating, budget, ceiling(source)),
1749 |project| shape.local.keeps(project),
1750 |cursor, limit| async move {
1751 calls.fetch_add(1, Ordering::Relaxed);
1752 let request = PageRequest { cursor, limit };
1753 source
1754 .source()
1755 .query_projects(&shape.pushed, &request)
1756 .await
1757 },
1758 )
1759 .await
1760}
1761
1762async fn fetch_labels(
1764 source: &ResolvedSource,
1765 start: &Resume,
1766 budget: u32,
1767 calls: &AtomicU32,
1768) -> Result<Fetched<Label>, SourceError> {
1769 walk(
1770 start,
1771 budget,
1772 page_size(false, budget, ceiling(source)),
1773 |_| true,
1774 |cursor, limit| async move {
1775 calls.fetch_add(1, Ordering::Relaxed);
1776 let request = PageRequest { cursor, limit };
1777 source.source().labels(&request).await
1778 },
1779 )
1780 .await
1781}
1782
1783async fn fetch_hits(
1785 source: &ResolvedSource,
1786 shape: &HitShape,
1787 start: &Resume,
1788 budget: u32,
1789 calls: &AtomicU32,
1790) -> Result<Fetched<Found>, SourceError> {
1791 let ceiling = ceiling(source);
1792 match shape.stream {
1793 StreamKind::Projects => {
1794 let compensating = shape.local_projects != LocalProjects::default();
1795 walk(
1796 start,
1797 budget,
1798 page_size(compensating, budget, ceiling),
1799 |found| match found {
1800 Found::Project(project) => shape.local_projects.keeps(project),
1801 Found::Task(_) => true,
1802 },
1803 |cursor, limit| async move {
1804 calls.fetch_add(1, Ordering::Relaxed);
1805 let request = PageRequest { cursor, limit };
1806 let page = source
1807 .source()
1808 .query_projects(&shape.projects, &request)
1809 .await?;
1810 Ok(Page {
1811 items: page.items.into_iter().map(Found::Project).collect(),
1812 next: page.next,
1813 })
1814 },
1815 )
1816 .await
1817 }
1818 StreamKind::Items | StreamKind::Tasks => {
1819 let compensating = shape.local_tasks != LocalTasks::default();
1820 walk(
1821 start,
1822 budget,
1823 page_size(compensating, budget, ceiling),
1824 |found| match found {
1825 Found::Task(task) => shape.local_tasks.keeps(task),
1826 Found::Project(_) => true,
1827 },
1828 |cursor, limit| async move {
1829 calls.fetch_add(1, Ordering::Relaxed);
1830 let request = PageRequest { cursor, limit };
1831 let page = source.source().query_tasks(&shape.tasks, &request).await?;
1832 Ok(Page {
1833 items: page.items.into_iter().map(Found::Task).collect(),
1834 next: page.next,
1835 })
1836 },
1837 )
1838 .await
1839 }
1840 }
1841}
1842
1843async fn forward_edges(
1845 source: &ResolvedSource,
1846 entity: Entity,
1847 id: &NativeId,
1848 request: &PageRequest,
1849) -> Result<Page<DependencyEdge>, SourceError> {
1850 match entity {
1851 Entity::Task => {
1852 source
1853 .source()
1854 .task_dependencies(id, Direction::DependsOn, request)
1855 .await
1856 }
1857 Entity::Project => {
1858 source
1859 .source()
1860 .project_dependencies(id, Direction::DependsOn, request)
1861 .await
1862 }
1863 }
1864}
1865
1866#[expect(
1875 clippy::too_many_arguments,
1876 reason = "every argument is one axis of one walk — the source, the item, the \
1877 direction, which of its two graphs, whether the reverse is emulated, where \
1878 to resume, how many rows to return and where to count calls. Grouping them \
1879 into a struct would name the same eight values one indirection further from \
1880 the loop that reads them."
1881)]
1882async fn fetch_edges(
1883 source: &ResolvedSource,
1884 native: &NativeId,
1885 direction: Direction,
1886 entity: Entity,
1887 emulating: bool,
1888 start: &Resume,
1889 budget: u32,
1890 calls: &AtomicU32,
1891) -> Result<Fetched<DependencyEdge>, SourceError> {
1892 let ceiling = ceiling(source);
1893 if !emulating {
1894 return walk(
1895 start,
1896 budget,
1897 page_size(false, budget, ceiling),
1898 |_| true,
1899 |cursor, limit| async move {
1900 calls.fetch_add(1, Ordering::Relaxed);
1901 let request = PageRequest { cursor, limit };
1902 match entity {
1903 Entity::Task => {
1904 source
1905 .source()
1906 .task_dependencies(native, direction, &request)
1907 .await
1908 }
1909 Entity::Project => {
1910 source
1911 .source()
1912 .project_dependencies(native, direction, &request)
1913 .await
1914 }
1915 }
1916 },
1917 )
1918 .await;
1919 }
1920
1921 walk(
1922 start,
1923 budget,
1924 ceiling,
1925 |_| true,
1926 |cursor, limit| async move {
1927 calls.fetch_add(1, Ordering::Relaxed);
1928 let request = PageRequest { cursor, limit };
1929 let (ids, next) = match entity {
1930 Entity::Task => {
1931 let page = source
1932 .source()
1933 .query_tasks(&TaskQuery::default(), &request)
1934 .await?;
1935 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
1936 (ids, page.next)
1937 }
1938 Entity::Project => {
1939 let page = source
1940 .source()
1941 .query_projects(&ProjectQuery::default(), &request)
1942 .await?;
1943 let ids: Vec<NativeId> =
1944 page.items.into_iter().map(|project| project.id).collect();
1945 (ids, page.next)
1946 }
1947 };
1948
1949 let mut edges = Vec::new();
1950 for id in ids {
1951 let mut inner: Option<Cursor> = None;
1952 loop {
1953 calls.fetch_add(1, Ordering::Relaxed);
1954 let request = PageRequest {
1955 cursor: inner.clone(),
1956 limit,
1957 };
1958 let page = forward_edges(source, entity, &id, &request).await?;
1959 fits(page.items.len(), limit)?;
1963 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
1964 if page.next.is_some() && page.next == inner {
1965 return Err(SourceError::Malformed {
1966 message: "the source returned the cursor it was given while its \
1967 forward edges were being scanned, so the scan would \
1968 never end"
1969 .to_owned(),
1970 });
1971 }
1972 match page.next {
1973 Some(cursor) => inner = Some(cursor),
1974 None => break,
1975 }
1976 }
1977 }
1978
1979 Ok(Page { items: edges, next })
1980 },
1981 )
1982 .await
1983}