[client] audio: handle USB transport backpressure

This commit is contained in:
Geoffrey McRae
2026-08-12 14:36:50 +10:00
parent 573204d44e
commit 173316a5a9
5 changed files with 67 additions and 18 deletions

View File

@@ -21,6 +21,8 @@
#ifndef _H_LG_CLIENT_USBREDIR_INTERFACE_ #ifndef _H_LG_CLIENT_USBREDIR_INTERFACE_
#define _H_LG_CLIENT_USBREDIR_INTERFACE_ #define _H_LG_CLIENT_USBREDIR_INTERFACE_
#include <stdbool.h>
struct usbredirparser; struct usbredirparser;
typedef struct LG_USBRedirDeviceOps typedef struct LG_USBRedirDeviceOps
@@ -31,9 +33,9 @@ typedef struct LG_USBRedirDeviceOps
/* Queue the device descriptors and connection announcement. */ /* Queue the device descriptors and connection announcement. */
void (*plug)(void * opaque, struct usbredirparser * parser); void (*plug)(void * opaque, struct usbredirparser * parser);
/* Queue time-sensitive device data immediately before parser output is /* Process device data immediately before parser output is flushed.
* flushed. */ * writable is false while previously queued output is backpressured. */
void (*process)(void * opaque); void (*process)(void * opaque, bool writable);
/* Stop all device activity. The bridge sends the disconnect packet. */ /* Stop all device activity. The bridge sends the disconnect packet. */
void (*unplug)(void * opaque); void (*unplug)(void * opaque);

View File

@@ -381,6 +381,7 @@ struct LG_USBAudio
uint64_t feedbackNextPacketTime; uint64_t feedbackNextPacketTime;
uint32_t feedbackQueueTarget; uint32_t feedbackQueueTarget;
atomic_uint feedbackValue; atomic_uint feedbackValue;
bool outputBackpressured;
atomic_bool recordStreaming; atomic_bool recordStreaming;
atomic_bool recordAccepting; atomic_bool recordAccepting;
@@ -924,6 +925,7 @@ static void resetDevice(LG_USBAudio * audio)
audio->playbackSampleRate = LG_USB_AUDIO_DEFAULT_SAMPLE_RATE; audio->playbackSampleRate = LG_USB_AUDIO_DEFAULT_SAMPLE_RATE;
audio->recordBatchPackets = 1; audio->recordBatchPackets = 1;
audio->recordSafetyPackets = 1; audio->recordSafetyPackets = 1;
audio->outputBackpressured = false;
atomic_store_explicit(&audio->recordSampleRate, atomic_store_explicit(&audio->recordSampleRate,
LG_USB_AUDIO_DEFAULT_SAMPLE_RATE, memory_order_relaxed); LG_USB_AUDIO_DEFAULT_SAMPLE_RATE, memory_order_relaxed);
atomic_store_explicit(&audio->recordRateQ16, 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) static void sendRecordRefill(LG_USBAudio * audio, uint64_t rateQ16)
{ {
const uint32_t count = min( const uint32_t count = min(
audio->recordRefillPackets, USB_AUDIO_RECORD_MAX_BATCH); audio->recordRefillPackets, audio->recordBatchPackets);
sendRecordPacketBatch(audio, count, rateQ16, true, 0, 0); sendRecordPacketBatch(audio, count, rateQ16, true, 0, 0);
audio->recordRefillPackets -= count; audio->recordRefillPackets -= count;
if (!audio->recordRefillPackets) if (!audio->recordRefillPackets)
@@ -1277,7 +1279,7 @@ static void sendRecordPackets(LG_USBAudio * audio)
} }
const uint32_t sendCount = (uint32_t)min( const uint32_t sendCount = (uint32_t)min(
count, (uint64_t)USB_AUDIO_RECORD_MAX_BATCH); count, (uint64_t)audio->recordBatchPackets);
const uint64_t catchupFrames = const uint64_t catchupFrames =
(audio->recordPacketPhase + rateQ16 * count) / (audio->recordPacketPhase + rateQ16 * count) /
USB_AUDIO_RECORD_RATE_DENOMINATOR; USB_AUDIO_RECORD_RATE_DENOMINATOR;
@@ -1823,11 +1825,29 @@ static void reportPlaybackDebug(LG_USBAudio * audio)
audio->playbackInvalidPackets = 0; audio->playbackInvalidPackets = 0;
} }
static void processDevice(void * opaque) static void processDevice(void * opaque, bool writable)
{ {
LG_USBAudio * audio = opaque; LG_USBAudio * audio = opaque;
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); sendFeedbackPacketsDue(audio);
flushPlayback(audio); flushPlayback(audio);
if (writable)
sendRecordPackets(audio); sendRecordPackets(audio);
reportPlaybackDebug(audio); reportPlaybackDebug(audio);
reportRecordDebug(audio); reportRecordDebug(audio);
@@ -1986,6 +2006,9 @@ uint64_t lgUsbAudio_processDelayNs(const LG_USBAudio * audio)
if (!audio) if (!audio)
return UINT64_MAX; return UINT64_MAX;
if (audio->outputBackpressured)
return UINT64_MAX;
const uint64_t now = streamTime(); const uint64_t now = streamTime();
uint64_t delay = UINT64_MAX; uint64_t delay = UINT64_MAX;

View File

@@ -376,6 +376,11 @@ int spiceSession_thread(void * opaque)
while (lgUsbRedir_disconnectPending(transport->usbRedir) && while (lgUsbRedir_disconnectPending(transport->usbRedir) &&
microtime() < deadline) microtime() < deadline)
{ {
if (!lgUsbRedir_process(transport->usbRedir))
{
DEBUG_WARN("Failed to process USB audio device removal");
break;
}
status = purespice_process(10); status = purespice_process(10);
if (status != PS_STATUS_RUN) if (status != PS_STATUS_RUN)
break; break;

View File

@@ -83,11 +83,12 @@ static int readUSBRedir(void * opaque, uint8_t * data, int count)
static int writeUSBRedir(void * opaque, uint8_t * data, int count) static int writeUSBRedir(void * opaque, uint8_t * data, int count)
{ {
LG_USBRedir * usbredir = opaque; LG_USBRedir * usbredir = opaque;
if (!usbredir->channel || if (!usbredir->channel)
!purespice_usbRedirWrite(usbredir->channel, data, count))
return -1; 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) static void logUSBRedir(void * opaque, int level, const char * message)
@@ -324,6 +325,14 @@ static bool processUSBRedir(LG_USBRedir * usbredir, bool recover)
return result; return result;
} }
bool result = flushUSBRedir(usbredir);
if (!result)
{
if (recover)
resetChannel(usbredir);
return false;
}
if (usbredir->disconnectPending) if (usbredir->disconnectPending)
{ {
if (usbredir->disconnectDeadline && if (usbredir->disconnectDeadline &&
@@ -334,10 +343,7 @@ static bool processUSBRedir(LG_USBRedir * usbredir, bool recover)
return true; return true;
} }
const bool result = flushUSBRedir(usbredir); return true;
if (!result && recover)
resetChannel(usbredir);
return result;
} }
const bool desired = atomic_load_explicit(&usbredir->desiredPlugged, 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; now + USB_REDIR_DISCONNECT_TIMEOUT_NS : 0;
usbredirparser_send_device_disconnect(usbredir->parser); usbredirparser_send_device_disconnect(usbredir->parser);
} }
result = flushUSBRedir(usbredir);
if (!result)
{
if (recover)
resetChannel(usbredir);
return false;
}
} }
if (usbredir->plugged && usbredir->deviceOps->process) 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) if (!result && recover)
resetChannel(usbredir); resetChannel(usbredir);
return result; return result;