commit 5cc5ac6f1da2ce992b2f96ae2728332bfe5229e5
parent bc21606bf9288440732b083590fc3a91ac4f1eb6
Author: Silas Brack <silasbrack@gmail.com>
Date: Fri, 26 Dec 2025 14:59:58 +0100
Add realtime SSE subscriptions on chats
Diffstat:
4 files changed, 365 insertions(+), 1 deletion(-)
diff --git a/networkmanager.cpp b/networkmanager.cpp
@@ -9,7 +9,18 @@
NetworkManager::NetworkManager(QObject *parent)
: QObject(parent)
, m_manager(new QNetworkAccessManager(this))
+ , m_sseReply(nullptr)
+ , m_currentSseChat(-1)
+ , m_reconnectTimer(new QTimer(this))
+ , m_reconnectAttempts(0)
{
+ m_reconnectTimer->setSingleShot(true);
+ connect(m_reconnectTimer, &QTimer::timeout, this, &NetworkManager::attemptReconnect);
+}
+
+NetworkManager::~NetworkManager()
+{
+ unsubscribeFromChat();
}
void NetworkManager::postChat(const QString& name)
@@ -390,3 +401,253 @@ void NetworkManager::onFetchMessagesFinished()
// Clean up
reply->deleteLater();
}
+
+// SSE Methods
+
+void NetworkManager::subscribeToChat(int chatId)
+{
+ if (chatId < 0) return;
+
+ // Unsubscribe from previous chat if any
+ if (m_currentSseChat != chatId) {
+ unsubscribeFromChat();
+ }
+
+ // Build SSE URL with filter
+ QString urlStr = QString(
+ "http://localhost:4000/api/records/v1/message/subscribe/*"
+ "?filter[chat_id][$eq]=%1"
+ ).arg(chatId);
+
+ QUrl url(urlStr);
+ QNetworkRequest request(url);
+
+ // SSE requires Accept header
+ request.setRawHeader("Accept", "text/event-stream");
+ request.setRawHeader("Cache-Control", "no-cache");
+
+ // Create persistent GET request
+ m_sseReply = m_manager->get(request);
+ m_currentSseChat = chatId;
+
+ // Connect signals
+ connect(m_sseReply, &QNetworkReply::readyRead,
+ this, &NetworkManager::onSseReadyRead);
+ connect(m_sseReply, &QNetworkReply::finished,
+ this, &NetworkManager::onSseFinished);
+ connect(m_sseReply,
+ QOverload<QNetworkReply::NetworkError>::of(&QNetworkReply::errorOccurred),
+ this, &NetworkManager::onSseError);
+
+ qDebug() << "SSE: Subscribed to chat" << chatId << "URL:" << urlStr;
+ emit sseConnected(chatId);
+ resetReconnectBackoff();
+}
+
+void NetworkManager::unsubscribeFromChat()
+{
+ if (!m_sseReply) return;
+
+ int previousChat = m_currentSseChat;
+
+ // Disconnect signals
+ disconnect(m_sseReply, nullptr, this, nullptr);
+
+ // Abort and cleanup
+ m_sseReply->abort();
+ m_sseReply->deleteLater();
+ m_sseReply = nullptr;
+
+ // Reset state
+ m_currentSseChat = -1;
+ m_sseBuffer.clear();
+ m_recentMessageIds.clear();
+ m_reconnectTimer->stop();
+
+ qDebug() << "SSE: Unsubscribed from chat" << previousChat;
+ emit sseDisconnected(previousChat, "User initiated");
+}
+
+void NetworkManager::onSseReadyRead()
+{
+ if (!m_sseReply) return;
+
+ // Read available data and append to buffer
+ QByteArray chunk = m_sseReply->readAll();
+ m_sseBuffer.append(chunk);
+
+ // Process complete events (delimited by \n\n)
+ while (m_sseBuffer.contains("\n\n")) {
+ int eventEnd = m_sseBuffer.indexOf("\n\n");
+ QByteArray eventData = m_sseBuffer.left(eventEnd);
+ m_sseBuffer.remove(0, eventEnd + 2); // Remove event + delimiter
+
+ // Parse the event
+ QString eventStr = QString::fromUtf8(eventData);
+ parseSSEEvent(eventStr);
+ }
+
+ // Prevent buffer from growing indefinitely
+ if (m_sseBuffer.size() > 65536) { // 64KB limit
+ qWarning() << "SSE: Buffer overflow, clearing";
+ m_sseBuffer.clear();
+ }
+}
+
+void NetworkManager::parseSSEEvent(const QString& eventData)
+{
+ QStringList lines = eventData.split('\n');
+
+ QString eventType;
+ QStringList dataLines;
+ QString eventId;
+
+ // Parse event fields
+ for (const QString& line : lines) {
+ if (line.startsWith("event:")) {
+ eventType = line.mid(6).trimmed();
+ }
+ else if (line.startsWith("data:")) {
+ dataLines.append(line.mid(5).trimmed());
+ }
+ else if (line.startsWith("id:")) {
+ eventId = line.mid(3).trimmed();
+ }
+ else if (line.startsWith(":")) {
+ // Comment, ignore
+ continue;
+ }
+ }
+
+ // Combine multi-line data
+ QString data = dataLines.join("\n");
+ if (data.isEmpty()) return;
+
+ qDebug() << "SSE: Received event, type:" << eventType << "data:" << data;
+
+ // Parse JSON data
+ QJsonDocument jsonDoc = QJsonDocument::fromJson(data.toUtf8());
+ if (!jsonDoc.isObject()) {
+ qWarning() << "SSE: Invalid JSON in event:" << data;
+ return;
+ }
+
+ QJsonObject jsonObj = jsonDoc.object();
+
+ // Trailbase wraps the actual data inside an "Insert" key
+ QVariantMap messageData;
+ if (jsonObj.contains("Insert")) {
+ messageData = jsonObj["Insert"].toObject().toVariantMap();
+ } else {
+ // Fallback to direct parsing if no wrapper
+ messageData = jsonObj.toVariantMap();
+ }
+
+ // Extract message ID for deduplication
+ int messageId = messageData.value("id", 0).toInt();
+ if (messageId <= 0) {
+ qWarning() << "SSE: Message without valid ID";
+ return;
+ }
+
+ // Check for duplicates
+ if (isDuplicateMessage(messageId)) {
+ qDebug() << "SSE: Ignoring duplicate message" << messageId;
+ return;
+ }
+
+ // Verify chat_id matches subscription
+ int messageChatId = messageData.value("chat_id", -1).toInt();
+ if (messageChatId != m_currentSseChat) {
+ qWarning() << "SSE: Message for wrong chat" << messageChatId
+ << "expected" << m_currentSseChat;
+ return;
+ }
+
+ // Cache message ID
+ addToRecentMessages(messageId);
+
+ // Emit to QML layer
+ qDebug() << "SSE: New message received:" << messageId;
+ emit messageReceived(messageData);
+}
+
+bool NetworkManager::isDuplicateMessage(int messageId)
+{
+ return m_recentMessageIds.contains(messageId);
+}
+
+void NetworkManager::addToRecentMessages(int messageId)
+{
+ m_recentMessageIds.insert(messageId);
+
+ // Keep cache bounded (FIFO eviction)
+ // QSet doesn't guarantee order, so when full, clear half
+ if (m_recentMessageIds.size() > 50) {
+ auto it = m_recentMessageIds.begin();
+ for (int i = 0; i < 25 && it != m_recentMessageIds.end(); ++i) {
+ it = m_recentMessageIds.erase(it);
+ }
+ }
+}
+
+void NetworkManager::onSseFinished()
+{
+ if (!m_sseReply) return;
+
+ int chatId = m_currentSseChat;
+ QString reason = "Connection closed by server";
+
+ qDebug() << "SSE: Connection finished for chat" << chatId;
+
+ // Cleanup
+ m_sseReply->deleteLater();
+ m_sseReply = nullptr;
+ m_sseBuffer.clear();
+
+ emit sseDisconnected(chatId, reason);
+
+ // Attempt reconnection if chat is still active
+ if (m_currentSseChat >= 0) {
+ scheduleReconnect();
+ } else {
+ m_currentSseChat = -1;
+ }
+}
+
+void NetworkManager::onSseError(QNetworkReply::NetworkError error)
+{
+ QString errorMsg = m_sseReply ? m_sseReply->errorString() : "Unknown error";
+ qWarning() << "SSE: Network error:" << error << errorMsg;
+
+ emit sseError(errorMsg);
+
+ // Reconnection will be handled by onSseFinished
+}
+
+void NetworkManager::attemptReconnect()
+{
+ if (m_currentSseChat < 0) return;
+
+ int chatId = m_currentSseChat;
+ qDebug() << "SSE: Attempting reconnect to chat" << chatId
+ << "attempt" << m_reconnectAttempts;
+
+ subscribeToChat(chatId);
+}
+
+void NetworkManager::scheduleReconnect()
+{
+ // Exponential backoff: 1s, 2s, 4s, 8s, 16s, max 30s
+ int delay = qMin(1000 * (1 << m_reconnectAttempts), 30000);
+ m_reconnectAttempts++;
+
+ qDebug() << "SSE: Scheduling reconnect in" << delay << "ms";
+ m_reconnectTimer->start(delay);
+}
+
+void NetworkManager::resetReconnectBackoff()
+{
+ m_reconnectAttempts = 0;
+ m_reconnectTimer->stop();
+}
diff --git a/networkmanager.h b/networkmanager.h
@@ -5,6 +5,8 @@
#include <QNetworkAccessManager>
#include <QNetworkReply>
#include <QVariantMap>
+#include <QTimer>
+#include <QSet>
class NetworkManager : public QObject
{
@@ -12,12 +14,17 @@ class NetworkManager : public QObject
public:
explicit NetworkManager(QObject *parent = nullptr);
+ ~NetworkManager();
Q_INVOKABLE void postChat(const QString& name);
Q_INVOKABLE void fetchChats(const QString& cursor = "");
Q_INVOKABLE void postMessage(int chatId, const QString& content);
Q_INVOKABLE void fetchMessages(int chatId, const QString& cursor = "");
+ // SSE methods
+ Q_INVOKABLE void subscribeToChat(int chatId);
+ Q_INVOKABLE void unsubscribeFromChat();
+
signals:
void chatCreated(QVariantMap chatData);
void chatsFetched(QVariantList chats, QString cursor);
@@ -25,14 +32,41 @@ signals:
void messagesFetched(QVariantList messages, QString cursor);
void requestFailed(QString error);
+ // SSE signals
+ void messageReceived(QVariantMap messageData);
+ void sseConnected(int chatId);
+ void sseDisconnected(int chatId, QString reason);
+ void sseError(QString error);
+
private slots:
void onPostChatFinished();
void onFetchChatsFinished();
void onPostMessageFinished();
void onFetchMessagesFinished();
+ // SSE slots
+ void onSseReadyRead();
+ void onSseFinished();
+ void onSseError(QNetworkReply::NetworkError error);
+ void attemptReconnect();
+
private:
+ // Helper methods for SSE
+ void parseSSEEvent(const QString& eventData);
+ bool isDuplicateMessage(int messageId);
+ void addToRecentMessages(int messageId);
+ void scheduleReconnect();
+ void resetReconnectBackoff();
+
QNetworkAccessManager* m_manager;
+
+ // SSE members
+ QNetworkReply* m_sseReply;
+ int m_currentSseChat;
+ QByteArray m_sseBuffer;
+ QSet<int> m_recentMessageIds;
+ QTimer* m_reconnectTimer;
+ int m_reconnectAttempts;
};
#endif // NETWORKMANAGER_H
diff --git a/qml/models/MessageListModel.qml b/qml/models/MessageListModel.qml
@@ -78,6 +78,44 @@ Item {
var isInitial = (model.count === 0)
root.loadMessages(messages, cursor, isInitial)
}
+
+ // SSE real-time message handler
+ function onMessageReceived(messageData) {
+ // Only process if message belongs to this chat
+ var messageChatId = messageData.chat_id
+ if (messageChatId !== root.chatId) {
+ console.log("SSE: Ignoring message for chat", messageChatId)
+ return
+ }
+
+ console.log("SSE: Adding real-time message to chat", root.chatId)
+
+ // Transform and add to model (prepend for newest-first)
+ model.insert(0, {
+ id: messageData.id || 0,
+ senderId: messageData.sender_id || 0,
+ text: messageData.content || "",
+ timestamp: messageData.sent_at || new Date().toISOString(),
+ isOutgoing: false,
+ isRead: false
+ })
+ }
+
+ function onSseConnected(chatId) {
+ if (chatId === root.chatId) {
+ console.log("SSE: Connected to chat", chatId)
+ }
+ }
+
+ function onSseDisconnected(chatId, reason) {
+ if (chatId === root.chatId) {
+ console.log("SSE: Disconnected from chat", chatId, "reason:", reason)
+ }
+ }
+
+ function onSseError(error) {
+ console.error("SSE: Error:", error)
+ }
}
onChatIdChanged: {
@@ -88,7 +126,22 @@ Item {
currentCursor = ""
hasMore = true
isLoadingMore = true
+
+ // Fetch initial messages
NetworkManager.fetchMessages(chatId, "")
+
+ // Subscribe to real-time updates
+ NetworkManager.subscribeToChat(chatId)
+ } else {
+ // Unsubscribe when no chat is active
+ NetworkManager.unsubscribeFromChat()
+ }
+ }
+
+ Component.onDestruction: {
+ // Ensure we unsubscribe when model is destroyed
+ if (chatId >= 0) {
+ NetworkManager.unsubscribeFromChat()
}
}
}
diff --git a/qml/views/MessageThreadView.qml b/qml/views/MessageThreadView.qml
@@ -9,6 +9,7 @@ Item {
property int chatId: -1
property bool isSending: false
+ property bool sseConnected: false
// Dynamic message list model
MessageListModel {
@@ -28,6 +29,18 @@ Item {
console.error("Failed:", error)
root.isSending = false
}
+
+ // Track SSE connection status
+ function onSseConnected(chatId) {
+ if (chatId === root.chatId) {
+ root.sseConnected = true
+ }
+ }
+ function onSseDisconnected(chatId, reason) {
+ if (chatId === root.chatId) {
+ root.sseConnected = false
+ }
+ }
}
ColumnLayout {
@@ -37,7 +50,10 @@ Item {
// Top bar
TopBar {
Layout.fillWidth: true
- title: chatId >= 0 ? "Chat " + chatId : "Chat"
+ title: {
+ var baseTitle = chatId >= 0 ? "Chat " + chatId : "Chat"
+ return root.sseConnected ? baseTitle + " • Online" : baseTitle
+ }
showBackButton: true
showMenuButton: true
onBackClicked: {