diff --git a/build.gradle b/build.gradle index 304060e4..eccbd0de 100644 --- a/build.gradle +++ b/build.gradle @@ -66,7 +66,6 @@ sourceSets { dependencies { // Apache Commons - compile 'commons-codec:commons-codec:1.10' compile 'commons-net:commons-net:3.3' // Apache HTTP @@ -88,6 +87,7 @@ dependencies { // Testing libraries testCompile 'junit:junit:4.11' + testCompile 'org.hamcrest:hamcrest:2.2' testCompile 'com.github.tomakehurst:wiremock:1.53' testCompile 'org.skyscreamer:jsonassert:1.2.3' testCompile 'org.mockito:mockito-core:3.2.4' diff --git a/examples/simple-console/src/main/java/com/snowplowanalytics/Main.java b/examples/simple-console/src/main/java/com/snowplowanalytics/Main.java index f84cf47d..a345bba9 100644 --- a/examples/simple-console/src/main/java/com/snowplowanalytics/Main.java +++ b/examples/simple-console/src/main/java/com/snowplowanalytics/Main.java @@ -18,18 +18,23 @@ import com.snowplowanalytics.snowplow.tracker.emitter.BatchEmitter; import com.snowplowanalytics.snowplow.tracker.emitter.Emitter; import com.snowplowanalytics.snowplow.tracker.emitter.RequestCallback; -import com.snowplowanalytics.snowplow.tracker.events.PageView; +import com.snowplowanalytics.snowplow.tracker.events.*; import com.snowplowanalytics.snowplow.tracker.http.HttpClientAdapter; import com.snowplowanalytics.snowplow.tracker.http.OkHttpClientAdapter; +import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; + import okhttp3.OkHttpClient; import java.util.List; +import java.util.Set; +import java.util.HashSet; import java.util.concurrent.TimeUnit; +import static java.util.Collections.singletonList; -public class Main { +import com.google.common.collect.ImmutableMap; - private static final int PAGEVIEW_COUNT = 10; +public class Main { public static String getUrlFromArgs(String[] args) { if (args == null || args.length < 1) { @@ -53,9 +58,10 @@ public static HttpClientAdapter getClient(String url) { } public static void main(String[] args) { + Set failedEventIds = new HashSet(); String collectorEndpoint = getUrlFromArgs(args); - System.out.println("Sending " + PAGEVIEW_COUNT + " events to " + collectorEndpoint); + System.out.println("Sending events to " + collectorEndpoint); // get the client adapter // this is used by the Java tracker to transmit events to the collector @@ -67,41 +73,107 @@ public static void main(String[] args) { String namespace = "demo"; // build an emitter, this is used by the tracker to batch and schedule transmission of events - Emitter emitter = BatchEmitter.builder() + BatchEmitter emitter = BatchEmitter.builder() .httpClientAdapter(okHttpClientAdapter) .requestCallback(new RequestCallback() { // let us know on successes (may be called multiple times) @Override - public void onSuccess(int successCount) { + public synchronized void onSuccess(int successCount) { System.out.println("Successfully sent " + successCount + " events"); } // let us know if something has gone wrong (may be called multiple times) @Override - public void onFailure(int successCount, List failedEvents) { + public synchronized void onFailure(int successCount, List failedEvents) { System.err.println("Successfully sent " + successCount + " events; failed to send " + failedEvents.size() + " events"); } }) - .bufferSize(1) // send an event every time one is given (no batching). In production this number should be higher, depending on the size/event volume + .bufferSize(4) // send an event every time one is given (no batching). In production this number should be higher, depending on the size/event volume .build(); // now we have the emitter, we need a tracker to turn our events into something a Snowplow collector can understand - Tracker tracker = new Tracker.TrackerBuilder(emitter, namespace, appId) - .base64(true) - .platform(DevicePlatform.ServerSideApp) - .build(); + final Tracker tracker = new Tracker.TrackerBuilder(emitter, namespace, appId) + .base64(true) + .platform(DevicePlatform.ServerSideApp) + .build(); - for (int i = 0; i < PAGEVIEW_COUNT; i++) { - // This is a sample page view event, many other event types (such as self-describing events) are available - PageView pageViewEvent = PageView.builder() - .pageTitle("Hello world " + i) - .pageUrl("https://www.snowplowanalytics.com") - .referrer("https://www.google.com") - .build(); - - tracker.track(pageViewEvent); // the .track method schedules the event for delivery to Snowplow - } + // This is an example of a custom context + List contexts = singletonList( + new SelfDescribingJson( + "iglu:com.snowplowanalytics.iglu/anything-c/jsonschema/1-0-0", + ImmutableMap.of("foo", "bar"))); + + // This is a sample page view event, many other event types (such as self-describing events) are available + PageView pageViewEvent = PageView.builder() + .pageTitle("Snowplow Analytics") + .pageUrl("https://www.snowplowanalytics.com") + .referrer("https://www.google.com") + .customContext(contexts) + .build(); + + tracker.track(pageViewEvent); // the .track method schedules the event for delivery to Snowplow + + EcommerceTransactionItem item = EcommerceTransactionItem.builder() + .itemId("order_id") + .sku("sku") + .price(1.0) + .quantity(2) + .name("name") + .category("category") + .currency("currency") + .customContext(contexts) + .build(); + + EcommerceTransaction ecommerceTransaction = EcommerceTransaction.builder() + .orderId("order_id") + .totalValue(1.0) + .affiliation("affiliation") + .taxValue(2.0) + .shipping(3.0) + .city("city") + .state("state") + .country("country") + .currency("currency") + .items(item) // EcommerceTransactionItem events are added to a parent EcommerceTransaction + .customContext(contexts) + .build(); + + tracker.track(ecommerceTransaction); // This will track two events + + // This is an example of a custom "Unsutrcutred" event based on a schema + Unstructured unstructured = Unstructured.builder() + .eventData(new SelfDescribingJson( + "iglu:com.snowplowanalytics.iglu/anything-a/jsonschema/1-0-0", + ImmutableMap.of("foo", "bar") + )) + .customContext(contexts) + .build(); + + tracker.track(unstructured); + + // This is an example of a ScreenView event which will be translated into an Unstructured event + ScreenView screenView = ScreenView.builder() + .name("name") + .id("id") + .customContext(contexts) + .build(); + + tracker.track(screenView); + + // This is an example of a Timing event which will be translated into an Unstructured event + Timing timing = Timing.builder() + .category("category") + .label("label") + .variable("variable") + .timing(10) + .customContext(contexts) + .build(); + + tracker.track(timing); + // Will close all threads and force send remaining events + // should be 1 left to flush, as we send 5 events with a bufferSize of 4 + emitter.close(); } } diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/Subject.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/Subject.java index 8b06b116..f591c93c 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/Subject.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/Subject.java @@ -44,6 +44,14 @@ private Subject(SubjectBuilder builder) { this.setDomainUserId(builder.domainUserId); } + /** + * Creates a new {@link Subject} object based on the map of another {@link Subject} object. + * @param subject The subject from which the map is copied. + */ + public Subject(Subject subject){ + this.standardPairs.putAll(subject.getSubject()); + } + /** * Builder for the Subject */ diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/Tracker.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/Tracker.java index be31db39..0f8cdf0e 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/Tracker.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/Tracker.java @@ -12,29 +12,18 @@ */ package com.snowplowanalytics.snowplow.tracker; -// Java -import java.util.*; - -// Google import com.google.common.base.Preconditions; -// This library -import com.snowplowanalytics.snowplow.tracker.constants.Constants; -import com.snowplowanalytics.snowplow.tracker.constants.Parameter; import com.snowplowanalytics.snowplow.tracker.emitter.Emitter; import com.snowplowanalytics.snowplow.tracker.events.*; -import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; -import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerEvent; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerParameters; public class Tracker { - private final String trackerVersion = Version.TRACKER; private Emitter emitter; private Subject subject; - private String appId; - private String namespace; - private DevicePlatform platform; - private boolean base64Encoded; + private final TrackerParameters parameters; /** * Creates a new Snowplow Tracker. @@ -50,12 +39,9 @@ private Tracker(TrackerBuilder builder) { Preconditions.checkArgument(!builder.namespace.isEmpty(), "namespace cannot be empty"); Preconditions.checkArgument(!builder.appId.isEmpty(), "appId cannot be empty"); + this.parameters = new TrackerParameters(builder.appId, builder.platform, builder.namespace, Version.TRACKER, builder.base64Encoded); this.emitter = builder.emitter; - this.namespace = builder.namespace; - this.appId = builder.appId; this.subject = builder.subject; - this.platform = builder.platform; - this.base64Encoded = builder.base64Encoded; } /** @@ -137,44 +123,6 @@ public void setSubject(Subject subject) { this.subject = subject; } - /** - * Sets the Trackers platform, defaults to a - * Server Side Application. - * - * @param platform the DevicePlatform - */ - public void setPlatform(DevicePlatform platform) { - this.platform = platform; - } - - /** - * Sets whether to base64 Encode custom contexts - * and unstructured events - * - * @param base64Encoded a boolean truth - */ - public void setBase64Encoded(boolean base64Encoded) { - this.base64Encoded = base64Encoded; - } - - /** - * Sets a new Application ID - * - * @param appId the new application id - */ - public void setAppId(String appId) { - this.appId = appId; - } - - /** - * Sets a new Tracker Namespace - * - * @param namespace the new tracker namespace - */ - public void setNamespace(String namespace) { - this.namespace = namespace; - } - // --- Getters /** @@ -195,49 +143,46 @@ public Subject getSubject() { * @return the tracker version that was set */ public String getTrackerVersion() { - return this.trackerVersion; + return this.parameters.getTrackerVersion(); } /** * @return the trackers namespace */ public String getNamespace() { - return this.namespace; + return this.parameters.getNamespace(); } /** * @return the trackers set Application ID */ public String getAppId() { - return this.appId; + return this.parameters.getAppId(); } /** * @return the base64 setting of the tracker */ public boolean getBase64Encoded() { - return this.base64Encoded; + return this.parameters.getBase64Encoded(); } /** * @return the Tracker platform */ public DevicePlatform getPlatform() { - return this.platform; + return this.parameters.getPlatform(); } - // --- Event Tracking Functions - /** - * Used for either Tracking a custom TrackerPayload or - * for re-sending a failed event. - * - * @param payload the payload to track + * @return the wrapper containing the Tracker parameters */ - public void track(TrackerPayload payload) { - this.emitter.emit(payload); + public TrackerParameters getParameters() { + return this.parameters; } + // --- Event Tracking Functions + /** * Handles tracking the different types of events that * the Tracker can encounter. @@ -245,89 +190,7 @@ public void track(TrackerPayload payload) { * @param event the event to track */ public void track(Event event) { - List context = event.getContext(); - Subject subject = event.getSubject(); - - // Figure out what type of event it is and track it! - Class eClass = event.getClass(); - if (eClass.equals(PageView.class) || eClass.equals(Structured.class)) { - this.addTrackerPayload((TrackerPayload) event.getPayload(), context, subject); - } else if (eClass.equals(EcommerceTransaction.class)) { - this.addTrackerPayload((TrackerPayload) event.getPayload(), context, subject); - - // Track each item individually - EcommerceTransaction ecommerceTransaction = (EcommerceTransaction) event; - for(EcommerceTransactionItem item : ecommerceTransaction.getItems()) { - item.setTimestamp(ecommerceTransaction.getTimestamp()); - this.addTrackerPayload(item.getPayload(), item.getContext(), item.getSubject()); - } - } else if (eClass.equals(Unstructured.class)) { - - // Need to set the Base64 rule for Unstructured events - Unstructured unstructured = (Unstructured) event; - unstructured.setBase64Encode(base64Encoded); - this.addTrackerPayload(unstructured.getPayload(), context, subject); - } else if (eClass.equals(Timing.class) || eClass.equals(ScreenView.class)) { - - // These are wrapper classes for Unstructured events; need to create Unstructured - // events from them and resend. - this.track(Unstructured.builder() - .eventData((SelfDescribingJson) event.getPayload()) - .customContext(context) - .deviceCreatedTimestamp(event.getDeviceCreatedTimestamp()) - .trueTimestamp(event.getTrueTimestamp()) - .eventId(event.getEventId()) - .subject(subject) - .build()); - } - } - - // --- Helpers - - /** - * Builds and Adds a finalised payload which is ready for sending. - * - * @param payload The raw event Payload - * @param contexts Custom context for the event - * @param eventSubject An optional event specific Subject - */ - private void addTrackerPayload(TrackerPayload payload, List contexts, Subject eventSubject) { - - // Add default parameters to the payload - payload.add(Parameter.PLATFORM, platform.toString()); - payload.add(Parameter.APP_ID, this.appId); - payload.add(Parameter.NAMESPACE, this.namespace); - payload.add(Parameter.TRACKER_VERSION, this.trackerVersion); - - // Build the final context and add it to the payload - if (contexts != null && contexts.size() > 0) { - SelfDescribingJson envelope = getFinalContext(contexts); - payload.addMap(envelope.getMap(), this.base64Encoded, Parameter.CONTEXT_ENCODED, Parameter.CONTEXT); - } - - // Add subject if available - if (eventSubject != null) { - payload.addMap(new HashMap<>(eventSubject.getSubject())); - } else if (this.subject != null) { - payload.addMap(new HashMap<>(this.subject.getSubject())); - } - - // Send the event! - this.emitter.emit(payload); - } - - /** - * Builds the final event context. - * - * @param contexts the base event context - * @return the final event context json with - * many contexts inside - */ - private SelfDescribingJson getFinalContext(List contexts) { - List contextMaps = new LinkedList<>(); - for (SelfDescribingJson selfDescribingJson : contexts) { - contextMaps.add(selfDescribingJson.getMap()); - } - return new SelfDescribingJson(Constants.SCHEMA_CONTEXTS, contextMaps); + // Emit the event + this.emitter.emit(new TrackerEvent(event, this.parameters, this.subject)); } } diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/Utils.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/Utils.java index a59c7ac0..ebcfdf3d 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/Utils.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/Utils.java @@ -12,23 +12,17 @@ */ package com.snowplowanalytics.snowplow.tracker; -// Java import java.nio.charset.Charset; import java.util.*; import java.net.URL; import java.net.URLEncoder; -// Jackson import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; -// Slf4j import org.slf4j.Logger; import org.slf4j.LoggerFactory; -// Apache -import static org.apache.commons.codec.binary.Base64.encodeBase64String; - /** * Provides basic Utilities for the Snowplow Tracker. */ @@ -111,7 +105,7 @@ public static String getTimezone() { * @return a Base64 encoded string */ public static String base64Encode(String string, Charset charset) { - return encodeBase64String(string.getBytes(charset)); + return Base64.getEncoder().encodeToString(string.getBytes(charset)); } /** @@ -121,7 +115,7 @@ public static String base64Encode(String string, Charset charset) { * @param map the map to process into a JSON String * @return the final JSON String */ - public static String mapToJSONString(Map map) { + public static String mapToJSONString(Map map) { String jString = ""; try { jString = objectMapper.writeValueAsString(map); @@ -137,7 +131,7 @@ public static String mapToJSONString(Map map) { * @param map The map to convert * @return the QueryString ready for sending */ - public static String mapToQueryString(Map map) { + public static String mapToQueryString(Map map) { StringBuilder sb = new StringBuilder(); for (String key : map.keySet()) { if (sb.length() > 0) { diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/AbstractEmitter.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/AbstractEmitter.java index e3fd3863..9e483f62 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/AbstractEmitter.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/AbstractEmitter.java @@ -12,18 +12,14 @@ */ package com.snowplowanalytics.snowplow.tracker.emitter; -// Java -import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -// Google import com.google.common.base.Preconditions; -// This library import com.snowplowanalytics.snowplow.tracker.http.HttpClientAdapter; -import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerEvent; /** * AbstractEmitter class which contains common elements to @@ -34,8 +30,6 @@ public abstract class AbstractEmitter implements Emitter { protected HttpClientAdapter httpClientAdapter; protected RequestCallback requestCallback; protected ExecutorService executor; - protected List buffer = new ArrayList<>(); - protected int bufferSize = 1; public static abstract class Builder> { @@ -102,28 +96,24 @@ protected AbstractEmitter(final Builder builder) { } /** - * Adds a payload to the buffer and checks whether we have reached the buffer - * limit yet. + * Adds an event to the buffer * - * @param payload an event payload + * @param event an event */ @Override - public abstract void emit(TrackerPayload payload); + public abstract void emit(TrackerEvent event); /** - * Customize the emitter buffer size to any valid integer greater than zero. - - * Will only effect the BatchEmitter + * Customize the emitter buffer size to any valid integer greater than zero. + * Has no effect on SimpleEmitter * * @param bufferSize number of events to collect before sending */ @Override - public void setBufferSize(final int bufferSize) { - Preconditions.checkArgument(bufferSize > 0, "bufferSize must be greater than 0"); - this.bufferSize = bufferSize; - } + public abstract void setBufferSize(final int bufferSize); /** - * When the buffer limit is reached sending of the buffer is initiated. + * Removes all events from the buffer and sends them */ @Override public abstract void flushBuffer(); @@ -134,19 +124,15 @@ public void setBufferSize(final int bufferSize) { * @return the buffer size */ @Override - public int getBufferSize() { - return this.bufferSize; - } + public abstract int getBufferSize(); /** - * Returns the List of Payloads that are in the buffer. + * Returns List of Events that are in the buffer. * - * @return the buffer payloads + * @return the buffered events */ @Override - public List getBuffer() { - return this.buffer; - } + public abstract List getBuffer(); /** * Sends a runnable to the executor service. diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/BatchEmitter.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/BatchEmitter.java index fcfda15f..c5024f9a 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/BatchEmitter.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/BatchEmitter.java @@ -12,34 +12,44 @@ */ package com.snowplowanalytics.snowplow.tracker.emitter; -// Java import java.io.Closeable; import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; -// Google import com.google.common.base.Preconditions; - -// Slf4j -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -// This library import com.snowplowanalytics.snowplow.tracker.constants.Constants; import com.snowplowanalytics.snowplow.tracker.constants.Parameter; -import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerEvent; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** - * An emitter that emit a batch of events in a single call + * An emitter that emit a batch of events in a single call * It uses the post method of under-laying http adapter */ public class BatchEmitter extends AbstractEmitter implements Closeable { private static final Logger LOGGER = LoggerFactory.getLogger(BatchEmitter.class); + private final Thread bufferConsumer; + private boolean isClosing = false; + + private int bufferSize = 1; + + // Queue for immediate buffering of events + private final BlockingQueue eventBuffer = new LinkedBlockingQueue<>(); + + // Queue for storing events until bufferSize is reached + private final BlockingQueue eventsToSend = new LinkedBlockingQueue<>(); + private final long closeTimeout = 5; public static abstract class Builder> extends AbstractEmitter.Builder { @@ -78,37 +88,113 @@ protected BatchEmitter(final Builder builder) { Preconditions.checkArgument(builder.bufferSize > 0, "bufferSize must be greater than 0"); this.bufferSize = builder.bufferSize; + + bufferConsumer = new Thread(getBufferConsumerRunnable()); + bufferConsumer.start(); } /** - * Adds a payload to the buffer and checks whether we have reached the buffer - * limit yet. + * Adds a TrackerEvent to the concurrent queue buffer * - * @param payload an event payload + * @param event an event */ @Override - public synchronized void emit(final TrackerPayload payload) { - buffer.add(payload); - if (buffer.size() >= bufferSize) { - flushBuffer(); + public void emit(final TrackerEvent event) { + boolean result = eventBuffer.offer(event); // Add to buffer and quickly return back to application + + if (!result) { + LOGGER.error("Unable to add event to emitter, emitter buffer is full"); } } - /** - * When the buffer limit is reached sending of the buffer is initiated. + /* + * Forces the events currently in the buffer to be sent */ + @Override public void flushBuffer() { - execute(getRequestRunnable(buffer)); - buffer = new ArrayList<>(); + // Drain immediate event buffer + while (true) { + TrackerEvent event = eventBuffer.poll(); + if (event == null) { + break; + } else { + eventsToSend.offer(event); + } + } + + drainBufferAndSend(); + } + + /** + * Returns List of Events that are in the buffer. + * + * @return the buffered events + */ + @Override + public List getBuffer() { + return eventsToSend.stream().collect(Collectors.toList()); + } + + /** + * Customize the emitter buffer size to any valid integer greater than zero. + * + * @param bufferSize number of events to collect before sending + */ + @Override + public void setBufferSize(final int bufferSize) { + Preconditions.checkArgument(bufferSize > 0, "bufferSize must be greater than 0"); + this.bufferSize = bufferSize; + } + + /** + * Gets the Emitter Buffer Size + * + * @return the buffer size + */ + @Override + public int getBufferSize() { + return this.bufferSize; + } + + /** + * Returns a Consumer for the concurrent queue buffer + * Consumes events onto another queue to be sent when bufferSize is reached + * + * @return the new Runnable object + */ + private Runnable getBufferConsumerRunnable() { + return new Runnable() { + @Override + public void run() { + while (true) { + try { + eventsToSend.put(eventBuffer.take()); + if (eventsToSend.size() >= bufferSize) { + drainBufferAndSend(); + } + } catch (InterruptedException ex) { + if (isClosing) { + return; + } + } + } + } + }; + } + + private void drainBufferAndSend() { + List events = new ArrayList<>(); + eventsToSend.drainTo(events); + execute(getRequestRunnable(events)); } /** * Returns a Runnable POST Request operation * * @param buffer the event buffer to be sent - * @return the new Callable object + * @return the new Runnable object */ - private Runnable getRequestRunnable(final List buffer) { + private Runnable getRequestRunnable(final List buffer) { return new Runnable() { @Override public void run() { @@ -133,7 +219,8 @@ public void run() { // Send the callback if available if (requestCallback != null) { if (failure != 0) { - requestCallback.onFailure(success, buffer); + requestCallback.onFailure(success, + buffer.stream().map(te -> te.getEvent()).collect(Collectors.toList())); } else { requestCallback.onSuccess(success); } @@ -145,13 +232,19 @@ public void run() { /** * Constructs the SelfDescribingJson to be sent to the endpoint * + * @param buffer the event buffer * @return the constructed POST payload */ - private SelfDescribingJson getFinalPost(final List buffer) { - final List toSendPayloads = new ArrayList<>(); - for (final TrackerPayload payload : buffer) { - payload.add(Parameter.DEVICE_SENT_TIMESTAMP, Long.toString(System.currentTimeMillis())); - toSendPayloads.add(payload.getMap()); + private SelfDescribingJson getFinalPost(final List buffer) { + final List> toSendPayloads = new ArrayList<>(); + final String sentTimestamp = Long.toString(System.currentTimeMillis()); + + for (TrackerEvent event : buffer) { + List payloads = event.getTrackerPayloads(); + for (TrackerPayload payload : payloads) { + payload.add(Parameter.DEVICE_SENT_TIMESTAMP, sentTimestamp); + toSendPayloads.add(payload.getMap()); + } } return new SelfDescribingJson(Constants.SCHEMA_PAYLOAD_DATA, toSendPayloads); @@ -162,7 +255,12 @@ private SelfDescribingJson getFinalPost(final List buffer) { */ @Override public void close() { - flushBuffer(); + isClosing = true; + + bufferConsumer.interrupt(); // Kill buffer consumer + flushBuffer(); // Attempt to send all reminaing events + + //Shutdown executor threadpool if (executor != null) { executor.shutdown(); try { diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/Emitter.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/Emitter.java index 0eab9e1c..e1cd535f 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/Emitter.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/Emitter.java @@ -12,23 +12,22 @@ */ package com.snowplowanalytics.snowplow.tracker.emitter; -// This library -import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; - import java.util.List; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerEvent; + /** * Emitter interface. */ public interface Emitter { /** - * Adds a payload to the buffer and checks whether + * Adds an event to the buffer and checks whether * we have reached the buffer limit yet. * - * @param payload an event payload + * @param event an event to be emitted */ - void emit(TrackerPayload payload); + void emit(TrackerEvent event); /** * Customize the emitter buffer size to any valid integer @@ -57,9 +56,9 @@ public interface Emitter { int getBufferSize(); /** - * Returns the List of Payloads that are in the buffer. + * Returns the List of Events that are in the buffer. * - * @return the buffer payloads + * @return the buffer events */ - List getBuffer(); + List getBuffer(); } diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/RequestCallback.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/RequestCallback.java index 00baeee4..deb9cf05 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/RequestCallback.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/RequestCallback.java @@ -12,11 +12,9 @@ */ package com.snowplowanalytics.snowplow.tracker.emitter; -// Java import java.util.List; -// This library -import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; +import com.snowplowanalytics.snowplow.tracker.events.Event; /** * Provides a callback interface for reporting counts of successfully sent @@ -34,10 +32,10 @@ public interface RequestCallback { /** * If all/some events failed then the count of successful - * events is returned along with all the failed Payloads. + * events is returned along with all the failed Events. * * @param successCount the successful count - * @param failedEvents the list of failed payloads + * @param failedEvents the list of failed events */ - void onFailure(int successCount, List failedEvents); + void onFailure(int successCount, List failedEvents); } diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/SimpleEmitter.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/SimpleEmitter.java index 823801bd..19c39c18 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/SimpleEmitter.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/emitter/SimpleEmitter.java @@ -12,17 +12,16 @@ */ package com.snowplowanalytics.snowplow.tracker.emitter; -// Java import java.util.ArrayList; import java.util.List; -// Slf4j import org.slf4j.Logger; import org.slf4j.LoggerFactory; -// This library +import com.snowplowanalytics.snowplow.tracker.payload.TrackerEvent; import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; import com.snowplowanalytics.snowplow.tracker.constants.Parameter; +import com.snowplowanalytics.snowplow.tracker.events.Event; /** * An emitter which sends events as soon as they are received via @@ -54,18 +53,20 @@ protected SimpleEmitter(final Builder builder) { } /** - * Adds a payload to the buffer and instantly sends it + * Adds an event to the buffer and instantly sends it * - * @param payload an event payload + * @param event an event */ @Override - public void emit(final TrackerPayload payload) { - execute(getRequestRunnable(payload)); + public void emit(final TrackerEvent event) { + execute(getRequestRunnable(event)); } /** - * When the buffer limit is reached sending of the buffer is initiated. + * Sends buffered events, but SimpleEmitter does not buffer events + * So has no effect */ + @Override public void flushBuffer() { // Do nothing! } @@ -73,32 +74,37 @@ public void flushBuffer() { /** * Returns a Runnable GET Request operation * - * @param payload the event to be sent + * @param event the event to be sent * @return the new Callable object */ - private Runnable getRequestRunnable(final TrackerPayload payload) { + private Runnable getRequestRunnable(final TrackerEvent event) { return new Runnable() { @Override public void run() { - payload.add(Parameter.DEVICE_SENT_TIMESTAMP, Long.toString(System.currentTimeMillis())); - final int code = httpClientAdapter.get(payload); - - // Process results int success = 0; int failure = 0; - if (!isSuccessfulSend(code)) { - LOGGER.error("SimpleEmitter failed to send {} events: code: {}", 1, code); - failure += 1; - } else { - LOGGER.debug("SimpleEmitter successfully sent {} events: code: {}", 1, code); - success += 1; + + List payloads = event.getTrackerPayloads(); + + for (TrackerPayload payload : payloads) { + payload.add(Parameter.DEVICE_SENT_TIMESTAMP, Long.toString(System.currentTimeMillis())); + final int code = httpClientAdapter.get(payload); + + // Process results + if (!isSuccessfulSend(code)) { + LOGGER.error("SimpleEmitter failed to send {} events: code: {}", 1, code); + failure += 1; + } else { + LOGGER.debug("SimpleEmitter successfully sent {} events: code: {}", 1, code); + success += 1; + } } // Send the callback if available if (requestCallback != null) { if (failure != 0) { - final List buffer = new ArrayList<>(); - buffer.add(payload); + final List buffer = new ArrayList<>(); + buffer.add(event.getEvent()); requestCallback.onFailure(success, buffer); } else { requestCallback.onSuccess(success); @@ -107,4 +113,38 @@ public void run() { } }; } + + /** + * Returns List of Events that are in the buffer. + * Always empty for SimpleEmitter + * + * @return the empty buffer + */ + @Override + public List getBuffer() { + return new ArrayList<>(); + } + + /** + * Customize the emitter buffer size to any valid integer greater than zero. + * Has no effect on SimpleEmitter + * + * @param bufferSize number of events to collect before sending + */ + @Override + public void setBufferSize(final int bufferSize) { + if (bufferSize != 1) { + LOGGER.debug("Noop. SimpleEmitter buffer size must always be 1."); + } + } + + /** + * Gets the Emitter Buffer Size - Will always be 1 for SimpleEmitter + * + * @return the buffer size + */ + @Override + public int getBufferSize() { + return 1; + } } diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/events/Event.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/events/Event.java index 2a65d54a..0c40230b 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/events/Event.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/events/Event.java @@ -12,10 +12,8 @@ */ package com.snowplowanalytics.snowplow.tracker.events; -// Java import java.util.List; -// This library import com.snowplowanalytics.snowplow.tracker.Subject; import com.snowplowanalytics.snowplow.tracker.payload.Payload; import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/events/ScreenView.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/events/ScreenView.java index d3210edc..034bd395 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/events/ScreenView.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/events/ScreenView.java @@ -12,10 +12,8 @@ */ package com.snowplowanalytics.snowplow.tracker.events; -// Google import com.google.common.base.Preconditions; -// This library import com.snowplowanalytics.snowplow.tracker.constants.Parameter; import com.snowplowanalytics.snowplow.tracker.constants.Constants; import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/events/Unstructured.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/events/Unstructured.java index 3c289575..5a5f34cd 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/events/Unstructured.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/events/Unstructured.java @@ -34,13 +34,13 @@ public static abstract class Builder> extends AbstractEvent private SelfDescribingJson eventData; /** - * @param eventData The properties of the event. Has two field: + * @param selfDescribingJson The properties of the event. Has two field: * A "data" field containing the event properties and * A "schema" field identifying the schema against which the data is validated * @return itself */ - public T eventData(SelfDescribingJson eventData) { - this.eventData = eventData; + public T eventData(SelfDescribingJson selfDescribingJson) { + this.eventData = selfDescribingJson; return self(); } diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/http/AbstractHttpClientAdapter.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/http/AbstractHttpClientAdapter.java index 0468aabd..ebd3c96c 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/http/AbstractHttpClientAdapter.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/http/AbstractHttpClientAdapter.java @@ -12,13 +12,8 @@ */ package com.snowplowanalytics.snowplow.tracker.http; -// Java -import java.util.Map; - -// Google import com.google.common.base.Preconditions; -// This library import com.snowplowanalytics.snowplow.tracker.constants.Constants; import com.snowplowanalytics.snowplow.tracker.Utils; import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; @@ -94,7 +89,6 @@ public int post(SelfDescribingJson payload) { * @param payload the TrackerPayload to send */ @Override - @SuppressWarnings("unchecked") public int get(TrackerPayload payload) { String url = this.url + "/i?" + Utils.mapToQueryString(payload.getMap()); return doGet(url); diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/http/ApacheHttpClientAdapter.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/http/ApacheHttpClientAdapter.java index 97e9d5ab..1499d294 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/http/ApacheHttpClientAdapter.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/http/ApacheHttpClientAdapter.java @@ -12,26 +12,18 @@ */ package com.snowplowanalytics.snowplow.tracker.http; -// Java -import java.util.Map; - -// Google import com.google.common.base.Preconditions; -// Apache import org.apache.http.HttpResponse; import org.apache.http.client.methods.HttpGet; import org.apache.http.client.methods.HttpPost; -import org.apache.http.client.utils.URIBuilder; import org.apache.http.entity.ContentType; import org.apache.http.entity.StringEntity; import org.apache.http.impl.client.CloseableHttpClient; -// Slf4j import org.slf4j.Logger; import org.slf4j.LoggerFactory; -// This library import com.snowplowanalytics.snowplow.tracker.constants.Constants; /** diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/Payload.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/Payload.java index 315ee8a0..8413a0b8 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/Payload.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/Payload.java @@ -12,7 +12,6 @@ */ package com.snowplowanalytics.snowplow.tracker.payload; -// Java import java.util.Map; /** @@ -48,14 +47,14 @@ public interface Payload { * @param typeEncoded The key that would be set if the encoding option was set to true * @param typeNotEncoded They key that would be set if the encoding option was set to false */ - void addMap(Map map, boolean base64Encoded, String typeEncoded, String typeNotEncoded); + void addMap(Map map, boolean base64Encoded, String typeEncoded, String typeNotEncoded); /** * Returns the Payload as a HashMap. * * @return A HashMap */ - Map getMap(); + Map getMap(); /** * Returns the byte size of a payload. diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/SelfDescribingJson.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/SelfDescribingJson.java index 18cc98d8..e7f88c85 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/SelfDescribingJson.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/SelfDescribingJson.java @@ -12,18 +12,14 @@ */ package com.snowplowanalytics.snowplow.tracker.payload; -// Java import java.util.LinkedHashMap; import java.util.Map; -// Google import com.google.common.base.Preconditions; -// Slf4j import org.slf4j.Logger; import org.slf4j.LoggerFactory; -// This library import com.snowplowanalytics.snowplow.tracker.Utils; import com.snowplowanalytics.snowplow.tracker.constants.Parameter; @@ -157,7 +153,7 @@ public void addMap(Map map) { @Deprecated @Override - public void addMap(Map map, boolean base64Encoded, String typeEncoded, String typeNotEncoded) { + public void addMap(Map map, boolean base64Encoded, String typeEncoded, String typeNotEncoded) { LOGGER.info("Payload: addMap(Map, boolean, String, String) method called - Doing nothing."); } @@ -167,7 +163,7 @@ public void addMap(Map map, boolean base64Encoded, String typeEncoded, String ty * @return A Map of all the key-value entries */ @Override - public Map getMap() { + public Map getMap() { return payload; } diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerEvent.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerEvent.java new file mode 100644 index 00000000..d6ef6eef --- /dev/null +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerEvent.java @@ -0,0 +1,164 @@ +/* + * Copyright (c) 2020-2020 Snowplow Analytics Ltd. All rights reserved. + * + * This program is licensed to you under the Apache License Version 2.0, + * and you may not use this file except in compliance with the Apache License Version 2.0. + * You may obtain a copy of the Apache License Version 2.0 at http://www.apache.org/licenses/LICENSE-2.0. + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the Apache License Version 2.0 is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the Apache License Version 2.0 for the specific language governing permissions and limitations there under. + */ +package com.snowplowanalytics.snowplow.tracker.payload; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; + +import com.snowplowanalytics.snowplow.tracker.Subject; +import com.snowplowanalytics.snowplow.tracker.constants.Constants; +import com.snowplowanalytics.snowplow.tracker.constants.Parameter; +import com.snowplowanalytics.snowplow.tracker.events.*; + +/** + * A TrackerEvent which allows the TrackerPayload to be filled later. The payload will be + * filled by the Emitter in the Emitter thread, using the getTrackerPayload() method. + */ +public class TrackerEvent { + + private final Event event; + private final TrackerParameters parameters; + private final Subject subject; + + public TrackerEvent(final Event event, final TrackerParameters parameters, final Subject subject) { + this.event = event; + this.parameters = parameters; + this.subject = subject; + } + + /** + * Returns the {@link Event} + * + * @return The {@link Event} + */ + public Event getEvent() { + return this.event; + } + + /** + * Converts a {@link Event} to a list of {@link TrackerPayload} and caches the values. + * Returns a list as some Events contain nested payloads (e.g. {@link EcommerceTransaction}) + * Adds fields to the {@link TrackerPayload} based on the type of the {@link Event}. + * + * @return The populated TrackerPayloads + */ + public List getTrackerPayloads() { + final List payloads = new ArrayList<>(); + final List contexts = event.getContext(); + final Subject subject = event.getSubject(); + + // Figure out what type of event it is + final Class eventClass = event.getClass(); + + if (eventClass.equals(Unstructured.class)) { + + // Need to set the Base64 rule for Unstructured events + final Unstructured unstructured = (Unstructured) event; + unstructured.setBase64Encode(this.parameters.getBase64Encoded()); + TrackerPayload payload = unstructured.getPayload(); + addTrackerParameters(payload); + addContextsAndSubject(contexts, subject, payload); + payloads.add(payload); + } else if (eventClass.equals(Timing.class) || eventClass.equals(ScreenView.class)) { + + // These are wrapper classes for Unstructured events; need to create + // Unstructured events from them and resend. + final Unstructured unstructured = Unstructured.builder() + .eventData((SelfDescribingJson) event.getPayload()) + .customContext(contexts) + .deviceCreatedTimestamp(event.getDeviceCreatedTimestamp()) + .trueTimestamp(event.getTrueTimestamp()) + .eventId(event.getEventId()) + .subject(subject) + .build(); + + unstructured.setBase64Encode(this.parameters.getBase64Encoded()); + TrackerPayload payload = unstructured.getPayload(); + addTrackerParameters(payload); + addContextsAndSubject(contexts, subject, payload); + payloads.add(payload); + } else if (eventClass.equals(EcommerceTransaction.class)) { + + final EcommerceTransaction ecommerceTransaction = (EcommerceTransaction) event; + TrackerPayload payload = ecommerceTransaction.getPayload(); + addTrackerParameters(payload); + addContextsAndSubject(contexts, subject, payload); + payloads.add(payload); + + // Track each item individually + for (final EcommerceTransactionItem item : ecommerceTransaction.getItems()) { + + item.setDeviceCreatedTimestamp(ecommerceTransaction.getDeviceCreatedTimestamp()); + TrackerPayload itemPayload = item.getPayload(); + addTrackerParameters(itemPayload); + addContextsAndSubject(item.getContext(), item.getSubject(), itemPayload); + payloads.add(itemPayload); + } + } else { + + // For all other events, simply get the payload + TrackerPayload payload = (TrackerPayload) event.getPayload(); + addTrackerParameters(payload); + addContextsAndSubject(contexts, subject, payload); + payloads.add(payload); + } + + return payloads; + } + + /** + * Adds the context and subject to the event payload + * + * @param contexts the base event context - can be null or empty + * @param subject the event subject - can be null + * @param payload the payload to add the contexts and subjects to + */ + private void addContextsAndSubject(final List contexts, final Subject subject, TrackerPayload payload) { + // Build the final context and add it to the payload + if (contexts != null && contexts.size() > 0) { + SelfDescribingJson envelope = getFinalContext(contexts); + payload.addMap(envelope.getMap(), this.parameters.getBase64Encoded(), Parameter.CONTEXT_ENCODED, Parameter.CONTEXT); + } + + // Add subject if available + if (subject != null) { + payload.addMap(new HashMap<>(subject.getSubject())); + } else if (this.subject != null) { + payload.addMap(new HashMap<>(this.subject.getSubject())); + } + } + + /** + * Builds the final event context. + * + * @param contexts the base event context + * @return the final event context json with many contexts inside + */ + private SelfDescribingJson getFinalContext(List contexts) { + List> contextMaps = new LinkedList<>(); + for (SelfDescribingJson selfDescribingJson : contexts) { + contextMaps.add(selfDescribingJson.getMap()); + } + return new SelfDescribingJson(Constants.SCHEMA_CONTEXTS, contextMaps); + } + + private void addTrackerParameters(TrackerPayload payload) { + payload.add(Parameter.PLATFORM, this.parameters.getPlatform().toString()); + payload.add(Parameter.APP_ID, this.parameters.getAppId()); + payload.add(Parameter.NAMESPACE, this.parameters.getNamespace()); + payload.add(Parameter.TRACKER_VERSION, this.parameters.getTrackerVersion()); + } +} diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerParameters.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerParameters.java new file mode 100644 index 00000000..f94abbf8 --- /dev/null +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerParameters.java @@ -0,0 +1,58 @@ +/* + * Copyright (c) 2020-2020 Snowplow Analytics Ltd. All rights reserved. + * + * This program is licensed to you under the Apache License Version 2.0, + * and you may not use this file except in compliance with the Apache License Version 2.0. + * You may obtain a copy of the Apache License Version 2.0 at http://www.apache.org/licenses/LICENSE-2.0. + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the Apache License Version 2.0 is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the Apache License Version 2.0 for the specific language governing permissions and limitations there under. + */ +package com.snowplowanalytics.snowplow.tracker.payload; + +import com.snowplowanalytics.snowplow.tracker.DevicePlatform; + +/** + * A TrackerEvent which allows the TrackerPayload to be filled later. The + * payload will be filled by the Emitter in the Emitter thread, using the + * getTrackerPayload() method. + */ +public class TrackerParameters { + + private final String trackerVersion; + private final String appId; + private final DevicePlatform platform; + private final String namespace; + private final boolean base64Encoded; + + public TrackerParameters(String appId, DevicePlatform platform, String namespace, String trackerVersion, + boolean base64Encoded) { + this.appId = appId; + this.platform = platform; + this.namespace = namespace; + this.trackerVersion = trackerVersion; + this.base64Encoded = base64Encoded; + } + + public boolean getBase64Encoded() { + return base64Encoded; + } + + public String getTrackerVersion() { + return trackerVersion; + } + + public String getNamespace() { + return namespace; + } + + public String getAppId() { + return appId; + } + + public DevicePlatform getPlatform() { + return platform; + } +} diff --git a/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerPayload.java b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerPayload.java index 08342b37..3a81ff1c 100644 --- a/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerPayload.java +++ b/src/main/java/com/snowplowanalytics/snowplow/tracker/payload/TrackerPayload.java @@ -12,16 +12,13 @@ */ package com.snowplowanalytics.snowplow.tracker.payload; -// Java import java.nio.charset.StandardCharsets; import java.util.LinkedHashMap; import java.util.Map; -// Slf4j import org.slf4j.Logger; import org.slf4j.LoggerFactory; -// This library import com.snowplowanalytics.snowplow.tracker.Utils; /** @@ -31,23 +28,22 @@ public class TrackerPayload implements Payload { private static final Logger LOGGER = LoggerFactory.getLogger(TrackerPayload.class); - private final LinkedHashMap payload = new LinkedHashMap<>(); + protected final Map payload = new LinkedHashMap<>(); /** - * Add a key-value pair to the payload: - * - Checks that the key is not null or empty - * - Checks that the value is not null or empty + * Add a key-value pair to the payload: - Checks that the key is not null or + * empty - Checks that the value is not null or empty * - * @param key The parameter key + * @param key The parameter key * @param value The parameter value as a String */ @Override - public void add(String key, String value) { + public void add(final String key, final String value) { if (key == null || key.isEmpty()) { LOGGER.error("Invalid key detected: {}", key); return; } - if (value == null || value.isEmpty()) { + if (value == null || value.isEmpty()) { LOGGER.info("null or empty value detected: {}", value); return; } @@ -56,40 +52,40 @@ public void add(String key, String value) { } /** - * Add all the mappings from the specified map. The effect is the equivalent to that of calling: - * - add(String key, String value) for each key value pair. + * Add all the mappings from the specified map. The effect is the equivalent to + * that of calling: - add(String key, String value) for each key value pair. * * @param map Key-Value pairs to be stored in this payload */ @Override - public void addMap(Map map) { + public void addMap(final Map map) { if (map == null) { LOGGER.debug("Map passed in is null, returning without adding map."); return; } LOGGER.debug("Adding new map: {}", map); - for (Map.Entry entry : map.entrySet()) { + for (final Map.Entry entry : map.entrySet()) { add(entry.getKey(), entry.getValue()); } } /** - * Add a map to the Payload with a key dependent on the base 64 encoding option you choose using the - * two keys provided. + * Add a map to the Payload with a key dependent on the base 64 encoding option + * you choose using the two keys provided. * - * @param map Map to be converted to a String and stored as a value - * @param base64Encoded The option you choose to encode the data - * @param typeEncoded The key that would be set if the encoding option was set to true + * @param map Map to be converted to a String and stored as a value + * @param base64Encoded The option you choose to encode the data + * @param typeEncoded The key that would be set if the encoding option was set to true * @param typeNotEncoded They key that would be set if the encoding option was set to false */ @Override - public void addMap(Map map, boolean base64Encoded, String typeEncoded, String typeNotEncoded) { + public void addMap(final Map map, final boolean base64Encoded, final String typeEncoded, final String typeNotEncoded) { if (map == null) { LOGGER.debug("Map passed in is null, returning nothing."); return; } - String mapString = Utils.mapToJSONString(map); + final String mapString = Utils.mapToJSONString(map); LOGGER.debug("Adding new map: {}", map); if (base64Encoded) { @@ -105,7 +101,7 @@ public void addMap(Map map, boolean base64Encoded, String typeEncoded, String ty * @return A Map of all the key-value entries */ @Override - public Map getMap() { + public Map getMap() { return payload; } @@ -120,8 +116,8 @@ public long getByteSize() { } /** - * Returns the Payload as a string. This is essentially the toString from the ObjectNode used - * to store the Payload. + * Returns the Payload as a string. This is essentially the toString from the + * ObjectNode used to store the Payload. * * @return A string value of the Payload. */ diff --git a/src/test/java/com/snowplowanalytics/snowplow/tracker/TrackerTest.java b/src/test/java/com/snowplowanalytics/snowplow/tracker/TrackerTest.java index bbc93ced..713aef74 100644 --- a/src/test/java/com/snowplowanalytics/snowplow/tracker/TrackerTest.java +++ b/src/test/java/com/snowplowanalytics/snowplow/tracker/TrackerTest.java @@ -12,32 +12,28 @@ */ package com.snowplowanalytics.snowplow.tracker; -// Java import java.util.*; import static java.util.Collections.singletonList; -// JUnit import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; + import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; -// Google import com.google.common.collect.ImmutableMap; -// Mockito import org.mockito.ArgumentCaptor; import org.mockito.Captor; import org.mockito.Mock; import org.mockito.runners.MockitoJUnitRunner; import static org.mockito.Mockito.*; -// This library import com.snowplowanalytics.snowplow.tracker.emitter.Emitter; import com.snowplowanalytics.snowplow.tracker.events.*; import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; -import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerEvent; @RunWith(MockitoJUnitRunner.class) public class TrackerTest { @@ -49,7 +45,7 @@ public class TrackerTest { Emitter emitter; @Captor - ArgumentCaptor captor; + ArgumentCaptor captor; Tracker tracker; private List contexts; @@ -102,10 +98,12 @@ public void testEcommerceEvent() { .build()); // Then - verify(emitter, times(2)).emit(captor.capture()); - List allValues = captor.getAllValues(); + verify(emitter, times(1)).emit(captor.capture()); + List allValues = captor.getAllValues(); - Map result1 = allValues.get(0).getMap(); + assertEquals(allValues.get(0).getTrackerPayloads().size(), 2); + + Map result1 = allValues.get(0).getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("e", "tr") .put("tr_cu", "currency") @@ -128,7 +126,7 @@ public void testEcommerceEvent() { .put("tr_st", "state") .build(), result1); - Map result2 = allValues.get(1).getMap(); + Map result2 = allValues.get(0).getTrackerPayloads().get(1).getMap(); assertEquals(ImmutableMap.builder() .put("ti_nm", "name") .put("ti_id", "order_id") @@ -167,7 +165,7 @@ public void testUnstructuredEventWithContext() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("p", "srv") .put("tv", Version.TRACKER) @@ -198,7 +196,7 @@ public void testUnstructuredEventWithoutContext() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("p", "srv") .put("tv", Version.TRACKER) @@ -227,7 +225,7 @@ public void testUnstructuredEventWithoutTrueTimestamp() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("p", "srv") .put("tv", Version.TRACKER) @@ -256,7 +254,7 @@ public void testTrackPageView() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("dtm", "123456") .put("ttm", "456789") @@ -274,6 +272,63 @@ public void testTrackPageView() { .build(), result); } + @Test + public void testTrackTwoEvents() { + // When + tracker.track(PageView.builder() + .pageUrl("url") + .pageTitle("title") + .referrer("referer") + .deviceCreatedTimestamp(123456) + .trueTimestamp(456789L) + .eventId("9783090a-dace-4c85-a75c-933b4596a6c5") + .build()); + + tracker.track(PageView.builder() + .pageUrl("url") + .pageTitle("title") + .referrer("referer") + .deviceCreatedTimestamp(123456) + .trueTimestamp(456789L) + .eventId("39139d43-ea13-4163-8559-adea258bf9c4") + .build()); + + // Then + verify(emitter, times(2)).emit(captor.capture()); + + Map result = captor.getAllValues().get(0).getTrackerPayloads().get(0).getMap(); + assertEquals(ImmutableMap.builder() + .put("dtm", "123456") + .put("ttm", "456789") + .put("tz", "Etc/UTC") + .put("e", "pv") + .put("page", "title") + .put("tv", Version.TRACKER) + .put("p", "srv") + .put("eid", "9783090a-dace-4c85-a75c-933b4596a6c5") + .put("tna", "AF003") + .put("aid", "cloudfront") + .put("refr", "referer") + .put("url", "url") + .build(), result); + + Map result2 = captor.getAllValues().get(1).getTrackerPayloads().get(0).getMap(); + assertEquals(ImmutableMap.builder() + .put("dtm", "123456") + .put("ttm", "456789") + .put("tz", "Etc/UTC") + .put("e", "pv") + .put("page", "title") + .put("tv", Version.TRACKER) + .put("p", "srv") + .put("eid", "39139d43-ea13-4163-8559-adea258bf9c4") + .put("tna", "AF003") + .put("aid", "cloudfront") + .put("refr", "referer") + .put("url", "url") + .build(), result2); + } + @Test public void testTrackScreenView() { // When @@ -288,7 +343,7 @@ public void testTrackScreenView() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("dtm", "123456") .put("ttm", "456789") @@ -317,7 +372,7 @@ public void testTrackScreenViewWithTimestamp() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("dtm", "123456") .put("ttm", "456789") @@ -346,7 +401,7 @@ public void testTrackScreenViewWithDefaultContextAndTimestamp() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("p", "srv") .put("tv", Version.TRACKER) @@ -378,7 +433,7 @@ public void testTrackTiming() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("p", "srv") .put("tv", Version.TRACKER) @@ -416,7 +471,7 @@ public void testTrackTimingWithSubject() { // Then verify(emitter).emit(captor.capture()); - Map result = captor.getValue().getMap(); + Map result = captor.getValue().getTrackerPayloads().get(0).getMap(); assertEquals(ImmutableMap.builder() .put("p", "srv") .put("ue_pr", "{\"schema\":\"iglu:com.snowplowanalytics.snowplow/unstruct_event/jsonschema/1-0-0\",\"data\":{\"schema\":\"iglu:com.snowplowanalytics.snowplow/timing/jsonschema/1-0-0\",\"data\":{\"category\":\"category\",\"label\":\"label\",\"timing\":10,\"variable\":\"variable\"}}}") @@ -447,8 +502,6 @@ public void testSetDefaultPlatform() throws Exception { .platform(DevicePlatform.Desktop) .build(); assertEquals(DevicePlatform.Desktop, tracker.getPlatform()); - tracker.setPlatform(DevicePlatform.ConnectedTV); - assertEquals(DevicePlatform.ConnectedTV, tracker.getPlatform()); } @Test @@ -473,23 +526,17 @@ public void testSetBase64Encoded() throws Exception { .base64(false) .build(); assertTrue(!tracker.getBase64Encoded()); - tracker.setBase64Encoded(true); - assertTrue(tracker.getBase64Encoded()); } @Test public void testSetAppId() throws Exception { Tracker tracker = new Tracker.TrackerBuilder(emitter, "AF003", "an-app-id").build(); assertEquals("an-app-id", tracker.getAppId()); - tracker.setAppId("cloudfront"); - assertEquals("cloudfront", tracker.getAppId()); } @Test public void testSetNamespace() throws Exception { Tracker tracker = new Tracker.TrackerBuilder(emitter, "namespace", "an-app-id").build(); assertEquals("namespace", tracker.getNamespace()); - tracker.setNamespace("cloudfront"); - assertEquals("cloudfront", tracker.getNamespace()); } } diff --git a/src/test/java/com/snowplowanalytics/snowplow/tracker/emitter/BatchEmitterTest.java b/src/test/java/com/snowplowanalytics/snowplow/tracker/emitter/BatchEmitterTest.java index 53c4e6cf..77b938d6 100644 --- a/src/test/java/com/snowplowanalytics/snowplow/tracker/emitter/BatchEmitterTest.java +++ b/src/test/java/com/snowplowanalytics/snowplow/tracker/emitter/BatchEmitterTest.java @@ -12,30 +12,31 @@ */ package com.snowplowanalytics.snowplow.tracker.emitter; -// Java import java.util.ArrayList; import java.util.List; import java.util.Map; -import java.util.UUID; -// Google import com.google.common.collect.Lists; -// JUnit import org.junit.Assert; import org.junit.Before; import org.junit.Rule; import org.junit.Test; import org.junit.rules.ExpectedException; -// Mockito import org.mockito.ArgumentCaptor; import static org.mockito.Mockito.*; -// This library +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.*; + +import com.snowplowanalytics.snowplow.tracker.DevicePlatform; import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerEvent; +import com.snowplowanalytics.snowplow.tracker.payload.TrackerParameters; import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload; import com.snowplowanalytics.snowplow.tracker.constants.Parameter; +import com.snowplowanalytics.snowplow.tracker.events.PageView; import com.snowplowanalytics.snowplow.tracker.http.HttpClientAdapter; public class BatchEmitterTest { @@ -56,51 +57,75 @@ public void setUp() throws Exception { } @Test - @SuppressWarnings("AssertEqualsBetweenInconvertibleTypes") - public void addToBuffer_withLess10Payloads_shouldNotFlushBuffer() throws Exception { + public void addToBuffer_withLess10Payloads_shouldNotEmptyBuffer() throws Exception { // Given ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(TrackerPayload.class); - List payloads = createPayloads(2); + List events = createEvents(2); // When - for (TrackerPayload payload : payloads) { - emitter.emit(payload); + for (TrackerEvent event : events) { + emitter.emit(event); } + Thread.sleep(500); + // Then - verify(emitter, never()).flushBuffer(); verify(httpClientAdapter, never()).get(argumentCaptor.capture()); Assert.assertEquals(2, emitter.getBuffer().size()); - Assert.assertEquals(payloads, emitter.getBuffer()); + Assert.assertEquals(events, emitter.getBuffer()); } @Test - @SuppressWarnings("AssertEqualsBetweenInconvertibleTypes") - public void addToBuffer_withMore10Payloads_shouldFlushBuffer() throws Exception { + public void addToBuffer_withMore10Payloads_shouldEmptyBuffer() throws Exception { // Given ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(SelfDescribingJson.class); - List payloads = createPayloads(10); + List events = createEvents(10); // When - for (TrackerPayload payload : payloads) { - emitter.emit(payload); + for (TrackerEvent event : events) { + emitter.emit(event); } Thread.sleep(500); // Then - verify(emitter).flushBuffer(); verify(httpClientAdapter).post(argumentCaptor.capture()); - List payloadMaps = new ArrayList<>(); - for (TrackerPayload payload : payloads) { - payloadMaps.add(payload.getMap()); + @SuppressWarnings("unchecked") + List> capturedPayload = (List>) argumentCaptor.getValue().getMap().get("data"); + + assertPayload(events, capturedPayload); + + Assert.assertEquals(0, emitter.getBuffer().size()); + } + + @Test + public void flushBuffer_shouldEmptyBuffer() throws Exception { + // Given + ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(SelfDescribingJson.class); + + List events = createEvents(2); + + // When + for (TrackerEvent event : events) { + emitter.emit(event); } - Assert.assertEquals(payloadMaps, argumentCaptor.getValue().getMap().get("data")); - Assert.assertTrue(emitter.getBuffer().size() == 0); + emitter.flushBuffer(); + + Thread.sleep(500); + + // Then + verify(httpClientAdapter).post(argumentCaptor.capture()); + + @SuppressWarnings("unchecked") + List> capturedPayload = (List>) argumentCaptor.getValue().getMap().get(Parameter.DATA); + + assertPayload(events, capturedPayload); + + Assert.assertEquals(0, emitter.getBuffer().size()); } @Test @@ -110,15 +135,14 @@ public void setBufferSize_WithNegativeValue_ThrowInvalidArgumentException() thro } @Test - @SuppressWarnings("unchecked") public void getFinalPost_shouldAddSTMParameter() throws Exception { // Given ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(SelfDescribingJson.class); - List payloads = createPayloads(10); + List events = createEvents(10); // When - for (TrackerPayload payload : payloads) { - emitter.emit(payload); + for (TrackerEvent event : events) { + emitter.emit(event); } Thread.sleep(500); @@ -126,23 +150,54 @@ public void getFinalPost_shouldAddSTMParameter() throws Exception { // Then verify(httpClientAdapter).post(argumentCaptor.capture()); - ArrayList> dataList = (ArrayList>) argumentCaptor.getValue().getMap().get(Parameter.DATA); - for (Map payloadMap : dataList) { + @SuppressWarnings("unchecked") + List> capturedPayload = (List>) argumentCaptor.getValue().getMap().get(Parameter.DATA); + + for (Map payloadMap : capturedPayload) { Assert.assertTrue(payloadMap.containsKey(Parameter.DEVICE_SENT_TIMESTAMP)); } } - private List createPayloads(int nbPayload) { - final List payloads = Lists.newArrayList(); - for (int i = 0; i < nbPayload; i++) { - payloads.add(createPayload()); + private List createEvents(int numEvents) { + final List payloads = Lists.newArrayList(); + for (int i = 0; i < numEvents; i++) { + payloads.add(createEvent()); } return payloads; } - private TrackerPayload createPayload() { - TrackerPayload payload = new TrackerPayload(); - payload.add("id", UUID.randomUUID().toString()); - return payload; + private TrackerEvent createEvent() { + PageView pv = PageView.builder() + .pageUrl("https://www.snowplowanalytics.com/") + .pageTitle("Snowplow") + .referrer("https://www.google.com/") + .build(); + + return new TrackerEvent(pv, new TrackerParameters("appId", DevicePlatform.ServerSideApp, "namespace", "0.0.0", false), null); + } + + private void assertPayload(List events, List> capturedPayload) { + List> eventPayloads = new ArrayList<>(); + for (TrackerEvent event : events) { + //All PageView events so we can get(0) from payloads + eventPayloads.add(event.getTrackerPayloads().get(0).getMap()); + } + + //Iterate through all captured payloads + for (Map capturedMap : capturedPayload) { + boolean matchFound = false; + for (Map eventMap : eventPayloads) { + //Find the matching events + if (capturedMap.get("eid") == eventMap.get("eid")) { + matchFound = true; + + //Assert that all the entries in the event are in the captured payload + //There might be extra entries in capturedMap, such as the STM parameter + //check for these addtional parameters in other tests + assertThat(eventMap.entrySet(), everyItem(is(in(capturedMap.entrySet())))); + } + } + assertThat(matchFound, is(true)); //Ensure every event was found + } } } diff --git a/src/test/java/com/snowplowanalytics/snowplow/tracker/http/HttpClientAdapterTest.java b/src/test/java/com/snowplowanalytics/snowplow/tracker/http/HttpClientAdapterTest.java index de437443..bad37a13 100644 --- a/src/test/java/com/snowplowanalytics/snowplow/tracker/http/HttpClientAdapterTest.java +++ b/src/test/java/com/snowplowanalytics/snowplow/tracker/http/HttpClientAdapterTest.java @@ -12,25 +12,20 @@ */ package com.snowplowanalytics.snowplow.tracker.http; -// Java import java.io.IOException; import java.util.Arrays; import java.util.Collection; import java.util.concurrent.TimeUnit; -// Google import com.google.common.collect.ImmutableMap; -// SquareUp import okhttp3.OkHttpClient; import okhttp3.mockwebserver.MockResponse; import okhttp3.mockwebserver.MockWebServer; import okhttp3.mockwebserver.RecordedRequest; -// Apache import org.apache.http.impl.client.HttpClients; -// JUnit import org.junit.Rule; import org.junit.Test; import org.junit.rules.ExpectedException; @@ -38,7 +33,6 @@ import org.junit.runners.Parameterized; import static org.junit.Assert.assertEquals; -// This library import com.snowplowanalytics.snowplow.tracker.payload.SelfDescribingJson; import com.snowplowanalytics.snowplow.tracker.payload.TrackerPayload;