1mod copy;
19mod fetch;
20mod join;
21mod local;
22mod resume;
23
24use std::collections::BTreeMap;
25use std::num::NonZeroU32;
26use std::sync::atomic::{AtomicU32, Ordering};
27
28use onetaskgraph_plugin_api::{
29 Capabilities, Cursor, DependencyEdge, Direction, Label, LabelFilter, NativeId, Page,
30 PageRequest, Project, ProjectFilter, ProjectQuery, SecretResolver, SourceError, SourceName,
31 StatusCategory, Task, TaskQuery, TextFields, TextQuery,
32};
33use schemars::JsonSchema;
34use serde::{Deserialize, Serialize};
35
36use crate::GlobalId;
37use crate::config::Config;
38use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
39use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
40
41use fetch::{Fetched, Stream, fits, merge, walk};
42use join::join_all;
43use local::{LocalProjects, LocalTasks};
44pub(crate) use resume::{Owed, Resumption, StreamState};
45use resume::{Resume, StreamKind};
46
47pub use copy::{CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy};
48pub use local::ProjectSelector;
49
50#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
55pub struct Qualified<T> {
56 pub id: GlobalId,
58 pub item: T,
60}
61
62#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
67pub struct QualifiedEdge {
68 pub from: QualifiedEndpoint,
73 pub to: QualifiedEndpoint,
75 pub kind: onetaskgraph_plugin_api::DependencyKind,
77}
78
79#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
81pub struct QualifiedEndpoint {
82 pub id: GlobalId,
84 pub kind: onetaskgraph_plugin_api::ItemKind,
86}
87
88impl std::fmt::Display for QualifiedEndpoint {
89 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
90 self.id.fmt(formatter)
91 }
92}
93
94#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
96#[serde(tag = "kind", rename_all = "kebab-case")]
97pub enum SearchHit {
98 Task(Qualified<Task>),
100 Project(Qualified<Project>),
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
106#[serde(rename_all = "kebab-case")]
107pub enum SearchKind {
108 Tasks,
110 Projects,
112 #[default]
114 Both,
115}
116
117#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
119pub struct SourceListing {
120 pub source: SourceName,
122 pub kind: String,
134 #[serde(flatten)]
136 pub state: SourceState,
137}
138
139#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
141#[serde(tag = "state", rename_all = "kebab-case")]
142pub enum SourceState {
143 Available {
145 capabilities: Capabilities,
147 },
148 Unavailable {
150 error: SourceError,
152 },
153}
154
155#[derive(Debug, Clone, PartialEq)]
157pub struct Paging {
158 pub limit: NonZeroU32,
160 pub token: Option<PageToken>,
162}
163
164#[derive(Debug, Clone, Default, PartialEq)]
166pub struct Filters {
167 pub text: Option<TextQuery>,
169 pub labels: LabelFilter,
171 pub statuses: Vec<StatusCategory>,
173}
174
175#[derive(Debug, Clone)]
177pub struct TaskRequest {
178 pub sources: Vec<SourceName>,
180 pub filters: Filters,
182 pub project: ProjectSelector,
184 pub paging: Paging,
186}
187
188#[derive(Debug, Clone)]
190pub struct ProjectRequest {
191 pub sources: Vec<SourceName>,
193 pub filters: Filters,
195 pub paging: Paging,
197}
198
199#[derive(Debug, Clone)]
201pub struct LabelRequest {
202 pub sources: Vec<SourceName>,
204 pub paging: Paging,
206}
207
208#[derive(Debug, Clone)]
210pub struct SearchRequest {
211 pub sources: Vec<SourceName>,
213 pub text: TextQuery,
215 pub kind: SearchKind,
217 pub paging: Paging,
219}
220
221#[derive(Debug, Clone)]
223pub struct DependencyRequest {
224 pub id: GlobalId,
226 pub direction: Direction,
228 pub paging: Paging,
230}
231
232#[derive(Debug, Clone, PartialEq, thiserror::Error)]
240pub enum EngineError {
241 #[error(
243 "no source named {name:?} is configured\n\
244 next: name one of the configured sources ({configured}), or add {name:?} under \
245 `sources` — `onetaskgraph sources list` shows what this configuration has."
246 )]
247 UnknownSource {
248 name: String,
250 configured: String,
252 },
253
254 #[error(
256 "{message}\n\
257 next: page with a token exactly as the previous page reported it, and against \
258 the same configuration — or drop `--page` to start the walk again."
259 )]
260 Token {
261 message: String,
263 },
264
265 #[error(
267 "no sources are configured\n\
268 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
269 prints what each plugin accepts."
270 )]
271 NoSources,
272
273 #[error(
275 "source {name} cannot be written: its plugin is {kind}, which has no write \
276 side\n\
277 next: copy into a source whose plugin can be written — `onetaskgraph sources \
278 list` reports each one's plugin."
279 )]
280 NotWritable {
281 name: String,
283 kind: String,
285 },
286
287 #[error(
289 "the destination source {name} could not be built: {error}\n\
290 next: fix that source — `onetaskgraph sources list` reports its state — then \
291 copy again."
292 )]
293 DestinationUnavailable {
294 name: String,
296 error: SourceError,
298 },
299
300 #[error(
302 "no item with the id {id}\n\
303 next: check the id, or list what is there — `onetaskgraph task list` and \
304 `onetaskgraph project list` report what the configured sources hold."
305 )]
306 NoSuchItem {
307 id: String,
309 },
310
311 #[error(
316 "{item} was copied from {origin}, which that destination no longer holds\n\
317 next: re-run with --recreate to create a new item there instead, or restore \
318 {origin}."
319 )]
320 StaleOrigin {
321 item: String,
323 origin: String,
325 },
326
327 #[error(
333 "source {name} could not do it: {error}\n\
334 next: fix what the source named above, then copy again."
335 )]
336 SourceRefused {
337 name: String,
339 error: SourceError,
341 },
342}
343
344pub enum ConfiguredSource {
351 Ready(ResolvedSource),
353 Unavailable(UnavailableSource),
355}
356
357impl ConfiguredSource {
358 #[must_use]
360 pub fn name(&self) -> &SourceName {
361 match self {
362 Self::Ready(source) => source.name(),
363 Self::Unavailable(source) => source.name(),
364 }
365 }
366}
367
368pub struct Engine {
370 sources: Vec<ConfiguredSource>,
372 selection: Vec<SourceName>,
374}
375
376impl Engine {
377 #[must_use]
385 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
386 let (ready, unavailable) = resolve_available(config, secrets);
387 Self::new(
388 ready
389 .into_iter()
390 .map(ConfiguredSource::Ready)
391 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
392 .collect(),
393 config.selected_sources(),
394 )
395 }
396
397 #[must_use]
400 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
401 Self { sources, selection }
402 }
403
404 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
406 self.sources.iter().filter_map(|source| match source {
407 ConfiguredSource::Ready(ready) => Some(ready),
408 ConfiguredSource::Unavailable(_) => None,
409 })
410 }
411
412 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
414 self.sources.iter().filter_map(|source| match source {
415 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
416 ConfiguredSource::Ready(_) => None,
417 })
418 }
419
420 #[must_use]
422 pub fn listing(&self) -> Vec<SourceListing> {
423 let mut listings: Vec<SourceListing> = self
424 .ready()
425 .map(|source| SourceListing {
426 source: source.name().clone(),
427 kind: source.kind().to_owned(),
428 state: SourceState::Available {
429 capabilities: source.source().capabilities(),
430 },
431 })
432 .chain(self.unavailable().map(|source| SourceListing {
433 source: source.name().clone(),
434 kind: source.kind().to_owned(),
435 state: SourceState::Unavailable {
436 error: source.error().clone(),
437 },
438 }))
439 .collect();
440 listings.sort_by(|left, right| left.source.cmp(&right.source));
441 listings
442 }
443
444 #[must_use]
450 pub fn has(&self, name: &SourceName) -> bool {
451 self.sources.iter().any(|source| source.name() == name)
452 }
453
454 pub async fn tasks(
462 &self,
463 request: &TaskRequest,
464 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
465 let mut names = self.resolve_selection(&request.sources)?;
466 if let ProjectSelector::Qualified(id) = &request.project {
471 self.known(&id.source)?;
472 names.retain(|name| name == &id.source);
473 }
474 let query = shape("task-list", &names, &(&request.filters, &request.project));
475 let states = resumption(
476 self,
477 request.paging.token.as_ref(),
478 &[StreamKind::Items],
479 &query,
480 )?;
481 let budget = request.paging.limit.get();
482
483 let mut answer = Answer::new();
484 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
485
486 let shapes: Vec<TaskShape> = ready
487 .iter()
488 .map(|source| {
489 shape_tasks(
490 &source.source().capabilities(),
491 &request.filters,
492 &project_filter(&request.project),
493 )
494 })
495 .collect();
496 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
497 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
498
499 let walks = ready
500 .iter()
501 .enumerate()
502 .map(|(index, source)| {
503 fetch_tasks(
504 source,
505 &shapes[index],
506 &starts[index],
507 budget,
508 &counters[index],
509 )
510 })
511 .collect();
512
513 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
514 answer.finish(
515 streams,
516 budget,
517 owed(&states),
518 &query,
519 |name, task: Task| Qualified {
520 id: GlobalId::new(name.clone(), task.id.clone()),
521 item: task,
522 },
523 )
524 }
525
526 pub async fn projects(
532 &self,
533 request: &ProjectRequest,
534 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
535 let names = self.resolve_selection(&request.sources)?;
536 let query = shape("project-list", &names, &request.filters);
537 let states = resumption(
538 self,
539 request.paging.token.as_ref(),
540 &[StreamKind::Items],
541 &query,
542 )?;
543 let budget = request.paging.limit.get();
544
545 let mut answer = Answer::new();
546 let mut with_projects = Vec::new();
552 for source in answer.split(self, &names) {
553 if source.source().capabilities().projects.is_native() {
554 with_projects.push(source);
555 } else {
556 answer.unreachable_predicate(source, Predicate::Project);
557 }
558 }
559 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
560
561 let shapes: Vec<ProjectShape> = ready
562 .iter()
563 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
564 .collect();
565 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
566 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
567
568 let walks = ready
569 .iter()
570 .enumerate()
571 .map(|(index, source)| {
572 fetch_projects(
573 source,
574 &shapes[index],
575 &starts[index],
576 budget,
577 &counters[index],
578 )
579 })
580 .collect();
581
582 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
583 answer.finish(
584 streams,
585 budget,
586 owed(&states),
587 &query,
588 |name, project: Project| Qualified {
589 id: GlobalId::new(name.clone(), project.id.clone()),
590 item: project,
591 },
592 )
593 }
594
595 pub async fn labels(
601 &self,
602 request: &LabelRequest,
603 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
604 let names = self.resolve_selection(&request.sources)?;
605 let query = shape("label-list", &names, &());
606 let states = resumption(
607 self,
608 request.paging.token.as_ref(),
609 &[StreamKind::Items],
610 &query,
611 )?;
612 let budget = request.paging.limit.get();
613
614 let mut answer = Answer::new();
615 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
616 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
617 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
618
619 let walks = ready
620 .iter()
621 .enumerate()
622 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
623 .collect();
624
625 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
626 answer.finish(
627 streams,
628 budget,
629 owed(&states),
630 &query,
631 |name, label: Label| Qualified {
632 id: GlobalId::new(name.clone(), label.id.clone()),
633 item: label,
634 },
635 )
636 }
637
638 pub async fn search(
644 &self,
645 request: &SearchRequest,
646 ) -> Result<QueryResponse<SearchHit>, EngineError> {
647 let names = self.resolve_selection(&request.sources)?;
648 let reads: &[StreamKind] = match request.kind {
652 SearchKind::Tasks => &[StreamKind::Tasks],
653 SearchKind::Projects => &[StreamKind::Projects],
654 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
655 };
656 let query = shape("search", &names, &request.text);
661 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
662 let budget = request.paging.limit.get();
663 let filters = Filters {
664 text: Some(request.text.clone()),
665 ..Filters::default()
666 };
667
668 let mut answer = Answer::new();
669
670 let mut ready = Vec::new();
673 let mut kinds = Vec::new();
674 let mut starts = Vec::new();
675 for source in answer.split(self, &names) {
676 let mut streams = Vec::new();
677 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
678 streams.push(StreamKind::Tasks);
679 }
680 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
681 if source.source().capabilities().projects.is_native() {
682 streams.push(StreamKind::Projects);
683 } else {
684 answer.unreachable_predicate(source, Predicate::Project);
685 }
686 }
687 for stream in streams {
688 if let Some(resume) = resume_at(&states, source.name(), stream) {
689 ready.push(source);
690 kinds.push(stream);
691 starts.push(resume);
692 }
693 }
694 }
695
696 let shapes: Vec<HitShape> = ready
697 .iter()
698 .zip(kinds.iter())
699 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
700 .collect();
701 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
702 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
703
704 let walks = ready
705 .iter()
706 .enumerate()
707 .map(|(index, source)| {
708 fetch_hits(
709 source,
710 &shapes[index],
711 &starts[index],
712 budget,
713 &counters[index],
714 )
715 })
716 .collect();
717
718 let streams =
719 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
720 answer.finish(
721 streams,
722 budget,
723 owed(&states),
724 &query,
725 |name, found: Found| match found {
726 Found::Task(task) => SearchHit::Task(Qualified {
727 id: GlobalId::new(name.clone(), task.id.clone()),
728 item: task,
729 }),
730 Found::Project(project) => SearchHit::Project(Qualified {
731 id: GlobalId::new(name.clone(), project.id.clone()),
732 item: project,
733 }),
734 },
735 )
736 }
737
738 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
745 let name = self.known(&id.source)?;
746 let mut answer = Answer::new();
747 let selected = answer.split(self, std::slice::from_ref(&name));
748 let Some(source) = selected.first() else {
749 return answer.nothing();
750 };
751 let found = source.source().get_task(&id.native).await;
752 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
753 answer.one(source, found, |task| Qualified {
754 id: qualified,
755 item: task,
756 })
757 }
758
759 pub async fn project(
765 &self,
766 id: &GlobalId,
767 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
768 let name = self.known(&id.source)?;
769 let mut answer = Answer::new();
770 let selected = answer.split(self, std::slice::from_ref(&name));
771 let Some(source) = selected.first() else {
772 return answer.nothing();
773 };
774 let found = source.source().get_project(&id.native).await;
775 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
776 answer.one(source, found, |project| Qualified {
777 id: qualified,
778 item: project,
779 })
780 }
781
782 pub async fn task_dependencies(
789 &self,
790 request: &DependencyRequest,
791 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
792 self.dependencies(request, Entity::Task).await
793 }
794
795 pub async fn project_dependencies(
801 &self,
802 request: &DependencyRequest,
803 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
804 self.dependencies(request, Entity::Project).await
805 }
806
807 async fn dependencies(
810 &self,
811 request: &DependencyRequest,
812 entity: Entity,
813 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
814 let name = self.known(&request.id.source)?;
815 let query = shape(
816 "dependencies",
817 std::slice::from_ref(&name),
818 &(entity, &request.id.native, request.direction),
819 );
820 let states = resumption(
821 self,
822 request.paging.token.as_ref(),
823 &[StreamKind::Items],
824 &query,
825 )?;
826 let budget = request.paging.limit.get();
827
828 let mut answer = Answer::new();
829 let (ready, starts) = walking(
830 answer.split(self, std::slice::from_ref(&name)),
831 &states,
832 StreamKind::Items,
833 );
834 let Some(source) = ready.first() else {
835 return answer.nothing();
836 };
837
838 let capabilities = source.source().capabilities();
839 let support = match entity {
840 Entity::Task => capabilities.task_dependencies,
841 Entity::Project => capabilities.project_dependencies,
842 };
843 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
847 let mut outcomes = Outcomes::default();
848 if request.direction == Direction::DependedOnBy {
849 if emulating {
850 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
851 } else {
852 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
853 }
854 }
855
856 let counters = vec![AtomicU32::new(0)];
857 let walked = fetch_edges(
858 source,
859 &request.id.native,
860 request.direction,
861 entity,
862 emulating,
863 &starts[0],
864 budget,
865 &counters[0],
866 )
867 .await;
868
869 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
870 answer.finish(
871 streams,
872 budget,
873 owed(&states),
874 &query,
875 |name, edge: DependencyEdge| QualifiedEdge {
876 from: qualify_endpoint(name, edge.from),
877 to: qualify_endpoint(name, edge.to),
878 kind: edge.kind,
879 },
880 )
881 }
882
883 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
885 if asked.is_empty() {
886 if self.selection.is_empty() {
887 return Err(EngineError::NoSources);
888 }
889 return Ok(self.selection.clone());
890 }
891 asked.iter().map(|name| self.known(name)).collect()
892 }
893
894 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
896 if self.has(name) {
897 return Ok(name.clone());
898 }
899 if self.sources.is_empty() {
900 return Err(EngineError::NoSources);
901 }
902 Err(EngineError::UnknownSource {
903 name: name.to_string(),
904 configured: self
905 .listing()
906 .iter()
907 .map(|listing| listing.source.to_string())
908 .collect::<Vec<_>>()
909 .join(", "),
910 })
911 }
912}
913
914fn qualify_endpoint(
915 source: &SourceName,
916 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
917) -> QualifiedEndpoint {
918 let kind = endpoint.kind;
919 let is_qualified = endpoint.is_qualified();
920 let endpoint_id = endpoint.into_id();
921 QualifiedEndpoint {
922 id: if is_qualified {
923 endpoint_id
924 .parse()
925 .expect("plugin-api validates qualified dependency endpoints")
926 } else {
927 GlobalId::new(source.clone(), NativeId(endpoint_id))
928 },
929 kind,
930 }
931}
932
933#[derive(Debug, Clone, Copy, PartialEq, Eq)]
935enum Entity {
936 Task,
938 Project,
940}
941
942enum Found {
944 Task(Task),
946 Project(Project),
948}
949
950#[derive(Debug, Clone, Copy, PartialEq, Eq)]
952enum Outcome {
953 PushedDown,
955 AppliedLocally,
957 Emulated,
959 Unavailable,
961}
962
963#[derive(Debug, Clone, Default, PartialEq)]
975struct Outcomes(BTreeMap<Predicate, Outcome>);
976
977impl Outcomes {
978 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
984 self.0.insert(predicate, outcome);
985 }
986
987 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
990 for predicate in predicates {
991 self.record(predicate, outcome);
992 }
993 }
994
995 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
997 self.0
998 .iter()
999 .filter(|(_, recorded)| **recorded == outcome)
1000 .map(|(predicate, _)| *predicate)
1001 .collect()
1002 }
1003}
1004
1005struct TaskShape {
1007 pushed: TaskQuery,
1009 local: LocalTasks,
1011 outcomes: Outcomes,
1013}
1014
1015struct ProjectShape {
1017 pushed: ProjectQuery,
1019 local: LocalProjects,
1021 outcomes: Outcomes,
1023}
1024
1025struct HitShape {
1027 stream: StreamKind,
1029 tasks: TaskQuery,
1031 projects: ProjectQuery,
1033 local_tasks: LocalTasks,
1035 local_projects: LocalProjects,
1037 outcomes: Outcomes,
1039}
1040
1041struct Answer {
1047 plans: Vec<SourcePlan>,
1049 errors: Vec<SourceFailure>,
1051}
1052
1053impl Answer {
1054 fn new() -> Self {
1055 Self {
1056 plans: Vec::new(),
1057 errors: Vec::new(),
1058 }
1059 }
1060
1061 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1066 let mut selected = Vec::new();
1067 for name in names {
1068 match engine.sources.iter().find(|source| source.name() == name) {
1069 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1070 Some(ConfiguredSource::Unavailable(source)) => {
1071 self.errors.push(source.failure());
1072 }
1073 None => {}
1074 }
1075 }
1076 selected
1077 }
1078
1079 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1081 let mut outcomes = Outcomes::default();
1082 outcomes.record(predicate, Outcome::Unavailable);
1083 self.plans.push(plan_for(source, outcomes, 0));
1084 }
1085
1086 fn collect<T>(
1088 &mut self,
1089 ready: &[&ResolvedSource],
1090 walked: Vec<Result<Fetched<T>, SourceError>>,
1091 counters: &[AtomicU32],
1092 outcomes: Vec<Outcomes>,
1093 ) -> Vec<Stream<T>> {
1094 let kinds = vec![StreamKind::Items; ready.len()];
1095 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1096 }
1097
1098 fn collect_streams<T>(
1100 &mut self,
1101 ready: &[&ResolvedSource],
1102 kinds: &[StreamKind],
1103 walked: Vec<Result<Fetched<T>, SourceError>>,
1104 counters: &[AtomicU32],
1105 outcomes: Vec<Outcomes>,
1106 ) -> Vec<Stream<T>> {
1107 let mut streams = Vec::new();
1108 for (index, result) in walked.into_iter().enumerate() {
1109 let source = ready[index];
1110 let pages = counters[index].load(Ordering::Relaxed);
1111 self.plans
1112 .push(plan_for(source, outcomes[index].clone(), pages));
1113 match result {
1114 Ok(fetched) => streams.push(Stream {
1115 source: source.name().clone(),
1116 kind: kinds[index],
1117 fetched,
1118 }),
1119 Err(error) => self.errors.push(SourceFailure {
1122 source: source.name().clone(),
1123 error,
1124 }),
1125 }
1126 }
1127 streams
1128 }
1129
1130 fn one<T, U>(
1132 mut self,
1133 source: &ResolvedSource,
1134 found: Result<Option<T>, SourceError>,
1135 qualify: impl FnOnce(T) -> U,
1136 ) -> Result<QueryResponse<U>, EngineError> {
1137 self.plans.push(plan_for(source, Outcomes::default(), 1));
1138 let items = match found {
1139 Ok(Some(item)) => vec![qualify(item)],
1140 Ok(None) => Vec::new(),
1141 Err(error) => {
1142 self.errors.push(SourceFailure {
1143 source: source.name().clone(),
1144 error,
1145 });
1146 Vec::new()
1147 }
1148 };
1149 Ok(QueryResponse {
1150 items,
1151 next: None,
1152 plan: QueryPlan {
1153 per_source: merge_plans(self.plans),
1154 },
1155 errors: self.errors,
1156 })
1157 }
1158
1159 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1161 Ok(QueryResponse {
1162 items: Vec::new(),
1163 next: None,
1164 plan: QueryPlan {
1165 per_source: merge_plans(self.plans),
1166 },
1167 errors: self.errors,
1168 })
1169 }
1170
1171 fn finish<T, U>(
1176 self,
1177 streams: Vec<Stream<T>>,
1178 budget: u32,
1179 first: Option<&Owed>,
1180 query: &str,
1181 qualify: impl Fn(&SourceName, T) -> U,
1182 ) -> Result<QueryResponse<U>, EngineError> {
1183 let (rows, states, owed) = merge(streams, budget, first);
1184 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1185 Ok(QueryResponse {
1186 items: rows
1187 .into_iter()
1188 .map(|(name, item)| qualify(&name, item))
1189 .collect(),
1190 next,
1191 plan: QueryPlan {
1192 per_source: merge_plans(self.plans),
1193 },
1194 errors: self.errors,
1195 })
1196 }
1197}
1198
1199fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1202 SourcePlan {
1203 source: source.name().clone(),
1204 kind: source.kind().to_owned(),
1205 pushed_down: outcomes.with(Outcome::PushedDown),
1206 applied_locally: outcomes.with(Outcome::AppliedLocally),
1207 emulated: outcomes.with(Outcome::Emulated),
1208 unavailable: outcomes.with(Outcome::Unavailable),
1209 pages_fetched: pages,
1210 }
1211}
1212
1213fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1218 let mut merged: Vec<SourcePlan> = Vec::new();
1219 for plan in plans {
1220 if let Some(existing) = merged
1221 .iter_mut()
1222 .find(|existing| existing.source == plan.source)
1223 {
1224 existing.pushed_down.extend(plan.pushed_down);
1225 existing.applied_locally.extend(plan.applied_locally);
1226 existing.emulated.extend(plan.emulated);
1227 existing.unavailable.extend(plan.unavailable);
1228 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1229 for list in [
1230 &mut existing.pushed_down,
1231 &mut existing.applied_locally,
1232 &mut existing.emulated,
1233 &mut existing.unavailable,
1234 ] {
1235 list.sort_unstable();
1236 list.dedup();
1237 }
1238 } else {
1239 merged.push(plan);
1240 }
1241 }
1242 merged
1243}
1244
1245fn walking<'a>(
1251 selected: Vec<&'a ResolvedSource>,
1252 states: &Option<Resumption>,
1253 kind: StreamKind,
1254) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1255 let mut ready = Vec::new();
1256 let mut starts = Vec::new();
1257 for source in selected {
1258 if let Some(resume) = resume_at(states, source.name(), kind) {
1259 ready.push(source);
1260 starts.push(resume);
1261 }
1262 }
1263 (ready, starts)
1264}
1265
1266fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1284 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1285 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1286}
1287
1288fn fingerprint(text: &str) -> String {
1290 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1291 for byte in text.as_bytes() {
1292 hash ^= u64::from(*byte);
1293 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1294 }
1295 format!("{hash:016x}")
1296}
1297
1298fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1303 document.as_ref()?.owed.as_ref()
1304}
1305
1306fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1308 match states {
1309 None => Some(Resume::default()),
1310 Some(document) => document
1311 .streams
1312 .iter()
1313 .find(|state| &state.source == source && state.stream == kind)
1314 .map(|state| state.resume.clone()),
1315 }
1316}
1317
1318fn resumption(
1349 engine: &Engine,
1350 token: Option<&PageToken>,
1351 reads: &[StreamKind],
1352 query: &str,
1353) -> Result<Option<Resumption>, EngineError> {
1354 let Some(document) = token.map(PageToken::decode) else {
1355 return Ok(None);
1356 };
1357
1358 if document.query != query {
1365 return Err(EngineError::Token {
1366 message: "this page token was written by a different query — resume the walk it \
1367 came from, or drop --page to start this one from the beginning"
1368 .to_owned(),
1369 });
1370 }
1371 let states = &document.streams;
1372
1373 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
1374 for state in states {
1375 if !reads.contains(&state.stream) {
1376 return Err(EngineError::Token {
1377 message: format!(
1378 "this page token resumes {}, which this command does not read — it \
1379 was written by a different query",
1380 state.stream.describe()
1381 ),
1382 });
1383 }
1384 let ceiling = engine
1385 .ready()
1386 .find(|source| source.name() == &state.source)
1387 .map(ceiling);
1388 if ceiling.is_none() && !engine.has(&state.source) {
1389 return Err(EngineError::Token {
1390 message: format!(
1391 "this page token resumes a source called {:?}, which this \
1392 configuration does not have",
1393 state.source.as_str()
1394 ),
1395 });
1396 }
1397 if let Some(ceiling) = ceiling
1398 && state.resume.skip >= ceiling
1399 {
1400 return Err(EngineError::Token {
1401 message: format!(
1402 "this page token resumes {} rows into a page of source {:?}, which \
1403 serves at most {ceiling}",
1404 state.resume.skip,
1405 state.source.as_str()
1406 ),
1407 });
1408 }
1409 if seen.contains(&(&state.source, state.stream)) {
1410 return Err(EngineError::Token {
1411 message: format!(
1412 "this page token gives source {:?} two places to resume from",
1413 state.source.as_str()
1414 ),
1415 });
1416 }
1417 seen.push((&state.source, state.stream));
1418 }
1419
1420 if let Some(owed) = &document.owed
1425 && !document
1426 .streams
1427 .iter()
1428 .any(|state| state.source == owed.source && state.stream == owed.stream)
1429 {
1430 return Err(EngineError::Token {
1431 message: format!(
1432 "this page token owes the next row to a stream it does not resume, \
1433 {:?}'s {}",
1434 owed.source.as_str(),
1435 owed.stream.describe()
1436 ),
1437 });
1438 }
1439
1440 Ok(Some(document))
1441}
1442
1443fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
1449 match selector {
1450 ProjectSelector::Any => ProjectFilter::Any,
1451 ProjectSelector::Orphans => ProjectFilter::Orphans,
1452 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
1453 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
1454 }
1455}
1456
1457fn text_predicates(fields: TextFields) -> Vec<Predicate> {
1459 match fields {
1460 TextFields::Title => vec![Predicate::SearchTitle],
1461 TextFields::Content => vec![Predicate::SearchContent],
1462 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
1463 }
1464}
1465
1466fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
1473 match fields {
1474 TextFields::Title => capabilities.search_title.is_native(),
1475 TextFields::Content => capabilities.search_content.is_native(),
1476 TextFields::TitleOrContent => {
1477 capabilities.search_title.is_native() && capabilities.search_content.is_native()
1478 }
1479 }
1480}
1481
1482fn shape_tasks(
1484 capabilities: &Capabilities,
1485 filters: &Filters,
1486 project: &ProjectFilter,
1487) -> TaskShape {
1488 let mut pushed = TaskQuery::default();
1489 let mut local = LocalTasks::default();
1490 let mut outcomes = Outcomes::default();
1491
1492 if !filters.labels.is_empty() {
1493 if capabilities.filter_by_label.is_native() {
1494 pushed.labels = filters.labels.clone();
1495 outcomes.record(Predicate::Label, Outcome::PushedDown);
1496 } else {
1497 local.labels = Some(filters.labels.clone());
1498 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1499 }
1500 }
1501 if !filters.statuses.is_empty() {
1502 if capabilities.filter_by_status.is_native() {
1503 pushed.statuses.clone_from(&filters.statuses);
1504 outcomes.record(Predicate::Status, Outcome::PushedDown);
1505 } else {
1506 local.statuses.clone_from(&filters.statuses);
1507 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1508 }
1509 }
1510 if let Some(text) = &filters.text {
1511 let predicates = text_predicates(text.fields);
1512 if searches_natively(capabilities, text.fields) {
1513 pushed.text = Some(text.clone());
1514 outcomes.record_all(predicates, Outcome::PushedDown);
1515 } else {
1516 local.text = Some(text.clone());
1517 outcomes.record_all(predicates, Outcome::AppliedLocally);
1518 }
1519 }
1520 match project {
1521 ProjectFilter::Any => {}
1522 ProjectFilter::Orphans => {
1523 if capabilities.orphan_tasks.is_native() {
1524 pushed.project = ProjectFilter::Orphans;
1525 outcomes.record(Predicate::Project, Outcome::PushedDown);
1526 } else {
1527 local.project = Some(ProjectFilter::Orphans);
1528 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1529 }
1530 }
1531 ProjectFilter::Is(id) => {
1532 if capabilities.projects.is_native() {
1533 pushed.project = ProjectFilter::Is(id.clone());
1534 outcomes.record(Predicate::Project, Outcome::PushedDown);
1535 } else {
1536 local.project = Some(ProjectFilter::Is(id.clone()));
1537 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
1538 }
1539 }
1540 }
1541
1542 TaskShape {
1543 pushed,
1544 local,
1545 outcomes,
1546 }
1547}
1548
1549fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
1551 let mut pushed = ProjectQuery::default();
1552 let mut local = LocalProjects::default();
1553 let mut outcomes = Outcomes::default();
1554
1555 if !filters.labels.is_empty() {
1556 if capabilities.filter_by_label.is_native() {
1557 pushed.labels = filters.labels.clone();
1558 outcomes.record(Predicate::Label, Outcome::PushedDown);
1559 } else {
1560 local.labels = Some(filters.labels.clone());
1561 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
1562 }
1563 }
1564 if !filters.statuses.is_empty() {
1565 if capabilities.filter_by_status.is_native() {
1566 pushed.statuses.clone_from(&filters.statuses);
1567 outcomes.record(Predicate::Status, Outcome::PushedDown);
1568 } else {
1569 local.statuses.clone_from(&filters.statuses);
1570 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
1571 }
1572 }
1573 if let Some(text) = &filters.text {
1574 let predicates = text_predicates(text.fields);
1575 if searches_natively(capabilities, text.fields) {
1576 pushed.text = Some(text.clone());
1577 outcomes.record_all(predicates, Outcome::PushedDown);
1578 } else {
1579 local.text = Some(text.clone());
1580 outcomes.record_all(predicates, Outcome::AppliedLocally);
1581 }
1582 }
1583
1584 ProjectShape {
1585 pushed,
1586 local,
1587 outcomes,
1588 }
1589}
1590
1591fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
1593 match stream {
1594 StreamKind::Projects => {
1595 let shaped = shape_projects(capabilities, filters);
1596 HitShape {
1597 stream,
1598 tasks: TaskQuery::default(),
1599 projects: shaped.pushed,
1600 local_tasks: LocalTasks::default(),
1601 local_projects: shaped.local,
1602 outcomes: shaped.outcomes,
1603 }
1604 }
1605 StreamKind::Items | StreamKind::Tasks => {
1606 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any);
1607 HitShape {
1608 stream,
1609 tasks: shaped.pushed,
1610 projects: ProjectQuery::default(),
1611 local_tasks: shaped.local,
1612 local_projects: LocalProjects::default(),
1613 outcomes: shaped.outcomes,
1614 }
1615 }
1616 }
1617}
1618
1619fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
1626 if compensating {
1627 ceiling
1628 } else {
1629 budget.min(ceiling)
1630 }
1631}
1632
1633fn ceiling(source: &ResolvedSource) -> u32 {
1635 source.source().capabilities().max_page_size.max(1)
1636}
1637
1638async fn fetch_tasks(
1640 source: &ResolvedSource,
1641 shape: &TaskShape,
1642 start: &Resume,
1643 budget: u32,
1644 calls: &AtomicU32,
1645) -> Result<Fetched<Task>, SourceError> {
1646 let compensating = shape.local != LocalTasks::default();
1647 walk(
1648 start,
1649 budget,
1650 page_size(compensating, budget, ceiling(source)),
1651 |task| shape.local.keeps(task),
1652 |cursor, limit| async move {
1653 calls.fetch_add(1, Ordering::Relaxed);
1654 let request = PageRequest { cursor, limit };
1655 source.source().query_tasks(&shape.pushed, &request).await
1656 },
1657 )
1658 .await
1659}
1660
1661async fn fetch_projects(
1663 source: &ResolvedSource,
1664 shape: &ProjectShape,
1665 start: &Resume,
1666 budget: u32,
1667 calls: &AtomicU32,
1668) -> Result<Fetched<Project>, SourceError> {
1669 let compensating = shape.local != LocalProjects::default();
1670 walk(
1671 start,
1672 budget,
1673 page_size(compensating, budget, ceiling(source)),
1674 |project| shape.local.keeps(project),
1675 |cursor, limit| async move {
1676 calls.fetch_add(1, Ordering::Relaxed);
1677 let request = PageRequest { cursor, limit };
1678 source
1679 .source()
1680 .query_projects(&shape.pushed, &request)
1681 .await
1682 },
1683 )
1684 .await
1685}
1686
1687async fn fetch_labels(
1689 source: &ResolvedSource,
1690 start: &Resume,
1691 budget: u32,
1692 calls: &AtomicU32,
1693) -> Result<Fetched<Label>, SourceError> {
1694 walk(
1695 start,
1696 budget,
1697 page_size(false, budget, ceiling(source)),
1698 |_| true,
1699 |cursor, limit| async move {
1700 calls.fetch_add(1, Ordering::Relaxed);
1701 let request = PageRequest { cursor, limit };
1702 source.source().labels(&request).await
1703 },
1704 )
1705 .await
1706}
1707
1708async fn fetch_hits(
1710 source: &ResolvedSource,
1711 shape: &HitShape,
1712 start: &Resume,
1713 budget: u32,
1714 calls: &AtomicU32,
1715) -> Result<Fetched<Found>, SourceError> {
1716 let ceiling = ceiling(source);
1717 match shape.stream {
1718 StreamKind::Projects => {
1719 let compensating = shape.local_projects != LocalProjects::default();
1720 walk(
1721 start,
1722 budget,
1723 page_size(compensating, budget, ceiling),
1724 |found| match found {
1725 Found::Project(project) => shape.local_projects.keeps(project),
1726 Found::Task(_) => true,
1727 },
1728 |cursor, limit| async move {
1729 calls.fetch_add(1, Ordering::Relaxed);
1730 let request = PageRequest { cursor, limit };
1731 let page = source
1732 .source()
1733 .query_projects(&shape.projects, &request)
1734 .await?;
1735 Ok(Page {
1736 items: page.items.into_iter().map(Found::Project).collect(),
1737 next: page.next,
1738 })
1739 },
1740 )
1741 .await
1742 }
1743 StreamKind::Items | StreamKind::Tasks => {
1744 let compensating = shape.local_tasks != LocalTasks::default();
1745 walk(
1746 start,
1747 budget,
1748 page_size(compensating, budget, ceiling),
1749 |found| match found {
1750 Found::Task(task) => shape.local_tasks.keeps(task),
1751 Found::Project(_) => true,
1752 },
1753 |cursor, limit| async move {
1754 calls.fetch_add(1, Ordering::Relaxed);
1755 let request = PageRequest { cursor, limit };
1756 let page = source.source().query_tasks(&shape.tasks, &request).await?;
1757 Ok(Page {
1758 items: page.items.into_iter().map(Found::Task).collect(),
1759 next: page.next,
1760 })
1761 },
1762 )
1763 .await
1764 }
1765 }
1766}
1767
1768async fn forward_edges(
1770 source: &ResolvedSource,
1771 entity: Entity,
1772 id: &NativeId,
1773 request: &PageRequest,
1774) -> Result<Page<DependencyEdge>, SourceError> {
1775 match entity {
1776 Entity::Task => {
1777 source
1778 .source()
1779 .task_dependencies(id, Direction::DependsOn, request)
1780 .await
1781 }
1782 Entity::Project => {
1783 source
1784 .source()
1785 .project_dependencies(id, Direction::DependsOn, request)
1786 .await
1787 }
1788 }
1789}
1790
1791#[expect(
1800 clippy::too_many_arguments,
1801 reason = "every argument is one axis of one walk — the source, the item, the \
1802 direction, which of its two graphs, whether the reverse is emulated, where \
1803 to resume, how many rows to return and where to count calls. Grouping them \
1804 into a struct would name the same eight values one indirection further from \
1805 the loop that reads them."
1806)]
1807async fn fetch_edges(
1808 source: &ResolvedSource,
1809 native: &NativeId,
1810 direction: Direction,
1811 entity: Entity,
1812 emulating: bool,
1813 start: &Resume,
1814 budget: u32,
1815 calls: &AtomicU32,
1816) -> Result<Fetched<DependencyEdge>, SourceError> {
1817 let ceiling = ceiling(source);
1818 if !emulating {
1819 return walk(
1820 start,
1821 budget,
1822 page_size(false, budget, ceiling),
1823 |_| true,
1824 |cursor, limit| async move {
1825 calls.fetch_add(1, Ordering::Relaxed);
1826 let request = PageRequest { cursor, limit };
1827 match entity {
1828 Entity::Task => {
1829 source
1830 .source()
1831 .task_dependencies(native, direction, &request)
1832 .await
1833 }
1834 Entity::Project => {
1835 source
1836 .source()
1837 .project_dependencies(native, direction, &request)
1838 .await
1839 }
1840 }
1841 },
1842 )
1843 .await;
1844 }
1845
1846 walk(
1847 start,
1848 budget,
1849 ceiling,
1850 |_| true,
1851 |cursor, limit| async move {
1852 calls.fetch_add(1, Ordering::Relaxed);
1853 let request = PageRequest { cursor, limit };
1854 let (ids, next) = match entity {
1855 Entity::Task => {
1856 let page = source
1857 .source()
1858 .query_tasks(&TaskQuery::default(), &request)
1859 .await?;
1860 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
1861 (ids, page.next)
1862 }
1863 Entity::Project => {
1864 let page = source
1865 .source()
1866 .query_projects(&ProjectQuery::default(), &request)
1867 .await?;
1868 let ids: Vec<NativeId> =
1869 page.items.into_iter().map(|project| project.id).collect();
1870 (ids, page.next)
1871 }
1872 };
1873
1874 let mut edges = Vec::new();
1875 for id in ids {
1876 let mut inner: Option<Cursor> = None;
1877 loop {
1878 calls.fetch_add(1, Ordering::Relaxed);
1879 let request = PageRequest {
1880 cursor: inner.clone(),
1881 limit,
1882 };
1883 let page = forward_edges(source, entity, &id, &request).await?;
1884 fits(page.items.len(), limit)?;
1888 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
1889 if page.next.is_some() && page.next == inner {
1890 return Err(SourceError::Malformed {
1891 message: "the source returned the cursor it was given while its \
1892 forward edges were being scanned, so the scan would \
1893 never end"
1894 .to_owned(),
1895 });
1896 }
1897 match page.next {
1898 Some(cursor) => inner = Some(cursor),
1899 None => break,
1900 }
1901 }
1902 }
1903
1904 Ok(Page { items: edges, next })
1905 },
1906 )
1907 .await
1908}