use std::collections::HashSet;
use std::fmt;
use std::path::{Component, Path};
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use crate::system_state::SimulationTime;
use super::SamplingInterval;
use super::error::StorageError;
pub(crate) const FORMAT_NAME: &str = "scientific-workflow-jsonl";
pub(crate) const FORMAT_VERSION: u32 = 3;
pub(crate) const PAYLOAD_ENCODING: &str = "json";
pub(crate) const RECORD_FRAMING: &str = "json_lines";
#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct RecordingMetadata {
pub(crate) format: String,
pub(crate) version: u32,
pub(crate) status: RecordingStatus,
pub(crate) records: RecordFormat,
pub(crate) time: TimeAxisMetadata,
#[serde(default, skip_serializing_if = "Map::is_empty")]
pub(crate) user_metadata: Map<String, Value>,
pub(crate) streams: Vec<StateStreamMetadata>,
}
impl RecordingMetadata {
pub(crate) fn running(
time: TimeAxisMetadata,
user_metadata: Map<String, Value>,
streams: Vec<StateStreamMetadata>,
) -> Self {
Self {
format: FORMAT_NAME.to_owned(),
version: FORMAT_VERSION,
status: RecordingStatus::Running,
records: RecordFormat::json_lines(),
time,
user_metadata,
streams,
}
}
pub(crate) fn validate(&self, path: &Path) -> Result<(), StorageError> {
if self.format != FORMAT_NAME {
return Err(invalid_metadata(
path,
format!("format must be `{FORMAT_NAME}`, got `{}`", self.format),
));
}
if self.version != FORMAT_VERSION {
return Err(StorageError::UnsupportedVersion {
path: path.to_path_buf(),
found: self.version,
supported: FORMAT_VERSION,
});
}
self.records.validate(path)?;
self.time.validate(path)?;
self.status.validate(path)?;
if self.streams.is_empty() {
return Err(invalid_metadata(
path,
"at least one output stream must be declared",
));
}
let mut names = HashSet::with_capacity(self.streams.len());
let mut directories = HashSet::with_capacity(self.streams.len());
for stream in &self.streams {
if !names.insert(stream.name.as_str()) {
return Err(StorageError::DuplicateStateStream {
stream: stream.name.clone(),
});
}
if !directories.insert(stream.directory.as_str()) {
return Err(invalid_metadata(
path,
format!(
"streams use the same output directory `{}`",
stream.directory
),
));
}
stream.validate(path)?;
}
Ok(())
}
pub(crate) fn stream(&self, name: &str) -> Option<&StateStreamMetadata> {
self.streams.iter().find(|stream| stream.name == name)
}
pub(crate) fn stream_mut(&mut self, name: &str) -> Option<&mut StateStreamMetadata> {
self.streams.iter_mut().find(|stream| stream.name == name)
}
}
#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
#[serde(tag = "state", rename_all = "snake_case", deny_unknown_fields)]
pub(crate) enum RecordingStatus {
Running,
Complete,
Failed {
message: String,
},
}
impl RecordingStatus {
fn validate(&self, path: &Path) -> Result<(), StorageError> {
if let Self::Failed { message } = self {
if message.trim().is_empty() {
return Err(invalid_metadata(
path,
"failed run status requires a non-empty message",
));
}
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct RecordFormat {
pub(crate) encoding: String,
pub(crate) framing: String,
}
impl RecordFormat {
fn json_lines() -> Self {
Self {
encoding: PAYLOAD_ENCODING.to_owned(),
framing: RECORD_FRAMING.to_owned(),
}
}
fn validate(&self, path: &Path) -> Result<(), StorageError> {
if self.encoding != PAYLOAD_ENCODING || self.framing != RECORD_FRAMING {
return Err(invalid_metadata(
path,
format!(
"record format must be `{PAYLOAD_ENCODING}` with `{RECORD_FRAMING}` framing"
),
));
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct TimeAxisMetadata {
pub(crate) iteration_name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) iteration_unit: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) physical_time_name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) physical_time_unit: Option<String>,
}
impl TimeAxisMetadata {
fn validate(&self, path: &Path) -> Result<(), StorageError> {
if self.iteration_name.trim().is_empty() {
return Err(invalid_metadata(
path,
"time.iteration_name must not be empty",
));
}
if self
.iteration_unit
.as_deref()
.is_some_and(|unit| unit.trim().is_empty())
{
return Err(invalid_metadata(
path,
"time.iteration_unit must not be empty when present",
));
}
if self
.physical_time_name
.as_deref()
.is_some_and(|name| name.trim().is_empty())
{
return Err(invalid_metadata(
path,
"time.physical_time_name must not be empty when present",
));
}
if self
.physical_time_unit
.as_deref()
.is_some_and(|unit| unit.trim().is_empty())
{
return Err(invalid_metadata(
path,
"time.physical_time_unit must not be empty when present",
));
}
if self.physical_time_unit.is_some() && self.physical_time_name.is_none() {
return Err(invalid_metadata(
path,
"time.physical_time_unit requires time.physical_time_name",
));
}
Ok(())
}
}
#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct StateStreamMetadata {
pub(crate) name: String,
pub(crate) directory: String,
pub(crate) sampling_interval: SamplingInterval,
pub(crate) fields: Vec<StateFieldMetadata>,
pub(crate) max_chunk_bytes: u64,
pub(crate) queue_bytes: u64,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub(crate) chunks: Vec<ChunkMetadata>,
}
impl StateStreamMetadata {
fn validate(&self, path: &Path) -> Result<(), StorageError> {
if self.name.trim().is_empty() {
return Err(invalid_metadata(path, "stream name must not be empty"));
}
validate_relative_path(path, "stream directory", &self.directory)?;
if self.max_chunk_bytes == 0 || self.queue_bytes == 0 {
return Err(invalid_metadata(
path,
format!("stream `{}` has a zero storage limit", self.name),
));
}
let mut fields = HashSet::with_capacity(self.fields.len());
for field in &self.fields {
field.validate(path, &self.name)?;
if !fields.insert(field.name.as_str()) {
return Err(invalid_metadata(
path,
format!(
"stream `{}` declares duplicate field `{}`",
self.name, field.name
),
));
}
}
let mut previous_last = None;
for (expected_ordinal, chunk) in self.chunks.iter().enumerate() {
chunk.validate(path, &self.name, expected_ordinal as u64)?;
if let Some(previous) = previous_last {
if chunk.first_iteration <= previous {
return Err(invalid_metadata(
path,
format!(
"stream `{}` chunk {} begins at iteration {}, not after {}",
self.name, chunk.ordinal, chunk.first_iteration, previous
),
));
}
}
previous_last = Some(chunk.last_iteration);
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct StateFieldMetadata {
pub(crate) name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) description: Option<String>,
}
impl StateFieldMetadata {
fn validate(&self, path: &Path, stream: &str) -> Result<(), StorageError> {
if self.name.trim().is_empty() {
return Err(invalid_metadata(
path,
format!("stream `{stream}` contains an empty field name"),
));
}
if self
.description
.as_deref()
.is_some_and(|description| description.trim().is_empty())
{
return Err(invalid_metadata(
path,
format!(
"stream `{stream}` field `{}` has an empty description",
self.name
),
));
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct ChunkMetadata {
pub(crate) ordinal: u64,
pub(crate) file: String,
pub(crate) records: u64,
pub(crate) bytes: u64,
pub(crate) checksum: String,
pub(crate) first_iteration: u64,
pub(crate) last_iteration: u64,
}
impl ChunkMetadata {
fn validate(
&self,
path: &Path,
stream: &str,
expected_ordinal: u64,
) -> Result<(), StorageError> {
if self.ordinal != expected_ordinal {
return Err(invalid_metadata(
path,
format!(
"stream `{stream}` expected chunk ordinal {expected_ordinal}, got {}",
self.ordinal
),
));
}
let expected_file = chunk_filename(self.ordinal);
if self.file != expected_file {
return Err(invalid_metadata(
path,
format!(
"stream `{stream}` chunk {} filename must be `{expected_file}`",
self.ordinal
),
));
}
validate_relative_path(path, "chunk file", &self.file)?;
if self.records == 0 || self.bytes == 0 {
return Err(invalid_metadata(
path,
format!("stream `{stream}` chunk {} is empty", self.ordinal),
));
}
if self.first_iteration > self.last_iteration {
return Err(invalid_metadata(
path,
format!(
"stream `{stream}` chunk {} iteration range {}..={} is reversed",
self.ordinal, self.first_iteration, self.last_iteration
),
));
}
if !valid_checksum(&self.checksum) {
return Err(invalid_metadata(
path,
format!(
"stream `{stream}` chunk {} has an invalid checksum",
self.ordinal
),
));
}
Ok(())
}
}
pub(crate) struct EncodedStateRecord {
time: SimulationTime,
bytes: Vec<u8>,
}
impl EncodedStateRecord {
pub(crate) fn new(time: SimulationTime, mut json: Vec<u8>) -> Self {
json.push(b'\n');
Self { time, bytes: json }
}
pub(crate) fn simulation_time(&self) -> SimulationTime {
self.time
}
pub(crate) fn len(&self) -> usize {
self.bytes.len()
}
pub(crate) fn bytes(&self) -> &[u8] {
&self.bytes
}
}
impl fmt::Debug for EncodedStateRecord {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("EncodedStateRecord")
.field("time", &self.time)
.field("bytes", &self.bytes.len())
.finish_non_exhaustive()
}
}
pub(crate) fn chunk_filename(ordinal: u64) -> String {
format!("chunk-{ordinal:06}.jsonl")
}
pub(crate) fn chunk_temp_filename(ordinal: u64) -> String {
format!("{}.tmp", chunk_filename(ordinal))
}
fn invalid_metadata(path: &Path, reason: impl Into<String>) -> StorageError {
StorageError::InvalidMetadata {
path: path.to_path_buf(),
reason: reason.into(),
}
}
fn validate_relative_path(
metadata_path: &Path,
label: &str,
value: &str,
) -> Result<(), StorageError> {
let path = Path::new(value);
if value.is_empty()
|| path.is_absolute()
|| !path
.components()
.all(|component| matches!(component, Component::Normal(_)))
{
return Err(invalid_metadata(
metadata_path,
format!("{label} `{value}` must be a safe relative path"),
));
}
Ok(())
}
fn valid_checksum(checksum: &str) -> bool {
let Some((algorithm, digest)) = checksum.split_once(':') else {
return false;
};
!algorithm.is_empty()
&& algorithm
.bytes()
.all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit())
&& !digest.is_empty()
&& digest
.bytes()
.all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
}