use crate::MwsMessage;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct HostDeliveryRequest {
pub delivery_id: String,
pub expires_at_unix_ms: i64,
pub destination: HostDeliveryDestination,
}
impl HostDeliveryRequest {
pub fn remaining(&self) -> std::time::Duration {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis())
.unwrap_or(u128::MAX);
let remaining = (self.expires_at_unix_ms.max(0) as u128).saturating_sub(now);
std::time::Duration::from_millis(remaining.min(u64::MAX as u128) as u64)
}
pub fn is_expired(&self) -> bool {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis())
.unwrap_or(u128::MAX);
self.expires_at_unix_ms <= 0 || now >= self.expires_at_unix_ms as u128
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(
tag = "kind",
rename_all = "camelCase",
rename_all_fields = "camelCase",
deny_unknown_fields
)]
pub enum HostDeliveryDestination {
AppConnection {
client_id: String,
message: MwsMessage,
},
NodeConnection {
connection_key: String,
message: MwsMessage,
},
Messaging {
request: crate::MessagingDeliveryRequest,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "status", rename_all = "camelCase", deny_unknown_fields)]
pub enum DeliveryOutcome {
Accepted {
boundary: DeliveryAcceptance,
},
Skipped {
reason: String,
},
Failed {
code: String,
message: String,
},
Unknown {
message: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub enum DeliveryAcceptance {
ConnectionQueue,
Provider,
}
impl DeliveryOutcome {
pub fn is_accepted(&self) -> bool {
matches!(self, Self::Accepted { .. })
}
pub fn is_handled(&self) -> bool {
matches!(self, Self::Accepted { .. } | Self::Skipped { .. })
}
pub fn failed(code: impl Into<String>, message: impl Into<String>) -> Self {
Self::Failed {
code: code.into(),
message: message.into(),
}
}
pub fn connection(accepted: bool) -> Self {
if accepted {
Self::Accepted {
boundary: DeliveryAcceptance::ConnectionQueue,
}
} else {
Self::Failed {
code: "UNREACHABLE".into(),
message: "Connection did not accept delivery".into(),
}
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(
tag = "kind",
rename_all = "camelCase",
rename_all_fields = "camelCase",
deny_unknown_fields
)]
pub enum HostConnectionControl {
DisconnectNode { connection_key: String },
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn acceptance_is_not_a_boolean_receipt() {
let result = DeliveryOutcome::connection(true);
assert_eq!(
serde_json::to_value(result).unwrap(),
serde_json::json!({
"status": "accepted", "boundary": "connectionQueue"
})
);
assert!(
serde_json::from_value::<DeliveryOutcome>(serde_json::json!({"delivered": true}))
.is_err()
);
}
#[test]
fn skipped_is_handled_but_never_accepted() {
let outcome = DeliveryOutcome::Skipped {
reason: "No provider message".into(),
};
assert!(outcome.is_handled());
assert!(!outcome.is_accepted());
let response = crate::MessagingDeliveryResponse {
outcome: outcome.clone(),
};
let json = serde_json::to_value(response).unwrap();
let decoded: crate::MessagingDeliveryResponse = serde_json::from_value(json).unwrap();
assert_eq!(decoded.outcome, outcome);
assert!(serde_json::from_value::<crate::MessagingDeliveryResponse>(
serde_json::json!({"delivered": true})
)
.is_err());
}
#[test]
fn obsolete_output_effects_are_rejected() {
assert!(serde_json::from_value::<crate::ServiceCoreOutput>(
serde_json::json!({"effects": []})
)
.is_err());
}
#[test]
fn unknown_is_not_failure_or_acceptance() {
let result = DeliveryOutcome::Unknown {
message: "Timed out after dispatch".into(),
};
assert!(!result.is_accepted());
let encoded = serde_json::to_value(&result).unwrap();
assert_eq!(
serde_json::from_value::<DeliveryOutcome>(encoded).unwrap(),
result
);
}
}