fix(notify): persist Docker queue stores and expose open errors (#8201)

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
Hauser
2026-09-28 21:24:12 +08:00
committed by GitHub
co-authored by zhi22915
parent 4a91c75152
commit 3bc36acf7f
8 changed files with 121 additions and 35 deletions
+3 -1
View File
@@ -106,7 +106,9 @@ RUN chmod +x /usr/bin/rustfs /entrypoint.sh
RUN addgroup -g 10001 -S rustfs && \
adduser -u 10001 -G rustfs -S rustfs -D && \
mkdir -p /data /logs && \
mkdir -p /data/.rustfs/events /data/.rustfs/audit /logs /opt/rustfs && \
ln -s /data/.rustfs/events /opt/rustfs/events && \
ln -s /data/.rustfs/audit /opt/rustfs/audit && \
chown -R rustfs:rustfs /data /logs && \
chmod 0750 /data /logs
+3 -1
View File
@@ -44,7 +44,9 @@ RUN set -eux; \
WORKDIR /app
RUN set -eux; \
mkdir -p /data /logs; \
mkdir -p /data/.rustfs/events /data/.rustfs/audit /logs /opt/rustfs; \
ln -s /data/.rustfs/events /opt/rustfs/events; \
ln -s /data/.rustfs/audit /opt/rustfs/audit; \
chown -R rustfs:rustfs /data /logs /app; \
chmod 0750 /data /logs
+3 -1
View File
@@ -110,7 +110,9 @@ RUN chmod +x /usr/bin/rustfs /entrypoint.sh
RUN groupadd -g 10001 rustfs && \
useradd -u 10001 -g rustfs -m -s /sbin/nologin rustfs && \
mkdir -p /data /logs && \
mkdir -p /data/.rustfs/events /data/.rustfs/audit /logs /opt/rustfs && \
ln -s /data/.rustfs/events /opt/rustfs/events && \
ln -s /data/.rustfs/audit /opt/rustfs/audit && \
chown -R rustfs:rustfs /data /logs && \
chmod 0750 /data /logs
+3 -1
View File
@@ -240,7 +240,9 @@ WORKDIR /app
# Prepare data/log directories with sane defaults
RUN set -eux; \
mkdir -p /data /logs; \
mkdir -p /data/.rustfs/events /data/.rustfs/audit /logs /opt/rustfs; \
ln -s /data/.rustfs/events /opt/rustfs/events; \
ln -s /data/.rustfs/audit /opt/rustfs/audit; \
chown -R rustfs:rustfs /data /logs /app; \
chmod 0750 /data /logs
+1 -1
View File
@@ -195,7 +195,7 @@ docker run -d --name rustfs -p 9000:9000 \
Notes:
- `RUSTFS_NOTIFY_ENABLE=true` enables the global notify module switch.
- For ARN `arn:rustfs:sqs::primary:webhook`, use instance-scoped env vars with `_PRIMARY`.
- If queue dir is omitted, default is `/opt/rustfs/events`; ensure it is writable by the container runtime user.
- If queue dir is omitted, the official image maps `/opt/rustfs/events` into the persistent `/data` volume so it is writable by the runtime user and pending events survive container recreation. Override `queue_dir` when using another deployment layout.
- `RUSTFS_NOTIFY_WEBHOOK_SKIP_TLS_VERIFY_PRIMARY` defaults to `false`; enabling it skips webhook TLS certificate verification, allows MITM attacks, and emits a startup warning. Prefer `RUSTFS_NOTIFY_WEBHOOK_CLIENT_CA_PRIMARY` for private CAs.
- Since `1.0.0-beta.11`, webhook endpoints on private or container networks
(`Docker Compose service names`, `host.docker.internal`, RFC 1918 addresses) are
+50 -24
View File
@@ -41,12 +41,20 @@ static STORE_OPEN_POOL: LazyLock<Result<rayon::ThreadPool, rayon::ThreadPoolBuil
.build()
});
enum StoreOpenFailure {
Error(String),
Panicked,
}
enum StoreOpenOutcome<E>
where
E: PluginEvent,
{
Accepted(SharedTarget<E>),
Rejected { panicked: bool, target: SharedTarget<E> },
Rejected {
failure: StoreOpenFailure,
target: SharedTarget<E>,
},
}
fn open_target_store<E>(target: SharedTarget<E>) -> StoreOpenOutcome<E>
@@ -55,8 +63,14 @@ where
{
match catch_unwind(AssertUnwindSafe(|| target.store().map(|store| store.open()))) {
Ok(None | Some(Ok(()))) => StoreOpenOutcome::Accepted(target),
Ok(Some(Err(_))) => StoreOpenOutcome::Rejected { panicked: false, target },
Err(_) => StoreOpenOutcome::Rejected { panicked: true, target },
Ok(Some(Err(error))) => StoreOpenOutcome::Rejected {
failure: StoreOpenFailure::Error(error.to_string()),
target,
},
Err(_) => StoreOpenOutcome::Rejected {
failure: StoreOpenFailure::Panicked,
target,
},
}
}
@@ -215,23 +229,27 @@ where
for outcome in outcomes {
match outcome {
StoreOpenOutcome::Accepted(target) => accepted.push(target),
StoreOpenOutcome::Rejected { panicked, target } => {
if panicked {
tracing::error!(
target_id = %target.id(),
reason = "store_open_panicked",
"Target queue store panicked while opening during runtime handoff"
);
} else {
tracing::error!(
target_id = %target.id(),
reason = "store_open_failed",
"Failed to open target queue store during runtime handoff"
);
}
failures.push(TargetActivationFailure {
detail: format!("{}: queue store open failed", target.id()),
});
StoreOpenOutcome::Rejected { failure, target } => {
let detail = match failure {
StoreOpenFailure::Error(error) => {
tracing::error!(
target_id = %target.id(),
error = %error,
reason = "store_open_failed",
"Failed to open target queue store during runtime handoff"
);
format!("{}: queue store open failed: {error}", target.id())
}
StoreOpenFailure::Panicked => {
tracing::error!(
target_id = %target.id(),
reason = "store_open_panicked",
"Target queue store panicked while opening during runtime handoff"
);
format!("{}: queue store open failed", target.id())
}
};
failures.push(TargetActivationFailure { detail });
rejected.push(target);
}
}
@@ -717,10 +735,18 @@ mod tests {
)));
let observer = target.clone();
let activation = adapter.activate_with_replay(vec![Box::new(target)]).await;
assert!(activation.targets.is_empty());
assert!(activation.replay_workers.is_empty());
let prepared = adapter.prepare_targets(vec![Box::new(target)]).await;
let (opened, rejected) = adapter.open_prepared_stores(prepared);
assert!(opened.targets.is_empty());
let failure = rejected
.failure_summary()
.expect("a rejected target should retain the store error");
assert!(failure.contains(invalid_base.to_string_lossy().as_ref()));
assert!(failure.contains("queue store open failed:"));
adapter
.close_prepared(rejected)
.await
.expect("a rejected target should close cleanly");
assert_eq!(observer.close_call_count(), 1);
}
+53 -5
View File
@@ -85,6 +85,34 @@ const FAILED_STORE_MAX_ENTRIES: usize = 10_000;
/// time, which is the instant the entry entered the failed store.
const FAILED_STORE_TTL: Duration = Duration::from_secs(72 * 60 * 60);
#[derive(Debug)]
struct IoErrorWithPath {
path: PathBuf,
source: std::io::Error,
}
impl std::fmt::Display for IoErrorWithPath {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(formatter, "{}: {}", self.path.display(), self.source)
}
}
impl std::error::Error for IoErrorWithPath {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
Some(&self.source)
}
}
fn io_error_with_path(path: &Path, error: std::io::Error) -> StoreError {
StoreError::Io(std::io::Error::new(
error.kind(),
IoErrorWithPath {
path: path.to_owned(),
source: error,
},
))
}
/// Writes payload to a temp file in the same directory, flushes the file to disk, then atomically
/// renames it onto final_path.
///
@@ -741,9 +769,9 @@ where
.fs_guard
.write()
.map_err(|_| StoreError::Internal("Failed to acquire write lock on store filesystem".to_string()))?;
std::fs::create_dir_all(&self.directory).map_err(StoreError::Io)?;
std::fs::create_dir_all(&self.directory).map_err(|error| io_error_with_path(&self.directory, error))?;
let dir_entries = std::fs::read_dir(&self.directory).map_err(StoreError::Io)?;
let dir_entries = std::fs::read_dir(&self.directory).map_err(|error| io_error_with_path(&self.directory, error))?;
let mut entries_map = self
.entries
.write()
@@ -751,8 +779,9 @@ where
self.pending_entries.store(0, Ordering::SeqCst);
entries_map.clear();
for entry in dir_entries {
let entry = entry.map_err(StoreError::Io)?;
let metadata = entry.metadata().map_err(StoreError::Io)?;
let entry = entry.map_err(|error| io_error_with_path(&self.directory, error))?;
let entry_path = entry.path();
let metadata = entry.metadata().map_err(|error| io_error_with_path(&entry_path, error))?;
if !metadata.is_file() {
continue;
}
@@ -780,7 +809,7 @@ where
continue;
}
let modified = metadata.modified().map_err(StoreError::Io)?;
let modified = metadata.modified().map_err(|error| io_error_with_path(&entry_path, error))?;
let unix_nano = modified.duration_since(UNIX_EPOCH).unwrap_or_default().as_nanos() as i64;
entries_map.insert(file_name, unix_nano);
}
@@ -1230,6 +1259,25 @@ mod tests {
std::env::temp_dir().join(format!("rustfs-targets-{name}-{}", Uuid::new_v4()))
}
#[test]
fn io_error_with_path_preserves_path_and_source_chain() {
let error = io_error_with_path(
Path::new("/data/.rustfs/events/webhook-primary"),
std::io::Error::from(std::io::ErrorKind::PermissionDenied),
);
assert!(error.to_string().contains("/data/.rustfs/events/webhook-primary"));
let wrapped = std::error::Error::source(&error).expect("store error should expose its I/O error");
let original = std::error::Error::source(wrapped).expect("path wrapper should expose the original I/O error");
assert_eq!(
original
.downcast_ref::<std::io::Error>()
.expect("source should remain an I/O error")
.kind(),
std::io::ErrorKind::PermissionDenied
);
}
#[test]
fn resolve_queue_store_compression_defaults_to_true() {
assert!(resolve_queue_store_compression_from_env_value(None));
+5 -1
View File
@@ -1264,7 +1264,11 @@ mod tests {
match result {
Ok(_) => panic!("expected open_target_queue_store to fail on file base path"),
Err(err) => assert!(err.to_string().contains("custom open context")),
Err(err) => {
let message = err.to_string();
assert!(message.contains("custom open context"));
assert!(message.contains(base.to_string_lossy().as_ref()));
}
}
let _ = fs::remove_file(base);
}