diff --git a/src/libraries/Common/src/System/Threading/UnixHandleAsyncContextUnsafeAccess.cs b/src/libraries/Common/src/System/Threading/UnixHandleAsyncContextUnsafeAccess.cs new file mode 100644 index 00000000000000..90ab82301c5bd9 --- /dev/null +++ b/src/libraries/Common/src/System/Threading/UnixHandleAsyncContextUnsafeAccess.cs @@ -0,0 +1,130 @@ +// Licensed to the .NET Foundation under one or more agreements. +// The .NET Foundation licenses this file to you under the MIT license. + +using System.Runtime.CompilerServices; +using System.Runtime.InteropServices; + +namespace System.Threading +{ + internal static class UnixHandleAsyncContextUnsafeAccess + { + private const string AsyncContextTypeName = "System.Threading.UnixHandleAsyncContext, System.Private.CoreLib"; + private const string OperationTypeName = "System.Threading.UnixHandleAsyncContext+Operation, System.Private.CoreLib"; + private const string DelegateOperationTypeName = "System.Threading.UnixHandleAsyncContext+DelegateOperation, System.Private.CoreLib"; + + public enum AsyncResult + { + Pending = 0, + Completed = 1, + Aborted = 2, + } + + public enum SyncResult + { + Completed = 1, + Aborted = 2, + TimedOut = 4, + } + + public enum OnCompletedResult + { + Completed = 1, + Aborted = 2, + Canceled = 3, + } + + [UnsafeAccessor(UnsafeAccessorKind.Constructor)] + [return: UnsafeAccessorType(AsyncContextTypeName)] + private static extern object CreateAsyncContext(SafeHandle handle); + + [UnsafeAccessor(UnsafeAccessorKind.Constructor)] + [return: UnsafeAccessorType(DelegateOperationTypeName)] + private static extern object CreateDelegateOperation(Func tryComplete, Action onCompleted); + + [UnsafeAccessor(UnsafeAccessorKind.Method, Name = "IsReadReady")] + private static extern bool IsReadReady( + [UnsafeAccessorType(AsyncContextTypeName)] object context, + out int observedSequenceNumber); + + [UnsafeAccessor(UnsafeAccessorKind.Method, Name = "IsWriteReady")] + private static extern bool IsWriteReady( + [UnsafeAccessorType(AsyncContextTypeName)] object context, + out int observedSequenceNumber); + + [UnsafeAccessor(UnsafeAccessorKind.Method, Name = "StartAsyncReadAsInt")] + private static extern int StartAsyncRead( + [UnsafeAccessorType(AsyncContextTypeName)] object context, + [UnsafeAccessorType(OperationTypeName)] object operation, + int observedSequenceNumber, + CancellationToken cancellationToken); + + [UnsafeAccessor(UnsafeAccessorKind.Method, Name = "StartAsyncWriteAsInt")] + private static extern int StartAsyncWrite( + [UnsafeAccessorType(AsyncContextTypeName)] object context, + [UnsafeAccessorType(OperationTypeName)] object operation, + int observedSequenceNumber, + CancellationToken cancellationToken); + + [UnsafeAccessor(UnsafeAccessorKind.Method, Name = "ReadAsInt")] + private static extern int ReadSync( + [UnsafeAccessorType(AsyncContextTypeName)] object context, + [UnsafeAccessorType(OperationTypeName)] object operation, + int observedSequenceNumber, + int timeout); + + [UnsafeAccessor(UnsafeAccessorKind.Method, Name = "WriteAsInt")] + private static extern int WriteSync( + [UnsafeAccessorType(AsyncContextTypeName)] object context, + [UnsafeAccessorType(OperationTypeName)] object operation, + int observedSequenceNumber, + int timeout); + + [UnsafeAccessor(UnsafeAccessorKind.Method, Name = "AbortAndDispose")] + private static extern bool AbortAndDispose( + [UnsafeAccessorType(AsyncContextTypeName)] object context); + + public sealed class AsyncContext + { + private readonly object _context; + + public AsyncContext(SafeHandle handle) + { + _context = CreateAsyncContext(handle); + } + + public bool IsReadReady(out int observedSequenceNumber) + => UnixHandleAsyncContextUnsafeAccess.IsReadReady(_context, out observedSequenceNumber); + + public bool IsWriteReady(out int observedSequenceNumber) + => UnixHandleAsyncContextUnsafeAccess.IsWriteReady(_context, out observedSequenceNumber); + + public AsyncResult StartAsyncRead(Operation operation, int observedSequenceNumber, CancellationToken cancellationToken) + => (AsyncResult)UnixHandleAsyncContextUnsafeAccess.StartAsyncRead(_context, operation.Instance, observedSequenceNumber, cancellationToken); + + public AsyncResult StartAsyncWrite(Operation operation, int observedSequenceNumber, CancellationToken cancellationToken) + => (AsyncResult)UnixHandleAsyncContextUnsafeAccess.StartAsyncWrite(_context, operation.Instance, observedSequenceNumber, cancellationToken); + + public SyncResult Read(Operation operation, int observedSequenceNumber, int timeout) + => (SyncResult)ReadSync(_context, operation.Instance, observedSequenceNumber, timeout); + + public SyncResult Write(Operation operation, int observedSequenceNumber, int timeout) + => (SyncResult)WriteSync(_context, operation.Instance, observedSequenceNumber, timeout); + + public bool AbortAndDispose() + => UnixHandleAsyncContextUnsafeAccess.AbortAndDispose(_context); + + public static Operation CreateOperation(Func tryComplete, Action onCompleted) + => new Operation(tryComplete, result => onCompleted((OnCompletedResult)result)); + + public readonly struct Operation + { + public object Instance { get; } + + internal Operation(Func tryComplete, Action onCompleted) + { + Instance = CreateDelegateOperation(tryComplete, onCompleted); + } + } + } + } +} diff --git a/src/libraries/System.IO.Ports/src/System.IO.Ports.csproj b/src/libraries/System.IO.Ports/src/System.IO.Ports.csproj index ceadbb41390e60..24a83c25ecb943 100644 --- a/src/libraries/System.IO.Ports/src/System.IO.Ports.csproj +++ b/src/libraries/System.IO.Ports/src/System.IO.Ports.csproj @@ -23,6 +23,8 @@ System.IO.Ports.SerialPort $([MSBuild]::GetTargetPlatformIdentifier('$(TargetFramework)')) true SR.PlatformNotSupported_IOPorts + + true @@ -142,6 +144,17 @@ System.IO.Ports.SerialPort Link="Common\Interop\Unix\Interop.Poll.Structs.cs" /> + + + + + + + + + + diff --git a/src/libraries/System.IO.Ports/src/System/IO/Ports/SafeSerialDeviceHandle.Unix.cs b/src/libraries/System.IO.Ports/src/System/IO/Ports/SafeSerialDeviceHandle.Unix.cs index 025643262b5708..d7e4eb4d526358 100644 --- a/src/libraries/System.IO.Ports/src/System/IO/Ports/SafeSerialDeviceHandle.Unix.cs +++ b/src/libraries/System.IO.Ports/src/System/IO/Ports/SafeSerialDeviceHandle.Unix.cs @@ -6,12 +6,19 @@ using System.IO; using System.Net.Sockets; using System.Runtime.InteropServices; +using System.Threading; using Microsoft.Win32.SafeHandles; namespace System.IO.Ports { - internal sealed class SafeSerialDeviceHandle : SafeHandleMinusOneIsInvalid + internal sealed partial class SafeSerialDeviceHandle : SafeHandleMinusOneIsInvalid { + // When the user calls Dispose, some operations for a pending read event might be in flight. + // If these use the handle, the user Dispose call won't actually release the handle immediately + // which causes opening the port after the Dispose to fail (EBUSY). + // DisposeLock guards these operations so that they can not happen concurrent with the Dispose. + private object DisposeLock => this; + public SafeSerialDeviceHandle() : base(ownsHandle: true) { } @@ -34,6 +41,44 @@ internal static SafeSerialDeviceHandle Open(string portName) return handle; } + // Gets the amount input data buffered by the handle. + partial void GetBufferedCount(ref int count); + + // Get the amount of bytes that can be read from the handle plus the amount buffered by the caller. + // When throwOnDispose is false, returns 0 when disposed instead of throwing. + internal int GetBytesToRead(int buffered, bool throwOnDispose = true) + { + lock (DisposeLock) + { + if (!throwOnDispose && IsClosed) + { + return 0; + } + + try + { + return buffered + BytesToRead; + } + catch (ObjectDisposedException) when (!throwOnDispose) + { + return 0; + } + } + } + + // Gets the amount of bytes that can be read from the handle. + private int BytesToRead + { + get + { + Debug.Assert(Monitor.IsEntered(DisposeLock)); + + int buffered = 0; + GetBufferedCount(ref buffered); + return Math.Max(Interop.Termios.TermiosGetAvailableBytes(this, true), 0) + buffered; + } + } + protected override bool ReleaseHandle() { Interop.Serial.Shutdown(handle, SocketShutdown.Both); diff --git a/src/libraries/System.IO.Ports/src/System/IO/Ports/SafeSerialDeviceHandle.UnixAsyncContext.cs b/src/libraries/System.IO.Ports/src/System/IO/Ports/SafeSerialDeviceHandle.UnixAsyncContext.cs new file mode 100644 index 00000000000000..8abb6fef64b457 --- /dev/null +++ b/src/libraries/System.IO.Ports/src/System/IO/Ports/SafeSerialDeviceHandle.UnixAsyncContext.cs @@ -0,0 +1,756 @@ +// Licensed to the .NET Foundation under one or more agreements. +// The .NET Foundation licenses this file to you under the MIT license. + +using System; +using System.Diagnostics; +using System.IO; +using System.Runtime.CompilerServices; +using System.Runtime.InteropServices; +using System.Threading; +using System.Threading.Tasks; +using System.Threading.Tasks.Sources; +using Microsoft.Win32.SafeHandles; +using UnixHandleAsyncContext = System.Threading.UnixHandleAsyncContextUnsafeAccess.AsyncContext; +using AsyncResult = System.Threading.UnixHandleAsyncContextUnsafeAccess.AsyncResult; +using SyncResult = System.Threading.UnixHandleAsyncContextUnsafeAccess.SyncResult; +using OnCompletedResult = System.Threading.UnixHandleAsyncContextUnsafeAccess.OnCompletedResult; + +namespace System.IO.Ports +{ + internal sealed partial class SafeSerialDeviceHandle : SafeHandleMinusOneIsInvalid + { + internal enum ReceiveThresholdResult + { + Success, + Eof, + Error, + Disposed, + } + + private UnixHandleAsyncContext? _asyncContext; + private ReadOperation? _cachedReadOp; + private WriteOperation? _cachedWriteOp; + private WaitReceiveThresholdOperation? _cachedWaitReceiveThresholdOp; + private byte[]? _thresholdBuffer; + private int _thresholdBufferCount; + + // This guards the _thresholdBuffer which different concurrent readers may try to use. + private object ReadLock => DisposeLock; // Use the same lock as the dispose lock (no potential ordering issues). + + partial void GetBufferedCount(ref int count) + => count = Volatile.Read(ref _thresholdBufferCount); + + internal void DiscardInBuffer() + { + lock (ReadLock) + { + // This may or may not work depending on hardware. + Interop.Termios.TermiosDiscard(this, Interop.Termios.Queue.ReceiveQueue); + + _thresholdBufferCount = 0; + } + } + + private unsafe bool ReadIntoThresholdBuffer(int threshold, out int bytesAvailable, out ReceiveThresholdResult result) + { + // Our caller is already holding the ReadLock. + Debug.Assert(Monitor.IsEntered(ReadLock)); + + int currentCount = _thresholdBufferCount; + + if (_thresholdBuffer == null || _thresholdBuffer.Length < threshold) + { + Array.Resize(ref _thresholdBuffer, threshold); + } + + fixed (byte* bufPtr = &_thresholdBuffer[currentCount]) + { + int bytesRead = Interop.Serial.Read(this, bufPtr, threshold - currentCount); + bytesAvailable = _thresholdBufferCount = currentCount + Math.Max(bytesRead, 0); + if (bytesRead < 0 && !IsPending(Interop.Sys.GetLastErrorInfo())) + { + result = ReceiveThresholdResult.Error; + return true; + } + bool eof = bytesRead == 0; + result = eof ? ReceiveThresholdResult.Eof : ReceiveThresholdResult.Success; + return eof || bytesAvailable >= threshold; + } + } + + private int DrainThresholdBuffer(Span destination) + { + int count = _thresholdBufferCount; + if (count == 0) + { + return 0; + } + + int toCopy = Math.Min(count, destination.Length); + _thresholdBuffer.AsSpan(0, toCopy).CopyTo(destination); + + int remaining = count - toCopy; + if (remaining > 0) + { + Buffer.BlockCopy(_thresholdBuffer!, toCopy, _thresholdBuffer!, 0, remaining); + } + _thresholdBufferCount = remaining; + + return toCopy; + } + + private UnixHandleAsyncContext AsyncContext + { + get + { + if (_asyncContext == null) + { + Interlocked.CompareExchange(ref _asyncContext, new UnixHandleAsyncContext(this), null); + } + return _asyncContext!; + } + } + + private ReadOperation RentReadOperation() + => Interlocked.Exchange(ref _cachedReadOp, null) ?? new ReadOperation(this); + + private WriteOperation RentWriteOperation() + => Interlocked.Exchange(ref _cachedWriteOp, null) ?? new WriteOperation(this); + + private WaitReceiveThresholdOperation RentWaitReceiveThresholdOperation() + => Interlocked.Exchange(ref _cachedWaitReceiveThresholdOp, null) ?? new WaitReceiveThresholdOperation(this); + + private void ReturnReadOperation(ReadOperation op) + { + op.Reset(); + Volatile.Write(ref _cachedReadOp, op); + } + + private void ReturnWriteOperation(WriteOperation op) + { + op.Reset(); + Volatile.Write(ref _cachedWriteOp, op); + } + + private void ReturnWaitReceiveThresholdOperation(WaitReceiveThresholdOperation op) + { + op.Reset(); + Volatile.Write(ref _cachedWaitReceiveThresholdOp, op); + } + + protected override void Dispose(bool disposing) + { + _asyncContext?.AbortAndDispose(); + lock (DisposeLock) + { + base.Dispose(disposing); + } + } + + internal unsafe int Read(Span buffer, int timeout, SerialStream? serialStream = null) + { + if (AsyncContext.IsReadReady(out int sequenceNumber) && + TryCompleteRead(buffer, out var result, serialStream)) + { + CheckIOResult(result.ErrorInfo); + return result.BytesRead; + } + + ReadOperation op = RentReadOperation(); + fixed (byte* bufPtr = &MemoryMarshal.GetReference(buffer)) + { + op.InitSync(bufPtr, buffer.Length, serialStream); + + SyncResult syncResult = AsyncContext.Read(op.Operation, sequenceNumber, timeout); + + if (syncResult == SyncResult.Completed) + { + int bytesRead = op.BytesRead; + Exception? exception = op.Exception; + + ReturnReadOperation(op); + + if (exception != null) + { + throw exception; + } + return bytesRead; + } + + if (syncResult == SyncResult.TimedOut) + { + ReturnReadOperation(op); + throw new TimeoutException(); + } + + throw new OperationCanceledException(); + } + } + + internal ValueTask ReadAsync(Memory destination, CancellationToken cancellationToken, SerialStream? serialStream = null) + { + if (AsyncContext.IsReadReady(out int sequenceNumber) && + TryCompleteRead(destination.Span, out var readResult, serialStream)) + { + CheckIOResult(readResult.ErrorInfo); + return new ValueTask(readResult.BytesRead); + } + + ReadOperation op = RentReadOperation(); + op.InitAsync(destination, cancellationToken, serialStream); + + AsyncResult result = AsyncContext.StartAsyncRead(op.Operation, sequenceNumber, cancellationToken); + + if (result == AsyncResult.Pending) + { + return new ValueTask(op, op.Version); + } + else if (result == AsyncResult.Completed) + { + int bytesRead = op.BytesRead; + Exception? exception = op.Exception; + + ReturnReadOperation(op); + + if (exception != null) + { + throw exception; + } + return new ValueTask(bytesRead); + } + + throw new OperationCanceledException(); + } + + internal unsafe void Write(ReadOnlySpan buffer, int timeout) + { + if (AsyncContext.IsWriteReady(out int sequenceNumber) && + TryCompleteWrite(buffer, out int bytesWritten, out Interop.ErrorInfo errorInfo)) + { + CheckIOResult(errorInfo); + return; + } + + WriteOperation op = RentWriteOperation(); + fixed (byte* bufPtr = &MemoryMarshal.GetReference(buffer)) + { + op.InitSync(bufPtr, buffer.Length); + + SyncResult syncResult = AsyncContext.Write(op.Operation, sequenceNumber, timeout); + + if (syncResult == SyncResult.Completed) + { + Exception? exception = op.Exception; + + ReturnWriteOperation(op); + + if (exception != null) + { + throw exception; + } + return; + } + + if (syncResult == SyncResult.TimedOut) + { + ReturnWriteOperation(op); + throw new TimeoutException(); + } + + throw new OperationCanceledException(); + } + } + + internal ValueTask WriteAsync(ReadOnlyMemory source, CancellationToken cancellationToken) + { + int bytesWritten = 0; + if (AsyncContext.IsWriteReady(out int sequenceNumber) && + TryCompleteWrite(source.Span, out bytesWritten, out Interop.ErrorInfo writeResult)) + { + CheckIOResult(writeResult); + return default; + } + + WriteOperation op = RentWriteOperation(); + op.InitAsync(source.Slice(bytesWritten), cancellationToken); + + AsyncResult result = AsyncContext.StartAsyncWrite(op.Operation, sequenceNumber, cancellationToken); + + if (result == AsyncResult.Pending) + { + return new ValueTask(op, op.Version); + } + else if (result == AsyncResult.Completed) + { + Exception? exception = op.Exception; + ReturnWriteOperation(op); + if (exception != null) + { + throw exception; + } + return default; + } + + throw new OperationCanceledException(); + } + + internal void WaitForReceiveThreshold(SerialStream serialStream) + { + if (AsyncContext.IsReadReady(out int sequenceNumber) && + TryCompleteWaitReceiveThreshold(serialStream.ReceivedBytesThreshold, out int bytesAvailable, out ReceiveThresholdResult thresholdResult)) + { + serialStream.OnReceiveThreshold(bytesAvailable, thresholdResult); + return; + } + + WaitReceiveThresholdOperation op = RentWaitReceiveThresholdOperation(); + op.Init(serialStream); + + AsyncResult result = AsyncContext.StartAsyncRead(op.Operation, sequenceNumber, CancellationToken.None); + + if (result == AsyncResult.Completed) + { + op.CompleteOperation(OnCompletedResult.Completed); + } + } + + private bool TryCompleteWaitReceiveThreshold(int receivedBytesThreshold, out int bytesAvailable, out ReceiveThresholdResult result) + { + lock (DisposeLock) + { + if (!IsClosed) + { + try + { + bytesAvailable = BytesToRead; + if (bytesAvailable >= receivedBytesThreshold) + { + result = ReceiveThresholdResult.Success; + return true; + } + + // Drain kernel data into the threshold buffer + // so epoll fires when new data arrives. + return ReadIntoThresholdBuffer(receivedBytesThreshold, out bytesAvailable, out result); + } + catch (ObjectDisposedException) + { } + } + + bytesAvailable = 0; + result = ReceiveThresholdResult.Disposed; + return true; + } + } + + private unsafe bool TryCompleteRead(Span buffer, out (int BytesRead, Interop.ErrorInfo ErrorInfo) result, SerialStream? serialStream = null) + { + fixed (byte* bufPtr = &MemoryMarshal.GetReference(buffer)) + { + return TryCompleteRead(bufPtr, buffer.Length, out result, serialStream); + } + } + + private unsafe bool TryCompleteRead(byte* buffer, int length, out (int BytesRead, Interop.ErrorInfo ErrorInfo) result, SerialStream? serialStream = null) + { + lock (ReadLock) + { + int fromBuffer = DrainThresholdBuffer(new Span(buffer, length)); + if (fromBuffer > 0) + { + buffer += fromBuffer; + length -= fromBuffer; + } + + int bytesRead = length == 0 ? 0 : Interop.Serial.Read(this, buffer, length); + if (bytesRead < 0 && fromBuffer == 0) + { + Interop.ErrorInfo errorInfo = Interop.Sys.GetLastErrorInfo(); + if (IsPending(errorInfo)) + { + result = default; + return false; + } + result = (-1, errorInfo); + return true; + } + + int totalRead = fromBuffer + Math.Max(bytesRead, 0); + if (totalRead > 0) + { + serialStream?.OnBytesRead(totalRead); + } + result = (totalRead, default); + return true; + } + } + + private unsafe bool TryCompleteWrite(ReadOnlySpan buffer, out int bytesWritten, out Interop.ErrorInfo errorInfo) + { + fixed (byte* bufPtr = &MemoryMarshal.GetReference(buffer)) + { + return TryCompleteWrite(bufPtr, buffer.Length, out bytesWritten, out errorInfo); + } + } + + private unsafe bool TryCompleteWrite(byte* buffer, int length, out int bytesWritten, out Interop.ErrorInfo errorInfo) + { + int totalBytesWritten = 0; + while (length > 0) + { + int written = Interop.Serial.Write(this, buffer, length); + if (written < 0) + { + errorInfo = Interop.Sys.GetLastErrorInfo(); + bytesWritten = totalBytesWritten; + if (!IsPending(errorInfo)) + { + return true; + } + return false; + } + + totalBytesWritten += written; + length -= written; + buffer += written; + } + + errorInfo = default; + bytesWritten = totalBytesWritten; + return true; + } + + private static void CheckIOResult(Interop.ErrorInfo errorInfo) + { + if (errorInfo.Error != Interop.Error.SUCCESS) + { + throw Interop.GetIOException(errorInfo); + } + } + + private static bool IsPending(Interop.ErrorInfo errorInfo) + => errorInfo.Error is Interop.Error.EAGAIN or Interop.Error.EWOULDBLOCK; + + private sealed unsafe class ReadOperation : IValueTaskSource + { + private readonly SafeSerialDeviceHandle _owner; + private readonly UnixHandleAsyncContext.Operation _callbackOp; + internal int BytesRead; + internal Exception? Exception; + private ManualResetValueTaskSourceCore _mrvtsc; + private SerialStream? _serialStream; + private Memory _buffer; + private byte* _syncBuffer; + private int _syncBufferLength; + private CancellationToken _cancellationToken; + + internal ReadOperation(SafeSerialDeviceHandle owner) + { + _owner = owner; + _callbackOp = UnixHandleAsyncContext.CreateOperation(TryCompleteOperation, OnCompleted); + } + + internal UnixHandleAsyncContext.Operation Operation => _callbackOp; + + internal short Version + => _mrvtsc.Version; + + internal void InitSync(byte* syncBuffer, int syncBufferLength, SerialStream? serialStream = null) + { + _serialStream = serialStream; + _syncBuffer = syncBuffer; + _syncBufferLength = syncBufferLength; + } + + internal void InitAsync(Memory buffer, CancellationToken cancellationToken, SerialStream? serialStream = null) + { + _serialStream = serialStream; + _buffer = buffer; + _cancellationToken = cancellationToken; + } + + internal void Reset() + { + _serialStream = null; + _buffer = default; + _syncBuffer = null; + _cancellationToken = default; + Exception = null; + _mrvtsc.Reset(); + } + + private bool TryCompleteOperation(SafeHandle handle) + { + (int BytesRead, Interop.ErrorInfo ErrorInfo) readResult; + + if (_syncBuffer != null) + { + Debug.Assert(_syncBufferLength > 0); + if (!_owner.TryCompleteRead(_syncBuffer, _syncBufferLength, out readResult, _serialStream)) + { + return false; + } + } + else + { + Span span = _buffer.Span; + Debug.Assert(!span.IsEmpty); + + fixed (byte* bufPtr = &MemoryMarshal.GetReference(span)) + { + if (!_owner.TryCompleteRead(bufPtr, span.Length, out readResult, _serialStream)) + { + return false; + } + } + } + + BytesRead = readResult.BytesRead; + if (readResult.ErrorInfo.Error != Interop.Error.SUCCESS) + { + Exception = Interop.GetIOException(readResult.ErrorInfo); + } + return true; + } + + private void OnCompleted(OnCompletedResult result) + { + if (result == OnCompletedResult.Completed) + { + if (Exception != null) + { + _mrvtsc.SetException(Exception); + } + else + { + _mrvtsc.SetResult(BytesRead); + } + } + else if (result == OnCompletedResult.Canceled) + { + _mrvtsc.SetException(new OperationCanceledException(_cancellationToken)); + } + else + { + Debug.Assert(result == OnCompletedResult.Aborted); + _mrvtsc.SetException(new OperationCanceledException()); + } + } + + ValueTaskSourceStatus IValueTaskSource.GetStatus(short token) + => _mrvtsc.GetStatus(token); + + void IValueTaskSource.OnCompleted(Action continuation, object? state, short token, ValueTaskSourceOnCompletedFlags flags) + => _mrvtsc.OnCompleted(continuation, state, token, flags); + + int IValueTaskSource.GetResult(short token) + { + bool canPool = _mrvtsc.GetStatus(token) != ValueTaskSourceStatus.Canceled; + try + { + return _mrvtsc.GetResult(token); + } + finally + { + if (canPool) + { + _owner.ReturnReadOperation(this); + } + } + } + } + + private sealed unsafe class WriteOperation : IValueTaskSource + { + private readonly SafeSerialDeviceHandle _owner; + private readonly UnixHandleAsyncContext.Operation _callbackOp; + internal Exception? Exception; + private ManualResetValueTaskSourceCore _mrvtsc; + private ReadOnlyMemory _buffer; + private byte* _syncBuffer; + private int _syncRemaining; + private CancellationToken _cancellationToken; + + internal WriteOperation(SafeSerialDeviceHandle owner) + { + _owner = owner; + _callbackOp = UnixHandleAsyncContext.CreateOperation(TryCompleteOperation, OnCompleted); + } + + internal UnixHandleAsyncContext.Operation Operation => _callbackOp; + + internal short Version + => _mrvtsc.Version; + + internal void InitSync(byte* syncBuffer, int syncRemaining) + { + _syncBuffer = syncBuffer; + _syncRemaining = syncRemaining; + } + + internal void InitAsync(ReadOnlyMemory buffer, CancellationToken cancellationToken) + { + _buffer = buffer; + _cancellationToken = cancellationToken; + } + + internal void Reset() + { + _buffer = default; + _syncBuffer = null; + _cancellationToken = default; + Exception = null; + _mrvtsc.Reset(); + } + + private bool TryCompleteOperation(SafeHandle handle) + { + Interop.ErrorInfo errorInfo; + + if (_syncBuffer != null) + { + Debug.Assert(_syncRemaining > 0); + + if (!_owner.TryCompleteWrite(_syncBuffer, _syncRemaining, out int bytesWritten, out errorInfo)) + { + _syncBuffer += bytesWritten; + _syncRemaining -= bytesWritten; + return false; + } + } + else + { + ReadOnlySpan span = _buffer.Span; + Debug.Assert(!span.IsEmpty); + + fixed (byte* bufPtr = &MemoryMarshal.GetReference(span)) + { + if (!_owner.TryCompleteWrite(bufPtr, span.Length, out int bytesWritten, out errorInfo)) + { + _buffer = _buffer.Slice(bytesWritten); + return false; + } + } + } + + if (errorInfo.Error != Interop.Error.SUCCESS) + { + Exception = Interop.GetIOException(errorInfo); + } + return true; + } + + private void OnCompleted(OnCompletedResult result) + { + if (result == OnCompletedResult.Completed) + { + if (Exception != null) + { + _mrvtsc.SetException(Exception); + } + else + { + _mrvtsc.SetResult(default); + } + } + else if (result == OnCompletedResult.Canceled) + { + _mrvtsc.SetException(new OperationCanceledException(_cancellationToken)); + } + else + { + Debug.Assert(result == OnCompletedResult.Aborted); + _mrvtsc.SetException(new OperationCanceledException()); + } + } + + ValueTaskSourceStatus IValueTaskSource.GetStatus(short token) + => _mrvtsc.GetStatus(token); + + void IValueTaskSource.OnCompleted(Action continuation, object? state, short token, ValueTaskSourceOnCompletedFlags flags) + => _mrvtsc.OnCompleted(continuation, state, token, flags); + + void IValueTaskSource.GetResult(short token) + { + bool canPool = _mrvtsc.GetStatus(token) != ValueTaskSourceStatus.Canceled; + try + { + _mrvtsc.GetResult(token); + } + finally + { + if (canPool) + { + _owner.ReturnWriteOperation(this); + } + } + } + } + + private sealed class WaitReceiveThresholdOperation + { + private readonly SafeSerialDeviceHandle _owner; + private readonly UnixHandleAsyncContext.Operation _callbackOp; + private SerialStream? _serialStream; + internal int BytesAvailable; + internal ReceiveThresholdResult ThresholdResult; + internal bool HasCompleted; + + internal WaitReceiveThresholdOperation(SafeSerialDeviceHandle owner) + { + _owner = owner; + _callbackOp = UnixHandleAsyncContext.CreateOperation(TryCompleteOperation, OnCompleted); + } + + internal UnixHandleAsyncContext.Operation Operation => _callbackOp; + + internal void Init(SerialStream serialStream) + { + _serialStream = serialStream; + } + + internal void Reset() + { + _serialStream = null; + } + + private bool TryCompleteOperation(SafeHandle handle) + { + SerialStream? serialStream = _serialStream; + if (serialStream == null) + { + return true; + } + + HasCompleted = _owner.TryCompleteWaitReceiveThreshold(serialStream.ReceivedBytesThreshold, out BytesAvailable, out ThresholdResult); + + // We'll enqueue at the end of the read queue in OnCompleted if we haven't completed. + // This enables other operations to read data instead of getting blocked on this operation waiting to reach the threshold. + return true; + } + + internal void CompleteOperation(OnCompletedResult result) + => OnCompleted(result); + + private void OnCompleted(OnCompletedResult result) + { + if (result == OnCompletedResult.Completed) + { + SerialStream? serialStream = _serialStream; + int bytesAvailable = BytesAvailable; + ReceiveThresholdResult thresholdResult = ThresholdResult; + _owner.ReturnWaitReceiveThresholdOperation(this); + if (HasCompleted) + { + serialStream?.OnReceiveThreshold(bytesAvailable, thresholdResult); + } + else + { + // We haven't met the threshold, wait again. + _owner.WaitForReceiveThreshold(serialStream); + } + } + } + } + } +} diff --git a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialPort.cs b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialPort.cs index db12734b0b2192..904dccbdbb1f74 100644 --- a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialPort.cs +++ b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialPort.cs @@ -199,7 +199,7 @@ public int BytesToRead { if (!IsOpen) throw new InvalidOperationException(SR.Port_not_open); - return _internalSerialStream.BytesToRead + CachedBytesToRead; // count the number of bytes we have in the internal buffer too. + return _internalSerialStream.GetBytesToRead(CachedBytesToRead); // count the number of bytes we have in the internal buffer too. } } @@ -458,6 +458,8 @@ public int ReceivedBytesThreshold if (IsOpen) { + _internalSerialStream.ReceivedBytesThreshold = value; + // fake the call to our event handler in case the threshold has been set lower // than how many bytes we currently have. SerialDataReceivedEventArgs args = new SerialDataReceivedEventArgs(SerialData.Chars); @@ -640,6 +642,8 @@ public void Open() _internalSerialStream.PinChanged += _pinChangedHandler; } + _internalSerialStream.ReceivedBytesThreshold = _receivedBytesThreshold; + if (_dataReceived != null) { _internalSerialStream.DataReceived += _dataReceivedHandler; @@ -1270,6 +1274,7 @@ private void CatchReceivedEvents(object src, SerialDataReceivedEventArgs e) if ((eventHandler != null) && (stream != null)) { + bool raiseSkipped = false; lock (stream) { // SerialStream might be closed between the time the event runner @@ -1280,7 +1285,12 @@ private void CatchReceivedEvents(object src, SerialDataReceivedEventArgs e) bool raiseEvent = false; try { - raiseEvent = stream.IsOpen && (SerialData.Eof == e.EventType || BytesToRead >= _receivedBytesThreshold); + if (!stream.IsOpen) + { + return; + } + raiseSkipped = SerialData.Chars == e.EventType && stream.GetBytesToRead(CachedBytesToRead, throwOnDispose: false) < _receivedBytesThreshold; + raiseEvent = SerialData.Eof == e.EventType || !raiseSkipped; } catch { @@ -1288,16 +1298,15 @@ private void CatchReceivedEvents(object src, SerialDataReceivedEventArgs e) } finally { - // ISSUE: This should be fired only when it wasn't already fired for the total number of bytes available - // Similarly as done in SerialStream.Linux (IOLoop) - // I.e: Let _receivedBytesThreshold be 8 - when we get an event when 7 bytes are available - // BytesToRead can change while we run this event and thus - // we virtually can get 2 events when 8th byte arrives - // I.e. we might want to add total bytes available as internal field in the args event if (raiseEvent) eventHandler(this, e); // here, do your reading, etc. } } + + if (raiseSkipped) + { + stream.OnRaiseCharsEventSkipped(); + } } } diff --git a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.Unix.cs b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.Unix.cs index d220cf95a832c1..eb6ee4acbbb91d 100644 --- a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.Unix.cs +++ b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.Unix.cs @@ -2,8 +2,6 @@ // The .NET Foundation licenses this file to you under the MIT license. using System; -using System.Collections; -using System.Collections.Generic; using System.Diagnostics; using System.IO.Ports; using System.Runtime.InteropServices; @@ -18,7 +16,7 @@ internal sealed partial class SerialStream : Stream private const int TimeoutResolution = 30; // time [ms] loop has to be idle before it stops private const int IOLoopIdleTimeout = 2000; - private bool _ioLoopFinished; + private volatile bool _disposed; private SafeSerialDeviceHandle _handle; private int _baudRate; @@ -31,20 +29,7 @@ internal sealed partial class SerialStream : Stream private readonly byte[] _tempBuf = new byte[1]; private Task _ioLoop; private readonly object _ioLoopLock = new object(); - // Use a Queue with locking instead of ConcurrentQueue because ConcurrentQueue preserves segments for - // observation when using TryPeek(). These segments will not clear out references after a dequeue - // and as a result they hold on to SerialStreamIORequest instances so that they cannot be GC'ed. - // This in turn means that any buffers that the client supplied are not eligible for GC either. - private readonly Queue _readQueue = new(); - private readonly object _readQueueLock = new(); - private readonly Queue _writeQueue = new(); - private readonly object _writeQueueLock = new(); - - private long _totalBytesRead; - private long TotalBytesAvailable => _totalBytesRead + BytesToRead; - private long _lastTotalBytesAvailable; - - // called when one character is received. + private SerialDataReceivedEventHandler _dataReceived; internal event SerialDataReceivedEventHandler DataReceived { @@ -55,7 +40,7 @@ internal event SerialDataReceivedEventHandler DataReceived if (wasNull) { - EnsureIOLoopRunning(); + DataReceiveEnable(); } } remove @@ -168,9 +153,18 @@ internal int BytesToWrite get { return Interop.Termios.TermiosGetAvailableBytes(_handle, false); } } - internal int BytesToRead + internal int GetBytesToRead(int buffered, bool throwOnDispose = true) { - get { return Interop.Termios.TermiosGetAvailableBytes(_handle, true); } + SafeSerialDeviceHandle? handle = _handle; + if (handle == null || !IsOpen) + { + if (throwOnDispose) + { + InternalResources.FileNotOpen(); + } + return 0; + } + return handle.GetBytesToRead(buffered, throwOnDispose); } internal bool CDHolding @@ -359,19 +353,6 @@ internal byte ParityReplace } #pragma warning restore CA1822 - private bool HasCancelledTasksToProcess - { - get => Volatile.Read(ref field); - set => Volatile.Write(ref field, value); - } - - internal void DiscardInBuffer() - { - if (_handle == null) InternalResources.FileNotOpen(); - // This may or may not work depending on hardware. - Interop.Termios.TermiosDiscard(_handle, Interop.Termios.Queue.ReceiveQueue); - } - internal void DiscardOutBuffer() { if (_handle == null) InternalResources.FileNotOpen(); @@ -399,11 +380,7 @@ public override void Flush() { if (_handle == null) InternalResources.FileNotOpen(); - SpinWait sw = default; - while (!IsWriteQueueEmpty()) - { - sw.SpinOnce(); - } + FlushWrites(); Interop.Termios.TermiosDrain(_handle); } @@ -419,173 +396,33 @@ public override int Read(byte[] array, int offset, int count) return Read(array, offset, count, ReadTimeout); } - internal int Read(byte[] array, int offset, int count, int timeout) - { - using (CancellationTokenSource cts = GetCancellationTokenSourceFromTimeout(timeout)) - { - Task t = ReadAsync(array, offset, count, cts?.Token ?? CancellationToken.None); - - try - { - return t.GetAwaiter().GetResult(); - } - catch (OperationCanceledException) - { - throw new TimeoutException(); - } - } - } - public override int EndRead(IAsyncResult asyncResult) - => EndReadWrite(asyncResult); - - public override Task ReadAsync(byte[] array, int offset, int count, CancellationToken cancellationToken) - { - CheckReadWriteArguments(array, offset, count); - - if (count == 0) - return Task.FromResult(0); // return immediately if no bytes requested; no need for overhead. - - Memory buffer = new Memory(array, offset, count); - SerialStreamReadRequest result = new SerialStreamReadRequest(this, cancellationToken, buffer); - lock (_readQueueLock) - { - _readQueue.Enqueue(result); - } - - EnsureIOLoopRunning(); - - return result.Task; - } - -#if !NETFRAMEWORK && !NETSTANDARD2_0 - public override ValueTask ReadAsync(Memory buffer, CancellationToken cancellationToken = default) - { - CheckHandle(); - - if (buffer.IsEmpty) - return new ValueTask(0); - - SerialStreamReadRequest result = new SerialStreamReadRequest(this, cancellationToken, buffer); - lock (_readQueueLock) - { - _readQueue.Enqueue(result); - } - - EnsureIOLoopRunning(); - - return new ValueTask(result.Task); - } -#endif - - public override Task WriteAsync(byte[] array, int offset, int count, CancellationToken cancellationToken) { - CheckWriteArguments(array, offset, count); - - if (count == 0) - return Task.CompletedTask; // return immediately if no bytes to write; no need for overhead. - - ReadOnlyMemory buffer = new ReadOnlyMemory(array, offset, count); - SerialStreamWriteRequest result = new SerialStreamWriteRequest(this, cancellationToken, buffer); - lock (_writeQueueLock) + try { - _writeQueue.Enqueue(result); + return TaskToAsyncResult.End(asyncResult); } - - EnsureIOLoopRunning(); - - return result.Task; - } - -#if !NETFRAMEWORK && !NETSTANDARD2_0 - public override ValueTask WriteAsync(ReadOnlyMemory buffer, CancellationToken cancellationToken = default) - { - CheckWriteArguments(); - - if (buffer.IsEmpty) - return ValueTask.CompletedTask; // return immediately if no bytes to write; no need for overhead. - - SerialStreamWriteRequest result = new SerialStreamWriteRequest(this, cancellationToken, buffer); - lock (_writeQueueLock) + catch (OperationCanceledException) { - _writeQueue.Enqueue(result); + throw new TimeoutException(); } - - EnsureIOLoopRunning(); - - return new ValueTask(result.Task); } -#endif public override IAsyncResult BeginRead(byte[] array, int offset, int numBytes, AsyncCallback userCallback, object stateObject) { return TaskToAsyncResult.Begin(ReadAsync(array, offset, numBytes), userCallback, stateObject); } - // Will wait `timeout` miliseconds or until reading or writing is possible - // If no operation is requested it will throw - // Returns event which has happened - private Interop.PollEvents PollEvents(int timeout, bool pollReadEvents, bool pollWriteEvents, out Interop.ErrorInfo? error) - { - if (!pollReadEvents && !pollWriteEvents) - { - Debug.Fail("This should not happen"); - throw new Exception(); - } - - Interop.PollEvents eventsToPoll = Interop.PollEvents.POLLERR; - - if (pollReadEvents) - { - eventsToPoll |= Interop.PollEvents.POLLIN; - } - - if (pollWriteEvents) - { - eventsToPoll |= Interop.PollEvents.POLLOUT; - } - - Interop.PollEvents events; - Interop.Error ret = Interop.Serial.Poll( - _handle, - eventsToPoll, - timeout, - out events); - - error = ret != Interop.Error.SUCCESS ? Interop.Sys.GetLastErrorInfo() : (Interop.ErrorInfo?)null; - return events; - } - - internal void Write(byte[] array, int offset, int count, int timeout) - { - using (CancellationTokenSource cts = GetCancellationTokenSourceFromTimeout(timeout)) - { - Task t = WriteAsync(array, offset, count, cts?.Token ?? CancellationToken.None); - - try - { - t.GetAwaiter().GetResult(); - } - catch (OperationCanceledException) - { - throw new TimeoutException(); - } - } - } - public override IAsyncResult BeginWrite(byte[] array, int offset, int count, AsyncCallback userCallback, object stateObject) { return TaskToAsyncResult.Begin(WriteAsync(array, offset, count), userCallback, stateObject); } public override void EndWrite(IAsyncResult asyncResult) - => EndReadWrite(asyncResult); - - private static int EndReadWrite(IAsyncResult asyncResult) { try { - return TaskToAsyncResult.End(asyncResult); + TaskToAsyncResult.End(asyncResult); } catch (OperationCanceledException) { @@ -667,9 +504,7 @@ internal SerialStream(string portName, int baudRate, Parity parity, int dataBits throw; } - _processReadDelegate = ProcessRead; - _processWriteDelegate = ProcessWrite; - _lastTotalBytesAvailable = TotalBytesAvailable; + OnCtor(); } private void EnsureIOLoopRunning() @@ -688,32 +523,9 @@ private void EnsureIOLoopRunning() } } - private void FinishPendingIORequests(Interop.ErrorInfo? error = null) - { - lock (_readQueueLock) - { - while (_readQueue.TryDequeue(out SerialStreamIORequest r)) - { - r.Complete(error.HasValue ? - Interop.GetIOException(error.Value) : - InternalResources.FileNotOpenException()); - } - } - - lock (_writeQueueLock) - { - while (_writeQueue.TryDequeue(out SerialStreamIORequest r)) - { - r.Complete(error.HasValue ? - Interop.GetIOException(error.Value) : - InternalResources.FileNotOpenException()); - } - } - } - protected override void Dispose(bool disposing) { - _ioLoopFinished = true; + _disposed = true; if (disposing) { @@ -769,276 +581,6 @@ private void RaiseDataReceivedEof() } } - // should return non-negative integer meaning numbers of bytes read/written (0 for errors) - private delegate int RequestProcessor(SerialStreamIORequest r); - private readonly RequestProcessor _processReadDelegate; - private readonly RequestProcessor _processWriteDelegate; - - private unsafe int ProcessRead(SerialStreamIORequest r) - { - SerialStreamReadRequest readRequest = (SerialStreamReadRequest)r; - Span buff = readRequest.Buffer.Span; - fixed (byte* bufPtr = buff) - { - // assumes dequeue-ing happens on a single thread - int numBytes = Interop.Serial.Read(_handle, bufPtr, buff.Length); - - if (numBytes < 0) - { - Interop.ErrorInfo lastError = Interop.Sys.GetLastErrorInfo(); - - // ignore EWOULDBLOCK since we handle timeout elsewhere - if (lastError.Error != Interop.Error.EWOULDBLOCK) - { - readRequest.Complete(Interop.GetIOException(lastError)); - } - } - else if (numBytes > 0) - { - readRequest.Complete(numBytes); - return numBytes; - } - else // numBytes == 0 - { - RaiseDataReceivedEof(); - } - } - - return 0; - } - - private unsafe int ProcessWrite(SerialStreamIORequest r) - { - SerialStreamWriteRequest writeRequest = (SerialStreamWriteRequest)r; - ReadOnlySpan buff = writeRequest.Buffer.Span; - fixed (byte* bufPtr = buff) - { - // assumes dequeue-ing happens on a single thread - int numBytes = Interop.Serial.Write(_handle, bufPtr, buff.Length); - - if (numBytes <= 0) - { - Interop.ErrorInfo lastError = Interop.Sys.GetLastErrorInfo(); - - // ignore EWOULDBLOCK since we handle timeout elsewhere - // numBytes == 0 means that there might be an error - if (lastError.Error != Interop.Error.SUCCESS && lastError.Error != Interop.Error.EWOULDBLOCK) - { - r.Complete(Interop.GetIOException(lastError)); - } - } - else - { - writeRequest.ProcessBytes(numBytes); - - if (writeRequest.Buffer.Length == 0) - { - writeRequest.Complete(); - } - - return numBytes; - } - } - - return 0; - } - - // returns number of bytes read/written - private static int DoIORequest(Queue q, object queueLock, RequestProcessor op) - { - // assumes dequeue-ing happens on a single thread - while (TryPeekNextRequest(out SerialStreamIORequest r)) - { - int ret = op(r); - Debug.Assert(ret >= 0); - - if (r.IsCompleted) - { - lock (queueLock) - { - q.TryDequeue(out _); - } - } - - return ret; - } - - return 0; - - bool TryPeekNextRequest(out SerialStreamIORequest r) - { - lock (queueLock) - { - while (q.TryPeek(out r)) - { - if (!r.IsCompleted) - { - return true; - } - q.TryDequeue(out _); - } - } - r = default; - return false; - } - } - - private void IOLoop() - { - bool eofReceived = false; - // we do not care about bytes we got before - only about changes - // loop just got started which means we just got request - bool lastIsIdle = false; - int ticksWhenIdleStarted = 0; - - Signals lastSignals = _pinChanged != null ? Interop.Termios.TermiosGetAllSignals(_handle) : Signals.Error; - - bool IsNoEventRegistered() => _dataReceived == null && _pinChanged == null; - - while (IsOpen && !eofReceived && !_ioLoopFinished) - { - if (HasCancelledTasksToProcess) - { - HasCancelledTasksToProcess = false; - RemoveCompletedTasks(_readQueue, _readQueueLock); - RemoveCompletedTasks(_writeQueue, _writeQueueLock); - } - - bool hasPendingReads = !IsReadQueueEmpty(); - bool hasPendingWrites = !IsWriteQueueEmpty(); - - bool hasPendingIO = hasPendingReads || hasPendingWrites; - bool isIdle = IsNoEventRegistered() && !hasPendingIO; - - if (!hasPendingIO) - { - if (isIdle) - { - if (!lastIsIdle) - { - // we've just started idling - ticksWhenIdleStarted = Environment.TickCount; - } - else if (Environment.TickCount - ticksWhenIdleStarted > IOLoopIdleTimeout) - { - // we are already idling for a while - // let's stop the loop until there is some work to do - - lock (_ioLoopLock) - { - // double check we are done under lock - if (IsNoEventRegistered() && IsReadQueueEmpty() && IsWriteQueueEmpty()) - { - _ioLoop = null; - break; - } - else - { - // to make sure timer restarts - lastIsIdle = false; - continue; - } - } - } - } - - Thread.Sleep(1); - } - else - { - Interop.PollEvents events = PollEvents(1, - pollReadEvents: hasPendingReads, - pollWriteEvents: hasPendingWrites, - out Interop.ErrorInfo? error); - - if (error.HasValue) - { - FinishPendingIORequests(error); - break; - } - - if (events.HasFlag(Interop.PollEvents.POLLNVAL) || - events.HasFlag(Interop.PollEvents.POLLERR)) - { - // bad descriptor or some other error we can't handle - FinishPendingIORequests(); - break; - } - - if (events.HasFlag(Interop.PollEvents.POLLIN)) - { - int bytesRead = DoIORequest(_readQueue, _readQueueLock, _processReadDelegate); - _totalBytesRead += bytesRead; - } - - if (events.HasFlag(Interop.PollEvents.POLLOUT)) - { - DoIORequest(_writeQueue, _writeQueueLock, _processWriteDelegate); - } - } - - // check if there is any new data (either already read or in the driver input) - // this event is private and handled inside of SerialPort - // which then throttles it with the threshold - long totalBytesAvailable = TotalBytesAvailable; - if (totalBytesAvailable > _lastTotalBytesAvailable) - { - _lastTotalBytesAvailable = totalBytesAvailable; - RaiseDataReceivedChars(); - } - - if (_pinChanged != null) - { - // Checking for changes could technically speaking be done by waiting with ioctl+TIOCMIWAIT - // This would require spinning new thread and also would potentially trigger events when - // user didn't have time to respond. - // Diffing seems like a better solution. - Signals current = Interop.Termios.TermiosGetAllSignals(_handle); - - // There is no really good action we can take when this errors so just ignore - // a sinle event. - if (current != Signals.Error && lastSignals != Signals.Error) - { - Signals changed = current ^ lastSignals; - if (changed != Signals.None) - { - NotifyPinChanges(changed); - } - } - - lastSignals = current; - } - - lastIsIdle = isIdle; - } - } - - private static void RemoveCompletedTasks(Queue queue, object queueLock) - { - // assumes dequeue-ing happens on a single thread - lock (queueLock) - { - while (queue.TryPeek(out var r) && r.IsCompleted) - queue.TryDequeue(out _); - } - } - - private bool IsReadQueueEmpty() - { - lock (_readQueueLock) - { - return _readQueue.Count == 0; - } - } - - private bool IsWriteQueueEmpty() - { - lock (_writeQueueLock) - { - return _writeQueue.Count == 0; - } - } - private void NotifyPinChanges(Signals signals) { if (signals.HasFlag(Signals.SignalCts)) @@ -1062,80 +604,9 @@ private void NotifyPinChanges(Signals signals) } } - private static CancellationTokenSource GetCancellationTokenSourceFromTimeout(int timeoutMs) - { - return timeoutMs == SerialPort.InfiniteTimeout ? - null : - new CancellationTokenSource(Math.Max(timeoutMs, TimeoutResolution)); - } - private static Exception GetLastIOError() { return Interop.GetIOException(Interop.Sys.GetLastErrorInfo()); } - - private abstract class SerialStreamIORequest : TaskCompletionSource - { - public bool IsCompleted => Task.IsCompleted; - private readonly SerialStream _parent; - private readonly CancellationTokenRegistration _cancellationTokenRegistration; - - protected SerialStreamIORequest(SerialStream parent, CancellationToken ct) - : base(TaskCreationOptions.RunContinuationsAsynchronously) - { - _parent = parent; - _cancellationTokenRegistration = ct.Register(s => - { - var request = (SerialStreamIORequest)s; - request.TrySetCanceled(); - request._parent.HasCancelledTasksToProcess = true; - }, this); - } - - internal void Complete(int numBytes) - { - TrySetResult(numBytes); - _cancellationTokenRegistration.Dispose(); - } - - internal void Complete(Exception exception) - { - TrySetException(exception); - _cancellationTokenRegistration.Dispose(); - } - } - - private sealed class SerialStreamReadRequest : SerialStreamIORequest - { - public Memory Buffer { get; } - - public SerialStreamReadRequest(SerialStream parent, CancellationToken ct, Memory buffer) - : base(parent, ct) - { - Buffer = buffer; - } - } - - private sealed class SerialStreamWriteRequest : SerialStreamIORequest - { - public ReadOnlyMemory Buffer { get; private set; } - - public SerialStreamWriteRequest(SerialStream parent, CancellationToken ct, ReadOnlyMemory buffer) - : base(parent, ct) - { - Buffer = buffer; - } - - internal void Complete() - { - Debug.Assert(Buffer.Length == 0); - Complete(Buffer.Length); - } - - internal void ProcessBytes(int numBytes) - { - Buffer = Buffer.Slice(numBytes); - } - } } } diff --git a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.UnixAsyncContext.cs b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.UnixAsyncContext.cs new file mode 100644 index 00000000000000..3bb2e38ac4238d --- /dev/null +++ b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.UnixAsyncContext.cs @@ -0,0 +1,239 @@ +// Licensed to the .NET Foundation under one or more agreements. +// The .NET Foundation licenses this file to you under the MIT license. + +using System; +using System.Diagnostics; +using System.IO.Ports; +using System.Threading; +using System.Threading.Tasks; +using Signals = Interop.Termios.Signals; + +namespace System.IO.Ports +{ + internal sealed partial class SerialStream : Stream + { + // Track how many bytes were available when we raised DataReceived. + // We use this to track how many bytes were read and start watching for new data when the user consumed enough. + private volatile int _cachedBytesAvailable; + private int _waitingForReceiveThreshold; + + internal int Read(byte[] array, int offset, int count, int timeout) + { + CheckReadWriteArguments(array, offset, count); + + if (count == 0) + return 0; + + return _handle.Read(new Span(array, offset, count), MapTimeout(timeout), this); + } + + public override Task ReadAsync(byte[] array, int offset, int count, CancellationToken cancellationToken) + { + CheckReadWriteArguments(array, offset, count); + + if (count == 0) + return Task.FromResult(0); + + return _handle.ReadAsync(new Memory(array, offset, count), cancellationToken, this).AsTask(); + } + + public override ValueTask ReadAsync(Memory buffer, CancellationToken cancellationToken = default) + { + CheckHandle(); + + if (buffer.IsEmpty) + return new ValueTask(0); + + return _handle.ReadAsync(buffer, cancellationToken, this); + } + + public override Task WriteAsync(byte[] array, int offset, int count, CancellationToken cancellationToken) + { + CheckWriteArguments(array, offset, count); + + if (count == 0) + return Task.CompletedTask; + + return _handle.WriteAsync(new ReadOnlyMemory(array, offset, count), cancellationToken).AsTask(); + } + + public override ValueTask WriteAsync(ReadOnlyMemory buffer, CancellationToken cancellationToken = default) + { + CheckWriteArguments(); + + if (buffer.IsEmpty) + return ValueTask.CompletedTask; + + return _handle.WriteAsync(buffer, cancellationToken); + } + + internal void Write(byte[] array, int offset, int count, int timeout) + { + CheckWriteArguments(array, offset, count); + + if (count == 0) + return; + + _handle.Write(new ReadOnlySpan(array, offset, count), MapTimeout(timeout)); + } + + private void DataReceiveEnable() + => EnsureWaitForReceiveThreshold(); + + private void FlushWrites() + { + // A zero-byte write doesn't actually write, but it won't complete until the preceding writes are handled. + _handle.Write(ReadOnlySpan.Empty, Timeout.Infinite); + } + +#pragma warning disable CA1822 + private void OnCtor() { } + private void FinishPendingIORequests() { } +#pragma warning restore CA1822 + + internal void DiscardInBuffer() + { + if (_handle == null) InternalResources.FileNotOpen(); + + _handle.DiscardInBuffer(); + + // Arming readable by reading (via OnBytesRead) might not happen due to clearing the readable data. + // Ensure we're armed. + if (_dataReceived != null) + { + EnsureWaitForReceiveThreshold(); + } + } + + internal void OnBytesRead(int bytesRead) + { + // Start watching for new data received event when the user consumed the previous data. + Debug.Assert(bytesRead > 0); + + int remaining = Interlocked.Add(ref _cachedBytesAvailable, -bytesRead); + // Clamp to zero: direct reads (without a prior OnReceiveThreshold) can make this negative. + if (remaining < 0) + { + _cachedBytesAvailable = 0; + } + + if (remaining < ReceivedBytesThreshold && _dataReceived != null) + { + EnsureWaitForReceiveThreshold(); + } + } + + private void EnsureWaitForReceiveThreshold() + { + SafeSerialDeviceHandle? handle = _handle; + if (handle != null && Interlocked.CompareExchange(ref _waitingForReceiveThreshold, 1, 0) == 0) + { + handle.WaitForReceiveThreshold(this); + } + } + + internal void OnReceiveThreshold(int bytesAvailable, SafeSerialDeviceHandle.ReceiveThresholdResult result) + { + Volatile.Write(ref _waitingForReceiveThreshold, 0); + + if (_dataReceived == null) + { + return; + } + + _cachedBytesAvailable = bytesAvailable; + + if (result == SafeSerialDeviceHandle.ReceiveThresholdResult.Success) + { + RaiseDataReceivedChars(); + } + // Emit Eof for errors too so that the user gets "an" event to trigger a read. + else if (result is SafeSerialDeviceHandle.ReceiveThresholdResult.Eof or SafeSerialDeviceHandle.ReceiveThresholdResult.Error) + { + RaiseDataReceivedEof(); + } + else if (result == SafeSerialDeviceHandle.ReceiveThresholdResult.Disposed) + { } + else + { + OnRaiseCharsEventSkipped(); + } + } + + internal void OnRaiseCharsEventSkipped() + { + // Unlikely: we didn't meet the threshold because it changed. + EnsureWaitForReceiveThreshold(); + } + + // This is the same loop as UnixPollLoop implementation, but limited to the pin polling. + private void IOLoop() + { + Signals lastSignals = Interop.Termios.TermiosGetAllSignals(_handle); + bool lastIsIdle = false; + int ticksWhenIdleStarted = 0; + + while (IsOpen && !_disposed) + { + bool isIdle = _pinChanged == null; + + if (isIdle) + { + if (!lastIsIdle) + { + ticksWhenIdleStarted = Environment.TickCount; + } + else if (Environment.TickCount - ticksWhenIdleStarted > IOLoopIdleTimeout) + { + lock (_ioLoopLock) + { + if (_pinChanged == null) + { + _ioLoop = null; + break; + } + else + { + lastIsIdle = false; + continue; + } + } + } + } + + Thread.Sleep(1); + + if (_pinChanged != null) + { + // Checking for changes could technically speaking be done by waiting with ioctl+TIOCMIWAIT + // This would require spinning new thread and also would potentially trigger events when + // user didn't have time to respond. + // Diffing seems like a better solution. + Signals current = Interop.Termios.TermiosGetAllSignals(_handle); + + // There is no really good action we can take when this errors so just ignore + // a sinle event. + if (current != Signals.Error && lastSignals != Signals.Error) + { + Signals changed = current ^ lastSignals; + if (changed != Signals.None) + { + NotifyPinChanges(changed); + } + } + + lastSignals = current; + } + + lastIsIdle = isIdle; + } + } + + private static int MapTimeout(int timeoutMs) + { + // SerialPort.InfiniteTimeout is -1, which maps to infinite for ReadSync/WriteSync. + // Positive values pass through. + return timeoutMs == SerialPort.InfiniteTimeout ? -1 : Math.Max(timeoutMs, TimeoutResolution); + } + } +} diff --git a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.UnixPollLoop.cs b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.UnixPollLoop.cs new file mode 100644 index 00000000000000..3a9c64fa44cfd4 --- /dev/null +++ b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.UnixPollLoop.cs @@ -0,0 +1,579 @@ +// Licensed to the .NET Foundation under one or more agreements. +// The .NET Foundation licenses this file to you under the MIT license. + +using System; +using System.Collections.Generic; +using System.Diagnostics; +using System.IO.Ports; +using System.Runtime.InteropServices; +using System.Threading; +using System.Threading.Tasks; +using Signals = Interop.Termios.Signals; + +namespace System.IO.Ports +{ + internal sealed partial class SerialStream : Stream + { + // Use a Queue with locking instead of ConcurrentQueue because ConcurrentQueue preserves segments for + // observation when using TryPeek(). These segments will not clear out references after a dequeue + // and as a result they hold on to SerialStreamIORequest instances so that they cannot be GC'ed. + // This in turn means that any buffers that the client supplied are not eligible for GC either. + private readonly Queue _readQueue = new(); + private readonly object _readQueueLock = new(); + private readonly Queue _writeQueue = new(); + private readonly object _writeQueueLock = new(); + + private long _totalBytesRead; + private long TotalBytesAvailable => _totalBytesRead + GetBytesToRead(buffered: 0); + private long _lastTotalBytesAvailable; + + private void DataReceiveEnable() + { + EnsureIOLoopRunning(); + } + + private bool HasCancelledTasksToProcess + { + get => Volatile.Read(ref field); + set => Volatile.Write(ref field, value); + } + internal void DiscardInBuffer() + { + if (_handle == null) InternalResources.FileNotOpen(); + // This may or may not work depending on hardware. + Interop.Termios.TermiosDiscard(_handle, Interop.Termios.Queue.ReceiveQueue); + } + + private void FlushWrites() + { + SpinWait sw = default; + while (!IsWriteQueueEmpty()) + { + sw.SpinOnce(); + } + } + + internal int Read(byte[] array, int offset, int count, int timeout) + { + using (CancellationTokenSource cts = GetCancellationTokenSourceFromTimeout(timeout)) + { + Task t = ReadAsync(array, offset, count, cts?.Token ?? CancellationToken.None); + + try + { + return t.GetAwaiter().GetResult(); + } + catch (OperationCanceledException) + { + throw new TimeoutException(); + } + } + } + + public override Task ReadAsync(byte[] array, int offset, int count, CancellationToken cancellationToken) + { + CheckReadWriteArguments(array, offset, count); + + if (count == 0) + return Task.FromResult(0); // return immediately if no bytes requested; no need for overhead. + + Memory buffer = new Memory(array, offset, count); + SerialStreamReadRequest result = new SerialStreamReadRequest(this, cancellationToken, buffer); + lock (_readQueueLock) + { + _readQueue.Enqueue(result); + } + + EnsureIOLoopRunning(); + + return result.Task; + } + +#if !NETFRAMEWORK && !NETSTANDARD2_0 + public override ValueTask ReadAsync(Memory buffer, CancellationToken cancellationToken = default) + { + CheckHandle(); + + if (buffer.IsEmpty) + return new ValueTask(0); + + SerialStreamReadRequest result = new SerialStreamReadRequest(this, cancellationToken, buffer); + lock (_readQueueLock) + { + _readQueue.Enqueue(result); + } + + EnsureIOLoopRunning(); + + return new ValueTask(result.Task); + } +#endif + + public override Task WriteAsync(byte[] array, int offset, int count, CancellationToken cancellationToken) + { + CheckWriteArguments(array, offset, count); + + if (count == 0) + return Task.CompletedTask; // return immediately if no bytes to write; no need for overhead. + + ReadOnlyMemory buffer = new ReadOnlyMemory(array, offset, count); + SerialStreamWriteRequest result = new SerialStreamWriteRequest(this, cancellationToken, buffer); + lock (_writeQueueLock) + { + _writeQueue.Enqueue(result); + } + + EnsureIOLoopRunning(); + + return result.Task; + } + +#if !NETFRAMEWORK && !NETSTANDARD2_0 + public override ValueTask WriteAsync(ReadOnlyMemory buffer, CancellationToken cancellationToken = default) + { + CheckWriteArguments(); + + if (buffer.IsEmpty) + return ValueTask.CompletedTask; // return immediately if no bytes to write; no need for overhead. + + SerialStreamWriteRequest result = new SerialStreamWriteRequest(this, cancellationToken, buffer); + lock (_writeQueueLock) + { + _writeQueue.Enqueue(result); + } + + EnsureIOLoopRunning(); + + return new ValueTask(result.Task); + } +#endif + + // Will wait `timeout` miliseconds or until reading or writing is possible + // If no operation is requested it will throw + // Returns event which has happened + private Interop.PollEvents PollEvents(int timeout, bool pollReadEvents, bool pollWriteEvents, out Interop.ErrorInfo? error) + { + if (!pollReadEvents && !pollWriteEvents) + { + Debug.Fail("This should not happen"); + throw new Exception(); + } + + Interop.PollEvents eventsToPoll = Interop.PollEvents.POLLERR; + + if (pollReadEvents) + { + eventsToPoll |= Interop.PollEvents.POLLIN; + } + + if (pollWriteEvents) + { + eventsToPoll |= Interop.PollEvents.POLLOUT; + } + + Interop.PollEvents events; + Interop.Error ret = Interop.Serial.Poll( + _handle, + eventsToPoll, + timeout, + out events); + + error = ret != Interop.Error.SUCCESS ? Interop.Sys.GetLastErrorInfo() : (Interop.ErrorInfo?)null; + return events; + } + + internal void Write(byte[] array, int offset, int count, int timeout) + { + using (CancellationTokenSource cts = GetCancellationTokenSourceFromTimeout(timeout)) + { + Task t = WriteAsync(array, offset, count, cts?.Token ?? CancellationToken.None); + + try + { + t.GetAwaiter().GetResult(); + } + catch (OperationCanceledException) + { + throw new TimeoutException(); + } + } + } + + private void OnCtor() + { + _processReadDelegate = ProcessRead; + _processWriteDelegate = ProcessWrite; + _lastTotalBytesAvailable = TotalBytesAvailable; + } + +#pragma warning disable CA1822 + internal void OnRaiseCharsEventSkipped() + { + } +#pragma warning restore CA1822 + + private void FinishPendingIORequests(Interop.ErrorInfo? error = null) + { + lock (_readQueueLock) + { + while (_readQueue.TryDequeue(out SerialStreamIORequest r)) + { + r.Complete(error.HasValue ? + Interop.GetIOException(error.Value) : + InternalResources.FileNotOpenException()); + } + } + + lock (_writeQueueLock) + { + while (_writeQueue.TryDequeue(out SerialStreamIORequest r)) + { + r.Complete(error.HasValue ? + Interop.GetIOException(error.Value) : + InternalResources.FileNotOpenException()); + } + } + } + + // should return non-negative integer meaning numbers of bytes read/written (0 for errors) + private delegate int RequestProcessor(SerialStreamIORequest r); + private RequestProcessor _processReadDelegate; + private RequestProcessor _processWriteDelegate; + + private unsafe int ProcessRead(SerialStreamIORequest r) + { + SerialStreamReadRequest readRequest = (SerialStreamReadRequest)r; + Span buff = readRequest.Buffer.Span; + fixed (byte* bufPtr = buff) + { + // assumes dequeue-ing happens on a single thread + int numBytes = Interop.Serial.Read(_handle, bufPtr, buff.Length); + + if (numBytes < 0) + { + Interop.ErrorInfo lastError = Interop.Sys.GetLastErrorInfo(); + + // ignore EWOULDBLOCK since we handle timeout elsewhere + if (lastError.Error != Interop.Error.EWOULDBLOCK) + { + readRequest.Complete(Interop.GetIOException(lastError)); + } + } + else if (numBytes > 0) + { + readRequest.Complete(numBytes); + return numBytes; + } + else // numBytes == 0 + { + RaiseDataReceivedEof(); + } + } + + return 0; + } + + private unsafe int ProcessWrite(SerialStreamIORequest r) + { + SerialStreamWriteRequest writeRequest = (SerialStreamWriteRequest)r; + ReadOnlySpan buff = writeRequest.Buffer.Span; + fixed (byte* bufPtr = buff) + { + // assumes dequeue-ing happens on a single thread + int numBytes = Interop.Serial.Write(_handle, bufPtr, buff.Length); + + if (numBytes <= 0) + { + Interop.ErrorInfo lastError = Interop.Sys.GetLastErrorInfo(); + + // ignore EWOULDBLOCK since we handle timeout elsewhere + // numBytes == 0 means that there might be an error + if (lastError.Error != Interop.Error.SUCCESS && lastError.Error != Interop.Error.EWOULDBLOCK) + { + r.Complete(Interop.GetIOException(lastError)); + } + } + else + { + writeRequest.ProcessBytes(numBytes); + + if (writeRequest.Buffer.Length == 0) + { + writeRequest.Complete(); + } + + return numBytes; + } + } + + return 0; + } + + // returns number of bytes read/written + private static int DoIORequest(Queue q, object queueLock, RequestProcessor op) + { + // assumes dequeue-ing happens on a single thread + while (TryPeekNextRequest(out SerialStreamIORequest r)) + { + int ret = op(r); + Debug.Assert(ret >= 0); + + if (r.IsCompleted) + { + lock (queueLock) + { + q.TryDequeue(out _); + } + } + + return ret; + } + + return 0; + + bool TryPeekNextRequest(out SerialStreamIORequest r) + { + lock (queueLock) + { + while (q.TryPeek(out r)) + { + if (!r.IsCompleted) + { + return true; + } + q.TryDequeue(out _); + } + } + r = default; + return false; + } + } + + private void IOLoop() + { + bool eofReceived = false; + // we do not care about bytes we got before - only about changes + // loop just got started which means we just got request + bool lastIsIdle = false; + int ticksWhenIdleStarted = 0; + + Signals lastSignals = _pinChanged != null ? Interop.Termios.TermiosGetAllSignals(_handle) : Signals.Error; + + bool IsNoEventRegistered() => _dataReceived == null && _pinChanged == null; + + while (IsOpen && !eofReceived && !_disposed) + { + if (HasCancelledTasksToProcess) + { + HasCancelledTasksToProcess = false; + RemoveCompletedTasks(_readQueue, _readQueueLock); + RemoveCompletedTasks(_writeQueue, _writeQueueLock); + } + + bool hasPendingReads = !IsReadQueueEmpty(); + bool hasPendingWrites = !IsWriteQueueEmpty(); + + bool hasPendingIO = hasPendingReads || hasPendingWrites; + bool isIdle = IsNoEventRegistered() && !hasPendingIO; + + if (!hasPendingIO) + { + if (isIdle) + { + if (!lastIsIdle) + { + // we've just started idling + ticksWhenIdleStarted = Environment.TickCount; + } + else if (Environment.TickCount - ticksWhenIdleStarted > IOLoopIdleTimeout) + { + // we are already idling for a while + // let's stop the loop until there is some work to do + + lock (_ioLoopLock) + { + // double check we are done under lock + if (IsNoEventRegistered() && IsReadQueueEmpty() && IsWriteQueueEmpty()) + { + _ioLoop = null; + break; + } + else + { + // to make sure timer restarts + lastIsIdle = false; + continue; + } + } + } + } + + Thread.Sleep(1); + } + else + { + Interop.PollEvents events = PollEvents(1, + pollReadEvents: hasPendingReads, + pollWriteEvents: hasPendingWrites, + out Interop.ErrorInfo? error); + + if (error.HasValue) + { + FinishPendingIORequests(error); + break; + } + + if (events.HasFlag(Interop.PollEvents.POLLNVAL) || + events.HasFlag(Interop.PollEvents.POLLERR)) + { + // bad descriptor or some other error we can't handle + FinishPendingIORequests(); + break; + } + + if (events.HasFlag(Interop.PollEvents.POLLIN)) + { + int bytesRead = DoIORequest(_readQueue, _readQueueLock, _processReadDelegate); + _totalBytesRead += bytesRead; + } + + if (events.HasFlag(Interop.PollEvents.POLLOUT)) + { + DoIORequest(_writeQueue, _writeQueueLock, _processWriteDelegate); + } + } + + // check if there is any new data (either already read or in the driver input) + // this event is private and handled inside of SerialPort + // which then throttles it with the threshold + long totalBytesAvailable = TotalBytesAvailable; + if (totalBytesAvailable > _lastTotalBytesAvailable) + { + _lastTotalBytesAvailable = totalBytesAvailable; + RaiseDataReceivedChars(); + } + + if (_pinChanged != null) + { + // Checking for changes could technically speaking be done by waiting with ioctl+TIOCMIWAIT + // This would require spinning new thread and also would potentially trigger events when + // user didn't have time to respond. + // Diffing seems like a better solution. + Signals current = Interop.Termios.TermiosGetAllSignals(_handle); + + // There is no really good action we can take when this errors so just ignore + // a sinle event. + if (current != Signals.Error && lastSignals != Signals.Error) + { + Signals changed = current ^ lastSignals; + if (changed != Signals.None) + { + NotifyPinChanges(changed); + } + } + + lastSignals = current; + } + + lastIsIdle = isIdle; + } + } + + private static void RemoveCompletedTasks(Queue queue, object queueLock) + { + // assumes dequeue-ing happens on a single thread + lock (queueLock) + { + while (queue.TryPeek(out var r) && r.IsCompleted) + queue.TryDequeue(out _); + } + } + + private bool IsReadQueueEmpty() + { + lock (_readQueueLock) + { + return _readQueue.Count == 0; + } + } + + private bool IsWriteQueueEmpty() + { + lock (_writeQueueLock) + { + return _writeQueue.Count == 0; + } + } + + private static CancellationTokenSource GetCancellationTokenSourceFromTimeout(int timeoutMs) + { + return timeoutMs == SerialPort.InfiniteTimeout ? + null : + new CancellationTokenSource(Math.Max(timeoutMs, TimeoutResolution)); + } + + private abstract class SerialStreamIORequest : TaskCompletionSource + { + public bool IsCompleted => Task.IsCompleted; + private readonly SerialStream _parent; + private readonly CancellationTokenRegistration _cancellationTokenRegistration; + + protected SerialStreamIORequest(SerialStream parent, CancellationToken ct) + : base(TaskCreationOptions.RunContinuationsAsynchronously) + { + _parent = parent; + _cancellationTokenRegistration = ct.Register(s => + { + var request = (SerialStreamIORequest)s; + request.TrySetCanceled(); + request._parent.HasCancelledTasksToProcess = true; + }, this); + } + + internal void Complete(int numBytes) + { + TrySetResult(numBytes); + _cancellationTokenRegistration.Dispose(); + } + + internal void Complete(Exception exception) + { + TrySetException(exception); + _cancellationTokenRegistration.Dispose(); + } + } + + private sealed class SerialStreamReadRequest : SerialStreamIORequest + { + public Memory Buffer { get; } + + public SerialStreamReadRequest(SerialStream parent, CancellationToken ct, Memory buffer) + : base(parent, ct) + { + Buffer = buffer; + } + } + + private sealed class SerialStreamWriteRequest : SerialStreamIORequest + { + public ReadOnlyMemory Buffer { get; private set; } + + public SerialStreamWriteRequest(SerialStream parent, CancellationToken ct, ReadOnlyMemory buffer) + : base(parent, ct) + { + Buffer = buffer; + } + + internal void Complete() + { + Debug.Assert(Buffer.Length == 0); + Complete(Buffer.Length); + } + + internal void ProcessBytes(int numBytes) + { + Buffer = Buffer.Slice(numBytes); + } + } + } +} diff --git a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.Windows.cs b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.Windows.cs index 27cc8a8c7a6b63..e1022ebf6a5ffd 100644 --- a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.Windows.cs +++ b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.Windows.cs @@ -525,21 +525,39 @@ internal bool DsrHolding } - // Fills comStat structure from an unmanaged function - // to determine the number of bytes waiting in the serial driver's internal receive buffer. - internal int BytesToRead + internal int GetBytesToRead(int buffered, bool throwOnDispose = true) { - get + SafeFileHandle? handle = _handle; + if (handle == null || !IsOpen) + { + if (throwOnDispose) + { + InternalResources.FileNotOpen(); + } + return 0; + } + + try { int errorCode = 0; // "ref" arguments need to have values, as opposed to "out" arguments - if (!Interop.Kernel32.ClearCommError(_handle, ref errorCode, ref _comStat)) + if (!Interop.Kernel32.ClearCommError(handle, ref errorCode, ref _comStat)) { throw Win32Marshal.GetExceptionForLastWin32Error(); } - return (int)_comStat.cbInQue; + return buffered + (int)_comStat.cbInQue; + } + catch (ObjectDisposedException) when (!throwOnDispose) + { + return 0; } } +#pragma warning disable CA1822 + internal void OnRaiseCharsEventSkipped() + { + } +#pragma warning restore CA1822 + // Fills comStat structure from an unmanaged function // to determine the number of bytes waiting in the serial driver's internal transmit buffer. internal int BytesToWrite diff --git a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.cs b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.cs index 82181c3e6b45d3..a55bcd11567136 100644 --- a/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.cs +++ b/src/libraries/System.IO.Ports/src/System/IO/Ports/SerialStream.cs @@ -22,6 +22,10 @@ internal sealed partial class SerialStream : Stream private bool _inBreak; private Handshake _handshake; + internal int ReceivedBytesThreshold { get; set; } = 1; + + internal int BytesToRead => GetBytesToRead(buffered: 0); + #pragma warning disable CS0067 // Events shared by Windows and Linux, on Linux we currently never call them // called when any runtime error occurs on the port (frame, overrun, parity, etc.) internal event SerialErrorReceivedEventHandler ErrorReceived; diff --git a/src/libraries/System.Private.CoreLib/ref/System.Private.CoreLib.ExtraApis.cs b/src/libraries/System.Private.CoreLib/ref/System.Private.CoreLib.ExtraApis.cs index 9df3fe330338ea..3e7f0c7c6287eb 100644 --- a/src/libraries/System.Private.CoreLib/ref/System.Private.CoreLib.ExtraApis.cs +++ b/src/libraries/System.Private.CoreLib/ref/System.Private.CoreLib.ExtraApis.cs @@ -49,9 +49,13 @@ public UnixHandleAsyncContext(System.Runtime.InteropServices.SafeHandle handle) public bool IsReadReady(out int observedSequenceNumber) { throw null; } public bool IsWriteReady(out int observedSequenceNumber) { throw null; } public System.Threading.UnixHandleAsyncContext.AsyncResult StartAsyncRead(System.Threading.UnixHandleAsyncContext.Operation operation, int observedSequenceNumber, System.Threading.CancellationToken cancellationToken) { throw null; } + public int StartAsyncReadAsInt(System.Threading.UnixHandleAsyncContext.Operation operation, int observedSequenceNumber, System.Threading.CancellationToken cancellationToken) { throw null; } public System.Threading.UnixHandleAsyncContext.AsyncResult StartAsyncWrite(System.Threading.UnixHandleAsyncContext.Operation operation, int observedSequenceNumber, System.Threading.CancellationToken cancellationToken) { throw null; } + public int StartAsyncWriteAsInt(System.Threading.UnixHandleAsyncContext.Operation operation, int observedSequenceNumber, System.Threading.CancellationToken cancellationToken) { throw null; } public System.Threading.UnixHandleAsyncContext.SyncResult Read(System.Threading.UnixHandleAsyncContext.Operation operation, int observedSequenceNumber, int timeout) { throw null; } + public int ReadAsInt(System.Threading.UnixHandleAsyncContext.Operation operation, int observedSequenceNumber, int timeout) { throw null; } public System.Threading.UnixHandleAsyncContext.SyncResult Write(System.Threading.UnixHandleAsyncContext.Operation operation, int observedSequenceNumber, int timeout) { throw null; } + public int WriteAsInt(System.Threading.UnixHandleAsyncContext.Operation operation, int observedSequenceNumber, int timeout) { throw null; } public bool AbortAndDispose() { throw null; } public enum AsyncResult { @@ -78,6 +82,12 @@ public abstract partial class Operation : System.Threading.IThreadPoolWorkItem protected virtual void ExecuteThreadPoolWorkItem() { } void System.Threading.IThreadPoolWorkItem.Execute() { } } + public sealed class DelegateOperation : System.Threading.UnixHandleAsyncContext.Operation + { + public DelegateOperation(System.Func tryComplete, System.Action onCompleted) { } + protected internal override bool TryCompleteOperation(System.Runtime.InteropServices.SafeHandle handle) { throw null; } + protected internal override void OnCompleted(System.Threading.UnixHandleAsyncContext.OnCompletedResult result) { } + } } } diff --git a/src/libraries/System.Private.CoreLib/src/System.Private.CoreLib.Shared.projitems b/src/libraries/System.Private.CoreLib/src/System.Private.CoreLib.Shared.projitems index abe0dcc14a7c0f..218f238cfbcd8b 100644 --- a/src/libraries/System.Private.CoreLib/src/System.Private.CoreLib.Shared.projitems +++ b/src/libraries/System.Private.CoreLib/src/System.Private.CoreLib.Shared.projitems @@ -2966,6 +2966,7 @@ + @@ -2984,6 +2985,7 @@ + diff --git a/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.DelegateOperation.cs b/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.DelegateOperation.cs new file mode 100644 index 00000000000000..01c58b753ea7dd --- /dev/null +++ b/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.DelegateOperation.cs @@ -0,0 +1,47 @@ +// Licensed to the .NET Foundation under one or more agreements. +// The .NET Foundation licenses this file to you under the MIT license. + +using System.Runtime.InteropServices; + +namespace System.Threading +{ + public sealed partial class UnixHandleAsyncContext + { + /// + /// An that delegates to caller-supplied callbacks. + /// This enables libraries that cannot subclass directly + /// (e.g. out-of-band packages) to implement operations. + /// + public sealed class DelegateOperation : Operation + { + private readonly Func _tryComplete; + private readonly Action _onCompleted; + + /// + /// Creates a new with the specified callbacks. + /// + /// + /// Performs the I/O operation. Returns if completed, + /// if pending (EWOULDBLOCK). + /// + /// + /// Called when the operation completes asynchronously. + /// The argument is the value cast to . + /// + public DelegateOperation(Func tryComplete, Action onCompleted) + { + ArgumentNullException.ThrowIfNull(tryComplete); + ArgumentNullException.ThrowIfNull(onCompleted); + + _tryComplete = tryComplete; + _onCompleted = onCompleted; + } + + protected internal override bool TryCompleteOperation(SafeHandle handle) + => _tryComplete(handle); + + protected internal override void OnCompleted(OnCompletedResult result) + => _onCompleted((int)result); + } + } +} diff --git a/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.PlatformNotSupported.cs b/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.PlatformNotSupported.cs index 5e054dbbf08f64..634bfd9dc643f3 100644 --- a/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.PlatformNotSupported.cs +++ b/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.PlatformNotSupported.cs @@ -21,21 +21,41 @@ public AsyncResult StartAsyncRead(Operation operation, int observedSequenceNumbe throw new PlatformNotSupportedException(); } + public int StartAsyncReadAsInt(Operation operation, int observedSequenceNumber, CancellationToken cancellationToken) + { + throw new PlatformNotSupportedException(); + } + public AsyncResult StartAsyncWrite(Operation operation, int observedSequenceNumber, CancellationToken cancellationToken) { throw new PlatformNotSupportedException(); } + public int StartAsyncWriteAsInt(Operation operation, int observedSequenceNumber, CancellationToken cancellationToken) + { + throw new PlatformNotSupportedException(); + } + public SyncResult Read(Operation operation, int observedSequenceNumber, int timeout) { throw new PlatformNotSupportedException(); } + public int ReadAsInt(Operation operation, int observedSequenceNumber, int timeout) + { + throw new PlatformNotSupportedException(); + } + public SyncResult Write(Operation operation, int observedSequenceNumber, int timeout) { throw new PlatformNotSupportedException(); } + public int WriteAsInt(Operation operation, int observedSequenceNumber, int timeout) + { + throw new PlatformNotSupportedException(); + } + public bool AbortAndDispose() { throw new PlatformNotSupportedException(); diff --git a/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.cs b/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.cs index 65a115591de0fc..d41ea4869d9acb 100644 --- a/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.cs +++ b/src/libraries/System.Private.CoreLib/src/System/Threading/UnixHandleAsyncContext.cs @@ -688,6 +688,20 @@ public bool AbortAndDispose() return aborted; } + // int-returning wrappers for UnsafeAccessor consumers that cannot use + // UnsafeAccessorTypeAttribute with value types (see dotnet/runtime#121655). + public int StartAsyncReadAsInt(Operation operation, int observedSequenceNumber, CancellationToken cancellationToken) + => (int)StartAsyncRead(operation, observedSequenceNumber, cancellationToken); + + public int StartAsyncWriteAsInt(Operation operation, int observedSequenceNumber, CancellationToken cancellationToken) + => (int)StartAsyncWrite(operation, observedSequenceNumber, cancellationToken); + + public int ReadAsInt(Operation operation, int observedSequenceNumber, int timeout) + => (int)Read(operation, observedSequenceNumber, timeout); + + public int WriteAsInt(Operation operation, int observedSequenceNumber, int timeout) + => (int)Write(operation, observedSequenceNumber, timeout); + // Called on the epoll thread, speculatively tries to process inline events and errors, // and returns any remaining events that remain to be processed. [MethodImpl(MethodImplOptions.AggressiveInlining)]