Browse Source

Rule engine: ack all rate limited failures (draft)

pull/8983/head
Sergey Matvienko 3 years ago
parent
commit
1b565e01ce
  1. 18
      application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java
  2. 4
      common/message/src/main/java/org/thingsboard/server/common/msg/queue/TbMsgCallback.java

18
application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java

@ -17,11 +17,13 @@ package org.thingsboard.server.service.queue;
import io.micrometer.core.instrument.Timer; import io.micrometer.core.instrument.Timer;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.exception.ApiUsageLimitsExceededException;
import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.RuleEngineException; import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.common.msg.queue.RuleNodeInfo; import org.thingsboard.server.common.msg.queue.RuleNodeInfo;
import org.thingsboard.server.common.msg.queue.TbMsgCallback; import org.thingsboard.server.common.msg.queue.TbMsgCallback;
import org.thingsboard.server.common.msg.tools.TbRateLimitsException;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
@ -57,8 +59,24 @@ public class TbMsgPackCallback implements TbMsgCallback {
ctx.onSuccess(id); ctx.onSuccess(id);
} }
@Override
public void onRateLimit(RuleEngineException e) {
log.debug("[{}] ON RATE LIMIT", id, e);
//TODO notify tenant on rate limit
if (failedMsgTimer != null) {
failedMsgTimer.record(System.currentTimeMillis() - startMsgProcessing, TimeUnit.MILLISECONDS);
}
ctx.onSuccess(id);
}
@Override @Override
public void onFailure(RuleEngineException e) { public void onFailure(RuleEngineException e) {
Throwable cause = e.getCause();
if (cause instanceof TbRateLimitsException || cause instanceof ApiUsageLimitsExceededException) {
onRateLimit(e);
return;
}
log.trace("[{}] ON FAILURE", id, e); log.trace("[{}] ON FAILURE", id, e);
if (failedMsgTimer != null) { if (failedMsgTimer != null) {
failedMsgTimer.record(System.currentTimeMillis() - startMsgProcessing, TimeUnit.MILLISECONDS); failedMsgTimer.record(System.currentTimeMillis() - startMsgProcessing, TimeUnit.MILLISECONDS);

4
common/message/src/main/java/org/thingsboard/server/common/msg/queue/TbMsgCallback.java

@ -39,6 +39,10 @@ public interface TbMsgCallback {
void onFailure(RuleEngineException e); void onFailure(RuleEngineException e);
default void onRateLimit(RuleEngineException e) {
onFailure(e);
};
/** /**
* Returns 'true' if rule engine is expecting the message to be processed, 'false' otherwise. * Returns 'true' if rule engine is expecting the message to be processed, 'false' otherwise.
* message may no longer be valid, if the message pack is already expired/canceled/failed. * message may no longer be valid, if the message pack is already expired/canceled/failed.

Loading…
Cancel
Save