From e2eba63a72c3d07abd11bf969a90ed6af0972c35 Mon Sep 17 00:00:00 2001 From: Jon Haddad Date: Tue, 9 Jun 2026 20:49:21 -0700 Subject: [PATCH] Add partial-set and multi-output differential compaction coverage Partial-set compactions (PartialSetDifferentialCompactionTest): only a subset of live sstables participates and the purge evaluator must consult overlapping non-participating sstables before dropping tombstones. Scenarios pin purge blocked by an overlapping non-participant holding older data, purge allowed with a disjoint non-participant, and tombstones retained over a shadowed non-participant left out of the compaction. All writes use explicit timestamps so purge decisions are deterministic. Multi-output compactions (MultiOutputDifferentialCompactionTest): a size-capped MaxSSTableSizeWriter (8KiB, compression disabled) forces writer switches at partition boundaries; 4-6 outputs per compaction are compared pairwise. Covers switch-decision parity and per-output finalization (first/last keys, promoted index, per-output stats). Harness extensions in DifferentialCompactionTester: - compactPath takes an input subset; outputs are identified by before/after descriptor diff; non-participating sstables stay live and untouched across the restore - TaskFactory hook so scenarios can shape the CompactionTask (e.g. override getCompactionAwareWriter for multi-output) identically for both pipelines - assertCursorMatchesIterator returns the captured output so scenarios can assert their mechanism actually fired: multi-output scenarios require >= 2 outputs, partial-set scenarios verify the purge outcome in the captured dump. Motivated by the first multi-output run being vacuously green: MaxSSTableSizeWriter switches on estimated on-disk (compressed) bytes and repetitive test padding compressed below the cap, so no switch ever happened. --- .../DifferentialCompactionTester.java | 99 +++++++--- ...MultiOutputDifferentialCompactionTest.java | 145 ++++++++++++++ .../PartialSetDifferentialCompactionTest.java | 186 ++++++++++++++++++ 3 files changed, 408 insertions(+), 22 deletions(-) create mode 100644 test/unit/org/apache/cassandra/db/compaction/differential/MultiOutputDifferentialCompactionTest.java create mode 100644 test/unit/org/apache/cassandra/db/compaction/differential/PartialSetDifferentialCompactionTest.java diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/DifferentialCompactionTester.java b/test/unit/org/apache/cassandra/db/compaction/differential/DifferentialCompactionTester.java index e724683995..4a451a5658 100644 --- a/test/unit/org/apache/cassandra/db/compaction/differential/DifferentialCompactionTester.java +++ b/test/unit/org/apache/cassandra/db/compaction/differential/DifferentialCompactionTester.java @@ -34,8 +34,6 @@ import java.util.SortedSet; import java.util.TreeMap; import java.util.TreeSet; -import com.google.common.collect.ImmutableSet; - import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.cql3.CQLTester; import org.apache.cassandra.db.ColumnFamilyStore; @@ -113,13 +111,37 @@ public abstract class DifferentialCompactionTester extends CQLTester final List sstables = new ArrayList<>(); } + /** Creates the CompactionTask for one differential run. MUST honor keepOriginals=true. */ + public interface TaskFactory + { + CompactionTask create(ColumnFamilyStore cfs, LifecycleTransaction txn, long gcBefore); + } + + public static final TaskFactory DEFAULT_TASK = (cfs, txn, gcBefore) -> new CompactionTask(cfs, txn, gcBefore, true); + /** * Runs both compaction paths over the current live sstables of the table and asserts * byte + logical equivalence. The byteDiffAllowlist contains component names (e.g. * "Statistics.db") that are permitted to differ at the byte level; logical equivalence * is always enforced. */ - protected void assertCursorMatchesIterator(ColumnFamilyStore cfs, Set byteDiffAllowlist) throws Exception + protected CapturedOutput assertCursorMatchesIterator(ColumnFamilyStore cfs, Set byteDiffAllowlist) throws Exception + { + return assertCursorMatchesIterator(cfs, cfs.getLiveSSTables(), byteDiffAllowlist, DEFAULT_TASK); + } + + /** + * Variant for partial-set compactions (inputs is a subset of the live sstables; the rest + * stay live and participate in purge-overlap decisions) and for custom CompactionTask + * shapes (e.g. multi-output writers via an overridden getCompactionAwareWriter). + */ + /** Returns the iterator-path capture so scenarios can assert structural expectations + * (e.g. multi-output scenarios MUST verify more than one sstable was produced — + * a scenario that does not exercise its mechanism passes vacuously). */ + protected CapturedOutput assertCursorMatchesIterator(ColumnFamilyStore cfs, + Set inputs, + Set byteDiffAllowlist, + TaskFactory taskFactory) throws Exception { long gcBefore = cfs.getDefaultGcBefore(FBUtilities.nowInSeconds()); Path scratch = Files.createTempDirectory("differential-compaction"); @@ -133,9 +155,19 @@ public abstract class DifferentialCompactionTester extends CQLTester DatabaseDescriptor.setSSTablePreemptiveOpenIntervalInMiB(-1); try { - CapturedOutput iterator = compactPath(cfs, false, gcBefore, scratch.resolve("iterator")); - CapturedOutput cursor = compactPath(cfs, true, gcBefore, scratch.resolve("cursor")); + CapturedOutput iterator = compactPath(cfs, inputs, false, gcBefore, scratch.resolve("iterator"), taskFactory); + // the input INSTANCES were replaced during restore; re-resolve the subset by descriptor + Set inputDescs = new HashSet<>(); + for (SSTableReader in : inputs) + inputDescs.add(in.descriptor); + Set reResolved = new HashSet<>(); + for (SSTableReader live : cfs.getLiveSSTables()) + if (inputDescs.contains(live.descriptor)) + reResolved.add(live); + assertEquals("input subset lost across restore", inputs.size(), reResolved.size()); + CapturedOutput cursor = compactPath(cfs, reResolved, true, gcBefore, scratch.resolve("cursor"), taskFactory); assertEquivalentOutputs(iterator, cursor, byteDiffAllowlist); + return iterator; } finally { @@ -144,36 +176,57 @@ public abstract class DifferentialCompactionTester extends CQLTester } /** - * Runs one compaction path over the current live sstables, captures the outputs, and - * restores the live set to the original inputs so the other path sees identical bytes. + * Runs one compaction path over the given input subset (non-participating live sstables + * stay live and feed purge-overlap decisions), captures the outputs, and restores the live + * set so the other path sees identical bytes. */ - protected CapturedOutput compactPath(ColumnFamilyStore cfs, boolean cursor, long gcBefore, Path scratch) throws Exception + protected CapturedOutput compactPath(ColumnFamilyStore cfs, + Set inputs, + boolean cursor, + long gcBefore, + Path scratch, + TaskFactory taskFactory) throws Exception { DatabaseDescriptor.setCursorCompactionEnabled(cursor); - Set inputs = ImmutableSet.copyOf(cfs.getLiveSSTables()); assertFalse("scenario produced no input sstables", inputs.isEmpty()); + Set liveBeforeDescs = new HashSet<>(); + int liveBeforeCount = 0; + for (SSTableReader live : cfs.getLiveSSTables()) + { + liveBeforeDescs.add(live.descriptor); + liveBeforeCount++; + } List inputDescriptors = new ArrayList<>(); for (SSTableReader in : inputs) + { + assertTrue("input is not live", liveBeforeDescs.contains(in.descriptor)); inputDescriptors.add(in.descriptor); + } + Set inputDescs = new HashSet<>(inputDescriptors); if (cursor) assertCursorPathWillRun(cfs, inputs, gcBefore); LifecycleTransaction txn = cfs.getTracker().tryModify(inputs, OperationType.COMPACTION); assertNotNull("unable to mark inputs compacting", txn); - new CompactionTask(cfs, txn, gcBefore, true /* keepOriginals */).execute(ActiveCompactionsTracker.NOOP); + taskFactory.create(cfs, txn, gcBefore).execute(ActiveCompactionsTracker.NOOP); - // IMPORTANT: with keepOriginals=true the originals (or early-open clones of them, with - // moved starts) may remain in the live set as DIFFERENT reader instances. Outputs must - // therefore be identified by descriptor, never by instance, or the harness would treat - // retained originals as outputs and delete the input files. - Set inputDescs = new HashSet<>(inputDescriptors); + // Outputs are identified by descriptor diff against the pre-compaction live set: + // with keepOriginals=true the originals (or early-open clones with moved starts) may + // remain live as DIFFERENT reader instances, and non-participating sstables are live + // throughout. Instance identity is never trusted here. List liveAfter = new ArrayList<>(cfs.getLiveSSTables()); List outputs = new ArrayList<>(); + List retainedInputClones = new ArrayList<>(); for (SSTableReader reader : liveAfter) - if (!inputDescs.contains(reader.descriptor)) + { + if (!liveBeforeDescs.contains(reader.descriptor)) outputs.add(reader); + else if (inputDescs.contains(reader.descriptor)) + retainedInputClones.add(reader); + // else: non-participant — leave completely untouched + } outputs.sort(Comparator.comparing(SSTableReader::getFirst)); CapturedOutput captured = new CapturedOutput(); @@ -186,11 +239,13 @@ public abstract class DifferentialCompactionTester extends CQLTester outputFiles.add(out.descriptor.fileFor(c).toPath()); } - // Restore: clear the whole live set (outputs AND any retained/moved-start originals), - // delete output files only, then reopen every input fresh from its descriptor so the - // second run sees pristine full-range readers identical to the first run's. - cfs.getTracker().removeUnsafe(new HashSet<>(liveAfter)); - for (SSTableReader reader : liveAfter) + // Restore: delist + release outputs and any retained input clones, delete output files + // only, then reopen every input fresh from its descriptor so the second run sees + // pristine full-range readers identical to the first run's. + Set toRemove = new HashSet<>(outputs); + toRemove.addAll(retainedInputClones); + cfs.getTracker().removeUnsafe(toRemove); + for (SSTableReader reader : toRemove) reader.selfRef().release(); for (Path f : outputFiles) Files.deleteIfExists(f); @@ -204,7 +259,7 @@ public abstract class DifferentialCompactionTester extends CQLTester reopened.add(SSTableReader.open(cfs, desc)); } cfs.getTracker().addInitialSSTables(reopened); - assertEquals("restore failed: live sstable count", inputs.size(), cfs.getLiveSSTables().size()); + assertEquals("restore failed: live sstable count", liveBeforeCount, cfs.getLiveSSTables().size()); return captured; } diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/MultiOutputDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/MultiOutputDifferentialCompactionTest.java new file mode 100644 index 0000000000..2b0cc57aca --- /dev/null +++ b/test/unit/org/apache/cassandra/db/compaction/differential/MultiOutputDifferentialCompactionTest.java @@ -0,0 +1,145 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.cassandra.db.compaction.differential; + +import java.util.Set; + +import org.junit.Test; + +import static org.junit.Assert.assertTrue; + +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.Directories; +import org.apache.cassandra.db.compaction.CompactionTask; +import org.apache.cassandra.db.compaction.writers.CompactionAwareWriter; +import org.apache.cassandra.db.compaction.writers.MaxSSTableSizeWriter; +import org.apache.cassandra.db.lifecycle.ILifecycleTransaction; +import org.apache.cassandra.db.lifecycle.LifecycleTransaction; +import org.apache.cassandra.io.sstable.format.SSTableReader; + +/** + * Multi-output compactions: a size-capped CompactionAwareWriter forces writer switches at + * partition boundaries, so one compaction produces several sstables. Both pipelines receive + * the same writer type (CursorCompactionPipeline and IteratorCompactionPipeline both call + * task.getCompactionAwareWriter); the differential assertion requires the SAME number of + * outputs with byte-identical components pairwise — divergent switch decisions would show up + * as a count mismatch, divergent per-output finalization (first/last keys, promoted index, + * stats) as component diffs. + */ +public class MultiOutputDifferentialCompactionTest extends DifferentialCompactionTester +{ + private static final Set ALLOWLIST = Set.of(); + + /** CompactionTask whose writer switches outputs once the current one exceeds maxBytes. */ + private static TaskFactory sizeCapped(long maxBytes) + { + return (cfs, txn, gcBefore) -> new CompactionTask(cfs, txn, gcBefore, true /* keepOriginals */) + { + @Override + public CompactionAwareWriter getCompactionAwareWriter(ColumnFamilyStore cfs, + Directories directories, + ILifecycleTransaction transaction, + Set nonExpiredSSTables) + { + return new MaxSSTableSizeWriter(cfs, directories, transaction, nonExpiredSSTables, + maxBytes, 0, true /* keepOriginals */); + } + }; + } + + /** Many small partitions split across several outputs. */ + @Test + public void manyPartitionsSplitAcrossOutputs() throws Exception + { + createTable("CREATE TABLE %s (pk bigint, ck bigint, v text, PRIMARY KEY (pk, ck)) " + + "WITH compression = {'enabled': 'false'}"); + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + cfs.disableAutoCompaction(); + + String padding = "x".repeat(100); + for (int round = 0; round < 2; round++) + { + for (long pk = 0; pk < 40; pk++) + for (long ck = 0; ck < 10; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?)", pk, ck, padding + round + "-" + ck); + flush(); + } + + // ~40 partitions * ~1KB each; 8KB cap => several outputs, switch points mid-stream + CapturedOutput out = assertCursorMatchesIterator(cfs, cfs.getLiveSSTables(), ALLOWLIST, sizeCapped(8 * 1024)); + assertTrue("scenario must produce multiple outputs to test anything, got " + out.sstables.size(), + out.sstables.size() >= 2); + } + + /** Tombstones and static rows crossing output boundaries. */ + @Test + public void tombstonesAndStaticsAcrossOutputs() throws Exception + { + createTable("CREATE TABLE %s (pk bigint, s text static, ck bigint, v text, PRIMARY KEY (pk, ck)) " + + "WITH gc_grace_seconds = 864000 AND compression = {'enabled': 'false'}"); + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + cfs.disableAutoCompaction(); + + String padding = "y".repeat(100); + for (long pk = 0; pk < 30; pk++) + { + if (pk % 3 == 0) + execute("INSERT INTO %s (pk, s, ck, v) VALUES (?, ?, ?, ?)", pk, "static" + pk, 0L, padding); + for (long ck = 0; ck < 10; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?)", pk, ck, padding + ck); + } + flush(); + + // second sstable: range tombstones inside many partitions + some partition deletes + for (long pk = 0; pk < 30; pk += 2) + execute("DELETE FROM %s WHERE pk = ? AND ck >= 3 AND ck < 7", pk); + execute("DELETE FROM %s WHERE pk = 5"); + execute("DELETE FROM %s WHERE pk = 15"); + flush(); + + CapturedOutput out = assertCursorMatchesIterator(cfs, cfs.getLiveSSTables(), ALLOWLIST, sizeCapped(8 * 1024)); + assertTrue("scenario must produce multiple outputs to test anything, got " + out.sstables.size(), + out.sstables.size() >= 2); + } + + /** One output ending exactly at a wide partition boundary, the next starting fresh. */ + @Test + public void widePartitionsForceFrequentSwitches() throws Exception + { + createTable("CREATE TABLE %s (pk bigint, ck bigint, v text, PRIMARY KEY (pk, ck)) " + + "WITH compression = {'enabled': 'false'}"); + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + cfs.disableAutoCompaction(); + + String padding = "z".repeat(300); + // each partition ~6KB > cap/2: nearly every partition triggers a switch decision + for (int round = 0; round < 2; round++) + { + for (long pk = 0; pk < 12; pk++) + for (long ck = 0; ck < 20; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?)", pk, ck, padding + round); + execute("DELETE FROM %s WHERE pk = ? AND ck >= 5 AND ck < 15", (long) round); + flush(); + } + + CapturedOutput out = assertCursorMatchesIterator(cfs, cfs.getLiveSSTables(), ALLOWLIST, sizeCapped(8 * 1024)); + assertTrue("scenario must produce multiple outputs to test anything, got " + out.sstables.size(), + out.sstables.size() >= 2); + } +} diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/PartialSetDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/PartialSetDifferentialCompactionTest.java new file mode 100644 index 0000000000..7bffeba1f6 --- /dev/null +++ b/test/unit/org/apache/cassandra/db/compaction/differential/PartialSetDifferentialCompactionTest.java @@ -0,0 +1,186 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.cassandra.db.compaction.differential; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +import org.junit.Test; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.io.sstable.format.SSTableReader; + +/** + * Partial-set compactions: only a subset of the live sstables participates, and the purge + * evaluator must consult overlapping NON-participating sstables before dropping tombstones + * (CompactionController.getPurgeEvaluator / CursorCompactor's shouldPurge). The differential + * assertion is agnostic to which outcome is correct — both paths must simply agree, byte for + * byte. All writes use explicit timestamps so purge decisions are deterministic. + */ +public class PartialSetDifferentialCompactionTest extends DifferentialCompactionTester +{ + private static final Set ALLOWLIST = Set.of(); + + private static String allJson(CapturedOutput out) + { + StringBuilder sb = new StringBuilder(); + for (CapturedSSTable s : out.sstables) + sb.append(s.json); + return sb.toString(); + } + + /** Flushes and returns the sstable that flush produced; tracks flush order explicitly. */ + private SSTableReader flushAndTrack(ColumnFamilyStore cfs, List flushed) throws Throwable + { + Set before = new HashSet<>(cfs.getLiveSSTables()); + flush(); + List fresh = new ArrayList<>(cfs.getLiveSSTables()); + fresh.removeAll(before); + if (fresh.size() != 1) + throw new AssertionError("expected exactly one new sstable from flush, got " + fresh.size()); + flushed.add(fresh.get(0)); + return fresh.get(0); + } + + /** + * Tombstones whose purge is BLOCKED by an overlapping non-participant: the partition also + * exists in a non-compacting sstable with OLDER timestamps, so the tombstone must survive + * the partial compaction in both paths. + */ + @Test + public void purgeBlockedByOverlappingNonParticipant() throws Throwable + { + createTable("CREATE TABLE %s (pk bigint, ck bigint, v text, PRIMARY KEY (pk, ck)) " + + "WITH gc_grace_seconds = 0"); + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + cfs.disableAutoCompaction(); + + List flushed = new ArrayList<>(); + // sstable A: data at ts=100 + for (long pk = 0; pk < 6; pk++) + for (long ck = 0; ck < 10; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?) USING TIMESTAMP 100", pk, ck, "a" + ck); + flushAndTrack(cfs, flushed); + + // sstable B: row + partition tombstones at ts=200 + for (long ck = 0; ck < 5; ck++) + execute("DELETE FROM %s USING TIMESTAMP 200 WHERE pk = 1 AND ck = ?", ck); + execute("DELETE FROM %s USING TIMESTAMP 200 WHERE pk = 2"); + flushAndTrack(cfs, flushed); + + // sstable C (non-participant): same partitions with OLDER data (ts=50) — purging the + // ts=200 tombstones would resurrect this data, so the evaluator must block it + for (long pk = 1; pk <= 2; pk++) + for (long ck = 0; ck < 10; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?) USING TIMESTAMP 50", pk, ck, "old" + ck); + flushAndTrack(cfs, flushed); + + Thread.sleep(1100); // deletion local times strictly in the past for both runs + + CapturedOutput out = assertCursorMatchesIterator(cfs, new HashSet<>(flushed.subList(0, 2)), ALLOWLIST, DEFAULT_TASK); + // non-vacuousness: the overlap must have actually BLOCKED the purge + assertTrue("expected ts=200 tombstones retained because of the overlapping non-participant", + allJson(out).contains("\"marked_deleted\":\"200\"")); + } + + /** + * Tombstones whose purge is ALLOWED despite a non-participant existing: the non-compacting + * sstable holds entirely different partitions, so nothing blocks the drop. Both paths must + * purge identically (and drop the emptied partitions). + */ + @Test + public void purgeAllowedWithDisjointNonParticipant() throws Throwable + { + createTable("CREATE TABLE %s (pk bigint, ck bigint, v text, PRIMARY KEY (pk, ck)) " + + "WITH gc_grace_seconds = 0"); + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + cfs.disableAutoCompaction(); + + List flushed = new ArrayList<>(); + // sstable A: data at ts=100 for pk 0-5 + for (long pk = 0; pk < 6; pk++) + for (long ck = 0; ck < 10; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?) USING TIMESTAMP 100", pk, ck, "a" + ck); + flushAndTrack(cfs, flushed); + + // sstable B: tombstones at ts=200 for pk 1-2 + for (long ck = 0; ck < 10; ck++) + execute("DELETE FROM %s USING TIMESTAMP 200 WHERE pk = 1 AND ck = ?", ck); + execute("DELETE FROM %s USING TIMESTAMP 200 WHERE pk = 2"); + flushAndTrack(cfs, flushed); + + // sstable C (non-participant): DIFFERENT partitions only (pk 100-105) + for (long pk = 100; pk < 106; pk++) + for (long ck = 0; ck < 10; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?) USING TIMESTAMP 50", pk, ck, "other" + ck); + flushAndTrack(cfs, flushed); + + Thread.sleep(1100); + + CapturedOutput out = assertCursorMatchesIterator(cfs, new HashSet<>(flushed.subList(0, 2)), ALLOWLIST, DEFAULT_TASK); + // non-vacuousness: with no overlap on those keys, the tombstones must actually purge + assertFalse("expected ts=200 tombstones purged (disjoint non-participant cannot block)", + allJson(out).contains("\"marked_deleted\":\"200\"")); + } + + /** + * Compacting newer sstables while the OLDEST (holding the shadowed data) stays out: + * tombstones must be retained because the overlapping non-participant holds data they + * still shadow. + */ + @Test + public void tombstonesRetainedOverShadowedNonParticipant() throws Throwable + { + createTable("CREATE TABLE %s (pk bigint, ck bigint, v text, PRIMARY KEY (pk, ck)) " + + "WITH gc_grace_seconds = 0"); + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + cfs.disableAutoCompaction(); + + List flushed = new ArrayList<>(); + // sstable A (will be the non-participant): the data being shadowed, ts=100 + for (long pk = 0; pk < 6; pk++) + for (long ck = 0; ck < 10; ck++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?) USING TIMESTAMP 100", pk, ck, "base" + ck); + flushAndTrack(cfs, flushed); + + // sstable B: tombstones at ts=200 over A's data + for (long ck = 0; ck < 10; ck++) + execute("DELETE FROM %s USING TIMESTAMP 200 WHERE pk = 3 AND ck = ?", ck); + execute("DELETE FROM %s USING TIMESTAMP 200 WHERE pk = 4"); + flushAndTrack(cfs, flushed); + + // sstable C: unrelated newer writes, merges with B + for (long pk = 0; pk < 6; pk++) + execute("INSERT INTO %s (pk, ck, v) VALUES (?, ?, ?) USING TIMESTAMP 300", pk, 20L, "new"); + flushAndTrack(cfs, flushed); + + Thread.sleep(1100); + + // compact B+C, leaving A (the shadowed data) out + CapturedOutput out = assertCursorMatchesIterator(cfs, new HashSet<>(flushed.subList(1, 3)), ALLOWLIST, DEFAULT_TASK); + // non-vacuousness: tombstones still shadow A's data and must survive + assertTrue("expected ts=200 tombstones retained over the shadowed non-participant", + allJson(out).contains("\"marked_deleted\":\"200\"")); + } +}