use crate::error::{BallError, Result};
use crate::negotiation::CommitPolicy;
use crate::participant::{Event, Field, Projection};
use crate::task::Task;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::{BTreeMap, BTreeSet};
pub use crate::plugin::native_fields::{field_wire_name, parse_field};
#[derive(Debug, Clone, Deserialize)]
pub struct DescribeResponse {
#[serde(default, deserialize_with = "lenient_events")]
pub subscriptions: Vec<Event>,
pub projection: ProjectionWire,
#[serde(default)]
pub retry_budget: Option<usize>,
#[serde(default)]
pub wants_context: bool,
}
pub const EVENT_CTX_SCHEMA_VERSION: u32 = 1;
#[derive(Debug, Clone, Serialize)]
pub struct EventCtxWire {
pub schema_version: u32,
pub event: String,
pub actor: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub repo: Option<String>,
pub overrides: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub task_before: Option<Value>,
#[serde(skip_serializing_if = "Option::is_none")]
pub commit: Option<String>,
}
impl EventCtxWire {
pub fn for_event(
event: &str,
actor: &str,
repo: Option<String>,
task_before: Option<&Task>,
commit: Option<&str>,
overrides: &[String],
) -> Result<String> {
let task_before = task_before
.map(serde_json::to_value)
.transpose()
.map_err(BallError::Json)?;
serde_json::to_string(&Self {
schema_version: EVENT_CTX_SCHEMA_VERSION,
event: event.to_string(),
actor: actor.to_string(),
repo,
overrides: overrides.to_vec(),
task_before,
commit: commit.map(str::to_string),
})
.map_err(BallError::Json)
}
}
fn lenient_events<'de, D>(d: D) -> std::result::Result<Vec<Event>, D::Error>
where
D: serde::Deserializer<'de>,
{
let raw = Vec::<Value>::deserialize(d)?;
Ok(raw
.into_iter()
.filter_map(|v| serde_json::from_value::<Event>(v).ok())
.collect())
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ProjectionWire {
#[serde(default)]
pub owns: Vec<String>,
#[serde(default)]
pub reads: Vec<String>,
#[serde(default)]
pub external_prefixes: Vec<String>,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ProposeResponse {
#[serde(default)]
pub ok: Option<ProposeOk>,
#[serde(default)]
pub conflict: Option<ProposeConflict>,
#[serde(default)]
pub reject: Option<ProposeReject>,
#[serde(flatten)]
pub extra: BTreeMap<String, Value>,
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct ProposeReject {
#[serde(default)]
pub reason: String,
}
#[derive(Debug, Clone, Deserialize)]
pub struct ProposeOk {
#[serde(default)]
pub task: Value,
#[serde(default)]
pub commit_policy: Option<CommitPolicyWire>,
}
#[derive(Debug, Clone, Deserialize)]
pub struct ProposeConflict {
#[serde(default)]
pub fields: Vec<String>,
#[serde(default)]
pub remote_view: Value,
#[serde(default)]
pub hint: Option<String>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(tag = "kind", rename_all = "lowercase")]
pub enum CommitPolicyWire {
Commit {
#[serde(default)]
message: Option<String>,
},
Batch {
tag: String,
},
Suppress,
}
impl CommitPolicyWire {
pub fn into_policy(self) -> CommitPolicy {
match self {
CommitPolicyWire::Commit { message } => CommitPolicy::Commit { message },
CommitPolicyWire::Batch { tag } => CommitPolicy::Batch { tag },
CommitPolicyWire::Suppress => CommitPolicy::Suppress,
}
}
}
impl ProjectionWire {
pub fn into_projection(self) -> Result<Projection> {
let owns: BTreeSet<Field> = self
.owns
.iter()
.map(|s| parse_field(s))
.collect::<Result<_>>()?;
let reads: BTreeSet<Field> = self
.reads
.iter()
.map(|s| parse_field(s))
.collect::<Result<_>>()?;
let external_prefixes: BTreeSet<String> = self.external_prefixes.into_iter().collect();
Ok(Projection {
owns,
reads,
external_prefixes,
})
}
}
#[cfg(test)]
#[path = "native_types_tests.rs"]
mod tests;