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, SecretResolver, SourceError, SourceName, StatusCategory, Task, TaskQuery,
39 TextFields, TextQuery,
40};
41use schemars::JsonSchema;
42use serde::{Deserialize, Serialize};
43
44use crate::GlobalId;
45use crate::config::Config;
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};
56pub(crate) use copy::malformed_links;
57pub use copy::{
58 BudgetSpent, CopyAction, CopyItems, CopyLink, CopyOutcome, CopyReport, CopyRequest, CopyScope,
59 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 paging: Paging,
229}
230
231#[derive(Debug, Clone)]
233pub struct ProjectRequest {
234 pub sources: Vec<SourceName>,
236 pub filters: Filters,
238 pub paging: Paging,
240}
241
242#[derive(Debug, Clone, Default, PartialEq)]
249pub struct DocumentFilters {
250 pub text: Option<TextQuery>,
252 pub labels: LabelFilter,
254}
255
256#[derive(Debug, Clone)]
258pub struct DocumentRequest {
259 pub sources: Vec<SourceName>,
261 pub filters: DocumentFilters,
263 pub project: ProjectSelector,
265 pub paging: Paging,
267}
268
269#[derive(Debug, Clone)]
271pub struct LabelRequest {
272 pub sources: Vec<SourceName>,
274 pub paging: Paging,
276}
277
278#[derive(Debug, Clone)]
280pub struct SearchRequest {
281 pub sources: Vec<SourceName>,
283 pub text: TextQuery,
285 pub kind: SearchKind,
287 pub paging: Paging,
289}
290
291#[derive(Debug, Clone)]
293pub struct DependencyRequest {
294 pub id: GlobalId,
296 pub direction: Direction,
298 pub paging: Paging,
300}
301
302#[derive(Debug, Clone, PartialEq, thiserror::Error)]
310pub enum EngineError {
311 #[error(
313 "no source named {name:?} is configured\n\
314 next: name one of the configured sources ({configured}), or add {name:?} under \
315 `sources` — `onetaskgraph sources list` shows what this configuration has."
316 )]
317 UnknownSource {
318 name: String,
320 configured: String,
322 },
323
324 #[error(
326 "{message}\n\
327 next: page with a token exactly as the previous page reported it, and against \
328 the same configuration — or drop `--page` to start the walk again."
329 )]
330 Token {
331 message: String,
333 },
334
335 #[error(
337 "no sources are configured\n\
338 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
339 prints what each plugin accepts."
340 )]
341 NoSources,
342
343 #[error(
345 "source {name} cannot be written: its plugin is {kind}, which has no write \
346 side\n\
347 next: copy into a source whose plugin can be written — `onetaskgraph sources \
348 list` reports each one's plugin."
349 )]
350 NotWritable {
351 name: String,
353 kind: String,
355 },
356
357 #[error(
364 "source {name} has no documents: its plugin is {kind}, which holds none\n\
365 next: name a source whose plugin has documents — `onetaskgraph sources list` \
366 reports each one's plugin and what it declares."
367 )]
368 NoDocuments {
369 name: String,
371 kind: String,
373 },
374
375 #[error(
380 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
381 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
382 list` reports each one's plugin."
383 )]
384 NoComments {
385 name: String,
387 kind: String,
389 },
390
391 #[error(
393 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
394 but not added to, edited or removed\n\
395 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
396 source whose plugin can be written — `onetaskgraph sources list` reports each \
397 one's plugin."
398 )]
399 CommentsNotWritable {
400 name: String,
402 kind: String,
404 },
405
406 #[error(
408 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
409 next: set the status in that source itself, or name a task of a source whose plugin \
410 can be written — `onetaskgraph sources list` reports each one's plugin."
411 )]
412 StatusNotWritable {
413 name: String,
415 kind: String,
417 },
418
419 #[error(
421 "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
422 write side\n\
423 next: set the key in that source itself, or name a {record} of a source whose plugin \
424 can be written — `onetaskgraph sources list` reports each one's plugin."
425 )]
426 MetadataNotWritable {
427 name: String,
429 kind: String,
431 record: MetadataRecord,
433 },
434
435 #[error(
437 "source {name} cannot write a priority: its plugin is {kind}, which has no write side\n\
438 next: set the priority in that source itself, or name a task of a source whose plugin \
439 can be written — `onetaskgraph sources list` reports each one's plugin."
440 )]
441 PriorityNotWritable {
442 name: String,
444 kind: String,
446 },
447
448 #[error(
450 "source {name} cannot write a task's content: its plugin is {kind}, which has no write \
451 side\n\
452 next: edit the content in that source itself, or name a task of a source whose plugin \
453 can be written — `onetaskgraph sources list` reports each one's plugin."
454 )]
455 ContentNotWritable {
456 name: String,
458 kind: String,
460 },
461
462 #[error(
464 "source {name} cannot update a task: its plugin is {kind}, which has no write side\n\
465 next: change the task in that source itself, or name a task of a source whose plugin \
466 can be written — `onetaskgraph sources list` reports each one's plugin."
467 )]
468 UpdateNotWritable {
469 name: String,
471 kind: String,
473 },
474
475 #[error(
477 "source {name} cannot create a {record}: its plugin is {kind}, which has no write side\n\
478 next: create it in a source whose plugin can be written — `onetaskgraph sources list` \
479 reports each one's plugin."
480 )]
481 NotCreatable {
482 name: String,
484 kind: String,
486 record: MetadataRecord,
488 },
489
490 #[error(
492 "source {name} cannot write a {record}'s rendering: its plugin is {kind}, which has no \
493 write side\n\
494 next: render it with --dry-run to read the result, or regenerate a {record} of a source \
495 whose plugin can be written."
496 )]
497 RenderingNotWritable {
498 name: String,
500 kind: String,
502 record: RenderedRecord,
504 },
505
506 #[error(
509 "supply every required answer to regenerate {id}: {} unanswered, and {reason}\n\
510 next: answer {} with --var NAME=VALUE or an answers file (--answers FILE), or run \
511 interactively to be asked.",
512 names.join(", "),
513 if names.len() == 1 { "it" } else { "each" }
514 )]
515 MissingAnswers {
516 id: String,
518 names: Vec<String>,
520 reason: String,
523 },
524
525 #[error("{error}")]
527 Template {
528 error: crate::template::TemplateError,
530 },
531
532 #[error(
534 "{record} {id} has no stored template answers: {reason}\n\
535 next: regenerate it with every required answer (`onetaskgraph {record} render {id} \
536 --var NAME=VALUE`), or read its provenance with `onetaskgraph {record} show {id}`."
537 )]
538 NoStoredAnswers {
539 record: RenderedRecord,
541 id: String,
543 reason: String,
545 },
546
547 #[error(
549 "{record} {id} records no template it was rendered from, and none was given\n\
550 next: name one with --template FILE or --template-loader FILE."
551 )]
552 NoTemplate {
553 record: RenderedRecord,
555 id: String,
557 },
558
559 #[error(
562 "{record} {id} records a template entry this product did not write — {problem}\n\
563 next: regenerate it with --template FILE or --template-loader FILE and every required \
564 answer, which records a fresh entry."
565 )]
566 MalformedProvenance {
567 record: RenderedRecord,
569 id: String,
571 problem: String,
573 },
574
575 #[error(
577 "{record} {id} was rendered from {reference:?}, which is not a readable file, so it \
578 cannot be re-read; a recorded reference is never turned into a location\n\
579 next: supply the template with --template-loader FILE (a loader document naming what \
580 to render), or name a template file with --template FILE."
581 )]
582 TemplateNotAFile {
583 record: RenderedRecord,
585 id: String,
587 reference: String,
589 },
590
591 #[error(
596 "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
597 be written to it: its plugin is {kind}, which declares priority unsupported\n\
598 next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
599 reports what each declares — or set the task's priority to none first; a \
600 github-projects source holds one once its configuration sets priority_mapping."
601 )]
602 NoPriority {
603 name: String,
605 kind: String,
607 task: String,
609 priority: onetaskgraph_plugin_api::Priority,
611 },
612
613 #[error(
615 "no project with the id {id}\n\
616 next: check the id, or list what is there — `onetaskgraph project list` reports every \
617 project the configured sources hold."
618 )]
619 NoSuchProject {
620 id: String,
622 },
623
624 #[error(
626 "no document with the id {id}\n\
627 next: check the id, or list what is there — `onetaskgraph document list` reports every \
628 document the configured sources hold."
629 )]
630 NoSuchDocument {
631 id: String,
633 },
634
635 #[error(
637 "no task with the id {id}\n\
638 next: check the id, or list what is there — `onetaskgraph task list` reports every \
639 task the configured sources hold."
640 )]
641 NoSuchTask {
642 id: String,
644 },
645
646 #[error(
648 "task {task} has no comment with the id {comment}\n\
649 next: list its comments — `onetaskgraph task comment list {task}` reports each \
650 one's id."
651 )]
652 NoSuchComment {
653 task: String,
655 comment: String,
657 },
658
659 #[error(
664 "source {name} could not be built: {error}\n\
665 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
666 command again."
667 )]
668 SourceUnavailable {
669 name: String,
671 error: SourceError,
673 },
674
675 #[error(
681 "source {name} could not do it: {error}\n\
682 next: fix what the source named above, then run the command again."
683 )]
684 SourceFailed {
685 name: String,
687 error: SourceError,
689 },
690
691 #[error(
693 "the destination source {name} could not be built: {error}\n\
694 next: fix that source — `onetaskgraph sources list` reports its state — then \
695 copy again."
696 )]
697 DestinationUnavailable {
698 name: String,
700 error: SourceError,
702 },
703
704 #[error(
706 "no item with the id {id}\n\
707 next: check the id, or list what is there — `onetaskgraph task list` and \
708 `onetaskgraph project list` report what the configured sources hold."
709 )]
710 NoSuchItem {
711 id: String,
713 },
714
715 #[error(
720 "{item} was copied from {origin}, which that destination no longer holds\n\
721 next: re-run with --recreate to create a new item there instead, or restore \
722 {origin}."
723 )]
724 StaleOrigin {
725 item: String,
727 origin: String,
729 },
730
731 #[error(
737 "{item} was last copied to {link}, which that destination no longer holds\n\
738 next: re-run with --recreate to create a new item there instead, or restore \
739 {link}."
740 )]
741 StaleLink {
742 item: String,
744 link: String,
746 },
747
748 #[error(
754 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
755 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
756 {id} on its own with `onetaskgraph task copy`.",
757 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
758 )]
759 NotAMember {
760 id: GlobalId,
762 projects: Vec<GlobalId>,
764 },
765
766 #[error(
775 "{item} depends on {member}, which this copy was not told to carry and which records \
776 no origin in {destination}\n\
777 next: name {member} with --member as well, record its {destination} id at \
778 onetaskgraph.origin, or copy the whole project without --member."
779 )]
780 UnrecordedMember {
781 item: GlobalId,
783 member: GlobalId,
785 destination: SourceName,
787 },
788
789 #[error(
795 "source {name} could not do it: {error}\n\
796 next: fix what the source named above, then copy again."
797 )]
798 SourceRefused {
799 name: String,
801 error: SourceError,
803 },
804
805 #[error(
813 "the copy failed and could not be undone.\n\
814 it failed because: {error}\n\
815 it could not be undone because: {refusal}\n\
816 so these still hold what it wrote: {left_behind}\n\
817 next: remove or put back those items, then copy again."
818 )]
819 CopyNotUndone {
820 error: Box<EngineError>,
822 left_behind: LeftBehind,
828 refusal: SourceError,
830 },
831}
832
833#[derive(Debug, Clone, PartialEq)]
842pub struct LeftBehind {
843 first: GlobalId,
845 rest: Vec<GlobalId>,
847}
848
849impl LeftBehind {
850 #[must_use]
852 pub fn new(first: GlobalId) -> Self {
853 Self {
854 first,
855 rest: Vec::new(),
856 }
857 }
858
859 pub fn push(&mut self, id: GlobalId) {
861 self.rest.push(id);
862 }
863
864 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
866 std::iter::once(&self.first).chain(self.rest.iter())
867 }
868}
869
870impl std::fmt::Display for LeftBehind {
871 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
873 write!(formatter, "{}", self.first)?;
874 for id in &self.rest {
875 write!(formatter, ", {id}")?;
876 }
877 Ok(())
878 }
879}
880
881pub enum ConfiguredSource {
888 Ready(ResolvedSource),
890 Unavailable(UnavailableSource),
892}
893
894impl ConfiguredSource {
895 #[must_use]
897 pub fn name(&self) -> &SourceName {
898 match self {
899 Self::Ready(source) => source.name(),
900 Self::Unavailable(source) => source.name(),
901 }
902 }
903}
904
905pub struct Engine {
907 sources: Vec<ConfiguredSource>,
909 selection: Vec<SourceName>,
911}
912
913impl Engine {
914 #[must_use]
922 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
923 let (ready, unavailable) = resolve_available(config, secrets);
924 Self::new(
925 ready
926 .into_iter()
927 .map(ConfiguredSource::Ready)
928 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
929 .collect(),
930 config.selected_sources(),
931 )
932 }
933
934 #[must_use]
937 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
938 Self { sources, selection }
939 }
940
941 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
943 self.sources.iter().filter_map(|source| match source {
944 ConfiguredSource::Ready(ready) => Some(ready),
945 ConfiguredSource::Unavailable(_) => None,
946 })
947 }
948
949 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
951 self.sources.iter().filter_map(|source| match source {
952 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
953 ConfiguredSource::Ready(_) => None,
954 })
955 }
956
957 #[must_use]
959 pub fn listing(&self) -> Vec<SourceListing> {
960 let mut listings: Vec<SourceListing> = self
961 .ready()
962 .map(|source| SourceListing {
963 source: source.name().clone(),
964 kind: source.kind().to_owned(),
965 state: SourceState::Available {
966 capabilities: source.source().capabilities(),
967 },
968 })
969 .chain(self.unavailable().map(|source| SourceListing {
970 source: source.name().clone(),
971 kind: source.kind().to_owned(),
972 state: SourceState::Unavailable {
973 error: source.error().clone(),
974 },
975 }))
976 .collect();
977 listings.sort_by(|left, right| left.source.cmp(&right.source));
978 listings
979 }
980
981 #[must_use]
987 pub fn has(&self, name: &SourceName) -> bool {
988 self.sources.iter().any(|source| source.name() == name)
989 }
990
991 pub async fn tasks(
999 &self,
1000 request: &TaskRequest,
1001 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1002 let mut names = self.resolve_selection(&request.sources)?;
1003 if let ProjectSelector::Qualified(id) = &request.project {
1008 self.known(&id.source)?;
1009 names.retain(|name| name == &id.source);
1010 }
1011 let query = shape(
1012 "task-list",
1013 &names,
1014 &(
1015 &request.filters,
1016 &request.project,
1017 &request.priorities,
1018 &request.commented_since,
1019 &request.metadata,
1020 &request.origin,
1021 ),
1022 );
1023 let states = resumption(
1024 self,
1025 request.paging.token.as_ref(),
1026 &[StreamKind::Items],
1027 &query,
1028 )?;
1029 let budget = request.paging.limit.get();
1030
1031 let mut answer = Answer::new();
1032 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1033
1034 let shapes: Vec<TaskShape> = ready
1035 .iter()
1036 .map(|source| {
1037 shape_tasks(
1038 &source.source().capabilities(),
1039 &request.filters,
1040 &project_filter(&request.project),
1041 &request.priorities,
1042 request.commented_since,
1043 &request.metadata,
1044 request.origin.as_ref(),
1045 )
1046 })
1047 .collect();
1048 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1049 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1050
1051 let walks = ready
1052 .iter()
1053 .enumerate()
1054 .map(|(index, source)| {
1055 fetch_tasks(
1056 source,
1057 &shapes[index],
1058 &starts[index],
1059 budget,
1060 &counters[index],
1061 )
1062 })
1063 .collect();
1064
1065 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1066 answer.finish(
1067 streams,
1068 budget,
1069 owed(&states),
1070 &query,
1071 |name, task: Task| {
1072 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1073 },
1074 )
1075 }
1076
1077 pub async fn projects(
1083 &self,
1084 request: &ProjectRequest,
1085 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1086 let names = self.resolve_selection(&request.sources)?;
1087 let query = shape("project-list", &names, &request.filters);
1088 let states = resumption(
1089 self,
1090 request.paging.token.as_ref(),
1091 &[StreamKind::Items],
1092 &query,
1093 )?;
1094 let budget = request.paging.limit.get();
1095
1096 let mut answer = Answer::new();
1097 let mut with_projects = Vec::new();
1103 for source in answer.split(self, &names) {
1104 if source.source().capabilities().projects.is_native() {
1105 with_projects.push(source);
1106 } else {
1107 answer.unreachable_predicate(source, Predicate::Project);
1108 }
1109 }
1110 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1111
1112 let shapes: Vec<ProjectShape> = ready
1113 .iter()
1114 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1115 .collect();
1116 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1117 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1118
1119 let walks = ready
1120 .iter()
1121 .enumerate()
1122 .map(|(index, source)| {
1123 fetch_projects(
1124 source,
1125 &shapes[index],
1126 &starts[index],
1127 budget,
1128 &counters[index],
1129 )
1130 })
1131 .collect();
1132
1133 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1134 answer.finish(
1135 streams,
1136 budget,
1137 owed(&states),
1138 &query,
1139 |name, project: Project| Qualified {
1140 id: GlobalId::new(name.clone(), project.id.clone()),
1141 item: project,
1142 },
1143 )
1144 }
1145
1146 pub async fn documents(
1158 &self,
1159 request: &DocumentRequest,
1160 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1161 let mut names = self.resolve_selection(&request.sources)?;
1162 if let ProjectSelector::Qualified(id) = &request.project {
1165 self.known(&id.source)?;
1166 names.retain(|name| name == &id.source);
1167 }
1168 let query = shape(
1169 "document-list",
1170 &names,
1171 &(&request.filters, &request.project),
1172 );
1173 let states = resumption(
1174 self,
1175 request.paging.token.as_ref(),
1176 &[StreamKind::Items],
1177 &query,
1178 )?;
1179 let budget = request.paging.limit.get();
1180
1181 let mut answer = Answer::new();
1182 let mut with_documents = Vec::new();
1183 for source in answer.split(self, &names) {
1184 if source.source().capabilities().documents.is_native() {
1185 with_documents.push(source);
1186 } else {
1187 answer.unreachable_predicate(source, Predicate::Document);
1188 }
1189 }
1190 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1191
1192 let shapes: Vec<DocumentShape> = ready
1193 .iter()
1194 .map(|source| {
1195 shape_documents(
1196 &source.source().capabilities(),
1197 &request.filters,
1198 &project_filter(&request.project),
1199 )
1200 })
1201 .collect();
1202 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1203 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1204
1205 let walks = ready
1206 .iter()
1207 .enumerate()
1208 .map(|(index, source)| {
1209 fetch_documents(
1210 source,
1211 &shapes[index],
1212 &starts[index],
1213 budget,
1214 &counters[index],
1215 )
1216 })
1217 .collect();
1218
1219 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1220 answer.finish(
1221 streams,
1222 budget,
1223 owed(&states),
1224 &query,
1225 |name, document: Document| Qualified {
1226 id: GlobalId::new(name.clone(), document.id.clone()),
1227 item: document,
1228 },
1229 )
1230 }
1231
1232 pub async fn labels(
1238 &self,
1239 request: &LabelRequest,
1240 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1241 let names = self.resolve_selection(&request.sources)?;
1242 let query = shape("label-list", &names, &());
1243 let states = resumption(
1244 self,
1245 request.paging.token.as_ref(),
1246 &[StreamKind::Items],
1247 &query,
1248 )?;
1249 let budget = request.paging.limit.get();
1250
1251 let mut answer = Answer::new();
1252 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1253 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1254 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1255
1256 let walks = ready
1257 .iter()
1258 .enumerate()
1259 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1260 .collect();
1261
1262 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1263 answer.finish(
1264 streams,
1265 budget,
1266 owed(&states),
1267 &query,
1268 |name, label: Label| Qualified {
1269 id: GlobalId::new(name.clone(), label.id.clone()),
1270 item: label,
1271 },
1272 )
1273 }
1274
1275 pub async fn search(
1281 &self,
1282 request: &SearchRequest,
1283 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1284 let names = self.resolve_selection(&request.sources)?;
1285 let reads: &[StreamKind] = match request.kind {
1289 SearchKind::Tasks => &[StreamKind::Tasks],
1290 SearchKind::Projects => &[StreamKind::Projects],
1291 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1292 };
1293 let query = shape("search", &names, &request.text);
1298 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1299 let budget = request.paging.limit.get();
1300 let filters = Filters {
1301 text: Some(request.text.clone()),
1302 ..Filters::default()
1303 };
1304
1305 let mut answer = Answer::new();
1306
1307 let mut ready = Vec::new();
1310 let mut kinds = Vec::new();
1311 let mut starts = Vec::new();
1312 for source in answer.split(self, &names) {
1313 let mut streams = Vec::new();
1314 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1315 streams.push(StreamKind::Tasks);
1316 }
1317 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1318 if source.source().capabilities().projects.is_native() {
1319 streams.push(StreamKind::Projects);
1320 } else {
1321 answer.unreachable_predicate(source, Predicate::Project);
1322 }
1323 }
1324 for stream in streams {
1325 if let Some(resume) = resume_at(&states, source.name(), stream) {
1326 ready.push(source);
1327 kinds.push(stream);
1328 starts.push(resume);
1329 }
1330 }
1331 }
1332
1333 let shapes: Vec<HitShape> = ready
1334 .iter()
1335 .zip(kinds.iter())
1336 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1337 .collect();
1338 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1339 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1340
1341 let walks = ready
1342 .iter()
1343 .enumerate()
1344 .map(|(index, source)| {
1345 fetch_hits(
1346 source,
1347 &shapes[index],
1348 &starts[index],
1349 budget,
1350 &counters[index],
1351 )
1352 })
1353 .collect();
1354
1355 let streams =
1356 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1357 answer.finish(
1358 streams,
1359 budget,
1360 owed(&states),
1361 &query,
1362 |name, found: Found| match found {
1363 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1364 GlobalId::new(name.clone(), task.id.clone()),
1365 task,
1366 )),
1367 Found::Project(project) => SearchHit::Project(Qualified {
1368 id: GlobalId::new(name.clone(), project.id.clone()),
1369 item: project,
1370 }),
1371 },
1372 )
1373 }
1374
1375 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1382 let name = self.known(&id.source)?;
1383 let mut answer = Answer::new();
1384 let selected = answer.split(self, std::slice::from_ref(&name));
1385 let Some(source) = selected.first() else {
1386 return answer.nothing();
1387 };
1388 let found = source.source().get_task(&id.native).await;
1389 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1390 answer.one(source, found, |task| {
1391 delivery::qualified_task(qualified, task)
1392 })
1393 }
1394
1395 pub async fn project(
1401 &self,
1402 id: &GlobalId,
1403 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1404 let name = self.known(&id.source)?;
1405 let mut answer = Answer::new();
1406 let selected = answer.split(self, std::slice::from_ref(&name));
1407 let Some(source) = selected.first() else {
1408 return answer.nothing();
1409 };
1410 let found = source.source().get_project(&id.native).await;
1411 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1412 answer.one(source, found, |project| Qualified {
1413 id: qualified,
1414 item: project,
1415 })
1416 }
1417
1418 pub async fn document(
1428 &self,
1429 id: &GlobalId,
1430 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1431 let name = self.known(&id.source)?;
1432 let mut answer = Answer::new();
1433 let selected = answer.split(self, std::slice::from_ref(&name));
1434 let Some(source) = selected.first() else {
1435 return answer.nothing();
1436 };
1437 if !source.source().capabilities().documents.is_native() {
1438 answer.unreachable_predicate(source, Predicate::Document);
1439 return answer.nothing();
1440 }
1441 let found = source.source().get_document(&id.native).await;
1442 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1443 answer.one(source, found, |document| Qualified {
1444 id: qualified,
1445 item: document,
1446 })
1447 }
1448
1449 pub async fn task_dependencies(
1456 &self,
1457 request: &DependencyRequest,
1458 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1459 self.dependencies(request, Entity::Task).await
1460 }
1461
1462 pub async fn project_dependencies(
1468 &self,
1469 request: &DependencyRequest,
1470 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1471 self.dependencies(request, Entity::Project).await
1472 }
1473
1474 async fn dependencies(
1477 &self,
1478 request: &DependencyRequest,
1479 entity: Entity,
1480 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1481 let name = self.known(&request.id.source)?;
1482 let query = shape(
1483 "dependencies",
1484 std::slice::from_ref(&name),
1485 &(entity, &request.id.native, request.direction),
1486 );
1487 let states = resumption(
1488 self,
1489 request.paging.token.as_ref(),
1490 &[StreamKind::Items],
1491 &query,
1492 )?;
1493 let budget = request.paging.limit.get();
1494
1495 let mut answer = Answer::new();
1496 let (ready, starts) = walking(
1497 answer.split(self, std::slice::from_ref(&name)),
1498 &states,
1499 StreamKind::Items,
1500 );
1501 let Some(source) = ready.first() else {
1502 return answer.nothing();
1503 };
1504
1505 let capabilities = source.source().capabilities();
1506 let support = match entity {
1507 Entity::Task => capabilities.task_dependencies,
1508 Entity::Project => capabilities.project_dependencies,
1509 };
1510 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1514 let mut outcomes = Outcomes::default();
1515 if request.direction == Direction::DependedOnBy {
1516 if emulating {
1517 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1518 } else {
1519 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1520 }
1521 }
1522
1523 let counters = vec![AtomicU32::new(0)];
1524 let walked = fetch_edges(
1525 source,
1526 &request.id.native,
1527 request.direction,
1528 entity,
1529 emulating,
1530 &starts[0],
1531 budget,
1532 &counters[0],
1533 )
1534 .await;
1535
1536 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1537 answer.finish(
1538 streams,
1539 budget,
1540 owed(&states),
1541 &query,
1542 |name, edge: DependencyEdge| QualifiedEdge {
1543 from: qualify_endpoint(name, edge.from),
1544 to: qualify_endpoint(name, edge.to),
1545 kind: edge.kind,
1546 },
1547 )
1548 }
1549
1550 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1552 if asked.is_empty() {
1553 if self.selection.is_empty() {
1554 return Err(EngineError::NoSources);
1555 }
1556 return Ok(self.selection.clone());
1557 }
1558 asked.iter().map(|name| self.known(name)).collect()
1559 }
1560
1561 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1563 if self.has(name) {
1564 return Ok(name.clone());
1565 }
1566 if self.sources.is_empty() {
1567 return Err(EngineError::NoSources);
1568 }
1569 Err(EngineError::UnknownSource {
1570 name: name.to_string(),
1571 configured: self
1572 .listing()
1573 .iter()
1574 .map(|listing| listing.source.to_string())
1575 .collect::<Vec<_>>()
1576 .join(", "),
1577 })
1578 }
1579}
1580
1581fn qualify_endpoint(
1582 source: &SourceName,
1583 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1584) -> QualifiedEndpoint {
1585 let kind = endpoint.kind;
1586 let is_qualified = endpoint.is_qualified();
1587 let endpoint_id = endpoint.into_id();
1588 QualifiedEndpoint {
1589 id: if is_qualified {
1590 endpoint_id
1591 .parse()
1592 .expect("plugin-api validates qualified dependency endpoints")
1593 } else {
1594 GlobalId::new(source.clone(), NativeId(endpoint_id))
1595 },
1596 kind,
1597 }
1598}
1599
1600#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1602enum Entity {
1603 Task,
1605 Project,
1607}
1608
1609enum Found {
1611 Task(Task),
1613 Project(Project),
1615}
1616
1617#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1619enum Outcome {
1620 PushedDown,
1622 AppliedLocally,
1624 Emulated,
1626 Unavailable,
1628}
1629
1630#[derive(Debug, Clone, Default, PartialEq)]
1642struct Outcomes(BTreeMap<Predicate, Outcome>);
1643
1644impl Outcomes {
1645 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1651 self.0.insert(predicate, outcome);
1652 }
1653
1654 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1657 for predicate in predicates {
1658 self.record(predicate, outcome);
1659 }
1660 }
1661
1662 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1664 self.0
1665 .iter()
1666 .filter(|(_, recorded)| **recorded == outcome)
1667 .map(|(predicate, _)| *predicate)
1668 .collect()
1669 }
1670}
1671
1672struct TaskShape {
1674 pushed: TaskQuery,
1676 local: LocalTasks,
1678 outcomes: Outcomes,
1680}
1681
1682struct ProjectShape {
1684 pushed: ProjectQuery,
1686 local: LocalProjects,
1688 outcomes: Outcomes,
1690}
1691
1692struct DocumentShape {
1694 pushed: DocumentQuery,
1696 local: LocalDocuments,
1698 outcomes: Outcomes,
1700}
1701
1702struct HitShape {
1704 stream: StreamKind,
1706 tasks: TaskQuery,
1708 projects: ProjectQuery,
1710 local_tasks: LocalTasks,
1712 local_projects: LocalProjects,
1714 outcomes: Outcomes,
1716}
1717
1718struct Answer {
1724 plans: Vec<SourcePlan>,
1726 errors: Vec<SourceFailure>,
1728}
1729
1730impl Answer {
1731 fn new() -> Self {
1732 Self {
1733 plans: Vec::new(),
1734 errors: Vec::new(),
1735 }
1736 }
1737
1738 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1743 let mut selected = Vec::new();
1744 for name in names {
1745 match engine.sources.iter().find(|source| source.name() == name) {
1746 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1747 Some(ConfiguredSource::Unavailable(source)) => {
1748 self.errors.push(source.failure());
1749 }
1750 None => {}
1751 }
1752 }
1753 selected
1754 }
1755
1756 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1758 let mut outcomes = Outcomes::default();
1759 outcomes.record(predicate, Outcome::Unavailable);
1760 self.plans.push(plan_for(source, outcomes, 0));
1761 }
1762
1763 fn collect<T>(
1765 &mut self,
1766 ready: &[&ResolvedSource],
1767 walked: Vec<Result<Fetched<T>, SourceError>>,
1768 counters: &[AtomicU32],
1769 outcomes: Vec<Outcomes>,
1770 ) -> Vec<Stream<T>> {
1771 let kinds = vec![StreamKind::Items; ready.len()];
1772 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1773 }
1774
1775 fn collect_streams<T>(
1777 &mut self,
1778 ready: &[&ResolvedSource],
1779 kinds: &[StreamKind],
1780 walked: Vec<Result<Fetched<T>, SourceError>>,
1781 counters: &[AtomicU32],
1782 outcomes: Vec<Outcomes>,
1783 ) -> Vec<Stream<T>> {
1784 let mut streams = Vec::new();
1785 for (index, result) in walked.into_iter().enumerate() {
1786 let source = ready[index];
1787 let pages = counters[index].load(Ordering::Relaxed);
1788 self.plans
1789 .push(plan_for(source, outcomes[index].clone(), pages));
1790 match result {
1791 Ok(fetched) => streams.push(Stream {
1792 source: source.name().clone(),
1793 kind: kinds[index],
1794 fetched,
1795 }),
1796 Err(error) => self.errors.push(SourceFailure {
1799 source: source.name().clone(),
1800 error,
1801 }),
1802 }
1803 }
1804 streams
1805 }
1806
1807 fn one<T, U>(
1809 mut self,
1810 source: &ResolvedSource,
1811 found: Result<Option<T>, SourceError>,
1812 qualify: impl FnOnce(T) -> U,
1813 ) -> Result<QueryResponse<U>, EngineError> {
1814 self.plans.push(plan_for(source, Outcomes::default(), 1));
1815 let items = match found {
1816 Ok(Some(item)) => vec![qualify(item)],
1817 Ok(None) => Vec::new(),
1818 Err(error) => {
1819 self.errors.push(SourceFailure {
1820 source: source.name().clone(),
1821 error,
1822 });
1823 Vec::new()
1824 }
1825 };
1826 Ok(QueryResponse {
1827 items,
1828 next: None,
1829 plan: QueryPlan {
1830 per_source: merge_plans(self.plans),
1831 },
1832 errors: self.errors,
1833 })
1834 }
1835
1836 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1838 Ok(QueryResponse {
1839 items: Vec::new(),
1840 next: None,
1841 plan: QueryPlan {
1842 per_source: merge_plans(self.plans),
1843 },
1844 errors: self.errors,
1845 })
1846 }
1847
1848 fn finish<T, U>(
1853 self,
1854 streams: Vec<Stream<T>>,
1855 budget: u32,
1856 first: Option<&Owed>,
1857 query: &str,
1858 qualify: impl Fn(&SourceName, T) -> U,
1859 ) -> Result<QueryResponse<U>, EngineError> {
1860 let (rows, states, owed) = merge(streams, budget, first);
1861 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1862 Ok(QueryResponse {
1863 items: rows
1864 .into_iter()
1865 .map(|(name, item)| qualify(&name, item))
1866 .collect(),
1867 next,
1868 plan: QueryPlan {
1869 per_source: merge_plans(self.plans),
1870 },
1871 errors: self.errors,
1872 })
1873 }
1874}
1875
1876fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1879 SourcePlan {
1880 source: source.name().clone(),
1881 kind: source.kind().to_owned(),
1882 pushed_down: outcomes.with(Outcome::PushedDown),
1883 applied_locally: outcomes.with(Outcome::AppliedLocally),
1884 emulated: outcomes.with(Outcome::Emulated),
1885 unavailable: outcomes.with(Outcome::Unavailable),
1886 pages_fetched: pages,
1887 }
1888}
1889
1890fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1895 let mut merged: Vec<SourcePlan> = Vec::new();
1896 for plan in plans {
1897 if let Some(existing) = merged
1898 .iter_mut()
1899 .find(|existing| existing.source == plan.source)
1900 {
1901 existing.pushed_down.extend(plan.pushed_down);
1902 existing.applied_locally.extend(plan.applied_locally);
1903 existing.emulated.extend(plan.emulated);
1904 existing.unavailable.extend(plan.unavailable);
1905 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1906 for list in [
1907 &mut existing.pushed_down,
1908 &mut existing.applied_locally,
1909 &mut existing.emulated,
1910 &mut existing.unavailable,
1911 ] {
1912 list.sort_unstable();
1913 list.dedup();
1914 }
1915 } else {
1916 merged.push(plan);
1917 }
1918 }
1919 merged
1920}
1921
1922fn walking<'a>(
1928 selected: Vec<&'a ResolvedSource>,
1929 states: &Option<Resumption>,
1930 kind: StreamKind,
1931) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1932 let mut ready = Vec::new();
1933 let mut starts = Vec::new();
1934 for source in selected {
1935 if let Some(resume) = resume_at(states, source.name(), kind) {
1936 ready.push(source);
1937 starts.push(resume);
1938 }
1939 }
1940 (ready, starts)
1941}
1942
1943fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1961 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1962 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1963}
1964
1965fn fingerprint(text: &str) -> String {
1967 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1968 for byte in text.as_bytes() {
1969 hash ^= u64::from(*byte);
1970 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1971 }
1972 format!("{hash:016x}")
1973}
1974
1975fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1980 document.as_ref()?.owed.as_ref()
1981}
1982
1983fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1985 match states {
1986 None => Some(Resume::default()),
1987 Some(document) => document
1988 .streams
1989 .iter()
1990 .find(|state| &state.source == source && state.stream == kind)
1991 .map(|state| state.resume.clone()),
1992 }
1993}
1994
1995fn resumption(
2026 engine: &Engine,
2027 token: Option<&PageToken>,
2028 reads: &[StreamKind],
2029 query: &str,
2030) -> Result<Option<Resumption>, EngineError> {
2031 let Some(document) = token.map(PageToken::decode) else {
2032 return Ok(None);
2033 };
2034
2035 if document.query != query {
2042 return Err(EngineError::Token {
2043 message: "this page token was written by a different query — resume the walk it \
2044 came from, or drop --page to start this one from the beginning"
2045 .to_owned(),
2046 });
2047 }
2048 let states = &document.streams;
2049
2050 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2051 for state in states {
2052 if !reads.contains(&state.stream) {
2053 return Err(EngineError::Token {
2054 message: format!(
2055 "this page token resumes {}, which this command does not read — it \
2056 was written by a different query",
2057 state.stream.describe()
2058 ),
2059 });
2060 }
2061 let ceiling = engine
2062 .ready()
2063 .find(|source| source.name() == &state.source)
2064 .map(ceiling);
2065 if ceiling.is_none() && !engine.has(&state.source) {
2066 return Err(EngineError::Token {
2067 message: format!(
2068 "this page token resumes a source called {:?}, which this \
2069 configuration does not have",
2070 state.source.as_str()
2071 ),
2072 });
2073 }
2074 if let Some(ceiling) = ceiling
2075 && state.resume.skip >= ceiling
2076 {
2077 return Err(EngineError::Token {
2078 message: format!(
2079 "this page token resumes {} rows into a page of source {:?}, which \
2080 serves at most {ceiling}",
2081 state.resume.skip,
2082 state.source.as_str()
2083 ),
2084 });
2085 }
2086 if seen.contains(&(&state.source, state.stream)) {
2087 return Err(EngineError::Token {
2088 message: format!(
2089 "this page token gives source {:?} two places to resume from",
2090 state.source.as_str()
2091 ),
2092 });
2093 }
2094 seen.push((&state.source, state.stream));
2095 }
2096
2097 if let Some(owed) = &document.owed
2102 && !document
2103 .streams
2104 .iter()
2105 .any(|state| state.source == owed.source && state.stream == owed.stream)
2106 {
2107 return Err(EngineError::Token {
2108 message: format!(
2109 "this page token owes the next row to a stream it does not resume, \
2110 {:?}'s {}",
2111 owed.source.as_str(),
2112 owed.stream.describe()
2113 ),
2114 });
2115 }
2116
2117 Ok(Some(document))
2118}
2119
2120fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2126 match selector {
2127 ProjectSelector::Any => ProjectFilter::Any,
2128 ProjectSelector::Orphans => ProjectFilter::Orphans,
2129 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2130 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2131 }
2132}
2133
2134fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2136 match fields {
2137 TextFields::Title => vec![Predicate::SearchTitle],
2138 TextFields::Content => vec![Predicate::SearchContent],
2139 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2140 }
2141}
2142
2143fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2150 match fields {
2151 TextFields::Title => capabilities.search_title.is_native(),
2152 TextFields::Content => capabilities.search_content.is_native(),
2153 TextFields::TitleOrContent => {
2154 capabilities.search_title.is_native() && capabilities.search_content.is_native()
2155 }
2156 }
2157}
2158
2159fn shape_tasks(
2161 capabilities: &Capabilities,
2162 filters: &Filters,
2163 project: &ProjectFilter,
2164 priorities: &[Priority],
2165 commented_since: Option<DateTime<Utc>>,
2166 metadata: &[MetadataMatch],
2167 origin: Option<&GlobalId>,
2168) -> TaskShape {
2169 let mut pushed = TaskQuery::default();
2170 let mut local = LocalTasks::default();
2171 let mut outcomes = Outcomes::default();
2172
2173 if !filters.labels.is_empty() {
2174 if capabilities.filter_by_label.is_native() {
2175 pushed.labels = filters.labels.clone();
2176 outcomes.record(Predicate::Label, Outcome::PushedDown);
2177 } else {
2178 local.labels = Some(filters.labels.clone());
2179 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2180 }
2181 }
2182 if !filters.statuses.is_empty() {
2183 if capabilities.filter_by_status.is_native() {
2184 pushed.statuses.clone_from(&filters.statuses);
2185 outcomes.record(Predicate::Status, Outcome::PushedDown);
2186 } else {
2187 local.statuses.clone_from(&filters.statuses);
2188 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2189 }
2190 }
2191 if !priorities.is_empty() {
2192 if capabilities.filter_by_priority.is_native() {
2193 pushed.priorities = priorities.to_vec();
2194 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2195 } else {
2196 local.priorities = priorities.to_vec();
2197 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2198 }
2199 }
2200 if let Some(since) = commented_since {
2201 if capabilities.filter_by_comment_activity.is_native() {
2202 pushed.commented_since = Some(since);
2203 outcomes.record(Predicate::CommentedSince, Outcome::PushedDown);
2204 } else {
2205 local.commented_since = Some(since);
2206 outcomes.record(Predicate::CommentedSince, Outcome::AppliedLocally);
2207 }
2208 }
2209 if !metadata.is_empty() {
2210 if capabilities.filter_by_metadata.is_native() {
2211 pushed.metadata = metadata.to_vec();
2212 outcomes.record(Predicate::Metadata, Outcome::PushedDown);
2213 } else {
2214 local.metadata = metadata.to_vec();
2215 outcomes.record(Predicate::Metadata, Outcome::AppliedLocally);
2216 }
2217 }
2218 if let Some(origin) = origin {
2219 if capabilities.filter_by_origin.is_native() {
2220 pushed.origin = Some(origin.to_string());
2223 outcomes.record(Predicate::Origin, Outcome::PushedDown);
2224 } else {
2225 local.origin = Some(origin.clone());
2226 outcomes.record(Predicate::Origin, Outcome::AppliedLocally);
2227 }
2228 }
2229 if let Some(text) = &filters.text {
2230 let predicates = text_predicates(text.fields);
2231 if searches_natively(capabilities, text.fields) {
2232 pushed.text = Some(text.clone());
2233 outcomes.record_all(predicates, Outcome::PushedDown);
2234 } else {
2235 local.text = Some(text.clone());
2236 outcomes.record_all(predicates, Outcome::AppliedLocally);
2237 }
2238 }
2239 match project {
2240 ProjectFilter::Any => {}
2241 ProjectFilter::Orphans => {
2242 if capabilities.orphan_tasks.is_native() {
2243 pushed.project = ProjectFilter::Orphans;
2244 outcomes.record(Predicate::Project, Outcome::PushedDown);
2245 } else {
2246 local.project = Some(ProjectFilter::Orphans);
2247 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2248 }
2249 }
2250 ProjectFilter::Is(id) => {
2251 if capabilities.projects.is_native() {
2252 pushed.project = ProjectFilter::Is(id.clone());
2253 outcomes.record(Predicate::Project, Outcome::PushedDown);
2254 } else {
2255 local.project = Some(ProjectFilter::Is(id.clone()));
2256 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2257 }
2258 }
2259 }
2260
2261 TaskShape {
2262 pushed,
2263 local,
2264 outcomes,
2265 }
2266}
2267
2268fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2270 let mut pushed = ProjectQuery::default();
2271 let mut local = LocalProjects::default();
2272 let mut outcomes = Outcomes::default();
2273
2274 if !filters.labels.is_empty() {
2275 if capabilities.filter_by_label.is_native() {
2276 pushed.labels = filters.labels.clone();
2277 outcomes.record(Predicate::Label, Outcome::PushedDown);
2278 } else {
2279 local.labels = Some(filters.labels.clone());
2280 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2281 }
2282 }
2283 if !filters.statuses.is_empty() {
2284 if capabilities.filter_by_status.is_native() {
2285 pushed.statuses.clone_from(&filters.statuses);
2286 outcomes.record(Predicate::Status, Outcome::PushedDown);
2287 } else {
2288 local.statuses.clone_from(&filters.statuses);
2289 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2290 }
2291 }
2292 if let Some(text) = &filters.text {
2293 let predicates = text_predicates(text.fields);
2294 if searches_natively(capabilities, text.fields) {
2295 pushed.text = Some(text.clone());
2296 outcomes.record_all(predicates, Outcome::PushedDown);
2297 } else {
2298 local.text = Some(text.clone());
2299 outcomes.record_all(predicates, Outcome::AppliedLocally);
2300 }
2301 }
2302
2303 ProjectShape {
2304 pushed,
2305 local,
2306 outcomes,
2307 }
2308}
2309
2310fn shape_documents(
2315 capabilities: &Capabilities,
2316 filters: &DocumentFilters,
2317 project: &ProjectFilter,
2318) -> DocumentShape {
2319 let mut pushed = DocumentQuery::default();
2320 let mut local = LocalDocuments::default();
2321 let mut outcomes = Outcomes::default();
2322
2323 if !filters.labels.is_empty() {
2324 if capabilities.filter_by_label.is_native() {
2325 pushed.labels = filters.labels.clone();
2326 outcomes.record(Predicate::Label, Outcome::PushedDown);
2327 } else {
2328 local.labels = Some(filters.labels.clone());
2329 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2330 }
2331 }
2332 if let Some(text) = &filters.text {
2333 let predicates = text_predicates(text.fields);
2334 if searches_natively(capabilities, text.fields) {
2335 pushed.text = Some(text.clone());
2336 outcomes.record_all(predicates, Outcome::PushedDown);
2337 } else {
2338 local.text = Some(text.clone());
2339 outcomes.record_all(predicates, Outcome::AppliedLocally);
2340 }
2341 }
2342 match project {
2343 ProjectFilter::Any => {}
2344 ProjectFilter::Orphans => {
2345 if capabilities.orphan_tasks.is_native() {
2346 pushed.project = ProjectFilter::Orphans;
2347 outcomes.record(Predicate::Project, Outcome::PushedDown);
2348 } else {
2349 local.project = Some(ProjectFilter::Orphans);
2350 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2351 }
2352 }
2353 ProjectFilter::Is(id) => {
2354 if capabilities.projects.is_native() {
2355 pushed.project = ProjectFilter::Is(id.clone());
2356 outcomes.record(Predicate::Project, Outcome::PushedDown);
2357 } else {
2358 local.project = Some(ProjectFilter::Is(id.clone()));
2359 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2360 }
2361 }
2362 }
2363
2364 DocumentShape {
2365 pushed,
2366 local,
2367 outcomes,
2368 }
2369}
2370
2371fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2373 match stream {
2374 StreamKind::Projects => {
2375 let shaped = shape_projects(capabilities, filters);
2376 HitShape {
2377 stream,
2378 tasks: TaskQuery::default(),
2379 projects: shaped.pushed,
2380 local_tasks: LocalTasks::default(),
2381 local_projects: shaped.local,
2382 outcomes: shaped.outcomes,
2383 }
2384 }
2385 StreamKind::Items | StreamKind::Tasks => {
2386 let shaped = shape_tasks(
2387 capabilities,
2388 filters,
2389 &ProjectFilter::Any,
2390 &[],
2391 None,
2392 &[],
2393 None,
2394 );
2395 HitShape {
2396 stream,
2397 tasks: shaped.pushed,
2398 projects: ProjectQuery::default(),
2399 local_tasks: shaped.local,
2400 local_projects: LocalProjects::default(),
2401 outcomes: shaped.outcomes,
2402 }
2403 }
2404 }
2405}
2406
2407fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2414 if compensating {
2415 ceiling
2416 } else {
2417 budget.min(ceiling)
2418 }
2419}
2420
2421fn ceiling(source: &ResolvedSource) -> u32 {
2423 source.source().capabilities().max_page_size.max(1)
2424}
2425
2426async fn fetch_tasks(
2428 source: &ResolvedSource,
2429 shape: &TaskShape,
2430 start: &Resume,
2431 budget: u32,
2432 calls: &AtomicU32,
2433) -> Result<Fetched<Task>, SourceError> {
2434 let compensating = shape.local != LocalTasks::default();
2435 walk(
2436 start,
2437 budget,
2438 page_size(compensating, budget, ceiling(source)),
2439 |task| shape.local.keeps(task),
2440 |cursor, limit| async move {
2441 calls.fetch_add(1, Ordering::Relaxed);
2442 let request = PageRequest { cursor, limit };
2443 let page = source.source().query_tasks(&shape.pushed, &request).await?;
2444 match shape.local.commented_since {
2445 None => Ok(page),
2446 Some(since) => comment::commented_since(source, &shape.local, page, since).await,
2447 }
2448 },
2449 )
2450 .await
2451}
2452
2453async fn fetch_projects(
2455 source: &ResolvedSource,
2456 shape: &ProjectShape,
2457 start: &Resume,
2458 budget: u32,
2459 calls: &AtomicU32,
2460) -> Result<Fetched<Project>, SourceError> {
2461 let compensating = shape.local != LocalProjects::default();
2462 walk(
2463 start,
2464 budget,
2465 page_size(compensating, budget, ceiling(source)),
2466 |project| shape.local.keeps(project),
2467 |cursor, limit| async move {
2468 calls.fetch_add(1, Ordering::Relaxed);
2469 let request = PageRequest { cursor, limit };
2470 source
2471 .source()
2472 .query_projects(&shape.pushed, &request)
2473 .await
2474 },
2475 )
2476 .await
2477}
2478
2479async fn fetch_documents(
2481 source: &ResolvedSource,
2482 shape: &DocumentShape,
2483 start: &Resume,
2484 budget: u32,
2485 calls: &AtomicU32,
2486) -> Result<Fetched<Document>, SourceError> {
2487 let compensating = shape.local != LocalDocuments::default();
2488 walk(
2489 start,
2490 budget,
2491 page_size(compensating, budget, ceiling(source)),
2492 |document| shape.local.keeps(document),
2493 |cursor, limit| async move {
2494 calls.fetch_add(1, Ordering::Relaxed);
2495 let request = PageRequest { cursor, limit };
2496 source
2497 .source()
2498 .query_documents(&shape.pushed, &request)
2499 .await
2500 },
2501 )
2502 .await
2503}
2504
2505async fn fetch_labels(
2507 source: &ResolvedSource,
2508 start: &Resume,
2509 budget: u32,
2510 calls: &AtomicU32,
2511) -> Result<Fetched<Label>, SourceError> {
2512 walk(
2513 start,
2514 budget,
2515 page_size(false, budget, ceiling(source)),
2516 |_| true,
2517 |cursor, limit| async move {
2518 calls.fetch_add(1, Ordering::Relaxed);
2519 let request = PageRequest { cursor, limit };
2520 source.source().labels(&request).await
2521 },
2522 )
2523 .await
2524}
2525
2526async fn fetch_hits(
2528 source: &ResolvedSource,
2529 shape: &HitShape,
2530 start: &Resume,
2531 budget: u32,
2532 calls: &AtomicU32,
2533) -> Result<Fetched<Found>, SourceError> {
2534 let ceiling = ceiling(source);
2535 match shape.stream {
2536 StreamKind::Projects => {
2537 let compensating = shape.local_projects != LocalProjects::default();
2538 walk(
2539 start,
2540 budget,
2541 page_size(compensating, budget, ceiling),
2542 |found| match found {
2543 Found::Project(project) => shape.local_projects.keeps(project),
2544 Found::Task(_) => true,
2545 },
2546 |cursor, limit| async move {
2547 calls.fetch_add(1, Ordering::Relaxed);
2548 let request = PageRequest { cursor, limit };
2549 let page = source
2550 .source()
2551 .query_projects(&shape.projects, &request)
2552 .await?;
2553 Ok(Page {
2554 items: page.items.into_iter().map(Found::Project).collect(),
2555 next: page.next,
2556 })
2557 },
2558 )
2559 .await
2560 }
2561 StreamKind::Items | StreamKind::Tasks => {
2562 let compensating = shape.local_tasks != LocalTasks::default();
2563 walk(
2564 start,
2565 budget,
2566 page_size(compensating, budget, ceiling),
2567 |found| match found {
2568 Found::Task(task) => shape.local_tasks.keeps(task),
2569 Found::Project(_) => true,
2570 },
2571 |cursor, limit| async move {
2572 calls.fetch_add(1, Ordering::Relaxed);
2573 let request = PageRequest { cursor, limit };
2574 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2575 Ok(Page {
2576 items: page.items.into_iter().map(Found::Task).collect(),
2577 next: page.next,
2578 })
2579 },
2580 )
2581 .await
2582 }
2583 }
2584}
2585
2586async fn forward_edges(
2588 source: &ResolvedSource,
2589 entity: Entity,
2590 id: &NativeId,
2591 request: &PageRequest,
2592) -> Result<Page<DependencyEdge>, SourceError> {
2593 match entity {
2594 Entity::Task => {
2595 source
2596 .source()
2597 .task_dependencies(id, Direction::DependsOn, request)
2598 .await
2599 }
2600 Entity::Project => {
2601 source
2602 .source()
2603 .project_dependencies(id, Direction::DependsOn, request)
2604 .await
2605 }
2606 }
2607}
2608
2609#[expect(
2618 clippy::too_many_arguments,
2619 reason = "every argument is one axis of one walk — the source, the item, the \
2620 direction, which of its two graphs, whether the reverse is emulated, where \
2621 to resume, how many rows to return and where to count calls. Grouping them \
2622 into a struct would name the same eight values one indirection further from \
2623 the loop that reads them."
2624)]
2625async fn fetch_edges(
2626 source: &ResolvedSource,
2627 native: &NativeId,
2628 direction: Direction,
2629 entity: Entity,
2630 emulating: bool,
2631 start: &Resume,
2632 budget: u32,
2633 calls: &AtomicU32,
2634) -> Result<Fetched<DependencyEdge>, SourceError> {
2635 let ceiling = ceiling(source);
2636 if !emulating {
2637 return walk(
2638 start,
2639 budget,
2640 page_size(false, budget, ceiling),
2641 |_| true,
2642 |cursor, limit| async move {
2643 calls.fetch_add(1, Ordering::Relaxed);
2644 let request = PageRequest { cursor, limit };
2645 match entity {
2646 Entity::Task => {
2647 source
2648 .source()
2649 .task_dependencies(native, direction, &request)
2650 .await
2651 }
2652 Entity::Project => {
2653 source
2654 .source()
2655 .project_dependencies(native, direction, &request)
2656 .await
2657 }
2658 }
2659 },
2660 )
2661 .await;
2662 }
2663
2664 walk(
2665 start,
2666 budget,
2667 ceiling,
2668 |_| true,
2669 |cursor, limit| async move {
2670 calls.fetch_add(1, Ordering::Relaxed);
2671 let request = PageRequest { cursor, limit };
2672 let (ids, next) = match entity {
2673 Entity::Task => {
2674 let page = source
2675 .source()
2676 .query_tasks(&TaskQuery::default(), &request)
2677 .await?;
2678 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2679 (ids, page.next)
2680 }
2681 Entity::Project => {
2682 let page = source
2683 .source()
2684 .query_projects(&ProjectQuery::default(), &request)
2685 .await?;
2686 let ids: Vec<NativeId> =
2687 page.items.into_iter().map(|project| project.id).collect();
2688 (ids, page.next)
2689 }
2690 };
2691
2692 let mut edges = Vec::new();
2693 for id in ids {
2694 let mut inner: Option<Cursor> = None;
2695 loop {
2696 calls.fetch_add(1, Ordering::Relaxed);
2697 let request = PageRequest {
2698 cursor: inner.clone(),
2699 limit,
2700 };
2701 let page = forward_edges(source, entity, &id, &request).await?;
2702 fits(page.items.len(), limit)?;
2706 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2707 unrepeated(
2708 page.next.as_ref(),
2709 inner.as_ref(),
2710 "its forward edges were being scanned",
2711 )?;
2712 match page.next {
2713 Some(cursor) => inner = Some(cursor),
2714 None => break,
2715 }
2716 }
2717 }
2718
2719 Ok(Page { items: edges, next })
2720 },
2721 )
2722 .await
2723}