From 01c379e6b40b4f10a4d67ad737639679343b6bbc Mon Sep 17 00:00:00 2001 From: birschick-bq Date: Tue, 6 Oct 2026 13:19:25 -0700 Subject: [PATCH] THRIFT-5830: Add per-call transport support for netstd Client: cpp,netstd Add opt-in per-call transports for generated asynchronous clients, isolating HTTP and layered transport state while preserving existing shared-protocol constructors. Preserve caller-configured HTTP timeouts and propagate the linked cancellation token through generated request and response operations. Add focused transport, lifecycle, concurrency, timeout, generator, and tutorial coverage. Co-Authored-By: GitHub Copilot --- .github/agents/openspec.agent.md | 65 +++ .github/workflows/copilot-setup-steps.yml | 26 ++ .../src/thrift/generate/t_netstd_generator.cc | 33 +- lib/netstd/Makefile.am | 2 + lib/netstd/README.md | 24 + .../Transports/TBaseClientPerCallTests.cs | 423 ++++++++++++++++++ .../Transports/THttpTransportPerCallTests.cs | 393 ++++++++++++++++ .../Transports/THttpTransportTests.cs | 40 ++ lib/netstd/Thrift/TBaseClient.cs | 252 ++++++++++- .../Thrift/Transport/Client/THttpTransport.cs | 175 ++++++-- .../Transport/ITPerCallTransportProvider.cs | 47 ++ .../Transport/Layered/TBufferedTransport.cs | 48 +- .../Transport/Layered/TFramedTransport.cs | 48 +- lib/netstd/openspec/changes/archive/.gitkeep | 0 .../.openspec.yaml | 2 + .../2026-10-06-per-call-transport/design.md | 39 ++ .../2026-10-06-per-call-transport/proposal.md | 28 ++ .../specs/per-call-transport/spec.md | 89 ++++ .../2026-10-06-per-call-transport/tasks.md | 42 ++ lib/netstd/openspec/config.yaml | 33 ++ lib/netstd/openspec/specs/.gitkeep | 0 .../openspec/specs/per-call-transport/spec.md | 89 ++++ lib/netstd/per-call-transport.md | 30 ++ tutorial/netstd/Client/Program.cs | 61 +++ tutorial/netstd/README.md | 7 + tutorial/netstd/smoketest.sh | 14 + 26 files changed, 1969 insertions(+), 41 deletions(-) create mode 100644 .github/agents/openspec.agent.md create mode 100644 .github/workflows/copilot-setup-steps.yml create mode 100644 lib/netstd/Tests/Thrift.Tests/Transports/TBaseClientPerCallTests.cs create mode 100644 lib/netstd/Tests/Thrift.Tests/Transports/THttpTransportPerCallTests.cs create mode 100644 lib/netstd/Thrift/Transport/ITPerCallTransportProvider.cs create mode 100644 lib/netstd/openspec/changes/archive/.gitkeep create mode 100644 lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/.openspec.yaml create mode 100644 lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/design.md create mode 100644 lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/proposal.md create mode 100644 lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/specs/per-call-transport/spec.md create mode 100644 lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/tasks.md create mode 100644 lib/netstd/openspec/config.yaml create mode 100644 lib/netstd/openspec/specs/.gitkeep create mode 100644 lib/netstd/openspec/specs/per-call-transport/spec.md create mode 100644 lib/netstd/per-call-transport.md diff --git a/.github/agents/openspec.agent.md b/.github/agents/openspec.agent.md new file mode 100644 index 00000000000..2d8c52578da --- /dev/null +++ b/.github/agents/openspec.agent.md @@ -0,0 +1,65 @@ +--- +name: OpenSpec +description: "Manages OpenSpec changes, specs, and workflows using the OpenSpec CLI. Use this agent for proposing changes, exploring ideas, validating artifacts, checking status, and archiving completed work." +tools: + - "execute" + - "read" + - "search" + - "edit" +--- + + + +# OpenSpec Agent + +You are a specialized agent for managing OpenSpec workflows. Before using the `openspec` CLI, run `openspec --version`. If it is unavailable, install it with `npm install -g @fission-ai/openspec@1.14.1`. + +## What is OpenSpec? + +OpenSpec is a structured change management system for codebases. It organizes work into **changes** with planning artifacts (proposals, specs, designs, tasks) that guide implementation. + +## Available Commands + +### Agent-Compatible CLI Commands (prefer `--json` for structured output) + +| Command | Purpose | +|---------|---------| +| `openspec list [--json]` | List all changes and specs | +| `openspec show [--json]` | View a specific change or spec | +| `openspec validate [--all] [--json]` | Validate changes and specs for issues | +| `openspec status [--change ] [--json]` | Show artifact progress for a change | +| `openspec instructions [artifact] [--change ] [--json]` | Get next-step instructions for an artifact | +| `openspec templates [--json]` | List available templates | +| `openspec schemas [--json]` | List available workflow schemas | +| `openspec archive --json [--yes]` | Archive a completed change; use `--yes` only after confirming all tasks are complete | + +### Interactive CLI Commands (use when prompted by the user) + +| Command | Purpose | +|---------|---------| +| `openspec init` | Initialize OpenSpec in the project | +| `openspec update` | Update OpenSpec configuration and artifacts | +| `openspec view` | Interactive dashboard | +| `openspec config` | View or modify settings | + +## Workflow + +Run OpenSpec commands from `lib/netstd`, where this change's OpenSpec project is located. + +1. Run `openspec list --json` to see active changes and specs. +2. Run `openspec status --change --json` for the selected change. +3. Run `openspec instructions [artifact] --change --json` for the next-step instructions. +4. Run `openspec validate --json` before completing. + +## Key Directories + +- `lib/netstd/openspec/` — OpenSpec directory +- `lib/netstd/openspec/changes/` — Active changes and their artifacts +- `lib/netstd/openspec/config.yaml` — Project configuration + +## Best Practices + +- Always use `--json` when parsing output programmatically +- Run `openspec validate` after creating or modifying artifacts +- Check `openspec status` before starting work to understand the current state +- When archiving, ensure all tasks are completed and validated first \ No newline at end of file diff --git a/.github/workflows/copilot-setup-steps.yml b/.github/workflows/copilot-setup-steps.yml new file mode 100644 index 00000000000..e8fd0b4df30 --- /dev/null +++ b/.github/workflows/copilot-setup-steps.yml @@ -0,0 +1,26 @@ +# Generated by OpenSpec for GitHub Copilot coding agent support. + +name: "Copilot Setup Steps" + +# Run manually to verify the setup; Copilot runs these steps when starting agent work. +on: + workflow_dispatch: + +jobs: + # The job MUST be called `copilot-setup-steps` for Copilot coding agent to pick it up. + copilot-setup-steps: + runs-on: ubuntu-latest + timeout-minutes: 10 + + permissions: + contents: read + + steps: + - name: Checkout code + uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + + - name: Install OpenSpec CLI + run: npm install -g @fission-ai/openspec@1.14.1 + + - name: Verify OpenSpec CLI + run: openspec --version \ No newline at end of file diff --git a/compiler/cpp/src/thrift/generate/t_netstd_generator.cc b/compiler/cpp/src/thrift/generate/t_netstd_generator.cc index 09a999d9885..fc4bc925501 100644 --- a/compiler/cpp/src/thrift/generate/t_netstd_generator.cc +++ b/compiler/cpp/src/thrift/generate/t_netstd_generator.cc @@ -2432,13 +2432,30 @@ void t_netstd_generator::generate_service_client(ostream& out, t_service* tservi << indent() << "{" << '\n'; indent_up(); - out << indent() << "public Client(TProtocol protocol) : this(protocol, protocol)" << '\n' + out << indent() << "/// Initializes the client with a shared protocol." << '\n' + << indent() << "/// Use this constructor when the application already owns and configures the protocol. It preserves shared-protocol behavior; use the transport and protocol-factory constructor to opt into per-call transports." << '\n' + << indent() << "/// The shared protocol used for input and output." << '\n' + << indent() << "public Client(TProtocol protocol) : this(protocol, protocol)" << '\n' << indent() << "{" << '\n' << indent() << "}" << '\n' << '\n' + << indent() << "/// Initializes the client with shared input and output protocols." << '\n' + << indent() << "/// This constructor remains available for callers that manage protocol instances directly and preserves their shared-transport semantics. For concurrent per-call operations, use a supporting transport with the protocol-factory constructor." << '\n' + << indent() << "/// The shared input protocol." << '\n' + << indent() << "/// The shared output protocol." << '\n' << indent() << "public Client(TProtocol inputProtocol, TProtocol outputProtocol) : base(inputProtocol, outputProtocol)" << '\n' << indent() << "{" << '\n' << indent() << "}" << '\n' + << '\n' + << indent() << "/// Initializes the client with protocols created for a supplied transport." << '\n' + << indent() << "/// When supports , generated asynchronous calls use a separate transport and protocol state per call. Existing protocol constructors retain shared-protocol behavior." << '\n' + << indent() << "/// The shared transport or per-call transport provider." << '\n' + << indent() << "/// The factory used to create input protocols." << '\n' + << indent() << "/// The factory used to create output protocols." << '\n' + << indent() << "public Client(TTransport transport, TProtocolFactory inputProtocolFactory, TProtocolFactory outputProtocolFactory)" << '\n' + << indent() << " : base(transport, inputProtocolFactory, outputProtocolFactory)" << '\n' + << indent() << "{" << '\n' + << indent() << "}" << '\n' << '\n'; vector functions = tservice->get_functions(); @@ -2448,23 +2465,31 @@ void t_netstd_generator::generate_service_client(ostream& out, t_service* tservi { string raw_func_name = (*functions_iterator)->get_name(); string function_name = raw_func_name + (add_async_postfix ? "Async" : ""); + string call_cancellation_token = tmp("callCancellationToken"); // async generate_deprecation_attribute(out, (*functions_iterator)->annotations_); out << indent() << "public async " << function_signature_async(*functions_iterator, "") << '\n' << indent() << "{" << '\n'; indent_up(); + bool returns_value = !(*functions_iterator)->is_oneway() && !(*functions_iterator)->get_returntype()->is_void(); + out << indent() << (returns_value ? "return await " : "await ") + << "ExecutePerCallAsync(async (" << call_cancellation_token << ") =>" << '\n' + << indent() << "{" << '\n'; + indent_up(); out << indent() << "await send_" << function_name << "("; string call_args = argument_list((*functions_iterator)->get_arglist(),false); if(! call_args.empty()) { out << call_args << ", "; } - out << CANCELLATION_TOKEN_NAME << ");" << '\n'; + out << call_cancellation_token << ");" << '\n'; if(! (*functions_iterator)->is_oneway()) { - out << indent() << ((*functions_iterator)->get_returntype()->is_void() ? "" : "return ") - << "await recv_" << function_name << "(" << CANCELLATION_TOKEN_NAME << ");" << '\n'; + out << indent() << (returns_value ? "return " : "") + << "await recv_" << function_name << "(" << call_cancellation_token << ");" << '\n'; } indent_down(); + out << indent() << "}, " << CANCELLATION_TOKEN_NAME << ");" << '\n'; + indent_down(); out << indent() << "}" << '\n' << '\n'; // async send diff --git a/lib/netstd/Makefile.am b/lib/netstd/Makefile.am index d811c9ad8c3..00fb4794adf 100644 --- a/lib/netstd/Makefile.am +++ b/lib/netstd/Makefile.am @@ -74,6 +74,8 @@ dist-hook: EXTRA_DIST = \ README.md \ + per-call-transport.md \ + openspec \ Directory.Build.props \ Benchmarks/Thrift.Benchmarks \ Tests/codegen/run-NetStd-Codegen-Tests.ps1 \ diff --git a/lib/netstd/README.md b/lib/netstd/README.md index a3623e1c599..b472ff60233 100644 --- a/lib/netstd/README.md +++ b/lib/netstd/README.md @@ -14,6 +14,30 @@ The library ships as two packages so that non-web projects no longer pull in the *in addition to* `ApacheThrift` only when you host Thrift over ASP.NET Core. Existing code keeps compiling unchanged once the package reference is added. +# Per-call HTTP client calls + +For concurrent asynchronous calls, construct the generated client with a transport that +implements `ITPerCallTransportProvider` and protocol factories. `THttpTransport`, +`TBufferedTransport`, and `TFramedTransport` over a supporting transport provide this capability. Each generated +high-level call then gets independent transport and protocol state while the HTTP connection +pool is shared: + +```csharp +var transport = new THttpTransport(new Uri("http://localhost:9090"), new TConfiguration()); +var protocolFactory = new TBinaryProtocol.Factory(); +using var client = new Calculator.Client(transport, protocolFactory, protocolFactory); + +var results = await Task.WhenAll(Enumerable.Range(1, 8) + .Select(value => client.add(value, value, cancellationToken))); +``` + +The existing constructors that accept `TProtocol` instances remain the shared-protocol path +for applications that manage protocol instances directly. Per-call behavior is opt-in; a +transport without the provider capability continues to use shared behavior. The per-call scope +covers generated high-level asynchronous methods that perform a complete RPC. Calling generated +`send_*` and `recv_*` methods separately continues to use the shared protocols and must not be +interleaved on the same client. + # Build the library ## How to build on Windows diff --git a/lib/netstd/Tests/Thrift.Tests/Transports/TBaseClientPerCallTests.cs b/lib/netstd/Tests/Thrift.Tests/Transports/TBaseClientPerCallTests.cs new file mode 100644 index 00000000000..3206e9f0884 --- /dev/null +++ b/lib/netstd/Tests/Thrift.Tests/Transports/TBaseClientPerCallTests.cs @@ -0,0 +1,423 @@ +// Licensed to the Apache Software Foundation(ASF) under one +// or more contributor license agreements.See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership.The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +using System; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.VisualStudio.TestTools.UnitTesting; +using Thrift.Protocol; +using Thrift.Transport; +using Thrift.Transport.Client; + +namespace Thrift.Tests.Transports +{ + [TestClass] + public class TBaseClientPerCallTests + { + [TestMethod] + public async Task Existing_protocol_constructor_keeps_shared_protocols() + { + var transport = new TMemoryBufferTransport(null); + var protocol = new TBinaryProtocol.Factory().GetProtocol(transport); + using var client = new TestClient(protocol); + + await client.RunAsync(() => + { + Assert.AreSame(protocol, client.InputProtocol); + Assert.AreSame(protocol, client.OutputProtocol); + return Task.CompletedTask; + }); + } + + [TestMethod] + public void Provider_constructor_rejects_null_transport() + { + try + { + _ = new TestClient((TTransport)null!); + Assert.Fail("Expected the provider constructor to reject a null transport."); + } + catch (ArgumentNullException exception) + { + Assert.AreEqual("transport", exception.ParamName); + } + } + + [TestMethod] + public void Per_call_operation_rejects_null_delegate() + { + using var rootTransport = new MemoryPerCallTransport(); + using var client = new TestClient(rootTransport); + + try + { + client.RunAsync((Func)null!); + Assert.Fail("Expected the per-call API to reject a null operation."); + } + catch (ArgumentNullException exception) + { + Assert.AreEqual("operation", exception.ParamName); + } + } + + [TestMethod] + public async Task Per_call_provider_rejects_null_transport_result() + { + using var rootTransport = new NullPerCallTransport(); + using var client = new TestClient(rootTransport); + + try + { + await client.RunAsync(() => Task.CompletedTask); + Assert.Fail("Expected the client to reject a null transport result."); + } + catch (InvalidOperationException exception) + { + Assert.AreEqual("The per-call transport provider returned null.", exception.Message); + } + } + + [TestMethod] + public async Task Provider_constructor_creates_protocols_for_each_call() + { + using var rootTransport = new MemoryPerCallTransport(); + using var client = new TestClient(rootTransport); + TTransport? firstCallTransport = null; + TTransport? secondCallTransport = null; + + await client.RunAsync(() => + { + firstCallTransport = client.InputProtocol.Transport; + Assert.AreSame(firstCallTransport, client.OutputProtocol.Transport); + return Task.CompletedTask; + }); + + await client.RunAsync(() => + { + secondCallTransport = client.InputProtocol.Transport; + Assert.AreSame(secondCallTransport, client.OutputProtocol.Transport); + return Task.CompletedTask; + }); + + Assert.AreNotSame(rootTransport, firstCallTransport); + Assert.AreNotSame(firstCallTransport, secondCallTransport); + } + + [TestMethod] + public async Task Per_call_timeout_cancels_operation_token_without_canceling_caller_token() + { + using var rootTransport = new MemoryPerCallTransport(TimeSpan.FromMilliseconds(100)); + using var client = new TestClient(rootTransport); + using var callerCancellation = new CancellationTokenSource(); + var operationStarted = NewSignal(); + var operationToken = CancellationToken.None; + + var operation = client.RunAsync(async cancellationToken => + { + operationToken = cancellationToken; + operationStarted.TrySetResult(true); + await Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken); + return 1; + }, callerCancellation.Token); + + await operationStarted.Task.WaitAsync(TimeSpan.FromSeconds(5)); + var completed = await Task.WhenAny(operation, Task.Delay(TimeSpan.FromSeconds(2))); + if (completed != operation) + { + callerCancellation.Cancel(); + await CaptureTransportExceptionAsync(() => operation); + Assert.Fail("The per-call timeout did not cancel the operation token."); + } + var exception = await CaptureTransportExceptionAsync(() => operation); + + Assert.AreEqual(TTransportException.ExceptionType.Interrupted, exception.Type); + Assert.IsTrue(operationToken.IsCancellationRequested); + Assert.IsFalse(callerCancellation.IsCancellationRequested); + } + + [TestMethod] + public async Task Caller_token_is_combined_with_per_call_timeout() + { + using var rootTransport = new MemoryPerCallTransport(TimeSpan.FromSeconds(5)); + using var client = new TestClient(rootTransport); + using var callerCancellation = new CancellationTokenSource(); + var operationStarted = NewSignal(); + var operationToken = CancellationToken.None; + + var operation = client.RunAsync(async cancellationToken => + { + operationToken = cancellationToken; + operationStarted.TrySetResult(true); + await Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken); + return 1; + }, callerCancellation.Token); + + await operationStarted.Task.WaitAsync(TimeSpan.FromSeconds(5)); + callerCancellation.Cancel(); + var exception = await CaptureTransportExceptionAsync(() => operation); + + Assert.AreEqual(TTransportException.ExceptionType.Interrupted, exception.Type); + Assert.IsTrue(operationToken.IsCancellationRequested); + Assert.IsTrue(callerCancellation.IsCancellationRequested); + } + + [TestMethod] + public async Task Generated_rpc_uses_the_linked_per_call_token_for_send() + { + using var rootTransport = new BlockingPerCallTransport(TimeSpan.FromMilliseconds(100)); + var protocolFactory = new TBinaryProtocol.Factory(); + using var client = new global::ThriftTest.ThriftTest.Client(rootTransport, + protocolFactory, protocolFactory); + using var callerCancellation = new CancellationTokenSource(); + + var operation = client.testVoid(callerCancellation.Token); + var completed = await Task.WhenAny(operation, Task.Delay(TimeSpan.FromSeconds(2))); + if (completed != operation) + { + callerCancellation.Cancel(); + await CaptureTransportExceptionAsync(() => operation); + Assert.Fail("The generated RPC did not use the per-call timeout token."); + } + + var exception = await CaptureTransportExceptionAsync(() => operation); + Assert.AreEqual(TTransportException.ExceptionType.Interrupted, exception.Type); + Assert.IsFalse(callerCancellation.IsCancellationRequested); + } + + [TestMethod] + public async Task Buffered_wrapper_delegates_per_call_transport_creation() + { + using var rootTransport = new MemoryPerCallTransport(); + var sharedWrapper = new TBufferedTransport(rootTransport); + using var client = new TestClient(sharedWrapper); + TBufferedTransport? callWrapper = null; + TTransport? callUnderlyingTransport = null; + + await client.RunAsync(() => + { + callWrapper = client.InputProtocol.Transport as TBufferedTransport; + Assert.IsNotNull(callWrapper); + Assert.AreSame(callWrapper, client.OutputProtocol.Transport); + callUnderlyingTransport = callWrapper.UnderlyingTransport; + return Task.CompletedTask; + }); + + Assert.AreNotSame(sharedWrapper, callWrapper); + Assert.AreNotSame(rootTransport, callUnderlyingTransport); + } + + [TestMethod] + public async Task Buffered_wrapper_without_per_call_support_keeps_shared_protocols() + { + var sharedTransport = new TBufferedTransport(new TMemoryBufferTransport(null)); + using var client = new TestClient(sharedTransport); + + await client.RunAsync(() => + { + Assert.AreSame(sharedTransport, client.InputProtocol.Transport); + Assert.AreSame(sharedTransport, client.OutputProtocol.Transport); + return Task.CompletedTask; + }); + } + + [TestMethod] + public async Task Framed_wrapper_delegates_per_call_transport_creation() + { + using var rootTransport = new MemoryPerCallTransport(); + var sharedWrapper = new TFramedTransport(rootTransport); + using var client = new TestClient(sharedWrapper); + TFramedTransport? callWrapper = null; + TTransport? callUnderlyingTransport = null; + + await client.RunAsync(() => + { + callWrapper = client.InputProtocol.Transport as TFramedTransport; + Assert.IsNotNull(callWrapper); + Assert.AreSame(callWrapper, client.OutputProtocol.Transport); + callUnderlyingTransport = callWrapper.InnerTransport; + return Task.CompletedTask; + }); + + Assert.AreNotSame(sharedWrapper, callWrapper); + Assert.AreNotSame(rootTransport, callUnderlyingTransport); + Assert.AreEqual(rootTransport.PerCallTimeout, sharedWrapper.PerCallTimeout); + } + + [TestMethod] + public async Task Framed_wrapper_without_per_call_support_keeps_shared_protocols() + { + var sharedTransport = new TFramedTransport(new TMemoryBufferTransport(null)); + using var client = new TestClient(sharedTransport); + + await client.RunAsync(() => + { + Assert.AreSame(sharedTransport, client.InputProtocol.Transport); + Assert.AreSame(sharedTransport, client.OutputProtocol.Transport); + return Task.CompletedTask; + }); + } + + [TestMethod] + public async Task Buffered_wrapper_delegates_through_framed_transport() + { + using var rootTransport = new MemoryPerCallTransport(); + var sharedFramedTransport = new TFramedTransport(rootTransport); + var sharedBufferedTransport = new TBufferedTransport(sharedFramedTransport); + using var client = new TestClient(sharedBufferedTransport); + TBufferedTransport? callBufferedTransport = null; + TFramedTransport? callFramedTransport = null; + TTransport? callUnderlyingTransport = null; + + await client.RunAsync(() => + { + callBufferedTransport = client.InputProtocol.Transport as TBufferedTransport; + Assert.IsNotNull(callBufferedTransport); + Assert.AreSame(callBufferedTransport, client.OutputProtocol.Transport); + callFramedTransport = callBufferedTransport.UnderlyingTransport as TFramedTransport; + Assert.IsNotNull(callFramedTransport); + callUnderlyingTransport = callFramedTransport.InnerTransport; + return Task.CompletedTask; + }); + + Assert.AreNotSame(sharedBufferedTransport, callBufferedTransport); + Assert.AreNotSame(sharedFramedTransport, callFramedTransport); + Assert.AreNotSame(rootTransport, callUnderlyingTransport); + Assert.AreEqual(rootTransport.PerCallTimeout, sharedBufferedTransport.PerCallTimeout); + } + + [TestMethod] + public void Layered_transport_base_does_not_advertise_per_call_support() + { + Assert.IsFalse(typeof(ITPerCallTransportProvider).IsAssignableFrom(typeof(TLayeredTransport))); + } + + private sealed class MemoryPerCallTransport : TMemoryBufferTransport, ITPerCallTransportProvider + { + public MemoryPerCallTransport(TimeSpan? perCallTimeout = null) + : base(null) + { + PerCallTimeout = perCallTimeout ?? Timeout.InfiniteTimeSpan; + } + + public TimeSpan PerCallTimeout { get; } + + public bool SupportsPerCallTransport => true; + + public Task CreatePerCallTransportAsync(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + return Task.FromResult(new TMemoryBufferTransport(Configuration)); + } + } + + private sealed class BlockingPerCallTransport : TMemoryBufferTransport, ITPerCallTransportProvider + { + public BlockingPerCallTransport(TimeSpan perCallTimeout) + : base(null) + { + PerCallTimeout = perCallTimeout; + } + + public TimeSpan PerCallTimeout { get; } + + public bool SupportsPerCallTransport => true; + + public Task CreatePerCallTransportAsync(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + return Task.FromResult(new BlockingMemoryBufferTransport(Configuration)); + } + } + + private sealed class NullPerCallTransport : TMemoryBufferTransport, ITPerCallTransportProvider + { + public NullPerCallTransport() + : base(null) + { + } + + public TimeSpan PerCallTimeout => Timeout.InfiniteTimeSpan; + + public bool SupportsPerCallTransport => true; + + public Task CreatePerCallTransportAsync(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + return Task.FromResult(null!); + } + } + + private sealed class BlockingMemoryBufferTransport : TMemoryBufferTransport + { + public BlockingMemoryBufferTransport(TConfiguration configuration) + : base(configuration) + { + } + + public override Task WriteAsync(byte[] buffer, int offset, int count, + CancellationToken cancellationToken) + { + return Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken); + } + } + + private sealed class TestClient : TBaseClient, System.IDisposable + { + private static readonly TProtocolFactory ProtocolFactory = new TBinaryProtocol.Factory(); + + public TestClient(TProtocol protocol) + : base(protocol, protocol) + { + } + + public TestClient(TTransport transport) + : base(transport, ProtocolFactory, ProtocolFactory) + { + } + + public Task RunAsync(System.Func operation) + { + return ExecutePerCallAsync(operation, default); + } + + public Task RunAsync(Func> operation, + CancellationToken cancellationToken) + { + return ExecutePerCallAsync(operation, cancellationToken); + } + } + + private static TaskCompletionSource NewSignal() + { + return new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + } + + private static async Task CaptureTransportExceptionAsync(Func operation) + { + try + { + await operation(); + } + catch (TTransportException exception) + { + return exception; + } + + throw new InvalidOperationException("Expected a transport interruption."); + } + } +} \ No newline at end of file diff --git a/lib/netstd/Tests/Thrift.Tests/Transports/THttpTransportPerCallTests.cs b/lib/netstd/Tests/Thrift.Tests/Transports/THttpTransportPerCallTests.cs new file mode 100644 index 00000000000..9c6b482926c --- /dev/null +++ b/lib/netstd/Tests/Thrift.Tests/Transports/THttpTransportPerCallTests.cs @@ -0,0 +1,393 @@ +// Licensed to the Apache Software Foundation(ASF) under one +// or more contributor license agreements.See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership.The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +using System; +using System.Buffers.Binary; +using System.Net; +using System.Net.Http; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.VisualStudio.TestTools.UnitTesting; +using Thrift.Protocol; +using Thrift.Transport; +using Thrift.Transport.Client; + +namespace Thrift.Tests.Transports +{ + [TestClass] + public class THttpTransportPerCallTests + { + [TestMethod] + public Task Interleaved_http_calls_keep_request_and_response_state_isolated() + { + return RunInterleavedCallsAsync(useBufferedTransport: false); + } + + [TestMethod] + public Task Buffered_transport_read_ahead_stays_with_its_call() + { + return RunInterleavedCallsAsync(useBufferedTransport: true); + } + + [TestMethod] + public Task Framed_http_calls_keep_frame_buffers_per_call() + { + return RunInterleavedCallsAsync(useBufferedTransport: false, useFramedTransport: true); + } + + [TestMethod] + public Task Buffered_over_framed_http_calls_keep_each_layer_per_call() + { + return RunInterleavedCallsAsync(useBufferedTransport: true, useFramedTransport: true); + } + + [TestMethod] + public async Task Successful_call_disposes_its_response_content() + { + var content = new TrackingContent(new byte[] { 30, 31 }); + var handler = new RequestHandler((_, _) => + { + return Task.FromResult(new HttpResponseMessage(HttpStatusCode.OK) { Content = content }); + }); + using var httpClient = new HttpClient(handler); + using var rootTransport = new THttpTransport(httpClient, null, new Uri("http://localhost")); + using var client = new TestClient(rootTransport); + + var result = await SendAndReceiveAsync(client, 30, CancellationToken.None); + + CollectionAssert.AreEqual(new byte[] { 30, 31 }, result); + Assert.IsTrue(content.IsDisposed); + } + + [TestMethod] + public async Task Failed_call_releases_resources_and_leaves_sibling_usable() + { + var handler = new RequestHandler((payload, _) => + { + if (payload == 40) + { + throw new HttpRequestException("Expected request failure."); + } + return Task.FromResult(Response(payload)); + }); + using var httpClient = new HttpClient(handler); + using var rootTransport = new THttpTransport(httpClient, null, new Uri("http://localhost")); + using var client = new TestClient(rootTransport); + + var exception = await CaptureExceptionAsync( + () => SendAndReceiveAsync(client, 40, CancellationToken.None)); + Assert.AreEqual(TTransportException.ExceptionType.Unknown, exception.Type); + + var result = await SendAndReceiveAsync(client, 41, CancellationToken.None); + CollectionAssert.AreEqual(new byte[] { 41, 42 }, result); + } + + [TestMethod] + public async Task Caller_cancellation_interrupts_only_its_call() + { + var started = NewSignal(); + var handler = new RequestHandler(async (payload, cancellationToken) => + { + if (payload == 50) + { + started.TrySetResult(true); + await Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken); + } + return Response(payload); + }); + using var httpClient = new HttpClient(handler); + using var rootTransport = new THttpTransport(httpClient, null, new Uri("http://localhost")); + using var client = new TestClient(rootTransport); + using var cancellation = new CancellationTokenSource(); + + var interruptedCall = SendAndReceiveAsync(client, 50, cancellation.Token); + await started.Task.WaitAsync(TimeSpan.FromSeconds(5)); + var siblingResult = await SendAndReceiveAsync(client, 51, CancellationToken.None) + .WaitAsync(TimeSpan.FromSeconds(5)); + cancellation.Cancel(); + + var exception = await CaptureExceptionAsync(() => interruptedCall); + Assert.AreEqual(TTransportException.ExceptionType.Interrupted, exception.Type); + CollectionAssert.AreEqual(new byte[] { 51, 52 }, siblingResult); + } + + [TestMethod] + public async Task Configured_timeout_interrupts_request_without_canceling_sibling() + { + var started = NewSignal(); + var handler = new RequestHandler(async (payload, cancellationToken) => + { + if (payload == 60) + { + started.TrySetResult(true); + await Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken); + } + return Response(payload); + }); + using var httpClient = new HttpClient(handler); + using var rootTransport = new THttpTransport(httpClient, null, new Uri("http://localhost")) + { + ConnectTimeout = 1000 + }; + using var client = new TestClient(rootTransport); + + var timedOutCall = SendAndReceiveAsync(client, 60, CancellationToken.None); + await started.Task.WaitAsync(TimeSpan.FromSeconds(5)); + var siblingResult = await SendAndReceiveAsync(client, 61, CancellationToken.None) + .WaitAsync(TimeSpan.FromSeconds(5)); + + var exception = await CaptureExceptionAsync( + () => timedOutCall.WaitAsync(TimeSpan.FromSeconds(5))); + Assert.AreEqual(TTransportException.ExceptionType.Interrupted, exception.Type); + CollectionAssert.AreEqual(new byte[] { 61, 62 }, siblingResult); + } + + [TestMethod] + public async Task Root_disposal_waits_for_active_per_call_transport() + { + var started = NewSignal(); + var release = NewSignal(); + var handler = new RequestHandler(async (payload, cancellationToken) => + { + started.TrySetResult(true); + await release.Task.WaitAsync(cancellationToken); + return Response(payload); + }); + var httpClient = new HttpClient(handler); + var rootTransport = new THttpTransport(httpClient, null, new Uri("http://localhost")); + var client = new TestClient(rootTransport); + + var call = SendAndReceiveAsync(client, 70, CancellationToken.None); + await started.Task.WaitAsync(TimeSpan.FromSeconds(5)); + rootTransport.Dispose(); + release.TrySetResult(true); + + var result = await call.WaitAsync(TimeSpan.FromSeconds(5)); + CollectionAssert.AreEqual(new byte[] { 70, 71 }, result); + client.Dispose(); + await CaptureExceptionAsync( + () => httpClient.GetAsync("http://localhost")); + httpClient.Dispose(); + } + + private static async Task RunInterleavedCallsAsync(bool useBufferedTransport, bool useFramedTransport = false) + { + var handler = new InterleavingHandler(useFramedTransport); + using var httpClient = new HttpClient(handler); + using var rootTransport = new THttpTransport(httpClient, null, new Uri("http://localhost")); + TTransport transport = useFramedTransport ? new TFramedTransport(rootTransport) : rootTransport; + if (useBufferedTransport) + { + transport = new TBufferedTransport(transport); + } + using var client = new TestClient(transport); + using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + + try + { + var firstCall = SendAndReceiveAsync(client, 10, timeout.Token); + await handler.FirstRequestStarted.WaitAsync(timeout.Token); + var secondCall = SendAndReceiveAsync(client, 20, timeout.Token); + + var results = await Task.WhenAll(firstCall, secondCall).WaitAsync(timeout.Token); + CollectionAssert.AreEqual(new byte[] { 10, 11 }, results[0]); + CollectionAssert.AreEqual(new byte[] { 20, 21 }, results[1]); + } + finally + { + handler.ReleaseFirstRequest(); + } + } + + private static Task SendAndReceiveAsync(TestClient client, byte requestByte, + CancellationToken cancellationToken) + { + return client.RunAsync(async () => + { + var request = new[] { requestByte }; + var outputTransport = client.OutputProtocol.Transport; + await outputTransport.WriteAsync(request, 0, request.Length, cancellationToken); + await outputTransport.FlushAsync(cancellationToken); + + var response = new byte[2]; + await client.InputProtocol.Transport.ReadAllAsync(response, 0, response.Length, cancellationToken); + return response; + }, cancellationToken); + } + + private static HttpResponseMessage Response(byte payload) + { + return new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new ByteArrayContent(new[] { payload, (byte)(payload + 1) }) + }; + } + + private static TaskCompletionSource NewSignal() + { + return new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + } + + private static async Task CaptureExceptionAsync(Func operation) + where TException : Exception + { + try + { + await operation(); + } + catch (TException exception) + { + return exception; + } + + throw new InvalidOperationException($"Expected {typeof(TException).Name}."); + } + + private sealed class RequestHandler : HttpMessageHandler + { + private readonly Func> _responseFactory; + + public RequestHandler(Func> responseFactory) + { + _responseFactory = responseFactory; + } + + protected override async Task SendAsync(HttpRequestMessage request, + CancellationToken cancellationToken) + { + var content = request.Content ?? throw new InvalidOperationException("The request content was missing."); + var body = await content.ReadAsByteArrayAsync(cancellationToken); + if (body.Length == 0) + { + throw new InvalidOperationException("The request body was empty."); + } + return await _responseFactory(body[0], cancellationToken); + } + } + + private sealed class TrackingContent : ByteArrayContent + { + public TrackingContent(byte[] content) + : base(content) + { + } + + public bool IsDisposed { get; private set; } + + protected override void Dispose(bool disposing) + { + IsDisposed = true; + base.Dispose(disposing); + } + } + + private sealed class InterleavingHandler : HttpMessageHandler + { + private readonly TaskCompletionSource _firstRequestStarted = + new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + private readonly TaskCompletionSource _releaseFirstRequest = + new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + private readonly bool _useFramedTransport; + private int _requestCount; + + public InterleavingHandler(bool useFramedTransport) + { + _useFramedTransport = useFramedTransport; + } + + public Task FirstRequestStarted => _firstRequestStarted.Task; + + public void ReleaseFirstRequest() + { + _releaseFirstRequest.TrySetResult(true); + } + + protected override async Task SendAsync(HttpRequestMessage request, + CancellationToken cancellationToken) + { + var requestNumber = Interlocked.Increment(ref _requestCount); + if (requestNumber == 1) + { + _firstRequestStarted.TrySetResult(true); + await _releaseFirstRequest.Task.WaitAsync(cancellationToken); + } + + byte[] requestBody; + try + { + var content = request.Content ?? throw new InvalidOperationException("The request content was missing."); + requestBody = await content.ReadAsByteArrayAsync(cancellationToken); + } + finally + { + if (requestNumber == 2) + { + _releaseFirstRequest.TrySetResult(true); + } + } + + var payloadOffset = 0; + if (_useFramedTransport) + { + if (requestBody.Length < sizeof(int)) + { + throw new InvalidOperationException("The framed request was shorter than its header."); + } + + var frameLength = BinaryPrimitives.ReadInt32BigEndian(requestBody.AsSpan(0, sizeof(int))); + if (frameLength != requestBody.Length - sizeof(int)) + { + throw new InvalidOperationException("The framed request length did not match its payload."); + } + + payloadOffset = sizeof(int); + } + + var requestValue = requestBody[payloadOffset]; + var responsePayload = new[] { requestValue, (byte)(requestValue + 1) }; + if (_useFramedTransport) + { + var framedResponse = new byte[sizeof(int) + responsePayload.Length]; + BinaryPrimitives.WriteInt32BigEndian(framedResponse.AsSpan(0, sizeof(int)), responsePayload.Length); + Buffer.BlockCopy(responsePayload, 0, framedResponse, sizeof(int), responsePayload.Length); + responsePayload = framedResponse; + } + + return new HttpResponseMessage(HttpStatusCode.OK) + { + Content = new ByteArrayContent(responsePayload) + }; + } + } + + private sealed class TestClient : TBaseClient, IDisposable + { + private static readonly TProtocolFactory ProtocolFactory = new TBinaryProtocol.Factory(); + + public TestClient(TTransport transport) + : base(transport, ProtocolFactory, ProtocolFactory) + { + } + + public Task RunAsync(Func> operation, + CancellationToken cancellationToken) + { + return ExecutePerCallAsync(operation, cancellationToken); + } + } + } +} \ No newline at end of file diff --git a/lib/netstd/Tests/Thrift.Tests/Transports/THttpTransportTests.cs b/lib/netstd/Tests/Thrift.Tests/Transports/THttpTransportTests.cs index 2a2d884203a..d47021730e2 100644 --- a/lib/netstd/Tests/Thrift.Tests/Transports/THttpTransportTests.cs +++ b/lib/netstd/Tests/Thrift.Tests/Transports/THttpTransportTests.cs @@ -16,6 +16,10 @@ // under the License. using System.Net.Http; +using System.Net.Http.Headers; +using System; +using System.Threading; +using System.Threading.Tasks; using Microsoft.VisualStudio.TestTools.UnitTesting; using Thrift.Transport.Client; @@ -36,5 +40,41 @@ public void THttpTransport_Uses_Configured_ConnectionTimeout_Test() Assert.IsTrue(client.Timeout.TotalMilliseconds == 5000); Assert.IsTrue(httpClientTransport.ConnectTimeout == 5000); } + + [TestMethod] + public void THttpTransport_Preserves_Provided_HttpClient_Timeout() + { + var timeout = TimeSpan.FromSeconds(7); + using var client = new HttpClient { Timeout = timeout }; + using var transport = new THttpTransport(client, null); + + Assert.AreEqual(timeout, client.Timeout); + Assert.AreEqual(timeout, transport.PerCallTimeout); + Assert.AreEqual((int)timeout.TotalMilliseconds, transport.ConnectTimeout); + } + + [TestMethod] + public void THttpTransport_Preserves_Infinite_Provided_HttpClient_Timeout() + { + using var client = new HttpClient { Timeout = Timeout.InfiniteTimeSpan }; + using var transport = new THttpTransport(client, null); + + Assert.AreEqual(Timeout.InfiniteTimeSpan, client.Timeout); + Assert.AreEqual(Timeout.InfiniteTimeSpan, transport.PerCallTimeout); + } + + [TestMethod] + public async Task THttpTransport_PerCallTransport_Preserves_ContentType() + { + using var client = new HttpClient(); + using var transport = new THttpTransport(client, null) + { + ContentType = new MediaTypeHeaderValue("application/custom-thrift") + }; + + using var callTransport = (THttpTransport)await transport.CreatePerCallTransportAsync(default); + + Assert.AreEqual(transport.ContentType, callTransport.ContentType); + } } } diff --git a/lib/netstd/Thrift/TBaseClient.cs b/lib/netstd/Thrift/TBaseClient.cs index aefb54dd5a1..26c23d925e0 100644 --- a/lib/netstd/Thrift/TBaseClient.cs +++ b/lib/netstd/Thrift/TBaseClient.cs @@ -19,6 +19,7 @@ using System.Threading; using System.Threading.Tasks; using Thrift.Protocol; +using Thrift.Transport; namespace Thrift { @@ -32,6 +33,10 @@ public abstract class TBaseClient { private readonly TProtocol _inputProtocol; private readonly TProtocol _outputProtocol; + private readonly TProtocolFactory _inputProtocolFactory; + private readonly TProtocolFactory _outputProtocolFactory; + private readonly ITPerCallTransportProvider _perCallTransportProvider; + private readonly AsyncLocal _perCallProtocols = new AsyncLocal(); private bool _isDisposed; private int _seqId; public readonly Guid ClientId = Guid.NewGuid(); @@ -42,13 +47,217 @@ protected TBaseClient(TProtocol inputProtocol, TProtocol outputProtocol) _outputProtocol = outputProtocol ?? throw new ArgumentNullException(nameof(outputProtocol)); } - public TProtocol InputProtocol => _inputProtocol; + /// + /// Initializes a client with protocols created from the supplied transport and protocol factories. + /// + /// The shared transport, which may provide a new transport for each call. + /// The factory used to create input protocols. + /// The factory used to create output protocols. + /// A required constructor argument is null. + /// A protocol factory returns null. + protected TBaseClient(TTransport transport, TProtocolFactory inputProtocolFactory, + TProtocolFactory outputProtocolFactory) + { + transport = transport ?? throw new ArgumentNullException(nameof(transport)); + _inputProtocolFactory = inputProtocolFactory ?? throw new ArgumentNullException(nameof(inputProtocolFactory)); + _outputProtocolFactory = outputProtocolFactory ?? throw new ArgumentNullException(nameof(outputProtocolFactory)); + _perCallTransportProvider = transport as ITPerCallTransportProvider; + if (_perCallTransportProvider != null && !_perCallTransportProvider.SupportsPerCallTransport) + { + _perCallTransportProvider = null; + } + + _inputProtocol = CreateProtocol(_inputProtocolFactory, transport); + try + { + _outputProtocol = CreateProtocol(_outputProtocolFactory, transport); + } + catch + { + _inputProtocol.Dispose(); + throw; + } + } + + /// + /// Gets the input protocol for the current call, or the shared input protocol when no per-call transport is active. + /// + public TProtocol InputProtocol => _perCallProtocols.Value?.InputProtocol ?? _inputProtocol; - public TProtocol OutputProtocol => _outputProtocol; + /// + /// Gets the output protocol for the current call, or the shared output protocol when no per-call transport is active. + /// + public TProtocol OutputProtocol => _perCallProtocols.Value?.OutputProtocol ?? _outputProtocol; + /// + /// Gets the next client sequence identifier. + /// public int SeqId { - get { return ++_seqId; } + get { return Interlocked.Increment(ref _seqId); } + } + + /// + /// Executes an operation with protocols scoped to a new transport when the shared transport supports per-call transports. + /// + /// The operation to execute. + /// The token used to cancel per-call transport creation and opening. + /// A task that represents the asynchronous operation. + protected Task ExecutePerCallAsync(Func operation, CancellationToken cancellationToken) + { + operation = operation ?? throw new ArgumentNullException(nameof(operation)); + + return ExecutePerCallAsync(new Func(_ => operation()), cancellationToken); + } + + /// + /// Executes an operation using the cancellation token linked to the caller and per-call timeout. + /// + /// The operation to execute. + /// The caller's cancellation token. + /// A task that represents the asynchronous operation. + protected Task ExecutePerCallAsync(Func operation, + CancellationToken cancellationToken) + { + return ExecutePerCallScopeAsync(operation, cancellationToken); + } + + private async Task ExecutePerCallScopeAsync(Func operation, + CancellationToken cancellationToken) + { + operation = operation ?? throw new ArgumentNullException(nameof(operation)); + + if (_perCallTransportProvider == null) + { + await operation(cancellationToken); + return; + } + + using (var callCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken)) + { + var timeout = _perCallTransportProvider.PerCallTimeout; + if (timeout < TimeSpan.Zero && timeout != Timeout.InfiniteTimeSpan) + { + throw new InvalidOperationException("Per-call timeout must be non-negative or infinite."); + } + + if (timeout == TimeSpan.Zero) + { + callCancellation.Cancel(); + } + else if (timeout != Timeout.InfiniteTimeSpan) + { + callCancellation.CancelAfter(timeout); + } + + await ExecutePerCallCoreAsync(operation, callCancellation.Token); + } + } + + private async Task ExecutePerCallCoreAsync(Func operation, + CancellationToken cancellationToken) + { + TTransport transport; + try + { + transport = await _perCallTransportProvider.CreatePerCallTransportAsync(cancellationToken); + } + catch (OperationCanceledException exception) + { + throw new TTransportException(TTransportException.ExceptionType.Interrupted, exception.Message, exception); + } + + transport = transport ?? throw new InvalidOperationException("The per-call transport provider returned null."); + + TProtocol inputProtocol = null; + TProtocol outputProtocol = null; + try + { + inputProtocol = CreateProtocol(_inputProtocolFactory, transport); + outputProtocol = CreateProtocol(_outputProtocolFactory, transport); + var previousProtocols = _perCallProtocols.Value; + _perCallProtocols.Value = new PerCallProtocols(inputProtocol, outputProtocol); + try + { + if (!inputProtocol.Transport.IsOpen) + { + await inputProtocol.Transport.OpenAsync(cancellationToken); + } + if (!ReferenceEquals(inputProtocol.Transport, outputProtocol.Transport) && !outputProtocol.Transport.IsOpen) + { + await outputProtocol.Transport.OpenAsync(cancellationToken); + } + await operation(cancellationToken); + } + catch (OperationCanceledException exception) + { + throw new TTransportException(TTransportException.ExceptionType.Interrupted, exception.Message, exception); + } + finally + { + _perCallProtocols.Value = previousProtocols; + } + } + finally + { + try + { + outputProtocol?.Dispose(); + } + finally + { + try + { + if (!ReferenceEquals(inputProtocol, outputProtocol)) + { + inputProtocol?.Dispose(); + } + } + finally + { + if (inputProtocol == null && outputProtocol == null) + { + transport.Dispose(); + } + } + } + } + } + + /// + /// Executes an operation with protocols scoped to a new transport when the shared transport supports per-call transports. + /// + /// The type of result returned by the operation. + /// The operation to execute. + /// The token used to cancel per-call transport creation and opening. + /// A task containing the operation result. + protected async Task ExecutePerCallAsync(Func> operation, + CancellationToken cancellationToken) + { + operation = operation ?? throw new ArgumentNullException(nameof(operation)); + + return await ExecutePerCallAsync(new Func>( + _ => operation()), cancellationToken); + } + + /// + /// Executes a result-producing operation using the cancellation token linked to the caller and per-call timeout. + /// + /// The type of result returned by the operation. + /// The operation to execute. + /// The caller's cancellation token. + /// A task containing the operation result. + protected async Task ExecutePerCallAsync( + Func> operation, CancellationToken cancellationToken) + { + operation = operation ?? throw new ArgumentNullException(nameof(operation)); + + TResult result = default(TResult); + await ExecutePerCallScopeAsync(new Func(async callCancellationToken => + { + result = await operation(callCancellationToken); + }), cancellationToken); + return result; } public virtual async Task OpenTransportAsync() @@ -78,14 +287,43 @@ protected virtual void Dispose(bool disposing) { if (!_isDisposed) { - if (disposing) + try { - _inputProtocol?.Dispose(); - _outputProtocol?.Dispose(); + if (disposing) + { + try + { + _inputProtocol?.Dispose(); + } + finally + { + _outputProtocol?.Dispose(); + } + } } + finally + { + _isDisposed = true; + } + } + } + + private static TProtocol CreateProtocol(TProtocolFactory factory, TTransport transport) + { + return factory.GetProtocol(transport) ?? + throw new InvalidOperationException("A protocol factory returned null."); + } + + private sealed class PerCallProtocols + { + public PerCallProtocols(TProtocol inputProtocol, TProtocol outputProtocol) + { + InputProtocol = inputProtocol; + OutputProtocol = outputProtocol; } - _isDisposed = true; + public TProtocol InputProtocol { get; } + public TProtocol OutputProtocol { get; } } } } diff --git a/lib/netstd/Thrift/Transport/Client/THttpTransport.cs b/lib/netstd/Thrift/Transport/Client/THttpTransport.cs index 439777eb390..7a36999f855 100644 --- a/lib/netstd/Thrift/Transport/Client/THttpTransport.cs +++ b/lib/netstd/Thrift/Transport/Client/THttpTransport.cs @@ -32,16 +32,19 @@ namespace Thrift.Transport.Client { // ReSharper disable once InconsistentNaming - public class THttpTransport : TEndpointTransport + public class THttpTransport : TEndpointTransport, ITPerCallTransportProvider { private readonly X509Certificate[] _certificates; private readonly Uri _uri; private int _connectTimeout = 30000; // Timeouts in milliseconds private HttpClient _httpClient; + private readonly HttpClientLease _httpClientLease; private Stream _inputStream; private MemoryStream _outputStream = new MemoryStream(); + private HttpResponseMessage _responseMessage; private bool _isDisposed; + private int _leaseReleased; public THttpTransport(Uri uri, TConfiguration config, IDictionary customRequestHeaders = null, string userAgent = null) : this(uri, config, Enumerable.Empty(), customRequestHeaders, userAgent) @@ -62,6 +65,7 @@ public THttpTransport(Uri uri, TConfiguration config, IEnumerableuse->dispose per flush) later _httpClient = CreateClient(customRequestHeaders); ConfigureClient(_httpClient); + _httpClientLease = new HttpClientLease(_httpClient); } /// @@ -87,7 +91,18 @@ public THttpTransport(HttpClient httpClient, TConfiguration config, Uri uri = nu if (!string.IsNullOrEmpty(userAgent)) UserAgent = userAgent; - ConfigureClient(_httpClient); + ConfigureClient(_httpClient, configureTimeout: false); + _httpClientLease = new HttpClientLease(_httpClient); + } + + private THttpTransport(HttpClientLease httpClientLease, TConfiguration config, Uri uri) + : base(config) + { + _httpClientLease = httpClientLease; + _httpClient = httpClientLease.Client; + _uri = uri; + _connectTimeout = (int)_httpClient.Timeout.TotalMilliseconds; + UserAgent = _httpClient.DefaultRequestHeaders.UserAgent.ToString(); } // According to RFC 2616 section 3.8, the "User-Agent" header may not carry a version number @@ -115,33 +130,55 @@ public int ConnectTimeout public MediaTypeHeaderValue ContentType { get; set; } - public override Task OpenAsync(CancellationToken cancellationToken) - { - cancellationToken.ThrowIfCancellationRequested(); - return Task.CompletedTask; - } + /// + /// Gets the configured timeout for a complete client call using this transport. + /// + public TimeSpan PerCallTimeout => _httpClient?.Timeout ?? TimeSpan.FromMilliseconds(_connectTimeout); - public override void Close() + /// + /// Gets whether this transport can create an independent transport for each call. + /// + public bool SupportsPerCallTransport => true; + + /// + /// Creates a transport with independent request and response buffers for one client call. + /// + /// Token used to cancel transport creation. + /// A new transport owned by the caller. + public Task CreatePerCallTransportAsync(CancellationToken cancellationToken) { - if (_inputStream != null) + cancellationToken.ThrowIfCancellationRequested(); + if (_isDisposed || _httpClient == null) { - _inputStream.Dispose(); - _inputStream = null; + throw new ObjectDisposedException(nameof(THttpTransport)); } - if (_outputStream != null) + _httpClientLease.AddReference(); + try { - _outputStream.Dispose(); - _outputStream = null; + return Task.FromResult(new THttpTransport(_httpClientLease, Configuration, _uri) + { + ContentType = ContentType + }); } - - if (_httpClient != null) + catch { - _httpClient.Dispose(); - _httpClient = null; + _httpClientLease.Release(); + throw; } } + public override Task OpenAsync(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + return Task.CompletedTask; + } + + public override void Close() + { + DisposeResources(); + } + public override async ValueTask ReadAsync(byte[] buffer, int offset, int length, CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); @@ -218,9 +255,9 @@ private HttpClient CreateClient(IDictionary customRequestHeaders return httpClient; } - private void ConfigureClient(HttpClient httpClient) + private void ConfigureClient(HttpClient httpClient, bool configureTimeout = true) { - if (_connectTimeout > 0) + if (configureTimeout && _connectTimeout > 0) { httpClient.Timeout = TimeSpan.FromMilliseconds(_connectTimeout); } @@ -245,13 +282,15 @@ public override async Task FlushAsync(CancellationToken cancellationToken) { contentStream.Headers.ContentType = ContentType ?? new MediaTypeHeaderValue(@"application/x-thrift"); - var response = (await _httpClient.PostAsync(_uri, contentStream, cancellationToken)).EnsureSuccessStatusCode(); - + var response = await _httpClient.PostAsync(_uri, contentStream, cancellationToken); _inputStream?.Dispose(); + _responseMessage?.Dispose(); + _responseMessage = response; + response.EnsureSuccessStatusCode(); #if NET5_0_OR_GREATER - _inputStream = await response.Content.ReadAsStreamAsync(cancellationToken); + _inputStream = await _responseMessage.Content.ReadAsStreamAsync(cancellationToken); #else - _inputStream = await response.Content.ReadAsStreamAsync(); + _inputStream = await _responseMessage.Content.ReadAsStreamAsync(); #endif if (_inputStream.CanSeek) { @@ -288,15 +327,95 @@ public override async Task FlushAsync(CancellationToken cancellationToken) protected override void Dispose(bool disposing) { if (!_isDisposed) + DisposeResources(); + } + + private void DisposeResources() + { + if (_isDisposed) { - if (disposing) + return; + } + + _isDisposed = true; + try + { + _inputStream?.Dispose(); + } + finally + { + try { - _inputStream?.Dispose(); _outputStream?.Dispose(); - _httpClient?.Dispose(); + } + finally + { + try + { + _responseMessage?.Dispose(); + } + finally + { + _inputStream = null; + _outputStream = null; + _responseMessage = null; + _httpClient = null; + if (Interlocked.Exchange(ref _leaseReleased, 1) == 0) + { + _httpClientLease.Release(); + } + } + } + } + } + + private sealed class HttpClientLease + { + private readonly object _gate = new object(); + private int _referenceCount = 1; + private bool _isDisposed; + + public HttpClientLease(HttpClient client) + { + Client = client ?? throw new ArgumentNullException(nameof(client)); + } + + public HttpClient Client { get; } + + public void AddReference() + { + lock (_gate) + { + if (_isDisposed) + { + throw new ObjectDisposedException(nameof(HttpClientLease)); + } + _referenceCount++; + } + } + + public void Release() + { + var disposeClient = false; + lock (_gate) + { + if (_referenceCount == 0) + { + return; + } + _referenceCount--; + if (_referenceCount == 0) + { + _isDisposed = true; + disposeClient = true; + } + } + + if (disposeClient) + { + Client.Dispose(); } } - _isDisposed = true; } } } diff --git a/lib/netstd/Thrift/Transport/ITPerCallTransportProvider.cs b/lib/netstd/Thrift/Transport/ITPerCallTransportProvider.cs new file mode 100644 index 00000000000..3bfa2893e95 --- /dev/null +++ b/lib/netstd/Thrift/Transport/ITPerCallTransportProvider.cs @@ -0,0 +1,47 @@ +// Licensed to the Apache Software Foundation(ASF) under one +// or more contributor license agreements.See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership.The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +using System.Threading; +using System.Threading.Tasks; +using System; + +namespace Thrift.Transport +{ + /// + /// Provides a mechanism to create a separate transport for each client call. + /// The caller owns each returned transport and must dispose it when the call completes. + /// + public interface ITPerCallTransportProvider + { + /// + /// Gets whether this transport can create an independent transport for each call. + /// + bool SupportsPerCallTransport { get; } + + /// + /// Gets the maximum duration for a complete per-call operation, or if unbounded. + /// + TimeSpan PerCallTimeout { get; } + + /// + /// Creates a new transport for the current client call. + /// + /// Token used to cancel transport creation. + /// A new transport owned by the caller. + Task CreatePerCallTransportAsync(CancellationToken cancellationToken); + } +} \ No newline at end of file diff --git a/lib/netstd/Thrift/Transport/Layered/TBufferedTransport.cs b/lib/netstd/Thrift/Transport/Layered/TBufferedTransport.cs index 977dcbfd119..80c32988278 100644 --- a/lib/netstd/Thrift/Transport/Layered/TBufferedTransport.cs +++ b/lib/netstd/Thrift/Transport/Layered/TBufferedTransport.cs @@ -24,7 +24,7 @@ namespace Thrift.Transport { // ReSharper disable once InconsistentNaming - public class TBufferedTransport : TLayeredTransport + public class TBufferedTransport : TLayeredTransport, ITPerCallTransportProvider { private readonly int DesiredBufferSize; private readonly Client.TMemoryBufferTransport ReadBuffer; @@ -67,6 +67,52 @@ public TTransport UnderlyingTransport } } + /// + /// Gets the maximum duration of a complete operation delegated to the underlying transport. + /// + public TimeSpan PerCallTimeout + { + get + { + if (!(InnerTransport is ITPerCallTransportProvider provider)) + { + throw new System.NotSupportedException("The underlying transport does not support per-call transports."); + } + return provider.PerCallTimeout; + } + } + + /// + /// Gets whether the underlying transport supports per-call transports. + /// + public bool SupportsPerCallTransport => + InnerTransport is ITPerCallTransportProvider provider && provider.SupportsPerCallTransport; + + /// + /// Creates a buffered transport around a new per-call underlying transport. + /// + /// Token used to cancel underlying transport creation. + /// A new buffered transport owned by the caller. + /// The underlying transport does not support per-call transports. + public async Task CreatePerCallTransportAsync(CancellationToken cancellationToken) + { + if (!(InnerTransport is ITPerCallTransportProvider provider) || !provider.SupportsPerCallTransport) + { + throw new System.NotSupportedException("The underlying transport does not support per-call transports."); + } + + var transport = await provider.CreatePerCallTransportAsync(cancellationToken); + try + { + return new TBufferedTransport(transport, DesiredBufferSize); + } + catch + { + transport?.Dispose(); + throw; + } + } + public override bool IsOpen => !IsDisposed && InnerTransport.IsOpen; public override async Task OpenAsync(CancellationToken cancellationToken) diff --git a/lib/netstd/Thrift/Transport/Layered/TFramedTransport.cs b/lib/netstd/Thrift/Transport/Layered/TFramedTransport.cs index faa3fa6e90f..83c43766be7 100644 --- a/lib/netstd/Thrift/Transport/Layered/TFramedTransport.cs +++ b/lib/netstd/Thrift/Transport/Layered/TFramedTransport.cs @@ -28,7 +28,7 @@ namespace Thrift.Transport { // ReSharper disable once InconsistentNaming - public class TFramedTransport : TLayeredTransport + public class TFramedTransport : TLayeredTransport, ITPerCallTransportProvider { private const int HeaderSize = 4; private readonly byte[] HeaderBuf = new byte[HeaderSize]; @@ -53,6 +53,52 @@ public TFramedTransport(TTransport transport) InitWriteBuffer(); } + /// + /// Gets the maximum duration of a complete operation delegated to the underlying transport. + /// + public TimeSpan PerCallTimeout + { + get + { + if (!(InnerTransport is ITPerCallTransportProvider provider)) + { + throw new NotSupportedException("The underlying transport does not support per-call transports."); + } + return provider.PerCallTimeout; + } + } + + /// + /// Gets whether the underlying transport supports per-call transports. + /// + public bool SupportsPerCallTransport => + InnerTransport is ITPerCallTransportProvider provider && provider.SupportsPerCallTransport; + + /// + /// Creates a framed transport with independent frame buffers around a new per-call underlying transport. + /// + /// Token used to cancel underlying transport creation. + /// A new framed transport owned by the caller. + /// The underlying transport does not support per-call transports. + public async Task CreatePerCallTransportAsync(CancellationToken cancellationToken) + { + if (!(InnerTransport is ITPerCallTransportProvider provider) || !provider.SupportsPerCallTransport) + { + throw new NotSupportedException("The underlying transport does not support per-call transports."); + } + + var transport = await provider.CreatePerCallTransportAsync(cancellationToken); + try + { + return new TFramedTransport(transport); + } + catch + { + transport?.Dispose(); + throw; + } + } + public override bool IsOpen => !IsDisposed && InnerTransport.IsOpen; public override async Task OpenAsync(CancellationToken cancellationToken) diff --git a/lib/netstd/openspec/changes/archive/.gitkeep b/lib/netstd/openspec/changes/archive/.gitkeep new file mode 100644 index 00000000000..e69de29bb2d diff --git a/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/.openspec.yaml b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/.openspec.yaml new file mode 100644 index 00000000000..aa32566bf64 --- /dev/null +++ b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-10-06 diff --git a/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/design.md b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/design.md new file mode 100644 index 00000000000..c59c2f06d81 --- /dev/null +++ b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/design.md @@ -0,0 +1,39 @@ +# Design + +## Context + +`THttpTransport` currently stores its output and input streams on the transport instance. `FlushAsync` posts the shared output stream, replaces it during cleanup, and `Dispose` closes both streams and the client. `TFramedTransport` also owns read/write frame buffers. Generated service methods execute send, flush, and receive in one asynchronous method, while generated clients currently retain fixed input/output protocols. See proposal.md for motivation and spec.md for the required behavior. + +## Goals / Non-Goals + +**Goals:** +- Add an explicit construction path for generated clients to create a per-call transport and protocols. +- Keep each generated high-level call's transport, buffers, response, and protocol state alive from serialization through response consumption. +- Let transport wrappers explicitly delegate the capability and preserve the current constructors as shared-mode behavior. +- Apply one caller cancellation/timeout scope across the whole call and clean up only that call's resources. + +**Non-Goals:** +- Change the wire protocol or make existing protocol-based constructors opt into per-call behavior. +- Make a manually split sequence of generated `send_*` and `recv_*` methods a per-call operation; those low-level methods retain their existing shared-protocol behavior. +- Make wrappers that do not opt into per-call support appear capable automatically. + +## Decisions + +- **Use explicit transport capability and additive constructors.** Add a per-call provider contract for a transport to create a call-scoped transport. Add a `TBaseClient` constructor path that accepts the provider and protocol factories, and emit corresponding constructors for generated clients. Preserve the existing protocol constructors and document that they remain the shared-mode path. This avoids changing behavior for existing callers; silently changing every client to per-call operation was rejected because it would alter resource and concurrency semantics. +- **Scope the full generated call in `TBaseClient`.** The generated convenience method will execute its existing send, flush, and receive sequence through one base-client per-call helper. The helper passes its linked caller/deadline token into the generated callback so serialization, flush, and response receive share one deadline. Allocate call identifiers atomically so concurrent calls through one client cannot collide. Keep direct low-level `send_*` and `recv_*` APIs on the legacy protocol path, as they do not represent one complete call lifetime. +- **Isolate mutable HTTP state, not the connection pool.** Each HTTP call owns its request body, response stream, and any per-call wrapper buffers. Reuse the thread-safe `HttpClient` owned by the root `THttpTransport`; a call-scoped transport must not dispose that shared client. Preserve the timeout of an `HttpClient` supplied by the caller, while internally created clients retain the transport's default timeout. Coordinate root disposal with active call leases so the client is not closed while a call still uses it. Creating and disposing a separate `HttpClient` per call was rejected because it discards connection pooling and complicates ownership. +- **Delegate wrappers explicitly.** `TBufferedTransport` and `TFramedTransport` each implement the provider and recreate themselves around a per-call underlying transport, preserving their configured buffer/frame behavior and delegating the timeout. `TLayeredTransport` does not implement the provider: it cannot recreate arbitrary derived wrapper types or their configuration, and would falsely advertise support for external subclasses. This follows the existing transport-composition model rather than adding special handling in generated clients for each wrapper type. +- **Bound the whole call with one linked cancellation scope.** Link the caller token with the configured transport timeout before acquiring/using call resources, and pass the linked token through protocol I/O and transport operations. Translate timeout/cancellation using the existing interrupted transport exception behavior; canceling one call must not cancel another. +- **Keep resource cleanup local and idempotent.** Dispose call-owned protocols, streams, content, and wrapper transports on success, fault, or cancellation; do not dispose root-owned shared resources from a call. Ensure each lease is released once, including failures during provider creation or protocol construction. + +## Risks / Trade-offs + +- [A custom wrapper may contain state that cannot be recreated per call] → Require explicit provider delegation and test built-in wrappers; retain shared mode for unsupported wrappers. +- [Generated high-level methods and direct `send_*`/`recv_*` calls have different concurrency guarantees] → Document the boundary on the new and existing generated APIs and exercise the supported high-level workflow in the tutorial. +- [Root transport disposal can race with active operations] → Track active call leases and defer disposal of root-owned HTTP resources until outstanding calls release them. +- [Timing-based concurrency tests may be flaky] → Use a controllable HTTP handler/server and synchronization gates to force request/response interleavings and timeout paths. +- [A layered wrapper may fail to delegate a newly supported inner wrapper] → Test both direct framed delegation and buffered-over-framed composition through the tutorial smoke path. + +## Migration Plan + +The change is additive. Existing clients and constructors continue to use shared transport semantics. Users opt into per-call behavior by constructing a generated client with a supporting transport and protocol factories. Rollback removes the opt-in capability and generated constructor path while retaining existing constructors, as described in proposal.md; no wire-format or default-mode migration is needed. \ No newline at end of file diff --git a/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/proposal.md b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/proposal.md new file mode 100644 index 00000000000..26555a39ec7 --- /dev/null +++ b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/proposal.md @@ -0,0 +1,28 @@ +# Proposal + +## Why + +Concurrent or interleaved asynchronous calls through a shared `THttpTransport` can reset or dispose the class-level output buffer while another call is still sending it, producing failures such as the closed-stream exception reported by THRIFT-5830. Callers need an opt-in way to isolate transport state per call without changing the established shared-buffer behavior of existing clients. + +## What Changes + +- Add per-call transport semantics that isolate each call's buffers and protocol state, with `THttpTransport` as the initial built-in implementation and wrapping transports able to support or delegate the capability. +- Keep existing constructors and shared-buffer semantics as the compatibility default; add the generated-client construction support needed to opt into per-call behavior. +- Define safe completion, exception, cancellation, timeout, and disposal behavior for per-call resources. +- Preserve a caller-configured `HttpClient.Timeout` and propagate the linked per-call deadline through generated request send and response receive operations. +- Add focused transport/client tests, a working netstd tutorial example and test coverage, and usage documentation in the netstd README. +- Document new public/protected API and generator-emitted API, including when existing constructors remain appropriate. + +## Capabilities + +### New Capabilities +- `per-call-transport`: Opt-in per-call isolation for asynchronous client transport operations, including safe lifecycle and compatibility with shared transport behavior. + +### Modified Capabilities + +## Impact + +- Affected implementation: `Thrift/Transport/Client/THttpTransport.cs`, `Thrift/Transport/Layered/TFramedTransport.cs`, `Thrift/TBaseClient.cs`, and generated client constructors in `compiler/cpp/src/thrift/generate/t_netstd_generator.cc`. +- Affected validation and docs: netstd transport/client tests, `tutorial/netstd` client and tutorial checks, and `lib/netstd/README.md`. +- Affected teams: netstd library maintainers, compiler/generator maintainers, and tutorial/test maintainers. +- Rollback: remove the opt-in per-call API and its generated constructor path while retaining the existing shared-mode constructors and behavior; revert the accompanying tutorial and documentation changes. No default behavior or wire format is intended to change. \ No newline at end of file diff --git a/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/specs/per-call-transport/spec.md b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/specs/per-call-transport/spec.md new file mode 100644 index 00000000000..4f6270564a8 --- /dev/null +++ b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/specs/per-call-transport/spec.md @@ -0,0 +1,89 @@ +# Spec Delta + +## Purpose + +Allows generated asynchronous clients to opt into independent transport state for each call, so overlapping calls do not interfere with request or response data while existing shared-transport clients remain compatible. + +## ADDED Requirements + +### Requirement: Per-call transport is opt-in +The system MUST allow clients to use per-call transport semantics when the supplied transport supports or delegates that capability, while preserving existing shared-transport behavior for clients that do not opt in. + +#### Scenario: Existing client retains shared transport behavior +- **GIVEN** a client is created with the existing protocol-based constructors +- **WHEN** the client performs a call over a transport without per-call semantics enabled +- **THEN** the call uses the existing shared transport and protocol behavior + +#### Scenario: Client opts into per-call transport behavior +- **GIVEN** a client is constructed with a transport that supports per-call semantics +- **WHEN** the client starts an asynchronous call +- **THEN** that call uses transport and protocol state scoped to the call + +#### Scenario: Wrapper delegates per-call capability +- **GIVEN** a transport wrapper delegates per-call operations to an underlying transport that supports them +- **WHEN** a client uses the wrapper for an asynchronous call +- **THEN** the call receives the same isolation guarantees as when using the underlying capability directly + +#### Scenario: Framed wrapper delegates per-call capability +- **GIVEN** a `TFramedTransport` wraps a transport that supports per-call operations +- **WHEN** a client uses the framed wrapper for an asynchronous call +- **THEN** the call uses a new framed wrapper with independent frame buffers around its own per-call underlying transport + +#### Scenario: Layered wrappers compose per-call capabilities +- **GIVEN** buffered and framed wrappers are layered over a transport that supports per-call operations +- **WHEN** a client performs an asynchronous call through the layered transport +- **THEN** each wrapper layer is recreated around the call's underlying transport and delegates its timeout + +### Requirement: Concurrent calls keep request and response state isolated +The system MUST prevent one per-call operation from modifying, replacing, or disposing another operation's request buffers, response buffers, or protocol state, and each reply MUST remain associated with its originating call. + +#### Scenario: Interleaved HTTP calls complete independently +- **GIVEN** two asynchronous calls use the same per-call-enabled client through an HTTP transport +- **WHEN** their request sends and response reads overlap +- **THEN** each call sends its complete request and receives the response associated with that request without a closed-stream failure or cross-call data + +#### Scenario: Per-call read-ahead does not consume another response +- **GIVEN** concurrent calls use a transport or wrapper with read-ahead buffering +- **WHEN** either call reads ahead while the other call is active +- **THEN** buffered bytes remain associated with their originating call + +### Requirement: Per-call resources are cleaned up on every terminal path +The system MUST keep resources needed by a call alive until that call has finished consuming its response, then release resources owned by that call exactly once whether the call succeeds, fails, or is interrupted. + +#### Scenario: Successful call releases its owned resources +- **GIVEN** a per-call operation completes its request and response successfully +- **WHEN** the call finishes +- **THEN** its call-owned buffers and resources are released without disposing resources still in use by another call + +#### Scenario: Failed call releases resources and leaves other calls usable +- **GIVEN** one per-call operation fails while another call is active +- **WHEN** failure cleanup runs +- **THEN** only resources owned by the failed call are released and the other call can continue + +#### Scenario: Interrupted call releases resources +- **GIVEN** a per-call operation is canceled or reaches its configured timeout +- **WHEN** the operation terminates +- **THEN** its call-owned resources are released and no resource remains disposed while another call still uses it + +### Requirement: Cancellation and timeouts cover the per-call operation +The system MUST propagate cancellation and enforce the configured timeout across the asynchronous per-call transport operation, and termination of one call MUST NOT cancel unrelated calls. + +#### Scenario: Caller cancellation reaches an in-flight request +- **GIVEN** a per-call operation is sending or consuming a response +- **WHEN** its caller cancels the operation +- **THEN** the operation terminates with the transport's established interrupted-operation behavior and releases its resources + +#### Scenario: Configured timeout bounds an in-flight call +- **GIVEN** a per-call operation exceeds the configured transport timeout while sending or consuming its response +- **WHEN** the timeout expires +- **THEN** that operation terminates with the transport's established interrupted-operation behavior and does not affect other calls + +#### Scenario: Caller-provided HTTP timeout is preserved +- **GIVEN** an `HttpClient` has a caller-configured timeout before it is supplied to `THttpTransport` +- **WHEN** a generated client performs a per-call operation through that transport +- **THEN** the supplied timeout value remains unchanged and bounds the complete operation + +#### Scenario: Per-call deadline reaches the generated RPC +- **GIVEN** a generated high-level RPC runs through a per-call-enabled transport +- **WHEN** its linked caller/deadline token is canceled during request serialization, flush, or response receive +- **THEN** the active phase observes that token and the RPC terminates with the interrupted-operation behavior \ No newline at end of file diff --git a/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/tasks.md b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/tasks.md new file mode 100644 index 00000000000..44ef9cb89b9 --- /dev/null +++ b/lib/netstd/openspec/changes/archive/2026-10-06-per-call-transport/tasks.md @@ -0,0 +1,42 @@ +# Tasks + +## 1. Opt-In API and Compatibility + +- [x] 1.1 Add failing client tests for the `Existing client retains shared transport behavior`, `Client opts into per-call transport behavior`, and `Wrapper delegates per-call capability` scenarios; verify they fail because no per-call construction path exists. +- [x] 1.2 Add the per-call provider contract and `TBaseClient` call context/protocol-factory construction path, then emit matching constructors from `t_netstd_generator.cc`; verify the focused client tests and generated-client compile tests pass. +- [x] 1.3 Refactor constructor and provider ownership as needed, and add XML documentation to new public/protected APIs and generator output explaining when the existing shared-protocol constructors remain appropriate; verify generated API documentation and compile tests. + +## 2. Concurrent Call Isolation + +- [x] 2.1 Add deterministic failing tests for `Interleaved HTTP calls complete independently` and `Per-call read-ahead does not consume another response`, using one client and controlled request/response interleaving; verify the tests reproduce shared-buffer or cross-call state interference. +- [x] 2.2 Scope generated high-level send/flush/receive operations to one call context, isolate HTTP request/response streams and wrapper buffers, and allocate call identifiers without races; verify both concurrency scenarios pass, including response-to-request association. +- [x] 2.3 Refactor call-context and wrapper delegation paths without changing shared-mode behavior; rerun the focused concurrency and existing transport tests. + +## 3. Resource Lifetime, Cancellation, and Timeouts + +- [x] 3.1 Add failing lifecycle tests for `Successful call releases its owned resources`, `Failed call releases resources and leaves other calls usable`, `Interrupted call releases resources`, `Caller cancellation reaches an in-flight request`, and `Configured timeout bounds an in-flight call`; verify each exercises its respective terminal path. +- [x] 3.2 Implement exception-safe per-call cleanup, root HTTP-client lease coordination, and a linked caller-cancellation/transport-timeout scope across send and response consumption; verify the lifecycle and timeout tests pass and cancellation of one call leaves another active. +- [x] 3.3 Refactor ownership and timeout handling to release each call-owned resource once on success, fault, and interruption; rerun all focused lifecycle and cancellation tests. + +## 4. Tutorial and Usage Documentation + +- [x] 4.1 Add an automated tutorial smoke assertion for the per-call client workflow; verify it fails before the tutorial exposes per-call mode. +- [x] 4.2 Add a per-call usage example to `tutorial/netstd` and document construction, opt-in behavior, and the shared-mode constructor choice in `lib/netstd/README.md`; verify the example builds and the documented invocation runs. +- [x] 4.3 Refactor the example and documentation for consistency with generated API comments, including the boundary for direct `send_*`/`recv_*` use; verify the tutorial smoke test passes. + +## 5. Cross-Target Integration + +- [x] 5.1 Build and run the focused transport/client test suite for `netstandard2.0`, `netstandard2.1`, `net8.0`, `net9.0`, and `net10.0`; verify the new API and existing constructors compile on every supported target. +- [x] 5.2 Run the netstd tutorial smoke workflow and the generator's relevant compile tests together; verify the tutorial works end-to-end and generated clients preserve shared-mode compatibility while supporting per-call calls. + +## 6. Framed Wrapper Delegation + +- [x] 6.1 Add failing tests for direct framed-provider delegation and buffered-over-framed composition; verify they fail before framed transports advertise per-call support. +- [x] 6.2 Implement per-call provider and timeout delegation on `TFramedTransport`, recreating its frame buffers around a new inner transport; verify the focused wrapper tests and framed HTTP handler test pass. +- [x] 6.3 Confirm `TLayeredTransport` does not advertise the provider to unsupported subclasses, update wrapper documentation as needed, and rerun the focused transport and tutorial checks. + +## 7. Whole-Call Timeout Propagation + +- [x] 7.1 Add failing tests for caller-provided `HttpClient.Timeout` preservation and linked-token delivery to a generated RPC; verify both regressions before implementation. +- [x] 7.2 Preserve supplied HTTP-client timeout values and pass the combined caller/deadline token into generated send and receive phases; verify the focused timeout tests pass. +- [x] 7.3 Refactor timeout validation and overload compatibility, update the main capability spec, and verify generated-client builds plus all focused transport tests. \ No newline at end of file diff --git a/lib/netstd/openspec/config.yaml b/lib/netstd/openspec/config.yaml new file mode 100644 index 00000000000..5a2b1cd81a2 --- /dev/null +++ b/lib/netstd/openspec/config.yaml @@ -0,0 +1,33 @@ +schema: spec-driven +githubCopilot: + cloudAgent: true +context: | + Tech stack: + - netstandard2.1 + - netstandard2.0 + - net8.0 + - net9.0 + - net10.0 + - C# + TDD discipline (non-negotiable): + - Follow Red-Green-Refactor strictly. + - For every behavior in the specs, write failing tests FIRST, then implement. + Task conventions: + - Group tasks under `## N. Group Name` headings. + - Every task must be a checkbox in the form `- [ ] N.M Task description`. + - Each task must be small, focused, reference the relevant spec scenario, and sequence: tests → implementation → refactor. + - Check off tasks as you complete them one-by-one or a section at a time. + +rules: + proposal: + - Include rollback plan + - Identify affected teams + specs: + - Use Given/When/Then format + - Reference existing patterns before inventing new ones + +operations: + apply: + guidance: + - Run focused tests before the full suite + - Prefer small, focused diffs; keep each implementation change under ~30 lines when possible diff --git a/lib/netstd/openspec/specs/.gitkeep b/lib/netstd/openspec/specs/.gitkeep new file mode 100644 index 00000000000..e69de29bb2d diff --git a/lib/netstd/openspec/specs/per-call-transport/spec.md b/lib/netstd/openspec/specs/per-call-transport/spec.md new file mode 100644 index 00000000000..713b1d92345 --- /dev/null +++ b/lib/netstd/openspec/specs/per-call-transport/spec.md @@ -0,0 +1,89 @@ +# per-call-transport Specification + +## Purpose + +Allows generated asynchronous clients to opt into independent transport state for each call, so overlapping calls do not interfere with request or response data while existing shared-transport clients remain compatible. + +## Requirements + +### Requirement: Per-call transport is opt-in +The system MUST allow clients to use per-call transport semantics when the supplied transport supports or delegates that capability, while preserving existing shared-transport behavior for clients that do not opt in. + +#### Scenario: Existing client retains shared transport behavior +- **GIVEN** a client is created with the existing protocol-based constructors +- **WHEN** the client performs a call over a transport without per-call semantics enabled +- **THEN** the call uses the existing shared transport and protocol behavior + +#### Scenario: Client opts into per-call transport behavior +- **GIVEN** a client is constructed with a transport that supports per-call semantics +- **WHEN** the client starts an asynchronous call +- **THEN** that call uses transport and protocol state scoped to the call + +#### Scenario: Wrapper delegates per-call capability +- **GIVEN** a transport wrapper delegates per-call operations to an underlying transport that supports them +- **WHEN** a client uses the wrapper for an asynchronous call +- **THEN** the call receives the same isolation guarantees as when using the underlying capability directly + +#### Scenario: Framed wrapper delegates per-call capability +- **GIVEN** a `TFramedTransport` wraps a transport that supports per-call operations +- **WHEN** a client uses the framed wrapper for an asynchronous call +- **THEN** the call uses a new framed wrapper with independent frame buffers around its own per-call underlying transport + +#### Scenario: Layered wrappers compose per-call capabilities +- **GIVEN** buffered and framed wrappers are layered over a transport that supports per-call operations +- **WHEN** a client performs an asynchronous call through the layered transport +- **THEN** each wrapper layer is recreated around the call's underlying transport and delegates its timeout + +### Requirement: Concurrent calls keep request and response state isolated +The system MUST prevent one per-call operation from modifying, replacing, or disposing another operation's request buffers, response buffers, or protocol state, and each reply MUST remain associated with its originating call. + +#### Scenario: Interleaved HTTP calls complete independently +- **GIVEN** two asynchronous calls use the same per-call-enabled client through an HTTP transport +- **WHEN** their request sends and response reads overlap +- **THEN** each call sends its complete request and receives the response associated with that request without a closed-stream failure or cross-call data + +#### Scenario: Per-call read-ahead does not consume another response +- **GIVEN** concurrent calls use a transport or wrapper with read-ahead buffering +- **WHEN** either call reads ahead while the other call is active +- **THEN** buffered bytes remain associated with their originating call + +### Requirement: Per-call resources are cleaned up on every terminal path +The system MUST keep resources needed by a call alive until that call has finished consuming its response, then release resources owned by that call exactly once whether the call succeeds, fails, or is interrupted. + +#### Scenario: Successful call releases its owned resources +- **GIVEN** a per-call operation completes its request and response successfully +- **WHEN** the call finishes +- **THEN** its call-owned buffers and resources are released without disposing resources still in use by another call + +#### Scenario: Failed call releases resources and leaves other calls usable +- **GIVEN** one per-call operation fails while another call is active +- **WHEN** failure cleanup runs +- **THEN** only resources owned by the failed call are released and the other call can continue + +#### Scenario: Interrupted call releases resources +- **GIVEN** a per-call operation is canceled or reaches its configured timeout +- **WHEN** the operation terminates +- **THEN** its call-owned resources are released and no resource remains disposed while another call still uses it + +### Requirement: Cancellation and timeouts cover the per-call operation +The system MUST propagate cancellation and enforce the configured timeout across the asynchronous per-call transport operation, and termination of one call MUST NOT cancel unrelated calls. + +#### Scenario: Caller cancellation reaches an in-flight request +- **GIVEN** a per-call operation is sending or consuming a response +- **WHEN** its caller cancels the operation +- **THEN** the operation terminates with the transport's established interrupted-operation behavior and releases its resources + +#### Scenario: Configured timeout bounds an in-flight call +- **GIVEN** a per-call operation exceeds the configured transport timeout while sending or consuming its response +- **WHEN** the timeout expires +- **THEN** that operation terminates with the transport's established interrupted-operation behavior and does not affect other calls + +#### Scenario: Caller-provided HTTP timeout is preserved +- **GIVEN** an `HttpClient` has a caller-configured timeout before it is supplied to `THttpTransport` +- **WHEN** a generated client performs a per-call operation through that transport +- **THEN** the supplied timeout value remains unchanged and bounds the complete operation + +#### Scenario: Per-call deadline reaches the generated RPC +- **GIVEN** a generated high-level RPC runs through a per-call-enabled transport +- **WHEN** its linked caller/deadline token is canceled during request serialization, flush, or response receive +- **THEN** the active phase observes that token and the RPC terminates with the interrupted-operation behavior \ No newline at end of file diff --git a/lib/netstd/per-call-transport.md b/lib/netstd/per-call-transport.md new file mode 100644 index 00000000000..91cbadaf2ce --- /dev/null +++ b/lib/netstd/per-call-transport.md @@ -0,0 +1,30 @@ +# Per-Call Transport + +Create an extension to existing transports so that per-call semantics are enabled such that +asynchronous and interleaved calls are safe and can implement a read-ahead buffering protocol. + +## TDD Discipline (non-negotiable) +- Follow Red–Green–Refactor. +- For every behavior in the specs, write failing tests FIRST, then implement. +- Task lists must explicitly sequence: tests → implementation → refactor. + +## Features + +### per-call-transport + +- add a new mechanism for per-call transport to support the existing THttpTransport and any wrapping transports that can support or delegate the per-call semantics. +- reference the observed bug from https://issues.apache.org/jira/browse/THRIFT-5830 where interleaved and asynchronous call can cause exceptions due to the shared buffers being created and disposed before the previous call can retrieve its contents. +- When using the THttpTransport, interleaved asynchronous calls can cause the following exception: System.AggregateException : One or more errors occurred. Couldn't connect to server: System.Net.Http.HttpRequestException: Error while copying content to a stream. ---> System.ObjectDisposedException: Cannot access a closed Stream. +- Essentially, the THttpTransport._outputStream can get disposed while another asynchronous call is running. This is due to allocating the `_outputStream` at the class level. The stream should be allocated on a call-by-call basis. +- backward compatibility must be maintained for the shared buffer semantics. +- ensure race conditions are handled safely. +- ensure disposing buffers and clients are done safely in both the happy path and any situation where an exception may interrupt the execution path. +- ensure the time-outs are handled properly from end-to-end. +- ensure TDD practices are followed. +- update the thrift\compiler\cpp\src\thrift\generate\t_netstd_generator.cc with any required changes to the TBaseClient constructors. +- ensure any new public and protected members have a documentation comment, including the t_netstd_generator.cc additions. +- if t_netstd_generator.cc has additions, update the documentation of the existing code to clarify why one would use the previous versions. +- add a tutorial addition for the per-call transport usage (test to ensure tutorial works). +- update the netstd/README.md with per-call documentation and usage. + +## Constraints and Non-Goals diff --git a/tutorial/netstd/Client/Program.cs b/tutorial/netstd/Client/Program.cs index 29e3f3c3597..544c915b2d2 100644 --- a/tutorial/netstd/Client/Program.cs +++ b/tutorial/netstd/Client/Program.cs @@ -65,6 +65,8 @@ will diplay help information Client -tr: -bf: -pr: [-mc:] [-multiplex] will run client with specified arguments (tcp transport and binary protocol by default) and with 1 client + Client -tr:http -per-call [-mc:] + will exercise concurrent calls through a per-call-enabled HTTP client Options: -tr (transport): @@ -85,6 +87,8 @@ will run client with specified arguments (tcp transport and binary protocol by d -multiplex - adds multiplexed protocol + -per-call - enables concurrent HTTP calls with per-call transport state + -mc (multiple clients): - number of multiple clients to connect to server (max 100, default 1) @@ -130,6 +134,20 @@ private static async Task RunAsync(string[] args, CancellationToken cancel if (Logger.IsEnabled(LogLevel.Information)) Logger.LogInformation("Multiplex {mplex}", mplex); + if (args.Contains("-per-call")) + { + if (transport != Transport.Http) + throw new ArgumentException("-per-call requires -tr:http."); + if (mplex) + throw new ArgumentException("-per-call cannot be combined with -multiplex."); + + var perCallTasks = new Task[numClients]; + for (int i = 0; i < numClients; i++) + perCallTasks[i] = RunPerCallClientAsync(MakeTransport(args), MakeProtocolFactory(args), cancellationToken); + + return (await Task.WhenAll(perCallTasks)).All(succeeded => succeeded); + } + var tasks = new Task[numClients]; for (int i = 0; i < numClients; i++) { @@ -309,6 +327,49 @@ private static TProtocol MakeProtocol(string[] args, TTransport transport) }; } + private static TProtocolFactory MakeProtocolFactory(string[] args) + { + Protocol selectedProtocol = GetProtocol(args); + return selectedProtocol switch + { + Protocol.Binary => new TBinaryProtocol.Factory(), + Protocol.Compact => new TCompactProtocol.Factory(), + Protocol.Json => new TJsonProtocol.Factory(), + _ => throw new Exception("unhandled protocol"), + }; + } + + private static async Task RunPerCallClientAsync(TTransport transport, + TProtocolFactory protocolFactory, CancellationToken cancellationToken) + { + try + { + if (!(transport is ITPerCallTransportProvider)) + throw new InvalidOperationException("The selected transport does not support per-call operations."); + + using var client = new Calculator.Client(transport, protocolFactory, protocolFactory); + await client.OpenTransportAsync(cancellationToken); + + var sums = await Task.WhenAll(Enumerable.Range(1, 8) + .Select(value => client.add(value, value, cancellationToken))); + for (int index = 0; index < sums.Length; index++) + { + var expected = (index + 1) * 2; + if (sums[index] != expected) + throw new InvalidOperationException($"Concurrent add returned {sums[index]}, expected {expected}."); + } + + await ExecuteCalculatorClientOperations(client, cancellationToken); + Logger.LogInformation("PER_CALL_WORKFLOW_OK"); + return true; + } + catch (Exception exception) + { + Logger.LogError("Per-call tutorial workflow failed: {exception}", exception); + return false; + } + } + private static async Task RunClientAsync(TProtocol protocol, bool multiplex, CancellationToken cancellationToken) { try diff --git a/tutorial/netstd/README.md b/tutorial/netstd/README.md index 8301e2702ed..04b1f83d240 100644 --- a/tutorial/netstd/README.md +++ b/tutorial/netstd/README.md @@ -77,6 +77,9 @@ Usage: Client -tr: -pr: -mc: will run client with specified arguments (tcp transport and binary protocol by default) + Client -tr:http -per-call [-mc:] + will run concurrent HTTP calls with a per-call-enabled generated client per task + Options: -tr (transport): @@ -96,12 +99,16 @@ Options: json - json protocol will be used multiplexed - multiplexed protocol will be used + -per-call: + enables per-call transport state for concurrent high-level asynchronous HTTP calls + -mc (multiple clients): - number of multiple clients to connect to server (max 100, default 1) Sample: Client -tr:tcp -pr:binary -mc:10 + Client -tr:http -per-call -mc:2 Remarks: diff --git a/tutorial/netstd/smoketest.sh b/tutorial/netstd/smoketest.sh index f786b084a20..b061d18726a 100755 --- a/tutorial/netstd/smoketest.sh +++ b/tutorial/netstd/smoketest.sh @@ -36,8 +36,13 @@ CASES=( "-tr:tcptls" "-tr:namedpipe" "-tr:http" + "-tr:http -per-call" ) +if [ "${THRIFT_TUTORIAL_PER_CALL_ONLY:-0}" = "1" ]; then + CASES=("-tr:http -per-call") +fi + cd "$(dirname "$0")" || exit 1 SERVER="" @@ -137,6 +142,15 @@ run_case() { show_logs return 1 fi + case "$args" in + *-per-call*) + if ! grep -q "PER_CALL_WORKFLOW_OK" "$CLOG"; then + echo "FAIL $args -mc:$CLIENTS: per-call workflow marker was not written" + show_logs + return 1 + fi + ;; + esac echo "PASS $args -mc:$CLIENTS" }