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
134 changes: 134 additions & 0 deletions src/NatsDistributedCache/CacheEntryBinarySerializer.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,134 @@
using System.Buffers;
using System.Buffers.Binary;
using NATS.Client.Core;

namespace CodeCargo.Nats.DistributedCache;

/// <summary>
/// Serializes and deserializes <see cref="CacheEntry"/> using a compact binary framing that avoids the
/// base64 inflation and JSON overhead of the previous envelope.
/// </summary>
/// <remarks>
/// Wire format (little-endian):
/// <code>
/// [version:1][flags:1][absExpTicks:8?][sldExpTicks:8?][payload...]
/// </code>
/// The 8-byte expiration fields are only present when their corresponding flag bit is set, so an entry
/// without expiration carries just a 2-byte header. Any buffer whose leading version byte does not match
/// <see cref="FormatVersion"/> — including pre-existing JSON entries, whose first byte is '{' (0x7B) —
/// deserializes to <c>null</c> and is treated as a cache miss.
/// </remarks>
internal sealed class CacheEntryBinarySerializer : INatsSerialize<CacheEntry>, INatsDeserialize<CacheEntry>
{
/// <summary>
/// Current on-wire format version. Bump when the framing changes so that entries written by an older
/// version are treated as cache misses instead of being misinterpreted.
/// </summary>
internal const byte FormatVersion = 1;

private const byte HasAbsoluteExpiration = 0b0000_0001;
private const byte HasSlidingExpiration = 0b0000_0010;
private const byte KnownFlags = HasAbsoluteExpiration | HasSlidingExpiration;

/// <summary>
/// Gets the shared, stateless serializer instance.
/// </summary>
public static CacheEntryBinarySerializer Default { get; } = new();

/// <inheritdoc />
public void Serialize(IBufferWriter<byte> bufferWriter, CacheEntry value)
{
var hasAbsolute = value.AbsoluteExpiration.HasValue;
var hasSliding = value.SlidingExpirationTicks.HasValue;

var flags = (byte)0;
if (hasAbsolute)
{
flags |= HasAbsoluteExpiration;
}

if (hasSliding)
{
flags |= HasSlidingExpiration;
}

var headerLength = 2 + (hasAbsolute ? 8 : 0) + (hasSliding ? 8 : 0);
var header = bufferWriter.GetSpan(headerLength);
header[0] = FormatVersion;
header[1] = flags;

var offset = 2;
if (hasAbsolute)
{
BinaryPrimitives.WriteInt64LittleEndian(header.Slice(offset), value.AbsoluteExpiration!.Value.UtcTicks);
offset += 8;
}

if (hasSliding)
{
BinaryPrimitives.WriteInt64LittleEndian(header.Slice(offset), value.SlidingExpirationTicks!.Value);
offset += 8;
}

bufferWriter.Advance(offset);

if (value.Data is { Length: > 0 } data)
{
bufferWriter.Write(data);
}
}

/// <inheritdoc />
public CacheEntry? Deserialize(in ReadOnlySequence<byte> buffer)
{
var reader = new SequenceReader<byte>(buffer);

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

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;
}

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;
}

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

long? slidingExpirationTicks = null;
if ((flags & HasSlidingExpiration) != 0)
{
if (!reader.TryReadLittleEndian(out long slidingTicks))
{
return null;
}

slidingExpirationTicks = slidingTicks;
}

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

return new CacheEntry
{
AbsoluteExpiration = absoluteExpiration,
SlidingExpirationTicks = slidingExpirationTicks,
Data = data,
};
}
}
17 changes: 2 additions & 15 deletions src/NatsDistributedCache/NatsCache.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
using System.Buffers;
using System.Text.Json.Serialization;
using Microsoft.Extensions.Caching.Distributed;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
Expand All @@ -15,32 +14,20 @@ namespace CodeCargo.Nats.DistributedCache;
/// </summary>
public class CacheEntry
{
[JsonPropertyName("absexp")]
public DateTimeOffset? AbsoluteExpiration { get; set; }

[JsonPropertyName("sldexp")]
public long? SlidingExpirationTicks { get; set; }

[JsonPropertyName("data")]
public byte[]? Data { get; set; }
}

/// <summary>
/// JsonSerializerContext for CacheEntry
/// </summary>
[JsonSerializable(typeof(CacheEntry))]
public partial class CacheEntryJsonContext : JsonSerializerContext
{
}

/// <summary>
/// Distributed cache implementation using NATS Key-Value Store.
/// </summary>
public partial class NatsCache : IBufferDistributedCache
{
// Static JSON serializer for CacheEntry
private static readonly NatsJsonContextSerializer<CacheEntry> CacheEntrySerializer =
new(CacheEntryJsonContext.Default);
// Compact binary serializer for the CacheEntry envelope (replaces the previous JSON+base64 format).
private static readonly CacheEntryBinarySerializer CacheEntrySerializer = CacheEntryBinarySerializer.Default;

private readonly string _bucketName;
private readonly INatsCacheKeyEncoder _keyEncoder;
Expand Down
1 change: 1 addition & 0 deletions src/NatsDistributedCache/NatsDistributedCache.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@
.snk into keys/ locally flips the build to signed and will break the UnitTests compile. -->
<ItemGroup Condition="'$(SignAssembly)' != 'true'">
<InternalsVisibleTo Include="UnitTests" />
<InternalsVisibleTo Include="Benchmarks" />
</ItemGroup>

</Project>
11 changes: 6 additions & 5 deletions test/IntegrationTests/Cache/HybridCacheSetAndRemoveTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -70,12 +70,13 @@ public async Task HybridCacheSerializesDateTime()
var serializedDateString = Encoding.ASCII.GetString(writer.WrittenSpan.ToArray());

var kvStore = await NatsConnection.CreateKeyValueStoreContext().GetStoreAsync("cache");
NatsJsonContextSerializer<CacheEntry> cacheEntrySerializer = new(CacheEntryJsonContext.Default);
var kvEntry = await kvStore.GetEntryAsync(key, serializer: cacheEntrySerializer);
Assert.NotNull(kvEntry.Value?.Data);
var storedDateString = Encoding.ASCII.GetString(kvEntry.Value.Data);
var kvEntry = await kvStore.GetEntryAsync<byte[]>(key);
Assert.NotNull(kvEntry.Value);
var storedDateString = Encoding.ASCII.GetString(kvEntry.Value);

// HybridCache adds additional data to the front of the serialized value, so we're matching only the relevant data
// The compact binary envelope stores the payload raw (no base64/JSON), so the HybridCache-serialized
// value appears verbatim in the stored bytes. HybridCache also adds framing of its own to the front,
// so we match only the relevant portion.
Assert.Contains(serializedDateString, storedDateString);

// Assert - date is deserialized as expected
Expand Down
Loading
Loading