diff --git a/pom.xml b/pom.xml
index f0a6e42..59341da 100644
--- a/pom.xml
+++ b/pom.xml
@@ -85,7 +85,7 @@
jaxb-impl
2.2.11
-
+
org.springframework.boot
spring-boot-starter-test
@@ -108,7 +108,6 @@
1.0
compile
-
diff --git a/src/main/java/com/uci/inbound/cdac/CDACConverter.java b/src/main/java/com/uci/inbound/cdac/CDACConverter.java
index 29da7e6..a429861 100644
--- a/src/main/java/com/uci/inbound/cdac/CDACConverter.java
+++ b/src/main/java/com/uci/inbound/cdac/CDACConverter.java
@@ -6,6 +6,8 @@
import com.uci.inbound.utils.XMsgProcessingUtil;
import com.uci.dao.repository.XMessageRepository;
import com.uci.utils.kafka.SimpleProducer;
+import com.uci.utils.kafka.RecordProducer;
+
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
@@ -38,7 +40,7 @@ public class CDACConverter {
private CdacBulkSmsAdapter cdacBulkSmsAdapter;
@Autowired
- public SimpleProducer kafkaProducer;
+ public RecordProducer kafkaProducer;
@Autowired
public XMessageRepository xmsgRepo;
diff --git a/src/main/java/com/uci/inbound/diksha/web/DikshaWebController.java b/src/main/java/com/uci/inbound/diksha/web/DikshaWebController.java
index 0f67fbe..100d068 100644
--- a/src/main/java/com/uci/inbound/diksha/web/DikshaWebController.java
+++ b/src/main/java/com/uci/inbound/diksha/web/DikshaWebController.java
@@ -8,6 +8,9 @@
import com.uci.dao.repository.XMessageRepository;
import com.uci.utils.BotService;
import com.uci.utils.kafka.SimpleProducer;
+import com.uci.utils.kafka.RecordProducer;
+
+import io.opentelemetry.api.trace.Tracer;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
@@ -35,13 +38,16 @@ public class DikshaWebController {
private SunbirdWebPortalAdapter sunbirdWebPortalAdapter;
@Autowired
- public SimpleProducer kafkaProducer;
+ public RecordProducer kafkaProducer;
@Autowired
public XMessageRepository xmsgRepo;
@Autowired
public BotService botService;
+
+ @Autowired
+ public Tracer tracer;
@RequestMapping(value = "/web", method = RequestMethod.POST, consumes = MediaType.APPLICATION_JSON_VALUE)
public void dikshaWeb(@RequestBody SunbirdWebMessage message) throws JsonProcessingException, JAXBException {
@@ -59,6 +65,7 @@ public void dikshaWeb(@RequestBody SunbirdWebMessage message) throws JsonProcess
.topicSuccess(inboundProcessed)
.kafkaProducer(kafkaProducer)
.botService(botService)
+ .tracer(tracer)
.build()
.process();
}
diff --git a/src/main/java/com/uci/inbound/health/HealthController.java b/src/main/java/com/uci/inbound/health/HealthController.java
index 6d9c3b2..66f89f4 100644
--- a/src/main/java/com/uci/inbound/health/HealthController.java
+++ b/src/main/java/com/uci/inbound/health/HealthController.java
@@ -6,6 +6,10 @@
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.uci.dao.service.HealthService;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.Tracer;
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.Scope;
import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
@@ -41,4 +45,38 @@ public ResponseEntity statusCheck() throws JsonProcessingException, IO
return ResponseEntity.ok(jsonNode);
}
+
+// @Autowired
+// private Tracer tracer;
+
+ @RequestMapping(value = "/test-lightstep", method = RequestMethod.GET, produces = { "application/json", "text/json" })
+ public ResponseEntity test() throws JsonProcessingException, IOException {
+
+// log.error("Health API called");
+//
+// Span span = tracer.spanBuilder("main-inbound-kafka").startSpan();
+//
+// int data;
+// try (Scope scope = span.makeCurrent()) {
+// Span childSpan = tracer.spanBuilder("child-inbound-kafka")
+// .setParent(Context.current().with(span))
+// .startSpan();
+//
+ ObjectMapper mapper = new ObjectMapper();
+
+ JsonNode jsonNode = mapper.readTree("{\"id\":\"api.content.health\",\"ver\":\"3.0\",\"ts\":\"2021-06-26T22:47:05Z+05:30\",\"params\":{\"resmsgid\":\"859fee0c-94d6-4a0d-b786-2025d763b78a\",\"msgid\":null,\"err\":null,\"status\":\"successful\",\"errmsg\":null},\"responseCode\":\"OK\",\"result\":{\"checks\":[{\"name\":\"redis cache\",\"healthy\":true},{\"name\":\"graph db\",\"healthy\":true},{\"name\":\"cassandra db\",\"healthy\":true}],\"healthy\":true}}");
+
+ int i = 0;
+ while(i <= 1000) {
+ log.info("Value of i: "+i);
+ i++;
+ }
+
+ return ResponseEntity.ok(jsonNode);
+// } finally {
+// span.end();
+// }
+
+
+ }
}
diff --git a/src/main/java/com/uci/inbound/incoming/GupShupWhatsappConverter.java b/src/main/java/com/uci/inbound/incoming/GupShupWhatsappConverter.java
index 8f33470..cad2851 100644
--- a/src/main/java/com/uci/inbound/incoming/GupShupWhatsappConverter.java
+++ b/src/main/java/com/uci/inbound/incoming/GupShupWhatsappConverter.java
@@ -8,6 +8,10 @@
import com.uci.utils.BotService;
import com.uci.inbound.utils.XMsgProcessingUtil;
import com.uci.utils.kafka.SimpleProducer;
+import com.uci.utils.kafka.RecordProducer;
+
+import io.opentelemetry.api.trace.Tracer;
+
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.http.MediaType;
@@ -37,13 +41,16 @@ public class GupShupWhatsappConverter {
private GupShupWhatsappAdapter gupShupWhatsappAdapter;
@Autowired
- public SimpleProducer kafkaProducer;
+ public RecordProducer kafkaProducer;
@Autowired
public XMessageRepository xmsgRepository;
@Autowired
public BotService botService;
+
+ @Autowired
+ public Tracer tracer;
@RequestMapping(value = "/whatsApp", method = RequestMethod.POST, consumes = MediaType.APPLICATION_FORM_URLENCODED_VALUE)
public void gupshupWhatsApp(@Valid GSWhatsAppMessage message) throws JsonProcessingException, JAXBException {
@@ -61,6 +68,7 @@ public void gupshupWhatsApp(@Valid GSWhatsAppMessage message) throws JsonProcess
.topicSuccess(inboundProcessed)
.kafkaProducer(kafkaProducer)
.botService(botService)
+ .tracer(tracer)
.build()
.process();
}
diff --git a/src/main/java/com/uci/inbound/netcore/NetcoreWhatsappConverter.java b/src/main/java/com/uci/inbound/netcore/NetcoreWhatsappConverter.java
index ba5a200..7ace01e 100644
--- a/src/main/java/com/uci/inbound/netcore/NetcoreWhatsappConverter.java
+++ b/src/main/java/com/uci/inbound/netcore/NetcoreWhatsappConverter.java
@@ -6,6 +6,9 @@
import com.uci.inbound.utils.XMsgProcessingUtil;
import com.uci.dao.repository.XMessageRepository;
import com.uci.utils.kafka.SimpleProducer;
+import com.uci.utils.kafka.RecordProducer;
+
+import io.opentelemetry.api.trace.Tracer;
import lombok.extern.slf4j.Slf4j;
import com.uci.utils.BotService;
import org.springframework.beans.factory.annotation.Autowired;
@@ -35,20 +38,20 @@ public class NetcoreWhatsappConverter {
private NetcoreWhatsappAdapter netcoreWhatsappAdapter;
@Autowired
- public SimpleProducer kafkaProducer;
+ public RecordProducer kafkaProducer;
@Autowired
public XMessageRepository xmsgRepo;
@Autowired
public BotService botService;
+
+ @Autowired
+ public Tracer tracer;
@RequestMapping(value = "/whatsApp", method = RequestMethod.POST, consumes = MediaType.APPLICATION_JSON_VALUE)
public void netcoreWhatsApp(@RequestBody NetcoreMessageFormat message) throws JsonProcessingException, JAXBException {
-
- System.out.println(message.toString());
-
- netcoreWhatsappAdapter = NetcoreWhatsappAdapter.builder()
+ netcoreWhatsappAdapter = NetcoreWhatsappAdapter.builder()
.botservice(botService)
.build();
@@ -60,6 +63,7 @@ public void netcoreWhatsApp(@RequestBody NetcoreMessageFormat message) throws Js
.topicSuccess(inboundProcessed)
.kafkaProducer(kafkaProducer)
.botService(botService)
+ .tracer(tracer)
.build()
.process();
}
diff --git a/src/main/java/com/uci/inbound/utils/XMsgProcessingUtil.java b/src/main/java/com/uci/inbound/utils/XMsgProcessingUtil.java
index 9cfc0ff..3ea3d68 100644
--- a/src/main/java/com/uci/inbound/utils/XMsgProcessingUtil.java
+++ b/src/main/java/com/uci/inbound/utils/XMsgProcessingUtil.java
@@ -9,6 +9,17 @@
import com.uci.dao.utils.XMessageDAOUtils;
import com.uci.utils.BotService;
import com.uci.utils.kafka.SimpleProducer;
+import com.uci.utils.kafka.RecordProducer;
+
+import io.opentelemetry.api.GlobalOpenTelemetry;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.StatusCode;
+import io.opentelemetry.api.trace.Tracer;
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.Scope;
+import io.opentelemetry.context.propagation.ContextPropagators;
+import io.opentelemetry.context.propagation.TextMapPropagator;
+import io.opentelemetry.context.propagation.TextMapSetter;
import lombok.Builder;
import lombok.extern.slf4j.Slf4j;
import messagerosa.core.model.SenderReceiverInfo;
@@ -16,10 +27,16 @@
import reactor.core.publisher.Mono;
import javax.xml.bind.JAXBException;
+
+import org.slf4j.Logger;
+import org.springframework.beans.factory.annotation.Autowired;
+
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Comparator;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
import java.util.function.Consumer;
import java.util.function.Function;
@@ -27,198 +44,310 @@
@Builder
public class XMsgProcessingUtil {
- AbstractProvider adapter;
- CommonMessage inboundMessage;
- SimpleProducer kafkaProducer;
- XMessageRepository xMsgRepo;
- String topicSuccess;
- String topicFailure;
- BotService botService;
-
-
- public void process() throws JsonProcessingException {
-
- log.info("incoming message {}", new ObjectMapper().writeValueAsString(inboundMessage));
- try {
- adapter.convertMessageToXMsg(inboundMessage)
- .doOnError(genericError("Error in converting to XMessage by Adapter"))
- .subscribe(xmsg -> {
- getAppName(xmsg.getPayload().getText(), xmsg.getFrom())
- .subscribe(appName -> {
- xmsg.setApp(appName);
- XMessageDAO currentMessageToBeInserted = XMessageDAOUtils.convertXMessageToDAO(xmsg);
- if (isCurrentMessageNotAReply(xmsg)) {
- String whatsappId = xmsg.getMessageId().getChannelMessageId();
- getLatestXMessage(xmsg.getFrom().getUserID(), XMessage.MessageState.REPLIED)
- .doOnError(genericError("Error in getting last message"))
- .subscribe(new Consumer() {
- @Override
- public void accept(XMessageDAO previousMessage) {
- previousMessage.setMessageId(whatsappId);
- xMsgRepo.save(previousMessage)
- .doOnError(genericError("Error in saving previous message"))
- .subscribe(new Consumer() {
- @Override
- public void accept(XMessageDAO updatedPreviousMessage) {
- xMsgRepo.insert(currentMessageToBeInserted)
- .doOnError(genericError("Error in inserting current message"))
- .subscribe(insertedMessage -> {
- sendEventToKafka(xmsg);
- });
- }
- });
- }
- });
- } else {
- xMsgRepo.insert(currentMessageToBeInserted)
- .doOnError(genericError("Error in inserting current message"))
- .subscribe(xMessageDAO -> {
- sendEventToKafka(xmsg);
- });
- }
- });
-
- });
-
- } catch (JAXBException e) {
- log.info("Error Message: "+e.getMessage());
- e.printStackTrace();
- }
- }
-
- private Consumer genericError(String s) {
- return c -> {
- log.error(s + "::" + c.getMessage());
- };
- }
-
- private boolean isCurrentMessageNotAReply(XMessage xmsg) {
- return !xmsg.getMessageState().equals(XMessage.MessageState.REPLIED);
- }
-
- private void sendEventToKafka(XMessage xmsg) {
- String xmessage = null;
- try {
- xmessage = xmsg.toXML();
- } catch (JAXBException e) {
- kafkaProducer.send(topicFailure, inboundMessage.toString());
- }
- kafkaProducer.send(topicSuccess, xmessage);
- }
-
- private Mono getLatestXMessage(String userID, XMessage.MessageState messageState) {
- LocalDateTime yesterday = LocalDateTime.now().minusDays(1L);
- return xMsgRepo
- .findAllByFromIdAndTimestampAfter(userID, yesterday)
- .doOnError(genericError(String.format("Unable to find previous Message for userID %s", userID)))
- .collectList()
- .map(xMessageDAOS -> {
- if (xMessageDAOS.size() > 0) {
- List filteredList = new ArrayList<>();
- for (XMessageDAO xMessageDAO : xMessageDAOS) {
- if (xMessageDAO.getMessageState().equals(messageState.name()))
- filteredList.add(xMessageDAO);
- }
- if (filteredList.size() > 0) {
- filteredList.sort(Comparator.comparing(XMessageDAO::getTimestamp));
- }
-
- return xMessageDAOS.get(0);
- }
- return new XMessageDAO();
- });
- }
-
- private Mono getAppName(String text, SenderReceiverInfo from) {
- LocalDateTime yesterday = LocalDateTime.now().minusDays(1L);
- if (text.equals("")) {
- try {
- return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name()).map(new Function() {
- @Override
- public String apply(XMessageDAO xMessageLast) {
- return xMessageLast.getApp();
- }
- }).doOnError(genericError("Error in getting latest xmessage"));
- } catch (Exception e2) {
- return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name()).map(new Function() {
- @Override
- public String apply(XMessageDAO xMessageLast) {
- return xMessageLast.getApp();
- }
- }).doOnError(genericError("Error in getting latest xmessage - catch"));
- }
- } else {
- try {
- log.error("getCampaignFromStartingMessage text: "+text);
- return botService.getCampaignFromStartingMessage(text)
- .flatMap(new Function>() {
- @Override
- public Mono apply(String appName1) {
- if (appName1 == null || appName1.equals("")) {
- try {
- return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name()).map(new Function() {
- @Override
- public String apply(XMessageDAO xMessageLast) {
- return (xMessageLast.getApp() == null || xMessageLast.getApp().isEmpty()) ? "finalAppName" : xMessageLast.getApp();
- }
- }).doOnError(genericError("Error in getting latest xmessage when app name empty"));
- } catch (Exception e2) {
- return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name()).map(new Function() {
- @Override
- public String apply(XMessageDAO xMessageLast) {
- return (xMessageLast.getApp() == null || xMessageLast.getApp().isEmpty()) ? "finalAppName" : xMessageLast.getApp();
- }
- }).doOnError(genericError("Error in getting latest xmessage when app name empty - catch"));
- }
- }
- return (appName1 == null || appName1.isEmpty()) ? Mono.just("finalAppName") : Mono.just(appName1);
- }
- });
- } catch (Exception e) {
- log.error("Exception in getCampaignFromStartingMessage :"+e.getMessage());
- try {
- return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name()).map(new Function() {
- @Override
- public String apply(XMessageDAO xMessageLast) {
- return xMessageLast.getApp();
- }
- }).doOnError(genericError("Error in getting latest xmessage when exception in getCampaignFromStartingMessage"));
- } catch (Exception e2) {
- return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name()).map(new Function() {
- @Override
- public String apply(XMessageDAO xMessageLast) {
- return xMessageLast.getApp();
- }
- }).doOnError(genericError("Error in getting latest xmessage when exception in getCampaignFromStartingMessage - catch"));
- }
- }
- }
- }
-
- private Mono getLatestXMessage(String userID, LocalDateTime yesterday, String messageState) {
- return xMsgRepo.findAllByUserIdAndTimestampAfter(userID, yesterday)
- .collectList()
- .map(new Function, XMessageDAO>() {
- @Override
- public XMessageDAO apply(List xMessageDAOS) {
- if (xMessageDAOS.size() > 0) {
- List filteredList = new ArrayList<>();
- for (XMessageDAO xMessageDAO : xMessageDAOS) {
- if (xMessageDAO.getMessageState().equals(XMessage.MessageState.SENT.name()))
- filteredList.add(xMessageDAO);
- }
- if (filteredList.size() > 0) {
- filteredList.sort(new Comparator() {
- @Override
- public int compare(XMessageDAO o1, XMessageDAO o2) {
- return o1.getTimestamp().compareTo(o2.getTimestamp());
- }
- });
- }
- return xMessageDAOS.get(0);
- }
- return new XMessageDAO();
- }
- });
- }
+ AbstractProvider adapter;
+ CommonMessage inboundMessage;
+ RecordProducer kafkaProducer;
+ XMessageRepository xMsgRepo;
+ String topicSuccess;
+ String topicFailure;
+ BotService botService;
+ Tracer tracer;
+
+ public void process() throws JsonProcessingException {
+
+ int data;
+ log.info("incoming message {}", new ObjectMapper().writeValueAsString(inboundMessage));
+ try {
+ adapter.convertMessageToXMsg(inboundMessage)
+ .doOnError(genericError("convertMessageToXMsg"))
+ .subscribe(xmsg -> {
+ Span childSpan1 = createChildSpan("getAppName");
+ getAppName(xmsg.getPayload().getText(), xmsg.getFrom())
+ .doOnError(genericError("getAppName"))
+ .subscribe(appName -> {
+ childSpan1.end();
+ xmsg.setApp(appName);
+ Span childSpan2 = createChildSpan("convertXMessageToDAO");
+ XMessageDAO currentMessageToBeInserted = XMessageDAOUtils
+ .convertXMessageToDAO(xmsg);
+ childSpan2.end();
+ if (isCurrentMessageNotAReply(xmsg)) {
+ String whatsappId = xmsg.getMessageId().getChannelMessageId();
+ getLatestXMessage(xmsg.getFrom().getUserID(), XMessage.MessageState.REPLIED)
+ .doOnError(genericError("getLatestXMessage"))
+ .subscribe(new Consumer() {
+ @Override
+ public void accept(XMessageDAO previousMessage) {
+ previousMessage.setMessageId(whatsappId);
+ xMsgRepo.save(previousMessage)
+ .doOnError(genericError(
+ "updatePreviousXMessage"))
+ .subscribe(new Consumer() {
+ @Override
+ public void accept(
+ XMessageDAO updatedPreviousMessage) {
+ xMsgRepo.insert(currentMessageToBeInserted)
+ .doOnError(genericError(
+ "insertXmessage"))
+ .subscribe(insertedMessage -> {
+ sendEventToKafka(xmsg, Context.current());
+ });
+ }
+ });
+ }
+ });
+ } else {
+ Span childSpan3 = createChildSpan("insertXmessage");
+ xMsgRepo.insert(currentMessageToBeInserted)
+ .doOnError(
+ genericError("insertXmessage"))
+ .subscribe(xMessageDAO -> {
+ childSpan3.end();
+ Span childSpan4 = createChildSpan("sendEventToKafka");
+ sendEventToKafka(xmsg, Context.current());
+ childSpan4.end();
+ });
+ }
+ });
+
+ });
+
+ } catch (JAXBException e) {
+ e.printStackTrace();
+ genericException(e.getMessage());
+ } catch (Throwable e) {
+ genericException(e.getMessage());
+ }
+ }
+
+ /**
+ * Create Child Span with current context & parent span
+ * @param spanName
+ * @param context
+ * @param parentSpan
+ * @return childSpan
+ */
+ private Span createChildSpan(String spanName, Context context, Span parentSpan) {
+ String prefix = "inboundSpan-";
+ return tracer.spanBuilder(prefix + spanName).setParent(context.with(parentSpan)).startSpan();
+ }
+
+ /**
+ * Create Child Span
+ * @param spanName
+ * @return childSpan
+ */
+ private Span createChildSpan(String spanName) {
+ String prefix = "inbound-";
+ return tracer.spanBuilder(prefix + spanName).startSpan();
+ }
+
+ private void propagateContext(Context currectContext, XMessage xmsg) {
+ log.info("current context: " + currectContext);
+// ContextPropagators propagators = GlobalOpenTelemetry.getPropagators();
+// TextMapPropagator textMapPropagator = propagators.getTextMapPropagator();
+
+ Map map = new HashMap();
+ map.put("from", "inbound");
+ GlobalOpenTelemetry.getPropagators().getTextMapPropagator().inject(currectContext, xmsg, null);
+ }
+
+ /**
+ * Log Exceptions & if span exists, add error to span
+ * @param eMsg
+ * @param span
+ */
+ private void genericException(String eMsg, Span span) {
+ eMsg = "Exception: " + eMsg;
+ log.error(eMsg);
+ if(span != null) {
+ span.setStatus(StatusCode.ERROR, "Exception: " + eMsg);
+ span.end();
+ }
+ }
+
+ /**
+ * Log Exceptions
+ * @param eMsg
+ */
+ private void genericException(String eMsg) {
+ eMsg = "Exception: " + eMsg;
+ log.error(eMsg);
+ }
+
+ /**
+ * Log Exception & if span exists, add error to span
+ * @param s
+ * @param span
+ * @return
+ */
+ private Consumer genericError(String s, Span span) {
+ return c -> {
+ String msg = "Error in " + s + "::" + c.getMessage();
+ log.error(msg);
+ if (span != null) {
+ log.info("generic message - span");
+ span.setStatus(StatusCode.ERROR, msg);
+ span.end();
+ }
+ };
+ }
+
+ /**
+ * Log Exception
+ * @param s
+ * @return
+ */
+ private Consumer genericError(String s) {
+ return c -> {
+ String msg = "Error in " + s + "::" + c.getMessage();
+ log.error(msg);
+ };
+ }
+
+ private boolean isCurrentMessageNotAReply(XMessage xmsg) {
+ return !xmsg.getMessageState().equals(XMessage.MessageState.REPLIED);
+ }
+
+ private void sendEventToKafka(XMessage xmsg, Context currentContext) {
+ String xmessage = null;
+ try {
+ xmessage = xmsg.toXML();
+ } catch (JAXBException e) {
+ kafkaProducer.send(topicFailure, inboundMessage.toString(), currentContext);
+ }
+ kafkaProducer.send(topicSuccess, xmessage, currentContext);
+ }
+
+ private Mono getLatestXMessage(String userID, XMessage.MessageState messageState) {
+ LocalDateTime yesterday = LocalDateTime.now().minusDays(1L);
+ return xMsgRepo.findAllByFromIdAndTimestampAfter(userID, yesterday)
+ .doOnError(genericError(String.format("finding previous Message for userID %s", userID), null))
+ .collectList().map(xMessageDAOS -> {
+ if (xMessageDAOS.size() > 0) {
+ List filteredList = new ArrayList<>();
+ for (XMessageDAO xMessageDAO : xMessageDAOS) {
+ if (xMessageDAO.getMessageState().equals(messageState.name()))
+ filteredList.add(xMessageDAO);
+ }
+ if (filteredList.size() > 0) {
+ filteredList.sort(Comparator.comparing(XMessageDAO::getTimestamp));
+ }
+
+ return xMessageDAOS.get(0);
+ }
+ return new XMessageDAO();
+ });
+ }
+
+ private Mono getAppName(String text, SenderReceiverInfo from) {
+ LocalDateTime yesterday = LocalDateTime.now().minusDays(1L);
+ if (text.equals("")) {
+ try {
+ return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name())
+ .map(new Function() {
+ @Override
+ public String apply(XMessageDAO xMessageLast) {
+ return xMessageLast.getApp();
+ }
+ }).doOnError(genericError("getLatestXMessage", null));
+ } catch (Exception e2) {
+ return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name())
+ .map(new Function() {
+ @Override
+ public String apply(XMessageDAO xMessageLast) {
+ return xMessageLast.getApp();
+ }
+ }).doOnError(genericError("getLatestXMessage - catch", null));
+ }
+ } else {
+ try {
+ log.error("getCampaignFromStartingMessage text: " + text);
+ return botService.getCampaignFromStartingMessage(text)
+ .flatMap(new Function>() {
+ @Override
+ public Mono apply(String appName1) {
+ if (appName1 == null || appName1.equals("")) {
+ try {
+ return getLatestXMessage(from.getUserID(), yesterday,
+ XMessage.MessageState.SENT.name())
+ .map(new Function() {
+ @Override
+ public String apply(XMessageDAO xMessageLast) {
+ return (xMessageLast.getApp() == null
+ || xMessageLast.getApp().isEmpty())
+ ? "finalAppName"
+ : xMessageLast.getApp();
+ }
+ }).doOnError(genericError(
+ "getLatestXMessage when appName empty", null));
+ } catch (Exception e2) {
+ return getLatestXMessage(from.getUserID(), yesterday,
+ XMessage.MessageState.SENT.name())
+ .map(new Function() {
+ @Override
+ public String apply(XMessageDAO xMessageLast) {
+ return (xMessageLast.getApp() == null
+ || xMessageLast.getApp().isEmpty())
+ ? "finalAppName"
+ : xMessageLast.getApp();
+ }
+ }).doOnError(genericError(
+ "getLatestXMessage when appName empty - catch", null));
+ }
+ }
+ return (appName1 == null || appName1.isEmpty()) ? Mono.just("finalAppName")
+ : Mono.just(appName1);
+ }
+ });
+ } catch (Exception e) {
+ log.error("Exception in getCampaignFromStartingMessage :" + e.getMessage());
+ try {
+ return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name())
+ .map(new Function() {
+ @Override
+ public String apply(XMessageDAO xMessageLast) {
+ return xMessageLast.getApp();
+ }
+ }).doOnError(genericError(
+ "getLatestXMessage when exception in getCampaignFromStartingMessage", null));
+ } catch (Exception e2) {
+ return getLatestXMessage(from.getUserID(), yesterday, XMessage.MessageState.SENT.name())
+ .map(new Function() {
+ @Override
+ public String apply(XMessageDAO xMessageLast) {
+ return xMessageLast.getApp();
+ }
+ }).doOnError(genericError(
+ "getLatestXMessage when exception in getCampaignFromStartingMessage - catch", null));
+ }
+ }
+ }
+ }
+
+ private Mono getLatestXMessage(String userID, LocalDateTime yesterday, String messageState) {
+ return xMsgRepo.findAllByUserIdAndTimestampAfter(userID, yesterday).collectList()
+ .map(new Function, XMessageDAO>() {
+ @Override
+ public XMessageDAO apply(List xMessageDAOS) {
+ if (xMessageDAOS.size() > 0) {
+ List filteredList = new ArrayList<>();
+ for (XMessageDAO xMessageDAO : xMessageDAOS) {
+ if (xMessageDAO.getMessageState().equals(XMessage.MessageState.SENT.name()))
+ filteredList.add(xMessageDAO);
+ }
+ if (filteredList.size() > 0) {
+ filteredList.sort(new Comparator() {
+ @Override
+ public int compare(XMessageDAO o1, XMessageDAO o2) {
+ return o1.getTimestamp().compareTo(o2.getTimestamp());
+ }
+ });
+ }
+ return xMessageDAOS.get(0);
+ }
+ return new XMessageDAO();
+ }
+ });
+ }
}
\ No newline at end of file
diff --git a/src/main/java/com/uci/inbound/xmsg/XMessageController.java b/src/main/java/com/uci/inbound/xmsg/XMessageController.java
index 213f98b..9bfa6b5 100644
--- a/src/main/java/com/uci/inbound/xmsg/XMessageController.java
+++ b/src/main/java/com/uci/inbound/xmsg/XMessageController.java
@@ -77,4 +77,12 @@ public void deleteAllByUserIdBeforeHours(@PathVariable("userid") String userid)
xMsgRepo.delete(xMsgDao).subscribe();
});
}
+
+ @RequestMapping(value = "/user/dataByUserId/{userId}", method = RequestMethod.DELETE, produces = {
+ "application/json", "text/json" })
+ public void deleteAllByMobile(@PathVariable("userId") String userid) {
+ xMsgRepo.findAllByUserId(userid).subscribe(xMsgDao -> {
+ xMsgRepo.delete(xMsgDao).subscribe();
+ });
+ }
}
diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties
index 25244e0..d668d4b 100644
--- a/src/main/resources/application.properties
+++ b/src/main/resources/application.properties
@@ -49,4 +49,11 @@ spring.r2dbc.password=${FORMS_DB_PASSWORD}
caffeine.cache.max.size=0
caffeine.cache.exprie.duration.seconds=${CAFFEINE_CACHE_EXPIRE_DURATION:#{300}}
+#Opentelemetry Lighstep Config
+opentelemetry.lightstep.tracer=${LS_TRACER_NAME}
+opentelemetry.lightstep.tracer.version=${LS_TRACER_VERSION}
+opentelemetry.lightstep.service=${LS_SERVICE_NAME}
+opentelemetry.lightstep.access.token=${LS_ACCESS_TOKEN}
+opentelemetry.lightstep.end.point=${LS_END_POINT}
+