yumeng 1 неделя назад
Родитель
Сommit
932ba32262

+ 9 - 0
src/main/java/com/adx/tencent/AppConfiguration.java

@@ -7,6 +7,7 @@ import com.adx.tencent.config.AppProperties;
 import com.adx.tencent.conversionsync.ConversionSyncRunner;
 import com.adx.tencent.conversionsync.ConversionSyncRunner;
 import com.adx.tencent.conversionsync.ConversionSyncService;
 import com.adx.tencent.conversionsync.ConversionSyncService;
 import com.adx.tencent.conversionsync.ConversionBackfillJobService;
 import com.adx.tencent.conversionsync.ConversionBackfillJobService;
+import com.adx.tencent.conversionsync.ManualTencentCallbackService;
 import com.adx.tencent.conversionsync.RetryService;
 import com.adx.tencent.conversionsync.RetryService;
 import com.adx.tencent.httpapi.MediaPlacement;
 import com.adx.tencent.httpapi.MediaPlacement;
 import com.adx.tencent.honor.client.HonorClient;
 import com.adx.tencent.honor.client.HonorClient;
@@ -346,6 +347,14 @@ public class AppConfiguration {
     }
     }
 
 
     @Bean
     @Bean
+    public ManualTencentCallbackService manualTencentCallbackService(@Nullable TiDBColdStore coldStore,
+                                                                     RedisHotStore hotStore,
+                                                                     TencentClient tencentClient) {
+        if (coldStore == null) return null;
+        return new ManualTencentCallbackService(coldStore, hotStore, tencentClient);
+    }
+
+    @Bean
     public TagEventSyncService tagEventSyncService(@Nullable TiDBColdStore coldStore,
     public TagEventSyncService tagEventSyncService(@Nullable TiDBColdStore coldStore,
                                                    StringRedisTemplate redisTemplate) {
                                                    StringRedisTemplate redisTemplate) {
         if (coldStore == null) return null;
         if (coldStore == null) return null;

+ 232 - 0
src/main/java/com/adx/tencent/conversionsync/ManualTencentCallbackService.java

@@ -0,0 +1,232 @@
+package com.adx.tencent.conversionsync;
+
+import com.adx.tencent.baidu.model.PaymentInfo;
+import com.adx.tencent.storage.RedisHotStore;
+import com.adx.tencent.storage.TiDBColdStore;
+import com.adx.tencent.storage.model.BidRecord;
+import com.adx.tencent.storage.model.MediaCallbackRecord;
+import com.adx.tencent.tencent.TencentClient;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.web.server.ResponseStatusException;
+
+import java.time.Instant;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.Map;
+import java.util.Set;
+
+import static org.springframework.http.HttpStatus.BAD_REQUEST;
+import static org.springframework.http.HttpStatus.NOT_FOUND;
+
+public class ManualTencentCallbackService {
+
+    private static final Logger log = LoggerFactory.getLogger(ManualTencentCallbackService.class);
+
+    private final TiDBColdStore coldStore;
+    private final RedisHotStore hotStore;
+    private final TencentClient tencentClient;
+
+    public ManualTencentCallbackService(TiDBColdStore coldStore,
+                                        RedisHotStore hotStore,
+                                        TencentClient tencentClient) {
+        this.coldStore = coldStore;
+        this.hotStore = hotStore;
+        this.tencentClient = tencentClient;
+    }
+
+    public Map<String, Object> replayDeductedCallback(long id) {
+        MediaCallbackRecord original = coldStore.getMediaCallbackById(id);
+        if (original == null) {
+            throw new ResponseStatusException(NOT_FOUND, "callback record not found");
+        }
+        if (!"tencent".equals(original.getMedia())) {
+            throw new ResponseStatusException(BAD_REQUEST, "callback record is not tencent");
+        }
+        if (!MediaCallbackRecord.DISPATCH_STATUS_DEDUCTED.equals(original.getDispatchStatus())) {
+            throw new ResponseStatusException(BAD_REQUEST, "callback record is not deducted");
+        }
+
+        MediaCallbackRecord replay = replayOnce(original);
+        replay.setId(original.getId());
+        coldStore.updateMediaCallback(replay);
+
+        Map<String, Object> result = new LinkedHashMap<>();
+        result.put("id", id);
+        result.put("dedupeKey", replay.getDedupeKey());
+        result.put("before", snapshot(original));
+        result.put("after", snapshot(replay));
+        return result;
+    }
+
+    private MediaCallbackRecord replayOnce(MediaCallbackRecord original) {
+        if (original.getRequestBody() != null && !original.getRequestBody().isBlank()) {
+            try {
+                MediaCallbackRecord replay = tencentClient.retryCallback(original);
+                replay.setDedupeKey(original.getDedupeKey());
+                replay.setAttempt(Math.max(original.getAttempt(), 1) + 1);
+                if (replay.getTrackingVersion() == null || replay.getTrackingVersion().isBlank()) {
+                    replay.setTrackingVersion(original.getTrackingVersion());
+                }
+                return replay;
+            } catch (Exception e) {
+                MediaCallbackRecord failed = copyForFailure(original);
+                failed.setErrorMessage(e.getMessage());
+                return failed;
+            }
+        }
+
+        BidRecord bid = hotStore.findBidByQk(original.getQk());
+        if (bid == null) {
+            bid = coldStore.getBidByQk(original.getQk());
+        }
+        if (bid == null) {
+            throw new ResponseStatusException(NOT_FOUND, "bid not found for qk=" + original.getQk());
+        }
+
+        DedupeInfo dedupeInfo = parseDedupeInfo(original);
+        String actionType = dedupeInfo.actionType != null && !dedupeInfo.actionType.isBlank()
+                ? dedupeInfo.actionType
+                : tencentClient.conversionActionType(original.getEventType());
+
+        patchBidForReplay(bid, original, dedupeInfo);
+
+        PaymentInfo payment = new PaymentInfo();
+        payment.setQk(original.getQk());
+        payment.setDate(dedupeInfo.date);
+        payment.setDeviceId(dedupeInfo.deviceId);
+        payment.setAct(original.getEventType());
+        payment.setPayment(original.getPurchase());
+        payment.setGmv(original.getPurchase());
+
+        MediaCallbackRecord replay = tencentClient.sendConversionEvent(
+                bid,
+                payment,
+                resolvePlatform(bid),
+                actionType
+        );
+        replay.setDedupeKey(original.getDedupeKey());
+        replay.setAttempt(Math.max(original.getAttempt(), 1) + 1);
+        if (replay.getTrackingVersion() == null || replay.getTrackingVersion().isBlank()) {
+            replay.setTrackingVersion(original.getTrackingVersion());
+        }
+        replay.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT);
+        replay.setCreatedAt(Instant.now());
+        log.info("[ManualTencentCallback] manual replay sent | qk={} | dedupeKey={} | ok={} | status={}",
+                original.getQk(), original.getDedupeKey(), replay.isOk(), replay.getStatus());
+        return replay;
+    }
+
+    private void patchBidForReplay(BidRecord bid, MediaCallbackRecord original, DedupeInfo dedupeInfo) {
+        if ((bid.getTagId() == null || bid.getTagId().isBlank()) && dedupeInfo.tagId != null && !dedupeInfo.tagId.isBlank()) {
+            bid.setTagId(dedupeInfo.tagId);
+        }
+
+        Map<String, String> params = bid.getMediaParams();
+        if (params == null) {
+            params = new LinkedHashMap<>();
+            bid.setMediaParams(params);
+        } else if (!(params instanceof LinkedHashMap<String, String>)) {
+            params = new LinkedHashMap<>(params);
+            bid.setMediaParams(params);
+        }
+        if ((params.get("callback") == null || params.get("callback").isBlank())
+                && original.getCallbackUrl() != null && !original.getCallbackUrl().isBlank()) {
+            params.put("callback", original.getCallbackUrl());
+        }
+        if ((params.get(BidRecord.TRACKING_VERSION_KEY) == null || params.get(BidRecord.TRACKING_VERSION_KEY).isBlank())
+                && original.getTrackingVersion() != null && !original.getTrackingVersion().isBlank()) {
+            params.put(BidRecord.TRACKING_VERSION_KEY, original.getTrackingVersion());
+        }
+    }
+
+    private static MediaCallbackRecord copyForFailure(MediaCallbackRecord original) {
+        MediaCallbackRecord failed = new MediaCallbackRecord();
+        failed.setDedupeKey(original.getDedupeKey());
+        failed.setMedia(original.getMedia());
+        failed.setQk(original.getQk());
+        failed.setCallbackUrl(original.getCallbackUrl());
+        failed.setRequestBody(original.getRequestBody());
+        failed.setEventType(original.getEventType());
+        failed.setEventTimeMs(Instant.now().toEpochMilli());
+        failed.setPurchase(original.getPurchase());
+        failed.setStatus(original.getStatus());
+        failed.setOk(false);
+        failed.setAttempt(Math.max(original.getAttempt(), 1) + 1);
+        failed.setResponseBody(original.getResponseBody());
+        failed.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT);
+        failed.setTrackingVersion(original.getTrackingVersion());
+        failed.setCreatedAt(Instant.now());
+        return failed;
+    }
+
+    private static DedupeInfo parseDedupeInfo(MediaCallbackRecord record) {
+        DedupeInfo info = new DedupeInfo();
+        String dedupeKey = record.getDedupeKey();
+        if (dedupeKey == null || dedupeKey.isBlank()) {
+            return info;
+        }
+
+        String[] parts = dedupeKey.split(":");
+        if (parts.length < 7) {
+            return info;
+        }
+        info.date = parts[2];
+        info.deviceId = parts[3];
+        info.tagId = parts[4];
+        info.actionType = joinTail(parts, 6);
+        return info;
+    }
+
+    private static String joinTail(String[] parts, int start) {
+        if (parts.length <= start) {
+            return null;
+        }
+        StringBuilder builder = new StringBuilder();
+        for (int i = start; i < parts.length; i++) {
+            if (i > start) {
+                builder.append(':');
+            }
+            builder.append(parts[i]);
+        }
+        return builder.toString();
+    }
+
+    private static String resolvePlatform(BidRecord bid) {
+        if (bid.getPlatform() != null && !bid.getPlatform().isBlank()) {
+            return bid.getPlatform().toLowerCase();
+        }
+        if (bid.getMediaParams() == null) {
+            return "android";
+        }
+        String platform = bid.getMediaParams().get("baidu_platform");
+        if (platform != null && !platform.isBlank()) {
+            return platform.toLowerCase();
+        }
+        String idfa = bid.getMediaParams().get("idfa");
+        if (idfa != null && !idfa.isBlank()) {
+            return "ios";
+        }
+        return "android";
+    }
+
+    private static Map<String, Object> snapshot(MediaCallbackRecord record) {
+        Map<String, Object> data = new LinkedHashMap<>();
+        data.put("id", record.getId());
+        data.put("dispatchStatus", record.getDispatchStatus());
+        data.put("attempt", record.getAttempt());
+        data.put("status", record.getStatus());
+        data.put("ok", record.isOk());
+        data.put("errorMessage", record.getErrorMessage());
+        data.put("callbackUrl", record.getCallbackUrl());
+        data.put("trackingVersion", record.getTrackingVersion());
+        return data;
+    }
+
+    private static class DedupeInfo {
+        private String date;
+        private String deviceId;
+        private String tagId;
+        private String actionType;
+    }
+}

+ 24 - 1
src/main/java/com/adx/tencent/httpapi/AdminController.java

@@ -5,6 +5,7 @@ import com.adx.tencent.baidu.model.ConversionQuery;
 import com.adx.tencent.conversionsync.ConversionBackfillJobService;
 import com.adx.tencent.conversionsync.ConversionBackfillJobService;
 import com.adx.tencent.conversionsync.ConversionSyncRunner;
 import com.adx.tencent.conversionsync.ConversionSyncRunner;
 import com.adx.tencent.conversionsync.ConversionSyncService;
 import com.adx.tencent.conversionsync.ConversionSyncService;
+import com.adx.tencent.conversionsync.ManualTencentCallbackService;
 import com.adx.tencent.conversionsync.RetryService;
 import com.adx.tencent.conversionsync.RetryService;
 import org.springframework.http.HttpStatus;
 import org.springframework.http.HttpStatus;
 import org.springframework.lang.Nullable;
 import org.springframework.lang.Nullable;
@@ -26,17 +27,20 @@ public class AdminController {
     private final ConversionBackfillJobService backfillJobService;
     private final ConversionBackfillJobService backfillJobService;
     private final ConversionSyncService conversionSyncer;
     private final ConversionSyncService conversionSyncer;
     private final RetryService retryService;
     private final RetryService retryService;
+    private final ManualTencentCallbackService manualTencentCallbackService;
 
 
     public AdminController(ConversionClient conversionClient,
     public AdminController(ConversionClient conversionClient,
                            @Nullable ConversionSyncRunner conversionSyncRunner,
                            @Nullable ConversionSyncRunner conversionSyncRunner,
                            @Nullable ConversionBackfillJobService backfillJobService,
                            @Nullable ConversionBackfillJobService backfillJobService,
                            @Nullable ConversionSyncService conversionSyncer,
                            @Nullable ConversionSyncService conversionSyncer,
-                           @Nullable RetryService retryService) {
+                           @Nullable RetryService retryService,
+                           @Nullable ManualTencentCallbackService manualTencentCallbackService) {
         this.conversionClient = conversionClient;
         this.conversionClient = conversionClient;
         this.conversionSyncRunner = conversionSyncRunner;
         this.conversionSyncRunner = conversionSyncRunner;
         this.backfillJobService = backfillJobService;
         this.backfillJobService = backfillJobService;
         this.conversionSyncer = conversionSyncer;
         this.conversionSyncer = conversionSyncer;
         this.retryService = retryService;
         this.retryService = retryService;
+        this.manualTencentCallbackService = manualTencentCallbackService;
     }
     }
 
 
     @PostMapping("/conversions/query")
     @PostMapping("/conversions/query")
@@ -86,12 +90,31 @@ public class AdminController {
         return retryService.retryTencentCallbacks(limit);
         return retryService.retryTencentCallbacks(limit);
     }
     }
 
 
+    @PostMapping("/callbacks/tencent/manual-replay")
+    public Object manualReplayTencentCallback(@RequestBody Map<String, Object> body) {
+        if (manualTencentCallbackService == null) {
+            throw new ResponseStatusException(HttpStatus.SERVICE_UNAVAILABLE,
+                    "manual tencent callback service is not configured");
+        }
+        long id = toLong(body.get("id"));
+        if (id <= 0) {
+            throw new ResponseStatusException(HttpStatus.BAD_REQUEST, "id is required");
+        }
+        return manualTencentCallbackService.replayDeductedCallback(id);
+    }
+
     private static int toInt(Object v) {
     private static int toInt(Object v) {
         if (v == null) return 0;
         if (v == null) return 0;
         if (v instanceof Number n) return n.intValue();
         if (v instanceof Number n) return n.intValue();
         try { return Integer.parseInt(v.toString()); } catch (NumberFormatException e) { return 0; }
         try { return Integer.parseInt(v.toString()); } catch (NumberFormatException e) { return 0; }
     }
     }
 
 
+    private static long toLong(Object v) {
+        if (v == null) return 0L;
+        if (v instanceof Number n) return n.longValue();
+        try { return Long.parseLong(v.toString()); } catch (NumberFormatException e) { return 0L; }
+    }
+
     private static List<Integer> parseOffsets(Map<String, Object> body) {
     private static List<Integer> parseOffsets(Map<String, Object> body) {
         if (body == null || body.isEmpty()) {
         if (body == null || body.isEmpty()) {
             return defaultBackfillOffsets();
             return defaultBackfillOffsets();

+ 26 - 0
src/main/java/com/adx/tencent/storage/TiDBColdStore.java

@@ -379,6 +379,32 @@ public class TiDBColdStore {
         }
         }
     }
     }
 
 
+    public MediaCallbackRecord getMediaCallbackById(long id) {
+        if (id <= 0) return null;
+        try (SqlSession session = sqlSessionFactory.openSession(true)) {
+            MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
+            return mapper.selectById(id);
+        }
+    }
+
+    public void updateMediaCallback(MediaCallbackRecord record) {
+        if (record == null || record.getId() == null || record.getId() <= 0) {
+            throw new IllegalArgumentException("media callback id is required");
+        }
+        int attempt = record.getAttempt() <= 0 ? 1 : record.getAttempt();
+        record.setAttempt(attempt);
+        if (record.getDispatchStatus() == null || record.getDispatchStatus().isBlank()) {
+            record.setDispatchStatus(MediaCallbackRecord.DISPATCH_STATUS_SENT);
+        }
+        if (record.getTrackingVersion() == null || record.getTrackingVersion().isBlank()) {
+            record.setTrackingVersion("v1");
+        }
+        try (SqlSession session = sqlSessionFactory.openSession(true)) {
+            MediaCallbackMapper mapper = session.getMapper(MediaCallbackMapper.class);
+            mapper.updateById(record);
+        }
+    }
+
     // ─── 查询:广告位回传方式配置 ──────────────────────────────────────────────
     // ─── 查询:广告位回传方式配置 ──────────────────────────────────────────────
 
 
     public List<TagEventRecord> listAllTagEvents() {
     public List<TagEventRecord> listAllTagEvents() {

+ 4 - 0
src/main/java/com/adx/tencent/storage/mapper/MediaCallbackMapper.java

@@ -41,4 +41,8 @@ public interface MediaCallbackMapper {
 
 
     List<MediaCallbackRecord> selectPendingCallbacks(@Param("media") String media,
     List<MediaCallbackRecord> selectPendingCallbacks(@Param("media") String media,
                                                      @Param("limit") int limit);
                                                      @Param("limit") int limit);
+
+    MediaCallbackRecord selectById(@Param("id") long id);
+
+    void updateById(MediaCallbackRecord record);
 }
 }

+ 3 - 0
src/main/java/com/adx/tencent/storage/model/MediaCallbackRecord.java

@@ -7,6 +7,7 @@ public class MediaCallbackRecord {
     public static final String DISPATCH_STATUS_SENT = "SENT";
     public static final String DISPATCH_STATUS_SENT = "SENT";
     public static final String DISPATCH_STATUS_DEDUCTED = "DEDUCTED";
     public static final String DISPATCH_STATUS_DEDUCTED = "DEDUCTED";
 
 
+    @JsonProperty("id")           private Long id;
     @JsonProperty("dedupeKey")    private String dedupeKey;
     @JsonProperty("dedupeKey")    private String dedupeKey;
     @JsonProperty("media")        private String media;
     @JsonProperty("media")        private String media;
     @JsonProperty("qk")           private String qk;
     @JsonProperty("qk")           private String qk;
@@ -24,6 +25,8 @@ public class MediaCallbackRecord {
     @JsonProperty("trackingVersion") private String trackingVersion;
     @JsonProperty("trackingVersion") private String trackingVersion;
     @JsonProperty("createdAt")    private Instant createdAt;
     @JsonProperty("createdAt")    private Instant createdAt;
 
 
+    public Long getId() { return id; }
+    public void setId(Long v) { this.id = v; }
     public String getDedupeKey() { return dedupeKey; }
     public String getDedupeKey() { return dedupeKey; }
     public void setDedupeKey(String v) { this.dedupeKey = v; }
     public void setDedupeKey(String v) { this.dedupeKey = v; }
     public String getMedia() { return media; }
     public String getMedia() { return media; }

+ 27 - 1
src/main/resources/mapper/MediaCallbackMapper.xml

@@ -3,6 +3,7 @@
         "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
         "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
 <mapper namespace="com.adx.tencent.storage.mapper.MediaCallbackMapper">
 <mapper namespace="com.adx.tencent.storage.mapper.MediaCallbackMapper">
     <resultMap id="mediaCallbackResultMap" type="com.adx.tencent.storage.model.MediaCallbackRecord">
     <resultMap id="mediaCallbackResultMap" type="com.adx.tencent.storage.model.MediaCallbackRecord">
+        <id property="id" column="id"/>
         <result property="dedupeKey" column="dedupe_key"/>
         <result property="dedupeKey" column="dedupe_key"/>
         <result property="media" column="media"/>
         <result property="media" column="media"/>
         <result property="qk" column="qk"/>
         <result property="qk" column="qk"/>
@@ -98,7 +99,7 @@
     </insert>
     </insert>
 
 
     <select id="selectPendingCallbacks" resultMap="mediaCallbackResultMap">
     <select id="selectPendingCallbacks" resultMap="mediaCallbackResultMap">
-        SELECT m.dedupe_key, m.media, m.qk, m.callback_url, m.event_type, m.event_time_ms,
+        SELECT m.id, m.dedupe_key, m.media, m.qk, m.callback_url, m.event_type, m.event_time_ms,
                m.purchase, m.status, m.ok AS ok_val, m.attempt, m.response_body, m.error_message,
                m.purchase, m.status, m.ok AS ok_val, m.attempt, m.response_body, m.error_message,
                m.request_body, m.dispatch_status, m.tracking_version, m.created_at
                m.request_body, m.dispatch_status, m.tracking_version, m.created_at
         FROM tencent_media_callbacks m
         FROM tencent_media_callbacks m
@@ -118,4 +119,29 @@
         ORDER BY m.created_at ASC
         ORDER BY m.created_at ASC
         LIMIT #{limit}
         LIMIT #{limit}
     </select>
     </select>
+
+    <select id="selectById" resultMap="mediaCallbackResultMap">
+        SELECT id, dedupe_key, media, qk, callback_url, event_type, event_time_ms,
+               purchase, status, ok AS ok_val, attempt, response_body, error_message,
+               request_body, dispatch_status, tracking_version, created_at
+        FROM tencent_media_callbacks
+        WHERE id = #{id}
+        LIMIT 1
+    </select>
+
+    <update id="updateById" parameterType="com.adx.tencent.storage.model.MediaCallbackRecord">
+        UPDATE tencent_media_callbacks
+        SET callback_url = #{callbackUrl},
+            event_time_ms = #{eventTimeMs},
+            purchase = #{purchase},
+            status = #{status},
+            ok = #{ok},
+            attempt = #{attempt},
+            response_body = #{responseBody},
+            error_message = #{errorMessage},
+            request_body = #{requestBody},
+            dispatch_status = #{dispatchStatus},
+            tracking_version = #{trackingVersion}
+        WHERE id = #{id}
+    </update>
 </mapper>
 </mapper>