Files
llink/cpp/src/store.cpp
T

823 lines
25 KiB
C++

#include "store.h"
#include "networkmanager.h"
#include <QDebug>
#include <QJsonArray>
#include <QJsonDocument>
#include <QJsonObject>
#include <QNetworkReply>
// ────────────────────────────────────────────────────────────────
// 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<Particle> 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<Stream> 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<Particle> &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<Particle> &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<Network> &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<Stream> 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<Stream> result;
result.reserve(static_cast<int>(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<int>(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<Particle> &Store::particlesForStream(const QString &streamId) const
{
static const QList<Particle> 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<Particle> 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<QString, QVariant> &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<QString, QVariant> &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<QString> 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);
}