#include "store.h" #include "networkmanager.h" #include #include #include #include #include // ──────────────────────────────────────────────────────────────── // Operation // ──────────────────────────────────────────────────────────────── Operation::Operation(QNetworkReply *reply, QObject *parent) : QObject(parent) { connect(this, &Operation::finished, this, &QObject::deleteLater); connect(reply, &QNetworkReply::finished, this, [this, reply]() { reply->deleteLater(); if (reply->error() != QNetworkReply::NoError) { QByteArray body = reply->readAll(); QJsonDocument doc = QJsonDocument::fromJson(body); QString message = doc.isObject() ? doc.object().value("message").toString() : QString(); if (message.isEmpty()) { message = reply->errorString(); } qWarning() << message; emit failed(message); emit finished(); return; } QByteArray body = reply->readAll(); QJsonDocument doc = body.isEmpty() ? QJsonDocument() : QJsonDocument::fromJson(body); emit success(doc); emit finished(); }); } // ──────────────────────────────────────────────────────────────── // Store — construction // ──────────────────────────────────────────────────────────────── Store::Store(QObject *parent) : QObject(parent) { } // ──────────────────────────────────────────────────────────────── // JSON parsing // ──────────────────────────────────────────────────────────────── Human Store::parseHuman(const QJsonObject &obj) { Human h; h.id = obj["id"].toString(); h.email = obj["email"].toString(); h.emailPrefix = obj["email_prefix"].toString(); return h; } Network Store::parseNetwork(const QJsonObject &obj) { Network n; n.id = obj["id"].toString(); n.name = obj["name"].toString(); n.admin = parseHuman(obj["admin_human"].toObject()); n.openStreamCount = obj["open_stream_count"].toInt(); n.openStreamCapacity = obj["open_stream_capacity"].toInt(); const QJsonArray humansArr = obj["humans"].toArray(); for (const QJsonValue &v : humansArr) { n.members.push_back(parseHuman(v.toObject())); } // Streams are parsed at a higher level (populateFromStartupData) so we can // also extract particles per stream. Individual fetch calls parse streams inline. return n; } Stream Store::parseStream(const QJsonObject &obj) { Stream s; s.id = obj["id"].toString(); s.name = obj["name"].toString(); s.description = obj["description"].toString(); s.isOpen = (obj["status"].toString() == "open"); s.unseenCount = obj["unseen_count"].toInt(); const QJsonArray membersArr = obj["members"].toArray(); for (const QJsonValue &v : membersArr) { s.memberEmails.push_back(v.toString()); } return s; } Particle Store::parseParticle(const QJsonObject &obj) { Particle p; p.id = obj["id"].toString(); p.type = obj["type"].toString(); p.createdByEmail = obj["created_by_email"].toString(); p.seen = obj["seen"].toBool(); const QJsonObject dataObj = obj["data"].toObject(); for (auto it = dataObj.begin(); it != dataObj.end(); ++it) { p.data[it.key()] = it.value().toVariant(); } const QJsonArray acksArr = obj["acks"].toArray(); for (const QJsonValue &v : acksArr) { // acks can be objects with { email, acked_at } or plain strings if (v.isObject()) { p.ackedByEmails.push_back(v.toObject()["email"].toString()); } else { p.ackedByEmails.push_back(v.toString()); } } p.updatedAt = QDateTime::fromString(obj["updated_at"].toString(), Qt::ISODate); p.createdAt = QDateTime::fromString(obj["created_at"].toString(), Qt::ISODate); return p; } void Store::populateFromStartupData(const QJsonObject &root) { m_networks.clear(); m_particles.clear(); m_networkIndex.clear(); m_streamToNetwork.clear(); const QJsonArray networksArr = root["networks"].toArray(); for (const QJsonValue &netVal : networksArr) { QJsonObject netObj = netVal.toObject(); Network network = parseNetwork(netObj); const QJsonArray streamsArr = netObj["streams"].toArray(); for (const QJsonValue &streamVal : streamsArr) { QJsonObject streamObj = streamVal.toObject(); Stream stream = parseStream(streamObj); network.streams.push_back(stream); m_streamToNetwork[stream.id] = network.id; // Parse particles embedded in stream QList particleList; const QJsonArray particlesArr = streamObj["particles"].toArray(); for (const QJsonValue &partVal : particlesArr) { particleList.append(parseParticle(partVal.toObject())); } m_particles[stream.id] = particleList; } m_networkIndex[network.id] = m_networks.size(); m_networks.append(network); } } // ──────────────────────────────────────────────────────────────── // State mutation helpers // ──────────────────────────────────────────────────────────────── void Store::upsertNetwork(const Network &network) { if (m_networkIndex.contains(network.id)) { int idx = m_networkIndex[network.id]; std::vector existingStreams = m_networks[idx].streams; m_networks[idx] = network; if (network.streams.empty()) { m_networks[idx].streams = existingStreams; } } else { m_networkIndex[network.id] = m_networks.size(); m_networks.append(network); } } void Store::upsertStream(const QString &networkId, const Stream &stream) { if (!m_networkIndex.contains(networkId)) { return; } auto &streams = m_networks[m_networkIndex[networkId]].streams; for (size_t i = 0; i < streams.size(); ++i) { if (streams[i].id == stream.id) { streams[i] = stream; m_streamToNetwork[stream.id] = networkId; return; } } streams.push_back(stream); m_streamToNetwork[stream.id] = networkId; } void Store::upsertParticle(const QString &streamId, const Particle &particle) { QList &list = m_particles[streamId]; for (int i = 0; i < list.size(); ++i) { if (list[i].id == particle.id) { list[i] = particle; return; } } list.append(particle); } void Store::removeParticleFromState(const QString &particleId) { for (auto it = m_particles.begin(); it != m_particles.end(); ++it) { QList &list = it.value(); for (int i = 0; i < list.size(); ++i) { if (list[i].id == particleId) { list.removeAt(i); return; } } } } QString Store::findStreamForParticle(const QString &particleId) const { for (auto it = m_particles.constBegin(); it != m_particles.constEnd(); ++it) { for (const auto &p : it.value()) { if (p.id == particleId) { return it.key(); } } } return QString(); } // ──────────────────────────────────────────────────────────────── // Accessors // ──────────────────────────────────────────────────────────────── const QList &Store::networks() const { return m_networks; } int Store::networkCount() const { return m_networks.size(); } const Network *Store::networkById(const QString &id) const { auto it = m_networkIndex.find(id); if (it == m_networkIndex.end()) { return nullptr; } return &m_networks[it.value()]; } QList Store::streamsForNetwork(const QString &networkId) const { auto it = m_networkIndex.find(networkId); if (it == m_networkIndex.end()) { return {}; } const auto &streams = m_networks[it.value()].streams; QList result; result.reserve(static_cast(streams.size())); for (const auto &s : streams) { result.append(s); } return result; } int Store::streamCount(const QString &networkId) const { auto it = m_networkIndex.find(networkId); if (it == m_networkIndex.end()) { return 0; } return static_cast(m_networks[it.value()].streams.size()); } const Stream *Store::streamById(const QString &id) const { auto netIt = m_streamToNetwork.find(id); if (netIt == m_streamToNetwork.end()) { return nullptr; } auto idxIt = m_networkIndex.find(netIt.value()); if (idxIt == m_networkIndex.end()) { return nullptr; } const auto &streams = m_networks[idxIt.value()].streams; for (const auto &s : streams) { if (s.id == id) { return &s; } } return nullptr; } const QList &Store::particlesForStream(const QString &streamId) const { static const QList empty; auto it = m_particles.find(streamId); if (it == m_particles.end()) { return empty; } return it.value(); } int Store::particleCount(const QString &streamId) const { return m_particles.value(streamId).size(); } const Particle *Store::particleById(const QString &id) const { for (auto it = m_particles.constBegin(); it != m_particles.constEnd(); ++it) { for (const auto &p : it.value()) { if (p.id == id) { return &p; } } } return nullptr; } // ──────────────────────────────────────────────────────────────── // Bootstrap // ──────────────────────────────────────────────────────────────── Operation *Store::loadStartupData() { QNetworkReply *reply = NetworkManager::instance().get("/startup"); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this](const QJsonDocument &doc) { if (!doc.isObject()) { qWarning() << "Invalid startup data"; return; } populateFromStartupData(doc.object()); emit startupDataLoaded(); emit networksChanged(); }); return op; } // ──────────────────────────────────────────────────────────────── // Network operations // ──────────────────────────────────────────────────────────────── Operation *Store::createNetwork(const QString &name) { QJsonObject body; body["name"] = name; QNetworkReply *reply = NetworkManager::instance().post("/networks", QJsonDocument(body)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this](const QJsonDocument &doc) { if (doc.isObject()) { upsertNetwork(parseNetwork(doc.object())); emit networksChanged(); } }); return op; } Operation *Store::fetchNetworks() { QNetworkReply *reply = NetworkManager::instance().get("/networks"); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this](const QJsonDocument &doc) { const QJsonArray arr = doc.array(); for (const QJsonValue &v : arr) { upsertNetwork(parseNetwork(v.toObject())); } emit networksChanged(); }); return op; } Operation *Store::fetchNetwork(const QString &networkId) { QNetworkReply *reply = NetworkManager::instance().get(QString("/networks/%1").arg(networkId)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, networkId](const QJsonDocument &doc) { QJsonObject obj = doc.object(); Network network = parseNetwork(obj); const QJsonArray streamsArr = obj["streams"].toArray(); for (const QJsonValue &sv : streamsArr) { Stream stream = parseStream(sv.toObject()); network.streams.push_back(stream); m_streamToNetwork[stream.id] = network.id; } upsertNetwork(network); emit networksChanged(); }); return op; } Operation *Store::addNetworkMembers(const QString &networkId, const QStringList &emails) { QJsonObject body; QJsonArray arr; for (const QString &e : emails) { arr.append(e); } body["email_addresses"] = arr; QNetworkReply *reply = NetworkManager::instance().post( QString("/networks/%1/members").arg(networkId), QJsonDocument(body)); return new Operation(reply, this); } Operation *Store::removeNetworkMember(const QString &networkId, const QString &email) { QNetworkReply *reply = NetworkManager::instance().del( QString("/networks/%1/members/%2").arg(networkId, email)); return new Operation(reply, this); } Operation *Store::setStreamCapacity(const QString &networkId, int capacity) { QJsonObject body; body["capacity"] = capacity; QNetworkReply *reply = NetworkManager::instance().put( QString("/networks/%1/capacity").arg(networkId), QJsonDocument(body)); return new Operation(reply, this); } // ──────────────────────────────────────────────────────────────── // Stream operations // ──────────────────────────────────────────────────────────────── Operation *Store::createStream(const QString &networkId, const QString &name, const QString &description, const QString &visibility, const QStringList &members) { QJsonObject body; body["name"] = name; body["description"] = description; body["visibility"] = visibility; QJsonArray membersArr; for (const QString &m : members) { membersArr.append(m); } body["members"] = membersArr; QNetworkReply *reply = NetworkManager::instance().post( QString("/networks/%1/streams").arg(networkId), QJsonDocument(body)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, networkId](const QJsonDocument &doc) { if (doc.isObject()) { Stream stream = parseStream(doc.object()); upsertStream(networkId, stream); emit streamsChanged(networkId); } }); return op; } Operation *Store::fetchStream(const QString &streamId) { QNetworkReply *reply = NetworkManager::instance().get(QString("/streams/%1").arg(streamId)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, streamId](const QJsonDocument &doc) { QJsonObject obj = doc.object(); Stream stream = parseStream(obj); QString networkId = m_streamToNetwork.value(streamId); if (!networkId.isEmpty()) { upsertStream(networkId, stream); emit streamsChanged(networkId); } // Parse particles if present const QJsonArray particlesArr = obj["particles"].toArray(); QList particleList; for (const QJsonValue &pv : particlesArr) { particleList.append(parseParticle(pv.toObject())); } m_particles[streamId] = particleList; emit particlesChanged(streamId); }); return op; } Operation *Store::updateStream(const QString &streamId, const QString &name, const QString &description) { QJsonObject body; if (!name.isEmpty()) { body["name"] = name; } if (!description.isEmpty()) { body["description"] = description; } QNetworkReply *reply = NetworkManager::instance().patch( QString("/streams/%1").arg(streamId), QJsonDocument(body)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, streamId](const QJsonDocument &doc) { if (doc.isObject()) { Stream stream = parseStream(doc.object()); QString networkId = m_streamToNetwork.value(streamId); if (!networkId.isEmpty()) { upsertStream(networkId, stream); emit streamsChanged(networkId); } } }); return op; } Operation *Store::openStream(const QString &streamId) { QNetworkReply *reply = NetworkManager::instance().post( QString("/streams/%1/open").arg(streamId), QJsonDocument()); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, streamId](const QJsonDocument &) { QString networkId = m_streamToNetwork.value(streamId); if (m_networkIndex.contains(networkId)) { auto &streams = m_networks[m_networkIndex[networkId]].streams; for (auto &s : streams) { if (s.id == streamId) { s.isOpen = true; break; } } emit streamsChanged(networkId); } }); return op; } Operation *Store::closeStream(const QString &streamId) { QNetworkReply *reply = NetworkManager::instance().post( QString("/streams/%1/close").arg(streamId), QJsonDocument()); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, streamId](const QJsonDocument &) { QString networkId = m_streamToNetwork.value(streamId); if (m_networkIndex.contains(networkId)) { auto &streams = m_networks[m_networkIndex[networkId]].streams; for (auto &s : streams) { if (s.id == streamId) { s.isOpen = false; break; } } emit streamsChanged(networkId); } }); return op; } Operation *Store::addStreamMembers(const QString &streamId, const QStringList &emails) { QJsonObject body; QJsonArray arr; for (const QString &e : emails) { arr.append(e); } body["emails"] = arr; QNetworkReply *reply = NetworkManager::instance().post( QString("/streams/%1/members").arg(streamId), QJsonDocument(body)); return new Operation(reply, this); } Operation *Store::removeStreamMembers(const QString &streamId, const QStringList &emails) { QJsonObject body; QJsonArray arr; for (const QString &e : emails) arr.append(e); body["emails"] = arr; QNetworkReply *reply = NetworkManager::instance().del( QString("/streams/%1/members").arg(streamId), QJsonDocument(body)); return new Operation(reply, this); } // ──────────────────────────────────────────────────────────────── // Particle operations // ──────────────────────────────────────────────────────────────── Operation *Store::createParticle(const QString &streamId, const QString &type, const QMap &data) { QJsonObject body; body["type"] = type; QJsonObject dataObj; for (auto it = data.constBegin(); it != data.constEnd(); ++it) dataObj[it.key()] = QJsonValue::fromVariant(it.value()); body["data"] = dataObj; QNetworkReply *reply = NetworkManager::instance().post( QString("/streams/%1/particles").arg(streamId), QJsonDocument(body)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, streamId](const QJsonDocument &doc) { if (doc.isObject()) { Particle particle = parseParticle(doc.object()); upsertParticle(streamId, particle); emit particlesChanged(streamId); } }); return op; } Operation *Store::fetchParticle(const QString &particleId) { QNetworkReply *reply = NetworkManager::instance().get(QString("/particles/%1").arg(particleId)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, particleId](const QJsonDocument &doc) { if (doc.isObject()) { Particle particle = parseParticle(doc.object()); QString streamId = findStreamForParticle(particleId); if (!streamId.isEmpty()) { upsertParticle(streamId, particle); emit particlesChanged(streamId); } } }); return op; } Operation *Store::updateParticle(const QString &particleId, const QMap &data) { QJsonObject body; QJsonObject dataObj; for (auto it = data.constBegin(); it != data.constEnd(); ++it) dataObj[it.key()] = QJsonValue::fromVariant(it.value()); body["data"] = dataObj; QNetworkReply *reply = NetworkManager::instance().patch( QString("/particles/%1").arg(particleId), QJsonDocument(body)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, particleId](const QJsonDocument &doc) { if (doc.isObject()) { Particle particle = parseParticle(doc.object()); QString streamId = findStreamForParticle(particleId); if (!streamId.isEmpty()) { upsertParticle(streamId, particle); emit particlesChanged(streamId); } } }); return op; } Operation *Store::deleteParticle(const QString &particleId) { QNetworkReply *reply = NetworkManager::instance().del(QString("/particles/%1").arg(particleId)); auto *op = new Operation(reply, this); connect(op, &Operation::success, this, [this, particleId](const QJsonDocument &) { QString streamId = findStreamForParticle(particleId); removeParticleFromState(particleId); if (!streamId.isEmpty()) emit particlesChanged(streamId); }); return op; } Operation *Store::markParticleSeen(const QString &particleId) { // Optimistic local update for (auto it = m_particles.begin(); it != m_particles.end(); ++it) { for (int i = 0; i < it.value().size(); ++i) { if (it.value()[i].id == particleId && !it.value()[i].seen) { it.value()[i].seen = true; QString streamId = it.key(); // Decrement the parent stream's unseen count QString networkId = m_streamToNetwork.value(streamId); if (m_networkIndex.contains(networkId)) { auto &streams = m_networks[m_networkIndex[networkId]].streams; for (auto &s : streams) { if (s.id == streamId && s.unseenCount > 0) { s.unseenCount--; break; } } } emit particleUpdated(streamId, particleId); emit streamUpdated(streamId); break; } } } QNetworkReply *reply = NetworkManager::instance().post( QString("/particles/%1/seen").arg(particleId), QJsonDocument()); return new Operation(reply, this); } Operation *Store::markParticlesSeen(const QStringList &particleIds) { // Optimistic local update QSet affectedStreams; for (const QString &pid : particleIds) { for (auto it = m_particles.begin(); it != m_particles.end(); ++it) { for (int i = 0; i < it.value().size(); ++i) { if (it.value()[i].id == pid) { it.value()[i].seen = true; affectedStreams.insert(it.key()); break; } } } } for (const QString &streamId : affectedStreams) emit particlesChanged(streamId); QJsonObject body; QJsonArray ids; for (const QString &id : particleIds) ids.append(id); body["particle_ids"] = ids; QNetworkReply *reply = NetworkManager::instance().post("/particles/seen", QJsonDocument(body)); return new Operation(reply, this); } void Store::setCurrentUserEmail(const QString &email) { m_currentUserEmail = email; } QString Store::currentUserEmail() const { return m_currentUserEmail; } Operation *Store::ackParticle(const QString &particleId) { // Optimistic local update for (auto it = m_particles.begin(); it != m_particles.end(); ++it) { for (int i = 0; i < it.value().size(); ++i) { if (it.value()[i].id == particleId) { auto &emails = it.value()[i].ackedByEmails; if (!m_currentUserEmail.isEmpty() && std::find(emails.begin(), emails.end(), m_currentUserEmail) == emails.end()) { emails.push_back(m_currentUserEmail); emit particleUpdated(it.key(), particleId); } break; } } } QNetworkReply *reply = NetworkManager::instance().post( QString("/particles/%1/ack").arg(particleId), QJsonDocument()); return new Operation(reply, this); }