Skip to content

Commit 1749b17

Browse files
committed
- Added failsSafeXXX() methods to Log
- Use of log.failsSafeXXX() calls to catch exceptions (https://redhat.atlassian.net/browse/JGRP-3034)
1 parent d5868f7 commit 1749b17

7 files changed

Lines changed: 125 additions & 58 deletions

File tree

src/org/jgroups/logging/Log.java

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,4 +54,49 @@ public interface Log {
5454
void setLevel(String level);
5555

5656
String getLevel();
57+
58+
default void failSafeWarn(String msg) {
59+
try {
60+
warn(msg);
61+
}
62+
catch(Throwable t) {
63+
}
64+
}
65+
default void failSafeWarn(String format, Object ... args) {
66+
try {
67+
warn(format, args);
68+
}
69+
catch(Throwable t) {
70+
}
71+
}
72+
73+
default void failSafeWarn(String msg, Throwable throwable) {
74+
try {
75+
warn(msg, throwable);
76+
}
77+
catch(Throwable t) {
78+
}
79+
}
80+
81+
default void failSafeError(String msg) {
82+
try {
83+
error(msg);
84+
}
85+
catch(Throwable t) {}
86+
}
87+
88+
default void failSafeError(String format, Object ... args) {
89+
try {
90+
error(format, args);
91+
}
92+
catch(Throwable t) {
93+
}
94+
}
95+
96+
default void failSafeError(String msg, Throwable throwable) {
97+
try {
98+
error(msg, throwable);
99+
}
100+
catch(Throwable t) {}
101+
}
57102
}

src/org/jgroups/protocols/ReliableMulticast.java

Lines changed: 19 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -842,18 +842,25 @@ protected void removeAndDeliver(Buffer<Message> win, Entry e, Address sender, bo
842842
AtomicInteger adders=win.getAdders();
843843
if(adders.getAndIncrement() != 0)
844844
return;
845-
boolean remove_msgs=!loopback;
846-
int cap=max_batch_size > 0 && max_batch_size < win.capacity()? max_batch_size : win.capacity();
847-
AsciiString cl=cluster != null? cluster : getTransport().getClusterNameAscii();
845+
848846
MessageBatch b=null;
849-
if(reuse_message_batches) {
850-
b=cached_batches.get(sender);
851-
if(b == null)
852-
b=cached_batches.computeIfAbsent(sender, __ -> new MessageBatch(cap).dest(null).sender(sender)
853-
.cluster(cl).mcast(true));
847+
boolean remove_msgs=!loopback;
848+
try {
849+
int cap=max_batch_size > 0 && max_batch_size < win.capacity()? max_batch_size : win.capacity();
850+
AsciiString cl=cluster != null? cluster : getTransport().getClusterNameAscii();
851+
if(reuse_message_batches) {
852+
b=cached_batches.get(sender);
853+
if(b == null)
854+
b=cached_batches.computeIfAbsent(sender, __ -> new MessageBatch(cap).dest(null).sender(sender)
855+
.cluster(cl).mcast(true));
856+
}
857+
else
858+
b=new MessageBatch(cap).dest(null).sender(sender).cluster(cl).mcast(true);
859+
}
860+
catch(Throwable t) {
861+
adders.set(0); // so others can remove/deliver msgs (https://redhat.atlassian.net/browse/JGRP-3034)
862+
throw t;
854863
}
855-
else
856-
b=new MessageBatch(cap).dest(null).sender(sender).cluster(cl).mcast(true);
857864
MessageBatch batch=b;
858865
Supplier<MessageBatch> batch_creator=() -> batch;
859866
MessageBatch mb=null;
@@ -866,7 +873,7 @@ protected void removeAndDeliver(Buffer<Message> win, Entry e, Address sender, bo
866873
batch.determineMode();
867874
}
868875
catch(Throwable t) {
869-
log.error("failed removing messages from table for " + sender, t);
876+
log.failSafeError("failed removing messages from table for " + sender, t);
870877
}
871878
int size=batch.size();
872879
if(size > 0) {
@@ -949,7 +956,7 @@ protected void deliverBatch(MessageBatch batch, Entry entry) {
949956
batch.reset(); // doesn't null messages in the batch
950957
}
951958
catch(Throwable t) {
952-
log.error(Util.getMessage("FailedToDeliverMsg"), local_addr, "batch", batch, t);
959+
log.failSafeError(Util.getMessage("FailedToDeliverMsg"), local_addr, "batch", batch, t);
953960
}
954961
}
955962

src/org/jgroups/protocols/ReliableUnicast.java

Lines changed: 17 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -892,18 +892,23 @@ protected void removeAndDeliver(Entry entry, Address sender, AsciiString cluster
892892
AtomicInteger adders=buf.getAdders();
893893
if(adders.getAndIncrement() != 0)
894894
return;
895-
896-
AsciiString cl=cluster != null? cluster : getTransport().getClusterNameAscii();
897-
int cap=Math.max(Math.max(Math.max(buf.size(), max_batch_size), min_size), DEFAULT_INITIAL_CAPACITY);
898895
MessageBatch b=null;
899-
if(reuse_message_batches) {
900-
b=cached_batches.get(sender);
901-
if(b == null)
902-
b=cached_batches.computeIfAbsent(sender, __ -> new MessageBatch(cap).dest(local_addr)
903-
.sender(sender).cluster(cl).incr(DEFAULT_INCREMENT));
896+
try {
897+
AsciiString cl=cluster != null? cluster : getTransport().getClusterNameAscii();
898+
int cap=Math.max(Math.max(Math.max(buf.size(), max_batch_size), min_size), DEFAULT_INITIAL_CAPACITY);
899+
if(reuse_message_batches) {
900+
b=cached_batches.get(sender);
901+
if(b == null)
902+
b=cached_batches.computeIfAbsent(sender, __ -> new MessageBatch(cap).dest(local_addr)
903+
.sender(sender).cluster(cl).incr(DEFAULT_INCREMENT));
904+
}
905+
else
906+
b=new MessageBatch(cap).dest(local_addr).sender(sender).cluster(cl).incr(DEFAULT_INCREMENT);
907+
}
908+
catch(Throwable t) {
909+
adders.set(0); // so others can remove/deliver msgs (https://redhat.atlassian.net/browse/JGRP-3034)
910+
throw t;
904911
}
905-
else
906-
b=new MessageBatch(cap).dest(local_addr).sender(sender).cluster(cl).incr(DEFAULT_INCREMENT);
907912
MessageBatch batch=b;
908913
Supplier<MessageBatch> batch_creator=() -> batch;
909914
MessageBatch mb=null;
@@ -914,7 +919,7 @@ protected void removeAndDeliver(Entry entry, Address sender, AsciiString cluster
914919
batch_creator, BATCH_ACCUMULATOR);
915920
}
916921
catch(Throwable t) {
917-
log.error("%s: failed removing messages from table for %s: %s", local_addr, sender, t);
922+
log.failSafeError("%s: failed removing messages from table for %s: %s", local_addr, sender, t);
918923
}
919924
if(!batch.isEmpty()) {
920925
// batch is guaranteed to NOT contain any OOB messages as the drop_oob_msgs_filter above removed them
@@ -1188,7 +1193,7 @@ protected void deliverBatch(MessageBatch batch, Entry entry, Address original_de
11881193
avg_delivery_batch_size.add(batch.size());
11891194
}
11901195
catch(Throwable t) {
1191-
log.warn(Util.getMessage("FailedToDeliverMsg"), local_addr, "batch", batch, t);
1196+
log.failSafeWarn(Util.getMessage("FailedToDeliverMsg"), local_addr, "batch", batch, t);
11921197
}
11931198
}
11941199

src/org/jgroups/protocols/UNICAST3.java

Lines changed: 19 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -975,18 +975,23 @@ protected void removeAndDeliver(Table<Message> win, Address sender, AsciiString
975975
AtomicInteger adders=win.getAdders();
976976
if(adders.getAndIncrement() != 0)
977977
return;
978-
979-
AsciiString cl=cluster != null? cluster : getTransport().getClusterNameAscii();
980-
int cap=Math.max(Math.max(Math.max(win.size(), max_batch_size), min_size), DEFAULT_INITIAL_CAPACITY);
981978
MessageBatch b=null;
982-
if(reuse_message_batches) {
983-
b=cached_batches.get(sender);
984-
if(b == null)
985-
b=cached_batches.computeIfAbsent(sender, __ -> new MessageBatch(cap).dest(local_addr)
986-
.sender(sender).cluster(cl).incr(DEFAULT_INCREMENT));
979+
try {
980+
AsciiString cl=cluster != null? cluster : getTransport().getClusterNameAscii();
981+
int cap=Math.max(Math.max(Math.max(win.size(), max_batch_size), min_size), DEFAULT_INITIAL_CAPACITY);
982+
if(reuse_message_batches) {
983+
b=cached_batches.get(sender);
984+
if(b == null)
985+
b=cached_batches.computeIfAbsent(sender, __ -> new MessageBatch(cap).dest(local_addr)
986+
.sender(sender).cluster(cl).incr(DEFAULT_INCREMENT));
987+
}
988+
else
989+
b=new MessageBatch(cap).dest(local_addr).sender(sender).cluster(cl).incr(DEFAULT_INCREMENT);
990+
}
991+
catch(Throwable t) {
992+
adders.set(0); // so others can remove/deliver msgs (https://redhat.atlassian.net/browse/JGRP-3034)
993+
throw t;
987994
}
988-
else
989-
b=new MessageBatch(cap).dest(local_addr).sender(sender).cluster(cl).incr(DEFAULT_INCREMENT);
990995
MessageBatch batch=b;
991996
Supplier<MessageBatch> batch_creator=() -> batch;
992997
MessageBatch mb=null;
@@ -997,7 +1002,8 @@ protected void removeAndDeliver(Table<Message> win, Address sender, AsciiString
9971002
batch_creator, BATCH_ACCUMULATOR);
9981003
}
9991004
catch(Throwable t) {
1000-
log.error("%s: failed removing messages from table for %s: %s", local_addr, sender, t);
1005+
// will not throw an exception
1006+
log.failSafeError("%s: failed removing messages from table for %s: %s", local_addr, sender, t);
10011007
}
10021008
if(!batch.isEmpty()) {
10031009
// batch is guaranteed to NOT contain any OOB messages as the drop_oob_msgs_filter above removed them
@@ -1172,7 +1178,7 @@ protected void deliverMessage(final Message msg, final Address sender, final lon
11721178
up_prot.up(msg);
11731179
}
11741180
catch(Throwable t) {
1175-
log.warn(Util.getMessage("FailedToDeliverMsg"), local_addr, msg.isFlagSet(OOB) ?
1181+
log.failSafeWarn(Util.getMessage("FailedToDeliverMsg"), local_addr, msg.isFlagSet(OOB) ?
11761182
"OOB message" : "message", msg, t);
11771183
}
11781184
}
@@ -1197,7 +1203,7 @@ protected void deliverBatch(MessageBatch batch) {
11971203
avg_delivery_batch_size.add(batch.size());
11981204
}
11991205
catch(Throwable t) {
1200-
log.warn(Util.getMessage("FailedToDeliverMsg"), local_addr, "batch", batch, t);
1206+
log.failSafeWarn(Util.getMessage("FailedToDeliverMsg"), local_addr, "batch", batch, t);
12011207
}
12021208
}
12031209

src/org/jgroups/protocols/pbcast/NAKACK2.java

Lines changed: 18 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -905,18 +905,24 @@ protected void removeAndDeliver(Table<Message> buf, Address sender, boolean loop
905905
AtomicInteger adders=buf.getAdders();
906906
if(adders.getAndIncrement() != 0)
907907
return;
908-
boolean remove_msgs=discard_delivered_msgs && !loopback;
909-
AsciiString cl=cluster != null? cluster : getTransport().getClusterNameAscii();
910-
int cap=Math.max(Math.max(Math.max(buf.size(), max_batch_size), min_size), DEFAULT_INITIAL_CAPACITY);
911908
MessageBatch b=null;
912-
if(reuse_message_batches) {
913-
b=cached_batches.get(sender);
914-
if(b == null)
915-
b=cached_batches.computeIfAbsent(sender, __ -> new MessageBatch(cap).dest(null).sender(sender)
916-
.cluster(cl).mcast(true).incr(DEFAULT_INCREMENT));
909+
boolean remove_msgs=discard_delivered_msgs && !loopback;
910+
try {
911+
AsciiString cl=cluster != null? cluster : getTransport().getClusterNameAscii();
912+
int cap=Math.max(Math.max(Math.max(buf.size(), max_batch_size), min_size), DEFAULT_INITIAL_CAPACITY);
913+
if(reuse_message_batches) {
914+
b=cached_batches.get(sender);
915+
if(b == null)
916+
b=cached_batches.computeIfAbsent(sender, __ -> new MessageBatch(cap).dest(null).sender(sender)
917+
.cluster(cl).mcast(true).incr(DEFAULT_INCREMENT));
918+
}
919+
else
920+
b=new MessageBatch(cap).dest(null).sender(sender).cluster(cl).mcast(true).incr(DEFAULT_INCREMENT);
921+
}
922+
catch(Throwable t) {
923+
adders.set(0); // so others can remove/deliver msgs (https://redhat.atlassian.net/browse/JGRP-3034)
924+
throw t;
917925
}
918-
else
919-
b=new MessageBatch(cap).dest(null).sender(sender).cluster(cl).mcast(true).incr(DEFAULT_INCREMENT);
920926
MessageBatch batch=b;
921927
Supplier<MessageBatch> batch_creator=() -> batch;
922928
MessageBatch mb=null;
@@ -928,7 +934,7 @@ protected void removeAndDeliver(Table<Message> buf, Address sender, boolean loop
928934
batch_creator, BATCH_ACCUMULATOR);
929935
}
930936
catch(Throwable t) {
931-
log.error("failed removing messages from table for " + sender, t);
937+
log.failSafeWarn("failed removing messages from table for " + sender, t);
932938
}
933939
int size=batch.size();
934940
if(size > 0) {
@@ -1003,7 +1009,7 @@ protected void deliverBatch(MessageBatch batch) {
10031009
batch.reset(); // doesn't null the messages in the batch
10041010
}
10051011
catch(Throwable t) {
1006-
log.error(Util.getMessage("FailedToDeliverMsg"), local_addr, "batch", batch, t);
1012+
log.failSafeError(Util.getMessage("FailedToDeliverMsg"), local_addr, "batch", batch, t);
10071013
}
10081014
}
10091015

src/org/jgroups/util/MaxOneThreadPerSender.java

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -273,14 +273,11 @@ public void run() {
273273
tp.passBatchUp(mb, !loopback, !loopback);
274274
}
275275
catch(Throwable t) {
276-
try {
277-
log.error("failed processing batch", t);
278-
}
279-
catch(Throwable tt) {
280-
// e.g. an OOME raised by log.error() above: don't terminate with running still set to true, or else
281-
// no further messages from that sender would ever be delivered: https://redhat.atlassian.net/browse/JGRP-3032
282-
;
283-
}
276+
// Will not throw an exception, e.g. an OOME raised by log.error() above: don't terminate with
277+
// entry.adders > 0, or else no further messages from that sender would ever be delivered:
278+
// https://redhat.atlassian.net/browse/JGRP-3032.
279+
// NPE due to null 'log' is impossible (log is guaranteed to be non-bull)
280+
log.failSafeError("failed processing batch", t);
284281
}
285282
}
286283
}

src/org/jgroups/util/SubmitToThreadPool.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
import org.jgroups.stack.MessageProcessingPolicy;
88

99
import java.util.Iterator;
10+
import java.util.Objects;
1011

1112
/**
1213
* Default message processing policy. Submits all received messages and batches to the thread pool
@@ -21,7 +22,7 @@ public class SubmitToThreadPool implements MessageProcessingPolicy {
2122

2223
public void init(TP transport) {
2324
this.tp=transport;
24-
this.log=tp.getLog();
25+
this.log=Objects.requireNonNull(tp.getLog());
2526
}
2627

2728
public boolean loopback(Message msg, boolean oob) {

0 commit comments

Comments
 (0)