-
Notifications
You must be signed in to change notification settings - Fork 68
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
[redis-rs][core] Move connection refresh to the background #2915
base: main
Are you sure you want to change the base?
Changes from 3 commits
bf2cbe2
f850d99
30b35f2
b2caf01
4e6535f
2d93b4a
d7ff41e
bf867cc
a12837a
7b84fb0
be41dc5
e87e0cb
029aafc
66a3e39
7edb54f
2e569a5
78a966f
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -2,14 +2,16 @@ use crate::cluster_async::ConnectionFuture; | |
use crate::cluster_routing::{Route, ShardAddrs, SlotAddr}; | ||
use crate::cluster_slotmap::{ReadFromReplicaStrategy, SlotMap, SlotMapValue}; | ||
use crate::cluster_topology::TopologyHash; | ||
use dashmap::DashMap; | ||
use dashmap::{DashMap, DashSet}; | ||
use futures::FutureExt; | ||
use rand::seq::IteratorRandom; | ||
use std::net::IpAddr; | ||
use std::sync::atomic::Ordering; | ||
use std::sync::Arc; | ||
use telemetrylib::Telemetry; | ||
|
||
use tokio::task::JoinHandle; | ||
|
||
/// Count the number of connections in a connections_map object | ||
macro_rules! count_connections { | ||
($conn_map:expr) => {{ | ||
|
@@ -121,6 +123,11 @@ pub(crate) enum ConnectionType { | |
|
||
pub(crate) struct ConnectionsMap<Connection>(pub(crate) DashMap<String, ClusterNode<Connection>>); | ||
|
||
pub(crate) struct RefreshState<Connection> { | ||
pub handle: JoinHandle<()>, // The currect running refresh task | ||
pub node_conn: Option<ClusterNode<Connection>>, // The refreshed connection after the task is done | ||
} | ||
|
||
impl<Connection> std::fmt::Display for ConnectionsMap<Connection> { | ||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { | ||
for item in self.0.iter() { | ||
|
@@ -139,6 +146,13 @@ pub(crate) struct ConnectionsContainer<Connection> { | |
pub(crate) slot_map: SlotMap, | ||
read_from_replica_strategy: ReadFromReplicaStrategy, | ||
topology_hash: TopologyHash, | ||
|
||
// Holds all the failed addresses that started a refresh task. | ||
pub(crate) refresh_addresses_started: DashSet<String>, | ||
// Follow the refresh ops on the connections | ||
pub(crate) refresh_operations: DashMap<String, RefreshState<Connection>>, | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. as we talked - instead of using RefreshState, use a ConnectionState that internally holds the operation / connection (same as is reconnectingConnection). For example -
or
the refresh_connections can be called only for user/management connection or for both, so you should make sure this solution covers all cases |
||
// Holds all the refreshed addresses that are ready to be inserted into the connection_map | ||
pub(crate) refresh_addresses_done: DashSet<String>, | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I don't see a good reason for using DashSet or DashMap, these structs are only going to be used from a single point without concurrency |
||
} | ||
|
||
impl<Connection> Drop for ConnectionsContainer<Connection> { | ||
|
@@ -155,6 +169,9 @@ impl<Connection> Default for ConnectionsContainer<Connection> { | |
slot_map: Default::default(), | ||
read_from_replica_strategy: ReadFromReplicaStrategy::AlwaysFromPrimary, | ||
topology_hash: 0, | ||
refresh_addresses_started: DashSet::new(), | ||
refresh_operations: DashMap::new(), | ||
refresh_addresses_done: DashSet::new(), | ||
} | ||
} | ||
} | ||
|
@@ -182,6 +199,9 @@ where | |
slot_map, | ||
read_from_replica_strategy, | ||
topology_hash, | ||
refresh_addresses_started: DashSet::new(), | ||
refresh_operations: DashMap::new(), | ||
refresh_addresses_done: DashSet::new(), | ||
} | ||
} | ||
|
||
|
@@ -572,6 +592,9 @@ mod tests { | |
connection_map, | ||
read_from_replica_strategy: ReadFromReplicaStrategy::AZAffinity("use-1a".to_string()), | ||
topology_hash: 0, | ||
refresh_addresses_started: DashSet::new(), | ||
refresh_operations: DashMap::new(), | ||
refresh_addresses_done: DashSet::new(), | ||
} | ||
} | ||
|
||
|
@@ -628,6 +651,9 @@ mod tests { | |
connection_map, | ||
read_from_replica_strategy: strategy, | ||
topology_hash: 0, | ||
refresh_addresses_started: DashSet::new(), | ||
refresh_operations: DashMap::new(), | ||
refresh_addresses_done: DashSet::new(), | ||
} | ||
} | ||
|
||
|
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
For example: