Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
91 changes: 49 additions & 42 deletions miner-apps/jd-client/src/lib/channel_manager/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -566,38 +566,11 @@ impl ChannelManager {
// todo: let start downstream accept channel manager as `Arc`, instead of clone
let this = Arc::new(self);

// Wait for initial template and prevhash before accepting connections
let fallback_token = fallback_coordinator.token();
loop {
let has_required_data = this
.last_future_template
.with(|template| template.is_some())
.map_err(JDCError::shutdown)?
&& this
.last_new_prev_hash
.with(|prev_hash| prev_hash.is_some())
.map_err(JDCError::shutdown)?;

if has_required_data {
info!("Required template data received, ready to accept connections");
break;
}

warn!("Waiting for initial template and prevhash from Template Provider...");
warn!("Is the Bitcoin node undergoing IBD?");
select! {
_ = cancellation_token.cancelled() => {
info!("Channel Manager: received shutdown while waiting for templates");
return Ok(());
}
_ = fallback_token.cancelled() => {
info!("Channel Manager: received fallback while waiting for templates");
return Ok(());
}
_ = tokio::time::sleep(std::time::Duration::from_secs(1)) => {}
}
}

// Bind before spawning, so that a bind failure is a startup failure the caller can
// propagate. Everything after this point (waiting for the first template, then accepting)
// runs in a background task, so bootstrap is never held hostage by the Template Provider.
info!("Starting downstream server at {listening_address}");
let server = TcpListener::bind(listening_address).await.map_err(|e| {
error!(error = ?e, "Failed to bind downstream server at {listening_address}");
Expand All @@ -609,6 +582,45 @@ impl ChannelManager {
// for this accept loop to stop before attempting to re-bind the same port.
let fallback_handler = fallback_coordinator.register();
task_manager.spawn(async move {
// Wait for initial template and prevhash before accepting connections
let ready = loop {
let has_required_data = match (
this.last_future_template.with(|template| template.is_some()),
this.last_new_prev_hash.with(|prev_hash| prev_hash.is_some()),
) {
(Ok(has_template), Ok(has_prev_hash)) => has_template && has_prev_hash,
_ => {
error!("Channel Manager: shared state poisoned while waiting for templates");
cancellation_token.cancel();
break false;
}
};

if has_required_data {
info!("Required template data received, ready to accept connections");
break true;
}

warn!("Waiting for initial template and prevhash from Template Provider...");
warn!("Is the Bitcoin node undergoing IBD?");
select! {
_ = cancellation_token.cancelled() => {
info!("Channel Manager: received shutdown while waiting for templates");
break false;
}
_ = fallback_token.cancelled() => {
info!("Channel Manager: received fallback while waiting for templates");
break false;
}
_ = tokio::time::sleep(std::time::Duration::from_secs(1)) => {}
}
};

if !ready {
fallback_handler.done();
return;
}

loop {
select! {
_ = cancellation_token.cancelled() => {
Expand Down Expand Up @@ -705,6 +717,9 @@ impl ChannelManager {
}
}
info!("Downstream server: Unified loop break");
// Release the port before signalling the fallback coordinator, so a subsequent
// fallback can re-bind the same listening address.
drop(server);
fallback_handler.done();
});
Ok(())
Expand All @@ -721,18 +736,8 @@ impl ChannelManager {
fallback_coordinator: FallbackCoordinator,
task_manager: Arc<TaskManager>,
coinbase_outputs: Vec<TxOut>,
) {
if let Err(e) = self.coinbase_output_constraints(coinbase_outputs).await {
error!(error = ?e, "Failed to send CoinbaseOutputConstraints message to TP");
if let Action::Shutdown = e.action {
warn!(
error_kind = ?e.kind,
"CoinbaseOutputConstraints requested shutdown; cancelling global token"
);
cancellation_token.cancel();
}
return;
}
) -> JDCResult<(), error::ChannelManager> {
self.coinbase_output_constraints(coinbase_outputs).await?;

task_manager.spawn(async move {
// we just spawned a new task that's relevant to fallback coordination
Expand Down Expand Up @@ -830,6 +835,8 @@ impl ChannelManager {
// signal fallback coordinator that this task has completed its cleanup
fallback_handler.done();
});

Ok(())
}

// Removes a downstream entry from the Channel Manager’s state.
Expand Down
31 changes: 31 additions & 0 deletions miner-apps/jd-client/src/lib/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@ use stratum_apps::{
};
use tokio::time::error::Elapsed;

use crate::config::ConfigJDCMode;

pub type JDCResult<T, Owner> = Result<T, JDCError<Owner>>;

#[derive(Debug)]
Expand All @@ -58,6 +60,9 @@ pub struct Upstream;
#[derive(Debug)]
pub struct Downstream;

#[derive(Debug)]
pub struct JobDeclaratorClient;

#[derive(Debug)]
pub struct JDCError<Owner> {
pub kind: JDCErrorKind,
Expand All @@ -73,6 +78,12 @@ pub enum Action {
Shutdown,
}

impl Action {
pub fn is_shutdown(self) -> bool {
matches!(self, Self::Shutdown)
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LoopControl {
Continue,
Expand All @@ -85,12 +96,14 @@ impl CanDisconnect for ChannelManager {}
impl CanFallback for Upstream {}
impl CanFallback for JobDeclarator {}
impl CanFallback for ChannelManager {}
impl CanFallback for JobDeclaratorClient {}

impl CanShutdown for ChannelManager {}
impl CanShutdown for TemplateProvider {}
impl CanShutdown for Downstream {}
impl CanShutdown for Upstream {}
impl CanShutdown for JobDeclarator {}
impl CanShutdown for JobDeclaratorClient {}

impl<O> JDCError<O> {
pub fn log<E: Into<JDCErrorKind>>(kind: E) -> Self {
Expand Down Expand Up @@ -259,6 +272,14 @@ pub enum JDCErrorKind {
InvalidKey,
/// Upstream not found
UpstreamNotFound,
/// Cannot determine Bitcoin data directory
InvalidBitcoinDataDir,
/// No upstream specified for pooled mining
NoUpstreamConfig(ConfigJDCMode),
/// Invalid coinbase output in config
InvalidCoinbaseOutput,
/// Cannot initialize monitoring tasks
MonitoringServerError(String),
}

impl std::error::Error for JDCErrorKind {}
Expand Down Expand Up @@ -402,6 +423,16 @@ impl fmt::Display for JDCErrorKind {
CouldNotInitiateSystem => write!(f, "Could not initiate subsystem"),
InvalidKey => write!(f, "Invalid key used during noise handshake"),
UpstreamNotFound => write!(f, "Upstream not found"),
InvalidBitcoinDataDir => write!(
f,
"Could not determine Bitcoin data directory. Please set data_dir in config."
),
NoUpstreamConfig(mode) => write!(
f,
"No upstreams configured for {mode:?} mode - at least one upstream is required"
),
InvalidCoinbaseOutput => write!(f, "Invalid coinbase output in config"),
MonitoringServerError(e) => write!(f, "Failed to initialize monitoring tasks: `{e:?}`"),
}
}
}
Expand Down
Loading
Loading