use serde::Serialize;
use sha2::{Digest, Sha256};
use crate::projection::{
ProjectionEventSelector, ProjectionExpression, ProjectionPartition, ProjectionProgramId,
ProjectionValueType,
};
use super::bind::{MutationEventBinding, MutationInputBinding};
use super::canonical::canonical_json_bytes;
use super::program::{MutationProgram, MutationProgramId};
use super::MutationProgramError;
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum MutationHandlerPlacement {
EventualLocal,
EventualRemote,
Direct,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
pub struct MutationHandlerRegistration {
owner: String,
epoch: String,
placement: MutationHandlerPlacement,
#[serde(skip)]
partition: ProjectionPartition,
#[serde(skip)]
binding: MutationEventBinding,
name: String,
version: u64,
}
impl MutationHandlerRegistration {
pub fn try_new(
name: impl Into<String>,
version: u64,
owner: impl Into<String>,
epoch: impl Into<String>,
placement: MutationHandlerPlacement,
partition: ProjectionPartition,
binding: MutationEventBinding,
) -> Result<Self, MutationProgramError> {
let name = super::expression::non_empty(name.into(), "handler name")?;
let owner = super::expression::non_empty(owner.into(), "handler owner")?;
let epoch = super::expression::non_empty(epoch.into(), "handler epoch")?;
if version == 0 {
return Err(MutationProgramError::ZeroVersion("handler version"));
}
Ok(Self {
owner,
epoch,
placement,
partition,
binding,
name,
version,
})
}
pub fn name(&self) -> &str {
&self.name
}
pub fn version(&self) -> u64 {
self.version
}
pub fn owner(&self) -> &str {
&self.owner
}
pub fn epoch(&self) -> &str {
&self.epoch
}
pub fn placement(&self) -> MutationHandlerPlacement {
self.placement
}
pub fn partition(&self) -> &ProjectionPartition {
&self.partition
}
pub fn binding(&self) -> &MutationEventBinding {
&self.binding
}
pub fn uniqueness_key(&self) -> MutationHandlerUniquenessKey {
MutationHandlerUniquenessKey {
owner: self.owner.clone(),
event_name: self.binding.selector().event_name().to_owned(),
event_version: self.binding.selector().event_version(),
body_fingerprint: self.binding.selector().body_fingerprint().to_owned(),
epoch: self.epoch.clone(),
}
}
pub fn target_models(&self) -> Vec<String> {
let mut models = self
.binding
.program()
.operations()
.iter()
.map(|operation| operation.target().model().to_owned())
.collect::<Vec<_>>();
models.sort();
models.dedup();
models
}
pub fn digest(&self) -> Result<[u8; 32], MutationProgramError> {
#[derive(Serialize)]
struct DigestBody<'a> {
name: &'a str,
version: u64,
owner: &'a str,
epoch: &'a str,
placement: MutationHandlerPlacement,
mutation_program_id: String,
selector_event: &'a str,
selector_version: u64,
selector_fingerprint: &'a str,
}
let program_id = self.binding.program().id()?;
let body = DigestBody {
name: &self.name,
version: self.version,
owner: &self.owner,
epoch: &self.epoch,
placement: self.placement,
mutation_program_id: program_id.to_string(),
selector_event: self.binding.selector().event_name(),
selector_version: self.binding.selector().event_version(),
selector_fingerprint: self.binding.selector().body_fingerprint(),
};
let bytes = canonical_json_bytes(&body)?;
let mut digest = Sha256::new();
digest.update(b"distributed.mutation-handler/v1\0");
digest.update((bytes.len() as u64).to_be_bytes());
digest.update(&bytes);
Ok(digest.finalize().into())
}
pub fn to_projection_program(
&self,
) -> Result<crate::projection::ProjectionProgram, MutationProgramError> {
self.binding.to_projection_program(
self.name.clone(),
self.version,
self.partition.clone(),
format!("{}-arm", self.name),
)
}
}
#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
pub struct MutationHandlerUniquenessKey {
pub owner: String,
pub event_name: String,
pub event_version: u64,
pub body_fingerprint: String,
pub epoch: String,
}
#[derive(Clone, Debug, Default)]
pub struct MutationHandlerCatalog {
registrations: Vec<MutationHandlerRegistration>,
}
impl MutationHandlerCatalog {
pub fn new() -> Self {
Self {
registrations: Vec::new(),
}
}
pub fn register(
&mut self,
registration: MutationHandlerRegistration,
) -> Result<(), MutationProgramError> {
let key = registration.uniqueness_key();
if self
.registrations
.iter()
.any(|existing| existing.uniqueness_key() == key)
{
return Err(MutationProgramError::InvalidOperation {
operation: registration.name().to_owned(),
reason: format!(
"duplicate binding for owner `{}` event `{}` v{} epoch `{}`",
key.owner, key.event_name, key.event_version, key.epoch
),
});
}
for model in registration.target_models() {
for existing in &self.registrations {
if existing.epoch() != registration.epoch() {
continue;
}
if !existing.target_models().iter().any(|item| item == &model) {
continue;
}
let direct_overlap = matches!(
(existing.placement(), registration.placement()),
(
MutationHandlerPlacement::Direct,
MutationHandlerPlacement::Direct
) | (
MutationHandlerPlacement::Direct,
MutationHandlerPlacement::EventualLocal
| MutationHandlerPlacement::EventualRemote
) | (
MutationHandlerPlacement::EventualLocal
| MutationHandlerPlacement::EventualRemote,
MutationHandlerPlacement::Direct
)
);
if direct_overlap
|| (existing.owner() != registration.owner()
&& existing.placement() == registration.placement())
{
return Err(MutationProgramError::InvalidOperation {
operation: registration.name().to_owned(),
reason: format!(
"dual writer for model `{model}` epoch `{}` between `{}` and `{}`",
registration.epoch(),
existing.name(),
registration.name()
),
});
}
}
}
self.registrations.push(registration);
Ok(())
}
pub fn registrations(&self) -> &[MutationHandlerRegistration] {
&self.registrations
}
pub fn for_selector(
&self,
selector: &ProjectionEventSelector,
) -> Vec<&MutationHandlerRegistration> {
self.registrations
.iter()
.filter(|registration| registration.binding().selector() == selector)
.collect()
}
}
#[derive(Clone, Debug)]
pub struct CustomMutationHandler {
pub name: String,
pub owner: String,
pub epoch: String,
pub placement: MutationHandlerPlacement,
pub selector: ProjectionEventSelector,
pub allowed_programs: Vec<MutationProgramId>,
}
impl CustomMutationHandler {
pub fn new(
name: impl Into<String>,
owner: impl Into<String>,
epoch: impl Into<String>,
placement: MutationHandlerPlacement,
selector: ProjectionEventSelector,
allowed_programs: Vec<MutationProgramId>,
) -> Self {
Self {
name: name.into(),
owner: owner.into(),
epoch: epoch.into(),
placement,
selector,
allowed_programs,
}
}
pub fn is_portable(&self) -> bool {
false
}
}
pub fn portable_binding(
selector: ProjectionEventSelector,
program: MutationProgram,
field_pairs: &[(&[&str], &[&str], ProjectionValueType)],
) -> Result<MutationEventBinding, MutationProgramError> {
let inputs = field_pairs
.iter()
.map(|(input, body, value_type)| {
super::bind::body_field_binding(
input.iter().copied(),
body.iter().copied(),
value_type.clone(),
)
})
.collect::<Result<Vec<_>, _>>()?;
MutationEventBinding::try_new(selector, inputs, program)
}
pub fn bindings_from_expressions(
pairs: Vec<(Vec<String>, ProjectionExpression)>,
) -> Result<Vec<MutationInputBinding>, MutationProgramError> {
pairs
.into_iter()
.map(|(path, expression)| MutationInputBinding::try_new(path, expression))
.collect()
}
#[allow(dead_code)]
fn _projection_program_id_type(_: ProjectionProgramId) {}