diff --git a/src/container/srv_target.c b/src/container/srv_target.c index 7e9b13f804a..d2256d3343e 100644 --- a/src/container/srv_target.c +++ b/src/container/srv_target.c @@ -63,10 +63,11 @@ agg_rate_ctl(void *arg) return -1; /* - * XXX temporary workaround: EC aggregation needs to be paused during rebuilding - * to avoid the race between EC rebuild and EC aggregation. + * EC aggregation needs to be paused by the rebuild EC agg barrier to avoid + * the race between EC rebuild and EC aggregation. **/ - if (ds_pool_is_rebuilding(pool) && cont->sc_ec_agg_active && !param->ap_vos_agg) + if (!param->ap_vos_agg && cont->sc_ec_agg_active && + ds_pool_child_ec_agg_paused(cont->sc_pool)) return -1; /* When system is idle or under space pressure, let aggregation run in tight mode */ @@ -207,10 +208,10 @@ cont_aggregate_runnable(struct ds_cont_child *cont, struct sched_request *req, return false; } - if (ds_pool_is_rebuilding(pool) && !vos_agg) { - D_DEBUG(DB_EPC, DF_CONT ": skip EC aggregation during rebuild %d, %d.\n", + if (ds_pool_child_ec_agg_paused(cont->sc_pool) && !vos_agg) { + D_DEBUG(DB_EPC, DF_CONT ": skip EC aggregation during pause gate %u.\n", DP_CONT(cont->sc_pool->spc_uuid, cont->sc_uuid), - atomic_load(&pool->sp_rebuilding), atomic_load(&pool->sp_rebuild_enum)); + atomic_load(&cont->sc_pool->spc_ec_agg_pause_gate)); return false; } @@ -511,10 +512,8 @@ cont_aggregate_interval(struct ds_cont_child *cont, cont_aggregate_cb_t cb, if (dss_ult_exiting(req)) break; - /* sleep 18 seconds for EC aggregation ULT if the pool is in rebuilding, - * if no space pressure. - */ - if (ds_pool_is_rebuilding(cont->sc_pool->spc_pool) && !param->ap_vos_agg && + /* Sleep longer while the rebuild EC agg barrier is active if no space pressure. */ + if (ds_pool_child_ec_agg_paused(cont->sc_pool) && !param->ap_vos_agg && msecs != 200) msecs = 18000; @@ -964,11 +963,11 @@ ds_cont_child_reset_ec_agg_eph_all(struct ds_pool_child *pool_child) #define WAIT_EC_PAUSE_MAX 600 -void +int ds_cont_child_wait_ec_agg_pause(struct ds_pool_child *pool_child, int wait_timeout) { uint64_t start_time = daos_wallclock_secs(); - int wait_intv = 10; + int next_warn = 60; int waited = 0; D_DEBUG(DB_MD, DF_UUID "[%d]: wait for pausing EC aggregation\n", @@ -981,35 +980,31 @@ ds_cont_child_wait_ec_agg_pause(struct ds_pool_child *pool_child, int wait_timeo struct ds_cont_child *coc; bool paused = true; - /* Wait for pausing aggregation - * XXX: There is no global barrier so we always wait for at least 10 seconds to - * lower the chance that remote targets are still running EC aggregation. - */ - if (wait_intv > wait_timeout - waited) - wait_intv = wait_timeout - waited; + if (!ds_pool_child_ec_agg_paused(pool_child)) + return -DER_OP_CANCELED; - dss_sleep(wait_intv * 1000); d_list_for_each_entry(coc, &pool_child->spc_cont_list, sc_link) { if (ds_cont_child_ec_aggregating(coc)) { - /* Aggregation is active on this container */ paused = false; break; } } if (paused) - return; + return 0; waited = daos_wallclock_secs() - start_time; if (waited >= wait_timeout) { D_WARN("can't pause EC aggregation after %d seconds\n", waited); - return; /* XXX what can I do? */ + return -DER_TIMEDOUT; } - if (waited % 60 == 0) { + if (waited >= next_warn) { D_WARN(DF_UUID "[%d]: waited %d secs for EC aggregation to pause\n", DP_UUID(pool_child->spc_uuid), dss_get_module_info()->dmi_tgt_id, waited); + next_warn += 60; } + dss_sleep(1000); } } diff --git a/src/include/daos/rpc.h b/src/include/daos/rpc.h index 030e9a078dd..0e93a10e42a 100644 --- a/src/include/daos/rpc.h +++ b/src/include/daos/rpc.h @@ -71,7 +71,7 @@ enum daos_module_id { #define DAOS_POOL_VERSION 7 #define DAOS_CONT_VERSION 9 #define DAOS_OBJ_VERSION 10 -#define DAOS_REBUILD_VERSION 5 +#define DAOS_REBUILD_VERSION 6 #define DAOS_RSVC_VERSION 5 #define DAOS_RDB_VERSION 5 #define DAOS_RDBT_VERSION 3 diff --git a/src/include/daos_srv/container.h b/src/include/daos_srv/container.h index f7c84d807e1..76b6afaec9b 100644 --- a/src/include/daos_srv/container.h +++ b/src/include/daos_srv/container.h @@ -228,7 +228,7 @@ ds_cont_child_destroy(uuid_t pool_uuid, uuid_t cont_uuid); void ds_cont_child_reset_ec_agg_eph_all(struct ds_pool_child *pool_child); -void +int ds_cont_child_wait_ec_agg_pause(struct ds_pool_child *pool_child, int wait_timeout); /** initialize a csummer based on container properties. Will retrieve the diff --git a/src/include/daos_srv/pool.h b/src/include/daos_srv/pool.h index dc4a7c747d7..84956edd059 100644 --- a/src/include/daos_srv/pool.h +++ b/src/include/daos_srv/pool.h @@ -183,6 +183,10 @@ struct ds_pool_child { ABT_eventual spc_ref_eventual; uint64_t spc_no_storage : 1; /* The pool shard has no storage. */ + ATOMIC uint32_t spc_ec_agg_pause_gate; + ATOMIC uint64_t spc_ec_agg_pause_term; + /* Debug state set by RB_OP_REBUILD START after the global EC agg pause. */ + ATOMIC uint64_t spc_rebuild_ec_agg_paused_hlc; uint32_t spc_reint_mode; uint32_t *spc_state; /* Pointer to ds_pool->sp_states[i] */ @@ -212,6 +216,26 @@ ds_pool_is_rebuilding(struct ds_pool *pool) return (atomic_load(&pool->sp_rebuilding) > 0 || atomic_load(&pool->sp_rebuild_enum) > 0); } +static inline bool +ds_pool_child_ec_agg_paused(struct ds_pool_child *child) +{ + return atomic_load(&child->spc_ec_agg_pause_gate) != 0; +} + +static inline bool +ds_pool_child_rebuild_started(struct ds_pool_child *child) +{ + return atomic_load(&child->spc_rebuild_ec_agg_paused_hlc) != 0; +} + +static inline bool +ds_pool_child_ec_agg_token_match(struct ds_pool_child *child, uint64_t leader_term, + uint32_t rebuild_gen) +{ + return atomic_load(&child->spc_ec_agg_pause_gate) == rebuild_gen && + atomic_load(&child->spc_ec_agg_pause_term) <= leader_term; +} + /* encode metadata RPC operation key: HLC time first, in network order, for keys sorted by time. * allocates the byte-stream, caller must free with D_FREE(). */ diff --git a/src/object/srv_ec_aggregate.c b/src/object/srv_ec_aggregate.c index cec63f5fe5a..c5475e81f5d 100644 --- a/src/object/srv_ec_aggregate.c +++ b/src/object/srv_ec_aggregate.c @@ -1340,7 +1340,7 @@ agg_peer_check_avail(struct ec_agg_param *agg_param, struct ec_agg_entry *entry) int i; int rc; - if (ds_pool_is_rebuilding(agg_param->ap_cont->sc_pool->spc_pool)) { + if (ds_pool_child_ec_agg_paused(agg_param->ap_cont->sc_pool)) { rc = -DER_OP_CANCELED; /* We currently pause EC aggregation for rebuild, so just cancel the * aggregation for the current stripe. It means the following peer status @@ -1426,8 +1426,15 @@ agg_peer_update_ult(void *arg) if (unlikely(DAOS_FAIL_CHECK(DAOS_FORCE_EC_AGG_PEER_FAIL))) D_GOTO(out, rc = -DER_TIMEDOUT); - agg_param = container_of(entry, struct ec_agg_param, ap_agg_entry); + D_ASSERTF(!ds_pool_child_rebuild_started(agg_param->ap_cont->sc_pool), + DF_UUID " EC aggregation peer update after rebuild START: " + "rank %u tgt %d " DF_UOID " stripe " DF_U64 " global_pause_done " DF_X64 + "\n", + DP_UUID(agg_param->ap_pool_info.api_pool_uuid), dss_self_rank(), + dss_get_module_info()->dmi_tgt_id, DP_UOID(entry->ae_oid), + entry->ae_cur_stripe.as_stripenum, + atomic_load(&agg_param->ap_cont->sc_pool->spc_rebuild_ec_agg_paused_hlc)); rc = agg_peer_check_avail(agg_param, entry); if (rc != 0) D_GOTO(out, rc); @@ -1939,9 +1946,10 @@ agg_process_stripe(struct ec_agg_param *agg_param, struct ec_agg_entry *entry) /* avoid race between EC aggregation and rebuild scanner */ agg_param->ap_cont->sc_ec_agg_updates++; - if (ds_pool_is_rebuilding(agg_param->ap_cont->sc_pool->spc_pool)) { + if (ds_pool_child_ec_agg_paused(agg_param->ap_cont->sc_pool)) { rc = -DER_OP_CANCELED; - DL_INFO(rc, DF_UOID " abort as rebuild started", DP_UOID(entry->ae_oid)); + DL_INFO(rc, DF_UOID " abort as EC aggregation pause gate is set", + DP_UOID(entry->ae_oid)); update_vos = false; goto out; } @@ -2366,10 +2374,10 @@ ec_aggregate_yield(struct ec_agg_param *agg_param) { int rc; - if (ds_pool_is_rebuilding(agg_param->ap_pool_info.api_pool)) { - D_INFO(DF_UUID ": abort ec aggregation, sp_rebuilding %d\n", + if (ds_pool_child_ec_agg_paused(agg_param->ap_cont->sc_pool)) { + D_INFO(DF_UUID ": abort EC aggregation, pause gate %u\n", DP_UUID(agg_param->ap_pool_info.api_pool->sp_uuid), - atomic_load(&agg_param->ap_pool_info.api_pool->sp_rebuilding)); + atomic_load(&agg_param->ap_cont->sc_pool->spc_ec_agg_pause_gate)); return true; } @@ -2578,15 +2586,15 @@ agg_iterate_pre_cb(daos_handle_t ih, vos_iter_entry_t *entry, D_ASSERT(agg_param->ap_initialized); - /* If rebuild started, abort it to save conflict window with rebuild + /* If the rebuild EC agg barrier is active, abort to save conflict window with rebuild * (see obj_inflight_io_check()). */ - if (ds_pool_is_rebuilding(agg_param->ap_pool_info.api_pool)) { + if (ds_pool_child_ec_agg_paused(agg_param->ap_cont->sc_pool)) { rc = -DER_OP_CANCELED; - DL_INFO(rc, DF_CONT " abort as rebuild started, sp_rebuilding %d.", + DL_INFO(rc, DF_CONT " abort as EC aggregation pause gate %u is set.", DP_CONT(agg_param->ap_pool_info.api_pool_uuid, agg_param->ap_pool_info.api_cont_uuid), - atomic_load(&agg_param->ap_pool_info.api_pool->sp_rebuilding)); + atomic_load(&agg_param->ap_cont->sc_pool->spc_ec_agg_pause_gate)); return rc; } @@ -2887,7 +2895,7 @@ cont_ec_aggregate_cb(struct ds_cont_child *cont, daos_epoch_range_t *epr, /* clear the flag before next turn's cont_aggregate_runnable(), to save conflict * window with rebuild (see obj_inflight_io_check()). */ - if (ds_pool_is_rebuilding(cont->sc_pool->spc_pool)) + if (ds_pool_child_ec_agg_paused(cont->sc_pool)) cont->sc_ec_agg_active = 0; if (rc == 0) { diff --git a/src/object/srv_obj.c b/src/object/srv_obj.c index 7ff04660fa9..3166567f100 100644 --- a/src/object/srv_obj.c +++ b/src/object/srv_obj.c @@ -2664,6 +2664,20 @@ ds_obj_ec_agg_handler(crt_rpc_t *rpc) } D_ASSERT(ioc.ioc_coc != NULL); + if (ds_pool_child_ec_agg_paused(ioc.ioc_coc->sc_pool)) { + rc = -DER_OP_CANCELED; + DL_INFO(rc, DF_CONT " reject EC aggregation update while rebuild pause gate is set", + DP_CONT(ioc.ioc_coc->sc_pool_uuid, ioc.ioc_coc->sc_uuid)); + goto out; + } + D_ASSERTF(!ds_pool_child_rebuild_started(ioc.ioc_coc->sc_pool), + DF_CONT " EC aggregation handler after rebuild START: " + "rank %u tgt %d " DF_UOID " stripe " DF_U64 " epoch " DF_X64 "-" DF_X64 + " global_pause_done " DF_X64 "\n", + DP_CONT(ioc.ioc_coc->sc_pool_uuid, ioc.ioc_coc->sc_uuid), dss_self_rank(), + dss_get_module_info()->dmi_tgt_id, DP_UOID(oea->ea_oid), oea->ea_stripenum, + oea->ea_epoch_range.epr_lo, oea->ea_epoch_range.epr_hi, + atomic_load(&ioc.ioc_coc->sc_pool->spc_rebuild_ec_agg_paused_hlc)); ioc.ioc_coc->sc_ec_agg_updates++; dkey = (daos_key_t *)&oea->ea_dkey; diff --git a/src/rebuild/rebuild_internal.h b/src/rebuild/rebuild_internal.h index dae19ae345f..509dfb02f20 100644 --- a/src/rebuild/rebuild_internal.h +++ b/src/rebuild/rebuild_internal.h @@ -87,22 +87,17 @@ struct rebuild_tgt_pool_tracker { /* new layout version for upgrade rebuild */ uint32_t rt_new_layout_ver; - unsigned int rt_lead_puller_running:1, - rt_abort:1, - /* re-report #rebuilt cnt per master change */ - rt_re_report:1, - rt_finishing:1, - rt_scan_done:1, - rt_global_scan_done:1, - rt_global_done:1; + unsigned int rt_lead_puller_running : 1, rt_abort : 1, + /* re-report #rebuilt cnt per master change */ + rt_re_report : 1, rt_finishing : 1, rt_ec_agg_paused : 1, rt_scan_done : 1, + rt_global_scan_done : 1, rt_global_done : 1; }; struct rebuild_server_status { d_rank_t rank; double last_update; uint32_t dtx_resync_version; - uint32_t scan_done:1, - pull_done:1; + uint32_t ec_agg_paused : 1, scan_done : 1, pull_done : 1; }; /* Track the rebuild status globally */ @@ -325,11 +320,8 @@ struct rebuild_iv { uint32_t riv_master_rank; uint32_t riv_ver; uint32_t riv_rebuild_gen; - uint32_t riv_global_done:1, - riv_global_scan_done:1, - riv_scan_done:1, - riv_pull_done:1, - riv_sync:1; + uint32_t riv_global_done : 1, riv_global_scan_done : 1, riv_scan_done : 1, + riv_pull_done : 1, riv_sync : 1, riv_ec_agg_paused : 1; int32_t riv_status; uint32_t riv_bukid; /* bucket ID */ }; diff --git a/src/rebuild/scan.c b/src/rebuild/scan.c index d96a6bddd8e..5d14253ef59 100644 --- a/src/rebuild/scan.c +++ b/src/rebuild/scan.c @@ -281,17 +281,18 @@ rebuild_cont_iter_cb(daos_handle_t ih, d_iov_t *key_iov, static int rpt_wait_rebuild_epoch(struct rebuild_tgt_pool_tracker *rpt) { - int wait_cnt = 0; - int wait_intv = 200; /* milliseconds */ - int wait_cnt_max = 180; + int wait_cnt = 0; - while (rpt->rt_stable_epoch == 0 && wait_cnt++ < wait_cnt_max) - dss_sleep(wait_intv); + while (rpt->rt_stable_epoch == 0 && !rpt->rt_abort && !rpt->rt_finishing && + !rpt->rt_global_done && wait_cnt++ < 180) + dss_sleep(200); if (rpt->rt_stable_epoch != 0) return 0; + if (!rpt->rt_abort && !rpt->rt_finishing && !rpt->rt_global_done) + return -DER_TIMEDOUT; - return -DER_TIMEDOUT; + return -DER_SHUTDOWN; } static void @@ -1162,8 +1163,10 @@ rebuild_scanner(void *data) child = ds_pool_child_lookup(rpt->rt_pool_uuid); if (child == NULL) D_GOTO(out, rc = -DER_NONEXIST); - - ds_cont_child_wait_ec_agg_pause(child, rebuild_wait_ec_pause); + if (rpt->rt_rebuild_op == RB_OP_REBUILD) { + D_ASSERT(rpt->rt_stable_epoch != 0); + atomic_store(&child->spc_rebuild_ec_agg_paused_hlc, rpt->rt_stable_epoch); + } /* There maybe orphan DTX entries after DTX resync, let's cleanup before rebuild scan. */ rc = dtx_cleanup_orphan(rpt->rt_pool_uuid, rpt->rt_pool->sp_dtx_resync_version); @@ -1228,6 +1231,34 @@ rebuild_scanner(void *data) return rc; } +static int +rebuild_ec_agg_pause_one(void *data) +{ + struct rebuild_tgt_pool_tracker *rpt = data; + struct ds_pool_child *child; + uint64_t current_term; + uint32_t current_gen; + int rc; + + child = ds_pool_child_lookup(rpt->rt_pool_uuid); + if (child == NULL) + return -DER_NONEXIST; + + current_term = atomic_load(&child->spc_ec_agg_pause_term); + current_gen = atomic_load(&child->spc_ec_agg_pause_gate); + if (current_term > rpt->rt_leader_term || + (current_term == rpt->rt_leader_term && current_gen > rpt->rt_rebuild_gen)) { + ds_pool_child_put(child); + return -DER_STALE; + } + + atomic_store(&child->spc_ec_agg_pause_term, rpt->rt_leader_term); + atomic_store(&child->spc_ec_agg_pause_gate, rpt->rt_rebuild_gen); + rc = ds_cont_child_wait_ec_agg_pause(child, rebuild_wait_ec_pause); + ds_pool_child_put(child); + return rc; +} + /** * Wait for pool map and setup global status, then spawn scanners for all * service xsteams @@ -1275,6 +1306,19 @@ rebuild_scan_leader(void *data) } do_scan: + if (rpt->rt_rebuild_op == RB_OP_REBUILD) { + rc = ds_pool_thread_collective( + rpt->rt_pool_uuid, PO_COMP_ST_NEW | PO_COMP_ST_DOWN | PO_COMP_ST_DOWNOUT, + rebuild_ec_agg_pause_one, rpt, 0); + if (rc != 0) + D_GOTO(out, rc); + + rpt->rt_ec_agg_paused = 1; + rc = rpt_wait_rebuild_epoch(rpt); + if (rc != 0) + D_GOTO(out, rc); + } + D_INFO(DF_RB " scan collective begin\n", DP_RB_RPT(rpt)); rc = ds_pool_thread_collective(rpt->rt_pool_uuid, PO_COMP_ST_NEW | PO_COMP_ST_DOWN | PO_COMP_ST_DOWNOUT, rebuild_scanner, rpt, @@ -1303,7 +1347,6 @@ rebuild_scan_leader(void *data) DL_INFO(rc, DF_RB " scan leader done", DP_RB_RPT(rpt)); rpt_put(rpt); } - /* Scan the local target and generate rebuild object list */ void rebuild_tgt_scan_handler(crt_rpc_t *rpc) @@ -1326,8 +1369,15 @@ rebuild_tgt_scan_handler(crt_rpc_t *rpc) DL_ERROR(rc, DF_RB " cannot find pool", DP_RB_RSI(rsi)); D_GOTO(out_put, rc); } + atomic_fetch_add(&pool->sp_rebuilding, 1); + if (rsi->rsi_leader_term < pool->sp_iv_ns->iv_master_term || + (rsi->rsi_leader_term == pool->sp_iv_ns->iv_master_term && + rsi->rsi_master_rank != pool->sp_iv_ns->iv_master_rank)) + D_GOTO(out_put, rc = -DER_STALE); + ds_pool_iv_ns_update(pool, rsi->rsi_master_rank, rsi->rsi_leader_term); + /* If PS leader has been changed, and rebuild version is also increased * due to adding new failure targets for rebuild, let's abort previous * rebuild. diff --git a/src/rebuild/srv.c b/src/rebuild/srv.c index 3371d34dc54..0d14cc23234 100644 --- a/src/rebuild/srv.c +++ b/src/rebuild/srv.c @@ -227,6 +227,17 @@ is_rebuild_global_scan_done(struct rebuild_global_pool_tracker *rgt) return true; } +static bool +is_rebuild_global_ec_agg_paused(struct rebuild_global_pool_tracker *rgt) +{ + int i; + + for (i = 0; i < rgt->rgt_servers_number; i++) + if (!rgt->rgt_servers[i].ec_agg_paused) + return false; + return true; +} + static bool is_rebuild_global_done(struct rebuild_global_pool_tracker *rgt) { @@ -254,6 +265,7 @@ is_rebuild_phase_mostly_done(int engines_done_ct, int engines_total_ct) #define SCAN_DONE 0x1 #define PULL_DONE 0x2 +#define EC_AGG_PAUSED 0x4 static void servers_sop_swap(void *array, int a, int b) @@ -335,6 +347,10 @@ rebuild_leader_set_status(struct rebuild_global_pool_tracker *rgt, D_DEBUG(DB_REBUILD, DF_RB " rank %d is pull_done", DP_RB_RGT(rgt), rank); status->pull_done = 1; } + if ((flags & EC_AGG_PAUSED) && !status->ec_agg_paused) { + D_DEBUG(DB_REBUILD, DF_RB " rank %d paused EC aggregation", DP_RB_RGT(rgt), rank); + status->ec_agg_paused = 1; + } } static void @@ -525,6 +541,24 @@ rebuild_global_status_update(struct rebuild_global_pool_tracker *rgt, struct rebuild_iv *iv) { rebuild_leader_set_update_time(rgt, iv->riv_rank); + if (iv->riv_stable_epoch != 0) { + if (rgt->rgt_stable_epoch == 0) { + rgt->rgt_stable_epoch = iv->riv_stable_epoch; + D_INFO(DF_RB ": recovered stable epoch " DF_X64 " from rank %u\n", + DP_RB_RGT(rgt), rgt->rgt_stable_epoch, iv->riv_rank); + } else if (rgt->rgt_stable_epoch != iv->riv_stable_epoch) { + D_ERROR(DF_RB ": stable epoch mismatch " DF_X64 "/" DF_X64 + " from rank %u\n", + DP_RB_RGT(rgt), rgt->rgt_stable_epoch, iv->riv_stable_epoch, + iv->riv_rank); + rgt->rgt_status.rs_errno = -DER_STALE; + rgt->rgt_status.rs_fail_rank = iv->riv_rank; + rgt->rgt_abort = 1; + } + } + if (iv->riv_ec_agg_paused) + rebuild_leader_set_status(rgt, iv->riv_rank, iv->riv_dtx_resyc_version, + EC_AGG_PAUSED); D_DEBUG(DB_REBUILD, DF_RB ": iv rank %d scan_done %d pull_done %d resync dtx %u\n", DP_RB_RGT(rgt), iv->riv_rank, iv->riv_scan_done, iv->riv_pull_done, @@ -1152,7 +1186,7 @@ rebuild_leader_status_check(struct ds_pool *pool, uint32_t op, if (rgt_up_dom_included(rgt, dom) == false) rebuild_leader_set_status(rgt, dom->do_comp.co_rank, RB_DTX_RESYNC_VER_SKIP, - SCAN_DONE | PULL_DONE); + EC_AGG_PAUSED | SCAN_DONE | PULL_DONE); } ABT_rwlock_unlock(pool->sp_lock); map_ranks_fini(&rank_list); @@ -1175,6 +1209,13 @@ rebuild_leader_status_check(struct ds_pool *pool, uint32_t op, goto done; } + if (rgt->rgt_opc == RB_OP_REBUILD && rgt->rgt_stable_epoch == 0 && + is_rebuild_global_ec_agg_paused(rgt)) { + rgt->rgt_stable_epoch = d_hlc_get(); + D_INFO(DF_RB " EC aggregation paused globally at " DF_X64 "\n", + DP_RB_RGT(rgt), rgt->rgt_stable_epoch); + } + /* leader send status to target engines if still running */ if (!rgt->rgt_abort && !is_rebuild_global_done(rgt) && myrank == pool->sp_iv_ns->iv_master_rank) @@ -1491,7 +1532,6 @@ rebuild_scan_broadcast(struct ds_pool *pool, struct rebuild_global_pool_tracker DL_ERROR(rc, DF_RB " failed to create scan broadcast request", DP_RB_RGT(rgt)); D_GOTO(out, rc); } - rsi = crt_req_get(rpc); uuid_copy(rsi->rsi_pool_uuid, pool->sp_uuid); rsi->rsi_ns_id = pool->sp_iv_ns->iv_ns_id; @@ -1515,7 +1555,8 @@ rebuild_scan_broadcast(struct ds_pool *pool, struct rebuild_global_pool_tracker DL_ERROR(rc, DF_RB " scan broadcast send failed.", DP_RB_RGT(rgt)); rgt->rgt_init_scan = 1; - rgt->rgt_stable_epoch = rso->rso_stable_epoch; + if (rebuild_op != RB_OP_REBUILD) + rgt->rgt_stable_epoch = rso->rso_stable_epoch; DL_INFO(rc, DF_RB " got stable/reclaim epoch " DF_X64 "/" DF_X64, DP_RB_RGT(rgt), rgt->rgt_stable_epoch, rgt->rgt_reclaim_epoch); crt_req_decref(rpc); @@ -2971,8 +3012,7 @@ static int rebuild_fini_one(void *arg) { struct rebuild_tgt_pool_tracker *rpt = arg; - struct rebuild_pool_tls *pool_tls; - struct ds_pool_child *dpc; + struct rebuild_pool_tls *pool_tls; pool_tls = rebuild_pool_tls_lookup(rpt->rt_pool_uuid, rpt->rt_rebuild_ver, rpt->rt_rebuild_gen); @@ -2982,14 +3022,24 @@ rebuild_fini_one(void *arg) rebuild_pool_tls_destroy(pool_tls); /* close the opened local ds_cont on main XS */ D_ASSERT(dss_get_module_info()->dmi_xs_id != 0); + return 0; +} + +static int +rebuild_ec_agg_fini_one(void *arg) +{ + struct rebuild_tgt_pool_tracker *rpt = arg; + struct ds_pool_child *dpc; dpc = ds_pool_child_lookup(rpt->rt_pool_uuid); - /* The ds_pool_child is already stopped */ if (dpc == NULL) return 0; - D_DEBUG(DB_REBUILD, DF_RB ": rebuild fini for stable epoch " DF_U64 "\n", DP_RB_RPT(rpt), - rpt->rt_stable_epoch); + if (ds_pool_child_ec_agg_token_match(dpc, rpt->rt_leader_term, rpt->rt_rebuild_gen)) { + atomic_store(&dpc->spc_ec_agg_pause_gate, 0); + atomic_store(&dpc->spc_ec_agg_pause_term, 0); + atomic_store(&dpc->spc_rebuild_ec_agg_paused_hlc, 0); + } ds_pool_child_put(dpc); return 0; } @@ -3035,7 +3085,10 @@ rebuild_tgt_fini(struct rebuild_tgt_pool_tracker *rpt) DL_WARN(rc, DF_RB " rebuild fini one failed", DP_RB_RPT(rpt)); /* destroy the migrate_tls of 0-xstream */ ds_migrate_stop(rpt->rt_pool, rpt->rt_rebuild_ver, rpt->rt_rebuild_gen); - /* No one should access rpt after rebuild_fini_one. */ + rc = dss_task_collective(rebuild_ec_agg_fini_one, rpt, 0); + if (rc != 0) + DL_WARN(rc, DF_RB " EC aggregation gate fini failed", DP_RB_RPT(rpt)); + /* No one should access rpt after the final gate cleanup. */ D_INFO(DF_RB " Finalized rebuild\n", DP_RB_RPT(rpt)); rpt_delete(rpt); } @@ -3106,6 +3159,8 @@ rebuild_tgt_status_check_ult(void *arg) rpt->rt_reported_size; } iv.riv_status = status.status; + iv.riv_ec_agg_paused = rpt->rt_ec_agg_paused; + iv.riv_stable_epoch = rpt->rt_stable_epoch; if (status.scanning == 0 || rpt->rt_abort || status.status != 0) { iv.riv_scan_done = 1;