Compare commits
11 Commits
a1fc600d84
...
497aa227af
| Author | SHA1 | Date | |
|---|---|---|---|
| 497aa227af | |||
| 4b15f643f6 | |||
| a470e0bcdb | |||
| 5a00708a79 | |||
| a5592ed533 | |||
| 20f45b2aaf | |||
| ca2d8019a1 | |||
| f57edca5a8 | |||
| 9ff5216495 | |||
| 5674853628 | |||
| 655ca30e0b |
@@ -286,43 +286,31 @@ public ValueTask<ulong> AppendAsync(string subject, ReadOnlyMemory<byte> payload
|
|||||||
|
|
||||||
### FileStore
|
### FileStore
|
||||||
|
|
||||||
`FileStore` appends messages to a JSONL file (`messages.jsonl`) and keeps a full in-memory index (`Dictionary<ulong, StoredMessage>`) identical in structure to `MemStore`. It is not production-safe for several reasons:
|
`FileStore` now persists messages into block files via `MsgBlock`, keeps a live in-memory message cache for load paths, and maintains a compact metadata index (`Dictionary<ulong, StoredMessageIndex>`) plus a per-subject last-sequence map for hot-path lookups such as `LoadLastBySubjectAsync`. Headers and payloads are stored separately and remain separate across snapshot, restore, block rewrite, and crash recovery. On startup, any legacy `messages.jsonl` file is migrated into block storage before recovery continues.
|
||||||
|
|
||||||
- **No locking**: `AppendAsync`, `LoadAsync`, `GetStateAsync`, and `TrimToMaxMessages` are not synchronized. Concurrent access from `StreamManager.Capture` and `PullConsumerEngine.FetchAsync` is unsafe.
|
|
||||||
- **Per-write file I/O**: Each `AppendAsync` calls `File.AppendAllTextAsync`, issuing a separate file open/write/close per message.
|
|
||||||
- **Full rewrite on trim**: `TrimToMaxMessages` calls `RewriteDataFile()`, which rewrites the entire file from the in-memory index. This is O(n) in message count and blocking.
|
|
||||||
- **Full in-memory index**: The in-memory dictionary holds every undeleted message payload; there is no paging or streaming read path.
|
|
||||||
|
|
||||||
```csharp
|
```csharp
|
||||||
// FileStore.cs
|
// FileStore.cs
|
||||||
public void TrimToMaxMessages(ulong maxMessages)
|
private readonly Dictionary<ulong, StoredMessage> _messages = new();
|
||||||
|
private readonly Dictionary<ulong, StoredMessageIndex> _messageIndexes = new();
|
||||||
|
private readonly Dictionary<string, ulong> _lastSequenceBySubject = new(StringComparer.Ordinal);
|
||||||
|
|
||||||
|
public ValueTask<StoredMessage?> LoadLastBySubjectAsync(string subject, CancellationToken ct)
|
||||||
{
|
{
|
||||||
while ((ulong)_messages.Count > maxMessages)
|
if (_lastSequenceBySubject.TryGetValue(subject, out var sequence)
|
||||||
|
&& _messages.TryGetValue(sequence, out var match))
|
||||||
{
|
{
|
||||||
var first = _messages.Keys.Min();
|
return ValueTask.FromResult<StoredMessage?>(match);
|
||||||
_messages.Remove(first);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
RewriteDataFile();
|
return ValueTask.FromResult<StoredMessage?>(null);
|
||||||
}
|
|
||||||
|
|
||||||
private void RewriteDataFile()
|
|
||||||
{
|
|
||||||
var lines = new List<string>(_messages.Count);
|
|
||||||
foreach (var message in _messages.OrderBy(kv => kv.Key).Select(kv => kv.Value))
|
|
||||||
{
|
|
||||||
lines.Add(JsonSerializer.Serialize(new FileRecord
|
|
||||||
{
|
|
||||||
Sequence = message.Sequence,
|
|
||||||
Subject = message.Subject,
|
|
||||||
PayloadBase64 = Convert.ToBase64String(message.Payload.ToArray()),
|
|
||||||
}));
|
|
||||||
}
|
|
||||||
File.WriteAllLines(_dataFilePath, lines);
|
|
||||||
}
|
}
|
||||||
```
|
```
|
||||||
|
|
||||||
The Go reference (`filestore.go`) uses block-based binary storage with S2 compression, per-block indexes, and memory-mapped I/O. This implementation shares none of those properties.
|
The current implementation is still materially simpler than Go `filestore.go`:
|
||||||
|
|
||||||
|
- **No synchronization**: `FileStore` still exposes unsynchronized mutation and read paths. It is safe only under the current test and single-process usage assumptions.
|
||||||
|
- **Payloads still stay resident**: the compact index removes duplicate payload ownership for metadata-heavy operations, but `_messages` still retains live payload bytes in memory for direct load paths.
|
||||||
|
- **No Go-equivalent block index stack**: there is no per-block subject tree, mmap-backed read path, or Go-style cache/compaction parity. Deletes and trims rely on tombstones plus later block maintenance rather than Go's full production filestore behavior.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -445,7 +433,7 @@ The following features are present in the Go reference (`golang/nats-server/serv
|
|||||||
- **Ephemeral consumers**: `ConsumerManager.CreateOrUpdate` requires a non-empty `DurableName`. There is no support for unnamed ephemeral consumers.
|
- **Ephemeral consumers**: `ConsumerManager.CreateOrUpdate` requires a non-empty `DurableName`. There is no support for unnamed ephemeral consumers.
|
||||||
- **Push delivery over the NATS wire**: Push consumers enqueue `PushFrame` objects into an in-memory queue. No MSG is written to any connected NATS client's TCP socket.
|
- **Push delivery over the NATS wire**: Push consumers enqueue `PushFrame` objects into an in-memory queue. No MSG is written to any connected NATS client's TCP socket.
|
||||||
- **Consumer filter subject enforcement**: `FilterSubject` is stored on `ConsumerConfig` but is never applied in `PullConsumerEngine.FetchAsync`. All messages in the stream are returned regardless of filter.
|
- **Consumer filter subject enforcement**: `FilterSubject` is stored on `ConsumerConfig` but is never applied in `PullConsumerEngine.FetchAsync`. All messages in the stream are returned regardless of filter.
|
||||||
- **FileStore production safety**: No locking, per-write file I/O, full-rewrite-on-trim, and full in-memory index make `FileStore` unsuitable for production use.
|
- **FileStore production safety**: `FileStore` now uses block files and compact metadata indexes, but it still lacks synchronization and Go-level block indexing, so it remains unsuitable for production use.
|
||||||
- **RAFT persistence and networking**: `RaftNode` log entries are not persisted across restarts. Replication uses direct in-process method calls; there is no network transport for multi-server consensus.
|
- **RAFT persistence and networking**: `RaftNode` log entries are not persisted across restarts. Replication uses direct in-process method calls; there is no network transport for multi-server consensus.
|
||||||
- **Cross-server replication**: Mirror and source coordinators work only within one `StreamManager` in one process. Messages published on a remote server are not replicated.
|
- **Cross-server replication**: Mirror and source coordinators work only within one `StreamManager` in one process. Messages published on a remote server are not replicated.
|
||||||
- **Duplicate message window**: `PublishPreconditions` tracks message IDs for deduplication but there is no configurable `DuplicateWindow` TTL to expire old IDs.
|
- **Duplicate message window**: `PublishPreconditions` tracks message IDs for deduplication but there is no configurable `DuplicateWindow` TTL to expire old IDs.
|
||||||
@@ -460,4 +448,4 @@ The following features are present in the Go reference (`golang/nats-server/serv
|
|||||||
- [Configuration Overview](../Configuration/Overview.md)
|
- [Configuration Overview](../Configuration/Overview.md)
|
||||||
- [Protocol Overview](../Protocol/Overview.md)
|
- [Protocol Overview](../Protocol/Overview.md)
|
||||||
|
|
||||||
<!-- Last verified against codebase: 2026-02-23 -->
|
<!-- Last verified against codebase: 2026-03-13 -->
|
||||||
|
|||||||
+46
-38
@@ -1,6 +1,6 @@
|
|||||||
# Go vs .NET NATS Server — Benchmark Comparison
|
# Go vs .NET NATS Server — Benchmark Comparison
|
||||||
|
|
||||||
Benchmark run: 2026-03-13 10:16 AM America/Indiana/Indianapolis. Both servers ran on the same machine using the benchmark project README command (`dotnet test tests/NATS.Server.Benchmark.Tests --filter "Category=Benchmark" -v normal --logger "console;verbosity=detailed"`). Test parallelization remained disabled inside the benchmark assembly.
|
Benchmark run: 2026-03-13 11:41 AM America/Indiana/Indianapolis. Both servers ran on the same machine using the benchmark project README command (`dotnet test tests/NATS.Server.Benchmark.Tests --filter "Category=Benchmark" -v normal --logger "console;verbosity=detailed"`). Test parallelization remained disabled inside the benchmark assembly.
|
||||||
|
|
||||||
**Environment:** Apple M4, .NET SDK 10.0.101, benchmark README command run in the benchmark project's default `Debug` configuration, Go toolchain installed, Go reference server built from `golang/nats-server/`.
|
**Environment:** Apple M4, .NET SDK 10.0.101, benchmark README command run in the benchmark project's default `Debug` configuration, Go toolchain installed, Go reference server built from `golang/nats-server/`.
|
||||||
|
|
||||||
@@ -13,27 +13,27 @@ Benchmark run: 2026-03-13 10:16 AM America/Indiana/Indianapolis. Both servers ra
|
|||||||
|
|
||||||
| Payload | Go msg/s | Go MB/s | .NET msg/s | .NET MB/s | Ratio (.NET/Go) |
|
| Payload | Go msg/s | Go MB/s | .NET msg/s | .NET MB/s | Ratio (.NET/Go) |
|
||||||
|---------|----------|---------|------------|-----------|-----------------|
|
|---------|----------|---------|------------|-----------|-----------------|
|
||||||
| 16 B | 2,258,647 | 34.5 | 1,275,230 | 19.5 | 0.56x |
|
| 16 B | 2,223,690 | 33.9 | 1,341,067 | 20.5 | 0.60x |
|
||||||
| 128 B | 2,251,274 | 274.8 | 1,661,668 | 202.8 | 0.74x |
|
| 128 B | 2,218,308 | 270.8 | 1,577,523 | 192.6 | 0.71x |
|
||||||
|
|
||||||
### Publisher + Subscriber (1:1)
|
### Publisher + Subscriber (1:1)
|
||||||
|
|
||||||
| Payload | Go msg/s | Go MB/s | .NET msg/s | .NET MB/s | Ratio (.NET/Go) |
|
| Payload | Go msg/s | Go MB/s | .NET msg/s | .NET MB/s | Ratio (.NET/Go) |
|
||||||
|---------|----------|---------|------------|-----------|-----------------|
|
|---------|----------|---------|------------|-----------|-----------------|
|
||||||
| 16 B | 296,374 | 4.5 | 875,105 | 13.4 | **2.95x** |
|
| 16 B | 292,711 | 4.5 | 862,381 | 13.2 | **2.95x** |
|
||||||
| 16 KB | 32,111 | 501.7 | 30,030 | 469.2 | 0.94x |
|
| 16 KB | 32,890 | 513.9 | 28,906 | 451.7 | 0.88x |
|
||||||
|
|
||||||
### Fan-Out (1 Publisher : 4 Subscribers)
|
### Fan-Out (1 Publisher : 4 Subscribers)
|
||||||
|
|
||||||
| Payload | Go msg/s | Go MB/s | .NET msg/s | .NET MB/s | Ratio (.NET/Go) |
|
| Payload | Go msg/s | Go MB/s | .NET msg/s | .NET MB/s | Ratio (.NET/Go) |
|
||||||
|---------|----------|---------|------------|-----------|-----------------|
|
|---------|----------|---------|------------|-----------|-----------------|
|
||||||
| 128 B | 2,387,889 | 291.5 | 1,780,888 | 217.4 | 0.75x |
|
| 128 B | 2,945,790 | 359.6 | 1,858,235 | 226.8 | 0.63x |
|
||||||
|
|
||||||
### Multi-Publisher / Multi-Subscriber (4P x 4S)
|
### Multi-Publisher / Multi-Subscriber (4P x 4S)
|
||||||
|
|
||||||
| Payload | Go msg/s | Go MB/s | .NET msg/s | .NET MB/s | Ratio (.NET/Go) |
|
| Payload | Go msg/s | Go MB/s | .NET msg/s | .NET MB/s | Ratio (.NET/Go) |
|
||||||
|---------|----------|---------|------------|-----------|-----------------|
|
|---------|----------|---------|------------|-----------|-----------------|
|
||||||
| 128 B | 1,079,112 | 131.7 | 953,596 | 116.4 | 0.88x |
|
| 128 B | 2,123,480 | 259.2 | 1,392,249 | 170.0 | 0.66x |
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -43,13 +43,13 @@ Benchmark run: 2026-03-13 10:16 AM America/Indiana/Indianapolis. Both servers ra
|
|||||||
|
|
||||||
| Payload | Go msg/s | .NET msg/s | Ratio | Go P50 (us) | .NET P50 (us) | Go P99 (us) | .NET P99 (us) |
|
| Payload | Go msg/s | .NET msg/s | Ratio | Go P50 (us) | .NET P50 (us) | Go P99 (us) | .NET P99 (us) |
|
||||||
|---------|----------|------------|-------|-------------|---------------|-------------|---------------|
|
|---------|----------|------------|-------|-------------|---------------|-------------|---------------|
|
||||||
| 128 B | 8,506 | 7,182 | 0.84x | 114.9 | 135.2 | 161.2 | 189.8 |
|
| 128 B | 8,386 | 7,014 | 0.84x | 115.8 | 139.0 | 175.5 | 193.0 |
|
||||||
|
|
||||||
### 10 Clients, 2 Services (Queue Group)
|
### 10 Clients, 2 Services (Queue Group)
|
||||||
|
|
||||||
| Payload | Go msg/s | .NET msg/s | Ratio | Go P50 (us) | .NET P50 (us) | Go P99 (us) | .NET P99 (us) |
|
| Payload | Go msg/s | .NET msg/s | Ratio | Go P50 (us) | .NET P50 (us) | Go P99 (us) | .NET P99 (us) |
|
||||||
|---------|----------|------------|-------|-------------|---------------|-------------|---------------|
|
|---------|----------|------------|-------|-------------|---------------|-------------|---------------|
|
||||||
| 16 B | 26,610 | 22,533 | 0.85x | 367.7 | 425.3 | 487.4 | 622.5 |
|
| 16 B | 26,470 | 23,478 | 0.89x | 370.2 | 410.6 | 486.0 | 592.8 |
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -57,10 +57,10 @@ Benchmark run: 2026-03-13 10:16 AM America/Indiana/Indianapolis. Both servers ra
|
|||||||
|
|
||||||
| Mode | Payload | Storage | Go msg/s | .NET msg/s | Ratio (.NET/Go) |
|
| Mode | Payload | Storage | Go msg/s | .NET msg/s | Ratio (.NET/Go) |
|
||||||
|------|---------|---------|----------|------------|-----------------|
|
|------|---------|---------|----------|------------|-----------------|
|
||||||
| Synchronous | 16 B | Memory | 13,756 | 9,954 | 0.72x |
|
| Synchronous | 16 B | Memory | 14,812 | 12,134 | 0.82x |
|
||||||
| Async (batch) | 128 B | File | 171,761 | 50,711 | 0.30x |
|
| Async (batch) | 128 B | File | 148,156 | 57,479 | 0.39x |
|
||||||
|
|
||||||
> **Note:** Async file-store publish remains the largest JetStream gap at 0.30x. The bottleneck is still the storage write path and the remaining managed allocation pressure around persisted message state.
|
> **Note:** Async file-store publish remains well below parity at 0.39x, but it is still materially better than the older 0.30x snapshot that motivated this FileStore round.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -68,10 +68,10 @@ Benchmark run: 2026-03-13 10:16 AM America/Indiana/Indianapolis. Both servers ra
|
|||||||
|
|
||||||
| Mode | Go msg/s | .NET msg/s | Ratio (.NET/Go) |
|
| Mode | Go msg/s | .NET msg/s | Ratio (.NET/Go) |
|
||||||
|------|----------|------------|-----------------|
|
|------|----------|------------|-----------------|
|
||||||
| Ordered ephemeral consumer | 135,704 | 107,168 | 0.79x |
|
| Ordered ephemeral consumer | 572,941 | 101,944 | 0.18x |
|
||||||
| Durable consumer fetch | 533,441 | 375,652 | 0.70x |
|
| Durable consumer fetch | 599,204 | 338,265 | 0.56x |
|
||||||
|
|
||||||
> **Note:** Ordered-consumer results in this run are much closer to parity than earlier snapshots. That suggests prior Go-side variance was material; `.NET` throughput is still clustered around ~107K msg/s.
|
> **Note:** Ordered-consumer throughput remains the clearest JetStream hotspot after this round. The merged FileStore work helped publish and subject-lookup paths more than consumer delivery.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -81,19 +81,27 @@ Benchmark run: 2026-03-13 10:16 AM America/Indiana/Indianapolis. Both servers ra
|
|||||||
|
|
||||||
| Benchmark | .NET msg/s | .NET MB/s | Alloc |
|
| Benchmark | .NET msg/s | .NET MB/s | Alloc |
|
||||||
|-----------|------------|-----------|-------|
|
|-----------|------------|-----------|-------|
|
||||||
| SubList Exact Match (128 subjects) | 17,746,607 | 236.9 | 0.00 B/op |
|
| SubList Exact Match (128 subjects) | 16,497,186 | 220.3 | 0.00 B/op |
|
||||||
| SubList Wildcard Match | 18,811,278 | 251.2 | 0.00 B/op |
|
| SubList Wildcard Match | 16,147,367 | 215.6 | 0.00 B/op |
|
||||||
| SubList Queue Match | 20,624,510 | 157.4 | 0.00 B/op |
|
| SubList Queue Match | 15,582,052 | 118.9 | 0.00 B/op |
|
||||||
| SubList Remote Interest | 264,725 | 4.3 | 0.00 B/op |
|
| SubList Remote Interest | 259,940 | 4.2 | 0.00 B/op |
|
||||||
|
|
||||||
### Parser
|
### Parser
|
||||||
|
|
||||||
| Benchmark | Ops/s | MB/s | Alloc |
|
| Benchmark | Ops/s | MB/s | Alloc |
|
||||||
|-----------|-------|------|-------|
|
|-----------|-------|------|-------|
|
||||||
| Parser PING | 5,598,176 | 32.0 | 0.0 B/op |
|
| Parser PING | 6,283,578 | 36.0 | 0.0 B/op |
|
||||||
| Parser PUB | 2,701,645 | 103.1 | 40.0 B/op |
|
| Parser PUB | 2,712,550 | 103.5 | 40.0 B/op |
|
||||||
| Parser HPUB | 2,177,745 | 116.3 | 40.0 B/op |
|
| Parser HPUB | 2,338,555 | 124.9 | 40.0 B/op |
|
||||||
| Parser PUB split payload | 1,702,439 | 64.9 | 176.0 B/op |
|
| Parser PUB split payload | 2,043,813 | 78.0 | 176.0 B/op |
|
||||||
|
|
||||||
|
### FileStore
|
||||||
|
|
||||||
|
| Benchmark | Ops/s | MB/s | Alloc |
|
||||||
|
|-----------|-------|------|-------|
|
||||||
|
| FileStore AppendAsync (128B) | 244,089 | 29.8 | 1552.9 B/op |
|
||||||
|
| FileStore LoadLastBySubject (hot) | 12,784,127 | 780.3 | 0.0 B/op |
|
||||||
|
| FileStore PurgeEx+Trim | 332 | 0.0 | 5440792.9 B/op |
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -101,25 +109,25 @@ Benchmark run: 2026-03-13 10:16 AM America/Indiana/Indianapolis. Both servers ra
|
|||||||
|
|
||||||
| Category | Ratio Range | Assessment |
|
| Category | Ratio Range | Assessment |
|
||||||
|----------|-------------|------------|
|
|----------|-------------|------------|
|
||||||
| Pub-only throughput | 0.56x–0.74x | Mixed — 128 B is solid, 16 B still trails materially |
|
| Pub-only throughput | 0.60x–0.71x | Mixed; still behind Go |
|
||||||
| Pub/sub (small payload) | **2.95x** | .NET outperforms Go decisively |
|
| Pub/sub (small payload) | **2.95x** | .NET outperforms Go decisively |
|
||||||
| Pub/sub (large payload) | 0.94x | Near parity |
|
| Pub/sub (large payload) | 0.88x | Close, but below parity |
|
||||||
| Fan-out | 0.75x | Good improvement; still limited by serial delivery |
|
| Fan-out | 0.63x | Still materially behind Go |
|
||||||
| Multi pub/sub | 0.88x | Close to parity in this run |
|
| Multi pub/sub | 0.66x | Meaningful gap remains |
|
||||||
| Request/reply latency | 0.84x–0.85x | Good |
|
| Request/reply latency | 0.84x–0.89x | Good |
|
||||||
| JetStream sync publish | 0.72x | Good |
|
| JetStream sync publish | 0.82x | Strong |
|
||||||
| JetStream async file publish | 0.30x | Storage write path still dominates |
|
| JetStream async file publish | 0.39x | Improved versus older snapshots, still storage-bound |
|
||||||
| JetStream ordered consume | 0.79x | Much closer to parity in this run |
|
| JetStream ordered consume | 0.18x | Highest-priority JetStream gap |
|
||||||
| JetStream durable fetch | 0.70x | Good |
|
| JetStream durable fetch | 0.56x | Regressed from prior snapshot |
|
||||||
|
|
||||||
### Key Observations
|
### Key Observations
|
||||||
|
|
||||||
1. **Small-payload 1:1 pub/sub still beats Go by ~3x** (875K vs 296K msg/s). The direct write path continues to pay off when message fanout is simple and payloads are tiny.
|
1. **Small-payload 1:1 pub/sub is back to a large `.NET` lead in this final run** at 2.95x (862K vs 293K msg/s). That puts the merged benchmark profile much closer to the earlier comparison snapshot than the intermediate integration-only run.
|
||||||
2. **Fan-out and multi pub/sub both improved in this run** to 0.75x and 0.88x respectively. The remaining gap is still consistent with Go's more naturally parallel fanout model.
|
2. **Async file-store publish is still materially better than the older 0.30x baseline** at 0.39x (57.5K vs 148.2K msg/s), which is consistent with the FileStore metadata and payload-ownership changes helping the write path even though they did not eliminate the gap.
|
||||||
3. **Ordered consumer moved up to 0.79x** (107K vs 136K msg/s). That is materially stronger than earlier runs and suggests previous Go-side variance was distorting the comparison more than the `.NET` consumer path itself.
|
3. **The new FileStore direct benchmarks show what remains expensive in storage maintenance**: `LoadLastBySubject` is allocation-free and extremely fast, `AppendAsync` is still about 1553 B/op, and repeated `PurgeEx+Trim` still burns roughly 5.4 MB/op.
|
||||||
4. **Durable fetch remains solid at 0.70x**. The Round 6 fetch-path work is still holding, but there is room left in consumer dispatch and storage reads.
|
4. **Ordered consumer throughput remains the largest JetStream gap at 0.18x** (102K vs 573K msg/s). That is better than the intermediate 0.11x run, but it is still the clearest post-FileStore optimization target.
|
||||||
5. **Async file-store publish is still the largest server-level gap at 0.30x**. The storage layer remains the highest-value runtime target after parser and SubList hot-path cleanup.
|
5. **Durable fetch regressed to 0.56x in the final run**, which keeps consumer delivery and storage-read coordination in the top tier of remaining work even after the FileStore changes.
|
||||||
6. **The new SubList microbenchmarks show effectively zero temporary allocation per operation** for exact, wildcard, queue, and remote-interest lookups in the current implementation. Parser contiguous hot paths also remain small and stable, while split-payload `PUB` still pays a higher copy cost.
|
6. **Parser and SubList microbenchmarks remain stable and low-allocation**. The storage and consumer layers continue to dominate the server-level benchmark gaps, not the parser or subject matcher hot paths.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
|||||||
@@ -95,9 +95,6 @@ public sealed class CompiledFilter
|
|||||||
|
|
||||||
public sealed class PullConsumerEngine
|
public sealed class PullConsumerEngine
|
||||||
{
|
{
|
||||||
// Reusable fetch buffer to avoid per-fetch List allocation.
|
|
||||||
// Go reference: consumer.go — reuses slice backing array across fetch calls.
|
|
||||||
[ThreadStatic] private static List<StoredMessage>? t_fetchBuf;
|
|
||||||
// Go: consumer.go — cluster-wide pending pull request tracking keyed by reply subject.
|
// Go: consumer.go — cluster-wide pending pull request tracking keyed by reply subject.
|
||||||
// Reference: golang/nats-server/server/consumer.go waitingRequestsPending / proposeWaitingRequest
|
// Reference: golang/nats-server/server/consumer.go waitingRequestsPending / proposeWaitingRequest
|
||||||
private readonly ConcurrentDictionary<string, PullWaitingRequest> _clusterPending =
|
private readonly ConcurrentDictionary<string, PullWaitingRequest> _clusterPending =
|
||||||
@@ -158,11 +155,7 @@ public sealed class PullConsumerEngine
|
|||||||
public async ValueTask<PullFetchBatch> FetchAsync(StreamHandle stream, ConsumerHandle consumer, PullFetchRequest request, CancellationToken ct)
|
public async ValueTask<PullFetchBatch> FetchAsync(StreamHandle stream, ConsumerHandle consumer, PullFetchRequest request, CancellationToken ct)
|
||||||
{
|
{
|
||||||
var batch = Math.Max(request.Batch, 1);
|
var batch = Math.Max(request.Batch, 1);
|
||||||
// Use thread-static buffer to avoid per-fetch List allocation.
|
var messages = new List<StoredMessage>(batch);
|
||||||
// Results are snapshot'd into PullFetchBatch before the buffer is reused.
|
|
||||||
var messages = t_fetchBuf ??= new List<StoredMessage>();
|
|
||||||
messages.Clear();
|
|
||||||
if (messages.Capacity < batch) messages.Capacity = batch;
|
|
||||||
|
|
||||||
// Go: consumer.go — enforce ExpiresMs timeout on pull fetch requests.
|
// Go: consumer.go — enforce ExpiresMs timeout on pull fetch requests.
|
||||||
// When ExpiresMs > 0, create a linked CancellationTokenSource that fires
|
// When ExpiresMs > 0, create a linked CancellationTokenSource that fires
|
||||||
@@ -313,6 +306,7 @@ public sealed class PullConsumerEngine
|
|||||||
await Task.Delay(5, ct).ConfigureAwait(false);
|
await Task.Delay(5, ct).ConfigureAwait(false);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
ct.ThrowIfCancellationRequested();
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -28,6 +28,8 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
// In-memory cache: keyed by sequence number. This is the primary data structure
|
// In-memory cache: keyed by sequence number. This is the primary data structure
|
||||||
// for reads and queries. The blocks are the on-disk persistence layer.
|
// for reads and queries. The blocks are the on-disk persistence layer.
|
||||||
private readonly Dictionary<ulong, StoredMessage> _messages = new();
|
private readonly Dictionary<ulong, StoredMessage> _messages = new();
|
||||||
|
private readonly Dictionary<ulong, StoredMessageIndex> _messageIndexes = new();
|
||||||
|
private readonly Dictionary<string, ulong> _lastSequenceBySubject = new(StringComparer.Ordinal);
|
||||||
|
|
||||||
// Block-based storage: the active (writable) block and sealed blocks.
|
// Block-based storage: the active (writable) block and sealed blocks.
|
||||||
private readonly List<MsgBlock> _blocks = [];
|
private readonly List<MsgBlock> _blocks = [];
|
||||||
@@ -141,20 +143,15 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
// We keep _messages for LoadAsync/RemoveAsync but avoid double payload.ToArray().
|
// We keep _messages for LoadAsync/RemoveAsync but avoid double payload.ToArray().
|
||||||
var persistedPayload = TransformForPersist(payload.Span);
|
var persistedPayload = TransformForPersist(payload.Span);
|
||||||
var storedPayload = _noTransform ? persistedPayload : payload.ToArray();
|
var storedPayload = _noTransform ? persistedPayload : payload.ToArray();
|
||||||
_messages[_last] = new StoredMessage
|
TrackMessage(new StoredMessage
|
||||||
{
|
{
|
||||||
Sequence = _last,
|
Sequence = _last,
|
||||||
Subject = subject,
|
Subject = subject,
|
||||||
Payload = storedPayload,
|
Payload = storedPayload,
|
||||||
TimestampUtc = now,
|
TimestampUtc = now,
|
||||||
};
|
});
|
||||||
_generation++;
|
_generation++;
|
||||||
|
|
||||||
_messageCount++;
|
|
||||||
_totalBytes += (ulong)payload.Length;
|
|
||||||
if (_messageCount == 1)
|
|
||||||
_firstSeq = _last;
|
|
||||||
|
|
||||||
// Go: register TTL only when TTL > 0.
|
// Go: register TTL only when TTL > 0.
|
||||||
if (_options.MaxAgeMs > 0)
|
if (_options.MaxAgeMs > 0)
|
||||||
RegisterTtl(_last, timestamp, (long)_options.MaxAgeMs * 1_000_000L);
|
RegisterTtl(_last, timestamp, (long)_options.MaxAgeMs * 1_000_000L);
|
||||||
@@ -195,11 +192,13 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
|
|
||||||
public ValueTask<StoredMessage?> LoadLastBySubjectAsync(string subject, CancellationToken ct)
|
public ValueTask<StoredMessage?> LoadLastBySubjectAsync(string subject, CancellationToken ct)
|
||||||
{
|
{
|
||||||
var match = _messages.Values
|
if (_lastSequenceBySubject.TryGetValue(subject, out var sequence)
|
||||||
.Where(m => string.Equals(m.Subject, subject, StringComparison.Ordinal))
|
&& _messages.TryGetValue(sequence, out var match))
|
||||||
.OrderByDescending(m => m.Sequence)
|
{
|
||||||
.FirstOrDefault();
|
return ValueTask.FromResult<StoredMessage?>(match);
|
||||||
return ValueTask.FromResult(match);
|
}
|
||||||
|
|
||||||
|
return ValueTask.FromResult<StoredMessage?>(null);
|
||||||
}
|
}
|
||||||
|
|
||||||
public ValueTask<IReadOnlyList<StoredMessage>> ListAsync(CancellationToken ct)
|
public ValueTask<IReadOnlyList<StoredMessage>> ListAsync(CancellationToken ct)
|
||||||
@@ -212,21 +211,11 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
|
|
||||||
public ValueTask<bool> RemoveAsync(ulong sequence, CancellationToken ct)
|
public ValueTask<bool> RemoveAsync(ulong sequence, CancellationToken ct)
|
||||||
{
|
{
|
||||||
if (!_messages.TryGetValue(sequence, out var msg))
|
if (!RemoveTrackedMessage(sequence, preserveHighWaterMark: false))
|
||||||
return ValueTask.FromResult(false);
|
return ValueTask.FromResult(false);
|
||||||
|
|
||||||
_messages.Remove(sequence);
|
|
||||||
_generation++;
|
_generation++;
|
||||||
|
|
||||||
// Incremental state tracking.
|
|
||||||
_messageCount--;
|
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
|
||||||
if (sequence == _firstSeq)
|
|
||||||
_firstSeq = _messages.Count == 0 ? 0UL : _messages.Keys.Min();
|
|
||||||
|
|
||||||
if (sequence == _last)
|
|
||||||
_last = _messages.Count == 0 ? 0UL : _messages.Keys.Max();
|
|
||||||
|
|
||||||
// Soft-delete in the block that contains this sequence.
|
// Soft-delete in the block that contains this sequence.
|
||||||
DeleteInBlock(sequence);
|
DeleteInBlock(sequence);
|
||||||
|
|
||||||
@@ -266,6 +255,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
{
|
{
|
||||||
Sequence = x.Sequence,
|
Sequence = x.Sequence,
|
||||||
Subject = x.Subject,
|
Subject = x.Subject,
|
||||||
|
HeadersBase64 = x.RawHeaders.IsEmpty ? null : Convert.ToBase64String(x.RawHeaders.Span),
|
||||||
PayloadBase64 = Convert.ToBase64String(TransformForPersist(x.Payload.Span)),
|
PayloadBase64 = Convert.ToBase64String(TransformForPersist(x.Payload.Span)),
|
||||||
TimestampUtc = x.TimestampUtc,
|
TimestampUtc = x.TimestampUtc,
|
||||||
})
|
})
|
||||||
@@ -276,10 +266,13 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
public ValueTask RestoreSnapshotAsync(ReadOnlyMemory<byte> snapshot, CancellationToken ct)
|
public ValueTask RestoreSnapshotAsync(ReadOnlyMemory<byte> snapshot, CancellationToken ct)
|
||||||
{
|
{
|
||||||
_messages.Clear();
|
_messages.Clear();
|
||||||
|
_messageIndexes.Clear();
|
||||||
|
_lastSequenceBySubject.Clear();
|
||||||
_last = 0;
|
_last = 0;
|
||||||
_messageCount = 0;
|
_messageCount = 0;
|
||||||
_totalBytes = 0;
|
_totalBytes = 0;
|
||||||
_firstSeq = 0;
|
_firstSeq = 0;
|
||||||
|
_first = 0;
|
||||||
|
|
||||||
// Dispose existing blocks and clean files.
|
// Dispose existing blocks and clean files.
|
||||||
DisposeAllBlocks();
|
DisposeAllBlocks();
|
||||||
@@ -292,25 +285,26 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
{
|
{
|
||||||
foreach (var record in records)
|
foreach (var record in records)
|
||||||
{
|
{
|
||||||
|
var restoredHeaders = string.IsNullOrEmpty(record.HeadersBase64)
|
||||||
|
? ReadOnlyMemory<byte>.Empty
|
||||||
|
: Convert.FromBase64String(record.HeadersBase64);
|
||||||
var restoredPayload = RestorePayload(Convert.FromBase64String(record.PayloadBase64 ?? string.Empty));
|
var restoredPayload = RestorePayload(Convert.FromBase64String(record.PayloadBase64 ?? string.Empty));
|
||||||
var message = new StoredMessage
|
var message = new StoredMessage
|
||||||
{
|
{
|
||||||
Sequence = record.Sequence,
|
Sequence = record.Sequence,
|
||||||
Subject = record.Subject ?? string.Empty,
|
Subject = record.Subject ?? string.Empty,
|
||||||
|
RawHeaders = restoredHeaders,
|
||||||
Payload = restoredPayload,
|
Payload = restoredPayload,
|
||||||
TimestampUtc = record.TimestampUtc,
|
TimestampUtc = record.TimestampUtc,
|
||||||
};
|
};
|
||||||
_messages[record.Sequence] = message;
|
_messages[record.Sequence] = message;
|
||||||
_last = Math.Max(_last, record.Sequence);
|
_last = Math.Max(_last, record.Sequence);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Recompute incremental state from restored messages.
|
|
||||||
_messageCount = (ulong)_messages.Count;
|
|
||||||
_totalBytes = (ulong)_messages.Values.Sum(m => (long)m.Payload.Length);
|
|
||||||
_firstSeq = _messages.Count > 0 ? _messages.Keys.Min() : 0UL;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
RebuildIndexesFromMessages();
|
||||||
|
|
||||||
// Write all messages to fresh blocks.
|
// Write all messages to fresh blocks.
|
||||||
RewriteBlocks();
|
RewriteBlocks();
|
||||||
return ValueTask.CompletedTask;
|
return ValueTask.CompletedTask;
|
||||||
@@ -332,14 +326,10 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
var trimmed = false;
|
var trimmed = false;
|
||||||
while ((ulong)_messages.Count > maxMessages)
|
while ((ulong)_messages.Count > maxMessages)
|
||||||
{
|
{
|
||||||
var first = _messages.Keys.Min();
|
var first = _firstSeq;
|
||||||
if (_messages.TryGetValue(first, out var msg))
|
if (first == 0 || !RemoveTrackedMessage(first, preserveHighWaterMark: true))
|
||||||
{
|
break;
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
|
||||||
_messageCount--;
|
|
||||||
}
|
|
||||||
|
|
||||||
_messages.Remove(first);
|
|
||||||
DeleteInBlock(first);
|
DeleteInBlock(first);
|
||||||
trimmed = true;
|
trimmed = true;
|
||||||
}
|
}
|
||||||
@@ -347,7 +337,6 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
if (!trimmed)
|
if (!trimmed)
|
||||||
return;
|
return;
|
||||||
|
|
||||||
_firstSeq = _messages.Count > 0 ? _messages.Keys.Min() : 0UL;
|
|
||||||
_generation++;
|
_generation++;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -377,36 +366,19 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
var now = DateTime.UtcNow;
|
var now = DateTime.UtcNow;
|
||||||
var timestamp = new DateTimeOffset(now).ToUnixTimeMilliseconds() * 1_000_000L;
|
var timestamp = new DateTimeOffset(now).ToUnixTimeMilliseconds() * 1_000_000L;
|
||||||
|
|
||||||
// Combine headers and payload (headers precede the body in NATS wire format).
|
var headers = hdr is { Length: > 0 } ? hdr : [];
|
||||||
byte[] combined;
|
var payload = msg ?? [];
|
||||||
if (hdr is { Length: > 0 })
|
var persistedPayload = TransformForPersist(payload);
|
||||||
{
|
TrackMessage(new StoredMessage
|
||||||
combined = new byte[hdr.Length + (msg?.Length ?? 0)];
|
|
||||||
hdr.CopyTo(combined, 0);
|
|
||||||
msg?.CopyTo(combined, hdr.Length);
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
combined = msg ?? [];
|
|
||||||
}
|
|
||||||
|
|
||||||
var persistedPayload = TransformForPersist(combined.AsSpan());
|
|
||||||
var stored = new StoredMessage
|
|
||||||
{
|
{
|
||||||
Sequence = _last,
|
Sequence = _last,
|
||||||
Subject = subject,
|
Subject = subject,
|
||||||
Payload = combined,
|
RawHeaders = headers,
|
||||||
|
Payload = payload,
|
||||||
TimestampUtc = now,
|
TimestampUtc = now,
|
||||||
};
|
});
|
||||||
_messages[_last] = stored;
|
|
||||||
_generation++;
|
_generation++;
|
||||||
|
|
||||||
// Incremental state tracking.
|
|
||||||
_messageCount++;
|
|
||||||
_totalBytes += (ulong)combined.Length;
|
|
||||||
if (_messageCount == 1)
|
|
||||||
_firstSeq = _last;
|
|
||||||
|
|
||||||
// Determine effective TTL: per-message ttl (ns) takes priority over MaxAgeMs.
|
// Determine effective TTL: per-message ttl (ns) takes priority over MaxAgeMs.
|
||||||
// Go: filestore.go:6830 — if msg.ttl > 0 use it, else use cfg.MaxAge.
|
// Go: filestore.go:6830 — if msg.ttl > 0 use it, else use cfg.MaxAge.
|
||||||
var effectiveTtlNs = ttl > 0 ? ttl : (_options.MaxAgeMs > 0 ? (long)_options.MaxAgeMs * 1_000_000L : 0L);
|
var effectiveTtlNs = ttl > 0 ? ttl : (_options.MaxAgeMs > 0 ? (long)_options.MaxAgeMs * 1_000_000L : 0L);
|
||||||
@@ -415,16 +387,16 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
EnsureActiveBlock();
|
EnsureActiveBlock();
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
_activeBlock!.WriteAt(_last, subject, ReadOnlyMemory<byte>.Empty, persistedPayload, timestamp);
|
_activeBlock!.WriteAt(_last, subject, headers, persistedPayload, timestamp);
|
||||||
}
|
}
|
||||||
catch (InvalidOperationException)
|
catch (InvalidOperationException)
|
||||||
{
|
{
|
||||||
RotateBlock();
|
RotateBlock();
|
||||||
_activeBlock!.WriteAt(_last, subject, ReadOnlyMemory<byte>.Empty, persistedPayload, timestamp);
|
_activeBlock!.WriteAt(_last, subject, headers, persistedPayload, timestamp);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Go: filestore.go:4443 (setupWriteCache) — record write in bounded cache manager.
|
// Go: filestore.go:4443 (setupWriteCache) — record write in bounded cache manager.
|
||||||
_writeCache.TrackWrite(_activeBlock!.BlockId, persistedPayload.Length);
|
_writeCache.TrackWrite(_activeBlock!.BlockId, headers.Length + persistedPayload.Length);
|
||||||
|
|
||||||
// Signal the background flush loop to coalesce and flush pending writes.
|
// Signal the background flush loop to coalesce and flush pending writes.
|
||||||
_flushSignal.Writer.TryWrite(0);
|
_flushSignal.Writer.TryWrite(0);
|
||||||
@@ -443,11 +415,14 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
{
|
{
|
||||||
var count = (ulong)_messages.Count;
|
var count = (ulong)_messages.Count;
|
||||||
_messages.Clear();
|
_messages.Clear();
|
||||||
|
_messageIndexes.Clear();
|
||||||
|
_lastSequenceBySubject.Clear();
|
||||||
_generation++;
|
_generation++;
|
||||||
_last = 0;
|
_last = 0;
|
||||||
_messageCount = 0;
|
_messageCount = 0;
|
||||||
_totalBytes = 0;
|
_totalBytes = 0;
|
||||||
_firstSeq = 0;
|
_firstSeq = 0;
|
||||||
|
_first = 0;
|
||||||
|
|
||||||
DisposeAllBlocks();
|
DisposeAllBlocks();
|
||||||
CleanBlockFiles();
|
CleanBlockFiles();
|
||||||
@@ -470,42 +445,53 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
if (string.IsNullOrEmpty(subject) && keep == 0 && seq == 0)
|
if (string.IsNullOrEmpty(subject) && keep == 0 && seq == 0)
|
||||||
return Purge();
|
return Purge();
|
||||||
|
|
||||||
// Collect all messages matching the subject (with wildcard support) at or below seq, ordered by sequence.
|
var upperBound = seq == 0 ? _last : Math.Min(seq, _last);
|
||||||
var candidates = _messages.Values
|
if (upperBound == 0)
|
||||||
.Where(m => SubjectMatchesFilter(m.Subject, subject))
|
|
||||||
.Where(m => seq == 0 || m.Sequence <= seq)
|
|
||||||
.OrderBy(m => m.Sequence)
|
|
||||||
.ToList();
|
|
||||||
|
|
||||||
if (candidates.Count == 0)
|
|
||||||
return 0;
|
return 0;
|
||||||
|
|
||||||
// Keep the newest `keep` messages; purge the rest.
|
ulong candidateCount = 0;
|
||||||
var toRemove = keep > 0 && (ulong)candidates.Count > keep
|
for (var current = _messageCount == 0 ? 0UL : _first; current != 0 && current <= upperBound; current++)
|
||||||
? candidates.Take(candidates.Count - (int)keep).ToList()
|
|
||||||
: (keep == 0 ? candidates : []);
|
|
||||||
|
|
||||||
if (toRemove.Count == 0)
|
|
||||||
return 0;
|
|
||||||
|
|
||||||
foreach (var msg in toRemove)
|
|
||||||
{
|
{
|
||||||
_messages.Remove(msg.Sequence);
|
if (!_messages.TryGetValue(current, out var message))
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
continue;
|
||||||
_messageCount--;
|
|
||||||
DeleteInBlock(msg.Sequence);
|
if (!SubjectMatchesFilter(message.Subject, subject))
|
||||||
|
continue;
|
||||||
|
|
||||||
|
candidateCount++;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (candidateCount == 0)
|
||||||
|
return 0;
|
||||||
|
|
||||||
|
var targetRemoveCount = keep > 0
|
||||||
|
? (candidateCount > keep ? candidateCount - keep : 0UL)
|
||||||
|
: candidateCount;
|
||||||
|
|
||||||
|
if (targetRemoveCount == 0)
|
||||||
|
return 0;
|
||||||
|
|
||||||
|
ulong removed = 0;
|
||||||
|
for (var current = _messageCount == 0 ? 0UL : _first; current != 0 && current <= upperBound && removed < targetRemoveCount; current++)
|
||||||
|
{
|
||||||
|
if (!_messages.TryGetValue(current, out var message))
|
||||||
|
continue;
|
||||||
|
|
||||||
|
if (!SubjectMatchesFilter(message.Subject, subject))
|
||||||
|
continue;
|
||||||
|
|
||||||
|
if (RemoveTrackedMessage(current, preserveHighWaterMark: true))
|
||||||
|
{
|
||||||
|
DeleteInBlock(current);
|
||||||
|
removed++;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (removed == 0)
|
||||||
|
return 0;
|
||||||
|
|
||||||
_generation++;
|
_generation++;
|
||||||
_firstSeq = _messages.Count > 0 ? _messages.Keys.Min() : 0UL;
|
return removed;
|
||||||
|
|
||||||
// Update _last if required.
|
|
||||||
if (_messages.Count == 0)
|
|
||||||
_last = 0;
|
|
||||||
else if (!_messages.ContainsKey(_last))
|
|
||||||
_last = _messages.Keys.Max();
|
|
||||||
|
|
||||||
return (ulong)toRemove.Count;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -524,13 +510,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
|
|
||||||
foreach (var s in toRemove)
|
foreach (var s in toRemove)
|
||||||
{
|
{
|
||||||
if (_messages.TryGetValue(s, out var msg))
|
RemoveTrackedMessage(s, preserveHighWaterMark: true);
|
||||||
{
|
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
|
||||||
_messageCount--;
|
|
||||||
}
|
|
||||||
|
|
||||||
_messages.Remove(s);
|
|
||||||
DeleteInBlock(s);
|
DeleteInBlock(s);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -545,11 +525,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
}
|
}
|
||||||
else
|
else
|
||||||
{
|
{
|
||||||
if (!_messages.ContainsKey(_last))
|
_first = _firstSeq;
|
||||||
_last = _messages.Keys.Max();
|
|
||||||
// Update _first to reflect the real first message.
|
|
||||||
_first = _messages.Keys.Min();
|
|
||||||
_firstSeq = _first;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return (ulong)toRemove.Length;
|
return (ulong)toRemove.Length;
|
||||||
@@ -566,11 +542,14 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
{
|
{
|
||||||
// Truncate to nothing.
|
// Truncate to nothing.
|
||||||
_messages.Clear();
|
_messages.Clear();
|
||||||
|
_messageIndexes.Clear();
|
||||||
|
_lastSequenceBySubject.Clear();
|
||||||
_generation++;
|
_generation++;
|
||||||
_last = 0;
|
_last = 0;
|
||||||
_messageCount = 0;
|
_messageCount = 0;
|
||||||
_totalBytes = 0;
|
_totalBytes = 0;
|
||||||
_firstSeq = 0;
|
_firstSeq = 0;
|
||||||
|
_first = 0;
|
||||||
DisposeAllBlocks();
|
DisposeAllBlocks();
|
||||||
CleanBlockFiles();
|
CleanBlockFiles();
|
||||||
return;
|
return;
|
||||||
@@ -579,13 +558,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
var toRemove = _messages.Keys.Where(k => k > seq).ToArray();
|
var toRemove = _messages.Keys.Where(k => k > seq).ToArray();
|
||||||
foreach (var s in toRemove)
|
foreach (var s in toRemove)
|
||||||
{
|
{
|
||||||
if (_messages.TryGetValue(s, out var msg))
|
RemoveTrackedMessage(s, preserveHighWaterMark: false);
|
||||||
{
|
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
|
||||||
_messageCount--;
|
|
||||||
}
|
|
||||||
|
|
||||||
_messages.Remove(s);
|
|
||||||
DeleteInBlock(s);
|
DeleteInBlock(s);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -594,8 +567,12 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
|
|
||||||
// Update _last to the new highest existing sequence (or seq if it exists,
|
// Update _last to the new highest existing sequence (or seq if it exists,
|
||||||
// or the highest below seq).
|
// or the highest below seq).
|
||||||
_last = _messages.Count == 0 ? 0 : _messages.Keys.Max();
|
if (_messageCount == 0)
|
||||||
_firstSeq = _messages.Count > 0 ? _messages.Keys.Min() : 0UL;
|
{
|
||||||
|
_last = 0;
|
||||||
|
_first = 0;
|
||||||
|
_firstSeq = 0;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -865,6 +842,119 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private void TrackMessage(StoredMessage message)
|
||||||
|
{
|
||||||
|
_messages[message.Sequence] = message;
|
||||||
|
_messageIndexes[message.Sequence] = message.ToIndex();
|
||||||
|
_lastSequenceBySubject[message.Subject] = message.Sequence;
|
||||||
|
_messageCount++;
|
||||||
|
_totalBytes += (ulong)(message.RawHeaders.Length + message.Payload.Length);
|
||||||
|
|
||||||
|
if (_messageCount == 1)
|
||||||
|
{
|
||||||
|
_first = message.Sequence;
|
||||||
|
_firstSeq = message.Sequence;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private bool RemoveTrackedMessage(ulong sequence, bool preserveHighWaterMark)
|
||||||
|
{
|
||||||
|
if (!_messages.Remove(sequence, out var message))
|
||||||
|
return false;
|
||||||
|
|
||||||
|
_messageIndexes.Remove(sequence);
|
||||||
|
_messageCount--;
|
||||||
|
_totalBytes -= (ulong)(message.RawHeaders.Length + message.Payload.Length);
|
||||||
|
UpdateLastSequenceForSubject(message.Subject, sequence);
|
||||||
|
|
||||||
|
if (_messageCount == 0)
|
||||||
|
{
|
||||||
|
if (preserveHighWaterMark)
|
||||||
|
_first = _last + 1;
|
||||||
|
else
|
||||||
|
_last = 0;
|
||||||
|
|
||||||
|
_firstSeq = 0;
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (sequence == _firstSeq)
|
||||||
|
AdvanceFirstSequence(sequence + 1);
|
||||||
|
|
||||||
|
if (!preserveHighWaterMark && sequence == _last)
|
||||||
|
_last = FindPreviousLiveSequence(sequence);
|
||||||
|
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
private void AdvanceFirstSequence(ulong start)
|
||||||
|
{
|
||||||
|
var candidate = start;
|
||||||
|
while (!_messageIndexes.ContainsKey(candidate) && candidate <= _last)
|
||||||
|
candidate++;
|
||||||
|
|
||||||
|
if (candidate <= _last)
|
||||||
|
{
|
||||||
|
_first = candidate;
|
||||||
|
_firstSeq = candidate;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
_first = _last + 1;
|
||||||
|
_firstSeq = 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
private ulong FindPreviousLiveSequence(ulong startExclusive)
|
||||||
|
{
|
||||||
|
if (_messageCount == 0 || startExclusive == 0)
|
||||||
|
return 0;
|
||||||
|
|
||||||
|
for (var seq = startExclusive - 1; ; seq--)
|
||||||
|
{
|
||||||
|
if (_messageIndexes.ContainsKey(seq))
|
||||||
|
return seq;
|
||||||
|
|
||||||
|
if (seq == 0)
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void UpdateLastSequenceForSubject(string subject, ulong removedSequence)
|
||||||
|
{
|
||||||
|
if (!_lastSequenceBySubject.TryGetValue(subject, out var currentLast) || currentLast != removedSequence)
|
||||||
|
return;
|
||||||
|
|
||||||
|
for (var seq = removedSequence - 1; ; seq--)
|
||||||
|
{
|
||||||
|
if (_messageIndexes.TryGetValue(seq, out var candidate) && string.Equals(candidate.Subject, subject, StringComparison.Ordinal))
|
||||||
|
{
|
||||||
|
_lastSequenceBySubject[subject] = seq;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (seq == 0)
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
_lastSequenceBySubject.Remove(subject);
|
||||||
|
}
|
||||||
|
|
||||||
|
private void RebuildIndexesFromMessages()
|
||||||
|
{
|
||||||
|
_messageIndexes.Clear();
|
||||||
|
_lastSequenceBySubject.Clear();
|
||||||
|
_messageCount = 0;
|
||||||
|
_totalBytes = 0;
|
||||||
|
_firstSeq = 0;
|
||||||
|
_first = 0;
|
||||||
|
|
||||||
|
foreach (var message in _messages.OrderBy(kv => kv.Key).Select(kv => kv.Value))
|
||||||
|
TrackMessage(message);
|
||||||
|
|
||||||
|
if (_messageCount == 0 && _last > 0)
|
||||||
|
_first = _last + 1;
|
||||||
|
}
|
||||||
|
|
||||||
// -------------------------------------------------------------------------
|
// -------------------------------------------------------------------------
|
||||||
// Subject matching helper
|
// Subject matching helper
|
||||||
// -------------------------------------------------------------------------
|
// -------------------------------------------------------------------------
|
||||||
@@ -1056,9 +1146,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
CleanBlockFiles();
|
CleanBlockFiles();
|
||||||
|
|
||||||
_last = _messages.Count == 0 ? 0UL : _messages.Keys.Max();
|
_last = _messages.Count == 0 ? 0UL : _messages.Keys.Max();
|
||||||
_firstSeq = _messages.Count > 0 ? _messages.Keys.Min() : 0UL;
|
RebuildIndexesFromMessages();
|
||||||
_messageCount = (ulong)_messages.Count;
|
|
||||||
_totalBytes = (ulong)_messages.Values.Sum(m => (long)m.Payload.Length);
|
|
||||||
|
|
||||||
foreach (var message in _messages.OrderBy(kv => kv.Key).Select(kv => kv.Value))
|
foreach (var message in _messages.OrderBy(kv => kv.Key).Select(kv => kv.Value))
|
||||||
{
|
{
|
||||||
@@ -1068,12 +1156,12 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
EnsureActiveBlock();
|
EnsureActiveBlock();
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
_activeBlock!.WriteAt(message.Sequence, message.Subject, ReadOnlyMemory<byte>.Empty, persistedPayload, timestamp);
|
_activeBlock!.WriteAt(message.Sequence, message.Subject, message.RawHeaders, persistedPayload, timestamp);
|
||||||
}
|
}
|
||||||
catch (InvalidOperationException)
|
catch (InvalidOperationException)
|
||||||
{
|
{
|
||||||
RotateBlock();
|
RotateBlock();
|
||||||
_activeBlock!.WriteAt(message.Sequence, message.Subject, ReadOnlyMemory<byte>.Empty, persistedPayload, timestamp);
|
_activeBlock!.WriteAt(message.Sequence, message.Subject, message.RawHeaders, persistedPayload, timestamp);
|
||||||
}
|
}
|
||||||
|
|
||||||
if (_activeBlock!.IsSealed)
|
if (_activeBlock!.IsSealed)
|
||||||
@@ -1168,15 +1256,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Sync _first from _messages; if empty, set to _last+1 (watermark).
|
// Sync _first from _messages; if empty, set to _last+1 (watermark).
|
||||||
if (_messages.Count > 0)
|
RebuildIndexesFromMessages();
|
||||||
_first = _messages.Keys.Min();
|
|
||||||
else if (_last > 0)
|
|
||||||
_first = _last + 1;
|
|
||||||
|
|
||||||
// Recompute incremental state from recovered messages.
|
|
||||||
_messageCount = (ulong)_messages.Count;
|
|
||||||
_totalBytes = (ulong)_messages.Values.Sum(m => (long)m.Payload.Length);
|
|
||||||
_firstSeq = _messages.Count > 0 ? _messages.Keys.Min() : 0UL;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -1207,6 +1287,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
{
|
{
|
||||||
Sequence = record.Sequence,
|
Sequence = record.Sequence,
|
||||||
Subject = record.Subject,
|
Subject = record.Subject,
|
||||||
|
RawHeaders = record.Headers,
|
||||||
Payload = originalPayload,
|
Payload = originalPayload,
|
||||||
TimestampUtc = DateTimeOffset.FromUnixTimeMilliseconds(record.Timestamp / 1_000_000L).UtcDateTime,
|
TimestampUtc = DateTimeOffset.FromUnixTimeMilliseconds(record.Timestamp / 1_000_000L).UtcDateTime,
|
||||||
};
|
};
|
||||||
@@ -1370,20 +1451,10 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
// Reference: golang/nats-server/server/filestore.go:expireMsgs — dmap-based removal.
|
// Reference: golang/nats-server/server/filestore.go:expireMsgs — dmap-based removal.
|
||||||
foreach (var seq in expired)
|
foreach (var seq in expired)
|
||||||
{
|
{
|
||||||
if (_messages.Remove(seq, out var msg))
|
RemoveTrackedMessage(seq, preserveHighWaterMark: true);
|
||||||
{
|
|
||||||
_messageCount--;
|
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
|
||||||
}
|
|
||||||
|
|
||||||
DeleteInBlock(seq);
|
DeleteInBlock(seq);
|
||||||
}
|
}
|
||||||
|
|
||||||
if (_messages.Count > 0)
|
|
||||||
_firstSeq = _messages.Keys.Min();
|
|
||||||
else
|
|
||||||
_firstSeq = 0;
|
|
||||||
|
|
||||||
_generation++;
|
_generation++;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1407,14 +1478,8 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
|
|
||||||
foreach (var sequence in expired)
|
foreach (var sequence in expired)
|
||||||
{
|
{
|
||||||
if (_messages.Remove(sequence, out var msg))
|
RemoveTrackedMessage(sequence, preserveHighWaterMark: true);
|
||||||
{
|
|
||||||
_messageCount--;
|
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
_firstSeq = _messages.Count > 0 ? _messages.Keys.Min() : 0UL;
|
|
||||||
RewriteBlocks();
|
RewriteBlocks();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1676,27 +1741,13 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
public bool RemoveMsg(ulong seq)
|
public bool RemoveMsg(ulong seq)
|
||||||
{
|
{
|
||||||
if (!_messages.Remove(seq, out var msg))
|
if (!RemoveTrackedMessage(seq, preserveHighWaterMark: true))
|
||||||
return false;
|
return false;
|
||||||
|
|
||||||
_generation++;
|
_generation++;
|
||||||
_messageCount--;
|
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
|
||||||
|
|
||||||
// Go: filestore.go — LastSeq (lmb.last.seq) is a high-water mark and is
|
// Go: filestore.go — LastSeq (lmb.last.seq) is a high-water mark and is
|
||||||
// never decremented on removal. Only FirstSeq advances when the first
|
// never decremented on removal. Only FirstSeq advances when the first
|
||||||
// live message is removed.
|
// live message is removed.
|
||||||
if (_messages.Count == 0)
|
|
||||||
{
|
|
||||||
_first = _last + 1; // All gone — next first would be after last
|
|
||||||
_firstSeq = 0;
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
_first = _messages.Keys.Min();
|
|
||||||
_firstSeq = _first;
|
|
||||||
}
|
|
||||||
|
|
||||||
DeleteInBlock(seq);
|
DeleteInBlock(seq);
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
@@ -1709,23 +1760,10 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
public bool EraseMsg(ulong seq)
|
public bool EraseMsg(ulong seq)
|
||||||
{
|
{
|
||||||
if (!_messages.Remove(seq, out var msg))
|
if (!RemoveTrackedMessage(seq, preserveHighWaterMark: true))
|
||||||
return false;
|
return false;
|
||||||
|
|
||||||
_generation++;
|
_generation++;
|
||||||
_messageCount--;
|
|
||||||
_totalBytes -= (ulong)msg.Payload.Length;
|
|
||||||
|
|
||||||
if (_messages.Count == 0)
|
|
||||||
{
|
|
||||||
_first = _last + 1;
|
|
||||||
_firstSeq = 0;
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
_first = _messages.Keys.Min();
|
|
||||||
_firstSeq = _first;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Secure erase: overwrite payload bytes with random data before marking deleted.
|
// Secure erase: overwrite payload bytes with random data before marking deleted.
|
||||||
// Reference: golang/nats-server/server/filestore.go:5890 (eraseMsg).
|
// Reference: golang/nats-server/server/filestore.go:5890 (eraseMsg).
|
||||||
@@ -1813,6 +1851,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
sm ??= new StoreMsg();
|
sm ??= new StoreMsg();
|
||||||
sm.Clear();
|
sm.Clear();
|
||||||
sm.Subject = stored.Subject;
|
sm.Subject = stored.Subject;
|
||||||
|
sm.Header = stored.RawHeaders.Length > 0 ? stored.RawHeaders.ToArray() : null;
|
||||||
sm.Data = stored.Payload.Length > 0 ? stored.Payload.ToArray() : null;
|
sm.Data = stored.Payload.Length > 0 ? stored.Payload.ToArray() : null;
|
||||||
sm.Sequence = stored.Sequence;
|
sm.Sequence = stored.Sequence;
|
||||||
sm.Timestamp = new DateTimeOffset(stored.TimestampUtc).ToUnixTimeMilliseconds() * 1_000_000L;
|
sm.Timestamp = new DateTimeOffset(stored.TimestampUtc).ToUnixTimeMilliseconds() * 1_000_000L;
|
||||||
@@ -1833,7 +1872,9 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
sm ??= new StoreMsg();
|
sm ??= new StoreMsg();
|
||||||
sm.Clear();
|
sm.Clear();
|
||||||
sm.Subject = record.Subject;
|
sm.Subject = record.Subject;
|
||||||
sm.Data = record.Payload.Length > 0 ? record.Payload.ToArray() : null;
|
sm.Header = record.Headers.Length > 0 ? record.Headers.ToArray() : null;
|
||||||
|
var originalPayload = RestorePayload(record.Payload.Span);
|
||||||
|
sm.Data = originalPayload.Length > 0 ? originalPayload : null;
|
||||||
sm.Sequence = record.Sequence;
|
sm.Sequence = record.Sequence;
|
||||||
sm.Timestamp = record.Timestamp;
|
sm.Timestamp = record.Timestamp;
|
||||||
return sm;
|
return sm;
|
||||||
@@ -1887,6 +1928,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
sm ??= new StoreMsg();
|
sm ??= new StoreMsg();
|
||||||
sm.Clear();
|
sm.Clear();
|
||||||
sm.Subject = match.Subject;
|
sm.Subject = match.Subject;
|
||||||
|
sm.Header = match.RawHeaders.Length > 0 ? match.RawHeaders.ToArray() : null;
|
||||||
sm.Data = match.Payload.Length > 0 ? match.Payload.ToArray() : null;
|
sm.Data = match.Payload.Length > 0 ? match.Payload.ToArray() : null;
|
||||||
sm.Sequence = match.Sequence;
|
sm.Sequence = match.Sequence;
|
||||||
sm.Timestamp = new DateTimeOffset(match.TimestampUtc).ToUnixTimeMilliseconds() * 1_000_000L;
|
sm.Timestamp = new DateTimeOffset(match.TimestampUtc).ToUnixTimeMilliseconds() * 1_000_000L;
|
||||||
@@ -1917,6 +1959,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
sm ??= new StoreMsg();
|
sm ??= new StoreMsg();
|
||||||
sm.Clear();
|
sm.Clear();
|
||||||
sm.Subject = found.Value.Subject;
|
sm.Subject = found.Value.Subject;
|
||||||
|
sm.Header = found.Value.RawHeaders.Length > 0 ? found.Value.RawHeaders.ToArray() : null;
|
||||||
sm.Data = found.Value.Payload.Length > 0 ? found.Value.Payload.ToArray() : null;
|
sm.Data = found.Value.Payload.Length > 0 ? found.Value.Payload.ToArray() : null;
|
||||||
sm.Sequence = found.Key;
|
sm.Sequence = found.Key;
|
||||||
sm.Timestamp = new DateTimeOffset(found.Value.TimestampUtc).ToUnixTimeMilliseconds() * 1_000_000L;
|
sm.Timestamp = new DateTimeOffset(found.Value.TimestampUtc).ToUnixTimeMilliseconds() * 1_000_000L;
|
||||||
@@ -2041,20 +2084,9 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
if (_stopped)
|
if (_stopped)
|
||||||
throw new ObjectDisposedException(nameof(FileStore), "Store has been stopped.");
|
throw new ObjectDisposedException(nameof(FileStore), "Store has been stopped.");
|
||||||
|
|
||||||
// Combine headers and payload, same as StoreMsg.
|
var headers = hdr is { Length: > 0 } ? hdr : [];
|
||||||
byte[] combined;
|
var payload = msg ?? [];
|
||||||
if (hdr is { Length: > 0 })
|
var persistedPayload = TransformForPersist(payload);
|
||||||
{
|
|
||||||
combined = new byte[hdr.Length + msg.Length];
|
|
||||||
hdr.CopyTo(combined, 0);
|
|
||||||
msg.CopyTo(combined, hdr.Length);
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
combined = msg;
|
|
||||||
}
|
|
||||||
|
|
||||||
var persistedPayload = TransformForPersist(combined.AsSpan());
|
|
||||||
// Recover UTC DateTime from caller-supplied Unix nanosecond timestamp.
|
// Recover UTC DateTime from caller-supplied Unix nanosecond timestamp.
|
||||||
var storedUtc = DateTimeOffset.FromUnixTimeMilliseconds(ts / 1_000_000L).UtcDateTime;
|
var storedUtc = DateTimeOffset.FromUnixTimeMilliseconds(ts / 1_000_000L).UtcDateTime;
|
||||||
|
|
||||||
@@ -2062,10 +2094,11 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
{
|
{
|
||||||
Sequence = seq,
|
Sequence = seq,
|
||||||
Subject = subject,
|
Subject = subject,
|
||||||
Payload = combined,
|
RawHeaders = headers,
|
||||||
|
Payload = payload,
|
||||||
TimestampUtc = storedUtc,
|
TimestampUtc = storedUtc,
|
||||||
};
|
};
|
||||||
_messages[seq] = stored;
|
TrackMessage(stored);
|
||||||
_generation++;
|
_generation++;
|
||||||
|
|
||||||
// Go: update _last to the high-water mark — do not decrement.
|
// Go: update _last to the high-water mark — do not decrement.
|
||||||
@@ -2078,16 +2111,16 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
EnsureActiveBlock();
|
EnsureActiveBlock();
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
_activeBlock!.WriteAt(seq, subject, ReadOnlyMemory<byte>.Empty, persistedPayload, ts);
|
_activeBlock!.WriteAt(seq, subject, headers, persistedPayload, ts);
|
||||||
}
|
}
|
||||||
catch (InvalidOperationException)
|
catch (InvalidOperationException)
|
||||||
{
|
{
|
||||||
RotateBlock();
|
RotateBlock();
|
||||||
_activeBlock!.WriteAt(seq, subject, ReadOnlyMemory<byte>.Empty, persistedPayload, ts);
|
_activeBlock!.WriteAt(seq, subject, headers, persistedPayload, ts);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Go: filestore.go:4443 (setupWriteCache) — record write in bounded cache manager.
|
// Go: filestore.go:4443 (setupWriteCache) — record write in bounded cache manager.
|
||||||
_writeCache.TrackWrite(_activeBlock!.BlockId, persistedPayload.Length);
|
_writeCache.TrackWrite(_activeBlock!.BlockId, headers.Length + persistedPayload.Length);
|
||||||
|
|
||||||
if (_activeBlock!.IsSealed)
|
if (_activeBlock!.IsSealed)
|
||||||
RotateBlock();
|
RotateBlock();
|
||||||
@@ -2113,6 +2146,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
sm ??= new StoreMsg();
|
sm ??= new StoreMsg();
|
||||||
sm.Clear();
|
sm.Clear();
|
||||||
sm.Subject = stored.Subject;
|
sm.Subject = stored.Subject;
|
||||||
|
sm.Header = stored.RawHeaders.Length > 0 ? stored.RawHeaders.ToArray() : null;
|
||||||
sm.Data = stored.Payload.Length > 0 ? stored.Payload.ToArray() : null;
|
sm.Data = stored.Payload.Length > 0 ? stored.Payload.ToArray() : null;
|
||||||
sm.Sequence = stored.Sequence;
|
sm.Sequence = stored.Sequence;
|
||||||
sm.Timestamp = new DateTimeOffset(stored.TimestampUtc).ToUnixTimeMilliseconds() * 1_000_000L;
|
sm.Timestamp = new DateTimeOffset(stored.TimestampUtc).ToUnixTimeMilliseconds() * 1_000_000L;
|
||||||
@@ -2365,6 +2399,7 @@ public sealed class FileStore : IStreamStore, IAsyncDisposable, IDisposable
|
|||||||
{
|
{
|
||||||
public ulong Sequence { get; init; }
|
public ulong Sequence { get; init; }
|
||||||
public string? Subject { get; init; }
|
public string? Subject { get; init; }
|
||||||
|
public string? HeadersBase64 { get; init; }
|
||||||
public string? PayloadBase64 { get; init; }
|
public string? PayloadBase64 { get; init; }
|
||||||
public DateTime TimestampUtc { get; init; }
|
public DateTime TimestampUtc { get; init; }
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ public sealed class StoredMessage
|
|||||||
public ulong Sequence { get; init; }
|
public ulong Sequence { get; init; }
|
||||||
public string Subject { get; init; } = string.Empty;
|
public string Subject { get; init; } = string.Empty;
|
||||||
public ReadOnlyMemory<byte> Payload { get; init; }
|
public ReadOnlyMemory<byte> Payload { get; init; }
|
||||||
|
internal ReadOnlyMemory<byte> RawHeaders { get; init; }
|
||||||
public DateTime TimestampUtc { get; init; } = DateTime.UtcNow;
|
public DateTime TimestampUtc { get; init; } = DateTime.UtcNow;
|
||||||
public string? Account { get; init; }
|
public string? Account { get; init; }
|
||||||
public bool Redelivered { get; init; }
|
public bool Redelivered { get; init; }
|
||||||
@@ -18,4 +19,7 @@ public sealed class StoredMessage
|
|||||||
/// Convenience accessor for the Nats-Msg-Id header value, used by source deduplication.
|
/// Convenience accessor for the Nats-Msg-Id header value, used by source deduplication.
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public string? MsgId => Headers is not null && Headers.TryGetValue("Nats-Msg-Id", out var id) ? id : null;
|
public string? MsgId => Headers is not null && Headers.TryGetValue("Nats-Msg-Id", out var id) ? id : null;
|
||||||
|
|
||||||
|
internal StoredMessageIndex ToIndex()
|
||||||
|
=> new(Sequence, Subject, Payload.Length, TimestampUtc);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,7 @@
|
|||||||
|
namespace NATS.Server.JetStream.Storage;
|
||||||
|
|
||||||
|
public readonly record struct StoredMessageIndex(
|
||||||
|
ulong Sequence,
|
||||||
|
string Subject,
|
||||||
|
int PayloadLength,
|
||||||
|
DateTime TimestampUtc);
|
||||||
@@ -29,6 +29,9 @@ public sealed class ConnzHandler(NatsServer server)
|
|||||||
{
|
{
|
||||||
var clients = server.GetClients().ToArray();
|
var clients = server.GetClients().ToArray();
|
||||||
connInfos.AddRange(clients.Select(c => BuildConnInfo(c, now, opts)));
|
connInfos.AddRange(clients.Select(c => BuildConnInfo(c, now, opts)));
|
||||||
|
|
||||||
|
// Include MQTT adapter connections
|
||||||
|
connInfos.AddRange(server.GetMqttAdapters().Select(a => BuildMqttConnInfo(a, now)));
|
||||||
}
|
}
|
||||||
|
|
||||||
// Collect closed connections from the ring buffer
|
// Collect closed connections from the ring buffer
|
||||||
@@ -254,6 +257,21 @@ public sealed class ConnzHandler(NatsServer server)
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static ConnInfo BuildMqttConnInfo(Mqtt.MqttNatsClientAdapter adapter, DateTime now)
|
||||||
|
{
|
||||||
|
return new ConnInfo
|
||||||
|
{
|
||||||
|
Cid = adapter.Id,
|
||||||
|
Kind = "Client",
|
||||||
|
Type = "mqtt",
|
||||||
|
Start = now, // MQTT adapters don't track start time yet
|
||||||
|
LastActivity = now,
|
||||||
|
NumSubs = (uint)adapter.Subscriptions.Count,
|
||||||
|
Account = adapter.Account?.Name ?? "",
|
||||||
|
MqttClient = adapter.MqttClientId,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
private static ConnzOptions ParseQueryParams(HttpContext ctx)
|
private static ConnzOptions ParseQueryParams(HttpContext ctx)
|
||||||
{
|
{
|
||||||
var q = ctx.Request.Query;
|
var q = ctx.Request.Query;
|
||||||
|
|||||||
@@ -23,10 +23,18 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
private bool _connected;
|
private bool _connected;
|
||||||
private bool _willCleared;
|
private bool _willCleared;
|
||||||
private MqttConnectInfo _connectInfo;
|
private MqttConnectInfo _connectInfo;
|
||||||
|
private readonly Dictionary<string, string> _topicToSid = new(StringComparer.Ordinal);
|
||||||
|
private int _nextSid;
|
||||||
|
|
||||||
/// <summary>Auth result after successful CONNECT (populated for AuthService path).</summary>
|
/// <summary>Auth result after successful CONNECT (populated for AuthService path).</summary>
|
||||||
public AuthResult? AuthResult { get; private set; }
|
public AuthResult? AuthResult { get; private set; }
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The adapter that bridges this MQTT connection to the NATS SubList routing.
|
||||||
|
/// Created on CONNECT when running with a NatsServer router; null in test-only mode.
|
||||||
|
/// </summary>
|
||||||
|
public MqttNatsClientAdapter? Adapter { get; private set; }
|
||||||
|
|
||||||
public string ClientId => _clientId;
|
public string ClientId => _clientId;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -142,11 +150,7 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
// Publish will message if not cleanly disconnected
|
// Publish will message if not cleanly disconnected
|
||||||
if (_connected && !_willCleared && _connectInfo.WillTopic != null)
|
if (_connected && !_willCleared && _connectInfo.WillTopic != null)
|
||||||
{
|
{
|
||||||
await _listener.PublishAsync(
|
RoutePublish(_connectInfo.WillTopic, _connectInfo.WillMessage ?? []);
|
||||||
_connectInfo.WillTopic,
|
|
||||||
Encoding.UTF8.GetString(_connectInfo.WillMessage ?? []),
|
|
||||||
this,
|
|
||||||
CancellationToken.None);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -276,6 +280,15 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
|
|
||||||
AuthResult = authResult;
|
AuthResult = authResult;
|
||||||
|
|
||||||
|
// Create MqttNatsClientAdapter for cross-protocol routing (when running with NatsServer)
|
||||||
|
if (_listener.AllocateClientId != null)
|
||||||
|
{
|
||||||
|
var adapterId = _listener.AllocateClientId();
|
||||||
|
Adapter = new MqttNatsClientAdapter(this, adapterId);
|
||||||
|
Adapter.Account = _listener.ResolveAccount?.Invoke(authResult.AccountName);
|
||||||
|
_listener.RegisterMqttAdapter(Adapter);
|
||||||
|
}
|
||||||
|
|
||||||
// Duplicate client-id takeover
|
// Duplicate client-id takeover
|
||||||
_listener.TakeoverExistingConnection(_clientId, this);
|
_listener.TakeoverExistingConnection(_clientId, this);
|
||||||
|
|
||||||
@@ -313,20 +326,18 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
switch (publishInfo.QoS)
|
switch (publishInfo.QoS)
|
||||||
{
|
{
|
||||||
case 0:
|
case 0:
|
||||||
await _listener.PublishAsync(publishInfo.Topic,
|
RoutePublish(publishInfo.Topic, publishInfo.Payload);
|
||||||
Encoding.UTF8.GetString(publishInfo.Payload.Span), this, ct);
|
|
||||||
break;
|
break;
|
||||||
|
|
||||||
case 1:
|
case 1:
|
||||||
_listener.RecordPendingPublish(_clientId, publishInfo.PacketId, publishInfo.Topic,
|
_listener.RecordPendingPublish(_clientId, publishInfo.PacketId, publishInfo.Topic,
|
||||||
Encoding.UTF8.GetString(publishInfo.Payload.Span));
|
Encoding.UTF8.GetString(publishInfo.Payload.Span));
|
||||||
await WriteBinaryAsync(MqttPacketWriter.WritePubAck(publishInfo.PacketId), ct);
|
await WriteBinaryAsync(MqttPacketWriter.WritePubAck(publishInfo.PacketId), ct);
|
||||||
await _listener.PublishAsync(publishInfo.Topic,
|
RoutePublish(publishInfo.Topic, publishInfo.Payload);
|
||||||
Encoding.UTF8.GetString(publishInfo.Payload.Span), this, ct);
|
|
||||||
break;
|
break;
|
||||||
|
|
||||||
case 2:
|
case 2:
|
||||||
// QoS 2 step 1: store and send PUBREC
|
// QoS 2 step 1: store and send PUBREC (delivery deferred to PUBREL)
|
||||||
_listener.RecordPendingPublish(_clientId, publishInfo.PacketId, publishInfo.Topic,
|
_listener.RecordPendingPublish(_clientId, publishInfo.PacketId, publishInfo.Topic,
|
||||||
Encoding.UTF8.GetString(publishInfo.Payload.Span));
|
Encoding.UTF8.GetString(publishInfo.Payload.Span));
|
||||||
await WriteBinaryAsync(MqttPacketWriter.WritePubRec(publishInfo.PacketId), ct);
|
await WriteBinaryAsync(MqttPacketWriter.WritePubRec(publishInfo.PacketId), ct);
|
||||||
@@ -343,6 +354,24 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Routes a published message through the NATS SubList (cross-protocol) if a router
|
||||||
|
/// and adapter are available, otherwise falls back to MQTT-only fan-out.
|
||||||
|
/// </summary>
|
||||||
|
private void RoutePublish(string mqttTopic, ReadOnlyMemory<byte> payload)
|
||||||
|
{
|
||||||
|
if (_listener.Router != null && Adapter != null)
|
||||||
|
{
|
||||||
|
var natsSubject = MqttTopicMapper.MqttToNats(mqttTopic);
|
||||||
|
_listener.Router.ProcessMessage(natsSubject, null, default, payload, Adapter);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
// Test-only fallback: MQTT-only fan-out
|
||||||
|
_ = _listener.PublishAsync(mqttTopic, Encoding.UTF8.GetString(payload.Span), this, CancellationToken.None);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private void HandlePubAck(MqttControlPacket packet)
|
private void HandlePubAck(MqttControlPacket packet)
|
||||||
{
|
{
|
||||||
if (packet.Payload.Length < 2) return;
|
if (packet.Payload.Length < 2) return;
|
||||||
@@ -362,7 +391,13 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
if (packet.Payload.Length < 2) return;
|
if (packet.Payload.Length < 2) return;
|
||||||
var packetId = (ushort)((packet.Payload.Span[0] << 8) | packet.Payload.Span[1]);
|
var packetId = (ushort)((packet.Payload.Span[0] << 8) | packet.Payload.Span[1]);
|
||||||
|
|
||||||
// QoS 2 step 2: deliver the stored message and send PUBCOMP
|
// QoS 2 step 2: deliver the stored message, then ack and send PUBCOMP
|
||||||
|
var pending = _listener.GetPendingPublish(_clientId, packetId);
|
||||||
|
if (pending != null)
|
||||||
|
{
|
||||||
|
RoutePublish(pending.Topic, Encoding.UTF8.GetBytes(pending.Payload));
|
||||||
|
}
|
||||||
|
|
||||||
_listener.AckPendingPublish(_clientId, packetId);
|
_listener.AckPendingPublish(_clientId, packetId);
|
||||||
await WriteBinaryAsync(MqttPacketWriter.WritePubComp(packetId), ct);
|
await WriteBinaryAsync(MqttPacketWriter.WritePubComp(packetId), ct);
|
||||||
}
|
}
|
||||||
@@ -383,8 +418,31 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
for (var i = 0; i < subscribeInfo.Filters.Count; i++)
|
for (var i = 0; i < subscribeInfo.Filters.Count; i++)
|
||||||
{
|
{
|
||||||
var (topicFilter, requestedQoS) = subscribeInfo.Filters[i];
|
var (topicFilter, requestedQoS) = subscribeInfo.Filters[i];
|
||||||
_listener.RegisterSubscription(this, topicFilter);
|
|
||||||
|
if (Adapter != null)
|
||||||
|
{
|
||||||
|
// Route through SubList for cross-protocol delivery
|
||||||
|
var natsSubject = MqttTopicMapper.MqttToNats(topicFilter);
|
||||||
|
var sid = $"$MQTT_{Interlocked.Increment(ref _nextSid)}";
|
||||||
|
Adapter.AddSubscription(natsSubject, sid);
|
||||||
|
_topicToSid[topicFilter] = sid;
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
// Test-only fallback: MQTT-only subscription
|
||||||
|
_listener.RegisterSubscription(this, topicFilter);
|
||||||
|
}
|
||||||
|
|
||||||
grantedQoS[i] = Math.Min(requestedQoS, (byte)2);
|
grantedQoS[i] = Math.Min(requestedQoS, (byte)2);
|
||||||
|
|
||||||
|
// Deliver retained messages for this topic filter
|
||||||
|
var retained = _listener.GetRetainedMessage(topicFilter);
|
||||||
|
if (retained != null)
|
||||||
|
{
|
||||||
|
var retainedPayload = Encoding.UTF8.GetBytes(retained);
|
||||||
|
await WriteBinaryAsync(
|
||||||
|
MqttPacketWriter.WritePublish(topicFilter, retainedPayload, qos: 0, retain: true, packetId: 0), ct);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
await WriteBinaryAsync(MqttPacketWriter.WriteSubAck(subscribeInfo.PacketId, grantedQoS), ct);
|
await WriteBinaryAsync(MqttPacketWriter.WriteSubAck(subscribeInfo.PacketId, grantedQoS), ct);
|
||||||
@@ -395,7 +453,16 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
var unsubInfo = MqttBinaryDecoder.ParseUnsubscribe(packet.Payload.Span, packet.Flags);
|
var unsubInfo = MqttBinaryDecoder.ParseUnsubscribe(packet.Payload.Span, packet.Flags);
|
||||||
|
|
||||||
foreach (var filter in unsubInfo.Filters)
|
foreach (var filter in unsubInfo.Filters)
|
||||||
_listener.UnregisterSubscription(this, filter);
|
{
|
||||||
|
if (Adapter != null && _topicToSid.Remove(filter, out var sid))
|
||||||
|
{
|
||||||
|
Adapter.RemoveSubscription(sid);
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
_listener.UnregisterSubscription(this, filter);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
await WriteBinaryAsync(MqttPacketWriter.WriteUnsubAck(unsubInfo.PacketId), ct);
|
await WriteBinaryAsync(MqttPacketWriter.WriteUnsubAck(unsubInfo.PacketId), ct);
|
||||||
}
|
}
|
||||||
@@ -427,6 +494,13 @@ public sealed class MqttConnection : IAsyncDisposable
|
|||||||
|
|
||||||
public async ValueTask DisposeAsync()
|
public async ValueTask DisposeAsync()
|
||||||
{
|
{
|
||||||
|
// Clean up adapter subscriptions and unregister from listener
|
||||||
|
if (Adapter != null)
|
||||||
|
{
|
||||||
|
Adapter.RemoveAllSubscriptions();
|
||||||
|
_listener.UnregisterMqttAdapter(Adapter);
|
||||||
|
}
|
||||||
|
|
||||||
_listener.Unregister(this);
|
_listener.Unregister(this);
|
||||||
_writeGate.Dispose();
|
_writeGate.Dispose();
|
||||||
await _stream.DisposeAsync();
|
await _stream.DisposeAsync();
|
||||||
|
|||||||
@@ -24,6 +24,8 @@ public sealed class MqttListener : IAsyncDisposable
|
|||||||
private readonly ConcurrentDictionary<string, MqttSessionState> _sessions = new(StringComparer.Ordinal);
|
private readonly ConcurrentDictionary<string, MqttSessionState> _sessions = new(StringComparer.Ordinal);
|
||||||
private readonly ConcurrentDictionary<string, MqttConnection> _clientIdMap = new(StringComparer.Ordinal);
|
private readonly ConcurrentDictionary<string, MqttConnection> _clientIdMap = new(StringComparer.Ordinal);
|
||||||
private readonly ConcurrentDictionary<string, string> _retainedMessages = new(StringComparer.Ordinal);
|
private readonly ConcurrentDictionary<string, string> _retainedMessages = new(StringComparer.Ordinal);
|
||||||
|
private readonly IMessageRouter? _router;
|
||||||
|
private readonly ConcurrentDictionary<ulong, MqttNatsClientAdapter> _mqttAdapters = new();
|
||||||
private MqttStreamInitializer? _streamInitializer;
|
private MqttStreamInitializer? _streamInitializer;
|
||||||
private MqttConsumerManager? _mqttConsumerManager;
|
private MqttConsumerManager? _mqttConsumerManager;
|
||||||
private TcpListener? _listener;
|
private TcpListener? _listener;
|
||||||
@@ -63,7 +65,8 @@ public sealed class MqttListener : IAsyncDisposable
|
|||||||
AuthService? authService,
|
AuthService? authService,
|
||||||
MqttOptions mqttOptions,
|
MqttOptions mqttOptions,
|
||||||
MqttStreamInitializer? streamInitializer = null,
|
MqttStreamInitializer? streamInitializer = null,
|
||||||
MqttConsumerManager? mqttConsumerManager = null)
|
MqttConsumerManager? mqttConsumerManager = null,
|
||||||
|
IMessageRouter? router = null)
|
||||||
{
|
{
|
||||||
_host = host;
|
_host = host;
|
||||||
_port = port;
|
_port = port;
|
||||||
@@ -73,6 +76,7 @@ public sealed class MqttListener : IAsyncDisposable
|
|||||||
_requiredPassword = mqttOptions.Password;
|
_requiredPassword = mqttOptions.Password;
|
||||||
_streamInitializer = streamInitializer;
|
_streamInitializer = streamInitializer;
|
||||||
_mqttConsumerManager = mqttConsumerManager;
|
_mqttConsumerManager = mqttConsumerManager;
|
||||||
|
_router = router;
|
||||||
|
|
||||||
// Build TLS options if configured
|
// Build TLS options if configured
|
||||||
if (mqttOptions.HasTls)
|
if (mqttOptions.HasTls)
|
||||||
@@ -91,6 +95,49 @@ public sealed class MqttListener : IAsyncDisposable
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
internal MqttConsumerManager? ConsumerManager => _mqttConsumerManager;
|
internal MqttConsumerManager? ConsumerManager => _mqttConsumerManager;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// The message router for cross-protocol delivery (MQTT→NATS SubList routing).
|
||||||
|
/// Null when running in test-only mode without NatsServer.
|
||||||
|
/// </summary>
|
||||||
|
internal IMessageRouter? Router => _router;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Delegate to allocate a server-unique client ID for MQTT adapters.
|
||||||
|
/// Set by NatsServer after construction.
|
||||||
|
/// </summary>
|
||||||
|
internal Func<ulong>? AllocateClientId { get; set; }
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Delegate to resolve an Account by name for MQTT adapters.
|
||||||
|
/// Set by NatsServer after construction.
|
||||||
|
/// </summary>
|
||||||
|
internal Func<string?, Auth.Account?>? ResolveAccount { get; set; }
|
||||||
|
|
||||||
|
internal void RegisterMqttAdapter(MqttNatsClientAdapter adapter)
|
||||||
|
=> _mqttAdapters[adapter.Id] = adapter;
|
||||||
|
|
||||||
|
internal void UnregisterMqttAdapter(MqttNatsClientAdapter adapter)
|
||||||
|
=> _mqttAdapters.TryRemove(adapter.Id, out _);
|
||||||
|
|
||||||
|
internal IEnumerable<MqttNatsClientAdapter> GetMqttAdapters()
|
||||||
|
=> _mqttAdapters.Values;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Looks up a specific pending publish by client ID and packet ID.
|
||||||
|
/// Used by QoS 2 PUBREL to retrieve the stored message for delivery.
|
||||||
|
/// </summary>
|
||||||
|
internal MqttPendingPublish? GetPendingPublish(string clientId, int packetId)
|
||||||
|
{
|
||||||
|
if (string.IsNullOrWhiteSpace(clientId) || packetId <= 0)
|
||||||
|
return null;
|
||||||
|
|
||||||
|
if (_sessions.TryGetValue(clientId, out var session)
|
||||||
|
&& session.Pending.TryGetValue(packetId, out var pending))
|
||||||
|
return pending;
|
||||||
|
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
public Task StartAsync(CancellationToken ct)
|
public Task StartAsync(CancellationToken ct)
|
||||||
{
|
{
|
||||||
var linked = CancellationTokenSource.CreateLinkedTokenSource(ct, _cts.Token);
|
var linked = CancellationTokenSource.CreateLinkedTokenSource(ct, _cts.Token);
|
||||||
@@ -280,6 +327,7 @@ public sealed class MqttListener : IAsyncDisposable
|
|||||||
_sessions.Clear();
|
_sessions.Clear();
|
||||||
_clientIdMap.Clear();
|
_clientIdMap.Clear();
|
||||||
_retainedMessages.Clear();
|
_retainedMessages.Clear();
|
||||||
|
_mqttAdapters.Clear();
|
||||||
_cts.Dispose();
|
_cts.Dispose();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -122,6 +122,12 @@ public sealed class NatsServer : IMessageRouter, ISubListAccess, IDisposable
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
public int? MqttListenerPort => _mqttListener?.Port;
|
public int? MqttListenerPort => _mqttListener?.Port;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Returns all active MQTT client adapters for monitoring (/connz).
|
||||||
|
/// </summary>
|
||||||
|
public IEnumerable<Mqtt.MqttNatsClientAdapter> GetMqttAdapters()
|
||||||
|
=> _mqttListener?.GetMqttAdapters() ?? [];
|
||||||
|
|
||||||
public Account SystemAccount => _systemAccount;
|
public Account SystemAccount => _systemAccount;
|
||||||
public string ServerNKey { get; }
|
public string ServerNKey { get; }
|
||||||
public InternalEventSystem? EventSystem => _eventSystem;
|
public InternalEventSystem? EventSystem => _eventSystem;
|
||||||
@@ -937,7 +943,10 @@ public sealed class NatsServer : IMessageRouter, ISubListAccess, IDisposable
|
|||||||
_authService,
|
_authService,
|
||||||
mqttOptions,
|
mqttOptions,
|
||||||
mqttStreamInit,
|
mqttStreamInit,
|
||||||
mqttConsumerMgr);
|
mqttConsumerMgr,
|
||||||
|
router: this);
|
||||||
|
_mqttListener.AllocateClientId = () => Interlocked.Increment(ref _nextClientId);
|
||||||
|
_mqttListener.ResolveAccount = name => GetOrCreateAccount(name ?? Auth.Account.GlobalAccountName);
|
||||||
await _mqttListener.StartAsync(linked.Token);
|
await _mqttListener.StartAsync(linked.Token);
|
||||||
}
|
}
|
||||||
if (_jetStreamService != null)
|
if (_jetStreamService != null)
|
||||||
|
|||||||
@@ -0,0 +1,219 @@
|
|||||||
|
using System.Text;
|
||||||
|
using MQTTnet;
|
||||||
|
using MQTTnet.Client;
|
||||||
|
using MQTTnet.Protocol;
|
||||||
|
using NATS.Client.Core;
|
||||||
|
using NATS.E2E.Tests.Infrastructure;
|
||||||
|
|
||||||
|
namespace NATS.E2E.Tests;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// End-to-end tests for MQTT 3.1.1 interop using MQTTnet client library.
|
||||||
|
/// Verifies binary MQTT protocol, cross-protocol MQTT↔NATS messaging, and QoS 1.
|
||||||
|
/// </summary>
|
||||||
|
[Collection("E2E-Mqtt")]
|
||||||
|
public class MqttE2ETests(MqttServerFixture fixture)
|
||||||
|
{
|
||||||
|
[Fact]
|
||||||
|
public async Task MqttE2E_ConnectPublishSubscribe()
|
||||||
|
{
|
||||||
|
using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(15));
|
||||||
|
|
||||||
|
var factory = new MqttFactory();
|
||||||
|
using var subscriber = factory.CreateMqttClient();
|
||||||
|
using var publisher = factory.CreateMqttClient();
|
||||||
|
|
||||||
|
var subOpts = new MqttClientOptionsBuilder()
|
||||||
|
.WithTcpServer("127.0.0.1", fixture.MqttPort)
|
||||||
|
.WithClientId("e2e-mqttnet-sub")
|
||||||
|
.WithProtocolVersion(MQTTnet.Formatter.MqttProtocolVersion.V311)
|
||||||
|
.Build();
|
||||||
|
|
||||||
|
var pubOpts = new MqttClientOptionsBuilder()
|
||||||
|
.WithTcpServer("127.0.0.1", fixture.MqttPort)
|
||||||
|
.WithClientId("e2e-mqttnet-pub")
|
||||||
|
.WithProtocolVersion(MQTTnet.Formatter.MqttProtocolVersion.V311)
|
||||||
|
.Build();
|
||||||
|
|
||||||
|
await subscriber.ConnectAsync(subOpts, cts.Token);
|
||||||
|
await publisher.ConnectAsync(pubOpts, cts.Token);
|
||||||
|
|
||||||
|
var received = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously);
|
||||||
|
subscriber.ApplicationMessageReceivedAsync += e =>
|
||||||
|
{
|
||||||
|
var payload = Encoding.UTF8.GetString(e.ApplicationMessage.PayloadSegment);
|
||||||
|
received.TrySetResult(payload);
|
||||||
|
return Task.CompletedTask;
|
||||||
|
};
|
||||||
|
|
||||||
|
await subscriber.SubscribeAsync(
|
||||||
|
factory.CreateSubscribeOptionsBuilder()
|
||||||
|
.WithTopicFilter("test/mqttnet/e2e")
|
||||||
|
.Build(),
|
||||||
|
cts.Token);
|
||||||
|
|
||||||
|
// Small delay to let subscription propagate
|
||||||
|
await Task.Delay(100, cts.Token);
|
||||||
|
|
||||||
|
await publisher.PublishAsync(
|
||||||
|
new MqttApplicationMessageBuilder()
|
||||||
|
.WithTopic("test/mqttnet/e2e")
|
||||||
|
.WithPayload("hello-mqttnet")
|
||||||
|
.Build(),
|
||||||
|
cts.Token);
|
||||||
|
|
||||||
|
var msg = await received.Task.WaitAsync(cts.Token);
|
||||||
|
msg.ShouldBe("hello-mqttnet");
|
||||||
|
|
||||||
|
await subscriber.DisconnectAsync(cancellationToken: cts.Token);
|
||||||
|
await publisher.DisconnectAsync(cancellationToken: cts.Token);
|
||||||
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public async Task MqttE2E_CrossProtocol_MqttPublish_NatsSubscribe()
|
||||||
|
{
|
||||||
|
using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(15));
|
||||||
|
|
||||||
|
// NATS subscriber on "sensor.temp" (MQTT topic "sensor/temp" maps to NATS subject "sensor.temp")
|
||||||
|
await using var natsConn = fixture.CreateNatsClient();
|
||||||
|
await natsConn.ConnectAsync();
|
||||||
|
|
||||||
|
var natsReceived = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously);
|
||||||
|
_ = Task.Run(async () =>
|
||||||
|
{
|
||||||
|
await foreach (var msg in natsConn.SubscribeAsync<string>("sensor.temp", cancellationToken: cts.Token))
|
||||||
|
{
|
||||||
|
natsReceived.TrySetResult(msg.Data ?? "");
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}, cts.Token);
|
||||||
|
|
||||||
|
// Small delay to let NATS subscription register
|
||||||
|
await Task.Delay(200, cts.Token);
|
||||||
|
|
||||||
|
// MQTT publisher on "sensor/temp"
|
||||||
|
var factory = new MqttFactory();
|
||||||
|
using var mqttPub = factory.CreateMqttClient();
|
||||||
|
var pubOpts = new MqttClientOptionsBuilder()
|
||||||
|
.WithTcpServer("127.0.0.1", fixture.MqttPort)
|
||||||
|
.WithClientId("e2e-cross-mqtt-pub")
|
||||||
|
.WithProtocolVersion(MQTTnet.Formatter.MqttProtocolVersion.V311)
|
||||||
|
.Build();
|
||||||
|
|
||||||
|
await mqttPub.ConnectAsync(pubOpts, cts.Token);
|
||||||
|
await mqttPub.PublishAsync(
|
||||||
|
new MqttApplicationMessageBuilder()
|
||||||
|
.WithTopic("sensor/temp")
|
||||||
|
.WithPayload("22.5")
|
||||||
|
.Build(),
|
||||||
|
cts.Token);
|
||||||
|
|
||||||
|
var result = await natsReceived.Task.WaitAsync(cts.Token);
|
||||||
|
result.ShouldBe("22.5");
|
||||||
|
|
||||||
|
await mqttPub.DisconnectAsync(cancellationToken: cts.Token);
|
||||||
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public async Task MqttE2E_CrossProtocol_NatsPublish_MqttSubscribe()
|
||||||
|
{
|
||||||
|
using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(15));
|
||||||
|
|
||||||
|
// MQTT subscriber on "sensor/humidity"
|
||||||
|
var factory = new MqttFactory();
|
||||||
|
using var mqttSub = factory.CreateMqttClient();
|
||||||
|
var subOpts = new MqttClientOptionsBuilder()
|
||||||
|
.WithTcpServer("127.0.0.1", fixture.MqttPort)
|
||||||
|
.WithClientId("e2e-cross-mqtt-sub")
|
||||||
|
.WithProtocolVersion(MQTTnet.Formatter.MqttProtocolVersion.V311)
|
||||||
|
.Build();
|
||||||
|
|
||||||
|
await mqttSub.ConnectAsync(subOpts, cts.Token);
|
||||||
|
|
||||||
|
var mqttReceived = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously);
|
||||||
|
mqttSub.ApplicationMessageReceivedAsync += e =>
|
||||||
|
{
|
||||||
|
var payload = Encoding.UTF8.GetString(e.ApplicationMessage.PayloadSegment);
|
||||||
|
mqttReceived.TrySetResult(payload);
|
||||||
|
return Task.CompletedTask;
|
||||||
|
};
|
||||||
|
|
||||||
|
await mqttSub.SubscribeAsync(
|
||||||
|
factory.CreateSubscribeOptionsBuilder()
|
||||||
|
.WithTopicFilter("sensor/humidity")
|
||||||
|
.Build(),
|
||||||
|
cts.Token);
|
||||||
|
|
||||||
|
// Small delay to let subscription propagate through SubList
|
||||||
|
await Task.Delay(200, cts.Token);
|
||||||
|
|
||||||
|
// NATS publisher on "sensor.humidity"
|
||||||
|
await using var natsConn = fixture.CreateNatsClient();
|
||||||
|
await natsConn.ConnectAsync();
|
||||||
|
await natsConn.PublishAsync("sensor.humidity", "65%", cancellationToken: cts.Token);
|
||||||
|
|
||||||
|
var result = await mqttReceived.Task.WaitAsync(cts.Token);
|
||||||
|
result.ShouldBe("65%");
|
||||||
|
|
||||||
|
await mqttSub.DisconnectAsync(cancellationToken: cts.Token);
|
||||||
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public async Task MqttE2E_Qos1_PubAck()
|
||||||
|
{
|
||||||
|
using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(15));
|
||||||
|
|
||||||
|
var factory = new MqttFactory();
|
||||||
|
using var subscriber = factory.CreateMqttClient();
|
||||||
|
using var publisher = factory.CreateMqttClient();
|
||||||
|
|
||||||
|
var subOpts = new MqttClientOptionsBuilder()
|
||||||
|
.WithTcpServer("127.0.0.1", fixture.MqttPort)
|
||||||
|
.WithClientId("e2e-qos1-sub")
|
||||||
|
.WithProtocolVersion(MQTTnet.Formatter.MqttProtocolVersion.V311)
|
||||||
|
.Build();
|
||||||
|
|
||||||
|
var pubOpts = new MqttClientOptionsBuilder()
|
||||||
|
.WithTcpServer("127.0.0.1", fixture.MqttPort)
|
||||||
|
.WithClientId("e2e-qos1-pub")
|
||||||
|
.WithProtocolVersion(MQTTnet.Formatter.MqttProtocolVersion.V311)
|
||||||
|
.Build();
|
||||||
|
|
||||||
|
await subscriber.ConnectAsync(subOpts, cts.Token);
|
||||||
|
await publisher.ConnectAsync(pubOpts, cts.Token);
|
||||||
|
|
||||||
|
var received = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously);
|
||||||
|
subscriber.ApplicationMessageReceivedAsync += e =>
|
||||||
|
{
|
||||||
|
var payload = Encoding.UTF8.GetString(e.ApplicationMessage.PayloadSegment);
|
||||||
|
received.TrySetResult(payload);
|
||||||
|
return Task.CompletedTask;
|
||||||
|
};
|
||||||
|
|
||||||
|
await subscriber.SubscribeAsync(
|
||||||
|
factory.CreateSubscribeOptionsBuilder()
|
||||||
|
.WithTopicFilter(f => f.WithTopic("test/qos1").WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce))
|
||||||
|
.Build(),
|
||||||
|
cts.Token);
|
||||||
|
|
||||||
|
await Task.Delay(100, cts.Token);
|
||||||
|
|
||||||
|
// Publish with QoS 1 — MQTTnet will expect PUBACK from server
|
||||||
|
var pubResult = await publisher.PublishAsync(
|
||||||
|
new MqttApplicationMessageBuilder()
|
||||||
|
.WithTopic("test/qos1")
|
||||||
|
.WithPayload("qos1-payload")
|
||||||
|
.WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce)
|
||||||
|
.Build(),
|
||||||
|
cts.Token);
|
||||||
|
|
||||||
|
// MQTTnet throws if PUBACK not received, so reaching here means server sent PUBACK
|
||||||
|
pubResult.ReasonCode.ShouldBe(MqttClientPublishReasonCode.Success);
|
||||||
|
|
||||||
|
var msg = await received.Task.WaitAsync(cts.Token);
|
||||||
|
msg.ShouldBe("qos1-payload");
|
||||||
|
|
||||||
|
await subscriber.DisconnectAsync(cancellationToken: cts.Token);
|
||||||
|
await publisher.DisconnectAsync(cancellationToken: cts.Token);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,151 @@
|
|||||||
|
using System.Diagnostics;
|
||||||
|
using NATS.Server.JetStream.Storage;
|
||||||
|
using Xunit.Abstractions;
|
||||||
|
|
||||||
|
namespace NATS.Server.Benchmark.Tests.JetStream;
|
||||||
|
|
||||||
|
[Collection("Benchmark-JetStream")]
|
||||||
|
public class FileStoreAppendBenchmarks(ITestOutputHelper output)
|
||||||
|
{
|
||||||
|
[Fact]
|
||||||
|
[Trait("Category", "Benchmark")]
|
||||||
|
public async Task FileStore_AppendAsync_128B_Throughput()
|
||||||
|
{
|
||||||
|
var payload = new byte[128];
|
||||||
|
var dir = CreateDirectory("append");
|
||||||
|
var opts = CreateOptions(dir);
|
||||||
|
|
||||||
|
try
|
||||||
|
{
|
||||||
|
await using var store = new FileStore(opts);
|
||||||
|
await MeasureAsync("FileStore AppendAsync (128B)", operations: 20_000, payload.Length,
|
||||||
|
i => store.AppendAsync($"bench.append.{i % 8}", payload, default).AsTask());
|
||||||
|
}
|
||||||
|
finally
|
||||||
|
{
|
||||||
|
DeleteDirectory(dir);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
[Trait("Category", "Benchmark")]
|
||||||
|
public void FileStore_LoadLastBySubject_Throughput()
|
||||||
|
{
|
||||||
|
var payload = new byte[64];
|
||||||
|
var dir = CreateDirectory("load-last");
|
||||||
|
var opts = CreateOptions(dir);
|
||||||
|
|
||||||
|
try
|
||||||
|
{
|
||||||
|
using var store = new FileStore(opts);
|
||||||
|
for (var i = 0; i < 25_000; i++)
|
||||||
|
store.StoreMsg($"bench.subject.{i % 16}", null, payload, 0L);
|
||||||
|
|
||||||
|
Measure("FileStore LoadLastBySubject (hot)", operations: 50_000, payload.Length,
|
||||||
|
() =>
|
||||||
|
{
|
||||||
|
var loaded = store.LoadLastBySubjectAsync("bench.subject.7", default).GetAwaiter().GetResult();
|
||||||
|
if (loaded is null || loaded.Payload.Length != payload.Length)
|
||||||
|
throw new InvalidOperationException("LoadLastBySubjectAsync returned an unexpected result.");
|
||||||
|
});
|
||||||
|
}
|
||||||
|
finally
|
||||||
|
{
|
||||||
|
DeleteDirectory(dir);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
[Trait("Category", "Benchmark")]
|
||||||
|
public void FileStore_PurgeEx_Trim_Overhead()
|
||||||
|
{
|
||||||
|
var payload = new byte[96];
|
||||||
|
var dir = CreateDirectory("purge-trim");
|
||||||
|
var opts = CreateOptions(dir);
|
||||||
|
|
||||||
|
try
|
||||||
|
{
|
||||||
|
using var store = new FileStore(opts);
|
||||||
|
for (var i = 0; i < 12_000; i++)
|
||||||
|
store.StoreMsg($"bench.purge.{i % 6}", null, payload, 0L);
|
||||||
|
|
||||||
|
Measure("FileStore PurgeEx+Trim", operations: 2_000, payload.Length,
|
||||||
|
() =>
|
||||||
|
{
|
||||||
|
store.PurgeEx("bench.purge.1", 0, 8);
|
||||||
|
store.TrimToMaxMessages(10_000);
|
||||||
|
store.StoreMsg("bench.purge.1", null, payload, 0L);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
finally
|
||||||
|
{
|
||||||
|
DeleteDirectory(dir);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private async Task MeasureAsync(string name, int operations, int payloadSize, Func<int, Task> action)
|
||||||
|
{
|
||||||
|
GC.Collect();
|
||||||
|
GC.WaitForPendingFinalizers();
|
||||||
|
GC.Collect();
|
||||||
|
|
||||||
|
var beforeAlloc = GC.GetAllocatedBytesForCurrentThread();
|
||||||
|
var sw = Stopwatch.StartNew();
|
||||||
|
|
||||||
|
for (var i = 0; i < operations; i++)
|
||||||
|
await action(i);
|
||||||
|
|
||||||
|
sw.Stop();
|
||||||
|
WriteResult(name, operations, (long)operations * payloadSize, sw.Elapsed, GC.GetAllocatedBytesForCurrentThread() - beforeAlloc);
|
||||||
|
}
|
||||||
|
|
||||||
|
private void Measure(string name, int operations, int payloadSize, Action action)
|
||||||
|
{
|
||||||
|
GC.Collect();
|
||||||
|
GC.WaitForPendingFinalizers();
|
||||||
|
GC.Collect();
|
||||||
|
|
||||||
|
var beforeAlloc = GC.GetAllocatedBytesForCurrentThread();
|
||||||
|
var sw = Stopwatch.StartNew();
|
||||||
|
|
||||||
|
for (var i = 0; i < operations; i++)
|
||||||
|
action();
|
||||||
|
|
||||||
|
sw.Stop();
|
||||||
|
WriteResult(name, operations, (long)operations * payloadSize, sw.Elapsed, GC.GetAllocatedBytesForCurrentThread() - beforeAlloc);
|
||||||
|
}
|
||||||
|
|
||||||
|
private void WriteResult(string name, int operations, long totalBytes, TimeSpan elapsed, long allocatedBytes)
|
||||||
|
{
|
||||||
|
var opsPerSecond = operations / elapsed.TotalSeconds;
|
||||||
|
var megabytesPerSecond = totalBytes / elapsed.TotalSeconds / (1024.0 * 1024.0);
|
||||||
|
var bytesPerOperation = allocatedBytes / (double)operations;
|
||||||
|
|
||||||
|
output.WriteLine($"=== {name} ===");
|
||||||
|
output.WriteLine($"Ops: {opsPerSecond:N0} ops/s");
|
||||||
|
output.WriteLine($"Data: {megabytesPerSecond:F1} MB/s");
|
||||||
|
output.WriteLine($"Alloc: {bytesPerOperation:F1} B/op");
|
||||||
|
output.WriteLine($"Elapsed: {elapsed.TotalMilliseconds:F0} ms");
|
||||||
|
output.WriteLine("");
|
||||||
|
}
|
||||||
|
|
||||||
|
private static string CreateDirectory(string suffix)
|
||||||
|
=> Path.Combine(Path.GetTempPath(), $"nats-js-filestore-bench-{suffix}-{Guid.NewGuid():N}");
|
||||||
|
|
||||||
|
private static FileStoreOptions CreateOptions(string dir)
|
||||||
|
{
|
||||||
|
Directory.CreateDirectory(dir);
|
||||||
|
|
||||||
|
return new FileStoreOptions
|
||||||
|
{
|
||||||
|
Directory = dir,
|
||||||
|
BlockSizeBytes = 256 * 1024,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
private static void DeleteDirectory(string dir)
|
||||||
|
{
|
||||||
|
if (Directory.Exists(dir))
|
||||||
|
Directory.Delete(dir, recursive: true);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -45,6 +45,9 @@ Use `-v normal` or `--logger "console;verbosity=detailed"` to see the comparison
|
|||||||
| `MultiClientLatencyTests` | `RequestReply_10Clients2Services_16B` | Request/reply latency, 10 concurrent clients, 2 queue-group services |
|
| `MultiClientLatencyTests` | `RequestReply_10Clients2Services_16B` | Request/reply latency, 10 concurrent clients, 2 queue-group services |
|
||||||
| `SyncPublishTests` | `JSSyncPublish_16B_MemoryStore` | JetStream synchronous publish, memory-backed stream |
|
| `SyncPublishTests` | `JSSyncPublish_16B_MemoryStore` | JetStream synchronous publish, memory-backed stream |
|
||||||
| `AsyncPublishTests` | `JSAsyncPublish_128B_FileStore` | JetStream async batch publish, file-backed stream |
|
| `AsyncPublishTests` | `JSAsyncPublish_128B_FileStore` | JetStream async batch publish, file-backed stream |
|
||||||
|
| `FileStoreAppendBenchmarks` | `FileStore_AppendAsync_128B_Throughput` | FileStore direct append throughput, 128-byte payload |
|
||||||
|
| `FileStoreAppendBenchmarks` | `FileStore_LoadLastBySubject_Throughput` | FileStore hot-path subject index lookup throughput |
|
||||||
|
| `FileStoreAppendBenchmarks` | `FileStore_PurgeEx_Trim_Overhead` | FileStore purge/trim maintenance overhead under repeated updates |
|
||||||
| `OrderedConsumerTests` | `JSOrderedConsumer_Throughput` | JetStream ordered ephemeral consumer read throughput |
|
| `OrderedConsumerTests` | `JSOrderedConsumer_Throughput` | JetStream ordered ephemeral consumer read throughput |
|
||||||
| `DurableConsumerFetchTests` | `JSDurableFetch_Throughput` | JetStream durable consumer fetch-in-batches throughput |
|
| `DurableConsumerFetchTests` | `JSDurableFetch_Throughput` | JetStream durable consumer fetch-in-batches throughput |
|
||||||
|
|
||||||
|
|||||||
@@ -15,4 +15,27 @@ public class FileStoreTests
|
|||||||
await using var recovered = new FileStore(new FileStoreOptions { Directory = dir.FullName });
|
await using var recovered = new FileStore(new FileStoreOptions { Directory = dir.FullName });
|
||||||
(await recovered.GetStateAsync(default)).Messages.ShouldBe((ulong)1);
|
(await recovered.GetStateAsync(default)).Messages.ShouldBe((ulong)1);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public async Task Snapshot_round_trip_preserves_headers_and_payload_separately()
|
||||||
|
{
|
||||||
|
var srcDir = Directory.CreateTempSubdirectory();
|
||||||
|
var dstDir = Directory.CreateTempSubdirectory();
|
||||||
|
|
||||||
|
await using var src = new FileStore(new FileStoreOptions { Directory = srcDir.FullName });
|
||||||
|
var hdr = "NATS/1.0\r\nX-Test: two\r\n\r\n"u8.ToArray();
|
||||||
|
var msg = "payload-two"u8.ToArray();
|
||||||
|
|
||||||
|
var (seq, _) = src.StoreMsg("events.a", hdr, msg, 0L);
|
||||||
|
var snapshot = await src.CreateSnapshotAsync(default);
|
||||||
|
|
||||||
|
await using var dst = new FileStore(new FileStoreOptions { Directory = dstDir.FullName });
|
||||||
|
await dst.RestoreSnapshotAsync(snapshot, default);
|
||||||
|
|
||||||
|
var loaded = dst.LoadMsg(seq, null);
|
||||||
|
loaded.Header.ShouldNotBeNull();
|
||||||
|
loaded.Header.ShouldBe(hdr);
|
||||||
|
loaded.Data.ShouldNotBeNull();
|
||||||
|
loaded.Data.ShouldBe(msg);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+28
@@ -0,0 +1,28 @@
|
|||||||
|
using NATS.Server.JetStream.Storage;
|
||||||
|
|
||||||
|
namespace NATS.Server.JetStream.Tests.JetStream.Storage;
|
||||||
|
|
||||||
|
public sealed class FileStoreOptimizationGuardTests
|
||||||
|
{
|
||||||
|
[Fact]
|
||||||
|
public async Task PurgeEx_updates_last_by_subject_after_recovery()
|
||||||
|
{
|
||||||
|
var dir = Directory.CreateTempSubdirectory();
|
||||||
|
|
||||||
|
await using (var store = new FileStore(new FileStoreOptions { Directory = dir.FullName }))
|
||||||
|
{
|
||||||
|
store.StoreMsg("events.a", null, "one"u8.ToArray(), 0L);
|
||||||
|
store.StoreMsg("events.a", null, "two"u8.ToArray(), 0L);
|
||||||
|
store.StoreMsg("events.b", null, "other"u8.ToArray(), 0L);
|
||||||
|
store.PurgeEx("events.a", 0, 1);
|
||||||
|
await store.FlushAllPending();
|
||||||
|
}
|
||||||
|
|
||||||
|
await using var recovered = new FileStore(new FileStoreOptions { Directory = dir.FullName });
|
||||||
|
var last = await recovered.LoadLastBySubjectAsync("events.a", default);
|
||||||
|
|
||||||
|
last.ShouldNotBeNull();
|
||||||
|
last.Sequence.ShouldBe(2UL);
|
||||||
|
last.Payload.ToArray().ShouldBe("two"u8.ToArray());
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -181,7 +181,7 @@ public sealed class FileStoreTtlTests : IDisposable
|
|||||||
|
|
||||||
// Go: TestFileStoreStoreMsg — filestore.go storeMsg with headers
|
// Go: TestFileStoreStoreMsg — filestore.go storeMsg with headers
|
||||||
[Fact]
|
[Fact]
|
||||||
public async Task StoreMsg_WithHeaders_CombinesHeadersAndPayload()
|
public async Task StoreMsg_WithHeaders_KeepsPayloadSeparateFromHeaders()
|
||||||
{
|
{
|
||||||
await using var store = CreateStore(sub: "storemsg-headers");
|
await using var store = CreateStore(sub: "storemsg-headers");
|
||||||
|
|
||||||
@@ -192,10 +192,10 @@ public sealed class FileStoreTtlTests : IDisposable
|
|||||||
seq.ShouldBe(1UL);
|
seq.ShouldBe(1UL);
|
||||||
ts.ShouldBeGreaterThan(0L);
|
ts.ShouldBeGreaterThan(0L);
|
||||||
|
|
||||||
// The stored payload should be the combination of headers + body.
|
// The stored payload should remain the message body only.
|
||||||
var loaded = await store.LoadAsync(seq, default);
|
var loaded = await store.LoadAsync(seq, default);
|
||||||
loaded.ShouldNotBeNull();
|
loaded.ShouldNotBeNull();
|
||||||
loaded!.Payload.Length.ShouldBe(hdr.Length + body.Length);
|
loaded!.Payload.ToArray().ShouldBe(body);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Go: TestFileStoreStoreMsgPerMsgTtl — filestore.go per-message TTL override
|
// Go: TestFileStoreStoreMsgPerMsgTtl — filestore.go per-message TTL override
|
||||||
|
|||||||
@@ -532,4 +532,22 @@ public sealed class StoreInterfaceTests
|
|||||||
lastMsg = s.LoadLastMsg("foo", null);
|
lastMsg = s.LoadLastMsg("foo", null);
|
||||||
lastMsg.Sequence.ShouldBe(2UL);
|
lastMsg.Sequence.ShouldBe(2UL);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public void FileStore_LoadMsg_preserves_headers_separately_from_payload()
|
||||||
|
{
|
||||||
|
var dir = Directory.CreateTempSubdirectory();
|
||||||
|
using var store = new FileStore(new FileStoreOptions { Directory = dir.FullName });
|
||||||
|
|
||||||
|
var hdr = "NATS/1.0\r\nX-Test: one\r\n\r\n"u8.ToArray();
|
||||||
|
var msg = "payload"u8.ToArray();
|
||||||
|
|
||||||
|
var (seq, _) = store.StoreMsg("foo", hdr, msg, 0L);
|
||||||
|
var loaded = store.LoadMsg(seq, null);
|
||||||
|
|
||||||
|
loaded.Header.ShouldNotBeNull();
|
||||||
|
loaded.Header.ShouldBe(hdr);
|
||||||
|
loaded.Data.ShouldNotBeNull();
|
||||||
|
loaded.Data.ShouldBe(msg);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,4 +15,22 @@ public class JetStreamStoreIndexTests
|
|||||||
var last = await store.LoadLastBySubjectAsync("orders.created", default);
|
var last = await store.LoadLastBySubjectAsync("orders.created", default);
|
||||||
last!.Payload.Span.SequenceEqual("3"u8).ShouldBeTrue();
|
last!.Payload.Span.SequenceEqual("3"u8).ShouldBeTrue();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
[Fact]
|
||||||
|
public async Task FileStore_trim_to_zero_preserves_high_water_mark_for_empty_state()
|
||||||
|
{
|
||||||
|
var dir = Directory.CreateTempSubdirectory();
|
||||||
|
await using var store = new FileStore(new FileStoreOptions { Directory = dir.FullName });
|
||||||
|
|
||||||
|
await store.AppendAsync("orders.created", "1"u8.ToArray(), default);
|
||||||
|
await store.AppendAsync("orders.updated", "2"u8.ToArray(), default);
|
||||||
|
await store.AppendAsync("orders.created", "3"u8.ToArray(), default);
|
||||||
|
|
||||||
|
store.TrimToMaxMessages(0);
|
||||||
|
|
||||||
|
var state = await store.GetStateAsync(default);
|
||||||
|
state.Messages.ShouldBe(0UL);
|
||||||
|
state.LastSeq.ShouldBe(3UL);
|
||||||
|
state.FirstSeq.ShouldBe(4UL);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user