fix: join the RNS loop thread before exit and keep the response slot per-request
This commit is contained in:
parent
39c4c4d6a5
commit
234b75651f
17 changed files with 84 additions and 11 deletions
1
cache/23913df1956165db1fde57c54afd29ca3fe771bc27a64e989b395f01fc9ba0e9
vendored
Normal file
1
cache/23913df1956165db1fde57c54afd29ca3fe771bc27a64e989b395f01fc9ba0e9
vendored
Normal file
|
|
@ -0,0 +1 @@
|
||||||
|
{"raw":"0f0018d0b427deefc9e827d8f905811ec679055626d3978c15abe954935d2f3e0eb4b6ac55db9f0318110b4463818f60e91d7596c2a32b033f0b0e37f17ea7da2690bf092bf719e2bed7111127c107abe61985","sent_at":1.790665185e9}
|
||||||
1
cache/4892c693602a28c2f4f30931f544453549589423759c216e35c31bee4c5413fc
vendored
Normal file
1
cache/4892c693602a28c2f4f30931f544453549589423759c216e35c31bee4c5413fc
vendored
Normal file
|
|
@ -0,0 +1 @@
|
||||||
|
{"raw":"0f007dcb759264a15f6ba66c335a2f3e8ace05739b67970ea95221dfb2c4e87b395bd124fae0f8c71e0f5a26f28c4c393f7695e2434e223e5d152b14a6b689e9c448b03fb6f2cab9e8fafa73efba7d6df5e58b","sent_at":1.790666247e9}
|
||||||
1
cache/b053eb65854a98dc4037b10914488d3062fb0d19bcab66e0c6f68063b488a168
vendored
Normal file
1
cache/b053eb65854a98dc4037b10914488d3062fb0d19bcab66e0c6f68063b488a168
vendored
Normal file
|
|
@ -0,0 +1 @@
|
||||||
|
{"raw":"0f001bac2103482a7558d16836454bee96a4059bc6f761475deb9d31eb16c966e95aeaded51b951f7b4e71be6f9ce139bce6fd0de49ff495d48717b3fc4f99a76063d9fa0892f40ea53f5f1ed430ac6804de55","sent_at":1.790665991e9}
|
||||||
1
cache/c3529808c5c5a39bba63a0996781692de3d681a1bfe32c48e1dad0ed37c4b3eb
vendored
Normal file
1
cache/c3529808c5c5a39bba63a0996781692de3d681a1bfe32c48e1dad0ed37c4b3eb
vendored
Normal file
|
|
@ -0,0 +1 @@
|
||||||
|
{"raw":"0f0041f7951ed2f222f77ee7a5def7be3b2d05ee03605065c6ca397b44e77add0187344e3c523b0c1999e0ac7987def42d47c6358b09de502cc3158040391e21cafd42fe8099371189158099517e5abd3dccbb","sent_at":1.790666418e9}
|
||||||
|
|
@ -154,7 +154,12 @@ enum Command {
|
||||||
|
|
||||||
fn main() {
|
fn main() {
|
||||||
let cli = Cli::parse();
|
let cli = Cli::parse();
|
||||||
if let Err(e) = run(cli) {
|
let result = run(cli);
|
||||||
|
// The RNS carrier owns a loop thread that must be joined before the
|
||||||
|
// process exits, or it races the teardown of the statics it calls into.
|
||||||
|
#[cfg(feature = "rns")]
|
||||||
|
fumi::rns::stop();
|
||||||
|
if let Err(e) = result {
|
||||||
eprintln!("error: {e}");
|
eprintln!("error: {e}");
|
||||||
hint(e.downcast_ref::<fumi::error::Error>());
|
hint(e.downcast_ref::<fumi::error::Error>());
|
||||||
std::process::exit(1);
|
std::process::exit(1);
|
||||||
|
|
|
||||||
|
|
@ -41,6 +41,10 @@ static RNS::Link active_link({RNS::Type::NONE});
|
||||||
static microStore::FileSystem filesystem{microStore::Adapters::UniversalFileSystem()};
|
static microStore::FileSystem filesystem{microStore::Adapters::UniversalFileSystem()};
|
||||||
|
|
||||||
static volatile bool running = false;
|
static volatile bool running = false;
|
||||||
|
// Joinable, never detached: the loop thread must be joined by
|
||||||
|
// smolmail_rns_stop before the process tears down statics, or it keeps
|
||||||
|
// calling reticulum.loop() while their destructors run.
|
||||||
|
static std::thread loop_thread;
|
||||||
|
|
||||||
// The single response slot (plan: one request in flight at a time).
|
// The single response slot (plan: one request in flight at a time).
|
||||||
static std::mutex slot_mutex;
|
static std::mutex slot_mutex;
|
||||||
|
|
@ -94,9 +98,11 @@ static void on_failed(const RNS::RequestReceipt& receipt) {
|
||||||
slot_cv.notify_all();
|
slot_cv.notify_all();
|
||||||
}
|
}
|
||||||
|
|
||||||
// The unwrapped large-response payload, filled in by smolmail_rns_request.
|
// The unwrapped large-response payload, filled in per request below. It is
|
||||||
static RNS::Bytes unwrapped;
|
// deliberately local: a static here kept the last resource-path response
|
||||||
|
// alive past its request, so every later bare response -- an empty fetch
|
||||||
|
// page, an ack -- was misread as that stale page and the client looped
|
||||||
|
// FETCH/DELETE pairs forever.
|
||||||
static void loop_thread_main() {
|
static void loop_thread_main() {
|
||||||
while (running) {
|
while (running) {
|
||||||
reticulum.loop();
|
reticulum.loop();
|
||||||
|
|
@ -139,7 +145,7 @@ extern "C" int smolmail_rns_start(const char* storage_dir,
|
||||||
reticulum.start();
|
reticulum.start();
|
||||||
|
|
||||||
running = true;
|
running = true;
|
||||||
std::thread(loop_thread_main).detach();
|
loop_thread = std::thread(loop_thread_main);
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -256,6 +262,7 @@ extern "C" int smolmail_rns_request(const uint8_t* request, size_t request_len,
|
||||||
// raw remainder to handle_response on that path). The wrapper is
|
// raw remainder to handle_response on that path). The wrapper is
|
||||||
// unambiguous: a smolmail status is a single byte below 16, and the
|
// unambiguous: a smolmail status is a single byte below 16, and the
|
||||||
// msgpack bin headers are 0xC4..0xC6, so only those are unwrapped.
|
// msgpack bin headers are 0xC4..0xC6, so only those are unwrapped.
|
||||||
|
RNS::Bytes unwrapped;
|
||||||
const RNS::Bytes& payload = [&]() -> const RNS::Bytes& {
|
const RNS::Bytes& payload = [&]() -> const RNS::Bytes& {
|
||||||
if (slot_response.size() > 0 && slot_response.data()[0] >= 0xC4
|
if (slot_response.size() > 0 && slot_response.data()[0] >= 0xC4
|
||||||
&& slot_response.data()[0] <= 0xC6) {
|
&& slot_response.data()[0] <= 0xC6) {
|
||||||
|
|
@ -286,3 +293,23 @@ extern "C" void smolmail_rns_close(void) {
|
||||||
active_link = RNS::Link({RNS::Type::NONE});
|
active_link = RNS::Link({RNS::Type::NONE});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
extern "C" void smolmail_rns_stop(void) {
|
||||||
|
if (!running) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
// Order matters: halt and join the loop thread first, so nothing is
|
||||||
|
// inside Reticulum, Transport or the filesystem while they are torn
|
||||||
|
// down; only then take the link, the interface and the instance apart.
|
||||||
|
running = false;
|
||||||
|
if (loop_thread.joinable()) {
|
||||||
|
loop_thread.join();
|
||||||
|
}
|
||||||
|
if (active_link) {
|
||||||
|
active_link.teardown();
|
||||||
|
active_link = RNS::Link({RNS::Type::NONE});
|
||||||
|
}
|
||||||
|
RNS::Transport::deregister_interface(udp_interface);
|
||||||
|
udp_interface.stop();
|
||||||
|
reticulum = RNS::Reticulum({RNS::Type::NONE});
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -52,6 +52,14 @@ int smolmail_rns_request(const uint8_t *request, size_t request_len,
|
||||||
* minutes (upstream spec sec 13.10). */
|
* minutes (upstream spec sec 13.10). */
|
||||||
void smolmail_rns_close(void);
|
void smolmail_rns_close(void);
|
||||||
|
|
||||||
|
/* Stops the Reticulum stack: halts and joins the loop thread, tears the link
|
||||||
|
* down, deregisters and stops the UDP interface, and releases the
|
||||||
|
* Reticulum instance. The loop thread must be joined before the process (or
|
||||||
|
* embedding host) tears down statics, or it keeps calling into Reticulum
|
||||||
|
* while their destructors run. A no-op when the stack is not running.
|
||||||
|
* The stack can be started again afterwards. */
|
||||||
|
void smolmail_rns_stop(void);
|
||||||
|
|
||||||
#ifdef __cplusplus
|
#ifdef __cplusplus
|
||||||
}
|
}
|
||||||
#endif
|
#endif
|
||||||
|
|
|
||||||
|
|
@ -42,6 +42,8 @@ mod inner {
|
||||||
) -> c_int;
|
) -> c_int;
|
||||||
|
|
||||||
pub fn smolmail_rns_close();
|
pub fn smolmail_rns_close();
|
||||||
|
|
||||||
|
pub fn smolmail_rns_stop();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -150,6 +152,14 @@ pub fn close() {
|
||||||
unsafe { inner::smolmail_rns_close() };
|
unsafe { inner::smolmail_rns_close() };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Stops the Reticulum stack: joins the loop thread, tears the link down,
|
||||||
|
/// stops the UDP interface and releases the instance. Must run before the
|
||||||
|
/// embedding process tears down statics; a no-op when never started, and
|
||||||
|
/// the stack can be started again afterwards.
|
||||||
|
pub fn stop() {
|
||||||
|
unsafe { inner::smolmail_rns_stop() };
|
||||||
|
}
|
||||||
|
|
||||||
/// The destination hash length, for callers checking address shapes.
|
/// The destination hash length, for callers checking address shapes.
|
||||||
pub const DEST_LEN: usize = 16;
|
pub const DEST_LEN: usize = 16;
|
||||||
const _: () = assert!(DEST_LEN == 16 && KEY_LEN == 32);
|
const _: () = assert!(DEST_LEN == 16 && KEY_LEN == 32);
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,10 @@
|
||||||
//! the C shim, dispatched through the same operations as TCP.
|
//! the C shim, dispatched through the same operations as TCP.
|
||||||
//!
|
//!
|
||||||
//! The carrier starts lazily, on the first `smol+rns://` dial, so local-only
|
//! The carrier starts lazily, on the first `smol+rns://` dial, so local-only
|
||||||
//! commands stay usable offline and start instantly. A client creates no
|
//! commands stay usable offline and start instantly, and stops via `stop()`,
|
||||||
|
//! which every embedding process must call before it tears down: the
|
||||||
|
//! Reticulum loop thread otherwise races the destruction of the statics it
|
||||||
|
//! is calling into. A client creates no
|
||||||
//! Reticulum identity (upstream spec sec 13.8): the shim never calls
|
//! Reticulum identity (upstream spec sec 13.8): the shim never calls
|
||||||
//! `Link::identify`, and the flags below configure only the storage path and
|
//! `Link::identify`, and the flags below configure only the storage path and
|
||||||
//! the UDP interface.
|
//! the UDP interface.
|
||||||
|
|
@ -11,6 +14,7 @@ pub mod ffi;
|
||||||
pub mod transport;
|
pub mod transport;
|
||||||
|
|
||||||
use std::path::Path;
|
use std::path::Path;
|
||||||
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
use std::sync::{Mutex, OnceLock};
|
use std::sync::{Mutex, OnceLock};
|
||||||
|
|
||||||
use crate::error::Error;
|
use crate::error::Error;
|
||||||
|
|
@ -44,7 +48,7 @@ impl Default for RnsConfig {
|
||||||
}
|
}
|
||||||
|
|
||||||
static CONFIG: OnceLock<RnsConfig> = OnceLock::new();
|
static CONFIG: OnceLock<RnsConfig> = OnceLock::new();
|
||||||
static STARTED: OnceLock<()> = OnceLock::new();
|
static STARTED: AtomicBool = AtomicBool::new(false);
|
||||||
static START_MUTEX: Mutex<()> = Mutex::new(());
|
static START_MUTEX: Mutex<()> = Mutex::new(());
|
||||||
|
|
||||||
/// Records the carrier configuration; the first dial wins, as in a CLI the
|
/// Records the carrier configuration; the first dial wins, as in a CLI the
|
||||||
|
|
@ -53,19 +57,33 @@ pub fn configure(config: RnsConfig) {
|
||||||
let _ = CONFIG.set(config);
|
let _ = CONFIG.set(config);
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Starts the Reticulum stack at most once, on the first RNS dial.
|
/// Starts the Reticulum stack at most once between calls to `stop`, on the
|
||||||
|
/// first RNS dial.
|
||||||
pub fn ensure_started() -> Result<(), Error> {
|
pub fn ensure_started() -> Result<(), Error> {
|
||||||
if STARTED.get().is_some() {
|
if STARTED.load(Ordering::Acquire) {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
let _guard = START_MUTEX.lock().unwrap();
|
let _guard = START_MUTEX.lock().unwrap();
|
||||||
if STARTED.get().is_some() {
|
if STARTED.load(Ordering::Acquire) {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
let config = CONFIG.get().cloned().unwrap_or_default();
|
let config = CONFIG.get().cloned().unwrap_or_default();
|
||||||
std::fs::create_dir_all(Path::new(&config.storage_dir))
|
std::fs::create_dir_all(Path::new(&config.storage_dir))
|
||||||
.map_err(|e| Error::Other(format!("cannot create {}: {e}", config.storage_dir)))?;
|
.map_err(|e| Error::Other(format!("cannot create {}: {e}", config.storage_dir)))?;
|
||||||
ffi::start(&config)?;
|
ffi::start(&config)?;
|
||||||
let _ = STARTED.set(());
|
STARTED.store(true, Ordering::Release);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Stops the Reticulum stack and joins its loop thread. The shim's loop
|
||||||
|
/// thread otherwise outlives the caller and races the teardown of its statics
|
||||||
|
/// at process exit — a nondeterministic hang the RNS carrier hit on a
|
||||||
|
/// first-run provisioning pass. Embedders call this when they are done with
|
||||||
|
/// the carrier; the CLI calls it before exiting. A no-op when never started,
|
||||||
|
/// and the stack starts again on the next dial.
|
||||||
|
pub fn stop() {
|
||||||
|
if !STARTED.swap(false, Ordering::AcqRel) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
ffi::stop();
|
||||||
|
}
|
||||||
|
|
|
||||||
BIN
hashlist_store/index.dat
Normal file
BIN
hashlist_store/index.dat
Normal file
Binary file not shown.
BIN
hashlist_store/seg0.dat
Normal file
BIN
hashlist_store/seg0.dat
Normal file
Binary file not shown.
BIN
hashlist_store/seg1.dat
Normal file
BIN
hashlist_store/seg1.dat
Normal file
Binary file not shown.
BIN
known_store/index.dat
Normal file
BIN
known_store/index.dat
Normal file
Binary file not shown.
BIN
known_store/seg0.dat
Normal file
BIN
known_store/seg0.dat
Normal file
Binary file not shown.
BIN
path_store/index.dat
Normal file
BIN
path_store/index.dat
Normal file
Binary file not shown.
BIN
path_store/seg0.dat
Normal file
BIN
path_store/seg0.dat
Normal file
Binary file not shown.
1
transport_identity
Normal file
1
transport_identity
Normal file
|
|
@ -0,0 +1 @@
|
||||||
|
pqαΥι„Β›KΚ¬^}ƒΑ¨Τ<>0¬λ%Ίlnn<6E>„µJ<C2B5>yσ:f}pU¤<<3C>ο<EFBFBD>~BZδfαΏΟΝ£<CE9D>μ/H¬
|
||||||
Loading…
Add table
Add a link
Reference in a new issue