diff --git a/.github/scripts/bench_logs_to_pages.py b/.github/scripts/bench_logs_to_pages.py
index 235294797e..e890d1d1c6 100644
--- a/.github/scripts/bench_logs_to_pages.py
+++ b/.github/scripts/bench_logs_to_pages.py
@@ -60,6 +60,7 @@
YAML_USED_NAMES = (
"db_bench-fillrandom.yaml",
"db_bench-fillseq.yaml",
+ "db_bench-fillseq-cspp.yaml",
)
ENGINE_LABELS = {
"zipkeyonly": "ToplingDB zipkeyonly",
@@ -137,7 +138,7 @@ def format_iec(num_bytes: int) -> str:
return f"{n:.1f}{units[idx]}"
-SHM_WORKLOADS = ("fillrandom", "fillseq")
+SHM_WORKLOADS = ("fillrandom", "fillseq", "fillseq-cspp")
SHM_WORKLOAD_LABELS = SHM_SUITE_LABELS
@@ -160,7 +161,14 @@ def load_shm_usages(eng_dir: Path) -> Dict[str, Optional[Dict[str, int]]]:
return out
-RSS_WORKLOADS = ("fillrandom", "fillseq", "fillrandom-omit", "fillseq-omit")
+RSS_WORKLOADS = (
+ "fillrandom",
+ "fillseq",
+ "fillrandom-omit",
+ "fillseq-omit",
+ "fillseq-cspp",
+ "fillseq-cspp-omit",
+)
def parse_rss_usage(text: str) -> Optional[int]:
@@ -240,6 +248,10 @@ def _bytes(eng: str, wl: str, key: str) -> Optional[int]:
rows_html = []
for wl in SHM_WORKLOADS:
+ if wl == "fillseq-cspp" and all(
+ _bytes(e, wl, "allocated_bytes") is None for e in ENGINES
+ ):
+ continue
cells = [f"
{html.escape(SHM_WORKLOAD_LABELS.get(wl, wl))} | "]
for e in ENGINES:
b = _bytes(e, wl, "allocated_bytes")
@@ -633,6 +645,26 @@ def build_db_bench_compare(
LAZY_ENGINES = ("zipkeyonly", "zipkeyvalue", "rocksdb-v8.10")
+def _has_topling_fillseq_cspp(engines: Dict[str, Any]) -> bool:
+ return any(
+ bool((engines.get(e) or {}).get("db_bench_fillseq_cspp"))
+ for e in TOPLING_ENGINES
+ )
+
+
+def _db_bench_by_engine(
+ engines: Dict[str, Any],
+ *,
+ topling_key: str,
+ rocks_key: str = "db_bench",
+) -> Dict[str, List[Dict[str, str]]]:
+ out: Dict[str, List[Dict[str, str]]] = {}
+ for e in ENGINES:
+ key = topling_key if e in TOPLING_ENGINES else rocks_key
+ out[e] = (engines.get(e) or {}).get(key) or []
+ return out
+
+
def _hl(text: str, kind: str) -> str:
"""Color a short phrase: kind is 'faster' (green) or 'slower' (red)."""
return f'{html.escape(text)}'
@@ -762,11 +794,51 @@ def _cost_ratio_cell(baseline: Optional[float], subject: Optional[float]) -> str
)
_CSPP_METRICS_LOW = (
"Elapsed time",
- "write us/op",
- "read us/op",
)
+def _memtablerep_elapsed_display(
+ mmap: Dict[str, str], bench: str, elapsed_raw: str
+) -> str:
+ """Append write/read us/op onto the Elapsed time cell."""
+ write_us = mmap.get(f"{bench}|write us/op")
+ read_us = mmap.get(f"{bench}|read us/op")
+ extras: List[str] = []
+ if write_us and read_us:
+ extras.append(f"write {write_us} us/op")
+ extras.append(f"read {read_us} us/op")
+ elif write_us:
+ extras.append(f"{write_us} us/op")
+ elif read_us:
+ extras.append(f"{read_us} us/op")
+ if not extras:
+ return elapsed_raw or "—"
+ base = elapsed_raw or "—"
+ return f"{base} ({', '.join(extras)})"
+
+
+def _fold_memtablerep_usop(rows: List[Dict[str, str]]) -> List[Dict[str, str]]:
+ """Drop standalone us/op rows; fold them into Elapsed time."""
+ mmap = _metric_map(rows)
+ out: List[Dict[str, str]] = []
+ for row in rows:
+ metric = row["metric"]
+ if metric in ("write us/op", "read us/op"):
+ continue
+ if metric == "Elapsed time":
+ out.append(
+ {
+ **row,
+ "value": _memtablerep_elapsed_display(
+ mmap, row["benchmark"], row["value"]
+ ),
+ }
+ )
+ else:
+ out.append(row)
+ return out
+
+
def build_memtablerep_compare(
cspp_rows: List[Dict[str, str]],
skiplist_topling: List[Dict[str, str]],
@@ -806,6 +878,11 @@ def build_memtablerep_compare(
_metric_number(c_raw),
_metric_number(r_raw),
)
+ if metric == "Elapsed time":
+ o_raw = _memtablerep_elapsed_display(offset_skiplist, bench, o_raw)
+ c_raw = _memtablerep_elapsed_display(cspp, bench, c_raw)
+ t_raw = _memtablerep_elapsed_display(skip_t, bench, t_raw)
+ r_raw = _memtablerep_elapsed_display(skip_r, bench, r_raw)
if metric in _CSPP_METRICS_HIGH:
offset_skiplist_ratio = _throughput_ratio_cell(r_n, o_n)
cspp_ratio = _throughput_ratio_cell(r_n, c_n)
@@ -860,13 +937,22 @@ def _load_engine_logs(log_root: Path) -> Dict[str, Dict[str, Any]]:
)
omit_fr_rows: List[Dict[str, str]] = []
omit_fs_rows: List[Dict[str, str]] = []
+ omit_fs_cspp_rows: List[Dict[str, str]] = []
+ fs_cspp_path = eng_dir / "db_bench-fillseq-cspp.log"
+ fs_cspp_rows: List[Dict[str, str]] = []
+ if fs_cspp_path.is_file():
+ fs_cspp_rows = parse_db_bench(
+ fs_cspp_path.read_text(encoding="utf-8", errors="replace")
+ )
if eng == "rocksdb-v8.10":
# Reuse readseq×3 from the main fill* suites (no separate omit/scan pass).
omit_fr_rows = _readseq_rows(fr_rows)
omit_fs_rows = _readseq_rows(db_rows)
+ omit_fs_cspp_rows = omit_fs_rows
else:
omit_fr = eng_dir / "db_bench-fillrandom-omit.log"
omit_fs = eng_dir / "db_bench-fillseq-omit.log"
+ omit_fs_cspp = eng_dir / "db_bench-fillseq-cspp-omit.log"
if omit_fr.is_file():
omit_fr_rows = parse_db_bench(
omit_fr.read_text(encoding="utf-8", errors="replace")
@@ -875,6 +961,10 @@ def _load_engine_logs(log_root: Path) -> Dict[str, Dict[str, Any]]:
omit_fs_rows = parse_db_bench(
omit_fs.read_text(encoding="utf-8", errors="replace")
)
+ if omit_fs_cspp.is_file():
+ omit_fs_cspp_rows = parse_db_bench(
+ omit_fs_cspp.read_text(encoding="utf-8", errors="replace")
+ )
skiplist_rows: List[Dict[str, str]] = []
cspp_rows: List[Dict[str, str]] = []
offset_skiplist_rows: List[Dict[str, str]] = []
@@ -893,8 +983,10 @@ def _load_engine_logs(log_root: Path) -> Dict[str, Dict[str, Any]]:
result[eng] = {
"db_bench": db_rows,
"db_bench_fillrandom": fr_rows,
+ "db_bench_fillseq_cspp": fs_cspp_rows,
"db_bench_omit_fillrandom": omit_fr_rows,
"db_bench_omit_fillseq": omit_fs_rows,
+ "db_bench_omit_fillseq_cspp": omit_fs_cspp_rows,
"memtablerep_skiplist": skiplist_rows,
"memtablerep_cspp": cspp_rows,
"memtablerep_OffsetSkipList": offset_skiplist_rows,
@@ -1117,6 +1209,15 @@ def _build_per_engine_details(engines_data: Dict[str, Any]) -> str:
detail_parts.append(
_table(db_bench_detail_keys, data["db_bench"], db_bench_detail_keys)
)
+ if data.get("db_bench_fillseq_cspp"):
+ detail_parts.append("db_bench (fillseq suite, CSPP)
")
+ detail_parts.append(
+ _table(
+ db_bench_detail_keys,
+ data["db_bench_fillseq_cspp"],
+ db_bench_detail_keys,
+ )
+ )
if eng in TOPLING_ENGINES:
if data.get("db_bench_omit_fillrandom"):
detail_parts.append(
@@ -1140,12 +1241,23 @@ def _build_per_engine_details(engines_data: Dict[str, Any]) -> str:
db_bench_detail_keys,
)
)
+ if data.get("db_bench_omit_fillseq_cspp"):
+ detail_parts.append(
+ "db_bench omit lazy-load (fillseq CSPP DB)
"
+ )
+ detail_parts.append(
+ _table(
+ db_bench_detail_keys,
+ data["db_bench_omit_fillseq_cspp"],
+ db_bench_detail_keys,
+ )
+ )
if data.get("memtablerep_skiplist"):
detail_parts.append("memtablerep_bench (skiplist)
")
detail_parts.append(
_table(
["benchmark", "metric", "value"],
- data["memtablerep_skiplist"],
+ _fold_memtablerep_usop(data["memtablerep_skiplist"]),
["benchmark", "metric", "value"],
)
)
@@ -1154,7 +1266,7 @@ def _build_per_engine_details(engines_data: Dict[str, Any]) -> str:
detail_parts.append(
_table(
["benchmark", "metric", "value"],
- data["memtablerep_cspp"],
+ _fold_memtablerep_usop(data["memtablerep_cspp"]),
["benchmark", "metric", "value"],
)
)
@@ -1165,7 +1277,7 @@ def _build_per_engine_details(engines_data: Dict[str, Any]) -> str:
detail_parts.append(
_table(
["benchmark", "metric", "value"],
- data["memtablerep_OffsetSkipList"],
+ _fold_memtablerep_usop(data["memtablerep_OffsetSkipList"]),
["benchmark", "metric", "value"],
)
)
@@ -1208,22 +1320,30 @@ def emit(args: argparse.Namespace) -> None:
"db_bench-fillrandom.log",
"db_bench-fillrandom-omit.log",
"db_bench-fillseq-omit.log",
+ "db_bench-fillseq-cspp.log",
+ "db_bench-fillseq-cspp-omit.log",
"memtablerep_bench-skiplist.log",
"memtablerep_bench-cspp.log",
"memtablerep_bench-OffsetSkipList.log",
"shm_usage.txt",
"shm_usage-fillrandom.txt",
"shm_usage-fillseq.txt",
+ "shm_usage-fillseq-cspp.txt",
"rss_usage-fillrandom.txt",
"rss_usage-fillseq.txt",
"rss_usage-fillrandom-omit.txt",
"rss_usage-fillseq-omit.txt",
+ "rss_usage-fillseq-cspp.txt",
+ "rss_usage-fillseq-cspp-omit.txt",
"statm_series-fillrandom.txt",
"statm_series-fillseq.txt",
+ "statm_series-fillseq-cspp.txt",
"time-fillrandom.txt",
"time-fillseq.txt",
"time-fillrandom-omit.txt",
"time-fillseq-omit.txt",
+ "time-fillseq-cspp.txt",
+ "time-fillseq-cspp-omit.txt",
"bench_settings.txt",
*YAML_USED_NAMES,
"engine-meta.json",
@@ -1360,12 +1480,18 @@ def emit(args: argparse.Namespace) -> None:
"db_bench_fillrandom": engines_data.get(eng, {}).get(
"db_bench_fillrandom", []
),
+ "db_bench_fillseq_cspp": engines_data.get(eng, {}).get(
+ "db_bench_fillseq_cspp", []
+ ),
"db_bench_omit_fillrandom": engines_data.get(eng, {}).get(
"db_bench_omit_fillrandom", []
),
"db_bench_omit_fillseq": engines_data.get(eng, {}).get(
"db_bench_omit_fillseq", []
),
+ "db_bench_omit_fillseq_cspp": engines_data.get(eng, {}).get(
+ "db_bench_omit_fillseq_cspp", []
+ ),
"shm_usage": engines_data.get(eng, {}).get("shm_usage")
or {wl: None for wl in SHM_WORKLOADS},
"rss_usage": rss_by_eng.get(eng)
@@ -1457,12 +1583,17 @@ def _render_latest_section(
eng_rss_raw = engines.get(e, {}).get("rss_usage") or {}
rss_data[e] = {wl: v for wl, v in eng_rss_raw.items()}
if e in ROCKSDB_ENGINES:
- for src, dst in (("fillrandom", "fillrandom-omit"), ("fillseq", "fillseq-omit")):
+ for src, dst in (
+ ("fillrandom", "fillrandom-omit"),
+ ("fillseq", "fillseq-omit"),
+ ("fillseq-cspp", "fillseq-cspp-omit"),
+ ):
if rss_data[e].get(dst) is None and rss_data[e].get(src) is not None:
rss_data[e][dst] = rss_data[e][src]
if (
rss_data[e].get("fillrandom-omit") is not None
or rss_data[e].get("fillseq-omit") is not None
+ or rss_data[e].get("fillseq-cspp-omit") is not None
):
rss_derived_engines.add(e)
if pages_root is not None:
@@ -1484,6 +1615,26 @@ def _render_latest_section(
omit_fs_table = build_lazy_load_compare(
{e: engines.get(e, {}).get("db_bench_omit_fillseq") or [] for e in LAZY_ENGINES}
)
+ fs_cspp_compare = ""
+ omit_fs_cspp_block = ""
+ if _has_topling_fillseq_cspp(engines):
+ fs_cspp_compare = (
+ "Comparison: db_bench fillseq suite (CSPP) (perf)
\n"
+ 'Same as fillrandom, except ToplingDB fillseq uses CSPP. '
+ "RocksDB fillseq benefits from shortcuts: trivial_move on "
+ "non-overlapping SSTs; refit level skips zstd on L6: faster, "
+ "larger size. Seqno-zeroing compact still runs.
\n"
+ f"{build_db_bench_compare(_db_bench_by_engine(engines, topling_key='db_bench_fillseq_cspp'))}"
+ )
+ omit_fs_cspp_block = (
+ "scan-omit-value on data from fillseq (CSPP)
\n"
+ + build_lazy_load_compare(
+ {
+ e: engines.get(e, {}).get("db_bench_omit_fillseq_cspp") or []
+ for e in LAZY_ENGINES
+ }
+ )
+ )
t_eng = engines.get("zipkeyonly") or {}
r_eng = engines.get("rocksdb-v8.10") or {}
@@ -1533,12 +1684,14 @@ def _render_latest_section(
Comparison: db_bench fillseq suite (perf)
Same as fillrandom, except ToplingDB fillseq uses OffsetSkipList (fillrandom still uses CSPP). RocksDB fillseq benefits from shortcuts: trivial_move on non-overlapping SSTs; refit level skips zstd on L6: faster, larger size. Seqno-zeroing compact still runs.
{db_compare_fs}
+ {fs_cspp_compare}
Lazy load demo (scan; RocksDB v8.10 baseline)
zipkey* needs an extra omit pass: scan_omit_key/value enables lazy value load (no real value load). RocksDB has no lazy load, so the baseline is readseq×3 already present in the main fill* suite (no extra pass). RocksDB nextwithkey cells are =readseq. master omitted here (v8.10 is the stronger RocksDB baseline). {_color_sign()}.
scan-omit-value on data from fillrandom
{omit_fr_table}
scan-omit-value on data from fillseq
{omit_fs_table}
+ {omit_fs_cspp_block}
memtablerep_bench: OffsetSkipList and CSPP vs skiplist
Focus: {_hl('OffsetSkipList / CSPP (ToplingDB)', 'faster')} vs skiplist. Baseline = RocksDB v8.10 skiplist. {_color_sign()}.
{memtablerep_compare}
diff --git a/.github/scripts/bench_pages_common.py b/.github/scripts/bench_pages_common.py
index 0504eec037..6d7d74ca8d 100644
--- a/.github/scripts/bench_pages_common.py
+++ b/.github/scripts/bench_pages_common.py
@@ -76,6 +76,7 @@ def stage_window_rss_bytes(
SUITE_READRANDOM = (
("fillrandom", "db_bench_fillrandom", "fillrandom-readrandom"),
("fillseq", "db_bench", "fillseq-readrandom"),
+ ("fillseq-cspp", "db_bench_fillseq_cspp", "fillseq-cspp-readrandom"),
)
@@ -104,6 +105,7 @@ def attach_suite_readrandom_rss(
SHM_SUITE_LABELS = {
"fillrandom": "fillrandom suite",
"fillseq": "fillseq suite",
+ "fillseq-cspp": "fillseq suite (CSPP)",
}
RSS_WORKLOAD_ORDER = (
"fillrandom",
@@ -112,6 +114,9 @@ def attach_suite_readrandom_rss(
"fillseq",
"fillseq-readrandom",
"fillseq-omit",
+ "fillseq-cspp",
+ "fillseq-cspp-readrandom",
+ "fillseq-cspp-omit",
)
RSS_WORKLOAD_LABELS = {
"fillrandom": "fillrandom suite peak",
@@ -120,6 +125,9 @@ def attach_suite_readrandom_rss(
"fillseq-readrandom": "fillseq suite readrandom",
"fillrandom-omit": "fillrandom scan-omit-value",
"fillseq-omit": "fillseq scan-omit-value",
+ "fillseq-cspp": "fillseq suite (CSPP) peak",
+ "fillseq-cspp-readrandom": "fillseq suite (CSPP) readrandom",
+ "fillseq-cspp-omit": "fillseq scan-omit-value (CSPP)",
}
RSS_WORKLOAD_TIPS = {
"fillrandom-readrandom": (
@@ -128,6 +136,9 @@ def attach_suite_readrandom_rss(
"fillseq-readrandom": (
"peak RSS during the readrandom stage of the fillseq suite"
),
+ "fillseq-cspp-readrandom": (
+ "peak RSS during the readrandom stage of the fillseq suite (CSPP)"
+ ),
"fillrandom-omit": (
"restart process with reuse db data of fillrandom, "
"scan without access value, benefited by lazy load value (ToplingDB feature)"
@@ -136,6 +147,10 @@ def attach_suite_readrandom_rss(
"restart process with reuse db data of fillseq, "
"scan without access value, benefited by lazy load value (ToplingDB feature)"
),
+ "fillseq-cspp-omit": (
+ "restart process with reuse db data of fillseq (CSPP), "
+ "scan without access value, benefited by lazy load value (ToplingDB feature)"
+ ),
}
@@ -445,6 +460,8 @@ def combine_db_bench_logs(engine_raw: Path) -> None:
"db_bench-fillrandom-omit.log",
"db_bench.log",
"db_bench-fillseq-omit.log",
+ "db_bench-fillseq-cspp.log",
+ "db_bench-fillseq-cspp-omit.log",
)
chunks = [
(engine_raw / name).read_bytes().rstrip(b"\n")
@@ -504,6 +521,7 @@ def build_rss_svg_section(
for suite, bench_key in [
("fillrandom", "db_bench_fillrandom"),
("fillseq", "db_bench"),
+ ("fillseq-cspp", "db_bench_fillseq_cspp"),
]:
series_path = eng_dir / f"statm_series-{suite}.txt"
if not series_path.is_file():
diff --git a/.github/scripts/test_bench_rss_series.py b/.github/scripts/test_bench_rss_series.py
index 88b3275ac2..c8188705ea 100755
--- a/.github/scripts/test_bench_rss_series.py
+++ b/.github/scripts/test_bench_rss_series.py
@@ -301,6 +301,97 @@ def check_pages_contract(mod, variant: str) -> None:
assert "dcompact bench →" in home
assert "offloads most CPU and memory cost" in home
assert "not RocksDB CompactionService" not in home
+ assert "Comparison: db_bench fillseq suite (CSPP)" not in home
+ assert "fillseq suite (CSPP)" not in home
+ assert "db_bench (fillseq suite, CSPP)" not in result_html
+
+
+def check_fillseq_cspp_pages(mod) -> None:
+ """ToplingDB CSPP fillseq is a full twin of the OffsetSkipList fillseq suite."""
+ with tempfile.TemporaryDirectory() as tmp:
+ tmp_path = Path(tmp)
+ log_root = tmp_path / "logs"
+ emit_out = tmp_path / "emit"
+ site = tmp_path / "site"
+ _write_min_logs(log_root)
+ bench_body = (
+ _DB_BENCH_LINE
+ + "readrandom : 1.0 micros/op 1000 ops/sec 1.0 seconds "
+ "1000 operations; x\n"
+ )
+ for eng in ("zipkeyonly", "zipkeyvalue"):
+ eng_dir = log_root / eng
+ (eng_dir / "db_bench-fillseq-omit.log").write_text(
+ "$ fillseq-omit\n"
+ "nextwithkey : 1.0 micros/op 1000 ops/sec 1.0 seconds "
+ "1000 operations; x\n",
+ encoding="utf-8",
+ )
+ (eng_dir / "db_bench-fillseq-cspp.log").write_text(
+ "$ fillseq-cspp\n" + bench_body, encoding="utf-8"
+ )
+ (eng_dir / "db_bench-fillseq-cspp-omit.log").write_text(
+ "$ fillseq-cspp-omit\n"
+ "nextwithkey : 1.0 micros/op 1000 ops/sec 1.0 seconds "
+ "1000 operations; x\n",
+ encoding="utf-8",
+ )
+ (eng_dir / "statm_series-fillseq-cspp.txt").write_text(
+ _STATM_SERIES, encoding="utf-8"
+ )
+ (eng_dir / "shm_usage-fillseq-cspp.txt").write_text(
+ "apparent_bytes=1000\nallocated_bytes=2000\n", encoding="utf-8"
+ )
+ (eng_dir / "rss_usage-fillseq-cspp.txt").write_text(
+ "max_rss_bytes=4096\n", encoding="utf-8"
+ )
+ (eng_dir / "rss_usage-fillseq-cspp-omit.txt").write_text(
+ "max_rss_bytes=2048\n", encoding="utf-8"
+ )
+ emit_args = argparse.Namespace(
+ variant="plain",
+ run_id="cspp-fillseq",
+ log_root=str(log_root),
+ engine_meta_root=None,
+ actions_run_url="",
+ out=str(emit_out),
+ )
+ mod.emit(emit_args)
+ run_dirs = list((emit_out / "runs").iterdir())
+ assert len(run_dirs) == 1, run_dirs
+ result_html = (run_dirs[0] / "index.html").read_text(encoding="utf-8")
+ assert "db_bench (fillseq suite)" in result_html
+ assert "db_bench (fillseq suite, CSPP)" in result_html
+ assert "db_bench omit lazy-load (fillseq DB)" in result_html
+ assert "db_bench omit lazy-load (fillseq CSPP DB)" in result_html
+ combined = (
+ run_dirs[0] / "raw" / "zipkeyonly" / "db_bench-all.log"
+ ).read_text(encoding="utf-8")
+ assert combined.index("$ fillseq-omit\n") < combined.index("$ fillseq-cspp\n")
+ assert combined.index("$ fillseq-cspp\n") < combined.index(
+ "$ fillseq-cspp-omit\n"
+ )
+ meta = json.loads((emit_out / "run-meta.json").read_text(encoding="utf-8"))
+ zko_rss = meta["engines"]["zipkeyonly"]["rss_usage"]
+ assert zko_rss.get("fillseq-readrandom") is not None
+ assert zko_rss.get("fillseq-cspp-readrandom") is not None
+ assert zko_rss.get("fillseq-cspp") == 4096
+ assert zko_rss.get("fillseq-cspp-omit") == 2048
+ mod.merge(
+ argparse.Namespace(
+ merge_into=str(site),
+ from_dir=str(emit_out),
+ variant="plain",
+ )
+ )
+ home = (site / "index.html").read_text(encoding="utf-8")
+ assert "Comparison: db_bench fillseq suite (perf)" in home
+ assert "Comparison: db_bench fillseq suite (CSPP) (perf)" in home
+ assert "scan-omit-value on data from fillseq" in home
+ assert "scan-omit-value on data from fillseq (CSPP)" in home
+ assert "fillseq suite (CSPP) peak" in home
+ assert "fillseq suite (CSPP) readrandom" in home
+ assert "fillseq scan-omit-value (CSPP)" in home
def check_dcompact_home_nav(mod) -> None:
@@ -596,6 +687,11 @@ def main() -> int:
("db_bench-fillrandom-omit.log", "$ fillrandom-omit\nomit output\n"),
("db_bench.log", "$ fillseq\nfillseq output\n"),
("db_bench-fillseq-omit.log", "$ fillseq-omit\nomit output\n"),
+ ("db_bench-fillseq-cspp.log", "$ fillseq-cspp\ncspp output\n"),
+ (
+ "db_bench-fillseq-cspp-omit.log",
+ "$ fillseq-cspp-omit\ncspp omit output\n",
+ ),
)
for name, content in source_logs:
(eng_raw / name).write_text(content, encoding="utf-8")
@@ -692,6 +788,7 @@ def main() -> int:
assert "TestOS" in html
check_readrandom_highlight(mod)
check_suite_readrandom_peak(mod)
+ check_fillseq_cspp_pages(mod)
if name == "bench_dcompact_pages":
check_dcompact_home_nav(mod)
check_dcompact_rss_row_tips(mod)
diff --git a/.github/workflows/db_bench-avx512-run.yml b/.github/workflows/db_bench-avx512-run.yml
index a11469afd1..9bc697ca6f 100644
--- a/.github/workflows/db_bench-avx512-run.yml
+++ b/.github/workflows/db_bench-avx512-run.yml
@@ -206,50 +206,59 @@ jobs:
record_shm fillrandom
rm -rf "$DB_PATH"
- # Pass 2: fillseq — OffsetSkipList; prefix 6 zipkeyonly (keep L6).
- prepare_db
- yaml_fs="${logdir}/db_bench-fillseq.yaml"
- python3 .github/scripts/graft_bench_yaml.py \
- --prefix-level-writers 6 zipkeyonly \
- --target-file-size-base 128M \
- --target-file-size-multiplier 1 \
- --memtable-factory '"${offset_skiplist}"' \
- --out "$yaml_fs" \
- "$yaml"
- args=(
- -json "$yaml_fs"
- -num=100000000
- -key_size=8
- -value_size="${VALUE_SIZE}"
- -batch_size=1000
- -benchmarks=fillseq,flush,compact,readseq,readseq,readseq,readrandom
- -enable_zero_copy
- -progress_reports=false
- -compact_target_level=6
- )
- echo '$' "$TOPLING/bin/db_bench" "${args[@]}" >"${logdir}/db_bench.log"
- /usr/bin/time -f 'max_rss_kb=%M' -o "${logdir}/time-fillseq.txt" -- \
- "$TOPLING/bin/db_bench" "${args[@]}" >>"${logdir}/db_bench.log" 2>&1
- cat "${logdir}/db_bench.log"
- args_omit_fs=(
- -json "$yaml_fs"
- -num=100000000
- -key_size=8
- -value_size="${VALUE_SIZE}"
- -batch_size=1000
- -benchmarks=nextwithkey,nextwithkey,nextwithkey,readseq,readseq,readseq
- -scan_omit_key -scan_omit_value
- -use_existing_db=1
- -enable_zero_copy
- -progress_reports=false
- )
- echo '$' "$TOPLING/bin/db_bench" "${args_omit_fs[@]}" >"${logdir}/db_bench-fillseq-omit.log"
- /usr/bin/time -f 'max_rss_kb=%M' -o "${logdir}/time-fillseq-omit.txt" -- \
- "$TOPLING/bin/db_bench" "${args_omit_fs[@]}" >>"${logdir}/db_bench-fillseq-omit.log" 2>&1
- cat "${logdir}/db_bench-fillseq-omit.log"
- record_rss fillseq
- record_rss fillseq-omit
- record_shm fillseq
+ # Pass 2: fillseq — OffsetSkipList then CSPP; prefix 6 zipkeyonly (keep L6).
+ run_topling_fillseq() {
+ local factory="$1"
+ local tag="$2"
+ local yaml_fs="${logdir}/db_bench-${tag}.yaml"
+ local main_log="${logdir}/db_bench.log"
+ [ "$tag" = "fillseq" ] || main_log="${logdir}/db_bench-${tag}.log"
+ local omit_log="${logdir}/db_bench-${tag}-omit.log"
+ prepare_db
+ python3 .github/scripts/graft_bench_yaml.py \
+ --prefix-level-writers 6 zipkeyonly \
+ --target-file-size-base 128M \
+ --target-file-size-multiplier 1 \
+ --memtable-factory '"${'"${factory}"'}"' \
+ --out "$yaml_fs" \
+ "$yaml"
+ args=(
+ -json "$yaml_fs"
+ -num=100000000
+ -key_size=8
+ -value_size="${VALUE_SIZE}"
+ -batch_size=1000
+ -benchmarks=fillseq,flush,compact,readseq,readseq,readseq,readrandom
+ -enable_zero_copy
+ -progress_reports=false
+ -compact_target_level=6
+ )
+ echo '$' "$TOPLING/bin/db_bench" "${args[@]}" >"$main_log"
+ /usr/bin/time -f 'max_rss_kb=%M' -o "${logdir}/time-${tag}.txt" -- \
+ "$TOPLING/bin/db_bench" "${args[@]}" >>"$main_log" 2>&1
+ cat "$main_log"
+ args_omit_fs=(
+ -json "$yaml_fs"
+ -num=100000000
+ -key_size=8
+ -value_size="${VALUE_SIZE}"
+ -batch_size=1000
+ -benchmarks=nextwithkey,nextwithkey,nextwithkey,readseq,readseq,readseq
+ -scan_omit_key -scan_omit_value
+ -use_existing_db=1
+ -enable_zero_copy
+ -progress_reports=false
+ )
+ echo '$' "$TOPLING/bin/db_bench" "${args_omit_fs[@]}" >"$omit_log"
+ /usr/bin/time -f 'max_rss_kb=%M' -o "${logdir}/time-${tag}-omit.txt" -- \
+ "$TOPLING/bin/db_bench" "${args_omit_fs[@]}" >>"$omit_log" 2>&1
+ cat "$omit_log"
+ record_rss "$tag"
+ record_rss "${tag}-omit"
+ record_shm "$tag"
+ }
+ run_topling_fillseq offset_skiplist fillseq
+ run_topling_fillseq cspp fillseq-cspp
if [ "$run_memtable" = "1" ]; then
mt=(
-benchmarks=fillrandom,readrandom
diff --git a/.github/workflows/db_bench-run.yml b/.github/workflows/db_bench-run.yml
index ccc59c0430..fdc6c1cd45 100644
--- a/.github/workflows/db_bench-run.yml
+++ b/.github/workflows/db_bench-run.yml
@@ -255,54 +255,64 @@ jobs:
record_shm fillrandom
rm -rf "$DB_PATH"
- # Pass 2: fillseq — OffsetSkipList; prefix 6 zipkeyonly (keep L6).
- prepare_db
- yaml_fs="${logdir}/db_bench-fillseq.yaml"
- python3 .github/scripts/graft_bench_yaml.py \
- --prefix-level-writers 6 zipkeyonly \
- --target-file-size-base 128M \
- --target-file-size-multiplier 1 \
- --memtable-factory '"${offset_skiplist}"' \
- --out "$yaml_fs" \
- "$yaml"
- args=(
- -json "$yaml_fs"
- -num="${NUM}"
- -key_size=8
- -value_size="${VALUE_SIZE}"
- -batch_size=1000
- -benchmarks=fillseq,flush,compact,readseq,readseq,readseq,readrandom
- -enable_zero_copy
- -progress_reports=false
- -report_bench_start_time
- -compact_target_level=6
- )
- echo '$' "$TOPLING/bin/db_bench" "${args[@]}" >"${logdir}/db_bench.log"
- .github/scripts/run_sample_statm_fdcache.sh "${logdir}/statm_series-fillseq.txt" "${logdir}/time-fillseq.txt" \
- "$TOPLING/bin/db_bench" "${args[@]}" \
- >>"${logdir}/db_bench.log" 2>&1
- cat "${logdir}/db_bench.log"
- save_db_log fillseq
- args_omit_fs=(
- -json "$yaml_fs"
- -num="${NUM}"
- -key_size=8
- -value_size="${VALUE_SIZE}"
- -batch_size=1000
- -benchmarks=nextwithkey,nextwithkey,nextwithkey,readseq,readseq,readseq
- -scan_omit_key -scan_omit_value
- -use_existing_db=1
- -enable_zero_copy
- -progress_reports=false
- )
- echo '$' "$TOPLING/bin/db_bench" "${args_omit_fs[@]}" >"${logdir}/db_bench-fillseq-omit.log"
- /usr/bin/time -f 'max_rss_kb=%M' -o "${logdir}/time-fillseq-omit.txt" -- \
- "$TOPLING/bin/db_bench" "${args_omit_fs[@]}" >>"${logdir}/db_bench-fillseq-omit.log" 2>&1
- cat "${logdir}/db_bench-fillseq-omit.log"
- save_db_log fillseq-omit
- record_rss fillseq
- record_rss fillseq-omit
- record_shm fillseq
+ # Pass 2: fillseq — OffsetSkipList then CSPP; prefix 6 zipkeyonly (keep L6).
+ run_topling_fillseq() {
+ local factory="$1"
+ local tag="$2"
+ local yaml_fs="${logdir}/db_bench-${tag}.yaml"
+ local main_log="${logdir}/db_bench.log"
+ [ "$tag" = "fillseq" ] || main_log="${logdir}/db_bench-${tag}.log"
+ local omit_log="${logdir}/db_bench-${tag}-omit.log"
+ prepare_db
+ python3 .github/scripts/graft_bench_yaml.py \
+ --prefix-level-writers 6 zipkeyonly \
+ --target-file-size-base 128M \
+ --target-file-size-multiplier 1 \
+ --memtable-factory '"${'"${factory}"'}"' \
+ --out "$yaml_fs" \
+ "$yaml"
+ args=(
+ -json "$yaml_fs"
+ -num="${NUM}"
+ -key_size=8
+ -value_size="${VALUE_SIZE}"
+ -batch_size=1000
+ -benchmarks=fillseq,flush,compact,readseq,readseq,readseq,readrandom
+ -enable_zero_copy
+ -progress_reports=false
+ -report_bench_start_time
+ -compact_target_level=6
+ )
+ echo '$' "$TOPLING/bin/db_bench" "${args[@]}" >"$main_log"
+ .github/scripts/run_sample_statm_fdcache.sh \
+ "${logdir}/statm_series-${tag}.txt" "${logdir}/time-${tag}.txt" \
+ "$TOPLING/bin/db_bench" "${args[@]}" \
+ >>"$main_log" 2>&1
+ cat "$main_log"
+ save_db_log "$tag"
+ args_omit_fs=(
+ -json "$yaml_fs"
+ -num="${NUM}"
+ -key_size=8
+ -value_size="${VALUE_SIZE}"
+ -batch_size=1000
+ -benchmarks=nextwithkey,nextwithkey,nextwithkey,readseq,readseq,readseq
+ -scan_omit_key -scan_omit_value
+ -use_existing_db=1
+ -enable_zero_copy
+ -progress_reports=false
+ )
+ echo '$' "$TOPLING/bin/db_bench" "${args_omit_fs[@]}" >"$omit_log"
+ /usr/bin/time -f 'max_rss_kb=%M' -o "${logdir}/time-${tag}-omit.txt" -- \
+ "$TOPLING/bin/db_bench" "${args_omit_fs[@]}" >>"$omit_log" 2>&1
+ cat "$omit_log"
+ save_db_log "${tag}-omit"
+ record_rss "$tag"
+ record_rss "${tag}-omit"
+ record_shm "$tag"
+ }
+ run_topling_fillseq offset_skiplist fillseq
+ run_topling_fillseq cspp fillseq-cspp
if [ "$run_memtable" = "1" ]; then
mt=(
-benchmarks=fillrandom,readrandom
diff --git a/AGENTS.md b/AGENTS.md
new file mode 100644
index 0000000000..398bb90949
--- /dev/null
+++ b/AGENTS.md
@@ -0,0 +1,5 @@
+# Formatting
+
+- Keep each single statement on one line when the complete line, including indentation, fits within 100 columns. This includes calls, declarations, assignments, returns, and `if` / `while` / `for` headers; do not collapse block bodies.
+- If that line would exceed 100 columns, wrap it to 80 columns per line, including indentation. Preserve the surrounding continuation-indent style.
+- Format only uncommitted added or modified code. Preserve untouched existing code; do not run whole-file formatting.
diff --git a/Makefile b/Makefile
index 20649fb0fe..9ab6500523 100644
--- a/Makefile
+++ b/Makefile
@@ -148,6 +148,13 @@ endif
include make_config.mk
PLATFORM_CCFLAGS := $(filter-out -fno-builtin-memcmp, ${PLATFORM_CCFLAGS})
PLATFORM_CXXFLAGS := $(filter-out -fno-builtin-memcmp, ${PLATFORM_CXXFLAGS})
+# Nothing in this library is meant to be interposed -- it is the implementation,
+# not a layer someone else overrides. Under default interposition semantics a
+# call to a same-DSO definition still has to go through the PLT anyway, and
+# under -flto it also keeps the inliner from using the whole-program view. The
+# flag has to be on the compile line; adding it to LDFLAGS alone does nothing
+# (measured).
+PLATFORM_CXXFLAGS += -fno-semantic-interposition
# defined in make_config.mk
ROCKSDB_FULL_VERSION := ${ROCKSDB_MAJOR}.${ROCKSDB_MINOR}.${ROCKSDB_PATCH}
@@ -244,6 +251,11 @@ OPTION_lto := lto-0
ifeq ($(USE_LTO), 1)
ifeq (${DEBUG_LEVEL},0)
CXXFLAGS += -flto
+ # Lets headers hand the LTO inliner an explicit mandate where it is
+ # wanted (see table/get_context.h): GCC defines no macro of its own
+ # for -flto, and `always_inline` on a cross-TU body is a hard error
+ # without it.
+ CXXFLAGS += -DTOPLINGDB_HAVE_LTO
LDFLAGS += -flto=auto -fuse-linker-plugin
OPTION_lto := lto-$(if $(filter 1,${USE_LTO}),1,0)
endif
@@ -451,6 +463,9 @@ ifndef WITH_TOPLING_ROCKS
# default 1
WITH_TOPLING_ROCKS := 1
endif
+ifeq ($(filter 0,${DEBUG_LEVEL})$(wildcard sideplugin/topling-rocks/src/table/top_patent_algo.cc),)
+ override WITH_TOPLING_ROCKS := 0
+endif
ifeq (${WITH_TOPLING_ROCKS},1)
ifneq (,$(wildcard sideplugin/topling-rocks))
@@ -493,6 +508,9 @@ endif
# allow override by env or cmd line
WITH_CSPP_MEMTABLE ?= 1
+ifeq ($(filter 0,${DEBUG_LEVEL})$(wildcard sideplugin/cspp-memtable/cspp_memtable.cc),)
+ override WITH_CSPP_MEMTABLE := 0
+endif
ifeq (${WITH_CSPP_MEMTABLE}${WITH_TOPLING_ROCKS},10)
$(error "When WITH_CSPP_MEMTABLE is 1, WITH_TOPLING_ROCKS must be 1 also")
@@ -1883,6 +1901,9 @@ db_bench_rls: $(OBJ_DIR)/tools/db_bench.o $(BENCH_OBJECTS) $(TESTUTIL) $(LIBRARY
$(AM_LINK)
endif
+crash_recover_bench: $(OBJ_DIR)/tools/crash_recover_bench.o $(LIBRARY)
+ $(AM_LINK)
+
trace_analyzer: $(OBJ_DIR)/tools/trace_analyzer.o $(ANALYZE_OBJECTS) $(TOOLS_LIBRARY) $(LIBRARY)
$(AM_LINK)
@@ -2068,6 +2089,9 @@ db_dynamic_level_test: $(OBJ_DIR)/db/db_dynamic_level_test.o $(TEST_LIBRARY) $(L
db_flush_test: $(OBJ_DIR)/db/db_flush_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_LINK)
+db_memtable_convert_test: $(OBJ_DIR)/db/db_memtable_convert_test.o $(TEST_LIBRARY) $(LIBRARY)
+ $(AM_LINK)
+
db_inplace_update_test: $(OBJ_DIR)/db/db_inplace_update_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_LINK)
@@ -2131,6 +2155,9 @@ db_universal_compaction_test: $(OBJ_DIR)/db/db_universal_compaction_test.o $(TES
db_wal_test: $(OBJ_DIR)/db/db_wal_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_LINK)
+db_cspp_crash_safe_test: $(OBJ_DIR)/db/db_cspp_crash_safe_test.o $(TEST_LIBRARY) $(LIBRARY)
+ $(AM_LINK)
+
db_io_failure_test: $(OBJ_DIR)/db/db_io_failure_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_LINK)
diff --git a/README-zh_cn.md b/README-zh_cn.md
index fc17d6b5a4..699f6a1305 100644
--- a/README-zh_cn.md
+++ b/README-zh_cn.md
@@ -37,6 +37,11 @@ ToplingDB 兼容 RocksDB API 的同时,增加了很多非常重要的功能与
1. 内置 Prometheus 指标的支持,这是在[内嵌 Http](https://github.com/topling/rockside/wiki/WebView) 中实现的
1. 修复了很多 RocksDB 的 bug,我们已将其中易于合并到 RocksDB 的很多修复与改进给上游 RocksDB 发了 [Pull Request](https://github.com/facebook/rocksdb/pulls?q=is%3Apr+author%3Arockeet)
+## 进程崩溃后的恢复
+恢复机制、配置方式及适用边界见 [MemTable Crash-Safe Recovery](https://github.com/topling/rockside/wiki/Crash-Safe-Recovery)。
+底层数据结构为何同时支持读侧无等待与进程崩溃后的读取,见 [读侧无等待与 Crash-Safe 的同构性](https://github.com/topling/rockside/wiki/Wait-Free-Reads-and-Crash-Safe)。
+异常退出后 `DB::Open` 的耗时见 [crash_recover_bench.md](tools/crash_recover_bench.md)。
+
## ToplingDB 云原生数据库服务
1. [MyTopling](https://github.com/topling/mytopling)(MySQL on ToplingDB), [阿里云上的 MyTopling](https://market.aliyun.com/products?k=mytopling)
1. [Todis](https://github.com/topling/todis)(Redis on ToplingDB)
diff --git a/README.md b/README.md
index 639040abbb..ea287b94cf 100644
--- a/README.md
+++ b/README.md
@@ -39,6 +39,11 @@ ToplingDB has much more key features than RocksDB:
1. Builtin Prometheus metrics support, this is based on [Embedded Http Server](https://github.com/topling/sideplugin-wiki-en/wiki/WebView)
1. Many bugfixes for RocksDB, a small part of such fixes was [Pull Requested](https://github.com/facebook/rocksdb/pulls?q=is%3Apr+author%3Arockeet) to [upstream RocksDB](https://github.com/facebook/rocksdb)
+## Crash-safe recovery
+See the [crash-safe recovery guide](https://github.com/topling/sideplugin-wiki-en/wiki/Crash-Safe-Recovery) for the recovery mechanism, configuration, and limitations.
+For the underlying data-structure principles, see [The Isomorphism Between Wait-Free Reads and Crash Safety](https://github.com/topling/sideplugin-wiki-en/wiki/Wait-Free-Reads-and-Crash-Safe).
+`DB::Open` after an abnormal exit is timed in [crash_recover_bench.md](tools/crash_recover_bench.md).
+
## ToplingDB cloud native DB services
1. [MyTopling](https://github.com/topling/mytopling)(MySQL on ToplingDB), [MyTopling on aliyun](https://market.aliyun.com/products?k=mytopling)
1. [Todis](https://github.com/topling/todis)(Redis on ToplingDB)
diff --git a/db/column_family.cc b/db/column_family.cc
index 2ea7cb8579..f69a82d778 100644
--- a/db/column_family.cc
+++ b/db/column_family.cc
@@ -40,6 +40,7 @@
#include "util/autovector.h"
#include "util/cast_util.h"
#include "util/compression.h"
+#include
namespace ROCKSDB_NAMESPACE {
@@ -468,6 +469,15 @@ ColumnFamilyOptions SanitizeOptions(const ImmutableDBOptions& db_options,
}
#endif
+ if (result.min_write_buffer_number_to_merge > 1 &&
+ result.memtable_factory->SupportConvertToSST()) {
+ ROCKS_LOG_WARN(db_options.logger,
+ "ConvertToSST converts each memtable separately; "
+ "min_write_buffer_number_to_merge > 1 is incompatible "
+ "and is sanitized to 1");
+ result.min_write_buffer_number_to_merge = 1;
+ }
+
return result;
}
@@ -1164,7 +1174,12 @@ uint64_t ColumnFamilyData::GetLiveSstFilesSize() const {
void ColumnFamilyData::PrepareNewMemtableInBackground(
const MutableCFOptions& mutable_cf_options) {
- #if !defined(ROCKSDB_UNIT_TEST)
+ bool use_cache = true;
+ TEST_SYNC_POINT_CALLBACK(
+ "ColumnFamilyData::PrepareNewMemtableInBackground:UseCache", &use_cache);
+ if (!use_cache) {
+ return;
+ }
{
std::lock_guard lk(precreated_memtable_mutex_);
if (precreated_memtable_list_.full()) {
@@ -1172,11 +1187,13 @@ void ColumnFamilyData::PrepareNewMemtableInBackground(
return;
}
}
- auto beg = ioptions_.clock->NowNanos();
+ auto beg = terark::qtime::now();
+ uint64_t number = ioptions_.memtable_factory->SupportCrashSafe()
+ ? dummy_versions_->version_set()->NewFileNumber() : 0;
auto tab = new MemTable(internal_comparator_, ioptions_, mutable_cf_options,
- write_buffer_manager_, 0/*earliest_seq*/, id_);
- auto end = ioptions_.clock->NowNanos();
- RecordInHistogram(ioptions_.stats, MEMTAB_CONSTRUCT_NANOS, end - beg);
+ write_buffer_manager_, 0/*earliest_seq*/, id_, number);
+ auto end = terark::qtime::now();
+ RecordInHistogram(ioptions_.stats, MEMTAB_CONSTRUCT_NANOS, (end - beg).ns());
{
std::lock_guard lk(precreated_memtable_mutex_);
if (LIKELY(!precreated_memtable_list_.full())) {
@@ -1191,38 +1208,100 @@ void ColumnFamilyData::PrepareNewMemtableInBackground(
"precreated_memtable_list_ is full, discard the newly created memtab");
delete tab;
}
- #endif
}
MemTable* ColumnFamilyData::ConstructNewMemtable(
const MutableCFOptions& mutable_cf_options, SequenceNumber earliest_seq) {
MemTable* tab = nullptr;
- #if !defined(ROCKSDB_UNIT_TEST)
- {
+ bool use_cache = true;
+ TEST_SYNC_POINT_CALLBACK("ColumnFamilyData::ConstructNewMemtable:UseCache",
+ &use_cache);
+ if (use_cache) {
std::lock_guard lk(precreated_memtable_mutex_);
if (!precreated_memtable_list_.empty()) {
tab = precreated_memtable_list_.front().release();
precreated_memtable_list_.pop_front();
}
}
- #endif
if (tab) {
tab->SetCreationSeq(earliest_seq);
tab->SetEarliestSequenceNumber(earliest_seq);
} else {
- #if !defined(ROCKSDB_UNIT_TEST)
- auto beg = ioptions_.clock->NowNanos();
- #endif
+ auto beg = terark::qtime::now();
+ // dummy_versions_ remains alive for the lifetime of this CF, unlike current_.
+ uint64_t number = ioptions_.memtable_factory->SupportCrashSafe()
+ ? dummy_versions_->version_set()->NewFileNumber() : 0;
tab = new MemTable(internal_comparator_, ioptions_, mutable_cf_options,
- write_buffer_manager_, earliest_seq, id_);
- #if !defined(ROCKSDB_UNIT_TEST)
- auto end = ioptions_.clock->NowNanos();
- RecordInHistogram(ioptions_.stats, MEMTAB_CONSTRUCT_NANOS, end - beg);
- #endif
+ write_buffer_manager_, earliest_seq, id_, number);
+ auto end = terark::qtime::now();
+ RecordInHistogram(ioptions_.stats, MEMTAB_CONSTRUCT_NANOS, (end - beg).ns());
}
return tab;
}
+MemTable* ColumnFamilyData::PeekPrecreatedMemtable() {
+ std::lock_guard lk(precreated_memtable_mutex_);
+ return precreated_memtable_list_.empty()
+ ? nullptr : precreated_memtable_list_.front().get();
+}
+
+void ColumnFamilyData::ApplyMemTableFileEdit(const VersionEdit& edit) {
+ ROCKSDB_ASSERT_EQ(edit.GetColumnFamily(), id_);
+ if (edit.IsColumnFamilyDrop()) {
+ memtable_files_.clear();
+ return;
+ }
+ for (uint64_t number : edit.GetMemTableFileDeletions()) {
+ memtable_files_.erase(number);
+ }
+ for (uint64_t number : edit.GetMemTableFileAdditions()) {
+ memtable_files_.insert(number);
+ }
+}
+
+void ColumnFamilyData::AddMemTableFileEdits(VersionEdit* edit) {
+ if (!ioptions_.memtable_crash_safe_recover) return;
+ std::lock_guard lk(precreated_memtable_mutex_);
+ auto add = [&](MemTable* mem) {
+ if (!mem->IsFileRegistered()) {
+ edit->AddMemTableFile(mem->GetFileNumber());
+ edit->SetMemTableFileTracking();
+ }
+ };
+ add(mem_);
+ for (size_t i = 0; i < precreated_memtable_list_.size(); ++i) {
+ add(precreated_memtable_list_[i].get());
+ }
+}
+
+void ColumnFamilyData::PublishRegisteredMemTables() {
+ if (!ioptions_.memtable_crash_safe_recover) return;
+ TEST_SYNC_POINT("FlushJob::MemTableCache:BeforePublish");
+ {
+ std::lock_guard lk(precreated_memtable_mutex_);
+ auto publish = [&](MemTable* mem) {
+ if (memtable_files_.count(mem->GetFileNumber())) {
+ mem->MarkFileRegistered();
+ }
+ };
+ publish(mem_);
+ for (size_t i = 0; i < precreated_memtable_list_.size(); ++i) {
+ publish(precreated_memtable_list_[i].get());
+ }
+ }
+ TEST_SYNC_POINT("FlushJob::MemTableCache:AfterPublish");
+}
+
+void ColumnFamilyData::AddMemTableFileNumbers(std::vector* live) {
+ if (!ioptions_.memtable_factory->SupportCrashSafe()) return;
+ if (mem_ != nullptr) live->push_back(mem_->GetFileNumber());
+ imm_.AddMemTableFileNumbers(live);
+ std::lock_guard lk(precreated_memtable_mutex_);
+ for (size_t i = 0; i < precreated_memtable_list_.size(); ++i) {
+ live->push_back(precreated_memtable_list_[i]->GetFileNumber());
+ }
+}
+
void ColumnFamilyData::CreateNewMemtable(
const MutableCFOptions& mutable_cf_options, SequenceNumber earliest_seq) {
if (mem_ != nullptr) {
@@ -1473,6 +1552,12 @@ void ColumnFamilyData::ResetThreadLocalSuperVersions() {
Status ColumnFamilyData::ValidateOptions(
const DBOptions& db_options, const ColumnFamilyOptions& cf_options) {
+ if (db_options.memtable_crash_safe_recover &&
+ !cf_options.memtable_factory->SupportCrashSafe()) {
+ return Status::InvalidArgument(
+ "memtable_crash_safe_recover requires a FileMmap memtable factory",
+ cf_options.memtable_factory->Name());
+ }
Status s;
s = CheckCompressionSupported(cf_options);
if (s.ok() && db_options.allow_concurrent_memtable_write) {
diff --git a/db/column_family.h b/db/column_family.h
index 367b94160e..47b5229aaa 100644
--- a/db/column_family.h
+++ b/db/column_family.h
@@ -373,6 +373,13 @@ class ColumnFamilyData {
uint64_t OldestLogToKeep();
void PrepareNewMemtableInBackground(const MutableCFOptions&);
+ MemTable* PeekPrecreatedMemtable();
+ // DB mutex must be held for registration and live-file collection.
+ const std::set& GetMemTableFiles() const { return memtable_files_; }
+ void ApplyMemTableFileEdit(const VersionEdit& edit);
+ void AddMemTableFileEdits(VersionEdit* edit);
+ void AddMemTableFileNumbers(std::vector* live);
+ void PublishRegisteredMemTables();
// See Memtable constructor for explanation of earliest_seq param.
MemTable* ConstructNewMemtable(const MutableCFOptions& mutable_cf_options,
@@ -612,14 +619,13 @@ class ColumnFamilyData {
WriteBufferManager* write_buffer_manager_;
- #if !defined(ROCKSDB_UNIT_TEST)
// precreated_memtable_list_.size() is normally 1
terark::fixed_circular_queue, 4> precreated_memtable_list_;
std::mutex precreated_memtable_mutex_;
- #endif
MemTable* mem_;
MemTableList imm_;
+ std::set memtable_files_; // Protected by DB mutex.
SuperVersion* super_version_;
// An ordinal representing the current SuperVersion. Updated by
diff --git a/db/compaction/compaction_iterator.cc b/db/compaction/compaction_iterator.cc
index aeaa7d3f39..66c80d4344 100644
--- a/db/compaction/compaction_iterator.cc
+++ b/db/compaction/compaction_iterator.cc
@@ -1269,7 +1269,6 @@ void CompactionIterator::DecideOutputLevel() {
}
}
-ROCKSDB_FLATTEN
void CompactionIterator::PrepareOutput() {
if (Valid()) {
if (LIKELY(!is_range_del_)) {
diff --git a/db/compaction/compaction_job.cc b/db/compaction/compaction_job.cc
index 9d0cc0690e..25368c8523 100644
--- a/db/compaction/compaction_job.cc
+++ b/db/compaction/compaction_job.cc
@@ -932,7 +932,46 @@ if (stats_) {
uint64_t expected =
compaction_stats_.stats.num_input_records - num_input_range_del;
uint64_t actual = compaction_job_stats_->num_input_records;
- if (expected != actual) {
+ auto can_verify_record_count = [&] {
+ auto* c = compact_->compaction;
+ auto* cfd = c->column_family_data();
+ auto* tc = cfd->table_cache();
+ const auto* cf_options = c->mutable_cf_options();
+ const ReadOptions ro(Env::IOActivity::kCompaction);
+ for (const auto& input : *c->inputs()) {
+ for (const auto* file : input.files) {
+ bool supported;
+ if (auto* reader = file->fd.table_reader) {
+ TEST_SYNC_POINT("CompactionJob::VerifyRecordCount:PinnedReader");
+ supported = reader->IsNumEntriesExact();
+ } else {
+ TEST_SYNC_POINT("CompactionJob::VerifyRecordCount:FindTable");
+ TableCache::TypedHandle* handle = nullptr;
+ Status s = tc->FindTable(
+ ro, file_options_, cfd->internal_comparator(), *file,
+ &handle, cf_options->block_protection_bytes_per_key,
+ cf_options->prefix_extractor);
+ supported = true;
+ if (s.ok()) {
+ supported = tc->GetTableReaderFromHandle(handle)->IsNumEntriesExact();
+ tc->ReleaseHandle(handle);
+ }
+ TEST_SYNC_POINT_CALLBACK("CompactionJob::VerifyRecordCount:FindTableStatus", &s);
+ if (!s.ok()) {
+ return true; // Keep the original mismatch error.
+ }
+ }
+ if (!supported) {
+ TEST_SYNC_POINT("CompactionJob::VerifyRecordCount:Unsupported");
+ return false;
+ }
+ }
+ }
+ TEST_SYNC_POINT("CompactionJob::VerifyRecordCount:Supported");
+ return true;
+ };
+ if (expected != actual && can_verify_record_count()) {
+ TEST_SYNC_POINT("CompactionJob::VerifyRecordCount:Mismatch");
std::string msg =
"Total number of input records: " + std::to_string(expected) +
", but processed " + std::to_string(actual) + " records.";
diff --git a/db/db_compaction_test.cc b/db/db_compaction_test.cc
index 4975b1cef3..d23232a865 100644
--- a/db/db_compaction_test.cc
+++ b/db/db_compaction_test.cc
@@ -10131,9 +10131,16 @@ TEST_F(DBCompactionTest, VerifyRecordCount) {
*(bool*)stop_ptr = true;
}
});
+ int supported = 0, mismatches = 0;
+ SyncPoint::GetInstance()->SetCallBack(
+ "CompactionJob::VerifyRecordCount:Supported", [&](void*) { supported++; });
+ SyncPoint::GetInstance()->SetCallBack(
+ "CompactionJob::VerifyRecordCount:Mismatch", [&](void*) { mismatches++; });
SyncPoint::GetInstance()->EnableProcessing();
Status s = db_->CompactRange(CompactRangeOptions(), nullptr, nullptr);
+ ASSERT_GT(supported, 0);
+ ASSERT_GT(mismatches, 0);
ASSERT_TRUE(s.IsCorruption());
const char* expect =
"Compaction number of input keys does not match number of keys "
diff --git a/db/db_cspp_crash_safe_test.cc b/db/db_cspp_crash_safe_test.cc
new file mode 100644
index 0000000000..fd5fede282
--- /dev/null
+++ b/db/db_cspp_crash_safe_test.cc
@@ -0,0 +1,5870 @@
+// Copyright (c) 2026-present, Topling Inc.
+// Crash-safe leftover recover: Convert + WAL tail, sync-point injection.
+
+#include
+#include
+#include
+
+#include
+#include
+#include
+#include
+#include
+#include
+#include
+#include
+#include