Compare commits

..

2 Commits

Author SHA1 Message Date
Richard Holland
e904290e4e jsontx stuff 2026-08-05 12:43:25 +10:00
Richard Holland
723103150c init jsontxsig amendment 2026-07-31 12:36:39 +10:00
16 changed files with 735 additions and 1726 deletions

View File

@@ -77,11 +77,6 @@ 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

View File

@@ -124,6 +124,9 @@ 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)

View File

@@ -35,6 +35,7 @@ class Xrpl(ConanFile):
'soci/4.0.3@xahaud/stable',
'xxhash/0.8.2',
'zlib/1.3.1',
'fmt/12.1.0',
]
tool_requires = [
@@ -191,6 +192,7 @@ class Xrpl(ConanFile):
'sqlite3::sqlite',
'xxhash::xxhash',
'zlib::zlib',
'fmt::fmt',
]
if self.options.rocksdb:
libxrpl.requires.append('rocksdb::librocksdb')

View File

@@ -0,0 +1,557 @@
//------------------------------------------------------------------------------
/*
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

View File

@@ -34,6 +34,7 @@
// 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)

View File

@@ -153,6 +153,7 @@ 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)
@@ -293,6 +294,7 @@ 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)

View File

@@ -630,6 +630,7 @@ 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

View File

@@ -49,6 +49,8 @@ TxFormats::TxFormats()
{sfNetworkID, soeOPTIONAL},
{sfHookParameters, soeOPTIONAL},
{sfHookName, soeOPTIONAL},
{sfTime, soeOPTIONAL},
{sfJsonTxDelta, soeOPTIONAL},
};
#pragma push_macro("UNWRAP")

File diff suppressed because it is too large Load Diff

View File

@@ -1,351 +0,0 @@
//------------------------------------------------------------------------------
/*
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

View File

@@ -22,6 +22,7 @@
#include <xrpld/core/JobQueue.h>
#include <xrpld/net/InfoSub.h>
#include <boost/asio/io_service.hpp>
namespace ripple {
@@ -38,17 +39,16 @@ 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,
// 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);
Logs& logs);
} // namespace ripple

View File

@@ -122,20 +122,12 @@ 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,
this,
shared_from_this(),
strPath,
std::placeholders::_1,
std::placeholders::_2),
@@ -401,12 +393,8 @@ 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 (hasContentLength)
if (boost::regex_match(strHeader, smMatch, reSize))
return beast::lexicalCast<std::size_t>(
std::string(smMatch[1]), maxResponseSize_);
return maxResponseSize_;
@@ -457,24 +445,22 @@ public:
JLOG(j_.trace()) << "Read error: " << mShutdown.message();
invokeComplete(mShutdown);
return;
}
// 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);
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);
}
}
}
// Call cancel the deadline timer and invoke the completion routine.
@@ -530,7 +516,6 @@ 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;

View File

@@ -1585,10 +1585,6 @@ 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 "
@@ -1752,7 +1748,6 @@ rpcClient(
}
{
//@@start blocking-request
boost::asio::io_service isService;
RPCCall::fromNetwork(
isService,
@@ -1776,7 +1771,6 @@ rpcClient(
headers);
isService.run(); // This blocks until there are no more
// outstanding async calls.
//@@end blocking-request
}
if (jvOutput.isMember("result"))
{
@@ -1887,21 +1881,15 @@ fromNetwork(
// Send request
// 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);
// Number of bytes to try to receive if no
// Content-Length header received
constexpr auto RPC_REPLY_MAX_BYTES = 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,
@@ -1926,7 +1914,6 @@ fromNetwork(
std::placeholders::_3,
j),
j);
//@@end async-request
}
} // namespace RPCCall

View File

@@ -24,30 +24,29 @@
#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, public std::enable_shared_from_this<RPCSubImp>
class RPCSubImp : public RPCSub
{
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,
std::size_t maxQueueSize)
Logs& logs)
: 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)
{
@@ -79,26 +78,14 @@ public:
{
std::lock_guard sl(mLock);
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;
// 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();
// }
auto jm = broadcast ? j_.debug() : j_.info();
JLOG(jm) << "RPCCall::fromNetwork push: " << jvObj;
@@ -110,7 +97,10 @@ public:
// Start a sending thread.
JLOG(j_.info()) << "RPCCall::fromNetwork start";
startSendingJob();
mSending = m_jobQueue.addJob(
jtCLIENT_SUBSCRIBE, "RPCSub::sendThread", [this]() {
sendThread();
});
}
}
@@ -131,66 +121,48 @@ public:
}
private:
// 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();
});
}
// XXX Could probably create a bunch of send jobs in a single get of the
// lock.
void
sendThread()
{
// 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;
Json::Value jvEvent;
bool bSend;
try
do
{
{
// Obtain the lock to manipulate the queue and change sending.
std::lock_guard sl(mLock);
while (!mDeque.empty() && dispatched < maxInFlight)
if (mDeque.empty())
{
mSending = false;
bSend = false;
}
else
{
auto const [seq, env] = mDeque.front();
mDeque.pop_front();
Json::Value jvEvent = env;
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(
io_service,
m_io_service,
mIp,
mPort,
mUsername,
@@ -201,51 +173,21 @@ private:
mSSL,
true,
logs_);
++dispatched;
}
catch (const std::exception& e)
{
JLOG(j_.info())
<< "RPCCall::fromNetwork exception: " << e.what();
}
}
// 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();
} while (bSend);
}
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;
@@ -258,15 +200,8 @@ 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_;
@@ -282,21 +217,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,
std::size_t maxQueueSize)
Logs& logs)
{
return std::make_shared<RPCSubImp>(
std::ref(source),
std::ref(io_service),
std::ref(jobQueue),
strUrl,
strUsername,
strPassword,
logs,
maxQueueSize);
logs);
}
} // namespace ripple

View File

@@ -26,10 +26,17 @@
#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
@@ -83,7 +90,8 @@ doInject(RPC::JsonContext& context)
}
// {
// tx_blob: <string> XOR tx_json: <object>,
// tx_blob: <string> XOR tx_json: <object>
// XOR { tx: <json text>, signature: <hex> },
// secret: <secret>
// }
Json::Value
@@ -91,7 +99,18 @@ doSubmit(RPC::JsonContext& context)
{
context.loadType = Resource::feeMediumBurdenRPC;
if (!context.params.isMember(jss::tx_blob))
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)
{
auto const failType = getFailHard(context);
@@ -119,18 +138,73 @@ doSubmit(RPC::JsonContext& context)
Json::Value jvResult;
auto ret = strUnHex(context.params[jss::tx_blob].asString());
std::optional<Blob> ret;
if (!isJsonTx)
{
ret = strUnHex(context.params[jss::tx_blob].asString());
if (!ret || !ret->size())
return rpcError(rpcINVALID_PARAMS);
SerialIter sitTrans(makeSlice(*ret));
if (!ret || !ret->size())
return rpcError(rpcINVALID_PARAMS);
}
std::shared_ptr<STTx const> stTx;
try
{
stTx = std::make_shared<STTx const>(std::ref(sitTrans));
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");
}
}
catch (std::exception& e)
{
@@ -141,7 +215,9 @@ doSubmit(RPC::JsonContext& context)
}
{
if (!context.app.checkSigs())
// 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)
forceValidity(
context.app.getHashRouter(),
stTx->getTransactionID(),

View File

@@ -76,6 +76,7 @@ doSubscribe(RPC::JsonContext& context)
{
auto rspSub = make_RPCSub(
context.app.getOPs(),
context.app.getIOService(),
context.app.getJobQueue(),
strUrl,
strUsername,