mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(svc): authenticate to nats with the configured token (#2581)
This commit is contained in:
@@ -17,6 +17,7 @@ pub struct ServiceConfig {
|
||||
pub shard_count: u32,
|
||||
pub listen_addr: SocketAddr,
|
||||
pub nats_url: String,
|
||||
pub nats_auth_token: Option<String>,
|
||||
pub cache_max_entries: u64,
|
||||
pub cache_ttl: Duration,
|
||||
pub cache_hard_ttl: Duration,
|
||||
@@ -105,6 +106,8 @@ impl ServiceConfig {
|
||||
let nats_url = optional_from(&get, "FLUXER_SVC_NATS_URL")
|
||||
.unwrap_or_else(|| "nats://127.0.0.1:4222".to_owned());
|
||||
|
||||
let nats_auth_token = optional_from(&get, "FLUXER_NATS_AUTH_TOKEN");
|
||||
|
||||
let cache_ttl_ms = optional_from(&get, "FLUXER_SVC_CACHE_TTL_MS")
|
||||
.map(|v| v.parse::<u64>())
|
||||
.transpose()?
|
||||
@@ -165,6 +168,7 @@ impl ServiceConfig {
|
||||
shard_count,
|
||||
listen_addr: format!("{listen_host}:{listen_port}").parse()?,
|
||||
nats_url,
|
||||
nats_auth_token,
|
||||
cache_max_entries: optional_from(&get, "FLUXER_SVC_CACHE_MAX_ENTRIES")
|
||||
.map(|v| v.parse::<u64>())
|
||||
.transpose()?
|
||||
|
||||
@@ -33,7 +33,9 @@ where
|
||||
{
|
||||
init_tracing();
|
||||
let config = config::ServiceConfig::from_env()?;
|
||||
let transport = transport::NatsTransport::connect(&config.nats_url).await?;
|
||||
let transport =
|
||||
transport::NatsTransport::connect(&config.nats_url, config.nats_auth_token.as_deref())
|
||||
.await?;
|
||||
tracing::info!(
|
||||
service = config.service_name,
|
||||
mode = ?config.mode,
|
||||
|
||||
@@ -868,6 +868,7 @@ mod tests {
|
||||
shard_count: 1,
|
||||
listen_addr: "127.0.0.1:0".parse().unwrap(),
|
||||
nats_url: "memory".to_owned(),
|
||||
nats_auth_token: None,
|
||||
cache_max_entries: 100,
|
||||
cache_ttl: Duration::from_secs(30),
|
||||
cache_hard_ttl: Duration::from_secs(600),
|
||||
|
||||
@@ -417,6 +417,7 @@ mod tests {
|
||||
shard_count: 1,
|
||||
listen_addr: "127.0.0.1:0".parse().unwrap(),
|
||||
nats_url: "memory".to_owned(),
|
||||
nats_auth_token: None,
|
||||
cache_max_entries: 100,
|
||||
cache_ttl: Duration::from_secs(30),
|
||||
cache_hard_ttl: Duration::from_secs(600),
|
||||
|
||||
@@ -77,7 +77,7 @@ pub struct NatsMessage {
|
||||
}
|
||||
|
||||
impl NatsTransport {
|
||||
pub async fn connect(url: &str) -> anyhow::Result<Self> {
|
||||
pub async fn connect(url: &str, auth_token: Option<&str>) -> anyhow::Result<Self> {
|
||||
let reconnect_notify = Arc::new(Notify::new());
|
||||
|
||||
let event_notify = reconnect_notify.clone();
|
||||
@@ -112,6 +112,11 @@ impl NatsTransport {
|
||||
}
|
||||
});
|
||||
|
||||
let options = match auth_token {
|
||||
Some(token) => options.token(token.to_owned()),
|
||||
None => options,
|
||||
};
|
||||
|
||||
let client = options
|
||||
.subscription_capacity(NATS_SUBSCRIPTION_CAPACITY)
|
||||
.connect(url)
|
||||
|
||||
Reference in New Issue
Block a user