mirror of
https://github.com/cirruslabs/softnet.git
synced 2026-10-01 04:21:54 +02:00
Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
18f5a3338e | ||
|
|
0c4327dd71 | ||
|
|
28bb29df4a | ||
|
|
d079057ecf | ||
|
|
237aec1df5 | ||
|
|
cdc2a508aa | ||
|
|
54b419fd18 |
Generated
+49
-141
@@ -246,12 +246,6 @@ version = "0.5.2"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "23b62fc65de8e4e7f52534fb52b0f3ed04746ae267519eef2a83941e8085068b"
|
||||
|
||||
[[package]]
|
||||
name = "atomic-waker"
|
||||
version = "1.1.2"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0"
|
||||
|
||||
[[package]]
|
||||
name = "autocfg"
|
||||
version = "1.4.0"
|
||||
@@ -374,9 +368,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "clap"
|
||||
version = "4.6.3"
|
||||
version = "4.6.5"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "0fb99565819980999fb7b4a1796046a5c949e6d4ff132cf5fadf5a641e20d776"
|
||||
checksum = "301b56658598e48f3648647ac6fc887be7e7108eddfa4e9b63fcf3ec58c0cadf"
|
||||
dependencies = [
|
||||
"clap_builder",
|
||||
"clap_derive",
|
||||
@@ -384,9 +378,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "clap_builder"
|
||||
version = "4.6.2"
|
||||
version = "4.6.5"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "f09628afdcc538b57f3c6341e9c8e9970f18e4a481690a64974d7023bd33548b"
|
||||
checksum = "94a65403d1a1bd28f7dc68eb8506e8874808ee5eecb59298de588e2e1407a078"
|
||||
dependencies = [
|
||||
"anstream",
|
||||
"anstyle",
|
||||
@@ -396,14 +390,14 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "clap_derive"
|
||||
version = "4.6.3"
|
||||
version = "4.6.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "32f2392eae7f16557a3d727ef3a12e57b2b2ca6f98566a5f4fb41ffe305df077"
|
||||
checksum = "d012d2b9d65aca7f18f4d9878a045bc17899bba951561ba5ec3c2ba1eed9a061"
|
||||
dependencies = [
|
||||
"heck",
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"syn 2.0.117",
|
||||
"syn 3.0.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -895,25 +889,6 @@ version = "0.31.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "32085ea23f3234fc7846555e85283ba4de91e21016dc0455a16286d87a292d64"
|
||||
|
||||
[[package]]
|
||||
name = "h2"
|
||||
version = "0.4.13"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "2f44da3a8150a6703ed5d34e164b875fd14c2cdab9af1252a9a1020bde2bdc54"
|
||||
dependencies = [
|
||||
"atomic-waker",
|
||||
"bytes",
|
||||
"fnv",
|
||||
"futures-core",
|
||||
"futures-sink",
|
||||
"http 1.1.0",
|
||||
"indexmap",
|
||||
"slab",
|
||||
"tokio",
|
||||
"tokio-util",
|
||||
"tracing",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "hash32"
|
||||
version = "0.3.1"
|
||||
@@ -1091,7 +1066,6 @@ dependencies = [
|
||||
"bytes",
|
||||
"futures-channel",
|
||||
"futures-util",
|
||||
"h2",
|
||||
"http 1.1.0",
|
||||
"http-body",
|
||||
"httparse",
|
||||
@@ -1102,22 +1076,6 @@ dependencies = [
|
||||
"want",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "hyper-rustls"
|
||||
version = "0.27.7"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "e3c93eb611681b207e1fe55d5a71ecf91572ec8a6705cdb6857f7d8d5242cf58"
|
||||
dependencies = [
|
||||
"http 1.1.0",
|
||||
"hyper",
|
||||
"hyper-util",
|
||||
"rustls",
|
||||
"rustls-pki-types",
|
||||
"tokio",
|
||||
"tokio-rustls",
|
||||
"tower-service",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "hyper-tls"
|
||||
version = "0.6.0"
|
||||
@@ -1335,9 +1293,9 @@ checksum = "aa2f047c0a98b2f299aa5d6d7088443570faae494e9ae1305e48be000c9e0eb1"
|
||||
|
||||
[[package]]
|
||||
name = "ipnet"
|
||||
version = "2.12.0"
|
||||
version = "2.12.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "d98f6fed1fde3f8c21bc40a1abb88dd75e67924f9cffc3ef95607bad8017f8e2"
|
||||
checksum = "6a756c3fac73139e83f14c2d742155dd2b78d3ee56597b419a0579b7bdd6dd78"
|
||||
|
||||
[[package]]
|
||||
name = "iri-string"
|
||||
@@ -1452,9 +1410,9 @@ checksum = "09edd9e8b54e49e587e4f6295a7d29c3ea94d469cb40ab8ca70b288248a81db2"
|
||||
|
||||
[[package]]
|
||||
name = "libc"
|
||||
version = "0.2.188"
|
||||
version = "0.2.189"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "22053b6a34f84abc97f9129e61334f40174659a1b9bd18c970b83db6a9a6348b"
|
||||
checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2"
|
||||
|
||||
[[package]]
|
||||
name = "linux-raw-sys"
|
||||
@@ -2066,12 +2024,10 @@ dependencies = [
|
||||
"futures-channel",
|
||||
"futures-core",
|
||||
"futures-util",
|
||||
"h2",
|
||||
"http 1.1.0",
|
||||
"http-body",
|
||||
"http-body-util",
|
||||
"hyper",
|
||||
"hyper-rustls",
|
||||
"hyper-tls",
|
||||
"hyper-util",
|
||||
"js-sys",
|
||||
@@ -2094,20 +2050,6 @@ dependencies = [
|
||||
"web-sys",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "ring"
|
||||
version = "0.17.14"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a4689e6c2294d81e88dc6261c768b63bc4fcdb852be6d1352498b114f61383b7"
|
||||
dependencies = [
|
||||
"cc",
|
||||
"cfg-if",
|
||||
"getrandom 0.2.15",
|
||||
"libc",
|
||||
"untrusted",
|
||||
"windows-sys 0.52.0",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustc-demangle"
|
||||
version = "0.1.24"
|
||||
@@ -2149,19 +2091,6 @@ dependencies = [
|
||||
"windows-sys 0.59.0",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustls"
|
||||
version = "0.23.37"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "758025cb5fccfd3bc2fd74708fd4682be41d99e5dff73c377c0646c6012c73a4"
|
||||
dependencies = [
|
||||
"once_cell",
|
||||
"rustls-pki-types",
|
||||
"rustls-webpki",
|
||||
"subtle",
|
||||
"zeroize",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustls-pemfile"
|
||||
version = "2.2.0"
|
||||
@@ -2180,17 +2109,6 @@ dependencies = [
|
||||
"zeroize",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustls-webpki"
|
||||
version = "0.103.13"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e"
|
||||
dependencies = [
|
||||
"ring",
|
||||
"rustls-pki-types",
|
||||
"untrusted",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustversion"
|
||||
version = "1.0.21"
|
||||
@@ -2258,9 +2176,9 @@ checksum = "61697e0a1c7e512e84a621326239844a24d8207b4669b41bc18b32ea5cbf988b"
|
||||
|
||||
[[package]]
|
||||
name = "sentry"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "d631477761f57c76148456e55e80e9a479ff3fa4c65b2b4a0c3acf1167fd4638"
|
||||
checksum = "63207365db50cb817402f0ec8097d6dad3d644ef3fffd6ba30c75b3077e514da"
|
||||
dependencies = [
|
||||
"cfg_aliases",
|
||||
"httpdate",
|
||||
@@ -2271,6 +2189,7 @@ dependencies = [
|
||||
"sentry-contexts",
|
||||
"sentry-core",
|
||||
"sentry-debug-images",
|
||||
"sentry-log",
|
||||
"sentry-panic",
|
||||
"sentry-tracing",
|
||||
"tokio",
|
||||
@@ -2279,9 +2198,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "sentry-actix"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "119ede4e37790ec04e8a14073c8414f0a4b2648402856c0563495a9a5cf54d57"
|
||||
checksum = "4a0eb1ed18478fb48db007aeaa672a3de176055357feb768eee0c3471802d0ce"
|
||||
dependencies = [
|
||||
"actix-http",
|
||||
"actix-web",
|
||||
@@ -2292,9 +2211,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "sentry-anyhow"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "d8f73466a403f4da78c7e576049d6044b31d2b16dfcc26b51705f318ed74b804"
|
||||
checksum = "6265521f1b724f709bc07747ecb01c5289b974cb70b348026df01780280fc136"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"sentry-backtrace",
|
||||
@@ -2303,9 +2222,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "sentry-backtrace"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "448b0981fbde6cdc9eb087ba3dc01035a253fef9a0d9e79aaf198fc26acb2e64"
|
||||
checksum = "e3bcc2497c2327998146207b7600599ef7592233920f4d4e1d41ddeeba0e4210"
|
||||
dependencies = [
|
||||
"backtrace",
|
||||
"regex",
|
||||
@@ -2314,9 +2233,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "sentry-contexts"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9e5e909d02170ba6d1dc5ebd05bef7999280f5d960e3eaee5d1a18d88d41b334"
|
||||
checksum = "1cb04cba225b38f59d08f9e3c29ab7259d9cc94fbb2c46d613a51cb0a4a7ee05"
|
||||
dependencies = [
|
||||
"hostname",
|
||||
"libc",
|
||||
@@ -2328,9 +2247,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "sentry-core"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ecf0b1a4a4e9ec88395b52e6fa4868b95564c1fa96b26a0606f5f6288a0f7149"
|
||||
checksum = "48759bc392fb5e3b36b4efcb485f51074c0ba8479711cd5c8cf79e86f9f75d41"
|
||||
dependencies = [
|
||||
"rand 0.9.4",
|
||||
"sentry-types",
|
||||
@@ -2341,19 +2260,30 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "sentry-debug-images"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "f81d7749c57fc78ed52134889e8121da874065bc40da788d5798cce0f8c19f15"
|
||||
checksum = "0a5bd325059b70b21ca6e42b0f2409821272dddc1cf37923de37269af55389b5"
|
||||
dependencies = [
|
||||
"findshlibs",
|
||||
"sentry-core",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "sentry-panic"
|
||||
version = "0.48.5"
|
||||
name = "sentry-log"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "0b7d93d6ecb55d2251c5fc084c55c03a2bc68904918d76117e602c999e92f00f"
|
||||
checksum = "f9553a2f0f75bc8550137ebc753e8d2161af63ff780ca59983935f8df8f01655"
|
||||
dependencies = [
|
||||
"bitflags 2.9.4",
|
||||
"log",
|
||||
"sentry-core",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "sentry-panic"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a685d22f672b1562ebeff3f2affceb5960595648cc936dae73da3e1d6a55fbd1"
|
||||
dependencies = [
|
||||
"sentry-backtrace",
|
||||
"sentry-core",
|
||||
@@ -2361,9 +2291,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "sentry-tracing"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "58d26379236c4ef97eaf081bca3c7bf0ef6f06b3aa881eca2ee8e8dbc923e5b8"
|
||||
checksum = "8aa7d8e0db4cccddbac01ddabd5344ad03b7c2f954b3ff850cecb1d9cf5d4758"
|
||||
dependencies = [
|
||||
"bitflags 2.9.4",
|
||||
"sentry-backtrace",
|
||||
@@ -2374,9 +2304,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "sentry-types"
|
||||
version = "0.48.5"
|
||||
version = "0.49.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "2239099d47e76857b0825a182ecdb90159add7aa1097b60247c6e7ca6bf49c6a"
|
||||
checksum = "fab146d15a30ab4897a95fd15d6a038b8d7fbf2661c5f8b13738ce8df7a095b7"
|
||||
dependencies = [
|
||||
"debugid",
|
||||
"hex",
|
||||
@@ -2446,9 +2376,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "serial_test"
|
||||
version = "3.5.0"
|
||||
version = "4.0.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "699f4197115b8a7e7ff19c9a315a4bd6fffec26cc4626ef45ecaea389e081c6d"
|
||||
checksum = "a6df5ed973ad8d834e09f824f9e9f449af6b9a3745f78dec7cc752770bd3bf11"
|
||||
dependencies = [
|
||||
"futures-executor",
|
||||
"futures-util",
|
||||
@@ -2460,13 +2390,13 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "serial_test_derive"
|
||||
version = "3.5.0"
|
||||
version = "4.0.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "94e153fc76e1c6a068703d6d29c508a0b15c061c4b7e43da59cc097bc342673c"
|
||||
checksum = "a22144e767da4ddd8416dbf383700542ffd8a5dc493dfecedfe1fe3ad03c98ae"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"syn 2.0.117",
|
||||
"syn 3.0.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -2611,12 +2541,6 @@ version = "0.11.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f"
|
||||
|
||||
[[package]]
|
||||
name = "subtle"
|
||||
version = "2.6.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292"
|
||||
|
||||
[[package]]
|
||||
name = "syn"
|
||||
version = "1.0.109"
|
||||
@@ -2827,16 +2751,6 @@ dependencies = [
|
||||
"tokio",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "tokio-rustls"
|
||||
version = "0.26.4"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61"
|
||||
dependencies = [
|
||||
"rustls",
|
||||
"tokio",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "tokio-util"
|
||||
version = "0.7.15"
|
||||
@@ -2998,12 +2912,6 @@ version = "0.2.6"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ebc1c04c71510c7f702b52b7c350734c9ff1295c464a03335b00bb84fc54f853"
|
||||
|
||||
[[package]]
|
||||
name = "untrusted"
|
||||
version = "0.9.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1"
|
||||
|
||||
[[package]]
|
||||
name = "ureq"
|
||||
version = "3.0.12"
|
||||
|
||||
+1
-1
@@ -32,7 +32,7 @@ prefix-trie = "0"
|
||||
ipnet = "2"
|
||||
oslog = "0.2.0"
|
||||
log = "0.4.29"
|
||||
serial_test = "3"
|
||||
serial_test = "4"
|
||||
coarsetime = "0.1.37"
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
|
||||
@@ -30,6 +30,33 @@ And assumes that:
|
||||
|
||||
...otherwise it's possible for two VMs to receive an identical IP-address from the macOS built-in DHCP-server (even in the presence of Softnet's packet filtering) and thus bypass the protections offered by Softnet.
|
||||
|
||||
### Stateful flow authorization
|
||||
|
||||
Stateful `in`/`out` rules use a bounded authorization cache that records the
|
||||
direction and exact transport tuple of policy-approved flows, allowing matching
|
||||
return traffic without treating it as a new flow and thus requiring a separate
|
||||
policy entry. It is not a complete TCP connection tracker: endpoint transport
|
||||
stacks remain responsible for validating sequence numbers, receive windows,
|
||||
resets, and application-level traffic.
|
||||
|
||||
This cache deliberately favors security, bounded resource use, and a simple
|
||||
implementation over availability. Only policy-authorized initiator traffic
|
||||
renews an entry; return traffic does not. Softnet does not maintain fairness
|
||||
quotas, eviction heuristics, or complete TCP lifecycle state.
|
||||
|
||||
A packet admitted by a stateful rule is denied when Softnet cannot represent its
|
||||
flow, including when the cache is full. High flow churn, long idle connections,
|
||||
or ambiguous retransmissions may therefore interrupt networking and require the
|
||||
affected VM to reconnect.
|
||||
|
||||
For TCP, a bare TCP SYN on an existing tuple is deliberately returned to policy
|
||||
because Softnet cannot distinguish a retransmission from tuple reuse without
|
||||
tracking TCP sequence state. If authorized, it may replace the tuple's previous
|
||||
cache lifetime; this can reduce availability but cannot grant traffic that
|
||||
policy did not permit.
|
||||
|
||||
For ICMP, stateful flow authorization supports only echo requests and replies.
|
||||
|
||||
## Installing
|
||||
|
||||
For proper functioning, Softnet binary requires two things:
|
||||
@@ -57,6 +84,6 @@ The supported methods are `softnet.policy.get` and `softnet.policy.set`. A compl
|
||||
|
||||
Every request must include a non-null string (at most 256 bytes) or non-negative integer `id`; notifications are rejected so policy changes always have an acknowledgment. Policy updates are atomic: all rules are parsed and a new prefix map is built before the active policy changes. Longest-prefix matching and block precedence for identical rules are preserved. Rules are normalized and deduplicated. A policy update may contain at most 4096 combined allow/block rules, and a request frame may not exceed 1 MiB.
|
||||
|
||||
When the effective normalized policy changes, Softnet clears existing conntrack state so the new policy applies to established flows immediately. This may interrupt active connections. Repeating an equivalent normalized policy is a no-op and preserves conntrack state.
|
||||
When the normalized allow or block policy changes, Softnet clears the flow table so the new policy applies to established flows immediately. This may interrupt active connections. Repeating the same normalized policy is a no-op and preserves the flow table.
|
||||
|
||||
Use `block=["0.0.0.0/0"]` with specific allow rules for a default-deny policy. Closing the control socket leaves the last accepted policy active.
|
||||
Use `block=["0.0.0.0/0"]` with specific allow rules for a default-deny egress policy. Closing the control socket leaves the last accepted policy active.
|
||||
|
||||
+82
-2
@@ -1,18 +1,20 @@
|
||||
use dhcproto::Decodable;
|
||||
use dhcproto::v4::{DhcpOption, MessageType, OptionCode};
|
||||
use dhcproto::v4::{DhcpOption, HType, Message, MessageType, Opcode, OptionCode};
|
||||
use smoltcp::wire::Ipv4Address;
|
||||
use std::collections::HashSet;
|
||||
use std::time::Duration;
|
||||
|
||||
#[derive(Default)]
|
||||
pub struct DhcpSnooper {
|
||||
vm_mac_address: [u8; 6],
|
||||
vm_lease: Option<Lease>,
|
||||
uncertainty_duration: Duration,
|
||||
}
|
||||
|
||||
impl DhcpSnooper {
|
||||
pub fn new(uncertainty_duration: Duration) -> Self {
|
||||
pub fn new(uncertainty_duration: Duration, vm_mac_address: [u8; 6]) -> Self {
|
||||
DhcpSnooper {
|
||||
vm_mac_address,
|
||||
uncertainty_duration,
|
||||
..Default::default()
|
||||
}
|
||||
@@ -26,6 +28,14 @@ impl DhcpSnooper {
|
||||
Err(_) => return,
|
||||
};
|
||||
|
||||
// Decoded DHCP replies may be broadcast[1], so additionally validate the BOOTP client
|
||||
// hardware address to avoid acting on another VM's lease transition
|
||||
//
|
||||
// [1]: https://datatracker.ietf.org/doc/html/rfc2131#section-4.1
|
||||
if !message_matches_bootp_client(&message, Opcode::BootReply, self.vm_mac_address) {
|
||||
return;
|
||||
}
|
||||
|
||||
match message.opts().msg_type() {
|
||||
Some(MessageType::Ack) => {
|
||||
let lease_time = match message.opts().get(OptionCode::AddressLeaseTime) {
|
||||
@@ -63,6 +73,11 @@ impl DhcpSnooper {
|
||||
&self.vm_lease
|
||||
}
|
||||
|
||||
pub(crate) fn address_and_dns_ips(&self) -> Option<(Ipv4Address, HashSet<Ipv4Address>)> {
|
||||
let lease = self.vm_lease.as_ref().filter(|lease| lease.valid())?;
|
||||
Some((lease.address(), lease.dns_ips.clone()))
|
||||
}
|
||||
|
||||
pub fn valid_dns_target(&self, addr: &Ipv4Address) -> bool {
|
||||
if let Some(lease) = &self.vm_lease {
|
||||
return lease.dns_ips.contains(addr);
|
||||
@@ -100,3 +115,68 @@ impl Lease {
|
||||
self.address == address && self.valid()
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn message_matches_bootp_client(
|
||||
message: &Message,
|
||||
opcode: Opcode,
|
||||
mac: [u8; 6],
|
||||
) -> bool {
|
||||
message.opcode() == opcode
|
||||
&& message.htype() == HType::Eth
|
||||
&& message.hlen() == mac.len() as u8
|
||||
&& message.chaddr() == mac
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{DhcpSnooper, Lease};
|
||||
use dhcproto::v4::{DhcpOption, Message, MessageType, Opcode};
|
||||
use dhcproto::{Encodable, Encoder};
|
||||
use smoltcp::wire::Ipv4Address;
|
||||
use std::collections::HashSet;
|
||||
use std::time::Duration;
|
||||
|
||||
const VM_MAC: [u8; 6] = [0x02, 0x00, 0x00, 0x00, 0x00, 0x01];
|
||||
const OTHER_MAC: [u8; 6] = [0x02, 0x00, 0x00, 0x00, 0x00, 0x02];
|
||||
const OLD_ADDRESS: Ipv4Address = Ipv4Address::new(192, 168, 64, 2);
|
||||
|
||||
#[test]
|
||||
fn processes_replies_only_for_matching_client() {
|
||||
// Start with an active lease
|
||||
let mut snooper = DhcpSnooper::new(Duration::ZERO, VM_MAC);
|
||||
snooper.set_lease(Some(Lease::new(
|
||||
OLD_ADDRESS,
|
||||
Duration::from_secs(600),
|
||||
HashSet::new(),
|
||||
)));
|
||||
|
||||
// Ignore a NAK for another client
|
||||
let mut message = Message::new(
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
&OTHER_MAC,
|
||||
);
|
||||
message.set_opcode(Opcode::BootReply);
|
||||
message
|
||||
.opts_mut()
|
||||
.insert(DhcpOption::MessageType(MessageType::Nak));
|
||||
|
||||
let mut encoded = Vec::new();
|
||||
message.encode(&mut Encoder::new(&mut encoded)).unwrap();
|
||||
|
||||
snooper.register_dhcp_reply(&encoded);
|
||||
|
||||
assert_eq!(snooper.lease().as_ref().unwrap().address(), OLD_ADDRESS);
|
||||
|
||||
// Process a NAK for the matching client
|
||||
message.set_chaddr(&VM_MAC);
|
||||
encoded.clear();
|
||||
message.encode(&mut Encoder::new(&mut encoded)).unwrap();
|
||||
|
||||
snooper.register_dhcp_reply(&encoded);
|
||||
|
||||
assert!(snooper.lease().is_none());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,386 +0,0 @@
|
||||
#[path = "conntrack_tcp.rs"]
|
||||
mod tcp;
|
||||
#[path = "conntrack_udp.rs"]
|
||||
mod udp;
|
||||
|
||||
use coarsetime::{Duration, Instant};
|
||||
use smoltcp::wire::{IpProtocol, Ipv4Address, Ipv4Packet};
|
||||
use std::collections::HashMap;
|
||||
|
||||
const MAX_FLOWS: usize = 4096;
|
||||
const MAX_VM_INITIATED_FLOWS: usize = 1024;
|
||||
const MAX_EMBRYONIC_FLOWS_PER_SOURCE: usize = 256;
|
||||
const SWEEP_INTERVAL: Duration = Duration::from_secs(1);
|
||||
|
||||
/// Tracks permission to use an exact host/VM flow tuple.
|
||||
///
|
||||
/// This deliberately does not replace either endpoint's transport stack:
|
||||
/// TCP sequence/window validation remains the host's and VM's responsibility.
|
||||
/// Its security boundary is flow initiation: only an authorized first packet
|
||||
/// can create an entry, and VM traffic to the host must match the reverse.
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct Conntrack {
|
||||
flows: HashMap<FlowKey, Flow>,
|
||||
next_sweep: Instant,
|
||||
}
|
||||
|
||||
pub(crate) enum ConntrackResult {
|
||||
Allowed,
|
||||
New(PendingFlow),
|
||||
Denied,
|
||||
}
|
||||
|
||||
pub(crate) struct PendingFlow {
|
||||
key: FlowKey,
|
||||
flow: Flow,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
||||
enum FlowKey {
|
||||
Tcp {
|
||||
host_addr: Ipv4Address,
|
||||
host_port: u16,
|
||||
vm_addr: Ipv4Address,
|
||||
vm_port: u16,
|
||||
},
|
||||
Udp {
|
||||
host_addr: Ipv4Address,
|
||||
host_port: u16,
|
||||
vm_addr: Ipv4Address,
|
||||
vm_port: u16,
|
||||
},
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
enum Initiator {
|
||||
Host,
|
||||
Vm,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
struct Flow {
|
||||
initiator: Initiator,
|
||||
state: FlowState,
|
||||
last_seen: Instant,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
enum FlowState {
|
||||
Tcp(tcp::State),
|
||||
Udp(udp::State),
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
enum Direction {
|
||||
FromHost,
|
||||
FromVm,
|
||||
}
|
||||
|
||||
impl Conntrack {
|
||||
pub(crate) fn new() -> Self {
|
||||
let now = Instant::recent();
|
||||
Self {
|
||||
flows: HashMap::new(),
|
||||
next_sweep: now + SWEEP_INTERVAL,
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn inspect_from_host(&mut self, packet: &Ipv4Packet<&[u8]>) -> ConntrackResult {
|
||||
self.inspect(packet, Direction::FromHost, Instant::recent())
|
||||
}
|
||||
|
||||
pub(crate) fn inspect_from_vm(&mut self, packet: &Ipv4Packet<&[u8]>) -> ConntrackResult {
|
||||
self.inspect(packet, Direction::FromVm, Instant::recent())
|
||||
}
|
||||
|
||||
pub(crate) fn commit(&mut self, mut pending: PendingFlow) -> bool {
|
||||
let now = Instant::recent();
|
||||
pending.flow.last_seen = now;
|
||||
self.insert(pending.key, pending.flow, now)
|
||||
}
|
||||
|
||||
pub(crate) fn tick(&mut self) {
|
||||
let now = Instant::recent();
|
||||
if now >= self.next_sweep {
|
||||
self.expire(now);
|
||||
self.next_sweep = now + SWEEP_INTERVAL;
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn clear(&mut self) {
|
||||
self.flows.clear();
|
||||
}
|
||||
|
||||
fn inspect(
|
||||
&mut self,
|
||||
packet: &Ipv4Packet<&[u8]>,
|
||||
direction: Direction,
|
||||
now: Instant,
|
||||
) -> ConntrackResult {
|
||||
// Later fragments do not contain the transport header needed to bind them
|
||||
// to a permitted flow. Fail closed instead of admitting them by IP alone.
|
||||
if packet.more_frags() || packet.frag_offset() != 0 {
|
||||
return ConntrackResult::Denied;
|
||||
}
|
||||
|
||||
match packet.next_header() {
|
||||
IpProtocol::Tcp => self.inspect_tcp(packet, direction, now),
|
||||
IpProtocol::Udp => self.inspect_udp(packet, direction, now),
|
||||
_ => ConntrackResult::Denied,
|
||||
}
|
||||
}
|
||||
|
||||
fn insert(&mut self, key: FlowKey, flow: Flow, now: Instant) -> bool {
|
||||
self.expire(now);
|
||||
if self.flows.len() >= MAX_FLOWS {
|
||||
return false;
|
||||
}
|
||||
if flow.initiator == Initiator::Vm
|
||||
&& self
|
||||
.flows
|
||||
.values()
|
||||
.filter(|flow| flow.initiator == Initiator::Vm)
|
||||
.count()
|
||||
>= MAX_VM_INITIATED_FLOWS
|
||||
{
|
||||
return false;
|
||||
}
|
||||
if let Some(initiator_addr) = embryonic_source(key, flow)
|
||||
&& self
|
||||
.flows
|
||||
.iter()
|
||||
.filter(|(key, flow)| embryonic_source(**key, **flow) == Some(initiator_addr))
|
||||
.count()
|
||||
>= MAX_EMBRYONIC_FLOWS_PER_SOURCE
|
||||
{
|
||||
return false;
|
||||
}
|
||||
self.flows.insert(key, flow);
|
||||
true
|
||||
}
|
||||
|
||||
fn expire(&mut self, now: Instant) {
|
||||
self.flows
|
||||
.retain(|_, flow| now.duration_since(flow.last_seen) < flow.timeout());
|
||||
}
|
||||
}
|
||||
|
||||
impl Flow {
|
||||
fn timeout(&self) -> Duration {
|
||||
match self.state {
|
||||
FlowState::Tcp(state) => state.timeout(),
|
||||
FlowState::Udp(state) => state.timeout(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn embryonic_source(key: FlowKey, flow: Flow) -> Option<Ipv4Address> {
|
||||
match (key, flow.initiator, flow.state) {
|
||||
(
|
||||
FlowKey::Tcp { host_addr, .. },
|
||||
Initiator::Host,
|
||||
FlowState::Tcp(tcp::State::SynSent | tcp::State::SynReceived),
|
||||
) => Some(host_addr),
|
||||
(
|
||||
FlowKey::Tcp { vm_addr, .. },
|
||||
Initiator::Vm,
|
||||
FlowState::Tcp(tcp::State::SynSent | tcp::State::SynReceived),
|
||||
) => Some(vm_addr),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
fn oriented_transport_key(
|
||||
packet: &Ipv4Packet<&[u8]>,
|
||||
direction: Direction,
|
||||
make_key: impl FnOnce(Ipv4Address, u16, Ipv4Address, u16) -> FlowKey,
|
||||
src_port: u16,
|
||||
dst_port: u16,
|
||||
) -> FlowKey {
|
||||
match direction {
|
||||
Direction::FromHost => make_key(packet.src_addr(), src_port, packet.dst_addr(), dst_port),
|
||||
Direction::FromVm => make_key(packet.dst_addr(), dst_port, packet.src_addr(), src_port),
|
||||
}
|
||||
}
|
||||
|
||||
fn initiator(direction: Direction) -> Initiator {
|
||||
match direction {
|
||||
Direction::FromHost => Initiator::Host,
|
||||
Direction::FromVm => Initiator::Vm,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::test_support::{HOST, TcpFlags, VM, inspect_from_host, inspect_from_vm, tcp_packet};
|
||||
use super::{Conntrack, ConntrackResult, MAX_EMBRYONIC_FLOWS_PER_SOURCE};
|
||||
use smoltcp::wire::{Ipv4Address, Ipv4Packet};
|
||||
|
||||
#[test]
|
||||
fn fragments_fail_closed() {
|
||||
let mut packet = tcp_packet(HOST, 49152, VM, 22, TcpFlags::SYN);
|
||||
Ipv4Packet::new_unchecked(packet.as_mut_slice()).set_more_frags(true);
|
||||
|
||||
let mut tracker = Conntrack::new();
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &packet),
|
||||
ConntrackResult::Denied
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn clear_removes_tracked_flows() {
|
||||
let packet = tcp_packet(HOST, 49152, VM, 22, TcpFlags::SYN);
|
||||
let mut tracker = Conntrack::new();
|
||||
|
||||
let ConntrackResult::New(pending) = inspect_from_host(&mut tracker, &packet) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
assert_eq!(tracker.flows.len(), 1);
|
||||
|
||||
tracker.clear();
|
||||
assert!(tracker.flows.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn embryonic_tcp_limit_is_per_source() {
|
||||
let mut tracker = Conntrack::new();
|
||||
|
||||
for index in 0..MAX_EMBRYONIC_FLOWS_PER_SOURCE {
|
||||
let syn = tcp_packet(HOST, 10000 + index as u16, VM, 22, TcpFlags::SYN);
|
||||
let ConntrackResult::New(pending) = inspect_from_host(&mut tracker, &syn) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
}
|
||||
|
||||
let over_limit = tcp_packet(HOST, 20000, VM, 22, TcpFlags::SYN);
|
||||
let ConntrackResult::New(pending) = inspect_from_host(&mut tracker, &over_limit) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(!tracker.commit(pending));
|
||||
|
||||
let other_host = Ipv4Address::new(192, 168, 64, 3);
|
||||
let other_source = tcp_packet(other_host, 20000, VM, 22, TcpFlags::SYN);
|
||||
let ConntrackResult::New(pending) = inspect_from_host(&mut tracker, &other_source) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
|
||||
let syn_ack = tcp_packet(VM, 22, HOST, 10000, TcpFlags::SYN_ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &syn_ack),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
let ack = tcp_packet(HOST, 10000, VM, 22, TcpFlags::ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &ack),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
|
||||
let ConntrackResult::New(pending) = inspect_from_host(&mut tracker, &over_limit) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod test_support {
|
||||
use super::{Conntrack, ConntrackResult};
|
||||
use smoltcp::wire::{IpProtocol, Ipv4Address, Ipv4Packet, TcpPacket, UdpPacket};
|
||||
|
||||
pub(super) const HOST: Ipv4Address = Ipv4Address::new(192, 168, 64, 1);
|
||||
pub(super) const VM: Ipv4Address = Ipv4Address::new(192, 168, 64, 2);
|
||||
|
||||
pub(super) fn inspect_from_host(tracker: &mut Conntrack, bytes: &[u8]) -> ConntrackResult {
|
||||
let packet = Ipv4Packet::new_checked(bytes).unwrap();
|
||||
tracker.inspect_from_host(&packet)
|
||||
}
|
||||
|
||||
pub(super) fn inspect_from_vm(tracker: &mut Conntrack, bytes: &[u8]) -> ConntrackResult {
|
||||
let packet = Ipv4Packet::new_checked(bytes).unwrap();
|
||||
tracker.inspect_from_vm(&packet)
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
pub(super) struct TcpFlags {
|
||||
syn: bool,
|
||||
ack: bool,
|
||||
rst: bool,
|
||||
}
|
||||
|
||||
impl TcpFlags {
|
||||
pub(super) const SYN: Self = Self {
|
||||
syn: true,
|
||||
ack: false,
|
||||
rst: false,
|
||||
};
|
||||
pub(super) const SYN_ACK: Self = Self {
|
||||
syn: true,
|
||||
ack: true,
|
||||
rst: false,
|
||||
};
|
||||
pub(super) const ACK: Self = Self {
|
||||
syn: false,
|
||||
ack: true,
|
||||
rst: false,
|
||||
};
|
||||
pub(super) const RST: Self = Self {
|
||||
syn: false,
|
||||
ack: false,
|
||||
rst: true,
|
||||
};
|
||||
}
|
||||
|
||||
pub(super) fn tcp_packet(
|
||||
src_addr: Ipv4Address,
|
||||
src_port: u16,
|
||||
dst_addr: Ipv4Address,
|
||||
dst_port: u16,
|
||||
flags: TcpFlags,
|
||||
) -> Vec<u8> {
|
||||
let mut bytes = ipv4_packet(src_addr, dst_addr, IpProtocol::Tcp, 20);
|
||||
let mut tcp = TcpPacket::new_unchecked(&mut bytes[20..]);
|
||||
tcp.set_src_port(src_port);
|
||||
tcp.set_dst_port(dst_port);
|
||||
tcp.set_header_len(20);
|
||||
tcp.set_syn(flags.syn);
|
||||
tcp.set_ack(flags.ack);
|
||||
tcp.set_rst(flags.rst);
|
||||
bytes
|
||||
}
|
||||
|
||||
pub(super) fn udp_packet(
|
||||
src_addr: Ipv4Address,
|
||||
src_port: u16,
|
||||
dst_addr: Ipv4Address,
|
||||
dst_port: u16,
|
||||
) -> Vec<u8> {
|
||||
let mut bytes = ipv4_packet(src_addr, dst_addr, IpProtocol::Udp, 8);
|
||||
let mut udp = UdpPacket::new_unchecked(&mut bytes[20..]);
|
||||
udp.set_src_port(src_port);
|
||||
udp.set_dst_port(dst_port);
|
||||
udp.set_len(8);
|
||||
bytes
|
||||
}
|
||||
|
||||
fn ipv4_packet(
|
||||
src_addr: Ipv4Address,
|
||||
dst_addr: Ipv4Address,
|
||||
protocol: IpProtocol,
|
||||
payload_len: usize,
|
||||
) -> Vec<u8> {
|
||||
let mut bytes = vec![0; 20 + payload_len];
|
||||
let total_len = bytes.len() as u16;
|
||||
let mut ipv4 = Ipv4Packet::new_unchecked(bytes.as_mut_slice());
|
||||
ipv4.set_version(4);
|
||||
ipv4.set_header_len(20);
|
||||
ipv4.set_total_len(total_len);
|
||||
ipv4.set_next_header(protocol);
|
||||
ipv4.set_src_addr(src_addr);
|
||||
ipv4.set_dst_addr(dst_addr);
|
||||
bytes
|
||||
}
|
||||
}
|
||||
@@ -1,257 +0,0 @@
|
||||
use super::{
|
||||
Conntrack, ConntrackResult, Direction, Flow, FlowKey, FlowState, PendingFlow, initiator,
|
||||
oriented_transport_key,
|
||||
};
|
||||
use coarsetime::{Duration, Instant};
|
||||
use smoltcp::wire::{Ipv4Packet, TcpPacket};
|
||||
|
||||
const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(60);
|
||||
const ESTABLISHED_TIMEOUT: Duration = Duration::from_secs(5 * 24 * 60 * 60);
|
||||
const CLOSING_TIMEOUT: Duration = Duration::from_secs(2 * 60);
|
||||
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub(super) enum State {
|
||||
SynSent,
|
||||
SynReceived,
|
||||
Established,
|
||||
Closing { host_fin: bool, vm_fin: bool },
|
||||
}
|
||||
|
||||
impl State {
|
||||
pub(super) fn timeout(self) -> Duration {
|
||||
match self {
|
||||
Self::SynSent | Self::SynReceived => HANDSHAKE_TIMEOUT,
|
||||
Self::Established => ESTABLISHED_TIMEOUT,
|
||||
Self::Closing { .. } => CLOSING_TIMEOUT,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Conntrack {
|
||||
pub(super) fn inspect_tcp(
|
||||
&mut self,
|
||||
packet: &Ipv4Packet<&[u8]>,
|
||||
direction: Direction,
|
||||
now: Instant,
|
||||
) -> ConntrackResult {
|
||||
let Ok(tcp) = TcpPacket::new_checked(packet.payload()) else {
|
||||
return ConntrackResult::Denied;
|
||||
};
|
||||
if tcp.src_port() == 0 || tcp.dst_port() == 0 {
|
||||
return ConntrackResult::Denied;
|
||||
}
|
||||
|
||||
let key = oriented_transport_key(
|
||||
packet,
|
||||
direction,
|
||||
|host_addr, host_port, vm_addr, vm_port| FlowKey::Tcp {
|
||||
host_addr,
|
||||
host_port,
|
||||
vm_addr,
|
||||
vm_port,
|
||||
},
|
||||
tcp.src_port(),
|
||||
tcp.dst_port(),
|
||||
);
|
||||
|
||||
if let Some(flow) = self.flows.get_mut(&key) {
|
||||
let from_initiator = matches!(
|
||||
(flow.initiator, direction),
|
||||
(super::Initiator::Host, Direction::FromHost)
|
||||
| (super::Initiator::Vm, Direction::FromVm)
|
||||
);
|
||||
let FlowState::Tcp(state) = &mut flow.state else {
|
||||
return ConntrackResult::Denied;
|
||||
};
|
||||
|
||||
if tcp.rst() {
|
||||
self.flows.remove(&key);
|
||||
return ConntrackResult::Allowed;
|
||||
}
|
||||
|
||||
let allowed = match *state {
|
||||
State::SynSent if from_initiator => is_initial_syn(&tcp),
|
||||
State::SynSent => {
|
||||
if tcp.syn() && tcp.ack() && !tcp.fin() {
|
||||
*state = State::SynReceived;
|
||||
true
|
||||
} else {
|
||||
false
|
||||
}
|
||||
}
|
||||
State::SynReceived if from_initiator => {
|
||||
if is_initial_syn(&tcp) {
|
||||
true
|
||||
} else if tcp.ack() && !tcp.syn() {
|
||||
*state = if tcp.fin() {
|
||||
State::Closing {
|
||||
host_fin: matches!(direction, Direction::FromHost),
|
||||
vm_fin: matches!(direction, Direction::FromVm),
|
||||
}
|
||||
} else {
|
||||
State::Established
|
||||
};
|
||||
true
|
||||
} else {
|
||||
false
|
||||
}
|
||||
}
|
||||
State::SynReceived => tcp.syn() && tcp.ack() && !tcp.fin(),
|
||||
State::Established | State::Closing { .. }
|
||||
if is_initial_syn(&tcp) && !from_initiator =>
|
||||
{
|
||||
false
|
||||
}
|
||||
State::Established | State::Closing { .. } if is_initial_syn(&tcp) => {
|
||||
*state = State::SynSent;
|
||||
true
|
||||
}
|
||||
State::Established if tcp.fin() => {
|
||||
*state = State::Closing {
|
||||
host_fin: matches!(direction, Direction::FromHost),
|
||||
vm_fin: matches!(direction, Direction::FromVm),
|
||||
};
|
||||
true
|
||||
}
|
||||
State::Closing { .. } if tcp.fin() => {
|
||||
if let State::Closing { host_fin, vm_fin } = state {
|
||||
match direction {
|
||||
Direction::FromHost => *host_fin = true,
|
||||
Direction::FromVm => *vm_fin = true,
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
_ => true,
|
||||
};
|
||||
|
||||
if allowed {
|
||||
flow.last_seen = now;
|
||||
}
|
||||
return if allowed {
|
||||
ConntrackResult::Allowed
|
||||
} else {
|
||||
ConntrackResult::Denied
|
||||
};
|
||||
}
|
||||
|
||||
if !is_initial_syn(&tcp) {
|
||||
return ConntrackResult::Denied;
|
||||
}
|
||||
|
||||
ConntrackResult::New(PendingFlow {
|
||||
key,
|
||||
flow: Flow {
|
||||
initiator: initiator(direction),
|
||||
state: FlowState::Tcp(State::SynSent),
|
||||
last_seen: now,
|
||||
},
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
fn is_initial_syn(tcp: &TcpPacket<&[u8]>) -> bool {
|
||||
tcp.syn() && !tcp.ack() && !tcp.fin() && !tcp.rst()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::super::test_support::{
|
||||
HOST, TcpFlags, VM, inspect_from_host, inspect_from_vm, tcp_packet,
|
||||
};
|
||||
use super::super::{Conntrack, ConntrackResult};
|
||||
|
||||
#[test]
|
||||
fn tcp_flows_are_oriented() {
|
||||
let mut tracker = Conntrack::new();
|
||||
|
||||
let vm_syn = tcp_packet(VM, 22, HOST, 49152, TcpFlags::SYN);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &vm_syn),
|
||||
ConntrackResult::New(_)
|
||||
));
|
||||
|
||||
let host_syn = tcp_packet(HOST, 49152, VM, 22, TcpFlags::SYN);
|
||||
let ConntrackResult::New(pending) = inspect_from_host(&mut tracker, &host_syn) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
|
||||
let premature_vm_ack = tcp_packet(VM, 22, HOST, 49152, TcpFlags::ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &premature_vm_ack),
|
||||
ConntrackResult::Denied
|
||||
));
|
||||
|
||||
let wrong_vm_reply = tcp_packet(VM, 22, HOST, 49153, TcpFlags::SYN_ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &wrong_vm_reply),
|
||||
ConntrackResult::Denied
|
||||
));
|
||||
|
||||
let vm_syn_ack = tcp_packet(VM, 22, HOST, 49152, TcpFlags::SYN_ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &vm_syn_ack),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
|
||||
let premature_vm_ack = tcp_packet(VM, 22, HOST, 49152, TcpFlags::ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &premature_vm_ack),
|
||||
ConntrackResult::Denied
|
||||
));
|
||||
|
||||
let host_ack = tcp_packet(HOST, 49152, VM, 22, TcpFlags::ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &host_ack),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &premature_vm_ack),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn vm_can_initiate_tcp() {
|
||||
let mut tracker = Conntrack::new();
|
||||
let vm_syn = tcp_packet(VM, 49152, HOST, 22, TcpFlags::SYN);
|
||||
let host_syn_ack = tcp_packet(HOST, 22, VM, 49152, TcpFlags::SYN_ACK);
|
||||
let vm_ack = tcp_packet(VM, 49152, HOST, 22, TcpFlags::ACK);
|
||||
|
||||
let ConntrackResult::New(pending) = inspect_from_vm(&mut tracker, &vm_syn) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &host_syn_ack),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &vm_ack),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tcp_rst_removes_permission() {
|
||||
let mut tracker = Conntrack::new();
|
||||
let host_syn = tcp_packet(HOST, 49152, VM, 22, TcpFlags::SYN);
|
||||
let vm_rst = tcp_packet(VM, 22, HOST, 49152, TcpFlags::RST);
|
||||
let vm_ack = tcp_packet(VM, 22, HOST, 49152, TcpFlags::ACK);
|
||||
|
||||
let ConntrackResult::New(pending) = inspect_from_host(&mut tracker, &host_syn) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &vm_rst),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &vm_ack),
|
||||
ConntrackResult::Denied
|
||||
));
|
||||
}
|
||||
}
|
||||
@@ -1,131 +0,0 @@
|
||||
use super::{
|
||||
Conntrack, ConntrackResult, Direction, Flow, FlowKey, FlowState, Initiator, PendingFlow,
|
||||
initiator, oriented_transport_key,
|
||||
};
|
||||
use coarsetime::{Duration, Instant};
|
||||
use smoltcp::wire::{Ipv4Packet, UdpPacket};
|
||||
|
||||
const UNREPLIED_TIMEOUT: Duration = Duration::from_secs(30);
|
||||
const REPLIED_TIMEOUT: Duration = Duration::from_secs(3 * 60);
|
||||
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub(super) struct State {
|
||||
replied: bool,
|
||||
}
|
||||
|
||||
impl State {
|
||||
pub(super) fn timeout(self) -> Duration {
|
||||
if self.replied {
|
||||
REPLIED_TIMEOUT
|
||||
} else {
|
||||
UNREPLIED_TIMEOUT
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Conntrack {
|
||||
pub(super) fn inspect_udp(
|
||||
&mut self,
|
||||
packet: &Ipv4Packet<&[u8]>,
|
||||
direction: Direction,
|
||||
now: Instant,
|
||||
) -> ConntrackResult {
|
||||
let Ok(udp) = UdpPacket::new_checked(packet.payload()) else {
|
||||
return ConntrackResult::Denied;
|
||||
};
|
||||
if udp.src_port() == 0 || udp.dst_port() == 0 {
|
||||
return ConntrackResult::Denied;
|
||||
}
|
||||
|
||||
let key = oriented_transport_key(
|
||||
packet,
|
||||
direction,
|
||||
|host_addr, host_port, vm_addr, vm_port| FlowKey::Udp {
|
||||
host_addr,
|
||||
host_port,
|
||||
vm_addr,
|
||||
vm_port,
|
||||
},
|
||||
udp.src_port(),
|
||||
udp.dst_port(),
|
||||
);
|
||||
|
||||
if let Some(flow) = self.flows.get_mut(&key) {
|
||||
let FlowState::Udp(state) = &mut flow.state else {
|
||||
return ConntrackResult::Denied;
|
||||
};
|
||||
|
||||
let is_reply = matches!(
|
||||
(flow.initiator, direction),
|
||||
(Initiator::Host, Direction::FromVm) | (Initiator::Vm, Direction::FromHost)
|
||||
);
|
||||
state.replied |= is_reply;
|
||||
flow.last_seen = now;
|
||||
return ConntrackResult::Allowed;
|
||||
}
|
||||
|
||||
ConntrackResult::New(PendingFlow {
|
||||
key,
|
||||
flow: Flow {
|
||||
initiator: initiator(direction),
|
||||
state: FlowState::Udp(State { replied: false }),
|
||||
last_seen: now,
|
||||
},
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::super::test_support::{HOST, VM, inspect_from_host, inspect_from_vm, udp_packet};
|
||||
use super::super::{Conntrack, ConntrackResult};
|
||||
|
||||
#[test]
|
||||
fn udp_reply_requires_an_exact_host_request() {
|
||||
let mut tracker = Conntrack::new();
|
||||
|
||||
let vm_datagram = udp_packet(VM, 5353, HOST, 50000);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &vm_datagram),
|
||||
ConntrackResult::New(_)
|
||||
));
|
||||
|
||||
let host_datagram = udp_packet(HOST, 50000, VM, 5353);
|
||||
let ConntrackResult::New(pending) = inspect_from_host(&mut tracker, &host_datagram) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &vm_datagram),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
|
||||
let wrong_vm_datagram = udp_packet(VM, 5353, HOST, 50001);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &wrong_vm_datagram),
|
||||
ConntrackResult::New(_)
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn explicitly_allowed_vm_udp_gets_only_its_reply() {
|
||||
let mut tracker = Conntrack::new();
|
||||
let vm_dns = udp_packet(VM, 53000, HOST, 53);
|
||||
let host_dns = udp_packet(HOST, 53, VM, 53000);
|
||||
|
||||
let ConntrackResult::New(pending) = inspect_from_vm(&mut tracker, &vm_dns) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &host_dns),
|
||||
ConntrackResult::Allowed
|
||||
));
|
||||
|
||||
let unsolicited_host_udp = udp_packet(HOST, 53, VM, 53001);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &unsolicited_host_udp),
|
||||
ConntrackResult::New(_)
|
||||
));
|
||||
}
|
||||
}
|
||||
+55
-53
@@ -1,4 +1,4 @@
|
||||
use super::{Rule, Rules, build_rules, rule_count};
|
||||
use super::{Rule, Rules};
|
||||
use anyhow::{Context, Result, bail};
|
||||
use jsonrpsee_types::{
|
||||
ErrorObjectOwned, Id, Request, Response, ResponsePayload,
|
||||
@@ -60,7 +60,7 @@ impl Policy {
|
||||
|
||||
let allow = parse_rules(allow)?;
|
||||
let block = parse_rules(block)?;
|
||||
let rules = build_rules(self.gateway_ip, &allow, &block);
|
||||
let rules = Rules::new(self.gateway_ip, &allow, &block);
|
||||
|
||||
Ok(PolicyUpdate {
|
||||
rules,
|
||||
@@ -71,14 +71,13 @@ impl Policy {
|
||||
|
||||
fn apply(&mut self, update: PolicyUpdate) -> Option<Rules> {
|
||||
// Build and validate everything before updating any active state. The packet filter
|
||||
// observes either the old PrefixMap or the complete new one.
|
||||
if self.allow == update.allow && self.block == update.block {
|
||||
return None;
|
||||
}
|
||||
// observes either the old rule set or the complete new one.
|
||||
let changed = self.allow != update.allow || self.block != update.block;
|
||||
|
||||
self.allow = update.allow;
|
||||
self.block = update.block;
|
||||
Some(update.rules)
|
||||
|
||||
changed.then_some(update.rules)
|
||||
}
|
||||
|
||||
fn result(&self, rule_count: usize) -> Value {
|
||||
@@ -88,7 +87,7 @@ impl Policy {
|
||||
|
||||
impl PolicyUpdate {
|
||||
fn result(&self) -> Value {
|
||||
policy_result(&self.allow, &self.block, rule_count(&self.rules))
|
||||
policy_result(&self.allow, &self.block, self.rules.len())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -116,7 +115,7 @@ fn parse_rules(rules: Vec<String>) -> std::result::Result<Vec<Rule>, ErrorObject
|
||||
Ok(normalize_rules(parsed))
|
||||
}
|
||||
|
||||
fn normalize_rules(mut rules: Vec<Rule>) -> Vec<Rule> {
|
||||
pub(super) fn normalize_rules(mut rules: Vec<Rule>) -> Vec<Rule> {
|
||||
rules.iter_mut().for_each(|rule| *rule = rule.normalized());
|
||||
rules.sort_by_key(ToString::to_string);
|
||||
rules.dedup();
|
||||
@@ -276,8 +275,7 @@ impl Control {
|
||||
continue;
|
||||
}
|
||||
|
||||
let (response, update) =
|
||||
handle_request(&self.policy, rule_count(rules), &line[..newline]);
|
||||
let (response, update) = handle_request(&self.policy, rules, &line[..newline]);
|
||||
self.enqueue(response)?;
|
||||
|
||||
if let Some(update) = update
|
||||
@@ -353,11 +351,7 @@ struct SetParams {
|
||||
block: Vec<String>,
|
||||
}
|
||||
|
||||
fn handle_request(
|
||||
policy: &Policy,
|
||||
rule_count: usize,
|
||||
line: &[u8],
|
||||
) -> (Value, Option<PolicyUpdate>) {
|
||||
fn handle_request(policy: &Policy, rules: &Rules, line: &[u8]) -> (Value, Option<PolicyUpdate>) {
|
||||
let value = match serde_json::from_slice::<Value>(line) {
|
||||
Ok(value) => value,
|
||||
Err(_) => {
|
||||
@@ -394,7 +388,7 @@ fn handle_request(
|
||||
"softnet.policy.get does not accept parameters",
|
||||
))
|
||||
} else {
|
||||
Ok(policy.result(rule_count))
|
||||
Ok(policy.result(rules.len()))
|
||||
}
|
||||
}
|
||||
"softnet.policy.set" => {
|
||||
@@ -540,11 +534,9 @@ fn validate_control_fd(control_fd: RawFd) -> Result<()> {
|
||||
mod tests {
|
||||
use super::{
|
||||
Control, INVALID_PARAMS, INVALID_REQUEST, MAX_PENDING_RESPONSE_BYTES, MAX_REQUEST_BYTES,
|
||||
MAX_RULES, METHOD_NOT_FOUND, PARSE_ERROR, Policy, build_rules, handle_request, rule_count,
|
||||
MAX_RULES, METHOD_NOT_FOUND, PARSE_ERROR, Policy, handle_request,
|
||||
};
|
||||
use crate::proxy::{Action, Rule, Rules};
|
||||
use ipnet::Ipv4Net;
|
||||
use prefix_trie::PrefixMap;
|
||||
use crate::proxy::{Direction, PolicyDecision, Rule, Rules};
|
||||
use serde_json::{Value, json};
|
||||
use smoltcp::wire::Ipv4Address;
|
||||
use std::fs::File;
|
||||
@@ -552,7 +544,6 @@ mod tests {
|
||||
use std::net::{Shutdown, TcpListener};
|
||||
use std::os::fd::{AsRawFd, RawFd};
|
||||
use std::os::unix::net::{UnixDatagram, UnixStream};
|
||||
use std::str::FromStr;
|
||||
use std::time::Duration;
|
||||
|
||||
struct TestPolicy {
|
||||
@@ -562,7 +553,7 @@ mod tests {
|
||||
|
||||
impl TestPolicy {
|
||||
fn result(&self) -> Value {
|
||||
self.state.result(rule_count(&self.rules))
|
||||
self.state.result(self.rules.len())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -584,7 +575,7 @@ mod tests {
|
||||
let block = rules(block);
|
||||
|
||||
TestPolicy {
|
||||
rules: build_rules(gateway_ip, &allow, &block),
|
||||
rules: Rules::new(gateway_ip, &allow, &block),
|
||||
state: Policy::new(gateway_ip, allow, block),
|
||||
}
|
||||
}
|
||||
@@ -603,7 +594,7 @@ mod tests {
|
||||
}
|
||||
|
||||
fn raw_request(policy: &mut TestPolicy, line: &[u8]) -> Value {
|
||||
let (response, update) = handle_request(&policy.state, rule_count(&policy.rules), line);
|
||||
let (response, update) = handle_request(&policy.state, &policy.rules, line);
|
||||
if let Some(update) = update
|
||||
&& let Some(rules) = policy.state.apply(update)
|
||||
{
|
||||
@@ -652,14 +643,16 @@ mod tests {
|
||||
assert_eq!(response["result"]["ruleCount"], 2);
|
||||
|
||||
assert_eq!(
|
||||
policy.rules.get(&Ipv4Net::from_str("10.0.0.0/8").unwrap()),
|
||||
Some(&vec![(Action::Block, "10.0.0.0/8".parse().unwrap())])
|
||||
policy
|
||||
.rules
|
||||
.policy_decision(Ipv4Address::new(10, 0, 0, 1), Direction::Out),
|
||||
Some(PolicyDecision::Block)
|
||||
);
|
||||
assert_eq!(
|
||||
policy
|
||||
.rules
|
||||
.get(&Ipv4Net::from_str("192.168.64.1/32").unwrap()),
|
||||
Some(&vec![(Action::Block, "192.168.64.1/32".parse().unwrap())])
|
||||
.policy_decision(Ipv4Address::new(192, 168, 64, 1), Direction::Out),
|
||||
Some(PolicyDecision::Block)
|
||||
);
|
||||
}
|
||||
|
||||
@@ -712,12 +705,14 @@ mod tests {
|
||||
);
|
||||
assert_eq!(response["result"]["block"], json!(["in 10.0.0.0/8"]));
|
||||
assert_eq!(response["result"]["ruleCount"], 2);
|
||||
let target = Ipv4Address::new(10, 1, 2, 3);
|
||||
assert_eq!(
|
||||
policy.rules.get(&Ipv4Net::from_str("10.0.0.0/8").unwrap()),
|
||||
Some(&vec![
|
||||
(Action::Block, "in 10.0.0.0/8".parse::<Rule>().unwrap()),
|
||||
(Action::Allow, "out 10.0.0.0/8".parse::<Rule>().unwrap()),
|
||||
])
|
||||
policy.rules.policy_decision(target, Direction::In),
|
||||
Some(PolicyDecision::Block)
|
||||
);
|
||||
assert_eq!(
|
||||
policy.rules.policy_decision(target, Direction::Out),
|
||||
Some(PolicyDecision::AllowStateful)
|
||||
);
|
||||
}
|
||||
|
||||
@@ -879,7 +874,7 @@ mod tests {
|
||||
.set_read_timeout(Some(Duration::from_secs(1)))
|
||||
.unwrap();
|
||||
let mut control = control(server.as_raw_fd()).unwrap();
|
||||
let mut rules = PrefixMap::new();
|
||||
let mut rules = Rules::default();
|
||||
|
||||
client
|
||||
.write_all(
|
||||
@@ -905,29 +900,36 @@ mod tests {
|
||||
assert_eq!(lines[1]["result"]["allow"], json!(["@host"]));
|
||||
assert_eq!(control.policy.allow, vec!["@host".parse().unwrap()]);
|
||||
|
||||
let before = control.policy.result(rule_count(&rules));
|
||||
let before = control.policy.result(rules.len());
|
||||
drop(client);
|
||||
assert!(!control.service(&mut rules).unwrap());
|
||||
assert_eq!(control.policy.result(rule_count(&rules)), before);
|
||||
assert_eq!(control.policy.result(rules.len()), before);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn repeated_normalized_policy_is_not_reported_as_changed() {
|
||||
fn policy_change_detection_uses_normalized_policy() {
|
||||
let (mut client, server) = UnixStream::pair().unwrap();
|
||||
let mut control = control(server.as_raw_fd()).unwrap();
|
||||
let mut rules = PrefixMap::new();
|
||||
let mut rules = Rules::default();
|
||||
|
||||
client
|
||||
.write_all(b"{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"softnet.policy.set\",\"params\":{\"allow\":[\"in 10.1.2.3/8\"],\"block\":[]}}\n")
|
||||
.write_all(b"{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"softnet.policy.set\",\"params\":{\"allow\":[\"out 10.1.2.3/8\"],\"block\":[\"in 10.0.0.0/8\",\"out 10.0.0.0/8\"]}}\n")
|
||||
.unwrap();
|
||||
assert!(control.service(&mut rules).unwrap());
|
||||
assert!(control.policy_changed());
|
||||
|
||||
client
|
||||
.write_all(b"{\"jsonrpc\":\"2.0\",\"id\":2,\"method\":\"softnet.policy.set\",\"params\":{\"allow\":[\"in 10.0.0.0/8\"],\"block\":[]}}\n")
|
||||
.write_all(b"{\"jsonrpc\":\"2.0\",\"id\":2,\"method\":\"softnet.policy.set\",\"params\":{\"allow\":[\"out 10.0.0.0/8\"],\"block\":[\"in 10.0.0.0/8\",\"out 10.0.0.0/8\"]}}\n")
|
||||
.unwrap();
|
||||
assert!(control.service(&mut rules).unwrap());
|
||||
assert!(!control.policy_changed());
|
||||
|
||||
client
|
||||
.write_all(b"{\"jsonrpc\":\"2.0\",\"id\":3,\"method\":\"softnet.policy.set\",\"params\":{\"allow\":[],\"block\":[\"in 10.0.0.0/8\",\"out 10.0.0.0/8\"]}}\n")
|
||||
.unwrap();
|
||||
assert!(control.service(&mut rules).unwrap());
|
||||
assert!(control.policy_changed());
|
||||
assert!(control.policy.allow.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -937,7 +939,7 @@ mod tests {
|
||||
.set_read_timeout(Some(Duration::from_secs(1)))
|
||||
.unwrap();
|
||||
let mut control = control(server.as_raw_fd()).unwrap();
|
||||
let mut rules = PrefixMap::new();
|
||||
let mut rules = Rules::default();
|
||||
|
||||
client
|
||||
.write_all(b"{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"softnet.policy.set\",\"params\":{\"allow\":[\"10.0.0.0/8\"],\"block\":[]}}\n")
|
||||
@@ -961,7 +963,7 @@ mod tests {
|
||||
.set_read_timeout(Some(Duration::from_secs(1)))
|
||||
.unwrap();
|
||||
let mut control = control(server.as_raw_fd()).unwrap();
|
||||
let mut rules = PrefixMap::new();
|
||||
let mut rules = Rules::default();
|
||||
|
||||
control.input = vec![b'x'; MAX_REQUEST_BYTES + 1];
|
||||
control.process_input(&mut rules).unwrap();
|
||||
@@ -992,8 +994,8 @@ mod tests {
|
||||
.set_read_timeout(Some(Duration::from_secs(1)))
|
||||
.unwrap();
|
||||
let mut control = control(server.as_raw_fd()).unwrap();
|
||||
let mut rules = PrefixMap::new();
|
||||
let before = control.policy.result(rule_count(&rules));
|
||||
let mut rules = Rules::default();
|
||||
let before = control.policy.result(rules.len());
|
||||
|
||||
client
|
||||
.write_all(
|
||||
@@ -1001,7 +1003,7 @@ mod tests {
|
||||
)
|
||||
.unwrap();
|
||||
assert!(control.service(&mut rules).unwrap());
|
||||
assert_eq!(control.policy.result(rule_count(&rules)), before);
|
||||
assert_eq!(control.policy.result(rules.len()), before);
|
||||
assert!(control.output.is_empty());
|
||||
|
||||
client.write_all(b"\"block\":[]}}\n").unwrap();
|
||||
@@ -1045,7 +1047,7 @@ mod tests {
|
||||
fn pipelined_policy_responses_stop_before_queue_overflow() {
|
||||
let (_client, server) = UnixStream::pair().unwrap();
|
||||
let mut control = control(server.as_raw_fd()).unwrap();
|
||||
let mut rules = PrefixMap::new();
|
||||
let mut rules = Rules::default();
|
||||
let allow = (0..MAX_RULES)
|
||||
.map(|index| format!("10.{}.{}.0/24", index / 256, index % 256))
|
||||
.collect::<Vec<_>>();
|
||||
@@ -1077,8 +1079,8 @@ mod tests {
|
||||
fn response_queue_overflow_does_not_apply_a_policy_update() {
|
||||
let (_client, server) = UnixStream::pair().unwrap();
|
||||
let mut control = control(server.as_raw_fd()).unwrap();
|
||||
let mut rules = PrefixMap::new();
|
||||
let before = control.policy.result(rule_count(&rules));
|
||||
let mut rules = Rules::default();
|
||||
let before = control.policy.result(rules.len());
|
||||
|
||||
control.output = vec![b'x'; MAX_PENDING_RESPONSE_BYTES];
|
||||
control.input = b"{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"softnet.policy.set\",\"params\":{\"allow\":[\"10.0.0.0/8\"],\"block\":[]}}\n".to_vec();
|
||||
@@ -1089,15 +1091,15 @@ mod tests {
|
||||
.to_string()
|
||||
.contains("control response queue exceeded")
|
||||
);
|
||||
assert_eq!(control.policy.result(rule_count(&rules)), before);
|
||||
assert_eq!(control.policy.result(rules.len()), before);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn response_backpressure_stops_consuming_policy_updates() {
|
||||
let (mut client, server) = UnixStream::pair().unwrap();
|
||||
let mut control = control(server.as_raw_fd()).unwrap();
|
||||
let mut rules = PrefixMap::new();
|
||||
let before = control.policy.result(rule_count(&rules));
|
||||
let mut rules = Rules::default();
|
||||
let before = control.policy.result(rules.len());
|
||||
|
||||
control.output = vec![b'x'; MAX_REQUEST_BYTES];
|
||||
assert!(control.flush().unwrap());
|
||||
@@ -1108,7 +1110,7 @@ mod tests {
|
||||
.unwrap();
|
||||
|
||||
assert!(control.service(&mut rules).unwrap());
|
||||
assert_eq!(control.policy.result(rule_count(&rules)), before);
|
||||
assert_eq!(control.policy.result(rules.len()), before);
|
||||
assert!(control.input.is_empty());
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,173 @@
|
||||
use super::{FlowDirection, FlowKey, FlowMatch, FlowTable};
|
||||
use coarsetime::{Duration, Instant};
|
||||
use smoltcp::wire::{Icmpv4Message, Icmpv4Packet, Ipv4Address, Ipv4Packet};
|
||||
|
||||
const ECHO_TIMEOUT: Duration = Duration::from_secs(30);
|
||||
|
||||
impl FlowTable {
|
||||
pub(super) fn inspect_icmp(
|
||||
&mut self,
|
||||
ipv4_pkt: &Ipv4Packet<&[u8]>,
|
||||
direction: FlowDirection,
|
||||
now: Instant,
|
||||
) -> FlowMatch {
|
||||
let Ok(icmp) = Icmpv4Packet::new_checked(ipv4_pkt.payload()) else {
|
||||
return FlowMatch::Denied;
|
||||
};
|
||||
if !icmp.verify_checksum() {
|
||||
return FlowMatch::Denied;
|
||||
}
|
||||
|
||||
match (icmp.msg_type(), icmp.msg_code()) {
|
||||
(Icmpv4Message::EchoRequest, 0) => {
|
||||
self.inspect_echo(ipv4_pkt, direction, icmp.echo_ident(), true, now)
|
||||
}
|
||||
(Icmpv4Message::EchoReply, 0) => {
|
||||
self.inspect_echo(ipv4_pkt, direction, icmp.echo_ident(), false, now)
|
||||
}
|
||||
// Other ICMP messages follow normal source policy: treating errors that quote
|
||||
// tracked tuples as RELATED requires validation beyond this exact-tuple table
|
||||
//
|
||||
// Potential degradation of PMTU discovery and traceroute is an accepted tradeoff.
|
||||
_ => FlowMatch::Untracked,
|
||||
}
|
||||
}
|
||||
|
||||
fn inspect_echo(
|
||||
&mut self,
|
||||
ipv4_pkt: &Ipv4Packet<&[u8]>,
|
||||
direction: FlowDirection,
|
||||
ident: u16,
|
||||
is_request: bool,
|
||||
now: Instant,
|
||||
) -> FlowMatch {
|
||||
// We need to preserve the request direction so opposite-direction echo flows remain distinct
|
||||
let initiating_direction = if is_request {
|
||||
direction
|
||||
} else {
|
||||
direction.opposite()
|
||||
};
|
||||
|
||||
let key = FlowKey::icmp_echo(
|
||||
ipv4_pkt.src_addr(),
|
||||
ipv4_pkt.dst_addr(),
|
||||
direction,
|
||||
initiating_direction,
|
||||
ident,
|
||||
);
|
||||
|
||||
if let Some(matched) = self.match_existing_flow(key, direction, now, ECHO_TIMEOUT) {
|
||||
return matched;
|
||||
}
|
||||
|
||||
if is_request {
|
||||
FlowMatch::candidate(key, initiating_direction, now, ECHO_TIMEOUT)
|
||||
} else {
|
||||
// Only replies matching an admitted request are tracked
|
||||
FlowMatch::Untracked
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl FlowKey {
|
||||
fn icmp_echo(
|
||||
src_addr: Ipv4Address,
|
||||
dst_addr: Ipv4Address,
|
||||
direction: FlowDirection,
|
||||
initiating_direction: FlowDirection,
|
||||
ident: u16,
|
||||
) -> Self {
|
||||
let (host_addr, vm_addr) = direction.host_vm_pair(src_addr, dst_addr);
|
||||
Self::IcmpEcho {
|
||||
host_addr,
|
||||
vm_addr,
|
||||
ident,
|
||||
initiating_direction,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::super::test_support::{HOST, VM, inspect_from_host, inspect_from_vm, ipv4_packet};
|
||||
use super::super::{FlowMatch, FlowTable};
|
||||
use smoltcp::wire::{Icmpv4Message, Icmpv4Packet, IpProtocol, Ipv4Address};
|
||||
|
||||
#[test]
|
||||
fn echo_request_and_reply_are_tracked_in_both_directions() {
|
||||
let mut tracker = FlowTable::new();
|
||||
let vm_request = echo_packet(VM, HOST, Icmpv4Message::EchoRequest, 7, 1);
|
||||
|
||||
let FlowMatch::Candidate(pending) = inspect_from_vm(&mut tracker, &vm_request) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
|
||||
let host_reply = echo_packet(HOST, VM, Icmpv4Message::EchoReply, 7, 1);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &host_reply),
|
||||
FlowMatch::Allowed
|
||||
));
|
||||
|
||||
let wrong_ident = echo_packet(HOST, VM, Icmpv4Message::EchoReply, 8, 1);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &wrong_ident),
|
||||
FlowMatch::Untracked
|
||||
));
|
||||
|
||||
// The opposite direction is a distinct flow even with the same identifier.
|
||||
let host_request = echo_packet(HOST, VM, Icmpv4Message::EchoRequest, 7, 1);
|
||||
let FlowMatch::Candidate(pending) = inspect_from_host(&mut tracker, &host_request) else {
|
||||
panic!("expected a new flow");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
|
||||
let vm_reply = echo_packet(VM, HOST, Icmpv4Message::EchoReply, 7, 1);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &vm_reply),
|
||||
FlowMatch::Allowed
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invalid_echo_is_denied_and_unsolicited_echo_is_untracked() {
|
||||
let mut tracker = FlowTable::new();
|
||||
let unsolicited = echo_packet(HOST, VM, Icmpv4Message::EchoReply, 7, 1);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &unsolicited),
|
||||
FlowMatch::Untracked
|
||||
));
|
||||
|
||||
let mut wrong_code = echo_packet(HOST, VM, Icmpv4Message::EchoRequest, 7, 1);
|
||||
wrong_code[21] = 1;
|
||||
Icmpv4Packet::new_unchecked(&mut wrong_code[20..]).fill_checksum();
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &wrong_code),
|
||||
FlowMatch::Untracked
|
||||
));
|
||||
|
||||
let mut bad_checksum = echo_packet(HOST, VM, Icmpv4Message::EchoRequest, 7, 1);
|
||||
bad_checksum[27] ^= 1;
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &bad_checksum),
|
||||
FlowMatch::Denied
|
||||
));
|
||||
}
|
||||
|
||||
fn echo_packet(
|
||||
src_addr: Ipv4Address,
|
||||
dst_addr: Ipv4Address,
|
||||
message: Icmpv4Message,
|
||||
ident: u16,
|
||||
sequence: u16,
|
||||
) -> Vec<u8> {
|
||||
let mut bytes = ipv4_packet(src_addr, dst_addr, IpProtocol::Icmp, 8);
|
||||
let mut icmp = Icmpv4Packet::new_unchecked(&mut bytes[20..]);
|
||||
icmp.set_msg_type(message);
|
||||
icmp.set_msg_code(0);
|
||||
icmp.set_echo_ident(ident);
|
||||
icmp.set_echo_seq_no(sequence);
|
||||
icmp.fill_checksum();
|
||||
bytes
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,394 @@
|
||||
mod icmp;
|
||||
mod tcp;
|
||||
mod udp;
|
||||
|
||||
use coarsetime::{Duration, Instant};
|
||||
use smoltcp::wire::{IpProtocol, Ipv4Address, Ipv4Packet};
|
||||
use std::collections::HashMap;
|
||||
|
||||
const MAX_FLOWS: usize = 32_768;
|
||||
const SWEEP_INTERVAL: Duration = Duration::from_secs(1);
|
||||
|
||||
/// A bounded cache of exact-tuple permissions for return traffic.
|
||||
///
|
||||
/// The table establishes only the direction in which a flow was authorized.
|
||||
/// This is deliberately not a TCP state machine. Endpoint transport stacks remain
|
||||
/// responsible for TCP handshakes, teardown, sequence numbers, and receive windows,
|
||||
/// as well as ICMP echo sequences.
|
||||
///
|
||||
/// Capacity is intentionally enforced with one fail-closed limit per VM.
|
||||
/// Exhaustion may deny further networking for that VM; this availability
|
||||
/// tradeoff is accepted to avoid fairness quotas, admission scans, and eviction.
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct FlowTable {
|
||||
flows: HashMap<FlowKey, Flow>,
|
||||
next_sweep: Instant,
|
||||
}
|
||||
|
||||
/// A proposed flow-table entry awaiting policy authorization.
|
||||
pub(crate) struct PendingFlow {
|
||||
key: FlowKey,
|
||||
flow: Flow,
|
||||
}
|
||||
|
||||
/// Metadata stored for an admitted flow.
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
struct Flow {
|
||||
initiating_direction: FlowDirection,
|
||||
expires_at: Instant,
|
||||
}
|
||||
|
||||
/// The result of classifying a packet against the flow table.
|
||||
pub(crate) enum FlowMatch {
|
||||
/// The packet is exact-tuple return traffic for an admitted flow.
|
||||
Allowed,
|
||||
|
||||
/// Policy must authorize this initiator-side packet before committing it.
|
||||
Candidate(PendingFlow),
|
||||
|
||||
/// The packet is malformed or violates a transport invariant that we deliberately enforce.
|
||||
Denied,
|
||||
|
||||
/// The flow table has no stateful interpretation for this packet.
|
||||
Untracked,
|
||||
}
|
||||
|
||||
impl FlowMatch {
|
||||
fn candidate(
|
||||
key: FlowKey,
|
||||
initiating_direction: FlowDirection,
|
||||
now: Instant,
|
||||
timeout: Duration,
|
||||
) -> Self {
|
||||
Self::Candidate(PendingFlow {
|
||||
key,
|
||||
flow: Flow {
|
||||
initiating_direction,
|
||||
expires_at: now + timeout,
|
||||
},
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// The canonical identity used to index a flow-table entry.
|
||||
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
||||
enum FlowKey {
|
||||
IcmpEcho {
|
||||
host_addr: Ipv4Address,
|
||||
vm_addr: Ipv4Address,
|
||||
ident: u16,
|
||||
initiating_direction: FlowDirection,
|
||||
},
|
||||
Tcp {
|
||||
host_addr: Ipv4Address,
|
||||
host_port: u16,
|
||||
vm_addr: Ipv4Address,
|
||||
vm_port: u16,
|
||||
},
|
||||
Udp {
|
||||
host_addr: Ipv4Address,
|
||||
host_port: u16,
|
||||
vm_addr: Ipv4Address,
|
||||
vm_port: u16,
|
||||
},
|
||||
}
|
||||
|
||||
/// A packet direction across the host–VM boundary.
|
||||
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
|
||||
pub(super) enum FlowDirection {
|
||||
FromHost,
|
||||
FromVm,
|
||||
}
|
||||
|
||||
impl FlowTable {
|
||||
pub(crate) fn new() -> Self {
|
||||
Self {
|
||||
flows: HashMap::new(),
|
||||
next_sweep: Instant::recent() + SWEEP_INTERVAL,
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn inspect(
|
||||
&mut self,
|
||||
ipv4_pkt: &Ipv4Packet<&[u8]>,
|
||||
direction: FlowDirection,
|
||||
) -> FlowMatch {
|
||||
self.inspect_at(ipv4_pkt, direction, Instant::recent())
|
||||
}
|
||||
|
||||
fn inspect_at(
|
||||
&mut self,
|
||||
ipv4_pkt: &Ipv4Packet<&[u8]>,
|
||||
direction: FlowDirection,
|
||||
now: Instant,
|
||||
) -> FlowMatch {
|
||||
// Perform lazy garbage collection of expired flow entries
|
||||
self.sweep_if_due(now);
|
||||
|
||||
// Later fragments do not contain enough transport information to bind
|
||||
// them to an exact flow, so packet policy must decide
|
||||
if ipv4_pkt.more_frags() || ipv4_pkt.frag_offset() != 0 {
|
||||
return FlowMatch::Untracked;
|
||||
}
|
||||
|
||||
match ipv4_pkt.next_header() {
|
||||
IpProtocol::Icmp => self.inspect_icmp(ipv4_pkt, direction, now),
|
||||
IpProtocol::Tcp => self.inspect_tcp(ipv4_pkt, direction, now),
|
||||
IpProtocol::Udp => self.inspect_udp(ipv4_pkt, direction, now),
|
||||
_ => FlowMatch::Untracked,
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn commit(&mut self, pending: PendingFlow) -> bool {
|
||||
let PendingFlow { key, flow } = pending;
|
||||
|
||||
// Reject new entries when full while allowing existing entries to be updated
|
||||
if !self.flows.contains_key(&key) && self.flows.len() >= MAX_FLOWS {
|
||||
return false;
|
||||
}
|
||||
|
||||
self.flows.insert(key, flow);
|
||||
|
||||
true
|
||||
}
|
||||
|
||||
pub(crate) fn clear(&mut self) {
|
||||
self.flows.clear();
|
||||
}
|
||||
|
||||
fn sweep_if_due(&mut self, now: Instant) {
|
||||
if now < self.next_sweep {
|
||||
return;
|
||||
}
|
||||
|
||||
self.flows.retain(|_, flow| now < flow.expires_at);
|
||||
|
||||
self.next_sweep = now + SWEEP_INTERVAL;
|
||||
}
|
||||
|
||||
fn match_existing_flow(
|
||||
&self,
|
||||
key: FlowKey,
|
||||
direction: FlowDirection,
|
||||
now: Instant,
|
||||
timeout: Duration,
|
||||
) -> Option<FlowMatch> {
|
||||
let existing = self.flows.get(&key).copied()?;
|
||||
|
||||
// Treat expired entries as missing even before the next sweep
|
||||
if now >= existing.expires_at {
|
||||
return None;
|
||||
}
|
||||
|
||||
if direction == existing.initiating_direction {
|
||||
Some(FlowMatch::candidate(
|
||||
key,
|
||||
existing.initiating_direction,
|
||||
now,
|
||||
timeout,
|
||||
))
|
||||
} else {
|
||||
Some(FlowMatch::Allowed)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl FlowDirection {
|
||||
fn host_vm_pair<T>(self, src: T, dst: T) -> (T, T) {
|
||||
match self {
|
||||
Self::FromHost => (src, dst),
|
||||
Self::FromVm => (dst, src),
|
||||
}
|
||||
}
|
||||
|
||||
fn opposite(self) -> Self {
|
||||
match self {
|
||||
Self::FromHost => Self::FromVm,
|
||||
Self::FromVm => Self::FromHost,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::test_support::{HOST, VM, inspect_from_host, ipv4_packet, udp_packet};
|
||||
use super::{
|
||||
Duration, FlowDirection, FlowMatch, FlowTable, Instant, MAX_FLOWS, SWEEP_INTERVAL,
|
||||
};
|
||||
use smoltcp::wire::{IpProtocol, Ipv4Packet};
|
||||
|
||||
#[test]
|
||||
fn all_ipv4_fragments_are_untracked() {
|
||||
for (more_fragments, offset) in [(true, 0), (false, 8)] {
|
||||
let mut bytes = udp_packet(HOST, 50_000, VM, 53);
|
||||
let mut packet = Ipv4Packet::new_unchecked(bytes.as_mut_slice());
|
||||
packet.set_frag_offset(offset);
|
||||
packet.set_more_frags(more_fragments);
|
||||
|
||||
let mut tracker = FlowTable::new();
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &bytes),
|
||||
FlowMatch::Untracked
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn expired_tuple_is_not_matched_before_the_next_sweep() {
|
||||
let mut tracker = FlowTable::new();
|
||||
let datagram = udp_packet(HOST, 50_000, VM, 53);
|
||||
let FlowMatch::Candidate(pending) = inspect_from_host(&mut tracker, &datagram) else {
|
||||
panic!("expected a candidate");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
|
||||
let now = Instant::recent();
|
||||
let flow = tracker.flows.values_mut().next().unwrap();
|
||||
flow.expires_at = now - Duration::from_secs(1);
|
||||
tracker.next_sweep = now + SWEEP_INTERVAL;
|
||||
|
||||
let reply = udp_packet(VM, 53, HOST, 50_000);
|
||||
assert!(matches!(
|
||||
tracker.inspect_at(
|
||||
&Ipv4Packet::new_checked(reply.as_slice()).unwrap(),
|
||||
FlowDirection::FromVm,
|
||||
now,
|
||||
),
|
||||
FlowMatch::Candidate(_)
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn global_limit_rejects_new_tuple_but_allows_replacement() {
|
||||
let mut tracker = FlowTable::new();
|
||||
for port in 10_000..10_000 + MAX_FLOWS as u16 {
|
||||
let datagram = udp_packet(HOST, port, VM, 53);
|
||||
let FlowMatch::Candidate(pending) = inspect_from_host(&mut tracker, &datagram) else {
|
||||
panic!("expected a candidate");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
}
|
||||
|
||||
let over_limit = udp_packet(HOST, 50_000, VM, 53);
|
||||
let FlowMatch::Candidate(pending) = inspect_from_host(&mut tracker, &over_limit) else {
|
||||
panic!("expected a candidate");
|
||||
};
|
||||
assert!(!tracker.commit(pending));
|
||||
|
||||
let replacement = udp_packet(HOST, 10_000, VM, 53);
|
||||
let FlowMatch::Candidate(pending) = inspect_from_host(&mut tracker, &replacement) else {
|
||||
panic!("expected a candidate");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unsupported_protocol_is_untracked() {
|
||||
let bytes = ipv4_packet(HOST, VM, IpProtocol::Unknown(253), 0);
|
||||
let mut tracker = FlowTable::new();
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &bytes),
|
||||
FlowMatch::Untracked
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod test_support {
|
||||
use super::{FlowDirection, FlowMatch, FlowTable};
|
||||
use smoltcp::wire::{IpProtocol, Ipv4Address, Ipv4Packet, TcpPacket, UdpPacket};
|
||||
|
||||
pub(super) const HOST: Ipv4Address = Ipv4Address::new(192, 168, 64, 1);
|
||||
pub(super) const VM: Ipv4Address = Ipv4Address::new(192, 168, 64, 2);
|
||||
|
||||
pub(super) fn inspect_from_host(tracker: &mut FlowTable, bytes: &[u8]) -> FlowMatch {
|
||||
let packet = Ipv4Packet::new_checked(bytes).unwrap();
|
||||
tracker.inspect(&packet, FlowDirection::FromHost)
|
||||
}
|
||||
|
||||
pub(super) fn inspect_from_vm(tracker: &mut FlowTable, bytes: &[u8]) -> FlowMatch {
|
||||
let packet = Ipv4Packet::new_checked(bytes).unwrap();
|
||||
tracker.inspect(&packet, FlowDirection::FromVm)
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
pub(super) struct TcpFlags {
|
||||
syn: bool,
|
||||
ack: bool,
|
||||
rst: bool,
|
||||
fin: bool,
|
||||
}
|
||||
|
||||
impl TcpFlags {
|
||||
pub(super) const SYN: Self = Self {
|
||||
syn: true,
|
||||
ack: false,
|
||||
rst: false,
|
||||
fin: false,
|
||||
};
|
||||
pub(super) const SYN_ACK: Self = Self {
|
||||
syn: true,
|
||||
ack: true,
|
||||
rst: false,
|
||||
fin: false,
|
||||
};
|
||||
pub(super) const ACK: Self = Self {
|
||||
syn: false,
|
||||
ack: true,
|
||||
rst: false,
|
||||
fin: false,
|
||||
};
|
||||
}
|
||||
|
||||
pub(super) fn tcp_packet(
|
||||
src_addr: Ipv4Address,
|
||||
src_port: u16,
|
||||
dst_addr: Ipv4Address,
|
||||
dst_port: u16,
|
||||
flags: TcpFlags,
|
||||
) -> Vec<u8> {
|
||||
let mut bytes = ipv4_packet(src_addr, dst_addr, IpProtocol::Tcp, 20);
|
||||
let mut tcp = TcpPacket::new_unchecked(&mut bytes[20..]);
|
||||
tcp.set_src_port(src_port);
|
||||
tcp.set_dst_port(dst_port);
|
||||
tcp.set_header_len(20);
|
||||
tcp.set_syn(flags.syn);
|
||||
tcp.set_ack(flags.ack);
|
||||
tcp.set_rst(flags.rst);
|
||||
tcp.set_fin(flags.fin);
|
||||
tcp.set_window_len(u16::MAX);
|
||||
bytes
|
||||
}
|
||||
|
||||
pub(super) fn udp_packet(
|
||||
src_addr: Ipv4Address,
|
||||
src_port: u16,
|
||||
dst_addr: Ipv4Address,
|
||||
dst_port: u16,
|
||||
) -> Vec<u8> {
|
||||
let mut bytes = ipv4_packet(src_addr, dst_addr, IpProtocol::Udp, 8);
|
||||
let mut udp = UdpPacket::new_unchecked(&mut bytes[20..]);
|
||||
udp.set_src_port(src_port);
|
||||
udp.set_dst_port(dst_port);
|
||||
udp.set_len(8);
|
||||
bytes
|
||||
}
|
||||
|
||||
pub(super) fn ipv4_packet(
|
||||
src_addr: Ipv4Address,
|
||||
dst_addr: Ipv4Address,
|
||||
protocol: IpProtocol,
|
||||
payload_len: usize,
|
||||
) -> Vec<u8> {
|
||||
let mut bytes = vec![0; 20 + payload_len];
|
||||
let total_len = bytes.len() as u16;
|
||||
let mut ipv4 = Ipv4Packet::new_unchecked(bytes.as_mut_slice());
|
||||
ipv4.set_version(4);
|
||||
ipv4.set_header_len(20);
|
||||
ipv4.set_total_len(total_len);
|
||||
ipv4.set_next_header(protocol);
|
||||
ipv4.set_src_addr(src_addr);
|
||||
ipv4.set_dst_addr(dst_addr);
|
||||
bytes
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,156 @@
|
||||
use super::{FlowDirection, FlowKey, FlowMatch, FlowTable};
|
||||
use coarsetime::{Duration, Instant};
|
||||
use smoltcp::wire::{Ipv4Address, Ipv4Packet, TcpPacket};
|
||||
|
||||
const SYN_TIMEOUT: Duration = Duration::from_secs(60);
|
||||
const TCP_TIMEOUT: Duration = Duration::from_secs(5 * 60);
|
||||
|
||||
impl FlowTable {
|
||||
pub(super) fn inspect_tcp(
|
||||
&mut self,
|
||||
ipv4_pkt: &Ipv4Packet<&[u8]>,
|
||||
direction: FlowDirection,
|
||||
now: Instant,
|
||||
) -> FlowMatch {
|
||||
let Ok(tcp) = TcpPacket::new_checked(ipv4_pkt.payload()) else {
|
||||
return FlowMatch::Denied;
|
||||
};
|
||||
if tcp.src_port() == 0 || tcp.dst_port() == 0 {
|
||||
return FlowMatch::Denied;
|
||||
}
|
||||
|
||||
let key = FlowKey::tcp(
|
||||
direction,
|
||||
(ipv4_pkt.src_addr(), tcp.src_port()),
|
||||
(ipv4_pkt.dst_addr(), tcp.dst_port()),
|
||||
);
|
||||
|
||||
// A bare SYN may be either a retransmission or a new connection reusing
|
||||
// the tuple in either direction. Always return it to policy, and do not
|
||||
// replace an existing permission until the candidate is committed.
|
||||
let is_initial_syn = tcp.syn() && !tcp.ack() && !tcp.fin() && !tcp.rst();
|
||||
|
||||
if is_initial_syn {
|
||||
return FlowMatch::candidate(key, direction, now, SYN_TIMEOUT);
|
||||
}
|
||||
|
||||
self.match_existing_flow(key, direction, now, TCP_TIMEOUT)
|
||||
.unwrap_or(FlowMatch::Untracked)
|
||||
}
|
||||
}
|
||||
|
||||
impl FlowKey {
|
||||
fn tcp(direction: FlowDirection, src: (Ipv4Address, u16), dst: (Ipv4Address, u16)) -> Self {
|
||||
let ((host_addr, host_port), (vm_addr, vm_port)) = direction.host_vm_pair(src, dst);
|
||||
Self::Tcp {
|
||||
host_addr,
|
||||
host_port,
|
||||
vm_addr,
|
||||
vm_port,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::super::test_support::{
|
||||
HOST, TcpFlags, VM, inspect_from_host, inspect_from_vm, tcp_packet,
|
||||
};
|
||||
use super::super::{FlowMatch, FlowTable};
|
||||
|
||||
fn admit_host_syn(tracker: &mut FlowTable) {
|
||||
let syn = tcp_packet(HOST, 49_152, VM, 22, TcpFlags::SYN);
|
||||
let FlowMatch::Candidate(pending) = inspect_from_host(tracker, &syn) else {
|
||||
panic!("expected a candidate");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn admitted_syn_creates_an_exact_reverse_permission() {
|
||||
let mut tracker = FlowTable::new();
|
||||
admit_host_syn(&mut tracker);
|
||||
let original_expiry = tracker.flows.values().next().unwrap().expires_at;
|
||||
|
||||
let reply = tcp_packet(VM, 22, HOST, 49_152, TcpFlags::ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &reply),
|
||||
FlowMatch::Allowed
|
||||
));
|
||||
assert_eq!(
|
||||
tracker.flows.values().next().unwrap().expires_at,
|
||||
original_expiry,
|
||||
"return traffic must not refresh a tuple"
|
||||
);
|
||||
|
||||
let initiator_data = tcp_packet(HOST, 49_152, VM, 22, TcpFlags::ACK);
|
||||
let FlowMatch::Candidate(pending) = inspect_from_host(&mut tracker, &initiator_data) else {
|
||||
panic!("expected an initiator candidate");
|
||||
};
|
||||
assert_eq!(
|
||||
tracker.flows.values().next().unwrap().expires_at,
|
||||
original_expiry,
|
||||
"inspection must not refresh a tuple before policy accepts it"
|
||||
);
|
||||
assert!(tracker.commit(pending));
|
||||
assert!(tracker.flows.values().next().unwrap().expires_at > original_expiry);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn new_tuple_requires_a_clean_syn() {
|
||||
let mut tracker = FlowTable::new();
|
||||
|
||||
for flags in [TcpFlags::ACK, TcpFlags::SYN_ACK] {
|
||||
let packet = tcp_packet(HOST, 49_152, VM, 22, flags);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &packet),
|
||||
FlowMatch::Untracked
|
||||
));
|
||||
}
|
||||
|
||||
let syn = tcp_packet(HOST, 49_152, VM, 22, TcpFlags::SYN);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &syn),
|
||||
FlowMatch::Candidate(_)
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reply_requires_the_exact_tuple() {
|
||||
let mut tracker = FlowTable::new();
|
||||
admit_host_syn(&mut tracker);
|
||||
|
||||
let wrong_port = tcp_packet(VM, 22, HOST, 49_153, TcpFlags::ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &wrong_port),
|
||||
FlowMatch::Untracked
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reverse_bare_syn_is_a_new_policy_candidate() {
|
||||
let mut tracker = FlowTable::new();
|
||||
admit_host_syn(&mut tracker);
|
||||
|
||||
let reverse_syn = tcp_packet(VM, 22, HOST, 49_152, TcpFlags::SYN);
|
||||
let FlowMatch::Candidate(pending) = inspect_from_vm(&mut tracker, &reverse_syn) else {
|
||||
panic!("reverse SYN must return to policy");
|
||||
};
|
||||
|
||||
// Inspection is provisional: a policy rejection leaves the admitted
|
||||
// flow untouched.
|
||||
let old_flow_reply = tcp_packet(VM, 22, HOST, 49_152, TcpFlags::ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &old_flow_reply),
|
||||
FlowMatch::Allowed
|
||||
));
|
||||
|
||||
assert!(tracker.commit(pending));
|
||||
|
||||
let reverse_reply = tcp_packet(HOST, 49_152, VM, 22, TcpFlags::SYN_ACK);
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &reverse_reply),
|
||||
FlowMatch::Allowed
|
||||
));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,92 @@
|
||||
use super::{FlowDirection, FlowKey, FlowMatch, FlowTable};
|
||||
use coarsetime::{Duration, Instant};
|
||||
use smoltcp::wire::{Ipv4Address, Ipv4Packet, UdpPacket};
|
||||
|
||||
const UDP_TIMEOUT: Duration = Duration::from_secs(30);
|
||||
|
||||
impl FlowTable {
|
||||
pub(super) fn inspect_udp(
|
||||
&mut self,
|
||||
ipv4_pkt: &Ipv4Packet<&[u8]>,
|
||||
direction: FlowDirection,
|
||||
now: Instant,
|
||||
) -> FlowMatch {
|
||||
let Ok(udp) = UdpPacket::new_checked(ipv4_pkt.payload()) else {
|
||||
return FlowMatch::Denied;
|
||||
};
|
||||
if udp.dst_port() == 0 {
|
||||
return FlowMatch::Denied;
|
||||
}
|
||||
|
||||
// Deliberate policy, not merely a packet-format validation:
|
||||
// RFC 8085, §5.1 says UDP senders SHOULD NOT use source port zero.
|
||||
//
|
||||
// We enforce this recommendation to retain source-port entropy
|
||||
// and protection against off-path packet injection.
|
||||
if udp.src_port() == 0 {
|
||||
return FlowMatch::Denied;
|
||||
}
|
||||
|
||||
let key = FlowKey::udp(
|
||||
direction,
|
||||
(ipv4_pkt.src_addr(), udp.src_port()),
|
||||
(ipv4_pkt.dst_addr(), udp.dst_port()),
|
||||
);
|
||||
|
||||
if let Some(matched) = self.match_existing_flow(key, direction, now, UDP_TIMEOUT) {
|
||||
return matched;
|
||||
}
|
||||
|
||||
FlowMatch::candidate(key, direction, now, UDP_TIMEOUT)
|
||||
}
|
||||
}
|
||||
|
||||
impl FlowKey {
|
||||
fn udp(direction: FlowDirection, src: (Ipv4Address, u16), dst: (Ipv4Address, u16)) -> Self {
|
||||
let ((host_addr, host_port), (vm_addr, vm_port)) = direction.host_vm_pair(src, dst);
|
||||
Self::Udp {
|
||||
host_addr,
|
||||
host_port,
|
||||
vm_addr,
|
||||
vm_port,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::super::test_support::{HOST, VM, inspect_from_host, inspect_from_vm, udp_packet};
|
||||
use super::super::{FlowMatch, FlowTable};
|
||||
|
||||
#[test]
|
||||
fn rejects_zero_source_port_per_rfc_8085() {
|
||||
let mut tracker = FlowTable::new();
|
||||
let datagram = udp_packet(HOST, 0, VM, 5353);
|
||||
|
||||
assert!(matches!(
|
||||
inspect_from_host(&mut tracker, &datagram),
|
||||
FlowMatch::Denied
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn udp_reply_requires_an_exact_admitted_request() {
|
||||
let mut tracker = FlowTable::new();
|
||||
let request = udp_packet(HOST, 50_000, VM, 5353);
|
||||
let reply = udp_packet(VM, 5353, HOST, 50_000);
|
||||
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &reply),
|
||||
FlowMatch::Candidate(_)
|
||||
));
|
||||
|
||||
let FlowMatch::Candidate(pending) = inspect_from_host(&mut tracker, &request) else {
|
||||
panic!("expected a candidate");
|
||||
};
|
||||
assert!(tracker.commit(pending));
|
||||
assert!(matches!(
|
||||
inspect_from_vm(&mut tracker, &reply),
|
||||
FlowMatch::Allowed
|
||||
));
|
||||
}
|
||||
}
|
||||
+135
-56
@@ -1,8 +1,12 @@
|
||||
use crate::proxy::conntrack::ConntrackResult;
|
||||
use crate::dhcp_snooper::message_matches_bootp_client;
|
||||
use crate::proxy::flows::{FlowDirection, FlowMatch};
|
||||
use crate::proxy::udp_packet_helper::UdpPacketHelper;
|
||||
use crate::proxy::{Action, Direction, Proxy, Rule, select_rules};
|
||||
use crate::proxy::{Direction, PolicyDecision, Proxy};
|
||||
use anyhow::{Context, Result};
|
||||
use smoltcp::wire::{EthernetFrame, EthernetProtocol, Ipv4Packet, UdpPacket};
|
||||
use dhcproto::Decodable;
|
||||
use dhcproto::v4::Opcode;
|
||||
use smoltcp::phy::ChecksumCapabilities;
|
||||
use smoltcp::wire::{EthernetFrame, EthernetProtocol, Ipv4Packet, Ipv4Repr, UdpPacket};
|
||||
|
||||
impl Proxy<'_> {
|
||||
pub(crate) fn process_frame_from_host(&mut self, frame: &EthernetFrame<&[u8]>) -> Result<()> {
|
||||
@@ -13,7 +17,7 @@ impl Proxy<'_> {
|
||||
|
||||
// Snoop bootpd(8) replies from the host to
|
||||
// figure out the IP assigned to the VM
|
||||
if frame.dst_addr() == self.vm_mac_address {
|
||||
if frame.dst_addr() == self.vm_mac_address || frame.dst_addr().is_broadcast() {
|
||||
self.snoop(frame);
|
||||
}
|
||||
|
||||
@@ -41,66 +45,68 @@ impl Proxy<'_> {
|
||||
match frame.ethertype() {
|
||||
EthernetProtocol::Arp => Some(()),
|
||||
EthernetProtocol::Ipv4 => {
|
||||
let ipv4_pkt = Ipv4Packet::new_checked(frame.payload()).ok()?;
|
||||
let ipv4_pkt = Ipv4Packet::new_unchecked(frame.payload());
|
||||
Ipv4Repr::parse(&ipv4_pkt, &ChecksumCapabilities::ignored()).ok()?;
|
||||
|
||||
self.allowed_from_host_ipv4(&ipv4_pkt)
|
||||
}
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
fn allowed_from_host_ipv4(&mut self, ipv4_pkt: &Ipv4Packet<&[u8]>) -> Option<()> {
|
||||
if !self.stateful_policy {
|
||||
pub(super) fn allowed_from_host_ipv4(&mut self, ipv4_pkt: &Ipv4Packet<&[u8]>) -> Option<()> {
|
||||
// Backwards compatibility with Softnet consumers that only use stateless rules
|
||||
if self.flows.is_none() {
|
||||
return Some(());
|
||||
}
|
||||
|
||||
if let Some(rules) = select_rules(&self.rules, ipv4_pkt.src_addr(), Direction::In) {
|
||||
// DHCP is required to maintain the VM's lease and must bypass user-specified rules
|
||||
if self.is_allowed_dhcp_response(ipv4_pkt) {
|
||||
return Some(());
|
||||
}
|
||||
|
||||
// Only process packets addressed to the VM's current IP
|
||||
let Some(lease) = self.dhcp_snooper.lease() else {
|
||||
return None;
|
||||
};
|
||||
if !lease.is_valid_for(ipv4_pkt.dst_addr()) {
|
||||
return None;
|
||||
}
|
||||
|
||||
// Stateless rules decide each packet immediately, without consulting the conntrack
|
||||
if let Some((action, _)) = rules
|
||||
.iter()
|
||||
.find(|(_, rule)| matches!(rule, Rule::Stateless(_)))
|
||||
{
|
||||
return (*action == Action::Allow).then_some(());
|
||||
}
|
||||
|
||||
// Existing connections follow conntrack; new ones require an inbound stateful allow rule
|
||||
return match self.conntrack.inspect_from_host(ipv4_pkt) {
|
||||
ConntrackResult::Allowed => Some(()),
|
||||
ConntrackResult::Denied => None,
|
||||
ConntrackResult::New(pending) => {
|
||||
let allow_new = rules.iter().any(|(action, rule)| {
|
||||
*action == Action::Allow
|
||||
&& matches!(
|
||||
rule,
|
||||
Rule::Stateful {
|
||||
direction: Direction::In,
|
||||
..
|
||||
}
|
||||
)
|
||||
});
|
||||
|
||||
if !allow_new {
|
||||
return None;
|
||||
}
|
||||
|
||||
self.conntrack.commit(pending).then_some(())
|
||||
}
|
||||
};
|
||||
// DHCP is required to maintain the VM's lease and must bypass user-specified rules
|
||||
if self.is_allowed_dhcp_response(ipv4_pkt) {
|
||||
return Some(());
|
||||
}
|
||||
|
||||
Some(())
|
||||
// Consult the flow table before evaluating inbound policy
|
||||
// so established flows are not treated as new traffic
|
||||
let pending = if self
|
||||
.dhcp_snooper
|
||||
.lease()
|
||||
.as_ref()
|
||||
.is_some_and(|lease| lease.is_valid_for(ipv4_pkt.dst_addr()))
|
||||
{
|
||||
match self
|
||||
.flows
|
||||
.as_mut()?
|
||||
.inspect(ipv4_pkt, FlowDirection::FromHost)
|
||||
{
|
||||
FlowMatch::Allowed => return Some(()),
|
||||
FlowMatch::Denied => return None,
|
||||
FlowMatch::Candidate(pending) => Some(pending),
|
||||
FlowMatch::Untracked => None,
|
||||
}
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
// The flow is either pending or untracked, evaluate it against inbound policy
|
||||
match self
|
||||
.rules
|
||||
.policy_decision(ipv4_pkt.src_addr(), Direction::In)
|
||||
{
|
||||
// Return traffic was handled above; enforce explicit inbound blocks here
|
||||
Some(PolicyDecision::Block) => None,
|
||||
|
||||
// Stateless policy is outbound-only; fail closed if this invariant is violated
|
||||
Some(PolicyDecision::AllowStateless) => None,
|
||||
|
||||
// Untracked packets cannot satisfy stateful policy
|
||||
Some(PolicyDecision::AllowStateful) => self.admit_with_tracking(pending?),
|
||||
|
||||
// No inbound rule matched, so allow by default. Track the flow when needed
|
||||
// so its reply is not treated as a new outbound flow
|
||||
None => {
|
||||
self.admit_with_tracking_if_stateful(pending, ipv4_pkt.src_addr(), Direction::Out)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn snoop(&mut self, frame: &EthernetFrame<&[u8]>) {
|
||||
@@ -122,7 +128,13 @@ impl Proxy<'_> {
|
||||
Err(_) => return,
|
||||
};
|
||||
|
||||
let address_and_dns_ips_saved = self.dhcp_snooper.address_and_dns_ips();
|
||||
self.dhcp_snooper.register_dhcp_reply(udp_pkt.payload());
|
||||
if address_and_dns_ips_saved != self.dhcp_snooper.address_and_dns_ips()
|
||||
&& let Some(flows) = &mut self.flows
|
||||
{
|
||||
flows.clear();
|
||||
}
|
||||
}
|
||||
|
||||
fn is_allowed_dhcp_response(&self, ipv4_pkt: &Ipv4Packet<&[u8]>) -> bool {
|
||||
@@ -132,8 +144,75 @@ impl Proxy<'_> {
|
||||
return false;
|
||||
}
|
||||
|
||||
UdpPacket::new_checked(ipv4_pkt.payload())
|
||||
.map(|udp_pkt| udp_pkt.is_dhcp_response())
|
||||
.unwrap_or(false)
|
||||
let Ok(udp_pkt) = UdpPacket::new_checked(ipv4_pkt.payload()) else {
|
||||
return false;
|
||||
};
|
||||
|
||||
// Require the standard DHCP server and client ports
|
||||
if !udp_pkt.is_dhcp_response() {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Require the BOOTP client hardware address to match this VM
|
||||
// (symmetric with is_allowed_dhcp_request / #191 on the VM→host path)
|
||||
let mut decoder = dhcproto::v4::Decoder::new(udp_pkt.payload());
|
||||
let Ok(message) = dhcproto::v4::Message::decode(&mut decoder) else {
|
||||
return false;
|
||||
};
|
||||
|
||||
message_matches_bootp_client(&message, Opcode::BootReply, self.vm_mac_address.0)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use crate::dhcp_snooper::message_matches_bootp_client;
|
||||
use dhcproto::Decodable;
|
||||
use dhcproto::v4::{DhcpOption, Message, MessageType, Opcode};
|
||||
use dhcproto::{Encodable, Encoder};
|
||||
use smoltcp::wire::Ipv4Address;
|
||||
|
||||
const VM_MAC: [u8; 6] = [0x02, 0x00, 0x00, 0x00, 0x00, 0x01];
|
||||
const OTHER_MAC: [u8; 6] = [0x02, 0x00, 0x00, 0x00, 0x00, 0x02];
|
||||
|
||||
#[test]
|
||||
fn dhcp_boot_reply_chaddr_must_match_vm() {
|
||||
let own = encode_boot_reply(VM_MAC);
|
||||
let foreign = encode_boot_reply(OTHER_MAC);
|
||||
|
||||
let mut dec = dhcproto::v4::Decoder::new(&own);
|
||||
let own_msg = Message::decode(&mut dec).unwrap();
|
||||
let mut dec = dhcproto::v4::Decoder::new(&foreign);
|
||||
let foreign_msg = Message::decode(&mut dec).unwrap();
|
||||
|
||||
assert!(message_matches_bootp_client(
|
||||
&own_msg,
|
||||
Opcode::BootReply,
|
||||
VM_MAC
|
||||
));
|
||||
assert!(!message_matches_bootp_client(
|
||||
&foreign_msg,
|
||||
Opcode::BootReply,
|
||||
VM_MAC
|
||||
));
|
||||
}
|
||||
|
||||
fn encode_boot_reply(chaddr: [u8; 6]) -> Vec<u8> {
|
||||
let mut message = Message::new(
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
Ipv4Address::new(192, 168, 64, 2),
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
&chaddr,
|
||||
);
|
||||
message.set_opcode(Opcode::BootReply);
|
||||
message
|
||||
.opts_mut()
|
||||
.insert(DhcpOption::MessageType(MessageType::Ack));
|
||||
message.opts_mut().insert(DhcpOption::AddressLeaseTime(600));
|
||||
|
||||
let mut encoded = Vec::new();
|
||||
message.encode(&mut Encoder::new(&mut encoded)).unwrap();
|
||||
encoded
|
||||
}
|
||||
}
|
||||
|
||||
+137
-60
@@ -1,6 +1,6 @@
|
||||
mod conntrack;
|
||||
mod control;
|
||||
mod exposed_port;
|
||||
mod flows;
|
||||
mod host;
|
||||
mod port_forwarder;
|
||||
mod rule;
|
||||
@@ -14,15 +14,15 @@ use crate::host::NetType;
|
||||
use crate::poller::Poller;
|
||||
use crate::vm::VM;
|
||||
use anyhow::Result;
|
||||
use conntrack::Conntrack;
|
||||
use control::Control;
|
||||
use control::{Control, normalize_rules};
|
||||
pub use exposed_port::ExposedPort;
|
||||
use flows::{FlowTable, PendingFlow};
|
||||
use ipnet::Ipv4Net;
|
||||
use mac_address::MacAddress;
|
||||
use port_forwarder::PortForwarder;
|
||||
pub use rule::{Direction, Rule, Target};
|
||||
pub(crate) use rules::{Action, Rules, build_rules, has_stateful_rules, rule_count, select_rules};
|
||||
use smoltcp::wire::EthernetFrame;
|
||||
pub(crate) use rules::{PolicyDecision, Rules};
|
||||
use smoltcp::wire::{EthernetFrame, Ipv4Address};
|
||||
use std::io::ErrorKind;
|
||||
use std::os::unix::io::{AsRawFd, RawFd};
|
||||
use std::time::Duration;
|
||||
@@ -35,9 +35,8 @@ pub struct Proxy<'proxy> {
|
||||
vm_mac_address: smoltcp::wire::EthernetAddress,
|
||||
dhcp_snooper: DhcpSnooper,
|
||||
rules: Rules,
|
||||
stateful_policy: bool,
|
||||
control: Option<Control>,
|
||||
conntrack: Conntrack,
|
||||
flows: Option<FlowTable>,
|
||||
enobufs_encountered: bool,
|
||||
port_forwarder: PortForwarder,
|
||||
}
|
||||
@@ -52,6 +51,9 @@ impl Proxy<'_> {
|
||||
exposed_ports: Vec<ExposedPort>,
|
||||
control_fd: Option<RawFd>,
|
||||
) -> Result<Proxy<'proxy>> {
|
||||
let allow = normalize_rules(allow);
|
||||
let block = normalize_rules(block);
|
||||
|
||||
let vm = VM::new(vm_fd)?;
|
||||
let host = Host::new(
|
||||
vm_net_type,
|
||||
@@ -70,19 +72,21 @@ impl Proxy<'_> {
|
||||
poller_timeout,
|
||||
)?;
|
||||
|
||||
let rules = build_rules(host.gateway_ip, &allow, &block);
|
||||
let stateful_policy = has_stateful_rules(&rules);
|
||||
let rules = Rules::new(host.gateway_ip, &allow, &block);
|
||||
|
||||
// Any stateful rule enables flow inspection for the whole VM, including
|
||||
// traffic admitted through implicit global, gateway, and DNS fallbacks
|
||||
let flows = rules.has_stateful().then(FlowTable::new);
|
||||
|
||||
Ok(Proxy {
|
||||
vm,
|
||||
host,
|
||||
poller,
|
||||
vm_mac_address: smoltcp::wire::EthernetAddress(vm_mac_address.bytes()),
|
||||
dhcp_snooper: DhcpSnooper::new(poller_timeout),
|
||||
dhcp_snooper: DhcpSnooper::new(poller_timeout, vm_mac_address.bytes()),
|
||||
rules,
|
||||
stateful_policy,
|
||||
control,
|
||||
conntrack: Conntrack::new(),
|
||||
flows,
|
||||
enobufs_encountered: false,
|
||||
port_forwarder: PortForwarder::new(exposed_ports),
|
||||
})
|
||||
@@ -104,11 +108,13 @@ impl Proxy<'_> {
|
||||
loop {
|
||||
let (vm_readable, host_readable, interrupt) = self.poller.wait()?;
|
||||
|
||||
// Update coarse time for DHCP snooping and conntrack
|
||||
coarsetime::Instant::update();
|
||||
// kqueue does not report peer disconnects for Unix datagram sockets.
|
||||
if !self.vm.is_connected()? {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Expire stale flows before processing packets
|
||||
self.conntrack.tick();
|
||||
// Update coarse time for DHCP snooping and flows
|
||||
coarsetime::Instant::update();
|
||||
|
||||
// Service control on every wake (including timeouts) so a bounded read or a pending
|
||||
// response continues making progress even when no new edge is generated.
|
||||
@@ -143,7 +149,7 @@ impl Proxy<'_> {
|
||||
loop {
|
||||
match self.vm.read(buf) {
|
||||
Ok(n) => {
|
||||
// Update coarse time for DHCP snooping and conntrack
|
||||
// Update coarse time for DHCP snooping and flows
|
||||
coarsetime::Instant::update();
|
||||
|
||||
if let Ok(frame) = EthernetFrame::new_checked(&buf[..n]) {
|
||||
@@ -171,7 +177,7 @@ impl Proxy<'_> {
|
||||
loop {
|
||||
match self.host.read(batch, bufs) {
|
||||
Ok(pktcnt) => {
|
||||
// Update coarse time for DHCP snooping and conntrack
|
||||
// Update coarse time for DHCP snooping and flows
|
||||
coarsetime::Instant::update();
|
||||
|
||||
for buf in batch.packet_sized_bufs(bufs).take(pktcnt) {
|
||||
@@ -206,9 +212,9 @@ impl Proxy<'_> {
|
||||
}
|
||||
};
|
||||
|
||||
// Invalidate tracked flows whenever the policy changes
|
||||
if control.policy_changed() {
|
||||
self.stateful_policy = has_stateful_rules(&self.rules);
|
||||
self.conntrack.clear();
|
||||
self.flows = self.rules.has_stateful().then(FlowTable::new);
|
||||
}
|
||||
|
||||
if keep_open {
|
||||
@@ -225,19 +231,44 @@ impl Proxy<'_> {
|
||||
log::warn!("failed to shut down Softnet control socket: {err:#}");
|
||||
}
|
||||
}
|
||||
|
||||
/// Commits the pending flow, rejecting the packet if the table cannot store it.
|
||||
fn admit_with_tracking(&mut self, pending: PendingFlow) -> Option<()> {
|
||||
self.flows.as_mut()?.commit(pending).then_some(())
|
||||
}
|
||||
|
||||
/// Commits a pending flow for trackable packets; untracked packets proceed without one.
|
||||
fn admit_with_tracking_if_trackable(&mut self, pending: Option<PendingFlow>) -> Option<()> {
|
||||
match pending {
|
||||
Some(pending) => self.admit_with_tracking(pending),
|
||||
None => Some(()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Commits a pending flow when the return-direction rule is stateful.
|
||||
fn admit_with_tracking_if_stateful(
|
||||
&mut self,
|
||||
pending: Option<PendingFlow>,
|
||||
peer_addr: Ipv4Address,
|
||||
return_direction: Direction,
|
||||
) -> Option<()> {
|
||||
if self.rules.is_stateful(peer_addr, return_direction) {
|
||||
self.admit_with_tracking_if_trackable(pending)
|
||||
} else {
|
||||
Some(())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use crate::NetType;
|
||||
use crate::dhcp_snooper::Lease;
|
||||
use crate::proxy::{Action, Proxy, Rule, Target};
|
||||
use ipnet::Ipv4Net;
|
||||
use crate::proxy::Proxy;
|
||||
use mac_address::MacAddress;
|
||||
use nix::sys::socket::{AddressFamily, SockFlag, SockType, socketpair};
|
||||
use prefix_trie::PrefixMap;
|
||||
use serial_test::serial;
|
||||
use smoltcp::wire::{Ipv4Address, Ipv4Packet};
|
||||
use smoltcp::wire::{IpProtocol, Ipv4Address, Ipv4Packet, UdpPacket};
|
||||
use std::collections::HashSet;
|
||||
use std::os::fd::AsRawFd;
|
||||
use std::str::FromStr;
|
||||
@@ -249,13 +280,7 @@ mod tests {
|
||||
let vm_ip = Ipv4Address::from_str("192.168.0.2").unwrap();
|
||||
let mut proxy = create_proxy(vm_ip, vec!["66.66.0.0/16"], vec!["66.66.0.0/16"]);
|
||||
|
||||
assert_eq!(
|
||||
proxy.rules,
|
||||
PrefixMap::<Ipv4Net, Vec<(Action, Rule)>>::from_iter(vec![(
|
||||
Ipv4Net::from_str("66.66.0.0/16").unwrap(),
|
||||
vec![(Action::Block, "66.66.0.0/16".parse().unwrap(),)]
|
||||
),])
|
||||
);
|
||||
assert_eq!(proxy.rules.len(), 1);
|
||||
|
||||
assert!(allowed_from_vm_ipv4(&mut proxy, vm_ip, "66.66.66.66").is_none());
|
||||
}
|
||||
@@ -266,19 +291,7 @@ mod tests {
|
||||
let vm_ip = Ipv4Address::from_str("192.168.0.2").unwrap();
|
||||
let mut proxy = create_proxy(vm_ip, vec!["33.33.33.33/32"], vec!["33.33.33.0/24"]);
|
||||
|
||||
assert_eq!(
|
||||
proxy.rules,
|
||||
PrefixMap::<Ipv4Net, Vec<(Action, Rule)>>::from_iter(vec![
|
||||
(
|
||||
Ipv4Net::from_str("33.33.33.33/32").unwrap(),
|
||||
vec![(Action::Allow, "33.33.33.33/32".parse().unwrap(),)]
|
||||
),
|
||||
(
|
||||
Ipv4Net::from_str("33.33.33.0/24").unwrap(),
|
||||
vec![(Action::Block, "33.33.33.0/24".parse().unwrap(),)]
|
||||
),
|
||||
])
|
||||
);
|
||||
assert_eq!(proxy.rules.len(), 2);
|
||||
|
||||
assert!(allowed_from_vm_ipv4(&mut proxy, vm_ip, "33.33.33.32").is_none());
|
||||
assert!(allowed_from_vm_ipv4(&mut proxy, vm_ip, "33.33.33.33").is_some());
|
||||
@@ -291,22 +304,7 @@ mod tests {
|
||||
let vm_ip = Ipv4Address::from_str("192.168.0.2").unwrap();
|
||||
let mut proxy = create_proxy(vm_ip, vec!["@host"], vec!["0.0.0.0/0"]);
|
||||
|
||||
assert_eq!(
|
||||
proxy.rules,
|
||||
PrefixMap::from_iter(vec![
|
||||
(
|
||||
proxy.host.gateway_ip.into(),
|
||||
vec![(
|
||||
Action::Allow,
|
||||
Rule::Stateless(Target::Prefix(proxy.host.gateway_ip.into())),
|
||||
)],
|
||||
),
|
||||
(
|
||||
Ipv4Net::from_str("0.0.0.0/0").unwrap(),
|
||||
vec![(Action::Block, "0.0.0.0/0".parse().unwrap(),)]
|
||||
),
|
||||
])
|
||||
);
|
||||
assert_eq!(proxy.rules.len(), 2);
|
||||
|
||||
// Access to global IPs should be disallowed because of --block=0.0.0.0/0
|
||||
assert!(allowed_from_vm_ipv4(&mut proxy, vm_ip, "8.8.8.8").is_none());
|
||||
@@ -316,6 +314,63 @@ mod tests {
|
||||
assert!(allowed_from_vm_ipv4(&mut proxy, vm_ip, &gateway_ip).is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial]
|
||||
fn test_bare_default_block_applies_in_both_directions_in_stateful_mode() {
|
||||
let vm_ip = Ipv4Address::new(192, 168, 0, 2);
|
||||
let fallback_peer = Ipv4Address::new(203, 0, 113, 1);
|
||||
let explicitly_allowed_peer = Ipv4Address::new(192, 0, 2, 1);
|
||||
let mut proxy = create_proxy(vm_ip, vec!["in 192.0.2.0/24"], vec!["0.0.0.0/0"]);
|
||||
|
||||
let fallback_request = udp_packet(fallback_peer, 40_000, vm_ip, 1_234);
|
||||
let fallback_request = Ipv4Packet::new_checked(fallback_request.as_slice()).unwrap();
|
||||
assert!(proxy.allowed_from_host_ipv4(&fallback_request).is_none());
|
||||
|
||||
let fallback_reply = udp_packet(vm_ip, 1_234, fallback_peer, 40_000);
|
||||
let fallback_reply = Ipv4Packet::new_checked(fallback_reply.as_slice()).unwrap();
|
||||
assert!(proxy.allowed_from_vm_ipv4(fallback_reply).is_none());
|
||||
|
||||
let allowed_request = udp_packet(explicitly_allowed_peer, 40_000, vm_ip, 1_234);
|
||||
let allowed_request = Ipv4Packet::new_checked(allowed_request.as_slice()).unwrap();
|
||||
assert!(proxy.allowed_from_host_ipv4(&allowed_request).is_some());
|
||||
|
||||
let allowed_reply = udp_packet(vm_ip, 1_234, explicitly_allowed_peer, 40_000);
|
||||
let allowed_reply = Ipv4Packet::new_checked(allowed_reply.as_slice()).unwrap();
|
||||
assert!(proxy.allowed_from_vm_ipv4(allowed_reply).is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial]
|
||||
fn test_directional_egress_block_does_not_block_reply_to_unmatched_inbound_flow() {
|
||||
let vm_ip = Ipv4Address::new(192, 168, 0, 2);
|
||||
let peer = Ipv4Address::new(203, 0, 113, 1);
|
||||
let mut proxy = create_proxy(vm_ip, vec![], vec!["out 203.0.113.0/24"]);
|
||||
|
||||
let request = udp_packet(peer, 40_000, vm_ip, 1_234);
|
||||
let request = Ipv4Packet::new_checked(request.as_slice()).unwrap();
|
||||
assert!(proxy.allowed_from_host_ipv4(&request).is_some());
|
||||
|
||||
let reply = udp_packet(vm_ip, 1_234, peer, 40_000);
|
||||
let reply = Ipv4Packet::new_checked(reply.as_slice()).unwrap();
|
||||
assert!(proxy.allowed_from_vm_ipv4(reply).is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[serial]
|
||||
fn test_directional_ingress_block_does_not_block_reply_to_bare_outbound_allow() {
|
||||
let vm_ip = Ipv4Address::new(192, 168, 0, 2);
|
||||
let peer = Ipv4Address::new(203, 0, 113, 1);
|
||||
let mut proxy = create_proxy(vm_ip, vec!["203.0.113.0/24"], vec!["in 203.0.113.0/24"]);
|
||||
|
||||
let request = udp_packet(vm_ip, 1_234, peer, 40_000);
|
||||
let request = Ipv4Packet::new_checked(request.as_slice()).unwrap();
|
||||
assert!(proxy.allowed_from_vm_ipv4(request).is_some());
|
||||
|
||||
let reply = udp_packet(peer, 40_000, vm_ip, 1_234);
|
||||
let reply = Ipv4Packet::new_checked(reply.as_slice()).unwrap();
|
||||
assert!(proxy.allowed_from_host_ipv4(&reply).is_some());
|
||||
}
|
||||
|
||||
fn create_proxy<'test>(vm_ip: Ipv4Address, allow: Vec<&str>, block: Vec<&str>) -> Proxy<'test> {
|
||||
let (vm_fd, _) = socketpair(
|
||||
AddressFamily::Unix,
|
||||
@@ -363,4 +418,26 @@ mod tests {
|
||||
|
||||
proxy.allowed_from_vm_ipv4(ipv4_pkt)
|
||||
}
|
||||
|
||||
fn udp_packet(
|
||||
src_addr: Ipv4Address,
|
||||
src_port: u16,
|
||||
dst_addr: Ipv4Address,
|
||||
dst_port: u16,
|
||||
) -> Vec<u8> {
|
||||
let mut bytes = vec![0; 28];
|
||||
let mut ipv4 = Ipv4Packet::new_unchecked(bytes.as_mut_slice());
|
||||
ipv4.set_version(4);
|
||||
ipv4.set_header_len(20);
|
||||
ipv4.set_total_len(28);
|
||||
ipv4.set_next_header(IpProtocol::Udp);
|
||||
ipv4.set_src_addr(src_addr);
|
||||
ipv4.set_dst_addr(dst_addr);
|
||||
|
||||
let mut udp = UdpPacket::new_unchecked(ipv4.payload_mut());
|
||||
udp.set_src_port(src_port);
|
||||
udp.set_dst_port(dst_port);
|
||||
udp.set_len(8);
|
||||
bytes
|
||||
}
|
||||
}
|
||||
|
||||
+6
-19
@@ -24,13 +24,6 @@ pub enum Direction {
|
||||
}
|
||||
|
||||
impl Rule {
|
||||
pub(super) fn target(self) -> Target {
|
||||
match self {
|
||||
Rule::Stateless(target) => target,
|
||||
Rule::Stateful { target, .. } => target,
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn normalized(self) -> Self {
|
||||
match self {
|
||||
Rule::Stateless(target) => Rule::Stateless(target.normalized()),
|
||||
@@ -137,18 +130,12 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn displays_normalized_rules() {
|
||||
assert_eq!(
|
||||
"in 10.1.2.3/8"
|
||||
.parse::<Rule>()
|
||||
.unwrap()
|
||||
.normalized()
|
||||
.to_string(),
|
||||
"in 10.0.0.0/8"
|
||||
);
|
||||
assert_eq!(
|
||||
"@host".parse::<Rule>().unwrap().normalized().to_string(),
|
||||
"@host"
|
||||
);
|
||||
for (input, expected) in [("in 10.1.2.3/8", "in 10.0.0.0/8"), ("@host", "@host")] {
|
||||
assert_eq!(
|
||||
input.parse::<Rule>().unwrap().normalized().to_string(),
|
||||
expected
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
+224
-147
@@ -9,194 +9,271 @@ pub(crate) enum Action {
|
||||
Allow,
|
||||
}
|
||||
|
||||
pub(crate) type Rules = PrefixMap<Ipv4Net, Vec<(Action, Rule)>>;
|
||||
|
||||
pub(crate) fn select_rules(
|
||||
rules: &Rules,
|
||||
address: Ipv4Address,
|
||||
direction: Direction,
|
||||
) -> Option<&[(Action, Rule)]> {
|
||||
if rules.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut prefix = Ipv4Net::from(address);
|
||||
loop {
|
||||
let (matched_prefix, matched_rules) = rules.get_lpm(&prefix)?;
|
||||
if matched_rules.iter().any(|(_, rule)| match rule {
|
||||
Rule::Stateless(_) => true,
|
||||
Rule::Stateful {
|
||||
direction: rule_direction,
|
||||
..
|
||||
} => *rule_direction == direction,
|
||||
}) {
|
||||
return Some(matched_rules.as_slice());
|
||||
}
|
||||
prefix = matched_prefix.supernet()?;
|
||||
}
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub(crate) enum PolicyDecision {
|
||||
Block,
|
||||
AllowStateless,
|
||||
AllowStateful,
|
||||
}
|
||||
|
||||
pub(crate) fn build_rules(host_address: Ipv4Address, allow: &[Rule], block: &[Rule]) -> Rules {
|
||||
let mut rules = PrefixMap::new();
|
||||
|
||||
for &rule in allow {
|
||||
insert_rule(&mut rules, rule, Action::Allow, host_address);
|
||||
}
|
||||
|
||||
for &rule in block {
|
||||
insert_rule(&mut rules, rule, Action::Block, host_address);
|
||||
}
|
||||
|
||||
rules
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
enum Mode {
|
||||
#[default]
|
||||
Legacy,
|
||||
Stateful,
|
||||
}
|
||||
|
||||
pub(crate) fn rule_count(rules: &Rules) -> usize {
|
||||
rules.into_iter().map(|(_, rules)| rules.len()).sum()
|
||||
#[derive(Default)]
|
||||
pub(crate) struct Rules {
|
||||
mode: Mode,
|
||||
inbound: PrefixMap<Ipv4Net, Action>,
|
||||
outbound: PrefixMap<Ipv4Net, Action>,
|
||||
}
|
||||
|
||||
pub(crate) fn has_stateful_rules(rules: &Rules) -> bool {
|
||||
rules.into_iter().any(|(_, rules)| {
|
||||
rules
|
||||
impl Rules {
|
||||
pub(crate) fn new(host_address: Ipv4Address, allow: &[Rule], block: &[Rule]) -> Self {
|
||||
// Preserve legacy behavior for bare-only policies. Once a directional rule
|
||||
// is present, compile the whole policy using directional semantics.
|
||||
let mode = if allow
|
||||
.iter()
|
||||
.any(|(_, rule)| matches!(rule, Rule::Stateful { .. }))
|
||||
})
|
||||
}
|
||||
.chain(block)
|
||||
.any(|rule| matches!(rule, Rule::Stateful { .. }))
|
||||
{
|
||||
Mode::Stateful
|
||||
} else {
|
||||
Mode::Legacy
|
||||
};
|
||||
|
||||
fn insert_rule(rules: &mut Rules, rule: Rule, mut action: Action, host_address: Ipv4Address) {
|
||||
let prefix = match rule.target() {
|
||||
Target::Prefix(prefix) => prefix,
|
||||
Target::Host => host_address.into(),
|
||||
};
|
||||
let rule = match rule {
|
||||
Rule::Stateless(_) => Rule::Stateless(Target::Prefix(prefix)),
|
||||
Rule::Stateful { direction, .. } => Rule::Stateful {
|
||||
direction,
|
||||
target: Target::Prefix(prefix),
|
||||
},
|
||||
};
|
||||
let mut rules = Self {
|
||||
mode,
|
||||
..Self::default()
|
||||
};
|
||||
|
||||
let prefix_rules = rules.entry(prefix).or_default();
|
||||
|
||||
// SECURITY: blocking rules must always take precedence
|
||||
// over allowing rules when prefixes are identical
|
||||
if let Some(existing) = prefix_rules
|
||||
.iter_mut()
|
||||
.find(|(_, existing_rule)| *existing_rule == rule)
|
||||
{
|
||||
if existing.0 == Action::Block {
|
||||
action = Action::Block;
|
||||
for &rule in allow {
|
||||
rules.insert(rule, Action::Allow, host_address);
|
||||
}
|
||||
*existing = (action, rule);
|
||||
} else {
|
||||
prefix_rules.push((action, rule));
|
||||
|
||||
// SECURITY: blocking rules must always take precedence
|
||||
// over allowing rules when the rules are identical.
|
||||
for &rule in block {
|
||||
rules.insert(rule, Action::Block, host_address);
|
||||
}
|
||||
|
||||
rules
|
||||
}
|
||||
|
||||
pub(crate) fn policy_decision(
|
||||
&self,
|
||||
address: Ipv4Address,
|
||||
direction: Direction,
|
||||
) -> Option<PolicyDecision> {
|
||||
match (self.select(address, direction)?, self.mode) {
|
||||
(Action::Block, _) => Some(PolicyDecision::Block),
|
||||
(Action::Allow, Mode::Legacy) => Some(PolicyDecision::AllowStateless),
|
||||
(Action::Allow, Mode::Stateful) => Some(PolicyDecision::AllowStateful),
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn is_stateful(&self, address: Ipv4Address, direction: Direction) -> bool {
|
||||
self.mode == Mode::Stateful && self.select(address, direction).is_some()
|
||||
}
|
||||
|
||||
pub(crate) fn len(&self) -> usize {
|
||||
self.inbound.len() + self.outbound.len()
|
||||
}
|
||||
|
||||
pub(crate) fn has_stateful(&self) -> bool {
|
||||
self.mode == Mode::Stateful
|
||||
}
|
||||
|
||||
fn select(&self, address: Ipv4Address, direction: Direction) -> Option<Action> {
|
||||
let entries = match direction {
|
||||
Direction::In => &self.inbound,
|
||||
Direction::Out => &self.outbound,
|
||||
};
|
||||
|
||||
entries
|
||||
.get_lpm(&Ipv4Net::from(address))
|
||||
.map(|(_, action)| *action)
|
||||
}
|
||||
|
||||
fn insert(&mut self, rule: Rule, action: Action, host_address: Ipv4Address) {
|
||||
match rule {
|
||||
Rule::Stateless(target) => {
|
||||
// Bare rules apply in both directions in stateful mode
|
||||
if self.mode == Mode::Stateful {
|
||||
self.insert_direction(Direction::In, target, action, host_address);
|
||||
}
|
||||
|
||||
self.insert_direction(Direction::Out, target, action, host_address);
|
||||
}
|
||||
Rule::Stateful { direction, target } => {
|
||||
self.insert_direction(direction, target, action, host_address);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn insert_direction(
|
||||
&mut self,
|
||||
direction: Direction,
|
||||
target: Target,
|
||||
action: Action,
|
||||
host_address: Ipv4Address,
|
||||
) {
|
||||
let prefix = match target {
|
||||
Target::Prefix(prefix) => prefix,
|
||||
Target::Host => host_address.into(),
|
||||
};
|
||||
let entries = match direction {
|
||||
Direction::In => &mut self.inbound,
|
||||
Direction::Out => &mut self.outbound,
|
||||
};
|
||||
|
||||
entries.insert(prefix, action);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{Action, has_stateful_rules, insert_rule, rule_count, select_rules};
|
||||
use crate::proxy::{Direction, Rule, Target};
|
||||
use ipnet::Ipv4Net;
|
||||
use prefix_trie::PrefixMap;
|
||||
use super::{Action, Mode, PolicyDecision, Rules};
|
||||
use crate::proxy::Direction;
|
||||
use smoltcp::wire::Ipv4Address;
|
||||
use std::str::FromStr;
|
||||
|
||||
const HOST: Ipv4Address = Ipv4Address::new(192, 168, 64, 1);
|
||||
|
||||
fn stateful_rules() -> Rules {
|
||||
Rules {
|
||||
mode: Mode::Stateful,
|
||||
..Rules::default()
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_policy_precedence() {
|
||||
let host = Ipv4Address::new(192, 168, 64, 1);
|
||||
let target = Ipv4Address::new(10, 0, 0, 1);
|
||||
let mut rules = PrefixMap::new();
|
||||
let mut rules = stateful_rules();
|
||||
|
||||
insert_rule(
|
||||
&mut rules,
|
||||
"0.0.0.0/0".parse().unwrap(),
|
||||
Action::Block,
|
||||
host,
|
||||
);
|
||||
assert!(!has_stateful_rules(&rules));
|
||||
|
||||
insert_rule(
|
||||
&mut rules,
|
||||
"in 10.0.0.0/8".parse().unwrap(),
|
||||
Action::Allow,
|
||||
host,
|
||||
);
|
||||
assert!(has_stateful_rules(&rules));
|
||||
rules.insert("0.0.0.0/0".parse().unwrap(), Action::Block, HOST);
|
||||
rules.insert("in 10.0.0.0/8".parse().unwrap(), Action::Allow, HOST);
|
||||
|
||||
assert_eq!(
|
||||
select_rules(&rules, target, Direction::In),
|
||||
Some(
|
||||
&[(
|
||||
Action::Allow,
|
||||
Rule::Stateful {
|
||||
direction: Direction::In,
|
||||
target: Target::Prefix(Ipv4Net::from_str("10.0.0.0/8").unwrap()),
|
||||
}
|
||||
)][..]
|
||||
)
|
||||
rules.policy_decision(target, Direction::In),
|
||||
Some(PolicyDecision::AllowStateful)
|
||||
);
|
||||
assert_eq!(
|
||||
select_rules(&rules, target, Direction::Out),
|
||||
Some(&[(Action::Block, "0.0.0.0/0".parse().unwrap())][..])
|
||||
rules.policy_decision(target, Direction::Out),
|
||||
Some(PolicyDecision::Block)
|
||||
);
|
||||
|
||||
insert_rule(
|
||||
&mut rules,
|
||||
"10.0.0.1/32".parse().unwrap(),
|
||||
Action::Allow,
|
||||
host,
|
||||
);
|
||||
rules.insert("10.0.0.1/32".parse().unwrap(), Action::Allow, HOST);
|
||||
assert_eq!(
|
||||
select_rules(&rules, target, Direction::Out).unwrap().len(),
|
||||
1
|
||||
);
|
||||
|
||||
insert_rule(
|
||||
&mut rules,
|
||||
"out 10.0.0.1/32".parse().unwrap(),
|
||||
Action::Allow,
|
||||
host,
|
||||
);
|
||||
assert_eq!(
|
||||
select_rules(&rules, target, Direction::Out).unwrap().len(),
|
||||
2
|
||||
rules.policy_decision(target, Direction::Out),
|
||||
Some(PolicyDecision::AllowStateful)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_directional_rules_share_a_prefix() {
|
||||
let host = Ipv4Address::new(192, 168, 64, 1);
|
||||
let mut rules = PrefixMap::new();
|
||||
fn test_directional_rules_at_same_prefix_are_independent() {
|
||||
let mut rules = stateful_rules();
|
||||
|
||||
for (target, action) in [
|
||||
("in @host", Action::Allow),
|
||||
("out @host", Action::Allow),
|
||||
("in @host", Action::Block),
|
||||
] {
|
||||
insert_rule(&mut rules, target.parse().unwrap(), action, host);
|
||||
rules.insert(target.parse().unwrap(), action, HOST);
|
||||
}
|
||||
|
||||
assert_eq!(
|
||||
select_rules(&rules, host, Direction::In),
|
||||
Some(
|
||||
&[
|
||||
(
|
||||
Action::Block,
|
||||
Rule::Stateful {
|
||||
direction: Direction::In,
|
||||
target: Target::Prefix(host.into()),
|
||||
}
|
||||
),
|
||||
(
|
||||
Action::Allow,
|
||||
Rule::Stateful {
|
||||
direction: Direction::Out,
|
||||
target: Target::Prefix(host.into()),
|
||||
}
|
||||
),
|
||||
][..]
|
||||
)
|
||||
rules.policy_decision(HOST, Direction::In),
|
||||
Some(PolicyDecision::Block)
|
||||
);
|
||||
assert_eq!(rule_count(&rules), 2);
|
||||
assert_eq!(
|
||||
rules.policy_decision(HOST, Direction::Out),
|
||||
Some(PolicyDecision::AllowStateful)
|
||||
);
|
||||
assert_eq!(rules.len(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_stateless_rules_are_outbound_only() {
|
||||
let target = Ipv4Address::new(10, 1, 2, 3);
|
||||
let mut rules = Rules::default();
|
||||
|
||||
rules.insert("10.0.0.0/8".parse().unwrap(), Action::Block, HOST);
|
||||
|
||||
assert!(rules.policy_decision(target, Direction::In).is_none());
|
||||
assert_eq!(
|
||||
rules.policy_decision(target, Direction::Out),
|
||||
Some(PolicyDecision::Block)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_inbound_selection_uses_more_specific_bare_rule() {
|
||||
let target = Ipv4Address::new(10, 1, 2, 3);
|
||||
let mut rules = stateful_rules();
|
||||
|
||||
rules.insert("10.1.0.0/16".parse().unwrap(), Action::Allow, HOST);
|
||||
rules.insert("in 10.0.0.0/8".parse().unwrap(), Action::Block, HOST);
|
||||
|
||||
assert_eq!(
|
||||
rules.policy_decision(target, Direction::In),
|
||||
Some(PolicyDecision::AllowStateful)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_block_wins_over_allow_at_same_outbound_prefix() {
|
||||
let target = Ipv4Address::new(10, 1, 2, 3);
|
||||
|
||||
for (allow, block) in [
|
||||
("10.0.0.0/8", "out 10.0.0.0/8"),
|
||||
("out 10.0.0.0/8", "10.0.0.0/8"),
|
||||
] {
|
||||
let mut rules = stateful_rules();
|
||||
rules.insert(allow.parse().unwrap(), Action::Allow, HOST);
|
||||
rules.insert(block.parse().unwrap(), Action::Block, HOST);
|
||||
|
||||
assert_eq!(
|
||||
rules.policy_decision(target, Direction::Out),
|
||||
Some(PolicyDecision::Block)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_directional_rule_makes_bare_rules_stateful() {
|
||||
let allow = "out @host".parse().unwrap();
|
||||
let block = "0.0.0.0/0".parse().unwrap();
|
||||
let rules = Rules::new(HOST, &[allow], &[block]);
|
||||
|
||||
assert_eq!(rules.len(), 3);
|
||||
assert!(rules.has_stateful());
|
||||
assert_eq!(
|
||||
rules.policy_decision(HOST, Direction::In),
|
||||
Some(PolicyDecision::Block)
|
||||
);
|
||||
assert_eq!(
|
||||
rules.policy_decision(HOST, Direction::Out),
|
||||
Some(PolicyDecision::AllowStateful)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_return_tracking_applies_to_all_rules_in_stateful_mode() {
|
||||
let stateful_target = Ipv4Address::new(10, 1, 2, 3);
|
||||
let stateless_target = Ipv4Address::new(192, 0, 2, 1);
|
||||
let mut rules = stateful_rules();
|
||||
|
||||
rules.insert("0.0.0.0/0".parse().unwrap(), Action::Block, HOST);
|
||||
rules.insert("out 10.0.0.0/8".parse().unwrap(), Action::Block, HOST);
|
||||
|
||||
assert!(rules.is_stateful(stateful_target, Direction::Out));
|
||||
assert!(rules.is_stateful(stateless_target, Direction::Out));
|
||||
assert!(rules.is_stateful(stateful_target, Direction::In));
|
||||
|
||||
rules.insert("10.0.0.0/8".parse().unwrap(), Action::Allow, HOST);
|
||||
assert!(rules.is_stateful(stateful_target, Direction::Out));
|
||||
}
|
||||
}
|
||||
|
||||
+158
-66
@@ -1,14 +1,19 @@
|
||||
use crate::dhcp_snooper::Lease;
|
||||
use crate::proxy::conntrack::ConntrackResult;
|
||||
use crate::dhcp_snooper::{Lease, message_matches_bootp_client};
|
||||
use crate::proxy::flows::{FlowDirection, FlowMatch};
|
||||
use crate::proxy::udp_packet_helper::UdpPacketHelper;
|
||||
use crate::proxy::{Action, Direction, Proxy, Rule, select_rules};
|
||||
use crate::proxy::{Direction, PolicyDecision, Proxy};
|
||||
use anyhow::Context;
|
||||
use anyhow::Result;
|
||||
use dhcproto::Decodable;
|
||||
use dhcproto::v4::Opcode;
|
||||
use smoltcp::phy::ChecksumCapabilities;
|
||||
use smoltcp::wire::{
|
||||
ArpOperation, ArpPacket, ArpRepr, EthernetFrame, EthernetProtocol, IpProtocol, Ipv4Address,
|
||||
Ipv4Packet, UdpPacket,
|
||||
Ipv4Packet, Ipv4Repr, UdpPacket,
|
||||
};
|
||||
|
||||
const IPV4_HEADER_LEN_WITHOUT_OPTIONS: u8 = 20;
|
||||
|
||||
impl Proxy<'_> {
|
||||
pub(crate) fn process_frame_from_vm(&mut self, frame: EthernetFrame<&[u8]>) -> Result<()> {
|
||||
if self.allowed_from_vm(&frame).is_none() {
|
||||
@@ -33,7 +38,14 @@ impl Proxy<'_> {
|
||||
self.allowed_from_vm_arp(arp_pkt)
|
||||
}
|
||||
EthernetProtocol::Ipv4 => {
|
||||
let ipv4_pkt = Ipv4Packet::new_checked(frame.payload()).ok()?;
|
||||
let ipv4_pkt = Ipv4Packet::new_unchecked(frame.payload());
|
||||
Ipv4Repr::parse(&ipv4_pkt, &ChecksumCapabilities::ignored()).ok()?;
|
||||
|
||||
// Reject IPv4 options because source routing could bypass destination-based policy
|
||||
if ipv4_pkt.header_len() != IPV4_HEADER_LEN_WITHOUT_OPTIONS {
|
||||
return None;
|
||||
}
|
||||
|
||||
self.allowed_from_vm_ipv4(ipv4_pkt)
|
||||
}
|
||||
_ => None,
|
||||
@@ -49,58 +61,60 @@ impl Proxy<'_> {
|
||||
if let Some(lease) = &self.dhcp_snooper.lease()
|
||||
&& lease.is_valid_for(ipv4_pkt.src_addr())
|
||||
{
|
||||
// Unicast DHCP renewal is required to maintain the VM's lease
|
||||
// and must bypass user-specified rules
|
||||
if is_allowed_dhcp_request(
|
||||
&ipv4_pkt,
|
||||
Some(self.host.gateway_ip),
|
||||
self.vm_mac_address,
|
||||
self.dhcp_snooper.lease(),
|
||||
) {
|
||||
return Some(());
|
||||
}
|
||||
|
||||
// Consult the flow table before evaluating outbound policy
|
||||
// so established flows are not treated as new traffic
|
||||
let pending = match self
|
||||
.flows
|
||||
.as_mut()
|
||||
.map(|flows| flows.inspect(&ipv4_pkt, FlowDirection::FromVm))
|
||||
.unwrap_or(FlowMatch::Untracked)
|
||||
{
|
||||
FlowMatch::Allowed => return Some(()),
|
||||
FlowMatch::Denied => return None,
|
||||
FlowMatch::Candidate(pending) => Some(pending),
|
||||
FlowMatch::Untracked => None,
|
||||
};
|
||||
|
||||
// The flow is either pending or untracked, evaluate it against outbound policy
|
||||
let dst_addr = ipv4_pkt.dst_addr();
|
||||
|
||||
// Filter traffic based on user-specified rules first
|
||||
if let Some(rules) = select_rules(&self.rules, dst_addr, Direction::Out) {
|
||||
// DHCP is required to maintain the VM's lease and must bypass user-specified rules
|
||||
if is_allowed_dhcp_request(&ipv4_pkt, Some(self.host.gateway_ip)) {
|
||||
return Some(());
|
||||
match self.rules.policy_decision(dst_addr, Direction::Out) {
|
||||
// Return traffic was handled above; enforce explicit outbound blocks here
|
||||
Some(PolicyDecision::Block) => return None,
|
||||
|
||||
// Track statelessly allowed traffic only when needed so its reply is not
|
||||
// treated as a new inbound flow
|
||||
Some(PolicyDecision::AllowStateless) => {
|
||||
return self.admit_with_tracking_if_stateful(pending, dst_addr, Direction::In);
|
||||
}
|
||||
|
||||
if let Some((action, _)) = rules
|
||||
.iter()
|
||||
.find(|(_, rule)| matches!(rule, Rule::Stateless(_)))
|
||||
{
|
||||
return match action {
|
||||
Action::Allow => Some(()),
|
||||
Action::Block => None,
|
||||
};
|
||||
}
|
||||
// Untracked packets cannot satisfy stateful policy
|
||||
Some(PolicyDecision::AllowStateful) => return self.admit_with_tracking(pending?),
|
||||
|
||||
return match self.conntrack.inspect_from_vm(&ipv4_pkt) {
|
||||
ConntrackResult::Allowed => Some(()),
|
||||
ConntrackResult::Denied => None,
|
||||
ConntrackResult::New(pending) => {
|
||||
let allow_new = rules.iter().any(|(action, rule)| {
|
||||
*action == Action::Allow
|
||||
&& matches!(
|
||||
rule,
|
||||
Rule::Stateful {
|
||||
direction: Direction::Out,
|
||||
..
|
||||
}
|
||||
)
|
||||
});
|
||||
|
||||
if !allow_new {
|
||||
return None;
|
||||
}
|
||||
|
||||
self.conntrack.commit(pending).then_some(())
|
||||
}
|
||||
};
|
||||
// No outbound rule matched; apply the built-in fallbacks below
|
||||
None => {}
|
||||
}
|
||||
|
||||
// When no user-specified rules matched, simply allow all global traffic
|
||||
if ip_network::IpNetwork::from(dst_addr).is_global() {
|
||||
return Some(());
|
||||
return self.admit_with_tracking_if_trackable(pending);
|
||||
}
|
||||
|
||||
// Additionally, allow communication with the host,
|
||||
// otherwise things like SSH to a VM won't work
|
||||
if ipv4_pkt.dst_addr() == self.host.gateway_ip {
|
||||
return Some(());
|
||||
if dst_addr == self.host.gateway_ip {
|
||||
return self.admit_with_tracking_if_trackable(pending);
|
||||
}
|
||||
|
||||
// Additionally, allow DNS requests to DNS-servers
|
||||
@@ -108,17 +122,20 @@ impl Proxy<'_> {
|
||||
if ipv4_pkt.next_header() == IpProtocol::Udp {
|
||||
let udp_pkt = UdpPacket::new_checked(ipv4_pkt.payload()).ok()?;
|
||||
|
||||
if udp_pkt.is_dns_request()
|
||||
&& self.dhcp_snooper.valid_dns_target(&ipv4_pkt.dst_addr())
|
||||
{
|
||||
return Some(());
|
||||
if udp_pkt.is_dns_request() && self.dhcp_snooper.valid_dns_target(&dst_addr) {
|
||||
return self.admit_with_tracking_if_trackable(pending);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Allow outgoing DHCP requests to the bootpd(8) broadcast address,
|
||||
// otherwise DHCP snooper will never be populated
|
||||
if is_allowed_dhcp_request(&ipv4_pkt, None) {
|
||||
if is_allowed_dhcp_request(
|
||||
&ipv4_pkt,
|
||||
None,
|
||||
self.vm_mac_address,
|
||||
self.dhcp_snooper.lease(),
|
||||
) {
|
||||
return Some(());
|
||||
}
|
||||
|
||||
@@ -129,7 +146,20 @@ impl Proxy<'_> {
|
||||
fn is_allowed_dhcp_request(
|
||||
ipv4_pkt: &Ipv4Packet<&[u8]>,
|
||||
unicast_target: Option<Ipv4Address>,
|
||||
vm_mac_address: smoltcp::wire::EthernetAddress,
|
||||
lease: &Option<Lease>,
|
||||
) -> bool {
|
||||
// Require the source address to be either:
|
||||
// * covered by the VM's current lease
|
||||
// * unspecified on the broadcast DHCP path
|
||||
let src_addr = ipv4_pkt.src_addr();
|
||||
let src_has_valid_lease = lease
|
||||
.as_ref()
|
||||
.is_some_and(|lease| lease.is_valid_for(src_addr));
|
||||
if !src_has_valid_lease && !(unicast_target.is_none() && src_addr.is_unspecified()) {
|
||||
return false;
|
||||
}
|
||||
|
||||
let dst_addr = ipv4_pkt.dst_addr();
|
||||
|
||||
// Keep the common path cheap and inspect UDP only for a permitted DHCP target
|
||||
@@ -145,7 +175,18 @@ fn is_allowed_dhcp_request(
|
||||
return false;
|
||||
};
|
||||
|
||||
udp_pkt.is_dhcp_request()
|
||||
// Require the standard DHCP client and server ports
|
||||
if !udp_pkt.is_dhcp_request() {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Require the BOOTP client hardware address to match this VM
|
||||
let mut decoder = dhcproto::v4::Decoder::new(udp_pkt.payload());
|
||||
let Ok(message) = dhcproto::v4::Message::decode(&mut decoder) else {
|
||||
return false;
|
||||
};
|
||||
|
||||
message_matches_bootp_client(&message, Opcode::BootRequest, vm_mac_address.0)
|
||||
}
|
||||
|
||||
fn vm_arp_allowed(
|
||||
@@ -186,6 +227,8 @@ fn vm_arp_allowed(
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use crate::dhcp_snooper::Lease;
|
||||
use dhcproto::v4::{DhcpOption, Message, MessageType};
|
||||
use dhcproto::{Encodable, Encoder};
|
||||
use smoltcp::wire::{
|
||||
ArpHardware, ArpOperation, ArpPacket, EthernetAddress, EthernetProtocol, IpProtocol,
|
||||
Ipv4Address, Ipv4Packet, UdpPacket,
|
||||
@@ -193,15 +236,31 @@ mod tests {
|
||||
use std::collections::HashSet;
|
||||
use std::time::Duration;
|
||||
|
||||
#[test]
|
||||
fn test_allowed_dhcp_request_targets() {
|
||||
let gateway = Ipv4Address::new(192, 168, 64, 1);
|
||||
let other = Ipv4Address::new(192, 168, 64, 2);
|
||||
const VM_MAC: EthernetAddress = EthernetAddress([0x02, 0x00, 0x00, 0x00, 0x00, 0x01]);
|
||||
|
||||
assert!(allowed_dhcp_request(Ipv4Address::BROADCAST, None));
|
||||
assert!(allowed_dhcp_request(gateway, Some(gateway)));
|
||||
assert!(!allowed_dhcp_request(gateway, None));
|
||||
assert!(!allowed_dhcp_request(other, Some(gateway)));
|
||||
#[test]
|
||||
fn test_allowed_dhcp_request_policy() {
|
||||
let gateway = Ipv4Address::new(192, 168, 64, 1);
|
||||
let lease_ip = Ipv4Address::new(192, 168, 64, 2);
|
||||
let other = Ipv4Address::new(192, 168, 64, 3);
|
||||
let no_lease = None;
|
||||
let lease = Some(Lease::new(
|
||||
lease_ip,
|
||||
Duration::from_secs(600),
|
||||
HashSet::new(),
|
||||
));
|
||||
let initial = |src, chaddr| {
|
||||
allowed_dhcp_request(src, Ipv4Address::BROADCAST, None, chaddr, &no_lease)
|
||||
};
|
||||
let renewal = |src, dst| allowed_dhcp_request(src, dst, Some(gateway), VM_MAC.0, &lease);
|
||||
let other_mac = [0x02, 0x00, 0x00, 0x00, 0x00, 0x02];
|
||||
|
||||
assert!(initial(Ipv4Address::UNSPECIFIED, VM_MAC.0));
|
||||
assert!(renewal(lease_ip, gateway));
|
||||
assert!(!renewal(other, gateway));
|
||||
assert!(!renewal(Ipv4Address::UNSPECIFIED, gateway));
|
||||
assert!(!renewal(lease_ip, other));
|
||||
assert!(!initial(Ipv4Address::UNSPECIFIED, other_mac));
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -303,21 +362,54 @@ mod tests {
|
||||
buf
|
||||
}
|
||||
|
||||
fn allowed_dhcp_request(dst_addr: Ipv4Address, unicast_target: Option<Ipv4Address>) -> bool {
|
||||
let mut buf = vec![0; 28];
|
||||
fn allowed_dhcp_request(
|
||||
src_addr: Ipv4Address,
|
||||
dst_addr: Ipv4Address,
|
||||
unicast_target: Option<Ipv4Address>,
|
||||
chaddr: [u8; 6],
|
||||
lease: &Option<Lease>,
|
||||
) -> bool {
|
||||
let mut buf = dhcp_request(chaddr);
|
||||
let mut ipv4_pkt = Ipv4Packet::new_unchecked(buf.as_mut_slice());
|
||||
ipv4_pkt.set_src_addr(src_addr);
|
||||
ipv4_pkt.set_dst_addr(dst_addr);
|
||||
|
||||
let ipv4_pkt = Ipv4Packet::new_checked(buf.as_slice()).unwrap();
|
||||
super::is_allowed_dhcp_request(&ipv4_pkt, unicast_target, VM_MAC, lease)
|
||||
}
|
||||
|
||||
fn dhcp_request(chaddr: [u8; 6]) -> Vec<u8> {
|
||||
let mut message = Message::new(
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
Ipv4Address::UNSPECIFIED,
|
||||
&chaddr,
|
||||
);
|
||||
message
|
||||
.opts_mut()
|
||||
.insert(DhcpOption::MessageType(MessageType::Discover));
|
||||
|
||||
let mut dhcp_payload = Vec::new();
|
||||
message
|
||||
.encode(&mut Encoder::new(&mut dhcp_payload))
|
||||
.unwrap();
|
||||
|
||||
let total_len = 20 + 8 + dhcp_payload.len();
|
||||
let mut buf = vec![0; total_len];
|
||||
let mut ipv4_pkt = Ipv4Packet::new_unchecked(buf.as_mut_slice());
|
||||
ipv4_pkt.set_version(4);
|
||||
ipv4_pkt.set_header_len(20);
|
||||
ipv4_pkt.set_total_len(28);
|
||||
ipv4_pkt.set_total_len(total_len as u16);
|
||||
ipv4_pkt.set_next_header(IpProtocol::Udp);
|
||||
ipv4_pkt.set_dst_addr(dst_addr);
|
||||
ipv4_pkt.set_src_addr(Ipv4Address::UNSPECIFIED);
|
||||
ipv4_pkt.set_dst_addr(Ipv4Address::BROADCAST);
|
||||
|
||||
let mut udp_pkt = UdpPacket::new_unchecked(ipv4_pkt.payload_mut());
|
||||
udp_pkt.set_src_port(68);
|
||||
udp_pkt.set_dst_port(67);
|
||||
udp_pkt.set_len(8);
|
||||
|
||||
let ipv4_pkt = Ipv4Packet::new_checked(buf.as_slice()).unwrap();
|
||||
super::is_allowed_dhcp_request(&ipv4_pkt, unicast_target)
|
||||
udp_pkt.set_len((8 + dhcp_payload.len()) as u16);
|
||||
udp_pkt.payload_mut().copy_from_slice(&dhcp_payload);
|
||||
buf
|
||||
}
|
||||
}
|
||||
|
||||
@@ -26,6 +26,14 @@ impl VM {
|
||||
pub fn read(&self, buf: &mut [u8]) -> std::io::Result<usize> {
|
||||
self.sock.recv(buf)
|
||||
}
|
||||
|
||||
pub fn is_connected(&self) -> io::Result<bool> {
|
||||
match self.sock.peer_addr() {
|
||||
Ok(_) => Ok(true),
|
||||
Err(error) if error.kind() == io::ErrorKind::NotConnected => Ok(false),
|
||||
Err(error) => Err(error),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn duplicate_vm_fd(vm_fd: RawFd) -> Result<RawFd> {
|
||||
@@ -115,10 +123,12 @@ impl AsRawFd for VM {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::VM;
|
||||
use polling::{Event, Events, PollMode, Poller};
|
||||
use std::fs::File;
|
||||
use std::net::UdpSocket;
|
||||
use std::os::fd::AsRawFd;
|
||||
use std::os::unix::net::{UnixDatagram, UnixStream};
|
||||
use std::time::Duration;
|
||||
|
||||
#[test]
|
||||
fn test_new_rejects_negative_fd() {
|
||||
@@ -187,4 +197,62 @@ mod tests {
|
||||
|
||||
assert!(socket_fd_is_open);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_connected_socket_has_peer() {
|
||||
let (socket, _peer) = UnixDatagram::pair().unwrap();
|
||||
let vm = VM::new(socket.as_raw_fd()).unwrap();
|
||||
|
||||
assert!(vm.is_connected().unwrap());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_disconnected_peer_is_detected_without_kqueue_event() {
|
||||
let (socket, peer) = UnixDatagram::pair().unwrap();
|
||||
let vm = VM::new(socket.as_raw_fd()).unwrap();
|
||||
let poller = Poller::new().unwrap();
|
||||
let mut events = Events::new();
|
||||
|
||||
unsafe {
|
||||
poller
|
||||
.add_with_mode(vm.as_raw_fd(), Event::readable(0), PollMode::Edge)
|
||||
.unwrap();
|
||||
}
|
||||
drop(peer);
|
||||
|
||||
poller
|
||||
.wait(&mut events, Some(Duration::from_millis(20)))
|
||||
.unwrap();
|
||||
|
||||
assert!(events.is_empty());
|
||||
assert!(!vm.is_connected().unwrap());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_disconnected_peer_is_detected_when_another_socket_wakes_kqueue() {
|
||||
let (socket, peer) = UnixDatagram::pair().unwrap();
|
||||
let (host, host_peer) = UnixDatagram::pair().unwrap();
|
||||
let vm = VM::new(socket.as_raw_fd()).unwrap();
|
||||
let poller = Poller::new().unwrap();
|
||||
let mut events = Events::new();
|
||||
|
||||
unsafe {
|
||||
poller
|
||||
.add_with_mode(vm.as_raw_fd(), Event::readable(0), PollMode::Edge)
|
||||
.unwrap();
|
||||
poller
|
||||
.add_with_mode(host.as_raw_fd(), Event::readable(1), PollMode::Edge)
|
||||
.unwrap();
|
||||
}
|
||||
drop(peer);
|
||||
host_peer.send(&[1]).unwrap();
|
||||
|
||||
poller
|
||||
.wait(&mut events, Some(Duration::from_millis(20)))
|
||||
.unwrap();
|
||||
|
||||
assert!(events.iter().any(|event| event.key == 1));
|
||||
assert!(!events.iter().any(|event| event.key == 0));
|
||||
assert!(!vm.is_connected().unwrap());
|
||||
}
|
||||
}
|
||||
|
||||
+63
-27
@@ -8,18 +8,21 @@ use softnet::NetType;
|
||||
use softnet::proxy::ExposedPort;
|
||||
use softnet::proxy::Proxy;
|
||||
use softnet::proxy::Rule;
|
||||
use std::borrow::Cow;
|
||||
use std::env;
|
||||
use std::os::raw::c_int;
|
||||
use std::os::unix::io::RawFd;
|
||||
use std::os::unix::process::CommandExt;
|
||||
use std::process::{Command, ExitCode};
|
||||
use system_configuration::core_foundation::base::TCFType;
|
||||
use system_configuration::core_foundation::boolean::CFBoolean;
|
||||
use system_configuration::core_foundation::dictionary::CFDictionary;
|
||||
use system_configuration::core_foundation::number::CFNumber;
|
||||
use system_configuration::core_foundation::string::CFString;
|
||||
use system_configuration::preferences::SCPreferences;
|
||||
use system_configuration::sys::preferences::{SCPreferencesCommitChanges, SCPreferencesSetValue};
|
||||
use system_configuration::sys::preferences::{
|
||||
SCPreferencesApplyChanges, SCPreferencesCommitChanges, SCPreferencesLock,
|
||||
SCPreferencesSetValue, SCPreferencesUnlock,
|
||||
};
|
||||
use uzers::{get_current_groupname, get_current_username, get_effective_uid};
|
||||
|
||||
#[derive(Parser, Debug)]
|
||||
@@ -59,14 +62,14 @@ struct Args {
|
||||
|
||||
#[clap(
|
||||
long,
|
||||
help = "Comma-separated list of rules for allowing traffic.\n\n\
|
||||
Rule forms:\n\n\
|
||||
* TARGET: traffic between TARGET and the VM in either direction\n\
|
||||
help = "Comma-separated rules for allowing traffic, in the following forms:\n\n\
|
||||
* TARGET: traffic sent from the VM to TARGET; reverse traffic is not filtered by this rule\n\
|
||||
* in TARGET: flows initiated from TARGET to the VM\n\
|
||||
* out TARGET: flows initiated from the VM to TARGET\n\n\
|
||||
Targets are:\n\n\
|
||||
* IPv4 CIDRs\n\
|
||||
* @host, which matches the vmnet bridge gateway IP\n\n\
|
||||
Directional rules make bare TARGET rules stateful in both directions.\n\n\
|
||||
When used with --block, the longest prefix match wins. If an identical rule is both \
|
||||
allowed and blocked, blocking takes precedence.\n\n\
|
||||
--allow=0.0.0.0/0 additionally disables bridge isolation, even when \
|
||||
@@ -84,21 +87,21 @@ struct Args {
|
||||
|
||||
#[clap(
|
||||
long,
|
||||
help = "Comma-separated list of rules for blocking traffic.\n\n\
|
||||
Rule forms:\n\n\
|
||||
* TARGET: traffic between TARGET and the VM in either direction\n\
|
||||
help = "Comma-separated rules for blocking traffic, in the following forms:\n\n\
|
||||
* TARGET: traffic sent from the VM to TARGET; reverse traffic is not filtered by this rule\n\
|
||||
* in TARGET: flows initiated from TARGET to the VM\n\
|
||||
* out TARGET: flows initiated from the VM to TARGET\n\n\
|
||||
Targets are:\n\n\
|
||||
* IPv4 CIDRs\n\
|
||||
* @host, which matches the vmnet bridge gateway IP\n\n\
|
||||
Directional rules make bare TARGET rules stateful in both directions.\n\n\
|
||||
When used with --allow, the longest prefix match wins. If an identical rule is both \
|
||||
allowed and blocked, blocking takes precedence.\n\n\
|
||||
Examples:\n\n\
|
||||
* --block=0.0.0.0/0 — establish a stateless default-deny policy\n\
|
||||
* --block=\"in @host\" — block stateful flows initiated from @host\n\
|
||||
* --block=0.0.0.0/0 — establish a stateless default-deny egress policy\n\
|
||||
* --block=\"out @host\" — block stateful flows initiated toward @host\n\
|
||||
* --block=\"out 66.66.66.0/24\" — block stateful flows initiated toward this CIDR\n\
|
||||
* --block=\"in @host,out 66.66.66.0/24\" — multiple rules may be comma-separated",
|
||||
* --block=\"out @host,out 66.66.66.0/24\" — multiple rules may be comma-separated",
|
||||
value_name = "comma-separated rules",
|
||||
use_value_delimiter = true,
|
||||
action = clap::ArgAction::Set
|
||||
@@ -130,10 +133,10 @@ fn main() -> ExitCode {
|
||||
}
|
||||
|
||||
// Initialize Sentry
|
||||
let _sentry = sentry::init(sentry::ClientOptions {
|
||||
release: option_env!("CIRRUS_TAG").map(|tag| Cow::from(format!("softnet@{tag}"))),
|
||||
..Default::default()
|
||||
});
|
||||
let _sentry = sentry::init(
|
||||
sentry::ClientOptions::default()
|
||||
.maybe_release(option_env!("CIRRUS_TAG").map(|tag| format!("softnet@{tag}"))),
|
||||
);
|
||||
|
||||
// Enrich future events with Cirrus CI-specific tags
|
||||
if let Ok(tags) = env::var("CIRRUS_SENTRY_TAGS") {
|
||||
@@ -216,8 +219,8 @@ fn try_main() -> anyhow::Result<()> {
|
||||
));
|
||||
}
|
||||
|
||||
// Set bootpd(8) min/max lease time while still having the root privileges
|
||||
set_bootpd_lease_time(args.bootpd_lease_time);
|
||||
// Configure bootpd(8) while still having the root privileges
|
||||
configure_bootpd(args.bootpd_lease_time)?;
|
||||
|
||||
// Initialize the proxy while still having the root privileges
|
||||
let mut proxy = Proxy::new(
|
||||
@@ -269,26 +272,59 @@ fn sudo_escalation_works() -> bool {
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
fn set_bootpd_lease_time(lease_time: u32) {
|
||||
fn configure_bootpd(lease_time: u32) -> anyhow::Result<()> {
|
||||
let prefs = SCPreferences::group(
|
||||
&CFString::new("softnet"),
|
||||
&CFString::new("com.apple.InternetSharing.default.plist"),
|
||||
);
|
||||
|
||||
let bootpd_dict = CFDictionary::from_CFType_pairs(&[(
|
||||
CFString::new("DHCPLeaseTimeSecs"),
|
||||
CFNumber::from(lease_time as i32),
|
||||
)]);
|
||||
let bootpd_dict = CFDictionary::from_CFType_pairs(&[
|
||||
(
|
||||
CFString::new("DHCPLeaseTimeSecs"),
|
||||
CFNumber::from(lease_time as i32).as_CFType(),
|
||||
),
|
||||
(
|
||||
CFString::new("dhcp_ignore_client_identifier"),
|
||||
CFBoolean::true_value().as_CFType(),
|
||||
),
|
||||
]);
|
||||
|
||||
unsafe {
|
||||
SCPreferencesSetValue(
|
||||
prefs.as_concrete_TypeRef(),
|
||||
CFString::new("bootpd").as_concrete_TypeRef(),
|
||||
bootpd_dict.as_concrete_TypeRef().cast(),
|
||||
let prefs = prefs.as_concrete_TypeRef();
|
||||
anyhow::ensure!(
|
||||
SCPreferencesLock(prefs, 1) != 0,
|
||||
"failed to lock bootpd preferences"
|
||||
);
|
||||
|
||||
SCPreferencesCommitChanges(prefs.as_concrete_TypeRef());
|
||||
let result = (|| -> anyhow::Result<()> {
|
||||
anyhow::ensure!(
|
||||
SCPreferencesSetValue(
|
||||
prefs,
|
||||
CFString::new("bootpd").as_concrete_TypeRef(),
|
||||
bootpd_dict.as_concrete_TypeRef().cast(),
|
||||
) != 0,
|
||||
"failed to set bootpd preferences"
|
||||
);
|
||||
|
||||
anyhow::ensure!(
|
||||
SCPreferencesCommitChanges(prefs) != 0,
|
||||
"failed to commit bootpd preferences"
|
||||
);
|
||||
|
||||
anyhow::ensure!(
|
||||
SCPreferencesApplyChanges(prefs) != 0,
|
||||
"failed to apply bootpd preferences"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
})();
|
||||
|
||||
let unlocked = SCPreferencesUnlock(prefs) != 0;
|
||||
result?;
|
||||
anyhow::ensure!(unlocked, "failed to unlock bootpd preferences");
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
Reference in New Issue
Block a user