From a5aacd2053085b2d52a64f53e5cad7712dc8699b Mon Sep 17 00:00:00 2001 From: selezenart Date: Sun, 13 Sep 2026 16:44:57 +0200 Subject: [PATCH] fix(activity): scope indexed transfers to the workspace The Graph activity query is sent with the shared server wallet alongside the workspace's user wallets, so a stored observation also carries transfers made for other workspaces. Settlement and uncertain counts were SQL-scoped to the workspace while the unmatched transfer count was taken from the whole observation, so Payment proof reported one workspace's 11 settlements next to 40 unmatched transfers drawn from every workspace. Filter the observed transfers to the ones this workspace can own before they are matched, listed or counted: the sender is a workspace payer wallet, or the transaction hash is already recorded on a workspace settlement, job or attempt. The attempt hash keeps a stranded server payment visible as unmatched evidence. --- packages/storage-postgres/src/jobs.ts | 136 ++++++++++++++------ packages/storage-postgres/test/jobs.test.ts | 64 +++++++++ 2 files changed, 158 insertions(+), 42 deletions(-) diff --git a/packages/storage-postgres/src/jobs.ts b/packages/storage-postgres/src/jobs.ts index 7751e68..e1f3125 100644 --- a/packages/storage-postgres/src/jobs.ts +++ b/packages/storage-postgres/src/jobs.ts @@ -566,56 +566,63 @@ export class JobLedger { } async activity(workspaceId: string): Promise { - const [observation, settlements, uncertain, recordedTransfers, activityJobs] = - await Promise.all([ - this.#pool.query<{ - freshness: string; - coverage_note: string; - observed_at: Date; - payload: unknown; - }>( - `SELECT freshness, coverage_note, observed_at, payload FROM wallet_activity_observations + const [ + observation, + settlements, + uncertain, + recordedTransfers, + activityJobs, + workspaceWalletRows, + workspaceHashRows, + ] = await Promise.all([ + this.#pool.query<{ + freshness: string; + coverage_note: string; + observed_at: Date; + payload: unknown; + }>( + `SELECT freshness, coverage_note, observed_at, payload FROM wallet_activity_observations WHERE workspace_id = $1 ORDER BY observation_id DESC LIMIT 1`, - [workspaceId], - ), - this.#pool.query<{ count: string }>( - `SELECT count(*)::text AS count FROM ( + [workspaceId], + ), + this.#pool.query<{ count: string }>( + `SELECT count(*)::text AS count FROM ( SELECT j.business_intent_id FROM resumable_jobs j JOIN settlements s ON s.business_intent_id = j.business_intent_id WHERE j.workspace_id = $1 ) recorded`, - [workspaceId], - ), - this.#pool.query<{ count: string }>( - `SELECT count(*)::text AS count FROM ( + [workspaceId], + ), + this.#pool.query<{ count: string }>( + `SELECT count(*)::text AS count FROM ( SELECT j.business_intent_id FROM resumable_jobs j JOIN business_intents i ON i.business_intent_id = j.business_intent_id WHERE j.workspace_id = $1 AND i.state = 'UNKNOWN' ) uncertain`, - [workspaceId], - ), - this.#pool.query<{ - transaction_hash: string; - transfer_log_index: number; - job_id: string; - }>( - `SELECT s.transaction_hash, s.transfer_log_index, j.job_id + [workspaceId], + ), + this.#pool.query<{ + transaction_hash: string; + transfer_log_index: number; + job_id: string; + }>( + `SELECT s.transaction_hash, s.transfer_log_index, j.job_id FROM settlements s JOIN resumable_jobs j ON j.business_intent_id = s.business_intent_id WHERE j.workspace_id = $1`, - [workspaceId], - ), - this.#pool.query<{ - job_id: string; - business_intent_id: string; - payment_state: JobView['payment_state']; - payment_mode: PaymentMode; - transaction_hash: string | null; - recipient: string; - amount_atomic: string; - transfer_log_index: number | null; - }>( - `SELECT j.job_id, j.business_intent_id, i.state AS payment_state, j.payment_mode, + [workspaceId], + ), + this.#pool.query<{ + job_id: string; + business_intent_id: string; + payment_state: JobView['payment_state']; + payment_mode: PaymentMode; + transaction_hash: string | null; + recipient: string; + amount_atomic: string; + transfer_log_index: number | null; + }>( + `SELECT j.job_id, j.business_intent_id, i.state AS payment_state, j.payment_mode, COALESCE(s.transaction_hash, j.payment_transaction_hash, latest.provider_transaction_hash) AS transaction_hash, j.supplier_quote->>'recipient' AS recipient, j.supplier_quote->>'amount_atomic' AS amount_atomic, @@ -633,11 +640,56 @@ export class JobLedger { WHERE j.workspace_id = $1 ORDER BY j.updated_at DESC, j.job_id ASC LIMIT 100`, - [workspaceId], - ), - ]); + [workspaceId], + ), + this.#pool.query<{ payer_wallet: string }>( + `SELECT DISTINCT lower(payer_wallet) AS payer_wallet + FROM resumable_jobs + WHERE workspace_id = $1 AND payer_wallet IS NOT NULL`, + [workspaceId], + ), + this.#pool.query<{ transaction_hash: string }>( + `SELECT DISTINCT lower(transaction_hash) AS transaction_hash FROM ( + SELECT workspace_settlement.transaction_hash + FROM settlements workspace_settlement + JOIN resumable_jobs workspace_job + ON workspace_job.business_intent_id = workspace_settlement.business_intent_id + WHERE workspace_job.workspace_id = $1 + UNION + SELECT workspace_job.payment_transaction_hash + FROM resumable_jobs workspace_job + WHERE workspace_job.workspace_id = $1 + AND workspace_job.payment_transaction_hash IS NOT NULL + UNION + SELECT workspace_attempt.provider_transaction_hash + FROM attempts workspace_attempt + JOIN resumable_jobs workspace_job + ON workspace_job.business_intent_id = workspace_attempt.business_intent_id + WHERE workspace_job.workspace_id = $1 + AND workspace_attempt.provider_transaction_hash IS NOT NULL + ) workspace_transaction_hashes`, + [workspaceId], + ), + ]); const row = observation.rows[0]; - const indexedTransfers = row ? parseActivityTransfers(row.payload) : []; + const observedTransfers = row ? parseActivityTransfers(row.payload) : []; + // The Graph is queried with the shared server wallet as well as this + // workspace's user wallets, so the raw observation also carries transfers + // that belong to other workspaces. Evidence is workspace-scoped: keep only + // transfers this workspace can own, either by payer wallet or by a + // transaction hash the workspace already recorded on a settlement, a job or + // an attempt. + const workspaceWallets = new Set( + workspaceWalletRows.rows.map((wallet) => wallet.payer_wallet.toLowerCase()), + ); + const workspaceHashes = new Set( + workspaceHashRows.rows.map((hash) => hash.transaction_hash.toLowerCase()), + ); + const indexedTransfers = observedTransfers.filter( + (transfer) => + workspaceHashes.has(transfer.transaction_hash.toLowerCase()) || + (transfer.sender !== undefined && workspaceWallets.has(transfer.sender.toLowerCase())), + ); const recordedByTransfer = new Map( recordedTransfers.rows.map((settlement) => [ activityTransferKey(settlement.transaction_hash, settlement.transfer_log_index), diff --git a/packages/storage-postgres/test/jobs.test.ts b/packages/storage-postgres/test/jobs.test.ts index b47ba74..b462901 100644 --- a/packages/storage-postgres/test/jobs.test.ts +++ b/packages/storage-postgres/test/jobs.test.ts @@ -54,6 +54,9 @@ describe('JobLedger delivery recovery', () => { it('matches indexed transfers to workspace settlements and surfaces unmatched activity', async () => { const recordedHash = `0x${'a'.repeat(64)}`; const unmatchedHash = `0x${'b'.repeat(64)}`; + const foreignHash = `0x${'d'.repeat(64)}`; + const workspaceWallet = '0x2222222222222222222222222222222222222222'; + const foreignWallet = '0x3333333333333333333333333333333333333333'; const pool = { async query(sql: string) { if (sql.includes('FROM wallet_activity_observations')) { @@ -75,6 +78,14 @@ describe('JobLedger delivery recovery', () => { { transaction_hash: unmatchedHash, log_index: 4, + sender: workspaceWallet, + recipient: failedJob.supplier_quote.recipient, + amount_atomic: failedJob.supplier_quote.amount_atomic, + }, + { + transaction_hash: foreignHash, + log_index: 9, + sender: foreignWallet, recipient: failedJob.supplier_quote.recipient, amount_atomic: failedJob.supplier_quote.amount_atomic, }, @@ -95,6 +106,10 @@ describe('JobLedger delivery recovery', () => { return { rows: [{ count: '1' }] }; } if (sql.includes("i.state = 'UNKNOWN'")) return { rows: [{ count: '0' }] }; + if (sql.includes('AS payer_wallet')) return { rows: [{ payer_wallet: workspaceWallet }] }; + if (sql.includes('workspace_transaction_hashes')) { + return { rows: [{ transaction_hash: recordedHash }] }; + } return { rows: [] }; }, }; @@ -119,6 +134,52 @@ describe('JobLedger delivery recovery', () => { }); }); + it('excludes shared server wallet transfers that belong to another workspace', async () => { + const foreignHash = `0x${'e'.repeat(64)}`; + const serverWallet = '0x4444444444444444444444444444444444444444'; + const pool = { + async query(sql: string) { + if (sql.includes('FROM wallet_activity_observations')) { + return { + rows: [ + { + freshness: 'FRESH', + coverage_note: 'indexed', + observed_at: new Date('2026-09-07T12:00:00.000Z'), + payload: { + transfers: [ + { + transaction_hash: foreignHash, + log_index: 1, + sender: serverWallet, + recipient: failedJob.supplier_quote.recipient, + amount_atomic: failedJob.supplier_quote.amount_atomic, + }, + ], + }, + }, + ], + }; + } + if (sql.includes('SELECT count(*)::text AS count FROM (') && sql.includes('recorded')) { + return { rows: [{ count: '1' }] }; + } + if (sql.includes("i.state = 'UNKNOWN'")) return { rows: [{ count: '0' }] }; + return { rows: [] }; + }, + }; + const ledger = new JobLedger(pool as never, { + now: () => new Date('2026-09-07T12:01:00.000Z'), + nextAttemptId: () => 'unused', + }); + + await expect(ledger.activity('workspace-unit')).resolves.toMatchObject({ + recorded_settlement_count: 1, + unmatched_transfer_count: 0, + transfers: [], + }); + }); + it('projects every site outcome with its Graph match status', async () => { const indexedHash = `0x${'c'.repeat(64)}`; const rejectedJob = { @@ -167,6 +228,9 @@ describe('JobLedger delivery recovery', () => { } if (sql.includes('FROM settlements s')) return { rows: [] }; if (sql.includes("i.state = 'UNKNOWN'")) return { rows: [{ count: '0' }] }; + if (sql.includes('workspace_transaction_hashes')) { + return { rows: [{ transaction_hash: indexedHash }] }; + } if (sql.includes('SELECT j.job_id, j.business_intent_id')) { return { rows: [failedJobActivity, rejectedJob] }; }