Repository navigation
Expand file tree
/
Copy pathRuntimeJournalNativeStorageTests.cs
More file actions
182 lines (167 loc) · 9.48 KB
/
Copy pathRuntimeJournalNativeStorageTests.cs
File metadata and controls
182 lines (167 loc) · 9.48 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
#pragma warning disable ORLEANSEXP005
using System.Buffers;
using KeyLoad.Core.Features.ClusterRouting.Contracts;
using KeyLoad.Orleans;
using Orleans.Journaling;
using Orleans.Storage;
namespace KeyLoad.UnitTests.Features.ClusterRouting;
[RuntimeJournalNativeDataSource]
[NotInParallel]
internal sealed class RuntimeJournalNativeStorageTests(RuntimeJournalNativeFixture fixture)
{
private const int FirstPatternIndex = 0;
private const int BinaryReplayPayloadBytes = 98_304;
private const string KindProperty = "kind";
private const string DurableKind = "durable";
private const string StaleProperty = "stale";
private const string StaleValue = "1";
private const string IndexProperty = "index";
private const string OwnerProperty = RuntimeJournalProtocol.OwnerProperty;
private const string FirstIndexValue = "one";
private const string SecondIndexValue = "two";
private const string InvalidNulProperty = "bad\0key";
[Test]
public async Task NativeStorageCreatesWritesReplaysReplacesAndDeletesCommittedContent()
{
var name = NewName("content");
var storage = fixture.Provider.CreateStorage(new(name));
await Assert.That(await storage.CreateIfNotExistsAsync(new Dictionary<string, string> { [KindProperty] = DurableKind })).IsTrue();
var content = Enumerable.Range(0, 80_000).Select(index => (byte)(index % byte.MaxValue)).ToArray();
await storage.AppendAsync(new ReadOnlySequence<byte>(content), CancellationToken.None);
using var replay = new RuntimeJournalNativeReadAccumulator();
await storage.ReadAsync(replay, CancellationToken.None);
await Assert.That(replay.IsCompleted).IsTrue();
await Assert.That(replay.ToArray().SequenceEqual(content)).IsTrue();
await Assert.That(replay.Metadata?.FormatKey).IsEqualTo(RuntimeJournalStoragePolicy.BinaryFormat);
var replacement = new byte[] { 9, 4, 7 };
await storage.ReplaceAsync(new ReadOnlySequence<byte>(replacement), CancellationToken.None);
using var replaced = new RuntimeJournalNativeReadAccumulator();
await storage.ReadAsync(replaced, CancellationToken.None);
await Assert.That(replaced.ToArray().SequenceEqual(replacement)).IsTrue();
await storage.DeleteAsync(CancellationToken.None);
await Assert.That(await storage.GetMetadataAsync()).IsNull();
}
[Test]
public async Task NativeOrleansBinaryReplaysOneEntryAcrossReadPageBoundaries()
{
var grainId = NewName("binary-replay");
var grain = fixture.Cluster.Client.GetGrain<IRuntimeJournalReplayGrain>(grainId);
var expected = Enumerable.Range(FirstPatternIndex, BinaryReplayPayloadBytes)
.Select(index => (byte)(index % byte.MaxValue)).ToArray();
await grain.SetAsync(expected);
var originalActivation = await grain.GetActivationTokenAsync();
await grain.DeactivateAsync();
await WaitForActivationChangeAsync(grain, originalActivation);
var recovered = fixture.Cluster.Client.GetGrain<IRuntimeJournalReplayGrain>(grainId);
var actual = await recovered.ReadAsync();
await Assert.That(actual).IsNotNull();
await Assert.That(actual!.SequenceEqual(expected)).IsTrue();
var beforeDeleteActivation = await recovered.GetActivationTokenAsync();
await recovered.DeleteAsync();
await recovered.DeactivateAsync();
await WaitForActivationChangeAsync(recovered, beforeDeleteActivation);
await Assert.That(await recovered.ReadAsync()).IsNull();
}
[Test]
public async Task NativeMetadataCasReturnsConflictAndFencesStaleBodyWriter()
{
var name = NewName("fence");
var first = fixture.Provider.CreateStorage(new(name));
await Assert.That(await first.CreateIfNotExistsAsync()).IsTrue();
var stale = fixture.Provider.CreateStorage(new(name));
await Assert.That(await stale.GetMetadataAsync()).IsNotNull();
var metadata = await first.GetMetadataAsync();
var updated = await first.UpdateMetadataAsync(
new Dictionary<string, string> { [OwnerProperty] = "worker-a" },
expectedETag: metadata!.ETag);
await Assert.That(updated).IsNotNull();
await Assert.That(updated!.ETag).IsNotEqualTo(metadata.ETag);
await Assert.That((await first.GetMetadataAsync())!.ETag).IsEqualTo(updated.ETag);
await Assert.That(await first.UpdateMetadataAsync(new Dictionary<string, string> { [StaleProperty] = StaleValue },
expectedETag: metadata.ETag)).IsNull();
await Assert.ThrowsExactlyAsync<InconsistentStateException>(async () =>
await stale.AppendAsync(new ReadOnlySequence<byte>(new byte[] { 1 }), CancellationToken.None));
await Assert.ThrowsExactlyAsync<InconsistentStateException>(async () =>
await first.AppendAsync(new ReadOnlySequence<byte>(new byte[] { 2 }), CancellationToken.None));
var current = fixture.Provider.CreateStorage(new(name));
await Assert.That(await current.GetMetadataAsync()).IsNotNull();
await current.AppendAsync(new ReadOnlySequence<byte>(new byte[] { 3 }), CancellationToken.None);
await current.DeleteAsync(CancellationToken.None);
await Assert.That(await first.GetMetadataAsync()).IsNull();
var replacement = fixture.Provider.CreateStorage(new(name));
await Assert.That(await replacement.CreateIfNotExistsAsync()).IsTrue();
await Assert.ThrowsExactlyAsync<InconsistentStateException>(async () =>
await first.AppendAsync(new ReadOnlySequence<byte>(new byte[] { 4 }), CancellationToken.None));
await replacement.DeleteAsync(CancellationToken.None);
}
[Test]
public async Task NativeCatalogAppliesOrdinalPrefixAndInclusiveRangeWithMetadata()
{
var prefix = NewName("catalog");
var firstName = prefix + "/a";
var secondName = prefix + "/b";
var outsideName = NewName("elsewhere");
var first = fixture.Provider.CreateStorage(new(firstName));
var second = fixture.Provider.CreateStorage(new(secondName));
var outside = fixture.Provider.CreateStorage(new(outsideName));
await first.CreateIfNotExistsAsync(new Dictionary<string, string> { [IndexProperty] = FirstIndexValue });
await second.CreateIfNotExistsAsync(new Dictionary<string, string> { [IndexProperty] = SecondIndexValue });
await outside.CreateIfNotExistsAsync();
var options = new JournalCatalogListOptions { Prefix = new(prefix), MinId = new(firstName), MaxId = new(secondName), IncludeMetadata = true };
var entries = new List<JournalCatalogEntry>();
await foreach (var entry in fixture.Catalog.ListAsync(options))
{
entries.Add(entry);
}
await Assert.That(entries.Count).IsEqualTo(2);
await Assert.That(entries.All(entry => entry.Metadata?.FormatKey == RuntimeJournalStoragePolicy.BinaryFormat)).IsTrue();
await Assert.That(entries.Select(entry => entry.Id.Value).ToHashSet(StringComparer.Ordinal))
.IsEquivalentTo(new[] { firstName, secondName });
await first.DeleteAsync(CancellationToken.None);
await second.DeleteAsync(CancellationToken.None);
await outside.DeleteAsync(CancellationToken.None);
}
[Test]
public async Task NativeMetadataRejectsNulKeysWithoutChangingCommittedMetadata()
{
var storage = fixture.Provider.CreateStorage(new(NewName("metadata")));
await storage.CreateIfNotExistsAsync();
var before = await storage.GetMetadataAsync();
await Assert.ThrowsExactlyAsync<KeyLoadException>(async () =>
await storage.UpdateMetadataAsync(new Dictionary<string, string> { [InvalidNulProperty] = "value" }));
var after = await storage.GetMetadataAsync();
await Assert.That(after?.ETag).IsEqualTo(before?.ETag);
await Assert.That(after?.Properties).IsEmpty();
await storage.DeleteAsync(CancellationToken.None);
}
[Test]
public async Task NativeOversizeWriteFailsBeforeChangingTheCommittedJournal()
{
var storage = fixture.Provider.CreateStorage(new(NewName("capacity")));
await storage.CreateIfNotExistsAsync();
var tooLarge = new byte[fixture.JournalOptions.Value.MaximumJournalBytes + 1];
var failure = await Assert.ThrowsExactlyAsync<KeyLoadException>(async () =>
await storage.ReplaceAsync(new ReadOnlySequence<byte>(tooLarge), CancellationToken.None));
await Assert.That(failure!.Code).IsEqualTo(ErrorCode.ResourceExhausted);
using var replay = new RuntimeJournalNativeReadAccumulator();
await storage.ReadAsync(replay, CancellationToken.None);
await Assert.That(replay.ToArray()).IsEmpty();
await storage.DeleteAsync(CancellationToken.None);
}
private static string NewName(string prefix) => $"native/{prefix}/{Guid.NewGuid():N}";
private async Task WaitForActivationChangeAsync(IRuntimeJournalReplayGrain grain, string previousToken)
{
using var deadline = new CancellationTokenSource(fixture.TimingOptions.Value.CompletionTimeout, TimeProvider.System);
while (true)
{
deadline.Token.ThrowIfCancellationRequested();
var current = await grain.GetActivationTokenAsync().WaitAsync(deadline.Token);
if (current != previousToken)
{
return;
}
await Task.Delay(fixture.TimingOptions.Value.PollInterval, TimeProvider.System, deadline.Token);
}
}
}
#pragma warning restore ORLEANSEXP005