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
98 changes: 98 additions & 0 deletions src/accounts/AccountMergeDetector.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,32 @@ export interface MergeEventPayload {
mergedAt: Date;
}

/**
* A custody account managed by the detector. Custody accounts are watched
* accounts whose funds are held on behalf of a recipient and which may be
* merged into a destination account.
*/
export interface CustodyAccount {
/** The custody account address */
address: string;
/** Optional human-readable label */
label?: string;
/** Optional asset the custody account is expected to hold */
asset?: { code: string; issuer: string };
/** When the custody account was registered */
registeredAt: Date;
}

/** Payload emitted on custody account lifecycle events. */
export interface CustodyAccountEventPayload {
account: CustodyAccount;
at: Date;
}

export class AccountMergeDetector extends EventEmitter {
private watchedAccounts = new Set<string>();
private mergeCache = new Map<string, string>(); // source -> destination mapping
private custodyAccounts = new Map<string, CustodyAccount>();
private streamActive = false;
private checkInterval: NodeJS.Timeout | null = null;

Expand Down Expand Up @@ -89,6 +112,72 @@ export class AccountMergeDetector extends EventEmitter {
this.watchedAccounts.delete(accountId);
}

/**
* Register a custody account and begin watching it for merges.
* Emits "custody:registered" with the created custody account.
*/
registerCustodyAccount(
address: string,
options: { label?: string; asset?: { code: string; issuer: string } } = {},
): CustodyAccount {
const existing = this.custodyAccounts.get(address);
if (existing) {
return existing;
}

const account: CustodyAccount = {
address,
label: options.label,
asset: options.asset,
registeredAt: new Date(),
};

this.custodyAccounts.set(address, account);
this.watchAccount(address);

this.emit("custody:registered", {
account,
at: account.registeredAt,
} satisfies CustodyAccountEventPayload);

return account;
}

/**
* Remove a custody account from management and stop watching it.
* Emits "custody:removed" when an account was actually removed.
*/
removeCustodyAccount(address: string): boolean {
const account = this.custodyAccounts.get(address);
if (!account) {
return false;
}

this.custodyAccounts.delete(address);
this.unwatchAccount(address);

this.emit("custody:removed", {
account,
at: new Date(),
} satisfies CustodyAccountEventPayload);

return true;
}

/**
* Retrieve a managed custody account by address.
*/
getCustodyAccount(address: string): CustodyAccount | undefined {
return this.custodyAccounts.get(address);
}

/**
* List all managed custody accounts.
*/
listCustodyAccounts(): CustodyAccount[] {
return Array.from(this.custodyAccounts.values());
}

/**
* Check if an account has been merged and resolve the final destination.
* Supports recursive merge chains up to depth 5.
Expand Down Expand Up @@ -200,6 +289,15 @@ export class AccountMergeDetector extends EventEmitter {
mergedAt: event.timestamp,
} satisfies MergeEventPayload);

// If the merged source was a managed custody account, emit a lifecycle event
const custodyAccount = this.custodyAccounts.get(sourceAccount);
if (custodyAccount) {
this.emit("custody:merged", {
account: custodyAccount,
at: event.timestamp,
} satisfies CustodyAccountEventPayload);
}

// Notify the client to reroute recipients
try {
// The client will handle rerouting via rerouteRecipient method
Expand Down
194 changes: 194 additions & 0 deletions src/accounts/AccountSignerWeightCalculator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,64 @@ export interface SignerWeightResult {
missingWeight: number;
}

/**
* A single custody account signer entry, as returned by the custody account
* management helpers.
*/
export interface CustodySigner {
/** The signer's public key (G… address, pre-auth tx, or hash(x)). */
key: string;
/** The signing weight assigned to this signer. */
weight: number;
}

/**
* A snapshot of a custody account's signer configuration and thresholds.
*/
export interface CustodyAccount {
/** The Stellar account G… address. */
accountId: string;
/** The account's current signers. */
signers: CustodySigner[];
/** The account's threshold configuration. */
thresholds: {
low: number;
medium: number;
high: number;
};
}

/**
* Event names emitted by the custody account management helpers during
* lifecycle operations.
*/
export type CustodyAccountEvent =
| "signer:added"
| "signer:removed"
| "signer:updated"
| "threshold:updated"
| "account:loaded";

/**
* Payload delivered to custody account event listeners.
*/
export interface CustodyAccountEventPayload {
/** The account the event pertains to. */
accountId: string;
/** The lifecycle event that occurred. */
event: CustodyAccountEvent;
/** The signer affected by the event, when applicable. */
signer?: CustodySigner;
/** The threshold level affected by the event, when applicable. */
thresholdLevel?: ThresholdLevel;
/** The previous value before the change, when applicable. */
previousValue?: number;
/** The new value after the change, when applicable. */
newValue?: number;
}

export type CustodyAccountEventListener = (payload: CustodyAccountEventPayload) => void;

// ---------------------------------------------------------------------------
// Cache entry
// ---------------------------------------------------------------------------
Expand All @@ -49,6 +107,7 @@ export class AccountSignerWeightCalculator {
/** Cache TTL in milliseconds (default 30 seconds). */
private readonly cacheTtlMs: number;
private readonly cache = new Map<string, CacheEntry>();
private readonly listeners = new Set<CustodyAccountEventListener>();

constructor(horizonUrl: string, cacheTtlMs = 30_000) {
this.server = new Horizon.Server(horizonUrl, { allowHttp: horizonUrl.startsWith("http://") });
Expand Down Expand Up @@ -124,10 +183,145 @@ export class AccountSignerWeightCalculator {
return result.sufficient;
}

// --------------------------------------------------------------------------
// Custody account management helpers
// --------------------------------------------------------------------------

/**
* Load a custody account snapshot (signers + thresholds) from Horizon.
* Emits an `account:loaded` event on success.
*/
async loadCustodyAccount(accountId: string): Promise<CustodyAccount> {
const record = await this._loadAccount(accountId);
const account = this._toCustodyAccount(accountId, record);
this._emit({ accountId, event: "account:loaded" });
return account;
}

/**
* Add a signer to a custody account snapshot. If the signer already exists
* its weight is updated instead. Emits `signer:added` or `signer:updated`.
*/
addSigner(account: CustodyAccount, signer: CustodySigner): CustodyAccount {
const existing = account.signers.find((s) => s.key === signer.key);
let next: CustodyAccount;

if (existing) {
next = {
...account,
signers: account.signers.map((s) => (s.key === signer.key ? { ...signer } : s)),
};
this._emit({
accountId: account.accountId,
event: "signer:updated",
signer: { ...signer },
previousValue: existing.weight,
newValue: signer.weight,
});
} else {
next = { ...account, signers: [...account.signers, { ...signer }] };
this._emit({
accountId: account.accountId,
event: "signer:added",
signer: { ...signer },
newValue: signer.weight,
});
}

return next;
}

/**
* Remove a signer from a custody account snapshot by public key.
* Emits `signer:removed` when a signer was actually removed.
*/
removeSigner(account: CustodyAccount, signerKey: string): CustodyAccount {
const existing = account.signers.find((s) => s.key === signerKey);
if (!existing) {
return account;
}

const next: CustodyAccount = {
...account,
signers: account.signers.filter((s) => s.key !== signerKey),
};

this._emit({
accountId: account.accountId,
event: "signer:removed",
signer: { ...existing },
previousValue: existing.weight,
});

return next;
}

/**
* Update a threshold level on a custody account snapshot.
* Emits `threshold:updated` when the value changes.
*/
updateThreshold(
account: CustodyAccount,
level: ThresholdLevel,
value: number,
): CustodyAccount {
const previousValue = account.thresholds[level];
if (previousValue === value) {
return account;
}

const next: CustodyAccount = {
...account,
thresholds: { ...account.thresholds, [level]: value },
};

this._emit({
accountId: account.accountId,
event: "threshold:updated",
thresholdLevel: level,
previousValue,
newValue: value,
});

return next;
}

/**
* Register a listener for custody account lifecycle events.
* Returns an unsubscribe function.
*/
onCustodyAccountEvent(listener: CustodyAccountEventListener): () => void {
this.listeners.add(listener);
return () => {
this.listeners.delete(listener);
};
}

// --------------------------------------------------------------------------
// Private helpers
// --------------------------------------------------------------------------

private _emit(payload: CustodyAccountEventPayload): void {
for (const listener of this.listeners) {
listener(payload);
}
}

private _toCustodyAccount(
accountId: string,
record: Horizon.AccountResponse,
): CustodyAccount {
return {
accountId,
signers: record.signers.map((s) => ({ key: s.key, weight: s.weight })),
thresholds: {
low: record.thresholds.low_threshold,
medium: record.thresholds.med_threshold,
high: record.thresholds.high_threshold,
},
};
}

private async _loadAccount(accountId: string): Promise<Horizon.AccountResponse> {
const now = Date.now();
const entry = this.cache.get(accountId);
Expand Down
31 changes: 31 additions & 0 deletions src/audit/AuditTrailHasher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,11 @@ import * as crypto from 'crypto';
// Use node crypto webcrypto subtle
const subtle = crypto.webcrypto.subtle;

export type AuditTrailListener = (entry: AuditChainEntry) => void;

export class AuditTrailHasher {
private entries: AuditChainEntry[] = [];
private listeners: Set<AuditTrailListener> = new Set();

constructor(entries: AuditChainEntry[] = []) {
this.entries = [...entries];
Expand All @@ -22,6 +25,33 @@ export class AuditTrailHasher {
return hashArray.map(b => b.toString(16).padStart(2, '0')).join('');
}

/**
* Subscribes to audit trail events. Returns an unsubscribe function.
*/
onAppend(listener: AuditTrailListener): () => void {
this.listeners.add(listener);
return () => {
this.listeners.delete(listener);
};
}

/**
* Removes a previously registered audit trail listener.
*/
offAppend(listener: AuditTrailListener): void {
this.listeners.delete(listener);
}

private emitAppend(entry: AuditChainEntry): void {
for (const listener of this.listeners) {
try {
listener(entry);
} catch {
// Listener errors must not break the audit chain.
}
}
}

/**
* Appends a new event to the audit trail
*/
Expand All @@ -35,6 +65,7 @@ export class AuditTrailHasher {

const entry: AuditChainEntry = { event, hash, prevHash, index };
this.entries.push(entry);
this.emitAppend(entry);
return entry;
}

Expand Down
Loading