Skip to content

Commit d50001d

Browse files
committed
Preserve strict native inspection and reject corrupt stores before provider writes
1 parent 366d3a8 commit d50001d

38 files changed

Lines changed: 1002 additions & 150 deletions

File tree

‎docs/ADR/ADR-060-native-internal-serialization.md‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,3 +93,44 @@ derived collection codecs fail closed before generated allocation, even below an
9393
object-typed field. Extending this closure requires a native shape specification,
9494
preflight/count/reference tests and homogeneous compatibility qualification.
9595
This is an internal concrete contract, not arbitrary Orleans codec compatibility.
96+
97+
## Accepted qualification repair contract
98+
99+
R11/R12 preserve AC-IS002/004/007 and the existing formats. Actual GitHub
100+
run37120641864 revealed nullable metadata construction, native reader exhaustion
101+
classification and stale fixtures. Native graph annotations must use .NET10's
102+
already-unwrapped Nullable<T> metadata; only genuine Orleans Reader buffer
103+
exhaustion joins coded corruption, while programming/session exceptions escape.
104+
105+
R12F validates the complete recoverable journal/checkpoint prefix before opening
106+
ZoneTree's native provider. Under the already acquired node owner lock, use the
107+
same owned journal handle and bounded current codecs to verify complete frames,
108+
checksums, sequence, checkpoint metadata/footer and record semantics without apply,
109+
truncation or tree writes. Leave a permitted incomplete current tail untouched
110+
during preflight; ordinary ordered recovery alone applies/truncates it afterward.
111+
Unsupported complete legacy frames and complete corruption must fail without
112+
changing journal/identity/provider files. Reset the journal position before
113+
ordinary recovery; keep startup preflight distinct from acknowledged write gates.
114+
115+
StorageRecovery owns the initializer and new preflight helper; UnitTests owns
116+
NativeStoreOpenPreflight regressions using real files plus unchanged historical
117+
file-preservation assertions. Root owns integration/docs; wire worker owns reader
118+
normalization; no shared runtime/phase-file overlap. Verify valid checkpoint/tail,
119+
torn-tail recovery, late complete corruption and legacy rejection through actual
120+
GitHub normal/scalar/recovery/RF3 suites. No migration or old-store conversion is
121+
introduced; rollback restores prior binaries for matching stores. Extra startup
122+
validation cost requires actual recovery/performance evidence and cannot count
123+
as a speed improvement. All fault assertions and numeric budgets remain.
124+
125+
R12G preserves the existing strict cold-term read contract: the exact stored
126+
ReplicaEntry must pass the same bounded replica inspection as ordinary stored
127+
entry reads before its term is observed. Pass the owning configuration's
128+
MaxAppendEntries through the term reader; preserve its scoped storage identity,
129+
cut cache, index/term checks and lookup counts. ClusterReplication owns the two
130+
reader/caller files; the existing genuine-file malformed nested-operation
131+
recovery test is the acceptance oracle. Unknown nested authority fields remain
132+
corruption; ordinary bounded persistence evolution does not weaken this check.
133+
The borrowed storage span is copied for inspection only on a cold miss because
134+
the profile inspector requires owned ReadOnlyMemory for its borrowed proxies;
135+
no proxy or storage buffer escapes the read gate. Warm cut observations still
136+
perform no point lookup. Retain this ownership cost in performance accounting.

‎src/KeyLoad.Abstractions/Features/InternalSerialization/NativePayloadSyntax.cs‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,21 @@ namespace KeyLoad.Features.InternalSerialization;
77
// Official native headers/scalars are inspected before generated decoding allocates values.
88
internal static class NativePayloadSyntax
99
{
10+
// Orleans 10.3.1 Reader uses this distinct throw site for exhausted/truncated buffers.
11+
// Other InvalidOperationException failures, including session/codec invariants, must escape.
12+
internal static bool IsReaderBufferFailure(Exception exception)
13+
{
14+
if (exception is not InvalidOperationException)
15+
{
16+
return false;
17+
}
18+
var method = exception.TargetSite;
19+
return exception.Message == "Insufficient data present in buffer."
20+
&& method?.Name == "ThrowInsufficientData"
21+
&& method.DeclaringType is { IsGenericType: true } declaring
22+
&& declaring.GetGenericTypeDefinition() == typeof(Reader<>);
23+
}
24+
1025
internal static void Validate(ReadOnlySpan<byte> bytes, SerializerSession session, Type? expectedRootType = null, bool exactRoot = false)
1126
{
1227
var reader = Reader.Create(bytes, session);

‎src/KeyLoad.Abstractions/Features/InternalSerialization/NativeSerialization.cs‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -155,5 +155,6 @@ internal static long Measure<T>(T value, SerializerSession session)
155155
private static bool IsMalformed(Exception exception)
156156
=> exception is SerializerException
157157
or ArgumentException or IndexOutOfRangeException or OverflowException or InvalidCastException
158-
or FormatException or EndOfStreamException or TypeLoadException;
158+
or FormatException or EndOfStreamException or TypeLoadException
159+
|| NativePayloadSyntax.IsReaderBufferFailure(exception);
159160
}

‎src/KeyLoad.Replication/Features/ClusterReplication/DurableReplicaLog.cs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ private long Term(long index)
5858
{
5959
throw Errors.Fail(ErrorCode.NotFound, ReplicaPersistence.MissingEntry);
6060
}
61-
var result = ReplicaTermObservationReader.Read(store, index, state.Term, termObservation);
61+
var result = ReplicaTermObservationReader.Read(store, index, state.Term, configuration.MaxAppendEntries, termObservation);
6262
termObservation = result.Observation;
6363
return result.Term;
6464
}

‎src/KeyLoad.Replication/Features/ClusterReplication/ReplicaInspectionFields.cs‎

Lines changed: 37 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,55 @@
11
using System.Buffers;
2+
using System.Collections.Immutable;
3+
using Microsoft.Extensions.DependencyInjection;
24
using Orleans.Serialization.Buffers;
35
using Orleans.Serialization.Codecs;
6+
using Orleans.Serialization.Session;
47
using Orleans.Serialization.WireProtocol;
58

69
namespace KeyLoad.Replication;
710

811
internal static class ReplicaInspectionFields
912
{
13+
// CodecProvider searches generic definitions, so closed profile codecs require explicit resolution.
14+
// This finite private dispatch shares the original session and never admits arbitrary codecs.
15+
private static readonly Dictionary<Type, Type> ProfileCodecs = new()
16+
{
17+
[typeof(ReadOnlyMemory<byte>)] = typeof(ReplicaBorrowedBytesCodec),
18+
[typeof(ImmutableArray<ReplicaEntry>)] = typeof(ReplicaEntriesInspectionCodec),
19+
[typeof(ImmutableArray<string>)] = typeof(ReplicaVotersInspectionCodec),
20+
[typeof(ReplicaEntry[])] = typeof(ReplicaEntryArrayInspectionCodec),
21+
[typeof(string[])] = typeof(ReplicaVoterArrayInspectionCodec),
22+
[typeof(string)] = typeof(ReplicaBorrowedStringCodec),
23+
[typeof(ReplicaSnapshot)] = typeof(ReplicaSnapshotInspectionCodec),
24+
[typeof(ReplicatedOperation)] = typeof(ReplicaOperationInspectionCodec)
25+
};
26+
private static readonly HashSet<Type> Scalars =
27+
[
28+
typeof(int), typeof(long), typeof(Guid), typeof(DateTimeOffset),
29+
typeof(OperationKind), typeof(ErrorCode), typeof(ErrorCode?)
30+
];
31+
1032
internal static string Sender<TInput>(ref Reader<TInput> reader)
1133
=> ReplicaInspectionBuffers.For(reader.Session).Sender(Read<string, TInput>(ref reader, 0));
1234

1335
internal static T Read<T, TInput>(ref Reader<TInput> reader, uint delta)
1436
{
1537
var field = reader.ReadFieldHeader();
1638
Require(field, delta, ExpectedFieldType<T>());
17-
return reader.Session.CodecProvider.GetCodec<T>().ReadValue(ref reader, field)!;
39+
return ProfileCodec<T>(reader.Session).ReadValue(ref reader, field)!;
40+
}
41+
42+
private static IFieldCodec<T> ProfileCodec<T>(SerializerSession session)
43+
{
44+
if (ProfileCodecs.TryGetValue(typeof(T), out var codec))
45+
{
46+
return (IFieldCodec<T>)session.CodecProvider.Services.GetRequiredService(codec);
47+
}
48+
if (!Scalars.Contains(typeof(T)))
49+
{
50+
throw Errors.Fail(ErrorCode.Corruption, ReplicaPersistence.InvalidEncoding);
51+
}
52+
return session.CodecProvider.GetCodec<T>();
1853
}
1954

2055
// These two current ungenerated contracts use Orleans' backing-integer codec.
@@ -63,5 +98,5 @@ internal static void EndBase<TInput>(ref Reader<TInput> reader)
6398

6499
internal static void Write<T, TBufferWriter>(ref Writer<TBufferWriter> writer, uint delta, T value)
65100
where TBufferWriter : IBufferWriter<byte>
66-
=> writer.Session.CodecProvider.GetCodec<T>().WriteField(ref writer, delta, typeof(T), value);
101+
=> ProfileCodec<T>(writer.Session).WriteField(ref writer, delta, typeof(T), value);
67102
}

‎src/KeyLoad.Replication/Features/ClusterReplication/ReplicaNativeInspection.cs‎

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,8 +52,13 @@ private static void Configure(ISerializerBuilder builder)
5252
Register<ReplicaBorrowedBytesCodec>(builder);
5353
Register<ReplicaEntryInspectionCodec>(builder);
5454
Register<ReplicaOperationInspectionCodec>(builder);
55-
Register<ReplicaEntryArrayInspectionCodec>(builder);
56-
Register<ReplicaVoterArrayInspectionCodec>(builder);
55+
builder.Services.AddSingleton(provider => new ReplicaEntryArrayInspectionCodec(provider.GetRequiredService<ReplicaEntryInspectionCodec>()));
56+
builder.Services.AddSingleton(provider => new ReplicaVoterArrayInspectionCodec(provider.GetRequiredService<ReplicaBorrowedStringCodec>()));
57+
builder.Configure(options =>
58+
{
59+
options.FieldCodecs.Add(typeof(ReplicaEntryArrayInspectionCodec));
60+
options.FieldCodecs.Add(typeof(ReplicaVoterArrayInspectionCodec));
61+
});
5762
Register<ReplicaEntriesInspectionCodec>(builder);
5863
Register<ReplicaVotersInspectionCodec>(builder);
5964
Register<ReplicaEntryBatchInspectionCodec>(builder);

‎src/KeyLoad.Replication/Features/ClusterReplication/ReplicaTermObservation.cs‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ internal readonly struct ReplicaTermObservationRead(long term, ReplicaTermObserv
3232

3333
internal static class ReplicaTermObservationReader
3434
{
35-
internal static ReplicaTermObservationRead Read(IAtomicStore store, long index, long hardStateTerm,
35+
internal static ReplicaTermObservationRead Read(IAtomicStore store, long index, long hardStateTerm, int maximumEntries,
3636
ReplicaTermObservation? current)
3737
{
3838
return store.Read<ReplicaTermObservationRead>(view =>
@@ -55,7 +55,7 @@ internal static ReplicaTermObservationRead Read(IAtomicStore store, long index,
5555
long term = 0;
5656
var found = view.ReadValue(ReplicaProtocol.EntryStorageKey(index), bytes =>
5757
{
58-
var entry = ReplicaProtocolCodec.Deserialize<ReplicaEntry>(bytes);
58+
var entry = ReplicaProtocolCodec.DeserializeStored<ReplicaEntry>(bytes.ToArray(), maximumEntries);
5959
if (entry.Index != index || entry.Term <= 0 || entry.Term > hardStateTerm)
6060
{
6161
throw Errors.Fail(ErrorCode.Corruption, ReplicaProtocol.CorruptLog);
Lines changed: 69 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,69 @@
1+
using System.Buffers.Binary;
2+
using System.Security.Cryptography;
3+
using static KeyLoad.Storage.ZoneTree.ZoneTreePersistenceFormat;
4+
5+
namespace KeyLoad.Storage.ZoneTree;
6+
7+
// Verify the whole recoverable prefix before opening a provider which can rewrite its metadata.
8+
// The caller holds the owner lock and retains this same journal handle through ordered recovery.
9+
internal static class ZoneTreeJournalPreflight
10+
{
11+
internal static void Validate(FileStream journal, ZoneTreeStoreOptions options, int identityVersion)
12+
{
13+
try
14+
{
15+
journal.Position = 0;
16+
var header = new byte[HeaderLength];
17+
var position = ReadCheckpoint(journal, options, header);
18+
while (journal.Length - journal.Position >= HeaderLength)
19+
{
20+
if (!ReadFrame(journal, options, identityVersion, header, ref position))
21+
{
22+
break;
23+
}
24+
}
25+
}
26+
finally
27+
{
28+
// Only ordinary recovery applies verified records and truncates an incomplete current tail.
29+
journal.Position = 0;
30+
}
31+
}
32+
33+
private static long ReadCheckpoint(FileStream journal, ZoneTreeStoreOptions options, byte[] header)
34+
{
35+
if (journal.Length < HeaderLength)
36+
{
37+
return 0;
38+
}
39+
journal.ReadExactly(header);
40+
journal.Position = 0;
41+
if (BinaryPrimitives.ReadUInt64LittleEndian(header) != CheckpointMagic)
42+
{
43+
return 0;
44+
}
45+
// Normal recovery also supplies an apply callback. A no-op preserves those same bounds;
46+
// the complete journal may exceed the standalone snapshot limit because it contains a tail.
47+
return ZoneTreeCheckpointReader.Read(journal, options, static _ => { }).Position;
48+
}
49+
50+
private static bool ReadFrame(FileStream journal, ZoneTreeStoreOptions options, int identityVersion,
51+
byte[] header, ref long position)
52+
{
53+
journal.ReadExactly(header);
54+
var (length, sequence) = ZoneTreeJournalRecovery.ValidateHeader(header, position, options.MaxFrameBytes, identityVersion);
55+
if (journal.Length - journal.Position < length)
56+
{
57+
return false;
58+
}
59+
var payload = new byte[length];
60+
journal.ReadExactly(payload);
61+
if (!CryptographicOperations.FixedTimeEquals(SHA256.HashData(payload), header.AsSpan(ChecksumOffset, ChecksumLength)))
62+
{
63+
throw Errors.Fail(ErrorCode.Corruption, JournalChecksumInvalid);
64+
}
65+
_ = ZoneTreeJournalCodec.Deserialize(payload);
66+
position = sequence;
67+
return true;
68+
}
69+
}

‎src/KeyLoad.Storage.ZoneTree/Features/StorageRecovery/ZoneTreeStoreInitializer.cs‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,9 @@ internal static void Open(ZoneTreeStoreRuntime runtime)
1212
runtime.Ownership = new FileStream(Path.Combine(runtime.Options.Directory, OwnerLockFileName),
1313
FileMode.OpenOrCreate, FileAccess.ReadWrite, FileShare.None);
1414
runtime.Identity = ZoneTreeIdentityFile.Open(runtime.Options, runtime.Ownership);
15-
runtime.Tree = ZoneTreeTreeFactory.Open(runtime.Options);
1615
runtime.Journal = ZoneTreeStoreFiles.OpenJournal(runtime.Options);
16+
ZoneTreeJournalPreflight.Validate(runtime.Journal, runtime.Options, runtime.Identity.FormatVersion);
17+
runtime.Tree = ZoneTreeTreeFactory.Open(runtime.Options);
1718
ZoneTreeJournalRecovery.Recover(runtime);
1819
runtime.Maintainer = runtime.Tree.CreateMaintainer();
1920
ZoneTreeCheckpointReclaimer.Reclaim(runtime.Options.Directory);
@@ -31,6 +32,7 @@ internal static void OpenExisting(ZoneTreeStoreRuntime runtime, Guid expectedNod
3132
FileMode.Open, FileAccess.ReadWrite, FileShare.None);
3233
runtime.Identity = ZoneTreeIdentityFile.OpenExisting(runtime.Options, expectedNodeId);
3334
runtime.Journal = ZoneTreeStoreFiles.OpenJournal(runtime.Options, FileMode.Open);
35+
ZoneTreeJournalPreflight.Validate(runtime.Journal, runtime.Options, runtime.Identity.FormatVersion);
3436
runtime.Tree = ZoneTreeTreeFactory.Open(runtime.Options, requireExisting: true);
3537
ZoneTreeJournalRecovery.Recover(runtime);
3638
}

‎tests/KeyLoad.RecoveryTests/Features/ClusterReplication/ReplicaSnapshotMetadataRecoveryTests.cs‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,6 @@ namespace KeyLoad.RecoveryTests;
66
/// <summary>AC-REP-004: real private-transfer metadata failures preserve the canonical cut and permit safe retry.</summary>
77
internal sealed class ReplicaSnapshotMetadataRecoveryTests
88
{
9-
private static readonly byte[] TruncatedManifest = "{"u8.ToArray();
10-
119
/// <summary>A foreign-incarnation descriptor is rejected before private files or canonical state are changed.</summary>
1210
[Test]
1311
public async Task ForeignIncarnationSnapshotIsRejectedWithoutEffectsAndValidDescriptorCanRetry()
@@ -42,7 +40,9 @@ public async Task TruncatedPrivateManifestIsDiscardedAndVerifiedTransferCanRetry
4240
var receiver = trial.Receiver();
4341
await Assert.That(receiver.Begin(trial.Image)).IsEqualTo(0);
4442
var directory = Path.GetDirectoryName(trial.TargetImagePath(trial.Image))!;
45-
await File.WriteAllBytesAsync(Path.Combine(directory, ReplicaProtocol.IncomingManifest), TruncatedManifest,
43+
var manifestPath = Path.Combine(directory, ReplicaProtocol.IncomingManifest);
44+
var manifest = await File.ReadAllBytesAsync(manifestPath, TestContext.Current!.Execution.CancellationToken);
45+
await File.WriteAllBytesAsync(manifestPath, manifest[..^1],
4646
TestContext.Current!.Execution.CancellationToken);
4747

4848
receiver.Recover();

0 commit comments

Comments
 (0)