Skip to content

Commit 401bec3

Browse files
Juan Pablo Diaz-VazJuan Pablo Diaz-Vaz
authored andcommitted
SUMOK-23 Transformation errors handling properly
1 parent b390502 commit 401bec3

2 files changed

Lines changed: 71 additions & 0 deletions

File tree

src/main/java/com/mcplusa/sumologic/CloudWatchMessageModelSumologicTransformer.java

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,9 @@
1212
import java.io.ByteArrayInputStream;
1313
import java.io.BufferedReader;
1414
import java.io.InputStreamReader;
15+
import java.nio.ByteBuffer;
16+
import java.nio.charset.Charset;
17+
import java.nio.charset.CharsetDecoder;
1518
import java.util.zip.GZIPInputStream;
1619

1720
/**
@@ -39,6 +42,12 @@ public String fromClass(SimpleKinesisMessageModel message) {
3942
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
4043
byte[] decodedRecord = record.getData().array();
4144
String stringifiedRecord = decompressGzip(decodedRecord);
45+
46+
if (stringifiedRecord == null) {
47+
LOG.error("Unable to decompress the record: "+new String(record.getData().array()));
48+
LOG.error("Not attempting to transform into a Message Model");
49+
return null;
50+
}
4251

4352
return new SimpleKinesisMessageModel(stringifiedRecord);
4453
}
@@ -59,5 +68,19 @@ public static String decompressGzip(byte[] compressedData) {
5968
return null;
6069
}
6170
}
71+
72+
public static String byteBufferToString(ByteBuffer buffer){
73+
String data = "";
74+
CharsetDecoder decoder = Charset.forName("UTF-8").newDecoder();
75+
try{
76+
int old_position = buffer.position();
77+
data = decoder.decode(buffer).toString();
78+
buffer.position(old_position);
79+
}catch (Exception e){
80+
e.printStackTrace();
81+
return "";
82+
}
83+
return data;
84+
}
6285

6386
}
Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
1+
package com.mcplusa.sumologic;
2+
3+
import org.junit.Assert;
4+
import org.junit.Before;
5+
import org.junit.Rule;
6+
import org.junit.Test;
7+
import org.junit.Ignore;
8+
9+
import com.amazonaws.services.kinesis.model.Record;
10+
import com.mcplusa.sumologic.CloudWatchMessageModelSumologicTransformer;
11+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
12+
13+
import java.io.IOException;
14+
import java.nio.charset.Charset;
15+
import java.nio.charset.CharsetEncoder;
16+
import java.nio.CharBuffer;
17+
import java.nio.ByteBuffer;
18+
19+
20+
public class CloudWatchMessageModelSumologicTransformerTest {
21+
public static Charset charset = Charset.forName("UTF-8");
22+
public static CharsetEncoder encoder = charset.newEncoder();
23+
24+
@Test
25+
public void theTransformerShouldFailGracefullyWhenUnableToTransform () {
26+
CloudWatchMessageModelSumologicTransformer transfomer = new CloudWatchMessageModelSumologicTransformer();
27+
28+
String randomData = "Some random string without GZIP compression";
29+
ByteBuffer bufferedData = null;
30+
try {
31+
bufferedData = encoder.encode(CharBuffer.wrap(randomData));
32+
} catch (Exception e) {
33+
Assert.fail("Getting error: "+e.getMessage());
34+
}
35+
36+
Record mockedRecord = new Record();
37+
mockedRecord.setData(bufferedData);
38+
39+
SimpleKinesisMessageModel messageModel = null;
40+
try {
41+
messageModel = transfomer.toClass(mockedRecord);
42+
} catch (IOException e) {
43+
Assert.fail("Getting error while transforming: "+e.getMessage());
44+
}
45+
46+
Assert.assertNull(messageModel);
47+
}
48+
}

0 commit comments

Comments
 (0)