use serde::{Deserialize, Serialize};
use std::fmt;
#[derive(Debug, Clone)]
pub enum IpcError {
UnknownMessageType(String),
ActorNotFound(String),
SerializationError(String),
TargetBusy,
ConnectionClosed,
ProtocolError(String),
IoError(String),
Timeout,
RateLimited {
retry_after_ms: u64,
},
ShuttingDown,
UnsupportedProtocolVersion {
received: u8,
min_supported: u8,
max_supported: u8,
},
}
impl fmt::Display for IpcError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::UnknownMessageType(t) => write!(f, "Unknown message type: {t}"),
Self::ActorNotFound(a) => write!(f, "Actor not found: {a}"),
Self::SerializationError(e) => write!(f, "Serialization error: {e}"),
Self::TargetBusy => write!(f, "Target actor inbox is full"),
Self::ConnectionClosed => write!(f, "Connection closed"),
Self::ProtocolError(e) => write!(f, "Protocol error: {e}"),
Self::IoError(e) => write!(f, "I/O error: {e}"),
Self::Timeout => write!(f, "Request timeout"),
Self::RateLimited { retry_after_ms } => {
write!(f, "Rate limit exceeded, retry after {retry_after_ms}ms")
}
Self::ShuttingDown => write!(f, "Server is shutting down"),
Self::UnsupportedProtocolVersion {
received,
min_supported,
max_supported,
} => {
write!(
f,
"Unsupported protocol version {received:#04x}, supported range: {min_supported:#04x}-{max_supported:#04x}"
)
}
}
}
}
impl std::error::Error for IpcError {}
impl From<serde_json::Error> for IpcError {
fn from(err: serde_json::Error) -> Self {
Self::SerializationError(err.to_string())
}
}
impl From<std::io::Error> for IpcError {
fn from(err: std::io::Error) -> Self {
Self::IoError(err.to_string())
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcEnvelope {
pub correlation_id: String,
pub target: String,
pub message_type: String,
pub payload: serde_json::Value,
#[serde(default)]
pub expects_reply: bool,
#[serde(default)]
pub expects_stream: bool,
#[serde(default = "default_response_timeout")]
pub response_timeout_ms: u64,
}
const fn default_response_timeout() -> u64 {
30_000
}
impl IpcEnvelope {
#[must_use]
pub fn new(
target: impl Into<String>,
message_type: impl Into<String>,
payload: serde_json::Value,
) -> Self {
use mti::prelude::*;
Self {
correlation_id: "req".create_type_id::<V7>().to_string(),
target: target.into(),
message_type: message_type.into(),
payload,
expects_reply: false,
expects_stream: false,
response_timeout_ms: default_response_timeout(),
}
}
#[must_use]
pub fn new_request(
target: impl Into<String>,
message_type: impl Into<String>,
payload: serde_json::Value,
) -> Self {
use mti::prelude::*;
Self {
correlation_id: "req".create_type_id::<V7>().to_string(),
target: target.into(),
message_type: message_type.into(),
payload,
expects_reply: true,
expects_stream: false,
response_timeout_ms: default_response_timeout(),
}
}
#[must_use]
pub fn new_request_with_timeout(
target: impl Into<String>,
message_type: impl Into<String>,
payload: serde_json::Value,
timeout_ms: u64,
) -> Self {
use mti::prelude::*;
Self {
correlation_id: "req".create_type_id::<V7>().to_string(),
target: target.into(),
message_type: message_type.into(),
payload,
expects_reply: true,
expects_stream: false,
response_timeout_ms: timeout_ms,
}
}
#[must_use]
pub fn new_stream_request(
target: impl Into<String>,
message_type: impl Into<String>,
payload: serde_json::Value,
) -> Self {
use mti::prelude::*;
Self {
correlation_id: "str".create_type_id::<V7>().to_string(),
target: target.into(),
message_type: message_type.into(),
payload,
expects_reply: false,
expects_stream: true,
response_timeout_ms: default_response_timeout(),
}
}
#[must_use]
pub fn new_stream_request_with_timeout(
target: impl Into<String>,
message_type: impl Into<String>,
payload: serde_json::Value,
timeout_ms: u64,
) -> Self {
use mti::prelude::*;
Self {
correlation_id: "str".create_type_id::<V7>().to_string(),
target: target.into(),
message_type: message_type.into(),
payload,
expects_reply: false,
expects_stream: true,
response_timeout_ms: timeout_ms,
}
}
#[must_use]
pub fn with_correlation_id(
correlation_id: impl Into<String>,
target: impl Into<String>,
message_type: impl Into<String>,
payload: serde_json::Value,
) -> Self {
Self {
correlation_id: correlation_id.into(),
target: target.into(),
message_type: message_type.into(),
payload,
expects_reply: false,
expects_stream: false,
response_timeout_ms: default_response_timeout(),
}
}
#[must_use]
pub fn with_correlation_id_request(
correlation_id: impl Into<String>,
target: impl Into<String>,
message_type: impl Into<String>,
payload: serde_json::Value,
timeout_ms: u64,
) -> Self {
Self {
correlation_id: correlation_id.into(),
target: target.into(),
message_type: message_type.into(),
payload,
expects_reply: true,
expects_stream: false,
response_timeout_ms: timeout_ms,
}
}
#[must_use]
pub const fn expects_reply(&self) -> bool {
self.expects_reply
}
#[must_use]
pub const fn expects_stream(&self) -> bool {
self.expects_stream
}
#[must_use]
pub const fn response_timeout(&self) -> std::time::Duration {
std::time::Duration::from_millis(self.response_timeout_ms)
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcResponse {
pub correlation_id: String,
pub success: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error_code: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub payload: Option<serde_json::Value>,
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcStreamFrame {
pub correlation_id: String,
pub sequence: u32,
#[serde(default)]
pub is_final: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error_code: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub payload: Option<serde_json::Value>,
}
impl IpcStreamFrame {
#[must_use]
pub fn data(
correlation_id: impl Into<String>,
sequence: u32,
payload: serde_json::Value,
) -> Self {
Self {
correlation_id: correlation_id.into(),
sequence,
is_final: false,
error: None,
error_code: None,
payload: Some(payload),
}
}
#[must_use]
pub fn final_frame(
correlation_id: impl Into<String>,
sequence: u32,
payload: Option<serde_json::Value>,
) -> Self {
Self {
correlation_id: correlation_id.into(),
sequence,
is_final: true,
error: None,
error_code: None,
payload,
}
}
#[must_use]
pub fn error(
correlation_id: impl Into<String>,
sequence: u32,
error: impl Into<String>,
) -> Self {
Self {
correlation_id: correlation_id.into(),
sequence,
is_final: true,
error: Some(error.into()),
error_code: Some("STREAM_ERROR".to_string()),
payload: None,
}
}
#[must_use]
pub fn error_with_code(
correlation_id: impl Into<String>,
sequence: u32,
error_code: impl Into<String>,
error: impl Into<String>,
) -> Self {
Self {
correlation_id: correlation_id.into(),
sequence,
is_final: true,
error: Some(error.into()),
error_code: Some(error_code.into()),
payload: None,
}
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcSubscribeRequest {
pub correlation_id: String,
pub message_types: Vec<String>,
}
impl IpcSubscribeRequest {
#[must_use]
#[allow(dead_code)]
pub fn new(message_types: Vec<String>) -> Self {
use mti::prelude::*;
Self {
correlation_id: "sub".create_type_id::<V7>().to_string(),
message_types,
}
}
#[must_use]
#[allow(dead_code)]
pub fn with_correlation_id(
correlation_id: impl Into<String>,
message_types: Vec<String>,
) -> Self {
Self {
correlation_id: correlation_id.into(),
message_types,
}
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcUnsubscribeRequest {
pub correlation_id: String,
pub message_types: Vec<String>,
}
impl IpcUnsubscribeRequest {
#[must_use]
#[allow(dead_code)]
pub fn new(message_types: Vec<String>) -> Self {
use mti::prelude::*;
Self {
correlation_id: "unsub".create_type_id::<V7>().to_string(),
message_types,
}
}
#[must_use]
#[allow(dead_code)]
pub fn unsubscribe_all() -> Self {
use mti::prelude::*;
Self {
correlation_id: "unsub".create_type_id::<V7>().to_string(),
message_types: Vec::new(),
}
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcPushNotification {
pub notification_id: String,
pub message_type: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub source_actor: Option<String>,
pub payload: serde_json::Value,
pub timestamp_ms: u64,
}
impl IpcPushNotification {
#[must_use]
pub fn new(
message_type: impl Into<String>,
source_actor: Option<String>,
payload: serde_json::Value,
) -> Self {
use mti::prelude::*;
use std::time::{SystemTime, UNIX_EPOCH};
let timestamp_ms = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| u64::try_from(d.as_millis()).unwrap_or(u64::MAX));
Self {
notification_id: "push".create_type_id::<V7>().to_string(),
message_type: message_type.into(),
source_actor,
payload,
timestamp_ms,
}
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcSubscriptionResponse {
pub correlation_id: String,
pub success: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
pub subscribed_types: Vec<String>,
}
impl IpcSubscriptionResponse {
#[must_use]
pub fn success(correlation_id: impl Into<String>, subscribed_types: Vec<String>) -> Self {
Self {
correlation_id: correlation_id.into(),
success: true,
error: None,
subscribed_types,
}
}
#[must_use]
pub fn error(correlation_id: impl Into<String>, error: impl Into<String>) -> Self {
Self {
correlation_id: correlation_id.into(),
success: false,
error: Some(error.into()),
subscribed_types: Vec::new(),
}
}
}
impl IpcResponse {
#[must_use]
pub fn success(correlation_id: impl Into<String>, payload: Option<serde_json::Value>) -> Self {
Self {
correlation_id: correlation_id.into(),
success: true,
error: None,
error_code: None,
payload,
}
}
#[must_use]
pub fn error(correlation_id: impl Into<String>, err: &IpcError) -> Self {
let (error_code, error_message) = match err {
IpcError::UnknownMessageType(_) => ("UNKNOWN_MESSAGE_TYPE", err.to_string()),
IpcError::ActorNotFound(_) => ("ACTOR_NOT_FOUND", err.to_string()),
IpcError::SerializationError(_) => ("SERIALIZATION_ERROR", err.to_string()),
IpcError::TargetBusy => ("TARGET_BUSY", err.to_string()),
IpcError::ConnectionClosed => ("CONNECTION_CLOSED", err.to_string()),
IpcError::ProtocolError(_) => ("PROTOCOL_ERROR", err.to_string()),
IpcError::IoError(_) => ("IO_ERROR", err.to_string()),
IpcError::Timeout => ("TIMEOUT", err.to_string()),
IpcError::RateLimited { .. } => ("RATE_LIMITED", err.to_string()),
IpcError::ShuttingDown => ("SHUTTING_DOWN", err.to_string()),
IpcError::UnsupportedProtocolVersion { .. } => {
("UNSUPPORTED_PROTOCOL_VERSION", err.to_string())
}
};
Self {
correlation_id: correlation_id.into(),
success: false,
error: Some(error_message),
error_code: Some(error_code.to_string()),
payload: None,
}
}
#[must_use]
pub fn error_with_message(
correlation_id: impl Into<String>,
error_code: impl Into<String>,
message: impl Into<String>,
) -> Self {
Self {
correlation_id: correlation_id.into(),
success: false,
error: Some(message.into()),
error_code: Some(error_code.into()),
payload: None,
}
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcDiscoverRequest {
pub correlation_id: String,
#[serde(default = "default_true")]
pub include_actors: bool,
#[serde(default = "default_true")]
pub include_message_types: bool,
}
const fn default_true() -> bool {
true
}
impl IpcDiscoverRequest {
#[must_use]
pub fn new() -> Self {
use mti::prelude::*;
Self {
correlation_id: "disc".create_type_id::<V7>().to_string(),
include_actors: true,
include_message_types: true,
}
}
#[must_use]
pub fn actors_only() -> Self {
use mti::prelude::*;
Self {
correlation_id: "disc".create_type_id::<V7>().to_string(),
include_actors: true,
include_message_types: false,
}
}
#[must_use]
pub fn message_types_only() -> Self {
use mti::prelude::*;
Self {
correlation_id: "disc".create_type_id::<V7>().to_string(),
include_actors: false,
include_message_types: true,
}
}
#[must_use]
pub fn with_correlation_id(
correlation_id: impl Into<String>,
include_actors: bool,
include_message_types: bool,
) -> Self {
Self {
correlation_id: correlation_id.into(),
include_actors,
include_message_types,
}
}
}
impl Default for IpcDiscoverRequest {
fn default() -> Self {
Self::new()
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct ActorInfo {
pub name: String,
pub ern: String,
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct ProtocolVersionInfo {
pub current: u8,
pub min_supported: u8,
pub max_supported: u8,
pub description: String,
pub capabilities: ProtocolCapabilities,
}
#[derive(Clone, Debug, Default)]
pub struct ProtocolCapabilities {
capabilities: std::collections::HashSet<crate::common::ipc::protocol::ProtocolCapability>,
}
impl ProtocolCapabilities {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn all() -> Self {
use crate::common::ipc::protocol::ProtocolCapability;
Self::with_capabilities([
ProtocolCapability::MessagePack,
ProtocolCapability::Streaming,
ProtocolCapability::Push,
ProtocolCapability::Discovery,
])
}
#[must_use]
pub fn with_capabilities<I>(iter: I) -> Self
where
I: IntoIterator<Item = crate::common::ipc::protocol::ProtocolCapability>,
{
Self {
capabilities: iter.into_iter().collect(),
}
}
pub fn insert(&mut self, capability: crate::common::ipc::protocol::ProtocolCapability) {
self.capabilities.insert(capability);
}
#[must_use]
pub fn messagepack(&self) -> bool {
use crate::common::ipc::protocol::ProtocolCapability;
self.capabilities.contains(&ProtocolCapability::MessagePack)
}
#[must_use]
pub fn streaming(&self) -> bool {
use crate::common::ipc::protocol::ProtocolCapability;
self.capabilities.contains(&ProtocolCapability::Streaming)
}
#[must_use]
pub fn push(&self) -> bool {
use crate::common::ipc::protocol::ProtocolCapability;
self.capabilities.contains(&ProtocolCapability::Push)
}
#[must_use]
pub fn discovery(&self) -> bool {
use crate::common::ipc::protocol::ProtocolCapability;
self.capabilities.contains(&ProtocolCapability::Discovery)
}
}
impl Serialize for ProtocolCapabilities {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
use serde::ser::SerializeStruct;
let mut state = serializer.serialize_struct("ProtocolCapabilities", 4)?;
state.serialize_field("messagepack", &self.messagepack())?;
state.serialize_field("streaming", &self.streaming())?;
state.serialize_field("push", &self.push())?;
state.serialize_field("discovery", &self.discovery())?;
state.end()
}
}
impl<'de> Deserialize<'de> for ProtocolCapabilities {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
use crate::common::ipc::protocol::ProtocolCapability;
use serde::de::{MapAccess, Visitor};
struct CapabilitiesVisitor;
impl<'de> Visitor<'de> for CapabilitiesVisitor {
type Value = ProtocolCapabilities;
fn expecting(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
formatter.write_str("a map of capability flags")
}
fn visit_map<M>(self, mut map: M) -> Result<Self::Value, M::Error>
where
M: MapAccess<'de>,
{
let mut caps = ProtocolCapabilities::new();
while let Some(key) = map.next_key::<&str>()? {
let enabled: bool = map.next_value()?;
if enabled {
match key {
"messagepack" => caps.insert(ProtocolCapability::MessagePack),
"streaming" => caps.insert(ProtocolCapability::Streaming),
"push" => caps.insert(ProtocolCapability::Push),
"discovery" => caps.insert(ProtocolCapability::Discovery),
_ => {} }
}
}
Ok(caps)
}
}
deserializer.deserialize_map(CapabilitiesVisitor)
}
}
#[derive(Serialize, Deserialize, Clone, Debug)]
pub struct IpcDiscoverResponse {
pub correlation_id: String,
pub success: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub protocol_version: Option<ProtocolVersionInfo>,
#[serde(skip_serializing_if = "Option::is_none")]
pub actors: Option<Vec<ActorInfo>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub message_types: Option<Vec<String>>,
}
impl ProtocolVersionInfo {
#[must_use]
pub fn current() -> Self {
use crate::common::ipc::protocol::{
ProtocolCapability, MAX_SUPPORTED_VERSION, MIN_SUPPORTED_VERSION, PROTOCOL_VERSION,
};
#[cfg(not(feature = "ipc-messagepack"))]
let capabilities = ProtocolCapabilities::with_capabilities([
ProtocolCapability::Streaming,
ProtocolCapability::Push,
ProtocolCapability::Discovery,
]);
#[cfg(feature = "ipc-messagepack")]
let capabilities = ProtocolCapabilities::with_capabilities([
ProtocolCapability::Streaming,
ProtocolCapability::Push,
ProtocolCapability::Discovery,
ProtocolCapability::MessagePack,
]);
Self {
current: PROTOCOL_VERSION,
min_supported: MIN_SUPPORTED_VERSION,
max_supported: MAX_SUPPORTED_VERSION,
description: "v2 (multi-format, streaming, push, discovery)".to_string(),
capabilities,
}
}
}
impl IpcDiscoverResponse {
#[must_use]
pub fn success(
correlation_id: impl Into<String>,
actors: Option<Vec<ActorInfo>>,
message_types: Option<Vec<String>>,
) -> Self {
Self {
correlation_id: correlation_id.into(),
success: true,
error: None,
protocol_version: Some(ProtocolVersionInfo::current()),
actors,
message_types,
}
}
#[must_use]
pub fn success_without_version(
correlation_id: impl Into<String>,
actors: Option<Vec<ActorInfo>>,
message_types: Option<Vec<String>>,
) -> Self {
Self {
correlation_id: correlation_id.into(),
success: true,
error: None,
protocol_version: None,
actors,
message_types,
}
}
#[must_use]
pub fn error(correlation_id: impl Into<String>, error: impl Into<String>) -> Self {
Self {
correlation_id: correlation_id.into(),
success: false,
error: Some(error.into()),
protocol_version: None,
actors: None,
message_types: None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_ipc_envelope_creation() {
let envelope = IpcEnvelope::new(
"price_service",
"PriceUpdate",
serde_json::json!({ "symbol": "AAPL", "price": 150.25 }),
);
assert!(envelope.correlation_id.starts_with("req_"));
assert_eq!(envelope.target, "price_service");
assert_eq!(envelope.message_type, "PriceUpdate");
assert!(!envelope.expects_reply);
assert_eq!(envelope.response_timeout_ms, 30_000);
}
#[test]
fn test_ipc_envelope_request() {
let envelope = IpcEnvelope::new_request(
"price_service",
"GetPrice",
serde_json::json!({ "symbol": "AAPL" }),
);
assert!(envelope.correlation_id.starts_with("req_"));
assert_eq!(envelope.target, "price_service");
assert_eq!(envelope.message_type, "GetPrice");
assert!(envelope.expects_reply);
assert_eq!(envelope.response_timeout_ms, 30_000);
}
#[test]
fn test_ipc_envelope_request_with_timeout() {
let envelope = IpcEnvelope::new_request_with_timeout(
"price_service",
"GetPrice",
serde_json::json!({ "symbol": "AAPL" }),
5000,
);
assert!(envelope.expects_reply);
assert_eq!(envelope.response_timeout_ms, 5000);
assert_eq!(
envelope.response_timeout(),
std::time::Duration::from_secs(5)
);
}
#[test]
fn test_ipc_envelope_serialization() {
let envelope = IpcEnvelope::with_correlation_id(
"test_123",
"actor",
"Ping",
serde_json::json!({ "value": 42 }),
);
let json = serde_json::to_string(&envelope).unwrap();
let deserialized: IpcEnvelope = serde_json::from_str(&json).unwrap();
assert_eq!(deserialized.correlation_id, "test_123");
assert_eq!(deserialized.target, "actor");
assert_eq!(deserialized.message_type, "Ping");
assert!(!deserialized.expects_reply);
}
#[test]
fn test_ipc_envelope_deserialization_defaults() {
let json = r#"{
"correlation_id": "test_123",
"target": "actor",
"message_type": "Ping",
"payload": {}
}"#;
let deserialized: IpcEnvelope = serde_json::from_str(json).unwrap();
assert!(!deserialized.expects_reply);
assert_eq!(deserialized.response_timeout_ms, 30_000);
}
#[test]
fn test_ipc_envelope_deserialization_with_expects_reply() {
let json = r#"{
"correlation_id": "test_123",
"target": "actor",
"message_type": "Query",
"payload": {"q": "test"},
"expects_reply": true,
"response_timeout_ms": 5000
}"#;
let deserialized: IpcEnvelope = serde_json::from_str(json).unwrap();
assert!(deserialized.expects_reply);
assert_eq!(deserialized.response_timeout_ms, 5000);
}
#[test]
fn test_ipc_response_success() {
let response = IpcResponse::success("test_123", Some(serde_json::json!({ "ok": true })));
assert!(response.success);
assert!(response.error.is_none());
assert!(response.payload.is_some());
}
#[test]
fn test_ipc_response_error() {
let err = IpcError::ActorNotFound("test_actor".to_string());
let response = IpcResponse::error("test_123", &err);
assert!(!response.success);
assert_eq!(response.error_code, Some("ACTOR_NOT_FOUND".to_string()));
assert!(response.error.is_some());
assert!(response.payload.is_none());
}
#[test]
fn test_ipc_error_display() {
let err = IpcError::UnknownMessageType("MyMessage".to_string());
assert_eq!(err.to_string(), "Unknown message type: MyMessage");
let err = IpcError::TargetBusy;
assert_eq!(err.to_string(), "Target actor inbox is full");
}
#[test]
fn test_ipc_subscribe_request_new() {
let request =
IpcSubscribeRequest::new(vec!["PriceUpdate".to_string(), "OrderStatus".to_string()]);
assert!(request.correlation_id.starts_with("sub_"));
assert_eq!(request.message_types.len(), 2);
assert!(request.message_types.contains(&"PriceUpdate".to_string()));
assert!(request.message_types.contains(&"OrderStatus".to_string()));
}
#[test]
fn test_ipc_subscribe_request_with_correlation_id() {
let request =
IpcSubscribeRequest::with_correlation_id("custom_123", vec!["MyMessage".to_string()]);
assert_eq!(request.correlation_id, "custom_123");
assert_eq!(request.message_types.len(), 1);
}
#[test]
fn test_ipc_subscribe_request_serialization() {
let request = IpcSubscribeRequest::with_correlation_id(
"test_sub",
vec!["TypeA".to_string(), "TypeB".to_string()],
);
let json = serde_json::to_string(&request).unwrap();
let deserialized: IpcSubscribeRequest = serde_json::from_str(&json).unwrap();
assert_eq!(deserialized.correlation_id, "test_sub");
assert_eq!(deserialized.message_types, request.message_types);
}
#[test]
fn test_ipc_unsubscribe_request_new() {
let request = IpcUnsubscribeRequest::new(vec!["PriceUpdate".to_string()]);
assert!(request.correlation_id.starts_with("unsub_"));
assert_eq!(request.message_types.len(), 1);
}
#[test]
fn test_ipc_unsubscribe_request_unsubscribe_all() {
let request = IpcUnsubscribeRequest::unsubscribe_all();
assert!(request.correlation_id.starts_with("unsub_"));
assert!(request.message_types.is_empty());
}
#[test]
fn test_ipc_unsubscribe_request_serialization() {
let request = IpcUnsubscribeRequest::unsubscribe_all();
let json = serde_json::to_string(&request).unwrap();
let deserialized: IpcUnsubscribeRequest = serde_json::from_str(&json).unwrap();
assert!(deserialized.message_types.is_empty());
}
#[test]
fn test_ipc_push_notification_new() {
let notification = IpcPushNotification::new(
"PriceUpdate",
Some("price_service".to_string()),
serde_json::json!({ "symbol": "AAPL", "price": 150.25 }),
);
assert!(notification.notification_id.starts_with("push_"));
assert_eq!(notification.message_type, "PriceUpdate");
assert_eq!(notification.source_actor, Some("price_service".to_string()));
assert!(notification.timestamp_ms > 0);
}
#[test]
fn test_ipc_push_notification_no_source() {
let notification = IpcPushNotification::new(
"SystemEvent",
None,
serde_json::json!({ "event": "startup" }),
);
assert_eq!(notification.message_type, "SystemEvent");
assert!(notification.source_actor.is_none());
}
#[test]
fn test_ipc_push_notification_serialization() {
let notification = IpcPushNotification::new(
"TestMessage",
Some("test_actor".to_string()),
serde_json::json!({ "data": 123 }),
);
let json = serde_json::to_string(¬ification).unwrap();
let deserialized: IpcPushNotification = serde_json::from_str(&json).unwrap();
assert_eq!(deserialized.message_type, "TestMessage");
assert_eq!(deserialized.source_actor, Some("test_actor".to_string()));
assert_eq!(deserialized.payload["data"], 123);
}
#[test]
fn test_ipc_subscription_response_success() {
let response = IpcSubscriptionResponse::success(
"sub_123",
vec!["TypeA".to_string(), "TypeB".to_string()],
);
assert!(response.success);
assert_eq!(response.correlation_id, "sub_123");
assert!(response.error.is_none());
assert_eq!(response.subscribed_types.len(), 2);
}
#[test]
fn test_ipc_subscription_response_error() {
let response = IpcSubscriptionResponse::error("sub_456", "Connection not registered");
assert!(!response.success);
assert_eq!(response.correlation_id, "sub_456");
assert_eq!(
response.error,
Some("Connection not registered".to_string())
);
assert!(response.subscribed_types.is_empty());
}
#[test]
fn test_ipc_subscription_response_serialization() {
let response =
IpcSubscriptionResponse::success("test_corr", vec!["PriceUpdate".to_string()]);
let json = serde_json::to_string(&response).unwrap();
let deserialized: IpcSubscriptionResponse = serde_json::from_str(&json).unwrap();
assert!(deserialized.success);
assert_eq!(deserialized.correlation_id, "test_corr");
assert_eq!(
deserialized.subscribed_types,
vec!["PriceUpdate".to_string()]
);
}
}