|
|
@@ -3,6 +3,9 @@ package com.adx.tencent.vivo.service;
|
|
|
import com.adx.tencent.vivo.client.VivoClient;
|
|
|
import com.adx.tencent.vivo.model.VivoMediaCallbackRecord;
|
|
|
import com.adx.tencent.vivo.store.VivoColdStore;
|
|
|
+import com.fasterxml.jackson.databind.JsonNode;
|
|
|
+import com.fasterxml.jackson.databind.ObjectMapper;
|
|
|
+import com.fasterxml.jackson.databind.node.ObjectNode;
|
|
|
|
|
|
import java.time.Instant;
|
|
|
import java.util.ArrayList;
|
|
|
@@ -12,11 +15,13 @@ public class VivoRetryService {
|
|
|
|
|
|
private final VivoColdStore coldStore;
|
|
|
private final VivoClient vivoClient;
|
|
|
+ private final ObjectMapper objectMapper;
|
|
|
private final int defaultLimit;
|
|
|
|
|
|
- public VivoRetryService(VivoColdStore coldStore, VivoClient vivoClient, int defaultLimit) {
|
|
|
+ public VivoRetryService(VivoColdStore coldStore, VivoClient vivoClient, ObjectMapper objectMapper, int defaultLimit) {
|
|
|
this.coldStore = coldStore;
|
|
|
this.vivoClient = vivoClient;
|
|
|
+ this.objectMapper = objectMapper;
|
|
|
this.defaultLimit = defaultLimit > 0 ? defaultLimit : 100;
|
|
|
}
|
|
|
|
|
|
@@ -45,8 +50,70 @@ public class VivoRetryService {
|
|
|
return result;
|
|
|
}
|
|
|
|
|
|
+ public ReplaySrcIdResult replayFailedCallbacksReplacingSrcId(String oldSrcId, String newSrcId, int limit) {
|
|
|
+ String normalizedOld = requireSrcId(oldSrcId, "oldSrcId");
|
|
|
+ String normalizedNew = requireSrcId(newSrcId, "newSrcId");
|
|
|
+ List<VivoMediaCallbackRecord> pending = coldStore.failedMediaCallbacksBySrcId(
|
|
|
+ normalizedOld, limit <= 0 ? defaultLimit : limit);
|
|
|
+
|
|
|
+ ReplaySrcIdResult result = new ReplaySrcIdResult();
|
|
|
+ result.oldSrcId = normalizedOld;
|
|
|
+ result.newSrcId = normalizedNew;
|
|
|
+ result.fetched = pending.size();
|
|
|
+ for (VivoMediaCallbackRecord record : pending) {
|
|
|
+ String replacedBody = null;
|
|
|
+ try {
|
|
|
+ replacedBody = replaceSrcId(record.getRequestBody(), normalizedOld, normalizedNew);
|
|
|
+ if (replacedBody == null) {
|
|
|
+ result.skipped++;
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ VivoMediaCallbackRecord replay = vivoClient.replayCallback(record, replacedBody);
|
|
|
+ coldStore.updateMediaCallbackReplayResult(replay);
|
|
|
+ if (replay.isOk()) result.sent++;
|
|
|
+ else result.failed++;
|
|
|
+ } catch (Exception e) {
|
|
|
+ VivoMediaCallbackRecord failed = copy(record);
|
|
|
+ failed.setAttempt(record.getAttempt() + 1);
|
|
|
+ if (replacedBody != null) {
|
|
|
+ failed.setRequestBody(replacedBody);
|
|
|
+ }
|
|
|
+ failed.setOk(false);
|
|
|
+ failed.setErrorMessage(e.getMessage());
|
|
|
+ failed.setCreatedAt(Instant.now());
|
|
|
+ coldStore.updateMediaCallbackReplayResult(failed);
|
|
|
+ result.failed++;
|
|
|
+ }
|
|
|
+ }
|
|
|
+ return result;
|
|
|
+ }
|
|
|
+
|
|
|
+ private String replaceSrcId(String requestBody, String oldSrcId, String newSrcId) throws Exception {
|
|
|
+ if (requestBody == null || requestBody.isBlank()) {
|
|
|
+ return null;
|
|
|
+ }
|
|
|
+ JsonNode root = objectMapper.readTree(requestBody);
|
|
|
+ if (!(root instanceof ObjectNode objectNode)) {
|
|
|
+ return null;
|
|
|
+ }
|
|
|
+ JsonNode srcId = objectNode.get("srcId");
|
|
|
+ if (srcId == null || !oldSrcId.equals(srcId.asText())) {
|
|
|
+ return null;
|
|
|
+ }
|
|
|
+ objectNode.put("srcId", newSrcId);
|
|
|
+ return objectMapper.writeValueAsString(objectNode);
|
|
|
+ }
|
|
|
+
|
|
|
+ private static String requireSrcId(String value, String field) {
|
|
|
+ if (value == null || value.isBlank()) {
|
|
|
+ throw new IllegalArgumentException(field + " is required");
|
|
|
+ }
|
|
|
+ return value.trim();
|
|
|
+ }
|
|
|
+
|
|
|
private static VivoMediaCallbackRecord copy(VivoMediaCallbackRecord src) {
|
|
|
VivoMediaCallbackRecord copy = new VivoMediaCallbackRecord();
|
|
|
+ copy.setId(src.getId());
|
|
|
copy.setDedupeKey(src.getDedupeKey());
|
|
|
copy.setQk(src.getQk());
|
|
|
copy.setTagId(src.getTagId());
|
|
|
@@ -72,4 +139,13 @@ public class VivoRetryService {
|
|
|
public int sent;
|
|
|
public int failed;
|
|
|
}
|
|
|
+
|
|
|
+ public static class ReplaySrcIdResult {
|
|
|
+ public String oldSrcId;
|
|
|
+ public String newSrcId;
|
|
|
+ public int fetched;
|
|
|
+ public int sent;
|
|
|
+ public int failed;
|
|
|
+ public int skipped;
|
|
|
+ }
|
|
|
}
|