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 streamstreams.close- Close stream (owner only)streams.subscribe- Subscribe to stream updatesstreams.unsubscribe- Unsubscribe from streamstreams.post- Post message to streamstreams.query- Query stream data/groupsstreams.search- Reverse search stream historystreams.change- Change stream permissions (owner only)streams.info.request- Request stream metadatastreams.info.reply- Stream metadata responsestreams.updates.{stream_name}- Real-time stream updatesstreams.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:
FM checks if stream exists on disk/archive
If exists, verify owner identity and resume
If new, create stream with requesting client as owner
Load configuration or use defaults
Initialize in-memory structures
Send response with stream info
Close Stream¶
Purpose: Persist stream and remove from memory Initiator: Owner or automatic (no subscribers) Process:
Verify owner authorization (if manual)
Notify all subscribers of closure
Persist in-memory data to disk
Remove from active streams
Send confirmation response
6.2 Data Operations¶
Post Message¶
Purpose: Append new message to stream Initiator: Clients with write permission Process:
Validate client has write permission
Validate message format and size
Generate sequence ID and timestamp
Add to in-memory stream
Trigger storage if memory limits exceeded
Notify all subscribers
Send confirmation with sequence ID
Query Stream¶
Purpose: Retrieve historical messages Initiator: Any client with read access Process:
Determine query scope (memory/disk/archive)
For memory: direct access
For disk: read persisted files
For archive: decompress and read
Apply filters and limits
Return matching messages
Reverse Search¶
Purpose: Search backwards through stream history Initiator: Any client with read access Process:
Start from most recent data in memory
If not found, search disk storage
If still not found, search archives in reverse chronological order
Return first matches found
6.3 Subscription Operations¶
Subscribe¶
Purpose: Receive real-time updates from stream Process:
Validate stream exists and client has read access
Add client to subscriber list
Set up NATS subscription for stream updates
Send current stream status
Unsubscribe¶
Purpose: Stop receiving stream updates Process:
Remove client from subscriber list
Clean up NATS subscription
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:
Verify owner authorization
Update in-memory permissions
Notify affected clients of changes
Log permission changes
Stream Information¶
Purpose: Get stream metadata and statistics Process:
Return current stream status
Include size, message count, activity timestamps
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 sizestream_message.max_metadata_keys- Maximum number of metadata keysstream_message.max_metadata_key_length- Maximum metadata key lengthstream_message.max_metadata_value_size- Maximum metadata value sizestream.max_name_length- Maximum stream name lengthstream.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.
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 sizestream_message.max_metadata_keys- Maximum number of metadata keysstream_message.max_metadata_key_length- Maximum metadata key lengthstream_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 sizestream.default_memory_limit_mb- Default memory limitstream.default_disk_limit_mb- Default disk limitstream.max_name_length- Maximum stream name lengthstream.min_name_length- Minimum stream name lengthstorage.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)¶
DONE Implement
Configurationclass (shared component)DONE Implement
StreamMessageclass with Configuration integrationDONE Implement
StreamPolicyclass (policy decision maker)DONE Implement
StreamPermissionsclassDONE Implement request/response structures in
stream_requests.hDONE Create basic unit tests for data validation and serialization
Phase 2: Stream Core Logic (Week 2)¶
Implement
Streamclass with in-memory operationsImplement
StreamProtocolfor message serializationCreate comprehensive unit tests for stream operations
Phase 3: Storage Layer (Week 3)¶
Implement
FileStreamStoragefor disk persistenceImplement
ZipStreamArchivefor compressionAdd storage integration tests
Phase 4: Management Layer (Week 4)¶
Implement
StreamManagercoordinating all operationsIntegration tests for full system without NATS
Phase 5: NATS Integration (Week 5)¶
Implement
StreamsNATSHandlerAdd NATS protocol handling
End-to-end integration tests
Phase 6: FileManipulator Integration (Week 6)¶
Create integration point in FileManipulator
Add streams configuration to FM config system
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
Configurationclass: Centralized, fail-fast configuration managementEliminated 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→StreamMessageto avoid naming conflictsConfiguration integration: All validation limits come from
Configuration::getInstance()Complete implementation: Fully functional with JSON serialization, validation, utilities
Immutable-style builders:
withContent()andwithMetadata()methodsRich functionality: Metadata manipulation, size estimation, comparison operators
3. StreamPolicy Class (Renamed from StreamConfig)¶
Renamed:
StreamConfig→StreamPolicyto better reflect its purposePolicy-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.