diff --git a/src/client.ts b/src/client.ts index 3cbb992..c4d6e55 100644 --- a/src/client.ts +++ b/src/client.ts @@ -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}. */ @@ -489,6 +490,7 @@ export interface StellarSplitClientConfig { /** Retry settings applied per RPC call (maxRetries, baseDelayMs, etc.). */ retry?: Partial; }; + offlineQueue?: import("./offlineQueue.js").OfflineQueueConfig; /** * Optional configuration for the CLOSED/OPEN/HALF_OPEN circuit breaker * (src/resilience/CircuitBreaker.ts) guarding the transaction-submission @@ -709,6 +711,7 @@ export class StellarSplitClient extends TypedEventEmitter { * 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 | null = null; private _sorobanFeatureDetector: SorobanFeatureDetector; @@ -1263,6 +1266,14 @@ export class StellarSplitClient extends TypedEventEmitter { * @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 { const { NetworkSwitcher } = await import("./network/NetworkSwitcher.js"); return NetworkSwitcher.switchTo(network, this); @@ -5281,7 +5292,9 @@ export class StellarSplitClient extends TypedEventEmitter { * @throws {Error} If the method fails. */ async checkRPCHealth(): Promise { - return checkRPCHealth(this.server); + const health = await checkRPCHealth(this.server); + if (health.status !== "down" && this._offlineQueue?.config?.enabled) { void this._offlineQueue.drain(); } + return health; } /** @@ -8230,6 +8243,38 @@ export class StellarSplitClient extends TypedEventEmitter { * @returns The transaction hash. * @throws {Error} If the method fails. */ + async release(invoiceId: string): Promise { + 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 { + 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 { const startTime = Date.now(); try { diff --git a/src/offlineQueue.ts b/src/offlineQueue.ts new file mode 100644 index 0000000..c07e1c4 --- /dev/null +++ b/src/offlineQueue.ts @@ -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 { + 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); + } + } + } +} diff --git a/test/offlineQueue.test.ts b/test/offlineQueue.test.ts new file mode 100644 index 0000000..38b31d6 --- /dev/null +++ b/test/offlineQueue.test.ts @@ -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, + 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"); + }); +});