From d04fe12a9085c809b49d5924cd598ec134bdbc11 Mon Sep 17 00:00:00 2001 From: Agash Date: Tue, 22 Sep 2026 12:05:19 +0200 Subject: [PATCH 1/4] feat: accept response metadata on CallAsync --- ObsWebSocket.Core/ObsWebSocketClient.cs | 25 +++++++++++-- .../IWebSocketMessageSerializer.cs | 21 +++++++++-- .../Serialization/JsonMessageSerializer.cs | 30 +++++++++++----- .../Serialization/MsgPackMessageSerializer.cs | 13 +++++-- ObsWebSocket.Tests/ObsWebSocketClientTests.cs | 13 +++---- ObsWebSocket.Tests/SerializerBehaviorTests.cs | 36 +++++++++++++++++++ 6 files changed, 117 insertions(+), 21 deletions(-) diff --git a/ObsWebSocket.Core/ObsWebSocketClient.cs b/ObsWebSocket.Core/ObsWebSocketClient.cs index 8cfa6e6..39f5972 100644 --- a/ObsWebSocket.Core/ObsWebSocketClient.cs +++ b/ObsWebSocket.Core/ObsWebSocketClient.cs @@ -377,6 +377,11 @@ await SendMessageAsync( /// Supplying it from your own JsonSerializerContext keeps the call AOT safe and /// avoids hand building a . /// + /// + /// Metadata for , for a type this library does not know. + /// Without it the response is resolved from this library's context, which has no entry for a + /// type it did not generate. + /// /// Optional override for the request timeout. /// A token to cancel the operation. /// The response data. @@ -388,14 +393,16 @@ public async Task CallRequiredAsync( string requestType, object? requestData = null, JsonTypeInfo? requestTypeInfo = null, + JsonTypeInfo? responseTypeInfo = null, int? timeoutMs = null, CancellationToken cancellationToken = default ) where TResponse : class => - await CallAsync( + await CallAsync( requestType, requestData, requestTypeInfo, + responseTypeInfo, timeoutMs, cancellationToken ) @@ -416,6 +423,11 @@ await CallAsync( /// Supplying it from your own JsonSerializerContext keeps the call AOT safe and /// avoids hand building a . /// + /// + /// Metadata for , for a type this library does not know. + /// Without it the response is resolved from this library's context, which has no entry for a + /// type it did not generate. + /// /// Optional timeout in milliseconds to wait for the response. Defaults to . /// A token to cancel the asynchronous operation. /// @@ -430,6 +442,7 @@ await CallAsync( string requestType, object? requestData = null, JsonTypeInfo? requestTypeInfo = null, + JsonTypeInfo? responseTypeInfo = null, int? timeoutMs = null, CancellationToken cancellationToken = default ) @@ -509,7 +522,7 @@ await SendMessageAsync( return typeof(TResponse) == typeof(object) ? null : RequireConnection() - .Serializer.DeserializePayload(response.ResponseData); + .Serializer.DeserializePayload(response.ResponseData, responseTypeInfo); } catch (Exception ex) when (ex is not OperationCanceledException || cancellationToken.IsCancellationRequested) @@ -549,6 +562,11 @@ await SendMessageAsync( /// Supplying it from your own JsonSerializerContext keeps the call AOT safe and /// avoids hand building a . /// + /// + /// Metadata for , for a type this library does not know. + /// Without it the response is resolved from this library's context, which has no entry for a + /// type it did not generate. + /// /// Optional timeout in milliseconds to wait for the response. Defaults to . /// A token to cancel the asynchronous operation. /// @@ -563,6 +581,7 @@ await SendMessageAsync( string requestType, object? requestData = null, JsonTypeInfo? requestTypeInfo = null, + JsonTypeInfo? responseTypeInfo = null, int? timeoutMs = null, CancellationToken cancellationToken = default ) @@ -614,7 +633,7 @@ await SendMessageAsync( ProcessResponseStatus(response.RequestStatus, requestType, requestId); TResponse? result = RequireConnection() - .Serializer.DeserializeValuePayload(response.ResponseData); + .Serializer.DeserializeValuePayload(response.ResponseData, responseTypeInfo); return (!result.HasValue && Nullable.GetUnderlyingType(typeof(TResponse)) == null) ? throw new ObsWebSocketException( $"Null deserialization for non-nullable value type '{typeof(TResponse).Name}'." diff --git a/ObsWebSocket.Core/Serialization/IWebSocketMessageSerializer.cs b/ObsWebSocket.Core/Serialization/IWebSocketMessageSerializer.cs index a857e8f..9383be8 100644 --- a/ObsWebSocket.Core/Serialization/IWebSocketMessageSerializer.cs +++ b/ObsWebSocket.Core/Serialization/IWebSocketMessageSerializer.cs @@ -1,3 +1,4 @@ +using System.Text.Json.Serialization.Metadata; using ObsWebSocket.Core.Protocol; namespace ObsWebSocket.Core.Serialization; @@ -94,7 +95,15 @@ CancellationToken cancellationToken /// /// Thrown when a payload is present but cannot be deserialized into . /// - TPayload? DeserializePayload(object? rawPayloadData) + /// + /// Metadata for , for a type this library does not know. + /// resolves it from this library's own context. A format that does not + /// read JSON metadata ignores it. + /// + TPayload? DeserializePayload( + object? rawPayloadData, + JsonTypeInfo? typeInfo = null + ) where TPayload : class; /// @@ -107,7 +116,15 @@ CancellationToken cancellationToken /// /// Thrown when a payload is present but cannot be deserialized into . /// - TPayload? DeserializeValuePayload(object? rawPayloadData) + /// + /// Metadata for , for a type this library does not know. + /// resolves it from this library's own context. A format that does not + /// read JSON metadata ignores it. + /// + TPayload? DeserializeValuePayload( + object? rawPayloadData, + JsonTypeInfo? typeInfo = null + ) where TPayload : struct; /// diff --git a/ObsWebSocket.Core/Serialization/JsonMessageSerializer.cs b/ObsWebSocket.Core/Serialization/JsonMessageSerializer.cs index c574b63..7f21226 100644 --- a/ObsWebSocket.Core/Serialization/JsonMessageSerializer.cs +++ b/ObsWebSocket.Core/Serialization/JsonMessageSerializer.cs @@ -174,8 +174,11 @@ public Task SerializeAsync( } /// - public TPayload? DeserializePayload(object? rawPayloadData) - where TPayload : class => DeserializePayloadCore(rawPayloadData); + public TPayload? DeserializePayload( + object? rawPayloadData, + JsonTypeInfo? typeInfo = null + ) + where TPayload : class => DeserializePayloadCore(rawPayloadData, typeInfo); /// public bool TryDeserializePayload(object? rawPayloadData, out TPayload? payload) @@ -198,7 +201,10 @@ public bool TryDeserializePayload(object? rawPayloadData, out TPayload } } - private TPayload? DeserializePayloadCore(object? rawPayloadData) + private TPayload? DeserializePayloadCore( + object? rawPayloadData, + JsonTypeInfo? callerTypeInfo = null + ) where TPayload : class { if ( @@ -314,8 +320,10 @@ out JsonElement dataElement (object)new RequestBatchResponsePayload(batchRequestId, mappedResults); } + // Metadata from the caller's own context is the only way a type this library has + // never heard of can be read, since s_options resolves from the generated context. JsonTypeInfo typeInfo = - (JsonTypeInfo)s_options.GetTypeInfo(typeof(TPayload)); + callerTypeInfo ?? (JsonTypeInfo)s_options.GetTypeInfo(typeof(TPayload)); return jsonElement.Deserialize(typeInfo); } catch (Exception ex) when (ex is not ObsWebSocketSerializationException) @@ -325,8 +333,11 @@ out JsonElement dataElement } /// - public TPayload? DeserializeValuePayload(object? rawPayloadData) - where TPayload : struct => DeserializeValuePayloadCore(rawPayloadData); + public TPayload? DeserializeValuePayload( + object? rawPayloadData, + JsonTypeInfo? typeInfo = null + ) + where TPayload : struct => DeserializeValuePayloadCore(rawPayloadData, typeInfo); /// public bool TryDeserializeValuePayload(object? rawPayloadData, out TPayload? payload) @@ -349,7 +360,10 @@ public bool TryDeserializeValuePayload(object? rawPayloadData, out TPa } } - private TPayload? DeserializeValuePayloadCore(object? rawPayloadData) + private TPayload? DeserializeValuePayloadCore( + object? rawPayloadData, + JsonTypeInfo? callerTypeInfo = null + ) where TPayload : struct { if ( @@ -377,7 +391,7 @@ is not null // Deserialize will return default(TPayload) if JSON is null, which is valid for nullable structs, // but might be undesirable for non-nullable ones (though caught earlier if JSON is explicitly null). JsonTypeInfo typeInfo = - (JsonTypeInfo)s_options.GetTypeInfo(typeof(TPayload)); + callerTypeInfo ?? (JsonTypeInfo)s_options.GetTypeInfo(typeof(TPayload)); return jsonElement.Deserialize(typeInfo); } catch (Exception ex) when (ex is not ObsWebSocketSerializationException) diff --git a/ObsWebSocket.Core/Serialization/MsgPackMessageSerializer.cs b/ObsWebSocket.Core/Serialization/MsgPackMessageSerializer.cs index a0f8480..7b5b964 100644 --- a/ObsWebSocket.Core/Serialization/MsgPackMessageSerializer.cs +++ b/ObsWebSocket.Core/Serialization/MsgPackMessageSerializer.cs @@ -1,4 +1,5 @@ using System.Buffers; +using System.Text.Json.Serialization.Metadata; using MessagePack; using MessagePack.Resolvers; using Microsoft.Extensions.Logging; @@ -107,7 +108,11 @@ public Task SerializeAsync( } /// - public TPayload? DeserializePayload(object? rawPayloadData) + /// MessagePack resolves its own contracts, so is unused. + public TPayload? DeserializePayload( + object? rawPayloadData, + JsonTypeInfo? typeInfo = null + ) where TPayload : class => DeserializePayloadCore(rawPayloadData); /// @@ -163,7 +168,11 @@ public bool TryDeserializePayload(object? rawPayloadData, out TPayload } /// - public TPayload? DeserializeValuePayload(object? rawPayloadData) + /// MessagePack resolves its own contracts, so is unused. + public TPayload? DeserializeValuePayload( + object? rawPayloadData, + JsonTypeInfo? typeInfo = null + ) where TPayload : struct => DeserializeValuePayloadCore(rawPayloadData); /// diff --git a/ObsWebSocket.Tests/ObsWebSocketClientTests.cs b/ObsWebSocket.Tests/ObsWebSocketClientTests.cs index 5c75bd9..eea6527 100644 --- a/ObsWebSocket.Tests/ObsWebSocketClientTests.cs +++ b/ObsWebSocket.Tests/ObsWebSocketClientTests.cs @@ -1,6 +1,7 @@ using System.Collections.Concurrent; using System.Net.WebSockets; using System.Text.Json; +using System.Text.Json.Serialization.Metadata; using Moq; using ObsWebSocket.Core; using ObsWebSocket.Core.Networking; @@ -226,14 +227,14 @@ CancellationToken ct // Mock DeserializePayload for the batch response structure _ = mockSerializer .Setup(s => - s.DeserializePayload>(It.IsAny()) + s.DeserializePayload( + It.IsAny(), + It.IsAny>?>() + ) ) .Returns( - (object? data) => - { - // Simulate the deserialization accurately - return data is RequestBatchResponsePayload typedData ? typedData : null; - } + (object? data, JsonTypeInfo>? _) => + data is RequestBatchResponsePayload typedData ? typedData : null ); // Act diff --git a/ObsWebSocket.Tests/SerializerBehaviorTests.cs b/ObsWebSocket.Tests/SerializerBehaviorTests.cs index 8f4fc66..18cbda1 100644 --- a/ObsWebSocket.Tests/SerializerBehaviorTests.cs +++ b/ObsWebSocket.Tests/SerializerBehaviorTests.cs @@ -1,6 +1,7 @@ using System.Buffers; using System.Reflection; using System.Text.Json; +using System.Text.Json.Serialization; using MessagePack; using Microsoft.Extensions.Logging.Abstractions; using ObsWebSocket.Core; @@ -22,6 +23,33 @@ private static JsonMessageSerializer CreateJsonSerializer() => private static MsgPackMessageSerializer CreateMsgPackSerializer() => new(NullLogger.Instance); + [TestMethod] + public void JsonSerializer_ConsumerType_WithMetadata_Deserializes() + { + JsonElement payload = JsonDocument + .Parse("""{"answer":42,"label":"forty two"}""") + .RootElement.Clone(); + + ConsumerResponse? result = CreateJsonSerializer() + .DeserializePayload(payload, ConsumerContext.Default.ConsumerResponse); + + Assert.IsNotNull(result); + Assert.AreEqual(42, result.Answer); + Assert.AreEqual("forty two", result.Label); + } + + [TestMethod] + public void JsonSerializer_ConsumerType_WithoutMetadata_Throws() + { + JsonElement payload = JsonDocument + .Parse("""{"answer":42,"label":"forty two"}""") + .RootElement.Clone(); + + _ = Assert.ThrowsExactly(() => + CreateJsonSerializer().DeserializePayload(payload) + ); + } + [TestMethod] public void JsonSerializer_DeserializePayload_SceneStubExtensionData_IsAvailable() { @@ -827,3 +855,11 @@ private static bool HasFormatter(IFormatterResolver resolver, Type targetType) return formatter is not null; } } + +/// A response type this library does not generate, standing in for a consumer's own. +internal sealed record ConsumerResponse(int Answer, string Label); + +// The consumer's own context carries its own naming policy, which is what governs the read. +[JsonSourceGenerationOptions(PropertyNamingPolicy = JsonKnownNamingPolicy.CamelCase)] +[JsonSerializable(typeof(ConsumerResponse))] +internal sealed partial class ConsumerContext : JsonSerializerContext; From 6552d7420c89b271b7dadd2899caae84439b4249 Mon Sep 17 00:00:00 2001 From: Agash Date: Tue, 22 Sep 2026 12:30:04 +0200 Subject: [PATCH 2/4] feat: retry requests OBS refuses as not ready --- ObsWebSocket.Core/IObsReconnectDelays.cs | 22 + ObsWebSocket.Core/ObsWebSocketClient.cs | 147 ++++++- .../ObsWebSocketClientBuilder.cs | 6 +- .../ObsWebSocketClientOptions.cs | 28 ++ ObsWebSocket.Core/ObsWebSocketResilience.cs | 64 +-- ...ObsWebSocketServiceCollectionExtensions.cs | 5 +- ObsWebSocket.Core/ReconnectDelays.cs | 14 +- ObsWebSocket.Tests/ClientContractTests.cs | 387 ++++++++++++++++++ ObsWebSocket.Tests/ReadmeCompileCheck.cs | 4 +- ObsWebSocket.Tests/TestUtils.cs | 10 +- README.md | 28 +- 11 files changed, 661 insertions(+), 54 deletions(-) create mode 100644 ObsWebSocket.Core/IObsReconnectDelays.cs diff --git a/ObsWebSocket.Core/IObsReconnectDelays.cs b/ObsWebSocket.Core/IObsReconnectDelays.cs new file mode 100644 index 0000000..79134cc --- /dev/null +++ b/ObsWebSocket.Core/IObsReconnectDelays.cs @@ -0,0 +1,22 @@ +namespace ObsWebSocket.Core; + +/// +/// Supplies the delay before each reconnect attempt. +/// +/// +/// Register an implementation after AddObsWebSocketClient to replace the backoff curve +/// built from . The connection loop still decides how many +/// attempts to make and which failures are fatal. +/// +public interface IObsReconnectDelays +{ + /// + /// Returns the delay to wait before the retry following . + /// + /// Zero-based index of the retry about to be made. + /// A token to cancel the operation. + ValueTask GetDelayAsync( + int retryIndex, + CancellationToken cancellationToken = default + ); +} diff --git a/ObsWebSocket.Core/ObsWebSocketClient.cs b/ObsWebSocket.Core/ObsWebSocketClient.cs index 39f5972..84a8afe 100644 --- a/ObsWebSocket.Core/ObsWebSocketClient.cs +++ b/ObsWebSocket.Core/ObsWebSocketClient.cs @@ -17,6 +17,8 @@ using ObsWebSocket.Core.Protocol.Events; using ObsWebSocket.Core.Protocol.Generated; using ObsWebSocket.Core.Serialization; +using Polly; +using Polly.Registry; namespace ObsWebSocket.Core; @@ -65,13 +67,17 @@ public sealed partial class ObsWebSocketClient : IAsyncDisposable /// Creates the underlying sockets. /// Source of time for timeouts and backoff. /// The instruments to record to. + /// Resolves a pipeline registered under . + /// Supplies the reconnect backoff curve. public ObsWebSocketClient( ILogger logger, ObsSerializerFactory serializerFactory, IOptions options, IWebSocketConnectionFactory? connectionFactory = null, TimeProvider? timeProvider = null, - ObsWebSocketMetrics? metrics = null + ObsWebSocketMetrics? metrics = null, + ResiliencePipelineProvider? pipelines = null, + IObsReconnectDelays? reconnectDelays = null ) { _logger = logger ?? throw new ArgumentNullException(nameof(logger)); @@ -81,12 +87,56 @@ public ObsWebSocketClient( _connectionFactory = connectionFactory ?? new WebSocketConnectionFactory(); _timeProvider = timeProvider ?? TimeProvider.System; _metrics = metrics ?? ObsWebSocketMetrics.Shared; + _pipelines = pipelines; + _reconnectDelays = reconnectDelays; + } + + /// + /// Runs through the NotReady pipeline, which retries only the + /// status OBS documents as retryable. Each attempt sends a fresh request, because OBS pairs a + /// response to the id it was sent with. + /// + private async Task ThroughNotReadyPipelineAsync( + Func> operation, + CancellationToken cancellationToken + ) + { + // Only a pipeline the application registered wins. The fallback is built from this + // client's own options, which for a named client are its own. + ResiliencePipeline pipeline = + _pipelines?.TryGetPipeline( + ObsWebSocketResilience.NotReadyPipelineKey, + out ResiliencePipeline? registered + ) == true + ? registered! + : _notReadyFallback ?? BuildNotReadyFallback(); + + return await pipeline + .ExecuteAsync(async ct => await operation(ct).ConfigureAwait(false), cancellationToken) + .ConfigureAwait(false); + } + + private ResiliencePipeline? _notReadyFallback; + + private ResiliencePipeline BuildNotReadyFallback() + { + NotReadyRetryOptions retry = _options.Value.NotReadyRetry; + if (!retry.Enabled) + { + return _notReadyFallback = ResiliencePipeline.Empty; + } + + ResiliencePipelineBuilder builder = new() { TimeProvider = _timeProvider }; + _ = builder.AddRetry(ObsWebSocketResilience.CreateNotReadyRetryOptions(retry)); + return _notReadyFallback = builder.Build(); } #endregion #region Fields internal readonly ILogger _logger; + private readonly ResiliencePipelineProvider? _pipelines; + private readonly IObsReconnectDelays? _reconnectDelays; private readonly ObsSerializerFactory _serializerFactory; internal readonly IOptions _options; private readonly IWebSocketConnectionFactory _connectionFactory; @@ -438,13 +488,35 @@ await CallAsync( /// Thrown if the client is not connected. /// Thrown if the operation is cancelled via the or the request times out. /// Thrown if is null or empty. - public async Task CallAsync( + public Task CallAsync( string requestType, object? requestData = null, JsonTypeInfo? requestTypeInfo = null, JsonTypeInfo? responseTypeInfo = null, int? timeoutMs = null, CancellationToken cancellationToken = default + ) + where TResponse : class => + ThroughNotReadyPipelineAsync( + ct => + CallOnceAsync( + requestType, + requestData, + requestTypeInfo, + responseTypeInfo, + timeoutMs, + ct + ), + cancellationToken + ); + + private async Task CallOnceAsync( + string requestType, + object? requestData, + JsonTypeInfo? requestTypeInfo, + JsonTypeInfo? responseTypeInfo, + int? timeoutMs, + CancellationToken cancellationToken ) where TResponse : class { @@ -577,13 +649,35 @@ await SendMessageAsync( /// Thrown if the client is not connected. /// Thrown if the operation is cancelled via the or the request times out. /// Thrown if is null or empty. - public async Task CallAsyncValue( + public Task CallAsyncValue( string requestType, object? requestData = null, JsonTypeInfo? requestTypeInfo = null, JsonTypeInfo? responseTypeInfo = null, int? timeoutMs = null, CancellationToken cancellationToken = default + ) + where TResponse : struct => + ThroughNotReadyPipelineAsync( + ct => + CallValueOnceAsync( + requestType, + requestData, + requestTypeInfo, + responseTypeInfo, + timeoutMs, + ct + ), + cancellationToken + ); + + private async Task CallValueOnceAsync( + string requestType, + object? requestData, + JsonTypeInfo? requestTypeInfo, + JsonTypeInfo? responseTypeInfo, + int? timeoutMs, + CancellationToken cancellationToken ) where TResponse : struct { @@ -662,12 +756,55 @@ await SendMessageAsync( } /// - public async Task>> CallBatchAsync( + public Task>> CallBatchAsync( IEnumerable requests, RequestBatchExecutionType? executionType = null, bool? haltOnFailure = null, int? timeoutMs = null, CancellationToken cancellationToken = default + ) => + ThroughNotReadyPipelineAsync( + async ct => + { + List> results = await CallBatchOnceAsync( + requests, + executionType, + haltOnFailure, + timeoutMs, + ct + ) + .ConfigureAwait(false); + + // A batch OBS was not ready for comes back as results, not as a failed call: every + // entry carries NotReady. Raising it is what lets the pipeline see it. + if ( + _options.Value.NotReadyRetry.Enabled + && results.Count > 0 + && results.TrueForAll(result => + result.RequestStatus.Code == (int)RequestStatusCode.NotReady + ) + ) + { + throw new ObsWebSocketRequestException( + "OBS is not ready to perform the request.", + "RequestBatch", + string.Empty, + new RequestStatus(false, (int)RequestStatusCode.NotReady, null), + null + ); + } + + return results; + }, + cancellationToken + ); + + private async Task>> CallBatchOnceAsync( + IEnumerable requests, + RequestBatchExecutionType? executionType, + bool? haltOnFailure, + int? timeoutMs, + CancellationToken cancellationToken ) { ArgumentNullException.ThrowIfNull(requests); @@ -876,7 +1013,7 @@ CancellationToken externalCancellationToken ) { int attempt = 0; - ReconnectDelays reconnectDelays = new(settings); + IObsReconnectDelays reconnectDelays = _reconnectDelays ?? new ReconnectDelays(settings); Debug.Assert(_clientLifetimeCts != null); CancellationToken clientLifetimeToken = _clientLifetimeCts.Token; diff --git a/ObsWebSocket.Core/ObsWebSocketClientBuilder.cs b/ObsWebSocket.Core/ObsWebSocketClientBuilder.cs index 92241ce..b611cfc 100644 --- a/ObsWebSocket.Core/ObsWebSocketClientBuilder.cs +++ b/ObsWebSocket.Core/ObsWebSocketClientBuilder.cs @@ -112,16 +112,16 @@ key is null } /// - /// Registers the reconnect pipeline for this client. + /// Registers the NotReady retry pipeline for this client. /// /// The client to configure. /// The same builder, for chaining. - public static IObsWebSocketClientBuilder WithReconnectPipeline( + public static IObsWebSocketClientBuilder WithNotReadyPipeline( this IObsWebSocketClientBuilder builder ) { ArgumentNullException.ThrowIfNull(builder); - _ = builder.Services.AddObsWebSocketReconnectPipeline(); + _ = builder.Services.AddObsWebSocketNotReadyPipeline(); return builder; } diff --git a/ObsWebSocket.Core/ObsWebSocketClientOptions.cs b/ObsWebSocket.Core/ObsWebSocketClientOptions.cs index 67ffd34..c2f5998 100644 --- a/ObsWebSocket.Core/ObsWebSocketClientOptions.cs +++ b/ObsWebSocket.Core/ObsWebSocketClientOptions.cs @@ -74,6 +74,12 @@ public sealed class ObsWebSocketClientOptions /// public int MaxReconnectDelayMs { get; set; } = 60000; + /// + /// Retries requests OBS refuses with NotReady (207), which it does while changing + /// scene collection or shutting down. Off by default. + /// + public NotReadyRetryOptions NotReadyRetry { get; set; } = new(); + /// /// Gets or sets the largest inbound message the client will assemble, in bytes. /// Defaults to (64 MiB). @@ -95,3 +101,25 @@ public sealed class ObsWebSocketClientOptions public int MaxIncomingMessageBytes { get; set; } = ObsWebSocketClient.DefaultMaxIncomingMessageBytes; } + +/// +/// Retry behaviour for requests OBS refuses with NotReady. +/// +/// +/// OBS rejects the request before the handler runs, so nothing is partially applied and a +/// mutation is as safe to resend as a read. +/// +public sealed class NotReadyRetryOptions +{ + /// Whether to retry. Defaults to . + public bool Enabled { get; set; } + + /// Retries after the first refusal. Defaults to 3. + public int MaxRetryAttempts { get; set; } = 3; + + /// Delay before the first retry, in milliseconds. Defaults to 250. + public int InitialDelayMs { get; set; } = 250; + + /// Ceiling on the delay, in milliseconds. Defaults to 2000. + public int MaxDelayMs { get; set; } = 2000; +} diff --git a/ObsWebSocket.Core/ObsWebSocketResilience.cs b/ObsWebSocket.Core/ObsWebSocketResilience.cs index 64becde..9145c24 100644 --- a/ObsWebSocket.Core/ObsWebSocketResilience.cs +++ b/ObsWebSocket.Core/ObsWebSocketResilience.cs @@ -1,39 +1,39 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; +using ObsWebSocket.Core.Protocol.Generated; using Polly; using Polly.Retry; namespace ObsWebSocket.Core; /// -/// The resilience pipeline governing connection attempts. +/// Resilience pipelines the client executes through. /// /// -/// The pipeline is registered by name, so an application can replace the reconnect policy -/// wholesale by registering its own pipeline under after -/// calling AddObsWebSocketClient, instead of being limited to the reconnect options. +/// Register a pipeline under the same key after AddObsWebSocketClient to replace the +/// default. Reconnect is not a pipeline: a clean disconnect is not an exception, so the +/// connection loop owns that control flow and takes only its delay curve from +/// . /// public static class ObsWebSocketResilience { - /// Key the reconnect pipeline is registered under. - public const string ReconnectPipelineKey = "obs-websocket-reconnect"; + /// Key the NotReady retry pipeline is registered under. + public const string NotReadyPipelineKey = "obs-websocket-not-ready"; /// - /// Registers the default reconnect pipeline. Delays grow by - /// , are capped at - /// , and carry jitter so several - /// clients recovering from one outage do not retry in lockstep. + /// Registers the default NotReady retry pipeline, described by + /// . /// /// The service collection to register into. /// The same collection, for chaining. - public static IServiceCollection AddObsWebSocketReconnectPipeline( + public static IServiceCollection AddObsWebSocketNotReadyPipeline( this IServiceCollection services ) { ArgumentNullException.ThrowIfNull(services); _ = services.AddResiliencePipeline( - ReconnectPipelineKey, + NotReadyPipelineKey, static (builder, context) => { ObsWebSocketClientOptions options = context @@ -43,7 +43,10 @@ this IServiceCollection services builder.TimeProvider = context.ServiceProvider.GetService() ?? TimeProvider.System; - _ = builder.AddRetry(CreateRetryOptions(options)); + if (options.NotReadyRetry.Enabled) + { + _ = builder.AddRetry(CreateNotReadyRetryOptions(options.NotReadyRetry)); + } } ); @@ -51,30 +54,37 @@ this IServiceCollection services } /// - /// Builds the retry strategy described by the reconnect options. + /// Builds the retry strategy for requests OBS refuses as not ready. /// /// Options describing the retry behaviour. - internal static RetryStrategyOptions CreateRetryOptions(ObsWebSocketClientOptions options) + internal static RetryStrategyOptions CreateNotReadyRetryOptions(NotReadyRetryOptions options) { ArgumentNullException.ThrowIfNull(options); - return CreateRetryOptions( - options.ReconnectBackoffMultiplier, - options.InitialReconnectDelayMs, - options.MaxReconnectDelayMs - ); + + return new RetryStrategyOptions + { + MaxRetryAttempts = Math.Max(options.MaxRetryAttempts, 1), + BackoffType = DelayBackoffType.Exponential, + UseJitter = true, + Delay = TimeSpan.FromMilliseconds(Math.Max(options.InitialDelayMs, 0)), + MaxDelay = TimeSpan.FromMilliseconds( + Math.Max(options.MaxDelayMs, options.InitialDelayMs) + ), + ShouldHandle = static args => + ValueTask.FromResult( + args.Outcome.Exception is ObsWebSocketRequestException request + && request.StatusCode == RequestStatusCode.NotReady + ), + }; } /// - /// Builds the retry strategy from the three values that describe the backoff curve. + /// Builds the retry strategy describing the reconnect backoff curve. /// - /// - /// Taken as values rather than as an options object so that a live connection can build its - /// strategy from the settings it was established with, which do not change underneath it. - /// /// Growth applied per attempt. /// Delay before the first retry. /// Ceiling on the delay. - internal static RetryStrategyOptions CreateRetryOptions( + internal static RetryStrategyOptions CreateReconnectRetryOptions( double backoffMultiplier, int initialDelayMs, int maxDelayMs @@ -86,8 +96,6 @@ int maxDelayMs return new RetryStrategyOptions { - // The connection loop decides how many attempts to make, because it also decides - // which failures are fatal. This strategy supplies the delay between them. MaxRetryAttempts = int.MaxValue, UseJitter = true, Delay = TimeSpan.FromMilliseconds(initialMs), diff --git a/ObsWebSocket.Core/ObsWebSocketServiceCollectionExtensions.cs b/ObsWebSocket.Core/ObsWebSocketServiceCollectionExtensions.cs index c221d55..1975100 100644 --- a/ObsWebSocket.Core/ObsWebSocketServiceCollectionExtensions.cs +++ b/ObsWebSocket.Core/ObsWebSocketServiceCollectionExtensions.cs @@ -4,6 +4,7 @@ using Microsoft.Extensions.Options; using ObsWebSocket.Core.Networking; using ObsWebSocket.Core.Serialization; +using Polly.Registry; namespace ObsWebSocket.Core; @@ -149,7 +150,9 @@ private static ObsWebSocketClient Create(IServiceProvider services, string? name options, services.GetRequiredService(), services.GetRequiredService(), - services.GetRequiredService() + services.GetRequiredService(), + services.GetService>(), + services.GetService() ); } } diff --git a/ObsWebSocket.Core/ReconnectDelays.cs b/ObsWebSocket.Core/ReconnectDelays.cs index ee641e1..929dac9 100644 --- a/ObsWebSocket.Core/ReconnectDelays.cs +++ b/ObsWebSocket.Core/ReconnectDelays.cs @@ -8,11 +8,9 @@ namespace ObsWebSocket.Core; /// /// /// The connection loop owns attempt counting and decides which failures are fatal, because a -/// clean disconnect is not an exception and so cannot drive a retry strategy. This type asks the -/// strategy only for the delay, which keeps the backoff curve, its cap, and its jitter in one -/// place that an application can replace. +/// clean disconnect is not an exception and so cannot drive a retry strategy. /// -internal sealed class ReconnectDelays +internal sealed class ReconnectDelays : IObsReconnectDelays { private readonly RetryStrategyOptions? _strategy; private readonly TimeSpan _fixedDelay; @@ -28,7 +26,11 @@ public ReconnectDelays(ObsWebSocketClientOptions options) { ArgumentNullException.ThrowIfNull(options); - _strategy = ObsWebSocketResilience.CreateRetryOptions(options); + _strategy = ObsWebSocketResilience.CreateReconnectRetryOptions( + options.ReconnectBackoffMultiplier, + options.InitialReconnectDelayMs, + options.MaxReconnectDelayMs + ); _fixedDelay = TimeSpan.FromMilliseconds(options.InitialReconnectDelayMs); } @@ -38,7 +40,7 @@ public ReconnectDelays(Networking.ObsConnectionSettings settings) { ArgumentNullException.ThrowIfNull(settings); - _strategy = ObsWebSocketResilience.CreateRetryOptions( + _strategy = ObsWebSocketResilience.CreateReconnectRetryOptions( settings.ReconnectBackoffMultiplier, settings.InitialReconnectDelayMs, settings.MaxReconnectDelayMs diff --git a/ObsWebSocket.Tests/ClientContractTests.cs b/ObsWebSocket.Tests/ClientContractTests.cs index d989f1d..471d1d4 100644 --- a/ObsWebSocket.Tests/ClientContractTests.cs +++ b/ObsWebSocket.Tests/ClientContractTests.cs @@ -3,6 +3,7 @@ using System.Net.WebSockets; using System.Text; using System.Text.Json; +using System.Text.Json.Serialization.Metadata; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; @@ -16,6 +17,7 @@ using ObsWebSocket.Core.Protocol.Events; using ObsWebSocket.Core.Protocol.Generated; using ObsWebSocket.Core.Serialization; +using Polly.Registry; namespace ObsWebSocket.Tests; @@ -688,4 +690,389 @@ public void Dispose() } #endregion + + [TestMethod] + public async Task ACallRefusedAsNotReadyIsRetriedWhenEnabled() + { + int sends = 0; + ( + ObsWebSocketClient client, + Mock mockSerializer, + Mock mockConnection + ) = TestUtils.SetupConnectedClientForceState( + configureOptions: o => + { + o.NotReadyRetry.Enabled = true; + o.NotReadyRetry.MaxRetryAttempts = 3; + o.NotReadyRetry.InitialDelayMs = 0; + o.NotReadyRetry.MaxDelayMs = 0; + } + ); + + await using ObsWebSocketClient owned = client; + + _ = mockConnection + .Setup(ws => + ws.SendAsync( + It.IsAny>(), + It.IsAny(), + true, + It.IsAny() + ) + ) + .Callback( + (ReadOnlyMemory buffer, WebSocketMessageType _, bool _, CancellationToken _) => + { + sends++; + string requestId = ReadRequestId(buffer); + bool ready = sends > 2; + _ = TestUtils.SimulateIncomingResponse( + client, + requestId, + new RequestResponsePayload( + "GetVersion", + requestId, + new RequestStatus( + ready, + (int)(ready ? RequestStatusCode.Success : RequestStatusCode.NotReady), + null + ), + null + ) + ); + } + ) + .Returns(ValueTask.CompletedTask); + + _ = await client.CallAsync("GetVersion"); + + Assert.AreEqual(3, sends); + } + + [TestMethod] + public async Task ACallRefusedAsNotReadyThrowsWhenRetryIsOff() + { + int sends = 0; + ( + ObsWebSocketClient client, + Mock mockSerializer, + Mock mockConnection + ) = TestUtils.SetupConnectedClientForceState(); + + await using ObsWebSocketClient owned = client; + + _ = mockConnection + .Setup(ws => + ws.SendAsync( + It.IsAny>(), + It.IsAny(), + true, + It.IsAny() + ) + ) + .Callback( + (ReadOnlyMemory buffer, WebSocketMessageType _, bool _, CancellationToken _) => + { + sends++; + string requestId = ReadRequestId(buffer); + _ = TestUtils.SimulateIncomingResponse( + client, + requestId, + new RequestResponsePayload( + "GetVersion", + requestId, + new RequestStatus(false, (int)RequestStatusCode.NotReady, null), + null + ) + ); + } + ) + .Returns(ValueTask.CompletedTask); + + ObsWebSocketRequestException error = + await Assert.ThrowsExactlyAsync( + () => client.CallAsync("GetVersion") + ); + + Assert.AreEqual(RequestStatusCode.NotReady, error.StatusCode); + Assert.AreEqual(1, sends); + } + + [TestMethod] + public async Task AnotherFailureCodeIsNotRetried() + { + int sends = 0; + ( + ObsWebSocketClient client, + Mock mockSerializer, + Mock mockConnection + ) = TestUtils.SetupConnectedClientForceState( + configureOptions: o => + { + o.NotReadyRetry.Enabled = true; + o.NotReadyRetry.InitialDelayMs = 0; + o.NotReadyRetry.MaxDelayMs = 0; + } + ); + + await using ObsWebSocketClient owned = client; + + _ = mockConnection + .Setup(ws => + ws.SendAsync( + It.IsAny>(), + It.IsAny(), + true, + It.IsAny() + ) + ) + .Callback( + (ReadOnlyMemory buffer, WebSocketMessageType _, bool _, CancellationToken _) => + { + sends++; + string requestId = ReadRequestId(buffer); + _ = TestUtils.SimulateIncomingResponse( + client, + requestId, + new RequestResponsePayload( + "GetVersion", + requestId, + new RequestStatus( + false, + (int)RequestStatusCode.ResourceNotFound, + null + ), + null + ) + ); + } + ) + .Returns(ValueTask.CompletedTask); + + _ = await Assert.ThrowsExactlyAsync( + () => client.CallAsync("GetVersion") + ); + + Assert.AreEqual(1, sends); + } + + private static string ReadRequestId(ReadOnlyMemory buffer) + { + using JsonDocument document = JsonDocument.Parse(buffer); + return document.RootElement.GetProperty("d").GetProperty("requestId").GetString()!; + } + + [TestMethod] + public async Task AValueCallRefusedAsNotReadyIsRetried() + { + int sends = 0; + ( + ObsWebSocketClient client, + Mock mockSerializer, + Mock mockConnection + ) = TestUtils.SetupConnectedClientForceState(configureOptions: o => + { + o.NotReadyRetry.Enabled = true; + o.NotReadyRetry.InitialDelayMs = 0; + o.NotReadyRetry.MaxDelayMs = 0; + }); + + await using ObsWebSocketClient owned = client; + + _ = mockSerializer + .Setup(x => + x.DeserializeValuePayload( + It.IsAny(), + It.IsAny?>() + ) + ) + .Returns(JsonDocument.Parse("{\"ok\":true}").RootElement.Clone()); + + RespondWith( + client, + mockConnection, + () => ++sends > 1 ? RequestStatusCode.Success : RequestStatusCode.NotReady + ); + + JsonElement? answer = await client.CallAsyncValue("GetStats"); + + Assert.IsNotNull(answer); + Assert.AreEqual(2, sends); + } + + [TestMethod] + public async Task ABatchRefusedAsNotReadyIsRetried() + { + int sends = 0; + ( + ObsWebSocketClient client, + Mock mockSerializer, + Mock mockConnection + ) = TestUtils.SetupConnectedClientForceState(configureOptions: o => + { + o.NotReadyRetry.Enabled = true; + o.NotReadyRetry.InitialDelayMs = 0; + o.NotReadyRetry.MaxDelayMs = 0; + }); + + await using ObsWebSocketClient owned = client; + + _ = mockSerializer + .Setup(x => + x.DeserializePayload( + It.IsAny(), + It.IsAny>?>() + ) + ) + .Returns( + (object? data, JsonTypeInfo>? _) => + data as RequestBatchResponsePayload + ); + + _ = mockConnection + .Setup(ws => + ws.SendAsync( + It.IsAny>(), + It.IsAny(), + true, + It.IsAny() + ) + ) + .Callback( + ( + ReadOnlyMemory buffer, + WebSocketMessageType _, + bool _, + CancellationToken _ + ) => + { + sends++; + string batchId = ReadRequestId(buffer); + bool ready = sends > 1; + _ = TestUtils.SimulateIncomingResponse( + client, + batchId, + new RequestBatchResponsePayload( + batchId, + [ + new RequestResponsePayload( + "GetVersion", + batchId + "_0", + new RequestStatus( + ready, + (int)( + ready + ? RequestStatusCode.Success + : RequestStatusCode.NotReady + ), + null + ), + null + ), + ] + ) + ); + } + ) + .Returns(ValueTask.CompletedTask); + + List> results = await client.CallBatchAsync([ + new BatchRequestItem("GetVersion", null), + ]); + + Assert.AreEqual(1, results.Count); + Assert.AreEqual(2, sends); + } + + [TestMethod] + public void ARegisteredPipelineIsResolvable() + { + ServiceCollection services = new(); + _ = services.AddLogging(); + _ = services + .AddObsWebSocketClient(o => o.ServerUri = new Uri("ws://localhost:4455")) + .WithNotReadyPipeline(); + + using ServiceProvider provider = services.BuildServiceProvider(); + + ResiliencePipelineProvider pipelines = provider.GetRequiredService< + ResiliencePipelineProvider + >(); + + Assert.IsNotNull(pipelines.GetPipeline(ObsWebSocketResilience.NotReadyPipelineKey)); + } + + [TestMethod] + public async Task ARegisteredReconnectDelaySourceIsUsed() + { + StubReconnectDelays delays = new(); + ServiceCollection services = new(); + _ = services.AddLogging(); + _ = services.AddSingleton(delays); + _ = services.AddObsWebSocketClient(o => + { + o.ServerUri = new Uri("ws://127.0.0.1:1"); + o.InitialReconnectDelayMs = 0; + o.MaxReconnectAttempts = 2; + }); + + await using ServiceProvider provider = services.BuildServiceProvider(); + ObsWebSocketClient client = provider.GetRequiredService(); + + _ = await Assert.ThrowsExactlyAsync(() => client.ConnectAsync()); + + Assert.IsGreaterThan(0, delays.Calls); + } + + private sealed class StubReconnectDelays : IObsReconnectDelays + { + public int Calls { get; private set; } + + public ValueTask GetDelayAsync( + int retryIndex, + CancellationToken cancellationToken = default + ) + { + Calls++; + return ValueTask.FromResult(TimeSpan.Zero); + } + } + + private static void RespondWith( + ObsWebSocketClient client, + Mock mockConnection, + Func next + ) => + mockConnection + .Setup(ws => + ws.SendAsync( + It.IsAny>(), + It.IsAny(), + true, + It.IsAny() + ) + ) + .Callback( + ( + ReadOnlyMemory buffer, + WebSocketMessageType _, + bool _, + CancellationToken _ + ) => + { + string requestId = ReadRequestId(buffer); + RequestStatusCode code = next(); + _ = TestUtils.SimulateIncomingResponse( + client, + requestId, + new RequestResponsePayload( + "GetStats", + requestId, + new RequestStatus(code == RequestStatusCode.Success, (int)code, null), + JsonDocument.Parse("{\"ok\":true}").RootElement.Clone() + ) + ); + } + ) + .Returns(ValueTask.CompletedTask); } diff --git a/ObsWebSocket.Tests/ReadmeCompileCheck.cs b/ObsWebSocket.Tests/ReadmeCompileCheck.cs index 71e887d..f831396 100644 --- a/ObsWebSocket.Tests/ReadmeCompileCheck.cs +++ b/ObsWebSocket.Tests/ReadmeCompileCheck.cs @@ -539,7 +539,7 @@ Microsoft.Extensions.Hosting.IHostApplicationBuilder builder .AddObsWebSocketClient("obs") .WithAutoConnect() .WithHealthCheck() - .WithReconnectPipeline(); + .WithNotReadyPipeline(); } internal static void TelemetryAndKeyedRegistration(IServiceCollection services) @@ -551,7 +551,7 @@ internal static void TelemetryAndKeyedRegistration(IServiceCollection services) _ = services.AddObsWebSocketClient("booth", o => o.ServerUri = new Uri("ws://booth:4455")); _ = ObsWebSocketDiagnostics.ActivitySourceName; _ = ObsWebSocketDiagnostics.MeterName; - _ = ObsWebSocketResilience.ReconnectPipelineKey; + _ = ObsWebSocketResilience.NotReadyPipelineKey; } internal static async Task ScreenshotsAsync(ObsWebSocketClient client, CancellationToken ct) diff --git a/ObsWebSocket.Tests/TestUtils.cs b/ObsWebSocket.Tests/TestUtils.cs index 1225cfe..021a190 100644 --- a/ObsWebSocket.Tests/TestUtils.cs +++ b/ObsWebSocket.Tests/TestUtils.cs @@ -209,14 +209,20 @@ internal static ( ObsWebSocketClient client, Mock mockSerializer, Mock mockConnection - ) SetupConnectedClientForceState(TimeProvider? timeProvider = null) + ) SetupConnectedClientForceState( + TimeProvider? timeProvider = null, + Action? configureOptions = null + ) { ( ObsWebSocketClient? client, Mock? mockConnection, Mock? mockSerializer, _ - ) = BuildMockedClientInfrastructure(timeProvider: timeProvider); + ) = BuildMockedClientInfrastructure( + configureOptions: configureOptions, + timeProvider: timeProvider + ); CancellationTokenSource lifetime = new(); SetPrivateField(client, "_clientLifetimeCts", lifetime); diff --git a/README.md b/README.md index e9468ce..0793eac 100644 --- a/README.md +++ b/README.md @@ -655,17 +655,31 @@ Reconnect delays grow by `ReconnectBackoffMultiplier`, are capped at `MaxReconne carry jitter so several clients recovering from one outage do not retry in lockstep. Authentication failures are not retried. -`WithReconnectPipeline()` registers the default pipeline explicitly, which is useful when a host has -its own resilience configuration: +Reconnect is not a Polly pipeline. A clean disconnect is not an exception, so the connection loop +owns attempt counting and takes only the delay from `IObsReconnectDelays`. Register your own +implementation after `AddObsWebSocketClient` to replace the curve. + +## Retrying NotReady + +OBS answers `NotReady` (207) while changing scene collection or shutting down, and documents it as +retryable. It rejects the request before the handler runs, so a mutation is as safe to resend as a +read. Off by default: ```csharp -builder.AddObsWebSocketClient("obs") - .WithAutoConnect() - .WithReconnectPipeline(); +builder.AddObsWebSocketClient("obs", o => +{ + o.NotReadyRetry.Enabled = true; + o.NotReadyRetry.MaxRetryAttempts = 5; +}); ``` -To replace the policy, register your own pipeline under -`ObsWebSocketResilience.ReconnectPipelineKey` after adding the client. +Each attempt sends a fresh request id, because OBS pairs a response to the id it was sent with. +Only 207 is retried. To replace the policy, register a pipeline under +`ObsWebSocketResilience.NotReadyPipelineKey` after adding the client: + +```csharp +builder.AddObsWebSocketClient("obs").WithNotReadyPipeline(); +``` ## Telemetry From be41582e78c7441ffcb8694a25f9bd67f06b18fd Mon Sep 17 00:00:00 2001 From: Agash Date: Tue, 22 Sep 2026 14:54:07 +0200 Subject: [PATCH 3/4] ci: open a pull request when the protocol moves --- .github/workflows/protocol-refresh.yml | 106 +++++++++++++++++++++++++ 1 file changed, 106 insertions(+) create mode 100644 .github/workflows/protocol-refresh.yml diff --git a/.github/workflows/protocol-refresh.yml b/.github/workflows/protocol-refresh.yml new file mode 100644 index 0000000..4ba9f17 --- /dev/null +++ b/.github/workflows/protocol-refresh.yml @@ -0,0 +1,106 @@ +# yaml-language-server: $schema=https://json.schemastore.org/github-workflow.json + +name: Protocol refresh + +on: + workflow_dispatch: + schedule: + - cron: "23 4 * * 1" + +permissions: + contents: read + +jobs: + refresh: + name: Refresh the protocol definition + runs-on: ubuntu-latest + permissions: + contents: write + pull-requests: write + + steps: + - name: Checkout code + uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7 + + - name: Setup .NET SDK + uses: actions/setup-dotnet@a98b56852c35b8e3190ac28c8c2271da59106c68 # v6 + with: + dotnet-version: | + 10.x + 11.0.100-rc.1.26425.128 + + - name: Read the pinned revision + id: pinned + run: | + set -euo pipefail + { + echo "repository=$(jq -r .repository protocol.lock.json)" + echo "path=$(jq -r .path protocol.lock.json)" + echo "commit=$(jq -r .commit protocol.lock.json)" + } >> "$GITHUB_OUTPUT" + + - name: Find the upstream revision + id: upstream + env: + GH_TOKEN: ${{ github.token }} + REPOSITORY: ${{ steps.pinned.outputs.repository }} + FILE_PATH: ${{ steps.pinned.outputs.path }} + run: | + set -euo pipefail + slug="${REPOSITORY#https://github.com/}" + # The newest commit touching the protocol file, not the newest commit on master: + # most upstream commits do not change it. + commit=$(gh api "repos/${slug}/commits?path=${FILE_PATH}&per_page=1" --jq '.[0].sha') + echo "commit=${commit}" >> "$GITHUB_OUTPUT" + + - name: Refresh and regenerate + if: steps.upstream.outputs.commit != steps.pinned.outputs.commit + run: > + dotnet build ObsWebSocket.Core + -t:RefreshObsProtocol + -p:ObsProtocolCommit=${{ steps.upstream.outputs.commit }} + -p:ObsCodegenForceRegeneration=true + + - name: Build and test + if: steps.upstream.outputs.commit != steps.pinned.outputs.commit + run: | + set -euo pipefail + dotnet build ObsWebSocket.sln --configuration Release + dotnet test --project "${{ github.workspace }}/ObsWebSocket.Tests/ObsWebSocket.Tests.csproj" \ + --configuration Release --no-build -- --filter "TestCategory!=Integration" + + - name: Open a pull request + if: steps.upstream.outputs.commit != steps.pinned.outputs.commit + env: + GH_TOKEN: ${{ github.token }} + NEW_COMMIT: ${{ steps.upstream.outputs.commit }} + OLD_COMMIT: ${{ steps.pinned.outputs.commit }} + REPOSITORY: ${{ steps.pinned.outputs.repository }} + run: | + set -euo pipefail + if [ -z "$(git status --porcelain)" ]; then + echo "The refresh produced no diff." + exit 0 + fi + + branch="protocol-refresh/${NEW_COMMIT:0:12}" + if git ls-remote --exit-code --heads origin "${branch}" >/dev/null 2>&1; then + echo "${branch} already exists." + exit 0 + fi + + git config user.name "github-actions[bot]" + git config user.email "41898282+github-actions[bot]@users.noreply.github.com" + git switch -c "${branch}" + git add -A + git commit -m "build: refresh the protocol definition" + git push origin "${branch}" + + gh pr create \ + --title "build: refresh the protocol definition" \ + --body "Pinned revision moves from \`${OLD_COMMIT}\` to \`${NEW_COMMIT}\`. + + ${REPOSITORY}/compare/${OLD_COMMIT}...${NEW_COMMIT} + + Generated with \`ObsCodegenForceRegeneration=true\`, so a change to the definition alone still regenerates. Check the live OBS run before merging: it is the only thing that catches a payload mapped to the wrong field." \ + --head "${branch}" From 62b7fbbb3eb1b263f57ba98fe03af41f13439d9b Mon Sep 17 00:00:00 2001 From: Agash Date: Tue, 22 Sep 2026 14:54:07 +0200 Subject: [PATCH 4/4] build: check the stub types against the obs sources --- CONTRIBUTING.md | 15 ++ .../Generation/Emitter.PayloadHandles.cs | 2 +- ObsWebSocket.Example/Worker.cs | 4 +- .../ObsWebSocket.StubAudit.csproj | 12 ++ ObsWebSocket.StubAudit/Program.cs | 192 ++++++++++++++++++ ObsWebSocket.Tests/ClientContractTests.cs | 69 ++++--- ObsWebSocket.sln | 14 ++ README.md | 10 +- 8 files changed, 280 insertions(+), 38 deletions(-) create mode 100644 ObsWebSocket.StubAudit/ObsWebSocket.StubAudit.csproj create mode 100644 ObsWebSocket.StubAudit/Program.cs diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index d7c8cb3..e977d80 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -74,6 +74,21 @@ To contribute code, you'll need to set up a local development environment: ``` Exits non-zero on the first failed check. +## Checking the stub types + +The stub types have no schema behind them: `protocol.json` types 14 array fields as +`Array`, so their shapes live only as C++ in obs-websocket. `ObsWebSocket.StubAudit` +diffs ours against those sources and reports fields OBS emits that no stub declares, stub fields +nothing emits, and numerics narrower than the C type behind them. + +It needs local clones and is run by hand, so the package builds without an obs-studio checkout: + +```bash +dotnet run --project ObsWebSocket.StubAudit -- /path/to/obs-websocket /path/to/obs-studio +``` + +Exits non-zero when it finds a missing or narrow field. + ## Pull Request Process 🚀 1. **Fork the repository** and create your branch from `master`. diff --git a/ObsWebSocket.Codegen.Tasks/Generation/Emitter.PayloadHandles.cs b/ObsWebSocket.Codegen.Tasks/Generation/Emitter.PayloadHandles.cs index 198d968..929c2bc 100644 --- a/ObsWebSocket.Codegen.Tasks/Generation/Emitter.PayloadHandles.cs +++ b/ObsWebSocket.Codegen.Tasks/Generation/Emitter.PayloadHandles.cs @@ -11,7 +11,7 @@ internal static partial class Emitter /// /// /// These are the handles that cost nothing. An event announcing a scene change already says - /// which scene, by uuid, so acting on it needs no lookup and the result is immune to a rename + /// which scene, by uuid, so acting on it costs no extra request and survives a rename /// that happens between the event arriving and the next request going out. Without them the /// caller reads e.EventData.SceneName and addresses the scene by name again, which is /// the round trip and the race the uuid was there to avoid. diff --git a/ObsWebSocket.Example/Worker.cs b/ObsWebSocket.Example/Worker.cs index 0acd54a..784d4d4 100644 --- a/ObsWebSocket.Example/Worker.cs +++ b/ObsWebSocket.Example/Worker.cs @@ -539,7 +539,7 @@ await _obsClient.Inputs.SetInputTextAsync( // of the loop, no handler bookkeeping, and cancellation ends it cleanly. The // classic events on the client are untouched and still work alongside this. // - // The event already says which scene, by uuid, so acting on it needs no lookup. + // The event carries the scene uuid, so acting on it costs no extra request. // Reading SceneName back off it and addressing the scene by name would add a round // trip and reintroduce the rename race the uuid exists to close. int seconds = @@ -4087,7 +4087,7 @@ CancellationToken cancellationToken ); } - // A miss is worth showing: the lookup already fetched the list, so the client can name what + // The lookup already fetched the list, so a miss can name what // does exist. OBS itself can only answer ResourceNotFound and the name you gave it. try { diff --git a/ObsWebSocket.StubAudit/ObsWebSocket.StubAudit.csproj b/ObsWebSocket.StubAudit/ObsWebSocket.StubAudit.csproj new file mode 100644 index 0000000..0cb00ff --- /dev/null +++ b/ObsWebSocket.StubAudit/ObsWebSocket.StubAudit.csproj @@ -0,0 +1,12 @@ + + + Exe + net10.0 + enable + enable + false + + + + + diff --git a/ObsWebSocket.StubAudit/Program.cs b/ObsWebSocket.StubAudit/Program.cs new file mode 100644 index 0000000..f8f9158 --- /dev/null +++ b/ObsWebSocket.StubAudit/Program.cs @@ -0,0 +1,192 @@ +using System.Reflection; +using System.Text.Json.Serialization; +using System.Text.RegularExpressions; +using ObsWebSocket.Core.Protocol.Common; + +if (args.Length != 2) +{ + Console.Error.WriteLine( + "usage: ObsWebSocket.StubAudit " + ); + return 2; +} + +string obsWebSocket = args[0]; +string obsStudio = args[1]; + +string[] sources = +[ + "src/utils/Obs_ArrayHelper.cpp", + "src/utils/Obs_ObjectHelper.cpp", + "src/utils/Obs_VolumeMeter.cpp", + "src/requesthandler/RequestHandler_Canvases.cpp", + "src/requesthandler/RequestHandler_Ui.cpp", +]; + +Dictionary emitted = new(StringComparer.Ordinal); +Regex assignment = new( + """\w+\s*\[\s*"(?[A-Za-z0-9_]+)"\s*\]\s*=\s*(?[^;]+);""", + RegexOptions.Compiled +); + +foreach (string relative in sources) +{ + string path = Path.Combine(obsWebSocket, relative); + if (!File.Exists(path)) + { + Console.Error.WriteLine($"missing: {path}"); + return 2; + } + + foreach (Match match in assignment.Matches(File.ReadAllText(path))) + { + string field = match.Groups["field"].Value; + if (field.StartsWith("OBS_", StringComparison.Ordinal)) + { + continue; + } + + emitted[field] = match.Groups["expr"].Value.Trim(); + } +} + +// Stub fields are the subject. Response fields are here only so a response field emitted by the +// same helper is not reported as a missing stub field; they are generated from protocol.json. +Dictionary declared = new(StringComparer.Ordinal); +Dictionary stubFields = new(StringComparer.Ordinal); + +foreach (Type type in typeof(SceneStub).Assembly.GetTypes().Where(t => t.IsPublic)) +{ + bool isStub = type.Name.EndsWith("Stub", StringComparison.Ordinal); + bool isResponse = type.Namespace?.EndsWith(".Responses", StringComparison.Ordinal) == true; + if (!isStub && !isResponse) + { + continue; + } + + foreach (PropertyInfo property in type.GetProperties()) + { + if (property.GetCustomAttribute() is not null) + { + continue; + } + + string name = + property.GetCustomAttribute()?.Name ?? property.Name; + declared[name] = (type, property.PropertyType); + if (isStub) + { + stubFields[name] = (type, property.PropertyType); + } + } +} + +// Return types of the libobs accessors, so a field's width can be compared against the C type. +Dictionary returns = new(StringComparer.Ordinal); +Regex declaration = new( + @"EXPORT\s+(?[A-Za-z_][A-Za-z0-9_ \*]*?)\s+(?obs_[A-Za-z0-9_]+)\s*\(", + RegexOptions.Compiled +); +foreach ( + string header in Directory.EnumerateFiles( + Path.Combine(obsStudio, "libobs"), + "*.h", + SearchOption.AllDirectories + ) +) +{ + foreach (Match match in declaration.Matches(File.ReadAllText(header))) + { + returns.TryAdd(match.Groups["name"].Value, match.Groups["type"].Value.Trim()); + } +} + +Dictionary widths = new(StringComparer.Ordinal) +{ + ["int64_t"] = 64, + ["uint64_t"] = 64, + ["long long"] = 64, + ["uint32_t"] = 32, + ["int32_t"] = 32, + ["int"] = 32, + ["size_t"] = 64, +}; + +List missing = []; +List extra = []; +List narrow = []; + +foreach ((string field, string expr) in emitted.OrderBy(pair => pair.Key, StringComparer.Ordinal)) +{ + if (!declared.TryGetValue(field, out (Type Owner, Type Declared) property)) + { + missing.Add($"{field} <- {expr}"); + continue; + } + + // An explicit cast or arithmetic means the C type is not what reaches the wire. + if (expr.Contains("(double)", StringComparison.Ordinal) || expr.Contains('/')) + { + continue; + } + + Match call = Regex.Match(expr, @"\b(obs_[A-Za-z0-9_]+)\s*\("); + if (!call.Success || !returns.TryGetValue(call.Groups[1].Value, out string? cType)) + { + continue; + } + + if (!widths.TryGetValue(cType, out int cWidth)) + { + continue; + } + + Type target = Nullable.GetUnderlyingType(property.Declared) ?? property.Declared; + int managedWidth = Type.GetTypeCode(target) switch + { + TypeCode.Int32 or TypeCode.UInt32 => 32, + TypeCode.Int64 or TypeCode.UInt64 => 64, + TypeCode.Double => 64, + _ => 0, + }; + + bool unsignedC = cType.StartsWith('u') || cType == "size_t"; + if (managedWidth != 0 && (managedWidth < cWidth || (unsignedC && managedWidth == cWidth))) + { + narrow.Add( + $"{property.Owner.Name}.{field}: {target.Name} holds {cType} from {call.Groups[1].Value}" + ); + } +} + +foreach ( + (string field, (Type owner, _)) in stubFields.OrderBy(pair => pair.Key, StringComparer.Ordinal) +) +{ + if (!emitted.ContainsKey(field)) + { + extra.Add($"{owner.Name}.{field}"); + } +} + +Report("Emitted by OBS, absent from the stubs", missing); +Report("Declared on a stub, never emitted", extra); +Report("Narrower than the C type behind them", narrow); + +Console.WriteLine( + $"{emitted.Count} emitted field(s), {stubFields.Count} stub field(s), " + + $"{missing.Count} missing, {extra.Count} extra, {narrow.Count} narrow." +); + +return missing.Count + narrow.Count == 0 ? 0 : 1; + +static void Report(string title, List entries) +{ + Console.WriteLine($"## {title}: {entries.Count}"); + foreach (string entry in entries) + { + Console.WriteLine($" {entry}"); + } + + Console.WriteLine(); +} diff --git a/ObsWebSocket.Tests/ClientContractTests.cs b/ObsWebSocket.Tests/ClientContractTests.cs index 471d1d4..54faf1e 100644 --- a/ObsWebSocket.Tests/ClientContractTests.cs +++ b/ObsWebSocket.Tests/ClientContractTests.cs @@ -699,15 +699,13 @@ public async Task ACallRefusedAsNotReadyIsRetriedWhenEnabled() ObsWebSocketClient client, Mock mockSerializer, Mock mockConnection - ) = TestUtils.SetupConnectedClientForceState( - configureOptions: o => - { - o.NotReadyRetry.Enabled = true; - o.NotReadyRetry.MaxRetryAttempts = 3; - o.NotReadyRetry.InitialDelayMs = 0; - o.NotReadyRetry.MaxDelayMs = 0; - } - ); + ) = TestUtils.SetupConnectedClientForceState(configureOptions: o => + { + o.NotReadyRetry.Enabled = true; + o.NotReadyRetry.MaxRetryAttempts = 3; + o.NotReadyRetry.InitialDelayMs = 0; + o.NotReadyRetry.MaxDelayMs = 0; + }); await using ObsWebSocketClient owned = client; @@ -721,7 +719,12 @@ Mock mockConnection ) ) .Callback( - (ReadOnlyMemory buffer, WebSocketMessageType _, bool _, CancellationToken _) => + ( + ReadOnlyMemory buffer, + WebSocketMessageType _, + bool _, + CancellationToken _ + ) => { sends++; string requestId = ReadRequestId(buffer); @@ -734,7 +737,9 @@ Mock mockConnection requestId, new RequestStatus( ready, - (int)(ready ? RequestStatusCode.Success : RequestStatusCode.NotReady), + (int)( + ready ? RequestStatusCode.Success : RequestStatusCode.NotReady + ), null ), null @@ -771,7 +776,12 @@ Mock mockConnection ) ) .Callback( - (ReadOnlyMemory buffer, WebSocketMessageType _, bool _, CancellationToken _) => + ( + ReadOnlyMemory buffer, + WebSocketMessageType _, + bool _, + CancellationToken _ + ) => { sends++; string requestId = ReadRequestId(buffer); @@ -790,8 +800,8 @@ Mock mockConnection .Returns(ValueTask.CompletedTask); ObsWebSocketRequestException error = - await Assert.ThrowsExactlyAsync( - () => client.CallAsync("GetVersion") + await Assert.ThrowsExactlyAsync(() => + client.CallAsync("GetVersion") ); Assert.AreEqual(RequestStatusCode.NotReady, error.StatusCode); @@ -806,14 +816,12 @@ public async Task AnotherFailureCodeIsNotRetried() ObsWebSocketClient client, Mock mockSerializer, Mock mockConnection - ) = TestUtils.SetupConnectedClientForceState( - configureOptions: o => - { - o.NotReadyRetry.Enabled = true; - o.NotReadyRetry.InitialDelayMs = 0; - o.NotReadyRetry.MaxDelayMs = 0; - } - ); + ) = TestUtils.SetupConnectedClientForceState(configureOptions: o => + { + o.NotReadyRetry.Enabled = true; + o.NotReadyRetry.InitialDelayMs = 0; + o.NotReadyRetry.MaxDelayMs = 0; + }); await using ObsWebSocketClient owned = client; @@ -827,7 +835,12 @@ Mock mockConnection ) ) .Callback( - (ReadOnlyMemory buffer, WebSocketMessageType _, bool _, CancellationToken _) => + ( + ReadOnlyMemory buffer, + WebSocketMessageType _, + bool _, + CancellationToken _ + ) => { sends++; string requestId = ReadRequestId(buffer); @@ -837,11 +850,7 @@ Mock mockConnection new RequestResponsePayload( "GetVersion", requestId, - new RequestStatus( - false, - (int)RequestStatusCode.ResourceNotFound, - null - ), + new RequestStatus(false, (int)RequestStatusCode.ResourceNotFound, null), null ) ); @@ -849,8 +858,8 @@ Mock mockConnection ) .Returns(ValueTask.CompletedTask); - _ = await Assert.ThrowsExactlyAsync( - () => client.CallAsync("GetVersion") + _ = await Assert.ThrowsExactlyAsync(() => + client.CallAsync("GetVersion") ); Assert.AreEqual(1, sends); diff --git a/ObsWebSocket.sln b/ObsWebSocket.sln index 597d0b6..0129734 100644 --- a/ObsWebSocket.sln +++ b/ObsWebSocket.sln @@ -15,6 +15,8 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Solution Items", "Solution .github\workflows\build.yml = .github\workflows\build.yml EndProjectSection EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "ObsWebSocket.StubAudit", "ObsWebSocket.StubAudit\ObsWebSocket.StubAudit.csproj", "{5AA84821-F49F-4236-AA3E-F1F0279099B6}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -61,6 +63,18 @@ Global {8B73365F-A31D-4F3D-5139-545126A1720D}.Release|x64.Build.0 = Release|Any CPU {8B73365F-A31D-4F3D-5139-545126A1720D}.Release|x86.ActiveCfg = Release|Any CPU {8B73365F-A31D-4F3D-5139-545126A1720D}.Release|x86.Build.0 = Release|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Debug|Any CPU.Build.0 = Debug|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Debug|x64.ActiveCfg = Debug|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Debug|x64.Build.0 = Debug|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Debug|x86.ActiveCfg = Debug|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Debug|x86.Build.0 = Debug|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Release|Any CPU.ActiveCfg = Release|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Release|Any CPU.Build.0 = Release|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Release|x64.ActiveCfg = Release|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Release|x64.Build.0 = Release|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Release|x86.ActiveCfg = Release|Any CPU + {5AA84821-F49F-4236-AA3E-F1F0279099B6}.Release|x86.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE diff --git a/README.md b/README.md index 0793eac..f9551f1 100644 --- a/README.md +++ b/README.md @@ -21,7 +21,7 @@ install it and the public surface can still change between versions. Enable the server under *Tools > WebSocket Server Settings*. -Two different compatibility claims are worth separating: +Two compatibility claims, which are not the same thing: - **Base protocol**: OBS Studio 28 or newer, which is where obs-websocket v5 arrived. Connecting, identifying, events and the long-standing requests work against any of those. @@ -150,7 +150,7 @@ ObsWebSocketResourceNotFoundException: No scene named 'Intor'. Available: 'Intro ### Handles from events and responses -Events and responses that carry a uuid expose a handle for it, so acting on one needs no lookup: +Events and responses that carry a uuid expose a handle for it, so acting on one costs no extra request: ```csharp client.Scenes.CurrentProgramSceneChanged += async (_, e) => @@ -170,7 +170,7 @@ SceneItemOperations logo = await client.Scene("Intro").ItemAsync("Logo", cancell await logo.SetEnabledAsync(false, ct); await logo.Scene.GetItemListAsync(ct); -await client.Scene("Intro").Item(3).SetIndexAsync(0, ct); // an id needs no lookup +await client.Scene("Intro").Item(3).SetIndexAsync(0, ct); // an id resolves directly ``` `Item(long)` and `Filter(string)` send nothing, since an id and a filter name are the whole @@ -705,8 +705,8 @@ One activity per request, and one per batch rather than per item. Instruments ar | `obsws.events.dropped` | Events discarded because an event stream's consumer fell behind. | | `obsws.messages.dropped` | Inbound messages discarded without being dispatched. | -The last two are worth wiring up if you rely on events: streams drop the oldest event when full, and -the receive loop ignores a message it cannot read. Both are deliberate and otherwise invisible. +Watch the last two if you rely on events. Streams drop the oldest event when full and the receive +loop ignores a message it cannot read, and neither is visible any other way. Timeouts and reconnect delays run on an injectable `TimeProvider`, so tests can drive them with `FakeTimeProvider`.