From 3b0a2606ee269e166f43f2adf505d0aacabaec8f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Beno=C3=AEt=20CORTIER?= Date: Tue, 18 Aug 2026 03:29:36 +0900 Subject: [PATCH 1/2] feat(agent): expose active package policy Expose the validated active package-broker policy through the shared authenticated GET /v1/policy route. Return a structured unavailable error without leaking policy source or file-security details. This requires now-policy-api and now-policy-server-template 0.4.0 from Devolutions/now-libraries#93 before the change can ship. Issue: Devolutions/now-libraries#93 Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- crates/now-package-broker/src/auth.rs | 19 +- crates/now-package-broker/src/server/mod.rs | 238 ++++++++++++++++++-- 2 files changed, 238 insertions(+), 19 deletions(-) diff --git a/crates/now-package-broker/src/auth.rs b/crates/now-package-broker/src/auth.rs index bc770d149..65b5515d4 100644 --- a/crates/now-package-broker/src/auth.rs +++ b/crates/now-package-broker/src/auth.rs @@ -31,6 +31,10 @@ impl PipeClient { /// unauthenticated work a connection flood can trigger. pub(crate) fn from_connected_pipe(server: &NamedPipeServer) -> anyhow::Result { let process_id = connected_pipe_client_process_id(server).context("failed to query pipe client process id")?; + Self::from_process_id(process_id) + } + + fn from_process_id(process_id: u32) -> anyhow::Result { let process = Process::get_by_pid(process_id, PROCESS_QUERY_LIMITED_INFORMATION) .with_context(|| format!("failed to open pipe client process {process_id}"))?; let executable_path = process @@ -50,6 +54,11 @@ impl PipeClient { }) } + #[cfg(test)] + pub(crate) fn from_current_process() -> anyhow::Result { + Self::from_process_id(std::process::id()) + } + /// Security identifier of the authenticated pipe client user, captured at connect. pub(crate) fn user_sid(&self) -> &Sid { &self.user_sid @@ -61,7 +70,7 @@ impl PipeClient { skip_signature_validation: bool, ) -> anyhow::Result<()> { self.validate_client_context(&request.client)?; - self.validate_signature(skip_signature_validation) + self.validate_connection(skip_signature_validation) } pub(crate) fn validate_status_request( @@ -70,7 +79,7 @@ impl PipeClient { skip_signature_validation: bool, ) -> anyhow::Result<()> { self.validate_client_context(&request.client)?; - self.validate_signature(skip_signature_validation) + self.validate_connection(skip_signature_validation) } pub(crate) fn validate_cancel_request( @@ -79,7 +88,7 @@ impl PipeClient { skip_signature_validation: bool, ) -> anyhow::Result<()> { self.validate_client_context(&request.client)?; - self.validate_signature(skip_signature_validation) + self.validate_connection(skip_signature_validation) } fn validate_client_context(&self, client: &ClientContext) -> anyhow::Result<()> { @@ -87,7 +96,7 @@ impl PipeClient { self.validate_executable_path(&client.client_executable_path) } - fn validate_signature(&self, skip_signature_validation: bool) -> anyhow::Result<()> { + pub(crate) fn validate_connection(&self, skip_signature_validation: bool) -> anyhow::Result<()> { if signature_validation_skipped(skip_signature_validation) { warn!("DEBUG MODE: Skipping package broker client signature validation"); return Ok(()); @@ -364,7 +373,7 @@ mod tests { user_sid: client_user_sid(), }; - assert!(client.validate_signature(true).is_err()); + assert!(client.validate_connection(true).is_err()); } } diff --git a/crates/now-package-broker/src/server/mod.rs b/crates/now-package-broker/src/server/mod.rs index 6f1e65b20..bd5640097 100644 --- a/crates/now-package-broker/src/server/mod.rs +++ b/crates/now-package-broker/src/server/mod.rs @@ -11,8 +11,8 @@ use now_policy_api::{ CancelRequest, CancelResponse, CancelResponseKind, CapabilitiesResponse, CapabilitiesResponseKind, Decision, DecisionInfo, Elevation, ErrorCode, ErrorResponse, EvaluationResponse, EvaluationResponseKind, ExecutionResponse, ExecutionResponseKind, HealthResponse, HealthResponseKind, HealthStatus, ManagerCapability, ManagerName, - OperationStatus, OperationSubmission, PackageRequest, Scope, StatusRequest, StatusResponse, StatusResponseKind, - Transport, + OperationStatus, OperationSubmission, PackageRequest, PolicyResponse, PolicyResponseKind, Scope, StatusRequest, + StatusResponse, StatusResponseKind, Transport, }; use now_policy_server_template::{MAX_REQUEST_BODY_BYTES, PackageBrokerServer, SharedPackageBrokerServer}; use tracing::{info, trace, warn}; @@ -115,6 +115,17 @@ impl PackageBrokerServer for BrokerConnection { self.state.capabilities(self.client.user_sid()).await } + async fn policy(&self) -> Result { + self.client + .validate_connection(self.state.skip_signature_validation) + .map_err(|error| { + warn!(error = format!("{error:#}"), "Rejected package broker policy request"); + error_response(ErrorCode::Unauthorized, "pipe client authentication failed") + })?; + + self.state.policy_response() + } + async fn evaluate(&self, request: PackageRequest) -> Result { self.client .validate_request(&request, self.state.skip_signature_validation) @@ -163,6 +174,25 @@ impl PackageBrokerServer for BrokerConnection { } impl BrokerState { + fn active_policy(&self) -> Result, ErrorResponse> { + let guard = self.policy.read().expect("policy lock poisoned"); + guard + .as_ref() + .map(Arc::clone) + .ok_or_else(|| error_response(ErrorCode::BrokerPaused, "policy file is unavailable or corrupted")) + } + + fn policy_response(&self) -> Result { + let policy = self.active_policy()?; + + Ok(PolicyResponse { + response_kind: PolicyResponseKind, + response_version: api_version(), + server: server_context(), + policy: (*policy).clone(), + }) + } + async fn health(&self) -> HealthResponse { let policy_guard = self.policy.read().expect("policy lock poisoned"); let (status, policy_id) = match policy_guard.as_ref() { @@ -419,18 +449,7 @@ impl BrokerState { } let received_at = Utc::now(); - let policy = { - let guard = self.policy.read().expect("policy lock poisoned"); - match guard.as_ref() { - Some(policy) => Arc::clone(policy), - None => { - return Err(error_response( - ErrorCode::BrokerPaused, - "policy file is unavailable or corrupted", - )); - } - } - }; + let policy = self.active_policy()?; if let Some(reason) = policy_validity_failure(&policy, received_at) { warn!(%reason, "Rejecting request: policy outside validity window"); @@ -521,12 +540,15 @@ mod tests { use std::sync::atomic::{AtomicUsize, Ordering}; + use axum::body::{Body, to_bytes}; + use axum::http::{Method, Request, StatusCode}; use chrono::Utc; use now_policy::{ PackageBrokerPolicy, PolicyEnforcement, PolicyMetadata, PolicySchemaUri, ResourceId, RulePrecedence, SemanticVersion, }; use now_policy_api as api; + use tower_service::Service as _; use super::*; use crate::executor::{ExecutionOutput, OperationCanceled, ProcessStartedCallback}; @@ -601,6 +623,194 @@ mod tests { } } + fn shared_state(policy: Option) -> Arc { + let mut state = state(); + state.policy = RwLock::new(policy.map(Arc::new)); + Arc::new(state) + } + + async fn route_request(state: Arc, method: Method, uri: &str) -> axum::response::Response { + let client = PipeClient::from_current_process().expect("capture current test process"); + let mut router = build_router_for_client(state, client); + router + .call( + Request::builder() + .method(method) + .uri(uri) + .body(Body::empty()) + .expect("valid test request"), + ) + .await + .expect("router is infallible") + } + + async fn response_json(response: axum::response::Response) -> serde_json::Value { + let body = to_bytes(response.into_body(), usize::MAX) + .await + .expect("read response body"); + serde_json::from_slice(&body).expect("response is valid JSON") + } + + #[cfg(feature = "dev-skip-broker-signature")] + #[tokio::test] + async fn policy_route_serializes_active_policy_with_empty_rules() { + let expected = permissive_policy(); + let response = route_request(shared_state(Some(expected.clone())), Method::GET, "/v1/policy").await; + + assert_eq!(response.status(), StatusCode::OK); + assert_eq!(response.headers().get("content-type").unwrap(), "application/json"); + + let response: PolicyResponse = + serde_json::from_value(response_json(response).await).expect("deserialize policy response"); + assert_eq!(response.response_kind, PolicyResponseKind); + assert_eq!(&*response.response_version, api::API_VERSION_STR); + assert_eq!(response.server.transport, Transport::HttpNamedPipe); + assert_eq!( + serde_json::to_value(response.policy).unwrap(), + serde_json::to_value(expected).unwrap() + ); + } + + #[cfg(feature = "dev-skip-broker-signature")] + #[tokio::test] + async fn policy_route_serializes_full_policy_matches_and_constraints() { + let expected = + now_policy::schema::parse_policy_json(include_str!("../assets/samples/corporate-allowlist.policy.json")) + .expect("sample policy is valid"); + let response = route_request(shared_state(Some(expected.clone())), Method::GET, "/v1/policy").await; + + assert_eq!(response.status(), StatusCode::OK); + + let response: PolicyResponse = + serde_json::from_value(response_json(response).await).expect("deserialize policy response"); + assert_eq!( + serde_json::to_value(response.policy).unwrap(), + serde_json::to_value(expected).unwrap() + ); + } + + #[cfg(feature = "dev-skip-broker-signature")] + #[tokio::test] + async fn policy_route_returns_structured_service_unavailable_without_active_policy() { + let response = route_request(shared_state(None), Method::GET, "/v1/policy").await; + + assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE); + + let body = response_json(response).await; + let error: ErrorResponse = serde_json::from_value(body.clone()).expect("deserialize error response"); + assert_eq!(error.code, ErrorCode::BrokerPaused); + assert_eq!(error.message, "policy file is unavailable or corrupted"); + assert!(error.details.is_empty()); + assert!(body.get("Policy").is_none()); + } + + #[cfg(not(feature = "dev-skip-broker-signature"))] + #[tokio::test] + async fn policy_route_rejects_unsigned_client() { + let response = route_request(shared_state(Some(permissive_policy())), Method::GET, "/v1/policy").await; + + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + + let body = response_json(response).await; + let error: ErrorResponse = serde_json::from_value(body.clone()).expect("deserialize error response"); + assert_eq!(error.code, ErrorCode::Unauthorized); + assert_eq!(error.message, "pipe client authentication failed"); + assert!(body.get("Policy").is_none()); + } + + #[cfg(feature = "dev-skip-broker-signature")] + #[tokio::test] + async fn policy_route_preserves_existing_routes_and_method_restrictions() { + let state = shared_state(Some(permissive_policy())); + + for uri in ["/v1/health", "/v1/capabilities"] { + let response = route_request(Arc::clone(&state), Method::GET, uri).await; + assert_eq!(response.status(), StatusCode::OK, "unexpected status for {uri}"); + } + + let response = route_request(Arc::clone(&state), Method::HEAD, "/v1/policy").await; + assert_eq!(response.status(), StatusCode::OK); + assert!( + to_bytes(response.into_body(), usize::MAX) + .await + .expect("read HEAD response") + .is_empty() + ); + + for method in [ + Method::POST, + Method::PUT, + Method::PATCH, + Method::DELETE, + Method::OPTIONS, + Method::TRACE, + Method::CONNECT, + ] { + let response = route_request(Arc::clone(&state), method.clone(), "/v1/policy").await; + assert_eq!( + response.status(), + StatusCode::METHOD_NOT_ALLOWED, + "unexpected status for {method}" + ); + } + + let response = route_request(state, Method::GET, "/v1/not-a-route").await; + assert_eq!(response.status(), StatusCode::NOT_FOUND); + } + + #[test] + fn concurrent_policy_replacement_returns_only_complete_snapshots() { + let policy_a = permissive_policy(); + let mut policy_b = + now_policy::schema::parse_policy_json(include_str!("../assets/samples/corporate-allowlist.policy.json")) + .expect("sample policy is valid"); + policy_b.metadata.id = ResourceId::from("replacement-policy"); + policy_b.metadata.revision = 42; + + let current_policy_json = serde_json::to_value(&policy_a).unwrap(); + let replacement_policy_json = serde_json::to_value(&policy_b).unwrap(); + let policy_a = Arc::new(policy_a); + let policy_b = Arc::new(policy_b); + let state = shared_state(None); + *state.policy.write().expect("policy lock") = Some(Arc::clone(&policy_a)); + + const READER_COUNT: usize = 4; + const ITERATIONS: usize = 1_000; + let barrier = Arc::new(std::sync::Barrier::new(READER_COUNT + 1)); + + std::thread::scope(|scope| { + for _ in 0..READER_COUNT { + let state = Arc::clone(&state); + let barrier = Arc::clone(&barrier); + let current_policy_json = ¤t_policy_json; + let replacement_policy_json = &replacement_policy_json; + scope.spawn(move || { + barrier.wait(); + for _ in 0..ITERATIONS { + let response = state.policy_response().expect("active policy response"); + let actual = serde_json::to_value(response.policy).unwrap(); + assert!( + actual == *current_policy_json || actual == *replacement_policy_json, + "response mixed two policy snapshots" + ); + std::thread::yield_now(); + } + }); + } + + barrier.wait(); + for index in 0..ITERATIONS { + let replacement = if index % 2 == 0 { + Arc::clone(&policy_b) + } else { + Arc::clone(&policy_a) + }; + *state.policy.write().expect("policy lock") = Some(replacement); + std::thread::yield_now(); + } + }); + } + fn request() -> PackageRequest { PackageRequest { request_kind: api::PackageRequestKind, From 4c6dc53ac0419de67cebacf07dc105dc249341ed Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Beno=C3=AEt=20CORTIER?= Date: Tue, 18 Aug 2026 13:38:09 +0900 Subject: [PATCH 2/2] fix(agent): hide policy source details Return a generic policy-unavailable message so clients cannot infer whether the active policy is file-backed, missing, or corrupt. Issue: Devolutions/now-libraries#93 Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- crates/now-package-broker/src/server/mod.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/now-package-broker/src/server/mod.rs b/crates/now-package-broker/src/server/mod.rs index bd5640097..9e40ac524 100644 --- a/crates/now-package-broker/src/server/mod.rs +++ b/crates/now-package-broker/src/server/mod.rs @@ -179,7 +179,7 @@ impl BrokerState { guard .as_ref() .map(Arc::clone) - .ok_or_else(|| error_response(ErrorCode::BrokerPaused, "policy file is unavailable or corrupted")) + .ok_or_else(|| error_response(ErrorCode::BrokerPaused, "active policy is unavailable")) } fn policy_response(&self) -> Result { @@ -699,7 +699,7 @@ mod tests { let body = response_json(response).await; let error: ErrorResponse = serde_json::from_value(body.clone()).expect("deserialize error response"); assert_eq!(error.code, ErrorCode::BrokerPaused); - assert_eq!(error.message, "policy file is unavailable or corrupted"); + assert_eq!(error.message, "active policy is unavailable"); assert!(error.details.is_empty()); assert!(body.get("Policy").is_none()); }