mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 12:26:37 +00:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 637d4f5e45 | |||
| a39ce768cf | |||
| d3515f5109 | |||
| 8c0eaf225d |
@@ -78,12 +78,6 @@ test-group = 'embedded-test-ports'
|
||||
filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persists_durable_job_progress)'
|
||||
test-group = 'ecstore-serial-flaky'
|
||||
|
||||
# The durable ILM decommission regressions build isolated multi-pool stores and
|
||||
# deliberately take source or target disks offline while checking fencing.
|
||||
[[profile.default.overrides]]
|
||||
filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))'
|
||||
test-group = 'ecstore-serial-flaky'
|
||||
|
||||
# Serialize the bucket-incarnation / lifecycle-fence tests. They drive
|
||||
# init_bucket_metadata_sys and bucket_metadata_sys_of, i.e. process-global
|
||||
# OnceLock state that serial_test's #[serial] cannot protect across nextest's
|
||||
@@ -196,10 +190,6 @@ test-group = 'embedded-test-ports'
|
||||
filter = 'package(rustfs-ecstore) & test(manual_transition_page_checkpoint_persists_durable_job_progress)'
|
||||
test-group = 'ecstore-serial-flaky'
|
||||
|
||||
[[profile.ci.overrides]]
|
||||
filter = 'package(rustfs-ecstore) & (test(decommission_migrates_and_verifies_registered_durable_ilm_records) | test(decommission_durable_ilm_target_read_error_is_not_masked_by_peer_success) | test(decommission_durable_ilm_terminal_receipt_recovers_failed_source_cleanup) | test(decommission_durable_ilm_receipt_pagination_fails_closed_on_second_page) | test(decommission_durable_ilm_recovery_keeps_multiple_active_sources))'
|
||||
test-group = 'ecstore-serial-flaky'
|
||||
|
||||
# Serialize the bucket-incarnation / lifecycle-fence tests under the ci profile
|
||||
# too (see the matching default-profile override near the top). No retries.
|
||||
[[profile.ci.overrides]]
|
||||
|
||||
@@ -14,9 +14,10 @@
|
||||
|
||||
name: "Schedule Failure Issue"
|
||||
description: >-
|
||||
Open (or update) a tracking issue when a scheduled workflow run fails.
|
||||
Open (or update) a tracking issue when a scheduled workflow run fails or
|
||||
does not complete normally.
|
||||
Dedupes by workflow name: if an open issue titled
|
||||
"[scheduled-failure] <workflow name>" already exists, the failure is
|
||||
"[scheduled-failure] <workflow name>" already exists, the result is
|
||||
appended as a comment; otherwise a new issue is created. This is the
|
||||
single alerting mechanism for all scheduled pipelines (backlog#1149 ci-8).
|
||||
|
||||
@@ -38,6 +39,26 @@ inputs:
|
||||
Set to an empty string to skip labeling.
|
||||
required: false
|
||||
default: "infrastructure"
|
||||
source-run-id:
|
||||
description: "Run ID to report. Defaults to the current workflow run."
|
||||
required: false
|
||||
default: ${{ github.run_id }}
|
||||
source-run-attempt:
|
||||
description: "Run attempt to report. Defaults to the current attempt."
|
||||
required: false
|
||||
default: ${{ github.run_attempt }}
|
||||
source-event:
|
||||
description: "Trigger event of the run being reported."
|
||||
required: false
|
||||
default: ${{ github.event_name }}
|
||||
source-ref-name:
|
||||
description: "Ref name of the run being reported."
|
||||
required: false
|
||||
default: ${{ github.ref_name }}
|
||||
source-sha:
|
||||
description: "Commit SHA of the run being reported."
|
||||
required: false
|
||||
default: ${{ github.sha }}
|
||||
|
||||
runs:
|
||||
using: "composite"
|
||||
@@ -48,17 +69,21 @@ runs:
|
||||
GH_TOKEN: ${{ inputs.github-token }}
|
||||
WORKFLOW_NAME: ${{ inputs.workflow-name }}
|
||||
ISSUE_LABEL: ${{ inputs.label }}
|
||||
SOURCE_RUN_ID: ${{ inputs.source-run-id }}
|
||||
SOURCE_RUN_ATTEMPT: ${{ inputs.source-run-attempt }}
|
||||
SOURCE_EVENT: ${{ inputs.source-event }}
|
||||
SOURCE_REF_NAME: ${{ inputs.source-ref-name }}
|
||||
SOURCE_SHA: ${{ inputs.source-sha }}
|
||||
run: |
|
||||
set -euo pipefail
|
||||
|
||||
title="[scheduled-failure] ${WORKFLOW_NAME}"
|
||||
run_url="${GITHUB_SERVER_URL}/${GITHUB_REPOSITORY}/actions/runs/${GITHUB_RUN_ID}"
|
||||
run_url="${GITHUB_SERVER_URL}/${GITHUB_REPOSITORY}/actions/runs/${SOURCE_RUN_ID}"
|
||||
|
||||
# Failed job names for this run attempt. The alert job runs while the
|
||||
# run as a whole is still in progress, so inspect the jobs that have
|
||||
# already completed with a non-success conclusion.
|
||||
# Inspect the reported run attempt. It can be the current in-workflow
|
||||
# failure or a completed run observed by the external watchdog.
|
||||
failed_jobs="$(gh api \
|
||||
"repos/${GITHUB_REPOSITORY}/actions/runs/${GITHUB_RUN_ID}/attempts/${GITHUB_RUN_ATTEMPT}/jobs" \
|
||||
"repos/${GITHUB_REPOSITORY}/actions/runs/${SOURCE_RUN_ID}/attempts/${SOURCE_RUN_ATTEMPT}/jobs" \
|
||||
--paginate \
|
||||
--jq '.jobs[]
|
||||
| select(.conclusion == "failure" or .conclusion == "timed_out" or .conclusion == "cancelled")
|
||||
@@ -68,13 +93,13 @@ runs:
|
||||
fi
|
||||
|
||||
body="$(cat <<EOF
|
||||
Scheduled run of **${WORKFLOW_NAME}** failed.
|
||||
Scheduled run of **${WORKFLOW_NAME}** did not complete successfully.
|
||||
|
||||
- Run: ${run_url} (attempt ${GITHUB_RUN_ATTEMPT})
|
||||
- Event: \`${GITHUB_EVENT_NAME}\`
|
||||
- Ref: \`${GITHUB_REF_NAME}\` @ \`${GITHUB_SHA}\`
|
||||
- Run: ${run_url} (attempt ${SOURCE_RUN_ATTEMPT})
|
||||
- Event: \`${SOURCE_EVENT}\`
|
||||
- Ref: \`${SOURCE_REF_NAME}\` @ \`${SOURCE_SHA}\`
|
||||
|
||||
Failed jobs:
|
||||
Non-success jobs:
|
||||
${failed_jobs}
|
||||
EOF
|
||||
)"
|
||||
|
||||
@@ -1032,3 +1032,23 @@ jobs:
|
||||
|
||||
echo "🎉 Released $TAG successfully!"
|
||||
echo "📄 Release URL: ${{ needs.create-release.outputs.release_url }}"
|
||||
|
||||
alert-on-failure:
|
||||
name: Alert on scheduled failure
|
||||
needs: [build-check, prepare-platform-matrix, build-rustfs, build-summary]
|
||||
if: >-
|
||||
always() && github.event_name == 'schedule' &&
|
||||
(contains(needs.*.result, 'failure') || contains(needs.*.result, 'cancelled'))
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 10
|
||||
permissions:
|
||||
contents: read
|
||||
issues: write
|
||||
steps:
|
||||
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
- name: Open or update failure-tracking issue
|
||||
uses: ./.github/actions/schedule-failure-issue
|
||||
with:
|
||||
github-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
@@ -1032,3 +1032,37 @@ jobs:
|
||||
path: artifacts/s3tests-single/**
|
||||
if-no-files-found: ignore
|
||||
retention-days: 3
|
||||
|
||||
alert-on-failure:
|
||||
name: Alert on scheduled failure
|
||||
needs:
|
||||
- typos
|
||||
- quick-checks
|
||||
- test-and-lint
|
||||
- test-ilm-integration-serial
|
||||
- test-and-lint-rio-v2
|
||||
- test-and-lint-protocols
|
||||
- build-rustfs-debug-binary
|
||||
- build-rustfs-debug-binary-rio-v2
|
||||
- uring-integration
|
||||
- e2e-tests
|
||||
- e2e-full
|
||||
- e2e-tests-rio-v2
|
||||
- s3-implemented-tests
|
||||
- s3-lifecycle-behavior-tests
|
||||
if: >-
|
||||
always() && github.event_name == 'schedule' &&
|
||||
(contains(needs.*.result, 'failure') || contains(needs.*.result, 'cancelled'))
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 10
|
||||
permissions:
|
||||
contents: read
|
||||
issues: write
|
||||
steps:
|
||||
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
- name: Open or update failure-tracking issue
|
||||
uses: ./.github/actions/schedule-failure-issue
|
||||
with:
|
||||
github-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
@@ -121,3 +121,21 @@ jobs:
|
||||
cargo nextest run --run-ignored ignored-only --no-tests=fail \
|
||||
-p "$INTEROP_PACKAGE" --features "$INTEROP_FEATURES" \
|
||||
-E "$INTEROP_FILTER"
|
||||
|
||||
alert-on-failure:
|
||||
name: Alert on scheduled failure
|
||||
needs: [minio-interop]
|
||||
if: always() && github.event_name == 'schedule' && contains(needs.*.result, 'failure')
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 10
|
||||
permissions:
|
||||
contents: read
|
||||
issues: write
|
||||
steps:
|
||||
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
- name: Open or update failure-tracking issue
|
||||
uses: ./.github/actions/schedule-failure-issue
|
||||
with:
|
||||
github-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
@@ -194,3 +194,23 @@ jobs:
|
||||
|
||||
- name: Run HA leader failover live checks (three-node Raft cluster in Docker)
|
||||
run: bash scripts/test/vault_ha_kms_live.sh
|
||||
|
||||
alert-on-failure:
|
||||
name: Alert on scheduled failure
|
||||
needs: [build, kms-vault-lane, kms-vault-ha-failover]
|
||||
if: >-
|
||||
always() && github.event_name == 'schedule' &&
|
||||
(contains(needs.*.result, 'failure') || contains(needs.*.result, 'cancelled'))
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 10
|
||||
permissions:
|
||||
contents: read
|
||||
issues: write
|
||||
steps:
|
||||
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
- name: Open or update failure-tracking issue
|
||||
uses: ./.github/actions/schedule-failure-issue
|
||||
with:
|
||||
github-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
# 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.
|
||||
|
||||
name: Scheduled Validation Watchdog
|
||||
|
||||
on:
|
||||
workflow_run:
|
||||
workflows:
|
||||
- "Security Audit"
|
||||
- "Build and Release"
|
||||
- "Continuous Integration"
|
||||
- "coverage"
|
||||
- "e2e-nightly"
|
||||
- "e2e-s3tests"
|
||||
- "Fuzz"
|
||||
- "mint"
|
||||
- "minio-interop"
|
||||
- "Nightly GNU Build"
|
||||
- "Performance A/B"
|
||||
- "Runner Hygiene"
|
||||
types: [completed]
|
||||
|
||||
permissions:
|
||||
contents: read
|
||||
|
||||
jobs:
|
||||
alert-on-incomplete-run:
|
||||
name: Alert on incomplete scheduled run
|
||||
if: >-
|
||||
github.event.workflow_run.event == 'schedule' &&
|
||||
github.event.workflow_run.conclusion != 'success' &&
|
||||
github.event.workflow_run.conclusion != 'failure'
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 10
|
||||
permissions:
|
||||
actions: read
|
||||
contents: read
|
||||
issues: write
|
||||
steps:
|
||||
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
- name: Open or update incomplete-run issue
|
||||
uses: ./.github/actions/schedule-failure-issue
|
||||
with:
|
||||
github-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
workflow-name: ${{ github.event.workflow_run.name }}
|
||||
source-run-id: ${{ github.event.workflow_run.id }}
|
||||
source-run-attempt: ${{ github.event.workflow_run.run_attempt }}
|
||||
source-event: ${{ github.event.workflow_run.event }}
|
||||
source-ref-name: ${{ github.event.workflow_run.head_branch }}
|
||||
source-sha: ${{ github.event.workflow_run.head_sha }}
|
||||
@@ -9136,7 +9136,6 @@ mod tests {
|
||||
assert_eq!(loaded.report.scanned, 37);
|
||||
assert_eq!(loaded.report.eligible, 11);
|
||||
assert_eq!(loaded.report.enqueued, 5);
|
||||
assert_eq!(loaded.cursor_revision, Some(1));
|
||||
assert!(loaded.lease_expires_at_unix_nanos > 0);
|
||||
let token = loaded
|
||||
.report
|
||||
@@ -9155,30 +9154,6 @@ mod tests {
|
||||
assert_eq!(admission.lease_id, loaded.lease_id);
|
||||
assert_eq!(admission.lease_expires_at_unix_nanos, loaded.lease_expires_at_unix_nanos);
|
||||
|
||||
let mut same_marker_report = report.clone();
|
||||
same_marker_report.scanned += 1;
|
||||
persist_manual_transition_page_checkpoint(
|
||||
&checkpoint_options,
|
||||
&same_marker_report,
|
||||
Some("logs/page-end".to_string()),
|
||||
Some("opaque-next-version".to_string()),
|
||||
)
|
||||
.await
|
||||
.expect("same-marker version checkpoint should persist through the durable progress sink");
|
||||
let same_marker_checkpointed = load_manual_transition_job_record(ecstore.clone(), job_id)
|
||||
.await
|
||||
.expect("same-marker version checkpoint should reload");
|
||||
assert_eq!(same_marker_checkpointed.cursor_revision, Some(2));
|
||||
let (_, version_marker) = decode_manual_transition_continuation_token(
|
||||
same_marker_checkpointed
|
||||
.report
|
||||
.continuation_token
|
||||
.as_deref()
|
||||
.expect("same-marker version checkpoint should persist a cursor"),
|
||||
)
|
||||
.expect("same-marker version cursor should decode");
|
||||
assert_eq!(version_marker.as_deref(), Some("opaque-next-version"));
|
||||
|
||||
create_test_bucket(&ecstore, &bucket).await;
|
||||
let lifecycle_xml = format!(
|
||||
r#"<?xml version="1.0" encoding="UTF-8"?>
|
||||
@@ -9235,7 +9210,6 @@ mod tests {
|
||||
assert_eq!(checkpointed.report.scanned, 1000);
|
||||
assert_eq!(checkpointed.report.eligible, 1000);
|
||||
assert_eq!(checkpointed.report.dry_run_eligible, 1000);
|
||||
assert_eq!(checkpointed.cursor_revision, Some(3));
|
||||
let token = checkpointed
|
||||
.report
|
||||
.continuation_token
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -24,10 +24,6 @@ use crate::bucket::lifecycle::bucket_lifecycle_ops::{
|
||||
ManualTransitionQueueSnapshot, ManualTransitionRunOptions, ManualTransitionRunReport,
|
||||
};
|
||||
use crate::bucket::lifecycle::config_boundary;
|
||||
use crate::bucket::lifecycle::durable_namespace::{
|
||||
MANUAL_TRANSITION_JOB_NAMESPACE, MANUAL_TRANSITION_SCOPE_NAMESPACE, MANUAL_TRANSITION_TASK_NAMESPACE,
|
||||
MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE,
|
||||
};
|
||||
use crate::disk::RUSTFS_META_BUCKET;
|
||||
use crate::error::{Error, Result as EcstoreResult};
|
||||
use crate::object_api::ObjectOptions;
|
||||
@@ -38,10 +34,10 @@ use crate::store::ECStore;
|
||||
pub const MANUAL_TRANSITION_JOB_SCHEMA: &str = "rustfs-manual-transition-job-v1";
|
||||
pub const MANUAL_TRANSITION_TASK_SCHEMA: &str = "rustfs-manual-transition-task-v1";
|
||||
pub const MANUAL_TRANSITION_WORKER_RESULT_SCHEMA: &str = "rustfs-manual-transition-worker-result-v1";
|
||||
pub const MANUAL_TRANSITION_JOB_RECORD_PREFIX: &str = MANUAL_TRANSITION_JOB_NAMESPACE.prefix;
|
||||
pub const MANUAL_TRANSITION_SCOPE_RECORD_PREFIX: &str = MANUAL_TRANSITION_SCOPE_NAMESPACE.prefix;
|
||||
pub const MANUAL_TRANSITION_TASK_PREFIX: &str = MANUAL_TRANSITION_TASK_NAMESPACE.prefix;
|
||||
pub const MANUAL_TRANSITION_WORKER_RESULT_PREFIX: &str = MANUAL_TRANSITION_WORKER_RESULT_NAMESPACE.prefix;
|
||||
pub const MANUAL_TRANSITION_JOB_RECORD_PREFIX: &str = "ilm/manual-transition/jobs";
|
||||
pub const MANUAL_TRANSITION_SCOPE_RECORD_PREFIX: &str = "ilm/manual-transition/scopes";
|
||||
pub const MANUAL_TRANSITION_TASK_PREFIX: &str = "ilm/manual-transition/tasks";
|
||||
pub const MANUAL_TRANSITION_WORKER_RESULT_PREFIX: &str = "ilm/manual-transition/results";
|
||||
pub const MAX_MANUAL_TRANSITION_JOB_RECORD_SIZE: usize = 64 * 1024;
|
||||
pub const MAX_MANUAL_TRANSITION_TASK_RECORD_SIZE: usize = 16 * 1024;
|
||||
pub const MAX_MANUAL_TRANSITION_WORKER_RESULT_RECORD_SIZE: usize = 8 * 1024;
|
||||
@@ -199,8 +195,6 @@ pub struct ManualTransitionJobRecord {
|
||||
pub updated_at_unix_nanos: i128,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub completed_at_unix_nanos: Option<i128>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub cursor_revision: Option<u64>,
|
||||
pub report: ManualTransitionRunReport,
|
||||
pub queue_snapshot: ManualTransitionQueueSnapshot,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
@@ -230,7 +224,6 @@ impl ManualTransitionJobRecord {
|
||||
created_at_unix_nanos: now,
|
||||
updated_at_unix_nanos: now,
|
||||
completed_at_unix_nanos: None,
|
||||
cursor_revision: Some(0),
|
||||
report: ManualTransitionRunReport {
|
||||
bucket: bucket.to_string(),
|
||||
prefix: options.prefix.clone(),
|
||||
@@ -245,7 +238,7 @@ impl ManualTransitionJobRecord {
|
||||
|
||||
pub fn complete(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||
self.scan_completed = true;
|
||||
self.merge_scan_report(&report);
|
||||
self.report.merge_scan_report_preserving_worker(&report);
|
||||
self.queue_snapshot = queue_snapshot;
|
||||
self.error = None;
|
||||
self.mark_terminal_if_worker_drained();
|
||||
@@ -323,7 +316,7 @@ impl ManualTransitionJobRecord {
|
||||
}
|
||||
}
|
||||
self.queue_snapshot = queue_snapshot;
|
||||
self.advance_updated_at();
|
||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||
self.mark_terminal_if_worker_drained();
|
||||
}
|
||||
|
||||
@@ -366,14 +359,14 @@ impl ManualTransitionJobRecord {
|
||||
self.report.tier_failure = scan_tier_failure.saturating_add(transition_failed);
|
||||
self.report.tier_failure_by_reason = scan_tier_failure_by_reason;
|
||||
self.queue_snapshot = queue_snapshot;
|
||||
self.advance_updated_at();
|
||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||
self.mark_terminal_if_worker_drained();
|
||||
true
|
||||
}
|
||||
|
||||
pub fn mark_cancel_requested(&mut self) {
|
||||
self.cancel_requested = true;
|
||||
self.advance_updated_at();
|
||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||
}
|
||||
|
||||
pub fn claim_recovery_lease(&mut self, owner_id: impl Into<String>, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||
@@ -387,7 +380,7 @@ impl ManualTransitionJobRecord {
|
||||
pub fn abandon_recovery_lease(&mut self, lease_id: Uuid) {
|
||||
if self.state == ManualTransitionJobState::Running && self.lease_id == lease_id {
|
||||
self.lease_expires_at_unix_nanos = 0;
|
||||
self.advance_updated_at();
|
||||
self.updated_at_unix_nanos = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -405,7 +398,7 @@ impl ManualTransitionJobRecord {
|
||||
|
||||
pub fn renew_lease(&mut self, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||
let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||
self.updated_at_unix_nanos = self.updated_at_unix_nanos.saturating_add(1).max(now);
|
||||
self.updated_at_unix_nanos = now;
|
||||
self.lease_expires_at_unix_nanos = manual_transition_job_lease_expires_at(now);
|
||||
self.queue_snapshot = queue_snapshot;
|
||||
}
|
||||
@@ -445,18 +438,11 @@ impl ManualTransitionJobRecord {
|
||||
|
||||
pub fn update_running_progress(&mut self, report: ManualTransitionRunReport, queue_snapshot: ManualTransitionQueueSnapshot) {
|
||||
if self.state == ManualTransitionJobState::Running {
|
||||
self.merge_scan_report(&report);
|
||||
self.report.merge_scan_report_preserving_worker(&report);
|
||||
self.renew_lease(queue_snapshot);
|
||||
}
|
||||
}
|
||||
|
||||
fn merge_scan_report(&mut self, report: &ManualTransitionRunReport) {
|
||||
if self.report.continuation_token != report.continuation_token {
|
||||
self.cursor_revision = Some(self.cursor_revision.unwrap_or(0).saturating_add(1));
|
||||
}
|
||||
self.report.merge_scan_report_preserving_worker(report);
|
||||
}
|
||||
|
||||
pub fn mark_unknown_if_unowned(&mut self) {
|
||||
if self.state == ManualTransitionJobState::Running {
|
||||
self.state = ManualTransitionJobState::Unknown;
|
||||
@@ -477,13 +463,9 @@ impl ManualTransitionJobRecord {
|
||||
}
|
||||
|
||||
fn mark_updated_terminal(&mut self) {
|
||||
self.advance_updated_at();
|
||||
self.completed_at_unix_nanos = Some(self.updated_at_unix_nanos);
|
||||
}
|
||||
|
||||
fn advance_updated_at(&mut self) {
|
||||
let now = OffsetDateTime::now_utc().unix_timestamp_nanos();
|
||||
self.updated_at_unix_nanos = self.updated_at_unix_nanos.saturating_add(1).max(now);
|
||||
self.updated_at_unix_nanos = now;
|
||||
self.completed_at_unix_nanos = Some(now);
|
||||
}
|
||||
|
||||
fn mark_terminal_if_worker_drained(&mut self) {
|
||||
@@ -1127,8 +1109,7 @@ pub fn manual_transition_scope_record_object_name(scope_key: &str) -> Result<Str
|
||||
pub async fn save_manual_transition_job_record(api: Arc<ECStore>, job: &ManualTransitionJobRecord) -> EcstoreResult<()> {
|
||||
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
|
||||
let data = job.encode().map_err(manual_transition_job_store_error)?;
|
||||
config_boundary::save_config(api.clone(), &object, data.clone()).await?;
|
||||
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||
config_boundary::save_config(api, &object, data).await
|
||||
}
|
||||
|
||||
pub async fn load_manual_transition_job_record(api: Arc<ECStore>, job_id: Uuid) -> EcstoreResult<ManualTransitionJobRecord> {
|
||||
@@ -1161,9 +1142,9 @@ pub async fn save_manual_transition_job_record_if_current(
|
||||
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
|
||||
let data = job.encode().map_err(manual_transition_job_store_error)?;
|
||||
config_boundary::save_config_with_opts_quiet(
|
||||
api.clone(),
|
||||
api,
|
||||
&object,
|
||||
data.clone(),
|
||||
data,
|
||||
&ObjectOptions {
|
||||
max_parity: true,
|
||||
http_preconditions: Some(HTTPPreconditions {
|
||||
@@ -1173,8 +1154,7 @@ pub async fn save_manual_transition_job_record_if_current(
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||
.await
|
||||
}
|
||||
|
||||
/// Applies a job-record mutation with optimistic concurrency control.
|
||||
@@ -1612,9 +1592,9 @@ pub async fn save_manual_transition_scope_admission_if_absent(
|
||||
let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?;
|
||||
let data = serde_json::to_vec(admission).map_err(Error::other)?;
|
||||
config_boundary::save_config_with_opts(
|
||||
api.clone(),
|
||||
api,
|
||||
&object,
|
||||
data.clone(),
|
||||
data,
|
||||
&ObjectOptions {
|
||||
max_parity: true,
|
||||
http_preconditions: Some(HTTPPreconditions {
|
||||
@@ -1624,8 +1604,7 @@ pub async fn save_manual_transition_scope_admission_if_absent(
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||
.await
|
||||
}
|
||||
|
||||
pub async fn load_manual_transition_scope_admission(
|
||||
@@ -1663,9 +1642,9 @@ pub async fn save_manual_transition_scope_admission_if_current(
|
||||
let object = manual_transition_scope_record_object_name(&admission.scope_key).map_err(manual_transition_job_store_error)?;
|
||||
let data = serde_json::to_vec(admission).map_err(Error::other)?;
|
||||
match config_boundary::save_config_with_opts(
|
||||
api.clone(),
|
||||
api,
|
||||
&object,
|
||||
data.clone(),
|
||||
data,
|
||||
&ObjectOptions {
|
||||
max_parity: true,
|
||||
http_preconditions: Some(HTTPPreconditions {
|
||||
@@ -1681,8 +1660,7 @@ pub async fn save_manual_transition_scope_admission_if_current(
|
||||
Err(Error::PreconditionFailed)
|
||||
}
|
||||
result => result,
|
||||
}?;
|
||||
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn claim_manual_transition_scope_admission(
|
||||
@@ -1975,15 +1953,13 @@ pub async fn delete_manual_transition_scope_admission_if_current(
|
||||
job_id: Uuid,
|
||||
lease_id: Uuid,
|
||||
) -> EcstoreResult<bool> {
|
||||
let (admission, etag) = match load_manual_transition_scope_admission_with_etag(api.clone(), scope_key).await {
|
||||
Ok((admission, etag)) if admission.job_id == job_id && admission.lease_id == lease_id => (admission, etag),
|
||||
let etag = match load_manual_transition_scope_admission_with_etag(api.clone(), scope_key).await {
|
||||
Ok((admission, etag)) if admission.job_id == job_id && admission.lease_id == lease_id => etag,
|
||||
Ok(_) => return Ok(false),
|
||||
Err(Error::ConfigNotFound) => return Ok(true),
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
let object = manual_transition_scope_record_object_name(scope_key).map_err(manual_transition_job_store_error)?;
|
||||
let data = serde_json::to_vec(&admission).map_err(Error::other)?;
|
||||
api.record_durable_ilm_decommission_terminal(&object, &data).await?;
|
||||
match config_boundary::delete_config_if_match(api, &object, &etag).await {
|
||||
Ok(()) | Err(Error::ConfigNotFound) => Ok(true),
|
||||
Err(Error::PreconditionFailed) => Ok(false),
|
||||
@@ -2601,25 +2577,6 @@ mod tests {
|
||||
assert!(decoded.report.tier_failure_by_reason.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_record_decodes_legacy_cursor_without_revision() {
|
||||
let options = ManualTransitionRunOptions::default();
|
||||
let record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER);
|
||||
let encoded = record.encode().expect("job record should encode");
|
||||
let mut value: serde_json::Value = serde_json::from_slice(&encoded).expect("encoded job should be json");
|
||||
value["job"]
|
||||
.as_object_mut()
|
||||
.expect("job should be object")
|
||||
.remove("cursor_revision");
|
||||
let record_bytes = serde_json::to_vec(&value["job"]).expect("legacy job should encode");
|
||||
value["content_sha256"] = serde_json::Value::String(hex_sha256(&record_bytes, ToOwned::to_owned));
|
||||
let legacy = serde_json::to_vec(&value).expect("legacy envelope should encode");
|
||||
|
||||
let decoded = ManualTransitionJobRecord::decode(record.job_id, &legacy).expect("legacy job should decode");
|
||||
|
||||
assert_eq!(decoded.cursor_revision, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_transition_job_record_rejects_unknown_report_fields() {
|
||||
let options = ManualTransitionRunOptions::default();
|
||||
|
||||
@@ -16,7 +16,6 @@ pub mod bucket_lifecycle_audit;
|
||||
pub mod bucket_lifecycle_ops;
|
||||
mod config_boundary;
|
||||
pub mod core;
|
||||
mod durable_namespace;
|
||||
pub mod evaluator;
|
||||
pub mod manual_transition_job;
|
||||
mod metadata_boundary;
|
||||
@@ -32,8 +31,3 @@ pub mod tier_free_version_recovery;
|
||||
pub mod tier_last_day_stats;
|
||||
pub mod tier_sweeper;
|
||||
pub mod transition_transaction;
|
||||
|
||||
pub(crate) use durable_namespace::{
|
||||
DurableIlmRecordCheckpoint, ILM_META_PREFIX, ValidatedDurableIlmRecord, classify_durable_ilm_record,
|
||||
validate_durable_ilm_record,
|
||||
};
|
||||
|
||||
@@ -20,7 +20,6 @@ use tokio_util::sync::CancellationToken;
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use crate::bucket::lifecycle::config_boundary;
|
||||
use crate::bucket::lifecycle::durable_namespace::TIER_DELETE_JOURNAL_NAMESPACE;
|
||||
use crate::bucket::lifecycle::runtime_boundary;
|
||||
use crate::bucket::lifecycle::tier_sweeper::{
|
||||
Jentry, TierDeleteJournalState, TierDeleteSourceIdentity,
|
||||
@@ -50,7 +49,7 @@ const TIER_DELETE_JOURNAL_VERSION: u8 = 2;
|
||||
const TIER_DELETE_JOURNAL_EXACT_VERSION: u8 = 3;
|
||||
const TIER_DELETE_JOURNAL_STATE_VERSION: u8 = 4;
|
||||
const TIER_DELETE_JOURNAL_TRANSACTION_VERSION: u8 = 5;
|
||||
pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = TIER_DELETE_JOURNAL_NAMESPACE.prefix;
|
||||
pub(crate) const TIER_DELETE_JOURNAL_PREFIX: &str = "ilm/tier-delete-journal/";
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
@@ -433,21 +432,6 @@ async fn process_committed_tier_delete_journal_entry(api: Arc<ECStore>, je: &Jen
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
let path = tier_delete_journal_object_name(je);
|
||||
let data = encode_tier_delete_journal_entry(je).map_err(std::io::Error::other)?;
|
||||
let target_pool_indices = api
|
||||
.record_durable_ilm_decommission_terminal_target_pools(&path, &data)
|
||||
.await
|
||||
.map_err(std::io::Error::other)?;
|
||||
if let Some(target_pool_indices) = target_pool_indices {
|
||||
for target_pool_idx in target_pool_indices {
|
||||
match config_boundary::delete_config(api.pools[target_pool_idx].clone(), &path).await {
|
||||
Ok(()) | Err(Error::ConfigNotFound) => {}
|
||||
Err(err) => return Err(std::io::Error::other(err)),
|
||||
}
|
||||
}
|
||||
return Ok(());
|
||||
}
|
||||
remove_tier_delete_journal_entry(api, je).await
|
||||
}
|
||||
|
||||
|
||||
@@ -21,7 +21,6 @@ use tracing::{debug, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::bucket::lifecycle::config_boundary;
|
||||
use crate::bucket::lifecycle::durable_namespace::TRANSITION_TRANSACTION_NAMESPACE;
|
||||
use crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE;
|
||||
use crate::bucket::lifecycle::tier_sweeper::{
|
||||
delete_confirmed_transition_candidate_exact_with_lease_idempotent,
|
||||
@@ -43,7 +42,7 @@ const TRANSITION_TRANSACTION_RECOVERY_INTERVAL: Duration = Duration::from_secs(6
|
||||
const TRANSITION_TRANSACTION_RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
|
||||
pub const TRANSITION_TRANSACTION_SCHEMA: &str = "rustfs-transition-transaction-v1";
|
||||
pub const TRANSITION_TRANSACTION_PREFIX: &str = "ilm/transition-transactions";
|
||||
pub const TRANSITION_TRANSACTION_RECORD_PREFIX: &str = TRANSITION_TRANSACTION_NAMESPACE.prefix;
|
||||
pub const TRANSITION_TRANSACTION_RECORD_PREFIX: &str = "ilm/transition-transactions/records";
|
||||
pub const MAX_TRANSITION_TRANSACTION_SIZE: usize = 64 * 1024;
|
||||
|
||||
pub type Result<T> = std::result::Result<T, TransitionTransactionError>;
|
||||
@@ -585,8 +584,7 @@ pub(crate) async fn save_transition_transaction_record(
|
||||
let object =
|
||||
transition_transaction_record_object_name(transaction.transaction_id).map_err(transition_transaction_store_error)?;
|
||||
let data = transaction.encode().map_err(transition_transaction_store_error)?;
|
||||
config_boundary::save_config(api.clone(), &object, data.clone()).await?;
|
||||
api.record_durable_ilm_decommission_progress(&object, &data).await
|
||||
config_boundary::save_config(api, &object, data).await
|
||||
}
|
||||
|
||||
pub(crate) async fn load_transition_transaction_record(
|
||||
@@ -598,14 +596,8 @@ pub(crate) async fn load_transition_transaction_record(
|
||||
TransitionTransaction::decode(transaction_id, &data).map_err(transition_transaction_store_error)
|
||||
}
|
||||
|
||||
pub(crate) async fn delete_transition_transaction_record(
|
||||
api: Arc<ECStore>,
|
||||
transaction: &TransitionTransaction,
|
||||
) -> EcstoreResult<()> {
|
||||
let object =
|
||||
transition_transaction_record_object_name(transaction.transaction_id).map_err(transition_transaction_store_error)?;
|
||||
let data = transaction.encode().map_err(transition_transaction_store_error)?;
|
||||
api.record_durable_ilm_decommission_terminal(&object, &data).await?;
|
||||
pub(crate) async fn delete_transition_transaction_record(api: Arc<ECStore>, transaction_id: Uuid) -> EcstoreResult<()> {
|
||||
let object = transition_transaction_record_object_name(transaction_id).map_err(transition_transaction_store_error)?;
|
||||
match config_boundary::delete_config(api, &object).await {
|
||||
Ok(()) | Err(Error::ConfigNotFound) => Ok(()),
|
||||
Err(err) => Err(err),
|
||||
@@ -821,7 +813,7 @@ pub async fn finalize_missing_transition_transaction_for_operator(
|
||||
if probe != TransitionOperatorProbe::Missing {
|
||||
return Err(TransitionOperatorError::CandidateNotMissing(probe));
|
||||
}
|
||||
delete_transition_transaction_record(api, &transaction)
|
||||
delete_transition_transaction_record(api, transaction_id)
|
||||
.await
|
||||
.map_err(TransitionOperatorError::Store)
|
||||
}
|
||||
@@ -857,22 +849,22 @@ pub async fn process_transition_transaction_record(
|
||||
match transaction.state {
|
||||
TransitionTransactionState::Uploaded => {
|
||||
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||
}
|
||||
TransitionTransactionState::CleanupPending => match local_commit_matches_transaction(api.clone(), transaction).await {
|
||||
Ok(true) => {
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
}
|
||||
Ok(false) => {
|
||||
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||
}
|
||||
Err(err) if transition_source_is_missing(&err) => {
|
||||
delete_transition_remote_candidate(api.clone(), transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||
}
|
||||
Err(err) => Err(err),
|
||||
@@ -880,7 +872,7 @@ pub async fn process_transition_transaction_record(
|
||||
TransitionTransactionState::LocalCommitStarted => {
|
||||
match local_commit_matches_transaction(api.clone(), transaction).await {
|
||||
Ok(true) => {
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
}
|
||||
Ok(false) => Ok(TransitionTransactionRecoveryOutcome::Retained),
|
||||
@@ -889,7 +881,7 @@ pub async fn process_transition_transaction_record(
|
||||
}
|
||||
}
|
||||
TransitionTransactionState::AbortedNoRemote | TransitionTransactionState::Committed => {
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
}
|
||||
TransitionTransactionState::UploadOutcomeUnknown => recover_unknown_upload_outcome(api, transaction).await,
|
||||
@@ -915,7 +907,7 @@ async fn recover_unknown_upload_outcome(
|
||||
.map_err(Error::other)?
|
||||
{
|
||||
TransitionCandidateProbe::Missing => {
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RecordDeleted)
|
||||
}
|
||||
TransitionCandidateProbe::UnversionedPresent => {
|
||||
@@ -933,7 +925,7 @@ async fn recover_unknown_upload_outcome(
|
||||
)
|
||||
.await
|
||||
.map_err(Error::other)?;
|
||||
delete_transition_transaction_record(api, transaction).await?;
|
||||
delete_transition_transaction_record(api, transaction.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||
}
|
||||
TransitionCandidateProbe::VersionedPresent(version_id) => {
|
||||
@@ -966,7 +958,7 @@ async fn cleanup_recovered_unknown_upload_candidate(
|
||||
.map_err(transition_transaction_store_error)?;
|
||||
save_transition_transaction_record(api.clone(), &cleanup).await?;
|
||||
delete_transition_remote_candidate(api.clone(), &cleanup).await?;
|
||||
delete_transition_transaction_record(api, &cleanup).await?;
|
||||
delete_transition_transaction_record(api, cleanup.transaction_id).await?;
|
||||
Ok(TransitionTransactionRecoveryOutcome::RemoteCandidateDeleted)
|
||||
}
|
||||
|
||||
|
||||
@@ -406,25 +406,6 @@ where
|
||||
Ok(data)
|
||||
}
|
||||
|
||||
pub(crate) async fn read_config_limited_preserve_empty<S>(api: Arc<S>, file: &str, max_bytes: usize) -> Result<Vec<u8>>
|
||||
where
|
||||
S: EcstoreObjectIO,
|
||||
{
|
||||
let (data, _obj) = read_config_limited_preserve_empty_with_metadata(api, file, max_bytes).await?;
|
||||
Ok(data)
|
||||
}
|
||||
|
||||
pub(crate) async fn read_config_limited_preserve_empty_with_metadata<S>(
|
||||
api: Arc<S>,
|
||||
file: &str,
|
||||
max_bytes: usize,
|
||||
) -> Result<(Vec<u8>, ObjectInfo)>
|
||||
where
|
||||
S: EcstoreObjectIO,
|
||||
{
|
||||
read_config_with_metadata_inner(api, file, &ObjectOptions::default(), true, Some(max_bytes)).await
|
||||
}
|
||||
|
||||
/// Read an existing config object without treating an empty payload as absent.
|
||||
/// Callers that validate their own payload format need to distinguish corruption
|
||||
/// from `ConfigNotFound`.
|
||||
|
||||
+47
-1610
File diff suppressed because it is too large
Load Diff
@@ -37,7 +37,7 @@ use crate::bucket::lifecycle::{
|
||||
transition_transaction::{
|
||||
TransitionRemoteVersion, TransitionSourceIdentity, TransitionSourceVersionMode, TransitionTransaction,
|
||||
TransitionTransactionInit, TransitionTransactionState, delete_transition_transaction_record,
|
||||
load_transition_transaction_record, save_transition_transaction_record,
|
||||
save_transition_transaction_record,
|
||||
},
|
||||
};
|
||||
use crate::bucket::quota::reservation;
|
||||
@@ -4246,12 +4246,7 @@ fn record_transition_uploaded_save_attempt(transaction: &TransitionTransaction,
|
||||
|
||||
async fn delete_transition_transaction_if_available(api: Option<&Arc<ECStore>>, transaction_id: Uuid) -> Result<()> {
|
||||
if let Some(api) = api {
|
||||
let transaction = match load_transition_transaction_record(api.clone(), transaction_id).await {
|
||||
Ok(transaction) => transaction,
|
||||
Err(Error::ConfigNotFound) => return Ok(()),
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
return delete_transition_transaction_record(api.clone(), &transaction).await;
|
||||
return delete_transition_transaction_record(api.clone(), transaction_id).await;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -15,6 +15,20 @@ from pathlib import Path
|
||||
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[1]
|
||||
SCHEDULED_ALERT_WORKFLOWS = (
|
||||
".github/workflows/audit.yml",
|
||||
".github/workflows/build.yml",
|
||||
".github/workflows/ci.yml",
|
||||
".github/workflows/coverage.yml",
|
||||
".github/workflows/e2e-replication-nightly.yml",
|
||||
".github/workflows/e2e-s3tests.yml",
|
||||
".github/workflows/fuzz.yml",
|
||||
".github/workflows/mint.yml",
|
||||
".github/workflows/minio-interop.yml",
|
||||
".github/workflows/nightly-gnu.yml",
|
||||
".github/workflows/performance-ab.yml",
|
||||
".github/workflows/runner-hygiene.yml",
|
||||
)
|
||||
|
||||
|
||||
def words(value: str) -> set[str]:
|
||||
@@ -252,6 +266,72 @@ def check_profile_definitions(root: Path) -> list[str]:
|
||||
return errors
|
||||
|
||||
|
||||
def check_scheduled_alerts(root: Path) -> list[str]:
|
||||
errors: list[str] = []
|
||||
for relative in SCHEDULED_ALERT_WORKFLOWS:
|
||||
path = root / relative
|
||||
try:
|
||||
lines = path.read_text().splitlines()
|
||||
except FileNotFoundError:
|
||||
errors.append(f"{relative}: missing scheduled validation workflow")
|
||||
continue
|
||||
|
||||
try:
|
||||
start = lines.index(" alert-on-failure:") + 1
|
||||
except ValueError:
|
||||
errors.append(f"{relative}: missing alert-on-failure job")
|
||||
continue
|
||||
end = next(
|
||||
(index for index in range(start, len(lines)) if re.fullmatch(r" [A-Za-z0-9_-]+:", lines[index])),
|
||||
len(lines),
|
||||
)
|
||||
job = "\n".join(line.split("#", 1)[0] for line in lines[start:end])
|
||||
required = (
|
||||
"always()",
|
||||
"github.event_name == 'schedule'",
|
||||
"contains(needs.*.result, 'failure')",
|
||||
"issues: write",
|
||||
"uses: ./.github/actions/schedule-failure-issue",
|
||||
"github-token: ${{ secrets.GITHUB_TOKEN }}",
|
||||
)
|
||||
missing = [token for token in required if token not in job]
|
||||
if missing:
|
||||
errors.append(f"{relative}: alert-on-failure missing {', '.join(missing)}")
|
||||
|
||||
watchdog_path = root / ".github/workflows/scheduled-validation-watchdog.yml"
|
||||
try:
|
||||
watchdog = "\n".join(
|
||||
line.split("#", 1)[0] for line in watchdog_path.read_text().splitlines()
|
||||
)
|
||||
except FileNotFoundError:
|
||||
errors.append(".github/workflows/scheduled-validation-watchdog.yml: missing completion watchdog")
|
||||
return errors
|
||||
for relative in SCHEDULED_ALERT_WORKFLOWS:
|
||||
path = root / relative
|
||||
if not path.is_file():
|
||||
continue
|
||||
source = path.read_text()
|
||||
match = re.search(r"^name:\s*[\"']?([^\"'\n]+)", source, re.MULTILINE)
|
||||
if not match:
|
||||
errors.append(f"{relative}: missing workflow name")
|
||||
elif f'- "{match.group(1).strip()}"' not in watchdog:
|
||||
errors.append(f"{relative}: missing from scheduled completion watchdog")
|
||||
required = (
|
||||
"github.event.workflow_run.event == 'schedule'",
|
||||
"github.event.workflow_run.conclusion != 'success'",
|
||||
"github.event.workflow_run.conclusion != 'failure'",
|
||||
"workflow-name: ${{ github.event.workflow_run.name }}",
|
||||
"source-run-id: ${{ github.event.workflow_run.id }}",
|
||||
"source-run-attempt: ${{ github.event.workflow_run.run_attempt }}",
|
||||
)
|
||||
missing = [token for token in required if token not in watchdog]
|
||||
if missing:
|
||||
errors.append(
|
||||
".github/workflows/scheduled-validation-watchdog.yml: missing " + ", ".join(missing)
|
||||
)
|
||||
return errors
|
||||
|
||||
|
||||
def check_profile_listing(root: Path, profile: str, listing: Path) -> list[str]:
|
||||
try:
|
||||
expected_digest = profile_selection(root, profile)
|
||||
@@ -281,6 +361,7 @@ def validate(root: Path) -> list[str]:
|
||||
errors.extend(check_runner_selection(root))
|
||||
errors.extend(check_s3_tests_runner(root))
|
||||
errors.extend(check_profile_definitions(root))
|
||||
errors.extend(check_scheduled_alerts(root))
|
||||
return errors
|
||||
|
||||
|
||||
@@ -363,6 +444,7 @@ class SelfTests(unittest.TestCase):
|
||||
mock.patch(__name__ + ".check_fuzz_targets", return_value=[]),
|
||||
mock.patch(__name__ + ".check_runner_selection", return_value=[]),
|
||||
mock.patch(__name__ + ".check_profile_definitions", return_value=[]),
|
||||
mock.patch(__name__ + ".check_scheduled_alerts", return_value=[]),
|
||||
):
|
||||
self.assertEqual(len(validate(root)), 1)
|
||||
|
||||
@@ -413,6 +495,54 @@ class SelfTests(unittest.TestCase):
|
||||
with mock.patch.object(sys, "platform", "linux"):
|
||||
self.assertEqual(len(check_profile_listing(root, "e2e-full", listing)), 1)
|
||||
|
||||
def test_scheduled_alerts_require_completion_watchdog(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
root = Path(tmp)
|
||||
alert = (
|
||||
" alert-on-failure:\n"
|
||||
" if: always() && github.event_name == 'schedule' && "
|
||||
"contains(needs.*.result, 'failure')\n"
|
||||
" permissions:\n"
|
||||
" issues: write\n"
|
||||
" steps:\n"
|
||||
" - uses: ./.github/actions/schedule-failure-issue\n"
|
||||
" with:\n"
|
||||
" github-token: ${{ secrets.GITHUB_TOKEN }}\n"
|
||||
)
|
||||
names: list[str] = []
|
||||
for relative in SCHEDULED_ALERT_WORKFLOWS:
|
||||
path = root / relative
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
names.append(path.stem)
|
||||
path.write_text(f'name: "{path.stem}"\n{alert}')
|
||||
watchdog = root / ".github/workflows/scheduled-validation-watchdog.yml"
|
||||
watchdog.write_text(
|
||||
"\n".join(f'- "{name}"' for name in names)
|
||||
+ "\ngithub.event.workflow_run.event == 'schedule'\n"
|
||||
+ "github.event.workflow_run.conclusion != 'success'\n"
|
||||
+ "github.event.workflow_run.conclusion != 'failure'\n"
|
||||
+ "workflow-name: ${{ github.event.workflow_run.name }}\n"
|
||||
+ "source-run-id: ${{ github.event.workflow_run.id }}\n"
|
||||
+ "source-run-attempt: ${{ github.event.workflow_run.run_attempt }}\n"
|
||||
)
|
||||
self.assertEqual(check_scheduled_alerts(root), [])
|
||||
|
||||
first = root / SCHEDULED_ALERT_WORKFLOWS[0]
|
||||
mutations = (
|
||||
("contains(needs.*.result, 'failure')", "false"),
|
||||
("issues: write", "issues: read"),
|
||||
("uses: ./.github/actions/schedule-failure-issue", "uses: actions/checkout@v7"),
|
||||
("github-token: ${{ secrets.GITHUB_TOKEN }}", "github-token: missing"),
|
||||
)
|
||||
for required, replacement in mutations:
|
||||
original = first.read_text()
|
||||
first.write_text(original.replace(required, replacement))
|
||||
self.assertEqual(len(check_scheduled_alerts(root)), 1)
|
||||
first.write_text(original)
|
||||
|
||||
watchdog.write_text(watchdog.read_text().replace(f'- "{names[0]}"\n', ""))
|
||||
self.assertEqual(len(check_scheduled_alerts(root)), 1)
|
||||
|
||||
def main() -> int:
|
||||
if sys.argv[1:] == ["--self-test"]:
|
||||
suite = unittest.defaultTestLoader.loadTestsFromTestCase(SelfTests)
|
||||
@@ -436,7 +566,7 @@ def main() -> int:
|
||||
for error in errors:
|
||||
print(f"ERROR: {error}", file=sys.stderr)
|
||||
return 1
|
||||
print("OK: e2e modules, runner selection, fuzz matrices, profiles, and bounded diagnostics are wired")
|
||||
print("OK: e2e modules, runner selection, fuzz matrices, profiles, and scheduled alerts are wired")
|
||||
return 0
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user