diff --git a/src/main/java/com/avaje/ebean/annotation/ChangeLog.java b/src/main/java/com/avaje/ebean/annotation/ChangeLog.java
index 97e01639c..e5798cb4e 100644
--- a/src/main/java/com/avaje/ebean/annotation/ChangeLog.java
+++ b/src/main/java/com/avaje/ebean/annotation/ChangeLog.java
@@ -8,7 +8,7 @@ import java.lang.annotation.Target;
/**
* Marks an entity bean as being included in the change logging.
*/
-@Target({ ElementType.TYPE })
+@Target({ElementType.TYPE})
@Retention(RetentionPolicy.RUNTIME)
public @interface ChangeLog {
@@ -18,4 +18,14 @@ public @interface ChangeLog {
*/
boolean excludeInserts() default false;
+ /**
+ * When specified only include update requests that have at least one
+ * of the given properties as a dirty property.
+ *
+ * This provides a way to filter requests to include in the change log such that
+ * only updates that include at least one of the given properties is included
+ * in the change log.
+ *
+ */
+ String[] updatesThatInclude() default {};
}
diff --git a/src/main/java/com/avaje/ebean/bean/EntityBeanIntercept.java b/src/main/java/com/avaje/ebean/bean/EntityBeanIntercept.java
index 2a4c21a4f..42d67e7c8 100644
--- a/src/main/java/com/avaje/ebean/bean/EntityBeanIntercept.java
+++ b/src/main/java/com/avaje/ebean/bean/EntityBeanIntercept.java
@@ -581,7 +581,29 @@ public final class EntityBeanIntercept implements Serializable {
}
}
}
-
+
+ /**
+ * Return true if any of the given property names are dirty.
+ */
+ public boolean hasDirtyProperty(Set propertyNames) {
+
+ String[] names = owner._ebean_getPropertyNames();
+ int len = getPropertyLength();
+ for (int i = 0; i < len; i++) {
+ if (changedProps != null && changedProps[i]) {
+ // the property has been changed on this bean
+ if (propertyNames.contains(names[i])) {
+ return true;
+ }
+ } else if (embeddedDirty != null && embeddedDirty[i]) {
+ if (propertyNames.contains(names[i])) {
+ return true;
+ }
+ }
+ }
+ return false;
+ }
+
/**
* Return a map of dirty properties with their new and old values.
*/
diff --git a/src/main/java/com/avaje/ebean/event/BeanPersistRequest.java b/src/main/java/com/avaje/ebean/event/BeanPersistRequest.java
index cd142bff8..8c6be01d6 100644
--- a/src/main/java/com/avaje/ebean/event/BeanPersistRequest.java
+++ b/src/main/java/com/avaje/ebean/event/BeanPersistRequest.java
@@ -1,12 +1,12 @@
package com.avaje.ebean.event;
-import java.util.Map;
-import java.util.Set;
-
import com.avaje.ebean.EbeanServer;
import com.avaje.ebean.Transaction;
import com.avaje.ebean.ValuePair;
+import java.util.Map;
+import java.util.Set;
+
/**
* Holds the information available for a bean persist (insert, update or
* delete).
@@ -34,9 +34,30 @@ public interface BeanPersistRequest {
/**
* For an update this is the set of properties that where updated.
+ *
+ * Note that hasDirtyProperty() is a more efficient check than this method and
+ * should be preferred if it satisfies the requirement.
+ *
*/
Set getUpdatedProperties();
+ /**
+ * Return true for an update request if at least one of dirty properties is contained
+ * in the given set of property names.
+ *
+ * This method will produce less GC compared with getUpdatedProperties() and should
+ * be preferred if it satisfies the requirement.
+ *
+ *
+ * Note that this method is used by the default ChangeLogFilter mechanism for when
+ * the @ChangeLog updatesThatInclude attribute has been specified.
+ *
+ *
+ * @param propertyNames a set of property names which we are checking to see if at least
+ * one of them is dirty.
+ */
+ boolean hasDirtyProperty(Set propertyNames);
+
/**
* Returns the bean being inserted updated or deleted.
*/
@@ -45,6 +66,6 @@ public interface BeanPersistRequest {
/**
* Returns a map of the properties that have changed and their new and old values.
*/
- Map getUpdatedValues();
+ Map getUpdatedValues();
}
diff --git a/src/main/java/com/avaje/ebean/plugin/SpiServerPlugin.java b/src/main/java/com/avaje/ebean/plugin/SpiServerPlugin.java
new file mode 100644
index 000000000..a622f89d9
--- /dev/null
+++ b/src/main/java/com/avaje/ebean/plugin/SpiServerPlugin.java
@@ -0,0 +1,17 @@
+package com.avaje.ebean.plugin;
+
+/**
+ * A 'plugin' that wants to be configured on startup so it can use features of the EbeanServer itself.
+ */
+public interface SpiServerPlugin {
+
+ /**
+ * Configure the plugin.
+ */
+ void configure(SpiServer server);
+
+ /**
+ * Called just before the server starts indicating if it is coming up in online mode.
+ */
+ void online(boolean online);
+}
diff --git a/src/main/java/com/avaje/ebeaninternal/server/changelog/BulkJsonBuilder.java b/src/main/java/com/avaje/ebeaninternal/server/changelog/BulkJsonBuilder.java
new file mode 100644
index 000000000..18c8bfb02
--- /dev/null
+++ b/src/main/java/com/avaje/ebeaninternal/server/changelog/BulkJsonBuilder.java
@@ -0,0 +1,173 @@
+package com.avaje.ebeaninternal.server.changelog;
+
+import com.avaje.ebean.ValuePair;
+import com.avaje.ebean.event.changelog.BeanChange;
+import com.avaje.ebean.event.changelog.ChangeSet;
+import com.avaje.ebean.event.changelog.ChangeType;
+import com.avaje.ebean.text.json.JsonContext;
+import com.avaje.ebean.text.json.JsonScalar;
+import com.fasterxml.jackson.core.JsonFactory;
+import com.fasterxml.jackson.core.JsonGenerator;
+
+import java.io.IOException;
+import java.io.Writer;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Builds JSON appropriate for loading into ElasticS via the bulk API.
+ */
+public class BulkJsonBuilder {
+
+ protected final JsonFactory jsonFactory = new JsonFactory();
+
+ protected final String indexName;
+
+ protected final String indexType;
+
+ protected final JsonContext json;
+
+ protected BulkJsonBuilder(JsonContext json, String indexName, String indexType) {
+ this.json = json;
+ this.indexName = indexName;
+ this.indexType = indexType;
+ }
+
+ /**
+ * Write the change set into Elastic bulk API JSON form (so contains special new line
+ * characters and bulk API header etc.
+ */
+ public void writeJson(ChangeSet changeSet, Writer writer) throws IOException {
+
+ JsonGenerator generator = jsonFactory.createGenerator(writer);
+
+ List changes = changeSet.getChanges();
+ for (int i = 0; i < changes.size(); i++) {
+ write(generator, changes.get(i), changeSet, i);
+ }
+
+ generator.flush();
+ generator.close();
+ }
+
+ /**
+ * Write the bean change as a single JSON document for storage into Elastic.
+ *
+ * Note that for ease of search/use we effectively denormalise by including the transaction header
+ * information in each bean document.
+ *
+ */
+ protected void write(JsonGenerator gen, BeanChange bean, ChangeSet changeSet, int position) throws IOException {
+
+ writeBulkHeader(gen, changeSet, position);
+ writeBeanChange(gen, bean, changeSet);
+ writeBeanChangeEnd(gen);
+ }
+
+ /**
+ * Write the bean change as JSON document containing the transaction header details.
+ */
+ protected void writeBeanChange(JsonGenerator gen, BeanChange bean, ChangeSet changeSet) throws IOException {
+
+ gen.writeStartObject();
+
+ writeBeanTransactionDetails(gen, changeSet);
+
+ gen.writeStringField("object", bean.getTable());
+ gen.writeStringField("objectId", bean.getId().toString());
+ gen.writeStringField("change", bean.getType().getCode());
+ gen.writeNumberField("eventTime", bean.getEventTime());
+
+ writeBeanValues(gen, bean);
+
+ gen.writeEndObject();
+ }
+
+ /**
+ * For Elastic bulk we append raw new line character.
+ */
+ protected void writeBeanChangeEnd(JsonGenerator gen) throws IOException {
+ gen.writeRawValue("\n");
+ }
+
+ /**
+ * Denormalise by writing the transaction header details.
+ */
+ protected void writeBeanTransactionDetails(JsonGenerator gen, ChangeSet changeSet) throws IOException {
+
+ gen.writeStringField("txnId", changeSet.getTxnId());
+ gen.writeStringField("txnState", changeSet.getTxnState().getCode());
+ gen.writeNumberField("txnBatch", changeSet.getTxnBatch());
+ gen.writeStringField("userId", changeSet.getUserId());
+ String userIpAddress = changeSet.getUserIpAddress();
+ if (userIpAddress != null) {
+ gen.writeStringField("userIpAddress", userIpAddress);
+ }
+ String userContext = changeSet.getUserContext();
+ if (userContext != null) {
+ gen.writeStringField("userContext", userContext);
+ }
+ }
+
+ /**
+ * Write the elastic bulk API header.
+ */
+ protected void writeBulkHeader(JsonGenerator gen, ChangeSet changeSet, int position) throws IOException {
+
+ // we index with an 'id' value so that we can process/reprocess the JSON and
+ // avoid duplicates being inserted. Appending the batch and position give us
+ // an effectively unique id value for the change
+ String uid = changeSet.getTxnId() + "_" + changeSet.getTxnBatch() + "." + position;
+
+ // the 'header' for elastic bulk API
+ gen.writeStartObject();
+ gen.writeFieldName("index");
+ gen.writeStartObject();
+ gen.writeStringField("_index", indexName);
+ gen.writeStringField("_type", indexType);
+ gen.writeStringField("_id", uid);
+ gen.writeEndObject();
+ gen.writeEndObject();
+ gen.writeRawValue("\n");
+ }
+
+ /**
+ * For insert and update write the new/old values.
+ */
+ protected void writeBeanValues(JsonGenerator gen, BeanChange bean) throws IOException {
+ if (bean.getType() != ChangeType.DELETE) {
+ gen.writeFieldName("values");
+ gen.writeStartObject();
+ // use JsonScalar as it knows how to encode all the scalar
+ // property types that Ebean supports (Java8, Joda etc)
+ JsonScalar scalarWriter = json.getScalar(gen);
+ writeValuePairs(bean, scalarWriter, gen);
+ gen.writeEndObject();
+ }
+ }
+
+ /**
+ * Write all the value pairs suppressing null values.
+ *
+ * We are intentionally keeping the same new/old structure for both inserts and updates.
+ *
+ */
+ protected void writeValuePairs(BeanChange bean, JsonScalar scalarWriter, JsonGenerator gen) throws IOException {
+
+ for (Map.Entry entry : bean.getValues().entrySet()) {
+ gen.writeFieldName(entry.getKey());
+ gen.writeStartObject();
+ ValuePair value = entry.getValue();
+ Object newValue = value.getNewValue();
+ if (newValue != null) {
+ scalarWriter.write("new", newValue);
+ }
+ Object oldValue = value.getOldValue();
+ if (oldValue != null) {
+ scalarWriter.write("old", oldValue);
+ }
+ gen.writeEndObject();
+ }
+ }
+
+}
diff --git a/src/main/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogPrepare.java b/src/main/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogPrepare.java
new file mode 100644
index 000000000..6edcf6ca3
--- /dev/null
+++ b/src/main/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogPrepare.java
@@ -0,0 +1,23 @@
+package com.avaje.ebeaninternal.server.changelog;
+
+import com.avaje.ebean.event.changelog.ChangeLogPrepare;
+import com.avaje.ebean.event.changelog.ChangeSet;
+
+/**
+ * Placeholder/default implementation that does not do anything.
+ *
+ * Generally an implementation should be provided that reads context
+ * information such as user id and user ip address etc and sets that
+ * on the changeSet.
+ *
+ */
+public class DefaultChangeLogPrepare implements ChangeLogPrepare {
+
+ /**
+ * Just return true to send change set through to the logger.
+ */
+ @Override
+ public boolean prepare(ChangeSet changeSet) {
+ return true;
+ }
+}
diff --git a/src/main/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogRegister.java b/src/main/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogRegister.java
new file mode 100644
index 000000000..af7968708
--- /dev/null
+++ b/src/main/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogRegister.java
@@ -0,0 +1,88 @@
+package com.avaje.ebeaninternal.server.changelog;
+
+import com.avaje.ebean.annotation.ChangeLog;
+import com.avaje.ebean.event.BeanPersistRequest;
+import com.avaje.ebean.event.changelog.ChangeLogFilter;
+import com.avaje.ebean.event.changelog.ChangeLogRegister;
+
+import java.util.HashSet;
+import java.util.Set;
+
+/**
+ * Default implementation of ChangeLogRegister.
+ */
+public class DefaultChangeLogRegister implements ChangeLogRegister {
+
+ private static final BasicFilter INCLUDE_INSERTS = new BasicFilter(true);
+
+ private static final BasicFilter EXCLUDE_INSERTS = new BasicFilter(false);
+
+
+ @Override
+ public ChangeLogFilter getChangeFilter(Class> beanType) {
+
+ ChangeLog changeLog = beanType.getAnnotation(ChangeLog.class);
+ if (changeLog == null) {
+ return null;
+ }
+
+ String[] updatesThatInclude = changeLog.updatesThatInclude();
+ if (updatesThatInclude.length == 0) {
+ return changeLog.excludeInserts() ? EXCLUDE_INSERTS : INCLUDE_INSERTS;
+ }
+
+ Set updateProps = new HashSet();
+ for (int i = 0; i < updatesThatInclude.length; i++) {
+ updateProps.add(updatesThatInclude[i]);
+ }
+
+ return new UpdateFilter(!changeLog.excludeInserts(), updateProps);
+ }
+
+ /**
+ * Basic filter that only handles include inserts flag.
+ */
+ protected static class BasicFilter implements ChangeLogFilter {
+
+ final boolean includeInserts;
+
+ BasicFilter(boolean includeInserts) {
+ this.includeInserts = includeInserts;
+ }
+
+ @Override
+ public boolean includeInsert(BeanPersistRequest> insertRequest) {
+ return includeInserts;
+ }
+
+ @Override
+ public boolean includeUpdate(BeanPersistRequest> updateRequest) {
+ return true;
+ }
+
+ @Override
+ public boolean includeDelete(BeanPersistRequest> deleteRequest) {
+ return true;
+ }
+ }
+
+ /**
+ * Filter that takes into account a set of properties to check for updates
+ * as well as the include inserts flag.
+ */
+ protected static class UpdateFilter extends BasicFilter {
+
+ final Set updateProperties;
+
+ UpdateFilter(boolean includeInserts, Set updateProperties) {
+ super(includeInserts);
+ this.updateProperties = updateProperties;
+ }
+
+ @Override
+ public boolean includeUpdate(BeanPersistRequest> updateRequest) {
+ return updateRequest.hasDirtyProperty(updateProperties);
+ }
+
+ }
+}
diff --git a/src/main/java/com/avaje/ebeaninternal/server/changelog/ElasticChangeLogListener.java b/src/main/java/com/avaje/ebeaninternal/server/changelog/ElasticChangeLogListener.java
new file mode 100644
index 000000000..08537d29b
--- /dev/null
+++ b/src/main/java/com/avaje/ebeaninternal/server/changelog/ElasticChangeLogListener.java
@@ -0,0 +1,80 @@
+package com.avaje.ebeaninternal.server.changelog;
+
+import com.avaje.ebean.event.changelog.ChangeLogListener;
+import com.avaje.ebean.event.changelog.ChangeSet;
+import com.avaje.ebean.plugin.SpiServer;
+import com.avaje.ebean.plugin.SpiServerPlugin;
+import org.jetbrains.annotations.NotNull;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.io.StringWriter;
+
+/**
+ * Logs the change sets in a Elastic Search Bulk API JSON format.
+ */
+public class ElasticChangeLogListener implements ChangeLogListener, SpiServerPlugin {
+
+ protected static final Logger appLog = LoggerFactory.getLogger(ElasticChangeLogListener.class);
+
+ protected static final Logger fileLog = LoggerFactory.getLogger("org.avaje.ebean.ElasticChangeLog");
+
+ protected BulkJsonBuilder jsonBuilder;
+
+ public ElasticChangeLogListener() {
+ }
+
+ /**
+ * Configure the underlying JSON handler.
+ */
+ @Override
+ public void configure(SpiServer server) {
+ jsonBuilder = new BulkJsonBuilder(server.json(), "changelog", "changelog");
+ }
+
+ @Override
+ public void online(boolean online) {
+ // We don't care online or offline in this case
+ // we might if we setup for network sending etc
+ }
+
+ @Override
+ public void log(ChangeSet changeSet) {
+
+ try {
+ // I'm pretty sure I'm going to change this to use a FileWriter (but apache kafta could be good here too)
+ // This buffer could get really big and to me normal logging (without a Writer) is not that well suited to
+ // this problem so ... works but lets do better here and write direct to files (without the buffer issue)
+ StringWriter writer = new StringWriter(getBufferSize(changeSet));
+ jsonBuilder.writeJson(changeSet, writer);
+ String json = writer.toString();
+
+ fileLog.info("Sending txnId:{} txnState:{} txnBatch:{} \n {}", changeSet.getTxnId(), changeSet.getTxnState(), changeSet.getTxnBatch(), json);
+
+ } catch (IOException e) {
+ String msg = extractErrorMessage(e);
+ fileLog.error("Exception sending txnId:{} txnState:{} txnBatch:{} error:{}", changeSet.getTxnId(), changeSet.getTxnState(), changeSet.getTxnBatch(), msg);
+ appLog.error("Exception sending changeSet "+changeSet.toString(), e);
+ }
+ }
+
+ protected int getBufferSize(ChangeSet changeSet) {
+ // a rough guess, could get smarter here
+ return Math.min(400 * changeSet.size(), 3000);
+ }
+
+ /**
+ * Extract an error message that does not have new line characters and hence safe to go into
+ * our log which contains the payloads (that we want to be able to extract easily later on).
+ */
+ @NotNull
+ protected String extractErrorMessage(Exception e) {
+
+ String msg = e.toString();
+ msg = msg.replace('\r','|');
+ msg = msg.replace('\n', '|');
+ return msg;
+ }
+
+}
diff --git a/src/main/java/com/avaje/ebeaninternal/server/core/DefaultServer.java b/src/main/java/com/avaje/ebeaninternal/server/core/DefaultServer.java
index e37b51334..d081ec385 100644
--- a/src/main/java/com/avaje/ebeaninternal/server/core/DefaultServer.java
+++ b/src/main/java/com/avaje/ebeaninternal/server/core/DefaultServer.java
@@ -17,6 +17,7 @@ import com.avaje.ebean.event.BeanQueryAdapter;
import com.avaje.ebean.meta.MetaInfoManager;
import com.avaje.ebean.plugin.SpiBeanType;
import com.avaje.ebean.plugin.SpiServer;
+import com.avaje.ebean.plugin.SpiServerPlugin;
import com.avaje.ebean.text.csv.CsvReader;
import com.avaje.ebean.text.json.JsonContext;
import com.avaje.ebeaninternal.api.LoadBeanRequest;
@@ -135,6 +136,8 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer {
private final CQueryEngine cqueryEngine;
+ private final List serverPlugins;
+
private DdlGenerator ddlGenerator;
private final ExpressionFactory expressionFactory;
@@ -237,14 +240,24 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer {
this.beanLoader = new DefaultBeanLoader(this);
this.jsonContext = config.createJsonContext(this);
+ this.serverPlugins = config.getPlugins();
// load normal plugins late and call setup on all
loadAndInitializePlugins(config.getServerConfig());
+
+ configureServerPlugins();
// Register with the JVM Shutdown hook
ShutdownManager.registerEbeanServer(this);
}
+ private void configureServerPlugins() {
+
+ for (SpiServerPlugin plugin : serverPlugins) {
+ plugin.configure(this);
+ }
+ }
+
protected void loadAndInitializePlugins(ServerConfig config) {
List spiPlugins = new ArrayList();
@@ -283,6 +296,9 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer {
for (SpiEbeanPlugin plugin : ebeanPlugins) {
plugin.execute(online);
}
+ for (SpiServerPlugin plugin : serverPlugins) {
+ plugin.online(online);
+ }
}
@Override
diff --git a/src/main/java/com/avaje/ebeaninternal/server/core/InternalConfiguration.java b/src/main/java/com/avaje/ebeaninternal/server/core/InternalConfiguration.java
index 2ee44017b..23f562a6c 100644
--- a/src/main/java/com/avaje/ebeaninternal/server/core/InternalConfiguration.java
+++ b/src/main/java/com/avaje/ebeaninternal/server/core/InternalConfiguration.java
@@ -1,25 +1,27 @@
package com.avaje.ebeaninternal.server.core;
-import javax.sql.DataSource;
-
-import com.avaje.ebean.config.dbplatform.DbHistorySupport;
-import com.avaje.ebeaninternal.server.deploy.generatedproperty.GeneratedPropertyFactory;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
import com.avaje.ebean.ExpressionFactory;
import com.avaje.ebean.cache.ServerCacheManager;
import com.avaje.ebean.config.ExternalTransactionManager;
import com.avaje.ebean.config.ServerConfig;
import com.avaje.ebean.config.dbplatform.DatabasePlatform;
+import com.avaje.ebean.config.dbplatform.DbHistorySupport;
+import com.avaje.ebean.event.changelog.ChangeLogListener;
+import com.avaje.ebean.event.changelog.ChangeLogPrepare;
+import com.avaje.ebean.event.changelog.ChangeLogRegister;
+import com.avaje.ebean.plugin.SpiServerPlugin;
import com.avaje.ebean.text.json.JsonContext;
import com.avaje.ebeaninternal.api.SpiBackgroundExecutor;
import com.avaje.ebeaninternal.api.SpiEbeanServer;
import com.avaje.ebeaninternal.server.autofetch.AutoFetchManager;
import com.avaje.ebeaninternal.server.autofetch.AutoFetchManagerFactory;
+import com.avaje.ebeaninternal.server.changelog.DefaultChangeLogPrepare;
+import com.avaje.ebeaninternal.server.changelog.DefaultChangeLogRegister;
+import com.avaje.ebeaninternal.server.changelog.ElasticChangeLogListener;
import com.avaje.ebeaninternal.server.cluster.ClusterManager;
import com.avaje.ebeaninternal.server.deploy.BeanDescriptorManager;
import com.avaje.ebeaninternal.server.deploy.DeployOrmXml;
+import com.avaje.ebeaninternal.server.deploy.generatedproperty.GeneratedPropertyFactory;
import com.avaje.ebeaninternal.server.deploy.parse.DeployCreateProperties;
import com.avaje.ebeaninternal.server.deploy.parse.DeployInherit;
import com.avaje.ebeaninternal.server.deploy.parse.DeployUtil;
@@ -42,7 +44,12 @@ import com.avaje.ebeaninternal.server.transaction.TransactionScopeManager;
import com.avaje.ebeaninternal.server.type.DefaultTypeManager;
import com.avaje.ebeaninternal.server.type.TypeManager;
import com.fasterxml.jackson.core.JsonFactory;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import javax.sql.DataSource;
+import java.util.ArrayList;
+import java.util.List;
import java.util.Map;
/**
@@ -93,6 +100,11 @@ public class InternalConfiguration {
private final JsonFactory jsonFactory;
+ /**
+ * List of plugins (that ultimately the DefaultServer configures late in construction).
+ */
+ private final List plugins = new ArrayList();
+
public InternalConfiguration(XmlConfig xmlConfig, ClusterManager clusterManager,
ServerCacheManager cacheManager, SpiBackgroundExecutor backgroundExecutor,
ServerConfig serverConfig, BootupClasses bootupClasses) {
@@ -139,6 +151,45 @@ public class InternalConfiguration {
}
+ /**
+ * Check if this is a SpiServerPlugin and if so 'collect' it to give the complete list
+ * later on the DefaultServer for late call to configure().
+ */
+ public T plugin(T maybePlugin) {
+ if (maybePlugin instanceof SpiServerPlugin) {
+ plugins.add((SpiServerPlugin)maybePlugin);
+ }
+ return maybePlugin;
+ }
+
+ /**
+ * Return the list of plugins we collected during construction.
+ */
+ public List getPlugins() {
+ return plugins;
+ }
+
+ /**
+ * Return the ChangeLogPrepare to use with a default implementation if none defined.
+ */
+ public ChangeLogPrepare changeLogPrepare(ChangeLogPrepare prepare) {
+ return plugin((prepare != null) ? prepare : new DefaultChangeLogPrepare());
+ }
+
+ /**
+ * Return the ChangeLogRegister to use with a default implementation if none defined.
+ */
+ public ChangeLogRegister changeLogRegister(ChangeLogRegister register) {
+ return plugin((register != null) ? register : new DefaultChangeLogRegister());
+ }
+
+ /**
+ * Return the ChangeLogListener to use with a default implementation if none defined.
+ */
+ public ChangeLogListener changeLogListener(ChangeLogListener listener) {
+ return plugin((listener != null) ? listener : new ElasticChangeLogListener());
+ }
+
/**
* For 'As Of' queries return the number of bind variables per predicate.
*/
diff --git a/src/main/java/com/avaje/ebeaninternal/server/core/PersistRequestBean.java b/src/main/java/com/avaje/ebeaninternal/server/core/PersistRequestBean.java
index c1d675ed7..e53c735d3 100644
--- a/src/main/java/com/avaje/ebeaninternal/server/core/PersistRequestBean.java
+++ b/src/main/java/com/avaje/ebeaninternal/server/core/PersistRequestBean.java
@@ -215,6 +215,14 @@ public final class PersistRequestBean extends PersistRequest implements BeanP
return intercept.getDirtyPropertyNames();
}
+ /**
+ * Return true if any of the given property names are dirty.
+ */
+ @Override
+ public boolean hasDirtyProperty(Set propertyNames) {
+ return intercept.hasDirtyProperty(propertyNames);
+ }
+
@Override
public Map getUpdatedValues() {
return intercept.getDirtyValues();
diff --git a/src/main/java/com/avaje/ebeaninternal/server/deploy/BeanDescriptorManager.java b/src/main/java/com/avaje/ebeaninternal/server/deploy/BeanDescriptorManager.java
index 5b5536cc9..caca836f1 100644
--- a/src/main/java/com/avaje/ebeaninternal/server/deploy/BeanDescriptorManager.java
+++ b/src/main/java/com/avaje/ebeaninternal/server/deploy/BeanDescriptorManager.java
@@ -209,9 +209,9 @@ public class BeanDescriptorManager implements BeanDescriptorMap {
this.reflectFactory = createReflectionFactory();
this.transientProperties = new TransientProperties();
- this.changeLogPrepare = bootupClasses.getChangeLogPrepare();
- this.changeLogListener = bootupClasses.getChangeLogListener();
- this.changeLogRegister = bootupClasses.getChangeLogRegister();
+ this.changeLogPrepare = config.changeLogPrepare(bootupClasses.getChangeLogPrepare());
+ this.changeLogListener = config.changeLogListener(bootupClasses.getChangeLogListener());
+ this.changeLogRegister = config.changeLogRegister(bootupClasses.getChangeLogRegister());
}
/**
diff --git a/src/test/java/com/avaje/ebean/bean/EntityBeanInterceptTest.java b/src/test/java/com/avaje/ebean/bean/EntityBeanInterceptTest.java
new file mode 100644
index 000000000..4c6b775db
--- /dev/null
+++ b/src/test/java/com/avaje/ebean/bean/EntityBeanInterceptTest.java
@@ -0,0 +1,44 @@
+package com.avaje.ebean.bean;
+
+import com.avaje.ebean.BaseTestCase;
+import com.avaje.ebean.Ebean;
+import com.avaje.tests.model.basic.Customer;
+import com.avaje.tests.model.basic.ResetBasicData;
+import org.junit.Test;
+
+import java.sql.Date;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+public class EntityBeanInterceptTest extends BaseTestCase {
+
+ @Test
+ public void testHasDirtyProperty() throws Exception {
+
+
+ ResetBasicData.reset();
+
+ List list = Ebean.find(Customer.class).findList();
+
+ Set propertyNames = new HashSet();
+ propertyNames.add("name");
+ propertyNames.add("status");
+
+
+ Customer customer = list.get(0);
+ EntityBeanIntercept ebi = ((EntityBean)customer)._ebean_getIntercept();
+
+ assertFalse(ebi.hasDirtyProperty(propertyNames));
+
+ customer.setAnniversary(new Date(System.currentTimeMillis()));
+ assertFalse(ebi.hasDirtyProperty(propertyNames));
+
+ customer.setStatus(Customer.Status.ACTIVE);
+ assertTrue(ebi.hasDirtyProperty(propertyNames));
+
+ }
+}
\ No newline at end of file
diff --git a/src/test/java/com/avaje/ebean/event/BeanFindControllerTest.java b/src/test/java/com/avaje/ebean/event/BeanFindControllerTest.java
index e2bb32d5b..d966179e2 100644
--- a/src/test/java/com/avaje/ebean/event/BeanFindControllerTest.java
+++ b/src/test/java/com/avaje/ebean/event/BeanFindControllerTest.java
@@ -22,7 +22,7 @@ public class BeanFindControllerTest extends BaseTestCase {
ServerConfig config = new ServerConfig();
- config.setName("h2other");
+ config.setName("h2otherfind");
config.loadFromProperties();
config.setDdlGenerate(true);
config.setDdlRun(true);
diff --git a/src/test/java/com/avaje/ebeaninternal/server/changelog/BulkJsonBuilderTest.java b/src/test/java/com/avaje/ebeaninternal/server/changelog/BulkJsonBuilderTest.java
new file mode 100644
index 000000000..e38025376
--- /dev/null
+++ b/src/test/java/com/avaje/ebeaninternal/server/changelog/BulkJsonBuilderTest.java
@@ -0,0 +1,30 @@
+package com.avaje.ebeaninternal.server.changelog;
+
+import com.avaje.ebean.BaseTestCase;
+import com.avaje.ebean.Ebean;
+import com.avaje.ebean.text.json.JsonContext;
+import org.junit.Test;
+
+import java.io.StringWriter;
+
+
+public class BulkJsonBuilderTest extends BaseTestCase {
+
+ Helper helper = new Helper();
+
+ @Test
+ public void testToJson() throws Exception {
+
+ JsonContext jsonContext = Ebean.getDefaultServer().json();
+ BulkJsonBuilder builder = new BulkJsonBuilder(jsonContext, "changelog2", "changelog2");
+
+
+ StringWriter buffer = new StringWriter();
+ builder.writeJson(helper.createChangeSet("ABCD", 10), buffer);
+ System.out.println(buffer.toString());
+
+ buffer = new StringWriter();
+ builder.writeJson(helper.createChangeSet("ABCD-2", 15), buffer);
+ System.out.println(buffer.toString());
+ }
+}
\ No newline at end of file
diff --git a/src/test/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogRegisterTest.java b/src/test/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogRegisterTest.java
new file mode 100644
index 000000000..85a10bf72
--- /dev/null
+++ b/src/test/java/com/avaje/ebeaninternal/server/changelog/DefaultChangeLogRegisterTest.java
@@ -0,0 +1,40 @@
+package com.avaje.ebeaninternal.server.changelog;
+
+import com.avaje.ebean.BaseTestCase;
+import com.avaje.ebean.event.changelog.ChangeLogFilter;
+import com.avaje.tests.model.basic.Address;
+import com.avaje.tests.model.basic.Contact;
+import com.avaje.tests.model.basic.Country;
+import com.avaje.tests.model.basic.Customer;
+import org.junit.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertTrue;
+
+public class DefaultChangeLogRegisterTest extends BaseTestCase {
+
+ DefaultChangeLogRegister register = new DefaultChangeLogRegister();
+
+ @Test
+ public void test() {
+
+ assertNull(register.getChangeFilter(Address.class));
+
+ ChangeLogFilter changeFilter = register.getChangeFilter(Customer.class);
+ DefaultChangeLogRegister.UpdateFilter updateFilter = (DefaultChangeLogRegister.UpdateFilter)changeFilter;
+ assertFalse(updateFilter.includeInserts);
+ assertThat(updateFilter.updateProperties).containsExactly("name", "status");
+
+ changeFilter = register.getChangeFilter(Contact.class);
+ DefaultChangeLogRegister.BasicFilter contactFilter = (DefaultChangeLogRegister.BasicFilter)changeFilter;
+ assertTrue(contactFilter.includeInserts);
+
+ changeFilter = register.getChangeFilter(Country.class);
+ DefaultChangeLogRegister.BasicFilter countryFilter = (DefaultChangeLogRegister.BasicFilter)changeFilter;
+ assertFalse(countryFilter.includeInserts);
+
+ }
+
+}
diff --git a/src/test/java/com/avaje/ebeaninternal/server/changelog/ElasticChangeLogListenerTest.java b/src/test/java/com/avaje/ebeaninternal/server/changelog/ElasticChangeLogListenerTest.java
new file mode 100644
index 000000000..1ce5ba59f
--- /dev/null
+++ b/src/test/java/com/avaje/ebeaninternal/server/changelog/ElasticChangeLogListenerTest.java
@@ -0,0 +1,26 @@
+package com.avaje.ebeaninternal.server.changelog;
+
+import com.avaje.ebean.BaseTestCase;
+import com.avaje.ebean.Ebean;
+import com.avaje.ebean.EbeanServer;
+import com.avaje.ebean.event.changelog.ChangeSet;
+import org.junit.Test;
+
+public class ElasticChangeLogListenerTest extends BaseTestCase {
+
+ Helper helper = new Helper();
+
+ @Test
+ public void test() {
+
+ ElasticChangeLogListener changeLogListener = new ElasticChangeLogListener();
+
+ EbeanServer defaultServer = Ebean.getDefaultServer();
+ changeLogListener.configure(defaultServer.getPluginApi());
+
+ ChangeSet changeSet = helper.createChangeSet("INT-001", 13);
+
+ changeLogListener.log(changeSet);
+ }
+
+}
\ No newline at end of file
diff --git a/src/test/java/com/avaje/ebeaninternal/server/changelog/Helper.java b/src/test/java/com/avaje/ebeaninternal/server/changelog/Helper.java
new file mode 100644
index 000000000..280759bbe
--- /dev/null
+++ b/src/test/java/com/avaje/ebeaninternal/server/changelog/Helper.java
@@ -0,0 +1,68 @@
+package com.avaje.ebeaninternal.server.changelog;
+
+import com.avaje.ebean.ValuePair;
+import com.avaje.ebean.event.changelog.BeanChange;
+import com.avaje.ebean.event.changelog.ChangeSet;
+import com.avaje.ebean.event.changelog.ChangeType;
+import com.avaje.ebean.event.changelog.TxnState;
+import org.jetbrains.annotations.NotNull;
+
+import java.sql.Timestamp;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+public class Helper {
+
+ public ChangeSet createChangeSet(String txnId, long startId) {
+
+ ChangeSet cs = new ChangeSet();
+ cs.setTxnId(txnId);
+ cs.setTxnState(TxnState.COMMITTED);
+ cs.setTxnBatch(0);
+ cs.setSource("myApp");
+ cs.setUserId("user234");
+ cs.setUserIpAddress("123.4.5.6");
+ cs.setUserContext("user defined input");
+
+ List changes = cs.getChanges();
+
+ changes.add(createInsert(startId));
+ changes.add(createUpdate(startId));
+ changes.add(createDelete(startId));
+
+ return cs;
+ }
+
+ @NotNull
+ private BeanChange createInsert(long startId) {
+ Map values = new LinkedHashMap();
+ values.put("name", new ValuePair("rob", null));
+ values.put("modified", new ValuePair(new Timestamp(System.currentTimeMillis()), null));
+
+ BeanChange bean = new BeanChange("mytable", startId+1, ChangeType.INSERT, null);
+ bean.setValues(values);
+ return bean;
+ }
+
+ @NotNull
+ private BeanChange createUpdate(long startId) {
+ Map values = new LinkedHashMap();
+ values.put("name", new ValuePair("jim", "steve"));
+ values.put("nowHasVal", new ValuePair("wasNull", null));
+ values.put("nowNull", new ValuePair(null, "hadVal"));
+
+ values.put("modified", new ValuePair(new Timestamp(System.currentTimeMillis()), null));
+
+ BeanChange bean = new BeanChange("mytable", startId+2, ChangeType.UPDATE, null);
+ bean.setValues(values);
+ return bean;
+ }
+
+ @NotNull
+ private BeanChange createDelete(long startId) {
+ return new BeanChange("mytable", startId+3, ChangeType.DELETE, new HashMap());
+ }
+
+}
diff --git a/src/test/java/com/avaje/tests/model/basic/Contact.java b/src/test/java/com/avaje/tests/model/basic/Contact.java
index 01539010b..5227ca0eb 100644
--- a/src/test/java/com/avaje/tests/model/basic/Contact.java
+++ b/src/test/java/com/avaje/tests/model/basic/Contact.java
@@ -10,8 +10,10 @@ import javax.persistence.OneToMany;
import javax.persistence.Version;
import com.avaje.ebean.annotation.CacheStrategy;
+import com.avaje.ebean.annotation.ChangeLog;
import com.avaje.ebean.annotation.CreatedTimestamp;
+@ChangeLog
@Entity
@CacheStrategy(naturalKey="email")
public class Contact {
diff --git a/src/test/java/com/avaje/tests/model/basic/Country.java b/src/test/java/com/avaje/tests/model/basic/Country.java
index 5563f1267..499e54ab8 100644
--- a/src/test/java/com/avaje/tests/model/basic/Country.java
+++ b/src/test/java/com/avaje/tests/model/basic/Country.java
@@ -7,10 +7,12 @@ import javax.validation.constraints.Size;
import com.avaje.ebean.annotation.CacheStrategy;
import com.avaje.ebean.annotation.CacheTuning;
+import com.avaje.ebean.annotation.ChangeLog;
/**
* Country entity bean.
*/
+@ChangeLog(excludeInserts = true)
@CacheStrategy(readOnly=true,warmingQuery="order by name")
@CacheTuning(maxSize=500)
@Entity
diff --git a/src/test/java/com/avaje/tests/model/basic/Customer.java b/src/test/java/com/avaje/tests/model/basic/Customer.java
index 344df084d..12ef8d2d1 100644
--- a/src/test/java/com/avaje/tests/model/basic/Customer.java
+++ b/src/test/java/com/avaje/tests/model/basic/Customer.java
@@ -1,5 +1,6 @@
package com.avaje.tests.model.basic;
+import com.avaje.ebean.annotation.ChangeLog;
import com.avaje.ebean.annotation.EnumValue;
import com.avaje.ebean.annotation.JsonIgnore;
import com.avaje.ebean.annotation.Where;
@@ -20,6 +21,7 @@ import java.util.concurrent.locks.ReentrantLock;
/**
* Customer entity bean.
*/
+@ChangeLog(excludeInserts = true, updatesThatInclude = {"name","status"})
@Entity
@Table(name = "o_customer")
public class Customer extends BasicDomain {
diff --git a/src/test/java/com/avaje/tests/model/basic/EBasicChangeLog.java b/src/test/java/com/avaje/tests/model/basic/EBasicChangeLog.java
index 91ca2ff79..65190e3c2 100644
--- a/src/test/java/com/avaje/tests/model/basic/EBasicChangeLog.java
+++ b/src/test/java/com/avaje/tests/model/basic/EBasicChangeLog.java
@@ -15,7 +15,7 @@ import java.sql.Timestamp;
/**
*
*/
-@ChangeLog
+@ChangeLog(updatesThatInclude = {"name","shortDescription"})
@Entity
public class EBasicChangeLog {
diff --git a/src/test/java/com/avaje/tests/model/basic/Order.java b/src/test/java/com/avaje/tests/model/basic/Order.java
index feb1df336..56ed489a3 100644
--- a/src/test/java/com/avaje/tests/model/basic/Order.java
+++ b/src/test/java/com/avaje/tests/model/basic/Order.java
@@ -22,6 +22,7 @@ import javax.persistence.Version;
import javax.validation.constraints.NotNull;
import javax.xml.bind.annotation.XmlType;
+import com.avaje.ebean.annotation.ChangeLog;
import com.avaje.ebean.annotation.Formula;
import com.avaje.ebean.annotation.WhenCreated;
import com.avaje.ebean.annotation.Where;
@@ -29,6 +30,7 @@ import com.avaje.ebean.annotation.Where;
/**
* Order entity bean.
*/
+@ChangeLog
@Entity
@Table(name = "o_order")
public class Order implements Serializable {
diff --git a/src/test/resources/ebean.properties b/src/test/resources/ebean.properties
index 4d273cc87..86562c599 100644
--- a/src/test/resources/ebean.properties
+++ b/src/test/resources/ebean.properties
@@ -80,6 +80,11 @@ datasource.h2other.password=
datasource.h2other.databaseUrl=jdbc:h2:mem:h2other;DB_CLOSE_DELAY=-1
datasource.h2other.databaseDriver=org.h2.Driver
+datasource.h2otherfind.username=sa
+datasource.h2otherfind.password=
+datasource.h2otherfind.databaseUrl=jdbc:h2:mem:h2otherfind;DB_CLOSE_DELAY=-1
+datasource.h2otherfind.databaseDriver=org.h2.Driver
+
datasource.h2ebasicver.username=sa
datasource.h2ebasicver.password=
datasource.h2ebasicver.databaseUrl=jdbc:h2:mem:h2ebasicver;DB_CLOSE_DELAY=-1
diff --git a/src/test/resources/logback-test.xml b/src/test/resources/logback-test.xml
index 6fb0362e4..920c2eb92 100644
--- a/src/test/resources/logback-test.xml
+++ b/src/test/resources/logback-test.xml
@@ -9,7 +9,7 @@
- log/ebean.log
+ log/application.log
TRACE
@@ -22,6 +22,21 @@
+
+ log/changeLog.log
+
+ log/changeLog.log.%d{yyyy-MM-dd}
+ 90
+
+
+ -- %d{HH:mm:ss.SSS} %msg%n
+
+
+
+
+
+
+