coinprices/LocalPriceManager.cpp¶
Namespaces¶
| Name |
|---|
| sgns |
Source code¶
#include "LocalPriceManager.hpp"
#include "PriceFreshness.hpp"
#include <cassert>
#include <future>
#include <optional>
#include <set>
#include <utility>
namespace sgns
{
LocalPriceManager::LocalPriceManager( std::shared_ptr<IPriceSource> coinGeckoTier,
std::shared_ptr<IPriceSource> gnusServiceTier,
std::chrono::milliseconds coalescingWindow,
Clock now )
: ioc_( std::make_shared<boost::asio::io_context>() ),
work_( std::make_unique<boost::asio::executor_work_guard<boost::asio::io_context::executor_type>>(
ioc_->get_executor() ) ),
strand_( ioc_->get_executor() ),
coinGeckoTier_( std::move( coinGeckoTier ) ),
gnusServiceTier_( std::move( gnusServiceTier ) ),
coalescingWindow_( coalescingWindow ),
now_( std::move( now ) )
{
// The work guard is held for the manager's whole lifetime (no natural
// idle point โ timers and posted handlers come and go); without it
// run() returns the instant the queue drains.
thread_ = std::thread( [ioc = ioc_]() { ioc->run(); } );
}
LocalPriceManager::~LocalPriceManager()
{
// Shutdown ordering (drain-then-join): the shutdown post runs on the
// strand while the guard still holds run() alive; resetting the guard
// unblocks run() after the queue drains; join precedes member
// destruction so handlers referencing `this` are always safe.
// ioc_->stop() is NEVER called โ it would abandon queued handlers
// (the HttpStubServer uses stop() only because it has no waiters).
boost::asio::post( strand_, [this]() {
m_logger->info( "LocalPriceManager shutting down" );
// Resolve parked waiters BEFORE reset/join: this is what keeps
// blocked GetQuotes callers from hanging at shutdown
// (RESEARCH ยง1 landmine 1). Cancel each open window's timer and
// fail its waiters โ the walk can no longer run for them.
for ( auto &[currency, window] : windows_ )
{
if ( window.timer )
{
boost::system::error_code ec;
window.timer->cancel( ec );
}
for ( auto &waiter : window.waiters )
{
waiter.done.set_value(
outcome::failure( PriceFetchFailure{ PriceFetchError::NetworkError, 0 } ) );
}
}
windows_.clear();
} );
work_->reset();
thread_.join();
}
PriceResult<std::vector<PriceQuote>> LocalPriceManager::GetQuotes( const std::vector<std::string> &ids,
const std::string ¤cy )
{
// Cheap self-deadlock guard: GetQuotes parks on a future that only a
// strand handler can satisfy โ calling it from the manager's own
// runner thread would deadlock (debug-asserted, Doxygen-documented).
assert( !strand_.running_in_this_thread() );
if ( ids.empty() )
{
// Facade parity (PriceHttpClient.cpp:49-52): empty input fails
// synchronously without posting.
return outcome::failure( PriceFetchFailure{ PriceFetchError::EmptyInput, 0 } );
}
std::promise<PriceResult<std::vector<PriceQuote>>> promise;
auto future = promise.get_future();
boost::asio::post( strand_,
[this, ids, currency, p = std::move( promise )]() mutable {
HandleRequestOnStrand( ids, currency, std::move( p ) );
} );
return future.get(); // D-01: blocking bridge โ caller parks until the chain resolves
}
void LocalPriceManager::HandleRequestOnStrand( const std::vector<std::string> &ids,
const std::string ¤cy,
std::promise<PriceResult<std::vector<PriceQuote>>> done )
{
// D-06 fresh/miss split โ modeled on the GeniusNode::GetCoinprice
// miss-collection loop with ClassifyFreshness as the only freshness
// authority (bands come from PriceFreshness.hpp constants).
std::vector<PriceQuote> immediate;
std::vector<std::string> misses;
const auto currentTime = now_();
for ( const auto &id : ids )
{
const auto currencyIt = cache_.find( currency );
if ( currencyIt != cache_.end() )
{
const auto idIt = currencyIt->second.find( id );
if ( idIt != currencyIt->second.end()
&& ClassifyFreshness( idIt->second.timestamp, currentTime ) == FreshnessBand::Fresh )
{
// Copy-on-serve (P-9): the stored entry keeps its
// fetch-time source and timestamp for diagnostics; only
// the served copy is rewritten to LocalCache.
PriceQuote served = idIt->second;
served.source = PriceSource::LocalCache;
immediate.push_back( std::move( served ) );
continue;
}
}
misses.push_back( id );
}
if ( misses.empty() )
{
// LPM-01: every requested id fresh in L1 โ zero network.
m_logger->debug( "GetQuotes served {} id(s) entirely from L1 (currency {})", immediate.size(), currency );
done.set_value( outcome::success( std::move( immediate ) ) );
return;
}
// D-05: misses join the currency's pending window; the first miss
// arms the ~50ms coalescing timer. Ids already fetched in the
// current walk are NOT re-added because dispatch removes the window
// from windows_ BEFORE walking (a new miss creates a fresh window โ
// D-07).
PendingWindow &window = windows_[currency];
if ( window.timer == nullptr )
{
window.timer = std::make_unique<boost::asio::steady_timer>( *ioc_ );
window.timer->expires_after( coalescingWindow_ );
// The timer is constructed from *ioc_, NOT from the strand โ a
// bare lambda would therefore NOT be strand-bound;
// bind_executor( strand_, ... ) puts the callback on the strand
// explicitly. The load-bearing safety invariant is that exactly
// ONE thread runs ioc_: single runner thread + strand together
// serialize all window state. A second runner thread on ioc_ โ
// not strand membership โ is what would break the
// timer-to-dispatch serialization.
window.timer->async_wait( boost::asio::bind_executor(
strand_, [this, currency]( const boost::system::error_code & ) { DispatchBatchOnStrand( currency ); } ) );
}
window.ids.insert( misses.begin(), misses.end() ); // set union (D-05)
window.waiters.push_back( PendingWaiter{ ids, std::move( immediate ), std::move( done ) } );
// The caller (including the FIRST caller โ D-05's accepted cost)
// blocks on its future until the window fires and the walk lands.
}
void LocalPriceManager::DispatchBatchOnStrand( const std::string ¤cy )
{
auto windowIt = windows_.find( currency );
if ( windowIt == windows_.end() )
{
return; // already dispatched (or shutdown cleared it)
}
// The window closes the moment dispatch begins: requests arriving
// during the walk land in a NEW window (D-05/D-07).
PendingWindow batch = std::move( windowIt->second );
windows_.erase( windowIt );
const std::vector<std::string> missIds( batch.ids.begin(), batch.ids.end() );
if ( missIds.empty() )
{
return;
}
// The fallback chain walk (LPM-04, D-09), inline on the strand:
// a FAILED tier leaves `remaining` untouched (wholesale escalation
// of the full miss-set); a successful-but-partial tier shrinks it by
// exactly the ids it returned (gap-chase); completed ids are never
// re-added because they left `remaining`.
std::set<std::string> remaining( batch.ids.begin(), batch.ids.end() );
std::optional<PriceFetchFailure> lastTierFailure;
std::vector<PriceQuote> fetched; // quotes returned by the walk โ served with TIER source
// Tier 1 โ CoinGecko direct. Always has ids at entry (the batch
// exists because a miss joined it).
{
auto r1 = coinGeckoTier_->FetchPrices( std::vector<std::string>( remaining.begin(), remaining.end() ),
currency );
if ( r1 )
{
StoreInL1( r1.value() );
fetched.insert( fetched.end(), r1.value().begin(), r1.value().end() );
for ( const auto "e : r1.value() )
{
remaining.erase( quote.asset );
}
m_logger->debug( "tier 1 (CoinGecko) served {} id(s)", r1.value().size() );
}
else
{
m_logger->warn( "tier 1 (CoinGecko) failed: {}", r1.error().Message() );
lastTierFailure = r1.error();
}
}
// Tier 2 โ token.gnus.ai: ONLY for ids tier 1 did not serve.
if ( !remaining.empty() )
{
auto r2 = gnusServiceTier_->FetchPrices( std::vector<std::string>( remaining.begin(), remaining.end() ),
currency );
if ( r2 )
{
StoreInL1( r2.value() );
fetched.insert( fetched.end(), r2.value().begin(), r2.value().end() );
for ( const auto "e : r2.value() )
{
remaining.erase( quote.asset );
}
m_logger->debug( "tier 2 (token.gnus.ai) served {} id(s)", r2.value().size() );
}
else
{
m_logger->warn( "tier 2 (token.gnus.ai) failed: {}", r2.error().Message() );
lastTierFailure = r2.error();
}
}
// Per-waiter resolution per the serving-source rule, extended with
// band-aware L1 lookups (D-10/D-11) for requested ids neither the
// immediate set nor either tier covered โ exactly the ids both
// network tiers failed to refresh (last-known-good territory).
const auto now = now_();
int resolvedFresh = 0, resolvedStale = 0, resolvedFailed = 0;
for ( auto &waiter : batch.waiters )
{
std::vector<PriceQuote> assembled = std::move( waiter.immediate );
// The walk's returned quotes serve with the TIER's source as
// received (no rewrite โ fetch provenance preserved).
for ( const auto "e : fetched )
{
for ( const auto &id : waiter.requestedIds )
{
if ( quote.asset == id )
{
assembled.push_back( quote );
break;
}
}
}
for ( const auto &id : waiter.requestedIds )
{
bool covered = false;
for ( const auto "e : assembled )
{
if ( quote.asset == id )
{
covered = true;
break;
}
}
if ( covered )
{
continue; // fresh-immediate or tier-returned โ nothing to look up
}
const auto currencyIt = cache_.find( currency );
if ( currencyIt == cache_.end() )
{
continue;
}
const auto idIt = currencyIt->second.find( id );
if ( idIt == currencyIt->second.end() )
{
continue;
}
const auto band = ClassifyFreshness( idIt->second.timestamp, now );
if ( band == FreshnessBand::Fresh )
{
// Rare but reachable: an overlapping earlier batch
// refreshed the entry after this request classified it
// a miss.
PriceQuote served = idIt->second;
served.source = PriceSource::LocalCache;
assembled.push_back( std::move( served ) );
}
else if ( band == FreshnessBand::StaleButUsable )
{
// LKG serve (D-10): only reachable when both tiers
// failed to refresh the id โ a tier success would have
// stored a fresh entry (structural FRESH-02 compliance).
// stale=true, timestamp UNTOUCHED (D-11).
PriceQuote served = idIt->second;
served.source = PriceSource::LocalCache;
served.stale = true;
assembled.push_back( std::move( served ) );
}
// Unavailable band: never served (FRESH-02) โ skip.
}
if ( assembled.empty() )
{
// Nothing servable for this waiter: surface the last tier
// failure, or NoDataFound when the tiers succeeded but
// returned nothing for the requested ids.
if ( lastTierFailure.has_value() )
{
waiter.done.set_value( outcome::failure( *lastTierFailure ) );
}
else
{
waiter.done.set_value(
outcome::failure( PriceFetchFailure{ PriceFetchError::NoDataFound, 0 } ) );
}
++resolvedFailed;
continue;
}
// A non-empty immediate set with failed tiers still resolves
// SUCCESS with the partial data (GetCoinprice's
// continue-with-what-we-have).
if ( lastTierFailure.has_value() )
{
++resolvedStale; // partial success over a failed walk
}
else
{
++resolvedFresh;
}
waiter.done.set_value( outcome::success( std::move( assembled ) ) );
}
// Chain trace (RESEARCH ยง8.7): the serving tier is otherwise
// indistinguishable when envelope quotes carry source=CoinGecko.
m_logger->info( "Dispatched batch of {} id(s), {} waiter(s): {}/{} unresolved after tiers, {} fresh / {} stale-or-partial / {} failed",
missIds.size(),
batch.waiters.size(),
remaining.size(),
missIds.size(),
resolvedFresh,
resolvedStale,
resolvedFailed );
}
void LocalPriceManager::StoreInL1( const std::vector<PriceQuote> "es )
{
// Keeps fetch-time source and timestamp verbatim. Only called on tier
// SUCCESS โ tier failures never write L1 entries (no negative
// caching, PITFALLS #14).
for ( const auto "e : quotes )
{
cache_[quote.currency][quote.asset] = quote;
}
}
} // namespace sgns
Updated on 2026-10-09 at 04:02:50 +0000