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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 28 additions & 5 deletions app/src/main/java/com/openipc/pixelpilot/WfbLinkManager.java
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,10 @@ public Map<String, UsbDevice> getAttachedAdapters() {

public synchronized void refreshAdapters() {
Map<String, UsbDevice> attachedAdapters = getAttachedAdapters();
if (attachedAdapters == null) {
Log.e(TAG, "Could not read the usb device filter, skipping adapter refresh.");
return;
}

boolean missingPermissions = false;
android.hardware.usb.UsbManager usbManager =
Expand All @@ -141,8 +145,13 @@ public synchronized void refreshAdapters() {
if (!usbManager.hasPermission(entry.getValue())) {
binding.tvMessage.setVisibility(View.VISIBLE);
binding.tvMessage.setText("No permission for wifi adapter(s) " + entry.getValue().getDeviceName());
// Android 14 refuses to deliver a PendingIntent built from an implicit
// intent to a runtime registered receiver, so the permission result never
// arrives unless the package is set explicitly.
Intent permissionIntent = new Intent(WfbLinkManager.ACTION_USB_PERMISSION);
permissionIntent.setPackage(context.getPackageName());
PendingIntent pendingIntent = PendingIntent.getBroadcast(context, 0,
new Intent(WfbLinkManager.ACTION_USB_PERMISSION), PendingIntent.FLAG_IMMUTABLE);
permissionIntent, PendingIntent.FLAG_IMMUTABLE);
usbManager.requestPermission(entry.getValue(), pendingIntent);
missingPermissions = true;
}
Expand All @@ -164,16 +173,27 @@ public synchronized void refreshAdapters() {
}

// Starts newly attached adapters.
boolean startFailed = false;
for (Map.Entry<String, UsbDevice> entry : attachedAdapters.entrySet()) {
if (activeWifiAdapters.containsKey(entry.getKey())) {
continue;
}
startAdapter(entry.getValue());
activeWifiAdapters.put(entry.getKey(), entry.getValue());
// Only track it as active if it actually came up, otherwise a failed adapter
// is never retried on the next refresh.
if (startAdapter(entry.getValue())) {
activeWifiAdapters.put(entry.getKey(), entry.getValue());
} else {
startFailed = true;
}
}

if (activeWifiAdapters.isEmpty()) {
String text = "No compatible wifi adapter found.";
// Now that a failed start no longer counts as active, an empty map covers two
// different problems, and blaming the filter for both sends people looking in
// the wrong place.
String text = startFailed
? "Wifi adapter found but could not be started - see the log."
: "No compatible wifi adapter found.";
binding.tvMessage.setText(text);
binding.tvMessage.setVisibility(View.VISIBLE);

Expand Down Expand Up @@ -217,7 +237,10 @@ public synchronized boolean startAdapter(UsbDevice dev) {
String text = "Starting wfb-ng channel " + wifiChannel + " with " + String.format(
"[%04X", dev.getVendorId()) + ":" + String.format("%04X]", dev.getProductId());
binding.tvMessage.setText(text);
wfbLink.start(wifiChannel, bandWidth.getValue(), dev);
if (!wfbLink.start(wifiChannel, bandWidth.getValue(), dev)) {
binding.tvMessage.setText("Could not open wifi adapter " + dev.getDeviceName());
return false;
}
return true;
}
}
30 changes: 21 additions & 9 deletions app/wfbngrtl8812/src/main/cpp/WfbngLink.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -34,12 +34,6 @@
#undef TAG
#define TAG "pixelpilot"

#define CRASH() \
do { \
int *i = 0; \
*i = 42; \
} while (0)

std::string generate_random_string(size_t length) {
const std::string characters = "abcdefghijklmnopqrstuvwxyz";
std::random_device rd;
Expand Down Expand Up @@ -92,6 +86,7 @@ void WfbngLink::initAgg() {
int WfbngLink::run(JNIEnv *env, jobject context, jint wifiChannel, jint bw, jint fd) {
int r;
libusb_context *ctx = NULL;
clear_stop_request(fd);
txFrame = std::make_shared<TxFrame>();

r = libusb_set_option(NULL, LIBUSB_OPTION_NO_DEVICE_DISCOVERY);
Expand Down Expand Up @@ -138,6 +133,14 @@ int WfbngLink::run(JNIEnv *env, jobject context, jint wifiChannel, jint bw, jint
return -1;
}

if (stop_requested(fd)) {
__android_log_print(ANDROID_LOG_WARN, TAG, "stop requested for fd=%d before bring-up, aborting", fd);
rtl_devices.erase(fd);
libusb_release_interface(dev_handle, 0);
libusb_exit(ctx);
return -1;
}

uint8_t *video_channel_id_be8 = reinterpret_cast<uint8_t *>(&video_channel_id_be);
uint8_t *udp_channel_id_be8 = reinterpret_cast<uint8_t *>(&udp_channel_id_be);
uint8_t *mavlink_channel_id_be8 = reinterpret_cast<uint8_t *>(&mavlink_channel_id_be);
Expand Down Expand Up @@ -246,7 +249,12 @@ int WfbngLink::run(JNIEnv *env, jobject context, jint wifiChannel, jint bw, jint

// Blocking RX loop on this thread; devourer pumps the libusb events
// itself. Returns once StopRxLoop() is called.
current_device->StartRxLoop(packetProcessor);
if (stop_requested(fd)) {
__android_log_print(
ANDROID_LOG_WARN, TAG, "stop requested for fd=%d during bring-up, not entering the rx loop", fd);
} else {
current_device->StartRxLoop(packetProcessor);
}
} catch (const std::runtime_error &error) {
__android_log_print(ANDROID_LOG_ERROR, TAG, "runtime_error: %s", error.what());
txFrame->stop();
Expand Down Expand Up @@ -282,9 +290,13 @@ int WfbngLink::run(JNIEnv *env, jobject context, jint wifiChannel, jint bw, jint
}

void WfbngLink::stop(JNIEnv *env, jobject context, jint fd) {
// Recorded first, and whether or not the device exists yet: run() may not have got as far
// as creating it, and StartRxLoop() would clear the flag set below anyway.
note_stop_requested(fd);
if (rtl_devices.find(fd) == rtl_devices.end()) {
__android_log_print(ANDROID_LOG_ERROR, TAG, "rtl_devices.find(%d) == rtl_devices.end()", fd);
CRASH();
// Happens when the adapter was already gone by the time the stop arrived, e.g. it
// was unplugged or the hub re-enumerated it. Nothing left to stop.
__android_log_print(ANDROID_LOG_WARN, TAG, "stop: no rtl device for fd=%d, already gone", fd);
return;
Comment thread
qodo-free-for-open-source-projects[bot] marked this conversation as resolved.
}
auto dev = rtl_devices.at(fd).get();
Expand Down
27 changes: 27 additions & 0 deletions app/wfbngrtl8812/src/main/cpp/WfbngLink.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ extern "C" {
#include <list>
#include <map>
#include <mutex>
#include <set>
#include <thread>
#include <vector> // Added for std::vector

Expand Down Expand Up @@ -59,6 +60,32 @@ class WfbngLink {
bool stbc_enabled{true};

std::map<int, std::shared_ptr<IRtlDevice>> rtl_devices;

// Set by stop() and read by run(). A StopRxLoop() only takes effect once the RX loop is
// running: RtlJaguarDevice::StartRxLoop() clears should_stop on entry, so a stop that
// lands anywhere before that - including the whole chip bring-up in InitWrite(), which
// is the longest part of run() - is thrown away, and run() then blocks in a loop nobody
// asked for. Recorded here instead, so run() can see it at the points where the flag
// itself cannot be trusted.
std::mutex stop_requested_mutex;
std::set<int> stop_requested_fds;

void note_stop_requested(int fd) {
std::lock_guard<std::mutex> lock(stop_requested_mutex);
stop_requested_fds.insert(fd);
}

// Cleared at the start of run(): fd numbers are reused, so a request left over from a
// previous session on the same number must not abort the new one.
void clear_stop_request(int fd) {
std::lock_guard<std::mutex> lock(stop_requested_mutex);
stop_requested_fds.erase(fd);
}

bool stop_requested(int fd) {
std::lock_guard<std::mutex> lock(stop_requested_mutex);
return stop_requested_fds.count(fd) > 0;
}
std::unique_ptr<std::thread> link_quality_thread{nullptr};
bool should_clear_stats{false};
FecChangeController fec;
Expand Down
103 changes: 91 additions & 12 deletions app/wfbngrtl8812/src/main/java/com/openipc/wfbngrtl8812/WfbNgLink.java
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import androidx.annotation.Keep;

import java.util.HashMap;
import java.util.Iterator;
import java.util.Map;
import java.util.Timer;
import java.util.TimerTask;
Expand Down Expand Up @@ -91,31 +92,108 @@ public void nativeSetUseStbc(int use) {
nativeSetUseStbc(nativeWfbngLink, use);
}

public synchronized void start(int wifiChannel, int bandWidth, UsbDevice usbDevice) {
public synchronized boolean start(int wifiChannel, int bandWidth, UsbDevice usbDevice) {
// linkThreads.put() below overwrites the entry for a device, which would orphan an
// older thread so that stopAll() never joins it and the interface is never released.
Thread existing = linkThreads.get(usbDevice);
if (existing != null && existing.isAlive()) {
Log.w(TAG, "wfb-ng already running on " + usbDevice.getDeviceName()
+ ", not starting a second");
return true;
}
if (existing != null) {
// A thread that outlived its join has since finished, so the connection stop()
// deliberately left open can go now. Otherwise the put() below would drop the
// last reference to it and leak the fd.
UsbDeviceConnection stale = linkConns.remove(usbDevice);
if (stale != null) {
Log.d(TAG, "releasing the usb connection left behind on " + usbDevice.getDeviceName());
stale.close();
}
linkThreads.remove(usbDevice);
}
Log.d(TAG, "wfb-ng monitoring on " + usbDevice.getDeviceName() + " using wifi channel " + wifiChannel);
UsbManager usbManager = (UsbManager) context.getSystemService(Context.USB_SERVICE);
// Returns null when the permission was revoked or the device disappeared between
// the permission check and here, which is easy to hit on a re-enumerating hub.
UsbDeviceConnection usbDeviceConnection = usbManager.openDevice(usbDevice);
if (usbDeviceConnection == null) {
Log.e(TAG, "Could not open " + usbDevice.getDeviceName() + " (no permission or already gone)");
return false;
}
int fd = usbDeviceConnection.getFileDescriptor();
if (fd < 0) {
Log.e(TAG, "Invalid file descriptor for " + usbDevice.getDeviceName());
usbDeviceConnection.close();
return false;
}
Thread t = new Thread(() -> nativeRun(nativeWfbngLink, context, wifiChannel, bandWidth, fd));
t.setName("wfb-" + usbDevice.getDeviceName().split("/dev/bus/usb/")[1]);
t.setName(threadNameFor(usbDevice));
linkThreads.put(usbDevice, t);
linkConns.put(usbDevice, usbDeviceConnection);
linkThreads.get(usbDevice).start();
t.start();
Log.d(TAG, "wfb-ng thread on " + usbDevice.getDeviceName() + " started.");
return true;
}

private static String threadNameFor(UsbDevice usbDevice) {
String name = usbDevice.getDeviceName();
String[] parts = name.split("/dev/bus/usb/");
return "wfb-" + (parts.length > 1 ? parts[1] : name);
}

/**
* The RX loop is joined so the USB interface is released before anything reopens it, but
* these calls come from Activity lifecycle callbacks on the main thread, where an
* unbounded join is a five second ANR waiting to happen. StopRxLoop() only breaks the
* receive loop - the thread then still has to stop the TX frame and the adaptive link,
* power the chip down, release the interface and exit libusb - so the wait has to be
* generous, but bounded.
*/
private static final long JOIN_TIMEOUT_MS = 3000;

/**
* @return true if the thread is gone and its usb connection can be released.
*/
private static boolean joinBounded(Thread t, String what) throws InterruptedException {
if (t == null) {
return true;
}
t.join(JOIN_TIMEOUT_MS);
if (t.isAlive()) {
// Nothing else may be torn down while the thread is still inside devourer.
// libusb_wrap_sys_device() keeps the fd it is given rather than duplicating it,
// so closing the UsbDeviceConnection now pulls the fd out from under a libusb
// that is still polling it. The kernel cancels the URBs on close but libusb
// never reaps them - op_handle_events() looks at POLLERR and not POLLNVAL - so
// poll() returns immediately, forever, and the loop spins on one core waiting
// for a transfer count that will never drop. Both map entries stay too, so
// start() can still see the thread and refuse a second RX loop on the device.
Log.e(TAG, "wfb-ng thread on " + what + " did not stop within " + JOIN_TIMEOUT_MS
+ "ms, leaving it and its usb connection in place");
return false;
}
Log.d(TAG, "wfb-ng thread on " + what + " done.");
return true;
}

public synchronized void stopAll() throws InterruptedException {
for (Map.Entry<UsbDevice, UsbDeviceConnection> entry : linkConns.entrySet()) {
nativeStop(nativeWfbngLink, context, entry.getValue().getFileDescriptor());
}
for (Map.Entry<UsbDevice, UsbDeviceConnection> entry : linkConns.entrySet()) {
Thread t = linkThreads.get(entry.getKey());
if (t != null) {
t.join();
Iterator<Map.Entry<UsbDevice, UsbDeviceConnection>> it = linkConns.entrySet().iterator();
while (it.hasNext()) {
Map.Entry<UsbDevice, UsbDeviceConnection> entry = it.next();
UsbDevice dev = entry.getKey();
if (!joinBounded(linkThreads.get(dev), dev.getDeviceName())) {
continue;
}
Log.d(TAG, "wfb-ng thread on " + entry.getKey().getDeviceName() + " done.");
// Without close() every attach/detach cycle leaks the fd openDevice() handed
// out, until the process runs out - but only once nothing is using it.
linkThreads.remove(dev);
it.remove();
entry.getValue().close();
}
linkThreads.clear();
}

public synchronized void stop(UsbDevice dev) throws InterruptedException {
Expand All @@ -125,11 +203,12 @@ public synchronized void stop(UsbDevice dev) throws InterruptedException {
}
int fd = conn.getFileDescriptor();
nativeStop(nativeWfbngLink, context, fd);
Thread t = linkThreads.get(dev);
if (t != null) {
t.join();
if (!joinBounded(linkThreads.get(dev), dev.getDeviceName())) {
return;
}
linkThreads.remove(dev);
linkConns.remove(dev);
conn.close();
}

public void SetWfbNGStatsChanged(final WfbNGStatsChanged callback) {
Expand Down