use std::net::{Shutdown, TcpStream};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::thread::JoinHandle;
use std::time::Duration;
use kevy_replicate::replica::ReplicaClient;
use kevy_rt::ReplicaInboxSender;
use crate::replica_runner_events::drain_client;
use crate::state::ReplicaProgress;
const RECONNECT_BACKOFF: Duration = Duration::from_millis(250);
pub(crate) struct ReplicaRunner {
handle: Option<JoinHandle<()>>,
stop: Arc<AtomicBool>,
socket: Arc<Mutex<Option<TcpStream>>>,
}
impl ReplicaRunner {
pub(crate) fn spawn(
upstream_addr: (std::net::IpAddr, u16),
replica_id: String,
sender: ReplicaInboxSender,
runner_slot: usize,
progress: Arc<ReplicaProgress>,
) -> Self {
Self::spawn_target(
upstream_addr,
replica_id,
Target::PerShard(sender),
runner_slot,
progress,
)
}
pub(crate) fn spawn_routed(
upstream_addr: (std::net::IpAddr, u16),
replica_id: String,
senders: Vec<ReplicaInboxSender>,
runner_slot: usize,
progress: Arc<ReplicaProgress>,
) -> Self {
Self::spawn_target(
upstream_addr,
replica_id,
Target::Routed(senders),
runner_slot,
progress,
)
}
fn spawn_target(
upstream_addr: (std::net::IpAddr, u16),
replica_id: String,
target: Target,
runner_slot: usize,
progress: Arc<ReplicaProgress>,
) -> Self {
let stop = Arc::new(AtomicBool::new(false));
let stop_thread = stop.clone();
let socket: Arc<Mutex<Option<TcpStream>>> = Arc::new(Mutex::new(None));
let socket_thread = socket.clone();
let handle = std::thread::Builder::new()
.name(format!("kevy-replica-{replica_id}"))
.spawn(move || {
run_loop(
upstream_addr,
replica_id,
target,
stop_thread,
socket_thread,
runner_slot,
progress,
);
})
.expect("spawn replica runner thread");
Self { handle: Some(handle), stop, socket }
}
#[allow(dead_code)] pub(crate) fn shutdown(mut self) {
self.signal_stop();
if let Some(h) = self.handle.take() {
let _ = h.join();
}
}
fn signal_stop(&self) {
self.stop.store(true, Ordering::Relaxed);
if let Ok(guard) = self.socket.lock()
&& let Some(s) = guard.as_ref()
{
let _ = s.shutdown(Shutdown::Both);
}
}
}
impl Drop for ReplicaRunner {
fn drop(&mut self) {
self.signal_stop();
if let Some(h) = self.handle.take() {
let _ = h.join();
}
}
}
enum Target {
PerShard(ReplicaInboxSender),
Routed(Vec<ReplicaInboxSender>),
}
#[allow(clippy::too_many_arguments)]
fn one_session(
upstream_addr: (std::net::IpAddr, u16),
replica_id: &str,
target: &Target,
stop: &Arc<AtomicBool>,
socket_slot: &Arc<Mutex<Option<TcpStream>>>,
runner_slot: usize,
progress: &Arc<ReplicaProgress>,
from_offset: u64,
data_gen: &mut u64,
) -> u64 {
match ReplicaClient::connect_at(
upstream_addr,
replica_id,
*data_gen,
from_offset,
Duration::from_secs(5),
) {
Ok(mut client) => {
drain_session(&mut client, target, stop, socket_slot, runner_slot, progress, data_gen)
}
Err(e) => {
eprintln!(
"kevy: replica runner '{replica_id}' connect to \
{upstream_addr:?} failed: {e}; retrying in \
{RECONNECT_BACKOFF:?}"
);
from_offset
}
}
}
fn run_loop(
upstream_addr: (std::net::IpAddr, u16),
replica_id: String,
target: Target,
stop: Arc<AtomicBool>,
socket_slot: Arc<Mutex<Option<TcpStream>>>,
runner_slot: usize,
progress: Arc<ReplicaProgress>,
) {
let mut from_offset: u64 = 0;
let mut data_gen: u64 = 0;
while !stop.load(Ordering::Relaxed) {
from_offset = one_session(
upstream_addr,
&replica_id,
&target,
&stop,
&socket_slot,
runner_slot,
&progress,
from_offset,
&mut data_gen,
);
if !stop.load(Ordering::Relaxed) {
std::thread::sleep(RECONNECT_BACKOFF);
}
}
}
fn drain_session(
client: &mut ReplicaClient,
target: &Target,
stop: &Arc<AtomicBool>,
socket_slot: &Mutex<Option<TcpStream>>,
runner_slot: usize,
progress: &Arc<ReplicaProgress>,
data_gen: &mut u64,
) -> u64 {
set_socket_slot(socket_slot, client.socket_handle().ok());
crate::replica_trace::trace_session_start(runner_slot, client, *data_gen);
let from_offset = match target {
Target::PerShard(sender) => {
drain_client(client, sender, stop, runner_slot, progress, data_gen)
}
Target::Routed(senders) => crate::replica_runner_routed::drain_client_routed(
client,
senders,
stop,
runner_slot,
progress,
data_gen,
),
};
set_socket_slot(socket_slot, None);
from_offset
}
fn set_socket_slot(slot: &Mutex<Option<TcpStream>>, value: Option<TcpStream>) {
if let Ok(mut guard) = slot.lock() {
*guard = value;
}
}