diff --git a/src/java/org/apache/cassandra/db/Column.java b/src/java/org/apache/cassandra/db/Column.java index 9ba3820d9a..8b9d5818af 100644 --- a/src/java/org/apache/cassandra/db/Column.java +++ b/src/java/org/apache/cassandra/db/Column.java @@ -26,7 +26,9 @@ import java.util.Collection; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.utils.ByteBufferUtil; @@ -232,5 +234,19 @@ public class Column implements IColumn { return !isMarkedForDelete(); } + + protected void validateName(CFMetaData metadata) throws MarshalException + { + AbstractType nameValidator = metadata.cfType == ColumnFamilyType.Super ? metadata.subcolumnComparator : metadata.comparator; + nameValidator.validate(name()); + } + + public void validateFields(CFMetaData metadata) throws MarshalException + { + validateName(metadata); + AbstractType valueValidator = metadata.getValueValidator(name()); + if (valueValidator != null) + valueValidator.validate(value()); + } } diff --git a/src/java/org/apache/cassandra/db/ColumnFamily.java b/src/java/org/apache/cassandra/db/ColumnFamily.java index bf154cf52d..f1018ee6ab 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamily.java +++ b/src/java/org/apache/cassandra/db/ColumnFamily.java @@ -35,6 +35,7 @@ import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.filter.QueryPath; import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.io.IColumnSerializer; import org.apache.cassandra.io.util.IIterableColumns; import org.apache.cassandra.utils.FBUtilities; @@ -424,4 +425,18 @@ public class ColumnFamily implements IColumnContainer, IIterableColumns remove(column.name()); addColumn(column.deepCopy()); } + + /** + * Goes over all columns and check the fields are valid (as far as we can + * tell). + * This is used to detect corruption after deserialization. + */ + public void validateColumnFields() throws MarshalException + { + CFMetaData metadata = metadata(); + for (IColumn column : getSortedColumns()) + { + column.validateFields(metadata); + } + } } diff --git a/src/java/org/apache/cassandra/db/DeletedColumn.java b/src/java/org/apache/cassandra/db/DeletedColumn.java index 71d5ce6304..3ee1837620 100644 --- a/src/java/org/apache/cassandra/db/DeletedColumn.java +++ b/src/java/org/apache/cassandra/db/DeletedColumn.java @@ -23,6 +23,8 @@ import java.nio.ByteBuffer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.CFMetaData; +import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.utils.ByteBufferUtil; public class DeletedColumn extends Column @@ -62,4 +64,14 @@ public class DeletedColumn extends Column { return new DeletedColumn(ByteBufferUtil.clone(name), ByteBufferUtil.clone(value), timestamp); } + + @Override + public void validateFields(CFMetaData metadata) throws MarshalException + { + validateName(metadata); + if (value().remaining() != 4) + throw new MarshalException("A tombstone value should be 4 bytes long"); + if (getLocalDeletionTime() < 0) + throw new MarshalException("The local deletion time should not be negative"); + } } diff --git a/src/java/org/apache/cassandra/db/ExpiringColumn.java b/src/java/org/apache/cassandra/db/ExpiringColumn.java index b72bcae7b3..a56fcdd2fc 100644 --- a/src/java/org/apache/cassandra/db/ExpiringColumn.java +++ b/src/java/org/apache/cassandra/db/ExpiringColumn.java @@ -24,7 +24,9 @@ import java.security.MessageDigest; import org.apache.log4j.Logger; +import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.utils.ByteBufferUtil; @@ -135,4 +137,14 @@ public class ExpiringColumn extends Column throw new IllegalStateException("column is not marked for delete"); } } + + @Override + public void validateFields(CFMetaData metadata) throws MarshalException + { + super.validateFields(metadata); + if (timeToLive <= 0) + throw new MarshalException("A column TTL should be > 0"); + if (localExpirationTime < 0) + throw new MarshalException("The local expiration time should not be negative"); + } } diff --git a/src/java/org/apache/cassandra/db/IColumn.java b/src/java/org/apache/cassandra/db/IColumn.java index 61bb6783c3..8d120fa34b 100644 --- a/src/java/org/apache/cassandra/db/IColumn.java +++ b/src/java/org/apache/cassandra/db/IColumn.java @@ -22,7 +22,9 @@ import java.nio.ByteBuffer; import java.security.MessageDigest; import java.util.Collection; +import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.utils.FBUtilities; public interface IColumn @@ -45,6 +47,7 @@ public interface IColumn public void updateDigest(MessageDigest digest); public int getLocalDeletionTime(); // for tombstone GC, so int is sufficient granularity public String getString(AbstractType comparator); + public void validateFields(CFMetaData metadata) throws MarshalException; /** clones the column, making copies of any underlying byte buffers */ IColumn deepCopy(); diff --git a/src/java/org/apache/cassandra/db/SuperColumn.java b/src/java/org/apache/cassandra/db/SuperColumn.java index abdadbb096..4434339ba2 100644 --- a/src/java/org/apache/cassandra/db/SuperColumn.java +++ b/src/java/org/apache/cassandra/db/SuperColumn.java @@ -32,7 +32,9 @@ import java.util.concurrent.atomic.AtomicLong; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.io.IColumnSerializer; import org.apache.cassandra.io.util.ColumnSortedMap; import org.apache.cassandra.io.util.DataOutputBuffer; @@ -306,6 +308,15 @@ public class SuperColumn implements IColumn, IColumnContainer { throw new UnsupportedOperationException("This operation is unsupported on super columns."); } + + public void validateFields(CFMetaData metadata) throws MarshalException + { + metadata.comparator.validate(name()); + for (IColumn column : getSubColumns()) + { + column.validateFields(metadata); + } + } } class SuperColumnSerializer implements IColumnSerializer diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java b/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java index d379902eb1..14b9acb967 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java @@ -34,6 +34,7 @@ import org.apache.cassandra.db.ColumnFamily; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.IColumn; import org.apache.cassandra.db.columniterator.IColumnIterator; +import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.io.util.BufferedRandomAccessFile; import org.apache.cassandra.utils.Filter; @@ -55,6 +56,8 @@ public class SSTableIdentityIterator implements Comparable