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
6 changes: 2 additions & 4 deletions examples/Actor.Next/03-PubSub/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -17,18 +17,16 @@
using Dapr.Actors.Next.Examples.PubSub;
using Dapr.Actors.Next.Streams;
using Dapr.Client;
using Dapr.Messaging.PublishSubscribe.Extensions;

var builder = WebApplication.CreateBuilder(args);
builder.Services.AddDaprActors();
builder.Services.AddDaprActorStreams();
builder.Services.AddDaprPubSubClient();
builder.Services.AddSingleton(_ => new DaprClientBuilder().Build());

var app = builder.Build();

// This sample registers the same subscription declared by [Subscribe] on
// RestockingCartActor.OnRestock so the hosted stream service opens it at startup.
// Register the subscription declared by [Subscribe] on RestockingCartActor.OnRestock.
// AddDaprActorStreams supplies the Dapr.Messaging streaming client.
app.Services.GetRequiredService<ActorStreamSubscriptionRegistry>().Add(
new ActorStreamSubscription(
RestockingCartNames.PubsubName,
Expand Down
2 changes: 1 addition & 1 deletion examples/Actor.Next/03-PubSub/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ One `[Subscribe]` attribute lets a cart react to restock events from the rest of

Tutorial: [Part 3 - Dynamic pub/sub with actors](../../../docs/dotnet-actorsnext/tutorial/part-3.md).

The old SDK pattern needed a separate subscriber service plus hand-rolled routing, retry, and idempotency. Here the stream runner forwards each event through the normal actor invocation path and only acknowledges the pub/sub message after the actor turn commits.
The old SDK pattern needed a separate subscriber service plus hand-rolled routing, retry, and idempotency. Here `AddDaprActorStreams` uses the Dapr.Messaging streaming client, while the stream runner forwards each event through the normal actor invocation path and only acknowledges the pub/sub message after the actor turn commits.

The local app exposes a small HTTP API that marks cart items as waiting for stock, publishes an inventory restock event through Dapr pub/sub, reads cart state, checks item availability, and clears the sample carts between runs.

Expand Down
10 changes: 10 additions & 0 deletions src/Dapr.Actors.Next.Streams/ActorStreamSubscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -53,5 +53,15 @@ public void Validate()
ArgumentException.ThrowIfNullOrWhiteSpace(ActorType);
ArgumentException.ThrowIfNullOrWhiteSpace(MethodName);
ArgumentException.ThrowIfNullOrWhiteSpace(RouteBy);

if (MessageTimeout <= TimeSpan.Zero)
{
throw new ArgumentOutOfRangeException(nameof(MessageTimeout), MessageTimeout, "The message timeout must be positive.");
}

if (MaximumQueuedMessages is <= 0)
{
throw new ArgumentOutOfRangeException(nameof(MaximumQueuedMessages), MaximumQueuedMessages, "The maximum queued message count must be positive.");
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -29,16 +29,25 @@ public sealed class ActorStreamSubscriptionHostedService(
/// <inheritdoc />
public async Task StartAsync(CancellationToken cancellationToken)
{
foreach (var subscription in registry.Subscriptions)
try
{
var handle = await subscriber.SubscribeAsync(subscription, cancellationToken).ConfigureAwait(false);
subscriptions.Add(handle);
logger.LogInformation(
"Opened actor stream subscription {PubsubName}/{Topic} for {ActorType}.{MethodName}.",
subscription.PubsubName,
subscription.Topic,
subscription.ActorType,
subscription.MethodName);
foreach (var subscription in registry.Subscriptions)
{
cancellationToken.ThrowIfCancellationRequested();
var handle = await subscriber.SubscribeAsync(subscription, cancellationToken).ConfigureAwait(false);
subscriptions.Add(handle);
logger.LogInformation(
"Opened actor stream subscription {PubsubName}/{Topic} for {ActorType}.{MethodName}.",
subscription.PubsubName,
subscription.Topic,
subscription.ActorType,
subscription.MethodName);
}
}
catch
{
await DisposeAsync().ConfigureAwait(false);
throw;
}
}

Expand All @@ -51,9 +60,9 @@ public async Task StopAsync(CancellationToken cancellationToken)
/// <inheritdoc />
public async ValueTask DisposeAsync()
{
foreach (var subscription in subscriptions)
for (var index = subscriptions.Count - 1; index >= 0; index--)
{
await subscription.DisposeAsync().ConfigureAwait(false);
await subscriptions[index].DisposeAsync().ConfigureAwait(false);
}

subscriptions.Clear();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.DependencyInjection.Extensions;
using Microsoft.Extensions.Hosting;
using Dapr.Messaging.PublishSubscribe;
using Dapr.Messaging.PublishSubscribe.Extensions;

namespace Dapr.Actors.Next.Streams;

Expand All @@ -29,6 +31,13 @@ public static IServiceCollection AddDaprActorStreams(this IServiceCollection ser
{
ArgumentNullException.ThrowIfNull(services);

// Actor streams use the Dapr.Messaging streaming client directly. Keep this
// registration self-contained while preserving an explicitly configured client.
if (!services.Any(static descriptor => descriptor.ServiceType == typeof(DaprPublishSubscribeClient)))
{
services.AddDaprPubSubClient();
}

services.TryAddSingleton<ActorStreamSubscriptionRegistry>();
services.TryAddSingleton<IActorStreamSubscriptionRegistry>(sp => sp.GetRequiredService<ActorStreamSubscriptionRegistry>());
services.TryAddSingleton<ActorStreamRoutingKeyExtractor>();
Expand Down
46 changes: 46 additions & 0 deletions test/Dapr.Actors.Next.Streams.Test/StreamsTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,17 @@ public void Routing_key_rejects_empty_missing_and_non_scalar_values()
Assert.Throws<ArgumentException>(() => extractor.ExtractActorId(Subscription, new ActorStreamEvent("1", "pubsub", "topic", ReadOnlyMemory<byte>.Empty, new Dictionary<string, string>())));
}

[Fact]
public void Subscription_validation_rejects_invalid_delivery_options()
{
Assert.Throws<ArgumentOutOfRangeException>(() =>
(Subscription with { MessageTimeout = TimeSpan.Zero }).Validate());
Assert.Throws<ArgumentOutOfRangeException>(() =>
(Subscription with { MaximumQueuedMessages = 0 }).Validate());
Assert.Throws<ArgumentOutOfRangeException>(() =>
(Subscription with { MaximumQueuedMessages = -1 }).Validate());
}

[MinimumDaprRuntimeFact("1.18")]
public void Cloudevents_attribute_lookup_is_case_insensitive()
{
Expand Down Expand Up @@ -242,6 +253,23 @@ public async Task Hosted_service_opens_and_disposes_registered_subscriptions()
Assert.All(client.Disposables, disposable => Assert.True(disposable.Disposed));
}

[Fact]
public async Task Hosted_service_disposes_opened_subscriptions_when_startup_fails()
{
var client = new FakePublishSubscribeClient { FailOnSubscription = 2 };
var registry = new ActorStreamSubscriptionRegistry()
.Add(Subscription)
.Add(Subscription with { Topic = "other" });
var service = new ActorStreamSubscriptionHostedService(
registry,
new DaprMessagingActorStreamSubscriber(client, Runner(new FakeInvocationClient())),
NullLogger<ActorStreamSubscriptionHostedService>.Instance);

await Assert.ThrowsAsync<InvalidOperationException>(() => service.StartAsync(CancellationToken.None));
Assert.Single(client.Disposables);
Assert.True(client.Disposables[0].Disposed);
}

[MinimumDaprRuntimeFact("1.18")]
public void Invalid_route_is_poison()
{
Expand Down Expand Up @@ -273,6 +301,17 @@ public void Service_collection_extension_registers_stream_services()
Assert.NotNull(provider.GetRequiredService<ActorStreamSubscriptionRunner>());
Assert.NotNull(provider.GetRequiredService<DaprMessagingActorStreamSubscriber>());
Assert.IsType<DefaultActorStreamFailureClassifier>(provider.GetRequiredService<IActorStreamFailureClassifier>());
Assert.Equal(1, services.Count(descriptor => descriptor.ServiceType == typeof(DaprPublishSubscribeClient)));
}

[Fact]
public void Service_collection_extension_adds_messaging_client_when_not_preconfigured()
{
var services = new ServiceCollection();

services.AddDaprActorStreams();

Assert.Contains(services, descriptor => descriptor.ServiceType == typeof(DaprPublishSubscribeClient));
}

private static ActorStreamSubscriptionRunner Runner(FakeInvocationClient client) =>
Expand Down Expand Up @@ -338,13 +377,20 @@ public FakePublishSubscribeClient()

public TopicMessageHandler? Handler { get; private set; }

public int? FailOnSubscription { get; init; }

public override Task<IAsyncDisposable> SubscribeAsync(
string pubSubName,
string topicName,
DaprSubscriptionOptions options,
TopicMessageHandler messageHandler,
CancellationToken cancellationToken = default)
{
if (FailOnSubscription == Subscriptions.Count + 1)
{
throw new InvalidOperationException("Subscription failed.");
}

Subscriptions.Add((pubSubName, topicName, options));
Handler = messageHandler;
var disposable = new TrackingAsyncDisposable();
Expand Down
Loading