From e5fce6c7f83fc1d912987cb880590cd312297920 Mon Sep 17 00:00:00 2001 From: Surabhi Date: Fri, 3 Dec 2021 10:06:51 +0530 Subject: [PATCH 1/3] lightstep opentelemetry - tracing on bot messages --- pom.xml | 3 +- .../com/uci/inbound/AppConfigInbound.java | 5 + .../uci/inbound/health/HealthController.java | 38 ++ .../netcore/NetcoreWhatsappConverter.java | 6 + .../uci/inbound/utils/XMsgProcessingUtil.java | 518 +++++++++++------- src/main/resources/application.properties | 7 + 6 files changed, 381 insertions(+), 196 deletions(-) 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/AppConfigInbound.java b/src/main/java/com/uci/inbound/AppConfigInbound.java index dfa101a..f08ab6c 100644 --- a/src/main/java/com/uci/inbound/AppConfigInbound.java +++ b/src/main/java/com/uci/inbound/AppConfigInbound.java @@ -1,9 +1,14 @@ package com.uci.inbound; +import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import com.uci.dao.service.HealthService; +import com.lightstep.opentelemetry.launcher.OpenTelemetryConfiguration; + +import io.opentelemetry.api.GlobalOpenTelemetry; +import io.opentelemetry.api.trace.Tracer; @Configuration public class AppConfigInbound { 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/netcore/NetcoreWhatsappConverter.java b/src/main/java/com/uci/inbound/netcore/NetcoreWhatsappConverter.java index ba5a200..b06e37c 100644 --- a/src/main/java/com/uci/inbound/netcore/NetcoreWhatsappConverter.java +++ b/src/main/java/com/uci/inbound/netcore/NetcoreWhatsappConverter.java @@ -6,6 +6,8 @@ import com.uci.inbound.utils.XMsgProcessingUtil; import com.uci.dao.repository.XMessageRepository; import com.uci.utils.kafka.SimpleProducer; + +import io.opentelemetry.api.trace.Tracer; import lombok.extern.slf4j.Slf4j; import com.uci.utils.BotService; import org.springframework.beans.factory.annotation.Autowired; @@ -42,6 +44,9 @@ public class NetcoreWhatsappConverter { @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 { @@ -60,6 +65,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..c332276 100644 --- a/src/main/java/com/uci/inbound/utils/XMsgProcessingUtil.java +++ b/src/main/java/com/uci/inbound/utils/XMsgProcessingUtil.java @@ -9,6 +9,16 @@ import com.uci.dao.utils.XMessageDAOUtils; import com.uci.utils.BotService; import com.uci.utils.kafka.SimpleProducer; + +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 +26,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 +43,312 @@ @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; + SimpleProducer kafkaProducer; + XMessageRepository xMsgRepo; + String topicSuccess; + String topicFailure; + BotService botService; + Tracer tracer; + + public void process() throws JsonProcessingException { + Span rootSpan = tracer.spanBuilder("inbound-processMessage").startSpan(); + + int data; + log.info("incoming message {}", new ObjectMapper().writeValueAsString(inboundMessage)); + try (Scope scope = rootSpan.makeCurrent()) { + Context currentContext = Context.current(); + Span childSpan1 = createChildSpan("convertMessageToXMsg", currentContext, rootSpan); + adapter.convertMessageToXMsg(inboundMessage) + .doOnError(genericError("convertMessageToXMsg", childSpan1)) + .subscribe(xmsg -> { + childSpan1.end(); + Span childSpan2 = createChildSpan("getAppName", currentContext, rootSpan); + getAppName(xmsg.getPayload().getText(), xmsg.getFrom()) + .doOnError(genericError("getAppName", childSpan2)) + .subscribe(appName -> { + childSpan2.end(); + xmsg.setApp(appName); + XMessageDAO currentMessageToBeInserted = XMessageDAOUtils + .convertXMessageToDAO(xmsg); + if (isCurrentMessageNotAReply(xmsg)) { + Span childSpan3 = createChildSpan("getLatestXMessage", currentContext, + rootSpan); + String whatsappId = xmsg.getMessageId().getChannelMessageId(); + getLatestXMessage(xmsg.getFrom().getUserID(), XMessage.MessageState.REPLIED) + .doOnError(genericError("getLatestXMessage", childSpan3)) + .subscribe(new Consumer() { + @Override + public void accept(XMessageDAO previousMessage) { + childSpan3.end(); + Span childSpan4 = createChildSpan( + "updatePreviousXMessage", currentContext, + rootSpan); + previousMessage.setMessageId(whatsappId); + xMsgRepo.save(previousMessage) + .doOnError(genericError( + "updatePreviousXMessage", childSpan4)) + .subscribe(new Consumer() { + @Override + public void accept( + XMessageDAO updatedPreviousMessage) { + childSpan4.end(); + Span childSpan5 = createChildSpan( + "insertXmessage", + currentContext, rootSpan); + xMsgRepo.insert(currentMessageToBeInserted) + .doOnError(genericError( + "insertXmessage", childSpan5)) + .subscribe(insertedMessage -> { + childSpan5.end(); + Span childSpan6 = createChildSpan( + "sendEventToKafka", + currentContext, rootSpan); +// log.info("current context l1: " +// + currentContext); + GlobalOpenTelemetry.getPropagators() + .getTextMapPropagator() + .inject(currentContext, + xmsg, null); + sendEventToKafka(xmsg); + childSpan6.end(); + rootSpan.end(); + }); + } + }); + } + }); + } else { + Span childSpan3 = createChildSpan("insertXmessage", currentContext, + rootSpan); + xMsgRepo.insert(currentMessageToBeInserted) + .doOnError( + genericError("insertXmessage", childSpan3)) + .subscribe(xMessageDAO -> { + childSpan3.end(); + Span childSpan4 = createChildSpan("sendEventToKafka", + currentContext, rootSpan); +// log.info("current context l2: " + currentContext); + GlobalOpenTelemetry.getPropagators().getTextMapPropagator() + .inject(currentContext, xmsg, null); + sendEventToKafka(xmsg); +// log.info("current context lc2: " + currentContext); + childSpan4.end(); + rootSpan.end(); + }); + } + }); + + }); + + } catch (JAXBException e) { + e.printStackTrace(); + genericException(e.getMessage(), rootSpan); + } catch (Throwable e) { + genericException(e.getMessage(), rootSpan); + } finally { +// rootSpan.end(); + } + } + + /** + * 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 = "inbound-"; + return tracer.spanBuilder(prefix + spanName).setParent(context.with(parentSpan)).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 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(); + } + }; + } + + 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("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/resources/application.properties b/src/main/resources/application.properties index 5d3367d..b6b1930 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -48,4 +48,11 @@ spring.r2dbc.password=${FORMS_DB_PASSWORD} caffeine.cache.max.size=${CAFFEINE_CACHE_MAX_SIZE:#{1000}} caffeine.cache.exprie.duration.seconds=${CAFFEINE_CACHE_EXPIRE_DURATION:#{300}} +#Opentelemetry Lighstep Config +opentelemetry.lightstep.tracer=${OPENTELEMETERY_LIGHTSTEP_TRACER} +opentelemetry.lightstep.tracer.version=${OPENTELEMETERY_LIGHTSTEP_TRACER_VERSION} +opentelemetry.lightstep.service=${OPENTELEMETERY_LIGHTSTEP_SERVICE} +opentelemetry.lightstep.access.token=${OPENTELEMETERY_LIGHTSTEP_ACCESS_TOKEN} +opentelemetry.lightstep.end.point=${OPENTELEMETERY_LIGHTSTEP_END_POINT} + From fa597100aa8481d3d0064ea68baf020d2d9d1695 Mon Sep 17 00:00:00 2001 From: Surabhi Date: Fri, 3 Dec 2021 10:10:04 +0530 Subject: [PATCH 2/3] lightstap tracer in gupshup & sunbird adapter builder --- src/main/java/com/uci/inbound/AppConfigInbound.java | 5 ----- .../com/uci/inbound/diksha/web/DikshaWebController.java | 6 ++++++ .../com/uci/inbound/incoming/GupShupWhatsappConverter.java | 7 +++++++ 3 files changed, 13 insertions(+), 5 deletions(-) diff --git a/src/main/java/com/uci/inbound/AppConfigInbound.java b/src/main/java/com/uci/inbound/AppConfigInbound.java index f08ab6c..dfa101a 100644 --- a/src/main/java/com/uci/inbound/AppConfigInbound.java +++ b/src/main/java/com/uci/inbound/AppConfigInbound.java @@ -1,14 +1,9 @@ package com.uci.inbound; -import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import com.uci.dao.service.HealthService; -import com.lightstep.opentelemetry.launcher.OpenTelemetryConfiguration; - -import io.opentelemetry.api.GlobalOpenTelemetry; -import io.opentelemetry.api.trace.Tracer; @Configuration public class AppConfigInbound { 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..0d1b130 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,8 @@ import com.uci.dao.repository.XMessageRepository; import com.uci.utils.BotService; import com.uci.utils.kafka.SimpleProducer; + +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; @@ -42,6 +44,9 @@ public class DikshaWebController { @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 +64,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/incoming/GupShupWhatsappConverter.java b/src/main/java/com/uci/inbound/incoming/GupShupWhatsappConverter.java index 8f33470..e70a54e 100644 --- a/src/main/java/com/uci/inbound/incoming/GupShupWhatsappConverter.java +++ b/src/main/java/com/uci/inbound/incoming/GupShupWhatsappConverter.java @@ -8,6 +8,9 @@ import com.uci.utils.BotService; import com.uci.inbound.utils.XMsgProcessingUtil; import com.uci.utils.kafka.SimpleProducer; + +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; @@ -44,6 +47,9 @@ public class GupShupWhatsappConverter { @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 +67,7 @@ public void gupshupWhatsApp(@Valid GSWhatsAppMessage message) throws JsonProcess .topicSuccess(inboundProcessed) .kafkaProducer(kafkaProducer) .botService(botService) + .tracer(tracer) .build() .process(); } From a2ec661444c9d6089688943b28232f3c1188b866 Mon Sep 17 00:00:00 2001 From: Surabhi Date: Tue, 14 Dec 2021 10:22:36 +0530 Subject: [PATCH 3/3] for log4j2 vulnaribility --- src/main/resources/application.properties | 1 + 1 file changed, 1 insertion(+) diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index b6b1930..a8b3474 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -55,4 +55,5 @@ opentelemetry.lightstep.service=${OPENTELEMETERY_LIGHTSTEP_SERVICE} opentelemetry.lightstep.access.token=${OPENTELEMETERY_LIGHTSTEP_ACCESS_TOKEN} opentelemetry.lightstep.end.point=${OPENTELEMETERY_LIGHTSTEP_END_POINT} +log4j2.formatMsgNoLookups=true