1mod comment;
19mod copy;
20mod delivery;
21mod fetch;
22mod join;
23mod local;
24mod metadata;
25mod narrow;
26mod rendered;
27mod resume;
28mod update;
29
30use std::collections::BTreeMap;
31use std::num::NonZeroU32;
32use std::sync::atomic::{AtomicU32, Ordering};
33
34use chrono::{DateTime, Utc};
35use onetaskgraph_plugin_api::{
36 Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
37 MetadataMatch, MetadataRecord, NativeId, Page, PageRequest, Priority, Project, ProjectFilter,
38 ProjectQuery, Repository, SecretResolver, SourceError, SourceName, StatusCategory, Task,
39 TaskQuery, TextFields, TextQuery,
40};
41use schemars::JsonSchema;
42use serde::{Deserialize, Serialize};
43
44use crate::GlobalId;
45use crate::config::{Config, Placement, Routes};
46use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
47use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
48
49use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
50use join::join_all;
51use local::{LocalDocuments, LocalProjects, LocalTasks};
52pub(crate) use resume::{Owed, Resumption, StreamState};
53use resume::{Resume, StreamKind};
54
55pub use comment::{CommentList, DeletedComment, TaskDetail, TaskDetails};
56pub(crate) use copy::malformed_links;
57pub use copy::{
58 BudgetSpent, CopyAction, CopyItems, CopyLink, CopyLookup, CopyOutcome, CopyReport, CopyRequest,
59 CopyScope, CopyVia, MatchBy, NoCounterpart, Spent,
60};
61pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
62pub use local::ProjectSelector;
63pub use metadata::MetadataSet;
64pub use narrow::{TaskContentSet, TaskPrioritySet};
65pub use rendered::{
66 Body, DocumentCreate, Regenerated, Regeneration, RenderRequest, RenderTemplate, RenderedRecord,
67 TaskCreate, TaskCreated, TemplateAnswers, UnusedAnswers,
68};
69pub use update::TaskUpdated;
70
71#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
76pub struct Qualified<T> {
77 pub id: GlobalId,
79 pub item: T,
81}
82
83#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
88pub struct QualifiedEdge {
89 pub from: QualifiedEndpoint,
94 pub to: QualifiedEndpoint,
96 pub kind: onetaskgraph_plugin_api::DependencyKind,
98}
99
100#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
102pub struct QualifiedEndpoint {
103 pub id: GlobalId,
105 pub kind: onetaskgraph_plugin_api::ItemKind,
107}
108
109impl std::fmt::Display for QualifiedEndpoint {
110 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
111 self.id.fmt(formatter)
112 }
113}
114
115#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
117#[serde(tag = "kind", rename_all = "kebab-case")]
118pub enum SearchHit {
119 Task(Qualified<Task>),
121 Project(Qualified<Project>),
123}
124
125#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
127#[serde(rename_all = "kebab-case")]
128pub enum SearchKind {
129 Tasks,
131 Projects,
133 #[default]
135 Both,
136}
137
138#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
140pub struct SourceListing {
141 pub source: SourceName,
143 pub kind: String,
155 #[serde(flatten)]
157 pub state: SourceState,
158}
159
160#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
162#[serde(tag = "state", rename_all = "kebab-case")]
163pub enum SourceState {
164 Available {
166 capabilities: Capabilities,
168 },
169 Unavailable {
171 error: SourceError,
173 },
174}
175
176#[derive(Debug, Clone, PartialEq)]
178pub struct Paging {
179 pub limit: NonZeroU32,
181 pub token: Option<PageToken>,
183}
184
185#[derive(Debug, Clone, Default, PartialEq)]
187pub struct Filters {
188 pub text: Option<TextQuery>,
190 pub labels: LabelFilter,
192 pub statuses: Vec<StatusCategory>,
194}
195
196#[derive(Debug, Clone)]
198pub struct TaskRequest {
199 pub sources: Vec<SourceName>,
201 pub filters: Filters,
203 pub project: ProjectSelector,
205 pub priorities: Vec<Priority>,
211 pub commented_since: Option<DateTime<Utc>>,
218 pub metadata: Vec<MetadataMatch>,
224 pub origin: Option<GlobalId>,
227 pub include_members: bool,
235 pub paging: Paging,
237}
238
239#[derive(Debug, Clone)]
241pub struct ProjectRequest {
242 pub sources: Vec<SourceName>,
244 pub filters: Filters,
246 pub paging: Paging,
248}
249
250#[derive(Debug, Clone, Default, PartialEq)]
257pub struct DocumentFilters {
258 pub text: Option<TextQuery>,
260 pub labels: LabelFilter,
262}
263
264#[derive(Debug, Clone)]
266pub struct DocumentRequest {
267 pub sources: Vec<SourceName>,
269 pub filters: DocumentFilters,
271 pub project: ProjectSelector,
273 pub paging: Paging,
275}
276
277#[derive(Debug, Clone)]
279pub struct LabelRequest {
280 pub sources: Vec<SourceName>,
282 pub paging: Paging,
284}
285
286#[derive(Debug, Clone)]
288pub struct SearchRequest {
289 pub sources: Vec<SourceName>,
291 pub text: TextQuery,
293 pub kind: SearchKind,
295 pub paging: Paging,
297}
298
299#[derive(Debug, Clone)]
301pub struct DependencyRequest {
302 pub id: GlobalId,
304 pub direction: Direction,
306 pub paging: Paging,
308}
309
310#[derive(Debug, Clone, PartialEq, thiserror::Error)]
318pub enum EngineError {
319 #[error(
321 "no source named {name:?} is configured\n\
322 next: name one of the configured sources ({configured}), or add {name:?} under \
323 `sources` — `onetaskgraph sources list` shows what this configuration has."
324 )]
325 UnknownSource {
326 name: String,
328 configured: String,
330 },
331
332 #[error(
334 "{message}\n\
335 next: page with a token exactly as the previous page reported it, and against \
336 the same configuration — or drop `--page` to start the walk again."
337 )]
338 Token {
339 message: String,
341 },
342
343 #[error(
345 "no sources are configured\n\
346 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
347 prints what each plugin accepts."
348 )]
349 NoSources,
350
351 #[error(
353 "source {name} cannot be written: its plugin is {kind}, which has no write \
354 side\n\
355 next: copy into a source whose plugin can be written — `onetaskgraph sources \
356 list` reports each one's plugin."
357 )]
358 NotWritable {
359 name: String,
361 kind: String,
363 },
364
365 #[error(
372 "source {name} has no documents: its plugin is {kind}, which holds none\n\
373 next: name a source whose plugin has documents — `onetaskgraph sources list` \
374 reports each one's plugin and what it declares."
375 )]
376 NoDocuments {
377 name: String,
379 kind: String,
381 },
382
383 #[error(
388 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
389 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
390 list` reports each one's plugin."
391 )]
392 NoComments {
393 name: String,
395 kind: String,
397 },
398
399 #[error(
401 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
402 but not added to, edited or removed\n\
403 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
404 source whose plugin can be written — `onetaskgraph sources list` reports each \
405 one's plugin."
406 )]
407 CommentsNotWritable {
408 name: String,
410 kind: String,
412 },
413
414 #[error(
416 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
417 next: set the status in that source itself, or name a task of a source whose plugin \
418 can be written — `onetaskgraph sources list` reports each one's plugin."
419 )]
420 StatusNotWritable {
421 name: String,
423 kind: String,
425 },
426
427 #[error(
429 "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
430 write side\n\
431 next: set the key in that source itself, or name a {record} of a source whose plugin \
432 can be written — `onetaskgraph sources list` reports each one's plugin."
433 )]
434 MetadataNotWritable {
435 name: String,
437 kind: String,
439 record: MetadataRecord,
441 },
442
443 #[error(
445 "source {name} cannot write a priority: its plugin is {kind}, which has no write side\n\
446 next: set the priority in that source itself, or name a task of a source whose plugin \
447 can be written — `onetaskgraph sources list` reports each one's plugin."
448 )]
449 PriorityNotWritable {
450 name: String,
452 kind: String,
454 },
455
456 #[error(
458 "source {name} cannot write a task's content: its plugin is {kind}, which has no write \
459 side\n\
460 next: edit the content in that source itself, or name a task of a source whose plugin \
461 can be written — `onetaskgraph sources list` reports each one's plugin."
462 )]
463 ContentNotWritable {
464 name: String,
466 kind: String,
468 },
469
470 #[error(
472 "source {name} cannot update a task: its plugin is {kind}, which has no write side\n\
473 next: change the task in that source itself, or name a task of a source whose plugin \
474 can be written — `onetaskgraph sources list` reports each one's plugin."
475 )]
476 UpdateNotWritable {
477 name: String,
479 kind: String,
481 },
482
483 #[error(
485 "source {name} cannot create a {record}: its plugin is {kind}, which has no write side\n\
486 next: create it in a source whose plugin can be written — `onetaskgraph sources list` \
487 reports each one's plugin."
488 )]
489 NotCreatable {
490 name: String,
492 kind: String,
494 record: MetadataRecord,
496 },
497
498 #[error(
500 "source {name} cannot write a {record}'s rendering: its plugin is {kind}, which has no \
501 write side\n\
502 next: render it with --dry-run to read the result, or regenerate a {record} of a source \
503 whose plugin can be written."
504 )]
505 RenderingNotWritable {
506 name: String,
508 kind: String,
510 record: RenderedRecord,
512 },
513
514 #[error(
517 "supply every required answer to regenerate {id}: {} unanswered, and {reason}\n\
518 next: answer {} with --var NAME=VALUE or an answers file (--answers FILE), or run \
519 interactively to be asked.",
520 names.join(", "),
521 if names.len() == 1 { "it" } else { "each" }
522 )]
523 MissingAnswers {
524 id: String,
526 names: Vec<String>,
528 reason: String,
531 },
532
533 #[error("{error}")]
535 Template {
536 error: crate::template::TemplateError,
538 },
539
540 #[error(
542 "{record} {id} has no stored template answers: {reason}\n\
543 next: regenerate it with every required answer (`onetaskgraph {record} render {id} \
544 --var NAME=VALUE`), or read its provenance with `onetaskgraph {record} show {id}`."
545 )]
546 NoStoredAnswers {
547 record: RenderedRecord,
549 id: String,
551 reason: String,
553 },
554
555 #[error(
557 "{record} {id} records no template it was rendered from, and none was given\n\
558 next: name one with --template FILE or --template-loader FILE."
559 )]
560 NoTemplate {
561 record: RenderedRecord,
563 id: String,
565 },
566
567 #[error(
570 "{record} {id} records a template entry this product did not write — {problem}\n\
571 next: regenerate it with --template FILE or --template-loader FILE and every required \
572 answer, which records a fresh entry."
573 )]
574 MalformedProvenance {
575 record: RenderedRecord,
577 id: String,
579 problem: String,
581 },
582
583 #[error(
585 "{record} {id} was rendered from {reference:?}, which is not a readable file, so it \
586 cannot be re-read; a recorded reference is never turned into a location\n\
587 next: supply the template with --template-loader FILE (a loader document naming what \
588 to render), or name a template file with --template FILE."
589 )]
590 TemplateNotAFile {
591 record: RenderedRecord,
593 id: String,
595 reference: String,
597 },
598
599 #[error(
604 "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
605 be written to it: its plugin is {kind}, which declares priority unsupported\n\
606 next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
607 reports what each declares — or set the task's priority to none first; a \
608 github-projects source holds one once its configuration sets priority_mapping."
609 )]
610 NoPriority {
611 name: String,
613 kind: String,
615 task: String,
617 priority: onetaskgraph_plugin_api::Priority,
619 },
620
621 #[error(
623 "no project with the id {id}\n\
624 next: check the id, or list what is there — `onetaskgraph project list` reports every \
625 project the configured sources hold."
626 )]
627 NoSuchProject {
628 id: String,
630 },
631
632 #[error(
634 "no document with the id {id}\n\
635 next: check the id, or list what is there — `onetaskgraph document list` reports every \
636 document the configured sources hold."
637 )]
638 NoSuchDocument {
639 id: String,
641 },
642
643 #[error(
645 "no task with the id {id}\n\
646 next: check the id, or list what is there — `onetaskgraph task list` reports every \
647 task the configured sources hold."
648 )]
649 NoSuchTask {
650 id: String,
652 },
653
654 #[error(
656 "task {task} has no comment with the id {comment}\n\
657 next: list its comments — `onetaskgraph task comment list {task}` reports each \
658 one's id."
659 )]
660 NoSuchComment {
661 task: String,
663 comment: String,
665 },
666
667 #[error(
672 "source {name} could not be built: {error}\n\
673 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
674 command again."
675 )]
676 SourceUnavailable {
677 name: String,
679 error: SourceError,
681 },
682
683 #[error(
689 "source {name} could not do it: {error}\n\
690 next: fix what the source named above, then run the command again."
691 )]
692 SourceFailed {
693 name: String,
695 error: SourceError,
697 },
698
699 #[error(
701 "the destination source {name} could not be built: {error}\n\
702 next: fix that source — `onetaskgraph sources list` reports its state — then \
703 copy again."
704 )]
705 DestinationUnavailable {
706 name: String,
708 error: SourceError,
710 },
711
712 #[error(
714 "no item with the id {id}\n\
715 next: check the id, or list what is there — `onetaskgraph task list` and \
716 `onetaskgraph project list` report what the configured sources hold."
717 )]
718 NoSuchItem {
719 id: String,
721 },
722
723 #[error(
728 "{item} was copied from {origin}, which that destination no longer holds\n\
729 next: re-run with --recreate to create a new item there instead, or restore \
730 {origin}."
731 )]
732 StaleOrigin {
733 item: String,
735 origin: String,
737 },
738
739 #[error(
745 "{item} was last copied to {link}, which that destination no longer holds\n\
746 next: re-run with --recreate to create a new item there instead, or restore \
747 {link}."
748 )]
749 StaleLink {
750 item: String,
752 link: String,
754 },
755
756 #[error(
759 "--create cannot be given with {flag}: --create asserts the destination holds no \
760 counterpart, so there is nothing for {flag} to look for\n\
761 next: drop {flag} to create each item without looking, or drop --create to look."
762 )]
763 CreateWith {
764 flag: CopyLookup,
766 },
767
768 #[error(
774 "{item} already records a counterpart at the destination, {carrier}, and --create \
775 asserts it has none\n\
776 next: copy it without --create, which updates {carrier}."
777 )]
778 CreateCarried {
779 item: GlobalId,
781 carrier: GlobalId,
783 },
784
785 #[error(
791 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
792 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
793 {id} on its own with `onetaskgraph task copy`.",
794 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
795 )]
796 NotAMember {
797 id: GlobalId,
799 projects: Vec<GlobalId>,
801 },
802
803 #[error(
812 "{item} depends on {member}, which this copy was not told to carry and which records \
813 no origin in {destination}\n\
814 next: name {member} with --member as well, record its {destination} id at \
815 onetaskgraph.origin, or copy the whole project without --member."
816 )]
817 UnrecordedMember {
818 item: GlobalId,
820 member: GlobalId,
822 destination: SourceName,
824 },
825
826 #[error(
833 "{item} was copied to {counterpart}, but its repositories now route it to {route}\n\
834 next: move or remove {counterpart} by hand and copy again, or restore {item}'s \
835 repositories so it routes to {source} again.",
836 source = .counterpart.source
837 )]
838 Misrouted {
839 item: GlobalId,
841 counterpart: GlobalId,
843 route: SourceName,
845 },
846
847 #[error(
853 "source {name} could not do it: {error}\n\
854 next: fix what the source named above, then copy again."
855 )]
856 SourceRefused {
857 name: String,
859 error: SourceError,
861 },
862
863 #[error(
871 "the copy failed and could not be undone.\n\
872 it failed because: {error}\n\
873 it could not be undone because: {refusal}\n\
874 so these still hold what it wrote: {left_behind}\n\
875 next: remove or put back those items, then copy again."
876 )]
877 CopyNotUndone {
878 error: Box<EngineError>,
880 left_behind: LeftBehind,
886 refusal: SourceError,
888 },
889}
890
891#[derive(Debug, Clone, PartialEq)]
900pub struct LeftBehind {
901 first: GlobalId,
903 rest: Vec<GlobalId>,
905}
906
907impl LeftBehind {
908 #[must_use]
910 pub fn new(first: GlobalId) -> Self {
911 Self {
912 first,
913 rest: Vec::new(),
914 }
915 }
916
917 pub fn push(&mut self, id: GlobalId) {
919 self.rest.push(id);
920 }
921
922 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
924 std::iter::once(&self.first).chain(self.rest.iter())
925 }
926}
927
928impl std::fmt::Display for LeftBehind {
929 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
931 write!(formatter, "{}", self.first)?;
932 for id in &self.rest {
933 write!(formatter, ", {id}")?;
934 }
935 Ok(())
936 }
937}
938
939pub enum ConfiguredSource {
946 Ready(ResolvedSource),
948 Unavailable(UnavailableSource),
950}
951
952impl ConfiguredSource {
953 #[must_use]
955 pub fn name(&self) -> &SourceName {
956 match self {
957 Self::Ready(source) => source.name(),
958 Self::Unavailable(source) => source.name(),
959 }
960 }
961}
962
963pub struct Engine {
965 sources: Vec<ConfiguredSource>,
967 selection: Vec<SourceName>,
969 routes: Routes,
972}
973
974impl Engine {
975 #[must_use]
983 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
984 let (ready, unavailable) = resolve_available(config, secrets);
985 Self::new(
986 ready
987 .into_iter()
988 .map(ConfiguredSource::Ready)
989 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
990 .collect(),
991 config.selected_sources(),
992 )
993 .with_routes(config.routes())
994 }
995
996 #[must_use]
999 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
1000 Self {
1001 sources,
1002 selection,
1003 routes: Routes::default(),
1004 }
1005 }
1006
1007 #[must_use]
1012 pub fn with_routes(mut self, routes: Routes) -> Self {
1013 self.routes = routes;
1014 self
1015 }
1016
1017 #[must_use]
1019 pub fn place(&self, source: &SourceName, repositories: &[Repository]) -> Placement {
1020 self.routes.place(source, repositories)
1021 }
1022
1023 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
1025 self.sources.iter().filter_map(|source| match source {
1026 ConfiguredSource::Ready(ready) => Some(ready),
1027 ConfiguredSource::Unavailable(_) => None,
1028 })
1029 }
1030
1031 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
1033 self.sources.iter().filter_map(|source| match source {
1034 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
1035 ConfiguredSource::Ready(_) => None,
1036 })
1037 }
1038
1039 #[must_use]
1041 pub fn listing(&self) -> Vec<SourceListing> {
1042 let mut listings: Vec<SourceListing> = self
1043 .ready()
1044 .map(|source| SourceListing {
1045 source: source.name().clone(),
1046 kind: source.kind().to_owned(),
1047 state: SourceState::Available {
1048 capabilities: source.source().capabilities(),
1049 },
1050 })
1051 .chain(self.unavailable().map(|source| SourceListing {
1052 source: source.name().clone(),
1053 kind: source.kind().to_owned(),
1054 state: SourceState::Unavailable {
1055 error: source.error().clone(),
1056 },
1057 }))
1058 .collect();
1059 listings.sort_by(|left, right| left.source.cmp(&right.source));
1060 listings
1061 }
1062
1063 #[must_use]
1069 pub fn has(&self, name: &SourceName) -> bool {
1070 self.sources.iter().any(|source| source.name() == name)
1071 }
1072
1073 pub async fn tasks(
1081 &self,
1082 request: &TaskRequest,
1083 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1084 let mut names = self.resolve_selection(&request.sources)?;
1085 let mut members: Vec<GlobalId> = Vec::new();
1090 let mut unread: Vec<SourceFailure> = Vec::new();
1091 if let ProjectSelector::Qualified(id) = &request.project {
1092 self.known(&id.source)?;
1093 names.retain(|name| name == &id.source);
1094 if request.include_members {
1097 (members, unread) = self.members(id).await;
1098 names.extend(members.iter().map(|member| member.source.clone()));
1099 }
1100 }
1101 let selector = |name: &SourceName| -> ProjectSelector {
1102 members
1103 .iter()
1104 .find(|member| &member.source == name)
1105 .map_or_else(
1106 || request.project.clone(),
1107 |member| ProjectSelector::Qualified(member.clone()),
1108 )
1109 };
1110 let query = shape(
1111 "task-list",
1112 &names,
1113 &(
1114 &request.filters,
1115 &request.project,
1116 &request.priorities,
1117 &request.commented_since,
1118 &request.metadata,
1119 &request.origin,
1120 &members,
1121 ),
1122 );
1123 let states = resumption(
1124 self,
1125 request.paging.token.as_ref(),
1126 &[StreamKind::Items],
1127 &query,
1128 )?;
1129 let budget = request.paging.limit.get();
1130
1131 let mut answer = Answer::new();
1132 answer.errors.extend(unread);
1133 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1134
1135 let shapes: Vec<TaskShape> = ready
1136 .iter()
1137 .map(|source| {
1138 shape_tasks(
1139 &source.source().capabilities(),
1140 &request.filters,
1141 &project_filter(&selector(source.name())),
1142 &request.priorities,
1143 request.commented_since,
1144 &request.metadata,
1145 request.origin.as_ref(),
1146 )
1147 })
1148 .collect();
1149 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1150 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1151
1152 let walks = ready
1153 .iter()
1154 .enumerate()
1155 .map(|(index, source)| {
1156 fetch_tasks(
1157 source,
1158 &shapes[index],
1159 &starts[index],
1160 budget,
1161 &counters[index],
1162 )
1163 })
1164 .collect();
1165
1166 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1167 answer.finish(
1168 streams,
1169 budget,
1170 owed(&states),
1171 &query,
1172 |name, task: Task| {
1173 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1174 },
1175 )
1176 }
1177
1178 async fn members(&self, home: &GlobalId) -> (Vec<GlobalId>, Vec<SourceFailure>) {
1189 let mut members = Vec::new();
1190 let mut unread = Vec::new();
1191 let Some(source) = self.ready().find(|source| source.name() == &home.source) else {
1192 return (members, unread);
1193 };
1194 let held = match source.source().get_project(&home.native).await {
1195 Ok(Some(held)) => held,
1196 Ok(None) => {
1197 unread.push(SourceFailure {
1198 source: home.source.clone(),
1199 error: SourceError::Refused {
1200 message: format!(
1201 "{home} names no project {} holds, so its member projects cannot \
1202 be read",
1203 home.source
1204 ),
1205 },
1206 });
1207 return (members, unread);
1208 }
1209 Err(error) => {
1210 unread.push(SourceFailure {
1211 source: home.source.clone(),
1212 error,
1213 });
1214 return (members, unread);
1215 }
1216 };
1217 let named = match copy::members_of(home, &held.metadata) {
1218 Ok(named) => named,
1219 Err(message) => {
1222 unread.push(SourceFailure {
1223 source: home.source.clone(),
1224 error: SourceError::Malformed { message },
1225 });
1226 return (members, unread);
1227 }
1228 };
1229 for member in named {
1230 if !self.has(&member.source) {
1231 unread.push(SourceFailure {
1232 source: member.source.clone(),
1233 error: SourceError::Config {
1234 message: format!(
1235 "{home} names {member} as a member project, but no source named \
1236 {} is configured",
1237 member.source
1238 ),
1239 },
1240 });
1241 continue;
1242 }
1243 let Some(there) = self.ready().find(|source| source.name() == &member.source) else {
1246 members.push(member);
1247 continue;
1248 };
1249 match there.source().get_project(&member.native).await {
1250 Ok(Some(_)) => members.push(member),
1251 Ok(None) => unread.push(SourceFailure {
1252 source: member.source.clone(),
1253 error: SourceError::Refused {
1254 message: format!(
1255 "{home} names {member} as a member project, and {} holds no such \
1256 project",
1257 member.source
1258 ),
1259 },
1260 }),
1261 Err(error) => unread.push(SourceFailure {
1262 source: member.source.clone(),
1263 error,
1264 }),
1265 }
1266 }
1267 (members, unread)
1268 }
1269
1270 pub async fn projects(
1276 &self,
1277 request: &ProjectRequest,
1278 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1279 let names = self.resolve_selection(&request.sources)?;
1280 let query = shape("project-list", &names, &request.filters);
1281 let states = resumption(
1282 self,
1283 request.paging.token.as_ref(),
1284 &[StreamKind::Items],
1285 &query,
1286 )?;
1287 let budget = request.paging.limit.get();
1288
1289 let mut answer = Answer::new();
1290 let mut with_projects = Vec::new();
1296 for source in answer.split(self, &names) {
1297 if source.source().capabilities().projects.is_native() {
1298 with_projects.push(source);
1299 } else {
1300 answer.unreachable_predicate(source, Predicate::Project);
1301 }
1302 }
1303 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1304
1305 let shapes: Vec<ProjectShape> = ready
1306 .iter()
1307 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1308 .collect();
1309 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1310 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1311
1312 let walks = ready
1313 .iter()
1314 .enumerate()
1315 .map(|(index, source)| {
1316 fetch_projects(
1317 source,
1318 &shapes[index],
1319 &starts[index],
1320 budget,
1321 &counters[index],
1322 )
1323 })
1324 .collect();
1325
1326 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1327 answer.finish(
1328 streams,
1329 budget,
1330 owed(&states),
1331 &query,
1332 |name, project: Project| Qualified {
1333 id: GlobalId::new(name.clone(), project.id.clone()),
1334 item: project,
1335 },
1336 )
1337 }
1338
1339 pub async fn documents(
1351 &self,
1352 request: &DocumentRequest,
1353 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1354 let mut names = self.resolve_selection(&request.sources)?;
1355 if let ProjectSelector::Qualified(id) = &request.project {
1358 self.known(&id.source)?;
1359 names.retain(|name| name == &id.source);
1360 }
1361 let query = shape(
1362 "document-list",
1363 &names,
1364 &(&request.filters, &request.project),
1365 );
1366 let states = resumption(
1367 self,
1368 request.paging.token.as_ref(),
1369 &[StreamKind::Items],
1370 &query,
1371 )?;
1372 let budget = request.paging.limit.get();
1373
1374 let mut answer = Answer::new();
1375 let mut with_documents = Vec::new();
1376 for source in answer.split(self, &names) {
1377 if source.source().capabilities().documents.is_native() {
1378 with_documents.push(source);
1379 } else {
1380 answer.unreachable_predicate(source, Predicate::Document);
1381 }
1382 }
1383 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1384
1385 let shapes: Vec<DocumentShape> = ready
1386 .iter()
1387 .map(|source| {
1388 shape_documents(
1389 &source.source().capabilities(),
1390 &request.filters,
1391 &project_filter(&request.project),
1392 )
1393 })
1394 .collect();
1395 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1396 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1397
1398 let walks = ready
1399 .iter()
1400 .enumerate()
1401 .map(|(index, source)| {
1402 fetch_documents(
1403 source,
1404 &shapes[index],
1405 &starts[index],
1406 budget,
1407 &counters[index],
1408 )
1409 })
1410 .collect();
1411
1412 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1413 answer.finish(
1414 streams,
1415 budget,
1416 owed(&states),
1417 &query,
1418 |name, document: Document| Qualified {
1419 id: GlobalId::new(name.clone(), document.id.clone()),
1420 item: document,
1421 },
1422 )
1423 }
1424
1425 pub async fn labels(
1431 &self,
1432 request: &LabelRequest,
1433 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1434 let names = self.resolve_selection(&request.sources)?;
1435 let query = shape("label-list", &names, &());
1436 let states = resumption(
1437 self,
1438 request.paging.token.as_ref(),
1439 &[StreamKind::Items],
1440 &query,
1441 )?;
1442 let budget = request.paging.limit.get();
1443
1444 let mut answer = Answer::new();
1445 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1446 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1447 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1448
1449 let walks = ready
1450 .iter()
1451 .enumerate()
1452 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1453 .collect();
1454
1455 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1456 answer.finish(
1457 streams,
1458 budget,
1459 owed(&states),
1460 &query,
1461 |name, label: Label| Qualified {
1462 id: GlobalId::new(name.clone(), label.id.clone()),
1463 item: label,
1464 },
1465 )
1466 }
1467
1468 pub async fn search(
1474 &self,
1475 request: &SearchRequest,
1476 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1477 let names = self.resolve_selection(&request.sources)?;
1478 let reads: &[StreamKind] = match request.kind {
1482 SearchKind::Tasks => &[StreamKind::Tasks],
1483 SearchKind::Projects => &[StreamKind::Projects],
1484 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1485 };
1486 let query = shape("search", &names, &request.text);
1491 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1492 let budget = request.paging.limit.get();
1493 let filters = Filters {
1494 text: Some(request.text.clone()),
1495 ..Filters::default()
1496 };
1497
1498 let mut answer = Answer::new();
1499
1500 let mut ready = Vec::new();
1503 let mut kinds = Vec::new();
1504 let mut starts = Vec::new();
1505 for source in answer.split(self, &names) {
1506 let mut streams = Vec::new();
1507 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1508 streams.push(StreamKind::Tasks);
1509 }
1510 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1511 if source.source().capabilities().projects.is_native() {
1512 streams.push(StreamKind::Projects);
1513 } else {
1514 answer.unreachable_predicate(source, Predicate::Project);
1515 }
1516 }
1517 for stream in streams {
1518 if let Some(resume) = resume_at(&states, source.name(), stream) {
1519 ready.push(source);
1520 kinds.push(stream);
1521 starts.push(resume);
1522 }
1523 }
1524 }
1525
1526 let shapes: Vec<HitShape> = ready
1527 .iter()
1528 .zip(kinds.iter())
1529 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1530 .collect();
1531 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1532 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1533
1534 let walks = ready
1535 .iter()
1536 .enumerate()
1537 .map(|(index, source)| {
1538 fetch_hits(
1539 source,
1540 &shapes[index],
1541 &starts[index],
1542 budget,
1543 &counters[index],
1544 )
1545 })
1546 .collect();
1547
1548 let streams =
1549 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1550 answer.finish(
1551 streams,
1552 budget,
1553 owed(&states),
1554 &query,
1555 |name, found: Found| match found {
1556 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1557 GlobalId::new(name.clone(), task.id.clone()),
1558 task,
1559 )),
1560 Found::Project(project) => SearchHit::Project(Qualified {
1561 id: GlobalId::new(name.clone(), project.id.clone()),
1562 item: project,
1563 }),
1564 },
1565 )
1566 }
1567
1568 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1575 let name = self.known(&id.source)?;
1576 let mut answer = Answer::new();
1577 let selected = answer.split(self, std::slice::from_ref(&name));
1578 let Some(source) = selected.first() else {
1579 return answer.nothing();
1580 };
1581 let found = source.source().get_task(&id.native).await;
1582 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1583 answer.one(source, found, |task| {
1584 delivery::qualified_task(qualified, task)
1585 })
1586 }
1587
1588 pub async fn project(
1594 &self,
1595 id: &GlobalId,
1596 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1597 let name = self.known(&id.source)?;
1598 let mut answer = Answer::new();
1599 let selected = answer.split(self, std::slice::from_ref(&name));
1600 let Some(source) = selected.first() else {
1601 return answer.nothing();
1602 };
1603 let found = source.source().get_project(&id.native).await;
1604 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1605 answer.one(source, found, |project| Qualified {
1606 id: qualified,
1607 item: project,
1608 })
1609 }
1610
1611 pub async fn document(
1621 &self,
1622 id: &GlobalId,
1623 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1624 let name = self.known(&id.source)?;
1625 let mut answer = Answer::new();
1626 let selected = answer.split(self, std::slice::from_ref(&name));
1627 let Some(source) = selected.first() else {
1628 return answer.nothing();
1629 };
1630 if !source.source().capabilities().documents.is_native() {
1631 answer.unreachable_predicate(source, Predicate::Document);
1632 return answer.nothing();
1633 }
1634 let found = source.source().get_document(&id.native).await;
1635 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1636 answer.one(source, found, |document| Qualified {
1637 id: qualified,
1638 item: document,
1639 })
1640 }
1641
1642 pub async fn task_dependencies(
1649 &self,
1650 request: &DependencyRequest,
1651 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1652 self.dependencies(request, Entity::Task).await
1653 }
1654
1655 pub async fn project_dependencies(
1661 &self,
1662 request: &DependencyRequest,
1663 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1664 self.dependencies(request, Entity::Project).await
1665 }
1666
1667 async fn dependencies(
1670 &self,
1671 request: &DependencyRequest,
1672 entity: Entity,
1673 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1674 let name = self.known(&request.id.source)?;
1675 let query = shape(
1676 "dependencies",
1677 std::slice::from_ref(&name),
1678 &(entity, &request.id.native, request.direction),
1679 );
1680 let states = resumption(
1681 self,
1682 request.paging.token.as_ref(),
1683 &[StreamKind::Items],
1684 &query,
1685 )?;
1686 let budget = request.paging.limit.get();
1687
1688 let mut answer = Answer::new();
1689 let (ready, starts) = walking(
1690 answer.split(self, std::slice::from_ref(&name)),
1691 &states,
1692 StreamKind::Items,
1693 );
1694 let Some(source) = ready.first() else {
1695 return answer.nothing();
1696 };
1697
1698 let capabilities = source.source().capabilities();
1699 let support = match entity {
1700 Entity::Task => capabilities.task_dependencies,
1701 Entity::Project => capabilities.project_dependencies,
1702 };
1703 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1707 let mut outcomes = Outcomes::default();
1708 if request.direction == Direction::DependedOnBy {
1709 if emulating {
1710 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1711 } else {
1712 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1713 }
1714 }
1715
1716 let counters = vec![AtomicU32::new(0)];
1717 let walked = fetch_edges(
1718 source,
1719 &request.id.native,
1720 request.direction,
1721 entity,
1722 emulating,
1723 &starts[0],
1724 budget,
1725 &counters[0],
1726 )
1727 .await;
1728
1729 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1730 answer.finish(
1731 streams,
1732 budget,
1733 owed(&states),
1734 &query,
1735 |name, edge: DependencyEdge| QualifiedEdge {
1736 from: qualify_endpoint(name, edge.from),
1737 to: qualify_endpoint(name, edge.to),
1738 kind: edge.kind,
1739 },
1740 )
1741 }
1742
1743 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1745 if asked.is_empty() {
1746 if self.selection.is_empty() {
1747 return Err(EngineError::NoSources);
1748 }
1749 return Ok(self.selection.clone());
1750 }
1751 asked.iter().map(|name| self.known(name)).collect()
1752 }
1753
1754 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1756 if self.has(name) {
1757 return Ok(name.clone());
1758 }
1759 if self.sources.is_empty() {
1760 return Err(EngineError::NoSources);
1761 }
1762 Err(EngineError::UnknownSource {
1763 name: name.to_string(),
1764 configured: self
1765 .listing()
1766 .iter()
1767 .map(|listing| listing.source.to_string())
1768 .collect::<Vec<_>>()
1769 .join(", "),
1770 })
1771 }
1772}
1773
1774fn qualify_endpoint(
1775 source: &SourceName,
1776 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1777) -> QualifiedEndpoint {
1778 let kind = endpoint.kind;
1779 let is_qualified = endpoint.is_qualified();
1780 let endpoint_id = endpoint.into_id();
1781 QualifiedEndpoint {
1782 id: if is_qualified {
1783 endpoint_id
1784 .parse()
1785 .expect("plugin-api validates qualified dependency endpoints")
1786 } else {
1787 GlobalId::new(source.clone(), NativeId(endpoint_id))
1788 },
1789 kind,
1790 }
1791}
1792
1793#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1795enum Entity {
1796 Task,
1798 Project,
1800}
1801
1802enum Found {
1804 Task(Task),
1806 Project(Project),
1808}
1809
1810#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1812enum Outcome {
1813 PushedDown,
1815 AppliedLocally,
1817 Emulated,
1819 Unavailable,
1821}
1822
1823#[derive(Debug, Clone, Default, PartialEq)]
1835struct Outcomes(BTreeMap<Predicate, Outcome>);
1836
1837impl Outcomes {
1838 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1844 self.0.insert(predicate, outcome);
1845 }
1846
1847 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1850 for predicate in predicates {
1851 self.record(predicate, outcome);
1852 }
1853 }
1854
1855 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1857 self.0
1858 .iter()
1859 .filter(|(_, recorded)| **recorded == outcome)
1860 .map(|(predicate, _)| *predicate)
1861 .collect()
1862 }
1863}
1864
1865struct TaskShape {
1867 pushed: TaskQuery,
1869 local: LocalTasks,
1871 outcomes: Outcomes,
1873}
1874
1875struct ProjectShape {
1877 pushed: ProjectQuery,
1879 local: LocalProjects,
1881 outcomes: Outcomes,
1883}
1884
1885struct DocumentShape {
1887 pushed: DocumentQuery,
1889 local: LocalDocuments,
1891 outcomes: Outcomes,
1893}
1894
1895struct HitShape {
1897 stream: StreamKind,
1899 tasks: TaskQuery,
1901 projects: ProjectQuery,
1903 local_tasks: LocalTasks,
1905 local_projects: LocalProjects,
1907 outcomes: Outcomes,
1909}
1910
1911struct Answer {
1917 plans: Vec<SourcePlan>,
1919 errors: Vec<SourceFailure>,
1921}
1922
1923impl Answer {
1924 fn new() -> Self {
1925 Self {
1926 plans: Vec::new(),
1927 errors: Vec::new(),
1928 }
1929 }
1930
1931 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1936 let mut selected = Vec::new();
1937 for name in names {
1938 match engine.sources.iter().find(|source| source.name() == name) {
1939 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1940 Some(ConfiguredSource::Unavailable(source)) => {
1941 self.errors.push(source.failure());
1942 }
1943 None => {}
1944 }
1945 }
1946 selected
1947 }
1948
1949 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1951 let mut outcomes = Outcomes::default();
1952 outcomes.record(predicate, Outcome::Unavailable);
1953 self.plans.push(plan_for(source, outcomes, 0));
1954 }
1955
1956 fn collect<T>(
1958 &mut self,
1959 ready: &[&ResolvedSource],
1960 walked: Vec<Result<Fetched<T>, SourceError>>,
1961 counters: &[AtomicU32],
1962 outcomes: Vec<Outcomes>,
1963 ) -> Vec<Stream<T>> {
1964 let kinds = vec![StreamKind::Items; ready.len()];
1965 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1966 }
1967
1968 fn collect_streams<T>(
1970 &mut self,
1971 ready: &[&ResolvedSource],
1972 kinds: &[StreamKind],
1973 walked: Vec<Result<Fetched<T>, SourceError>>,
1974 counters: &[AtomicU32],
1975 outcomes: Vec<Outcomes>,
1976 ) -> Vec<Stream<T>> {
1977 let mut streams = Vec::new();
1978 for (index, result) in walked.into_iter().enumerate() {
1979 let source = ready[index];
1980 let pages = counters[index].load(Ordering::Relaxed);
1981 self.plans
1982 .push(plan_for(source, outcomes[index].clone(), pages));
1983 match result {
1984 Ok(fetched) => streams.push(Stream {
1985 source: source.name().clone(),
1986 kind: kinds[index],
1987 fetched,
1988 }),
1989 Err(error) => self.errors.push(SourceFailure {
1992 source: source.name().clone(),
1993 error,
1994 }),
1995 }
1996 }
1997 streams
1998 }
1999
2000 fn one<T, U>(
2002 self,
2003 source: &ResolvedSource,
2004 found: Result<Option<T>, SourceError>,
2005 qualify: impl FnOnce(T) -> U,
2006 ) -> Result<QueryResponse<U>, EngineError> {
2007 Ok(self.one_response(source, found, qualify))
2008 }
2009
2010 fn one_response<T, U>(
2012 mut self,
2013 source: &ResolvedSource,
2014 found: Result<Option<T>, SourceError>,
2015 qualify: impl FnOnce(T) -> U,
2016 ) -> QueryResponse<U> {
2017 self.plans.push(plan_for(source, Outcomes::default(), 1));
2018 let items = match found {
2019 Ok(Some(item)) => vec![qualify(item)],
2020 Ok(None) => Vec::new(),
2021 Err(error) => {
2022 self.errors.push(SourceFailure {
2023 source: source.name().clone(),
2024 error,
2025 });
2026 Vec::new()
2027 }
2028 };
2029 QueryResponse {
2030 items,
2031 next: None,
2032 plan: QueryPlan {
2033 per_source: merge_plans(self.plans),
2034 },
2035 errors: self.errors,
2036 }
2037 }
2038
2039 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
2041 Ok(QueryResponse {
2042 items: Vec::new(),
2043 next: None,
2044 plan: QueryPlan {
2045 per_source: merge_plans(self.plans),
2046 },
2047 errors: self.errors,
2048 })
2049 }
2050
2051 fn finish<T, U>(
2056 self,
2057 streams: Vec<Stream<T>>,
2058 budget: u32,
2059 first: Option<&Owed>,
2060 query: &str,
2061 qualify: impl Fn(&SourceName, T) -> U,
2062 ) -> Result<QueryResponse<U>, EngineError> {
2063 let (rows, states, owed) = merge(streams, budget, first);
2064 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
2065 Ok(QueryResponse {
2066 items: rows
2067 .into_iter()
2068 .map(|(name, item)| qualify(&name, item))
2069 .collect(),
2070 next,
2071 plan: QueryPlan {
2072 per_source: merge_plans(self.plans),
2073 },
2074 errors: self.errors,
2075 })
2076 }
2077}
2078
2079fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
2082 SourcePlan {
2083 source: source.name().clone(),
2084 kind: source.kind().to_owned(),
2085 pushed_down: outcomes.with(Outcome::PushedDown),
2086 applied_locally: outcomes.with(Outcome::AppliedLocally),
2087 emulated: outcomes.with(Outcome::Emulated),
2088 unavailable: outcomes.with(Outcome::Unavailable),
2089 pages_fetched: pages,
2090 }
2091}
2092
2093fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
2098 let mut merged: Vec<SourcePlan> = Vec::new();
2099 for plan in plans {
2100 if let Some(existing) = merged
2101 .iter_mut()
2102 .find(|existing| existing.source == plan.source)
2103 {
2104 existing.pushed_down.extend(plan.pushed_down);
2105 existing.applied_locally.extend(plan.applied_locally);
2106 existing.emulated.extend(plan.emulated);
2107 existing.unavailable.extend(plan.unavailable);
2108 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
2109 for list in [
2110 &mut existing.pushed_down,
2111 &mut existing.applied_locally,
2112 &mut existing.emulated,
2113 &mut existing.unavailable,
2114 ] {
2115 list.sort_unstable();
2116 list.dedup();
2117 }
2118 } else {
2119 merged.push(plan);
2120 }
2121 }
2122 merged
2123}
2124
2125fn walking<'a>(
2131 selected: Vec<&'a ResolvedSource>,
2132 states: &Option<Resumption>,
2133 kind: StreamKind,
2134) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
2135 let mut ready = Vec::new();
2136 let mut starts = Vec::new();
2137 for source in selected {
2138 if let Some(resume) = resume_at(states, source.name(), kind) {
2139 ready.push(source);
2140 starts.push(resume);
2141 }
2142 }
2143 (ready, starts)
2144}
2145
2146fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
2164 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
2165 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
2166}
2167
2168fn fingerprint(text: &str) -> String {
2170 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
2171 for byte in text.as_bytes() {
2172 hash ^= u64::from(*byte);
2173 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
2174 }
2175 format!("{hash:016x}")
2176}
2177
2178fn owed(document: &Option<Resumption>) -> Option<&Owed> {
2183 document.as_ref()?.owed.as_ref()
2184}
2185
2186fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
2188 match states {
2189 None => Some(Resume::default()),
2190 Some(document) => document
2191 .streams
2192 .iter()
2193 .find(|state| &state.source == source && state.stream == kind)
2194 .map(|state| state.resume.clone()),
2195 }
2196}
2197
2198fn resumption(
2229 engine: &Engine,
2230 token: Option<&PageToken>,
2231 reads: &[StreamKind],
2232 query: &str,
2233) -> Result<Option<Resumption>, EngineError> {
2234 let Some(document) = token.map(PageToken::decode) else {
2235 return Ok(None);
2236 };
2237
2238 if document.query != query {
2245 return Err(EngineError::Token {
2246 message: "this page token was written by a different query — resume the walk it \
2247 came from, or drop --page to start this one from the beginning"
2248 .to_owned(),
2249 });
2250 }
2251 let states = &document.streams;
2252
2253 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2254 for state in states {
2255 if !reads.contains(&state.stream) {
2256 return Err(EngineError::Token {
2257 message: format!(
2258 "this page token resumes {}, which this command does not read — it \
2259 was written by a different query",
2260 state.stream.describe()
2261 ),
2262 });
2263 }
2264 let ceiling = engine
2265 .ready()
2266 .find(|source| source.name() == &state.source)
2267 .map(ceiling);
2268 if ceiling.is_none() && !engine.has(&state.source) {
2269 return Err(EngineError::Token {
2270 message: format!(
2271 "this page token resumes a source called {:?}, which this \
2272 configuration does not have",
2273 state.source.as_str()
2274 ),
2275 });
2276 }
2277 if let Some(ceiling) = ceiling
2278 && state.resume.skip >= ceiling
2279 {
2280 return Err(EngineError::Token {
2281 message: format!(
2282 "this page token resumes {} rows into a page of source {:?}, which \
2283 serves at most {ceiling}",
2284 state.resume.skip,
2285 state.source.as_str()
2286 ),
2287 });
2288 }
2289 if seen.contains(&(&state.source, state.stream)) {
2290 return Err(EngineError::Token {
2291 message: format!(
2292 "this page token gives source {:?} two places to resume from",
2293 state.source.as_str()
2294 ),
2295 });
2296 }
2297 seen.push((&state.source, state.stream));
2298 }
2299
2300 if let Some(owed) = &document.owed
2305 && !document
2306 .streams
2307 .iter()
2308 .any(|state| state.source == owed.source && state.stream == owed.stream)
2309 {
2310 return Err(EngineError::Token {
2311 message: format!(
2312 "this page token owes the next row to a stream it does not resume, \
2313 {:?}'s {}",
2314 owed.source.as_str(),
2315 owed.stream.describe()
2316 ),
2317 });
2318 }
2319
2320 Ok(Some(document))
2321}
2322
2323fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2329 match selector {
2330 ProjectSelector::Any => ProjectFilter::Any,
2331 ProjectSelector::Orphans => ProjectFilter::Orphans,
2332 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2333 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2334 }
2335}
2336
2337fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2339 match fields {
2340 TextFields::Title => vec![Predicate::SearchTitle],
2341 TextFields::Content => vec![Predicate::SearchContent],
2342 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2343 }
2344}
2345
2346fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2353 match fields {
2354 TextFields::Title => capabilities.search_title.is_native(),
2355 TextFields::Content => capabilities.search_content.is_native(),
2356 TextFields::TitleOrContent => {
2357 capabilities.search_title.is_native() && capabilities.search_content.is_native()
2358 }
2359 }
2360}
2361
2362fn shape_tasks(
2364 capabilities: &Capabilities,
2365 filters: &Filters,
2366 project: &ProjectFilter,
2367 priorities: &[Priority],
2368 commented_since: Option<DateTime<Utc>>,
2369 metadata: &[MetadataMatch],
2370 origin: Option<&GlobalId>,
2371) -> TaskShape {
2372 let mut pushed = TaskQuery::default();
2373 let mut local = LocalTasks::default();
2374 let mut outcomes = Outcomes::default();
2375
2376 if !filters.labels.is_empty() {
2377 if capabilities.filter_by_label.is_native() {
2378 pushed.labels = filters.labels.clone();
2379 outcomes.record(Predicate::Label, Outcome::PushedDown);
2380 } else {
2381 local.labels = Some(filters.labels.clone());
2382 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2383 }
2384 }
2385 if !filters.statuses.is_empty() {
2386 if capabilities.filter_by_status.is_native() {
2387 pushed.statuses.clone_from(&filters.statuses);
2388 outcomes.record(Predicate::Status, Outcome::PushedDown);
2389 } else {
2390 local.statuses.clone_from(&filters.statuses);
2391 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2392 }
2393 }
2394 if !priorities.is_empty() {
2395 if capabilities.filter_by_priority.is_native() {
2396 pushed.priorities = priorities.to_vec();
2397 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2398 } else {
2399 local.priorities = priorities.to_vec();
2400 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2401 }
2402 }
2403 if let Some(since) = commented_since {
2404 if capabilities.filter_by_comment_activity.is_native() {
2405 pushed.commented_since = Some(since);
2406 outcomes.record(Predicate::CommentedSince, Outcome::PushedDown);
2407 } else {
2408 local.commented_since = Some(since);
2409 outcomes.record(Predicate::CommentedSince, Outcome::AppliedLocally);
2410 }
2411 }
2412 if !metadata.is_empty() {
2413 if capabilities.filter_by_metadata.is_native() {
2414 pushed.metadata = metadata.to_vec();
2415 outcomes.record(Predicate::Metadata, Outcome::PushedDown);
2416 } else {
2417 local.metadata = metadata.to_vec();
2418 outcomes.record(Predicate::Metadata, Outcome::AppliedLocally);
2419 }
2420 }
2421 if let Some(origin) = origin {
2422 if capabilities.filter_by_origin.is_native() {
2423 pushed.origin = Some(origin.to_string());
2426 outcomes.record(Predicate::Origin, Outcome::PushedDown);
2427 } else {
2428 local.origin = Some(origin.clone());
2429 outcomes.record(Predicate::Origin, Outcome::AppliedLocally);
2430 }
2431 }
2432 if let Some(text) = &filters.text {
2433 let predicates = text_predicates(text.fields);
2434 if searches_natively(capabilities, text.fields) {
2435 pushed.text = Some(text.clone());
2436 outcomes.record_all(predicates, Outcome::PushedDown);
2437 } else {
2438 local.text = Some(text.clone());
2439 outcomes.record_all(predicates, Outcome::AppliedLocally);
2440 }
2441 }
2442 match project {
2443 ProjectFilter::Any => {}
2444 ProjectFilter::Orphans => {
2445 if capabilities.orphan_tasks.is_native() {
2446 pushed.project = ProjectFilter::Orphans;
2447 outcomes.record(Predicate::Project, Outcome::PushedDown);
2448 } else {
2449 local.project = Some(ProjectFilter::Orphans);
2450 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2451 }
2452 }
2453 ProjectFilter::Is(id) => {
2454 if capabilities.projects.is_native() {
2455 pushed.project = ProjectFilter::Is(id.clone());
2456 outcomes.record(Predicate::Project, Outcome::PushedDown);
2457 } else {
2458 local.project = Some(ProjectFilter::Is(id.clone()));
2459 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2460 }
2461 }
2462 }
2463
2464 TaskShape {
2465 pushed,
2466 local,
2467 outcomes,
2468 }
2469}
2470
2471fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2473 let mut pushed = ProjectQuery::default();
2474 let mut local = LocalProjects::default();
2475 let mut outcomes = Outcomes::default();
2476
2477 if !filters.labels.is_empty() {
2478 if capabilities.filter_by_label.is_native() {
2479 pushed.labels = filters.labels.clone();
2480 outcomes.record(Predicate::Label, Outcome::PushedDown);
2481 } else {
2482 local.labels = Some(filters.labels.clone());
2483 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2484 }
2485 }
2486 if !filters.statuses.is_empty() {
2487 if capabilities.filter_by_status.is_native() {
2488 pushed.statuses.clone_from(&filters.statuses);
2489 outcomes.record(Predicate::Status, Outcome::PushedDown);
2490 } else {
2491 local.statuses.clone_from(&filters.statuses);
2492 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2493 }
2494 }
2495 if let Some(text) = &filters.text {
2496 let predicates = text_predicates(text.fields);
2497 if searches_natively(capabilities, text.fields) {
2498 pushed.text = Some(text.clone());
2499 outcomes.record_all(predicates, Outcome::PushedDown);
2500 } else {
2501 local.text = Some(text.clone());
2502 outcomes.record_all(predicates, Outcome::AppliedLocally);
2503 }
2504 }
2505
2506 ProjectShape {
2507 pushed,
2508 local,
2509 outcomes,
2510 }
2511}
2512
2513fn shape_documents(
2518 capabilities: &Capabilities,
2519 filters: &DocumentFilters,
2520 project: &ProjectFilter,
2521) -> DocumentShape {
2522 let mut pushed = DocumentQuery::default();
2523 let mut local = LocalDocuments::default();
2524 let mut outcomes = Outcomes::default();
2525
2526 if !filters.labels.is_empty() {
2527 if capabilities.filter_by_label.is_native() {
2528 pushed.labels = filters.labels.clone();
2529 outcomes.record(Predicate::Label, Outcome::PushedDown);
2530 } else {
2531 local.labels = Some(filters.labels.clone());
2532 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2533 }
2534 }
2535 if let Some(text) = &filters.text {
2536 let predicates = text_predicates(text.fields);
2537 if searches_natively(capabilities, text.fields) {
2538 pushed.text = Some(text.clone());
2539 outcomes.record_all(predicates, Outcome::PushedDown);
2540 } else {
2541 local.text = Some(text.clone());
2542 outcomes.record_all(predicates, Outcome::AppliedLocally);
2543 }
2544 }
2545 match project {
2546 ProjectFilter::Any => {}
2547 ProjectFilter::Orphans => {
2548 if capabilities.orphan_tasks.is_native() {
2549 pushed.project = ProjectFilter::Orphans;
2550 outcomes.record(Predicate::Project, Outcome::PushedDown);
2551 } else {
2552 local.project = Some(ProjectFilter::Orphans);
2553 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2554 }
2555 }
2556 ProjectFilter::Is(id) => {
2557 if capabilities.projects.is_native() {
2558 pushed.project = ProjectFilter::Is(id.clone());
2559 outcomes.record(Predicate::Project, Outcome::PushedDown);
2560 } else {
2561 local.project = Some(ProjectFilter::Is(id.clone()));
2562 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2563 }
2564 }
2565 }
2566
2567 DocumentShape {
2568 pushed,
2569 local,
2570 outcomes,
2571 }
2572}
2573
2574fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2576 match stream {
2577 StreamKind::Projects => {
2578 let shaped = shape_projects(capabilities, filters);
2579 HitShape {
2580 stream,
2581 tasks: TaskQuery::default(),
2582 projects: shaped.pushed,
2583 local_tasks: LocalTasks::default(),
2584 local_projects: shaped.local,
2585 outcomes: shaped.outcomes,
2586 }
2587 }
2588 StreamKind::Items | StreamKind::Tasks => {
2589 let shaped = shape_tasks(
2590 capabilities,
2591 filters,
2592 &ProjectFilter::Any,
2593 &[],
2594 None,
2595 &[],
2596 None,
2597 );
2598 HitShape {
2599 stream,
2600 tasks: shaped.pushed,
2601 projects: ProjectQuery::default(),
2602 local_tasks: shaped.local,
2603 local_projects: LocalProjects::default(),
2604 outcomes: shaped.outcomes,
2605 }
2606 }
2607 }
2608}
2609
2610fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2617 if compensating {
2618 ceiling
2619 } else {
2620 budget.min(ceiling)
2621 }
2622}
2623
2624fn ceiling(source: &ResolvedSource) -> u32 {
2626 source.source().capabilities().max_page_size.max(1)
2627}
2628
2629async fn fetch_tasks(
2631 source: &ResolvedSource,
2632 shape: &TaskShape,
2633 start: &Resume,
2634 budget: u32,
2635 calls: &AtomicU32,
2636) -> Result<Fetched<Task>, SourceError> {
2637 let compensating = shape.local != LocalTasks::default();
2638 walk(
2639 start,
2640 budget,
2641 page_size(compensating, budget, ceiling(source)),
2642 |task| shape.local.keeps(task),
2643 |cursor, limit| async move {
2644 calls.fetch_add(1, Ordering::Relaxed);
2645 let request = PageRequest { cursor, limit };
2646 let page = source.source().query_tasks(&shape.pushed, &request).await?;
2647 match shape.local.commented_since {
2648 None => Ok(page),
2649 Some(since) => comment::commented_since(source, &shape.local, page, since).await,
2650 }
2651 },
2652 )
2653 .await
2654}
2655
2656async fn fetch_projects(
2658 source: &ResolvedSource,
2659 shape: &ProjectShape,
2660 start: &Resume,
2661 budget: u32,
2662 calls: &AtomicU32,
2663) -> Result<Fetched<Project>, SourceError> {
2664 let compensating = shape.local != LocalProjects::default();
2665 walk(
2666 start,
2667 budget,
2668 page_size(compensating, budget, ceiling(source)),
2669 |project| shape.local.keeps(project),
2670 |cursor, limit| async move {
2671 calls.fetch_add(1, Ordering::Relaxed);
2672 let request = PageRequest { cursor, limit };
2673 source
2674 .source()
2675 .query_projects(&shape.pushed, &request)
2676 .await
2677 },
2678 )
2679 .await
2680}
2681
2682async fn fetch_documents(
2684 source: &ResolvedSource,
2685 shape: &DocumentShape,
2686 start: &Resume,
2687 budget: u32,
2688 calls: &AtomicU32,
2689) -> Result<Fetched<Document>, SourceError> {
2690 let compensating = shape.local != LocalDocuments::default();
2691 walk(
2692 start,
2693 budget,
2694 page_size(compensating, budget, ceiling(source)),
2695 |document| shape.local.keeps(document),
2696 |cursor, limit| async move {
2697 calls.fetch_add(1, Ordering::Relaxed);
2698 let request = PageRequest { cursor, limit };
2699 source
2700 .source()
2701 .query_documents(&shape.pushed, &request)
2702 .await
2703 },
2704 )
2705 .await
2706}
2707
2708async fn fetch_labels(
2710 source: &ResolvedSource,
2711 start: &Resume,
2712 budget: u32,
2713 calls: &AtomicU32,
2714) -> Result<Fetched<Label>, SourceError> {
2715 walk(
2716 start,
2717 budget,
2718 page_size(false, budget, ceiling(source)),
2719 |_| true,
2720 |cursor, limit| async move {
2721 calls.fetch_add(1, Ordering::Relaxed);
2722 let request = PageRequest { cursor, limit };
2723 source.source().labels(&request).await
2724 },
2725 )
2726 .await
2727}
2728
2729async fn fetch_hits(
2731 source: &ResolvedSource,
2732 shape: &HitShape,
2733 start: &Resume,
2734 budget: u32,
2735 calls: &AtomicU32,
2736) -> Result<Fetched<Found>, SourceError> {
2737 let ceiling = ceiling(source);
2738 match shape.stream {
2739 StreamKind::Projects => {
2740 let compensating = shape.local_projects != LocalProjects::default();
2741 walk(
2742 start,
2743 budget,
2744 page_size(compensating, budget, ceiling),
2745 |found| match found {
2746 Found::Project(project) => shape.local_projects.keeps(project),
2747 Found::Task(_) => true,
2748 },
2749 |cursor, limit| async move {
2750 calls.fetch_add(1, Ordering::Relaxed);
2751 let request = PageRequest { cursor, limit };
2752 let page = source
2753 .source()
2754 .query_projects(&shape.projects, &request)
2755 .await?;
2756 Ok(Page {
2757 items: page.items.into_iter().map(Found::Project).collect(),
2758 next: page.next,
2759 })
2760 },
2761 )
2762 .await
2763 }
2764 StreamKind::Items | StreamKind::Tasks => {
2765 let compensating = shape.local_tasks != LocalTasks::default();
2766 walk(
2767 start,
2768 budget,
2769 page_size(compensating, budget, ceiling),
2770 |found| match found {
2771 Found::Task(task) => shape.local_tasks.keeps(task),
2772 Found::Project(_) => true,
2773 },
2774 |cursor, limit| async move {
2775 calls.fetch_add(1, Ordering::Relaxed);
2776 let request = PageRequest { cursor, limit };
2777 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2778 Ok(Page {
2779 items: page.items.into_iter().map(Found::Task).collect(),
2780 next: page.next,
2781 })
2782 },
2783 )
2784 .await
2785 }
2786 }
2787}
2788
2789async fn forward_edges(
2791 source: &ResolvedSource,
2792 entity: Entity,
2793 id: &NativeId,
2794 request: &PageRequest,
2795) -> Result<Page<DependencyEdge>, SourceError> {
2796 match entity {
2797 Entity::Task => {
2798 source
2799 .source()
2800 .task_dependencies(id, Direction::DependsOn, request)
2801 .await
2802 }
2803 Entity::Project => {
2804 source
2805 .source()
2806 .project_dependencies(id, Direction::DependsOn, request)
2807 .await
2808 }
2809 }
2810}
2811
2812#[expect(
2821 clippy::too_many_arguments,
2822 reason = "every argument is one axis of one walk — the source, the item, the \
2823 direction, which of its two graphs, whether the reverse is emulated, where \
2824 to resume, how many rows to return and where to count calls. Grouping them \
2825 into a struct would name the same eight values one indirection further from \
2826 the loop that reads them."
2827)]
2828async fn fetch_edges(
2829 source: &ResolvedSource,
2830 native: &NativeId,
2831 direction: Direction,
2832 entity: Entity,
2833 emulating: bool,
2834 start: &Resume,
2835 budget: u32,
2836 calls: &AtomicU32,
2837) -> Result<Fetched<DependencyEdge>, SourceError> {
2838 let ceiling = ceiling(source);
2839 if !emulating {
2840 return walk(
2841 start,
2842 budget,
2843 page_size(false, budget, ceiling),
2844 |_| true,
2845 |cursor, limit| async move {
2846 calls.fetch_add(1, Ordering::Relaxed);
2847 let request = PageRequest { cursor, limit };
2848 match entity {
2849 Entity::Task => {
2850 source
2851 .source()
2852 .task_dependencies(native, direction, &request)
2853 .await
2854 }
2855 Entity::Project => {
2856 source
2857 .source()
2858 .project_dependencies(native, direction, &request)
2859 .await
2860 }
2861 }
2862 },
2863 )
2864 .await;
2865 }
2866
2867 walk(
2868 start,
2869 budget,
2870 ceiling,
2871 |_| true,
2872 |cursor, limit| async move {
2873 calls.fetch_add(1, Ordering::Relaxed);
2874 let request = PageRequest { cursor, limit };
2875 let (ids, next) = match entity {
2876 Entity::Task => {
2877 let page = source
2878 .source()
2879 .query_tasks(&TaskQuery::default(), &request)
2880 .await?;
2881 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2882 (ids, page.next)
2883 }
2884 Entity::Project => {
2885 let page = source
2886 .source()
2887 .query_projects(&ProjectQuery::default(), &request)
2888 .await?;
2889 let ids: Vec<NativeId> =
2890 page.items.into_iter().map(|project| project.id).collect();
2891 (ids, page.next)
2892 }
2893 };
2894
2895 let mut edges = Vec::new();
2896 for id in ids {
2897 let mut inner: Option<Cursor> = None;
2898 loop {
2899 calls.fetch_add(1, Ordering::Relaxed);
2900 let request = PageRequest {
2901 cursor: inner.clone(),
2902 limit,
2903 };
2904 let page = forward_edges(source, entity, &id, &request).await?;
2905 fits(page.items.len(), limit)?;
2909 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2910 unrepeated(
2911 page.next.as_ref(),
2912 inner.as_ref(),
2913 "its forward edges were being scanned",
2914 )?;
2915 match page.next {
2916 Some(cursor) => inner = Some(cursor),
2917 None => break,
2918 }
2919 }
2920 }
2921
2922 Ok(Page { items: edges, next })
2923 },
2924 )
2925 .await
2926}