From 81791229bd23b4da5da812f7d6735a26953b4ffc Mon Sep 17 00:00:00 2001 From: Simon Shanks Date: Wed, 16 Sep 2026 17:20:02 +0100 Subject: [PATCH] add asynchronous I/O cancellation ability --- docs/README.md | 21 ++ .../Connection/ConnectionCancellationTests.cs | 156 ++++++++++++++ .../ConnectionSynchronousFailureTests.cs | 38 ++++ kx/c.cs | 197 ++++++++++++++---- 4 files changed, 372 insertions(+), 40 deletions(-) create mode 100644 kx.Test/Connection/ConnectionCancellationTests.cs create mode 100644 kx.Test/Connection/ConnectionSynchronousFailureTests.cs diff --git a/docs/README.md b/docs/README.md index 531b5f3..0063955 100644 --- a/docs/README.md +++ b/docs/README.md @@ -260,6 +260,27 @@ Asynchronous I/O does not cause q to execute a request in parallel and does not Do not interleave reads and writes for multiple request/response exchanges on the same connection, because a response could be associated with the wrong request. Use a separate connection for each concurrently active exchange, or serialize access to a shared connection. +#### Asynchronous I/O cancellation (send and receive timeouts) + +The asynchronous methods also have overloads that accept a `CancellationToken`: + +```c# +using (var cancellation = new CancellationTokenSource(TimeSpan.FromSeconds(10))) +{ + await connection.knAsync("2 + 3".ToCharArray(), cancellation.Token); + object result = await connection.kAsync(cancellation.Token); +} +``` + +Cancellation-token overloads are available for `kAsync`, `k0Async`, every `ksAsync` overload, `knAsync`, and `krAsync`. +The existing overloads remain available and behave as though `CancellationToken.None` was supplied. + +`SendTimeout` and `ReceiveTimeout` apply only to synchronous I/O; they do not impose a deadline on `ReadAsync` or `WriteAsync`. +Use a `CancellationTokenSource`, optionally with `CancelAfter` or a `TimeSpan` timeout, to place a deadline on an asynchronous exchange. + +If cancellation interrupts an asynchronous read or write, the connection is closed before `OperationCanceledException` is rethrown. +Do not reuse that connection: cancellation may have interrupted a partially transmitted or partially received q IPC message. + ## Accessing items of arrays We can access elements using the `at` method: diff --git a/kx.Test/Connection/ConnectionCancellationTests.cs b/kx.Test/Connection/ConnectionCancellationTests.cs new file mode 100644 index 0000000..7847749 --- /dev/null +++ b/kx.Test/Connection/ConnectionCancellationTests.cs @@ -0,0 +1,156 @@ +using System; +using System.IO; +using System.Threading; +using System.Threading.Tasks; +using NUnit.Framework; + +namespace kx.Test.Connection +{ + [TestFixture] + public class ConnectionCancellationTests + { + [Test] + public async Task CancellingAsyncWriteClosesConnection() + { + using (var stream = new CancellableStream()) + using (var connection = new c(stream)) + using (var cancellation = new CancellationTokenSource()) + { + Task write = connection.ksAsync("test", cancellation.Token); + await stream.WriteStarted.Task; + + cancellation.Cancel(); + + Assert.CatchAsync(async () => await write); + Assert.AreEqual(cancellation.Token, stream.WriteCancellationToken); + Assert.IsTrue(stream.IsDisposed); + } + } + + [Test] + public async Task CancellingAsyncReadClosesConnection() + { + using (var stream = new CancellableStream()) + using (var connection = new c(stream)) + using (var cancellation = new CancellationTokenSource()) + { + Task read = connection.k0Async(cancellation.Token); + await stream.ReadStarted.Task; + + cancellation.Cancel(); + + Assert.CatchAsync(async () => await read); + Assert.AreEqual(cancellation.Token, stream.ReadCancellationToken); + Assert.IsTrue(stream.IsDisposed); + } + } + + private sealed class CancellableStream : Stream + { + internal TaskCompletionSource ReadStarted { get; } = + new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + internal TaskCompletionSource WriteStarted { get; } = + new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + internal CancellationToken ReadCancellationToken { get; private set; } + + internal CancellationToken WriteCancellationToken { get; private set; } + + internal bool IsDisposed { get; private set; } + + public override bool CanRead => true; + + public override bool CanSeek => false; + + public override bool CanWrite => true; + + public override long Length => throw new NotSupportedException(); + + public override long Position + { + get => throw new NotSupportedException(); + set => throw new NotSupportedException(); + } + + public override void Flush() + { + } + + public override int Read(byte[] buffer, int offset, int count) + { + throw new NotSupportedException(); + } + + public override Task ReadAsync( + byte[] buffer, + int offset, + int count, + CancellationToken cancellationToken) + { + ReadCancellationToken = cancellationToken; + ReadStarted.TrySetResult(true); + return WaitForReadCancellation(cancellationToken); + } + + public override long Seek(long offset, SeekOrigin origin) + { + throw new NotSupportedException(); + } + + public override void SetLength(long value) + { + throw new NotSupportedException(); + } + + public override void Write(byte[] buffer, int offset, int count) + { + throw new NotSupportedException(); + } + + public override Task WriteAsync( + byte[] buffer, + int offset, + int count, + CancellationToken cancellationToken) + { + WriteCancellationToken = cancellationToken; + WriteStarted.TrySetResult(true); + return Task.Delay(Timeout.Infinite, cancellationToken); + } + +#if NETCOREAPP3_1_OR_GREATER + public override ValueTask ReadAsync( + Memory buffer, + CancellationToken cancellationToken = default) + { + ReadCancellationToken = cancellationToken; + ReadStarted.TrySetResult(true); + return new ValueTask(WaitForReadCancellation(cancellationToken)); + } + + public override ValueTask WriteAsync( + ReadOnlyMemory buffer, + CancellationToken cancellationToken = default) + { + WriteCancellationToken = cancellationToken; + WriteStarted.TrySetResult(true); + return new ValueTask(Task.Delay(Timeout.Infinite, cancellationToken)); + } +#endif + + protected override void Dispose(bool disposing) + { + IsDisposed = true; + base.Dispose(disposing); + } + + private static async Task WaitForReadCancellation( + CancellationToken cancellationToken) + { + await Task.Delay(Timeout.Infinite, cancellationToken); + return 0; + } + } + } +} diff --git a/kx.Test/Connection/ConnectionSynchronousFailureTests.cs b/kx.Test/Connection/ConnectionSynchronousFailureTests.cs new file mode 100644 index 0000000..9a1f3d7 --- /dev/null +++ b/kx.Test/Connection/ConnectionSynchronousFailureTests.cs @@ -0,0 +1,38 @@ +using System.IO; +using Moq; +using NUnit.Framework; + +namespace kx.Test.Connection +{ + [TestFixture] + public class ConnectionSynchronousFailureTests + { + [Test] + public void SynchronousWriteFailureClosesConnection() + { + var stream = new Mock(); + stream.Setup(s => s.Write(It.IsAny(), It.IsAny(), It.IsAny())) + .Throws(new IOException("write timeout")); + + using (var connection = new c(stream.Object)) + { + Assert.Throws(() => connection.ks("test")); + stream.Verify(s => s.Close(), Times.Once); + } + } + + [Test] + public void SynchronousReadFailureClosesConnection() + { + var stream = new Mock(); + stream.Setup(s => s.Read(It.IsAny(), It.IsAny(), It.IsAny())) + .Throws(new IOException("read timeout")); + + using (var connection = new c(stream.Object)) + { + Assert.Throws(() => connection.k0()); + stream.Verify(s => s.Close(), Times.Once); + } + } + } +} diff --git a/kx/c.cs b/kx/c.cs index 76dcefb..a8db7a8 100644 --- a/kx/c.cs +++ b/kx/c.cs @@ -6,6 +6,7 @@ using System.Security.Authentication; using System.Security.Cryptography.X509Certificates; using System.Text; +using System.Threading; using System.Threading.Tasks; namespace kx @@ -577,9 +578,21 @@ protected virtual void Dispose(bool disposing){ /// /// Deserialised response to request. /// - public async Task kAsync() + public Task kAsync() { - await k0Async().ConfigureAwait(false); + return kAsync(CancellationToken.None); + } + + /// + /// Reads an incoming message from the remote KDB+ process async. + /// + /// The token used to cancel the asynchronous operation. + /// + /// Deserialised response to request. + /// + public async Task kAsync(CancellationToken cancellationToken) + { + await k0Async(cancellationToken).ConfigureAwait(false); return r(); } @@ -734,16 +747,25 @@ public void k0() /// /// Waits for an async message and read header. /// - public async Task k0Async() + public Task k0Async() + { + return k0Async(CancellationToken.None); + } + + /// + /// Waits for an async message and reads its header. + /// + /// The token used to cancel the asynchronous operation. + public async Task k0Async(CancellationToken cancellationToken) { _readBuffer = new byte[8]; - await ReadAsync(_readBuffer).ConfigureAwait(false); + await ReadAsync(_readBuffer, cancellationToken).ConfigureAwait(false); ParseHeader(); _readPosition = 4; _readBuffer = new byte[ri() - 8]; - await ReadAsync(_readBuffer).ConfigureAwait(false); + await ReadAsync(_readBuffer, cancellationToken).ConfigureAwait(false); if (IsCompressed) { @@ -760,23 +782,44 @@ public async Task k0Async() /// Sends an async message to the remote KDB+ process with a specified object parameter. /// /// The object parameter. - public async Task ksAsync(object x) + public Task ksAsync(object x) + { + return ksAsync(x, CancellationToken.None); + } + + /// + /// Sends an async message to the remote KDB+ process with a specified object parameter. + /// + /// The object parameter. + /// The token used to cancel the asynchronous operation. + public async Task ksAsync(object x, CancellationToken cancellationToken) + { + await wAsync(0, x, cancellationToken).ConfigureAwait(false); + } + + /// + /// Sends an async message to the remote KDB+ process with a specified expression. + /// + /// The expression to send. + /// parameter was null. + public Task ksAsync(string s) { - await wAsync(0, x).ConfigureAwait(false); + return ksAsync(s, CancellationToken.None); } /// /// Sends an async message to the remote KDB+ process with a specified expression. /// /// The expression to send. + /// The token used to cancel the asynchronous operation. /// parameter was null. - public async Task ksAsync(string s) + public async Task ksAsync(string s, CancellationToken cancellationToken) { if (s == null) { throw new ArgumentNullException(nameof(s)); } - await wAsync(0, s.ToCharArray()).ConfigureAwait(false); + await wAsync(0, s.ToCharArray(), cancellationToken).ConfigureAwait(false); } /// @@ -786,7 +829,20 @@ public async Task ksAsync(string s) /// The expression to send. /// The object parameter to send. /// parameter was null. - public async Task ksAsync(string s, object x) + public Task ksAsync(string s, object x) + { + return ksAsync(s, x, CancellationToken.None); + } + + /// + /// Sends an async message to the remote KDB+ process with a specified expression + /// and object parameter. + /// + /// The expression to send. + /// The object parameter to send. + /// The token used to cancel the asynchronous operation. + /// parameter was null. + public async Task ksAsync(string s, object x, CancellationToken cancellationToken) { if (s == null) { @@ -798,7 +854,7 @@ public async Task ksAsync(string s, object x) x }; - await wAsync(0, array).ConfigureAwait(false); + await wAsync(0, array, cancellationToken).ConfigureAwait(false); } /// @@ -809,7 +865,21 @@ public async Task ksAsync(string s, object x) /// The first object parameter to send. /// The second object parameter to send. /// parameter was null. - public async Task ksAsync(string s, object x, object y) + public Task ksAsync(string s, object x, object y) + { + return ksAsync(s, x, y, CancellationToken.None); + } + + /// + /// Sends an async message to the remote KDB+ process with a specified expression + /// and object parameters. + /// + /// The expression to send. + /// The first object parameter to send. + /// The second object parameter to send. + /// The token used to cancel the asynchronous operation. + /// parameter was null. + public async Task ksAsync(string s, object x, object y, CancellationToken cancellationToken) { if (s == null) { @@ -822,7 +892,7 @@ public async Task ksAsync(string s, object x, object y) y }; - await wAsync(0, array).ConfigureAwait(false); + await wAsync(0, array, cancellationToken).ConfigureAwait(false); } /// @@ -907,10 +977,21 @@ public void kn(object x) /// Sends an async message to the remote KDB+ process with a specified object parameter. /// /// The object parameter. - public async Task knAsync(object x) + public Task knAsync(object x) { - await wAsync(1, x).ConfigureAwait(false); + return knAsync(x, CancellationToken.None); } + + /// + /// Sends an async message to the remote KDB+ process with a specified object parameter. + /// + /// The object parameter. + /// The token used to cancel the asynchronous operation. + public async Task knAsync(object x, CancellationToken cancellationToken) + { + await wAsync(1, x, cancellationToken).ConfigureAwait(false); + } + /// /// Sends a response message to the remote KDB+ process. /// @@ -930,9 +1011,22 @@ public void kr(object x) /// /// This should be called only during processing of an incoming sync message. /// - public async Task krAsync(object x) + public Task krAsync(object x) + { + return krAsync(x, CancellationToken.None); + } + + /// + /// Sends a response message to the remote KDB+ process. + /// + /// The response message to send. + /// The token used to cancel the asynchronous operation. + /// + /// This should be called only during processing of an incoming sync message. + /// + public async Task krAsync(object x, CancellationToken cancellationToken) { - await wAsync(2, x).ConfigureAwait(false); + await wAsync(2, x, cancellationToken).ConfigureAwait(false); } /// @@ -1074,13 +1168,32 @@ protected void Write(byte[] bytes, int number) /// /// The byte array to be writtern to the client stream. /// The number of bytes to be written to the client stream. - protected async Task WriteAsync(byte[] bytes, int number) + protected Task WriteAsync(byte[] bytes, int number) + { + return WriteAsync(bytes, number, CancellationToken.None); + } + + /// + /// Writes a specified byte array directly to the underlying client stream asynchronously. + /// + /// The byte array to be written to the client stream. + /// The number of bytes to be written to the client stream. + /// The token used to cancel the asynchronous operation. + protected async Task WriteAsync(byte[] bytes, int number, CancellationToken cancellationToken) { + try + { #if NETSTANDARD2_1_OR_GREATER || NETCOREAPP2_1_OR_GREATER - await _clientStream.WriteAsync(bytes.AsMemory(0, number)).ConfigureAwait(false); + await _clientStream.WriteAsync(bytes.AsMemory(0, number), cancellationToken).ConfigureAwait(false); #else - await _clientStream.WriteAsync(bytes,0,number).ConfigureAwait(false); + await _clientStream.WriteAsync(bytes,0,number,cancellationToken).ConfigureAwait(false); #endif + } + catch (OperationCanceledException) + { + Close(); + throw; + } } /// @@ -2107,14 +2220,10 @@ private object r() } } - private async Task wAsync(int i, object x) + private async Task wAsync(int i, object x, CancellationToken cancellationToken) { byte[] buffer = Serialize(i, x); -#if NETSTANDARD2_1_OR_GREATER || NETCOREAPP2_1_OR_GREATER - await _clientStream.WriteAsync(buffer.AsMemory(0, buffer.Length)).ConfigureAwait(false); -#else - await _clientStream.WriteAsync(buffer,0,buffer.Length).ConfigureAwait(false); -#endif + await WriteAsync(buffer, buffer.Length, cancellationToken).ConfigureAwait(false); } private void read(byte[] b) @@ -2146,29 +2255,37 @@ private void read(byte[] b) } } - private async Task ReadAsync(byte[] b) + private async Task ReadAsync(byte[] b, CancellationToken cancellationToken) { - int k = 0; - int j = b.Length; - while (true) + try { - if (k < j) + int k = 0; + int j = b.Length; + while (true) { - int i; + if (k < j) + { + int i; #if NETSTANDARD2_1_OR_GREATER || NETCOREAPP2_1_OR_GREATER - if ((i = await _clientStream.ReadAsync(b.AsMemory(k, Math.Min(_maxBufferSize, j - k))).ConfigureAwait(false)) == 0) + if ((i = await _clientStream.ReadAsync(b.AsMemory(k, Math.Min(_maxBufferSize, j - k)),cancellationToken).ConfigureAwait(false)) == 0) #else - if ((i = await _clientStream.ReadAsync(b,k,Math.Min(_maxBufferSize,j-k)).ConfigureAwait(false)) == 0) + if ((i = await _clientStream.ReadAsync(b,k,Math.Min(_maxBufferSize,j-k),cancellationToken).ConfigureAwait(false)) == 0) #endif - { - break; + { + break; + } + k += i; + continue; } - k += i; - continue; + return; } - return; + throw new KException("read"); + } + catch (OperationCanceledException) + { + Close(); + throw; } - throw new KException("read"); } private static int ns(string s)