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} +