From 1f36d8dc93953912d0c1ec7bdeb8e2adcf4471f6 Mon Sep 17 00:00:00 2001 From: Tom Deseyn Date: Tue, 10 Feb 2026 11:27:34 +0100 Subject: [PATCH] Protocol: change send message methods to be void instead of ValueTask. The sends are backed by an unbound Channel, they always complete. --- src/Tmds.DBus.Protocol/IMessageStream.cs | 2 +- src/Tmds.DBus.Protocol/InnerConnection.cs | 46 ++++++------------- src/Tmds.DBus.Protocol/MessageStream.cs | 12 +---- .../PairedConnection.cs | 6 +-- 4 files changed, 21 insertions(+), 45 deletions(-) diff --git a/src/Tmds.DBus.Protocol/IMessageStream.cs b/src/Tmds.DBus.Protocol/IMessageStream.cs index 5337180e..acda675c 100644 --- a/src/Tmds.DBus.Protocol/IMessageStream.cs +++ b/src/Tmds.DBus.Protocol/IMessageStream.cs @@ -6,7 +6,7 @@ interface IMessageStream void ReceiveMessages(MessageReceivedHandler handler, T state); - ValueTask TrySendMessageAsync(MessageBuffer message); + bool TrySendMessage(MessageBuffer message); void BecomeMonitor(); diff --git a/src/Tmds.DBus.Protocol/InnerConnection.cs b/src/Tmds.DBus.Protocol/InnerConnection.cs index 2a38ef92..5b8bcffb 100644 --- a/src/Tmds.DBus.Protocol/InnerConnection.cs +++ b/src/Tmds.DBus.Protocol/InnerConnection.cs @@ -247,7 +247,7 @@ public async ValueTask ConnectAsync(Stream stream, string? userId, CancellationT { MyValueTaskSource vts = new(); - await CallMethodAsync( + CallMethod( message: CreateHelloMessage(), static (Exception? exception, Message message, object? state) => { @@ -265,7 +265,7 @@ await CallMethodAsync( { vtsState.SetResult(null); } - }, vts).ConfigureAwait(false); + }, vts); return await new ValueTask(vts, token: 0).ConfigureAwait(false); @@ -496,7 +496,7 @@ public void Dispose() _disconnectedTcs?.SetResult(GetWaitForDisconnectException()); } - private ValueTask CallMethodAsync(MessageBuffer message, MessageReceivedHandler returnHandler, object? state) + private void CallMethod(MessageBuffer message, MessageReceivedHandler returnHandler, object? state) { MessageHandlerDelegate fn = static (Exception? exception, Message message, object? state1, object? state2, object? state3) => { @@ -504,10 +504,10 @@ private ValueTask CallMethodAsync(MessageBuffer message, MessageReceivedHandler }; MessageHandler handler = new(fn, returnHandler, state); - return CallMethodAsync(message, handler); + CallMethod(message, handler); } - private async ValueTask CallMethodAsync(MessageBuffer message, MessageHandler handler) + private void CallMethod(MessageBuffer message, MessageHandler handler) { bool messageSent = false; try @@ -528,7 +528,7 @@ private async ValueTask CallMethodAsync(MessageBuffer message, MessageHandler ha } } - messageSent = await _messageStream!.TrySendMessageAsync(message).ConfigureAwait(false); + messageSent = _messageStream!.TrySendMessage(message); } finally { @@ -574,7 +574,7 @@ public async Task CallMethodAsync(MessageBuffer message, MessageValueReade MyValueTaskSource vts = new(); MessageHandler handler = new(fn, valueReader, vts, state); - await CallMethodAsync(message, handler).ConfigureAwait(false); + CallMethod(message, handler); return await new ValueTask(vts, 0).ConfigureAwait(false); } @@ -583,8 +583,7 @@ public async Task CallMethodAsync(MessageBuffer message) { MyValueTaskSource vts = new(); - await CallMethodAsync(message, - static (Exception? exception, Message message, object? state) => CompleteCallValueTaskSource(exception, message, state), vts).ConfigureAwait(false); + CallMethod(message, static (Exception? exception, Message message, object? state) => CompleteCallValueTaskSource(exception, message, state), vts); await new ValueTask(vts, 0).ConfigureAwait(false); } @@ -750,19 +749,13 @@ public async ValueTask AddMatchAsync(SynchronizationContext? syn }; _pendingCalls.Add(addMatchMessage.Serial, new(fn, matchMaker)); + + _messageStream!.TrySendMessage(addMatchMessage); } } if (subscribe) { - if (addMatchMessage is not null) - { - if (!await _messageStream!.TrySendMessageAsync(addMatchMessage).ConfigureAwait(false)) - { - addMatchMessage.ReturnToPool(); - } - } - try { await matchMaker.AddMatchTask!.ConfigureAwait(false); @@ -995,7 +988,6 @@ internal void InvokeHandler(Message message) private async void RemoveObserver(MatchMaker matchMaker, Observer observer) { string ruleString = matchMaker.RuleString; - bool sendMessage = false; lock (_gate) { @@ -1007,23 +999,16 @@ private async void RemoveObserver(MatchMaker matchMaker, Observer observer) if (_matchMakers.TryGetValue(ruleString, out _)) { matchMaker.Observers.Remove(observer); - sendMessage = matchMaker.AddMatchTcs is not null && matchMaker.HasSubscribers; + bool sendMessage = matchMaker.AddMatchTcs is not null && matchMaker.HasSubscribers; if (sendMessage) { _matchMakers.Remove(ruleString); + var message = CreateRemoveMatchMessage(); + _messageStream!.TrySendMessage(message); } } } - if (sendMessage) - { - var message = CreateRemoveMatchMessage(); - if (!await _messageStream!.TrySendMessageAsync(message).ConfigureAwait(false)) - { - message.ReturnToPool(); - } - } - MessageBuffer CreateRemoveMatchMessage() { using var writer = GetMessageWriter(); @@ -1261,10 +1246,9 @@ private static bool IsEqual(ReadOnlySpan lhs, ReadOnlySpan rhs) public MessageWriter GetMessageWriter() => _parentDBusConnection.GetMessageWriter(); - public async void SendMessage(MessageBuffer message) + public void SendMessage(MessageBuffer message) { - bool messageSent = await _messageStream!.TrySendMessageAsync(message).ConfigureAwait(false); - if (!messageSent) + if (!_messageStream!.TrySendMessage(message)) { message.ReturnToPool(); } diff --git a/src/Tmds.DBus.Protocol/MessageStream.cs b/src/Tmds.DBus.Protocol/MessageStream.cs index 5c0d6a1c..bbe06647 100644 --- a/src/Tmds.DBus.Protocol/MessageStream.cs +++ b/src/Tmds.DBus.Protocol/MessageStream.cs @@ -409,16 +409,8 @@ int CopyBuffer(ReadOnlySequence src, Memory dst) } } - public async ValueTask TrySendMessageAsync(MessageBuffer message) - { - while (await _messageWriter.WaitToWriteAsync().ConfigureAwait(false)) - { - if (_messageWriter.TryWrite(message)) - return true; - } - - return false; - } + public bool TrySendMessage(MessageBuffer message) + => _messageWriter.TryWrite(message); public void Close(Exception closeReason) => CloseCore(closeReason); diff --git a/test/Tmds.DBus.Protocol.Tests/PairedConnection.cs b/test/Tmds.DBus.Protocol.Tests/PairedConnection.cs index 86345f2d..84954851 100644 --- a/test/Tmds.DBus.Protocol.Tests/PairedConnection.cs +++ b/test/Tmds.DBus.Protocol.Tests/PairedConnection.cs @@ -82,16 +82,16 @@ public async void ReceiveMessages(IMessageStream.MessageReceivedHandler ha } } - public ValueTask TrySendMessageAsync(MessageBuffer message) + public bool TrySendMessage(MessageBuffer message) { _writeQueue.Enqueue(message); _writeSemaphore.Release(); - return ValueTask.FromResult(true); + return true; } public void Close(Exception? closeReason = null) { - TrySendMessageAsync(null!); // Use null as EOF. + TrySendMessage(null!); // Use null as EOF. } public void BecomeMonitor()