use std::collections::BTreeMap;
use onetaskgraph_plugin_api::{
DependencyEdge, DependencyEndpoint, DependencyKind, Direction, ItemKind, ItemWrite, 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::local::ProjectSelector;
use super::{Engine, EngineError, Filters, Paging, 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,
},
}
#[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>,
}
#[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(NativeId),
Create,
}
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),
}
}
}
struct Planned {
source: GlobalId,
item: Item,
edges: Vec<DependencyEdge>,
target: Target,
}
enum Item {
Task(Box<Task>),
Project(Box<Project>),
}
impl Item {
fn id(&self) -> &NativeId {
match self {
Self::Task(task) => &task.id,
Self::Project(project) => &project.id,
}
}
fn kind(&self) -> ItemKind {
match self {
Self::Task(_) => ItemKind::Task,
Self::Project(_) => ItemKind::Project,
}
}
}
impl Engine {
pub async fn copy(&self, request: &CopyRequest) -> Result<CopyReport, EngineError> {
let destination = self.writable(&request.destination)?;
match request.scope {
CopyScope::Tasks => Ok(CopyReport {
items: self
.copy_items(
destination,
request,
ItemKind::Task,
request.items.as_slice(),
None,
)
.await?,
}),
CopyScope::Projects { tasks } => {
let mut items = Vec::new();
for id in request.items.as_slice() {
items.extend(self.copy_project(destination, request, id, tasks).await?);
}
Ok(CopyReport { items })
}
}
}
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_project(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
id: &GlobalId,
tasks: bool,
) -> Result<Vec<CopyOutcome>, EngineError> {
let members = if tasks {
self.project_members(id).await?
} else {
Vec::new()
};
let project_plan = self
.plan(destination, request, ItemKind::Project, id)
.await?;
let mut copied = vec![id.clone()];
copied.extend(members.iter().cloned());
let mut known = BTreeMap::new();
if let Target::Update(target) = &project_plan.target {
known.insert(id.to_string(), target.clone());
}
for member in &members {
let member_plan = self
.plan(destination, request, ItemKind::Task, member)
.await?;
if let Target::Update(target) = member_plan.target {
known.insert(member.to_string(), target);
}
}
let project_was_unchanged = if let Target::Update(target) = &project_plan.target {
let edges = mapped_edges(
&project_plan.edges,
&id.source,
destination,
&copied,
&known,
);
!edges.iter().any(Option::is_none)
&& !self
.changes(destination, &project_plan, target, &None, &resolved(&edges))
.await?
} else {
false
};
let mut outcomes = self
.copy_items(
destination,
request,
ItemKind::Project,
std::slice::from_ref(id),
None,
)
.await?;
if !tasks {
return Ok(outcomes);
}
let project = outcomes.first().and_then(CopyOutcome::destination).cloned();
let task_outcomes = self
.copy_items(
destination,
request,
ItemKind::Task,
&members,
project.as_ref().map(|project| project.native.clone()),
)
.await?;
outcomes.extend(task_outcomes);
if let Some(project) = project {
if !request.dry_run {
let planned = self
.plan(destination, request, ItemKind::Project, id)
.await?;
let mut copied = vec![id.clone()];
copied.extend(members.iter().cloned());
let written: BTreeMap<String, NativeId> = outcomes
.iter()
.filter_map(|outcome| {
outcome.destination().map(|destination| {
(outcome.source.to_string(), destination.native.clone())
})
})
.collect();
let edges =
mapped_edges(&planned.edges, &id.source, destination, &copied, &written);
self.write(
destination,
&planned,
Some(project.native.clone()),
self.filed(destination, &planned, None).await?,
&resolved(&edges),
)
.await?;
}
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> {
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();
loop {
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(),
});
}
members.extend(response.items.into_iter().map(|task| task.id));
match response.next {
Some(token) => request.paging.token = Some(token),
None => return Ok(members),
}
}
}
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 = None;
loop {
let page: Page<Task> = destination
.source()
.query_tasks(&TaskQuery::default(), &request_for(destination, cursor))
.await
.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()),
},
});
}
match page.next {
Some(next) => cursor = Some(next),
None => return Ok(orphans),
}
}
}
async fn copy_items(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
kind: ItemKind,
items: &[GlobalId],
project: Option<NativeId>,
) -> Result<Vec<CopyOutcome>, EngineError> {
let mut planned = Vec::new();
for id in items {
planned.push(self.plan(destination, request, kind, id).await?);
}
let copied: Vec<GlobalId> = planned.iter().map(|item| item.source.clone()).collect();
let mut written: BTreeMap<String, NativeId> = BTreeMap::new();
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 deferred = 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) {
deferred.push(index);
}
let outcome = self
.land(destination, request, item, filed[index].clone(), &edges)
.await?;
if let Some(id) = outcome.destination() {
written.insert(item.source.to_string(), id.native.clone());
}
outcomes.push(outcome);
}
if request.dry_run {
return Ok(outcomes);
}
for index in deferred {
let item = &planned[index];
let Some(id) = outcomes[index].destination().cloned() else {
continue;
};
let edges = mapped_edges(
&item.edges,
&item.source.source,
destination,
&copied,
&written,
);
self.write(
destination,
item,
Some(id.native),
filed[index].clone(),
&resolved(&edges),
)
.await?;
}
Ok(outcomes)
}
async fn plan(
&self,
destination: &ResolvedSource,
request: &CopyRequest,
kind: ItemKind,
id: &GlobalId,
) -> Result<Planned, EngineError> {
let source = self.readable(&id.source)?;
let item = match kind {
ItemKind::Task => source
.source()
.get_task(&id.native)
.await
.map_err(|error| refused(source, error))?
.map(|task| Item::Task(Box::new(task))),
ItemKind::Project => source
.source()
.get_project(&id.native)
.await
.map_err(|error| refused(source, error))?
.map(|project| Item::Project(Box::new(project))),
}
.ok_or_else(|| EngineError::NoSuchItem { id: id.to_string() })?;
let edges = forward_edges(source, &id.native, item.kind()).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.kind()).await? {
return Ok(Target::Update(origin.native));
}
if !request.recreate {
return Err(EngineError::StaleOrigin {
item: id.to_string(),
origin: origin.to_string(),
});
}
}
if let Some(found) = self
.scan(destination, item.kind(), &Wanted::Origin(id.to_string()))
.await?
{
return Ok(Target::Update(found));
}
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.kind(), &wanted).await?
{
return Ok(Target::Update(found));
}
Ok(Target::Create)
}
async fn scan(
&self,
destination: &ResolvedSource,
kind: ItemKind,
wanted: &Wanted,
) -> Result<Option<NativeId>, EngineError> {
let mut cursor = None;
loop {
let next = match kind {
ItemKind::Task => {
let page = destination
.source()
.query_tasks(&TaskQuery::default(), &request_for(destination, cursor))
.await
.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
}
ItemKind::Project => {
let page = destination
.source()
.query_projects(&ProjectQuery::default(), &request_for(destination, cursor))
.await
.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
}
};
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>],
) -> Result<CopyOutcome, EngineError> {
let target = match &item.target {
Target::Update(id) => Some(id.clone()),
Target::Create => None,
};
let edges = resolved(edges);
let qualified = |native: NativeId| GlobalId::new(destination.name().clone(), native);
if let Some(id) = &target
&& !self
.changes(destination, item, id, &project, &edges)
.await?
{
return Ok(CopyOutcome {
source: item.source.clone(),
action: CopyAction::Unchanged {
destination: qualified(id.clone()),
},
});
}
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 },
},
});
}
let updating = target.is_some();
let written = qualified(
self.write(destination, item, target, project, &edges)
.await?,
);
Ok(CopyOutcome {
source: item.source.clone(),
action: if updating {
CopyAction::Updated {
destination: written,
}
} else {
CopyAction::Created {
destination: Some(written),
}
},
})
}
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).await,
(Item::Task(_), filed) => Ok(filed),
(Item::Project(_), _) => Ok(None),
}
}
async fn counterpart(
&self,
destination: &ResolvedSource,
item: &Planned,
task: &Task,
) -> Result<Option<NativeId>, EngineError> {
let Some(project) = &task.project else {
return Ok(None);
};
let qualified = GlobalId::new(item.source.source.clone(), project.clone());
let found = self
.scan(
destination,
ItemKind::Project,
&Wanted::Origin(qualified.to_string()),
)
.await?;
Ok(Some(found.unwrap_or_else(|| project.clone())))
}
async fn changes(
&self,
destination: &ResolvedSource,
item: &Planned,
target: &NativeId,
project: &Option<NativeId>,
edges: &[DependencyEdge],
) -> Result<bool, EngineError> {
let held = match &item.item {
Item::Task(_) => destination
.source()
.get_task(target)
.await
.map_err(|error| refused(destination, error))?
.map(|task| Item::Task(Box::new(task))),
Item::Project(_) => destination
.source()
.get_project(target)
.await
.map_err(|error| refused(destination, error))?
.map(|project| Item::Project(Box::new(project))),
};
let Some(held) = held else {
return Ok(true);
};
let outgoing = outgoing(item, target.clone(), project.clone());
if !same(&held, &outgoing) {
return Ok(true);
}
let at_destination = forward_edges(destination, target, item.item.kind()).await?;
Ok(!same_edges(&at_destination, edges))
}
async fn write(
&self,
destination: &ResolvedSource,
item: &Planned,
target: Option<NativeId>,
project: Option<NativeId>,
edges: &[DependencyEdge],
) -> Result<NativeId, EngineError> {
let suggested = target.clone().unwrap_or_else(|| item.item.id().clone());
match outgoing(item, suggested, project) {
Item::Task(task) => destination
.source()
.write_task(&ItemWrite {
target,
item: *task,
depends_on: edges.to_vec(),
})
.await
.map_err(|error| refused(destination, error)),
Item::Project(project) => destination
.source()
.write_project(&ItemWrite {
target,
item: *project,
depends_on: edges.to_vec(),
})
.await
.map_err(|error| refused(destination, error)),
}
}
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<onetaskgraph_plugin_api::Cursor>,
) -> PageRequest {
PageRequest {
cursor,
limit: source.source().capabilities().max_page_size.max(1),
}
}
fn refused(source: &ResolvedSource, error: SourceError) -> EngineError {
EngineError::SourceRefused {
name: source.name().to_string(),
error,
}
}
async fn forward_edges(
source: &ResolvedSource,
id: &NativeId,
kind: ItemKind,
) -> Result<Vec<DependencyEdge>, EngineError> {
let mut edges = Vec::new();
let mut cursor = None;
loop {
let page = match kind {
ItemKind::Task => {
source
.source()
.task_dependencies(id, Direction::DependsOn, &request_for(source, cursor))
.await
}
ItemKind::Project => {
source
.source()
.project_dependencies(id, Direction::DependsOn, &request_for(source, cursor))
.await
}
}
.map_err(|error| refused(source, error))?;
edges.extend(page.items);
match page.next {
Some(next) => cursor = Some(next),
None => return Ok(edges),
}
}
}
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),
}
}
fn outgoing(item: &Planned, id: NativeId, project: Option<NativeId>) -> Item {
let origin = item.source.to_string();
match &item.item {
Item::Task(task) => Item::Task(Box::new(Task {
id,
url: 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,
created_at: None,
updated_at: None,
metadata: carried(&project.metadata, &origin),
..(**project).clone()
})),
}
}
fn carried(metadata: &BTreeMap<String, Value>, origin: &str) -> BTreeMap<String, Value> {
let mut carried = metadata.clone();
carried.remove(Repository::METADATA_KEY);
carried.remove(DependencyEdge::RECORDED_KEY);
carried.insert(
GlobalId::ORIGIN_KEY.to_owned(),
Value::String(origin.to_owned()),
);
carried
}
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
}
_ => 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: ItemKind,
) -> Result<bool, EngineError> {
let found = match kind {
ItemKind::Task => destination
.source()
.get_task(id)
.await
.map_err(|error| refused(destination, error))?
.is_some(),
ItemKind::Project => destination
.source()
.get_project(id)
.await
.map_err(|error| refused(destination, error))?
.is_some(),
};
Ok(found)
}