@@ -13,7 +13,7 @@ use std::fmt;
1313use std:: future:: Future ;
1414#[ cfg( test) ]
1515use std:: panic:: RefUnwindSafe ;
16- use std:: sync:: atomic:: { AtomicU64 , AtomicUsize , Ordering } ;
16+ use std:: sync:: atomic:: { AtomicU64 , Ordering } ;
1717use std:: sync:: { Arc , Mutex } ;
1818use std:: time:: Duration ;
1919
@@ -45,6 +45,7 @@ use vss_client::util::storable_builder::{EntropySource, StorableBuilder};
4545use crate :: entropy:: NodeEntropy ;
4646use crate :: io:: utils:: check_namespace_key_validity;
4747use crate :: lnurl_auth:: LNURL_AUTH_HARDENED_CHILD_INDEX ;
48+ use crate :: runtime:: StoreRuntime ;
4849
4950type CustomRetryPolicy = FilteredRetryPolicy <
5051 JitteredRetryPolicy <
@@ -85,13 +86,8 @@ pub struct VssStore {
8586 // Version counter to ensure that writes are applied in the correct order. It is assumed that read and list
8687 // operations aren't sensitive to the order of execution.
8788 next_version : AtomicU64 ,
88- // A VSS-internal runtime we use to avoid any deadlocks we could hit when waiting on a spawned
89- // blocking task to finish while the blocked thread had acquired the reactor. In particular,
90- // this works around a previously-hit case where a concurrent call to
91- // `PeerManager::process_pending_events` -> `ChannelManager::get_and_clear_pending_msg_events`
92- // would deadlock when trying to acquire sync `Mutex` locks that are held by the thread
93- // currently being blocked waiting on the VSS operation to finish.
94- internal_runtime : Option < tokio:: runtime:: Runtime > ,
89+ // A VSS-internal runtime that drives VSS I/O independently from the node runtime.
90+ internal_runtime : Option < Arc < StoreRuntime > > ,
9591}
9692
9793impl VssStore {
@@ -100,52 +96,46 @@ impl VssStore {
10096 header_provider : Arc < dyn VssHeaderProvider > ,
10197 ) -> io:: Result < Self > {
10298 let next_version = AtomicU64 :: new ( 1 ) ;
103- let internal_runtime = tokio:: runtime:: Builder :: new_multi_thread ( )
104- . enable_all ( )
105- . thread_name_fn ( || {
106- static ATOMIC_ID : AtomicUsize = AtomicUsize :: new ( 0 ) ;
107- let id = ATOMIC_ID . fetch_add ( 1 , Ordering :: SeqCst ) ;
108- format ! ( "ldk-node-vss-runtime-{}" , id)
109- } )
110- . worker_threads ( INTERNAL_RUNTIME_WORKERS )
111- . max_blocking_threads ( INTERNAL_RUNTIME_WORKERS )
112- . build ( )
113- . map_err ( |e| {
114- io:: Error :: new ( io:: ErrorKind :: Other , format ! ( "Failed to build VSS runtime: {}" , e) )
115- } ) ?;
99+ let internal_runtime =
100+ Arc :: new ( StoreRuntime :: new ( "ldk-node-vss-runtime" , INTERNAL_RUNTIME_WORKERS , "VSS" ) ?) ;
116101
117102 let ( data_encryption_key, obfuscation_master_key) =
118103 derive_data_encryption_and_obfuscation_keys ( & vss_seed) ;
119104 let key_obfuscator = KeyObfuscator :: new ( obfuscation_master_key) ;
105+ let setup_key_obfuscator = KeyObfuscator :: new ( obfuscation_master_key) ;
120106
121107 let mut entropy_seed = [ 0u8 ; 32 ] ;
122108 getrandom:: fill ( & mut entropy_seed) . expect ( "Failed to generate random bytes" ) ;
123109 let entropy_source = RandomBytes :: new ( entropy_seed) ;
110+ let setup_entropy_source = RandomBytes :: new ( entropy_seed) ;
124111
125- let sync_retry_policy = retry_policy ( ) ;
126- let blocking_client = VssClient :: new_with_headers (
112+ let setup_retry_policy = retry_policy ( ) ;
113+ let setup_client = VssClient :: new_with_headers (
127114 base_url. clone ( ) ,
128- sync_retry_policy ,
129- header_provider . clone ( ) ,
115+ setup_retry_policy ,
116+ Arc :: clone ( & header_provider ) ,
130117 ) ;
131118
132- let runtime_handle = internal_runtime. handle ( ) ;
133- let schema_version = tokio:: task:: block_in_place ( || {
119+ let async_retry_policy = retry_policy ( ) ;
120+ let async_client =
121+ VssClient :: new_with_headers ( base_url, async_retry_policy, header_provider) ;
122+
123+ let setup_store_id = store_id. clone ( ) ;
124+ let runtime_handle = internal_runtime. handle ( ) . clone ( ) ;
125+ let schema_version = std:: thread:: spawn ( move || {
134126 runtime_handle. block_on ( async {
135127 determine_and_write_schema_version (
136- & blocking_client ,
137- & store_id ,
128+ & setup_client ,
129+ & setup_store_id ,
138130 data_encryption_key,
139- & key_obfuscator ,
140- & entropy_source ,
131+ & setup_key_obfuscator ,
132+ & setup_entropy_source ,
141133 )
142134 . await
143135 } )
144- } ) ?;
145-
146- let async_retry_policy = retry_policy ( ) ;
147- let async_client =
148- VssClient :: new_with_headers ( base_url, async_retry_policy, header_provider) ;
136+ } )
137+ . join ( )
138+ . map_err ( |_| io:: Error :: new ( io:: ErrorKind :: Other , "VSS schema setup task panicked" ) ) ??;
149139
150140 let inner = Arc :: new ( VssStoreInner :: new (
151141 schema_version,
@@ -158,6 +148,10 @@ impl VssStore {
158148
159149 Ok ( Self { inner, next_version, internal_runtime : Some ( internal_runtime) } )
160150 }
151+
152+ fn internal_runtime ( & self ) -> Arc < StoreRuntime > {
153+ Arc :: clone ( self . internal_runtime . as_ref ( ) . expect ( "VSS runtime must be available" ) )
154+ }
161155 /// Returns a [`VssStoreBuilder`] allowing to build a [`VssStore`].
162156 pub fn builder (
163157 node_entropy : NodeEntropy , vss_url : String , store_id : String , network : Network ,
@@ -200,10 +194,16 @@ impl KVStore for VssStore {
200194 let secondary_namespace = secondary_namespace. to_string ( ) ;
201195 let key = key. to_string ( ) ;
202196 let inner = Arc :: clone ( & self . inner ) ;
197+ let runtime = self . internal_runtime ( ) ;
203198 async move {
204- inner
205- . read_internal ( & inner. async_client , primary_namespace, secondary_namespace, key)
206- . await
199+ let task = runtime. spawn ( async move {
200+ inner
201+ . read_internal ( & inner. async_client , primary_namespace, secondary_namespace, key)
202+ . await
203+ } ) ;
204+ task. await . map_err ( |e| {
205+ io:: Error :: new ( io:: ErrorKind :: Other , format ! ( "VSS runtime task failed: {}" , e) )
206+ } ) ?
207207 }
208208 }
209209 fn write (
@@ -215,19 +215,25 @@ impl KVStore for VssStore {
215215 let secondary_namespace = secondary_namespace. to_string ( ) ;
216216 let key = key. to_string ( ) ;
217217 let inner = Arc :: clone ( & self . inner ) ;
218+ let runtime = self . internal_runtime ( ) ;
218219 async move {
219- inner
220- . write_internal (
221- & inner. async_client ,
222- inner_lock_ref,
223- locking_key,
224- version,
225- primary_namespace,
226- secondary_namespace,
227- key,
228- buf,
229- )
230- . await
220+ let task = runtime. spawn ( async move {
221+ inner
222+ . write_internal (
223+ & inner. async_client ,
224+ inner_lock_ref,
225+ locking_key,
226+ version,
227+ primary_namespace,
228+ secondary_namespace,
229+ key,
230+ buf,
231+ )
232+ . await
233+ } ) ;
234+ task. await . map_err ( |e| {
235+ io:: Error :: new ( io:: ErrorKind :: Other , format ! ( "VSS runtime task failed: {}" , e) )
236+ } ) ?
231237 }
232238 }
233239 fn remove (
@@ -239,6 +245,7 @@ impl KVStore for VssStore {
239245 let secondary_namespace = secondary_namespace. to_string ( ) ;
240246 let key = key. to_string ( ) ;
241247 let inner = Arc :: clone ( & self . inner ) ;
248+ let runtime = self . internal_runtime ( ) ;
242249 let fut = async move {
243250 inner
244251 . remove_internal (
@@ -254,10 +261,15 @@ impl KVStore for VssStore {
254261 } ;
255262 async move {
256263 if lazy {
257- tokio:: task:: spawn ( async move { fut. await } ) ;
264+ runtime. spawn ( async move {
265+ let _ = fut. await ;
266+ } ) ;
258267 Ok ( ( ) )
259268 } else {
260- fut. await
269+ let task = runtime. spawn ( fut) ;
270+ task. await . map_err ( |e| {
271+ io:: Error :: new ( io:: ErrorKind :: Other , format ! ( "VSS runtime task failed: {}" , e) )
272+ } ) ?
261273 }
262274 }
263275 }
@@ -267,16 +279,27 @@ impl KVStore for VssStore {
267279 let primary_namespace = primary_namespace. to_string ( ) ;
268280 let secondary_namespace = secondary_namespace. to_string ( ) ;
269281 let inner = Arc :: clone ( & self . inner ) ;
282+ let runtime = self . internal_runtime ( ) ;
270283 async move {
271- inner. list_internal ( & inner. async_client , primary_namespace, secondary_namespace) . await
284+ let task = runtime. spawn ( async move {
285+ inner
286+ . list_internal ( & inner. async_client , primary_namespace, secondary_namespace)
287+ . await
288+ } ) ;
289+ task. await . map_err ( |e| {
290+ io:: Error :: new ( io:: ErrorKind :: Other , format ! ( "VSS runtime task failed: {}" , e) )
291+ } ) ?
272292 }
273293 }
274294}
275295
276296impl Drop for VssStore {
277297 fn drop ( & mut self ) {
278- let internal_runtime = self . internal_runtime . take ( ) ;
279- tokio:: task:: block_in_place ( move || drop ( internal_runtime) ) ;
298+ if let Some ( runtime) = self . internal_runtime . take ( ) {
299+ if let Ok ( runtime) = Arc :: try_unwrap ( runtime) {
300+ runtime. shutdown_background ( ) ;
301+ }
302+ }
280303 }
281304}
282305
0 commit comments