From 173316a5a93692b390dc62b13c4148b3dd6d186c Mon Sep 17 00:00:00 2001 From: Geoffrey McRae Date: Wed, 12 Aug 2026 14:36:50 +1000 Subject: [PATCH] [client] audio: handle USB transport backpressure --- client/include/interface/usbredir.h | 8 ++++--- client/src/usb_audio.c | 33 +++++++++++++++++++++---- client/transports/SPICE/session.c | 5 ++++ client/transports/SPICE/usbredir.c | 37 ++++++++++++++++++++++------- repos/PureSpice | 2 +- 5 files changed, 67 insertions(+), 18 deletions(-) diff --git a/client/include/interface/usbredir.h b/client/include/interface/usbredir.h index 8ceb8dbd..9817731b 100644 --- a/client/include/interface/usbredir.h +++ b/client/include/interface/usbredir.h @@ -21,6 +21,8 @@ #ifndef _H_LG_CLIENT_USBREDIR_INTERFACE_ #define _H_LG_CLIENT_USBREDIR_INTERFACE_ +#include + struct usbredirparser; typedef struct LG_USBRedirDeviceOps @@ -31,9 +33,9 @@ typedef struct LG_USBRedirDeviceOps /* Queue the device descriptors and connection announcement. */ void (*plug)(void * opaque, struct usbredirparser * parser); - /* Queue time-sensitive device data immediately before parser output is - * flushed. */ - void (*process)(void * opaque); + /* Process device data immediately before parser output is flushed. + * writable is false while previously queued output is backpressured. */ + void (*process)(void * opaque, bool writable); /* Stop all device activity. The bridge sends the disconnect packet. */ void (*unplug)(void * opaque); diff --git a/client/src/usb_audio.c b/client/src/usb_audio.c index 69f85217..45a4b3f5 100644 --- a/client/src/usb_audio.c +++ b/client/src/usb_audio.c @@ -381,6 +381,7 @@ struct LG_USBAudio uint64_t feedbackNextPacketTime; uint32_t feedbackQueueTarget; atomic_uint feedbackValue; + bool outputBackpressured; atomic_bool recordStreaming; atomic_bool recordAccepting; @@ -924,6 +925,7 @@ static void resetDevice(LG_USBAudio * audio) audio->playbackSampleRate = LG_USB_AUDIO_DEFAULT_SAMPLE_RATE; audio->recordBatchPackets = 1; audio->recordSafetyPackets = 1; + audio->outputBackpressured = false; atomic_store_explicit(&audio->recordSampleRate, LG_USB_AUDIO_DEFAULT_SAMPLE_RATE, memory_order_relaxed); atomic_store_explicit(&audio->recordRateQ16, @@ -1220,7 +1222,7 @@ static void sendRecordPacketBatch(LG_USBAudio * audio, uint32_t count, static void sendRecordRefill(LG_USBAudio * audio, uint64_t rateQ16) { const uint32_t count = min( - audio->recordRefillPackets, USB_AUDIO_RECORD_MAX_BATCH); + audio->recordRefillPackets, audio->recordBatchPackets); sendRecordPacketBatch(audio, count, rateQ16, true, 0, 0); audio->recordRefillPackets -= count; if (!audio->recordRefillPackets) @@ -1277,7 +1279,7 @@ static void sendRecordPackets(LG_USBAudio * audio) } const uint32_t sendCount = (uint32_t)min( - count, (uint64_t)USB_AUDIO_RECORD_MAX_BATCH); + count, (uint64_t)audio->recordBatchPackets); const uint64_t catchupFrames = (audio->recordPacketPhase + rateQ16 * count) / USB_AUDIO_RECORD_RATE_DENOMINATOR; @@ -1823,12 +1825,30 @@ static void reportPlaybackDebug(LG_USBAudio * audio) audio->playbackInvalidPackets = 0; } -static void processDevice(void * opaque) +static void processDevice(void * opaque, bool writable) { LG_USBAudio * audio = opaque; - sendFeedbackPacketsDue(audio); + const uint64_t now = streamTime(); + if (!writable) + { + if (atomic_load_explicit( + &audio->recordStreaming, memory_order_acquire)) + audio->recordNextPacketTime = + now + USB_AUDIO_RECORD_PACKET_INTERVAL_NS; + audio->outputBackpressured = true; + } + else if (audio->outputBackpressured) + { + audio->recordNextPacketTime = + now + USB_AUDIO_RECORD_PACKET_INTERVAL_NS; + audio->outputBackpressured = false; + } + + if (writable) + sendFeedbackPacketsDue(audio); flushPlayback(audio); - sendRecordPackets(audio); + if (writable) + sendRecordPackets(audio); reportPlaybackDebug(audio); reportRecordDebug(audio); } @@ -1986,6 +2006,9 @@ uint64_t lgUsbAudio_processDelayNs(const LG_USBAudio * audio) if (!audio) return UINT64_MAX; + if (audio->outputBackpressured) + return UINT64_MAX; + const uint64_t now = streamTime(); uint64_t delay = UINT64_MAX; diff --git a/client/transports/SPICE/session.c b/client/transports/SPICE/session.c index 7b18280a..ce94655a 100644 --- a/client/transports/SPICE/session.c +++ b/client/transports/SPICE/session.c @@ -376,6 +376,11 @@ int spiceSession_thread(void * opaque) while (lgUsbRedir_disconnectPending(transport->usbRedir) && microtime() < deadline) { + if (!lgUsbRedir_process(transport->usbRedir)) + { + DEBUG_WARN("Failed to process USB audio device removal"); + break; + } status = purespice_process(10); if (status != PS_STATUS_RUN) break; diff --git a/client/transports/SPICE/usbredir.c b/client/transports/SPICE/usbredir.c index 9734f6dc..433c9629 100644 --- a/client/transports/SPICE/usbredir.c +++ b/client/transports/SPICE/usbredir.c @@ -83,11 +83,12 @@ static int readUSBRedir(void * opaque, uint8_t * data, int count) static int writeUSBRedir(void * opaque, uint8_t * data, int count) { LG_USBRedir * usbredir = opaque; - if (!usbredir->channel || - !purespice_usbRedirWrite(usbredir->channel, data, count)) + if (!usbredir->channel) return -1; - return count; + const ssize_t written = purespice_usbRedirWriteNonblocking( + usbredir->channel, data, count); + return written >= 0 && written <= count ? (int)written : -1; } static void logUSBRedir(void * opaque, int level, const char * message) @@ -324,6 +325,14 @@ static bool processUSBRedir(LG_USBRedir * usbredir, bool recover) return result; } + bool result = flushUSBRedir(usbredir); + if (!result) + { + if (recover) + resetChannel(usbredir); + return false; + } + if (usbredir->disconnectPending) { if (usbredir->disconnectDeadline && @@ -334,10 +343,7 @@ static bool processUSBRedir(LG_USBRedir * usbredir, bool recover) return true; } - const bool result = flushUSBRedir(usbredir); - if (!result && recover) - resetChannel(usbredir); - return result; + return true; } const bool desired = atomic_load_explicit(&usbredir->desiredPlugged, @@ -359,12 +365,25 @@ static bool processUSBRedir(LG_USBRedir * usbredir, bool recover) now + USB_REDIR_DISCONNECT_TIMEOUT_NS : 0; usbredirparser_send_device_disconnect(usbredir->parser); } + + result = flushUSBRedir(usbredir); + if (!result) + { + if (recover) + resetChannel(usbredir); + return false; + } } if (usbredir->plugged && usbredir->deviceOps->process) - usbredir->deviceOps->process(usbredir->deviceOpaque); + { + const bool writable = + usbredirparser_has_data_to_write(usbredir->parser) == 0 && + purespice_usbRedirWritable(usbredir->channel); + usbredir->deviceOps->process(usbredir->deviceOpaque, writable); + } - const bool result = flushUSBRedir(usbredir); + result = flushUSBRedir(usbredir); if (!result && recover) resetChannel(usbredir); return result; diff --git a/repos/PureSpice b/repos/PureSpice index 5ddda1e4..b21f8b20 160000 --- a/repos/PureSpice +++ b/repos/PureSpice @@ -1 +1 @@ -Subproject commit 5ddda1e49e6459d6da369296bff374dc53269f85 +Subproject commit b21f8b20e2ed00c5f29ce5ae28f8dd9156fc4d22