Skip to content
Closed
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
6 changes: 6 additions & 0 deletions app/src/androidTest/java/com/proxyagent/app/e2e/E2EConfig.kt
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,12 @@ object E2EConfig {
quicTransportFactory = quicFactory,
quicDialTimeoutMs = 15000,
tcpWarmPoolSize = tcpWarmPool,
// The testserver, and every tunnel target it drives, lives on the
// emulator's host-loopback alias — an RFC1918 address, which the
// device-side target policy refuses by default and is right to: on a
// real handset that range is the owner's home network. The exception
// belongs here, in the harness that needs it, and nowhere else.
targetPolicyAllowCidrs = listOf(testserverHost),
)

/** Single-line per-tunnel summary for embedding in assertion failures.
Expand Down
106 changes: 97 additions & 9 deletions app/src/main/java/com/proxyagent/app/nativeagent/NativeProxyAgent.kt
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,21 @@ class NativeProxyAgent {
// QUIC factory chosen above also receives this — keep them
// in sync from the host.
val networkProfile: NetworkProfile = NetworkProfile.LOW_100,
// Destinations this device refuses to dial on a buyer's behalf. The
// address classes are fixed and unconditional — see [TargetPolicy] for
// why a consumer handset has no legitimate buyer destination inside its
// own network — so the only knobs are exceptions and extra refusals.
//
// targetPolicyAllowCidrs exempts ranges the embedder knows are its own:
// a test rig on 10.0.2.2, a lab network. Entries are "10.0.0.0/8" or a
// bare address, and a malformed one is dropped and logged rather than
// thrown, because failing towards refusing more is the safe direction
// inside somebody else's application.
val targetPolicyAllowCidrs: List<String> = emptyList(),
// Ports refused outright, on top of the address classes. Empty by
// default: which ports stop working is a commercial decision, not a
// technical one.
val targetPolicyDenyPorts: Set<Int> = emptySet(),
) {
fun hasDirectRegistrator(): Boolean =
!registratorHost.isNullOrBlank() && registratorPort > 0
Expand Down Expand Up @@ -996,6 +1011,18 @@ internal class Uplink(
private val writeLock = Any()
private val shuttingDown = AtomicBoolean(false)

// The device-side target policy. Built once per uplink and consulted on
// the resolved address of every target, immediately before connect — the
// last point at which anyone knows what this handset will actually dial.
private val targetPolicy =
TargetPolicy.of(cfg.targetPolicyAllowCidrs, cfg.targetPolicyDenyPorts).also {
val bad = TargetPolicy.invalidEntries(cfg.targetPolicyAllowCidrs)
if (bad.isNotEmpty()) {
agent.logWarn("target policy: ignoring malformed allow-list entries",
"entries" to bad.joinToString(","))
}
}

private val openExecutor = Executors.newCachedThreadPool(daemonFactory("uplink-open"))
private val bridgeExecutor = Executors.newCachedThreadPool(daemonFactory("uplink-bridge"))

Expand Down Expand Up @@ -1175,6 +1202,7 @@ internal class Uplink(
// matcher — we deliberately diverge for parity in the UI.
agent.logInfo("opening tunnel", "target" to target, "transport" to "quic")
opened = true
val targetAddr = resolveTargetOrRefuse(host, port)
val sock = Socket()
sock.tcpNoDelay = true
// Kernel keepalive so a silently-dead target (NAT drop, Wi-Fi↔
Expand All @@ -1184,7 +1212,7 @@ internal class Uplink(
try { sock.keepAlive = true } catch (_: Throwable) {}
try { sock.receiveBufferSize = cfg.networkProfile.tuning().tcp.socketBufferBytes } catch (_: Throwable) {}
try { sock.sendBufferSize = cfg.networkProfile.tuning().tcp.socketBufferBytes } catch (_: Throwable) {}
sock.connect(InetSocketAddress(dns.resolve(host), port),
sock.connect(InetSocketAddress(targetAddr, port),
NativeProxyAgent.TARGET_DIAL_TIMEOUT_MS.toInt())
// NOTE: this is a plain java.net.Socket, so we can only set bare
// SO_KEEPALIVE (above) — no safe read-only fd path to tune the
Expand All @@ -1200,14 +1228,17 @@ internal class Uplink(
bridgeStreams(streamIn, stream.output, sock)
} catch (t: Throwable) {
val msg = t.message ?: t.javaClass.simpleName
agent.logWarn("quic tunnel failed", "error" to msg)
val refused = isPolicyRefusal(t)
if (!refused) {
agent.logWarn("quic tunnel failed", "error" to msg)
}
try { stream.close() } catch (_: Throwable) {}
try { targetSock?.close() } catch (_: Throwable) {}
// Keep the host's activeTunnels parser balanced: it counts
// "tunnel closed", not "quic tunnel failed", so emit one here
// when we'd already logged "opening tunnel".
if (opened) agent.logInfo("tunnel closed", "transport" to "quic", "reason" to "failed")
noteOpenFailure(msg)
if (!refused) noteOpenFailure(msg)
} finally {
agent.decTunnels()
tunnelPermits.release()
Expand Down Expand Up @@ -1452,9 +1483,17 @@ internal class Uplink(
}
} catch (t: Throwable) {
val msg = t.message ?: t.javaClass.simpleName
agent.logWarn("tunnel open failed", "target" to target, "error" to msg)
reportOpenFail(token, msg)
noteOpenFailure(msg)
val refused = isPolicyRefusal(t)
if (!refused) {
agent.logWarn("tunnel open failed", "target" to target, "error" to msg)
}
// The server is told the open failed either way; it has no use for
// the resolved address, and a buyer able to read one back out of a
// failure would have a probe for the device's own network.
reportOpenFail(token, if (refused) "target refused by device policy" else msg)
// A refusal is a correct outcome, not a sick data plane — see
// noteOpenFailure.
if (!refused) noteOpenFailure(msg)
// Balance the host's log-parsed activeTunnels gauge: we already
// emitted "opening tunnel" above, and the parser only decrements
// on "tunnel closed" — without this a failed open would leave
Expand All @@ -1468,10 +1507,49 @@ internal class Uplink(
}
}

/** Resolves [host] and refuses the result if the target policy says so.
*
* Every target dial in this class goes through here, and it returns the
* same [InetAddress] the connect will use — resolving once and checking
* that exact answer is the point. Checking the hostname, or resolving a
* second time inside connect(), would leave the gap this exists to close.
*/
private fun resolveTargetOrRefuse(host: String, port: Int): InetAddress {
val addr = dns.resolve(host)
val reason = targetPolicy.refusalReason(addr, port) ?: return addr
agent.logWarn("target refused by policy on the device",
"target_host" to host, "resolved" to (addr.hostAddress ?: "?"),
"port" to port, "reason" to reason)
throw TargetRefusedException(reason)
}

/** True when [t] or anything it wraps is a policy refusal.
*
* The refusal is raised inside a Future and re-wrapped on the way out, so
* the identity has to be looked for down the cause chain rather than on
* the throwable in hand.
*/
private fun isPolicyRefusal(t: Throwable?): Boolean {
var cur = t
var hops = 0
while (cur != null && hops < 8) {
if (cur is TargetRefusedException) return true
cur = cur.cause
hops++
}
return false
}

/** Record a failed tunnel open and self-heal if the data plane looks
* wedged. An EMFILE is fatal on its own; otherwise we wait for a
* sustained back-to-back streak so a few transient dial timeouts on a
* flaky link don't needlessly bounce a working session. */
* flaky link don't needlessly bounce a working session.
*
* A policy refusal must never reach this. It is not a symptom of a sick
* data plane, and counting it would hand a buyer a way to take a handset
* off the network: aim forty opens at a blocked address and the streak
* trips selfHeal. Callers filter with [isPolicyRefusal] before calling.
*/
private fun noteOpenFailure(msg: String) {
if (shuttingDown.get()) return
val emfile = msg.contains(NativeProxyAgent.EMFILE_MARKER, ignoreCase = true)
Expand Down Expand Up @@ -1544,6 +1622,13 @@ internal class Uplink(
try {
// Dial target + take data conn in parallel — same as Go.
val futureTarget = bridgeExecutor.submit<SocketChannel> {
// Resolve and refuse FIRST, before a descriptor exists. This used
// to sit down in connect(), which meant a refusal threw with the
// channel already open and nothing to close it — and since a
// buyer chooses the target, that is a remotely driven descriptor
// leak on someone's phone, which is the harm the policy exists
// to prevent rather than cause.
val targetAddr = resolveTargetOrRefuse(host, port)
val ch = SocketChannel.open()
ch.socket().tcpNoDelay = true
// Kernel keepalive so a silently-dead target is reaped by
Expand All @@ -1555,7 +1640,7 @@ internal class Uplink(
try { ch.socket().receiveBufferSize = cfg.networkProfile.tuning().tcp.socketBufferBytes } catch (_: Throwable) {}
try { ch.socket().sendBufferSize = cfg.networkProfile.tuning().tcp.socketBufferBytes } catch (_: Throwable) {}
ch.socket().connect(
InetSocketAddress(dns.resolve(host), port),
InetSocketAddress(targetAddr, port),
NativeProxyAgent.TARGET_DIAL_TIMEOUT_MS.toInt(),
)
// Tune keepalive idle/interval/count on the raw fd. On the
Expand Down Expand Up @@ -1618,13 +1703,16 @@ internal class Uplink(
// initiated streams. But if the server emits an OPEN over the
// control channel anyway, fulfill it by opening a fresh stream.
val q = quic ?: throw IOException("no quic session")
// Before the stream is opened, so a refusal costs neither a stream nor a
// descriptor.
val targetAddr = resolveTargetOrRefuse(host, port)
val stream = q.openStream()
val targetSock = Socket().apply {
tcpNoDelay = true
try { keepAlive = true } catch (_: Throwable) {}
try { receiveBufferSize = cfg.networkProfile.tuning().tcp.socketBufferBytes } catch (_: Throwable) {}
try { sendBufferSize = cfg.networkProfile.tuning().tcp.socketBufferBytes } catch (_: Throwable) {}
connect(InetSocketAddress(dns.resolve(host), port),
connect(InetSocketAddress(targetAddr, port),
NativeProxyAgent.TARGET_DIAL_TIMEOUT_MS.toInt())
}
try {
Expand Down
Loading
Loading