1mod comment;
19mod copy;
20mod delivery;
21mod fetch;
22mod join;
23mod local;
24mod metadata;
25mod narrow;
26mod resume;
27
28use std::collections::BTreeMap;
29use std::num::NonZeroU32;
30use std::sync::atomic::{AtomicU32, Ordering};
31
32use onetaskgraph_plugin_api::{
33 Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
34 MetadataRecord, NativeId, Page, PageRequest, Priority, Project, ProjectFilter, ProjectQuery,
35 SecretResolver, SourceError, SourceName, StatusCategory, Task, TaskQuery, TextFields,
36 TextQuery,
37};
38use schemars::JsonSchema;
39use serde::{Deserialize, Serialize};
40
41use crate::GlobalId;
42use crate::config::Config;
43use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
44use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
45
46use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
47use join::join_all;
48use local::{LocalDocuments, LocalProjects, LocalTasks};
49pub(crate) use resume::{Owed, Resumption, StreamState};
50use resume::{Resume, StreamKind};
51
52pub use comment::{CommentList, DeletedComment, TaskDetail};
53pub use copy::{
54 BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
55 Spent,
56};
57pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
58pub use local::ProjectSelector;
59pub use metadata::MetadataSet;
60pub use narrow::{TaskContentSet, TaskPrioritySet};
61
62#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
67pub struct Qualified<T> {
68 pub id: GlobalId,
70 pub item: T,
72}
73
74#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
79pub struct QualifiedEdge {
80 pub from: QualifiedEndpoint,
85 pub to: QualifiedEndpoint,
87 pub kind: onetaskgraph_plugin_api::DependencyKind,
89}
90
91#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
93pub struct QualifiedEndpoint {
94 pub id: GlobalId,
96 pub kind: onetaskgraph_plugin_api::ItemKind,
98}
99
100impl std::fmt::Display for QualifiedEndpoint {
101 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
102 self.id.fmt(formatter)
103 }
104}
105
106#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
108#[serde(tag = "kind", rename_all = "kebab-case")]
109pub enum SearchHit {
110 Task(Qualified<Task>),
112 Project(Qualified<Project>),
114}
115
116#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
118#[serde(rename_all = "kebab-case")]
119pub enum SearchKind {
120 Tasks,
122 Projects,
124 #[default]
126 Both,
127}
128
129#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
131pub struct SourceListing {
132 pub source: SourceName,
134 pub kind: String,
146 #[serde(flatten)]
148 pub state: SourceState,
149}
150
151#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
153#[serde(tag = "state", rename_all = "kebab-case")]
154pub enum SourceState {
155 Available {
157 capabilities: Capabilities,
159 },
160 Unavailable {
162 error: SourceError,
164 },
165}
166
167#[derive(Debug, Clone, PartialEq)]
169pub struct Paging {
170 pub limit: NonZeroU32,
172 pub token: Option<PageToken>,
174}
175
176#[derive(Debug, Clone, Default, PartialEq)]
178pub struct Filters {
179 pub text: Option<TextQuery>,
181 pub labels: LabelFilter,
183 pub statuses: Vec<StatusCategory>,
185}
186
187#[derive(Debug, Clone)]
189pub struct TaskRequest {
190 pub sources: Vec<SourceName>,
192 pub filters: Filters,
194 pub project: ProjectSelector,
196 pub priorities: Vec<Priority>,
202 pub paging: Paging,
204}
205
206#[derive(Debug, Clone)]
208pub struct ProjectRequest {
209 pub sources: Vec<SourceName>,
211 pub filters: Filters,
213 pub paging: Paging,
215}
216
217#[derive(Debug, Clone, Default, PartialEq)]
224pub struct DocumentFilters {
225 pub text: Option<TextQuery>,
227 pub labels: LabelFilter,
229}
230
231#[derive(Debug, Clone)]
233pub struct DocumentRequest {
234 pub sources: Vec<SourceName>,
236 pub filters: DocumentFilters,
238 pub project: ProjectSelector,
240 pub paging: Paging,
242}
243
244#[derive(Debug, Clone)]
246pub struct LabelRequest {
247 pub sources: Vec<SourceName>,
249 pub paging: Paging,
251}
252
253#[derive(Debug, Clone)]
255pub struct SearchRequest {
256 pub sources: Vec<SourceName>,
258 pub text: TextQuery,
260 pub kind: SearchKind,
262 pub paging: Paging,
264}
265
266#[derive(Debug, Clone)]
268pub struct DependencyRequest {
269 pub id: GlobalId,
271 pub direction: Direction,
273 pub paging: Paging,
275}
276
277#[derive(Debug, Clone, PartialEq, thiserror::Error)]
285pub enum EngineError {
286 #[error(
288 "no source named {name:?} is configured\n\
289 next: name one of the configured sources ({configured}), or add {name:?} under \
290 `sources` — `onetaskgraph sources list` shows what this configuration has."
291 )]
292 UnknownSource {
293 name: String,
295 configured: String,
297 },
298
299 #[error(
301 "{message}\n\
302 next: page with a token exactly as the previous page reported it, and against \
303 the same configuration — or drop `--page` to start the walk again."
304 )]
305 Token {
306 message: String,
308 },
309
310 #[error(
312 "no sources are configured\n\
313 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
314 prints what each plugin accepts."
315 )]
316 NoSources,
317
318 #[error(
320 "source {name} cannot be written: its plugin is {kind}, which has no write \
321 side\n\
322 next: copy into a source whose plugin can be written — `onetaskgraph sources \
323 list` reports each one's plugin."
324 )]
325 NotWritable {
326 name: String,
328 kind: String,
330 },
331
332 #[error(
339 "source {name} has no documents: its plugin is {kind}, which holds none\n\
340 next: name a source whose plugin has documents — `onetaskgraph sources list` \
341 reports each one's plugin and what it declares."
342 )]
343 NoDocuments {
344 name: String,
346 kind: String,
348 },
349
350 #[error(
355 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
356 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
357 list` reports each one's plugin."
358 )]
359 NoComments {
360 name: String,
362 kind: String,
364 },
365
366 #[error(
368 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
369 but not added to, edited or removed\n\
370 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
371 source whose plugin can be written — `onetaskgraph sources list` reports each \
372 one's plugin."
373 )]
374 CommentsNotWritable {
375 name: String,
377 kind: String,
379 },
380
381 #[error(
383 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
384 next: set the status in that source itself, or name a task of a source whose plugin \
385 can be written — `onetaskgraph sources list` reports each one's plugin."
386 )]
387 StatusNotWritable {
388 name: String,
390 kind: String,
392 },
393
394 #[error(
396 "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
397 write side\n\
398 next: set the key in that source itself, or name a {record} of a source whose plugin \
399 can be written — `onetaskgraph sources list` reports each one's plugin."
400 )]
401 MetadataNotWritable {
402 name: String,
404 kind: String,
406 record: MetadataRecord,
408 },
409
410 #[error(
412 "source {name} cannot write a priority: its plugin is {kind}, which has no write side\n\
413 next: set the priority in that source itself, or name a task of a source whose plugin \
414 can be written — `onetaskgraph sources list` reports each one's plugin."
415 )]
416 PriorityNotWritable {
417 name: String,
419 kind: String,
421 },
422
423 #[error(
425 "source {name} cannot write a task's content: its plugin is {kind}, which has no write \
426 side\n\
427 next: edit the content in that source itself, or name a task of a source whose plugin \
428 can be written — `onetaskgraph sources list` reports each one's plugin."
429 )]
430 ContentNotWritable {
431 name: String,
433 kind: String,
435 },
436
437 #[error(
442 "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
443 be written to it: its plugin is {kind}, which declares priority unsupported\n\
444 next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
445 reports what each declares — or set the task's priority to none first; a \
446 github-projects source holds one once its configuration sets priority_mapping."
447 )]
448 NoPriority {
449 name: String,
451 kind: String,
453 task: String,
455 priority: onetaskgraph_plugin_api::Priority,
457 },
458
459 #[error(
461 "no project with the id {id}\n\
462 next: check the id, or list what is there — `onetaskgraph project list` reports every \
463 project the configured sources hold."
464 )]
465 NoSuchProject {
466 id: String,
468 },
469
470 #[error(
472 "no document with the id {id}\n\
473 next: check the id, or list what is there — `onetaskgraph document list` reports every \
474 document the configured sources hold."
475 )]
476 NoSuchDocument {
477 id: String,
479 },
480
481 #[error(
483 "no task with the id {id}\n\
484 next: check the id, or list what is there — `onetaskgraph task list` reports every \
485 task the configured sources hold."
486 )]
487 NoSuchTask {
488 id: String,
490 },
491
492 #[error(
494 "task {task} has no comment with the id {comment}\n\
495 next: list its comments — `onetaskgraph task comment list {task}` reports each \
496 one's id."
497 )]
498 NoSuchComment {
499 task: String,
501 comment: String,
503 },
504
505 #[error(
510 "source {name} could not be built: {error}\n\
511 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
512 command again."
513 )]
514 SourceUnavailable {
515 name: String,
517 error: SourceError,
519 },
520
521 #[error(
527 "source {name} could not do it: {error}\n\
528 next: fix what the source named above, then run the command again."
529 )]
530 SourceFailed {
531 name: String,
533 error: SourceError,
535 },
536
537 #[error(
539 "the destination source {name} could not be built: {error}\n\
540 next: fix that source — `onetaskgraph sources list` reports its state — then \
541 copy again."
542 )]
543 DestinationUnavailable {
544 name: String,
546 error: SourceError,
548 },
549
550 #[error(
552 "no item with the id {id}\n\
553 next: check the id, or list what is there — `onetaskgraph task list` and \
554 `onetaskgraph project list` report what the configured sources hold."
555 )]
556 NoSuchItem {
557 id: String,
559 },
560
561 #[error(
566 "{item} was copied from {origin}, which that destination no longer holds\n\
567 next: re-run with --recreate to create a new item there instead, or restore \
568 {origin}."
569 )]
570 StaleOrigin {
571 item: String,
573 origin: String,
575 },
576
577 #[error(
583 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
584 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
585 {id} on its own with `onetaskgraph task copy`.",
586 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
587 )]
588 NotAMember {
589 id: GlobalId,
591 projects: Vec<GlobalId>,
593 },
594
595 #[error(
604 "{item} depends on {member}, which this copy was not told to carry and which records \
605 no origin in {destination}\n\
606 next: name {member} with --member as well, record its {destination} id at \
607 onetaskgraph.origin, or copy the whole project without --member."
608 )]
609 UnrecordedMember {
610 item: GlobalId,
612 member: GlobalId,
614 destination: SourceName,
616 },
617
618 #[error(
624 "source {name} could not do it: {error}\n\
625 next: fix what the source named above, then copy again."
626 )]
627 SourceRefused {
628 name: String,
630 error: SourceError,
632 },
633
634 #[error(
642 "the copy failed and could not be undone.\n\
643 it failed because: {error}\n\
644 it could not be undone because: {refusal}\n\
645 so the destination still holds: {left_behind}\n\
646 next: remove those items at the destination, then copy again."
647 )]
648 CopyNotUndone {
649 error: Box<EngineError>,
651 left_behind: LeftBehind,
657 refusal: SourceError,
659 },
660}
661
662#[derive(Debug, Clone, PartialEq)]
671pub struct LeftBehind {
672 first: GlobalId,
674 rest: Vec<GlobalId>,
676}
677
678impl LeftBehind {
679 #[must_use]
681 pub fn new(first: GlobalId) -> Self {
682 Self {
683 first,
684 rest: Vec::new(),
685 }
686 }
687
688 pub fn push(&mut self, id: GlobalId) {
690 self.rest.push(id);
691 }
692
693 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
695 std::iter::once(&self.first).chain(self.rest.iter())
696 }
697}
698
699impl std::fmt::Display for LeftBehind {
700 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
702 write!(formatter, "{}", self.first)?;
703 for id in &self.rest {
704 write!(formatter, ", {id}")?;
705 }
706 Ok(())
707 }
708}
709
710pub enum ConfiguredSource {
717 Ready(ResolvedSource),
719 Unavailable(UnavailableSource),
721}
722
723impl ConfiguredSource {
724 #[must_use]
726 pub fn name(&self) -> &SourceName {
727 match self {
728 Self::Ready(source) => source.name(),
729 Self::Unavailable(source) => source.name(),
730 }
731 }
732}
733
734pub struct Engine {
736 sources: Vec<ConfiguredSource>,
738 selection: Vec<SourceName>,
740}
741
742impl Engine {
743 #[must_use]
751 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
752 let (ready, unavailable) = resolve_available(config, secrets);
753 Self::new(
754 ready
755 .into_iter()
756 .map(ConfiguredSource::Ready)
757 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
758 .collect(),
759 config.selected_sources(),
760 )
761 }
762
763 #[must_use]
766 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
767 Self { sources, selection }
768 }
769
770 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
772 self.sources.iter().filter_map(|source| match source {
773 ConfiguredSource::Ready(ready) => Some(ready),
774 ConfiguredSource::Unavailable(_) => None,
775 })
776 }
777
778 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
780 self.sources.iter().filter_map(|source| match source {
781 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
782 ConfiguredSource::Ready(_) => None,
783 })
784 }
785
786 #[must_use]
788 pub fn listing(&self) -> Vec<SourceListing> {
789 let mut listings: Vec<SourceListing> = self
790 .ready()
791 .map(|source| SourceListing {
792 source: source.name().clone(),
793 kind: source.kind().to_owned(),
794 state: SourceState::Available {
795 capabilities: source.source().capabilities(),
796 },
797 })
798 .chain(self.unavailable().map(|source| SourceListing {
799 source: source.name().clone(),
800 kind: source.kind().to_owned(),
801 state: SourceState::Unavailable {
802 error: source.error().clone(),
803 },
804 }))
805 .collect();
806 listings.sort_by(|left, right| left.source.cmp(&right.source));
807 listings
808 }
809
810 #[must_use]
816 pub fn has(&self, name: &SourceName) -> bool {
817 self.sources.iter().any(|source| source.name() == name)
818 }
819
820 pub async fn tasks(
828 &self,
829 request: &TaskRequest,
830 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
831 let mut names = self.resolve_selection(&request.sources)?;
832 if let ProjectSelector::Qualified(id) = &request.project {
837 self.known(&id.source)?;
838 names.retain(|name| name == &id.source);
839 }
840 let query = shape(
841 "task-list",
842 &names,
843 &(&request.filters, &request.project, &request.priorities),
844 );
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 (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
855
856 let shapes: Vec<TaskShape> = ready
857 .iter()
858 .map(|source| {
859 shape_tasks(
860 &source.source().capabilities(),
861 &request.filters,
862 &project_filter(&request.project),
863 &request.priorities,
864 )
865 })
866 .collect();
867 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
868 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
869
870 let walks = ready
871 .iter()
872 .enumerate()
873 .map(|(index, source)| {
874 fetch_tasks(
875 source,
876 &shapes[index],
877 &starts[index],
878 budget,
879 &counters[index],
880 )
881 })
882 .collect();
883
884 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
885 answer.finish(
886 streams,
887 budget,
888 owed(&states),
889 &query,
890 |name, task: Task| {
891 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
892 },
893 )
894 }
895
896 pub async fn projects(
902 &self,
903 request: &ProjectRequest,
904 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
905 let names = self.resolve_selection(&request.sources)?;
906 let query = shape("project-list", &names, &request.filters);
907 let states = resumption(
908 self,
909 request.paging.token.as_ref(),
910 &[StreamKind::Items],
911 &query,
912 )?;
913 let budget = request.paging.limit.get();
914
915 let mut answer = Answer::new();
916 let mut with_projects = Vec::new();
922 for source in answer.split(self, &names) {
923 if source.source().capabilities().projects.is_native() {
924 with_projects.push(source);
925 } else {
926 answer.unreachable_predicate(source, Predicate::Project);
927 }
928 }
929 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
930
931 let shapes: Vec<ProjectShape> = ready
932 .iter()
933 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
934 .collect();
935 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
936 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
937
938 let walks = ready
939 .iter()
940 .enumerate()
941 .map(|(index, source)| {
942 fetch_projects(
943 source,
944 &shapes[index],
945 &starts[index],
946 budget,
947 &counters[index],
948 )
949 })
950 .collect();
951
952 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
953 answer.finish(
954 streams,
955 budget,
956 owed(&states),
957 &query,
958 |name, project: Project| Qualified {
959 id: GlobalId::new(name.clone(), project.id.clone()),
960 item: project,
961 },
962 )
963 }
964
965 pub async fn documents(
977 &self,
978 request: &DocumentRequest,
979 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
980 let mut names = self.resolve_selection(&request.sources)?;
981 if let ProjectSelector::Qualified(id) = &request.project {
984 self.known(&id.source)?;
985 names.retain(|name| name == &id.source);
986 }
987 let query = shape(
988 "document-list",
989 &names,
990 &(&request.filters, &request.project),
991 );
992 let states = resumption(
993 self,
994 request.paging.token.as_ref(),
995 &[StreamKind::Items],
996 &query,
997 )?;
998 let budget = request.paging.limit.get();
999
1000 let mut answer = Answer::new();
1001 let mut with_documents = Vec::new();
1002 for source in answer.split(self, &names) {
1003 if source.source().capabilities().documents.is_native() {
1004 with_documents.push(source);
1005 } else {
1006 answer.unreachable_predicate(source, Predicate::Document);
1007 }
1008 }
1009 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1010
1011 let shapes: Vec<DocumentShape> = ready
1012 .iter()
1013 .map(|source| {
1014 shape_documents(
1015 &source.source().capabilities(),
1016 &request.filters,
1017 &project_filter(&request.project),
1018 )
1019 })
1020 .collect();
1021 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1022 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1023
1024 let walks = ready
1025 .iter()
1026 .enumerate()
1027 .map(|(index, source)| {
1028 fetch_documents(
1029 source,
1030 &shapes[index],
1031 &starts[index],
1032 budget,
1033 &counters[index],
1034 )
1035 })
1036 .collect();
1037
1038 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1039 answer.finish(
1040 streams,
1041 budget,
1042 owed(&states),
1043 &query,
1044 |name, document: Document| Qualified {
1045 id: GlobalId::new(name.clone(), document.id.clone()),
1046 item: document,
1047 },
1048 )
1049 }
1050
1051 pub async fn labels(
1057 &self,
1058 request: &LabelRequest,
1059 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1060 let names = self.resolve_selection(&request.sources)?;
1061 let query = shape("label-list", &names, &());
1062 let states = resumption(
1063 self,
1064 request.paging.token.as_ref(),
1065 &[StreamKind::Items],
1066 &query,
1067 )?;
1068 let budget = request.paging.limit.get();
1069
1070 let mut answer = Answer::new();
1071 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1072 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1073 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1074
1075 let walks = ready
1076 .iter()
1077 .enumerate()
1078 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1079 .collect();
1080
1081 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1082 answer.finish(
1083 streams,
1084 budget,
1085 owed(&states),
1086 &query,
1087 |name, label: Label| Qualified {
1088 id: GlobalId::new(name.clone(), label.id.clone()),
1089 item: label,
1090 },
1091 )
1092 }
1093
1094 pub async fn search(
1100 &self,
1101 request: &SearchRequest,
1102 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1103 let names = self.resolve_selection(&request.sources)?;
1104 let reads: &[StreamKind] = match request.kind {
1108 SearchKind::Tasks => &[StreamKind::Tasks],
1109 SearchKind::Projects => &[StreamKind::Projects],
1110 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1111 };
1112 let query = shape("search", &names, &request.text);
1117 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1118 let budget = request.paging.limit.get();
1119 let filters = Filters {
1120 text: Some(request.text.clone()),
1121 ..Filters::default()
1122 };
1123
1124 let mut answer = Answer::new();
1125
1126 let mut ready = Vec::new();
1129 let mut kinds = Vec::new();
1130 let mut starts = Vec::new();
1131 for source in answer.split(self, &names) {
1132 let mut streams = Vec::new();
1133 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1134 streams.push(StreamKind::Tasks);
1135 }
1136 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1137 if source.source().capabilities().projects.is_native() {
1138 streams.push(StreamKind::Projects);
1139 } else {
1140 answer.unreachable_predicate(source, Predicate::Project);
1141 }
1142 }
1143 for stream in streams {
1144 if let Some(resume) = resume_at(&states, source.name(), stream) {
1145 ready.push(source);
1146 kinds.push(stream);
1147 starts.push(resume);
1148 }
1149 }
1150 }
1151
1152 let shapes: Vec<HitShape> = ready
1153 .iter()
1154 .zip(kinds.iter())
1155 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1156 .collect();
1157 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1158 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1159
1160 let walks = ready
1161 .iter()
1162 .enumerate()
1163 .map(|(index, source)| {
1164 fetch_hits(
1165 source,
1166 &shapes[index],
1167 &starts[index],
1168 budget,
1169 &counters[index],
1170 )
1171 })
1172 .collect();
1173
1174 let streams =
1175 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1176 answer.finish(
1177 streams,
1178 budget,
1179 owed(&states),
1180 &query,
1181 |name, found: Found| match found {
1182 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1183 GlobalId::new(name.clone(), task.id.clone()),
1184 task,
1185 )),
1186 Found::Project(project) => SearchHit::Project(Qualified {
1187 id: GlobalId::new(name.clone(), project.id.clone()),
1188 item: project,
1189 }),
1190 },
1191 )
1192 }
1193
1194 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1201 let name = self.known(&id.source)?;
1202 let mut answer = Answer::new();
1203 let selected = answer.split(self, std::slice::from_ref(&name));
1204 let Some(source) = selected.first() else {
1205 return answer.nothing();
1206 };
1207 let found = source.source().get_task(&id.native).await;
1208 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1209 answer.one(source, found, |task| {
1210 delivery::qualified_task(qualified, task)
1211 })
1212 }
1213
1214 pub async fn project(
1220 &self,
1221 id: &GlobalId,
1222 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1223 let name = self.known(&id.source)?;
1224 let mut answer = Answer::new();
1225 let selected = answer.split(self, std::slice::from_ref(&name));
1226 let Some(source) = selected.first() else {
1227 return answer.nothing();
1228 };
1229 let found = source.source().get_project(&id.native).await;
1230 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1231 answer.one(source, found, |project| Qualified {
1232 id: qualified,
1233 item: project,
1234 })
1235 }
1236
1237 pub async fn document(
1247 &self,
1248 id: &GlobalId,
1249 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1250 let name = self.known(&id.source)?;
1251 let mut answer = Answer::new();
1252 let selected = answer.split(self, std::slice::from_ref(&name));
1253 let Some(source) = selected.first() else {
1254 return answer.nothing();
1255 };
1256 if !source.source().capabilities().documents.is_native() {
1257 answer.unreachable_predicate(source, Predicate::Document);
1258 return answer.nothing();
1259 }
1260 let found = source.source().get_document(&id.native).await;
1261 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1262 answer.one(source, found, |document| Qualified {
1263 id: qualified,
1264 item: document,
1265 })
1266 }
1267
1268 pub async fn task_dependencies(
1275 &self,
1276 request: &DependencyRequest,
1277 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1278 self.dependencies(request, Entity::Task).await
1279 }
1280
1281 pub async fn project_dependencies(
1287 &self,
1288 request: &DependencyRequest,
1289 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1290 self.dependencies(request, Entity::Project).await
1291 }
1292
1293 async fn dependencies(
1296 &self,
1297 request: &DependencyRequest,
1298 entity: Entity,
1299 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1300 let name = self.known(&request.id.source)?;
1301 let query = shape(
1302 "dependencies",
1303 std::slice::from_ref(&name),
1304 &(entity, &request.id.native, request.direction),
1305 );
1306 let states = resumption(
1307 self,
1308 request.paging.token.as_ref(),
1309 &[StreamKind::Items],
1310 &query,
1311 )?;
1312 let budget = request.paging.limit.get();
1313
1314 let mut answer = Answer::new();
1315 let (ready, starts) = walking(
1316 answer.split(self, std::slice::from_ref(&name)),
1317 &states,
1318 StreamKind::Items,
1319 );
1320 let Some(source) = ready.first() else {
1321 return answer.nothing();
1322 };
1323
1324 let capabilities = source.source().capabilities();
1325 let support = match entity {
1326 Entity::Task => capabilities.task_dependencies,
1327 Entity::Project => capabilities.project_dependencies,
1328 };
1329 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1333 let mut outcomes = Outcomes::default();
1334 if request.direction == Direction::DependedOnBy {
1335 if emulating {
1336 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1337 } else {
1338 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1339 }
1340 }
1341
1342 let counters = vec![AtomicU32::new(0)];
1343 let walked = fetch_edges(
1344 source,
1345 &request.id.native,
1346 request.direction,
1347 entity,
1348 emulating,
1349 &starts[0],
1350 budget,
1351 &counters[0],
1352 )
1353 .await;
1354
1355 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1356 answer.finish(
1357 streams,
1358 budget,
1359 owed(&states),
1360 &query,
1361 |name, edge: DependencyEdge| QualifiedEdge {
1362 from: qualify_endpoint(name, edge.from),
1363 to: qualify_endpoint(name, edge.to),
1364 kind: edge.kind,
1365 },
1366 )
1367 }
1368
1369 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1371 if asked.is_empty() {
1372 if self.selection.is_empty() {
1373 return Err(EngineError::NoSources);
1374 }
1375 return Ok(self.selection.clone());
1376 }
1377 asked.iter().map(|name| self.known(name)).collect()
1378 }
1379
1380 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1382 if self.has(name) {
1383 return Ok(name.clone());
1384 }
1385 if self.sources.is_empty() {
1386 return Err(EngineError::NoSources);
1387 }
1388 Err(EngineError::UnknownSource {
1389 name: name.to_string(),
1390 configured: self
1391 .listing()
1392 .iter()
1393 .map(|listing| listing.source.to_string())
1394 .collect::<Vec<_>>()
1395 .join(", "),
1396 })
1397 }
1398}
1399
1400fn qualify_endpoint(
1401 source: &SourceName,
1402 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1403) -> QualifiedEndpoint {
1404 let kind = endpoint.kind;
1405 let is_qualified = endpoint.is_qualified();
1406 let endpoint_id = endpoint.into_id();
1407 QualifiedEndpoint {
1408 id: if is_qualified {
1409 endpoint_id
1410 .parse()
1411 .expect("plugin-api validates qualified dependency endpoints")
1412 } else {
1413 GlobalId::new(source.clone(), NativeId(endpoint_id))
1414 },
1415 kind,
1416 }
1417}
1418
1419#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1421enum Entity {
1422 Task,
1424 Project,
1426}
1427
1428enum Found {
1430 Task(Task),
1432 Project(Project),
1434}
1435
1436#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1438enum Outcome {
1439 PushedDown,
1441 AppliedLocally,
1443 Emulated,
1445 Unavailable,
1447}
1448
1449#[derive(Debug, Clone, Default, PartialEq)]
1461struct Outcomes(BTreeMap<Predicate, Outcome>);
1462
1463impl Outcomes {
1464 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1470 self.0.insert(predicate, outcome);
1471 }
1472
1473 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1476 for predicate in predicates {
1477 self.record(predicate, outcome);
1478 }
1479 }
1480
1481 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1483 self.0
1484 .iter()
1485 .filter(|(_, recorded)| **recorded == outcome)
1486 .map(|(predicate, _)| *predicate)
1487 .collect()
1488 }
1489}
1490
1491struct TaskShape {
1493 pushed: TaskQuery,
1495 local: LocalTasks,
1497 outcomes: Outcomes,
1499}
1500
1501struct ProjectShape {
1503 pushed: ProjectQuery,
1505 local: LocalProjects,
1507 outcomes: Outcomes,
1509}
1510
1511struct DocumentShape {
1513 pushed: DocumentQuery,
1515 local: LocalDocuments,
1517 outcomes: Outcomes,
1519}
1520
1521struct HitShape {
1523 stream: StreamKind,
1525 tasks: TaskQuery,
1527 projects: ProjectQuery,
1529 local_tasks: LocalTasks,
1531 local_projects: LocalProjects,
1533 outcomes: Outcomes,
1535}
1536
1537struct Answer {
1543 plans: Vec<SourcePlan>,
1545 errors: Vec<SourceFailure>,
1547}
1548
1549impl Answer {
1550 fn new() -> Self {
1551 Self {
1552 plans: Vec::new(),
1553 errors: Vec::new(),
1554 }
1555 }
1556
1557 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1562 let mut selected = Vec::new();
1563 for name in names {
1564 match engine.sources.iter().find(|source| source.name() == name) {
1565 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1566 Some(ConfiguredSource::Unavailable(source)) => {
1567 self.errors.push(source.failure());
1568 }
1569 None => {}
1570 }
1571 }
1572 selected
1573 }
1574
1575 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1577 let mut outcomes = Outcomes::default();
1578 outcomes.record(predicate, Outcome::Unavailable);
1579 self.plans.push(plan_for(source, outcomes, 0));
1580 }
1581
1582 fn collect<T>(
1584 &mut self,
1585 ready: &[&ResolvedSource],
1586 walked: Vec<Result<Fetched<T>, SourceError>>,
1587 counters: &[AtomicU32],
1588 outcomes: Vec<Outcomes>,
1589 ) -> Vec<Stream<T>> {
1590 let kinds = vec![StreamKind::Items; ready.len()];
1591 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1592 }
1593
1594 fn collect_streams<T>(
1596 &mut self,
1597 ready: &[&ResolvedSource],
1598 kinds: &[StreamKind],
1599 walked: Vec<Result<Fetched<T>, SourceError>>,
1600 counters: &[AtomicU32],
1601 outcomes: Vec<Outcomes>,
1602 ) -> Vec<Stream<T>> {
1603 let mut streams = Vec::new();
1604 for (index, result) in walked.into_iter().enumerate() {
1605 let source = ready[index];
1606 let pages = counters[index].load(Ordering::Relaxed);
1607 self.plans
1608 .push(plan_for(source, outcomes[index].clone(), pages));
1609 match result {
1610 Ok(fetched) => streams.push(Stream {
1611 source: source.name().clone(),
1612 kind: kinds[index],
1613 fetched,
1614 }),
1615 Err(error) => self.errors.push(SourceFailure {
1618 source: source.name().clone(),
1619 error,
1620 }),
1621 }
1622 }
1623 streams
1624 }
1625
1626 fn one<T, U>(
1628 mut self,
1629 source: &ResolvedSource,
1630 found: Result<Option<T>, SourceError>,
1631 qualify: impl FnOnce(T) -> U,
1632 ) -> Result<QueryResponse<U>, EngineError> {
1633 self.plans.push(plan_for(source, Outcomes::default(), 1));
1634 let items = match found {
1635 Ok(Some(item)) => vec![qualify(item)],
1636 Ok(None) => Vec::new(),
1637 Err(error) => {
1638 self.errors.push(SourceFailure {
1639 source: source.name().clone(),
1640 error,
1641 });
1642 Vec::new()
1643 }
1644 };
1645 Ok(QueryResponse {
1646 items,
1647 next: None,
1648 plan: QueryPlan {
1649 per_source: merge_plans(self.plans),
1650 },
1651 errors: self.errors,
1652 })
1653 }
1654
1655 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1657 Ok(QueryResponse {
1658 items: Vec::new(),
1659 next: None,
1660 plan: QueryPlan {
1661 per_source: merge_plans(self.plans),
1662 },
1663 errors: self.errors,
1664 })
1665 }
1666
1667 fn finish<T, U>(
1672 self,
1673 streams: Vec<Stream<T>>,
1674 budget: u32,
1675 first: Option<&Owed>,
1676 query: &str,
1677 qualify: impl Fn(&SourceName, T) -> U,
1678 ) -> Result<QueryResponse<U>, EngineError> {
1679 let (rows, states, owed) = merge(streams, budget, first);
1680 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1681 Ok(QueryResponse {
1682 items: rows
1683 .into_iter()
1684 .map(|(name, item)| qualify(&name, item))
1685 .collect(),
1686 next,
1687 plan: QueryPlan {
1688 per_source: merge_plans(self.plans),
1689 },
1690 errors: self.errors,
1691 })
1692 }
1693}
1694
1695fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1698 SourcePlan {
1699 source: source.name().clone(),
1700 kind: source.kind().to_owned(),
1701 pushed_down: outcomes.with(Outcome::PushedDown),
1702 applied_locally: outcomes.with(Outcome::AppliedLocally),
1703 emulated: outcomes.with(Outcome::Emulated),
1704 unavailable: outcomes.with(Outcome::Unavailable),
1705 pages_fetched: pages,
1706 }
1707}
1708
1709fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1714 let mut merged: Vec<SourcePlan> = Vec::new();
1715 for plan in plans {
1716 if let Some(existing) = merged
1717 .iter_mut()
1718 .find(|existing| existing.source == plan.source)
1719 {
1720 existing.pushed_down.extend(plan.pushed_down);
1721 existing.applied_locally.extend(plan.applied_locally);
1722 existing.emulated.extend(plan.emulated);
1723 existing.unavailable.extend(plan.unavailable);
1724 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1725 for list in [
1726 &mut existing.pushed_down,
1727 &mut existing.applied_locally,
1728 &mut existing.emulated,
1729 &mut existing.unavailable,
1730 ] {
1731 list.sort_unstable();
1732 list.dedup();
1733 }
1734 } else {
1735 merged.push(plan);
1736 }
1737 }
1738 merged
1739}
1740
1741fn walking<'a>(
1747 selected: Vec<&'a ResolvedSource>,
1748 states: &Option<Resumption>,
1749 kind: StreamKind,
1750) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1751 let mut ready = Vec::new();
1752 let mut starts = Vec::new();
1753 for source in selected {
1754 if let Some(resume) = resume_at(states, source.name(), kind) {
1755 ready.push(source);
1756 starts.push(resume);
1757 }
1758 }
1759 (ready, starts)
1760}
1761
1762fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1780 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1781 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1782}
1783
1784fn fingerprint(text: &str) -> String {
1786 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1787 for byte in text.as_bytes() {
1788 hash ^= u64::from(*byte);
1789 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1790 }
1791 format!("{hash:016x}")
1792}
1793
1794fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1799 document.as_ref()?.owed.as_ref()
1800}
1801
1802fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1804 match states {
1805 None => Some(Resume::default()),
1806 Some(document) => document
1807 .streams
1808 .iter()
1809 .find(|state| &state.source == source && state.stream == kind)
1810 .map(|state| state.resume.clone()),
1811 }
1812}
1813
1814fn resumption(
1845 engine: &Engine,
1846 token: Option<&PageToken>,
1847 reads: &[StreamKind],
1848 query: &str,
1849) -> Result<Option<Resumption>, EngineError> {
1850 let Some(document) = token.map(PageToken::decode) else {
1851 return Ok(None);
1852 };
1853
1854 if document.query != query {
1861 return Err(EngineError::Token {
1862 message: "this page token was written by a different query — resume the walk it \
1863 came from, or drop --page to start this one from the beginning"
1864 .to_owned(),
1865 });
1866 }
1867 let states = &document.streams;
1868
1869 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
1870 for state in states {
1871 if !reads.contains(&state.stream) {
1872 return Err(EngineError::Token {
1873 message: format!(
1874 "this page token resumes {}, which this command does not read — it \
1875 was written by a different query",
1876 state.stream.describe()
1877 ),
1878 });
1879 }
1880 let ceiling = engine
1881 .ready()
1882 .find(|source| source.name() == &state.source)
1883 .map(ceiling);
1884 if ceiling.is_none() && !engine.has(&state.source) {
1885 return Err(EngineError::Token {
1886 message: format!(
1887 "this page token resumes a source called {:?}, which this \
1888 configuration does not have",
1889 state.source.as_str()
1890 ),
1891 });
1892 }
1893 if let Some(ceiling) = ceiling
1894 && state.resume.skip >= ceiling
1895 {
1896 return Err(EngineError::Token {
1897 message: format!(
1898 "this page token resumes {} rows into a page of source {:?}, which \
1899 serves at most {ceiling}",
1900 state.resume.skip,
1901 state.source.as_str()
1902 ),
1903 });
1904 }
1905 if seen.contains(&(&state.source, state.stream)) {
1906 return Err(EngineError::Token {
1907 message: format!(
1908 "this page token gives source {:?} two places to resume from",
1909 state.source.as_str()
1910 ),
1911 });
1912 }
1913 seen.push((&state.source, state.stream));
1914 }
1915
1916 if let Some(owed) = &document.owed
1921 && !document
1922 .streams
1923 .iter()
1924 .any(|state| state.source == owed.source && state.stream == owed.stream)
1925 {
1926 return Err(EngineError::Token {
1927 message: format!(
1928 "this page token owes the next row to a stream it does not resume, \
1929 {:?}'s {}",
1930 owed.source.as_str(),
1931 owed.stream.describe()
1932 ),
1933 });
1934 }
1935
1936 Ok(Some(document))
1937}
1938
1939fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
1945 match selector {
1946 ProjectSelector::Any => ProjectFilter::Any,
1947 ProjectSelector::Orphans => ProjectFilter::Orphans,
1948 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
1949 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
1950 }
1951}
1952
1953fn text_predicates(fields: TextFields) -> Vec<Predicate> {
1955 match fields {
1956 TextFields::Title => vec![Predicate::SearchTitle],
1957 TextFields::Content => vec![Predicate::SearchContent],
1958 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
1959 }
1960}
1961
1962fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
1969 match fields {
1970 TextFields::Title => capabilities.search_title.is_native(),
1971 TextFields::Content => capabilities.search_content.is_native(),
1972 TextFields::TitleOrContent => {
1973 capabilities.search_title.is_native() && capabilities.search_content.is_native()
1974 }
1975 }
1976}
1977
1978fn shape_tasks(
1980 capabilities: &Capabilities,
1981 filters: &Filters,
1982 project: &ProjectFilter,
1983 priorities: &[Priority],
1984) -> TaskShape {
1985 let mut pushed = TaskQuery::default();
1986 let mut local = LocalTasks::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 !priorities.is_empty() {
2008 if capabilities.filter_by_priority.is_native() {
2009 pushed.priorities = priorities.to_vec();
2010 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2011 } else {
2012 local.priorities = priorities.to_vec();
2013 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2014 }
2015 }
2016 if let Some(text) = &filters.text {
2017 let predicates = text_predicates(text.fields);
2018 if searches_natively(capabilities, text.fields) {
2019 pushed.text = Some(text.clone());
2020 outcomes.record_all(predicates, Outcome::PushedDown);
2021 } else {
2022 local.text = Some(text.clone());
2023 outcomes.record_all(predicates, Outcome::AppliedLocally);
2024 }
2025 }
2026 match project {
2027 ProjectFilter::Any => {}
2028 ProjectFilter::Orphans => {
2029 if capabilities.orphan_tasks.is_native() {
2030 pushed.project = ProjectFilter::Orphans;
2031 outcomes.record(Predicate::Project, Outcome::PushedDown);
2032 } else {
2033 local.project = Some(ProjectFilter::Orphans);
2034 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2035 }
2036 }
2037 ProjectFilter::Is(id) => {
2038 if capabilities.projects.is_native() {
2039 pushed.project = ProjectFilter::Is(id.clone());
2040 outcomes.record(Predicate::Project, Outcome::PushedDown);
2041 } else {
2042 local.project = Some(ProjectFilter::Is(id.clone()));
2043 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2044 }
2045 }
2046 }
2047
2048 TaskShape {
2049 pushed,
2050 local,
2051 outcomes,
2052 }
2053}
2054
2055fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2057 let mut pushed = ProjectQuery::default();
2058 let mut local = LocalProjects::default();
2059 let mut outcomes = Outcomes::default();
2060
2061 if !filters.labels.is_empty() {
2062 if capabilities.filter_by_label.is_native() {
2063 pushed.labels = filters.labels.clone();
2064 outcomes.record(Predicate::Label, Outcome::PushedDown);
2065 } else {
2066 local.labels = Some(filters.labels.clone());
2067 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2068 }
2069 }
2070 if !filters.statuses.is_empty() {
2071 if capabilities.filter_by_status.is_native() {
2072 pushed.statuses.clone_from(&filters.statuses);
2073 outcomes.record(Predicate::Status, Outcome::PushedDown);
2074 } else {
2075 local.statuses.clone_from(&filters.statuses);
2076 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2077 }
2078 }
2079 if let Some(text) = &filters.text {
2080 let predicates = text_predicates(text.fields);
2081 if searches_natively(capabilities, text.fields) {
2082 pushed.text = Some(text.clone());
2083 outcomes.record_all(predicates, Outcome::PushedDown);
2084 } else {
2085 local.text = Some(text.clone());
2086 outcomes.record_all(predicates, Outcome::AppliedLocally);
2087 }
2088 }
2089
2090 ProjectShape {
2091 pushed,
2092 local,
2093 outcomes,
2094 }
2095}
2096
2097fn shape_documents(
2102 capabilities: &Capabilities,
2103 filters: &DocumentFilters,
2104 project: &ProjectFilter,
2105) -> DocumentShape {
2106 let mut pushed = DocumentQuery::default();
2107 let mut local = LocalDocuments::default();
2108 let mut outcomes = Outcomes::default();
2109
2110 if !filters.labels.is_empty() {
2111 if capabilities.filter_by_label.is_native() {
2112 pushed.labels = filters.labels.clone();
2113 outcomes.record(Predicate::Label, Outcome::PushedDown);
2114 } else {
2115 local.labels = Some(filters.labels.clone());
2116 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2117 }
2118 }
2119 if let Some(text) = &filters.text {
2120 let predicates = text_predicates(text.fields);
2121 if searches_natively(capabilities, text.fields) {
2122 pushed.text = Some(text.clone());
2123 outcomes.record_all(predicates, Outcome::PushedDown);
2124 } else {
2125 local.text = Some(text.clone());
2126 outcomes.record_all(predicates, Outcome::AppliedLocally);
2127 }
2128 }
2129 match project {
2130 ProjectFilter::Any => {}
2131 ProjectFilter::Orphans => {
2132 if capabilities.orphan_tasks.is_native() {
2133 pushed.project = ProjectFilter::Orphans;
2134 outcomes.record(Predicate::Project, Outcome::PushedDown);
2135 } else {
2136 local.project = Some(ProjectFilter::Orphans);
2137 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2138 }
2139 }
2140 ProjectFilter::Is(id) => {
2141 if capabilities.projects.is_native() {
2142 pushed.project = ProjectFilter::Is(id.clone());
2143 outcomes.record(Predicate::Project, Outcome::PushedDown);
2144 } else {
2145 local.project = Some(ProjectFilter::Is(id.clone()));
2146 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2147 }
2148 }
2149 }
2150
2151 DocumentShape {
2152 pushed,
2153 local,
2154 outcomes,
2155 }
2156}
2157
2158fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2160 match stream {
2161 StreamKind::Projects => {
2162 let shaped = shape_projects(capabilities, filters);
2163 HitShape {
2164 stream,
2165 tasks: TaskQuery::default(),
2166 projects: shaped.pushed,
2167 local_tasks: LocalTasks::default(),
2168 local_projects: shaped.local,
2169 outcomes: shaped.outcomes,
2170 }
2171 }
2172 StreamKind::Items | StreamKind::Tasks => {
2173 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any, &[]);
2174 HitShape {
2175 stream,
2176 tasks: shaped.pushed,
2177 projects: ProjectQuery::default(),
2178 local_tasks: shaped.local,
2179 local_projects: LocalProjects::default(),
2180 outcomes: shaped.outcomes,
2181 }
2182 }
2183 }
2184}
2185
2186fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2193 if compensating {
2194 ceiling
2195 } else {
2196 budget.min(ceiling)
2197 }
2198}
2199
2200fn ceiling(source: &ResolvedSource) -> u32 {
2202 source.source().capabilities().max_page_size.max(1)
2203}
2204
2205async fn fetch_tasks(
2207 source: &ResolvedSource,
2208 shape: &TaskShape,
2209 start: &Resume,
2210 budget: u32,
2211 calls: &AtomicU32,
2212) -> Result<Fetched<Task>, SourceError> {
2213 let compensating = shape.local != LocalTasks::default();
2214 walk(
2215 start,
2216 budget,
2217 page_size(compensating, budget, ceiling(source)),
2218 |task| shape.local.keeps(task),
2219 |cursor, limit| async move {
2220 calls.fetch_add(1, Ordering::Relaxed);
2221 let request = PageRequest { cursor, limit };
2222 source.source().query_tasks(&shape.pushed, &request).await
2223 },
2224 )
2225 .await
2226}
2227
2228async fn fetch_projects(
2230 source: &ResolvedSource,
2231 shape: &ProjectShape,
2232 start: &Resume,
2233 budget: u32,
2234 calls: &AtomicU32,
2235) -> Result<Fetched<Project>, SourceError> {
2236 let compensating = shape.local != LocalProjects::default();
2237 walk(
2238 start,
2239 budget,
2240 page_size(compensating, budget, ceiling(source)),
2241 |project| shape.local.keeps(project),
2242 |cursor, limit| async move {
2243 calls.fetch_add(1, Ordering::Relaxed);
2244 let request = PageRequest { cursor, limit };
2245 source
2246 .source()
2247 .query_projects(&shape.pushed, &request)
2248 .await
2249 },
2250 )
2251 .await
2252}
2253
2254async fn fetch_documents(
2256 source: &ResolvedSource,
2257 shape: &DocumentShape,
2258 start: &Resume,
2259 budget: u32,
2260 calls: &AtomicU32,
2261) -> Result<Fetched<Document>, SourceError> {
2262 let compensating = shape.local != LocalDocuments::default();
2263 walk(
2264 start,
2265 budget,
2266 page_size(compensating, budget, ceiling(source)),
2267 |document| shape.local.keeps(document),
2268 |cursor, limit| async move {
2269 calls.fetch_add(1, Ordering::Relaxed);
2270 let request = PageRequest { cursor, limit };
2271 source
2272 .source()
2273 .query_documents(&shape.pushed, &request)
2274 .await
2275 },
2276 )
2277 .await
2278}
2279
2280async fn fetch_labels(
2282 source: &ResolvedSource,
2283 start: &Resume,
2284 budget: u32,
2285 calls: &AtomicU32,
2286) -> Result<Fetched<Label>, SourceError> {
2287 walk(
2288 start,
2289 budget,
2290 page_size(false, budget, ceiling(source)),
2291 |_| true,
2292 |cursor, limit| async move {
2293 calls.fetch_add(1, Ordering::Relaxed);
2294 let request = PageRequest { cursor, limit };
2295 source.source().labels(&request).await
2296 },
2297 )
2298 .await
2299}
2300
2301async fn fetch_hits(
2303 source: &ResolvedSource,
2304 shape: &HitShape,
2305 start: &Resume,
2306 budget: u32,
2307 calls: &AtomicU32,
2308) -> Result<Fetched<Found>, SourceError> {
2309 let ceiling = ceiling(source);
2310 match shape.stream {
2311 StreamKind::Projects => {
2312 let compensating = shape.local_projects != LocalProjects::default();
2313 walk(
2314 start,
2315 budget,
2316 page_size(compensating, budget, ceiling),
2317 |found| match found {
2318 Found::Project(project) => shape.local_projects.keeps(project),
2319 Found::Task(_) => true,
2320 },
2321 |cursor, limit| async move {
2322 calls.fetch_add(1, Ordering::Relaxed);
2323 let request = PageRequest { cursor, limit };
2324 let page = source
2325 .source()
2326 .query_projects(&shape.projects, &request)
2327 .await?;
2328 Ok(Page {
2329 items: page.items.into_iter().map(Found::Project).collect(),
2330 next: page.next,
2331 })
2332 },
2333 )
2334 .await
2335 }
2336 StreamKind::Items | StreamKind::Tasks => {
2337 let compensating = shape.local_tasks != LocalTasks::default();
2338 walk(
2339 start,
2340 budget,
2341 page_size(compensating, budget, ceiling),
2342 |found| match found {
2343 Found::Task(task) => shape.local_tasks.keeps(task),
2344 Found::Project(_) => true,
2345 },
2346 |cursor, limit| async move {
2347 calls.fetch_add(1, Ordering::Relaxed);
2348 let request = PageRequest { cursor, limit };
2349 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2350 Ok(Page {
2351 items: page.items.into_iter().map(Found::Task).collect(),
2352 next: page.next,
2353 })
2354 },
2355 )
2356 .await
2357 }
2358 }
2359}
2360
2361async fn forward_edges(
2363 source: &ResolvedSource,
2364 entity: Entity,
2365 id: &NativeId,
2366 request: &PageRequest,
2367) -> Result<Page<DependencyEdge>, SourceError> {
2368 match entity {
2369 Entity::Task => {
2370 source
2371 .source()
2372 .task_dependencies(id, Direction::DependsOn, request)
2373 .await
2374 }
2375 Entity::Project => {
2376 source
2377 .source()
2378 .project_dependencies(id, Direction::DependsOn, request)
2379 .await
2380 }
2381 }
2382}
2383
2384#[expect(
2393 clippy::too_many_arguments,
2394 reason = "every argument is one axis of one walk — the source, the item, the \
2395 direction, which of its two graphs, whether the reverse is emulated, where \
2396 to resume, how many rows to return and where to count calls. Grouping them \
2397 into a struct would name the same eight values one indirection further from \
2398 the loop that reads them."
2399)]
2400async fn fetch_edges(
2401 source: &ResolvedSource,
2402 native: &NativeId,
2403 direction: Direction,
2404 entity: Entity,
2405 emulating: bool,
2406 start: &Resume,
2407 budget: u32,
2408 calls: &AtomicU32,
2409) -> Result<Fetched<DependencyEdge>, SourceError> {
2410 let ceiling = ceiling(source);
2411 if !emulating {
2412 return walk(
2413 start,
2414 budget,
2415 page_size(false, budget, ceiling),
2416 |_| true,
2417 |cursor, limit| async move {
2418 calls.fetch_add(1, Ordering::Relaxed);
2419 let request = PageRequest { cursor, limit };
2420 match entity {
2421 Entity::Task => {
2422 source
2423 .source()
2424 .task_dependencies(native, direction, &request)
2425 .await
2426 }
2427 Entity::Project => {
2428 source
2429 .source()
2430 .project_dependencies(native, direction, &request)
2431 .await
2432 }
2433 }
2434 },
2435 )
2436 .await;
2437 }
2438
2439 walk(
2440 start,
2441 budget,
2442 ceiling,
2443 |_| true,
2444 |cursor, limit| async move {
2445 calls.fetch_add(1, Ordering::Relaxed);
2446 let request = PageRequest { cursor, limit };
2447 let (ids, next) = match entity {
2448 Entity::Task => {
2449 let page = source
2450 .source()
2451 .query_tasks(&TaskQuery::default(), &request)
2452 .await?;
2453 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2454 (ids, page.next)
2455 }
2456 Entity::Project => {
2457 let page = source
2458 .source()
2459 .query_projects(&ProjectQuery::default(), &request)
2460 .await?;
2461 let ids: Vec<NativeId> =
2462 page.items.into_iter().map(|project| project.id).collect();
2463 (ids, page.next)
2464 }
2465 };
2466
2467 let mut edges = Vec::new();
2468 for id in ids {
2469 let mut inner: Option<Cursor> = None;
2470 loop {
2471 calls.fetch_add(1, Ordering::Relaxed);
2472 let request = PageRequest {
2473 cursor: inner.clone(),
2474 limit,
2475 };
2476 let page = forward_edges(source, entity, &id, &request).await?;
2477 fits(page.items.len(), limit)?;
2481 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2482 unrepeated(
2483 page.next.as_ref(),
2484 inner.as_ref(),
2485 "its forward edges were being scanned",
2486 )?;
2487 match page.next {
2488 Some(cursor) => inner = Some(cursor),
2489 None => break,
2490 }
2491 }
2492 }
2493
2494 Ok(Page { items: edges, next })
2495 },
2496 )
2497 .await
2498}