1mod comment;
19mod copy;
20mod delivery;
21mod fetch;
22mod join;
23mod local;
24mod metadata;
25mod narrow;
26mod rendered;
27mod resume;
28mod update;
29
30use std::collections::BTreeMap;
31use std::num::NonZeroU32;
32use std::sync::atomic::{AtomicU32, Ordering};
33
34use onetaskgraph_plugin_api::{
35 Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
36 MetadataRecord, NativeId, Page, PageRequest, Priority, Project, ProjectFilter, ProjectQuery,
37 SecretResolver, SourceError, SourceName, StatusCategory, Task, TaskQuery, TextFields,
38 TextQuery,
39};
40use schemars::JsonSchema;
41use serde::{Deserialize, Serialize};
42
43use crate::GlobalId;
44use crate::config::Config;
45use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
46use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
47
48use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
49use join::join_all;
50use local::{LocalDocuments, LocalProjects, LocalTasks};
51pub(crate) use resume::{Owed, Resumption, StreamState};
52use resume::{Resume, StreamKind};
53
54pub use comment::{CommentList, DeletedComment, TaskDetail};
55pub use copy::{
56 BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
57 Spent,
58};
59pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
60pub use local::ProjectSelector;
61pub use metadata::MetadataSet;
62pub use narrow::{TaskContentSet, TaskPrioritySet};
63pub use rendered::{
64 Body, DocumentCreate, Regenerated, Regeneration, RenderRequest, RenderTemplate, RenderedRecord,
65 TaskCreate, TaskCreated, TemplateAnswers, UnusedAnswers,
66};
67pub use update::TaskUpdated;
68
69#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
74pub struct Qualified<T> {
75 pub id: GlobalId,
77 pub item: T,
79}
80
81#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
86pub struct QualifiedEdge {
87 pub from: QualifiedEndpoint,
92 pub to: QualifiedEndpoint,
94 pub kind: onetaskgraph_plugin_api::DependencyKind,
96}
97
98#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
100pub struct QualifiedEndpoint {
101 pub id: GlobalId,
103 pub kind: onetaskgraph_plugin_api::ItemKind,
105}
106
107impl std::fmt::Display for QualifiedEndpoint {
108 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
109 self.id.fmt(formatter)
110 }
111}
112
113#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
115#[serde(tag = "kind", rename_all = "kebab-case")]
116pub enum SearchHit {
117 Task(Qualified<Task>),
119 Project(Qualified<Project>),
121}
122
123#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
125#[serde(rename_all = "kebab-case")]
126pub enum SearchKind {
127 Tasks,
129 Projects,
131 #[default]
133 Both,
134}
135
136#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
138pub struct SourceListing {
139 pub source: SourceName,
141 pub kind: String,
153 #[serde(flatten)]
155 pub state: SourceState,
156}
157
158#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
160#[serde(tag = "state", rename_all = "kebab-case")]
161pub enum SourceState {
162 Available {
164 capabilities: Capabilities,
166 },
167 Unavailable {
169 error: SourceError,
171 },
172}
173
174#[derive(Debug, Clone, PartialEq)]
176pub struct Paging {
177 pub limit: NonZeroU32,
179 pub token: Option<PageToken>,
181}
182
183#[derive(Debug, Clone, Default, PartialEq)]
185pub struct Filters {
186 pub text: Option<TextQuery>,
188 pub labels: LabelFilter,
190 pub statuses: Vec<StatusCategory>,
192}
193
194#[derive(Debug, Clone)]
196pub struct TaskRequest {
197 pub sources: Vec<SourceName>,
199 pub filters: Filters,
201 pub project: ProjectSelector,
203 pub priorities: Vec<Priority>,
209 pub paging: Paging,
211}
212
213#[derive(Debug, Clone)]
215pub struct ProjectRequest {
216 pub sources: Vec<SourceName>,
218 pub filters: Filters,
220 pub paging: Paging,
222}
223
224#[derive(Debug, Clone, Default, PartialEq)]
231pub struct DocumentFilters {
232 pub text: Option<TextQuery>,
234 pub labels: LabelFilter,
236}
237
238#[derive(Debug, Clone)]
240pub struct DocumentRequest {
241 pub sources: Vec<SourceName>,
243 pub filters: DocumentFilters,
245 pub project: ProjectSelector,
247 pub paging: Paging,
249}
250
251#[derive(Debug, Clone)]
253pub struct LabelRequest {
254 pub sources: Vec<SourceName>,
256 pub paging: Paging,
258}
259
260#[derive(Debug, Clone)]
262pub struct SearchRequest {
263 pub sources: Vec<SourceName>,
265 pub text: TextQuery,
267 pub kind: SearchKind,
269 pub paging: Paging,
271}
272
273#[derive(Debug, Clone)]
275pub struct DependencyRequest {
276 pub id: GlobalId,
278 pub direction: Direction,
280 pub paging: Paging,
282}
283
284#[derive(Debug, Clone, PartialEq, thiserror::Error)]
292pub enum EngineError {
293 #[error(
295 "no source named {name:?} is configured\n\
296 next: name one of the configured sources ({configured}), or add {name:?} under \
297 `sources` — `onetaskgraph sources list` shows what this configuration has."
298 )]
299 UnknownSource {
300 name: String,
302 configured: String,
304 },
305
306 #[error(
308 "{message}\n\
309 next: page with a token exactly as the previous page reported it, and against \
310 the same configuration — or drop `--page` to start the walk again."
311 )]
312 Token {
313 message: String,
315 },
316
317 #[error(
319 "no sources are configured\n\
320 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
321 prints what each plugin accepts."
322 )]
323 NoSources,
324
325 #[error(
327 "source {name} cannot be written: its plugin is {kind}, which has no write \
328 side\n\
329 next: copy into a source whose plugin can be written — `onetaskgraph sources \
330 list` reports each one's plugin."
331 )]
332 NotWritable {
333 name: String,
335 kind: String,
337 },
338
339 #[error(
346 "source {name} has no documents: its plugin is {kind}, which holds none\n\
347 next: name a source whose plugin has documents — `onetaskgraph sources list` \
348 reports each one's plugin and what it declares."
349 )]
350 NoDocuments {
351 name: String,
353 kind: String,
355 },
356
357 #[error(
362 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
363 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
364 list` reports each one's plugin."
365 )]
366 NoComments {
367 name: String,
369 kind: String,
371 },
372
373 #[error(
375 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
376 but not added to, edited or removed\n\
377 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
378 source whose plugin can be written — `onetaskgraph sources list` reports each \
379 one's plugin."
380 )]
381 CommentsNotWritable {
382 name: String,
384 kind: String,
386 },
387
388 #[error(
390 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
391 next: set the status in that source itself, or name a task of a source whose plugin \
392 can be written — `onetaskgraph sources list` reports each one's plugin."
393 )]
394 StatusNotWritable {
395 name: String,
397 kind: String,
399 },
400
401 #[error(
403 "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
404 write side\n\
405 next: set the key in that source itself, or name a {record} of a source whose plugin \
406 can be written — `onetaskgraph sources list` reports each one's plugin."
407 )]
408 MetadataNotWritable {
409 name: String,
411 kind: String,
413 record: MetadataRecord,
415 },
416
417 #[error(
419 "source {name} cannot write a priority: its plugin is {kind}, which has no write side\n\
420 next: set the priority in that source itself, or name a task of a source whose plugin \
421 can be written — `onetaskgraph sources list` reports each one's plugin."
422 )]
423 PriorityNotWritable {
424 name: String,
426 kind: String,
428 },
429
430 #[error(
432 "source {name} cannot write a task's content: its plugin is {kind}, which has no write \
433 side\n\
434 next: edit the content in that source itself, or name a task of a source whose plugin \
435 can be written — `onetaskgraph sources list` reports each one's plugin."
436 )]
437 ContentNotWritable {
438 name: String,
440 kind: String,
442 },
443
444 #[error(
446 "source {name} cannot update a task: its plugin is {kind}, which has no write side\n\
447 next: change the task in that source itself, or name a task of a source whose plugin \
448 can be written — `onetaskgraph sources list` reports each one's plugin."
449 )]
450 UpdateNotWritable {
451 name: String,
453 kind: String,
455 },
456
457 #[error(
459 "source {name} cannot create a {record}: its plugin is {kind}, which has no write side\n\
460 next: create it in a source whose plugin can be written — `onetaskgraph sources list` \
461 reports each one's plugin."
462 )]
463 NotCreatable {
464 name: String,
466 kind: String,
468 record: MetadataRecord,
470 },
471
472 #[error(
474 "source {name} cannot write a {record}'s rendering: its plugin is {kind}, which has no \
475 write side\n\
476 next: render it with --dry-run to read the result, or regenerate a {record} of a source \
477 whose plugin can be written."
478 )]
479 RenderingNotWritable {
480 name: String,
482 kind: String,
484 record: RenderedRecord,
486 },
487
488 #[error(
491 "supply every required answer to regenerate {id}: {} unanswered, and {reason}\n\
492 next: answer {} with --var NAME=VALUE or an answers file (--answers FILE), or run \
493 interactively to be asked.",
494 names.join(", "),
495 if names.len() == 1 { "it" } else { "each" }
496 )]
497 MissingAnswers {
498 id: String,
500 names: Vec<String>,
502 reason: String,
505 },
506
507 #[error("{error}")]
509 Template {
510 error: crate::template::TemplateError,
512 },
513
514 #[error(
516 "{record} {id} has no stored template answers: {reason}\n\
517 next: regenerate it with every required answer (`onetaskgraph {record} render {id} \
518 --var NAME=VALUE`), or read its provenance with `onetaskgraph {record} show {id}`."
519 )]
520 NoStoredAnswers {
521 record: RenderedRecord,
523 id: String,
525 reason: String,
527 },
528
529 #[error(
531 "{record} {id} records no template it was rendered from, and none was given\n\
532 next: name one with --template FILE or --template-loader FILE."
533 )]
534 NoTemplate {
535 record: RenderedRecord,
537 id: String,
539 },
540
541 #[error(
544 "{record} {id} records a template entry this product did not write — {problem}\n\
545 next: regenerate it with --template FILE or --template-loader FILE and every required \
546 answer, which records a fresh entry."
547 )]
548 MalformedProvenance {
549 record: RenderedRecord,
551 id: String,
553 problem: String,
555 },
556
557 #[error(
559 "{record} {id} was rendered from {reference:?}, which is not a readable file, so it \
560 cannot be re-read; a recorded reference is never turned into a location\n\
561 next: supply the template with --template-loader FILE (a loader document naming what \
562 to render), or name a template file with --template FILE."
563 )]
564 TemplateNotAFile {
565 record: RenderedRecord,
567 id: String,
569 reference: String,
571 },
572
573 #[error(
578 "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
579 be written to it: its plugin is {kind}, which declares priority unsupported\n\
580 next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
581 reports what each declares — or set the task's priority to none first; a \
582 github-projects source holds one once its configuration sets priority_mapping."
583 )]
584 NoPriority {
585 name: String,
587 kind: String,
589 task: String,
591 priority: onetaskgraph_plugin_api::Priority,
593 },
594
595 #[error(
597 "no project with the id {id}\n\
598 next: check the id, or list what is there — `onetaskgraph project list` reports every \
599 project the configured sources hold."
600 )]
601 NoSuchProject {
602 id: String,
604 },
605
606 #[error(
608 "no document with the id {id}\n\
609 next: check the id, or list what is there — `onetaskgraph document list` reports every \
610 document the configured sources hold."
611 )]
612 NoSuchDocument {
613 id: String,
615 },
616
617 #[error(
619 "no task with the id {id}\n\
620 next: check the id, or list what is there — `onetaskgraph task list` reports every \
621 task the configured sources hold."
622 )]
623 NoSuchTask {
624 id: String,
626 },
627
628 #[error(
630 "task {task} has no comment with the id {comment}\n\
631 next: list its comments — `onetaskgraph task comment list {task}` reports each \
632 one's id."
633 )]
634 NoSuchComment {
635 task: String,
637 comment: String,
639 },
640
641 #[error(
646 "source {name} could not be built: {error}\n\
647 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
648 command again."
649 )]
650 SourceUnavailable {
651 name: String,
653 error: SourceError,
655 },
656
657 #[error(
663 "source {name} could not do it: {error}\n\
664 next: fix what the source named above, then run the command again."
665 )]
666 SourceFailed {
667 name: String,
669 error: SourceError,
671 },
672
673 #[error(
675 "the destination source {name} could not be built: {error}\n\
676 next: fix that source — `onetaskgraph sources list` reports its state — then \
677 copy again."
678 )]
679 DestinationUnavailable {
680 name: String,
682 error: SourceError,
684 },
685
686 #[error(
688 "no item with the id {id}\n\
689 next: check the id, or list what is there — `onetaskgraph task list` and \
690 `onetaskgraph project list` report what the configured sources hold."
691 )]
692 NoSuchItem {
693 id: String,
695 },
696
697 #[error(
702 "{item} was copied from {origin}, which that destination no longer holds\n\
703 next: re-run with --recreate to create a new item there instead, or restore \
704 {origin}."
705 )]
706 StaleOrigin {
707 item: String,
709 origin: String,
711 },
712
713 #[error(
719 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
720 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
721 {id} on its own with `onetaskgraph task copy`.",
722 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
723 )]
724 NotAMember {
725 id: GlobalId,
727 projects: Vec<GlobalId>,
729 },
730
731 #[error(
740 "{item} depends on {member}, which this copy was not told to carry and which records \
741 no origin in {destination}\n\
742 next: name {member} with --member as well, record its {destination} id at \
743 onetaskgraph.origin, or copy the whole project without --member."
744 )]
745 UnrecordedMember {
746 item: GlobalId,
748 member: GlobalId,
750 destination: SourceName,
752 },
753
754 #[error(
760 "source {name} could not do it: {error}\n\
761 next: fix what the source named above, then copy again."
762 )]
763 SourceRefused {
764 name: String,
766 error: SourceError,
768 },
769
770 #[error(
778 "the copy failed and could not be undone.\n\
779 it failed because: {error}\n\
780 it could not be undone because: {refusal}\n\
781 so the destination still holds: {left_behind}\n\
782 next: remove those items at the destination, then copy again."
783 )]
784 CopyNotUndone {
785 error: Box<EngineError>,
787 left_behind: LeftBehind,
793 refusal: SourceError,
795 },
796}
797
798#[derive(Debug, Clone, PartialEq)]
807pub struct LeftBehind {
808 first: GlobalId,
810 rest: Vec<GlobalId>,
812}
813
814impl LeftBehind {
815 #[must_use]
817 pub fn new(first: GlobalId) -> Self {
818 Self {
819 first,
820 rest: Vec::new(),
821 }
822 }
823
824 pub fn push(&mut self, id: GlobalId) {
826 self.rest.push(id);
827 }
828
829 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
831 std::iter::once(&self.first).chain(self.rest.iter())
832 }
833}
834
835impl std::fmt::Display for LeftBehind {
836 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
838 write!(formatter, "{}", self.first)?;
839 for id in &self.rest {
840 write!(formatter, ", {id}")?;
841 }
842 Ok(())
843 }
844}
845
846pub enum ConfiguredSource {
853 Ready(ResolvedSource),
855 Unavailable(UnavailableSource),
857}
858
859impl ConfiguredSource {
860 #[must_use]
862 pub fn name(&self) -> &SourceName {
863 match self {
864 Self::Ready(source) => source.name(),
865 Self::Unavailable(source) => source.name(),
866 }
867 }
868}
869
870pub struct Engine {
872 sources: Vec<ConfiguredSource>,
874 selection: Vec<SourceName>,
876}
877
878impl Engine {
879 #[must_use]
887 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
888 let (ready, unavailable) = resolve_available(config, secrets);
889 Self::new(
890 ready
891 .into_iter()
892 .map(ConfiguredSource::Ready)
893 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
894 .collect(),
895 config.selected_sources(),
896 )
897 }
898
899 #[must_use]
902 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
903 Self { sources, selection }
904 }
905
906 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
908 self.sources.iter().filter_map(|source| match source {
909 ConfiguredSource::Ready(ready) => Some(ready),
910 ConfiguredSource::Unavailable(_) => None,
911 })
912 }
913
914 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
916 self.sources.iter().filter_map(|source| match source {
917 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
918 ConfiguredSource::Ready(_) => None,
919 })
920 }
921
922 #[must_use]
924 pub fn listing(&self) -> Vec<SourceListing> {
925 let mut listings: Vec<SourceListing> = self
926 .ready()
927 .map(|source| SourceListing {
928 source: source.name().clone(),
929 kind: source.kind().to_owned(),
930 state: SourceState::Available {
931 capabilities: source.source().capabilities(),
932 },
933 })
934 .chain(self.unavailable().map(|source| SourceListing {
935 source: source.name().clone(),
936 kind: source.kind().to_owned(),
937 state: SourceState::Unavailable {
938 error: source.error().clone(),
939 },
940 }))
941 .collect();
942 listings.sort_by(|left, right| left.source.cmp(&right.source));
943 listings
944 }
945
946 #[must_use]
952 pub fn has(&self, name: &SourceName) -> bool {
953 self.sources.iter().any(|source| source.name() == name)
954 }
955
956 pub async fn tasks(
964 &self,
965 request: &TaskRequest,
966 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
967 let mut names = self.resolve_selection(&request.sources)?;
968 if let ProjectSelector::Qualified(id) = &request.project {
973 self.known(&id.source)?;
974 names.retain(|name| name == &id.source);
975 }
976 let query = shape(
977 "task-list",
978 &names,
979 &(&request.filters, &request.project, &request.priorities),
980 );
981 let states = resumption(
982 self,
983 request.paging.token.as_ref(),
984 &[StreamKind::Items],
985 &query,
986 )?;
987 let budget = request.paging.limit.get();
988
989 let mut answer = Answer::new();
990 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
991
992 let shapes: Vec<TaskShape> = ready
993 .iter()
994 .map(|source| {
995 shape_tasks(
996 &source.source().capabilities(),
997 &request.filters,
998 &project_filter(&request.project),
999 &request.priorities,
1000 )
1001 })
1002 .collect();
1003 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1004 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1005
1006 let walks = ready
1007 .iter()
1008 .enumerate()
1009 .map(|(index, source)| {
1010 fetch_tasks(
1011 source,
1012 &shapes[index],
1013 &starts[index],
1014 budget,
1015 &counters[index],
1016 )
1017 })
1018 .collect();
1019
1020 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1021 answer.finish(
1022 streams,
1023 budget,
1024 owed(&states),
1025 &query,
1026 |name, task: Task| {
1027 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1028 },
1029 )
1030 }
1031
1032 pub async fn projects(
1038 &self,
1039 request: &ProjectRequest,
1040 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1041 let names = self.resolve_selection(&request.sources)?;
1042 let query = shape("project-list", &names, &request.filters);
1043 let states = resumption(
1044 self,
1045 request.paging.token.as_ref(),
1046 &[StreamKind::Items],
1047 &query,
1048 )?;
1049 let budget = request.paging.limit.get();
1050
1051 let mut answer = Answer::new();
1052 let mut with_projects = Vec::new();
1058 for source in answer.split(self, &names) {
1059 if source.source().capabilities().projects.is_native() {
1060 with_projects.push(source);
1061 } else {
1062 answer.unreachable_predicate(source, Predicate::Project);
1063 }
1064 }
1065 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1066
1067 let shapes: Vec<ProjectShape> = ready
1068 .iter()
1069 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1070 .collect();
1071 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1072 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1073
1074 let walks = ready
1075 .iter()
1076 .enumerate()
1077 .map(|(index, source)| {
1078 fetch_projects(
1079 source,
1080 &shapes[index],
1081 &starts[index],
1082 budget,
1083 &counters[index],
1084 )
1085 })
1086 .collect();
1087
1088 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1089 answer.finish(
1090 streams,
1091 budget,
1092 owed(&states),
1093 &query,
1094 |name, project: Project| Qualified {
1095 id: GlobalId::new(name.clone(), project.id.clone()),
1096 item: project,
1097 },
1098 )
1099 }
1100
1101 pub async fn documents(
1113 &self,
1114 request: &DocumentRequest,
1115 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1116 let mut names = self.resolve_selection(&request.sources)?;
1117 if let ProjectSelector::Qualified(id) = &request.project {
1120 self.known(&id.source)?;
1121 names.retain(|name| name == &id.source);
1122 }
1123 let query = shape(
1124 "document-list",
1125 &names,
1126 &(&request.filters, &request.project),
1127 );
1128 let states = resumption(
1129 self,
1130 request.paging.token.as_ref(),
1131 &[StreamKind::Items],
1132 &query,
1133 )?;
1134 let budget = request.paging.limit.get();
1135
1136 let mut answer = Answer::new();
1137 let mut with_documents = Vec::new();
1138 for source in answer.split(self, &names) {
1139 if source.source().capabilities().documents.is_native() {
1140 with_documents.push(source);
1141 } else {
1142 answer.unreachable_predicate(source, Predicate::Document);
1143 }
1144 }
1145 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1146
1147 let shapes: Vec<DocumentShape> = ready
1148 .iter()
1149 .map(|source| {
1150 shape_documents(
1151 &source.source().capabilities(),
1152 &request.filters,
1153 &project_filter(&request.project),
1154 )
1155 })
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_documents(
1165 source,
1166 &shapes[index],
1167 &starts[index],
1168 budget,
1169 &counters[index],
1170 )
1171 })
1172 .collect();
1173
1174 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1175 answer.finish(
1176 streams,
1177 budget,
1178 owed(&states),
1179 &query,
1180 |name, document: Document| Qualified {
1181 id: GlobalId::new(name.clone(), document.id.clone()),
1182 item: document,
1183 },
1184 )
1185 }
1186
1187 pub async fn labels(
1193 &self,
1194 request: &LabelRequest,
1195 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1196 let names = self.resolve_selection(&request.sources)?;
1197 let query = shape("label-list", &names, &());
1198 let states = resumption(
1199 self,
1200 request.paging.token.as_ref(),
1201 &[StreamKind::Items],
1202 &query,
1203 )?;
1204 let budget = request.paging.limit.get();
1205
1206 let mut answer = Answer::new();
1207 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1208 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1209 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1210
1211 let walks = ready
1212 .iter()
1213 .enumerate()
1214 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1215 .collect();
1216
1217 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1218 answer.finish(
1219 streams,
1220 budget,
1221 owed(&states),
1222 &query,
1223 |name, label: Label| Qualified {
1224 id: GlobalId::new(name.clone(), label.id.clone()),
1225 item: label,
1226 },
1227 )
1228 }
1229
1230 pub async fn search(
1236 &self,
1237 request: &SearchRequest,
1238 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1239 let names = self.resolve_selection(&request.sources)?;
1240 let reads: &[StreamKind] = match request.kind {
1244 SearchKind::Tasks => &[StreamKind::Tasks],
1245 SearchKind::Projects => &[StreamKind::Projects],
1246 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1247 };
1248 let query = shape("search", &names, &request.text);
1253 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1254 let budget = request.paging.limit.get();
1255 let filters = Filters {
1256 text: Some(request.text.clone()),
1257 ..Filters::default()
1258 };
1259
1260 let mut answer = Answer::new();
1261
1262 let mut ready = Vec::new();
1265 let mut kinds = Vec::new();
1266 let mut starts = Vec::new();
1267 for source in answer.split(self, &names) {
1268 let mut streams = Vec::new();
1269 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1270 streams.push(StreamKind::Tasks);
1271 }
1272 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1273 if source.source().capabilities().projects.is_native() {
1274 streams.push(StreamKind::Projects);
1275 } else {
1276 answer.unreachable_predicate(source, Predicate::Project);
1277 }
1278 }
1279 for stream in streams {
1280 if let Some(resume) = resume_at(&states, source.name(), stream) {
1281 ready.push(source);
1282 kinds.push(stream);
1283 starts.push(resume);
1284 }
1285 }
1286 }
1287
1288 let shapes: Vec<HitShape> = ready
1289 .iter()
1290 .zip(kinds.iter())
1291 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1292 .collect();
1293 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1294 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1295
1296 let walks = ready
1297 .iter()
1298 .enumerate()
1299 .map(|(index, source)| {
1300 fetch_hits(
1301 source,
1302 &shapes[index],
1303 &starts[index],
1304 budget,
1305 &counters[index],
1306 )
1307 })
1308 .collect();
1309
1310 let streams =
1311 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1312 answer.finish(
1313 streams,
1314 budget,
1315 owed(&states),
1316 &query,
1317 |name, found: Found| match found {
1318 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1319 GlobalId::new(name.clone(), task.id.clone()),
1320 task,
1321 )),
1322 Found::Project(project) => SearchHit::Project(Qualified {
1323 id: GlobalId::new(name.clone(), project.id.clone()),
1324 item: project,
1325 }),
1326 },
1327 )
1328 }
1329
1330 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1337 let name = self.known(&id.source)?;
1338 let mut answer = Answer::new();
1339 let selected = answer.split(self, std::slice::from_ref(&name));
1340 let Some(source) = selected.first() else {
1341 return answer.nothing();
1342 };
1343 let found = source.source().get_task(&id.native).await;
1344 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1345 answer.one(source, found, |task| {
1346 delivery::qualified_task(qualified, task)
1347 })
1348 }
1349
1350 pub async fn project(
1356 &self,
1357 id: &GlobalId,
1358 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1359 let name = self.known(&id.source)?;
1360 let mut answer = Answer::new();
1361 let selected = answer.split(self, std::slice::from_ref(&name));
1362 let Some(source) = selected.first() else {
1363 return answer.nothing();
1364 };
1365 let found = source.source().get_project(&id.native).await;
1366 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1367 answer.one(source, found, |project| Qualified {
1368 id: qualified,
1369 item: project,
1370 })
1371 }
1372
1373 pub async fn document(
1383 &self,
1384 id: &GlobalId,
1385 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1386 let name = self.known(&id.source)?;
1387 let mut answer = Answer::new();
1388 let selected = answer.split(self, std::slice::from_ref(&name));
1389 let Some(source) = selected.first() else {
1390 return answer.nothing();
1391 };
1392 if !source.source().capabilities().documents.is_native() {
1393 answer.unreachable_predicate(source, Predicate::Document);
1394 return answer.nothing();
1395 }
1396 let found = source.source().get_document(&id.native).await;
1397 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1398 answer.one(source, found, |document| Qualified {
1399 id: qualified,
1400 item: document,
1401 })
1402 }
1403
1404 pub async fn task_dependencies(
1411 &self,
1412 request: &DependencyRequest,
1413 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1414 self.dependencies(request, Entity::Task).await
1415 }
1416
1417 pub async fn project_dependencies(
1423 &self,
1424 request: &DependencyRequest,
1425 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1426 self.dependencies(request, Entity::Project).await
1427 }
1428
1429 async fn dependencies(
1432 &self,
1433 request: &DependencyRequest,
1434 entity: Entity,
1435 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1436 let name = self.known(&request.id.source)?;
1437 let query = shape(
1438 "dependencies",
1439 std::slice::from_ref(&name),
1440 &(entity, &request.id.native, request.direction),
1441 );
1442 let states = resumption(
1443 self,
1444 request.paging.token.as_ref(),
1445 &[StreamKind::Items],
1446 &query,
1447 )?;
1448 let budget = request.paging.limit.get();
1449
1450 let mut answer = Answer::new();
1451 let (ready, starts) = walking(
1452 answer.split(self, std::slice::from_ref(&name)),
1453 &states,
1454 StreamKind::Items,
1455 );
1456 let Some(source) = ready.first() else {
1457 return answer.nothing();
1458 };
1459
1460 let capabilities = source.source().capabilities();
1461 let support = match entity {
1462 Entity::Task => capabilities.task_dependencies,
1463 Entity::Project => capabilities.project_dependencies,
1464 };
1465 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1469 let mut outcomes = Outcomes::default();
1470 if request.direction == Direction::DependedOnBy {
1471 if emulating {
1472 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1473 } else {
1474 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1475 }
1476 }
1477
1478 let counters = vec![AtomicU32::new(0)];
1479 let walked = fetch_edges(
1480 source,
1481 &request.id.native,
1482 request.direction,
1483 entity,
1484 emulating,
1485 &starts[0],
1486 budget,
1487 &counters[0],
1488 )
1489 .await;
1490
1491 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1492 answer.finish(
1493 streams,
1494 budget,
1495 owed(&states),
1496 &query,
1497 |name, edge: DependencyEdge| QualifiedEdge {
1498 from: qualify_endpoint(name, edge.from),
1499 to: qualify_endpoint(name, edge.to),
1500 kind: edge.kind,
1501 },
1502 )
1503 }
1504
1505 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1507 if asked.is_empty() {
1508 if self.selection.is_empty() {
1509 return Err(EngineError::NoSources);
1510 }
1511 return Ok(self.selection.clone());
1512 }
1513 asked.iter().map(|name| self.known(name)).collect()
1514 }
1515
1516 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1518 if self.has(name) {
1519 return Ok(name.clone());
1520 }
1521 if self.sources.is_empty() {
1522 return Err(EngineError::NoSources);
1523 }
1524 Err(EngineError::UnknownSource {
1525 name: name.to_string(),
1526 configured: self
1527 .listing()
1528 .iter()
1529 .map(|listing| listing.source.to_string())
1530 .collect::<Vec<_>>()
1531 .join(", "),
1532 })
1533 }
1534}
1535
1536fn qualify_endpoint(
1537 source: &SourceName,
1538 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1539) -> QualifiedEndpoint {
1540 let kind = endpoint.kind;
1541 let is_qualified = endpoint.is_qualified();
1542 let endpoint_id = endpoint.into_id();
1543 QualifiedEndpoint {
1544 id: if is_qualified {
1545 endpoint_id
1546 .parse()
1547 .expect("plugin-api validates qualified dependency endpoints")
1548 } else {
1549 GlobalId::new(source.clone(), NativeId(endpoint_id))
1550 },
1551 kind,
1552 }
1553}
1554
1555#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1557enum Entity {
1558 Task,
1560 Project,
1562}
1563
1564enum Found {
1566 Task(Task),
1568 Project(Project),
1570}
1571
1572#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1574enum Outcome {
1575 PushedDown,
1577 AppliedLocally,
1579 Emulated,
1581 Unavailable,
1583}
1584
1585#[derive(Debug, Clone, Default, PartialEq)]
1597struct Outcomes(BTreeMap<Predicate, Outcome>);
1598
1599impl Outcomes {
1600 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1606 self.0.insert(predicate, outcome);
1607 }
1608
1609 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1612 for predicate in predicates {
1613 self.record(predicate, outcome);
1614 }
1615 }
1616
1617 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1619 self.0
1620 .iter()
1621 .filter(|(_, recorded)| **recorded == outcome)
1622 .map(|(predicate, _)| *predicate)
1623 .collect()
1624 }
1625}
1626
1627struct TaskShape {
1629 pushed: TaskQuery,
1631 local: LocalTasks,
1633 outcomes: Outcomes,
1635}
1636
1637struct ProjectShape {
1639 pushed: ProjectQuery,
1641 local: LocalProjects,
1643 outcomes: Outcomes,
1645}
1646
1647struct DocumentShape {
1649 pushed: DocumentQuery,
1651 local: LocalDocuments,
1653 outcomes: Outcomes,
1655}
1656
1657struct HitShape {
1659 stream: StreamKind,
1661 tasks: TaskQuery,
1663 projects: ProjectQuery,
1665 local_tasks: LocalTasks,
1667 local_projects: LocalProjects,
1669 outcomes: Outcomes,
1671}
1672
1673struct Answer {
1679 plans: Vec<SourcePlan>,
1681 errors: Vec<SourceFailure>,
1683}
1684
1685impl Answer {
1686 fn new() -> Self {
1687 Self {
1688 plans: Vec::new(),
1689 errors: Vec::new(),
1690 }
1691 }
1692
1693 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1698 let mut selected = Vec::new();
1699 for name in names {
1700 match engine.sources.iter().find(|source| source.name() == name) {
1701 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1702 Some(ConfiguredSource::Unavailable(source)) => {
1703 self.errors.push(source.failure());
1704 }
1705 None => {}
1706 }
1707 }
1708 selected
1709 }
1710
1711 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1713 let mut outcomes = Outcomes::default();
1714 outcomes.record(predicate, Outcome::Unavailable);
1715 self.plans.push(plan_for(source, outcomes, 0));
1716 }
1717
1718 fn collect<T>(
1720 &mut self,
1721 ready: &[&ResolvedSource],
1722 walked: Vec<Result<Fetched<T>, SourceError>>,
1723 counters: &[AtomicU32],
1724 outcomes: Vec<Outcomes>,
1725 ) -> Vec<Stream<T>> {
1726 let kinds = vec![StreamKind::Items; ready.len()];
1727 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1728 }
1729
1730 fn collect_streams<T>(
1732 &mut self,
1733 ready: &[&ResolvedSource],
1734 kinds: &[StreamKind],
1735 walked: Vec<Result<Fetched<T>, SourceError>>,
1736 counters: &[AtomicU32],
1737 outcomes: Vec<Outcomes>,
1738 ) -> Vec<Stream<T>> {
1739 let mut streams = Vec::new();
1740 for (index, result) in walked.into_iter().enumerate() {
1741 let source = ready[index];
1742 let pages = counters[index].load(Ordering::Relaxed);
1743 self.plans
1744 .push(plan_for(source, outcomes[index].clone(), pages));
1745 match result {
1746 Ok(fetched) => streams.push(Stream {
1747 source: source.name().clone(),
1748 kind: kinds[index],
1749 fetched,
1750 }),
1751 Err(error) => self.errors.push(SourceFailure {
1754 source: source.name().clone(),
1755 error,
1756 }),
1757 }
1758 }
1759 streams
1760 }
1761
1762 fn one<T, U>(
1764 mut self,
1765 source: &ResolvedSource,
1766 found: Result<Option<T>, SourceError>,
1767 qualify: impl FnOnce(T) -> U,
1768 ) -> Result<QueryResponse<U>, EngineError> {
1769 self.plans.push(plan_for(source, Outcomes::default(), 1));
1770 let items = match found {
1771 Ok(Some(item)) => vec![qualify(item)],
1772 Ok(None) => Vec::new(),
1773 Err(error) => {
1774 self.errors.push(SourceFailure {
1775 source: source.name().clone(),
1776 error,
1777 });
1778 Vec::new()
1779 }
1780 };
1781 Ok(QueryResponse {
1782 items,
1783 next: None,
1784 plan: QueryPlan {
1785 per_source: merge_plans(self.plans),
1786 },
1787 errors: self.errors,
1788 })
1789 }
1790
1791 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1793 Ok(QueryResponse {
1794 items: Vec::new(),
1795 next: None,
1796 plan: QueryPlan {
1797 per_source: merge_plans(self.plans),
1798 },
1799 errors: self.errors,
1800 })
1801 }
1802
1803 fn finish<T, U>(
1808 self,
1809 streams: Vec<Stream<T>>,
1810 budget: u32,
1811 first: Option<&Owed>,
1812 query: &str,
1813 qualify: impl Fn(&SourceName, T) -> U,
1814 ) -> Result<QueryResponse<U>, EngineError> {
1815 let (rows, states, owed) = merge(streams, budget, first);
1816 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1817 Ok(QueryResponse {
1818 items: rows
1819 .into_iter()
1820 .map(|(name, item)| qualify(&name, item))
1821 .collect(),
1822 next,
1823 plan: QueryPlan {
1824 per_source: merge_plans(self.plans),
1825 },
1826 errors: self.errors,
1827 })
1828 }
1829}
1830
1831fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1834 SourcePlan {
1835 source: source.name().clone(),
1836 kind: source.kind().to_owned(),
1837 pushed_down: outcomes.with(Outcome::PushedDown),
1838 applied_locally: outcomes.with(Outcome::AppliedLocally),
1839 emulated: outcomes.with(Outcome::Emulated),
1840 unavailable: outcomes.with(Outcome::Unavailable),
1841 pages_fetched: pages,
1842 }
1843}
1844
1845fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1850 let mut merged: Vec<SourcePlan> = Vec::new();
1851 for plan in plans {
1852 if let Some(existing) = merged
1853 .iter_mut()
1854 .find(|existing| existing.source == plan.source)
1855 {
1856 existing.pushed_down.extend(plan.pushed_down);
1857 existing.applied_locally.extend(plan.applied_locally);
1858 existing.emulated.extend(plan.emulated);
1859 existing.unavailable.extend(plan.unavailable);
1860 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1861 for list in [
1862 &mut existing.pushed_down,
1863 &mut existing.applied_locally,
1864 &mut existing.emulated,
1865 &mut existing.unavailable,
1866 ] {
1867 list.sort_unstable();
1868 list.dedup();
1869 }
1870 } else {
1871 merged.push(plan);
1872 }
1873 }
1874 merged
1875}
1876
1877fn walking<'a>(
1883 selected: Vec<&'a ResolvedSource>,
1884 states: &Option<Resumption>,
1885 kind: StreamKind,
1886) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1887 let mut ready = Vec::new();
1888 let mut starts = Vec::new();
1889 for source in selected {
1890 if let Some(resume) = resume_at(states, source.name(), kind) {
1891 ready.push(source);
1892 starts.push(resume);
1893 }
1894 }
1895 (ready, starts)
1896}
1897
1898fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1916 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1917 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1918}
1919
1920fn fingerprint(text: &str) -> String {
1922 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1923 for byte in text.as_bytes() {
1924 hash ^= u64::from(*byte);
1925 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1926 }
1927 format!("{hash:016x}")
1928}
1929
1930fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1935 document.as_ref()?.owed.as_ref()
1936}
1937
1938fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1940 match states {
1941 None => Some(Resume::default()),
1942 Some(document) => document
1943 .streams
1944 .iter()
1945 .find(|state| &state.source == source && state.stream == kind)
1946 .map(|state| state.resume.clone()),
1947 }
1948}
1949
1950fn resumption(
1981 engine: &Engine,
1982 token: Option<&PageToken>,
1983 reads: &[StreamKind],
1984 query: &str,
1985) -> Result<Option<Resumption>, EngineError> {
1986 let Some(document) = token.map(PageToken::decode) else {
1987 return Ok(None);
1988 };
1989
1990 if document.query != query {
1997 return Err(EngineError::Token {
1998 message: "this page token was written by a different query — resume the walk it \
1999 came from, or drop --page to start this one from the beginning"
2000 .to_owned(),
2001 });
2002 }
2003 let states = &document.streams;
2004
2005 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2006 for state in states {
2007 if !reads.contains(&state.stream) {
2008 return Err(EngineError::Token {
2009 message: format!(
2010 "this page token resumes {}, which this command does not read — it \
2011 was written by a different query",
2012 state.stream.describe()
2013 ),
2014 });
2015 }
2016 let ceiling = engine
2017 .ready()
2018 .find(|source| source.name() == &state.source)
2019 .map(ceiling);
2020 if ceiling.is_none() && !engine.has(&state.source) {
2021 return Err(EngineError::Token {
2022 message: format!(
2023 "this page token resumes a source called {:?}, which this \
2024 configuration does not have",
2025 state.source.as_str()
2026 ),
2027 });
2028 }
2029 if let Some(ceiling) = ceiling
2030 && state.resume.skip >= ceiling
2031 {
2032 return Err(EngineError::Token {
2033 message: format!(
2034 "this page token resumes {} rows into a page of source {:?}, which \
2035 serves at most {ceiling}",
2036 state.resume.skip,
2037 state.source.as_str()
2038 ),
2039 });
2040 }
2041 if seen.contains(&(&state.source, state.stream)) {
2042 return Err(EngineError::Token {
2043 message: format!(
2044 "this page token gives source {:?} two places to resume from",
2045 state.source.as_str()
2046 ),
2047 });
2048 }
2049 seen.push((&state.source, state.stream));
2050 }
2051
2052 if let Some(owed) = &document.owed
2057 && !document
2058 .streams
2059 .iter()
2060 .any(|state| state.source == owed.source && state.stream == owed.stream)
2061 {
2062 return Err(EngineError::Token {
2063 message: format!(
2064 "this page token owes the next row to a stream it does not resume, \
2065 {:?}'s {}",
2066 owed.source.as_str(),
2067 owed.stream.describe()
2068 ),
2069 });
2070 }
2071
2072 Ok(Some(document))
2073}
2074
2075fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2081 match selector {
2082 ProjectSelector::Any => ProjectFilter::Any,
2083 ProjectSelector::Orphans => ProjectFilter::Orphans,
2084 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2085 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2086 }
2087}
2088
2089fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2091 match fields {
2092 TextFields::Title => vec![Predicate::SearchTitle],
2093 TextFields::Content => vec![Predicate::SearchContent],
2094 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2095 }
2096}
2097
2098fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2105 match fields {
2106 TextFields::Title => capabilities.search_title.is_native(),
2107 TextFields::Content => capabilities.search_content.is_native(),
2108 TextFields::TitleOrContent => {
2109 capabilities.search_title.is_native() && capabilities.search_content.is_native()
2110 }
2111 }
2112}
2113
2114fn shape_tasks(
2116 capabilities: &Capabilities,
2117 filters: &Filters,
2118 project: &ProjectFilter,
2119 priorities: &[Priority],
2120) -> TaskShape {
2121 let mut pushed = TaskQuery::default();
2122 let mut local = LocalTasks::default();
2123 let mut outcomes = Outcomes::default();
2124
2125 if !filters.labels.is_empty() {
2126 if capabilities.filter_by_label.is_native() {
2127 pushed.labels = filters.labels.clone();
2128 outcomes.record(Predicate::Label, Outcome::PushedDown);
2129 } else {
2130 local.labels = Some(filters.labels.clone());
2131 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2132 }
2133 }
2134 if !filters.statuses.is_empty() {
2135 if capabilities.filter_by_status.is_native() {
2136 pushed.statuses.clone_from(&filters.statuses);
2137 outcomes.record(Predicate::Status, Outcome::PushedDown);
2138 } else {
2139 local.statuses.clone_from(&filters.statuses);
2140 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2141 }
2142 }
2143 if !priorities.is_empty() {
2144 if capabilities.filter_by_priority.is_native() {
2145 pushed.priorities = priorities.to_vec();
2146 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2147 } else {
2148 local.priorities = priorities.to_vec();
2149 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2150 }
2151 }
2152 if let Some(text) = &filters.text {
2153 let predicates = text_predicates(text.fields);
2154 if searches_natively(capabilities, text.fields) {
2155 pushed.text = Some(text.clone());
2156 outcomes.record_all(predicates, Outcome::PushedDown);
2157 } else {
2158 local.text = Some(text.clone());
2159 outcomes.record_all(predicates, Outcome::AppliedLocally);
2160 }
2161 }
2162 match project {
2163 ProjectFilter::Any => {}
2164 ProjectFilter::Orphans => {
2165 if capabilities.orphan_tasks.is_native() {
2166 pushed.project = ProjectFilter::Orphans;
2167 outcomes.record(Predicate::Project, Outcome::PushedDown);
2168 } else {
2169 local.project = Some(ProjectFilter::Orphans);
2170 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2171 }
2172 }
2173 ProjectFilter::Is(id) => {
2174 if capabilities.projects.is_native() {
2175 pushed.project = ProjectFilter::Is(id.clone());
2176 outcomes.record(Predicate::Project, Outcome::PushedDown);
2177 } else {
2178 local.project = Some(ProjectFilter::Is(id.clone()));
2179 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2180 }
2181 }
2182 }
2183
2184 TaskShape {
2185 pushed,
2186 local,
2187 outcomes,
2188 }
2189}
2190
2191fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2193 let mut pushed = ProjectQuery::default();
2194 let mut local = LocalProjects::default();
2195 let mut outcomes = Outcomes::default();
2196
2197 if !filters.labels.is_empty() {
2198 if capabilities.filter_by_label.is_native() {
2199 pushed.labels = filters.labels.clone();
2200 outcomes.record(Predicate::Label, Outcome::PushedDown);
2201 } else {
2202 local.labels = Some(filters.labels.clone());
2203 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2204 }
2205 }
2206 if !filters.statuses.is_empty() {
2207 if capabilities.filter_by_status.is_native() {
2208 pushed.statuses.clone_from(&filters.statuses);
2209 outcomes.record(Predicate::Status, Outcome::PushedDown);
2210 } else {
2211 local.statuses.clone_from(&filters.statuses);
2212 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2213 }
2214 }
2215 if let Some(text) = &filters.text {
2216 let predicates = text_predicates(text.fields);
2217 if searches_natively(capabilities, text.fields) {
2218 pushed.text = Some(text.clone());
2219 outcomes.record_all(predicates, Outcome::PushedDown);
2220 } else {
2221 local.text = Some(text.clone());
2222 outcomes.record_all(predicates, Outcome::AppliedLocally);
2223 }
2224 }
2225
2226 ProjectShape {
2227 pushed,
2228 local,
2229 outcomes,
2230 }
2231}
2232
2233fn shape_documents(
2238 capabilities: &Capabilities,
2239 filters: &DocumentFilters,
2240 project: &ProjectFilter,
2241) -> DocumentShape {
2242 let mut pushed = DocumentQuery::default();
2243 let mut local = LocalDocuments::default();
2244 let mut outcomes = Outcomes::default();
2245
2246 if !filters.labels.is_empty() {
2247 if capabilities.filter_by_label.is_native() {
2248 pushed.labels = filters.labels.clone();
2249 outcomes.record(Predicate::Label, Outcome::PushedDown);
2250 } else {
2251 local.labels = Some(filters.labels.clone());
2252 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2253 }
2254 }
2255 if let Some(text) = &filters.text {
2256 let predicates = text_predicates(text.fields);
2257 if searches_natively(capabilities, text.fields) {
2258 pushed.text = Some(text.clone());
2259 outcomes.record_all(predicates, Outcome::PushedDown);
2260 } else {
2261 local.text = Some(text.clone());
2262 outcomes.record_all(predicates, Outcome::AppliedLocally);
2263 }
2264 }
2265 match project {
2266 ProjectFilter::Any => {}
2267 ProjectFilter::Orphans => {
2268 if capabilities.orphan_tasks.is_native() {
2269 pushed.project = ProjectFilter::Orphans;
2270 outcomes.record(Predicate::Project, Outcome::PushedDown);
2271 } else {
2272 local.project = Some(ProjectFilter::Orphans);
2273 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2274 }
2275 }
2276 ProjectFilter::Is(id) => {
2277 if capabilities.projects.is_native() {
2278 pushed.project = ProjectFilter::Is(id.clone());
2279 outcomes.record(Predicate::Project, Outcome::PushedDown);
2280 } else {
2281 local.project = Some(ProjectFilter::Is(id.clone()));
2282 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2283 }
2284 }
2285 }
2286
2287 DocumentShape {
2288 pushed,
2289 local,
2290 outcomes,
2291 }
2292}
2293
2294fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2296 match stream {
2297 StreamKind::Projects => {
2298 let shaped = shape_projects(capabilities, filters);
2299 HitShape {
2300 stream,
2301 tasks: TaskQuery::default(),
2302 projects: shaped.pushed,
2303 local_tasks: LocalTasks::default(),
2304 local_projects: shaped.local,
2305 outcomes: shaped.outcomes,
2306 }
2307 }
2308 StreamKind::Items | StreamKind::Tasks => {
2309 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any, &[]);
2310 HitShape {
2311 stream,
2312 tasks: shaped.pushed,
2313 projects: ProjectQuery::default(),
2314 local_tasks: shaped.local,
2315 local_projects: LocalProjects::default(),
2316 outcomes: shaped.outcomes,
2317 }
2318 }
2319 }
2320}
2321
2322fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2329 if compensating {
2330 ceiling
2331 } else {
2332 budget.min(ceiling)
2333 }
2334}
2335
2336fn ceiling(source: &ResolvedSource) -> u32 {
2338 source.source().capabilities().max_page_size.max(1)
2339}
2340
2341async fn fetch_tasks(
2343 source: &ResolvedSource,
2344 shape: &TaskShape,
2345 start: &Resume,
2346 budget: u32,
2347 calls: &AtomicU32,
2348) -> Result<Fetched<Task>, SourceError> {
2349 let compensating = shape.local != LocalTasks::default();
2350 walk(
2351 start,
2352 budget,
2353 page_size(compensating, budget, ceiling(source)),
2354 |task| shape.local.keeps(task),
2355 |cursor, limit| async move {
2356 calls.fetch_add(1, Ordering::Relaxed);
2357 let request = PageRequest { cursor, limit };
2358 source.source().query_tasks(&shape.pushed, &request).await
2359 },
2360 )
2361 .await
2362}
2363
2364async fn fetch_projects(
2366 source: &ResolvedSource,
2367 shape: &ProjectShape,
2368 start: &Resume,
2369 budget: u32,
2370 calls: &AtomicU32,
2371) -> Result<Fetched<Project>, SourceError> {
2372 let compensating = shape.local != LocalProjects::default();
2373 walk(
2374 start,
2375 budget,
2376 page_size(compensating, budget, ceiling(source)),
2377 |project| shape.local.keeps(project),
2378 |cursor, limit| async move {
2379 calls.fetch_add(1, Ordering::Relaxed);
2380 let request = PageRequest { cursor, limit };
2381 source
2382 .source()
2383 .query_projects(&shape.pushed, &request)
2384 .await
2385 },
2386 )
2387 .await
2388}
2389
2390async fn fetch_documents(
2392 source: &ResolvedSource,
2393 shape: &DocumentShape,
2394 start: &Resume,
2395 budget: u32,
2396 calls: &AtomicU32,
2397) -> Result<Fetched<Document>, SourceError> {
2398 let compensating = shape.local != LocalDocuments::default();
2399 walk(
2400 start,
2401 budget,
2402 page_size(compensating, budget, ceiling(source)),
2403 |document| shape.local.keeps(document),
2404 |cursor, limit| async move {
2405 calls.fetch_add(1, Ordering::Relaxed);
2406 let request = PageRequest { cursor, limit };
2407 source
2408 .source()
2409 .query_documents(&shape.pushed, &request)
2410 .await
2411 },
2412 )
2413 .await
2414}
2415
2416async fn fetch_labels(
2418 source: &ResolvedSource,
2419 start: &Resume,
2420 budget: u32,
2421 calls: &AtomicU32,
2422) -> Result<Fetched<Label>, SourceError> {
2423 walk(
2424 start,
2425 budget,
2426 page_size(false, budget, ceiling(source)),
2427 |_| true,
2428 |cursor, limit| async move {
2429 calls.fetch_add(1, Ordering::Relaxed);
2430 let request = PageRequest { cursor, limit };
2431 source.source().labels(&request).await
2432 },
2433 )
2434 .await
2435}
2436
2437async fn fetch_hits(
2439 source: &ResolvedSource,
2440 shape: &HitShape,
2441 start: &Resume,
2442 budget: u32,
2443 calls: &AtomicU32,
2444) -> Result<Fetched<Found>, SourceError> {
2445 let ceiling = ceiling(source);
2446 match shape.stream {
2447 StreamKind::Projects => {
2448 let compensating = shape.local_projects != LocalProjects::default();
2449 walk(
2450 start,
2451 budget,
2452 page_size(compensating, budget, ceiling),
2453 |found| match found {
2454 Found::Project(project) => shape.local_projects.keeps(project),
2455 Found::Task(_) => true,
2456 },
2457 |cursor, limit| async move {
2458 calls.fetch_add(1, Ordering::Relaxed);
2459 let request = PageRequest { cursor, limit };
2460 let page = source
2461 .source()
2462 .query_projects(&shape.projects, &request)
2463 .await?;
2464 Ok(Page {
2465 items: page.items.into_iter().map(Found::Project).collect(),
2466 next: page.next,
2467 })
2468 },
2469 )
2470 .await
2471 }
2472 StreamKind::Items | StreamKind::Tasks => {
2473 let compensating = shape.local_tasks != LocalTasks::default();
2474 walk(
2475 start,
2476 budget,
2477 page_size(compensating, budget, ceiling),
2478 |found| match found {
2479 Found::Task(task) => shape.local_tasks.keeps(task),
2480 Found::Project(_) => true,
2481 },
2482 |cursor, limit| async move {
2483 calls.fetch_add(1, Ordering::Relaxed);
2484 let request = PageRequest { cursor, limit };
2485 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2486 Ok(Page {
2487 items: page.items.into_iter().map(Found::Task).collect(),
2488 next: page.next,
2489 })
2490 },
2491 )
2492 .await
2493 }
2494 }
2495}
2496
2497async fn forward_edges(
2499 source: &ResolvedSource,
2500 entity: Entity,
2501 id: &NativeId,
2502 request: &PageRequest,
2503) -> Result<Page<DependencyEdge>, SourceError> {
2504 match entity {
2505 Entity::Task => {
2506 source
2507 .source()
2508 .task_dependencies(id, Direction::DependsOn, request)
2509 .await
2510 }
2511 Entity::Project => {
2512 source
2513 .source()
2514 .project_dependencies(id, Direction::DependsOn, request)
2515 .await
2516 }
2517 }
2518}
2519
2520#[expect(
2529 clippy::too_many_arguments,
2530 reason = "every argument is one axis of one walk — the source, the item, the \
2531 direction, which of its two graphs, whether the reverse is emulated, where \
2532 to resume, how many rows to return and where to count calls. Grouping them \
2533 into a struct would name the same eight values one indirection further from \
2534 the loop that reads them."
2535)]
2536async fn fetch_edges(
2537 source: &ResolvedSource,
2538 native: &NativeId,
2539 direction: Direction,
2540 entity: Entity,
2541 emulating: bool,
2542 start: &Resume,
2543 budget: u32,
2544 calls: &AtomicU32,
2545) -> Result<Fetched<DependencyEdge>, SourceError> {
2546 let ceiling = ceiling(source);
2547 if !emulating {
2548 return walk(
2549 start,
2550 budget,
2551 page_size(false, budget, ceiling),
2552 |_| true,
2553 |cursor, limit| async move {
2554 calls.fetch_add(1, Ordering::Relaxed);
2555 let request = PageRequest { cursor, limit };
2556 match entity {
2557 Entity::Task => {
2558 source
2559 .source()
2560 .task_dependencies(native, direction, &request)
2561 .await
2562 }
2563 Entity::Project => {
2564 source
2565 .source()
2566 .project_dependencies(native, direction, &request)
2567 .await
2568 }
2569 }
2570 },
2571 )
2572 .await;
2573 }
2574
2575 walk(
2576 start,
2577 budget,
2578 ceiling,
2579 |_| true,
2580 |cursor, limit| async move {
2581 calls.fetch_add(1, Ordering::Relaxed);
2582 let request = PageRequest { cursor, limit };
2583 let (ids, next) = match entity {
2584 Entity::Task => {
2585 let page = source
2586 .source()
2587 .query_tasks(&TaskQuery::default(), &request)
2588 .await?;
2589 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2590 (ids, page.next)
2591 }
2592 Entity::Project => {
2593 let page = source
2594 .source()
2595 .query_projects(&ProjectQuery::default(), &request)
2596 .await?;
2597 let ids: Vec<NativeId> =
2598 page.items.into_iter().map(|project| project.id).collect();
2599 (ids, page.next)
2600 }
2601 };
2602
2603 let mut edges = Vec::new();
2604 for id in ids {
2605 let mut inner: Option<Cursor> = None;
2606 loop {
2607 calls.fetch_add(1, Ordering::Relaxed);
2608 let request = PageRequest {
2609 cursor: inner.clone(),
2610 limit,
2611 };
2612 let page = forward_edges(source, entity, &id, &request).await?;
2613 fits(page.items.len(), limit)?;
2617 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2618 unrepeated(
2619 page.next.as_ref(),
2620 inner.as_ref(),
2621 "its forward edges were being scanned",
2622 )?;
2623 match page.next {
2624 Some(cursor) => inner = Some(cursor),
2625 None => break,
2626 }
2627 }
2628 }
2629
2630 Ok(Page { items: edges, next })
2631 },
2632 )
2633 .await
2634}