This is an automated email from the ASF dual-hosted git repository. chibenwa pushed a commit to branch metric-lwt in repository https://gitbox.apache.org/repos/asf/james-project.git
commit 1c69317decc48f4f227f77a4d9431add82d1484b Author: Benoit TELLIER <[email protected]> AuthorDate: Tue Jul 21 21:24:56 2026 +0200 [METRICS] Mesure LWT time for UID and ModSeq --- mailbox/cassandra/pom.xml | 4 ++++ .../cassandra/mail/CassandraModSeqProvider.java | 21 +++++++++++++-------- .../cassandra/mail/CassandraUidProvider.java | 19 ++++++++++++------- .../cassandra/mail/CassandraMapperProvider.java | 6 ++++-- .../cassandra/mail/CassandraModSeqProviderTest.java | 7 +++++-- .../cassandra/mail/CassandraUidProviderTest.java | 7 +++++-- 6 files changed, 43 insertions(+), 21 deletions(-) diff --git a/mailbox/cassandra/pom.xml b/mailbox/cassandra/pom.xml index 76008a4e4c..2cc11e3a85 100644 --- a/mailbox/cassandra/pom.xml +++ b/mailbox/cassandra/pom.xml @@ -152,6 +152,10 @@ <groupId>${james.groupId}</groupId> <artifactId>james-server-util</artifactId> </dependency> + <dependency> + <groupId>${james.groupId}</groupId> + <artifactId>metrics-api</artifactId> + </dependency> <dependency> <groupId>${james.groupId}</groupId> <artifactId>metrics-tests</artifactId> diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraModSeqProvider.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraModSeqProvider.java index a5d0de7b4c..ea689e8a00 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraModSeqProvider.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraModSeqProvider.java @@ -46,6 +46,7 @@ import org.apache.james.mailbox.exception.MailboxException; import org.apache.james.mailbox.model.Mailbox; import org.apache.james.mailbox.model.MailboxId; import org.apache.james.mailbox.store.mail.ModSeqProvider; +import org.apache.james.metrics.api.MetricFactory; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.config.DriverExecutionProfile; @@ -61,6 +62,7 @@ import reactor.util.retry.RetryBackoffSpec; public class CassandraModSeqProvider implements ModSeqProvider { public static final String MOD_SEQ_CONDITION = "modSeqCondition"; + private static final String NEXT_MODSEQ_METRIC = "cassandra-nextModseq"; public static class ExceptionRelay extends RuntimeException { private final MailboxException underlying; @@ -94,9 +96,10 @@ public class CassandraModSeqProvider implements ModSeqProvider { private final DriverExecutionProfile lwtProfile; private final CassandraConfiguration cassandraConfiguration; private final DriverExecutionProfile readProfile; + private final MetricFactory metricFactory; @Inject - public CassandraModSeqProvider(CqlSession session, CassandraConfiguration cassandraConfiguration) { + public CassandraModSeqProvider(CqlSession session, CassandraConfiguration cassandraConfiguration, MetricFactory metricFactory) { this.cassandraAsyncExecutor = new CassandraAsyncExecutor(session); this.lwtProfile = JamesExecutionProfiles.getLWTProfile(session); this.insert = prepareInsert(session); @@ -107,6 +110,7 @@ public class CassandraModSeqProvider implements ModSeqProvider { .scheduler(Schedulers.parallel()); this.cassandraConfiguration = cassandraConfiguration; this.readProfile = ProfileLocator.READ.locateProfile(session, "MODSEQ"); + this.metricFactory = metricFactory; } private PreparedStatement prepareInsert(CqlSession session) { @@ -205,13 +209,14 @@ public class CassandraModSeqProvider implements ModSeqProvider { @Override public Mono<ModSeq> nextModSeqReactive(MailboxId mailboxId) { CassandraId cassandraId = (CassandraId) mailboxId; - return findHighestModSeq(cassandraId, Optional.of(lwtProfile)) - .flatMap(maybeHighestModSeq -> maybeHighestModSeq - .map(highestModSeq -> tryUpdateModSeq(cassandraId, highestModSeq)) - .orElseGet(() -> tryInsertModSeq(cassandraId, ModSeq.first()))) - .single() - .retryWhen(retrySpec) - .map(modSeq -> modSeq.add(cassandraConfiguration.getUidModseqIncrement())); + return Mono.from(metricFactory.decoratePublisherWithTimerMetric(NEXT_MODSEQ_METRIC, + findHighestModSeq(cassandraId, Optional.of(lwtProfile)) + .flatMap(maybeHighestModSeq -> maybeHighestModSeq + .map(highestModSeq -> tryUpdateModSeq(cassandraId, highestModSeq)) + .orElseGet(() -> tryInsertModSeq(cassandraId, ModSeq.first()))) + .single() + .retryWhen(retrySpec) + .map(modSeq -> modSeq.add(cassandraConfiguration.getUidModseqIncrement())))); } @Override diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraUidProvider.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraUidProvider.java index 3a12845669..fa39c5c7cc 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraUidProvider.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraUidProvider.java @@ -46,6 +46,7 @@ import org.apache.james.mailbox.exception.MailboxException; import org.apache.james.mailbox.model.Mailbox; import org.apache.james.mailbox.model.MailboxId; import org.apache.james.mailbox.store.mail.UidProvider; +import org.apache.james.metrics.api.MetricFactory; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.config.DriverExecutionProfile; @@ -61,6 +62,7 @@ import reactor.util.retry.RetryBackoffSpec; public class CassandraUidProvider implements UidProvider { private static final String CONDITION = "Condition"; + private static final String NEXT_UID_METRIC = "cassandra-nextUid"; private final CassandraAsyncExecutor executor; private final PreparedStatement insertStatement; @@ -70,9 +72,10 @@ public class CassandraUidProvider implements UidProvider { private final RetryBackoffSpec retrySpec; private final CassandraConfiguration cassandraConfiguration; private final DriverExecutionProfile readProfile; + private final MetricFactory metricFactory; @Inject - public CassandraUidProvider(CqlSession session, CassandraConfiguration cassandraConfiguration) { + public CassandraUidProvider(CqlSession session, CassandraConfiguration cassandraConfiguration, MetricFactory metricFactory) { this.executor = new CassandraAsyncExecutor(session); this.lwtProfile = JamesExecutionProfiles.getLWTProfile(session); this.selectStatement = prepareSelect(session); @@ -83,6 +86,7 @@ public class CassandraUidProvider implements UidProvider { .scheduler(Schedulers.parallel()); this.cassandraConfiguration = cassandraConfiguration; this.readProfile = ProfileLocator.READ.locateProfile(session, "UID"); + this.metricFactory = metricFactory; } private PreparedStatement prepareSelect(CqlSession session) { @@ -127,12 +131,13 @@ public class CassandraUidProvider implements UidProvider { Mono<MessageUid> updateUid = findHighestUid(cassandraId, Optional.of(lwtProfile)) .flatMap(messageUid -> tryUpdateUid(cassandraId, messageUid)); - return updateUid - .switchIfEmpty(tryInsert(cassandraId)) - .switchIfEmpty(updateUid) - .single() - .retryWhen(retrySpec) - .map(uid -> uid.add(cassandraConfiguration.getUidModseqIncrement())); + return Mono.from(metricFactory.decoratePublisherWithTimerMetric(NEXT_UID_METRIC, + updateUid + .switchIfEmpty(tryInsert(cassandraId)) + .switchIfEmpty(updateUid) + .single() + .retryWhen(retrySpec) + .map(uid -> uid.add(cassandraConfiguration.getUidModseqIncrement())))); } @Override diff --git a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraMapperProvider.java b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraMapperProvider.java index 4f5c310b18..9d5d3ef1d0 100644 --- a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraMapperProvider.java +++ b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraMapperProvider.java @@ -42,6 +42,7 @@ import org.apache.james.mailbox.store.mail.MessageIdMapper; import org.apache.james.mailbox.store.mail.MessageMapper; import org.apache.james.mailbox.store.mail.UidProvider; import org.apache.james.mailbox.store.mail.model.MapperProvider; +import org.apache.james.metrics.tests.RecordingMetricFactory; import org.apache.james.utils.UpdatableTickingClock; import com.google.common.collect.ImmutableList; @@ -60,10 +61,11 @@ public class CassandraMapperProvider implements MapperProvider { public CassandraMapperProvider(CassandraCluster cassandra, CassandraConfiguration cassandraConfiguration) { this.cassandra = cassandra; - messageUidProvider = new CassandraUidProvider(this.cassandra.getConf(), cassandraConfiguration); + messageUidProvider = new CassandraUidProvider(this.cassandra.getConf(), cassandraConfiguration, new RecordingMetricFactory()); cassandraModSeqProvider = new CassandraModSeqProvider( this.cassandra.getConf(), - cassandraConfiguration); + cassandraConfiguration, + new RecordingMetricFactory()); updatableTickingClock = new UpdatableTickingClock(Instant.now()); mapperFactory = createMapperFactory(cassandraConfiguration, updatableTickingClock); } diff --git a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraModSeqProviderTest.java b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraModSeqProviderTest.java index a138c3a6e0..906c69fc4e 100644 --- a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraModSeqProviderTest.java +++ b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraModSeqProviderTest.java @@ -46,6 +46,7 @@ import org.apache.james.mailbox.cassandra.modules.CassandraModSeqDataDefinition; import org.apache.james.mailbox.model.Mailbox; import org.apache.james.mailbox.model.MailboxPath; import org.apache.james.mailbox.model.UidValidity; +import org.apache.james.metrics.tests.RecordingMetricFactory; import org.apache.james.util.concurrency.ConcurrentTestRunner; import org.assertj.core.api.SoftAssertions; import org.junit.jupiter.api.BeforeEach; @@ -69,7 +70,8 @@ class CassandraModSeqProviderTest { void setUp(CassandraCluster cassandra) { modSeqProvider = new CassandraModSeqProvider( cassandra.getConf(), - CassandraConfiguration.DEFAULT_CONFIGURATION); + CassandraConfiguration.DEFAULT_CONFIGURATION, + new RecordingMetricFactory()); MailboxPath path = new MailboxPath("gsoc", Username.of("ieugen"), "Trash"); mailbox = new Mailbox(path, UidValidity.of(1234), CASSANDRA_ID); } @@ -161,7 +163,8 @@ class CassandraModSeqProviderTest { modSeqProvider = new CassandraModSeqProvider(cassandra.getConf(), CassandraConfiguration.builder() .uidModseqIncrement(10) - .build()); + .build(), + new RecordingMetricFactory()); ModSeq modseq0 = modSeqProvider.highestModSeq(mailbox); ModSeq modseq1 = modSeqProvider.nextModSeq(mailbox); diff --git a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraUidProviderTest.java b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraUidProviderTest.java index 7ff5a30796..4241c608d8 100644 --- a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraUidProviderTest.java +++ b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/CassandraUidProviderTest.java @@ -35,6 +35,7 @@ import org.apache.james.mailbox.cassandra.modules.CassandraUidDataDefinition; import org.apache.james.mailbox.model.Mailbox; import org.apache.james.mailbox.model.MailboxPath; import org.apache.james.mailbox.model.UidValidity; +import org.apache.james.metrics.tests.RecordingMetricFactory; import org.apache.james.util.concurrency.ConcurrentTestRunner; import org.assertj.core.api.SoftAssertions; import org.junit.jupiter.api.BeforeEach; @@ -56,7 +57,8 @@ class CassandraUidProviderTest { void setUp(CassandraCluster cassandra) { uidProvider = new CassandraUidProvider( cassandra.getConf(), - CassandraConfiguration.DEFAULT_CONFIGURATION); + CassandraConfiguration.DEFAULT_CONFIGURATION, + new RecordingMetricFactory()); MailboxPath path = new MailboxPath("gsoc", Username.of("ieugen"), "Trash"); mailbox = new Mailbox(path, UidValidity.of(1234), CASSANDRA_ID); } @@ -119,7 +121,8 @@ class CassandraUidProviderTest { uidProvider = new CassandraUidProvider(cassandra.getConf(), CassandraConfiguration.builder() .uidModseqIncrement(10) - .build()); + .build(), + new RecordingMetricFactory()); Optional<MessageUid> uid0 = uidProvider.lastUid(mailbox); MessageUid uid1 = uidProvider.nextUid(mailbox); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
