#include "segmented_journal.h"
#include "crimson/common/config_proxy.h"
+#include "crimson/common/coroutine.h"
#include "crimson/os/seastore/logging.h"
SET_SUBSYS(seastore_journal);
return scan_last_segment(last_segment_id, last_header
).safe_then([this, FNAME, segments=std::move(segments)] {
INFO("dirty_tail={}, alloc_tail={}",
- trimmer.get_dirty_tail(),
- trimmer.get_alloc_tail());
+ get_dirty_tail(),
+ get_alloc_tail());
auto journal_tail = trimmer.get_journal_tail();
auto journal_tail_paddr = journal_tail.offset;
ceph_assert(journal_tail != JOURNAL_SEQ_NULL);
auto num_segments = segments.end() - from;
INFO("{} segments to replay", num_segments);
- auto ret = replay_segments_t(num_segments);
+ replay_segments.resize(num_segments);
std::transform(
- from, segments.end(), ret.begin(),
+ from, segments.end(), replay_segments.begin(),
[this](const auto &p) {
auto ret = journal_seq_t{
p.second.segment_seq,
};
return std::make_pair(ret, p.second);
});
- ret[0].first.offset = journal_tail_paddr;
+ replay_segments[0].first.offset = journal_tail_paddr;
return prep_replay_segments_fut(
- replay_ertr::ready_future_marker{},
- std::move(ret));
+ replay_ertr::ready_future_marker{});
});
}
SegmentedJournal::replay_segment(
journal_seq_t seq,
segment_header_t header,
- delta_handler_t &handler,
- replay_stats_t &stats)
+ scan_delta_handler_t &handler)
{
LOG_PREFIX(Journal::replay_segment);
INFO("starting at {} -- {}", seq, header);
return seastar::do_with(
scan_valid_records_cursor(seq),
SegmentManagerGroup::found_record_handler_t(
- [&handler, this, &stats](
+ [&handler, this](
record_locator_t locator,
const record_group_header_t& header,
const bufferlist& mdbuf)
[write_result=locator.write_result,
this,
FNAME,
- &handler,
- &stats](auto& record_deltas_list)
+ &handler](auto& record_deltas_list)
{
return crimson::do_for_each(
record_deltas_list,
[write_result,
this,
FNAME,
- &handler,
- &stats](record_deltas_t& record_deltas)
+ &handler](record_deltas_t& record_deltas)
{
++stats.num_records;
auto locator = record_locator_t{
return crimson::do_for_each(
record_deltas.deltas,
[locator,
- this,
- &handler,
- &stats](auto &p)
+ &handler](auto &p)
{
auto& modify_time = p.first;
auto& delta = p.second;
- return handler(
- locator,
- delta,
- trimmer.get_dirty_tail(),
- trimmer.get_alloc_tail(),
- modify_time
- ).safe_then([&stats, delta_type=delta.type](auto ret) {
- auto [is_applied, ext] = ret;
- if (is_applied) {
- // see Cache::replay_delta()
- assert(delta_type != extent_types_t::JOURNAL_TAIL);
- if (delta_type == extent_types_t::ALLOC_INFO) {
- ++stats.num_alloc_deltas;
- } else {
- ++stats.num_dirty_deltas;
- }
- }
- });
+ return handler(locator, delta, modify_time).discard_result();
});
});
});
);
}
+SegmentedJournal::replay_ret
+SegmentedJournal::scan_valid_record_delta(
+ scan_delta_handler_t &&delta_handler,
+ journal_seq_t tail)
+{
+ auto handler = std::move(delta_handler);
+ for (auto &[seq, header] : replay_segments) {
+ co_await replay_segment(seq, header, handler);
+ }
+}
+
SegmentedJournal::replay_ret SegmentedJournal::replay(
delta_handler_t &&delta_handler)
{
auto handler = std::move(delta_handler);
auto segment_headers = co_await sm_group.find_journal_segment_headers();
INFO("got {} segments", segment_headers.size());
- replay_stats_t stats;
- auto segments = co_await prep_replay_segments(std::move(segment_headers));
- for (auto &[seq, header] : segments) {
- co_await replay_segment(seq, header, handler, stats);
- }
+ co_await prep_replay_segments(std::move(segment_headers));
+ auto d_handler = [&handler, this](
+ const record_locator_t &locator,
+ const delta_info_t &delta,
+ sea_time_point modify_time) -> replay_ertr::future<bool> {
+ auto ret = co_await handler(
+ locator,
+ delta,
+ get_dirty_tail(),
+ get_alloc_tail(),
+ modify_time);
+ auto [is_applied, ext] = ret;
+ if (is_applied) {
+ // see Cache::replay_delta()
+ assert(delta.type != extent_types_t::JOURNAL_TAIL);
+ if (delta.type == extent_types_t::ALLOC_INFO) {
+ ++stats.num_alloc_deltas;
+ } else {
+ ++stats.num_dirty_deltas;
+ }
+ }
+ co_return true;
+ };
+ journal_seq_t tail = get_dirty_tail() <= get_alloc_tail() ?
+ get_dirty_tail() : get_alloc_tail();
+ co_return co_await scan_valid_record_delta(std::move(d_handler), tail);
INFO("replay done, record_groups={}, records={}, "
"alloc_deltas={}, dirty_deltas={}",
stats.num_record_groups,
WritePipeline* write_pipeline = nullptr;
/// return ordered vector of segments to replay
- using replay_segments_t = std::vector<
- std::pair<journal_seq_t, segment_header_t>>;
- using prep_replay_segments_fut = replay_ertr::future<
- replay_segments_t>;
+ using prep_replay_segments_fut = replay_ertr::future<>;
prep_replay_segments_fut prep_replay_segments(
std::vector<std::pair<segment_id_t, segment_header_t>> segments);
std::size_t num_records = 0;
std::size_t num_alloc_deltas = 0;
std::size_t num_dirty_deltas = 0;
- };
+ } stats;
+
+ using replay_segments_t = std::vector<
+ std::pair<journal_seq_t, segment_header_t>>;
+ replay_segments_t replay_segments;
/// replays records starting at start through end of segment
replay_ertr::future<>
replay_segment(
journal_seq_t start, ///< [in] starting addr, seq
segment_header_t header, ///< [in] segment header
- delta_handler_t &delta_handler, ///< [in] processes deltas in order
- replay_stats_t &stats ///< [out] replay stats
+ scan_delta_handler_t &delta_handler ///< [in] processes deltas in order
);
journal_seq_t get_dirty_tail() const final {
replay_ret scan_valid_record_delta(
scan_delta_handler_t &&delta_handler,
- journal_seq_t tail) final {
- return replay_ertr::now();
- }
+ journal_seq_t tail) final;
};
}