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

class StreamMessage {
    std::map<std::string, nlohmann::json> 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

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<std::string> 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

class StreamPermissions {
    std::string owner_;              // client identifier
    std::set<std::string> writers_;  // clients with write access
    std::set<std::string> 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:

struct OperationResult {
    bool success;
    std::string error_message;
    std::optional<nlohmann::json> data;
}

5.3 Message Formats

Stream Open Request

{
    "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

{
    "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

{
    "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

{
    "success": true,
    "error_message": "",
    "data": {
        "sequence_id": 1251,
        "timestamp": "2024-01-15T14:30:15Z"
    }
}

Stream Query Request

{
    "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

{
    "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)

{
    "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

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

{
    "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:

class Stream {
    std::string name;
    std::string owner;
    StreamConfig config;
    StreamPermissions permissions;
    std::deque<Message> messages;
    std::set<std::string> subscribers;
    
    bool addMessage(const Message& msg);
    std::vector<Message> getMessages(QueryParams params);
    bool updatePermissions(const PermissionChange& change);
    void notifySubscribers(const StreamUpdate& update);
};

StreamManager Class Requirements:

class StreamManager {
    std::unordered_map<std::string, std::unique_ptr<Stream>> 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:

class NATSHandler {
    natsConnection* conn;
    std::unordered_map<std::string, natsSubscription*> 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:

class StreamsClient {
    std::unique_ptr<NATSConnection> 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:

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:

# 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

[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

class Configuration {
public:
    Configuration();
    
    template<typename T>
    T get(const std::string& key) const; // Throws if key doesn't exist
    
    template<typename T>
    void set(const std::string& key, const T& value);
    
    bool has(const std::string& key) const;
    std::vector<std::string> getAllKeys() const;
    static Configuration& getInstance();

private:
    std::map<std::string, std::any> 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

class StreamMessage {
public:
    using MetadataMap = std::map<std::string, nlohmann::json>;

    // 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

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<std::string> 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<std::string>& 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<std::string>& 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

class StreamPermissions {
private:
    std::string owner_;
    std::set<std::string> writers_;
    std::set<std::string> 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

class Stream {
private:
    std::string name_;
    StreamPolicy policy_;
    StreamPermissions permissions_;
    std::deque<StreamMessage> messages_;
    std::set<std::string> 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<StreamMessage> getMessages(const QueryParams& params) const;
    std::vector<StreamMessage> searchMessages(const SearchParams& params) const;
    
    // Subscription management
    bool addSubscriber(const std::string& client_id);
    bool removeSubscriber(const std::string& client_id);
    const std::set<std::string>& 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

class StreamManager {
private:
    std::unordered_map<std::string, std::unique_ptr<Stream>> active_streams_;
    std::unique_ptr<StreamStorage> storage_;
    std::unique_ptr<StreamArchive> archive_;
    mutable std::shared_mutex streams_mutex_;
    
public:
    StreamManager(std::unique_ptr<StreamStorage> storage,
                  std::unique_ptr<StreamArchive> 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<std::string> 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

class StreamStorage {
public:
    virtual ~StreamStorage() = default;
    
    // Stream persistence
    virtual bool persistStream(const Stream& stream) = 0;
    virtual std::unique_ptr<Stream> loadStream(const std::string& name) = 0;
    virtual bool deleteStream(const std::string& name) = 0;
    
    // Message operations
    virtual std::vector<StreamMessage> getMessages(const std::string& stream_name,
                                                  const QueryParams& params) = 0;
    virtual std::vector<StreamMessage> searchMessages(const std::string& stream_name,
                                                     const SearchParams& params) = 0;
    
    // Metadata
    virtual bool streamExists(const std::string& name) const = 0;
    virtual std::vector<std::string> 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

class StreamArchive {
public:
    virtual ~StreamArchive() = default;
    
    // Archive operations
    virtual bool archiveData(const std::string& stream_name,
                           const std::vector<StreamMessage>& messages,
                           const std::string& archive_period) = 0;
    virtual std::vector<StreamMessage> 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<std::string> 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

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

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<StreamInfo> 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

class StreamsNATSHandler {
private:
    natsConnection* connection_;
    StreamManager* stream_manager_;
    std::unordered_map<std::string, natsSubscription*> 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

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

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: MessageStreamMessage 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: StreamConfigStreamPolicy 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.