+using System.Text.Json.Serialization;
+namespace VNLib.Data.Caching.Extensions
+ public class ActiveServer
+ {
+ [JsonPropertyName("address")]
+ public string? HostName { get; set; }
+ [JsonPropertyName("server_id")]
+ public string? ServerId { get; set; }
+ [JsonPropertyName("ip_address")]
+ public string? Ip { get; set; }
+ }
+using System;
+using System.Net;
+using System.Text;
+using System.Security;
+using System.Text.Json;
+using System.Security.Cryptography;
+using System.Runtime.CompilerServices;
+using RestSharp;
+using VNLib.Utils.Memory;
+using VNLib.Utils.Logging;
+using VNLib.Utils.Extensions;
+using VNLib.Hashing;
+using VNLib.Hashing.IdentityUtility;
+using VNLib.Net.Http;
+using VNLib.Net.Rest.Client;
+using VNLib.Net.Messaging.FBM.Client;
+using VNLib.Net.Messaging.FBM;
+namespace VNLib.Data.Caching.Extensions
+ /// <summary>
+ /// Provides extension methods for FBM data caching using
+ /// cache servers and brokers
+ /// </summary>
+ public static class FBMDataCacheExtensions
+ {
+ /// <summary>
+ /// The websocket sub-protocol to use when connecting to cache servers
+ /// </summary>
+ public const string CACHE_WS_SUB_PROCOL = "object-cache";
+ /// <summary>
+ /// The default cache message header size
+ /// </summary>
+ public const int MAX_FBM_MESSAGE_HEADER_SIZE = 1024;
+ private static readonly IReadOnlyDictionary<string, string> BrokerJwtHeader = new Dictionary<string, string>()
+ {
+ { "alg", "ES384" }, //Must match alg name
+ { "typ", "JWT"}
+ };
+ private static readonly RestClientPool ClientPool = new(2,new RestClientOptions()
+ {
+ MaxTimeout = 10 * 1000,
+ FollowRedirects = false,
+ Encoding = Encoding.UTF8,
+ AutomaticDecompression = DecompressionMethods.All,
+ ThrowOnAnyError = true,
+ });
+ /// <summary>
+ /// The default hashing algorithm used to sign an verify connection
+ /// tokens
+ /// </summary>
+ public static readonly HashAlgorithmName CacheJwtAlgorithm = HashAlgorithmName.SHA384;
+ //using the es384 algorithm for signing (friendlyname is secp384r1)
+ /// <summary>
+ /// The default ECCurve used by the connection library
+ /// </summary>
+ public static readonly ECCurve CacheCurve = ECCurve.CreateFromFriendlyName("secp384r1");
+ /// <summary>
+ /// Gets a <see cref="FBMClientConfig"/> preconfigured object caching
+ /// protocl
+ /// </summary>
+ /// <param name="heap">The client buffer heap</param>
+ /// <param name="maxMessageSize">The maxium message size (in bytes)</param>
+ /// <param name="debugLog">An optional debug log</param>
+ /// <returns>A preconfigured <see cref="FBMClientConfig"/> for object caching</returns>
+ public static FBMClientConfig GetDefaultConfig(IUnmangedHeap heap, int maxMessageSize, ILogProvider? debugLog = null)
+ {
+ return new()
+ {
+ BufferHeap = heap,
+ MaxMessageSize = maxMessageSize * 2,
+ RecvBufferSize = maxMessageSize,
+ MessageBufferSize = maxMessageSize,
+ SubProtocol = CACHE_WS_SUB_PROCOL,
+ HeaderEncoding = Helpers.DefaultEncoding,
+ KeepAliveInterval = TimeSpan.FromSeconds(30),
+ DebugLog = debugLog
+ };
+ }
+ private class CacheConnectionConfig
+ {
+ public ECDsa ClientAlg { get; init; }
+ public ECDsa BrokerAlg { get; init; }
+ public string ServerChallenge { get; init; }
+ public string? NodeId { get; set; }
+ public Uri? BrokerAddress { get; set; }
+ public bool useTls { get; set; }
+ public ActiveServer[]? BrokerServers { get; set; }
+ public CacheConnectionConfig()
+ {
+ //Init the algorithms
+ ClientAlg = ECDsa.Create(CacheCurve);
+ BrokerAlg = ECDsa.Create(CacheCurve);
+ ServerChallenge = RandomHash.GetRandomBase32(24);
+ }
+ ~CacheConnectionConfig()
+ {
+ ClientAlg.Clear();
+ BrokerAlg.Clear();
+ }
+ }
+ /// <summary>
+ /// Contacts the cache broker to get a list of active servers to connect to
+ /// </summary>
+ /// <param name="brokerAddress">The broker server to connec to</param>
+ /// <param name="clientPrivKey">The private key used to sign messages sent to the broker</param>
+ /// <param name="brokerPubKey">The broker public key used to verify broker messages</param>
+ /// <param name="cancellationToken">A token to cancel the operationS</param>
+ /// <returns>The list of active servers</returns>
+ /// <exception cref="SecurityException"></exception>
+ /// <exception cref="ArgumentNullException"></exception>
+ public static async Task<ActiveServer[]?> ListServersAsync(Uri brokerAddress, ReadOnlyMemory<byte> clientPrivKey, ReadOnlyMemory<byte> brokerPubKey, CancellationToken cancellationToken = default)
+ {
+ using ECDsa client = ECDsa.Create(CacheCurve);
+ using ECDsa broker = ECDsa.Create(CacheCurve);
+ //Import client private key
+ client.ImportPkcs8PrivateKey(clientPrivKey.Span, out _);
+ //Broker public key to verify broker messages
+ broker.ImportSubjectPublicKeyInfo(brokerPubKey.Span, out _);
+ return await ListServersAsync(brokerAddress, client, broker, cancellationToken);
+ }
+ /// <summary>
+ /// Contacts the cache broker to get a list of active servers to connect to
+ /// </summary>
+ /// <param name="brokerAddress">The broker server to connec to</param>
+ /// <param name="clientAlg">The signature algorithm used to sign messages to the broker</param>
+ /// <param name="brokerAlg">The signature used to verify broker messages</param>
+ /// <param name="cancellationToken">A token to cancel the operationS</param>
+ /// <returns>The list of active servers</returns>
+ /// <exception cref="SecurityException"></exception>
+ /// <exception cref="ArgumentNullException"></exception>
+ public static async Task<ActiveServer[]?> ListServersAsync(Uri brokerAddress, ECDsa clientAlg, ECDsa brokerAlg, CancellationToken cancellationToken = default)
+ {
+ _ = brokerAddress ?? throw new ArgumentNullException(nameof(brokerAddress));
+ _ = clientAlg ?? throw new ArgumentNullException(nameof(clientAlg));
+ _ = brokerAlg ?? throw new ArgumentNullException(nameof(brokerAlg));
+ string jwtBody;
+ //Build request jwt
+ using (JsonWebToken requestJwt = new())
+ {
+ requestJwt.WriteHeader(BrokerJwtHeader);
+ requestJwt.InitPayloadClaim()
+ .AddClaim("iat", DateTimeOffset.UtcNow.ToUnixTimeMilliseconds())
+ .CommitClaims();
+ //sign the jwt
+ requestJwt.Sign(clientAlg, in CacheJwtAlgorithm, 512);
+ //Compile the jwt
+ jwtBody = requestJwt.Compile();
+ }
+ //New list request
+ RestRequest listRequest = new(brokerAddress, Method.Post);
+ //Add the jwt as a string to the request body
+ listRequest.AddStringBody(jwtBody, DataFormat.None);
+ listRequest.AddHeader("Content-Type", HttpHelpers.GetContentTypeString(ContentType.Text));
+ //Rent client
+ using ClientContract client = ClientPool.Lease();
+ //Exec list request
+ RestResponse response = await client.Resource.ExecuteAsync(listRequest, cancellationToken);
+ if (!response.IsSuccessful)
+ {
+ throw response.ErrorException!;
+ }
+ //Response is jwt
+ using JsonWebToken responseJwt = JsonWebToken.ParseRaw(response.RawBytes);
+ //Verify the jwt
+ if (!responseJwt.Verify(brokerAlg, in CacheJwtAlgorithm))
+ {
+ throw new SecurityException("Failed to verify the broker's challenge, cannot continue");
+ }
+ using JsonDocument doc = responseJwt.GetPayload();
+ return doc.RootElement.GetProperty("servers").Deserialize<ActiveServer[]>();
+ }
+ /// <summary>
+ /// Configures a connection to the remote cache server at the specified location
+ /// with proper authentication.
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="serverUri">The server's address</param>
+ /// <param name="signingKey">The pks8 format EC private key uesd to sign the message</param>
+ /// <param name="challenge">A challenge to send to the server</param>
+ /// <param name="nodeId">A token used to identify the current server's event queue on the remote server</param>
+ /// <param name="token">A token to cancel the connection operation</param>
+ /// <param name="useTls">Enables the secure websocket protocol</param>
+ /// <returns>A Task that completes when the connection has been established</returns>
+ /// <exception cref="ArgumentNullException"></exception>
+ public static Task ConnectAsync(this FBMClient client, string serverUri, ReadOnlyMemory<byte> signingKey, string challenge, string? nodeId, bool useTls, CancellationToken token = default)
+ {
+ //Sign the jwt
+ using ECDsa sigAlg = ECDsa.Create(CacheCurve);
+ //Import the signing key
+ sigAlg.ImportPkcs8PrivateKey(signingKey.Span, out _);
+ //Return without await because the alg is used to sign before this method returns and can be discarded
+ return ConnectAsync(client, serverUri, sigAlg, challenge, nodeId, useTls, token);
+ }
+ private static Task ConnectAsync(FBMClient client, string serverUri, ECDsa sigAlg, string challenge, string? nodeId, bool useTls, CancellationToken token = default)
+ {
+ _ = serverUri ?? throw new ArgumentNullException(nameof(serverUri));
+ _ = challenge ?? throw new ArgumentNullException(nameof(challenge));
+ //build ws uri
+ UriBuilder uriBuilder = new(serverUri)
+ {
+ Scheme = useTls ? "wss://" : "ws://"
+ };
+ string jwtMessage;
+ //Init jwt for connecting to server
+ using (JsonWebToken jwt = new())
+ {
+ jwt.WriteHeader(BrokerJwtHeader);
+ //Init claim
+ JwtPayload claim = jwt.InitPayloadClaim();
+ claim.AddClaim("challenge", challenge);
+ if (!string.IsNullOrWhiteSpace(nodeId))
+ {
+ /*
+ * The unique node id so the other nodes know to load the
+ * proper event queue for the current server
+ */
+ claim.AddClaim("server_id", nodeId);
+ }
+ claim.CommitClaims();
+ //Sign jwt
+ jwt.Sign(sigAlg, in CacheJwtAlgorithm, 512);
+ //Compile to string
+ jwtMessage = jwt.Compile();
+ }
+ //Set jwt as authorization header
+ client.ClientSocket.Headers[HttpRequestHeader.Authorization] = jwtMessage;
+ //Connect async
+ return client.ConnectAsync(uriBuilder.Uri, token);
+ }
+ /// <summary>
+ /// Registers the current server as active with the specified broker
+ /// </summary>
+ /// <param name="brokerAddress">The address of the broker to register with</param>
+ /// <param name="signingKey">The private key used to sign the message</param>
+ /// <param name="serverAddress">The local address of the current server used for discovery</param>
+ /// <param name="nodeId">The unique id to identify this server (for event queues)</param>
+ /// <param name="keepAliveToken">A unique security token used by the broker to authenticate itself</param>
+ /// <returns>A task that resolves when a successful registration is completed, raises exceptions otherwise</returns>
+ public static async Task ResgisterWithBrokerAsync(Uri brokerAddress, ReadOnlyMemory<byte> signingKey, string serverAddress, string nodeId, string keepAliveToken)
+ {
+ _ = brokerAddress ?? throw new ArgumentNullException(nameof(brokerAddress));
+ _ = serverAddress ?? throw new ArgumentNullException(nameof(serverAddress));
+ _ = keepAliveToken ?? throw new ArgumentNullException(nameof(keepAliveToken));
+ _ = nodeId ?? throw new ArgumentNullException(nameof(nodeId));
+ string requestData;
+ //Create the jwt for signed registration message
+ using (JsonWebToken jwt = new())
+ {
+ //Shared jwt header
+ jwt.WriteHeader(BrokerJwtHeader);
+ //build jwt claim
+ jwt.InitPayloadClaim()
+ .AddClaim("address", serverAddress)
+ .AddClaim("server_id", nodeId)
+ .AddClaim("token", keepAliveToken)
+ .CommitClaims();
+ //Sign the jwt
+ using (ECDsa sigAlg = ECDsa.Create(CacheCurve))
+ {
+ //Import the signing key
+ sigAlg.ImportPkcs8PrivateKey(signingKey.Span, out _);
+ jwt.Sign(sigAlg, in CacheJwtAlgorithm, 512);
+ }
+ //Compile and save
+ requestData = jwt.Compile();
+ }
+ //Create reg request message
+ RestRequest regRequest = new(brokerAddress);
+ regRequest.AddStringBody(requestData, DataFormat.None);
+ regRequest.AddHeader("Content-Type", "text/plain");
+ //Rent client
+ using ClientContract client = ClientPool.Lease();
+ //Exec the regitration request
+ RestResponse response = await client.Resource.ExecutePutAsync(regRequest);
+ if(!response.IsSuccessful)
+ {
+ throw response.ErrorException!;
+ }
+ }
+ private static readonly ConditionalWeakTable<FBMClient, CacheConnectionConfig> ClientCacheConfig = new();
+ /// <summary>
+ /// Imports the client signature algorithim's private key from its pkcs8 binary representation
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="pkcs8PrivateKey">Pkcs8 format private key</param>
+ /// <returns>Chainable fluent object</returns>
+ /// <exception cref="ArgumentException"></exception>
+ /// <exception cref="CryptographicException"></exception>
+ public static FBMClient ImportClientPrivateKey(this FBMClient client, ReadOnlySpan<byte> pkcs8PrivateKey)
+ {
+ CacheConnectionConfig conf = ClientCacheConfig.GetOrCreateValue(client);
+ conf.ClientAlg.ImportPkcs8PrivateKey(pkcs8PrivateKey, out _);
+ return client;
+ }
+ /// <summary>
+ /// Imports the public key used to verify broker server messages
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="spkiPublicKey">The subject-public-key-info formatted broker public key</param>
+ /// <returns>Chainable fluent object</returns>
+ /// <exception cref="ArgumentException"></exception>
+ /// <exception cref="CryptographicException"></exception>
+ public static FBMClient ImportBrokerPublicKey(this FBMClient client, ReadOnlySpan<byte> spkiPublicKey)
+ {
+ CacheConnectionConfig conf = ClientCacheConfig.GetOrCreateValue(client);
+ conf.BrokerAlg.ImportSubjectPublicKeyInfo(spkiPublicKey, out _);
+ return client;
+ }
+ /// <summary>
+ /// Specifies if all connections should be using TLS
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="useTls">A value that indicates if connections should use TLS</param>
+ /// <returns>Chainable fluent object</returns>
+ public static FBMClient UseTls(this FBMClient client, bool useTls)
+ {
+ CacheConnectionConfig conf = ClientCacheConfig.GetOrCreateValue(client);
+ conf.useTls = useTls;
+ return client;
+ }
+ /// <summary>
+ /// Specifies the broker address to discover cache nodes from
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="brokerAddress">The address of the server broker</param>
+ /// <returns>Chainable fluent object</returns>
+ /// <exception cref="ArgumentNullException"></exception>
+ public static FBMClient UseBroker(this FBMClient client, Uri brokerAddress)
+ {
+ CacheConnectionConfig conf = ClientCacheConfig.GetOrCreateValue(client);
+ conf.BrokerAddress = brokerAddress ?? throw new ArgumentNullException(nameof(brokerAddress));
+ return client;
+ }
+ /// <summary>
+ /// Specifies the current server's cluster node id. If this
+ /// is a server connection attempting to listen for changes on the
+ /// remote server, this id must be set and unique
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="nodeId">The cluster node id of the current server</param>
+ /// <returns>Chainable fluent object</returns>
+ /// <exception cref="ArgumentNullException"></exception>
+ public static FBMClient SetNodeId(this FBMClient client, string nodeId)
+ {
+ CacheConnectionConfig conf = ClientCacheConfig.GetOrCreateValue(client);
+ conf.NodeId = nodeId ?? throw new ArgumentNullException(nameof(nodeId));
+ return client;
+ }
+ /// <summary>
+ /// Discovers cache nodes in the broker configured for the current client.
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="token">A token to cancel the discovery</param>
+ /// <returns>A task the resolves the list of active servers on the broker server</returns>
+ public static Task<ActiveServer[]?> DiscoverNodesAsync(this FBMClientWorkerBase client, CancellationToken token = default)
+ {
+ return client.Client.DiscoverNodesAsync(token);
+ }
+ /// <summary>
+ /// Discovers cache nodes in the broker configured for the current client.
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="token">A token to cancel the discovery </param>
+ /// <returns>A task the resolves the list of active servers on the broker server</returns>
+ public static async Task<ActiveServer[]?> DiscoverNodesAsync(this FBMClient client, CancellationToken token = default)
+ {
+ CacheConnectionConfig conf = ClientCacheConfig.GetOrCreateValue(client);
+ //List servers async
+ ActiveServer[]? servers = await ListServersAsync(conf.BrokerAddress!, conf.ClientAlg, conf.BrokerAlg, token);
+ conf.BrokerServers = servers;
+ return servers;
+ }
+ /// <summary>
+ /// Connects the client to a remote cache server
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="server">The server to connect to</param>
+ /// <param name="token">A token to cancel the connection and/or wait operation</param>
+ /// <returns>A task that resolves when cancelled or when the connection is lost to the server</returns>
+ /// <exception cref="OperationCanceledException"></exception>
+ public static Task ConnectAndWaitForExitAsync(this FBMClientWorkerBase client, ActiveServer server, CancellationToken token = default)
+ {
+ return client.Client.ConnectAndWaitForExitAsync(server, token);
+ }
+ /// <summary>
+ /// Connects the client to a remote cache server
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="server">The server to connect to</param>
+ /// <param name="token">A token to cancel the connection and/or wait operation</param>
+ /// <returns>A task that resolves when cancelled or when the connection is lost to the server</returns>
+ /// <exception cref="OperationCanceledException"></exception>
+ public static async Task ConnectAndWaitForExitAsync(this FBMClient client, ActiveServer server, CancellationToken token = default)
+ {
+ CacheConnectionConfig conf = ClientCacheConfig.GetOrCreateValue(client);
+ //Connect to server (no server id because client not replication server)
+ await ConnectAsync(client, server.HostName!, conf.ClientAlg, conf.ServerChallenge, conf.NodeId, conf.useTls, token);
+ //Get task for cancellation
+ Task cancellation = token.WaitHandle.WaitAsync();
+ //Task for status handle
+ Task run = client.ConnectionStatusHandle.WaitAsync();
+ //Wait for cancellation or
+ _ = await Task.WhenAny(cancellation, run);
+ //Normal try to disconnect the socket
+ await client.DisconnectAsync(CancellationToken.None);
+ //Notify if cancelled
+ token.ThrowIfCancellationRequested();
+ }
+ /// <summary>
+ /// Selects a random server from a collection of active servers
+ /// </summary>
+ /// <param name="servers"></param>
+ /// <returns>A server selected at random</returns>
+ public static ActiveServer SelectRandom(this ICollection<ActiveServer> servers)
+ {
+ //select random server
+ int randServer = RandomNumberGenerator.GetInt32(0, servers.Count);
+ return servers.ElementAt(randServer);
+ }
+ }
+<Project Sdk="Microsoft.NET.Sdk">
+ <PropertyGroup>
+ <TargetFramework>net6.0</TargetFramework>
+ <ImplicitUsings>enable</ImplicitUsings>
+ <Nullable>enable</Nullable>
+ <PlatformTarget>x64</PlatformTarget>
+ <GenerateDocumentationFile>True</GenerateDocumentationFile>
+ <Version></Version>
+ <Authors>Vaughn Nugent</Authors>
+ <Copyright>Copyright © 2022 Vaughn Nugent</Copyright>
+ <PackageProjectUrl>www.vaughnnugent.com/resources</PackageProjectUrl>
+ <Platforms>AnyCPU;x64</Platforms>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Debug|AnyCPU'">
+ <CheckForOverflowUnderflow>True</CheckForOverflowUnderflow>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Debug|x64'">
+ <CheckForOverflowUnderflow>True</CheckForOverflowUnderflow>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Release|AnyCPU'">
+ <CheckForOverflowUnderflow>True</CheckForOverflowUnderflow>
+ <Deterministic>False</Deterministic>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Release|x64'">
+ <CheckForOverflowUnderflow>True</CheckForOverflowUnderflow>
+ <Deterministic>False</Deterministic>
+ </PropertyGroup>
+ <ItemGroup>
+ <PackageReference Include="ErrorProne.NET.CoreAnalyzers" Version="0.1.2">
+ <PrivateAssets>all</PrivateAssets>
+ <IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
+ </PackageReference>
+ <PackageReference Include="ErrorProne.NET.Structs" Version="0.1.2">
+ <PrivateAssets>all</PrivateAssets>
+ <IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
+ </PackageReference>
+ <PackageReference Include="RestSharp" Version="108.0.2" />
+ </ItemGroup>
+ <ItemGroup>
+ <ProjectReference Include="..\..\..\VNLib\Hashing\VNLib.Hashing.Portable.csproj" />
+ <ProjectReference Include="..\..\..\VNLib\Http\VNLib.Net.Http.csproj" />
+ <ProjectReference Include="..\..\VNLib.Net.Rest.Client\VNLib.Net.Rest.Client.csproj" />
+ <ProjectReference Include="..\VNLib.Net.Messaging.FBM\src\VNLib.Net.Messaging.FBM.csproj" />
+ </ItemGroup>
+namespace VNLib.Data.Caching.Global.Exceptions
+ public class CacheNotLoadedException : GlobalCacheException
+ {
+ public CacheNotLoadedException()
+ { }
+ public CacheNotLoadedException(string? message) : base(message)
+ { }
+ public CacheNotLoadedException(string? message, Exception? innerException) : base(message, innerException)
+ { }
+ }
+} \ No newline at end of file
+namespace VNLib.Data.Caching.Global.Exceptions
+ public class GlobalCacheException : Exception
+ {
+ public GlobalCacheException()
+ { }
+ public GlobalCacheException(string? message) : base(message)
+ { }
+ public GlobalCacheException(string? message, Exception? innerException) : base(message, innerException)
+ { }
+ }
+} \ No newline at end of file
+using VNLib.Data.Caching.Global.Exceptions;
+namespace VNLib.Data.Caching.Global
+ /// <summary>
+ /// A static library for caching data in-process or a remote data
+ /// cache
+ /// </summary>
+ public static class GlobalDataCache
+ {
+ private static IGlobalCacheProvider? CacheProvider;
+ private static CancellationTokenRegistration _reg;
+ private static readonly object CacheLock = new();
+ private static readonly Dictionary<string, WeakReference<object>> LocalCache = new();
+ /// <summary>
+ /// Gets a value that indicates if global cache is available
+ /// </summary>
+ public static bool IsAvailable => CacheProvider != null && CacheProvider.IsConnected;
+ /// <summary>
+ /// Sets the backing cache provider for the process-wide global cache
+ /// </summary>
+ /// <param name="cacheProvider">The cache provider instance</param>
+ /// <param name="statusToken">A token that represents the store's validity</param>
+ public static void SetProvider(IGlobalCacheProvider cacheProvider, CancellationToken statusToken)
+ {
+ CacheProvider = cacheProvider ?? throw new ArgumentNullException(nameof(cacheProvider));
+ //Remove cache provider when cache provider is no longer valid
+ _reg = statusToken.Register(Cleanup);
+ }
+ private static void Cleanup()
+ {
+ CacheProvider = null;
+ //Clear local cache
+ lock (CacheLock)
+ {
+ LocalCache.Clear();
+ }
+ _reg.Dispose();
+ }
+ /// <summary>
+ /// Asynchronously gets a value from the global cache provider
+ /// </summary>
+ /// <typeparam name="T"></typeparam>
+ /// <param name="key">The key identifying the object to recover from cache</param>
+ /// <returns>The value if found, or null if it does not exist in the store</returns>
+ /// <exception cref="GlobalCacheException"></exception>
+ /// <exception cref="CacheNotLoadedException"></exception>
+ public static async Task<T?> GetAsync<T>(string key) where T: class
+ {
+ //Check local cache first
+ lock (CacheLock)
+ {
+ if (LocalCache.TryGetValue(key, out WeakReference<object>? wr))
+ {
+ //Value is found
+ if(wr.TryGetTarget(out object? value))
+ {
+ //Value exists and is loaded to local cache
+ return (T)value;
+ }
+ //Value has been collected
+ else
+ {
+ //Remove the key from the table
+ LocalCache.Remove(key);
+ }
+ }
+ }
+ //get ref to local cache provider
+ IGlobalCacheProvider? prov = CacheProvider;
+ if(prov == null)
+ {
+ throw new CacheNotLoadedException("Global cache provider was not found");
+ }
+ //get the value from the store
+ T? val = await prov.GetAsync<T>(key);
+ //Only store the value if it was successfully found
+ if (val != null)
+ {
+ //Store in local cache
+ lock (CacheLock)
+ {
+ LocalCache[key] = new WeakReference<object>(val);
+ }
+ }
+ return val;
+ }
+ /// <summary>
+ /// Asynchronously sets (or updates) a cached value in the global cache
+ /// </summary>
+ /// <typeparam name="T"></typeparam>
+ /// <param name="key">The key identifying the object to recover from cache</param>
+ /// <param name="value">The value to set at the given key</param>
+ /// <returns>A task that completes when the update operation has compelted</returns>
+ /// <exception cref="GlobalCacheException"></exception>
+ /// <exception cref="CacheNotLoadedException"></exception>
+ public static async Task SetAsync<T>(string key, T value) where T : class
+ {
+ //Local record is stale, allow it to be loaded from cache next call to get
+ lock (CacheLock)
+ {
+ LocalCache.Remove(key);
+ }
+ //get ref to local cache provider
+ IGlobalCacheProvider? prov = CacheProvider;
+ if (prov == null)
+ {
+ throw new CacheNotLoadedException("Global cache provider was not found");
+ }
+ //set the value in the store
+ await prov.SetAsync<T>(key, value);
+ }
+ /// <summary>
+ /// Asynchronously deletes an item from cache by its key
+ /// </summary>
+ /// <param name="key"></param>
+ /// <returns>A task that completes when the delete operation has compelted</returns>
+ /// <exception cref="GlobalCacheException"></exception>
+ /// <exception cref="CacheNotLoadedException"></exception>
+ public static async Task DeleteAsync(string key)
+ {
+ //Delete from local cache
+ lock (CacheLock)
+ {
+ LocalCache.Remove(key);
+ }
+ //get ref to local cache provider
+ IGlobalCacheProvider? prov = CacheProvider;
+ if (prov == null)
+ {
+ throw new CacheNotLoadedException("Global cache provider was not found");
+ }
+ //Delete value from store
+ await prov.DeleteAsync(key);
+ }
+ }
+} \ No newline at end of file
+namespace VNLib.Data.Caching.Global
+ /// <summary>
+ /// An interface that a cache provoider must impelement to provide data caching to the
+ /// <see cref="GlobalDataCache"/> environment
+ /// </summary>
+ public interface IGlobalCacheProvider
+ {
+ /// <summary>
+ /// Gets a value that indicates if the cache provider is currently available
+ /// </summary>
+ public bool IsConnected { get; }
+ /// <summary>
+ /// Asynchronously gets a value from the backing cache store
+ /// </summary>
+ /// <typeparam name="T"></typeparam>
+ /// <param name="key">The key identifying the object to recover from cache</param>
+ /// <returns>The value if found, or null if it does not exist in the store</returns>
+ Task<T?> GetAsync<T>(string key);
+ /// <summary>
+ /// Asynchronously sets (or updates) a cached value in the backing cache store
+ /// </summary>
+ /// <typeparam name="T"></typeparam>
+ /// <param name="key">The key identifying the object to recover from cache</param>
+ /// <param name="value">The value to set at the given key</param>
+ /// <returns>A task that completes when the update operation has compelted</returns>
+ Task SetAsync<T>(string key, T value);
+ /// <summary>
+ /// Asynchronously deletes an item from cache by its key
+ /// </summary>
+ /// <param name="key">The key identifying the item to delete</param>
+ /// <returns>A task that completes when the delete operation has compelted</returns>
+ Task DeleteAsync(string key);
+ }
+} \ No newline at end of file
+<Project Sdk="Microsoft.NET.Sdk">
+ <PropertyGroup>
+ <TargetFramework>net6.0</TargetFramework>
+ <ImplicitUsings>enable</ImplicitUsings>
+ <Nullable>enable</Nullable>
+ <Authors>Vaughn Nugent</Authors>
+ <Copyright>Copyright © 2022 Vaughn Nugent</Copyright>
+ <Platforms>AnyCPU;ARM32;x64</Platforms>
+ <PlatformTarget>x64</PlatformTarget>
+ <GenerateDocumentationFile>True</GenerateDocumentationFile>
+ </PropertyGroup>
+ <ItemGroup>
+ <PackageReference Include="ErrorProne.NET.CoreAnalyzers" Version="0.1.2">
+ <PrivateAssets>all</PrivateAssets>
+ <IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
+ </PackageReference>
+ <PackageReference Include="ErrorProne.NET.Structs" Version="0.1.2">
+ <PrivateAssets>all</PrivateAssets>
+ <IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
+ </PackageReference>
+ </ItemGroup>
+namespace VNLib.Data.Caching.ObjectCache
+ /// <summary>
+ /// An event object that is passed when change events occur
+ /// </summary>
+ public class ChangeEvent
+ {
+ public readonly string CurrentId;
+ public readonly string? AlternateId;
+ public readonly bool Deleted;
+ internal ChangeEvent(string id, string? alternate, bool deleted)
+ {
+ CurrentId = id;
+ AlternateId = alternate;
+ Deleted = deleted;
+ }
+ }
+using System;
+using System.IO;
+using System.Threading;
+using System.Threading.Tasks;
+using VNLib.Utils.IO;
+using VNLib.Utils.Async;
+using VNLib.Utils.Memory;
+using VNLib.Utils.Logging;
+using VNLib.Utils.Extensions;
+using VNLib.Net.Messaging.FBM.Server;
+using static VNLib.Data.Caching.Constants;
+namespace VNLib.Data.Caching.ObjectCache
+ public delegate ReadOnlySpan<byte> GetBodyDataCallback<T>(T state);
+ /// <summary>
+ /// A <see cref="FBMListener"/> implementation of a <see cref="CacheListener"/>
+ /// </summary>
+ public class ObjectCacheStore : CacheListener, IDisposable
+ {
+ private readonly SemaphoreSlim StoreLock;
+ private bool disposedValue;
+ ///<inheritdoc/>
+ protected override ILogProvider Log { get; }
+ /// <summary>
+ /// A queue that stores update and delete events
+ /// </summary>
+ public AsyncQueue<ChangeEvent> EventQueue { get; }
+ /// <summary>
+ /// Initialzies a new <see cref="ObjectCacheStore"/>
+ /// </summary>
+ /// <param name="dir">The <see cref="DirectoryInfo"/> to store blob files to</param>
+ /// <param name="cacheMax"></param>
+ /// <param name="log"></param>
+ /// <param name="heap"></param>
+ /// <param name="singleReader">A value that indicates if a single thread is processing events</param>
+ public ObjectCacheStore(DirectoryInfo dir, int cacheMax, ILogProvider log, IUnmangedHeap heap, bool singleReader)
+ {
+ Log = log;
+ //We can use a single writer and single reader in this context
+ EventQueue = new(true, singleReader);
+ InitCache(dir, cacheMax, heap);
+ InitListener(heap);
+ StoreLock = new(1,1);
+ }
+ ///<inheritdoc/>
+ protected override async Task ProcessAsync(FBMContext context, object? userState, CancellationToken cancellationToken)
+ {
+ try
+ {
+ //Get the action header
+ string action = context.Method();
+ //Optional newid header
+ string? alternateId = context.NewObjectId();
+ switch (action)
+ {
+ case Actions.Get:
+ {
+ //Get the object-id header
+ string objectId = context.ObjectId();
+ //Take lock on store
+ using SemSlimReleaser rel = await StoreLock.GetReleaserAsync(cancellationToken: cancellationToken);
+ if (Cache!.TryGetValue(objectId, out MemoryHandle<byte>? data))
+ {
+ //Set the status code and write the buffered data to the response buffer
+ context.CloseResponse(ResponseCodes.Okay);
+ //Copy data to response buffer
+ context.Response.WriteBody(data.Span);
+ }
+ else
+ {
+ context.CloseResponse(ResponseCodes.NotFound);
+ }
+ }
+ break;
+ case Actions.AddOrUpdate:
+ {
+ //Get the object-id header
+ string objectId = context.ObjectId();
+ //Add/update a blob async
+ await AddOrUpdateBlobAsync(objectId, alternateId, static context => context.Request.BodyData, context);
+ //Notify update the event bus
+ await EventQueue.EnqueueAsync(new(objectId, alternateId, false), cancellationToken);
+ //Set status code
+ context.CloseResponse(ResponseCodes.Okay);
+ }
+ break;
+ case Actions.Delete:
+ {
+ //Get the object-id header
+ string objectId = context.ObjectId();
+ if (await DeleteItemAsync(objectId))
+ {
+ //Notify deleted
+ await EventQueue.EnqueueAsync(new(objectId, null, true), cancellationToken);
+ //Set status header
+ context.CloseResponse(ResponseCodes.Okay);
+ }
+ else
+ {
+ //Set status header
+ context.CloseResponse(ResponseCodes.NotFound);
+ }
+ }
+ break;
+ // event queue dequeue request
+ case Actions.Dequeue:
+ {
+ //If no event bus is registered, then this is not a legal command
+ if (userState is not AsyncQueue<ChangeEvent> eventBus)
+ {
+ context.CloseResponse(ResponseCodes.NotFound);
+ break;
+ }
+ //Wait for a new message to process
+ ChangeEvent ev = await eventBus.DequeueAsync(cancellationToken);
+ if (ev.Deleted)
+ {
+ context.CloseResponse("deleted");
+ context.Response.WriteHeader(ObjectId, ev.CurrentId);
+ }
+ else
+ {
+ //Changed
+ context.CloseResponse("modified");
+ context.Response.WriteHeader(ObjectId, ev.CurrentId);
+ //Set old id if an old id is set
+ if (ev.CurrentId != null)
+ {
+ context.Response.WriteHeader(NewObjectId, ev.AlternateId);
+ }
+ }
+ }
+ break;
+ }
+ }
+ catch (OperationCanceledException)
+ {
+ throw;
+ }
+ catch(Exception ex)
+ {
+ //Log error and set error status code
+ Log.Error(ex);
+ context.CloseResponse(ResponseCodes.Error);
+ }
+ }
+ /// <summary>
+ /// Asynchronously deletes a previously stored item
+ /// </summary>
+ /// <param name="id">The id of the object to delete</param>
+ /// <returns>A task that completes when the item has been deleted</returns>
+ public async Task<bool> DeleteItemAsync(string id)
+ {
+ using SemSlimReleaser rel = await StoreLock.GetReleaserAsync();
+ return Cache!.Remove(id);
+ }
+ /// <summary>
+ /// Asynchronously adds or updates an object in the store and optionally update's its id
+ /// </summary>
+ /// <param name="objectId">The current (or old) id of the object</param>
+ /// <param name="alternateId">An optional id to update the blob to</param>
+ /// <param name="bodyData">A callback that returns the data for the blob</param>
+ /// <param name="state">The state parameter to pass to the data callback</param>
+ /// <returns></returns>
+ public async Task AddOrUpdateBlobAsync<T>(string objectId, string? alternateId, GetBodyDataCallback<T> bodyData, T state)
+ {
+ MemoryHandle<byte>? blob;
+ //See if new/alt session id was specified
+ if (string.IsNullOrWhiteSpace(alternateId))
+ {
+ //Take lock on store
+ using SemSlimReleaser rel = await StoreLock.GetReleaserAsync();
+ //See if blob exists
+ if (!Cache!.TryGetValue(objectId, out blob))
+ {
+ //If not, create new blob and add to store
+ blob = Heap.AllocAndCopy(bodyData(state));
+ Cache.Add(objectId, blob);
+ }
+ else
+ {
+ //Reset the buffer state
+ blob.WriteAndResize(bodyData(state));
+ }
+ }
+ //Need to change the id of the record
+ else
+ {
+ //Take lock on store
+ using SemSlimReleaser rel = await StoreLock.GetReleaserAsync();
+ //Try to change the blob key
+ if (!Cache!.TryChangeKey(objectId, alternateId, out blob))
+ {
+ //Blob not found, create new blob
+ blob = Heap.AllocAndCopy(bodyData(state));
+ Cache.Add(alternateId, blob);
+ }
+ else
+ {
+ //Reset the buffer state
+ blob.WriteAndResize(bodyData(state));
+ }
+ }
+ }
+ ///<inheritdoc/>
+ protected virtual void Dispose(bool disposing)
+ {
+ if (!disposedValue)
+ {
+ if (disposing)
+ {
+ Cache?.Clear();
+ }
+ disposedValue = true;
+ }
+ }
+ ///<inheritdoc/>
+ public void Dispose()
+ {
+ // Do not change this code. Put cleanup code in 'Dispose(bool disposing)' method
+ Dispose(disposing: true);
+ GC.SuppressFinalize(this);
+ }
+ }
+<Project Sdk="Microsoft.NET.Sdk">
+ <PropertyGroup>
+ <TargetFramework>net6.0</TargetFramework>
+ <Platforms>AnyCPU;x64</Platforms>
+ <Authors>Vaughn Nugent</Authors>
+ <Copyright>Copyright © 2022 Vaughn Nugent</Copyright>
+ <Nullable>enable</Nullable>
+ <GenerateDocumentationFile>True</GenerateDocumentationFile>
+ <PackageProjectUrl>www.vaughnnugent.com/resources</PackageProjectUrl>
+ <AssemblyVersion></AssemblyVersion>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Debug|x64'">
+ <DocumentationFile></DocumentationFile>
+ <CheckForOverflowUnderflow>False</CheckForOverflowUnderflow>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Debug|AnyCPU'">
+ <CheckForOverflowUnderflow>False</CheckForOverflowUnderflow>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Release|AnyCPU'">
+ <CheckForOverflowUnderflow>False</CheckForOverflowUnderflow>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Release|x64'">
+ <CheckForOverflowUnderflow>False</CheckForOverflowUnderflow>
+ </PropertyGroup>
+ <ItemGroup>
+ <PackageReference Include="ErrorProne.NET.CoreAnalyzers" Version="0.1.2">
+ <PrivateAssets>all</PrivateAssets>
+ <IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
+ </PackageReference>
+ <PackageReference Include="ErrorProne.NET.Structs" Version="0.1.2">
+ <PrivateAssets>all</PrivateAssets>
+ <IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
+ </PackageReference>
+ </ItemGroup>
+ <ItemGroup>
+ <ProjectReference Include="..\..\..\VNLib\Utils\src\VNLib.Utils.csproj" />
+ <ProjectReference Include="..\VNLib.Data.Caching\src\VNLib.Data.Caching.csproj" />
+ </ItemGroup>
diff --git a/VNLib.Data.Caching/src/BlobCache.cs b/VNLib.Data.Caching/src/BlobCache.cs
new file mode 100644
index 0000000..89818be
--- /dev/null
+++ b/VNLib.Data.Caching/src/BlobCache.cs
@@ -0,0 +1,118 @@
+using System;
+using System.IO;
+using System.Linq;
+using System.Collections.Generic;
+using System.Diagnostics.CodeAnalysis;
+using VNLib.Utils.IO;
+using VNLib.Utils.Logging;
+using VNLib.Utils.Memory;
+using VNLib.Utils.Memory.Caching;
+#nullable enable
+namespace VNLib.Data.Caching
+ /// <summary>
+ /// A general purpose binary data storage
+ /// </summary>
+ public class BlobCache : LRUCache<string, MemoryHandle<byte>>
+ {
+ readonly IUnmangedHeap Heap;
+ readonly DirectoryInfo SwapDir;
+ readonly ILogProvider Log;
+ ///<inheritdoc/>
+ public override bool IsReadOnly { get; }
+ ///<inheritdoc/>
+ protected override int MaxCapacity { get; }
+ /// <summary>
+ /// Initializes a new <see cref="BlobCache"/> store
+ /// </summary>
+ /// <param name="swapDir">The <see cref="IsolatedStorageDirectory"/> to swap blob data to when cache</param>
+ /// <param name="maxCapacity">The maximum number of items to keep in memory</param>
+ /// <param name="log">A <see cref="ILogProvider"/> to write log data to</param>
+ /// <param name="heap">A <see cref="IUnmangedHeap"/> to allocate buffers and store <see cref="BlobItem"/> data in memory</param>
+ public BlobCache(DirectoryInfo swapDir, int maxCapacity, ILogProvider log, IUnmangedHeap heap)
+ :base(StringComparer.Ordinal)
+ {
+ IsReadOnly = false;
+ MaxCapacity = maxCapacity;
+ SwapDir = swapDir;
+ //Update the lookup table size
+ LookupTable.EnsureCapacity(maxCapacity);
+ //Set default heap if not specified
+ Heap = heap;
+ Log = log;
+ }
+ ///<inheritdoc/>
+ protected override bool CacheMiss(string key, [NotNullWhen(true)] out MemoryHandle<byte>? value)
+ {
+ value = null;
+ return false;
+ }
+ ///<inheritdoc/>
+ protected override void Evicted(KeyValuePair<string, MemoryHandle<byte>> evicted)
+ {
+ //Dispose the blob
+ evicted.Value.Dispose();
+ }
+ /// <summary>
+ /// If the <see cref="BlobItem"/> is found in the store, changes the key
+ /// that referrences the blob.
+ /// </summary>
+ /// <param name="currentKey">The key that currently referrences the blob in the store</param>
+ /// <param name="newKey">The new key that will referrence the blob</param>
+ /// <param name="blob">The <see cref="BlobItem"/> if its found in the store</param>
+ /// <returns>True if the record was found and the key was changes</returns>
+ public bool TryChangeKey(string currentKey, string newKey, [NotNullWhen(true)] out MemoryHandle<byte>? blob)
+ {
+ if (LookupTable.Remove(currentKey, out LinkedListNode<KeyValuePair<string, MemoryHandle<byte>>>? node))
+ {
+ //Remove the node from the ll
+ List.Remove(node);
+ //Update the node kvp
+ blob = node.Value.Value;
+ node.Value = new KeyValuePair<string, MemoryHandle<byte>>(newKey, blob);
+ //Add to end of list
+ List.AddLast(node);
+ //Re-add to lookup table with new key
+ LookupTable.Add(newKey, node);
+ return true;
+ }
+ blob = null;
+ return false;
+ }
+ /// <summary>
+ /// Removes the <see cref="BlobItem"/> from the store without disposing the blobl
+ /// </summary>
+ /// <param name="key">The key that referrences the <see cref="BlobItem"/> in the store</param>
+ /// <returns>A value indicating if the blob was removed</returns>
+ public override bool Remove(string key)
+ {
+ //Remove the item from the lookup table and if it exists, remove the node from the list
+ if (LookupTable.Remove(key, out LinkedListNode<KeyValuePair<string, MemoryHandle<byte>>>? node))
+ {
+ //Remove the new from the list
+ List.Remove(node);
+ //dispose the buffer
+ node.Value.Value.Dispose();
+ return true;
+ }
+ return false;
+ }
+ /// <summary>
+ /// Removes and disposes all blobl elements in cache (or in the backing store)
+ /// </summary>
+ public override void Clear()
+ {
+ foreach (MemoryHandle<byte> blob in List.Select(kp => kp.Value))
+ {
+ blob.Dispose();
+ }
+ base.Clear();
+ }
+ }
diff --git a/VNLib.Data.Caching/src/BlobItem.cs b/VNLib.Data.Caching/src/BlobItem.cs
new file mode 100644
index 0000000..a5630e9
--- /dev/null
+++ b/VNLib.Data.Caching/src/BlobItem.cs
@@ -0,0 +1,185 @@
+using System;
+using System.IO;
+using System.Threading;
+using System.Threading.Tasks;
+using VNLib.Utils;
+using VNLib.Utils.IO;
+using VNLib.Utils.Memory;
+using VNLib.Utils.Logging;
+using VNLib.Utils.Extensions;
+#nullable enable
+namespace VNLib.Data.Caching
+ /// <summary>
+ /// A general purpose binary storage item
+ /// </summary>
+ public class BlobItem //: VnDisposeable
+ {
+ /*
+ private static readonly JoinableTaskContext JTX = new();
+ private static readonly Semaphore CentralSwapLock = new(Environment.ProcessorCount, Environment.ProcessorCount);
+ private readonly VnMemoryStream _loadedData;
+ private bool _loaded;
+ /// <summary>
+ /// The time the blob was last modified
+ /// </summary>
+ public DateTimeOffset LastAccessed { get; private set; }
+ /// <summary>
+ /// Gets the current size of the file (in bytes) as an atomic operation
+ /// </summary>
+ public int FileSize => (int)_loadedData.Length;
+ /// <summary>
+ /// The operation synchronization lock
+ /// </summary>
+ public AsyncReaderWriterLock OpLock { get; }
+ /// <summary>
+ /// Initializes a new <see cref="BlobItem"/>
+ /// </summary>
+ /// <param name="heap">The heap to allocate buffers from</param>
+ internal BlobItem(IUnmangedHeap heap)
+ {
+ _loadedData = new(heap);
+ OpLock = new AsyncReaderWriterLock(JTX);
+ _loaded = true;
+ LastAccessed = DateTimeOffset.UtcNow;
+ }
+ ///<inheritdoc/>
+ protected override void Free()
+ {
+ _loadedData.Dispose();
+ OpLock.Dispose();
+ }
+ /// <summary>
+ /// Reads data from the internal buffer and copies it to the specified buffer.
+ /// Use the <see cref="FileSize"/> property to obtain the size of the internal buffer
+ /// </summary>
+ /// <param name="buffer">The buffer to copy data to</param>
+ /// <returns>When completed, the number of bytes copied to the buffer</returns>
+ public int Read(Span<byte> buffer)
+ {
+ //Make sure the blob has been swapped back into memory
+ if (!_loaded)
+ {
+ throw new InvalidOperationException("The blob was not loaded from the disk");
+ }
+ //Read all data from the buffer and write it to the output buffer
+ _loadedData.AsSpan().CopyTo(buffer);
+ //Update last-accessed
+ LastAccessed = DateTimeOffset.UtcNow;
+ return (int)_loadedData.Length;
+ }
+ /// <summary>
+ /// Overwrites the internal buffer with the contents of the supplied buffer
+ /// </summary>
+ /// <param name="buffer">The buffer containing data to store within the blob</param>
+ /// <returns>A <see cref="ValueTask"/> that completes when write access has been granted and copied</returns>
+ /// <exception cref="InvalidOperationException"></exception>
+ public void Write(ReadOnlySpan<byte> buffer)
+ {
+ //Make sure the blob has been swapped back into memory
+ if (!_loaded)
+ {
+ throw new InvalidOperationException("The blob was not loaded from the disk");
+ }
+ //Reset the buffer
+ _loadedData.SetLength(buffer.Length);
+ _loadedData.Seek(0, SeekOrigin.Begin);
+ _loadedData.Write(buffer);
+ LastAccessed = DateTimeOffset.UtcNow;
+ }
+ /// <summary>
+ /// Writes the contents of the memory buffer to its designated file on the disk
+ /// </summary>
+ /// <param name="heap">The heap to allocate buffers from</param>
+ /// <param name="swapDir">The <see cref="IsolatedStorageDirectory"/> that stores the file</param>
+ /// <param name="filename">The name of the file to write data do</param>
+ /// <param name="log">A log to write errors to</param>
+ /// <returns>A task that completes when the swap to disk is complete</returns>
+ internal async Task SwapToDiskAsync(IUnmangedHeap heap, DirectoryInfo swapDir, string filename, ILogProvider log)
+ {
+ try
+ {
+ //Wait for write lock
+ await using (AsyncReaderWriterLock.Releaser releaser = await OpLock.WriteLockAsync())
+ {
+ //Enter swap lock
+ await CentralSwapLock;
+ try
+ {
+ //Open swap file data stream
+ await using FileStream swapFile = swapDir.OpenFile(filename, FileMode.OpenOrCreate, FileAccess.ReadWrite, bufferSize: 8128);
+ //reset swap file
+ swapFile.SetLength(0);
+ //Seek loaded-data back to 0 before writing
+ _loadedData.Seek(0, SeekOrigin.Begin);
+ //Write loaded data to disk
+ await _loadedData.CopyToAsync(swapFile, 8128, heap);
+ }
+ finally
+ {
+ CentralSwapLock.Release();
+ }
+ //Release memory held by stream
+ _loadedData.SetLength(0);
+ //Clear loaded flag
+ _loaded = false;
+ LastAccessed = DateTimeOffset.UtcNow;
+ }
+ log.Debug("Blob {name} swapped to disk", filename);
+ }
+ catch(Exception ex)
+ {
+ log.Error(ex, "Blob swap to disk error");
+ }
+ }
+ /// <summary>
+ /// Reads the contents of the blob into a memory buffer from its designated file on disk
+ /// </summary>
+ /// <param name="heap">The heap to allocate buffers from</param>
+ /// <param name="swapDir">The <see cref="IsolatedStorageDirectory"/> that stores the file</param>
+ /// <param name="filename">The name of the file to write the blob data to</param>
+ /// <param name="log">A log to write errors to</param>
+ /// <returns>A task that completes when the swap from disk is complete</returns>
+ internal async Task SwapFromDiskAsync(IUnmangedHeap heap, DirectoryInfo swapDir, string filename, ILogProvider log)
+ {
+ try
+ {
+ //Wait for write lock
+ await using (AsyncReaderWriterLock.Releaser releaser = await OpLock.WriteLockAsync())
+ {
+ //Enter swap lock
+ await CentralSwapLock;
+ try
+ {
+ //Open swap file data stream
+ await using FileStream swapFile = swapDir.OpenFile(filename, FileMode.OpenOrCreate, FileAccess.Read, bufferSize:8128);
+ //Copy from disk to memory
+ await swapFile.CopyToAsync(_loadedData, 8128, heap);
+ }
+ finally
+ {
+ CentralSwapLock.Release();
+ }
+ //Set loaded flag
+ _loaded = true;
+ LastAccessed = DateTimeOffset.UtcNow;
+ }
+ log.Debug("Blob {name} swapped from disk", filename);
+ }
+ catch(Exception ex)
+ {
+ log.Error(ex, "Blob swap from disk error");
+ }
+ }
+ */
+ }
diff --git a/VNLib.Data.Caching/src/CacheListener.cs b/VNLib.Data.Caching/src/CacheListener.cs
new file mode 100644
index 0000000..18e9684
--- /dev/null
+++ b/VNLib.Data.Caching/src/CacheListener.cs
@@ -0,0 +1,40 @@
+using System;
+using System.IO;
+using VNLib.Utils.Memory;
+using VNLib.Net.Messaging.FBM.Server;
+namespace VNLib.Data.Caching
+ /// <summary>
+ /// A base implementation of a memory/disk LRU data cache FBM listener
+ /// </summary>
+ public abstract class CacheListener : FBMListenerBase
+ {
+ /// <summary>
+ /// The directory swap files will be stored
+ /// </summary>
+ public DirectoryInfo? Directory { get; private set; }
+ /// <summary>
+ /// The Cache store to access data blobs
+ /// </summary>
+ protected BlobCache? Cache { get; private set; }
+ /// <summary>
+ /// The <see cref="IUnmangedHeap"/> to allocate buffers from
+ /// </summary>
+ protected IUnmangedHeap? Heap { get; private set; }
+ /// <summary>
+ /// Initializes the <see cref="Cache"/> data store
+ /// </summary>
+ /// <param name="dir">The directory to swap cache records to</param>
+ /// <param name="cacheSize">The size of the LRU cache</param>
+ /// <param name="heap">The heap to allocate buffers from</param>
+ protected void InitCache(DirectoryInfo dir, int cacheSize, IUnmangedHeap heap)
+ {
+ Heap = heap;
+ Cache = new(dir, cacheSize, Log, Heap);
+ Directory = dir;
+ }
+ }
diff --git a/VNLib.Data.Caching/src/ClientExtensions.cs b/VNLib.Data.Caching/src/ClientExtensions.cs
new file mode 100644
index 0000000..18a1aa9
--- /dev/null
+++ b/VNLib.Data.Caching/src/ClientExtensions.cs
@@ -0,0 +1,310 @@
+using System;
+using System.IO;
+using System.Linq;
+using System.Buffers;
+using System.Text.Json;
+using System.Threading;
+using System.Threading.Tasks;
+using System.Collections.Generic;
+using System.Text.Json.Serialization;
+using System.Runtime.CompilerServices;
+using VNLib.Utils.Logging;
+using VNLib.Net.Messaging.FBM;
+using VNLib.Net.Messaging.FBM.Client;
+using VNLib.Net.Messaging.FBM.Server;
+using VNLib.Data.Caching.Exceptions;
+using static VNLib.Data.Caching.Constants;
+namespace VNLib.Data.Caching
+ /// <summary>
+ /// Provides caching extension methods for <see cref="FBMClient"/>
+ /// </summary>
+ public static class ClientExtensions
+ {
+ private static readonly JsonSerializerOptions LocalOptions = new()
+ {
+ DictionaryKeyPolicy = JsonNamingPolicy.CamelCase,
+ NumberHandling = JsonNumberHandling.Strict,
+ ReadCommentHandling = JsonCommentHandling.Disallow,
+ WriteIndented = false,
+ DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingNull,
+ IgnoreReadOnlyFields = true,
+ PropertyNameCaseInsensitive = true,
+ IncludeFields = false,
+ //Use small buffers
+ DefaultBufferSize = 128
+ };
+ private static readonly ConditionalWeakTable<FBMClient, SemaphoreSlim> GetLock = new();
+ private static readonly ConditionalWeakTable<FBMClient, SemaphoreSlim> UpdateLock = new();
+ private static SemaphoreSlim GetLockCtor(FBMClient client) => new (50);
+ private static SemaphoreSlim UpdateLockCtor(FBMClient client) => new (25);
+ /// <summary>
+ /// Gets an object from the server if it exists
+ /// </summary>
+ /// <typeparam name="T"></typeparam>
+ /// <param name="client"></param>
+ /// <param name="objectId">The id of the object to get</param>
+ /// <param name="cancellationToken">A token to cancel the operation</param>
+ /// <returns>A task that completes to return the results of the response payload</returns>
+ /// <exception cref="JsonException"></exception>
+ /// <exception cref="OutOfMemoryException"></exception>
+ /// <exception cref="InvalidStatusException"></exception>
+ /// <exception cref="ObjectDisposedException"></exception>
+ /// <exception cref="InvalidResponseException"></exception>
+ public static async Task<T?> GetObjectAsync<T>(this FBMClient client, string objectId, CancellationToken cancellationToken = default)
+ {
+ client.Config.DebugLog?.Debug("[DEBUG] Getting object {id}", objectId);
+ SemaphoreSlim getLock = GetLock.GetValue(client, GetLockCtor);
+ //Wait for entry
+ await getLock.WaitAsync(cancellationToken);
+ //Rent a new request
+ FBMRequest request = client.RentRequest();
+ try
+ {
+ //Set action as get/create
+ request.WriteHeader(HeaderCommand.Action, Actions.Get);
+ //Set session-id header
+ request.WriteHeader(Constants.ObjectId, objectId);
+ //Make request
+ using FBMResponse response = await client.SendAsync(request, cancellationToken);
+ response.ThrowIfNotSet();
+ //Get the status code
+ ReadOnlyMemory<char> status = response.Headers.FirstOrDefault(static a => a.Key == HeaderCommand.Status).Value;
+ if (status.Span.Equals(ResponseCodes.Okay, StringComparison.Ordinal))
+ {
+ return JsonSerializer.Deserialize<T>(response.ResponseBody, LocalOptions);
+ }
+ //Session may not exist on the server yet
+ if (status.Span.Equals(ResponseCodes.NotFound, StringComparison.Ordinal))
+ {
+ return default;
+ }
+ throw new InvalidStatusException("Invalid status code recived for object get request", status.ToString());
+ }
+ finally
+ {
+ getLock.Release();
+ client.ReturnRequest(request);
+ }
+ }
+ /// <summary>
+ /// Updates the state of the object, and optionally updates the ID of the object. The data
+ /// parameter is serialized, buffered, and streamed to the remote server
+ /// </summary>
+ /// <typeparam name="T"></typeparam>
+ /// <param name="client"></param>
+ /// <param name="objectId">The id of the object to update or replace</param>
+ /// <param name="newId">An optional parameter to specify a new ID for the old object</param>
+ /// <param name="data">The payload data to serialize and set as the data state of the session</param>
+ /// <param name="cancellationToken">A token to cancel the operation</param>
+ /// <returns>A task that resolves when the server responds</returns>
+ /// <exception cref="JsonException"></exception>
+ /// <exception cref="OutOfMemoryException"></exception>
+ /// <exception cref="InvalidStatusException"></exception>
+ /// <exception cref="ObjectDisposedException"></exception>
+ /// <exception cref="InvalidResponseException"></exception>
+ /// <exception cref="MessageTooLargeException"></exception>
+ /// <exception cref="ObjectNotFoundException"></exception>
+ public static async Task AddOrUpdateObjectAsync<T>(this FBMClient client, string objectId, string? newId, T data, CancellationToken cancellationToken = default)
+ {
+ client.Config.DebugLog?.Debug("[DEBUG] Updating object {id}, newid {nid}", objectId, newId);
+ SemaphoreSlim updateLock = UpdateLock.GetValue(client, UpdateLockCtor);
+ //Wait for entry
+ await updateLock.WaitAsync(cancellationToken);
+ //Rent a new request
+ FBMRequest request = client.RentRequest();
+ try
+ {
+ //Set action as get/create
+ request.WriteHeader(HeaderCommand.Action, Actions.AddOrUpdate);
+ //Set session-id header
+ request.WriteHeader(Constants.ObjectId, objectId);
+ //if new-id set, set the new-id header
+ if (!string.IsNullOrWhiteSpace(newId))
+ {
+ request.WriteHeader(Constants.NewObjectId, newId);
+ }
+ //Get the body writer for the message
+ IBufferWriter<byte> bodyWriter = request.GetBodyWriter();
+ //Write json data to the message
+ using (Utf8JsonWriter jsonWriter = new(bodyWriter))
+ {
+ JsonSerializer.Serialize(jsonWriter, data, LocalOptions);
+ }
+ //Make request
+ using FBMResponse response = await client.SendAsync(request, cancellationToken);
+ response.ThrowIfNotSet();
+ //Get the status code
+ ReadOnlyMemory<char> status = response.Headers.FirstOrDefault(static a => a.Key == HeaderCommand.Status).Value;
+ //Check status code
+ if (status.Span.Equals(ResponseCodes.Okay, StringComparison.OrdinalIgnoreCase))
+ {
+ return;
+ }
+ else if(status.Span.Equals(ResponseCodes.NotFound, StringComparison.OrdinalIgnoreCase))
+ {
+ throw new ObjectNotFoundException($"object {objectId} not found on remote server");
+ }
+ //Invalid status
+ throw new InvalidStatusException("Invalid status code recived for object upsert request", status.ToString());
+ }
+ finally
+ {
+ updateLock.Release();
+ //Return the request(clears data and reset)
+ client.ReturnRequest(request);
+ }
+ }
+ /// <summary>
+ /// Asynchronously deletes an object in the remote store
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="objectId">The id of the object to update or replace</param>
+ /// <param name="cancellationToken">A token to cancel the operation</param>
+ /// <returns>A task that resolves when the operation has completed</returns>
+ /// <exception cref="InvalidStatusException"></exception>
+ /// <exception cref="ObjectDisposedException"></exception>
+ /// <exception cref="InvalidResponseException"></exception>
+ /// <exception cref="ObjectNotFoundException"></exception>
+ public static async Task DeleteObjectAsync(this FBMClient client, string objectId, CancellationToken cancellationToken = default)
+ {
+ client.Config.DebugLog?.Debug("[DEBUG] Deleting object {id}", objectId);
+ SemaphoreSlim updateLock = UpdateLock.GetValue(client, UpdateLockCtor);
+ //Wait for entry
+ await updateLock.WaitAsync(cancellationToken);
+ //Rent a new request
+ FBMRequest request = client.RentRequest();
+ try
+ {
+ //Set action as delete
+ request.WriteHeader(HeaderCommand.Action, Actions.Delete);
+ //Set session-id header
+ request.WriteHeader(Constants.ObjectId, objectId);
+ //Make request
+ using FBMResponse response = await client.SendAsync(request, cancellationToken);
+ response.ThrowIfNotSet();
+ //Get the status code
+ ReadOnlyMemory<char> status = response.Headers.FirstOrDefault(static a => a.Key == HeaderCommand.Status).Value;
+ if (status.Span.Equals(ResponseCodes.Okay, StringComparison.Ordinal))
+ {
+ return;
+ }
+ else if(status.Span.Equals(ResponseCodes.NotFound, StringComparison.OrdinalIgnoreCase))
+ {
+ throw new ObjectNotFoundException($"object {objectId} not found on remote server");
+ }
+ throw new InvalidStatusException("Invalid status code recived for object get request", status.ToString());
+ }
+ finally
+ {
+ updateLock.Release();
+ client.ReturnRequest(request);
+ }
+ }
+ /// <summary>
+ /// Dequeues a change event from the server event queue for the current connection, or waits until a change happens
+ /// </summary>
+ /// <param name="client"></param>
+ /// <param name="cancellationToken">A token to cancel the deuque operation</param>
+ /// <returns>A <see cref="KeyValuePair{TKey, TValue}"/> that contains the modified object id and optionally its new id</returns>
+ public static async Task<WaitForChangeResult> WaitForChangeAsync(this FBMClient client, CancellationToken cancellationToken = default)
+ {
+ //Rent a new request
+ FBMRequest request = client.RentRequest();
+ try
+ {
+ //Set action as event dequeue to dequeue a change event
+ request.WriteHeader(HeaderCommand.Action, Actions.Dequeue);
+ //Make request
+ using FBMResponse response = await client.SendAsync(request, cancellationToken);
+ response.ThrowIfNotSet();
+ return new()
+ {
+ Status = response.Headers.FirstOrDefault(static a => a.Key == HeaderCommand.Status).Value.ToString(),
+ CurrentId = response.Headers.SingleOrDefault(static v => v.Key == Constants.ObjectId).Value.ToString(),
+ NewId = response.Headers.SingleOrDefault(static v => v.Key == Constants.NewObjectId).Value.ToString()
+ };
+ }
+ finally
+ {
+ client.ReturnRequest(request);
+ }
+ }
+ /// <summary>
+ /// Gets the Object-id for the request message, or throws an <see cref="InvalidOperationException"/> if not specified
+ /// </summary>
+ /// <param name="context"></param>
+ /// <returns>The id of the object requested</returns>
+ /// <exception cref="InvalidOperationException"></exception>
+ public static string ObjectId(this FBMContext context)
+ {
+ return context.Request.Headers.First(static kvp => kvp.Key == Constants.ObjectId).Value.ToString();
+ }
+ /// <summary>
+ /// Gets the new ID of the object if specified from the request. Null if the request did not specify an id update
+ /// </summary>
+ /// <param name="context"></param>
+ /// <returns>The new ID of the object if speicifed, null otherwise</returns>
+ public static string? NewObjectId(this FBMContext context)
+ {
+ return context.Request.Headers.FirstOrDefault(static kvp => kvp.Key == Constants.NewObjectId).Value.ToString();
+ }
+ /// <summary>
+ /// Gets the request method for the request
+ /// </summary>
+ /// <param name="context"></param>
+ /// <returns>The request method string</returns>
+ public static string Method(this FBMContext context)
+ {
+ return context.Request.Headers.First(static kvp => kvp.Key == HeaderCommand.Action).Value.ToString();
+ }
+ /// <summary>
+ /// Closes a response with a status code
+ /// </summary>
+ /// <param name="context"></param>
+ /// <param name="responseCode">The status code to send to the client</param>
+ public static void CloseResponse(this FBMContext context, string responseCode)
+ {
+ context.Response.WriteHeader(HeaderCommand.Status, responseCode);
+ }
+ /// <summary>
+ /// Initializes the worker for a reconnect policy and returns an object that can listen for changes
+ /// and configure the connection as necessary
+ /// </summary>
+ /// <param name="worker"></param>
+ /// <param name="retryDelay">The amount of time to wait between retries</param>
+ /// <param name="serverUri">The uri to reconnect the client to</param>
+ /// <returns>A <see cref="ClientRetryManager{T}"/> for listening for retry events</returns>
+ public static ClientRetryManager<T> SetReconnectPolicy<T>(this T worker, TimeSpan retryDelay, Uri serverUri) where T: IStatefulConnection
+ {
+ //Return new manager
+ return new (worker, retryDelay, serverUri);
+ }
+ }
+using System;
+using System.Threading.Tasks;
+using System.Security.Cryptography;
+using VNLib.Utils;
+using VNLib.Net.Messaging.FBM.Client;
+namespace VNLib.Data.Caching
+ /// <summary>
+ /// Manages a <see cref="FBMClientWorkerBase"/> reconnect policy
+ /// </summary>
+ public class ClientRetryManager<T> : VnDisposeable where T: IStatefulConnection
+ {
+ const int RetryRandMaxMsDelay = 1000;
+ private readonly TimeSpan RetryDelay;
+ private readonly T Client;
+ private readonly Uri ServerUri;
+ internal ClientRetryManager(T worker, TimeSpan delay, Uri serverUri)
+ {
+ this.Client = worker;
+ this.RetryDelay = delay;
+ this.ServerUri = serverUri;
+ //Register disconnect listener
+ worker.ConnectionClosed += Worker_Disconnected;
+ }
+ private void Worker_Disconnected(object? sender, EventArgs args)
+ {
+ //Exec retry on exit
+ _ = RetryAsync().ConfigureAwait(false);
+ }
+ /// <summary>
+ /// Raised before client is to be reconnected
+ /// </summary>
+ public event Action<T>? OnBeforeReconnect;
+ /// <summary>
+ /// Raised when the client fails to reconnect. Should return a value that instructs the
+ /// manager to reconnect
+ /// </summary>
+ public event Func<T, Exception, bool>? OnReconnectFailed;
+ async Task RetryAsync()
+ {
+ //Begin random delay with retry ms
+ int randomDelayMs = (int)RetryDelay.TotalMilliseconds;
+ //random delay to add to prevent retry-storm
+ randomDelayMs += RandomNumberGenerator.GetInt32(RetryRandMaxMsDelay);
+ //Retry loop
+ bool retry = true;
+ while (retry)
+ {
+ try
+ {
+ //Inform Listener for the retry
+ OnBeforeReconnect?.Invoke(Client);
+ //wait for delay before reconnecting
+ await Task.Delay(randomDelayMs);
+ //Reconnect async
+ await Client.ConnectAsync(ServerUri).ConfigureAwait(false);
+ break;
+ }
+ catch (Exception Ex)
+ {
+ //Invoke error handler, may be null, incase exit
+ retry = OnReconnectFailed?.Invoke(Client, Ex) ?? false;
+ }
+ }
+ }
+ ///<inheritdoc/>
+ protected override void Free()
+ {
+ //Unregister the event listener
+ Client.ConnectionClosed -= Worker_Disconnected;
+ }
+ }
+using System;
+using VNLib.Net.Messaging.FBM;
+namespace VNLib.Data.Caching
+ public static class Constants
+ {
+ /// <summary>
+ /// Contains constants the define actions
+ /// </summary>
+ public static class Actions
+ {
+ public const string Get= "g";
+ public const string AddOrUpdate = "u";
+ public const string Delete = "d";
+ public const string Dequeue = "dq";
+ }
+ /// <summary>
+ /// Containts constants for operation response codes
+ /// </summary>
+ public static class ResponseCodes
+ {
+ public const string Okay = "ok";
+ public const string Error = "err";
+ public const string NotFound = "nf";
+ }
+ public const HeaderCommand ObjectId = (HeaderCommand)0xAA;
+ public const HeaderCommand NewObjectId = (HeaderCommand)0xAB;
+ }
+using System;
+using VNLib.Net.Messaging.FBM;
+namespace VNLib.Data.Caching.Exceptions
+ /// <summary>
+ /// Raised when the response status code of an FBM Request message is not valid for
+ /// the specified request
+ /// </summary>
+ public class InvalidStatusException : InvalidResponseException
+ {
+ private readonly string? StatusCode;
+ /// <summary>
+ /// Initalizes a new <see cref="InvalidStatusException"/> with the specfied status code
+ /// </summary>
+ /// <param name="message"></param>
+ /// <param name="statusCode"></param>
+ public InvalidStatusException(string message, string statusCode):this(message)
+ {
+ this.StatusCode = statusCode;
+ }
+ ///<inheritdoc/>
+ public InvalidStatusException()
+ {
+ }
+ ///<inheritdoc/>
+ public InvalidStatusException(string message) : base(message)
+ {
+ }
+ ///<inheritdoc/>
+ public InvalidStatusException(string message, Exception innerException) : base(message, innerException)
+ {
+ }
+ ///<inheritdoc/>
+ public override string Message => $"InvalidStatusException: Status Code {StatusCode} \r\n {base.Message}";
+ }
+using System;
+using System.Runtime.Serialization;
+using VNLib.Net.Messaging.FBM;
+namespace VNLib.Data.Caching.Exceptions
+ /// <summary>
+ /// Raised when a request (or server response) calculates the size of the message to be too large to proccess
+ /// </summary>
+ public class MessageTooLargeException : FBMException
+ {
+ ///<inheritdoc/>
+ public MessageTooLargeException()
+ {}
+ ///<inheritdoc/>
+ public MessageTooLargeException(string message) : base(message)
+ {}
+ ///<inheritdoc/>
+ public MessageTooLargeException(string message, Exception innerException) : base(message, innerException)
+ {}
+ ///<inheritdoc/>
+ protected MessageTooLargeException(SerializationInfo info, StreamingContext context) : base(info, context)
+ {}
+ }
+using System;
+namespace VNLib.Data.Caching.Exceptions
+ /// <summary>
+ /// Raised when a command was executed on a desired object in the remote cache
+ /// but the object was not found
+ /// </summary>
+ public class ObjectNotFoundException : InvalidStatusException
+ {
+ internal ObjectNotFoundException()
+ {}
+ internal ObjectNotFoundException(string message) : base(message)
+ {}
+ internal ObjectNotFoundException(string message, string statusCode) : base(message, statusCode)
+ {}
+ internal ObjectNotFoundException(string message, Exception innerException) : base(message, innerException)
+ {}
+ }
+<Project Sdk="Microsoft.NET.Sdk">
+ <PropertyGroup>
+ <TargetFramework>net6.0</TargetFramework>
+ <Platforms>AnyCPU;x64</Platforms>
+ <Authors>Vaughn Nugent</Authors>
+ <Copyright>Copyright © 2022 Vaughn Nugent</Copyright>
+ <Version></Version>
+ <GenerateDocumentationFile>True</GenerateDocumentationFile>
+ <PlatformTarget>x64</PlatformTarget>
+ <Nullable>enable</Nullable>
+ </PropertyGroup>
+ <PropertyGroup Condition="'$(Configuration)|$(Platform)'=='Debug|x64'">
+ <DocumentationFile></DocumentationFile>
+ </PropertyGroup>
+ <ItemGroup>
+ <PackageReference Include="ErrorProne.NET.CoreAnalyzers" Version="0.1.2">
+ <PrivateAssets>all</PrivateAssets>
+ <IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
+ </PackageReference>
+ <PackageReference Include="ErrorProne.NET.Structs" Version="0.1.2">
+ <PrivateAssets>all</PrivateAssets>
+ <IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
+ </PackageReference>
+ </ItemGroup>
+ <ItemGroup>
+ <ProjectReference Include="..\..\..\..\VNLib\Utils\src\VNLib.Utils.csproj" />
+ <ProjectReference Include="..\..\VNLib.Net.Messaging.FBM\src\VNLib.Net.Messaging.FBM.csproj" />
+ </ItemGroup>
new file mode 100644
index 0000000..a309c7c
--- /dev/null
+++ b/VNLib.Data.Caching/src/WaitForChangeResult.cs
@@ -0,0 +1,21 @@
+namespace VNLib.Data.Caching
+ /// <summary>
+ /// The result of a cache server change event
+ /// </summary>
+ public readonly struct WaitForChangeResult
+ {
+ /// <summary>
+ /// The operation status code
+ /// </summary>
+ public readonly string Status { get; init; }
+ /// <summary>
+ /// The current (or old) id of the element that changed
+ /// </summary>
+ public readonly string CurrentId { get; init; }
+ /// <summary>
+ /// The new id of the element that changed
+ /// </summary>
+ public readonly string NewId { get; init; }
+ }