package com.dexels.replication.impl; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.Date; import java.util.HashMap; import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Optional; import java.util.Set; import java.util.function.Predicate; import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.dexels.immutable.api.ImmutableMessage; import com.dexels.immutable.api.ImmutableMessage.Trifunction; import com.dexels.immutable.factory.ImmutableFactory; import com.dexels.replication.api.ReplicationMessage; import com.dexels.replication.api.ReplicationMessageParser; public class ReplicationImmutableMessageImpl implements ReplicationMessage { private final static Logger logger = LoggerFactory.getLogger(ReplicationImmutableMessageImpl.class); private final Optional source; private final String transactionId; private final long timestamp; private final Operation operation; private final List primaryKeys; private Optional commitAction; private final Optional partition; private final Optional offset; private final ImmutableMessage immutableMessage; private final Optional paramMessage; public ReplicationImmutableMessageImpl(Optional source, Optional partition, Optional offset, String transactionId, Operation operation, long timestamp, Map values, Map types,Map submessage, Map> submessages, List primaryKeys, Optional commitAction, Optional paramMessage) { this.immutableMessage = ImmutableFactory.create(values, types, submessage, submessages); this.transactionId = transactionId; this.timestamp = timestamp; this.operation = operation; this.primaryKeys = Collections.unmodifiableList(primaryKeys); this.commitAction = commitAction; this.source = source; this.partition = partition; this.offset = offset; this.paramMessage = paramMessage; } public ReplicationImmutableMessageImpl(Optional source, Optional partition, Optional offset, String transactionId, Operation operation, long timestamp, ImmutableMessage parentMessage , List primaryKeys, Optional commitAction, Optional paramMessage) { this.immutableMessage = parentMessage; this.transactionId = transactionId; this.timestamp = timestamp; this.operation = operation; this.primaryKeys = Collections.unmodifiableList(primaryKeys); this.commitAction = commitAction; this.source = source; this.offset = offset; this.partition = partition; this.paramMessage = paramMessage; } public ReplicationImmutableMessageImpl(ReplicationMessage message1, ReplicationMessage message2, String key) { this.transactionId = null; // message1.transactionId()+"-"+message2.transactionId(); this.timestamp = System.currentTimeMillis(); this.operation = Operation.MERGE; this.primaryKeys = Collections.unmodifiableList(Arrays.asList(new String[]{key})); this.immutableMessage = message1.message().merge(message2.message(), Optional.empty()); this.source = Optional.empty(); this.partition = Optional.empty(); this.offset = Optional.empty(); this.paramMessage = Optional.empty(); } @Override public ImmutableMessage message() { return immutableMessage; } @Override public Set subMessageNames() { return message().subMessageNames(); } @Override public Set subMessageListNames() { return message().subMessageListNames(); } @Override public String queueKey() { if(primaryKeys.size()==0) { return "NO_KEY_PRESENT"; } return String.join(ReplicationMessage.KEYSEPARATOR, primaryKeys.stream().map(col->""+columnValue(col)).collect(Collectors.toList())); } @Override public byte[] toBytes(ReplicationMessageParser c) { return c.serialize(this); } public boolean equals(Object other) { if(!(other instanceof ReplicationMessage)) { return false; } return equalsByKey((ReplicationMessage)other); } public boolean equalsByKey(ReplicationMessage other) { String key = queueKey(); if(key==null) { return super.equals(other); } return key.equals(other.queueKey()); // return key.equals(((ReplicationImmutableMessageImpl)other).queueKey()); } public int hashCode() { String key = queueKey(); if(key==null) { return super.hashCode(); } return key.hashCode(); } public ReplicationImmutableMessageImpl(Throwable t) { logger.error("Creating failure replication message",t); t.printStackTrace(System.err); t.printStackTrace(System.out); this.transactionId = null; this.timestamp = -1; // this.status = null; this.operation = Operation.UPDATE; this.primaryKeys = Collections.emptyList(); this.immutableMessage = ImmutableFactory.create(Collections.emptyMap(), Collections.emptyMap(), Collections.emptyMap(), Collections.emptyMap()); this.source = Optional.empty(); this.partition = Optional.empty(); this.offset = Optional.empty(); this.paramMessage = Optional.empty(); } @SuppressWarnings("unchecked") public ReplicationImmutableMessageImpl(Map initial) { this.transactionId = (String) initial.get("TransactionId"); this.timestamp = (long) initial.get("Timestamp"); this.operation = Operation.valueOf((String) initial.get("Operation")); this.primaryKeys = Collections.unmodifiableList(((List) initial.get("PrimaryKeys"))); Map initialValues = Collections.unmodifiableMap((Map) initial.get("Columns")); this.immutableMessage = ImmutableFactory.create(initialValues, resolveTypesFromValues(initialValues) , Collections.emptyMap(), Collections.emptyMap()); this.source = Optional.empty(); this.partition = Optional.empty(); this.offset = Optional.empty(); this.paramMessage = Optional.empty(); } private Map resolveTypesFromValues(Map values) { Map t = new HashMap<>(); for (Entry e : values.entrySet()) { Object val = e.getValue(); if(val==null) { // t.put(e.getKey(), "unknown"); } else if(val instanceof Long) { t.put(e.getKey(), "long"); } else if(val instanceof Double) { t.put(e.getKey(), "double"); } else if(val instanceof Integer) { t.put(e.getKey(), "integer"); } else if(val instanceof Float) { t.put(e.getKey(), "float"); } else if(val instanceof Date) { t.put(e.getKey(), "date"); } else if(val instanceof Boolean) { t.put(e.getKey(), "boolean"); } else if(val instanceof String) { t.put(e.getKey(), "string"); } else { logger.warn("Unknown type::: {}",val.getClass()); t.put(e.getKey(), "string"); } } return t; } // // private Map combineTypes(Map typesa, Map typesb) { // HashMap combine = new HashMap<>(typesa); // for (Entry e : typesb.entrySet()) { // // TODO sanity check types? // combine.put(e.getKey(), e.getValue()); // } // return Collections.unmodifiableMap(combine); // } @Override public Map types() { return message().types(); } @Override public Map> toDataMap() { return message().toDataMap(); } @Override public Map valueMap(boolean ignoreNull, Set ignore) { return valueMap(ignoreNull,ignore,Collections.emptyList()); } // // private static BiFunction,Boolean> checkIgnoreList(Set ignoreList) { // return (item,currentPath)->{ // if(ignoreList.contains(item)) { // return false; // } // return true; // }; // } @Override public Map valueMap(boolean ignoreNull, Set ignore, List currentPath) { return message().valueMap(ignoreNull, ignore,currentPath); } @Override public boolean isErrorMessage() { return transactionId==null; } @Override public String transactionId() { return transactionId; } @Override public long timestamp() { return timestamp; } @Override public Operation operation() { return operation; } @Override public List primaryKeys() { return this.primaryKeys; } @Override public Set columnNames() { return message().columnNames(); } @Override public Object columnValue(String name) { return message().columnValue(name); } @Override public String columnType(String name) { return message().columnType(name); } @Override public String toString() { return "Operation: " + this.operation.toString() + " Ts: " + this.timestamp + "Transactionid: "+this.transactionId+" pk: "+primaryKeys+"Value:\n" +message().toString(); } @Override public void commit() { if(this.commitAction!=null && this.commitAction.isPresent()) { this.commitAction.get().run(); } } @Override public Optional> subMessages(String field) { return message().subMessages(field); } @Override public Optional subMessage(String field) { return message().subMessage(field); } @Override public ReplicationMessage withImmutableMessage(ImmutableMessage msg) { return new ReplicationImmutableMessageImpl(this.source, this.partition,this.offset, this.transactionId,this.operation,this.timestamp, msg , primaryKeys,this.commitAction,this.paramMessage); } @Override public ReplicationMessage withSubMessages(String field, List message) { return withImmutableMessage(message().withSubMessages(field, message)); } @Override public ReplicationMessage withSubMessage(String field, ImmutableMessage message) { return withImmutableMessage(message().withSubMessage(field, message)); } @Override public ReplicationMessage without(String columnName) { List prim = new LinkedList<>(primaryKeys()); prim.remove(columnName); return withImmutableMessage(message().without(columnName)).withPrimaryKeys(Collections.unmodifiableList(prim)); } @Override public ReplicationMessage without(List columns) { List prim = new LinkedList<>(primaryKeys()); prim.removeAll(columns); return withImmutableMessage(message().without(columns)).withPrimaryKeys(Collections.unmodifiableList(prim)); } @Override public ReplicationMessage rename(String columnName, String newName) { int keyIndex = this.primaryKeys.indexOf(columnName); List primary = this.primaryKeys; if(keyIndex!=-1) { primary = new ArrayList<>(this.primaryKeys); primary.set(keyIndex, newName); } return withImmutableMessage(message().rename(columnName,newName)).withPrimaryKeys(primary); } @Override public ReplicationMessage with(String key, Object value, String type) { return withImmutableMessage(message().with(key, value, type)); } @Override public ReplicationMessage withPrimaryKeys(List primary) { return new ReplicationImmutableMessageImpl(this.source, this.partition,this.offset, this.transactionId,this.operation,this.timestamp, message(),primary,this.commitAction,this.paramMessage); } @Override public String toFlatString(ReplicationMessageParser parser) { return parser.describe(this); } @Override public ReplicationMessage withoutSubMessages(String field) { return withImmutableMessage(message().withoutSubMessages(field)); } @Override public ReplicationMessage withoutSubMessage(String field) { return withImmutableMessage(message().withoutSubMessage(field)); } @Override public Map subMessageMap() { return message().subMessageMap(); } @Override public Map> subMessageListMap() { return message().subMessageListMap(); } @Override public ReplicationMessage merge(ReplicationMessage other, Optional> only) { return withImmutableMessage(message().merge(other.message(), only)); } @Override public ReplicationMessage withOnlySubMessages(List subMessages) { return withImmutableMessage(message().withOnlySubMessages(subMessages)); } @Override public ReplicationMessage withOnlyColumns(List columns) { return withImmutableMessage(message().withOnlyColumns(columns)); } public ReplicationMessage withAllSubMessageLists(Map> subMessageListMap) { return withImmutableMessage(message().withAllSubMessageLists(subMessageListMap)); } public ReplicationMessage withAllSubMessage(Map subMessageMap) { return withImmutableMessage(message().withAllSubMessage(subMessageMap)); } @Override public ReplicationMessage withAddedSubMessage(String field, ImmutableMessage message) { return withImmutableMessage(message().withAddedSubMessage(field,message)); } @Override public ReplicationMessage withoutSubMessageInList(String field,Predicate selector) { return withImmutableMessage(message().withoutSubMessageInList(field,selector)); } @Override public ReplicationMessage withOperation(Operation operation) { return new ReplicationImmutableMessageImpl(this.source, this.partition,this.offset,this.transactionId,operation, this.timestamp, message(),primaryKeys,this.commitAction,this.paramMessage); } @Override public Map flatValueMap(boolean ignoreNull, Set ignore, String prefix) { return message().flatValueMap(ignoreNull, ignore, prefix); } public Map flatValueMap(String prefix, Trifunction processType) { return message().flatValueMap(prefix, processType); } public boolean equalsToMessage(ReplicationMessage c) { Map other = c.flatValueMap(false,Collections.emptySet(), ""); final Map myMap = this.flatValueMap(false, Collections.emptySet(), ""); return myMap.equals(other); } @Override public ReplicationMessage now() { return atTime(new Date().getTime()); } @Override public ReplicationMessage atTime(long timestamp) { return new ReplicationImmutableMessageImpl(this.source, this.partition,this.offset,this.transactionId,this.operation, timestamp,this.message(), this.primaryKeys,commitAction,this.paramMessage); } @Override public Map values() { return message().values(); } public ReplicationMessage withCommitAction(Runnable commitAction) { return new ReplicationImmutableMessageImpl(this.source,this.partition,this.offset, this.transactionId, this.operation, timestamp, message() , this.primaryKeys,Optional.of(commitAction),this.paramMessage); } @Override public Optional source() { return this.source; } @Override public ReplicationMessage withSource(Optional source) { return new ReplicationImmutableMessageImpl(source,this.partition,this.offset, this.transactionId, this.operation, timestamp, message() , this.primaryKeys,this.commitAction,this.paramMessage); } @Override public Optional partition() { return this.partition; } @Override public Optional offset() { return this.offset; } @Override public ReplicationMessage withPartition(Optional partition) { return new ReplicationImmutableMessageImpl(source, partition, this.offset, this.transactionId, this.operation, timestamp, message() , this.primaryKeys,this.commitAction,this.paramMessage); // public ReplicationImmutableMessageImpl(Optional source, Optional partition, Optional offset, String transactionId, Operation operation, long timestamp, Map values, Map types,Map submessage, Map> submessages, List primaryKeys, Optional commitAction) { } @Override public ReplicationMessage withOffset(Optional offset) { return new ReplicationImmutableMessageImpl(source,this.partition,offset, this.transactionId, this.operation, timestamp, message() , this.primaryKeys,this.commitAction,this.paramMessage); } public Optional paramMessage() { return this.paramMessage; } }