package com.sw.dualscreen.socket; import android.util.Log; import org.json.JSONObject; import java.io.DataInputStream; import java.io.DataOutputStream; import java.net.InetSocketAddress; import java.net.Socket; import java.net.SocketTimeoutException; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; /** * TcpClient - Android ready *

* Features: * - length-prefix protocol (int length + bytes) * - separate read thread and write queue & writer thread * - auto-reconnect with exponential backoff * - heartbeat scheduler * - send queue with optional callback for send result * - auto send "auth" JSON after connection */ public class TcpClient { private static final String TAG = "TcpClient"; // configuration private final String serverIp; private final int serverPort; private final String clientId; // will be sent in auth message private final int connectTimeoutMs; private final long heartbeatIntervalMs; private final long heartbeatTimeoutMs; // socket + streams private Socket socket; private DataOutputStream out; private DataInputStream in; // threads & executors private final ExecutorService writerExecutor = Executors.newSingleThreadExecutor(r -> new Thread(r, "TcpClient-Writer")); private final ExecutorService readerExecutor = Executors.newSingleThreadExecutor(r -> new Thread(r, "TcpClient-Reader")); private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(r -> new Thread(r, "TcpClient-Scheduler")); private final ExecutorService connectExecutor = Executors.newSingleThreadExecutor(r -> new Thread(r, "TcpClient-Connect")); // send queue private final BlockingQueue sendQueue = new LinkedBlockingQueue<>(); // state private final AtomicBoolean running = new AtomicBoolean(false); private final AtomicBoolean connected = new AtomicBoolean(false); private final AtomicBoolean authSent = new AtomicBoolean(false); // reconnection/backoff private final long baseReconnectDelayMs = 1000; // 1s private final long maxReconnectDelayMs = 30_000; // 30s private final AtomicInteger reconnectAttempt = new AtomicInteger(0); // heartbeat task future private ScheduledFuture heartbeatFuture; // listener public interface Listener { void onConnected(); void onDisconnected(Exception e); void onMessage(JSONObject json); void onSendSuccess(JSONObject json); void onSendFailed(JSONObject json, Exception e); } private Listener listener; public void setListener(Listener l) { this.listener = l; } // ctor public TcpClient(String serverIp, int serverPort, String clientId, int connectTimeoutMs, long heartbeatIntervalMs, long heartbeatTimeoutMs) { this.serverIp = serverIp; this.serverPort = serverPort; this.clientId = clientId; this.connectTimeoutMs = connectTimeoutMs; this.heartbeatIntervalMs = heartbeatIntervalMs; this.heartbeatTimeoutMs = heartbeatTimeoutMs; } // start client (will attempt connect) public void start() { if (running.getAndSet(true)) return; scheduleConnect(0); // writer thread drains sendQueue writerExecutor.execute(this::writerLoop); } // stop client and cleanup public void stop() { running.set(false); cancelHeartbeat(); closeSocketQuiet(); writerExecutor.shutdownNow(); readerExecutor.shutdownNow(); scheduler.shutdownNow(); connectExecutor.shutdownNow(); sendQueue.clear(); } // send JSON (queued). Non-blocking. public void send(JSONObject json) { if (!running.get()) return; sendQueue.offer(json); } // AUTH shortcut (immediately send auth JSON) private void sendAuth() { try { JSONObject auth = new JSONObject(); auth.put("type", "auth"); auth.put("clientId", clientId); sendQueue.offer(auth); authSent.set(true); } catch (Exception ignored) { } } // writer thread loop (serializes sends) private void writerLoop() { while (running.get()) { try { JSONObject json = sendQueue.take(); // blocks if (connected.get() && out != null) { try { byte[] data = json.toString().getBytes(); synchronized (out) { out.writeInt(data.length); out.write(data); out.flush(); } if (listener != null) listener.onSendSuccess(json); } catch (Exception e) { if (listener != null) listener.onSendFailed(json, e); // on write failure, attempt reconnect safeCloseAndScheduleReconnect(e); } } else { // not connected: requeue it and wait for connection sendQueue.offer(json); Thread.sleep(500); // avoid busy loop } } catch (InterruptedException ignored) { break; } } } // reader loop (runs in readerExecutor) private void startReaderLoop() { readerExecutor.execute(() -> { try { while (running.get() && connected.get() && in != null) { int length; try { length = in.readInt(); // will throw SocketTimeoutException if set } catch (SocketTimeoutException ste) { // used to detect socket liveness; continue loop continue; } if (length <= 0 || length > 10 * 1024 * 1024) { // invalid length, break throw new RuntimeException("Invalid message length: " + length); } byte[] buf = new byte[length]; in.readFully(buf); String s = new String(buf); try { JSONObject json = new JSONObject(s); // update last seen time via heartbeat ack if needed if ("heartbeat".equals(json.optString("type"))) { // optionally respond or update time } else if ("auth_ok".equals(json.optString("type")) || "ack".equals(json.optString("type"))) { // ignore or process ack } else { if (listener != null) listener.onMessage(json); } } catch (Exception je) { Log.w(TAG, "Invalid JSON from server: " + s, je); } } } catch (Exception e) { if (running.get()) { safeCloseAndScheduleReconnect(e); } } }); } // schedule connect attempt with delay (ms) private void scheduleConnect(long delayMs) { connectExecutor.execute(() -> { try { if (delayMs > 0) Thread.sleep(delayMs); } catch (InterruptedException ignored) { } if (!running.get()) return; tryConnect(); }); } // connect logic private void tryConnect() { if (!running.get()) return; closeSocketQuiet(); // ensure closed try { Socket s = new Socket(); s.connect(new InetSocketAddress(serverIp, serverPort), connectTimeoutMs); s.setSoTimeout((int) Math.max(heartbeatTimeoutMs, 5_000)); socket = s; out = new DataOutputStream(socket.getOutputStream()); in = new DataInputStream(socket.getInputStream()); connected.set(true); reconnectAttempt.set(0); authSent.set(false); // start reader startReaderLoop(); // send auth immediately sendAuth(); // start heartbeat startHeartbeat(); if (listener != null) listener.onConnected(); Log.i(TAG, "Connected to " + serverIp + ":" + serverPort); } catch (Exception e) { Log.w(TAG, "Connect failed: " + e.getMessage()); scheduleReconnectWithBackoff(); } } // start heartbeat scheduler private void startHeartbeat() { cancelHeartbeat(); heartbeatFuture = scheduler.scheduleAtFixedRate(() -> { if (!running.get() || !connected.get()) return; try { JSONObject hb = new JSONObject(); hb.put("type", "heartbeat"); hb.put("time", System.currentTimeMillis()); sendQueue.offer(hb); } catch (Exception ignored) { } }, 0, heartbeatIntervalMs, TimeUnit.MILLISECONDS); } private void cancelHeartbeat() { if (heartbeatFuture != null && !heartbeatFuture.isCancelled()) { heartbeatFuture.cancel(true); heartbeatFuture = null; } } // close socket quietly and notify listener private void safeCloseAndScheduleReconnect(Exception cause) { closeSocketQuiet(); if (listener != null) listener.onDisconnected(cause); scheduleReconnectWithBackoff(); } private void scheduleReconnectWithBackoff() { int attempt = reconnectAttempt.incrementAndGet(); long delay = Math.min(maxReconnectDelayMsFromAttempt(attempt), maxReconnectDelayMs); Log.i(TAG, "Scheduling reconnect attempt " + attempt + " after " + delay + "ms"); scheduleConnect(delay); } // compute exponential backoff private long maxReconnectDelayMsFromAttempt(int attempt) { long d = baseReconnectDelayMs * (1L << Math.min(attempt, 30)); if (d < 0) d = maxReconnectDelayMs; return Math.min(d, maxReconnectDelayMs); } // close socket and streams private void closeSocketQuiet() { connected.set(false); cancelHeartbeat(); try { if (out != null) { out.close(); } } catch (Exception ignored) { } try { if (in != null) { in.close(); } } catch (Exception ignored) { } try { if (socket != null && !socket.isClosed()) { socket.close(); } } catch (Exception ignored) { } out = null; in = null; socket = null; } // helper: when manually call reconnect (immediately) public void reconnectNow() { scheduleConnect(0); } // helper to set immediate send of a JSON and wait (blocking) until it is queued (not until delivered) public boolean sendBlocking(JSONObject json, long timeoutMs) throws InterruptedException { return sendQueue.offer(json, timeoutMs, TimeUnit.MILLISECONDS); } }