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]
