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
550 changes: 73 additions & 477 deletions docs/WEBHOOK_SECURITY_FLOW.md

Large diffs are not rendered by default.

194 changes: 194 additions & 0 deletions src/accessControl.ts
Original file line number Diff line number Diff line change
Expand Up @@ -109,3 +109,197 @@ export class AclManager {
return `${resourceId}:${address}`;
}
}

export type WithdrawalStatus =
| "pending"
| "approved"
| "rejected"
| "executed";

export interface WithdrawalRequest {
id: string;
resourceId: string;
requester: string;
amount: string;
status: WithdrawalStatus;
approvals: string[];
rejections: string[];
createdAt: number;
updatedAt: number;
}

export interface WithdrawalApprovalOptions {
requiredApprovals?: number;
acl?: AclManager;
}

export type WithdrawalEventType =
| "submitted"
| "approved"
| "rejected"
| "executed";

export interface WithdrawalEvent {
type: WithdrawalEventType;
request: WithdrawalRequest;
actor: string;
timestamp: number;
}

export type WithdrawalEventListener = (event: WithdrawalEvent) => void;

/**
* Manages custody withdrawal approval workflows.
*
* A withdrawal request moves through a multi-step approval state machine:
* pending -> approved (once enough approvals are collected) -> executed,
* or pending -> rejected. Approvers must hold access to the resource
* when an AclManager is provided.
*/
export class WithdrawalApprovalWorkflow {
private readonly requiredApprovals: number;
private readonly acl?: AclManager;
private readonly requests = new Map<string, WithdrawalRequest>();
private readonly listeners = new Set<WithdrawalEventListener>();

constructor(options: WithdrawalApprovalOptions = {}) {
this.requiredApprovals = options.requiredApprovals ?? 1;
if (this.requiredApprovals < 1) {
throw new Error("requiredApprovals must be at least 1");
}
this.acl = options.acl;
}

/**
* Register a listener for withdrawal lifecycle events.
*
* @returns An unsubscribe function.
*/
onEvent(listener: WithdrawalEventListener): () => void {
this.listeners.add(listener);
return () => this.listeners.delete(listener);
}

/**
* Submit a new withdrawal request in the pending state.
*/
async submit(params: {
id: string;
resourceId: string;
requester: string;
amount: string;
}): Promise<WithdrawalRequest> {
if (this.requests.has(params.id)) {
throw new Error(`Withdrawal request ${params.id} already exists`);
}
const now = Date.now();
const request: WithdrawalRequest = {
id: params.id,
resourceId: params.resourceId,
requester: params.requester,
amount: params.amount,
status: "pending",
approvals: [],
rejections: [],
createdAt: now,
updatedAt: now,
};
this.requests.set(request.id, request);
this.emit("submitted", request, params.requester);
return request;
}

/**
* Approve a pending withdrawal request.
*/
async approve(id: string, approver: string): Promise<WithdrawalRequest> {
const request = this.getRequest(id);
if (request.status !== "pending") {
throw new Error(`Cannot approve request in status ${request.status}`);
}
if (approver === request.requester) {
throw new Error("Requester cannot approve their own withdrawal");
}
if (request.approvals.includes(approver)) {
throw new Error(`Approver ${approver} has already approved`);
}
if (this.acl && !(await this.acl.check(request.resourceId, approver))) {
throw new Error(`Approver ${approver} lacks access to resource`);
}

request.approvals.push(approver);
request.updatedAt = Date.now();
if (request.approvals.length >= this.requiredApprovals) {
request.status = "approved";
}
this.emit("approved", request, approver);
return request;
}

/**
* Reject a pending withdrawal request.
*/
async reject(id: string, rejector: string): Promise<WithdrawalRequest> {
const request = this.getRequest(id);
if (request.status !== "pending") {
throw new Error(`Cannot reject request in status ${request.status}`);
}
if (request.rejections.includes(rejector)) {
throw new Error(`Rejector ${rejector} has already rejected`);
}
if (this.acl && !(await this.acl.check(request.resourceId, rejector))) {
throw new Error(`Rejector ${rejector} lacks access to resource`);
}

request.rejections.push(rejector);
request.status = "rejected";
request.updatedAt = Date.now();
this.emit("rejected", request, rejector);
return request;
}

/**
* Execute an approved withdrawal request.
*/
async execute(id: string, executor: string): Promise<WithdrawalRequest> {
const request = this.getRequest(id);
if (request.status !== "approved") {
throw new Error(`Cannot execute request in status ${request.status}`);
}
request.status = "executed";
request.updatedAt = Date.now();
this.emit("executed", request, executor);
return request;
}

/**
* Retrieve a withdrawal request by id.
*/
get(id: string): WithdrawalRequest | undefined {
return this.requests.get(id);
}

private getRequest(id: string): WithdrawalRequest {
const request = this.requests.get(id);
if (!request) {
throw new Error(`Withdrawal request ${id} not found`);
}
return request;
}

private emit(
type: WithdrawalEventType,
request: WithdrawalRequest,
actor: string,
): void {
const event: WithdrawalEvent = {
type,
request,
actor,
timestamp: Date.now(),
};
for (const listener of this.listeners) {
listener(event);
}
}
}
159 changes: 159 additions & 0 deletions src/accountDataManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,9 @@
* Wraps `Operation.manageData()` with validation for the protocol's 64-byte
* key/value limits and 64-entry-per-account cap, so callers can store custom
* metadata alongside SDK state without hand-rolling raw manageData calls.
*
* Also provides SDK data migration utilities for moving data entries between
* accounts, with lifecycle event handling for observability.
*/

import {
Expand Down Expand Up @@ -36,6 +39,49 @@ export interface AccountDataManagerConfig {
networkPassphrase: string;
}

/** Options controlling an SDK data migration. */
export interface DataMigrationOptions {
/** Source account whose data entries are migrated. */
sourceAccountId: string;
/** Destination account that receives the migrated entries. */
destinationAccountId: string;
/** Secret key used to sign transactions on the source account. */
sourceSignerSecret: string;
/** Secret key used to sign transactions on the destination account. */
destinationSignerSecret: string;
/** Restrict the migration to these keys; defaults to all source entries. */
keys?: string[];
/** Delete migrated entries from the source account after copying. */
deleteSource?: boolean;
}

/** Per-key outcome of a migration run. */
export interface DataMigrationEntryResult {
key: string;
status: "migrated" | "skipped" | "failed";
error?: string;
}

/** Aggregate result of a migration run. */
export interface DataMigrationResult {
sourceAccountId: string;
destinationAccountId: string;
entries: DataMigrationEntryResult[];
migrated: number;
skipped: number;
failed: number;
}

/** Lifecycle events emitted during a migration. */
export type DataMigrationEvent =
| { type: "start"; sourceAccountId: string; destinationAccountId: string; total: number }
| { type: "progress"; key: string; index: number; total: number; status: DataMigrationEntryResult["status"] }
| { type: "complete"; result: DataMigrationResult }
| { type: "error"; key?: string; error: Error };

/** Listener invoked for each {@link DataMigrationEvent}. */
export type DataMigrationEventListener = (event: DataMigrationEvent) => void;

function byteLength(value: string): number {
return Buffer.byteLength(value, "utf8");
}
Expand All @@ -47,6 +93,7 @@ function byteLength(value: string): number {
export class AccountDataManager {
private readonly server: Horizon.Server;
private readonly networkPassphrase: string;
private readonly migrationListeners = new Set<DataMigrationEventListener>();

constructor(config: AccountDataManagerConfig) {
this.server = new Horizon.Server(config.horizonUrl);
Expand Down Expand Up @@ -102,6 +149,118 @@ export class AccountDataManager {
return result;
}

/**
* Subscribe to migration lifecycle events.
*
* @returns an unsubscribe function that removes the listener.
*/
onMigrationEvent(listener: DataMigrationEventListener): () => void {
this.migrationListeners.add(listener);
return () => {
this.migrationListeners.delete(listener);
};
}

/**
* Migrate data entries from a source account to a destination account.
*
* Copies each selected entry to the destination, optionally deleting it
* from the source, and emits `start`, `progress`, `complete`, and `error`
* lifecycle events. Per-key failures are captured in the result rather than
* aborting the whole run; a fatal error (e.g. source load failure) emits an
* `error` event and rejects.
*/
async migrateData(options: DataMigrationOptions): Promise<DataMigrationResult> {
const {
sourceAccountId,
destinationAccountId,
sourceSignerSecret,
destinationSignerSecret,
deleteSource = false,
} = options;

let sourceEntries: AccountDataMap;
try {
sourceEntries = await this.list(sourceAccountId);
} catch (err) {
const error = err instanceof Error ? err : new Error(String(err));
this.emitMigrationEvent({ type: "error", error });
throw error;
}

const keys = options.keys ?? Object.keys(sourceEntries);
const total = keys.length;
this.emitMigrationEvent({
type: "start",
sourceAccountId,
destinationAccountId,
total,
});

const entries: DataMigrationEntryResult[] = [];
let migrated = 0;
let skipped = 0;
let failed = 0;

for (let index = 0; index < keys.length; index++) {
const key = keys[index]!;
let status: DataMigrationEntryResult["status"];
let errorMessage: string | undefined;

if (!Object.prototype.hasOwnProperty.call(sourceEntries, key)) {
status = "skipped";
skipped++;
} else {
try {
await this.set(
destinationAccountId,
key,
sourceEntries[key]!,
destinationSignerSecret,
);
if (deleteSource) {
await this.delete(sourceAccountId, key, sourceSignerSecret);
}
status = "migrated";
migrated++;
} catch (err) {
status = "failed";
failed++;
errorMessage = err instanceof Error ? err.message : String(err);
this.emitMigrationEvent({
type: "error",
key,
error: err instanceof Error ? err : new Error(String(err)),
});
}
}

const entry: DataMigrationEntryResult = { key, status };
if (errorMessage !== undefined) {
entry.error = errorMessage;
}
entries.push(entry);
this.emitMigrationEvent({ type: "progress", key, index, total, status });
}

const result: DataMigrationResult = {
sourceAccountId,
destinationAccountId,
entries,
migrated,
skipped,
failed,
};
this.emitMigrationEvent({ type: "complete", result });
return result;
}

private emitMigrationEvent(event: DataMigrationEvent): void {
for (const listener of this.migrationListeners) {
listener(event);
}
}

private async validateEntry(accountId: string, key: string, value: string): Promise<void> {
if (byteLength(key) > MAX_DATA_ENTRY_BYTES) {
throw new DataEntryValidationError(
Expand Down
Loading