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: 14 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,20 @@ new StellarSplitClient(config: StellarSplitClientConfig)
|-------|-------------|
| `ProfilerSession` | Record SDK method timings during a session and produce a flame-graph-compatible report |

`ProfilerSession` instruments both SDK methods and each client's Soroban RPC
server, so RPC round-trips appear as frames nested inside the method that
issued them. `report()` emits a **speedscope v0.6** file (validated with ajv
against the published schema in the test suite); `exportJSON(path)` writes it
to disk for drag-and-drop into [speedscope.app](https://www.speedscope.app).

```ts
const profiler = new ProfilerSession({ name: "my-session" });
profiler.start();
await client.pay({ invoiceId, payer });
profiler.stop();
profiler.exportJSON("./profile.json");
```

### Webhook Validation

| Function | Returns | Description |
Expand Down
14 changes: 14 additions & 0 deletions src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -654,6 +654,17 @@ export class StellarSplitClient extends TypedEventEmitter<SplitClientEventMap> {
*/
private _effectiveRpcPoolSize = 0;
private _batcher: BatchedRpcClient | null = null;

/**
* Live client instances, so {@link ProfilerSession} can instrument each
* client's RPC server.
*
* Held on the class rather than patched per-instance because the profiler
* wraps the prototype once but every client owns its own `server` object.
*
* @internal
*/
static readonly _instances = new Set<StellarSplitClient>();
private _telemetryHookManager = new TelemetryHookManager();
private _timeoutManager: TimeoutManager | null = null;
private _traceIdManager = new TraceIdManager();
Expand Down Expand Up @@ -830,6 +841,9 @@ export class StellarSplitClient extends TypedEventEmitter<SplitClientEventMap> {
* @throws {Error} If the method fails.
*/
super();

// Register so ProfilerSession can wrap this client's RPC server.
StellarSplitClient._instances.add(this);
/**
* validateOrThrow
* @param params - The parameters for the method.
Expand Down
129 changes: 123 additions & 6 deletions src/profiler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,26 @@ export interface RpcCallTiming {
operation: string;
/** Duration in milliseconds. */
durationMs: number;
/** Unix timestamp (ms) when the RPC call started, for nesting. */
timestamp: number;
}

/**
* Methods on the Soroban RPC client that the profiler wraps to capture
* per-RPC timings nested inside the SDK method that issued them.
*/
const PROFILED_RPC_METHODS = [
"getTransaction",
"getLedgerEntries",
"simulateTransaction",
"sendTransaction",
"getAccount",
"getHealth",
"getLatestLedger",
"getEvents",
"getVersion",
] as const;

/** An aggregated snapshot of one start/stop recording cycle. */
export interface ProfileSession {
startedAt: number;
Expand Down Expand Up @@ -79,6 +97,7 @@ export interface SpeedscopeEventedProfile {
*/
export interface SpeedscopeProfile {
$schema: "https://www.speedscope.app/file-format-schema.json";
exporter: string;
version: "0.6.0";
name: string;
activeProfileIndex: number;
Expand Down Expand Up @@ -120,6 +139,12 @@ export class ProfilerSession {
private currentStartedAt = 0;
private currentStoppedAt = 0;
private originalMethods = new Map<string, Function>();
private originalRpcMethods = new Map<object, Map<string, Function>>();
/**
* Per-invocation RPC sink. Set by the method wrapper before delegating so
* that concurrent SDK calls do not cross-contaminate each other's timings.
*/
private activeRpcSink: RpcCallTiming[] | null = null;

constructor(options: ProfilerSessionOptions = {}) {
this.sessionName = options.name ?? "StellarSplit SDK";
Expand Down Expand Up @@ -152,8 +177,16 @@ export class ProfilerSession {
const startTs = Date.now();
const rpcCalls: RpcCallTiming[] = [];

// Route any RPC calls made during this invocation into our sink.
const previousSink = thisSession.activeRpcSink;
thisSession.activeRpcSink = rpcCalls;
const restoreSink = (): void => {
thisSession.activeRpcSink = previousSink;
};

const record = (success: boolean, error?: string): void => {
const durationMs = performance.now() - startTime;
restoreSink();
thisSession.currentEntries.push({
method: methodName,
durationMs,
Expand Down Expand Up @@ -190,11 +223,90 @@ export class ProfilerSession {
} as unknown as Function;
}

this.instrumentRpc();
this.currentEntries = [];
this.currentStartedAt = Date.now();
this.active = true;
}

/**
* Wrap each client's Soroban RPC server so individual RPC round-trips are
* timed and attributed to the SDK method that issued them.
*
* Patching the server (rather than each SDK call site) keeps the
* instrumentation to one place, so new SDK methods are profiled for free.
*/
private instrumentRpc(): void {
for (const client of StellarSplitClient._instances) {
const server = (client as unknown as { server?: Record<string, unknown> }).server;
if (!server || typeof server !== "object") continue;

const originals = new Map<string, Function>();
const thisSession = this;

for (const methodName of PROFILED_RPC_METHODS) {
const original = server[methodName];
if (typeof original !== "function") continue;
originals.set(methodName, original as Function);

server[methodName] = function (this: unknown, ...args: unknown[]) {
const start = performance.now();
const startedAt = Date.now();
const sink = thisSession.activeRpcSink;

const finish = (): void => {
// Only record when a sink is active, i.e. inside a profiled call.
if (!sink) return;
sink.push({
operation: methodName,
durationMs: performance.now() - start,
timestamp: startedAt,
});
};

let result: unknown;
try {
result = (original as Function).apply(this, args);
} catch (err: unknown) {
finish();
throw err;
}

if (result && typeof (result as Promise<unknown>).then === "function") {
return (result as Promise<unknown>).then(
(value: unknown) => {
finish();
return value;
},
(err: unknown) => {
finish();
return Promise.reject(err);
},
);
}

finish();
return result;
};
}

if (originals.size > 0) {
this.originalRpcMethods.set(server as object, originals);
}
}
}

/** Restore the RPC servers wrapped by {@link instrumentRpc}. */
private restoreRpc(): void {
for (const [server, originals] of this.originalRpcMethods.entries()) {
for (const [methodName, original] of originals.entries()) {
(server as Record<string, unknown>)[methodName] = original;
}
}
this.originalRpcMethods.clear();
this.activeRpcSink = null;
}

// -------------------------------------------------------------------------
// stop() — restore prototype and finalise session
// -------------------------------------------------------------------------
Expand Down Expand Up @@ -222,6 +334,7 @@ export class ProfilerSession {

this.active = false;
this.originalMethods.clear();
this.restoreRpc();

return this.getReport();
}
Expand Down Expand Up @@ -272,16 +385,18 @@ export class ProfilerSession {
const methodFrame = getFrame(entry.method);
events.push({ type: "O", at: relStart, frame: methodFrame });

// Emit nested RPC call events (synthesised, evenly distributed within
// the parent window so the flame graph is always valid)
// Emit nested RPC frames at their real offsets within the parent
// window. Clamping to the parent bounds keeps the flame graph valid
// even if a timer drifts slightly past the enclosing call.
if (entry.rpcCalls && entry.rpcCalls.length > 0) {
let cursor = relStart;
for (const rpc of entry.rpcCalls) {
const rpcFrame = getFrame(`rpc:${rpc.operation}`);
const rpcEnd = Math.min(cursor + rpc.durationMs, relEnd);
events.push({ type: "O", at: cursor, frame: rpcFrame });
const rpcStart = Math.max(relStart, rpc.timestamp - sessionStart);
const rpcEnd = Math.min(rpcStart + rpc.durationMs, relEnd);
if (rpcEnd <= rpcStart) continue;

events.push({ type: "O", at: rpcStart, frame: rpcFrame });
events.push({ type: "C", at: rpcEnd, frame: rpcFrame });
cursor = rpcEnd;
}
}

Expand Down Expand Up @@ -309,6 +424,8 @@ export class ProfilerSession {

return {
$schema: "https://www.speedscope.app/file-format-schema.json",
// Required by the speedscope schema; identifies the producing tool.
exporter: "@stellar-split/sdk",
version: "0.6.0",
name: this.sessionName,
activeProfileIndex: 0,
Expand Down
Loading
Loading