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 MetadataRecord, NativeId, Page, PageRequest, Priority, Project, ProjectFilter, ProjectQuery,
38 SecretResolver, SourceError, SourceName, StatusCategory, Task, TaskQuery, TextFields,
39 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 use copy::{
57 BudgetSpent, CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy,
58 Spent,
59};
60pub use delivery::{Delivered, DeliveryOutcome, TaskStatusSet, settled};
61pub use local::ProjectSelector;
62pub use metadata::MetadataSet;
63pub use narrow::{TaskContentSet, TaskPrioritySet};
64pub use rendered::{
65 Body, DocumentCreate, Regenerated, Regeneration, RenderRequest, RenderTemplate, RenderedRecord,
66 TaskCreate, TaskCreated, TemplateAnswers, UnusedAnswers,
67};
68pub use update::TaskUpdated;
69
70#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
75pub struct Qualified<T> {
76 pub id: GlobalId,
78 pub item: T,
80}
81
82#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
87pub struct QualifiedEdge {
88 pub from: QualifiedEndpoint,
93 pub to: QualifiedEndpoint,
95 pub kind: onetaskgraph_plugin_api::DependencyKind,
97}
98
99#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
101pub struct QualifiedEndpoint {
102 pub id: GlobalId,
104 pub kind: onetaskgraph_plugin_api::ItemKind,
106}
107
108impl std::fmt::Display for QualifiedEndpoint {
109 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
110 self.id.fmt(formatter)
111 }
112}
113
114#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
116#[serde(tag = "kind", rename_all = "kebab-case")]
117pub enum SearchHit {
118 Task(Qualified<Task>),
120 Project(Qualified<Project>),
122}
123
124#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
126#[serde(rename_all = "kebab-case")]
127pub enum SearchKind {
128 Tasks,
130 Projects,
132 #[default]
134 Both,
135}
136
137#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
139pub struct SourceListing {
140 pub source: SourceName,
142 pub kind: String,
154 #[serde(flatten)]
156 pub state: SourceState,
157}
158
159#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
161#[serde(tag = "state", rename_all = "kebab-case")]
162pub enum SourceState {
163 Available {
165 capabilities: Capabilities,
167 },
168 Unavailable {
170 error: SourceError,
172 },
173}
174
175#[derive(Debug, Clone, PartialEq)]
177pub struct Paging {
178 pub limit: NonZeroU32,
180 pub token: Option<PageToken>,
182}
183
184#[derive(Debug, Clone, Default, PartialEq)]
186pub struct Filters {
187 pub text: Option<TextQuery>,
189 pub labels: LabelFilter,
191 pub statuses: Vec<StatusCategory>,
193}
194
195#[derive(Debug, Clone)]
197pub struct TaskRequest {
198 pub sources: Vec<SourceName>,
200 pub filters: Filters,
202 pub project: ProjectSelector,
204 pub priorities: Vec<Priority>,
210 pub commented_since: Option<DateTime<Utc>>,
217 pub paging: Paging,
219}
220
221#[derive(Debug, Clone)]
223pub struct ProjectRequest {
224 pub sources: Vec<SourceName>,
226 pub filters: Filters,
228 pub paging: Paging,
230}
231
232#[derive(Debug, Clone, Default, PartialEq)]
239pub struct DocumentFilters {
240 pub text: Option<TextQuery>,
242 pub labels: LabelFilter,
244}
245
246#[derive(Debug, Clone)]
248pub struct DocumentRequest {
249 pub sources: Vec<SourceName>,
251 pub filters: DocumentFilters,
253 pub project: ProjectSelector,
255 pub paging: Paging,
257}
258
259#[derive(Debug, Clone)]
261pub struct LabelRequest {
262 pub sources: Vec<SourceName>,
264 pub paging: Paging,
266}
267
268#[derive(Debug, Clone)]
270pub struct SearchRequest {
271 pub sources: Vec<SourceName>,
273 pub text: TextQuery,
275 pub kind: SearchKind,
277 pub paging: Paging,
279}
280
281#[derive(Debug, Clone)]
283pub struct DependencyRequest {
284 pub id: GlobalId,
286 pub direction: Direction,
288 pub paging: Paging,
290}
291
292#[derive(Debug, Clone, PartialEq, thiserror::Error)]
300pub enum EngineError {
301 #[error(
303 "no source named {name:?} is configured\n\
304 next: name one of the configured sources ({configured}), or add {name:?} under \
305 `sources` — `onetaskgraph sources list` shows what this configuration has."
306 )]
307 UnknownSource {
308 name: String,
310 configured: String,
312 },
313
314 #[error(
316 "{message}\n\
317 next: page with a token exactly as the previous page reported it, and against \
318 the same configuration — or drop `--page` to start the walk again."
319 )]
320 Token {
321 message: String,
323 },
324
325 #[error(
327 "no sources are configured\n\
328 next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
329 prints what each plugin accepts."
330 )]
331 NoSources,
332
333 #[error(
335 "source {name} cannot be written: its plugin is {kind}, which has no write \
336 side\n\
337 next: copy into a source whose plugin can be written — `onetaskgraph sources \
338 list` reports each one's plugin."
339 )]
340 NotWritable {
341 name: String,
343 kind: String,
345 },
346
347 #[error(
354 "source {name} has no documents: its plugin is {kind}, which holds none\n\
355 next: name a source whose plugin has documents — `onetaskgraph sources list` \
356 reports each one's plugin and what it declares."
357 )]
358 NoDocuments {
359 name: String,
361 kind: String,
363 },
364
365 #[error(
370 "source {name} has no comments: its plugin is {kind}, whose tasks hold none\n\
371 next: name a task of a source whose plugin has comments — `onetaskgraph sources \
372 list` reports each one's plugin."
373 )]
374 NoComments {
375 name: String,
377 kind: String,
379 },
380
381 #[error(
383 "source {name} cannot be written: its plugin is {kind}, whose comments can be read \
384 but not added to, edited or removed\n\
385 next: list them with `onetaskgraph task comment list`, or comment on a task of a \
386 source whose plugin can be written — `onetaskgraph sources list` reports each \
387 one's plugin."
388 )]
389 CommentsNotWritable {
390 name: String,
392 kind: String,
394 },
395
396 #[error(
398 "source {name} cannot write a status: its plugin is {kind}, which has no write side\n\
399 next: set the status in that source itself, or name a task of a source whose plugin \
400 can be written — `onetaskgraph sources list` reports each one's plugin."
401 )]
402 StatusNotWritable {
403 name: String,
405 kind: String,
407 },
408
409 #[error(
411 "source {name} cannot write a {record}'s metadata: its plugin is {kind}, which has no \
412 write side\n\
413 next: set the key in that source itself, or name a {record} of a source whose plugin \
414 can be written — `onetaskgraph sources list` reports each one's plugin."
415 )]
416 MetadataNotWritable {
417 name: String,
419 kind: String,
421 record: MetadataRecord,
423 },
424
425 #[error(
427 "source {name} cannot write a priority: its plugin is {kind}, which has no write side\n\
428 next: set the priority in that source itself, or name a task of a source whose plugin \
429 can be written — `onetaskgraph sources list` reports each one's plugin."
430 )]
431 PriorityNotWritable {
432 name: String,
434 kind: String,
436 },
437
438 #[error(
440 "source {name} cannot write a task's content: its plugin is {kind}, which has no write \
441 side\n\
442 next: edit the content in that source itself, or name a task of a source whose plugin \
443 can be written — `onetaskgraph sources list` reports each one's plugin."
444 )]
445 ContentNotWritable {
446 name: String,
448 kind: String,
450 },
451
452 #[error(
454 "source {name} cannot update a task: its plugin is {kind}, which has no write side\n\
455 next: change the task in that source itself, or name a task of a source whose plugin \
456 can be written — `onetaskgraph sources list` reports each one's plugin."
457 )]
458 UpdateNotWritable {
459 name: String,
461 kind: String,
463 },
464
465 #[error(
467 "source {name} cannot create a {record}: its plugin is {kind}, which has no write side\n\
468 next: create it in a source whose plugin can be written — `onetaskgraph sources list` \
469 reports each one's plugin."
470 )]
471 NotCreatable {
472 name: String,
474 kind: String,
476 record: MetadataRecord,
478 },
479
480 #[error(
482 "source {name} cannot write a {record}'s rendering: its plugin is {kind}, which has no \
483 write side\n\
484 next: render it with --dry-run to read the result, or regenerate a {record} of a source \
485 whose plugin can be written."
486 )]
487 RenderingNotWritable {
488 name: String,
490 kind: String,
492 record: RenderedRecord,
494 },
495
496 #[error(
499 "supply every required answer to regenerate {id}: {} unanswered, and {reason}\n\
500 next: answer {} with --var NAME=VALUE or an answers file (--answers FILE), or run \
501 interactively to be asked.",
502 names.join(", "),
503 if names.len() == 1 { "it" } else { "each" }
504 )]
505 MissingAnswers {
506 id: String,
508 names: Vec<String>,
510 reason: String,
513 },
514
515 #[error("{error}")]
517 Template {
518 error: crate::template::TemplateError,
520 },
521
522 #[error(
524 "{record} {id} has no stored template answers: {reason}\n\
525 next: regenerate it with every required answer (`onetaskgraph {record} render {id} \
526 --var NAME=VALUE`), or read its provenance with `onetaskgraph {record} show {id}`."
527 )]
528 NoStoredAnswers {
529 record: RenderedRecord,
531 id: String,
533 reason: String,
535 },
536
537 #[error(
539 "{record} {id} records no template it was rendered from, and none was given\n\
540 next: name one with --template FILE or --template-loader FILE."
541 )]
542 NoTemplate {
543 record: RenderedRecord,
545 id: String,
547 },
548
549 #[error(
552 "{record} {id} records a template entry this product did not write — {problem}\n\
553 next: regenerate it with --template FILE or --template-loader FILE and every required \
554 answer, which records a fresh entry."
555 )]
556 MalformedProvenance {
557 record: RenderedRecord,
559 id: String,
561 problem: String,
563 },
564
565 #[error(
567 "{record} {id} was rendered from {reference:?}, which is not a readable file, so it \
568 cannot be re-read; a recorded reference is never turned into a location\n\
569 next: supply the template with --template-loader FILE (a loader document naming what \
570 to render), or name a template file with --template FILE."
571 )]
572 TemplateNotAFile {
573 record: RenderedRecord,
575 id: String,
577 reference: String,
579 },
580
581 #[error(
586 "source {name} cannot hold the field priority, so {task}'s priority {priority} cannot \
587 be written to it: its plugin is {kind}, which declares priority unsupported\n\
588 next: write to a source whose plugin holds a priority — `onetaskgraph sources list` \
589 reports what each declares — or set the task's priority to none first; a \
590 github-projects source holds one once its configuration sets priority_mapping."
591 )]
592 NoPriority {
593 name: String,
595 kind: String,
597 task: String,
599 priority: onetaskgraph_plugin_api::Priority,
601 },
602
603 #[error(
605 "no project with the id {id}\n\
606 next: check the id, or list what is there — `onetaskgraph project list` reports every \
607 project the configured sources hold."
608 )]
609 NoSuchProject {
610 id: String,
612 },
613
614 #[error(
616 "no document with the id {id}\n\
617 next: check the id, or list what is there — `onetaskgraph document list` reports every \
618 document the configured sources hold."
619 )]
620 NoSuchDocument {
621 id: String,
623 },
624
625 #[error(
627 "no task with the id {id}\n\
628 next: check the id, or list what is there — `onetaskgraph task list` reports every \
629 task the configured sources hold."
630 )]
631 NoSuchTask {
632 id: String,
634 },
635
636 #[error(
638 "task {task} has no comment with the id {comment}\n\
639 next: list its comments — `onetaskgraph task comment list {task}` reports each \
640 one's id."
641 )]
642 NoSuchComment {
643 task: String,
645 comment: String,
647 },
648
649 #[error(
654 "source {name} could not be built: {error}\n\
655 next: fix that source — `onetaskgraph sources list` reports its state — then run the \
656 command again."
657 )]
658 SourceUnavailable {
659 name: String,
661 error: SourceError,
663 },
664
665 #[error(
671 "source {name} could not do it: {error}\n\
672 next: fix what the source named above, then run the command again."
673 )]
674 SourceFailed {
675 name: String,
677 error: SourceError,
679 },
680
681 #[error(
683 "the destination source {name} could not be built: {error}\n\
684 next: fix that source — `onetaskgraph sources list` reports its state — then \
685 copy again."
686 )]
687 DestinationUnavailable {
688 name: String,
690 error: SourceError,
692 },
693
694 #[error(
696 "no item with the id {id}\n\
697 next: check the id, or list what is there — `onetaskgraph task list` and \
698 `onetaskgraph project list` report what the configured sources hold."
699 )]
700 NoSuchItem {
701 id: String,
703 },
704
705 #[error(
710 "{item} was copied from {origin}, which that destination no longer holds\n\
711 next: re-run with --recreate to create a new item there instead, or restore \
712 {origin}."
713 )]
714 StaleOrigin {
715 item: String,
717 origin: String,
719 },
720
721 #[error(
727 "{id} is not a task of {project}, so it cannot be copied as one of its members\n\
728 next: name a task `onetaskgraph task list --project {project}` reports, or copy \
729 {id} on its own with `onetaskgraph task copy`.",
730 project = .projects.iter().map(ToString::to_string).collect::<Vec<_>>().join(", ")
731 )]
732 NotAMember {
733 id: GlobalId,
735 projects: Vec<GlobalId>,
737 },
738
739 #[error(
748 "{item} depends on {member}, which this copy was not told to carry and which records \
749 no origin in {destination}\n\
750 next: name {member} with --member as well, record its {destination} id at \
751 onetaskgraph.origin, or copy the whole project without --member."
752 )]
753 UnrecordedMember {
754 item: GlobalId,
756 member: GlobalId,
758 destination: SourceName,
760 },
761
762 #[error(
768 "source {name} could not do it: {error}\n\
769 next: fix what the source named above, then copy again."
770 )]
771 SourceRefused {
772 name: String,
774 error: SourceError,
776 },
777
778 #[error(
786 "the copy failed and could not be undone.\n\
787 it failed because: {error}\n\
788 it could not be undone because: {refusal}\n\
789 so the destination still holds: {left_behind}\n\
790 next: remove those items at the destination, then copy again."
791 )]
792 CopyNotUndone {
793 error: Box<EngineError>,
795 left_behind: LeftBehind,
801 refusal: SourceError,
803 },
804}
805
806#[derive(Debug, Clone, PartialEq)]
815pub struct LeftBehind {
816 first: GlobalId,
818 rest: Vec<GlobalId>,
820}
821
822impl LeftBehind {
823 #[must_use]
825 pub fn new(first: GlobalId) -> Self {
826 Self {
827 first,
828 rest: Vec::new(),
829 }
830 }
831
832 pub fn push(&mut self, id: GlobalId) {
834 self.rest.push(id);
835 }
836
837 pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
839 std::iter::once(&self.first).chain(self.rest.iter())
840 }
841}
842
843impl std::fmt::Display for LeftBehind {
844 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
846 write!(formatter, "{}", self.first)?;
847 for id in &self.rest {
848 write!(formatter, ", {id}")?;
849 }
850 Ok(())
851 }
852}
853
854pub enum ConfiguredSource {
861 Ready(ResolvedSource),
863 Unavailable(UnavailableSource),
865}
866
867impl ConfiguredSource {
868 #[must_use]
870 pub fn name(&self) -> &SourceName {
871 match self {
872 Self::Ready(source) => source.name(),
873 Self::Unavailable(source) => source.name(),
874 }
875 }
876}
877
878pub struct Engine {
880 sources: Vec<ConfiguredSource>,
882 selection: Vec<SourceName>,
884}
885
886impl Engine {
887 #[must_use]
895 pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
896 let (ready, unavailable) = resolve_available(config, secrets);
897 Self::new(
898 ready
899 .into_iter()
900 .map(ConfiguredSource::Ready)
901 .chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
902 .collect(),
903 config.selected_sources(),
904 )
905 }
906
907 #[must_use]
910 pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
911 Self { sources, selection }
912 }
913
914 fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
916 self.sources.iter().filter_map(|source| match source {
917 ConfiguredSource::Ready(ready) => Some(ready),
918 ConfiguredSource::Unavailable(_) => None,
919 })
920 }
921
922 fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
924 self.sources.iter().filter_map(|source| match source {
925 ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
926 ConfiguredSource::Ready(_) => None,
927 })
928 }
929
930 #[must_use]
932 pub fn listing(&self) -> Vec<SourceListing> {
933 let mut listings: Vec<SourceListing> = self
934 .ready()
935 .map(|source| SourceListing {
936 source: source.name().clone(),
937 kind: source.kind().to_owned(),
938 state: SourceState::Available {
939 capabilities: source.source().capabilities(),
940 },
941 })
942 .chain(self.unavailable().map(|source| SourceListing {
943 source: source.name().clone(),
944 kind: source.kind().to_owned(),
945 state: SourceState::Unavailable {
946 error: source.error().clone(),
947 },
948 }))
949 .collect();
950 listings.sort_by(|left, right| left.source.cmp(&right.source));
951 listings
952 }
953
954 #[must_use]
960 pub fn has(&self, name: &SourceName) -> bool {
961 self.sources.iter().any(|source| source.name() == name)
962 }
963
964 pub async fn tasks(
972 &self,
973 request: &TaskRequest,
974 ) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
975 let mut names = self.resolve_selection(&request.sources)?;
976 if let ProjectSelector::Qualified(id) = &request.project {
981 self.known(&id.source)?;
982 names.retain(|name| name == &id.source);
983 }
984 let query = shape(
985 "task-list",
986 &names,
987 &(
988 &request.filters,
989 &request.project,
990 &request.priorities,
991 &request.commented_since,
992 ),
993 );
994 let states = resumption(
995 self,
996 request.paging.token.as_ref(),
997 &[StreamKind::Items],
998 &query,
999 )?;
1000 let budget = request.paging.limit.get();
1001
1002 let mut answer = Answer::new();
1003 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1004
1005 let shapes: Vec<TaskShape> = ready
1006 .iter()
1007 .map(|source| {
1008 shape_tasks(
1009 &source.source().capabilities(),
1010 &request.filters,
1011 &project_filter(&request.project),
1012 &request.priorities,
1013 request.commented_since,
1014 )
1015 })
1016 .collect();
1017 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1018 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1019
1020 let walks = ready
1021 .iter()
1022 .enumerate()
1023 .map(|(index, source)| {
1024 fetch_tasks(
1025 source,
1026 &shapes[index],
1027 &starts[index],
1028 budget,
1029 &counters[index],
1030 )
1031 })
1032 .collect();
1033
1034 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1035 answer.finish(
1036 streams,
1037 budget,
1038 owed(&states),
1039 &query,
1040 |name, task: Task| {
1041 delivery::qualified_task(GlobalId::new(name.clone(), task.id.clone()), task)
1042 },
1043 )
1044 }
1045
1046 pub async fn projects(
1052 &self,
1053 request: &ProjectRequest,
1054 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1055 let names = self.resolve_selection(&request.sources)?;
1056 let query = shape("project-list", &names, &request.filters);
1057 let states = resumption(
1058 self,
1059 request.paging.token.as_ref(),
1060 &[StreamKind::Items],
1061 &query,
1062 )?;
1063 let budget = request.paging.limit.get();
1064
1065 let mut answer = Answer::new();
1066 let mut with_projects = Vec::new();
1072 for source in answer.split(self, &names) {
1073 if source.source().capabilities().projects.is_native() {
1074 with_projects.push(source);
1075 } else {
1076 answer.unreachable_predicate(source, Predicate::Project);
1077 }
1078 }
1079 let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
1080
1081 let shapes: Vec<ProjectShape> = ready
1082 .iter()
1083 .map(|source| shape_projects(&source.source().capabilities(), &request.filters))
1084 .collect();
1085 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1086 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1087
1088 let walks = ready
1089 .iter()
1090 .enumerate()
1091 .map(|(index, source)| {
1092 fetch_projects(
1093 source,
1094 &shapes[index],
1095 &starts[index],
1096 budget,
1097 &counters[index],
1098 )
1099 })
1100 .collect();
1101
1102 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1103 answer.finish(
1104 streams,
1105 budget,
1106 owed(&states),
1107 &query,
1108 |name, project: Project| Qualified {
1109 id: GlobalId::new(name.clone(), project.id.clone()),
1110 item: project,
1111 },
1112 )
1113 }
1114
1115 pub async fn documents(
1127 &self,
1128 request: &DocumentRequest,
1129 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1130 let mut names = self.resolve_selection(&request.sources)?;
1131 if let ProjectSelector::Qualified(id) = &request.project {
1134 self.known(&id.source)?;
1135 names.retain(|name| name == &id.source);
1136 }
1137 let query = shape(
1138 "document-list",
1139 &names,
1140 &(&request.filters, &request.project),
1141 );
1142 let states = resumption(
1143 self,
1144 request.paging.token.as_ref(),
1145 &[StreamKind::Items],
1146 &query,
1147 )?;
1148 let budget = request.paging.limit.get();
1149
1150 let mut answer = Answer::new();
1151 let mut with_documents = Vec::new();
1152 for source in answer.split(self, &names) {
1153 if source.source().capabilities().documents.is_native() {
1154 with_documents.push(source);
1155 } else {
1156 answer.unreachable_predicate(source, Predicate::Document);
1157 }
1158 }
1159 let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
1160
1161 let shapes: Vec<DocumentShape> = ready
1162 .iter()
1163 .map(|source| {
1164 shape_documents(
1165 &source.source().capabilities(),
1166 &request.filters,
1167 &project_filter(&request.project),
1168 )
1169 })
1170 .collect();
1171 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1172 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1173
1174 let walks = ready
1175 .iter()
1176 .enumerate()
1177 .map(|(index, source)| {
1178 fetch_documents(
1179 source,
1180 &shapes[index],
1181 &starts[index],
1182 budget,
1183 &counters[index],
1184 )
1185 })
1186 .collect();
1187
1188 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1189 answer.finish(
1190 streams,
1191 budget,
1192 owed(&states),
1193 &query,
1194 |name, document: Document| Qualified {
1195 id: GlobalId::new(name.clone(), document.id.clone()),
1196 item: document,
1197 },
1198 )
1199 }
1200
1201 pub async fn labels(
1207 &self,
1208 request: &LabelRequest,
1209 ) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
1210 let names = self.resolve_selection(&request.sources)?;
1211 let query = shape("label-list", &names, &());
1212 let states = resumption(
1213 self,
1214 request.paging.token.as_ref(),
1215 &[StreamKind::Items],
1216 &query,
1217 )?;
1218 let budget = request.paging.limit.get();
1219
1220 let mut answer = Answer::new();
1221 let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
1222 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1223 let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
1224
1225 let walks = ready
1226 .iter()
1227 .enumerate()
1228 .map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
1229 .collect();
1230
1231 let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
1232 answer.finish(
1233 streams,
1234 budget,
1235 owed(&states),
1236 &query,
1237 |name, label: Label| Qualified {
1238 id: GlobalId::new(name.clone(), label.id.clone()),
1239 item: label,
1240 },
1241 )
1242 }
1243
1244 pub async fn search(
1250 &self,
1251 request: &SearchRequest,
1252 ) -> Result<QueryResponse<SearchHit>, EngineError> {
1253 let names = self.resolve_selection(&request.sources)?;
1254 let reads: &[StreamKind] = match request.kind {
1258 SearchKind::Tasks => &[StreamKind::Tasks],
1259 SearchKind::Projects => &[StreamKind::Projects],
1260 SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
1261 };
1262 let query = shape("search", &names, &request.text);
1267 let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
1268 let budget = request.paging.limit.get();
1269 let filters = Filters {
1270 text: Some(request.text.clone()),
1271 ..Filters::default()
1272 };
1273
1274 let mut answer = Answer::new();
1275
1276 let mut ready = Vec::new();
1279 let mut kinds = Vec::new();
1280 let mut starts = Vec::new();
1281 for source in answer.split(self, &names) {
1282 let mut streams = Vec::new();
1283 if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
1284 streams.push(StreamKind::Tasks);
1285 }
1286 if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
1287 if source.source().capabilities().projects.is_native() {
1288 streams.push(StreamKind::Projects);
1289 } else {
1290 answer.unreachable_predicate(source, Predicate::Project);
1291 }
1292 }
1293 for stream in streams {
1294 if let Some(resume) = resume_at(&states, source.name(), stream) {
1295 ready.push(source);
1296 kinds.push(stream);
1297 starts.push(resume);
1298 }
1299 }
1300 }
1301
1302 let shapes: Vec<HitShape> = ready
1303 .iter()
1304 .zip(kinds.iter())
1305 .map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
1306 .collect();
1307 let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
1308 let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
1309
1310 let walks = ready
1311 .iter()
1312 .enumerate()
1313 .map(|(index, source)| {
1314 fetch_hits(
1315 source,
1316 &shapes[index],
1317 &starts[index],
1318 budget,
1319 &counters[index],
1320 )
1321 })
1322 .collect();
1323
1324 let streams =
1325 answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
1326 answer.finish(
1327 streams,
1328 budget,
1329 owed(&states),
1330 &query,
1331 |name, found: Found| match found {
1332 Found::Task(task) => SearchHit::Task(delivery::qualified_task(
1333 GlobalId::new(name.clone(), task.id.clone()),
1334 task,
1335 )),
1336 Found::Project(project) => SearchHit::Project(Qualified {
1337 id: GlobalId::new(name.clone(), project.id.clone()),
1338 item: project,
1339 }),
1340 },
1341 )
1342 }
1343
1344 pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
1351 let name = self.known(&id.source)?;
1352 let mut answer = Answer::new();
1353 let selected = answer.split(self, std::slice::from_ref(&name));
1354 let Some(source) = selected.first() else {
1355 return answer.nothing();
1356 };
1357 let found = source.source().get_task(&id.native).await;
1358 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1359 answer.one(source, found, |task| {
1360 delivery::qualified_task(qualified, task)
1361 })
1362 }
1363
1364 pub async fn project(
1370 &self,
1371 id: &GlobalId,
1372 ) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
1373 let name = self.known(&id.source)?;
1374 let mut answer = Answer::new();
1375 let selected = answer.split(self, std::slice::from_ref(&name));
1376 let Some(source) = selected.first() else {
1377 return answer.nothing();
1378 };
1379 let found = source.source().get_project(&id.native).await;
1380 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1381 answer.one(source, found, |project| Qualified {
1382 id: qualified,
1383 item: project,
1384 })
1385 }
1386
1387 pub async fn document(
1397 &self,
1398 id: &GlobalId,
1399 ) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
1400 let name = self.known(&id.source)?;
1401 let mut answer = Answer::new();
1402 let selected = answer.split(self, std::slice::from_ref(&name));
1403 let Some(source) = selected.first() else {
1404 return answer.nothing();
1405 };
1406 if !source.source().capabilities().documents.is_native() {
1407 answer.unreachable_predicate(source, Predicate::Document);
1408 return answer.nothing();
1409 }
1410 let found = source.source().get_document(&id.native).await;
1411 let qualified = GlobalId::new(source.name().clone(), id.native.clone());
1412 answer.one(source, found, |document| Qualified {
1413 id: qualified,
1414 item: document,
1415 })
1416 }
1417
1418 pub async fn task_dependencies(
1425 &self,
1426 request: &DependencyRequest,
1427 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1428 self.dependencies(request, Entity::Task).await
1429 }
1430
1431 pub async fn project_dependencies(
1437 &self,
1438 request: &DependencyRequest,
1439 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1440 self.dependencies(request, Entity::Project).await
1441 }
1442
1443 async fn dependencies(
1446 &self,
1447 request: &DependencyRequest,
1448 entity: Entity,
1449 ) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
1450 let name = self.known(&request.id.source)?;
1451 let query = shape(
1452 "dependencies",
1453 std::slice::from_ref(&name),
1454 &(entity, &request.id.native, request.direction),
1455 );
1456 let states = resumption(
1457 self,
1458 request.paging.token.as_ref(),
1459 &[StreamKind::Items],
1460 &query,
1461 )?;
1462 let budget = request.paging.limit.get();
1463
1464 let mut answer = Answer::new();
1465 let (ready, starts) = walking(
1466 answer.split(self, std::slice::from_ref(&name)),
1467 &states,
1468 StreamKind::Items,
1469 );
1470 let Some(source) = ready.first() else {
1471 return answer.nothing();
1472 };
1473
1474 let capabilities = source.source().capabilities();
1475 let support = match entity {
1476 Entity::Task => capabilities.task_dependencies,
1477 Entity::Project => capabilities.project_dependencies,
1478 };
1479 let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
1483 let mut outcomes = Outcomes::default();
1484 if request.direction == Direction::DependedOnBy {
1485 if emulating {
1486 outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
1487 } else {
1488 outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
1489 }
1490 }
1491
1492 let counters = vec![AtomicU32::new(0)];
1493 let walked = fetch_edges(
1494 source,
1495 &request.id.native,
1496 request.direction,
1497 entity,
1498 emulating,
1499 &starts[0],
1500 budget,
1501 &counters[0],
1502 )
1503 .await;
1504
1505 let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
1506 answer.finish(
1507 streams,
1508 budget,
1509 owed(&states),
1510 &query,
1511 |name, edge: DependencyEdge| QualifiedEdge {
1512 from: qualify_endpoint(name, edge.from),
1513 to: qualify_endpoint(name, edge.to),
1514 kind: edge.kind,
1515 },
1516 )
1517 }
1518
1519 fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
1521 if asked.is_empty() {
1522 if self.selection.is_empty() {
1523 return Err(EngineError::NoSources);
1524 }
1525 return Ok(self.selection.clone());
1526 }
1527 asked.iter().map(|name| self.known(name)).collect()
1528 }
1529
1530 fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
1532 if self.has(name) {
1533 return Ok(name.clone());
1534 }
1535 if self.sources.is_empty() {
1536 return Err(EngineError::NoSources);
1537 }
1538 Err(EngineError::UnknownSource {
1539 name: name.to_string(),
1540 configured: self
1541 .listing()
1542 .iter()
1543 .map(|listing| listing.source.to_string())
1544 .collect::<Vec<_>>()
1545 .join(", "),
1546 })
1547 }
1548}
1549
1550fn qualify_endpoint(
1551 source: &SourceName,
1552 endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
1553) -> QualifiedEndpoint {
1554 let kind = endpoint.kind;
1555 let is_qualified = endpoint.is_qualified();
1556 let endpoint_id = endpoint.into_id();
1557 QualifiedEndpoint {
1558 id: if is_qualified {
1559 endpoint_id
1560 .parse()
1561 .expect("plugin-api validates qualified dependency endpoints")
1562 } else {
1563 GlobalId::new(source.clone(), NativeId(endpoint_id))
1564 },
1565 kind,
1566 }
1567}
1568
1569#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1571enum Entity {
1572 Task,
1574 Project,
1576}
1577
1578enum Found {
1580 Task(Task),
1582 Project(Project),
1584}
1585
1586#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1588enum Outcome {
1589 PushedDown,
1591 AppliedLocally,
1593 Emulated,
1595 Unavailable,
1597}
1598
1599#[derive(Debug, Clone, Default, PartialEq)]
1611struct Outcomes(BTreeMap<Predicate, Outcome>);
1612
1613impl Outcomes {
1614 fn record(&mut self, predicate: Predicate, outcome: Outcome) {
1620 self.0.insert(predicate, outcome);
1621 }
1622
1623 fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
1626 for predicate in predicates {
1627 self.record(predicate, outcome);
1628 }
1629 }
1630
1631 fn with(&self, outcome: Outcome) -> Vec<Predicate> {
1633 self.0
1634 .iter()
1635 .filter(|(_, recorded)| **recorded == outcome)
1636 .map(|(predicate, _)| *predicate)
1637 .collect()
1638 }
1639}
1640
1641struct TaskShape {
1643 pushed: TaskQuery,
1645 local: LocalTasks,
1647 outcomes: Outcomes,
1649}
1650
1651struct ProjectShape {
1653 pushed: ProjectQuery,
1655 local: LocalProjects,
1657 outcomes: Outcomes,
1659}
1660
1661struct DocumentShape {
1663 pushed: DocumentQuery,
1665 local: LocalDocuments,
1667 outcomes: Outcomes,
1669}
1670
1671struct HitShape {
1673 stream: StreamKind,
1675 tasks: TaskQuery,
1677 projects: ProjectQuery,
1679 local_tasks: LocalTasks,
1681 local_projects: LocalProjects,
1683 outcomes: Outcomes,
1685}
1686
1687struct Answer {
1693 plans: Vec<SourcePlan>,
1695 errors: Vec<SourceFailure>,
1697}
1698
1699impl Answer {
1700 fn new() -> Self {
1701 Self {
1702 plans: Vec::new(),
1703 errors: Vec::new(),
1704 }
1705 }
1706
1707 fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
1712 let mut selected = Vec::new();
1713 for name in names {
1714 match engine.sources.iter().find(|source| source.name() == name) {
1715 Some(ConfiguredSource::Ready(source)) => selected.push(source),
1716 Some(ConfiguredSource::Unavailable(source)) => {
1717 self.errors.push(source.failure());
1718 }
1719 None => {}
1720 }
1721 }
1722 selected
1723 }
1724
1725 fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
1727 let mut outcomes = Outcomes::default();
1728 outcomes.record(predicate, Outcome::Unavailable);
1729 self.plans.push(plan_for(source, outcomes, 0));
1730 }
1731
1732 fn collect<T>(
1734 &mut self,
1735 ready: &[&ResolvedSource],
1736 walked: Vec<Result<Fetched<T>, SourceError>>,
1737 counters: &[AtomicU32],
1738 outcomes: Vec<Outcomes>,
1739 ) -> Vec<Stream<T>> {
1740 let kinds = vec![StreamKind::Items; ready.len()];
1741 self.collect_streams(ready, &kinds, walked, counters, outcomes)
1742 }
1743
1744 fn collect_streams<T>(
1746 &mut self,
1747 ready: &[&ResolvedSource],
1748 kinds: &[StreamKind],
1749 walked: Vec<Result<Fetched<T>, SourceError>>,
1750 counters: &[AtomicU32],
1751 outcomes: Vec<Outcomes>,
1752 ) -> Vec<Stream<T>> {
1753 let mut streams = Vec::new();
1754 for (index, result) in walked.into_iter().enumerate() {
1755 let source = ready[index];
1756 let pages = counters[index].load(Ordering::Relaxed);
1757 self.plans
1758 .push(plan_for(source, outcomes[index].clone(), pages));
1759 match result {
1760 Ok(fetched) => streams.push(Stream {
1761 source: source.name().clone(),
1762 kind: kinds[index],
1763 fetched,
1764 }),
1765 Err(error) => self.errors.push(SourceFailure {
1768 source: source.name().clone(),
1769 error,
1770 }),
1771 }
1772 }
1773 streams
1774 }
1775
1776 fn one<T, U>(
1778 mut self,
1779 source: &ResolvedSource,
1780 found: Result<Option<T>, SourceError>,
1781 qualify: impl FnOnce(T) -> U,
1782 ) -> Result<QueryResponse<U>, EngineError> {
1783 self.plans.push(plan_for(source, Outcomes::default(), 1));
1784 let items = match found {
1785 Ok(Some(item)) => vec![qualify(item)],
1786 Ok(None) => Vec::new(),
1787 Err(error) => {
1788 self.errors.push(SourceFailure {
1789 source: source.name().clone(),
1790 error,
1791 });
1792 Vec::new()
1793 }
1794 };
1795 Ok(QueryResponse {
1796 items,
1797 next: None,
1798 plan: QueryPlan {
1799 per_source: merge_plans(self.plans),
1800 },
1801 errors: self.errors,
1802 })
1803 }
1804
1805 fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
1807 Ok(QueryResponse {
1808 items: Vec::new(),
1809 next: None,
1810 plan: QueryPlan {
1811 per_source: merge_plans(self.plans),
1812 },
1813 errors: self.errors,
1814 })
1815 }
1816
1817 fn finish<T, U>(
1822 self,
1823 streams: Vec<Stream<T>>,
1824 budget: u32,
1825 first: Option<&Owed>,
1826 query: &str,
1827 qualify: impl Fn(&SourceName, T) -> U,
1828 ) -> Result<QueryResponse<U>, EngineError> {
1829 let (rows, states, owed) = merge(streams, budget, first);
1830 let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
1831 Ok(QueryResponse {
1832 items: rows
1833 .into_iter()
1834 .map(|(name, item)| qualify(&name, item))
1835 .collect(),
1836 next,
1837 plan: QueryPlan {
1838 per_source: merge_plans(self.plans),
1839 },
1840 errors: self.errors,
1841 })
1842 }
1843}
1844
1845fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
1848 SourcePlan {
1849 source: source.name().clone(),
1850 kind: source.kind().to_owned(),
1851 pushed_down: outcomes.with(Outcome::PushedDown),
1852 applied_locally: outcomes.with(Outcome::AppliedLocally),
1853 emulated: outcomes.with(Outcome::Emulated),
1854 unavailable: outcomes.with(Outcome::Unavailable),
1855 pages_fetched: pages,
1856 }
1857}
1858
1859fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
1864 let mut merged: Vec<SourcePlan> = Vec::new();
1865 for plan in plans {
1866 if let Some(existing) = merged
1867 .iter_mut()
1868 .find(|existing| existing.source == plan.source)
1869 {
1870 existing.pushed_down.extend(plan.pushed_down);
1871 existing.applied_locally.extend(plan.applied_locally);
1872 existing.emulated.extend(plan.emulated);
1873 existing.unavailable.extend(plan.unavailable);
1874 existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
1875 for list in [
1876 &mut existing.pushed_down,
1877 &mut existing.applied_locally,
1878 &mut existing.emulated,
1879 &mut existing.unavailable,
1880 ] {
1881 list.sort_unstable();
1882 list.dedup();
1883 }
1884 } else {
1885 merged.push(plan);
1886 }
1887 }
1888 merged
1889}
1890
1891fn walking<'a>(
1897 selected: Vec<&'a ResolvedSource>,
1898 states: &Option<Resumption>,
1899 kind: StreamKind,
1900) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
1901 let mut ready = Vec::new();
1902 let mut starts = Vec::new();
1903 for source in selected {
1904 if let Some(resume) = resume_at(states, source.name(), kind) {
1905 ready.push(source);
1906 starts.push(resume);
1907 }
1908 }
1909 (ready, starts)
1910}
1911
1912fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
1930 let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
1931 fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
1932}
1933
1934fn fingerprint(text: &str) -> String {
1936 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1937 for byte in text.as_bytes() {
1938 hash ^= u64::from(*byte);
1939 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1940 }
1941 format!("{hash:016x}")
1942}
1943
1944fn owed(document: &Option<Resumption>) -> Option<&Owed> {
1949 document.as_ref()?.owed.as_ref()
1950}
1951
1952fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
1954 match states {
1955 None => Some(Resume::default()),
1956 Some(document) => document
1957 .streams
1958 .iter()
1959 .find(|state| &state.source == source && state.stream == kind)
1960 .map(|state| state.resume.clone()),
1961 }
1962}
1963
1964fn resumption(
1995 engine: &Engine,
1996 token: Option<&PageToken>,
1997 reads: &[StreamKind],
1998 query: &str,
1999) -> Result<Option<Resumption>, EngineError> {
2000 let Some(document) = token.map(PageToken::decode) else {
2001 return Ok(None);
2002 };
2003
2004 if document.query != query {
2011 return Err(EngineError::Token {
2012 message: "this page token was written by a different query — resume the walk it \
2013 came from, or drop --page to start this one from the beginning"
2014 .to_owned(),
2015 });
2016 }
2017 let states = &document.streams;
2018
2019 let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
2020 for state in states {
2021 if !reads.contains(&state.stream) {
2022 return Err(EngineError::Token {
2023 message: format!(
2024 "this page token resumes {}, which this command does not read — it \
2025 was written by a different query",
2026 state.stream.describe()
2027 ),
2028 });
2029 }
2030 let ceiling = engine
2031 .ready()
2032 .find(|source| source.name() == &state.source)
2033 .map(ceiling);
2034 if ceiling.is_none() && !engine.has(&state.source) {
2035 return Err(EngineError::Token {
2036 message: format!(
2037 "this page token resumes a source called {:?}, which this \
2038 configuration does not have",
2039 state.source.as_str()
2040 ),
2041 });
2042 }
2043 if let Some(ceiling) = ceiling
2044 && state.resume.skip >= ceiling
2045 {
2046 return Err(EngineError::Token {
2047 message: format!(
2048 "this page token resumes {} rows into a page of source {:?}, which \
2049 serves at most {ceiling}",
2050 state.resume.skip,
2051 state.source.as_str()
2052 ),
2053 });
2054 }
2055 if seen.contains(&(&state.source, state.stream)) {
2056 return Err(EngineError::Token {
2057 message: format!(
2058 "this page token gives source {:?} two places to resume from",
2059 state.source.as_str()
2060 ),
2061 });
2062 }
2063 seen.push((&state.source, state.stream));
2064 }
2065
2066 if let Some(owed) = &document.owed
2071 && !document
2072 .streams
2073 .iter()
2074 .any(|state| state.source == owed.source && state.stream == owed.stream)
2075 {
2076 return Err(EngineError::Token {
2077 message: format!(
2078 "this page token owes the next row to a stream it does not resume, \
2079 {:?}'s {}",
2080 owed.source.as_str(),
2081 owed.stream.describe()
2082 ),
2083 });
2084 }
2085
2086 Ok(Some(document))
2087}
2088
2089fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
2095 match selector {
2096 ProjectSelector::Any => ProjectFilter::Any,
2097 ProjectSelector::Orphans => ProjectFilter::Orphans,
2098 ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
2099 ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
2100 }
2101}
2102
2103fn text_predicates(fields: TextFields) -> Vec<Predicate> {
2105 match fields {
2106 TextFields::Title => vec![Predicate::SearchTitle],
2107 TextFields::Content => vec![Predicate::SearchContent],
2108 TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
2109 }
2110}
2111
2112fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
2119 match fields {
2120 TextFields::Title => capabilities.search_title.is_native(),
2121 TextFields::Content => capabilities.search_content.is_native(),
2122 TextFields::TitleOrContent => {
2123 capabilities.search_title.is_native() && capabilities.search_content.is_native()
2124 }
2125 }
2126}
2127
2128fn shape_tasks(
2130 capabilities: &Capabilities,
2131 filters: &Filters,
2132 project: &ProjectFilter,
2133 priorities: &[Priority],
2134 commented_since: Option<DateTime<Utc>>,
2135) -> TaskShape {
2136 let mut pushed = TaskQuery::default();
2137 let mut local = LocalTasks::default();
2138 let mut outcomes = Outcomes::default();
2139
2140 if !filters.labels.is_empty() {
2141 if capabilities.filter_by_label.is_native() {
2142 pushed.labels = filters.labels.clone();
2143 outcomes.record(Predicate::Label, Outcome::PushedDown);
2144 } else {
2145 local.labels = Some(filters.labels.clone());
2146 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2147 }
2148 }
2149 if !filters.statuses.is_empty() {
2150 if capabilities.filter_by_status.is_native() {
2151 pushed.statuses.clone_from(&filters.statuses);
2152 outcomes.record(Predicate::Status, Outcome::PushedDown);
2153 } else {
2154 local.statuses.clone_from(&filters.statuses);
2155 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2156 }
2157 }
2158 if !priorities.is_empty() {
2159 if capabilities.filter_by_priority.is_native() {
2160 pushed.priorities = priorities.to_vec();
2161 outcomes.record(Predicate::Priority, Outcome::PushedDown);
2162 } else {
2163 local.priorities = priorities.to_vec();
2164 outcomes.record(Predicate::Priority, Outcome::AppliedLocally);
2165 }
2166 }
2167 if let Some(since) = commented_since {
2168 if capabilities.filter_by_comment_activity.is_native() {
2169 pushed.commented_since = Some(since);
2170 outcomes.record(Predicate::CommentedSince, Outcome::PushedDown);
2171 } else {
2172 local.commented_since = Some(since);
2173 outcomes.record(Predicate::CommentedSince, Outcome::AppliedLocally);
2174 }
2175 }
2176 if let Some(text) = &filters.text {
2177 let predicates = text_predicates(text.fields);
2178 if searches_natively(capabilities, text.fields) {
2179 pushed.text = Some(text.clone());
2180 outcomes.record_all(predicates, Outcome::PushedDown);
2181 } else {
2182 local.text = Some(text.clone());
2183 outcomes.record_all(predicates, Outcome::AppliedLocally);
2184 }
2185 }
2186 match project {
2187 ProjectFilter::Any => {}
2188 ProjectFilter::Orphans => {
2189 if capabilities.orphan_tasks.is_native() {
2190 pushed.project = ProjectFilter::Orphans;
2191 outcomes.record(Predicate::Project, Outcome::PushedDown);
2192 } else {
2193 local.project = Some(ProjectFilter::Orphans);
2194 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2195 }
2196 }
2197 ProjectFilter::Is(id) => {
2198 if capabilities.projects.is_native() {
2199 pushed.project = ProjectFilter::Is(id.clone());
2200 outcomes.record(Predicate::Project, Outcome::PushedDown);
2201 } else {
2202 local.project = Some(ProjectFilter::Is(id.clone()));
2203 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2204 }
2205 }
2206 }
2207
2208 TaskShape {
2209 pushed,
2210 local,
2211 outcomes,
2212 }
2213}
2214
2215fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
2217 let mut pushed = ProjectQuery::default();
2218 let mut local = LocalProjects::default();
2219 let mut outcomes = Outcomes::default();
2220
2221 if !filters.labels.is_empty() {
2222 if capabilities.filter_by_label.is_native() {
2223 pushed.labels = filters.labels.clone();
2224 outcomes.record(Predicate::Label, Outcome::PushedDown);
2225 } else {
2226 local.labels = Some(filters.labels.clone());
2227 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2228 }
2229 }
2230 if !filters.statuses.is_empty() {
2231 if capabilities.filter_by_status.is_native() {
2232 pushed.statuses.clone_from(&filters.statuses);
2233 outcomes.record(Predicate::Status, Outcome::PushedDown);
2234 } else {
2235 local.statuses.clone_from(&filters.statuses);
2236 outcomes.record(Predicate::Status, Outcome::AppliedLocally);
2237 }
2238 }
2239 if let Some(text) = &filters.text {
2240 let predicates = text_predicates(text.fields);
2241 if searches_natively(capabilities, text.fields) {
2242 pushed.text = Some(text.clone());
2243 outcomes.record_all(predicates, Outcome::PushedDown);
2244 } else {
2245 local.text = Some(text.clone());
2246 outcomes.record_all(predicates, Outcome::AppliedLocally);
2247 }
2248 }
2249
2250 ProjectShape {
2251 pushed,
2252 local,
2253 outcomes,
2254 }
2255}
2256
2257fn shape_documents(
2262 capabilities: &Capabilities,
2263 filters: &DocumentFilters,
2264 project: &ProjectFilter,
2265) -> DocumentShape {
2266 let mut pushed = DocumentQuery::default();
2267 let mut local = LocalDocuments::default();
2268 let mut outcomes = Outcomes::default();
2269
2270 if !filters.labels.is_empty() {
2271 if capabilities.filter_by_label.is_native() {
2272 pushed.labels = filters.labels.clone();
2273 outcomes.record(Predicate::Label, Outcome::PushedDown);
2274 } else {
2275 local.labels = Some(filters.labels.clone());
2276 outcomes.record(Predicate::Label, Outcome::AppliedLocally);
2277 }
2278 }
2279 if let Some(text) = &filters.text {
2280 let predicates = text_predicates(text.fields);
2281 if searches_natively(capabilities, text.fields) {
2282 pushed.text = Some(text.clone());
2283 outcomes.record_all(predicates, Outcome::PushedDown);
2284 } else {
2285 local.text = Some(text.clone());
2286 outcomes.record_all(predicates, Outcome::AppliedLocally);
2287 }
2288 }
2289 match project {
2290 ProjectFilter::Any => {}
2291 ProjectFilter::Orphans => {
2292 if capabilities.orphan_tasks.is_native() {
2293 pushed.project = ProjectFilter::Orphans;
2294 outcomes.record(Predicate::Project, Outcome::PushedDown);
2295 } else {
2296 local.project = Some(ProjectFilter::Orphans);
2297 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2298 }
2299 }
2300 ProjectFilter::Is(id) => {
2301 if capabilities.projects.is_native() {
2302 pushed.project = ProjectFilter::Is(id.clone());
2303 outcomes.record(Predicate::Project, Outcome::PushedDown);
2304 } else {
2305 local.project = Some(ProjectFilter::Is(id.clone()));
2306 outcomes.record(Predicate::Project, Outcome::AppliedLocally);
2307 }
2308 }
2309 }
2310
2311 DocumentShape {
2312 pushed,
2313 local,
2314 outcomes,
2315 }
2316}
2317
2318fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
2320 match stream {
2321 StreamKind::Projects => {
2322 let shaped = shape_projects(capabilities, filters);
2323 HitShape {
2324 stream,
2325 tasks: TaskQuery::default(),
2326 projects: shaped.pushed,
2327 local_tasks: LocalTasks::default(),
2328 local_projects: shaped.local,
2329 outcomes: shaped.outcomes,
2330 }
2331 }
2332 StreamKind::Items | StreamKind::Tasks => {
2333 let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any, &[], None);
2334 HitShape {
2335 stream,
2336 tasks: shaped.pushed,
2337 projects: ProjectQuery::default(),
2338 local_tasks: shaped.local,
2339 local_projects: LocalProjects::default(),
2340 outcomes: shaped.outcomes,
2341 }
2342 }
2343 }
2344}
2345
2346fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
2353 if compensating {
2354 ceiling
2355 } else {
2356 budget.min(ceiling)
2357 }
2358}
2359
2360fn ceiling(source: &ResolvedSource) -> u32 {
2362 source.source().capabilities().max_page_size.max(1)
2363}
2364
2365async fn fetch_tasks(
2367 source: &ResolvedSource,
2368 shape: &TaskShape,
2369 start: &Resume,
2370 budget: u32,
2371 calls: &AtomicU32,
2372) -> Result<Fetched<Task>, SourceError> {
2373 let compensating = shape.local != LocalTasks::default();
2374 walk(
2375 start,
2376 budget,
2377 page_size(compensating, budget, ceiling(source)),
2378 |task| shape.local.keeps(task),
2379 |cursor, limit| async move {
2380 calls.fetch_add(1, Ordering::Relaxed);
2381 let request = PageRequest { cursor, limit };
2382 let page = source.source().query_tasks(&shape.pushed, &request).await?;
2383 match shape.local.commented_since {
2384 None => Ok(page),
2385 Some(since) => comment::commented_since(source, &shape.local, page, since).await,
2386 }
2387 },
2388 )
2389 .await
2390}
2391
2392async fn fetch_projects(
2394 source: &ResolvedSource,
2395 shape: &ProjectShape,
2396 start: &Resume,
2397 budget: u32,
2398 calls: &AtomicU32,
2399) -> Result<Fetched<Project>, SourceError> {
2400 let compensating = shape.local != LocalProjects::default();
2401 walk(
2402 start,
2403 budget,
2404 page_size(compensating, budget, ceiling(source)),
2405 |project| shape.local.keeps(project),
2406 |cursor, limit| async move {
2407 calls.fetch_add(1, Ordering::Relaxed);
2408 let request = PageRequest { cursor, limit };
2409 source
2410 .source()
2411 .query_projects(&shape.pushed, &request)
2412 .await
2413 },
2414 )
2415 .await
2416}
2417
2418async fn fetch_documents(
2420 source: &ResolvedSource,
2421 shape: &DocumentShape,
2422 start: &Resume,
2423 budget: u32,
2424 calls: &AtomicU32,
2425) -> Result<Fetched<Document>, SourceError> {
2426 let compensating = shape.local != LocalDocuments::default();
2427 walk(
2428 start,
2429 budget,
2430 page_size(compensating, budget, ceiling(source)),
2431 |document| shape.local.keeps(document),
2432 |cursor, limit| async move {
2433 calls.fetch_add(1, Ordering::Relaxed);
2434 let request = PageRequest { cursor, limit };
2435 source
2436 .source()
2437 .query_documents(&shape.pushed, &request)
2438 .await
2439 },
2440 )
2441 .await
2442}
2443
2444async fn fetch_labels(
2446 source: &ResolvedSource,
2447 start: &Resume,
2448 budget: u32,
2449 calls: &AtomicU32,
2450) -> Result<Fetched<Label>, SourceError> {
2451 walk(
2452 start,
2453 budget,
2454 page_size(false, budget, ceiling(source)),
2455 |_| true,
2456 |cursor, limit| async move {
2457 calls.fetch_add(1, Ordering::Relaxed);
2458 let request = PageRequest { cursor, limit };
2459 source.source().labels(&request).await
2460 },
2461 )
2462 .await
2463}
2464
2465async fn fetch_hits(
2467 source: &ResolvedSource,
2468 shape: &HitShape,
2469 start: &Resume,
2470 budget: u32,
2471 calls: &AtomicU32,
2472) -> Result<Fetched<Found>, SourceError> {
2473 let ceiling = ceiling(source);
2474 match shape.stream {
2475 StreamKind::Projects => {
2476 let compensating = shape.local_projects != LocalProjects::default();
2477 walk(
2478 start,
2479 budget,
2480 page_size(compensating, budget, ceiling),
2481 |found| match found {
2482 Found::Project(project) => shape.local_projects.keeps(project),
2483 Found::Task(_) => true,
2484 },
2485 |cursor, limit| async move {
2486 calls.fetch_add(1, Ordering::Relaxed);
2487 let request = PageRequest { cursor, limit };
2488 let page = source
2489 .source()
2490 .query_projects(&shape.projects, &request)
2491 .await?;
2492 Ok(Page {
2493 items: page.items.into_iter().map(Found::Project).collect(),
2494 next: page.next,
2495 })
2496 },
2497 )
2498 .await
2499 }
2500 StreamKind::Items | StreamKind::Tasks => {
2501 let compensating = shape.local_tasks != LocalTasks::default();
2502 walk(
2503 start,
2504 budget,
2505 page_size(compensating, budget, ceiling),
2506 |found| match found {
2507 Found::Task(task) => shape.local_tasks.keeps(task),
2508 Found::Project(_) => true,
2509 },
2510 |cursor, limit| async move {
2511 calls.fetch_add(1, Ordering::Relaxed);
2512 let request = PageRequest { cursor, limit };
2513 let page = source.source().query_tasks(&shape.tasks, &request).await?;
2514 Ok(Page {
2515 items: page.items.into_iter().map(Found::Task).collect(),
2516 next: page.next,
2517 })
2518 },
2519 )
2520 .await
2521 }
2522 }
2523}
2524
2525async fn forward_edges(
2527 source: &ResolvedSource,
2528 entity: Entity,
2529 id: &NativeId,
2530 request: &PageRequest,
2531) -> Result<Page<DependencyEdge>, SourceError> {
2532 match entity {
2533 Entity::Task => {
2534 source
2535 .source()
2536 .task_dependencies(id, Direction::DependsOn, request)
2537 .await
2538 }
2539 Entity::Project => {
2540 source
2541 .source()
2542 .project_dependencies(id, Direction::DependsOn, request)
2543 .await
2544 }
2545 }
2546}
2547
2548#[expect(
2557 clippy::too_many_arguments,
2558 reason = "every argument is one axis of one walk — the source, the item, the \
2559 direction, which of its two graphs, whether the reverse is emulated, where \
2560 to resume, how many rows to return and where to count calls. Grouping them \
2561 into a struct would name the same eight values one indirection further from \
2562 the loop that reads them."
2563)]
2564async fn fetch_edges(
2565 source: &ResolvedSource,
2566 native: &NativeId,
2567 direction: Direction,
2568 entity: Entity,
2569 emulating: bool,
2570 start: &Resume,
2571 budget: u32,
2572 calls: &AtomicU32,
2573) -> Result<Fetched<DependencyEdge>, SourceError> {
2574 let ceiling = ceiling(source);
2575 if !emulating {
2576 return walk(
2577 start,
2578 budget,
2579 page_size(false, budget, ceiling),
2580 |_| true,
2581 |cursor, limit| async move {
2582 calls.fetch_add(1, Ordering::Relaxed);
2583 let request = PageRequest { cursor, limit };
2584 match entity {
2585 Entity::Task => {
2586 source
2587 .source()
2588 .task_dependencies(native, direction, &request)
2589 .await
2590 }
2591 Entity::Project => {
2592 source
2593 .source()
2594 .project_dependencies(native, direction, &request)
2595 .await
2596 }
2597 }
2598 },
2599 )
2600 .await;
2601 }
2602
2603 walk(
2604 start,
2605 budget,
2606 ceiling,
2607 |_| true,
2608 |cursor, limit| async move {
2609 calls.fetch_add(1, Ordering::Relaxed);
2610 let request = PageRequest { cursor, limit };
2611 let (ids, next) = match entity {
2612 Entity::Task => {
2613 let page = source
2614 .source()
2615 .query_tasks(&TaskQuery::default(), &request)
2616 .await?;
2617 let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
2618 (ids, page.next)
2619 }
2620 Entity::Project => {
2621 let page = source
2622 .source()
2623 .query_projects(&ProjectQuery::default(), &request)
2624 .await?;
2625 let ids: Vec<NativeId> =
2626 page.items.into_iter().map(|project| project.id).collect();
2627 (ids, page.next)
2628 }
2629 };
2630
2631 let mut edges = Vec::new();
2632 for id in ids {
2633 let mut inner: Option<Cursor> = None;
2634 loop {
2635 calls.fetch_add(1, Ordering::Relaxed);
2636 let request = PageRequest {
2637 cursor: inner.clone(),
2638 limit,
2639 };
2640 let page = forward_edges(source, entity, &id, &request).await?;
2641 fits(page.items.len(), limit)?;
2645 edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
2646 unrepeated(
2647 page.next.as_ref(),
2648 inner.as_ref(),
2649 "its forward edges were being scanned",
2650 )?;
2651 match page.next {
2652 Some(cursor) => inner = Some(cursor),
2653 None => break,
2654 }
2655 }
2656 }
2657
2658 Ok(Page { items: edges, next })
2659 },
2660 )
2661 .await
2662}