Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 38 additions & 15 deletions src/NatsDistributedCache/CacheEntryBinarySerializer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -98,35 +98,66 @@ public void Serialize(IBufferWriter<byte> bufferWriter, CacheEntry value)
public CacheEntry? Deserialize(in ReadOnlySequence<byte> buffer)
{
var reader = new SequenceReader<byte>(buffer);
if (!TryReadHeader(ref reader, out var absoluteExpiration, out var slidingExpirationTicks))
{
// Unknown/legacy/corrupt framing: treat as a cache miss (see TryReadHeader).
return null;
}

var remaining = reader.UnreadSequence;
var data = remaining.IsEmpty ? Array.Empty<byte>() : remaining.ToArray();

return new CacheEntry
{
AbsoluteExpiration = absoluteExpiration,
SlidingExpirationTicks = slidingExpirationTicks,
Data = data,
};
}

/// <summary>
/// Reads and validates the fixed <see cref="CacheEntry"/> header (version, flags, and the optional
/// expiration fields), advancing <paramref name="reader"/> to the first payload byte. Returns
/// <see langword="false"/> for any unknown, legacy (e.g. JSON), or corrupt framing so callers can fail
/// closed to a cache miss instead of misinterpreting the remaining bytes. Shared by
/// <see cref="Deserialize"/> and the single-copy read path in
/// <see cref="BufferWritingCacheEntryDeserializer"/> so both agree, byte for byte, on what a valid entry
/// is and where its payload begins.
/// </summary>
internal static bool TryReadHeader(
ref SequenceReader<byte> reader,
out DateTimeOffset? absoluteExpiration,
out long? slidingExpirationTicks)
{
absoluteExpiration = null;
slidingExpirationTicks = null;

if (!reader.TryRead(out var version) || version != FormatVersion)
{
// Unknown or legacy (e.g. JSON) entry: treat as a cache miss.
return null;
return false;
}

if (!reader.TryRead(out var flags) || (flags & ~KnownFlags) != 0)
{
// Unknown flag bits: corrupt data, or a future format that mistakenly reused this version.
// Fail closed rather than misinterpreting the remaining bytes.
return null;
return false;
}

DateTimeOffset? absoluteExpiration = null;
if ((flags & HasAbsoluteExpiration) != 0)
{
if (!reader.TryReadLittleEndian(out long absoluteTicks) ||
absoluteTicks < DateTimeOffset.MinValue.Ticks ||
absoluteTicks > DateTimeOffset.MaxValue.Ticks)
{
// Missing or out-of-range ticks (corrupt entry): treat as a miss instead of throwing.
return null;
return false;
}

absoluteExpiration = new DateTimeOffset(absoluteTicks, TimeSpan.Zero);
}

long? slidingExpirationTicks = null;
if ((flags & HasSlidingExpiration) != 0)
{
if (!reader.TryReadLittleEndian(out long slidingTicks) ||
Expand All @@ -137,20 +168,12 @@ public void Serialize(IBufferWriter<byte> bufferWriter, CacheEntry value)
// a miss, like the absolute-ticks bounds check above. A valid sliding window is always
// a positive TimeSpan within (0, MaxTtlTicks]; the write path enforces the same ceiling,
// so any legitimately written value round-trips.
return null;
return false;
}

slidingExpirationTicks = slidingTicks;
}

var remaining = reader.UnreadSequence;
var data = remaining.IsEmpty ? Array.Empty<byte>() : remaining.ToArray();

return new CacheEntry
{
AbsoluteExpiration = absoluteExpiration,
SlidingExpirationTicks = slidingExpirationTicks,
Data = data,
};
return true;
}
}
172 changes: 172 additions & 0 deletions src/NatsDistributedCache/CacheEntryReadDeserializer.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
using System.Buffers;
using NATS.Client.Core;

namespace CodeCargo.Nats.DistributedCache;

/// <summary>
/// What a single <see cref="CacheEntry"/> read produced, as seen by <c>NatsCache</c>'s read core. A
/// <see langword="null"/> result (the deserializer returning <c>null</c>) means the stored bytes were
/// undeserializable — the same way <see cref="CacheEntryBinarySerializer.Deserialize"/> signals it with a
/// null entry.
/// </summary>
internal enum CacheEntryReadOutcome
{
/// <summary>
/// The payload was materialized into <see cref="CacheEntryReadResult.Entry"/>'s <see cref="CacheEntry.Data"/>.
/// The read core checks absolute expiry, refreshes a sliding TTL, then emits the payload — returning it for
/// the array path, or writing it to the destination for a sliding buffer read. Used for every array read
/// and for sliding entries on the buffer path.
/// </summary>
Materialized,

/// <summary>
/// Buffer fast path: the payload was written straight into the caller's <see cref="IBufferWriter{T}"/>.
/// A confirmed hit; the core only records it.
/// </summary>
PayloadWritten,

/// <summary>
/// Buffer fast path: the entry is past its absolute expiration. Nothing was written; the core evicts the
/// entry and reports a miss.
/// </summary>
AbsolutelyExpired,

/// <summary>
/// The caller's <see cref="IBufferWriter{T}"/> threw while the payload was being written (e.g. HybridCache's
/// payload quota). Not corrupt data: the exception is carried out so the core surfaces an error instead of
/// an undeserializable miss.
/// </summary>
DestinationFailure,
}

/// <summary>
/// Outcome of a single read produced by <see cref="CacheEntryReadDeserializer"/>. The reusable
/// <see cref="PayloadWritten"/> and <see cref="AbsolutelyExpired"/> singletons let the common buffer hit and
/// the expired path allocate nothing beyond the payload itself; only the materialized and destination-failure
/// outcomes carry per-read state.
/// </summary>
internal sealed class CacheEntryReadResult
{
internal static readonly CacheEntryReadResult PayloadWritten =
new(CacheEntryReadOutcome.PayloadWritten, entry: null, destinationFailure: null);

internal static readonly CacheEntryReadResult AbsolutelyExpired =
new(CacheEntryReadOutcome.AbsolutelyExpired, entry: null, destinationFailure: null);

private CacheEntryReadResult(CacheEntryReadOutcome outcome, CacheEntry? entry, Exception? destinationFailure)
{
Outcome = outcome;
Entry = entry;
DestinationFailure = destinationFailure;
}

internal CacheEntryReadOutcome Outcome { get; }

/// <summary>The decoded entry, non-null only for <see cref="CacheEntryReadOutcome.Materialized"/>.</summary>
internal CacheEntry? Entry { get; }

/// <summary>The writer exception, non-null only for <see cref="CacheEntryReadOutcome.DestinationFailure"/>.</summary>
internal Exception? DestinationFailure { get; }

internal static CacheEntryReadResult Materialized(CacheEntry entry) =>
new(CacheEntryReadOutcome.Materialized, entry, destinationFailure: null);

internal static CacheEntryReadResult DestinationFailed(Exception exception) =>
new(CacheEntryReadOutcome.DestinationFailure, entry: null, exception);
}

/// <summary>
/// The per-read <see cref="INatsDeserialize{T}"/> used by <c>NatsCache</c>'s unified read core.
/// </summary>
/// <remarks>
/// With a <see langword="null"/> destination it materializes the entry for the array read path (Get/Refresh).
/// With a destination it powers the single-copy <c>TryGetAsync(IBufferWriter&lt;byte&gt;)</c> path: on a genuine
/// hit with no sliding refresh it writes the payload straight into the caller's writer, so the read costs one
/// copy (transport buffer -> caller buffer) instead of two. Nothing is written for a miss (undeserializable
/// framing or absolute expiry) or when the entry needs its sliding TTL refreshed — that re-writes the value and
/// so needs the bytes, which are materialized for the core to handle. If the caller's writer throws mid-write
/// the exception is carried out as <see cref="CacheEntryReadOutcome.DestinationFailure"/> so the core surfaces
/// an error rather than an undeserializable miss; on a multi-segment payload any segments already written stay
/// committed, so the buffer contents are unspecified when a write fails.
/// </remarks>
internal sealed class CacheEntryReadDeserializer : INatsDeserialize<CacheEntryReadResult>
{
private readonly IBufferWriter<byte>? _destination;
private readonly TimeProvider _timeProvider;

internal CacheEntryReadDeserializer(IBufferWriter<byte>? destination, TimeProvider timeProvider)
{
_destination = destination;
_timeProvider = timeProvider;
}

/// <inheritdoc />
public CacheEntryReadResult? Deserialize(in ReadOnlySequence<byte> buffer)
{
// Array read path: no destination to stream into, so materialize the whole entry via the canonical
// deserializer and let the read core decide expiry/refresh/eviction.
if (_destination is null)
{
var entry = CacheEntryBinarySerializer.Default.Deserialize(buffer);
return entry is null ? null : CacheEntryReadResult.Materialized(entry);
}

var reader = new SequenceReader<byte>(buffer);
if (!CacheEntryBinarySerializer.TryReadHeader(ref reader, out var absoluteExpiration, out var slidingExpirationTicks))
{
// Unknown/legacy/corrupt framing: signal an undeserializable entry.
return null;
}

var remaining = reader.UnreadSequence;

if (slidingExpirationTicks.HasValue)
{
// Sliding entries have their TTL refreshed by re-writing the value, which needs the payload
// bytes, so this path cannot stream to the caller. Materialize and let the core handle it,
// exactly like the array path.
var data = remaining.IsEmpty ? Array.Empty<byte>() : remaining.ToArray();
return CacheEntryReadResult.Materialized(new CacheEntry
{
AbsoluteExpiration = absoluteExpiration,
SlidingExpirationTicks = slidingExpirationTicks,
Data = data,
});
}

if (NatsCache.IsAbsolutelyExpired(absoluteExpiration, _timeProvider.GetUtcNow()))
{
// Absolutely expired: write nothing. The core evicts and reports a miss. Expiry is evaluated
// here (the only point holding the payload) so the buffer stays untouched for an expired miss,
// using the same clock and predicate the core would.
return CacheEntryReadResult.AbsolutelyExpired;
}

// Genuine hit with no sliding refresh: copy the payload straight into the caller's buffer — the
// single copy that makes TryGet free of any intermediate array on the common path. A writer failure
// (e.g. HybridCache's payload quota) is not corrupt data, so carry it out for the core to surface.
try
{
WritePayload(remaining);
}
catch (Exception ex)
{
return CacheEntryReadResult.DestinationFailed(ex);
}

return CacheEntryReadResult.PayloadWritten;
}

private void WritePayload(in ReadOnlySequence<byte> payload)
{
// The enumerator is a struct, and it yields the single segment for a single-segment sequence, so this
// covers both without a special case.
foreach (var segment in payload)
{
if (!segment.IsEmpty)
{
_destination!.Write(segment.Span);
}
}
}
}
3 changes: 2 additions & 1 deletion src/NatsDistributedCache/NatsCache.Log.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,10 @@ private void LogException(Exception exception) =>
private void LogSwallowedException(Exception exception) =>
_logger.LogWarning(EventIds.Exception, exception, "NATS cache read failed in TryGetAsync; returning a cache miss");

private void LogUndeserializableEntry(string key) =>
private void LogUndeserializableEntry(string key, Exception? error) =>
_logger.LogDebug(
EventIds.UndeserializableEntry,
error,
"Cache entry for key {Key} could not be deserialized (legacy or corrupt format); returning a cache miss",
key);

Expand Down
Loading
Loading