Skip to content
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System;
using System.Diagnostics.CodeAnalysis;
using Microsoft.Extensions.Options;
using Orleans.Configuration;

Expand All @@ -25,8 +26,41 @@ public static IClientBuilder AddAzureQueueStreams(this IClientBuilder builder,
public static IClientBuilder AddAzureQueueStreams(this IClientBuilder builder,
string name, Action<OptionsBuilder<AzureQueueOptions>> configureOptions)
{
builder.AddAzureQueueStreams(name, b=>
b.ConfigureAzureQueue(configureOptions));
builder.AddAzureQueueStreams(name, b => b.ConfigureAzureQueue(configureOptions));
return builder;
}

/// <summary>
/// Configure cluster client to use Azure Queue persistent streams with JSON serialization.
/// This feature is experimental and subject to change in future updates.
/// </summary>
/// <param name="builder">The client builder.</param>
/// <param name="name">The stream provider name.</param>
/// <param name="configure">Configuration delegate for the JSON-enabled Azure Queue stream provider.</param>
/// <returns>The client builder for method chaining.</returns>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static IClientBuilder AddAzureQueueJsonStreams(this IClientBuilder builder,
string name,
Action<ClusterClientAzureQueueJsonStreamConfigurator> configure)
{
var configurator = new ClusterClientAzureQueueJsonStreamConfigurator(name, builder);
configure?.Invoke(configurator);
return builder;
}

/// <summary>
/// Configure cluster client to use Azure Queue persistent streams with JSON serialization and default settings.
/// This feature is experimental and subject to change in future updates.
/// </summary>
/// <param name="builder">The client builder.</param>
/// <param name="name">The stream provider name.</param>
/// <param name="configureOptions">Configuration delegate for Azure Queue options.</param>
/// <returns>The client builder for method chaining.</returns>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static IClientBuilder AddAzureQueueJsonStreams(this IClientBuilder builder,
string name, Action<OptionsBuilder<AzureQueueOptions>> configureOptions)
{
builder.AddAzureQueueJsonStreams(name, b => b.ConfigureAzureQueue(configureOptions));
return builder;
}
}
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System;
using System.Diagnostics.CodeAnalysis;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using Orleans.Configuration;
Expand All @@ -15,8 +16,7 @@ public static class SiloBuilderExtensions
public static ISiloBuilder AddAzureQueueStreams(this ISiloBuilder builder, string name,
Action<SiloAzureQueueStreamConfigurator> configure)
{
var configurator = new SiloAzureQueueStreamConfigurator(name,
configureServicesDelegate => builder.ConfigureServices(configureServicesDelegate));
var configurator = new SiloAzureQueueStreamConfigurator(name, configureServicesDelegate => builder.ConfigureServices(configureServicesDelegate));
configure?.Invoke(configurator);
return builder;
}
Expand All @@ -26,8 +26,39 @@ public static ISiloBuilder AddAzureQueueStreams(this ISiloBuilder builder, strin
/// </summary>
public static ISiloBuilder AddAzureQueueStreams(this ISiloBuilder builder, string name, Action<OptionsBuilder<AzureQueueOptions>> configureOptions)
{
builder.AddAzureQueueStreams(name, b =>
b.ConfigureAzureQueue(configureOptions));
builder.AddAzureQueueStreams(name, b => b.ConfigureAzureQueue(configureOptions));
return builder;
}

/// <summary>
/// Configure silo to use Azure Queue persistent streams with JSON serialization.
/// This feature is experimental and subject to change in future updates.
/// </summary>
/// <param name="builder">The silo builder.</param>
/// <param name="name">The stream provider name.</param>
/// <param name="configure">Configuration delegate for the JSON-enabled Azure Queue stream provider.</param>
/// <returns>The silo builder for method chaining.</returns>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static ISiloBuilder AddAzureQueueJsonStreams(this ISiloBuilder builder, string name,
Action<SiloAzureQueueJsonStreamConfigurator> configure)
{
var configurator = new SiloAzureQueueJsonStreamConfigurator(name, configureServicesDelegate => builder.ConfigureServices(configureServicesDelegate));
configure?.Invoke(configurator);
return builder;
}

/// <summary>
/// Configure silo to use Azure Queue persistent streams with JSON serialization and default settings.
/// This feature is experimental and subject to change in future updates.
/// </summary>
/// <param name="builder">The silo builder.</param>
/// <param name="name">The stream provider name.</param>
/// <param name="configureOptions">Configuration delegate for Azure Queue options.</param>
/// <returns>The silo builder for method chaining.</returns>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static ISiloBuilder AddAzureQueueJsonStreams(this ISiloBuilder builder, string name, Action<OptionsBuilder<AzureQueueOptions>> configureOptions)
{
builder.AddAzureQueueJsonStreams(name, b => b.ConfigureAzureQueue(configureOptions));
return builder;
}

Expand Down
Original file line number Diff line number Diff line change
@@ -1,10 +1,13 @@
using System;
using System.Diagnostics.CodeAnalysis;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.DependencyInjection.Extensions;
using Microsoft.Extensions.Options;
using Orleans.Providers.Streams.AzureQueue;
using Orleans.Configuration;
using Orleans.Serialization;
using Orleans.Streams;
using Orleans.Streaming.AzureStorage.Providers.Streams.AzureQueue.Json;

namespace Orleans.Hosting
{
Expand Down Expand Up @@ -82,4 +85,153 @@ public ClusterClientAzureQueueStreamConfigurator(string name, IClientBuilder bui
this.ConfigureDelegate(services => services.TryAddSingleton<IQueueDataAdapter<string, IBatchContainer>, AzureQueueDataAdapterV2>());
}
}

/// <summary>
/// Silo configurator interface for Azure Queue streams with JSON serialization.
/// This feature is experimental and subject to change in future updates.
/// </summary>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public interface ISiloAzureQueueJsonStreamConfigurator : IAzureQueueStreamConfigurator, ISiloPersistentStreamConfigurator { }

/// <summary>
/// Extension methods for JSON-enabled silo Azure Queue stream configurator.
/// </summary>
public static class SiloAzureQueueJsonStreamConfiguratorExtensions
{
/// <summary>
/// Configures the cache size for the JSON-enabled Azure Queue stream provider.
/// </summary>
/// <param name="configurator">The configurator.</param>
/// <param name="cacheSize">The cache size.</param>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static void ConfigureCacheSize(this ISiloAzureQueueJsonStreamConfigurator configurator, int cacheSize = SimpleQueueCacheOptions.DEFAULT_CACHE_SIZE)
{
configurator.Configure<SimpleQueueCacheOptions>(ob => ob.Configure(options => options.CacheSize = cacheSize));
}

/// <summary>
/// Configures JSON serializer options for the Azure Queue stream provider.
/// </summary>
/// <param name="configurator">The configurator.</param>
/// <param name="configureJsonOptions">Action to configure JSON serializer options.</param>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static void ConfigureJsonSerialization(this ISiloAzureQueueJsonStreamConfigurator configurator, Action<OrleansJsonSerializerOptions> configureJsonOptions)
{
configurator.Configure<OrleansJsonSerializerOptions>(options => options.Configure(configureJsonOptions));
}

/// <summary>
/// Configures the JSON data adapter behavior options.
/// </summary>
/// <param name="configurator">The configurator.</param>
/// <param name="configureAdapterOptions">Action to configure JSON data adapter options.</param>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static void ConfigureJsonAdapter(this ISiloAzureQueueJsonStreamConfigurator configurator, Action<AzureQueueJsonDataAdapterOptions> configureAdapterOptions)
{
configurator.Configure<AzureQueueJsonDataAdapterOptions>(options => options.Configure(configureAdapterOptions));
}
}

/// <summary>
/// Silo configurator for Azure Queue streams with JSON serialization support.
/// This configurator automatically sets up the JSON data adapter and required dependencies.
/// </summary>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public class SiloAzureQueueJsonStreamConfigurator : SiloPersistentStreamConfigurator, ISiloAzureQueueJsonStreamConfigurator
{
/// <summary>
/// Initializes a new instance of the <see cref="SiloAzureQueueJsonStreamConfigurator"/> class.
/// </summary>
/// <param name="name">The stream provider name.</param>
/// <param name="configureServicesDelegate">The delegate used to configure services.</param>
public SiloAzureQueueJsonStreamConfigurator(string name, Action<Action<IServiceCollection>> configureServicesDelegate)
: base(name, configureServicesDelegate, AzureQueueAdapterFactory.Create)
{
this.ConfigureComponent(AzureQueueOptionsValidator.Create);
this.ConfigureComponent(SimpleQueueCacheOptionsValidator.Create);

// Configure default queue names
this.ConfigureAzureQueue(ob => ob.PostConfigure<IOptions<ClusterOptions>>((op, clusterOp) =>
{
if (op.QueueNames == null || op.QueueNames?.Count == 0)
{
op.QueueNames =
AzureQueueStreamProviderUtils.GenerateDefaultAzureQueueNames(clusterOp.Value.ServiceId,
this.Name);
}
}));

this.Configure<OrleansJsonSerializerOptions>(options => { });
this.Configure<AzureQueueJsonDataAdapterOptions>(options => { });
this.ConfigureQueueDataAdapter(AzureQueueJsonDataAdapter.Create);
}
}

/// <summary>
/// Cluster client configurator interface for Azure Queue streams with JSON serialization.
/// This feature is experimental and subject to change in future updates.
/// </summary>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public interface IClusterClientAzureQueueJsonStreamConfigurator : IAzureQueueStreamConfigurator, IClusterClientPersistentStreamConfigurator { }

/// <summary>
/// Extension methods for JSON-enabled cluster client Azure Queue stream configurator.
/// </summary>
public static class ClusterClientAzureQueueJsonStreamConfiguratorExtensions
{
/// <summary>
/// Configures JSON serializer options for the Azure Queue stream provider.
/// </summary>
/// <param name="configurator">The configurator.</param>
/// <param name="configureJsonOptions">Action to configure JSON serializer options.</param>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static void ConfigureJsonSerialization(this IClusterClientAzureQueueJsonStreamConfigurator configurator, Action<OrleansJsonSerializerOptions> configureJsonOptions)
{
configurator.Configure<OrleansJsonSerializerOptions>(options => options.Configure(configureJsonOptions));
}

/// <summary>
/// Configures the JSON data adapter behavior options.
/// </summary>
/// <param name="configurator">The configurator.</param>
/// <param name="configureAdapterOptions">Action to configure JSON data adapter options.</param>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public static void ConfigureJsonAdapter(this IClusterClientAzureQueueJsonStreamConfigurator configurator, Action<AzureQueueJsonDataAdapterOptions> configureAdapterOptions)
{
configurator.Configure<AzureQueueJsonDataAdapterOptions>(options => options.Configure(configureAdapterOptions));
}
}

/// <summary>
/// Cluster client configurator for Azure Queue streams with JSON serialization support.
/// This configurator automatically sets up the JSON data adapter and required dependencies.
/// </summary>
[Experimental("StreamingJsonSerializationExperimental", UrlFormat = "https://github.com/dotnet/orleans/pull/9618")]
public class ClusterClientAzureQueueJsonStreamConfigurator : ClusterClientPersistentStreamConfigurator, IClusterClientAzureQueueJsonStreamConfigurator
{
/// <summary>
/// Initializes a new instance of the <see cref="ClusterClientAzureQueueJsonStreamConfigurator"/> class.
/// </summary>
/// <param name="name">The stream provider name.</param>
/// <param name="builder">The client builder.</param>
public ClusterClientAzureQueueJsonStreamConfigurator(string name, IClientBuilder builder)
: base(name, builder, AzureQueueAdapterFactory.Create)
{
this.ConfigureComponent(AzureQueueOptionsValidator.Create);

// Configure default queue names
this.ConfigureAzureQueue(ob => ob.PostConfigure<IOptions<ClusterOptions>>((op, clusterOp) =>
{
if (op.QueueNames == null || op.QueueNames?.Count == 0)
{
op.QueueNames =
AzureQueueStreamProviderUtils.GenerateDefaultAzureQueueNames(clusterOp.Value.ServiceId, this.Name);
}
}));

this.Configure<OrleansJsonSerializerOptions>(options => { });
this.Configure<AzureQueueJsonDataAdapterOptions>(options => { });
this.ConfigureQueueDataAdapter(AzureQueueJsonDataAdapter.Create);
}
}
}
Loading