Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
156 changes: 156 additions & 0 deletions kx.Test/Connection/ConnectionCancellationTests.cs
Original file line number Diff line number Diff line change
@@ -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<OperationCanceledException>(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<OperationCanceledException>(async () => await read);
Assert.AreEqual(cancellation.Token, stream.ReadCancellationToken);
Assert.IsTrue(stream.IsDisposed);
}
}

private sealed class CancellableStream : Stream
{
internal TaskCompletionSource<bool> ReadStarted { get; } =
new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);

internal TaskCompletionSource<bool> WriteStarted { get; } =
new TaskCompletionSource<bool>(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<int> 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<int> ReadAsync(
Memory<byte> buffer,
CancellationToken cancellationToken = default)
{
ReadCancellationToken = cancellationToken;
ReadStarted.TrySetResult(true);
return new ValueTask<int>(WaitForReadCancellation(cancellationToken));
}

public override ValueTask WriteAsync(
ReadOnlyMemory<byte> 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<int> WaitForReadCancellation(
CancellationToken cancellationToken)
{
await Task.Delay(Timeout.Infinite, cancellationToken);
return 0;
}
}
}
}
38 changes: 38 additions & 0 deletions kx.Test/Connection/ConnectionSynchronousFailureTests.cs
Original file line number Diff line number Diff line change
@@ -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>();
stream.Setup(s => s.Write(It.IsAny<byte[]>(), It.IsAny<int>(), It.IsAny<int>()))
.Throws(new IOException("write timeout"));

using (var connection = new c(stream.Object))
{
Assert.Throws<IOException>(() => connection.ks("test"));
stream.Verify(s => s.Close(), Times.Once);
}
}

[Test]
public void SynchronousReadFailureClosesConnection()
{
var stream = new Mock<Stream>();
stream.Setup(s => s.Read(It.IsAny<byte[]>(), It.IsAny<int>(), It.IsAny<int>()))
.Throws(new IOException("read timeout"));

using (var connection = new c(stream.Object))
{
Assert.Throws<IOException>(() => connection.k0());
stream.Verify(s => s.Close(), Times.Once);
}
}
}
}
Loading
Loading