feat(streaming): support compact JSON Azure Queue migration - #10507
Conversation
There was a problem hiding this comment.
Pull request overview
Adds an experimental Azure Queue streaming data adapter which uses OrleansJsonSerializer to support JSON payloads for migration scenarios, while preserving the existing binary format behavior via fallback and configuration. This fits into the Azure Queue streaming provider as an opt-in adapter + hosting configurators for silo/client.
Changes:
- Introduces
AzureQueueJsonDataAdapterwith configurable JSON/binary preference and fallback behavior. - Adds JSON-enabled Azure Queue stream configurators and
AddAzureQueueJsonStreams(...)hosting extensions for silo and client. - Adds migration-focused tests and updates the generated public API surface.
Show a summary per file
| File | Description |
|---|---|
| test/Extensions/Orleans.Azure.Tests/Streaming/AzureQueueJsonDataAdapterTests.cs | New tests covering JSON/binary behavior, fallback, and legacy message compatibility. |
| src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/Json/AzureQueueJsonDataAdapterOptions.cs | New experimental options for adapter preference and fallback behavior. |
| src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs | Adds the experimental JSON adapter implementation and DI factory method. |
| src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/AzureQueueStreamBuilder.cs | Adds JSON stream configurator types/extensions for silo and client registration. |
| src/Azure/Orleans.Streaming.AzureStorage/Hosting/SiloBuilderExtensions.cs | Adds AddAzureQueueJsonStreams(...) for silo builder. |
| src/Azure/Orleans.Streaming.AzureStorage/Hosting/ClientBuilderExtensions.cs | Adds AddAzureQueueJsonStreams(...) for client builder. |
| src/api/Azure/Orleans.Streaming.AzureStorage/Orleans.Streaming.AzureStorage.cs | Updates generated API surface for newly introduced experimental APIs. |
Review details
Tip
Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
- Files reviewed: 7/7 changed files
- Comments generated: 2
- Review effort level: Lite
There was a problem hiding this comment.
Review details
Suppressed comments (1)
test/Extensions/Orleans.Azure.Tests/Streaming/AzureQueueJsonDataAdapterTests.cs:54
- The test constructs
OrleansJsonSerializerwith a newOrleansJsonSerializerOptionsinstance, which bypasses DI post-configuration (ConfigureOrleansJsonSerializerOptions) that wires up Orleans' JSON binder/converters. That can make the tests diverge from real runtime behavior. Prefer using the configuredIOptions<OrleansJsonSerializerOptions>fromfixture.Services.
var serializer = this.fixture.Services.GetRequiredService<Serializer>();
var azureQueueDataAdapterV2 = new AzureQueueDataAdapterV2(serializer);
var jsonOrleansSerializer = new OrleansJsonSerializer(Options.Create(new OrleansJsonSerializerOptions()));
return new AzureQueueJsonDataAdapter(
jsonOrleansSerializer,
fallbackAdapter: azureQueueDataAdapterV2,
options ?? new AzureQueueJsonDataAdapterOptions(),
NullLogger<AzureQueueJsonDataAdapter>.Instance);
- Files reviewed: 7/7 changed files
- Comments generated: 1
- Review effort level: Lite
There was a problem hiding this comment.
Review details
Suppressed comments (4)
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:145
- ToQueueMessage materializes the events into a List and then later re-materializes into a List during JSON serialization. Caching the object list avoids the extra allocation/iteration in the JSON path (including when falling back after a failed binary attempt).
var eventList = events.ToList(); try { return _options.PreferJsonsrc/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:128
- The public constructor unnecessarily couples this adapter to AzureQueueDataAdapterV2 by taking a concrete fallbackAdapter type, even though the adapter stores it as an interface. This makes the experimental API less flexible (e.g., cannot supply an alternative IQueueDataAdapter implementation) and forces consumers to reference AzureQueueDataAdapterV2 specifically.
Consider changing the parameter type to IQueueDataAdapter<string, IBatchContainer> (or similar) and updating Create(...), tests, and the generated API surface accordingly.
public AzureQueueJsonDataAdapter( OrleansJsonSerializer jsonSerializer, AzureQueueDataAdapterV2 fallbackAdapter, AzureQueueJsonDataAdapterOptions options, ILogger<AzureQueueJsonDataAdapter> logger)src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:173
- FromQueueMessage treats whitespace-only messages as valid input (ThrowIfNullOrEmpty), but OrleansJsonSerializer.Deserialize returns null for whitespace (IsNullOrWhiteSpace). Using ThrowIfNullOrWhiteSpace aligns the guard with the serializer behavior and provides a clearer argument error for this case.
public IBatchContainer FromQueueMessage(string cloudMsg, long sequenceId) { ArgumentException.ThrowIfNullOrEmpty(cloudMsg, nameof(cloudMsg));test/Extensions/Orleans.Azure.Tests/Streaming/AzureQueueJsonDataAdapterTests.cs:217
- This assertion is overly broad: with fallback disabled, FromQueueMessage should fail specifically due to JSON deserialization of a base64 payload. Asserting the expected exception type makes the test more precise and helps catch unintended behavioral changes.
var jsonAdapterNoFallback = InitializeQueueJsonDataAdapter(new AzureQueueJsonDataAdapterOptions { EnableFallback = false }); Assert.ThrowsAny<Exception>(() => jsonAdapterNoFallback.FromQueueMessage(binaryMsg, token.SequenceNumber)); }- Files reviewed: 7/7 changed files
- Comments generated: 0 new
- Review effort level: Lite
Register provider-keyed JSON adapters with named options, preserve binary fallback behavior, and add generated API plus Orleans 7 binary and Orleans 3 JSON compatibility coverage. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 2cd8574b-4108-4bf5-b49a-4294c21435e1
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 2cd8574b-4108-4bf5-b49a-4294c21435e1
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Serialize the migration envelope without Orleans container metadata while preserving allow-listed polymorphic event and request-context graphs. Continue accepting legacy direct-container JSON and binary messages. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 2cd8574b-4108-4bf5-b49a-4294c21435e1
There was a problem hiding this comment.
Review details
Suppressed comments (3)
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:186
- If JSON deserialization fails and the fallback binary deserialization also fails, the original exception is lost and the caller only sees the fallback failure. Capturing both exceptions (for example via AggregateException) makes troubleshooting much easier.
catch (Exception ex) when (_options.EnableFallback)
{
if (_options.PreferJson)
{
_logger.LogDebug(ex, "Failed to deserialize cloud message using JSON, falling back to binary deserialization");
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:163
- The fallback path can hide the original serialization failure if the fallback attempt also throws (the caller will only see the fallback exception). Wrapping the fallback attempt in its own try/catch and rethrowing an AggregateException (or similar) preserves the original failure for diagnostics.
This issue also appears on line 182 of the same file.
catch (Exception ex) when (_options.EnableFallback)
{
if (_options.PreferJson)
{
_logger.LogDebug(ex, "JSON serialization failed for stream {StreamId}, falling back to binary serialization", streamId);
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:230
- This allocates a new List even though the method already has a List (and the result is immediately parsed and re-emitted). Serializing the existing List avoids an extra allocation and iteration.
var serializedEvents = _jsonSerializer.Serialize(events.Cast<object>().ToList(), typeof(List<object>));- Files reviewed: 7/7 changed files
- Comments generated: 0 new
- Review effort level: Lite
d0d3a09 to
a1ddb62
Compare
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 2cd8574b-4108-4bf5-b49a-4294c21435e1
There was a problem hiding this comment.
Review details
Suppressed comments (3)
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:192
FromQueueMessagefalls back to binary on any exception whenEnableFallbackis true. This can hide meaningful JSON errors (for example, an unsupported compact-envelopeversion) and rethrow aFormatExceptionfromConvert.FromBase64Stringinstead, making diagnosis harder and potentially masking forward-compatibility errors. Consider only falling back to binary when the failure indicates the payload is not valid JSON / not JSON-deserializable (JSON exceptions), and rethrow other errors.
}
catch (Exception ex) when (_options.EnableFallback)
{
if (_options.PreferJson)
{
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:310
TryDeserializeCompactJsoncurrently requiresrequestContextto be present and a JSON object. However,AzureQueueBatchContainerV2supports a null request-context (Dictionary<string, object>?) and compact envelopes coming from other producers could reasonably usenull(or omit the property) when no request context is present. With the current code, such messages will throw and (with fallback enabled) may surface as a base64 error instead of being handled as an empty/null request context.
{
throw new InvalidDataException("The Azure Queue JSON envelope property 'namespace' must be a String or Null.");
}
var keyElement = GetRequiredProperty(streamElement, "key", JsonValueKind.String);
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:164
ToQueueMessageeagerly materializeseventsinto aList<T>even whenPreferJsonis false and no fallback is needed. This adds an avoidable allocation/iteration on the binary-preferred path. Deferring theToList()until the JSON path (or until a fallback is actually required) keeps the common case cheaper.
/// </summary>
public string ToQueueMessage<T>(StreamId streamId, IEnumerable<T> events, StreamSequenceToken? token, Dictionary<string, object>? requestContext)
{
var eventList = events.ToList();
- Files reviewed: 7/7 changed files
- Comments generated: 0 new
- Review effort level: Lite
There was a problem hiding this comment.
Review details
Suppressed comments (2)
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:193
FromQueueMessagefalls back to the binary adapter for any JSON deserialization exception whenPreferJson=true. If the payload is valid JSON but indicates an unsupported envelope version (or otherwise fails JSON parsing/validation), the fallback path will attemptConvert.FromBase64Stringand typically throw aFormatException, obscuring the real error. Consider skipping binary fallback when the message is clearly JSON so that meaningful JSON errors (like unsupportedversion) are preserved.
catch (Exception ex) when (_options.EnableFallback)
{
if (_options.PreferJson)
{
_logger.LogDebug(ex, "Failed to deserialize cloud message using JSON, falling back to binary deserialization");
src/Azure/Orleans.Streaming.AzureStorage/Providers/Streams/AzureQueue/IAzureQueueDataAdapter.cs:133
- The public
AzureQueueJsonDataAdapterconstructor takes a concreteAzureQueueDataAdapterV2asfallbackAdapter, but the implementation stores it asIQueueDataAdapter<string, IBatchContainer>. Taking the interface in the public API makes the adapter easier to test/extend and avoids unnecessarily coupling callers to the built-in binary adapter type.
public AzureQueueJsonDataAdapter(
OrleansJsonSerializer jsonSerializer,
AzureQueueDataAdapterV2 fallbackAdapter,
AzureQueueJsonDataAdapterOptions options,
ILogger<AzureQueueJsonDataAdapter> logger)
- Files reviewed: 7/7 changed files
- Comments generated: 0 new
- Review effort level: Lite
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 2cd8574b-4108-4bf5-b49a-4294c21435e1
Summary
Continues and supersedes #9618 from a maintainer-owned fork because the original same-repository branch does not permit maintainer edits.
This preserves @DeagleGross's original contributor commits and authorship while replaying them onto current
main.StreamIdrepresentationWire format
{"version":1,"stream":{"namespace":"test-namespace","key":"00112233445566778899aabbccddeeff"},"events":["test-event"],"requestContext":{"key":"value"}}Compatibility
Orleans 3.x emits GUID stream keys in canonical
Nformat. Orleans 7+ stores GUID keys using those same UTF-8 bytes, so the modern reader reconstructs the identity directly from the raw key without a key-type discriminator. This also allows arbitrary UTF-8 modern stream keys—including GUID-shaped strings—to round-trip without coercion. Orleans 3.x can consume only keys which are validN-format GUIDs because its streaming API is GUID-based.The exact golden payload emitted by the Orleans 3.x producer in #10508 is covered, along with Unicode and GUID-shaped string keys, shared references, null request contexts, null namespaces, legacy direct-container JSON, Orleans 7 binary messages, and JSON/binary fallback paths.
Application event types must be permitted by the configured JSON serializer.
Microsoft Reviewers: Open in CodeFlow