Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 7 additions & 7 deletions .github/workflows/docker-build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -70,20 +70,20 @@ jobs:
sleep 1
done

# /api/health: status=ok, offer_count=0 on a fresh container.
# /api/health: status=ok, maker_count=0 on a fresh container.
health=$(curl -fsS http://127.0.0.1:3000/api/health)
echo "Health: $health"
if ! echo "$health" | jq -e '.status == "ok" and .offer_count == 0' > /dev/null; then
if ! echo "$health" | jq -e '.status == "ok" and .maker_count == 0 and .with_offer == 0' > /dev/null; then
echo "Health assertion failed"
docker logs marketd-test || true
exit 1
fi

# /api/offers: empty array on a fresh container.
offers=$(curl -fsS http://127.0.0.1:3000/api/offers)
echo "Offers: $offers"
if ! echo "$offers" | jq -e 'type == "array" and length == 0' > /dev/null; then
echo "Offers assertion failed"
# /api/makers: empty array on a fresh container.
makers=$(curl -fsS http://127.0.0.1:3000/api/makers)
echo "Makers: $makers"
if ! echo "$makers" | jq -e 'type == "array" and length == 0' > /dev/null; then
echo "Makers assertion failed"
docker logs marketd-test || true
exit 1
fi
Expand Down
38 changes: 30 additions & 8 deletions src/server.rs
Original file line number Diff line number Diff line change
@@ -1,38 +1,60 @@
use axum::{Router, extract::State, http::StatusCode, response::Json, routing::get};
use axum::{
Router,
extract::{Query, State},
http::StatusCode,
response::Json,
routing::get,
};
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use tower_http::{
cors::CorsLayer,
services::{ServeDir, ServeFile},
};

use crate::state::SharedStore;
use crate::state::{ApiMaker, ApiMakerState, SharedStore};

pub fn router(store: SharedStore, static_dir: String) -> Router {
let index = format!("{static_dir}/index.html");
let spa = ServeDir::new(&static_dir).not_found_service(ServeFile::new(index));

Router::new()
.route("/api/offers", get(get_offers))
.route("/api/makers", get(get_makers))
.route("/api/health", get(get_health))
.with_state(store)
.layer(CorsLayer::permissive())
.fallback_service(spa)
}

async fn get_offers(State(store): State<SharedStore>) -> Result<Json<Value>, StatusCode> {
let offers = store
#[derive(Serialize, Deserialize)]
struct MakerQueryParams {
state: Option<ApiMakerState>,
}

async fn get_makers(
State(store): State<SharedStore>,
Query(params): Query<MakerQueryParams>,
) -> Result<Json<Value>, StatusCode> {
let makers = store
.read()
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?
.offers
.makers
.clone();
Ok(Json(json!(offers)))
let makers = makers
.iter()
.filter(|m| params.state.as_ref().is_none_or(|state| state == &m.state))
.collect::<Vec<&ApiMaker>>();

Ok(Json(json!(makers)))
}

async fn get_health(State(store): State<SharedStore>) -> Json<Value> {
let s = store.read().unwrap();
let with_offer = s.makers.iter().filter(|m| m.offer.is_some()).count();
Json(json!({
"status": "ok",
"offer_count": s.offers.len(),
"maker_count": s.makers.len(),
"with_offer": with_offer,
"last_sync": s.last_sync,
}))
}
84 changes: 70 additions & 14 deletions src/state.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
use coinswap::{protocol::common_messages::Offer, taker::offers::MakerAddress};
use serde::Serialize;
use coinswap::{
protocol::common_messages::Offer,
taker::offers::{MakerOfferCandidate, MakerProtocol, MakerState},
};
use serde::{Deserialize, Serialize};
use std::sync::{Arc, RwLock};

#[derive(Serialize, Clone, Debug)]
Expand All @@ -17,10 +20,10 @@ pub struct ApiFidelityBond {
pub cert_sig: String,
}

/// Offer payload — populated only when the taker has successfully fetched
/// the maker's advertisement.
#[derive(Serialize, Clone, Debug)]
pub struct ApiOffer {
pub address: String,
pub timestamp: u64,
pub base_fee: u64,
pub amount_relative_fee_pct: f64,
pub time_relative_fee_pct: f64,
Expand All @@ -33,15 +36,11 @@ pub struct ApiOffer {
}

impl ApiOffer {
pub fn from_coinswap(offer: &Offer, address: &MakerAddress, timestamp: u64) -> Self {
pub fn from_coinswap(offer: &Offer) -> Self {
let bond = &offer.fidelity.bond;
let outpoint = bond.outpoint();
let txid = outpoint.txid.to_string();
let vout = outpoint.vout;

Self {
address: address.to_string(),
timestamp,
base_fee: offer.base_fee,
amount_relative_fee_pct: offer.amount_relative_fee_pct,
time_relative_fee_pct: offer.time_relative_fee_pct,
Expand All @@ -52,7 +51,10 @@ impl ApiOffer {
tweakable_point: offer.tweakable_point.to_string(),
fidelity_bond: ApiFidelityBond {
amount: bond.amount.to_sat(),
outpoint: ApiOutpoint { txid, vout },
outpoint: ApiOutpoint {
txid: outpoint.txid.to_string(),
vout: outpoint.vout,
},
lock_time: bond.lock_time.to_consensus_u32(),
cert_hash: offer.fidelity.cert_hash.to_string(),
cert_sig: hex::encode(offer.fidelity.cert_sig.serialize_der().as_ref()),
Expand All @@ -61,14 +63,68 @@ impl ApiOffer {
}
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ApiMakerState {
Good,
Unresponsive { retries: u8 },
Bad,
}

impl From<&MakerState> for ApiMakerState {
fn from(s: &MakerState) -> Self {
match s {
MakerState::Good => Self::Good,
MakerState::Unresponsive { retries } => Self::Unresponsive { retries: *retries },
MakerState::Bad => Self::Bad,
}
}
}

fn protocol_label(p: &MakerProtocol) -> &'static str {
match p {
MakerProtocol::Legacy => "legacy",
MakerProtocol::Taproot => "taproot",
MakerProtocol::Unified => "unified",
}
}

/// A maker known to the taker's offerbook — *whether or not* we have a
/// current offer for it. Bad and unresponsive makers come through here too,
/// with `offer: None`.
#[derive(Serialize, Clone, Debug)]
pub struct ApiMaker {
pub address: String,
pub state: ApiMakerState,
pub protocol: Option<&'static str>,
pub timestamp: u64,
pub last_offer_update_ts: Option<u64>,
pub next_offer_check_ts: Option<u64>,
pub offer: Option<ApiOffer>,
}

impl ApiMaker {
pub fn from_candidate(candidate: &MakerOfferCandidate, timestamp: u64) -> Self {
Self {
address: candidate.address.to_string(),
state: (&candidate.state).into(),
protocol: candidate.protocol.as_ref().map(protocol_label),
timestamp,
last_offer_update_ts: candidate.last_offer_update_ts,
next_offer_check_ts: candidate.next_offer_check_ts,
offer: candidate.offer.as_ref().map(ApiOffer::from_coinswap),
}
}
}

#[derive(Default, Clone)]
pub struct OfferStore {
pub offers: Vec<ApiOffer>,
pub struct MakerStore {
pub makers: Vec<ApiMaker>,
pub last_sync: Option<u64>,
}

pub type SharedStore = Arc<RwLock<OfferStore>>;
pub type SharedStore = Arc<RwLock<MakerStore>>;

pub fn new_store() -> SharedStore {
Arc::new(RwLock::new(OfferStore::default()))
Arc::new(RwLock::new(MakerStore::default()))
}
26 changes: 14 additions & 12 deletions src/sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ pub fn build_taker_config(cfg: &Config) -> TakerInitConfig {
}

pub fn sync_loop(init_config: TakerInitConfig, sync_interval_secs: u64, store: state::SharedStore) {
use coinswap::taker::{Taker, offers::MakerState};
use coinswap::taker::Taker;

if let Some(rpc) = init_config.rpc_config.as_ref() {
wait_for_tcp(&rpc.url, "Bitcoin RPC");
Expand Down Expand Up @@ -97,25 +97,27 @@ pub fn sync_loop(init_config: TakerInitConfig, sync_interval_secs: u64, store: s
.unwrap_or_default()
.as_secs();

let offers: Vec<state::ApiOffer> = book
// Include every maker the offerbook tracks — Good, Unresponsive, Bad.
// The API consumer decides what to display; we don't filter here.
let makers: Vec<state::ApiMaker> = book
.all_makers()
.into_iter()
.filter(|m| matches!(m.state, MakerState::Good))
.filter_map(|m| {
m.offer
.as_ref()
.map(|o| state::ApiOffer::from_coinswap(o, &m.address, timestamp))
})
.iter()
.map(|m| state::ApiMaker::from_candidate(m, timestamp))
.collect();

let count = offers.len();
let count = makers.len();
let with_offer = makers.iter().filter(|m| m.offer.is_some()).count();
{
let mut s = store.write().unwrap();
s.offers = offers;
s.makers = makers;
s.last_sync = Some(timestamp);
}

tracing::info!(count, "Sync done, sleeping {sync_interval_secs}s");
tracing::info!(
count,
with_offer,
"Sync done, sleeping {sync_interval_secs}s"
);
std::thread::sleep(Duration::from_secs(sync_interval_secs));
}
}
38 changes: 21 additions & 17 deletions tests/e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@
//! Spawns a funded, fidelity-bonded `MakerServer` in-process, lets it
//! broadcast its offer to the nostr relay, runs marketd's `sync_loop`
//! against the same relay, and asserts the offer reaches
//! `GET /api/offers` with the right fields.
//! `GET /api/makers` with the right fields.
//!
//! Both tests require `BITCOIND_EXE` and assume a nostr-rs-relay listening
//! on `ws://127.0.0.1:8000` (the docker entrypoint provides this).
Expand Down Expand Up @@ -91,15 +91,16 @@ fn taker_init_and_http_layer_against_regtest() {
assert_eq!(resp.status(), 200);
let body: Value = resp.into_json().expect("decode /api/health");
assert_eq!(body["status"], "ok");
assert_eq!(body["offer_count"], 0);
assert_eq!(body["maker_count"], 0);
assert_eq!(body["with_offer"], 0);

let resp = guard
.agent
.get(&guard.url("/api/offers"))
.get(&guard.url("/api/makers"))
.call()
.expect("GET /api/offers");
.expect("GET /api/makers");
assert_eq!(resp.status(), 200);
let offers: Value = resp.into_json().expect("decode /api/offers");
let offers: Value = resp.into_json().expect("decode /api/makers");
assert_eq!(offers.as_array().expect("array").len(), 0);
}

Expand Down Expand Up @@ -234,30 +235,32 @@ fn maker_offer_flows_through_to_api_offers() {
// ───────── HTTP layer ─────────
let guard = common::MarketdServerGuard::start(store.clone());

// Poll /api/offers until non-empty.
// Poll /api/makers until a maker with an offer appears.
let offer_deadline = Instant::now() + Duration::from_secs(120);
let offer = loop {
let maker_json = loop {
assert!(
Instant::now() < offer_deadline,
"marketd /api/offers stayed empty for 120s"
"marketd /api/makers never saw a maker with an offer (120s)"
);

let resp = guard
.agent
.get(&guard.url("/api/offers"))
.get(&guard.url("/api/makers"))
.call()
.expect("GET /api/offers");
let body: Value = resp.into_json().expect("decode /api/offers");
let arr = body.as_array().expect("offers JSON is array");
if let Some(first) = arr.first() {
break first.clone();
.expect("GET /api/makers");
let body: Value = resp.into_json().expect("decode /api/makers");
let arr = body.as_array().expect("makers JSON is array");
if let Some(m) = arr.iter().find(|m| !m["offer"].is_null()) {
break m.clone();
}
thread::sleep(Duration::from_secs(2));
};

tracing::info!(?offer, "marketd saw an offer");
tracing::info!(?maker_json, "marketd saw a maker with an offer");

assert_eq!(offer["address"], format!("127.0.0.1:{maker_net_port}"));
assert_eq!(maker_json["address"], format!("127.0.0.1:{maker_net_port}"));
assert_eq!(maker_json["state"]["kind"], "good");
let offer = &maker_json["offer"];
assert_eq!(offer["base_fee"], 1000);
assert_eq!(offer["min_size"], 10_000);
assert_eq!(offer["required_confirms"], 1);
Expand All @@ -281,7 +284,8 @@ fn maker_offer_flows_through_to_api_offers() {
.expect("GET /api/health");
let body: Value = resp.into_json().expect("decode /api/health");
assert_eq!(body["status"], "ok");
assert!(body["offer_count"].as_u64().unwrap() >= 1);
assert!(body["maker_count"].as_u64().unwrap() >= 1);
assert!(body["with_offer"].as_u64().unwrap() >= 1);
assert!(body["last_sync"].is_number());

// ───────── Teardown ─────────
Expand Down
Loading
Loading