package com.dexels.replication.api; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; import java.util.function.Predicate; import java.util.stream.Collectors; import org.osgi.annotation.versioning.ProviderType; import com.dexels.immutable.api.ImmutableMessage; @ProviderType public interface ReplicationMessage { public static final String KEYSEPARATOR = "<$>"; public static final String PRETTY_JSON = "PRETTY_JSON"; public String transactionId(); public Optional source(); public Optional partition(); public Optional offset(); public long timestamp(); public Operation operation(); public List primaryKeys(); public Set columnNames(); public Object columnValue(String name); public String columnType(String name); public boolean equals(Object o); public enum Operation { INSERT, UPDATE, DELETE, NONE, COMMIT, MERGE, INITIAL } public String queueKey(); public void commit(); boolean isErrorMessage(); public Map> toDataMap(); public Map valueMap(boolean ignoreNull, Set ignore); public Map valueMap(boolean ignoreNull, Set ignore, List currentPath); public Map flatValueMap(boolean ignoreNull, Set ignore, String prefix); // public Map flatValueMap(String prefix,Func3 processType); public boolean equalsToMessage(ReplicationMessage c); public boolean equalsByKey(ReplicationMessage c); public byte[] toBytes(ReplicationMessageParser c); public Map types(); public Optional> subMessages(String field); public Optional subMessage(String field); public ReplicationMessage withImmutableMessage(ImmutableMessage msg); public ReplicationMessage withSubMessages(String field, List message); public ReplicationMessage withSubMessage(String field, ImmutableMessage message); public ReplicationMessage withAddedSubMessage(String field, ImmutableMessage message); public ReplicationMessage withoutSubMessageInList(String field, Predicate s); public ReplicationMessage withoutSubMessages(String field); public ReplicationMessage withoutSubMessage(String field); public Set subMessageListNames(); public Set subMessageNames(); public ReplicationMessage without(String columnName); public ReplicationMessage without(List columns); public ReplicationMessage with(String key, Object value, String type); public ReplicationMessage withOnlyColumns(List columns); public ReplicationMessage withOnlySubMessages(List subMessages); public ReplicationMessage rename(String columnName, String newName); public ReplicationMessage withPrimaryKeys(List primary); public ReplicationMessage withSource(Optional primary); public ReplicationMessage withPartition(Optional partition); public ReplicationMessage withOffset(Optional offset); public ReplicationMessage now(); public ReplicationMessage atTime(long timestamp); public String toFlatString(ReplicationMessageParser parser); public ReplicationMessage merge(ReplicationMessage other, Optional> only); public static final boolean usePretty = System.getenv(PRETTY_JSON)!=null || System.getProperty(PRETTY_JSON)!=null; public static boolean usePrettyPrint() { return usePretty; } public Map subMessageMap(); public Map> subMessageListMap(); public ReplicationMessage withAllSubMessageLists(Map> subMessageListMap); public ReplicationMessage withAllSubMessage(Map subMessageMap); public ReplicationMessage withOperation(Operation operation); public Map values(); public ReplicationMessage withCommitAction(Runnable commitAction); public ImmutableMessage message(); public Optional paramMessage(); default public String combinedKey() { return primaryKeys().stream().map(k->columnValue(k).toString()).collect(Collectors.joining(KEYSEPARATOR)); } }