diff --git a/CHANGES.txt b/CHANGES.txt index cd2a14a50d..1227337073 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0 + * Add fqltool replay (CASSANDRA-14618) * Log keyspace in full query log (CASSANDRA-14656) * Transient Replication and Cheap Quorums (CASSANDRA-14404) * Log server-generated timestamp and nowInSeconds used by queries in FQL (CASSANDRA-14675) diff --git a/src/java/org/apache/cassandra/audit/FullQueryLogger.java b/src/java/org/apache/cassandra/audit/FullQueryLogger.java index c9f84476a6..9c1f472aa3 100644 --- a/src/java/org/apache/cassandra/audit/FullQueryLogger.java +++ b/src/java/org/apache/cassandra/audit/FullQueryLogger.java @@ -23,6 +23,7 @@ import java.util.List; import javax.annotation.Nullable; +import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; import com.google.common.primitives.Ints; @@ -151,7 +152,7 @@ public class FullQueryLogger extends BinLogAuditLogger implements IAuditLogger logRecord(wrappedQuery, binLog); } - static class Query extends AbstractLogEntry + public static class Query extends AbstractLogEntry { private final String query; @@ -181,7 +182,7 @@ public class FullQueryLogger extends BinLogAuditLogger implements IAuditLogger } } - static class Batch extends AbstractLogEntry + public static class Batch extends AbstractLogEntry { private final int weight; private final BatchStatement.Type batchType; diff --git a/src/java/org/apache/cassandra/service/QueryState.java b/src/java/org/apache/cassandra/service/QueryState.java index 2bd07abf33..26f58bfead 100644 --- a/src/java/org/apache/cassandra/service/QueryState.java +++ b/src/java/org/apache/cassandra/service/QueryState.java @@ -19,6 +19,7 @@ package org.apache.cassandra.service; import java.net.InetAddress; +import org.apache.cassandra.transport.ClientStat; import org.apache.cassandra.utils.FBUtilities; /** @@ -39,6 +40,13 @@ public class QueryState this.clientState = clientState; } + public QueryState(ClientState clientState, long timestamp, int nowInSeconds) + { + this(clientState); + this.timestamp = timestamp; + this.nowInSeconds = nowInSeconds; + } + /** * @return a QueryState object for internal C* calls (not limited by any kind of auth). */ diff --git a/src/java/org/apache/cassandra/tools/FullQueryLogTool.java b/src/java/org/apache/cassandra/tools/FullQueryLogTool.java index 0d170d9246..c1d4713513 100644 --- a/src/java/org/apache/cassandra/tools/FullQueryLogTool.java +++ b/src/java/org/apache/cassandra/tools/FullQueryLogTool.java @@ -31,7 +31,8 @@ import io.airlift.airline.ParseCommandUnrecognizedException; import io.airlift.airline.ParseOptionConversionException; import io.airlift.airline.ParseOptionMissingException; import io.airlift.airline.ParseOptionMissingValueException; -import org.apache.cassandra.tools.fqltool.Dump; +import org.apache.cassandra.tools.fqltool.commands.Dump; +import org.apache.cassandra.tools.fqltool.commands.Replay; import static com.google.common.base.Throwables.getStackTraceAsString; import static com.google.common.collect.Lists.newArrayList; @@ -42,7 +43,8 @@ public class FullQueryLogTool { List> commands = newArrayList( Help.class, - Dump.class + Dump.class, + Replay.class ); Cli.CliBuilder builder = Cli.builder("fqltool"); diff --git a/src/java/org/apache/cassandra/tools/fqltool/DriverResultSet.java b/src/java/org/apache/cassandra/tools/fqltool/DriverResultSet.java new file mode 100644 index 0000000000..6c4ee453d4 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/DriverResultSet.java @@ -0,0 +1,241 @@ +/* + * 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.tools.fqltool; + +import java.nio.ByteBuffer; +import java.util.Collections; +import java.util.Iterator; +import java.util.List; +import java.util.stream.Collectors; + +import com.google.common.collect.AbstractIterator; + +import com.datastax.driver.core.ColumnDefinitions; +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Row; +import org.apache.cassandra.utils.ByteBufferUtil; + + +/** + * Wraps a result set from the driver so that we can reuse the compare code when reading + * up a result set produced by ResultStore. + */ +public class DriverResultSet implements ResultHandler.ComparableResultSet +{ + private final ResultSet resultSet; + private final Throwable failureException; + + public DriverResultSet(ResultSet resultSet) + { + this(resultSet, null); + } + + private DriverResultSet(ResultSet res, Throwable failureException) + { + resultSet = res; + this.failureException = failureException; + } + + public static DriverResultSet failed(Throwable ex) + { + return new DriverResultSet(null, ex); + } + + public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() + { + if (wasFailed()) + return new DriverColumnDefinitions(null, true); + + return new DriverColumnDefinitions(resultSet.getColumnDefinitions()); + } + + public boolean wasFailed() + { + return failureException != null; + } + + public Throwable getFailureException() + { + return failureException; + } + + public Iterator iterator() + { + if (wasFailed()) + return Collections.emptyListIterator(); + return new AbstractIterator() + { + Iterator iter = resultSet.iterator(); + protected ResultHandler.ComparableRow computeNext() + { + if (iter.hasNext()) + return new DriverRow(iter.next()); + return endOfData(); + } + }; + } + + public static class DriverRow implements ResultHandler.ComparableRow + { + private final Row row; + + public DriverRow(Row row) + { + this.row = row; + } + + public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() + { + return new DriverColumnDefinitions(row.getColumnDefinitions()); + } + + public ByteBuffer getBytesUnsafe(int i) + { + return row.getBytesUnsafe(i); + } + + @Override + public boolean equals(Object oo) + { + if (!(oo instanceof ResultHandler.ComparableRow)) + return false; + + ResultHandler.ComparableRow o = (ResultHandler.ComparableRow)oo; + if (getColumnDefinitions().size() != o.getColumnDefinitions().size()) + return false; + + for (int j = 0; j < getColumnDefinitions().size(); j++) + { + ByteBuffer b1 = getBytesUnsafe(j); + ByteBuffer b2 = o.getBytesUnsafe(j); + + if (b1 != null && b2 != null && !b1.equals(b2)) + { + return false; + } + if (b1 == null && b2 != null || b2 == null && b1 != null) + { + return false; + } + } + return true; + } + + public String toString() + { + StringBuilder sb = new StringBuilder(); + List colDefs = getColumnDefinitions().asList(); + for (int i = 0; i < getColumnDefinitions().size(); i++) + { + ByteBuffer bb = getBytesUnsafe(i); + String row = bb != null ? ByteBufferUtil.bytesToHex(bb) : "NULL"; + sb.append(colDefs.get(i)).append(':').append(row).append(","); + } + return sb.toString(); + } + } + + public static class DriverColumnDefinitions implements ResultHandler.ComparableColumnDefinitions + { + private final ColumnDefinitions columnDefinitions; + private final boolean failed; + + public DriverColumnDefinitions(ColumnDefinitions columnDefinitions) + { + this(columnDefinitions, false); + } + + private DriverColumnDefinitions(ColumnDefinitions columnDefinitions, boolean failed) + { + this.columnDefinitions = columnDefinitions; + this.failed = failed; + } + + public List asList() + { + if (wasFailed()) + return Collections.emptyList(); + return columnDefinitions.asList().stream().map(DriverDefinition::new).collect(Collectors.toList()); + } + + public boolean wasFailed() + { + return failed; + } + + public int size() + { + return columnDefinitions.size(); + } + + public Iterator iterator() + { + return asList().iterator(); + } + + public boolean equals(Object oo) + { + if (!(oo instanceof ResultHandler.ComparableColumnDefinitions)) + return false; + + ResultHandler.ComparableColumnDefinitions o = (ResultHandler.ComparableColumnDefinitions)oo; + if (wasFailed() && o.wasFailed()) + return true; + + if (size() != o.size()) + return false; + + return asList().equals(o.asList()); + } + } + + public static class DriverDefinition implements ResultHandler.ComparableDefinition + { + private final ColumnDefinitions.Definition def; + + public DriverDefinition(ColumnDefinitions.Definition def) + { + this.def = def; + } + + public String getType() + { + return def.getType().toString(); + } + + public String getName() + { + return def.getName(); + } + + public boolean equals(Object oo) + { + if (!(oo instanceof ResultHandler.ComparableDefinition)) + return false; + + return def.equals(((DriverDefinition)oo).def); + } + + public String toString() + { + return getName() + ':' + getType(); + } + } + +} diff --git a/src/java/org/apache/cassandra/tools/fqltool/FQLQuery.java b/src/java/org/apache/cassandra/tools/fqltool/FQLQuery.java new file mode 100644 index 0000000000..6c0a6b9ad0 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/FQLQuery.java @@ -0,0 +1,278 @@ +/* + * 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.tools.fqltool; + +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; + +import com.google.common.primitives.Longs; + +import com.datastax.driver.core.BatchStatement; +import com.datastax.driver.core.ConsistencyLevel; +import com.datastax.driver.core.SimpleStatement; +import com.datastax.driver.core.Statement; +import org.apache.cassandra.audit.FullQueryLogger; +import org.apache.cassandra.cql3.QueryOptions; +import org.apache.cassandra.service.ClientState; +import org.apache.cassandra.service.QueryState; +import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.utils.binlog.BinLog; + +public abstract class FQLQuery implements Comparable +{ + public final long queryStartTime; + public final QueryOptions queryOptions; + public final int protocolVersion; + public final String keyspace; + public final long generatedTimestamp; + private final int generatedNowInSeconds; + + public FQLQuery(String keyspace, int protocolVersion, QueryOptions queryOptions, long queryStartTime, long generatedTimestamp, int generatedNowInSeconds) + { + this.queryStartTime = queryStartTime; + this.queryOptions = queryOptions; + this.protocolVersion = protocolVersion; + this.keyspace = keyspace; + this.generatedTimestamp = generatedTimestamp; + this.generatedNowInSeconds = generatedNowInSeconds; + } + + public abstract Statement toStatement(); + + /** + * used when storing the queries executed + */ + public abstract BinLog.ReleaseableWriteMarshallable toMarshallable(); + + public QueryState queryState() + { + ClientState clientState = keyspace != null ? ClientState.forInternalCalls(keyspace) : ClientState.forInternalCalls(); + + return new QueryState(clientState, generatedTimestamp, generatedNowInSeconds); + } + + public boolean equals(Object o) + { + if (this == o) return true; + if (!(o instanceof FQLQuery)) return false; + FQLQuery fqlQuery = (FQLQuery) o; + return queryStartTime == fqlQuery.queryStartTime && + protocolVersion == fqlQuery.protocolVersion && + generatedTimestamp == fqlQuery.generatedTimestamp && + generatedNowInSeconds == fqlQuery.generatedNowInSeconds && + Objects.equals(queryOptions.getValues(), fqlQuery.queryOptions.getValues()) && + Objects.equals(keyspace, fqlQuery.keyspace); + } + + public int hashCode() + { + return Objects.hash(queryStartTime, queryOptions, protocolVersion, keyspace, generatedTimestamp, generatedNowInSeconds); + } + + public int compareTo(FQLQuery other) + { + int cmp = Longs.compare(queryStartTime, other.queryStartTime); + if (cmp != 0) + return cmp; + cmp = Longs.compare(generatedTimestamp, other.generatedTimestamp); + if (cmp != 0) + return cmp; + + return Longs.compare(generatedNowInSeconds, other.generatedNowInSeconds); + } + + public String toString() + { + return "FQLQuery{" + + "queryStartTime=" + queryStartTime + + ", protocolVersion=" + protocolVersion + + ", keyspace='" + keyspace + '\'' + + ", generatedTimestamp=" + generatedTimestamp + + ", generatedNowInSeconds=" + generatedNowInSeconds + + '}'; + } + + public static class Single extends FQLQuery + { + public final String query; + public final List values; + + public Single(String keyspace, int protocolVersion, QueryOptions queryOptions, long queryStartTime, long generatedTimestamp, int generatedNowInSeconds, String queryString, List values) + { + super(keyspace, protocolVersion, queryOptions, queryStartTime, generatedTimestamp, generatedNowInSeconds); + this.query = queryString; + this.values = values; + } + + @Override + public String toString() + { + return String.format("%s%nQuery = %s, Values = %s", + super.toString(), + query, + values.stream().map(ByteBufferUtil::bytesToHex).collect(Collectors.joining(","))); + } + + public Statement toStatement() + { + SimpleStatement ss = new SimpleStatement(query, values.toArray()); + ss.setConsistencyLevel(ConsistencyLevel.valueOf(queryOptions.getConsistency().name())); + ss.setDefaultTimestamp(generatedTimestamp); + return ss; + } + + public BinLog.ReleaseableWriteMarshallable toMarshallable() + { + + return new FullQueryLogger.Query(query, queryOptions, queryState(), queryStartTime); + } + + public int compareTo(FQLQuery other) + { + int cmp = super.compareTo(other); + + if (cmp == 0) + { + if (other instanceof Batch) + return -1; + + Single singleQuery = (Single) other; + + cmp = query.compareTo(singleQuery.query); + if (cmp == 0) + { + if (values.size() != singleQuery.values.size()) + return values.size() - singleQuery.values.size(); + for (int i = 0; i < values.size(); i++) + { + cmp = values.get(i).compareTo(singleQuery.values.get(i)); + if (cmp != 0) + return cmp; + } + } + } + return cmp; + } + + public boolean equals(Object o) + { + if (this == o) return true; + if (!(o instanceof Single)) return false; + if (!super.equals(o)) return false; + Single single = (Single) o; + return Objects.equals(query, single.query) && + Objects.equals(values, single.values); + } + + public int hashCode() + { + return Objects.hash(super.hashCode(), query, values); + } + } + + public static class Batch extends FQLQuery + { + public final BatchStatement.Type batchType; + public final List queries; + + public Batch(String keyspace, int protocolVersion, QueryOptions queryOptions, long queryStartTime, long generatedTimestamp, int generatedNowInSeconds, BatchStatement.Type batchType, List queries, List> values) + { + super(keyspace, protocolVersion, queryOptions, queryStartTime, generatedTimestamp, generatedNowInSeconds); + this.batchType = batchType; + this.queries = new ArrayList<>(queries.size()); + for (int i = 0; i < queries.size(); i++) + this.queries.add(new Single(keyspace, protocolVersion, queryOptions, queryStartTime, generatedTimestamp, generatedNowInSeconds, queries.get(i), values.get(i))); + } + + public Statement toStatement() + { + BatchStatement bs = new BatchStatement(batchType); + + for (Single query : queries) + { + bs.add(new SimpleStatement(query.query, query.values.toArray())); + } + bs.setConsistencyLevel(ConsistencyLevel.valueOf(queryOptions.getConsistency().name())); + bs.setDefaultTimestamp(generatedTimestamp); // todo: set actual server side generated time + return bs; + } + + public int compareTo(FQLQuery other) + { + int cmp = super.compareTo(other); + + if (cmp == 0) + { + if (other instanceof Single) + return 1; + + Batch otherBatch = (Batch) other; + if (queries.size() != otherBatch.queries.size()) + return queries.size() - otherBatch.queries.size(); + for (int i = 0; i < queries.size(); i++) + { + cmp = queries.get(i).compareTo(otherBatch.queries.get(i)); + if (cmp != 0) + return cmp; + } + } + return cmp; + } + + public BinLog.ReleaseableWriteMarshallable toMarshallable() + { + List queryStrings = new ArrayList<>(); + List> values = new ArrayList<>(); + for (Single q : queries) + { + queryStrings.add(q.query); + values.add(q.values); + } + return new FullQueryLogger.Batch(org.apache.cassandra.cql3.statements.BatchStatement.Type.valueOf(batchType.name()), queryStrings, values, queryOptions, queryState(), queryStartTime); + } + + public String toString() + { + StringBuilder sb = new StringBuilder(super.toString()).append("\nbatch: ").append(batchType).append('\n'); + for (Single q : queries) + sb.append(q.toString()).append('\n'); + sb.append("end batch"); + return sb.toString(); + } + + public boolean equals(Object o) + { + if (this == o) return true; + if (!(o instanceof Batch)) return false; + if (!super.equals(o)) return false; + Batch batch = (Batch) o; + return batchType == batch.batchType && + Objects.equals(queries, batch.queries); + } + + public int hashCode() + { + return Objects.hash(super.hashCode(), batchType, queries); + } + } +} diff --git a/src/java/org/apache/cassandra/tools/fqltool/FQLQueryIterator.java b/src/java/org/apache/cassandra/tools/fqltool/FQLQueryIterator.java new file mode 100644 index 0000000000..390a52e1a1 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/FQLQueryIterator.java @@ -0,0 +1,72 @@ +/* + * 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.tools.fqltool; + +import java.util.PriorityQueue; + +import net.openhft.chronicle.queue.ExcerptTailer; +import org.apache.cassandra.utils.AbstractIterator; + +public class FQLQueryIterator extends AbstractIterator +{ + // use a priority queue to be able to sort the head of the query logs in memory + private final PriorityQueue pq; + private final ExcerptTailer tailer; + private final FQLQueryReader reader; + + /** + * Create an iterator over the FQLQueries in tailer + * + * Reads up to readAhead queries in to memory to be able to sort them (the files are mostly sorted already) + */ + public FQLQueryIterator(ExcerptTailer tailer, int readAhead) + { + assert readAhead > 0 : "readAhead needs to be > 0"; + reader = new FQLQueryReader(); + this.tailer = tailer; + pq = new PriorityQueue<>(readAhead); + for (int i = 0; i < readAhead; i++) + { + FQLQuery next = readNext(); + if (next != null) + pq.add(next); + else + break; + } + } + + protected FQLQuery computeNext() + { + FQLQuery q = pq.poll(); + if (q == null) + return endOfData(); + FQLQuery next = readNext(); + if (next != null) + pq.add(next); + return q; + } + + private FQLQuery readNext() + { + if (tailer.readDocument(reader)) + return reader.getQuery(); + return null; + } +} + diff --git a/src/java/org/apache/cassandra/tools/fqltool/FQLQueryReader.java b/src/java/org/apache/cassandra/tools/fqltool/FQLQueryReader.java new file mode 100644 index 0000000000..af77c595c3 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/FQLQueryReader.java @@ -0,0 +1,116 @@ +/* + * 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.tools.fqltool; + + +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.List; + +import com.datastax.driver.core.BatchStatement; +import io.netty.buffer.Unpooled; +import net.openhft.chronicle.core.io.IORuntimeException; +import net.openhft.chronicle.wire.ReadMarshallable; +import net.openhft.chronicle.wire.ValueIn; +import net.openhft.chronicle.wire.WireIn; +import org.apache.cassandra.cql3.QueryOptions; +import org.apache.cassandra.transport.ProtocolVersion; + +import static org.apache.cassandra.audit.FullQueryLogger.GENERATED_NOW_IN_SECONDS; +import static org.apache.cassandra.audit.FullQueryLogger.GENERATED_TIMESTAMP; +import static org.apache.cassandra.audit.FullQueryLogger.KEYSPACE; +import static org.apache.cassandra.audit.FullQueryLogger.PROTOCOL_VERSION; +import static org.apache.cassandra.audit.FullQueryLogger.QUERY_OPTIONS; +import static org.apache.cassandra.audit.FullQueryLogger.QUERY_START_TIME; +import static org.apache.cassandra.audit.FullQueryLogger.TYPE; +import static org.apache.cassandra.audit.FullQueryLogger.VERSION; +import static org.apache.cassandra.audit.FullQueryLogger.BATCH; +import static org.apache.cassandra.audit.FullQueryLogger.BATCH_TYPE; +import static org.apache.cassandra.audit.FullQueryLogger.QUERIES; +import static org.apache.cassandra.audit.FullQueryLogger.QUERY; +import static org.apache.cassandra.audit.FullQueryLogger.SINGLE_QUERY; +import static org.apache.cassandra.audit.FullQueryLogger.VALUES; + +public class FQLQueryReader implements ReadMarshallable +{ + private FQLQuery query; + + public void readMarshallable(WireIn wireIn) throws IORuntimeException + { + int currentVersion = wireIn.read(VERSION).int16(); + String type = wireIn.read(TYPE).text(); + long queryStartTime = wireIn.read(QUERY_START_TIME).int64(); + int protocolVersion = wireIn.read(PROTOCOL_VERSION).int32(); + QueryOptions queryOptions = QueryOptions.codec.decode(Unpooled.wrappedBuffer(wireIn.read(QUERY_OPTIONS).bytes()), ProtocolVersion.decode(protocolVersion)); + long generatedTimestamp = wireIn.read(GENERATED_TIMESTAMP).int64(); + int generatedNowInSeconds = wireIn.read(GENERATED_NOW_IN_SECONDS).int32(); + String keyspace = wireIn.read(KEYSPACE).text(); + + switch (type) + { + case SINGLE_QUERY: + String queryString = wireIn.read(QUERY).text(); + query = new FQLQuery.Single(keyspace, + protocolVersion, + queryOptions, + queryStartTime, + generatedTimestamp, + generatedNowInSeconds, + queryString, + queryOptions.getValues()); + break; + case BATCH: + BatchStatement.Type batchType = BatchStatement.Type.valueOf(wireIn.read(BATCH_TYPE).text()); + ValueIn in = wireIn.read(QUERIES); + int queryCount = in.int32(); + + List queries = new ArrayList<>(queryCount); + for (int i = 0; i < queryCount; i++) + queries.add(in.text()); + in = wireIn.read(VALUES); + int valueCount = in.int32(); + List> values = new ArrayList<>(valueCount); + for (int ii = 0; ii < valueCount; ii++) + { + List subValues = new ArrayList<>(); + values.add(subValues); + int numSubValues = in.int32(); + for (int zz = 0; zz < numSubValues; zz++) + subValues.add(ByteBuffer.wrap(in.bytes())); + } + query = new FQLQuery.Batch(keyspace, + protocolVersion, + queryOptions, + queryStartTime, + generatedTimestamp, + generatedNowInSeconds, + batchType, + queries, + values); + break; + default: + throw new RuntimeException("Unknown type: " + type); + } + } + + public FQLQuery getQuery() + { + return query; + } +} \ No newline at end of file diff --git a/src/java/org/apache/cassandra/tools/fqltool/QueryReplayer.java b/src/java/org/apache/cassandra/tools/fqltool/QueryReplayer.java new file mode 100644 index 0000000000..0c8382f4fb --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/QueryReplayer.java @@ -0,0 +1,167 @@ +/* + * 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.tools.fqltool; + +import java.io.Closeable; +import java.io.File; +import java.io.PrintStream; +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.function.Predicate; +import java.util.stream.Collectors; + +import javax.annotation.Nullable; + +import com.google.common.util.concurrent.FluentFuture; +import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; + +import com.codahale.metrics.MetricRegistry; +import com.codahale.metrics.Timer; +import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.ConsistencyLevel; +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.Session; +import com.datastax.driver.core.Statement; +import org.apache.cassandra.utils.FBUtilities; + +public class QueryReplayer implements Closeable +{ + private static final int PRINT_RATE = 5000; + private final ExecutorService es = Executors.newFixedThreadPool(1); + private final Iterator> queryIterator; + private final List targetClusters; + private final List> filters; + private final List sessions; + private final ResultHandler resultHandler; + private final MetricRegistry metrics = new MetricRegistry(); + private final boolean debug; + private final PrintStream out; + + public QueryReplayer(Iterator> queryIterator, + List targetHosts, + List resultPaths, + List> filters, + PrintStream out, + String queryFilePathString, + boolean debug) + { + this.queryIterator = queryIterator; + targetClusters = targetHosts.stream().map(h -> Cluster.builder().addContactPoint(h).build()).collect(Collectors.toList()); + this.filters = filters; + sessions = targetClusters.stream().map(Cluster::connect).collect(Collectors.toList()); + File queryFilePath = queryFilePathString != null ? new File(queryFilePathString) : null; + resultHandler = new ResultHandler(targetHosts, resultPaths, queryFilePath); + this.debug = debug; + this.out = out; + } + + public void replay() + { + while (queryIterator.hasNext()) + { + List queries = queryIterator.next(); + for (FQLQuery query : queries) + { + if (filters.stream().anyMatch(f -> !f.test(query))) + continue; + try (Timer.Context ctx = metrics.timer("queries").time()) + { + List> results = new ArrayList<>(sessions.size()); + Statement statement = query.toStatement(); + for (Session session : sessions) + { + try + { + if (query.keyspace != null && !query.keyspace.equals(session.getLoggedKeyspace())) + { + if (debug) + out.printf("Switching keyspace from %s to %s%n", session.getLoggedKeyspace(), query.keyspace); + session.execute("USE " + query.keyspace); + } + } + catch (Throwable t) + { + out.printf("USE %s failed: %s%n", query.keyspace, t.getMessage()); + } + if (debug) + { + out.println("Executing query:"); + out.println(query); + } + ListenableFuture future = session.executeAsync(statement); + results.add(handleErrors(future)); + } + + ListenableFuture> resultList = Futures.allAsList(results); + + Futures.addCallback(resultList, new FutureCallback>() + { + public void onSuccess(@Nullable List resultSets) + { + // note that the order of resultSets is signifcant here - resultSets.get(x) should + // be the result from a query against targetHosts.get(x) + resultHandler.handleResults(query, resultSets); + } + + public void onFailure(Throwable throwable) + { + throw new AssertionError("Errors should be handled in FQLQuery.execute", throwable); + } + }, es); + + FBUtilities.waitOnFuture(resultList); + } + catch (Throwable t) + { + out.printf("QUERY %s got exception: %s", query, t.getMessage()); + } + + Timer timer = metrics.timer("queries"); + if (timer.getCount() % PRINT_RATE == 0) + out.printf("%d queries, rate = %.2f%n", timer.getCount(), timer.getOneMinuteRate()); + } + } + } + + /** + * Make sure we catch any query errors + * + * On error, this creates a failed ComparableResultSet with the exception set to be able to store + * this fact in the result file and handle comparison of failed result sets. + */ + private static ListenableFuture handleErrors(ListenableFuture result) + { + FluentFuture fluentFuture = FluentFuture.from(result) + .transform(DriverResultSet::new, MoreExecutors.directExecutor()); + return fluentFuture.catching(Throwable.class, DriverResultSet::failed, MoreExecutors.directExecutor()); + } + + public void close() + { + sessions.forEach(Session::close); + targetClusters.forEach(Cluster::close); + } +} diff --git a/src/java/org/apache/cassandra/tools/fqltool/ResultComparator.java b/src/java/org/apache/cassandra/tools/fqltool/ResultComparator.java new file mode 100644 index 0000000000..4bbaf7a756 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/ResultComparator.java @@ -0,0 +1,116 @@ +/* + * 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.tools.fqltool; + + +import java.util.List; +import java.util.Objects; +import java.util.stream.Collectors; + +import com.google.common.collect.Streams; + +public class ResultComparator +{ + /** + * Compares the rows in rows + * the row at position x in rows will have come from host at position x in targetHosts + */ + public boolean compareRows(List targetHosts, FQLQuery query, List rows) + { + if (rows.size() < 2 || rows.stream().allMatch(Objects::isNull)) + return true; + + if (rows.stream().anyMatch(Objects::isNull)) + { + handleMismatch(targetHosts, query, rows); + return false; + } + + ResultHandler.ComparableRow ref = rows.get(0); + boolean equal = true; + for (int i = 1; i < rows.size(); i++) + { + ResultHandler.ComparableRow compare = rows.get(i); + if (!ref.equals(compare)) + equal = false; + } + if (!equal) + handleMismatch(targetHosts, query, rows); + return equal; + } + + /** + * Compares the column definitions + * + * the column definitions at position x in cds will have come from host at position x in targetHosts + */ + public boolean compareColumnDefinitions(List targetHosts, FQLQuery query, List cds) + { + if (cds.size() < 2) + return true; + + boolean equal = true; + List refDefs = cds.get(0).asList(); + for (int i = 1; i < cds.size(); i++) + { + List toCompare = cds.get(i).asList(); + if (!refDefs.equals(toCompare)) + equal = false; + } + if (!equal) + handleColumnDefMismatch(targetHosts, query, cds); + return equal; + } + + private void handleMismatch(List targetHosts, FQLQuery query, List rows) + { + System.out.println("MISMATCH:"); + System.out.println("Query = " + query); + System.out.println("Results:"); + System.out.println(Streams.zip(rows.stream(), targetHosts.stream(), (r, host) -> String.format("%s: %s%n", host, r == null ? "null" : r)).collect(Collectors.joining())); + } + + private void handleColumnDefMismatch(List targetHosts, FQLQuery query, List cds) + { + System.out.println("COLUMN DEFINITION MISMATCH:"); + System.out.println("Query = " + query); + System.out.println("Results: "); + System.out.println(Streams.zip(cds.stream(), targetHosts.stream(), (cd, host) -> String.format("%s: %s%n", host, columnDefinitionsString(cd))).collect(Collectors.joining())); + } + + private String columnDefinitionsString(ResultHandler.ComparableColumnDefinitions cd) + { + StringBuilder sb = new StringBuilder(); + if (cd == null) + sb.append("NULL"); + else if (cd.wasFailed()) + sb.append("FAILED"); + else + { + for (ResultHandler.ComparableDefinition def : cd) + { + sb.append(def.toString()); + } + } + return sb.toString(); + } + + + +} \ No newline at end of file diff --git a/src/java/org/apache/cassandra/tools/fqltool/ResultHandler.java b/src/java/org/apache/cassandra/tools/fqltool/ResultHandler.java new file mode 100644 index 0000000000..c7692319d9 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/ResultHandler.java @@ -0,0 +1,124 @@ +/* + * 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.tools.fqltool; + +import java.io.File; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; +import java.util.Objects; +import java.util.stream.Collectors; + +import com.google.common.annotations.VisibleForTesting; + +public class ResultHandler +{ + private final ResultStore resultStore; + private final ResultComparator resultComparator; + private final List targetHosts; + + public ResultHandler(List targetHosts, List resultPaths, File queryFilePath) + { + this.targetHosts = targetHosts; + resultStore = resultPaths != null ? new ResultStore(resultPaths, queryFilePath) : null; + resultComparator = new ResultComparator(); + } + + /** + * Since we can't iterate a ResultSet more than once, and we don't want to keep the entire result set in memory + * we feed the rows one-by-one to resultComparator and resultStore. + * + * results.get(x) should be the results from executing query against targetHosts.get(x) + */ + public void handleResults(FQLQuery query, List results) + { + for (int i = 0; i < targetHosts.size(); i++) + { + if (results.get(i).wasFailed()) + { + System.out.println("Query against "+targetHosts.get(i)+" failure:"); + System.out.println(query); + System.out.println("Message: "+results.get(i).getFailureException().getMessage()); + } + } + + List columnDefinitions = results.stream().map(ComparableResultSet::getColumnDefinitions).collect(Collectors.toList()); + resultComparator.compareColumnDefinitions(targetHosts, query, columnDefinitions); + if (resultStore != null) + resultStore.storeColumnDefinitions(query, columnDefinitions); + List> iters = results.stream().map(Iterable::iterator).collect(Collectors.toList()); + + while (true) + { + List rows = rows(iters); + resultComparator.compareRows(targetHosts, query, rows); + if (resultStore != null) + resultStore.storeRows(rows); + // all rows being null marks end of all resultsets, we need to call compareRows + // and storeRows once with everything null to mark that fact + if (rows.stream().allMatch(Objects::isNull)) + return; + } + } + + /** + * Get the first row from each of the iterators, if the iterator has run out, null will mark that in the list + */ + @VisibleForTesting + public static List rows(List> iters) + { + List rows = new ArrayList<>(iters.size()); + for (Iterator iter : iters) + { + if (iter.hasNext()) + rows.add(iter.next()); + else + rows.add(null); + } + return rows; + } + + public interface ComparableResultSet extends Iterable + { + public ComparableColumnDefinitions getColumnDefinitions(); + public boolean wasFailed(); + public Throwable getFailureException(); + } + + public interface ComparableColumnDefinitions extends Iterable + { + public List asList(); + public boolean wasFailed(); + public int size(); + } + + public interface ComparableDefinition + { + public String getType(); + public String getName(); + } + + public interface ComparableRow + { + public ByteBuffer getBytesUnsafe(int i); + public ComparableColumnDefinitions getColumnDefinitions(); + } + +} diff --git a/src/java/org/apache/cassandra/tools/fqltool/ResultStore.java b/src/java/org/apache/cassandra/tools/fqltool/ResultStore.java new file mode 100644 index 0000000000..6d6aaac94d --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/ResultStore.java @@ -0,0 +1,142 @@ +/* + * 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.tools.fqltool; + + +import java.io.File; +import java.nio.ByteBuffer; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.stream.Collectors; + +import net.openhft.chronicle.bytes.BytesStore; +import net.openhft.chronicle.core.io.Closeable; +import net.openhft.chronicle.queue.ChronicleQueue; +import net.openhft.chronicle.queue.ChronicleQueueBuilder; +import net.openhft.chronicle.queue.ExcerptAppender; +import net.openhft.chronicle.wire.ValueOut; +import org.apache.cassandra.utils.binlog.BinLog; + +/** + * see FQLReplayTest#readResultFile for how to read files produced by this class + */ +public class ResultStore +{ + private final List queues; + private final List appenders; + private final ChronicleQueue queryStoreQueue; + private final ExcerptAppender queryStoreAppender; + private final Set finishedHosts = new HashSet<>(); + + public ResultStore(List resultPaths, File queryFilePath) + { + queues = resultPaths.stream().map(path -> ChronicleQueueBuilder.single(path).build()).collect(Collectors.toList()); + appenders = queues.stream().map(ChronicleQueue::acquireAppender).collect(Collectors.toList()); + queryStoreQueue = queryFilePath != null ? ChronicleQueueBuilder.single(queryFilePath).build() : null; + queryStoreAppender = queryStoreQueue != null ? queryStoreQueue.acquireAppender() : null; + } + + /** + * Store the column definitions in cds + * + * the ColumnDefinitions at position x will get stored by the appender at position x + * + * Calling this method indicates that we are starting a new result set from a query, it must be called before + * calling storeRows. + * + */ + public void storeColumnDefinitions(FQLQuery query, List cds) + { + finishedHosts.clear(); + if (queryStoreAppender != null) + { + BinLog.ReleaseableWriteMarshallable writeMarshallableQuery = query.toMarshallable(); + queryStoreAppender.writeDocument(writeMarshallableQuery); + writeMarshallableQuery.release(); + } + for (int i = 0; i < cds.size(); i++) + { + ResultHandler.ComparableColumnDefinitions cd = cds.get(i); + appenders.get(i).writeDocument(wire -> + { + if (!cd.wasFailed()) + { + wire.write("type").text("column_definitions"); + wire.write("column_count").int32(cd.size()); + for (ResultHandler.ComparableDefinition d : cd.asList()) + { + ValueOut vo = wire.write("column_definition"); + vo.text(d.getName()); + vo.text(d.getType()); + } + } + else + { + wire.write("type").text("query_failed"); + } + }); + } + } + + /** + * Store rows + * + * the row at position x will get stored by appender at position x + * + * Before calling this for a new result set, storeColumnDefinitions must be called. + */ + public void storeRows(List rows) + { + for (int i = 0; i < rows.size(); i++) + { + ResultHandler.ComparableRow row = rows.get(i); + if (row == null && !finishedHosts.contains(i)) + { + appenders.get(i).writeDocument(wire -> wire.write("type").text("end_resultset")); + finishedHosts.add(i); + } + else if (row != null) + { + appenders.get(i).writeDocument(wire -> + { + { + wire.write("type").text("row"); + wire.write("row_column_count").int32(row.getColumnDefinitions().size()); + for (int jj = 0; jj < row.getColumnDefinitions().size(); jj++) + { + ByteBuffer bb = row.getBytesUnsafe(jj); + if (bb != null) + wire.write("column").bytes(BytesStore.wrap(bb)); + else + wire.write("column").bytes("NULL".getBytes()); + } + } + }); + } + } + } + + public void close() + { + queues.forEach(Closeable::close); + if (queryStoreQueue != null) + queryStoreQueue.close(); + } +} diff --git a/src/java/org/apache/cassandra/tools/fqltool/Dump.java b/src/java/org/apache/cassandra/tools/fqltool/commands/Dump.java similarity index 99% rename from src/java/org/apache/cassandra/tools/fqltool/Dump.java rename to src/java/org/apache/cassandra/tools/fqltool/commands/Dump.java index a8e75928cc..5c23d3e5c4 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/Dump.java +++ b/src/java/org/apache/cassandra/tools/fqltool/commands/Dump.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.tools.fqltool.commands; import java.io.File; import java.nio.BufferUnderflowException; diff --git a/src/java/org/apache/cassandra/tools/fqltool/commands/Replay.java b/src/java/org/apache/cassandra/tools/fqltool/commands/Replay.java new file mode 100644 index 0000000000..043ead8d41 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/fqltool/commands/Replay.java @@ -0,0 +1,148 @@ +/* + * 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.tools.fqltool.commands; + + +import java.io.File; +import java.util.ArrayList; +import java.util.List; +import java.util.function.Predicate; +import java.util.stream.Collectors; + +import com.google.common.annotations.VisibleForTesting; + +import io.airlift.airline.Arguments; +import io.airlift.airline.Command; +import io.airlift.airline.Option; +import net.openhft.chronicle.core.io.Closeable; +import net.openhft.chronicle.queue.ChronicleQueue; +import net.openhft.chronicle.queue.ChronicleQueueBuilder; + +import org.apache.cassandra.tools.fqltool.FQLQuery; +import org.apache.cassandra.tools.fqltool.FQLQueryIterator; +import org.apache.cassandra.tools.fqltool.QueryReplayer; +import org.apache.cassandra.utils.AbstractIterator; +import org.apache.cassandra.utils.MergeIterator; + +/** + * replay the contents of a list of paths containing full query logs + */ +@Command(name = "replay", description = "Replay full query logs") +public class Replay implements Runnable +{ + @Arguments(usage = " [...]", description = "Paths containing the full query logs to replay.", required = true) + private List arguments = new ArrayList<>(); + + @Option(title = "target", name = {"--target"}, description = "Hosts to replay the logs to, can be repeated to replay to more hosts.") + private List targetHosts; + + @Option(title = "results", name = { "--results"}, description = "Where to store the results of the queries, this should be a directory. Leave this option out to avoid storing results.") + private String resultPath; + + @Option(title = "keyspace", name = { "--keyspace"}, description = "Only replay queries against this keyspace and queries without keyspace set.") + private String keyspace; + + @Option(title = "debug", name = {"--debug"}, description = "Debug mode, print all queries executed.") + private boolean debug; + + @Option(title = "store_queries", name = {"--store-queries"}, description = "Path to store the queries executed. Stores queries in the same order as the result sets are in the result files. Requires --results") + private String queryStorePath; + + @Override + public void run() + { + try + { + List resultPaths = null; + if (resultPath != null) + { + File basePath = new File(resultPath); + if (!basePath.exists() || !basePath.isDirectory()) + { + System.err.println("The results path (" + basePath + ") should be an existing directory"); + System.exit(1); + } + resultPaths = targetHosts.stream().map(target -> new File(basePath, target)).collect(Collectors.toList()); + resultPaths.forEach(File::mkdir); + } + if (targetHosts.size() < 1) + { + System.err.println("You need to state at least one --target host to replay the query against"); + System.exit(1); + } + replay(keyspace, arguments, targetHosts, resultPaths, queryStorePath, debug); + } + catch (Exception e) + { + throw new RuntimeException(e); + } + } + + public static void replay(String keyspace, List arguments, List targetHosts, List resultPaths, String queryStorePath, boolean debug) + { + int readAhead = 200; // how many fql queries should we read in to memory to be able to sort them? + List readQueues = null; + List iterators = null; + List> filters = new ArrayList<>(); + + if (keyspace != null) + filters.add(fqlQuery -> fqlQuery.keyspace == null || fqlQuery.keyspace.equals(keyspace)); + + try + { + readQueues = arguments.stream().map(s -> ChronicleQueueBuilder.single(s).readOnly(true).build()).collect(Collectors.toList()); + iterators = readQueues.stream().map(ChronicleQueue::createTailer).map(tailer -> new FQLQueryIterator(tailer, readAhead)).collect(Collectors.toList()); + try (MergeIterator> iter = MergeIterator.get(iterators, FQLQuery::compareTo, new Reducer()); + QueryReplayer replayer = new QueryReplayer(iter, targetHosts, resultPaths, filters, System.out, queryStorePath, debug)) + { + replayer.replay(); + } + } + catch (Exception e) + { + throw new RuntimeException(e); + } + finally + { + if (iterators != null) + iterators.forEach(AbstractIterator::close); + if (readQueues != null) + readQueues.forEach(Closeable::close); + } + } + + @VisibleForTesting + public static class Reducer extends MergeIterator.Reducer> + { + List queries = new ArrayList<>(); + public void reduce(int idx, FQLQuery current) + { + queries.add(current); + } + + protected List getReduced() + { + return queries; + } + protected void onKeyChange() + { + queries.clear(); + } + } +} diff --git a/test/unit/org/apache/cassandra/tools/fqltool/FQLReplayTest.java b/test/unit/org/apache/cassandra/tools/fqltool/FQLReplayTest.java new file mode 100644 index 0000000000..a662699f7f --- /dev/null +++ b/test/unit/org/apache/cassandra/tools/fqltool/FQLReplayTest.java @@ -0,0 +1,760 @@ +/* + * 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.tools.fqltool; + +import java.io.File; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.file.Files; +import java.util.ArrayList; +import java.util.Collections; +import java.util.Iterator; +import java.util.List; +import java.util.Objects; +import java.util.Random; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.Collectors; + +import com.google.common.collect.AbstractIterator; +import com.google.common.collect.ImmutableList; +import com.google.common.collect.Iterables; +import com.google.common.collect.Lists; +import org.junit.Test; + +import net.openhft.chronicle.queue.ChronicleQueue; +import net.openhft.chronicle.queue.ChronicleQueueBuilder; +import net.openhft.chronicle.queue.ExcerptAppender; +import net.openhft.chronicle.queue.ExcerptTailer; +import net.openhft.chronicle.wire.ValueIn; +import org.apache.cassandra.audit.FullQueryLogger; +import org.apache.cassandra.cql3.QueryOptions; +import org.apache.cassandra.cql3.statements.BatchStatement; +import org.apache.cassandra.service.ClientState; +import org.apache.cassandra.service.QueryState; +import org.apache.cassandra.tools.Util; +import org.apache.cassandra.tools.fqltool.commands.Replay; +import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.utils.MergeIterator; +import org.apache.cassandra.utils.Pair; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; + +public class FQLReplayTest +{ + public FQLReplayTest() + { + Util.initDatabaseDescriptor(); + } + + @Test + public void testOrderedReplay() throws IOException + { + File f = generateQueries(100, true); + int queryCount = 0; + try (ChronicleQueue queue = ChronicleQueueBuilder.single(f).build(); + FQLQueryIterator iter = new FQLQueryIterator(queue.createTailer(), 101)) + { + long last = -1; + while (iter.hasNext()) + { + FQLQuery q = iter.next(); + assertTrue(q.queryStartTime >= last); + last = q.queryStartTime; + queryCount++; + } + } + assertEquals(100, queryCount); + } + @Test + public void testMergingIterator() throws IOException + { + File f = generateQueries(100, false); + File f2 = generateQueries(100, false); + int queryCount = 0; + try (ChronicleQueue queue = ChronicleQueueBuilder.single(f).build(); + ChronicleQueue queue2 = ChronicleQueueBuilder.single(f2).build(); + FQLQueryIterator iter = new FQLQueryIterator(queue.createTailer(), 101); + FQLQueryIterator iter2 = new FQLQueryIterator(queue2.createTailer(), 101); + MergeIterator> merger = MergeIterator.get(Lists.newArrayList(iter, iter2), FQLQuery::compareTo, new Replay.Reducer())) + { + long last = -1; + + while (merger.hasNext()) + { + List qs = merger.next(); + assertEquals(2, qs.size()); + assertEquals(0, qs.get(0).compareTo(qs.get(1))); + assertTrue(qs.get(0).queryStartTime >= last); + last = qs.get(0).queryStartTime; + queryCount++; + } + } + assertEquals(100, queryCount); + } + + @Test + public void testFQLQueryReader() throws IOException + { + FQLQueryReader reader = new FQLQueryReader(); + + try (ChronicleQueue queue = ChronicleQueueBuilder.single(generateQueries(1000, true)).build()) + { + ExcerptTailer tailer = queue.createTailer(); + int queryCount = 0; + while (tailer.readDocument(reader)) + { + assertNotNull(reader.getQuery()); + if (reader.getQuery() instanceof FQLQuery.Single) + { + assertTrue(reader.getQuery().keyspace == null || reader.getQuery().keyspace.equals("querykeyspace")); + } + else + { + assertEquals("someks", reader.getQuery().keyspace); + } + queryCount++; + } + assertEquals(1000, queryCount); + } + } + + @Test + public void testStoringResults() throws Throwable + { + File tmpDir = Files.createTempDirectory("results").toFile(); + File queryDir = Files.createTempDirectory("queries").toFile(); + + ResultHandler.ComparableResultSet res = createResultSet(10, 10, true); + ResultStore rs = new ResultStore(Collections.singletonList(tmpDir), queryDir); + try + { + FQLQuery query = new FQLQuery.Single("abc", 3, QueryOptions.DEFAULT, 12345, 11111, 22, "select * from abc", Collections.emptyList()); + rs.storeColumnDefinitions(query, Collections.singletonList(res.getColumnDefinitions())); + Iterator it = res.iterator(); + while (it.hasNext()) + { + List row = Collections.singletonList(it.next()); + rs.storeRows(row); + } + // this marks the end of the result set: + rs.storeRows(Collections.singletonList(null)); + } + finally + { + rs.close(); + } + + List> resultSets = readResultFile(tmpDir, queryDir); + assertEquals(1, resultSets.size()); + assertEquals(res, resultSets.get(0).right); + + } + + @Test + public void testCompareColumnDefinitions() + { + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); + ResultComparator rc = new ResultComparator(); + + List colDefs = new ArrayList<>(100); + List targetHosts = new ArrayList<>(100); + for (int i = 0; i < 100; i++) + { + targetHosts.add("host"+i); + colDefs.add(res.getColumnDefinitions()); + } + assertTrue(rc.compareColumnDefinitions(targetHosts, null, colDefs)); + colDefs.set(50, createResultSet(9, 9, false).getColumnDefinitions()); + assertFalse(rc.compareColumnDefinitions(targetHosts, null, colDefs)); + } + + @Test + public void testCompareEqualRows() + { + ResultComparator rc = new ResultComparator(); + + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); + ResultHandler.ComparableResultSet res2 = createResultSet(10, 10, false); + List toCompare = Lists.newArrayList(res, res2); + List> iters = toCompare.stream().map(Iterable::iterator).collect(Collectors.toList()); + + while (true) + { + List rows = ResultHandler.rows(iters); + assertTrue(rc.compareRows(Lists.newArrayList("eq1", "eq2"), null, rows)); + if (rows.stream().allMatch(Objects::isNull)) + break; + } + } + + @Test + public void testCompareRowsDifferentCount() + { + ResultComparator rc = new ResultComparator(); + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); + ResultHandler.ComparableResultSet res2 = createResultSet(10, 10, false); + List toCompare = Lists.newArrayList(res, res2, createResultSet(10, 11, false)); + List> iters = toCompare.stream().map(Iterable::iterator).collect(Collectors.toList()); + boolean foundMismatch = false; + while (true) + { + List rows = ResultHandler.rows(iters); + if (rows.stream().allMatch(Objects::isNull)) + break; + if (!rc.compareRows(Lists.newArrayList("eq1", "eq2", "diff"), null, rows)) + { + foundMismatch = true; + } + } + assertTrue(foundMismatch); + } + + @Test + public void testCompareRowsDifferentContent() + { + ResultComparator rc = new ResultComparator(); + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); + ResultHandler.ComparableResultSet res2 = createResultSet(10, 10, false); + List toCompare = Lists.newArrayList(res, res2, createResultSet(10, 10, true)); + List> iters = toCompare.stream().map(Iterable::iterator).collect(Collectors.toList()); + while (true) + { + List rows = ResultHandler.rows(iters); + if (rows.stream().allMatch(Objects::isNull)) + break; + assertFalse(rows.toString(), rc.compareRows(Lists.newArrayList("eq1", "eq2", "diff"), null, rows)); + } + } + + @Test + public void testCompareRowsDifferentColumnCount() + { + ResultComparator rc = new ResultComparator(); + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); + ResultHandler.ComparableResultSet res2 = createResultSet(10, 10, false); + List toCompare = Lists.newArrayList(res, res2, createResultSet(11, 10, false)); + List> iters = toCompare.stream().map(Iterable::iterator).collect(Collectors.toList()); + while (true) + { + List rows = ResultHandler.rows(iters); + if (rows.stream().allMatch(Objects::isNull)) + break; + assertFalse(rows.toString(), rc.compareRows(Lists.newArrayList("eq1", "eq2", "diff"), null, rows)); + } + } + + @Test + public void testResultHandler() throws IOException + { + List targetHosts = Lists.newArrayList("hosta", "hostb", "hostc"); + File tmpDir = Files.createTempDirectory("testresulthandler").toFile(); + File queryDir = Files.createTempDirectory("queries").toFile(); + List resultPaths = new ArrayList<>(); + targetHosts.forEach(host -> { File f = new File(tmpDir, host); f.mkdir(); resultPaths.add(f);}); + ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir); + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); + ResultHandler.ComparableResultSet res2 = createResultSet(10, 10, false); + ResultHandler.ComparableResultSet res3 = createResultSet(10, 10, false); + List toCompare = Lists.newArrayList(res, res2, res3); + rh.handleResults(new FQLQuery.Single("abcabc", 3, QueryOptions.DEFAULT, 1111, 2222, 3333, "select * from xyz", Collections.emptyList()), toCompare); + List> results1 = readResultFile(resultPaths.get(0), queryDir); + List> results2 = readResultFile(resultPaths.get(1), queryDir); + List> results3 = readResultFile(resultPaths.get(2), queryDir); + assertEquals(results1, results2); + assertEquals(results1, results3); + assertEquals(Iterables.getOnlyElement(results3).right, res); + } + + @Test + public void testResultHandlerWithDifference() throws IOException + { + List targetHosts = Lists.newArrayList("hosta", "hostb", "hostc"); + File tmpDir = Files.createTempDirectory("testresulthandler").toFile(); + File queryDir = Files.createTempDirectory("queries").toFile(); + List resultPaths = new ArrayList<>(); + targetHosts.forEach(host -> { File f = new File(tmpDir, host); f.mkdir(); resultPaths.add(f);}); + ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir); + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); + ResultHandler.ComparableResultSet res2 = createResultSet(10, 5, false); + ResultHandler.ComparableResultSet res3 = createResultSet(10, 10, false); + List toCompare = Lists.newArrayList(res, res2, res3); + rh.handleResults(new FQLQuery.Single("aaa", 3, QueryOptions.DEFAULT, 123123, 11111, 22222, "select * from abcabc", Collections.emptyList()), toCompare); + List> results1 = readResultFile(resultPaths.get(0), queryDir); + List> results2 = readResultFile(resultPaths.get(1), queryDir); + List> results3 = readResultFile(resultPaths.get(2), queryDir); + assertEquals(results1, results3); + assertEquals(results2.get(0).right, res2); + } + + @Test + public void testResultHandlerMultipleResultSets() throws IOException + { + List targetHosts = Lists.newArrayList("hosta", "hostb", "hostc"); + File tmpDir = Files.createTempDirectory("testresulthandler").toFile(); + File queryDir = Files.createTempDirectory("queries").toFile(); + List resultPaths = new ArrayList<>(); + targetHosts.forEach(host -> { File f = new File(tmpDir, host); f.mkdir(); resultPaths.add(f);}); + ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir); + List>> resultSets = new ArrayList<>(); + Random random = new Random(); + for (int i = 0; i < 10; i++) + { + List results = new ArrayList<>(); + List values = Collections.singletonList(ByteBufferUtil.bytes(i * 50)); + for (int jj = 0; jj < targetHosts.size(); jj++) + { + results.add(createResultSet(5, 1 + random.nextInt(10), true)); + } + FQLQuery q = new FQLQuery.Single("abc"+i, + 3, + QueryOptions.forInternalCalls(values), + i * 1000, + 12345, + 54321, + "select * from xyz where id = "+i, + values); + resultSets.add(Pair.create(q, results)); + } + for (int i = 0; i < resultSets.size(); i++) + rh.handleResults(resultSets.get(i).left, resultSets.get(i).right); + + for (int i = 0; i < targetHosts.size(); i++) + compareWithFile(resultPaths, queryDir, resultSets, i); + } + + @Test + public void testResultHandlerFailedQuery() throws IOException + { + List targetHosts = Lists.newArrayList("hosta", "hostb", "hostc", "hostd"); + File tmpDir = Files.createTempDirectory("testresulthandler").toFile(); + File queryDir = Files.createTempDirectory("queries").toFile(); + List resultPaths = new ArrayList<>(); + targetHosts.forEach(host -> { File f = new File(tmpDir, host); f.mkdir(); resultPaths.add(f);}); + ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir); + List>> resultSets = new ArrayList<>(); + Random random = new Random(); + for (int i = 0; i < 10; i++) + { + List results = new ArrayList<>(); + List values = Collections.singletonList(ByteBufferUtil.bytes(i * 50)); + for (int jj = 0; jj < targetHosts.size(); jj++) + { + results.add(createResultSet(5, 1 + random.nextInt(10), true)); + } + results.set(0, FakeResultSet.failed(new RuntimeException("testing abc"))); + results.set(3, FakeResultSet.failed(new RuntimeException("testing abc"))); + FQLQuery q = new FQLQuery.Single("abc"+i, + 3, + QueryOptions.forInternalCalls(values), + i * 1000, + i * 12345, + i * 54321, + "select * from xyz where id = "+i, + values); + resultSets.add(Pair.create(q, results)); + } + for (int i = 0; i < resultSets.size(); i++) + rh.handleResults(resultSets.get(i).left, resultSets.get(i).right); + + for (int i = 0; i < targetHosts.size(); i++) + compareWithFile(resultPaths, queryDir, resultSets, i); + } + + + @Test + public void testCompare() + { + FQLQuery q1 = new FQLQuery.Single("abc", 0, QueryOptions.DEFAULT, 123, 111, 222, "aaaa", Collections.emptyList()); + FQLQuery q2 = new FQLQuery.Single("abc", 0, QueryOptions.DEFAULT, 123, 111, 222,"aaaa", Collections.emptyList()); + + assertEquals(0, q1.compareTo(q2)); + assertEquals(0, q2.compareTo(q1)); + + FQLQuery q3 = new FQLQuery.Batch("abc", 0, QueryOptions.DEFAULT, 123, 111, 222, com.datastax.driver.core.BatchStatement.Type.UNLOGGED, Collections.emptyList(), Collections.emptyList()); + // single queries before batch queries + assertTrue(q1.compareTo(q3) < 0); + assertTrue(q3.compareTo(q1) > 0); + + // check that smaller query time + FQLQuery q4 = new FQLQuery.Single("abc", 0, QueryOptions.DEFAULT, 124, 111, 222, "aaaa", Collections.emptyList()); + assertTrue(q1.compareTo(q4) < 0); + assertTrue(q4.compareTo(q1) > 0); + + FQLQuery q5 = new FQLQuery.Batch("abc", 0, QueryOptions.DEFAULT, 124, 111, 222, com.datastax.driver.core.BatchStatement.Type.UNLOGGED, Collections.emptyList(), Collections.emptyList()); + assertTrue(q1.compareTo(q5) < 0); + assertTrue(q5.compareTo(q1) > 0); + + FQLQuery q6 = new FQLQuery.Single("abc", 0, QueryOptions.DEFAULT, 123, 111, 222, "aaaa", Collections.singletonList(ByteBufferUtil.bytes(10))); + FQLQuery q7 = new FQLQuery.Single("abc", 0, QueryOptions.DEFAULT, 123, 111, 222, "aaaa", Collections.emptyList()); + assertTrue(q6.compareTo(q7) > 0); + assertTrue(q7.compareTo(q6) < 0); + + FQLQuery q8 = new FQLQuery.Single("abc", 0, QueryOptions.DEFAULT, 123, 111, 222, "aaaa", Collections.singletonList(ByteBufferUtil.bytes("a"))); + FQLQuery q9 = new FQLQuery.Single("abc", 0, QueryOptions.DEFAULT, 123, 111, 222, "aaaa", Collections.singletonList(ByteBufferUtil.bytes("b"))); + assertTrue(q8.compareTo(q9) < 0); + assertTrue(q9.compareTo(q8) > 0); + } + + private File generateQueries(int count, boolean random) throws IOException + { + Random r = new Random(); + File dir = Files.createTempDirectory("chronicle").toFile(); + try (ChronicleQueue readQueue = ChronicleQueueBuilder.single(dir).build()) + { + ExcerptAppender appender = readQueue.acquireAppender(); + + for (int i = 0; i < count; i++) + { + long timestamp = random ? Math.abs(r.nextLong() % 10000) : i; + if (random ? r.nextBoolean() : i % 2 == 0) + { + String query = "abcdefghijklm " + i; + QueryState qs = r.nextBoolean() ? queryState() : queryState("querykeyspace"); + FullQueryLogger.Query q = new FullQueryLogger.Query(query, QueryOptions.DEFAULT, qs, timestamp); + appender.writeDocument(q); + q.release(); + } + else + { + int batchSize = random ? r.nextInt(99) + 1 : i + 1; + List queries = new ArrayList<>(batchSize); + List> values = new ArrayList<>(batchSize); + for (int jj = 0; jj < (random ? r.nextInt(batchSize) : 10); jj++) + { + queries.add("aaaaaa batch "+i+":"+jj); + values.add(Collections.emptyList()); + } + FullQueryLogger.Batch batch = new FullQueryLogger.Batch(BatchStatement.Type.UNLOGGED, + queries, + values, + QueryOptions.DEFAULT, + queryState("someks"), + timestamp); + appender.writeDocument(batch); + batch.release(); + } + } + } + return dir; + } + + private QueryState queryState() + { + return QueryState.forInternalCalls(); + } + + private QueryState queryState(String keyspace) + { + ClientState clientState = ClientState.forInternalCalls(keyspace); + return new QueryState(clientState); + } + + private static ResultHandler.ComparableResultSet createResultSet(int columnCount, int rowCount, boolean random) + { + List> columnDefs = new ArrayList<>(columnCount); + Random r = new Random(); + for (int i = 0; i < columnCount; i++) + { + columnDefs.add(Pair.create("a" + i, "int")); + } + List> rows = new ArrayList<>(); + for (int i = 0; i < rowCount; i++) + { + List row = new ArrayList<>(columnCount); + for (int jj = 0; jj < columnCount; jj++) + row.add(i + " col " + jj + (random ? r.nextInt() : "")); + rows.add(row); + } + return new FakeResultSet(columnDefs, rows); + } + + private static void compareWithFile(List dirs, File resultDir, List>> resultSets, int idx) + { + List> results1 = readResultFile(dirs.get(idx), resultDir); + for (int i = 0; i < results1.size(); i++) + { + assertEquals(results1.get(i).left, resultSets.get(i).left); + assertEquals(results1.get(i).right, resultSets.get(i).right.get(idx)); + } + } + + private static List> readResultFile(File dir, File queryDir) + { + List> resultSets = new ArrayList<>(); + try (ChronicleQueue q = ChronicleQueueBuilder.single(dir).build(); + ChronicleQueue queryQ = ChronicleQueueBuilder.single(queryDir).build()) + { + ExcerptTailer tailer = q.createTailer(); + ExcerptTailer queryTailer = queryQ.createTailer(); + List> columnDefinitions = new ArrayList<>(); + List> rowColumns = new ArrayList<>(); + AtomicBoolean allRowsRead = new AtomicBoolean(false); + AtomicBoolean failedQuery = new AtomicBoolean(false); + while (tailer.readDocument(wire -> { + String type = wire.read("type").text(); + if (type.equals("column_definitions")) + { + int columnCount = wire.read("column_count").int32(); + for (int i = 0; i < columnCount; i++) + { + ValueIn vi = wire.read("column_definition"); + String name = vi.text(); + String dataType = vi.text(); + columnDefinitions.add(Pair.create(name, dataType)); + } + } + else if (type.equals("row")) + { + int rowColumnCount = wire.read("row_column_count").int32(); + List r = new ArrayList<>(rowColumnCount); + for (int i = 0; i < rowColumnCount; i++) + { + byte[] b = wire.read("column").bytes(); + r.add(new String(b)); + } + rowColumns.add(r); + } + else if (type.equals("end_resultset")) + { + allRowsRead.set(true); + } + else if (type.equals("query_failed")) + { + failedQuery.set(true); + } + })) + { + if (allRowsRead.get()) + { + FQLQueryReader reader = new FQLQueryReader(); + queryTailer.readDocument(reader); + resultSets.add(Pair.create(reader.getQuery(), failedQuery.get() ? FakeResultSet.failed(new RuntimeException("failure")) + : new FakeResultSet(ImmutableList.copyOf(columnDefinitions), ImmutableList.copyOf(rowColumns)))); + allRowsRead.set(false); + failedQuery.set(false); + columnDefinitions.clear(); + rowColumns.clear(); + } + } + } + return resultSets; + } + + private static class FakeResultSet implements ResultHandler.ComparableResultSet + { + private final List> cdStrings; + private final List> rows; + private final Throwable ex; + + public FakeResultSet(List> cdStrings, List> rows) + { + this(cdStrings, rows, null); + } + + public FakeResultSet(List> cdStrings, List> rows, Throwable ex) + { + this.cdStrings = cdStrings; + this.rows = rows; + this.ex = ex; + } + + public static FakeResultSet failed(Throwable ex) + { + return new FakeResultSet(null, null, ex); + } + + public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() + { + return new FakeComparableColumnDefinitions(cdStrings, wasFailed()); + } + + public boolean wasFailed() + { + return getFailureException() != null; + } + + public Throwable getFailureException() + { + return ex; + } + + public Iterator iterator() + { + if (wasFailed()) + return Collections.emptyListIterator(); + return new AbstractIterator() + { + Iterator> iter = rows.iterator(); + protected ResultHandler.ComparableRow computeNext() + { + if (iter.hasNext()) + return new FakeComparableRow(iter.next(), cdStrings); + return endOfData(); + } + }; + } + + public boolean equals(Object o) + { + if (this == o) return true; + if (!(o instanceof FakeResultSet)) return false; + FakeResultSet that = (FakeResultSet) o; + if (wasFailed() && that.wasFailed()) + return true; + return Objects.equals(cdStrings, that.cdStrings) && + Objects.equals(rows, that.rows); + } + + public int hashCode() + { + return Objects.hash(cdStrings, rows); + } + + public String toString() + { + return "FakeResultSet{" + + "cdStrings=" + cdStrings + + ", rows=" + rows + + '}'; + } + } + + private static class FakeComparableRow implements ResultHandler.ComparableRow + { + private final List row; + private final List> cds; + + public FakeComparableRow(List row, List> cds) + { + this.row = row; + this.cds = cds; + } + + public ByteBuffer getBytesUnsafe(int i) + { + return ByteBufferUtil.bytes(row.get(i)); + } + + public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() + { + return new FakeComparableColumnDefinitions(cds, false); + } + + public boolean equals(Object other) + { + if (!(other instanceof FakeComparableRow)) + return false; + return row.equals(((FakeComparableRow)other).row); + } + + public String toString() + { + return row.toString(); + } + } + + private static class FakeComparableColumnDefinitions implements ResultHandler.ComparableColumnDefinitions + { + private final List defs; + private final boolean failed; + public FakeComparableColumnDefinitions(List> cds, boolean failed) + { + defs = cds != null ? cds.stream().map(FakeComparableDefinition::new).collect(Collectors.toList()) : null; + this.failed = failed; + } + + public List asList() + { + if (wasFailed()) + return Collections.emptyList(); + return defs; + } + + public boolean wasFailed() + { + return failed; + } + + public int size() + { + return defs.size(); + } + + public Iterator iterator() + { + if (wasFailed()) + return Collections.emptyListIterator(); + return new AbstractIterator() + { + Iterator iter = defs.iterator(); + protected ResultHandler.ComparableDefinition computeNext() + { + if (iter.hasNext()) + return iter.next(); + return endOfData(); + } + }; + } + public boolean equals(Object other) + { + if (!(other instanceof FakeComparableColumnDefinitions)) + return false; + return defs.equals(((FakeComparableColumnDefinitions)other).defs); + } + + public String toString() + { + return defs.toString(); + } + } + + private static class FakeComparableDefinition implements ResultHandler.ComparableDefinition + { + private final Pair p; + + public FakeComparableDefinition(Pair p) + { + this.p = p; + } + public String getType() + { + return p.right; + } + + public String getName() + { + return p.left; + } + + public boolean equals(Object other) + { + if (!(other instanceof FakeComparableDefinition)) + return false; + return p.equals(((FakeComparableDefinition)other).p); + } + + public String toString() + { + return getName() + ':' + getType(); + } + } +}