From d36afdf07216b13cbeab5ba2251c0e3404802f00 Mon Sep 17 00:00:00 2001 From: Ricardo Guilherme Schmidt <3esmit@gmail.com> Date: Fri, 17 Jul 2026 20:09:27 -0300 Subject: [PATCH] fix(amm): force fresh post-transaction reads --- apps/amm/src/SequencerClient.cpp | 18 +++ apps/amm/src/SequencerClient.h | 3 + apps/amm/tests/cpp/NewPositionRuntimeTest.cpp | 111 ++++++++++++++++-- 3 files changed, 120 insertions(+), 12 deletions(-) diff --git a/apps/amm/src/SequencerClient.cpp b/apps/amm/src/SequencerClient.cpp index 40254fe..54083ef 100644 --- a/apps/amm/src/SequencerClient.cpp +++ b/apps/amm/src/SequencerClient.cpp @@ -167,6 +167,10 @@ void SequencerClient::readAccount(const QString& accountId, [callback = std::move(callback), cached]() mutable { callback(cached); }); return; } + if (forceRefresh && m_activeReadIds.contains(accountId)) { + m_forcedWaiters[accountId].append(std::move(callback)); + return; + } const bool alreadyPending = m_waiters.contains(accountId); m_waiters[accountId].append(std::move(callback)); @@ -180,6 +184,7 @@ void SequencerClient::startPendingReads() while (m_activeReads < MAX_CONCURRENT_READS && !m_pending.isEmpty()) { const PendingRead pending = m_pending.dequeue(); ++m_activeReads; + m_activeReadIds.insert(pending.accountId); startRead(pending.accountId); } } @@ -226,7 +231,14 @@ void SequencerClient::completeRead(const QString& accountId, if (read.ok()) m_cache.insert(accountId, read); const QVector callbacks = m_waiters.take(accountId); + const QVector forcedCallbacks = m_forcedWaiters.take(accountId); + m_activeReadIds.remove(accountId); --m_activeReads; + if (!forcedCallbacks.isEmpty()) { + m_cache.remove(accountId); + m_waiters.insert(accountId, forcedCallbacks); + m_pending.enqueue({ accountId }); + } for (const AccountCallback& callback : callbacks) callback(read); startPendingReads(); @@ -279,7 +291,13 @@ void SequencerClient::cancelPendingReads() { decltype(m_waiters) waiters; waiters.swap(m_waiters); + for (auto iterator = m_forcedWaiters.cbegin(); + iterator != m_forcedWaiters.cend(); ++iterator) { + waiters[iterator.key()].append(iterator.value()); + } + m_forcedWaiters.clear(); m_pending.clear(); + m_activeReadIds.clear(); m_activeReads = 0; for (auto iterator = waiters.cbegin(); iterator != waiters.cend(); ++iterator) { const WalletAccountRead failed { iterator.key() }; diff --git a/apps/amm/src/SequencerClient.h b/apps/amm/src/SequencerClient.h index 5062999..cbcd789 100644 --- a/apps/amm/src/SequencerClient.h +++ b/apps/amm/src/SequencerClient.h @@ -5,6 +5,7 @@ #include #include #include +#include #include #include #include @@ -55,7 +56,9 @@ private: QByteArray m_authorization; QHash m_cache; QHash> m_waiters; + QHash> m_forcedWaiters; QQueue m_pending; + QSet m_activeReadIds; int m_activeReads = 0; quint64 m_generation = 0; }; diff --git a/apps/amm/tests/cpp/NewPositionRuntimeTest.cpp b/apps/amm/tests/cpp/NewPositionRuntimeTest.cpp index 9b92657..e83b255 100644 --- a/apps/amm/tests/cpp/NewPositionRuntimeTest.cpp +++ b/apps/amm/tests/cpp/NewPositionRuntimeTest.cpp @@ -5,6 +5,7 @@ #include #include +#include #include #include #include @@ -51,8 +52,39 @@ namespace { int requestCount() const { return m_requestCount; } void failNextRequest() { ++m_failuresRemaining; } + void holdResponses() { m_holdResponses = true; } + void releaseNextResponse() + { + if (m_heldResponses.isEmpty()) + return; + const HeldResponse held = m_heldResponses.takeFirst(); + respond(held.socket, held.fail); + } private: + struct HeldResponse { + QTcpSocket* socket = nullptr; + bool fail = false; + }; + + static void respond(QTcpSocket* socket, bool fail) + { + if (!socket) + return; + const QByteArray payload = QByteArrayLiteral( + R"({"jsonrpc":"2.0","id":1,"result":null})"); + QByteArray response = fail + ? QByteArrayLiteral( + "HTTP/1.1 500 Internal Server Error\r\nContent-Type: application/json\r\nContent-Length: ") + : QByteArrayLiteral( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: "); + response += QByteArray::number(payload.size()); + response += QByteArrayLiteral("\r\nConnection: close\r\n\r\n"); + response += payload; + socket->write(response); + socket->disconnectFromHost(); + } + void process(QTcpSocket* socket) { QByteArray& request = m_requests[socket]; @@ -77,27 +109,40 @@ namespace { const bool fail = m_failuresRemaining > 0; if (fail) --m_failuresRemaining; - const QByteArray payload = QByteArrayLiteral( - R"({"jsonrpc":"2.0","id":1,"result":null})"); - QByteArray response = fail - ? QByteArrayLiteral( - "HTTP/1.1 500 Internal Server Error\r\nContent-Type: application/json\r\nContent-Length: ") - : QByteArrayLiteral( - "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: "); - response += QByteArray::number(payload.size()); - response += QByteArrayLiteral("\r\nConnection: close\r\n\r\n"); - response += payload; - socket->write(response); - socket->disconnectFromHost(); m_requests.remove(socket); + if (m_holdResponses) { + m_heldResponses.append({ socket, fail }); + return; + } + respond(socket, fail); } QTcpServer m_server; QHash m_requests; int m_requestCount = 0; int m_failuresRemaining = 0; + bool m_holdResponses = false; + QVector m_heldResponses; }; + bool waitForRequestCount(const LocalRpcServer& server, int expected) + { + QElapsedTimer timer; + timer.start(); + while (server.requestCount() < expected && timer.elapsed() < 3000) + QCoreApplication::processEvents(QEventLoop::AllEvents, 10); + return server.requestCount() >= expected; + } + + bool waitForCallbacks(const bool& first, const bool& second) + { + QElapsedTimer timer; + timer.start(); + while ((!first || !second) && timer.elapsed() < 3000) + QCoreApplication::processEvents(QEventLoop::AllEvents, 10); + return first && second; + } + class FakeWallet final : public WalletProvider { public: WalletSession connect(const WalletPaths&) override @@ -481,5 +526,47 @@ int main(int argc, char** argv) "same-path endpoint update should clear cache and use new sequencer")) return 1; + LocalRpcServer forcedRefreshServer; + if (!expect(forcedRefreshServer.listen(), + "forced-refresh sequencer should listen")) + return 1; + forcedRefreshServer.holdResponses(); + if (!expect(sequencerConfig.resize(0) && sequencerConfig.seek(0), + "forced-refresh config should rewind")) + return 1; + sequencerConfig.write(QJsonDocument(QJsonObject { + { QStringLiteral("sequencer_addr"), forcedRefreshServer.endpoint() }, + }).toJson(QJsonDocument::Compact)); + sequencerConfig.flush(); + if (!expect(sequencer.configure(sequencerConfig.fileName()), + "forced-refresh sequencer should configure")) + return 1; + + bool ordinaryCompleted = false; + bool forcedCompleted = false; + sequencer.readAccounts({ holding.address }, false, + [&](QVector) { ordinaryCompleted = true; }); + if (!expect(waitForRequestCount(forcedRefreshServer, 1), + "ordinary read should reach sequencer")) + return 1; + sequencer.readAccounts({ holding.address }, true, + [&](QVector) { forcedCompleted = true; }); + QCoreApplication::processEvents(QEventLoop::AllEvents, 10); + if (!expect(forcedRefreshServer.requestCount() == 1, + "forced refresh should wait for active account read")) + return 1; + + forcedRefreshServer.releaseNextResponse(); + if (!expect(waitForRequestCount(forcedRefreshServer, 2), + "forced refresh should issue a follow-up account read")) + return 1; + if (!expect(ordinaryCompleted && !forcedCompleted, + "pre-refresh result must not complete forced callback")) + return 1; + forcedRefreshServer.releaseNextResponse(); + if (!expect(waitForCallbacks(ordinaryCompleted, forcedCompleted), + "forced follow-up read should complete")) + return 1; + return 0; }