diff --git a/build.xml b/build.xml
index 66dac1c..4ec5f3b 100644
--- a/build.xml
+++ b/build.xml
@@ -33,6 +33,7 @@
+
diff --git a/src/main/java/com/sumologic/client/CloudWatchLogsMessageModel.java b/src/main/java/com/sumologic/client/CloudWatchLogsMessageModel.java
new file mode 100644
index 0000000..75e7f23
--- /dev/null
+++ b/src/main/java/com/sumologic/client/CloudWatchLogsMessageModel.java
@@ -0,0 +1,119 @@
+package com.sumologic.client;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import org.apache.commons.lang.builder.ToStringBuilder;
+
+import com.sumologic.client.LogEvent;
+import com.fasterxml.jackson.annotation.JsonAnyGetter;
+import com.fasterxml.jackson.annotation.JsonAnySetter;
+import com.fasterxml.jackson.annotation.JsonIgnore;
+import com.fasterxml.jackson.annotation.JsonInclude;
+import com.fasterxml.jackson.annotation.JsonProperty;
+import com.fasterxml.jackson.annotation.JsonPropertyOrder;
+
+@JsonInclude(JsonInclude.Include.NON_NULL)
+@JsonPropertyOrder({
+"logEvents",
+"logGroup",
+"logStream",
+"messageType",
+"owner",
+"subscriptionFilters"
+})
+
+public class CloudWatchLogsMessageModel {
+
+ @JsonProperty("logEvents")
+ private List logEvents = new ArrayList();
+ @JsonProperty("logGroup")
+ private String logGroup;
+ @JsonProperty("logStream")
+ private String logStream;
+ @JsonProperty("messageType")
+ private String messageType;
+ @JsonProperty("owner")
+ private String owner;
+ @JsonProperty("subscriptionFilters")
+ private List subscriptionFilters = new ArrayList();
+ @JsonIgnore
+ private Map additionalProperties = new HashMap();
+
+ @JsonProperty("logEvents")
+ public List getLogEvents() {
+ return logEvents;
+ }
+
+ @JsonProperty("logEvents")
+ public void setLogEvents(List logEvents) {
+ this.logEvents = logEvents;
+ }
+
+ @JsonProperty("logGroup")
+ public String getLogGroup() {
+ return logGroup;
+ }
+
+ @JsonProperty("logGroup")
+ public void setLogGroup(String logGroup) {
+ this.logGroup = logGroup;
+ }
+
+ @JsonProperty("logStream")
+ public String getLogStream() {
+ return logStream;
+ }
+
+ @JsonProperty("logStream")
+ public void setLogStream(String logStream) {
+ this.logStream = logStream;
+ }
+
+ @JsonProperty("messageType")
+ public String getMessageType() {
+ return messageType;
+ }
+
+ @JsonProperty("messageType")
+ public void setMessageType(String messageType) {
+ this.messageType = messageType;
+ }
+
+ @JsonProperty("owner")
+ public String getOwner() {
+ return owner;
+ }
+
+ @JsonProperty("owner")
+ public void setOwner(String owner) {
+ this.owner = owner;
+ }
+
+ @JsonProperty("subscriptionFilters")
+ public List getSubscriptionFilters() {
+ return subscriptionFilters;
+ }
+
+ @JsonProperty("subscriptionFilters")
+ public void setSubscriptionFilters(List subscriptionFilters) {
+ this.subscriptionFilters = subscriptionFilters;
+ }
+
+ @JsonAnyGetter
+ public Map getAdditionalProperties() {
+ return this.additionalProperties;
+ }
+
+ @JsonAnySetter
+ public void setAdditionalProperty(String name, Object value) {
+ this.additionalProperties.put(name, value);
+ }
+
+ @Override
+ public String toString() {
+ return ToStringBuilder.reflectionToString(this);
+ }
+}
\ No newline at end of file
diff --git a/src/main/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformer.java b/src/main/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformer.java
index 78d9529..eb9072d 100644
--- a/src/main/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformer.java
+++ b/src/main/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformer.java
@@ -2,7 +2,11 @@
import java.io.IOException;
+import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
import com.amazonaws.services.kinesis.model.Record;
+import com.amazonaws.util.json.JSONArray;
+import com.amazonaws.util.json.JSONException;
+import com.amazonaws.util.json.JSONObject;
import com.sumologic.client.SimpleKinesisMessageModel;
import com.sumologic.client.implementations.SumologicTransformer;
@@ -13,18 +17,28 @@
import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.nio.ByteBuffer;
+import java.nio.CharBuffer;
+import java.nio.charset.CharacterCodingException;
import java.nio.charset.Charset;
import java.nio.charset.CharsetDecoder;
+import java.nio.charset.CharsetEncoder;
+import java.util.List;
import java.util.zip.GZIPInputStream;
+import com.fasterxml.jackson.databind.JsonMappingException;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.google.gson.Gson;
+
/**
- * A custom transfomer for {@link SimpleKinesisMessageModel} records in JSON format. The output is in a format
+ * A custom transfomer for {@link CloudWatchLogsMessageModel} records in JSON format. The output is in a format
* usable for insertions to Sumologic.
*/
-public class CloudWatchMessageModelSumologicTransformer implements
- SumologicTransformer {
+public class CloudWatchMessageModelSumologicTransformer
+ implements SumologicTransformer {
private static final Log LOG = LogFactory.getLog(CloudWatchMessageModelSumologicTransformer.class);
+
+ private static CharsetEncoder encoder = Charset.forName("UTF-8").newEncoder();
/**
* Creates a new KinesisMessageModelSumologicTransformer.
@@ -34,24 +48,70 @@ public CloudWatchMessageModelSumologicTransformer() {
}
@Override
- public String fromClass(SimpleKinesisMessageModel message) {
- return message.toString();
- }
+ public String fromClass(CloudWatchLogsMessageModel message) {
+ String jsonMessage = "";
+ JSONObject outputObject;
+
+ List logEvents = message.getLogEvents();
+ int logEventsSize = logEvents.size();
+ for (int i=0;i additionalProperties = new HashMap();
+
+ @JsonProperty("id")
+ public String getId() {
+ return id;
+ }
+
+ @JsonProperty("id")
+ public void setId(String id) {
+ this.id = id;
+ }
+
+ @JsonProperty("message")
+ public String getMessage() {
+ return message;
+ }
+
+ @JsonProperty("message")
+ public void setMessage(String message) {
+ this.message = message;
+ }
+
+ @JsonProperty("timestamp")
+ public Long getTimestamp() {
+ return timestamp;
+ }
+
+ @JsonProperty("timestamp")
+ public void setTimestamp(Long timestamp) {
+ this.timestamp = timestamp;
+ }
+
+ @JsonAnyGetter
+ public Map getAdditionalProperties() {
+ return this.additionalProperties;
+ }
+
+ @JsonAnySetter
+ public void setAdditionalProperty(String name, Object value) {
+ this.additionalProperties.put(name, value);
+ }
+
+}
\ No newline at end of file
diff --git a/src/main/java/com/sumologic/client/SumologicSender.java b/src/main/java/com/sumologic/client/SumologicSender.java
index 295234e..c0ee0c8 100644
--- a/src/main/java/com/sumologic/client/SumologicSender.java
+++ b/src/main/java/com/sumologic/client/SumologicSender.java
@@ -48,7 +48,7 @@ public boolean sendToSumologic(String data) throws IOException{
BoundRequestBuilder builder = null;
builder = this.clientPreparePost(url);
- byte[] compressedData = compressGzip(data);
+ byte[] compressedData = SumologicSender.compressGzip(data);
builder.setHeader("Content-Encoding", "gzip");
builder.setBody(compressedData);
@@ -82,7 +82,7 @@ public boolean sendToSumologic(String data) throws IOException{
}
}
- public byte[] compressGzip(String data) {
+ public static byte[] compressGzip(String data) {
if (data == null || data.length() == 0) {
return null;
}
diff --git a/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java b/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java
index 7f31a29..1588a77 100644
--- a/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java
+++ b/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java
@@ -22,7 +22,7 @@ public class CloudWatchMessageModelSumologicTransformerTest {
public static CharsetEncoder encoder = charset.newEncoder();
@Test
- public void theTransformerShouldFailGracefullyWhenUnableToTransform () {
+ public void theTransformerShouldFailGracefullyWhenUnableToCompress () {
CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer();
String randomData = "Some random string without GZIP compression";
@@ -36,13 +36,128 @@ public void theTransformerShouldFailGracefullyWhenUnableToTransform () {
Record mockedRecord = new Record();
mockedRecord.setData(bufferedData);
- SimpleKinesisMessageModel messageModel = null;
+ CloudWatchLogsMessageModel messageModel = transfomer.toClass(mockedRecord);
+
+
+ Assert.assertNull(messageModel);
+ }
+
+ @Test
+ public void theTransformerShouldSucceedWhenTransformingAProperJSON() {
+ CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer();
+
+ String jsonData = ""
+ +"{"
+ + "\"logEvents\": [{"
+ + "\"id\": \"3889492387492837492374982374897239847289374892\","
+ + "\"message\": \"1 23423532532 eni-ac9342k3492 10.1.1.75 66.175.209.17 123 123 17 1 76 1437755534 1437755549 ACCEPT OK\","
+ + "\"timestamp\": \"2342342342300\""
+ + "}],"
+ + "\"logGroup\": \"MyFirstVPC\","
+ + "\"logStream\": \"eni-ac6a7de4-all\","
+ + "\"messageType\": \"DATA_MESSAGE\","
+ + "\"owner\": \"2342352352\","
+ + "\"subscriptionFilters\": [\"MyFirstVPC\"]"
+ + "}"
+ +"";
+
+ byte[] compressData = SumologicSender.compressGzip(jsonData);
+
+ ByteBuffer bufferedData = null;
+ try {
+ bufferedData = ByteBuffer.wrap(compressData);
+ } catch (Exception e) {
+ Assert.fail("Getting error: "+e.getMessage());
+ }
+
+ Record mockedRecord = new Record();
+ mockedRecord.setData(bufferedData);
+
+ CloudWatchLogsMessageModel messageModel = transfomer.toClass(mockedRecord);
+
+ Assert.assertNotNull(messageModel);
+ }
+
+ @Test
+ public void theTransformerShouldFailWhenTransformingAJSONWithTrailingCommas() {
+ CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer();
+
+ String jsonData = ""
+ +"{"
+ + "\"logEvents\": [{"
+ + "\"id\": \"3889492387492837492374982374897239847289374892\","
+ + "\"message\": \"1 23423532532 eni-ac9342k3492 10.1.1.75 66.175.209.17 123 123 17 1 76 1437755534 1437755549 ACCEPT OK\","
+ + "\"timestamp\": \"2342342342300\""
+ + "}],"
+ + "\"logGroup\": \"MyFirstVPC\","
+ + "\"logStream\": \"eni-ac6a7de4-all\","
+ + "\"messageType\": \"DATA_MESSAGE\","
+ + "\"owner\": \"2342352352\","
+ + "\"subscriptionFilters\": [\"MyFirstVPC\"],"
+ + "}"
+ +"";
+
+ byte[] compressData = SumologicSender.compressGzip(jsonData);
+
+ ByteBuffer bufferedData = null;
try {
- messageModel = transfomer.toClass(mockedRecord);
- } catch (IOException e) {
- Assert.fail("Getting error while transforming: "+e.getMessage());
+ bufferedData = ByteBuffer.wrap(compressData);
+ } catch (Exception e) {
+ Assert.fail("Getting error: "+e.getMessage());
}
+ Record mockedRecord = new Record();
+ mockedRecord.setData(bufferedData);
+
+ CloudWatchLogsMessageModel messageModel = null;
+ messageModel = transfomer.toClass(mockedRecord);
+
Assert.assertNull(messageModel);
}
+
+ @Test
+ public void theTransfomerShouldSeparateBatchesOfLogs() {
+ CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer();
+
+ String jsonData = ""
+ +"{"
+ + "\"logEvents\": [{"
+ + "\"id\": \"3889492387492837492374982374897239847289374892\","
+ + "\"message\": \"1 23423532532 eni-ac9342k3492 10.1.1.75 66.175.209.17 123 123 17 1 76 1437755534 1437755549 ACCEPT OK\","
+ + "\"timestamp\": \"2342342342300\""
+ + "},"
+ + "{"
+ + "\"id\": \"3289429357928375892739857238975235235235\","
+ + "\"message\": \"1 23423516 eni-ac9342k3492 10.1.1.75 66.175.209.17 123 123 17 1 76 1437755534 1437755549 REJECT OK\","
+ + "\"timestamp\": \"2342352351616\""
+ + "}],"
+ + "\"logGroup\": \"MyFirstVPC\","
+ + "\"logStream\": \"eni-ac6a7de4-all\","
+ + "\"messageType\": \"DATA_MESSAGE\","
+ + "\"owner\": \"2342352352\","
+ + "\"subscriptionFilters\": [\"MyFirstVPC\"]"
+ + "}"
+ +"";
+
+ byte[] compressData = SumologicSender.compressGzip(jsonData);
+
+ ByteBuffer bufferedData = null;
+ try {
+ bufferedData = ByteBuffer.wrap(compressData);
+ } catch (Exception e) {
+ Assert.fail("Getting error: "+e.getMessage());
+ }
+
+ Record mockedRecord = new Record();
+ mockedRecord.setData(bufferedData);
+
+ CloudWatchLogsMessageModel messageModel = null;
+ messageModel = transfomer.toClass(mockedRecord);
+
+ String debatchedMessage = transfomer.fromClass(messageModel);
+ System.out.println(debatchedMessage);
+
+ String[] messages = debatchedMessage.split("\n");
+ Assert.assertTrue(messages.length == 2);
+ }
}
\ No newline at end of file
diff --git a/src/test/java/com/sumologic/client/SumologicSenderTest.java b/src/test/java/com/sumologic/client/SumologicSenderTest.java
index 8961ec4..1d9c54d 100644
--- a/src/test/java/com/sumologic/client/SumologicSenderTest.java
+++ b/src/test/java/com/sumologic/client/SumologicSenderTest.java
@@ -74,9 +74,7 @@ public void decompressGzipTest() {
String data = "a string of characters";
- SumologicSender sender = new SumologicSender(url);
-
- byte[] compressData = sender.compressGzip(data);
+ byte[] compressData = SumologicSender.compressGzip(data);
String result = CloudWatchMessageModelSumologicTransformer.decompressGzip(compressData);
Assert.assertTrue(data.equals(result));