use uuid::Uuid;
use crate::frame::frame_errors::{
ClientRoutesChangeEventParseError, ClusterChangeEventParseError, CqlEventParseError,
SchemaChangeEventParseError,
};
use crate::frame::server_event_type::{EventType, EventTypeV2};
use crate::frame::types;
use std::net::SocketAddr;
#[derive(Debug)]
#[expect(clippy::enum_variant_names)]
pub enum Event {
TopologyChange(TopologyChangeEvent),
StatusChange(StatusChangeEvent),
SchemaChange(SchemaChangeEvent),
}
#[derive(Debug)]
#[expect(clippy::enum_variant_names)]
#[non_exhaustive]
pub enum EventV2 {
TopologyChange(TopologyChangeEvent),
StatusChange(StatusChangeEvent),
SchemaChange(SchemaChangeEvent),
ClientRoutesChange(ClientRoutesChangeEvent),
}
#[derive(Debug)]
pub enum TopologyChangeEvent {
NewNode(SocketAddr),
RemovedNode(SocketAddr),
}
#[derive(Debug)]
pub enum StatusChangeEvent {
Up(SocketAddr),
Down(SocketAddr),
}
#[derive(Debug)]
#[expect(clippy::enum_variant_names)]
pub enum SchemaChangeEvent {
KeyspaceChange {
change_type: SchemaChangeType,
keyspace_name: String,
},
TableChange {
change_type: SchemaChangeType,
keyspace_name: String,
object_name: String,
},
TypeChange {
change_type: SchemaChangeType,
keyspace_name: String,
type_name: String,
},
FunctionChange {
change_type: SchemaChangeType,
keyspace_name: String,
function_name: String,
arguments: Vec<String>,
},
AggregateChange {
change_type: SchemaChangeType,
keyspace_name: String,
aggregate_name: String,
arguments: Vec<String>,
},
}
#[derive(Debug)]
pub enum SchemaChangeType {
Created,
Updated,
Dropped,
Invalid,
}
#[derive(Debug)]
#[non_exhaustive]
pub enum ClientRoutesChangeEvent {
UpdateNodes {
connection_ids: Vec<String>,
host_ids: Vec<Uuid>,
},
}
impl Event {
pub fn deserialize(buf: &mut &[u8]) -> Result<Self, CqlEventParseError> {
let event_type: EventType = types::read_string(buf)
.map_err(CqlEventParseError::EventTypeParseError)?
.parse()?;
match event_type {
EventType::TopologyChange => Ok(Self::TopologyChange(
TopologyChangeEvent::deserialize(buf)
.map_err(CqlEventParseError::TopologyChangeEventParseError)?,
)),
EventType::StatusChange => Ok(Self::StatusChange(
StatusChangeEvent::deserialize(buf)
.map_err(CqlEventParseError::StatusChangeEventParseError)?,
)),
EventType::SchemaChange => Ok(Self::SchemaChange(SchemaChangeEvent::deserialize(buf)?)),
}
}
}
impl EventV2 {
pub fn deserialize(buf: &mut &[u8]) -> Result<Self, CqlEventParseError> {
let event_type: EventTypeV2 = types::read_string(buf)
.map_err(CqlEventParseError::EventTypeParseError)?
.parse()?;
match event_type {
EventTypeV2::TopologyChange => Ok(Self::TopologyChange(
TopologyChangeEvent::deserialize(buf)
.map_err(CqlEventParseError::TopologyChangeEventParseError)?,
)),
EventTypeV2::StatusChange => Ok(Self::StatusChange(
StatusChangeEvent::deserialize(buf)
.map_err(CqlEventParseError::StatusChangeEventParseError)?,
)),
EventTypeV2::SchemaChange => {
Ok(Self::SchemaChange(SchemaChangeEvent::deserialize(buf)?))
}
EventTypeV2::ClientRoutesChange => Ok(Self::ClientRoutesChange(
ClientRoutesChangeEvent::deserialize(buf)
.map_err(CqlEventParseError::ClientRoutesChangeEventParseError)?,
)),
}
}
}
impl SchemaChangeEvent {
pub fn deserialize(buf: &mut &[u8]) -> Result<Self, SchemaChangeEventParseError> {
let type_of_change_string =
types::read_string(buf).map_err(SchemaChangeEventParseError::TypeOfChangeParseError)?;
let type_of_change = match type_of_change_string {
"CREATED" => SchemaChangeType::Created,
"UPDATED" => SchemaChangeType::Updated,
"DROPPED" => SchemaChangeType::Dropped,
_ => SchemaChangeType::Invalid,
};
let target =
types::read_string(buf).map_err(SchemaChangeEventParseError::TargetTypeParseError)?;
let keyspace_affected = types::read_string(buf)
.map_err(SchemaChangeEventParseError::AffectedKeyspaceParseError)?
.to_string();
match target {
"KEYSPACE" => Ok(Self::KeyspaceChange {
change_type: type_of_change,
keyspace_name: keyspace_affected,
}),
"TABLE" => {
let table_name = types::read_string(buf)
.map_err(SchemaChangeEventParseError::AffectedTargetNameParseError)?
.to_string();
Ok(Self::TableChange {
change_type: type_of_change,
keyspace_name: keyspace_affected,
object_name: table_name,
})
}
"TYPE" => {
let changed_type = types::read_string(buf)
.map_err(SchemaChangeEventParseError::AffectedTargetNameParseError)?
.to_string();
Ok(Self::TypeChange {
change_type: type_of_change,
keyspace_name: keyspace_affected,
type_name: changed_type,
})
}
"FUNCTION" => {
let function = types::read_string(buf)
.map_err(SchemaChangeEventParseError::AffectedTargetNameParseError)?
.to_string();
let number_of_arguments = types::read_short(buf).map_err(|err| {
SchemaChangeEventParseError::ArgumentCountParseError(err.into())
})?;
let mut argument_vector = Vec::with_capacity(number_of_arguments as usize);
for _ in 0..number_of_arguments {
argument_vector.push(
types::read_string(buf)
.map_err(SchemaChangeEventParseError::FunctionArgumentParseError)?
.to_string(),
);
}
Ok(Self::FunctionChange {
change_type: type_of_change,
keyspace_name: keyspace_affected,
function_name: function,
arguments: argument_vector,
})
}
"AGGREGATE" => {
let name = types::read_string(buf)
.map_err(SchemaChangeEventParseError::AffectedTargetNameParseError)?
.to_string();
let number_of_arguments = types::read_short(buf).map_err(|err| {
SchemaChangeEventParseError::ArgumentCountParseError(err.into())
})?;
let mut argument_vector = Vec::with_capacity(number_of_arguments as usize);
for _ in 0..number_of_arguments {
argument_vector.push(
types::read_string(buf)
.map_err(SchemaChangeEventParseError::FunctionArgumentParseError)?
.to_string(),
);
}
Ok(Self::AggregateChange {
change_type: type_of_change,
keyspace_name: keyspace_affected,
aggregate_name: name,
arguments: argument_vector,
})
}
_ => Err(SchemaChangeEventParseError::UnknownTargetOfSchemaChange(
target.to_string(),
)),
}
}
}
impl TopologyChangeEvent {
pub fn deserialize(buf: &mut &[u8]) -> Result<Self, ClusterChangeEventParseError> {
let type_of_change = types::read_string(buf)
.map_err(ClusterChangeEventParseError::TypeOfChangeParseError)?;
let addr =
types::read_inet(buf).map_err(ClusterChangeEventParseError::NodeAddressParseError)?;
match type_of_change {
"NEW_NODE" => Ok(Self::NewNode(addr)),
"REMOVED_NODE" => Ok(Self::RemovedNode(addr)),
_ => Err(ClusterChangeEventParseError::UnknownTypeOfChange(
type_of_change.to_string(),
)),
}
}
}
impl StatusChangeEvent {
pub fn deserialize(buf: &mut &[u8]) -> Result<Self, ClusterChangeEventParseError> {
let type_of_change = types::read_string(buf)
.map_err(ClusterChangeEventParseError::TypeOfChangeParseError)?;
let addr =
types::read_inet(buf).map_err(ClusterChangeEventParseError::NodeAddressParseError)?;
match type_of_change {
"UP" => Ok(Self::Up(addr)),
"DOWN" => Ok(Self::Down(addr)),
_ => Err(ClusterChangeEventParseError::UnknownTypeOfChange(
type_of_change.to_string(),
)),
}
}
}
impl ClientRoutesChangeEvent {
pub fn deserialize(buf: &mut &[u8]) -> Result<Self, ClientRoutesChangeEventParseError> {
let type_of_change = types::read_string(buf)
.map_err(ClientRoutesChangeEventParseError::TypeOfChangeParseError)?;
match type_of_change {
"UPDATE_NODES" => {
}
_ => {
return Err(ClientRoutesChangeEventParseError::UnknownTypeOfChange(
type_of_change.to_string(),
));
}
}
let connection_ids = types::read_string_list(buf)
.map_err(ClientRoutesChangeEventParseError::ConnectionIdsParseError)?;
let connection_ids_count = connection_ids.len();
let (host_ids_count, host_ids_str_iter) = types::read_string_list_iter(buf)
.map_err(ClientRoutesChangeEventParseError::HostIdsParseError)?;
if connection_ids_count != host_ids_count {
return Err(
ClientRoutesChangeEventParseError::ConnectionHostIdsLengthMismatch {
connection_ids_count,
host_ids_count,
},
);
}
let host_ids: Vec<Uuid> = host_ids_str_iter
.map(|r| {
let host_id_str =
r.map_err(ClientRoutesChangeEventParseError::HostIdsParseError)?;
Uuid::try_parse(host_id_str)
.map_err(ClientRoutesChangeEventParseError::HostIdsUuidParseError)
})
.collect::<Result<_, _>>()?;
Ok(Self::UpdateNodes {
connection_ids,
host_ids,
})
}
}