mirror of
https://github.com/logos-blockchain/lez-programs.git
synced 2026-07-30 02:33:24 +00:00
fix(amm): force fresh post-transaction reads
This commit is contained in:
parent
1f70ed51ba
commit
d36afdf072
@ -167,6 +167,10 @@ void SequencerClient::readAccount(const QString& accountId,
|
|||||||
[callback = std::move(callback), cached]() mutable { callback(cached); });
|
[callback = std::move(callback), cached]() mutable { callback(cached); });
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
if (forceRefresh && m_activeReadIds.contains(accountId)) {
|
||||||
|
m_forcedWaiters[accountId].append(std::move(callback));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
const bool alreadyPending = m_waiters.contains(accountId);
|
const bool alreadyPending = m_waiters.contains(accountId);
|
||||||
m_waiters[accountId].append(std::move(callback));
|
m_waiters[accountId].append(std::move(callback));
|
||||||
@ -180,6 +184,7 @@ void SequencerClient::startPendingReads()
|
|||||||
while (m_activeReads < MAX_CONCURRENT_READS && !m_pending.isEmpty()) {
|
while (m_activeReads < MAX_CONCURRENT_READS && !m_pending.isEmpty()) {
|
||||||
const PendingRead pending = m_pending.dequeue();
|
const PendingRead pending = m_pending.dequeue();
|
||||||
++m_activeReads;
|
++m_activeReads;
|
||||||
|
m_activeReadIds.insert(pending.accountId);
|
||||||
startRead(pending.accountId);
|
startRead(pending.accountId);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -226,7 +231,14 @@ void SequencerClient::completeRead(const QString& accountId,
|
|||||||
if (read.ok())
|
if (read.ok())
|
||||||
m_cache.insert(accountId, read);
|
m_cache.insert(accountId, read);
|
||||||
const QVector<AccountCallback> callbacks = m_waiters.take(accountId);
|
const QVector<AccountCallback> callbacks = m_waiters.take(accountId);
|
||||||
|
const QVector<AccountCallback> forcedCallbacks = m_forcedWaiters.take(accountId);
|
||||||
|
m_activeReadIds.remove(accountId);
|
||||||
--m_activeReads;
|
--m_activeReads;
|
||||||
|
if (!forcedCallbacks.isEmpty()) {
|
||||||
|
m_cache.remove(accountId);
|
||||||
|
m_waiters.insert(accountId, forcedCallbacks);
|
||||||
|
m_pending.enqueue({ accountId });
|
||||||
|
}
|
||||||
for (const AccountCallback& callback : callbacks)
|
for (const AccountCallback& callback : callbacks)
|
||||||
callback(read);
|
callback(read);
|
||||||
startPendingReads();
|
startPendingReads();
|
||||||
@ -279,7 +291,13 @@ void SequencerClient::cancelPendingReads()
|
|||||||
{
|
{
|
||||||
decltype(m_waiters) waiters;
|
decltype(m_waiters) waiters;
|
||||||
waiters.swap(m_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_pending.clear();
|
||||||
|
m_activeReadIds.clear();
|
||||||
m_activeReads = 0;
|
m_activeReads = 0;
|
||||||
for (auto iterator = waiters.cbegin(); iterator != waiters.cend(); ++iterator) {
|
for (auto iterator = waiters.cbegin(); iterator != waiters.cend(); ++iterator) {
|
||||||
const WalletAccountRead failed { iterator.key() };
|
const WalletAccountRead failed { iterator.key() };
|
||||||
|
|||||||
@ -5,6 +5,7 @@
|
|||||||
#include <QHash>
|
#include <QHash>
|
||||||
#include <QObject>
|
#include <QObject>
|
||||||
#include <QQueue>
|
#include <QQueue>
|
||||||
|
#include <QSet>
|
||||||
#include <QStringList>
|
#include <QStringList>
|
||||||
#include <QUrl>
|
#include <QUrl>
|
||||||
#include <QVector>
|
#include <QVector>
|
||||||
@ -55,7 +56,9 @@ private:
|
|||||||
QByteArray m_authorization;
|
QByteArray m_authorization;
|
||||||
QHash<QString, WalletAccountRead> m_cache;
|
QHash<QString, WalletAccountRead> m_cache;
|
||||||
QHash<QString, QVector<AccountCallback>> m_waiters;
|
QHash<QString, QVector<AccountCallback>> m_waiters;
|
||||||
|
QHash<QString, QVector<AccountCallback>> m_forcedWaiters;
|
||||||
QQueue<PendingRead> m_pending;
|
QQueue<PendingRead> m_pending;
|
||||||
|
QSet<QString> m_activeReadIds;
|
||||||
int m_activeReads = 0;
|
int m_activeReads = 0;
|
||||||
quint64 m_generation = 0;
|
quint64 m_generation = 0;
|
||||||
};
|
};
|
||||||
|
|||||||
@ -5,6 +5,7 @@
|
|||||||
|
|
||||||
#include <QCoreApplication>
|
#include <QCoreApplication>
|
||||||
#include <QDateTime>
|
#include <QDateTime>
|
||||||
|
#include <QElapsedTimer>
|
||||||
#include <QEventLoop>
|
#include <QEventLoop>
|
||||||
#include <QHash>
|
#include <QHash>
|
||||||
#include <QHostAddress>
|
#include <QHostAddress>
|
||||||
@ -51,8 +52,39 @@ namespace {
|
|||||||
|
|
||||||
int requestCount() const { return m_requestCount; }
|
int requestCount() const { return m_requestCount; }
|
||||||
void failNextRequest() { ++m_failuresRemaining; }
|
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:
|
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)
|
void process(QTcpSocket* socket)
|
||||||
{
|
{
|
||||||
QByteArray& request = m_requests[socket];
|
QByteArray& request = m_requests[socket];
|
||||||
@ -77,27 +109,40 @@ namespace {
|
|||||||
const bool fail = m_failuresRemaining > 0;
|
const bool fail = m_failuresRemaining > 0;
|
||||||
if (fail)
|
if (fail)
|
||||||
--m_failuresRemaining;
|
--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);
|
m_requests.remove(socket);
|
||||||
|
if (m_holdResponses) {
|
||||||
|
m_heldResponses.append({ socket, fail });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
respond(socket, fail);
|
||||||
}
|
}
|
||||||
|
|
||||||
QTcpServer m_server;
|
QTcpServer m_server;
|
||||||
QHash<QTcpSocket*, QByteArray> m_requests;
|
QHash<QTcpSocket*, QByteArray> m_requests;
|
||||||
int m_requestCount = 0;
|
int m_requestCount = 0;
|
||||||
int m_failuresRemaining = 0;
|
int m_failuresRemaining = 0;
|
||||||
|
bool m_holdResponses = false;
|
||||||
|
QVector<HeldResponse> 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 {
|
class FakeWallet final : public WalletProvider {
|
||||||
public:
|
public:
|
||||||
WalletSession connect(const WalletPaths&) override
|
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"))
|
"same-path endpoint update should clear cache and use new sequencer"))
|
||||||
return 1;
|
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<WalletAccountRead>) { ordinaryCompleted = true; });
|
||||||
|
if (!expect(waitForRequestCount(forcedRefreshServer, 1),
|
||||||
|
"ordinary read should reach sequencer"))
|
||||||
|
return 1;
|
||||||
|
sequencer.readAccounts({ holding.address }, true,
|
||||||
|
[&](QVector<WalletAccountRead>) { 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;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user