aboutsummaryrefslogtreecommitdiff
path: root/app
diff options
context:
space:
mode:
Diffstat (limited to 'app')
-rw-r--r--app/src/main/AndroidManifest.xml6
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/AudioSink.java74
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/AudioStream.java7
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/ControlMessages.java8
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/ControlStream.java54
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/Controller.java8
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/Devices.java13
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/Main.java18
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/Mirror.java63
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/MuxRecorder.java44
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/Server.java142
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/Session.java104
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/Sessions.java6
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/VideoSink.java55
-rw-r--r--app/src/main/java/invalid/lena/scrcpy/VideoStream.java7
-rw-r--r--app/src/test/java/invalid/lena/scrcpy/ControlMessagesTest.java6
-rw-r--r--app/src/test/java/invalid/lena/scrcpy/ControlStreamTest.java14
17 files changed, 408 insertions, 221 deletions
diff --git a/app/src/main/AndroidManifest.xml b/app/src/main/AndroidManifest.xml
index 3e1e4f3..540a499 100644
--- a/app/src/main/AndroidManifest.xml
+++ b/app/src/main/AndroidManifest.xml
@@ -2,9 +2,8 @@
<manifest xmlns:android="http://schemas.android.com/apk/res/android">
<uses-permission android:name="android.permission.INTERNET"/>
- <uses-permission android:name="android.permission.ACCESS_NETWORK_STATE"/>
<uses-permission android:name="android.permission.FOREGROUND_SERVICE"/>
- <uses-permission android:name="android.permission.FOREGROUND_SERVICE_DATA_SYNC"/>
+ <uses-permission android:name="android.permission.FOREGROUND_SERVICE_MEDIA_PLAYBACK"/>
<uses-permission android:name="android.permission.POST_NOTIFICATIONS"/>
<application
@@ -29,6 +28,7 @@
<activity
android:name=".Mirror"
android:exported="false"
+ android:launchMode="singleTask"
android:configChanges="orientation|screenSize|keyboardHidden"
android:screenOrientation="unspecified"
android:hardwareAccelerated="true"
@@ -42,7 +42,7 @@
<service
android:name=".Sessions"
android:exported="false"
- android:foregroundServiceType="dataSync"/>
+ android:foregroundServiceType="mediaPlayback"/>
</application>
diff --git a/app/src/main/java/invalid/lena/scrcpy/AudioSink.java b/app/src/main/java/invalid/lena/scrcpy/AudioSink.java
index 2b3e32d..7f09620 100644
--- a/app/src/main/java/invalid/lena/scrcpy/AudioSink.java
+++ b/app/src/main/java/invalid/lena/scrcpy/AudioSink.java
@@ -12,6 +12,8 @@ import android.os.HandlerThread;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.ByteOrder;
+import java.util.ArrayDeque;
+import java.util.Deque;
// Audio output: AudioTrack writing 48 kHz stereo 16-bit PCM. The
// upstream feed is either raw PCM (passthrough) or Opus packets
@@ -21,22 +23,32 @@ public final class AudioSink implements AudioFrames {
private static final int SAMPLE_RATE = 48_000;
private static final int CHANNEL_OUT = AudioFormat.CHANNEL_OUT_STEREO;
private static final int ENCODING = AudioFormat.ENCODING_PCM_16BIT;
+ private static final int MAX_PENDING_OPUS = 16;
// Defaults documented at <https://developer.android.com/reference/android/media/MediaCodec#CSD>.
private static final long DEFAULT_PRE_ROLL_NS = 80_000_000L;
private AudioTrack track;
- private MediaCodec opusCodec;
+ private volatile MediaCodec opusCodec;
private HandlerThread opusThread;
private Handler opusHandler;
private boolean opusConfigured;
private int fourcc;
+ private final Object lock = new Object();
+ private final Runnable onFatalError;
+ private final Deque<Integer> freeOpusInputs = new ArrayDeque<>();
+ private final Deque<byte[]> pendingOpus = new ArrayDeque<>(MAX_PENDING_OPUS);
+
private volatile boolean released;
public volatile long frames; // public read for the status overlay
private long droppedBytes;
private long lastDropLogMs;
+ public AudioSink(Runnable onFatalError) {
+ this.onFatalError = onFatalError;
+ }
+
@Override
public void start(int fourcc) {
this.fourcc = fourcc;
@@ -71,11 +83,14 @@ public final class AudioSink implements AudioFrames {
} catch (IOException e) {
throw new IllegalStateException("audio sink: no opus decoder", e);
}
- opusCodec.setCallback(new MediaCodec.Callback() {
- @Override public void onInputBufferAvailable(MediaCodec mc, int idx) {
- // Frames are queued synchronously from feed(); we
- // don't pull on this callback. This handler exists
- // so MediaCodec's async machinery is wired up.
+ opusCodec.setCallback(new MediaCodec.Callback() {
+ @Override public void onInputBufferAvailable(MediaCodec mc, int idx) {
+ synchronized (lock) {
+ if (released || mc != opusCodec) return;
+ byte[] packet = pendingOpus.pollFirst();
+ if (packet == null) freeOpusInputs.offerLast(idx);
+ else submitOpus(mc, idx, packet);
+ }
}
@Override public void onOutputBufferAvailable(MediaCodec mc, int idx,
MediaCodec.BufferInfo info) {
@@ -146,21 +161,33 @@ public final class AudioSink implements AudioFrames {
return;
}
if (!opusConfigured) return;
- // Encoded packet → MediaCodec input.
- int idx;
- try { idx = opusCodec.dequeueInputBuffer(0); }
- catch (IllegalStateException e) { return; }
- if (idx < 0) {
- // No input buffer right now; drop the packet. Opus is forgiving
- // about gaps for short stalls.
- return;
+ byte[] packet = new byte[len];
+ System.arraycopy(data, off, packet, 0, len);
+ synchronized (lock) {
+ if (released || opusCodec == null) return;
+ Integer idx = freeOpusInputs.pollFirst();
+ if (idx != null) {
+ submitOpus(opusCodec, idx, packet);
+ return;
+ }
+ if (pendingOpus.size() == MAX_PENDING_OPUS) pendingOpus.pollFirst();
+ pendingOpus.offerLast(packet);
}
+ }
+
+ // Must be called with lock held. Async MediaCodec input indices belong
+ // to the codec instance that delivered them.
+ private void submitOpus(MediaCodec codec, int idx, byte[] packet) {
try {
- ByteBuffer in = opusCodec.getInputBuffer(idx);
- if (in == null) return;
+ ByteBuffer in = codec.getInputBuffer(idx);
+ if (in == null || packet.length > in.capacity()) {
+ Log.e("audio sink: opus packet exceeds codec input (%d bytes)", packet.length);
+ if (onFatalError != null) onFatalError.run();
+ return;
+ }
in.clear();
- in.put(data, off, len);
- opusCodec.queueInputBuffer(idx, 0, len, 0, 0);
+ in.put(packet);
+ codec.queueInputBuffer(idx, 0, packet.length, 0, 0);
frames++;
} catch (IllegalStateException e) {
Log.w("audio sink: opus queueInputBuffer: %s", e);
@@ -194,9 +221,14 @@ public final class AudioSink implements AudioFrames {
@Override
public void release() {
- released = true;
- MediaCodec c = opusCodec;
- opusCodec = null;
+ MediaCodec c;
+ synchronized (lock) {
+ released = true;
+ c = opusCodec;
+ opusCodec = null;
+ freeOpusInputs.clear();
+ pendingOpus.clear();
+ }
HandlerThread ht = opusThread;
opusThread = null;
opusHandler = null;
diff --git a/app/src/main/java/invalid/lena/scrcpy/AudioStream.java b/app/src/main/java/invalid/lena/scrcpy/AudioStream.java
index 39b184d..7fe33f0 100644
--- a/app/src/main/java/invalid/lena/scrcpy/AudioStream.java
+++ b/app/src/main/java/invalid/lena/scrcpy/AudioStream.java
@@ -42,7 +42,12 @@ public final class AudioStream {
public void stop() {
stop = true;
- if (thread != null) thread.interrupt();
+ Thread t = thread;
+ if (t == null) return;
+ t.interrupt();
+ if (t == Thread.currentThread()) return;
+ try { t.join(1_000); }
+ catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}
public void run() {
diff --git a/app/src/main/java/invalid/lena/scrcpy/ControlMessages.java b/app/src/main/java/invalid/lena/scrcpy/ControlMessages.java
index bd983ed..aba2ca8 100644
--- a/app/src/main/java/invalid/lena/scrcpy/ControlMessages.java
+++ b/app/src/main/java/invalid/lena/scrcpy/ControlMessages.java
@@ -8,6 +8,8 @@ import java.nio.charset.StandardCharsets;
// unit-tested without android.* on the classpath.
public final class ControlMessages {
+ static final int MAX_CLIPBOARD_BYTES = 1 << 20;
+
public static final int TYPE_INJECT_KEYCODE = 0;
public static final int TYPE_INJECT_TEXT = 1;
public static final int TYPE_INJECT_TOUCH_EVENT = 2;
@@ -60,7 +62,13 @@ public final class ControlMessages {
}
public static byte[] setClipboard(long sequence, boolean paste, String text) {
+ if (text.length() > MAX_CLIPBOARD_BYTES) {
+ throw new IllegalArgumentException("clipboard text is too large");
+ }
byte[] data = text.getBytes(StandardCharsets.UTF_8);
+ if (data.length > MAX_CLIPBOARD_BYTES) {
+ throw new IllegalArgumentException("clipboard UTF-8 data is too large");
+ }
byte[] m = new byte[1 + 8 + 1 + 4 + data.length];
m[0] = TYPE_SET_CLIPBOARD;
Wire.writeBe64(m, 1, sequence);
diff --git a/app/src/main/java/invalid/lena/scrcpy/ControlStream.java b/app/src/main/java/invalid/lena/scrcpy/ControlStream.java
index eea6c7d..a231b92 100644
--- a/app/src/main/java/invalid/lena/scrcpy/ControlStream.java
+++ b/app/src/main/java/invalid/lena/scrcpy/ControlStream.java
@@ -5,6 +5,7 @@ import java.io.InputStream;
import java.io.OutputStream;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.atomic.AtomicBoolean;
// Bidirectional bridge to the scrcpy control socket.
//
@@ -36,7 +37,9 @@ public final class ControlStream {
private final InputStream in;
private final OutputStream out;
+ private final Runnable onFatalError;
private final LinkedBlockingDeque<byte[]> outbox = new LinkedBlockingDeque<>(MAX_QUEUED);
+ private final AtomicBoolean fatalReported = new AtomicBoolean();
private volatile InboundSink sink;
private Thread writer;
@@ -44,8 +47,13 @@ public final class ControlStream {
private volatile boolean stop;
public ControlStream(InputStream in, OutputStream out) {
+ this(in, out, null);
+ }
+
+ public ControlStream(InputStream in, OutputStream out, Runnable onFatalError) {
this.in = in;
this.out = out;
+ this.onFatalError = onFatalError;
}
// Wired after construction: the Controller that consumes inbound
@@ -64,8 +72,8 @@ public final class ControlStream {
public void stop() {
stop = true;
outbox.clear();
- if (writer != null) writer.interrupt();
- if (reader != null) reader.interrupt();
+ stopThread(writer);
+ stopThread(reader);
}
public void send(byte[] msg) {
@@ -79,14 +87,24 @@ public final class ControlStream {
if (outbox.remove(b)) break;
}
}
- if (!outbox.offerLast(msg)) {
- Log.w("control: outbox full, dropping msg type=%d", msg.length > 0 ? msg[0] & 0xff : -1);
- }
+ if (outbox.offerLast(msg)) return;
+
+ // A queue containing only state transitions is unhealthy, but
+ // silently dropping UP/CANCEL/key events leaves input stuck on the
+ // target. Fail the control stream visibly instead.
+ stop = true;
+ if (writer != null) writer.interrupt();
+ if (reader != null) reader.interrupt();
+ try { out.close(); } catch (IOException ignored) {}
+ try { in.close(); } catch (IOException ignored) {}
+ Log.e("control: outbox saturated with non-droppable events");
+ reportFatal();
}
// ---- writer ----
public void runWriter() {
+ if (writer == null) writer = Thread.currentThread();
try {
while (!stop) {
byte[] msg = outbox.takeFirst();
@@ -95,7 +113,10 @@ public final class ControlStream {
}
} catch (InterruptedException ignored) {
} catch (IOException e) {
- if (!stop) Log.e(e, "control writer");
+ if (!stop) {
+ Log.e(e, "control writer");
+ reportFatal();
+ }
} finally {
Log.i("control writer: end");
}
@@ -104,6 +125,7 @@ public final class ControlStream {
// ---- reader ----
public void runReader() {
+ if (reader == null) reader = Thread.currentThread();
try {
byte[] tmp = new byte[12];
while (!stop) {
@@ -144,11 +166,29 @@ public final class ControlStream {
}
}
} catch (IOException e) {
- if (!stop) Log.e(e, "control reader");
+ if (!stop) {
+ Log.e(e, "control reader");
+ reportFatal();
+ }
} catch (Exception e) {
Log.e(e, "control reader unexpected");
+ if (!stop) reportFatal();
} finally {
Log.i("control reader: end");
}
}
+
+ private void reportFatal() {
+ if (onFatalError != null && fatalReported.compareAndSet(false, true)) {
+ onFatalError.run();
+ }
+ }
+
+ private static void stopThread(Thread t) {
+ if (t == null) return;
+ t.interrupt();
+ if (t == Thread.currentThread()) return;
+ try { t.join(1_000); }
+ catch (InterruptedException e) { Thread.currentThread().interrupt(); }
+ }
}
diff --git a/app/src/main/java/invalid/lena/scrcpy/Controller.java b/app/src/main/java/invalid/lena/scrcpy/Controller.java
index 188ac34..242d6a5 100644
--- a/app/src/main/java/invalid/lena/scrcpy/Controller.java
+++ b/app/src/main/java/invalid/lena/scrcpy/Controller.java
@@ -170,7 +170,11 @@ public final class Controller implements ControlStream.InboundSink {
}
private void sendSetClipboard(String text, boolean paste) {
- sender.accept(ControlMessages.setClipboard(/* sequence */ 0L, paste, text));
- lastEvent = "clip " + text.length() + " chars";
+ try {
+ sender.accept(ControlMessages.setClipboard(/* sequence */ 0L, paste, text));
+ lastEvent = "clip " + text.length() + " chars";
+ } catch (IllegalArgumentException e) {
+ Log.w("clipboard not sent: %s", e.getMessage());
+ }
}
}
diff --git a/app/src/main/java/invalid/lena/scrcpy/Devices.java b/app/src/main/java/invalid/lena/scrcpy/Devices.java
index 2883025..c3e5141 100644
--- a/app/src/main/java/invalid/lena/scrcpy/Devices.java
+++ b/app/src/main/java/invalid/lena/scrcpy/Devices.java
@@ -6,6 +6,7 @@ import org.json.JSONArray;
import org.json.JSONObject;
import java.io.File;
+import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.util.ArrayList;
@@ -97,18 +98,14 @@ public final class Devices {
}
}
- public static void save(Context ctx, List<Device> devices) {
+ public static void save(Context ctx, List<Device> devices) throws IOException {
File f = new File(ctx.getFilesDir(), FILE);
- try {
- AtomicFiles.write(f, serialize(devices).getBytes(StandardCharsets.UTF_8));
- } catch (Exception e) {
- Log.e(e, "devices: save failed");
- }
+ AtomicFiles.write(f, serialize(devices).getBytes(StandardCharsets.UTF_8));
}
// Add or replace by host+port. Context-rooted; the in-place helper
// is the test seam.
- public static List<Device> upsert(Context ctx, Device d) {
+ public static synchronized List<Device> upsert(Context ctx, Device d) throws IOException {
List<Device> list = load(ctx);
list.removeIf(d::equals);
list.add(d);
@@ -117,7 +114,7 @@ public final class Devices {
}
// Remove the matching device (by host+port). Returns the updated list.
- public static List<Device> remove(Context ctx, Device d) {
+ public static synchronized List<Device> remove(Context ctx, Device d) throws IOException {
List<Device> list = load(ctx);
list.removeIf(d::equals);
save(ctx, list);
diff --git a/app/src/main/java/invalid/lena/scrcpy/Main.java b/app/src/main/java/invalid/lena/scrcpy/Main.java
index 104f57f..2f95375 100644
--- a/app/src/main/java/invalid/lena/scrcpy/Main.java
+++ b/app/src/main/java/invalid/lena/scrcpy/Main.java
@@ -83,12 +83,18 @@ public final class Main extends Activity {
new AlertDialog.Builder(this)
.setTitle(R.string.forget_device)
.setMessage(d.host + ":" + d.port)
- .setPositiveButton(android.R.string.ok, (dlg, w) -> {
- Log.i("forget device: %s", d);
- List<Devices.Device> updated = Devices.remove(this, d);
- adapter.clear();
- adapter.addAll(updated);
- adapter.notifyDataSetChanged();
+ .setPositiveButton(android.R.string.ok, (dlg, w) -> {
+ Log.i("forget device: %s", d);
+ try {
+ List<Devices.Device> updated = Devices.remove(this, d);
+ adapter.clear();
+ adapter.addAll(updated);
+ adapter.notifyDataSetChanged();
+ } catch (Exception e) {
+ Log.e(e, "forget device: save failed");
+ Toast.makeText(this, "could not save device list",
+ Toast.LENGTH_LONG).show();
+ }
})
.setNegativeButton(android.R.string.cancel, null)
.show();
diff --git a/app/src/main/java/invalid/lena/scrcpy/Mirror.java b/app/src/main/java/invalid/lena/scrcpy/Mirror.java
index 2c361da..2a9d390 100644
--- a/app/src/main/java/invalid/lena/scrcpy/Mirror.java
+++ b/app/src/main/java/invalid/lena/scrcpy/Mirror.java
@@ -60,6 +60,10 @@ public final class Mirror extends Activity {
private Devices.Device target;
private Session session;
private Surface currentSurface;
+ private boolean ownsCurrentSurface;
+ private volatile boolean destroyed;
+ private long sessionGeneration;
+ private boolean stoppingSession;
private State state = State.CONNECTING;
private int connectedW, connectedH;
@@ -138,6 +142,7 @@ public final class Mirror extends Activity {
try {
Adb a = Adb.getInstance(this);
runOnUiThread(() -> {
+ if (destroyed) return;
adb = a;
if (session == null && currentSurface != null) {
startSession(currentSurface);
@@ -146,6 +151,7 @@ public final class Mirror extends Activity {
} catch (Exception e) {
Log.e(e, "mirror: adb init");
runOnUiThread(() -> {
+ if (destroyed) return;
Toast.makeText(this, "adb init: " + e.getMessage(),
Toast.LENGTH_LONG).show();
finish();
@@ -161,7 +167,7 @@ public final class Mirror extends Activity {
@Override
public void onSurfaceTextureAvailable(SurfaceTexture st, int w, int h) {
Log.i("mirror: texture available %dx%d", w, h);
- attachSurface(new Surface(st), w, h);
+ attachSurface(new Surface(st), true, w, h);
}
@Override
public void onSurfaceTextureSizeChanged(SurfaceTexture st, int w, int h) {
@@ -188,7 +194,7 @@ public final class Mirror extends Activity {
@Override
public void surfaceCreated(SurfaceHolder holder) {
Log.i("mirror: surface created");
- attachSurface(holder.getSurface(), 0, 0);
+ attachSurface(holder.getSurface(), false, 0, 0);
}
@Override
public void surfaceChanged(SurfaceHolder holder, int format, int w, int h) {
@@ -204,30 +210,43 @@ public final class Mirror extends Activity {
// ---- session driver ----
- private void attachSurface(Surface s, int w, int h) {
+ private void attachSurface(Surface s, boolean owned, int w, int h) {
+ if (session != null && currentSurface != null) session.swapSurface(null);
+ releaseOwnedSurface();
currentSurface = s;
+ ownsCurrentSurface = owned;
if (session != null) {
session.swapSurface(s);
if (w > 0) session.setViewSize(w, h);
return;
}
if (adb == null) return; // adb-init thread will start the session
+ if (stoppingSession) return; // stop worker starts the latest target
if (state == State.DISCONNECTED) return; // wait for user to tap reconnect
startSession(s);
if (w > 0) session.setViewSize(w, h);
}
private void detachSurface() {
- currentSurface = null;
if (session != null) session.swapSurface(null);
+ releaseOwnedSurface();
+ currentSurface = null;
+ }
+
+ private void releaseOwnedSurface() {
+ if (ownsCurrentSurface && currentSurface != null) currentSurface.release();
+ ownsCurrentSurface = false;
}
private void startSession(Surface s) {
+ if (destroyed) return;
state = State.CONNECTING;
updateStatusBar();
+ long generation = ++sessionGeneration;
session = new Session(this, adb, target, s, new Session.Listener() {
@Override public void onConnected(int w, int h) {
runOnUiThread(() -> {
+ if (destroyed || generation != sessionGeneration) return;
state = State.CONNECTED;
connectedW = w; connectedH = h;
updateStatusBar();
@@ -236,12 +255,14 @@ public final class Mirror extends Activity {
@Override public void onReconnecting() {
Log.i("mirror: link lost, reconnecting");
runOnUiThread(() -> {
+ if (destroyed || generation != sessionGeneration) return;
state = State.CONNECTING;
updateStatusBar();
});
}
@Override public void onError(Throwable t) {
runOnUiThread(() -> {
+ if (destroyed || generation != sessionGeneration) return;
Toast.makeText(Mirror.this,
"session error: " + t.getMessage(), Toast.LENGTH_LONG).show();
});
@@ -249,6 +270,7 @@ public final class Mirror extends Activity {
@Override public void onStopped() {
Log.i("mirror: session stopped");
runOnUiThread(() -> {
+ if (destroyed || generation != sessionGeneration) return;
state = State.DISCONNECTED;
updateStatusBar();
});
@@ -286,10 +308,13 @@ public final class Mirror extends Activity {
private void reconnect() {
Log.i("mirror: reconnect tapped");
- Session old = session;
- session = null;
state = State.CONNECTING;
updateStatusBar();
+ if (stoppingSession) return;
+ stoppingSession = true;
+ Session old = session;
+ session = null;
+ sessionGeneration++;
// Stop the old session off the UI thread (teardown closes
// sockets and joins the server's log pump), THEN start the new
// one. Sequencing matters: both sessions share the singleton
@@ -298,6 +323,8 @@ public final class Mirror extends Activity {
new Thread(() -> {
if (old != null) old.stop();
runOnUiThread(() -> {
+ if (destroyed) return;
+ stoppingSession = false;
if (session != null) return; // another path already started one
if (adb == null || currentSurface == null) return;
startSession(currentSurface);
@@ -305,6 +332,24 @@ public final class Mirror extends Activity {
}, "session-stop").start();
}
+ @Override
+ protected void onNewIntent(Intent intent) {
+ super.onNewIntent(intent);
+ String host = intent.getStringExtra(EXTRA_HOST);
+ int port = intent.getIntExtra(EXTRA_PORT, -1);
+ if (host == null || port <= 0 || port > 65535) {
+ Log.w("mirror: ignoring bad replacement target host=%s port=%d", host, port);
+ return;
+ }
+ if (target.host.equals(host) && target.port == port) return;
+
+ setIntent(intent);
+ target = new Devices.Device(host, port);
+ connectedW = connectedH = 0;
+ if (overlayTarget != null) overlayTarget.setText("target: " + target);
+ reconnect();
+ }
+
private void updateStatusBar() {
if (statusText == null) return;
String s;
@@ -394,8 +439,11 @@ public final class Mirror extends Activity {
@Override
protected void onDestroy() {
- super.onDestroy();
+ destroyed = true;
+ sessionGeneration++;
ui.removeCallbacksAndMessages(null);
+ releaseOwnedSurface();
+ currentSurface = null;
Session s = session;
session = null;
if (s != null) {
@@ -404,6 +452,7 @@ public final class Mirror extends Activity {
new Thread(s::stop, "session-stop").start();
}
stopService(new Intent(this, Sessions.class));
+ super.onDestroy();
}
// Android 13+ requires runtime grant for POST_NOTIFICATIONS. The
diff --git a/app/src/main/java/invalid/lena/scrcpy/MuxRecorder.java b/app/src/main/java/invalid/lena/scrcpy/MuxRecorder.java
index b8d4aaf..6631f66 100644
--- a/app/src/main/java/invalid/lena/scrcpy/MuxRecorder.java
+++ b/app/src/main/java/invalid/lena/scrcpy/MuxRecorder.java
@@ -26,7 +26,7 @@ import java.nio.ByteBuffer;
// illegal). User can re-arm; a fresh output gets the new dimensions.
public final class MuxRecorder implements VideoRecorder {
- private enum State { IDLE, ARMED, RECORDING, ERROR }
+ private enum State { IDLE, ARMED, RECORDING }
private final Object lock = new Object();
@@ -66,8 +66,8 @@ public final class MuxRecorder implements VideoRecorder {
boolean wrote;
synchronized (lock) {
if (state == State.IDLE) return false;
- wrote = (state == State.RECORDING);
- closeMuxerLocked();
+ wrote = state == State.RECORDING && closeMuxerLocked();
+ if (!wrote) deleteOutputLocked();
bytes = bytesWritten;
state = State.IDLE;
outFile = null;
@@ -115,19 +115,21 @@ public final class MuxRecorder implements VideoRecorder {
synchronized (lock) {
switch (state) {
case IDLE:
- case ERROR:
return;
case ARMED:
if (!isKeyframe || csd == null) return;
if (!startMuxerLocked()) {
- state = State.ERROR;
+ failLocked("start");
+ return;
+ }
+ if (!writeSampleLocked(data, ptsUs, true)) {
+ failLocked("write");
return;
}
- writeSampleLocked(data, ptsUs, true);
state = State.RECORDING;
return;
case RECORDING:
- writeSampleLocked(data, ptsUs, isKeyframe);
+ if (!writeSampleLocked(data, ptsUs, isKeyframe)) failLocked("write");
}
}
}
@@ -163,8 +165,8 @@ public final class MuxRecorder implements VideoRecorder {
}
}
- private void writeSampleLocked(byte[] data, long ptsUs, boolean isKeyframe) {
- if (muxer == null) return;
+ private boolean writeSampleLocked(byte[] data, long ptsUs, boolean isKeyframe) {
+ if (muxer == null) return false;
// PTS is rebased to zero so the file is self-contained; some
// players choke on absolute PTS that doesn't start near zero.
if (firstPtsUs < 0) firstPtsUs = ptsUs;
@@ -181,18 +183,36 @@ public final class MuxRecorder implements VideoRecorder {
isKeyframe ? MediaCodec.BUFFER_FLAG_KEY_FRAME : 0);
muxer.writeSampleData(trackIdx, buf, info);
bytesWritten += data.length;
+ return true;
} catch (IllegalStateException e) {
Log.w("rec: writeSampleData: %s", e);
+ return false;
}
}
- private void closeMuxerLocked() {
+ private boolean closeMuxerLocked() {
MediaMuxer m = muxer;
muxer = null;
trackIdx = -1;
- if (m == null) return;
- try { m.stop(); } catch (Exception ignored) {}
+ if (m == null) return false;
+ boolean ok = true;
+ try { m.stop(); } catch (Exception e) { ok = false; Log.w("rec: muxer stop: %s", e); }
try { m.release(); } catch (Exception ignored) {}
+ return ok;
+ }
+
+ private void failLocked(String operation) {
+ Log.e("rec: %s failed, discarding output", operation);
+ closeMuxerLocked();
+ deleteOutputLocked();
+ state = State.IDLE;
+ outFile = null;
+ }
+
+ private void deleteOutputLocked() {
+ if (outFile != null && outFile.exists() && !outFile.delete()) {
+ Log.w("rec: could not delete incomplete output %s", outFile);
+ }
}
private static String mimeFor(int fourcc) {
diff --git a/app/src/main/java/invalid/lena/scrcpy/Server.java b/app/src/main/java/invalid/lena/scrcpy/Server.java
index 6ae752f..284dd6c 100644
--- a/app/src/main/java/invalid/lena/scrcpy/Server.java
+++ b/app/src/main/java/invalid/lena/scrcpy/Server.java
@@ -7,12 +7,11 @@ import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
import java.io.OutputStream;
+import java.net.ConnectException;
import java.nio.charset.StandardCharsets;
import java.security.SecureRandom;
import java.util.ArrayList;
import java.util.List;
-import java.util.concurrent.atomic.AtomicBoolean;
-import java.util.concurrent.atomic.AtomicReference;
import io.github.muntashirakon.adb.AdbStream;
@@ -38,21 +37,8 @@ public final class Server {
private static final String ASSET_JAR = "scrcpy-server.jar";
private static final String ASSET_VERSION = "scrcpy-server.version";
private static final int FILE_MODE = 0100644; // regular file, 0644
- // Server forks a CleanUp helper before opening its abstract sockets;
- // on slow emulators that takes several seconds, so the budget needs
- // to be generous. The deadline is per stream: in practice only the
- // first open waits (the server is still starting) and the other two
- // dial instantly, so a healthy bring-up stays well below the e2e
- // test deadline (test-rig/e2e.sh's E2E_DEADLINE, default 60 s).
- private static final long OPEN_DEADLINE_MS = 20_000;
- private static final int OPEN_BACKOFF_MS = 100;
- // Per-attempt timeout for adb.openAbstract. libadb-android's
- // AdbConnection.open() blocks on a naked stream.wait() with no
- // loop or timeout - a missed notification (response arrives before
- // we park) wedges the call forever. Wrap each attempt with a
- // timeout + interrupt so the loop can move on.
- private static final long OPEN_ATTEMPT_TIMEOUT_MS = 800;
-
+ private static final long LISTENER_DEADLINE_MS = 20_000;
+ private static final long LISTENER_RETRY_MS = 100;
public static final class Streams {
public final AdbStream videoAds, audioAds, controlAds;
public final InputStream videoIn, audioIn, controlIn;
@@ -84,6 +70,7 @@ public final class Server {
}
public Streams bringUp() throws Exception {
+ serverEof = false;
String version = readVersion();
long pushed = push();
Log.i("push %s bytes=%d", REMOTE_PATH, pushed);
@@ -93,24 +80,40 @@ public final class Server {
Log.i("spawn server ver=%s scid=%s", version, scid);
Log.i("cmdline: %s", cmd);
- shell = adb.openShell(cmd);
- shellPump = new Thread(() -> pump(shell.openInputStream()), "server-stdout");
- shellPump.setDaemon(true);
- shellPump.start();
+ AdbStream va = null, aa = null, ca = null;
+ boolean committed = false;
+ try {
+ shell = adb.openShell(cmd);
+ AdbStream shellRef = shell;
+ shellPump = new Thread(() -> pump(shellRef.openInputStream()), "server-stdout");
+ shellPump.setDaemon(true);
+ shellPump.start();
- AdbStream va = openAbstract(scid);
- AdbStream aa = openAbstract(scid);
- AdbStream ca = openAbstract(scid);
+ // These accepts are ordered. If one times out, the whole ADB
+ // connection is discarded by Session; retrying an individual
+ // open could shift video/audio/control onto the wrong sockets.
+ va = openAbstract(scid);
+ aa = openAbstract(scid);
+ ca = openAbstract(scid);
- InputStream vi = va.openInputStream();
- InputStream ai = aa.openInputStream();
- InputStream ci = ca.openInputStream();
- OutputStream co = ca.openOutputStream();
+ InputStream vi = va.openInputStream();
+ InputStream ai = aa.openInputStream();
+ InputStream ci = ca.openInputStream();
+ OutputStream co = ca.openOutputStream();
- String name = readDeviceMeta(vi);
- Log.i("device name=%s", name);
- streams = new Streams(va, aa, ca, vi, ai, ci, co, name, scid, version);
- return streams;
+ String name = readDeviceMeta(vi);
+ Log.i("device name=%s", name);
+ streams = new Streams(va, aa, ca, vi, ai, ci, co, name, scid, version);
+ committed = true;
+ return streams;
+ } finally {
+ if (!committed) {
+ closeQuietly(va);
+ closeQuietly(aa);
+ closeQuietly(ca);
+ closeShell();
+ }
+ }
}
// Idempotent. Closes the three media/control streams first (lets
@@ -127,6 +130,10 @@ public final class Server {
closeQuietly(streams.controlAds);
streams = null;
}
+ closeShell();
+ }
+
+ private void closeShell() {
Thread t = shellPump;
AdbStream s = shell;
shell = null;
@@ -198,67 +205,26 @@ public final class Server {
private AdbStream openAbstract(String scid) throws Exception {
String name = "scrcpy_" + scid;
- long deadline = System.currentTimeMillis() + OPEN_DEADLINE_MS;
- Throwable last = null;
- for (int attempt = 1; System.currentTimeMillis() < deadline; attempt++) {
- // If the shell stream closed (server died early - usually with
- // an error printed to stderr that our pump captured), bail
- // immediately. Otherwise we'd spin out the deadline on a
- // dead target.
+ long deadline = System.currentTimeMillis() + LISTENER_DEADLINE_MS;
+ ConnectException last = null;
+ while (System.currentTimeMillis() < deadline) {
if (serverEof) {
- throw new IOException("server exited before opening "
- + name + " (see 'server:' logs)",
- last instanceof Exception ? (Exception) last : null);
+ throw new IOException("server exited before opening " + name, last);
}
try {
- AdbStream s = openAbstractOnce(name, OPEN_ATTEMPT_TIMEOUT_MS);
- Log.i("openAbstract %s ok (attempt %d)", name, attempt);
- return s;
- } catch (Throwable t) {
- last = t;
- Thread.sleep(OPEN_BACKOFF_MS);
- }
- }
- throw new IOException("openAbstract " + name + " failed after "
- + (OPEN_DEADLINE_MS / 1000) + " s",
- last instanceof Exception ? (Exception) last : null);
- }
-
- // Wraps a single adb.openAbstract call with a hard timeout. The
- // upstream call can wedge forever on a missed notify in its naked
- // stream.wait(); we run it on a daemon thread and join with timeout,
- // interrupt on overrun, and let the caller retry. If the orphaned
- // call succeeds after we gave up, the stream must be closed:
- // keeping it would silently consume one of the scrcpy server's
- // three accepts and shift every later dial off by one.
- private AdbStream openAbstractOnce(String name, long timeoutMs) throws Exception {
- AtomicReference<AdbStream> result = new AtomicReference<>();
- AtomicReference<Throwable> err = new AtomicReference<>();
- AtomicBoolean abandoned = new AtomicBoolean();
- Thread t = new Thread(() -> {
- try {
- AdbStream s = adb.openAbstract(name);
- result.set(s);
- if (abandoned.get()) closeQuietly(s);
- } catch (Throwable ex) {
- err.set(ex);
+ AdbStream stream = adb.openAbstract(name);
+ Log.i("openAbstract %s ok", name);
+ return stream;
+ } catch (ConnectException e) {
+ // A rejected OPEN consumed no server accept. Retry while the
+ // server creates its listener; timeout failures remain fatal
+ // because their acceptance state is ambiguous.
+ last = e;
+ Thread.sleep(LISTENER_RETRY_MS);
}
- }, "openAbstract-" + name);
- t.setDaemon(true);
- t.start();
- t.join(timeoutMs);
- if (t.isAlive()) {
- abandoned.set(true);
- t.interrupt();
- // One side of the publish/abandon race closes the stream;
- // closeQuietly tolerates both doing it.
- closeQuietly(result.get());
- throw new IOException("openAbstract " + name + " timed out");
}
- Throwable ex = err.get();
- if (ex instanceof Exception) throw (Exception) ex;
- if (ex != null) throw new RuntimeException(ex);
- return result.get();
+ throw new IOException("server did not open " + name + " within "
+ + LISTENER_DEADLINE_MS + " ms", last);
}
private static String readDeviceMeta(InputStream in) throws IOException {
diff --git a/app/src/main/java/invalid/lena/scrcpy/Session.java b/app/src/main/java/invalid/lena/scrcpy/Session.java
index 774fdc1..947930b 100644
--- a/app/src/main/java/invalid/lena/scrcpy/Session.java
+++ b/app/src/main/java/invalid/lena/scrcpy/Session.java
@@ -238,48 +238,67 @@ public final class Session {
}
Log.i("adb connect ok");
- Server srv = new Server(ctx, adb);
- Server.Streams s = srv.bringUp();
-
- ControlStream cs = new ControlStream(s.controlIn, s.controlOut);
- Controller ctrl = new Controller(ctx, cs::send);
- cs.setInboundSink(ctrl);
-
- AudioSink ak = new AudioSink();
- AudioStream as = new AudioStream(s.audioIn, ak);
-
- VideoSink vk = new VideoSink(surface);
- final Controller ctrlRef = ctrl; // capture for SizeListener
- AtomicBoolean reported = new AtomicBoolean();
- VideoStream vs = new VideoStream(s.videoIn, vk, (w, h) -> {
- ctrlRef.setTargetSize(w, h);
- if (reported.compareAndSet(false, true) && listener != null) {
- listener.onConnected(w, h);
- }
- });
- vs.setOnEnd(this::onVideoEnded);
- MuxRecorder rec = new MuxRecorder();
- vs.setRecorder(rec);
+ Server srv = null;
+ ControlStream cs = null;
+ Controller ctrl = null;
+ AudioSink ak = null;
+ AudioStream as = null;
+ VideoSink vk = null;
+ VideoStream vs = null;
+ MuxRecorder rec = null;
+ boolean installed = false;
+ try {
+ srv = new Server(ctx, adb);
+ Server.Streams s = srv.bringUp();
+
+ cs = new ControlStream(s.controlIn, s.controlOut, this::onVideoEnded);
+ ctrl = new Controller(ctx, cs::send);
+ cs.setInboundSink(ctrl);
+
+ ak = new AudioSink(this::onVideoEnded);
+ as = new AudioStream(s.audioIn, ak);
+
+ vk = new VideoSink(surface, this::onVideoEnded);
+ Controller ctrlRef = ctrl;
+ AtomicBoolean reported = new AtomicBoolean();
+ vs = new VideoStream(s.videoIn, vk,
+ (w, h) -> reportConnected(ctrlRef, reported, w, h));
+ vs.setOnEnd(this::onVideoEnded);
+ rec = new MuxRecorder();
+ vs.setRecorder(rec);
+
+ // Start locals before publishing them. stop() either tears down
+ // a previously installed generation or marks this generation for
+ // rollback; it can never release objects that bringUp then starts.
+ cs.start();
+ as.start();
+ vs.start();
- synchronized (this) {
- if (stopped) {
- tearDownLocals(srv, cs, ctrl, ak, as, vk, vs, rec);
- throw new IOException("session: stopped during bring-up");
+ synchronized (this) {
+ if (stopped) throw new IOException("session: stopped during bring-up");
+ server = srv;
+ controlStream = cs;
+ controller = ctrl;
+ audioSink = ak;
+ audioStream = as;
+ videoSink = vk;
+ videoStream = vs;
+ recorder = rec;
+ if (pendingViewW > 0) ctrl.setViewSize(pendingViewW, pendingViewH);
+ installed = true;
}
- server = srv;
- controlStream = cs;
- controller = ctrl;
- audioSink = ak;
- audioStream = as;
- videoSink = vk;
- videoStream = vs;
- recorder = rec;
- if (pendingViewW > 0) ctrl.setViewSize(pendingViewW, pendingViewH);
+ } finally {
+ if (!installed) tearDownLocals(srv, cs, ctrl, ak, as, vk, vs, rec);
}
+ }
- cs.start();
- as.start();
- vs.start();
+ private synchronized void reportConnected(Controller ctrl, AtomicBoolean reported,
+ int w, int h) {
+ if (stopped) return;
+ ctrl.setTargetSize(w, h);
+ if (reported.compareAndSet(false, true) && listener != null) {
+ listener.onConnected(w, h);
+ }
}
// Fired on the video-reader thread when its read loop exits (EOF,
@@ -308,12 +327,15 @@ public final class Session {
AudioSink ak, AudioStream as,
VideoSink vk, VideoStream vs, MuxRecorder rec) {
if (rec != null) rec.close();
+ // Closing the owning ADB streams first unblocks readers. Join them
+ // before releasing their sinks so no callback can recreate resources
+ // after teardown.
+ if (srv != null) srv.close();
if (vs != null) vs.stop();
- if (vk != null) vk.release();
if (as != null) as.stop();
- if (ak != null) ak.release();
if (cs != null) cs.stop();
+ if (vk != null) vk.release();
+ if (ak != null) ak.release();
if (ctrl != null) ctrl.release();
- if (srv != null) srv.close();
}
}
diff --git a/app/src/main/java/invalid/lena/scrcpy/Sessions.java b/app/src/main/java/invalid/lena/scrcpy/Sessions.java
index 95c0a1f..91144be 100644
--- a/app/src/main/java/invalid/lena/scrcpy/Sessions.java
+++ b/app/src/main/java/invalid/lena/scrcpy/Sessions.java
@@ -16,8 +16,8 @@ import android.os.IBinder;
// briefly backgrounded (rotation, IME, swipe-to-home) without the
// scrcpy server tearing down.
//
-// Type is FOREGROUND_SERVICE_DATA_SYNC: we are pulling a continuous
-// data stream (encoded video + raw PCM) from another device.
+// Type is FOREGROUND_SERVICE_MEDIA_PLAYBACK: the app continuously plays
+// remote audio/video for as long as the user keeps a mirror session open.
//
// No binder API - the service does not own Session. Mirror owns it.
public final class Sessions extends Service {
@@ -41,7 +41,7 @@ public final class Sessions extends Service {
.setContentIntent(reopenIntent())
.setOngoing(true)
.build();
- startForeground(NOTIF_ID, n, ServiceInfo.FOREGROUND_SERVICE_TYPE_DATA_SYNC);
+ startForeground(NOTIF_ID, n, ServiceInfo.FOREGROUND_SERVICE_TYPE_MEDIA_PLAYBACK);
return START_NOT_STICKY;
}
diff --git a/app/src/main/java/invalid/lena/scrcpy/VideoSink.java b/app/src/main/java/invalid/lena/scrcpy/VideoSink.java
index 2d69339..ff8f449 100644
--- a/app/src/main/java/invalid/lena/scrcpy/VideoSink.java
+++ b/app/src/main/java/invalid/lena/scrcpy/VideoSink.java
@@ -20,8 +20,8 @@ import java.util.Iterator;
// plus a corresponding pool of free input buffer indices.
//
// Back-pressure policy: when no input buffer is free and the pending
-// queue is full, the oldest pending non-config frame is dropped. Config
-// frames (CSD) always survive - the decoder cannot start without them.
+// queue is full, replace the oldest frame of the same kind, or the oldest
+// frame overall. The queue remains bounded even if a peer floods CSD.
// Keyframes are not distinguishable here (the flag stays in VideoStream)
// so they drop like any delta frame; the picture heals at the next one.
//
@@ -36,8 +36,9 @@ public final class VideoSink implements VideoFrames {
private static final int MAX_PENDING = 8;
private volatile Surface surface;
+ private final Runnable onFatalError;
public volatile long frames; // public read for the status overlay
- private MediaCodec codec;
+ private volatile MediaCodec codec;
private HandlerThread handlerThread;
private Handler handler;
@@ -56,8 +57,9 @@ public final class VideoSink implements VideoFrames {
Frame(byte[] d, long pts, boolean cfg) { data = d; ptsUs = pts; isConfig = cfg; }
}
- public VideoSink(Surface surface) {
+ public VideoSink(Surface surface, Runnable onFatalError) {
this.surface = surface;
+ this.onFatalError = onFatalError;
}
// Swap the output Surface without rebuilding MediaCodec. Passing null
@@ -97,7 +99,7 @@ public final class VideoSink implements VideoFrames {
codec = MediaCodec.createDecoderByType(mime);
codec.setCallback(new MediaCodec.Callback() {
@Override public void onInputBufferAvailable(MediaCodec mc, int idx) {
- onFreeInput(idx);
+ onFreeInput(mc, idx);
}
@Override public void onOutputBufferAvailable(MediaCodec mc, int idx, MediaCodec.BufferInfo info) {
try {
@@ -118,6 +120,7 @@ public final class VideoSink implements VideoFrames {
}
@Override public void onError(MediaCodec mc, MediaCodec.CodecException e) {
Log.e(e, "video sink: codec error");
+ if (mc == codec && onFatalError != null) onFatalError.run();
}
@Override public void onOutputFormatChanged(MediaCodec mc, MediaFormat fmt) {
Log.i("video sink: output format %s", fmt);
@@ -137,20 +140,26 @@ public final class VideoSink implements VideoFrames {
if (released) return;
// Try to drain immediately if there's a free input.
while (!pending.isEmpty() && !freeInputs.isEmpty()) {
- submit(pending.pollFirst(), freeInputs.pollFirst());
+ submit(codec, pending.pollFirst(), freeInputs.pollFirst());
}
if (!freeInputs.isEmpty()) {
- submit(new Frame(data, ptsUs, isConfig), freeInputs.pollFirst());
+ submit(codec, new Frame(data, ptsUs, isConfig), freeInputs.pollFirst());
return;
}
- // Queue, with bounded drop policy on non-config frames.
- if (pending.size() >= MAX_PENDING && !isConfig) {
- // Drop the oldest non-config frame to avoid stalling
- // forever. Config frames must survive: the decoder
- // cannot start without its CSD.
+ // Keep the queue strictly bounded. A newer config frame replaces
+ // an older one; retaining every config packet lets a hostile peer
+ // turn the queue into an unbounded allocation sink.
+ if (pending.size() >= MAX_PENDING) {
+ boolean removed = false;
for (Iterator<Frame> it = pending.iterator(); it.hasNext(); ) {
- if (!it.next().isConfig) { it.remove(); break; }
+ Frame f = it.next();
+ if (f.isConfig == isConfig) {
+ it.remove();
+ removed = true;
+ break;
+ }
}
+ if (!removed) pending.pollFirst();
}
pending.offerLast(new Frame(data, ptsUs, isConfig));
}
@@ -196,10 +205,10 @@ public final class VideoSink implements VideoFrames {
}
// Internal - runs on the MediaCodec callback thread.
- private void onFreeInput(int idx) {
+ private void onFreeInput(MediaCodec mc, int idx) {
synchronized (lock) {
- if (released) return;
- if (!pending.isEmpty()) submit(pending.pollFirst(), idx);
+ if (released || mc != codec) return;
+ if (!pending.isEmpty()) submit(mc, pending.pollFirst(), idx);
else freeInputs.offerLast(idx);
}
}
@@ -207,15 +216,19 @@ public final class VideoSink implements VideoFrames {
// Must be called with `lock` held. codec can be null mid-reconfigure
// (teardownCodec runs unlocked); the frame is dropped like any other
// back-pressure casualty.
- private void submit(Frame f, int idx) {
- if (codec == null) return;
+ private void submit(MediaCodec mc, Frame f, int idx) {
+ if (mc == null || mc != codec) return;
try {
- ByteBuffer buf = codec.getInputBuffer(idx);
- if (buf == null) return;
+ ByteBuffer buf = mc.getInputBuffer(idx);
+ if (buf == null || f.data.length > buf.capacity()) {
+ Log.e("video sink: frame exceeds codec input (%d bytes)", f.data.length);
+ if (onFatalError != null) onFatalError.run();
+ return;
+ }
buf.clear();
buf.put(f.data);
int flags = f.isConfig ? MediaCodec.BUFFER_FLAG_CODEC_CONFIG : 0;
- codec.queueInputBuffer(idx, 0, f.data.length, f.ptsUs, flags);
+ mc.queueInputBuffer(idx, 0, f.data.length, f.ptsUs, flags);
} catch (IllegalStateException e) {
Log.w("video sink: queueInputBuffer: %s", e);
}
diff --git a/app/src/main/java/invalid/lena/scrcpy/VideoStream.java b/app/src/main/java/invalid/lena/scrcpy/VideoStream.java
index 4eb62d3..7ecb3f1 100644
--- a/app/src/main/java/invalid/lena/scrcpy/VideoStream.java
+++ b/app/src/main/java/invalid/lena/scrcpy/VideoStream.java
@@ -61,7 +61,12 @@ public final class VideoStream {
public void stop() {
stop = true;
- if (thread != null) thread.interrupt();
+ Thread t = thread;
+ if (t == null) return;
+ t.interrupt();
+ if (t == Thread.currentThread()) return;
+ try { t.join(1_000); }
+ catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}
public void setRecorder(VideoRecorder r) {
diff --git a/app/src/test/java/invalid/lena/scrcpy/ControlMessagesTest.java b/app/src/test/java/invalid/lena/scrcpy/ControlMessagesTest.java
index a84e18e..1890860 100644
--- a/app/src/test/java/invalid/lena/scrcpy/ControlMessagesTest.java
+++ b/app/src/test/java/invalid/lena/scrcpy/ControlMessagesTest.java
@@ -9,6 +9,12 @@ import java.nio.charset.StandardCharsets;
public class ControlMessagesTest {
+ @Test(expected = IllegalArgumentException.class)
+ public void clipboardRejectsOversizedText() {
+ char[] chars = new char[ControlMessages.MAX_CLIPBOARD_BYTES + 1];
+ ControlMessages.setClipboard(0, false, new String(chars));
+ }
+
@Test
public void touchByteLayout() {
// ACTION_DOWN=0, single finger, target 1080x2400, x=100, y=200, full pressure
diff --git a/app/src/test/java/invalid/lena/scrcpy/ControlStreamTest.java b/app/src/test/java/invalid/lena/scrcpy/ControlStreamTest.java
index cafff81..8573d3b 100644
--- a/app/src/test/java/invalid/lena/scrcpy/ControlStreamTest.java
+++ b/app/src/test/java/invalid/lena/scrcpy/ControlStreamTest.java
@@ -7,6 +7,7 @@ import org.junit.Test;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.io.DataOutputStream;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
@@ -149,4 +150,17 @@ public class ControlStreamTest {
System.arraycopy(bytes, bytes.length - 14, tail, 0, 14);
for (int i = 0; i < 14; i++) assertEquals(key[i], tail[i]);
}
+
+ @Test
+ public void nonDroppableOverflowFailsStream() {
+ AtomicBoolean failed = new AtomicBoolean();
+ ControlStream cs = new ControlStream(
+ new ByteArrayInputStream(new byte[0]), new ByteArrayOutputStream(),
+ () -> failed.set(true));
+
+ byte[] key = ControlMessages.keycode(0, 29, 0, 0);
+ for (int i = 0; i < 257; i++) cs.send(key);
+
+ assertEquals(true, failed.get());
+ }
}