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
47 changes: 46 additions & 1 deletion src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import type { CircuitStateChangeLogEvent } from "./resilience/CircuitBreaker.js"
import { InvoiceStateMachine } from "./state/InvoiceStateMachine.js";
import type { StateMachineConfig } from "./types/state.js";
import { RpcLoadBalancer } from "./rpc/RpcLoadBalancer.js";
import { OfflineQueue, type QueuedOperation } from "./offlineQueue.js";
import type { EndpointConfig, RpcLoadBalancerOptions } from "./rpc/RpcLoadBalancer.js";

/** Events emitted by {@link StellarSplitClient}. */
Expand Down Expand Up @@ -489,6 +490,7 @@ export interface StellarSplitClientConfig {
/** Retry settings applied per RPC call (maxRetries, baseDelayMs, etc.). */
retry?: Partial<ResilientRetryConfig>;
};
offlineQueue?: import("./offlineQueue.js").OfflineQueueConfig;
/**
* Optional configuration for the CLOSED/OPEN/HALF_OPEN circuit breaker
* (src/resilience/CircuitBreaker.ts) guarding the transaction-submission
Expand Down Expand Up @@ -709,6 +711,7 @@ export class StellarSplitClient extends TypedEventEmitter<SplitClientEventMap> {
* configs keep working unchanged; enable via `advancedCircuitBreaker`.
*/
private _advancedCircuitBreaker: AdvancedCircuitBreaker | null = null;
private _offlineQueue: import("./offlineQueue.js").OfflineQueue | null = null;
/** Optimistic UI cache for Invoice reads during a pending pay() call. */
private _optimisticCache: OptimisticCache<Invoice> | null = null;
private _sorobanFeatureDetector: SorobanFeatureDetector;
Expand Down Expand Up @@ -1263,6 +1266,14 @@ export class StellarSplitClient extends TypedEventEmitter<SplitClientEventMap> {
* @returns The result of the method.
* @throws {Error} If the method fails.
*/
getOfflineQueue(): import("./offlineQueue.js").QueuedOperation[] {
return this._offlineQueue ? this._offlineQueue.getQueue() : [];
}

clearOfflineQueue(): void {
if (this._offlineQueue) this._offlineQueue.clear();
}

async switchTo(network: "mainnet" | "testnet" | "futurenet"): Promise<void> {
const { NetworkSwitcher } = await import("./network/NetworkSwitcher.js");
return NetworkSwitcher.switchTo(network, this);
Expand Down Expand Up @@ -5281,7 +5292,9 @@ export class StellarSplitClient extends TypedEventEmitter<SplitClientEventMap> {
* @throws {Error} If the method fails.
*/
async checkRPCHealth(): Promise<RPCHealth> {
return checkRPCHealth(this.server);
const health = await checkRPCHealth(this.server);
if (health.status !== "down" && this._offlineQueue?.config?.enabled) { void this._offlineQueue.drain(); }
return health;
}

/**
Expand Down Expand Up @@ -8230,6 +8243,38 @@ export class StellarSplitClient extends TypedEventEmitter<SplitClientEventMap> {
* @returns The transaction hash.
* @throws {Error} If the method fails.
*/
async release(invoiceId: string): Promise<TxResult> {
if (this._offlineQueue?.config?.enabled) {
try {
const health = await this.checkRPCHealth();
if (health.status === "down") {
this._offlineQueue.enqueue("release", [invoiceId]);
return { txHash: "queued" };
}
} catch (e) {
this._offlineQueue.enqueue("release", [invoiceId]);
return { txHash: "queued" };
}
}
throw new Error("Not implemented");
}

async cancel(invoiceId: string): Promise<TxResult> {
if (this._offlineQueue?.config?.enabled) {
try {
const health = await this.checkRPCHealth();
if (health.status === "down") {
this._offlineQueue.enqueue("cancel", [invoiceId]);
return { txHash: "queued" };
}
} catch (e) {
this._offlineQueue.enqueue("cancel", [invoiceId]);
return { txHash: "queued" };
}
}
throw new Error("Not implemented");
}

async cancelAction(caller: string, actionId: string): Promise<TxResult> {
const startTime = Date.now();
try {
Expand Down
102 changes: 102 additions & 0 deletions src/offlineQueue.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
import type { StellarSplitClient } from "./client.js";

export interface QueuedOperation {
id: string;
method: string;
args: any[];
timestamp: number;
}

export interface OfflineQueueConfig {
enabled: boolean;
maxQueueSize: number;
persistToStorage: boolean;
}

export class OfflineQueue {
private queue: QueuedOperation[] = [];
private config: OfflineQueueConfig;
private client: StellarSplitClient;
private readonly storageKey = "stellar_split_offline_queue";

constructor(client: StellarSplitClient, config: OfflineQueueConfig) {
this.client = client;
this.config = config;
if (this.config.persistToStorage) {
this.loadFromStorage();
}
}

public enqueue(method: string, args: any[]): void {
if (!this.config.enabled) return;

if (this.queue.length >= this.config.maxQueueSize) {
this.queue.shift();
}

const operation: QueuedOperation = {
id: crypto.randomUUID ? crypto.randomUUID() : Math.random().toString(36).substring(2),
method,
args,
timestamp: Date.now(),
};

this.queue.push(operation);

if (this.config.persistToStorage) {
this.saveToStorage();
}
}

public getQueue(): QueuedOperation[] {
return [...this.queue];
}

public clear(): void {
this.queue = [];
if (this.config.persistToStorage) {
this.saveToStorage();
}
}

public async drain(): Promise<void> {
if (this.queue.length === 0) return;

const operationsToDrain = [...this.queue];
this.clear();

for (const op of operationsToDrain) {
try {
// @ts-ignore
await this.client[op.method](...op.args);
} catch (error) {
this.client.emit("queue:operation:failed", { operation: op, error });
}
}

this.client.emit("queue:drained", undefined);
}

private loadFromStorage(): void {
if (typeof localStorage !== "undefined") {
try {
const stored = localStorage.getItem(this.storageKey);
if (stored) {
this.queue = JSON.parse(stored);
}
} catch (e) {
console.warn("Failed to load offline queue from storage", e);
}
}
}

private saveToStorage(): void {
if (typeof localStorage !== "undefined") {
try {
localStorage.setItem(this.storageKey, JSON.stringify(this.queue));
} catch (e) {
console.warn("Failed to save offline queue to storage", e);
}
}
}
}
84 changes: 84 additions & 0 deletions test/offlineQueue.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
import { expect } from "chai";
import { StellarSplitClient } from "../src/client.js";
import * as health from "../src/health.js";
import sinon from "sinon";

describe("OfflineQueue", () => {
let client: StellarSplitClient;
let healthStub: sinon.SinonStub;

beforeEach(() => {
if (typeof localStorage === "undefined") {
(global as any).localStorage = {
store: {} as Record<string, string>,
getItem(key: string) { return this.store[key] || null; },
setItem(key: string, value: string) { this.store[key] = value; },
clear() { this.store = {}; }
};
} else {
localStorage.clear();
}

client = new StellarSplitClient({
rpcUrl: "http://localhost:8000",
networkPassphrase: "Test",
contractId: "C123",
offlineQueue: {
enabled: true,
maxQueueSize: 2,
persistToStorage: true,
}
});

healthStub = sinon.stub(health, "checkRPCHealth").resolves({ status: "down", latencyMs: 0, blockHeight: 0, timestamp: 0 });
});

afterEach(() => {
healthStub.restore();
sinon.restore();
});

it("fills on RPC failure and enforces max size", async () => {
await client.pay({ invoiceId: "1", amount: 100n, payer: "A" });
await client.pay({ invoiceId: "2", amount: 200n, payer: "B" });
await client.pay({ invoiceId: "3", amount: 300n, payer: "C" });

const q = client.getOfflineQueue();
expect(q.length).to.equal(2);
expect(q[0].args[0].invoiceId).to.equal("2");
expect(q[1].args[0].invoiceId).to.equal("3");
});

it("persists to localStorage", async () => {
await client.pay({ invoiceId: "1", amount: 100n, payer: "A" });

const client2 = new StellarSplitClient({
rpcUrl: "http://localhost:8000",
networkPassphrase: "Test",
contractId: "C123",
offlineQueue: {
enabled: true,
maxQueueSize: 5,
persistToStorage: true,
}
});

const q = client2.getOfflineQueue();
expect(q.length).to.equal(1);
expect(q[0].args[0].invoiceId).to.equal("1");
});

it("drains on recovery", async () => {
await client.pay({ invoiceId: "1", amount: 100n, payer: "A" });

healthStub.resolves({ status: "ok", latencyMs: 10, blockHeight: 1, timestamp: 1 });
const payStub = sinon.stub(client, "pay").resolves({ txHash: "success" });

await client.checkRPCHealth();

const q = client.getOfflineQueue();
expect(q.length).to.equal(0);
expect(payStub.calledOnce).to.be.true;
expect(payStub.firstCall.args[0].invoiceId).to.equal("1");
});
});
Loading