Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@

public class CassandraSchemaVersionManager {
public static final SchemaVersion MIN_VERSION = new SchemaVersion(12);
public static final SchemaVersion MAX_VERSION = new SchemaVersion(15);
public static final SchemaVersion MAX_VERSION = new SchemaVersion(16);
public static final SchemaVersion DEFAULT_VERSION = MIN_VERSION;

private static final Logger LOGGER = LoggerFactory.getLogger(CassandraSchemaVersionManager.class);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@
import static org.apache.james.mailbox.cassandra.table.CassandraMessageIds.MESSAGE_ID;
import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.ATTACHMENTS;
import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.BODY_CONTENT;
import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.BODY_OCTECTS;
import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.BODY_START_OCTET;
import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.FULL_CONTENT_OCTETS;
import static org.apache.james.mailbox.cassandra.table.CassandraMessageV3Table.HEADER_CONTENT;
Expand Down Expand Up @@ -158,7 +157,6 @@ private PreparedStatement prepareInsert(CqlSession session) {
.set(setColumn(INTERNAL_DATE, bindMarker(INTERNAL_DATE)),
setColumn(BODY_START_OCTET, bindMarker(BODY_START_OCTET)),
setColumn(FULL_CONTENT_OCTETS, bindMarker(FULL_CONTENT_OCTETS)),
setColumn(BODY_OCTECTS, bindMarker(BODY_OCTECTS)),
setColumn(BODY_CONTENT, bindMarker(BODY_CONTENT)),
setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT)),
prepend(ATTACHMENTS, bindMarker(ATTACHMENTS)))
Expand All @@ -184,7 +182,6 @@ public Mono<Void> save(MessageRepresentation message) {
.setInstant(INTERNAL_DATE, message.getInternalDate().toInstant())
.setInt(BODY_START_OCTET, message.getBodyStartOctet())
.setLong(FULL_CONTENT_OCTETS, message.getSize())
.setLong(BODY_OCTECTS, message.getSize() - message.getBodyStartOctet())
.setString(BODY_CONTENT, message.getBodyId().asString())
.setString(HEADER_CONTENT, message.getHeaderId().asString());

Expand Down Expand Up @@ -227,7 +224,6 @@ private BoundStatement boundWriteStatement(MailboxMessage message, Tuple2<BlobId
.setInstant(INTERNAL_DATE, message.getInternalDate().toInstant())
.setInt(BODY_START_OCTET, (int) (message.getHeaderOctets()))
.setLong(FULL_CONTENT_OCTETS, message.getFullContentOctets())
.setLong(BODY_OCTECTS, message.getBodyOctets())
.setString(BODY_CONTENT, pair.getT2().asString())
.setString(HEADER_CONTENT, pair.getT1().asString())
.setExecutionProfile(writeProfile);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,6 @@
import static org.apache.james.mailbox.cassandra.table.Flag.FLAGGED;
import static org.apache.james.mailbox.cassandra.table.Flag.RECENT;
import static org.apache.james.mailbox.cassandra.table.Flag.SEEN;
import static org.apache.james.mailbox.cassandra.table.Flag.USER;
import static org.apache.james.mailbox.cassandra.table.Flag.USER_FLAGS;
import static org.apache.james.mailbox.cassandra.table.MessageIdToImapUid.MOD_SEQ;
import static org.apache.james.util.ReactorUtils.publishIfPresent;
Expand Down Expand Up @@ -117,6 +116,7 @@ T get(Supplier<T> initializer) {
private final BlobId.Factory blobIdFactory;
private final PreparedStatement delete;
private final PreparedStatement insert;
private final PreparedStatement updateDenormalizedFields;
private final PreparedStatement select;
private final PreparedStatement selectAll;
private final PreparedStatement selectAllUids;
Expand Down Expand Up @@ -144,6 +144,7 @@ public CassandraMessageIdDAO(CqlSession session, BlobId.Factory blobIdFactory) {
this.delete = prepareDelete(session);
this.insert = prepareInsert(session);
this.update = prepareUpdate(session);
this.updateDenormalizedFields = prepareUpdateDenormalizedFields(session);
this.select = prepareSelect(session);
this.selectAll = prepareSelectAll(session);
this.selectAllUids = prepareSelectAllUids(session);
Expand Down Expand Up @@ -179,7 +180,6 @@ private PreparedStatement prepareInsert(CqlSession session) {
setColumn(FLAGGED, bindMarker(FLAGGED)),
setColumn(RECENT, bindMarker(RECENT)),
setColumn(SEEN, bindMarker(SEEN)),
setColumn(USER, bindMarker(USER)),
setColumn(INTERNAL_DATE, bindMarker(INTERNAL_DATE)),
setColumn(SAVE_DATE, bindMarker(SAVE_DATE)),
setColumn(BODY_START_OCTET, bindMarker(BODY_START_OCTET)),
Expand All @@ -191,6 +191,17 @@ private PreparedStatement prepareInsert(CqlSession session) {
.build());
}

private PreparedStatement prepareUpdateDenormalizedFields(CqlSession session) {
return session.prepare(update(TABLE_NAME)
.set(setColumn(INTERNAL_DATE, bindMarker(INTERNAL_DATE)),
setColumn(BODY_START_OCTET, bindMarker(BODY_START_OCTET)),
setColumn(FULL_CONTENT_OCTETS, bindMarker(FULL_CONTENT_OCTETS)),
setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT)))
.where(column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)),
column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID)))
.build());
}

private PreparedStatement prepareUpdate(CqlSession session) {
return session.prepare(update(TABLE_NAME)
.set(setColumn(MOD_SEQ, bindMarker(MOD_SEQ)),
Expand All @@ -200,7 +211,6 @@ private PreparedStatement prepareUpdate(CqlSession session) {
setColumn(FLAGGED, bindMarker(FLAGGED)),
setColumn(RECENT, bindMarker(RECENT)),
setColumn(SEEN, bindMarker(SEEN)),
setColumn(USER, bindMarker(USER)),
append(USER_FLAGS, bindMarker(ADDED_USERS_FLAGS)),
remove(USER_FLAGS, bindMarker(REMOVED_USERS_FLAGS)))
.where(column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)),
Expand Down Expand Up @@ -297,7 +307,6 @@ private PreparedStatement prepareSelectMetadataRange(CqlSession session) {
RECENT,
SEEN,
FLAGGED,
USER,
USER_FLAGS,
MOD_SEQ)
.where(column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)),
Expand Down Expand Up @@ -359,7 +368,6 @@ public Mono<Void> insert(CassandraMessageMetadata metadata) {
.setBoolean(FLAGGED, flags.contains(Flag.FLAGGED))
.setBoolean(RECENT, flags.contains(Flag.RECENT))
.setBoolean(SEEN, flags.contains(Flag.SEEN))
.setBoolean(USER, flags.contains(Flag.USER))
.setInstant(INTERNAL_DATE, metadata.getInternalDate().get().toInstant())
.setInstant(SAVE_DATE, metadata.getSaveDate().map(Date::toInstant).orElse(null))
.setInt(BODY_START_OCTET, Math.toIntExact(metadata.getBodyStartOctet().get()))
Expand All @@ -369,6 +377,17 @@ public Mono<Void> insert(CassandraMessageMetadata metadata) {
.build());
}

public Mono<Void> updateDenormalizedFields(CassandraId mailboxId, MessageUid uid, Date internalDate,
int bodyStartOctet, long size, BlobId headerContent) {
return cassandraAsyncExecutor.executeVoid(updateDenormalizedFields.bind()
.setUuid(MAILBOX_ID, mailboxId.asUuid())
.setLong(IMAP_UID, uid.asLong())
.setInstant(INTERNAL_DATE, internalDate.toInstant())
.setInt(BODY_START_OCTET, bodyStartOctet)
.setLong(FULL_CONTENT_OCTETS, size)
.setString(HEADER_CONTENT, headerContent.asString()));
}

public Mono<Void> updateMetadata(ComposedMessageId composedMessageId, UpdatedFlags updatedFlags) {
return cassandraAsyncExecutor.executeVoid(updateBoundStatement(composedMessageId, updatedFlags));
}
Expand Down Expand Up @@ -409,11 +428,6 @@ private BoundStatement updateBoundStatement(ComposedMessageId id, UpdatedFlags u
} else {
statementBuilder.unset(SEEN);
}
if (updatedFlags.isChanged(Flag.USER)) {
statementBuilder.setBoolean(USER, updatedFlags.isModifiedToSet(Flag.USER));
} else {
statementBuilder.unset(USER);
}
Sets.SetView<String> removedFlags = Sets.difference(
ImmutableSet.copyOf(updatedFlags.getOldFlags().getUserFlags()),
ImmutableSet.copyOf(updatedFlags.getNewFlags().getUserFlags()));
Expand Down Expand Up @@ -647,7 +661,6 @@ Mono<Void> insertNullInternalDateAndHeaderContent(CassandraMessageMetadata metad
.setBoolean(FLAGGED, flags.contains(Flag.FLAGGED))
.setBoolean(RECENT, flags.contains(Flag.RECENT))
.setBoolean(SEEN, flags.contains(Flag.SEEN))
.setBoolean(USER, flags.contains(Flag.USER))
.setInstant(INTERNAL_DATE, null)
.setInt(BODY_START_OCTET, 0)
.setLong(FULL_CONTENT_OCTETS, 0)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,11 +40,11 @@
import static org.apache.james.mailbox.cassandra.table.Flag.FLAGGED;
import static org.apache.james.mailbox.cassandra.table.Flag.RECENT;
import static org.apache.james.mailbox.cassandra.table.Flag.SEEN;
import static org.apache.james.mailbox.cassandra.table.Flag.USER;
import static org.apache.james.mailbox.cassandra.table.Flag.USER_FLAGS;
import static org.apache.james.mailbox.cassandra.table.MessageIdToImapUid.MOD_SEQ;
import static org.apache.james.mailbox.cassandra.table.MessageIdToImapUid.TABLE_NAME;
import static org.apache.james.mailbox.cassandra.table.MessageIdToImapUid.THREAD_ID;
import static org.apache.james.util.ReactorUtils.publishIfPresent;

import java.time.Duration;
import java.util.Date;
Expand Down Expand Up @@ -87,6 +87,7 @@

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;

public class CassandraMessageIdToImapUidDAO {
private static final String MOD_SEQ_CONDITION = "modSeqCondition";
Expand All @@ -98,6 +99,7 @@ public class CassandraMessageIdToImapUidDAO {
private final PreparedStatement delete;
private final PreparedStatement insert;
private final PreparedStatement update;
private final PreparedStatement updateDenormalizedFields;
private final PreparedStatement selectAll;
private final PreparedStatement select;
private final PreparedStatement listStatement;
Expand All @@ -116,6 +118,7 @@ public CassandraMessageIdToImapUidDAO(CqlSession session, BlobId.Factory blobIdF
this.cassandraConfiguration = cassandraConfiguration;
this.delete = prepareDelete();
this.insert = prepareInsert();
this.updateDenormalizedFields = prepareUpdateDenormalizedFields();
this.update = prepareUpdate();
this.selectAll = prepareSelectAll();
this.select = prepareSelect();
Expand Down Expand Up @@ -145,7 +148,6 @@ private PreparedStatement prepareInsert() {
.value(FLAGGED, bindMarker(FLAGGED))
.value(RECENT, bindMarker(RECENT))
.value(SEEN, bindMarker(SEEN))
.value(USER, bindMarker(USER))
.value(USER_FLAGS, bindMarker(USER_FLAGS))
.value(INTERNAL_DATE, bindMarker(INTERNAL_DATE))
.value(SAVE_DATE, bindMarker(SAVE_DATE))
Expand All @@ -164,7 +166,6 @@ private PreparedStatement prepareInsert() {
setColumn(FLAGGED, bindMarker(FLAGGED)),
setColumn(RECENT, bindMarker(RECENT)),
setColumn(SEEN, bindMarker(SEEN)),
setColumn(USER, bindMarker(USER)),
setColumn(INTERNAL_DATE, bindMarker(INTERNAL_DATE)),
setColumn(SAVE_DATE, bindMarker(SAVE_DATE)),
setColumn(BODY_START_OCTET, bindMarker(BODY_START_OCTET)),
Expand All @@ -178,6 +179,18 @@ private PreparedStatement prepareInsert() {
}
}

private PreparedStatement prepareUpdateDenormalizedFields() {
return session.prepare(QueryBuilder.update(TABLE_NAME)
.set(setColumn(INTERNAL_DATE, bindMarker(INTERNAL_DATE)),
setColumn(BODY_START_OCTET, bindMarker(BODY_START_OCTET)),
setColumn(FULL_CONTENT_OCTETS, bindMarker(FULL_CONTENT_OCTETS)),
setColumn(HEADER_CONTENT, bindMarker(HEADER_CONTENT)))
.where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID)),
column(MAILBOX_ID).isEqualTo(bindMarker(MAILBOX_ID)),
column(IMAP_UID).isEqualTo(bindMarker(IMAP_UID)))
.build());
}

private PreparedStatement prepareUpdate() {
Update update = QueryBuilder.update(TABLE_NAME)
.set(setColumn(MOD_SEQ, bindMarker(MOD_SEQ)),
Expand All @@ -186,8 +199,7 @@ private PreparedStatement prepareUpdate() {
setColumn(DRAFT, bindMarker(DRAFT)),
setColumn(FLAGGED, bindMarker(FLAGGED)),
setColumn(RECENT, bindMarker(RECENT)),
setColumn(SEEN, bindMarker(SEEN)),
setColumn(USER, bindMarker(USER)))
setColumn(SEEN, bindMarker(SEEN)))
.append(USER_FLAGS, bindMarker(ADDED_USERS_FLAGS))
.remove(USER_FLAGS, bindMarker(REMOVED_USERS_FLAGS))
.where(column(MESSAGE_ID).isEqualTo(bindMarker(MESSAGE_ID)),
Expand Down Expand Up @@ -251,7 +263,6 @@ public Mono<Void> insert(CassandraMessageMetadata metadata) {
.setBoolean(FLAGGED, flags.contains(Flag.FLAGGED))
.setBoolean(RECENT, flags.contains(Flag.RECENT))
.setBoolean(SEEN, flags.contains(Flag.SEEN))
.setBoolean(USER, flags.contains(Flag.USER))
.setInstant(INTERNAL_DATE, metadata.getInternalDate().get().toInstant())
.setInstant(SAVE_DATE, metadata.getSaveDate().map(Date::toInstant).orElse(null))
.setInt(BODY_START_OCTET, Math.toIntExact(metadata.getBodyStartOctet().get()))
Expand All @@ -261,6 +272,18 @@ public Mono<Void> insert(CassandraMessageMetadata metadata) {
.build());
}

public Mono<Void> updateDenormalizedFields(CassandraMessageId messageId, CassandraId mailboxId, MessageUid uid,
Date internalDate, int bodyStartOctet, long size, BlobId headerContent) {
return cassandraAsyncExecutor.executeVoid(updateDenormalizedFields.bind()
.setUuid(MESSAGE_ID, messageId.get())
.setUuid(MAILBOX_ID, mailboxId.asUuid())
.setLong(IMAP_UID, uid.asLong())
.setInstant(INTERNAL_DATE, internalDate.toInstant())
.setInt(BODY_START_OCTET, bodyStartOctet)
.setLong(FULL_CONTENT_OCTETS, size)
.setString(HEADER_CONTENT, headerContent.asString()));
}

public Mono<Boolean> updateMetadata(ComposedMessageId id, UpdatedFlags updatedFlags, ModSeq previousModeq) {
if (cassandraConfiguration.isMessageWriteStrongConsistency()) {
return cassandraAsyncExecutor.executeReturnApplied(updateBoundStatement(id, updatedFlags, previousModeq));
Expand Down Expand Up @@ -306,11 +329,6 @@ private BoundStatement updateBoundStatement(ComposedMessageId id, UpdatedFlags u
} else {
statementBuilder.unset(SEEN);
}
if (updatedFlags.isChanged(Flag.USER)) {
statementBuilder.setBoolean(USER, updatedFlags.isModifiedToSet(Flag.USER));
} else {
statementBuilder.unset(USER);
}
Sets.SetView<String> removedFlags = Sets.difference(
ImmutableSet.copyOf(updatedFlags.getOldFlags().getUserFlags()),
ImmutableSet.copyOf(updatedFlags.getNewFlags().getUserFlags()));
Expand All @@ -335,7 +353,8 @@ private BoundStatement updateBoundStatement(ComposedMessageId id, UpdatedFlags u

public Flux<CassandraMessageMetadata> retrieve(CassandraMessageId messageId, Optional<CassandraId> mailboxId, JamesExecutionProfiles.ConsistencyChoice readConsistencyChoice) {
return cassandraAsyncExecutor.executeRows(setExecutionProfileIfNeeded(selectStatement(messageId, mailboxId), readConsistencyChoice))
.map(this::toComposedMessageIdWithMetadata);
.map(this::toComposedMessageIdWithMetadata)
.handle(publishIfPresent());
}

@VisibleForTesting
Expand All @@ -346,12 +365,22 @@ public Flux<CassandraMessageMetadata> retrieve(CassandraMessageId messageId, Opt
public Flux<CassandraMessageMetadata> retrieveAllMessages() {
return cassandraAsyncExecutor.executeRows(listStatement.bind()
.setTimeout(Duration.ofDays(1)))
.map(this::toComposedMessageIdWithMetadata);
.map(this::toComposedMessageIdWithMetadata)
.handle(publishIfPresent());
}

private CassandraMessageMetadata toComposedMessageIdWithMetadata(Row row) {
private Optional<CassandraMessageMetadata> toComposedMessageIdWithMetadata(Row row) {
final CassandraMessageId messageId = CassandraMessageId.Factory.of(row.getUuid(MESSAGE_ID));
return CassandraMessageMetadata.builder()
if (row.get(MOD_SEQ, Long.class) == null) {
// Out of order updates with concurrent deletes can result in the row being partially deleted
// We filter out such records, and cleanup them.
// TODO Test INTERNAL_DATE instead once schema version 16 is enforced: unlike MOD_SEQ it also catches rows resurrected by a flag update.
delete(messageId, CassandraId.of(row.getUuid(MAILBOX_ID)))
.subscribeOn(Schedulers.parallel())
.subscribe();
return Optional.empty();
}
return Optional.of(CassandraMessageMetadata.builder()
.ids(ComposedMessageIdWithMetaData.builder()
.composedMessageId(new ComposedMessageId(
CassandraId.of(row.getUuid(MAILBOX_ID)),
Expand All @@ -369,7 +398,7 @@ private CassandraMessageMetadata toComposedMessageIdWithMetadata(Row row) {
.size(row.get(FULL_CONTENT_OCTETS, Long.class))
.headerContent(Optional.ofNullable(row.getString(HEADER_CONTENT))
.map(blobIdFactory::parse))
.build();
.build());
}

private ThreadId getThreadIdFromRow(Row row, MessageId messageId) {
Expand Down Expand Up @@ -421,7 +450,6 @@ Mono<Void> insertNullInternalDateAndHeaderContent(CassandraMessageMetadata metad
.setBoolean(FLAGGED, flags.contains(Flag.FLAGGED))
.setBoolean(RECENT, flags.contains(Flag.RECENT))
.setBoolean(SEEN, flags.contains(Flag.SEEN))
.setBoolean(USER, flags.contains(Flag.USER))
.setInstant(INTERNAL_DATE, null)
.setInt(BODY_START_OCTET, 0)
.setLong(FULL_CONTENT_OCTETS, 0)
Expand Down
Loading