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, ProjectCreate, Regenerated, Regeneration, RenderRequest, RenderTemplate,
67 RenderedRecord, 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(
501 "source {name} cannot write a {record}'s rendering: its plugin is {kind}, which has no \
502 write side\n\
503 next: render it with --dry-run to read the result, or regenerate a {record} of a source \
504 whose plugin can be written."
505 )]
506 RenderingNotWritable {
507 name: String,
509 kind: String,
511 record: RenderedRecord,
513 },
514
515 #[error(
518 "supply every required answer to regenerate {id}: {} unanswered, and {reason}\n\
519 next: answer {} with --var NAME=VALUE or an answers file (--answers FILE), or run \
520 interactively to be asked.",
521 names.join(", "),
522 if names.len() == 1 { "it" } else { "each" }
523 )]
524 MissingAnswers {
525 id: String,
527 names: Vec<String>,
529 reason: String,
532 },
533
534 #[error("{error}")]
536 Template {
537 error: crate::template::TemplateError,
539 },
540
541 #[error(
543 "{record} {id} has no stored template answers: {reason}\n\
544 next: regenerate it with every required answer (`onetaskgraph {record} render {id} \
545 --var NAME=VALUE`), or read its provenance with `onetaskgraph {record} show {id}`."
546 )]
547 NoStoredAnswers {
548 record: RenderedRecord,
550 id: String,
552 reason: String,
554 },
555
556 #[error(
558 "{record} {id} records no template it was rendered from, and none was given\n\
559 next: name one with --template FILE or --template-loader FILE."
560 )]
561 NoTemplate {
562 record: RenderedRecord,
564 id: String,
566 },
567
568 #[error(
571 "{record} {id} records a template entry this product did not write — {problem}\n\
572 next: regenerate it with --template FILE or --template-loader FILE and every required \
573 answer, which records a fresh entry."
574 )]
575 MalformedProvenance {
576 record: RenderedRecord,
578 id: String,
580 problem: String,
582 },
583
584 #[error(
586 "{record} {id} was rendered from {reference:?}, which is not a readable file, so it \
587 cannot be re-read; a recorded reference is never turned into a location\n\
588 next: supply the template with --template-loader FILE (a loader document naming what \
589 to render), or name a template file with --template FILE."
590 )]
591 TemplateNotAFile {
592 record: RenderedRecord,
594 id: String,
596 reference: String,
598 },
599
600 #[error(
605 "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
606 be written to it: its plugin is {kind}, which declares priority unsupported\n\
607 next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
608 reports what each declares — or set the task's priority to none first; a \
609 github-projects source holds one once its configuration sets priority_mapping."
610 )]
611 NoPriority {
612 name: String,
614 kind: String,
616 task: String,
618 priority: onetaskgraph_plugin_api::Priority,
620 },
621
622 #[error(
624 "no project with the id {id}\n\
625 next: check the id, or list what is there — `onetaskgraph project list` reports every \
626 project the configured sources hold."
627 )]
628 NoSuchProject {
629 id: String,
631 },
632
633 #[error(
635 "no document with the id {id}\n\
636 next: check the id, or list what is there — `onetaskgraph document list` reports every \
637 document the configured sources hold."
638 )]
639 NoSuchDocument {
640 id: String,
642 },
643
644 #[error(
646 "no task with the id {id}\n\
647 next: check the id, or list what is there — `onetaskgraph task list` reports every \
648 task the configured sources hold."
649 )]
650 NoSuchTask {
651 id: String,
653 },
654
655 #[error(
657 "task {task} has no comment with the id {comment}\n\
658 next: list its comments — `onetaskgraph task comment list {task}` reports each \
659 one's id."
660 )]
661 NoSuchComment {
662 task: String,
664 comment: String,
666 },
667
668 #[error(
673 "source {name} could not be built: {error}\n\
674 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
675 command again."
676 )]
677 SourceUnavailable {
678 name: String,
680 error: SourceError,
682 },
683
684 #[error(
690 "source {name} could not do it: {error}\n\
691 next: fix what the source named above, then run the command again."
692 )]
693 SourceFailed {
694 name: String,
696 error: SourceError,
698 },
699
700 #[error(
702 "the destination source {name} could not be built: {error}\n\
703 next: fix that source — `onetaskgraph sources list` reports its state — then \
704 copy again."
705 )]
706 DestinationUnavailable {
707 name: String,
709 error: SourceError,
711 },
712
713 #[error(
715 "no item with the id {id}\n\
716 next: check the id, or list what is there — `onetaskgraph task list` and \
717 `onetaskgraph project list` report what the configured sources hold."
718 )]
719 NoSuchItem {
720 id: String,
722 },
723
724 #[error(
729 "{item} was copied from {origin}, which that destination no longer holds\n\
730 next: re-run with --recreate to create a new item there instead, or restore \
731 {origin}."
732 )]
733 StaleOrigin {
734 item: String,
736 origin: String,
738 },
739
740 #[error(
746 "{item} was last copied to {link}, which that destination no longer holds\n\
747 next: re-run with --recreate to create a new item there instead, or restore \
748 {link}."
749 )]
750 StaleLink {
751 item: String,
753 link: String,
755 },
756
757 #[error(
760 "--create cannot be given with {flag}: --create asserts the destination holds no \
761 counterpart, so there is nothing for {flag} to look for\n\
762 next: drop {flag} to create each item without looking, or drop --create to look."
763 )]
764 CreateWith {
765 flag: CopyLookup,
767 },
768
769 #[error(
775 "{item} already records a counterpart at the destination, {carrier}, and --create \
776 asserts it has none\n\
777 next: copy it without --create, which updates {carrier}."
778 )]
779 CreateCarried {
780 item: GlobalId,
782 carrier: GlobalId,
784 },
785
786 #[error(
792 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
793 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
794 {id} on its own with `onetaskgraph task copy`.",
795 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
796 )]
797 NotAMember {
798 id: GlobalId,
800 projects: Vec<GlobalId>,
802 },
803
804 #[error(
813 "{item} depends on {member}, which this copy was not told to carry and which records \
814 no origin in {destination}\n\
815 next: name {member} with --member as well, record its {destination} id at \
816 onetaskgraph.origin, or copy the whole project without --member."
817 )]
818 UnrecordedMember {
819 item: GlobalId,
821 member: GlobalId,
823 destination: SourceName,
825 },
826
827 #[error(
834 "{item} was copied to {counterpart}, but its repositories now route it to {route}\n\
835 next: move or remove {counterpart} by hand and copy again, or restore {item}'s \
836 repositories so it routes to {source} again.",
837 source = .counterpart.source
838 )]
839 Misrouted {
840 item: GlobalId,
842 counterpart: GlobalId,
844 route: SourceName,
846 },
847
848 #[error(
854 "source {name} could not do it: {error}\n\
855 next: fix what the source named above, then copy again."
856 )]
857 SourceRefused {
858 name: String,
860 error: SourceError,
862 },
863
864 #[error(
872 "the copy failed and could not be undone.\n\
873 it failed because: {error}\n\
874 it could not be undone because: {refusal}\n\
875 so these still hold what it wrote: {left_behind}\n\
876 next: remove or put back those items, then copy again."
877 )]
878 CopyNotUndone {
879 error: Box<EngineError>,
881 left_behind: LeftBehind,
887 refusal: SourceError,
889 },
890}
891
892#[derive(Debug, Clone, PartialEq)]
901pub struct LeftBehind {
902 first: GlobalId,
904 rest: Vec<GlobalId>,
906}
907
908impl LeftBehind {
909 #[must_use]
911 pub fn new(first: GlobalId) -> Self {
912 Self {
913 first,
914 rest: Vec::new(),
915 }
916 }
917
918 pub fn push(&mut self, id: GlobalId) {
920 self.rest.push(id);
921 }
922
923 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
925 std::iter::once(&self.first).chain(self.rest.iter())
926 }
927}
928
929impl std::fmt::Display for LeftBehind {
930 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
932 write!(formatter, "{}", self.first)?;
933 for id in &self.rest {
934 write!(formatter, ", {id}")?;
935 }
936 Ok(())
937 }
938}
939
940pub enum ConfiguredSource {
947 Ready(ResolvedSource),
949 Unavailable(UnavailableSource),
951}
952
953impl ConfiguredSource {
954 #[must_use]
956 pub fn name(&self) -> &SourceName {
957 match self {
958 Self::Ready(source) => source.name(),
959 Self::Unavailable(source) => source.name(),
960 }
961 }
962}
963
964pub struct Engine {
966 sources: Vec<ConfiguredSource>,
968 selection: Vec<SourceName>,
970 routes: Routes,
973}
974
975impl Engine {
976 #[must_use]
984 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
985 let (ready, unavailable) = resolve_available(config, secrets);
986 Self::new(
987 ready
988 .into_iter()
989 .map(ConfiguredSource::Ready)
990 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
991 .collect(),
992 config.selected_sources(),
993 )
994 .with_routes(config.routes())
995 }
996
997 #[must_use]
1000 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
1001 Self {
1002 sources,
1003 selection,
1004 routes: Routes::default(),
1005 }
1006 }
1007
1008 #[must_use]
1013 pub fn with_routes(mut self, routes: Routes) -> Self {
1014 self.routes = routes;
1015 self
1016 }
1017
1018 #[must_use]
1020 pub fn place(&self, source: &SourceName, repositories: &[Repository]) -> Placement {
1021 self.routes.place(source, repositories)
1022 }
1023
1024 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
1026 self.sources.iter().filter_map(|source| match source {
1027 ConfiguredSource::Ready(ready) => Some(ready),
1028 ConfiguredSource::Unavailable(_) => None,
1029 })
1030 }
1031
1032 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
1034 self.sources.iter().filter_map(|source| match source {
1035 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
1036 ConfiguredSource::Ready(_) => None,
1037 })
1038 }
1039
1040 #[must_use]
1042 pub fn listing(&self) -> Vec<SourceListing> {
1043 let mut listings: Vec<SourceListing> = self
1044 .ready()
1045 .map(|source| SourceListing {
1046 source: source.name().clone(),
1047 kind: source.kind().to_owned(),
1048 state: SourceState::Available {
1049 capabilities: source.source().capabilities(),
1050 },
1051 })
1052 .chain(self.unavailable().map(|source| SourceListing {
1053 source: source.name().clone(),
1054 kind: source.kind().to_owned(),
1055 state: SourceState::Unavailable {
1056 error: source.error().clone(),
1057 },
1058 }))
1059 .collect();
1060 listings.sort_by(|left, right| left.source.cmp(&right.source));
1061 listings
1062 }
1063
1064 #[must_use]
1070 pub fn has(&self, name: &SourceName) -> bool {
1071 self.sources.iter().any(|source| source.name() == name)
1072 }
1073
1074 pub async fn end_command(&self) -> Result<(), EngineError> {
1094 let mut first = None;
1095 for source in self.ready() {
1096 if let Err(error) = source.source().end_command().await
1097 && first.is_none()
1098 {
1099 first = Some(EngineError::SourceFailed {
1100 name: source.name().to_string(),
1101 error,
1102 });
1103 }
1104 }
1105 first.map_or(Ok(()), Err)
1106 }
1107
1108 pub async fn tasks(
1116 &self,
1117 request: &TaskRequest,
1118 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1119 let mut names = self.resolve_selection(&request.sources)?;
1120 let mut members: Vec<GlobalId> = Vec::new();
1125 let mut unread: Vec<SourceFailure> = Vec::new();
1126 if let ProjectSelector::Qualified(id) = &request.project {
1127 self.known(&id.source)?;
1128 names.retain(|name| name == &id.source);
1129 if request.include_members {
1132 (members, unread) = self.members(id).await;
1133 names.extend(members.iter().map(|member| member.source.clone()));
1134 }
1135 }
1136 let selector = |name: &SourceName| -> ProjectSelector {
1137 members
1138 .iter()
1139 .find(|member| &member.source == name)
1140 .map_or_else(
1141 || request.project.clone(),
1142 |member| ProjectSelector::Qualified(member.clone()),
1143 )
1144 };
1145 let query = shape(
1146 "task-list",
1147 &names,
1148 &(
1149 &request.filters,
1150 &request.project,
1151 &request.priorities,
1152 &request.commented_since,
1153 &request.metadata,
1154 &request.origin,
1155 &members,
1156 ),
1157 );
1158 let states = resumption(
1159 self,
1160 request.paging.token.as_ref(),
1161 &[StreamKind::Items],
1162 &query,
1163 )?;
1164 let budget = request.paging.limit.get();
1165
1166 let mut answer = Answer::new();
1167 answer.errors.extend(unread);
1168 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1169
1170 let shapes: Vec<TaskShape> = ready
1171 .iter()
1172 .map(|source| {
1173 shape_tasks(
1174 &source.source().capabilities(),
1175 &request.filters,
1176 &project_filter(&selector(source.name())),
1177 &request.priorities,
1178 request.commented_since,
1179 &request.metadata,
1180 request.origin.as_ref(),
1181 )
1182 })
1183 .collect();
1184 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1185 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1186
1187 let walks = ready
1188 .iter()
1189 .enumerate()
1190 .map(|(index, source)| {
1191 fetch_tasks(
1192 source,
1193 &shapes[index],
1194 &starts[index],
1195 budget,
1196 &counters[index],
1197 )
1198 })
1199 .collect();
1200
1201 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1202 answer.finish(
1203 streams,
1204 budget,
1205 owed(&states),
1206 &query,
1207 |name, task: Task| {
1208 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1209 },
1210 )
1211 }
1212
1213 async fn members(&self, home: &GlobalId) -> (Vec<GlobalId>, Vec<SourceFailure>) {
1224 let mut members = Vec::new();
1225 let mut unread = Vec::new();
1226 let Some(source) = self.ready().find(|source| source.name() == &home.source) else {
1227 return (members, unread);
1228 };
1229 let held = match source.source().get_project(&home.native).await {
1230 Ok(Some(held)) => held,
1231 Ok(None) => {
1232 unread.push(SourceFailure {
1233 source: home.source.clone(),
1234 error: SourceError::Refused {
1235 message: format!(
1236 "{home} names no project {} holds, so its member projects cannot \
1237 be read",
1238 home.source
1239 ),
1240 },
1241 });
1242 return (members, unread);
1243 }
1244 Err(error) => {
1245 unread.push(SourceFailure {
1246 source: home.source.clone(),
1247 error,
1248 });
1249 return (members, unread);
1250 }
1251 };
1252 let named = match copy::members_of(home, &held.metadata) {
1253 Ok(named) => named,
1254 Err(message) => {
1257 unread.push(SourceFailure {
1258 source: home.source.clone(),
1259 error: SourceError::Malformed { message },
1260 });
1261 return (members, unread);
1262 }
1263 };
1264 for member in named {
1265 if !self.has(&member.source) {
1266 unread.push(SourceFailure {
1267 source: member.source.clone(),
1268 error: SourceError::Config {
1269 message: format!(
1270 "{home} names {member} as a member project, but no source named \
1271 {} is configured",
1272 member.source
1273 ),
1274 },
1275 });
1276 continue;
1277 }
1278 let Some(there) = self.ready().find(|source| source.name() == &member.source) else {
1281 members.push(member);
1282 continue;
1283 };
1284 match there.source().get_project(&member.native).await {
1285 Ok(Some(_)) => members.push(member),
1286 Ok(None) => unread.push(SourceFailure {
1287 source: member.source.clone(),
1288 error: SourceError::Refused {
1289 message: format!(
1290 "{home} names {member} as a member project, and {} holds no such \
1291 project",
1292 member.source
1293 ),
1294 },
1295 }),
1296 Err(error) => unread.push(SourceFailure {
1297 source: member.source.clone(),
1298 error,
1299 }),
1300 }
1301 }
1302 (members, unread)
1303 }
1304
1305 pub async fn projects(
1311 &self,
1312 request: &ProjectRequest,
1313 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1314 let names = self.resolve_selection(&request.sources)?;
1315 let query = shape("project-list", &names, &request.filters);
1316 let states = resumption(
1317 self,
1318 request.paging.token.as_ref(),
1319 &[StreamKind::Items],
1320 &query,
1321 )?;
1322 let budget = request.paging.limit.get();
1323
1324 let mut answer = Answer::new();
1325 let mut with_projects = Vec::new();
1331 for source in answer.split(self, &names) {
1332 if source.source().capabilities().projects.is_native() {
1333 with_projects.push(source);
1334 } else {
1335 answer.unreachable_predicate(source, Predicate::Project);
1336 }
1337 }
1338 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1339
1340 let shapes: Vec<ProjectShape> = ready
1341 .iter()
1342 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1343 .collect();
1344 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1345 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1346
1347 let walks = ready
1348 .iter()
1349 .enumerate()
1350 .map(|(index, source)| {
1351 fetch_projects(
1352 source,
1353 &shapes[index],
1354 &starts[index],
1355 budget,
1356 &counters[index],
1357 )
1358 })
1359 .collect();
1360
1361 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1362 answer.finish(
1363 streams,
1364 budget,
1365 owed(&states),
1366 &query,
1367 |name, project: Project| Qualified {
1368 id: GlobalId::new(name.clone(), project.id.clone()),
1369 item: project,
1370 },
1371 )
1372 }
1373
1374 pub async fn documents(
1386 &self,
1387 request: &DocumentRequest,
1388 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1389 let mut names = self.resolve_selection(&request.sources)?;
1390 if let ProjectSelector::Qualified(id) = &request.project {
1393 self.known(&id.source)?;
1394 names.retain(|name| name == &id.source);
1395 }
1396 let query = shape(
1397 "document-list",
1398 &names,
1399 &(&request.filters, &request.project),
1400 );
1401 let states = resumption(
1402 self,
1403 request.paging.token.as_ref(),
1404 &[StreamKind::Items],
1405 &query,
1406 )?;
1407 let budget = request.paging.limit.get();
1408
1409 let mut answer = Answer::new();
1410 let mut with_documents = Vec::new();
1411 for source in answer.split(self, &names) {
1412 if source.source().capabilities().documents.is_native() {
1413 with_documents.push(source);
1414 } else {
1415 answer.unreachable_predicate(source, Predicate::Document);
1416 }
1417 }
1418 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1419
1420 let shapes: Vec<DocumentShape> = ready
1421 .iter()
1422 .map(|source| {
1423 shape_documents(
1424 &source.source().capabilities(),
1425 &request.filters,
1426 &project_filter(&request.project),
1427 )
1428 })
1429 .collect();
1430 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1431 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1432
1433 let walks = ready
1434 .iter()
1435 .enumerate()
1436 .map(|(index, source)| {
1437 fetch_documents(
1438 source,
1439 &shapes[index],
1440 &starts[index],
1441 budget,
1442 &counters[index],
1443 )
1444 })
1445 .collect();
1446
1447 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1448 answer.finish(
1449 streams,
1450 budget,
1451 owed(&states),
1452 &query,
1453 |name, document: Document| Qualified {
1454 id: GlobalId::new(name.clone(), document.id.clone()),
1455 item: document,
1456 },
1457 )
1458 }
1459
1460 pub async fn labels(
1466 &self,
1467 request: &LabelRequest,
1468 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1469 let names = self.resolve_selection(&request.sources)?;
1470 let query = shape("label-list", &names, &());
1471 let states = resumption(
1472 self,
1473 request.paging.token.as_ref(),
1474 &[StreamKind::Items],
1475 &query,
1476 )?;
1477 let budget = request.paging.limit.get();
1478
1479 let mut answer = Answer::new();
1480 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1481 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1482 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1483
1484 let walks = ready
1485 .iter()
1486 .enumerate()
1487 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1488 .collect();
1489
1490 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1491 answer.finish(
1492 streams,
1493 budget,
1494 owed(&states),
1495 &query,
1496 |name, label: Label| Qualified {
1497 id: GlobalId::new(name.clone(), label.id.clone()),
1498 item: label,
1499 },
1500 )
1501 }
1502
1503 pub async fn search(
1509 &self,
1510 request: &SearchRequest,
1511 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1512 let names = self.resolve_selection(&request.sources)?;
1513 let reads: &[StreamKind] = match request.kind {
1517 SearchKind::Tasks => &[StreamKind::Tasks],
1518 SearchKind::Projects => &[StreamKind::Projects],
1519 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1520 };
1521 let query = shape("search", &names, &request.text);
1526 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1527 let budget = request.paging.limit.get();
1528 let filters = Filters {
1529 text: Some(request.text.clone()),
1530 ..Filters::default()
1531 };
1532
1533 let mut answer = Answer::new();
1534
1535 let mut ready = Vec::new();
1538 let mut kinds = Vec::new();
1539 let mut starts = Vec::new();
1540 for source in answer.split(self, &names) {
1541 let mut streams = Vec::new();
1542 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1543 streams.push(StreamKind::Tasks);
1544 }
1545 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1546 if source.source().capabilities().projects.is_native() {
1547 streams.push(StreamKind::Projects);
1548 } else {
1549 answer.unreachable_predicate(source, Predicate::Project);
1550 }
1551 }
1552 for stream in streams {
1553 if let Some(resume) = resume_at(&states, source.name(), stream) {
1554 ready.push(source);
1555 kinds.push(stream);
1556 starts.push(resume);
1557 }
1558 }
1559 }
1560
1561 let shapes: Vec<HitShape> = ready
1562 .iter()
1563 .zip(kinds.iter())
1564 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1565 .collect();
1566 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1567 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1568
1569 let walks = ready
1570 .iter()
1571 .enumerate()
1572 .map(|(index, source)| {
1573 fetch_hits(
1574 source,
1575 &shapes[index],
1576 &starts[index],
1577 budget,
1578 &counters[index],
1579 )
1580 })
1581 .collect();
1582
1583 let streams =
1584 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1585 answer.finish(
1586 streams,
1587 budget,
1588 owed(&states),
1589 &query,
1590 |name, found: Found| match found {
1591 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1592 GlobalId::new(name.clone(), task.id.clone()),
1593 task,
1594 )),
1595 Found::Project(project) => SearchHit::Project(Qualified {
1596 id: GlobalId::new(name.clone(), project.id.clone()),
1597 item: project,
1598 }),
1599 },
1600 )
1601 }
1602
1603 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1610 let name = self.known(&id.source)?;
1611 let mut answer = Answer::new();
1612 let selected = answer.split(self, std::slice::from_ref(&name));
1613 let Some(source) = selected.first() else {
1614 return answer.nothing();
1615 };
1616 let found = source.source().get_task(&id.native).await;
1617 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1618 answer.one(source, found, |task| {
1619 delivery::qualified_task(qualified, task)
1620 })
1621 }
1622
1623 pub async fn project(
1629 &self,
1630 id: &GlobalId,
1631 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1632 let name = self.known(&id.source)?;
1633 let mut answer = Answer::new();
1634 let selected = answer.split(self, std::slice::from_ref(&name));
1635 let Some(source) = selected.first() else {
1636 return answer.nothing();
1637 };
1638 let found = source.source().get_project(&id.native).await;
1639 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1640 answer.one(source, found, |project| Qualified {
1641 id: qualified,
1642 item: project,
1643 })
1644 }
1645
1646 pub async fn document(
1656 &self,
1657 id: &GlobalId,
1658 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1659 let name = self.known(&id.source)?;
1660 let mut answer = Answer::new();
1661 let selected = answer.split(self, std::slice::from_ref(&name));
1662 let Some(source) = selected.first() else {
1663 return answer.nothing();
1664 };
1665 if !source.source().capabilities().documents.is_native() {
1666 answer.unreachable_predicate(source, Predicate::Document);
1667 return answer.nothing();
1668 }
1669 let found = source.source().get_document(&id.native).await;
1670 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1671 answer.one(source, found, |document| Qualified {
1672 id: qualified,
1673 item: document,
1674 })
1675 }
1676
1677 pub async fn task_dependencies(
1684 &self,
1685 request: &DependencyRequest,
1686 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1687 self.dependencies(request, Entity::Task).await
1688 }
1689
1690 pub async fn project_dependencies(
1696 &self,
1697 request: &DependencyRequest,
1698 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1699 self.dependencies(request, Entity::Project).await
1700 }
1701
1702 async fn dependencies(
1705 &self,
1706 request: &DependencyRequest,
1707 entity: Entity,
1708 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1709 let name = self.known(&request.id.source)?;
1710 let query = shape(
1711 "dependencies",
1712 std::slice::from_ref(&name),
1713 &(entity, &request.id.native, request.direction),
1714 );
1715 let states = resumption(
1716 self,
1717 request.paging.token.as_ref(),
1718 &[StreamKind::Items],
1719 &query,
1720 )?;
1721 let budget = request.paging.limit.get();
1722
1723 let mut answer = Answer::new();
1724 let (ready, starts) = walking(
1725 answer.split(self, std::slice::from_ref(&name)),
1726 &states,
1727 StreamKind::Items,
1728 );
1729 let Some(source) = ready.first() else {
1730 return answer.nothing();
1731 };
1732
1733 let capabilities = source.source().capabilities();
1734 let support = match entity {
1735 Entity::Task => capabilities.task_dependencies,
1736 Entity::Project => capabilities.project_dependencies,
1737 };
1738 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1742 let mut outcomes = Outcomes::default();
1743 if request.direction == Direction::DependedOnBy {
1744 if emulating {
1745 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1746 } else {
1747 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1748 }
1749 }
1750
1751 let counters = vec![AtomicU32::new(0)];
1752 let walked = fetch_edges(
1753 source,
1754 &request.id.native,
1755 request.direction,
1756 entity,
1757 emulating,
1758 &starts[0],
1759 budget,
1760 &counters[0],
1761 )
1762 .await;
1763
1764 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1765 answer.finish(
1766 streams,
1767 budget,
1768 owed(&states),
1769 &query,
1770 |name, edge: DependencyEdge| QualifiedEdge {
1771 from: qualify_endpoint(name, edge.from),
1772 to: qualify_endpoint(name, edge.to),
1773 kind: edge.kind,
1774 },
1775 )
1776 }
1777
1778 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1780 if asked.is_empty() {
1781 if self.selection.is_empty() {
1782 return Err(EngineError::NoSources);
1783 }
1784 return Ok(self.selection.clone());
1785 }
1786 asked.iter().map(|name| self.known(name)).collect()
1787 }
1788
1789 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1791 if self.has(name) {
1792 return Ok(name.clone());
1793 }
1794 if self.sources.is_empty() {
1795 return Err(EngineError::NoSources);
1796 }
1797 Err(EngineError::UnknownSource {
1798 name: name.to_string(),
1799 configured: self
1800 .listing()
1801 .iter()
1802 .map(|listing| listing.source.to_string())
1803 .collect::<Vec<_>>()
1804 .join(", "),
1805 })
1806 }
1807}
1808
1809fn qualify_endpoint(
1810 source: &SourceName,
1811 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1812) -> QualifiedEndpoint {
1813 let kind = endpoint.kind;
1814 let is_qualified = endpoint.is_qualified();
1815 let endpoint_id = endpoint.into_id();
1816 QualifiedEndpoint {
1817 id: if is_qualified {
1818 endpoint_id
1819 .parse()
1820 .expect("plugin-api validates qualified dependency endpoints")
1821 } else {
1822 GlobalId::new(source.clone(), NativeId(endpoint_id))
1823 },
1824 kind,
1825 }
1826}
1827
1828#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1830enum Entity {
1831 Task,
1833 Project,
1835}
1836
1837enum Found {
1839 Task(Task),
1841 Project(Project),
1843}
1844
1845#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1847enum Outcome {
1848 PushedDown,
1850 AppliedLocally,
1852 Emulated,
1854 Unavailable,
1856}
1857
1858#[derive(Debug, Clone, Default, PartialEq)]
1870struct Outcomes(BTreeMap<Predicate, Outcome>);
1871
1872impl Outcomes {
1873 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1879 self.0.insert(predicate, outcome);
1880 }
1881
1882 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1885 for predicate in predicates {
1886 self.record(predicate, outcome);
1887 }
1888 }
1889
1890 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1892 self.0
1893 .iter()
1894 .filter(|(_, recorded)| **recorded == outcome)
1895 .map(|(predicate, _)| *predicate)
1896 .collect()
1897 }
1898}
1899
1900struct TaskShape {
1902 pushed: TaskQuery,
1904 local: LocalTasks,
1906 outcomes: Outcomes,
1908}
1909
1910struct ProjectShape {
1912 pushed: ProjectQuery,
1914 local: LocalProjects,
1916 outcomes: Outcomes,
1918}
1919
1920struct DocumentShape {
1922 pushed: DocumentQuery,
1924 local: LocalDocuments,
1926 outcomes: Outcomes,
1928}
1929
1930struct HitShape {
1932 stream: StreamKind,
1934 tasks: TaskQuery,
1936 projects: ProjectQuery,
1938 local_tasks: LocalTasks,
1940 local_projects: LocalProjects,
1942 outcomes: Outcomes,
1944}
1945
1946struct Answer {
1952 plans: Vec<SourcePlan>,
1954 errors: Vec<SourceFailure>,
1956}
1957
1958impl Answer {
1959 fn new() -> Self {
1960 Self {
1961 plans: Vec::new(),
1962 errors: Vec::new(),
1963 }
1964 }
1965
1966 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1971 let mut selected = Vec::new();
1972 for name in names {
1973 match engine.sources.iter().find(|source| source.name() == name) {
1974 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1975 Some(ConfiguredSource::Unavailable(source)) => {
1976 self.errors.push(source.failure());
1977 }
1978 None => {}
1979 }
1980 }
1981 selected
1982 }
1983
1984 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1986 let mut outcomes = Outcomes::default();
1987 outcomes.record(predicate, Outcome::Unavailable);
1988 self.plans.push(plan_for(source, outcomes, 0));
1989 }
1990
1991 fn collect<T>(
1993 &mut self,
1994 ready: &[&ResolvedSource],
1995 walked: Vec<Result<Fetched<T>, SourceError>>,
1996 counters: &[AtomicU32],
1997 outcomes: Vec<Outcomes>,
1998 ) -> Vec<Stream<T>> {
1999 let kinds = vec![StreamKind::Items; ready.len()];
2000 self.collect_streams(ready, &kinds, walked, counters, outcomes)
2001 }
2002
2003 fn collect_streams<T>(
2005 &mut self,
2006 ready: &[&ResolvedSource],
2007 kinds: &[StreamKind],
2008 walked: Vec<Result<Fetched<T>, SourceError>>,
2009 counters: &[AtomicU32],
2010 outcomes: Vec<Outcomes>,
2011 ) -> Vec<Stream<T>> {
2012 let mut streams = Vec::new();
2013 for (index, result) in walked.into_iter().enumerate() {
2014 let source = ready[index];
2015 let pages = counters[index].load(Ordering::Relaxed);
2016 self.plans
2017 .push(plan_for(source, outcomes[index].clone(), pages));
2018 match result {
2019 Ok(fetched) => streams.push(Stream {
2020 source: source.name().clone(),
2021 kind: kinds[index],
2022 fetched,
2023 }),
2024 Err(error) => self.errors.push(SourceFailure {
2027 source: source.name().clone(),
2028 error,
2029 }),
2030 }
2031 }
2032 streams
2033 }
2034
2035 fn one<T, U>(
2037 self,
2038 source: &ResolvedSource,
2039 found: Result<Option<T>, SourceError>,
2040 qualify: impl FnOnce(T) -> U,
2041 ) -> Result<QueryResponse<U>, EngineError> {
2042 Ok(self.one_response(source, found, qualify))
2043 }
2044
2045 fn one_response<T, U>(
2047 mut self,
2048 source: &ResolvedSource,
2049 found: Result<Option<T>, SourceError>,
2050 qualify: impl FnOnce(T) -> U,
2051 ) -> QueryResponse<U> {
2052 self.plans.push(plan_for(source, Outcomes::default(), 1));
2053 let items = match found {
2054 Ok(Some(item)) => vec![qualify(item)],
2055 Ok(None) => Vec::new(),
2056 Err(error) => {
2057 self.errors.push(SourceFailure {
2058 source: source.name().clone(),
2059 error,
2060 });
2061 Vec::new()
2062 }
2063 };
2064 QueryResponse {
2065 items,
2066 next: None,
2067 plan: QueryPlan {
2068 per_source: merge_plans(self.plans),
2069 },
2070 errors: self.errors,
2071 }
2072 }
2073
2074 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
2076 Ok(QueryResponse {
2077 items: Vec::new(),
2078 next: None,
2079 plan: QueryPlan {
2080 per_source: merge_plans(self.plans),
2081 },
2082 errors: self.errors,
2083 })
2084 }
2085
2086 fn finish<T, U>(
2091 self,
2092 streams: Vec<Stream<T>>,
2093 budget: u32,
2094 first: Option<&Owed>,
2095 query: &str,
2096 qualify: impl Fn(&SourceName, T) -> U,
2097 ) -> Result<QueryResponse<U>, EngineError> {
2098 let (rows, states, owed) = merge(streams, budget, first);
2099 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
2100 Ok(QueryResponse {
2101 items: rows
2102 .into_iter()
2103 .map(|(name, item)| qualify(&name, item))
2104 .collect(),
2105 next,
2106 plan: QueryPlan {
2107 per_source: merge_plans(self.plans),
2108 },
2109 errors: self.errors,
2110 })
2111 }
2112}
2113
2114fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
2117 SourcePlan {
2118 source: source.name().clone(),
2119 kind: source.kind().to_owned(),
2120 pushed_down: outcomes.with(Outcome::PushedDown),
2121 applied_locally: outcomes.with(Outcome::AppliedLocally),
2122 emulated: outcomes.with(Outcome::Emulated),
2123 unavailable: outcomes.with(Outcome::Unavailable),
2124 pages_fetched: pages,
2125 }
2126}
2127
2128fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
2133 let mut merged: Vec<SourcePlan> = Vec::new();
2134 for plan in plans {
2135 if let Some(existing) = merged
2136 .iter_mut()
2137 .find(|existing| existing.source == plan.source)
2138 {
2139 existing.pushed_down.extend(plan.pushed_down);
2140 existing.applied_locally.extend(plan.applied_locally);
2141 existing.emulated.extend(plan.emulated);
2142 existing.unavailable.extend(plan.unavailable);
2143 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
2144 for list in [
2145 &mut existing.pushed_down,
2146 &mut existing.applied_locally,
2147 &mut existing.emulated,
2148 &mut existing.unavailable,
2149 ] {
2150 list.sort_unstable();
2151 list.dedup();
2152 }
2153 } else {
2154 merged.push(plan);
2155 }
2156 }
2157 merged
2158}
2159
2160fn walking<'a>(
2166 selected: Vec<&'a ResolvedSource>,
2167 states: &Option<Resumption>,
2168 kind: StreamKind,
2169) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
2170 let mut ready = Vec::new();
2171 let mut starts = Vec::new();
2172 for source in selected {
2173 if let Some(resume) = resume_at(states, source.name(), kind) {
2174 ready.push(source);
2175 starts.push(resume);
2176 }
2177 }
2178 (ready, starts)
2179}
2180
2181fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
2199 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
2200 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
2201}
2202
2203fn fingerprint(text: &str) -> String {
2205 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
2206 for byte in text.as_bytes() {
2207 hash ^= u64::from(*byte);
2208 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
2209 }
2210 format!("{hash:016x}")
2211}
2212
2213fn owed(document: &Option<Resumption>) -> Option<&Owed> {
2218 document.as_ref()?.owed.as_ref()
2219}
2220
2221fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
2223 match states {
2224 None => Some(Resume::default()),
2225 Some(document) => document
2226 .streams
2227 .iter()
2228 .find(|state| &state.source == source && state.stream == kind)
2229 .map(|state| state.resume.clone()),
2230 }
2231}
2232
2233fn resumption(
2264 engine: &Engine,
2265 token: Option<&PageToken>,
2266 reads: &[StreamKind],
2267 query: &str,
2268) -> Result<Option<Resumption>, EngineError> {
2269 let Some(document) = token.map(PageToken::decode) else {
2270 return Ok(None);
2271 };
2272
2273 if document.query != query {
2280 return Err(EngineError::Token {
2281 message: "this page token was written by a different query — resume the walk it \
2282 came from, or drop --page to start this one from the beginning"
2283 .to_owned(),
2284 });
2285 }
2286 let states = &document.streams;
2287
2288 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2289 for state in states {
2290 if !reads.contains(&state.stream) {
2291 return Err(EngineError::Token {
2292 message: format!(
2293 "this page token resumes {}, which this command does not read — it \
2294 was written by a different query",
2295 state.stream.describe()
2296 ),
2297 });
2298 }
2299 let ceiling = engine
2300 .ready()
2301 .find(|source| source.name() == &state.source)
2302 .map(ceiling);
2303 if ceiling.is_none() && !engine.has(&state.source) {
2304 return Err(EngineError::Token {
2305 message: format!(
2306 "this page token resumes a source called {:?}, which this \
2307 configuration does not have",
2308 state.source.as_str()
2309 ),
2310 });
2311 }
2312 if let Some(ceiling) = ceiling
2313 && state.resume.skip >= ceiling
2314 {
2315 return Err(EngineError::Token {
2316 message: format!(
2317 "this page token resumes {} rows into a page of source {:?}, which \
2318 serves at most {ceiling}",
2319 state.resume.skip,
2320 state.source.as_str()
2321 ),
2322 });
2323 }
2324 if seen.contains(&(&state.source, state.stream)) {
2325 return Err(EngineError::Token {
2326 message: format!(
2327 "this page token gives source {:?} two places to resume from",
2328 state.source.as_str()
2329 ),
2330 });
2331 }
2332 seen.push((&state.source, state.stream));
2333 }
2334
2335 if let Some(owed) = &document.owed
2340 && !document
2341 .streams
2342 .iter()
2343 .any(|state| state.source == owed.source && state.stream == owed.stream)
2344 {
2345 return Err(EngineError::Token {
2346 message: format!(
2347 "this page token owes the next row to a stream it does not resume, \
2348 {:?}'s {}",
2349 owed.source.as_str(),
2350 owed.stream.describe()
2351 ),
2352 });
2353 }
2354
2355 Ok(Some(document))
2356}
2357
2358fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2364 match selector {
2365 ProjectSelector::Any => ProjectFilter::Any,
2366 ProjectSelector::Orphans => ProjectFilter::Orphans,
2367 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2368 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2369 }
2370}
2371
2372fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2374 match fields {
2375 TextFields::Title => vec![Predicate::SearchTitle],
2376 TextFields::Content => vec![Predicate::SearchContent],
2377 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2378 }
2379}
2380
2381fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2388 match fields {
2389 TextFields::Title => capabilities.search_title.is_native(),
2390 TextFields::Content => capabilities.search_content.is_native(),
2391 TextFields::TitleOrContent => {
2392 capabilities.search_title.is_native() && capabilities.search_content.is_native()
2393 }
2394 }
2395}
2396
2397fn shape_tasks(
2399 capabilities: &Capabilities,
2400 filters: &Filters,
2401 project: &ProjectFilter,
2402 priorities: &[Priority],
2403 commented_since: Option<DateTime<Utc>>,
2404 metadata: &[MetadataMatch],
2405 origin: Option<&GlobalId>,
2406) -> TaskShape {
2407 let mut pushed = TaskQuery::default();
2408 let mut local = LocalTasks::default();
2409 let mut outcomes = Outcomes::default();
2410
2411 if !filters.labels.is_empty() {
2412 if capabilities.filter_by_label.is_native() {
2413 pushed.labels = filters.labels.clone();
2414 outcomes.record(Predicate::Label, Outcome::PushedDown);
2415 } else {
2416 local.labels = Some(filters.labels.clone());
2417 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2418 }
2419 }
2420 if !filters.statuses.is_empty() {
2421 if capabilities.filter_by_status.is_native() {
2422 pushed.statuses.clone_from(&filters.statuses);
2423 outcomes.record(Predicate::Status, Outcome::PushedDown);
2424 } else {
2425 local.statuses.clone_from(&filters.statuses);
2426 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2427 }
2428 }
2429 if !priorities.is_empty() {
2430 if capabilities.filter_by_priority.is_native() {
2431 pushed.priorities = priorities.to_vec();
2432 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2433 } else {
2434 local.priorities = priorities.to_vec();
2435 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2436 }
2437 }
2438 if let Some(since) = commented_since {
2439 if capabilities.filter_by_comment_activity.is_native() {
2440 pushed.commented_since = Some(since);
2441 outcomes.record(Predicate::CommentedSince, Outcome::PushedDown);
2442 } else {
2443 local.commented_since = Some(since);
2444 outcomes.record(Predicate::CommentedSince, Outcome::AppliedLocally);
2445 }
2446 }
2447 if !metadata.is_empty() {
2448 if capabilities.filter_by_metadata.is_native() {
2449 pushed.metadata = metadata.to_vec();
2450 outcomes.record(Predicate::Metadata, Outcome::PushedDown);
2451 } else {
2452 local.metadata = metadata.to_vec();
2453 outcomes.record(Predicate::Metadata, Outcome::AppliedLocally);
2454 }
2455 }
2456 if let Some(origin) = origin {
2457 if capabilities.filter_by_origin.is_native() {
2458 pushed.origin = Some(origin.to_string());
2461 outcomes.record(Predicate::Origin, Outcome::PushedDown);
2462 } else {
2463 local.origin = Some(origin.clone());
2464 outcomes.record(Predicate::Origin, Outcome::AppliedLocally);
2465 }
2466 }
2467 if let Some(text) = &filters.text {
2468 let predicates = text_predicates(text.fields);
2469 if searches_natively(capabilities, text.fields) {
2470 pushed.text = Some(text.clone());
2471 outcomes.record_all(predicates, Outcome::PushedDown);
2472 } else {
2473 local.text = Some(text.clone());
2474 outcomes.record_all(predicates, Outcome::AppliedLocally);
2475 }
2476 }
2477 match project {
2478 ProjectFilter::Any => {}
2479 ProjectFilter::Orphans => {
2480 if capabilities.orphan_tasks.is_native() {
2481 pushed.project = ProjectFilter::Orphans;
2482 outcomes.record(Predicate::Project, Outcome::PushedDown);
2483 } else {
2484 local.project = Some(ProjectFilter::Orphans);
2485 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2486 }
2487 }
2488 ProjectFilter::Is(id) => {
2489 if capabilities.projects.is_native() {
2490 pushed.project = ProjectFilter::Is(id.clone());
2491 outcomes.record(Predicate::Project, Outcome::PushedDown);
2492 } else {
2493 local.project = Some(ProjectFilter::Is(id.clone()));
2494 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2495 }
2496 }
2497 }
2498
2499 TaskShape {
2500 pushed,
2501 local,
2502 outcomes,
2503 }
2504}
2505
2506fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2508 let mut pushed = ProjectQuery::default();
2509 let mut local = LocalProjects::default();
2510 let mut outcomes = Outcomes::default();
2511
2512 if !filters.labels.is_empty() {
2513 if capabilities.filter_by_label.is_native() {
2514 pushed.labels = filters.labels.clone();
2515 outcomes.record(Predicate::Label, Outcome::PushedDown);
2516 } else {
2517 local.labels = Some(filters.labels.clone());
2518 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2519 }
2520 }
2521 if !filters.statuses.is_empty() {
2522 if capabilities.filter_by_status.is_native() {
2523 pushed.statuses.clone_from(&filters.statuses);
2524 outcomes.record(Predicate::Status, Outcome::PushedDown);
2525 } else {
2526 local.statuses.clone_from(&filters.statuses);
2527 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2528 }
2529 }
2530 if let Some(text) = &filters.text {
2531 let predicates = text_predicates(text.fields);
2532 if searches_natively(capabilities, text.fields) {
2533 pushed.text = Some(text.clone());
2534 outcomes.record_all(predicates, Outcome::PushedDown);
2535 } else {
2536 local.text = Some(text.clone());
2537 outcomes.record_all(predicates, Outcome::AppliedLocally);
2538 }
2539 }
2540
2541 ProjectShape {
2542 pushed,
2543 local,
2544 outcomes,
2545 }
2546}
2547
2548fn shape_documents(
2553 capabilities: &Capabilities,
2554 filters: &DocumentFilters,
2555 project: &ProjectFilter,
2556) -> DocumentShape {
2557 let mut pushed = DocumentQuery::default();
2558 let mut local = LocalDocuments::default();
2559 let mut outcomes = Outcomes::default();
2560
2561 if !filters.labels.is_empty() {
2562 if capabilities.filter_by_label.is_native() {
2563 pushed.labels = filters.labels.clone();
2564 outcomes.record(Predicate::Label, Outcome::PushedDown);
2565 } else {
2566 local.labels = Some(filters.labels.clone());
2567 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2568 }
2569 }
2570 if let Some(text) = &filters.text {
2571 let predicates = text_predicates(text.fields);
2572 if searches_natively(capabilities, text.fields) {
2573 pushed.text = Some(text.clone());
2574 outcomes.record_all(predicates, Outcome::PushedDown);
2575 } else {
2576 local.text = Some(text.clone());
2577 outcomes.record_all(predicates, Outcome::AppliedLocally);
2578 }
2579 }
2580 match project {
2581 ProjectFilter::Any => {}
2582 ProjectFilter::Orphans => {
2583 if capabilities.orphan_tasks.is_native() {
2584 pushed.project = ProjectFilter::Orphans;
2585 outcomes.record(Predicate::Project, Outcome::PushedDown);
2586 } else {
2587 local.project = Some(ProjectFilter::Orphans);
2588 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2589 }
2590 }
2591 ProjectFilter::Is(id) => {
2592 if capabilities.projects.is_native() {
2593 pushed.project = ProjectFilter::Is(id.clone());
2594 outcomes.record(Predicate::Project, Outcome::PushedDown);
2595 } else {
2596 local.project = Some(ProjectFilter::Is(id.clone()));
2597 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2598 }
2599 }
2600 }
2601
2602 DocumentShape {
2603 pushed,
2604 local,
2605 outcomes,
2606 }
2607}
2608
2609fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2611 match stream {
2612 StreamKind::Projects => {
2613 let shaped = shape_projects(capabilities, filters);
2614 HitShape {
2615 stream,
2616 tasks: TaskQuery::default(),
2617 projects: shaped.pushed,
2618 local_tasks: LocalTasks::default(),
2619 local_projects: shaped.local,
2620 outcomes: shaped.outcomes,
2621 }
2622 }
2623 StreamKind::Items | StreamKind::Tasks => {
2624 let shaped = shape_tasks(
2625 capabilities,
2626 filters,
2627 &ProjectFilter::Any,
2628 &[],
2629 None,
2630 &[],
2631 None,
2632 );
2633 HitShape {
2634 stream,
2635 tasks: shaped.pushed,
2636 projects: ProjectQuery::default(),
2637 local_tasks: shaped.local,
2638 local_projects: LocalProjects::default(),
2639 outcomes: shaped.outcomes,
2640 }
2641 }
2642 }
2643}
2644
2645fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2652 if compensating {
2653 ceiling
2654 } else {
2655 budget.min(ceiling)
2656 }
2657}
2658
2659fn ceiling(source: &ResolvedSource) -> u32 {
2661 source.source().capabilities().max_page_size.max(1)
2662}
2663
2664async fn fetch_tasks(
2666 source: &ResolvedSource,
2667 shape: &TaskShape,
2668 start: &Resume,
2669 budget: u32,
2670 calls: &AtomicU32,
2671) -> Result<Fetched<Task>, SourceError> {
2672 let compensating = shape.local != LocalTasks::default();
2673 walk(
2674 start,
2675 budget,
2676 page_size(compensating, budget, ceiling(source)),
2677 |task| shape.local.keeps(task),
2678 |cursor, limit| async move {
2679 calls.fetch_add(1, Ordering::Relaxed);
2680 let request = PageRequest { cursor, limit };
2681 let page = source.source().query_tasks(&shape.pushed, &request).await?;
2682 match shape.local.commented_since {
2683 None => Ok(page),
2684 Some(since) => comment::commented_since(source, &shape.local, page, since).await,
2685 }
2686 },
2687 )
2688 .await
2689}
2690
2691async fn fetch_projects(
2693 source: &ResolvedSource,
2694 shape: &ProjectShape,
2695 start: &Resume,
2696 budget: u32,
2697 calls: &AtomicU32,
2698) -> Result<Fetched<Project>, SourceError> {
2699 let compensating = shape.local != LocalProjects::default();
2700 walk(
2701 start,
2702 budget,
2703 page_size(compensating, budget, ceiling(source)),
2704 |project| shape.local.keeps(project),
2705 |cursor, limit| async move {
2706 calls.fetch_add(1, Ordering::Relaxed);
2707 let request = PageRequest { cursor, limit };
2708 source
2709 .source()
2710 .query_projects(&shape.pushed, &request)
2711 .await
2712 },
2713 )
2714 .await
2715}
2716
2717async fn fetch_documents(
2719 source: &ResolvedSource,
2720 shape: &DocumentShape,
2721 start: &Resume,
2722 budget: u32,
2723 calls: &AtomicU32,
2724) -> Result<Fetched<Document>, SourceError> {
2725 let compensating = shape.local != LocalDocuments::default();
2726 walk(
2727 start,
2728 budget,
2729 page_size(compensating, budget, ceiling(source)),
2730 |document| shape.local.keeps(document),
2731 |cursor, limit| async move {
2732 calls.fetch_add(1, Ordering::Relaxed);
2733 let request = PageRequest { cursor, limit };
2734 source
2735 .source()
2736 .query_documents(&shape.pushed, &request)
2737 .await
2738 },
2739 )
2740 .await
2741}
2742
2743async fn fetch_labels(
2745 source: &ResolvedSource,
2746 start: &Resume,
2747 budget: u32,
2748 calls: &AtomicU32,
2749) -> Result<Fetched<Label>, SourceError> {
2750 walk(
2751 start,
2752 budget,
2753 page_size(false, budget, ceiling(source)),
2754 |_| true,
2755 |cursor, limit| async move {
2756 calls.fetch_add(1, Ordering::Relaxed);
2757 let request = PageRequest { cursor, limit };
2758 source.source().labels(&request).await
2759 },
2760 )
2761 .await
2762}
2763
2764async fn fetch_hits(
2766 source: &ResolvedSource,
2767 shape: &HitShape,
2768 start: &Resume,
2769 budget: u32,
2770 calls: &AtomicU32,
2771) -> Result<Fetched<Found>, SourceError> {
2772 let ceiling = ceiling(source);
2773 match shape.stream {
2774 StreamKind::Projects => {
2775 let compensating = shape.local_projects != LocalProjects::default();
2776 walk(
2777 start,
2778 budget,
2779 page_size(compensating, budget, ceiling),
2780 |found| match found {
2781 Found::Project(project) => shape.local_projects.keeps(project),
2782 Found::Task(_) => true,
2783 },
2784 |cursor, limit| async move {
2785 calls.fetch_add(1, Ordering::Relaxed);
2786 let request = PageRequest { cursor, limit };
2787 let page = source
2788 .source()
2789 .query_projects(&shape.projects, &request)
2790 .await?;
2791 Ok(Page {
2792 items: page.items.into_iter().map(Found::Project).collect(),
2793 next: page.next,
2794 })
2795 },
2796 )
2797 .await
2798 }
2799 StreamKind::Items | StreamKind::Tasks => {
2800 let compensating = shape.local_tasks != LocalTasks::default();
2801 walk(
2802 start,
2803 budget,
2804 page_size(compensating, budget, ceiling),
2805 |found| match found {
2806 Found::Task(task) => shape.local_tasks.keeps(task),
2807 Found::Project(_) => true,
2808 },
2809 |cursor, limit| async move {
2810 calls.fetch_add(1, Ordering::Relaxed);
2811 let request = PageRequest { cursor, limit };
2812 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2813 Ok(Page {
2814 items: page.items.into_iter().map(Found::Task).collect(),
2815 next: page.next,
2816 })
2817 },
2818 )
2819 .await
2820 }
2821 }
2822}
2823
2824async fn forward_edges(
2826 source: &ResolvedSource,
2827 entity: Entity,
2828 id: &NativeId,
2829 request: &PageRequest,
2830) -> Result<Page<DependencyEdge>, SourceError> {
2831 match entity {
2832 Entity::Task => {
2833 source
2834 .source()
2835 .task_dependencies(id, Direction::DependsOn, request)
2836 .await
2837 }
2838 Entity::Project => {
2839 source
2840 .source()
2841 .project_dependencies(id, Direction::DependsOn, request)
2842 .await
2843 }
2844 }
2845}
2846
2847#[expect(
2856 clippy::too_many_arguments,
2857 reason = "every argument is one axis of one walk — the source, the item, the \
2858 direction, which of its two graphs, whether the reverse is emulated, where \
2859 to resume, how many rows to return and where to count calls. Grouping them \
2860 into a struct would name the same eight values one indirection further from \
2861 the loop that reads them."
2862)]
2863async fn fetch_edges(
2864 source: &ResolvedSource,
2865 native: &NativeId,
2866 direction: Direction,
2867 entity: Entity,
2868 emulating: bool,
2869 start: &Resume,
2870 budget: u32,
2871 calls: &AtomicU32,
2872) -> Result<Fetched<DependencyEdge>, SourceError> {
2873 let ceiling = ceiling(source);
2874 if !emulating {
2875 return walk(
2876 start,
2877 budget,
2878 page_size(false, budget, ceiling),
2879 |_| true,
2880 |cursor, limit| async move {
2881 calls.fetch_add(1, Ordering::Relaxed);
2882 let request = PageRequest { cursor, limit };
2883 match entity {
2884 Entity::Task => {
2885 source
2886 .source()
2887 .task_dependencies(native, direction, &request)
2888 .await
2889 }
2890 Entity::Project => {
2891 source
2892 .source()
2893 .project_dependencies(native, direction, &request)
2894 .await
2895 }
2896 }
2897 },
2898 )
2899 .await;
2900 }
2901
2902 walk(
2903 start,
2904 budget,
2905 ceiling,
2906 |_| true,
2907 |cursor, limit| async move {
2908 calls.fetch_add(1, Ordering::Relaxed);
2909 let request = PageRequest { cursor, limit };
2910 let (ids, next) = match entity {
2911 Entity::Task => {
2912 let page = source
2913 .source()
2914 .query_tasks(&TaskQuery::default(), &request)
2915 .await?;
2916 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2917 (ids, page.next)
2918 }
2919 Entity::Project => {
2920 let page = source
2921 .source()
2922 .query_projects(&ProjectQuery::default(), &request)
2923 .await?;
2924 let ids: Vec<NativeId> =
2925 page.items.into_iter().map(|project| project.id).collect();
2926 (ids, page.next)
2927 }
2928 };
2929
2930 let mut edges = Vec::new();
2931 for id in ids {
2932 let mut inner: Option<Cursor> = None;
2933 loop {
2934 calls.fetch_add(1, Ordering::Relaxed);
2935 let request = PageRequest {
2936 cursor: inner.clone(),
2937 limit,
2938 };
2939 let page = forward_edges(source, entity, &id, &request).await?;
2940 fits(page.items.len(), limit)?;
2944 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2945 unrepeated(
2946 page.next.as_ref(),
2947 inner.as_ref(),
2948 "its forward edges were being scanned",
2949 )?;
2950 match page.next {
2951 Some(cursor) => inner = Some(cursor),
2952 None => break,
2953 }
2954 }
2955 }
2956
2957 Ok(Page { items: edges, next })
2958 },
2959 )
2960 .await
2961}