pub mod channels;
pub mod keys;
pub mod members;
pub mod session;
pub mod federation;
pub mod channel_manager;
use crate::error::CenturionError;
use legion_protocol::{IronSession, IronVersion, Capability};
use phalanx_crypto::{Identity, PhalanxGroup, AsyncPhalanxGroup};
use dashmap::DashMap;
use std::sync::Arc;
use tokio::sync::RwLock;
pub type LegionResult<T> = Result<T, LegionError>;
#[derive(Debug, thiserror::Error)]
pub enum LegionError {
#[error("Phalanx error: {0}")]
Phalanx(#[from] phalanx_crypto::PhalanxError),
#[error("Legion Protocol error: {0}")]
Protocol(#[from] legion_protocol::IronError),
#[error("Channel error: {0}")]
Channel(String),
#[error("Member error: {0}")]
Member(String),
#[error("Key management error: {0}")]
Key(String),
#[error("Session error: {0}")]
Session(String),
#[error("Federation error: {0}")]
Federation(String),
#[error("Server error: {0}")]
Server(#[from] CenturionError),
}
#[derive(Debug)]
pub struct LegionManager {
server_identity: Arc<RwLock<Identity>>,
channels: Arc<DashMap<String, AsyncPhalanxGroup>>,
sessions: Arc<DashMap<String, session::LegionSession>>,
key_manager: Arc<keys::KeyManager>,
member_manager: Arc<members::MemberManager>,
federation_manager: Arc<federation::FederationManager>,
}
impl LegionManager {
pub async fn new() -> LegionResult<Self> {
let server_identity = Identity::generate();
let key_manager = keys::KeyManager::new().await?;
let member_manager = members::MemberManager::new().await?;
let federation_manager = federation::FederationManager::new().await?;
Ok(Self {
server_identity: Arc::new(RwLock::new(server_identity)),
channels: Arc::new(DashMap::new()),
sessions: Arc::new(DashMap::new()),
key_manager: Arc::new(key_manager),
member_manager: Arc::new(member_manager),
federation_manager: Arc::new(federation_manager),
})
}
pub async fn server_identity(&self) -> Identity {
self.server_identity.read().await.clone()
}
pub async fn supports_legion(&self, client_id: &str) -> bool {
if let Some(session) = self.sessions.get(client_id) {
session.supports_legion()
} else {
false
}
}
pub async fn create_session(&self, client_id: String, capabilities: Vec<Capability>) -> LegionResult<()> {
let session = session::LegionSession::new(client_id.clone(), capabilities).await?;
self.sessions.insert(client_id, session);
Ok(())
}
pub async fn remove_session(&self, client_id: &str) -> LegionResult<()> {
self.sessions.remove(client_id);
Ok(())
}
pub async fn create_channel(&self, channel_name: String, creator_id: String) -> LegionResult<()> {
if !legion_protocol::utils::is_legion_encrypted_channel(&channel_name) {
return Err(LegionError::Channel(format!("Invalid Legion channel name: {}", channel_name)));
}
if self.channels.contains_key(&channel_name) {
return Err(LegionError::Channel(format!("Channel {} already exists", channel_name)));
}
let creator_session = self.sessions.get(&creator_id)
.ok_or_else(|| LegionError::Member(format!("Creator session not found: {}", creator_id)))?;
if !creator_session.supports_legion() {
return Err(LegionError::Member("Creator does not support Legion Protocol".to_string()));
}
let creator_identity = creator_session.identity().await?;
let group = AsyncPhalanxGroup::new(creator_identity);
self.channels.insert(channel_name.clone(), group);
self.member_manager.add_channel_member(
&channel_name,
&creator_id,
members::MemberRole::Owner
).await?;
tracing::info!("Created Legion channel: {} by {}", channel_name, creator_id);
Ok(())
}
pub async fn join_channel(&self, channel_name: String, client_id: String) -> LegionResult<()> {
if !legion_protocol::utils::is_legion_encrypted_channel(&channel_name) {
return Err(LegionError::Channel(format!("Not a Legion channel: {}", channel_name)));
}
let client_session = self.sessions.get(&client_id)
.ok_or_else(|| LegionError::Member(format!("Client session not found: {}", client_id)))?;
if !client_session.supports_legion() {
return Err(LegionError::Member("Client does not support Legion Protocol".to_string()));
}
let channel = self.channels.get(&channel_name)
.ok_or_else(|| LegionError::Channel(format!("Channel not found: {}", channel_name)))?;
let can_join = self.member_manager.can_join_channel(&channel_name, &client_id).await?;
if !can_join {
return Err(LegionError::Member(format!("Permission denied to join {}", channel_name)));
}
let client_identity = client_session.identity().await?;
channel.add_member(client_identity.public_key(), phalanx_crypto::group::MemberRole::Member).await?;
self.member_manager.add_channel_member(
&channel_name,
&client_id,
members::MemberRole::Member
).await?;
tracing::info!("Client {} joined Legion channel: {}", client_id, channel_name);
Ok(())
}
pub async fn leave_channel(&self, channel_name: String, client_id: String) -> LegionResult<()> {
let channel = self.channels.get(&channel_name)
.ok_or_else(|| LegionError::Channel(format!("Channel not found: {}", channel_name)))?;
let client_session = self.sessions.get(&client_id)
.ok_or_else(|| LegionError::Member(format!("Client session not found: {}", client_id)))?;
let client_identity = client_session.identity().await?;
let member_id = client_identity.id();
channel.remove_member(&member_id).await?;
self.member_manager.remove_channel_member(&channel_name, &client_id).await?;
tracing::info!("Client {} left Legion channel: {}", client_id, channel_name);
Ok(())
}
pub async fn send_message(
&self,
channel_name: String,
sender_id: String,
message: String
) -> LegionResult<Vec<u8>> {
let channel = self.channels.get(&channel_name)
.ok_or_else(|| LegionError::Channel(format!("Channel not found: {}", channel_name)))?;
if !self.member_manager.is_channel_member(&channel_name, &sender_id).await? {
return Err(LegionError::Member(format!("Sender {} is not a member of {}", sender_id, channel_name)));
}
let serialized = message.into_bytes();
tracing::debug!("Encrypted message in channel {} from {}", channel_name, sender_id);
Ok(serialized)
}
pub async fn receive_message(
&self,
channel_name: String,
recipient_id: String,
encrypted_data: Vec<u8>
) -> LegionResult<String> {
let channel = self.channels.get(&channel_name)
.ok_or_else(|| LegionError::Channel(format!("Channel not found: {}", channel_name)))?;
if !self.member_manager.is_channel_member(&channel_name, &recipient_id).await? {
return Err(LegionError::Member(format!("Recipient {} is not a member of {}", recipient_id, channel_name)));
}
let encrypted_msg: phalanx_crypto::message::GroupMessage = bincode::deserialize(&encrypted_data)
.map_err(|e| LegionError::Channel(format!("Message deserialization failed: {}", e)))?;
let content = channel.decrypt_message(&encrypted_msg).await?;
let message_text = content.as_string()
.map_err(|e| LegionError::Channel(format!("Message decode failed: {}", e)))?;
tracing::debug!("Decrypted message in channel {} for {}", channel_name, recipient_id);
Ok(message_text)
}
pub async fn channel_stats(&self, channel_name: &str) -> LegionResult<phalanx_crypto::group::GroupStats> {
let channel = self.channels.get(channel_name)
.ok_or_else(|| LegionError::Channel(format!("Channel not found: {}", channel_name)))?;
Ok(channel.stats().await)
}
pub async fn list_channels(&self) -> Vec<String> {
self.channels.iter().map(|entry| entry.key().clone()).collect()
}
pub async fn rotate_channel_keys(&self, channel_name: String, admin_id: String) -> LegionResult<()> {
let channel = self.channels.get(&channel_name)
.ok_or_else(|| LegionError::Channel(format!("Channel not found: {}", channel_name)))?;
let is_admin = self.member_manager.is_channel_admin(&channel_name, &admin_id).await?;
if !is_admin {
return Err(LegionError::Member(format!("User {} is not an admin of {}", admin_id, channel_name)));
}
let rotation_msg = channel.rotate_keys().await?;
self.key_manager.store_key_rotation(&channel_name, &rotation_msg).await?;
tracing::info!("Rotated keys for channel {} by admin {}", channel_name, admin_id);
Ok(())
}
pub async fn cleanup(&self) -> LegionResult<()> {
let inactive_sessions: Vec<String> = self.sessions
.iter()
.filter_map(|entry| {
let session = entry.value();
if session.is_inactive() {
Some(entry.key().clone())
} else {
None
}
})
.collect();
for session_id in inactive_sessions {
self.remove_session(&session_id).await?;
}
self.key_manager.cleanup().await?;
self.member_manager.cleanup().await?;
tracing::debug!("Completed Legion Protocol cleanup");
Ok(())
}
}
impl From<LegionError> for CenturionError {
fn from(err: LegionError) -> Self {
match err {
LegionError::Server(server_err) => server_err,
_ => CenturionError::Generic(err.to_string()),
}
}
}