@@ -17,7 +17,6 @@ public class NtfyConnectionImpl implements NtfyConnection {
1717 private final String hostName ;
1818 private final ObjectMapper mapper = new ObjectMapper ();
1919
20-
2120 public NtfyConnectionImpl () {
2221 Dotenv dotenv = Dotenv .load ();
2322 hostName = Objects .requireNonNull (dotenv .get ("HOST_NAME" ));
@@ -31,7 +30,7 @@ public NtfyConnectionImpl(String hostName) {
3130 public CompletableFuture <Void > send (String message ) {
3231 HttpRequest httpRequest = HttpRequest .newBuilder ()
3332 .POST (HttpRequest .BodyPublishers .ofString (message ))
34- .header ("Cache" , "no" )
33+ .header ("Cache-Control " , "no" )
3534 .uri (URI .create (hostName + "/mytopic" ))
3635 .build ();
3736
@@ -52,9 +51,15 @@ public void receive (Consumer < NtfyMessageDto > messageHandler) {
5251
5352 http .sendAsync (httpRequest , HttpResponse .BodyHandlers .ofLines ())
5453 .thenAccept (response -> response .body ()
55- .map (s ->
56- mapper .readValue (s , NtfyMessageDto .class ))
57- .filter (message -> message .event ().equals ("message" ))
54+ .map (s -> {
55+ try {
56+ return mapper .readValue (s , NtfyMessageDto .class );
57+ } catch (Exception e ) {
58+ System .out .println ("Failed to parse message" );
59+ return null ;
60+ }
61+ })
62+ .filter (message -> message !=null && message .event ().equals ("message" ))
5863 .peek (System .out ::println )
5964 .forEach (messageHandler ));
6065 }
0 commit comments