use std::sync::Arc;
use async_trait::async_trait;
use crate::case::{ClaimError, TaskStore};
use crate::core::{CaseId, StoreError, Task, TaskId, TaskState, TenantId, Timestamp, Withheld};
use crate::journal::payload;
use super::KeyRing;
#[derive(Debug)]
pub struct SealedTasks {
inner: Arc<dyn TaskStore>,
keys: Arc<dyn KeyRing>,
tenant: TenantId,
}
impl SealedTasks {
#[must_use]
pub fn wrap(inner: Arc<dyn TaskStore>, keys: Arc<dyn KeyRing>, tenant: TenantId) -> Arc<Self> {
super::assert_serves(inner.tenant(), &tenant, "task");
Arc::new(Self {
inner,
keys,
tenant,
})
}
fn scope_for(&self, task: &Task) -> String {
task.case.map_or_else(
|| super::scope(&self.tenant, &task.run.to_string()),
|c| super::scope(&self.tenant, &c.to_string()),
)
}
pub(super) fn aad(tenant: &TenantId, id: TaskId) -> String {
format!("task:{tenant}:{id}")
}
async fn opened(&self, mut task: Task) -> Result<Task, StoreError> {
if task.withheld != Some(Withheld::Sealed) {
return Ok(task);
}
let aad = Self::aad(&self.tenant, task.id);
let open = async |envelope: Option<Vec<u8>>| -> Result<Opening, StoreError> {
let Some(envelope) = envelope else {
return Ok(Opening::Undecodable);
};
Ok(
match super::envelope::open_or_erased(self.keys.as_ref(), aad.as_bytes(), &envelope)
.await
.map_err(|e| StoreError::Backend(e.to_string()))?
{
Some(plain) => Opening::Clear(plain),
None => Opening::Erased,
},
)
};
let mut withheld = None;
match open(payload::unwrap(&task.justification.proposed_action)).await? {
Opening::Clear(plain) => match serde_json::from_slice(&plain) {
Ok(value) => task.justification.proposed_action = value,
Err(_) => withheld = Some(worse(withheld, Withheld::Undecodable)),
},
Opening::Erased => withheld = Some(worse(withheld, Withheld::Erased)),
Opening::Undecodable => withheld = Some(worse(withheld, Withheld::Undecodable)),
}
for item in &mut task.justification.evidence {
match open(payload::unwrap_text(item.peek())).await? {
Opening::Clear(plain) => match String::from_utf8(plain) {
Ok(text) => *item = item.clone().map(|_| text),
Err(_) => withheld = Some(worse(withheld, Withheld::Undecodable)),
},
Opening::Erased => withheld = Some(worse(withheld, Withheld::Erased)),
Opening::Undecodable => withheld = Some(worse(withheld, Withheld::Undecodable)),
}
}
task.withheld = withheld;
Ok(task)
}
async fn opened_all(&self, tasks: Vec<Task>) -> Result<Vec<Task>, StoreError> {
let mut out = Vec::with_capacity(tasks.len());
for task in tasks {
out.push(self.opened(task).await?);
}
Ok(out)
}
}
#[async_trait]
impl TaskStore for SealedTasks {
fn tenant(&self) -> &str {
self.tenant.as_str()
}
async fn open(&self, task: &Task) -> Result<Task, StoreError> {
let plain = crate::core::canon::to_bytes(&task.justification.proposed_action)
.map_err(|e| StoreError::Backend(format!("a proposal would not serialise: {e}")))?;
let envelope = super::envelope::seal(
self.keys.as_ref(),
&self.scope_for(task),
Self::aad(&self.tenant, task.id).as_bytes(),
&plain,
)
.await
.map_err(|e| StoreError::Backend(format!("sealing a proposal failed: {e}")))?;
let mut sealed = task.clone();
sealed.justification.proposed_action = payload::wrap(&envelope);
sealed.withheld = Some(Withheld::Sealed);
for item in &mut sealed.justification.evidence {
let wrapped = super::envelope::seal(
self.keys.as_ref(),
&self.scope_for(task),
Self::aad(&self.tenant, task.id).as_bytes(),
item.peek().as_bytes(),
)
.await
.map_err(|e| StoreError::Backend(format!("sealing a task's evidence failed: {e}")))?;
*item = item.clone().map(|_| payload::wrap_text(&wrapped));
}
let written = self.inner.open(&sealed).await?;
self.opened(written).await
}
async fn task(&self, id: TaskId) -> Result<Option<Task>, StoreError> {
let found = self.inner.task(id).await?;
Ok(match found {
Some(task) => Some(self.opened(task).await?),
None => None,
})
}
async fn claim(&self, id: TaskId, actor: &str, roles: &[String]) -> Result<Task, ClaimError> {
let claimed = self.inner.claim(id, actor, roles).await?;
Ok(self.opened(claimed).await?)
}
async fn take_over(
&self,
id: TaskId,
from: &str,
actor: &str,
roles: &[String],
) -> Result<Task, ClaimError> {
let taken = self.inner.take_over(id, from, actor, roles).await?;
Ok(self.opened(taken).await?)
}
async fn release(&self, id: TaskId, actor: &str) -> Result<(), ClaimError> {
self.inner.release(id, actor).await
}
async fn set_state(&self, id: TaskId, state: TaskState) -> Result<bool, StoreError> {
self.inner.set_state(id, state).await
}
async fn withdraw_run(
&self,
run: crate::core::RunId,
awaited: &[TaskId],
) -> Result<usize, StoreError> {
self.inner.withdraw_run(run, awaited).await
}
async fn escalate(&self, id: TaskId) -> Result<Task, StoreError> {
let escalated = self.inner.escalate(id).await?;
self.opened(escalated).await
}
async fn queue(&self, roles: &[String], limit: usize) -> Result<Vec<Task>, StoreError> {
let tasks = self.inner.queue(roles, limit).await?;
self.opened_all(tasks).await
}
async fn for_case(&self, case: CaseId) -> Result<Vec<Task>, StoreError> {
let tasks = self.inner.for_case(case).await?;
self.opened_all(tasks).await
}
async fn open_count(&self) -> Result<u64, StoreError> {
self.inner.open_count().await
}
async fn overdue(&self, now: Timestamp, limit: usize) -> Result<Vec<Task>, StoreError> {
let tasks = self.inner.overdue(now, limit).await?;
self.opened_all(tasks).await
}
}
enum Opening {
Clear(Vec<u8>),
Erased,
Undecodable,
}
const fn worse(so_far: Option<Withheld>, this: Withheld) -> Withheld {
match (so_far, this) {
(Some(Withheld::Undecodable), _) | (_, Withheld::Undecodable) => Withheld::Undecodable,
_ => this,
}
}