Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 69 additions & 5 deletions crates/storage/src/coverage_index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ const MAX_COVERAGE_INDEX_V2_BUCKET_READ_THREADS: usize = 8;
const MAX_COVERAGE_INDEX_V2_DELTA_GET_THREADS: usize = 8;
const MAX_EXACT_EVM_LOGS_V2_DELTA_OBJECTS: usize = 128;
const MAX_SEMANTIC_FALLBACK_V2_DELTA_OBJECTS: usize = 128;
const MAX_EMPTY_COVERAGE_COALESCE_V2_DELTA_OBJECTS: usize = 128;
const COVERAGE_INDEX_V2_BUCKET_READ_CACHE_TTL: Duration = Duration::from_secs(300);
const MAX_COVERAGE_INDEX_V2_BUCKET_READ_CACHE_ENTRIES: usize = 512;
const MAX_COVERAGE_INDEX_V2_BUCKET_READ_CACHE_BYTES: usize = 256 * 1024 * 1024;
Expand Down Expand Up @@ -2028,12 +2029,12 @@ where
S: ObjectStore,
{
let probe = empty_coverage_coalescing_probe(entry)?;
let probe_buckets = coverage_index_v2_entry_buckets(&probe);
if empty_coverage_coalescing_over_budget(object_store, &probe_buckets)? {
return Ok(None);
}
let mut existing_entries = preexisting_entries.to_vec();
read_entries_for_v2_buckets(
object_store,
coverage_index_v2_entry_buckets(&probe),
&mut existing_entries,
)?;
read_entries_for_v2_buckets(object_store, probe_buckets, &mut existing_entries)?;
let mut existing_manifest = Manifest {
entries: existing_entries,
};
Expand Down Expand Up @@ -2108,6 +2109,69 @@ where

type CoalescedEmptyV2Entry = (ManifestEntry, Vec<ManifestEntry>, BTreeSet<String>);

fn empty_coverage_coalescing_over_budget<S>(
object_store: &S,
buckets: &BTreeSet<CoverageIndexV2Bucket>,
) -> Result<bool, DatalensError>
where
S: ObjectStore,
{
for bucket in buckets {
let over_budget_key = coverage_index_v2_bucket_recent_key(object_store, bucket);
if is_recent_over_budget_v2_read(&over_budget_key) {
enqueue_v2_compaction_for_over_budget_bucket(object_store, bucket);
return Ok(true);
}
if v2_bucket_pending_delta_count_over_budget(
object_store,
bucket,
MAX_EMPTY_COVERAGE_COALESCE_V2_DELTA_OBJECTS,
)? {
mark_recent_over_budget_v2_read(over_budget_key);
enqueue_v2_compaction_for_over_budget_bucket(object_store, bucket);
log::warn!(
"storage coverage index v2 empty coalesce skipped over-budget bucket chain_key={} scope={} bucket={}-{} max_delta_objects={}",
bucket.chain_key,
bucket.scope,
bucket.bucket_start,
bucket.bucket_end,
MAX_EMPTY_COVERAGE_COALESCE_V2_DELTA_OBJECTS
);
return Ok(true);
}
}
Ok(false)
}

fn v2_bucket_pending_delta_count_over_budget<S>(
object_store: &S,
bucket: &CoverageIndexV2Bucket,
max_delta_objects: usize,
) -> Result<bool, DatalensError>
where
S: ObjectStore,
{
let mut compacted_delta_keys = BTreeSet::new();
let mut start_after_delta_key = None;
if let Some((_, head)) = latest_v2_snapshot_head_object(object_store, bucket)? {
if !head.included_delta_high_watermark.is_empty() {
start_after_delta_key = Some(head.included_delta_high_watermark);
} else {
let snapshot = read_v2_snapshot(object_store, bucket, &head.snapshot_key)?;
compacted_delta_keys.extend(snapshot.compacted_delta_keys);
}
}
Ok(list_v2_delta_object_metadata_for_bucket(
object_store,
bucket,
&compacted_delta_keys,
max_delta_objects.checked_add(1),
start_after_delta_key.as_deref(),
)?
.len()
> max_delta_objects)
}

fn can_merge_empty_coverage(left: &ManifestEntry, right: &ManifestEntry) -> bool {
left.object_key.is_none()
&& right.object_key.is_none()
Expand Down
72 changes: 72 additions & 0 deletions crates/storage/tests/manifest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5082,6 +5082,78 @@ fn test_coverage_index_v2_many_empty_deltas_coalesce_to_one_pending_delta() {
);
}

#[test]
fn test_coverage_index_v2_empty_coalesce_skips_over_budget_bucket() {
let storage = LocalStorage::new(temp_storage_root("coverage-index-v2-empty-over-budget"));
let chain = test_chain();
let selector = DatasetSelector::all();
let rows = DatasetRows::new(DatasetKey::evm_blocks(), QueryRows::EvmBlocks(Vec::new()))
.expect("dataset rows");
let scope = coverage_index_v2_exact_scope(
&DatasetKey::evm_blocks(),
"block",
&selector,
ManifestFinalityLevel::Safe,
);

for block in 1..=129 {
let range = LedgerRange::blocks(block, block).expect("valid range");
let entry = empty_manifest_entry(
&chain,
DatasetKey::evm_blocks(),
&selector,
range.clone(),
ManifestFinalityLevel::Safe,
);
write_coverage_index_v2_delta(
&storage,
&chain,
&scope,
&range,
&format!("{block:04}"),
vec![entry],
);
}

storage
.write_rows(StorageWriteRequest {
chain: &chain,
dataset_key: DatasetKey::evm_blocks(),
selector: &selector,
range: LedgerRange::blocks(130, 130).expect("valid range"),
rows: &rows,
finality_level: FinalityLevel::Safe,
record_empty_coverage: true,
})
.expect("write empty coverage");

let delta_prefix = coverage_index_v2_delta_prefix(
&chain,
&scope,
&LedgerRange::blocks(1, 1).expect("valid range"),
);
assert_eq!(
storage
.object_store()
.list(&delta_prefix)
.expect("coverage index v2 delta list")
.len(),
130,
"over-budget foreground write should append instead of coalescing the hot bucket"
);
assert_eq!(
storage
.covered_ranges(
&chain,
&DatasetKey::evm_blocks(),
&selector,
LedgerRange::blocks(1, 130).expect("valid range"),
)
.expect("covered ranges"),
vec![LedgerRange::blocks(1, 130).expect("valid range")]
);
}

#[test]
fn test_empty_coverage_does_not_coalesce_over_real_data_object() {
let object_store = CountingObjectStore::new(LocalObjectStore::new(temp_storage_root(
Expand Down
Loading