#include "redis_publisher.hpp" #include #include using namespace sw::redis; using json = nlohmann::json; // ------- helpers de conversão: funcionam em qualquer versão ------- namespace { // overload para retorno std::string (versões novas) inline std::string id_to_string(const std::string& id) { return id; } // overload para retorno OptionalString (versões antigas) inline std::string id_to_string(const sw::redis::OptionalString& id) { return id ? *id : std::string{}; } // Coloca um valor (string) em um json aninhado por path "a.b.c" inline void json_put_path(json& root, const std::string& path, const std::string& value, const std::string& sep) { size_t start = 0, pos; json* cur = &root; while ((pos = path.find(sep, start)) != std::string::npos) { auto key = path.substr(start, pos - start); if (!cur->contains(key) || !(*cur)[key].is_object()) { (*cur)[key] = json::object(); } cur = &(*cur)[key]; start = pos + sep.size(); } auto leaf = path.substr(start); // tenta parsear números/bool/json; se falhar, mantém string try { // Se for um literal JSON válido (ex: 1, 3.14, true, {"x":1}, [1,2]) json parsed = json::parse(value); (*cur)[leaf] = parsed; } catch (...) { // fallback: string (*cur)[leaf] = value; } } // Deep-merge recursivo: objetos são mesclados; tipos escalares/substituições diretas inline void json_deep_merge(json& dst, const json& src) { if (dst.is_object() && src.is_object()) { for (auto it = src.begin(); it != src.end(); ++it) { if (dst.contains(it.key())) { json_deep_merge(dst[it.key()], it.value()); } else { dst[it.key()] = it.value(); } } } else { // src substitui dst dst = src; } } } // namespace RedisPublisher::RedisPublisher(const std::string& uri) { try { r_ = std::make_unique(uri); r_->ping(); // testa conexão std::cout << "[REDIS] conectado.\n"; } catch (const Error& e) { std::cerr << "[REDIS] falha ao conectar: " << e.what() << "\n"; } } bool RedisPublisher::set_json(const std::string& key, const std::string& json, int ttl_sec) { try { if (!r_) return false; r_->set(key, json); if (ttl_sec > 0) r_->expire(key, std::chrono::seconds(ttl_sec)); return true; } catch (const TimeoutError& e) { std::cerr << "[REDIS] timeout SET: " << e.what() << "\n"; return false; } catch (const Error& e) { std::cerr << "[REDIS] erro SET: " << e.what() << "\n"; return false; } } std::string RedisPublisher::xadd(const std::string& stream, const std::vector>& fields, size_t maxlen) { try { if (!r_) return {}; if (maxlen > 0) { // MAXLEN ~ (aproximado) — API antiga/atual aceita estes args auto id = r_->xadd(stream, "*", fields.begin(), fields.end(), /*approx*/ true, static_cast(maxlen)); return id_to_string(id); } else { auto id = r_->xadd(stream, "*", fields.begin(), fields.end()); return id_to_string(id); } } catch (const sw::redis::Error& e) { std::cerr << "[REDIS] erro XADD: " << e.what() << "\n"; return {}; } } bool RedisPublisher::hset_fields(const std::string& key, const std::vector>& fields, int ttl_sec) { try { if (!r_) return false; // HSET com múltiplos campos de uma vez (iteradores) r_->hset(key, fields.begin(), fields.end()); if (ttl_sec > 0) r_->expire(key, std::chrono::seconds(ttl_sec)); return true; } catch (const sw::redis::Error& e) { std::cerr << "[REDIS] erro HSET: " << e.what() << "\n"; return false; } } bool RedisPublisher::set_json_merge_paths(const std::string& key, const std::vector>& flat_fields, int ttl_sec, const std::string& sep) { try { if (!r_) return false; // Tentativa com controle de concorrência (optimistic locking) for (int attempt = 0; attempt < 5; ++attempt) { r_->watch(key); // Lê JSON atual (pode não existir) auto cur_opt = r_->get(key); json cur = json::object(); if (cur_opt) { try { cur = json::parse(*cur_opt); } catch (...) { cur = json::object(); } // se não for JSON, corrige } // Aplica o patch “flat” nos caminhos for (const auto& kv : flat_fields) { // aceita também "a__b" => "a.b" se quiser std::string path = kv.first; // se quiser compat com "__" como no seu Python: // substitua "__" por "." // (descomente a linha abaixo se quiser esse comportamento aqui tbm) // for (size_t pos = 0; (pos = path.find("__", pos)) != std::string::npos; pos += 1) path.replace(pos, 2, "."); json_put_path(cur, path, kv.second, sep); } auto dump = cur.dump(); // Transação: SET (+ EXPIRE opcional) atômicos auto tx = r_->transaction(); tx.set(key, dump); if (ttl_sec > 0) tx.expire(key, std::chrono::seconds(ttl_sec)); try { tx.exec(); // sucesso se NÃO lançar r_->unwatch(); return true; } catch (const sw::redis::WatchError&) { r_->unwatch(); // conflito de WATCH: tenta de novo continue; // volta pro loop de retry } catch (const sw::redis::Error& e) { r_->unwatch(); std::cerr << "[REDIS] exec falhou: " << e.what() << "\n"; return false; } // conflito: retry } r_->unwatch(); return false; } catch (const sw::redis::Error& e) { std::cerr << "[REDIS] erro set_json_merge_paths: " << e.what() << "\n"; return false; } } bool RedisPublisher::set_json_merge_object(const std::string& key, const json& patch, int ttl_sec) { try { if (!r_) return false; for (int attempt = 0; attempt < 5; ++attempt) { r_->watch(key); auto cur_opt = r_->get(key); json cur = json::object(); if (cur_opt) { try { cur = json::parse(*cur_opt); } catch (...) { cur = json::object(); } } json_deep_merge(cur, patch); auto dump = cur.dump(); auto tx = r_->transaction(); tx.set(key, dump); if (ttl_sec > 0) tx.expire(key, std::chrono::seconds(ttl_sec)); try { tx.exec(); // sucesso se NÃO lançar r_->unwatch(); return true; } catch (const sw::redis::WatchError&) { r_->unwatch(); // conflito de WATCH: tenta de novo continue; // volta pro loop de retry } catch (const sw::redis::Error& e) { r_->unwatch(); std::cerr << "[REDIS] exec falhou: " << e.what() << "\n"; return false; } } r_->unwatch(); return false; } catch (const sw::redis::Error& e) { std::cerr << "[REDIS] erro set_json_merge_object: " << e.what() << "\n"; return false; } }