@@ -182,6 +182,17 @@ impl NamespaceStore {
182182 namespace : NamespaceName ,
183183 restore_option : RestoreOption ,
184184 ) -> anyhow:: Result < ( ) > {
185+ // Reset destroys the namespace's data and writes nothing to the metastore, so the fence
186+ // check is made here, under the namespace's transition lock: a fence command either
187+ // finished before this check or starts after the reset (section 3.3, lifecycle).
188+ let _transition = self
189+ . inner
190+ . fences
191+ . controller ( & namespace)
192+ . begin_transition ( )
193+ . await ;
194+ self . check_lifecycle ( & namespace) ?;
195+
185196 // The process for reseting is as follow:
186197 // - get a lock on the namespace entry, if the entry exists, then it's a lock on the entry,
187198 // if it doesn't exist, insert an empty entry and take a lock on it
@@ -221,18 +232,26 @@ impl NamespaceStore {
221232 Box :: new ( move |op| {
222233 let this = this. clone ( ) ;
223234 tokio:: spawn ( async move {
224- match op {
225- ResetOp :: Reset ( ns) => {
226- tracing:: info!( "received reset signal for: {ns}" ) ;
227- if let Err ( e) = this. reset ( ns. clone ( ) , RestoreOption :: Latest ) . await {
228- tracing:: error!( "error resetting namespace `{ns}`: {e}" ) ;
229- }
230- }
231- }
235+ let _ = this. handle_reset_op ( op) . await ;
232236 } ) ;
233237 } )
234238 }
235239
240+ /// A reset requested by a replica's replicator. A namespace whose fence denies lifecycle
241+ /// work is not reset: the refusal is logged and returned.
242+ async fn handle_reset_op ( & self , op : ResetOp ) -> anyhow:: Result < ( ) > {
243+ match op {
244+ ResetOp :: Reset ( ns) => {
245+ tracing:: info!( "received reset signal for: {ns}" ) ;
246+ let result = self . reset ( ns. clone ( ) , RestoreOption :: Latest ) . await ;
247+ if let Err ( e) = & result {
248+ tracing:: error!( "error resetting namespace `{ns}`: {e}" ) ;
249+ }
250+ result
251+ }
252+ }
253+ }
254+
236255 pub async fn fork (
237256 & self ,
238257 from : NamespaceName ,
@@ -245,14 +264,23 @@ impl NamespaceStore {
245264 }
246265
247266 // The destination is refused before anything is stored for it when it is being created
248- // as a migration target or its fence state is unknown.
267+ // as a migration target, its fence state is unknown, or its fence denies lifecycle work
268+ // (an existing fenced namespace, whose directory the fork would otherwise replace).
249269 self . inner . fences . check_available ( & to) ?;
270+ self . check_lifecycle ( & to) ?;
250271
251272 // check that the source namespace exists
252273 if !self . inner . metadata . exists ( & from) . await {
253274 return Err ( crate :: error:: Error :: NamespaceDoesntExist ( from. to_string ( ) ) ) ;
254275 }
255276
277+ // A fork reads the source's data without a read lease, so it runs under the source's
278+ // transition lock and checks the source's gate under it: a fence command on the source
279+ // (a write or read fence) either finished before this check, and the fork is refused,
280+ // or waits for the fork to finish (section 3.3, fork as source).
281+ let _from_transition = self . inner . fences . controller ( & from) . begin_transition ( ) . await ;
282+ self . check_lifecycle ( & from) ?;
283+
256284 let to_entry = self
257285 . inner
258286 . store
@@ -262,6 +290,9 @@ impl NamespaceStore {
262290 if to_lock. is_some ( ) {
263291 return Err ( crate :: error:: Error :: NamespaceAlreadyExist ( to. to_string ( ) ) ) ;
264292 }
293+ // With the destination's entry held, a fence command cannot load the destination, so
294+ // the check cannot go stale before the fork has stored and flushed its config.
295+ self . check_lifecycle ( & to) ?;
265296
266297 // FIXME: we could potentially delete the namespace while trying to fork it
267298 if !self . inner . metadata . exists ( & from) . await {
@@ -462,8 +493,10 @@ impl NamespaceStore {
462493 db_config : DatabaseConfig ,
463494 ) -> crate :: Result < ( ) > {
464495 // A name that is being created as a migration target, or whose fence state is unknown,
465- // is refused before anything is stored for it.
496+ // is refused before anything is stored for it; so is a name whose fence denies lifecycle
497+ // work (creating over an existing record, with or without a restore).
466498 self . inner . fences . check_available ( & namespace) ?;
499+ self . check_lifecycle ( & namespace) ?;
467500 if let Some ( shared_schema_name) = & db_config. shared_schema_name {
468501 // we hold a lock for the duration of the namespace creation
469502 let _lock = self
@@ -552,6 +585,16 @@ impl NamespaceStore {
552585 & self . inner . metadata
553586 }
554587
588+ /// Refuse generic lifecycle and configuration work on `namespace` (config mutation,
589+ /// delete, reset, fork on either side, create over an existing record, restore, dump load,
590+ /// shared-schema linking, schema migration) while its fence denies it
591+ /// (`docs/NAMESPACE_FENCE.md` section 3.3), without loading the namespace. A name without
592+ /// fence state is not refused here: the existing checks apply to it. Paths that persist
593+ /// through the metastore are refused again inside its transaction.
594+ pub ( crate ) fn check_lifecycle ( & self , namespace : & NamespaceName ) -> crate :: Result < ( ) > {
595+ Ok ( self . inner . fences . check_lifecycle ( namespace) ?)
596+ }
597+
555598 /// Run one fence command on its namespace, including the drain it starts
556599 /// (`docs/NAMESPACE_FENCE.md` sections 5.3 and 8). `AcquireSourceWriteFence` loads the
557600 /// namespace first, so that its connection manager and replication log are registered with
@@ -947,6 +990,7 @@ pub(crate) mod fence_tests {
947990
948991 use super :: * ;
949992 use crate :: config:: MetaStoreConfig ;
993+ use crate :: connection:: Connection as _;
950994 use crate :: namespace:: configurator:: { BaseNamespaceConfig , PrimaryConfig , PrimaryConfigurator } ;
951995 use crate :: namespace:: fence:: command:: { FenceCommand , FenceRequest } ;
952996 use crate :: namespace:: fence:: outcome:: { FenceDetail , FenceOutcome } ;
@@ -1148,4 +1192,208 @@ pub(crate) mod fence_tests {
11481192 store. destroy ( "ns" . into ( ) , false ) . await . unwrap ( ) ;
11491193 assert ! ( store. inner. fences. get( & "ns" . into( ) ) . is_none( ) ) ;
11501194 }
1195+
1196+ fn release ( ns : & ' static str , command_id : u128 ) -> FenceRequest {
1197+ FenceRequest {
1198+ namespace : ns. into ( ) ,
1199+ operation_id : OP ,
1200+ command_id : Uuid :: from_u128 ( command_id) ,
1201+ expected_state : FenceState :: SourceDraining ,
1202+ expected_revision : 1 ,
1203+ command : FenceCommand :: ReleaseSourceWriteFence ,
1204+ }
1205+ }
1206+
1207+ /// Create `ns` holding a table `t` with one row.
1208+ async fn create_with_row ( store : & NamespaceStore , ns : & ' static str ) {
1209+ store
1210+ . create ( ns. into ( ) , RestoreOption :: Latest , Default :: default ( ) )
1211+ . await
1212+ . unwrap ( ) ;
1213+ let conn = store
1214+ . with ( ns. into ( ) , |ns| ns. db . connection_maker ( ) )
1215+ . await
1216+ . unwrap ( )
1217+ . create ( )
1218+ . await
1219+ . unwrap ( ) ;
1220+ tokio:: task:: spawn_blocking ( move || {
1221+ conn. with_raw ( |c| c. execute_batch ( "create table t (x); insert into t values (1);" ) )
1222+ } )
1223+ . await
1224+ . unwrap ( )
1225+ . unwrap ( ) ;
1226+ }
1227+
1228+ /// The number of rows in `ns`'s table `t`, or the error reading it.
1229+ async fn rows ( store : & NamespaceStore , ns : & ' static str ) -> rusqlite:: Result < i64 > {
1230+ let conn = store
1231+ . with ( ns. into ( ) , |ns| ns. db . connection_maker ( ) )
1232+ . await
1233+ . unwrap ( )
1234+ . create ( )
1235+ . await
1236+ . unwrap ( ) ;
1237+ tokio:: task:: spawn_blocking ( move || {
1238+ conn. with_raw ( |c| c. query_row ( "select count(*) from t" , ( ) , |r| r. get ( 0 ) ) )
1239+ } )
1240+ . await
1241+ . unwrap ( )
1242+ }
1243+
1244+ #[ track_caller]
1245+ fn assert_fenced ( result : crate :: Result < ( ) > , outcome : FenceOutcome ) {
1246+ match result {
1247+ Err ( Error :: NamespaceFence ( e) ) => assert_eq ! ( e. outcome( ) , outcome, "{e}" ) ,
1248+ other => panic ! ( "expected {outcome}, got {other:?}" ) ,
1249+ }
1250+ }
1251+
1252+ #[ track_caller]
1253+ fn assert_fenced_anyhow ( result : anyhow:: Result < ( ) > , outcome : FenceOutcome ) {
1254+ match result {
1255+ Err ( e) => match e. downcast_ref :: < Error > ( ) {
1256+ Some ( Error :: NamespaceFence ( e) ) => assert_eq ! ( e. outcome( ) , outcome, "{e}" ) ,
1257+ _ => panic ! ( "expected {outcome}, got {e:?}" ) ,
1258+ } ,
1259+ Ok ( ( ) ) => panic ! ( "expected {outcome}, got Ok" ) ,
1260+ }
1261+ }
1262+
1263+ /// Reset, which destroys the namespace's data and recreates it, is refused while the
1264+ /// namespace is fenced, both called directly and as the replicator's reset callback does.
1265+ #[ tokio:: test( flavor = "multi_thread" ) ]
1266+ async fn reset_refused_while_fenced ( ) {
1267+ let tmp = tempdir ( ) . unwrap ( ) ;
1268+ let store = open_store ( tmp. path ( ) ) . await ;
1269+ create_with_row ( & store, "ns" ) . await ;
1270+ let fence = store. inner . fences . controller ( & "ns" . into ( ) ) ;
1271+ fence
1272+ . apply_command ( store. meta_store ( ) , acquire ( "ns" ) , ctx ( ) )
1273+ . await
1274+ . unwrap ( ) ;
1275+
1276+ assert_fenced_anyhow (
1277+ store. reset ( "ns" . into ( ) , RestoreOption :: Latest ) . await ,
1278+ FenceOutcome :: MigrationWriteFenced ,
1279+ ) ;
1280+ assert_fenced_anyhow (
1281+ store. handle_reset_op ( ResetOp :: Reset ( "ns" . into ( ) ) ) . await ,
1282+ FenceOutcome :: MigrationWriteFenced ,
1283+ ) ;
1284+ // The namespace was not touched and still serves reads.
1285+ assert_eq ! ( rows( & store, "ns" ) . await . unwrap( ) , 1 ) ;
1286+
1287+ // Once the fence is released, reset works as before, and its data is gone.
1288+ fence
1289+ . apply_command ( store. meta_store ( ) , release ( "ns" , 2 ) , ctx ( ) )
1290+ . await
1291+ . unwrap ( ) ;
1292+ assert_eq ! ( fence. gate( ) . state( ) , FenceState :: Released ) ;
1293+ store
1294+ . handle_reset_op ( ResetOp :: Reset ( "ns" . into ( ) ) )
1295+ . await
1296+ . unwrap ( ) ;
1297+ assert ! ( rows( & store, "ns" ) . await . is_err( ) ) ;
1298+ }
1299+
1300+ /// Fork is lifecycle work on both sides: a fenced source is not copied, and a fenced
1301+ /// destination (whose directory a fork would replace) is not overwritten. Create over a
1302+ /// fenced name, delete and config mutation, including linking the namespace to a shared
1303+ /// schema, are refused too.
1304+ #[ tokio:: test( flavor = "multi_thread" ) ]
1305+ async fn lifecycle_refused_while_fenced ( ) {
1306+ let tmp = tempdir ( ) . unwrap ( ) ;
1307+ let store = open_store ( tmp. path ( ) ) . await ;
1308+ create_with_row ( & store, "src" ) . await ;
1309+ create_with_row ( & store, "other" ) . await ;
1310+ let fence = store. inner . fences . controller ( & "src" . into ( ) ) ;
1311+ fence
1312+ . apply_command ( store. meta_store ( ) , acquire ( "src" ) , ctx ( ) )
1313+ . await
1314+ . unwrap ( ) ;
1315+ let write_fenced = FenceOutcome :: MigrationWriteFenced ;
1316+
1317+ // Fork with the fenced namespace as the source: nothing is created.
1318+ assert_fenced (
1319+ store
1320+ . fork ( "src" . into ( ) , "copy" . into ( ) , Default :: default ( ) , None )
1321+ . await ,
1322+ write_fenced,
1323+ ) ;
1324+ assert ! ( !store. exists( & "copy" . into( ) ) . await ) ;
1325+ assert ! ( !tmp. path( ) . join( "dbs" ) . join( "copy" ) . exists( ) ) ;
1326+ // Fork onto the fenced namespace: its data and config are untouched.
1327+ let config_before = store. config_store ( "src" . into ( ) ) . await . unwrap ( ) . get ( ) ;
1328+ assert_fenced (
1329+ store
1330+ . fork (
1331+ "other" . into ( ) ,
1332+ "src" . into ( ) ,
1333+ DatabaseConfig {
1334+ block_reason : Some ( "fork" . into ( ) ) ,
1335+ ..Default :: default ( )
1336+ } ,
1337+ None ,
1338+ )
1339+ . await ,
1340+ write_fenced,
1341+ ) ;
1342+ assert_eq ! ( rows( & store, "src" ) . await . unwrap( ) , 1 ) ;
1343+ let config_after = store. config_store ( "src" . into ( ) ) . await . unwrap ( ) . get ( ) ;
1344+ assert_eq ! ( config_after. block_reason, config_before. block_reason) ;
1345+ assert_eq ! ( config_after. block_writes, config_before. block_writes) ;
1346+
1347+ // Create over it, with or without a restore, and delete.
1348+ assert_fenced (
1349+ store
1350+ . create ( "src" . into ( ) , RestoreOption :: Latest , Default :: default ( ) )
1351+ . await ,
1352+ write_fenced,
1353+ ) ;
1354+ assert_fenced ( store. destroy ( "src" . into ( ) , false ) . await , write_fenced) ;
1355+
1356+ // Config mutation, including linking the namespace to a shared schema, is refused in the
1357+ // metastore transaction that would store it.
1358+ let handle = store. config_store ( "src" . into ( ) ) . await . unwrap ( ) ;
1359+ assert_fenced (
1360+ handle
1361+ . store ( DatabaseConfig {
1362+ block_reason : Some ( "changed" . into ( ) ) ,
1363+ ..Default :: default ( )
1364+ } )
1365+ . await ,
1366+ write_fenced,
1367+ ) ;
1368+ assert_fenced (
1369+ handle
1370+ . store ( DatabaseConfig {
1371+ shared_schema_name : Some ( "other" . into ( ) ) ,
1372+ ..Default :: default ( )
1373+ } )
1374+ . await ,
1375+ write_fenced,
1376+ ) ;
1377+ assert_eq ! ( handle. get( ) . block_reason, config_before. block_reason) ;
1378+ assert ! ( handle. get( ) . shared_schema_name. is_none( ) ) ;
1379+ assert_eq ! ( rows( & store, "src" ) . await . unwrap( ) , 1 ) ;
1380+
1381+ // After release the same operations follow the existing policy again.
1382+ fence
1383+ . apply_command ( store. meta_store ( ) , release ( "src" , 2 ) , ctx ( ) )
1384+ . await
1385+ . unwrap ( ) ;
1386+ store
1387+ . fork ( "src" . into ( ) , "copy" . into ( ) , Default :: default ( ) , None )
1388+ . await
1389+ . unwrap ( ) ;
1390+ assert_eq ! ( rows( & store, "copy" ) . await . unwrap( ) , 1 ) ;
1391+ assert ! ( matches!(
1392+ store
1393+ . create( "src" . into( ) , RestoreOption :: Latest , Default :: default ( ) )
1394+ . await ,
1395+ Err ( Error :: NamespaceAlreadyExist ( _) )
1396+ ) ) ;
1397+ store. destroy ( "src" . into ( ) , false ) . await . unwrap ( ) ;
1398+ }
11511399}
0 commit comments