diff --git a/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java b/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java index 8024087de..ef327d69e 100644 --- a/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java +++ b/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java @@ -53,7 +53,7 @@ public class RemoteCacheEvent implements BinaryWritable { @Override public String toString() { - return "clearAll:" + clearAll + " caches:" + clearCaches; + return "CacheEvent[ clearAll:" + clearAll + " caches:" + clearCaches + "]"; } public static RemoteCacheEvent readBinaryMessage(BinaryReadContext dataInput) throws IOException { diff --git a/src/main/java/io/ebeaninternal/server/core/InternalConfiguration.java b/src/main/java/io/ebeaninternal/server/core/InternalConfiguration.java index f3b537217..a2090e2e0 100644 --- a/src/main/java/io/ebeaninternal/server/core/InternalConfiguration.java +++ b/src/main/java/io/ebeaninternal/server/core/InternalConfiguration.java @@ -106,7 +106,7 @@ public class InternalConfiguration { private static final Logger logger = LoggerFactory.getLogger(InternalConfiguration.class); - private final TableModState tableModState = new TableModState(); + private final TableModState tableModState; private final boolean online; @@ -169,6 +169,7 @@ public class InternalConfiguration { this.online = online; this.serverConfig = serverConfig; this.clockService = new ClockService(serverConfig.getClock()); + this.tableModState = new TableModState(clockService); this.logManager = initLogManager(); this.docStoreFactory = initDocStoreFactory(serverConfig.service(DocStoreFactory.class)); this.jsonFactory = serverConfig.getJsonFactory(); diff --git a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptor.java b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptor.java index 12ffbcb19..2cf11026b 100644 --- a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptor.java +++ b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptor.java @@ -1517,13 +1517,6 @@ public class BeanDescriptor implements BeanType, STreeType { cacheHelp.handleBulkUpdate(tableIUD); } - /** - * Handle a delete by id request adding an cache change into the changeSet. - */ - public void cacheHandleDeleteById(Object id, CacheChangeSet changeSet) { - cacheHelp.handleDelete(id, changeSet); - } - /** * Remove a bean from the cache given its Id. */ diff --git a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorManager.java b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorManager.java index a82974fa4..a517dd25c 100644 --- a/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorManager.java +++ b/src/main/java/io/ebeaninternal/server/deploy/BeanDescriptorManager.java @@ -447,8 +447,7 @@ public class BeanDescriptorManager implements BeanDescriptorMap { } /** - * For SQL based modifications we need to invalidate appropriate parts of the - * cache. + * For SQL based modifications we need to invalidate appropriate parts of the cache. */ public void cacheNotify(TransactionEventTable.TableIUD tableIUD) { diff --git a/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java b/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java index 13810e4fe..b4da3e13d 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java +++ b/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java @@ -94,6 +94,7 @@ public class BeanPersistIds implements BinaryWritable { @Override public String toString() { StringBuilder sb = new StringBuilder(); + sb.append("BeanIds["); if (beanDescriptor != null) { sb.append(beanDescriptor.getFullName()); } else { @@ -102,6 +103,7 @@ public class BeanPersistIds implements BinaryWritable { if (ids != null) { sb.append(" ids:").append(ids); } + sb.append("]"); return sb.toString(); } diff --git a/src/main/java/io/ebeaninternal/server/transaction/DeleteByIdMap.java b/src/main/java/io/ebeaninternal/server/transaction/DeleteByIdMap.java index 1bf84c0c2..edbe8a2c9 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/DeleteByIdMap.java +++ b/src/main/java/io/ebeaninternal/server/transaction/DeleteByIdMap.java @@ -7,7 +7,6 @@ import io.ebeaninternal.server.deploy.BeanDescriptor; import io.ebeanservice.docstore.api.DocStoreUpdates; import io.ebeanservice.docstore.api.support.DocStoreDeleteEvent; -import java.io.Serializable; import java.util.Collection; import java.util.LinkedHashMap; import java.util.List; @@ -22,7 +21,7 @@ public final class DeleteByIdMap { @Override public String toString() { - return beanMap.toString(); + return "DeleteById[" + beanMap.values() + "]"; } public void notifyCache(CacheChangeSet changeSet) { diff --git a/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java b/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java index 0234a6aba..8e6316a58 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java +++ b/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java @@ -20,6 +20,11 @@ public class RemoteTableMod implements BinaryWritable { this.tables = tables; } + @Override + public String toString() { + return "TableMod[" + timestamp + "; " + tables + "]"; + } + public long getTimestamp() { return timestamp; } diff --git a/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java b/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java index c89a255d1..8e328e2a2 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java +++ b/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java @@ -52,6 +52,10 @@ public class RemoteTransactionEvent implements Runnable, BinaryWritable { @Override public String toString() { StringBuilder sb = new StringBuilder(100); + sb.append("TransEvent["); + if (remoteTableMod != null) { + sb.append(remoteTableMod); + } if (!beanPersistList.isEmpty()) { sb.append(beanPersistList); } @@ -59,8 +63,9 @@ public class RemoteTransactionEvent implements Runnable, BinaryWritable { sb.append(tableList); } if (deleteByIdMap != null) { - sb.append(deleteByIdMap.values()); + sb.append(deleteByIdMap); } + sb.append("]"); return sb.toString(); } diff --git a/src/main/java/io/ebeaninternal/server/transaction/TableModState.java b/src/main/java/io/ebeaninternal/server/transaction/TableModState.java index b876a4ed0..443d35fd1 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/TableModState.java +++ b/src/main/java/io/ebeaninternal/server/transaction/TableModState.java @@ -4,6 +4,7 @@ import io.ebean.cache.QueryCacheEntry; import io.ebean.cache.QueryCacheEntryValidate; import io.ebean.cache.ServerCacheNotification; import io.ebean.cache.ServerCacheNotify; +import io.ebeaninternal.server.core.ClockService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -21,8 +22,14 @@ public class TableModState implements QueryCacheEntryValidate, ServerCacheNotify private static final Logger log = LoggerFactory.getLogger("io.ebean.cache.TABLEMOD"); + private final ClockService clockService; + private Map tableModStamp = new ConcurrentHashMap<>(); + public TableModState(ClockService clockService) { + this.clockService = clockService; + } + /** * Set the modified timestamp on the tables that have been touched. */ @@ -59,12 +66,28 @@ public class TableModState implements QueryCacheEntryValidate, ServerCacheNotify /** * Update the table modification timestamps based on remote table modification events. + *

+ * Generally this is used with distributed caches (Hazelcast, Ignite etc) via topic. + *

*/ @Override public void notify(ServerCacheNotification notification) { - // TODO: Change to use ClockService ... - long modifyTimestamp = notification.getModifyTimestamp(); - touch(notification.getDependentTables(), modifyTimestamp); + // use local clock - for slightly more aggressive invalidation (as later) + // that removes any concern regarding clock syncing across cluster + touch(notification.getDependentTables(), clockService.nowMillis()); + } + + /** + * Update from Remote transaction event. + *

+ * Generally this is used with Clustering (ebean-cluster, k8scache). + *

+ */ + public void notify(RemoteTableMod tableMod) { + + // use local clock - for slightly more aggressive invalidation (as later) + // that removes any concern regarding clock syncing across cluster + touch(tableMod.getTables(), clockService.nowMillis()); } } diff --git a/src/main/java/io/ebeaninternal/server/transaction/TransactionManager.java b/src/main/java/io/ebeaninternal/server/transaction/TransactionManager.java index 09f24c867..f330e8cce 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/TransactionManager.java +++ b/src/main/java/io/ebeaninternal/server/transaction/TransactionManager.java @@ -468,6 +468,11 @@ public class TransactionManager implements SpiTransactionManager { clusterLogger.debug("processing {}", remoteEvent); } + RemoteTableMod tableMod = remoteEvent.getRemoteTableMod(); + if (tableMod != null) { + tableModState.notify(tableMod); + } + List tableIUDList = remoteEvent.getTableIUDList(); if (tableIUDList != null) { for (TableIUD tableIUD : tableIUDList) { diff --git a/src/test/java/io/ebeaninternal/server/cache/DefaultCacheHolderTest.java b/src/test/java/io/ebeaninternal/server/cache/DefaultCacheHolderTest.java index 046ae15d6..e7a7966b7 100644 --- a/src/test/java/io/ebeaninternal/server/cache/DefaultCacheHolderTest.java +++ b/src/test/java/io/ebeaninternal/server/cache/DefaultCacheHolderTest.java @@ -4,11 +4,14 @@ import io.ebean.cache.ServerCacheFactory; import io.ebean.cache.ServerCacheOptions; import io.ebean.cache.ServerCacheType; import io.ebean.config.ServerConfig; +import io.ebeaninternal.server.core.ClockService; import io.ebeaninternal.server.transaction.TableModState; import org.tests.model.basic.Contact; import org.tests.model.basic.Customer; import org.junit.Test; +import java.time.Clock; + import static org.assertj.core.api.Assertions.assertThat; @@ -22,7 +25,7 @@ public class DefaultCacheHolderTest { private CacheManagerOptions options() { return new CacheManagerOptions(null, new ServerConfig(), true) .with(defaultOptions, defaultOptions) - .with(cacheFactory, new TableModState()); + .with(cacheFactory, new TableModState(new ClockService(Clock.systemUTC()))); } diff --git a/src/test/java/io/ebeaninternal/server/transaction/TableModStateTest.java b/src/test/java/io/ebeaninternal/server/transaction/TableModStateTest.java index 5abcb1820..1fb80996c 100644 --- a/src/test/java/io/ebeaninternal/server/transaction/TableModStateTest.java +++ b/src/test/java/io/ebeaninternal/server/transaction/TableModStateTest.java @@ -1,7 +1,9 @@ package io.ebeaninternal.server.transaction; +import io.ebeaninternal.server.core.ClockService; import org.junit.Test; +import java.time.Clock; import java.util.Collections; import java.util.HashSet; import java.util.Set; @@ -10,7 +12,7 @@ import static org.junit.Assert.*; public class TableModStateTest { - private TableModState tableModState = new TableModState(); + private TableModState tableModState = new TableModState(new ClockService(Clock.systemUTC())); @Test public void isValid() {