diff --git a/.github/workflows/connect-profile-memory-acceptance.yml b/.github/workflows/connect-profile-memory-acceptance.yml new file mode 100644 index 000000000..2070f2069 --- /dev/null +++ b/.github/workflows/connect-profile-memory-acceptance.yml @@ -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 diff --git a/docs/architecture/compat-cleanup-register.md b/docs/architecture/compat-cleanup-register.md index 2643f7a1a..c02543f42 100644 --- a/docs/architecture/compat-cleanup-register.md +++ b/docs/architecture/compat-cleanup-register.md @@ -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. diff --git a/protocol/agent/v1/fixtures/heartbeat/MANIFEST.sha256 b/protocol/agent/v1/fixtures/heartbeat/MANIFEST.sha256 index f5271d17f..de8e25c91 100644 --- a/protocol/agent/v1/fixtures/heartbeat/MANIFEST.sha256 +++ b/protocol/agent/v1/fixtures/heartbeat/MANIFEST.sha256 @@ -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 diff --git a/protocol/agent/v1/fixtures/heartbeat/producer-capabilities.json b/protocol/agent/v1/fixtures/heartbeat/producer-capabilities.json index 716d651f2..27571c0fa 100644 --- a/protocol/agent/v1/fixtures/heartbeat/producer-capabilities.json +++ b/protocol/agent/v1/fixtures/heartbeat/producer-capabilities.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", diff --git a/protocol/agent/v1/fixtures/manifest.sha256 b/protocol/agent/v1/fixtures/manifest.sha256 index 4ffe6e680..199ae5862 100644 --- a/protocol/agent/v1/fixtures/manifest.sha256 +++ b/protocol/agent/v1/fixtures/manifest.sha256 @@ -1,7 +1,7 @@ a499742c06046e9cf2f781fc17d8b4cf95e018ab45fa41ee5944ba15b0e2902c auth 9b54e39584a442270ae486e8378df1d6f168c45d0c59baba9459b46bb212d2df version 159d116e965b51771f1687fc6ec9bdfbf2a4a0e6d6aa26d3f1046b3dac2697dd registration -fe1287eac16832bed88320254245ba07a985e6ae7e66e1d1570e324ffa972873 heartbeat +9b3862f5160968a2ff281afb3ea519a3d49ae1bc3d767153d70f309d2a7457a7 heartbeat cc6c3f3bf5e6a938f2e00b8ce42d8fa7bfb974f4eb2927101d4ede4ebdc10dd5 inventory 0207458ec9b97b368e2f347da112511a3aabd27d752890cacda4fa2408e7f19e object-performance a464d4c695d76dd985c0f2abf68c2135c080221098ad2427022f85243a7bf6ae network-performance diff --git a/rustfs/src/connect/diagnostics/job.rs b/rustfs/src/connect/diagnostics/job.rs index 62466dd3c..2db29c11f 100644 --- a/rustfs/src/connect/diagnostics/job.rs +++ b/rustfs/src/connect/diagnostics/job.rs @@ -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 { + 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; diff --git a/rustfs/src/connect/diagnostics/mod.rs b/rustfs/src/connect/diagnostics/mod.rs index e4ec5241f..671e366e1 100644 --- a/rustfs/src/connect/diagnostics/mod.rs +++ b/rustfs/src/connect/diagnostics/mod.rs @@ -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, diff --git a/rustfs/src/connect/diagnostics/profile_memory.rs b/rustfs/src/connect/diagnostics/profile_memory.rs index 0ff38a78e..2cb23789f 100644 --- a/rustfs/src/connect/diagnostics/profile_memory.rs +++ b/rustfs/src/connect/diagnostics/profile_memory.rs @@ -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, diff --git a/rustfs/src/connect/heartbeat.rs b/rustfs/src/connect/heartbeat.rs index 247d8242b..27d70a625 100644 --- a/rustfs/src/connect/heartbeat.rs +++ b/rustfs/src/connect/heartbeat.rs @@ -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) diff --git a/rustfs/tests/connect_heartbeat.rs b/rustfs/tests/connect_heartbeat.rs index 7b5b1034e..16f38276c 100644 --- a/rustfs/tests/connect_heartbeat.rs +++ b/rustfs/tests/connect_heartbeat.rs @@ -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(); diff --git a/rustfs/tests/connect_profile_memory.rs b/rustfs/tests/connect_profile_memory.rs index 42e25dbd6..442b8338e 100644 --- a/rustfs/tests/connect_profile_memory.rs +++ b/rustfs/tests/connect_profile_memory.rs @@ -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>>, name: &str) -> Vec 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;