mirror of
https://github.com/Xahau/xahaud.git
synced 2026-08-25 17:20:53 +00:00
Compare commits
21 Commits
jsontx
...
subscripti
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
93910f6295 | ||
|
|
15181211df | ||
|
|
c1488a7f0b | ||
|
|
12e8e79b1f | ||
|
|
d214344616 | ||
|
|
db7ecfb86a | ||
|
|
280659c32b | ||
|
|
827177d009 | ||
|
|
67eebc686a | ||
|
|
60401fcb40 | ||
|
|
2d7b17ae1f | ||
|
|
159e00570d | ||
|
|
202353ac8c | ||
|
|
12cdcbcecf | ||
|
|
1855350b65 | ||
|
|
41db993a20 | ||
|
|
7241fed6d0 | ||
|
|
a6dfb40413 | ||
|
|
3e0d6b9cd2 | ||
|
|
a8388e48a4 | ||
|
|
49908096d5 |
@@ -77,6 +77,11 @@ test.ledger > xrpld.app
|
||||
test.ledger > xrpld.core
|
||||
test.ledger > xrpld.ledger
|
||||
test.ledger > xrpl.protocol
|
||||
test.net > test.toplevel
|
||||
test.net > xrpl.basics
|
||||
test.net > xrpld.core
|
||||
test.net > xrpld.net
|
||||
test.net > xrpl.json
|
||||
test.nodestore > test.jtx
|
||||
test.nodestore > test.toplevel
|
||||
test.nodestore > test.unit_test
|
||||
|
||||
@@ -124,9 +124,6 @@ find_package(date REQUIRED)
|
||||
find_package(xxHash REQUIRED)
|
||||
find_package(magic_enum REQUIRED)
|
||||
|
||||
find_package(fmt REQUIRED)
|
||||
target_link_libraries(ripple_libs INTERFACE fmt::fmt)
|
||||
|
||||
include(deps/WasmEdge)
|
||||
if(TARGET nudb::core)
|
||||
set(nudb nudb::core)
|
||||
|
||||
@@ -95,16 +95,8 @@ if [[ "$4" == "" ]]; then
|
||||
echo "Non GH, local building, no Action runner magic"
|
||||
else
|
||||
# GH Action, runner
|
||||
if [[ "$(git rev-parse --abbrev-ref HEAD)" == "release" ]]; then
|
||||
echo "building on the release branch... placing it in builds/candidate"
|
||||
mkdir /data/builds/candidate
|
||||
cp /io/release-build/xahaud /data/builds/candidate/$(date +%Y).$(date +%-m).$(date +%-d)-$(git rev-parse --abbrev-ref HEAD)+$4
|
||||
cp /io/release-build/release.info /data/builds/candidate/$(date +%Y).$(date +%-m).$(date +%-d)-$(git rev-parse --abbrev-ref HEAD)+$4.releaseinfo
|
||||
else
|
||||
echo "building non-release branch, placing it in builds root"
|
||||
cp /io/release-build/xahaud /data/builds/$(date +%Y).$(date +%-m).$(date +%-d)-$(git rev-parse --abbrev-ref HEAD)+$4
|
||||
cp /io/release-build/release.info /data/builds/$(date +%Y).$(date +%-m).$(date +%-d)-$(git rev-parse --abbrev-ref HEAD)+$4.releaseinfo
|
||||
fi
|
||||
cp /io/release-build/xahaud /data/builds/$(date +%Y).$(date +%-m).$(date +%-d)-$(git rev-parse --abbrev-ref HEAD)+$4
|
||||
cp /io/release-build/release.info /data/builds/$(date +%Y).$(date +%-m).$(date +%-d)-$(git rev-parse --abbrev-ref HEAD)+$4.releaseinfo
|
||||
echo "Published build to: http://build.xahau.tech/"
|
||||
echo $(date +%Y).$(date +%-m).$(date +%-d)-$(git rev-parse --abbrev-ref HEAD)+$4
|
||||
fi
|
||||
|
||||
@@ -35,7 +35,6 @@ class Xrpl(ConanFile):
|
||||
'soci/4.0.3@xahaud/stable',
|
||||
'xxhash/0.8.2',
|
||||
'zlib/1.3.1',
|
||||
'fmt/12.1.0',
|
||||
]
|
||||
|
||||
tool_requires = [
|
||||
@@ -192,7 +191,6 @@ class Xrpl(ConanFile):
|
||||
'sqlite3::sqlite',
|
||||
'xxhash::xxhash',
|
||||
'zlib::zlib',
|
||||
'fmt::fmt',
|
||||
]
|
||||
if self.options.rocksdb:
|
||||
libxrpl.requires.append('rocksdb::librocksdb')
|
||||
|
||||
@@ -1,557 +0,0 @@
|
||||
//------------------------------------------------------------------------------
|
||||
/*
|
||||
This file is part of rippled: https://github.com/ripple/rippled
|
||||
Copyright (c) 2012-2014 Ripple Labs Inc.
|
||||
|
||||
Permission to use, copy, modify, and/or distribute this software for any
|
||||
purpose with or without fee is hereby granted, provided that the above
|
||||
copyright notice and this permission notice appear in all copies.
|
||||
|
||||
THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES
|
||||
WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF
|
||||
MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR
|
||||
ANY SPECIAL , DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES
|
||||
WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN
|
||||
ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF
|
||||
OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE.
|
||||
*/
|
||||
//==============================================================================
|
||||
|
||||
#ifndef RIPPLE_PROTOCOL_JSONTXSIGNATURES_H_INCLUDED
|
||||
#define RIPPLE_PROTOCOL_JSONTXSIGNATURES_H_INCLUDED
|
||||
|
||||
#include <xrpl/json/json_reader.h>
|
||||
#include <xrpl/json/json_writer.h>
|
||||
#include <xrpl/protocol/PublicKey.h>
|
||||
#include <xrpl/protocol/SField.h>
|
||||
#include <xrpl/protocol/STParsedJSON.h>
|
||||
#include <boost/algorithm/string.hpp>
|
||||
#include <algorithm>
|
||||
#include <cctype>
|
||||
#include <cmath>
|
||||
#include <cstdint>
|
||||
#include <fmt/format.h>
|
||||
#include <limits>
|
||||
#include <map>
|
||||
#include <functional>
|
||||
#include <unordered_map>
|
||||
|
||||
namespace ripple {
|
||||
|
||||
//------------------------------------------------------------------------------
|
||||
// jsontx: plaintext-JSON signing support (featureJsonTx)
|
||||
//
|
||||
// The delta is attacker-controlled: it arrives over the wire beside a binary
|
||||
// transaction and is not covered by the signature it helps reconstruct. Every
|
||||
// bound below is therefore explicit, and unsanitize_jsontx accepts only the
|
||||
// exact encoding sanitize_jsontx would have produced.
|
||||
//------------------------------------------------------------------------------
|
||||
|
||||
static constexpr std::size_t jsontx_max_text = 8192; // canonical and original
|
||||
static constexpr std::size_t jsontx_max_diff = 1024; // delta bytes
|
||||
static constexpr std::size_t jsontx_max_ops = 256; // delta instructions
|
||||
static constexpr std::size_t jsontx_min_copy = 4; // encoder match threshold
|
||||
static constexpr std::size_t jsontx_max_cand = 64; // encoder candidate cap
|
||||
|
||||
// Case-insensitive field-name -> canonical SField. Built once from
|
||||
// SField::knownCodeToField, the same table doServerDefinitions publishes, using
|
||||
// its serializability filter (useful, binary, non-pseudo). sfInvalid if
|
||||
// unknown.
|
||||
static SField const&
|
||||
jsontx_field(std::string const& name)
|
||||
{
|
||||
static auto const tbl = [] {
|
||||
std::unordered_map<std::string, SField const*> m;
|
||||
for (auto const& [code, f] : SField::knownCodeToField)
|
||||
if (f->isUseful() && f->isBinary() && f->fieldType < 10000 &&
|
||||
!f->fieldName.empty())
|
||||
m.emplace(boost::algorithm::to_lower_copy(f->fieldName), f);
|
||||
return m;
|
||||
}();
|
||||
|
||||
auto const i = tbl.find(boost::algorithm::to_lower_copy(name));
|
||||
return i == tbl.end() ? sfInvalid : *i->second;
|
||||
}
|
||||
|
||||
// Civil calendar arithmetic (Howard Hinnant's algorithm), used in both
|
||||
// directions. Pure integer maths - no strptime, no timegm, no locale, no
|
||||
// tzdata - because these conversions decide whether a signature verifies and
|
||||
// so must give the same answer on every node forever.
|
||||
static constexpr std::int64_t
|
||||
jsontx_days(int y, unsigned m, unsigned d) // days from 1970-01-01
|
||||
{
|
||||
y -= m <= 2;
|
||||
std::int64_t const era = (y >= 0 ? y : y - 399) / 400;
|
||||
unsigned const yoe = static_cast<unsigned>(y - era * 400);
|
||||
unsigned const doy = (153 * (m + (m > 2 ? -3 : 9)) + 2) / 5 + d - 1;
|
||||
unsigned const doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
|
||||
return era * 146097 + doe - 719468;
|
||||
}
|
||||
|
||||
static constexpr std::int64_t jsontx_epoch_day = jsontx_days(2000, 1, 1);
|
||||
|
||||
// 9999-12-31T23:59:59.999Z: past this toISOString() switches to expanded years
|
||||
// (+275760-09-13T...) and the fixed 24 character shape no longer holds
|
||||
static constexpr std::uint64_t jsontx_max_time =
|
||||
(static_cast<std::uint64_t>(jsontx_days(9999, 12, 31) - jsontx_epoch_day) *
|
||||
86400 +
|
||||
86399) *
|
||||
1000 +
|
||||
999;
|
||||
|
||||
static_assert(jsontx_epoch_day == 10957); // matches chrono.h epoch_offset
|
||||
|
||||
// Strict Date().toISOString() -> milliseconds since the ripple epoch. Exactly
|
||||
// YYYY-MM-DDTHH:MM:SS.sssZ, always UTC, always three fractional digits.
|
||||
static std::uint64_t
|
||||
jsontx_iso(std::string const& s)
|
||||
{
|
||||
static constexpr char pat[] = "0000-00-00T00:00:00.000Z";
|
||||
if (s.size() != 24)
|
||||
throw std::runtime_error("jsontx: Time must be an ISO 8601 instant");
|
||||
for (std::size_t i = 0; i < 24; ++i)
|
||||
if (pat[i] == '0' ? !std::isdigit(static_cast<unsigned char>(s[i]))
|
||||
: s[i] != pat[i])
|
||||
throw std::runtime_error("jsontx: malformed Time");
|
||||
|
||||
auto const n = [&s](std::size_t i, std::size_t c) {
|
||||
int v = 0;
|
||||
while (c--)
|
||||
v = v * 10 + (s[i++] - '0');
|
||||
return v;
|
||||
};
|
||||
int const y = n(0, 4), mo = n(5, 2), d = n(8, 2), h = n(11, 2),
|
||||
mi = n(14, 2), se = n(17, 2), ms = n(20, 3);
|
||||
if (mo < 1 || mo > 12)
|
||||
throw std::runtime_error("jsontx: Time month out of range");
|
||||
bool const leap = (y % 4 == 0 && y % 100 != 0) || y % 400 == 0;
|
||||
int const dim =
|
||||
mo == 2 ? (leap ? 29 : 28) : ((mo % 2 == 1) == (mo <= 7) ? 31 : 30);
|
||||
// 60 is rejected: JS cannot emit a leap second and the ledger cannot
|
||||
// represent one
|
||||
if (d < 1 || d > dim || h > 23 || mi > 59 || se > 59)
|
||||
throw std::runtime_error("jsontx: Time out of range");
|
||||
|
||||
std::int64_t const t = (jsontx_days(y, mo, d) - jsontx_epoch_day) * 86400 +
|
||||
h * 3600 + mi * 60 + se;
|
||||
if (t < 0)
|
||||
throw std::runtime_error("jsontx: Time precedes the ripple epoch");
|
||||
return static_cast<std::uint64_t>(t) * 1000 + ms;
|
||||
}
|
||||
|
||||
// The exact inverse. Total over [0, jsontx_max_time] and injective, so sfTime
|
||||
// and its ISO spelling are two views of one value and the delta carries
|
||||
// nothing for the field.
|
||||
static std::string
|
||||
jsontx_iso_str(std::uint64_t ms)
|
||||
{
|
||||
if (ms > jsontx_max_time)
|
||||
throw std::runtime_error("jsontx: Time out of range");
|
||||
std::int64_t const z =
|
||||
static_cast<std::int64_t>(ms / 86400000) + jsontx_epoch_day + 719468;
|
||||
unsigned const tod = static_cast<unsigned>(ms / 1000 % 86400);
|
||||
std::int64_t const era = (z >= 0 ? z : z - 146096) / 146097;
|
||||
unsigned const doe = static_cast<unsigned>(z - era * 146097);
|
||||
unsigned const yoe = (doe - doe / 1460 + doe / 36524 - doe / 146096) / 365;
|
||||
unsigned const doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
|
||||
unsigned const mp = (5 * doy + 2) / 153;
|
||||
unsigned const d = doy - (153 * mp + 2) / 5 + 1;
|
||||
unsigned const m = mp + (mp < 10 ? 3 : -9);
|
||||
return fmt::format(
|
||||
"{:04}-{:02}-{:02}T{:02}:{:02}:{:02}.{:03}Z",
|
||||
static_cast<std::int64_t>(yoe) + era * 400 + (m <= 2),
|
||||
m,
|
||||
d,
|
||||
tod / 3600,
|
||||
tod / 60 % 60,
|
||||
tod % 60,
|
||||
static_cast<unsigned>(ms % 1000));
|
||||
}
|
||||
|
||||
//------------------------------------------------------------------------------
|
||||
// jsoncpp compatibility
|
||||
//
|
||||
// This is xrpl's vendored jsoncpp, which predates JSON_HAS_INT64: Json::Int is
|
||||
// int and Json::UInt is unsigned int, both 32 bit, and there is no isInt64 /
|
||||
// asInt64 / asUInt64. Reader::decodeNumber yields intValue or uintValue only
|
||||
// while the digits fit in 32 bits; on overflow, and for any token carrying a
|
||||
// '.' or an exponent, it falls through to decodeDouble and the number arrives
|
||||
// as a realValue.
|
||||
//
|
||||
// So a large integer is not lost, but it is no longer held as an integer, and
|
||||
// the digits the signer wrote are recoverable only while the double is an
|
||||
// exact integer view of them. That holds to 2^53; above it consecutive doubles
|
||||
// are more than 1 apart and distinct decimal integers collapse onto the same
|
||||
// double. Past that point the value is refused rather than guessed at.
|
||||
//------------------------------------------------------------------------------
|
||||
|
||||
static constexpr double jsontx_exact_max = 9007199254740992.0; // 2^53
|
||||
|
||||
// Exact integer view of a json number. False if v is not a number at all, or
|
||||
// is one this build cannot reproduce digit for digit.
|
||||
static bool
|
||||
jsontx_exact(Json::Value const& v, std::int64_t& out)
|
||||
{
|
||||
switch (v.type())
|
||||
{
|
||||
case Json::intValue:
|
||||
out = v.asInt();
|
||||
return true;
|
||||
case Json::uintValue:
|
||||
out = static_cast<std::int64_t>(v.asUInt());
|
||||
return true;
|
||||
case Json::realValue:
|
||||
break;
|
||||
default: // including booleanValue, which isIntegral() would admit
|
||||
return false;
|
||||
}
|
||||
double const d = v.asDouble();
|
||||
if (!std::isfinite(d) || d != std::trunc(d) || d < -jsontx_exact_max ||
|
||||
d > jsontx_exact_max)
|
||||
return false;
|
||||
out = static_cast<std::int64_t>(d);
|
||||
return true;
|
||||
}
|
||||
|
||||
// STI_UINT64 renders as a fixed 16 digit hex string and so needs the value
|
||||
// unsigned. A negative is refused rather than wrapped: the old
|
||||
// static_cast<std::uint64_t>(v.asDouble()) was undefined for a negative or
|
||||
// oversized double, and a silent wrap would let -1 and 18446744073709551615
|
||||
// canonicalize to the same bytes.
|
||||
static std::uint64_t
|
||||
jsontx_u64(Json::Value const& v)
|
||||
{
|
||||
std::int64_t n = 0;
|
||||
if (!jsontx_exact(v, n) || n < 0)
|
||||
throw std::runtime_error(
|
||||
"jsontx: UInt64 must be an exact non-negative integer");
|
||||
return static_cast<std::uint64_t>(n);
|
||||
}
|
||||
|
||||
// Renders a javascript number as an exact integer. Anything this build cannot
|
||||
// reproduce digit for digit - a fraction, an infinity, a magnitude past 2^53 -
|
||||
// is rejected outright: shortest-round-trip rendering of a double is not
|
||||
// portable enough to sit in a consensus preimage, and nothing in a transaction
|
||||
// needs one. Fractional amounts arrive as strings, which is what the ledger
|
||||
// wants anyway.
|
||||
static std::string
|
||||
jsontx_num(Json::Value const& v)
|
||||
{
|
||||
std::int64_t n = 0;
|
||||
if (!jsontx_exact(v, n))
|
||||
throw std::runtime_error("jsontx: number must be an exact integer");
|
||||
return std::to_string(n);
|
||||
}
|
||||
|
||||
// Returns { sanitized, diff }. `sanitized` is the canonical form: whitespace
|
||||
// stripped, field names capitalized to their xahau spelling, members reordered
|
||||
// by field code, numbers reformatted per field type. `diff` is a binary delta
|
||||
// which, applied to `sanitized` by unsanitize_jsontx, reproduces `raw` byte for
|
||||
// byte. Throws on anything it cannot canonicalize.
|
||||
static std::pair<std::string, std::string>
|
||||
sanitize_jsontx(std::string_view raw)
|
||||
{
|
||||
if (raw.size() > jsontx_max_text)
|
||||
throw std::runtime_error("jsontx: document too large");
|
||||
|
||||
Json::Value jv;
|
||||
if (Json::Reader r; !r.parse(raw.data(), raw.data() + raw.size(), jv) ||
|
||||
!jv.isObject())
|
||||
throw std::runtime_error("jsontx: malformed json");
|
||||
|
||||
// (a plain recursive lambda; deducing-this would drop the std::function)
|
||||
std::function<
|
||||
void(Json::Value const&, SerializedTypeID, SField const*, std::string&)>
|
||||
emit = [&](Json::Value const& v,
|
||||
SerializedTypeID ty,
|
||||
SField const* fld,
|
||||
std::string& o) {
|
||||
if (v.isObject())
|
||||
{
|
||||
auto keys = v.getMemberNames();
|
||||
o += '{';
|
||||
if (ty == STI_OBJECT) // keys are xahau fields
|
||||
{
|
||||
std::vector<std::pair<SField const*, std::string>> ks;
|
||||
for (auto const& k : keys)
|
||||
{
|
||||
auto const& f = jsontx_field(k);
|
||||
if (f == sfInvalid)
|
||||
throw std::runtime_error(
|
||||
"jsontx: unknown field '" + k + "'");
|
||||
ks.emplace_back(&f, k);
|
||||
}
|
||||
std::sort(
|
||||
ks.begin(), ks.end(), [](auto const& a, auto const& b) {
|
||||
return a.first->fieldCode < b.first->fieldCode;
|
||||
});
|
||||
for (std::size_t n = 0; n < ks.size(); ++n)
|
||||
{
|
||||
auto const& [f, k] = ks[n];
|
||||
if (n && f == ks[n - 1].first) // e.g. "Fee" and "fee"
|
||||
throw std::runtime_error(
|
||||
"jsontx: duplicate field '" + f->fieldName +
|
||||
"'");
|
||||
if (o.back() != '{')
|
||||
o += ',';
|
||||
o += Json::valueToQuotedString(f->fieldName.c_str()) +
|
||||
':';
|
||||
emit(v[k], f->fieldType, f, o);
|
||||
}
|
||||
}
|
||||
else // amount / issue style subobject: lexicographic, quoted
|
||||
{
|
||||
std::sort(keys.begin(), keys.end());
|
||||
for (auto const& k : keys)
|
||||
{
|
||||
if (o.back() != '{')
|
||||
o += ',';
|
||||
o += Json::valueToQuotedString(k.c_str()) + ':';
|
||||
emit(v[k], STI_NOTPRESENT, nullptr, o);
|
||||
}
|
||||
}
|
||||
o += '}';
|
||||
}
|
||||
else if (v.isArray())
|
||||
{
|
||||
o += '[';
|
||||
for (auto const& e : v)
|
||||
{
|
||||
if (o.back() != '[')
|
||||
o += ',';
|
||||
emit(e, STI_OBJECT, nullptr, o);
|
||||
}
|
||||
o += ']';
|
||||
}
|
||||
else if (v.isString())
|
||||
{
|
||||
// Time is spelled Date().toISOString() in the preimage and
|
||||
// stored as an sfTime u64 of milliseconds. Re-emitting the
|
||||
// round-tripped spelling rather than the input is what makes
|
||||
// the canonical form a fixed point: any string that is not
|
||||
// exactly what jsontx_iso_str produces is rejected here.
|
||||
if (fld && *fld == sfTime)
|
||||
o += Json::valueToQuotedString(
|
||||
jsontx_iso_str(jsontx_iso(v.asString())).c_str());
|
||||
else
|
||||
o += Json::valueToQuotedString(v.asCString());
|
||||
}
|
||||
else if (v.isBool())
|
||||
o += v.asBool() ? "true" : "false";
|
||||
else if (v.isNull())
|
||||
throw std::runtime_error("jsontx: null value");
|
||||
else if (ty == STI_UINT8 || ty == STI_UINT16 || ty == STI_UINT32)
|
||||
o += jsontx_num(v); // small ints stay bare
|
||||
else if (ty == STI_UINT64)
|
||||
o += '"' + fmt::format("{:016X}", jsontx_u64(v)) + '"';
|
||||
else // amounts, u64, everything else the ledger wants as a string
|
||||
o += Json::valueToQuotedString(jsontx_num(v).c_str());
|
||||
};
|
||||
|
||||
std::string out;
|
||||
emit(jv, STI_OBJECT, nullptr, out);
|
||||
|
||||
// Greedy copy/insert delta over `raw`, sourcing from `out`.
|
||||
// op 0x00 <varint len> <bytes> literal
|
||||
// op 0x01 <varint off> <varint len> copy from sanitized
|
||||
std::string diff, lit;
|
||||
auto const gram = [](std::string_view s, std::size_t i) {
|
||||
return std::uint32_t(std::uint8_t(s[i])) << 24 |
|
||||
std::uint32_t(std::uint8_t(s[i + 1])) << 16 |
|
||||
std::uint32_t(std::uint8_t(s[i + 2])) << 8 |
|
||||
std::uint32_t(std::uint8_t(s[i + 3]));
|
||||
};
|
||||
auto const varint = [](std::string& o, std::uint64_t v) {
|
||||
do
|
||||
{
|
||||
std::uint8_t const c = v & 0x7F;
|
||||
v >>= 7;
|
||||
o += static_cast<char>(c | (v ? 0x80 : 0));
|
||||
} while (v);
|
||||
};
|
||||
auto const flush = [&] {
|
||||
if (lit.empty())
|
||||
return;
|
||||
diff += char(0);
|
||||
varint(diff, lit.size());
|
||||
diff += lit;
|
||||
lit.clear();
|
||||
};
|
||||
|
||||
// This encoder is normative - unsanitize_jsontx only accepts its exact
|
||||
// output - so it must emit identical bytes on every node and every stdlib.
|
||||
// An unordered container would not: bucket order for equal keys is
|
||||
// unspecified, and with a candidate cap that changes which match wins.
|
||||
std::map<std::uint32_t, std::vector<std::uint32_t>> idx;
|
||||
for (std::size_t i = 0; i + jsontx_min_copy <= out.size(); ++i)
|
||||
idx[gram(out, i)].push_back(i);
|
||||
|
||||
for (std::size_t i = 0; i < raw.size();)
|
||||
{
|
||||
std::size_t bo = 0, bl = 0;
|
||||
if (i + jsontx_min_copy <= raw.size())
|
||||
if (auto const it = idx.find(gram(raw, i)); it != idx.end())
|
||||
{
|
||||
std::size_t tried = 0;
|
||||
for (auto const off : it->second) // ascending, so ties keep
|
||||
{ // the lowest offset
|
||||
if (++tried > jsontx_max_cand)
|
||||
break;
|
||||
std::size_t l = 0;
|
||||
while (i + l < raw.size() && off + l < out.size() &&
|
||||
out[off + l] == raw[i + l])
|
||||
++l;
|
||||
if (l > bl)
|
||||
bl = l, bo = off;
|
||||
}
|
||||
}
|
||||
if (bl >= jsontx_min_copy)
|
||||
{
|
||||
flush();
|
||||
diff += char(1);
|
||||
varint(diff, bo);
|
||||
varint(diff, bl);
|
||||
i += bl;
|
||||
}
|
||||
else
|
||||
lit += raw[i++];
|
||||
}
|
||||
flush();
|
||||
|
||||
return {std::move(out), std::move(diff)};
|
||||
}
|
||||
|
||||
// Applies an UNTRUSTED delta to a canonical form the node derived itself.
|
||||
// Copies read only from `sanitized`, never from the output being built, so a
|
||||
// short delta cannot expand geometrically. Offsets and lengths are range
|
||||
// checked before use, varints are length- and minimality-bounded, and the two
|
||||
// encodings the encoder can never emit - an unmerged literal run, and a copy
|
||||
// abutting the previous copy in the source - are rejected. Throws on anything
|
||||
// else.
|
||||
static std::string
|
||||
unsanitize_jsontx(std::string_view sanitized, std::string_view diff)
|
||||
{
|
||||
if (sanitized.size() > jsontx_max_text || diff.size() > jsontx_max_diff)
|
||||
throw std::runtime_error("jsontx: oversize delta input");
|
||||
|
||||
std::string out;
|
||||
std::size_t p = 0, ops = 0, prevEnd = 0;
|
||||
int prev = -1;
|
||||
|
||||
auto const varint = [&](std::uint64_t max) -> std::uint64_t {
|
||||
std::uint64_t v = 0;
|
||||
for (int s = 0; s <= 21; s += 7) // four bytes; caps far under a shift
|
||||
{ // wide enough to be undefined
|
||||
if (p >= diff.size())
|
||||
throw std::runtime_error("jsontx: truncated delta");
|
||||
std::uint8_t const c = diff[p++];
|
||||
v |= std::uint64_t(c & 0x7F) << s;
|
||||
if (c & 0x80)
|
||||
continue;
|
||||
if (s && !(c & 0x7F))
|
||||
throw std::runtime_error("jsontx: non-minimal varint");
|
||||
if (v > max)
|
||||
throw std::runtime_error("jsontx: delta value out of range");
|
||||
return v;
|
||||
}
|
||||
throw std::runtime_error("jsontx: overlong varint");
|
||||
};
|
||||
|
||||
while (p < diff.size())
|
||||
{
|
||||
if (++ops > jsontx_max_ops)
|
||||
throw std::runtime_error("jsontx: too many delta ops");
|
||||
|
||||
std::uint8_t const op = diff[p++];
|
||||
if (op > 1)
|
||||
throw std::runtime_error("jsontx: unknown delta op");
|
||||
|
||||
std::size_t n = 0;
|
||||
if (op == 0) // literal
|
||||
{
|
||||
if (prev == 0)
|
||||
throw std::runtime_error("jsontx: unmerged literal run");
|
||||
n = varint(jsontx_max_text);
|
||||
if (n == 0 || n > diff.size() - p)
|
||||
throw std::runtime_error("jsontx: bad literal length");
|
||||
if (out.size() + n > jsontx_max_text)
|
||||
throw std::runtime_error("jsontx: delta expands too far");
|
||||
out += diff.substr(p, n);
|
||||
p += n;
|
||||
}
|
||||
else // copy from the canonical form
|
||||
{
|
||||
auto const off = varint(sanitized.size());
|
||||
n = varint(sanitized.size() - off);
|
||||
if (n < jsontx_min_copy)
|
||||
throw std::runtime_error("jsontx: undersize copy");
|
||||
if (prev == 1 && off == prevEnd)
|
||||
throw std::runtime_error("jsontx: unmerged copy run");
|
||||
if (out.size() + n > jsontx_max_text)
|
||||
throw std::runtime_error("jsontx: delta expands too far");
|
||||
out += sanitized.substr(off, n);
|
||||
prevEnd = off + n;
|
||||
}
|
||||
prev = op;
|
||||
}
|
||||
|
||||
if (out.empty())
|
||||
throw std::runtime_error("jsontx: empty delta");
|
||||
return out;
|
||||
}
|
||||
|
||||
// The complete untrusted-side check, in one place so the RPC path and the
|
||||
// relay/consensus path cannot drift. Takes the transaction exactly as it came
|
||||
// off the wire and returns the reconstructed preimage.
|
||||
static std::string
|
||||
jsontx_verify(STTx const& stx, std::string_view diff)
|
||||
{
|
||||
if (!stx.isFieldPresent(sfTxnSignature) || stx.isFieldPresent(sfSigners))
|
||||
throw std::runtime_error("jsontx: expects a lone TxnSignature");
|
||||
|
||||
// Out comes everything the signer did not have in front of them: the
|
||||
// signature, and (once the field exists) the delta carrier. SigningPubKey
|
||||
// stays. It being inside the preimage is what binds key to signature and
|
||||
// stops a third party re-signing a captured preimage under their own key.
|
||||
auto txj = stx.STObject::getJson(JsonOptions::none);
|
||||
txj.removeMember(sfTxnSignature.fieldName);
|
||||
|
||||
// sfTime is a u64 of milliseconds on the wire and an ISO 8601 instant in
|
||||
// the preimage. The two are a bijection over the representable range, so
|
||||
// this is a rewrite rather than a reconstruction and the delta carries
|
||||
// nothing for the field.
|
||||
if (stx.isFieldPresent(sfTime))
|
||||
txj[sfTime.fieldName] = jsontx_iso_str(stx.getFieldU64(sfTime));
|
||||
|
||||
auto const pkb = stx.getSigningPubKey();
|
||||
if (publicKeyType(makeSlice(pkb)) != KeyType::ed25519)
|
||||
throw std::runtime_error("jsontx: SigningPubKey must be ed25519");
|
||||
|
||||
// canonical form, derived only from data the node has already validated
|
||||
auto const san = sanitize_jsontx(Json::FastWriter{}.write(txj)).first;
|
||||
|
||||
// reconstruct the signed preimage under the caps above
|
||||
auto const raw = unsanitize_jsontx(san, diff);
|
||||
|
||||
// Bind the preimage to the transaction. This is the load-bearing check,
|
||||
// not a sanity check: a delta of pure literals can reconstruct ANY text,
|
||||
// so without it any ed25519 signature the key ever produced over anything
|
||||
// at all would authorise this transaction. Comparing the delta too - not
|
||||
// just the canonical form - makes the delta a pure function of the
|
||||
// preimage, which rules out a second delta reconstructing the same bytes
|
||||
// and yielding a second valid transaction id.
|
||||
auto const [san2, diff2] = sanitize_jsontx(raw);
|
||||
if (san2 != san || diff2 != diff)
|
||||
throw std::runtime_error("jsontx: preimage does not match transaction");
|
||||
|
||||
if (!verify(
|
||||
PublicKey(makeSlice(pkb)),
|
||||
makeSlice(raw),
|
||||
makeSlice(stx.getFieldVL(sfTxnSignature))))
|
||||
throw std::runtime_error("jsontx: signature does not verify");
|
||||
|
||||
return raw;
|
||||
}
|
||||
|
||||
} // namespace ripple
|
||||
#endif
|
||||
@@ -34,7 +34,6 @@
|
||||
// If you add an amendment here, then do not forget to increment `numFeatures`
|
||||
// in include/xrpl/protocol/Feature.h.
|
||||
|
||||
XRPL_FEATURE(JsonTx, Supported::yes, VoteBehavior::DefaultNo)
|
||||
XRPL_FIX (HookMap, Supported::yes, VoteBehavior::DefaultYes)
|
||||
XRPL_FIX (GuardDepth32, Supported::yes, VoteBehavior::DefaultNo)
|
||||
XRPL_FEATURE(NamedHooks, Supported::yes, VoteBehavior::DefaultNo)
|
||||
|
||||
@@ -153,7 +153,6 @@ TYPED_SFIELD(sfOutstandingAmount, UINT64, 25, SField::sMD_BaseTen|SFie
|
||||
TYPED_SFIELD(sfMPTAmount, UINT64, 26, SField::sMD_BaseTen|SField::sMD_Default)
|
||||
TYPED_SFIELD(sfIssuerNode, UINT64, 27)
|
||||
TYPED_SFIELD(sfSubjectNode, UINT64, 28)
|
||||
TYPED_SFIELD(sfTime, UINT64, 96)
|
||||
TYPED_SFIELD(sfTouchCount, UINT64, 97)
|
||||
TYPED_SFIELD(sfAccountIndex, UINT64, 98)
|
||||
TYPED_SFIELD(sfAccountCount, UINT64, 99)
|
||||
@@ -294,7 +293,6 @@ TYPED_SFIELD(sfAssetClass, VL, 29)
|
||||
TYPED_SFIELD(sfProvider, VL, 30)
|
||||
TYPED_SFIELD(sfMPTokenMetadata, VL, 31)
|
||||
TYPED_SFIELD(sfCredentialType, VL, 32)
|
||||
TYPED_SFIELD(sfJsonTxDelta, VL, 96)
|
||||
TYPED_SFIELD(sfHookName, VL, 97)
|
||||
TYPED_SFIELD(sfRemarkValue, VL, 98)
|
||||
TYPED_SFIELD(sfRemarkName, VL, 99)
|
||||
|
||||
@@ -630,7 +630,6 @@ JSS(server_status); // out: NetworkOPs
|
||||
JSS(server_version); // out: NetworkOPs
|
||||
JSS(settle_delay); // out: AccountChannels
|
||||
JSS(severity); // in: LogLevel
|
||||
JSS(sig);
|
||||
JSS(signature); // out: NetworkOPs, ChannelAuthorize
|
||||
JSS(signature_verified); // out: ChannelVerify
|
||||
JSS(signing_key); // out: NetworkOPs
|
||||
|
||||
@@ -49,8 +49,6 @@ TxFormats::TxFormats()
|
||||
{sfNetworkID, soeOPTIONAL},
|
||||
{sfHookParameters, soeOPTIONAL},
|
||||
{sfHookName, soeOPTIONAL},
|
||||
{sfTime, soeOPTIONAL},
|
||||
{sfJsonTxDelta, soeOPTIONAL},
|
||||
};
|
||||
|
||||
#pragma push_macro("UNWRAP")
|
||||
|
||||
1187
src/test/net/HTTPClient_test.cpp
Normal file
1187
src/test/net/HTTPClient_test.cpp
Normal file
File diff suppressed because it is too large
Load Diff
351
src/test/net/RPCSub_test.cpp
Normal file
351
src/test/net/RPCSub_test.cpp
Normal file
@@ -0,0 +1,351 @@
|
||||
//------------------------------------------------------------------------------
|
||||
/*
|
||||
This file is part of rippled: https://github.com/ripple/rippled
|
||||
Copyright (c) 2024 Ripple Labs Inc.
|
||||
|
||||
Permission to use, copy, modify, and/or distribute this software for any
|
||||
purpose with or without fee is hereby granted, provided that the above
|
||||
copyright notice and this permission notice appear in all copies.
|
||||
|
||||
THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES
|
||||
WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF
|
||||
MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR
|
||||
ANY SPECIAL , DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES
|
||||
WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN
|
||||
ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF
|
||||
OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE.
|
||||
*/
|
||||
//==============================================================================
|
||||
|
||||
#include <test/jtx.h>
|
||||
#include <xrpld/core/Job.h>
|
||||
#include <xrpld/core/JobQueue.h>
|
||||
#include <xrpld/net/RPCSub.h>
|
||||
#include <xrpl/json/json_value.h>
|
||||
|
||||
#include <boost/asio.hpp>
|
||||
#include <boost/asio/ip/tcp.hpp>
|
||||
|
||||
#include <atomic>
|
||||
#include <chrono>
|
||||
#include <memory>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
|
||||
namespace ripple {
|
||||
namespace test {
|
||||
|
||||
// Minimal HTTP endpoint that counts received webhook POSTs and replies
|
||||
// with a configurable status. Responses are EOF-delimited (no
|
||||
// Content-Length) and the socket is closed right after writing — the
|
||||
// exact shape that triggered the original handleData EOF-completion
|
||||
// leak. So these tests exercise RPCSub flow control AND the HTTPClient
|
||||
// EOF fix end to end: if either regressed, delivery would stall and the
|
||||
// expected count would never be reached within the timeout.
|
||||
class MockWebhookEndpoint
|
||||
{
|
||||
boost::asio::io_service ios_;
|
||||
std::unique_ptr<boost::asio::io_service::work> work_;
|
||||
boost::asio::ip::tcp::acceptor acceptor_;
|
||||
std::thread thread_;
|
||||
unsigned short port_;
|
||||
|
||||
std::atomic<int> received_{0};
|
||||
std::atomic<int> status_{200};
|
||||
std::atomic<int> delayMs_{0};
|
||||
|
||||
public:
|
||||
MockWebhookEndpoint()
|
||||
: work_(std::make_unique<boost::asio::io_service::work>(ios_))
|
||||
, acceptor_(
|
||||
ios_,
|
||||
boost::asio::ip::tcp::endpoint(
|
||||
boost::asio::ip::address::from_string("127.0.0.1"),
|
||||
0))
|
||||
{
|
||||
port_ = acceptor_.local_endpoint().port();
|
||||
accept();
|
||||
thread_ = std::thread([this] { ios_.run(); });
|
||||
}
|
||||
|
||||
~MockWebhookEndpoint()
|
||||
{
|
||||
work_.reset();
|
||||
boost::system::error_code ec;
|
||||
acceptor_.close(ec);
|
||||
ios_.stop();
|
||||
if (thread_.joinable())
|
||||
thread_.join();
|
||||
}
|
||||
|
||||
unsigned short
|
||||
port() const
|
||||
{
|
||||
return port_;
|
||||
}
|
||||
|
||||
int
|
||||
received() const
|
||||
{
|
||||
return received_;
|
||||
}
|
||||
|
||||
void
|
||||
setStatus(int s)
|
||||
{
|
||||
status_ = s;
|
||||
}
|
||||
|
||||
// Delay each reply so delivery is deterministically slower than the
|
||||
// microsecond-fast enqueue loop — keeps the deque full for the
|
||||
// queue-cap drop test regardless of scheduling.
|
||||
void
|
||||
setResponseDelay(int ms)
|
||||
{
|
||||
delayMs_ = ms;
|
||||
}
|
||||
|
||||
private:
|
||||
void
|
||||
accept()
|
||||
{
|
||||
auto sock = std::make_shared<boost::asio::ip::tcp::socket>(ios_);
|
||||
acceptor_.async_accept(*sock, [this, sock](auto ec) {
|
||||
if (ec)
|
||||
return;
|
||||
handle(sock);
|
||||
accept();
|
||||
});
|
||||
}
|
||||
|
||||
void
|
||||
handle(std::shared_ptr<boost::asio::ip::tcp::socket> sock)
|
||||
{
|
||||
auto buf = std::make_shared<boost::asio::streambuf>();
|
||||
boost::asio::async_read_until(
|
||||
*sock, *buf, "\r\n\r\n", [this, sock, buf](auto ec, std::size_t) {
|
||||
if (ec)
|
||||
return;
|
||||
|
||||
++received_;
|
||||
|
||||
auto const delay = delayMs_.load();
|
||||
if (delay > 0)
|
||||
{
|
||||
auto timer =
|
||||
std::make_shared<boost::asio::steady_timer>(ios_);
|
||||
timer->expires_from_now(std::chrono::milliseconds(delay));
|
||||
timer->async_wait(
|
||||
[this, sock, timer](auto) { reply(sock); });
|
||||
}
|
||||
else
|
||||
{
|
||||
reply(sock);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
void
|
||||
reply(std::shared_ptr<boost::asio::ip::tcp::socket> sock)
|
||||
{
|
||||
// EOF-delimited reply: no Content-Length, close after writing.
|
||||
// This is the realistic failing-webhook shape.
|
||||
auto resp = std::make_shared<std::string>(
|
||||
"HTTP/1.0 " + std::to_string(status_.load()) +
|
||||
" Reply\r\n\r\n{\"result\":{}}");
|
||||
boost::asio::async_write(
|
||||
*sock, boost::asio::buffer(*resp), [sock, resp](auto, std::size_t) {
|
||||
boost::system::error_code ig;
|
||||
sock->shutdown(boost::asio::ip::tcp::socket::shutdown_both, ig);
|
||||
sock->close(ig);
|
||||
});
|
||||
}
|
||||
};
|
||||
|
||||
//------------------------------------------------------------------------------
|
||||
|
||||
class RPCSub_test : public beast::unit_test::suite
|
||||
{
|
||||
// Generous ceiling: the instrumented Debug (coverage) build is much
|
||||
// slower than Release, so timeouts are sized for that, not Release.
|
||||
template <class Cond>
|
||||
bool
|
||||
waitFor(Cond cond, std::chrono::seconds timeout = std::chrono::seconds{30})
|
||||
{
|
||||
auto const deadline = std::chrono::steady_clock::now() + timeout;
|
||||
while (!cond() && std::chrono::steady_clock::now() < deadline)
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(10));
|
||||
return cond();
|
||||
}
|
||||
|
||||
std::shared_ptr<RPCSub>
|
||||
makeSub(
|
||||
jtx::Env& env,
|
||||
MockWebhookEndpoint& ep,
|
||||
std::size_t maxQueueSize = 16384)
|
||||
{
|
||||
return make_RPCSub(
|
||||
env.app().getOPs(),
|
||||
env.app().getJobQueue(),
|
||||
"http://127.0.0.1:" + std::to_string(ep.port()) + "/",
|
||||
"",
|
||||
"",
|
||||
env.app().logs(),
|
||||
maxQueueSize);
|
||||
}
|
||||
|
||||
// True once no RPCSub sending job is queued or running. sendThread
|
||||
// captures a raw `this`, so the RPCSub must not be destroyed while a
|
||||
// job is still in flight — wait on this before letting the sub die.
|
||||
bool
|
||||
sendingIdle(jtx::Env& env)
|
||||
{
|
||||
return env.app().getJobQueue().getJobCountTotal(jtCLIENT_SUBSCRIBE) ==
|
||||
0;
|
||||
}
|
||||
|
||||
// Wait for all events to reach the endpoint AND the sending job to
|
||||
// finish, so the sub can be torn down without racing sendThread.
|
||||
void
|
||||
drainAndSettle(jtx::Env& env, MockWebhookEndpoint& ep, int expected)
|
||||
{
|
||||
bool const delivered =
|
||||
waitFor([&] { return ep.received() >= expected; });
|
||||
bool const idle = waitFor([&] { return sendingIdle(env); });
|
||||
log << " drainAndSettle: received=" << ep.received() << "/" << expected
|
||||
<< " idle=" << idle << std::endl;
|
||||
BEAST_EXPECT(delivered);
|
||||
BEAST_EXPECT(idle);
|
||||
}
|
||||
|
||||
void
|
||||
send(std::shared_ptr<RPCSub> const& sub, int n)
|
||||
{
|
||||
Json::Value ev(Json::objectValue);
|
||||
ev["n"] = n;
|
||||
sub->send(ev, false);
|
||||
}
|
||||
|
||||
void
|
||||
testDelivery()
|
||||
{
|
||||
testcase("Webhook events are delivered");
|
||||
|
||||
using namespace jtx;
|
||||
Env env{*this};
|
||||
MockWebhookEndpoint ep;
|
||||
|
||||
static constexpr int N = 10;
|
||||
{
|
||||
auto sub = makeSub(env, ep);
|
||||
for (int i = 0; i < N; ++i)
|
||||
send(sub, i);
|
||||
drainAndSettle(env, ep, N);
|
||||
}
|
||||
|
||||
BEAST_EXPECT(ep.received() == N);
|
||||
}
|
||||
|
||||
void
|
||||
testErrorsDoNotStall()
|
||||
{
|
||||
testcase("Delivery continues when endpoint returns HTTP 500");
|
||||
|
||||
// The original bug (xrpld #6341): an endpoint returning errors
|
||||
// without Content-Length never completed, stalling delivery to
|
||||
// ALL subscribers. Here every response is a 500 with no
|
||||
// Content-Length (EOF-delimited) — all N must still arrive.
|
||||
using namespace jtx;
|
||||
Env env{*this};
|
||||
MockWebhookEndpoint ep;
|
||||
ep.setStatus(500);
|
||||
|
||||
static constexpr int N = 10;
|
||||
{
|
||||
auto sub = makeSub(env, ep);
|
||||
for (int i = 0; i < N; ++i)
|
||||
send(sub, i);
|
||||
drainAndSettle(env, ep, N);
|
||||
}
|
||||
|
||||
BEAST_EXPECT(ep.received() == N);
|
||||
}
|
||||
|
||||
void
|
||||
testRestartAfterDrain()
|
||||
{
|
||||
testcase("Sending restarts after the queue drains");
|
||||
|
||||
// After a batch drains, sendThread clears mSending and returns.
|
||||
// A later send() must start a fresh sending job; if mSending were
|
||||
// left set (the #6341 failure mode) the second burst would never
|
||||
// be delivered.
|
||||
using namespace jtx;
|
||||
Env env{*this};
|
||||
MockWebhookEndpoint ep;
|
||||
|
||||
{
|
||||
auto sub = makeSub(env, ep);
|
||||
|
||||
// First burst, then wait for the sending job to fully drain
|
||||
// and exit (mSending cleared) — deterministically, not via a
|
||||
// sleep.
|
||||
for (int i = 0; i < 5; ++i)
|
||||
send(sub, i);
|
||||
drainAndSettle(env, ep, 5);
|
||||
|
||||
// Second burst must start a fresh sending job.
|
||||
for (int i = 5; i < 10; ++i)
|
||||
send(sub, i);
|
||||
drainAndSettle(env, ep, 10);
|
||||
}
|
||||
|
||||
BEAST_EXPECT(ep.received() == 10);
|
||||
}
|
||||
|
||||
void
|
||||
testQueueCapDrops()
|
||||
{
|
||||
testcase("Events past the queue cap are dropped");
|
||||
|
||||
// With a tiny cap, pushing far more events than delivery can keep
|
||||
// up with forces send() down the drop path: enqueue is microsecond
|
||||
// -fast while each (delayed) HTTP delivery is a full round-trip, so
|
||||
// the deque sits at the cap and excess events are dropped. The
|
||||
// delay makes "delivery slower than enqueue" hold regardless of
|
||||
// scheduling, so this isn't timing-dependent. We just need some
|
||||
// delivered (cap works) and some dropped (drop path exercised).
|
||||
using namespace jtx;
|
||||
Env env{*this};
|
||||
MockWebhookEndpoint ep;
|
||||
ep.setResponseDelay(50);
|
||||
|
||||
static constexpr int pushed = 50;
|
||||
{
|
||||
auto sub = makeSub(env, ep, /*maxQueueSize*/ 2);
|
||||
for (int i = 0; i < pushed; ++i)
|
||||
send(sub, i);
|
||||
BEAST_EXPECT(waitFor([&] { return sendingIdle(env); }));
|
||||
}
|
||||
|
||||
log << " queue cap: received " << ep.received() << "/" << pushed
|
||||
<< std::endl;
|
||||
BEAST_EXPECT(ep.received() > 0);
|
||||
BEAST_EXPECT(ep.received() < pushed);
|
||||
}
|
||||
|
||||
public:
|
||||
void
|
||||
run() override
|
||||
{
|
||||
testDelivery();
|
||||
testErrorsDoNotStall();
|
||||
testRestartAfterDrain();
|
||||
testQueueCapDrops();
|
||||
}
|
||||
};
|
||||
|
||||
BEAST_DEFINE_TESTSUITE(RPCSub, net, ripple);
|
||||
|
||||
} // namespace test
|
||||
} // namespace ripple
|
||||
@@ -22,7 +22,6 @@
|
||||
|
||||
#include <xrpld/core/JobQueue.h>
|
||||
#include <xrpld/net/InfoSub.h>
|
||||
#include <boost/asio/io_service.hpp>
|
||||
|
||||
namespace ripple {
|
||||
|
||||
@@ -39,16 +38,17 @@ protected:
|
||||
explicit RPCSub(InfoSub::Source& source);
|
||||
};
|
||||
|
||||
// VFALCO Why is the io_service needed?
|
||||
std::shared_ptr<RPCSub>
|
||||
make_RPCSub(
|
||||
InfoSub::Source& source,
|
||||
boost::asio::io_service& io_service,
|
||||
JobQueue& jobQueue,
|
||||
std::string const& strUrl,
|
||||
std::string const& strUsername,
|
||||
std::string const& strPassword,
|
||||
Logs& logs);
|
||||
Logs& logs,
|
||||
// Max events buffered before new ones are dropped. Configurable so
|
||||
// tests can exercise the drop path without queueing the full default.
|
||||
std::size_t maxQueueSize = 16384);
|
||||
|
||||
} // namespace ripple
|
||||
|
||||
|
||||
@@ -122,12 +122,20 @@ public:
|
||||
mComplete = complete;
|
||||
mTimeout = timeout;
|
||||
|
||||
// Bind a non-owning `this` (not shared_from_this()) into mBuild.
|
||||
// mBuild is a member, so capturing a shared_ptr to self here would
|
||||
// form a reference cycle (this -> mBuild -> shared_ptr<this>) that
|
||||
// never breaks, leaking the object and its socket FD after the
|
||||
// request completes. mBuild is only ever invoked from
|
||||
// handleRequest(), which always runs inside an async handler that
|
||||
// already holds a shared_from_this(), so the object is guaranteed
|
||||
// alive whenever mBuild fires — a raw `this` is safe.
|
||||
request(
|
||||
bSSL,
|
||||
deqSites,
|
||||
std::bind(
|
||||
&HTTPClientImp::makeGet,
|
||||
shared_from_this(),
|
||||
this,
|
||||
strPath,
|
||||
std::placeholders::_1,
|
||||
std::placeholders::_2),
|
||||
@@ -393,8 +401,12 @@ public:
|
||||
if (boost::regex_match(strHeader, smMatch, reBody)) // we got some body
|
||||
mBody = smMatch[1];
|
||||
|
||||
bool const hasContentLength =
|
||||
boost::regex_match(strHeader, smMatch, reSize);
|
||||
mReceivedContentLength = hasContentLength;
|
||||
|
||||
std::size_t const responseSize = [&] {
|
||||
if (boost::regex_match(strHeader, smMatch, reSize))
|
||||
if (hasContentLength)
|
||||
return beast::lexicalCast<std::size_t>(
|
||||
std::string(smMatch[1]), maxResponseSize_);
|
||||
return maxResponseSize_;
|
||||
@@ -445,22 +457,24 @@ public:
|
||||
JLOG(j_.trace()) << "Read error: " << mShutdown.message();
|
||||
|
||||
invokeComplete(mShutdown);
|
||||
return;
|
||||
}
|
||||
else
|
||||
{
|
||||
if (mShutdown)
|
||||
{
|
||||
JLOG(j_.trace()) << "Complete.";
|
||||
}
|
||||
else
|
||||
{
|
||||
mResponse.commit(bytes_transferred);
|
||||
std::string strBody{
|
||||
{std::istreambuf_iterator<char>(&mResponse)},
|
||||
std::istreambuf_iterator<char>()};
|
||||
invokeComplete(ecResult, mStatus, mBody + strBody);
|
||||
}
|
||||
}
|
||||
|
||||
// Either the read completed normally or it ended at EOF. EOF is a
|
||||
// successful completion for EOF-delimited responses, but it is an
|
||||
// error when the server promised a Content-Length and closed early.
|
||||
JLOG(j_.trace()) << "Complete.";
|
||||
|
||||
mResponse.commit(bytes_transferred);
|
||||
std::string strBody{
|
||||
{std::istreambuf_iterator<char>(&mResponse)},
|
||||
std::istreambuf_iterator<char>()};
|
||||
|
||||
auto completeEc = ecResult;
|
||||
if (completeEc == boost::asio::error::eof && !mReceivedContentLength)
|
||||
completeEc.clear();
|
||||
|
||||
invokeComplete(completeEc, mStatus, mBody + strBody);
|
||||
}
|
||||
|
||||
// Call cancel the deadline timer and invoke the completion routine.
|
||||
@@ -516,6 +530,7 @@ private:
|
||||
boost::asio::streambuf mHeader;
|
||||
boost::asio::streambuf mResponse;
|
||||
std::string mBody;
|
||||
bool mReceivedContentLength = false;
|
||||
const unsigned short mPort;
|
||||
std::size_t const maxResponseSize_;
|
||||
int mStatus;
|
||||
|
||||
@@ -1585,6 +1585,10 @@ struct RPCCallImp
|
||||
// callbackFuncP.
|
||||
|
||||
// Receive reply
|
||||
if (ecResult)
|
||||
Throw<std::runtime_error>(
|
||||
"RPC transport error: " + ecResult.message());
|
||||
|
||||
if (strData.empty())
|
||||
Throw<std::runtime_error>(
|
||||
"no response from server. Please "
|
||||
@@ -1748,6 +1752,7 @@ rpcClient(
|
||||
}
|
||||
|
||||
{
|
||||
//@@start blocking-request
|
||||
boost::asio::io_service isService;
|
||||
RPCCall::fromNetwork(
|
||||
isService,
|
||||
@@ -1771,6 +1776,7 @@ rpcClient(
|
||||
headers);
|
||||
isService.run(); // This blocks until there are no more
|
||||
// outstanding async calls.
|
||||
//@@end blocking-request
|
||||
}
|
||||
if (jvOutput.isMember("result"))
|
||||
{
|
||||
@@ -1881,15 +1887,21 @@ fromNetwork(
|
||||
|
||||
// Send request
|
||||
|
||||
// Number of bytes to try to receive if no
|
||||
// Content-Length header received
|
||||
constexpr auto RPC_REPLY_MAX_BYTES = megabytes(256);
|
||||
// Number of bytes to try to receive if no Content-Length header is
|
||||
// received. Webhook event deliveries ("event") ignore the response
|
||||
// body, so a missing Content-Length must not pre-allocate the full
|
||||
// 256MB RPC reply budget per in-flight delivery (maxInFlight can be
|
||||
// 32 -> 8GB). Cap those small; genuine RPC replies (CLI) keep the
|
||||
// large budget.
|
||||
auto const RPC_REPLY_MAX_BYTES =
|
||||
(strMethod == "event") ? megabytes(1) : megabytes(256);
|
||||
|
||||
using namespace std::chrono_literals;
|
||||
// auto constexpr RPC_NOTIFY = 10min; // Wietse: lolwut 10 minutes for one
|
||||
// HTTP call?
|
||||
auto constexpr RPC_NOTIFY = 30s;
|
||||
|
||||
//@@start async-request
|
||||
HTTPClient::request(
|
||||
bSSL,
|
||||
io_service,
|
||||
@@ -1914,6 +1926,7 @@ fromNetwork(
|
||||
std::placeholders::_3,
|
||||
j),
|
||||
j);
|
||||
//@@end async-request
|
||||
}
|
||||
|
||||
} // namespace RPCCall
|
||||
|
||||
@@ -24,29 +24,30 @@
|
||||
#include <xrpl/basics/contract.h>
|
||||
#include <xrpl/json/to_string.h>
|
||||
#include <deque>
|
||||
#include <memory>
|
||||
|
||||
namespace ripple {
|
||||
|
||||
// Subscription object for JSON-RPC
|
||||
class RPCSubImp : public RPCSub
|
||||
class RPCSubImp : public RPCSub, public std::enable_shared_from_this<RPCSubImp>
|
||||
{
|
||||
public:
|
||||
RPCSubImp(
|
||||
InfoSub::Source& source,
|
||||
boost::asio::io_service& io_service,
|
||||
JobQueue& jobQueue,
|
||||
std::string const& strUrl,
|
||||
std::string const& strUsername,
|
||||
std::string const& strPassword,
|
||||
Logs& logs)
|
||||
Logs& logs,
|
||||
std::size_t maxQueueSize)
|
||||
: RPCSub(source)
|
||||
, m_io_service(io_service)
|
||||
, m_jobQueue(jobQueue)
|
||||
, mUrl(strUrl)
|
||||
, mSSL(false)
|
||||
, mUsername(strUsername)
|
||||
, mPassword(strPassword)
|
||||
, mSending(false)
|
||||
, maxQueueSize_(maxQueueSize)
|
||||
, j_(logs.journal("RPCSub"))
|
||||
, logs_(logs)
|
||||
{
|
||||
@@ -78,14 +79,26 @@ public:
|
||||
{
|
||||
std::lock_guard sl(mLock);
|
||||
|
||||
// Wietse: we're not going to limit this, this is admin-port only, scale
|
||||
// accordingly Dropping events just like this results in inconsistent
|
||||
// data on the receiving end if (mDeque.size() >= eventQueueMax)
|
||||
// {
|
||||
// // Drop the previous event.
|
||||
// JLOG(j_.warn()) << "RPCCall::fromNetwork drop";
|
||||
// mDeque.pop_back();
|
||||
// }
|
||||
if (mDeque.size() >= maxQueueSize_)
|
||||
{
|
||||
// Always advance mSeq so consumers can detect the gap, but
|
||||
// rate-limit the log: a hopelessly behind endpoint drops on
|
||||
// every send() and would otherwise flood the log. Warn on
|
||||
// the first drop of a run and then once per dropLogInterval.
|
||||
if (mDropped++ % dropLogInterval == 0)
|
||||
{
|
||||
JLOG(j_.warn())
|
||||
<< "RPCCall::fromNetwork drop: queue full ("
|
||||
<< mDeque.size() << "), seq=" << mSeq
|
||||
<< ", endpoint=" << mIp << ", dropped=" << mDropped;
|
||||
}
|
||||
++mSeq;
|
||||
return;
|
||||
}
|
||||
|
||||
// Endpoint caught up enough to accept again; reset so the next
|
||||
// overflow burst logs its first drop immediately.
|
||||
mDropped = 0;
|
||||
|
||||
auto jm = broadcast ? j_.debug() : j_.info();
|
||||
JLOG(jm) << "RPCCall::fromNetwork push: " << jvObj;
|
||||
@@ -97,10 +110,7 @@ public:
|
||||
// Start a sending thread.
|
||||
JLOG(j_.info()) << "RPCCall::fromNetwork start";
|
||||
|
||||
mSending = m_jobQueue.addJob(
|
||||
jtCLIENT_SUBSCRIBE, "RPCSub::sendThread", [this]() {
|
||||
sendThread();
|
||||
});
|
||||
startSendingJob();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -121,48 +131,66 @@ public:
|
||||
}
|
||||
|
||||
private:
|
||||
// XXX Could probably create a bunch of send jobs in a single get of the
|
||||
// lock.
|
||||
// Maximum concurrent HTTP deliveries per batch. Bounds file
|
||||
// descriptor usage while still allowing parallel delivery to
|
||||
// capable endpoints. With a 1024 FD process limit shared across
|
||||
// peers, clients, and the node store, 32 per subscriber is a
|
||||
// meaningful but survivable chunk even with multiple subscribers.
|
||||
static constexpr int maxInFlight = 32;
|
||||
|
||||
// Log one drop warning per this many drops while the queue stays
|
||||
// full, to avoid flooding the log on a persistently behind endpoint.
|
||||
static constexpr std::size_t dropLogInterval = 1000;
|
||||
|
||||
// Schedule a sending job. Must be called under mLock. The job holds a
|
||||
// weak_ptr and re-locks it on entry, so the RPCSub is kept alive for
|
||||
// the duration of the batch even if it is unsubscribed (and would
|
||||
// otherwise be destroyed) concurrently — sendThread dereferences this
|
||||
// only via that strong ref. mDeque events are delivered until the sub
|
||||
// is gone, after which weak.lock() fails and the job is a no-op.
|
||||
void
|
||||
startSendingJob()
|
||||
{
|
||||
std::weak_ptr<RPCSubImp> weak = weak_from_this();
|
||||
mSending = m_jobQueue.addJob(
|
||||
jtCLIENT_SUBSCRIBE, "RPCSub::sendThread", [weak]() {
|
||||
if (auto self = weak.lock())
|
||||
self->sendThread();
|
||||
});
|
||||
}
|
||||
|
||||
void
|
||||
sendThread()
|
||||
{
|
||||
Json::Value jvEvent;
|
||||
bool bSend;
|
||||
// Process exactly ONE batch per job, then re-queue if more events
|
||||
// remain, rather than draining the whole backlog in a single job.
|
||||
// A local io_service's .run() blocks this worker thread for the
|
||||
// batch (up to the per-request timeout), so re-queueing between
|
||||
// batches keeps one slow/hung subscriber from monopolising a
|
||||
// job-queue worker and starving consensus/ledger/RPC work.
|
||||
//
|
||||
// mSending must be cleared under the lock on every non-requeue
|
||||
// exit path; if it ever stays set without a job in flight, send()
|
||||
// sees mSending == true and never restarts us, stalling the queue
|
||||
// forever — the original bug (xrpld issue #6341).
|
||||
boost::asio::io_service io_service;
|
||||
int dispatched = 0;
|
||||
|
||||
do
|
||||
try
|
||||
{
|
||||
{
|
||||
// Obtain the lock to manipulate the queue and change sending.
|
||||
std::lock_guard sl(mLock);
|
||||
|
||||
if (mDeque.empty())
|
||||
{
|
||||
mSending = false;
|
||||
bSend = false;
|
||||
}
|
||||
else
|
||||
while (!mDeque.empty() && dispatched < maxInFlight)
|
||||
{
|
||||
auto const [seq, env] = mDeque.front();
|
||||
|
||||
mDeque.pop_front();
|
||||
|
||||
jvEvent = env;
|
||||
Json::Value jvEvent = env;
|
||||
jvEvent["seq"] = seq;
|
||||
|
||||
bSend = true;
|
||||
}
|
||||
}
|
||||
|
||||
// Send outside of the lock.
|
||||
if (bSend)
|
||||
{
|
||||
// XXX Might not need this in a try.
|
||||
try
|
||||
{
|
||||
JLOG(j_.info()) << "RPCCall::fromNetwork: " << mIp;
|
||||
|
||||
RPCCall::fromNetwork(
|
||||
m_io_service,
|
||||
io_service,
|
||||
mIp,
|
||||
mPort,
|
||||
mUsername,
|
||||
@@ -173,21 +201,51 @@ private:
|
||||
mSSL,
|
||||
true,
|
||||
logs_);
|
||||
}
|
||||
catch (const std::exception& e)
|
||||
{
|
||||
JLOG(j_.info())
|
||||
<< "RPCCall::fromNetwork exception: " << e.what();
|
||||
++dispatched;
|
||||
}
|
||||
}
|
||||
} while (bSend);
|
||||
|
||||
// dispatched is always > 0 here (send() only starts a job
|
||||
// after enqueuing, and the re-queue below only fires with a
|
||||
// non-empty deque), but guard anyway so an empty batch can't
|
||||
// log/spin — it falls straight through to clear mSending.
|
||||
if (dispatched > 0)
|
||||
{
|
||||
JLOG(j_.info()) << "RPCCall::fromNetwork: " << mIp
|
||||
<< " dispatching " << dispatched << " events";
|
||||
|
||||
io_service.run();
|
||||
}
|
||||
}
|
||||
catch (std::exception const& e)
|
||||
{
|
||||
// Bail rather than re-queue: a persistently failing endpoint
|
||||
// would otherwise spin the job queue. mSending is reset so the
|
||||
// next send() restarts delivery.
|
||||
JLOG(j_.warn()) << "RPCSub::sendThread exception: " << e.what();
|
||||
std::lock_guard sl(mLock);
|
||||
mSending = false;
|
||||
return;
|
||||
}
|
||||
catch (...)
|
||||
{
|
||||
JLOG(j_.warn()) << "RPCSub::sendThread unknown exception";
|
||||
std::lock_guard sl(mLock);
|
||||
mSending = false;
|
||||
return;
|
||||
}
|
||||
|
||||
// Batch complete: re-queue for the next one (mSending stays set)
|
||||
// or clear mSending if the queue drained — both under the lock to
|
||||
// avoid a lost-wakeup race with send().
|
||||
std::lock_guard sl(mLock);
|
||||
if (mDeque.empty())
|
||||
mSending = false;
|
||||
else
|
||||
startSendingJob();
|
||||
}
|
||||
|
||||
private:
|
||||
// Wietse: we're not going to limit this, this is admin-port only, scale
|
||||
// accordingly enum { eventQueueMax = 32 };
|
||||
|
||||
boost::asio::io_service& m_io_service;
|
||||
JobQueue& m_jobQueue;
|
||||
|
||||
std::string mUrl;
|
||||
@@ -200,8 +258,15 @@ private:
|
||||
|
||||
int mSeq; // Next id to allocate.
|
||||
|
||||
std::size_t mDropped = 0; // Consecutive drops while queue is full.
|
||||
|
||||
bool mSending; // Sending threead is active.
|
||||
|
||||
// Maximum queued events before dropping. The default (16384) is a
|
||||
// ~10-minute buffer at 100+ events/ledger; a hopelessly behind
|
||||
// endpoint trips it and consumers detect the gap via the seq field.
|
||||
std::size_t const maxQueueSize_;
|
||||
|
||||
std::deque<std::pair<int, Json::Value>> mDeque;
|
||||
|
||||
beast::Journal const j_;
|
||||
@@ -217,21 +282,21 @@ RPCSub::RPCSub(InfoSub::Source& source) : InfoSub(source, Consumer())
|
||||
std::shared_ptr<RPCSub>
|
||||
make_RPCSub(
|
||||
InfoSub::Source& source,
|
||||
boost::asio::io_service& io_service,
|
||||
JobQueue& jobQueue,
|
||||
std::string const& strUrl,
|
||||
std::string const& strUsername,
|
||||
std::string const& strPassword,
|
||||
Logs& logs)
|
||||
Logs& logs,
|
||||
std::size_t maxQueueSize)
|
||||
{
|
||||
return std::make_shared<RPCSubImp>(
|
||||
std::ref(source),
|
||||
std::ref(io_service),
|
||||
std::ref(jobQueue),
|
||||
strUrl,
|
||||
strUsername,
|
||||
strPassword,
|
||||
logs);
|
||||
logs,
|
||||
maxQueueSize);
|
||||
}
|
||||
|
||||
} // namespace ripple
|
||||
|
||||
@@ -26,17 +26,10 @@
|
||||
#include <xrpld/rpc/GRPCHandlers.h>
|
||||
#include <xrpld/rpc/detail/RPCHelpers.h>
|
||||
#include <xrpld/rpc/detail/TransactionSign.h>
|
||||
#include <xrpl/json/json_reader.h>
|
||||
#include <xrpl/json/json_writer.h>
|
||||
#include <xrpl/protocol/ErrorCodes.h>
|
||||
#include <xrpl/protocol/PublicKey.h>
|
||||
#include <xrpl/protocol/RPCErr.h>
|
||||
#include <xrpl/protocol/SField.h>
|
||||
#include <xrpl/protocol/STParsedJSON.h>
|
||||
#include <xrpl/resource/Fees.h>
|
||||
|
||||
#include <xrpl/protocol/JSONTxSignatures.h>
|
||||
|
||||
namespace ripple {
|
||||
|
||||
static NetworkOPs::FailHard
|
||||
@@ -90,8 +83,7 @@ doInject(RPC::JsonContext& context)
|
||||
}
|
||||
|
||||
// {
|
||||
// tx_blob: <string> XOR tx_json: <object>
|
||||
// XOR { tx: <json text>, signature: <hex> },
|
||||
// tx_blob: <string> XOR tx_json: <object>,
|
||||
// secret: <secret>
|
||||
// }
|
||||
Json::Value
|
||||
@@ -99,18 +91,7 @@ doSubmit(RPC::JsonContext& context)
|
||||
{
|
||||
context.loadType = Resource::feeMediumBurdenRPC;
|
||||
|
||||
bool const hasJsonTx = context.ledgerMaster.getCurrentLedger()->rules().enabled(featureJsonTx);
|
||||
|
||||
bool const isJsonTx = !context.params.isMember(jss::tx_blob) &&
|
||||
context.params.isMember(jss::tx) &&
|
||||
context.params.isMember(jss::sig);
|
||||
|
||||
if (isJsonTx && !hasJsonTx)
|
||||
return RPC::make_error(
|
||||
rpcNOT_SUPPORTED, "JsonTx is not enabled yet.");
|
||||
|
||||
|
||||
if (!context.params.isMember(jss::tx_blob) && !isJsonTx)
|
||||
if (!context.params.isMember(jss::tx_blob))
|
||||
{
|
||||
auto const failType = getFailHard(context);
|
||||
|
||||
@@ -138,73 +119,18 @@ doSubmit(RPC::JsonContext& context)
|
||||
|
||||
Json::Value jvResult;
|
||||
|
||||
std::optional<Blob> ret;
|
||||
if (!isJsonTx)
|
||||
{
|
||||
ret = strUnHex(context.params[jss::tx_blob].asString());
|
||||
auto ret = strUnHex(context.params[jss::tx_blob].asString());
|
||||
|
||||
if (!ret || !ret->size())
|
||||
return rpcError(rpcINVALID_PARAMS);
|
||||
}
|
||||
if (!ret || !ret->size())
|
||||
return rpcError(rpcINVALID_PARAMS);
|
||||
|
||||
SerialIter sitTrans(makeSlice(*ret));
|
||||
|
||||
std::shared_ptr<STTx const> stTx;
|
||||
|
||||
try
|
||||
{
|
||||
if (!isJsonTx)
|
||||
{
|
||||
SerialIter sitTrans(makeSlice(*ret));
|
||||
stTx = std::make_shared<STTx const>(std::ref(sitTrans));
|
||||
}
|
||||
else
|
||||
{
|
||||
std::string const raw = context.params[jss::tx].asString();
|
||||
auto const [san, diff] = sanitize_jsontx(raw);
|
||||
auto const sig = strUnHex(context.params[jss::sig].asString());
|
||||
if (!sig || sig->empty())
|
||||
throw std::runtime_error("JsonTx: bad signature");
|
||||
|
||||
Json::Value jv;
|
||||
if (Json::Reader r; !r.parse(san, jv))
|
||||
throw std::runtime_error("JsonTx: unparsable canonical form");
|
||||
|
||||
// The preimage carries the key but not the signature over itself.
|
||||
for (auto const& n :
|
||||
{sfTxnSignature.fieldName, sfSigners.fieldName})
|
||||
if (jv.isMember(n))
|
||||
throw std::runtime_error(
|
||||
"JsonTx: " + n + " must not appear in tx");
|
||||
if (!jv.isMember(sfSigningPubKey.fieldName))
|
||||
throw std::runtime_error("JsonTx: tx must carry SigningPubKey");
|
||||
|
||||
// Hand the parser the u64 rather than teaching STUInt64 a second
|
||||
// spelling; the ISO form only ever exists in the preimage.
|
||||
std::optional<std::uint64_t> ms;
|
||||
if (jv.isMember(sfTime.fieldName))
|
||||
{
|
||||
ms = jsontx_iso(jv[sfTime.fieldName].asString());
|
||||
jv.removeMember(sfTime.fieldName);
|
||||
}
|
||||
|
||||
STParsedJSONObject parsed("tx_json", jv);
|
||||
if (!parsed.object)
|
||||
throw std::runtime_error(
|
||||
parsed.error[jss::error_message].asString());
|
||||
if (ms)
|
||||
parsed.object->setFieldU64(sfTime, *ms);
|
||||
parsed.object->setFieldVL(sfTxnSignature, *sig);
|
||||
stTx = std::make_shared<STTx const>(std::move(*parsed.object));
|
||||
|
||||
// Round-trip the binary codec, then run the exact check a relaying
|
||||
// node will run, so this path cannot accept anything the network
|
||||
// would later reject.
|
||||
Serializer s;
|
||||
stTx->add(s);
|
||||
SerialIter si(s.slice());
|
||||
STTx const rt{si};
|
||||
if (jsontx_verify(rt, diff) != raw)
|
||||
throw std::runtime_error("JsonTx: does not round-trip");
|
||||
}
|
||||
stTx = std::make_shared<STTx const>(std::ref(sitTrans));
|
||||
}
|
||||
catch (std::exception& e)
|
||||
{
|
||||
@@ -215,9 +141,7 @@ doSubmit(RPC::JsonContext& context)
|
||||
}
|
||||
|
||||
{
|
||||
// JsonTx signs the plaintext preimage rather than the binary one, so
|
||||
// the binary TxnSignature check is satisfied out of band above.
|
||||
if (!context.app.checkSigs() || isJsonTx)
|
||||
if (!context.app.checkSigs())
|
||||
forceValidity(
|
||||
context.app.getHashRouter(),
|
||||
stTx->getTransactionID(),
|
||||
|
||||
@@ -76,7 +76,6 @@ doSubscribe(RPC::JsonContext& context)
|
||||
{
|
||||
auto rspSub = make_RPCSub(
|
||||
context.app.getOPs(),
|
||||
context.app.getIOService(),
|
||||
context.app.getJobQueue(),
|
||||
strUrl,
|
||||
strUsername,
|
||||
|
||||
Reference in New Issue
Block a user