feat(connect): sample memory within the running service (#8091)

* feat(connect): sample memory within the running service

* test(connect): add official memory service acceptance

* test(connect): consume the final memory job without cloning
This commit is contained in:
Chris
2026-09-24 03:25:28 +08:00
committed by GitHub
parent 9ab034df62
commit be01e513be
11 changed files with 637 additions and 8 deletions
@@ -0,0 +1,201 @@
# Copyright 2026 RustFS Team
# SPDX-License-Identifier: Apache-2.0
name: Connect service memory profile acceptance
on:
workflow_dispatch:
inputs:
build_run_id:
description: Successful main-branch Build and Release run ID
required: true
type: string
artifact_id:
description: Linux x86_64 GNU artifact from that run
required: true
type: string
source_sha:
description: Exact RustFS source commit
required: true
type: string
artifact_digest:
description: GitHub artifact digest including sha256 prefix
required: true
type: string
binary_sha256:
description: Independently verified rustfs binary SHA-256
required: true
type: string
connect_sha:
description: Exact Connect memory acceptance harness commit
required: true
type: string
permissions:
actions: read
contents: read
jobs:
profile-memory:
name: Verify signed service memory profile over mTLS
runs-on: ubuntu-latest
timeout-minutes: 45
steps:
- name: Checkout exact RustFS source
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
path: rustfs-source
persist-credentials: false
ref: ${{ inputs.source_sha }}
- name: Checkout exact Connect harness
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
repository: rustfs/connect
path: connect-harness
persist-credentials: false
ref: ${{ inputs.connect_sha }}
token: ${{ secrets.PF_TESTING_GH_TOKEN }}
- name: Set up Node.js
uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0
with:
node-version: 22.22.2
cache: npm
cache-dependency-path: connect-harness/web/package-lock.json
- name: Verify official source and artifact identity
shell: bash
env:
GH_TOKEN: ${{ github.token }}
BUILD_RUN_ID: ${{ inputs.build_run_id }}
ARTIFACT_ID: ${{ inputs.artifact_id }}
SOURCE_SHA: ${{ inputs.source_sha }}
ARTIFACT_DIGEST: ${{ inputs.artifact_digest }}
BINARY_SHA256: ${{ inputs.binary_sha256 }}
CONNECT_SHA: ${{ inputs.connect_sha }}
run: |
set -euo pipefail
[[ "$GITHUB_REPOSITORY" == rustfs/rustfs ]]
[[ "$BUILD_RUN_ID" =~ ^[1-9][0-9]*$ && "$ARTIFACT_ID" =~ ^[1-9][0-9]*$ ]]
[[ "$SOURCE_SHA" =~ ^[0-9a-f]{40}$ && "$CONNECT_SHA" =~ ^[0-9a-f]{40}$ ]]
[[ "$ARTIFACT_DIGEST" =~ ^sha256:[0-9a-f]{64}$ ]]
[[ "$BINARY_SHA256" =~ ^[0-9a-f]{64}$ ]]
[[ $(git -C rustfs-source rev-parse HEAD) == "$SOURCE_SHA" ]]
[[ $(git -C connect-harness rev-parse HEAD) == "$CONNECT_SHA" ]]
[[ $(git -C rustfs-source remote get-url origin) == https://github.com/rustfs/rustfs ]]
[[ $(git -C connect-harness remote get-url origin) == https://github.com/rustfs/connect ]]
run=$(gh api "repos/rustfs/rustfs/actions/runs/${BUILD_RUN_ID}")
jq -e --arg source "$SOURCE_SHA" '
.head_sha == $source and .head_branch == "main"
and .head_repository.full_name == "rustfs/rustfs"
and .name == "Build and Release" and .path == ".github/workflows/build.yml"
and .status == "completed" and .conclusion == "success"
' <<<"$run" >/dev/null
jobs=$(gh api --paginate --slurp "repos/rustfs/rustfs/actions/runs/${BUILD_RUN_ID}/jobs?per_page=100")
jq -e '
[.[].jobs[] | select(.name | test("^Build RustFS \\(linux-x86_64-gnu, [a-z0-9-]+, x86_64-unknown-linux-gnu, false, linux, pyroscope\\)$"))] as $matches
| ($matches | length) == 1 and $matches[0].conclusion == "success"
' <<<"$jobs" >/dev/null
artifact=$(gh api "repos/rustfs/rustfs/actions/artifacts/${ARTIFACT_ID}")
jq -e --argjson run "$BUILD_RUN_ID" --arg source "$SOURCE_SHA" --arg digest "$ARTIFACT_DIGEST" '
.workflow_run.id == $run and .workflow_run.head_sha == $source
and .name == ("rustfs-linux-x86_64-gnu-dev-" + $source[0:7])
and .digest == $digest and .expired == false
and (.expires_at | fromdateiso8601) > now
and .size_in_bytes > 0 and .size_in_bytes <= 2147483648
' <<<"$artifact" >/dev/null
jq -r '.expires_at' <<<"$artifact" > artifact-expires-at
gh api "repos/rustfs/rustfs/actions/artifacts/${ARTIFACT_ID}/zip" > official-artifact.zip
[[ $(stat --format=%s official-artifact.zip) == $(jq -r '.size_in_bytes' <<<"$artifact") ]]
printf '%s official-artifact.zip\n' "${ARTIFACT_DIGEST#sha256:}" | sha256sum --check --strict
- name: Extract only the verified official service binary
shell: bash
env:
SOURCE_SHA: ${{ inputs.source_sha }}
BINARY_SHA256: ${{ inputs.binary_sha256 }}
run: |
set -euo pipefail
python3 - <<'PY'
import io
import os
import pathlib
import shutil
import stat
import zipfile
def regular(entry):
mode = entry.external_attr >> 16
return not entry.is_dir() and stat.S_IFMT(mode) in (0, stat.S_IFREG)
package = f"rustfs-linux-x86_64-gnu-dev-{os.environ['SOURCE_SHA'][:7]}.zip"
with zipfile.ZipFile('official-artifact.zip') as outer:
entries = outer.infolist()
assert 1 <= len(entries) <= 16
assert len({entry.filename for entry in entries}) == len(entries)
assert all(regular(entry) and pathlib.PurePosixPath(entry.filename).name == entry.filename for entry in entries)
assert sum(entry.file_size for entry in entries) <= 2147483648
assert outer.testzip() is None
with zipfile.ZipFile(io.BytesIO(outer.read(package))) as inner:
entries = inner.infolist()
assert {entry.filename for entry in entries} == {'rustfs', 'rustfs-cli'} and len(entries) == 2
assert all(regular(entry) and 0 < entry.file_size <= 1073741824 for entry in entries)
assert inner.testzip() is None
pathlib.Path('binary').mkdir()
with inner.open('rustfs') as source, open('binary/rustfs', 'xb') as destination:
shutil.copyfileobj(source, destination)
PY
chmod 0755 binary/rustfs
printf '%s binary/rustfs\n' "$BINARY_SHA256" | sha256sum --check --strict
file binary/rustfs | grep -Eq 'ELF 64-bit LSB.*x86-64'
- name: Run the real service job under S3 workload
shell: bash
env:
SOURCE_SHA: ${{ inputs.source_sha }}
run: |
set -euo pipefail
git -C connect-harness diff --exit-code
printf '%s\n' "$SOURCE_SHA" >connect-harness/tests/e2e/connected/rustfs-ref
changed=$(git -C connect-harness diff --name-only)
[[ -z "$changed" || "$changed" == tests/e2e/connected/rustfs-ref ]]
npm --prefix connect-harness/web ci
connect-harness/web/node_modules/.bin/playwright install --with-deps chromium
make -C connect-harness e2e-connected-dispatch-check
(cd connect-harness && ./tests/connectivity/profile-memory-evidence.test.sh)
RUSTFS_BINARY="$GITHUB_WORKSPACE/binary/rustfs" \
RUSTFS_WORKTREE="$GITHUB_WORKSPACE/rustfs-source" \
CONNECT_E2E_PROFILE_MEMORY_EVIDENCE="$GITHUB_WORKSPACE/profile-memory-evidence.json" \
make -C connect-harness e2e-connected E2E_SCENARIO=profile-memory
- name: Bind and validate sanitized evidence
shell: bash
env:
BUILD_RUN_ID: ${{ inputs.build_run_id }}
ARTIFACT_ID: ${{ inputs.artifact_id }}
ARTIFACT_DIGEST: ${{ inputs.artifact_digest }}
BINARY_SHA256: ${{ inputs.binary_sha256 }}
CONNECT_SHA: ${{ inputs.connect_sha }}
SOURCE_SHA: ${{ inputs.source_sha }}
run: |
set -euo pipefail
jq -e --arg connect "$CONNECT_SHA" --arg source "$SOURCE_SHA" --arg binary "$BINARY_SHA256" '
.provenance.connectSha == $connect and .provenance.sourceSha == $source
and .provenance.binarySha256 == $binary
' profile-memory-evidence.json >/dev/null
jq --arg run "$GITHUB_RUN_ID" --arg build "$BUILD_RUN_ID" \
--arg artifact "$ARTIFACT_ID" --arg digest "$ARTIFACT_DIGEST" \
--arg expiry "$(cat artifact-expires-at)" '
.provenance += {workflowRunId:$run,buildRunId:$build,artifactId:$artifact,artifactDigest:$digest,artifactExpiresAt:$expiry}
' profile-memory-evidence.json > profile-memory-evidence.bound.json
mv profile-memory-evidence.bound.json profile-memory-evidence.json
node connect-harness/tests/connectivity/profile-memory-evidence.mjs profile-memory-evidence.json --bound
- name: Upload only verified sanitized evidence
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
with:
name: connect-profile-memory-evidence-${{ github.run_id }}
path: profile-memory-evidence.json
retention-days: 14
if-no-files-found: error
@@ -11,6 +11,8 @@
## Open Items
- `connect-894` pending Connect heartbeats: replay the exact preceding producer-capability list, with or without jobs, when upgrading to the memory-service capability. Preserve request ID, sequence, and persisted body rather than inserting the new capability into a retry. Remove after upgrades from the pre-service-memory capability set are unsupported.
- `backlog-2539` administrator erasure-set scope: decode historical pending intents whose empty bucket list omitted the all-buckets marker. Normalize that representation to the explicit scope before replay or checkpoint comparison. Remove after all pre-marker administrator ErasureSet intents are retired; deployment verification must confirm that no such pending records remain on coordinator disks.
- `backlog-2519` retained admin heal reports: keep the schema-1 terminal as the commit and replay fence, with a bounded versioned report in a separate namespace on the same disk. Older rollback readers can still query the terminal and suppress replay, but cannot expose its outcome; newer readers mark missing outcomes as unavailable. Remove the legacy marker and missing-report adapter only after all supported direct-upgrade and rollback readers understand the report format and retained schema-1-only receipts have expired.
@@ -1,6 +1,6 @@
975c1ca53eefeef6766a6fc0b3d3281f7408255342b0686e5e2aee5ad055414c duplicate.json
963529a38a02849c6c2acc6d72668dca9f63218b49c89fae41a451b584850411 overflow.json
6afec3806c98fb55d477455b43935798d8fdb7f9ae39bd9f6637465539e5c02b producer-capabilities.json
60300d053dd3746a3aceb327315e8c5abe292fb6505dd02a64b8d6279fb56c09 producer-capabilities.json
e3adeee1c8a19aa17e70894896fb79c072e3785bea3611b93c11e79f039ed5af stale.json
22cc7337f545271ebf11ca6b76e7e479d62c51278ac001a8d870b07d0fa9884e unknown.json
b06a72ff1d964efe65e996797705ec2ef72574a813aec5d8e9aff26c899f1352 valid.json
@@ -17,6 +17,7 @@
"logs.capture@1",
"profile.cpu@1",
"profile.memory@1",
"profile.memory.service@1",
"profile.threads@1",
"telemetry.record@1",
"telemetry.otlp@1",
@@ -40,6 +41,7 @@
"logs.capture@1",
"profile.cpu@1",
"profile.memory@1",
"profile.memory.service@1",
"profile.threads@1",
"telemetry.record@1",
"telemetry.otlp@1",
+1 -1
View File
@@ -1,7 +1,7 @@
a499742c06046e9cf2f781fc17d8b4cf95e018ab45fa41ee5944ba15b0e2902c auth
9b54e39584a442270ae486e8378df1d6f168c45d0c59baba9459b46bb212d2df version
159d116e965b51771f1687fc6ec9bdfbf2a4a0e6d6aa26d3f1046b3dac2697dd registration
fe1287eac16832bed88320254245ba07a985e6ae7e66e1d1570e324ffa972873 heartbeat
9b3862f5160968a2ff281afb3ea519a3d49ae1bc3d767153d70f309d2a7457a7 heartbeat
cc6c3f3bf5e6a938f2e00b8ce42d8fa7bfb974f4eb2927101d4ede4ebdc10dd5 inventory
0207458ec9b97b368e2f347da112511a3aabd27d752890cacda4fa2408e7f19e object-performance
a464d4c695d76dd985c0f2abf68c2135c080221098ad2427022f85243a7bf6ae network-performance
+71 -4
View File
@@ -29,6 +29,7 @@ use thiserror::Error;
use tokio_util::sync::CancellationToken;
use uuid::{Uuid, Variant, Version};
use super::profile_memory::MEMORY_PROFILE_SERVICE_CAPABILITY;
use super::{
CPU_PROFILE_CAPABILITY, DRIVE_CAPABILITY, DRIVE_SCHEMA_VERSION, DriveOutcome, DrivePerformanceError, DrivePerformanceRequest,
DriveProvenance, LocalDriveConsent, LocalNetworkConsent, LocalProfileConsent, LocalTopConsent, MAX_NETWORK_TRAFFIC_BYTES,
@@ -45,6 +46,7 @@ use crate::connect::DeviceIdentity;
const PROTOCOL_VERSION: &str = "v1";
const PROFILE_CPU_JOB_TYPE: &str = "profile.cpu";
const PROFILE_THREADS_JOB_TYPE: &str = "profile.threads";
const PROFILE_MEMORY_JOB_TYPE: &str = "profile.memory";
const PERFORMANCE_DRIVE_JOB_TYPE: &str = "performance.drive";
const PERFORMANCE_NETWORK_JOB_TYPE: &str = "performance.network";
const TOP_API_JOB_TYPE: &str = "top.api";
@@ -56,6 +58,7 @@ const MAX_FUTURE_SKEW_SECONDS: i64 = 300;
const MAX_OUTPUT_BYTES: u64 = 524_288;
const MAX_MEMORY_BYTES: u64 = 64 * 1024 * 1024;
const MAX_CPU_MILLIS: u64 = 30_000;
const MAX_MEMORY_PROFILE_CPU_MILLIS: u64 = 5_000;
const MAX_NETWORK_CPU_MILLIS: u64 = 5_000;
const MIN_NETWORK_MEMORY_BYTES: u64 = 1_048_576;
const DRIVE_TARGET_ALIAS: &str = "drive-1";
@@ -75,6 +78,7 @@ const MIN_TOP_RPC_MEMORY_BYTES: u64 = 1_048_576;
enum DiagnosticJobKind {
ProfileCpu,
ProfileThreads,
ProfileMemory,
PerformanceDrive,
PerformanceNetwork,
TopApi,
@@ -424,6 +428,11 @@ impl DiagnosticJobEnvelope {
{
return Err(DiagnosticJobError::LimitExceeded);
}
if kind == DiagnosticJobKind::ProfileMemory
&& (self.limits.max_cpu_millis > MAX_MEMORY_PROFILE_CPU_MILLIS || self.limits.max_memory_bytes != MAX_MEMORY_BYTES)
{
return Err(DiagnosticJobError::LimitExceeded);
}
match (kind, self.parameters.traffic_bytes) {
(DiagnosticJobKind::PerformanceNetwork, Some(1..=MAX_NETWORK_TRAFFIC_BYTES)) => {
if self.limits.max_cpu_millis > MAX_NETWORK_CPU_MILLIS || self.limits.max_memory_bytes < MIN_NETWORK_MEMORY_BYTES
@@ -478,6 +487,11 @@ impl DiagnosticJobEnvelope {
(PROFILE_THREADS_JOB_TYPE, [capability], PROFILE_SCHEMA_VERSION) if capability == THREAD_PROFILE_CAPABILITY => {
Ok(DiagnosticJobKind::ProfileThreads)
}
(PROFILE_MEMORY_JOB_TYPE, [capability], PROFILE_SCHEMA_VERSION)
if capability == MEMORY_PROFILE_SERVICE_CAPABILITY =>
{
Ok(DiagnosticJobKind::ProfileMemory)
}
(PERFORMANCE_DRIVE_JOB_TYPE, [capability], DRIVE_SCHEMA_VERSION) if capability == DRIVE_CAPABILITY => {
Ok(DiagnosticJobKind::PerformanceDrive)
}
@@ -518,6 +532,7 @@ pub async fn execute_diagnostic_job(
match envelope.kind()? {
DiagnosticJobKind::ProfileCpu => execute_profile_cpu_job(envelope, nonce, identity, provenance, cancel).await,
DiagnosticJobKind::ProfileThreads => execute_profile_threads_job(envelope, nonce, identity, provenance, cancel).await,
DiagnosticJobKind::ProfileMemory => execute_profile_memory_job(envelope, nonce, identity, provenance, cancel).await,
DiagnosticJobKind::PerformanceDrive => {
let Some(scratch_root) = runtime_drive_scratch_root() else {
return Ok(failed_drive_execution(&envelope.job_id, "SOURCE_UNAVAILABLE"));
@@ -742,6 +757,54 @@ async fn execute_performance_network_job(
})
}
async fn execute_profile_memory_job(
envelope: DiagnosticJobEnvelope,
nonce: [u8; 32],
identity: &DeviceIdentity,
provenance: ProfileProvenance,
cancel: &CancellationToken,
) -> Result<DiagnosticJobExecution, DiagnosticJobError> {
let expire = parse_time(&envelope.expire_time)?;
let consent_expire = parse_time(&envelope.parameters.consent_expires_at)?;
let request = ProfileCaptureRequest {
organization_name: envelope.organization_name,
cluster_name: envelope.cluster_name,
device_name: envelope.device_name,
run_uid: envelope.job_id.clone(),
artifact_uid: envelope.parameters.artifact_uid,
schema_version: envelope.schema_version,
// The service capability gates dispatch; schema-v1 exports retain the
// original allocation-aggregate capability and signed artifact format.
capability: super::MEMORY_PROFILE_CAPABILITY.to_owned(),
consent: LocalProfileConsent {
consent_uid: envelope.parameters.consent_uid,
policy_revision: envelope.parameters.consent_policy_revision,
expires_at_unix: consent_expire.timestamp(),
confirmed: true,
},
produced_at_unix: Utc::now().timestamp(),
expires_at_unix: expire.timestamp(),
nonce,
duration: Duration::from_millis(envelope.parameters.duration_millis),
sample_period: Duration::from_micros(envelope.parameters.sample_period_micros),
provenance,
};
let export = super::export_memory_profile(&request, identity, cancel)
.await
.map_err(capture_failure)?;
if export.archive_bytes.len() > usize::try_from(envelope.limits.max_output_bytes).unwrap_or(usize::MAX) {
return Err(DiagnosticJobError::LimitExceeded);
}
Ok(DiagnosticJobExecution {
job_id: envelope.job_id,
outcome: "SUCCEEDED".to_owned(),
reason: "COMPLETE".to_owned(),
artifact_uid: Some(export.artifact_uid),
artifact_sha256: Some(export.archive_sha256),
artifact_bytes: Some(export.archive_bytes),
})
}
async fn execute_profile_threads_job(
envelope: DiagnosticJobEnvelope,
nonce: [u8; 32],
@@ -1220,6 +1283,7 @@ mod tests {
("logs.capture@1", &["logs"][..]),
("profile.cpu@1", &["profile"][..]),
("profile.memory@1", &["profile"][..]),
("profile.memory.service@1", &[][..]),
("profile.threads@1", &["profile"][..]),
("telemetry.record@1", &["telemetry", "record"][..]),
("telemetry.otlp@1", &["telemetry", "otlp"][..]),
@@ -1240,13 +1304,16 @@ mod tests {
let connect = command.find_subcommand("connect").expect("connect command");
for (capability, path) in expected {
if path.is_empty() {
// Network probes use the authenticated service dispatcher and
// locally resolved peers, not a standalone CLI command.
let mut job = envelope();
job.job_type = PERFORMANCE_NETWORK_JOB_TYPE.to_owned();
let (job_type, kind) = match capability {
NETWORK_CAPABILITY => (PERFORMANCE_NETWORK_JOB_TYPE, DiagnosticJobKind::PerformanceNetwork),
MEMORY_PROFILE_SERVICE_CAPABILITY => (PROFILE_MEMORY_JOB_TYPE, DiagnosticJobKind::ProfileMemory),
_ => panic!("{capability} is missing service dispatch"),
};
job.job_type = job_type.to_owned();
job.required_capabilities = vec![capability.to_owned()];
job.schema_version = NETWORK_SCHEMA_VERSION;
assert_eq!(job.kind(), Ok(DiagnosticJobKind::PerformanceNetwork));
assert_eq!(job.kind(), Ok(kind));
continue;
}
let mut command = connect;
+1
View File
@@ -48,6 +48,7 @@ pub const CONNECT_DIAGNOSTIC_CAPABILITIES: &[&str] = &[
logs::LOGS_CAPABILITY,
profile_cpu::CPU_PROFILE_CAPABILITY,
profile_cpu::MEMORY_PROFILE_CAPABILITY,
profile_memory::MEMORY_PROFILE_SERVICE_CAPABILITY,
profile_cpu::THREAD_PROFILE_CAPABILITY,
trace_record::TELEMETRY_RECORD_CAPABILITY,
trace_record::TELEMETRY_OTLP_CAPABILITY,
@@ -34,6 +34,9 @@ use super::profile_cpu::{
const MAX_ALLOCATOR_STATS_BYTES: usize = 262_144;
/// Service-job negotiation is separate from the unchanged memory result schema.
pub const MEMORY_PROFILE_SERVICE_CAPABILITY: &str = "profile.memory.service@1";
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct AllocationSnapshot {
total_allocated_bytes: u64,
+9 -1
View File
@@ -111,7 +111,15 @@ impl PendingHeartbeat {
super::diagnostics::CPU_PROFILE_CAPABILITY,
]
|| self.capabilities == heartbeat_capabilities(false)
|| self.capabilities == heartbeat_capabilities(true))
|| self.capabilities == heartbeat_capabilities(true)
// RUSTFS_COMPAT_TODO(connect-894) Remove after upgrades from the pre-service-memory capability set are unsupported.
// Preserve exact pending requests; never add the capability to a retry.
|| [false, true].into_iter().any(|job_capable| {
heartbeat_capabilities(job_capable)
.iter()
.filter(|capability| capability.as_str() != "profile.memory.service@1")
.eq(self.capabilities.iter())
}))
&& self.sequence <= MAX_SEQUENCE
&& self.coarse_node_summary.is_valid()
&& is_exact_utc_seconds(&self.client_time)
+111
View File
@@ -544,6 +544,7 @@ async fn sends_only_l0_fields_and_accepts_additive_response_fields() {
"logs.capture@1",
"profile.cpu@1",
"profile.memory@1",
"profile.memory.service@1",
"profile.threads@1",
"telemetry.record@1",
"telemetry.otlp@1",
@@ -642,6 +643,116 @@ 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));
}
fn pre_service_memory_capabilities(job_capable: bool) -> Vec<&'static str> {
let mut capabilities = vec!["heartbeat", "diagnostics.policy.v1", "inventory.environment@1"];
if job_capable {
capabilities.push("jobs");
}
// Freeze the preceding release, independently of today's advertisement.
capabilities.extend([
"performance.client@1",
"performance.drive@1",
"performance.network@1",
"performance.object@1",
"performance.siteReplication@1",
"logs.capture@1",
"profile.cpu@1",
"profile.memory@1",
"profile.threads@1",
"telemetry.record@1",
"telemetry.otlp@1",
"telemetry.replay@1",
"top.api@1",
"top.disk@1",
"top.locks@1",
"top.net@1",
"top.rpc@1",
"inspect.object@1",
]);
capabilities
}
fn pending_heartbeat_with_capabilities(capabilities: &[&str]) -> Value {
json!({
"protocolVersion": "v1",
"requestId": "550e8400-e29b-41d4-a716-446655440000",
"agentVersion": format!("rustfs-agent/{}", env!("CARGO_PKG_VERSION")),
"capabilities": capabilities,
"sequence": 0,
"clientTime": "2026-08-22T01:02:03Z",
"coarseNodeSummary": {"total": 1, "healthy": 1, "degraded": 0}
})
}
#[tokio::test]
async fn restart_replays_exact_pre_service_memory_capabilities_with_and_without_jobs() {
for job_capable in [false, true] {
let pki = TestPki::new();
let server = server(&pki, vec![Reply::ok("2026-08-22T01:02:03Z")]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let config = config(&temp, &pki, &server);
fs::create_dir_all(config.state_path.parent().expect("state directory")).expect("create state directory");
let pending = pending_heartbeat_with_capabilities(&pre_service_memory_capabilities(job_capable));
let state = json!({"nextSequence": 0, "pending": pending});
fs::write(&config.state_path, serde_json::to_vec(&state).expect("heartbeat state JSON")).expect("write heartbeat state");
private_mode(&config.state_path);
let runtime = spawn_heartbeat_runtime(Some(config), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
HeartbeatStatus::Online { .. }
));
runtime.shutdown().await;
let seen = server.seen.lock().expect("seen lock");
assert_eq!(seen.len(), 1);
assert_eq!(seen[0], pending);
}
}
#[tokio::test]
async fn pre_service_memory_compatibility_does_not_accept_changed_capabilities() {
for job_capable in [false, true] {
for mutation in ["unknown", "missing", "duplicate", "reordered"] {
let pki = TestPki::new();
let server = server(&pki, vec![]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let config = config(&temp, &pki, &server);
fs::create_dir_all(config.state_path.parent().expect("state directory")).expect("create state directory");
let mut capabilities = pre_service_memory_capabilities(job_capable);
match mutation {
"unknown" => capabilities.push("shell.exec@1"),
"missing" => {
capabilities.pop();
}
"duplicate" => capabilities.push("profile.memory@1"),
"reordered" => capabilities.swap(5, 6),
_ => unreachable!(),
}
let state = json!({"nextSequence": 0, "pending": pending_heartbeat_with_capabilities(&capabilities)});
fs::write(&config.state_path, serde_json::to_vec(&state).expect("heartbeat state JSON"))
.expect("write heartbeat state");
private_mode(&config.state_path);
let runtime = spawn_heartbeat_runtime(Some(config), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
assert!(
matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
HeartbeatStatus::Failed { reason } if reason.contains("violates the protocol invariants")
),
"accepted altered persisted capabilities: {mutation}, jobs={job_capable}"
);
assert!(server.seen.lock().expect("seen lock").is_empty());
runtime.shutdown().await;
}
}
}
#[tokio::test]
async fn restart_sends_the_last_durable_diagnostic_receipt_outbound() {
let pki = TestPki::new();
+235 -1
View File
@@ -30,6 +30,8 @@ use std::sync::Mutex;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use base64_simd::URL_SAFE_NO_PAD;
use chrono::{SecondsFormat, Utc};
use ed25519_dalek::{Signer as _, SigningKey};
use p256::ecdsa::signature::Verifier as _;
use p256::ecdsa::{Signature, VerifyingKey};
use p256::pkcs8::DecodePublicKey as _;
@@ -37,7 +39,13 @@ use profile_cpu::{
LocalProfileConsent, MEMORY_PROFILE_CAPABILITY, ProfileCaptureRequest, ProfileError, ProfileProvenance,
save_signed_profile_export,
};
use profile_memory::{AllocationProfileSource, export_memory_profile, export_memory_profile_from, parse_allocator_stats};
use profile_memory::{
AllocationProfileSource, MEMORY_PROFILE_SERVICE_CAPABILITY, export_memory_profile, export_memory_profile_from,
parse_allocator_stats,
};
use rustfs::connect::{
DiagnosticJobEnvelope, DiagnosticJobError, DiagnosticJobTarget, TrustedDiagnosticJobSigner, execute_diagnostic_job,
};
use sha2::{Digest as _, Sha256};
use tokio::sync::Mutex as AsyncMutex;
use tokio_util::sync::CancellationToken;
@@ -117,6 +125,232 @@ fn archive_entry(archive: &mut ZipArchive<Cursor<Vec<u8>>>, name: &str) -> Vec<u
bytes
}
fn service_job() -> DiagnosticJobEnvelope {
let request = request();
let now = Utc::now();
serde_json::from_value(serde_json::json!({
"jobId": request.run_uid,
"protocolVersion": "v1",
"jobType": "profile.memory",
"schemaVersion": 1,
"organizationName": request.organization_name,
"clusterName": request.cluster_name,
"deviceName": request.device_name,
"createTime": now.to_rfc3339_opts(SecondsFormat::Secs, true),
"expireTime": (now + chrono::Duration::seconds(30)).to_rfc3339_opts(SecondsFormat::Secs, true),
"nonce": URL_SAFE_NO_PAD.encode_to_string(request.nonce),
"requiredCapabilities": [MEMORY_PROFILE_SERVICE_CAPABILITY],
"authorization": {
"actorType": "BROWSER_USER",
"actorName": "users/019e3ae0-0000-7000-8000-000000000018",
"requestId": "123e4567-e89b-42d3-a456-426614174001"
},
"limits": {
"timeoutSeconds": 30,
"maxOutputBytes": 524288,
"maxMemoryBytes": 67108864,
"maxCpuMillis": 5000
},
"parameters": {
"artifactUid": request.artifact_uid,
"consentUid": request.consent.consent_uid,
"consentPolicyRevision": request.consent.policy_revision,
"consentExpiresAt": (now + chrono::Duration::seconds(60)).to_rfc3339_opts(SecondsFormat::Secs, true),
"durationMillis": 1000,
"samplePeriodMicros": 10000
},
"signature": {"algorithm": "Ed25519", "keyId": "", "value": ""}
}))
.expect("service job fixture")
}
fn sign_service_job(job: DiagnosticJobEnvelope) -> (DiagnosticJobEnvelope, TrustedDiagnosticJobSigner) {
let key = SigningKey::from_bytes(&[0x17; 32]);
let key_id = hex(&Sha256::digest(key.verifying_key().as_bytes()));
let signer = TrustedDiagnosticJobSigner::new(key_id.clone(), *key.verifying_key().as_bytes()).expect("trusted signer");
let signature = URL_SAFE_NO_PAD.encode_to_string(key.sign(&job.signing_payload().expect("job payload")).to_bytes());
let mut value = serde_json::to_value(job).expect("job JSON");
value["signature"]["keyId"] = key_id.into();
value["signature"]["value"] = signature.into();
(serde_json::from_value(value).expect("signed job"), signer)
}
fn service_target(job: &DiagnosticJobEnvelope) -> DiagnosticJobTarget {
DiagnosticJobTarget {
organization_name: job.organization_name.clone(),
cluster_name: job.cluster_name.clone(),
device_name: job.device_name.clone(),
}
}
#[test]
fn service_job_requires_the_independently_negotiated_capability_and_signed_bounds() {
let (job, signer) = sign_service_job(service_job());
let target = service_target(&job);
signer.verify(&job, &target, Utc::now()).expect("bounded service job");
for capabilities in [
vec![MEMORY_PROFILE_CAPABILITY.to_owned()],
vec![
MEMORY_PROFILE_SERVICE_CAPABILITY.to_owned(),
MEMORY_PROFILE_CAPABILITY.to_owned(),
],
vec![],
] {
let mut candidate = job.clone();
candidate.required_capabilities = capabilities;
let (candidate, signer) = sign_service_job(candidate);
assert_eq!(signer.verify(&candidate, &target, Utc::now()), Err(DiagnosticJobError::Unsupported));
}
let mut limits = vec![job.limits.clone(); 4];
limits[0].max_cpu_millis = 5001;
limits[1].max_memory_bytes = 67108863;
limits[2].max_output_bytes = 524289;
limits[3].timeout_seconds = 31;
for limit in limits {
let mut candidate = job.clone();
candidate.limits = limit;
let (candidate, signer) = sign_service_job(candidate);
assert_eq!(signer.verify(&candidate, &target, Utc::now()), Err(DiagnosticJobError::LimitExceeded));
}
let mut tampered = job.clone();
tampered.parameters.sample_period_micros += 1;
assert_eq!(signer.verify(&tampered, &target, Utc::now()), Err(DiagnosticJobError::SignatureInvalid));
let mut wrong_target = target.clone();
wrong_target.organization_name.push('x');
assert_eq!(signer.verify(&job, &wrong_target, Utc::now()), Err(DiagnosticJobError::TargetMismatch));
wrong_target = target.clone();
wrong_target.device_name.push('x');
assert_eq!(signer.verify(&job, &wrong_target, Utc::now()), Err(DiagnosticJobError::TargetMismatch));
let mut expired = job;
expired.parameters.consent_expires_at =
(Utc::now() - chrono::Duration::seconds(1)).to_rfc3339_opts(SecondsFormat::Secs, true);
let (expired, signer) = sign_service_job(expired);
assert_eq!(signer.verify(&expired, &target, Utc::now()), Err(DiagnosticJobError::Expired));
}
#[tokio::test]
async fn signed_service_memory_job_exports_the_same_process_allocator_with_original_bindings() {
let _guard = TEST_PROFILE_LOCK.lock().await;
let (job, signer) = sign_service_job(service_job());
let verified = signer
.verify(&job, &service_target(&job), Utc::now())
.expect("verified memory job");
let identity = connect::DeviceIdentity::generate();
let provenance = rustfs::connect::ProfileProvenance::new("c".repeat(40), "d".repeat(64), "1.0.0-rc.6", vec![]);
// This invokes the service dispatcher directly in this mimalloc-using process,
// not the standalone CLI or a fixture allocation source.
let execution = execute_diagnostic_job(verified, &identity, provenance, &CancellationToken::new()).await;
#[cfg(target_os = "windows")]
{
assert_eq!(execution, Err(DiagnosticJobError::ProfileSourceUnavailable));
}
#[cfg(not(target_os = "windows"))]
{
let execution = execution.expect("service allocator capture");
assert_eq!(execution.job_id, job.job_id);
assert_eq!(execution.outcome, "SUCCEEDED");
assert_eq!(execution.reason, "COMPLETE");
assert_eq!(execution.artifact_uid.as_deref(), Some(job.parameters.artifact_uid.as_str()));
let bytes = execution.artifact_bytes.expect("signed artifact");
assert_eq!(execution.artifact_sha256, Some(hex(&Sha256::digest(&bytes))));
let mut archive = ZipArchive::new(Cursor::new(bytes)).expect("profile archive");
let envelope_bytes = archive_entry(&mut archive, "envelope.json");
let envelope: serde_json::Value = serde_json::from_slice(&envelope_bytes).expect("envelope JSON");
let result_bytes = archive_entry(&mut archive, "result.json");
let result: serde_json::Value = serde_json::from_slice(&result_bytes).expect("result JSON");
assert_eq!(result["toolId"], "profile.memory");
assert_eq!(result["schemaVersion"], 1);
assert_eq!(result["capability"], MEMORY_PROFILE_CAPABILITY);
assert_eq!(result["runUid"], job.job_id);
assert_eq!(result["provenance"]["repository"], "rustfs/rustfs");
assert_eq!(result["provenance"]["sourceCommit"], "c".repeat(40));
assert_eq!(result["provenance"]["executableSha256"], "d".repeat(64));
assert_eq!(result["data"]["scope"], "ALLOCATION_AGGREGATES");
assert!(result["data"]["allocatedBytes"].is_u64());
assert!(result["data"]["allocationCount"].is_u64());
assert_eq!(envelope["organizationName"], job.organization_name);
assert_eq!(envelope["clusterName"], job.cluster_name);
assert_eq!(envelope["deviceName"], job.device_name);
assert_eq!(envelope["artifactUid"], job.parameters.artifact_uid);
assert_eq!(envelope["runUid"], job.job_id);
assert_eq!(envelope["consentUid"], job.parameters.consent_uid);
assert_eq!(envelope["policyRevision"], job.parameters.consent_policy_revision);
assert_eq!(envelope["classification"], "L3");
assert_eq!(envelope["nonce"], job.nonce);
assert_eq!(envelope["payload"]["sha256"], hex(&Sha256::digest(&result_bytes)));
let signature: serde_json::Value =
serde_json::from_slice(&archive_entry(&mut archive, "envelope.sig")).expect("signature JSON");
let raw = URL_SAFE_NO_PAD
.decode_to_vec(signature["value"].as_str().expect("signature"))
.expect("base64 signature");
let signature = Signature::from_slice(&raw).expect("ES256 signature");
let public = VerifyingKey::from_public_key_der(&identity.public_key_der()).expect("device public key");
let mut input = b"rustfs-diagnostic-envelope-v1\0".to_vec();
input.extend_from_slice(&envelope_bytes);
public.verify(&input, &signature).expect("original device signed artifact");
}
}
#[tokio::test]
async fn cancelled_service_memory_job_produces_no_artifact() {
let (job, signer) = sign_service_job(service_job());
let verified = signer
.verify(&job, &service_target(&job), Utc::now())
.expect("verified memory job");
let cancel = CancellationToken::new();
cancel.cancel();
assert_eq!(
execute_diagnostic_job(
verified,
&connect::DeviceIdentity::generate(),
rustfs::connect::ProfileProvenance::new("c".repeat(40), "d".repeat(64), "1.0.0-rc.6", vec![]),
&cancel,
)
.await,
Err(DiagnosticJobError::Cancelled)
);
}
#[tokio::test]
async fn service_memory_job_honors_cancellation_during_capture_and_the_signed_output_limit() {
let _guard = TEST_PROFILE_LOCK.lock().await;
let mut job = service_job();
job.parameters.sample_period_micros = 500_000;
let (job, signer) = sign_service_job(job);
let verified = signer
.verify(&job, &service_target(&job), Utc::now())
.expect("verified memory job");
let identity = connect::DeviceIdentity::generate();
let provenance = rustfs::connect::ProfileProvenance::new("c".repeat(40), "d".repeat(64), "1.0.0-rc.6", vec![]);
let cancel = CancellationToken::new();
let execution = execute_diagnostic_job(verified, &identity, provenance.clone(), &cancel);
let cancellation = async {
// The first join branch polls the collector through its first snapshot
// into the bounded wait before this branch cancels the same token.
cancel.cancel();
};
let (result, ()) = tokio::join!(biased; execution, cancellation);
#[cfg(not(target_os = "windows"))]
assert_eq!(result, Err(DiagnosticJobError::Cancelled));
#[cfg(target_os = "windows")]
assert_eq!(result, Err(DiagnosticJobError::ProfileSourceUnavailable));
let mut job = service_job();
job.limits.max_output_bytes = 1;
let (job, signer) = sign_service_job(job);
let verified = signer
.verify(&job, &service_target(&job), Utc::now())
.expect("bounded output memory job");
let result = execute_diagnostic_job(verified, &identity, provenance, &CancellationToken::new()).await;
#[cfg(not(target_os = "windows"))]
assert_eq!(result, Err(DiagnosticJobError::LimitExceeded));
#[cfg(target_os = "windows")]
assert_eq!(result, Err(DiagnosticJobError::ProfileSourceUnavailable));
}
#[tokio::test]
async fn real_mimalloc_profile_produces_a_signed_three_file_export() {
let _guard = TEST_PROFILE_LOCK.lock().await;