diff --git a/redisstream/RedisStreamMQ.cpp b/redisstream/RedisStreamMQ.cpp index 707ed9d..ebccb0f 100644 --- a/redisstream/RedisStreamMQ.cpp +++ b/redisstream/RedisStreamMQ.cpp @@ -686,6 +686,17 @@ private: } if (r->type == REDIS_REPLY_STRING) { + std::string message_id = r->str + ? std::string(r->str, r->len) + : std::string(); + std::cout << "[REDIS_STREAM][SEND_OK]" + << " stream=" << stream + << " id=" << message_id + << " key=" << key + << " tag=" << tag + << " enc=" << enc + << " body_len=" << body.size() + << std::endl; freeReplyObject(r); return REDIS_XADD_OK; } @@ -1023,13 +1034,19 @@ private: return NULL; } - void ack(const std::string& stream, const std::string& id) + bool ack(const std::string& stream, const std::string& id) { redisReply* r = conn_.command("XACK %b %b %b", stream.data(), stream.size(), group_.data(), group_.size(), id.data(), id.size()); - if (r) freeReplyObject(r); + if (!r) { + return false; + } + + bool acked = (r->type == REDIS_REPLY_INTEGER && r->integer > 0); + freeReplyObject(r); + return acked; } void handleOne(const Sub& sub, const std::string& id, redisReply* fieldsReply) @@ -1054,7 +1071,16 @@ private: rocketmq::ConsumeStatus ret = cb(msg); if (ret == rocketmq::CONSUME_SUCCESS) { - ack(sub.stream, id); + bool acked = ack(sub.stream, id); + std::cout << "[REDIS_STREAM][CONSUME_OK]" + << " stream=" << sub.stream + << " id=" << id + << " topic=" << sub.topic + << " tag=" << msg.getTags() + << " key=" << msg.getKeys() + << " body_len=" << msg.getBody().size() + << " ack=" << (acked ? "OK" : "FAIL") + << std::endl; } else { std::cout << "[REDIS_STREAM][RETRY_LATER] stream=" << sub.stream << " id=" << id << std::endl;