1mod comment;
19mod copy;
20mod delivery;
21mod fetch;
22mod join;
23mod local;
24mod resume;
25
26use std::collections::BTreeMap;
27use std::num::NonZeroU32;
28use std::sync::atomic::{AtomicU32, Ordering};
29
30use onetaskgraph_plugin_api::{
31 Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
32 NativeId, Page, PageRequest, Project, ProjectFilter, ProjectQuery, SecretResolver, SourceError,
33 SourceName, StatusCategory, Task, TaskQuery, TextFields, TextQuery,
34};
35use schemars::JsonSchema;
36use serde::{Deserialize, Serialize};
37
38use crate::GlobalId;
39use crate::config::Config;
40use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
41use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
42
43use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
44use join::join_all;
45use local::{LocalDocuments, LocalProjects, LocalTasks};
46pub(crate) use resume::{Owed, Resumption, StreamState};
47use resume::{Resume, StreamKind};
48
49pub use comment::{CommentList, DeletedComment, TaskDetail};
50pub use copy::{
51 BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
52 Spent,
53};
54pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
55pub use local::ProjectSelector;
56
57#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
62pub struct Qualified<T> {
63 pub id: GlobalId,
65 pub item: T,
67}
68
69#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
74pub struct QualifiedEdge {
75 pub from: QualifiedEndpoint,
80 pub to: QualifiedEndpoint,
82 pub kind: onetaskgraph_plugin_api::DependencyKind,
84}
85
86#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
88pub struct QualifiedEndpoint {
89 pub id: GlobalId,
91 pub kind: onetaskgraph_plugin_api::ItemKind,
93}
94
95impl std::fmt::Display for QualifiedEndpoint {
96 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
97 self.id.fmt(formatter)
98 }
99}
100
101#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
103#[serde(tag = "kind", rename_all = "kebab-case")]
104pub enum SearchHit {
105 Task(Qualified<Task>),
107 Project(Qualified<Project>),
109}
110
111#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
113#[serde(rename_all = "kebab-case")]
114pub enum SearchKind {
115 Tasks,
117 Projects,
119 #[default]
121 Both,
122}
123
124#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
126pub struct SourceListing {
127 pub source: SourceName,
129 pub kind: String,
141 #[serde(flatten)]
143 pub state: SourceState,
144}
145
146#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
148#[serde(tag = "state", rename_all = "kebab-case")]
149pub enum SourceState {
150 Available {
152 capabilities: Capabilities,
154 },
155 Unavailable {
157 error: SourceError,
159 },
160}
161
162#[derive(Debug, Clone, PartialEq)]
164pub struct Paging {
165 pub limit: NonZeroU32,
167 pub token: Option<PageToken>,
169}
170
171#[derive(Debug, Clone, Default, PartialEq)]
173pub struct Filters {
174 pub text: Option<TextQuery>,
176 pub labels: LabelFilter,
178 pub statuses: Vec<StatusCategory>,
180}
181
182#[derive(Debug, Clone)]
184pub struct TaskRequest {
185 pub sources: Vec<SourceName>,
187 pub filters: Filters,
189 pub project: ProjectSelector,
191 pub paging: Paging,
193}
194
195#[derive(Debug, Clone)]
197pub struct ProjectRequest {
198 pub sources: Vec<SourceName>,
200 pub filters: Filters,
202 pub paging: Paging,
204}
205
206#[derive(Debug, Clone, Default, PartialEq)]
213pub struct DocumentFilters {
214 pub text: Option<TextQuery>,
216 pub labels: LabelFilter,
218}
219
220#[derive(Debug, Clone)]
222pub struct DocumentRequest {
223 pub sources: Vec<SourceName>,
225 pub filters: DocumentFilters,
227 pub project: ProjectSelector,
229 pub paging: Paging,
231}
232
233#[derive(Debug, Clone)]
235pub struct LabelRequest {
236 pub sources: Vec<SourceName>,
238 pub paging: Paging,
240}
241
242#[derive(Debug, Clone)]
244pub struct SearchRequest {
245 pub sources: Vec<SourceName>,
247 pub text: TextQuery,
249 pub kind: SearchKind,
251 pub paging: Paging,
253}
254
255#[derive(Debug, Clone)]
257pub struct DependencyRequest {
258 pub id: GlobalId,
260 pub direction: Direction,
262 pub paging: Paging,
264}
265
266#[derive(Debug, Clone, PartialEq, thiserror::Error)]
274pub enum EngineError {
275 #[error(
277 "no source named {name:?} is configured\n\
278 next: name one of the configured sources ({configured}), or add {name:?} under \
279 `sources` — `onetaskgraph sources list` shows what this configuration has."
280 )]
281 UnknownSource {
282 name: String,
284 configured: String,
286 },
287
288 #[error(
290 "{message}\n\
291 next: page with a token exactly as the previous page reported it, and against \
292 the same configuration — or drop `--page` to start the walk again."
293 )]
294 Token {
295 message: String,
297 },
298
299 #[error(
301 "no sources are configured\n\
302 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
303 prints what each plugin accepts."
304 )]
305 NoSources,
306
307 #[error(
309 "source {name} cannot be written: its plugin is {kind}, which has no write \
310 side\n\
311 next: copy into a source whose plugin can be written — `onetaskgraph sources \
312 list` reports each one's plugin."
313 )]
314 NotWritable {
315 name: String,
317 kind: String,
319 },
320
321 #[error(
328 "source {name} has no documents: its plugin is {kind}, which holds none\n\
329 next: name a source whose plugin has documents — `onetaskgraph sources list` \
330 reports each one's plugin and what it declares."
331 )]
332 NoDocuments {
333 name: String,
335 kind: String,
337 },
338
339 #[error(
344 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
345 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
346 list` reports each one's plugin."
347 )]
348 NoComments {
349 name: String,
351 kind: String,
353 },
354
355 #[error(
357 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
358 but not added to, edited or removed\n\
359 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
360 source whose plugin can be written — `onetaskgraph sources list` reports each \
361 one's plugin."
362 )]
363 CommentsNotWritable {
364 name: String,
366 kind: String,
368 },
369
370 #[error(
372 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
373 next: set the status in that source itself, or name a task of a source whose plugin \
374 can be written — `onetaskgraph sources list` reports each one's plugin."
375 )]
376 StatusNotWritable {
377 name: String,
379 kind: String,
381 },
382
383 #[error(
385 "no task with the id {id}\n\
386 next: check the id, or list what is there — `onetaskgraph task list` reports every \
387 task the configured sources hold."
388 )]
389 NoSuchTask {
390 id: String,
392 },
393
394 #[error(
396 "task {task} has no comment with the id {comment}\n\
397 next: list its comments — `onetaskgraph task comment list {task}` reports each \
398 one's id."
399 )]
400 NoSuchComment {
401 task: String,
403 comment: String,
405 },
406
407 #[error(
412 "source {name} could not be built: {error}\n\
413 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
414 command again."
415 )]
416 SourceUnavailable {
417 name: String,
419 error: SourceError,
421 },
422
423 #[error(
429 "source {name} could not do it: {error}\n\
430 next: fix what the source named above, then run the command again."
431 )]
432 SourceFailed {
433 name: String,
435 error: SourceError,
437 },
438
439 #[error(
441 "the destination source {name} could not be built: {error}\n\
442 next: fix that source — `onetaskgraph sources list` reports its state — then \
443 copy again."
444 )]
445 DestinationUnavailable {
446 name: String,
448 error: SourceError,
450 },
451
452 #[error(
454 "no item with the id {id}\n\
455 next: check the id, or list what is there — `onetaskgraph task list` and \
456 `onetaskgraph project list` report what the configured sources hold."
457 )]
458 NoSuchItem {
459 id: String,
461 },
462
463 #[error(
468 "{item} was copied from {origin}, which that destination no longer holds\n\
469 next: re-run with --recreate to create a new item there instead, or restore \
470 {origin}."
471 )]
472 StaleOrigin {
473 item: String,
475 origin: String,
477 },
478
479 #[error(
485 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
486 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
487 {id} on its own with `onetaskgraph task copy`.",
488 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
489 )]
490 NotAMember {
491 id: GlobalId,
493 projects: Vec<GlobalId>,
495 },
496
497 #[error(
506 "{item} depends on {member}, which this copy was not told to carry and which records \
507 no origin in {destination}\n\
508 next: name {member} with --member as well, record its {destination} id at \
509 onetaskgraph.origin, or copy the whole project without --member."
510 )]
511 UnrecordedMember {
512 item: GlobalId,
514 member: GlobalId,
516 destination: SourceName,
518 },
519
520 #[error(
526 "source {name} could not do it: {error}\n\
527 next: fix what the source named above, then copy again."
528 )]
529 SourceRefused {
530 name: String,
532 error: SourceError,
534 },
535
536 #[error(
544 "the copy failed and could not be undone.\n\
545 it failed because: {error}\n\
546 it could not be undone because: {refusal}\n\
547 so the destination still holds: {left_behind}\n\
548 next: remove those items at the destination, then copy again."
549 )]
550 CopyNotUndone {
551 error: Box<EngineError>,
553 left_behind: LeftBehind,
559 refusal: SourceError,
561 },
562}
563
564#[derive(Debug, Clone, PartialEq)]
573pub struct LeftBehind {
574 first: GlobalId,
576 rest: Vec<GlobalId>,
578}
579
580impl LeftBehind {
581 #[must_use]
583 pub fn new(first: GlobalId) -> Self {
584 Self {
585 first,
586 rest: Vec::new(),
587 }
588 }
589
590 pub fn push(&mut self, id: GlobalId) {
592 self.rest.push(id);
593 }
594
595 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
597 std::iter::once(&self.first).chain(self.rest.iter())
598 }
599}
600
601impl std::fmt::Display for LeftBehind {
602 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
604 write!(formatter, "{}", self.first)?;
605 for id in &self.rest {
606 write!(formatter, ", {id}")?;
607 }
608 Ok(())
609 }
610}
611
612pub enum ConfiguredSource {
619 Ready(ResolvedSource),
621 Unavailable(UnavailableSource),
623}
624
625impl ConfiguredSource {
626 #[must_use]
628 pub fn name(&self) -> &SourceName {
629 match self {
630 Self::Ready(source) => source.name(),
631 Self::Unavailable(source) => source.name(),
632 }
633 }
634}
635
636pub struct Engine {
638 sources: Vec<ConfiguredSource>,
640 selection: Vec<SourceName>,
642}
643
644impl Engine {
645 #[must_use]
653 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
654 let (ready, unavailable) = resolve_available(config, secrets);
655 Self::new(
656 ready
657 .into_iter()
658 .map(ConfiguredSource::Ready)
659 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
660 .collect(),
661 config.selected_sources(),
662 )
663 }
664
665 #[must_use]
668 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
669 Self { sources, selection }
670 }
671
672 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
674 self.sources.iter().filter_map(|source| match source {
675 ConfiguredSource::Ready(ready) => Some(ready),
676 ConfiguredSource::Unavailable(_) => None,
677 })
678 }
679
680 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
682 self.sources.iter().filter_map(|source| match source {
683 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
684 ConfiguredSource::Ready(_) => None,
685 })
686 }
687
688 #[must_use]
690 pub fn listing(&self) -> Vec<SourceListing> {
691 let mut listings: Vec<SourceListing> = self
692 .ready()
693 .map(|source| SourceListing {
694 source: source.name().clone(),
695 kind: source.kind().to_owned(),
696 state: SourceState::Available {
697 capabilities: source.source().capabilities(),
698 },
699 })
700 .chain(self.unavailable().map(|source| SourceListing {
701 source: source.name().clone(),
702 kind: source.kind().to_owned(),
703 state: SourceState::Unavailable {
704 error: source.error().clone(),
705 },
706 }))
707 .collect();
708 listings.sort_by(|left, right| left.source.cmp(&right.source));
709 listings
710 }
711
712 #[must_use]
718 pub fn has(&self, name: &SourceName) -> bool {
719 self.sources.iter().any(|source| source.name() == name)
720 }
721
722 pub async fn tasks(
730 &self,
731 request: &TaskRequest,
732 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
733 let mut names = self.resolve_selection(&request.sources)?;
734 if let ProjectSelector::Qualified(id) = &request.project {
739 self.known(&id.source)?;
740 names.retain(|name| name == &id.source);
741 }
742 let query = shape("task-list", &names, &(&request.filters, &request.project));
743 let states = resumption(
744 self,
745 request.paging.token.as_ref(),
746 &[StreamKind::Items],
747 &query,
748 )?;
749 let budget = request.paging.limit.get();
750
751 let mut answer = Answer::new();
752 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
753
754 let shapes: Vec<TaskShape> = ready
755 .iter()
756 .map(|source| {
757 shape_tasks(
758 &source.source().capabilities(),
759 &request.filters,
760 &project_filter(&request.project),
761 )
762 })
763 .collect();
764 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
765 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
766
767 let walks = ready
768 .iter()
769 .enumerate()
770 .map(|(index, source)| {
771 fetch_tasks(
772 source,
773 &shapes[index],
774 &starts[index],
775 budget,
776 &counters[index],
777 )
778 })
779 .collect();
780
781 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
782 answer.finish(
783 streams,
784 budget,
785 owed(&states),
786 &query,
787 |name, task: Task| {
788 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
789 },
790 )
791 }
792
793 pub async fn projects(
799 &self,
800 request: &ProjectRequest,
801 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
802 let names = self.resolve_selection(&request.sources)?;
803 let query = shape("project-list", &names, &request.filters);
804 let states = resumption(
805 self,
806 request.paging.token.as_ref(),
807 &[StreamKind::Items],
808 &query,
809 )?;
810 let budget = request.paging.limit.get();
811
812 let mut answer = Answer::new();
813 let mut with_projects = Vec::new();
819 for source in answer.split(self, &names) {
820 if source.source().capabilities().projects.is_native() {
821 with_projects.push(source);
822 } else {
823 answer.unreachable_predicate(source, Predicate::Project);
824 }
825 }
826 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
827
828 let shapes: Vec<ProjectShape> = ready
829 .iter()
830 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
831 .collect();
832 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
833 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
834
835 let walks = ready
836 .iter()
837 .enumerate()
838 .map(|(index, source)| {
839 fetch_projects(
840 source,
841 &shapes[index],
842 &starts[index],
843 budget,
844 &counters[index],
845 )
846 })
847 .collect();
848
849 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
850 answer.finish(
851 streams,
852 budget,
853 owed(&states),
854 &query,
855 |name, project: Project| Qualified {
856 id: GlobalId::new(name.clone(), project.id.clone()),
857 item: project,
858 },
859 )
860 }
861
862 pub async fn documents(
874 &self,
875 request: &DocumentRequest,
876 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
877 let mut names = self.resolve_selection(&request.sources)?;
878 if let ProjectSelector::Qualified(id) = &request.project {
881 self.known(&id.source)?;
882 names.retain(|name| name == &id.source);
883 }
884 let query = shape(
885 "document-list",
886 &names,
887 &(&request.filters, &request.project),
888 );
889 let states = resumption(
890 self,
891 request.paging.token.as_ref(),
892 &[StreamKind::Items],
893 &query,
894 )?;
895 let budget = request.paging.limit.get();
896
897 let mut answer = Answer::new();
898 let mut with_documents = Vec::new();
899 for source in answer.split(self, &names) {
900 if source.source().capabilities().documents.is_native() {
901 with_documents.push(source);
902 } else {
903 answer.unreachable_predicate(source, Predicate::Document);
904 }
905 }
906 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
907
908 let shapes: Vec<DocumentShape> = ready
909 .iter()
910 .map(|source| {
911 shape_documents(
912 &source.source().capabilities(),
913 &request.filters,
914 &project_filter(&request.project),
915 )
916 })
917 .collect();
918 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
919 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
920
921 let walks = ready
922 .iter()
923 .enumerate()
924 .map(|(index, source)| {
925 fetch_documents(
926 source,
927 &shapes[index],
928 &starts[index],
929 budget,
930 &counters[index],
931 )
932 })
933 .collect();
934
935 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
936 answer.finish(
937 streams,
938 budget,
939 owed(&states),
940 &query,
941 |name, document: Document| Qualified {
942 id: GlobalId::new(name.clone(), document.id.clone()),
943 item: document,
944 },
945 )
946 }
947
948 pub async fn labels(
954 &self,
955 request: &LabelRequest,
956 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
957 let names = self.resolve_selection(&request.sources)?;
958 let query = shape("label-list", &names, &());
959 let states = resumption(
960 self,
961 request.paging.token.as_ref(),
962 &[StreamKind::Items],
963 &query,
964 )?;
965 let budget = request.paging.limit.get();
966
967 let mut answer = Answer::new();
968 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
969 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
970 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
971
972 let walks = ready
973 .iter()
974 .enumerate()
975 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
976 .collect();
977
978 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
979 answer.finish(
980 streams,
981 budget,
982 owed(&states),
983 &query,
984 |name, label: Label| Qualified {
985 id: GlobalId::new(name.clone(), label.id.clone()),
986 item: label,
987 },
988 )
989 }
990
991 pub async fn search(
997 &self,
998 request: &SearchRequest,
999 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1000 let names = self.resolve_selection(&request.sources)?;
1001 let reads: &[StreamKind] = match request.kind {
1005 SearchKind::Tasks => &[StreamKind::Tasks],
1006 SearchKind::Projects => &[StreamKind::Projects],
1007 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1008 };
1009 let query = shape("search", &names, &request.text);
1014 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1015 let budget = request.paging.limit.get();
1016 let filters = Filters {
1017 text: Some(request.text.clone()),
1018 ..Filters::default()
1019 };
1020
1021 let mut answer = Answer::new();
1022
1023 let mut ready = Vec::new();
1026 let mut kinds = Vec::new();
1027 let mut starts = Vec::new();
1028 for source in answer.split(self, &names) {
1029 let mut streams = Vec::new();
1030 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1031 streams.push(StreamKind::Tasks);
1032 }
1033 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1034 if source.source().capabilities().projects.is_native() {
1035 streams.push(StreamKind::Projects);
1036 } else {
1037 answer.unreachable_predicate(source, Predicate::Project);
1038 }
1039 }
1040 for stream in streams {
1041 if let Some(resume) = resume_at(&states, source.name(), stream) {
1042 ready.push(source);
1043 kinds.push(stream);
1044 starts.push(resume);
1045 }
1046 }
1047 }
1048
1049 let shapes: Vec<HitShape> = ready
1050 .iter()
1051 .zip(kinds.iter())
1052 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1053 .collect();
1054 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1055 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1056
1057 let walks = ready
1058 .iter()
1059 .enumerate()
1060 .map(|(index, source)| {
1061 fetch_hits(
1062 source,
1063 &shapes[index],
1064 &starts[index],
1065 budget,
1066 &counters[index],
1067 )
1068 })
1069 .collect();
1070
1071 let streams =
1072 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1073 answer.finish(
1074 streams,
1075 budget,
1076 owed(&states),
1077 &query,
1078 |name, found: Found| match found {
1079 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1080 GlobalId::new(name.clone(), task.id.clone()),
1081 task,
1082 )),
1083 Found::Project(project) => SearchHit::Project(Qualified {
1084 id: GlobalId::new(name.clone(), project.id.clone()),
1085 item: project,
1086 }),
1087 },
1088 )
1089 }
1090
1091 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1098 let name = self.known(&id.source)?;
1099 let mut answer = Answer::new();
1100 let selected = answer.split(self, std::slice::from_ref(&name));
1101 let Some(source) = selected.first() else {
1102 return answer.nothing();
1103 };
1104 let found = source.source().get_task(&id.native).await;
1105 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1106 answer.one(source, found, |task| {
1107 delivery::qualified_task(qualified, task)
1108 })
1109 }
1110
1111 pub async fn project(
1117 &self,
1118 id: &GlobalId,
1119 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1120 let name = self.known(&id.source)?;
1121 let mut answer = Answer::new();
1122 let selected = answer.split(self, std::slice::from_ref(&name));
1123 let Some(source) = selected.first() else {
1124 return answer.nothing();
1125 };
1126 let found = source.source().get_project(&id.native).await;
1127 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1128 answer.one(source, found, |project| Qualified {
1129 id: qualified,
1130 item: project,
1131 })
1132 }
1133
1134 pub async fn document(
1144 &self,
1145 id: &GlobalId,
1146 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1147 let name = self.known(&id.source)?;
1148 let mut answer = Answer::new();
1149 let selected = answer.split(self, std::slice::from_ref(&name));
1150 let Some(source) = selected.first() else {
1151 return answer.nothing();
1152 };
1153 if !source.source().capabilities().documents.is_native() {
1154 answer.unreachable_predicate(source, Predicate::Document);
1155 return answer.nothing();
1156 }
1157 let found = source.source().get_document(&id.native).await;
1158 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1159 answer.one(source, found, |document| Qualified {
1160 id: qualified,
1161 item: document,
1162 })
1163 }
1164
1165 pub async fn task_dependencies(
1172 &self,
1173 request: &DependencyRequest,
1174 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1175 self.dependencies(request, Entity::Task).await
1176 }
1177
1178 pub async fn project_dependencies(
1184 &self,
1185 request: &DependencyRequest,
1186 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1187 self.dependencies(request, Entity::Project).await
1188 }
1189
1190 async fn dependencies(
1193 &self,
1194 request: &DependencyRequest,
1195 entity: Entity,
1196 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1197 let name = self.known(&request.id.source)?;
1198 let query = shape(
1199 "dependencies",
1200 std::slice::from_ref(&name),
1201 &(entity, &request.id.native, request.direction),
1202 );
1203 let states = resumption(
1204 self,
1205 request.paging.token.as_ref(),
1206 &[StreamKind::Items],
1207 &query,
1208 )?;
1209 let budget = request.paging.limit.get();
1210
1211 let mut answer = Answer::new();
1212 let (ready, starts) = walking(
1213 answer.split(self, std::slice::from_ref(&name)),
1214 &states,
1215 StreamKind::Items,
1216 );
1217 let Some(source) = ready.first() else {
1218 return answer.nothing();
1219 };
1220
1221 let capabilities = source.source().capabilities();
1222 let support = match entity {
1223 Entity::Task => capabilities.task_dependencies,
1224 Entity::Project => capabilities.project_dependencies,
1225 };
1226 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1230 let mut outcomes = Outcomes::default();
1231 if request.direction == Direction::DependedOnBy {
1232 if emulating {
1233 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1234 } else {
1235 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1236 }
1237 }
1238
1239 let counters = vec![AtomicU32::new(0)];
1240 let walked = fetch_edges(
1241 source,
1242 &request.id.native,
1243 request.direction,
1244 entity,
1245 emulating,
1246 &starts[0],
1247 budget,
1248 &counters[0],
1249 )
1250 .await;
1251
1252 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1253 answer.finish(
1254 streams,
1255 budget,
1256 owed(&states),
1257 &query,
1258 |name, edge: DependencyEdge| QualifiedEdge {
1259 from: qualify_endpoint(name, edge.from),
1260 to: qualify_endpoint(name, edge.to),
1261 kind: edge.kind,
1262 },
1263 )
1264 }
1265
1266 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1268 if asked.is_empty() {
1269 if self.selection.is_empty() {
1270 return Err(EngineError::NoSources);
1271 }
1272 return Ok(self.selection.clone());
1273 }
1274 asked.iter().map(|name| self.known(name)).collect()
1275 }
1276
1277 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1279 if self.has(name) {
1280 return Ok(name.clone());
1281 }
1282 if self.sources.is_empty() {
1283 return Err(EngineError::NoSources);
1284 }
1285 Err(EngineError::UnknownSource {
1286 name: name.to_string(),
1287 configured: self
1288 .listing()
1289 .iter()
1290 .map(|listing| listing.source.to_string())
1291 .collect::<Vec<_>>()
1292 .join(", "),
1293 })
1294 }
1295}
1296
1297fn qualify_endpoint(
1298 source: &SourceName,
1299 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1300) -> QualifiedEndpoint {
1301 let kind = endpoint.kind;
1302 let is_qualified = endpoint.is_qualified();
1303 let endpoint_id = endpoint.into_id();
1304 QualifiedEndpoint {
1305 id: if is_qualified {
1306 endpoint_id
1307 .parse()
1308 .expect("plugin-api validates qualified dependency endpoints")
1309 } else {
1310 GlobalId::new(source.clone(), NativeId(endpoint_id))
1311 },
1312 kind,
1313 }
1314}
1315
1316#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1318enum Entity {
1319 Task,
1321 Project,
1323}
1324
1325enum Found {
1327 Task(Task),
1329 Project(Project),
1331}
1332
1333#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1335enum Outcome {
1336 PushedDown,
1338 AppliedLocally,
1340 Emulated,
1342 Unavailable,
1344}
1345
1346#[derive(Debug, Clone, Default, PartialEq)]
1358struct Outcomes(BTreeMap<Predicate, Outcome>);
1359
1360impl Outcomes {
1361 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1367 self.0.insert(predicate, outcome);
1368 }
1369
1370 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1373 for predicate in predicates {
1374 self.record(predicate, outcome);
1375 }
1376 }
1377
1378 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1380 self.0
1381 .iter()
1382 .filter(|(_, recorded)| **recorded == outcome)
1383 .map(|(predicate, _)| *predicate)
1384 .collect()
1385 }
1386}
1387
1388struct TaskShape {
1390 pushed: TaskQuery,
1392 local: LocalTasks,
1394 outcomes: Outcomes,
1396}
1397
1398struct ProjectShape {
1400 pushed: ProjectQuery,
1402 local: LocalProjects,
1404 outcomes: Outcomes,
1406}
1407
1408struct DocumentShape {
1410 pushed: DocumentQuery,
1412 local: LocalDocuments,
1414 outcomes: Outcomes,
1416}
1417
1418struct HitShape {
1420 stream: StreamKind,
1422 tasks: TaskQuery,
1424 projects: ProjectQuery,
1426 local_tasks: LocalTasks,
1428 local_projects: LocalProjects,
1430 outcomes: Outcomes,
1432}
1433
1434struct Answer {
1440 plans: Vec<SourcePlan>,
1442 errors: Vec<SourceFailure>,
1444}
1445
1446impl Answer {
1447 fn new() -> Self {
1448 Self {
1449 plans: Vec::new(),
1450 errors: Vec::new(),
1451 }
1452 }
1453
1454 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1459 let mut selected = Vec::new();
1460 for name in names {
1461 match engine.sources.iter().find(|source| source.name() == name) {
1462 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1463 Some(ConfiguredSource::Unavailable(source)) => {
1464 self.errors.push(source.failure());
1465 }
1466 None => {}
1467 }
1468 }
1469 selected
1470 }
1471
1472 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1474 let mut outcomes = Outcomes::default();
1475 outcomes.record(predicate, Outcome::Unavailable);
1476 self.plans.push(plan_for(source, outcomes, 0));
1477 }
1478
1479 fn collect<T>(
1481 &mut self,
1482 ready: &[&ResolvedSource],
1483 walked: Vec<Result<Fetched<T>, SourceError>>,
1484 counters: &[AtomicU32],
1485 outcomes: Vec<Outcomes>,
1486 ) -> Vec<Stream<T>> {
1487 let kinds = vec![StreamKind::Items; ready.len()];
1488 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1489 }
1490
1491 fn collect_streams<T>(
1493 &mut self,
1494 ready: &[&ResolvedSource],
1495 kinds: &[StreamKind],
1496 walked: Vec<Result<Fetched<T>, SourceError>>,
1497 counters: &[AtomicU32],
1498 outcomes: Vec<Outcomes>,
1499 ) -> Vec<Stream<T>> {
1500 let mut streams = Vec::new();
1501 for (index, result) in walked.into_iter().enumerate() {
1502 let source = ready[index];
1503 let pages = counters[index].load(Ordering::Relaxed);
1504 self.plans
1505 .push(plan_for(source, outcomes[index].clone(), pages));
1506 match result {
1507 Ok(fetched) => streams.push(Stream {
1508 source: source.name().clone(),
1509 kind: kinds[index],
1510 fetched,
1511 }),
1512 Err(error) => self.errors.push(SourceFailure {
1515 source: source.name().clone(),
1516 error,
1517 }),
1518 }
1519 }
1520 streams
1521 }
1522
1523 fn one<T, U>(
1525 mut self,
1526 source: &ResolvedSource,
1527 found: Result<Option<T>, SourceError>,
1528 qualify: impl FnOnce(T) -> U,
1529 ) -> Result<QueryResponse<U>, EngineError> {
1530 self.plans.push(plan_for(source, Outcomes::default(), 1));
1531 let items = match found {
1532 Ok(Some(item)) => vec![qualify(item)],
1533 Ok(None) => Vec::new(),
1534 Err(error) => {
1535 self.errors.push(SourceFailure {
1536 source: source.name().clone(),
1537 error,
1538 });
1539 Vec::new()
1540 }
1541 };
1542 Ok(QueryResponse {
1543 items,
1544 next: None,
1545 plan: QueryPlan {
1546 per_source: merge_plans(self.plans),
1547 },
1548 errors: self.errors,
1549 })
1550 }
1551
1552 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1554 Ok(QueryResponse {
1555 items: Vec::new(),
1556 next: None,
1557 plan: QueryPlan {
1558 per_source: merge_plans(self.plans),
1559 },
1560 errors: self.errors,
1561 })
1562 }
1563
1564 fn finish<T, U>(
1569 self,
1570 streams: Vec<Stream<T>>,
1571 budget: u32,
1572 first: Option<&Owed>,
1573 query: &str,
1574 qualify: impl Fn(&SourceName, T) -> U,
1575 ) -> Result<QueryResponse<U>, EngineError> {
1576 let (rows, states, owed) = merge(streams, budget, first);
1577 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1578 Ok(QueryResponse {
1579 items: rows
1580 .into_iter()
1581 .map(|(name, item)| qualify(&name, item))
1582 .collect(),
1583 next,
1584 plan: QueryPlan {
1585 per_source: merge_plans(self.plans),
1586 },
1587 errors: self.errors,
1588 })
1589 }
1590}
1591
1592fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1595 SourcePlan {
1596 source: source.name().clone(),
1597 kind: source.kind().to_owned(),
1598 pushed_down: outcomes.with(Outcome::PushedDown),
1599 applied_locally: outcomes.with(Outcome::AppliedLocally),
1600 emulated: outcomes.with(Outcome::Emulated),
1601 unavailable: outcomes.with(Outcome::Unavailable),
1602 pages_fetched: pages,
1603 }
1604}
1605
1606fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1611 let mut merged: Vec<SourcePlan> = Vec::new();
1612 for plan in plans {
1613 if let Some(existing) = merged
1614 .iter_mut()
1615 .find(|existing| existing.source == plan.source)
1616 {
1617 existing.pushed_down.extend(plan.pushed_down);
1618 existing.applied_locally.extend(plan.applied_locally);
1619 existing.emulated.extend(plan.emulated);
1620 existing.unavailable.extend(plan.unavailable);
1621 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1622 for list in [
1623 &mut existing.pushed_down,
1624 &mut existing.applied_locally,
1625 &mut existing.emulated,
1626 &mut existing.unavailable,
1627 ] {
1628 list.sort_unstable();
1629 list.dedup();
1630 }
1631 } else {
1632 merged.push(plan);
1633 }
1634 }
1635 merged
1636}
1637
1638fn walking<'a>(
1644 selected: Vec<&'a ResolvedSource>,
1645 states: &Option<Resumption>,
1646 kind: StreamKind,
1647) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1648 let mut ready = Vec::new();
1649 let mut starts = Vec::new();
1650 for source in selected {
1651 if let Some(resume) = resume_at(states, source.name(), kind) {
1652 ready.push(source);
1653 starts.push(resume);
1654 }
1655 }
1656 (ready, starts)
1657}
1658
1659fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1677 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1678 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1679}
1680
1681fn fingerprint(text: &str) -> String {
1683 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1684 for byte in text.as_bytes() {
1685 hash ^= u64::from(*byte);
1686 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1687 }
1688 format!("{hash:016x}")
1689}
1690
1691fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1696 document.as_ref()?.owed.as_ref()
1697}
1698
1699fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1701 match states {
1702 None => Some(Resume::default()),
1703 Some(document) => document
1704 .streams
1705 .iter()
1706 .find(|state| &state.source == source && state.stream == kind)
1707 .map(|state| state.resume.clone()),
1708 }
1709}
1710
1711fn resumption(
1742 engine: &Engine,
1743 token: Option<&PageToken>,
1744 reads: &[StreamKind],
1745 query: &str,
1746) -> Result<Option<Resumption>, EngineError> {
1747 let Some(document) = token.map(PageToken::decode) else {
1748 return Ok(None);
1749 };
1750
1751 if document.query != query {
1758 return Err(EngineError::Token {
1759 message: "this page token was written by a different query — resume the walk it \
1760 came from, or drop --page to start this one from the beginning"
1761 .to_owned(),
1762 });
1763 }
1764 let states = &document.streams;
1765
1766 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
1767 for state in states {
1768 if !reads.contains(&state.stream) {
1769 return Err(EngineError::Token {
1770 message: format!(
1771 "this page token resumes {}, which this command does not read — it \
1772 was written by a different query",
1773 state.stream.describe()
1774 ),
1775 });
1776 }
1777 let ceiling = engine
1778 .ready()
1779 .find(|source| source.name() == &state.source)
1780 .map(ceiling);
1781 if ceiling.is_none() && !engine.has(&state.source) {
1782 return Err(EngineError::Token {
1783 message: format!(
1784 "this page token resumes a source called {:?}, which this \
1785 configuration does not have",
1786 state.source.as_str()
1787 ),
1788 });
1789 }
1790 if let Some(ceiling) = ceiling
1791 && state.resume.skip >= ceiling
1792 {
1793 return Err(EngineError::Token {
1794 message: format!(
1795 "this page token resumes {} rows into a page of source {:?}, which \
1796 serves at most {ceiling}",
1797 state.resume.skip,
1798 state.source.as_str()
1799 ),
1800 });
1801 }
1802 if seen.contains(&(&state.source, state.stream)) {
1803 return Err(EngineError::Token {
1804 message: format!(
1805 "this page token gives source {:?} two places to resume from",
1806 state.source.as_str()
1807 ),
1808 });
1809 }
1810 seen.push((&state.source, state.stream));
1811 }
1812
1813 if let Some(owed) = &document.owed
1818 && !document
1819 .streams
1820 .iter()
1821 .any(|state| state.source == owed.source && state.stream == owed.stream)
1822 {
1823 return Err(EngineError::Token {
1824 message: format!(
1825 "this page token owes the next row to a stream it does not resume, \
1826 {:?}'s {}",
1827 owed.source.as_str(),
1828 owed.stream.describe()
1829 ),
1830 });
1831 }
1832
1833 Ok(Some(document))
1834}
1835
1836fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
1842 match selector {
1843 ProjectSelector::Any => ProjectFilter::Any,
1844 ProjectSelector::Orphans => ProjectFilter::Orphans,
1845 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
1846 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
1847 }
1848}
1849
1850fn text_predicates(fields: TextFields) -> Vec<Predicate> {
1852 match fields {
1853 TextFields::Title => vec![Predicate::SearchTitle],
1854 TextFields::Content => vec![Predicate::SearchContent],
1855 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
1856 }
1857}
1858
1859fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
1866 match fields {
1867 TextFields::Title => capabilities.search_title.is_native(),
1868 TextFields::Content => capabilities.search_content.is_native(),
1869 TextFields::TitleOrContent => {
1870 capabilities.search_title.is_native() && capabilities.search_content.is_native()
1871 }
1872 }
1873}
1874
1875fn shape_tasks(
1877 capabilities: &Capabilities,
1878 filters: &Filters,
1879 project: &ProjectFilter,
1880) -> TaskShape {
1881 let mut pushed = TaskQuery::default();
1882 let mut local = LocalTasks::default();
1883 let mut outcomes = Outcomes::default();
1884
1885 if !filters.labels.is_empty() {
1886 if capabilities.filter_by_label.is_native() {
1887 pushed.labels = filters.labels.clone();
1888 outcomes.record(Predicate::Label, Outcome::PushedDown);
1889 } else {
1890 local.labels = Some(filters.labels.clone());
1891 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1892 }
1893 }
1894 if !filters.statuses.is_empty() {
1895 if capabilities.filter_by_status.is_native() {
1896 pushed.statuses.clone_from(&filters.statuses);
1897 outcomes.record(Predicate::Status, Outcome::PushedDown);
1898 } else {
1899 local.statuses.clone_from(&filters.statuses);
1900 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1901 }
1902 }
1903 if let Some(text) = &filters.text {
1904 let predicates = text_predicates(text.fields);
1905 if searches_natively(capabilities, text.fields) {
1906 pushed.text = Some(text.clone());
1907 outcomes.record_all(predicates, Outcome::PushedDown);
1908 } else {
1909 local.text = Some(text.clone());
1910 outcomes.record_all(predicates, Outcome::AppliedLocally);
1911 }
1912 }
1913 match project {
1914 ProjectFilter::Any => {}
1915 ProjectFilter::Orphans => {
1916 if capabilities.orphan_tasks.is_native() {
1917 pushed.project = ProjectFilter::Orphans;
1918 outcomes.record(Predicate::Project, Outcome::PushedDown);
1919 } else {
1920 local.project = Some(ProjectFilter::Orphans);
1921 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1922 }
1923 }
1924 ProjectFilter::Is(id) => {
1925 if capabilities.projects.is_native() {
1926 pushed.project = ProjectFilter::Is(id.clone());
1927 outcomes.record(Predicate::Project, Outcome::PushedDown);
1928 } else {
1929 local.project = Some(ProjectFilter::Is(id.clone()));
1930 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1931 }
1932 }
1933 }
1934
1935 TaskShape {
1936 pushed,
1937 local,
1938 outcomes,
1939 }
1940}
1941
1942fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
1944 let mut pushed = ProjectQuery::default();
1945 let mut local = LocalProjects::default();
1946 let mut outcomes = Outcomes::default();
1947
1948 if !filters.labels.is_empty() {
1949 if capabilities.filter_by_label.is_native() {
1950 pushed.labels = filters.labels.clone();
1951 outcomes.record(Predicate::Label, Outcome::PushedDown);
1952 } else {
1953 local.labels = Some(filters.labels.clone());
1954 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1955 }
1956 }
1957 if !filters.statuses.is_empty() {
1958 if capabilities.filter_by_status.is_native() {
1959 pushed.statuses.clone_from(&filters.statuses);
1960 outcomes.record(Predicate::Status, Outcome::PushedDown);
1961 } else {
1962 local.statuses.clone_from(&filters.statuses);
1963 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1964 }
1965 }
1966 if let Some(text) = &filters.text {
1967 let predicates = text_predicates(text.fields);
1968 if searches_natively(capabilities, text.fields) {
1969 pushed.text = Some(text.clone());
1970 outcomes.record_all(predicates, Outcome::PushedDown);
1971 } else {
1972 local.text = Some(text.clone());
1973 outcomes.record_all(predicates, Outcome::AppliedLocally);
1974 }
1975 }
1976
1977 ProjectShape {
1978 pushed,
1979 local,
1980 outcomes,
1981 }
1982}
1983
1984fn shape_documents(
1989 capabilities: &Capabilities,
1990 filters: &DocumentFilters,
1991 project: &ProjectFilter,
1992) -> DocumentShape {
1993 let mut pushed = DocumentQuery::default();
1994 let mut local = LocalDocuments::default();
1995 let mut outcomes = Outcomes::default();
1996
1997 if !filters.labels.is_empty() {
1998 if capabilities.filter_by_label.is_native() {
1999 pushed.labels = filters.labels.clone();
2000 outcomes.record(Predicate::Label, Outcome::PushedDown);
2001 } else {
2002 local.labels = Some(filters.labels.clone());
2003 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2004 }
2005 }
2006 if let Some(text) = &filters.text {
2007 let predicates = text_predicates(text.fields);
2008 if searches_natively(capabilities, text.fields) {
2009 pushed.text = Some(text.clone());
2010 outcomes.record_all(predicates, Outcome::PushedDown);
2011 } else {
2012 local.text = Some(text.clone());
2013 outcomes.record_all(predicates, Outcome::AppliedLocally);
2014 }
2015 }
2016 match project {
2017 ProjectFilter::Any => {}
2018 ProjectFilter::Orphans => {
2019 if capabilities.orphan_tasks.is_native() {
2020 pushed.project = ProjectFilter::Orphans;
2021 outcomes.record(Predicate::Project, Outcome::PushedDown);
2022 } else {
2023 local.project = Some(ProjectFilter::Orphans);
2024 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2025 }
2026 }
2027 ProjectFilter::Is(id) => {
2028 if capabilities.projects.is_native() {
2029 pushed.project = ProjectFilter::Is(id.clone());
2030 outcomes.record(Predicate::Project, Outcome::PushedDown);
2031 } else {
2032 local.project = Some(ProjectFilter::Is(id.clone()));
2033 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2034 }
2035 }
2036 }
2037
2038 DocumentShape {
2039 pushed,
2040 local,
2041 outcomes,
2042 }
2043}
2044
2045fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2047 match stream {
2048 StreamKind::Projects => {
2049 let shaped = shape_projects(capabilities, filters);
2050 HitShape {
2051 stream,
2052 tasks: TaskQuery::default(),
2053 projects: shaped.pushed,
2054 local_tasks: LocalTasks::default(),
2055 local_projects: shaped.local,
2056 outcomes: shaped.outcomes,
2057 }
2058 }
2059 StreamKind::Items | StreamKind::Tasks => {
2060 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any);
2061 HitShape {
2062 stream,
2063 tasks: shaped.pushed,
2064 projects: ProjectQuery::default(),
2065 local_tasks: shaped.local,
2066 local_projects: LocalProjects::default(),
2067 outcomes: shaped.outcomes,
2068 }
2069 }
2070 }
2071}
2072
2073fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2080 if compensating {
2081 ceiling
2082 } else {
2083 budget.min(ceiling)
2084 }
2085}
2086
2087fn ceiling(source: &ResolvedSource) -> u32 {
2089 source.source().capabilities().max_page_size.max(1)
2090}
2091
2092async fn fetch_tasks(
2094 source: &ResolvedSource,
2095 shape: &TaskShape,
2096 start: &Resume,
2097 budget: u32,
2098 calls: &AtomicU32,
2099) -> Result<Fetched<Task>, SourceError> {
2100 let compensating = shape.local != LocalTasks::default();
2101 walk(
2102 start,
2103 budget,
2104 page_size(compensating, budget, ceiling(source)),
2105 |task| shape.local.keeps(task),
2106 |cursor, limit| async move {
2107 calls.fetch_add(1, Ordering::Relaxed);
2108 let request = PageRequest { cursor, limit };
2109 source.source().query_tasks(&shape.pushed, &request).await
2110 },
2111 )
2112 .await
2113}
2114
2115async fn fetch_projects(
2117 source: &ResolvedSource,
2118 shape: &ProjectShape,
2119 start: &Resume,
2120 budget: u32,
2121 calls: &AtomicU32,
2122) -> Result<Fetched<Project>, SourceError> {
2123 let compensating = shape.local != LocalProjects::default();
2124 walk(
2125 start,
2126 budget,
2127 page_size(compensating, budget, ceiling(source)),
2128 |project| shape.local.keeps(project),
2129 |cursor, limit| async move {
2130 calls.fetch_add(1, Ordering::Relaxed);
2131 let request = PageRequest { cursor, limit };
2132 source
2133 .source()
2134 .query_projects(&shape.pushed, &request)
2135 .await
2136 },
2137 )
2138 .await
2139}
2140
2141async fn fetch_documents(
2143 source: &ResolvedSource,
2144 shape: &DocumentShape,
2145 start: &Resume,
2146 budget: u32,
2147 calls: &AtomicU32,
2148) -> Result<Fetched<Document>, SourceError> {
2149 let compensating = shape.local != LocalDocuments::default();
2150 walk(
2151 start,
2152 budget,
2153 page_size(compensating, budget, ceiling(source)),
2154 |document| shape.local.keeps(document),
2155 |cursor, limit| async move {
2156 calls.fetch_add(1, Ordering::Relaxed);
2157 let request = PageRequest { cursor, limit };
2158 source
2159 .source()
2160 .query_documents(&shape.pushed, &request)
2161 .await
2162 },
2163 )
2164 .await
2165}
2166
2167async fn fetch_labels(
2169 source: &ResolvedSource,
2170 start: &Resume,
2171 budget: u32,
2172 calls: &AtomicU32,
2173) -> Result<Fetched<Label>, SourceError> {
2174 walk(
2175 start,
2176 budget,
2177 page_size(false, budget, ceiling(source)),
2178 |_| true,
2179 |cursor, limit| async move {
2180 calls.fetch_add(1, Ordering::Relaxed);
2181 let request = PageRequest { cursor, limit };
2182 source.source().labels(&request).await
2183 },
2184 )
2185 .await
2186}
2187
2188async fn fetch_hits(
2190 source: &ResolvedSource,
2191 shape: &HitShape,
2192 start: &Resume,
2193 budget: u32,
2194 calls: &AtomicU32,
2195) -> Result<Fetched<Found>, SourceError> {
2196 let ceiling = ceiling(source);
2197 match shape.stream {
2198 StreamKind::Projects => {
2199 let compensating = shape.local_projects != LocalProjects::default();
2200 walk(
2201 start,
2202 budget,
2203 page_size(compensating, budget, ceiling),
2204 |found| match found {
2205 Found::Project(project) => shape.local_projects.keeps(project),
2206 Found::Task(_) => true,
2207 },
2208 |cursor, limit| async move {
2209 calls.fetch_add(1, Ordering::Relaxed);
2210 let request = PageRequest { cursor, limit };
2211 let page = source
2212 .source()
2213 .query_projects(&shape.projects, &request)
2214 .await?;
2215 Ok(Page {
2216 items: page.items.into_iter().map(Found::Project).collect(),
2217 next: page.next,
2218 })
2219 },
2220 )
2221 .await
2222 }
2223 StreamKind::Items | StreamKind::Tasks => {
2224 let compensating = shape.local_tasks != LocalTasks::default();
2225 walk(
2226 start,
2227 budget,
2228 page_size(compensating, budget, ceiling),
2229 |found| match found {
2230 Found::Task(task) => shape.local_tasks.keeps(task),
2231 Found::Project(_) => true,
2232 },
2233 |cursor, limit| async move {
2234 calls.fetch_add(1, Ordering::Relaxed);
2235 let request = PageRequest { cursor, limit };
2236 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2237 Ok(Page {
2238 items: page.items.into_iter().map(Found::Task).collect(),
2239 next: page.next,
2240 })
2241 },
2242 )
2243 .await
2244 }
2245 }
2246}
2247
2248async fn forward_edges(
2250 source: &ResolvedSource,
2251 entity: Entity,
2252 id: &NativeId,
2253 request: &PageRequest,
2254) -> Result<Page<DependencyEdge>, SourceError> {
2255 match entity {
2256 Entity::Task => {
2257 source
2258 .source()
2259 .task_dependencies(id, Direction::DependsOn, request)
2260 .await
2261 }
2262 Entity::Project => {
2263 source
2264 .source()
2265 .project_dependencies(id, Direction::DependsOn, request)
2266 .await
2267 }
2268 }
2269}
2270
2271#[expect(
2280 clippy::too_many_arguments,
2281 reason = "every argument is one axis of one walk — the source, the item, the \
2282 direction, which of its two graphs, whether the reverse is emulated, where \
2283 to resume, how many rows to return and where to count calls. Grouping them \
2284 into a struct would name the same eight values one indirection further from \
2285 the loop that reads them."
2286)]
2287async fn fetch_edges(
2288 source: &ResolvedSource,
2289 native: &NativeId,
2290 direction: Direction,
2291 entity: Entity,
2292 emulating: bool,
2293 start: &Resume,
2294 budget: u32,
2295 calls: &AtomicU32,
2296) -> Result<Fetched<DependencyEdge>, SourceError> {
2297 let ceiling = ceiling(source);
2298 if !emulating {
2299 return walk(
2300 start,
2301 budget,
2302 page_size(false, budget, ceiling),
2303 |_| true,
2304 |cursor, limit| async move {
2305 calls.fetch_add(1, Ordering::Relaxed);
2306 let request = PageRequest { cursor, limit };
2307 match entity {
2308 Entity::Task => {
2309 source
2310 .source()
2311 .task_dependencies(native, direction, &request)
2312 .await
2313 }
2314 Entity::Project => {
2315 source
2316 .source()
2317 .project_dependencies(native, direction, &request)
2318 .await
2319 }
2320 }
2321 },
2322 )
2323 .await;
2324 }
2325
2326 walk(
2327 start,
2328 budget,
2329 ceiling,
2330 |_| true,
2331 |cursor, limit| async move {
2332 calls.fetch_add(1, Ordering::Relaxed);
2333 let request = PageRequest { cursor, limit };
2334 let (ids, next) = match entity {
2335 Entity::Task => {
2336 let page = source
2337 .source()
2338 .query_tasks(&TaskQuery::default(), &request)
2339 .await?;
2340 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2341 (ids, page.next)
2342 }
2343 Entity::Project => {
2344 let page = source
2345 .source()
2346 .query_projects(&ProjectQuery::default(), &request)
2347 .await?;
2348 let ids: Vec<NativeId> =
2349 page.items.into_iter().map(|project| project.id).collect();
2350 (ids, page.next)
2351 }
2352 };
2353
2354 let mut edges = Vec::new();
2355 for id in ids {
2356 let mut inner: Option<Cursor> = None;
2357 loop {
2358 calls.fetch_add(1, Ordering::Relaxed);
2359 let request = PageRequest {
2360 cursor: inner.clone(),
2361 limit,
2362 };
2363 let page = forward_edges(source, entity, &id, &request).await?;
2364 fits(page.items.len(), limit)?;
2368 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2369 unrepeated(
2370 page.next.as_ref(),
2371 inner.as_ref(),
2372 "its forward edges were being scanned",
2373 )?;
2374 match page.next {
2375 Some(cursor) => inner = Some(cursor),
2376 None => break,
2377 }
2378 }
2379 }
2380
2381 Ok(Page { items: edges, next })
2382 },
2383 )
2384 .await
2385}