do {
if (!backfill_listener().budget_available()) {
- post_event(RequestWaiting{});
+ if (backfill_state().progress_tracker->tracked_objects_completed()) {
+ INFODPP("backfill budget unavailable with nothing in flight, "
+ "entering BudgetBlocked", pg());
+ post_event(RequestBudgetBlocked{});
+ } else {
+ post_event(RequestWaiting{});
+ }
return;
} else if (should_rescan_replicas(backfill_state().peer_backfill_info,
primary_bi)) {
return discard_event();
}
+// -- BudgetBlocked
+BackfillState::BudgetBlocked::BudgetBlocked(my_context ctx)
+ : my_base(ctx)
+{
+ LOG_PREFIX(BackfillState::BudgetBlocked::BudgetBlocked);
+ DEBUGDPP("budget unavailable with nothing in flight, "
+ "waiting for throttle slot to become available", pg());
+ backfill_listener().request_budget_retry();
+}
+
+boost::statechart::result
+BackfillState::BudgetBlocked::react(SuspendBackfill evt)
+{
+ LOG_PREFIX(BackfillState::BudgetBlocked::react::SuspendBackfill);
+ DEBUGDPP("suspended within BudgetBlocked", pg());
+ backfill_state().on_suspended();
+ return discard_event();
+}
+
+boost::statechart::result
+BackfillState::BudgetBlocked::react(Triggered evt)
+{
+ LOG_PREFIX(BackfillState::BudgetBlocked::react::Triggered);
+ ceph_assert(backfill_state().is_suspended());
+ if (backfill_state().on_resumed()) {
+ DEBUGDPP("Backfill resumed, going Enqueuing", pg());
+ return transit<Enqueuing>();
+ }
+ return discard_event();
+}
+
+boost::statechart::result
+BackfillState::BudgetBlocked::react(BudgetAvailable evt)
+{
+ LOG_PREFIX(BackfillState::BudgetBlocked::react::BudgetAvailable);
+ DEBUGDPP("BudgetBlocked::react() on BudgetAvailable", pg());
+ if (!backfill_state().is_suspended()) {
+ return transit<Enqueuing>();
+ } else {
+ DEBUGDPP("backfill suspended, not going Enqueuing", pg());
+ backfill_state().go_enqueuing_on_resume();
+ }
+ return discard_event();
+}
+
// -- Done
BackfillState::Done::Done(my_context ctx)
: my_base(ctx)
struct SuspendBackfill : sc::event<SuspendBackfill> {
};
+ struct RequestBudgetBlocked : sc::event<RequestBudgetBlocked> {
+ };
+
+ struct BudgetAvailable : sc::event<BudgetAvailable> {
+ };
+
private:
// internal events
struct RequestPrimaryScanning : sc::event<RequestPrimaryScanning> {
struct ReplicasScanning;
struct Waiting;
struct Done;
+ struct BudgetBlocked;
struct BackfillMachine : sc::state_machine<BackfillMachine, Initial> {
BackfillMachine(BackfillState& backfill_state,
sc::transition<RequestPrimaryScanning, PrimaryScanning>,
sc::transition<RequestReplicasScanning, ReplicasScanning>,
sc::transition<RequestWaiting, Waiting>,
+ sc::transition<RequestBudgetBlocked, BudgetBlocked>,
sc::transition<sc::event_base, Crashed>>;
explicit Enqueuing(my_context);
}
};
+ struct BudgetBlocked : sc::state<BudgetBlocked, BackfillMachine>,
+ StateHelper<BudgetBlocked> {
+ using reactions = boost::mpl::list<
+ sc::custom_reaction<BudgetAvailable>,
+ sc::custom_reaction<SuspendBackfill>,
+ sc::custom_reaction<Triggered>,
+ sc::transition<sc::event_base, Crashed>>;
+ explicit BudgetBlocked(my_context ctx);
+ sc::result react(BudgetAvailable);
+ sc::result react(SuspendBackfill);
+ sc::result react(Triggered);
+ };
+
BackfillState(BackfillListener& backfill_listener,
std::unique_ptr<PeeringFacade> peering_state,
std::unique_ptr<PGFacade> pg);
virtual void backfilled() = 0;
+ virtual void request_budget_retry() = 0;
+
virtual ~BackfillListener() = default;
};
#if FMT_VERSION >= 90000
template <> struct fmt::formatter<crimson::osd::BackfillState::PGFacade>
: fmt::ostream_formatter {};
+template <> struct fmt::formatter<crimson::osd::BackfillState::BudgetBlocked>
+ : fmt::ostream_formatter {};
#endif
if (!added)
return;
peering_state.prepare_backfill_for_missing(obj, v, peers);
+ // release the budget_retry_releaser now that we're dispatching a real
+ // push -- recover_object_with_throttle will acquire its own slot
+ budget_retry_releaser.reset();
std::ignore = recover_object_with_throttle(obj, v).\
handle_exception_interruptible([] (auto) {
ceph_abort_msg("got exception on backfill's push");
{
LOG_PREFIX(PGRecovery::update_peers_last_backfill);
DEBUGDPP("obj={} v={} target={}", *pg->get_dpp(), obj, v, target);
+ // release the budget_retry_releaser now that we're dispatching work
+ // (same as enqueue_push) -- slot was held to guarantee Enqueuing
+ // found budget available, drops don't need their own throttle slot
+ budget_retry_releaser.reset();
// allocate a pair if target is seen for the first time
auto& req = backfill_drop_requests[target];
if (!req) {
return ss.throttle_available();
}
+PGRecovery::interruptible_future<>
+PGRecovery::do_request_budget_retry()
+{
+ LOG_PREFIX(PGRecovery::do_request_budget_retry);
+ DEBUGDPP("budget unavailable with nothing in flight, "
+ "waiting for throttle slot", *pg->get_dpp());
+ auto releaser = co_await get_backfill_throttle();
+ budget_retry_in_flight = false;
+ if (!backfill_state) {
+ DEBUGDPP("backfill_state is null, skipping BudgetAvailable "
+ "(pg cleaned or interval changed)", *pg->get_dpp());
+ co_return;
+ }
+ // hold the releaser so the slot stays acquired until Enqueuing
+ budget_retry_releaser.emplace(std::move(releaser));
+ if (backfill_state->is_triggered()) {
+ backfill_state->post_event(
+ BackfillState::BudgetAvailable{}.intrusive_from_this());
+ } else {
+ backfill_state->process_event(
+ BackfillState::BudgetAvailable{}.intrusive_from_this());
+ }
+}
+
+void PGRecovery::request_budget_retry()
+{
+ LOG_PREFIX(PGRecovery::request_budget_retry);
+ if (budget_retry_in_flight) {
+ DEBUGDPP("budget retry already in flight, skipping", *pg->get_dpp());
+ return;
+ }
+ budget_retry_in_flight = true;
+ std::ignore = do_request_budget_retry();
+}
+
void PGRecovery::on_pg_clean()
{
replica_scan_throttle_releasers.clear();
+ budget_retry_releaser.reset();
+ budget_retry_in_flight = false;
backfill_state.reset();
}
LOG_PREFIX(PGRecovery::on_activate_complete);
DEBUGDPP("backfill_state={}", *pg->get_dpp(), fmt::ptr(backfill_state.get()));
replica_scan_throttle_releasers.clear();
+ budget_retry_releaser.reset();
+ budget_retry_in_flight = false;
backfill_state.reset();
}
return backend->recover_object(soid, need);
}
+ interruptible_future<> do_request_budget_retry();
+
// backfill begin
std::unique_ptr<crimson::osd::BackfillState> backfill_state;
std::map<pg_shard_t,
MURef<MOSDPGBackfillRemove>> backfill_drop_requests;
std::map<pg_shard_t,
OperationThrottler::ThrottleReleaser> replica_scan_throttle_releasers;
-
+ std::optional<OperationThrottler::ThrottleReleaser> budget_retry_releaser;
template <class EventT>
void start_backfill_recovery(
const EventT& evt);
void update_peers_last_backfill(
const hobject_t& new_last_backfill) final;
bool budget_available() const final;
+ void request_budget_retry() final;
template <typename T>
void start_peering_event_operation_listener(T &&evt, float delay = 0);
friend crimson::osd::BackfillState::PGFacade;
friend crimson::osd::PG;
+ bool budget_retry_in_flight = false;
// backfill end
};