diff --git a/CI.md b/CI.md index 538c33cc04..1eb087f2a6 100644 --- a/CI.md +++ b/CI.md @@ -50,8 +50,8 @@ Manually dispatch `Integration Tests` on the candidate branch with an [{"environment":"ubuntu-docker-rootful","installer":"binaries","testsuite":"policy-advisor"}] ``` -This runs the `mechanistic-proposal` and `policy-local` conformance tests in the -installed-artifact suite. The artifact run must contain the candidate CLI and +This runs the `mechanistic-proposal`, `new-hostname-proposal`, and +`policy-local` conformance tests in the installed-artifact suite. The artifact run must contain the candidate CLI and gateway binaries and runtime images. This manual run does not replace the required PR E2E gate. diff --git a/architecture/sandbox.md b/architecture/sandbox.md index b9a2b88603..b1169ca64e 100644 --- a/architecture/sandbox.md +++ b/architecture/sandbox.md @@ -269,8 +269,9 @@ Input mediation, DNS/TCP authorization, and outer-fence enforcement are identical in both modes. The selected mode is emitted in the sandbox qualification output (`seccomp_listener_mode`). -DNS uses an exact sandbox-local resolver at `127.0.0.53:53`. The driver sets the -nameserver and permits an unprivileged bind to port 53. UDP and TCP DNS requests +DNS uses an exact sandbox-local resolver at `127.0.0.53:53`. The driver sets that +nameserver without search domains, so names reach policy DNS as the workload +wrote them, and permits an unprivileged bind to port 53. UDP and TCP DNS requests are forwarded through the supervisor, which applies hostname-based DNS policy. The Podman driver supplies that resolver configuration as a driver-owned, read-only secret mounted at `/etc/resolv.conf`; the workload remains on @@ -323,10 +324,12 @@ preserves this lifetime rule; persistent responses remain eligible for reuse. An explicit `protocol: tcp` endpoint with a valid DNS hostname opts into native DNS and transparent TCP when the selected runtime advertises that substrate. Hostless `allowed_ips` and literal-IP selectors remain available only to the -legacy explicit-proxy path when `protocol` is omitted. The shared supervisor -answers only eligible DNS names, returns an epoch-scoped synthetic address, and +legacy explicit-proxy path when `protocol` is omitted. For an eligible DNS +name, the shared supervisor returns an epoch-scoped synthetic address and publishes the expiring name, endpoint, ports, policy generation, and validated -real addresses as one correlation. A connection to that synthetic address is +real addresses as one correlation. A name absent from policy receives at most a +contract-free observation address for policy advisor proposals; see +[Security Policy](security-policy.md). A connection to that synthetic address is captured before the bypass fence, mapped back to its workload process, authorized through the same egress pipeline, and dialed only through the pinned addresses. Omitted protocol endpoints retain explicit-proxy behavior. diff --git a/architecture/security-policy.md b/architecture/security-policy.md index 98946df833..07dc57cdda 100644 --- a/architecture/security-policy.md +++ b/architecture/security-policy.md @@ -282,6 +282,19 @@ because it changes the effective access model for every sandbox on the gateway. The policy advisor pipeline turns observed denials into draft policy recommendations. There are two proposers (sandbox-side mechanistic mapper, agent-authored via `policy.local`); the gateway is the single referee. +For a new DNS name absent from policy, the supervisor can +publish a bounded, short-lived synthetic observation mapping without querying +an upstream resolver. It carries only the name to the subsequent TCP mediation +step, which supplies the destination port and verified process identity for a +denial summary. Observation mappings have no endpoint contracts or real pinned +addresses, cannot authorize a relay, and become stale on policy generation +change. Transparent TCP re-acquires the mapping at the generation that +authorized the connection, so a reload between DNS and authorization fails +closed instead of resolving the name again. Every unknown name still produces a +DNS denial event. Reserved names, fail-closed quarantine, and an exhausted +observation budget (a quarter of each address family's synthetic pool) keep the +plain DNS refusal. This mechanistic observation path does not depend on the +agent-authored proposal setting. When enabled, L7 `policy_denied` responses include both structured `next_steps` and a short `agent_guidance` string so generic agents can continue through the proposal loop instead of treating the denial as terminal. diff --git a/crates/openshell-conformance/src/lib.rs b/crates/openshell-conformance/src/lib.rs index 189e41c4fe..b0e4295651 100644 --- a/crates/openshell-conformance/src/lib.rs +++ b/crates/openshell-conformance/src/lib.rs @@ -24,8 +24,8 @@ use tokio::time::sleep; use self::executor::{CliExecutionError, CliExecutor, ProcessCli}; pub use scenarios::{ - MECHANISTIC_PROPOSAL_SCENARIO, POLICY_LOCAL_SCENARIO, SANDBOX_LIFECYCLE_SCENARIO, - SMOKE_SCENARIO, + MECHANISTIC_PROPOSAL_SCENARIO, NEW_HOSTNAME_PROPOSAL_SCENARIO, POLICY_LOCAL_SCENARIO, + SANDBOX_LIFECYCLE_SCENARIO, SMOKE_SCENARIO, }; /// An installed conformance scenario. @@ -48,6 +48,7 @@ const SCENARIOS: &[Scenario] = &[ SMOKE_SCENARIO, SANDBOX_LIFECYCLE_SCENARIO, MECHANISTIC_PROPOSAL_SCENARIO, + NEW_HOSTNAME_PROPOSAL_SCENARIO, POLICY_LOCAL_SCENARIO, ]; diff --git a/crates/openshell-conformance/src/scenarios/mod.rs b/crates/openshell-conformance/src/scenarios/mod.rs index 3802d49de5..af6280fe1f 100644 --- a/crates/openshell-conformance/src/scenarios/mod.rs +++ b/crates/openshell-conformance/src/scenarios/mod.rs @@ -7,6 +7,8 @@ mod policy_behavior; mod sandbox_lifecycle; mod smoke; -pub use policy_behavior::{MECHANISTIC_PROPOSAL_SCENARIO, POLICY_LOCAL_SCENARIO}; +pub use policy_behavior::{ + MECHANISTIC_PROPOSAL_SCENARIO, NEW_HOSTNAME_PROPOSAL_SCENARIO, POLICY_LOCAL_SCENARIO, +}; pub use sandbox_lifecycle::SANDBOX_LIFECYCLE_SCENARIO; pub use smoke::SMOKE_SCENARIO; diff --git a/crates/openshell-conformance/src/scenarios/policy_behavior.rs b/crates/openshell-conformance/src/scenarios/policy_behavior.rs index 13cb9574a5..e26b89e45c 100644 --- a/crates/openshell-conformance/src/scenarios/policy_behavior.rs +++ b/crates/openshell-conformance/src/scenarios/policy_behavior.rs @@ -37,6 +37,21 @@ pub const MECHANISTIC_PROPOSAL_SCENARIO: Scenario = Scenario { run: run_mechanistic_proposal, }; +pub const NEW_HOSTNAME_PROPOSAL_SCENARIO: Scenario = Scenario { + name: "new-hostname-proposal", + description: "Turn a denied TCP open to a hostname absent from policy into a scoped draft.", + run: run_new_hostname_proposal, +}; + +const EMPTY_NETWORK_POLICY: &[u8] = br"version: 1 +filesystem_policy: + include_workdir: true + read_only: [/usr, /bin, /lib, /lib64, /proc, /dev/urandom, /app, /etc, /var/log] + read_write: [/sandbox, /tmp, /dev/null] +landlock: { compatibility: best_effort } +network_policies: {} +"; + fn run_policy_local(runner: &mut OpenShellRunner) -> ScenarioFuture<'_> { Box::pin(async move { let name = format!("ct-{}-pl", runner.id()); @@ -279,16 +294,7 @@ fn run_mechanistic_proposal(runner: &mut OpenShellRunner) -> ScenarioFuture<'_> Box::pin(async move { let mut policy = NamedTempFile::new().map_err(|error| error.to_string())?; policy - .write_all( - br"version: 1 -filesystem_policy: - include_workdir: true - read_only: [/usr, /bin, /lib, /lib64, /proc, /dev/urandom, /app, /etc, /var/log] - read_write: [/sandbox, /tmp, /dev/null] -landlock: { compatibility: best_effort } -network_policies: {} -", - ) + .write_all(EMPTY_NETWORK_POLICY) .map_err(|error| error.to_string())?; let policy_path = policy .path() @@ -345,60 +351,139 @@ network_policies: {} return Err(probe.failure_diagnostic("TCP open is denied before any upstream dial")); } - let started = Instant::now(); - loop { - let draft = runner - .step("mechanistic-draft") - .description("a single scoped mechanistic draft appears") - .with_timeout(COMMAND_TIMEOUT) - .run(&["rule", "get", &name]) - .await - .map_err(|error| error.to_string())?; - if !draft.success() { - if started.elapsed() >= PROPOSAL_TIMEOUT { - return Err(draft.failure_diagnostic("the reviewer inbox is readable")); - } - sleep(POLL_INTERVAL).await; - continue; + await_mechanistic_draft( + runner, + &name, + &ExpectedDraft { + rule: "allow_1_1_1_1_443", + endpoint: "1.1.1.1:443", + binary: &binary, + }, + probe.stderr(), + ) + .await + }) +} + +fn run_new_hostname_proposal(runner: &mut OpenShellRunner) -> ScenarioFuture<'_> { + Box::pin(async move { + let mut policy = NamedTempFile::new().map_err(|error| error.to_string())?; + policy + .write_all(EMPTY_NETWORK_POLICY) + .map_err(|error| error.to_string())?; + let policy_path = policy + .path() + .to_str() + .ok_or("temporary policy path is not UTF-8")?; + let name = format!("ct-{}-nh", runner.id()); + create_sandbox(runner, &name, Some(policy_path)).await?; + let binary = sandbox_bash_path(runner, &name).await?; + + let probe = runner + .step("denied-new-hostname") + .description("Bash cannot connect to pypi.org:80 before approval") + .with_timeout(COMMAND_TIMEOUT) + .run(&[ + "sandbox", + "exec", + "--name", + &name, + "--no-tty", + "--", + "bash", + "-c", + "if exec 3<>/dev/tcp/pypi.org/80; then echo UNEXPECTED_ALLOWED; exit 1; else echo DENIED; fi", + ]) + .await + .map_err(|error| error.to_string())?; + probe.require_success()?; + if !probe.stdout().lines().any(|line| line == "DENIED") { + return Err(probe.failure_diagnostic("new hostname stays denied before approval")); + } + + await_mechanistic_draft( + runner, + &name, + &ExpectedDraft { + rule: "allow_pypi_org_80", + endpoint: "pypi.org:80", + binary: &binary, + }, + probe.stderr(), + ) + .await + }) +} + +/// The single L4 mechanistic draft expected for one denied endpoint. +struct ExpectedDraft<'a> { + rule: &'a str, + endpoint: &'a str, + binary: &'a str, +} + +/// Poll the reviewer inbox until a draft appears, tolerating transient read +/// failures, then require that it is the single expected draft. +async fn await_mechanistic_draft( + runner: &OpenShellRunner, + sandbox: &str, + expected: &ExpectedDraft<'_>, + probe_stderr: &str, +) -> Result<(), String> { + let started = Instant::now(); + loop { + let draft = runner + .step("mechanistic-draft") + .description("a single scoped mechanistic draft appears") + .with_timeout(COMMAND_TIMEOUT) + .run(&["rule", "get", sandbox]) + .await + .map_err(|error| error.to_string())?; + if !draft.success() { + if started.elapsed() >= PROPOSAL_TIMEOUT { + return Err(draft.failure_diagnostic("the reviewer inbox is readable")); } - if !draft.stdout().contains("Chunk:") { - if started.elapsed() >= PROPOSAL_TIMEOUT { - return Err(draft.failure_diagnostic(&format!( - "one mechanistic draft for 1.1.1.1:443 and {binary}; probe stderr:\n{}", - probe.stderr() - ))); - } - sleep(POLL_INTERVAL).await; - continue; + sleep(POLL_INTERVAL).await; + continue; + } + if !draft.stdout().contains("Chunk:") { + if started.elapsed() >= PROPOSAL_TIMEOUT { + return Err(draft.failure_diagnostic(&format!( + "one mechanistic draft for {} and {}; probe stderr:\n{probe_stderr}", + expected.endpoint, expected.binary + ))); } - assert_mechanistic_draft(draft.stdout(), &binary) - .map_err(|error| draft.failure_diagnostic(&error))?; - return Ok(()); + sleep(POLL_INTERVAL).await; + continue; } - }) + return assert_mechanistic_draft(draft.stdout(), expected) + .map_err(|error| draft.failure_diagnostic(&error)); + } } -fn assert_mechanistic_draft(output: &str, binary: &str) -> Result<(), String> { +fn assert_mechanistic_draft(output: &str, expected: &ExpectedDraft<'_>) -> Result<(), String> { let fields = output.lines().map(str::trim).collect::>(); let field = |name: &str| { fields .iter() .find_map(|line| line.strip_prefix(name).map(str::trim)) }; + let endpoints = format!("{} [L4]", expected.endpoint); if fields .iter() .filter(|line| line.starts_with("Chunk:")) .count() != 1 || !matches!(field("Status:"), Some("pending" | "approved")) - || field("Rule:") != Some("allow_1_1_1_1_443") - || field("Binary:") != Some(binary) - || field("Binaries:") != Some(binary) - || field("Endpoints:") != Some("1.1.1.1:443 [L4]") - || !field("Rationale:").is_some_and(|value| value.contains("1.1.1.1:443")) + || field("Rule:") != Some(expected.rule) + || field("Binary:") != Some(expected.binary) + || field("Binaries:") != Some(expected.binary) + || field("Endpoints:") != Some(endpoints.as_str()) + || !field("Rationale:").is_some_and(|value| value.contains(expected.endpoint)) { return Err(format!( - "expected one pending or approved L4 mechanistic draft scoped to {binary} and 1.1.1.1:443" + "expected one pending or approved L4 mechanistic draft scoped to {} and {}", + expected.binary, expected.endpoint )); } Ok(()) @@ -462,11 +547,36 @@ async fn create_sandbox( #[cfg(test)] mod tests { - use super::assert_mechanistic_draft; + use super::{ExpectedDraft, assert_mechanistic_draft}; + + const EXPECTED: ExpectedDraft<'static> = ExpectedDraft { + rule: "allow_pypi_org_80", + endpoint: "pypi.org:80", + binary: "/usr/bin/bash", + }; + + fn draft(binary: &str) -> String { + format!( + "Chunk: id\nStatus: pending\nRule: allow_pypi_org_80\nBinary: {binary}\nRationale: Allow {binary} to connect to pypi.org:80 (HTTP).\nEndpoints: pypi.org:80 [L4]\nBinaries: {binary}\n" + ) + } + + #[test] + fn draft_assertion_accepts_a_hostname_scoped_draft() { + assert!(assert_mechanistic_draft(&draft("/usr/bin/bash"), &EXPECTED).is_ok()); + } #[test] fn draft_assertion_rejects_unrelated_binary() { - let draft = "Chunk: id\nStatus: pending\nRule: allow_1_1_1_1_443\nBinary: /usr/bin/sh\nRationale: Allow sh to connect to 1.1.1.1:443.\nEndpoints: 1.1.1.1:443 [L4]\nBinaries: /usr/bin/sh\n"; - assert!(assert_mechanistic_draft(draft, "/usr/bin/bash").is_err()); + assert!(assert_mechanistic_draft(&draft("/usr/bin/sh"), &EXPECTED).is_err()); + } + + #[test] + fn draft_assertion_rejects_fields_spread_across_drafts() { + let drafts = format!( + "{}Chunk: other\nStatus: pending\nRule: allow_1_1_1_1_443\n", + draft("/usr/bin/bash") + ); + assert!(assert_mechanistic_draft(&drafts, &EXPECTED).is_err()); } } diff --git a/crates/openshell-driver-docker/src/lib.rs b/crates/openshell-driver-docker/src/lib.rs index c76c4c838b..1f4c91f5e4 100644 --- a/crates/openshell-driver-docker/src/lib.rs +++ b/crates/openshell-driver-docker/src/lib.rs @@ -5796,6 +5796,10 @@ fn build_container_create_body_for_image( }), network_mode: Some("none".to_string()), dns: Some(vec!["127.0.0.53".to_string()]), + // Docker otherwise copies the host's search domains. Like the + // Podman, Kubernetes, and VM resolvers, send names to policy DNS + // exactly as written so short names never expand to host domains. + dns_search: Some(vec![".".to_string()]), tmpfs: Some(HashMap::from([( openshell_sandbox_backend::SUPERVISOR_CA_RUNTIME_DIR.to_string(), format!( diff --git a/crates/openshell-driver-docker/src/tests.rs b/crates/openshell-driver-docker/src/tests.rs index aa97d0ff09..c5036f4547 100644 --- a/crates/openshell-driver-docker/src/tests.rs +++ b/crates/openshell-driver-docker/src/tests.rs @@ -1278,6 +1278,7 @@ fn container_creation_uses_inspected_immutable_image() { ); assert_eq!(host.network_mode.as_deref(), Some("none")); assert_eq!(host.dns, Some(vec!["127.0.0.53".to_string()])); + assert_eq!(host.dns_search, Some(vec![".".to_string()])); } #[test] @@ -2714,6 +2715,11 @@ fn build_container_create_body_disables_docker_networking() { ); assert_eq!(host_config.extra_hosts, None); assert_eq!(host_config.dns, Some(vec!["127.0.0.53".to_string()])); + assert_eq!( + host_config.dns_search, + Some(vec![".".to_string()]), + "host search domains must not expand workload names before policy DNS" + ); } #[test] diff --git a/crates/openshell-supervisor-network/src/opa.rs b/crates/openshell-supervisor-network/src/opa.rs index 9db54fa900..c883651aaf 100644 --- a/crates/openshell-supervisor-network/src/opa.rs +++ b/crates/openshell-supervisor-network/src/opa.rs @@ -73,6 +73,8 @@ pub struct MatchedEndpoint { pub struct PolicyDnsEligibilitySnapshot { pub endpoints: Vec, pub generation: u64, + /// The generation is a fail-closed quarantine, so every name is refused. + pub fail_closed: bool, } /// Atomic policy result used to authorize and materialize one egress request. @@ -645,7 +647,7 @@ impl OpaEngine { /// The owned endpoint records and generation are captured while holding /// the engine lock, so reloads cannot mix data from one generation with /// the generation number of another. Fail-closed quarantine produces an - /// empty snapshot for its quarantine generation. + /// empty snapshot, marked `fail_closed`, for its quarantine generation. pub fn policy_dns_eligibility_snapshot(&self) -> Result { let mut engine = self .engine @@ -662,6 +664,7 @@ impl OpaEngine { return Ok(PolicyDnsEligibilitySnapshot { endpoints: Vec::new(), generation, + fail_closed: true, }); } @@ -679,6 +682,7 @@ impl OpaEngine { Ok(PolicyDnsEligibilitySnapshot { endpoints, generation, + fail_closed: false, }) } @@ -4033,6 +4037,7 @@ process: let snapshot = engine.policy_dns_eligibility_snapshot().unwrap(); assert_eq!(snapshot.generation, engine.current_generation()); + assert!(!snapshot.fail_closed); assert_eq!(snapshot.endpoints.len(), 4); assert_eq!(snapshot.endpoints[0].policy_name, "dns_transport"); assert_eq!(snapshot.endpoints[0].endpoint_index, 0); @@ -4068,6 +4073,7 @@ process: assert_eq!(snapshot.generation, generation); assert!(snapshot.endpoints.is_empty()); + assert!(snapshot.fail_closed); } #[test] diff --git a/crates/openshell-supervisor-network/src/policy_dns/mod.rs b/crates/openshell-supervisor-network/src/policy_dns/mod.rs index a57caec2e7..ee1a5cc122 100644 --- a/crates/openshell-supervisor-network/src/policy_dns/mod.rs +++ b/crates/openshell-supervisor-network/src/policy_dns/mod.rs @@ -149,7 +149,13 @@ impl PolicyDnsService { "policy_dns_ineligible", "Policy DNS refused a name that is not eligible in the active policy", ); - return Err(PolicyDnsError::Ineligible); + // DNS alone cannot identify the requesting executable or intended + // port. Outside quarantine, stage a contract-free observation so + // the later TCP attempt can be denied with both and proposed. + if snapshot.fail_closed || !observation_eligible(&normalized_name) { + return Err(PolicyDnsError::Ineligible); + } + return self.stage_observation(&normalized_name, family, snapshot.generation, now); } // The trusted resolver is invoked only after the immutable snapshot @@ -221,58 +227,110 @@ impl PolicyDnsService { ttl, contracts, }; - let record = match self + let record = self.publish_guarded( + &normalized_name, + family, + &endpoint_context, + snapshot.generation, + |current_generation| self.store.publish(request, current_generation, now), + )?; + emit_mapping_publication(&record); + Ok(SyntheticAnswer { + address: record.synthetic_address, + ttl, + mapping_id: record.mapping_id, + mapping_generation: record.mapping_generation, + policy_generation: record.policy_generation, + }) + } + + /// Stage a contract-free observation for a name absent from policy. No + /// upstream DNS request or egress grant is made. + fn stage_observation( + &self, + name: &NormalizedName, + family: AddressFamily, + policy_generation: u64, + now: Instant, + ) -> Result { + let record = self + .publish_guarded(name, family, &[], policy_generation, |current_generation| { + self.store.publish_observation( + name.clone(), + family, + policy_generation, + current_generation, + now, + ) + }) + .map_err(|error| match error { + // Without capacity, keep the refusal unknown names received + // before observations existed. + PolicyDnsError::Publish( + PublishError::ObservationBudgetExhausted | PublishError::PoolExhausted, + ) => PolicyDnsError::Ineligible, + error => error, + })?; + emit_observation_publication(&record); + Ok(SyntheticAnswer { + address: record.synthetic_address, + ttl: MAX_MAPPING_TTL, + mapping_id: record.mapping_id, + mapping_generation: record.mapping_generation, + policy_generation: record.policy_generation, + }) + } + + /// Publish only while `policy_generation` is current. The store rejects + /// invalid or stale mappings, exhausted capacity, and a poisoned lock; + /// every rejection stays observable. + fn publish_guarded( + &self, + name: &NormalizedName, + family: AddressFamily, + endpoint_context: &[PolicyEndpointId], + policy_generation: u64, + publish: impl FnOnce(u64) -> Result, + ) -> Result { + match self .policy - .with_current_generation(snapshot.generation, |current_generation| { - self.store.publish(request, current_generation, now) - }) { - Ok(Some(Ok(record))) => record, + .with_current_generation(policy_generation, publish) + { + Ok(Some(Ok(record))) => Ok(record), Ok(Some(Err(error))) => { - // InvalidMapping is unreachable for the well-formed request - // assembled above, and LockPoisoned requires a prior panic - // while holding the store lock. Keep both defensive outcomes - // observable because the store API intentionally rejects them. emit_dns_failure( - &normalized_name, + name, family, - &endpoint_context, - snapshot.generation, + endpoint_context, + policy_generation, publication_failure_detail(error), "Policy DNS resolved-endpoint mapping publication failed", ); - return Err(PolicyDnsError::Publish(error)); + Err(PolicyDnsError::Publish(error)) } Ok(None) => { emit_dns_failure( - &normalized_name, + name, family, - &endpoint_context, - snapshot.generation, + endpoint_context, + policy_generation, "policy_dns_publication_stale_generation", "Policy DNS discarded a stale resolved-endpoint mapping", ); - return Err(PolicyDnsError::StalePolicy); + Err(PolicyDnsError::StalePolicy) } Err(error) => { emit_dns_failure( - &normalized_name, + name, family, - &endpoint_context, - snapshot.generation, + endpoint_context, + policy_generation, "policy_dns_publication_generation_check_failed", "Policy DNS could not validate the active policy generation before publication", ); - return Err(PolicyDnsError::Policy(error.to_string())); + Err(PolicyDnsError::Policy(error.to_string())) } - }; - emit_mapping_publication(&record); - Ok(SyntheticAnswer { - address: record.synthetic_address, - ttl, - mapping_id: record.mapping_id, - mapping_generation: record.mapping_generation, - policy_generation: record.policy_generation, - }) + } } pub(crate) fn store(&self) -> &Arc { @@ -280,6 +338,13 @@ impl PolicyDnsService { } } +/// Reserved names keep the plain refusal and never become proposals. +fn observation_eligible(name: &NormalizedName) -> bool { + name.as_str() != "localhost" + && !openshell_core::net::is_known_metadata_hostname(name.as_str()) + && !is_host_gateway_alias(name.as_str()) +} + struct EligibleEndpoint { endpoint_id: PolicyEndpointId, ports: Vec, @@ -412,19 +477,35 @@ fn clamp_mapping_ttl(ttl: Duration) -> Duration { ttl.max(MIN_MAPPING_TTL).min(MAX_MAPPING_TTL) } +/// A policy DNS decision concerns the queried name, not a connection to it, +/// so the endpoint carries no port. +fn dns_query_endpoint(name: &NormalizedName) -> Endpoint { + Endpoint { + domain: Some(name.as_str().to_string()), + ip: None, + port: None, + } +} + fn emit_dns_denial(name: &NormalizedName, detail: &str, message: &str) { - ocsf_emit!( - NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) - .activity(ActivityId::Refuse) - .action(ActionId::Denied) - .disposition(DispositionId::Blocked) - .severity(SeverityId::Medium) - .status(StatusId::Failure) - .dst_endpoint(Endpoint::from_domain(name.as_str(), 53)) - .status_detail(detail) - .message(message) - .build() - ); + ocsf_emit!(build_dns_denial_event(name, detail, message)); +} + +fn build_dns_denial_event( + name: &NormalizedName, + detail: &str, + message: &str, +) -> openshell_ocsf::OcsfEvent { + NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + .activity(ActivityId::Refuse) + .action(ActionId::Denied) + .disposition(DispositionId::Blocked) + .severity(SeverityId::Medium) + .status(StatusId::Failure) + .dst_endpoint(dns_query_endpoint(name)) + .status_detail(detail) + .message(message) + .build() } fn resolver_failure_detail(error: &resolver::ResolveError) -> &'static str { @@ -445,6 +526,7 @@ fn publication_failure_detail(error: PublishError) -> &'static str { PublishError::StalePolicy => "policy_dns_publication_stale_generation", PublishError::InvalidMapping => "policy_dns_publication_invalid_mapping", PublishError::PoolExhausted => "policy_dns_publication_pool_exhausted", + PublishError::ObservationBudgetExhausted => "policy_dns_observation_budget_exhausted", PublishError::LockPoisoned => "policy_dns_publication_store_unavailable", } } @@ -472,7 +554,7 @@ fn build_dns_failure_event( .disposition(DispositionId::Blocked) .severity(SeverityId::Low) .status(StatusId::Failure) - .dst_endpoint(Endpoint::from_domain(name.as_str(), 53)) + .dst_endpoint(dns_query_endpoint(name)) .status_detail(detail) .unmapped("normalized_name", name.as_str()) .unmapped("address_family", family.as_str()) @@ -507,6 +589,30 @@ fn emit_mapping_publication(record: &ResolvedEndpointRecord) { ocsf_emit!(build_mapping_publication_event(record)); } +fn emit_observation_publication(record: &ResolvedEndpointRecord) { + ocsf_emit!(build_observation_publication_event(record)); +} + +fn build_observation_publication_event( + record: &ResolvedEndpointRecord, +) -> openshell_ocsf::OcsfEvent { + ConfigStateChangeBuilder::new(openshell_ocsf::ctx::ctx()) + .severity(SeverityId::Informational) + .status(StatusId::Success) + .state(StateId::Enabled, "observation") + .unmapped("normalized_domain", record.normalized_name.as_str()) + .unmapped("address_family", format!("{:?}", record.family)) + .unmapped("synthetic_ip", record.synthetic_address.to_string()) + .unmapped("policy_generation", record.policy_generation) + .unmapped("mapping_generation", record.mapping_generation) + .unmapped("mapping_id", record.mapping_id.to_string()) + .message(format!( + "Policy DNS staged unapproved name {} synthetic={} for TCP policy review", + record.normalized_name, record.synthetic_address + )) + .build() +} + fn build_mapping_publication_event(record: &ResolvedEndpointRecord) -> openshell_ocsf::OcsfEvent { let approved_real_ip_candidates = record .contracts @@ -648,15 +754,272 @@ process: { run_as_user: sandbox, run_as_group: sandbox } "; #[tokio::test] - async fn refuses_ineligible_name_before_upstream_resolution() { + async fn stages_unknown_name_without_upstream_resolution() { let service = service(BASE_POLICY, vec!["8.8.8.8".parse().unwrap()]); - let result = service + let answer = service .answer_query("other.example", AddressFamily::Ipv4, Instant::now()) - .await; - assert!(matches!(result, Err(PolicyDnsError::Ineligible))); + .await + .expect("unapproved name receives a local synthetic address"); + assert!( + service + .store + .lookup_intent( + answer.address, + 80, + service.policy.current_generation(), + Instant::now() + ) + .is_ok() + ); + assert_eq!(service.resolver.calls.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn unknown_name_correlates_without_authorizing_a_port() { + let service = service(BASE_POLICY, vec!["8.8.8.8".parse().unwrap()]); + let now = Instant::now(); + let answer = service + .answer_query("PyPI.org", AddressFamily::Ipv4, now) + .await + .expect("synthetic observation address"); + assert_eq!(service.resolver.calls.load(Ordering::SeqCst), 0); + let generation = service.policy.current_generation(); + assert!(matches!( + service.store.lookup(answer.address, 80, generation, now), + Err(MappingLookupError::PortMismatch) + )); + let intent = service + .store + .lookup_intent(answer.address, 80, generation, now) + .expect("name available for a later denied connection"); + assert_eq!(intent.record.normalized_name.as_str(), "pypi.org"); + assert!(intent.record.is_observation()); + assert!(matches!( + service + .store + .lookup_intent(answer.address, 80, generation + 1, now), + Err(MappingLookupError::StalePolicy) + )); + } + + #[tokio::test] + async fn reserved_names_and_quarantine_are_refused_without_spending_the_budget() { + let service = service_with_gateway( + BASE_POLICY, + Vec::new(), + Some(IpAddr::V4(Ipv4Addr::new(172, 17, 0, 1))), + ); + let now = Instant::now(); + for name in [ + "localhost", + "metadata.google.internal", + "host.openshell.internal", + ] { + assert!( + matches!( + service.answer_query(name, AddressFamily::Ipv4, now).await, + Err(PolicyDnsError::Ineligible) + ), + "{name} must stay refused" + ); + } + + service + .policy + .enter_fail_closed("invalid candidate") + .unwrap(); + assert!(matches!( + service + .answer_query("pypi.org", AddressFamily::Ipv4, now) + .await, + Err(PolicyDnsError::Ineligible) + )); + service.policy.exit_fail_closed().unwrap(); + + // This pool admits two observations; none were spent above. + for name in ["pypi.org", "files.pythonhosted.org"] { + service + .answer_query(name, AddressFamily::Ipv4, now) + .await + .expect("observation budget remains available"); + } assert_eq!(service.resolver.calls.load(Ordering::SeqCst), 0); } + #[tokio::test] + async fn exhausted_observation_budget_refuses_only_new_unknown_names() { + let service = service(BASE_POLICY, vec!["203.0.113.8".parse().unwrap()]); + let now = Instant::now(); + let first = service + .answer_query("first.example", AddressFamily::Ipv4, now) + .await + .unwrap(); + service + .answer_query("second.example", AddressFamily::Ipv4, now) + .await + .unwrap(); + + assert!(matches!( + service + .answer_query("third.example", AddressFamily::Ipv4, now) + .await, + Err(PolicyDnsError::Ineligible) + )); + let refreshed = service + .answer_query("first.example", AddressFamily::Ipv4, now) + .await + .expect("a staged name keeps its observation"); + assert_eq!(refreshed.address, first.address); + service + .answer_query("db.example", AddressFamily::Ipv4, now) + .await + .expect("policy-backed names keep their reserved capacity"); + } + + struct OcsfCapture(Arc>>); + + impl tracing_subscriber::Layer for OcsfCapture { + fn on_event(&self, _: &tracing::Event<'_>, _: tracing_subscriber::layer::Context<'_, S>) { + if let Some(event) = openshell_ocsf::tracing_layers::clone_current_event() { + self.0 + .lock() + .unwrap() + .push(serde_json::to_value(&event).unwrap()); + } + } + } + + #[tokio::test] + async fn unknown_names_keep_the_ineligible_denial_event() { + use tracing_subscriber::layer::SubscriberExt as _; + + const CHILD: &str = "OPENSHELL_TEST_POLICY_DNS_EVENTS_CHILD"; + // Tracing callsite interest is process-wide, so concurrent tests with + // other subscribers can disable this thread's capture. Exercise the + // service in an isolated test process with the same executable. + if std::env::var_os(CHILD).is_none() { + let output = std::process::Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "policy_dns::tests::unknown_names_keep_the_ineligible_denial_event", + "--nocapture", + ]) + .env(CHILD, "1") + .output() + .unwrap(); + let stdout = String::from_utf8_lossy(&output.stdout); + assert!( + output.status.success() && stdout.contains("1 passed"), + "{stdout}\n{}", + String::from_utf8_lossy(&output.stderr) + ); + return; + } + + let service = service(BASE_POLICY, Vec::new()); + let events = Arc::new(std::sync::Mutex::new(Vec::new())); + let _subscriber = tracing::subscriber::set_default( + tracing_subscriber::registry().with(OcsfCapture(Arc::clone(&events))), + ); + let drain = || { + events + .lock() + .unwrap() + .drain(..) + .map(|event| { + event["status_detail"] + .as_str() + .or_else(|| event["state"].as_str()) + .unwrap_or_default() + .to_string() + }) + .collect::>() + }; + let now = Instant::now(); + + for name in ["first.example", "second.example"] { + service + .answer_query(name, AddressFamily::Ipv4, now) + .await + .unwrap(); + assert_eq!(drain(), ["policy_dns_ineligible", "observation"], "{name}"); + } + assert!( + service + .answer_query("third.example", AddressFamily::Ipv4, now) + .await + .is_err() + ); + assert_eq!( + drain(), + [ + "policy_dns_ineligible", + "policy_dns_observation_budget_exhausted" + ] + ); + service + .policy + .enter_fail_closed("invalid candidate") + .unwrap(); + assert!( + service + .answer_query("fourth.example", AddressFamily::Ipv4, now) + .await + .is_err() + ); + assert_eq!(drain(), ["policy_dns_ineligible"]); + } + + #[test] + fn dns_denial_names_the_query_without_a_port() { + let event = build_dns_denial_event( + &NormalizedName::parse("blocked.invalid").unwrap(), + "policy_dns_ineligible", + "Policy DNS refused a name that is not eligible in the active policy", + ); + + assert_eq!( + event.format_shorthand(), + "NET:REFUSE [MED] DENIED blocked.invalid [reason:policy_dns_ineligible]" + ); + assert_eq!( + serde_json::to_value(event).unwrap()["dst_endpoint"], + serde_json::json!({"domain": "blocked.invalid"}) + ); + } + + #[test] + fn observation_event_correlates_the_name_with_its_synthetic_address() { + let store = ResolvedEndpointStore::new( + StoreConfig::new( + SyntheticPools::new( + Ipv4Addr::new(198, 18, 0, 1)..=Ipv4Addr::new(198, 18, 0, 8), + "fd00:1::1".parse::().unwrap() + ..="fd00:1::8".parse::().unwrap(), + ) + .unwrap(), + 16, + ) + .unwrap(), + ); + let record = store + .publish_observation( + NormalizedName::parse("pypi.org").unwrap(), + AddressFamily::Ipv4, + 3, + 3, + Instant::now(), + ) + .unwrap(); + + let json = serde_json::to_value(build_observation_publication_event(&record)).unwrap(); + + assert_eq!(json["state"], "observation"); + assert_eq!(json["unmapped"]["normalized_domain"], "pypi.org"); + assert_eq!(json["unmapped"]["synthetic_ip"], "198.18.0.1"); + assert_eq!(json["unmapped"]["policy_generation"], 3); + } + #[tokio::test] async fn eligible_nxdomain_fails_without_publishing_a_mapping() { let policy = Arc::new( @@ -999,8 +1362,10 @@ process: { run_as_user: sandbox, run_as_group: sandbox } assert_eq!(json["severity"], "Low"); assert_eq!(json["status"], "Failure"); assert_eq!(json["status_detail"], "policy_dns_upstream_nxdomain"); - assert_eq!(json["dst_endpoint"]["domain"], "db.example"); - assert_eq!(json["dst_endpoint"]["port"], 53); + assert_eq!( + json["dst_endpoint"], + serde_json::json!({"domain": "db.example"}) + ); assert_eq!(json["unmapped"]["normalized_name"], "db.example"); assert_eq!(json["unmapped"]["address_family"], "ipv4"); assert_eq!(json["unmapped"]["policy_generation"], 7); diff --git a/crates/openshell-supervisor-network/src/policy_dns/store.rs b/crates/openshell-supervisor-network/src/policy_dns/store.rs index 439cdcfa36..eeb5d4b85f 100644 --- a/crates/openshell-supervisor-network/src/policy_dns/store.rs +++ b/crates/openshell-supervisor-network/src/policy_dns/store.rs @@ -17,6 +17,11 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::time::{Duration, Instant}; use uuid::Uuid; +/// Unapproved-name observations may hold at most this fraction of each +/// address family's identities; the rest stay available to policy-backed +/// names. Allocations are never reused, even after their TTL expires. +const OBSERVATION_POOL_DIVISOR: usize = 4; + #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] pub(crate) struct PolicyEndpointId { pub(crate) policy_name: String, @@ -58,6 +63,12 @@ pub(crate) struct ResolvedEndpointRecord { } impl ResolvedEndpointRecord { + /// An observation correlates a name absent from policy with a later TCP + /// attempt. It carries no endpoint contract and never authorizes egress. + pub(crate) fn is_observation(&self) -> bool { + self.contracts.is_empty() + } + pub(crate) fn allowed_ports(&self) -> BTreeSet { self.contracts .iter() @@ -206,6 +217,8 @@ pub(crate) enum PublishError { InvalidMapping, #[error("synthetic address pool is exhausted")] PoolExhausted, + #[error("unapproved-name observation budget is exhausted")] + ObservationBudgetExhausted, #[error("resolved endpoint store lock was poisoned")] LockPoisoned, } @@ -228,11 +241,26 @@ pub(crate) enum MappingLookupError { LockPoisoned, } +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +enum AllocationIdentity { + /// Contract-free correlation for a name absent from policy. + Observation, + /// Digest of every compatible endpoint contract. + Contracts([u8; 32]), +} + #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] struct AllocationKey { normalized_name: NormalizedName, family: AddressFamily, - allocation_identity: [u8; 32], + identity: AllocationIdentity, +} + +struct Publication { + key: AllocationKey, + policy_generation: u64, + ttl: Duration, + contracts: Vec, } struct StoreState { @@ -299,32 +327,106 @@ impl ResolvedEndpointStore { return Err(PublishError::InvalidMapping); } - let key = AllocationKey { - normalized_name: request.normalized_name.clone(), - family: request.family, - allocation_identity: request.allocation_identity, - }; + self.publish_checked( + Publication { + key: AllocationKey { + normalized_name: request.normalized_name, + family: request.family, + identity: AllocationIdentity::Contracts(request.allocation_identity), + }, + policy_generation: request.policy_generation, + ttl: request.ttl, + contracts: request.contracts, + }, + now, + ) + } + + /// Correlate a denied hostname with a later TCP attempt without resolving + /// it upstream or granting any port. Only the intent lookup may read this + /// record; the authorization lookup requires an endpoint contract. + pub(crate) fn publish_observation( + &self, + normalized_name: NormalizedName, + family: AddressFamily, + policy_generation: u64, + current_policy_generation: u64, + now: Instant, + ) -> Result { + if policy_generation != current_policy_generation { + return Err(PublishError::StalePolicy); + } + self.publish_checked( + Publication { + key: AllocationKey { + normalized_name, + family, + identity: AllocationIdentity::Observation, + }, + policy_generation, + ttl: super::MAX_MAPPING_TTL, + contracts: Vec::new(), + }, + now, + ) + } + + fn observation_budget(&self, family: AddressFamily) -> usize { + (self + .config + .pools + .capacity(family) + .min(self.config.max_mappings) + / OBSERVATION_POOL_DIVISOR) + .max(1) + } + + fn publish_checked( + &self, + publication: Publication, + now: Instant, + ) -> Result { + let Publication { + key, + policy_generation, + ttl, + contracts, + } = publication; + let family = key.family; let mut state = self.state.write().map_err(|_| PublishError::LockPoisoned)?; let synthetic_address = if let Some(address) = state.allocations.get(&key) { *address } else { + if key.identity == AllocationIdentity::Observation + && state + .allocations + .keys() + .filter(|allocated| { + allocated.identity == AllocationIdentity::Observation + && allocated.family == family + }) + .count() + >= self.observation_budget(family) + { + return Err(PublishError::ObservationBudgetExhausted); + } if state.allocations.len() >= self.config.max_mappings { return Err(PublishError::PoolExhausted); } let address = - allocate_address(&mut state, request.family).ok_or(PublishError::PoolExhausted)?; - state.allocations.insert(key, address); + allocate_address(&mut state, family).ok_or(PublishError::PoolExhausted)?; + state.allocations.insert(key.clone(), address); let allocated = state .allocations .keys() - .filter(|key| key.family == request.family) + .filter(|allocated| allocated.family == family) .count(); let capacity = self .config .pools - .capacity(request.family) + .capacity(family) .min(self.config.max_mappings); - let emitted = match request.family { + let emitted = match family { AddressFamily::Ipv4 => &self.ipv4_pool_high_water_emitted, AddressFamily::Ipv6 => &self.ipv6_pool_high_water_emitted, }; @@ -332,9 +434,7 @@ impl ResolvedEndpointStore { && !emitted.swap(true, Ordering::Relaxed) { openshell_ocsf::ocsf_emit!(build_pool_high_water_event( - allocated, - capacity, - request.family, + allocated, capacity, family, )); } address @@ -346,7 +446,7 @@ impl ResolvedEndpointStore { if state .records .get(&synthetic_address) - .is_some_and(|record| record.policy_generation > request.policy_generation) + .is_some_and(|record| record.policy_generation > policy_generation) { return Err(PublishError::StalePolicy); } @@ -354,14 +454,14 @@ impl ResolvedEndpointStore { state.next_mapping_generation = state.next_mapping_generation.saturating_add(1); let record = ResolvedEndpointRecord { synthetic_address, - normalized_name: request.normalized_name, - family: request.family, - policy_generation: request.policy_generation, + normalized_name: key.normalized_name, + family, + policy_generation, mapping_generation: state.next_mapping_generation, mapping_id: Uuid::new_v4(), created_at: now, - expires_at: now + request.ttl, - contracts: request.contracts, + expires_at: now + ttl, + contracts, }; state.expired_allocations.remove(&synthetic_address); state.records.insert(synthetic_address, record.clone()); @@ -379,19 +479,7 @@ impl ResolvedEndpointStore { .state .read() .map_err(|_| MappingLookupError::LockPoisoned)?; - let Some(record) = state.records.get(&synthetic_address) else { - return if state.expired_allocations.contains(&synthetic_address) { - Err(MappingLookupError::Expired) - } else { - Err(MappingLookupError::Missing) - }; - }; - if now >= record.expires_at { - return Err(MappingLookupError::Expired); - } - if record.policy_generation != current_policy_generation { - return Err(MappingLookupError::StalePolicy); - } + let record = live_record(&state, synthetic_address, current_policy_generation, now)?; if !record .contracts .iter() @@ -405,6 +493,35 @@ impl ResolvedEndpointStore { }) } + /// Look up the hostname of an observation-only mapping for a policy + /// decision. Ordinary mappings still enforce their authored port scope. + /// The caller must never use this lookup to establish an upstream relay. + pub(crate) fn lookup_intent( + &self, + synthetic_address: IpAddr, + port: u16, + current_policy_generation: u64, + now: Instant, + ) -> Result { + let state = self + .state + .read() + .map_err(|_| MappingLookupError::LockPoisoned)?; + let record = live_record(&state, synthetic_address, current_policy_generation, now)?; + if !record.is_observation() + && !record + .contracts + .iter() + .any(|contract| contract.port == port) + { + return Err(MappingLookupError::PortMismatch); + } + Ok(MappingLookup { + record: record.clone(), + port, + }) + } + /// Remove expired active records without freeing their synthetic identity. pub(crate) fn expire(&self, now: Instant) -> Result { let mut state = self @@ -424,6 +541,28 @@ impl ResolvedEndpointStore { } } +fn live_record( + state: &StoreState, + synthetic_address: IpAddr, + current_policy_generation: u64, + now: Instant, +) -> Result<&ResolvedEndpointRecord, MappingLookupError> { + let Some(record) = state.records.get(&synthetic_address) else { + return if state.expired_allocations.contains(&synthetic_address) { + Err(MappingLookupError::Expired) + } else { + Err(MappingLookupError::Missing) + }; + }; + if now >= record.expires_at { + return Err(MappingLookupError::Expired); + } + if record.policy_generation != current_policy_generation { + return Err(MappingLookupError::StalePolicy); + } + Ok(record) +} + fn build_pool_high_water_event( allocated: usize, capacity: usize, @@ -497,6 +636,123 @@ mod tests { } } + fn observation( + store: &ResolvedEndpointStore, + name: &str, + family: AddressFamily, + now: Instant, + ) -> Result { + store.publish_observation(NormalizedName::parse(name).unwrap(), family, 1, 1, now) + } + + #[test] + fn observation_budget_preserves_capacity_for_policy_backed_names() { + let store = store(8); + let now = Instant::now(); + let first = observation(&store, "unknown.example", AddressFamily::Ipv4, now).unwrap(); + assert!(first.is_observation()); + assert!(matches!( + observation(&store, "another.example", AddressFamily::Ipv4, now), + Err(PublishError::ObservationBudgetExhausted) + )); + let refreshed = observation(&store, "unknown.example", AddressFamily::Ipv4, now).unwrap(); + assert_eq!(refreshed.synthetic_address, first.synthetic_address); + let approved = store + .publish(request("db.example", 1, Duration::from_secs(10)), 1, now) + .unwrap(); + assert_ne!(first.synthetic_address, approved.synthetic_address); + } + + #[test] + fn observation_budget_is_counted_per_address_family() { + let store = store(8); + let now = Instant::now(); + observation(&store, "first.example", AddressFamily::Ipv4, now).unwrap(); + assert!(matches!( + observation(&store, "second.example", AddressFamily::Ipv4, now), + Err(PublishError::ObservationBudgetExhausted) + )); + let ipv6 = observation(&store, "second.example", AddressFamily::Ipv6, now).unwrap(); + assert!(ipv6.synthetic_address.is_ipv6()); + } + + #[test] + fn production_observation_budget_is_a_quarter_of_each_family_pool() { + let pools = SyntheticPools::new( + Ipv4Addr::new(198, 18, 0, 2)..=Ipv4Addr::new(198, 18, 1, 255), + "fd23:6f70:656e::".parse().unwrap()..="fd23:6f70:656e::1ff".parse().unwrap(), + ) + .unwrap(); + let store = ResolvedEndpointStore::new(StoreConfig::new(pools, 1024).unwrap()); + let now = Instant::now(); + + for index in 0..127 { + observation( + &store, + &format!("unknown-{index}.example"), + AddressFamily::Ipv4, + now, + ) + .unwrap(); + } + assert!(matches!( + observation(&store, "unknown-127.example", AddressFamily::Ipv4, now), + Err(PublishError::ObservationBudgetExhausted) + )); + for index in 0..383 { + store + .publish( + request(&format!("host-{index}.example"), 1, Duration::from_secs(5)), + 1, + now, + ) + .unwrap(); + } + assert!(matches!( + store.publish( + request("host-383.example", 1, Duration::from_secs(5)), + 1, + now + ), + Err(PublishError::PoolExhausted) + )); + } + + #[test] + fn only_the_intent_lookup_exposes_an_observation() { + let store = store(8); + let now = Instant::now(); + let observed = observation(&store, "unknown.example", AddressFamily::Ipv4, now).unwrap(); + let approved = store + .publish(request("db.example", 1, Duration::from_mins(1)), 1, now) + .unwrap(); + + assert!(matches!( + store.lookup(observed.synthetic_address, 443, 1, now), + Err(MappingLookupError::PortMismatch) + )); + let intent = store + .lookup_intent(observed.synthetic_address, 443, 1, now) + .unwrap(); + assert_eq!(intent.record.normalized_name.as_str(), "unknown.example"); + assert!(intent.pinned_addresses().is_empty()); + assert!(matches!( + store.lookup_intent(observed.synthetic_address, 443, 2, now), + Err(MappingLookupError::StalePolicy) + )); + assert!(matches!( + store.lookup_intent(approved.synthetic_address, 3306, 1, now), + Err(MappingLookupError::PortMismatch) + )); + + let expired_at = now + crate::policy_dns::MAX_MAPPING_TTL; + assert_eq!(store.expire(expired_at).unwrap(), 1); + assert!(matches!( + store.lookup_intent(observed.synthetic_address, 443, 1, expired_at), + Err(MappingLookupError::Expired) + )); + } + #[test] fn refresh_retains_synthetic_identity_and_changes_mapping_generation() { let store = store(2); diff --git a/crates/openshell-supervisor-network/src/policy_dns/wire.rs b/crates/openshell-supervisor-network/src/policy_dns/wire.rs index b9b79213d4..fe6a01c3a4 100644 --- a/crates/openshell-supervisor-network/src/policy_dns/wire.rs +++ b/crates/openshell-supervisor-network/src/policy_dns/wire.rs @@ -289,11 +289,26 @@ process: { run_as_user: sandbox, run_as_group: sandbox } } #[tokio::test] - async fn ineligible_query_is_refused_without_upstream_call() { + async fn unknown_query_gets_observation_address_without_upstream_call() { let service = service(); let wire = handle_udp_query(&service, &request("other.example.", RecordType::A)) .await .unwrap(); + let response = Message::from_vec(&wire).unwrap(); + assert_eq!(response.metadata.response_code, ResponseCode::NoError); + assert!(matches!(response.answers[0].data, RData::A(_))); + assert_eq!(service.resolver.calls.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn unknown_query_is_refused_after_the_observation_budget() { + let service = service(); + handle_udp_query(&service, &request("first.example.", RecordType::A)) + .await + .unwrap(); + let wire = handle_udp_query(&service, &request("second.example.", RecordType::A)) + .await + .unwrap(); assert_eq!( Message::from_vec(&wire).unwrap().metadata.response_code, ResponseCode::Refused diff --git a/crates/openshell-supervisor-network/src/proxy.rs b/crates/openshell-supervisor-network/src/proxy.rs index f23d8138af..5f224b474e 100644 --- a/crates/openshell-supervisor-network/src/proxy.rs +++ b/crates/openshell-supervisor-network/src/proxy.rs @@ -10,9 +10,9 @@ mod relay; use crate::identity::{BinaryIdentityCache, SuppliedIdentityError}; use crate::l7::tls::ProxyTlsState; use crate::opa::{NetworkAction, OpaEngine, PolicyGenerationGuard}; -#[cfg(target_os = "linux")] +#[cfg(any(target_os = "linux", test))] use crate::policy_dns::PolicyEndpointId; -use crate::policy_dns::{MappingLookupError, ResolvedEndpointStore}; +use crate::policy_dns::{MappingLookup, MappingLookupError, ResolvedEndpointStore}; use crate::policy_local::{POLICY_LOCAL_HOST, PolicyLocalContext}; use crate::upstream_proxy::{self, UpstreamProxyConfig}; use futures::{FutureExt as _, StreamExt as _, stream::FuturesUnordered}; @@ -77,8 +77,8 @@ enum ProxyAcceptError { } use self::destination::{ - DestinationDenial, DestinationDenialKind, DestinationRequest, build_pinned_validation_plan, - build_validation_plan, validate_destination, + DestinationDenial, DestinationDenialKind, DestinationRequest, DestinationValidationPlan, + build_pinned_validation_plan, build_validation_plan, validate_destination, }; use self::egress::{ EgressDecision, EgressIntent, EndpointDecision, IdentityUnavailableReason, L7ConfigSnapshot, @@ -614,6 +614,7 @@ async fn preauthorize_transparent_open( if destination.port() != 80 || !has_policy_local { emit_staged_transparent_denial( destination, + None, &binary_identity, "sandbox-local policy API requires port 80 and an active context", "transparent_tcp_policy_local_invalid_destination", @@ -637,6 +638,7 @@ async fn preauthorize_transparent_open( if let Err(denial) = identity_check { emit_staged_transparent_denial( destination, + None, &binary_identity, "sandbox-local policy API requires a verified workload identity", "transparent_tcp_policy_local_identity_unavailable", @@ -657,20 +659,23 @@ async fn preauthorize_transparent_open( }), )); } - let host = match transparent_destination_host(destination, policy_dns_store, opa_engine) { - Ok(host) => host, - Err(error) => { - warn!(%destination, %error, "Denied staged transparent connection"); - emit_staged_transparent_denial( - destination, - &binary_identity, - &error.to_string(), - "transparent_tcp_mapping_denied", - ); - let _ = completion.send(TcpOpenDecision::Denied(TcpOpenDenial::InvalidDestination)); - return None; - } - }; + let TransparentTarget { host, intent } = + match resolve_transparent_target(destination, policy_dns_store, opa_engine) { + Ok(target) => target, + Err(error) => { + warn!(%destination, %error, "Denied staged transparent connection"); + emit_staged_transparent_denial( + destination, + None, + &binary_identity, + &error.to_string(), + "transparent_tcp_mapping_denied", + ); + let _ = completion.send(TcpOpenDecision::Denied(TcpOpenDenial::InvalidDestination)); + return None; + } + }; + let mapped_host = intent.is_some().then_some(host.as_str()); let supplied_authorization = authorize_supplied_identity_with_denial( opa_engine, identity_cache, @@ -693,7 +698,13 @@ async fn preauthorize_transparent_open( }, ); warn!(%destination, %reason, "Denied staged transparent connection"); - emit_staged_transparent_denial(destination, &binary_identity, reason, status_detail); + emit_staged_transparent_denial( + destination, + mapped_host, + &binary_identity, + reason, + status_detail, + ); if supplied_authorization.denial.is_none() && !is_always_blocked_ip(destination.ip()) && let Some(binary) = decision.binary.as_ref() @@ -711,12 +722,53 @@ async fn preauthorize_transparent_open( let _ = completion.send(TcpOpenDecision::Denied(denial)); return None; } + // A synthetic destination may dial only the addresses pinned for the + // generation that produced this decision. Re-acquiring the mapping at that + // generation keeps a policy reload between correlation and authorization + // from falling back to resolving the name again. Observation records pin + // no addresses, so they never reach an upstream dial. + let pinned_plan = match &intent { + None => None, + Some(intent) => match policy_dns_store + .ok_or(MappingLookupError::Missing) + .and_then(|store| pinned_transparent_plan(store, destination, &decision)) + { + Ok(plan) => Some(plan), + Err(error) => { + let (reason, status_detail) = match error { + _ if intent.record.is_observation() => ( + "unapproved DNS observation cannot authorize egress".to_string(), + "transparent_tcp_observation_denied", + ), + MappingLookupError::InvalidMapping => ( + "policy DNS produced an invalid pinned destination".to_string(), + "transparent_tcp_destination_denied", + ), + error => ( + format!("transparent destination mapping is unavailable: {error}"), + "transparent_tcp_mapping_denied", + ), + }; + warn!(%destination, %reason, "Denied staged transparent connection"); + emit_staged_transparent_denial( + destination, + mapped_host, + &binary_identity, + &reason, + status_detail, + ); + let _ = completion.send(TcpOpenDecision::Denied(TcpOpenDenial::InvalidDestination)); + return None; + } + }, + }; if let Err(denial) = hydrate_destination_plan(&mut decision, backend_host_gateway, trusted_host_gateway) { warn!(%destination, reason = %denial.reason, "Denied staged transparent destination"); emit_staged_transparent_denial( destination, + mapped_host, &binary_identity, &denial.reason, "transparent_tcp_destination_denied", @@ -724,26 +776,7 @@ async fn preauthorize_transparent_open( let _ = completion.send(TcpOpenDecision::Denied(TcpOpenDenial::InvalidDestination)); return None; } - if let Some(mapping) = policy_dns_store.and_then(|store| { - store - .lookup( - destination.ip(), - destination.port(), - opa_engine.current_generation(), - std::time::Instant::now(), - ) - .ok() - }) { - let Ok(plan) = build_pinned_validation_plan(mapping.pinned_addresses()) else { - emit_staged_transparent_denial( - destination, - &binary_identity, - "policy DNS produced an invalid pinned destination", - "transparent_tcp_destination_denied", - ); - let _ = completion.send(TcpOpenDecision::Denied(TcpOpenDenial::InvalidDestination)); - return None; - }; + if let Some(plan) = pinned_plan { decision.endpoint.destination = Some(plan); } let plan = decision @@ -764,6 +797,7 @@ async fn preauthorize_transparent_open( warn!(%destination, reason = %denial.reason, "Denied staged transparent destination"); emit_staged_transparent_denial( destination, + mapped_host, &binary_identity, &denial.reason, "transparent_tcp_destination_denied", @@ -788,10 +822,27 @@ async fn preauthorize_transparent_open( fn emit_staged_transparent_denial( destination: SocketAddr, + mapped_host: Option<&str>, identity: &Result, reason: &str, status_detail: &'static str, ) { + ocsf_emit!(build_staged_transparent_denial_event( + destination, + mapped_host, + identity, + reason, + status_detail, + )); +} + +fn build_staged_transparent_denial_event( + destination: SocketAddr, + mapped_host: Option<&str>, + identity: &Result, + reason: &str, + status_detail: &'static str, +) -> openshell_ocsf::OcsfEvent { let (binary, ancestors, cmdline) = identity.as_ref().map_or_else( |_| ("-".to_string(), "-".to_string(), "-".to_string()), |identity| { @@ -812,43 +863,87 @@ fn emit_staged_transparent_denial( ) }, ); - ocsf_emit!( - NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) - .activity(ActivityId::Open) - .action(ActionId::Denied) - .disposition(DispositionId::Blocked) - .severity(SeverityId::Medium) - .status(StatusId::Failure) - .dst_endpoint(Endpoint::from_ip(destination.ip(), destination.port())) - .actor_process(Process::from_bypass(&binary, "-", &ancestors).with_cmd_line(&cmdline)) - .message(format!("Transparent TCP denied before relay: {reason}")) - .status_detail(status_detail) - .build() - ); + NetworkActivityBuilder::new(openshell_ocsf::ctx::ctx()) + .activity(ActivityId::Open) + .action(ActionId::Denied) + .disposition(DispositionId::Blocked) + .severity(SeverityId::Medium) + .status(StatusId::Failure) + .dst_endpoint(transparent_destination_endpoint(destination, mapped_host)) + .actor_process(Process::from_bypass(&binary, "-", &ancestors).with_cmd_line(&cmdline)) + .message(format!("Transparent TCP denied before relay: {reason}")) + .status_detail(status_detail) + .build() +} + +/// Name the policy DNS host an operator recognizes while keeping the +/// synthetic address the workload dialed. +fn transparent_destination_endpoint( + destination: SocketAddr, + mapped_host: Option<&str>, +) -> Endpoint { + let mut endpoint = Endpoint::from_ip(destination.ip(), destination.port()); + endpoint.domain = mapped_host.map(str::to_string); + endpoint +} + +/// Logical destination recovered for a staged transparent TCP open. +struct TransparentTarget { + host: String, + /// Policy DNS correlation for a synthetic destination, read at the + /// generation current when the open arrived. `None` outside the pools. + intent: Option, } -fn transparent_destination_host( +fn resolve_transparent_target( destination: SocketAddr, policy_dns_store: Option<&Arc>, opa_engine: &OpaEngine, -) -> Result { +) -> Result { let Some(store) = policy_dns_store else { - return Ok(destination.ip().to_string()); + return Ok(TransparentTarget { + host: destination.ip().to_string(), + intent: None, + }); }; - match store.lookup( + match store.lookup_intent( destination.ip(), destination.port(), opa_engine.current_generation(), std::time::Instant::now(), ) { - Ok(mapping) => Ok(mapping.record.normalized_name.as_str().to_string()), - Err(MappingLookupError::Missing) => Ok(destination.ip().to_string()), + Ok(mapping) => Ok(TransparentTarget { + host: mapping.record.normalized_name.as_str().to_string(), + intent: Some(mapping), + }), + Err(MappingLookupError::Missing) => Ok(TransparentTarget { + host: destination.ip().to_string(), + intent: None, + }), Err(error) => Err(miette::miette!( "transparent destination mapping is unavailable: {error}" )), } } +/// Re-acquire a policy DNS mapping at the generation that produced +/// `decision` and pin its addresses. Observation records, stale mappings, and +/// unmapped ports fail closed. +fn pinned_transparent_plan( + store: &ResolvedEndpointStore, + destination: SocketAddr, + decision: &EgressDecision, +) -> std::result::Result { + let mapping = store.lookup( + destination.ip(), + destination.port(), + decision.policy_generation, + std::time::Instant::now(), + )?; + build_pinned_validation_plan(mapping.pinned_addresses()) + .map_err(|_| MappingLookupError::InvalidMapping) +} + fn valid_policy_local_request(method: &str, target: &str, request_headers: &str) -> bool { if method == "CONNECT" || !target.starts_with('/') { return false; @@ -982,7 +1077,7 @@ async fn handle_transparent_tcp_connection( let workload_addr = client.peer_addr().into_diagnostic()?; let original = original_destination(&client).into_diagnostic()?; let current_generation = opa_engine.current_generation(); - let mapping = match store.lookup( + let mapping = match store.lookup_intent( original.ip(), original.port(), current_generation, @@ -990,7 +1085,7 @@ async fn handle_transparent_tcp_connection( ) { Ok(mapping) => mapping, Err(error) => { - emit_transparent_mapping_denial(workload_addr, original, error); + emit_transparent_mapping_denial(workload_addr, original, None, error); emit_activity(&activity_tx, true, "transparent_tcp_mapping"); return Ok(()); } @@ -1026,6 +1121,17 @@ async fn handle_transparent_tcp_connection( return Ok(()); } + if mapping.record.is_observation() { + emit_transparent_mapping_denial( + workload_addr, + original, + Some(&host), + MappingLookupError::EndpointMismatch, + ); + emit_activity(&activity_tx, true, "transparent_tcp_mapping"); + return Ok(()); + } + // Authorization may race a policy reload. Re-pin the exact generation // that produced the decision, then reacquire the DNS mapping against that // generation before correlating endpoint identity or constructing a @@ -1034,7 +1140,12 @@ async fn handle_transparent_tcp_connection( let Ok(generation_guard) = relay::pin_policy_generation(&opa_engine, decision.policy_generation) else { - emit_transparent_mapping_denial(workload_addr, original, MappingLookupError::StalePolicy); + emit_transparent_mapping_denial( + workload_addr, + original, + Some(&host), + MappingLookupError::StalePolicy, + ); emit_activity(&activity_tx, true, "transparent_tcp_mapping"); return Ok(()); }; @@ -1046,7 +1157,7 @@ async fn handle_transparent_tcp_connection( ) { Ok(mapping) => mapping, Err(error) => { - emit_transparent_mapping_denial(workload_addr, original, error); + emit_transparent_mapping_denial(workload_addr, original, Some(&host), error); emit_activity(&activity_tx, true, "transparent_tcp_mapping"); return Ok(()); } @@ -1315,6 +1426,7 @@ fn original_destination(stream: &TcpStream) -> std::io::Result { fn emit_transparent_mapping_denial( workload: SocketAddr, original: SocketAddr, + mapped_host: Option<&str>, error: MappingLookupError, ) { let detail = match error { @@ -1333,7 +1445,7 @@ fn emit_transparent_mapping_denial( .disposition(DispositionId::Blocked) .severity(SeverityId::Medium) .status(StatusId::Failure) - .dst_endpoint(Endpoint::from_ip(original.ip(), original.port())) + .dst_endpoint(transparent_destination_endpoint(original, mapped_host)) .src_endpoint_addr(workload.ip(), workload.port()) .message(format!("Transparent TCP denied: {error}")) .status_detail(detail) @@ -2315,7 +2427,8 @@ async fn handle_mediated_connection( (None, None) } else { let host = - transparent_destination_host(destination, policy_dns_store.as_ref(), &opa_engine)?; + resolve_transparent_target(destination, policy_dns_store.as_ref(), &opa_engine)? + .host; let (decision, connector) = transparent .authorization .map_or((None, None), |(decision, connector)| { @@ -7103,6 +7216,359 @@ process: assert!(denial_rx.try_recv().is_err(), "exactly one mapper event"); } + const POLICY_DNS_OPEN_POLICY: &str = r" +network_policies: + database: + name: database + endpoints: + - { host: db.example, port: 5432, protocol: tcp } + binaries: [{ path: /usr/bin/curl }] + pypi: + name: pypi + endpoints: + - { host: pypi.org, port: 80 } + binaries: [{ path: /usr/bin/curl }] +filesystem_policy: { include_workdir: true, read_only: [], read_write: [] } +landlock: { compatibility: best_effort } +process: { run_as_user: sandbox, run_as_group: sandbox } +"; + + /// Like the production pools, skip the reserved sandbox-local address. + fn policy_dns_test_store() -> Arc { + Arc::new(ResolvedEndpointStore::new( + crate::policy_dns::StoreConfig::new( + crate::policy_dns::SyntheticPools::new( + Ipv4Addr::new(198, 18, 0, 2)..=Ipv4Addr::new(198, 18, 0, 9), + "fd00:1::1".parse::().unwrap() + ..="fd00:1::8".parse::().unwrap(), + ) + .unwrap(), + 16, + ) + .unwrap(), + )) + } + + fn publish_mapping( + store: &ResolvedEndpointStore, + name: &str, + policy_name: &str, + port: u16, + generation: u64, + ) -> crate::policy_dns::ResolvedEndpointRecord { + store + .publish( + crate::policy_dns::PublishRequest { + normalized_name: crate::policy_dns::NormalizedName::parse(name).unwrap(), + family: crate::policy_dns::AddressFamily::Ipv4, + allocation_identity: [1; 32], + policy_generation: generation, + ttl: std::time::Duration::from_secs(30), + contracts: vec![crate::policy_dns::ResolvedPortContract { + endpoint_id: PolicyEndpointId { + policy_name: policy_name.to_string(), + endpoint_index: 0, + }, + port, + destination_plan: DestinationValidationPlan { + address_authorization: + destination::AddressAuthorization::ExactDeclaredHost, + }, + pinned_addresses: vec!["203.0.113.8".parse().unwrap()], + }], + }, + generation, + std::time::Instant::now(), + ) + .unwrap() + } + + fn publish_observation( + store: &ResolvedEndpointStore, + name: &str, + generation: u64, + ) -> crate::policy_dns::ResolvedEndpointRecord { + store + .publish_observation( + crate::policy_dns::NormalizedName::parse(name).unwrap(), + crate::policy_dns::AddressFamily::Ipv4, + generation, + generation, + std::time::Instant::now(), + ) + .unwrap() + } + + fn staged_curl_open( + destination: SocketAddr, + generation: u64, + ) -> ( + PendingTcpOpen, + tokio::sync::oneshot::Receiver, + ) { + let (stream, _peer) = tokio::io::duplex(64); + let (decision, completion) = tokio::sync::oneshot::channel(); + ( + PendingTcpOpen { + stream: Box::new(stream), + binary_identity: Ok(ContractBinaryIdentity { + executable: ContractExecutableIdentity { + path: PathBuf::from("/usr/bin/curl"), + digest: Some("00".repeat(32).parse().unwrap()), + }, + ancestors: Vec::new(), + cmdline_paths: Vec::new(), + }), + destination, + socket: openshell_isolation_interface::contract::NetworkSocketMetadata { + socket_cookie: 7, + nonblocking: false, + process_generation: 1, + }, + policy_generation: generation, + timing: MediationTiming::default(), + decision, + }, + completion, + ) + } + + #[tokio::test] + async fn staged_transparent_open_dials_only_pinned_policy_dns_addresses() { + let engine = OpaEngine::from_strings( + include_str!("../data/sandbox-policy.rego"), + POLICY_DNS_OPEN_POLICY, + ) + .unwrap(); + let store = policy_dns_test_store(); + let record = publish_mapping( + &store, + "db.example", + "database", + 5432, + engine.current_generation(), + ); + let (open, completion) = staged_curl_open( + SocketAddr::new(record.synthetic_address, 5432), + engine.current_generation(), + ); + + let (_, _, _, transparent) = preauthorize_transparent_open( + open, + Some(&store), + &engine, + &BinaryIdentityCache::new(), + None, + None, + false, + None, + ) + .await + .expect("policy-backed mapping is admitted"); + + assert_eq!(completion.await.unwrap(), TcpOpenDecision::RelayReady); + let (_, connector) = transparent + .and_then(|open| open.authorization) + .expect("transparent authorization"); + assert_eq!(connector.addrs(), &["203.0.113.8:5432".parse().unwrap()]); + } + + #[tokio::test] + async fn staged_transparent_open_proposes_a_denied_observation_hostname() { + let engine = OpaEngine::from_strings( + include_str!("../data/sandbox-policy.rego"), + POLICY_DNS_OPEN_POLICY, + ) + .unwrap(); + let store = policy_dns_test_store(); + let record = publish_observation(&store, "unknown.example", engine.current_generation()); + let (open, completion) = staged_curl_open( + SocketAddr::new(record.synthetic_address, 443), + engine.current_generation(), + ); + let (denial_tx, mut denial_rx) = mpsc::unbounded_channel(); + + assert!( + preauthorize_transparent_open( + open, + Some(&store), + &engine, + &BinaryIdentityCache::new(), + None, + None, + false, + Some(&denial_tx), + ) + .await + .is_none() + ); + + assert_eq!( + completion.await.unwrap(), + TcpOpenDecision::Denied(TcpOpenDenial::PolicyDenied) + ); + let event = denial_rx.try_recv().expect("denial is sent to the mapper"); + assert_eq!(event.host, "unknown.example"); + assert_eq!(event.port, 443); + assert_eq!(event.binary, "/usr/bin/curl"); + } + + #[tokio::test] + async fn staged_transparent_open_never_relays_an_observation_that_policy_allows() { + let engine = OpaEngine::from_strings( + include_str!("../data/sandbox-policy.rego"), + POLICY_DNS_OPEN_POLICY, + ) + .unwrap(); + let identity_cache = BinaryIdentityCache::new(); + let (denial_tx, mut denial_rx) = mpsc::unbounded_channel(); + let preauthorize = |store: Arc, destination: SocketAddr| { + let (open, completion) = staged_curl_open(destination, engine.current_generation()); + let engine = &engine; + let identity_cache = &identity_cache; + let denial_tx = &denial_tx; + async move { + let admitted = preauthorize_transparent_open( + open, + Some(&store), + engine, + identity_cache, + None, + None, + false, + Some(denial_tx), + ) + .await + .is_some(); + (admitted, completion.await.unwrap()) + } + }; + + // The same policy admits curl to pypi.org:80 through a mapping with + // an endpoint contract. + let mapped_store = policy_dns_test_store(); + let mapped = publish_mapping( + &mapped_store, + "pypi.org", + "pypi", + 80, + engine.current_generation(), + ); + assert_eq!( + preauthorize(mapped_store, SocketAddr::new(mapped.synthetic_address, 80)).await, + (true, TcpOpenDecision::RelayReady) + ); + + let observed_store = policy_dns_test_store(); + let observed = + publish_observation(&observed_store, "pypi.org", engine.current_generation()); + assert_eq!( + preauthorize( + observed_store, + SocketAddr::new(observed.synthetic_address, 80) + ) + .await, + ( + false, + TcpOpenDecision::Denied(TcpOpenDenial::InvalidDestination) + ) + ); + assert!( + denial_rx.try_recv().is_err(), + "an allowed intent is not a proposal" + ); + } + + #[test] + fn pinned_plan_requires_the_mapping_of_the_deciding_generation() { + let engine = OpaEngine::from_strings( + include_str!("../data/sandbox-policy.rego"), + POLICY_DNS_OPEN_POLICY, + ) + .unwrap(); + let store = policy_dns_test_store(); + let mapped = publish_mapping( + &store, + "db.example", + "database", + 5432, + engine.current_generation(), + ); + let observed = publish_observation(&store, "pypi.org", engine.current_generation()); + let (open, _) = staged_curl_open( + SocketAddr::new(mapped.synthetic_address, 5432), + engine.current_generation(), + ); + let decide = |host: &str, port| { + authorize_supplied_identity( + &engine, + &BinaryIdentityCache::new(), + EgressIntent::connect(host.to_string(), port), + &open.binary_identity, + ) + }; + let mapped_destination = SocketAddr::new(mapped.synthetic_address, 5432); + + pinned_transparent_plan(&store, mapped_destination, &decide("db.example", 5432)) + .expect("the deciding generation pins its addresses"); + assert!(matches!( + pinned_transparent_plan( + &store, + SocketAddr::new(observed.synthetic_address, 80), + &decide("pypi.org", 80), + ), + Err(MappingLookupError::PortMismatch) + )); + + // A reload between DNS correlation and authorization yields a decision + // from a newer generation than the DNS answer. It must not fall back + // to resolving the name again. + engine + .reload( + include_str!("../data/sandbox-policy.rego"), + POLICY_DNS_OPEN_POLICY, + ) + .unwrap(); + let reloaded = decide("db.example", 5432); + assert!(matches!(reloaded.action, NetworkAction::Allow { .. })); + assert!(matches!( + pinned_transparent_plan(&store, mapped_destination, &reloaded), + Err(MappingLookupError::StalePolicy) + )); + } + + #[test] + fn staged_denial_names_the_mapped_host_and_keeps_the_synthetic_address() { + let (open, _) = staged_curl_open("198.18.0.2:80".parse().unwrap(), 1); + let mapped = build_staged_transparent_denial_event( + open.destination, + Some("blocked.invalid"), + &open.binary_identity, + "endpoint blocked.invalid:80 is not allowed by any policy", + "transparent_tcp_policy_denied", + ); + assert_eq!( + mapped.format_shorthand(), + "NET:OPEN [MED] DENIED /usr/bin/curl(0) -> blocked.invalid:80 [reason:transparent_tcp_policy_denied]" + ); + assert_eq!( + serde_json::to_value(mapped).unwrap()["dst_endpoint"], + serde_json::json!({"domain": "blocked.invalid", "ip": "198.18.0.2", "port": 80}) + ); + + let unmapped = build_staged_transparent_denial_event( + "203.0.113.7:443".parse().unwrap(), + None, + &open.binary_identity, + "endpoint 203.0.113.7:443 is not allowed by any policy", + "transparent_tcp_policy_denied", + ); + assert_eq!( + serde_json::to_value(unmapped).unwrap()["dst_endpoint"], + serde_json::json!({"ip": "203.0.113.7", "port": 443}) + ); + } + #[tokio::test] async fn staged_policy_local_open_reaches_the_sandbox_scoped_api() { let engine = Arc::new( diff --git a/docs/sandboxes/policies.mdx b/docs/sandboxes/policies.mdx index 37eb13abb1..86d3b541e1 100644 --- a/docs/sandboxes/policies.mdx +++ b/docs/sandboxes/policies.mdx @@ -153,11 +153,11 @@ network_policies: - path: /usr/bin/psql ``` -OpenShell answers DNS only for hostnames eligible under an active `protocol: tcp` endpoint. It validates upstream answers against destination and SSRF controls, returns a supervisor-owned synthetic address, and records the validated real addresses. When the application connects to the synthetic address, OpenShell recovers the hostname and port, evaluates the calling process against the current policy generation, and dials only an address pinned by that DNS result. +OpenShell uses policy DNS for hostnames eligible under TCP-carried endpoints, including `protocol: tcp` and inspected HTTP-based protocols such as `protocol: rest`. It validates upstream answers against destination and SSRF controls, returns a supervisor-owned synthetic address, and records the validated real addresses. A direct HTTP or HTTPS connection to that address enters transparent TCP capture; the supervisor then applies the endpoint's L4 or L7 policy before dialing a pinned address. Hostnames absent from policy also receive a contract-free synthetic observation address so a denied connection can produce a mechanistic proposal. Sandbox resolvers apply no DNS search domains, so a workload must request a hostname exactly as the policy names it. Clients that use the explicit HTTP proxy send the target hostname in the proxy request and do not need a policy DNS address for that target. Treat that hostname as a connection-routing constraint, not an application-authority boundary. OpenShell does not inspect TLS SNI, HTTP `Host`, or another protocol-level destination inside a `protocol: tcp` stream. If the approved hostname reaches compatible shared infrastructure, a client may be able to select another tenant, virtual host, or service behind the same front door. Use native TCP only when the trust boundary includes every destination that the shared infrastructure can expose; use an inspected protocol when application authority must remain constrained. -Applications must honor the returned DNS TTL and resolve the hostname again before reconnecting after that TTL expires. A client that caches the synthetic address indefinitely can receive a connection failure after the mapping expires. Docker and Podman currently advertise only IPv4 egress for this feature, so OpenShell returns an empty successful answer for AAAA queries and lets dual-stack clients use the working A record. +Applications must honor the returned DNS TTL and resolve the hostname again before reconnecting after that TTL expires. A client that caches the synthetic address indefinitely can receive a connection failure after the mapping expires. Synthetic address identities are not reused within a supervisor lifetime, even after the active DNS record expires; the production pool has up to 512 IPv4 addresses and 512 IPv6 addresses, shared across policy-backed names and unknown-host observations. Observations can hold at most a quarter of each pool. Docker and Podman currently advertise only IPv4 egress for this feature, so OpenShell returns an empty successful answer for AAAA queries and lets dual-stack clients use the working A record. DNS resolution does not authorize a connection by itself. Unknown names, wrong ports, stale mappings, disallowed destination addresses, and binaries outside the matching policy fail closed. Applications cannot inherit access by connecting directly to a real IP returned by an upstream resolver. diff --git a/docs/sandboxes/policy-advisor.mdx b/docs/sandboxes/policy-advisor.mdx index 83413e21e1..56e71f2bed 100644 --- a/docs/sandboxes/policy-advisor.mdx +++ b/docs/sandboxes/policy-advisor.mdx @@ -131,6 +131,8 @@ OpenShell has two proposal paths: | Mechanistic mapper | Aggregated denial summaries from the sandbox. | Groups by host, port, and binary. If L7 request samples are available, it can draft REST method and path rules. Otherwise it drafts an L4 endpoint. | | Agent-authored proposal | The in-sandbox agent, using `policy.local`. | Usually a REST `addRule` with exact host, port, binary, method, and path from the structured denial. It can also omit `protocol` for endpoint-only access through the explicit proxy. | +For a new hostname absent from policy, policy DNS does not query an upstream resolver. It gives the process a synthetic address with a 30-second observation record that cannot reach the destination. The resulting TCP attempt records the actual hostname, port, and verified binary for a mechanistic proposal, even when the agent proposal surface is disabled. For example, `curl http://pypi.org/` can create a pending `pypi.org:80` rule for `/usr/bin/curl`. The first request is still denied. After approval and policy reload, retry the request so DNS can resolve the now-authorized endpoint. Each unknown lookup still appears in the OCSF log as a `policy_dns_ineligible` DNS denial. Policy DNS refuses reserved loopback, metadata, and host-gateway names outright, and refuses every unknown name while the sandbox policy is quarantined. Observations can use at most a quarter of each synthetic address pool, up to 128 unique unknown names per address family for the supervisor's lifetime. After that, policy DNS refuses further unknown names, and they do not produce proposals until the supervisor restarts. + ### How proposal provenance works OpenShell tracks whether an endpoint and binary came from policy advisor. This is internal provenance; it is not a policy YAML field that authors set. Think of each marker as answering “who introduced this identity?” rather than “what traffic does this allow?” diff --git a/e2e/policy-advisor/README.md b/e2e/policy-advisor/README.md index 9e29d1ec5e..a035a9547e 100644 --- a/e2e/policy-advisor/README.md +++ b/e2e/policy-advisor/README.md @@ -55,18 +55,19 @@ contents write on the repository. The test auto-resolves the token from ## Conformance coverage -The `mechanistic-proposal` and `policy-local` conformance scenarios check draft -generation and use `policy.local` to inspect policy, submit a narrow permission -request, and read the resulting proposal. Run them against a configured gateway +The `mechanistic-proposal` and `new-hostname-proposal` conformance scenarios +check draft generation for a denied IP address and for a hostname absent from +policy. The `policy-local` scenario uses `policy.local` to inspect policy, submit +a narrow permission request, and read the resulting proposal. Run them against a configured gateway with `--openshell-bin` pointing to the CLI under test: ```bash -openshell-conformance run mechanistic-proposal policy-local --openshell-bin target/debug/openshell +openshell-conformance run mechanistic-proposal new-hostname-proposal policy-local --openshell-bin target/debug/openshell ``` Run `openshell-conformance list` to see all scenario names. A manual `Integration Tests` workflow run can select the `policy-advisor` testsuite to -run only these two scenarios against an installed candidate. Set +run only these three scenarios against an installed candidate. Set `artifact-run-id` to the candidate build's workflow run ID and `test-matrix` to: ```json diff --git a/tests/suites/conformance/cli/tests/policy_advisor.rs b/tests/suites/conformance/cli/tests/policy_advisor.rs index 56bf18ba87..7da2186dd7 100644 --- a/tests/suites/conformance/cli/tests/policy_advisor.rs +++ b/tests/suites/conformance/cli/tests/policy_advisor.rs @@ -4,7 +4,8 @@ //! Installed-artifact policy advisor conformance. use openshell_conformance::{ - MECHANISTIC_PROPOSAL_SCENARIO, OpenShellRunner, POLICY_LOCAL_SCENARIO, Scenario, + MECHANISTIC_PROPOSAL_SCENARIO, NEW_HOSTNAME_PROPOSAL_SCENARIO, OpenShellRunner, + POLICY_LOCAL_SCENARIO, Scenario, }; async fn run(scenario: &'static Scenario) { @@ -25,6 +26,11 @@ async fn mechanistic_proposal() { run(&MECHANISTIC_PROPOSAL_SCENARIO).await; } +#[tokio::test] +async fn new_hostname_proposal() { + run(&NEW_HOSTNAME_PROPOSAL_SCENARIO).await; +} + #[tokio::test] async fn policy_local() { run(&POLICY_LOCAL_SCENARIO).await;