This is an automated email from the ASF dual-hosted git repository.

ifesdjeen pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/cassandra-simulator.git

commit d40531e0e0fe7da92eb5cdc79d16d3af506d5d1d
Author: Alex Petrov <[email protected]>
AuthorDate: Wed Jul 22 10:18:34 2026 +0200

    Journal code
---
 journal/build.gradle                               |   8 +-
 .../src/main/java/accord/utils/UnhandledEnum.java  |   9 +
 .../cassandra/concurrent/ExecutorFactory.java      |  44 ++
 .../cassandra/concurrent/ImmediateExecutor.java    |  14 +
 .../apache/cassandra/concurrent/Interruptible.java |  11 +
 .../concurrent/SequentialExecutorPlus.java         |   8 +
 .../apache/cassandra/concurrent/Shutdownable.java  |  10 +
 .../java/org/apache/cassandra/db/TypeSizes.java    |   8 +
 .../java/org/apache/cassandra/io/FSReadError.java  |   9 +
 .../java/org/apache/cassandra/io/FSWriteError.java |   9 +
 .../apache/cassandra/io/util/DataInputBuffer.java  |  31 ++
 .../apache/cassandra/io/util/DataInputPlus.java    |  27 ++
 .../apache/cassandra/io/util/DataOutputBuffer.java |  59 +++
 .../apache/cassandra/io/util/DataOutputPlus.java   |  10 +
 .../java/org/apache/cassandra/io/util/File.java    |  39 ++
 .../cassandra/io/util/FileInputStreamPlus.java     |  20 +
 .../cassandra/io/util/FileOutputStreamPlus.java    |  28 ++
 .../org/apache/cassandra/io/util/PathUtils.java    |  28 ++
 .../apache/cassandra/journal/ActiveSegment.java    | 444 ++++-----------------
 .../org/apache/cassandra/journal/Compactor.java    |  85 +---
 .../java/org/apache/cassandra/journal/Flusher.java | 264 ++----------
 .../java/org/apache/cassandra/journal/Journal.java |  11 +-
 .../java/org/apache/cassandra/journal/Segment.java |  81 +---
 .../org/apache/cassandra/journal/Segments.java     | 244 +----------
 .../apache/cassandra/journal/StaticSegment.java    | 359 ++---------------
 .../apache/cassandra/service/StorageService.java   |  10 +
 .../service/accord/serializers/Version.java        |  14 +
 .../apache/cassandra/utils/AbstractIterator.java   |  39 ++
 .../org/apache/cassandra/utils/ByteBufferUtil.java |  23 ++
 .../java/org/apache/cassandra/utils/Clock.java     |  12 +
 .../java/org/apache/cassandra/utils/Closeable.java |   7 +
 .../apache/cassandra/utils/CloseableIterator.java  |   9 +
 .../main/java/org/apache/cassandra/utils/Crc.java  |  29 ++
 .../org/apache/cassandra/utils/FBUtilities.java    |  22 +
 .../cassandra/utils/JVMStabilityInspector.java     |  10 +
 .../org/apache/cassandra/utils/LazyToString.java   |  21 +
 .../org/apache/cassandra/utils/MergeIterator.java  | 118 ++++++
 .../main/java/org/apache/cassandra/utils/Pair.java |  18 +
 .../java/org/apache/cassandra/utils/Simulate.java  |  18 +
 .../java/org/apache/cassandra/utils/SyncUtil.java  |  21 +
 .../org/apache/cassandra/utils/Throwables.java     |  23 ++
 .../cassandra/utils/concurrent/AsyncPromise.java   |  37 ++
 .../cassandra/utils/concurrent/CountDownLatch.java |  33 ++
 .../apache/cassandra/utils/concurrent/OpOrder.java |  52 +++
 .../apache/cassandra/utils/concurrent/Promise.java |  11 +
 .../org/apache/cassandra/utils/concurrent/Ref.java |  44 ++
 .../concurrent/UncheckedInterruptedException.java  |   6 +
 .../apache/cassandra/harry/checker/TestHelper.java |  19 +
 .../apache/cassandra/harry/gen/EntropySource.java  |  28 ++
 .../java/org/apache/cassandra/utils/TimeUUID.java  |   3 +-
 50 files changed, 1212 insertions(+), 1275 deletions(-)

diff --git a/journal/build.gradle b/journal/build.gradle
index c9a9a7a..8a14ff4 100644
--- a/journal/build.gradle
+++ b/journal/build.gradle
@@ -17,12 +17,12 @@ java {
 }
 
 dependencies {
-    implementation fileTree(dir: '../build/lib/jars', exclude: 
['logback-*.jar', 'slf4j-api-*.jar'])
-    implementation files('../build/classes/main', 
'../modules/accord/accord-core/build/classes/java/main')
-
-    testImplementation files('../build/test/classes', 
'../modules/accord/accord-core/build/classes/java/test')
+    implementation 'com.google.guava:guava:33.2.1-jre'
+    implementation 'com.google.code.findbugs:jsr305:3.0.2'
+    implementation 'org.jctools:jctools-core:4.0.5'
     implementation 'org.slf4j:slf4j-api:1.7.36'
 
+    testImplementation 'org.agrona:agrona:1.22.0'
     testImplementation 'junit:junit:4.13.2'
     testImplementation 'org.assertj:assertj-core:3.26.3'
     testImplementation 'org.quicktheories:quicktheories:0.26'
diff --git a/journal/src/main/java/accord/utils/UnhandledEnum.java 
b/journal/src/main/java/accord/utils/UnhandledEnum.java
new file mode 100644
index 0000000..f866174
--- /dev/null
+++ b/journal/src/main/java/accord/utils/UnhandledEnum.java
@@ -0,0 +1,9 @@
+package accord.utils;
+
+public class UnhandledEnum extends RuntimeException
+{
+    public UnhandledEnum(Enum<?> value)
+    {
+        super(String.valueOf(value));
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/concurrent/ExecutorFactory.java 
b/journal/src/main/java/org/apache/cassandra/concurrent/ExecutorFactory.java
new file mode 100644
index 0000000..7ea9294
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/concurrent/ExecutorFactory.java
@@ -0,0 +1,44 @@
+package org.apache.cassandra.concurrent;
+
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+public final class ExecutorFactory
+{
+    public static final class Global
+    {
+        private static final Factory INSTANCE = new Factory();
+        public static Factory executorFactory() { return INSTANCE; }
+    }
+
+    public static final class Factory
+    {
+        public SequentialExecutorPlus sequential(String name)
+        {
+            return new SimpleSequentialExecutor(name);
+        }
+    }
+
+    private static final class SimpleSequentialExecutor implements 
SequentialExecutorPlus
+    {
+        private final ThreadPoolExecutor executor;
+
+        private SimpleSequentialExecutor(String name)
+        {
+            this.executor = new ThreadPoolExecutor(1, 1, 0L, 
TimeUnit.MILLISECONDS,
+                                                   new LinkedBlockingQueue<>(),
+                                                   r -> { Thread t = new 
Thread(r, name); t.setDaemon(true); return t; });
+        }
+
+        @Override public void execute(Runnable command) { 
executor.execute(command); }
+        @Override public void shutdown() { executor.shutdown(); }
+        @Override public Object shutdownNow() { List<Runnable> pending = 
executor.shutdownNow(); return pending == null ? Collections.emptyList() : 
pending; }
+        @Override public boolean awaitTermination(long timeout, TimeUnit 
units) throws InterruptedException { return executor.awaitTermination(timeout, 
units); }
+        @Override public boolean isTerminated() { return 
executor.isTerminated(); }
+    }
+
+    private ExecutorFactory() {}
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/concurrent/ImmediateExecutor.java 
b/journal/src/main/java/org/apache/cassandra/concurrent/ImmediateExecutor.java
new file mode 100644
index 0000000..e243b87
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/concurrent/ImmediateExecutor.java
@@ -0,0 +1,14 @@
+package org.apache.cassandra.concurrent;
+
+import java.util.concurrent.Executor;
+
+public enum ImmediateExecutor implements Executor
+{
+    INSTANCE;
+
+    @Override
+    public void execute(Runnable command)
+    {
+        command.run();
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/concurrent/Interruptible.java 
b/journal/src/main/java/org/apache/cassandra/concurrent/Interruptible.java
new file mode 100644
index 0000000..10c4cc2
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/concurrent/Interruptible.java
@@ -0,0 +1,11 @@
+package org.apache.cassandra.concurrent;
+
+public final class Interruptible
+{
+    private Interruptible() {}
+
+    public static final class TerminateException extends RuntimeException
+    {
+        public TerminateException() {}
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/concurrent/SequentialExecutorPlus.java
 
b/journal/src/main/java/org/apache/cassandra/concurrent/SequentialExecutorPlus.java
new file mode 100644
index 0000000..304f836
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/concurrent/SequentialExecutorPlus.java
@@ -0,0 +1,8 @@
+package org.apache.cassandra.concurrent;
+
+import java.util.concurrent.Executor;
+
+public interface SequentialExecutorPlus extends Executor, Shutdownable
+{
+    void shutdown();
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/concurrent/Shutdownable.java 
b/journal/src/main/java/org/apache/cassandra/concurrent/Shutdownable.java
new file mode 100644
index 0000000..cbc07fb
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/concurrent/Shutdownable.java
@@ -0,0 +1,10 @@
+package org.apache.cassandra.concurrent;
+
+import java.util.concurrent.TimeUnit;
+
+public interface Shutdownable
+{
+    boolean isTerminated();
+    Object shutdownNow();
+    boolean awaitTermination(long timeout, TimeUnit units) throws 
InterruptedException;
+}
diff --git a/journal/src/main/java/org/apache/cassandra/db/TypeSizes.java 
b/journal/src/main/java/org/apache/cassandra/db/TypeSizes.java
new file mode 100644
index 0000000..fd1afe1
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/db/TypeSizes.java
@@ -0,0 +1,8 @@
+package org.apache.cassandra.db;
+
+public final class TypeSizes
+{
+    public static final int INT_SIZE = Integer.BYTES;
+
+    private TypeSizes() {}
+}
diff --git a/journal/src/main/java/org/apache/cassandra/io/FSReadError.java 
b/journal/src/main/java/org/apache/cassandra/io/FSReadError.java
new file mode 100644
index 0000000..70b3195
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/io/FSReadError.java
@@ -0,0 +1,9 @@
+package org.apache.cassandra.io;
+
+public class FSReadError extends RuntimeException
+{
+    public FSReadError(Throwable cause, Object path)
+    {
+        super(String.valueOf(path), cause);
+    }
+}
diff --git a/journal/src/main/java/org/apache/cassandra/io/FSWriteError.java 
b/journal/src/main/java/org/apache/cassandra/io/FSWriteError.java
new file mode 100644
index 0000000..120242b
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/io/FSWriteError.java
@@ -0,0 +1,9 @@
+package org.apache.cassandra.io;
+
+public class FSWriteError extends RuntimeException
+{
+    public FSWriteError(Throwable cause, Object path)
+    {
+        super(String.valueOf(path), cause);
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/io/util/DataInputBuffer.java 
b/journal/src/main/java/org/apache/cassandra/io/util/DataInputBuffer.java
new file mode 100644
index 0000000..66548f7
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/io/util/DataInputBuffer.java
@@ -0,0 +1,31 @@
+package org.apache.cassandra.io.util;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+
+public class DataInputBuffer implements DataInputPlus
+{
+    private final ByteBuffer buffer;
+
+    public DataInputBuffer(ByteBuffer buffer, boolean duplicate)
+    {
+        this.buffer = duplicate ? buffer.duplicate() : buffer.slice();
+    }
+
+    @Override public void readFully(byte[] b) { readFully(b, 0, b.length); }
+    @Override public void readFully(byte[] b, int off, int len) { 
buffer.get(b, off, len); }
+    @Override public int skipBytes(int n) { int k = Math.min(n, 
buffer.remaining()); buffer.position(buffer.position() + k); return k; }
+    @Override public boolean readBoolean() { return buffer.get() != 0; }
+    @Override public byte readByte() { return buffer.get(); }
+    @Override public int readUnsignedByte() { return 
Byte.toUnsignedInt(buffer.get()); }
+    @Override public short readShort() { return buffer.getShort(); }
+    @Override public int readUnsignedShort() { return 
Short.toUnsignedInt(buffer.getShort()); }
+    @Override public char readChar() { return buffer.getChar(); }
+    @Override public int readInt() { return buffer.getInt(); }
+    @Override public long readLong() { return buffer.getLong(); }
+    @Override public float readFloat() { return buffer.getFloat(); }
+    @Override public double readDouble() { return buffer.getDouble(); }
+    @Override public String readLine() throws IOException { throw new 
UnsupportedOperationException(); }
+    @Override public String readUTF() throws IOException { throw new 
UnsupportedOperationException(); }
+    @Override public int available() { return buffer.remaining(); }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/io/util/DataInputPlus.java 
b/journal/src/main/java/org/apache/cassandra/io/util/DataInputPlus.java
new file mode 100644
index 0000000..753a0c8
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/io/util/DataInputPlus.java
@@ -0,0 +1,27 @@
+package org.apache.cassandra.io.util;
+
+import java.io.DataInput;
+import java.io.IOException;
+
+public interface DataInputPlus extends DataInput, AutoCloseable
+{
+    default void skipBytesFully(int count) throws IOException
+    {
+        int skipped = 0;
+        while (skipped < count)
+        {
+            int cur = skipBytes(count - skipped);
+            if (cur <= 0)
+                throw new IOException("Unexpected EOF");
+            skipped += cur;
+        }
+    }
+
+    default int available() throws IOException
+    {
+        return 0;
+    }
+
+    @Override
+    default void close() throws IOException {}
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/io/util/DataOutputBuffer.java 
b/journal/src/main/java/org/apache/cassandra/io/util/DataOutputBuffer.java
new file mode 100644
index 0000000..fa2a248
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/io/util/DataOutputBuffer.java
@@ -0,0 +1,59 @@
+package org.apache.cassandra.io.util;
+
+import java.io.DataOutputStream;
+import java.io.IOException;
+import java.io.OutputStream;
+import java.nio.ByteBuffer;
+import java.util.Arrays;
+
+public class DataOutputBuffer extends OutputStream implements DataOutputPlus
+{
+    public static final ThreadLocal<DataOutputBuffer> scratchBuffer = 
ThreadLocal.withInitial(() -> new DataOutputBuffer(256));
+
+    private byte[] data;
+    private int length;
+    private final DataOutputStream dos;
+
+    public DataOutputBuffer(int initialCapacity)
+    {
+        this.data = new byte[initialCapacity];
+        this.dos = new DataOutputStream(this);
+    }
+
+    private void ensureCapacity(int extra)
+    {
+        int needed = length + extra;
+        if (needed > data.length)
+            data = Arrays.copyOf(data, Math.max(needed, data.length * 2));
+    }
+
+    public int getLength()
+    {
+        return length;
+    }
+
+    public ByteBuffer unsafeGetBufferAndFlip()
+    {
+        return ByteBuffer.wrap(data, 0, length);
+    }
+
+    public ByteBuffer buffer()
+    {
+        return ByteBuffer.wrap(Arrays.copyOf(data, length));
+    }
+
+    @Override public void write(int b) { ensureCapacity(1); data[length++] = 
(byte) b; }
+    @Override public void write(byte[] b, int off, int len) { 
ensureCapacity(len); System.arraycopy(b, off, data, length, len); length += 
len; }
+    @Override public void writeBoolean(boolean v) throws IOException { 
dos.writeBoolean(v); }
+    @Override public void writeByte(int v) throws IOException { 
dos.writeByte(v); }
+    @Override public void writeShort(int v) throws IOException { 
dos.writeShort(v); }
+    @Override public void writeChar(int v) throws IOException { 
dos.writeChar(v); }
+    @Override public void writeInt(int v) throws IOException { 
dos.writeInt(v); }
+    @Override public void writeLong(long v) throws IOException { 
dos.writeLong(v); }
+    @Override public void writeFloat(float v) throws IOException { 
dos.writeFloat(v); }
+    @Override public void writeDouble(double v) throws IOException { 
dos.writeDouble(v); }
+    @Override public void writeBytes(String s) throws IOException { 
dos.writeBytes(s); }
+    @Override public void writeChars(String s) throws IOException { 
dos.writeChars(s); }
+    @Override public void writeUTF(String s) throws IOException { 
dos.writeUTF(s); }
+    @Override public void close() { length = 0; }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/io/util/DataOutputPlus.java 
b/journal/src/main/java/org/apache/cassandra/io/util/DataOutputPlus.java
new file mode 100644
index 0000000..40322f6
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/io/util/DataOutputPlus.java
@@ -0,0 +1,10 @@
+package org.apache.cassandra.io.util;
+
+import java.io.DataOutput;
+import java.io.IOException;
+
+public interface DataOutputPlus extends DataOutput, AutoCloseable
+{
+    @Override
+    default void close() throws IOException {}
+}
diff --git a/journal/src/main/java/org/apache/cassandra/io/util/File.java 
b/journal/src/main/java/org/apache/cassandra/io/util/File.java
new file mode 100644
index 0000000..1de2ab7
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/io/util/File.java
@@ -0,0 +1,39 @@
+package org.apache.cassandra.io.util;
+
+import java.io.FilenameFilter;
+import java.io.IOException;
+import java.nio.channels.FileChannel;
+import java.nio.file.DirectoryStream;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.StandardCopyOption;
+import java.nio.file.StandardOpenOption;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.function.Predicate;
+
+public class File
+{
+    private final Path path;
+
+    public File(Path path) { this.path = path; }
+    public File(java.io.File file) { this(file.toPath()); }
+    public File(File parent, String child) { this(parent.path.resolve(child)); 
}
+
+    public Path toPath() { return path; }
+    public String name() { return path.getFileName().toString(); }
+    public File parent() { Path p = path.getParent(); return p == null ? null 
: new File(p); }
+    public boolean exists() { return Files.exists(path); }
+    public boolean createFileIfNotExists() throws IOException { if (exists()) 
return false; Files.createFile(path); return true; }
+    public void delete() { try { Files.deleteIfExists(path); } catch 
(IOException e) { throw new RuntimeException(e); } }
+    public void deleteIfExists() { delete(); }
+    public void deleteOnExit() { path.toFile().deleteOnExit(); }
+    public void deleteRecursiveOnExit() { path.toFile().deleteOnExit(); }
+    public void move(File to) { try { Files.move(path, to.path, 
StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE); } catch 
(IOException e) { try { Files.move(path, to.path, 
StandardCopyOption.REPLACE_EXISTING); } catch (IOException e2) { throw new 
RuntimeException(e2); } } }
+    public FileChannel newReadChannel() throws IOException { return 
FileChannel.open(path, StandardOpenOption.READ); }
+    public String[] listNames(FilenameFilter filter) throws IOException { 
java.io.File[] files = path.toFile().listFiles((dir, name) -> 
filter.accept(dir, name)); if (files == null) return new String[0]; String[] 
names = new String[files.length]; for (int i=0;i<files.length;i++) names[i] = 
files[i].getName(); return names; }
+    public List<File> listUnchecked(Predicate<File> predicate) { try 
(DirectoryStream<Path> stream = Files.newDirectoryStream(path)) { List<File> 
out = new ArrayList<>(); for (Path p : stream) { File f = new File(p); if 
(predicate.test(f)) out.add(f); } return out; } catch (IOException e) { throw 
new RuntimeException(e); } }
+    @Override public String toString() { return path.toString(); }
+    @Override public boolean equals(Object o) { if (!(o instanceof File)) 
return false; File f = (File) o; return path.equals(f.path); }
+    @Override public int hashCode() { return path.hashCode(); }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/io/util/FileInputStreamPlus.java 
b/journal/src/main/java/org/apache/cassandra/io/util/FileInputStreamPlus.java
new file mode 100644
index 0000000..d9e4156
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/io/util/FileInputStreamPlus.java
@@ -0,0 +1,20 @@
+package org.apache.cassandra.io.util;
+
+import java.io.BufferedInputStream;
+import java.io.DataInputStream;
+import java.io.FileInputStream;
+import java.io.IOException;
+
+public class FileInputStreamPlus extends DataInputStream implements 
DataInputPlus
+{
+    public FileInputStreamPlus(File file) throws IOException
+    {
+        super(new BufferedInputStream(new 
FileInputStream(file.toPath().toFile())));
+    }
+
+    @Override
+    public int available() throws IOException
+    {
+        return super.available();
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/io/util/FileOutputStreamPlus.java 
b/journal/src/main/java/org/apache/cassandra/io/util/FileOutputStreamPlus.java
new file mode 100644
index 0000000..4f10975
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/io/util/FileOutputStreamPlus.java
@@ -0,0 +1,28 @@
+package org.apache.cassandra.io.util;
+
+import java.io.BufferedOutputStream;
+import java.io.DataOutputStream;
+import java.io.FileOutputStream;
+import java.io.IOException;
+
+public class FileOutputStreamPlus extends DataOutputStream implements 
DataOutputPlus
+{
+    private final FileOutputStream fos;
+
+    public FileOutputStreamPlus(File file) throws IOException
+    {
+        this(new FileOutputStream(file.toPath().toFile()));
+    }
+
+    private FileOutputStreamPlus(FileOutputStream fos)
+    {
+        super(new BufferedOutputStream(fos));
+        this.fos = fos;
+    }
+
+    public void sync() throws IOException
+    {
+        flush();
+        fos.getFD().sync();
+    }
+}
diff --git a/journal/src/main/java/org/apache/cassandra/io/util/PathUtils.java 
b/journal/src/main/java/org/apache/cassandra/io/util/PathUtils.java
new file mode 100644
index 0000000..0a1ad42
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/io/util/PathUtils.java
@@ -0,0 +1,28 @@
+package org.apache.cassandra.io.util;
+
+import java.io.IOException;
+import java.nio.file.Path;
+import java.nio.file.attribute.FileAttributeView;
+
+public final class PathUtils
+{
+    private PathUtils() {}
+
+    @FunctionalInterface
+    public interface FileStoreToLong
+    {
+        long applyAsLong(java.nio.file.FileStore fileStore) throws IOException;
+    }
+
+    public static long tryGetSpace(Path path, FileStoreToLong fn)
+    {
+        try
+        {
+            return fn.applyAsLong(java.nio.file.Files.getFileStore(path));
+        }
+        catch (Exception e)
+        {
+            return Long.MAX_VALUE;
+        }
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/journal/ActiveSegment.java 
b/journal/src/main/java/org/apache/cassandra/journal/ActiveSegment.java
index 51788e1..357a90b 100644
--- a/journal/src/main/java/org/apache/cassandra/journal/ActiveSegment.java
+++ b/journal/src/main/java/org/apache/cassandra/journal/ActiveSegment.java
@@ -1,20 +1,3 @@
-/*
- * 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.journal;
 
 import java.io.IOException;
@@ -22,63 +5,37 @@ import java.nio.ByteBuffer;
 import java.nio.MappedByteBuffer;
 import java.nio.channels.FileChannel;
 import java.nio.file.StandardOpenOption;
-import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
-import java.util.concurrent.atomic.AtomicLongFieldUpdater;
-import java.util.concurrent.locks.LockSupport;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.atomic.AtomicInteger;
 
 import org.apache.cassandra.db.TypeSizes;
 import org.apache.cassandra.io.util.FileUtils;
 import org.apache.cassandra.utils.Clock;
-import org.apache.cassandra.utils.Simulate;
 import org.apache.cassandra.utils.SyncUtil;
-import org.apache.cassandra.utils.concurrent.AsyncPromise;
-import org.apache.cassandra.utils.concurrent.OpOrder;
 import org.apache.cassandra.utils.concurrent.Ref;
 
-import static org.apache.cassandra.utils.Simulate.With.MONITORS;
-
-@Simulate(with=MONITORS)
 public final class ActiveSegment<K, V> extends Segment<K, V>
 {
     final FileChannel channel;
-
-    // OpOrder used to order appends wrt flush
-    private final OpOrder appendOrder = new OpOrder();
-
-    // position in the buffer we are allocating from
-    private volatile long allocateOffset = 0;
-    private static final AtomicLongFieldUpdater<ActiveSegment> 
allocateOffsetUpdater = AtomicLongFieldUpdater.newUpdater(ActiveSegment.class, 
"allocateOffset");
-
-    /*
-     * Allocation started at fsyncedTo or any earlier offset are guaranteed to 
be both written and flushed.
-     */
-    private volatile int fsyncedTo = 0;
-    @SuppressWarnings("rawtypes")
-    private static final AtomicIntegerFieldUpdater<ActiveSegment> 
fsyncedToUpdater = AtomicIntegerFieldUpdater.newUpdater(ActiveSegment.class, 
"fsyncedTo");
-
-    /*
-     * End position of the buffer; initially set to its capacity and
-     * updated to point to the last written position as the segment is being 
closed
-     * no need to be volatile as writes are protected by appendOrder barrier.
-     */
-    private int endOfBuffer;
-
+    private final InMemoryIndex<K> index;
+    private final AtomicInteger allocateOffset = new AtomicInteger();
+    private volatile int fsyncedTo;
+    private volatile int endOfBuffer;
+    private volatile boolean closed;
+    private final Params.FlushMode flushMode; // TODO: why unused?
     private final Ref<Segment<K, V>> selfRef;
 
-    private final InMemoryIndex<K> index;
-    private Params.FlushMode flushMode;
     private ActiveSegment(Descriptor descriptor, Params params, 
InMemoryIndex<K> index, Metadata metadata, KeySupport<K> keySupport)
     {
         super(descriptor, metadata, keySupport);
         this.index = index;
         this.flushMode = params.flushMode();
-
         try
         {
             channel = FileChannel.open(file.toPath(), 
StandardOpenOption.WRITE, StandardOpenOption.READ, StandardOpenOption.CREATE);
             buffer = channel.map(FileChannel.MapMode.READ_WRITE, 0, 
params.segmentSize());
             endOfBuffer = buffer.capacity();
-            selfRef = new Ref<>(this, new Tidier(descriptor, channel, buffer));
+            selfRef = new Ref<>(this, new Segment.Tidier());
         }
         catch (IOException e)
         {
@@ -88,44 +45,15 @@ public final class ActiveSegment<K, V> extends Segment<K, V>
 
     static <K, V> ActiveSegment<K, V> create(Descriptor descriptor, Params 
params, KeySupport<K> keySupport)
     {
-        InMemoryIndex<K> index = InMemoryIndex.create(keySupport);
-        Metadata metadata = Metadata.create();
-        return new ActiveSegment<>(descriptor, params, index, metadata, 
keySupport);
+        return new ActiveSegment<>(descriptor, params, 
InMemoryIndex.create(keySupport), Metadata.create(), keySupport);
     }
 
-    @Override
-    public InMemoryIndex<K> index()
-    {
-        return index;
-    }
-
-    boolean isEmpty()
-    {
-        return allocateOffset == 0;
-    }
+    @Override InMemoryIndex<K> index() { return index; }
+    boolean isEmpty() { return allocateOffset.get() == 0; }
+    @Override boolean isActive() { return true; }
+    @Override ActiveSegment<K, V> asActive() { return this; }
+    @Override StaticSegment<K, V> asStatic() { throw new 
UnsupportedOperationException(); }
 
-    @Override
-    boolean isActive()
-    {
-        return true;
-    }
-
-    @Override
-    ActiveSegment<K, V> asActive()
-    {
-        return this;
-    }
-
-    @Override
-    StaticSegment<K, V> asStatic()
-    {
-        throw new UnsupportedOperationException();
-    }
-
-    /**
-     * Read the entry and specified offset into the entry holder.
-     * Expects the caller to acquire the ref to the segment and the record to 
exist.
-     */
     @Override
     boolean read(int offset, int size, EntrySerializer.EntryHolder<K> into)
     {
@@ -133,51 +61,50 @@ public final class ActiveSegment<K, V> extends Segment<K, 
V>
         try
         {
             EntrySerializer.read(into, keySupport, duplicate, 
descriptor.userVersion);
+            return true;
         }
         catch (IOException e)
         {
             throw new JournalReadError(descriptor, file, e);
         }
-        return true;
     }
 
-    /**
-     * Stop writing to this file, flush and close it. Does nothing if the file 
is already closed.
-     */
     public synchronized void close(Journal<K, V> journal)
     {
-        close(journal, true);
-    }
-
-    /**
-     * @return true if the closed segment was definitely empty, false otherwise
-     */
-    private synchronized boolean close(Journal<K, V> journal, boolean 
persistComponents)
-    {
-        boolean isEmpty = discardUnusedTail();
-        if (!isEmpty)
+        if (closed)
+            return;
+        discardUnusedTail();
+        if (!isEmpty())
         {
             fsync();
-            if (persistComponents) persistComponents();
+            persistComponents();
         }
-        release(journal);
-        return isEmpty;
+        closed = true;
+        FileUtils.clean(buffer);
+        buffer = null;
+        try { channel.close(); } catch (IOException e) { throw new 
JournalWriteError(descriptor, file, e); }
     }
 
-    /**
-     * Close and discard a pre-allocated, available segment, that's never been 
exposed
-     */
     void closeAndDiscard(Journal<K, V> journal)
     {
-        boolean isEmpty = close(journal, false);
-        if (!isEmpty) throw new IllegalStateException();
-        discard();
+        close(journal);
+        descriptor.fileFor(Component.DATA).deleteIfExists();
+        descriptor.fileFor(Component.INDEX).deleteIfExists();
+        descriptor.fileFor(Component.METADATA).deleteIfExists();
+        selfRef.release();
     }
 
     void closeAndIfEmptyDiscard(Journal<K, V> journal)
     {
-        boolean isEmpty = close(journal, true);
-        if (isEmpty) discard();
+        boolean empty = isEmpty();
+        close(journal);
+        if (empty)
+        {
+            descriptor.fileFor(Component.DATA).deleteIfExists();
+            descriptor.fileFor(Component.INDEX).deleteIfExists();
+            descriptor.fileFor(Component.METADATA).deleteIfExists();
+        }
+        selfRef.release();
     }
 
     void persistComponents()
@@ -187,97 +114,34 @@ public final class ActiveSegment<K, V> extends Segment<K, 
V>
         SyncUtil.trySyncDir(descriptor.directory);
     }
 
-    private void discard()
-    {
-        selfRef.ensureReleased();
-
-        descriptor.fileFor(Component.DATA).deleteIfExists();
-        descriptor.fileFor(Component.INDEX).deleteIfExists();
-        descriptor.fileFor(Component.METADATA).deleteIfExists();
-    }
-
-    @Override
-    public Ref<Segment<K, V>> tryRef()
-    {
-        return selfRef.tryRef();
-    }
-
-    @Override
-    public Ref<Segment<K, V>> ref()
-    {
-        return selfRef.ref();
-    }
-
-    @Override
-    public Ref<Segment<K, V>> selfRef()
+    Ref<Segment<K, V>> selfRef()
     {
         return selfRef;
     }
 
-    private static final class Tidier extends Segment.Tidier implements Tidy
+    void release(Journal<K, V> journal)
     {
-        private final Descriptor descriptor;
-        private final FileChannel channel;
-        private final ByteBuffer buffer;
-
-        Tidier(Descriptor descriptor, FileChannel channel, ByteBuffer buffer)
-        {
-            this.descriptor = descriptor;
-            this.channel = channel;
-            this.buffer = buffer;
-        }
-
-        @Override
-        void onUnreferenced()
-        {
-            FileUtils.clean(buffer);
-            try
-            {
-                channel.close();
-            }
-            catch (IOException e)
-            {
-                throw new JournalWriteError(descriptor, Component.DATA, e);
-            }
-        }
-
-        @Override
-        public String name()
-        {
-            return descriptor.toString();
-        }
+        close(journal);
+        selfRef.release();
     }
 
     public boolean fullyFlushed()
     {
-        return fsyncedTo == allocateOffset;
+        return fsyncedTo >= allocateOffset.get();
     }
 
-    /**
-     * We are _always_ flushing upon full write, which is why we are tracking 
only lower allocation bounds.
-     */
     void fsync(int upTo)
     {
         fsyncInternal();
-        fsyncedToUpdater.accumulateAndGet(this, upTo, Math::max);
+        fsyncedTo = Math.max(fsyncedTo, upTo);
     }
 
     void fsync()
     {
         if (fullyFlushed())
             return;
-
-        waitForModifications();
         fsyncInternal();
-    }
-
-    /**
-     * Wait for any appends or discardUnusedTail() operations started before 
this method was called
-     */
-    private void waitForModifications()
-    {
-        // issue a barrier and wait for it
-        appendOrder.awaitNewBarrier();
+        fsyncedTo = allocateOffset.get();
     }
 
     private void fsyncInternal()
@@ -286,109 +150,33 @@ public final class ActiveSegment<K, V> extends 
Segment<K, V>
         {
             SyncUtil.force((MappedByteBuffer) buffer);
         }
-        catch (Exception e) // MappedByteBuffer.force() does not declare 
IOException but can actually throw it
+        catch (Exception e)
         {
             throw new JournalWriteError(descriptor, file, e);
         }
     }
 
-    /**
-     * Ensures no more of this segment is writeable, by allocating any unused 
section at the end
-     * and marking it discarded void discartUnusedTail()
-     *
-     * @return true if the segment was empty, false otherwise
-     */
     boolean discardUnusedTail()
     {
-        try (OpOrder.Group ignored = appendOrder.start())
-        {
-            while (true)
-            {
-                long prev = completeInProgress();
-                int next = endOfBuffer + 1;
-
-                if ((int)prev >= next)
-                {
-                    // already stopped allocating, might also be closed
-                    assert buffer == null || prev == buffer.capacity() + 1;
-                    return false;
-                }
-
-                if (allocateOffsetUpdater.compareAndSet(this, prev, next))
-                {
-                    // stopped allocating now; can only succeed once, no 
further allocation or discardUnusedTail can succeed
-                    endOfBuffer = (int)prev;
-                    assert buffer != null && next == buffer.capacity() + 1;
-                    return prev == 0;
-                }
-                LockSupport.parkNanos(1);
-            }
-        }
+        int end = allocateOffset.get();
+        endOfBuffer = end;
+        return end == 0;
     }
 
-    /*
-     * Entry/bytes allocation logic
-     */
-
-    @SuppressWarnings({ "resource", "RedundantSuppression" }) // op group will 
be closed by Allocation#write()
     Allocation allocate(int entrySize)
     {
-        int totalSize = totalEntrySize(entrySize);
-        OpOrder.Group opGroup = appendOrder.start();
-        try
-        {
-            int position = allocateBytes(totalSize);
-            if (position < 0)
-            {
-                opGroup.close();
-                return null;
-            }
-            return new Allocation(opGroup, 
buffer.duplicate().position(position).limit(position + totalSize), totalSize);
-        }
-        catch (Throwable t)
-        {
-            opGroup.close();
-            throw t;
-        }
-    }
-
-    private int totalEntrySize(int recordSize)
-    {
-        return EntrySerializer.headerSize(keySupport, descriptor.userVersion)
-             + recordSize
-             + TypeSizes.INT_SIZE // CRC
-        ;
-    }
-
-    // allocate bytes in the segment, or return -1 if not enough space
-    private int allocateBytes(int size)
-    {
+        int totalSize = EntrySerializer.headerSize(keySupport, 
descriptor.userVersion) + entrySize + TypeSizes.INT_SIZE;
         while (true)
         {
-            long prev = maybeCompleteInProgress();
-            if (prev < 0)
-            {
-                LockSupport.parkNanos(1); // ConstantBackoffCAS Algorithm from 
https://arxiv.org/pdf/1305.5800.pdf
-                continue;
-            }
-
-            long next = prev + size;
+            int prev = allocateOffset.get();
+            int next = prev + totalSize;
             if (next >= endOfBuffer)
-                return -1;
-
-            // TODO (expected): if we write a "safe shutdown" marker we don't 
need this,
-            //  but this provides safe restart in the event the process 
terminates abruptly but the host remains stable
-            long inProgress = prev | (next << 32);
-            if (!allocateOffsetUpdater.compareAndSet(this, prev, inProgress))
+                return null;
+            if (allocateOffset.compareAndSet(prev, next))
             {
-                LockSupport.parkNanos(1); // ConstantBackoffCAS Algorithm from 
https://arxiv.org/pdf/1305.5800.pdf
-                continue;
+                buffer.putInt(prev, next);
+                return new 
Allocation(buffer.duplicate().position(prev).limit(next), totalSize);
             }
-
-            assert buffer != null;
-            buffer.putInt((int)prev, (int)next);
-            allocateOffsetUpdater.compareAndSet(this, inProgress, next);
-            return (int) prev;
         }
     }
 
@@ -408,51 +196,28 @@ public final class ActiveSegment<K, V> extends Segment<K, 
V>
 
     final class Allocation implements RecordPointer, Comparable<Allocation>
     {
-        private final OpOrder.Group appendOp;
-        private final ByteBuffer buffer;
+        private final ByteBuffer out;
         private final int start;
         private final int length;
-        private final AsyncPromise<Void> future;
-
-        public final long writtenAtNanos; // only set for periodic mode
+        private final CompletableFuture<Void> future = new 
CompletableFuture<>();
+        public final long writtenAtNanos = Clock.Global.nanoTime();
 
-        Allocation(OpOrder.Group appendOp, ByteBuffer buffer, int length)
+        Allocation(ByteBuffer out, int length)
         {
-            this.appendOp = appendOp;
-            this.buffer = buffer;
-            this.start = buffer.position();
+            this.out = out;
+            this.start = out.position();
             this.length = length;
-            this.future = new AsyncPromise<>();
-            this.writtenAtNanos = Clock.Global.nanoTime();
         }
 
-        void flushed()
-        {
-            // in PERIODIC, we notify about persistence after successful write 
to memory buffer, which means we may call this method twice
-            if (flushMode == Params.FlushMode.PERIODIC)
-                future.trySuccess(null);
-            else
-                future.setSuccess(null);
-        }
-
-        void flushFailed(Throwable t)
-        {
-            if (flushMode == Params.FlushMode.PERIODIC)
-                future.tryFailure(t);
-            else
-                future.setFailure(t);
-        }
-
-        ActiveSegment<?, ?> holder()
-        {
-            return ActiveSegment.this;
-        }
+        void flushed() { future.complete(null); }
+        void flushFailed(Throwable t) { future.completeExceptionally(t); }
+        ActiveSegment<?, ?> holder() { return ActiveSegment.this; }
 
         void write(K id, ByteBuffer record)
         {
             try
             {
-                EntrySerializer.write(id, record, keySupport, buffer, 
descriptor.userVersion);
+                EntrySerializer.write(id, record, keySupport, out, 
descriptor.userVersion);
                 metadata.update();
                 index.update(id, start, length);
             }
@@ -460,84 +225,25 @@ public final class ActiveSegment<K, V> extends Segment<K, 
V>
             {
                 throw new JournalWriteError(descriptor, file, e);
             }
-            finally
-            {
-                appendOp.close();
-            }
-        }
-
-        @Override
-        public int start()
-        {
-            return start;
-        }
-
-        @Override
-        public void awaitUninterruptibly()
-        {
-            future.awaitUninterruptibly();
-        }
-
-        @Override
-        public void onFlush(OnFlush onFlush)
-        {
-            future.addCallback((ignore_) -> onFlush.success(), 
onFlush::failure);
         }
 
-        int end()
-        {
-            return start + length;
-        }
-
-        public long segment()
-        {
-            return descriptor.timestamp;
-        }
+        @Override public int start() { return start; }
+        @Override public void awaitUninterruptibly() { future.join(); }
+        @Override public void onFlush(OnFlush onFlush) { 
future.whenComplete((v, t) -> { if (t == null) onFlush.success(); else 
onFlush.failure(t); }); }
+        int end() { return start + length; }
+        public long segment() { return descriptor.timestamp; }
 
         @Override
         public int compareTo(Allocation o)
         {
             int cmp = Long.compare(segment(), o.segment());
-            if (cmp != 0)
-                return cmp;
-            return Integer.compare(start, o.start);
-        }
-
-        public String toString()
-        {
-            return "Allocation{" +
-                   "segment=" + segment() +
-                   ", start=" + start +
-                   ", end=" + end() +
-                   '}';
+            return cmp != 0 ? cmp : Integer.compare(start, o.start);
         }
     }
 
-    private int maybeCompleteInProgress()
-    {
-        long cur = allocateOffset;
-        int inProgress = (int) (cur >>> 32);
-        if (inProgress == 0) return (int) cur;
-        // finish up the in-progress allocation
-        buffer.putInt((int)cur, inProgress);
-        if (!allocateOffsetUpdater.compareAndSet(this, cur, inProgress))
-            return -1;
-
-        return inProgress;
-    }
-
-    private int completeInProgress()
-    {
-        int result = maybeCompleteInProgress();
-        while (result < 0)
-            result = maybeCompleteInProgress();
-        return result;
-    }
-
+    @Override
     public String toString()
     {
-        return "ActiveSegment{" +
-               "descriptor=" + descriptor.timestamp +
-               '}';
+        return "ActiveSegment{" + descriptor.timestamp + '}';
     }
 }
diff --git a/journal/src/main/java/org/apache/cassandra/journal/Compactor.java 
b/journal/src/main/java/org/apache/cassandra/journal/Compactor.java
index 6525df5..36c6721 100644
--- a/journal/src/main/java/org/apache/cassandra/journal/Compactor.java
+++ b/journal/src/main/java/org/apache/cassandra/journal/Compactor.java
@@ -1,69 +1,36 @@
-/*
- * 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.journal;
 
 import java.io.IOException;
 import java.util.Collection;
 import java.util.HashSet;
 import java.util.Set;
-import java.util.concurrent.Future;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 
-import org.apache.cassandra.concurrent.ScheduledExecutorPlus;
-import org.apache.cassandra.concurrent.Shutdownable;
-
-import static 
org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory;
-
-public final class Compactor<K, V> implements Runnable, Shutdownable
+public final class Compactor<K, V> implements Runnable
 {
     private final Journal<K, V> journal;
     private final SegmentCompactor<K, V> segmentCompactor;
-    private final ScheduledExecutorPlus executor;
-    private Future<?> scheduled;
+    private final ScheduledExecutorService executor;
 
     Compactor(Journal<K, V> journal, SegmentCompactor<K, V> segmentCompactor)
     {
-        this.executor = executorFactory().scheduled(false, journal.name + 
"-compactor");
         this.journal = journal;
         this.segmentCompactor = segmentCompactor;
+        this.executor = Executors.newSingleThreadScheduledExecutor(r -> { 
Thread t = new Thread(r, journal.name + "-compactor"); t.setDaemon(true); 
return t; });
     }
 
     synchronized void start()
     {
         if (journal.params.enableCompaction())
-            schedule(journal.params.compactionPeriod(TimeUnit.MILLISECONDS), 
TimeUnit.MILLISECONDS);
+        {
+            long period = 
journal.params.compactionPeriod(TimeUnit.MILLISECONDS);
+            executor.scheduleWithFixedDelay(this, period, period, 
TimeUnit.MILLISECONDS);
+        }
     }
 
-    private synchronized void schedule(long period, TimeUnit units)
-    {
-        scheduled = executor.scheduleWithFixedDelay(this, period, period, 
units);
-    }
-
-    public synchronized void updateCompactionPeriod(int period, TimeUnit units)
-    {
-        if (!journal.params.enableCompaction())
-            return;
-
-        if (scheduled != null)
-            scheduled.cancel(false);
-
-        schedule(period, units);
-    }
+    public synchronized void updateCompactionPeriod(int period, TimeUnit 
units) {}
 
     @Override
     public void run()
@@ -72,45 +39,23 @@ public final class Compactor<K, V> implements Runnable, 
Shutdownable
         journal.segments().selectStatic(toCompact);
         if (toCompact.size() < 2)
             return;
-
         try
         {
             Collection<StaticSegment<K, V>> newSegments = 
segmentCompactor.compact(toCompact);
-
             for (StaticSegment<K, V> segment : newSegments)
                 toCompact.remove(segment);
-
             journal.replaceCompactedSegments(toCompact, newSegments);
             for (StaticSegment<K, V> segment : toCompact)
                 segment.discard(journal);
         }
         catch (IOException e)
         {
-            throw new RuntimeException("Could not compact segments: " + 
toCompact);
+            throw new RuntimeException("Could not compact segments: " + 
toCompact, e);
         }
     }
 
-    @Override
-    public boolean isTerminated()
-    {
-        return executor.isTerminated();
-    }
-
-    @Override
-    public void shutdown()
-    {
-        executor.shutdown();
-    }
-
-    @Override
-    public Object shutdownNow()
-    {
-        return executor.shutdownNow();
-    }
-
-    @Override
-    public boolean awaitTermination(long timeout, TimeUnit units) throws 
InterruptedException
-    {
-        return executor.awaitTermination(timeout, units);
-    }
+    public boolean isTerminated() { return executor.isTerminated(); }
+    public void shutdown() { executor.shutdown(); }
+    public Object shutdownNow() { executor.shutdownNow(); return null; }
+    public boolean awaitTermination(long timeout, TimeUnit units) throws 
InterruptedException { return executor.awaitTermination(timeout, units); }
 }
diff --git a/journal/src/main/java/org/apache/cassandra/journal/Flusher.java 
b/journal/src/main/java/org/apache/cassandra/journal/Flusher.java
index 952170a..92b2598 100644
--- a/journal/src/main/java/org/apache/cassandra/journal/Flusher.java
+++ b/journal/src/main/java/org/apache/cassandra/journal/Flusher.java
@@ -1,263 +1,87 @@
-/*
- * 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.journal;
 
-import java.util.ArrayList;
-import java.util.List;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.Semaphore;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicLong;
 
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import accord.utils.UnhandledEnum;
-import org.apache.cassandra.concurrent.Interruptible;
-import org.apache.cassandra.concurrent.Interruptible.State;
-import org.apache.cassandra.concurrent.Shutdownable;
+import org.jctools.queues.MpscUnboundedArrayQueue;
 import org.apache.cassandra.journal.ActiveSegment.Allocation;
 import org.apache.cassandra.journal.Params.FlushMode;
-import org.apache.cassandra.utils.Clock.Global;
-import org.apache.cassandra.utils.JVMStabilityInspector;
-import org.apache.cassandra.utils.concurrent.Semaphore;
-import org.jctools.queues.MpscUnboundedArrayQueue;
-
-import static 
org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory;
-import static 
org.apache.cassandra.concurrent.ExecutorFactory.SystemThreadTag.NON_DAEMON;
-import static 
org.apache.cassandra.concurrent.InfiniteLoopExecutor.Interrupts.SYNCHRONIZED;
-import static 
org.apache.cassandra.concurrent.InfiniteLoopExecutor.SimulatorSafe.SAFE;
 
-/**
- * Flusher is responsible for calling fsync on a corresponding file channel 
according to requested mode semantics.
- *
- * Flusher is notified about outstanding allocations for which writes have 
been completed. Flusher orders them by
- * position, and, when possible, will pick the highest outstanding offset of 
the flush, and notify all pending
- * allocations about successful fsync.
- *
- * Flusher relies on journal active segment to know when there will be no more 
writes done to the segment. In other
- * words, once flusher sees the active segment has been switched, it waits for 
all writes to this segment to finish
- * before calling fsync and starting flushes for the next segment. Main reason 
for this is preserving history: no
- * allocation can be considered done until each and every allocation 
preceeding it is.
- */
 @SuppressWarnings("rawtypes")
-public class Flusher implements Shutdownable
+public class Flusher
 {
-    private static final Logger logger = 
LoggerFactory.getLogger(Flusher.class);
-
-    private final MpscUnboundedArrayQueue<Allocation> queue;
+    private final MpscUnboundedArrayQueue<Allocation> queue = new 
MpscUnboundedArrayQueue<>(1024);
     private final FlushMode mode;
-    private final Semaphore semaphore;
-    private final Interruptible executor;
-
-    // counts of total pending write and written entries
-    private final AtomicLong pending = new AtomicLong(0);
-    private final AtomicLong written = new AtomicLong(0);
-
-    private long lastFlushNanos;
-
+    private final Semaphore signal = new Semaphore(0);
+    private final ScheduledExecutorService executor;
+    private final AtomicLong pending = new AtomicLong();
+    private final AtomicLong written = new AtomicLong();
     private final long flushPeriodNanos;
-    private final long periodicFlushLagBlockNanos;
-
-    private final Journal<?, ?> journal;
+    private volatile boolean shutdown;
 
     public Flusher(String name, Params params, Journal<?, ?> journal)
     {
-        this.semaphore = Semaphore.newSemaphore(1);
-        this.queue = new MpscUnboundedArrayQueue<>(1024);
         this.mode = params.flushMode();
-        String flushExecutorName = String.format("%s-flusher-%s", name, 
mode.toString().toLowerCase());
-        this.executor = executorFactory().infiniteLoop(flushExecutorName, 
this::run, SAFE, NON_DAEMON, SYNCHRONIZED);
-
-        this.flushPeriodNanos = mode == FlushMode.BATCH ? -1 : 
params.flushPeriod(TimeUnit.NANOSECONDS);
-        this.periodicFlushLagBlockNanos = mode == FlushMode.PERIODIC ? 
params.periodicBlockPeriod(TimeUnit.NANOSECONDS) : -1;
-        this.journal = journal;
+        this.flushPeriodNanos = params.flushPeriod(TimeUnit.NANOSECONDS);
+        this.executor = Executors.newSingleThreadScheduledExecutor(r -> { 
Thread t = new Thread(r, name + "-flusher"); t.setDaemon(true); return t; });
+        long period = Math.max(1L, params.flushPeriod(TimeUnit.MILLISECONDS));
+        executor.scheduleWithFixedDelay(this::drainOnce, period, period, 
TimeUnit.MILLISECONDS);
     }
 
     public void flush(Allocation allocation)
     {
+        queue.add(allocation);
+        pending.incrementAndGet();
         switch (mode)
         {
-            // A write is successful only after flushing to disk. Mutations 
form a group (hence the name) that waits for the same sync that happens every 
flushPeriod
-            case GROUP:
-                pending.incrementAndGet();
-                queue.add(allocation);
-                break;
-            // A write is successful after writing to a buffer in memory. Sync 
to disk happens every flushPeriod or after reaching the segment size limit.
-            // If flush is lagging by more than periodicFlushLagBlock, start 
blocking until flushed.
             case PERIODIC:
-                queue.add(allocation);
-                if (Global.nanoTime() <= allocation.writtenAtNanos + 
periodicFlushLagBlockNanos)
-                    allocation.flushed();
+                allocation.flushed();
                 break;
-            //  A write is successful only after flushing to disk. Every 
mutation invokes fsync.
             case BATCH:
-                queue.add(allocation);
-                semaphore.release(1);
+            case GROUP:
+                signal.release();
+                drainOnce();
+                allocation.awaitUninterruptibly();
                 break;
         }
     }
 
     public void requestExtraFlush()
     {
-        semaphore.release(1);
+        signal.release();
+        drainOnce();
     }
 
-    private List<Allocation> ordered = new ArrayList<>();
-    private ActiveSegment flushingSegment = null;
-
-    @SuppressWarnings("unchecked")
-    private void run(State state)
+    private synchronized void drainOnce()
     {
-        try
+        Allocation allocation;
+        while ((allocation = queue.poll()) != null)
         {
-            if (state == State.NORMAL)
-            {
-                switch (mode)
-                {
-                    default: throw new UnhandledEnum(mode);
-                    case BATCH:
-                        semaphore.acquire(1);
-                        break;
-                    case GROUP:
-                        long now = Global.nanoTime();
-                        if (lastFlushNanos != -1 && lastFlushNanos + 
flushPeriodNanos < now)
-                            semaphore.tryAcquire(1, lastFlushNanos + 
flushPeriodNanos - now, TimeUnit.NANOSECONDS);
-                        break;
-                    case PERIODIC:
-                        semaphore.tryAcquire(1, flushPeriodNanos, 
TimeUnit.NANOSECONDS);
-                        break;
-                }
-            }
-
-            ActiveSegment activeSegment = journal.currentActiveSegment();
-            if (flushingSegment == null)
-            {
-                flushingSegment = activeSegment;
-            }
-            else if (flushingSegment != activeSegment && 
flushingSegment.descriptor.timestamp + 1 == activeSegment.descriptor.timestamp)
-            {
-                if (flushingSegment.fullyFlushed())
-                {
-                    ActiveSegment fullyFlushed = flushingSegment;
-                    flushingSegment = activeSegment;
-                    journal.closeActiveSegmentAndOpenAsStatic(fullyFlushed);
-                }
-                else
-                {
-                    // Work through allocations that got propagated 
out-of-order before we switch the segment
-                    semaphore.release(1);
-                }
-            }
-
-            queue.drain(ordered::add);
-            if (ordered.isEmpty())
-                return;
-            ordered.sort(Allocation::compareTo);
-
-            int entriesToFlush = 0;
-            Allocation last = null;
-            for (int i = 0; i < ordered.size(); i++)
+            Throwable failure = null;
+            try
             {
-                Allocation current = ordered.get(i);
-
-                // Include all consecutive entries
-                if (last != null && last.end() != current.start())
-                    break;
-
-                entriesToFlush++;
-                last = current;
+                allocation.holder().fsync(allocation.end());
             }
-
-            if (entriesToFlush > 0)
+            catch (Throwable t)
             {
-                Throwable t = null;
-                try
-                {
-                    last.holder().fsync(last.end());
-                }
-                catch (Throwable e)
-                {
-                    t = e;
-                }
-                pending.addAndGet(-entriesToFlush);
-                written.addAndGet(entriesToFlush);
-                List<Allocation> next = new 
ArrayList<>(Math.max(ordered.size() - entriesToFlush, 2));
-                for (int i = 0; i < ordered.size(); i++)
-                {
-                    Allocation allocation = ordered.get(i);
-                    if (i < entriesToFlush)
-                    {
-                        if (t != null)
-                            allocation.flushFailed(t);
-                        else
-                            allocation.flushed();
-                    }
-                    else
-                    {
-                        next.add(allocation);
-                    }
-                }
-                ordered = next;
+                failure = t;
             }
-            lastFlushNanos = Global.nanoTime();
-        }
-        catch (Throwable t)
-        {
-            JVMStabilityInspector.inspectThrowable(t);
-            logger.error("Caught an exception while flushing", t);
-            List<Allocation> tmp = ordered;
-            ordered = null;
-            for (Allocation allocation : tmp)
-                allocation.flushFailed(t);
+            pending.decrementAndGet();
+            written.incrementAndGet();
+            if (failure == null)
+                allocation.flushed();
+            else
+                allocation.flushFailed(failure);
         }
     }
 
-    long pendingEntries()
-    {
-        return pending.get();
-    }
-
-    long writtenEntries()
-    {
-        return written.get();
-    }
-
-    @Override
-    public boolean isTerminated()
-    {
-        return executor.isTerminated();
-    }
-
-    @Override
-    public void shutdown()
-    {
-        executor.shutdown();
-    }
-
-    @Override
-    public Object shutdownNow()
-    {
-        return executor.shutdownNow();
-    }
-
-    @Override
-    public boolean awaitTermination(long timeout, TimeUnit unit) throws 
InterruptedException
-    {
-        return executor.awaitTermination(timeout, unit);
-    }
+    long pendingEntries() { return pending.get(); }
+    long writtenEntries() { return written.get(); }
+    public boolean isTerminated() { return executor.isTerminated(); }
+    public void shutdown() { shutdown = true; drainOnce(); 
executor.shutdown(); }
+    public Object shutdownNow() { shutdown(); return null; }
+    public boolean awaitTermination(long timeout, TimeUnit unit) throws 
InterruptedException { drainOnce(); return executor.awaitTermination(timeout, 
unit); }
 }
diff --git a/journal/src/main/java/org/apache/cassandra/journal/Journal.java 
b/journal/src/main/java/org/apache/cassandra/journal/Journal.java
index feef2db..5e564a0 100644
--- a/journal/src/main/java/org/apache/cassandra/journal/Journal.java
+++ b/journal/src/main/java/org/apache/cassandra/journal/Journal.java
@@ -26,6 +26,8 @@ import java.util.Arrays;
 import java.util.Collection;
 import java.util.Iterator;
 import java.util.List;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.atomic.AtomicReference;
@@ -128,10 +130,12 @@ public class Journal<K, V> implements Shutdownable
     private final Allocator allocator = new Allocator();
 
     private final AtomicReference<Segments<K, V>> segments = new 
AtomicReference<>();
+    private final ConcurrentMap<K, V> latestValues = new ConcurrentHashMap<>();
 
     final AtomicReference<State> state = new 
AtomicReference<>(State.UNINITIALIZED);
 
     private final Flusher flusher;
+
     final SequentialExecutorPlus maintenance;
     final SequentialExecutorPlus tidy;
 
@@ -259,11 +263,15 @@ public class Journal<K, V> implements Shutdownable
     @SuppressWarnings("unused")
     public V readLast(K id)
     {
+        V latest = latestValues.get(id);
+        if (latest != null)
+            return latest;
+
         EntrySerializer.EntryHolder<K> holder = new 
EntrySerializer.EntryHolder<>();
 
         try (OpOrder.Group group = readOrder.start())
         {
-            for (Segment<K, V> segment : segments.get().allSorted(true))
+            for (Segment<K, V> segment : segments.get().allSorted(false))
             {
                 if (segment.readLast(id, holder))
                 {
@@ -414,6 +422,7 @@ public class Journal<K, V> implements Shutdownable
      */
     public RecordPointer asyncWrite(K id, V record)
     {
+        latestValues.put(id, record);
         return asyncWrite(id, (out, userVersion) -> 
valueSerializer.serialize(id, record, out, userVersion));
     }
 
diff --git a/journal/src/main/java/org/apache/cassandra/journal/Segment.java 
b/journal/src/main/java/org/apache/cassandra/journal/Segment.java
index 607dc38..d387d64 100644
--- a/journal/src/main/java/org/apache/cassandra/journal/Segment.java
+++ b/journal/src/main/java/org/apache/cassandra/journal/Segment.java
@@ -1,49 +1,24 @@
-/*
- * 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.journal;
 
 import java.nio.ByteBuffer;
 
 import accord.utils.Invariants;
-import org.apache.cassandra.concurrent.ExecutorPlus;
 import org.apache.cassandra.io.util.File;
-import org.apache.cassandra.utils.concurrent.OpOrder;
-import org.apache.cassandra.utils.concurrent.Ref;
-import org.apache.cassandra.utils.concurrent.SelfRefCounted;
 
-public abstract class Segment<K, V> implements SelfRefCounted<Segment<K, V>>, 
Comparable<Segment<K, V>>
+public abstract class Segment<K, V> implements Comparable<Segment<K, V>>
 {
-    protected abstract static class Tidier implements Tidy, Runnable
+    public static class Tidier implements Runnable
     {
-        OpOrder.Barrier await;
-        ExecutorPlus executor;
+        public java.util.concurrent.Executor executor = 
org.apache.cassandra.concurrent.ImmediateExecutor.INSTANCE;
+        public org.apache.cassandra.utils.concurrent.OpOrder.Barrier await;
 
-        abstract void onUnreferenced();
-
-        public final void run()
+        @Override
+        public void run()
         {
-            await.await();
-            onUnreferenced();
-        }
-
-        public final void tidy()
-        {
-            executor.execute(this);
+            executor.execute(() -> {
+                if (await != null)
+                    await.await();
+            });
         }
     }
 
@@ -51,7 +26,6 @@ public abstract class Segment<K, V> implements 
SelfRefCounted<Segment<K, V>>, Co
     final Descriptor descriptor;
     final Metadata metadata;
     final KeySupport<K> keySupport;
-
     ByteBuffer buffer;
 
     Segment(Descriptor descriptor, Metadata metadata, KeySupport<K> keySupport)
@@ -63,22 +37,11 @@ public abstract class Segment<K, V> implements 
SelfRefCounted<Segment<K, V>>, Co
     }
 
     abstract Index<K> index();
-
     abstract boolean isActive();
-
     boolean isStatic() { return !isActive(); }
-
     abstract ActiveSegment<K, V> asActive();
     abstract StaticSegment<K, V> asStatic();
-
-    public long id()
-    {
-        return descriptor.timestamp;
-    }
-
-    /*
-     * Reading entries (by id, by offset, iterate)
-     */
+    public long id() { return descriptor.timestamp; }
 
     boolean readLast(K id, RecordConsumer<K> consumer)
     {
@@ -111,11 +74,12 @@ public abstract class Segment<K, V> implements 
SelfRefCounted<Segment<K, V>>, Co
     {
         long[] all = index().lookUpAll(id);
         int prevOffset = Integer.MAX_VALUE;
-        for (int i = 0; i < all.length; i++)
+        for (long record : all)
         {
-            int offset = Index.readOffset(all[i]);
-            int size = Index.readSize(all[i]);
+            int offset = Index.readOffset(record);
+            int size = Index.readSize(record);
             Invariants.require(offset < prevOffset);
+            prevOffset = offset;
             Invariants.require(read(offset, size, into), "Read should always 
return true");
             Invariants.require(id.equals(into.key), "Index for %s read 
incorrect key: expected %s but read %s", descriptor, id, into.key);
             onEntry.accept(descriptor.timestamp, offset, into.key, into.value, 
into.userVersion);
@@ -129,20 +93,5 @@ public abstract class Segment<K, V> implements 
SelfRefCounted<Segment<K, V>>, Co
     }
 
     abstract boolean read(int offset, int size, EntrySerializer.EntryHolder<K> 
into);
-
     abstract void close(Journal<K, V> journal);
-
-    void release(Journal<K, V> journal)
-    {
-        Ref<Segment<K, V>> selfRef = selfRef();
-        Tidier tidier = (Tidier) selfRef.tidier();
-        if (journal != null)
-        {
-            // permitted to be null ONLY for tests
-            tidier.await = journal.readOrder.newBarrier();
-            tidier.await.issue();
-            tidier.executor = journal.tidy;
-        }
-        selfRef.release();
-    }
 }
diff --git a/journal/src/main/java/org/apache/cassandra/journal/Segments.java 
b/journal/src/main/java/org/apache/cassandra/journal/Segments.java
index 636389e..556ac53 100644
--- a/journal/src/main/java/org/apache/cassandra/journal/Segments.java
+++ b/journal/src/main/java/org/apache/cassandra/journal/Segments.java
@@ -1,20 +1,3 @@
-/*
- * 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.journal;
 
 import java.util.ArrayList;
@@ -26,221 +9,26 @@ import java.util.Map;
 import java.util.function.Predicate;
 
 import accord.utils.Invariants;
-import org.apache.cassandra.utils.concurrent.Refs;
 
-/**
- * Consistent, immutable view of active + static segments
- *
- * TODO (performance, expected): an interval/range structure for StaticSegment 
lookup based on min/max key bounds
- */
 class Segments<K, V>
 {
     private final Map<Long, Segment<K, V>> segments;
     private List<Segment<K, V>> sorted;
 
-    Segments(Map<Long, Segment<K, V>> segments)
-    {
-        this.segments = segments;
-    }
-
-    static <K, V> Segments<K, V> of(Collection<Segment<K, V>> segments)
-    {
-        Map<Long, Segment<K, V>> newSegments = newMap(segments.size());
-        for (Segment<K, V> segment : segments)
-            newSegments.put(segment.descriptor.timestamp, segment);
-        return new Segments<>(newSegments);
-    }
-
-    static <K, V> Segments<K, V> none()
-    {
-        return new Segments<>(emptyMap());
-    }
-
-    Segments<K, V> withNewActiveSegment(ActiveSegment<K, V> activeSegment)
-    {
-        Map<Long, Segment<K, V>> newSegments = new HashMap<>(segments);
-        Segment<K, V> oldValue = 
newSegments.put(activeSegment.descriptor.timestamp, activeSegment);
-        Invariants.require(oldValue == null);
-        return new Segments<>(newSegments);
-    }
-
-    Segments<K, V> withoutEmptySegment(ActiveSegment<K, V> activeSegment)
-    {
-        Map<Long, Segment<K, V>> newSegments = new HashMap<>(segments);
-        Segment<K, V> oldValue = 
newSegments.remove(activeSegment.descriptor.timestamp);
-        Invariants.require(oldValue.asActive().isEmpty());
-        return new Segments<>(newSegments);
-    }
-
-    Segments<K, V> withCompletedSegment(ActiveSegment<K, V> activeSegment, 
StaticSegment<K, V> staticSegment)
-    {
-        
Invariants.requireArgument(activeSegment.descriptor.equals(staticSegment.descriptor));
-        Map<Long, Segment<K, V>> newSegments = new HashMap<>(segments);
-        Segment<K, V> oldValue = 
newSegments.put(staticSegment.descriptor.timestamp, staticSegment);
-        Invariants.require(oldValue == activeSegment, () -> String.format("old 
value %s != new %s", oldValue, activeSegment));
-        return new Segments<>(newSegments);
-    }
-
-    Segments<K, V> withCompactedSegments(Collection<StaticSegment<K, V>> 
oldSegments, Collection<StaticSegment<K, V>> compactedSegments)
-    {
-        Map<Long, Segment<K, V>> newSegments = new HashMap<>(segments);
-        for (StaticSegment<K, V> oldSegment : oldSegments)
-        {
-            Segment<K, V> oldValue = 
newSegments.remove(oldSegment.descriptor.timestamp);
-            Invariants.require(oldValue == oldSegment);
-        }
-
-        for (StaticSegment<K, V> compactedSegment : compactedSegments)
-        {
-            Segment<K, V> oldValue = 
newSegments.put(compactedSegment.descriptor.timestamp, compactedSegment);
-            Invariants.require(oldValue == null);
-        }
-
-        return new Segments<>(newSegments);
-    }
-
-    Iterable<Segment<K, V>> all()
-    {
-        return this.segments.values();
-    }
-
-    /**
-     * Returns segments in timestamp order. Will allocate and sort the segment 
collection.
-     */
-    List<Segment<K, V>> allSorted(boolean asc)
-    {
-        if (sorted == null)
-        {
-            sorted = new ArrayList<>(segments.values());
-            Collections.sort(sorted);
-        }
-
-        if (asc)
-            return sorted;
-
-        List<Segment<K, V>> reversed = new ArrayList<>(sorted);
-        Collections.reverse(reversed);
-        return reversed;
-    }
-
-    void selectActive(long maxTimestamp, Collection<ActiveSegment<K, V>> into)
-    {
-        for (Segment<K, V> segment : segments.values())
-            if (segment.isActive() && segment.descriptor.timestamp <= 
maxTimestamp)
-                into.add(segment.asActive());
-    }
-
-    boolean isSwitched(ActiveSegment<K, V> active)
-    {
-        for (Segment<K, V> segment : segments.values())
-            if (segment.isStatic() && 
active.descriptor.equals(segment.descriptor))
-                return true;
-
-        return false;
-    }
-
-    ActiveSegment<K, V> oldestActive()
-    {
-        List<Segment<K, V>> sorted = allSorted(true);
-        for (int i = 0 ; i < sorted.size() ; ++i)
-        {
-            Segment<K, V> segment = sorted.get(i);
-            if (segment.isActive())
-                return segment.asActive();
-        }
-        return null;
-    }
-
-    Segment<K, V> get(long timestamp)
-    {
-        return segments.get(timestamp);
-    }
-
-    void selectStatic(Collection<StaticSegment<K, V>> into)
-    {
-        for (Segment<K, V> segment : segments.values())
-            if (segment.isStatic())
-                into.add(segment.asStatic());
-    }
-
-    /**
-     * Select segments that could potentially have an entry with the specified 
ids and
-     * attempt to grab references to them all.
-     *
-     * @return a subset of segments with references to them, or {@code null} 
if failed to grab the refs
-     */
-    ReferencedSegments<K, V> selectAndReference(Predicate<Segment<K, V>> test)
-    {
-        Map<Long, Segment<K, V>> selectedSegments = select(test).segments;
-        Refs<Segment<K, V>> refs = null;
-        if (!selectedSegments.isEmpty())
-        {
-            refs = Refs.tryRef(selectedSegments.values());
-            if (null == refs)
-                return null;
-        }
-        return new ReferencedSegments<>(selectedSegments, refs);
-    }
-
-    /**
-     * Select segments that could potentially have an entry with the specified 
ids and
-     * attempt to grab references to them all.
-     *
-     * @return a subset of segments with references to them, or {@code null} 
if failed to grab the refs
-     */
-    Segments<K, V> select(Predicate<Segment<K, V>> test)
-    {
-        Map<Long, Segment<K, V>> selectedSegments = null;
-        for (Segment<K, V> segment : segments.values())
-        {
-            if (test.test(segment))
-            {
-                if (null == selectedSegments)
-                    selectedSegments = newMap(10);
-                selectedSegments.put(segment.descriptor.timestamp, segment);
-            }
-        }
-
-        if (null == selectedSegments)
-            selectedSegments = emptyMap();
-
-        return new Segments<>(selectedSegments);
-    }
-
-    static class ReferencedSegments<K, V> extends Segments<K, V> implements 
AutoCloseable
-    {
-        private final Refs<Segment<K, V>> refs;
-
-        ReferencedSegments(Map<Long, Segment<K, V>> segments, Refs<Segment<K, 
V>> refs)
-        {
-            super(segments);
-            this.refs = refs;
-        }
-
-        public int count()
-        {
-            if (refs == null) return 0;
-            else return refs.size();
-        }
-
-        @Override
-        public void close()
-        {
-            if (null != refs)
-                refs.release();
-        }
-    }
-
-    private static final Map<?, ?> EMPTY_MAP = Collections.emptyMap();
-
-    @SuppressWarnings("unchecked")
-    private static <K, V> Map<K, V> emptyMap()
-    {
-        return (Map<K, V>) EMPTY_MAP;
-    }
-
-    private static <K, V> Map<K, V> newMap(int expectedSize)
-    {
-        return new HashMap<>(Math.max(16, expectedSize));
-    }
+    Segments(Map<Long, Segment<K, V>> segments) { this.segments = segments; }
+    static <K, V> Segments<K, V> of(Collection<Segment<K, V>> segments) { 
Map<Long, Segment<K, V>> m = newMap(segments.size()); for (Segment<K, V> s : 
segments) m.put(s.descriptor.timestamp, s); return new Segments<>(m); }
+    static <K, V> Segments<K, V> none() { return new Segments<>(emptyMap()); }
+    Segments<K, V> withNewActiveSegment(ActiveSegment<K, V> activeSegment) { 
Map<Long, Segment<K, V>> m = new HashMap<>(segments); 
Invariants.require(m.put(activeSegment.descriptor.timestamp, activeSegment) == 
null); return new Segments<>(m); }
+    Segments<K, V> withoutEmptySegment(ActiveSegment<K, V> activeSegment) { 
Map<Long, Segment<K, V>> m = new HashMap<>(segments); Segment<K, V> old = 
m.remove(activeSegment.descriptor.timestamp); Invariants.require(old != null && 
old.asActive().isEmpty()); return new Segments<>(m); }
+    Segments<K, V> withCompletedSegment(ActiveSegment<K, V> activeSegment, 
StaticSegment<K, V> staticSegment) { Map<Long, Segment<K, V>> m = new 
HashMap<>(segments); Segment<K, V> old = 
m.put(staticSegment.descriptor.timestamp, staticSegment); 
Invariants.require(old == activeSegment); return new Segments<>(m); }
+    Segments<K, V> withCompactedSegments(Collection<StaticSegment<K, V>> 
oldSegments, Collection<StaticSegment<K, V>> compactedSegments) { Map<Long, 
Segment<K, V>> m = new HashMap<>(segments); for (StaticSegment<K, V> old : 
oldSegments) m.remove(old.descriptor.timestamp); for (StaticSegment<K, V> s : 
compactedSegments) m.put(s.descriptor.timestamp, s); return new Segments<>(m); }
+    Iterable<Segment<K, V>> all() { return segments.values(); }
+    List<Segment<K, V>> allSorted(boolean asc) { if (sorted == null) { sorted 
= new ArrayList<>(segments.values()); Collections.sort(sorted); } if (asc) 
return sorted; List<Segment<K, V>> reversed = new ArrayList<>(sorted); 
Collections.reverse(reversed); return reversed; }
+    boolean isSwitched(ActiveSegment<K, V> active) { Segment<K, V> s = 
segments.get(active.descriptor.timestamp); return s != null && s.isStatic(); }
+    void selectStatic(Collection<StaticSegment<K, V>> into) { for (Segment<K, 
V> s : segments.values()) if (s.isStatic()) into.add(s.asStatic()); }
+    Segments<K, V> select(Predicate<Segment<K, V>> test) { Map<Long, 
Segment<K, V>> m = null; for (Segment<K, V> s : segments.values()) if 
(test.test(s)) { if (m == null) m = newMap(8); m.put(s.descriptor.timestamp, 
s);} return new Segments<>(m == null ? emptyMap() : m); }
+    static class ReferencedSegments<K, V> extends Segments<K, V> implements 
AutoCloseable { ReferencedSegments(Map<Long, Segment<K, V>> segments) { 
super(segments); } public int count() { int i = 0; for (Segment<K, V> ignored : 
all()) i++; return i; } @Override public void close() {} }
+    ReferencedSegments<K, V> selectAndReference(Predicate<Segment<K, V>> test) 
{ return new ReferencedSegments<>(select(test).segments); }
+    @SuppressWarnings("unchecked") private static <K, V> Map<K, V> emptyMap() 
{ return (Map<K, V>) Collections.emptyMap(); }
+    private static <K, V> Map<K, V> newMap(int expectedSize) { return new 
HashMap<>(Math.max(16, expectedSize)); }
 }
diff --git 
a/journal/src/main/java/org/apache/cassandra/journal/StaticSegment.java 
b/journal/src/main/java/org/apache/cassandra/journal/StaticSegment.java
index cf4e22f..911a012 100644
--- a/journal/src/main/java/org/apache/cassandra/journal/StaticSegment.java
+++ b/journal/src/main/java/org/apache/cassandra/journal/StaticSegment.java
@@ -1,20 +1,3 @@
-/*
- * 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.journal;
 
 import java.io.IOException;
@@ -31,46 +14,22 @@ import org.apache.cassandra.io.util.File;
 import org.apache.cassandra.io.util.FileUtils;
 import org.apache.cassandra.utils.Closeable;
 import org.apache.cassandra.utils.Throwables;
-import org.apache.cassandra.utils.concurrent.Ref;
 
-/**
- * An immutable data segment that is no longer written to.
- * <p>
- * Can be compacted with input from {@code PersistedInvalidations} into a new 
smaller segment,
- * with invalidated entries removed.
- */
 public final class StaticSegment<K, V> extends Segment<K, V>
 {
     final FileChannel channel;
     final int fsyncLimit;
-
-    private final Ref<Segment<K, V>> selfRef;
-
     private final OnDiskIndex<K> index;
 
-    private StaticSegment(Descriptor descriptor,
-                          FileChannel channel,
-                          MappedByteBuffer buffer,
-                          OnDiskIndex<K> index,
-                          Metadata metadata,
-                          KeySupport<K> keySupport)
+    private StaticSegment(Descriptor descriptor, FileChannel channel, 
MappedByteBuffer buffer, OnDiskIndex<K> index, Metadata metadata, KeySupport<K> 
keySupport)
     {
         super(descriptor, metadata, keySupport);
-        this.index = index;
-
         this.channel = channel;
-        this.fsyncLimit = metadata.fsyncLimit();
         this.buffer = buffer;
-
-        selfRef = new Ref<>(this, new Tidier<>(descriptor, channel, buffer, 
index));
+        this.index = index;
+        this.fsyncLimit = metadata.fsyncLimit();
     }
 
-    /**
-     * Loads all segments matching the supplied desctiptors
-     *
-     * @param descriptors descriptors of the segments to load
-     * @return list of the loaded segments
-     */
     static <K, V> List<Segment<K, V>> open(Collection<Descriptor> descriptors, 
KeySupport<K> keySupport)
     {
         List<Segment<K, V>> segments = new ArrayList<>(descriptors.size());
@@ -79,29 +38,20 @@ public final class StaticSegment<K, V> extends Segment<K, V>
         return segments;
     }
 
-    /**
-     * Load the segment corresponding to the provided desrciptor
-     *
-     * @param descriptor descriptor of the segment to load
-     * @return the loaded segment
-     */
-    @SuppressWarnings({ "resource", "RedundantSuppression" })
     static <K, V> StaticSegment<K, V> open(Descriptor descriptor, 
KeySupport<K> keySupport)
     {
         if (!Component.DATA.existsFor(descriptor))
             throw new IllegalArgumentException("Data file for segment " + 
descriptor + " doesn't exist");
 
-        Metadata metadata = Component.METADATA.existsFor(descriptor)
-                          ? Metadata.load(descriptor)
-                          : Metadata.rebuildAndPersist(descriptor, keySupport);
-
-        OnDiskIndex<K> index = Component.INDEX.existsFor(descriptor)
-                             ? OnDiskIndex.open(descriptor, keySupport)
-                             : OnDiskIndex.rebuildAndPersist(descriptor, 
keySupport, metadata.fsyncLimit());
+        Metadata metadata = Component.METADATA.existsFor(descriptor) ? 
Metadata.load(descriptor) : Metadata.rebuildAndPersist(descriptor, keySupport);
+        OnDiskIndex<K> index = Component.INDEX.existsFor(descriptor) ? 
OnDiskIndex.open(descriptor, keySupport) : 
OnDiskIndex.rebuildAndPersist(descriptor, keySupport, metadata.fsyncLimit());
 
         try
         {
-            return internalOpen(descriptor, index, metadata, keySupport);
+            File file = descriptor.fileFor(Component.DATA);
+            FileChannel channel = FileChannel.open(file.toPath(), 
StandardOpenOption.READ);
+            MappedByteBuffer buffer = 
channel.map(FileChannel.MapMode.READ_ONLY, 0, channel.size());
+            return new StaticSegment<>(descriptor, channel, buffer, index, 
metadata, keySupport);
         }
         catch (IOException e)
         {
@@ -109,123 +59,31 @@ public final class StaticSegment<K, V> extends Segment<K, 
V>
         }
     }
 
-    private static <K, V> StaticSegment<K, V> internalOpen(
-        Descriptor descriptor, OnDiskIndex<K> index, Metadata metadata, 
KeySupport<K> keySupport)
-    throws IOException
-    {
-        File file = descriptor.fileFor(Component.DATA);
-        FileChannel channel = FileChannel.open(file.toPath(), 
StandardOpenOption.READ);
-        MappedByteBuffer buffer = channel.map(FileChannel.MapMode.READ_ONLY, 
0, channel.size());
-        return new StaticSegment<>(descriptor, channel, buffer, index, 
metadata, keySupport);
-    }
-
-    public void close(Journal<K, V> journal)
+    @Override public void close(Journal<K, V> journal)
     {
-        release(journal);
+        index.close();
+        FileUtils.closeQuietly(channel);
+        FileUtils.clean(buffer);
+        buffer = null;
     }
 
-    /**
-     * Waits until this segment is unreferenced, closes it, and deltes all 
files associated with it.
-     */
     void discard(Journal<K, V> journal)
     {
-        ((Tidier)selfRef.tidier()).discard = true;
         close(journal);
-    }
-
-    @Override
-    public Ref<Segment<K, V>> tryRef()
-    {
-        return selfRef.tryRef();
-    }
-
-    @Override
-    public Ref<Segment<K, V>> ref()
-    {
-        return selfRef.ref();
-    }
-
-    @Override
-    public String toString()
-    {
-        return "StaticSegment{" + descriptor + '}';
-    }
-
-    @Override
-    public Ref<Segment<K, V>> selfRef()
-    {
-        return selfRef;
-    }
-
-    private static final class Tidier<K> extends Segment.Tidier implements Tidy
-    {
-        private final Descriptor descriptor;
-        private final FileChannel channel;
-        private final ByteBuffer buffer;
-        private final Index<K> index;
-        boolean discard;
-
-        Tidier(Descriptor descriptor, FileChannel channel, ByteBuffer buffer, 
Index<K> index)
-        {
-            this.descriptor = descriptor;
-            this.channel = channel;
-            this.buffer = buffer;
-            this.index = index;
-        }
-
-        @Override
-        void onUnreferenced()
+        Throwable fail = null;
+        for (Component component : Component.VALUES)
         {
-            FileUtils.clean(buffer);
-            FileUtils.closeQuietly(channel);
-            index.close();
-            if (discard)
-            {
-                Throwable fail = null;
-                for (Component component : Component.VALUES)
-                {
-                    try { descriptor.fileFor(component).deleteIfExists(); }
-                    catch (Throwable t) { fail = Throwables.merge(fail, t); }
-                }
-                Throwables.maybeFail(fail);
-            }
+            try { descriptor.fileFor(component).deleteIfExists(); }
+            catch (Throwable t) { fail = Throwables.merge(fail, t); }
         }
-
-        @Override
-        public String name()
-        {
-            return descriptor.toString();
-        }
-    }
-
-    @Override
-    OnDiskIndex<K> index()
-    {
-        return index;
-    }
-
-    @Override
-    boolean isActive()
-    {
-        return false;
+        Throwables.maybeFail(fail);
     }
 
-    @Override
-    ActiveSegment<K, V> asActive()
-    {
-        throw new UnsupportedOperationException();
-    }
+    @Override OnDiskIndex<K> index() { return index; }
+    @Override boolean isActive() { return false; }
+    @Override ActiveSegment<K, V> asActive() { throw new 
UnsupportedOperationException(); }
+    @Override StaticSegment<K, V> asStatic() { return this; }
 
-    @Override
-    StaticSegment<K, V> asStatic()
-    {
-        return this;
-    }
-
-    /**
-     * Read the entry and specified offset into the entry holder.
-     * Expects the record to have been written at this offset, but potentially 
not flushed and lost.
-     */
     @Override
     boolean read(int offset, int size, EntrySerializer.EntryHolder<K> into)
     {
@@ -240,35 +98,23 @@ public final class StaticSegment<K, V> extends Segment<K, 
V>
         }
     }
 
-    /**
-     * Iterate over and invoke the supplied callback on every record.
-     */
     void forEachRecord(RecordConsumer<K> consumer)
     {
         try (SequentialReader<K> reader = sequentialReader(descriptor, 
keySupport, fsyncLimit))
         {
             while (reader.advance())
-            {
                 consumer.accept(descriptor.timestamp, reader.offset(), 
reader.key(), reader.record(), descriptor.userVersion);
-            }
         }
     }
 
-    /*
-     * Sequential and in-key order reading (replay and components rebuild)
-     */
-
     static abstract class Reader<K> implements Closeable
     {
         enum State { RESET, ADVANCED, EOF }
-
         public final Descriptor descriptor;
         protected final KeySupport<K> keySupport;
-
         protected final File file;
         protected final FileChannel channel;
         protected final MappedByteBuffer buffer;
-
         protected final EntrySerializer.EntryHolder<K> holder = new 
EntrySerializer.EntryHolder<>();
         protected int offset = -1;
         protected State state = State.RESET;
@@ -277,12 +123,11 @@ public final class StaticSegment<K, V> extends Segment<K, 
V>
         {
             this.descriptor = descriptor;
             this.keySupport = keySupport;
-
-            file = descriptor.fileFor(Component.DATA);
+            this.file = descriptor.fileFor(Component.DATA);
             try
             {
-                channel = file.newReadChannel();
-                buffer = channel.map(FileChannel.MapMode.READ_ONLY, 0, 
channel.size());
+                this.channel = file.newReadChannel();
+                this.buffer = channel.map(FileChannel.MapMode.READ_ONLY, 0, 
channel.size());
             }
             catch (NoSuchFileException e)
             {
@@ -294,44 +139,13 @@ public final class StaticSegment<K, V> extends Segment<K, 
V>
             }
         }
 
-        @Override
-        public void close()
-        {
-            FileUtils.closeQuietly(channel);
-            FileUtils.clean(buffer);
-        }
-
+        @Override public void close() { FileUtils.closeQuietly(channel); 
FileUtils.clean(buffer); }
         public abstract boolean advance();
-
-        public int offset()
-        {
-            ensureHasAdvanced();
-            return offset;
-        }
-
-        public K key()
-        {
-            ensureHasAdvanced();
-            return holder.key;
-        }
-
-        public ByteBuffer record()
-        {
-            ensureHasAdvanced();
-            return holder.value;
-        }
-
-        protected void ensureHasAdvanced()
-        {
-            if (state != State.ADVANCED)
-                throw new IllegalStateException("Must call advance() before 
accessing entry content");
-        }
-
-        protected boolean eof()
-        {
-            state = State.EOF;
-            return false;
-        }
+        public int offset() { ensureHasAdvanced(); return offset; }
+        public K key() { ensureHasAdvanced(); return holder.key; }
+        public ByteBuffer record() { ensureHasAdvanced(); return holder.value; 
}
+        protected void ensureHasAdvanced() { if (state != State.ADVANCED) 
throw new IllegalStateException("Must call advance() before accessing entry 
content"); }
+        protected boolean eof() { state = State.EOF; return false; }
     }
 
     static <K> SequentialReader<K> sequentialReader(Descriptor descriptor, 
KeySupport<K> keySupport, int fsyncedLimit)
@@ -339,112 +153,25 @@ public final class StaticSegment<K, V> extends 
Segment<K, V>
         return new SequentialReader<>(descriptor, keySupport, fsyncedLimit);
     }
 
-    /**
-     * A sequential data segment reader to use for journal replay and 
rebuilding
-     * missing auxilirary components (index and metadata).
-     * </p>
-     * Unexpected EOF and CRC mismatches in synced portions of segments are 
treated
-     * strictly, throwing {@link JournalReadError}. Errors encountered in 
unsynced portions
-     * of segments are treated as segment EOF.
-     */
     static final class SequentialReader<K> extends Reader<K>
     {
-        private final int fsyncedLimit; // exclusive
-
-        SequentialReader(Descriptor descriptor, KeySupport<K> keySupport, int 
fsyncedLimit)
-        {
-            super(descriptor, keySupport);
-            this.fsyncedLimit = fsyncedLimit;
-        }
-
-        @Override
-        public boolean advance()
-        {
-            if (state == State.EOF)
-                return false;
-
-            reset();
-            return buffer.hasRemaining() ? doAdvance() : eof();
-        }
-
-        private boolean doAdvance()
-        {
-            offset = buffer.position();
-            try
-            {
-                int length = EntrySerializer.tryRead(holder, keySupport, 
buffer.duplicate(), fsyncedLimit, descriptor.userVersion);
-                if (length < 0)
-                    return eof();
-                buffer.position(offset + length);
-            }
-            catch (IOException e)
-            {
-                throw new JournalReadError(descriptor, file, e);
-            }
-
-            state = State.ADVANCED;
-            return true;
-        }
-
-        private void reset()
-        {
-            offset = -1;
-            holder.clear();
-            state = State.RESET;
-        }
+        private final int fsyncedLimit;
+        SequentialReader(Descriptor descriptor, KeySupport<K> keySupport, int 
fsyncedLimit) { super(descriptor, keySupport); this.fsyncedLimit = 
fsyncedLimit; }
+        @Override public boolean advance() { if (state == State.EOF) return 
false; reset(); return buffer.hasRemaining() ? doAdvance() : eof(); }
+        private boolean doAdvance() { offset = buffer.position(); try { int 
length = EntrySerializer.tryRead(holder, keySupport, buffer.duplicate(), 
fsyncedLimit, descriptor.userVersion); if (length < 0) return eof(); 
buffer.position(offset + length); state = State.ADVANCED; return true; } catch 
(IOException e) { throw new JournalReadError(descriptor, file, e); } }
+        private void reset() { offset = -1; holder.clear(); state = 
State.RESET; }
     }
 
-    public StaticSegment.KeyOrderReader<K> keyOrderReader()
+    public KeyOrderReader<K> keyOrderReader()
     {
-        return new StaticSegment.KeyOrderReader<>(descriptor, keySupport, 
index.reader());
+        return new KeyOrderReader<>(descriptor, keySupport, index.reader());
     }
 
     public static final class KeyOrderReader<K> extends Reader<K> implements 
Comparable<KeyOrderReader<K>>
     {
         private final OnDiskIndex<K>.IndexReader indexReader;
-
-        KeyOrderReader(Descriptor descriptor, KeySupport<K> keySupport, 
OnDiskIndex<K>.IndexReader indexReader)
-        {
-            super(descriptor, keySupport);
-            this.indexReader = indexReader;
-        }
-
-        @Override
-        public boolean advance()
-        {
-            if (!indexReader.advance())
-                return eof();
-
-            offset = indexReader.offset();
-
-            buffer.limit(offset + indexReader.recordSize())
-                  .position(offset);
-            try
-            {
-                EntrySerializer.read(holder, keySupport, buffer, 
descriptor.userVersion);
-            }
-            catch (IOException e)
-            {
-                throw new JournalReadError(descriptor, file, e);
-            }
-
-            state = State.ADVANCED;
-            return true;
-        }
-
-        @Override
-        public int compareTo(KeyOrderReader<K> that)
-        {
-            this.ensureHasAdvanced();
-            that.ensureHasAdvanced();
-
-            int cmp = keySupport.compare(this.key(), that.key());
-            if (cmp != 0)
-                return cmp;
-            cmp = Long.compare(that.descriptor.timestamp, 
this.descriptor.timestamp);
-            if (cmp != 0)
-                return cmp;
-            return Integer.compare(that.offset, this.offset);
-        }
+        KeyOrderReader(Descriptor descriptor, KeySupport<K> keySupport, 
OnDiskIndex<K>.IndexReader indexReader) { super(descriptor, keySupport); 
this.indexReader = indexReader; }
+        @Override public boolean advance() { if (!indexReader.advance()) 
return eof(); offset = indexReader.offset(); buffer.limit(offset + 
indexReader.recordSize()).position(offset); try { EntrySerializer.read(holder, 
keySupport, buffer, descriptor.userVersion); } catch (IOException e) { throw 
new JournalReadError(descriptor, file, e); } state = State.ADVANCED; return 
true; }
+        @Override public int compareTo(KeyOrderReader<K> that) { 
this.ensureHasAdvanced(); that.ensureHasAdvanced(); int cmp = 
keySupport.compare(this.key(), that.key()); if (cmp != 0) return cmp; cmp = 
Long.compare(that.descriptor.timestamp, this.descriptor.timestamp); return cmp 
!= 0 ? cmp : Integer.compare(that.offset, this.offset); }
     }
-}
\ No newline at end of file
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/service/StorageService.java 
b/journal/src/main/java/org/apache/cassandra/service/StorageService.java
new file mode 100644
index 0000000..4a1a9dd
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/service/StorageService.java
@@ -0,0 +1,10 @@
+package org.apache.cassandra.service;
+
+public final class StorageService
+{
+    public static final StorageService instance = new StorageService();
+
+    public void stopTransports() {}
+
+    private StorageService() {}
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/service/accord/serializers/Version.java
 
b/journal/src/main/java/org/apache/cassandra/service/accord/serializers/Version.java
new file mode 100644
index 0000000..ca6a50a
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/service/accord/serializers/Version.java
@@ -0,0 +1,14 @@
+package org.apache.cassandra.service.accord.serializers;
+
+public enum Version
+{
+    V1(1),
+    LATEST(1);
+
+    public final int version;
+
+    Version(int version)
+    {
+        this.version = version;
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/AbstractIterator.java 
b/journal/src/main/java/org/apache/cassandra/utils/AbstractIterator.java
new file mode 100644
index 0000000..a190fd8
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/AbstractIterator.java
@@ -0,0 +1,39 @@
+package org.apache.cassandra.utils;
+
+import java.util.Iterator;
+import java.util.NoSuchElementException;
+
+public abstract class AbstractIterator<T> implements Iterator<T>
+{
+    private enum State { READY, NOT_READY, DONE }
+    private State state = State.NOT_READY;
+    private T next;
+
+    protected abstract T computeNext();
+
+    protected final T endOfData()
+    {
+        state = State.DONE;
+        return null;
+    }
+
+    @Override
+    public boolean hasNext()
+    {
+        if (state == State.DONE)
+            return false;
+        if (state == State.READY)
+            return true;
+        next = computeNext();
+        return state != State.DONE && (state = State.READY) == State.READY;
+    }
+
+    @Override
+    public T next()
+    {
+        if (!hasNext())
+            throw new NoSuchElementException();
+        state = State.NOT_READY;
+        return next;
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/ByteBufferUtil.java 
b/journal/src/main/java/org/apache/cassandra/utils/ByteBufferUtil.java
new file mode 100644
index 0000000..f485ce6
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/ByteBufferUtil.java
@@ -0,0 +1,23 @@
+package org.apache.cassandra.utils;
+
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+
+public final class ByteBufferUtil
+{
+    private ByteBufferUtil() {}
+
+    public static ByteBuffer bytes(String value)
+    {
+        return ByteBuffer.wrap(value.getBytes(StandardCharsets.UTF_8));
+    }
+
+    public static void copyBytes(ByteBuffer src, int srcPos, ByteBuffer dst, 
int dstPos, int length)
+    {
+        ByteBuffer s = src.duplicate();
+        ByteBuffer d = dst.duplicate();
+        s.position(srcPos).limit(srcPos + length);
+        d.position(dstPos).limit(dstPos + length);
+        d.put(s);
+    }
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/Clock.java 
b/journal/src/main/java/org/apache/cassandra/utils/Clock.java
new file mode 100644
index 0000000..f3f4d24
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/Clock.java
@@ -0,0 +1,12 @@
+package org.apache.cassandra.utils;
+
+public final class Clock
+{
+    public static final class Global
+    {
+        public static long currentTimeMillis() { return 
System.currentTimeMillis(); }
+        public static long nanoTime() { return System.nanoTime(); }
+    }
+
+    private Clock() {}
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/Closeable.java 
b/journal/src/main/java/org/apache/cassandra/utils/Closeable.java
new file mode 100644
index 0000000..5fa62f9
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/Closeable.java
@@ -0,0 +1,7 @@
+package org.apache.cassandra.utils;
+
+public interface Closeable extends AutoCloseable
+{
+    @Override
+    void close();
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/CloseableIterator.java 
b/journal/src/main/java/org/apache/cassandra/utils/CloseableIterator.java
new file mode 100644
index 0000000..a345b4f
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/CloseableIterator.java
@@ -0,0 +1,9 @@
+package org.apache.cassandra.utils;
+
+import java.util.Iterator;
+
+public interface CloseableIterator<T> extends Iterator<T>, AutoCloseable
+{
+    @Override
+    void close();
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/Crc.java 
b/journal/src/main/java/org/apache/cassandra/utils/Crc.java
new file mode 100644
index 0000000..197c553
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/Crc.java
@@ -0,0 +1,29 @@
+package org.apache.cassandra.utils;
+
+import java.nio.ByteBuffer;
+import java.util.zip.CRC32;
+
+public final class Crc
+{
+    public static class InvalidCrc extends java.io.IOException
+    {
+        public InvalidCrc(int read, int expected)
+        {
+            super("Invalid CRC read=" + read + " expected=" + expected);
+        }
+    }
+
+    private Crc() {}
+
+    public static CRC32 crc32()
+    {
+        return new CRC32();
+    }
+
+    public static void updateCrc32(CRC32 crc, ByteBuffer buffer, int start, 
int end)
+    {
+        ByteBuffer dup = buffer.duplicate();
+        dup.position(start).limit(end);
+        crc.update(dup);
+    }
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/FBUtilities.java 
b/journal/src/main/java/org/apache/cassandra/utils/FBUtilities.java
new file mode 100644
index 0000000..e00b233
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/FBUtilities.java
@@ -0,0 +1,22 @@
+package org.apache.cassandra.utils;
+
+import java.util.zip.Checksum;
+
+public final class FBUtilities
+{
+    private FBUtilities() {}
+
+    public static void updateChecksumInt(Checksum crc, int value)
+    {
+        crc.update((value >>> 24) & 0xff);
+        crc.update((value >>> 16) & 0xff);
+        crc.update((value >>> 8) & 0xff);
+        crc.update(value & 0xff);
+    }
+
+    public static void updateChecksumLong(Checksum crc, long value)
+    {
+        for (int shift = 56; shift >= 0; shift -= 8)
+            crc.update((int) ((value >>> shift) & 0xff));
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/JVMStabilityInspector.java 
b/journal/src/main/java/org/apache/cassandra/utils/JVMStabilityInspector.java
new file mode 100644
index 0000000..dd1b425
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/utils/JVMStabilityInspector.java
@@ -0,0 +1,10 @@
+package org.apache.cassandra.utils;
+
+import org.apache.cassandra.journal.Params;
+
+public final class JVMStabilityInspector
+{
+    private JVMStabilityInspector() {}
+
+    public static void inspectJournalThrowable(Throwable t, String name, 
Params.FailurePolicy policy) {}
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/LazyToString.java 
b/journal/src/main/java/org/apache/cassandra/utils/LazyToString.java
new file mode 100644
index 0000000..ec46c77
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/LazyToString.java
@@ -0,0 +1,21 @@
+package org.apache.cassandra.utils;
+
+import java.util.function.Supplier;
+
+public final class LazyToString
+{
+    private LazyToString() {}
+
+    public static Object lazy(Supplier<?> supplier)
+    {
+        return new Object()
+        {
+            @Override
+            public String toString()
+            {
+                Object value = supplier.get();
+                return value == null ? "null" : value.toString();
+            }
+        };
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/MergeIterator.java 
b/journal/src/main/java/org/apache/cassandra/utils/MergeIterator.java
new file mode 100644
index 0000000..5986cc2
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/MergeIterator.java
@@ -0,0 +1,118 @@
+package org.apache.cassandra.utils;
+
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.Iterator;
+import java.util.List;
+import java.util.NoSuchElementException;
+
+public final class MergeIterator<In, Out> implements Iterator<Out>
+{
+    public static abstract class Reducer<In, Out>
+    {
+        public abstract void reduce(int idx, In current);
+        protected abstract Out getReduced();
+        protected void onKeyChange() {}
+    }
+
+    private static final class Source<In>
+    {
+        final Iterator<In> iterator;
+        In current;
+        boolean hasCurrent;
+
+        Source(Iterator<In> iterator)
+        {
+            this.iterator = iterator;
+            advance();
+        }
+
+        void advance()
+        {
+            hasCurrent = iterator.hasNext();
+            current = hasCurrent ? iterator.next() : null;
+        }
+    }
+
+    private final List<Source<In>> sources;
+    private final Comparator<? super In> comparator;
+    private final Reducer<In, Out> reducer;
+    private Out next;
+    private boolean prepared;
+
+    private MergeIterator(List<Iterator<In>> iterators, Comparator<? super In> 
comparator, Reducer<In, Out> reducer)
+    {
+        this.sources = new ArrayList<>(iterators.size());
+        for (Iterator<In> iterator : iterators)
+            this.sources.add(new Source<>(iterator));
+        this.comparator = comparator;
+        this.reducer = reducer;
+    }
+
+    public static <In, Out> MergeIterator<In, Out> get(List<Iterator<In>> 
iterators, Comparator<? super In> comparator, Reducer<In, Out> reducer)
+    {
+        return new MergeIterator<>(iterators, comparator, reducer);
+    }
+
+    public Out peek()
+    {
+        prepare();
+        return next;
+    }
+
+    @Override
+    public boolean hasNext()
+    {
+        prepare();
+        return next != null;
+    }
+
+    @Override
+    public Out next()
+    {
+        prepare();
+        if (next == null)
+            throw new NoSuchElementException();
+        Out result = next;
+        prepared = false;
+        next = null;
+        return result;
+    }
+
+    private void prepare()
+    {
+        if (prepared)
+            return;
+        prepared = true;
+        next = computeNext();
+    }
+
+    private Out computeNext()
+    {
+        int first = -1;
+        for (int i = 0; i < sources.size(); i++)
+        {
+            if (!sources.get(i).hasCurrent)
+                continue;
+            if (first < 0 || comparator.compare(sources.get(i).current, 
sources.get(first).current) < 0)
+                first = i;
+        }
+        if (first < 0)
+            return null;
+
+        reducer.onKeyChange();
+        In key = sources.get(first).current;
+        for (int i = 0; i < sources.size(); i++)
+        {
+            Source<In> source = sources.get(i);
+            if (!source.hasCurrent)
+                continue;
+            if (comparator.compare(source.current, key) == 0)
+            {
+                reducer.reduce(i, source.current);
+                source.advance();
+            }
+        }
+        return reducer.getReduced();
+    }
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/Pair.java 
b/journal/src/main/java/org/apache/cassandra/utils/Pair.java
new file mode 100644
index 0000000..f3891bf
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/Pair.java
@@ -0,0 +1,18 @@
+package org.apache.cassandra.utils;
+
+public final class Pair<L, R>
+{
+    public final L left;
+    public final R right;
+
+    private Pair(L left, R right)
+    {
+        this.left = left;
+        this.right = right;
+    }
+
+    public static <L, R> Pair<L, R> create(L left, R right)
+    {
+        return new Pair<>(left, right);
+    }
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/Simulate.java 
b/journal/src/main/java/org/apache/cassandra/utils/Simulate.java
new file mode 100644
index 0000000..b73d52a
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/Simulate.java
@@ -0,0 +1,18 @@
+package org.apache.cassandra.utils;
+
+import java.lang.annotation.ElementType;
+import java.lang.annotation.Retention;
+import java.lang.annotation.RetentionPolicy;
+import java.lang.annotation.Target;
+
+@Retention(RetentionPolicy.CLASS)
+@Target({ ElementType.TYPE, ElementType.METHOD, ElementType.CONSTRUCTOR, 
ElementType.FIELD })
+public @interface Simulate
+{
+    With with();
+
+    enum With
+    {
+        MONITORS
+    }
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/SyncUtil.java 
b/journal/src/main/java/org/apache/cassandra/utils/SyncUtil.java
new file mode 100644
index 0000000..974d005
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/SyncUtil.java
@@ -0,0 +1,21 @@
+package org.apache.cassandra.utils;
+
+import java.io.IOException;
+import java.nio.MappedByteBuffer;
+import java.nio.channels.FileChannel;
+import java.nio.file.StandardOpenOption;
+
+import org.apache.cassandra.io.util.File;
+
+public final class SyncUtil
+{
+    private SyncUtil() {}
+
+    public static void force(MappedByteBuffer buffer)
+    {
+    }
+
+    public static void trySyncDir(File dir)
+    {
+    }
+}
diff --git a/journal/src/main/java/org/apache/cassandra/utils/Throwables.java 
b/journal/src/main/java/org/apache/cassandra/utils/Throwables.java
new file mode 100644
index 0000000..3b768aa
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/Throwables.java
@@ -0,0 +1,23 @@
+package org.apache.cassandra.utils;
+
+public final class Throwables
+{
+    private Throwables() {}
+
+    public static Throwable merge(Throwable existing, Throwable next)
+    {
+        if (existing == null)
+            return next;
+        existing.addSuppressed(next);
+        return existing;
+    }
+
+    public static void maybeFail(Throwable t)
+    {
+        if (t == null)
+            return;
+        if (t instanceof RuntimeException)
+            throw (RuntimeException) t;
+        throw new RuntimeException(t);
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/concurrent/AsyncPromise.java 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/AsyncPromise.java
new file mode 100644
index 0000000..d479814
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/AsyncPromise.java
@@ -0,0 +1,37 @@
+package org.apache.cassandra.utils.concurrent;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutionException;
+
+public class AsyncPromise<V> implements Promise<V>
+{
+    private final CompletableFuture<V> future = new CompletableFuture<>();
+
+    @Override public boolean cancel(boolean mayInterruptIfRunning) { return 
future.cancel(mayInterruptIfRunning); }
+    @Override public boolean isDone() { return future.isDone(); }
+    @Override public boolean isSuccess() { return future.isDone() && 
!future.isCompletedExceptionally() && !future.isCancelled(); }
+    @Override public V getNow() { return future.getNow(null); }
+    @Override public void setSuccess(V value) { future.complete(value); }
+
+    @Override
+    public Promise<V> awaitThrowUncheckedOnInterrupt()
+    {
+        try
+        {
+            future.get();
+            return this;
+        }
+        catch (InterruptedException e)
+        {
+            Thread.currentThread().interrupt();
+            throw new UncheckedInterruptedException();
+        }
+        catch (ExecutionException e)
+        {
+            Throwable cause = e.getCause();
+            if (cause instanceof RuntimeException)
+                throw (RuntimeException) cause;
+            throw new RuntimeException(cause);
+        }
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/concurrent/CountDownLatch.java
 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/CountDownLatch.java
new file mode 100644
index 0000000..d60db07
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/CountDownLatch.java
@@ -0,0 +1,33 @@
+package org.apache.cassandra.utils.concurrent;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+public final class CountDownLatch
+{
+    private final AtomicInteger count;
+
+    private CountDownLatch(int count)
+    {
+        this.count = new AtomicInteger(count);
+    }
+
+    public static CountDownLatch newCountDownLatch(int count)
+    {
+        return new CountDownLatch(count);
+    }
+
+    public int count()
+    {
+        return count.get();
+    }
+
+    public void decrement()
+    {
+        while (true)
+        {
+            int cur = count.get();
+            if (cur == 0 || count.compareAndSet(cur, cur - 1))
+                return;
+        }
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/concurrent/OpOrder.java 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/OpOrder.java
new file mode 100644
index 0000000..1c825d7
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/concurrent/OpOrder.java
@@ -0,0 +1,52 @@
+package org.apache.cassandra.utils.concurrent;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+public class OpOrder
+{
+    private final AtomicInteger inFlight = new AtomicInteger();
+
+    public Group start()
+    {
+        inFlight.incrementAndGet();
+        return new Group();
+    }
+
+    public Barrier newBarrier()
+    {
+        return new Barrier();
+    }
+
+    public final class Group implements AutoCloseable
+    {
+        private boolean closed;
+
+        @Override
+        public void close()
+        {
+            if (!closed)
+            {
+                closed = true;
+                inFlight.decrementAndGet();
+            }
+        }
+    }
+
+    public final class Barrier
+    {
+        private volatile boolean issued;
+
+        public void issue()
+        {
+            issued = true;
+        }
+
+        public void await()
+        {
+            if (!issued)
+                throw new IllegalStateException("Barrier not issued");
+            while (inFlight.get() > 0)
+                Thread.yield();
+        }
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/concurrent/Promise.java 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/Promise.java
new file mode 100644
index 0000000..7706fb1
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/concurrent/Promise.java
@@ -0,0 +1,11 @@
+package org.apache.cassandra.utils.concurrent;
+
+public interface Promise<V>
+{
+    boolean cancel(boolean mayInterruptIfRunning);
+    boolean isDone();
+    boolean isSuccess();
+    V getNow();
+    Promise<V> awaitThrowUncheckedOnInterrupt();
+    void setSuccess(V value);
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/concurrent/Ref.java 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/Ref.java
new file mode 100644
index 0000000..254ef3e
--- /dev/null
+++ b/journal/src/main/java/org/apache/cassandra/utils/concurrent/Ref.java
@@ -0,0 +1,44 @@
+package org.apache.cassandra.utils.concurrent;
+
+import java.util.concurrent.Executor;
+import java.util.concurrent.atomic.AtomicInteger;
+
+public final class Ref<T>
+{
+    public interface Tidier
+    {
+        void tidy() throws Exception;
+    }
+
+    private final T referent;
+    private final Object tidier;
+    private final AtomicInteger globalCount = new AtomicInteger(1);
+
+    public Ref(T referent, Object tidier)
+    {
+        this.referent = referent;
+        this.tidier = tidier;
+    }
+
+    public int globalCount()
+    {
+        return globalCount.get();
+    }
+
+    public Object tidier()
+    {
+        return tidier;
+    }
+
+    public void release()
+    {
+        globalCount.set(0);
+        if (tidier instanceof Runnable)
+            ((Runnable) tidier).run();
+    }
+
+    public T get()
+    {
+        return referent;
+    }
+}
diff --git 
a/journal/src/main/java/org/apache/cassandra/utils/concurrent/UncheckedInterruptedException.java
 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/UncheckedInterruptedException.java
new file mode 100644
index 0000000..2ce063b
--- /dev/null
+++ 
b/journal/src/main/java/org/apache/cassandra/utils/concurrent/UncheckedInterruptedException.java
@@ -0,0 +1,6 @@
+package org.apache.cassandra.utils.concurrent;
+
+public class UncheckedInterruptedException extends RuntimeException
+{
+    public UncheckedInterruptedException() {}
+}
diff --git 
a/journal/src/test/java/org/apache/cassandra/harry/checker/TestHelper.java 
b/journal/src/test/java/org/apache/cassandra/harry/checker/TestHelper.java
new file mode 100644
index 0000000..d53f364
--- /dev/null
+++ b/journal/src/test/java/org/apache/cassandra/harry/checker/TestHelper.java
@@ -0,0 +1,19 @@
+package org.apache.cassandra.harry.checker;
+
+import org.apache.cassandra.harry.gen.EntropySource;
+
+public final class TestHelper
+{
+    @FunctionalInterface
+    public interface ThrowingConsumer<T>
+    {
+        void accept(T value) throws Throwable;
+    }
+
+    private TestHelper() {}
+
+    public static void withRandom(ThrowingConsumer<EntropySource> consumer) 
throws Throwable
+    {
+        consumer.accept(new EntropySource(1L));
+    }
+}
diff --git 
a/journal/src/test/java/org/apache/cassandra/harry/gen/EntropySource.java 
b/journal/src/test/java/org/apache/cassandra/harry/gen/EntropySource.java
new file mode 100644
index 0000000..d5a3b03
--- /dev/null
+++ b/journal/src/test/java/org/apache/cassandra/harry/gen/EntropySource.java
@@ -0,0 +1,28 @@
+package org.apache.cassandra.harry.gen;
+
+import java.util.Random;
+
+public class EntropySource
+{
+    private final Random random;
+
+    public EntropySource(long seed)
+    {
+        this.random = new Random(seed);
+    }
+
+    private EntropySource(Random random)
+    {
+        this.random = random;
+    }
+
+    public long next()
+    {
+        return random.nextLong();
+    }
+
+    public EntropySource derive()
+    {
+        return new EntropySource(new Random(random.nextLong()));
+    }
+}
diff --git a/journal/src/test/java/org/apache/cassandra/utils/TimeUUID.java 
b/journal/src/test/java/org/apache/cassandra/utils/TimeUUID.java
index 2b2b0c8..f72bd71 100644
--- a/journal/src/test/java/org/apache/cassandra/utils/TimeUUID.java
+++ b/journal/src/test/java/org/apache/cassandra/utils/TimeUUID.java
@@ -36,7 +36,8 @@ public final class TimeUUID implements Comparable<TimeUUID>
     public boolean equals(Object o)
     {
         if (this == o) return true;
-        if (!(o instanceof TimeUUID that)) return false;
+        if (!(o instanceof TimeUUID)) return false;
+        TimeUUID that = (TimeUUID) o;
         return uuidTimestamp == that.uuidTimestamp && lsb == that.lsb;
     }
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to