windjammer 0.48.0

A simple language inspired by Go, Ruby, and Elixir that transpiles to Rust - 80% of Rust's power with 20% of the complexity
Documentation
// Direct messaging (1-to-1 chat) for wschat

use std::collections::HashMap
use std::sync::{Arc, RwLock}
use std::time
use std::log

use ./message::{ClientMessage, ServerMessage, InternalMessage}
use ./storage::MessageStore

// Direct message manager for handling 1-to-1 conversations
pub struct DirectMessageManager {
    // Map of user_id -> list of connected user_ids they're chatting with
    active_conversations: Arc<RwLock<HashMap<string, Vec<string>>>>,
    
    // Message store for persistence
    store: Arc<MessageStore>,
}

impl DirectMessageManager {
    pub fn new(store: Arc<MessageStore>) -> Self {
        DirectMessageManager {
            active_conversations: Arc::new(RwLock::new(HashMap::new())),
            store: store,
        }
    }
    
    pub async fn send_direct_message(
        self,
        from_user: string,
        to_user: string,
        text: string
    ) -> Result<(), Error> {
        let timestamp = time.now_unix()
        
        // Create internal message
        let msg = InternalMessage {
            from_user: from_user.clone(),
            to_room: format!("dm:{}:{}", self.get_conversation_id(from_user.clone(), to_user.clone())),
            text: text.clone(),
            timestamp: timestamp,
        }
        
        // Store message in database
        self.store.store_message(msg).await?
        
        // Track active conversation
        self.add_to_active_conversations(from_user.clone(), to_user.clone())
        
        log.info("Direct message sent", {
            "from": from_user,
            "to": to_user,
            "length": text.len(),
        })
        
        Ok(())
    }
    
    pub async fn get_conversation_history(
        self,
        user1: string,
        user2: string,
        limit: int,
        before_timestamp: Option<int>
    ) -> Result<Vec<HistoryMessage>, Error> {
        let conversation_id = self.get_conversation_id(user1, user2)
        
        self.store.get_history(
            format!("dm:{}", conversation_id),
            limit,
            before_timestamp
        ).await
    }
    
    pub async fn get_user_conversations(
        self,
        user_id: string
    ) -> Result<Vec<ConversationInfo>, Error> {
        // Get all conversations for this user from database
        let query = "
            SELECT DISTINCT 
                CASE 
                    WHEN user_id = $1 THEN other_user_id
                    ELSE user_id 
                END as other_user,
                MAX(timestamp) as last_message_time,
                COUNT(*) as message_count
            FROM (
                SELECT 
                    CASE 
                        WHEN user_id < $1 THEN user_id 
                        ELSE $1 
                    END as user_id,
                    CASE 
                        WHEN user_id < $1 THEN $1 
                        ELSE user_id 
                    END as other_user_id,
                    timestamp
                FROM messages
                WHERE room_id LIKE 'dm:%' 
                  AND (user_id = $1 OR room_id LIKE '%' || $1 || '%')
            )
            GROUP BY other_user
            ORDER BY last_message_time DESC
        "
        
        let rows = self.store.pool.query_all(query, vec![user_id]).await?
        
        let mut conversations = vec![]
        
        for row in rows {
            let info = ConversationInfo {
                other_user: db.get_string(row, "other_user"),
                last_message_time: db.get_int(row, "last_message_time"),
                message_count: db.get_int(row, "message_count"),
                unread_count: 0,  // TODO: Track unread messages,
            }
            conversations.push(info)
        }
        
        Ok(conversations)
    }
    
    pub async fn mark_as_read(
        self,
        user_id: string,
        other_user_id: string
    ) -> Result<(), Error> {
        // TODO: Implement read receipts
        // For now, just log
        log.debug("Mark as read", {
            "user": user_id,
            "conversation_with": other_user_id,
        })
        
        Ok(())
    }
    
    pub async fn delete_conversation(
        self,
        user1: string,
        user2: string
    ) -> Result<int, Error> {
        let conversation_id = self.get_conversation_id(user1.clone(), user2.clone())
        
        let query = "DELETE FROM messages WHERE room_id = $1"
        let result = self.store.pool.execute(
            query,
            vec![format!("dm:{}", conversation_id)]
        ).await?
        
        // Remove from active conversations
        self.remove_from_active_conversations(user1, user2)
        
        log.info("Conversation deleted", {
            "conversation_id": conversation_id,
            "messages_deleted": result.rows_affected(),
        })
        
        Ok(result.rows_affected())
    }
    
    pub async fn get_unread_count(
        self,
        user_id: string
    ) -> Result<int, Error> {
        // TODO: Implement read tracking
        // For now, return 0
        Ok(0)
    }
    
    pub fn is_conversation_active(
        self,
        user1: string,
        user2: string
    ) -> bool {
        let conversations = self.active_conversations.read()
        
        if let Some(user_convos) = conversations.get(&user1) {
            if user_convos.contains(&user2) {
                return true
            }
        }
        
        if let Some(user_convos) = conversations.get(&user2) {
            if user_convos.contains(&user1) {
                return true
            }
        }
        
        false
    }
    
    // Private helpers
    
    fn get_conversation_id(self, user1: string, user2: string) -> string {
        // Create deterministic conversation ID (alphabetically sorted)
        if user1 < user2 {
            format!("{}:{}", user1, user2)
        } else {
            format!("{}:{}", user2, user1)
        }
    }
    
    fn add_to_active_conversations(self, user1: string, user2: string) {
        let mut conversations = self.active_conversations.write()
        
        // Add user2 to user1's conversations
        conversations.entry(user1.clone())
            .or_insert(vec![])
            .push(user2.clone())
        
        // Add user1 to user2's conversations
        conversations.entry(user2)
            .or_insert(vec![])
            .push(user1)
    }
    
    fn remove_from_active_conversations(self, user1: string, user2: string) {
        let mut conversations = self.active_conversations.write()
        
        // Remove user2 from user1's conversations
        if let Some(user_convos) = conversations.get_mut(&user1) {
            user_convos.retain(|u| u != &user2)
        }
        
        // Remove user1 from user2's conversations
        if let Some(user_convos) = conversations.get_mut(&user2) {
            user_convos.retain(|u| u != &user1)
        }
    }
}

@derive(Debug, Clone, Serialize, Deserialize)
pub struct ConversationInfo {
    pub other_user: string,
    pub last_message_time: int,
    pub message_count: int,
    pub unread_count: int,
}

@derive(Debug, Clone, Serialize, Deserialize)
pub struct DirectMessageRequest {
    pub to_user: string,
    pub text: string,
}

@derive(Debug, Clone, Serialize, Deserialize)
pub struct DirectMessageResponse {
    pub success: bool,
    pub message_id: Option<int>,
    pub error: Option<string>,
}