use chrono::{DateTime, Utc};
use onetaskgraph_plugin_api::{
Asset, Comment, CommentBody, Cursor, Document, NativeId, NewComment, Page, PageRequest,
SourceError, SourceName, Task, TaskDetailRead, TaskQuery,
};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use super::assets;
use super::fetch::{fits, unrepeated};
use super::{Answer, ConfiguredSource, Engine, EngineError, Qualified, delivery};
use crate::GlobalId;
use crate::plan::{QueryPlan, QueryResponse, SourceFailure};
use crate::resolve::ResolvedSource;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct CommentList {
pub comments: Vec<Comment>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct DeletedComment {
pub deleted: NativeId,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct TaskDetail {
#[serde(flatten)]
pub response: QueryResponse<Qualified<Task>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub comments: Option<Vec<Comment>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub assets: Option<Vec<Asset>>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct DocumentDetail {
#[serde(flatten)]
pub response: QueryResponse<Qualified<Document>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub assets: Option<Vec<Asset>>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
pub struct TaskDetails {
pub details: Vec<TaskDetail>,
}
impl Engine {
pub async fn task_detail(&self, id: &GlobalId) -> Result<TaskDetail, EngineError> {
let name = self.known(&id.source)?;
let mut details = self
.source_details(&name, std::slice::from_ref(id), true)
.await;
Ok(details.remove(0))
}
pub async fn task_details(&self, ids: &[GlobalId], comments: bool) -> TaskDetails {
let mut details: Vec<Option<TaskDetail>> = vec![None; ids.len()];
let mut named: Vec<&SourceName> = Vec::new();
for id in ids {
if !named.contains(&&id.source) {
named.push(&id.source);
}
}
for name in named {
let at: Vec<usize> = (0..ids.len())
.filter(|index| ids[*index].source == *name)
.collect();
let asked: Vec<GlobalId> = at.iter().map(|index| ids[*index].clone()).collect();
let answered = match self.known(name) {
Ok(name) => self.source_details(&name, &asked, comments).await,
Err(refusal) => asked
.iter()
.map(|_| TaskDetail {
response: failed_response(
name,
SourceError::Config {
message: refusal.to_string(),
},
),
comments: None,
assets: None,
})
.collect(),
};
for ((index, id), mut detail) in at.into_iter().zip(&asked).zip(answered) {
if detail.response.items.is_empty() && detail.response.errors.is_empty() {
detail.response.errors.push(SourceFailure {
source: id.source.clone(),
error: SourceError::Refused {
message: no_such_task(id).to_string(),
},
});
}
details[index] = Some(detail);
}
}
TaskDetails {
details: details.into_iter().flatten().collect(),
}
}
async fn source_details(
&self,
name: &SourceName,
ids: &[GlobalId],
comments: bool,
) -> Vec<TaskDetail> {
let mut answer = Answer::new();
let selected = answer.split(self, std::slice::from_ref(name));
let Some(source) = selected.first() else {
return ids
.iter()
.map(|_| TaskDetail {
response: QueryResponse {
items: Vec::new(),
next: None,
plan: QueryPlan::default(),
errors: answer.errors.clone(),
},
comments: None,
assets: None,
})
.collect();
};
let commented = comments && source.source().capabilities().comments.is_native();
let page = PageRequest {
cursor: None,
limit: source.source().capabilities().max_page_size.max(1),
};
let natives: Vec<NativeId> = ids.iter().map(|id| id.native.clone()).collect();
let mut read = source
.source()
.get_task_details(&natives, commented.then_some(&page))
.await;
if read.len() != ids.len() {
let message = format!(
"the source answered {} task details for the {} ids it was asked for",
read.len(),
ids.len()
);
read = ids
.iter()
.map(|_| {
Err(SourceError::Malformed {
message: message.clone(),
})
})
.collect();
}
let mut details = Vec::with_capacity(ids.len());
for (id, result) in ids.iter().zip(read) {
let (found, first) = match result {
Ok(Some(TaskDetailRead { task, comments })) => (Ok(Some(task)), comments),
Ok(None) => (Ok(None), None),
Err(error) => (Err(error), None),
};
let qualified = GlobalId::new(source.name().clone(), id.native.clone());
let mut response = Answer::new().one_response(source, found, |task| {
delivery::qualified_task(qualified, task)
});
let comments = if commented && !response.items.is_empty() {
let walked = match first {
Some(Ok(Some(first))) => walk_from(source, &id.native, first, page.limit).await,
Some(Ok(None)) => Ok(None),
Some(Err(error)) => Err(error),
None => walk(source, &id.native).await,
};
match walked {
Ok(comments) => comments,
Err(error) => {
response.errors.push(SourceFailure {
source: source.name().clone(),
error,
});
None
}
}
} else {
None
};
let assets = match response.items.first() {
Some(task) => match source.source().task_assets(&id.native).await {
Ok(listed) => Some(assets::ordered(listed, task.item.content.as_deref())),
Err(error) => {
response.errors.push(SourceFailure {
source: source.name().clone(),
error,
});
None
}
},
None => None,
};
details.push(TaskDetail {
response,
comments,
assets,
});
}
details
}
async fn assets_of<T>(
&self,
response: &mut QueryResponse<Qualified<T>>,
content: impl Fn(&T) -> Option<&str>,
owner: assets::Owner,
) -> Option<Vec<Asset>> {
let item = response.items.first()?;
let source = self
.ready()
.find(|source| source.name() == &item.id.source)?;
let listed = match owner {
assets::Owner::Document => source.source().document_assets(&item.id.native).await,
assets::Owner::Task => source.source().task_assets(&item.id.native).await,
};
match listed {
Ok(listed) => Some(assets::ordered(listed, content(&item.item))),
Err(error) => {
response.errors.push(SourceFailure {
source: source.name().clone(),
error,
});
None
}
}
}
pub async fn task_without_comments(&self, id: &GlobalId) -> Result<TaskDetail, EngineError> {
let mut response = self.task(id).await?;
let assets = self
.assets_of(
&mut response,
|task: &Task| task.content.as_deref(),
assets::Owner::Task,
)
.await;
Ok(TaskDetail {
response,
comments: None,
assets,
})
}
pub async fn document_detail(&self, id: &GlobalId) -> Result<DocumentDetail, EngineError> {
let mut response = self.document(id).await?;
let assets = self
.assets_of(
&mut response,
|document: &Document| document.content.as_deref(),
assets::Owner::Document,
)
.await;
Ok(DocumentDetail { response, assets })
}
pub async fn comments(&self, task: &GlobalId) -> Result<CommentList, EngineError> {
let source = self.commented(&task.source)?;
match walk(source, &task.native).await {
Ok(Some(comments)) => Ok(CommentList { comments }),
Ok(None) => Err(no_such_task(task)),
Err(error) => Err(failed(source, error)),
}
}
pub async fn add_comment(
&self,
task: &GlobalId,
comment: &NewComment,
) -> Result<Comment, EngineError> {
let source = self.writable_comments(&task.source)?;
match source.source().add_comment(&task.native, comment).await {
Ok(Some(added)) => Ok(added),
Ok(None) => Err(no_such_task(task)),
Err(error) => Err(failed(source, error)),
}
}
pub async fn edit_comment(
&self,
task: &GlobalId,
comment: &NativeId,
body: &CommentBody,
) -> Result<Comment, EngineError> {
let source = self.writable_comments(&task.source)?;
match source
.source()
.edit_comment(&task.native, comment, body)
.await
{
Ok(Some(edited)) => Ok(edited),
Ok(None) => Err(missing(source, task, comment).await),
Err(error) => Err(failed(source, error)),
}
}
pub async fn delete_comment(
&self,
task: &GlobalId,
comment: &NativeId,
) -> Result<DeletedComment, EngineError> {
let source = self.writable_comments(&task.source)?;
match source.source().delete_comment(&task.native, comment).await {
Ok(Some(deleted)) => Ok(DeletedComment { deleted }),
Ok(None) => Err(missing(source, task, comment).await),
Err(error) => Err(failed(source, error)),
}
}
fn configured(&self, name: &SourceName) -> Option<&ConfiguredSource> {
self.sources.iter().find(|source| source.name() == name)
}
fn commented(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
let name = self.known(name)?;
match self.configured(&name) {
Some(ConfiguredSource::Ready(source)) => {
if source.source().capabilities().comments.is_native() {
Ok(source)
} else {
Err(EngineError::NoComments {
name: name.to_string(),
kind: source.kind().to_owned(),
})
}
}
Some(ConfiguredSource::Unavailable(source)) => Err(EngineError::SourceUnavailable {
name: name.to_string(),
error: source.error().clone(),
}),
None => Err(EngineError::NoSources),
}
}
fn writable_comments(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
let source = self.commented(name)?;
if source.source().writes().is_supported() {
return Ok(source);
}
Err(EngineError::CommentsNotWritable {
name: source.name().to_string(),
kind: source.kind().to_owned(),
})
}
}
async fn walk(
source: &ResolvedSource,
task: &NativeId,
) -> Result<Option<Vec<Comment>>, SourceError> {
let limit = source.source().capabilities().max_page_size.max(1);
let request = PageRequest {
cursor: None,
limit,
};
let Some(first) = source.source().task_comments(task, &request).await? else {
return Ok(None);
};
walk_from(source, task, first, limit).await
}
async fn walk_from(
source: &ResolvedSource,
task: &NativeId,
first: Page<Comment>,
limit: u32,
) -> Result<Option<Vec<Comment>>, SourceError> {
let mut comments = Vec::new();
let mut cursor: Option<Cursor> = None;
let mut page = first;
loop {
fits(page.items.len(), limit)?;
unrepeated(
page.next.as_ref(),
cursor.as_ref(),
"walking a task's comments",
)?;
comments.extend(page.items);
let Some(next) = page.next else {
return Ok(Some(comments));
};
cursor = Some(next);
let request = PageRequest {
cursor: cursor.clone(),
limit,
};
let Some(read) = source.source().task_comments(task, &request).await? else {
return Ok(None);
};
page = read;
}
}
pub(super) async fn commented_since(
source: &ResolvedSource,
local: &super::local::LocalTasks,
page: Page<Task>,
since: DateTime<Utc>,
) -> Result<Page<Task>, SourceError> {
if !source.source().capabilities().comments.is_native() {
return Ok(Page {
items: Vec::new(),
next: page.next,
});
}
let query = TaskQuery {
commented_since: Some(since),
..TaskQuery::default()
};
let mut kept = Vec::new();
for task in page.items {
if local.keeps(&task) && any_comment_matches(source, &task.id, &query).await? {
kept.push(task);
}
}
Ok(Page {
items: kept,
next: page.next,
})
}
async fn any_comment_matches(
source: &ResolvedSource,
task: &NativeId,
query: &TaskQuery,
) -> Result<bool, SourceError> {
let limit = source.source().capabilities().max_page_size.max(1);
let mut cursor: Option<Cursor> = None;
loop {
let request = PageRequest {
cursor: cursor.clone(),
limit,
};
let Some(page) = source.source().task_comments(task, &request).await? else {
return Ok(false);
};
fits(page.items.len(), limit)?;
unrepeated(
page.next.as_ref(),
cursor.as_ref(),
"walking a task's comments",
)?;
if query.comments_match(&page.items) {
return Ok(true);
}
match page.next {
Some(next) => cursor = Some(next),
None => return Ok(false),
}
}
}
async fn missing(source: &ResolvedSource, task: &GlobalId, comment: &NativeId) -> EngineError {
match source.source().get_task(&task.native).await {
Ok(Some(_)) => EngineError::NoSuchComment {
task: task.to_string(),
comment: comment.to_string(),
},
Ok(None) => no_such_task(task),
Err(error) => failed(source, error),
}
}
fn failed_response(source: &SourceName, error: SourceError) -> QueryResponse<Qualified<Task>> {
QueryResponse {
items: Vec::new(),
next: None,
plan: QueryPlan::default(),
errors: vec![SourceFailure {
source: source.clone(),
error,
}],
}
}
fn no_such_task(task: &GlobalId) -> EngineError {
EngineError::NoSuchTask {
id: task.to_string(),
}
}
fn failed(source: &ResolvedSource, error: SourceError) -> EngineError {
EngineError::SourceFailed {
name: source.name().to_string(),
error,
}
}