use std::collections::{BTreeMap, BTreeSet};
use onetaskgraph_plugin_api::{
Capabilities, Cursor, DependencyEdge, DependencyEndpoint, DependencyKind, Direction, Document,
DocumentQuery, ItemKind, ItemWrite, Location, MetadataKey, MetadataMatch, Metering, NativeId,
Page, PageRequest, Project, ProjectQuery, Repository, SourceError, SourceName, StatusCategory,
Task, TaskQuery, TaskRef, TextFields, TextQuery,
};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::GlobalId;
use crate::resolve::ResolvedSource;
use crate::template::{Sha256Digest, TemplateProvenance, body_digest};
use super::delivery::{Delivered, targets};
use super::fetch::{fits, unrepeated};
use super::local::ProjectSelector;
use super::narrow::holds_priority;
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, PartialEq, Eq)]
pub enum CopyScope {
Tasks,
Projects {
tasks: bool,
},
Members(CopyItems),
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,
#[serde(default, skip_serializing_if = "nothing_to_report")]
#[schemars(!skip_serializing_if)]
pub delivers_rewritten: u64,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
#[schemars(!skip_serializing_if)]
pub delivered: Vec<Delivered>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub spent: Option<Spent>,
}
fn nothing_to_report(figure: &u64) -> bool {
*figure == 0
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
pub struct Spent {
pub requests: u64,
pub budgets: Vec<BudgetSpent>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
pub struct BudgetSpent {
pub budget: String,
pub unit: String,
pub amount: u64,
pub lower_bound: bool,
}
#[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,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "kebab-case")]
pub enum CopyVia {
Link,
Origin,
Scan,
Match,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "kebab-case")]
pub enum NoCounterpart {
#[default]
Created,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "kebab-case")]
pub enum CopyLink {
Recorded,
Unchanged,
Unrecorded,
}
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>,
#[serde(default)]
via: NoCounterpart,
#[serde(default, skip_serializing_if = "Option::is_none")]
link: Option<CopyLink>,
},
Updated {
destination: GlobalId,
via: CopyVia,
#[serde(default, skip_serializing_if = "Option::is_none")]
link: Option<CopyLink>,
},
Unchanged {
destination: GlobalId,
via: CopyVia,
#[serde(default, skip_serializing_if = "Option::is_none")]
link: Option<CopyLink>,
},
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 found_by(&self) -> Option<CopyVia> {
match self {
Self::Updated { via, .. } | Self::Unchanged { via, .. } => Some(*via),
Self::Created { .. } | Self::Orphaned { .. } => None,
}
}
#[must_use]
pub fn link(&self) -> Option<CopyLink> {
match self {
Self::Created { link, .. }
| Self::Updated { link, .. }
| Self::Unchanged { link, .. } => *link,
Self::Orphaned { .. } => None,
}
}
fn link_slot(&mut self) -> Option<&mut Option<CopyLink>> {
match self {
Self::Created { link, .. }
| Self::Updated { link, .. }
| Self::Unchanged { link, .. } => Some(link),
Self::Orphaned { .. } => None,
}
}
#[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()
}
}
struct Pointing<'a> {
edges: &'a [Option<DependencyEdge>],
delivers: &'a [TaskRef],
}
enum Target {
Update {
id: NativeId,
found: Found,
},
Create,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Found {
Link,
Origin,
Scan,
Match,
}
impl Found {
fn via(self) -> CopyVia {
match self {
Self::Link => CopyVia::Link,
Self::Origin => CopyVia::Origin,
Self::Scan => CopyVia::Scan,
Self::Match => CopyVia::Match,
}
}
}
impl Target {
fn found(&self) -> Option<Found> {
match self {
Self::Update { found, .. } => Some(*found),
Self::Create => None,
}
}
}
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),
}
}
fn task_query(&self, capabilities: &Capabilities) -> TaskQuery {
match self {
Self::Origin(id) if capabilities.filter_by_origin.is_native() => TaskQuery {
origin: Some(id.clone()),
..TaskQuery::default()
},
Self::Title(title) if capabilities.search_title.is_native() && has_words(title) => {
TaskQuery {
text: Some(TextQuery {
terms: title.clone(),
fields: TextFields::Title,
}),
..TaskQuery::default()
}
}
Self::Metadata(key, Value::String(value))
if capabilities.filter_by_metadata.is_native() && has_words(value) =>
{
MetadataMatch::new(key.clone(), Vec::new(), value.clone())
.map(|wanted| TaskQuery {
metadata: vec![wanted],
..TaskQuery::default()
})
.unwrap_or_default()
}
_ => TaskQuery::default(),
}
}
}
fn has_words(terms: &str) -> bool {
terms.chars().any(char::is_alphanumeric)
}
#[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>,
links: Vec<Unlink>,
}
struct Unlink {
item: GlobalId,
level: Level,
before: Option<Value>,
}
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);
}
}
#[derive(Default)]
struct Running {
resolvable: Vec<GlobalId>,
counterparts: BTreeMap<String, NativeId>,
deferred: Vec<Deferred>,
references: Counted,
filings: BTreeMap<String, Option<NativeId>>,
delivers_rewritten: u64,
landed: Vec<LandedTask>,
linking: Vec<Linking>,
}
struct Linking {
item: GlobalId,
level: Level,
found: Option<Found>,
destination: GlobalId,
held: Option<Value>,
}
struct LandedTask {
destination: GlobalId,
origin: SourceName,
delivers: Vec<TaskRef>,
before: Vec<GlobalId>,
category: StatusCategory,
}
struct Deliverer {
destination: GlobalId,
category: StatusCategory,
now: Vec<GlobalId>,
before: Vec<GlobalId>,
}
struct Planned {
source: GlobalId,
item: Item,
edges: Vec<DependencyEdge>,
target: Target,
held: Option<Prior>,
}
#[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 metered = vec![destination];
for id in request.items.as_slice() {
if let Some(source) = self.ready().find(|source| source.name() == &id.source)
&& !metered.iter().any(|held| held.name() == source.name())
{
metered.push(source);
}
}
let before = readings(&metered).await;
let mut journal = Journal::default();
match self.copy_all(destination, request, &mut journal).await {
Ok((mut report, deliverers)) => {
for deliverer in &deliverers {
report.delivered.extend(
self.deliver(
&deliverer.destination,
deliverer.category,
&deliverer.now,
&deliverer.before,
)
.await,
);
}
report.spent = spent_between(&before, &readings(&metered).await);
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, Vec<Deliverer>), EngineError> {
let mut running = Running::default();
let items = match &request.scope {
CopyScope::Tasks | CopyScope::Documents => {
let kind = if request.scope == CopyScope::Documents {
Level::Document
} else {
Level::Task
};
running
.resolvable
.extend(request.items.as_slice().iter().cloned());
let mut planned = Vec::new();
for id in request.items.as_slice() {
planned.push(self.plan(destination, request, kind, id).await?);
}
self.copy_items(destination, request, planned, None, &mut running, journal)
.await?
}
CopyScope::Projects { tasks } => {
let mut projects = Vec::new();
for id in request.items.as_slice() {
let members = if *tasks {
self.project_members(id).await?
} else {
Vec::new()
};
running.resolvable.push(id.clone());
running.resolvable.extend(members.iter().cloned());
projects.push((id.clone(), members));
}
self.copy_projects(destination, request, &projects, &[], &mut running, journal)
.await?
}
CopyScope::Members(named) => {
let (projects, unrecorded) = self
.named_members(destination, request.items.as_slice(), named, &mut running)
.await?;
self.copy_projects(
destination,
request,
&projects,
&unrecorded,
&mut running,
journal,
)
.await?
}
};
let references = running.references;
let delivers_rewritten = running.delivers_rewritten;
let deliverers: Vec<Deliverer> = running
.landed
.iter()
.map(|task| Deliverer {
destination: task.destination.clone(),
category: task.category,
now: targets(
&resolved_entries(&mapped_delivers(
&task.delivers,
&task.origin,
destination,
&running.resolvable,
&running.counterparts,
)),
destination.name(),
),
before: task.before.clone(),
})
.collect();
let linking = std::mem::take(&mut running.linking);
self.repair(destination, request, running, journal).await?;
let mut items = items;
self.link(destination, linking, &mut items, journal).await?;
Ok((
CopyReport {
items,
references_rewritten: references.rewritten,
references_unresolved: references.unresolved,
references_ambiguous: references.ambiguous,
delivers_rewritten,
delivered: Vec::new(),
spent: None,
},
deliverers,
))
}
async fn repair(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
running: Running,
journal: &mut Journal,
) -> Result<(), EngineError> {
if request.dry_run {
return Ok(());
}
let Running {
resolvable,
counterparts,
deferred,
..
} = running;
for entry in deferred {
let edges = mapped_edges(
&entry.item.edges,
&entry.item.source.source,
destination,
&resolvable,
&counterparts,
);
let delivers = delivers_of(&entry.item, destination, &resolvable, &counterparts);
self.write(
destination,
&entry.item,
Some(entry.destination),
entry.filed,
&resolved(&edges),
&resolved_entries(&delivers),
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.links.iter().rev() {
if let Err(problem) = self.unlink(entry).await {
match &mut unrestored {
Some((left_behind, _)) => left_behind.push(entry.item.clone()),
None => unrestored = Some((LeftBehind::new(entry.item.clone()), problem)),
}
}
}
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,
},
}
}
async fn link(
&self,
destination: &ResolvedSource,
linking: Vec<Linking>,
items: &mut [CopyOutcome],
journal: &mut Journal,
) -> Result<(), EngineError> {
let mut landed = items
.iter_mut()
.filter(|outcome| !matches!(outcome.action, CopyAction::Orphaned { .. }));
for entry in linking {
let slot = landed
.next()
.filter(|outcome| outcome.source == entry.item)
.and_then(|outcome| outcome.action.link_slot())
.expect("every item a copy landed is reported, in the order it landed");
*slot = Some(self.linked(destination, entry, journal).await?);
}
Ok(())
}
async fn linked(
&self,
destination: &ResolvedSource,
entry: Linking,
journal: &mut Journal,
) -> Result<CopyLink, EngineError> {
if entry.found == Some(Found::Origin) {
return Ok(CopyLink::Unchanged);
}
let wanted = Value::String(entry.destination.to_string());
let mut links = well_formed(entry.held.as_ref());
if links.get(destination.name().as_str()) == Some(&wanted) {
return Ok(CopyLink::Unchanged);
}
let source = self.readable(&entry.item.source)?;
if !source.source().writes().is_supported() {
return Ok(CopyLink::Unrecorded);
}
links.insert(destination.name().to_string(), wanted);
journal.links.push(Unlink {
item: entry.item.clone(),
level: entry.level,
before: entry.held,
});
match set_link(
source,
entry.level,
&entry.item.native,
&Value::Object(links),
)
.await
{
Ok(true) => Ok(CopyLink::Recorded),
Ok(false) | Err(SourceError::Refused { .. }) => {
journal.links.pop();
Ok(CopyLink::Unrecorded)
}
Err(SourceError::Malformed { .. }) => Ok(CopyLink::Unrecorded),
Err(error) => Err(refused(source, error)),
}
}
async fn unlink(&self, entry: &Unlink) -> Result<(), SourceError> {
let as_source_error = |error: EngineError| match error {
EngineError::SourceRefused { error, .. } => error,
other => SourceError::Refused {
message: other.to_string(),
},
};
let source = self.readable(&entry.item.source).map_err(as_source_error)?;
if let Some(before) = &entry.before {
return set_link(source, entry.level, &entry.item.native, before)
.await
.map(|_| ());
}
let Some(mut held) = self
.prior(source, entry.level, &entry.item.native)
.await
.map_err(as_source_error)?
else {
return Ok(());
};
metadata_of(&mut held.item).remove(MetadataKey::COPIES_KEY);
restore(source, &entry.item.native, &held).await
}
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)
}
async fn copy_projects(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
projects: &[(GlobalId, Vec<GlobalId>)],
unrecorded: &[GlobalId],
running: &mut Running,
journal: &mut Journal,
) -> Result<Vec<CopyOutcome>, EngineError> {
let mut plans = Vec::new();
for (id, members) in projects {
let project = self.plan(destination, request, Level::Project, id).await?;
let mut tasks = Vec::new();
for member in members {
tasks.push(self.plan(destination, request, Level::Task, member).await?);
}
plans.push((project, tasks));
}
for (project, tasks) in &plans {
for item in std::iter::once(project).chain(tasks) {
unrecorded_far_end(item, destination, unrecorded)?;
}
}
let mut outcomes = Vec::new();
for (project, tasks) in plans {
outcomes.extend(
self.copy_project(destination, request, project, tasks, running, journal)
.await?,
);
}
Ok(outcomes)
}
async fn copy_project(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
project: Planned,
tasks: Vec<Planned>,
running: &mut Running,
journal: &mut Journal,
) -> Result<Vec<CopyOutcome>, EngineError> {
let (carries_tasks, walks_orphans) = match request.scope {
CopyScope::Projects { tasks } => (tasks, tasks),
CopyScope::Members(_) => (true, false),
CopyScope::Tasks | CopyScope::Documents => (false, false),
};
let id = project.source.clone();
let members: Vec<GlobalId> = tasks.iter().map(|task| task.source.clone()).collect();
let mut before_members = running.counterparts.clone();
if let Target::Update { id: target, .. } = &project.target {
before_members.insert(id.to_string(), target.clone());
}
for item in std::iter::once(&project).chain(&tasks) {
if let Target::Update { id: target, .. } = &item.target {
running
.counterparts
.insert(item.source.to_string(), target.clone());
}
}
let unchanged = match &project.target {
Target::Update { id: target, .. } => {
let settled = mapped_edges(
&project.edges,
&id.source,
destination,
&running.resolvable,
&running.counterparts,
);
let first = mapped_edges(
&project.edges,
&id.source,
destination,
&running.resolvable,
&before_members,
);
let held = project.held.as_ref();
let unchanged_with = |edges: &[Option<DependencyEdge>]| {
!changes(
held,
&project,
target,
&None,
&resolved(edges),
&[],
destination.name(),
)
};
Some(
(!settled.iter().any(Option::is_none) && unchanged_with(&settled))
|| unchanged_with(&first),
)
}
Target::Create => None,
};
let created = matches!(project.target, Target::Create);
let mut outcomes = self
.copy_items(destination, request, vec![project], None, running, journal)
.await?;
if let (Some(unchanged), Some(landed), Some(via)) = (
unchanged,
outcomes[0].destination().cloned(),
outcomes[0].action.found_by(),
) {
let link = outcomes[0].action.link();
outcomes[0].action = if unchanged {
CopyAction::Unchanged {
destination: landed,
via,
link,
}
} else {
CopyAction::Updated {
destination: landed,
via,
link,
}
};
}
if !carries_tasks {
return Ok(outcomes);
}
let filed = outcomes.first().and_then(CopyOutcome::destination).cloned();
outcomes.extend(
self.copy_items(
destination,
request,
tasks,
filed.as_ref().map(|project| project.native.clone()),
running,
journal,
)
.await?,
);
if walks_orphans
&& !created
&& let Some(filed) = filed
{
outcomes.extend(
self.orphans(destination, &id, &filed.native, &members)
.await?,
);
}
Ok(outcomes)
}
async fn named_members(
&self,
destination: &ResolvedSource,
projects: &[GlobalId],
named: &CopyItems,
running: &mut Running,
) -> Result<(Vec<(GlobalId, Vec<GlobalId>)>, Vec<GlobalId>), EngineError> {
let mut carried = Vec::new();
let mut unrecorded = Vec::new();
for project in projects {
let held = self.project_member_tasks(project).await?;
let mut members: Vec<GlobalId> = Vec::new();
for id in named.as_slice() {
if held.iter().any(|task| &task.id == id) && !members.contains(id) {
members.push(id.clone());
}
}
for task in held {
if members.contains(&task.id) {
continue;
}
match origin_of(&task.item.metadata)
.filter(|origin| &origin.source == destination.name())
{
Some(origin) => {
running
.counterparts
.insert(task.id.to_string(), origin.native);
running.resolvable.push(task.id);
}
None => unrecorded.push(task.id),
}
}
running.resolvable.push(project.clone());
running.resolvable.extend(members.iter().cloned());
carried.push((project.clone(), members));
}
if let Some(stray) = named
.as_slice()
.iter()
.find(|id| !carried.iter().any(|(_, members)| members.contains(id)))
{
return Err(EngineError::NotAMember {
id: stray.clone(),
projects: projects.to_vec(),
});
}
Ok((carried, unrecorded))
}
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()),
priorities: Vec::new(),
commented_since: None,
metadata: Vec::new(),
origin: None,
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));
if made.rewritten > 0
&& let Some(provenance) = restamped(&document.metadata, content, &rewritten)
{
document
.metadata
.insert(TemplateProvenance::KEY.to_owned(), provenance.to_value());
}
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),
}
}
}
async fn copy_items(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
mut planned: Vec<Planned>,
project: Option<NativeId>,
running: &mut Running,
journal: &mut Journal,
) -> Result<Vec<CopyOutcome>, EngineError> {
if planned
.iter()
.any(|item| item.item.level() == Level::Document)
{
self.rewrite_references(destination, &mut planned, &mut running.references)
.await?;
}
for item in &planned {
if let Target::Update { id, .. } = &item.target {
running
.counterparts
.insert(item.source.to_string(), id.clone());
}
if let Item::Task(task) = &item.item {
running.delivers_rewritten +=
members_named(&task.delivers, &item.source.source, &running.resolvable);
}
}
let mut filed = Vec::new();
for item in &planned {
filed.push(
self.filed(destination, item, project.clone(), running)
.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,
&running.resolvable,
&running.counterparts,
);
let delivers = delivers_of(
item,
destination,
&running.resolvable,
&running.counterparts,
);
if edges.iter().any(Option::is_none) || delivers.iter().any(Option::is_none) {
unresolved.push(index);
}
let resolvable_delivers = resolved_entries(&delivers);
let (outcome, prior) = self
.land(
destination,
request,
item,
filed[index].clone(),
Pointing {
edges: &edges,
delivers: &resolvable_delivers,
},
journal,
)
.await?;
if let Some(id) = outcome.destination() {
running
.counterparts
.insert(item.source.to_string(), id.native.clone());
if !request.dry_run {
running.linking.push(Linking {
item: item.source.clone(),
level: item.item.level(),
found: item.target.found(),
destination: id.clone(),
held: described(&item.item)
.1
.get(MetadataKey::COPIES_KEY)
.cloned(),
});
}
}
if !request.dry_run
&& let (Item::Task(task), Some(landed)) = (&item.item, outcome.destination())
{
running.landed.push(LandedTask {
destination: landed.clone(),
origin: item.source.source.clone(),
delivers: task.delivers.clone(),
before: match prior.as_ref().map(|prior| &prior.item) {
Some(Item::Task(held)) => targets(&held.delivers, destination.name()),
_ => Vec::new(),
},
category: task.status.category,
});
}
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();
running.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() })?;
if let Item::Task(task) = &item {
holds_priority(destination, &id.to_string(), task.priority)?;
}
let edges = forward_edges(source, &id.native, item.level()).await?;
let (target, held) = self.target(destination, request, id, &item).await?;
Ok(Planned {
source: id.clone(),
item,
edges,
target,
held,
})
}
async fn target(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
id: &GlobalId,
item: &Item,
) -> Result<(Target, Option<Prior>), EngineError> {
let (title, metadata) = described(item);
if let Some(link) = link_of(metadata, destination.name()) {
match self.prior(destination, item.level(), &link.native).await? {
Some(held) if origin_of(described(&held.item).1).as_ref() == Some(id) => {
return Ok((
Target::Update {
id: link.native,
found: Found::Link,
},
Some(held),
));
}
Some(_) => {}
None if request.recreate => {}
None => {
return Err(EngineError::StaleLink {
item: id.to_string(),
link: link.to_string(),
});
}
}
}
if let Some(origin) = origin_of(metadata)
&& &origin.source == destination.name()
{
if let Some(held) = self
.prior(destination, item.level(), &origin.native)
.await?
{
return Ok((
Target::Update {
id: origin.native,
found: Found::Origin,
},
Some(held),
));
}
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?
{
let held = self.prior(destination, item.level(), &found).await?;
return Ok((
Target::Update {
id: found,
found: Found::Scan,
},
held,
));
}
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?
{
let held = self.prior(destination, item.level(), &found).await?;
return Ok((
Target::Update {
id: found,
found: Found::Match,
},
held,
));
}
Ok((Target::Create, None))
}
async fn scan(
&self,
destination: &ResolvedSource,
kind: Level,
wanted: &Wanted,
) -> Result<Option<NativeId>, EngineError> {
let mut cursor: Option<Cursor> = None;
let tasks = wanted.task_query(&destination.source().capabilities());
loop {
let asked = cursor.clone();
let request = request_for(destination, cursor);
let next = match kind {
Level::Task => {
let page = destination
.source()
.query_tasks(&tasks, &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>,
pointing: Pointing<'_>,
journal: &mut Journal,
) -> Result<(CopyOutcome, Option<Prior>), EngineError> {
let Pointing { edges, delivers } = pointing;
let (target, found) = match &item.target {
Target::Update { id, found } => (Some(id.clone()), Some(*found)),
Target::Create => (None, None),
};
let reached = |destination: Option<GlobalId>, changed: bool| match (found, destination) {
(Some(via), Some(destination)) if changed => CopyAction::Updated {
destination,
via: via.via(),
link: None,
},
(Some(via), Some(destination)) => CopyAction::Unchanged {
destination,
via: via.via(),
link: None,
},
(_, destination) => CopyAction::Created {
destination,
via: NoCounterpart::Created,
link: None,
},
};
let prior = item.held.clone();
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,
delivers,
destination.name(),
)
{
return Ok((
CopyOutcome {
source: item.source.clone(),
action: reached(Some(qualified(id.clone())), false),
},
prior,
));
}
if request.dry_run {
return Ok((
CopyOutcome {
source: item.source.clone(),
action: reached(target.map(qualified), true),
},
prior,
));
}
let written = qualified(
self.write(
destination,
item,
target,
project,
&edges,
delivers,
prior.clone(),
journal,
)
.await?,
);
Ok((
CopyOutcome {
source: item.source.clone(),
action: reached(Some(written), true),
},
prior,
))
}
async fn filed(
&self,
destination: &ResolvedSource,
item: &Planned,
project: Option<NativeId>,
running: &mut Running,
) -> Result<Option<NativeId>, EngineError> {
match (&item.item, project) {
(Item::Task(task), None) => {
self.counterpart(destination, item, task.project.as_ref(), running)
.await
}
(Item::Document(document), None) => {
self.counterpart(destination, item, document.project.as_ref(), running)
.await
}
(Item::Task(_) | Item::Document(_), filed) => Ok(filed),
(Item::Project(_), _) => Ok(None),
}
}
async fn counterpart(
&self,
destination: &ResolvedSource,
item: &Planned,
project: Option<&NativeId>,
running: &mut Running,
) -> Result<Option<NativeId>, EngineError> {
let Some(project) = project else {
return Ok(None);
};
let qualified = GlobalId::new(item.source.source.clone(), project.clone()).to_string();
if let Some(landed) = running.counterparts.get(&qualified) {
return Ok(Some(landed.clone()));
}
if let Some(looked) = running.filings.get(&qualified) {
return Ok(looked.clone());
}
let found = match self
.linked_project(destination, &item.source.source, project, &qualified)
.await?
{
Some(linked) => Some(linked),
None => {
self.scan(
destination,
Level::Project,
&Wanted::Origin(qualified.clone()),
)
.await?
}
};
let filed = Some(found.unwrap_or_else(|| project.clone()));
running.filings.insert(qualified, filed.clone());
Ok(filed)
}
async fn linked_project(
&self,
destination: &ResolvedSource,
source: &SourceName,
project: &NativeId,
qualified: &str,
) -> Result<Option<NativeId>, EngineError> {
let source = self.readable(source)?;
if !source.source().capabilities().projects.is_native() {
return Ok(None);
}
let Some(held) = source
.source()
.get_project(project)
.await
.map_err(|error| refused(source, error))?
else {
return Ok(None);
};
let Some(link) = link_of(&held.metadata, destination.name()) else {
return Ok(None);
};
let there = destination
.source()
.get_project(&link.native)
.await
.map_err(|error| refused(destination, error))?;
Ok(there
.filter(|there| {
origin_of(&there.metadata).is_some_and(|origin| origin.to_string() == qualified)
})
.map(|_| link.native))
}
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],
delivers: &[TaskRef],
prior: Option<Prior>,
journal: &mut Journal,
) -> Result<NativeId, EngineError> {
let created_kind = item.item.level();
let suggested = target
.clone()
.unwrap_or_else(|| created_id(&item.item, project.as_ref()));
let origin = recorded(item, prior.as_ref());
let landing = outgoing(item, suggested, project, &origin, delivers, prior.as_ref());
if let (Some(id), Some(prior)) = (target.clone(), prior) {
journal.record(Undo::Updated { id, prior });
}
let landed = match landing {
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");
pub(super) async fn readings(sources: &[&ResolvedSource]) -> Vec<Option<Metering>> {
let mut read = Vec::with_capacity(sources.len());
for source in sources {
read.push(source.source().metering().await.ok().flatten());
}
read
}
pub(super) fn spent_between(
before: &[Option<Metering>],
after: &[Option<Metering>],
) -> Option<Spent> {
let mut metered = false;
let mut requests = 0_u64;
let mut budgets: BTreeMap<(String, String), (u64, u64)> = BTreeMap::new();
for (before, after) in before.iter().zip(after) {
let (Some(before), Some(after)) = (before, after) else {
continue;
};
let Some((sent, spent)) = difference(before, after) else {
continue;
};
metered = true;
requests = requests.saturating_add(sent);
for (key, measured, modelled) in spent {
let total = budgets.entry(key).or_default();
total.0 = total.0.saturating_add(measured).saturating_add(modelled);
total.1 = total.1.saturating_add(modelled);
}
}
metered.then(|| Spent {
requests,
budgets: budgets
.into_iter()
.map(|((budget, unit), (amount, modelled))| BudgetSpent {
budget,
unit,
amount,
lower_bound: modelled > 0,
})
.collect(),
})
}
type BudgetDifference = ((String, String), u64, u64);
fn difference(before: &Metering, after: &Metering) -> Option<(u64, Vec<BudgetDifference>)> {
let sent = after.requests.checked_sub(before.requests)?;
for reading in [before, after] {
let mut named = BTreeSet::new();
for budget in &reading.budgets {
if budget.budget.is_empty()
|| budget.unit.is_empty()
|| !named.insert((&budget.budget, &budget.unit))
{
return None;
}
}
}
let held = |from: &Metering, budget: &onetaskgraph_plugin_api::Metered| {
from.budgets
.iter()
.find(|held| held.budget == budget.budget && held.unit == budget.unit)
.cloned()
};
if before
.budgets
.iter()
.any(|budget| held(after, budget).is_none())
{
return None;
}
let mut spent = Vec::with_capacity(after.budgets.len());
for budget in &after.budgets {
let earlier = held(before, budget);
let measured = budget
.measured
.checked_sub(earlier.as_ref().map_or(0, |held| held.measured))?;
let modelled = budget
.modelled
.checked_sub(earlier.as_ref().map_or(0, |held| held.modelled))?;
spent.push((
(budget.budget.clone(), budget.unit.clone()),
measured,
modelled,
));
}
Some((sent, spent))
}
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],
delivers: &[TaskRef],
destination: &SourceName,
) -> bool {
let Some(held) = held else {
return true;
};
let outgoing = outgoing(
item,
target.clone(),
project.clone(),
&recorded(item, Some(held)),
delivers,
Some(held),
);
!same(&held.item, &outgoing, destination) || !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(|_| ()),
}
}
async fn set_link(
source: &ResolvedSource,
level: Level,
id: &NativeId,
value: &Value,
) -> Result<bool, SourceError> {
let key = MetadataKey::copies();
Ok(match level {
Level::Task => source
.source()
.set_task_metadata(id, &key, value)
.await?
.is_some(),
Level::Project => source
.source()
.set_project_metadata(id, &key, value)
.await?
.is_some(),
Level::Document => source
.source()
.set_document_metadata(id, &key, value)
.await?
.is_some(),
})
}
fn metadata_of(item: &mut Item) -> &mut BTreeMap<String, Value> {
match item {
Item::Task(task) => &mut task.metadata,
Item::Project(project) => &mut project.metadata,
Item::Document(document) => &mut document.metadata,
}
}
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 restamped(
metadata: &BTreeMap<String, Value>,
read: &str,
written: &str,
) -> Option<TemplateProvenance> {
let mut provenance = TemplateProvenance::read(metadata).ok().flatten()?;
if provenance.body_digest.as_str() != body_digest(read) {
return None;
}
provenance.body_digest =
Sha256Digest::parse(body_digest(written)).expect("`body_digest` spells a digest");
Some(provenance)
}
fn origin_of(metadata: &BTreeMap<String, Value>) -> Option<GlobalId> {
metadata
.get(GlobalId::ORIGIN_KEY)?
.as_str()?
.parse::<GlobalId>()
.ok()
}
fn link_of(metadata: &BTreeMap<String, Value>, destination: &SourceName) -> Option<GlobalId> {
let linked = metadata
.get(MetadataKey::COPIES_KEY)?
.as_object()?
.get(destination.as_str())?
.as_str()?
.parse::<GlobalId>()
.ok()?;
(&linked.source == destination).then_some(linked)
}
fn well_formed(value: Option<&Value>) -> serde_json::Map<String, Value> {
let Some(Value::Object(held)) = value else {
return serde_json::Map::new();
};
held.iter()
.filter(|(destination, linked)| {
linked
.as_str()
.and_then(|linked| linked.parse::<GlobalId>().ok())
.is_some_and(|linked| linked.source.as_str() == destination.as_str())
})
.map(|(destination, linked)| (destination.clone(), linked.clone()))
.collect()
}
pub(crate) fn malformed_links(value: &Value) -> Option<String> {
let Value::Object(held) = value else {
return Some(format!(
"{} holds an object of destination source names to qualified ids, not {value}",
MetadataKey::COPIES_KEY
));
};
let wrong = held.len() - well_formed(Some(value)).len();
(wrong > 0).then(|| {
format!(
"{} holds an object of destination source names to qualified ids of that source, \
and {wrong} of its entries is not one: {value}",
MetadataKey::COPIES_KEY
)
})
}
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,
delivers: &[TaskRef],
held: Option<&Prior>,
) -> Item {
let own = held.and_then(|held| described(&held.item).1.get(MetadataKey::COPIES_KEY));
let carried = |metadata: &BTreeMap<String, Value>| carried(metadata, origin, own);
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),
delivers: delivers.to_vec(),
delivered_by: match held.map(|held| &held.item) {
Some(Item::Task(held)) => held.delivered_by.clone(),
_ => Vec::new(),
},
..(**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),
..(**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),
..(**document).clone()
})),
}
}
fn created_id(item: &Item, filed: Option<&NativeId>) -> NativeId {
if let (Item::Task(task), Some(filed)) = (item, filed)
&& let Some(own) = &task.project
&& let Some(rest) = task
.id
.as_str()
.strip_prefix(own.as_str())
.and_then(|rest| rest.strip_prefix('/'))
&& !rest.is_empty()
{
return NativeId(format!("{}/{rest}", filed.as_str()));
}
item.id().clone()
}
fn carried(
metadata: &BTreeMap<String, Value>,
origin: &Origin,
own: Option<&Value>,
) -> BTreeMap<String, Value> {
let mut carried = metadata.clone();
carried.remove(Repository::METADATA_KEY);
carried.remove(DependencyEdge::RECORDED_KEY);
carried.remove(TaskRef::DELIVERS_KEY);
carried.remove(TaskRef::DELIVERED_BY_KEY);
carried.remove(GlobalId::ORIGIN_KEY);
carried.remove(MetadataKey::COPIES_KEY);
if let Some(own) = own {
carried.insert(MetadataKey::COPIES_KEY.to_owned(), own.clone());
}
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, destination: &SourceName) -> bool {
match (held, outgoing) {
(Item::Task(held), Item::Task(outgoing)) => {
targets(&held.delivers, destination) == targets(&outgoing.delivers, destination)
&& targets(&held.delivered_by, destination)
== targets(&outgoing.delivered_by, destination)
&& held.title == outgoing.title
&& held.content == outgoing.content
&& held.status == outgoing.status
&& held.priority == outgoing.priority
&& 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() && copied.contains(&far) {
written.get(&far.to_string()).map(|native| native.0.clone())
} else if edge.to.is_qualified() || origin == destination.name() {
Some(edge.to.id().to_owned())
} 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 delivers_of(
item: &Planned,
destination: &ResolvedSource,
copied: &[GlobalId],
written: &BTreeMap<String, NativeId>,
) -> Vec<Option<TaskRef>> {
match &item.item {
Item::Task(task) => mapped_delivers(
&task.delivers,
&item.source.source,
destination,
copied,
written,
),
Item::Project(_) | Item::Document(_) => Vec::new(),
}
}
fn mapped_delivers(
entries: &[TaskRef],
origin: &SourceName,
destination: &ResolvedSource,
copied: &[GlobalId],
written: &BTreeMap<String, NativeId>,
) -> Vec<Option<TaskRef>> {
entries
.iter()
.map(|entry| {
let qualified = entry.in_source(origin);
let Ok(far) = qualified.as_str().parse::<GlobalId>() else {
return Some(qualified);
};
if !copied.contains(&far) {
return Some(qualified);
}
let native = written.get(&far.to_string())?;
Some(if native.as_str().contains(':') {
TaskRef::qualified(destination.name(), native)
} else {
TaskRef::new(native.as_str())
.unwrap_or_else(|_| TaskRef::qualified(destination.name(), native))
})
})
.collect()
}
fn members_named(entries: &[TaskRef], origin: &SourceName, copied: &[GlobalId]) -> u64 {
let named = entries
.iter()
.filter(|entry| {
entry
.in_source(origin)
.as_str()
.parse::<GlobalId>()
.is_ok_and(|far| copied.contains(&far))
})
.count();
u64::try_from(named).unwrap_or(u64::MAX)
}
fn resolved_entries(entries: &[Option<TaskRef>]) -> Vec<TaskRef> {
entries.iter().flatten().cloned().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()
}
fn unrecorded_far_end(
item: &Planned,
destination: &ResolvedSource,
unrecorded: &[GlobalId],
) -> Result<(), EngineError> {
if &item.source.source == destination.name() {
return Ok(());
}
for edge in &item.edges {
if edge.to.is_qualified() {
continue;
}
let far = GlobalId::new(
item.source.source.clone(),
NativeId(edge.to.id().to_owned()),
);
if unrecorded.contains(&far) {
return Err(EngineError::UnrecordedMember {
item: item.source.clone(),
member: far,
destination: destination.name().clone(),
});
}
}
Ok(())
}