signer/sidecar.mjs
v0.1.0 · 96.3 KB
/** @experimental Optional Phoenix signing process. No trading decisions; every policy value is supplied by the caller. */
import crypto from 'node:crypto';
import fs from 'node:fs';
import http from 'node:http';
import path from 'node:path';
import { fileURLToPath } from 'node:url';
import {
buildCancelAllIx,
buildCancelOrdersByIdIx,
buildPlaceLimitOrderIx,
buildPlaceMarketOrderIx,
buildPlacePostOnlyOrderIx,
getPhoenixTraderSubaccountAddress,
OrderFlags,
SelfTradeBehavior,
Side,
} from '@ellipsis-labs/rise';
import {
address,
appendTransactionMessageInstructions,
compileTransaction,
createKeyPairSignerFromBytes,
createTransactionMessage,
getBase58Decoder,
getBase58Encoder,
getBase64EncodedWireTransaction,
getSignatureFromTransaction,
pipe,
setTransactionMessageFeePayer,
setTransactionMessageFeePayerSigner,
setTransactionMessageLifetimeUsingBlockhash,
signTransactionMessageWithSigners,
} from '@solana/kit';
export const VERSION = 'markpaper-phx-signer/1';
export const PROGRAM_ADDRESS = 'EtrnLzgbS7nMMy5fbD42kXiUzGg8XQzJ972Xtk1cjWih'; // privacy-allow: public mainnet protocol address, rise-public and Solana compute budget
export const LOG_AUTHORITY = 'GdxfTLSsdSY37G6fZoYtdGDSfgFnbT2EmRpuePZxWShS'; // privacy-allow: public mainnet protocol address, rise-public and Solana compute budget
export const GLOBAL_CONFIGURATION = '2zskx2iyCvb6Stg7RBZkt1f6MrF4dpYtMG3yMvKwqtUZ'; // privacy-allow: public mainnet protocol address, rise-public and Solana compute budget
export const COMPUTE_BUDGET_PROGRAM = 'ComputeBudget111111111111111111111111111111'; // privacy-allow: public mainnet protocol address, rise-public and Solana compute budget
export const U64_MAX = (1n << 64n) - 1n;
const BID_SEQ_MIN = 1n << 63n;
const MAX_CU = 1400000; // Solana documented maximum compute units per transaction
export const MAX_TX_BYTES = 1232;
export const MAX_CANCEL_IDS = 30;
const b58enc = getBase58Encoder();
const b58dec = getBase58Decoder();
const BASE58_RE = /^[1-9A-HJ-NP-Za-km-z]+$/;
const QUOTE_LOTS_PER_USD = 1000000n;
export const PREF_DISABLE_COLLATERAL_SWEEP = 1 << 0;
export const PREF_DISABLE_POSITION_AUTHORITY_SWAP = 1 << 1;
const SPOT_FLAG_IS_ACTIVE = 1 << 0;
const SPOT_FLAG_DISABLE_POSITION_AUTHORITY_SWAP = 1 << 2;
export class HttpError extends Error {
constructor(status, code, message, extra = null) {
super(message);
this.status = status;
this.code = code;
this.extra = extra;
}
}
export class ConfigError extends Error {}
const bad = (message) => new HttpError(400, 'bad_request', message);
export const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
export function toJson(v) {
return JSON.stringify(v, (_k, x) => (typeof x === 'bigint' ? x.toString() : x));
}
function sha256(data) {
return crypto.createHash('sha256').update(data).digest();
}
export function pubkeyOrNull(s) {
if (typeof s !== 'string' || s.length < 32 || s.length > 44 || !BASE58_RE.test(s)) return null;
try {
return b58enc.encode(s).length === 32 ? s : null;
} catch {
return null;
}
}
export function isSignatureString(s) {
if (typeof s !== 'string' || s.length < 64 || s.length > 90 || !BASE58_RE.test(s)) return false;
try {
return b58enc.encode(s).length === 64;
} catch {
return false;
}
}
export function parseU64(v, name, { min = 0n } = {}) {
if (typeof v !== 'string') throw bad(`${name}: expected an exact u64 decimal string, ${typeof v}`);
if (!/^(0|[1-9]\d{0,19})$/.test(v)) throw bad(`${name}: expected an exact u64 decimal string (${v.slice(0, 40)})`);
const x = BigInt(v);
if (x > U64_MAX) throw bad(`${name}: expected an exact u64 decimal string`);
if (x < min) throw bad(`${name}: must be at least ${min}`);
return x;
}
export function maskUrl(u) {
try {
const x = new URL(u);
const hidden = (x.pathname && x.pathname !== '/') || x.search || x.username || x.password;
return `${x.protocol}//${x.host}${hidden ? '/***' : ''}`;
} catch {
return '<invalid url>';
}
}
function makeSanitizer(secrets) {
const list = secrets.filter((s) => typeof s === 'string' && s.length >= 8);
return (msg) => {
let s = String(msg ?? '');
for (const x of list) s = s.split(x).join('***');
return s.replace(/(api[-_]?key=)[^&\s"']+/gi, '$1***');
};
}
function envInt(env, name, def, min, max) {
if (def === undefined && (env[name] === undefined || String(env[name]).trim() === ''))
throw new ConfigError('Required setting: ' + name);
const raw = env[name];
if (raw === undefined || raw === '') return def;
if (!/^-?\d+$/.test(String(raw).trim())) throw new ConfigError(`${name}: expected an integer, received ${raw}`);
const v = Number(String(raw).trim());
if (!Number.isSafeInteger(v) || v < min || v > max) throw new ConfigError(`${name}: ${v} [${min}, ${max}]`);
return v;
}
const MULT_SCALE = 10000n;
function envMultBps(env, name, maxWhole) {
const raw = env[name];
if (raw === undefined || String(raw).trim() === '') return 0n;
const s = String(raw).trim();
const m = /^(\d{1,9})(?:\.(\d{1,4}))?$/.exec(s);
if (!m)
throw new ConfigError(`${name}: expected a nonnegative decimal with at most 4 fractional digits, received ${raw}`);
const bps = BigInt(m[1]) * MULT_SCALE + BigInt((m[2] ?? '').padEnd(4, '0') || '0');
if (bps > BigInt(maxWhole) * MULT_SCALE)
throw new ConfigError(`${name}: ${s} exceeds the maximum whole value ${maxWhole}`);
return bps;
}
const multToString = (bps) => {
const frac = (bps % MULT_SCALE).toString().padStart(4, '0').replace(/0+$/, '');
return `${bps / MULT_SCALE}${frac ? `.${frac}` : ''}`;
};
function envBool(env, name, def) {
const raw = env[name];
if (raw === undefined || raw === '') return def;
const s = String(raw).trim().toLowerCase();
if (['1', 'true', 'yes', 'on'].includes(s)) return true;
if (['0', 'false', 'no', 'off'].includes(s)) return false;
throw new ConfigError(`${name}: true/false, ${raw}`);
}
export function isLoopbackHost(h) {
if (h === '::1') return true;
const m = /^127\.(\d{1,3})\.(\d{1,3})\.(\d{1,3})$/.exec(h);
return !!m && m.slice(1).every((o) => Number(o) <= 255);
}
function isLoopbackPeer(a) {
if (typeof a !== 'string') return false;
return isLoopbackHost(a.startsWith('::ffff:') ? a.slice(7) : a);
}
function envUrl(env, name, def, { required }) {
const raw = (env[name] ?? '').trim() || def;
if (!raw) {
if (required) throw new ConfigError(`${name}: expected a nonnegative USD decimal with at most 6 fractional digits`);
return null;
}
let u;
try {
u = new URL(raw);
} catch {
throw new ConfigError(`${name}: URL`);
}
const host = u.hostname.replace(/^\[|\]$/g, '');
if (u.protocol !== 'https:' && !(u.protocol === 'http:' && isLoopbackHost(host))) {
throw new ConfigError(`${name}: https (http loopback)`);
}
return raw.replace(/\/+$/, '');
}
export function parseConfig(env) {
const botPubkey = pubkeyOrNull((env.PHX_SIGNER_BOT_PUBKEY ?? '').trim());
if (!botPubkey) throw new ConfigError('PHX_SIGNER_BOT_PUBKEY must be a canonical 32-byte base58 public key');
const authority = pubkeyOrNull((env.PHX_SIGNER_AUTHORITY ?? '').trim());
if (!authority) throw new ConfigError('PHX_SIGNER_AUTHORITY must be a canonical owner public key');
if (authority === botPubkey) {
throw new ConfigError('Owner and delegate must differ; use a delegated position authority');
}
const token = env.PHX_SIGNER_TOKEN ?? '';
if (token.length < 32) throw new ConfigError('PHX_SIGNER_TOKEN must contain at least 32 characters');
if (/\s/.test(token)) throw new ConfigError('PHX_SIGNER_TOKEN cannot contain whitespace');
const bind = (env.PHX_SIGNER_BIND ?? '127.0.0.1').trim() || '127.0.0.1';
if (!isLoopbackHost(bind)) throw new ConfigError(`PHX_SIGNER_BIND=${bind}: loopback-IP (127.x.x.x ::1)`);
const rpcUrl = envUrl(env, 'PHX_SIGNER_RPC_URL', '', { required: true });
const apiUrl = envUrl(env, 'PHX_SIGNER_API_URL', 'https://perp-api.phoenix.trade', { required: true });
const priorityMin = BigInt(envInt(env, 'PHX_SIGNER_PRIORITY_FEE_MIN', undefined, 0, Number.MAX_SAFE_INTEGER));
const priorityMax = BigInt(envInt(env, 'PHX_SIGNER_PRIORITY_FEE_MAX', undefined, 0, Number.MAX_SAFE_INTEGER));
if (priorityMin > priorityMax) throw new ConfigError('Priority minimum exceeds priority maximum');
const iocOffsetMax = envInt(env, 'PHX_SIGNER_IOC_SLOT_OFFSET_MAX', undefined, 1, 150);
const readonly = envBool(env, 'PHX_SIGNER_READONLY', true);
const journalPath = (env.PHX_SIGNER_JOURNAL ?? '').trim() || null;
if (!readonly && !journalPath) {
throw new ConfigError('Writes require a durable intent journal; set PHX_SIGNER_JOURNAL or use readonly mode');
}
const intentTtlMs = envInt(env, 'PHX_SIGNER_INTENT_TTL_MS', undefined, 60000, 24 * 3600000);
const tombstoneTtlMs = envInt(env, 'PHX_SIGNER_INTENT_TOMBSTONE_MS', undefined, 60000, 30 * 24 * 3600000);
if (tombstoneTtlMs < intentTtlMs)
throw new ConfigError('Intent tombstone retention must be at least intent retention');
const maxNotionalUsd = envInt(env, 'PHX_SIGNER_MAX_ORDER_NOTIONAL_USD', 0, 0, 1000000000);
const collateralMultBps = envMultBps(env, 'PHX_SIGNER_MAX_ORDER_COLLATERAL_MULT', 1000);
const collateralMaxAgeMs = envInt(env, 'PHX_SIGNER_COLLATERAL_MAX_AGE_MS', undefined, 5000, 3600000);
const collateralLazy = envBool(env, 'PHX_SIGNER_COLLATERAL_LAZY', true);
const collateralRefreshMs = envInt(env, 'PHX_SIGNER_COLLATERAL_REFRESH_MS', undefined, 1000, 600000);
if (collateralMultBps > 0n && collateralRefreshMs >= collateralMaxAgeMs) {
throw new ConfigError('Collateral refresh interval must be shorter than its allowed age');
}
const clientOrderNamespace = env.PHX_SIGNER_CLIENT_ORDER_NAMESPACE;
if (typeof clientOrderNamespace !== 'string' || !clientOrderNamespace.trim())
throw new ConfigError('PHX_SIGNER_CLIENT_ORDER_NAMESPACE required');
return {
clientOrderNamespace,
botPubkey,
authority,
traderPdaIndex: envInt(env, 'PHX_SIGNER_TRADER_PDA_INDEX', 0, 0, 255),
subaccountIndex: envInt(env, 'PHX_SIGNER_SUBACCOUNT_INDEX', 0, 0, 255),
rpcUrl,
apiUrl,
token,
bind,
port: envInt(env, 'PHX_SIGNER_PORT', undefined, 1, 65535),
readonly,
cu: {
order: envInt(env, 'PHX_SIGNER_CU_ORDER', undefined, 50000, MAX_CU),
cancel: envInt(env, 'PHX_SIGNER_CU_CANCEL', undefined, 50000, MAX_CU),
cancelPerId: envInt(env, 'PHX_SIGNER_CU_CANCEL_PER_ID', undefined, 0, 100000),
},
priority: {
min: priorityMin,
max: priorityMax,
percentile: envInt(env, 'PHX_SIGNER_PRIORITY_FEE_PERCENTILE', undefined, 0, 100),
},
iocSlotOffset: envInt(env, 'PHX_SIGNER_IOC_SLOT_OFFSET', undefined, 1, iocOffsetMax),
iocSlotOffsetMax: iocOffsetMax,
httpWaitMs: envInt(env, 'PHX_SIGNER_HTTP_WAIT_MS', undefined, 50, 60000),
intentTtlMs,
tombstoneTtlMs,
maxInflight: envInt(env, 'PHX_SIGNER_MAX_INFLIGHT', undefined, 1, 1024),
journalPath,
skipPreflight: envBool(env, 'PHX_SIGNER_SKIP_PREFLIGHT', false),
rpcTimeoutMs: envInt(env, 'PHX_SIGNER_RPC_TIMEOUT_MS', undefined, 100, 60000),
pollMs: envInt(env, 'PHX_SIGNER_CONFIRM_POLL_MS', undefined, 1, 10000),
rebroadcastMs: envInt(env, 'PHX_SIGNER_REBROADCAST_MS', undefined, 1, 30000),
confirmMaxMs: envInt(env, 'PHX_SIGNER_CONFIRM_MAX_MS', undefined, 100, 600000),
expiryMarginBlocks: envInt(env, 'PHX_SIGNER_EXPIRY_MARGIN_BLOCKS', undefined, 0, 300),
expiryRecheckMs: envInt(env, 'PHX_SIGNER_EXPIRY_RECHECK_MS', undefined, 0, 60000),
resolveMs: envInt(env, 'PHX_SIGNER_RESOLVE_MS', undefined, 10, 600000),
metaRefreshMs: envInt(env, 'PHX_SIGNER_META_REFRESH_MS', undefined, 1000, 24 * 3600000),
headerShareMs: envInt(env, 'PHX_SIGNER_HEADER_SHARE_MS', undefined, 0, 600000),
collateralLazy: envBool(env, 'PHX_SIGNER_COLLATERAL_LAZY', true),
headerHeartbeatMs: envInt(env, 'PHX_SIGNER_HEADER_HEARTBEAT_MS', undefined, 1000, 3600000),
collateralLazyMaxAgeMs: envInt(env, 'PHX_SIGNER_COLLATERAL_LAZY_MAX_AGE_MS', undefined, 0, 600000),
solRefreshMs: envInt(env, 'PHX_SIGNER_SOL_REFRESH_MS', undefined, 5000, 3600000),
solMaxAgeMs: envInt(env, 'PHX_SIGNER_SOL_MAX_AGE_MS', undefined, 10000, 24 * 3600000),
globalConfigTtlMs: envInt(env, 'PHX_SIGNER_GLOBAL_CONFIG_TTL_MS', undefined, 0, 24 * 3600000),
minSolLamports: BigInt(envInt(env, 'PHX_SIGNER_MIN_SOL_LAMPORTS', undefined, 0, Number.MAX_SAFE_INTEGER)),
maxOrderNotionalQuoteLots: BigInt(maxNotionalUsd) * QUOTE_LOTS_PER_USD,
maxOrderCollateralMultBps: collateralMultBps,
collateralMaxAgeMs,
collateralRefreshMs,
};
}
export function takeSecretFromEnv(env) {
const s = env.PHX_SIGNER_BOT_SECRET_KEY;
delete env.PHX_SIGNER_BOT_SECRET_KEY;
return typeof s === 'string' ? s.trim() : '';
}
export const CREDENTIAL_NAME = 'markpaper-phx-delegate';
export function takeBotSecret(env, { readFile = fs.readFileSync } = {}) {
const fromEnv = takeSecretFromEnv(env);
const dir = env.CREDENTIALS_DIRECTORY;
if (dir) {
let s = '';
try {
s = String(readFile(path.join(dir, CREDENTIAL_NAME), 'utf8')).trim();
} catch (e) {
if (e?.code !== 'ENOENT')
throw new ConfigError(`Cannot read credential ${CREDENTIAL_NAME}: ${e?.code ?? 'read failed'}`);
}
if (s) return { secret: s, source: 'credential', alsoInEnv: !!fromEnv };
}
return { secret: fromEnv, source: fromEnv ? 'env' : 'none', alsoInEnv: false };
}
export async function loadBotSigner(secretB58, expectedPubkey) {
if (!secretB58) throw new ConfigError('Delegate key material is missing');
let bytes;
try {
bytes = new Uint8Array(b58enc.encode(secretB58));
} catch {
throw new ConfigError('Delegate key material is not valid base58');
}
try {
if (bytes.length !== 64)
throw new ConfigError(`Delegate key requires 64 bytes (seed and public half); received ${bytes.length}`);
let signer;
try {
signer = await createKeyPairSignerFromBytes(bytes);
} catch {
throw new ConfigError('Delegate key public half does not match its private seed');
}
if (signer.address !== expectedPubkey) {
throw new ConfigError(
`Delegate public key ${signer.address} differs from PHX_SIGNER_BOT_PUBKEY ${expectedPubkey}`,
);
}
return signer;
} finally {
bytes.fill(0);
}
}
export const TRADER_DISCRIMINANT = sha256('account:trader').subarray(0, 8);
const HDR = {
AUTHORITY: 56,
COLLATERAL: 88,
FLAGS: 96,
MAX_POSITIONS: 112,
PREFS: 116,
POSITION_AUTHORITY: 120,
PDA_INDEX: 154,
SUB_INDEX: 155,
LEN: 224,
};
export function parseTraderHeader(data, owner) {
if (owner !== PROGRAM_ADDRESS)
return { ok: false, why: `Account owner ${owner ?? 'absent'} differs from Phoenix program ${PROGRAM_ADDRESS}` };
if (!(data instanceof Uint8Array) || data.length < HDR.LEN)
return { ok: false, why: `Trader header requires ${HDR.LEN} bytes; received ${data?.length ?? 0}` };
if (!Buffer.from(data.subarray(0, 8)).equals(TRADER_DISCRIMINANT))
return { ok: false, why: 'Invalid account:trader discriminator' };
const dv = new DataView(data.buffer, data.byteOffset, data.byteLength);
return {
ok: true,
authority: b58dec.decode(data.subarray(HDR.AUTHORITY, HDR.AUTHORITY + 32)),
positionAuthority: b58dec.decode(data.subarray(HDR.POSITION_AUTHORITY, HDR.POSITION_AUTHORITY + 32)),
collateralQuoteLots: dv.getBigInt64(HDR.COLLATERAL, true).toString(),
flags: dv.getUint32(HDR.FLAGS, true),
maxPositions: dv.getUint32(HDR.MAX_POSITIONS, true),
preferenceBits: dv.getUint32(HDR.PREFS, true),
pdaIndex: data[HDR.PDA_INDEX],
subaccountIndex: data[HDR.SUB_INDEX],
};
}
export function describePreferences(bits) {
const b = Number(bits) >>> 0;
return {
bits: b,
disableCollateralSweep: (b & PREF_DISABLE_COLLATERAL_SWEEP) !== 0,
disablePositionAuthoritySwap: (b & PREF_DISABLE_POSITION_AUTHORITY_SWAP) !== 0,
reservedBits: b & ~(PREF_DISABLE_COLLATERAL_SWEEP | PREF_DISABLE_POSITION_AUTHORITY_SWAP),
};
}
export const GLOBAL_CONFIG_DISCRIMINANT = sha256('account:global_configuration').subarray(0, 8);
const GC = {
ACCOUNT_KEY: 8,
PERP_ASSET_MAP: 360,
GTI_HEADER: 392,
ATB_HEADER: 424,
QUOTE_DECIMALS: 505,
SPOT_META: 776,
SPOT_FLAGS: 776 + 104,
LEN: 1104,
};
export function parseGlobalConfig(data, owner) {
if (owner !== PROGRAM_ADDRESS)
return {
ok: false,
why: `Global configuration owner ${owner ?? 'absent'} differs from Phoenix program ${PROGRAM_ADDRESS}`,
};
if (!(data instanceof Uint8Array) || data.length < GC.LEN)
return { ok: false, why: `Global configuration requires ${GC.LEN} bytes; received ${data?.length ?? 0}` };
if (!Buffer.from(data.subarray(0, 8)).equals(GLOBAL_CONFIG_DISCRIMINANT))
return { ok: false, why: 'Invalid account:global_configuration discriminator' };
const key = (o) => b58dec.decode(data.subarray(o, o + 32));
const spotFlags = data[GC.SPOT_FLAGS];
return {
ok: true,
accountKey: key(GC.ACCOUNT_KEY),
perpAssetMap: key(GC.PERP_ASSET_MAP),
globalTraderIndexHeader: key(GC.GTI_HEADER),
activeTraderBufferHeader: key(GC.ATB_HEADER),
quoteDecimals: data[GC.QUOTE_DECIMALS],
nativeSol: {
flags: spotFlags,
active: (spotFlags & SPOT_FLAG_IS_ACTIVE) !== 0,
positionAuthoritySwapDisabled: (spotFlags & SPOT_FLAG_DISABLE_POSITION_AUTHORITY_SWAP) !== 0,
},
};
}
export function checkViewKeysAgainstChain(keys, gc) {
const p = [];
if (!gc?.ok) return [gc?.why ?? 'Global configuration unavailable'];
if (gc.accountKey !== GLOBAL_CONFIGURATION)
p.push(`GlobalConfig.account_key ${gc.accountKey} ≠ ${GLOBAL_CONFIGURATION}`);
if (keys.perpAssetMap !== gc.perpAssetMap)
p.push(`REST perpAssetMap ${keys.perpAssetMap} differs from on-chain key ${gc.perpAssetMap}`);
if (keys.globalTraderIndex[0] !== gc.globalTraderIndexHeader)
p.push(
`REST globalTraderIndex[0] ${keys.globalTraderIndex[0]} differs from on-chain key ${gc.globalTraderIndexHeader}`,
);
if (keys.activeTraderBuffer[0] !== gc.activeTraderBufferHeader)
p.push(
`REST activeTraderBuffer[0] ${keys.activeTraderBuffer[0]} differs from on-chain key ${gc.activeTraderBufferHeader}`,
);
if (gc.quoteDecimals !== 6) p.push(`Quote decimals ${gc.quoteDecimals} differs from the supported protocol scale 6`);
return p;
}
export function keyRiskOf(prefs, gc) {
const traderOptOut = prefs ? prefs.disablePositionAuthoritySwap : null;
const exchangeKill = gc?.ok ? gc.nativeSol.positionAuthoritySwapDisabled : null;
const featureActive = gc?.ok ? gc.nativeSol.active : null;
let positionAuthoritySwap;
if (traderOptOut === true) positionAuthoritySwap = 'closed_by_trader';
else if (exchangeKill === true) positionAuthoritySwap = 'closed_by_exchange';
else if (traderOptOut === null || exchangeKill === null) positionAuthoritySwap = 'unknown';
else positionAuthoritySwap = featureActive ? 'open' : 'open_when_feature_active';
const warning =
positionAuthoritySwap === 'open' || positionAuthoritySwap === 'open_when_feature_active'
? 'position_authority_swap_enabled'
: null;
return {
positionAuthoritySwap,
warning,
traderOptOut,
exchangeKillSwitch: exchangeKill,
nativeSolCollateralActive: featureActive,
};
}
export function describeCapabilities(flags) {
const HOT = 1,
LIMIT = 2,
MARKET = 4,
RISK = 8,
DEPOSIT = 16,
WITHDRAW = 32;
const has = (m) => (flags & m) === m;
const hot = has(HOT);
const frozen = has(MARKET) && !has(RISK) && !has(WITHDRAW) && !has(DEPOSIT);
const reduceOnly = has(MARKET) && !has(RISK) && has(WITHDRAW) && has(DEPOSIT);
const state =
flags === 0
? 'uninitialized'
: frozen
? hot
? 'frozen_hot'
: 'frozen_cold'
: reduceOnly
? hot
? 'reduce_only_hot'
: 'reduce_only_cold'
: hot
? 'hot_active'
: 'cold';
return {
flags,
state,
hot,
reservedBits: flags & ~63,
placeLimitOrder: { immediate: hot && has(LIMIT), viaColdActivation: !hot && has(LIMIT) },
placeMarketOrder: { immediate: has(MARKET) },
riskIncreasingTrade: { immediate: has(MARKET) && has(RISK) },
riskReducingTrade: { immediate: has(MARKET) },
depositCollateral: { immediate: has(DEPOSIT) },
withdrawCollateral: { immediate: has(WITHDRAW) },
};
}
export function checkBinding(header, cfg) {
const problems = [];
if (!header?.ok) problems.push(header?.why ?? 'Trader header unavailable');
else {
if (header.authority !== cfg.authority)
problems.push(`Trader authority ${header.authority} differs from configured owner ${cfg.authority}`);
if (header.positionAuthority !== cfg.botPubkey) {
problems.push(
header.positionAuthority === header.authority
? 'Delegation absent: position authority is the owner'
: `Position authority ${header.positionAuthority} differs from configured delegate ${cfg.botPubkey}`,
);
}
if (header.pdaIndex !== cfg.traderPdaIndex || header.subaccountIndex !== cfg.subaccountIndex) {
problems.push(
`Trader PDA/subaccount indices ${header.pdaIndex}/${header.subaccountIndex} differ from configured indices ${cfg.traderPdaIndex}/${cfg.subaccountIndex}`,
);
}
const caps = describeCapabilities(header.flags);
if (caps.state === 'uninitialized') problems.push('Trader account is uninitialized');
}
return { ok: problems.length === 0, problems };
}
function reqAddr(v, name) {
const a = pubkeyOrNull(v);
if (!a) throw new Error(`${name}: expected a canonical 32-byte base58 public key`);
return a;
}
export function parseExchangeView(json) {
if (!json || typeof json !== 'object') throw new Error('Exchange view must be an object with a nonempty market list');
const k = json.keys;
if (!k || typeof k !== 'object') throw new Error('Exchange view keys missing');
const globalConfig = reqAddr(k.globalConfig, 'keys.globalConfig');
if (globalConfig !== GLOBAL_CONFIGURATION)
throw new Error(
`REST keys.globalConfig ${globalConfig} differs from the Phoenix mainnet configuration ${GLOBAL_CONFIGURATION}`,
);
const perpAssetMap = reqAddr(k.perpAssetMap, 'keys.perpAssetMap');
const arr = (v, name) => {
if (!Array.isArray(v) || v.length === 0 || v.length > 16)
throw new Error(`${name}: expected between 1 and 16 public keys`);
return v.map((x, i) => reqAddr(x, `${name}[${i}]`));
};
const keys = {
globalConfig,
perpAssetMap,
globalTraderIndex: arr(k.globalTraderIndex, 'keys.globalTraderIndex'),
activeTraderBuffer: arr(k.activeTraderBuffer, 'keys.activeTraderBuffer'),
};
if (!Array.isArray(json.markets) || json.markets.length === 0)
throw new Error('Exchange view must be an object with a nonempty market list');
const markets = new Map();
const skipped = [];
for (const m of json.markets) {
const symbol = m?.symbol;
if (typeof symbol !== 'string' || !/^[A-Za-z0-9._-]{1,32}$/.test(symbol)) {
skipped.push({ symbol: String(symbol).slice(0, 40), why: 'Invalid market symbol' });
continue;
}
if (markets.has(symbol)) throw new Error(`Duplicate symbol in exchange view: ${symbol}`);
try {
const tick = m.tickSize;
const tickOk =
(typeof tick === 'number' && Number.isSafeInteger(tick) && tick > 0) ||
(typeof tick === 'string' && /^[1-9]\d{0,19}$/.test(tick));
if (!tickOk) throw new Error(`tickSize ${tick}`);
const bld = m.baseLotsDecimals;
if (!Number.isInteger(bld) || bld < -12 || bld > 18) throw new Error(`baseLotsDecimals ${bld}`);
if (typeof m.isolatedOnly !== 'boolean') throw new Error('isolatedOnly');
if (typeof m.marketStatus !== 'string') throw new Error('marketStatus');
markets.set(symbol, {
symbol,
assetId: m.assetId,
status: m.marketStatus,
marketPubkey: reqAddr(m.marketPubkey, 'marketPubkey'),
splinePubkey: reqAddr(m.splinePubkey, 'splinePubkey'),
tickSize: String(tick),
baseLotsDecimals: bld,
isolatedOnly: m.isolatedOnly,
});
} catch (e) {
skipped.push({ symbol, why: e.message });
}
}
if (markets.size === 0) throw new Error('Exchange view must be an object with a nonempty market list');
return { keys, markets, skipped };
}
const INTENT_RE = /^[A-Za-z0-9][A-Za-z0-9._:-]{7,127}$/;
function strictKeys(b, allowed, what) {
if (!b || typeof b !== 'object' || Array.isArray(b)) throw bad(`${what}: expected a JSON object`);
for (const k of Object.keys(b)) if (!allowed.includes(k)) throw bad(`${what}: unsupported field ${k}`);
}
function reqIntent(b, required) {
if (b.intentId === undefined && !required) return null;
if (typeof b.intentId !== 'string' || !INTENT_RE.test(b.intentId))
throw bad('intentId required: 8..128 characters from [A-Za-z0-9._:-]');
return b.intentId;
}
function reqSymbol(b) {
if (typeof b.symbol !== 'string' || !/^[A-Za-z0-9._-]{1,32}$/.test(b.symbol))
throw bad('symbol must match a Phoenix market symbol');
return b.symbol;
}
export function normalizeOrderRequest(b, cfg, { requireIntent = true } = {}) {
strictKeys(
b,
['intentId', 'symbol', 'side', 'kind', 'priceInTicks', 'numBaseLots', 'reduceOnly', 'lastValidSlotOffset'],
'order',
);
const intentId = reqIntent(b, requireIntent);
const symbol = reqSymbol(b);
if (b.side !== 'bid' && b.side !== 'ask') throw bad("side: 'bid' | 'ask'");
if (!['limit', 'postOnly', 'ioc'].includes(b.kind)) throw bad("kind: 'limit' | 'postOnly' | 'ioc'");
if (typeof b.reduceOnly !== 'boolean') throw bad('reduceOnly must be explicitly true or false');
const priceInTicks = parseU64(b.priceInTicks, 'priceInTicks', { min: 1n });
const numBaseLots = parseU64(b.numBaseLots, 'numBaseLots', { min: 1n });
let lastValidSlotOffset = null;
if (b.kind === 'ioc') {
const o = b.lastValidSlotOffset ?? cfg.iocSlotOffset;
if (!Number.isInteger(o) || o < 1 || o > cfg.iocSlotOffsetMax)
throw bad(`lastValidSlotOffset must be an integer between 1 and ${cfg.iocSlotOffsetMax}`);
lastValidSlotOffset = o;
} else if (b.lastValidSlotOffset !== undefined && b.lastValidSlotOffset !== null) {
throw bad('lastValidSlotOffset is only valid for kind=ioc');
}
return {
op: 'order',
intentId,
symbol,
side: b.side,
kind: b.kind,
priceInTicks,
numBaseLots,
reduceOnly: b.reduceOnly,
lastValidSlotOffset,
};
}
export function normalizeCancelRequest(b, _cfg, { requireIntent = true } = {}) {
strictKeys(b, ['intentId', 'symbol', 'ids'], 'cancel');
const intentId = reqIntent(b, requireIntent);
const symbol = reqSymbol(b);
if (!Array.isArray(b.ids) || b.ids.length === 0 || b.ids.length > MAX_CANCEL_IDS)
throw bad(
`ids must contain between 1 and ${MAX_CANCEL_IDS} exact {priceInTicks, seq} identities; transaction wire size is limited to ${MAX_TX_BYTES} bytes`,
);
const seen = new Set();
const ids = b.ids.map((x, i) => {
strictKeys(x, ['priceInTicks', 'seq'], `ids[${i}]`);
const priceInTicks = parseU64(x.priceInTicks, `ids[${i}].priceInTicks`, { min: 1n });
const seq = parseU64(x.seq, `ids[${i}].seq`);
const key = `${priceInTicks}|${seq}`;
if (seen.has(key)) throw bad(`ids[${i}]: duplicate order identity ${key}`);
seen.add(key);
return { priceInTicks, seq };
});
return { op: 'cancel', intentId, symbol, ids };
}
export function normalizeCancelAllRequest(b, _cfg, { requireIntent = true } = {}) {
strictKeys(b, ['intentId', 'symbol'], 'cancelAll');
return { op: 'cancelAll', intentId: reqIntent(b, requireIntent), symbol: reqSymbol(b) };
}
const NORMALIZERS = {
order: normalizeOrderRequest,
cancel: normalizeCancelRequest,
cancelAll: normalizeCancelAllRequest,
};
export function fingerprintOf(req) {
const { intentId: _omit, ...rest } = req;
return sha256(toJson(rest)).toString('hex');
}
export function clientOrderIdFor(intentId, namespace) {
if (!intentId) return 0n;
if (typeof namespace !== 'string' || !namespace)
throw new TypeError('Caller-defined client order namespace required');
return BigInt('0x' + sha256(`${namespace}:${intentId}`).subarray(0, 16).toString('hex'));
}
export function buildOrderPacket(req, { clientOrderId = 0n, lastValidSlot = null } = {}) {
const side = req.side === 'bid' ? Side.Bid : Side.Ask;
const orderFlags = req.reduceOnly ? OrderFlags.ReduceOnly : OrderFlags.None;
if (req.kind === 'limit') {
return {
side,
priceInTicks: req.priceInTicks,
numBaseLots: req.numBaseLots,
selfTradeBehavior: SelfTradeBehavior.CancelProvide,
matchLimit: null,
clientOrderId,
lastValidSlot: null,
orderFlags,
cancelExisting: false,
};
}
if (req.kind === 'postOnly') {
return {
side,
priceInTicks: req.priceInTicks,
numBaseLots: req.numBaseLots,
clientOrderId,
slide: false,
lastValidSlot: null,
orderFlags,
cancelExisting: false,
};
}
if (req.kind === 'ioc') {
if (typeof lastValidSlot !== 'bigint') throw new Error('IoC requires an explicit lastValidSlot');
return {
side,
priceInTicks: req.priceInTicks,
numBaseLots: req.numBaseLots,
numQuoteLots: null,
minBaseLotsToFill: 0n,
minQuoteLotsToFill: 0n,
selfTradeBehavior: SelfTradeBehavior.CancelProvide,
matchLimit: null,
clientOrderId,
lastValidSlot,
orderFlags,
cancelExisting: false,
};
}
throw new Error(`Unsupported order kind: ${req.kind}`);
}
export function buildPhoenixInstruction(
req,
{ keys, market, signer, traderAccount, lastValidSlot = null, clientOrderNamespace },
) {
const common = {
programAddress: address(PROGRAM_ADDRESS),
logAuthorityAddress: address(LOG_AUTHORITY),
globalConfigurationAddress: address(keys.globalConfig),
traderAccount: address(traderAccount),
perpAssetMap: address(keys.perpAssetMap),
globalTraderIndex: keys.globalTraderIndex.map((a) => address(a)),
activeTraderBuffer: keys.activeTraderBuffer.map((a) => address(a)),
orderbook: address(market.marketPubkey),
splineCollection: address(market.splinePubkey),
};
if (req.op === 'order') {
const orderPacket = buildOrderPacket(req, {
clientOrderId: clientOrderIdFor(req.intentId, clientOrderNamespace),
lastValidSlot,
});
const p = { ...common, trader: address(signer), orderPacket };
if (req.kind === 'limit') return buildPlaceLimitOrderIx(p);
if (req.kind === 'postOnly') return buildPlacePostOnlyOrderIx(p);
return buildPlaceMarketOrderIx(p);
}
if (req.op === 'cancel') {
return buildCancelOrdersByIdIx({
...common,
traderWallet: address(signer),
orderIds: req.ids.map((x) => ({
nodePointer: 0,
orderId: { priceInTicks: x.priceInTicks, orderSequenceNumber: x.seq },
})),
});
}
if (req.op === 'cancelAll') return buildCancelAllIx({ ...common, traderWallet: address(signer) });
throw new Error(`Unsupported operation: ${req.op}`);
}
export function computeUnitLimitIx(units) {
const data = new Uint8Array(5);
data[0] = 2;
new DataView(data.buffer).setUint32(1, units, true);
return { programAddress: address(COMPUTE_BUDGET_PROGRAM), accounts: [], data };
}
export function computeUnitPriceIx(microLamports) {
const data = new Uint8Array(9);
data[0] = 3;
new DataView(data.buffer).setBigUint64(1, BigInt(microLamports), true);
return { programAddress: address(COMPUTE_BUDGET_PROGRAM), accounts: [], data };
}
export function cuLimitFor(req, cfg) {
if (req.op === 'order') return cfg.cu.order;
if (req.op === 'cancel') {
const extra = Math.max(0, (Array.isArray(req.ids) ? req.ids.length : 1) - 1) * (cfg.cu.cancelPerId ?? 0);
return Math.min(MAX_CU, cfg.cu.cancel + extra);
}
return cfg.cu.cancel;
}
export function wireBytesOf(msg) {
return Buffer.from(getBase64EncodedWireTransaction(compileTransaction(msg)), 'base64').length;
}
const SIZE_PROBE_BLOCKHASH = { blockhash: '11111111111111111111111111111111', lastValidBlockHeight: 0n };
export function writeTxBytes(req, { keys, market, signer, traderAccount, cfg }) {
const lastValidSlot = req.kind === 'ioc' ? 1n : null;
const ix = buildPhoenixInstruction(req, {
keys,
market,
signer,
traderAccount,
lastValidSlot,
clientOrderNamespace: cfg.clientOrderNamespace,
});
return wireBytesOf(
buildTransactionMessage({
feePayer: signer,
blockhash: SIZE_PROBE_BLOCKHASH,
cuLimit: cuLimitFor(req, cfg),
cuPrice: 1n,
ix,
}),
);
}
export function buildTransactionMessage({ feePayerSigner = null, feePayer = null, blockhash, cuLimit, cuPrice, ix }) {
const ixs = [computeUnitLimitIx(cuLimit)];
if (cuPrice > 0n) ixs.push(computeUnitPriceIx(cuPrice));
ixs.push(ix);
return pipe(
createTransactionMessage({ version: 0 }),
(m) =>
feePayerSigner
? setTransactionMessageFeePayerSigner(feePayerSigner, m)
: setTransactionMessageFeePayer(address(feePayer), m),
(m) =>
setTransactionMessageLifetimeUsingBlockhash(
{ blockhash: blockhash.blockhash, lastValidBlockHeight: BigInt(blockhash.lastValidBlockHeight) },
m,
),
(m) => appendTransactionMessageInstructions(ixs, m),
);
}
export function decodeCpiResponse(bytes, side = null) {
const b = bytes instanceof Uint8Array ? bytes : Uint8Array.from(bytes);
if (b.length !== 64) throw new Error(`MatchingEngineCPIResponse requires 64 bytes; received ${b.length}`);
const dv = new DataView(b.buffer, b.byteOffset, b.byteLength);
const f = (i) => dv.getBigUint64(i * 8, true);
const [priceInTicks, seq, quoteIn, baseIn, quoteOut, baseOut, quotePosted, basePosted] = [0, 1, 2, 3, 4, 5, 6, 7].map(
f,
);
const cpi = {
priceInTicks: priceInTicks.toString(),
seq: seq.toString(),
baseIn: baseIn.toString(),
baseOut: baseOut.toString(),
quoteIn: quoteIn.toString(),
quoteOut: quoteOut.toString(),
quotePosted: quotePosted.toString(),
basePosted: basePosted.toString(),
};
if (basePosted > 0n && priceInTicks === 0n)
throw new Error('returnData reports a posted order with an empty order identity');
if (side === 'bid' || side === 'ask') {
const bid = side === 'bid';
if (bid ? baseIn !== 0n || quoteOut !== 0n : baseOut !== 0n || quoteIn !== 0n)
throw new Error(`returnData has fills for the wrong side: ${side}`);
if (basePosted > 0n && seq >= BID_SEQ_MIN !== bid)
throw new Error(`Posted sequence ${seq} disagrees with order side ${side}`);
const filledBase = bid ? baseOut : baseIn;
const filledQuote = bid ? quoteIn : quoteOut;
if ((filledBase === 0n) !== (filledQuote === 0n)) throw new Error('returnData base and quote fills disagree');
cpi.filledBaseLots = (bid ? baseOut : baseIn).toString();
cpi.filledQuoteLots = (bid ? quoteIn : quoteOut).toString();
}
return cpi;
}
function cpiFromReturnData(rd, side) {
if (!rd) return { cpi: null, cpiError: 'Missing transaction returnData' };
if (rd.programId !== PROGRAM_ADDRESS)
return { cpi: null, cpiError: `returnData belongs to a different program: ${rd.programId}` };
const data = Array.isArray(rd.data) ? rd.data : [rd.data, 'base64'];
if (data[1] !== 'base64' || typeof data[0] !== 'string')
return { cpi: null, cpiError: 'returnData is not base64 encoded' };
try {
return { cpi: decodeCpiResponse(Buffer.from(data[0], 'base64'), side), cpiError: null };
} catch (e) {
return { cpi: null, cpiError: e.message };
}
}
const EVENT_REASONS = [
'TooManyLimitOrders',
'PostOnlyCross',
'InvalidOrderPacket',
'TiFInvalid',
'OutsideExecutionPriceBand',
];
const LOG_REASONS = [
[/PostOnly order would cross the book/, 'PostOnlyCross'],
[/Reduce-only order would increase exposure/, 'ReduceOnlyWouldIncreaseExposure'],
[/Invalid position side: None/, 'ReduceOnlyNoPosition'],
[new RegExp(`\\b(${EVENT_REASONS.join('|')})\\b`), null],
];
const NOISE_LOGS = [
/^Phoenix Eternal:/,
/^Attempting to /,
/^Successfully /,
/^Failed to perform MatchingEngine::place_order$/,
/^[1-9A-HJ-NP-Za-km-z]{32,44}$/,
];
export function parsePhoenixLogs(logs) {
const out = {
reason: null,
message: null,
activated: false,
roClamped: null,
coldNoop: false,
cancelling: null,
cancelNotFound: [],
};
const lines = (Array.isArray(logs) ? logs : [])
.filter((l) => typeof l === 'string' && l.startsWith('Program log: '))
.map((l) => l.slice(13));
for (const l of lines) {
if (!out.reason) {
for (const [re, name] of LOG_REASONS) {
const m = re.exec(l);
if (m) {
out.reason = name ?? m[1];
break;
}
}
}
if (l === 'Successfully activated selected trader') out.activated = true;
let m =
/Reduce-only order size is greater than trader position\. Original order size: BaseLots\((\d+)\), new order size: BaseLots\((\d+)\)/.exec(
l,
);
if (m) out.roClamped = { requestedLots: m[1], effectiveLots: m[2] };
if (/cancel requested for cold trader; operation is a no-op/.test(l)) out.coldNoop = true;
m = /^Cancelling (\d+) orders/.exec(l);
if (m) out.cancelling = m[1];
m =
/^Failed to cancel order FIFOOrderId \{ price_in_ticks: Ticks\((\d+)\), order_sequence_number: (\d+) \}: (\w+)/.exec(
l,
);
if (m) out.cancelNotFound.push({ priceInTicks: m[1], seq: m[2], reason: m[3] });
}
const meaningful = lines.filter((l) => !NOISE_LOGS.some((re) => re.test(l)));
out.message = meaningful.length ? meaningful[meaningful.length - 1].slice(0, 300) : null;
return out;
}
export function scanEventReasons(innerInstructions) {
if (!Array.isArray(innerInstructions)) return null;
for (const g of innerInstructions) {
for (const ii of g?.instructions ?? []) {
if (typeof ii?.data !== 'string') continue;
let bytes;
try {
bytes = Buffer.from(b58enc.encode(ii.data));
} catch {
continue;
}
for (const r of EVENT_REASONS) {
const i = bytes.indexOf(Buffer.from(r));
if (i >= 0 && (i + r.length === bytes.length || bytes[i + r.length] === 0)) return r;
}
}
}
return null;
}
function errName(err) {
if (!err) return null;
if (typeof err === 'string') return err;
const ie = err.InstructionError;
if (Array.isArray(ie)) {
const e = ie[1];
if (typeof e === 'string') return e;
if (e && typeof e === 'object' && 'Custom' in e) return `Custom:${e.Custom}`;
return JSON.stringify(e);
}
return Object.keys(err)[0] ?? JSON.stringify(err);
}
const tail = (logs, n = 14) => (Array.isArray(logs) ? logs.slice(-n).map((l) => String(l).slice(0, 300)) : []);
export function execFromTransaction(tx) {
const meta = tx?.meta;
if (!meta || typeof meta !== 'object') return null;
return {
slot: tx?.slot != null ? String(tx.slot) : null,
err: meta.err ?? null,
logs: meta.logMessages ?? [],
returnData: meta.returnData ?? null,
innerInstructions: meta.innerInstructions ?? [],
unitsConsumed: meta.computeUnitsConsumed != null ? String(meta.computeUnitsConsumed) : null,
};
}
export function execFromSimulation(sim) {
const v = sim?.value ?? sim ?? {};
return {
slot: sim?.context?.slot != null ? String(sim.context.slot) : null,
err: v.err ?? null,
logs: v.logs ?? [],
returnData: v.returnData ?? null,
innerInstructions: v.innerInstructions ?? [],
unitsConsumed: v.unitsConsumed != null ? String(v.unitsConsumed) : null,
};
}
function inferContext(exec) {
const l = (exec.logs ?? []).join('\n');
if (/Phoenix Eternal: Place Market Order/.test(l)) return { op: 'order', kind: 'ioc', side: null };
if (/Phoenix Eternal: PlaceLimitOrder/.test(l)) return { op: 'order', kind: 'limit', side: null };
if (/Phoenix Eternal: Cancel Orders By ID/.test(l)) return { op: 'cancel' };
if (/Phoenix Eternal: Cancel All/.test(l)) return { op: 'cancelAll' };
return { op: 'unknown' };
}
export function classifyExecution(exec, ctx = {}) {
const c = ctx.op ? ctx : inferContext(exec);
const info = parsePhoenixLogs(exec.logs);
const eventReason = scanEventReasons(exec.innerInstructions);
const base = {
slot: exec.slot ?? null,
err: exec.err ?? null,
cpi: null,
logsTail: tail(exec.logs),
unitsConsumed: exec.unitsConsumed ?? null,
};
if (info.activated && !exec.err) base.activated = true;
if (info.roClamped && !exec.err) base.roClamped = info.roClamped;
if (exec.err) {
const en = errName(exec.err);
const cuOut =
en === 'ComputationalBudgetExceeded' ||
(Array.isArray(exec.logs) && exec.logs.some((l) => /exceeded CUs meter/.test(String(l))));
return {
...base,
status: 'failed',
rejectReason: cuOut ? 'ComputationalBudgetExceeded' : (info.reason ?? eventReason ?? en),
rejectMessage: info.message,
};
}
if (c.op === 'order') {
const { cpi, cpiError } = cpiFromReturnData(exec.returnData, c.side ?? null);
const r = { ...base, status: 'confirmed', cpi };
if (cpiError) r.cpiError = cpiError;
if (cpi && cpi.basePosted === '0' && cpi.baseIn === '0' && cpi.baseOut === '0') {
r.rejectReason = info.reason ?? eventReason ?? (c.kind === 'ioc' ? 'NoFill' : 'NotPostedNotFilled');
if (info.message) r.rejectMessage = info.message;
}
return r;
}
if (c.op === 'cancel' || c.op === 'cancelAll') {
if (info.coldNoop) {
const notFound = Array.isArray(c.ids)
? c.ids.map((x) => ({ priceInTicks: String(x.priceInTicks), seq: String(x.seq), reason: 'ColdNoop' }))
: [];
return {
...base,
status: 'confirmed',
rejectReason: 'ColdNoop',
rejectMessage: info.message,
cancel: { effective: false, cancelling: '0', notFound, coldNoop: true },
};
}
return {
...base,
status: 'confirmed',
cancel: { effective: true, cancelling: info.cancelling, notFound: info.cancelNotFound, coldNoop: false },
};
}
return { ...base, status: 'confirmed' };
}
export class RpcError extends Error {
constructor(message, { code = null, data = null, httpStatus = null } = {}) {
super(message);
this.code = code;
this.data = data;
this.httpStatus = httpStatus;
}
}
export function createJsonRpc(url, { timeoutMs, fetchImpl = globalThis.fetch }) {
if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 1) throw new TypeError('Explicit RPC timeout required');
let id = 0;
async function call(method, params) {
let res;
try {
res = await fetchImpl(url, {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: toJson({ jsonrpc: '2.0', id: ++id, method, params }),
signal: AbortSignal.timeout(timeoutMs),
});
} catch (e) {
throw new RpcError(
`${method}: ${e?.name === 'TimeoutError' ? 'RPC request timed out' : 'RPC transport failed'} (${e?.cause?.code ?? e?.message ?? e})`,
);
}
const text = await res.text();
if (!res.ok) throw new RpcError(`${method}: HTTP ${res.status}`, { httpStatus: res.status });
let j;
try {
j = JSON.parse(text);
} catch {
throw new RpcError(`${method}: response is not valid JSON`);
}
if (j.error)
throw new RpcError(`${method}: ${j.error.message ?? 'RPC returned an error without a message'}`, {
code: j.error.code ?? null,
data: j.error.data ?? null,
});
return j.result;
}
const c = { commitment: 'confirmed' };
return {
call,
getSlot: () => call('getSlot', [c]),
getBlockHeight: () => call('getBlockHeight', [c]),
getLatestBlockhash: () => call('getLatestBlockhash', [c]).then((r) => r.value),
getBalance: (a) => call('getBalance', [a, c]),
getAccountInfo: (a, slice) =>
call('getAccountInfo', [a, { ...c, encoding: 'base64', ...(slice ? { dataSlice: slice } : {}) }]),
getRecentPrioritizationFees: (accounts) => call('getRecentPrioritizationFees', [accounts]),
sendTransaction: (wire, { skipPreflight }) =>
call('sendTransaction', [
wire,
{ encoding: 'base64', skipPreflight, preflightCommitment: 'confirmed', maxRetries: 0 },
]),
getSignatureStatus: (sig, history) =>
call('getSignatureStatuses', [[sig], { searchTransactionHistory: history }]).then((r) => r?.value?.[0] ?? null),
getTransaction: (sig) =>
call('getTransaction', [sig, { ...c, encoding: 'json', maxSupportedTransactionVersion: 0 }]),
simulateTransaction: (wire) =>
call('simulateTransaction', [
wire,
{ ...c, encoding: 'base64', sigVerify: false, replaceRecentBlockhash: true, innerInstructions: true },
]),
};
}
export function createRestClient(apiUrl, { timeoutMs, fetchImpl = globalThis.fetch }) {
if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 1) throw new TypeError('Explicit REST timeout required');
return {
async getExchange() {
const res = await fetchImpl(`${apiUrl}/v1/view/exchange`, {
headers: { accept: 'application/json' },
signal: AbortSignal.timeout(timeoutMs),
});
if (!res.ok) throw new Error(`/v1/view/exchange: HTTP ${res.status}`);
return res.json();
},
};
}
export function openJournal(file, { fsImpl = fs } = {}) {
fsImpl.mkdirSync(path.dirname(file), { recursive: true, mode: 0o700 });
let fd = fsImpl.openSync(file, 'a', 0o600);
let closed = false;
const ensureOpen = () => {
if (closed) throw new Error('Intent journal is closed');
if (fd === null) fd = fsImpl.openSync(file, 'a', 0o600);
return fd;
};
return {
path: file,
append(rec) {
const f = ensureOpen();
fsImpl.writeSync(f, toJson(rec) + '\n');
fsImpl.fsyncSync(f);
},
load() {
// A missing history proof must never be interpreted as an empty spent-intent set.
const text = fsImpl.readFileSync(file, 'utf8');
const out = [];
let lineNumber = 0;
for (const line of text.split('\n')) {
lineNumber++;
if (!line.trim()) continue;
try {
out.push(JSON.parse(line));
} catch {
throw new Error(`Intent journal is corrupt at line ${lineNumber}; startup refused`);
}
}
return out;
},
rewrite(records) {
const tmp = `${file}.tmp`;
const tfd = fsImpl.openSync(tmp, 'w', 0o600);
try {
fsImpl.writeSync(tfd, records.map((r) => toJson(r) + '\n').join(''));
fsImpl.fsyncSync(tfd);
} finally {
fsImpl.closeSync(tfd);
}
const old = fd;
fd = null;
if (old !== null) {
try {
fsImpl.closeSync(old);
} catch {}
}
try {
fsImpl.renameSync(tmp, file);
} finally {
try {
ensureOpen();
} catch {
fd = null;
}
}
try {
const dfd = fsImpl.openSync(path.dirname(file), 'r');
try {
fsImpl.fsyncSync(dfd);
} finally {
fsImpl.closeSync(dfd);
}
} catch {}
},
get isOpen() {
return fd !== null;
},
close() {
closed = true;
const old = fd;
fd = null;
if (old === null) return;
try {
fsImpl.closeSync(old);
} catch {}
},
};
}
function withTimeout(promise, ms) {
let t;
return Promise.race([
promise,
new Promise((r) => {
t = setTimeout(r, ms);
t.unref?.();
}),
]).finally(() => clearTimeout(t));
}
const isConfirmed = (st) => st && (st.confirmationStatus === 'confirmed' || st.confirmationStatus === 'finalized');
const monoDefault = () => globalThis.performance.now();
export function pinMarkets(pins, markets) {
const conflicts = [];
for (const [symbol, m] of markets) {
const cur = {
marketPubkey: m.marketPubkey,
splinePubkey: m.splinePubkey,
tickSize: m.tickSize,
baseLotsDecimals: m.baseLotsDecimals,
};
const was = pins.get(symbol);
if (!was) {
pins.set(symbol, cur);
continue;
}
for (const k of Object.keys(cur))
if (was[k] !== cur[k]) conflicts.push({ symbol, field: k, was: String(was[k]), now: String(cur[k]) });
}
return conflicts;
}
const absBig = (x) => (x < 0n ? -x : x);
const fmtUsd = (q) =>
`${q < 0n ? '-' : ''}$${absBig(q) / QUOTE_LOTS_PER_USD}.${(absBig(q) % QUOTE_LOTS_PER_USD).toString().padStart(6, '0').slice(0, 2)}`;
const usdString = (q) => {
const cents = (absBig(q) % QUOTE_LOTS_PER_USD) / 10000n;
return `${q < 0n ? '-' : ''}${absBig(q) / QUOTE_LOTS_PER_USD}${cents ? `.${cents.toString().padStart(2, '0')}` : ''}`;
};
export function notionalCapMode(cfg) {
const fixed = cfg.maxOrderNotionalQuoteLots > 0n;
const coll = (cfg.maxOrderCollateralMultBps ?? 0n) > 0n;
return fixed && coll ? 'min' : fixed ? 'fixed' : coll ? 'collateral' : 'off';
}
export function collateralCapQuoteLots(collateralQuoteLots, multBps) {
const c = BigInt(collateralQuoteLots);
return c <= 0n ? 0n : (c * multBps) / MULT_SCALE;
}
export function describeNotionalCap(cfg) {
const mode = notionalCapMode(cfg);
const fixed = cfg.maxOrderNotionalQuoteLots > 0n ? fmtUsd(cfg.maxOrderNotionalQuoteLots) : 'off';
const mult = multToString(cfg.maxOrderCollateralMultBps);
return 'mode=' + mode + ', fixed=' + fixed + ', collateral multiplier=' + mult;
}
export const exemptFromCollateralCap = (req) => req.op === 'order' && req.reduceOnly === true;
export function createSidecar({
config: cfg,
signer = null,
rpc,
api,
journal = null,
now = Date.now,
mono = monoDefault,
log = defaultLog,
}) {
if (!cfg.readonly) {
if (
!journal ||
journal.isOpen === false ||
!['load', 'append', 'rewrite', 'close'].every((method) => typeof journal[method] === 'function')
)
throw new ConfigError('Writable mode requires an open durable intent journal');
if (!signer || String(signer.address) !== cfg.botPubkey || typeof signer.signTransactions !== 'function')
throw new ConfigError('Writable mode requires a signer matching the configured delegate key');
}
const sanitize = makeSanitizer([cfg.rpcUrl, cfg.token]);
const store = new Map();
const bySig = new Map();
const used = new Map();
const usedBySig = new Map();
const feeCache = new Map();
const meta = { data: null, loadedAt: null, error: 'Market metadata not yet loaded', skipped: [], chainProblems: [] };
const pins = new Map();
const conflicts = new Map();
let traderPda = null;
let shuttingDown = false;
let metaTimer = null;
let pruneTimer = null;
let resolveTimer = null;
let resolving = false;
let healthCache = null;
let collateralTimer = null;
let solTimer = null;
const collat = { quoteLots: null, atMono: null, slot: null, error: 'Collateral not yet read', unknownLogged: false };
const hdrShare = {
header: null,
dataBase64: null,
slot: null,
owner: null,
exists: null,
lamports: null,
atMono: null,
error: 'Trader header not yet read',
};
const hdrAgeMs = () => (hdrShare.atMono === null ? null : mono() - hdrShare.atMono);
const hdrFresh = (maxAgeMs) => {
const age = hdrAgeMs();
return hdrShare.header && hdrShare.error === null && age !== null && age <= maxAgeMs ? hdrShare : null;
};
const solCache = { lamports: null, atMono: null, error: null };
async function refreshSol() {
try {
const b = await rpc.getBalance(cfg.botPubkey);
solCache.lamports = BigInt(b?.value ?? 0);
solCache.atMono = mono();
solCache.error = null;
return true;
} catch (e) {
solCache.error = sanitize(e?.message ?? e);
return false;
}
}
const say = (m) => log(sanitize(m));
const gcCache = { value: null, atMono: null };
async function readGlobalConfig() {
if (
cfg.globalConfigTtlMs > 0 &&
gcCache.value &&
gcCache.atMono !== null &&
mono() - gcCache.atMono <= cfg.globalConfigTtlMs
)
return gcCache.value;
const r = await rpc.getAccountInfo(GLOBAL_CONFIGURATION, { offset: 0, length: GC.LEN });
const v = r?.value;
if (!v) return { ok: false, why: 'Global configuration unavailable' };
const parsed = parseGlobalConfig(new Uint8Array(Buffer.from(v.data[0], 'base64')), v.owner);
if (parsed?.ok) {
gcCache.value = parsed;
gcCache.atMono = mono();
}
return parsed;
}
async function refreshMeta() {
try {
const parsed = parseExchangeView(await api.getExchange());
const problems = checkViewKeysAgainstChain(parsed.keys, await readGlobalConfig());
meta.chainProblems = problems;
if (problems.length) throw new Error(`Exchange view contradicts on-chain configuration: ${problems.join('; ')}`);
for (const c of pinMarkets(pins, parsed.markets)) {
if (!conflicts.has(c.symbol)) conflicts.set(c.symbol, []);
const list = conflicts.get(c.symbol);
if (list.some((x) => x.field === c.field && x.now === c.now)) continue;
list.push({ field: c.field, was: c.was, now: c.now });
say(
`[meta] Identity conflict for ${c.symbol}.${c.field}: ${c.was} changed to ${c.now}; writes for this symbol are blocked`,
);
}
meta.data = parsed;
meta.loadedAt = now();
meta.error = null;
if (parsed.skipped.length && toJson(parsed.skipped) !== toJson(meta.skipped))
say(`[meta] Skipped invalid markets: ${toJson(parsed.skipped)}`);
meta.skipped = parsed.skipped;
} catch (e) {
meta.error = sanitize(e?.message ?? e);
say(`[meta] Refresh failed (${meta.data ? 'retaining previous cache' : 'cache empty'}): ${meta.error}`);
}
}
function resolveMarket(symbol) {
if (!meta.data) throw new HttpError(503, 'meta_unavailable', `Market metadata unavailable: ${meta.error}`);
const market = meta.data.markets.get(symbol);
if (!market) throw new HttpError(400, 'unknown_symbol', `Symbol ${symbol} is absent from the exchange view`);
const c = conflicts.get(symbol);
if (c)
throw new HttpError(
503,
'meta_conflict',
`Market identity changed for ${symbol}: ${c.map((x) => `${x.field} ${x.was}→${x.now}`).join(', ')}`,
);
return { keys: meta.data.keys, market: { ...market, ...pins.get(symbol) } };
}
const multOn = () => cfg.maxOrderCollateralMultBps > 0n;
const multLabel = () => `${multToString(cfg.maxOrderCollateralMultBps)} `;
function noteCollateral(header, slot) {
if (!header?.ok) return;
const s = slot != null ? BigInt(slot) : null;
if (s !== null && collat.slot !== null && s < collat.slot) return;
const v = BigInt(header.collateralQuoteLots);
const had = collat.quoteLots !== null && collat.atMono !== null && mono() - collat.atMono <= cfg.collateralMaxAgeMs;
collat.quoteLots = v;
collat.atMono = mono();
if (s !== null) collat.slot = s;
collat.error = null;
if (multOn() && collat.unknownLogged && !had) {
say(
`[cap] Collateral ${fmtUsd(v)}; order cap ${fmtUsd(collateralCapQuoteLots(v, cfg.maxOrderCollateralMultBps))} (${multLabel()})`,
);
}
collat.unknownLogged = false;
}
let headerInFlight = null;
function readHeaderFromChain() {
if (headerInFlight) return headerInFlight;
headerInFlight = readHeaderOnce().finally(() => {
headerInFlight = null;
});
return headerInFlight;
}
async function readHeaderOnce() {
if (!traderPda) {
hdrShare.error = 'Trader PDA has not been derived';
return false;
}
try {
const r = await rpc.getAccountInfo(traderPda, { offset: 0, length: HDR.LEN });
const v = r?.value;
const slot = r?.context?.slot ?? null;
const header = v
? parseTraderHeader(new Uint8Array(Buffer.from(v.data[0], 'base64')), v.owner)
: { ok: false, why: `Trader account ${traderPda} does not exist` };
if (!header.ok) {
hdrShare.error = header.why;
collat.error = header.why;
return false;
}
if (slot != null && hdrShare.slot != null && BigInt(slot) < BigInt(hdrShare.slot)) return true;
hdrShare.header = header;
hdrShare.dataBase64 = v.data[0];
hdrShare.slot = slot;
hdrShare.owner = v.owner ?? null;
hdrShare.exists = true;
hdrShare.lamports = String(v.lamports);
hdrShare.atMono = mono();
hdrShare.error = null;
noteCollateral(header, slot);
return true;
} catch (e) {
hdrShare.error = sanitize(e?.message ?? e);
collat.error = hdrShare.error;
return false;
}
}
async function headerWithin(maxAgeMs) {
const hit = hdrFresh(maxAgeMs);
if (hit) return hit;
await readHeaderFromChain();
return hdrFresh(maxAgeMs);
}
async function refreshCollateral() {
return readHeaderFromChain();
}
function needsLazyCollateral(req) {
if (!cfg.collateralLazy || !multOn() || exemptFromCollateralCap(req)) return false;
if (req.op !== 'order') return false;
return collat.atMono === null || mono() - collat.atMono > cfg.collateralLazyMaxAgeMs;
}
function freshCollateral() {
if (collat.quoteLots === null || collat.atMono === null) return null;
const ageMs = mono() - collat.atMono;
return ageMs <= cfg.collateralMaxAgeMs ? { quoteLots: collat.quoteLots, ageMs } : null;
}
function checkNotional(req, market) {
if (req.op !== 'order') return;
const q = req.priceInTicks * BigInt(market.tickSize) * req.numBaseLots;
const fixed = cfg.maxOrderNotionalQuoteLots;
if (fixed > 0n && q > fixed) {
throw new HttpError(
400,
'notional_cap',
`Order notional ${fmtUsd(q)} exceeds the configured fixed cap ${fmtUsd(fixed)} (PHX_SIGNER_MAX_ORDER_NOTIONAL_USD)`,
);
}
if (!multOn() || exemptFromCollateralCap(req)) return;
const fc = freshCollateral();
if (!fc) {
const ago =
collat.atMono === null
? 'no successful read'
: `last read ${Math.round((mono() - collat.atMono) / 1000)} seconds ago`;
const why = `Collateral is unavailable or stale (${ago}${collat.error ? `; ${collat.error}` : ''}; permitted age ${Math.round(cfg.collateralMaxAgeMs / 1000)} seconds)`;
if (!collat.unknownLogged) {
collat.unknownLogged = true;
say(
`[cap] ${why}; risk-increasing orders are blocked under multiplier ${multLabel()}; reduce-only orders remain eligible`,
);
}
throw new HttpError(
400,
'notional_cap_unknown',
`${why}; collateral multiplier ${multLabel()} blocks this risk-increasing order; reduce-only orders are exempt`,
);
}
const cap = collateralCapQuoteLots(fc.quoteLots, cfg.maxOrderCollateralMultBps);
if (q > cap) {
throw new HttpError(
400,
'notional_cap',
`Order notional ${fmtUsd(q)} exceeds collateral cap ${fmtUsd(cap)} (multiplier ${multLabel()}, collateral ${fmtUsd(fc.quoteLots)})`,
);
}
}
function capView() {
const mode = notionalCapMode(cfg);
const fixed = cfg.maxOrderNotionalQuoteLots;
const fc = multOn() ? freshCollateral() : null;
const unknown = multOn() && !fc;
const parts = [];
if (fixed > 0n) parts.push(fixed);
if (fc) parts.push(collateralCapQuoteLots(fc.quoteLots, cfg.maxOrderCollateralMultBps));
const eff = unknown || parts.length === 0 ? null : parts.reduce((a, b) => (b < a ? b : a));
return {
maxOrderNotionalUsd: eff === null ? null : usdString(eff),
maxOrderNotionalMode: mode,
maxOrderCollateralMult: multOn() ? Number(cfg.maxOrderCollateralMultBps) / Number(MULT_SCALE) : null,
maxOrderNotionalFixedUsd: fixed > 0n ? usdString(fixed) : null,
maxOrderNotionalUnknown: unknown,
maxOrderCollateralAgeMs:
multOn() && collat.atMono !== null ? Math.max(0, Math.round(mono() - collat.atMono)) : null,
};
}
async function priorityFee(accounts) {
const key = accounts.join(',');
const c = feeCache.get(key);
if (c && now() - c.at < 10000) return c.value;
let v = c?.value ?? cfg.priority.min;
try {
const arr = await rpc.getRecentPrioritizationFees(accounts);
const vals = (Array.isArray(arr) ? arr : [])
.map((x) => BigInt(x.prioritizationFee ?? 0))
.sort((a, b) => (a < b ? -1 : a > b ? 1 : 0));
if (vals.length) v = vals[Math.min(vals.length - 1, Math.floor((cfg.priority.percentile / 100) * vals.length))];
} catch (e) {
say(`[fee] Fee estimate unavailable: ${e.message}; using configured minimum ${v}`);
}
if (v < cfg.priority.min) v = cfg.priority.min;
if (v > cfg.priority.max) v = cfg.priority.max;
feeCache.set(key, { at: now(), value: v });
return v;
}
function snapshot(e) {
const common = {
intentId: e.intentId,
op: e.op,
symbol: e.symbol,
signature: e.signature ?? null,
lastValidBlockHeight: e.lastValidBlockHeight ?? null,
};
if (e.lastValidSlot) common.lastValidSlot = e.lastValidSlot;
if (e.result) return { ...common, ...e.result, signature: e.signature ?? null };
return { ...common, status: 'unknown', slot: null, err: null, cpi: null, pending: true };
}
function answer(e) {
if (!e.signature && !e.result) e.answeredUnsigned = true;
return snapshot(e);
}
function journalRec(e, t) {
const r = {
v: 1,
t,
intentId: e.intentId,
fp: e.fp,
op: e.op,
symbol: e.symbol,
side: e.side ?? null,
kind: e.kind ?? null,
sig: e.signature ?? null,
lvbh: e.lastValidBlockHeight ?? null,
lvs: e.lastValidSlot ?? null,
};
if (e.ids) r.ids = e.ids;
return r;
}
function finalize(e, result) {
if (e.result) return false;
e.result = result;
e.state = 'final';
e.finalAt = now();
e.wire = null;
if (journal) {
try {
journal.append({ ...journalRec(e, 'final'), result, ts: e.finalAt });
} catch (err) {
say(
`[journal] Final result append failed for ${e.intentId}: ${err.message}; the prior signed record remains authoritative`,
);
}
}
const c = result.cpi;
const extra = c
? ` filled=${c.filledBaseLots ?? '?'} posted=${c.basePosted}${c.basePosted !== '0' ? ` id=${c.priceInTicks}/${c.seq}` : ''}`
: '';
const nf = result.cancel?.notFound?.length ? ` notFound=${result.cancel.notFound.length}` : '';
const why = result.rejectReason
? ` (${result.rejectReason}${result.notSent && result.rejectMessage ? `: ${result.rejectMessage}` : ''})`
: '';
say(
`← ${e.intentId} ${result.status}${why}${extra}${nf}${e.signature ? ` sig=${e.signature}` : ''}${result.slot ? ` slot=${result.slot}` : ''}`,
);
return true;
}
const notSent = (why) => ({
status: 'failed',
slot: null,
err: { notSent: why },
cpi: null,
rejectReason: 'NotSent',
rejectMessage: why,
notSent: true,
});
const expired = () => ({
status: 'failed',
slot: null,
err: { expired: 'blockhash expired, transaction not found' },
cpi: null,
rejectReason: 'Expired',
expired: true,
});
async function sendWire(wire, skipPreflight) {
try {
await rpc.sendTransaction(wire, { skipPreflight });
return { ok: true };
} catch (e) {
if (!skipPreflight && e?.code === -32002) {
const d = e.data ?? {};
if (d.err === 'AlreadyProcessed') return { ok: true };
return {
preflightRejected: true,
err: d.err ?? { preflight: e.message },
logs: d.logs ?? [],
unitsConsumed: d.unitsConsumed ?? null,
};
}
return { ok: false, error: sanitize(e?.message ?? e) };
}
}
async function fetchAndClassify(e) {
let tx;
try {
tx = await rpc.getTransaction(e.signature);
} catch {
return null;
}
const ex = tx ? execFromTransaction(tx) : null;
if (!ex) return null;
return classifyExecution(ex, e);
}
function probe(e, { history }) {
if (e.result) return Promise.resolve({ final: true });
if (!e.probing)
e.probing = probeOnce(e, { history }).finally(() => {
e.probing = null;
});
return e.probing;
}
async function probeOnce(e, { history }) {
const found = async () => {
const r = await fetchAndClassify(e);
if (r) {
finalize(e, r);
return { final: true };
}
return { pending: true, found: true };
};
let st;
try {
st = await rpc.getSignatureStatus(e.signature, history);
e.lastRpcError = null;
} catch (err) {
e.lastRpcError = sanitize(err?.message ?? err);
return { pending: true };
}
e.lastSeen = st ? (st.confirmationStatus ?? 'processed') : null;
if (isConfirmed(st)) return found();
if (st) {
e.expiryMissAt = null;
return { pending: true, found: true };
}
let h;
try {
h = BigInt(await rpc.getBlockHeight());
} catch (err) {
e.lastRpcError = sanitize(err?.message ?? err);
return { pending: true, notSeen: true, blockhashAlive: null };
}
const lvbh = e.lastValidBlockHeight != null ? BigInt(e.lastValidBlockHeight) : null;
if (lvbh === null) return { pending: true, notSeen: true, blockhashAlive: null };
if (h <= lvbh + BigInt(cfg.expiryMarginBlocks)) {
e.expiryMissAt = null;
return { pending: true, notSeen: true, blockhashAlive: h <= lvbh };
}
let st2 = st;
if (!history) {
try {
st2 = await rpc.getSignatureStatus(e.signature, true);
} catch (err) {
e.lastRpcError = sanitize(err?.message ?? err);
return { pending: true, notSeen: true, blockhashAlive: false };
}
}
if (isConfirmed(st2)) return found();
if (st2) {
e.expiryMissAt = null;
return { pending: true, found: true };
}
const t = mono();
if (e.expiryMissAt == null) {
e.expiryMissAt = t;
return { pending: true, notSeen: true, blockhashAlive: false };
}
if (t - e.expiryMissAt < cfg.expiryRecheckMs) return { pending: true, notSeen: true, blockhashAlive: false };
finalize(e, expired());
return { final: true };
}
async function confirmLoop(e, { rebroadcast }) {
const deadline = mono() + cfg.confirmMaxMs;
let lastSend = mono();
for (;;) {
await sleep(cfg.pollMs);
if (e.result) return;
const p = await probe(e, { history: !rebroadcast });
if (e.result) return;
if (rebroadcast && e.wire && p.notSeen && p.blockhashAlive !== false && mono() - lastSend >= cfg.rebroadcastMs) {
lastSend = mono();
const r = await sendWire(e.wire, true);
if (!r.ok) say(`[send] ${e.signature}: ${r.error}`);
}
if (mono() > deadline) {
e.state = 'stale';
e.wire = null;
say(
`[confirm] ${e.intentId} sig=${e.signature}: confirmation wait ${cfg.confirmMaxMs}ms ended; background resolution interval ${cfg.resolveMs}ms; query GET /tx/:sig for evidence`,
);
return;
}
}
}
async function runIntent(e, req, keys, market) {
try {
let lastValidSlot = null;
if (req.kind === 'ioc') {
lastValidSlot = BigInt(await rpc.getSlot()) + BigInt(req.lastValidSlotOffset);
e.lastValidSlot = lastValidSlot.toString();
}
const ix = buildPhoenixInstruction(req, {
keys,
market,
signer: signer.address,
traderAccount: traderPda,
lastValidSlot,
clientOrderNamespace: cfg.clientOrderNamespace,
});
const cuLimit = cuLimitFor(req, cfg);
const cuPrice = await priorityFee([traderPda, market.marketPubkey]);
const bh = await rpc.getLatestBlockhash();
if (!bh?.blockhash || bh.lastValidBlockHeight == null)
throw new Error('getLatestBlockhash returned an empty value');
const signDeadlineMs = Math.floor(cfg.httpWaitMs / 2);
if (e.answeredUnsigned) throw new Error('Unknown-without-signature response already sent; signing is forbidden');
if (mono() - e.createdMono > signDeadlineMs)
throw new Error(`Signing deadline of ${signDeadlineMs}ms elapsed before signing started`);
const msg = buildTransactionMessage({ feePayerSigner: signer, blockhash: bh, cuLimit, cuPrice, ix });
const signed = await signTransactionMessageWithSigners(msg);
const wireLen = Buffer.from(getBase64EncodedWireTransaction(signed), 'base64').length;
if (wireLen > MAX_TX_BYTES)
throw new Error(`tx_too_large: wire transaction has ${wireLen} bytes; protocol maximum is ${MAX_TX_BYTES}`);
if (e.answeredUnsigned)
throw new Error(
'Unknown-without-signature response sent while signing; signed transaction discarded before broadcast',
);
e.signature = getSignatureFromTransaction(signed);
e.wire = getBase64EncodedWireTransaction(signed);
e.lastValidBlockHeight = String(bh.lastValidBlockHeight);
if (journal) {
try {
journal.append({ ...journalRec(e, 'signed'), ts: e.createdAt });
} catch (err) {
const sig = e.signature;
e.signature = null;
e.wire = null;
finalize(
e,
notSent(
`Signed transaction was not broadcast because journal append failed: ${err.message}; discarded signature ${sig}`,
),
);
return;
}
}
bySig.set(e.signature, e);
say(
`→ ${e.intentId} ${describe(req)} cu=${cuLimit} price=${cuPrice}µ lvbh=${e.lastValidBlockHeight}${e.lastValidSlot ? ` lvs=${e.lastValidSlot}` : ''} sig=${e.signature}`,
);
const first = await sendWire(e.wire, cfg.skipPreflight);
if (first.preflightRejected) {
const r = classifyExecution(
{
slot: null,
err: first.err,
logs: first.logs,
returnData: null,
innerInstructions: [],
unitsConsumed: first.unitsConsumed != null ? String(first.unitsConsumed) : null,
},
e,
);
finalize(e, { ...r, status: 'failed', preflight: true });
return;
}
if (!first.ok) say(`[send] ${e.signature}: ${first.error}; transaction outcome remains unresolved`);
await confirmLoop(e, { rebroadcast: true });
} catch (err) {
if (!e.signature) finalize(e, notSent(sanitize(err?.message ?? err)));
else {
e.state = 'stale';
e.wire = null;
say(
`[intent] ${e.intentId} sig=${e.signature}: ${sanitize(err?.message ?? err)}; retained for background resolution`,
);
}
} finally {
if (e.result) e.wire = null;
}
}
function describe(req) {
if (req.op === 'order')
return `${req.symbol} ${req.side} ${req.kind} px=${req.priceInTicks} lots=${req.numBaseLots} ro=${req.reduceOnly}${req.kind === 'ioc' ? ` +${req.lastValidSlotOffset}sl` : ''}`;
if (req.op === 'cancel')
return `${req.symbol} cancel ${req.ids.map((x) => `${x.priceInTicks}/${x.seq}`).join(',')}`;
return `${req.symbol} cancelAll`;
}
const inflight = () => [...store.values()].filter((e) => e.state !== 'final' && e.state !== 'stale').length;
const staleCount = () => [...store.values()].filter((e) => e.state === 'stale' && !e.result).length;
function assertTxFits(req, keys, market, signerAddress) {
if (req.op !== 'cancel') return;
const bytes = writeTxBytes(req, { keys, market, signer: signerAddress, traderAccount: traderPda, cfg });
if (bytes > MAX_TX_BYTES) {
throw new HttpError(
400,
'tx_too_large',
`Cancel batch of ${req.ids.length} IDs has ${bytes} wire bytes; protocol maximum is ${MAX_TX_BYTES}; split the batch`,
);
}
}
async function submit(op, body) {
if (cfg.readonly) throw new HttpError(403, 'readonly', 'PHX_SIGNER_READONLY=true: writes disabled');
if (!signer) throw new HttpError(503, 'no_signer', 'Delegate key not loaded; simulation only');
if (shuttingDown) throw new HttpError(503, 'shutting_down', 'Signer is shutting down');
if (!traderPda) throw new HttpError(503, 'not_ready', 'Signer is not ready');
const req = NORMALIZERS[op](body, cfg);
const fp = fingerprintOf(req);
if (needsLazyCollateral(req)) await readHeaderFromChain().catch(() => {});
const existing = store.get(req.intentId);
if (existing) {
if (existing.fp !== fp)
throw new HttpError(
409,
'intent_conflict',
`intentId ${req.intentId} was already used with a different request`,
);
if (existing.done) await withTimeout(existing.done, cfg.httpWaitMs);
if (!existing.result && existing.state === 'stale' && existing.signature) {
await withTimeout(
probe(existing, { history: true }).catch(() => null),
cfg.httpWaitMs,
);
}
return answer(existing);
}
const tomb = used.get(req.intentId);
if (tomb) {
if (tomb.fp !== fp)
throw new HttpError(
409,
'intent_conflict',
`intentId ${req.intentId} was already used with a different request`,
);
throw new HttpError(
410,
'intent_forgotten',
`intentId ${req.intentId} is retained as a spent-intent tombstone (${tomb.status}${tomb.rejectReason ? `/${tomb.rejectReason}` : ''}); full result retention ended; query GET /tx/:sig`,
{ signature: tomb.sig, status: tomb.status },
);
}
const { keys, market } = resolveMarket(req.symbol);
checkNotional(req, market);
assertTxFits(req, keys, market, signer.address);
if (inflight() >= cfg.maxInflight) throw new HttpError(503, 'busy', 'Configured maximum inflight intents reached');
const e = {
intentId: req.intentId,
fp,
op: req.op,
symbol: req.symbol,
side: req.side,
kind: req.kind,
ids:
req.op === 'cancel'
? req.ids.map((x) => ({ priceInTicks: String(x.priceInTicks), seq: String(x.seq) }))
: undefined,
createdAt: now(),
createdMono: mono(),
state: 'building',
signature: null,
wire: null,
result: null,
};
store.set(req.intentId, e);
e.done = runIntent(e, req, keys, market);
await withTimeout(e.done, cfg.httpWaitMs);
return answer(e);
}
async function simulate(op, body) {
if (!traderPda) throw new HttpError(503, 'not_ready', 'Signer is not ready');
const req = NORMALIZERS[op](body, cfg, { requireIntent: false });
if (needsLazyCollateral(req)) await readHeaderFromChain().catch(() => {});
const { keys, market } = resolveMarket(req.symbol);
checkNotional(req, market);
let lastValidSlot = null;
if (req.kind === 'ioc') lastValidSlot = BigInt(await rpc.getSlot()) + BigInt(req.lastValidSlotOffset);
assertTxFits(req, keys, market, cfg.botPubkey);
const ix = buildPhoenixInstruction(req, {
keys,
market,
signer: cfg.botPubkey,
traderAccount: traderPda,
lastValidSlot,
clientOrderNamespace: cfg.clientOrderNamespace,
});
const cuLimit = cuLimitFor(req, cfg);
const cuPrice = await priorityFee([traderPda, market.marketPubkey]);
const bh = await rpc.getLatestBlockhash();
const msg = buildTransactionMessage({ feePayer: cfg.botPubkey, blockhash: bh, cuLimit, cuPrice, ix });
const wire = getBase64EncodedWireTransaction(compileTransaction(msg));
const sim = await rpc.simulateTransaction(wire);
const r = classifyExecution(execFromSimulation(sim), req);
return {
simulated: true,
sent: false,
op: req.op,
symbol: req.symbol,
...r,
logs: execFromSimulation(sim).logs,
cuLimit,
cuPrice: cuPrice.toString(),
lastValidSlot: lastValidSlot?.toString() ?? null,
wireBytes: Buffer.from(wire, 'base64').length,
};
}
async function txStatus(sig) {
if (!isSignatureString(sig)) throw bad('Expected a canonical 64-byte base58 Solana signature');
const e = bySig.get(sig);
if (e) {
if (!e.result) {
try {
await probe(e, { history: true });
} catch {}
}
if (e.result) return snapshot(e);
return {
...snapshot(e),
status: 'unknown',
found: !!e.lastSeen,
confirmationStatus: e.lastSeen ?? null,
...(e.lastRpcError ? { rpcError: e.lastRpcError } : {}),
};
}
const tombId = usedBySig.get(sig);
const tomb = tombId ? used.get(tombId) : null;
const ctx = tomb ? { op: tomb.op, side: tomb.side, kind: tomb.kind } : {};
const who = tomb ? { ours: true, intentId: tombId, forgotten: true } : { ours: false };
let st;
try {
st = await rpc.getSignatureStatus(sig, true);
} catch (err) {
return { signature: sig, status: 'unknown', found: null, rpcError: sanitize(err.message), ...who };
}
if (isConfirmed(st)) {
let tx = null;
try {
tx = await rpc.getTransaction(sig);
} catch {
tx = null;
}
const ex = tx ? execFromTransaction(tx) : null;
if (ex) return { signature: sig, ...classifyExecution(ex, ctx), found: true, ...who };
}
return {
signature: sig,
slot: null,
err: null,
cpi: null,
status: 'unknown',
found: !!st,
confirmationStatus: st?.confirmationStatus ?? null,
...who,
};
}
async function account(pda, offset, length) {
if (!pubkeyOrNull(pda)) throw bad('Expected a canonical 32-byte base58 public key');
if (cfg.headerShareMs > 0 && traderPda && pda === traderPda && offset === 0 && length === HDR.LEN) {
let hit = hdrFresh(cfg.headerShareMs);
if (!hit) {
await readHeaderFromChain();
hit = hdrFresh(cfg.headerShareMs);
}
if (hit) {
return {
address: pda,
slot: hit.slot != null ? String(hit.slot) : null,
exists: !!hit.exists,
owner: hit.owner,
lamports: hit.lamports,
offset,
length,
dataBase64: hit.dataBase64,
fromCache: true,
ageMs: hdrAgeMs(),
};
}
}
const r = await rpc.getAccountInfo(pda, { offset, length });
const v = r?.value ?? null;
return {
address: pda,
slot: r?.context?.slot != null ? String(r.context.slot) : null,
exists: !!v,
owner: v?.owner ?? null,
lamports: v ? String(v.lamports) : null,
offset,
length,
dataBase64: v ? v.data[0] : null,
};
}
async function health() {
if (healthCache && now() - healthCache.at < 3000) return healthCache.value;
const out = {
version: VERSION,
pubkey: cfg.botPubkey,
authority: cfg.authority,
traderPda,
traderPdaIndex: cfg.traderPdaIndex,
subaccountIndex: cfg.subaccountIndex,
program: PROGRAM_ADDRESS,
readonly: cfg.readonly,
signerLoaded: !!signer,
positionAuthority: null,
bindingOk: false,
bindingProblems: [],
capabilities: null,
preferences: null,
keyRisk: null,
collateralQuoteLots: null,
solLamports: null,
solLow: null,
slot: null,
rpcOk: false,
rpcError: null,
meta: {
ok: !!meta.data,
markets: meta.data?.markets.size ?? 0,
ageMs: meta.loadedAt ? now() - meta.loadedAt : null,
error: meta.error,
skipped: meta.skipped.length,
chainProblems: meta.chainProblems,
conflicts: Object.fromEntries(conflicts),
},
inflight: inflight(),
stale: staleCount(),
intentsCached: store.size,
tombstones: used.size,
journal: journal?.path ?? null,
journalOpen: journal ? journal.isOpen !== false : null,
cu: cfg.cu,
priorityFeeMicroLamports: {
min: cfg.priority.min.toString(),
max: cfg.priority.max.toString(),
percentile: cfg.priority.percentile,
},
skipPreflight: cfg.skipPreflight,
};
const solStale = solCache.atMono === null || mono() - solCache.atMono > cfg.solRefreshMs;
const [hit, , gcr] = await Promise.all([
headerWithin(cfg.collateralMaxAgeMs),
solStale ? refreshSol() : Promise.resolve(true),
readGlobalConfig().then(
(v) => ({ status: 'fulfilled', value: v }),
(e) => ({ status: 'rejected', reason: e }),
),
]);
const rpcAgeLimit = (cfg.collateralLazy ? cfg.headerHeartbeatMs : cfg.collateralRefreshMs) * 2;
out.rpcOk = !!hit && hdrShare.error === null && hdrAgeMs() !== null && hdrAgeMs() <= rpcAgeLimit;
out.rpcError =
hdrShare.error ??
solCache.error ??
(gcr.status === 'rejected' ? sanitize(gcr.reason?.message ?? gcr.reason) : null);
out.headerAgeMs = hdrAgeMs();
if (hdrShare.slot != null) out.slot = String(hdrShare.slot);
const solAge = solCache.atMono === null ? null : mono() - solCache.atMono;
out.solAgeMs = solAge;
if (solCache.lamports !== null && solAge !== null && solAge <= cfg.solMaxAgeMs) {
out.solLamports = solCache.lamports.toString();
out.solLow = solCache.lamports < cfg.minSolLamports;
}
if (hit) {
const header = hit.header;
const b = checkBinding(header, cfg);
out.bindingOk = b.ok;
out.bindingProblems = b.problems;
out.positionAuthority = header.positionAuthority;
out.capabilities = describeCapabilities(header.flags);
out.preferences = describePreferences(header.preferenceBits);
out.collateralQuoteLots = header.collateralQuoteLots;
} else {
out.bindingProblems = [`Trader binding cannot be verified from RPC: ${hdrShare.error ?? 'header unavailable'}`];
}
Object.assign(out, capView());
out.keyRisk = keyRiskOf(out.preferences, gcr.status === 'fulfilled' ? gcr.value : null);
out.ok = out.rpcOk && out.bindingOk && out.meta.ok && !!signer;
healthCache = { at: now(), value: out };
return out;
}
function tombOf(e, ts) {
return {
fp: e.fp,
sig: e.signature ?? null,
status: e.result?.status ?? 'unknown',
rejectReason: e.result?.rejectReason ?? null,
op: e.op,
side: e.side ?? null,
kind: e.kind ?? null,
ts,
};
}
function bury(intentId, t) {
used.set(intentId, t);
if (t.sig) usedBySig.set(t.sig, intentId);
}
function prune() {
const t = now();
for (const [id, e] of store) {
if (e.state === 'final' && e.result && t - (e.finalAt ?? e.createdAt) > cfg.intentTtlMs) {
store.delete(id);
if (e.signature) bySig.delete(e.signature);
bury(id, tombOf(e, t));
}
}
for (const [id, u] of used) {
if (t - u.ts > cfg.tombstoneTtlMs) {
used.delete(id);
if (u.sig) usedBySig.delete(u.sig);
}
}
}
function journalRecords() {
const recs = [];
for (const e of store.values()) {
if (e.signature) recs.push({ ...journalRec(e, 'signed'), ts: e.createdAt });
if (e.result) recs.push({ ...journalRec(e, 'final'), result: e.result, ts: e.finalAt });
}
for (const [intentId, u] of used)
recs.push({
v: 1,
t: 'used',
intentId,
fp: u.fp,
sig: u.sig,
status: u.status,
rejectReason: u.rejectReason,
op: u.op,
side: u.side,
kind: u.kind,
ts: u.ts,
});
return recs;
}
function compact() {
prune();
if (!journal) return;
try {
journal.rewrite(journalRecords());
} catch (err) {
say(`[journal] Compaction failed; previous journal retained: ${err.message}`);
}
}
async function resolveStale() {
if (resolving) return;
resolving = true;
try {
for (const e of [...store.values()]) {
if (shuttingDown) break;
if (e.result || !e.signature || e.state !== 'stale') continue;
try {
await probe(e, { history: true });
} catch (err) {
say(`[resolve] ${e.intentId}: ${sanitize(err?.message ?? err)}`);
}
}
} finally {
resolving = false;
}
}
function recoverFromJournal() {
if (!journal) return 0;
const t = now();
const groups = new Map();
for (const r of journal.load()) {
if (
r?.v !== 1 ||
!/^[A-Za-z0-9._:-]{8,128}$/.test(r.intentId ?? '') ||
!['signed', 'final', 'used'].includes(r.t) ||
typeof r.fp !== 'string' ||
!r.fp ||
!['order', 'cancel'].includes(r.op) ||
!Number.isSafeInteger(r.ts) ||
r.ts < 0 ||
(r.sig != null && !isSignatureString(r.sig)) ||
(r.t === 'signed' && (!r.sig || !/^[0-9]+$/.test(r.lvbh ?? ''))) ||
(r.t === 'final' && !['confirmed', 'failed'].includes(r.result?.status))
)
throw new Error('Intent journal contains an invalid record; startup refused');
if (!groups.has(r.intentId)) groups.set(r.intentId, []);
if (groups.get(r.intentId).some((prior) => prior.fp !== r.fp))
throw new Error('Intent journal reuses an intent identity with conflicting requests; startup refused');
groups.get(r.intentId).push(r);
}
let n = 0;
for (const [intentId, recs] of groups) {
const final = recs.find((r) => r.t === 'final');
const signedRec = recs.find((r) => r.t === 'signed');
const tombRec = recs.find((r) => r.t === 'used');
if (!final && !signedRec) {
if (tombRec && t - (tombRec.ts ?? 0) <= cfg.tombstoneTtlMs) {
bury(intentId, {
fp: tombRec.fp,
sig: tombRec.sig ?? null,
status: tombRec.status ?? 'unknown',
rejectReason: tombRec.rejectReason ?? null,
op: tombRec.op,
side: tombRec.side ?? null,
kind: tombRec.kind ?? null,
ts: tombRec.ts ?? t,
});
}
continue;
}
if (final && t - (final.ts ?? 0) > cfg.intentTtlMs) {
if (t - (final.ts ?? 0) <= cfg.tombstoneTtlMs) {
bury(intentId, {
fp: final.fp,
sig: final.sig ?? null,
status: final.result?.status ?? 'unknown',
rejectReason: final.result?.rejectReason ?? null,
op: final.op,
side: final.side ?? null,
kind: final.kind ?? null,
ts: final.ts ?? t,
});
}
continue;
}
const first = signedRec ?? final;
const e = {
intentId,
fp: first.fp,
op: first.op,
symbol: first.symbol,
side: first.side ?? undefined,
kind: first.kind ?? undefined,
ids: first.ids ?? undefined,
createdAt: first.ts ?? t,
createdMono: mono(),
state: 'recovering',
signature: null,
wire: null,
result: null,
recovered: true,
};
for (const r of recs) {
if (r.sig) e.signature = r.sig;
if (r.lvbh) e.lastValidBlockHeight = r.lvbh;
if (r.lvs) e.lastValidSlot = r.lvs;
}
if (final) {
e.result = final.result;
e.state = 'final';
e.finalAt = final.ts ?? t;
}
store.set(intentId, e);
if (e.signature) bySig.set(e.signature, e);
n++;
}
for (const e of store.values()) {
if (e.recovered && !e.result) {
if (!e.signature) {
e.result = notSent('Recovered intent has no signed transaction; it was never broadcast');
e.state = 'final';
e.finalAt = t;
} else {
say(
`[journal] Recovering unresolved intent ${e.intentId} sig=${e.signature}; signed bytes will not be rebuilt`,
);
e.done = confirmLoop(e, { rebroadcast: false }).catch((err) => {
e.state = 'stale';
say(`[journal] ${e.intentId}: ${sanitize(err?.message ?? err)}`);
});
}
}
}
compact();
return n;
}
async function start() {
traderPda = String(
await getPhoenixTraderSubaccountAddress({
authority: address(cfg.authority),
traderPdaIndex: cfg.traderPdaIndex,
subaccountIndex: cfg.subaccountIndex,
phoenixProgramAddress: address(PROGRAM_ADDRESS),
}),
);
const recovered = recoverFromJournal();
if (recovered) say(`[journal] Recovered ${recovered} intents and ${used.size} spent-intent tombstones`);
await refreshMeta();
const schedule = () => {
metaTimer = setTimeout(
async () => {
await refreshMeta();
if (!shuttingDown) schedule();
},
meta.data ? cfg.metaRefreshMs : cfg.metaRefreshMs,
);
metaTimer.unref?.();
};
schedule();
pruneTimer = setInterval(compact, cfg.resolveMs);
pruneTimer.unref?.();
resolveTimer = setInterval(() => {
resolveStale().catch(() => {});
}, cfg.resolveMs);
resolveTimer.unref?.();
const headerRead = await refreshCollateral();
if (multOn()) {
if (headerRead) {
say(
`[cap] ${fmtUsd(collat.quoteLots)} ${fmtUsd(collateralCapQuoteLots(collat.quoteLots, cfg.maxOrderCollateralMultBps))} (${multLabel()})`,
);
} else {
say(
`[cap] Initial collateral read failed: ${collat.error}; retry interval ${Math.round(cfg.collateralRefreshMs / 1000)} seconds`,
);
}
} else if (!headerRead) {
say(
`[hdr] Initial trader header read failed: ${hdrShare.error}; retry interval ${Math.round(cfg.collateralRefreshMs / 1000)} seconds`,
);
}
const heartbeatForWarn = cfg.collateralLazy ? cfg.headerHeartbeatMs : cfg.collateralRefreshMs;
if (cfg.headerShareMs > 0 && cfg.headerShareMs < heartbeatForWarn) {
say(
`[hdr] Header share age ${Math.round(cfg.headerShareMs / 1000)} seconds is shorter than heartbeat ${Math.round(heartbeatForWarn / 1000)} seconds; /account may trigger additional reads`,
);
}
const heartbeatMs = cfg.collateralLazy ? cfg.headerHeartbeatMs : cfg.collateralRefreshMs;
collateralTimer = setInterval(() => {
refreshCollateral().catch(() => {});
}, heartbeatMs);
collateralTimer.unref?.();
if (!(await refreshSol()))
say(
`[sol] Initial delegate SOL read failed: ${solCache.error}; retry interval ${Math.round(cfg.solRefreshMs / 1000)} seconds`,
);
solTimer = setInterval(() => {
refreshSol().catch(() => {});
}, cfg.solRefreshMs);
solTimer.unref?.();
return { traderPda };
}
async function stop({ graceMs }) {
if (!Number.isSafeInteger(graceMs) || graceMs < 0) throw new TypeError('Explicit shutdown grace required');
shuttingDown = true;
clearTimeout(metaTimer);
clearInterval(pruneTimer);
clearInterval(resolveTimer);
clearInterval(collateralTimer);
clearInterval(solTimer);
const until = mono() + graceMs;
while (inflight() > 0 && mono() < until) await sleep(100);
journal?.close();
}
return {
start,
stop,
submit,
simulate,
txStatus,
account,
health,
refreshMeta,
refreshCollateral,
capView,
resolveStale,
compact,
get traderPda() {
return traderPda;
},
_store: store,
_used: used,
};
}
function send(res, code, body) {
res.writeHead(code, { 'content-type': 'application/json; charset=utf-8', 'cache-control': 'no-store' });
res.end(toJson(body));
}
async function readJsonBody(req) {
let s = '';
for await (const c of req) {
s += c;
if (s.length > 16 * 1024) throw new HttpError(413, 'too_large', 'Request body exceeds 16 KiB');
}
if (!s) throw bad('Request body is empty');
try {
return JSON.parse(s);
} catch {
throw bad('Request body is not valid JSON');
}
}
export function createHttpHandler(sc, cfg, { log = defaultLog } = {}) {
const tokenDigest = sha256(cfg.token);
const sanitize = makeSanitizer([cfg.rpcUrl, cfg.token]);
const authorized = (h) => {
if (typeof h !== 'string' || !h.startsWith('Bearer ')) return false;
return crypto.timingSafeEqual(sha256(h.slice(7)), tokenDigest);
};
const WRITE = { '/order': 'order', '/cancel': 'cancel', '/cancelAll': 'cancelAll' };
const SIM = { '/simulate/order': 'order', '/simulate/cancel': 'cancel', '/simulate/cancelAll': 'cancelAll' };
return async (req, res) => {
try {
if (!isLoopbackPeer(req.socket?.remoteAddress))
return send(res, 403, { error: 'forbidden', message: 'Only loopback requests are accepted' });
if (req.headers.origin)
return send(res, 403, {
error: 'forbidden',
message: 'Browser Origin requests are forbidden on the signer service',
});
if (!authorized(req.headers.authorization)) return send(res, 401, { error: 'unauthorized' });
const url = new URL(req.url, 'http://127.0.0.1');
const p = url.pathname;
if (req.method === 'GET' && p === '/health') return send(res, 200, await sc.health());
if (req.method === 'POST' && WRITE[p]) return send(res, 200, await sc.submit(WRITE[p], await readJsonBody(req)));
if (req.method === 'POST' && SIM[p]) return send(res, 200, await sc.simulate(SIM[p], await readJsonBody(req)));
let m = /^\/tx\/([^/]+)$/.exec(p);
if (req.method === 'GET' && m) return send(res, 200, await sc.txStatus(m[1]));
m = /^\/account\/([^/]+)$/.exec(p);
if (req.method === 'GET' && m) {
const num = (k, d, max) => {
const v = url.searchParams.get(k);
if (v === null || v === '') return d;
if (!/^\d{1,7}$/.test(v) || Number(v) > max) throw bad(`${k}: 0..${max}`);
return Number(v);
};
const offset = num('offset', 0, 10000000);
const length = num('length', 224, 16384);
if (length < 1) throw bad('length ≥ 1');
return send(res, 200, await sc.account(m[1], offset, length));
}
return send(res, 404, { error: 'not_found' });
} catch (e) {
const known = e instanceof HttpError;
const status = known ? e.status : 500;
if (!known) log(sanitize(`[http] ${req.method} ${req.url}: ${e?.stack ?? e}`));
else if (status >= 500) log(sanitize(`[http] ${req.method} ${req.url}: ${status} ${e.code}: ${e.message}`));
return send(res, status, {
error: known ? e.code : 'internal',
message: sanitize(e?.message ?? e),
...(known && e.extra ? e.extra : {}),
});
}
};
}
function defaultLog(m) {
console.log(`[${new Date().toISOString()}] ${m}`);
}
export async function main(env = process.env) {
let secretInfo;
try {
secretInfo = takeBotSecret(env);
} catch (e) {
defaultLog(`FATAL: ${e.message}`);
process.exit(1);
}
let cfg;
try {
cfg = parseConfig(env);
} catch (e) {
defaultLog(`FATAL: ${e.message}`);
process.exit(1);
}
let signer;
try {
signer = await loadBotSigner(secretInfo.secret, cfg.botPubkey);
} catch (e) {
defaultLog(`FATAL: ${e.message}`);
process.exit(1);
}
secretInfo.secret = null;
if (secretInfo.alsoInEnv)
defaultLog(
`Delegate secret loaded from credential ${CREDENTIAL_NAME}; a duplicate PHX_SIGNER_BOT_SECRET_KEY remains in the environment`,
);
const rpc = createJsonRpc(cfg.rpcUrl, { timeoutMs: cfg.rpcTimeoutMs });
const api = createRestClient(cfg.apiUrl, { timeoutMs: cfg.rpcTimeoutMs });
let journal = null;
if (cfg.journalPath) {
try {
journal = openJournal(cfg.journalPath);
} catch (e) {
defaultLog(`FATAL: cannot open intent journal ${cfg.journalPath}: ${e.message}`);
process.exit(1);
}
} else {
defaultLog('No intent journal: readonly mode only');
}
const sc = createSidecar({ config: cfg, signer, rpc, api, journal });
const { traderPda } = await sc.start();
const server = http.createServer(createHttpHandler(sc, cfg));
server.requestTimeout = 60000;
server.headersTimeout = 10000;
await new Promise((resolve, reject) => {
server.once('error', reject);
server.listen(cfg.port, cfg.bind, resolve);
});
defaultLog(
`Phoenix signer ready: delegate=${cfg.botPubkey}, owner=${cfg.authority}, trader=${traderPda} (${cfg.traderPdaIndex}/${cfg.subaccountIndex}), ` +
`rpc ${maskUrl(cfg.rpcUrl)}, api ${cfg.apiUrl}, ${cfg.bind}:${cfg.port}, readonly=${cfg.readonly}, ` +
`cu order/cancel ${cfg.cu.order}/${cfg.cu.cancel}, priority fee maximum ${cfg.priority.max}µ, preflight=${!cfg.skipPreflight}, journal=${cfg.journalPath ?? 'none'}, ` +
`key source=${secretInfo.source}, notional cap=${describeNotionalCap(cfg)}`,
);
try {
const h = await sc.health();
defaultLog(
`health: bindingOk=${h.bindingOk}${h.bindingProblems.length ? ` (${h.bindingProblems.join('; ')})` : ''} state=${h.capabilities?.state ?? '?'} prefs=${h.preferences?.bits ?? '?'} sol=${h.solLamports ?? '?'} meta=${h.meta.markets} rpcOk=${h.rpcOk} notionalCap=${h.maxOrderNotionalUsd !== null ? `$${h.maxOrderNotionalUsd}` : h.maxOrderNotionalUnknown ? 'unknown: collateral not read' : 'off'} (${h.maxOrderNotionalMode})`,
);
if (h.keyRisk?.warning) {
defaultLog(
`keyRisk=${h.keyRisk.warning}: position authority can use the native SOL swap feature (trader opt-out=${h.keyRisk.traderOptOut}, ` +
`exchange kill switch=${h.keyRisk.exchangeKillSwitch}, native SOL feature active=${h.keyRisk.nativeSolCollateralActive}); see signer README for scope`,
);
}
} catch (e) {
defaultLog(`Startup health read failed: ${e.message}`);
}
let stopping = false;
const shutdown = async (sig) => {
if (stopping) return;
stopping = true;
defaultLog(`${sig}: stopping HTTP intake and waiting for in-flight intents`);
server.close();
await sc.stop({ graceMs: cfg.httpWaitMs });
process.exit(0);
};
process.on('SIGTERM', () => shutdown('SIGTERM'));
process.on('SIGINT', () => shutdown('SIGINT'));
}
if (process.argv[1] && path.resolve(process.argv[1]) === fileURLToPath(import.meta.url)) {
main().catch((e) => {
defaultLog(`FATAL: ${e?.message ?? e}`);
process.exit(1);
});
}