Files
rustfs/crates/config/src/constants/targets.rs
T
escapecode a80699b6dd feat: add an opt-in NATS JetStream publish path for the notify and audit targets (#4634)
feat(targets): add an opt-in NATS JetStream publish path for the notify and audit targets

The NATS notify and audit targets publish through NATS Core, which returns
before the server has durably accepted the message. A broker restart or a
connection drop between the publish and the flush loses the event, even though
the send queue has already cleared it, and no acknowledgement gates that clear.

An opt-in JetStream publish path clears a queued event only after the server
returns a durable PublishAck, so delivery is at-least-once across a broker
restart or a reconnect. It applies to both the notify and audit NATS targets, is
off by default, and is byte-identical to the NATS Core path when disabled.

The path includes durable store-and-forward, a stable dedup id sent as the
Nats-Msg-Id header so a replayed event is collapsed by the stream duplicate
window, pre-flight stream validation, and a bounded failed-events store for
terminally-failed and retry-exhausted events. Three configuration keys per
target select it: JETSTREAM_ENABLE, JETSTREAM_STREAM_NAME, and
JETSTREAM_ACK_TIMEOUT_SECS, under the RUSTFS_NOTIFY_NATS_ and RUSTFS_AUDIT_NATS_
prefixes.

The on-disk batch filename separator changes from colon to underscore so
batch names are valid on Windows filesystems, with transparent read-back
of files written under the previous separator. The migration affects the
shared queue store for every target type and lands with this feature
because the store gains its first Windows-exercised paths here.

Co-authored-by: houseme <housemecn@gmail.com>
2026-07-14 15:36:14 +08:00

149 lines
6.9 KiB
Rust

// 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.
pub const WEBHOOK_ENDPOINT: &str = "endpoint";
pub const WEBHOOK_AUTH_TOKEN: &str = "auth_token";
pub const WEBHOOK_CLIENT_CERT: &str = "client_cert";
pub const WEBHOOK_CLIENT_KEY: &str = "client_key";
pub const WEBHOOK_CLIENT_CA: &str = "client_ca";
pub const WEBHOOK_SKIP_TLS_VERIFY: &str = "skip_tls_verify";
pub const WEBHOOK_BATCH_SIZE: &str = "batch_size";
pub const WEBHOOK_QUEUE_LIMIT: &str = "queue_limit";
pub const WEBHOOK_QUEUE_DIR: &str = "queue_dir";
pub const WEBHOOK_MAX_RETRY: &str = "max_retry";
pub const WEBHOOK_RETRY_INTERVAL: &str = "retry_interval";
pub const WEBHOOK_HTTP_TIMEOUT: &str = "http_timeout";
pub const MQTT_BROKER: &str = "broker";
pub const MQTT_TOPIC: &str = "topic";
pub const MQTT_QOS: &str = "qos";
pub const MQTT_USERNAME: &str = "username";
pub const MQTT_PASSWORD: &str = "password";
pub const MQTT_RECONNECT_INTERVAL: &str = "reconnect_interval";
pub const MQTT_KEEP_ALIVE_INTERVAL: &str = "keep_alive_interval";
pub const MQTT_QUEUE_DIR: &str = "queue_dir";
pub const MQTT_QUEUE_LIMIT: &str = "queue_limit";
pub const MQTT_TLS_POLICY: &str = "tls_policy";
pub const MQTT_TLS_CA: &str = "tls_ca";
pub const MQTT_TLS_CLIENT_CERT: &str = "tls_client_cert";
pub const MQTT_TLS_CLIENT_KEY: &str = "tls_client_key";
pub const MQTT_TLS_TRUST_LEAF_AS_CA: &str = "tls_trust_leaf_as_ca";
pub const MQTT_WS_PATH_ALLOWLIST: &str = "ws_path_allowlist";
pub const KAFKA_BROKERS: &str = "brokers";
pub const KAFKA_TOPIC: &str = "topic";
pub const KAFKA_ACKS: &str = "acks";
pub const KAFKA_QUEUE_DIR: &str = "queue_dir";
pub const KAFKA_QUEUE_LIMIT: &str = "queue_limit";
pub const KAFKA_TLS_ENABLE: &str = "tls_enable";
pub const KAFKA_TLS_CA: &str = "tls_ca";
pub const KAFKA_TLS_CLIENT_CERT: &str = "tls_client_cert";
pub const KAFKA_TLS_CLIENT_KEY: &str = "tls_client_key";
pub const KAFKA_SASL_ENABLE: &str = "sasl_enable";
pub const KAFKA_SASL_MECHANISM: &str = "sasl_mechanism";
pub const KAFKA_SASL_USERNAME: &str = "sasl_username";
pub const KAFKA_SASL_PASSWORD: &str = "sasl_password";
pub const AMQP_URL: &str = "url";
pub const AMQP_EXCHANGE: &str = "exchange";
pub const AMQP_ROUTING_KEY: &str = "routing_key";
pub const AMQP_MANDATORY: &str = "mandatory";
pub const AMQP_PERSISTENT: &str = "persistent";
pub const AMQP_USERNAME: &str = "username";
pub const AMQP_PASSWORD: &str = "password";
pub const AMQP_TLS_CA: &str = "tls_ca";
pub const AMQP_TLS_CLIENT_CERT: &str = "tls_client_cert";
pub const AMQP_TLS_CLIENT_KEY: &str = "tls_client_key";
pub const AMQP_QUEUE_DIR: &str = "queue_dir";
pub const AMQP_QUEUE_LIMIT: &str = "queue_limit";
pub const NATS_ADDRESS: &str = "address";
pub const NATS_SUBJECT: &str = "subject";
pub const NATS_USERNAME: &str = "username";
pub const NATS_PASSWORD: &str = "password";
pub const NATS_TOKEN: &str = "token";
pub const NATS_CREDENTIALS_FILE: &str = "credentials_file";
pub const NATS_TLS_CA: &str = "tls_ca";
pub const NATS_TLS_CLIENT_CERT: &str = "tls_client_cert";
pub const NATS_TLS_CLIENT_KEY: &str = "tls_client_key";
pub const NATS_TLS_REQUIRED: &str = "tls_required";
pub const NATS_QUEUE_DIR: &str = "queue_dir";
pub const NATS_QUEUE_LIMIT: &str = "queue_limit";
pub const NATS_JETSTREAM_ENABLE: &str = "jetstream_enable";
pub const NATS_JETSTREAM_STREAM_NAME: &str = "jetstream_stream_name";
pub const NATS_JETSTREAM_ACK_TIMEOUT_SECS: &str = "jetstream_ack_timeout_secs";
pub const NATS_JETSTREAM_ACK_TIMEOUT_DEFAULT_SECS: u64 = 30;
pub const NATS_JETSTREAM_ACK_TIMEOUT_MIN_SECS: u64 = 10;
pub const NATS_JETSTREAM_ACK_TIMEOUT_MAX_SECS: u64 = 120;
pub const PULSAR_BROKER: &str = "broker";
pub const PULSAR_TOPIC: &str = "topic";
pub const PULSAR_AUTH_TOKEN: &str = "auth_token";
pub const PULSAR_USERNAME: &str = "username";
pub const PULSAR_PASSWORD: &str = "password";
pub const PULSAR_TLS_CA: &str = "tls_ca";
pub const PULSAR_TLS_ALLOW_INSECURE: &str = "tls_allow_insecure";
pub const PULSAR_TLS_HOSTNAME_VERIFICATION: &str = "tls_hostname_verification";
pub const PULSAR_QUEUE_DIR: &str = "queue_dir";
pub const PULSAR_QUEUE_LIMIT: &str = "queue_limit";
pub const BASE_DSN_STRING: &str = "dsn_string";
pub const MYSQL_DSN_STRING: &str = BASE_DSN_STRING;
pub const MYSQL_TABLE: &str = "table";
pub const MYSQL_FORMAT: &str = "format";
pub const MYSQL_TLS_CA: &str = "tls_ca";
pub const MYSQL_TLS_CLIENT_CERT: &str = "tls_client_cert";
pub const MYSQL_TLS_CLIENT_KEY: &str = "tls_client_key";
pub const MYSQL_QUEUE_DIR: &str = "queue_dir";
pub const MYSQL_QUEUE_LIMIT: &str = "queue_limit";
pub const MYSQL_MAX_OPEN_CONNECTIONS: &str = "max_open_connections";
pub const REDIS_URL: &str = "url";
pub const REDIS_CHANNEL: &str = "channel";
pub const REDIS_USERNAME: &str = "username";
pub const REDIS_PASSWORD: &str = "password";
pub const REDIS_KEEP_ALIVE_INTERVAL: &str = "keep_alive_interval";
pub const REDIS_QUEUE_DIR: &str = "queue_dir";
pub const REDIS_QUEUE_LIMIT: &str = "queue_limit";
pub const REDIS_MAX_RETRY_ATTEMPTS: &str = "max_retry_attempts";
pub const REDIS_RECONNECT_RETRY_ATTEMPTS: &str = "reconnect_retry_attempts";
pub const REDIS_MIN_RETRY_DELAY: &str = "min_retry_delay";
pub const REDIS_MAX_RETRY_DELAY: &str = "max_retry_delay";
pub const REDIS_CONNECTION_TIMEOUT: &str = "connection_timeout";
pub const REDIS_RESPONSE_TIMEOUT: &str = "response_timeout";
pub const REDIS_PIPELINE_BUFFER_SIZE: &str = "pipeline_buffer_size";
pub const REDIS_TLS_POLICY: &str = "tls_policy";
pub const REDIS_TLS_CA: &str = "tls_ca";
pub const REDIS_TLS_CLIENT_CERT: &str = "tls_client_cert";
pub const REDIS_TLS_CLIENT_KEY: &str = "tls_client_key";
pub const REDIS_TLS_ALLOW_INSECURE: &str = "tls_allow_insecure";
pub const POSTGRES_DSN_STRING: &str = BASE_DSN_STRING;
pub const POSTGRES_TABLE: &str = "table";
pub const POSTGRES_FORMAT: &str = "format";
pub const POSTGRES_TLS_REQUIRED: &str = "tls_required";
pub const POSTGRES_TLS_CA: &str = "tls_ca";
pub const POSTGRES_TLS_CLIENT_CERT: &str = "tls_client_cert";
pub const POSTGRES_TLS_CLIENT_KEY: &str = "tls_client_key";
pub const POSTGRES_QUEUE_DIR: &str = "queue_dir";
pub const POSTGRES_QUEUE_LIMIT: &str = "queue_limit";
/// Environment variable controlling whether target queue files are Snappy-compressed.
/// Applies to both notify and audit target queue stores.
pub const ENV_TARGET_STORE_COMPRESS: &str = "RUSTFS_TARGET_STORE_COMPRESS";
/// Queue-store compression is enabled by default to reduce disk footprint.
pub const DEFAULT_TARGET_STORE_COMPRESS: bool = true;