use std::sync::Arc;
use async_trait::async_trait;
use crate::core::{CaseId, Digest, Operator, RunId, StoreError, Timestamp};
use crate::export::Selection;
use crate::journal::{Checkpoint, JournalStore};
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Disclosure {
pub id: String,
pub runs: Vec<RunId>,
pub cases: Vec<CaseId>,
pub recipient: String,
pub sealed: bool,
pub checkpoint: Checkpoint,
pub package: Digest,
pub by: Operator,
pub at: Timestamp,
}
impl Disclosure {
#[must_use]
pub fn covers(&self, cases: &[CaseId], runs: &[RunId]) -> bool {
self.cases.iter().any(|c| cases.contains(c)) || self.runs.iter().any(|r| runs.contains(r))
}
#[must_use]
pub fn after_erasure(&self) -> String {
let reach = if self.sealed {
"the copy carried sealed payloads, and those sealed under the erased key become \
unopenable; anything travelling unsealed, and the fields the journal keeps in the \
clear, stay readable in the copy"
} else {
"a plaintext copy this erasure cannot reach"
};
format!(
"disclosed to {} on {} by {} — {reach} (named from the operator's disclosure \
register, an unchained row)",
self.recipient,
self.at,
self.by.actor()
)
}
}
#[async_trait]
pub trait DisclosureRegister: Send + Sync + std::fmt::Debug {
async fn record(&self, act: &Disclosure) -> Result<(), StoreError>;
async fn disclosures(
&self,
cases: &[CaseId],
runs: &[RunId],
) -> Result<Vec<Disclosure>, StoreError>;
}
#[derive(Debug, Clone)]
pub struct Request {
pub selection: Selection,
pub recipient: String,
pub by: Operator,
pub at: Timestamp,
}
#[derive(Debug, thiserror::Error)]
pub enum DiscloseError {
#[error("the package could not be written: {0}")]
Write(#[from] std::io::Error),
#[error("the disclosure could not be recorded, so nothing was delivered: {0}")]
Unrecorded(StoreError),
}
pub async fn disclose(
store: &Arc<dyn JournalStore>,
cases: &Arc<dyn crate::case::CaseStore>,
register: &dyn DisclosureRegister,
request: &Request,
destination: &std::path::Path,
) -> Result<Disclosure, DiscloseError> {
#[allow(clippy::disallowed_methods)]
let id = format!("dsc_{}", ulid::Ulid::generate());
receivable(destination).map_err(DiscloseError::Write)?;
let staging = staging_path(destination, &id);
let result = stage_and_record(store, cases, register, request, &id, &staging).await;
match result {
Ok(act) => match std::fs::rename(&staging, destination) {
Ok(()) => Ok(act),
Err(e) => {
let _ = std::fs::remove_file(&staging);
Err(DiscloseError::Write(e))
}
},
Err(e) => {
let _ = std::fs::remove_file(&staging);
Err(e)
}
}
}
async fn stage_and_record(
store: &Arc<dyn JournalStore>,
cases: &Arc<dyn crate::case::CaseStore>,
register: &dyn DisclosureRegister,
request: &Request,
id: &str,
staging: &std::path::Path,
) -> Result<Disclosure, DiscloseError> {
let file = private_file(staging)?;
let mut out = Digesting::new(std::io::BufWriter::new(file));
let package =
crate::export::package_to_jsonl(store, cases, &request.selection, &mut out).await?;
let (buffered, package_digest) = out.finish()?;
buffered
.into_inner()
.map_err(std::io::IntoInnerError::into_error)?
.sync_all()?;
let act = Disclosure {
id: id.to_owned(),
runs: package.runs,
cases: package.cases,
recipient: request.recipient.clone(),
sealed: package.sealed,
checkpoint: package.checkpoint,
package: package_digest,
by: request.by.clone(),
at: request.at,
};
register
.record(&act)
.await
.map_err(DiscloseError::Unrecorded)?;
Ok(act)
}
fn receivable(destination: &std::path::Path) -> std::io::Result<()> {
use std::io::{Error, ErrorKind};
if destination.is_dir() {
return Err(Error::new(
ErrorKind::InvalidInput,
format!(
"{} is a directory, not a file a package can be written to",
destination.display()
),
));
}
let parent = match destination.parent() {
Some(p) if !p.as_os_str().is_empty() => p,
_ => std::path::Path::new("."),
};
if !parent.is_dir() {
return Err(Error::new(
ErrorKind::NotFound,
format!(
"{} is not a directory a package can be written into",
parent.display()
),
));
}
Ok(())
}
fn staging_path(destination: &std::path::Path, id: &str) -> std::path::PathBuf {
let name = destination
.file_name()
.map_or_else(|| "package".into(), |n| n.to_string_lossy().into_owned());
destination.with_file_name(format!(".{name}.{id}.partial"))
}
fn private_file(path: &std::path::Path) -> std::io::Result<std::fs::File> {
let mut options = std::fs::OpenOptions::new();
options.write(true).create_new(true);
#[cfg(unix)]
std::os::unix::fs::OpenOptionsExt::mode(&mut options, 0o600);
options.open(path)
}
struct Digesting<W: std::io::Write> {
inner: W,
hasher: sha2::Sha256,
}
impl<W: std::io::Write> Digesting<W> {
fn new(inner: W) -> Self {
use sha2::Digest as _;
Self {
inner,
hasher: sha2::Sha256::new(),
}
}
fn finish(mut self) -> std::io::Result<(W, Digest)> {
use sha2::Digest as _;
self.inner.flush()?;
let bytes: [u8; 32] = self.hasher.finalize().into();
Ok((self.inner, Digest::from_bytes(bytes)))
}
}
impl<W: std::io::Write> std::io::Write for Digesting<W> {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
use sha2::Digest as _;
let n = self.inner.write(buf)?;
self.hasher.update(&buf[..n]);
Ok(n)
}
fn flush(&mut self) -> std::io::Result<()> {
self.inner.flush()
}
}