use std::collections::{BTreeMap, BTreeSet};
use onetaskgraph_plugin_api::{
Cursor, DependencyEdge, DependencyEndpoint, DependencyKind, Direction, Document, DocumentQuery,
ItemKind, ItemWrite, Location, NativeId, Page, PageRequest, Project, ProjectQuery, Repository,
SourceError, SourceName, Task, TaskQuery,
};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::GlobalId;
use crate::resolve::ResolvedSource;
use super::fetch::{fits, unrepeated};
use super::local::ProjectSelector;
use super::{
DocumentFilters, DocumentRequest, Engine, EngineError, Filters, LeftBehind, Paging, Qualified,
TaskRequest,
};
#[derive(Debug, Clone)]
pub struct CopyRequest {
pub items: CopyItems,
pub scope: CopyScope,
pub destination: SourceName,
pub match_by: Option<MatchBy>,
pub recreate: bool,
pub dry_run: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CopyItems(Vec<GlobalId>);
impl CopyItems {
#[must_use]
pub fn new(items: Vec<GlobalId>) -> Option<Self> {
(!items.is_empty()).then_some(Self(items))
}
#[must_use]
pub fn as_slice(&self) -> &[GlobalId] {
&self.0
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CopyScope {
Tasks,
Projects {
tasks: bool,
},
Documents,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MatchBy {
Title,
Metadata(String),
}
impl MatchBy {
#[must_use]
pub fn parse(key: &str) -> Self {
if key == "title" {
Self::Title
} else {
Self::Metadata(key.to_owned())
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct CopyReport {
pub items: Vec<CopyOutcome>,
#[serde(default, skip_serializing_if = "nothing_to_report")]
#[schemars(!skip_serializing_if)]
pub references_rewritten: u64,
#[serde(default, skip_serializing_if = "nothing_to_report")]
#[schemars(!skip_serializing_if)]
pub references_unresolved: u64,
#[serde(default, skip_serializing_if = "nothing_to_report")]
#[schemars(!skip_serializing_if)]
pub references_ambiguous: u64,
}
fn nothing_to_report(figure: &u64) -> bool {
*figure == 0
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
struct Counted {
rewritten: u64,
unresolved: u64,
ambiguous: u64,
}
impl Counted {
fn add(&mut self, other: Self) {
self.rewritten += other.rewritten;
self.unresolved += other.unresolved;
self.ambiguous += other.ambiguous;
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct CopyOutcome {
pub source: GlobalId,
#[serde(flatten)]
pub action: CopyAction,
}
impl CopyOutcome {
#[must_use]
pub fn destination(&self) -> Option<&GlobalId> {
self.action.destination()
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(tag = "action", rename_all = "kebab-case")]
pub enum CopyAction {
Created {
destination: Option<GlobalId>,
},
Updated {
destination: GlobalId,
},
Unchanged {
destination: GlobalId,
},
Orphaned {
destination: GlobalId,
},
}
impl CopyAction {
#[must_use]
pub fn destination(&self) -> Option<&GlobalId> {
match self {
Self::Created { destination } => destination.as_ref(),
Self::Updated { destination }
| Self::Unchanged { destination }
| Self::Orphaned { destination } => Some(destination),
}
}
#[must_use]
pub fn name(&self) -> String {
serde_json::to_value(self).expect("a contract enum serialises")["action"]
.as_str()
.expect("an internally tagged enum carries its tag")
.to_owned()
}
}
enum Target {
Update {
id: NativeId,
found: Found,
},
Create,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Found {
Origin,
Search,
}
enum Wanted {
Origin(String),
Title(String),
Metadata(String, Value),
}
impl Wanted {
fn found(&self, title: &str, metadata: &BTreeMap<String, Value>) -> bool {
match self {
Self::Origin(id) => {
metadata.get(GlobalId::ORIGIN_KEY) == Some(&Value::String(id.clone()))
}
Self::Title(wanted) => title == wanted,
Self::Metadata(key, value) => metadata.get(key) == Some(value),
}
}
}
#[derive(Clone)]
struct Prior {
item: Item,
edges: Vec<DependencyEdge>,
}
struct Deferred {
item: Planned,
filed: Option<NativeId>,
destination: NativeId,
prior: Option<Prior>,
}
enum Undo {
Created {
kind: Level,
id: NativeId,
},
Updated {
id: NativeId,
prior: Prior,
},
}
impl Undo {
fn id(&self) -> &NativeId {
match self {
Self::Created { id, .. } | Self::Updated { id, .. } => id,
}
}
fn kind(&self) -> Level {
match self {
Self::Created { kind, .. } => *kind,
Self::Updated { prior, .. } => prior.item.level(),
}
}
}
#[derive(Default)]
struct Journal {
entries: Vec<Undo>,
}
impl Journal {
fn record(&mut self, entry: Undo) {
if self
.entries
.iter()
.any(|held| held.kind() == entry.kind() && held.id() == entry.id())
{
return;
}
self.entries.push(entry);
}
}
struct Planned {
source: GlobalId,
item: Item,
edges: Vec<DependencyEdge>,
target: Target,
}
#[derive(Clone)]
enum Item {
Task(Box<Task>),
Project(Box<Project>),
Document(Box<Document>),
}
impl Item {
fn id(&self) -> &NativeId {
match self {
Self::Task(task) => &task.id,
Self::Project(project) => &project.id,
Self::Document(document) => &document.id,
}
}
fn level(&self) -> Level {
match self {
Self::Task(_) => Level::Task,
Self::Project(_) => Level::Project,
Self::Document(_) => Level::Document,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
enum Level {
Task,
Project,
Document,
}
#[derive(Clone)]
struct Referent {
id: GlobalId,
origin: Option<GlobalId>,
level: Level,
location: String,
}
impl Referent {
fn keys(&self) -> Vec<String> {
let mut keys = vec![self.id.to_string()];
if let Some(origin) = &self.origin {
let recorded = origin.to_string();
if recorded != keys[0] {
keys.push(recorded);
}
}
keys
}
}
enum Resolution {
Rewrite(String),
NoCounterpart,
Ambiguous,
}
#[derive(Default)]
struct Counterparts {
by_origin: BTreeMap<(Level, String), Vec<Held>>,
}
type Held = (NativeId, Option<String>);
impl Counterparts {
fn note(
&mut self,
level: Level,
id: &NativeId,
location: Option<&Location>,
metadata: &BTreeMap<String, Value>,
) {
let Some(origin) = origin_of(metadata) else {
return;
};
self.by_origin
.entry((level, origin.to_string()))
.or_default()
.push((id.clone(), located(location)));
}
fn resolve(&self, referent: &Referent) -> Resolution {
let mut candidates: Vec<&Held> = Vec::new();
for key in referent.keys() {
for record in self
.by_origin
.get(&(referent.level, key))
.into_iter()
.flatten()
{
if !candidates.iter().any(|held| held.0 == record.0) {
candidates.push(record);
}
}
}
match candidates.as_slice() {
[] => Resolution::NoCounterpart,
[(_, location)] => location
.clone()
.map_or(Resolution::NoCounterpart, Resolution::Rewrite),
_ => Resolution::Ambiguous,
}
}
}
impl Engine {
pub async fn copy(&self, request: &CopyRequest) -> Result<CopyReport, EngineError> {
let destination = self.writable(&request.destination)?;
if request.scope == CopyScope::Documents {
documentary(destination)?;
}
let mut journal = Journal::default();
match self.copy_all(destination, request, &mut journal).await {
Ok(report) => Ok(report),
Err(error) => Err(self.undo(destination, journal, error).await),
}
}
async fn copy_all(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
journal: &mut Journal,
) -> Result<CopyReport, EngineError> {
let mut written: BTreeMap<String, NativeId> = BTreeMap::new();
let mut deferred: Vec<Deferred> = Vec::new();
let mut references = Counted::default();
let mut membership = Vec::new();
let mut copied = Vec::new();
match request.scope {
CopyScope::Tasks | CopyScope::Documents => {
copied.extend(request.items.as_slice().iter().cloned());
}
CopyScope::Projects { tasks } => {
for id in request.items.as_slice() {
let members = if tasks {
self.project_members(id).await?
} else {
Vec::new()
};
copied.push(id.clone());
copied.extend(members.iter().cloned());
membership.push((id.clone(), members));
}
}
}
let items = match request.scope {
CopyScope::Tasks | CopyScope::Documents => {
self.copy_items(
destination,
request,
match request.scope {
CopyScope::Documents => Level::Document,
_ => Level::Task,
},
request.items.as_slice(),
None,
&copied,
&mut written,
&mut deferred,
journal,
&mut references,
)
.await?
}
CopyScope::Projects { tasks } => {
let mut items = Vec::new();
for (id, members) in &membership {
items.extend(
self.copy_project(
destination,
request,
id,
members,
tasks,
&copied,
&mut written,
&mut deferred,
journal,
&mut references,
)
.await?,
);
}
items
}
};
self.repair(destination, request, &copied, &written, deferred, journal)
.await?;
Ok(CopyReport {
items,
references_rewritten: references.rewritten,
references_unresolved: references.unresolved,
references_ambiguous: references.ambiguous,
})
}
async fn repair(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
copied: &[GlobalId],
written: &BTreeMap<String, NativeId>,
deferred: Vec<Deferred>,
journal: &mut Journal,
) -> Result<(), EngineError> {
if request.dry_run {
return Ok(());
}
for entry in deferred {
let edges = mapped_edges(
&entry.item.edges,
&entry.item.source.source,
destination,
copied,
written,
);
self.write(
destination,
&entry.item,
Some(entry.destination),
entry.filed,
&resolved(&edges),
entry.prior,
journal,
)
.await?;
}
Ok(())
}
async fn undo(
&self,
destination: &ResolvedSource,
journal: Journal,
error: EngineError,
) -> EngineError {
let created: Vec<(Level, &NativeId)> = journal
.entries
.iter()
.filter_map(|entry| match entry {
Undo::Created { kind, id } => Some((*kind, id)),
Undo::Updated { .. } => None,
})
.collect();
let mut unrestored: Option<(LeftBehind, SourceError)> = None;
for entry in journal.entries.iter().rev() {
let outcome = match entry {
Undo::Created { kind, id } => remove(destination, *kind, id).await,
Undo::Updated { id, prior, .. } if !created.contains(&(prior.item.level(), id)) => {
restore(destination, id, prior).await
}
Undo::Updated { .. } => Ok(()),
};
if let Err(problem) = outcome {
let id = GlobalId::new(destination.name().clone(), entry.id().clone());
match &mut unrestored {
Some((left_behind, _)) => left_behind.push(id),
None => unrestored = Some((LeftBehind::new(id), problem)),
}
}
}
match unrestored {
None => error,
Some((left_behind, refusal)) => EngineError::CopyNotUndone {
error: Box::new(error),
left_behind,
refusal,
},
}
}
fn writable(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
let name = self.known(name)?;
if let Some(unavailable) = self.unavailable().find(|source| source.name() == &name) {
return Err(EngineError::DestinationUnavailable {
name: name.to_string(),
error: unavailable.error().clone(),
});
}
let source = self
.ready()
.find(|source| source.name() == &name)
.ok_or(EngineError::NoSources)?;
if !source.source().writes().is_supported() {
return Err(EngineError::NotWritable {
name: name.to_string(),
kind: source.kind().to_owned(),
});
}
Ok(source)
}
#[allow(clippy::too_many_arguments)]
async fn copy_project(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
id: &GlobalId,
members: &[GlobalId],
tasks: bool,
copied: &[GlobalId],
written: &mut BTreeMap<String, NativeId>,
deferred: &mut Vec<Deferred>,
journal: &mut Journal,
references: &mut Counted,
) -> Result<Vec<CopyOutcome>, EngineError> {
let project_plan = self.plan(destination, request, Level::Project, id).await?;
let mut known = BTreeMap::new();
if let Target::Update { id: target, .. } = &project_plan.target {
known.insert(id.to_string(), target.clone());
}
for member in members {
let member_plan = self.plan(destination, request, Level::Task, member).await?;
if let Target::Update { id: target, .. } = member_plan.target {
known.insert(member.to_string(), target);
}
}
let project_was_unchanged = if let Target::Update { id: target, .. } = &project_plan.target
{
let edges = mapped_edges(&project_plan.edges, &id.source, destination, copied, &known);
let held = self.prior(destination, Level::Project, target).await?;
!edges.iter().any(Option::is_none)
&& !changes(
held.as_ref(),
&project_plan,
target,
&None,
&resolved(&edges),
)
} else {
false
};
let mut outcomes = self
.copy_items(
destination,
request,
Level::Project,
std::slice::from_ref(id),
None,
copied,
written,
deferred,
journal,
references,
)
.await?;
if !tasks {
return Ok(outcomes);
}
let project = outcomes.first().and_then(CopyOutcome::destination).cloned();
let task_outcomes = self
.copy_items(
destination,
request,
Level::Task,
members,
project.as_ref().map(|project| project.native.clone()),
copied,
written,
deferred,
journal,
references,
)
.await?;
outcomes.extend(task_outcomes);
if let Some(project) = project {
if project_was_unchanged {
outcomes[0].action = CopyAction::Unchanged {
destination: project.clone(),
};
}
outcomes.extend(
self.orphans(destination, id, &project.native, members)
.await?,
);
}
Ok(outcomes)
}
async fn project_members(&self, project: &GlobalId) -> Result<Vec<GlobalId>, EngineError> {
Ok(self
.project_member_tasks(project)
.await?
.into_iter()
.map(|task| task.id)
.collect())
}
async fn project_member_tasks(
&self,
project: &GlobalId,
) -> Result<Vec<Qualified<Task>>, EngineError> {
let mut request = TaskRequest {
sources: vec![project.source.clone()],
filters: Filters::default(),
project: ProjectSelector::Qualified(project.clone()),
paging: Paging {
limit: PROJECT_PAGE,
token: None,
},
};
let mut members = Vec::new();
let misbehaved = |error| EngineError::SourceRefused {
name: project.source.to_string(),
error,
};
loop {
let asked = request.paging.token.clone();
let response = self.tasks(&request).await?;
if let Some(failure) = response.errors.first() {
return Err(EngineError::SourceRefused {
name: failure.source.to_string(),
error: failure.error.clone(),
});
}
unrepeated(
response.next.as_ref(),
asked.as_ref(),
"the tasks of a project were being read for a copy",
)
.map_err(misbehaved)?;
members.extend(response.items);
match response.next {
Some(token) => request.paging.token = Some(token),
None => return Ok(members),
}
}
}
async fn project_documents(
&self,
project: &GlobalId,
) -> Result<Vec<Qualified<Document>>, EngineError> {
let mut request = DocumentRequest {
sources: vec![project.source.clone()],
filters: DocumentFilters::default(),
project: ProjectSelector::Qualified(project.clone()),
paging: Paging {
limit: PROJECT_PAGE,
token: None,
},
};
let mut held = Vec::new();
let misbehaved = |error| EngineError::SourceRefused {
name: project.source.to_string(),
error,
};
loop {
let asked = request.paging.token.clone();
let response = self.documents(&request).await?;
if let Some(failure) = response.errors.first() {
return Err(EngineError::SourceRefused {
name: failure.source.to_string(),
error: failure.error.clone(),
});
}
unrepeated(
response.next.as_ref(),
asked.as_ref(),
"the documents of a project were being read for a copy",
)
.map_err(misbehaved)?;
held.extend(response.items);
match response.next {
Some(token) => request.paging.token = Some(token),
None => return Ok(held),
}
}
}
async fn rewrite_references(
&self,
destination: &ResolvedSource,
planned: &mut [Planned],
counts: &mut Counted,
) -> Result<(), EngineError> {
let mut by_project: BTreeMap<String, Vec<Referent>> = BTreeMap::new();
let mut named: Vec<Vec<Referent>> = Vec::new();
for item in planned.iter() {
named.push(self.named_referents(item, &mut by_project).await?);
}
let mut levels: Vec<Level> = Vec::new();
for referent in named.iter().flatten() {
if !levels.contains(&referent.level) {
levels.push(referent.level);
}
}
if levels.is_empty() {
return Ok(());
}
let counterparts = self.counterparts(destination, &levels).await?;
for (item, referents) in planned.iter_mut().zip(named) {
let Item::Document(document) = &mut item.item else {
continue;
};
let Some(content) = &document.content else {
continue;
};
let (rewritten, made) = substitute(content, &table_for(&referents, &counterparts));
document.content = Some(rewritten);
counts.add(made);
}
Ok(())
}
async fn named_referents(
&self,
item: &Planned,
by_project: &mut BTreeMap<String, Vec<Referent>>,
) -> Result<Vec<Referent>, EngineError> {
let Item::Document(document) = &item.item else {
return Ok(Vec::new());
};
let (Some(content), Some(project)) =
(document.content.as_deref(), document.project.as_ref())
else {
return Ok(Vec::new());
};
if content.is_empty() {
return Ok(Vec::new());
}
let project = GlobalId::new(item.source.source.clone(), project.clone());
let key = project.to_string();
if !by_project.contains_key(&key) {
let read = self.referents(&project).await?;
by_project.insert(key.clone(), read);
}
Ok(by_project[&key]
.iter()
.filter(|referent| referent.level != Level::Document || referent.id != item.source)
.filter(|referent| holds(content, &referent.location))
.cloned()
.collect())
}
async fn referents(&self, project: &GlobalId) -> Result<Vec<Referent>, EngineError> {
let source = self.readable(&project.source)?;
let mut referents = Vec::new();
if let Some(held) = source
.source()
.get_project(&project.native)
.await
.map_err(|error| refused(source, error))?
{
note(
&mut referents,
project.clone(),
Level::Project,
held.location.as_ref(),
&held.metadata,
);
}
for task in self.project_member_tasks(project).await? {
note(
&mut referents,
task.id,
Level::Task,
task.item.location.as_ref(),
&task.item.metadata,
);
}
for document in self.project_documents(project).await? {
note(
&mut referents,
document.id,
Level::Document,
document.item.location.as_ref(),
&document.item.metadata,
);
}
Ok(referents)
}
async fn counterparts(
&self,
destination: &ResolvedSource,
levels: &[Level],
) -> Result<Counterparts, EngineError> {
let mut found = Counterparts::default();
for level in levels {
let mut asked_before: BTreeSet<String> = BTreeSet::new();
let mut cursor: Option<Cursor> = None;
loop {
if let Some(next) = &cursor
&& !asked_before.insert(next.0.clone())
{
return Err(refused(
destination,
SourceError::Malformed {
message: "the source returned a cursor it had already been \
given while the destination was being walked for the \
records a document's references name, so the walk \
would never end"
.to_owned(),
},
));
}
let asked = cursor.clone();
let request = request_for(destination, cursor);
let next = match level {
Level::Task => {
let page = destination
.source()
.query_tasks(&TaskQuery::default(), &request)
.await
.map_err(|error| refused(destination, error))?;
fits(page.items.len(), request.limit)
.map_err(|error| refused(destination, error))?;
for task in &page.items {
found.note(*level, &task.id, task.location.as_ref(), &task.metadata);
}
page.next
}
Level::Project => {
let page = destination
.source()
.query_projects(&ProjectQuery::default(), &request)
.await
.map_err(|error| refused(destination, error))?;
fits(page.items.len(), request.limit)
.map_err(|error| refused(destination, error))?;
for project in &page.items {
found.note(
*level,
&project.id,
project.location.as_ref(),
&project.metadata,
);
}
page.next
}
Level::Document => {
let page = destination
.source()
.query_documents(&DocumentQuery::default(), &request)
.await
.map_err(|error| refused(destination, error))?;
fits(page.items.len(), request.limit)
.map_err(|error| refused(destination, error))?;
for document in &page.items {
found.note(
*level,
&document.id,
document.location.as_ref(),
&document.metadata,
);
}
page.next
}
};
unrepeated(
next.as_ref(),
asked.as_ref(),
"the destination was being walked for the records a document's \
references name",
)
.map_err(|error| refused(destination, error))?;
match next {
Some(next) => cursor = Some(next),
None => break,
}
}
}
Ok(found)
}
async fn orphans(
&self,
destination: &ResolvedSource,
project: &GlobalId,
at_destination: &NativeId,
copied: &[GlobalId],
) -> Result<Vec<CopyOutcome>, EngineError> {
let mut orphans = Vec::new();
let mut cursor: Option<Cursor> = None;
loop {
let asked = cursor.clone();
let request = request_for(destination, cursor);
let page: Page<Task> = destination
.source()
.query_tasks(&TaskQuery::default(), &request)
.await
.map_err(|error| refused(destination, error))?;
fits(page.items.len(), request.limit).map_err(|error| refused(destination, error))?;
for task in &page.items {
if task.project.as_ref() != Some(at_destination) {
continue;
}
let Some(origin) = origin_of(&task.metadata) else {
continue;
};
if origin.source != project.source || copied.contains(&origin) {
continue;
}
orphans.push(CopyOutcome {
source: origin,
action: CopyAction::Orphaned {
destination: GlobalId::new(destination.name().clone(), task.id.clone()),
},
});
}
unrepeated(
page.next.as_ref(),
asked.as_ref(),
"the destination was being read for items the copy left behind",
)
.map_err(|error| refused(destination, error))?;
match page.next {
Some(next) => cursor = Some(next),
None => return Ok(orphans),
}
}
}
#[allow(clippy::too_many_arguments)]
async fn copy_items(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
kind: Level,
items: &[GlobalId],
project: Option<NativeId>,
copied: &[GlobalId],
written: &mut BTreeMap<String, NativeId>,
deferred: &mut Vec<Deferred>,
journal: &mut Journal,
references: &mut Counted,
) -> Result<Vec<CopyOutcome>, EngineError> {
let mut planned = Vec::new();
for id in items {
planned.push(self.plan(destination, request, kind, id).await?);
}
if kind == Level::Document {
self.rewrite_references(destination, &mut planned, references)
.await?;
}
for item in &planned {
if let Target::Update { id, .. } = &item.target {
written.insert(item.source.to_string(), id.clone());
}
}
let mut filed = Vec::new();
for item in &planned {
filed.push(self.filed(destination, item, project.clone()).await?);
}
let mut outcomes = Vec::new();
let mut unresolved = Vec::new();
let mut priors = Vec::new();
for (index, item) in planned.iter().enumerate() {
let edges = mapped_edges(
&item.edges,
&item.source.source,
destination,
copied,
written,
);
if edges.iter().any(Option::is_none) {
unresolved.push(index);
}
let (outcome, prior) = self
.land(
destination,
request,
item,
filed[index].clone(),
&edges,
journal,
)
.await?;
if let Some(id) = outcome.destination() {
written.insert(item.source.to_string(), id.native.clone());
}
outcomes.push(outcome);
priors.push(prior);
}
if !request.dry_run {
for (index, item) in planned.into_iter().enumerate() {
if !unresolved.contains(&index) {
continue;
}
let id = outcomes[index]
.destination()
.expect("a copy that writes lands every item it planned")
.clone();
deferred.push(Deferred {
item,
filed: filed[index].clone(),
destination: id.native,
prior: priors[index].clone(),
});
}
}
Ok(outcomes)
}
async fn plan(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
kind: Level,
id: &GlobalId,
) -> Result<Planned, EngineError> {
let source = self.readable(&id.source)?;
if kind == Level::Document {
documentary(source)?;
}
let item = match kind {
Level::Task => source
.source()
.get_task(&id.native)
.await
.map_err(|error| refused(source, error))?
.map(|task| Item::Task(Box::new(task))),
Level::Project => source
.source()
.get_project(&id.native)
.await
.map_err(|error| refused(source, error))?
.map(|project| Item::Project(Box::new(project))),
Level::Document => source
.source()
.get_document(&id.native)
.await
.map_err(|error| refused(source, error))?
.map(|document| Item::Document(Box::new(document))),
}
.ok_or_else(|| EngineError::NoSuchItem { id: id.to_string() })?;
let edges = forward_edges(source, &id.native, item.level()).await?;
let target = self.target(destination, request, id, &item).await?;
Ok(Planned {
source: id.clone(),
item,
edges,
target,
})
}
async fn target(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
id: &GlobalId,
item: &Item,
) -> Result<Target, EngineError> {
let (title, metadata) = described(item);
if let Some(origin) = origin_of(metadata)
&& &origin.source == destination.name()
{
if exists(destination, &origin.native, item.level()).await? {
return Ok(Target::Update {
id: origin.native,
found: Found::Origin,
});
}
if !request.recreate {
return Err(EngineError::StaleOrigin {
item: id.to_string(),
origin: origin.to_string(),
});
}
}
if let Some(found) = self
.scan(destination, item.level(), &Wanted::Origin(id.to_string()))
.await?
{
return Ok(Target::Update {
id: found,
found: Found::Search,
});
}
let wanted = match &request.match_by {
Some(MatchBy::Title) => Some(Wanted::Title(title.to_owned())),
Some(MatchBy::Metadata(key)) => metadata
.get(key)
.map(|value| Wanted::Metadata(key.clone(), value.clone())),
None => None,
};
if let Some(wanted) = wanted
&& let Some(found) = self.scan(destination, item.level(), &wanted).await?
{
return Ok(Target::Update {
id: found,
found: Found::Search,
});
}
Ok(Target::Create)
}
async fn scan(
&self,
destination: &ResolvedSource,
kind: Level,
wanted: &Wanted,
) -> Result<Option<NativeId>, EngineError> {
let mut cursor: Option<Cursor> = None;
loop {
let asked = cursor.clone();
let request = request_for(destination, cursor);
let next = match kind {
Level::Task => {
let page = destination
.source()
.query_tasks(&TaskQuery::default(), &request)
.await
.map_err(|error| refused(destination, error))?;
fits(page.items.len(), request.limit)
.map_err(|error| refused(destination, error))?;
for task in &page.items {
if wanted.found(&task.title, &task.metadata) {
return Ok(Some(task.id.clone()));
}
}
page.next
}
Level::Project => {
let page = destination
.source()
.query_projects(&ProjectQuery::default(), &request)
.await
.map_err(|error| refused(destination, error))?;
fits(page.items.len(), request.limit)
.map_err(|error| refused(destination, error))?;
for project in &page.items {
if wanted.found(&project.title, &project.metadata) {
return Ok(Some(project.id.clone()));
}
}
page.next
}
Level::Document => {
let page = destination
.source()
.query_documents(&DocumentQuery::default(), &request)
.await
.map_err(|error| refused(destination, error))?;
fits(page.items.len(), request.limit)
.map_err(|error| refused(destination, error))?;
for document in &page.items {
if wanted.found(&document.title, &document.metadata) {
return Ok(Some(document.id.clone()));
}
}
page.next
}
};
unrepeated(
next.as_ref(),
asked.as_ref(),
"the destination was being scanned for the item to update",
)
.map_err(|error| refused(destination, error))?;
match next {
Some(next) => cursor = Some(next),
None => return Ok(None),
}
}
}
async fn land(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
item: &Planned,
project: Option<NativeId>,
edges: &[Option<DependencyEdge>],
journal: &mut Journal,
) -> Result<(CopyOutcome, Option<Prior>), EngineError> {
let target = match &item.target {
Target::Update { id, .. } => Some(id.clone()),
Target::Create => None,
};
let prior = match &target {
Some(id) => self.prior(destination, item.item.level(), id).await?,
None => None,
};
let edges = resolved(edges);
let qualified = |native: NativeId| GlobalId::new(destination.name().clone(), native);
if let Some(id) = &target
&& !changes(prior.as_ref(), item, id, &project, &edges)
{
return Ok((
CopyOutcome {
source: item.source.clone(),
action: CopyAction::Unchanged {
destination: qualified(id.clone()),
},
},
prior,
));
}
if request.dry_run {
return Ok((
CopyOutcome {
source: item.source.clone(),
action: match target {
Some(id) => CopyAction::Updated {
destination: qualified(id),
},
None => CopyAction::Created { destination: None },
},
},
prior,
));
}
let updating = target.is_some();
let written = qualified(
self.write(
destination,
item,
target,
project,
&edges,
prior.clone(),
journal,
)
.await?,
);
Ok((
CopyOutcome {
source: item.source.clone(),
action: if updating {
CopyAction::Updated {
destination: written,
}
} else {
CopyAction::Created {
destination: Some(written),
}
},
},
prior,
))
}
async fn filed(
&self,
destination: &ResolvedSource,
item: &Planned,
project: Option<NativeId>,
) -> Result<Option<NativeId>, EngineError> {
match (&item.item, project) {
(Item::Task(task), None) => {
self.counterpart(destination, item, task.project.as_ref())
.await
}
(Item::Document(document), None) => {
self.counterpart(destination, item, document.project.as_ref())
.await
}
(Item::Task(_) | Item::Document(_), filed) => Ok(filed),
(Item::Project(_), _) => Ok(None),
}
}
async fn counterpart(
&self,
destination: &ResolvedSource,
item: &Planned,
project: Option<&NativeId>,
) -> Result<Option<NativeId>, EngineError> {
let Some(project) = project else {
return Ok(None);
};
let qualified = GlobalId::new(item.source.source.clone(), project.clone());
let found = self
.scan(
destination,
Level::Project,
&Wanted::Origin(qualified.to_string()),
)
.await?;
Ok(Some(found.unwrap_or_else(|| project.clone())))
}
async fn prior(
&self,
destination: &ResolvedSource,
kind: Level,
id: &NativeId,
) -> Result<Option<Prior>, EngineError> {
let held = match kind {
Level::Task => destination
.source()
.get_task(id)
.await
.map_err(|error| refused(destination, error))?
.map(|task| Item::Task(Box::new(task))),
Level::Project => destination
.source()
.get_project(id)
.await
.map_err(|error| refused(destination, error))?
.map(|project| Item::Project(Box::new(project))),
Level::Document => destination
.source()
.get_document(id)
.await
.map_err(|error| refused(destination, error))?
.map(|document| Item::Document(Box::new(document))),
};
let Some(item) = held else {
return Ok(None);
};
let edges = forward_edges(destination, id, kind).await?;
Ok(Some(Prior { item, edges }))
}
#[allow(clippy::too_many_arguments)]
async fn write(
&self,
destination: &ResolvedSource,
item: &Planned,
target: Option<NativeId>,
project: Option<NativeId>,
edges: &[DependencyEdge],
prior: Option<Prior>,
journal: &mut Journal,
) -> Result<NativeId, EngineError> {
let created_kind = item.item.level();
let suggested = target.clone().unwrap_or_else(|| item.item.id().clone());
let origin = recorded(item, prior.as_ref());
if let (Some(id), Some(prior)) = (target.clone(), prior) {
journal.record(Undo::Updated { id, prior });
}
let landed = match outgoing(item, suggested, project, &origin) {
Item::Task(task) => destination
.source()
.write_task(&ItemWrite {
target: target.clone(),
item: *task,
depends_on: edges.to_vec(),
})
.await
.map_err(|error| refused(destination, error))?,
Item::Project(project) => destination
.source()
.write_project(&ItemWrite {
target: target.clone(),
item: *project,
depends_on: edges.to_vec(),
})
.await
.map_err(|error| refused(destination, error))?,
Item::Document(document) => destination
.source()
.write_document(&ItemWrite {
target: target.clone(),
item: *document,
depends_on: Vec::new(),
})
.await
.map_err(|error| refused(destination, error))?,
};
if target.is_none() {
journal.record(Undo::Created {
kind: created_kind,
id: landed.clone(),
});
}
Ok(landed)
}
fn readable(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
let name = self.known(name)?;
if let Some(unavailable) = self.unavailable().find(|source| source.name() == &name) {
return Err(EngineError::SourceRefused {
name: name.to_string(),
error: unavailable.error().clone(),
});
}
self.ready()
.find(|source| source.name() == &name)
.ok_or(EngineError::NoSources)
}
}
const PROJECT_PAGE: std::num::NonZeroU32 = std::num::NonZeroU32::new(50).expect("50 is not zero");
fn request_for(source: &ResolvedSource, cursor: Option<Cursor>) -> PageRequest {
PageRequest {
cursor,
limit: source.source().capabilities().max_page_size.max(1),
}
}
fn changes(
held: Option<&Prior>,
item: &Planned,
target: &NativeId,
project: &Option<NativeId>,
edges: &[DependencyEdge],
) -> bool {
let Some(held) = held else {
return true;
};
let outgoing = outgoing(
item,
target.clone(),
project.clone(),
&recorded(item, Some(held)),
);
!same(&held.item, &outgoing) || !same_edges(&held.edges, edges)
}
async fn remove(
destination: &ResolvedSource,
kind: Level,
id: &NativeId,
) -> Result<(), SourceError> {
match kind {
Level::Task => destination.source().delete_task(id).await,
Level::Project => destination.source().delete_project(id).await,
Level::Document => destination.source().delete_document(id).await,
}
}
async fn restore(
destination: &ResolvedSource,
id: &NativeId,
prior: &Prior,
) -> Result<(), SourceError> {
match &prior.item {
Item::Task(task) => destination
.source()
.write_task(&ItemWrite {
target: Some(id.clone()),
item: (**task).clone(),
depends_on: prior.edges.clone(),
})
.await
.map(|_| ()),
Item::Project(project) => destination
.source()
.write_project(&ItemWrite {
target: Some(id.clone()),
item: (**project).clone(),
depends_on: prior.edges.clone(),
})
.await
.map(|_| ()),
Item::Document(document) => destination
.source()
.write_document(&ItemWrite {
target: Some(id.clone()),
item: (**document).clone(),
depends_on: Vec::new(),
})
.await
.map(|_| ()),
}
}
fn documentary(source: &ResolvedSource) -> Result<(), EngineError> {
if source.source().capabilities().documents.is_native() {
return Ok(());
}
Err(EngineError::NoDocuments {
name: source.name().to_string(),
kind: source.kind().to_owned(),
})
}
fn refused(source: &ResolvedSource, error: SourceError) -> EngineError {
EngineError::SourceRefused {
name: source.name().to_string(),
error,
}
}
async fn forward_edges(
source: &ResolvedSource,
id: &NativeId,
kind: Level,
) -> Result<Vec<DependencyEdge>, EngineError> {
if kind == Level::Document {
return Ok(Vec::new());
}
let mut edges = Vec::new();
let mut cursor: Option<Cursor> = None;
loop {
let asked = cursor.clone();
let request = request_for(source, cursor);
let page = match kind {
Level::Task | Level::Document => {
source
.source()
.task_dependencies(id, Direction::DependsOn, &request)
.await
}
Level::Project => {
source
.source()
.project_dependencies(id, Direction::DependsOn, &request)
.await
}
}
.map_err(|error| refused(source, error))?;
fits(page.items.len(), request.limit).map_err(|error| refused(source, error))?;
edges.extend(page.items);
unrepeated(
page.next.as_ref(),
asked.as_ref(),
"an item's dependencies were being read for a copy",
)
.map_err(|error| refused(source, error))?;
match page.next {
Some(next) => cursor = Some(next),
None => return Ok(edges),
}
}
}
fn located(location: Option<&Location>) -> Option<String> {
let (Location::Path(held) | Location::Url(held)) = location?;
(!held.is_empty()).then(|| held.clone())
}
fn note(
into: &mut Vec<Referent>,
id: GlobalId,
level: Level,
location: Option<&Location>,
metadata: &BTreeMap<String, Value>,
) {
if let Some(location) = located(location) {
into.push(Referent {
id,
origin: origin_of(metadata),
level,
location,
});
}
}
fn table_for(referents: &[Referent], counterparts: &Counterparts) -> Vec<(String, Resolution)> {
let mut table: Vec<(String, Resolution)> = Vec::new();
for referent in referents {
if let Some(held) = table
.iter_mut()
.find(|(location, _)| location == &referent.location)
{
held.1 = Resolution::Ambiguous;
continue;
}
table.push((referent.location.clone(), counterparts.resolve(referent)));
}
table.sort_by_key(|(location, _)| std::cmp::Reverse(location.len()));
table
}
fn holds(content: &str, location: &str) -> bool {
(0..content.len()).any(|at| delimited_at(content, at, location))
}
fn delimited_at(content: &str, at: usize, location: &str) -> bool {
if !content.is_char_boundary(at) || !content[at..].starts_with(location) {
return false;
}
let before = content[..at].chars().next_back();
let after = content[at + location.len()..].chars().next();
stops_a_location(before) && stops_a_location(after)
}
fn stops_a_location(character: Option<char>) -> bool {
match character {
None => true,
Some(character) => character.is_whitespace() || "`\"'()[]{}<>|,;".contains(character),
}
}
fn substitute(content: &str, table: &[(String, Resolution)]) -> (String, Counted) {
let mut written = String::with_capacity(content.len());
let mut counts = Counted::default();
let mut at = 0;
while at < content.len() {
if let Some((location, resolution)) = table
.iter()
.find(|(location, _)| delimited_at(content, at, location))
{
match resolution {
Resolution::Rewrite(there) => {
written.push_str(there);
counts.rewritten += 1;
}
Resolution::NoCounterpart => {
written.push_str(location);
counts.unresolved += 1;
}
Resolution::Ambiguous => {
written.push_str(location);
counts.unresolved += 1;
counts.ambiguous += 1;
}
}
at += location.len();
continue;
}
let character = content[at..]
.chars()
.next()
.expect("a character at a boundary this walk only ever lands on");
written.push(character);
at += character.len_utf8();
}
(written, counts)
}
fn origin_of(metadata: &BTreeMap<String, Value>) -> Option<GlobalId> {
metadata
.get(GlobalId::ORIGIN_KEY)?
.as_str()?
.parse::<GlobalId>()
.ok()
}
fn described(item: &Item) -> (&str, &BTreeMap<String, Value>) {
match item {
Item::Task(task) => (&task.title, &task.metadata),
Item::Project(project) => (&project.title, &project.metadata),
Item::Document(document) => (&document.title, &document.metadata),
}
}
fn outgoing(item: &Planned, id: NativeId, project: Option<NativeId>, origin: &Origin) -> Item {
match &item.item {
Item::Task(task) => Item::Task(Box::new(Task {
id,
url: None,
location: None,
created_at: None,
updated_at: None,
project,
metadata: carried(&task.metadata, origin),
..(**task).clone()
})),
Item::Project(project) => Item::Project(Box::new(Project {
id,
url: None,
location: None,
created_at: None,
updated_at: None,
metadata: carried(&project.metadata, origin),
..(**project).clone()
})),
Item::Document(document) => Item::Document(Box::new(Document {
id,
url: None,
location: None,
created_at: None,
updated_at: None,
project,
metadata: carried(&document.metadata, origin),
..(**document).clone()
})),
}
}
fn carried(metadata: &BTreeMap<String, Value>, origin: &Origin) -> BTreeMap<String, Value> {
let mut carried = metadata.clone();
carried.remove(Repository::METADATA_KEY);
carried.remove(DependencyEdge::RECORDED_KEY);
carried.remove(GlobalId::ORIGIN_KEY);
let held = match origin {
Origin::Records(id) => Some(Value::String(id.to_string())),
Origin::Keeps(held) => held.clone(),
};
if let Some(held) = held {
carried.insert(GlobalId::ORIGIN_KEY.to_owned(), held);
}
carried
}
enum Origin {
Records(GlobalId),
Keeps(Option<Value>),
}
fn recorded(item: &Planned, held: Option<&Prior>) -> Origin {
if let Target::Update {
found: Found::Origin,
..
} = &item.target
{
return Origin::Keeps(
held.and_then(|held| described(&held.item).1.get(GlobalId::ORIGIN_KEY).cloned()),
);
}
Origin::Records(item.source.clone())
}
fn same(held: &Item, outgoing: &Item) -> bool {
match (held, outgoing) {
(Item::Task(held), Item::Task(outgoing)) => {
held.title == outgoing.title
&& held.content == outgoing.content
&& held.status == outgoing.status
&& held.labels == outgoing.labels
&& held.project == outgoing.project
&& held.metadata == outgoing.metadata
&& held.repositories == outgoing.repositories
}
(Item::Project(held), Item::Project(outgoing)) => {
held.title == outgoing.title
&& held.content == outgoing.content
&& held.status == outgoing.status
&& held.labels == outgoing.labels
&& held.metadata == outgoing.metadata
&& held.repositories == outgoing.repositories
}
(Item::Document(held), Item::Document(outgoing)) => {
held.title == outgoing.title
&& held.content == outgoing.content
&& held.labels == outgoing.labels
&& held.project == outgoing.project
&& held.metadata == outgoing.metadata
&& held.repositories == outgoing.repositories
}
_ => false,
}
}
fn same_edges(held: &[DependencyEdge], outgoing: &[DependencyEdge]) -> bool {
let ends = |edges: &[DependencyEdge]| {
let mut ends: Vec<(String, ItemKind, DependencyKind)> = edges
.iter()
.map(|edge| (edge.to.id().to_owned(), edge.to.kind, edge.kind))
.collect();
ends.sort_by(|left, right| left.0.cmp(&right.0));
ends
};
ends(held) == ends(outgoing)
}
fn mapped_edges(
edges: &[DependencyEdge],
origin: &SourceName,
destination: &ResolvedSource,
copied: &[GlobalId],
written: &BTreeMap<String, NativeId>,
) -> Vec<Option<DependencyEdge>> {
edges
.iter()
.map(|edge| {
let far = GlobalId::new(origin.clone(), NativeId(edge.to.id().to_owned()));
let id = if let Some(native) = names(&edge.to, destination.name()) {
Some(native)
} else if edge.to.is_qualified() || origin == destination.name() {
Some(edge.to.id().to_owned())
} else if copied.contains(&far) {
written.get(&far.to_string()).map(|native| native.0.clone())
} else {
Some(far.to_string())
}?;
DependencyEndpoint::new(id, edge.to.kind)
.ok()
.map(|to| DependencyEdge {
from: edge.from.clone(),
to,
kind: edge.kind,
})
})
.collect()
}
fn names(endpoint: &DependencyEndpoint, destination: &SourceName) -> Option<String> {
if !endpoint.is_qualified() {
return None;
}
let id: GlobalId = endpoint.id().parse().ok()?;
(&id.source == destination).then_some(id.native.0)
}
fn resolved(edges: &[Option<DependencyEdge>]) -> Vec<DependencyEdge> {
edges.iter().flatten().cloned().collect()
}
async fn exists(
destination: &ResolvedSource,
id: &NativeId,
kind: Level,
) -> Result<bool, EngineError> {
let found = match kind {
Level::Task => destination
.source()
.get_task(id)
.await
.map_err(|error| refused(destination, error))?
.is_some(),
Level::Project => destination
.source()
.get_project(id)
.await
.map_err(|error| refused(destination, error))?
.is_some(),
Level::Document => destination
.source()
.get_document(id)
.await
.map_err(|error| refused(destination, error))?
.is_some(),
};
Ok(found)
}