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, TaskDetails};
56pub(crate) use copy::malformed_links;
57pub use copy::{
58 BudgetSpent, CopyAction, CopyItems, CopyLink, CopyLookup, CopyOutcome, CopyReport, CopyRequest,
59 CopyScope, CopyVia, MatchBy, NoCounterpart, Spent,
60};
61pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
62pub use local::ProjectSelector;
63pub use metadata::MetadataSet;
64pub use narrow::{TaskContentSet, TaskPrioritySet};
65pub use rendered::{
66 Body, DocumentCreate, Regenerated, Regeneration, RenderRequest, RenderTemplate, RenderedRecord,
67 TaskCreate, TaskCreated, TemplateAnswers, UnusedAnswers,
68};
69pub use update::TaskUpdated;
70
71#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
76pub struct Qualified<T> {
77 pub id: GlobalId,
79 pub item: T,
81}
82
83#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
88pub struct QualifiedEdge {
89 pub from: QualifiedEndpoint,
94 pub to: QualifiedEndpoint,
96 pub kind: onetaskgraph_plugin_api::DependencyKind,
98}
99
100#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
102pub struct QualifiedEndpoint {
103 pub id: GlobalId,
105 pub kind: onetaskgraph_plugin_api::ItemKind,
107}
108
109impl std::fmt::Display for QualifiedEndpoint {
110 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
111 self.id.fmt(formatter)
112 }
113}
114
115#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
117#[serde(tag = "kind", rename_all = "kebab-case")]
118pub enum SearchHit {
119 Task(Qualified<Task>),
121 Project(Qualified<Project>),
123}
124
125#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
127#[serde(rename_all = "kebab-case")]
128pub enum SearchKind {
129 Tasks,
131 Projects,
133 #[default]
135 Both,
136}
137
138#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
140pub struct SourceListing {
141 pub source: SourceName,
143 pub kind: String,
155 #[serde(flatten)]
157 pub state: SourceState,
158}
159
160#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
162#[serde(tag = "state", rename_all = "kebab-case")]
163pub enum SourceState {
164 Available {
166 capabilities: Capabilities,
168 },
169 Unavailable {
171 error: SourceError,
173 },
174}
175
176#[derive(Debug, Clone, PartialEq)]
178pub struct Paging {
179 pub limit: NonZeroU32,
181 pub token: Option<PageToken>,
183}
184
185#[derive(Debug, Clone, Default, PartialEq)]
187pub struct Filters {
188 pub text: Option<TextQuery>,
190 pub labels: LabelFilter,
192 pub statuses: Vec<StatusCategory>,
194}
195
196#[derive(Debug, Clone)]
198pub struct TaskRequest {
199 pub sources: Vec<SourceName>,
201 pub filters: Filters,
203 pub project: ProjectSelector,
205 pub priorities: Vec<Priority>,
211 pub commented_since: Option<DateTime<Utc>>,
218 pub metadata: Vec<MetadataMatch>,
224 pub origin: Option<GlobalId>,
227 pub 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(
751 "--create cannot be given with {flag}: --create asserts the destination holds no \
752 counterpart, so there is nothing for {flag} to look for\n\
753 next: drop {flag} to create each item without looking, or drop --create to look."
754 )]
755 CreateWith {
756 flag: CopyLookup,
758 },
759
760 #[error(
766 "{item} already records a counterpart at the destination, {carrier}, and --create \
767 asserts it has none\n\
768 next: copy it without --create, which updates {carrier}."
769 )]
770 CreateCarried {
771 item: GlobalId,
773 carrier: GlobalId,
775 },
776
777 #[error(
783 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
784 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
785 {id} on its own with `onetaskgraph task copy`.",
786 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
787 )]
788 NotAMember {
789 id: GlobalId,
791 projects: Vec<GlobalId>,
793 },
794
795 #[error(
804 "{item} depends on {member}, which this copy was not told to carry and which records \
805 no origin in {destination}\n\
806 next: name {member} with --member as well, record its {destination} id at \
807 onetaskgraph.origin, or copy the whole project without --member."
808 )]
809 UnrecordedMember {
810 item: GlobalId,
812 member: GlobalId,
814 destination: SourceName,
816 },
817
818 #[error(
824 "source {name} could not do it: {error}\n\
825 next: fix what the source named above, then copy again."
826 )]
827 SourceRefused {
828 name: String,
830 error: SourceError,
832 },
833
834 #[error(
842 "the copy failed and could not be undone.\n\
843 it failed because: {error}\n\
844 it could not be undone because: {refusal}\n\
845 so these still hold what it wrote: {left_behind}\n\
846 next: remove or put back those items, then copy again."
847 )]
848 CopyNotUndone {
849 error: Box<EngineError>,
851 left_behind: LeftBehind,
857 refusal: SourceError,
859 },
860}
861
862#[derive(Debug, Clone, PartialEq)]
871pub struct LeftBehind {
872 first: GlobalId,
874 rest: Vec<GlobalId>,
876}
877
878impl LeftBehind {
879 #[must_use]
881 pub fn new(first: GlobalId) -> Self {
882 Self {
883 first,
884 rest: Vec::new(),
885 }
886 }
887
888 pub fn push(&mut self, id: GlobalId) {
890 self.rest.push(id);
891 }
892
893 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
895 std::iter::once(&self.first).chain(self.rest.iter())
896 }
897}
898
899impl std::fmt::Display for LeftBehind {
900 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
902 write!(formatter, "{}", self.first)?;
903 for id in &self.rest {
904 write!(formatter, ", {id}")?;
905 }
906 Ok(())
907 }
908}
909
910pub enum ConfiguredSource {
917 Ready(ResolvedSource),
919 Unavailable(UnavailableSource),
921}
922
923impl ConfiguredSource {
924 #[must_use]
926 pub fn name(&self) -> &SourceName {
927 match self {
928 Self::Ready(source) => source.name(),
929 Self::Unavailable(source) => source.name(),
930 }
931 }
932}
933
934pub struct Engine {
936 sources: Vec<ConfiguredSource>,
938 selection: Vec<SourceName>,
940}
941
942impl Engine {
943 #[must_use]
951 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
952 let (ready, unavailable) = resolve_available(config, secrets);
953 Self::new(
954 ready
955 .into_iter()
956 .map(ConfiguredSource::Ready)
957 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
958 .collect(),
959 config.selected_sources(),
960 )
961 }
962
963 #[must_use]
966 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
967 Self { sources, selection }
968 }
969
970 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
972 self.sources.iter().filter_map(|source| match source {
973 ConfiguredSource::Ready(ready) => Some(ready),
974 ConfiguredSource::Unavailable(_) => None,
975 })
976 }
977
978 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
980 self.sources.iter().filter_map(|source| match source {
981 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
982 ConfiguredSource::Ready(_) => None,
983 })
984 }
985
986 #[must_use]
988 pub fn listing(&self) -> Vec<SourceListing> {
989 let mut listings: Vec<SourceListing> = self
990 .ready()
991 .map(|source| SourceListing {
992 source: source.name().clone(),
993 kind: source.kind().to_owned(),
994 state: SourceState::Available {
995 capabilities: source.source().capabilities(),
996 },
997 })
998 .chain(self.unavailable().map(|source| SourceListing {
999 source: source.name().clone(),
1000 kind: source.kind().to_owned(),
1001 state: SourceState::Unavailable {
1002 error: source.error().clone(),
1003 },
1004 }))
1005 .collect();
1006 listings.sort_by(|left, right| left.source.cmp(&right.source));
1007 listings
1008 }
1009
1010 #[must_use]
1016 pub fn has(&self, name: &SourceName) -> bool {
1017 self.sources.iter().any(|source| source.name() == name)
1018 }
1019
1020 pub async fn tasks(
1028 &self,
1029 request: &TaskRequest,
1030 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1031 let mut names = self.resolve_selection(&request.sources)?;
1032 if let ProjectSelector::Qualified(id) = &request.project {
1037 self.known(&id.source)?;
1038 names.retain(|name| name == &id.source);
1039 }
1040 let query = shape(
1041 "task-list",
1042 &names,
1043 &(
1044 &request.filters,
1045 &request.project,
1046 &request.priorities,
1047 &request.commented_since,
1048 &request.metadata,
1049 &request.origin,
1050 ),
1051 );
1052 let states = resumption(
1053 self,
1054 request.paging.token.as_ref(),
1055 &[StreamKind::Items],
1056 &query,
1057 )?;
1058 let budget = request.paging.limit.get();
1059
1060 let mut answer = Answer::new();
1061 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1062
1063 let shapes: Vec<TaskShape> = ready
1064 .iter()
1065 .map(|source| {
1066 shape_tasks(
1067 &source.source().capabilities(),
1068 &request.filters,
1069 &project_filter(&request.project),
1070 &request.priorities,
1071 request.commented_since,
1072 &request.metadata,
1073 request.origin.as_ref(),
1074 )
1075 })
1076 .collect();
1077 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1078 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1079
1080 let walks = ready
1081 .iter()
1082 .enumerate()
1083 .map(|(index, source)| {
1084 fetch_tasks(
1085 source,
1086 &shapes[index],
1087 &starts[index],
1088 budget,
1089 &counters[index],
1090 )
1091 })
1092 .collect();
1093
1094 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1095 answer.finish(
1096 streams,
1097 budget,
1098 owed(&states),
1099 &query,
1100 |name, task: Task| {
1101 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1102 },
1103 )
1104 }
1105
1106 pub async fn projects(
1112 &self,
1113 request: &ProjectRequest,
1114 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1115 let names = self.resolve_selection(&request.sources)?;
1116 let query = shape("project-list", &names, &request.filters);
1117 let states = resumption(
1118 self,
1119 request.paging.token.as_ref(),
1120 &[StreamKind::Items],
1121 &query,
1122 )?;
1123 let budget = request.paging.limit.get();
1124
1125 let mut answer = Answer::new();
1126 let mut with_projects = Vec::new();
1132 for source in answer.split(self, &names) {
1133 if source.source().capabilities().projects.is_native() {
1134 with_projects.push(source);
1135 } else {
1136 answer.unreachable_predicate(source, Predicate::Project);
1137 }
1138 }
1139 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1140
1141 let shapes: Vec<ProjectShape> = ready
1142 .iter()
1143 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1144 .collect();
1145 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1146 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1147
1148 let walks = ready
1149 .iter()
1150 .enumerate()
1151 .map(|(index, source)| {
1152 fetch_projects(
1153 source,
1154 &shapes[index],
1155 &starts[index],
1156 budget,
1157 &counters[index],
1158 )
1159 })
1160 .collect();
1161
1162 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1163 answer.finish(
1164 streams,
1165 budget,
1166 owed(&states),
1167 &query,
1168 |name, project: Project| Qualified {
1169 id: GlobalId::new(name.clone(), project.id.clone()),
1170 item: project,
1171 },
1172 )
1173 }
1174
1175 pub async fn documents(
1187 &self,
1188 request: &DocumentRequest,
1189 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1190 let mut names = self.resolve_selection(&request.sources)?;
1191 if let ProjectSelector::Qualified(id) = &request.project {
1194 self.known(&id.source)?;
1195 names.retain(|name| name == &id.source);
1196 }
1197 let query = shape(
1198 "document-list",
1199 &names,
1200 &(&request.filters, &request.project),
1201 );
1202 let states = resumption(
1203 self,
1204 request.paging.token.as_ref(),
1205 &[StreamKind::Items],
1206 &query,
1207 )?;
1208 let budget = request.paging.limit.get();
1209
1210 let mut answer = Answer::new();
1211 let mut with_documents = Vec::new();
1212 for source in answer.split(self, &names) {
1213 if source.source().capabilities().documents.is_native() {
1214 with_documents.push(source);
1215 } else {
1216 answer.unreachable_predicate(source, Predicate::Document);
1217 }
1218 }
1219 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1220
1221 let shapes: Vec<DocumentShape> = ready
1222 .iter()
1223 .map(|source| {
1224 shape_documents(
1225 &source.source().capabilities(),
1226 &request.filters,
1227 &project_filter(&request.project),
1228 )
1229 })
1230 .collect();
1231 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1232 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1233
1234 let walks = ready
1235 .iter()
1236 .enumerate()
1237 .map(|(index, source)| {
1238 fetch_documents(
1239 source,
1240 &shapes[index],
1241 &starts[index],
1242 budget,
1243 &counters[index],
1244 )
1245 })
1246 .collect();
1247
1248 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1249 answer.finish(
1250 streams,
1251 budget,
1252 owed(&states),
1253 &query,
1254 |name, document: Document| Qualified {
1255 id: GlobalId::new(name.clone(), document.id.clone()),
1256 item: document,
1257 },
1258 )
1259 }
1260
1261 pub async fn labels(
1267 &self,
1268 request: &LabelRequest,
1269 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1270 let names = self.resolve_selection(&request.sources)?;
1271 let query = shape("label-list", &names, &());
1272 let states = resumption(
1273 self,
1274 request.paging.token.as_ref(),
1275 &[StreamKind::Items],
1276 &query,
1277 )?;
1278 let budget = request.paging.limit.get();
1279
1280 let mut answer = Answer::new();
1281 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1282 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1283 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1284
1285 let walks = ready
1286 .iter()
1287 .enumerate()
1288 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1289 .collect();
1290
1291 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1292 answer.finish(
1293 streams,
1294 budget,
1295 owed(&states),
1296 &query,
1297 |name, label: Label| Qualified {
1298 id: GlobalId::new(name.clone(), label.id.clone()),
1299 item: label,
1300 },
1301 )
1302 }
1303
1304 pub async fn search(
1310 &self,
1311 request: &SearchRequest,
1312 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1313 let names = self.resolve_selection(&request.sources)?;
1314 let reads: &[StreamKind] = match request.kind {
1318 SearchKind::Tasks => &[StreamKind::Tasks],
1319 SearchKind::Projects => &[StreamKind::Projects],
1320 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1321 };
1322 let query = shape("search", &names, &request.text);
1327 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1328 let budget = request.paging.limit.get();
1329 let filters = Filters {
1330 text: Some(request.text.clone()),
1331 ..Filters::default()
1332 };
1333
1334 let mut answer = Answer::new();
1335
1336 let mut ready = Vec::new();
1339 let mut kinds = Vec::new();
1340 let mut starts = Vec::new();
1341 for source in answer.split(self, &names) {
1342 let mut streams = Vec::new();
1343 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1344 streams.push(StreamKind::Tasks);
1345 }
1346 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1347 if source.source().capabilities().projects.is_native() {
1348 streams.push(StreamKind::Projects);
1349 } else {
1350 answer.unreachable_predicate(source, Predicate::Project);
1351 }
1352 }
1353 for stream in streams {
1354 if let Some(resume) = resume_at(&states, source.name(), stream) {
1355 ready.push(source);
1356 kinds.push(stream);
1357 starts.push(resume);
1358 }
1359 }
1360 }
1361
1362 let shapes: Vec<HitShape> = ready
1363 .iter()
1364 .zip(kinds.iter())
1365 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1366 .collect();
1367 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1368 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1369
1370 let walks = ready
1371 .iter()
1372 .enumerate()
1373 .map(|(index, source)| {
1374 fetch_hits(
1375 source,
1376 &shapes[index],
1377 &starts[index],
1378 budget,
1379 &counters[index],
1380 )
1381 })
1382 .collect();
1383
1384 let streams =
1385 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1386 answer.finish(
1387 streams,
1388 budget,
1389 owed(&states),
1390 &query,
1391 |name, found: Found| match found {
1392 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1393 GlobalId::new(name.clone(), task.id.clone()),
1394 task,
1395 )),
1396 Found::Project(project) => SearchHit::Project(Qualified {
1397 id: GlobalId::new(name.clone(), project.id.clone()),
1398 item: project,
1399 }),
1400 },
1401 )
1402 }
1403
1404 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1411 let name = self.known(&id.source)?;
1412 let mut answer = Answer::new();
1413 let selected = answer.split(self, std::slice::from_ref(&name));
1414 let Some(source) = selected.first() else {
1415 return answer.nothing();
1416 };
1417 let found = source.source().get_task(&id.native).await;
1418 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1419 answer.one(source, found, |task| {
1420 delivery::qualified_task(qualified, task)
1421 })
1422 }
1423
1424 pub async fn project(
1430 &self,
1431 id: &GlobalId,
1432 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1433 let name = self.known(&id.source)?;
1434 let mut answer = Answer::new();
1435 let selected = answer.split(self, std::slice::from_ref(&name));
1436 let Some(source) = selected.first() else {
1437 return answer.nothing();
1438 };
1439 let found = source.source().get_project(&id.native).await;
1440 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1441 answer.one(source, found, |project| Qualified {
1442 id: qualified,
1443 item: project,
1444 })
1445 }
1446
1447 pub async fn document(
1457 &self,
1458 id: &GlobalId,
1459 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1460 let name = self.known(&id.source)?;
1461 let mut answer = Answer::new();
1462 let selected = answer.split(self, std::slice::from_ref(&name));
1463 let Some(source) = selected.first() else {
1464 return answer.nothing();
1465 };
1466 if !source.source().capabilities().documents.is_native() {
1467 answer.unreachable_predicate(source, Predicate::Document);
1468 return answer.nothing();
1469 }
1470 let found = source.source().get_document(&id.native).await;
1471 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1472 answer.one(source, found, |document| Qualified {
1473 id: qualified,
1474 item: document,
1475 })
1476 }
1477
1478 pub async fn task_dependencies(
1485 &self,
1486 request: &DependencyRequest,
1487 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1488 self.dependencies(request, Entity::Task).await
1489 }
1490
1491 pub async fn project_dependencies(
1497 &self,
1498 request: &DependencyRequest,
1499 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1500 self.dependencies(request, Entity::Project).await
1501 }
1502
1503 async fn dependencies(
1506 &self,
1507 request: &DependencyRequest,
1508 entity: Entity,
1509 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1510 let name = self.known(&request.id.source)?;
1511 let query = shape(
1512 "dependencies",
1513 std::slice::from_ref(&name),
1514 &(entity, &request.id.native, request.direction),
1515 );
1516 let states = resumption(
1517 self,
1518 request.paging.token.as_ref(),
1519 &[StreamKind::Items],
1520 &query,
1521 )?;
1522 let budget = request.paging.limit.get();
1523
1524 let mut answer = Answer::new();
1525 let (ready, starts) = walking(
1526 answer.split(self, std::slice::from_ref(&name)),
1527 &states,
1528 StreamKind::Items,
1529 );
1530 let Some(source) = ready.first() else {
1531 return answer.nothing();
1532 };
1533
1534 let capabilities = source.source().capabilities();
1535 let support = match entity {
1536 Entity::Task => capabilities.task_dependencies,
1537 Entity::Project => capabilities.project_dependencies,
1538 };
1539 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1543 let mut outcomes = Outcomes::default();
1544 if request.direction == Direction::DependedOnBy {
1545 if emulating {
1546 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1547 } else {
1548 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1549 }
1550 }
1551
1552 let counters = vec![AtomicU32::new(0)];
1553 let walked = fetch_edges(
1554 source,
1555 &request.id.native,
1556 request.direction,
1557 entity,
1558 emulating,
1559 &starts[0],
1560 budget,
1561 &counters[0],
1562 )
1563 .await;
1564
1565 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1566 answer.finish(
1567 streams,
1568 budget,
1569 owed(&states),
1570 &query,
1571 |name, edge: DependencyEdge| QualifiedEdge {
1572 from: qualify_endpoint(name, edge.from),
1573 to: qualify_endpoint(name, edge.to),
1574 kind: edge.kind,
1575 },
1576 )
1577 }
1578
1579 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1581 if asked.is_empty() {
1582 if self.selection.is_empty() {
1583 return Err(EngineError::NoSources);
1584 }
1585 return Ok(self.selection.clone());
1586 }
1587 asked.iter().map(|name| self.known(name)).collect()
1588 }
1589
1590 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1592 if self.has(name) {
1593 return Ok(name.clone());
1594 }
1595 if self.sources.is_empty() {
1596 return Err(EngineError::NoSources);
1597 }
1598 Err(EngineError::UnknownSource {
1599 name: name.to_string(),
1600 configured: self
1601 .listing()
1602 .iter()
1603 .map(|listing| listing.source.to_string())
1604 .collect::<Vec<_>>()
1605 .join(", "),
1606 })
1607 }
1608}
1609
1610fn qualify_endpoint(
1611 source: &SourceName,
1612 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1613) -> QualifiedEndpoint {
1614 let kind = endpoint.kind;
1615 let is_qualified = endpoint.is_qualified();
1616 let endpoint_id = endpoint.into_id();
1617 QualifiedEndpoint {
1618 id: if is_qualified {
1619 endpoint_id
1620 .parse()
1621 .expect("plugin-api validates qualified dependency endpoints")
1622 } else {
1623 GlobalId::new(source.clone(), NativeId(endpoint_id))
1624 },
1625 kind,
1626 }
1627}
1628
1629#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1631enum Entity {
1632 Task,
1634 Project,
1636}
1637
1638enum Found {
1640 Task(Task),
1642 Project(Project),
1644}
1645
1646#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1648enum Outcome {
1649 PushedDown,
1651 AppliedLocally,
1653 Emulated,
1655 Unavailable,
1657}
1658
1659#[derive(Debug, Clone, Default, PartialEq)]
1671struct Outcomes(BTreeMap<Predicate, Outcome>);
1672
1673impl Outcomes {
1674 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1680 self.0.insert(predicate, outcome);
1681 }
1682
1683 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1686 for predicate in predicates {
1687 self.record(predicate, outcome);
1688 }
1689 }
1690
1691 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1693 self.0
1694 .iter()
1695 .filter(|(_, recorded)| **recorded == outcome)
1696 .map(|(predicate, _)| *predicate)
1697 .collect()
1698 }
1699}
1700
1701struct TaskShape {
1703 pushed: TaskQuery,
1705 local: LocalTasks,
1707 outcomes: Outcomes,
1709}
1710
1711struct ProjectShape {
1713 pushed: ProjectQuery,
1715 local: LocalProjects,
1717 outcomes: Outcomes,
1719}
1720
1721struct DocumentShape {
1723 pushed: DocumentQuery,
1725 local: LocalDocuments,
1727 outcomes: Outcomes,
1729}
1730
1731struct HitShape {
1733 stream: StreamKind,
1735 tasks: TaskQuery,
1737 projects: ProjectQuery,
1739 local_tasks: LocalTasks,
1741 local_projects: LocalProjects,
1743 outcomes: Outcomes,
1745}
1746
1747struct Answer {
1753 plans: Vec<SourcePlan>,
1755 errors: Vec<SourceFailure>,
1757}
1758
1759impl Answer {
1760 fn new() -> Self {
1761 Self {
1762 plans: Vec::new(),
1763 errors: Vec::new(),
1764 }
1765 }
1766
1767 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1772 let mut selected = Vec::new();
1773 for name in names {
1774 match engine.sources.iter().find(|source| source.name() == name) {
1775 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1776 Some(ConfiguredSource::Unavailable(source)) => {
1777 self.errors.push(source.failure());
1778 }
1779 None => {}
1780 }
1781 }
1782 selected
1783 }
1784
1785 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1787 let mut outcomes = Outcomes::default();
1788 outcomes.record(predicate, Outcome::Unavailable);
1789 self.plans.push(plan_for(source, outcomes, 0));
1790 }
1791
1792 fn collect<T>(
1794 &mut self,
1795 ready: &[&ResolvedSource],
1796 walked: Vec<Result<Fetched<T>, SourceError>>,
1797 counters: &[AtomicU32],
1798 outcomes: Vec<Outcomes>,
1799 ) -> Vec<Stream<T>> {
1800 let kinds = vec![StreamKind::Items; ready.len()];
1801 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1802 }
1803
1804 fn collect_streams<T>(
1806 &mut self,
1807 ready: &[&ResolvedSource],
1808 kinds: &[StreamKind],
1809 walked: Vec<Result<Fetched<T>, SourceError>>,
1810 counters: &[AtomicU32],
1811 outcomes: Vec<Outcomes>,
1812 ) -> Vec<Stream<T>> {
1813 let mut streams = Vec::new();
1814 for (index, result) in walked.into_iter().enumerate() {
1815 let source = ready[index];
1816 let pages = counters[index].load(Ordering::Relaxed);
1817 self.plans
1818 .push(plan_for(source, outcomes[index].clone(), pages));
1819 match result {
1820 Ok(fetched) => streams.push(Stream {
1821 source: source.name().clone(),
1822 kind: kinds[index],
1823 fetched,
1824 }),
1825 Err(error) => self.errors.push(SourceFailure {
1828 source: source.name().clone(),
1829 error,
1830 }),
1831 }
1832 }
1833 streams
1834 }
1835
1836 fn one<T, U>(
1838 self,
1839 source: &ResolvedSource,
1840 found: Result<Option<T>, SourceError>,
1841 qualify: impl FnOnce(T) -> U,
1842 ) -> Result<QueryResponse<U>, EngineError> {
1843 Ok(self.one_response(source, found, qualify))
1844 }
1845
1846 fn one_response<T, U>(
1848 mut self,
1849 source: &ResolvedSource,
1850 found: Result<Option<T>, SourceError>,
1851 qualify: impl FnOnce(T) -> U,
1852 ) -> QueryResponse<U> {
1853 self.plans.push(plan_for(source, Outcomes::default(), 1));
1854 let items = match found {
1855 Ok(Some(item)) => vec![qualify(item)],
1856 Ok(None) => Vec::new(),
1857 Err(error) => {
1858 self.errors.push(SourceFailure {
1859 source: source.name().clone(),
1860 error,
1861 });
1862 Vec::new()
1863 }
1864 };
1865 QueryResponse {
1866 items,
1867 next: None,
1868 plan: QueryPlan {
1869 per_source: merge_plans(self.plans),
1870 },
1871 errors: self.errors,
1872 }
1873 }
1874
1875 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1877 Ok(QueryResponse {
1878 items: Vec::new(),
1879 next: None,
1880 plan: QueryPlan {
1881 per_source: merge_plans(self.plans),
1882 },
1883 errors: self.errors,
1884 })
1885 }
1886
1887 fn finish<T, U>(
1892 self,
1893 streams: Vec<Stream<T>>,
1894 budget: u32,
1895 first: Option<&Owed>,
1896 query: &str,
1897 qualify: impl Fn(&SourceName, T) -> U,
1898 ) -> Result<QueryResponse<U>, EngineError> {
1899 let (rows, states, owed) = merge(streams, budget, first);
1900 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1901 Ok(QueryResponse {
1902 items: rows
1903 .into_iter()
1904 .map(|(name, item)| qualify(&name, item))
1905 .collect(),
1906 next,
1907 plan: QueryPlan {
1908 per_source: merge_plans(self.plans),
1909 },
1910 errors: self.errors,
1911 })
1912 }
1913}
1914
1915fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1918 SourcePlan {
1919 source: source.name().clone(),
1920 kind: source.kind().to_owned(),
1921 pushed_down: outcomes.with(Outcome::PushedDown),
1922 applied_locally: outcomes.with(Outcome::AppliedLocally),
1923 emulated: outcomes.with(Outcome::Emulated),
1924 unavailable: outcomes.with(Outcome::Unavailable),
1925 pages_fetched: pages,
1926 }
1927}
1928
1929fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1934 let mut merged: Vec<SourcePlan> = Vec::new();
1935 for plan in plans {
1936 if let Some(existing) = merged
1937 .iter_mut()
1938 .find(|existing| existing.source == plan.source)
1939 {
1940 existing.pushed_down.extend(plan.pushed_down);
1941 existing.applied_locally.extend(plan.applied_locally);
1942 existing.emulated.extend(plan.emulated);
1943 existing.unavailable.extend(plan.unavailable);
1944 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1945 for list in [
1946 &mut existing.pushed_down,
1947 &mut existing.applied_locally,
1948 &mut existing.emulated,
1949 &mut existing.unavailable,
1950 ] {
1951 list.sort_unstable();
1952 list.dedup();
1953 }
1954 } else {
1955 merged.push(plan);
1956 }
1957 }
1958 merged
1959}
1960
1961fn walking<'a>(
1967 selected: Vec<&'a ResolvedSource>,
1968 states: &Option<Resumption>,
1969 kind: StreamKind,
1970) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1971 let mut ready = Vec::new();
1972 let mut starts = Vec::new();
1973 for source in selected {
1974 if let Some(resume) = resume_at(states, source.name(), kind) {
1975 ready.push(source);
1976 starts.push(resume);
1977 }
1978 }
1979 (ready, starts)
1980}
1981
1982fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
2000 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
2001 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
2002}
2003
2004fn fingerprint(text: &str) -> String {
2006 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
2007 for byte in text.as_bytes() {
2008 hash ^= u64::from(*byte);
2009 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
2010 }
2011 format!("{hash:016x}")
2012}
2013
2014fn owed(document: &Option<Resumption>) -> Option<&Owed> {
2019 document.as_ref()?.owed.as_ref()
2020}
2021
2022fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
2024 match states {
2025 None => Some(Resume::default()),
2026 Some(document) => document
2027 .streams
2028 .iter()
2029 .find(|state| &state.source == source && state.stream == kind)
2030 .map(|state| state.resume.clone()),
2031 }
2032}
2033
2034fn resumption(
2065 engine: &Engine,
2066 token: Option<&PageToken>,
2067 reads: &[StreamKind],
2068 query: &str,
2069) -> Result<Option<Resumption>, EngineError> {
2070 let Some(document) = token.map(PageToken::decode) else {
2071 return Ok(None);
2072 };
2073
2074 if document.query != query {
2081 return Err(EngineError::Token {
2082 message: "this page token was written by a different query — resume the walk it \
2083 came from, or drop --page to start this one from the beginning"
2084 .to_owned(),
2085 });
2086 }
2087 let states = &document.streams;
2088
2089 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2090 for state in states {
2091 if !reads.contains(&state.stream) {
2092 return Err(EngineError::Token {
2093 message: format!(
2094 "this page token resumes {}, which this command does not read — it \
2095 was written by a different query",
2096 state.stream.describe()
2097 ),
2098 });
2099 }
2100 let ceiling = engine
2101 .ready()
2102 .find(|source| source.name() == &state.source)
2103 .map(ceiling);
2104 if ceiling.is_none() && !engine.has(&state.source) {
2105 return Err(EngineError::Token {
2106 message: format!(
2107 "this page token resumes a source called {:?}, which this \
2108 configuration does not have",
2109 state.source.as_str()
2110 ),
2111 });
2112 }
2113 if let Some(ceiling) = ceiling
2114 && state.resume.skip >= ceiling
2115 {
2116 return Err(EngineError::Token {
2117 message: format!(
2118 "this page token resumes {} rows into a page of source {:?}, which \
2119 serves at most {ceiling}",
2120 state.resume.skip,
2121 state.source.as_str()
2122 ),
2123 });
2124 }
2125 if seen.contains(&(&state.source, state.stream)) {
2126 return Err(EngineError::Token {
2127 message: format!(
2128 "this page token gives source {:?} two places to resume from",
2129 state.source.as_str()
2130 ),
2131 });
2132 }
2133 seen.push((&state.source, state.stream));
2134 }
2135
2136 if let Some(owed) = &document.owed
2141 && !document
2142 .streams
2143 .iter()
2144 .any(|state| state.source == owed.source && state.stream == owed.stream)
2145 {
2146 return Err(EngineError::Token {
2147 message: format!(
2148 "this page token owes the next row to a stream it does not resume, \
2149 {:?}'s {}",
2150 owed.source.as_str(),
2151 owed.stream.describe()
2152 ),
2153 });
2154 }
2155
2156 Ok(Some(document))
2157}
2158
2159fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2165 match selector {
2166 ProjectSelector::Any => ProjectFilter::Any,
2167 ProjectSelector::Orphans => ProjectFilter::Orphans,
2168 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2169 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2170 }
2171}
2172
2173fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2175 match fields {
2176 TextFields::Title => vec![Predicate::SearchTitle],
2177 TextFields::Content => vec![Predicate::SearchContent],
2178 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2179 }
2180}
2181
2182fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2189 match fields {
2190 TextFields::Title => capabilities.search_title.is_native(),
2191 TextFields::Content => capabilities.search_content.is_native(),
2192 TextFields::TitleOrContent => {
2193 capabilities.search_title.is_native() && capabilities.search_content.is_native()
2194 }
2195 }
2196}
2197
2198fn shape_tasks(
2200 capabilities: &Capabilities,
2201 filters: &Filters,
2202 project: &ProjectFilter,
2203 priorities: &[Priority],
2204 commented_since: Option<DateTime<Utc>>,
2205 metadata: &[MetadataMatch],
2206 origin: Option<&GlobalId>,
2207) -> TaskShape {
2208 let mut pushed = TaskQuery::default();
2209 let mut local = LocalTasks::default();
2210 let mut outcomes = Outcomes::default();
2211
2212 if !filters.labels.is_empty() {
2213 if capabilities.filter_by_label.is_native() {
2214 pushed.labels = filters.labels.clone();
2215 outcomes.record(Predicate::Label, Outcome::PushedDown);
2216 } else {
2217 local.labels = Some(filters.labels.clone());
2218 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2219 }
2220 }
2221 if !filters.statuses.is_empty() {
2222 if capabilities.filter_by_status.is_native() {
2223 pushed.statuses.clone_from(&filters.statuses);
2224 outcomes.record(Predicate::Status, Outcome::PushedDown);
2225 } else {
2226 local.statuses.clone_from(&filters.statuses);
2227 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2228 }
2229 }
2230 if !priorities.is_empty() {
2231 if capabilities.filter_by_priority.is_native() {
2232 pushed.priorities = priorities.to_vec();
2233 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2234 } else {
2235 local.priorities = priorities.to_vec();
2236 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2237 }
2238 }
2239 if let Some(since) = commented_since {
2240 if capabilities.filter_by_comment_activity.is_native() {
2241 pushed.commented_since = Some(since);
2242 outcomes.record(Predicate::CommentedSince, Outcome::PushedDown);
2243 } else {
2244 local.commented_since = Some(since);
2245 outcomes.record(Predicate::CommentedSince, Outcome::AppliedLocally);
2246 }
2247 }
2248 if !metadata.is_empty() {
2249 if capabilities.filter_by_metadata.is_native() {
2250 pushed.metadata = metadata.to_vec();
2251 outcomes.record(Predicate::Metadata, Outcome::PushedDown);
2252 } else {
2253 local.metadata = metadata.to_vec();
2254 outcomes.record(Predicate::Metadata, Outcome::AppliedLocally);
2255 }
2256 }
2257 if let Some(origin) = origin {
2258 if capabilities.filter_by_origin.is_native() {
2259 pushed.origin = Some(origin.to_string());
2262 outcomes.record(Predicate::Origin, Outcome::PushedDown);
2263 } else {
2264 local.origin = Some(origin.clone());
2265 outcomes.record(Predicate::Origin, Outcome::AppliedLocally);
2266 }
2267 }
2268 if let Some(text) = &filters.text {
2269 let predicates = text_predicates(text.fields);
2270 if searches_natively(capabilities, text.fields) {
2271 pushed.text = Some(text.clone());
2272 outcomes.record_all(predicates, Outcome::PushedDown);
2273 } else {
2274 local.text = Some(text.clone());
2275 outcomes.record_all(predicates, Outcome::AppliedLocally);
2276 }
2277 }
2278 match project {
2279 ProjectFilter::Any => {}
2280 ProjectFilter::Orphans => {
2281 if capabilities.orphan_tasks.is_native() {
2282 pushed.project = ProjectFilter::Orphans;
2283 outcomes.record(Predicate::Project, Outcome::PushedDown);
2284 } else {
2285 local.project = Some(ProjectFilter::Orphans);
2286 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2287 }
2288 }
2289 ProjectFilter::Is(id) => {
2290 if capabilities.projects.is_native() {
2291 pushed.project = ProjectFilter::Is(id.clone());
2292 outcomes.record(Predicate::Project, Outcome::PushedDown);
2293 } else {
2294 local.project = Some(ProjectFilter::Is(id.clone()));
2295 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2296 }
2297 }
2298 }
2299
2300 TaskShape {
2301 pushed,
2302 local,
2303 outcomes,
2304 }
2305}
2306
2307fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2309 let mut pushed = ProjectQuery::default();
2310 let mut local = LocalProjects::default();
2311 let mut outcomes = Outcomes::default();
2312
2313 if !filters.labels.is_empty() {
2314 if capabilities.filter_by_label.is_native() {
2315 pushed.labels = filters.labels.clone();
2316 outcomes.record(Predicate::Label, Outcome::PushedDown);
2317 } else {
2318 local.labels = Some(filters.labels.clone());
2319 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2320 }
2321 }
2322 if !filters.statuses.is_empty() {
2323 if capabilities.filter_by_status.is_native() {
2324 pushed.statuses.clone_from(&filters.statuses);
2325 outcomes.record(Predicate::Status, Outcome::PushedDown);
2326 } else {
2327 local.statuses.clone_from(&filters.statuses);
2328 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2329 }
2330 }
2331 if let Some(text) = &filters.text {
2332 let predicates = text_predicates(text.fields);
2333 if searches_natively(capabilities, text.fields) {
2334 pushed.text = Some(text.clone());
2335 outcomes.record_all(predicates, Outcome::PushedDown);
2336 } else {
2337 local.text = Some(text.clone());
2338 outcomes.record_all(predicates, Outcome::AppliedLocally);
2339 }
2340 }
2341
2342 ProjectShape {
2343 pushed,
2344 local,
2345 outcomes,
2346 }
2347}
2348
2349fn shape_documents(
2354 capabilities: &Capabilities,
2355 filters: &DocumentFilters,
2356 project: &ProjectFilter,
2357) -> DocumentShape {
2358 let mut pushed = DocumentQuery::default();
2359 let mut local = LocalDocuments::default();
2360 let mut outcomes = Outcomes::default();
2361
2362 if !filters.labels.is_empty() {
2363 if capabilities.filter_by_label.is_native() {
2364 pushed.labels = filters.labels.clone();
2365 outcomes.record(Predicate::Label, Outcome::PushedDown);
2366 } else {
2367 local.labels = Some(filters.labels.clone());
2368 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2369 }
2370 }
2371 if let Some(text) = &filters.text {
2372 let predicates = text_predicates(text.fields);
2373 if searches_natively(capabilities, text.fields) {
2374 pushed.text = Some(text.clone());
2375 outcomes.record_all(predicates, Outcome::PushedDown);
2376 } else {
2377 local.text = Some(text.clone());
2378 outcomes.record_all(predicates, Outcome::AppliedLocally);
2379 }
2380 }
2381 match project {
2382 ProjectFilter::Any => {}
2383 ProjectFilter::Orphans => {
2384 if capabilities.orphan_tasks.is_native() {
2385 pushed.project = ProjectFilter::Orphans;
2386 outcomes.record(Predicate::Project, Outcome::PushedDown);
2387 } else {
2388 local.project = Some(ProjectFilter::Orphans);
2389 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2390 }
2391 }
2392 ProjectFilter::Is(id) => {
2393 if capabilities.projects.is_native() {
2394 pushed.project = ProjectFilter::Is(id.clone());
2395 outcomes.record(Predicate::Project, Outcome::PushedDown);
2396 } else {
2397 local.project = Some(ProjectFilter::Is(id.clone()));
2398 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2399 }
2400 }
2401 }
2402
2403 DocumentShape {
2404 pushed,
2405 local,
2406 outcomes,
2407 }
2408}
2409
2410fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2412 match stream {
2413 StreamKind::Projects => {
2414 let shaped = shape_projects(capabilities, filters);
2415 HitShape {
2416 stream,
2417 tasks: TaskQuery::default(),
2418 projects: shaped.pushed,
2419 local_tasks: LocalTasks::default(),
2420 local_projects: shaped.local,
2421 outcomes: shaped.outcomes,
2422 }
2423 }
2424 StreamKind::Items | StreamKind::Tasks => {
2425 let shaped = shape_tasks(
2426 capabilities,
2427 filters,
2428 &ProjectFilter::Any,
2429 &[],
2430 None,
2431 &[],
2432 None,
2433 );
2434 HitShape {
2435 stream,
2436 tasks: shaped.pushed,
2437 projects: ProjectQuery::default(),
2438 local_tasks: shaped.local,
2439 local_projects: LocalProjects::default(),
2440 outcomes: shaped.outcomes,
2441 }
2442 }
2443 }
2444}
2445
2446fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2453 if compensating {
2454 ceiling
2455 } else {
2456 budget.min(ceiling)
2457 }
2458}
2459
2460fn ceiling(source: &ResolvedSource) -> u32 {
2462 source.source().capabilities().max_page_size.max(1)
2463}
2464
2465async fn fetch_tasks(
2467 source: &ResolvedSource,
2468 shape: &TaskShape,
2469 start: &Resume,
2470 budget: u32,
2471 calls: &AtomicU32,
2472) -> Result<Fetched<Task>, SourceError> {
2473 let compensating = shape.local != LocalTasks::default();
2474 walk(
2475 start,
2476 budget,
2477 page_size(compensating, budget, ceiling(source)),
2478 |task| shape.local.keeps(task),
2479 |cursor, limit| async move {
2480 calls.fetch_add(1, Ordering::Relaxed);
2481 let request = PageRequest { cursor, limit };
2482 let page = source.source().query_tasks(&shape.pushed, &request).await?;
2483 match shape.local.commented_since {
2484 None => Ok(page),
2485 Some(since) => comment::commented_since(source, &shape.local, page, since).await,
2486 }
2487 },
2488 )
2489 .await
2490}
2491
2492async fn fetch_projects(
2494 source: &ResolvedSource,
2495 shape: &ProjectShape,
2496 start: &Resume,
2497 budget: u32,
2498 calls: &AtomicU32,
2499) -> Result<Fetched<Project>, SourceError> {
2500 let compensating = shape.local != LocalProjects::default();
2501 walk(
2502 start,
2503 budget,
2504 page_size(compensating, budget, ceiling(source)),
2505 |project| shape.local.keeps(project),
2506 |cursor, limit| async move {
2507 calls.fetch_add(1, Ordering::Relaxed);
2508 let request = PageRequest { cursor, limit };
2509 source
2510 .source()
2511 .query_projects(&shape.pushed, &request)
2512 .await
2513 },
2514 )
2515 .await
2516}
2517
2518async fn fetch_documents(
2520 source: &ResolvedSource,
2521 shape: &DocumentShape,
2522 start: &Resume,
2523 budget: u32,
2524 calls: &AtomicU32,
2525) -> Result<Fetched<Document>, SourceError> {
2526 let compensating = shape.local != LocalDocuments::default();
2527 walk(
2528 start,
2529 budget,
2530 page_size(compensating, budget, ceiling(source)),
2531 |document| shape.local.keeps(document),
2532 |cursor, limit| async move {
2533 calls.fetch_add(1, Ordering::Relaxed);
2534 let request = PageRequest { cursor, limit };
2535 source
2536 .source()
2537 .query_documents(&shape.pushed, &request)
2538 .await
2539 },
2540 )
2541 .await
2542}
2543
2544async fn fetch_labels(
2546 source: &ResolvedSource,
2547 start: &Resume,
2548 budget: u32,
2549 calls: &AtomicU32,
2550) -> Result<Fetched<Label>, SourceError> {
2551 walk(
2552 start,
2553 budget,
2554 page_size(false, budget, ceiling(source)),
2555 |_| true,
2556 |cursor, limit| async move {
2557 calls.fetch_add(1, Ordering::Relaxed);
2558 let request = PageRequest { cursor, limit };
2559 source.source().labels(&request).await
2560 },
2561 )
2562 .await
2563}
2564
2565async fn fetch_hits(
2567 source: &ResolvedSource,
2568 shape: &HitShape,
2569 start: &Resume,
2570 budget: u32,
2571 calls: &AtomicU32,
2572) -> Result<Fetched<Found>, SourceError> {
2573 let ceiling = ceiling(source);
2574 match shape.stream {
2575 StreamKind::Projects => {
2576 let compensating = shape.local_projects != LocalProjects::default();
2577 walk(
2578 start,
2579 budget,
2580 page_size(compensating, budget, ceiling),
2581 |found| match found {
2582 Found::Project(project) => shape.local_projects.keeps(project),
2583 Found::Task(_) => true,
2584 },
2585 |cursor, limit| async move {
2586 calls.fetch_add(1, Ordering::Relaxed);
2587 let request = PageRequest { cursor, limit };
2588 let page = source
2589 .source()
2590 .query_projects(&shape.projects, &request)
2591 .await?;
2592 Ok(Page {
2593 items: page.items.into_iter().map(Found::Project).collect(),
2594 next: page.next,
2595 })
2596 },
2597 )
2598 .await
2599 }
2600 StreamKind::Items | StreamKind::Tasks => {
2601 let compensating = shape.local_tasks != LocalTasks::default();
2602 walk(
2603 start,
2604 budget,
2605 page_size(compensating, budget, ceiling),
2606 |found| match found {
2607 Found::Task(task) => shape.local_tasks.keeps(task),
2608 Found::Project(_) => true,
2609 },
2610 |cursor, limit| async move {
2611 calls.fetch_add(1, Ordering::Relaxed);
2612 let request = PageRequest { cursor, limit };
2613 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2614 Ok(Page {
2615 items: page.items.into_iter().map(Found::Task).collect(),
2616 next: page.next,
2617 })
2618 },
2619 )
2620 .await
2621 }
2622 }
2623}
2624
2625async fn forward_edges(
2627 source: &ResolvedSource,
2628 entity: Entity,
2629 id: &NativeId,
2630 request: &PageRequest,
2631) -> Result<Page<DependencyEdge>, SourceError> {
2632 match entity {
2633 Entity::Task => {
2634 source
2635 .source()
2636 .task_dependencies(id, Direction::DependsOn, request)
2637 .await
2638 }
2639 Entity::Project => {
2640 source
2641 .source()
2642 .project_dependencies(id, Direction::DependsOn, request)
2643 .await
2644 }
2645 }
2646}
2647
2648#[expect(
2657 clippy::too_many_arguments,
2658 reason = "every argument is one axis of one walk — the source, the item, the \
2659 direction, which of its two graphs, whether the reverse is emulated, where \
2660 to resume, how many rows to return and where to count calls. Grouping them \
2661 into a struct would name the same eight values one indirection further from \
2662 the loop that reads them."
2663)]
2664async fn fetch_edges(
2665 source: &ResolvedSource,
2666 native: &NativeId,
2667 direction: Direction,
2668 entity: Entity,
2669 emulating: bool,
2670 start: &Resume,
2671 budget: u32,
2672 calls: &AtomicU32,
2673) -> Result<Fetched<DependencyEdge>, SourceError> {
2674 let ceiling = ceiling(source);
2675 if !emulating {
2676 return walk(
2677 start,
2678 budget,
2679 page_size(false, budget, ceiling),
2680 |_| true,
2681 |cursor, limit| async move {
2682 calls.fetch_add(1, Ordering::Relaxed);
2683 let request = PageRequest { cursor, limit };
2684 match entity {
2685 Entity::Task => {
2686 source
2687 .source()
2688 .task_dependencies(native, direction, &request)
2689 .await
2690 }
2691 Entity::Project => {
2692 source
2693 .source()
2694 .project_dependencies(native, direction, &request)
2695 .await
2696 }
2697 }
2698 },
2699 )
2700 .await;
2701 }
2702
2703 walk(
2704 start,
2705 budget,
2706 ceiling,
2707 |_| true,
2708 |cursor, limit| async move {
2709 calls.fetch_add(1, Ordering::Relaxed);
2710 let request = PageRequest { cursor, limit };
2711 let (ids, next) = match entity {
2712 Entity::Task => {
2713 let page = source
2714 .source()
2715 .query_tasks(&TaskQuery::default(), &request)
2716 .await?;
2717 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2718 (ids, page.next)
2719 }
2720 Entity::Project => {
2721 let page = source
2722 .source()
2723 .query_projects(&ProjectQuery::default(), &request)
2724 .await?;
2725 let ids: Vec<NativeId> =
2726 page.items.into_iter().map(|project| project.id).collect();
2727 (ids, page.next)
2728 }
2729 };
2730
2731 let mut edges = Vec::new();
2732 for id in ids {
2733 let mut inner: Option<Cursor> = None;
2734 loop {
2735 calls.fetch_add(1, Ordering::Relaxed);
2736 let request = PageRequest {
2737 cursor: inner.clone(),
2738 limit,
2739 };
2740 let page = forward_edges(source, entity, &id, &request).await?;
2741 fits(page.items.len(), limit)?;
2745 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2746 unrepeated(
2747 page.next.as_ref(),
2748 inner.as_ref(),
2749 "its forward edges were being scanned",
2750 )?;
2751 match page.next {
2752 Some(cursor) => inner = Some(cursor),
2753 None => break,
2754 }
2755 }
2756 }
2757
2758 Ok(Page { items: edges, next })
2759 },
2760 )
2761 .await
2762}