From e76f1cde646cc646bec6b2395856a90e1d1e4a19 Mon Sep 17 00:00:00 2001 From: Juan Pablo Diaz-Vaz Date: Fri, 24 Jul 2015 12:36:08 -0300 Subject: [PATCH 1/3] SUMOK-28 Added verification of JSON in transformer for Cloudwatch --- build.xml | 1 + ...WatchMessageModelSumologicTransformer.java | 22 +++++- .../com/sumologic/client/SumologicSender.java | 4 +- ...hMessageModelSumologicTransformerTest.java | 74 +++++++++++++++++++ .../sumologic/client/SumologicSenderTest.java | 4 +- 5 files changed, 98 insertions(+), 7 deletions(-) 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/CloudWatchMessageModelSumologicTransformer.java b/src/main/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformer.java index 78d9529..0370e07 100644 --- a/src/main/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformer.java +++ b/src/main/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformer.java @@ -17,6 +17,8 @@ import java.nio.charset.CharsetDecoder; import java.util.zip.GZIPInputStream; +import com.google.gson.Gson; + /** * A custom transfomer for {@link SimpleKinesisMessageModel} records in JSON format. The output is in a format * usable for insertions to Sumologic. @@ -44,8 +46,14 @@ public SimpleKinesisMessageModel toClass(Record record) throws IOException { String stringifiedRecord = decompressGzip(decodedRecord); if (stringifiedRecord == null) { - LOG.error("Unable to decompress the record: "+new String(record.getData().array())); - LOG.error("Not attempting to transform into a Message Model"); + LOG.error("Unable to decompress the record: "+new String(record.getData().array()) + +"\nNot attempting to transform into a Message Model"); + return null; + } + + if (!verifyJSON(stringifiedRecord)) { + LOG.error("The record is not a valid JSON: "+stringifiedRecord + +"\nNot attempting to transform into a Message Model"); return null; } @@ -83,4 +91,14 @@ public static String byteBufferToString(ByteBuffer buffer){ return data; } + private static final Gson gson = new Gson(); + public static boolean verifyJSON(String json) { + try { + gson.fromJson(json, Object.class); + return true; + } catch(com.google.gson.JsonSyntaxException ex) { + return false; + } + } + } 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..9d9f087 100644 --- a/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java +++ b/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java @@ -21,6 +21,7 @@ public class CloudWatchMessageModelSumologicTransformerTest { public static Charset charset = Charset.forName("UTF-8"); public static CharsetEncoder encoder = charset.newEncoder(); + @Ignore @Test public void theTransformerShouldFailGracefullyWhenUnableToTransform () { CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer(); @@ -45,4 +46,77 @@ public void theTransformerShouldFailGracefullyWhenUnableToTransform () { Assert.assertNull(messageModel); } + + @Test + public void theTransformerShouldSucceedWhenTransformingAProperJSON() { + CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer(); + + String jsonData = "[" + +"{" + + "\"id\": \"55b25585730b5e5bcd53c580\"," + + "\"index\": 0," + + "\"guid\": \"f3bf1c45-306d-4799-801f-c6d16acca931\"," + + "\"isActive\": false," + + "\"picture\": \"http://placehold.it/32x32\"" + + "}" + +"]"; + + 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); + + SimpleKinesisMessageModel messageModel = null; + try { + messageModel = transfomer.toClass(mockedRecord); + } catch (IOException e) { + Assert.fail("Getting error while transforming: "+e.getMessage()); + } + + Assert.assertNotNull(messageModel); + Assert.assertTrue(messageModel.getData().equals(jsonData)); + } + + @Test + public void theTransformerShouldSucceedWhenTransformingAJSONWithTrailingCommas() { + CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer(); + + String jsonData = "[" + +"{" + + "\"id\": \"55b25585730b5e5bcd53c580\"," + + "\"index\": 0," + + "\"guid\": \"f3bf1c45-306d-4799-801f-c6d16acca931\"," + + "\"isActive\": false," + + "\"picture\": \"http://placehold.it/32x32\"," + + "}" + +"]"; + + 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); + + SimpleKinesisMessageModel messageModel = null; + try { + messageModel = transfomer.toClass(mockedRecord); + } catch (IOException e) { + Assert.fail("Getting error while transforming: "+e.getMessage()); + } + + Assert.assertNull(messageModel); + } } \ 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)); From 8b4c87d36076a407ebc138f39685859ba7415325 Mon Sep 17 00:00:00 2001 From: Juan Pablo Diaz-Vaz Date: Fri, 24 Jul 2015 12:38:36 -0300 Subject: [PATCH 2/3] Readded test ignored --- .../client/CloudWatchMessageModelSumologicTransformerTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java b/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java index 9d9f087..410549a 100644 --- a/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java +++ b/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java @@ -21,7 +21,6 @@ public class CloudWatchMessageModelSumologicTransformerTest { public static Charset charset = Charset.forName("UTF-8"); public static CharsetEncoder encoder = charset.newEncoder(); - @Ignore @Test public void theTransformerShouldFailGracefullyWhenUnableToTransform () { CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer(); From af18b587ef7ad8c47895725d6bd2b30ddce42658 Mon Sep 17 00:00:00 2001 From: Juan Pablo Diaz-Vaz Date: Fri, 24 Jul 2015 17:13:07 -0300 Subject: [PATCH 3/3] SUMOK-26 SUMOK-27 CloudWatch logs decomposed and header added --- .../client/CloudWatchLogsMessageModel.java | 119 ++++++++++++++++++ ...WatchMessageModelSumologicTransformer.java | 78 ++++++++++-- .../java/com/sumologic/client/LogEvent.java | 71 +++++++++++ ...hMessageModelSumologicTransformerTest.java | 116 +++++++++++------ 4 files changed, 334 insertions(+), 50 deletions(-) create mode 100644 src/main/java/com/sumologic/client/CloudWatchLogsMessageModel.java create mode 100644 src/main/java/com/sumologic/client/LogEvent.java 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 0370e07..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,20 +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. @@ -36,12 +48,41 @@ 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/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java b/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java index 410549a..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,12 +36,8 @@ public void theTransformerShouldFailGracefullyWhenUnableToTransform () { Record mockedRecord = new Record(); mockedRecord.setData(bufferedData); - SimpleKinesisMessageModel messageModel = null; - try { - messageModel = transfomer.toClass(mockedRecord); - } catch (IOException e) { - Assert.fail("Getting error while transforming: "+e.getMessage()); - } + CloudWatchLogsMessageModel messageModel = transfomer.toClass(mockedRecord); + Assert.assertNull(messageModel); } @@ -50,15 +46,20 @@ public void theTransformerShouldFailGracefullyWhenUnableToTransform () { public void theTransformerShouldSucceedWhenTransformingAProperJSON() { CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer(); - String jsonData = "[" + String jsonData = "" +"{" - + "\"id\": \"55b25585730b5e5bcd53c580\"," - + "\"index\": 0," - + "\"guid\": \"f3bf1c45-306d-4799-801f-c6d16acca931\"," - + "\"isActive\": false," - + "\"picture\": \"http://placehold.it/32x32\"" + + "\"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); @@ -72,31 +73,30 @@ public void theTransformerShouldSucceedWhenTransformingAProperJSON() { Record mockedRecord = new Record(); mockedRecord.setData(bufferedData); - SimpleKinesisMessageModel messageModel = null; - try { - messageModel = transfomer.toClass(mockedRecord); - } catch (IOException e) { - Assert.fail("Getting error while transforming: "+e.getMessage()); - } + CloudWatchLogsMessageModel messageModel = transfomer.toClass(mockedRecord); Assert.assertNotNull(messageModel); - Assert.assertTrue(messageModel.getData().equals(jsonData)); } @Test - public void theTransformerShouldSucceedWhenTransformingAJSONWithTrailingCommas() { + public void theTransformerShouldFailWhenTransformingAJSONWithTrailingCommas() { CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer(); - String jsonData = "[" - +"{" - + "\"id\": \"55b25585730b5e5bcd53c580\"," - + "\"index\": 0," - + "\"guid\": \"f3bf1c45-306d-4799-801f-c6d16acca931\"," - + "\"isActive\": false," - + "\"picture\": \"http://placehold.it/32x32\"," - + "}" - +"]"; - + 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; @@ -109,13 +109,55 @@ public void theTransformerShouldSucceedWhenTransformingAJSONWithTrailingCommas() Record mockedRecord = new Record(); mockedRecord.setData(bufferedData); - SimpleKinesisMessageModel messageModel = null; + 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 { - 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()); } - Assert.assertNull(messageModel); + 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