use crate::sync::error::{SyncError, SyncResult, UPDATE_REQUIRED_MESSAGE};
use crate::sync::types::SyncMessage;
use appcore_core::{CoreCompatibilityPolicy, CoreIdentity};
pub const SYNC_WIRE_SCHEMA_V1: &str = "appcore.sync.v1";
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct SyncEnvelopeV1 {
pub schema: String,
pub source_identity: CoreIdentity,
pub message: SyncMessage,
}
#[derive(serde::Serialize)]
struct SyncEnvelopeV1Ref<'a> {
schema: &'static str,
source_identity: &'a CoreIdentity,
message: &'a SyncMessage,
}
impl SyncEnvelopeV1 {
pub fn new(source_identity: CoreIdentity, message: SyncMessage) -> SyncResult<Self> {
validate_source_binding(&source_identity, &message)?;
Ok(Self {
schema: SYNC_WIRE_SCHEMA_V1.to_string(),
source_identity,
message,
})
}
pub fn validate_for(&self, local_identity: &CoreIdentity) -> SyncResult<()> {
if self.schema != SYNC_WIRE_SCHEMA_V1 {
return Err(SyncError::InvalidSyncMessage(
"unsupported sync wire schema",
));
}
if self.source_identity.runtime.node_id != self.message.source_node_id {
return Err(SyncError::InvalidSyncMessage(
"source identity does not match message node",
));
}
let policy = CoreCompatibilityPolicy {
require_same_cluster: true,
required_capability: None,
};
local_identity
.ensure_compatible(&self.source_identity, &policy, &[])
.map_err(|_| SyncError::IncompatiblePeer)
}
}
pub fn encode_sync_envelope_v1(
source_identity: &CoreIdentity,
message: &SyncMessage,
) -> SyncResult<String> {
validate_source_binding(source_identity, message)?;
let envelope = SyncEnvelopeV1Ref {
schema: SYNC_WIRE_SCHEMA_V1,
source_identity,
message,
};
serde_json::to_string(&envelope)
.map_err(|_| SyncError::InvalidSyncMessage("sync wire serialization failed"))
}
fn validate_source_binding(
source_identity: &CoreIdentity,
message: &SyncMessage,
) -> SyncResult<()> {
if source_identity.runtime.node_id != message.source_node_id {
return Err(SyncError::InvalidSyncMessage(
"source identity does not match message node",
));
}
Ok(())
}
pub fn decode_sync_envelope(input: &str) -> SyncResult<SyncEnvelopeV1> {
if input.is_empty() {
return Err(SyncError::EmptyRequestBody);
}
if !input.trim_start().starts_with('{') {
return Err(SyncError::InvalidSyncMessage(UPDATE_REQUIRED_MESSAGE));
}
let envelope = serde_json::from_str::<SyncEnvelopeV1>(input)
.map_err(|_| SyncError::InvalidSyncMessage("invalid sync wire envelope"))?;
if envelope.schema != SYNC_WIRE_SCHEMA_V1 {
return Err(SyncError::InvalidSyncMessage(UPDATE_REQUIRED_MESSAGE));
}
if envelope.source_identity.runtime.node_id != envelope.message.source_node_id {
return Err(SyncError::InvalidSyncMessage(
"source identity does not match message node",
));
}
Ok(envelope)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::sync::{MAX_SYNC_BATCH_PAYLOAD_BYTES, MAX_SYNC_REQUEST_BODY_BYTES};
use appcore_core::{
AppFamily, AppId, ClusterId, CoreId, CoreKind, InstanceId, NodeId, ProtocolVersion,
RuntimeContractVersion, RuntimeIdentity, SyncGroup, TenantId,
};
fn identity(tenant: &str, node: &str) -> CoreIdentity {
CoreIdentity {
tenant_id: TenantId::new(tenant).unwrap(),
cluster_id: ClusterId::new("cluster-a").unwrap(),
core_id: CoreId::new(format!("core-{node}")).unwrap(),
instance_id: InstanceId::new(format!("instance-{node}")).unwrap(),
kind: CoreKind::new("replica").unwrap(),
protocol_version: ProtocolVersion::new(1),
runtime: RuntimeIdentity {
app_id: AppId::new("app-a").unwrap(),
app_family: AppFamily::new("family-a").unwrap(),
sync_group: SyncGroup::new("cluster-a").unwrap(),
runtime_contract: RuntimeContractVersion::new(1),
node_id: NodeId::new(node).unwrap(),
},
}
}
fn message() -> SyncMessage {
SyncMessage {
batch_id: "batch-1".to_string(),
source_node_id: NodeId::new("node-a").unwrap(),
sequence_start: 1,
sequence_end: 1,
event_count: 1,
events_hash: "hash".to_string(),
created_at_ms: 10,
previous_batch_hash: None,
events: vec![b"event".to_vec()],
}
}
#[test]
fn v1_encoding_matches_golden_fixture() {
let encoded = encode_sync_envelope_v1(&identity("tenant-a", "node-a"), &message())
.expect("v1 envelope");
assert_eq!(encoded, include_str!("fixtures/sync-wire-v1.json").trim());
assert!(decode_sync_envelope(&encoded).is_ok());
}
#[test]
fn maximum_raw_batch_fits_the_bounded_http_envelope() {
let message = SyncMessage::new(
"batch-max".to_string(),
NodeId::new("node-a").unwrap(),
1,
1,
10,
None,
vec![vec![u8::MAX; MAX_SYNC_BATCH_PAYLOAD_BYTES]],
);
let encoded = encode_sync_envelope_v1(&identity("tenant-a", "node-a"), &message).unwrap();
assert!(encoded.len() <= MAX_SYNC_REQUEST_BODY_BYTES);
}
#[test]
fn borrowed_v1_encoding_matches_the_owned_contract() {
let source_identity = identity("tenant-a", "node-a");
let mut message = message();
message.events = vec!["Olá 日本語 العربية".as_bytes().to_vec(), vec![0, 1, 255]];
message.event_count = 2;
message.sequence_end = 2;
message.previous_batch_hash = Some("previous-hash".to_string());
let owned = SyncEnvelopeV1::new(source_identity.clone(), message.clone()).unwrap();
assert_eq!(
encode_sync_envelope_v1(&source_identity, &message).unwrap(),
serde_json::to_string(&owned).unwrap()
);
}
#[test]
fn v1_rejects_source_node_mismatch() {
assert!(matches!(
SyncEnvelopeV1::new(identity("tenant-a", "node-b"), message()),
Err(SyncError::InvalidSyncMessage(_))
));
assert!(matches!(
encode_sync_envelope_v1(&identity("tenant-a", "node-b"), &message()),
Err(SyncError::InvalidSyncMessage(_))
));
}
#[test]
fn v1_rejects_incompatible_tenant() {
let envelope = SyncEnvelopeV1::new(identity("tenant-a", "node-a"), message()).unwrap();
assert_eq!(
envelope.validate_for(&identity("tenant-b", "node-b")),
Err(SyncError::IncompatiblePeer)
);
}
#[test]
fn decoder_rejects_unversioned_wire_with_update_wall() {
assert_eq!(
decode_sync_envelope("batch-1\nnode-a\n1\n1\n"),
Err(SyncError::InvalidSyncMessage(UPDATE_REQUIRED_MESSAGE))
);
}
}