1mod comment;
19mod copy;
20mod delivery;
21mod fetch;
22mod join;
23mod local;
24mod metadata;
25mod narrow;
26mod rendered;
27mod resume;
28
29use std::collections::BTreeMap;
30use std::num::NonZeroU32;
31use std::sync::atomic::{AtomicU32, Ordering};
32
33use onetaskgraph_plugin_api::{
34 Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
35 MetadataRecord, NativeId, Page, PageRequest, Priority, Project, ProjectFilter, ProjectQuery,
36 SecretResolver, SourceError, SourceName, StatusCategory, Task, TaskQuery, TextFields,
37 TextQuery,
38};
39use schemars::JsonSchema;
40use serde::{Deserialize, Serialize};
41
42use crate::GlobalId;
43use crate::config::Config;
44use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
45use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
46
47use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
48use join::join_all;
49use local::{LocalDocuments, LocalProjects, LocalTasks};
50pub(crate) use resume::{Owed, Resumption, StreamState};
51use resume::{Resume, StreamKind};
52
53pub use comment::{CommentList, DeletedComment, TaskDetail};
54pub use copy::{
55 BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
56 Spent,
57};
58pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
59pub use local::ProjectSelector;
60pub use metadata::MetadataSet;
61pub use narrow::{TaskContentSet, TaskPrioritySet};
62pub use rendered::{
63 Body, DocumentCreate, Regenerated, Regeneration, RenderRequest, RenderTemplate, RenderedRecord,
64 TaskCreate, TaskCreated, TemplateAnswers, UnusedAnswers,
65};
66
67#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
72pub struct Qualified<T> {
73 pub id: GlobalId,
75 pub item: T,
77}
78
79#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
84pub struct QualifiedEdge {
85 pub from: QualifiedEndpoint,
90 pub to: QualifiedEndpoint,
92 pub kind: onetaskgraph_plugin_api::DependencyKind,
94}
95
96#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
98pub struct QualifiedEndpoint {
99 pub id: GlobalId,
101 pub kind: onetaskgraph_plugin_api::ItemKind,
103}
104
105impl std::fmt::Display for QualifiedEndpoint {
106 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
107 self.id.fmt(formatter)
108 }
109}
110
111#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
113#[serde(tag = "kind", rename_all = "kebab-case")]
114pub enum SearchHit {
115 Task(Qualified<Task>),
117 Project(Qualified<Project>),
119}
120
121#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
123#[serde(rename_all = "kebab-case")]
124pub enum SearchKind {
125 Tasks,
127 Projects,
129 #[default]
131 Both,
132}
133
134#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
136pub struct SourceListing {
137 pub source: SourceName,
139 pub kind: String,
151 #[serde(flatten)]
153 pub state: SourceState,
154}
155
156#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
158#[serde(tag = "state", rename_all = "kebab-case")]
159pub enum SourceState {
160 Available {
162 capabilities: Capabilities,
164 },
165 Unavailable {
167 error: SourceError,
169 },
170}
171
172#[derive(Debug, Clone, PartialEq)]
174pub struct Paging {
175 pub limit: NonZeroU32,
177 pub token: Option<PageToken>,
179}
180
181#[derive(Debug, Clone, Default, PartialEq)]
183pub struct Filters {
184 pub text: Option<TextQuery>,
186 pub labels: LabelFilter,
188 pub statuses: Vec<StatusCategory>,
190}
191
192#[derive(Debug, Clone)]
194pub struct TaskRequest {
195 pub sources: Vec<SourceName>,
197 pub filters: Filters,
199 pub project: ProjectSelector,
201 pub priorities: Vec<Priority>,
207 pub paging: Paging,
209}
210
211#[derive(Debug, Clone)]
213pub struct ProjectRequest {
214 pub sources: Vec<SourceName>,
216 pub filters: Filters,
218 pub paging: Paging,
220}
221
222#[derive(Debug, Clone, Default, PartialEq)]
229pub struct DocumentFilters {
230 pub text: Option<TextQuery>,
232 pub labels: LabelFilter,
234}
235
236#[derive(Debug, Clone)]
238pub struct DocumentRequest {
239 pub sources: Vec<SourceName>,
241 pub filters: DocumentFilters,
243 pub project: ProjectSelector,
245 pub paging: Paging,
247}
248
249#[derive(Debug, Clone)]
251pub struct LabelRequest {
252 pub sources: Vec<SourceName>,
254 pub paging: Paging,
256}
257
258#[derive(Debug, Clone)]
260pub struct SearchRequest {
261 pub sources: Vec<SourceName>,
263 pub text: TextQuery,
265 pub kind: SearchKind,
267 pub paging: Paging,
269}
270
271#[derive(Debug, Clone)]
273pub struct DependencyRequest {
274 pub id: GlobalId,
276 pub direction: Direction,
278 pub paging: Paging,
280}
281
282#[derive(Debug, Clone, PartialEq, thiserror::Error)]
290pub enum EngineError {
291 #[error(
293 "no source named {name:?} is configured\n\
294 next: name one of the configured sources ({configured}), or add {name:?} under \
295 `sources` — `onetaskgraph sources list` shows what this configuration has."
296 )]
297 UnknownSource {
298 name: String,
300 configured: String,
302 },
303
304 #[error(
306 "{message}\n\
307 next: page with a token exactly as the previous page reported it, and against \
308 the same configuration — or drop `--page` to start the walk again."
309 )]
310 Token {
311 message: String,
313 },
314
315 #[error(
317 "no sources are configured\n\
318 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
319 prints what each plugin accepts."
320 )]
321 NoSources,
322
323 #[error(
325 "source {name} cannot be written: its plugin is {kind}, which has no write \
326 side\n\
327 next: copy into a source whose plugin can be written — `onetaskgraph sources \
328 list` reports each one's plugin."
329 )]
330 NotWritable {
331 name: String,
333 kind: String,
335 },
336
337 #[error(
344 "source {name} has no documents: its plugin is {kind}, which holds none\n\
345 next: name a source whose plugin has documents — `onetaskgraph sources list` \
346 reports each one's plugin and what it declares."
347 )]
348 NoDocuments {
349 name: String,
351 kind: String,
353 },
354
355 #[error(
360 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
361 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
362 list` reports each one's plugin."
363 )]
364 NoComments {
365 name: String,
367 kind: String,
369 },
370
371 #[error(
373 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
374 but not added to, edited or removed\n\
375 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
376 source whose plugin can be written — `onetaskgraph sources list` reports each \
377 one's plugin."
378 )]
379 CommentsNotWritable {
380 name: String,
382 kind: String,
384 },
385
386 #[error(
388 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
389 next: set the status in that source itself, or name a task of a source whose plugin \
390 can be written — `onetaskgraph sources list` reports each one's plugin."
391 )]
392 StatusNotWritable {
393 name: String,
395 kind: String,
397 },
398
399 #[error(
401 "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
402 write side\n\
403 next: set the key in that source itself, or name a {record} of a source whose plugin \
404 can be written — `onetaskgraph sources list` reports each one's plugin."
405 )]
406 MetadataNotWritable {
407 name: String,
409 kind: String,
411 record: MetadataRecord,
413 },
414
415 #[error(
417 "source {name} cannot write a priority: its plugin is {kind}, which has no write side\n\
418 next: set the priority in that source itself, or name a task of a source whose plugin \
419 can be written — `onetaskgraph sources list` reports each one's plugin."
420 )]
421 PriorityNotWritable {
422 name: String,
424 kind: String,
426 },
427
428 #[error(
430 "source {name} cannot write a task's content: its plugin is {kind}, which has no write \
431 side\n\
432 next: edit the content in that source itself, or name a task of a source whose plugin \
433 can be written — `onetaskgraph sources list` reports each one's plugin."
434 )]
435 ContentNotWritable {
436 name: String,
438 kind: String,
440 },
441
442 #[error(
444 "source {name} cannot create a {record}: its plugin is {kind}, which has no write side\n\
445 next: create it in a source whose plugin can be written — `onetaskgraph sources list` \
446 reports each one's plugin."
447 )]
448 NotCreatable {
449 name: String,
451 kind: String,
453 record: MetadataRecord,
455 },
456
457 #[error(
459 "source {name} cannot write a {record}'s rendering: its plugin is {kind}, which has no \
460 write side\n\
461 next: render it with --dry-run to read the result, or regenerate a {record} of a source \
462 whose plugin can be written."
463 )]
464 RenderingNotWritable {
465 name: String,
467 kind: String,
469 record: RenderedRecord,
471 },
472
473 #[error(
476 "supply every required answer to regenerate {id}: {} unanswered, and {reason}\n\
477 next: answer {} with --var NAME=VALUE or an answers file (--answers FILE), or run \
478 interactively to be asked.",
479 names.join(", "),
480 if names.len() == 1 { "it" } else { "each" }
481 )]
482 MissingAnswers {
483 id: String,
485 names: Vec<String>,
487 reason: String,
490 },
491
492 #[error("{error}")]
494 Template {
495 error: crate::template::TemplateError,
497 },
498
499 #[error(
501 "{record} {id} has no stored template answers: {reason}\n\
502 next: regenerate it with every required answer (`onetaskgraph {record} render {id} \
503 --var NAME=VALUE`), or read its provenance with `onetaskgraph {record} show {id}`."
504 )]
505 NoStoredAnswers {
506 record: RenderedRecord,
508 id: String,
510 reason: String,
512 },
513
514 #[error(
516 "{record} {id} records no template it was rendered from, and none was given\n\
517 next: name one with --template FILE or --template-loader FILE."
518 )]
519 NoTemplate {
520 record: RenderedRecord,
522 id: String,
524 },
525
526 #[error(
529 "{record} {id} records a template entry this product did not write — {problem}\n\
530 next: regenerate it with --template FILE or --template-loader FILE and every required \
531 answer, which records a fresh entry."
532 )]
533 MalformedProvenance {
534 record: RenderedRecord,
536 id: String,
538 problem: String,
540 },
541
542 #[error(
544 "{record} {id} was rendered from {reference:?}, which is not a readable file, so it \
545 cannot be re-read; a recorded reference is never turned into a location\n\
546 next: supply the template with --template-loader FILE (a loader document naming what \
547 to render), or name a template file with --template FILE."
548 )]
549 TemplateNotAFile {
550 record: RenderedRecord,
552 id: String,
554 reference: String,
556 },
557
558 #[error(
563 "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
564 be written to it: its plugin is {kind}, which declares priority unsupported\n\
565 next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
566 reports what each declares — or set the task's priority to none first; a \
567 github-projects source holds one once its configuration sets priority_mapping."
568 )]
569 NoPriority {
570 name: String,
572 kind: String,
574 task: String,
576 priority: onetaskgraph_plugin_api::Priority,
578 },
579
580 #[error(
582 "no project with the id {id}\n\
583 next: check the id, or list what is there — `onetaskgraph project list` reports every \
584 project the configured sources hold."
585 )]
586 NoSuchProject {
587 id: String,
589 },
590
591 #[error(
593 "no document with the id {id}\n\
594 next: check the id, or list what is there — `onetaskgraph document list` reports every \
595 document the configured sources hold."
596 )]
597 NoSuchDocument {
598 id: String,
600 },
601
602 #[error(
604 "no task with the id {id}\n\
605 next: check the id, or list what is there — `onetaskgraph task list` reports every \
606 task the configured sources hold."
607 )]
608 NoSuchTask {
609 id: String,
611 },
612
613 #[error(
615 "task {task} has no comment with the id {comment}\n\
616 next: list its comments — `onetaskgraph task comment list {task}` reports each \
617 one's id."
618 )]
619 NoSuchComment {
620 task: String,
622 comment: String,
624 },
625
626 #[error(
631 "source {name} could not be built: {error}\n\
632 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
633 command again."
634 )]
635 SourceUnavailable {
636 name: String,
638 error: SourceError,
640 },
641
642 #[error(
648 "source {name} could not do it: {error}\n\
649 next: fix what the source named above, then run the command again."
650 )]
651 SourceFailed {
652 name: String,
654 error: SourceError,
656 },
657
658 #[error(
660 "the destination source {name} could not be built: {error}\n\
661 next: fix that source — `onetaskgraph sources list` reports its state — then \
662 copy again."
663 )]
664 DestinationUnavailable {
665 name: String,
667 error: SourceError,
669 },
670
671 #[error(
673 "no item with the id {id}\n\
674 next: check the id, or list what is there — `onetaskgraph task list` and \
675 `onetaskgraph project list` report what the configured sources hold."
676 )]
677 NoSuchItem {
678 id: String,
680 },
681
682 #[error(
687 "{item} was copied from {origin}, which that destination no longer holds\n\
688 next: re-run with --recreate to create a new item there instead, or restore \
689 {origin}."
690 )]
691 StaleOrigin {
692 item: String,
694 origin: String,
696 },
697
698 #[error(
704 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
705 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
706 {id} on its own with `onetaskgraph task copy`.",
707 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
708 )]
709 NotAMember {
710 id: GlobalId,
712 projects: Vec<GlobalId>,
714 },
715
716 #[error(
725 "{item} depends on {member}, which this copy was not told to carry and which records \
726 no origin in {destination}\n\
727 next: name {member} with --member as well, record its {destination} id at \
728 onetaskgraph.origin, or copy the whole project without --member."
729 )]
730 UnrecordedMember {
731 item: GlobalId,
733 member: GlobalId,
735 destination: SourceName,
737 },
738
739 #[error(
745 "source {name} could not do it: {error}\n\
746 next: fix what the source named above, then copy again."
747 )]
748 SourceRefused {
749 name: String,
751 error: SourceError,
753 },
754
755 #[error(
763 "the copy failed and could not be undone.\n\
764 it failed because: {error}\n\
765 it could not be undone because: {refusal}\n\
766 so the destination still holds: {left_behind}\n\
767 next: remove those items at the destination, then copy again."
768 )]
769 CopyNotUndone {
770 error: Box<EngineError>,
772 left_behind: LeftBehind,
778 refusal: SourceError,
780 },
781}
782
783#[derive(Debug, Clone, PartialEq)]
792pub struct LeftBehind {
793 first: GlobalId,
795 rest: Vec<GlobalId>,
797}
798
799impl LeftBehind {
800 #[must_use]
802 pub fn new(first: GlobalId) -> Self {
803 Self {
804 first,
805 rest: Vec::new(),
806 }
807 }
808
809 pub fn push(&mut self, id: GlobalId) {
811 self.rest.push(id);
812 }
813
814 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
816 std::iter::once(&self.first).chain(self.rest.iter())
817 }
818}
819
820impl std::fmt::Display for LeftBehind {
821 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
823 write!(formatter, "{}", self.first)?;
824 for id in &self.rest {
825 write!(formatter, ", {id}")?;
826 }
827 Ok(())
828 }
829}
830
831pub enum ConfiguredSource {
838 Ready(ResolvedSource),
840 Unavailable(UnavailableSource),
842}
843
844impl ConfiguredSource {
845 #[must_use]
847 pub fn name(&self) -> &SourceName {
848 match self {
849 Self::Ready(source) => source.name(),
850 Self::Unavailable(source) => source.name(),
851 }
852 }
853}
854
855pub struct Engine {
857 sources: Vec<ConfiguredSource>,
859 selection: Vec<SourceName>,
861}
862
863impl Engine {
864 #[must_use]
872 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
873 let (ready, unavailable) = resolve_available(config, secrets);
874 Self::new(
875 ready
876 .into_iter()
877 .map(ConfiguredSource::Ready)
878 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
879 .collect(),
880 config.selected_sources(),
881 )
882 }
883
884 #[must_use]
887 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
888 Self { sources, selection }
889 }
890
891 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
893 self.sources.iter().filter_map(|source| match source {
894 ConfiguredSource::Ready(ready) => Some(ready),
895 ConfiguredSource::Unavailable(_) => None,
896 })
897 }
898
899 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
901 self.sources.iter().filter_map(|source| match source {
902 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
903 ConfiguredSource::Ready(_) => None,
904 })
905 }
906
907 #[must_use]
909 pub fn listing(&self) -> Vec<SourceListing> {
910 let mut listings: Vec<SourceListing> = self
911 .ready()
912 .map(|source| SourceListing {
913 source: source.name().clone(),
914 kind: source.kind().to_owned(),
915 state: SourceState::Available {
916 capabilities: source.source().capabilities(),
917 },
918 })
919 .chain(self.unavailable().map(|source| SourceListing {
920 source: source.name().clone(),
921 kind: source.kind().to_owned(),
922 state: SourceState::Unavailable {
923 error: source.error().clone(),
924 },
925 }))
926 .collect();
927 listings.sort_by(|left, right| left.source.cmp(&right.source));
928 listings
929 }
930
931 #[must_use]
937 pub fn has(&self, name: &SourceName) -> bool {
938 self.sources.iter().any(|source| source.name() == name)
939 }
940
941 pub async fn tasks(
949 &self,
950 request: &TaskRequest,
951 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
952 let mut names = self.resolve_selection(&request.sources)?;
953 if let ProjectSelector::Qualified(id) = &request.project {
958 self.known(&id.source)?;
959 names.retain(|name| name == &id.source);
960 }
961 let query = shape(
962 "task-list",
963 &names,
964 &(&request.filters, &request.project, &request.priorities),
965 );
966 let states = resumption(
967 self,
968 request.paging.token.as_ref(),
969 &[StreamKind::Items],
970 &query,
971 )?;
972 let budget = request.paging.limit.get();
973
974 let mut answer = Answer::new();
975 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
976
977 let shapes: Vec<TaskShape> = ready
978 .iter()
979 .map(|source| {
980 shape_tasks(
981 &source.source().capabilities(),
982 &request.filters,
983 &project_filter(&request.project),
984 &request.priorities,
985 )
986 })
987 .collect();
988 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
989 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
990
991 let walks = ready
992 .iter()
993 .enumerate()
994 .map(|(index, source)| {
995 fetch_tasks(
996 source,
997 &shapes[index],
998 &starts[index],
999 budget,
1000 &counters[index],
1001 )
1002 })
1003 .collect();
1004
1005 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1006 answer.finish(
1007 streams,
1008 budget,
1009 owed(&states),
1010 &query,
1011 |name, task: Task| {
1012 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1013 },
1014 )
1015 }
1016
1017 pub async fn projects(
1023 &self,
1024 request: &ProjectRequest,
1025 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1026 let names = self.resolve_selection(&request.sources)?;
1027 let query = shape("project-list", &names, &request.filters);
1028 let states = resumption(
1029 self,
1030 request.paging.token.as_ref(),
1031 &[StreamKind::Items],
1032 &query,
1033 )?;
1034 let budget = request.paging.limit.get();
1035
1036 let mut answer = Answer::new();
1037 let mut with_projects = Vec::new();
1043 for source in answer.split(self, &names) {
1044 if source.source().capabilities().projects.is_native() {
1045 with_projects.push(source);
1046 } else {
1047 answer.unreachable_predicate(source, Predicate::Project);
1048 }
1049 }
1050 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1051
1052 let shapes: Vec<ProjectShape> = ready
1053 .iter()
1054 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1055 .collect();
1056 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1057 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1058
1059 let walks = ready
1060 .iter()
1061 .enumerate()
1062 .map(|(index, source)| {
1063 fetch_projects(
1064 source,
1065 &shapes[index],
1066 &starts[index],
1067 budget,
1068 &counters[index],
1069 )
1070 })
1071 .collect();
1072
1073 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1074 answer.finish(
1075 streams,
1076 budget,
1077 owed(&states),
1078 &query,
1079 |name, project: Project| Qualified {
1080 id: GlobalId::new(name.clone(), project.id.clone()),
1081 item: project,
1082 },
1083 )
1084 }
1085
1086 pub async fn documents(
1098 &self,
1099 request: &DocumentRequest,
1100 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1101 let mut names = self.resolve_selection(&request.sources)?;
1102 if let ProjectSelector::Qualified(id) = &request.project {
1105 self.known(&id.source)?;
1106 names.retain(|name| name == &id.source);
1107 }
1108 let query = shape(
1109 "document-list",
1110 &names,
1111 &(&request.filters, &request.project),
1112 );
1113 let states = resumption(
1114 self,
1115 request.paging.token.as_ref(),
1116 &[StreamKind::Items],
1117 &query,
1118 )?;
1119 let budget = request.paging.limit.get();
1120
1121 let mut answer = Answer::new();
1122 let mut with_documents = Vec::new();
1123 for source in answer.split(self, &names) {
1124 if source.source().capabilities().documents.is_native() {
1125 with_documents.push(source);
1126 } else {
1127 answer.unreachable_predicate(source, Predicate::Document);
1128 }
1129 }
1130 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1131
1132 let shapes: Vec<DocumentShape> = ready
1133 .iter()
1134 .map(|source| {
1135 shape_documents(
1136 &source.source().capabilities(),
1137 &request.filters,
1138 &project_filter(&request.project),
1139 )
1140 })
1141 .collect();
1142 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1143 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1144
1145 let walks = ready
1146 .iter()
1147 .enumerate()
1148 .map(|(index, source)| {
1149 fetch_documents(
1150 source,
1151 &shapes[index],
1152 &starts[index],
1153 budget,
1154 &counters[index],
1155 )
1156 })
1157 .collect();
1158
1159 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1160 answer.finish(
1161 streams,
1162 budget,
1163 owed(&states),
1164 &query,
1165 |name, document: Document| Qualified {
1166 id: GlobalId::new(name.clone(), document.id.clone()),
1167 item: document,
1168 },
1169 )
1170 }
1171
1172 pub async fn labels(
1178 &self,
1179 request: &LabelRequest,
1180 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1181 let names = self.resolve_selection(&request.sources)?;
1182 let query = shape("label-list", &names, &());
1183 let states = resumption(
1184 self,
1185 request.paging.token.as_ref(),
1186 &[StreamKind::Items],
1187 &query,
1188 )?;
1189 let budget = request.paging.limit.get();
1190
1191 let mut answer = Answer::new();
1192 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1193 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1194 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1195
1196 let walks = ready
1197 .iter()
1198 .enumerate()
1199 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1200 .collect();
1201
1202 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1203 answer.finish(
1204 streams,
1205 budget,
1206 owed(&states),
1207 &query,
1208 |name, label: Label| Qualified {
1209 id: GlobalId::new(name.clone(), label.id.clone()),
1210 item: label,
1211 },
1212 )
1213 }
1214
1215 pub async fn search(
1221 &self,
1222 request: &SearchRequest,
1223 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1224 let names = self.resolve_selection(&request.sources)?;
1225 let reads: &[StreamKind] = match request.kind {
1229 SearchKind::Tasks => &[StreamKind::Tasks],
1230 SearchKind::Projects => &[StreamKind::Projects],
1231 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1232 };
1233 let query = shape("search", &names, &request.text);
1238 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1239 let budget = request.paging.limit.get();
1240 let filters = Filters {
1241 text: Some(request.text.clone()),
1242 ..Filters::default()
1243 };
1244
1245 let mut answer = Answer::new();
1246
1247 let mut ready = Vec::new();
1250 let mut kinds = Vec::new();
1251 let mut starts = Vec::new();
1252 for source in answer.split(self, &names) {
1253 let mut streams = Vec::new();
1254 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1255 streams.push(StreamKind::Tasks);
1256 }
1257 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1258 if source.source().capabilities().projects.is_native() {
1259 streams.push(StreamKind::Projects);
1260 } else {
1261 answer.unreachable_predicate(source, Predicate::Project);
1262 }
1263 }
1264 for stream in streams {
1265 if let Some(resume) = resume_at(&states, source.name(), stream) {
1266 ready.push(source);
1267 kinds.push(stream);
1268 starts.push(resume);
1269 }
1270 }
1271 }
1272
1273 let shapes: Vec<HitShape> = ready
1274 .iter()
1275 .zip(kinds.iter())
1276 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1277 .collect();
1278 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1279 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1280
1281 let walks = ready
1282 .iter()
1283 .enumerate()
1284 .map(|(index, source)| {
1285 fetch_hits(
1286 source,
1287 &shapes[index],
1288 &starts[index],
1289 budget,
1290 &counters[index],
1291 )
1292 })
1293 .collect();
1294
1295 let streams =
1296 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1297 answer.finish(
1298 streams,
1299 budget,
1300 owed(&states),
1301 &query,
1302 |name, found: Found| match found {
1303 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1304 GlobalId::new(name.clone(), task.id.clone()),
1305 task,
1306 )),
1307 Found::Project(project) => SearchHit::Project(Qualified {
1308 id: GlobalId::new(name.clone(), project.id.clone()),
1309 item: project,
1310 }),
1311 },
1312 )
1313 }
1314
1315 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1322 let name = self.known(&id.source)?;
1323 let mut answer = Answer::new();
1324 let selected = answer.split(self, std::slice::from_ref(&name));
1325 let Some(source) = selected.first() else {
1326 return answer.nothing();
1327 };
1328 let found = source.source().get_task(&id.native).await;
1329 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1330 answer.one(source, found, |task| {
1331 delivery::qualified_task(qualified, task)
1332 })
1333 }
1334
1335 pub async fn project(
1341 &self,
1342 id: &GlobalId,
1343 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1344 let name = self.known(&id.source)?;
1345 let mut answer = Answer::new();
1346 let selected = answer.split(self, std::slice::from_ref(&name));
1347 let Some(source) = selected.first() else {
1348 return answer.nothing();
1349 };
1350 let found = source.source().get_project(&id.native).await;
1351 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1352 answer.one(source, found, |project| Qualified {
1353 id: qualified,
1354 item: project,
1355 })
1356 }
1357
1358 pub async fn document(
1368 &self,
1369 id: &GlobalId,
1370 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1371 let name = self.known(&id.source)?;
1372 let mut answer = Answer::new();
1373 let selected = answer.split(self, std::slice::from_ref(&name));
1374 let Some(source) = selected.first() else {
1375 return answer.nothing();
1376 };
1377 if !source.source().capabilities().documents.is_native() {
1378 answer.unreachable_predicate(source, Predicate::Document);
1379 return answer.nothing();
1380 }
1381 let found = source.source().get_document(&id.native).await;
1382 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1383 answer.one(source, found, |document| Qualified {
1384 id: qualified,
1385 item: document,
1386 })
1387 }
1388
1389 pub async fn task_dependencies(
1396 &self,
1397 request: &DependencyRequest,
1398 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1399 self.dependencies(request, Entity::Task).await
1400 }
1401
1402 pub async fn project_dependencies(
1408 &self,
1409 request: &DependencyRequest,
1410 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1411 self.dependencies(request, Entity::Project).await
1412 }
1413
1414 async fn dependencies(
1417 &self,
1418 request: &DependencyRequest,
1419 entity: Entity,
1420 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1421 let name = self.known(&request.id.source)?;
1422 let query = shape(
1423 "dependencies",
1424 std::slice::from_ref(&name),
1425 &(entity, &request.id.native, request.direction),
1426 );
1427 let states = resumption(
1428 self,
1429 request.paging.token.as_ref(),
1430 &[StreamKind::Items],
1431 &query,
1432 )?;
1433 let budget = request.paging.limit.get();
1434
1435 let mut answer = Answer::new();
1436 let (ready, starts) = walking(
1437 answer.split(self, std::slice::from_ref(&name)),
1438 &states,
1439 StreamKind::Items,
1440 );
1441 let Some(source) = ready.first() else {
1442 return answer.nothing();
1443 };
1444
1445 let capabilities = source.source().capabilities();
1446 let support = match entity {
1447 Entity::Task => capabilities.task_dependencies,
1448 Entity::Project => capabilities.project_dependencies,
1449 };
1450 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1454 let mut outcomes = Outcomes::default();
1455 if request.direction == Direction::DependedOnBy {
1456 if emulating {
1457 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1458 } else {
1459 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1460 }
1461 }
1462
1463 let counters = vec![AtomicU32::new(0)];
1464 let walked = fetch_edges(
1465 source,
1466 &request.id.native,
1467 request.direction,
1468 entity,
1469 emulating,
1470 &starts[0],
1471 budget,
1472 &counters[0],
1473 )
1474 .await;
1475
1476 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1477 answer.finish(
1478 streams,
1479 budget,
1480 owed(&states),
1481 &query,
1482 |name, edge: DependencyEdge| QualifiedEdge {
1483 from: qualify_endpoint(name, edge.from),
1484 to: qualify_endpoint(name, edge.to),
1485 kind: edge.kind,
1486 },
1487 )
1488 }
1489
1490 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1492 if asked.is_empty() {
1493 if self.selection.is_empty() {
1494 return Err(EngineError::NoSources);
1495 }
1496 return Ok(self.selection.clone());
1497 }
1498 asked.iter().map(|name| self.known(name)).collect()
1499 }
1500
1501 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1503 if self.has(name) {
1504 return Ok(name.clone());
1505 }
1506 if self.sources.is_empty() {
1507 return Err(EngineError::NoSources);
1508 }
1509 Err(EngineError::UnknownSource {
1510 name: name.to_string(),
1511 configured: self
1512 .listing()
1513 .iter()
1514 .map(|listing| listing.source.to_string())
1515 .collect::<Vec<_>>()
1516 .join(", "),
1517 })
1518 }
1519}
1520
1521fn qualify_endpoint(
1522 source: &SourceName,
1523 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1524) -> QualifiedEndpoint {
1525 let kind = endpoint.kind;
1526 let is_qualified = endpoint.is_qualified();
1527 let endpoint_id = endpoint.into_id();
1528 QualifiedEndpoint {
1529 id: if is_qualified {
1530 endpoint_id
1531 .parse()
1532 .expect("plugin-api validates qualified dependency endpoints")
1533 } else {
1534 GlobalId::new(source.clone(), NativeId(endpoint_id))
1535 },
1536 kind,
1537 }
1538}
1539
1540#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1542enum Entity {
1543 Task,
1545 Project,
1547}
1548
1549enum Found {
1551 Task(Task),
1553 Project(Project),
1555}
1556
1557#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1559enum Outcome {
1560 PushedDown,
1562 AppliedLocally,
1564 Emulated,
1566 Unavailable,
1568}
1569
1570#[derive(Debug, Clone, Default, PartialEq)]
1582struct Outcomes(BTreeMap<Predicate, Outcome>);
1583
1584impl Outcomes {
1585 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1591 self.0.insert(predicate, outcome);
1592 }
1593
1594 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1597 for predicate in predicates {
1598 self.record(predicate, outcome);
1599 }
1600 }
1601
1602 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1604 self.0
1605 .iter()
1606 .filter(|(_, recorded)| **recorded == outcome)
1607 .map(|(predicate, _)| *predicate)
1608 .collect()
1609 }
1610}
1611
1612struct TaskShape {
1614 pushed: TaskQuery,
1616 local: LocalTasks,
1618 outcomes: Outcomes,
1620}
1621
1622struct ProjectShape {
1624 pushed: ProjectQuery,
1626 local: LocalProjects,
1628 outcomes: Outcomes,
1630}
1631
1632struct DocumentShape {
1634 pushed: DocumentQuery,
1636 local: LocalDocuments,
1638 outcomes: Outcomes,
1640}
1641
1642struct HitShape {
1644 stream: StreamKind,
1646 tasks: TaskQuery,
1648 projects: ProjectQuery,
1650 local_tasks: LocalTasks,
1652 local_projects: LocalProjects,
1654 outcomes: Outcomes,
1656}
1657
1658struct Answer {
1664 plans: Vec<SourcePlan>,
1666 errors: Vec<SourceFailure>,
1668}
1669
1670impl Answer {
1671 fn new() -> Self {
1672 Self {
1673 plans: Vec::new(),
1674 errors: Vec::new(),
1675 }
1676 }
1677
1678 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1683 let mut selected = Vec::new();
1684 for name in names {
1685 match engine.sources.iter().find(|source| source.name() == name) {
1686 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1687 Some(ConfiguredSource::Unavailable(source)) => {
1688 self.errors.push(source.failure());
1689 }
1690 None => {}
1691 }
1692 }
1693 selected
1694 }
1695
1696 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1698 let mut outcomes = Outcomes::default();
1699 outcomes.record(predicate, Outcome::Unavailable);
1700 self.plans.push(plan_for(source, outcomes, 0));
1701 }
1702
1703 fn collect<T>(
1705 &mut self,
1706 ready: &[&ResolvedSource],
1707 walked: Vec<Result<Fetched<T>, SourceError>>,
1708 counters: &[AtomicU32],
1709 outcomes: Vec<Outcomes>,
1710 ) -> Vec<Stream<T>> {
1711 let kinds = vec![StreamKind::Items; ready.len()];
1712 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1713 }
1714
1715 fn collect_streams<T>(
1717 &mut self,
1718 ready: &[&ResolvedSource],
1719 kinds: &[StreamKind],
1720 walked: Vec<Result<Fetched<T>, SourceError>>,
1721 counters: &[AtomicU32],
1722 outcomes: Vec<Outcomes>,
1723 ) -> Vec<Stream<T>> {
1724 let mut streams = Vec::new();
1725 for (index, result) in walked.into_iter().enumerate() {
1726 let source = ready[index];
1727 let pages = counters[index].load(Ordering::Relaxed);
1728 self.plans
1729 .push(plan_for(source, outcomes[index].clone(), pages));
1730 match result {
1731 Ok(fetched) => streams.push(Stream {
1732 source: source.name().clone(),
1733 kind: kinds[index],
1734 fetched,
1735 }),
1736 Err(error) => self.errors.push(SourceFailure {
1739 source: source.name().clone(),
1740 error,
1741 }),
1742 }
1743 }
1744 streams
1745 }
1746
1747 fn one<T, U>(
1749 mut self,
1750 source: &ResolvedSource,
1751 found: Result<Option<T>, SourceError>,
1752 qualify: impl FnOnce(T) -> U,
1753 ) -> Result<QueryResponse<U>, EngineError> {
1754 self.plans.push(plan_for(source, Outcomes::default(), 1));
1755 let items = match found {
1756 Ok(Some(item)) => vec![qualify(item)],
1757 Ok(None) => Vec::new(),
1758 Err(error) => {
1759 self.errors.push(SourceFailure {
1760 source: source.name().clone(),
1761 error,
1762 });
1763 Vec::new()
1764 }
1765 };
1766 Ok(QueryResponse {
1767 items,
1768 next: None,
1769 plan: QueryPlan {
1770 per_source: merge_plans(self.plans),
1771 },
1772 errors: self.errors,
1773 })
1774 }
1775
1776 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1778 Ok(QueryResponse {
1779 items: Vec::new(),
1780 next: None,
1781 plan: QueryPlan {
1782 per_source: merge_plans(self.plans),
1783 },
1784 errors: self.errors,
1785 })
1786 }
1787
1788 fn finish<T, U>(
1793 self,
1794 streams: Vec<Stream<T>>,
1795 budget: u32,
1796 first: Option<&Owed>,
1797 query: &str,
1798 qualify: impl Fn(&SourceName, T) -> U,
1799 ) -> Result<QueryResponse<U>, EngineError> {
1800 let (rows, states, owed) = merge(streams, budget, first);
1801 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1802 Ok(QueryResponse {
1803 items: rows
1804 .into_iter()
1805 .map(|(name, item)| qualify(&name, item))
1806 .collect(),
1807 next,
1808 plan: QueryPlan {
1809 per_source: merge_plans(self.plans),
1810 },
1811 errors: self.errors,
1812 })
1813 }
1814}
1815
1816fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1819 SourcePlan {
1820 source: source.name().clone(),
1821 kind: source.kind().to_owned(),
1822 pushed_down: outcomes.with(Outcome::PushedDown),
1823 applied_locally: outcomes.with(Outcome::AppliedLocally),
1824 emulated: outcomes.with(Outcome::Emulated),
1825 unavailable: outcomes.with(Outcome::Unavailable),
1826 pages_fetched: pages,
1827 }
1828}
1829
1830fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1835 let mut merged: Vec<SourcePlan> = Vec::new();
1836 for plan in plans {
1837 if let Some(existing) = merged
1838 .iter_mut()
1839 .find(|existing| existing.source == plan.source)
1840 {
1841 existing.pushed_down.extend(plan.pushed_down);
1842 existing.applied_locally.extend(plan.applied_locally);
1843 existing.emulated.extend(plan.emulated);
1844 existing.unavailable.extend(plan.unavailable);
1845 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1846 for list in [
1847 &mut existing.pushed_down,
1848 &mut existing.applied_locally,
1849 &mut existing.emulated,
1850 &mut existing.unavailable,
1851 ] {
1852 list.sort_unstable();
1853 list.dedup();
1854 }
1855 } else {
1856 merged.push(plan);
1857 }
1858 }
1859 merged
1860}
1861
1862fn walking<'a>(
1868 selected: Vec<&'a ResolvedSource>,
1869 states: &Option<Resumption>,
1870 kind: StreamKind,
1871) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1872 let mut ready = Vec::new();
1873 let mut starts = Vec::new();
1874 for source in selected {
1875 if let Some(resume) = resume_at(states, source.name(), kind) {
1876 ready.push(source);
1877 starts.push(resume);
1878 }
1879 }
1880 (ready, starts)
1881}
1882
1883fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1901 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1902 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1903}
1904
1905fn fingerprint(text: &str) -> String {
1907 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1908 for byte in text.as_bytes() {
1909 hash ^= u64::from(*byte);
1910 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1911 }
1912 format!("{hash:016x}")
1913}
1914
1915fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1920 document.as_ref()?.owed.as_ref()
1921}
1922
1923fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1925 match states {
1926 None => Some(Resume::default()),
1927 Some(document) => document
1928 .streams
1929 .iter()
1930 .find(|state| &state.source == source && state.stream == kind)
1931 .map(|state| state.resume.clone()),
1932 }
1933}
1934
1935fn resumption(
1966 engine: &Engine,
1967 token: Option<&PageToken>,
1968 reads: &[StreamKind],
1969 query: &str,
1970) -> Result<Option<Resumption>, EngineError> {
1971 let Some(document) = token.map(PageToken::decode) else {
1972 return Ok(None);
1973 };
1974
1975 if document.query != query {
1982 return Err(EngineError::Token {
1983 message: "this page token was written by a different query — resume the walk it \
1984 came from, or drop --page to start this one from the beginning"
1985 .to_owned(),
1986 });
1987 }
1988 let states = &document.streams;
1989
1990 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
1991 for state in states {
1992 if !reads.contains(&state.stream) {
1993 return Err(EngineError::Token {
1994 message: format!(
1995 "this page token resumes {}, which this command does not read — it \
1996 was written by a different query",
1997 state.stream.describe()
1998 ),
1999 });
2000 }
2001 let ceiling = engine
2002 .ready()
2003 .find(|source| source.name() == &state.source)
2004 .map(ceiling);
2005 if ceiling.is_none() && !engine.has(&state.source) {
2006 return Err(EngineError::Token {
2007 message: format!(
2008 "this page token resumes a source called {:?}, which this \
2009 configuration does not have",
2010 state.source.as_str()
2011 ),
2012 });
2013 }
2014 if let Some(ceiling) = ceiling
2015 && state.resume.skip >= ceiling
2016 {
2017 return Err(EngineError::Token {
2018 message: format!(
2019 "this page token resumes {} rows into a page of source {:?}, which \
2020 serves at most {ceiling}",
2021 state.resume.skip,
2022 state.source.as_str()
2023 ),
2024 });
2025 }
2026 if seen.contains(&(&state.source, state.stream)) {
2027 return Err(EngineError::Token {
2028 message: format!(
2029 "this page token gives source {:?} two places to resume from",
2030 state.source.as_str()
2031 ),
2032 });
2033 }
2034 seen.push((&state.source, state.stream));
2035 }
2036
2037 if let Some(owed) = &document.owed
2042 && !document
2043 .streams
2044 .iter()
2045 .any(|state| state.source == owed.source && state.stream == owed.stream)
2046 {
2047 return Err(EngineError::Token {
2048 message: format!(
2049 "this page token owes the next row to a stream it does not resume, \
2050 {:?}'s {}",
2051 owed.source.as_str(),
2052 owed.stream.describe()
2053 ),
2054 });
2055 }
2056
2057 Ok(Some(document))
2058}
2059
2060fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2066 match selector {
2067 ProjectSelector::Any => ProjectFilter::Any,
2068 ProjectSelector::Orphans => ProjectFilter::Orphans,
2069 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2070 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2071 }
2072}
2073
2074fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2076 match fields {
2077 TextFields::Title => vec![Predicate::SearchTitle],
2078 TextFields::Content => vec![Predicate::SearchContent],
2079 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2080 }
2081}
2082
2083fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2090 match fields {
2091 TextFields::Title => capabilities.search_title.is_native(),
2092 TextFields::Content => capabilities.search_content.is_native(),
2093 TextFields::TitleOrContent => {
2094 capabilities.search_title.is_native() && capabilities.search_content.is_native()
2095 }
2096 }
2097}
2098
2099fn shape_tasks(
2101 capabilities: &Capabilities,
2102 filters: &Filters,
2103 project: &ProjectFilter,
2104 priorities: &[Priority],
2105) -> TaskShape {
2106 let mut pushed = TaskQuery::default();
2107 let mut local = LocalTasks::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 !filters.statuses.is_empty() {
2120 if capabilities.filter_by_status.is_native() {
2121 pushed.statuses.clone_from(&filters.statuses);
2122 outcomes.record(Predicate::Status, Outcome::PushedDown);
2123 } else {
2124 local.statuses.clone_from(&filters.statuses);
2125 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2126 }
2127 }
2128 if !priorities.is_empty() {
2129 if capabilities.filter_by_priority.is_native() {
2130 pushed.priorities = priorities.to_vec();
2131 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2132 } else {
2133 local.priorities = priorities.to_vec();
2134 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2135 }
2136 }
2137 if let Some(text) = &filters.text {
2138 let predicates = text_predicates(text.fields);
2139 if searches_natively(capabilities, text.fields) {
2140 pushed.text = Some(text.clone());
2141 outcomes.record_all(predicates, Outcome::PushedDown);
2142 } else {
2143 local.text = Some(text.clone());
2144 outcomes.record_all(predicates, Outcome::AppliedLocally);
2145 }
2146 }
2147 match project {
2148 ProjectFilter::Any => {}
2149 ProjectFilter::Orphans => {
2150 if capabilities.orphan_tasks.is_native() {
2151 pushed.project = ProjectFilter::Orphans;
2152 outcomes.record(Predicate::Project, Outcome::PushedDown);
2153 } else {
2154 local.project = Some(ProjectFilter::Orphans);
2155 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2156 }
2157 }
2158 ProjectFilter::Is(id) => {
2159 if capabilities.projects.is_native() {
2160 pushed.project = ProjectFilter::Is(id.clone());
2161 outcomes.record(Predicate::Project, Outcome::PushedDown);
2162 } else {
2163 local.project = Some(ProjectFilter::Is(id.clone()));
2164 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2165 }
2166 }
2167 }
2168
2169 TaskShape {
2170 pushed,
2171 local,
2172 outcomes,
2173 }
2174}
2175
2176fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2178 let mut pushed = ProjectQuery::default();
2179 let mut local = LocalProjects::default();
2180 let mut outcomes = Outcomes::default();
2181
2182 if !filters.labels.is_empty() {
2183 if capabilities.filter_by_label.is_native() {
2184 pushed.labels = filters.labels.clone();
2185 outcomes.record(Predicate::Label, Outcome::PushedDown);
2186 } else {
2187 local.labels = Some(filters.labels.clone());
2188 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2189 }
2190 }
2191 if !filters.statuses.is_empty() {
2192 if capabilities.filter_by_status.is_native() {
2193 pushed.statuses.clone_from(&filters.statuses);
2194 outcomes.record(Predicate::Status, Outcome::PushedDown);
2195 } else {
2196 local.statuses.clone_from(&filters.statuses);
2197 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2198 }
2199 }
2200 if let Some(text) = &filters.text {
2201 let predicates = text_predicates(text.fields);
2202 if searches_natively(capabilities, text.fields) {
2203 pushed.text = Some(text.clone());
2204 outcomes.record_all(predicates, Outcome::PushedDown);
2205 } else {
2206 local.text = Some(text.clone());
2207 outcomes.record_all(predicates, Outcome::AppliedLocally);
2208 }
2209 }
2210
2211 ProjectShape {
2212 pushed,
2213 local,
2214 outcomes,
2215 }
2216}
2217
2218fn shape_documents(
2223 capabilities: &Capabilities,
2224 filters: &DocumentFilters,
2225 project: &ProjectFilter,
2226) -> DocumentShape {
2227 let mut pushed = DocumentQuery::default();
2228 let mut local = LocalDocuments::default();
2229 let mut outcomes = Outcomes::default();
2230
2231 if !filters.labels.is_empty() {
2232 if capabilities.filter_by_label.is_native() {
2233 pushed.labels = filters.labels.clone();
2234 outcomes.record(Predicate::Label, Outcome::PushedDown);
2235 } else {
2236 local.labels = Some(filters.labels.clone());
2237 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2238 }
2239 }
2240 if let Some(text) = &filters.text {
2241 let predicates = text_predicates(text.fields);
2242 if searches_natively(capabilities, text.fields) {
2243 pushed.text = Some(text.clone());
2244 outcomes.record_all(predicates, Outcome::PushedDown);
2245 } else {
2246 local.text = Some(text.clone());
2247 outcomes.record_all(predicates, Outcome::AppliedLocally);
2248 }
2249 }
2250 match project {
2251 ProjectFilter::Any => {}
2252 ProjectFilter::Orphans => {
2253 if capabilities.orphan_tasks.is_native() {
2254 pushed.project = ProjectFilter::Orphans;
2255 outcomes.record(Predicate::Project, Outcome::PushedDown);
2256 } else {
2257 local.project = Some(ProjectFilter::Orphans);
2258 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2259 }
2260 }
2261 ProjectFilter::Is(id) => {
2262 if capabilities.projects.is_native() {
2263 pushed.project = ProjectFilter::Is(id.clone());
2264 outcomes.record(Predicate::Project, Outcome::PushedDown);
2265 } else {
2266 local.project = Some(ProjectFilter::Is(id.clone()));
2267 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2268 }
2269 }
2270 }
2271
2272 DocumentShape {
2273 pushed,
2274 local,
2275 outcomes,
2276 }
2277}
2278
2279fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2281 match stream {
2282 StreamKind::Projects => {
2283 let shaped = shape_projects(capabilities, filters);
2284 HitShape {
2285 stream,
2286 tasks: TaskQuery::default(),
2287 projects: shaped.pushed,
2288 local_tasks: LocalTasks::default(),
2289 local_projects: shaped.local,
2290 outcomes: shaped.outcomes,
2291 }
2292 }
2293 StreamKind::Items | StreamKind::Tasks => {
2294 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any, &[]);
2295 HitShape {
2296 stream,
2297 tasks: shaped.pushed,
2298 projects: ProjectQuery::default(),
2299 local_tasks: shaped.local,
2300 local_projects: LocalProjects::default(),
2301 outcomes: shaped.outcomes,
2302 }
2303 }
2304 }
2305}
2306
2307fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2314 if compensating {
2315 ceiling
2316 } else {
2317 budget.min(ceiling)
2318 }
2319}
2320
2321fn ceiling(source: &ResolvedSource) -> u32 {
2323 source.source().capabilities().max_page_size.max(1)
2324}
2325
2326async fn fetch_tasks(
2328 source: &ResolvedSource,
2329 shape: &TaskShape,
2330 start: &Resume,
2331 budget: u32,
2332 calls: &AtomicU32,
2333) -> Result<Fetched<Task>, SourceError> {
2334 let compensating = shape.local != LocalTasks::default();
2335 walk(
2336 start,
2337 budget,
2338 page_size(compensating, budget, ceiling(source)),
2339 |task| shape.local.keeps(task),
2340 |cursor, limit| async move {
2341 calls.fetch_add(1, Ordering::Relaxed);
2342 let request = PageRequest { cursor, limit };
2343 source.source().query_tasks(&shape.pushed, &request).await
2344 },
2345 )
2346 .await
2347}
2348
2349async fn fetch_projects(
2351 source: &ResolvedSource,
2352 shape: &ProjectShape,
2353 start: &Resume,
2354 budget: u32,
2355 calls: &AtomicU32,
2356) -> Result<Fetched<Project>, SourceError> {
2357 let compensating = shape.local != LocalProjects::default();
2358 walk(
2359 start,
2360 budget,
2361 page_size(compensating, budget, ceiling(source)),
2362 |project| shape.local.keeps(project),
2363 |cursor, limit| async move {
2364 calls.fetch_add(1, Ordering::Relaxed);
2365 let request = PageRequest { cursor, limit };
2366 source
2367 .source()
2368 .query_projects(&shape.pushed, &request)
2369 .await
2370 },
2371 )
2372 .await
2373}
2374
2375async fn fetch_documents(
2377 source: &ResolvedSource,
2378 shape: &DocumentShape,
2379 start: &Resume,
2380 budget: u32,
2381 calls: &AtomicU32,
2382) -> Result<Fetched<Document>, SourceError> {
2383 let compensating = shape.local != LocalDocuments::default();
2384 walk(
2385 start,
2386 budget,
2387 page_size(compensating, budget, ceiling(source)),
2388 |document| shape.local.keeps(document),
2389 |cursor, limit| async move {
2390 calls.fetch_add(1, Ordering::Relaxed);
2391 let request = PageRequest { cursor, limit };
2392 source
2393 .source()
2394 .query_documents(&shape.pushed, &request)
2395 .await
2396 },
2397 )
2398 .await
2399}
2400
2401async fn fetch_labels(
2403 source: &ResolvedSource,
2404 start: &Resume,
2405 budget: u32,
2406 calls: &AtomicU32,
2407) -> Result<Fetched<Label>, SourceError> {
2408 walk(
2409 start,
2410 budget,
2411 page_size(false, budget, ceiling(source)),
2412 |_| true,
2413 |cursor, limit| async move {
2414 calls.fetch_add(1, Ordering::Relaxed);
2415 let request = PageRequest { cursor, limit };
2416 source.source().labels(&request).await
2417 },
2418 )
2419 .await
2420}
2421
2422async fn fetch_hits(
2424 source: &ResolvedSource,
2425 shape: &HitShape,
2426 start: &Resume,
2427 budget: u32,
2428 calls: &AtomicU32,
2429) -> Result<Fetched<Found>, SourceError> {
2430 let ceiling = ceiling(source);
2431 match shape.stream {
2432 StreamKind::Projects => {
2433 let compensating = shape.local_projects != LocalProjects::default();
2434 walk(
2435 start,
2436 budget,
2437 page_size(compensating, budget, ceiling),
2438 |found| match found {
2439 Found::Project(project) => shape.local_projects.keeps(project),
2440 Found::Task(_) => true,
2441 },
2442 |cursor, limit| async move {
2443 calls.fetch_add(1, Ordering::Relaxed);
2444 let request = PageRequest { cursor, limit };
2445 let page = source
2446 .source()
2447 .query_projects(&shape.projects, &request)
2448 .await?;
2449 Ok(Page {
2450 items: page.items.into_iter().map(Found::Project).collect(),
2451 next: page.next,
2452 })
2453 },
2454 )
2455 .await
2456 }
2457 StreamKind::Items | StreamKind::Tasks => {
2458 let compensating = shape.local_tasks != LocalTasks::default();
2459 walk(
2460 start,
2461 budget,
2462 page_size(compensating, budget, ceiling),
2463 |found| match found {
2464 Found::Task(task) => shape.local_tasks.keeps(task),
2465 Found::Project(_) => true,
2466 },
2467 |cursor, limit| async move {
2468 calls.fetch_add(1, Ordering::Relaxed);
2469 let request = PageRequest { cursor, limit };
2470 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2471 Ok(Page {
2472 items: page.items.into_iter().map(Found::Task).collect(),
2473 next: page.next,
2474 })
2475 },
2476 )
2477 .await
2478 }
2479 }
2480}
2481
2482async fn forward_edges(
2484 source: &ResolvedSource,
2485 entity: Entity,
2486 id: &NativeId,
2487 request: &PageRequest,
2488) -> Result<Page<DependencyEdge>, SourceError> {
2489 match entity {
2490 Entity::Task => {
2491 source
2492 .source()
2493 .task_dependencies(id, Direction::DependsOn, request)
2494 .await
2495 }
2496 Entity::Project => {
2497 source
2498 .source()
2499 .project_dependencies(id, Direction::DependsOn, request)
2500 .await
2501 }
2502 }
2503}
2504
2505#[expect(
2514 clippy::too_many_arguments,
2515 reason = "every argument is one axis of one walk — the source, the item, the \
2516 direction, which of its two graphs, whether the reverse is emulated, where \
2517 to resume, how many rows to return and where to count calls. Grouping them \
2518 into a struct would name the same eight values one indirection further from \
2519 the loop that reads them."
2520)]
2521async fn fetch_edges(
2522 source: &ResolvedSource,
2523 native: &NativeId,
2524 direction: Direction,
2525 entity: Entity,
2526 emulating: bool,
2527 start: &Resume,
2528 budget: u32,
2529 calls: &AtomicU32,
2530) -> Result<Fetched<DependencyEdge>, SourceError> {
2531 let ceiling = ceiling(source);
2532 if !emulating {
2533 return walk(
2534 start,
2535 budget,
2536 page_size(false, budget, ceiling),
2537 |_| true,
2538 |cursor, limit| async move {
2539 calls.fetch_add(1, Ordering::Relaxed);
2540 let request = PageRequest { cursor, limit };
2541 match entity {
2542 Entity::Task => {
2543 source
2544 .source()
2545 .task_dependencies(native, direction, &request)
2546 .await
2547 }
2548 Entity::Project => {
2549 source
2550 .source()
2551 .project_dependencies(native, direction, &request)
2552 .await
2553 }
2554 }
2555 },
2556 )
2557 .await;
2558 }
2559
2560 walk(
2561 start,
2562 budget,
2563 ceiling,
2564 |_| true,
2565 |cursor, limit| async move {
2566 calls.fetch_add(1, Ordering::Relaxed);
2567 let request = PageRequest { cursor, limit };
2568 let (ids, next) = match entity {
2569 Entity::Task => {
2570 let page = source
2571 .source()
2572 .query_tasks(&TaskQuery::default(), &request)
2573 .await?;
2574 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2575 (ids, page.next)
2576 }
2577 Entity::Project => {
2578 let page = source
2579 .source()
2580 .query_projects(&ProjectQuery::default(), &request)
2581 .await?;
2582 let ids: Vec<NativeId> =
2583 page.items.into_iter().map(|project| project.id).collect();
2584 (ids, page.next)
2585 }
2586 };
2587
2588 let mut edges = Vec::new();
2589 for id in ids {
2590 let mut inner: Option<Cursor> = None;
2591 loop {
2592 calls.fetch_add(1, Ordering::Relaxed);
2593 let request = PageRequest {
2594 cursor: inner.clone(),
2595 limit,
2596 };
2597 let page = forward_edges(source, entity, &id, &request).await?;
2598 fits(page.items.len(), limit)?;
2602 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2603 unrepeated(
2604 page.next.as_ref(),
2605 inner.as_ref(),
2606 "its forward edges were being scanned",
2607 )?;
2608 match page.next {
2609 Some(cursor) => inner = Some(cursor),
2610 None => break,
2611 }
2612 }
2613 }
2614
2615 Ok(Page { items: edges, next })
2616 },
2617 )
2618 .await
2619}