Skip to content
Open
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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,10 @@ to docs, or any other relevant information.
you disable size enforcement by setting `DisablePayloadErrorLimit` to `true` on the worker.
### Added

- Added experimental SDK payload converter support for values and target types that expose
Temporal transfer type conversion hooks. This lets hook-aware types delegate
their wire representation to the configured payload converter, preserving SDK
behavior such as serialization contexts.
- Added the experimental `TemporalWorkerOptions.PatchActivationCallback`, allowing workers to
decide whether a first non-replay `Workflow.Patched` call should activate a patch during rolling
deployments.
Expand Down
15 changes: 10 additions & 5 deletions src/Temporalio/Client/TemporalClient.AsyncActivity.cs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ internal partial class Impl
/// <inheritdoc />
public override async Task HeartbeatAsyncActivityAsync(HeartbeatAsyncActivityInput input)
{
var converter = input.DataConverterOverride ?? Client.Options.DataConverter;
var converter = DataConverterForAsyncActivity(input.DataConverterOverride);
Payloads? details = null;
if (input.Options?.Details != null && input.Options.Details.Count > 0)
{
Expand Down Expand Up @@ -80,7 +80,7 @@ await converter.ToPayloadsAsync(input.Options.Details).ConfigureAwait(false),
/// <inheritdoc />
public override async Task CompleteAsyncActivityAsync(CompleteAsyncActivityInput input)
{
var converter = input.DataConverterOverride ?? Client.Options.DataConverter;
var converter = DataConverterForAsyncActivity(input.DataConverterOverride);
var result = await converter.ToPayloadAsync(input.Result).ConfigureAwait(false);
if (input.Activity is AsyncActivityHandle.IdReference idRef)
{
Expand Down Expand Up @@ -117,7 +117,7 @@ await Client.Connection.WorkflowService.RespondActivityTaskCompletedAsync(
/// <inheritdoc />
public override async Task FailAsyncActivityAsync(FailAsyncActivityInput input)
{
var converter = input.DataConverterOverride ?? Client.Options.DataConverter;
var converter = DataConverterForAsyncActivity(input.DataConverterOverride);
var failure = await converter.ToFailureAsync(input.Exception).ConfigureAwait(false);
Payloads? lastHeartbeatDetails = null;
if (input.Options?.LastHeartbeatDetails != null &&
Expand Down Expand Up @@ -170,7 +170,7 @@ await Client.Connection.WorkflowService.RespondActivityTaskFailedAsync(
public override async Task ReportCancellationAsyncActivityAsync(
ReportCancellationAsyncActivityInput input)
{
var converter = input.DataConverterOverride ?? Client.Options.DataConverter;
var converter = DataConverterForAsyncActivity(input.DataConverterOverride);
Payloads? details = null;
if (input.Options?.Details != null && input.Options.Details.Count > 0)
{
Expand Down Expand Up @@ -213,6 +213,11 @@ await Client.Connection.WorkflowService.RespondActivityTaskCanceledAsync(
throw new ArgumentException("Unrecognized activity reference type");
}
}

private DataConverter DataConverterForAsyncActivity(DataConverter? dataConverterOverride) =>
dataConverterOverride == null ?
Client.Options.DataConverter :
TemporalTransferTypePayloadConverter.Wrap(dataConverterOverride);
}
}
}
}
2 changes: 2 additions & 0 deletions src/Temporalio/Client/TemporalClient.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
using System;
using System.Linq;
using System.Threading.Tasks;
using Temporalio.Converters;

namespace Temporalio.Client
{
Expand Down Expand Up @@ -29,6 +30,7 @@ public TemporalClient(ITemporalConnection connection, TemporalClientOptions opti
}
}

options.DataConverter = TemporalTransferTypePayloadConverter.Wrap(options.DataConverter);
Connection = connection;
Options = options;
OutboundInterceptor = new Impl(this);
Expand Down
32 changes: 32 additions & 0 deletions src/Temporalio/Converters/ITemporalTransferTypeConverter.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
using System;

namespace Temporalio.Converters
{
/// <summary>
/// Converter for a type marked with <see cref="TemporalTransferTypeConverterAttribute"/>.
/// </summary>
/// <remarks>
/// This API is experimental and may change in a future release.
/// </remarks>
public interface ITemporalTransferTypeConverter
{
/// <summary>
/// Gets the transfer type handed to or read from the payload converter.
/// </summary>
Type TransferType { get; }

/// <summary>
/// Convert a value to the transfer type value that should be passed to the payload converter.
/// </summary>
/// <param name="value">Value to convert.</param>
/// <returns>Transfer type value to pass to the payload converter.</returns>
object? ToTransferType(object? value);

/// <summary>
/// Convert a transfer type value from the payload converter to the marked type.
/// </summary>
/// <param name="transferType">Transfer type value returned by the payload converter.</param>
/// <returns>Converted value.</returns>
object? FromTransferType(object? transferType);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
using System;

namespace Temporalio.Converters
{
/// <summary>
/// Marks a type as converting to and from a transfer type before payload conversion.
/// </summary>
/// <remarks>
/// This is used by the SDK payload converter to delegate hook-aware values to the configured
/// payload converter using the transfer type representation provided by
/// <see cref="ITemporalTransferTypeConverter"/>.
/// This API is experimental and may change in a future release.
/// </remarks>
[AttributeUsage(AttributeTargets.Class | AttributeTargets.Struct, Inherited = false)]
public sealed class TemporalTransferTypeConverterAttribute : Attribute
{
/// <summary>
/// Initializes a new instance of the
/// <see cref="TemporalTransferTypeConverterAttribute"/> class.
/// </summary>
/// <param name="converterType">Converter type implementing
/// <see cref="ITemporalTransferTypeConverter"/>.</param>
public TemporalTransferTypeConverterAttribute(Type converterType) =>
ConverterType = converterType;

/// <summary>
/// Gets the converter type.
/// </summary>
public Type ConverterType { get; }
}
}
128 changes: 128 additions & 0 deletions src/Temporalio/Converters/TemporalTransferTypePayloadConverter.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
using System;
using System.Collections.Concurrent;
using System.Reflection;
using Temporalio.Api.Common.V1;

namespace Temporalio.Converters
{
/// <summary>
/// Payload converter wrapper that applies Temporal transfer type hooks.
/// </summary>
internal sealed class TemporalTransferTypePayloadConverter :
IPayloadConverter,
IWithSerializationContext<IPayloadConverter>
{
private static readonly ConcurrentDictionary<Type, ITemporalTransferTypeConverter?> Converters = new();
private readonly IPayloadConverter inner;

private TemporalTransferTypePayloadConverter(IPayloadConverter inner) => this.inner = inner;

/// <summary>
/// Wrap a payload converter unless it is already wrapped.
/// </summary>
/// <param name="payloadConverter">Payload converter to wrap.</param>
/// <returns>Wrapped payload converter.</returns>
public static IPayloadConverter Wrap(IPayloadConverter payloadConverter) =>
payloadConverter is TemporalTransferTypePayloadConverter ?
payloadConverter : new TemporalTransferTypePayloadConverter(payloadConverter);

/// <summary>
/// Wrap the data converter's payload converter unless it is already wrapped.
/// </summary>
/// <param name="dataConverter">Data converter to wrap.</param>
/// <returns>Data converter with a wrapped payload converter.</returns>
public static DataConverter Wrap(DataConverter dataConverter)
{
var payloadConverter = Wrap(dataConverter.PayloadConverter);
return ReferenceEquals(payloadConverter, dataConverter.PayloadConverter) ?
dataConverter : dataConverter with { PayloadConverter = payloadConverter };
}

/// <inheritdoc />
public Payload ToPayload(object? value)
{
var converter = value == null ? null : Converters.GetOrAdd(value.GetType(), CreateConverter);
if (converter != null)
{
value = converter.ToTransferType(value);
}
return inner.ToPayload(value);
}

/// <inheritdoc />
public object? ToValue(Payload payload, Type type)
{
var converter = Converters.GetOrAdd(type, CreateConverter);
if (converter == null)
{
return inner.ToValue(payload, type);
}

var transferTypeValue = inner.ToValue(payload, converter.TransferType);
return converter.FromTransferType(transferTypeValue);
}

/// <inheritdoc/>
public IPayloadConverter WithSerializationContext(ISerializationContext context)
{
if (inner is not IWithSerializationContext<IPayloadConverter> withContext)
{
return this;
}

var contextInner = withContext.WithSerializationContext(context);
return ReferenceEquals(contextInner, inner) ? this : new TemporalTransferTypePayloadConverter(contextInner);
}

private static ITemporalTransferTypeConverter? CreateConverter(Type type)
{
var attr = type.GetCustomAttribute<TemporalTransferTypeConverterAttribute>(
inherit: false);
if (attr == null)
{
return null;
}

if (!typeof(ITemporalTransferTypeConverter).IsAssignableFrom(attr.ConverterType))
{
throw new InvalidOperationException(
$"Type {type} has a Temporal transfer type converter type " +
$"{attr.ConverterType} that does not implement {nameof(ITemporalTransferTypeConverter)}.");
}
if (attr.ConverterType.IsAbstract)
{
throw new InvalidOperationException(
$"Type {type} has an abstract Temporal transfer type converter type " +
$"{attr.ConverterType}.");
}
if (attr.ConverterType.ContainsGenericParameters)
{
throw new InvalidOperationException(
$"Type {type} has an open generic Temporal transfer type converter type " +
$"{attr.ConverterType}.");
}
if (!attr.ConverterType.IsValueType &&
attr.ConverterType.GetConstructor(Type.EmptyTypes) == null)
{
throw new InvalidOperationException(
$"Type {type} has a Temporal transfer type converter type " +
$"{attr.ConverterType} without a public parameterless constructor.");
}

if (Activator.CreateInstance(attr.ConverterType) is not ITemporalTransferTypeConverter converter)
Comment thread
tconley1428 marked this conversation as resolved.
Comment thread
tconley1428 marked this conversation as resolved.
{
throw new InvalidOperationException(
$"Type {type} has a Temporal transfer type converter type " +
$"{attr.ConverterType} that could not be instantiated.");
}
if (converter.TransferType == null)
{
throw new InvalidOperationException(
$"Type {type} has a Temporal transfer type converter type " +
$"{attr.ConverterType} with a null transfer type.");
}

return converter;
}
}
}
25 changes: 24 additions & 1 deletion tests/Temporalio.Tests/Common/PluginTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
using Temporalio.Client;
using Temporalio.Common;
using Temporalio.Converters;
using Temporalio.Tests.Converters;
using Temporalio.Worker;
using Temporalio.Workflows;
using Xunit.Abstractions;
Expand Down Expand Up @@ -205,6 +206,28 @@ public void TestSimplePlugin_Function()
Assert.NotNull(client.Options.DataConverter.PayloadCodec);
}

[Fact]
public void TestClientPlugin_DataConverter_WrapsTransferTypePayloadConverter()
{
var plugin = new SimplePlugin("SimplePlugin", new SimplePluginOptions()
{
DataConverterOption = new SimplePluginOptions.SimplePluginOption<DataConverter>(
(_) => new DataConverter(
new PayloadConverterTests.ContextStringPayloadConverter(),
new DefaultFailureConverter())),
});
var newOptions = (TemporalClientOptions)Client.Options.Clone();
newOptions.Plugins = new[] { plugin };

var client = new TemporalClient(Env.Client.Connection, newOptions);
var dataConverter = client.Options.DataConverter.WithSerializationContext(
new ISerializationContext.Workflow("default", "workflow-id"));
var payload = dataConverter.PayloadConverter.ToPayload(
new PayloadConverterTests.TransferTypeHookValue("payload-value"));

Assert.Equal("workflow-id:payload-value", payload.Data.ToStringUtf8());
}

[Fact]
public async Task TestSimplePlugin_RunContext()
{
Expand Down Expand Up @@ -236,4 +259,4 @@ public async Task TestSimplePlugin_RunContext()
Assert.Contains("Beginning", transitions);
Assert.Contains("Ending", transitions);
}
}
}
32 changes: 28 additions & 4 deletions tests/Temporalio.Tests/Converters/DataConverterTests.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
namespace Temporalio.Tests.Converters;

using System;
using Google.Protobuf;
using Temporalio.Api.Common.V1;
using Temporalio.Converters;
using Xunit;
Expand All @@ -16,17 +17,40 @@ public DataConverterTests(ITestOutputHelper output)
[Fact]
public void NewDataConverter_WithPayloadConverter_ProperlyInitializes()
{
var payloadConverter = new MyPayloadConverter();
var newConverter = DataConverter.Default with
{
PayloadConverter = new MyPayloadConverter(),
PayloadConverter = payloadConverter,
};
Assert.IsType<MyPayloadConverter>(newConverter.PayloadConverter);
Assert.Same(payloadConverter, newConverter.PayloadConverter);
Assert.Equal(
"payload",
newConverter.PayloadConverter.ToValue(
newConverter.PayloadConverter.ToPayload("payload"), typeof(string)));
}

[Fact]
public void NewDataConverter_EquivalentConverters_AreEqual()
{
var payloadConverter = new MyPayloadConverter();
var failureConverter = new DefaultFailureConverter();

Assert.Equal(
new DataConverter(payloadConverter, failureConverter),
new DataConverter(payloadConverter, failureConverter));
}

public class MyPayloadConverter : IPayloadConverter
{
public Payload ToPayload(object? value) => throw new NotImplementedException();
public Payload ToPayload(object? value) => new()
{
Metadata =
{
["encoding"] = ByteString.CopyFromUtf8("test/plain"),
},
Data = ByteString.CopyFromUtf8((string)value!),
};

public object? ToValue(Payload payload, Type type) => throw new NotImplementedException();
public object? ToValue(Payload payload, Type type) => payload.Data.ToStringUtf8();
}
}
Loading
Loading