# Streams System Documentation ## 1. Introduction The Streams system provides an append-only, persistent data streaming service built on top of NATS messaging. It enables multiple clients to write to named streams in a first-come-first-served basis, with real-time subscription capabilities and historical data querying. The system consists of: - **FileManipulator (FM)** service that manages stream persistence and operations - **Client libraries** for C++ and Python that interact via NATS messaging - **NATS messaging layer** for all client-server communication Key characteristics: - Append-only streams (no editing or deletion of messages) - Owner-based access control for write permissions - Multi-tier storage (memory → disk → compressed archives) - Real-time notifications and historical querying - No built-in authentication or encryption ## 2. Terms and Definitions ### Core Components - **Stream**: An append-only sequence of messages identified by name - **Message**: 0 or more lines of text with associated metadata - **Metadata**: Key-value pairs attached to each message (JSON values supported) - **FileManipulator (FM)**: Single service instance handling all stream persistence - **Owner**: Client that created/opened a stream, controls write permissions - **Subscriber**: Client receiving real-time updates from a stream ### Stream States - **Active**: Stream is open in memory, accepting writes and subscriptions - **Orphaned**: Owner disconnected, subscribers notified but stream remains active - **Closed**: No active subscribers, stream persisted and removed from memory - **Archived**: Older data compressed and stored on disk ### Storage Tiers - **Memory**: Current active portion of stream data - **Disk**: Persisted stream data in identical format to memory - **Archive**: Compressed (zip) older stream data, searchable but requires decompression ## 3. System Architecture ### 3.1 FileManipulator Service - Single-instance service (no clustering) - NATS client connecting to message broker - Manages all stream lifecycle operations - Handles three-tier storage system - Maintains in-memory permissions and stream metadata ### 3.2 Client Architecture - Clients connect via NATS messaging only - No direct communication with FileManipulator - Libraries provide abstraction over NATS protocol - Support for both C++ and Python implementations ### 3.3 Message Flow ``` Client → NATS → FileManipulator → Storage (Memory/Disk/Archive) ↓ NATS ← Subscribers (real-time updates) ``` ## 4. Data Structures ### 4.1 StreamMessage Format ```cpp class StreamMessage { std::map metadata_; std::string content_; std::string timestamp_; // ISO8601 string (system-generated) uint64_t sequence_id_; // stream-local, monotonic // Validation limits from Configuration: // stream_message.max_content_size_bytes // stream_message.max_metadata_keys // stream_message.max_metadata_key_length // stream_message.max_metadata_value_size } ``` ### 4.2 StreamPolicy Structure ```cpp class StreamPolicy { std::string name_; uint32_t max_size_mb_; // -1 = unlimited uint32_t memory_limit_mb_; // triggers storage at 80% threshold uint32_t disk_limit_mb_; // triggers archiving bool infinite_length_; // overrides size limits std::optional archive_date_; // ISO8601, cleanup threshold bool no_storage_; // memory-only mode bool no_archive_; // disk-only, no compression // Policy decisions: bool shouldTriggerStorage(size_t current_memory_mb) const; bool shouldTriggerArchive(size_t current_disk_mb) const; bool shouldCleanupArchive(const std::string& archive_timestamp) const; } ``` ### 4.3 StreamPermissions Structure ```cpp class StreamPermissions { std::string owner_; // client identifier std::set writers_; // clients with write access std::set readers_; // clients with read access // Access control methods: bool canWrite(const std::string& client_id) const; bool canRead(const std::string& client_id) const; bool isOwner(const std::string& client_id) const; } ``` ## 5. NATS Message Protocol ### 5.1 Subject Naming Convention - `streams.open` - Open/create stream - `streams.close` - Close stream (owner only) - `streams.subscribe` - Subscribe to stream updates - `streams.unsubscribe` - Unsubscribe from stream - `streams.post` - Post message to stream - `streams.query` - Query stream data/groups - `streams.search` - Reverse search stream history - `streams.change` - Change stream permissions (owner only) - `streams.info.request` - Request stream metadata - `streams.info.reply` - Stream metadata response - `streams.updates.{stream_name}` - Real-time stream updates - `streams.changes.{stream_name}` - Permission change notifications ### 5.2 Operation Result Format All operations return a standardised result structure: ```cpp struct OperationResult { bool success; std::string error_message; std::optional data; } ``` ### 5.3 Message Formats #### Stream Open Request ```json { "operation": "open", "stream_name": "string", "owner": "string", "policy": { "max_size_mb": 1000, "memory_limit_mb": 100, "disk_limit_mb": 500, "infinite_length": false, "archive_date": "2024-12-31T00:00:00Z", "no_storage": false, "no_archive": false } } ``` #### Stream Open Response ```json { "success": true, "error_message": "", "data": { "name": "string", "current_size_mb": 0.5, "message_count": 1250, "created": "2024-01-01T10:00:00Z", "last_activity": "2024-01-15T14:30:00Z" } } ``` #### Post Message Request ```json { "operation": "post", "stream_name": "string", "sender": "string", "message": { "metadata": { "sender": "client_id", "priority": "high", "tags": ["important", "notification"] }, "content": "Line 1\nLine 2\nLine 3" } } ``` #### Post Message Response ```json { "success": true, "error_message": "", "data": { "sequence_id": 1251, "timestamp": "2024-01-15T14:30:15Z" } } ``` #### Stream Query Request ```json { "operation": "query", "stream_name": "string", "query_type": "recent|range|search", "parameters": { "limit": 100, "from_sequence": 1000, "to_sequence": 1100, "search_term": "error", "reverse": true } } ``` #### Stream Update Notification ```json { "stream_name": "string", "update_type": "new_message|owner_disconnected|stream_closed|permission_changed", "data": { "sequence_id": 1251, "timestamp": "2024-01-15T14:30:15Z", "message": { /* full message object */ } } } ``` #### Permission Change Request (Owner Only) ```json { "operation": "change_permissions", "stream_name": "string", "owner": "string", "permissions": { "add_writers": ["client_2", "client_3"], "remove_writers": ["client_1"], "add_readers": ["client_4"], "remove_readers": [] } } ``` ## 6. Operations Specification ### 6.1 Stream Lifecycle Operations #### Open Stream **Purpose**: Create new stream or resume existing stream **Initiator**: Any client (becomes owner if new) **Preconditions**: Stream name is valid identifier **Process**: 1. FM checks if stream exists on disk/archive 2. If exists, verify owner identity and resume 3. If new, create stream with requesting client as owner 4. Load configuration or use defaults 5. Initialize in-memory structures 6. Send response with stream info #### Close Stream **Purpose**: Persist stream and remove from memory **Initiator**: Owner or automatic (no subscribers) **Process**: 1. Verify owner authorization (if manual) 2. Notify all subscribers of closure 3. Persist in-memory data to disk 4. Remove from active streams 5. Send confirmation response ### 6.2 Data Operations #### Post Message **Purpose**: Append new message to stream **Initiator**: Clients with write permission **Process**: 1. Validate client has write permission 2. Validate message format and size 3. Generate sequence ID and timestamp 4. Add to in-memory stream 5. Trigger storage if memory limits exceeded 6. Notify all subscribers 7. Send confirmation with sequence ID #### Query Stream **Purpose**: Retrieve historical messages **Initiator**: Any client with read access **Process**: 1. Determine query scope (memory/disk/archive) 2. For memory: direct access 3. For disk: read persisted files 4. For archive: decompress and read 5. Apply filters and limits 6. Return matching messages #### Reverse Search **Purpose**: Search backwards through stream history **Initiator**: Any client with read access **Process**: 1. Start from most recent data in memory 2. If not found, search disk storage 3. If still not found, search archives in reverse chronological order 4. Return first matches found ### 6.3 Subscription Operations #### Subscribe **Purpose**: Receive real-time updates from stream **Process**: 1. Validate stream exists and client has read access 2. Add client to subscriber list 3. Set up NATS subscription for stream updates 4. Send current stream status #### Unsubscribe **Purpose**: Stop receiving stream updates **Process**: 1. Remove client from subscriber list 2. Clean up NATS subscription 3. If no subscribers remain and owner disconnected, trigger stream closure ### 6.4 Administrative Operations #### Change Permissions **Purpose**: Grant/revoke read/write access **Initiator**: Stream owner only **Process**: 1. Verify owner authorization 2. Update in-memory permissions 3. Notify affected clients of changes 4. Log permission changes #### Stream Information **Purpose**: Get stream metadata and statistics **Process**: 1. Return current stream status 2. Include size, message count, activity timestamps 3. Include permission information (for owner) ## 7. Error Handling ### 7.1 Error Categories - **Authorization**: "Permission denied", "Stream owner required" - **Validation**: "Invalid message format", "Stream name invalid", "Content too large", "Too many metadata keys" - **Resource**: "Stream size limit exceeded", "Memory limit exceeded", "Disk storage full" - **State**: "Stream not found", "Stream closed", "Owner disconnected" - **System**: "Storage error", "Archive corruption detected", "Configuration error" ### 7.2 Configuration-Based Validation Errors All validation limits come from the Configuration system: - `stream_message.max_content_size_bytes` - Maximum message content size - `stream_message.max_metadata_keys` - Maximum number of metadata keys - `stream_message.max_metadata_key_length` - Maximum metadata key length - `stream_message.max_metadata_value_size` - Maximum metadata value size - `stream.max_name_length` - Maximum stream name length - `stream.min_name_length` - Minimum stream name length ### 7.3 Error Response Format ```json { "success": false, "error_message": "Stream 'logs-2024' does not exist", "data": { "error_code": "STREAM_NOT_FOUND", "suggested_action": "Create stream first or check stream name", "available_streams": ["logs-2023", "events-main"] } } ``` ## 8. Storage Management ### 8.1 Memory Management - **Trigger**: 80% of memory_limit_mb threshold (configurable via `storage.memory_threshold_percent`) - **Process**: Move oldest data to disk storage - **Format**: Identical to memory format for consistency - **Thread Safety**: All operations protected by stream-level mutexes - **Recovery**: Handle premature storage on memory pressure ### 8.2 Archive Management - **Trigger**: Disk storage size limits or age-based - **Format**: Compressed zip files with metadata - **Structure**: One archive per time period (configurable) - **Access**: Decompress on-demand for queries - **Compression**: Configurable level via `storage.archive_compression_level` ### 8.3 Cleanup Policies - **Archive Expiration**: Based on archive_date configuration in StreamPolicy - **Orphaned Streams**: Cleanup after configurable timeout (`stream.orphaned_timeout_hours`) - **Failed Operations**: Log and alert for manual intervention - **Storage Triggers**: Automatic based on policy thresholds ### 8.4 Thread Safety Model - **Stream-level locking**: Each Stream has its own mutex for operations - **Manager-level coordination**: StreamManager uses shared_mutex for stream map access - **Lock-free reads**: Where possible, use shared locks for read operations - **Deadlock prevention**: Consistent lock ordering and timeout mechanisms ## 9. Implementation Work Plan ### 9.1 Phase 1: Core Server Components (C++) #### 9.1.1 Stream Management Classes **Files**: `stream.hpp/.cpp`, `stream_manager.hpp/.cpp` **Stream Class Requirements**: ```cpp class Stream { std::string name; std::string owner; StreamConfig config; StreamPermissions permissions; std::deque messages; std::set subscribers; bool addMessage(const Message& msg); std::vector getMessages(QueryParams params); bool updatePermissions(const PermissionChange& change); void notifySubscribers(const StreamUpdate& update); }; ``` **StreamManager Class Requirements**: ```cpp class StreamManager { std::unordered_map> active_streams; Stream* openStream(const OpenRequest& request); bool closeStream(const std::string& name, const std::string& requester); bool postMessage(const PostRequest& request); QueryResult queryStream(const QueryRequest& request); }; ``` #### 9.1.2 Storage Layer Classes **Files**: `storage_manager.hpp/.cpp`, `archive_manager.hpp/.cpp` **StorageManager Requirements**: - Persist/load stream data to/from disk - Handle memory-to-disk transitions - Maintain file format consistency **ArchiveManager Requirements**: - Compress older data to zip archives - Decompress for historical queries - Manage archive cleanup policies #### 9.1.3 NATS Integration **Files**: `nats_handler.hpp/.cpp`, `message_protocol.hpp/.cpp` **NATSHandler Requirements**: ```cpp class NATSHandler { natsConnection* conn; std::unordered_map subscriptions; bool initialize(const std::string& nats_url); void handleStreamRequest(const std::string& subject, const std::string& data); void publishUpdate(const std::string& stream_name, const StreamUpdate& update); void subscribeToSubject(const std::string& subject, MessageCallback callback); }; ``` ### 9.2 Phase 2: Client Libraries #### 9.2.1 C++ Client Library **Files**: `streams_client.hpp/.cpp` **StreamsClient Requirements**: ```cpp class StreamsClient { std::unique_ptr nats; std::string client_id; bool openStream(const std::string& name, const StreamConfig& config = {}); bool postMessage(const std::string& stream, const Message& msg); bool subscribe(const std::string& stream, MessageCallback callback); QueryResult queryStream(const std::string& stream, const QueryParams& params); bool changePermissions(const std::string& stream, const PermissionChange& change); }; ``` #### 9.2.2 Python Client Library **Files**: `streams_client.py`, `message_types.py` **StreamsClient Requirements**: ```python class StreamsClient: def __init__(self, nats_url: str, client_id: str) def open_stream(self, name: str, config: StreamConfig = None) -> bool def post_message(self, stream: str, message: Message) -> bool def subscribe(self, stream: str, callback: Callable) -> bool def query_stream(self, stream: str, params: QueryParams) -> QueryResult def change_permissions(self, stream: str, change: PermissionChange) -> bool ``` ### 9.3 Phase 3: Integration and Testing #### 9.3.1 FileManipulator Integration - Integrate Stream system into existing FM service - Add configuration management for stream settings - Implement service startup/shutdown procedures #### 9.3.2 Testing Framework - Unit tests for all core classes - Integration tests for NATS communication - Performance tests for large streams and concurrent access - Storage/archive functionality testing #### 9.3.3 Documentation and Examples - API documentation for client libraries - Usage examples for common patterns - Deployment and configuration guides ## 10. Configuration ### 10.1 FileManipulator Configuration All configuration values are managed through the central Configuration class: ```ini # Stream message validation limits stream_message.max_content_size_bytes=1048576 stream_message.max_metadata_keys=50 stream_message.max_metadata_key_length=256 stream_message.max_metadata_value_size=4096 # Stream management settings stream.max_name_length=128 stream.min_name_length=1 stream.default_max_size_mb=1000 stream.default_memory_limit_mb=100 stream.default_disk_limit_mb=500 stream.orphaned_timeout_hours=24 stream.max_concurrent_streams=1000 # Storage and archiving policies storage.memory_threshold_percent=80 storage.archive_compression_level=6 storage.base_path=/var/lib/filemanipulator/streams storage.archive_path=/var/lib/filemanipulator/archives # NATS connection parameters nats.url=nats://localhost:4222 nats.reconnect_attempts=10 nats.reconnect_delay_ms=2000 nats.request_timeout_ms=5000 # Performance tuning performance.max_query_results=10000 performance.max_search_results=1000 performance.batch_size=100 # Security constraints security.max_stream_name_length=128 security.allowed_metadata_types=["string","number","boolean","array","object"] # Logging behaviour logging.log_all_operations=true logging.log_performance_metrics=true logging.log_level=INFO ``` ### 10.2 Client Configuration ```ini [client] nats_url=nats://localhost:4222 client_id_prefix=streams_client default_timeout_ms=5000 retry_attempts=3 ``` ## 11. Deployment Notes ### 11.1 Dependencies - **Server**: NATS C client library, zlib for compression, standard C++17 - **C++ Client**: NATS C++ client library - **Python Client**: NATS Python client (asyncio-nats) ### 11.2 Resource Requirements - Memory usage scales with active stream count and memory limits - Disk space for persistence and archives - Network bandwidth for real-time updates ### 11.3 Monitoring Points - Active stream count and memory usage - Message throughput per stream - Storage and archive operations - NATS connection health and message latency # Streams System Architecture & Development Plan ## Core Architecture Overview The streams system will be implemented as a modular component that can be integrated into the FileManipulator service. The design follows a layered architecture with clear separation of concerns. ## Shared Core Components ### `Configuration` Class **File**: `configuration.h/.cpp` (Shared across entire FileManipulator service) **Responsibility**: Centralized configuration management with fail-fast design ```cpp class Configuration { public: Configuration(); template T get(const std::string& key) const; // Throws if key doesn't exist template void set(const std::string& key, const T& value); bool has(const std::string& key) const; std::vector getAllKeys() const; static Configuration& getInstance(); private: std::map values_; void initializeDefaults(); // All magic numbers centralized here }; ``` **Key Design Principle**: Configuration throws `std::out_of_range` if you request a non-existent key. This forces developers to add any missing configuration values and prevents silent bugs from magic numbers. **Configuration Namespaces**: - `stream_message.*` - StreamMessage validation limits - `stream.*` - Stream management settings - `storage.*` - Storage and archiving policies - `nats.*` - NATS connection parameters - `performance.*` - Performance tuning values - `security.*` - Security constraints - `logging.*` - Logging behavior controls ## Class Hierarchy & Responsibilities ### 1. Data Layer Classes #### `StreamMessage` Class **File**: `stream_message.h/.cpp` **Responsibility**: Represent individual stream messages with comprehensive validation ```cpp class StreamMessage { public: using MetadataMap = std::map; // Constructors StreamMessage(); StreamMessage(const MetadataMap& metadata, const std::string& content); StreamMessage(const MetadataMap& metadata, const std::string& content, const std::string& timestamp, uint64_t sequence_id); // Getters (const methods for thread safety) const MetadataMap& getMetadata() const; const std::string& getContent() const; const std::string& getTimestamp() const; uint64_t getSequenceId() const; // System field setters (used by Stream class) void setTimestamp(const std::string& timestamp); void setSequenceId(uint64_t sequence_id); // Metadata manipulation OperationResult setMetadataValue(const std::string& key, const nlohmann::json& value); bool hasMetadataKey(const std::string& key) const; nlohmann::json getMetadataValue(const std::string& key) const; void removeMetadataKey(const std::string& key); // Validation and serialization OperationResult validate() const; // Uses Configuration for all limits std::string toJson() const; static OperationResult fromJson(const std::string& json, StreamMessage& out_message); // Utility methods static std::string generateTimestamp(); static bool isValidTimestamp(const std::string& timestamp); size_t estimateSize() const; // Comparison operators bool operator==(const StreamMessage& other) const; bool operator!=(const StreamMessage& other) const; bool operator<(const StreamMessage& other) const; // By sequence_id // Immutable-style builders StreamMessage withContent(const std::string& new_content) const; StreamMessage withMetadata(const MetadataMap& additional_metadata) const; private: MetadataMap metadata_; std::string content_; std::string timestamp_; uint64_t sequence_id_; // Validation helpers (use Configuration for all limits) OperationResult validateMetadata() const; OperationResult validateContent() const; OperationResult validateMetadataKey(const std::string& key) const; OperationResult validateMetadataValue(const nlohmann::json& value) const; }; ``` **Configuration Keys Used**: - `stream_message.max_content_size_bytes` - Maximum message content size - `stream_message.max_metadata_keys` - Maximum number of metadata keys - `stream_message.max_metadata_key_length` - Maximum metadata key length - `stream_message.max_metadata_value_size` - Maximum metadata value size #### `StreamPolicy` Class **File**: `stream_policy.h/.cpp` **Responsibility**: Make policy decisions about stream behaviour and resource management ```cpp class StreamPolicy { private: std::string name_; uint32_t max_size_mb_; uint32_t memory_limit_mb_; uint32_t disk_limit_mb_; bool infinite_length_; std::optional archive_date_; bool no_storage_; bool no_archive_; public: // Constructors StreamPolicy(); explicit StreamPolicy(const std::string& name); // Configuration management OperationResult setFromJson(const nlohmann::json& policy_json); OperationResult validate() const; nlohmann::json toJson() const; // Getters const std::string& getName() const; uint32_t getMaxSizeMB() const; uint32_t getMemoryLimitMB() const; uint32_t getDiskLimitMB() const; bool isInfiniteLength() const; const std::optional& getArchiveDate() const; bool isNoStorage() const; bool isNoArchive() const; // Setters with validation OperationResult setName(const std::string& name); OperationResult setMaxSizeMB(uint32_t max_size_mb); OperationResult setMemoryLimitMB(uint32_t memory_limit_mb); OperationResult setDiskLimitMB(uint32_t disk_limit_mb); void setInfiniteLength(bool infinite_length); OperationResult setArchiveDate(const std::optional& archive_date); void setNoStorage(bool no_storage); void setNoArchive(bool no_archive); // Policy decision methods bool shouldTriggerStorage(size_t current_memory_mb) const; bool shouldTriggerArchive(size_t current_disk_mb) const; bool shouldCleanupArchive(const std::string& archive_timestamp) const; bool allowsMoreData(size_t current_total_mb) const; // Utility methods void applyDefaults(); // Uses Configuration for default values bool isValid() const; private: // Validation helpers OperationResult validateName(const std::string& name) const; OperationResult validateSizeLimits() const; OperationResult validateArchiveDate(const std::string& date) const; }; ``` **Configuration Keys Used**: - `stream.default_max_size_mb` - Default maximum stream size - `stream.default_memory_limit_mb` - Default memory limit - `stream.default_disk_limit_mb` - Default disk limit - `stream.max_name_length` - Maximum stream name length - `stream.min_name_length` - Minimum stream name length - `storage.memory_threshold_percent` - Percentage threshold for storage trigger #### `StreamPermissions` Class **File**: `stream_permissions.hpp/.cpp` **Responsibility**: Manage access control for streams ```cpp class StreamPermissions { private: std::string owner_; std::set writers_; std::set readers_; public: // Permission management bool canWrite(const std::string& client_id) const; bool canRead(const std::string& client_id) const; bool isOwner(const std::string& client_id) const; void addWriter(const std::string& client_id); void removeWriter(const std::string& client_id); void addReader(const std::string& client_id); void removeReader(const std::string& client_id); }; ``` ### 2. Core Stream Management #### `Stream` Class **File**: `stream.hpp/.cpp` **Responsibility**: Manage individual stream lifecycle and operations ```cpp class Stream { private: std::string name_; StreamPolicy policy_; StreamPermissions permissions_; std::deque messages_; std::set subscribers_; uint64_t next_sequence_id_; std::mutex stream_mutex_; public: // Core stream operations Stream(const std::string& name, const std::string& owner, const StreamPolicy& policy); // Message operations bool addMessage(const StreamMessage& msg, const std::string& sender); std::vector getMessages(const QueryParams& params) const; std::vector searchMessages(const SearchParams& params) const; // Subscription management bool addSubscriber(const std::string& client_id); bool removeSubscriber(const std::string& client_id); const std::set& getSubscribers() const; // State management bool isEmpty() const; bool hasSubscribers() const; bool isOrphaned() const; size_t getCurrentSizeMB() const; uint64_t getMessageCount() const; // Permissions bool updatePermissions(const PermissionChange& change, const std::string& requester); const StreamPermissions& getPermissions() const; }; ``` #### `StreamManager` Class **File**: `stream_manager.hpp/.cpp` **Responsibility**: Coordinate all streams, handle lifecycle, routing ```cpp class StreamManager { private: std::unordered_map> active_streams_; std::unique_ptr storage_; std::unique_ptr archive_; mutable std::shared_mutex streams_mutex_; public: StreamManager(std::unique_ptr storage, std::unique_ptr archive); // Stream lifecycle Stream* openStream(const OpenStreamRequest& request); bool closeStream(const std::string& name, const std::string& requester); // Message operations bool postMessage(const PostMessageRequest& request); QueryResult queryStream(const QueryStreamRequest& request); SearchResult searchStream(const SearchStreamRequest& request); // Subscription management bool subscribe(const SubscribeRequest& request); bool unsubscribe(const UnsubscribeRequest& request); // Administrative bool changePermissions(const PermissionChangeRequest& request); StreamInfo getStreamInfo(const std::string& name) const; std::vector listStreams() const; // Cleanup and maintenance void performMaintenance(); void cleanupOrphanedStreams(); }; ``` ### 3. Storage Layer #### `StreamStorage` Interface & Implementation **File**: `stream_storage.hpp/.cpp` **Responsibility**: Handle persistent storage to disk ```cpp class StreamStorage { public: virtual ~StreamStorage() = default; // Stream persistence virtual bool persistStream(const Stream& stream) = 0; virtual std::unique_ptr loadStream(const std::string& name) = 0; virtual bool deleteStream(const std::string& name) = 0; // Message operations virtual std::vector getMessages(const std::string& stream_name, const QueryParams& params) = 0; virtual std::vector searchMessages(const std::string& stream_name, const SearchParams& params) = 0; // Metadata virtual bool streamExists(const std::string& name) const = 0; virtual std::vector listStoredStreams() const = 0; }; class FileStreamStorage : public StreamStorage { private: std::string base_path_; public: // Implementation of interface methods // File-based storage with JSON format }; ``` #### `StreamArchive` Interface & Implementation **File**: `stream_archive.hpp/.cpp` **Responsibility**: Handle compressed archive storage ```cpp class StreamArchive { public: virtual ~StreamArchive() = default; // Archive operations virtual bool archiveData(const std::string& stream_name, const std::vector& messages, const std::string& archive_period) = 0; virtual std::vector searchArchive(const std::string& stream_name, const SearchParams& params) = 0; virtual bool deleteArchive(const std::string& stream_name, const std::string& archive_period) = 0; // Archive management virtual std::vector listArchives(const std::string& stream_name) const = 0; virtual bool cleanupExpiredArchives(const std::string& archive_date) = 0; }; class ZipStreamArchive : public StreamArchive { private: std::string archive_path_; int compression_level_; public: // ZIP-based archive implementation }; ``` ### 4. Protocol Layer #### `StreamProtocol` Class **File**: `stream_protocol.hpp/.cpp` **Responsibility**: Handle NATS message serialization/deserialization ```cpp class StreamProtocol { public: // Request/Response serialization static std::string serializeOpenRequest(const OpenStreamRequest& request); static OpenStreamRequest deserializeOpenRequest(const std::string& json); static std::string serializeOpenResponse(const OpenStreamResponse& response); static OpenStreamResponse deserializeOpenResponse(const std::string& json); static std::string serializePostRequest(const PostMessageRequest& request); static PostMessageRequest deserializePostRequest(const std::string& json); // ... similar methods for all protocol messages // Update notifications static std::string serializeStreamUpdate(const StreamUpdate& update); static StreamUpdate deserializeStreamUpdate(const std::string& json); // Error handling static std::string serializeError(const StreamError& error); static StreamError deserializeError(const std::string& json); }; ``` #### Request/Response Data Structures **File**: `stream_requests.hpp` **Responsibility**: Define all protocol request/response structures ```cpp struct OpenStreamRequest { std::string operation; std::string stream_name; std::string owner; StreamPolicy policy; }; st ruct OpenStreamResponse { std::string status; std::string message; std::optional stream_info; }; struct PostMessageRequest { std::string operation; std::string stream_name; std::string sender; StreamMessage message; }; // ... other request/response structures ``` ### 5. NATS Integration #### `StreamsNATSHandler` Class **File**: `streams_nats_handler.hpp/.cpp` **Responsibility**: Handle NATS communication for streams ```cpp class StreamsNATSHandler { private: natsConnection* connection_; StreamManager* stream_manager_; std::unordered_map subscriptions_; public: StreamsNATSHandler(natsConnection* conn, StreamManager* manager); ~StreamsNATSHandler(); // NATS setup bool initialize(); void cleanup(); // Message handlers void handleStreamOpen(const std::string& data, const std::string& reply); void handleStreamClose(const std::string& data, const std::string& reply); void handleStreamPost(const std::string& data, const std::string& reply); void handleStreamQuery(const std::string& data, const std::string& reply); void handleStreamSubscribe(const std::string& data, const std::string& reply); void handleStreamUnsubscribe(const std::string& data, const std::string& reply); void handlePermissionChange(const std::string& data, const std::string& reply); void handleInfoRequest(const std::string& data, const std::string& reply); // Publishing bool publishStreamUpdate(const std::string& stream_name, const StreamUpdate& update); bool publishPermissionChange(const std::string& stream_name, const PermissionChange& change); bool sendResponse(const std::string& reply_subject, const std::string& response); }; ``` ### 6. Utility Classes #### `StreamsConfig` Class **File**: `streams_config.hpp/.cpp` **Responsibility**: Global streams system configuration ```cpp class StreamsConfig { private: uint32_t default_max_size_mb_; uint32_t default_memory_limit_mb_; uint32_t default_disk_limit_mb_; int archive_compression_level_; uint32_t orphaned_stream_timeout_hours_; uint32_t max_concurrent_streams_; public: StreamsConfig(); void loadFromFile(const std::string& config_file); void loadFromJson(const nlohmann::json& config); // Getters for all config values }; ``` #### `StreamsLogger` Class **File**: `streams_logger.hpp/.cpp` **Responsibility**: Logging specific to streams operations ```cpp class StreamsLogger { public: static void logStreamOperation(const std::string& operation, const std::string& stream_name, const std::string& client_id); static void logError(const std::string& operation, const std::string& error_msg); static void logPerformanceMetric(const std::string& metric, double value); }; ``` ## Development Phases ### Phase 1: Core Data Structures (Week 1) 1. **DONE** Implement `Configuration` class (shared component) 2. **DONE** Implement `StreamMessage` class with Configuration integration 3. **DONE** Implement `StreamPolicy` class (policy decision maker) 4. **DONE** Implement `StreamPermissions` class 5. **DONE** Implement request/response structures in `stream_requests.h` 6. **DONE** Create basic unit tests for data validation and serialization ### Phase 2: Stream Core Logic (Week 2) 1. Implement `Stream` class with in-memory operations 2. Implement `StreamProtocol` for message serialization 3. Create comprehensive unit tests for stream operations ### Phase 3: Storage Layer (Week 3) 1. Implement `FileStreamStorage` for disk persistence 2. Implement `ZipStreamArchive` for compression 3. Add storage integration tests ### Phase 4: Management Layer (Week 4) 1. Implement `StreamManager` coordinating all operations 2. Integration tests for full system without NATS ### Phase 5: NATS Integration (Week 5) 1. Implement `StreamsNATSHandler` 2. Add NATS protocol handling 3. End-to-end integration tests ### Phase 6: FileManipulator Integration (Week 6) 1. Create integration point in FileManipulator 2. Add streams configuration to FM config system 3. Full system testing and performance tuning ## File Structure ``` // Shared across entire FileManipulator service configuration.h/.cpp streams/ ├── core/ │ ├── stream_message.h/.cpp [DONE] │ ├── stream_policy.h/.cpp [DONE] │ ├── stream_permissions.h/.cpp [DONE] │ ├── stream.h/.cpp │ └── stream_manager.h/.cpp ├── storage/ │ ├── stream_storage.h/.cpp │ └── stream_archive.h/.cpp ├── protocol/ │ ├── stream_protocol.h/.cpp │ ├── stream_requests.h [DONE] │ └── streams_nats_handler.h/.cpp ├── utils/ │ └── streams_logger.h/.cpp └── tests/ ├── unit/ ├── integration/ └── performance/ ``` ## Key Changes Made ### 1. Configuration System - **Added `Configuration` class**: Centralized, fail-fast configuration management - **Eliminated all magic numbers**: All hardcoded values now come from configuration - **Shared component**: Used across entire FileManipulator service, not just streams - **Fail-fast design**: Throws exception if configuration key doesn't exist ### 2. StreamMessage Class (Renamed from Message) - **Renamed**: `Message` → `StreamMessage` to avoid naming conflicts - **Configuration integration**: All validation limits come from `Configuration::getInstance()` - **Complete implementation**: Fully functional with JSON serialization, validation, utilities - **Immutable-style builders**: `withContent()` and `withMetadata()` methods - **Rich functionality**: Metadata manipulation, size estimation, comparison operators ### 3. StreamPolicy Class (Renamed from StreamConfig) - **Renamed**: `StreamConfig` → `StreamPolicy` to better reflect its purpose - **Policy-oriented design**: Answers "what should I do?" questions rather than just holding data - **Configuration integration**: Uses central Configuration for defaults and validation limits This architecture provides clear separation of concerns, makes testing easier, and allows for gradual implementation and integration into the FileManipulator service.