diff --git a/protocol/agent/v1/fixtures/diagnostic-scheduler/MANIFEST.sha256 b/protocol/agent/v1/fixtures/diagnostic-scheduler/MANIFEST.sha256 index 0c8d77a44..c74e0ed28 100644 --- a/protocol/agent/v1/fixtures/diagnostic-scheduler/MANIFEST.sha256 +++ b/protocol/agent/v1/fixtures/diagnostic-scheduler/MANIFEST.sha256 @@ -1,2 +1,3 @@ 29c93b4c8c27164896c5197e1df4ba03175425ccf5b12dc8e3f3d8d6bf19dc51 accept-vectors.json +fef527fdd89aa87fa9f81394eab6484fc2ac05a31559e1589149716193b336fb receipt-vectors.json 219749ea7bdf73461108320e13ba14b604a1967e32139c5317a017ad95dab000 reject-vectors.json diff --git a/protocol/agent/v1/fixtures/diagnostic-scheduler/receipt-vectors.json b/protocol/agent/v1/fixtures/diagnostic-scheduler/receipt-vectors.json new file mode 100644 index 000000000..6abfcb3c3 --- /dev/null +++ b/protocol/agent/v1/fixtures/diagnostic-scheduler/receipt-vectors.json @@ -0,0 +1,69 @@ +{ + "protocolVersion": "v1", + "fixtureSet": "diagnostic-scheduler", + "classification": "L1", + "description": "Outbound-only diagnostic execution receipt ingestion. Connect binds the route to the authenticated device, retains exact replays, and rejects changed or expired claims.", + "fixture": "receipt-vectors", + "accepted": [ + { + "id": "diagnostic-receipt.success", + "request": { + "protocolVersion": "v1", + "requestId": "123e4567-e89b-42d3-a456-426614174000", + "receiptId": "123e4567-e89b-42d3-a456-426614174000", + "policyRevision": 7, + "toolId": "inventory.environment", + "intervalStartedAt": "2030-01-01T00:00:00Z", + "completedAt": "2030-01-01T00:00:01Z", + "outcome": "SUCCEEDED", + "attemptCount": 1, + "reason": null, + "resultSha256": "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + "resultBytes": 256 + }, + "expected": { "httpStatus": 200, "connectVisible": true } + }, + { + "id": "diagnostic-receipt.failure-after-three-attempts", + "request": { + "protocolVersion": "v1", + "requestId": "123e4567-e89b-42d3-a456-426614174001", + "receiptId": "123e4567-e89b-42d3-a456-426614174001", + "policyRevision": 8, + "toolId": "inventory.environment", + "intervalStartedAt": "2030-01-01T00:00:00Z", + "completedAt": "2030-01-01T00:00:07Z", + "outcome": "FAILED", + "attemptCount": 3, + "reason": "connect_diagnostic_inventory_unavailable", + "resultSha256": null, + "resultBytes": null + }, + "expected": { "httpStatus": 200, "connectVisible": true } + }, + { + "id": "diagnostic-receipt.disable-or-revoke-cancelled", + "request": { + "protocolVersion": "v1", + "requestId": "123e4567-e89b-42d3-a456-426614174002", + "receiptId": "123e4567-e89b-42d3-a456-426614174002", + "policyRevision": 8, + "toolId": "inventory.environment", + "intervalStartedAt": "2030-01-01T00:00:00Z", + "completedAt": "2030-01-01T00:00:01Z", + "outcome": "CANCELLED", + "attemptCount": 1, + "reason": "connect_diagnostic_cancelled", + "resultSha256": null, + "resultBytes": null + }, + "expected": { "httpStatus": 200, "connectVisible": true } + } + ], + "rejected": [ + { "id": "diagnostic-receipt.changed-replay", "expected": { "httpStatus": 409, "reason": "REQUEST_ID_REUSED" } }, + { "id": "diagnostic-receipt.expired", "expected": { "httpStatus": 400, "reason": "DIAGNOSTIC_RECEIPT_INVALID" } }, + { "id": "diagnostic-receipt.wrong-cluster", "expected": { "httpStatus": 403, "reason": "CLUSTER_MISMATCH" } }, + { "id": "diagnostic-receipt.expired-device", "expected": { "httpStatus": 401, "reason": "CREDENTIAL_EXPIRED" } } + ] +} diff --git a/protocol/agent/v1/fixtures/fixture-sets.json b/protocol/agent/v1/fixtures/fixture-sets.json index b141dd5da..151c95840 100644 --- a/protocol/agent/v1/fixtures/fixture-sets.json +++ b/protocol/agent/v1/fixtures/fixture-sets.json @@ -56,7 +56,7 @@ { "name": "diagnostic-scheduler", "status": "populated", - "purpose": "Consent-bound diagnostic policy scheduling, bounded retry, cancellation, restart deduplication, and local execution receipts." + "purpose": "Consent-bound diagnostic policy scheduling, bounded retry, cancellation, restart deduplication, and outbound Connect-visible execution receipts." }, { "name": "offline-enrollment", diff --git a/protocol/agent/v1/fixtures/manifest.sha256 b/protocol/agent/v1/fixtures/manifest.sha256 index 1ee4032da..30041d062 100644 --- a/protocol/agent/v1/fixtures/manifest.sha256 +++ b/protocol/agent/v1/fixtures/manifest.sha256 @@ -7,7 +7,7 @@ cc6c3f3bf5e6a938f2e00b8ce42d8fa7bfb974f4eb2927101d4ede4ebdc10dd5 inventory a464d4c695d76dd985c0f2abf68c2135c080221098ad2427022f85243a7bf6ae network-performance 0dc06ec2caa8764a51d44e4176959462aa5d0ed9fe5a38453a681ed9d83b8ad6 site-replication-performance fad21eec58d4d4d547893af9050cdc5eaf9d4f84b15fb423214b16be50e7ff7a telemetry -655b83debfcbb6e2d8ad1974aeb0ce2d42f42d651f97b465dd791e3697919240 diagnostic-scheduler +484da2e7599148d36e1ac19d0ae5d305c4a96da7a00495a8ea202d3bae240a88 diagnostic-scheduler a2557ab1f7f70fb86a8affc8451596c272478fd2d5500c367b8f9e8385d72578 offline-enrollment 8b6fd07759d1264489a8bb54d2dda4836ad0f59b5ee5d52966365918b3de4c9e bundle 42cc936dc8fe87335ebff5abe01e3ebf85bdada7f37b43bfd8add46435c9d6f9 redaction diff --git a/rustfs/src/connect/diagnostics/mod.rs b/rustfs/src/connect/diagnostics/mod.rs index 3334d248a..f98af6564 100644 --- a/rustfs/src/connect/diagnostics/mod.rs +++ b/rustfs/src/connect/diagnostics/mod.rs @@ -21,6 +21,7 @@ mod perf_object; mod profile_cpu; mod profile_memory; mod profile_threads; +mod receipt_delivery; mod schedule; mod top_api; mod top_disk; @@ -83,6 +84,7 @@ pub use profile_cpu::{ }; pub use profile_memory::export_memory_profile; pub use profile_threads::{capture_thread_profile, export_thread_profile}; +pub(crate) use receipt_delivery::{DiagnosticReceiptDelivery, DiagnosticReceiptSender}; pub use schedule::{ DiagnosticCollectionPolicy, DiagnosticReceipt, DiagnosticScheduleError, DiagnosticScheduleRuntime, DiagnosticScheduleStatus, ReceiptOutcome, run_local_environment_once, spawn_environment_schedule, diff --git a/rustfs/src/connect/diagnostics/receipt_delivery.rs b/rustfs/src/connect/diagnostics/receipt_delivery.rs new file mode 100644 index 000000000..f85cde976 --- /dev/null +++ b/rustfs/src/connect/diagnostics/receipt_delivery.rs @@ -0,0 +1,112 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::time::Duration; + +use serde::{Deserialize, Serialize}; + +use super::DiagnosticReceipt; +use crate::connect::config::HeartbeatConfig; +use crate::connect::heartbeat::HeartbeatError; +use crate::connect::telemetry::{TelemetryDelivery, TelemetryTransport}; + +const PROTOCOL_VERSION: &str = "v1"; + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct DiagnosticReceiptRequest<'a> { + protocol_version: &'static str, + request_id: &'a str, + #[serde(flatten)] + receipt: &'a DiagnosticReceipt, +} + +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +struct DiagnosticReceiptResponse { + accepted_version: String, +} + +pub(crate) enum DiagnosticReceiptDelivery { + Accepted, + Retry { retry_after: Option }, + AuthenticationStopped, + Rejected, +} + +pub(crate) struct DiagnosticReceiptSender { + transport: TelemetryTransport, +} + +impl DiagnosticReceiptSender { + pub(crate) fn new(config: HeartbeatConfig) -> Result { + Ok(Self { + transport: TelemetryTransport::new(config)?, + }) + } + + pub(crate) async fn send(&self, receipt: &DiagnosticReceipt) -> Result { + let request = DiagnosticReceiptRequest { + protocol_version: PROTOCOL_VERSION, + request_id: &receipt.receipt_id, + receipt, + }; + Ok(match self.transport.post("diagnosticExecutionReceipts", &request).await? { + TelemetryDelivery::Accepted { body, .. } => { + let response: DiagnosticReceiptResponse = serde_json::from_slice(&body).map_err(|_| HeartbeatError::Response)?; + if response.accepted_version != PROTOCOL_VERSION { + return Err(HeartbeatError::Response); + } + DiagnosticReceiptDelivery::Accepted + } + TelemetryDelivery::Retry { retry_after } => DiagnosticReceiptDelivery::Retry { retry_after }, + TelemetryDelivery::AuthenticationStopped { .. } => DiagnosticReceiptDelivery::AuthenticationStopped, + TelemetryDelivery::Rejected { .. } => DiagnosticReceiptDelivery::Rejected, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::connect::diagnostics::ReceiptOutcome; + + #[test] + fn request_is_the_versioned_frozen_receipt_with_receipt_id_as_idempotency_key() { + let receipt = DiagnosticReceipt { + receipt_id: "123e4567-e89b-42d3-a456-426614174001".to_owned(), + policy_revision: 8, + tool_id: "inventory.environment".to_owned(), + interval_started_at: "2030-01-01T00:00:00Z".to_owned(), + completed_at: "2030-01-01T00:00:07Z".to_owned(), + outcome: ReceiptOutcome::Failed, + attempt_count: 3, + reason: Some("connect_diagnostic_inventory_unavailable".to_owned()), + result_sha256: None, + result_bytes: None, + }; + let value = serde_json::to_value(DiagnosticReceiptRequest { + protocol_version: PROTOCOL_VERSION, + request_id: &receipt.receipt_id, + receipt: &receipt, + }) + .expect("receipt request"); + + assert_eq!(value["protocolVersion"], "v1"); + assert_eq!(value["requestId"], value["receiptId"]); + assert_eq!(value["attemptCount"], 3); + assert_eq!(value["outcome"], "FAILED"); + assert!(value.get("command").is_none()); + } +} diff --git a/rustfs/src/connect/diagnostics/schedule.rs b/rustfs/src/connect/diagnostics/schedule.rs index 1e4182353..59c4d686c 100644 --- a/rustfs/src/connect/diagnostics/schedule.rs +++ b/rustfs/src/connect/diagnostics/schedule.rs @@ -267,6 +267,7 @@ type Runner = Arc RunFuture + Send + Sync>; pub struct DiagnosticScheduleRuntime { status: watch::Receiver, + receipts: watch::Receiver>, task: JoinHandle<()>, } @@ -275,6 +276,10 @@ impl DiagnosticScheduleRuntime { self.status.clone() } + pub(crate) fn receipts(&self) -> watch::Receiver> { + self.receipts.clone() + } + pub async fn shutdown(self) { let _ = self.task.await; } @@ -331,15 +336,20 @@ fn spawn_schedule( runner: Runner, ) -> DiagnosticScheduleRuntime { let (status_tx, status_rx) = watch::channel(DiagnosticScheduleStatus::Waiting); + let (receipt_tx, receipt_rx) = watch::channel(None); let task = tokio::spawn(async move { - if let Err(error) = run_collection_schedule(&store, &mut policies, &shutdown, &runner, &status_tx).await { + if let Err(error) = run_collection_schedule(&store, &mut policies, &shutdown, &runner, &status_tx, &receipt_tx).await { let _ = status_tx.send(DiagnosticScheduleStatus::Failed { reason: error.to_string(), }); } let _ = status_tx.send(DiagnosticScheduleStatus::Stopped); }); - DiagnosticScheduleRuntime { status: status_rx, task } + DiagnosticScheduleRuntime { + status: status_rx, + receipts: receipt_rx, + task, + } } async fn run_collection_schedule( @@ -348,8 +358,12 @@ async fn run_collection_schedule( shutdown: &CancellationToken, runner: &Runner, status: &watch::Sender, + receipts: &watch::Sender>, ) -> Result<(), DiagnosticScheduleError> { let mut state = store.read().await?; + if state.last_receipt.is_some() { + let _ = receipts.send(state.last_receipt.clone()); + } if let Some(started_at) = state.active_interval_started_at.take() { let receipt = receipt( state.policy_revision.unwrap_or_default(), @@ -363,6 +377,7 @@ async fn run_collection_schedule( state.last_receipt = Some(receipt.clone()); store.write(state.clone()).await?; let _ = status.send(DiagnosticScheduleStatus::Receipt(receipt)); + let _ = receipts.send(state.last_receipt.clone()); } loop { @@ -550,6 +565,7 @@ async fn run_collection_schedule( state.last_receipt = Some(receipt.clone()); store.write(state.clone()).await?; let _ = status.send(DiagnosticScheduleStatus::Receipt(receipt)); + let _ = receipts.send(state.last_receipt.clone()); } Ok(()) } diff --git a/rustfs/src/connect/runtime.rs b/rustfs/src/connect/runtime.rs index d6cbcd00f..8a6838b93 100644 --- a/rustfs/src/connect/runtime.rs +++ b/rustfs/src/connect/runtime.rs @@ -26,7 +26,8 @@ use tokio_util::sync::CancellationToken; use super::client::{ClientError, ConnectClient, ConnectConfig, RotationAttempt}; use super::config::HeartbeatConfig; use super::diagnostics::{ - DiagnosticCollectionPolicy, DiagnosticScheduleRuntime, DiagnosticScheduleStatus, spawn_environment_schedule, + DiagnosticCollectionPolicy, DiagnosticReceipt, DiagnosticReceiptDelivery, DiagnosticReceiptSender, DiagnosticScheduleRuntime, + DiagnosticScheduleStatus, spawn_environment_schedule, }; use super::heartbeat::{CoarseNodeSummary, Delivery, HeartbeatError, HeartbeatSender, HeartbeatStateStore, HeartbeatStatus}; use super::inventory::{ @@ -40,6 +41,7 @@ pub struct HeartbeatRuntime { task: Option>, diagnostic_status: watch::Receiver, diagnostic_task: Option, + diagnostic_receipt_task: Option>, } impl HeartbeatRuntime { @@ -59,6 +61,9 @@ impl HeartbeatRuntime { if let Some(task) = self.diagnostic_task.take() { task.shutdown().await; } + if let Some(task) = self.diagnostic_receipt_task.take() { + let _ = task.await; + } } } @@ -123,6 +128,7 @@ where return Ok(None); } let sender = HeartbeatSender::new(config.clone())?; + let receipt_sender = DiagnosticReceiptSender::new(config.clone())?; let rotation = ConnectClient::new(ConnectConfig { endpoint: &config.endpoint, root_ca_pem: &config.root_ca_pem, @@ -143,6 +149,19 @@ where let diagnostic_task = spawn_environment_schedule(&state_root, policy_rx, shutdown.clone()).map_err(|_| HeartbeatError::StateConflict)?; let diagnostic_status = diagnostic_task.status(); + let receipt_status = diagnostic_task.receipts(); + let receipt_shutdown = shutdown.clone(); + let retry_schedule = schedule; + let diagnostic_receipt_task = tokio::spawn(async move { + run_diagnostic_receipt_delivery( + receipt_sender, + receipt_status, + retry_schedule.initial_backoff, + retry_schedule.max_backoff, + receipt_shutdown, + ) + .await; + }); let task = tokio::spawn(async move { let _lock = lock; let mut backoff = schedule.initial_backoff; @@ -240,9 +259,58 @@ where task: Some(task), diagnostic_status, diagnostic_task: Some(diagnostic_task), + diagnostic_receipt_task: Some(diagnostic_receipt_task), })) } +async fn run_diagnostic_receipt_delivery( + sender: DiagnosticReceiptSender, + mut receipts: watch::Receiver>, + initial_backoff: Duration, + max_backoff: Duration, + shutdown: CancellationToken, +) { + let mut last_accepted = None; + let mut backoff = initial_backoff; + loop { + let Some(receipt) = receipts.borrow_and_update().clone() else { + tokio::select! { + biased; + () = shutdown.cancelled() => break, + changed = receipts.changed() => if changed.is_err() { break; }, + } + continue; + }; + if last_accepted.as_deref() == Some(receipt.receipt_id.as_str()) { + tokio::select! { + biased; + () = shutdown.cancelled() => break, + changed = receipts.changed() => if changed.is_err() { break; }, + } + continue; + } + let delivery = tokio::select! { + biased; + () = shutdown.cancelled() => break, + delivery = sender.send(&receipt) => delivery, + }; + match delivery { + Ok(DiagnosticReceiptDelivery::Accepted) => { + last_accepted = Some(receipt.receipt_id.clone()); + backoff = initial_backoff; + } + Ok(DiagnosticReceiptDelivery::Retry { retry_after }) => { + let delay = retry_after.unwrap_or(backoff).clamp(initial_backoff, max_backoff); + backoff = backoff.saturating_mul(2).min(max_backoff); + if sleep_or_cancel(&shutdown, delay).await { + break; + } + } + Ok(DiagnosticReceiptDelivery::AuthenticationStopped | DiagnosticReceiptDelivery::Rejected) | Err(_) => break, + } + } +} + pub fn spawn_inventory_runtime( config: Option, schedule: InventorySchedule, @@ -595,6 +663,7 @@ mod tests { task: Some(heartbeat_task), diagnostic_status, diagnostic_task: None, + diagnostic_receipt_task: None, }; let inventory = InventoryRuntime { shutdown: inventory_shutdown, diff --git a/rustfs/tests/connect_heartbeat.rs b/rustfs/tests/connect_heartbeat.rs index 7280f8779..06d0f4909 100644 --- a/rustfs/tests/connect_heartbeat.rs +++ b/rustfs/tests/connect_heartbeat.rs @@ -218,16 +218,29 @@ async fn server(pki: &TestPki, replies: Vec) -> TestServer { let replies = replies.clone(); let seen = seen.clone(); async move { - assert_eq!(request.uri().path(), format!("/agent/clusters/{CLUSTER_UID}/heartbeats")); + let path = request.uri().path().to_owned(); + assert!( + path == format!("/agent/clusters/{CLUSTER_UID}/heartbeats") + || path == format!("/agent/clusters/{CLUSTER_UID}/diagnosticExecutionReceipts") + ); let body = request.into_body().collect().await.expect("request body").to_bytes(); seen.lock() .expect("seen lock") .push(serde_json::from_slice(&body).expect("request JSON")); - let reply = replies - .lock() - .expect("reply lock") - .pop_front() - .unwrap_or_else(|| Reply::error(StatusCode::SERVICE_UNAVAILABLE)); + let reply = if path.ends_with("/diagnosticExecutionReceipts") { + Reply { + status: StatusCode::OK, + body: json!({"acceptedVersion": "v1"}), + retry_after: None, + delay: Duration::ZERO, + } + } else { + replies + .lock() + .expect("reply lock") + .pop_front() + .unwrap_or_else(|| Reply::error(StatusCode::SERVICE_UNAVAILABLE)) + }; if !reply.delay.is_zero() { tokio::time::sleep(reply.delay).await; } @@ -602,6 +615,71 @@ async fn restart_replays_pending_request_then_advances_sequence() { assert_eq!(seen[1]["sequence"].as_u64(), seen[0]["sequence"].as_u64().map(|value| value + 1)); } +#[tokio::test] +async fn restart_sends_the_last_durable_diagnostic_receipt_outbound() { + let pki = TestPki::new(); + let server = server(&pki, vec![Reply::ok("2026-08-22T01:02:03Z")]).await; + let temp = tempfile::tempdir().expect("tempdir"); + let diagnostics = temp.path().join("diagnostics"); + fs::create_dir_all(&diagnostics).expect("diagnostics directory"); + fs::write( + diagnostics.join("schedule.json"), + serde_json::to_vec(&json!({ + "policyRevision": null, + "policyFingerprint": null, + "nextDueAt": null, + "activeIntervalStartedAt": null, + "lastReceipt": { + "receiptId": "123e4567-e89b-42d3-a456-426614174001", + "policyRevision": 8, + "toolId": "inventory.environment", + "intervalStartedAt": "2030-01-01T00:00:00Z", + "completedAt": "2030-01-01T00:00:07Z", + "outcome": "FAILED", + "attemptCount": 3, + "reason": "connect_diagnostic_inventory_unavailable", + "resultSha256": null, + "resultBytes": null + } + })) + .expect("schedule state"), + ) + .expect("write schedule state"); + private_mode(&diagnostics.join("schedule.json")); + let shutdown = CancellationToken::new(); + let runtime = spawn_heartbeat_runtime(Some(config(&temp, &pki, &server)), &shutdown, summary) + .expect("start runtime") + .expect("configured runtime"); + + tokio::time::timeout(Duration::from_secs(3), async { + loop { + if server + .seen + .lock() + .expect("seen lock") + .iter() + .any(|request| request["receiptId"] == "123e4567-e89b-42d3-a456-426614174001") + { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("receipt delivery"); + runtime.shutdown().await; + + let seen = server.seen.lock().expect("seen lock"); + let receipt = seen + .iter() + .find(|request| request.get("receiptId").is_some()) + .expect("receipt request"); + assert_eq!(receipt["protocolVersion"], "v1"); + assert_eq!(receipt["requestId"], receipt["receiptId"]); + assert_eq!(receipt["attemptCount"], 3); + assert!(receipt.get("command").is_none()); +} + #[tokio::test] async fn retry_after_is_respected_with_the_local_upper_bound() { let pki = TestPki::new();