Deliver diagnostic scheduler receipts to Connect (#7767)

feat(connect): deliver diagnostic scheduler receipts
This commit is contained in:
Chris
2026-09-14 01:34:26 +08:00
committed by GitHub
parent 1e53090cbb
commit 7111f8b44d
9 changed files with 358 additions and 11 deletions
@@ -1,2 +1,3 @@
29c93b4c8c27164896c5197e1df4ba03175425ccf5b12dc8e3f3d8d6bf19dc51 accept-vectors.json
fef527fdd89aa87fa9f81394eab6484fc2ac05a31559e1589149716193b336fb receipt-vectors.json
219749ea7bdf73461108320e13ba14b604a1967e32139c5317a017ad95dab000 reject-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" } }
]
}
+1 -1
View File
@@ -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",
+1 -1
View File
@@ -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
+2
View File
@@ -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,
@@ -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<Duration> },
AuthenticationStopped,
Rejected,
}
pub(crate) struct DiagnosticReceiptSender {
transport: TelemetryTransport,
}
impl DiagnosticReceiptSender {
pub(crate) fn new(config: HeartbeatConfig) -> Result<Self, HeartbeatError> {
Ok(Self {
transport: TelemetryTransport::new(config)?,
})
}
pub(crate) async fn send(&self, receipt: &DiagnosticReceipt) -> Result<DiagnosticReceiptDelivery, HeartbeatError> {
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());
}
}
+18 -2
View File
@@ -267,6 +267,7 @@ type Runner = Arc<dyn Fn(CancellationToken) -> RunFuture + Send + Sync>;
pub struct DiagnosticScheduleRuntime {
status: watch::Receiver<DiagnosticScheduleStatus>,
receipts: watch::Receiver<Option<DiagnosticReceipt>>,
task: JoinHandle<()>,
}
@@ -275,6 +276,10 @@ impl DiagnosticScheduleRuntime {
self.status.clone()
}
pub(crate) fn receipts(&self) -> watch::Receiver<Option<DiagnosticReceipt>> {
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<DiagnosticScheduleStatus>,
receipts: &watch::Sender<Option<DiagnosticReceipt>>,
) -> 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(())
}
+70 -1
View File
@@ -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<JoinHandle<()>>,
diagnostic_status: watch::Receiver<DiagnosticScheduleStatus>,
diagnostic_task: Option<DiagnosticScheduleRuntime>,
diagnostic_receipt_task: Option<JoinHandle<()>>,
}
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<Option<DiagnosticReceipt>>,
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<F, Fut>(
config: Option<HeartbeatConfig>,
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,
+84 -6
View File
@@ -218,16 +218,29 @@ async fn server(pki: &TestPki, replies: Vec<Reply>) -> 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();