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, Document, DocumentQuery, Label, LabelFilter,
30 NativeId, Page, PageRequest, Project, ProjectFilter, ProjectQuery, SecretResolver, SourceError,
31 SourceName, 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, unrepeated, walk};
42use join::join_all;
43use local::{LocalDocuments, LocalProjects, LocalTasks};
44pub(crate) use resume::{Owed, Resumption, StreamState};
45use resume::{Resume, StreamKind};
46
47pub use copy::{
48 BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
49 Spent,
50};
51pub use local::ProjectSelector;
52
53#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
58pub struct Qualified<T> {
59 pub id: GlobalId,
61 pub item: T,
63}
64
65#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
70pub struct QualifiedEdge {
71 pub from: QualifiedEndpoint,
76 pub to: QualifiedEndpoint,
78 pub kind: onetaskgraph_plugin_api::DependencyKind,
80}
81
82#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
84pub struct QualifiedEndpoint {
85 pub id: GlobalId,
87 pub kind: onetaskgraph_plugin_api::ItemKind,
89}
90
91impl std::fmt::Display for QualifiedEndpoint {
92 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
93 self.id.fmt(formatter)
94 }
95}
96
97#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
99#[serde(tag = "kind", rename_all = "kebab-case")]
100pub enum SearchHit {
101 Task(Qualified<Task>),
103 Project(Qualified<Project>),
105}
106
107#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
109#[serde(rename_all = "kebab-case")]
110pub enum SearchKind {
111 Tasks,
113 Projects,
115 #[default]
117 Both,
118}
119
120#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
122pub struct SourceListing {
123 pub source: SourceName,
125 pub kind: String,
137 #[serde(flatten)]
139 pub state: SourceState,
140}
141
142#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
144#[serde(tag = "state", rename_all = "kebab-case")]
145pub enum SourceState {
146 Available {
148 capabilities: Capabilities,
150 },
151 Unavailable {
153 error: SourceError,
155 },
156}
157
158#[derive(Debug, Clone, PartialEq)]
160pub struct Paging {
161 pub limit: NonZeroU32,
163 pub token: Option<PageToken>,
165}
166
167#[derive(Debug, Clone, Default, PartialEq)]
169pub struct Filters {
170 pub text: Option<TextQuery>,
172 pub labels: LabelFilter,
174 pub statuses: Vec<StatusCategory>,
176}
177
178#[derive(Debug, Clone)]
180pub struct TaskRequest {
181 pub sources: Vec<SourceName>,
183 pub filters: Filters,
185 pub project: ProjectSelector,
187 pub paging: Paging,
189}
190
191#[derive(Debug, Clone)]
193pub struct ProjectRequest {
194 pub sources: Vec<SourceName>,
196 pub filters: Filters,
198 pub paging: Paging,
200}
201
202#[derive(Debug, Clone, Default, PartialEq)]
209pub struct DocumentFilters {
210 pub text: Option<TextQuery>,
212 pub labels: LabelFilter,
214}
215
216#[derive(Debug, Clone)]
218pub struct DocumentRequest {
219 pub sources: Vec<SourceName>,
221 pub filters: DocumentFilters,
223 pub project: ProjectSelector,
225 pub paging: Paging,
227}
228
229#[derive(Debug, Clone)]
231pub struct LabelRequest {
232 pub sources: Vec<SourceName>,
234 pub paging: Paging,
236}
237
238#[derive(Debug, Clone)]
240pub struct SearchRequest {
241 pub sources: Vec<SourceName>,
243 pub text: TextQuery,
245 pub kind: SearchKind,
247 pub paging: Paging,
249}
250
251#[derive(Debug, Clone)]
253pub struct DependencyRequest {
254 pub id: GlobalId,
256 pub direction: Direction,
258 pub paging: Paging,
260}
261
262#[derive(Debug, Clone, PartialEq, thiserror::Error)]
270pub enum EngineError {
271 #[error(
273 "no source named {name:?} is configured\n\
274 next: name one of the configured sources ({configured}), or add {name:?} under \
275 `sources` — `onetaskgraph sources list` shows what this configuration has."
276 )]
277 UnknownSource {
278 name: String,
280 configured: String,
282 },
283
284 #[error(
286 "{message}\n\
287 next: page with a token exactly as the previous page reported it, and against \
288 the same configuration — or drop `--page` to start the walk again."
289 )]
290 Token {
291 message: String,
293 },
294
295 #[error(
297 "no sources are configured\n\
298 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
299 prints what each plugin accepts."
300 )]
301 NoSources,
302
303 #[error(
305 "source {name} cannot be written: its plugin is {kind}, which has no write \
306 side\n\
307 next: copy into a source whose plugin can be written — `onetaskgraph sources \
308 list` reports each one's plugin."
309 )]
310 NotWritable {
311 name: String,
313 kind: String,
315 },
316
317 #[error(
324 "source {name} has no documents: its plugin is {kind}, which holds none\n\
325 next: name a source whose plugin has documents — `onetaskgraph sources list` \
326 reports each one's plugin and what it declares."
327 )]
328 NoDocuments {
329 name: String,
331 kind: String,
333 },
334
335 #[error(
337 "the destination source {name} could not be built: {error}\n\
338 next: fix that source — `onetaskgraph sources list` reports its state — then \
339 copy again."
340 )]
341 DestinationUnavailable {
342 name: String,
344 error: SourceError,
346 },
347
348 #[error(
350 "no item with the id {id}\n\
351 next: check the id, or list what is there — `onetaskgraph task list` and \
352 `onetaskgraph project list` report what the configured sources hold."
353 )]
354 NoSuchItem {
355 id: String,
357 },
358
359 #[error(
364 "{item} was copied from {origin}, which that destination no longer holds\n\
365 next: re-run with --recreate to create a new item there instead, or restore \
366 {origin}."
367 )]
368 StaleOrigin {
369 item: String,
371 origin: String,
373 },
374
375 #[error(
381 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
382 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
383 {id} on its own with `onetaskgraph task copy`.",
384 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
385 )]
386 NotAMember {
387 id: GlobalId,
389 projects: Vec<GlobalId>,
391 },
392
393 #[error(
402 "{item} depends on {member}, which this copy was not told to carry and which records \
403 no origin in {destination}\n\
404 next: name {member} with --member as well, record its {destination} id at \
405 onetaskgraph.origin, or copy the whole project without --member."
406 )]
407 UnrecordedMember {
408 item: GlobalId,
410 member: GlobalId,
412 destination: SourceName,
414 },
415
416 #[error(
422 "source {name} could not do it: {error}\n\
423 next: fix what the source named above, then copy again."
424 )]
425 SourceRefused {
426 name: String,
428 error: SourceError,
430 },
431
432 #[error(
440 "the copy failed and could not be undone.\n\
441 it failed because: {error}\n\
442 it could not be undone because: {refusal}\n\
443 so the destination still holds: {left_behind}\n\
444 next: remove those items at the destination, then copy again."
445 )]
446 CopyNotUndone {
447 error: Box<EngineError>,
449 left_behind: LeftBehind,
455 refusal: SourceError,
457 },
458}
459
460#[derive(Debug, Clone, PartialEq)]
469pub struct LeftBehind {
470 first: GlobalId,
472 rest: Vec<GlobalId>,
474}
475
476impl LeftBehind {
477 #[must_use]
479 pub fn new(first: GlobalId) -> Self {
480 Self {
481 first,
482 rest: Vec::new(),
483 }
484 }
485
486 pub fn push(&mut self, id: GlobalId) {
488 self.rest.push(id);
489 }
490
491 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
493 std::iter::once(&self.first).chain(self.rest.iter())
494 }
495}
496
497impl std::fmt::Display for LeftBehind {
498 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
500 write!(formatter, "{}", self.first)?;
501 for id in &self.rest {
502 write!(formatter, ", {id}")?;
503 }
504 Ok(())
505 }
506}
507
508pub enum ConfiguredSource {
515 Ready(ResolvedSource),
517 Unavailable(UnavailableSource),
519}
520
521impl ConfiguredSource {
522 #[must_use]
524 pub fn name(&self) -> &SourceName {
525 match self {
526 Self::Ready(source) => source.name(),
527 Self::Unavailable(source) => source.name(),
528 }
529 }
530}
531
532pub struct Engine {
534 sources: Vec<ConfiguredSource>,
536 selection: Vec<SourceName>,
538}
539
540impl Engine {
541 #[must_use]
549 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
550 let (ready, unavailable) = resolve_available(config, secrets);
551 Self::new(
552 ready
553 .into_iter()
554 .map(ConfiguredSource::Ready)
555 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
556 .collect(),
557 config.selected_sources(),
558 )
559 }
560
561 #[must_use]
564 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
565 Self { sources, selection }
566 }
567
568 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
570 self.sources.iter().filter_map(|source| match source {
571 ConfiguredSource::Ready(ready) => Some(ready),
572 ConfiguredSource::Unavailable(_) => None,
573 })
574 }
575
576 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
578 self.sources.iter().filter_map(|source| match source {
579 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
580 ConfiguredSource::Ready(_) => None,
581 })
582 }
583
584 #[must_use]
586 pub fn listing(&self) -> Vec<SourceListing> {
587 let mut listings: Vec<SourceListing> = self
588 .ready()
589 .map(|source| SourceListing {
590 source: source.name().clone(),
591 kind: source.kind().to_owned(),
592 state: SourceState::Available {
593 capabilities: source.source().capabilities(),
594 },
595 })
596 .chain(self.unavailable().map(|source| SourceListing {
597 source: source.name().clone(),
598 kind: source.kind().to_owned(),
599 state: SourceState::Unavailable {
600 error: source.error().clone(),
601 },
602 }))
603 .collect();
604 listings.sort_by(|left, right| left.source.cmp(&right.source));
605 listings
606 }
607
608 #[must_use]
614 pub fn has(&self, name: &SourceName) -> bool {
615 self.sources.iter().any(|source| source.name() == name)
616 }
617
618 pub async fn tasks(
626 &self,
627 request: &TaskRequest,
628 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
629 let mut names = self.resolve_selection(&request.sources)?;
630 if let ProjectSelector::Qualified(id) = &request.project {
635 self.known(&id.source)?;
636 names.retain(|name| name == &id.source);
637 }
638 let query = shape("task-list", &names, &(&request.filters, &request.project));
639 let states = resumption(
640 self,
641 request.paging.token.as_ref(),
642 &[StreamKind::Items],
643 &query,
644 )?;
645 let budget = request.paging.limit.get();
646
647 let mut answer = Answer::new();
648 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
649
650 let shapes: Vec<TaskShape> = ready
651 .iter()
652 .map(|source| {
653 shape_tasks(
654 &source.source().capabilities(),
655 &request.filters,
656 &project_filter(&request.project),
657 )
658 })
659 .collect();
660 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
661 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
662
663 let walks = ready
664 .iter()
665 .enumerate()
666 .map(|(index, source)| {
667 fetch_tasks(
668 source,
669 &shapes[index],
670 &starts[index],
671 budget,
672 &counters[index],
673 )
674 })
675 .collect();
676
677 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
678 answer.finish(
679 streams,
680 budget,
681 owed(&states),
682 &query,
683 |name, task: Task| Qualified {
684 id: GlobalId::new(name.clone(), task.id.clone()),
685 item: task,
686 },
687 )
688 }
689
690 pub async fn projects(
696 &self,
697 request: &ProjectRequest,
698 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
699 let names = self.resolve_selection(&request.sources)?;
700 let query = shape("project-list", &names, &request.filters);
701 let states = resumption(
702 self,
703 request.paging.token.as_ref(),
704 &[StreamKind::Items],
705 &query,
706 )?;
707 let budget = request.paging.limit.get();
708
709 let mut answer = Answer::new();
710 let mut with_projects = Vec::new();
716 for source in answer.split(self, &names) {
717 if source.source().capabilities().projects.is_native() {
718 with_projects.push(source);
719 } else {
720 answer.unreachable_predicate(source, Predicate::Project);
721 }
722 }
723 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
724
725 let shapes: Vec<ProjectShape> = ready
726 .iter()
727 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
728 .collect();
729 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
730 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
731
732 let walks = ready
733 .iter()
734 .enumerate()
735 .map(|(index, source)| {
736 fetch_projects(
737 source,
738 &shapes[index],
739 &starts[index],
740 budget,
741 &counters[index],
742 )
743 })
744 .collect();
745
746 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
747 answer.finish(
748 streams,
749 budget,
750 owed(&states),
751 &query,
752 |name, project: Project| Qualified {
753 id: GlobalId::new(name.clone(), project.id.clone()),
754 item: project,
755 },
756 )
757 }
758
759 pub async fn documents(
771 &self,
772 request: &DocumentRequest,
773 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
774 let mut names = self.resolve_selection(&request.sources)?;
775 if let ProjectSelector::Qualified(id) = &request.project {
778 self.known(&id.source)?;
779 names.retain(|name| name == &id.source);
780 }
781 let query = shape(
782 "document-list",
783 &names,
784 &(&request.filters, &request.project),
785 );
786 let states = resumption(
787 self,
788 request.paging.token.as_ref(),
789 &[StreamKind::Items],
790 &query,
791 )?;
792 let budget = request.paging.limit.get();
793
794 let mut answer = Answer::new();
795 let mut with_documents = Vec::new();
796 for source in answer.split(self, &names) {
797 if source.source().capabilities().documents.is_native() {
798 with_documents.push(source);
799 } else {
800 answer.unreachable_predicate(source, Predicate::Document);
801 }
802 }
803 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
804
805 let shapes: Vec<DocumentShape> = ready
806 .iter()
807 .map(|source| {
808 shape_documents(
809 &source.source().capabilities(),
810 &request.filters,
811 &project_filter(&request.project),
812 )
813 })
814 .collect();
815 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
816 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
817
818 let walks = ready
819 .iter()
820 .enumerate()
821 .map(|(index, source)| {
822 fetch_documents(
823 source,
824 &shapes[index],
825 &starts[index],
826 budget,
827 &counters[index],
828 )
829 })
830 .collect();
831
832 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
833 answer.finish(
834 streams,
835 budget,
836 owed(&states),
837 &query,
838 |name, document: Document| Qualified {
839 id: GlobalId::new(name.clone(), document.id.clone()),
840 item: document,
841 },
842 )
843 }
844
845 pub async fn labels(
851 &self,
852 request: &LabelRequest,
853 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
854 let names = self.resolve_selection(&request.sources)?;
855 let query = shape("label-list", &names, &());
856 let states = resumption(
857 self,
858 request.paging.token.as_ref(),
859 &[StreamKind::Items],
860 &query,
861 )?;
862 let budget = request.paging.limit.get();
863
864 let mut answer = Answer::new();
865 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
866 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
867 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
868
869 let walks = ready
870 .iter()
871 .enumerate()
872 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
873 .collect();
874
875 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
876 answer.finish(
877 streams,
878 budget,
879 owed(&states),
880 &query,
881 |name, label: Label| Qualified {
882 id: GlobalId::new(name.clone(), label.id.clone()),
883 item: label,
884 },
885 )
886 }
887
888 pub async fn search(
894 &self,
895 request: &SearchRequest,
896 ) -> Result<QueryResponse<SearchHit>, EngineError> {
897 let names = self.resolve_selection(&request.sources)?;
898 let reads: &[StreamKind] = match request.kind {
902 SearchKind::Tasks => &[StreamKind::Tasks],
903 SearchKind::Projects => &[StreamKind::Projects],
904 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
905 };
906 let query = shape("search", &names, &request.text);
911 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
912 let budget = request.paging.limit.get();
913 let filters = Filters {
914 text: Some(request.text.clone()),
915 ..Filters::default()
916 };
917
918 let mut answer = Answer::new();
919
920 let mut ready = Vec::new();
923 let mut kinds = Vec::new();
924 let mut starts = Vec::new();
925 for source in answer.split(self, &names) {
926 let mut streams = Vec::new();
927 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
928 streams.push(StreamKind::Tasks);
929 }
930 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
931 if source.source().capabilities().projects.is_native() {
932 streams.push(StreamKind::Projects);
933 } else {
934 answer.unreachable_predicate(source, Predicate::Project);
935 }
936 }
937 for stream in streams {
938 if let Some(resume) = resume_at(&states, source.name(), stream) {
939 ready.push(source);
940 kinds.push(stream);
941 starts.push(resume);
942 }
943 }
944 }
945
946 let shapes: Vec<HitShape> = ready
947 .iter()
948 .zip(kinds.iter())
949 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
950 .collect();
951 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
952 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
953
954 let walks = ready
955 .iter()
956 .enumerate()
957 .map(|(index, source)| {
958 fetch_hits(
959 source,
960 &shapes[index],
961 &starts[index],
962 budget,
963 &counters[index],
964 )
965 })
966 .collect();
967
968 let streams =
969 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
970 answer.finish(
971 streams,
972 budget,
973 owed(&states),
974 &query,
975 |name, found: Found| match found {
976 Found::Task(task) => SearchHit::Task(Qualified {
977 id: GlobalId::new(name.clone(), task.id.clone()),
978 item: task,
979 }),
980 Found::Project(project) => SearchHit::Project(Qualified {
981 id: GlobalId::new(name.clone(), project.id.clone()),
982 item: project,
983 }),
984 },
985 )
986 }
987
988 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
995 let name = self.known(&id.source)?;
996 let mut answer = Answer::new();
997 let selected = answer.split(self, std::slice::from_ref(&name));
998 let Some(source) = selected.first() else {
999 return answer.nothing();
1000 };
1001 let found = source.source().get_task(&id.native).await;
1002 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1003 answer.one(source, found, |task| Qualified {
1004 id: qualified,
1005 item: task,
1006 })
1007 }
1008
1009 pub async fn project(
1015 &self,
1016 id: &GlobalId,
1017 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1018 let name = self.known(&id.source)?;
1019 let mut answer = Answer::new();
1020 let selected = answer.split(self, std::slice::from_ref(&name));
1021 let Some(source) = selected.first() else {
1022 return answer.nothing();
1023 };
1024 let found = source.source().get_project(&id.native).await;
1025 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1026 answer.one(source, found, |project| Qualified {
1027 id: qualified,
1028 item: project,
1029 })
1030 }
1031
1032 pub async fn document(
1042 &self,
1043 id: &GlobalId,
1044 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1045 let name = self.known(&id.source)?;
1046 let mut answer = Answer::new();
1047 let selected = answer.split(self, std::slice::from_ref(&name));
1048 let Some(source) = selected.first() else {
1049 return answer.nothing();
1050 };
1051 if !source.source().capabilities().documents.is_native() {
1052 answer.unreachable_predicate(source, Predicate::Document);
1053 return answer.nothing();
1054 }
1055 let found = source.source().get_document(&id.native).await;
1056 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1057 answer.one(source, found, |document| Qualified {
1058 id: qualified,
1059 item: document,
1060 })
1061 }
1062
1063 pub async fn task_dependencies(
1070 &self,
1071 request: &DependencyRequest,
1072 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1073 self.dependencies(request, Entity::Task).await
1074 }
1075
1076 pub async fn project_dependencies(
1082 &self,
1083 request: &DependencyRequest,
1084 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1085 self.dependencies(request, Entity::Project).await
1086 }
1087
1088 async fn dependencies(
1091 &self,
1092 request: &DependencyRequest,
1093 entity: Entity,
1094 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1095 let name = self.known(&request.id.source)?;
1096 let query = shape(
1097 "dependencies",
1098 std::slice::from_ref(&name),
1099 &(entity, &request.id.native, request.direction),
1100 );
1101 let states = resumption(
1102 self,
1103 request.paging.token.as_ref(),
1104 &[StreamKind::Items],
1105 &query,
1106 )?;
1107 let budget = request.paging.limit.get();
1108
1109 let mut answer = Answer::new();
1110 let (ready, starts) = walking(
1111 answer.split(self, std::slice::from_ref(&name)),
1112 &states,
1113 StreamKind::Items,
1114 );
1115 let Some(source) = ready.first() else {
1116 return answer.nothing();
1117 };
1118
1119 let capabilities = source.source().capabilities();
1120 let support = match entity {
1121 Entity::Task => capabilities.task_dependencies,
1122 Entity::Project => capabilities.project_dependencies,
1123 };
1124 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1128 let mut outcomes = Outcomes::default();
1129 if request.direction == Direction::DependedOnBy {
1130 if emulating {
1131 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1132 } else {
1133 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1134 }
1135 }
1136
1137 let counters = vec![AtomicU32::new(0)];
1138 let walked = fetch_edges(
1139 source,
1140 &request.id.native,
1141 request.direction,
1142 entity,
1143 emulating,
1144 &starts[0],
1145 budget,
1146 &counters[0],
1147 )
1148 .await;
1149
1150 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1151 answer.finish(
1152 streams,
1153 budget,
1154 owed(&states),
1155 &query,
1156 |name, edge: DependencyEdge| QualifiedEdge {
1157 from: qualify_endpoint(name, edge.from),
1158 to: qualify_endpoint(name, edge.to),
1159 kind: edge.kind,
1160 },
1161 )
1162 }
1163
1164 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1166 if asked.is_empty() {
1167 if self.selection.is_empty() {
1168 return Err(EngineError::NoSources);
1169 }
1170 return Ok(self.selection.clone());
1171 }
1172 asked.iter().map(|name| self.known(name)).collect()
1173 }
1174
1175 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1177 if self.has(name) {
1178 return Ok(name.clone());
1179 }
1180 if self.sources.is_empty() {
1181 return Err(EngineError::NoSources);
1182 }
1183 Err(EngineError::UnknownSource {
1184 name: name.to_string(),
1185 configured: self
1186 .listing()
1187 .iter()
1188 .map(|listing| listing.source.to_string())
1189 .collect::<Vec<_>>()
1190 .join(", "),
1191 })
1192 }
1193}
1194
1195fn qualify_endpoint(
1196 source: &SourceName,
1197 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1198) -> QualifiedEndpoint {
1199 let kind = endpoint.kind;
1200 let is_qualified = endpoint.is_qualified();
1201 let endpoint_id = endpoint.into_id();
1202 QualifiedEndpoint {
1203 id: if is_qualified {
1204 endpoint_id
1205 .parse()
1206 .expect("plugin-api validates qualified dependency endpoints")
1207 } else {
1208 GlobalId::new(source.clone(), NativeId(endpoint_id))
1209 },
1210 kind,
1211 }
1212}
1213
1214#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1216enum Entity {
1217 Task,
1219 Project,
1221}
1222
1223enum Found {
1225 Task(Task),
1227 Project(Project),
1229}
1230
1231#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1233enum Outcome {
1234 PushedDown,
1236 AppliedLocally,
1238 Emulated,
1240 Unavailable,
1242}
1243
1244#[derive(Debug, Clone, Default, PartialEq)]
1256struct Outcomes(BTreeMap<Predicate, Outcome>);
1257
1258impl Outcomes {
1259 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1265 self.0.insert(predicate, outcome);
1266 }
1267
1268 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1271 for predicate in predicates {
1272 self.record(predicate, outcome);
1273 }
1274 }
1275
1276 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1278 self.0
1279 .iter()
1280 .filter(|(_, recorded)| **recorded == outcome)
1281 .map(|(predicate, _)| *predicate)
1282 .collect()
1283 }
1284}
1285
1286struct TaskShape {
1288 pushed: TaskQuery,
1290 local: LocalTasks,
1292 outcomes: Outcomes,
1294}
1295
1296struct ProjectShape {
1298 pushed: ProjectQuery,
1300 local: LocalProjects,
1302 outcomes: Outcomes,
1304}
1305
1306struct DocumentShape {
1308 pushed: DocumentQuery,
1310 local: LocalDocuments,
1312 outcomes: Outcomes,
1314}
1315
1316struct HitShape {
1318 stream: StreamKind,
1320 tasks: TaskQuery,
1322 projects: ProjectQuery,
1324 local_tasks: LocalTasks,
1326 local_projects: LocalProjects,
1328 outcomes: Outcomes,
1330}
1331
1332struct Answer {
1338 plans: Vec<SourcePlan>,
1340 errors: Vec<SourceFailure>,
1342}
1343
1344impl Answer {
1345 fn new() -> Self {
1346 Self {
1347 plans: Vec::new(),
1348 errors: Vec::new(),
1349 }
1350 }
1351
1352 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1357 let mut selected = Vec::new();
1358 for name in names {
1359 match engine.sources.iter().find(|source| source.name() == name) {
1360 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1361 Some(ConfiguredSource::Unavailable(source)) => {
1362 self.errors.push(source.failure());
1363 }
1364 None => {}
1365 }
1366 }
1367 selected
1368 }
1369
1370 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1372 let mut outcomes = Outcomes::default();
1373 outcomes.record(predicate, Outcome::Unavailable);
1374 self.plans.push(plan_for(source, outcomes, 0));
1375 }
1376
1377 fn collect<T>(
1379 &mut self,
1380 ready: &[&ResolvedSource],
1381 walked: Vec<Result<Fetched<T>, SourceError>>,
1382 counters: &[AtomicU32],
1383 outcomes: Vec<Outcomes>,
1384 ) -> Vec<Stream<T>> {
1385 let kinds = vec![StreamKind::Items; ready.len()];
1386 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1387 }
1388
1389 fn collect_streams<T>(
1391 &mut self,
1392 ready: &[&ResolvedSource],
1393 kinds: &[StreamKind],
1394 walked: Vec<Result<Fetched<T>, SourceError>>,
1395 counters: &[AtomicU32],
1396 outcomes: Vec<Outcomes>,
1397 ) -> Vec<Stream<T>> {
1398 let mut streams = Vec::new();
1399 for (index, result) in walked.into_iter().enumerate() {
1400 let source = ready[index];
1401 let pages = counters[index].load(Ordering::Relaxed);
1402 self.plans
1403 .push(plan_for(source, outcomes[index].clone(), pages));
1404 match result {
1405 Ok(fetched) => streams.push(Stream {
1406 source: source.name().clone(),
1407 kind: kinds[index],
1408 fetched,
1409 }),
1410 Err(error) => self.errors.push(SourceFailure {
1413 source: source.name().clone(),
1414 error,
1415 }),
1416 }
1417 }
1418 streams
1419 }
1420
1421 fn one<T, U>(
1423 mut self,
1424 source: &ResolvedSource,
1425 found: Result<Option<T>, SourceError>,
1426 qualify: impl FnOnce(T) -> U,
1427 ) -> Result<QueryResponse<U>, EngineError> {
1428 self.plans.push(plan_for(source, Outcomes::default(), 1));
1429 let items = match found {
1430 Ok(Some(item)) => vec![qualify(item)],
1431 Ok(None) => Vec::new(),
1432 Err(error) => {
1433 self.errors.push(SourceFailure {
1434 source: source.name().clone(),
1435 error,
1436 });
1437 Vec::new()
1438 }
1439 };
1440 Ok(QueryResponse {
1441 items,
1442 next: None,
1443 plan: QueryPlan {
1444 per_source: merge_plans(self.plans),
1445 },
1446 errors: self.errors,
1447 })
1448 }
1449
1450 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1452 Ok(QueryResponse {
1453 items: Vec::new(),
1454 next: None,
1455 plan: QueryPlan {
1456 per_source: merge_plans(self.plans),
1457 },
1458 errors: self.errors,
1459 })
1460 }
1461
1462 fn finish<T, U>(
1467 self,
1468 streams: Vec<Stream<T>>,
1469 budget: u32,
1470 first: Option<&Owed>,
1471 query: &str,
1472 qualify: impl Fn(&SourceName, T) -> U,
1473 ) -> Result<QueryResponse<U>, EngineError> {
1474 let (rows, states, owed) = merge(streams, budget, first);
1475 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1476 Ok(QueryResponse {
1477 items: rows
1478 .into_iter()
1479 .map(|(name, item)| qualify(&name, item))
1480 .collect(),
1481 next,
1482 plan: QueryPlan {
1483 per_source: merge_plans(self.plans),
1484 },
1485 errors: self.errors,
1486 })
1487 }
1488}
1489
1490fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1493 SourcePlan {
1494 source: source.name().clone(),
1495 kind: source.kind().to_owned(),
1496 pushed_down: outcomes.with(Outcome::PushedDown),
1497 applied_locally: outcomes.with(Outcome::AppliedLocally),
1498 emulated: outcomes.with(Outcome::Emulated),
1499 unavailable: outcomes.with(Outcome::Unavailable),
1500 pages_fetched: pages,
1501 }
1502}
1503
1504fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1509 let mut merged: Vec<SourcePlan> = Vec::new();
1510 for plan in plans {
1511 if let Some(existing) = merged
1512 .iter_mut()
1513 .find(|existing| existing.source == plan.source)
1514 {
1515 existing.pushed_down.extend(plan.pushed_down);
1516 existing.applied_locally.extend(plan.applied_locally);
1517 existing.emulated.extend(plan.emulated);
1518 existing.unavailable.extend(plan.unavailable);
1519 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1520 for list in [
1521 &mut existing.pushed_down,
1522 &mut existing.applied_locally,
1523 &mut existing.emulated,
1524 &mut existing.unavailable,
1525 ] {
1526 list.sort_unstable();
1527 list.dedup();
1528 }
1529 } else {
1530 merged.push(plan);
1531 }
1532 }
1533 merged
1534}
1535
1536fn walking<'a>(
1542 selected: Vec<&'a ResolvedSource>,
1543 states: &Option<Resumption>,
1544 kind: StreamKind,
1545) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1546 let mut ready = Vec::new();
1547 let mut starts = Vec::new();
1548 for source in selected {
1549 if let Some(resume) = resume_at(states, source.name(), kind) {
1550 ready.push(source);
1551 starts.push(resume);
1552 }
1553 }
1554 (ready, starts)
1555}
1556
1557fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1575 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1576 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1577}
1578
1579fn fingerprint(text: &str) -> String {
1581 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1582 for byte in text.as_bytes() {
1583 hash ^= u64::from(*byte);
1584 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1585 }
1586 format!("{hash:016x}")
1587}
1588
1589fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1594 document.as_ref()?.owed.as_ref()
1595}
1596
1597fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1599 match states {
1600 None => Some(Resume::default()),
1601 Some(document) => document
1602 .streams
1603 .iter()
1604 .find(|state| &state.source == source && state.stream == kind)
1605 .map(|state| state.resume.clone()),
1606 }
1607}
1608
1609fn resumption(
1640 engine: &Engine,
1641 token: Option<&PageToken>,
1642 reads: &[StreamKind],
1643 query: &str,
1644) -> Result<Option<Resumption>, EngineError> {
1645 let Some(document) = token.map(PageToken::decode) else {
1646 return Ok(None);
1647 };
1648
1649 if document.query != query {
1656 return Err(EngineError::Token {
1657 message: "this page token was written by a different query — resume the walk it \
1658 came from, or drop --page to start this one from the beginning"
1659 .to_owned(),
1660 });
1661 }
1662 let states = &document.streams;
1663
1664 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
1665 for state in states {
1666 if !reads.contains(&state.stream) {
1667 return Err(EngineError::Token {
1668 message: format!(
1669 "this page token resumes {}, which this command does not read — it \
1670 was written by a different query",
1671 state.stream.describe()
1672 ),
1673 });
1674 }
1675 let ceiling = engine
1676 .ready()
1677 .find(|source| source.name() == &state.source)
1678 .map(ceiling);
1679 if ceiling.is_none() && !engine.has(&state.source) {
1680 return Err(EngineError::Token {
1681 message: format!(
1682 "this page token resumes a source called {:?}, which this \
1683 configuration does not have",
1684 state.source.as_str()
1685 ),
1686 });
1687 }
1688 if let Some(ceiling) = ceiling
1689 && state.resume.skip >= ceiling
1690 {
1691 return Err(EngineError::Token {
1692 message: format!(
1693 "this page token resumes {} rows into a page of source {:?}, which \
1694 serves at most {ceiling}",
1695 state.resume.skip,
1696 state.source.as_str()
1697 ),
1698 });
1699 }
1700 if seen.contains(&(&state.source, state.stream)) {
1701 return Err(EngineError::Token {
1702 message: format!(
1703 "this page token gives source {:?} two places to resume from",
1704 state.source.as_str()
1705 ),
1706 });
1707 }
1708 seen.push((&state.source, state.stream));
1709 }
1710
1711 if let Some(owed) = &document.owed
1716 && !document
1717 .streams
1718 .iter()
1719 .any(|state| state.source == owed.source && state.stream == owed.stream)
1720 {
1721 return Err(EngineError::Token {
1722 message: format!(
1723 "this page token owes the next row to a stream it does not resume, \
1724 {:?}'s {}",
1725 owed.source.as_str(),
1726 owed.stream.describe()
1727 ),
1728 });
1729 }
1730
1731 Ok(Some(document))
1732}
1733
1734fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
1740 match selector {
1741 ProjectSelector::Any => ProjectFilter::Any,
1742 ProjectSelector::Orphans => ProjectFilter::Orphans,
1743 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
1744 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
1745 }
1746}
1747
1748fn text_predicates(fields: TextFields) -> Vec<Predicate> {
1750 match fields {
1751 TextFields::Title => vec![Predicate::SearchTitle],
1752 TextFields::Content => vec![Predicate::SearchContent],
1753 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
1754 }
1755}
1756
1757fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
1764 match fields {
1765 TextFields::Title => capabilities.search_title.is_native(),
1766 TextFields::Content => capabilities.search_content.is_native(),
1767 TextFields::TitleOrContent => {
1768 capabilities.search_title.is_native() && capabilities.search_content.is_native()
1769 }
1770 }
1771}
1772
1773fn shape_tasks(
1775 capabilities: &Capabilities,
1776 filters: &Filters,
1777 project: &ProjectFilter,
1778) -> TaskShape {
1779 let mut pushed = TaskQuery::default();
1780 let mut local = LocalTasks::default();
1781 let mut outcomes = Outcomes::default();
1782
1783 if !filters.labels.is_empty() {
1784 if capabilities.filter_by_label.is_native() {
1785 pushed.labels = filters.labels.clone();
1786 outcomes.record(Predicate::Label, Outcome::PushedDown);
1787 } else {
1788 local.labels = Some(filters.labels.clone());
1789 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1790 }
1791 }
1792 if !filters.statuses.is_empty() {
1793 if capabilities.filter_by_status.is_native() {
1794 pushed.statuses.clone_from(&filters.statuses);
1795 outcomes.record(Predicate::Status, Outcome::PushedDown);
1796 } else {
1797 local.statuses.clone_from(&filters.statuses);
1798 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1799 }
1800 }
1801 if let Some(text) = &filters.text {
1802 let predicates = text_predicates(text.fields);
1803 if searches_natively(capabilities, text.fields) {
1804 pushed.text = Some(text.clone());
1805 outcomes.record_all(predicates, Outcome::PushedDown);
1806 } else {
1807 local.text = Some(text.clone());
1808 outcomes.record_all(predicates, Outcome::AppliedLocally);
1809 }
1810 }
1811 match project {
1812 ProjectFilter::Any => {}
1813 ProjectFilter::Orphans => {
1814 if capabilities.orphan_tasks.is_native() {
1815 pushed.project = ProjectFilter::Orphans;
1816 outcomes.record(Predicate::Project, Outcome::PushedDown);
1817 } else {
1818 local.project = Some(ProjectFilter::Orphans);
1819 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1820 }
1821 }
1822 ProjectFilter::Is(id) => {
1823 if capabilities.projects.is_native() {
1824 pushed.project = ProjectFilter::Is(id.clone());
1825 outcomes.record(Predicate::Project, Outcome::PushedDown);
1826 } else {
1827 local.project = Some(ProjectFilter::Is(id.clone()));
1828 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1829 }
1830 }
1831 }
1832
1833 TaskShape {
1834 pushed,
1835 local,
1836 outcomes,
1837 }
1838}
1839
1840fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
1842 let mut pushed = ProjectQuery::default();
1843 let mut local = LocalProjects::default();
1844 let mut outcomes = Outcomes::default();
1845
1846 if !filters.labels.is_empty() {
1847 if capabilities.filter_by_label.is_native() {
1848 pushed.labels = filters.labels.clone();
1849 outcomes.record(Predicate::Label, Outcome::PushedDown);
1850 } else {
1851 local.labels = Some(filters.labels.clone());
1852 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1853 }
1854 }
1855 if !filters.statuses.is_empty() {
1856 if capabilities.filter_by_status.is_native() {
1857 pushed.statuses.clone_from(&filters.statuses);
1858 outcomes.record(Predicate::Status, Outcome::PushedDown);
1859 } else {
1860 local.statuses.clone_from(&filters.statuses);
1861 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1862 }
1863 }
1864 if let Some(text) = &filters.text {
1865 let predicates = text_predicates(text.fields);
1866 if searches_natively(capabilities, text.fields) {
1867 pushed.text = Some(text.clone());
1868 outcomes.record_all(predicates, Outcome::PushedDown);
1869 } else {
1870 local.text = Some(text.clone());
1871 outcomes.record_all(predicates, Outcome::AppliedLocally);
1872 }
1873 }
1874
1875 ProjectShape {
1876 pushed,
1877 local,
1878 outcomes,
1879 }
1880}
1881
1882fn shape_documents(
1887 capabilities: &Capabilities,
1888 filters: &DocumentFilters,
1889 project: &ProjectFilter,
1890) -> DocumentShape {
1891 let mut pushed = DocumentQuery::default();
1892 let mut local = LocalDocuments::default();
1893 let mut outcomes = Outcomes::default();
1894
1895 if !filters.labels.is_empty() {
1896 if capabilities.filter_by_label.is_native() {
1897 pushed.labels = filters.labels.clone();
1898 outcomes.record(Predicate::Label, Outcome::PushedDown);
1899 } else {
1900 local.labels = Some(filters.labels.clone());
1901 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1902 }
1903 }
1904 if let Some(text) = &filters.text {
1905 let predicates = text_predicates(text.fields);
1906 if searches_natively(capabilities, text.fields) {
1907 pushed.text = Some(text.clone());
1908 outcomes.record_all(predicates, Outcome::PushedDown);
1909 } else {
1910 local.text = Some(text.clone());
1911 outcomes.record_all(predicates, Outcome::AppliedLocally);
1912 }
1913 }
1914 match project {
1915 ProjectFilter::Any => {}
1916 ProjectFilter::Orphans => {
1917 if capabilities.orphan_tasks.is_native() {
1918 pushed.project = ProjectFilter::Orphans;
1919 outcomes.record(Predicate::Project, Outcome::PushedDown);
1920 } else {
1921 local.project = Some(ProjectFilter::Orphans);
1922 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1923 }
1924 }
1925 ProjectFilter::Is(id) => {
1926 if capabilities.projects.is_native() {
1927 pushed.project = ProjectFilter::Is(id.clone());
1928 outcomes.record(Predicate::Project, Outcome::PushedDown);
1929 } else {
1930 local.project = Some(ProjectFilter::Is(id.clone()));
1931 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1932 }
1933 }
1934 }
1935
1936 DocumentShape {
1937 pushed,
1938 local,
1939 outcomes,
1940 }
1941}
1942
1943fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
1945 match stream {
1946 StreamKind::Projects => {
1947 let shaped = shape_projects(capabilities, filters);
1948 HitShape {
1949 stream,
1950 tasks: TaskQuery::default(),
1951 projects: shaped.pushed,
1952 local_tasks: LocalTasks::default(),
1953 local_projects: shaped.local,
1954 outcomes: shaped.outcomes,
1955 }
1956 }
1957 StreamKind::Items | StreamKind::Tasks => {
1958 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any);
1959 HitShape {
1960 stream,
1961 tasks: shaped.pushed,
1962 projects: ProjectQuery::default(),
1963 local_tasks: shaped.local,
1964 local_projects: LocalProjects::default(),
1965 outcomes: shaped.outcomes,
1966 }
1967 }
1968 }
1969}
1970
1971fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
1978 if compensating {
1979 ceiling
1980 } else {
1981 budget.min(ceiling)
1982 }
1983}
1984
1985fn ceiling(source: &ResolvedSource) -> u32 {
1987 source.source().capabilities().max_page_size.max(1)
1988}
1989
1990async fn fetch_tasks(
1992 source: &ResolvedSource,
1993 shape: &TaskShape,
1994 start: &Resume,
1995 budget: u32,
1996 calls: &AtomicU32,
1997) -> Result<Fetched<Task>, SourceError> {
1998 let compensating = shape.local != LocalTasks::default();
1999 walk(
2000 start,
2001 budget,
2002 page_size(compensating, budget, ceiling(source)),
2003 |task| shape.local.keeps(task),
2004 |cursor, limit| async move {
2005 calls.fetch_add(1, Ordering::Relaxed);
2006 let request = PageRequest { cursor, limit };
2007 source.source().query_tasks(&shape.pushed, &request).await
2008 },
2009 )
2010 .await
2011}
2012
2013async fn fetch_projects(
2015 source: &ResolvedSource,
2016 shape: &ProjectShape,
2017 start: &Resume,
2018 budget: u32,
2019 calls: &AtomicU32,
2020) -> Result<Fetched<Project>, SourceError> {
2021 let compensating = shape.local != LocalProjects::default();
2022 walk(
2023 start,
2024 budget,
2025 page_size(compensating, budget, ceiling(source)),
2026 |project| shape.local.keeps(project),
2027 |cursor, limit| async move {
2028 calls.fetch_add(1, Ordering::Relaxed);
2029 let request = PageRequest { cursor, limit };
2030 source
2031 .source()
2032 .query_projects(&shape.pushed, &request)
2033 .await
2034 },
2035 )
2036 .await
2037}
2038
2039async fn fetch_documents(
2041 source: &ResolvedSource,
2042 shape: &DocumentShape,
2043 start: &Resume,
2044 budget: u32,
2045 calls: &AtomicU32,
2046) -> Result<Fetched<Document>, SourceError> {
2047 let compensating = shape.local != LocalDocuments::default();
2048 walk(
2049 start,
2050 budget,
2051 page_size(compensating, budget, ceiling(source)),
2052 |document| shape.local.keeps(document),
2053 |cursor, limit| async move {
2054 calls.fetch_add(1, Ordering::Relaxed);
2055 let request = PageRequest { cursor, limit };
2056 source
2057 .source()
2058 .query_documents(&shape.pushed, &request)
2059 .await
2060 },
2061 )
2062 .await
2063}
2064
2065async fn fetch_labels(
2067 source: &ResolvedSource,
2068 start: &Resume,
2069 budget: u32,
2070 calls: &AtomicU32,
2071) -> Result<Fetched<Label>, SourceError> {
2072 walk(
2073 start,
2074 budget,
2075 page_size(false, budget, ceiling(source)),
2076 |_| true,
2077 |cursor, limit| async move {
2078 calls.fetch_add(1, Ordering::Relaxed);
2079 let request = PageRequest { cursor, limit };
2080 source.source().labels(&request).await
2081 },
2082 )
2083 .await
2084}
2085
2086async fn fetch_hits(
2088 source: &ResolvedSource,
2089 shape: &HitShape,
2090 start: &Resume,
2091 budget: u32,
2092 calls: &AtomicU32,
2093) -> Result<Fetched<Found>, SourceError> {
2094 let ceiling = ceiling(source);
2095 match shape.stream {
2096 StreamKind::Projects => {
2097 let compensating = shape.local_projects != LocalProjects::default();
2098 walk(
2099 start,
2100 budget,
2101 page_size(compensating, budget, ceiling),
2102 |found| match found {
2103 Found::Project(project) => shape.local_projects.keeps(project),
2104 Found::Task(_) => true,
2105 },
2106 |cursor, limit| async move {
2107 calls.fetch_add(1, Ordering::Relaxed);
2108 let request = PageRequest { cursor, limit };
2109 let page = source
2110 .source()
2111 .query_projects(&shape.projects, &request)
2112 .await?;
2113 Ok(Page {
2114 items: page.items.into_iter().map(Found::Project).collect(),
2115 next: page.next,
2116 })
2117 },
2118 )
2119 .await
2120 }
2121 StreamKind::Items | StreamKind::Tasks => {
2122 let compensating = shape.local_tasks != LocalTasks::default();
2123 walk(
2124 start,
2125 budget,
2126 page_size(compensating, budget, ceiling),
2127 |found| match found {
2128 Found::Task(task) => shape.local_tasks.keeps(task),
2129 Found::Project(_) => true,
2130 },
2131 |cursor, limit| async move {
2132 calls.fetch_add(1, Ordering::Relaxed);
2133 let request = PageRequest { cursor, limit };
2134 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2135 Ok(Page {
2136 items: page.items.into_iter().map(Found::Task).collect(),
2137 next: page.next,
2138 })
2139 },
2140 )
2141 .await
2142 }
2143 }
2144}
2145
2146async fn forward_edges(
2148 source: &ResolvedSource,
2149 entity: Entity,
2150 id: &NativeId,
2151 request: &PageRequest,
2152) -> Result<Page<DependencyEdge>, SourceError> {
2153 match entity {
2154 Entity::Task => {
2155 source
2156 .source()
2157 .task_dependencies(id, Direction::DependsOn, request)
2158 .await
2159 }
2160 Entity::Project => {
2161 source
2162 .source()
2163 .project_dependencies(id, Direction::DependsOn, request)
2164 .await
2165 }
2166 }
2167}
2168
2169#[expect(
2178 clippy::too_many_arguments,
2179 reason = "every argument is one axis of one walk — the source, the item, the \
2180 direction, which of its two graphs, whether the reverse is emulated, where \
2181 to resume, how many rows to return and where to count calls. Grouping them \
2182 into a struct would name the same eight values one indirection further from \
2183 the loop that reads them."
2184)]
2185async fn fetch_edges(
2186 source: &ResolvedSource,
2187 native: &NativeId,
2188 direction: Direction,
2189 entity: Entity,
2190 emulating: bool,
2191 start: &Resume,
2192 budget: u32,
2193 calls: &AtomicU32,
2194) -> Result<Fetched<DependencyEdge>, SourceError> {
2195 let ceiling = ceiling(source);
2196 if !emulating {
2197 return walk(
2198 start,
2199 budget,
2200 page_size(false, budget, ceiling),
2201 |_| true,
2202 |cursor, limit| async move {
2203 calls.fetch_add(1, Ordering::Relaxed);
2204 let request = PageRequest { cursor, limit };
2205 match entity {
2206 Entity::Task => {
2207 source
2208 .source()
2209 .task_dependencies(native, direction, &request)
2210 .await
2211 }
2212 Entity::Project => {
2213 source
2214 .source()
2215 .project_dependencies(native, direction, &request)
2216 .await
2217 }
2218 }
2219 },
2220 )
2221 .await;
2222 }
2223
2224 walk(
2225 start,
2226 budget,
2227 ceiling,
2228 |_| true,
2229 |cursor, limit| async move {
2230 calls.fetch_add(1, Ordering::Relaxed);
2231 let request = PageRequest { cursor, limit };
2232 let (ids, next) = match entity {
2233 Entity::Task => {
2234 let page = source
2235 .source()
2236 .query_tasks(&TaskQuery::default(), &request)
2237 .await?;
2238 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2239 (ids, page.next)
2240 }
2241 Entity::Project => {
2242 let page = source
2243 .source()
2244 .query_projects(&ProjectQuery::default(), &request)
2245 .await?;
2246 let ids: Vec<NativeId> =
2247 page.items.into_iter().map(|project| project.id).collect();
2248 (ids, page.next)
2249 }
2250 };
2251
2252 let mut edges = Vec::new();
2253 for id in ids {
2254 let mut inner: Option<Cursor> = None;
2255 loop {
2256 calls.fetch_add(1, Ordering::Relaxed);
2257 let request = PageRequest {
2258 cursor: inner.clone(),
2259 limit,
2260 };
2261 let page = forward_edges(source, entity, &id, &request).await?;
2262 fits(page.items.len(), limit)?;
2266 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2267 unrepeated(
2268 page.next.as_ref(),
2269 inner.as_ref(),
2270 "its forward edges were being scanned",
2271 )?;
2272 match page.next {
2273 Some(cursor) => inner = Some(cursor),
2274 None => break,
2275 }
2276 }
2277 }
2278
2279 Ok(Page { items: edges, next })
2280 },
2281 )
2282 .await
2283}