Skip to content

Latest commit

 

History

3 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Generic Messaging Queue (Mq)

A high-performance, persistent messaging queue implementation in .NET 9 with CRC32 integrity checking and asynchronous I/O operations.

Project Structure

The solution consists of three main projects:

Mq.Core

Core data structures and interfaces

  • Record.cs - Immutable record struct representing a timestamped message with payload
    • Uses readonly record struct for value equality semantics
    • Contains UTC timestamp in milliseconds and binary payload

Mq.Storage

Persistent storage layer with integrity checking

  • LogWriter.cs - Asynchronous log file writer with structured binary format
  • LogReader.cs - Asynchronous log file reader with data validation
  • Crc32Util.cs - CRC32 checksum utility for data integrity verification

Mq.Dev

Development and testing console application

  • Program.cs - Sample usage demonstrating write/read operations

Binary Log Format

The storage layer uses a custom binary format optimized for sequential access:

┌──────────┬──────────┬─────────────────┬─────────────────┬──────────────────┐
│  Length  │   CRC32  │  TimestampUtcMs │  Payload Length │     Payload      │
│  4 bytes │  4 bytes │     8 bytes     │     4 bytes     │   N bytes        │
└──────────┴──────────┴─────────────────┴─────────────────┴──────────────────┘
     │           │              └──────────────────┬──────────────────┘
     │           │                                 │
     │           └── CRC covers this portion ──────┘
     │
     └── Value = 4 (CRC) + 8 (ts) + 4 (payloadLen) + N (payload)

Format Details

  • Length Field: Total bytes following this field (CRC + timestamp + payload length + payload)
  • CRC32: Checksum covering timestamp, payload length, and payload data
  • Timestamp: UTC milliseconds since Unix epoch (Int64, little-endian)
  • Payload Length: Size of payload in bytes (Int32, little-endian)
  • Payload: Raw message data

Key Features

Data Integrity

  • CRC32 checksums on all records to detect corruption
  • Structured binary format with length prefixes for safe parsing
  • Validation during read operations with detailed error reporting

Performance Optimizations

  • Asynchronous I/O throughout the stack
  • Sequential file access patterns optimized for SSDs
  • 64KB buffer sizes for efficient disk operations
  • Memory-efficient streaming with IAsyncEnumerable<T>

Reliability

  • Append-only log structure prevents data loss
  • Concurrent read access while writing (FileShare.ReadWrite)
  • Proper resource disposal with IAsyncDisposable
  • Robust error handling for truncated or corrupted data

Current Implementation Status

✅ Completed Features:

  • Core record structure with value semantics
  • Binary log format with CRC32 integrity checking
  • Asynchronous log writer with buffering
  • Asynchronous log reader with streaming enumeration
  • Data validation and corruption detection
  • Sample console application demonstrating usage

🚧 Architecture Decisions:

  • Uses record struct for immutable, value-type semantics
  • Little-endian encoding for cross-platform compatibility
  • CRC32 for balance of performance and error detection
  • Async/await throughout for non-blocking I/O

Usage Example

var path = Path.Combine(AppContext.BaseDirectory, "data", "test.log");

// Writing records
await using (var writer = new LogWriter(path))
{
    var payload = Encoding.UTF8.GetBytes("Hello World");
    var record = new Record(DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(), payload);
    await writer.AppendAsync(record);
    await writer.FlushAsync();
}

// Reading records
await using (var reader = new LogReader(path))
{
    await foreach (var record in reader.ReadAllAsync())
    {
        var message = Encoding.UTF8.GetString(record.Payload.Span);
        Console.WriteLine($"{record.TimestampUtcMs}: {message}");
    }
}

Technical Requirements

  • .NET 9.0 - Latest LTS with performance improvements
  • System.IO.Hashing - For CRC32 implementation
  • Nullable reference types enabled for better null safety
  • Implicit usings for cleaner code

Next Steps

This foundation provides the core storage primitives for building a full messaging queue system. Future enhancements could include:

  • Message queue abstractions and pub/sub patterns
  • Index structures for efficient random access
  • Compression and encryption options
  • Distributed coordination and replication
  • Management APIs and monitoring tools

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages