From f7743fe3ab2f802977af8720ebd2a8fc1a32dd97 Mon Sep 17 00:00:00 2001 From: Benedict Elliott Smith Date: Mon, 8 Sep 2025 14:42:32 +0100 Subject: [PATCH] Lazy Virtual Tables Also Improve: - Searchable system_accord_debug.txn - Integrate txn_blocked_by deadline and depth filter with execution logic, to ensure we terminate promptly and get a best effort reply patch by Benedict; reviewed by Aleksey Yeschenko for CASSANDRA-20899 --- .../restrictions/StatementRestrictions.java | 12 +- .../cql3/statements/SelectStatement.java | 4 +- .../db/PartitionRangeReadCommand.java | 2 +- .../db/SinglePartitionReadCommand.java | 2 +- .../cassandra/db/marshal/TxnIdUtf8Type.java | 8 +- .../db/virtual/AbstractLazyVirtualTable.java | 767 +++++++++++ .../AbstractMutableLazyVirtualTable.java | 137 ++ .../db/virtual/AbstractVirtualTable.java | 5 +- .../db/virtual/AccordDebugKeyspace.java | 1167 +++++++---------- .../CollectionVirtualTableAdapter.java | 5 +- .../db/virtual/PartitionKeyStatsTable.java | 5 +- .../cassandra/db/virtual/VirtualTable.java | 24 +- .../cassandra/journal/InMemoryIndex.java | 11 + .../org/apache/cassandra/journal/Journal.java | 87 +- .../org/apache/cassandra/journal/Segment.java | 3 +- .../cassandra/journal/StaticSegment.java | 6 + .../service/accord/AccordExecutor.java | 5 +- .../service/accord/AccordJournal.java | 14 +- .../service/accord/AccordJournalTable.java | 11 +- .../service/accord/AccordService.java | 149 +-- .../accord/CommandStoreTxnBlockedGraph.java | 135 -- .../service/accord/DebugBlockedTxns.java | 249 ++++ .../service/accord/IAccordService.java | 16 - .../tools/StandaloneJournalUtil.java | 2 +- .../fuzz/topology/JournalGCTest.java | 4 +- .../service/accord/AccordJournalBurnTest.java | 2 +- .../org/apache/cassandra/cql3/CQLTester.java | 2 +- .../db/virtual/AccordDebugKeyspaceTest.java | 75 +- 28 files changed, 1855 insertions(+), 1054 deletions(-) create mode 100644 src/java/org/apache/cassandra/db/virtual/AbstractLazyVirtualTable.java create mode 100644 src/java/org/apache/cassandra/db/virtual/AbstractMutableLazyVirtualTable.java delete mode 100644 src/java/org/apache/cassandra/service/accord/CommandStoreTxnBlockedGraph.java create mode 100644 src/java/org/apache/cassandra/service/accord/DebugBlockedTxns.java diff --git a/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java b/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java index 3b40d7e412..8373f3eb88 100644 --- a/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java +++ b/src/java/org/apache/cassandra/cql3/restrictions/StatementRestrictions.java @@ -156,13 +156,13 @@ public final class StatementRestrictions return new StatementRestrictions(type, table, IndexHints.NONE, false); } - private StatementRestrictions(StatementType type, TableMetadata table, IndexHints indexHints, boolean allowFiltering) + private StatementRestrictions(StatementType type, TableMetadata table, IndexHints indexHints, boolean allowFilteringOfPrimaryKeys) { this.type = type; this.table = table; this.indexHints = indexHints; this.partitionKeyRestrictions = new PartitionKeyRestrictions(table.partitionKeyAsClusteringComparator()); - this.clusteringColumnsRestrictions = new ClusteringColumnRestrictions(table, allowFiltering); + this.clusteringColumnsRestrictions = new ClusteringColumnRestrictions(table, allowFilteringOfPrimaryKeys); this.nonPrimaryKeyRestrictions = RestrictionSet.empty(); this.notNullColumns = new HashSet<>(); } @@ -370,7 +370,7 @@ public final class StatementRestrictions } else { - if (!allowFiltering && requiresAllowFilteringIfNotSpecified(table)) + if (!allowFiltering && requiresAllowFilteringIfNotSpecified(table, false)) throw invalidRequest(allowFilteringMessage(state)); } @@ -381,14 +381,14 @@ public final class StatementRestrictions validateSecondaryIndexSelections(); } - public static boolean requiresAllowFilteringIfNotSpecified(TableMetadata metadata) + public static boolean requiresAllowFilteringIfNotSpecified(TableMetadata metadata, boolean isPrimaryKey) { if (!metadata.isVirtual()) return true; VirtualTable tableNullable = VirtualKeyspaceRegistry.instance.getTableNullable(metadata.id); assert tableNullable != null; - return !tableNullable.allowFilteringImplicitly(); + return isPrimaryKey ? !tableNullable.allowFilteringPrimaryKeysImplicitly() : !tableNullable.allowFilteringImplicitly(); } private void addRestriction(Restriction restriction, IndexRegistry indexRegistry, IndexHints indexHints) @@ -593,7 +593,7 @@ public final class StatementRestrictions // components must have a EQ. Only the last partition key component can be in IN relation. if (partitionKeyRestrictions.needFiltering()) { - if (!allowFiltering && !forView && !hasQueriableIndex && requiresAllowFilteringIfNotSpecified(table)) + if (!allowFiltering && !forView && !hasQueriableIndex && requiresAllowFilteringIfNotSpecified(table, true)) throw new InvalidRequestException(allowFilteringMessage(state)); isKeyRange = true; diff --git a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java index ecdce1437b..aa5500356d 100644 --- a/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/SelectStatement.java @@ -1474,7 +1474,7 @@ public class SelectStatement implements CQLStatement.SingleKeyspaceCqlStatement, boundNames, orderings, selectsOnlyStaticColumns, - parameters.allowFiltering || !requiresAllowFilteringIfNotSpecified(metadata), + parameters.allowFiltering || !requiresAllowFilteringIfNotSpecified(metadata, true), forView); } @@ -1700,7 +1700,7 @@ public class SelectStatement implements CQLStatement.SingleKeyspaceCqlStatement, { // We will potentially filter data if the row filter is not the identity and there isn't any index group // supporting all the expressions in the filter. - if (requiresAllowFilteringIfNotSpecified(table)) + if (requiresAllowFilteringIfNotSpecified(table, true)) checkFalse(restrictions.needFiltering(table), StatementRestrictions.REQUIRES_ALLOW_FILTERING_MESSAGE); } } diff --git a/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java b/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java index 25986314a2..f448bcf872 100644 --- a/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java +++ b/src/java/org/apache/cassandra/db/PartitionRangeReadCommand.java @@ -647,7 +647,7 @@ public class PartitionRangeReadCommand extends ReadCommand implements PartitionR public UnfilteredPartitionIterator executeLocally(ReadExecutionController executionController) { VirtualTable view = VirtualKeyspaceRegistry.instance.getTableNullable(metadata().id); - UnfilteredPartitionIterator resultIterator = view.select(dataRange, columnFilter(), rowFilter()); + UnfilteredPartitionIterator resultIterator = view.select(dataRange, columnFilter(), rowFilter(), limits()); return limits().filter(rowFilter().filter(resultIterator, nowInSec()), nowInSec(), selectsFullPartition()); } diff --git a/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java b/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java index 7c4864e6db..780a8c083f 100644 --- a/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java +++ b/src/java/org/apache/cassandra/db/SinglePartitionReadCommand.java @@ -1516,7 +1516,7 @@ public class SinglePartitionReadCommand extends ReadCommand implements SinglePar public UnfilteredPartitionIterator executeLocally(ReadExecutionController executionController) { VirtualTable view = VirtualKeyspaceRegistry.instance.getTableNullable(metadata().id); - UnfilteredPartitionIterator resultIterator = view.select(partitionKey, clusteringIndexFilter, columnFilter(), rowFilter()); + UnfilteredPartitionIterator resultIterator = view.select(partitionKey, clusteringIndexFilter, columnFilter(), rowFilter(), limits()); return limits().filter(rowFilter().filter(resultIterator, nowInSec()), nowInSec(), selectsFullPartition()); } diff --git a/src/java/org/apache/cassandra/db/marshal/TxnIdUtf8Type.java b/src/java/org/apache/cassandra/db/marshal/TxnIdUtf8Type.java index 784969ead7..1f333a0cd1 100644 --- a/src/java/org/apache/cassandra/db/marshal/TxnIdUtf8Type.java +++ b/src/java/org/apache/cassandra/db/marshal/TxnIdUtf8Type.java @@ -36,7 +36,7 @@ public class TxnIdUtf8Type extends PseudoUtf8Type { super.validate(value, accessor); String str = deserialize(value, accessor); - if (null == TxnId.tryParse(str)) + if (!str.isEmpty() && null == TxnId.tryParse(str)) throw new MarshalException("Invalid TxnId: " + str); } }; @@ -59,6 +59,12 @@ public class TxnIdUtf8Type extends PseudoUtf8Type { String leftStr = UTF8Serializer.instance.deserialize(left, accessorL); String rightStr = UTF8Serializer.instance.deserialize(right, accessorR); + if (leftStr.isEmpty() || rightStr.isEmpty()) + { + if (leftStr.isEmpty() && rightStr.isEmpty()) + return 0; + return leftStr.isEmpty() ? -1 : 1; + } TxnId leftId = TxnId.parse(leftStr); TxnId rightId = TxnId.parse(rightStr); return leftId.compareTo(rightId); diff --git a/src/java/org/apache/cassandra/db/virtual/AbstractLazyVirtualTable.java b/src/java/org/apache/cassandra/db/virtual/AbstractLazyVirtualTable.java new file mode 100644 index 0000000000..6377824417 --- /dev/null +++ b/src/java/org/apache/cassandra/db/virtual/AbstractLazyVirtualTable.java @@ -0,0 +1,767 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.cassandra.db.virtual; + +import java.nio.ByteBuffer; +import java.util.Arrays; +import java.util.Comparator; +import java.util.HashMap; +import java.util.Iterator; +import java.util.Map; +import java.util.NavigableMap; +import java.util.TreeMap; +import java.util.concurrent.TimeUnit; +import java.util.function.Consumer; +import java.util.function.Function; +import java.util.function.UnaryOperator; + +import javax.annotation.Nullable; + +import accord.utils.Invariants; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.db.Clusterable; +import org.apache.cassandra.db.Clustering; +import org.apache.cassandra.db.ClusteringPrefix; +import org.apache.cassandra.db.DataRange; +import org.apache.cassandra.db.DecoratedKey; +import org.apache.cassandra.db.DeletionTime; +import org.apache.cassandra.db.LivenessInfo; +import org.apache.cassandra.db.RegularAndStaticColumns; +import org.apache.cassandra.db.filter.ClusteringIndexFilter; +import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.DataLimits; +import org.apache.cassandra.db.filter.RowFilter; +import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.db.marshal.CompositeType; +import org.apache.cassandra.db.partitions.PartitionUpdate; +import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; +import org.apache.cassandra.db.rows.BTreeRow; +import org.apache.cassandra.db.rows.BufferCell; +import org.apache.cassandra.db.rows.ColumnData; +import org.apache.cassandra.db.rows.EncodingStats; +import org.apache.cassandra.db.rows.Row; +import org.apache.cassandra.db.rows.Unfiltered; +import org.apache.cassandra.db.rows.UnfilteredRowIterator; +import org.apache.cassandra.dht.AbstractBounds; +import org.apache.cassandra.dht.Bounds; +import org.apache.cassandra.exceptions.InvalidRequestException; +import org.apache.cassandra.exceptions.ReadTimeoutException; +import org.apache.cassandra.schema.ColumnMetadata; +import org.apache.cassandra.schema.SchemaConstants; +import org.apache.cassandra.schema.TableMetadata; +import org.apache.cassandra.service.ClientWarn; +import org.apache.cassandra.utils.BulkIterator; +import org.apache.cassandra.utils.Clock; +import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.btree.BTree; +import org.apache.cassandra.utils.btree.UpdateFunction; + +import static org.apache.cassandra.db.ClusteringPrefix.Kind.STATIC_CLUSTERING; +import static org.apache.cassandra.db.ConsistencyLevel.ONE; +import static org.apache.cassandra.utils.Clock.Global.nanoTime; + +/** + * An abstract virtual table implementation that builds the resultset on demand. + */ +public abstract class AbstractLazyVirtualTable implements VirtualTable +{ + public enum OnTimeout { BEST_EFFORT, FAIL } + + // in the special case where we know we have enough rows in the collector, throw this exception to terminate early + public static class InternalDoneException extends RuntimeException {} + // in the special case where we have timed out, throw this exception to terminate early + public static class InternalTimeoutException extends RuntimeException {} + + public static class FilterRange + { + final V min, max; + public FilterRange(V min, V max) + { + this.min = min; + this.max = max; + } + } + + public interface PartitionsCollector + { + DataRange dataRange(); + RowFilter rowFilter(); + ColumnFilter columnFilter(); + DataLimits limits(); + long nowInSeconds(); + long timestampMicros(); + long deadlineNanos(); + boolean isEmpty(); + + RowCollector row(Object... primaryKeys); + PartitionCollector partition(Object... partitionKeys); + UnfilteredPartitionIterator finish(); + + @Nullable Object[] singlePartitionKey(); + FilterRange filters(String column, Function translate, UnaryOperator exclusiveStart, UnaryOperator exclusiveEnd); + } + + public interface PartitionCollector + { + void collect(Consumer addTo); + } + + public interface RowsCollector + { + RowCollector add(Object... clusteringKeys); + } + + public interface RowCollector + { + default void lazyCollect(Consumer addToIfNeeded) { eagerCollect(addToIfNeeded); } + void eagerCollect(Consumer addToNow); + } + + public interface ColumnsCollector + { + /** + * equivalent to + * {@code + * if (value == null) add(columnName, null); + * else if (f1.apply(value) == null) add(columnName, f1.apply(value)); + * else add(columnName, f2.apply(f1.apply(value))); + * } + */ + ColumnsCollector add(String columnName, V1 value, Function f1, Function f2); + + default ColumnsCollector add(String columnName, V value, Function transform) + { + return add(columnName, value, Function.identity(), transform); + } + default ColumnsCollector add(String columnName, Object value) + { + return add(columnName, value, Function.identity()); + } + } + + public static class SimplePartitionsCollector implements PartitionsCollector + { + final TableMetadata metadata; + final boolean isSorted; + final boolean isSortedByPartitionKey; + + final Map columnLookup = new HashMap<>(); + final NavigableMap partitions; + + final DataRange dataRange; + final ColumnFilter columnFilter; + final RowFilter rowFilter; + final DataLimits limits; + + final long startedAtNanos = Clock.Global.nanoTime(); + final long deadlineNanos; + + final long nowInSeconds = Clock.Global.nowInSeconds(); + final long timestampMicros; + + int totalRowCount; + int lastFilteredTotalRowCount, lastFilteredPartitionCount; + + @Override public DataRange dataRange() { return dataRange; } + @Override public RowFilter rowFilter() { return rowFilter; } + @Override public ColumnFilter columnFilter() { return columnFilter; } + @Override public DataLimits limits() { return limits; } + @Override public long nowInSeconds() { return nowInSeconds; } + @Override public long timestampMicros() { return timestampMicros; } + @Override public long deadlineNanos() { return deadlineNanos; } + @Override public boolean isEmpty() { return totalRowCount == 0; } + + public SimplePartitionsCollector(TableMetadata metadata, boolean isSorted, boolean isSortedByPartitionKey, + DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) + { + this.metadata = metadata; + this.isSorted = isSorted; + this.isSortedByPartitionKey = isSortedByPartitionKey; + this.dataRange = dataRange; + this.columnFilter = columnFilter; + this.rowFilter = rowFilter; + this.limits = limits; + this.timestampMicros = FBUtilities.timestampMicros(); + this.deadlineNanos = startedAtNanos + DatabaseDescriptor.getReadRpcTimeout(TimeUnit.NANOSECONDS); + this.partitions = new TreeMap<>(dataRange.isReversed() ? DecoratedKey.comparator.reversed() : DecoratedKey.comparator); + for (ColumnMetadata cm : metadata.columns()) + columnLookup.put(cm.name.toString(), cm); + } + + public Object[] singlePartitionKey() + { + AbstractBounds bounds = dataRange().keyRange(); + if (!bounds.isStartInclusive() || !bounds.isEndInclusive() || !bounds.left.equals(bounds.right) || !(bounds.left instanceof DecoratedKey)) + return null; + + return composePartitionKeys((DecoratedKey) bounds.left, metadata); + } + + @Override + public PartitionCollector partition(Object ... partitionKeys) + { + int pkSize = metadata.partitionKeyColumns().size(); + if (pkSize != partitionKeys.length) + throw new IllegalArgumentException(); + + DecoratedKey partitionKey = decomposePartitionKeys(metadata, partitionKeys); + if (!dataRange.contains(partitionKey)) + return dropCks -> {}; + + return partitions.computeIfAbsent(partitionKey, SimplePartition::new); + } + + @Override + public UnfilteredPartitionIterator finish() + { + final Iterator partitions = this.partitions.values().iterator(); + return new UnfilteredPartitionIterator() + { + @Override public TableMetadata metadata() { return metadata; } + @Override public void close() {} + + @Override + public boolean hasNext() + { + return partitions.hasNext(); + } + + @Override + public UnfilteredRowIterator next() + { + SimplePartition partition = partitions.next(); + Iterator rows = partition.rows(); + + return new UnfilteredRowIterator() + { + @Override public TableMetadata metadata() { return metadata; } + @Override public boolean isReverseOrder() { return dataRange.isReversed(); } + @Override public RegularAndStaticColumns columns() { return columnFilter.fetchedColumns(); } + @Override public DecoratedKey partitionKey() { return partition.key; } + + @Override public Row staticRow() { return partition.staticRow(); } + @Override public boolean hasNext() { return rows.hasNext(); } + @Override public Unfiltered next() { return rows.next(); } + + @Override public void close() {} + @Override public DeletionTime partitionLevelDeletion() { return DeletionTime.LIVE; } + @Override public EncodingStats stats() { return EncodingStats.NO_STATS; } + }; + } + }; + } + + @Override + @Nullable + public FilterRange filters(String columnName, Function translate, UnaryOperator exclusiveStart, UnaryOperator exclusiveEnd) + { + ColumnMetadata column = columnLookup.get(columnName); + O min = null, max = null; + for (RowFilter.Expression expression : rowFilter().getExpressions()) + { + if (!expression.column().equals(column)) + continue; + + O bound = translate.apply((I)column.type.compose(expression.getIndexValue())); + switch (expression.operator()) + { + default: throw new InvalidRequestException("Operator " + expression.operator() + " not supported for txn_id"); + case EQ: min = max = bound; break; + case LTE: max = bound; break; + case LT: max = exclusiveEnd.apply(bound); break; + case GTE: min = bound; break; + case GT: min = exclusiveStart.apply(bound); break; + } + } + + return new FilterRange<>(min, max); + } + + @Override + public RowCollector row(Object... primaryKeys) + { + int pkSize = metadata.partitionKeyColumns().size(); + int ckSize = metadata.clusteringColumns().size(); + if (pkSize + ckSize != primaryKeys.length) + throw new IllegalArgumentException(); + + Object[] partitionKeyValues = new Object[pkSize]; + Object[] clusteringValues = new Object[ckSize]; + + System.arraycopy(primaryKeys, 0, partitionKeyValues, 0, pkSize); + System.arraycopy(primaryKeys, pkSize, clusteringValues, 0, ckSize); + + DecoratedKey partitionKey = decomposePartitionKeys(metadata, partitionKeyValues); + Clustering clustering = decomposeClusterings(metadata, clusteringValues); + + if (!dataRange.contains(partitionKey) || !dataRange.clusteringIndexFilter(partitionKey).selects(clustering)) + return drop -> {}; + + return partitions.computeIfAbsent(partitionKey, SimplePartition::new).row(clustering); + } + + private final class SimplePartition implements PartitionCollector, RowsCollector + { + private final DecoratedKey key; + // we assume no duplicate rows, and impose the condition lazily + private SimpleRow[] rows; + private int rowCount; + private SimpleRow staticRow; + private boolean dropRows; + + private SimplePartition(DecoratedKey key) + { + this.key = key; + this.rows = new SimpleRow[1]; + } + + @Override + public void collect(Consumer addTo) + { + addTo.accept(this); + } + + @Override + public RowCollector add(Object... clusteringKeys) + { + int ckSize = metadata.clusteringColumns().size(); + if (ckSize != clusteringKeys.length) + throw new IllegalArgumentException(); + + return row(decomposeClusterings(metadata, clusteringKeys)); + } + + RowCollector row(Clustering clustering) + { + if (nanoTime() > deadlineNanos) + throw new InternalTimeoutException(); + + if (dropRows || !dataRange.clusteringIndexFilter(key).selects(clustering)) + return drop -> {}; + + if (totalRowCount >= limits.count()) + { + boolean filter; + if (!isSortedByPartitionKey || lastFilteredPartitionCount == partitions.size()) + { + filter = totalRowCount / 2 >= Math.max(1024, limits.count()); + } + else + { + int rowsAddedSinceLastFiltered = totalRowCount - lastFilteredTotalRowCount; + int threshold = Math.max(32, Math.min(1024, lastFilteredTotalRowCount / 2)); + filter = lastFilteredTotalRowCount == 0 || rowsAddedSinceLastFiltered >= threshold; + } + + if (filter) + { + // first filter within each partition + for (SimplePartition partition : partitions.values()) + partition.truncate(limits.perPartitionCount()); + + // then drop any partitions that completely fall outside our limit + Iterator iter = partitions.descendingMap().values().iterator(); + SimplePartition last; + while (true) + { + SimplePartition next = last = iter.next(); + if (totalRowCount - next.rowCount < limits.count()) + break; + + iter.remove(); + totalRowCount -= next.rowCount; + if (next == this) + dropRows = true; + } + + // possibly truncate the last partition if it partially falls outside the limit + int overflow = Math.max(0, totalRowCount - limits.count()); + int newCount = last.truncate(last.rowCount - overflow); + lastFilteredTotalRowCount = totalRowCount; + lastFilteredPartitionCount = partitions.size(); + + if (isSortedByPartitionKey && totalRowCount - newCount >= limits.count()) + throw new InternalDoneException(); + + if (isSorted && totalRowCount >= limits.count()) + throw new InternalDoneException(); + + if (dropRows) + return drop -> {}; + } + } + + SimpleRow result = new SimpleRow(clustering); + if (clustering.kind() == STATIC_CLUSTERING) + { + Invariants.require(staticRow == null); + staticRow = result; + } + else + { + totalRowCount++; + if (rowCount == rows.length) + rows = Arrays.copyOf(rows, Math.max(8, rowCount * 2)); + rows[rowCount++] = result; + } + return result; + } + + void filterAndSort() + { + int newCount = 0; + for (int i = 0 ; i < rowCount; ++i) + { + if (rows[i].rowFilterIncludes()) + { + if (newCount != i) + rows[newCount] = rows[i]; + newCount++; + } + } + if (newCount != rowCount) + { + Arrays.fill(rows, newCount, rowCount, null); + totalRowCount -= (rowCount - newCount); + rowCount = newCount; + } + Arrays.sort(rows, 0, newCount, rowComparator()); + } + + int truncate(int newCount) + { + if (rowCount <= newCount) + return rowCount; + + filterAndSort(); + + if (rowCount <= newCount) + return rowCount; + + Arrays.fill(rows, newCount, rowCount, null); + totalRowCount -= (rowCount - newCount); + rowCount = newCount; + return newCount; + } + + private Comparator rowComparator() + { + Comparator cmp = dataRange.isReversed() ? metadata.comparator.reversed() : metadata.comparator; + return (a, b) -> cmp.compare(a.clustering, b.clustering); + } + + Row staticRow() + { + if (staticRow == null) + return null; + + return staticRow.materialiseAndFilter(); + } + + Iterator rows() + { + filterAndSort(); + return Arrays.stream(rows, 0, rowCount).map(SimpleRow::materialiseAndFilter).iterator(); + } + + private final class SimpleRow implements RowCollector + { + final Clustering clustering; + SomeColumns state; + + private SimpleRow(Clustering clustering) + { + this.clustering = clustering; + } + + @Override + public void lazyCollect(Consumer addToIfNeeded) + { + Invariants.require(state == null); + state = new LazyColumnsCollector(addToIfNeeded); + } + + @Override + public void eagerCollect(Consumer addToNow) + { + Invariants.require(state == null); + state = new EagerColumnsCollector(addToNow); + } + + boolean rowFilterIncludes() + { + return null != materialiseAndFilter(); + } + + Row materialiseAndFilter() + { + if (state == null) + return null; + + FilteredRow filtered = state.materialiseAndFilter(this); + state = filtered; + return filtered == null ? null : filtered.row; + } + + DecoratedKey partitionKey() + { + return SimplePartition.this.key; + } + + SimplePartitionsCollector collector() + { + return SimplePartitionsCollector.this; + } + } + } + + static abstract class SomeColumns + { + abstract FilteredRow materialiseAndFilter(SimplePartition.SimpleRow parent); + } + + static class LazyColumnsCollector extends SomeColumns + { + final Consumer lazy; + LazyColumnsCollector(Consumer lazy) + { + this.lazy = lazy; + } + + @Override + FilteredRow materialiseAndFilter(SimplePartition.SimpleRow parent) + { + return parent.collector().new EagerColumnsCollector(lazy).materialiseAndFilter(parent); + } + } + + class EagerColumnsCollector extends SomeColumns implements ColumnsCollector + { + Object[] columns = new Object[4]; + int columnCount; + + public EagerColumnsCollector(Consumer add) + { + add.accept(this); + } + + @Override + public ColumnsCollector add(String name, V1 v1, Function f1, Function f2) + { + if (v1 == null) + return this; + + ColumnMetadata cm = columnLookup.get(name); + if (cm == null) + throw new IllegalArgumentException("Unknown column name " + name); + + if (!columnFilter.fetches(cm)) + return this; + + V2 v2 = f1.apply(v1); + if (v2 == null) + return this; + + Object result = f2.apply(v2); + if (result == null) + return this; + + if (columnCount * 2 == columns.length) + columns = Arrays.copyOf(columns, columnCount * 4); + + columns[columnCount * 2] = cm; + columns[columnCount * 2 + 1] = result; + ++columnCount; + return this; + } + + @Override + FilteredRow materialiseAndFilter(SimplePartition.SimpleRow parent) + { + for (int i = 0 ; i < columnCount ; i++) + { + ColumnMetadata cm = (ColumnMetadata) columns[i * 2]; + Object value = columns[i * 2 + 1]; + ByteBuffer bb = value instanceof ByteBuffer ? (ByteBuffer)value : decompose(cm.type, value); + columns[i] = BufferCell.live(cm, timestampMicros, bb); + } + Arrays.sort(columns, 0, columnCount, (a, b) -> ColumnData.comparator.compare((BufferCell)a, (BufferCell)b)); + Object[] btree = BTree.build(BulkIterator.of(columns), columnCount, UpdateFunction.noOp); + BTreeRow row = BTreeRow.create(parent.clustering, LivenessInfo.create(timestampMicros, nowInSeconds), Row.Deletion.LIVE, btree); + if (!rowFilter.isSatisfiedBy(metadata, parent.partitionKey(), row, nowInSeconds)) + return null; + return new FilteredRow(row); + } + } + + static class FilteredRow extends SomeColumns + { + final Row row; + FilteredRow(Row row) + { + this.row = row; + } + + @Override + FilteredRow materialiseAndFilter(SimplePartition.SimpleRow parent) + { + return this; + } + } + } + + protected final TableMetadata metadata; + private final OnTimeout onTimeout; + private final Sorted sorted, sortedByPartitionKey; + + protected AbstractLazyVirtualTable(TableMetadata metadata, OnTimeout onTimeout, Sorted sorted) + { + this(metadata, onTimeout, sorted, sorted); + } + + protected AbstractLazyVirtualTable(TableMetadata metadata, OnTimeout onTimeout, Sorted sorted, Sorted sortedByPartitionKey) + { + if (!metadata.isVirtual()) + throw new IllegalArgumentException("Cannot instantiate a non-virtual table"); + + if (!metadata.keyspace.startsWith(SchemaConstants.ACCORD_KEYSPACE_NAME)) + { + // NOTE: there is nothing stopping other use cases from using this facility, but there was + // feedback on the ticket questioning the reliance on Accord integration tests for validating the API. + // If another use case wishes to use the facility, simply satisfy reviewers in this regard. See PR #4373 for details. + throw new IllegalArgumentException("This facility is only currently supported by Accord keyspaces"); + } + + this.metadata = metadata; + this.onTimeout = onTimeout; + this.sorted = sorted; + this.sortedByPartitionKey = sortedByPartitionKey; + } + + @Override + public TableMetadata metadata() + { + return metadata; + } + + public OnTimeout onTimeout() { return onTimeout; } + + protected PartitionsCollector collector(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) + { + boolean isSorted = isSorted(sorted, !dataRange.isReversed()); + boolean isSortedByPartitionKey = isSorted || isSorted(sortedByPartitionKey, !dataRange.isReversed()); + return new SimplePartitionsCollector(metadata, isSorted, isSortedByPartitionKey, dataRange, columnFilter, rowFilter, limits); + } + + + private static boolean isSorted(Sorted sorted, boolean asc) + { + return sorted == Sorted.SORTED || sorted == (asc ? Sorted.ASC : Sorted.DESC); + } + + protected abstract void collect(PartitionsCollector collector); + + @Override + public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) + { + return select(new DataRange(new Bounds<>(partitionKey, partitionKey), clusteringIndexFilter), columnFilter, rowFilter, limits); + } + + @Override + public final UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) + { + PartitionsCollector collector = collector(dataRange, columnFilter, rowFilter, limits); + try + { + collect(collector); + } + catch (InternalDoneException ignore) {} + catch (InternalTimeoutException ignore) + { + if (onTimeout != OnTimeout.BEST_EFFORT || collector.isEmpty()) + throw new ReadTimeoutException(ONE, 0, 1, false); + ClientWarn.instance.warn("Ran out of time. Returning best effort."); + } + return collector.finish(); + } + + @Override + public void apply(PartitionUpdate update) + { + throw new InvalidRequestException("Modification is not supported by table " + metadata); + } + + @Override + public void truncate() + { + throw new InvalidRequestException("Truncation is not supported by table " + metadata); + } + + @Override + public String toString() + { + return metadata().toString(); + } + + static Object[] composePartitionKeys(DecoratedKey decoratedKey, TableMetadata metadata) + { + if (metadata.partitionKeyColumns().size() == 1) + return new Object[] { metadata.partitionKeyType.compose(decoratedKey.getKey()) }; + + ByteBuffer[] split = ((CompositeType)metadata.partitionKeyType).split(decoratedKey.getKey()); + Object[] result = new Object[split.length]; + for (int i = 0 ; i < split.length ; ++i) + result[i] = metadata.partitionKeyColumns().get(i).type.compose(split[i]); + return result; + } + + static Object[] composeClusterings(ClusteringPrefix clustering, TableMetadata metadata) + { + Object[] result = new Object[clustering.size()]; + for (int i = 0 ; i < result.length ; ++i) + result[i] = metadata.clusteringColumns().get(i).type.compose(clustering.get(i), clustering.accessor()); + return result; + } + + private static ByteBuffer decompose(AbstractType type, Object value) + { + return type.decomposeUntyped(value); + } + + static DecoratedKey decomposePartitionKeys(TableMetadata metadata, Object... partitionKeys) + { + ByteBuffer partitionKey = partitionKeys.length == 1 + ? decompose(metadata.partitionKeyType, partitionKeys[0]) + : ((CompositeType) metadata.partitionKeyType).decompose(partitionKeys); + return metadata.partitioner.decorateKey(partitionKey); + } + + static Clustering decomposeClusterings(TableMetadata metadata, Object... clusteringKeys) + { + if (clusteringKeys.length == 0) + return Clustering.EMPTY; + + ByteBuffer[] clusteringByteBuffers = new ByteBuffer[clusteringKeys.length]; + for (int i = 0; i < clusteringKeys.length; i++) + { + if (clusteringKeys[i] instanceof ByteBuffer) clusteringByteBuffers[i] = (ByteBuffer) clusteringKeys[i]; + else clusteringByteBuffers[i] = decompose(metadata.clusteringColumns().get(i).type, clusteringKeys[i]); + } + return Clustering.make(clusteringByteBuffers); + } +} diff --git a/src/java/org/apache/cassandra/db/virtual/AbstractMutableLazyVirtualTable.java b/src/java/org/apache/cassandra/db/virtual/AbstractMutableLazyVirtualTable.java new file mode 100644 index 0000000000..1302f0e7bf --- /dev/null +++ b/src/java/org/apache/cassandra/db/virtual/AbstractMutableLazyVirtualTable.java @@ -0,0 +1,137 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.cassandra.db.virtual; + +import java.util.Iterator; + +import javax.annotation.Nullable; + +import accord.utils.Invariants; +import org.apache.cassandra.db.DeletionInfo; +import org.apache.cassandra.db.RangeTombstone; +import org.apache.cassandra.db.Slice; +import org.apache.cassandra.db.partitions.PartitionUpdate; +import org.apache.cassandra.db.rows.Cell; +import org.apache.cassandra.db.rows.ColumnData; +import org.apache.cassandra.db.rows.Row; +import org.apache.cassandra.exceptions.InvalidRequestException; +import org.apache.cassandra.schema.ColumnMetadata; +import org.apache.cassandra.schema.TableMetadata; + +import static org.apache.cassandra.cql3.statements.RequestValidations.invalidRequest; +import static org.apache.cassandra.db.ClusteringPrefix.Kind.STATIC_CLUSTERING; + +/** + * An abstract virtual table implementation that builds the resultset on demand and allows fine-grained source + * modification via INSERT/UPDATE, DELETE and TRUNCATE operations. + * + * Virtual table implementation need to be thread-safe has they can be called from different threads. + */ +public abstract class AbstractMutableLazyVirtualTable extends AbstractLazyVirtualTable +{ + protected AbstractMutableLazyVirtualTable(TableMetadata metadata, OnTimeout onTimeout, Sorted sorted) + { + super(metadata, onTimeout, sorted); + } + + protected AbstractMutableLazyVirtualTable(TableMetadata metadata, OnTimeout onTimeout, Sorted sorted, Sorted sortedByPartitionKey) + { + super(metadata, onTimeout, sorted, sortedByPartitionKey); + } + + protected void applyPartitionDeletion(Object[] partitionKeys) + { + throw invalidRequest("Partition deletion is not supported by table %s", metadata()); + } + + protected void applyRangeTombstone(Object[] partitionKey, Object[] start, boolean startInclusive, Object[] end, boolean endInclusive) + { + throw invalidRequest("Range deletion is not supported by table %s", metadata()); + } + + protected void applyRowDeletion(Object[] partitionKey, @Nullable Object[] clusteringKeys) + { + throw invalidRequest("Row deletion is not supported by table %s", metadata()); + } + + protected void applyRowUpdate(Object[] partitionKeys, @Nullable Object[] clusteringKeys, ColumnMetadata[] columns, Object[] values) + { + throw invalidRequest("Column modification is not supported by table %s", metadata()); + } + + private void applyRangeTombstone(Object[] pks, RangeTombstone rt) + { + Slice slice = rt.deletedSlice(); + Object[] starts = composeClusterings(slice.start(), metadata()); + Object[] ends = composeClusterings(slice.end(), metadata()); + applyRangeTombstone(pks, starts, slice.start().isInclusive(), ends, slice.end().isInclusive()); + } + + private void applyRow(Object[] pks, Row row) + { + Object[] cks = row.clustering().kind() == STATIC_CLUSTERING ? null : composeClusterings(row.clustering(), metadata()); + if (!row.deletion().isLive()) + { + applyRowDeletion(pks, cks); + } + else + { + ColumnMetadata[] columns = new ColumnMetadata[row.columnCount()]; + Object[] values = new Object[row.columnCount()]; + int i = 0; + for (ColumnData cd : row) + { + ColumnMetadata cm = cd.column(); + if (cm.isComplex()) + throw new InvalidRequestException(metadata() + " does not support complex column updates"); + Cell cell = (Cell)cd; + columns[i] = cm; + if (!cell.isTombstone()) + values[i] = cm.type.compose(cell.value(), cell.accessor()); + ++i; + } + Invariants.require(i == columns.length); + applyRowUpdate(pks, cks, columns, values); + } + } + + public void apply(PartitionUpdate update) + { + TableMetadata metadata = metadata(); + Object[] pks = composePartitionKeys(update.partitionKey(), metadata); + + DeletionInfo deletionInfo = update.deletionInfo(); + if (!deletionInfo.getPartitionDeletion().isLive()) + { + applyPartitionDeletion(pks); + } + else if (deletionInfo.hasRanges()) + { + Iterator iter = deletionInfo.rangeIterator(false); + while (iter.hasNext()) + applyRangeTombstone(pks, iter.next()); + } + else + { + for (Row row : update) + applyRow(pks, row); + if (!update.staticRow().isEmpty()) + applyRow(pks, update.staticRow()); + } + } +} diff --git a/src/java/org/apache/cassandra/db/virtual/AbstractVirtualTable.java b/src/java/org/apache/cassandra/db/virtual/AbstractVirtualTable.java index a32ea67ab6..2c4a04e5a0 100644 --- a/src/java/org/apache/cassandra/db/virtual/AbstractVirtualTable.java +++ b/src/java/org/apache/cassandra/db/virtual/AbstractVirtualTable.java @@ -29,6 +29,7 @@ import org.apache.cassandra.db.EmptyIterators; import org.apache.cassandra.db.PartitionPosition; import org.apache.cassandra.db.filter.ClusteringIndexFilter; import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.DataLimits; import org.apache.cassandra.db.filter.RowFilter; import org.apache.cassandra.db.partitions.AbstractUnfilteredPartitionIterator; import org.apache.cassandra.db.partitions.PartitionUpdate; @@ -76,7 +77,7 @@ public abstract class AbstractVirtualTable implements VirtualTable } @Override - public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter) + public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) { Partition partition = data(partitionKey).getPartition(partitionKey); @@ -89,7 +90,7 @@ public abstract class AbstractVirtualTable implements VirtualTable } @Override - public final UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter) + public final UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) { DataSet data = data(); diff --git a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java index 87c69688eb..1f796bdd15 100644 --- a/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java +++ b/src/java/org/apache/cassandra/db/virtual/AccordDebugKeyspace.java @@ -18,45 +18,36 @@ package org.apache.cassandra.db.virtual; import java.nio.ByteBuffer; -import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; +import java.util.Comparator; import java.util.Date; -import java.util.HashSet; -import java.util.Iterator; import java.util.LinkedHashMap; import java.util.List; import java.util.Locale; import java.util.Map; import java.util.Objects; import java.util.Optional; -import java.util.Set; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeoutException; import java.util.function.BiConsumer; import java.util.function.Function; +import java.util.function.UnaryOperator; import java.util.stream.Collectors; import javax.annotation.Nullable; -import com.google.common.collect.BoundType; -import com.google.common.collect.Range; -import com.google.common.collect.Sets; - import accord.coordinate.AbstractCoordination; import accord.coordinate.Coordination; import accord.coordinate.Coordinations; import accord.coordinate.PrepareRecovery; import accord.coordinate.tracking.AbstractTracker; +import accord.local.cfk.CommandsForKey.TxnInfo; import accord.primitives.RoutingKeys; import accord.utils.SortedListMap; -import org.apache.cassandra.cql3.Operator; -import org.apache.cassandra.db.EmptyIterators; -import org.apache.cassandra.db.filter.ClusteringIndexFilter; -import org.apache.cassandra.db.filter.ColumnFilter; -import org.apache.cassandra.db.filter.RowFilter; -import org.apache.cassandra.db.partitions.SingletonUnfilteredPartitionIterator; -import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; -import org.apache.cassandra.db.rows.UnfilteredRowIterator; +import org.apache.cassandra.db.PartitionPosition; +import org.apache.cassandra.db.marshal.CompositeType; +import org.apache.cassandra.db.marshal.TxnIdUtf8Type; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -91,10 +82,8 @@ import accord.primitives.Participants; import accord.primitives.ProgressToken; import accord.primitives.Route; import accord.primitives.SaveStatus; -import accord.primitives.Status; import accord.primitives.Timestamp; import accord.primitives.TxnId; -import accord.utils.Invariants; import accord.utils.UnhandledEnum; import accord.utils.async.AsyncChain; import accord.utils.async.AsyncChains; @@ -109,11 +98,11 @@ import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.db.marshal.Int32Type; import org.apache.cassandra.db.marshal.UTF8Type; -import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.LocalPartitioner; import org.apache.cassandra.dht.NormalizedRanges; import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.InvalidRequestException; +import org.apache.cassandra.schema.ColumnMetadata; import org.apache.cassandra.schema.Schema; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableMetadata; @@ -125,7 +114,9 @@ import org.apache.cassandra.service.accord.AccordJournal; import org.apache.cassandra.service.accord.AccordKeyspace; import org.apache.cassandra.service.accord.AccordService; import org.apache.cassandra.service.accord.AccordTracing; -import org.apache.cassandra.service.accord.CommandStoreTxnBlockedGraph; +import org.apache.cassandra.service.accord.DebugBlockedTxns; +import org.apache.cassandra.service.accord.IAccordService; +import org.apache.cassandra.service.accord.JournalKey; import org.apache.cassandra.service.accord.api.AccordAgent; import org.apache.cassandra.service.accord.api.TokenKey; import org.apache.cassandra.service.consensus.migration.ConsensusMigrationState; @@ -150,8 +141,12 @@ import static com.google.common.collect.ImmutableList.toImmutableList; import static java.lang.String.format; import static java.util.concurrent.TimeUnit.NANOSECONDS; import static org.apache.cassandra.cql3.statements.RequestValidations.invalidRequest; +import static org.apache.cassandra.db.virtual.AbstractLazyVirtualTable.OnTimeout.BEST_EFFORT; +import static org.apache.cassandra.db.virtual.AbstractLazyVirtualTable.OnTimeout.FAIL; +import static org.apache.cassandra.db.virtual.VirtualTable.Sorted.ASC; +import static org.apache.cassandra.db.virtual.VirtualTable.Sorted.SORTED; +import static org.apache.cassandra.db.virtual.VirtualTable.Sorted.UNSORTED; import static org.apache.cassandra.schema.SchemaConstants.VIRTUAL_ACCORD_DEBUG; -import static org.apache.cassandra.utils.Clock.Global.currentTimeMillis; import static org.apache.cassandra.utils.MonotonicClock.Global.approxTime; public class AccordDebugKeyspace extends VirtualKeyspace @@ -175,6 +170,8 @@ public class AccordDebugKeyspace extends VirtualKeyspace public static final String TXN_TRACES = "txn_traces"; public static final String TXN_OPS = "txn_ops"; + private static final Function TO_STRING = AccordDebugKeyspace::toStringOrNull; + public static final AccordDebugKeyspace instance = new AccordDebugKeyspace(); private AccordDebugKeyspace() @@ -202,7 +199,7 @@ public class AccordDebugKeyspace extends VirtualKeyspace } // TODO (desired): human readable packed key tracker (but requires loading Txn, so might be preferable to only do conditionally) - public static final class ExecutorsTable extends AbstractVirtualTable + public static final class ExecutorsTable extends AbstractLazyVirtualTable { private ExecutorsTable() { @@ -221,43 +218,44 @@ public class AccordDebugKeyspace extends VirtualKeyspace " keys_loading text,\n" + " keys_loading_for text,\n" + " PRIMARY KEY (executor_id, status, position, unique_position)" + - ')', UTF8Type.instance)); + ')', Int32Type.instance), FAIL, ASC); } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { AccordCommandStores commandStores = (AccordCommandStores) AccordService.instance().node().commandStores(); - SimpleDataSet ds = new SimpleDataSet(metadata()); - + // TODO (desired): we can easily also support sorted collection for DESC queries for (AccordExecutor executor : commandStores.executors()) { - int uniquePos = 0; int executorId = executor.executorId(); - AccordExecutor.TaskInfo prev = null; - for (AccordExecutor.TaskInfo info : executor.taskSnapshot()) - { - if (prev != null && info.status() == prev.status() && info.position() == prev.position()) ++uniquePos; - else uniquePos = 0; - prev = info; - PreLoadContext preLoadContext = info.preLoadContext(); - ds.row(executorId, Objects.toString(info.status()), info.position(), uniquePos) - .column("description", info.describe()) - .column("command_store_id", info.commandStoreId()) - .column("txn_id", preLoadContext == null ? null : toStringOrNull(preLoadContext.primaryTxnId())) - .column("txn_id_additional", preLoadContext == null ? null : toStringOrNull(preLoadContext.additionalTxnId())) - .column("keys", preLoadContext == null ? null : toStringOrNull(preLoadContext.keys())) - .column("keys_loading", preLoadContext == null ? null : toStringOrNull(preLoadContext.loadKeys())) - .column("keys_loading_for", preLoadContext == null ? null : toStringOrNull(preLoadContext.loadKeysFor())) - ; - } + collector.partition(executorId).collect(rows -> { + int uniquePos = 0; + AccordExecutor.TaskInfo prev = null; + for (AccordExecutor.TaskInfo info : executor.taskSnapshot()) + { + if (prev != null && info.status() == prev.status() && info.position() == prev.position()) ++uniquePos; + else uniquePos = 0; + prev = info; + PreLoadContext preLoadContext = info.preLoadContext(); + rows.add(info.status().name(), info.position(), uniquePos) + .lazyCollect(columns -> { + columns.add("description", info.describe()) + .add("command_store_id", info.commandStoreId()) + .add("txn_id", preLoadContext, PreLoadContext::primaryTxnId, TO_STRING) + .add("txn_id_additional", preLoadContext, PreLoadContext::additionalTxnId, TO_STRING) + .add("keys", preLoadContext, PreLoadContext::keys, TO_STRING) + .add("keys_loading", preLoadContext, PreLoadContext::loadKeys, TO_STRING) + .add("keys_loading_for", preLoadContext, PreLoadContext::loadKeysFor, TO_STRING); + }); + } + }); } - return ds; } } // TODO (desired): human readable packed key tracker (but requires loading Txn, so might be preferable to only do conditionally) - public static final class CoordinationsTable extends AbstractVirtualTable + public static final class CoordinationsTable extends AbstractLazyVirtualTable { private CoordinationsTable() { @@ -275,48 +273,37 @@ public class AccordDebugKeyspace extends VirtualKeyspace " replies text,\n" + " tracker text,\n" + " PRIMARY KEY (txn_id, kind, coordination_id)" + - ')', UTF8Type.instance)); + ')', TxnIdUtf8Type.instance), FAIL, UNSORTED); } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { Coordinations coordinations = AccordService.instance().node().coordinations(); - SimpleDataSet ds = new SimpleDataSet(metadata()); for (Coordination c : coordinations) { - ds.row(toStringOrNull(c.txnId()), c.kind().toString(), c.coordinationId()) - .column("nodes", toStringOrNull(c.nodes())) - .column("nodes_inflight", toStringOrNull(c.inflight())) - .column("nodes_contacted", toStringOrNull(c.contacted())) - .column("description", c.describe()) - .column("participants", toStringOrNull(c.scope())) - .column("replies", summarise(c.replies())) - .column("tracker", summarise(c.tracker())); + collector.row(toStringOrNull(c.txnId()), c.kind().toString(), c.coordinationId()) + .lazyCollect(columns -> { + columns.add("nodes", c, Coordination::nodes, TO_STRING) + .add("nodes_inflight", c, Coordination::inflight, TO_STRING) + .add("nodes_contacted", c, Coordination::contacted, TO_STRING) + .add("description", c, Coordination::describe, TO_STRING) + .add("participants", c, Coordination::scope, TO_STRING) + .add("replies", c, Coordination::replies, CoordinationsTable::summarise) + .add("tracker", c, Coordination::tracker, AbstractTracker::summariseTracker); + }); } - return ds; } - private static String summarise(@Nullable SortedListMap replies) + private static String summarise(SortedListMap replies) { - if (replies == null) - return null; return AbstractCoordination.summariseReplies(replies, 60); } - - private static String summarise(@Nullable AbstractTracker tracker) - { - if (tracker == null) - return null; - return tracker.summariseTracker(); - } } - - // TODO (desired): don't report null as "null" - public static final class CommandsForKeyTable extends AbstractVirtualTable implements AbstractVirtualTable.DataSet + private static abstract class AbstractCommandsForKeyTable extends AbstractLazyVirtualTable { - static class Entry + static class Entry implements Comparable { final int commandStoreId; final CommandsForKey cfk; @@ -326,7 +313,52 @@ public class AccordDebugKeyspace extends VirtualKeyspace this.commandStoreId = commandStoreId; this.cfk = cfk; } + + @Override + public int compareTo(Entry that) + { + return Integer.compare(this.commandStoreId, that.commandStoreId); + } } + + AbstractCommandsForKeyTable(TableMetadata metadata) + { + super(metadata, BEST_EFFORT, SORTED); + } + + abstract void collect(PartitionCollector partition, int commandStoreId, CommandsForKey cfk); + + @Override + public void collect(PartitionsCollector collector) + { + Object[] partitionKey = collector.singlePartitionKey(); + if (partitionKey == null) + throw new InvalidRequestException(metadata + " currently only supports querying single partitions"); + + TokenKey key = TokenKey.parse((String) partitionKey[0], DatabaseDescriptor.getPartitioner()); + + List cfks = new CopyOnWriteArrayList<>(); + CommandStores commandStores = AccordService.instance().node().commandStores(); + AccordService.getBlocking(commandStores.forEach("commands_for_key table query", RoutingKeys.of(key), Long.MIN_VALUE, Long.MAX_VALUE, safeStore -> { + SafeCommandsForKey safeCfk = safeStore.get(key); + CommandsForKey cfk = safeCfk.current(); + if (cfk == null) + return; + + cfks.add(new Entry(safeStore.commandStore().id(), cfk)); + })); + + if (cfks.isEmpty()) + return; + + cfks.sort(collector.dataRange().isReversed() ? Comparator.reverseOrder() : Comparator.naturalOrder()); + for (Entry entry : cfks) + collect(collector.partition(partitionKey[0]), entry.commandStoreId, entry.cfk); + } + } + + public static final class CommandsForKeyTable extends AbstractCommandsForKeyTable + { private CommandsForKeyTable() { super(parse(VIRTUAL_ACCORD_DEBUG, COMMANDS_FOR_KEY, @@ -347,65 +379,27 @@ public class AccordDebugKeyspace extends VirtualKeyspace } @Override - public DataSet data() + void collect(PartitionCollector partition, int commandStoreId, CommandsForKey cfk) { - return this; - } - - @Override - public boolean isEmpty() - { - return false; - } - - @Override - public Partition getPartition(DecoratedKey partitionKey) - { - String keyStr = UTF8Type.instance.compose(partitionKey.getKey()); - TokenKey key = TokenKey.parse(keyStr, DatabaseDescriptor.getPartitioner()); - - List cfks = new CopyOnWriteArrayList<>(); - CommandStores commandStores = AccordService.instance().node().commandStores(); - AccordService.getBlocking(commandStores.forEach("commands_for_key table query", RoutingKeys.of(key), Long.MIN_VALUE, Long.MAX_VALUE, safeStore -> { - SafeCommandsForKey safeCfk = safeStore.get(key); - CommandsForKey cfk = safeCfk.current(); - if (cfk == null) - return; - - cfks.add(new Entry(safeStore.commandStore().id(), cfk)); - })); - - if (cfks.isEmpty()) - return null; - - SimpleDataSet ds = new SimpleDataSet(metadata); - for (Entry e : cfks) - { - CommandsForKey cfk = e.cfk; + partition.collect(rows -> { for (int i = 0 ; i < cfk.size() ; ++i) { - CommandsForKey.TxnInfo txn = cfk.get(i); - ds.row(keyStr, e.commandStoreId, toStringOrNull(txn.plainTxnId())) - .column("ballot", toStringOrNull(txn.ballot())) - .column("deps_known_before", toStringOrNull(txn.depsKnownUntilExecuteAt())) - .column("flags", flags(txn)) - .column("execute_at", toStringOrNull(txn.plainExecuteAt())) - .column("missing", Arrays.toString(txn.missing())) - .column("status", toStringOrNull(txn.status())) - .column("status_overrides", txn.statusOverrides() == 0 ? null : ("0x" + Integer.toHexString(txn.statusOverrides()))); + TxnInfo txn = cfk.get(i); + rows.add(commandStoreId, txn.plainTxnId().toString()) + .lazyCollect(columns -> { + columns.add("ballot", txn.ballot(), AccordDebugKeyspace::toStringOrNull) + .add("deps_known_before", txn, TxnInfo::depsKnownUntilExecuteAt, TO_STRING) + .add("flags", txn, CommandsForKeyTable::flags) + .add("execute_at", txn, TxnInfo::plainExecuteAt, TO_STRING) + .add("missing", txn, TxnInfo::missing, Arrays::toString) + .add("status", txn, TxnInfo::status, TO_STRING) + .add("status_overrides", txn.statusOverrides() == 0 ? null : ("0x" + Integer.toHexString(txn.statusOverrides()))); + }); } - } - - return ds.getPartition(partitionKey); + }); } - @Override - public Iterator getPartitions(DataRange range) - { - throw new UnsupportedOperationException(); - } - - private static String flags(CommandsForKey.TxnInfo txn) + private static String flags(TxnInfo txn) { StringBuilder sb = new StringBuilder(); if (!txn.mayExecute()) @@ -426,21 +420,8 @@ public class AccordDebugKeyspace extends VirtualKeyspace } } - - // TODO (expected): test this table - public static final class CommandsForKeyUnmanagedTable extends AbstractVirtualTable implements AbstractVirtualTable.DataSet + public static final class CommandsForKeyUnmanagedTable extends AbstractCommandsForKeyTable { - static class Entry - { - final int commandStoreId; - final CommandsForKey cfk; - - Entry(int commandStoreId, CommandsForKey cfk) - { - this.commandStoreId = commandStoreId; - this.cfk = cfk; - } - } private CommandsForKeyUnmanagedTable() { super(parse(VIRTUAL_ACCORD_DEBUG, COMMANDS_FOR_KEY_UNMANAGED, @@ -456,70 +437,30 @@ public class AccordDebugKeyspace extends VirtualKeyspace } @Override - public DataSet data() + void collect(PartitionCollector partition, int commandStoreId, CommandsForKey cfk) { - return this; - } - - @Override - public boolean isEmpty() - { - return false; - } - - @Override - public Partition getPartition(DecoratedKey partitionKey) - { - String keyStr = UTF8Type.instance.compose(partitionKey.getKey()); - TokenKey key = TokenKey.parse(keyStr, DatabaseDescriptor.getPartitioner()); - - List cfks = new CopyOnWriteArrayList<>(); - CommandStores commandStores = AccordService.instance().node().commandStores(); - AccordService.getBlocking(commandStores.forEach("commands_for_key_unmanaged table query", RoutingKeys.of(key), Long.MIN_VALUE, Long.MAX_VALUE, safeStore -> { - SafeCommandsForKey safeCfk = safeStore.get(key); - CommandsForKey cfk = safeCfk.current(); - if (cfk == null) - return; - - cfks.add(new Entry(safeStore.commandStore().id(), cfk)); - })); - - if (cfks.isEmpty()) - return null; - - SimpleDataSet ds = new SimpleDataSet(metadata); - for (Entry e : cfks) - { - CommandsForKey cfk = e.cfk; + partition.collect(rows -> { for (int i = 0 ; i < cfk.unmanagedCount() ; ++i) { CommandsForKey.Unmanaged txn = cfk.getUnmanaged(i); - ds.row(keyStr, e.commandStoreId, toStringOrNull(txn.txnId)) - .column("waiting_until", toStringOrNull(txn.waitingUntil)) - .column("waiting_until_status", toStringOrNull(txn.pending)); + rows.add(commandStoreId, toStringOrNull(txn.txnId)) + .lazyCollect(columns -> { + columns.add("waiting_until", txn.waitingUntil, TO_STRING) + .add("waiting_until_status", txn.pending, TO_STRING); + }); } - } - - return ds.getPartition(partitionKey); - } - - @Override - public Iterator getPartitions(DataRange range) - { - throw new UnsupportedOperationException(); + }); } } - - public static final class DurabilityServiceTable extends AbstractVirtualTable + public static final class DurabilityServiceTable extends AbstractLazyVirtualTable { private DurabilityServiceTable() { super(parse(VIRTUAL_ACCORD_DEBUG, DURABILITY_SERVICE, "Accord per-Range Durability Service State", "CREATE TABLE %s (\n" + - " keyspace_name text,\n" + - " table_name text,\n" + + " table_id text,\n" + " token_start 'TokenUtf8Type',\n" + " token_end 'TokenUtf8Type',\n" + " last_started_at bigint,\n" + @@ -538,82 +479,77 @@ public class AccordDebugKeyspace extends VirtualKeyspace " current_splits int,\n" + " stopping boolean,\n" + " stopped boolean,\n" + - " PRIMARY KEY (keyspace_name, table_name, token_start)" + - ')', UTF8Type.instance)); + " PRIMARY KEY (table_id, token_start)" + + ')', UTF8Type.instance), FAIL, UNSORTED); } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { ShardDurability.ImmutableView view = ((AccordService) AccordService.instance()).shardDurability(); - SimpleDataSet ds = new SimpleDataSet(metadata()); while (view.advance()) { TableId tableId = (TableId) view.shard().range.start().prefix(); - TableMetadata tableMetadata = tableMetadata(tableId); - ds.row(keyspace(tableMetadata), table(tableId, tableMetadata), printToken(view.shard().range.start())) - .column("token_end", printToken(view.shard().range.end())) - .column("last_started_at", approxTime.translate().toMillisSinceEpoch(view.lastStartedAtMicros() * 1000)) - .column("cycle_started_at", approxTime.translate().toMillisSinceEpoch(view.cycleStartedAtMicros() * 1000)) - .column("retries", view.retries()) - .column("min", Objects.toString(view.min())) - .column("requested_by", Objects.toString(view.requestedBy())) - .column("active", Objects.toString(view.active())) - .column("waiting", Objects.toString(view.waiting())) - .column("node_offset", view.nodeOffset()) - .column("cycle_offset", view.cycleOffset()) - .column("active_index", view.activeIndex()) - .column("next_index", view.nextIndex()) - .column("next_to_index", view.toIndex()) - .column("end_index", view.cycleLength()) - .column("current_splits", view.currentSplits()) - .column("stopping", view.stopping()) - .column("stopped", view.stopped()) - ; + collector.row(tableId.toString(), printToken(view.shard().range.start())) + .eagerCollect(columns -> { + columns.add("token_end", printToken(view.shard().range.end())) + .add("last_started_at", approxTime.translate().toMillisSinceEpoch(view.lastStartedAtMicros() * 1000)) + .add("cycle_started_at", approxTime.translate().toMillisSinceEpoch(view.cycleStartedAtMicros() * 1000)) + .add("retries", view.retries()) + .add("min", Objects.toString(view.min())) + .add("requested_by", Objects.toString(view.requestedBy())) + .add("active", Objects.toString(view.active())) + .add("waiting", Objects.toString(view.waiting())) + .add("node_offset", view.nodeOffset()) + .add("cycle_offset", view.cycleOffset()) + .add("active_index", view.activeIndex()) + .add("next_index", view.nextIndex()) + .add("next_to_index", view.toIndex()) + .add("end_index", view.cycleLength()) + .add("current_splits", view.currentSplits()) + .add("stopping", view.stopping()) + .add("stopped", view.stopped()); + }); } - return ds; } } - public static final class DurableBeforeTable extends AbstractVirtualTable + public static final class DurableBeforeTable extends AbstractLazyVirtualTable { private DurableBeforeTable() { super(parse(VIRTUAL_ACCORD_DEBUG, DURABLE_BEFORE, "Accord Node's DurableBefore State", "CREATE TABLE %s (\n" + - " keyspace_name text,\n" + - " table_name text,\n" + + " table_id text,\n" + " token_start 'TokenUtf8Type',\n" + " token_end 'TokenUtf8Type',\n" + " quorum 'TxnIdUtf8Type',\n" + " universal 'TxnIdUtf8Type',\n" + - " PRIMARY KEY (keyspace_name, table_name, token_start)" + - ')', UTF8Type.instance)); + " PRIMARY KEY (table_id, token_start)" + + ')', UTF8Type.instance), FAIL, UNSORTED); } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { DurableBefore durableBefore = AccordService.instance().node().durableBefore(); - return durableBefore.foldlWithBounds( - (entry, ds, start, end) -> { + durableBefore.foldlWithBounds( + (entry, ignore, start, end) -> { TableId tableId = (TableId) start.prefix(); - TableMetadata tableMetadata = tableMetadata(tableId); - ds.row(keyspace(tableMetadata), table(tableId, tableMetadata), printToken(start)) - .column("token_end", printToken(end)) - .column("quorum", entry.quorumBefore.toString()) - .column("universal", entry.universalBefore.toString()); - return ds; - }, - new SimpleDataSet(metadata()), - ignore -> false - ); + collector.row(tableId.toString(), printToken(start)) + .lazyCollect(columns -> { + columns.add("token_end", end, AccordDebugKeyspace::printToken) + .add("quorum", entry.quorumBefore, TO_STRING) + .add("universal", entry.universalBefore, TO_STRING); + }); + return null; + }, null, ignore -> false); } } - public static final class ExecutorCacheTable extends AbstractVirtualTable + public static final class ExecutorCacheTable extends AbstractLazyVirtualTable { private ExecutorCacheTable() { @@ -626,77 +562,76 @@ public class AccordDebugKeyspace extends VirtualKeyspace " hits bigint,\n" + " misses bigint,\n" + " PRIMARY KEY (executor_id, scope)" + - ')', Int32Type.instance)); + ')', Int32Type.instance), FAIL, UNSORTED); } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { AccordCommandStores stores = (AccordCommandStores) AccordService.instance().node().commandStores(); - SimpleDataSet ds = new SimpleDataSet(metadata()); for (AccordExecutor executor : stores.executors()) { try (AccordExecutor.ExclusiveGlobalCaches cache = executor.lockCaches()) { - addRow(ds, executor.executorId(), "commands", cache.commands.statsSnapshot()); - addRow(ds, executor.executorId(), AccordKeyspace.COMMANDS_FOR_KEY, cache.commandsForKey.statsSnapshot()); + addRow(collector, executor.executorId(), "commands", cache.commands.statsSnapshot()); + addRow(collector, executor.executorId(), AccordKeyspace.COMMANDS_FOR_KEY, cache.commandsForKey.statsSnapshot()); } } - return ds; } - private static void addRow(SimpleDataSet ds, int executorId, String scope, AccordCache.ImmutableStats stats) + private static void addRow(PartitionsCollector collector, int executorId, String scope, AccordCache.ImmutableStats stats) { - ds.row(executorId, scope) - .column("queries", stats.hits + stats.misses) - .column("hits", stats.hits) - .column("misses", stats.misses); + collector.row(executorId, scope) + .eagerCollect(columns -> { + columns.add("queries", stats.hits + stats.misses) + .add("hits", stats.hits) + .add("misses", stats.misses); + }); } } - - public static final class MaxConflictsTable extends AbstractVirtualTable + public static final class MaxConflictsTable extends AbstractLazyVirtualTable { private MaxConflictsTable() { super(parse(VIRTUAL_ACCORD_DEBUG, MAX_CONFLICTS, "Accord per-CommandStore MaxConflicts State", "CREATE TABLE %s (\n" + - " keyspace_name text,\n" + - " table_name text,\n" + " command_store_id bigint,\n" + " token_start 'TokenUtf8Type',\n" + + " table_id text,\n" + " token_end 'TokenUtf8Type',\n" + " timestamp text,\n" + - " PRIMARY KEY (keyspace_name, table_name, command_store_id, token_start)" + - ')', UTF8Type.instance)); + " PRIMARY KEY (command_store_id, token_start)" + + ')', Int32Type.instance), FAIL, ASC); } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { CommandStores commandStores = AccordService.instance().node().commandStores(); - SimpleDataSet dataSet = new SimpleDataSet(metadata()); for (CommandStore commandStore : commandStores.all()) { int commandStoreId = commandStore.id(); MaxConflicts maxConflicts = commandStore.unsafeGetMaxConflicts(); TableId tableId = ((AccordCommandStore) commandStore).tableId(); - TableMetadata tableMetadata = tableMetadata(tableId); + String tableIdStr = tableId.toString(); - maxConflicts.foldlWithBounds( - (timestamp, ds, start, end) -> { - return ds.row(keyspace(tableMetadata), table(tableId, tableMetadata), commandStoreId, printToken(start)) - .column("token_end", printToken(end)) - .column("timestamp", timestamp.toString()) - ; - }, - dataSet, - ignore -> false - ); + collector.partition(commandStoreId).collect(rows -> { + maxConflicts.foldlWithBounds( + (timestamp, rs, start, end) -> { + rows.add(printToken(start)) + .lazyCollect(columns -> { + columns.add("token_end", end, AccordDebugKeyspace::printToken) + .add("table_id", tableIdStr) + .add("timestamp", timestamp, TO_STRING); + }); + return rows; + }, rows, ignore -> false + ); + }); } - return dataSet; } } @@ -789,18 +724,16 @@ public class AccordDebugKeyspace extends VirtualKeyspace } // TODO (desired): human readable packed key tracker (but requires loading Txn, so might be preferable to only do conditionally) - public static final class ProgressLogTable extends AbstractVirtualTable + public static final class ProgressLogTable extends AbstractLazyVirtualTable { private ProgressLogTable() { super(parse(VIRTUAL_ACCORD_DEBUG, PROGRESS_LOG, "Accord per-CommandStore ProgressLog State", "CREATE TABLE %s (\n" + - " keyspace_name text,\n" + - " table_name text,\n" + - " table_id text,\n" + " command_store_id int,\n" + " txn_id 'TxnIdUtf8Type',\n" + + " table_id text,\n" + // Timer + BaseTxnState " contact_everyone boolean,\n" + // WaitingState @@ -816,43 +749,46 @@ public class AccordDebugKeyspace extends VirtualKeyspace " home_progress text,\n" + " home_retry_counter int,\n" + " home_scheduled_at timestamp,\n" + - " PRIMARY KEY (keyspace_name, table_name, table_id, command_store_id, txn_id)" + - ')', UTF8Type.instance)); + " PRIMARY KEY (command_store_id, txn_id)" + + ')', Int32Type.instance), FAIL, ASC); } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { CommandStores commandStores = AccordService.instance().node().commandStores(); - SimpleDataSet ds = new SimpleDataSet(metadata()); for (CommandStore commandStore : commandStores.all()) { DefaultProgressLog.ImmutableView view = ((DefaultProgressLog) commandStore.unsafeProgressLog()).immutableView(); TableId tableId = ((AccordCommandStore)commandStore).tableId(); String tableIdStr = tableId.toString(); - TableMetadata tableMetadata = tableMetadata(tableId); - while (view.advance()) - { - ds.row(keyspace(tableMetadata), table(tableId, tableMetadata), tableIdStr, view.commandStoreId(), view.txnId().toString()) - .column("contact_everyone", view.contactEveryone()) - .column("waiting_is_uninitialised", view.isWaitingUninitialised()) - .column("waiting_blocked_until", view.waitingIsBlockedUntil().name()) - .column("waiting_home_satisfies", view.waitingHomeSatisfies().name()) - .column("waiting_progress", view.waitingProgress().name()) - .column("waiting_retry_counter", view.waitingRetryCounter()) - .column("waiting_packed_key_tracker_bits", Long.toBinaryString(view.waitingPackedKeyTrackerBits())) - .column("waiting_scheduled_at", toTimestamp(view.timerScheduledAt(TxnStateKind.Waiting))) - .column("home_phase", view.homePhase().name()) - .column("home_progress", view.homeProgress().name()) - .column("home_retry_counter", view.homeRetryCounter()) - .column("home_scheduled_at", toTimestamp(view.timerScheduledAt(TxnStateKind.Home))) - ; - } + collector.partition(commandStore.id()).collect(collect -> { + while (view.advance()) + { + // TODO (required): view should return an immutable per-row view so that we can call lazyAdd + collect.add(view.txnId().toString()) + .eagerCollect(columns -> { + columns.add("table_id", tableIdStr) + .add("contact_everyone", view.contactEveryone()) + .add("waiting_is_uninitialised", view.isWaitingUninitialised()) + .add("waiting_blocked_until", view.waitingIsBlockedUntil().name()) + .add("waiting_home_satisfies", view.waitingHomeSatisfies().name()) + .add("waiting_progress", view.waitingProgress().name()) + .add("waiting_retry_counter", view.waitingRetryCounter()) + .add("waiting_packed_key_tracker_bits", Long.toBinaryString(view.waitingPackedKeyTrackerBits())) + .add("waiting_scheduled_at", view.timerScheduledAt(TxnStateKind.Waiting), ProgressLogTable::toTimestamp) + .add("home_phase", view.homePhase().name()) + .add("home_progress", view.homeProgress().name()) + .add("home_retry_counter", view.homeRetryCounter()) + .add("home_scheduled_at", view.timerScheduledAt(TxnStateKind.Home), ProgressLogTable::toTimestamp); + }); + } + + }); } - return ds; } - private Date toTimestamp(Long deadline) + private static Date toTimestamp(Long deadline) { if (deadline == null) return null; @@ -862,19 +798,17 @@ public class AccordDebugKeyspace extends VirtualKeyspace } } - public static final class RedundantBeforeTable extends AbstractVirtualTable + public static final class RedundantBeforeTable extends AbstractLazyVirtualTable { private RedundantBeforeTable() { super(parse(VIRTUAL_ACCORD_DEBUG, REDUNDANT_BEFORE, "Accord per-CommandStore RedundantBefore State", "CREATE TABLE %s (\n" + - " keyspace_name text,\n" + - " table_name text,\n" + - " table_id text,\n" + - " token_start 'TokenUtf8Type',\n" + - " token_end 'TokenUtf8Type',\n" + " command_store_id int,\n" + + " token_start 'TokenUtf8Type',\n" + + " table_id text,\n" + + " token_end 'TokenUtf8Type',\n" + " start_epoch bigint,\n" + " end_epoch bigint,\n" + " gc_before 'TxnIdUtf8Type',\n" + @@ -888,97 +822,88 @@ public class AccordDebugKeyspace extends VirtualKeyspace " locally_witnessed 'TxnIdUtf8Type',\n" + " log_unavailable 'TxnIdUtf8Type',\n" + " unready 'TxnIdUtf8Type',\n" + - " stale_until 'TxnIdUtf8Type',\n" + - " PRIMARY KEY (keyspace_name, table_name, table_id, command_store_id, token_start)" + - ')', UTF8Type.instance)); + " stale_until_at_least 'TxnIdUtf8Type',\n" + + " PRIMARY KEY (command_store_id, token_start)" + + ')', Int32Type.instance), FAIL, ASC); } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { CommandStores commandStores = AccordService.instance().node().commandStores(); - SimpleDataSet dataSet = new SimpleDataSet(metadata()); for (CommandStore commandStore : commandStores.all()) { int commandStoreId = commandStore.id(); - TableId tableId = ((AccordCommandStore)commandStore).tableId(); - String tableIdStr = tableId.toString(); - TableMetadata tableMetadata = tableMetadata(tableId); - String keyspace = keyspace(tableMetadata); - String table = table(tableId, tableMetadata); - commandStore.unsafeGetRedundantBefore().foldl( - (entry, ds) -> { - ds.row(keyspace, table, tableIdStr, commandStoreId, printToken(entry.range.start())) - .column("token_end", printToken(entry.range.end())) - .column("start_epoch", entry.startEpoch) - .column("end_epoch", entry.endEpoch) - .column("gc_before", entry.maxBound(GC_BEFORE).toString()) - .column("shard_applied", entry.maxBound(SHARD_APPLIED).toString()) - .column("quorum_applied", entry.maxBound(QUORUM_APPLIED).toString()) - .column("locally_applied", entry.maxBound(LOCALLY_APPLIED).toString()) - .column("locally_durable_to_command_store", entry.maxBound(LOCALLY_DURABLE_TO_COMMAND_STORE).toString()) - .column("locally_durable_to_data_store", entry.maxBound(LOCALLY_DURABLE_TO_DATA_STORE).toString()) - .column("locally_redundant", entry.maxBound(LOCALLY_REDUNDANT).toString()) - .column("locally_synced", entry.maxBound(LOCALLY_SYNCED).toString()) - .column("locally_witnessed", entry.maxBound(LOCALLY_WITNESSED).toString()) - .column("log_unavailable", entry.maxBound(LOG_UNAVAILABLE).toString()) - .column("unready", entry.maxBound(UNREADY).toString()) - .column("stale_until", entry.staleUntilAtLeast != null ? entry.staleUntilAtLeast.toString() : null); - return ds; - }, - dataSet, - ignore -> false - ); + collector.partition(commandStoreId).collect(rows -> { + TableId tableId = ((AccordCommandStore)commandStore).tableId(); + String tableIdStr = tableId.toString(); + commandStore.unsafeGetRedundantBefore().foldl( + (entry, rs) -> { + rs.add(printToken(entry.range.start())).lazyCollect(columns -> { + columns.add("table_id", tableIdStr) + .add("token_end", entry.range.end(), AccordDebugKeyspace::printToken) + .add("start_epoch", entry.startEpoch) + .add("end_epoch", entry.endEpoch) + .add("gc_before", entry, e -> e.maxBound(GC_BEFORE), TO_STRING) + .add("shard_applied", entry, e -> e.maxBound(SHARD_APPLIED), TO_STRING) + .add("quorum_applied", entry, e -> e.maxBound(QUORUM_APPLIED), TO_STRING) + .add("locally_applied", entry, e -> e.maxBound(LOCALLY_APPLIED), TO_STRING) + .add("locally_durable_to_command_store", entry, e -> e.maxBound(LOCALLY_DURABLE_TO_COMMAND_STORE), TO_STRING) + .add("locally_durable_to_data_store", entry, e -> e.maxBound(LOCALLY_DURABLE_TO_DATA_STORE), TO_STRING) + .add("locally_redundant", entry, e -> e.maxBound(LOCALLY_REDUNDANT), TO_STRING) + .add("locally_synced", entry, e -> e.maxBound(LOCALLY_SYNCED), TO_STRING) + .add("locally_witnessed", entry, e -> e.maxBound(LOCALLY_WITNESSED), TO_STRING) + .add("log_unavailable", entry, e -> e.maxBound(LOG_UNAVAILABLE), TO_STRING) + .add("unready", entry, e -> e.maxBound(UNREADY), TO_STRING) + .add("stale_until_at_least", entry.staleUntilAtLeast, TO_STRING); + }); + return rs; + }, rows, ignore -> false + ); + }); } - return dataSet; } } - public static final class RejectBeforeTable extends AbstractVirtualTable + public static final class RejectBeforeTable extends AbstractLazyVirtualTable { private RejectBeforeTable() { super(parse(VIRTUAL_ACCORD_DEBUG, REJECT_BEFORE, "Accord per-CommandStore RejectBefore State", "CREATE TABLE %s (\n" + - " keyspace_name text,\n" + - " table_name text,\n" + - " table_id text,\n" + " command_store_id int,\n" + " token_start 'TokenUtf8Type',\n" + + " table_id text,\n" + " token_end 'TokenUtf8Type',\n" + " timestamp text,\n" + - " PRIMARY KEY (keyspace_name, table_name, table_id, command_store_id, token_start)" + - ')', UTF8Type.instance)); + " PRIMARY KEY (command_store_id, token_start)" + + ')', UTF8Type.instance), FAIL, ASC); } @Override - public DataSet data() + protected void collect(PartitionsCollector collector) { CommandStores commandStores = AccordService.instance().node().commandStores(); - SimpleDataSet dataSet = new SimpleDataSet(metadata()); for (CommandStore commandStore : commandStores.all()) { RejectBefore rejectBefore = commandStore.unsafeGetRejectBefore(); if (rejectBefore == null) continue; - TableId tableId = ((AccordCommandStore)commandStore).tableId(); - String tableIdStr = tableId.toString(); - TableMetadata tableMetadata = tableMetadata(tableId); - String keyspace = keyspace(tableMetadata); - String table = table(tableId, tableMetadata); - rejectBefore.foldlWithBounds( - (timestamp, ds, start, end) -> ds.row(keyspace, table, tableIdStr, commandStore.id(), printToken(start)) - .column("token_end", printToken(end)) - .column("timestamp", timestamp.toString()) - , - dataSet, - ignore -> false - ); + collector.partition(commandStore.id()).collect(rows -> { + TableId tableId = ((AccordCommandStore)commandStore).tableId(); + String tableIdStr = tableId.toString(); + rejectBefore.foldlWithBounds((timestamp, rs, start, end) -> { + rs.add(printToken(start)) + .lazyCollect(columns -> columns.add("table_id", tableIdStr) + .add("token_end", end, AccordDebugKeyspace::printToken) + .add("timestamp", timestamp, AccordDebugKeyspace::toStringOrNull)); + return rs; + }, rows, ignore -> false); + }); } - return dataSet; } } @@ -995,11 +920,11 @@ public class AccordDebugKeyspace extends VirtualKeyspace super(parse(VIRTUAL_ACCORD_DEBUG, TXN_TRACE, "Accord Transaction Trace Configuration", "CREATE TABLE %s (\n" + - " txn_id text,\n" + + " txn_id 'TxnIdUtf8Type',\n" + " event_type text,\n" + " permits int,\n" + " PRIMARY KEY (txn_id, event_type)" + - ')', UTF8Type.instance)); + ')', TxnIdUtf8Type.instance)); } @Override @@ -1056,21 +981,21 @@ public class AccordDebugKeyspace extends VirtualKeyspace } } - public static final class TxnTracesTable extends AbstractMutableVirtualTable + public static final class TxnTracesTable extends AbstractMutableLazyVirtualTable { private TxnTracesTable() { super(parse(VIRTUAL_ACCORD_DEBUG, TXN_TRACES, "Accord Transaction Traces", "CREATE TABLE %s (\n" + - " txn_id text,\n" + + " txn_id 'TxnIdUtf8Type',\n" + " event_type text,\n" + " id_micros bigint,\n" + " at_micros bigint,\n" + " command_store_id int,\n" + " message text,\n" + " PRIMARY KEY (txn_id, event_type, id_micros, at_micros)" + - ')', UTF8Type.instance)); + ')', TxnIdUtf8Type.instance), FAIL, UNSORTED, UNSORTED); } private AccordTracing tracing() @@ -1079,29 +1004,29 @@ public class AccordDebugKeyspace extends VirtualKeyspace } @Override - protected void applyPartitionDeletion(ColumnValues partitionKey) + protected void applyPartitionDeletion(Object[] partitionKeys) { - TxnId txnId = TxnId.parse(partitionKey.value(0)); + TxnId txnId = TxnId.parse((String)partitionKeys[0]); tracing().eraseEvents(txnId); } @Override - protected void applyRangeTombstone(ColumnValues partitionKey, Range range) + protected void applyRangeTombstone(Object[] partitionKeys, Object[] starts, boolean startInclusive, Object[] ends, boolean endInclusive) { - TxnId txnId = TxnId.parse(partitionKey.value(0)); - if (!range.hasLowerBound() || range.lowerBoundType() != BoundType.CLOSED) throw invalidRequest("May restrict deletion by at most one event_type"); - if (range.lowerEndpoint().size() != 1) throw invalidRequest("Deletion restricted by lower bound on id_micros or at_micros is unsupported"); - if (!range.hasUpperBound() || (range.upperBoundType() != BoundType.CLOSED && range.upperEndpoint().size() == 1)) throw invalidRequest("Range deletion must specify one event_type"); - if (!range.upperEndpoint().value(0).equals(range.lowerEndpoint().value(0))) throw invalidRequest("May restrict deletion by at most one event_type"); - if (range.upperEndpoint().size() > 2) throw invalidRequest("Deletion restricted by upper bound on at_micros is unsupported"); - TraceEventType eventType = parseEventType(range.lowerEndpoint().value(0)); - if (range.upperEndpoint().size() == 1) + TxnId txnId = TxnId.parse((String) partitionKeys[0]); + if (!startInclusive) throw invalidRequest("May restrict deletion by at most one event_type"); + if (starts.length != 1) throw invalidRequest("Deletion restricted by lower bound on id_micros or at_micros is unsupported"); + if (ends.length == 0 || (ends.length == 1 && !endInclusive)) throw invalidRequest("Range deletion must specify one event_type"); + if (!ends[0].equals(starts[0])) throw invalidRequest("May restrict deletion by at most one event_type"); + if (ends.length > 2) throw invalidRequest("Deletion restricted by upper bound on at_micros is unsupported"); + TraceEventType eventType = parseEventType((String) starts[0]); + if (ends.length == 1) { tracing().eraseEvents(txnId, eventType); } else { - long before = range.upperEndpoint().value(1); + long before = (Long)ends[1]; tracing().eraseEventsBefore(txnId, eventType, before); } } @@ -1113,36 +1038,118 @@ public class AccordDebugKeyspace extends VirtualKeyspace } @Override - public DataSet data() + public void collect(PartitionsCollector collector) { - SimpleDataSet dataSet = new SimpleDataSet(metadata()); tracing().forEach(id -> true, (txnId, eventType, permits, events) -> { events.forEach(e -> { - e.messages().forEach(m -> { - dataSet.row(txnId.toString(), eventType.name(), e.idMicros, NANOSECONDS.toMicros(m.atNanos - e.atNanos)) - .column("command_store_id", m.commandStoreId) - .column("message", m.message); - }); + if (e.messages().isEmpty()) + { + collector.row(txnId.toString(), eventType.name(), e.idMicros, 0L) + .eagerCollect(columns -> { + columns.add("message", ""); + }); + } + else + { + e.messages().forEach(m -> { + collector.row(txnId.toString(), eventType.name(), e.idMicros, NANOSECONDS.toMicros(m.atNanos - e.atNanos)) + .eagerCollect(columns -> { + columns.add("command_store_id", m.commandStoreId) + .add("message", m.message); + }); + }); + } }); }); - return dataSet; } } // TODO (desired): don't report null as "null" - public static final class TxnTable extends AbstractVirtualTable implements AbstractVirtualTable.DataSet + abstract static class AbstractJournalTable extends AbstractLazyVirtualTable { - static class Entry - { - final int commandStoreId; - final Command command; + static final CompositeType PK = CompositeType.getInstance(Int32Type.instance, UTF8Type.instance); - Entry(int commandStoreId, Command command) - { - this.commandStoreId = commandStoreId; - this.command = command; - } + AbstractJournalTable(TableMetadata metadata) + { + super(metadata, FAIL, ASC); } + + @Override + public boolean allowFilteringImplicitly() + { + return false; + } + + @Override + public boolean allowFilteringPrimaryKeysImplicitly() + { + return true; + } + + @Override + public void collect(PartitionsCollector collector) + { + AccordService accord; + { + IAccordService iaccord = AccordService.instance(); + if (!iaccord.isEnabled()) + return; + + accord = (AccordService) iaccord; + } + + DataRange dataRange = collector.dataRange(); + JournalKey min = toJournalKey(dataRange.startKey()), + max = toJournalKey(dataRange.stopKey()); + + if (min == null && max == null) + { + FilterRange filterTxnId = collector.filters("txn_id", TxnId::parse, UnaryOperator.identity(), UnaryOperator.identity()); + FilterRange filterCommandStoreId = collector.filters("command_store_id", UnaryOperator.identity(), i -> i + 1, i -> i - 1); + + int minCommandStoreId = filterCommandStoreId.min == null ? -1 : filterCommandStoreId.min; + int maxCommandStoreId = filterCommandStoreId.max == null ? Integer.MAX_VALUE : filterCommandStoreId.max; + + if (filterTxnId.min != null && filterTxnId.max != null && filterTxnId.min.equals(filterTxnId.max)) + { + TxnId txnId = filterTxnId.min; + accord.node().commandStores().forAllUnsafe(commandStore -> { + if (commandStore.id() < minCommandStoreId || commandStore.id() > maxCommandStoreId) + return; + + collect(collector, accord, new JournalKey(txnId, JournalKey.Type.COMMAND_DIFF, commandStore.id())); + }); + return; + } + + if (filterTxnId.min != null || filterTxnId.max != null || minCommandStoreId >= 0 || maxCommandStoreId < Integer.MAX_VALUE) + { + min = new JournalKey(filterTxnId.min == null ? TxnId.NONE : filterTxnId.min, JournalKey.Type.COMMAND_DIFF, Math.max(0, minCommandStoreId)); + max = new JournalKey(filterTxnId.max == null ? TxnId.MAX.withoutNonIdentityFlags() : filterTxnId.max, JournalKey.Type.COMMAND_DIFF, maxCommandStoreId); + } + } + + accord.journal().forEach(key -> collect(collector, accord, key), min, max, true); + } + + abstract void collect(PartitionsCollector collector, AccordService accord, JournalKey key); + + private static JournalKey toJournalKey(PartitionPosition position) + { + if (position.isMinimum()) + return null; + + if (!(position instanceof DecoratedKey)) + throw new InvalidRequestException("Cannot filter this table by partial partition key"); + + ByteBuffer[] keys = PK.split(((DecoratedKey) position).getKey()); + return new JournalKey(TxnId.parse(UTF8Type.instance.compose(keys[1])), JournalKey.Type.COMMAND_DIFF, Int32Type.instance.compose(keys[0])); + } + } + + // TODO (desired): don't report null as "null" + public static final class TxnTable extends AbstractJournalTable + { private TxnTable() { super(parse(VIRTUAL_ACCORD_DEBUG, TXN, @@ -1165,88 +1172,51 @@ public class AccordDebugKeyspace extends VirtualKeyspace " participants_has_touched text,\n" + " participants_executes text,\n" + " participants_waits_on text,\n" + - " PRIMARY KEY (txn_id, command_store_id)" + - ')', UTF8Type.instance)); + " PRIMARY KEY ((command_store_id, txn_id))" + + ')', PK)); } @Override - public DataSet data() + void collect(PartitionsCollector collector, AccordService accord, JournalKey key) { - return this; + if (key.type != JournalKey.Type.COMMAND_DIFF) + return; + + AccordCommandStore commandStore = (AccordCommandStore) accord.node().commandStores().forId(key.commandStoreId); + if (commandStore == null) + return; + + Command command = commandStore.loadCommand(key.id); + if (command == null) + return; + + collector.row(key.commandStoreId, key.id.toString()) + .lazyCollect(columns -> addColumns(command, columns)); } - @Override - public boolean isEmpty() + private static void addColumns(Command command, ColumnsCollector columns) { - return false; - } - - @Override - public Partition getPartition(DecoratedKey partitionKey) - { - String txnIdStr = UTF8Type.instance.compose(partitionKey.getKey()); - TxnId txnId = TxnId.parse(txnIdStr); - - List commands = new CopyOnWriteArrayList<>(); - AccordService.instance().node().commandStores().forAllUnsafe(store -> { - Command command = ((AccordCommandStore)store).loadCommand(txnId); - if (command != null) - commands.add(new Entry(store.id(), command)); - }); - - if (commands.isEmpty()) - return null; - - SimpleDataSet ds = new SimpleDataSet(metadata); - for (Entry e : commands) - { - Command command = e.command; - ds.row(txnIdStr, e.commandStoreId) - .column("save_status", toStringOrNull(command.saveStatus())) - .column("route", toStringOrNull(command.route())) - .column("participants_owns", toStr(command, StoreParticipants::owns, StoreParticipants::stillOwns)) - .column("participants_touches", toStr(command, StoreParticipants::touches, StoreParticipants::stillTouches)) - .column("participants_has_touched", toStringOrNull(command.participants().hasTouched())) - .column("participants_executes", toStr(command, StoreParticipants::executes, StoreParticipants::stillExecutes)) - .column("participants_waits_on", toStr(command, StoreParticipants::waitsOn, StoreParticipants::stillWaitsOn)) - .column("durability", toStringOrNull(command.durability())) - .column("execute_at", toStringOrNull(command.executeAt())) - .column("executes_at_least", toStringOrNull(command.executesAtLeast())) - .column("txn", toStringOrNull(command.partialTxn())) - .column("deps", toStringOrNull(command.partialDeps())) - .column("waiting_on", toStringOrNull(command.waitingOn())) - .column("writes", toStringOrNull(command.writes())) - .column("result", toStringOrNull(command.result())); - } - - return ds.getPartition(partitionKey); - } - - @Override - public Iterator getPartitions(DataRange range) - { - throw new UnsupportedOperationException(); + StoreParticipants participants = command.participants(); + columns.add("save_status", command.saveStatus(), TO_STRING) + .add("route", participants, StoreParticipants::route, TO_STRING) + .add("participants_owns", participants, p -> toStr(p, StoreParticipants::owns, StoreParticipants::stillOwns)) + .add("participants_touches", participants, p -> toStr(p, StoreParticipants::touches, StoreParticipants::stillTouches)) + .add("participants_has_touched", participants, StoreParticipants::hasTouched, TO_STRING) + .add("participants_executes", participants, p -> toStr(p, StoreParticipants::executes, StoreParticipants::stillExecutes)) + .add("participants_waits_on", participants, p -> toStr(p, StoreParticipants::waitsOn, StoreParticipants::stillWaitsOn)) + .add("durability", command, Command::durability, TO_STRING) + .add("execute_at", command, Command::executeAt, TO_STRING) + .add("executes_at_least", command, Command::executesAtLeast, TO_STRING) + .add("txn", command, Command::partialTxn, TO_STRING) + .add("deps", command, Command::partialDeps, TO_STRING) + .add("waiting_on", command, Command::waitingOn, TO_STRING) + .add("writes", command, Command::writes, TO_STRING) + .add("result", command, Command::result, TO_STRING); } } - public static final class JournalTable extends AbstractVirtualTable implements AbstractVirtualTable.DataSet + public static final class JournalTable extends AbstractJournalTable { - static class Entry - { - final int commandStoreId; - final long segment; - final int position; - final CommandChange.Builder builder; - - Entry(int commandStoreId, long segment, int position, CommandChange.Builder builder) - { - this.commandStoreId = commandStoreId; - this.segment = segment; - this.position = position; - this.builder = builder; - } - } - private JournalTable() { super(parse(VIRTUAL_ACCORD_DEBUG, JOURNAL, @@ -1270,68 +1240,40 @@ public class AccordDebugKeyspace extends VirtualKeyspace " participants_has_touched text,\n" + " participants_executes text,\n" + " participants_waits_on text,\n" + - " PRIMARY KEY (txn_id, command_store_id, segment, segment_position)" + - ')', UTF8Type.instance)); + " PRIMARY KEY ((command_store_id, txn_id), segment, segment_position)" + + ')', PK)); } @Override - public DataSet data() + void collect(PartitionsCollector collector, AccordService accord, JournalKey key) { - return this; - } - - @Override - public boolean isEmpty() - { - return false; - } - - @Override - public Partition getPartition(DecoratedKey partitionKey) - { - String txnIdStr = UTF8Type.instance.compose(partitionKey.getKey()); - TxnId txnId = TxnId.parse(txnIdStr); - - List entries = new ArrayList<>(); - AccordService.instance().node().commandStores().forAllUnsafe(store -> { - for (AccordJournal.DebugEntry e : ((AccordCommandStore)store).debugCommand(txnId)) - entries.add(new Entry(store.id(), e.segment, e.position, e.builder)); + AccordCommandStore commandStore = (AccordCommandStore) accord.node().commandStores().forId(key.commandStoreId); + collector.partition(key.commandStoreId, key.id.toString()).collect(rows -> { + for (AccordJournal.DebugEntry e : commandStore.debugCommand(key.id)) + { + CommandChange.Builder b = e.builder; + StoreParticipants participants = b.participants() != null ? b.participants() : StoreParticipants.empty(key.id); + rows.add(e.segment, e.position) + .lazyCollect(columns -> { + columns.add("save_status", b.saveStatus(), TO_STRING) + .add("route", participants, StoreParticipants::route, TO_STRING) + .add("participants_owns", participants, p -> toStr(p, StoreParticipants::owns, StoreParticipants::stillOwns)) + .add("participants_touches", participants, p -> toStr(p, StoreParticipants::touches, StoreParticipants::stillTouches)) + .add("participants_has_touched", participants, StoreParticipants::hasTouched, TO_STRING) + .add("participants_executes", participants, p -> toStr(p, StoreParticipants::executes, StoreParticipants::stillExecutes)) + .add("participants_waits_on", participants, p -> toStr(p, StoreParticipants::waitsOn, StoreParticipants::stillWaitsOn)) + .add("durability", b.durability(), TO_STRING) + .add("execute_at", b.executeAt(), TO_STRING) + .add("executes_at_least", b.executesAtLeast(), TO_STRING) + .add("txn", b.partialTxn(), TO_STRING) + .add("deps", b.partialDeps(), TO_STRING) + .add("writes", b.writes(), TO_STRING) + .add("result", b.result(), TO_STRING); + }); + } }); - - if (entries.isEmpty()) - return null; - - SimpleDataSet ds = new SimpleDataSet(metadata); - for (Entry e : entries) - { - CommandChange.Builder b = e.builder; - StoreParticipants participants = b.participants(); - if (participants == null) participants = StoreParticipants.empty(txnId); - ds.row(txnIdStr, e.commandStoreId, e.segment, e.position) - .column("save_status", toStringOrNull(b.saveStatus())) - .column("route", toStringOrNull(participants.route())) - .column("participants_owns", toStr(participants, StoreParticipants::owns, StoreParticipants::stillOwns)) - .column("participants_touches", toStr(participants, StoreParticipants::touches, StoreParticipants::stillTouches)) - .column("participants_has_touched", toStringOrNull(participants.hasTouched())) - .column("participants_executes", toStr(participants, StoreParticipants::executes, StoreParticipants::stillExecutes)) - .column("participants_waits_on", toStr(participants, StoreParticipants::waitsOn, StoreParticipants::stillWaitsOn)) - .column("durability", toStringOrNull(b.durability())) - .column("execute_at", toStringOrNull(b.executeAt())) - .column("executes_at_least", toStringOrNull(b.executesAtLeast())) - .column("txn", toStringOrNull(b.partialTxn())) - .column("deps", toStringOrNull(b.partialDeps())) - .column("writes", toStringOrNull(b.writes())) - .column("result", toStringOrNull(b.result())); - } - - return ds.getPartition(partitionKey); } - @Override - public Iterator getPartitions(DataRange range) - { - throw new UnsupportedOperationException(); - } } /** @@ -1346,7 +1288,7 @@ public class AccordDebugKeyspace extends VirtualKeyspace */ // Had to be separate from the "regular" journal table since it does not have segment and position, and command store id is inferred // TODO (required): add access control - public static final class TxnOpsTable extends AbstractMutableVirtualTable implements AbstractVirtualTable.DataSet + public static final class TxnOpsTable extends AbstractMutableLazyVirtualTable { // TODO (expected): test each of these operations enum Op { ERASE_VESTIGIAL, INVALIDATE, TRY_EXECUTE, FORCE_APPLY, FORCE_UPDATE, RECOVER, FETCH, RESET_PROGRESS_LOG } @@ -1359,41 +1301,21 @@ public class AccordDebugKeyspace extends VirtualKeyspace " command_store_id int,\n" + " op text," + " PRIMARY KEY (txn_id, command_store_id)" + - ')', UTF8Type.instance)); + ')', UTF8Type.instance), FAIL, UNSORTED); } @Override - public DataSet data() + protected void collect(PartitionsCollector collector) { throw new UnsupportedOperationException(TXN_OPS + " is a write-only table"); } @Override - public boolean isEmpty() + protected void applyRowUpdate(Object[] partitionKeys, Object[] clusteringKeys, ColumnMetadata[] columns, Object[] values) { - return true; - } - - @Override - public Partition getPartition(DecoratedKey partitionKey) - { - throw new UnsupportedOperationException(TXN_OPS + " is a write-only table"); - } - - @Override - public Iterator getPartitions(DataRange range) - { - throw new UnsupportedOperationException(TXN_OPS + " is a write-only table"); - } - - - @Override - protected void applyColumnUpdate(ColumnValues partitionKey, ColumnValues clusteringColumns, Optional columnValue) - { - TxnId txnId = TxnId.parse(partitionKey.value(0)); - int commandStoreId = clusteringColumns.value(0); - Invariants.require(columnValue.isPresent()); - Op op = Op.valueOf(columnValue.get().value()); + TxnId txnId = TxnId.parse((String) partitionKeys[0]); + int commandStoreId = (Integer) clusteringKeys[0]; + Op op = Op.valueOf((String)values[0]); switch (op) { default: throw new UnhandledEnum(op); @@ -1521,117 +1443,57 @@ public class AccordDebugKeyspace extends VirtualKeyspace } } - public static class TxnBlockedByTable extends AbstractVirtualTable + public static class TxnBlockedByTable extends AbstractLazyVirtualTable { - enum Reason { Self, Txn, Key } + enum Reason + {Self, Txn, Key} protected TxnBlockedByTable() { super(parse(VIRTUAL_ACCORD_DEBUG, TXN_BLOCKED_BY, - "Accord Transactions Blocked By Table" , + "Accord Transactions Blocked By Table", "CREATE TABLE %s (\n" + - " txn_id text,\n" + - " keyspace_name text,\n" + - " table_name text,\n" + + " txn_id 'TxnIdUtf8Type',\n" + " command_store_id int,\n" + " depth int,\n" + - " blocked_by text,\n" + - " reason text,\n" + + " blocked_by_key text,\n" + + " blocked_by_txn_id 'TxnIdUtf8Type',\n" + " save_status text,\n" + " execute_at text,\n" + - " key text,\n" + - " PRIMARY KEY (txn_id, keyspace_name, table_name, command_store_id, depth, blocked_by, reason)" + - ')', UTF8Type.instance)); + " PRIMARY KEY (txn_id, command_store_id, depth, blocked_by_key, blocked_by_txn_id)" + + ')', TxnIdUtf8Type.instance), BEST_EFFORT, ASC); } @Override - public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter) + protected void collect(PartitionsCollector collector) { - Partition partition = data(partitionKey, rowFilter).getPartition(partitionKey); + Object[] pks = collector.singlePartitionKey(); + if (pks == null) + throw new InvalidRequestException(metadata + " only supports single partition key queries"); - if (null == partition) - return EmptyIterators.unfilteredPartition(metadata); + FilterRange depthRange = collector.filters("depth", Function.identity(), i -> i + 1, i -> i - 1); + int maxDepth = depthRange.max == null ? Integer.MAX_VALUE : depthRange.max; - long now = currentTimeMillis(); - UnfilteredRowIterator rowIterator = partition.toRowIterator(metadata(), clusteringIndexFilter, columnFilter, now); - return new SingletonUnfilteredPartitionIterator(rowIterator); - } - - public DataSet data(DecoratedKey partitionKey, RowFilter rowFilter) - { - int maxDepth = Integer.MAX_VALUE; - if (rowFilter != null && rowFilter.getExpressions().size() > 0) - { - Invariants.require(rowFilter.getExpressions().size() == 1, "Only depth filter is supported"); - RowFilter.Expression expression = rowFilter.getExpressions().get(0); - Invariants.require(expression.column().name.toString().equals("depth"), "Only depth filter is supported, but got: %s", expression.column().name); - Invariants.require(expression.operator() == Operator.LT || expression.operator() == Operator.LTE, "Only < and <= queries are supported"); - if (expression.operator() == Operator.LT) - maxDepth = expression.getIndexValue().getInt(0); - else - maxDepth = expression.getIndexValue().getInt(0) + 1; - } - - TxnId id = TxnId.parse(UTF8Type.instance.compose(partitionKey.getKey())); - List shards = AccordService.instance().debugTxnBlockedGraph(id); - - SimpleDataSet ds = new SimpleDataSet(metadata()); - CommandStores commandStores = AccordService.instance().node().commandStores(); - for (CommandStoreTxnBlockedGraph shard : shards) - { - Set processed = new HashSet<>(); - process(ds, commandStores, shard, processed, id, 0, maxDepth, id, Reason.Self, null); - // everything was processed right? - if (!shard.txns.isEmpty() && !shard.txns.keySet().containsAll(processed)) - Invariants.expect(false, "Skipped txns: " + Sets.difference(shard.txns.keySet(), processed)); - } - - return ds; - } - - private void process(SimpleDataSet ds, CommandStores commandStores, CommandStoreTxnBlockedGraph shard, Set processed, TxnId userTxn, int depth, int maxDepth, TxnId txnId, Reason reason, Runnable onDone) - { - if (!processed.add(txnId)) - throw new IllegalStateException("Double processed " + txnId); - CommandStoreTxnBlockedGraph.TxnState txn = shard.txns.get(txnId); - if (txn == null) - { - Invariants.require(reason == Reason.Self, "Txn %s unknown for reason %s", txnId, reason); - return; - } - // was it applied? If so ignore it - if (reason != Reason.Self && txn.saveStatus.hasBeen(Status.Applied)) - return; - TableId tableId = tableId(shard.commandStoreId, commandStores); - TableMetadata tableMetadata = tableMetadata(tableId); - ds.row(userTxn.toString(), keyspace(tableMetadata), table(tableId, tableMetadata), - shard.commandStoreId, depth, reason == Reason.Self ? "" : txn.txnId.toString(), reason.name()); - ds.column("save_status", txn.saveStatus.name()); - if (txn.executeAt != null) - ds.column("execute_at", txn.executeAt.toString()); - if (onDone != null) - onDone.run(); - if (txn.isBlocked()) - { - for (TxnId blockedBy : txn.blockedBy) + TxnId txnId = TxnId.parse((String) pks[0]); + PartitionCollector partition = collector.partition(pks[0]); + partition.collect(rows -> { + try { - if (!processed.contains(blockedBy) && depth < maxDepth) - process(ds, commandStores, shard, processed, userTxn, depth + 1, maxDepth, blockedBy, Reason.Txn, null); + DebugBlockedTxns.visit(AccordService.instance(), txnId, maxDepth, collector.deadlineNanos(), txn -> { + String keyStr = txn.blockedViaKey == null ? "" : txn.blockedViaKey.toString(); + String txnIdStr = txn.txnId == null || txn.txnId.equals(txnId) ? "" : txn.txnId.toString(); + rows.add(txn.commandStoreId, txn.depth, keyStr, txnIdStr) + .eagerCollect(columns -> { + columns.add("save_status", txn.saveStatus, TO_STRING) + .add("execute_at", txn.executeAt, TO_STRING); + }); + }); } - - for (TokenKey blockedBy : txn.blockedByKey) + catch (TimeoutException e) { - TxnId blocking = shard.keys.get(blockedBy); - if (!processed.contains(blocking) && depth < maxDepth) - process(ds, commandStores, shard, processed, userTxn, depth + 1, maxDepth, blocking, Reason.Key, () -> ds.column("key", printToken(blockedBy))); + throw new InternalTimeoutException(); } - } - } - - @Override - public DataSet data() - { - throw new InvalidRequestException("Must select a single txn_id"); + }); } } @@ -1650,33 +1512,12 @@ public class AccordDebugKeyspace extends VirtualKeyspace return Schema.instance.getTableMetadata(tableId); } - private static String keyspace(TableMetadata metadata) - { - return metadata == null ? "Unknown" : metadata.keyspace; - } - - private static String table(TableId tableId, TableMetadata metadata) - { - return metadata == null ? tableId.toString() : metadata.name; - } - private static String printToken(RoutingKey routingKey) { TokenKey key = (TokenKey) routingKey; return key.token().getPartitioner().getTokenFactory().toString(key.token()); } - private static ByteBuffer sortToken(RoutingKey routingKey) - { - TokenKey key = (TokenKey) routingKey; - Token token = key.token(); - IPartitioner partitioner = token.getPartitioner(); - ByteBuffer out = ByteBuffer.allocate(partitioner.accordSerializedSize(token)); - partitioner.accordSerialize(token, out); - out.flip(); - return out; - } - private static TableMetadata parse(String keyspace, String table, String comment, String schema, AbstractType partitionKeyType) { return CreateTableStatement.parse(format(schema, table), keyspace) @@ -1686,13 +1527,11 @@ public class AccordDebugKeyspace extends VirtualKeyspace .build(); } - private static String toStr(Command command, Function> a, Function> b) - { - return toStr(command.participants(), a, b); - } - private static String toStr(StoreParticipants participants, Function> a, Function> b) { + if (participants == null) + return null; + Participants av = a.apply(participants); Participants bv = b.apply(participants); if (av == bv || av.equals(bv)) @@ -1710,6 +1549,14 @@ public class AccordDebugKeyspace extends VirtualKeyspace { if (o == null) return null; - return Objects.toString(o); + + try + { + return Objects.toString(o); + } + catch (Throwable t) + { + return "'; + } } } diff --git a/src/java/org/apache/cassandra/db/virtual/CollectionVirtualTableAdapter.java b/src/java/org/apache/cassandra/db/virtual/CollectionVirtualTableAdapter.java index cc311e1a26..5bff0fe3be 100644 --- a/src/java/org/apache/cassandra/db/virtual/CollectionVirtualTableAdapter.java +++ b/src/java/org/apache/cassandra/db/virtual/CollectionVirtualTableAdapter.java @@ -50,6 +50,7 @@ import org.apache.cassandra.db.DeletionTime; import org.apache.cassandra.db.EmptyIterators; import org.apache.cassandra.db.filter.ClusteringIndexFilter; import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.DataLimits; import org.apache.cassandra.db.filter.RowFilter; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.db.marshal.BooleanType; @@ -308,7 +309,7 @@ public class CollectionVirtualTableAdapter implements VirtualTable public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringFilter, ColumnFilter columnFilter, - RowFilter rowFilter) + RowFilter rowFilter, DataLimits limits) { if (!data.iterator().hasNext()) return EmptyIterators.unfilteredPartition(metadata); @@ -349,7 +350,7 @@ public class CollectionVirtualTableAdapter implements VirtualTable } @Override - public UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter) + public UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) { return createPartitionIterator(metadata, new AbstractIterator<>() { diff --git a/src/java/org/apache/cassandra/db/virtual/PartitionKeyStatsTable.java b/src/java/org/apache/cassandra/db/virtual/PartitionKeyStatsTable.java index 550743c6a7..d63ae194b4 100644 --- a/src/java/org/apache/cassandra/db/virtual/PartitionKeyStatsTable.java +++ b/src/java/org/apache/cassandra/db/virtual/PartitionKeyStatsTable.java @@ -42,6 +42,7 @@ import org.apache.cassandra.db.Slices; import org.apache.cassandra.db.context.CounterContext; import org.apache.cassandra.db.filter.ClusteringIndexFilter; import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.DataLimits; import org.apache.cassandra.db.filter.RowFilter; import org.apache.cassandra.db.marshal.CompositeType; import org.apache.cassandra.db.marshal.CounterColumnType; @@ -146,7 +147,7 @@ public class PartitionKeyStatsTable implements VirtualTable } @Override - public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter) + public UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) { if (clusteringIndexFilter.isReversed()) throw new InvalidRequestException(REVERSED_QUERY_ERROR); @@ -345,7 +346,7 @@ public class PartitionKeyStatsTable implements VirtualTable } @Override - public UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter) + public UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits) { throw new InvalidRequestException(UNSUPPORTED_RANGE_QUERY_ERROR); } diff --git a/src/java/org/apache/cassandra/db/virtual/VirtualTable.java b/src/java/org/apache/cassandra/db/virtual/VirtualTable.java index 770cb13983..3f05ed3a69 100644 --- a/src/java/org/apache/cassandra/db/virtual/VirtualTable.java +++ b/src/java/org/apache/cassandra/db/virtual/VirtualTable.java @@ -21,6 +21,7 @@ import org.apache.cassandra.db.DataRange; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.filter.ClusteringIndexFilter; import org.apache.cassandra.db.filter.ColumnFilter; +import org.apache.cassandra.db.filter.DataLimits; import org.apache.cassandra.db.filter.RowFilter; import org.apache.cassandra.db.partitions.PartitionUpdate; import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; @@ -31,6 +32,8 @@ import org.apache.cassandra.schema.TableMetadata; */ public interface VirtualTable { + enum Sorted { UNSORTED, ASC, DESC, SORTED } + /** * Returns the view name. * @@ -57,23 +60,25 @@ public interface VirtualTable /** * Selects the rows from a single partition. * - * @param partitionKey the partition key + * @param partitionKey the partition key * @param clusteringIndexFilter the clustering columns to selected - * @param columnFilter the selected columns - * @param rowFilter filter on which rows a given query should include or exclude + * @param columnFilter the selected columns + * @param rowFilter filter on which rows a given query should include or exclude + * @param limits result limits to apply * @return the rows corresponding to the requested data. */ - UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter); + UnfilteredPartitionIterator select(DecoratedKey partitionKey, ClusteringIndexFilter clusteringIndexFilter, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits); /** * Selects the rows from a range of partitions. * - * @param dataRange the range of data to retrieve + * @param dataRange the range of data to retrieve * @param columnFilter the selected columns - * @param rowFilter filter on which rows a given query should include or exclude + * @param rowFilter filter on which rows a given query should include or exclude + * @param limits * @return the rows corresponding to the requested data. */ - UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter); + UnfilteredPartitionIterator select(DataRange dataRange, ColumnFilter columnFilter, RowFilter rowFilter, DataLimits limits); /** * Truncates data from the underlying source, if supported. @@ -90,4 +95,9 @@ public interface VirtualTable { return true; } + + default boolean allowFilteringPrimaryKeysImplicitly() + { + return allowFilteringImplicitly(); + } } diff --git a/src/java/org/apache/cassandra/journal/InMemoryIndex.java b/src/java/org/apache/cassandra/journal/InMemoryIndex.java index 49fe4d1367..974767590a 100644 --- a/src/java/org/apache/cassandra/journal/InMemoryIndex.java +++ b/src/java/org/apache/cassandra/journal/InMemoryIndex.java @@ -18,6 +18,7 @@ package org.apache.cassandra.journal; import java.io.IOException; +import java.util.Iterator; import java.util.NavigableMap; import java.util.TreeMap; import java.util.concurrent.ConcurrentSkipListMap; @@ -126,6 +127,16 @@ public final class InMemoryIndex extends Index return lookUp(id); } + public Iterator keyIterator(@Nullable K min, @Nullable K max) + { + NavigableMap m; + if (min == null && max == null) m = index; + else if (min == null) m = index.headMap(max, true); + else if (max == null) m = index.tailMap(min, true); + else m = index.subMap(min, true, max, true); + return m.keySet().iterator(); + } + public void persist(Descriptor descriptor) { File tmpFile = descriptor.tmpFileFor(Component.INDEX); diff --git a/src/java/org/apache/cassandra/journal/Journal.java b/src/java/org/apache/cassandra/journal/Journal.java index d0793d39a7..2a63ab7cc4 100644 --- a/src/java/org/apache/cassandra/journal/Journal.java +++ b/src/java/org/apache/cassandra/journal/Journal.java @@ -934,11 +934,11 @@ public class Journal implements Shutdownable } /** - * Static segment iterator iterates all keys in _static_ segments in order. + * segment iterator iterates all keys in order. */ - public StaticSegmentKeyIterator staticSegmentKeyIterator(K min, K max) + public SegmentKeyIterator segmentKeyIterator(K min, K max, Predicate> include) { - return new StaticSegmentKeyIterator(min, max); + return new SegmentKeyIterator(min, max, include); } /** @@ -1000,53 +1000,36 @@ public class Journal implements Shutdownable } } - public class StaticSegmentKeyIterator implements CloseableIterator> + public class SegmentKeyIterator implements CloseableIterator> { private final ReferencedSegments segments; private final MergeIterator> iterator; - public StaticSegmentKeyIterator(K min, K max) + public SegmentKeyIterator(K min, K max, Predicate> include) { - this.segments = selectAndReference(s -> s.isStatic() - && s.asStatic().index().entryCount() > 0 + this.segments = selectAndReference(s -> include.test(s) && !s.isEmpty() && (min == null || keySupport.compare(s.index().lastId(), min) >= 0) && (max == null || keySupport.compare(s.index().firstId(), max) <= 0)); List> iterators = new ArrayList<>(segments.count()); for (Segment segment : segments.allSorted(true)) { - final StaticSegment staticSegment = (StaticSegment) segment; - final OnDiskIndex.IndexReader iter = staticSegment.index().reader(); - if (min != null) iter.seek(min); - if (max != null) iter.seekEnd(max); - if (!iter.hasNext()) - continue; - - iterators.add(new AbstractIterator<>() + if (segment.isStatic()) { - final Head head = new Head(staticSegment.descriptor.timestamp); - - @Override - protected Head computeNext() - { - if (!iter.hasNext()) - return endOfData(); - - K next = iter.next(); - while (next.equals(head.key)) - { - if (!iter.hasNext()) - return endOfData(); - - next = iter.next(); - } - - Invariants.require(!next.equals(head.key), - "%s == %s", next, head.key); - head.key = next; - return head; - } - }); + final StaticSegment staticSegment = (StaticSegment) segment; + final OnDiskIndex.IndexReader iter = staticSegment.index().reader(); + if (min != null) iter.seek(min); + if (max != null) iter.seekEnd(max); + if (iter.hasNext()) + iterators.add(keyIterator(segment.descriptor.timestamp, iter)); + } + else + { + final ActiveSegment activeSegment = (ActiveSegment) segment; + final Iterator iter = activeSegment.index().keyIterator(min, max); + if (iter.hasNext()) + iterators.add(keyIterator(segment.descriptor.timestamp, iter)); + } } this.iterator = MergeIterator.get(iterators, @@ -1077,6 +1060,34 @@ public class Journal implements Shutdownable }); } + private Iterator keyIterator(long segment, Iterator iter) + { + final Head head = new Head(segment); + return new AbstractIterator<>() + { + @Override + protected Head computeNext() + { + if (!iter.hasNext()) + return endOfData(); + + K next = iter.next(); + while (next.equals(head.key)) + { + if (!iter.hasNext()) + return endOfData(); + + next = iter.next(); + } + + Invariants.require(!next.equals(head.key), + "%s == %s", next, head.key); + head.key = next; + return head; + } + }; + } + @Override public void close() { diff --git a/src/java/org/apache/cassandra/journal/Segment.java b/src/java/org/apache/cassandra/journal/Segment.java index 3854f0ee27..1fda2f57c9 100644 --- a/src/java/org/apache/cassandra/journal/Segment.java +++ b/src/java/org/apache/cassandra/journal/Segment.java @@ -66,7 +66,8 @@ public abstract class Segment implements SelfRefCounted>, Co abstract boolean isActive(); abstract boolean isFlushed(long position); - boolean isStatic() { return !isActive(); } + public boolean isStatic() { return !isActive(); } + abstract boolean isEmpty(); abstract ActiveSegment asActive(); abstract StaticSegment asStatic(); diff --git a/src/java/org/apache/cassandra/journal/StaticSegment.java b/src/java/org/apache/cassandra/journal/StaticSegment.java index 35c987c8a4..bc425dda7e 100644 --- a/src/java/org/apache/cassandra/journal/StaticSegment.java +++ b/src/java/org/apache/cassandra/journal/StaticSegment.java @@ -246,6 +246,12 @@ public final class StaticSegment extends Segment return index.entryCount(); } + @Override + boolean isEmpty() + { + return entryCount() == 0; + } + @Override boolean isActive() { diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java index 637b158177..2fa8675563 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java @@ -1655,7 +1655,8 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor { - public enum Status { WAITING_TO_LOAD, SCANNING_RANGES, LOADING, WAITING_TO_RUN, RUNNING } + // sorted in name order for reporting to virtual tables + public enum Status { LOADING, RUNNING, SCANNING_RANGES, WAITING_TO_LOAD, WAITING_TO_RUN } final Status status; final int commandStoreId; @@ -1706,7 +1707,7 @@ public abstract class AccordExecutor implements CacheSize, LoadExecutor iter = new CloseableIterator<>() { final CloseableIterator> iter = journalTable.keyIterator(topologyUpdateKey(0L), - topologyUpdateKey(Timestamp.MAX_EPOCH)); + topologyUpdateKey(Timestamp.MAX_EPOCH), + true); TopologyImage prev = null; @Override @@ -571,9 +572,14 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier journalTable.forceCompaction(); } - public void forEach(Consumer consumer) + public void forEach(Consumer consumer, boolean includeActive) { - try (CloseableIterator> iter = journalTable.keyIterator(null, null)) + forEach(consumer, null, null, includeActive); + } + + public void forEach(Consumer consumer, @Nullable JournalKey min, @Nullable JournalKey max, boolean includeActive) + { + try (CloseableIterator> iter = journalTable.keyIterator(min, max, includeActive)) { while (iter.hasNext()) { @@ -610,7 +616,7 @@ public class AccordJournal implements accord.api.Journal, RangeSearcher.Supplier this.commandStore = commandStore; this.replayer = commandStore.replayer(); // Keys in the index are sorted by command store id, so index iteration will be sequential - this.iter = journalTable.keyIterator(new JournalKey(TxnId.NONE, COMMAND_DIFF, commandStore.id()), new JournalKey(TxnId.MAX.withoutNonIdentityFlags(), COMMAND_DIFF, commandStore.id())); + this.iter = journalTable.keyIterator(new JournalKey(TxnId.NONE, COMMAND_DIFF, commandStore.id()), new JournalKey(TxnId.MAX.withoutNonIdentityFlags(), COMMAND_DIFF, commandStore.id()), false); } boolean replay() diff --git a/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java b/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java index 984ab7ed62..17aff49f86 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java +++ b/src/java/org/apache/cassandra/service/accord/AccordJournalTable.java @@ -75,6 +75,7 @@ import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.journal.Journal; import org.apache.cassandra.journal.KeySupport; import org.apache.cassandra.journal.RecordConsumer; +import org.apache.cassandra.journal.Segment; import org.apache.cassandra.schema.ColumnMetadata; import org.apache.cassandra.service.RetryStrategy; import org.apache.cassandra.service.accord.AccordKeyspace.JournalColumns; @@ -457,9 +458,9 @@ public class AccordJournalTable implements RangeSearche } @SuppressWarnings("resource") // Auto-closeable iterator will release related resources - public CloseableIterator> keyIterator(@Nullable K min, @Nullable K max) + public CloseableIterator> keyIterator(@Nullable K min, @Nullable K max, boolean includeActive) { - return new JournalAndTableKeyIterator(min, max); + return new JournalAndTableKeyIterator(min, max, includeActive); } private class TableIterator extends AbstractIterator implements CloseableIterator @@ -515,12 +516,12 @@ public class AccordJournalTable implements RangeSearche private class JournalAndTableKeyIterator extends AbstractIterator> implements CloseableIterator> { final TableIterator tableIterator; - final Journal.StaticSegmentKeyIterator journalIterator; + final Journal.SegmentKeyIterator journalIterator; - private JournalAndTableKeyIterator(K min, K max) + private JournalAndTableKeyIterator(K min, K max, boolean includeActive) { this.tableIterator = new TableIterator(min, max); - this.journalIterator = journal.staticSegmentKeyIterator(min, max); + this.journalIterator = journal.segmentKeyIterator(min, max, includeActive ? ignore -> true : Segment::isStatic); } K prevFromTable = null; diff --git a/src/java/org/apache/cassandra/service/accord/AccordService.java b/src/java/org/apache/cassandra/service/accord/AccordService.java index 34b7dd28c4..6742f24843 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordService.java +++ b/src/java/org/apache/cassandra/service/accord/AccordService.java @@ -21,7 +21,6 @@ package org.apache.cassandra.service.accord; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; -import java.util.Collections; import java.util.HashSet; import java.util.List; import java.util.Set; @@ -40,6 +39,7 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.primitives.Ints; import accord.api.ConfigurationService.EpochReady; +import accord.primitives.Txn; import org.apache.cassandra.metrics.AccordReplicaMetrics; import org.apache.cassandra.service.accord.api.AccordViolationHandler; import org.apache.cassandra.utils.Clock; @@ -58,17 +58,11 @@ import accord.impl.DefaultRemoteListeners; import accord.impl.RequestCallbacks; import accord.impl.SizeOfIntersectionSorter; import accord.impl.progresslog.DefaultProgressLogs; -import accord.local.Command; -import accord.local.CommandStore; -import accord.local.CommandStores; import accord.local.Node; import accord.local.Node.Id; -import accord.local.PreLoadContext; -import accord.local.SafeCommand; import accord.local.ShardDistributor.EvenSplit; import accord.local.UniqueTimeService.AtomicUniqueTimeWithStaleReservation; import accord.local.cfk.CommandsForKey; -import accord.local.cfk.SafeCommandsForKey; import accord.local.durability.DurabilityService; import accord.local.durability.ShardDurability; import accord.messages.Reply; @@ -76,13 +70,9 @@ import accord.messages.Request; import accord.primitives.FullRoute; import accord.primitives.Keys; import accord.primitives.Ranges; -import accord.primitives.RoutingKeys; -import accord.primitives.SaveStatus; import accord.primitives.Seekable; import accord.primitives.Seekables; -import accord.primitives.Status; import accord.primitives.Timestamp; -import accord.primitives.Txn; import accord.primitives.TxnId; import accord.topology.Shard; import accord.topology.Topology; @@ -90,7 +80,6 @@ import accord.topology.TopologyManager; import accord.utils.DefaultRandom; import accord.utils.Invariants; import accord.utils.async.AsyncChain; -import accord.utils.async.AsyncChains; import accord.utils.async.AsyncResult; import accord.utils.async.AsyncResults; import org.apache.cassandra.concurrent.Shutdownable; @@ -115,7 +104,6 @@ import org.apache.cassandra.service.accord.api.AccordScheduler; import org.apache.cassandra.service.accord.api.AccordTimeService; import org.apache.cassandra.service.accord.api.AccordTopologySorter; import org.apache.cassandra.service.accord.api.CompositeTopologySorter; -import org.apache.cassandra.service.accord.api.TokenKey; import org.apache.cassandra.service.accord.api.TokenKey.KeyspaceSplitter; import org.apache.cassandra.service.accord.interop.AccordInteropAdapter.AccordInteropFactory; import org.apache.cassandra.service.accord.serializers.TableMetadatas; @@ -140,8 +128,6 @@ import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException; import static accord.api.Journal.TopologyUpdate; import static accord.api.ProtocolModifiers.Toggles.FastExec.MAY_BYPASS_SAFESTORE; -import static accord.local.LoadKeys.SYNC; -import static accord.local.LoadKeysFor.READ_WRITE; import static accord.local.durability.DurabilityService.SyncLocal.Self; import static accord.local.durability.DurabilityService.SyncRemote.All; import static accord.messages.SimpleReply.Ok; @@ -875,139 +861,6 @@ public class AccordService implements IAccordService, Shutdownable return node.id(); } - @Override - public List debugTxnBlockedGraph(TxnId txnId) - { - return getBlocking(loadDebug(txnId)); - } - - public AsyncChain> loadDebug(TxnId original) - { - CommandStores commandStores = node.commandStores(); - if (commandStores.count() == 0) - return AsyncChains.success(Collections.emptyList()); - int[] ids = commandStores.ids(); - List> chains = new ArrayList<>(ids.length); - for (int id : ids) - chains.add(loadDebug(original, commandStores.forId(id)).chain()); - return AsyncChains.allOf(chains); - } - - private AsyncResult loadDebug(TxnId txnId, CommandStore store) - { - CommandStoreTxnBlockedGraph.Builder state = new CommandStoreTxnBlockedGraph.Builder(store.id()); - populateAsync(state, store, txnId); - return state; - } - - private static void populate(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TxnId blockedBy) - { - if (safeStore.ifLoadedAndInitialised(blockedBy) != null) populateSync(state, safeStore, blockedBy); - else populateAsync(state, safeStore.commandStore(), blockedBy); - } - - private static void populateAsync(CommandStoreTxnBlockedGraph.Builder state, CommandStore store, TxnId txnId) - { - state.asyncTxns.incrementAndGet(); - store.execute(PreLoadContext.contextFor(txnId, "Populate txn_blocked_by"), in -> { - populateSync(state, (AccordSafeCommandStore) in, txnId); - if (0 == state.asyncTxns.decrementAndGet() && 0 == state.asyncKeys.get()) - state.complete(); - }); - } - - @Nullable - private static void populateSync(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TxnId txnId) - { - try - { - if (state.txns.containsKey(txnId)) - return; // could plausibly request same txn twice - - SafeCommand safeCommand = safeStore.unsafeGet(txnId); - Invariants.nonNull(safeCommand, "Txn %s is not in the cache", txnId); - if (safeCommand.current() == null || safeCommand.current().saveStatus() == SaveStatus.Uninitialised) - return; - - CommandStoreTxnBlockedGraph.TxnState cmdTxnState = populateSync(state, safeCommand.current()); - if (cmdTxnState.notBlocked()) - return; - - for (TxnId blockedBy : cmdTxnState.blockedBy) - { - if (!state.knows(blockedBy)) - populate(state, safeStore, blockedBy); - } - for (TokenKey blockedBy : cmdTxnState.blockedByKey) - { - if (!state.keys.containsKey(blockedBy)) - populate(state, safeStore, blockedBy, txnId, safeCommand.current().executeAt()); - } - } - catch (Throwable t) - { - state.tryFailure(t); - } - } - - private static void populate(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TokenKey blockedBy, TxnId txnId, Timestamp executeAt) - { - if (safeStore.ifLoadedAndInitialised(txnId) != null && safeStore.ifLoadedAndInitialised(blockedBy) != null) populateSync(state, safeStore, blockedBy, txnId, executeAt); - else populateAsync(state, safeStore.commandStore(), blockedBy, txnId, executeAt); - } - - private static void populateAsync(CommandStoreTxnBlockedGraph.Builder state, CommandStore commandStore, TokenKey blockedBy, TxnId txnId, Timestamp executeAt) - { - state.asyncKeys.incrementAndGet(); - commandStore.execute(PreLoadContext.contextFor(txnId, RoutingKeys.of(blockedBy.toUnseekable()), SYNC, READ_WRITE, "Populate txn_blocked_by"), in -> { - populateSync(state, (AccordSafeCommandStore) in, blockedBy, txnId, executeAt); - if (0 == state.asyncKeys.decrementAndGet() && 0 == state.asyncTxns.get()) - state.complete(); - }); - } - - private static void populateSync(CommandStoreTxnBlockedGraph.Builder state, AccordSafeCommandStore safeStore, TokenKey pk, TxnId txnId, Timestamp executeAt) - { - try - { - SafeCommandsForKey commandsForKey = safeStore.ifLoadedAndInitialised(pk); - TxnId blocking = commandsForKey.current().blockedOnTxnId(txnId, executeAt); - if (blocking instanceof CommandsForKey.TxnInfo) - blocking = ((CommandsForKey.TxnInfo) blocking).plainTxnId(); - state.keys.put(pk, blocking); - if (state.txns.containsKey(blocking)) - return; - populate(state, safeStore, blocking); - } - catch (Throwable t) - { - state.tryFailure(t); - } - } - - private static CommandStoreTxnBlockedGraph.TxnState populateSync(CommandStoreTxnBlockedGraph.Builder state, Command cmd) - { - CommandStoreTxnBlockedGraph.Builder.TxnBuilder cmdTxnState = state.txn(cmd.txnId(), cmd.executeAt(), cmd.saveStatus()); - if (!cmd.hasBeen(Status.Applied) && cmd.hasBeen(Status.Stable)) - { - // check blocking state - Command.WaitingOn waitingOn = cmd.asCommitted().waitingOn(); - waitingOn.waitingOn.reverseForEach(null, null, null, null, (i1, i2, i3, i4, i) -> { - if (i < waitingOn.txnIdCount()) - { - // blocked on txn - cmdTxnState.blockedBy.add(waitingOn.txnId(i)); - } - else - { - // blocked on key - cmdTxnState.blockedByKey.add((TokenKey) waitingOn.keys.get(i - waitingOn.txnIdCount())); - } - }); - } - return cmdTxnState.build(); - } - @Override public long minEpoch() { diff --git a/src/java/org/apache/cassandra/service/accord/CommandStoreTxnBlockedGraph.java b/src/java/org/apache/cassandra/service/accord/CommandStoreTxnBlockedGraph.java deleted file mode 100644 index 7d553dda7c..0000000000 --- a/src/java/org/apache/cassandra/service/accord/CommandStoreTxnBlockedGraph.java +++ /dev/null @@ -1,135 +0,0 @@ -/* - * 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.service.accord; - -import java.util.ArrayList; -import java.util.LinkedHashMap; -import java.util.LinkedHashSet; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.atomic.AtomicInteger; - -import com.google.common.collect.ImmutableList; -import com.google.common.collect.ImmutableMap; -import com.google.common.collect.ImmutableSet; - -import accord.primitives.SaveStatus; -import accord.primitives.Timestamp; -import accord.primitives.TxnId; -import accord.utils.async.AsyncResults; -import org.apache.cassandra.service.accord.api.TokenKey; - -public class CommandStoreTxnBlockedGraph -{ - public final int commandStoreId; - public final Map txns; - public final Map keys; - - public CommandStoreTxnBlockedGraph(Builder builder) - { - commandStoreId = builder.storeId; - txns = ImmutableMap.copyOf(builder.txns); - keys = ImmutableMap.copyOf(builder.keys); - } - - public static class TxnState - { - public final TxnId txnId; - public final Timestamp executeAt; - public final SaveStatus saveStatus; - public final List blockedBy; - public final Set blockedByKey; - - public TxnState(Builder.TxnBuilder builder) - { - txnId = builder.txnId; - executeAt = builder.executeAt; - saveStatus = builder.saveStatus; - blockedBy = ImmutableList.copyOf(builder.blockedBy); - blockedByKey = ImmutableSet.copyOf(builder.blockedByKey); - } - - public boolean isBlocked() - { - return !notBlocked(); - } - - public boolean notBlocked() - { - return blockedBy.isEmpty() && blockedByKey.isEmpty(); - } - } - - public static class Builder extends AsyncResults.SettableResult - { - final AtomicInteger asyncTxns = new AtomicInteger(), asyncKeys = new AtomicInteger(); - final int storeId; - final Map txns = new LinkedHashMap<>(); - final Map keys = new LinkedHashMap<>(); - - public Builder(int storeId) - { - this.storeId = storeId; - } - - boolean knows(TxnId id) - { - return txns.containsKey(id); - } - - public void complete() - { - trySuccess(build()); - } - - public CommandStoreTxnBlockedGraph build() - { - return new CommandStoreTxnBlockedGraph(this); - } - - public TxnBuilder txn(TxnId txnId, Timestamp executeAt, SaveStatus saveStatus) - { - return new TxnBuilder(txnId, executeAt, saveStatus); - } - - public class TxnBuilder - { - final TxnId txnId; - final Timestamp executeAt; - final SaveStatus saveStatus; - List blockedBy = new ArrayList<>(); - Set blockedByKey = new LinkedHashSet<>(); - - public TxnBuilder(TxnId txnId, Timestamp executeAt, SaveStatus saveStatus) - { - this.txnId = txnId; - this.executeAt = executeAt; - this.saveStatus = saveStatus; - } - - public TxnState build() - { - TxnState state = new TxnState(this); - txns.put(txnId, state); - return state; - } - } - } -} diff --git a/src/java/org/apache/cassandra/service/accord/DebugBlockedTxns.java b/src/java/org/apache/cassandra/service/accord/DebugBlockedTxns.java new file mode 100644 index 0000000000..a4ca0af5e1 --- /dev/null +++ b/src/java/org/apache/cassandra/service/accord/DebugBlockedTxns.java @@ -0,0 +1,249 @@ +/* + * 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.service.accord; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.Comparator; +import java.util.List; +import java.util.Objects; +import java.util.Queue; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.TimeoutException; +import java.util.function.Consumer; +import java.util.stream.Collectors; + +import javax.annotation.Nullable; + +import accord.api.RoutingKey; +import accord.local.Command; +import accord.local.CommandStore; +import accord.local.CommandStores; +import accord.local.PreLoadContext; +import accord.local.SafeCommandStore; +import accord.local.cfk.CommandsForKey; +import accord.local.cfk.SafeCommandsForKey; +import accord.primitives.RoutingKeys; +import accord.primitives.SaveStatus; +import accord.primitives.Status; +import accord.primitives.Timestamp; +import accord.primitives.TxnId; +import accord.utils.async.AsyncChain; +import accord.utils.async.AsyncChains; +import org.apache.cassandra.service.accord.api.TokenKey; +import org.apache.cassandra.utils.concurrent.Future; + +import static accord.local.LoadKeys.SYNC; +import static accord.local.LoadKeysFor.READ_WRITE; +import static java.util.Collections.emptyList; + +public class DebugBlockedTxns +{ + public static class Txn implements Comparable + { + public final int commandStoreId; + public final int depth; + public final TxnId txnId; + public final Timestamp executeAt; + public final SaveStatus saveStatus; + public final RoutingKey blockedViaKey; + public final List blockedBy; + public final List blockedByKey; + + public Txn(int commandStoreId, int depth, TxnId txnId, Timestamp executeAt, SaveStatus saveStatus, RoutingKey blockedViaKey, List blockedBy, List blockedByKey) + { + this.commandStoreId = commandStoreId; + this.depth = depth; + this.txnId = txnId; + this.executeAt = executeAt; + this.saveStatus = saveStatus; + this.blockedViaKey = blockedViaKey; + this.blockedBy = blockedBy; + this.blockedByKey = blockedByKey; + } + + public boolean isBlocked() + { + return !notBlocked(); + } + + public boolean notBlocked() + { + return blockedBy.isEmpty() && blockedByKey.isEmpty(); + } + + @Override + public int compareTo(Txn that) + { + int c = Integer.compare(this.commandStoreId, that.commandStoreId); + if (c == 0) c = Integer.compare(this.depth, that.depth); + if (c == 0) c = this.txnId.compareTo(that.txnId); + return c; + } + } + + final IAccordService service; + final Consumer visit; + final TxnId root; + final int maxDepth; + final Set visited = Collections.newSetFromMap(new ConcurrentHashMap<>()); + final ConcurrentLinkedQueue> queuedKeys = new ConcurrentLinkedQueue<>(); + final ConcurrentLinkedQueue> queuedTxn = new ConcurrentLinkedQueue<>(); + + public DebugBlockedTxns(IAccordService service, TxnId root, int maxDepth, Consumer visit) + { + this.service = service; + this.visit = visit; + this.root = root; + this.maxDepth = maxDepth; + } + + public static void visit(IAccordService accord, TxnId txnId, int maxDepth, long deadlineNanos, Consumer visit) throws TimeoutException + { + new DebugBlockedTxns(accord, txnId, maxDepth, visit).visit(deadlineNanos); + } + + private void visit(long deadlineNanos) throws TimeoutException + { + CommandStores commandStores = service.node().commandStores(); + if (commandStores.count() == 0) + return; + + int[] ids = commandStores.ids(); + List> chains = new ArrayList<>(ids.length); + for (int id : ids) + chains.add(visitRootTxnAsync(commandStores.forId(id), root)); + + List tmp = new ArrayList<>(); + Future> next = AccordService.toFuture(AsyncChains.allOf(chains)); + while (next != null) + { + if (!next.awaitUntilThrowUncheckedOnInterrupt(deadlineNanos)) + throw new TimeoutException(); + + next.rethrowIfFailed(); + List process = next.getNow().stream() + .filter(Objects::nonNull) + .sorted(Comparator.naturalOrder()) + .collect(Collectors.toList()); + + for (Txn txn : process) + visit.accept(txn); + + Future> awaitKeys = drainToFuture(queuedKeys, (List>)(List)tmp); + if (awaitKeys != null && !awaitKeys.awaitUntilThrowUncheckedOnInterrupt(deadlineNanos)) + throw new TimeoutException(); + + next = drainToFuture(queuedTxn, (List>)(List)tmp); + } + } + + private Future> drainToFuture(Queue> drain, List> tmp) + { + AsyncChain next; + while (null != (next = drain.poll())) + tmp.add(next); + if (tmp.isEmpty()) + return null; + Future> result = AccordService.toFuture(AsyncChains.allOf(List.copyOf(tmp))); + tmp.clear(); + return result; + } + + private AsyncChain visitRootTxnAsync(CommandStore commandStore, TxnId txnId) + { + return commandStore.chain(PreLoadContext.contextFor(txnId, "Populate txn_blocked_by"), safeStore -> { + Command command = safeStore.unsafeGetNoCleanup(txnId).current(); + if (command == null || command.saveStatus() == SaveStatus.Uninitialised) + return null; + return visitTxnSync(safeStore, command, command.executeAt(), null, 0); + }); + } + + private AsyncChain visitTxnAsync(CommandStore commandStore, TxnId txnId, Timestamp rootExecuteAt, @Nullable TokenKey byKey, int depth, boolean recurse) + { + return commandStore.chain(PreLoadContext.contextFor(txnId, "Populate txn_blocked_by"), safeStore -> { + Command command = safeStore.unsafeGetNoCleanup(txnId).current(); + if (command == null || command.saveStatus() == SaveStatus.Uninitialised) + return null; + return visitTxnSync(safeStore, command, rootExecuteAt, byKey, depth); + }); + } + + private Txn visitTxnSync(SafeCommandStore safeStore, Command command, Timestamp rootExecuteAt, @Nullable TokenKey byKey, int depth) + { + List waitingOnTxnId = new ArrayList<>(); + List waitingOnKey = new ArrayList<>(); + if (!command.hasBeen(Status.Applied) && command.hasBeen(Status.Stable)) + { + // check blocking state + Command.WaitingOn waitingOn = command.asCommitted().waitingOn(); + waitingOn.waitingOn.reverseForEach(null, null, null, null, (i1, i2, i3, i4, i) -> { + if (i < waitingOn.txnIdCount()) waitingOnTxnId.add(waitingOn.txnId(i)); + else waitingOnKey.add((TokenKey) waitingOn.keys.get(i - waitingOn.txnIdCount())); + }); + } + + CommandStore commandStore = safeStore.commandStore(); + if (depth < maxDepth) + { + for (TxnId waitingOn : waitingOnTxnId) + { + if (visited.add(waitingOn)) + queuedTxn.add(visitTxnAsync(commandStore, waitingOn, rootExecuteAt, null, depth + 1, true)); + } + for (TokenKey key : waitingOnKey) + { + if (visited.add(key)) + queuedKeys.add(visitKeysAsync(commandStore, key, rootExecuteAt, depth + 1)); + } + } + + return new Txn(commandStore.id(), depth, command.txnId(), command.executeAt(), command.saveStatus(), byKey, waitingOnTxnId, waitingOnKey); + } + + + private AsyncChain visitKeysAsync(CommandStore commandStore, TokenKey key, Timestamp rootExecuteAt, int depth) + { + return commandStore.chain(PreLoadContext.contextFor(RoutingKeys.of(key.toUnseekable()), SYNC, READ_WRITE, "Populate txn_blocked_by"), safeStore -> { + visitKeysSync(safeStore, key, rootExecuteAt, depth); + }); + } + + private void visitKeysSync(SafeCommandStore safeStore, TokenKey key, Timestamp rootExecuteAt, int depth) + { + SafeCommandsForKey commandsForKey = safeStore.ifLoadedAndInitialised(key); + TxnId blocking = commandsForKey.current().blockedOnTxnId(root, rootExecuteAt); + CommandStore commandStore = safeStore.commandStore(); + if (blocking == null) + { + queuedTxn.add(AsyncChains.success(new Txn(commandStore.id(), depth, null, null, null, key, emptyList(), emptyList()))); + } + else + { + // TODO (required): this type check should not be needed; release accord version that fixes it at origin + if (blocking instanceof CommandsForKey.TxnInfo) + blocking = ((CommandsForKey.TxnInfo) blocking).plainTxnId(); + boolean recurse = visited.add(blocking); + queuedTxn.add(visitTxnAsync(commandStore, blocking, rootExecuteAt, key, depth, recurse)); + } + } +} diff --git a/src/java/org/apache/cassandra/service/accord/IAccordService.java b/src/java/org/apache/cassandra/service/accord/IAccordService.java index 422155449b..37bf99739b 100644 --- a/src/java/org/apache/cassandra/service/accord/IAccordService.java +++ b/src/java/org/apache/cassandra/service/accord/IAccordService.java @@ -19,7 +19,6 @@ package org.apache.cassandra.service.accord; import java.util.Collection; -import java.util.Collections; import java.util.EnumSet; import java.util.List; import java.util.concurrent.TimeUnit; @@ -49,7 +48,6 @@ import accord.primitives.Keys; import accord.primitives.Ranges; import accord.primitives.Timestamp; import accord.primitives.Txn; -import accord.primitives.TxnId; import accord.topology.TopologyManager; import accord.utils.Invariants; import accord.utils.async.AsyncChain; @@ -178,8 +176,6 @@ public interface IAccordService Id nodeId(); - List debugTxnBlockedGraph(TxnId txnId); - long minEpoch(); void awaitDone(TableId id, long epoch); @@ -341,12 +337,6 @@ public interface IAccordService throw new UnsupportedOperationException(); } - @Override - public List debugTxnBlockedGraph(TxnId txnId) - { - return Collections.emptyList(); - } - @Override public long minEpoch() { @@ -551,12 +541,6 @@ public interface IAccordService return delegate.nodeId(); } - @Override - public List debugTxnBlockedGraph(TxnId txnId) - { - return delegate.debugTxnBlockedGraph(txnId); - } - @Override public long minEpoch() { diff --git a/src/java/org/apache/cassandra/tools/StandaloneJournalUtil.java b/src/java/org/apache/cassandra/tools/StandaloneJournalUtil.java index 93132c5467..a14770472c 100644 --- a/src/java/org/apache/cassandra/tools/StandaloneJournalUtil.java +++ b/src/java/org/apache/cassandra/tools/StandaloneJournalUtil.java @@ -274,7 +274,7 @@ public class StandaloneJournalUtil implements Runnable Map cache = new HashMap<>(); journal.start(null); - journal.forEach(key -> processKey(cache, journal, key, txnId, sinceTimestamp, untilTimestamp, skipAllErrors, skipExceptionTypes)); + journal.forEach(key -> processKey(cache, journal, key, txnId, sinceTimestamp, untilTimestamp, skipAllErrors, skipExceptionTypes), false); } private void processKey(Map redundantBeforeCache, AccordJournal journal, JournalKey key, Timestamp txnId, Timestamp minTimestamp, Timestamp maxTimestamp, boolean skipAllErrors, Set skipExceptionTypes) diff --git a/test/distributed/org/apache/cassandra/fuzz/topology/JournalGCTest.java b/test/distributed/org/apache/cassandra/fuzz/topology/JournalGCTest.java index 5ba2a7c5d3..e4c497657d 100644 --- a/test/distributed/org/apache/cassandra/fuzz/topology/JournalGCTest.java +++ b/test/distributed/org/apache/cassandra/fuzz/topology/JournalGCTest.java @@ -113,7 +113,7 @@ public class JournalGCTest extends FuzzTestBase ((AccordService) AccordService.instance()).journal().forEach((v) -> { if (v.type == JournalKey.Type.COMMAND_DIFF && (a.get() == null || v.id.compareTo(a.get()) > 0)) a.set(v.id); - }); + }, false); return a.get() == null ? "" : a.get().toString(); }); @@ -123,7 +123,7 @@ public class JournalGCTest extends FuzzTestBase ((AccordService) AccordService.instance()).journal().forEach((v) -> { if (v.type == JournalKey.Type.COMMAND_DIFF && v.id.compareTo(maxId) <= 0) a.incrementAndGet(); - }); + }, false); return a.get(); }, maximumId); diff --git a/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java b/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java index 42db56fb94..e074ed950e 100644 --- a/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java +++ b/test/distributed/org/apache/cassandra/service/accord/AccordJournalBurnTest.java @@ -332,7 +332,7 @@ public class AccordJournalBurnTest extends BurnTestBase private TreeMap read(CommandStores commandStores) { TreeMap result = new TreeMap<>(JournalKey.SUPPORT::compare); - try (CloseableIterator> iter = journalTable.keyIterator(null, null)) + try (CloseableIterator> iter = journalTable.keyIterator(null, null, false)) { JournalKey prev = null; while (iter.hasNext()) diff --git a/test/unit/org/apache/cassandra/cql3/CQLTester.java b/test/unit/org/apache/cassandra/cql3/CQLTester.java index 23ff085c77..413b1553fc 100644 --- a/test/unit/org/apache/cassandra/cql3/CQLTester.java +++ b/test/unit/org/apache/cassandra/cql3/CQLTester.java @@ -3185,7 +3185,7 @@ public abstract class CQLTester private static String formatValue(ByteBuffer bb, AbstractType type) { - if (bb == null) + if (bb == null || (!bb.hasRemaining() && type.isEmptyValueMeaningless())) return "null"; if (type instanceof CollectionType) diff --git a/test/unit/org/apache/cassandra/db/virtual/AccordDebugKeyspaceTest.java b/test/unit/org/apache/cassandra/db/virtual/AccordDebugKeyspaceTest.java index 024e343ed7..6cac2f25fc 100644 --- a/test/unit/org/apache/cassandra/db/virtual/AccordDebugKeyspaceTest.java +++ b/test/unit/org/apache/cassandra/db/virtual/AccordDebugKeyspaceTest.java @@ -19,6 +19,8 @@ package org.apache.cassandra.db.virtual; import java.util.Collections; +import java.util.ArrayList; +import java.util.List; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -35,6 +37,7 @@ import org.slf4j.LoggerFactory; import accord.api.ProtocolModifiers; import accord.messages.NoWaitRequest; +import accord.api.RoutingKey; import accord.primitives.Ranges; import accord.primitives.Routable; import accord.primitives.SaveStatus; @@ -47,6 +50,7 @@ import org.apache.cassandra.config.OptionaldPositiveInt; import org.apache.cassandra.config.YamlConfigurationLoader; import org.apache.cassandra.cql3.CQLTester; import org.apache.cassandra.cql3.UntypedResultSet; +import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.dht.Murmur3Partitioner.LongToken; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.Message; @@ -56,11 +60,14 @@ import org.apache.cassandra.schema.Schema; import org.apache.cassandra.schema.SchemaConstants; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.service.CassandraDaemon; +import org.apache.cassandra.service.accord.AccordCommandStore; import org.apache.cassandra.service.accord.AccordService; +import org.apache.cassandra.service.accord.IAccordService; import org.apache.cassandra.service.accord.TokenRange; import org.apache.cassandra.service.accord.api.TokenKey; import org.apache.cassandra.tcm.ClusterMetadata; -import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.transport.Dispatcher; +import org.apache.cassandra.utils.Clock; import org.apache.cassandra.utils.concurrent.Condition; import org.assertj.core.api.Assertions; import org.awaitility.Awaitility; @@ -87,6 +94,12 @@ public class AccordDebugKeyspaceTest extends CQLTester private static final String QUERY_TXN = String.format("SELECT txn_id, save_status FROM %s.%s WHERE txn_id=?", SchemaConstants.VIRTUAL_ACCORD_DEBUG, AccordDebugKeyspace.TXN); + private static final String QUERY_TXNS = + String.format("SELECT save_status FROM %s.%s WHERE command_store_id = ? LIMIT 5", SchemaConstants.VIRTUAL_ACCORD_DEBUG, AccordDebugKeyspace.TXN); + + private static final String QUERY_TXNS_SEARCH = + String.format("SELECT save_status FROM %s.%s WHERE command_store_id = ? AND txn_id > ? LIMIT 5", SchemaConstants.VIRTUAL_ACCORD_DEBUG, AccordDebugKeyspace.TXN); + private static final String QUERY_JOURNAL = String.format("SELECT txn_id, save_status FROM %s.%s WHERE txn_id=?", SchemaConstants.VIRTUAL_ACCORD_DEBUG, AccordDebugKeyspace.JOURNAL); @@ -234,12 +247,42 @@ public class AccordDebugKeyspaceTest extends CQLTester getBlocking(accord.node().coordinate(id, txn)); spinUntilSuccess(() -> assertRows(execute(QUERY_TXN_BLOCKED_BY, id.toString()), - row(id.toString(), KEYSPACE, tableName, anyInt(), 0, ByteBufferUtil.EMPTY_BYTE_BUFFER, "Self", any(), null, anyOf(SaveStatus.ReadyToExecute.name(), SaveStatus.Applying.name(), SaveStatus.Applied.name())))); - spinUntilSuccess(() -> assertRows(execute(QUERY_TXN, id.toString()), row(id.toString(), "Applied"))); + row(id.toString(), anyInt(), 0, "", "", any(), anyOf(SaveStatus.ReadyToExecute.name(), SaveStatus.Applying.name(), SaveStatus.Applied.name())))); + assertRows(execute(QUERY_TXN, id.toString()), row(id.toString(), "Applied")); assertRows(execute(QUERY_JOURNAL, id.toString()), row(id.toString(), "PreAccepted"), row(id.toString(), "Applying"), row(id.toString(), "Applied"), row(id.toString(), null)); assertRows(execute(QUERY_COMMANDS_FOR_KEY, keyStr), row(id.toString(), "APPLIED_DURABLE")); } + @Test + public void manyTxns() throws ExecutionException, InterruptedException + { + String tableName = createTable("CREATE TABLE %s (k int, c int, v int, PRIMARY KEY (k, c)) WITH transactional_mode = 'full'"); + AccordService accord = accord(); + List await = new ArrayList<>(); + Txn txn = createTxn(wrapInTxn(String.format("INSERT INTO %s.%s(k, c, v) VALUES (?, ?, ?)", KEYSPACE, tableName)), 0, 0, 0); + for (int i = 0 ; i < 100; ++i) + await.add(accord.coordinateAsync(0, 0, txn, ConsistencyLevel.QUORUM, new Dispatcher.RequestTime(Clock.Global.nanoTime()))); + + AccordCommandStore commandStore = (AccordCommandStore) accord.node().commandStores().unsafeForKey((RoutingKey) txn.keys().get(0).toUnseekable()); + await.forEach(IAccordService.IAccordResult::awaitAndGet); + + assertRows(execute(QUERY_TXNS, commandStore.id()), + row("Applied"), + row("Applied"), + row("Applied"), + row("Applied"), + row("Applied") + ); + + assertRows(execute(QUERY_TXNS_SEARCH, commandStore.id(), TxnId.NONE.toString()), + row("Applied"), + row("Applied"), + row("Applied"), + row("Applied"), + row("Applied") + ); + } + @Test public void inflight() throws ExecutionException, InterruptedException { @@ -263,11 +306,10 @@ public class AccordDebugKeyspaceTest extends CQLTester filter.preAccept.awaitThrowUncheckedOnInterrupt(); assertRows(execute(QUERY_TXN_BLOCKED_BY, id.toString()), - row(id.toString(), KEYSPACE, tableName, anyInt(), 0, ByteBufferUtil.EMPTY_BYTE_BUFFER, "Self", any(), null, anyOf(SaveStatus.PreAccepted.name(), SaveStatus.ReadyToExecute.name()))); - + row(id.toString(), anyInt(), 0, "", "", any(), anyOf(SaveStatus.PreAccepted.name(), SaveStatus.ReadyToExecute.name()))); filter.apply.awaitThrowUncheckedOnInterrupt(); assertRows(execute(QUERY_TXN_BLOCKED_BY, id.toString()), - row(id.toString(), KEYSPACE, tableName, anyInt(), 0, ByteBufferUtil.EMPTY_BYTE_BUFFER, "Self", any(), null, SaveStatus.ReadyToExecute.name())); + row(id.toString(), anyInt(), 0, "", "", any(), SaveStatus.ReadyToExecute.name())); } finally { @@ -299,12 +341,13 @@ public class AccordDebugKeyspaceTest extends CQLTester accord.node().coordinate(first, createTxn(insertTxn, 0, 0, 0, 0, 0)).beginAsResult(); filter.preAccept.awaitThrowUncheckedOnInterrupt(); - spinUntilSuccess(() ->assertRows(execute(QUERY_TXN_BLOCKED_BY, first.toString()), - row(first.toString(), KEYSPACE, tableName, anyInt(), 0, ByteBufferUtil.EMPTY_BYTE_BUFFER, "Self", any(), null, anyOf(SaveStatus.PreAccepted.name(), SaveStatus.ReadyToExecute.name())))); - + assertRows(execute(QUERY_TXN_BLOCKED_BY, first.toString()), + row(first.toString(), anyInt(), 0, "", any(), any(), anyOf(SaveStatus.PreAccepted.name(), SaveStatus.ReadyToExecute.name()))); filter.apply.awaitThrowUncheckedOnInterrupt(); - spinUntilSuccess(() -> assertRows(execute(QUERY_TXN_BLOCKED_BY, first.toString()), - row(first.toString(), KEYSPACE, tableName, anyInt(), 0, ByteBufferUtil.EMPTY_BYTE_BUFFER, "Self", anyNonNull(), null, SaveStatus.ReadyToExecute.name()))); + assertRows(execute(QUERY_TXN_BLOCKED_BY, first.toString()), + row(first.toString(), anyInt(), 0, "", any(), anyNonNull(), SaveStatus.ReadyToExecute.name())); + + filter.reset(); TxnId second = accord.node().nextTxnIdWithDefaultFlags(Txn.Kind.Write, Routable.Domain.Key); filter.reset(); @@ -319,12 +362,10 @@ public class AccordDebugKeyspaceTest extends CQLTester return rs.size() == 2; }); assertRows(execute(QUERY_TXN_BLOCKED_BY, second.toString()), - row(second.toString(), KEYSPACE, tableName, anyInt(), 0, ByteBufferUtil.EMPTY_BYTE_BUFFER, "Self", anyNonNull(), null, SaveStatus.Stable.name()), - row(second.toString(), KEYSPACE, tableName, anyInt(), 1, first.toString(), "Key", anyNonNull(), anyNonNull(), SaveStatus.ReadyToExecute.name())); - + row(second.toString(), anyInt(), 0, "", "", anyNonNull(), SaveStatus.Stable.name()), + row(second.toString(), anyInt(), 1, any(), first.toString(), anyNonNull(), SaveStatus.ReadyToExecute.name())); assertRows(execute(QUERY_TXN_BLOCKED_BY + " AND depth < 1", second.toString()), - row(second.toString(), KEYSPACE, tableName, anyInt(), 0, ByteBufferUtil.EMPTY_BYTE_BUFFER, "Self", anyNonNull(), null, SaveStatus.Stable.name())); - + row(second.toString(), anyInt(), 0, any(), "", anyNonNull(), SaveStatus.Stable.name())); } finally { @@ -463,4 +504,6 @@ public class AccordDebugKeyspaceTest extends CQLTester return !dropVerbs.contains(msg.verb()); } } + + } \ No newline at end of file