1mod comment;
19mod copy;
20mod fetch;
21mod join;
22mod local;
23mod resume;
24
25use std::collections::BTreeMap;
26use std::num::NonZeroU32;
27use std::sync::atomic::{AtomicU32, Ordering};
28
29use onetaskgraph_plugin_api::{
30 Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
31 NativeId, Page, PageRequest, Project, ProjectFilter, ProjectQuery, SecretResolver, SourceError,
32 SourceName, StatusCategory, Task, TaskQuery, TextFields, TextQuery,
33};
34use schemars::JsonSchema;
35use serde::{Deserialize, Serialize};
36
37use crate::GlobalId;
38use crate::config::Config;
39use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
40use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
41
42use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
43use join::join_all;
44use local::{LocalDocuments, LocalProjects, LocalTasks};
45pub(crate) use resume::{Owed, Resumption, StreamState};
46use resume::{Resume, StreamKind};
47
48pub use comment::{CommentList, DeletedComment, TaskDetail};
49pub use copy::{
50 BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
51 Spent,
52};
53pub use local::ProjectSelector;
54
55#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
60pub struct Qualified<T> {
61 pub id: GlobalId,
63 pub item: T,
65}
66
67#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
72pub struct QualifiedEdge {
73 pub from: QualifiedEndpoint,
78 pub to: QualifiedEndpoint,
80 pub kind: onetaskgraph_plugin_api::DependencyKind,
82}
83
84#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
86pub struct QualifiedEndpoint {
87 pub id: GlobalId,
89 pub kind: onetaskgraph_plugin_api::ItemKind,
91}
92
93impl std::fmt::Display for QualifiedEndpoint {
94 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
95 self.id.fmt(formatter)
96 }
97}
98
99#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
101#[serde(tag = "kind", rename_all = "kebab-case")]
102pub enum SearchHit {
103 Task(Qualified<Task>),
105 Project(Qualified<Project>),
107}
108
109#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
111#[serde(rename_all = "kebab-case")]
112pub enum SearchKind {
113 Tasks,
115 Projects,
117 #[default]
119 Both,
120}
121
122#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
124pub struct SourceListing {
125 pub source: SourceName,
127 pub kind: String,
139 #[serde(flatten)]
141 pub state: SourceState,
142}
143
144#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
146#[serde(tag = "state", rename_all = "kebab-case")]
147pub enum SourceState {
148 Available {
150 capabilities: Capabilities,
152 },
153 Unavailable {
155 error: SourceError,
157 },
158}
159
160#[derive(Debug, Clone, PartialEq)]
162pub struct Paging {
163 pub limit: NonZeroU32,
165 pub token: Option<PageToken>,
167}
168
169#[derive(Debug, Clone, Default, PartialEq)]
171pub struct Filters {
172 pub text: Option<TextQuery>,
174 pub labels: LabelFilter,
176 pub statuses: Vec<StatusCategory>,
178}
179
180#[derive(Debug, Clone)]
182pub struct TaskRequest {
183 pub sources: Vec<SourceName>,
185 pub filters: Filters,
187 pub project: ProjectSelector,
189 pub paging: Paging,
191}
192
193#[derive(Debug, Clone)]
195pub struct ProjectRequest {
196 pub sources: Vec<SourceName>,
198 pub filters: Filters,
200 pub paging: Paging,
202}
203
204#[derive(Debug, Clone, Default, PartialEq)]
211pub struct DocumentFilters {
212 pub text: Option<TextQuery>,
214 pub labels: LabelFilter,
216}
217
218#[derive(Debug, Clone)]
220pub struct DocumentRequest {
221 pub sources: Vec<SourceName>,
223 pub filters: DocumentFilters,
225 pub project: ProjectSelector,
227 pub paging: Paging,
229}
230
231#[derive(Debug, Clone)]
233pub struct LabelRequest {
234 pub sources: Vec<SourceName>,
236 pub paging: Paging,
238}
239
240#[derive(Debug, Clone)]
242pub struct SearchRequest {
243 pub sources: Vec<SourceName>,
245 pub text: TextQuery,
247 pub kind: SearchKind,
249 pub paging: Paging,
251}
252
253#[derive(Debug, Clone)]
255pub struct DependencyRequest {
256 pub id: GlobalId,
258 pub direction: Direction,
260 pub paging: Paging,
262}
263
264#[derive(Debug, Clone, PartialEq, thiserror::Error)]
272pub enum EngineError {
273 #[error(
275 "no source named {name:?} is configured\n\
276 next: name one of the configured sources ({configured}), or add {name:?} under \
277 `sources` — `onetaskgraph sources list` shows what this configuration has."
278 )]
279 UnknownSource {
280 name: String,
282 configured: String,
284 },
285
286 #[error(
288 "{message}\n\
289 next: page with a token exactly as the previous page reported it, and against \
290 the same configuration — or drop `--page` to start the walk again."
291 )]
292 Token {
293 message: String,
295 },
296
297 #[error(
299 "no sources are configured\n\
300 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
301 prints what each plugin accepts."
302 )]
303 NoSources,
304
305 #[error(
307 "source {name} cannot be written: its plugin is {kind}, which has no write \
308 side\n\
309 next: copy into a source whose plugin can be written — `onetaskgraph sources \
310 list` reports each one's plugin."
311 )]
312 NotWritable {
313 name: String,
315 kind: String,
317 },
318
319 #[error(
326 "source {name} has no documents: its plugin is {kind}, which holds none\n\
327 next: name a source whose plugin has documents — `onetaskgraph sources list` \
328 reports each one's plugin and what it declares."
329 )]
330 NoDocuments {
331 name: String,
333 kind: String,
335 },
336
337 #[error(
342 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
343 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
344 list` reports each one's plugin."
345 )]
346 NoComments {
347 name: String,
349 kind: String,
351 },
352
353 #[error(
355 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
356 but not added to, edited or removed\n\
357 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
358 source whose plugin can be written — `onetaskgraph sources list` reports each \
359 one's plugin."
360 )]
361 CommentsNotWritable {
362 name: String,
364 kind: String,
366 },
367
368 #[error(
370 "no task with the id {id}\n\
371 next: check the id, or list what is there — `onetaskgraph task list` reports every \
372 task the configured sources hold."
373 )]
374 NoSuchTask {
375 id: String,
377 },
378
379 #[error(
381 "task {task} has no comment with the id {comment}\n\
382 next: list its comments — `onetaskgraph task comment list {task}` reports each \
383 one's id."
384 )]
385 NoSuchComment {
386 task: String,
388 comment: String,
390 },
391
392 #[error(
397 "source {name} could not be built: {error}\n\
398 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
399 command again."
400 )]
401 SourceUnavailable {
402 name: String,
404 error: SourceError,
406 },
407
408 #[error(
414 "source {name} could not do it: {error}\n\
415 next: fix what the source named above, then run the command again."
416 )]
417 SourceFailed {
418 name: String,
420 error: SourceError,
422 },
423
424 #[error(
426 "the destination source {name} could not be built: {error}\n\
427 next: fix that source — `onetaskgraph sources list` reports its state — then \
428 copy again."
429 )]
430 DestinationUnavailable {
431 name: String,
433 error: SourceError,
435 },
436
437 #[error(
439 "no item with the id {id}\n\
440 next: check the id, or list what is there — `onetaskgraph task list` and \
441 `onetaskgraph project list` report what the configured sources hold."
442 )]
443 NoSuchItem {
444 id: String,
446 },
447
448 #[error(
453 "{item} was copied from {origin}, which that destination no longer holds\n\
454 next: re-run with --recreate to create a new item there instead, or restore \
455 {origin}."
456 )]
457 StaleOrigin {
458 item: String,
460 origin: String,
462 },
463
464 #[error(
470 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
471 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
472 {id} on its own with `onetaskgraph task copy`.",
473 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
474 )]
475 NotAMember {
476 id: GlobalId,
478 projects: Vec<GlobalId>,
480 },
481
482 #[error(
491 "{item} depends on {member}, which this copy was not told to carry and which records \
492 no origin in {destination}\n\
493 next: name {member} with --member as well, record its {destination} id at \
494 onetaskgraph.origin, or copy the whole project without --member."
495 )]
496 UnrecordedMember {
497 item: GlobalId,
499 member: GlobalId,
501 destination: SourceName,
503 },
504
505 #[error(
511 "source {name} could not do it: {error}\n\
512 next: fix what the source named above, then copy again."
513 )]
514 SourceRefused {
515 name: String,
517 error: SourceError,
519 },
520
521 #[error(
529 "the copy failed and could not be undone.\n\
530 it failed because: {error}\n\
531 it could not be undone because: {refusal}\n\
532 so the destination still holds: {left_behind}\n\
533 next: remove those items at the destination, then copy again."
534 )]
535 CopyNotUndone {
536 error: Box<EngineError>,
538 left_behind: LeftBehind,
544 refusal: SourceError,
546 },
547}
548
549#[derive(Debug, Clone, PartialEq)]
558pub struct LeftBehind {
559 first: GlobalId,
561 rest: Vec<GlobalId>,
563}
564
565impl LeftBehind {
566 #[must_use]
568 pub fn new(first: GlobalId) -> Self {
569 Self {
570 first,
571 rest: Vec::new(),
572 }
573 }
574
575 pub fn push(&mut self, id: GlobalId) {
577 self.rest.push(id);
578 }
579
580 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
582 std::iter::once(&self.first).chain(self.rest.iter())
583 }
584}
585
586impl std::fmt::Display for LeftBehind {
587 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
589 write!(formatter, "{}", self.first)?;
590 for id in &self.rest {
591 write!(formatter, ", {id}")?;
592 }
593 Ok(())
594 }
595}
596
597pub enum ConfiguredSource {
604 Ready(ResolvedSource),
606 Unavailable(UnavailableSource),
608}
609
610impl ConfiguredSource {
611 #[must_use]
613 pub fn name(&self) -> &SourceName {
614 match self {
615 Self::Ready(source) => source.name(),
616 Self::Unavailable(source) => source.name(),
617 }
618 }
619}
620
621pub struct Engine {
623 sources: Vec<ConfiguredSource>,
625 selection: Vec<SourceName>,
627}
628
629impl Engine {
630 #[must_use]
638 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
639 let (ready, unavailable) = resolve_available(config, secrets);
640 Self::new(
641 ready
642 .into_iter()
643 .map(ConfiguredSource::Ready)
644 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
645 .collect(),
646 config.selected_sources(),
647 )
648 }
649
650 #[must_use]
653 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
654 Self { sources, selection }
655 }
656
657 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
659 self.sources.iter().filter_map(|source| match source {
660 ConfiguredSource::Ready(ready) => Some(ready),
661 ConfiguredSource::Unavailable(_) => None,
662 })
663 }
664
665 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
667 self.sources.iter().filter_map(|source| match source {
668 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
669 ConfiguredSource::Ready(_) => None,
670 })
671 }
672
673 #[must_use]
675 pub fn listing(&self) -> Vec<SourceListing> {
676 let mut listings: Vec<SourceListing> = self
677 .ready()
678 .map(|source| SourceListing {
679 source: source.name().clone(),
680 kind: source.kind().to_owned(),
681 state: SourceState::Available {
682 capabilities: source.source().capabilities(),
683 },
684 })
685 .chain(self.unavailable().map(|source| SourceListing {
686 source: source.name().clone(),
687 kind: source.kind().to_owned(),
688 state: SourceState::Unavailable {
689 error: source.error().clone(),
690 },
691 }))
692 .collect();
693 listings.sort_by(|left, right| left.source.cmp(&right.source));
694 listings
695 }
696
697 #[must_use]
703 pub fn has(&self, name: &SourceName) -> bool {
704 self.sources.iter().any(|source| source.name() == name)
705 }
706
707 pub async fn tasks(
715 &self,
716 request: &TaskRequest,
717 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
718 let mut names = self.resolve_selection(&request.sources)?;
719 if let ProjectSelector::Qualified(id) = &request.project {
724 self.known(&id.source)?;
725 names.retain(|name| name == &id.source);
726 }
727 let query = shape("task-list", &names, &(&request.filters, &request.project));
728 let states = resumption(
729 self,
730 request.paging.token.as_ref(),
731 &[StreamKind::Items],
732 &query,
733 )?;
734 let budget = request.paging.limit.get();
735
736 let mut answer = Answer::new();
737 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
738
739 let shapes: Vec<TaskShape> = ready
740 .iter()
741 .map(|source| {
742 shape_tasks(
743 &source.source().capabilities(),
744 &request.filters,
745 &project_filter(&request.project),
746 )
747 })
748 .collect();
749 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
750 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
751
752 let walks = ready
753 .iter()
754 .enumerate()
755 .map(|(index, source)| {
756 fetch_tasks(
757 source,
758 &shapes[index],
759 &starts[index],
760 budget,
761 &counters[index],
762 )
763 })
764 .collect();
765
766 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
767 answer.finish(
768 streams,
769 budget,
770 owed(&states),
771 &query,
772 |name, task: Task| Qualified {
773 id: GlobalId::new(name.clone(), task.id.clone()),
774 item: task,
775 },
776 )
777 }
778
779 pub async fn projects(
785 &self,
786 request: &ProjectRequest,
787 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
788 let names = self.resolve_selection(&request.sources)?;
789 let query = shape("project-list", &names, &request.filters);
790 let states = resumption(
791 self,
792 request.paging.token.as_ref(),
793 &[StreamKind::Items],
794 &query,
795 )?;
796 let budget = request.paging.limit.get();
797
798 let mut answer = Answer::new();
799 let mut with_projects = Vec::new();
805 for source in answer.split(self, &names) {
806 if source.source().capabilities().projects.is_native() {
807 with_projects.push(source);
808 } else {
809 answer.unreachable_predicate(source, Predicate::Project);
810 }
811 }
812 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
813
814 let shapes: Vec<ProjectShape> = ready
815 .iter()
816 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
817 .collect();
818 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
819 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
820
821 let walks = ready
822 .iter()
823 .enumerate()
824 .map(|(index, source)| {
825 fetch_projects(
826 source,
827 &shapes[index],
828 &starts[index],
829 budget,
830 &counters[index],
831 )
832 })
833 .collect();
834
835 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
836 answer.finish(
837 streams,
838 budget,
839 owed(&states),
840 &query,
841 |name, project: Project| Qualified {
842 id: GlobalId::new(name.clone(), project.id.clone()),
843 item: project,
844 },
845 )
846 }
847
848 pub async fn documents(
860 &self,
861 request: &DocumentRequest,
862 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
863 let mut names = self.resolve_selection(&request.sources)?;
864 if let ProjectSelector::Qualified(id) = &request.project {
867 self.known(&id.source)?;
868 names.retain(|name| name == &id.source);
869 }
870 let query = shape(
871 "document-list",
872 &names,
873 &(&request.filters, &request.project),
874 );
875 let states = resumption(
876 self,
877 request.paging.token.as_ref(),
878 &[StreamKind::Items],
879 &query,
880 )?;
881 let budget = request.paging.limit.get();
882
883 let mut answer = Answer::new();
884 let mut with_documents = Vec::new();
885 for source in answer.split(self, &names) {
886 if source.source().capabilities().documents.is_native() {
887 with_documents.push(source);
888 } else {
889 answer.unreachable_predicate(source, Predicate::Document);
890 }
891 }
892 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
893
894 let shapes: Vec<DocumentShape> = ready
895 .iter()
896 .map(|source| {
897 shape_documents(
898 &source.source().capabilities(),
899 &request.filters,
900 &project_filter(&request.project),
901 )
902 })
903 .collect();
904 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
905 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
906
907 let walks = ready
908 .iter()
909 .enumerate()
910 .map(|(index, source)| {
911 fetch_documents(
912 source,
913 &shapes[index],
914 &starts[index],
915 budget,
916 &counters[index],
917 )
918 })
919 .collect();
920
921 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
922 answer.finish(
923 streams,
924 budget,
925 owed(&states),
926 &query,
927 |name, document: Document| Qualified {
928 id: GlobalId::new(name.clone(), document.id.clone()),
929 item: document,
930 },
931 )
932 }
933
934 pub async fn labels(
940 &self,
941 request: &LabelRequest,
942 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
943 let names = self.resolve_selection(&request.sources)?;
944 let query = shape("label-list", &names, &());
945 let states = resumption(
946 self,
947 request.paging.token.as_ref(),
948 &[StreamKind::Items],
949 &query,
950 )?;
951 let budget = request.paging.limit.get();
952
953 let mut answer = Answer::new();
954 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
955 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
956 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
957
958 let walks = ready
959 .iter()
960 .enumerate()
961 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
962 .collect();
963
964 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
965 answer.finish(
966 streams,
967 budget,
968 owed(&states),
969 &query,
970 |name, label: Label| Qualified {
971 id: GlobalId::new(name.clone(), label.id.clone()),
972 item: label,
973 },
974 )
975 }
976
977 pub async fn search(
983 &self,
984 request: &SearchRequest,
985 ) -> Result<QueryResponse<SearchHit>, EngineError> {
986 let names = self.resolve_selection(&request.sources)?;
987 let reads: &[StreamKind] = match request.kind {
991 SearchKind::Tasks => &[StreamKind::Tasks],
992 SearchKind::Projects => &[StreamKind::Projects],
993 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
994 };
995 let query = shape("search", &names, &request.text);
1000 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1001 let budget = request.paging.limit.get();
1002 let filters = Filters {
1003 text: Some(request.text.clone()),
1004 ..Filters::default()
1005 };
1006
1007 let mut answer = Answer::new();
1008
1009 let mut ready = Vec::new();
1012 let mut kinds = Vec::new();
1013 let mut starts = Vec::new();
1014 for source in answer.split(self, &names) {
1015 let mut streams = Vec::new();
1016 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1017 streams.push(StreamKind::Tasks);
1018 }
1019 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1020 if source.source().capabilities().projects.is_native() {
1021 streams.push(StreamKind::Projects);
1022 } else {
1023 answer.unreachable_predicate(source, Predicate::Project);
1024 }
1025 }
1026 for stream in streams {
1027 if let Some(resume) = resume_at(&states, source.name(), stream) {
1028 ready.push(source);
1029 kinds.push(stream);
1030 starts.push(resume);
1031 }
1032 }
1033 }
1034
1035 let shapes: Vec<HitShape> = ready
1036 .iter()
1037 .zip(kinds.iter())
1038 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1039 .collect();
1040 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1041 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1042
1043 let walks = ready
1044 .iter()
1045 .enumerate()
1046 .map(|(index, source)| {
1047 fetch_hits(
1048 source,
1049 &shapes[index],
1050 &starts[index],
1051 budget,
1052 &counters[index],
1053 )
1054 })
1055 .collect();
1056
1057 let streams =
1058 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1059 answer.finish(
1060 streams,
1061 budget,
1062 owed(&states),
1063 &query,
1064 |name, found: Found| match found {
1065 Found::Task(task) => SearchHit::Task(Qualified {
1066 id: GlobalId::new(name.clone(), task.id.clone()),
1067 item: task,
1068 }),
1069 Found::Project(project) => SearchHit::Project(Qualified {
1070 id: GlobalId::new(name.clone(), project.id.clone()),
1071 item: project,
1072 }),
1073 },
1074 )
1075 }
1076
1077 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1084 let name = self.known(&id.source)?;
1085 let mut answer = Answer::new();
1086 let selected = answer.split(self, std::slice::from_ref(&name));
1087 let Some(source) = selected.first() else {
1088 return answer.nothing();
1089 };
1090 let found = source.source().get_task(&id.native).await;
1091 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1092 answer.one(source, found, |task| Qualified {
1093 id: qualified,
1094 item: task,
1095 })
1096 }
1097
1098 pub async fn project(
1104 &self,
1105 id: &GlobalId,
1106 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1107 let name = self.known(&id.source)?;
1108 let mut answer = Answer::new();
1109 let selected = answer.split(self, std::slice::from_ref(&name));
1110 let Some(source) = selected.first() else {
1111 return answer.nothing();
1112 };
1113 let found = source.source().get_project(&id.native).await;
1114 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1115 answer.one(source, found, |project| Qualified {
1116 id: qualified,
1117 item: project,
1118 })
1119 }
1120
1121 pub async fn document(
1131 &self,
1132 id: &GlobalId,
1133 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1134 let name = self.known(&id.source)?;
1135 let mut answer = Answer::new();
1136 let selected = answer.split(self, std::slice::from_ref(&name));
1137 let Some(source) = selected.first() else {
1138 return answer.nothing();
1139 };
1140 if !source.source().capabilities().documents.is_native() {
1141 answer.unreachable_predicate(source, Predicate::Document);
1142 return answer.nothing();
1143 }
1144 let found = source.source().get_document(&id.native).await;
1145 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1146 answer.one(source, found, |document| Qualified {
1147 id: qualified,
1148 item: document,
1149 })
1150 }
1151
1152 pub async fn task_dependencies(
1159 &self,
1160 request: &DependencyRequest,
1161 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1162 self.dependencies(request, Entity::Task).await
1163 }
1164
1165 pub async fn project_dependencies(
1171 &self,
1172 request: &DependencyRequest,
1173 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1174 self.dependencies(request, Entity::Project).await
1175 }
1176
1177 async fn dependencies(
1180 &self,
1181 request: &DependencyRequest,
1182 entity: Entity,
1183 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1184 let name = self.known(&request.id.source)?;
1185 let query = shape(
1186 "dependencies",
1187 std::slice::from_ref(&name),
1188 &(entity, &request.id.native, request.direction),
1189 );
1190 let states = resumption(
1191 self,
1192 request.paging.token.as_ref(),
1193 &[StreamKind::Items],
1194 &query,
1195 )?;
1196 let budget = request.paging.limit.get();
1197
1198 let mut answer = Answer::new();
1199 let (ready, starts) = walking(
1200 answer.split(self, std::slice::from_ref(&name)),
1201 &states,
1202 StreamKind::Items,
1203 );
1204 let Some(source) = ready.first() else {
1205 return answer.nothing();
1206 };
1207
1208 let capabilities = source.source().capabilities();
1209 let support = match entity {
1210 Entity::Task => capabilities.task_dependencies,
1211 Entity::Project => capabilities.project_dependencies,
1212 };
1213 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1217 let mut outcomes = Outcomes::default();
1218 if request.direction == Direction::DependedOnBy {
1219 if emulating {
1220 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1221 } else {
1222 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1223 }
1224 }
1225
1226 let counters = vec![AtomicU32::new(0)];
1227 let walked = fetch_edges(
1228 source,
1229 &request.id.native,
1230 request.direction,
1231 entity,
1232 emulating,
1233 &starts[0],
1234 budget,
1235 &counters[0],
1236 )
1237 .await;
1238
1239 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1240 answer.finish(
1241 streams,
1242 budget,
1243 owed(&states),
1244 &query,
1245 |name, edge: DependencyEdge| QualifiedEdge {
1246 from: qualify_endpoint(name, edge.from),
1247 to: qualify_endpoint(name, edge.to),
1248 kind: edge.kind,
1249 },
1250 )
1251 }
1252
1253 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1255 if asked.is_empty() {
1256 if self.selection.is_empty() {
1257 return Err(EngineError::NoSources);
1258 }
1259 return Ok(self.selection.clone());
1260 }
1261 asked.iter().map(|name| self.known(name)).collect()
1262 }
1263
1264 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1266 if self.has(name) {
1267 return Ok(name.clone());
1268 }
1269 if self.sources.is_empty() {
1270 return Err(EngineError::NoSources);
1271 }
1272 Err(EngineError::UnknownSource {
1273 name: name.to_string(),
1274 configured: self
1275 .listing()
1276 .iter()
1277 .map(|listing| listing.source.to_string())
1278 .collect::<Vec<_>>()
1279 .join(", "),
1280 })
1281 }
1282}
1283
1284fn qualify_endpoint(
1285 source: &SourceName,
1286 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1287) -> QualifiedEndpoint {
1288 let kind = endpoint.kind;
1289 let is_qualified = endpoint.is_qualified();
1290 let endpoint_id = endpoint.into_id();
1291 QualifiedEndpoint {
1292 id: if is_qualified {
1293 endpoint_id
1294 .parse()
1295 .expect("plugin-api validates qualified dependency endpoints")
1296 } else {
1297 GlobalId::new(source.clone(), NativeId(endpoint_id))
1298 },
1299 kind,
1300 }
1301}
1302
1303#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1305enum Entity {
1306 Task,
1308 Project,
1310}
1311
1312enum Found {
1314 Task(Task),
1316 Project(Project),
1318}
1319
1320#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1322enum Outcome {
1323 PushedDown,
1325 AppliedLocally,
1327 Emulated,
1329 Unavailable,
1331}
1332
1333#[derive(Debug, Clone, Default, PartialEq)]
1345struct Outcomes(BTreeMap<Predicate, Outcome>);
1346
1347impl Outcomes {
1348 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1354 self.0.insert(predicate, outcome);
1355 }
1356
1357 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1360 for predicate in predicates {
1361 self.record(predicate, outcome);
1362 }
1363 }
1364
1365 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1367 self.0
1368 .iter()
1369 .filter(|(_, recorded)| **recorded == outcome)
1370 .map(|(predicate, _)| *predicate)
1371 .collect()
1372 }
1373}
1374
1375struct TaskShape {
1377 pushed: TaskQuery,
1379 local: LocalTasks,
1381 outcomes: Outcomes,
1383}
1384
1385struct ProjectShape {
1387 pushed: ProjectQuery,
1389 local: LocalProjects,
1391 outcomes: Outcomes,
1393}
1394
1395struct DocumentShape {
1397 pushed: DocumentQuery,
1399 local: LocalDocuments,
1401 outcomes: Outcomes,
1403}
1404
1405struct HitShape {
1407 stream: StreamKind,
1409 tasks: TaskQuery,
1411 projects: ProjectQuery,
1413 local_tasks: LocalTasks,
1415 local_projects: LocalProjects,
1417 outcomes: Outcomes,
1419}
1420
1421struct Answer {
1427 plans: Vec<SourcePlan>,
1429 errors: Vec<SourceFailure>,
1431}
1432
1433impl Answer {
1434 fn new() -> Self {
1435 Self {
1436 plans: Vec::new(),
1437 errors: Vec::new(),
1438 }
1439 }
1440
1441 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1446 let mut selected = Vec::new();
1447 for name in names {
1448 match engine.sources.iter().find(|source| source.name() == name) {
1449 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1450 Some(ConfiguredSource::Unavailable(source)) => {
1451 self.errors.push(source.failure());
1452 }
1453 None => {}
1454 }
1455 }
1456 selected
1457 }
1458
1459 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1461 let mut outcomes = Outcomes::default();
1462 outcomes.record(predicate, Outcome::Unavailable);
1463 self.plans.push(plan_for(source, outcomes, 0));
1464 }
1465
1466 fn collect<T>(
1468 &mut self,
1469 ready: &[&ResolvedSource],
1470 walked: Vec<Result<Fetched<T>, SourceError>>,
1471 counters: &[AtomicU32],
1472 outcomes: Vec<Outcomes>,
1473 ) -> Vec<Stream<T>> {
1474 let kinds = vec![StreamKind::Items; ready.len()];
1475 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1476 }
1477
1478 fn collect_streams<T>(
1480 &mut self,
1481 ready: &[&ResolvedSource],
1482 kinds: &[StreamKind],
1483 walked: Vec<Result<Fetched<T>, SourceError>>,
1484 counters: &[AtomicU32],
1485 outcomes: Vec<Outcomes>,
1486 ) -> Vec<Stream<T>> {
1487 let mut streams = Vec::new();
1488 for (index, result) in walked.into_iter().enumerate() {
1489 let source = ready[index];
1490 let pages = counters[index].load(Ordering::Relaxed);
1491 self.plans
1492 .push(plan_for(source, outcomes[index].clone(), pages));
1493 match result {
1494 Ok(fetched) => streams.push(Stream {
1495 source: source.name().clone(),
1496 kind: kinds[index],
1497 fetched,
1498 }),
1499 Err(error) => self.errors.push(SourceFailure {
1502 source: source.name().clone(),
1503 error,
1504 }),
1505 }
1506 }
1507 streams
1508 }
1509
1510 fn one<T, U>(
1512 mut self,
1513 source: &ResolvedSource,
1514 found: Result<Option<T>, SourceError>,
1515 qualify: impl FnOnce(T) -> U,
1516 ) -> Result<QueryResponse<U>, EngineError> {
1517 self.plans.push(plan_for(source, Outcomes::default(), 1));
1518 let items = match found {
1519 Ok(Some(item)) => vec![qualify(item)],
1520 Ok(None) => Vec::new(),
1521 Err(error) => {
1522 self.errors.push(SourceFailure {
1523 source: source.name().clone(),
1524 error,
1525 });
1526 Vec::new()
1527 }
1528 };
1529 Ok(QueryResponse {
1530 items,
1531 next: None,
1532 plan: QueryPlan {
1533 per_source: merge_plans(self.plans),
1534 },
1535 errors: self.errors,
1536 })
1537 }
1538
1539 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1541 Ok(QueryResponse {
1542 items: Vec::new(),
1543 next: None,
1544 plan: QueryPlan {
1545 per_source: merge_plans(self.plans),
1546 },
1547 errors: self.errors,
1548 })
1549 }
1550
1551 fn finish<T, U>(
1556 self,
1557 streams: Vec<Stream<T>>,
1558 budget: u32,
1559 first: Option<&Owed>,
1560 query: &str,
1561 qualify: impl Fn(&SourceName, T) -> U,
1562 ) -> Result<QueryResponse<U>, EngineError> {
1563 let (rows, states, owed) = merge(streams, budget, first);
1564 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1565 Ok(QueryResponse {
1566 items: rows
1567 .into_iter()
1568 .map(|(name, item)| qualify(&name, item))
1569 .collect(),
1570 next,
1571 plan: QueryPlan {
1572 per_source: merge_plans(self.plans),
1573 },
1574 errors: self.errors,
1575 })
1576 }
1577}
1578
1579fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1582 SourcePlan {
1583 source: source.name().clone(),
1584 kind: source.kind().to_owned(),
1585 pushed_down: outcomes.with(Outcome::PushedDown),
1586 applied_locally: outcomes.with(Outcome::AppliedLocally),
1587 emulated: outcomes.with(Outcome::Emulated),
1588 unavailable: outcomes.with(Outcome::Unavailable),
1589 pages_fetched: pages,
1590 }
1591}
1592
1593fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1598 let mut merged: Vec<SourcePlan> = Vec::new();
1599 for plan in plans {
1600 if let Some(existing) = merged
1601 .iter_mut()
1602 .find(|existing| existing.source == plan.source)
1603 {
1604 existing.pushed_down.extend(plan.pushed_down);
1605 existing.applied_locally.extend(plan.applied_locally);
1606 existing.emulated.extend(plan.emulated);
1607 existing.unavailable.extend(plan.unavailable);
1608 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1609 for list in [
1610 &mut existing.pushed_down,
1611 &mut existing.applied_locally,
1612 &mut existing.emulated,
1613 &mut existing.unavailable,
1614 ] {
1615 list.sort_unstable();
1616 list.dedup();
1617 }
1618 } else {
1619 merged.push(plan);
1620 }
1621 }
1622 merged
1623}
1624
1625fn walking<'a>(
1631 selected: Vec<&'a ResolvedSource>,
1632 states: &Option<Resumption>,
1633 kind: StreamKind,
1634) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1635 let mut ready = Vec::new();
1636 let mut starts = Vec::new();
1637 for source in selected {
1638 if let Some(resume) = resume_at(states, source.name(), kind) {
1639 ready.push(source);
1640 starts.push(resume);
1641 }
1642 }
1643 (ready, starts)
1644}
1645
1646fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1664 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1665 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1666}
1667
1668fn fingerprint(text: &str) -> String {
1670 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1671 for byte in text.as_bytes() {
1672 hash ^= u64::from(*byte);
1673 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1674 }
1675 format!("{hash:016x}")
1676}
1677
1678fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1683 document.as_ref()?.owed.as_ref()
1684}
1685
1686fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1688 match states {
1689 None => Some(Resume::default()),
1690 Some(document) => document
1691 .streams
1692 .iter()
1693 .find(|state| &state.source == source && state.stream == kind)
1694 .map(|state| state.resume.clone()),
1695 }
1696}
1697
1698fn resumption(
1729 engine: &Engine,
1730 token: Option<&PageToken>,
1731 reads: &[StreamKind],
1732 query: &str,
1733) -> Result<Option<Resumption>, EngineError> {
1734 let Some(document) = token.map(PageToken::decode) else {
1735 return Ok(None);
1736 };
1737
1738 if document.query != query {
1745 return Err(EngineError::Token {
1746 message: "this page token was written by a different query — resume the walk it \
1747 came from, or drop --page to start this one from the beginning"
1748 .to_owned(),
1749 });
1750 }
1751 let states = &document.streams;
1752
1753 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
1754 for state in states {
1755 if !reads.contains(&state.stream) {
1756 return Err(EngineError::Token {
1757 message: format!(
1758 "this page token resumes {}, which this command does not read — it \
1759 was written by a different query",
1760 state.stream.describe()
1761 ),
1762 });
1763 }
1764 let ceiling = engine
1765 .ready()
1766 .find(|source| source.name() == &state.source)
1767 .map(ceiling);
1768 if ceiling.is_none() && !engine.has(&state.source) {
1769 return Err(EngineError::Token {
1770 message: format!(
1771 "this page token resumes a source called {:?}, which this \
1772 configuration does not have",
1773 state.source.as_str()
1774 ),
1775 });
1776 }
1777 if let Some(ceiling) = ceiling
1778 && state.resume.skip >= ceiling
1779 {
1780 return Err(EngineError::Token {
1781 message: format!(
1782 "this page token resumes {} rows into a page of source {:?}, which \
1783 serves at most {ceiling}",
1784 state.resume.skip,
1785 state.source.as_str()
1786 ),
1787 });
1788 }
1789 if seen.contains(&(&state.source, state.stream)) {
1790 return Err(EngineError::Token {
1791 message: format!(
1792 "this page token gives source {:?} two places to resume from",
1793 state.source.as_str()
1794 ),
1795 });
1796 }
1797 seen.push((&state.source, state.stream));
1798 }
1799
1800 if let Some(owed) = &document.owed
1805 && !document
1806 .streams
1807 .iter()
1808 .any(|state| state.source == owed.source && state.stream == owed.stream)
1809 {
1810 return Err(EngineError::Token {
1811 message: format!(
1812 "this page token owes the next row to a stream it does not resume, \
1813 {:?}'s {}",
1814 owed.source.as_str(),
1815 owed.stream.describe()
1816 ),
1817 });
1818 }
1819
1820 Ok(Some(document))
1821}
1822
1823fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
1829 match selector {
1830 ProjectSelector::Any => ProjectFilter::Any,
1831 ProjectSelector::Orphans => ProjectFilter::Orphans,
1832 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
1833 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
1834 }
1835}
1836
1837fn text_predicates(fields: TextFields) -> Vec<Predicate> {
1839 match fields {
1840 TextFields::Title => vec![Predicate::SearchTitle],
1841 TextFields::Content => vec![Predicate::SearchContent],
1842 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
1843 }
1844}
1845
1846fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
1853 match fields {
1854 TextFields::Title => capabilities.search_title.is_native(),
1855 TextFields::Content => capabilities.search_content.is_native(),
1856 TextFields::TitleOrContent => {
1857 capabilities.search_title.is_native() && capabilities.search_content.is_native()
1858 }
1859 }
1860}
1861
1862fn shape_tasks(
1864 capabilities: &Capabilities,
1865 filters: &Filters,
1866 project: &ProjectFilter,
1867) -> TaskShape {
1868 let mut pushed = TaskQuery::default();
1869 let mut local = LocalTasks::default();
1870 let mut outcomes = Outcomes::default();
1871
1872 if !filters.labels.is_empty() {
1873 if capabilities.filter_by_label.is_native() {
1874 pushed.labels = filters.labels.clone();
1875 outcomes.record(Predicate::Label, Outcome::PushedDown);
1876 } else {
1877 local.labels = Some(filters.labels.clone());
1878 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1879 }
1880 }
1881 if !filters.statuses.is_empty() {
1882 if capabilities.filter_by_status.is_native() {
1883 pushed.statuses.clone_from(&filters.statuses);
1884 outcomes.record(Predicate::Status, Outcome::PushedDown);
1885 } else {
1886 local.statuses.clone_from(&filters.statuses);
1887 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1888 }
1889 }
1890 if let Some(text) = &filters.text {
1891 let predicates = text_predicates(text.fields);
1892 if searches_natively(capabilities, text.fields) {
1893 pushed.text = Some(text.clone());
1894 outcomes.record_all(predicates, Outcome::PushedDown);
1895 } else {
1896 local.text = Some(text.clone());
1897 outcomes.record_all(predicates, Outcome::AppliedLocally);
1898 }
1899 }
1900 match project {
1901 ProjectFilter::Any => {}
1902 ProjectFilter::Orphans => {
1903 if capabilities.orphan_tasks.is_native() {
1904 pushed.project = ProjectFilter::Orphans;
1905 outcomes.record(Predicate::Project, Outcome::PushedDown);
1906 } else {
1907 local.project = Some(ProjectFilter::Orphans);
1908 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1909 }
1910 }
1911 ProjectFilter::Is(id) => {
1912 if capabilities.projects.is_native() {
1913 pushed.project = ProjectFilter::Is(id.clone());
1914 outcomes.record(Predicate::Project, Outcome::PushedDown);
1915 } else {
1916 local.project = Some(ProjectFilter::Is(id.clone()));
1917 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1918 }
1919 }
1920 }
1921
1922 TaskShape {
1923 pushed,
1924 local,
1925 outcomes,
1926 }
1927}
1928
1929fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
1931 let mut pushed = ProjectQuery::default();
1932 let mut local = LocalProjects::default();
1933 let mut outcomes = Outcomes::default();
1934
1935 if !filters.labels.is_empty() {
1936 if capabilities.filter_by_label.is_native() {
1937 pushed.labels = filters.labels.clone();
1938 outcomes.record(Predicate::Label, Outcome::PushedDown);
1939 } else {
1940 local.labels = Some(filters.labels.clone());
1941 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1942 }
1943 }
1944 if !filters.statuses.is_empty() {
1945 if capabilities.filter_by_status.is_native() {
1946 pushed.statuses.clone_from(&filters.statuses);
1947 outcomes.record(Predicate::Status, Outcome::PushedDown);
1948 } else {
1949 local.statuses.clone_from(&filters.statuses);
1950 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1951 }
1952 }
1953 if let Some(text) = &filters.text {
1954 let predicates = text_predicates(text.fields);
1955 if searches_natively(capabilities, text.fields) {
1956 pushed.text = Some(text.clone());
1957 outcomes.record_all(predicates, Outcome::PushedDown);
1958 } else {
1959 local.text = Some(text.clone());
1960 outcomes.record_all(predicates, Outcome::AppliedLocally);
1961 }
1962 }
1963
1964 ProjectShape {
1965 pushed,
1966 local,
1967 outcomes,
1968 }
1969}
1970
1971fn shape_documents(
1976 capabilities: &Capabilities,
1977 filters: &DocumentFilters,
1978 project: &ProjectFilter,
1979) -> DocumentShape {
1980 let mut pushed = DocumentQuery::default();
1981 let mut local = LocalDocuments::default();
1982 let mut outcomes = Outcomes::default();
1983
1984 if !filters.labels.is_empty() {
1985 if capabilities.filter_by_label.is_native() {
1986 pushed.labels = filters.labels.clone();
1987 outcomes.record(Predicate::Label, Outcome::PushedDown);
1988 } else {
1989 local.labels = Some(filters.labels.clone());
1990 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1991 }
1992 }
1993 if let Some(text) = &filters.text {
1994 let predicates = text_predicates(text.fields);
1995 if searches_natively(capabilities, text.fields) {
1996 pushed.text = Some(text.clone());
1997 outcomes.record_all(predicates, Outcome::PushedDown);
1998 } else {
1999 local.text = Some(text.clone());
2000 outcomes.record_all(predicates, Outcome::AppliedLocally);
2001 }
2002 }
2003 match project {
2004 ProjectFilter::Any => {}
2005 ProjectFilter::Orphans => {
2006 if capabilities.orphan_tasks.is_native() {
2007 pushed.project = ProjectFilter::Orphans;
2008 outcomes.record(Predicate::Project, Outcome::PushedDown);
2009 } else {
2010 local.project = Some(ProjectFilter::Orphans);
2011 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2012 }
2013 }
2014 ProjectFilter::Is(id) => {
2015 if capabilities.projects.is_native() {
2016 pushed.project = ProjectFilter::Is(id.clone());
2017 outcomes.record(Predicate::Project, Outcome::PushedDown);
2018 } else {
2019 local.project = Some(ProjectFilter::Is(id.clone()));
2020 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2021 }
2022 }
2023 }
2024
2025 DocumentShape {
2026 pushed,
2027 local,
2028 outcomes,
2029 }
2030}
2031
2032fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2034 match stream {
2035 StreamKind::Projects => {
2036 let shaped = shape_projects(capabilities, filters);
2037 HitShape {
2038 stream,
2039 tasks: TaskQuery::default(),
2040 projects: shaped.pushed,
2041 local_tasks: LocalTasks::default(),
2042 local_projects: shaped.local,
2043 outcomes: shaped.outcomes,
2044 }
2045 }
2046 StreamKind::Items | StreamKind::Tasks => {
2047 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any);
2048 HitShape {
2049 stream,
2050 tasks: shaped.pushed,
2051 projects: ProjectQuery::default(),
2052 local_tasks: shaped.local,
2053 local_projects: LocalProjects::default(),
2054 outcomes: shaped.outcomes,
2055 }
2056 }
2057 }
2058}
2059
2060fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2067 if compensating {
2068 ceiling
2069 } else {
2070 budget.min(ceiling)
2071 }
2072}
2073
2074fn ceiling(source: &ResolvedSource) -> u32 {
2076 source.source().capabilities().max_page_size.max(1)
2077}
2078
2079async fn fetch_tasks(
2081 source: &ResolvedSource,
2082 shape: &TaskShape,
2083 start: &Resume,
2084 budget: u32,
2085 calls: &AtomicU32,
2086) -> Result<Fetched<Task>, SourceError> {
2087 let compensating = shape.local != LocalTasks::default();
2088 walk(
2089 start,
2090 budget,
2091 page_size(compensating, budget, ceiling(source)),
2092 |task| shape.local.keeps(task),
2093 |cursor, limit| async move {
2094 calls.fetch_add(1, Ordering::Relaxed);
2095 let request = PageRequest { cursor, limit };
2096 source.source().query_tasks(&shape.pushed, &request).await
2097 },
2098 )
2099 .await
2100}
2101
2102async fn fetch_projects(
2104 source: &ResolvedSource,
2105 shape: &ProjectShape,
2106 start: &Resume,
2107 budget: u32,
2108 calls: &AtomicU32,
2109) -> Result<Fetched<Project>, SourceError> {
2110 let compensating = shape.local != LocalProjects::default();
2111 walk(
2112 start,
2113 budget,
2114 page_size(compensating, budget, ceiling(source)),
2115 |project| shape.local.keeps(project),
2116 |cursor, limit| async move {
2117 calls.fetch_add(1, Ordering::Relaxed);
2118 let request = PageRequest { cursor, limit };
2119 source
2120 .source()
2121 .query_projects(&shape.pushed, &request)
2122 .await
2123 },
2124 )
2125 .await
2126}
2127
2128async fn fetch_documents(
2130 source: &ResolvedSource,
2131 shape: &DocumentShape,
2132 start: &Resume,
2133 budget: u32,
2134 calls: &AtomicU32,
2135) -> Result<Fetched<Document>, SourceError> {
2136 let compensating = shape.local != LocalDocuments::default();
2137 walk(
2138 start,
2139 budget,
2140 page_size(compensating, budget, ceiling(source)),
2141 |document| shape.local.keeps(document),
2142 |cursor, limit| async move {
2143 calls.fetch_add(1, Ordering::Relaxed);
2144 let request = PageRequest { cursor, limit };
2145 source
2146 .source()
2147 .query_documents(&shape.pushed, &request)
2148 .await
2149 },
2150 )
2151 .await
2152}
2153
2154async fn fetch_labels(
2156 source: &ResolvedSource,
2157 start: &Resume,
2158 budget: u32,
2159 calls: &AtomicU32,
2160) -> Result<Fetched<Label>, SourceError> {
2161 walk(
2162 start,
2163 budget,
2164 page_size(false, budget, ceiling(source)),
2165 |_| true,
2166 |cursor, limit| async move {
2167 calls.fetch_add(1, Ordering::Relaxed);
2168 let request = PageRequest { cursor, limit };
2169 source.source().labels(&request).await
2170 },
2171 )
2172 .await
2173}
2174
2175async fn fetch_hits(
2177 source: &ResolvedSource,
2178 shape: &HitShape,
2179 start: &Resume,
2180 budget: u32,
2181 calls: &AtomicU32,
2182) -> Result<Fetched<Found>, SourceError> {
2183 let ceiling = ceiling(source);
2184 match shape.stream {
2185 StreamKind::Projects => {
2186 let compensating = shape.local_projects != LocalProjects::default();
2187 walk(
2188 start,
2189 budget,
2190 page_size(compensating, budget, ceiling),
2191 |found| match found {
2192 Found::Project(project) => shape.local_projects.keeps(project),
2193 Found::Task(_) => true,
2194 },
2195 |cursor, limit| async move {
2196 calls.fetch_add(1, Ordering::Relaxed);
2197 let request = PageRequest { cursor, limit };
2198 let page = source
2199 .source()
2200 .query_projects(&shape.projects, &request)
2201 .await?;
2202 Ok(Page {
2203 items: page.items.into_iter().map(Found::Project).collect(),
2204 next: page.next,
2205 })
2206 },
2207 )
2208 .await
2209 }
2210 StreamKind::Items | StreamKind::Tasks => {
2211 let compensating = shape.local_tasks != LocalTasks::default();
2212 walk(
2213 start,
2214 budget,
2215 page_size(compensating, budget, ceiling),
2216 |found| match found {
2217 Found::Task(task) => shape.local_tasks.keeps(task),
2218 Found::Project(_) => true,
2219 },
2220 |cursor, limit| async move {
2221 calls.fetch_add(1, Ordering::Relaxed);
2222 let request = PageRequest { cursor, limit };
2223 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2224 Ok(Page {
2225 items: page.items.into_iter().map(Found::Task).collect(),
2226 next: page.next,
2227 })
2228 },
2229 )
2230 .await
2231 }
2232 }
2233}
2234
2235async fn forward_edges(
2237 source: &ResolvedSource,
2238 entity: Entity,
2239 id: &NativeId,
2240 request: &PageRequest,
2241) -> Result<Page<DependencyEdge>, SourceError> {
2242 match entity {
2243 Entity::Task => {
2244 source
2245 .source()
2246 .task_dependencies(id, Direction::DependsOn, request)
2247 .await
2248 }
2249 Entity::Project => {
2250 source
2251 .source()
2252 .project_dependencies(id, Direction::DependsOn, request)
2253 .await
2254 }
2255 }
2256}
2257
2258#[expect(
2267 clippy::too_many_arguments,
2268 reason = "every argument is one axis of one walk — the source, the item, the \
2269 direction, which of its two graphs, whether the reverse is emulated, where \
2270 to resume, how many rows to return and where to count calls. Grouping them \
2271 into a struct would name the same eight values one indirection further from \
2272 the loop that reads them."
2273)]
2274async fn fetch_edges(
2275 source: &ResolvedSource,
2276 native: &NativeId,
2277 direction: Direction,
2278 entity: Entity,
2279 emulating: bool,
2280 start: &Resume,
2281 budget: u32,
2282 calls: &AtomicU32,
2283) -> Result<Fetched<DependencyEdge>, SourceError> {
2284 let ceiling = ceiling(source);
2285 if !emulating {
2286 return walk(
2287 start,
2288 budget,
2289 page_size(false, budget, ceiling),
2290 |_| true,
2291 |cursor, limit| async move {
2292 calls.fetch_add(1, Ordering::Relaxed);
2293 let request = PageRequest { cursor, limit };
2294 match entity {
2295 Entity::Task => {
2296 source
2297 .source()
2298 .task_dependencies(native, direction, &request)
2299 .await
2300 }
2301 Entity::Project => {
2302 source
2303 .source()
2304 .project_dependencies(native, direction, &request)
2305 .await
2306 }
2307 }
2308 },
2309 )
2310 .await;
2311 }
2312
2313 walk(
2314 start,
2315 budget,
2316 ceiling,
2317 |_| true,
2318 |cursor, limit| async move {
2319 calls.fetch_add(1, Ordering::Relaxed);
2320 let request = PageRequest { cursor, limit };
2321 let (ids, next) = match entity {
2322 Entity::Task => {
2323 let page = source
2324 .source()
2325 .query_tasks(&TaskQuery::default(), &request)
2326 .await?;
2327 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2328 (ids, page.next)
2329 }
2330 Entity::Project => {
2331 let page = source
2332 .source()
2333 .query_projects(&ProjectQuery::default(), &request)
2334 .await?;
2335 let ids: Vec<NativeId> =
2336 page.items.into_iter().map(|project| project.id).collect();
2337 (ids, page.next)
2338 }
2339 };
2340
2341 let mut edges = Vec::new();
2342 for id in ids {
2343 let mut inner: Option<Cursor> = None;
2344 loop {
2345 calls.fetch_add(1, Ordering::Relaxed);
2346 let request = PageRequest {
2347 cursor: inner.clone(),
2348 limit,
2349 };
2350 let page = forward_edges(source, entity, &id, &request).await?;
2351 fits(page.items.len(), limit)?;
2355 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2356 unrepeated(
2357 page.next.as_ref(),
2358 inner.as_ref(),
2359 "its forward edges were being scanned",
2360 )?;
2361 match page.next {
2362 Some(cursor) => inner = Some(cursor),
2363 None => break,
2364 }
2365 }
2366 }
2367
2368 Ok(Page { items: edges, next })
2369 },
2370 )
2371 .await
2372}