1mod comment;
19mod copy;
20mod delivery;
21mod fetch;
22mod join;
23mod local;
24mod metadata;
25mod resume;
26
27use std::collections::BTreeMap;
28use std::num::NonZeroU32;
29use std::sync::atomic::{AtomicU32, Ordering};
30
31use onetaskgraph_plugin_api::{
32 Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
33 MetadataRecord, NativeId, Page, PageRequest, Project, ProjectFilter, ProjectQuery,
34 SecretResolver, SourceError, SourceName, StatusCategory, Task, TaskQuery, TextFields,
35 TextQuery,
36};
37use schemars::JsonSchema;
38use serde::{Deserialize, Serialize};
39
40use crate::GlobalId;
41use crate::config::Config;
42use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
43use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
44
45use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
46use join::join_all;
47use local::{LocalDocuments, LocalProjects, LocalTasks};
48pub(crate) use resume::{Owed, Resumption, StreamState};
49use resume::{Resume, StreamKind};
50
51pub use comment::{CommentList, DeletedComment, TaskDetail};
52pub use copy::{
53 BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
54 Spent,
55};
56pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
57pub use local::ProjectSelector;
58pub use metadata::MetadataSet;
59
60#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
65pub struct Qualified<T> {
66 pub id: GlobalId,
68 pub item: T,
70}
71
72#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
77pub struct QualifiedEdge {
78 pub from: QualifiedEndpoint,
83 pub to: QualifiedEndpoint,
85 pub kind: onetaskgraph_plugin_api::DependencyKind,
87}
88
89#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
91pub struct QualifiedEndpoint {
92 pub id: GlobalId,
94 pub kind: onetaskgraph_plugin_api::ItemKind,
96}
97
98impl std::fmt::Display for QualifiedEndpoint {
99 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
100 self.id.fmt(formatter)
101 }
102}
103
104#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
106#[serde(tag = "kind", rename_all = "kebab-case")]
107pub enum SearchHit {
108 Task(Qualified<Task>),
110 Project(Qualified<Project>),
112}
113
114#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
116#[serde(rename_all = "kebab-case")]
117pub enum SearchKind {
118 Tasks,
120 Projects,
122 #[default]
124 Both,
125}
126
127#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
129pub struct SourceListing {
130 pub source: SourceName,
132 pub kind: String,
144 #[serde(flatten)]
146 pub state: SourceState,
147}
148
149#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
151#[serde(tag = "state", rename_all = "kebab-case")]
152pub enum SourceState {
153 Available {
155 capabilities: Capabilities,
157 },
158 Unavailable {
160 error: SourceError,
162 },
163}
164
165#[derive(Debug, Clone, PartialEq)]
167pub struct Paging {
168 pub limit: NonZeroU32,
170 pub token: Option<PageToken>,
172}
173
174#[derive(Debug, Clone, Default, PartialEq)]
176pub struct Filters {
177 pub text: Option<TextQuery>,
179 pub labels: LabelFilter,
181 pub statuses: Vec<StatusCategory>,
183}
184
185#[derive(Debug, Clone)]
187pub struct TaskRequest {
188 pub sources: Vec<SourceName>,
190 pub filters: Filters,
192 pub project: ProjectSelector,
194 pub paging: Paging,
196}
197
198#[derive(Debug, Clone)]
200pub struct ProjectRequest {
201 pub sources: Vec<SourceName>,
203 pub filters: Filters,
205 pub paging: Paging,
207}
208
209#[derive(Debug, Clone, Default, PartialEq)]
216pub struct DocumentFilters {
217 pub text: Option<TextQuery>,
219 pub labels: LabelFilter,
221}
222
223#[derive(Debug, Clone)]
225pub struct DocumentRequest {
226 pub sources: Vec<SourceName>,
228 pub filters: DocumentFilters,
230 pub project: ProjectSelector,
232 pub paging: Paging,
234}
235
236#[derive(Debug, Clone)]
238pub struct LabelRequest {
239 pub sources: Vec<SourceName>,
241 pub paging: Paging,
243}
244
245#[derive(Debug, Clone)]
247pub struct SearchRequest {
248 pub sources: Vec<SourceName>,
250 pub text: TextQuery,
252 pub kind: SearchKind,
254 pub paging: Paging,
256}
257
258#[derive(Debug, Clone)]
260pub struct DependencyRequest {
261 pub id: GlobalId,
263 pub direction: Direction,
265 pub paging: Paging,
267}
268
269#[derive(Debug, Clone, PartialEq, thiserror::Error)]
277pub enum EngineError {
278 #[error(
280 "no source named {name:?} is configured\n\
281 next: name one of the configured sources ({configured}), or add {name:?} under \
282 `sources` — `onetaskgraph sources list` shows what this configuration has."
283 )]
284 UnknownSource {
285 name: String,
287 configured: String,
289 },
290
291 #[error(
293 "{message}\n\
294 next: page with a token exactly as the previous page reported it, and against \
295 the same configuration — or drop `--page` to start the walk again."
296 )]
297 Token {
298 message: String,
300 },
301
302 #[error(
304 "no sources are configured\n\
305 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
306 prints what each plugin accepts."
307 )]
308 NoSources,
309
310 #[error(
312 "source {name} cannot be written: its plugin is {kind}, which has no write \
313 side\n\
314 next: copy into a source whose plugin can be written — `onetaskgraph sources \
315 list` reports each one's plugin."
316 )]
317 NotWritable {
318 name: String,
320 kind: String,
322 },
323
324 #[error(
331 "source {name} has no documents: its plugin is {kind}, which holds none\n\
332 next: name a source whose plugin has documents — `onetaskgraph sources list` \
333 reports each one's plugin and what it declares."
334 )]
335 NoDocuments {
336 name: String,
338 kind: String,
340 },
341
342 #[error(
347 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
348 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
349 list` reports each one's plugin."
350 )]
351 NoComments {
352 name: String,
354 kind: String,
356 },
357
358 #[error(
360 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
361 but not added to, edited or removed\n\
362 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
363 source whose plugin can be written — `onetaskgraph sources list` reports each \
364 one's plugin."
365 )]
366 CommentsNotWritable {
367 name: String,
369 kind: String,
371 },
372
373 #[error(
375 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
376 next: set the status in that source itself, or name a task of a source whose plugin \
377 can be written — `onetaskgraph sources list` reports each one's plugin."
378 )]
379 StatusNotWritable {
380 name: String,
382 kind: String,
384 },
385
386 #[error(
388 "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
389 write side\n\
390 next: set the key in that source itself, or name a {record} of a source whose plugin \
391 can be written — `onetaskgraph sources list` reports each one's plugin."
392 )]
393 MetadataNotWritable {
394 name: String,
396 kind: String,
398 record: MetadataRecord,
400 },
401
402 #[error(
404 "no project with the id {id}\n\
405 next: check the id, or list what is there — `onetaskgraph project list` reports every \
406 project the configured sources hold."
407 )]
408 NoSuchProject {
409 id: String,
411 },
412
413 #[error(
415 "no document with the id {id}\n\
416 next: check the id, or list what is there — `onetaskgraph document list` reports every \
417 document the configured sources hold."
418 )]
419 NoSuchDocument {
420 id: String,
422 },
423
424 #[error(
426 "no task with the id {id}\n\
427 next: check the id, or list what is there — `onetaskgraph task list` reports every \
428 task the configured sources hold."
429 )]
430 NoSuchTask {
431 id: String,
433 },
434
435 #[error(
437 "task {task} has no comment with the id {comment}\n\
438 next: list its comments — `onetaskgraph task comment list {task}` reports each \
439 one's id."
440 )]
441 NoSuchComment {
442 task: String,
444 comment: String,
446 },
447
448 #[error(
453 "source {name} could not be built: {error}\n\
454 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
455 command again."
456 )]
457 SourceUnavailable {
458 name: String,
460 error: SourceError,
462 },
463
464 #[error(
470 "source {name} could not do it: {error}\n\
471 next: fix what the source named above, then run the command again."
472 )]
473 SourceFailed {
474 name: String,
476 error: SourceError,
478 },
479
480 #[error(
482 "the destination source {name} could not be built: {error}\n\
483 next: fix that source — `onetaskgraph sources list` reports its state — then \
484 copy again."
485 )]
486 DestinationUnavailable {
487 name: String,
489 error: SourceError,
491 },
492
493 #[error(
495 "no item with the id {id}\n\
496 next: check the id, or list what is there — `onetaskgraph task list` and \
497 `onetaskgraph project list` report what the configured sources hold."
498 )]
499 NoSuchItem {
500 id: String,
502 },
503
504 #[error(
509 "{item} was copied from {origin}, which that destination no longer holds\n\
510 next: re-run with --recreate to create a new item there instead, or restore \
511 {origin}."
512 )]
513 StaleOrigin {
514 item: String,
516 origin: String,
518 },
519
520 #[error(
526 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
527 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
528 {id} on its own with `onetaskgraph task copy`.",
529 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
530 )]
531 NotAMember {
532 id: GlobalId,
534 projects: Vec<GlobalId>,
536 },
537
538 #[error(
547 "{item} depends on {member}, which this copy was not told to carry and which records \
548 no origin in {destination}\n\
549 next: name {member} with --member as well, record its {destination} id at \
550 onetaskgraph.origin, or copy the whole project without --member."
551 )]
552 UnrecordedMember {
553 item: GlobalId,
555 member: GlobalId,
557 destination: SourceName,
559 },
560
561 #[error(
567 "source {name} could not do it: {error}\n\
568 next: fix what the source named above, then copy again."
569 )]
570 SourceRefused {
571 name: String,
573 error: SourceError,
575 },
576
577 #[error(
585 "the copy failed and could not be undone.\n\
586 it failed because: {error}\n\
587 it could not be undone because: {refusal}\n\
588 so the destination still holds: {left_behind}\n\
589 next: remove those items at the destination, then copy again."
590 )]
591 CopyNotUndone {
592 error: Box<EngineError>,
594 left_behind: LeftBehind,
600 refusal: SourceError,
602 },
603}
604
605#[derive(Debug, Clone, PartialEq)]
614pub struct LeftBehind {
615 first: GlobalId,
617 rest: Vec<GlobalId>,
619}
620
621impl LeftBehind {
622 #[must_use]
624 pub fn new(first: GlobalId) -> Self {
625 Self {
626 first,
627 rest: Vec::new(),
628 }
629 }
630
631 pub fn push(&mut self, id: GlobalId) {
633 self.rest.push(id);
634 }
635
636 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
638 std::iter::once(&self.first).chain(self.rest.iter())
639 }
640}
641
642impl std::fmt::Display for LeftBehind {
643 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
645 write!(formatter, "{}", self.first)?;
646 for id in &self.rest {
647 write!(formatter, ", {id}")?;
648 }
649 Ok(())
650 }
651}
652
653pub enum ConfiguredSource {
660 Ready(ResolvedSource),
662 Unavailable(UnavailableSource),
664}
665
666impl ConfiguredSource {
667 #[must_use]
669 pub fn name(&self) -> &SourceName {
670 match self {
671 Self::Ready(source) => source.name(),
672 Self::Unavailable(source) => source.name(),
673 }
674 }
675}
676
677pub struct Engine {
679 sources: Vec<ConfiguredSource>,
681 selection: Vec<SourceName>,
683}
684
685impl Engine {
686 #[must_use]
694 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
695 let (ready, unavailable) = resolve_available(config, secrets);
696 Self::new(
697 ready
698 .into_iter()
699 .map(ConfiguredSource::Ready)
700 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
701 .collect(),
702 config.selected_sources(),
703 )
704 }
705
706 #[must_use]
709 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
710 Self { sources, selection }
711 }
712
713 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
715 self.sources.iter().filter_map(|source| match source {
716 ConfiguredSource::Ready(ready) => Some(ready),
717 ConfiguredSource::Unavailable(_) => None,
718 })
719 }
720
721 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
723 self.sources.iter().filter_map(|source| match source {
724 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
725 ConfiguredSource::Ready(_) => None,
726 })
727 }
728
729 #[must_use]
731 pub fn listing(&self) -> Vec<SourceListing> {
732 let mut listings: Vec<SourceListing> = self
733 .ready()
734 .map(|source| SourceListing {
735 source: source.name().clone(),
736 kind: source.kind().to_owned(),
737 state: SourceState::Available {
738 capabilities: source.source().capabilities(),
739 },
740 })
741 .chain(self.unavailable().map(|source| SourceListing {
742 source: source.name().clone(),
743 kind: source.kind().to_owned(),
744 state: SourceState::Unavailable {
745 error: source.error().clone(),
746 },
747 }))
748 .collect();
749 listings.sort_by(|left, right| left.source.cmp(&right.source));
750 listings
751 }
752
753 #[must_use]
759 pub fn has(&self, name: &SourceName) -> bool {
760 self.sources.iter().any(|source| source.name() == name)
761 }
762
763 pub async fn tasks(
771 &self,
772 request: &TaskRequest,
773 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
774 let mut names = self.resolve_selection(&request.sources)?;
775 if let ProjectSelector::Qualified(id) = &request.project {
780 self.known(&id.source)?;
781 names.retain(|name| name == &id.source);
782 }
783 let query = shape("task-list", &names, &(&request.filters, &request.project));
784 let states = resumption(
785 self,
786 request.paging.token.as_ref(),
787 &[StreamKind::Items],
788 &query,
789 )?;
790 let budget = request.paging.limit.get();
791
792 let mut answer = Answer::new();
793 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
794
795 let shapes: Vec<TaskShape> = ready
796 .iter()
797 .map(|source| {
798 shape_tasks(
799 &source.source().capabilities(),
800 &request.filters,
801 &project_filter(&request.project),
802 )
803 })
804 .collect();
805 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
806 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
807
808 let walks = ready
809 .iter()
810 .enumerate()
811 .map(|(index, source)| {
812 fetch_tasks(
813 source,
814 &shapes[index],
815 &starts[index],
816 budget,
817 &counters[index],
818 )
819 })
820 .collect();
821
822 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
823 answer.finish(
824 streams,
825 budget,
826 owed(&states),
827 &query,
828 |name, task: Task| {
829 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
830 },
831 )
832 }
833
834 pub async fn projects(
840 &self,
841 request: &ProjectRequest,
842 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
843 let names = self.resolve_selection(&request.sources)?;
844 let query = shape("project-list", &names, &request.filters);
845 let states = resumption(
846 self,
847 request.paging.token.as_ref(),
848 &[StreamKind::Items],
849 &query,
850 )?;
851 let budget = request.paging.limit.get();
852
853 let mut answer = Answer::new();
854 let mut with_projects = Vec::new();
860 for source in answer.split(self, &names) {
861 if source.source().capabilities().projects.is_native() {
862 with_projects.push(source);
863 } else {
864 answer.unreachable_predicate(source, Predicate::Project);
865 }
866 }
867 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
868
869 let shapes: Vec<ProjectShape> = ready
870 .iter()
871 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
872 .collect();
873 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
874 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
875
876 let walks = ready
877 .iter()
878 .enumerate()
879 .map(|(index, source)| {
880 fetch_projects(
881 source,
882 &shapes[index],
883 &starts[index],
884 budget,
885 &counters[index],
886 )
887 })
888 .collect();
889
890 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
891 answer.finish(
892 streams,
893 budget,
894 owed(&states),
895 &query,
896 |name, project: Project| Qualified {
897 id: GlobalId::new(name.clone(), project.id.clone()),
898 item: project,
899 },
900 )
901 }
902
903 pub async fn documents(
915 &self,
916 request: &DocumentRequest,
917 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
918 let mut names = self.resolve_selection(&request.sources)?;
919 if let ProjectSelector::Qualified(id) = &request.project {
922 self.known(&id.source)?;
923 names.retain(|name| name == &id.source);
924 }
925 let query = shape(
926 "document-list",
927 &names,
928 &(&request.filters, &request.project),
929 );
930 let states = resumption(
931 self,
932 request.paging.token.as_ref(),
933 &[StreamKind::Items],
934 &query,
935 )?;
936 let budget = request.paging.limit.get();
937
938 let mut answer = Answer::new();
939 let mut with_documents = Vec::new();
940 for source in answer.split(self, &names) {
941 if source.source().capabilities().documents.is_native() {
942 with_documents.push(source);
943 } else {
944 answer.unreachable_predicate(source, Predicate::Document);
945 }
946 }
947 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
948
949 let shapes: Vec<DocumentShape> = ready
950 .iter()
951 .map(|source| {
952 shape_documents(
953 &source.source().capabilities(),
954 &request.filters,
955 &project_filter(&request.project),
956 )
957 })
958 .collect();
959 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
960 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
961
962 let walks = ready
963 .iter()
964 .enumerate()
965 .map(|(index, source)| {
966 fetch_documents(
967 source,
968 &shapes[index],
969 &starts[index],
970 budget,
971 &counters[index],
972 )
973 })
974 .collect();
975
976 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
977 answer.finish(
978 streams,
979 budget,
980 owed(&states),
981 &query,
982 |name, document: Document| Qualified {
983 id: GlobalId::new(name.clone(), document.id.clone()),
984 item: document,
985 },
986 )
987 }
988
989 pub async fn labels(
995 &self,
996 request: &LabelRequest,
997 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
998 let names = self.resolve_selection(&request.sources)?;
999 let query = shape("label-list", &names, &());
1000 let states = resumption(
1001 self,
1002 request.paging.token.as_ref(),
1003 &[StreamKind::Items],
1004 &query,
1005 )?;
1006 let budget = request.paging.limit.get();
1007
1008 let mut answer = Answer::new();
1009 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1010 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1011 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1012
1013 let walks = ready
1014 .iter()
1015 .enumerate()
1016 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1017 .collect();
1018
1019 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1020 answer.finish(
1021 streams,
1022 budget,
1023 owed(&states),
1024 &query,
1025 |name, label: Label| Qualified {
1026 id: GlobalId::new(name.clone(), label.id.clone()),
1027 item: label,
1028 },
1029 )
1030 }
1031
1032 pub async fn search(
1038 &self,
1039 request: &SearchRequest,
1040 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1041 let names = self.resolve_selection(&request.sources)?;
1042 let reads: &[StreamKind] = match request.kind {
1046 SearchKind::Tasks => &[StreamKind::Tasks],
1047 SearchKind::Projects => &[StreamKind::Projects],
1048 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1049 };
1050 let query = shape("search", &names, &request.text);
1055 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1056 let budget = request.paging.limit.get();
1057 let filters = Filters {
1058 text: Some(request.text.clone()),
1059 ..Filters::default()
1060 };
1061
1062 let mut answer = Answer::new();
1063
1064 let mut ready = Vec::new();
1067 let mut kinds = Vec::new();
1068 let mut starts = Vec::new();
1069 for source in answer.split(self, &names) {
1070 let mut streams = Vec::new();
1071 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1072 streams.push(StreamKind::Tasks);
1073 }
1074 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1075 if source.source().capabilities().projects.is_native() {
1076 streams.push(StreamKind::Projects);
1077 } else {
1078 answer.unreachable_predicate(source, Predicate::Project);
1079 }
1080 }
1081 for stream in streams {
1082 if let Some(resume) = resume_at(&states, source.name(), stream) {
1083 ready.push(source);
1084 kinds.push(stream);
1085 starts.push(resume);
1086 }
1087 }
1088 }
1089
1090 let shapes: Vec<HitShape> = ready
1091 .iter()
1092 .zip(kinds.iter())
1093 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1094 .collect();
1095 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1096 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1097
1098 let walks = ready
1099 .iter()
1100 .enumerate()
1101 .map(|(index, source)| {
1102 fetch_hits(
1103 source,
1104 &shapes[index],
1105 &starts[index],
1106 budget,
1107 &counters[index],
1108 )
1109 })
1110 .collect();
1111
1112 let streams =
1113 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1114 answer.finish(
1115 streams,
1116 budget,
1117 owed(&states),
1118 &query,
1119 |name, found: Found| match found {
1120 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1121 GlobalId::new(name.clone(), task.id.clone()),
1122 task,
1123 )),
1124 Found::Project(project) => SearchHit::Project(Qualified {
1125 id: GlobalId::new(name.clone(), project.id.clone()),
1126 item: project,
1127 }),
1128 },
1129 )
1130 }
1131
1132 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1139 let name = self.known(&id.source)?;
1140 let mut answer = Answer::new();
1141 let selected = answer.split(self, std::slice::from_ref(&name));
1142 let Some(source) = selected.first() else {
1143 return answer.nothing();
1144 };
1145 let found = source.source().get_task(&id.native).await;
1146 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1147 answer.one(source, found, |task| {
1148 delivery::qualified_task(qualified, task)
1149 })
1150 }
1151
1152 pub async fn project(
1158 &self,
1159 id: &GlobalId,
1160 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1161 let name = self.known(&id.source)?;
1162 let mut answer = Answer::new();
1163 let selected = answer.split(self, std::slice::from_ref(&name));
1164 let Some(source) = selected.first() else {
1165 return answer.nothing();
1166 };
1167 let found = source.source().get_project(&id.native).await;
1168 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1169 answer.one(source, found, |project| Qualified {
1170 id: qualified,
1171 item: project,
1172 })
1173 }
1174
1175 pub async fn document(
1185 &self,
1186 id: &GlobalId,
1187 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1188 let name = self.known(&id.source)?;
1189 let mut answer = Answer::new();
1190 let selected = answer.split(self, std::slice::from_ref(&name));
1191 let Some(source) = selected.first() else {
1192 return answer.nothing();
1193 };
1194 if !source.source().capabilities().documents.is_native() {
1195 answer.unreachable_predicate(source, Predicate::Document);
1196 return answer.nothing();
1197 }
1198 let found = source.source().get_document(&id.native).await;
1199 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1200 answer.one(source, found, |document| Qualified {
1201 id: qualified,
1202 item: document,
1203 })
1204 }
1205
1206 pub async fn task_dependencies(
1213 &self,
1214 request: &DependencyRequest,
1215 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1216 self.dependencies(request, Entity::Task).await
1217 }
1218
1219 pub async fn project_dependencies(
1225 &self,
1226 request: &DependencyRequest,
1227 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1228 self.dependencies(request, Entity::Project).await
1229 }
1230
1231 async fn dependencies(
1234 &self,
1235 request: &DependencyRequest,
1236 entity: Entity,
1237 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1238 let name = self.known(&request.id.source)?;
1239 let query = shape(
1240 "dependencies",
1241 std::slice::from_ref(&name),
1242 &(entity, &request.id.native, request.direction),
1243 );
1244 let states = resumption(
1245 self,
1246 request.paging.token.as_ref(),
1247 &[StreamKind::Items],
1248 &query,
1249 )?;
1250 let budget = request.paging.limit.get();
1251
1252 let mut answer = Answer::new();
1253 let (ready, starts) = walking(
1254 answer.split(self, std::slice::from_ref(&name)),
1255 &states,
1256 StreamKind::Items,
1257 );
1258 let Some(source) = ready.first() else {
1259 return answer.nothing();
1260 };
1261
1262 let capabilities = source.source().capabilities();
1263 let support = match entity {
1264 Entity::Task => capabilities.task_dependencies,
1265 Entity::Project => capabilities.project_dependencies,
1266 };
1267 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1271 let mut outcomes = Outcomes::default();
1272 if request.direction == Direction::DependedOnBy {
1273 if emulating {
1274 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1275 } else {
1276 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1277 }
1278 }
1279
1280 let counters = vec![AtomicU32::new(0)];
1281 let walked = fetch_edges(
1282 source,
1283 &request.id.native,
1284 request.direction,
1285 entity,
1286 emulating,
1287 &starts[0],
1288 budget,
1289 &counters[0],
1290 )
1291 .await;
1292
1293 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1294 answer.finish(
1295 streams,
1296 budget,
1297 owed(&states),
1298 &query,
1299 |name, edge: DependencyEdge| QualifiedEdge {
1300 from: qualify_endpoint(name, edge.from),
1301 to: qualify_endpoint(name, edge.to),
1302 kind: edge.kind,
1303 },
1304 )
1305 }
1306
1307 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1309 if asked.is_empty() {
1310 if self.selection.is_empty() {
1311 return Err(EngineError::NoSources);
1312 }
1313 return Ok(self.selection.clone());
1314 }
1315 asked.iter().map(|name| self.known(name)).collect()
1316 }
1317
1318 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1320 if self.has(name) {
1321 return Ok(name.clone());
1322 }
1323 if self.sources.is_empty() {
1324 return Err(EngineError::NoSources);
1325 }
1326 Err(EngineError::UnknownSource {
1327 name: name.to_string(),
1328 configured: self
1329 .listing()
1330 .iter()
1331 .map(|listing| listing.source.to_string())
1332 .collect::<Vec<_>>()
1333 .join(", "),
1334 })
1335 }
1336}
1337
1338fn qualify_endpoint(
1339 source: &SourceName,
1340 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1341) -> QualifiedEndpoint {
1342 let kind = endpoint.kind;
1343 let is_qualified = endpoint.is_qualified();
1344 let endpoint_id = endpoint.into_id();
1345 QualifiedEndpoint {
1346 id: if is_qualified {
1347 endpoint_id
1348 .parse()
1349 .expect("plugin-api validates qualified dependency endpoints")
1350 } else {
1351 GlobalId::new(source.clone(), NativeId(endpoint_id))
1352 },
1353 kind,
1354 }
1355}
1356
1357#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1359enum Entity {
1360 Task,
1362 Project,
1364}
1365
1366enum Found {
1368 Task(Task),
1370 Project(Project),
1372}
1373
1374#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1376enum Outcome {
1377 PushedDown,
1379 AppliedLocally,
1381 Emulated,
1383 Unavailable,
1385}
1386
1387#[derive(Debug, Clone, Default, PartialEq)]
1399struct Outcomes(BTreeMap<Predicate, Outcome>);
1400
1401impl Outcomes {
1402 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1408 self.0.insert(predicate, outcome);
1409 }
1410
1411 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1414 for predicate in predicates {
1415 self.record(predicate, outcome);
1416 }
1417 }
1418
1419 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1421 self.0
1422 .iter()
1423 .filter(|(_, recorded)| **recorded == outcome)
1424 .map(|(predicate, _)| *predicate)
1425 .collect()
1426 }
1427}
1428
1429struct TaskShape {
1431 pushed: TaskQuery,
1433 local: LocalTasks,
1435 outcomes: Outcomes,
1437}
1438
1439struct ProjectShape {
1441 pushed: ProjectQuery,
1443 local: LocalProjects,
1445 outcomes: Outcomes,
1447}
1448
1449struct DocumentShape {
1451 pushed: DocumentQuery,
1453 local: LocalDocuments,
1455 outcomes: Outcomes,
1457}
1458
1459struct HitShape {
1461 stream: StreamKind,
1463 tasks: TaskQuery,
1465 projects: ProjectQuery,
1467 local_tasks: LocalTasks,
1469 local_projects: LocalProjects,
1471 outcomes: Outcomes,
1473}
1474
1475struct Answer {
1481 plans: Vec<SourcePlan>,
1483 errors: Vec<SourceFailure>,
1485}
1486
1487impl Answer {
1488 fn new() -> Self {
1489 Self {
1490 plans: Vec::new(),
1491 errors: Vec::new(),
1492 }
1493 }
1494
1495 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1500 let mut selected = Vec::new();
1501 for name in names {
1502 match engine.sources.iter().find(|source| source.name() == name) {
1503 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1504 Some(ConfiguredSource::Unavailable(source)) => {
1505 self.errors.push(source.failure());
1506 }
1507 None => {}
1508 }
1509 }
1510 selected
1511 }
1512
1513 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1515 let mut outcomes = Outcomes::default();
1516 outcomes.record(predicate, Outcome::Unavailable);
1517 self.plans.push(plan_for(source, outcomes, 0));
1518 }
1519
1520 fn collect<T>(
1522 &mut self,
1523 ready: &[&ResolvedSource],
1524 walked: Vec<Result<Fetched<T>, SourceError>>,
1525 counters: &[AtomicU32],
1526 outcomes: Vec<Outcomes>,
1527 ) -> Vec<Stream<T>> {
1528 let kinds = vec![StreamKind::Items; ready.len()];
1529 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1530 }
1531
1532 fn collect_streams<T>(
1534 &mut self,
1535 ready: &[&ResolvedSource],
1536 kinds: &[StreamKind],
1537 walked: Vec<Result<Fetched<T>, SourceError>>,
1538 counters: &[AtomicU32],
1539 outcomes: Vec<Outcomes>,
1540 ) -> Vec<Stream<T>> {
1541 let mut streams = Vec::new();
1542 for (index, result) in walked.into_iter().enumerate() {
1543 let source = ready[index];
1544 let pages = counters[index].load(Ordering::Relaxed);
1545 self.plans
1546 .push(plan_for(source, outcomes[index].clone(), pages));
1547 match result {
1548 Ok(fetched) => streams.push(Stream {
1549 source: source.name().clone(),
1550 kind: kinds[index],
1551 fetched,
1552 }),
1553 Err(error) => self.errors.push(SourceFailure {
1556 source: source.name().clone(),
1557 error,
1558 }),
1559 }
1560 }
1561 streams
1562 }
1563
1564 fn one<T, U>(
1566 mut self,
1567 source: &ResolvedSource,
1568 found: Result<Option<T>, SourceError>,
1569 qualify: impl FnOnce(T) -> U,
1570 ) -> Result<QueryResponse<U>, EngineError> {
1571 self.plans.push(plan_for(source, Outcomes::default(), 1));
1572 let items = match found {
1573 Ok(Some(item)) => vec![qualify(item)],
1574 Ok(None) => Vec::new(),
1575 Err(error) => {
1576 self.errors.push(SourceFailure {
1577 source: source.name().clone(),
1578 error,
1579 });
1580 Vec::new()
1581 }
1582 };
1583 Ok(QueryResponse {
1584 items,
1585 next: None,
1586 plan: QueryPlan {
1587 per_source: merge_plans(self.plans),
1588 },
1589 errors: self.errors,
1590 })
1591 }
1592
1593 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1595 Ok(QueryResponse {
1596 items: Vec::new(),
1597 next: None,
1598 plan: QueryPlan {
1599 per_source: merge_plans(self.plans),
1600 },
1601 errors: self.errors,
1602 })
1603 }
1604
1605 fn finish<T, U>(
1610 self,
1611 streams: Vec<Stream<T>>,
1612 budget: u32,
1613 first: Option<&Owed>,
1614 query: &str,
1615 qualify: impl Fn(&SourceName, T) -> U,
1616 ) -> Result<QueryResponse<U>, EngineError> {
1617 let (rows, states, owed) = merge(streams, budget, first);
1618 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1619 Ok(QueryResponse {
1620 items: rows
1621 .into_iter()
1622 .map(|(name, item)| qualify(&name, item))
1623 .collect(),
1624 next,
1625 plan: QueryPlan {
1626 per_source: merge_plans(self.plans),
1627 },
1628 errors: self.errors,
1629 })
1630 }
1631}
1632
1633fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1636 SourcePlan {
1637 source: source.name().clone(),
1638 kind: source.kind().to_owned(),
1639 pushed_down: outcomes.with(Outcome::PushedDown),
1640 applied_locally: outcomes.with(Outcome::AppliedLocally),
1641 emulated: outcomes.with(Outcome::Emulated),
1642 unavailable: outcomes.with(Outcome::Unavailable),
1643 pages_fetched: pages,
1644 }
1645}
1646
1647fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1652 let mut merged: Vec<SourcePlan> = Vec::new();
1653 for plan in plans {
1654 if let Some(existing) = merged
1655 .iter_mut()
1656 .find(|existing| existing.source == plan.source)
1657 {
1658 existing.pushed_down.extend(plan.pushed_down);
1659 existing.applied_locally.extend(plan.applied_locally);
1660 existing.emulated.extend(plan.emulated);
1661 existing.unavailable.extend(plan.unavailable);
1662 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1663 for list in [
1664 &mut existing.pushed_down,
1665 &mut existing.applied_locally,
1666 &mut existing.emulated,
1667 &mut existing.unavailable,
1668 ] {
1669 list.sort_unstable();
1670 list.dedup();
1671 }
1672 } else {
1673 merged.push(plan);
1674 }
1675 }
1676 merged
1677}
1678
1679fn walking<'a>(
1685 selected: Vec<&'a ResolvedSource>,
1686 states: &Option<Resumption>,
1687 kind: StreamKind,
1688) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1689 let mut ready = Vec::new();
1690 let mut starts = Vec::new();
1691 for source in selected {
1692 if let Some(resume) = resume_at(states, source.name(), kind) {
1693 ready.push(source);
1694 starts.push(resume);
1695 }
1696 }
1697 (ready, starts)
1698}
1699
1700fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1718 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1719 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1720}
1721
1722fn fingerprint(text: &str) -> String {
1724 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1725 for byte in text.as_bytes() {
1726 hash ^= u64::from(*byte);
1727 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1728 }
1729 format!("{hash:016x}")
1730}
1731
1732fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1737 document.as_ref()?.owed.as_ref()
1738}
1739
1740fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1742 match states {
1743 None => Some(Resume::default()),
1744 Some(document) => document
1745 .streams
1746 .iter()
1747 .find(|state| &state.source == source && state.stream == kind)
1748 .map(|state| state.resume.clone()),
1749 }
1750}
1751
1752fn resumption(
1783 engine: &Engine,
1784 token: Option<&PageToken>,
1785 reads: &[StreamKind],
1786 query: &str,
1787) -> Result<Option<Resumption>, EngineError> {
1788 let Some(document) = token.map(PageToken::decode) else {
1789 return Ok(None);
1790 };
1791
1792 if document.query != query {
1799 return Err(EngineError::Token {
1800 message: "this page token was written by a different query — resume the walk it \
1801 came from, or drop --page to start this one from the beginning"
1802 .to_owned(),
1803 });
1804 }
1805 let states = &document.streams;
1806
1807 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
1808 for state in states {
1809 if !reads.contains(&state.stream) {
1810 return Err(EngineError::Token {
1811 message: format!(
1812 "this page token resumes {}, which this command does not read — it \
1813 was written by a different query",
1814 state.stream.describe()
1815 ),
1816 });
1817 }
1818 let ceiling = engine
1819 .ready()
1820 .find(|source| source.name() == &state.source)
1821 .map(ceiling);
1822 if ceiling.is_none() && !engine.has(&state.source) {
1823 return Err(EngineError::Token {
1824 message: format!(
1825 "this page token resumes a source called {:?}, which this \
1826 configuration does not have",
1827 state.source.as_str()
1828 ),
1829 });
1830 }
1831 if let Some(ceiling) = ceiling
1832 && state.resume.skip >= ceiling
1833 {
1834 return Err(EngineError::Token {
1835 message: format!(
1836 "this page token resumes {} rows into a page of source {:?}, which \
1837 serves at most {ceiling}",
1838 state.resume.skip,
1839 state.source.as_str()
1840 ),
1841 });
1842 }
1843 if seen.contains(&(&state.source, state.stream)) {
1844 return Err(EngineError::Token {
1845 message: format!(
1846 "this page token gives source {:?} two places to resume from",
1847 state.source.as_str()
1848 ),
1849 });
1850 }
1851 seen.push((&state.source, state.stream));
1852 }
1853
1854 if let Some(owed) = &document.owed
1859 && !document
1860 .streams
1861 .iter()
1862 .any(|state| state.source == owed.source && state.stream == owed.stream)
1863 {
1864 return Err(EngineError::Token {
1865 message: format!(
1866 "this page token owes the next row to a stream it does not resume, \
1867 {:?}'s {}",
1868 owed.source.as_str(),
1869 owed.stream.describe()
1870 ),
1871 });
1872 }
1873
1874 Ok(Some(document))
1875}
1876
1877fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
1883 match selector {
1884 ProjectSelector::Any => ProjectFilter::Any,
1885 ProjectSelector::Orphans => ProjectFilter::Orphans,
1886 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
1887 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
1888 }
1889}
1890
1891fn text_predicates(fields: TextFields) -> Vec<Predicate> {
1893 match fields {
1894 TextFields::Title => vec![Predicate::SearchTitle],
1895 TextFields::Content => vec![Predicate::SearchContent],
1896 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
1897 }
1898}
1899
1900fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
1907 match fields {
1908 TextFields::Title => capabilities.search_title.is_native(),
1909 TextFields::Content => capabilities.search_content.is_native(),
1910 TextFields::TitleOrContent => {
1911 capabilities.search_title.is_native() && capabilities.search_content.is_native()
1912 }
1913 }
1914}
1915
1916fn shape_tasks(
1918 capabilities: &Capabilities,
1919 filters: &Filters,
1920 project: &ProjectFilter,
1921) -> TaskShape {
1922 let mut pushed = TaskQuery::default();
1923 let mut local = LocalTasks::default();
1924 let mut outcomes = Outcomes::default();
1925
1926 if !filters.labels.is_empty() {
1927 if capabilities.filter_by_label.is_native() {
1928 pushed.labels = filters.labels.clone();
1929 outcomes.record(Predicate::Label, Outcome::PushedDown);
1930 } else {
1931 local.labels = Some(filters.labels.clone());
1932 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1933 }
1934 }
1935 if !filters.statuses.is_empty() {
1936 if capabilities.filter_by_status.is_native() {
1937 pushed.statuses.clone_from(&filters.statuses);
1938 outcomes.record(Predicate::Status, Outcome::PushedDown);
1939 } else {
1940 local.statuses.clone_from(&filters.statuses);
1941 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1942 }
1943 }
1944 if let Some(text) = &filters.text {
1945 let predicates = text_predicates(text.fields);
1946 if searches_natively(capabilities, text.fields) {
1947 pushed.text = Some(text.clone());
1948 outcomes.record_all(predicates, Outcome::PushedDown);
1949 } else {
1950 local.text = Some(text.clone());
1951 outcomes.record_all(predicates, Outcome::AppliedLocally);
1952 }
1953 }
1954 match project {
1955 ProjectFilter::Any => {}
1956 ProjectFilter::Orphans => {
1957 if capabilities.orphan_tasks.is_native() {
1958 pushed.project = ProjectFilter::Orphans;
1959 outcomes.record(Predicate::Project, Outcome::PushedDown);
1960 } else {
1961 local.project = Some(ProjectFilter::Orphans);
1962 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1963 }
1964 }
1965 ProjectFilter::Is(id) => {
1966 if capabilities.projects.is_native() {
1967 pushed.project = ProjectFilter::Is(id.clone());
1968 outcomes.record(Predicate::Project, Outcome::PushedDown);
1969 } else {
1970 local.project = Some(ProjectFilter::Is(id.clone()));
1971 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1972 }
1973 }
1974 }
1975
1976 TaskShape {
1977 pushed,
1978 local,
1979 outcomes,
1980 }
1981}
1982
1983fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
1985 let mut pushed = ProjectQuery::default();
1986 let mut local = LocalProjects::default();
1987 let mut outcomes = Outcomes::default();
1988
1989 if !filters.labels.is_empty() {
1990 if capabilities.filter_by_label.is_native() {
1991 pushed.labels = filters.labels.clone();
1992 outcomes.record(Predicate::Label, Outcome::PushedDown);
1993 } else {
1994 local.labels = Some(filters.labels.clone());
1995 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1996 }
1997 }
1998 if !filters.statuses.is_empty() {
1999 if capabilities.filter_by_status.is_native() {
2000 pushed.statuses.clone_from(&filters.statuses);
2001 outcomes.record(Predicate::Status, Outcome::PushedDown);
2002 } else {
2003 local.statuses.clone_from(&filters.statuses);
2004 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2005 }
2006 }
2007 if let Some(text) = &filters.text {
2008 let predicates = text_predicates(text.fields);
2009 if searches_natively(capabilities, text.fields) {
2010 pushed.text = Some(text.clone());
2011 outcomes.record_all(predicates, Outcome::PushedDown);
2012 } else {
2013 local.text = Some(text.clone());
2014 outcomes.record_all(predicates, Outcome::AppliedLocally);
2015 }
2016 }
2017
2018 ProjectShape {
2019 pushed,
2020 local,
2021 outcomes,
2022 }
2023}
2024
2025fn shape_documents(
2030 capabilities: &Capabilities,
2031 filters: &DocumentFilters,
2032 project: &ProjectFilter,
2033) -> DocumentShape {
2034 let mut pushed = DocumentQuery::default();
2035 let mut local = LocalDocuments::default();
2036 let mut outcomes = Outcomes::default();
2037
2038 if !filters.labels.is_empty() {
2039 if capabilities.filter_by_label.is_native() {
2040 pushed.labels = filters.labels.clone();
2041 outcomes.record(Predicate::Label, Outcome::PushedDown);
2042 } else {
2043 local.labels = Some(filters.labels.clone());
2044 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2045 }
2046 }
2047 if let Some(text) = &filters.text {
2048 let predicates = text_predicates(text.fields);
2049 if searches_natively(capabilities, text.fields) {
2050 pushed.text = Some(text.clone());
2051 outcomes.record_all(predicates, Outcome::PushedDown);
2052 } else {
2053 local.text = Some(text.clone());
2054 outcomes.record_all(predicates, Outcome::AppliedLocally);
2055 }
2056 }
2057 match project {
2058 ProjectFilter::Any => {}
2059 ProjectFilter::Orphans => {
2060 if capabilities.orphan_tasks.is_native() {
2061 pushed.project = ProjectFilter::Orphans;
2062 outcomes.record(Predicate::Project, Outcome::PushedDown);
2063 } else {
2064 local.project = Some(ProjectFilter::Orphans);
2065 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2066 }
2067 }
2068 ProjectFilter::Is(id) => {
2069 if capabilities.projects.is_native() {
2070 pushed.project = ProjectFilter::Is(id.clone());
2071 outcomes.record(Predicate::Project, Outcome::PushedDown);
2072 } else {
2073 local.project = Some(ProjectFilter::Is(id.clone()));
2074 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2075 }
2076 }
2077 }
2078
2079 DocumentShape {
2080 pushed,
2081 local,
2082 outcomes,
2083 }
2084}
2085
2086fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2088 match stream {
2089 StreamKind::Projects => {
2090 let shaped = shape_projects(capabilities, filters);
2091 HitShape {
2092 stream,
2093 tasks: TaskQuery::default(),
2094 projects: shaped.pushed,
2095 local_tasks: LocalTasks::default(),
2096 local_projects: shaped.local,
2097 outcomes: shaped.outcomes,
2098 }
2099 }
2100 StreamKind::Items | StreamKind::Tasks => {
2101 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any);
2102 HitShape {
2103 stream,
2104 tasks: shaped.pushed,
2105 projects: ProjectQuery::default(),
2106 local_tasks: shaped.local,
2107 local_projects: LocalProjects::default(),
2108 outcomes: shaped.outcomes,
2109 }
2110 }
2111 }
2112}
2113
2114fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2121 if compensating {
2122 ceiling
2123 } else {
2124 budget.min(ceiling)
2125 }
2126}
2127
2128fn ceiling(source: &ResolvedSource) -> u32 {
2130 source.source().capabilities().max_page_size.max(1)
2131}
2132
2133async fn fetch_tasks(
2135 source: &ResolvedSource,
2136 shape: &TaskShape,
2137 start: &Resume,
2138 budget: u32,
2139 calls: &AtomicU32,
2140) -> Result<Fetched<Task>, SourceError> {
2141 let compensating = shape.local != LocalTasks::default();
2142 walk(
2143 start,
2144 budget,
2145 page_size(compensating, budget, ceiling(source)),
2146 |task| shape.local.keeps(task),
2147 |cursor, limit| async move {
2148 calls.fetch_add(1, Ordering::Relaxed);
2149 let request = PageRequest { cursor, limit };
2150 source.source().query_tasks(&shape.pushed, &request).await
2151 },
2152 )
2153 .await
2154}
2155
2156async fn fetch_projects(
2158 source: &ResolvedSource,
2159 shape: &ProjectShape,
2160 start: &Resume,
2161 budget: u32,
2162 calls: &AtomicU32,
2163) -> Result<Fetched<Project>, SourceError> {
2164 let compensating = shape.local != LocalProjects::default();
2165 walk(
2166 start,
2167 budget,
2168 page_size(compensating, budget, ceiling(source)),
2169 |project| shape.local.keeps(project),
2170 |cursor, limit| async move {
2171 calls.fetch_add(1, Ordering::Relaxed);
2172 let request = PageRequest { cursor, limit };
2173 source
2174 .source()
2175 .query_projects(&shape.pushed, &request)
2176 .await
2177 },
2178 )
2179 .await
2180}
2181
2182async fn fetch_documents(
2184 source: &ResolvedSource,
2185 shape: &DocumentShape,
2186 start: &Resume,
2187 budget: u32,
2188 calls: &AtomicU32,
2189) -> Result<Fetched<Document>, SourceError> {
2190 let compensating = shape.local != LocalDocuments::default();
2191 walk(
2192 start,
2193 budget,
2194 page_size(compensating, budget, ceiling(source)),
2195 |document| shape.local.keeps(document),
2196 |cursor, limit| async move {
2197 calls.fetch_add(1, Ordering::Relaxed);
2198 let request = PageRequest { cursor, limit };
2199 source
2200 .source()
2201 .query_documents(&shape.pushed, &request)
2202 .await
2203 },
2204 )
2205 .await
2206}
2207
2208async fn fetch_labels(
2210 source: &ResolvedSource,
2211 start: &Resume,
2212 budget: u32,
2213 calls: &AtomicU32,
2214) -> Result<Fetched<Label>, SourceError> {
2215 walk(
2216 start,
2217 budget,
2218 page_size(false, budget, ceiling(source)),
2219 |_| true,
2220 |cursor, limit| async move {
2221 calls.fetch_add(1, Ordering::Relaxed);
2222 let request = PageRequest { cursor, limit };
2223 source.source().labels(&request).await
2224 },
2225 )
2226 .await
2227}
2228
2229async fn fetch_hits(
2231 source: &ResolvedSource,
2232 shape: &HitShape,
2233 start: &Resume,
2234 budget: u32,
2235 calls: &AtomicU32,
2236) -> Result<Fetched<Found>, SourceError> {
2237 let ceiling = ceiling(source);
2238 match shape.stream {
2239 StreamKind::Projects => {
2240 let compensating = shape.local_projects != LocalProjects::default();
2241 walk(
2242 start,
2243 budget,
2244 page_size(compensating, budget, ceiling),
2245 |found| match found {
2246 Found::Project(project) => shape.local_projects.keeps(project),
2247 Found::Task(_) => true,
2248 },
2249 |cursor, limit| async move {
2250 calls.fetch_add(1, Ordering::Relaxed);
2251 let request = PageRequest { cursor, limit };
2252 let page = source
2253 .source()
2254 .query_projects(&shape.projects, &request)
2255 .await?;
2256 Ok(Page {
2257 items: page.items.into_iter().map(Found::Project).collect(),
2258 next: page.next,
2259 })
2260 },
2261 )
2262 .await
2263 }
2264 StreamKind::Items | StreamKind::Tasks => {
2265 let compensating = shape.local_tasks != LocalTasks::default();
2266 walk(
2267 start,
2268 budget,
2269 page_size(compensating, budget, ceiling),
2270 |found| match found {
2271 Found::Task(task) => shape.local_tasks.keeps(task),
2272 Found::Project(_) => true,
2273 },
2274 |cursor, limit| async move {
2275 calls.fetch_add(1, Ordering::Relaxed);
2276 let request = PageRequest { cursor, limit };
2277 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2278 Ok(Page {
2279 items: page.items.into_iter().map(Found::Task).collect(),
2280 next: page.next,
2281 })
2282 },
2283 )
2284 .await
2285 }
2286 }
2287}
2288
2289async fn forward_edges(
2291 source: &ResolvedSource,
2292 entity: Entity,
2293 id: &NativeId,
2294 request: &PageRequest,
2295) -> Result<Page<DependencyEdge>, SourceError> {
2296 match entity {
2297 Entity::Task => {
2298 source
2299 .source()
2300 .task_dependencies(id, Direction::DependsOn, request)
2301 .await
2302 }
2303 Entity::Project => {
2304 source
2305 .source()
2306 .project_dependencies(id, Direction::DependsOn, request)
2307 .await
2308 }
2309 }
2310}
2311
2312#[expect(
2321 clippy::too_many_arguments,
2322 reason = "every argument is one axis of one walk — the source, the item, the \
2323 direction, which of its two graphs, whether the reverse is emulated, where \
2324 to resume, how many rows to return and where to count calls. Grouping them \
2325 into a struct would name the same eight values one indirection further from \
2326 the loop that reads them."
2327)]
2328async fn fetch_edges(
2329 source: &ResolvedSource,
2330 native: &NativeId,
2331 direction: Direction,
2332 entity: Entity,
2333 emulating: bool,
2334 start: &Resume,
2335 budget: u32,
2336 calls: &AtomicU32,
2337) -> Result<Fetched<DependencyEdge>, SourceError> {
2338 let ceiling = ceiling(source);
2339 if !emulating {
2340 return walk(
2341 start,
2342 budget,
2343 page_size(false, budget, ceiling),
2344 |_| true,
2345 |cursor, limit| async move {
2346 calls.fetch_add(1, Ordering::Relaxed);
2347 let request = PageRequest { cursor, limit };
2348 match entity {
2349 Entity::Task => {
2350 source
2351 .source()
2352 .task_dependencies(native, direction, &request)
2353 .await
2354 }
2355 Entity::Project => {
2356 source
2357 .source()
2358 .project_dependencies(native, direction, &request)
2359 .await
2360 }
2361 }
2362 },
2363 )
2364 .await;
2365 }
2366
2367 walk(
2368 start,
2369 budget,
2370 ceiling,
2371 |_| true,
2372 |cursor, limit| async move {
2373 calls.fetch_add(1, Ordering::Relaxed);
2374 let request = PageRequest { cursor, limit };
2375 let (ids, next) = match entity {
2376 Entity::Task => {
2377 let page = source
2378 .source()
2379 .query_tasks(&TaskQuery::default(), &request)
2380 .await?;
2381 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2382 (ids, page.next)
2383 }
2384 Entity::Project => {
2385 let page = source
2386 .source()
2387 .query_projects(&ProjectQuery::default(), &request)
2388 .await?;
2389 let ids: Vec<NativeId> =
2390 page.items.into_iter().map(|project| project.id).collect();
2391 (ids, page.next)
2392 }
2393 };
2394
2395 let mut edges = Vec::new();
2396 for id in ids {
2397 let mut inner: Option<Cursor> = None;
2398 loop {
2399 calls.fetch_add(1, Ordering::Relaxed);
2400 let request = PageRequest {
2401 cursor: inner.clone(),
2402 limit,
2403 };
2404 let page = forward_edges(source, entity, &id, &request).await?;
2405 fits(page.items.len(), limit)?;
2409 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2410 unrepeated(
2411 page.next.as_ref(),
2412 inner.as_ref(),
2413 "its forward edges were being scanned",
2414 )?;
2415 match page.next {
2416 Some(cursor) => inner = Some(cursor),
2417 None => break,
2418 }
2419 }
2420 }
2421
2422 Ok(Page { items: edges, next })
2423 },
2424 )
2425 .await
2426}