mod copy;
mod fetch;
mod join;
mod local;
mod resume;
use std::collections::BTreeMap;
use std::num::NonZeroU32;
use std::sync::atomic::{AtomicU32, Ordering};
use onetaskgraph_plugin_api::{
Capabilities, Cursor, DependencyEdge, Direction, Document, DocumentQuery, Label, LabelFilter,
NativeId, Page, PageRequest, Project, ProjectFilter, ProjectQuery, SecretResolver, SourceError,
SourceName, StatusCategory, Task, TaskQuery, TextFields, TextQuery,
};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use crate::GlobalId;
use crate::config::Config;
use crate::plan::{PageToken, Predicate, QueryPlan, QueryResponse, SourceFailure, SourcePlan};
use crate::resolve::{ResolvedSource, UnavailableSource, resolve_available};
use fetch::{Fetched, Stream, fits, merge, unrepeated, walk};
use join::join_all;
use local::{LocalDocuments, LocalProjects, LocalTasks};
pub(crate) use resume::{Owed, Resumption, StreamState};
use resume::{Resume, StreamKind};
pub use copy::{CopyAction, CopyItems, CopyOutcome, CopyReport, CopyRequest, CopyScope, MatchBy};
pub use local::ProjectSelector;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct Qualified<T> {
pub id: GlobalId,
pub item: T,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct QualifiedEdge {
pub from: QualifiedEndpoint,
pub to: QualifiedEndpoint,
pub kind: onetaskgraph_plugin_api::DependencyKind,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct QualifiedEndpoint {
pub id: GlobalId,
pub kind: onetaskgraph_plugin_api::ItemKind,
}
impl std::fmt::Display for QualifiedEndpoint {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.id.fmt(formatter)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(tag = "kind", rename_all = "kebab-case")]
pub enum SearchHit {
Task(Qualified<Task>),
Project(Qualified<Project>),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "kebab-case")]
pub enum SearchKind {
Tasks,
Projects,
#[default]
Both,
}
#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
pub struct SourceListing {
pub source: SourceName,
pub kind: String,
#[serde(flatten)]
pub state: SourceState,
}
#[derive(Debug, Clone, PartialEq, Serialize, JsonSchema)]
#[serde(tag = "state", rename_all = "kebab-case")]
pub enum SourceState {
Available {
capabilities: Capabilities,
},
Unavailable {
error: SourceError,
},
}
#[derive(Debug, Clone, PartialEq)]
pub struct Paging {
pub limit: NonZeroU32,
pub token: Option<PageToken>,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct Filters {
pub text: Option<TextQuery>,
pub labels: LabelFilter,
pub statuses: Vec<StatusCategory>,
}
#[derive(Debug, Clone)]
pub struct TaskRequest {
pub sources: Vec<SourceName>,
pub filters: Filters,
pub project: ProjectSelector,
pub paging: Paging,
}
#[derive(Debug, Clone)]
pub struct ProjectRequest {
pub sources: Vec<SourceName>,
pub filters: Filters,
pub paging: Paging,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct DocumentFilters {
pub text: Option<TextQuery>,
pub labels: LabelFilter,
}
#[derive(Debug, Clone)]
pub struct DocumentRequest {
pub sources: Vec<SourceName>,
pub filters: DocumentFilters,
pub project: ProjectSelector,
pub paging: Paging,
}
#[derive(Debug, Clone)]
pub struct LabelRequest {
pub sources: Vec<SourceName>,
pub paging: Paging,
}
#[derive(Debug, Clone)]
pub struct SearchRequest {
pub sources: Vec<SourceName>,
pub text: TextQuery,
pub kind: SearchKind,
pub paging: Paging,
}
#[derive(Debug, Clone)]
pub struct DependencyRequest {
pub id: GlobalId,
pub direction: Direction,
pub paging: Paging,
}
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
pub enum EngineError {
#[error(
"no source named {name:?} is configured\n\
next: name one of the configured sources ({configured}), or add {name:?} under \
`sources` — `onetaskgraph sources list` shows what this configuration has."
)]
UnknownSource {
name: String,
configured: String,
},
#[error(
"{message}\n\
next: page with a token exactly as the previous page reported it, and against \
the same configuration — or drop `--page` to start the walk again."
)]
Token {
message: String,
},
#[error(
"no sources are configured\n\
next: add one under `sources` in onetaskgraph.yaml — `onetaskgraph schema` \
prints what each plugin accepts."
)]
NoSources,
#[error(
"source {name} cannot be written: its plugin is {kind}, which has no write \
side\n\
next: copy into a source whose plugin can be written — `onetaskgraph sources \
list` reports each one's plugin."
)]
NotWritable {
name: String,
kind: String,
},
#[error(
"source {name} has no documents: its plugin is {kind}, which holds none\n\
next: name a source whose plugin has documents — `onetaskgraph sources list` \
reports each one's plugin and what it declares."
)]
NoDocuments {
name: String,
kind: String,
},
#[error(
"the destination source {name} could not be built: {error}\n\
next: fix that source — `onetaskgraph sources list` reports its state — then \
copy again."
)]
DestinationUnavailable {
name: String,
error: SourceError,
},
#[error(
"no item with the id {id}\n\
next: check the id, or list what is there — `onetaskgraph task list` and \
`onetaskgraph project list` report what the configured sources hold."
)]
NoSuchItem {
id: String,
},
#[error(
"{item} was copied from {origin}, which that destination no longer holds\n\
next: re-run with --recreate to create a new item there instead, or restore \
{origin}."
)]
StaleOrigin {
item: String,
origin: String,
},
#[error(
"source {name} could not do it: {error}\n\
next: fix what the source named above, then copy again."
)]
SourceRefused {
name: String,
error: SourceError,
},
#[error(
"the copy failed and could not be undone.\n\
it failed because: {error}\n\
it could not be undone because: {refusal}\n\
so the destination still holds: {left_behind}\n\
next: remove those items at the destination, then copy again."
)]
CopyNotUndone {
error: Box<EngineError>,
left_behind: LeftBehind,
refusal: SourceError,
},
}
#[derive(Debug, Clone, PartialEq)]
pub struct LeftBehind {
first: GlobalId,
rest: Vec<GlobalId>,
}
impl LeftBehind {
#[must_use]
pub fn new(first: GlobalId) -> Self {
Self {
first,
rest: Vec::new(),
}
}
pub fn push(&mut self, id: GlobalId) {
self.rest.push(id);
}
pub fn iter(&self) -> impl Iterator<Item = &GlobalId> {
std::iter::once(&self.first).chain(self.rest.iter())
}
}
impl std::fmt::Display for LeftBehind {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(formatter, "{}", self.first)?;
for id in &self.rest {
write!(formatter, ", {id}")?;
}
Ok(())
}
}
pub enum ConfiguredSource {
Ready(ResolvedSource),
Unavailable(UnavailableSource),
}
impl ConfiguredSource {
#[must_use]
pub fn name(&self) -> &SourceName {
match self {
Self::Ready(source) => source.name(),
Self::Unavailable(source) => source.name(),
}
}
}
pub struct Engine {
sources: Vec<ConfiguredSource>,
selection: Vec<SourceName>,
}
impl Engine {
#[must_use]
pub fn build(config: &Config, secrets: &dyn SecretResolver) -> Self {
let (ready, unavailable) = resolve_available(config, secrets);
Self::new(
ready
.into_iter()
.map(ConfiguredSource::Ready)
.chain(unavailable.into_iter().map(ConfiguredSource::Unavailable))
.collect(),
config.selected_sources(),
)
}
#[must_use]
pub fn new(sources: Vec<ConfiguredSource>, selection: Vec<SourceName>) -> Self {
Self { sources, selection }
}
fn ready(&self) -> impl Iterator<Item = &ResolvedSource> {
self.sources.iter().filter_map(|source| match source {
ConfiguredSource::Ready(ready) => Some(ready),
ConfiguredSource::Unavailable(_) => None,
})
}
fn unavailable(&self) -> impl Iterator<Item = &UnavailableSource> {
self.sources.iter().filter_map(|source| match source {
ConfiguredSource::Unavailable(unavailable) => Some(unavailable),
ConfiguredSource::Ready(_) => None,
})
}
#[must_use]
pub fn listing(&self) -> Vec<SourceListing> {
let mut listings: Vec<SourceListing> = self
.ready()
.map(|source| SourceListing {
source: source.name().clone(),
kind: source.kind().to_owned(),
state: SourceState::Available {
capabilities: source.source().capabilities(),
},
})
.chain(self.unavailable().map(|source| SourceListing {
source: source.name().clone(),
kind: source.kind().to_owned(),
state: SourceState::Unavailable {
error: source.error().clone(),
},
}))
.collect();
listings.sort_by(|left, right| left.source.cmp(&right.source));
listings
}
#[must_use]
pub fn has(&self, name: &SourceName) -> bool {
self.sources.iter().any(|source| source.name() == name)
}
pub async fn tasks(
&self,
request: &TaskRequest,
) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
let mut names = self.resolve_selection(&request.sources)?;
if let ProjectSelector::Qualified(id) = &request.project {
self.known(&id.source)?;
names.retain(|name| name == &id.source);
}
let query = shape("task-list", &names, &(&request.filters, &request.project));
let states = resumption(
self,
request.paging.token.as_ref(),
&[StreamKind::Items],
&query,
)?;
let budget = request.paging.limit.get();
let mut answer = Answer::new();
let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
let shapes: Vec<TaskShape> = ready
.iter()
.map(|source| {
shape_tasks(
&source.source().capabilities(),
&request.filters,
&project_filter(&request.project),
)
})
.collect();
let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
let walks = ready
.iter()
.enumerate()
.map(|(index, source)| {
fetch_tasks(
source,
&shapes[index],
&starts[index],
budget,
&counters[index],
)
})
.collect();
let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
answer.finish(
streams,
budget,
owed(&states),
&query,
|name, task: Task| Qualified {
id: GlobalId::new(name.clone(), task.id.clone()),
item: task,
},
)
}
pub async fn projects(
&self,
request: &ProjectRequest,
) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
let names = self.resolve_selection(&request.sources)?;
let query = shape("project-list", &names, &request.filters);
let states = resumption(
self,
request.paging.token.as_ref(),
&[StreamKind::Items],
&query,
)?;
let budget = request.paging.limit.get();
let mut answer = Answer::new();
let mut with_projects = Vec::new();
for source in answer.split(self, &names) {
if source.source().capabilities().projects.is_native() {
with_projects.push(source);
} else {
answer.unreachable_predicate(source, Predicate::Project);
}
}
let (ready, starts) = walking(with_projects, &states, StreamKind::Items);
let shapes: Vec<ProjectShape> = ready
.iter()
.map(|source| shape_projects(&source.source().capabilities(), &request.filters))
.collect();
let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
let walks = ready
.iter()
.enumerate()
.map(|(index, source)| {
fetch_projects(
source,
&shapes[index],
&starts[index],
budget,
&counters[index],
)
})
.collect();
let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
answer.finish(
streams,
budget,
owed(&states),
&query,
|name, project: Project| Qualified {
id: GlobalId::new(name.clone(), project.id.clone()),
item: project,
},
)
}
pub async fn documents(
&self,
request: &DocumentRequest,
) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
let mut names = self.resolve_selection(&request.sources)?;
if let ProjectSelector::Qualified(id) = &request.project {
self.known(&id.source)?;
names.retain(|name| name == &id.source);
}
let query = shape(
"document-list",
&names,
&(&request.filters, &request.project),
);
let states = resumption(
self,
request.paging.token.as_ref(),
&[StreamKind::Items],
&query,
)?;
let budget = request.paging.limit.get();
let mut answer = Answer::new();
let mut with_documents = Vec::new();
for source in answer.split(self, &names) {
if source.source().capabilities().documents.is_native() {
with_documents.push(source);
} else {
answer.unreachable_predicate(source, Predicate::Document);
}
}
let (ready, starts) = walking(with_documents, &states, StreamKind::Items);
let shapes: Vec<DocumentShape> = ready
.iter()
.map(|source| {
shape_documents(
&source.source().capabilities(),
&request.filters,
&project_filter(&request.project),
)
})
.collect();
let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
let walks = ready
.iter()
.enumerate()
.map(|(index, source)| {
fetch_documents(
source,
&shapes[index],
&starts[index],
budget,
&counters[index],
)
})
.collect();
let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
answer.finish(
streams,
budget,
owed(&states),
&query,
|name, document: Document| Qualified {
id: GlobalId::new(name.clone(), document.id.clone()),
item: document,
},
)
}
pub async fn labels(
&self,
request: &LabelRequest,
) -> Result<QueryResponse<Qualified<Label>>, EngineError> {
let names = self.resolve_selection(&request.sources)?;
let query = shape("label-list", &names, &());
let states = resumption(
self,
request.paging.token.as_ref(),
&[StreamKind::Items],
&query,
)?;
let budget = request.paging.limit.get();
let mut answer = Answer::new();
let (ready, starts) = walking(answer.split(self, &names), &states, StreamKind::Items);
let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
let outcomes: Vec<Outcomes> = ready.iter().map(|_| Outcomes::default()).collect();
let walks = ready
.iter()
.enumerate()
.map(|(index, source)| fetch_labels(source, &starts[index], budget, &counters[index]))
.collect();
let streams = answer.collect(&ready, join_all(walks).await, &counters, outcomes);
answer.finish(
streams,
budget,
owed(&states),
&query,
|name, label: Label| Qualified {
id: GlobalId::new(name.clone(), label.id.clone()),
item: label,
},
)
}
pub async fn search(
&self,
request: &SearchRequest,
) -> Result<QueryResponse<SearchHit>, EngineError> {
let names = self.resolve_selection(&request.sources)?;
let reads: &[StreamKind] = match request.kind {
SearchKind::Tasks => &[StreamKind::Tasks],
SearchKind::Projects => &[StreamKind::Projects],
SearchKind::Both => &[StreamKind::Tasks, StreamKind::Projects],
};
let query = shape("search", &names, &request.text);
let states = resumption(self, request.paging.token.as_ref(), reads, &query)?;
let budget = request.paging.limit.get();
let filters = Filters {
text: Some(request.text.clone()),
..Filters::default()
};
let mut answer = Answer::new();
let mut ready = Vec::new();
let mut kinds = Vec::new();
let mut starts = Vec::new();
for source in answer.split(self, &names) {
let mut streams = Vec::new();
if matches!(request.kind, SearchKind::Tasks | SearchKind::Both) {
streams.push(StreamKind::Tasks);
}
if matches!(request.kind, SearchKind::Projects | SearchKind::Both) {
if source.source().capabilities().projects.is_native() {
streams.push(StreamKind::Projects);
} else {
answer.unreachable_predicate(source, Predicate::Project);
}
}
for stream in streams {
if let Some(resume) = resume_at(&states, source.name(), stream) {
ready.push(source);
kinds.push(stream);
starts.push(resume);
}
}
}
let shapes: Vec<HitShape> = ready
.iter()
.zip(kinds.iter())
.map(|(source, kind)| shape_hits(&source.source().capabilities(), &filters, *kind))
.collect();
let counters: Vec<AtomicU32> = ready.iter().map(|_| AtomicU32::new(0)).collect();
let outcomes: Vec<Outcomes> = shapes.iter().map(|shape| shape.outcomes.clone()).collect();
let walks = ready
.iter()
.enumerate()
.map(|(index, source)| {
fetch_hits(
source,
&shapes[index],
&starts[index],
budget,
&counters[index],
)
})
.collect();
let streams =
answer.collect_streams(&ready, &kinds, join_all(walks).await, &counters, outcomes);
answer.finish(
streams,
budget,
owed(&states),
&query,
|name, found: Found| match found {
Found::Task(task) => SearchHit::Task(Qualified {
id: GlobalId::new(name.clone(), task.id.clone()),
item: task,
}),
Found::Project(project) => SearchHit::Project(Qualified {
id: GlobalId::new(name.clone(), project.id.clone()),
item: project,
}),
},
)
}
pub async fn task(&self, id: &GlobalId) -> Result<QueryResponse<Qualified<Task>>, EngineError> {
let name = self.known(&id.source)?;
let mut answer = Answer::new();
let selected = answer.split(self, std::slice::from_ref(&name));
let Some(source) = selected.first() else {
return answer.nothing();
};
let found = source.source().get_task(&id.native).await;
let qualified = GlobalId::new(source.name().clone(), id.native.clone());
answer.one(source, found, |task| Qualified {
id: qualified,
item: task,
})
}
pub async fn project(
&self,
id: &GlobalId,
) -> Result<QueryResponse<Qualified<Project>>, EngineError> {
let name = self.known(&id.source)?;
let mut answer = Answer::new();
let selected = answer.split(self, std::slice::from_ref(&name));
let Some(source) = selected.first() else {
return answer.nothing();
};
let found = source.source().get_project(&id.native).await;
let qualified = GlobalId::new(source.name().clone(), id.native.clone());
answer.one(source, found, |project| Qualified {
id: qualified,
item: project,
})
}
pub async fn document(
&self,
id: &GlobalId,
) -> Result<QueryResponse<Qualified<Document>>, EngineError> {
let name = self.known(&id.source)?;
let mut answer = Answer::new();
let selected = answer.split(self, std::slice::from_ref(&name));
let Some(source) = selected.first() else {
return answer.nothing();
};
if !source.source().capabilities().documents.is_native() {
answer.unreachable_predicate(source, Predicate::Document);
return answer.nothing();
}
let found = source.source().get_document(&id.native).await;
let qualified = GlobalId::new(source.name().clone(), id.native.clone());
answer.one(source, found, |document| Qualified {
id: qualified,
item: document,
})
}
pub async fn task_dependencies(
&self,
request: &DependencyRequest,
) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
self.dependencies(request, Entity::Task).await
}
pub async fn project_dependencies(
&self,
request: &DependencyRequest,
) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
self.dependencies(request, Entity::Project).await
}
async fn dependencies(
&self,
request: &DependencyRequest,
entity: Entity,
) -> Result<QueryResponse<QualifiedEdge>, EngineError> {
let name = self.known(&request.id.source)?;
let query = shape(
"dependencies",
std::slice::from_ref(&name),
&(entity, &request.id.native, request.direction),
);
let states = resumption(
self,
request.paging.token.as_ref(),
&[StreamKind::Items],
&query,
)?;
let budget = request.paging.limit.get();
let mut answer = Answer::new();
let (ready, starts) = walking(
answer.split(self, std::slice::from_ref(&name)),
&states,
StreamKind::Items,
);
let Some(source) = ready.first() else {
return answer.nothing();
};
let capabilities = source.source().capabilities();
let support = match entity {
Entity::Task => capabilities.task_dependencies,
Entity::Project => capabilities.project_dependencies,
};
let emulating = request.direction == Direction::DependedOnBy && !support.answers_reverse();
let mut outcomes = Outcomes::default();
if request.direction == Direction::DependedOnBy {
if emulating {
outcomes.record(Predicate::ReverseDependencies, Outcome::Emulated);
} else {
outcomes.record(Predicate::ReverseDependencies, Outcome::PushedDown);
}
}
let counters = vec![AtomicU32::new(0)];
let walked = fetch_edges(
source,
&request.id.native,
request.direction,
entity,
emulating,
&starts[0],
budget,
&counters[0],
)
.await;
let streams = answer.collect(&ready, vec![walked], &counters, vec![outcomes]);
answer.finish(
streams,
budget,
owed(&states),
&query,
|name, edge: DependencyEdge| QualifiedEdge {
from: qualify_endpoint(name, edge.from),
to: qualify_endpoint(name, edge.to),
kind: edge.kind,
},
)
}
fn resolve_selection(&self, asked: &[SourceName]) -> Result<Vec<SourceName>, EngineError> {
if asked.is_empty() {
if self.selection.is_empty() {
return Err(EngineError::NoSources);
}
return Ok(self.selection.clone());
}
asked.iter().map(|name| self.known(name)).collect()
}
fn known(&self, name: &SourceName) -> Result<SourceName, EngineError> {
if self.has(name) {
return Ok(name.clone());
}
if self.sources.is_empty() {
return Err(EngineError::NoSources);
}
Err(EngineError::UnknownSource {
name: name.to_string(),
configured: self
.listing()
.iter()
.map(|listing| listing.source.to_string())
.collect::<Vec<_>>()
.join(", "),
})
}
}
fn qualify_endpoint(
source: &SourceName,
endpoint: onetaskgraph_plugin_api::DependencyEndpoint,
) -> QualifiedEndpoint {
let kind = endpoint.kind;
let is_qualified = endpoint.is_qualified();
let endpoint_id = endpoint.into_id();
QualifiedEndpoint {
id: if is_qualified {
endpoint_id
.parse()
.expect("plugin-api validates qualified dependency endpoints")
} else {
GlobalId::new(source.clone(), NativeId(endpoint_id))
},
kind,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Entity {
Task,
Project,
}
enum Found {
Task(Task),
Project(Project),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Outcome {
PushedDown,
AppliedLocally,
Emulated,
Unavailable,
}
#[derive(Debug, Clone, Default, PartialEq)]
struct Outcomes(BTreeMap<Predicate, Outcome>);
impl Outcomes {
fn record(&mut self, predicate: Predicate, outcome: Outcome) {
self.0.insert(predicate, outcome);
}
fn record_all(&mut self, predicates: impl IntoIterator<Item = Predicate>, outcome: Outcome) {
for predicate in predicates {
self.record(predicate, outcome);
}
}
fn with(&self, outcome: Outcome) -> Vec<Predicate> {
self.0
.iter()
.filter(|(_, recorded)| **recorded == outcome)
.map(|(predicate, _)| *predicate)
.collect()
}
}
struct TaskShape {
pushed: TaskQuery,
local: LocalTasks,
outcomes: Outcomes,
}
struct ProjectShape {
pushed: ProjectQuery,
local: LocalProjects,
outcomes: Outcomes,
}
struct DocumentShape {
pushed: DocumentQuery,
local: LocalDocuments,
outcomes: Outcomes,
}
struct HitShape {
stream: StreamKind,
tasks: TaskQuery,
projects: ProjectQuery,
local_tasks: LocalTasks,
local_projects: LocalProjects,
outcomes: Outcomes,
}
struct Answer {
plans: Vec<SourcePlan>,
errors: Vec<SourceFailure>,
}
impl Answer {
fn new() -> Self {
Self {
plans: Vec::new(),
errors: Vec::new(),
}
}
fn split<'a>(&mut self, engine: &'a Engine, names: &[SourceName]) -> Vec<&'a ResolvedSource> {
let mut selected = Vec::new();
for name in names {
match engine.sources.iter().find(|source| source.name() == name) {
Some(ConfiguredSource::Ready(source)) => selected.push(source),
Some(ConfiguredSource::Unavailable(source)) => {
self.errors.push(source.failure());
}
None => {}
}
}
selected
}
fn unreachable_predicate(&mut self, source: &ResolvedSource, predicate: Predicate) {
let mut outcomes = Outcomes::default();
outcomes.record(predicate, Outcome::Unavailable);
self.plans.push(plan_for(source, outcomes, 0));
}
fn collect<T>(
&mut self,
ready: &[&ResolvedSource],
walked: Vec<Result<Fetched<T>, SourceError>>,
counters: &[AtomicU32],
outcomes: Vec<Outcomes>,
) -> Vec<Stream<T>> {
let kinds = vec![StreamKind::Items; ready.len()];
self.collect_streams(ready, &kinds, walked, counters, outcomes)
}
fn collect_streams<T>(
&mut self,
ready: &[&ResolvedSource],
kinds: &[StreamKind],
walked: Vec<Result<Fetched<T>, SourceError>>,
counters: &[AtomicU32],
outcomes: Vec<Outcomes>,
) -> Vec<Stream<T>> {
let mut streams = Vec::new();
for (index, result) in walked.into_iter().enumerate() {
let source = ready[index];
let pages = counters[index].load(Ordering::Relaxed);
self.plans
.push(plan_for(source, outcomes[index].clone(), pages));
match result {
Ok(fetched) => streams.push(Stream {
source: source.name().clone(),
kind: kinds[index],
fetched,
}),
Err(error) => self.errors.push(SourceFailure {
source: source.name().clone(),
error,
}),
}
}
streams
}
fn one<T, U>(
mut self,
source: &ResolvedSource,
found: Result<Option<T>, SourceError>,
qualify: impl FnOnce(T) -> U,
) -> Result<QueryResponse<U>, EngineError> {
self.plans.push(plan_for(source, Outcomes::default(), 1));
let items = match found {
Ok(Some(item)) => vec![qualify(item)],
Ok(None) => Vec::new(),
Err(error) => {
self.errors.push(SourceFailure {
source: source.name().clone(),
error,
});
Vec::new()
}
};
Ok(QueryResponse {
items,
next: None,
plan: QueryPlan {
per_source: merge_plans(self.plans),
},
errors: self.errors,
})
}
fn nothing<U>(self) -> Result<QueryResponse<U>, EngineError> {
Ok(QueryResponse {
items: Vec::new(),
next: None,
plan: QueryPlan {
per_source: merge_plans(self.plans),
},
errors: self.errors,
})
}
fn finish<T, U>(
self,
streams: Vec<Stream<T>>,
budget: u32,
first: Option<&Owed>,
query: &str,
qualify: impl Fn(&SourceName, T) -> U,
) -> Result<QueryResponse<U>, EngineError> {
let (rows, states, owed) = merge(streams, budget, first);
let next = (!states.is_empty()).then(|| PageToken::encode(query, owed, &states));
Ok(QueryResponse {
items: rows
.into_iter()
.map(|(name, item)| qualify(&name, item))
.collect(),
next,
plan: QueryPlan {
per_source: merge_plans(self.plans),
},
errors: self.errors,
})
}
}
fn plan_for(source: &ResolvedSource, outcomes: Outcomes, pages: u32) -> SourcePlan {
SourcePlan {
source: source.name().clone(),
kind: source.kind().to_owned(),
pushed_down: outcomes.with(Outcome::PushedDown),
applied_locally: outcomes.with(Outcome::AppliedLocally),
emulated: outcomes.with(Outcome::Emulated),
unavailable: outcomes.with(Outcome::Unavailable),
pages_fetched: pages,
}
}
fn merge_plans(plans: Vec<SourcePlan>) -> Vec<SourcePlan> {
let mut merged: Vec<SourcePlan> = Vec::new();
for plan in plans {
if let Some(existing) = merged
.iter_mut()
.find(|existing| existing.source == plan.source)
{
existing.pushed_down.extend(plan.pushed_down);
existing.applied_locally.extend(plan.applied_locally);
existing.emulated.extend(plan.emulated);
existing.unavailable.extend(plan.unavailable);
existing.pages_fetched = existing.pages_fetched.saturating_add(plan.pages_fetched);
for list in [
&mut existing.pushed_down,
&mut existing.applied_locally,
&mut existing.emulated,
&mut existing.unavailable,
] {
list.sort_unstable();
list.dedup();
}
} else {
merged.push(plan);
}
}
merged
}
fn walking<'a>(
selected: Vec<&'a ResolvedSource>,
states: &Option<Resumption>,
kind: StreamKind,
) -> (Vec<&'a ResolvedSource>, Vec<Resume>) {
let mut ready = Vec::new();
let mut starts = Vec::new();
for source in selected {
if let Some(resume) = resume_at(states, source.name(), kind) {
ready.push(source);
starts.push(resume);
}
}
(ready, starts)
}
fn shape(verb: &str, sources: &[SourceName], filters: &impl std::fmt::Debug) -> String {
let names: Vec<&str> = sources.iter().map(SourceName::as_str).collect();
fingerprint(&format!("{verb}|{names:?}|{filters:?}"))
}
fn fingerprint(text: &str) -> String {
let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
for byte in text.as_bytes() {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
}
format!("{hash:016x}")
}
fn owed(document: &Option<Resumption>) -> Option<&Owed> {
document.as_ref()?.owed.as_ref()
}
fn resume_at(states: &Option<Resumption>, source: &SourceName, kind: StreamKind) -> Option<Resume> {
match states {
None => Some(Resume::default()),
Some(document) => document
.streams
.iter()
.find(|state| &state.source == source && state.stream == kind)
.map(|state| state.resume.clone()),
}
}
fn resumption(
engine: &Engine,
token: Option<&PageToken>,
reads: &[StreamKind],
query: &str,
) -> Result<Option<Resumption>, EngineError> {
let Some(document) = token.map(PageToken::decode) else {
return Ok(None);
};
if document.query != query {
return Err(EngineError::Token {
message: "this page token was written by a different query — resume the walk it \
came from, or drop --page to start this one from the beginning"
.to_owned(),
});
}
let states = &document.streams;
let mut seen: Vec<(&SourceName, StreamKind)> = Vec::new();
for state in states {
if !reads.contains(&state.stream) {
return Err(EngineError::Token {
message: format!(
"this page token resumes {}, which this command does not read — it \
was written by a different query",
state.stream.describe()
),
});
}
let ceiling = engine
.ready()
.find(|source| source.name() == &state.source)
.map(ceiling);
if ceiling.is_none() && !engine.has(&state.source) {
return Err(EngineError::Token {
message: format!(
"this page token resumes a source called {:?}, which this \
configuration does not have",
state.source.as_str()
),
});
}
if let Some(ceiling) = ceiling
&& state.resume.skip >= ceiling
{
return Err(EngineError::Token {
message: format!(
"this page token resumes {} rows into a page of source {:?}, which \
serves at most {ceiling}",
state.resume.skip,
state.source.as_str()
),
});
}
if seen.contains(&(&state.source, state.stream)) {
return Err(EngineError::Token {
message: format!(
"this page token gives source {:?} two places to resume from",
state.source.as_str()
),
});
}
seen.push((&state.source, state.stream));
}
if let Some(owed) = &document.owed
&& !document
.streams
.iter()
.any(|state| state.source == owed.source && state.stream == owed.stream)
{
return Err(EngineError::Token {
message: format!(
"this page token owes the next row to a stream it does not resume, \
{:?}'s {}",
owed.source.as_str(),
owed.stream.describe()
),
});
}
Ok(Some(document))
}
fn project_filter(selector: &ProjectSelector) -> ProjectFilter {
match selector {
ProjectSelector::Any => ProjectFilter::Any,
ProjectSelector::Orphans => ProjectFilter::Orphans,
ProjectSelector::Native(id) => ProjectFilter::Is(id.clone()),
ProjectSelector::Qualified(id) => ProjectFilter::Is(id.native.clone()),
}
}
fn text_predicates(fields: TextFields) -> Vec<Predicate> {
match fields {
TextFields::Title => vec![Predicate::SearchTitle],
TextFields::Content => vec![Predicate::SearchContent],
TextFields::TitleOrContent => vec![Predicate::SearchTitle, Predicate::SearchContent],
}
}
fn searches_natively(capabilities: &Capabilities, fields: TextFields) -> bool {
match fields {
TextFields::Title => capabilities.search_title.is_native(),
TextFields::Content => capabilities.search_content.is_native(),
TextFields::TitleOrContent => {
capabilities.search_title.is_native() && capabilities.search_content.is_native()
}
}
}
fn shape_tasks(
capabilities: &Capabilities,
filters: &Filters,
project: &ProjectFilter,
) -> TaskShape {
let mut pushed = TaskQuery::default();
let mut local = LocalTasks::default();
let mut outcomes = Outcomes::default();
if !filters.labels.is_empty() {
if capabilities.filter_by_label.is_native() {
pushed.labels = filters.labels.clone();
outcomes.record(Predicate::Label, Outcome::PushedDown);
} else {
local.labels = Some(filters.labels.clone());
outcomes.record(Predicate::Label, Outcome::AppliedLocally);
}
}
if !filters.statuses.is_empty() {
if capabilities.filter_by_status.is_native() {
pushed.statuses.clone_from(&filters.statuses);
outcomes.record(Predicate::Status, Outcome::PushedDown);
} else {
local.statuses.clone_from(&filters.statuses);
outcomes.record(Predicate::Status, Outcome::AppliedLocally);
}
}
if let Some(text) = &filters.text {
let predicates = text_predicates(text.fields);
if searches_natively(capabilities, text.fields) {
pushed.text = Some(text.clone());
outcomes.record_all(predicates, Outcome::PushedDown);
} else {
local.text = Some(text.clone());
outcomes.record_all(predicates, Outcome::AppliedLocally);
}
}
match project {
ProjectFilter::Any => {}
ProjectFilter::Orphans => {
if capabilities.orphan_tasks.is_native() {
pushed.project = ProjectFilter::Orphans;
outcomes.record(Predicate::Project, Outcome::PushedDown);
} else {
local.project = Some(ProjectFilter::Orphans);
outcomes.record(Predicate::Project, Outcome::AppliedLocally);
}
}
ProjectFilter::Is(id) => {
if capabilities.projects.is_native() {
pushed.project = ProjectFilter::Is(id.clone());
outcomes.record(Predicate::Project, Outcome::PushedDown);
} else {
local.project = Some(ProjectFilter::Is(id.clone()));
outcomes.record(Predicate::Project, Outcome::AppliedLocally);
}
}
}
TaskShape {
pushed,
local,
outcomes,
}
}
fn shape_projects(capabilities: &Capabilities, filters: &Filters) -> ProjectShape {
let mut pushed = ProjectQuery::default();
let mut local = LocalProjects::default();
let mut outcomes = Outcomes::default();
if !filters.labels.is_empty() {
if capabilities.filter_by_label.is_native() {
pushed.labels = filters.labels.clone();
outcomes.record(Predicate::Label, Outcome::PushedDown);
} else {
local.labels = Some(filters.labels.clone());
outcomes.record(Predicate::Label, Outcome::AppliedLocally);
}
}
if !filters.statuses.is_empty() {
if capabilities.filter_by_status.is_native() {
pushed.statuses.clone_from(&filters.statuses);
outcomes.record(Predicate::Status, Outcome::PushedDown);
} else {
local.statuses.clone_from(&filters.statuses);
outcomes.record(Predicate::Status, Outcome::AppliedLocally);
}
}
if let Some(text) = &filters.text {
let predicates = text_predicates(text.fields);
if searches_natively(capabilities, text.fields) {
pushed.text = Some(text.clone());
outcomes.record_all(predicates, Outcome::PushedDown);
} else {
local.text = Some(text.clone());
outcomes.record_all(predicates, Outcome::AppliedLocally);
}
}
ProjectShape {
pushed,
local,
outcomes,
}
}
fn shape_documents(
capabilities: &Capabilities,
filters: &DocumentFilters,
project: &ProjectFilter,
) -> DocumentShape {
let mut pushed = DocumentQuery::default();
let mut local = LocalDocuments::default();
let mut outcomes = Outcomes::default();
if !filters.labels.is_empty() {
if capabilities.filter_by_label.is_native() {
pushed.labels = filters.labels.clone();
outcomes.record(Predicate::Label, Outcome::PushedDown);
} else {
local.labels = Some(filters.labels.clone());
outcomes.record(Predicate::Label, Outcome::AppliedLocally);
}
}
if let Some(text) = &filters.text {
let predicates = text_predicates(text.fields);
if searches_natively(capabilities, text.fields) {
pushed.text = Some(text.clone());
outcomes.record_all(predicates, Outcome::PushedDown);
} else {
local.text = Some(text.clone());
outcomes.record_all(predicates, Outcome::AppliedLocally);
}
}
match project {
ProjectFilter::Any => {}
ProjectFilter::Orphans => {
if capabilities.orphan_tasks.is_native() {
pushed.project = ProjectFilter::Orphans;
outcomes.record(Predicate::Project, Outcome::PushedDown);
} else {
local.project = Some(ProjectFilter::Orphans);
outcomes.record(Predicate::Project, Outcome::AppliedLocally);
}
}
ProjectFilter::Is(id) => {
if capabilities.projects.is_native() {
pushed.project = ProjectFilter::Is(id.clone());
outcomes.record(Predicate::Project, Outcome::PushedDown);
} else {
local.project = Some(ProjectFilter::Is(id.clone()));
outcomes.record(Predicate::Project, Outcome::AppliedLocally);
}
}
}
DocumentShape {
pushed,
local,
outcomes,
}
}
fn shape_hits(capabilities: &Capabilities, filters: &Filters, stream: StreamKind) -> HitShape {
match stream {
StreamKind::Projects => {
let shaped = shape_projects(capabilities, filters);
HitShape {
stream,
tasks: TaskQuery::default(),
projects: shaped.pushed,
local_tasks: LocalTasks::default(),
local_projects: shaped.local,
outcomes: shaped.outcomes,
}
}
StreamKind::Items | StreamKind::Tasks => {
let shaped = shape_tasks(capabilities, filters, &ProjectFilter::Any);
HitShape {
stream,
tasks: shaped.pushed,
projects: ProjectQuery::default(),
local_tasks: shaped.local,
local_projects: LocalProjects::default(),
outcomes: shaped.outcomes,
}
}
}
}
fn page_size(compensating: bool, budget: u32, ceiling: u32) -> u32 {
if compensating {
ceiling
} else {
budget.min(ceiling)
}
}
fn ceiling(source: &ResolvedSource) -> u32 {
source.source().capabilities().max_page_size.max(1)
}
async fn fetch_tasks(
source: &ResolvedSource,
shape: &TaskShape,
start: &Resume,
budget: u32,
calls: &AtomicU32,
) -> Result<Fetched<Task>, SourceError> {
let compensating = shape.local != LocalTasks::default();
walk(
start,
budget,
page_size(compensating, budget, ceiling(source)),
|task| shape.local.keeps(task),
|cursor, limit| async move {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest { cursor, limit };
source.source().query_tasks(&shape.pushed, &request).await
},
)
.await
}
async fn fetch_projects(
source: &ResolvedSource,
shape: &ProjectShape,
start: &Resume,
budget: u32,
calls: &AtomicU32,
) -> Result<Fetched<Project>, SourceError> {
let compensating = shape.local != LocalProjects::default();
walk(
start,
budget,
page_size(compensating, budget, ceiling(source)),
|project| shape.local.keeps(project),
|cursor, limit| async move {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest { cursor, limit };
source
.source()
.query_projects(&shape.pushed, &request)
.await
},
)
.await
}
async fn fetch_documents(
source: &ResolvedSource,
shape: &DocumentShape,
start: &Resume,
budget: u32,
calls: &AtomicU32,
) -> Result<Fetched<Document>, SourceError> {
let compensating = shape.local != LocalDocuments::default();
walk(
start,
budget,
page_size(compensating, budget, ceiling(source)),
|document| shape.local.keeps(document),
|cursor, limit| async move {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest { cursor, limit };
source
.source()
.query_documents(&shape.pushed, &request)
.await
},
)
.await
}
async fn fetch_labels(
source: &ResolvedSource,
start: &Resume,
budget: u32,
calls: &AtomicU32,
) -> Result<Fetched<Label>, SourceError> {
walk(
start,
budget,
page_size(false, budget, ceiling(source)),
|_| true,
|cursor, limit| async move {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest { cursor, limit };
source.source().labels(&request).await
},
)
.await
}
async fn fetch_hits(
source: &ResolvedSource,
shape: &HitShape,
start: &Resume,
budget: u32,
calls: &AtomicU32,
) -> Result<Fetched<Found>, SourceError> {
let ceiling = ceiling(source);
match shape.stream {
StreamKind::Projects => {
let compensating = shape.local_projects != LocalProjects::default();
walk(
start,
budget,
page_size(compensating, budget, ceiling),
|found| match found {
Found::Project(project) => shape.local_projects.keeps(project),
Found::Task(_) => true,
},
|cursor, limit| async move {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest { cursor, limit };
let page = source
.source()
.query_projects(&shape.projects, &request)
.await?;
Ok(Page {
items: page.items.into_iter().map(Found::Project).collect(),
next: page.next,
})
},
)
.await
}
StreamKind::Items | StreamKind::Tasks => {
let compensating = shape.local_tasks != LocalTasks::default();
walk(
start,
budget,
page_size(compensating, budget, ceiling),
|found| match found {
Found::Task(task) => shape.local_tasks.keeps(task),
Found::Project(_) => true,
},
|cursor, limit| async move {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest { cursor, limit };
let page = source.source().query_tasks(&shape.tasks, &request).await?;
Ok(Page {
items: page.items.into_iter().map(Found::Task).collect(),
next: page.next,
})
},
)
.await
}
}
}
async fn forward_edges(
source: &ResolvedSource,
entity: Entity,
id: &NativeId,
request: &PageRequest,
) -> Result<Page<DependencyEdge>, SourceError> {
match entity {
Entity::Task => {
source
.source()
.task_dependencies(id, Direction::DependsOn, request)
.await
}
Entity::Project => {
source
.source()
.project_dependencies(id, Direction::DependsOn, request)
.await
}
}
}
#[expect(
clippy::too_many_arguments,
reason = "every argument is one axis of one walk — the source, the item, the \
direction, which of its two graphs, whether the reverse is emulated, where \
to resume, how many rows to return and where to count calls. Grouping them \
into a struct would name the same eight values one indirection further from \
the loop that reads them."
)]
async fn fetch_edges(
source: &ResolvedSource,
native: &NativeId,
direction: Direction,
entity: Entity,
emulating: bool,
start: &Resume,
budget: u32,
calls: &AtomicU32,
) -> Result<Fetched<DependencyEdge>, SourceError> {
let ceiling = ceiling(source);
if !emulating {
return walk(
start,
budget,
page_size(false, budget, ceiling),
|_| true,
|cursor, limit| async move {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest { cursor, limit };
match entity {
Entity::Task => {
source
.source()
.task_dependencies(native, direction, &request)
.await
}
Entity::Project => {
source
.source()
.project_dependencies(native, direction, &request)
.await
}
}
},
)
.await;
}
walk(
start,
budget,
ceiling,
|_| true,
|cursor, limit| async move {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest { cursor, limit };
let (ids, next) = match entity {
Entity::Task => {
let page = source
.source()
.query_tasks(&TaskQuery::default(), &request)
.await?;
let ids: Vec<NativeId> = page.items.into_iter().map(|task| task.id).collect();
(ids, page.next)
}
Entity::Project => {
let page = source
.source()
.query_projects(&ProjectQuery::default(), &request)
.await?;
let ids: Vec<NativeId> =
page.items.into_iter().map(|project| project.id).collect();
(ids, page.next)
}
};
let mut edges = Vec::new();
for id in ids {
let mut inner: Option<Cursor> = None;
loop {
calls.fetch_add(1, Ordering::Relaxed);
let request = PageRequest {
cursor: inner.clone(),
limit,
};
let page = forward_edges(source, entity, &id, &request).await?;
fits(page.items.len(), limit)?;
edges.extend(page.items.into_iter().filter(|edge| &edge.to == native));
unrepeated(
page.next.as_ref(),
inner.as_ref(),
"its forward edges were being scanned",
)?;
match page.next {
Some(cursor) => inner = Some(cursor),
None => break,
}
}
}
Ok(Page { items: edges, next })
},
)
.await
}