Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -9,14 +9,24 @@

class DeltaEncoder {
private long prev = 0;
private long undo = 0;

long encode(long v) {
long r = v - prev;
prev = v;
return r;
}

void commit() {
undo = prev;
}

void revert() {
prev = undo;
}

void clear() {
prev = 0;
undo = 0;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@

package com.datadoghq.dogstatsd.http.serializer;

import java.util.ArrayList;
import java.util.HashMap;

class Interner<T> {
Expand All @@ -15,6 +16,8 @@ interface Encoder<T> {
}

private HashMap<T, Long> inner = new HashMap<>();
private ArrayList<T> undo = new ArrayList<>();

private long lastId = 0;
private final Encoder<T> encoder;
private final T empty;
Expand All @@ -38,12 +41,26 @@ long intern(T val) {
encoder.encode(val);

lastId++;
undo.add(val);
inner.put(val, lastId);
return lastId;
}

void commit() {
undo.clear();
}

void revert() {
for (T it : undo) {
inner.remove(it);
}
lastId -= undo.size();
undo.clear();
}

void clear() {
inner.clear();
undo.clear();
lastId = 0;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -234,7 +234,21 @@ void endMetric() {
m.clearDependentFields();
m.encodeDependentFields();
}

payload.put(record);

nameStrInterner.commit();
tagStrInterner.commit();
tagsInterner.commit();
resourceStrInterner.commit();
resourcesInterner.commit();
originInfoInterner.commit();

nameRefsDelta.commit();
tagsetRefsDelta.commit();
resourcesRefsDelta.commit();
originInfoRefsDelta.commit();
timestampsDelta.commit();
} finally {
resetMetric();
}
Expand All @@ -251,6 +265,20 @@ public void resetMetric() {
timestamps.clear();
values.clear();
counts.clear();

nameStrInterner.revert();
tagStrInterner.revert();
tagsInterner.revert();
resourceStrInterner.revert();
resourcesInterner.revert();
originInfoInterner.revert();

nameRefsDelta.revert();
tagsetRefsDelta.revert();
resourcesRefsDelta.revert();
originInfoRefsDelta.revert();
timestampsDelta.revert();

metricInProgress = null;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,15 @@
package com.datadoghq.dogstatsd.http.serializer;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertThrows;
import static org.junit.Assert.assertTrue;

import com.datadoghq.dogstatsd.Sketch;
import java.nio.BufferOverflowException;
import java.util.ArrayList;
import java.util.Arrays;
import org.junit.Test;
import org.junit.function.ThrowingRunnable;

public class PayloadBuilderTest {
@Test
Expand Down Expand Up @@ -370,4 +373,70 @@ public void handle(byte[] p) {
assertTrue(p.length <= b.maxPayloadSize);
}
}

@Test
// Make sure we leave builder in consistent state on error.
public void rollback() {
final ArrayList<byte[]> payloads1 = new ArrayList<>();
PayloadBuilder pb1 =
new PayloadBuilder(
new PayloadConsumer() {
@Override
public void handle(byte[] p) {
payloads1.add(p);
}
});

final ArrayList<byte[]> payloads2 = new ArrayList<>();
PayloadBuilder pb2 =
new PayloadBuilder(
new PayloadConsumer() {
@Override
public void handle(byte[] p) {
payloads2.add(p);
}
});

pb1.count("m.before")
.setTags(Arrays.asList(new String[] {"env:prod"}))
.addPoint(100, 1)
.close();
pb2.count("m.before")
.setTags(Arrays.asList(new String[] {"env:prod"}))
.addPoint(100, 1)
.close();

// 40k * 8 bytes = 320 KB > 256 KB payload limit.
final ScalarMetric huge =
pb2.gauge("m.huge").setTags(Arrays.asList(new String[] {"env:prod", "kind:huge"}));
for (int i = 0; i < 40000; i++) {
huge.addPoint(1000, 3.14);
}
assertThrows(
BufferOverflowException.class,
new ThrowingRunnable() {
@Override
public void run() {
huge.close();
}
});
pb2.resetMetric();

pb1.gauge("m.huge")
.setTags(Arrays.asList(new String[] {"env:prod", "kind:huge"}))
.addPoint(200, 2)
.close();
pb2.gauge("m.huge")
.setTags(Arrays.asList(new String[] {"env:prod", "kind:huge"}))
.addPoint(200, 2)
.close();

pb1.close();
pb2.close();

assertEquals(payloads1.size(), payloads2.size());
for (int i = 0; i < payloads1.size(); i++) {
TestUtil.assertPayload(payloads2.get(i), payloads1.get(i));
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -103,9 +103,14 @@ static String protodump(byte[] p) {
}

static String protodump(int[] p) {
return protodump(toBytes(p));
}

// int[] lets tests write unsigned byte values without conversion.
static byte[] toBytes(int[] p) {
byte[] pb = new byte[p.length];
for (int i = 0; i < p.length; i++) pb[i] = (byte) p[i];
return protodump(pb);
return pb;
}

static void formatTwoCols(Formatter out, String hl, String hr, String dl, String dr) {
Expand All @@ -129,9 +134,13 @@ static void formatTwoCols(Formatter out, String hl, String hr, String dl, String

// expected is int[] to be able to write unsigned byte values wtihout conversion.
static void assertPayload(byte[] got, int[] expected) {
assertPayload(got, toBytes(expected));
}

static void assertPayload(byte[] got, byte[] expected) {
boolean same = got.length == expected.length;
for (int i = 0; same && i < got.length; i++) {
same &= got[i] == (byte) expected[i];
same &= got[i] == expected[i];
}
if (!same) {
Formatter out = new Formatter();
Expand Down
Loading