diff --git a/src/agreements/agreement-activity.service.spec.ts b/src/agreements/agreement-activity.service.spec.ts new file mode 100644 index 0000000..31e601e --- /dev/null +++ b/src/agreements/agreement-activity.service.spec.ts @@ -0,0 +1,55 @@ +import { AgreementActivityService } from './agreement-activity.service'; + +describe('AgreementActivityService', () => { + it('inserts activity rows via supabase including optional state columns', async () => { + const insert = jest.fn().mockResolvedValue({ error: null }); + const from = jest.fn().mockReturnValue({ insert }); + const supabase = { getClient: () => ({ from }) } as never; + + const svc = new AgreementActivityService(supabase); + await svc.logActivity( + 'agr-1', + 'GWALLET', + 'dispute_opened', + { dispute_id: 'd1' }, + { previousState: 'active', newState: 'disputed' }, + ); + + expect(from).toHaveBeenCalledWith('agreement_activity'); + expect(insert).toHaveBeenCalledWith({ + agreement_id: 'agr-1', + actor_wallet: 'GWALLET', + action: 'dispute_opened', + details: { dispute_id: 'd1' }, + previous_state: 'active', + new_state: 'disputed', + }); + }); + + it('defaults previous_state/new_state to null when omitted', async () => { + const insert = jest.fn().mockResolvedValue({ error: null }); + const from = jest.fn().mockReturnValue({ insert }); + const supabase = { getClient: () => ({ from }) } as never; + + const svc = new AgreementActivityService(supabase); + await svc.logActivity('agr-1', 'G', 'created'); + + expect(insert).toHaveBeenCalledWith({ + agreement_id: 'agr-1', + actor_wallet: 'G', + action: 'created', + details: {}, + previous_state: null, + new_state: null, + }); + }); + + it('swallows insert errors without throwing', async () => { + const insert = jest.fn().mockResolvedValue({ error: { message: 'boom' } }); + const from = jest.fn().mockReturnValue({ insert }); + const supabase = { getClient: () => ({ from }) } as never; + + const svc = new AgreementActivityService(supabase); + await expect(svc.logActivity('agr-1', 'G', 'created')).resolves.toBeUndefined(); + }); +}); diff --git a/src/agreements/agreement-activity.service.ts b/src/agreements/agreement-activity.service.ts new file mode 100644 index 0000000..a014f25 --- /dev/null +++ b/src/agreements/agreement-activity.service.ts @@ -0,0 +1,48 @@ +import { Injectable, Logger } from '@nestjs/common'; +import { SupabaseService } from '../supabase/supabase.service'; + +export type ActivityStates = { + previousState?: string | null; + newState?: string | null; +}; + +/** + * Single shared writer for `agreement_activity` rows. + * All services must use this instead of private logActivity copies. + * + * Supports optional `previous_state` / `new_state` columns introduced by the + * activity-state logging work so status transitions stay queryable. + */ +@Injectable() +export class AgreementActivityService { + private readonly logger = new Logger(AgreementActivityService.name); + + constructor(private readonly supabase: SupabaseService) {} + + async logActivity( + agreementId: string, + actorWallet: string, + action: string, + details: Record = {}, + states: ActivityStates = {}, + ): Promise { + try { + const { error } = await this.supabase + .getClient() + .from('agreement_activity') + .insert({ + agreement_id: agreementId, + actor_wallet: actorWallet, + action, + details, + previous_state: states.previousState ?? null, + new_state: states.newState ?? null, + }); + if (error) { + this.logger.error(`logActivity insert failed: ${error.message}`); + } + } catch (e) { + this.logger.error('logActivity', e); + } + } +} diff --git a/src/agreements/agreement-lifecycle.spec.ts b/src/agreements/agreement-lifecycle.spec.ts index 052fc70..62b73c2 100644 --- a/src/agreements/agreement-lifecycle.spec.ts +++ b/src/agreements/agreement-lifecycle.spec.ts @@ -21,6 +21,7 @@ import type { EventEmitter2 } from '@nestjs/event-emitter'; import { validate } from 'class-validator'; import { plainToInstance } from 'class-transformer'; import { AgreementsService } from './agreements.service'; +import { AgreementActivityService } from './agreement-activity.service'; import { DisputesService } from '../disputes/disputes.service'; import { UpdateAgreementStatusDto } from './dto/update-status.dto'; import type { SupabaseService } from '../supabase/supabase.service'; @@ -392,9 +393,11 @@ describe('AgreementsService lifecycle enforcement (business rules)', () => { beforeEach(() => { db = new InMemoryDb(); emit = jest.fn(); + const activity = new AgreementActivityService(db as unknown as SupabaseService); service = new AgreementsService( db as unknown as SupabaseService, { emit } as unknown as EventEmitter2, + activity, ); }); @@ -768,8 +771,9 @@ describe('Dispute flows drive the agreement lifecycle', () => { db = new InMemoryDb(); emit = jest.fn(); const emitter = { emit } as unknown as EventEmitter2; - agreements = new AgreementsService(db as unknown as SupabaseService, emitter); - disputes = new DisputesService(db as unknown as SupabaseService, agreements, emitter); + const activity = new AgreementActivityService(db as unknown as SupabaseService); + agreements = new AgreementsService(db as unknown as SupabaseService, emitter, activity); + disputes = new DisputesService(db as unknown as SupabaseService, agreements, emitter, activity); db.insert('agreements', { id: AGREEMENT_ID, @@ -809,7 +813,17 @@ describe('Dispute flows drive the agreement lifecycle', () => { const disputeId = await openDispute(); expect(db.agreement(AGREEMENT_ID).status).toBe('disputed'); + // Shared side-effect path logs status_changed_to_* then dispute-specific action expect(db.activityFor(AGREEMENT_ID)).toEqual([ + expect.objectContaining({ + action: 'status_changed_to_disputed', + details: expect.objectContaining({ + status: 'disputed', + from: 'active', + to: 'disputed', + dispute_id: disputeId, + }), + }), expect.objectContaining({ action: 'dispute_opened', details: expect.objectContaining({ dispute_id: disputeId }), @@ -828,20 +842,22 @@ describe('Dispute flows drive the agreement lifecycle', () => { }); const timeline = db.activityFor(AGREEMENT_ID); - // Dispute lifecycle events live in the SAME agreement timeline (deduped logger). + // Shared side-effect path also writes status_changed_to_* around dispute actions. expect(timeline.map((a: Row) => a.action)).toEqual([ + 'status_changed_to_disputed', 'dispute_opened', 'dispute_resolver_assigned', + 'status_changed_to_resolved', 'dispute_resolved', ]); - expect(timeline[0]).toEqual( + expect(timeline.find((a: Row) => a.action === 'dispute_opened')).toEqual( expect.objectContaining({ action: 'dispute_opened', previous_state: 'active', new_state: 'disputed', }), ); - expect(timeline[2]).toEqual( + expect(timeline.find((a: Row) => a.action === 'dispute_resolved')).toEqual( expect.objectContaining({ action: 'dispute_resolved', previous_state: 'disputed', @@ -877,8 +893,10 @@ describe('Dispute flows drive the agreement lifecycle', () => { expect(db.agreement(AGREEMENT_ID).status).toBe('resolved'); expect(db.agreement(AGREEMENT_ID).completed_at).toBeDefined(); expect(db.activityFor(AGREEMENT_ID).map((a) => a.action)).toEqual([ + 'status_changed_to_disputed', 'dispute_opened', 'dispute_resolver_assigned', + 'status_changed_to_resolved', 'dispute_resolved', ]); diff --git a/src/agreements/agreements.module.ts b/src/agreements/agreements.module.ts index cf67f75..351fdd2 100644 --- a/src/agreements/agreements.module.ts +++ b/src/agreements/agreements.module.ts @@ -2,11 +2,12 @@ import { Module } from '@nestjs/common'; import { AuthModule } from '../auth/auth.module'; import { AgreementsController } from './agreements.controller'; import { AgreementsService } from './agreements.service'; +import { AgreementActivityService } from './agreement-activity.service'; @Module({ imports: [AuthModule], controllers: [AgreementsController], - providers: [AgreementsService], - exports: [AgreementsService], + providers: [AgreementsService, AgreementActivityService], + exports: [AgreementsService, AgreementActivityService], }) export class AgreementsModule {} diff --git a/src/agreements/agreements.service.ts b/src/agreements/agreements.service.ts index 5ef2eff..856ac5c 100644 --- a/src/agreements/agreements.service.ts +++ b/src/agreements/agreements.service.ts @@ -16,12 +16,14 @@ import { invalidTransitionMessage, milestonesSatisfyCompletion, } from './agreement-lifecycle'; +import { AgreementActivityService } from './agreement-activity.service'; @Injectable() export class AgreementsService { constructor( private readonly supabase: SupabaseService, private readonly eventEmitter: EventEmitter2, + private readonly activity: AgreementActivityService, ) {} private async walletForUserId(userId: string): Promise { @@ -142,7 +144,7 @@ export class AgreementsService { console.error('agreement_participants insert:', participantsError); } - await this.logActivity(agreement.id, dto.created_by, 'created', { + await this.activity.logActivity(agreement.id, dto.created_by, 'created', { title: dto.title, amount: dto.amount, }); @@ -153,7 +155,7 @@ export class AgreementsService { }>; for (let index = 0; index < createdMilestones.length; index++) { const milestone = createdMilestones[index]; - await this.logActivity( + await this.activity.logActivity( agreement.id, dto.created_by, 'milestone_created', @@ -193,7 +195,7 @@ export class AgreementsService { if (error) return { success: false, error: error.message }; - await this.logActivity(agreementId, dto.actor_wallet, 'contract_linked', { + await this.activity.logActivity(agreementId, dto.actor_wallet, 'contract_linked', { contract_id: dto.contract_id, }); return { success: true, error: null }; @@ -244,7 +246,7 @@ export class AgreementsService { if (error) return { success: false, error: error.message }; - await this.logActivity( + await this.activity.logActivity( agreementId, dto.actor_wallet, `status_changed_to_${dto.status}`, @@ -331,7 +333,7 @@ export class AgreementsService { if (updateError) return { success: false, error: updateError.message }; - await this.logActivity( + await this.activity.logActivity( agreementId, dto.actor_wallet, `milestone_${dto.status}`, @@ -449,6 +451,90 @@ export class AgreementsService { return { agreement: data, error: null }; } + /** + * Shared side-effect path for Agreement status changes. + * Used by DisputesService (and any other internal caller) so dispute-driven + * updates emit the same activity log shape and domain events as updateStatus. + */ + async applyStatusChange( + agreementId: string, + actorWallet: string, + toStatus: string, + options: { + fromStatus?: string; + enforceTransition?: boolean; + activityDetails?: Record; + } = {}, + ): Promise<{ success: boolean; error: string | null; fromStatus?: string }> { + const { data: current, error: fetchError } = await this.supabase + .getClient() + .from('agreements') + .select('status, title, amount, asset') + .eq('id', agreementId) + .single(); + + if (fetchError || !current) { + return { success: false, error: fetchError?.message || 'Agreement not found' }; + } + + const fromStatus = options.fromStatus ?? (current.status as string); + + if (options.enforceTransition && !canTransition(fromStatus, toStatus)) { + return { success: false, error: invalidTransitionMessage(fromStatus, toStatus) }; + } + + const updates: Record = { + status: toStatus, + updated_at: new Date().toISOString(), + }; + if (toStatus === 'funded') { + updates.funded_at = new Date().toISOString(); + } else if (toStatus === 'completed' || toStatus === 'resolved') { + updates.completed_at = new Date().toISOString(); + } + + const { error } = await this.supabase + .getClient() + .from('agreements') + .update(updates) + .eq('id', agreementId); + + if (error) return { success: false, error: error.message }; + + await this.activity.logActivity( + agreementId, + actorWallet, + `status_changed_to_${toStatus}`, + { + status: toStatus, + from: fromStatus, + to: toStatus, + ...(options.activityDetails ?? {}), + }, + { previousState: fromStatus, newState: toStatus }, + ); + + if (toStatus === 'funded') { + this.eventEmitter.emit(AGREEMENT_EVENTS.FUNDED, { + agreementId, + title: current.title, + amount: current.amount, + asset: current.asset ?? 'USDC', + fundedByWallet: actorWallet, + }); + } else if (toStatus === 'completed' || toStatus === 'resolved') { + this.eventEmitter.emit(AGREEMENT_EVENTS.COMPLETED, { + agreementId, + title: current.title, + totalAmount: current.amount, + asset: current.asset ?? 'USDC', + completedAt: new Date().toISOString(), + }); + } + + return { success: true, error: null, fromStatus }; + } + async getActivity(userId: string, agreementId: string) { await this.assertCanAccessAgreement(userId, agreementId); @@ -462,45 +548,4 @@ export class AgreementsService { if (error) return { activities: [], error: error.message }; return { activities: data ?? [], error: null }; } - - /** Optional state-transition metadata for an activity entry. */ - private async logActivity( - agreementId: string, - actorWallet: string, - action: string, - details: Record = {}, - states: { previousState?: string | null; newState?: string | null } = {}, - ) { - try { - await this.supabase - .getClient() - .from('agreement_activity') - .insert({ - agreement_id: agreementId, - actor_wallet: actorWallet, - action, - details, - previous_state: states.previousState ?? null, - new_state: states.newState ?? null, - }); - } catch (e) { - console.error('logAgreementActivity', e); - } - } - - /** - * Public, backward-compatible alias for {@link logActivity}. Kept for external - * callers that expect this name and used by other modules (e.g. DisputesService) - * so dispute lifecycle events land in the same agreement timeline without - * duplicating the logger. - */ - async logAgreementActivity( - agreementId: string, - actorWallet: string, - action: string, - details: Record = {}, - states: { previousState?: string | null; newState?: string | null } = {}, - ) { - return this.logActivity(agreementId, actorWallet, action, details, states); - } } diff --git a/src/agreements/dispute-side-effects.spec.ts b/src/agreements/dispute-side-effects.spec.ts new file mode 100644 index 0000000..650a652 --- /dev/null +++ b/src/agreements/dispute-side-effects.spec.ts @@ -0,0 +1,338 @@ +import { EventEmitter2 } from '@nestjs/event-emitter'; +import { AgreementsService } from './agreements.service'; +import { AgreementActivityService } from './agreement-activity.service'; +import { AGREEMENT_EVENTS } from '../common/events/agreement-events.constants'; +import { DisputesService } from '../disputes/disputes.service'; +import { DISPUTE_OPENED, DISPUTE_RESOLVED } from '../common/constants/notification-events'; + +type Row = Record; + +/** + * Minimal in-memory Supabase stub covering the chains used by + * AgreementsService.applyStatusChange / updateStatus and DisputesService open/resolve. + */ +function buildDb(seed: { + agreements: Row[]; + auth_users?: Row[]; + agreement_participants?: Row[]; + disputes?: Row[]; + dispute_resolutions?: Row[]; + agreement_activity?: Row[]; +}) { + const tables: Record = { + agreements: seed.agreements.map((r) => ({ ...r })), + auth_users: (seed.auth_users ?? []).map((r) => ({ ...r })), + agreement_participants: (seed.agreement_participants ?? []).map((r) => ({ ...r })), + disputes: (seed.disputes ?? []).map((r) => ({ ...r })), + dispute_resolutions: (seed.dispute_resolutions ?? []).map((r) => ({ ...r })), + agreement_activity: (seed.agreement_activity ?? []).map((r) => ({ ...r })), + }; + + const activityInserts: Row[] = []; + + function chain(table: string) { + const rows = tables[table] ?? []; + const filters: Array<(r: Row) => boolean> = []; + let mode: 'select' | 'insert' | 'update' = 'select'; + let payload: Row | Row[] | null = null; + let wantSingle = false; + let wantMaybe = false; + + const api: Record = {}; + const self = () => api; + + api.select = () => self(); + api.eq = (col: string, val: unknown) => { + filters.push((r) => r[col] === val); + return self(); + }; + api.in = (col: string, vals: unknown[]) => { + filters.push((r) => vals.includes(r[col])); + return self(); + }; + api.limit = () => self(); + api.order = () => self(); + api.insert = (data: Row | Row[]) => { + mode = 'insert'; + payload = data; + return self(); + }; + api.update = (data: Row) => { + mode = 'update'; + payload = data; + return self(); + }; + api.single = () => { + wantSingle = true; + return finalize(); + }; + api.maybeSingle = () => { + wantMaybe = true; + return finalize(); + }; + + const finalize = () => { + if (mode === 'insert') { + const items = Array.isArray(payload) ? payload : [payload as Row]; + const created = items.map((item, i) => ({ + id: (item.id as string) || `${table}-${tables[table].length + i + 1}`, + created_at: new Date().toISOString(), + updated_at: new Date().toISOString(), + ...item, + })); + tables[table].push(...created); + if (table === 'agreement_activity') activityInserts.push(...created); + const data = created.length === 1 ? created[0] : created; + return Promise.resolve({ data, error: null }); + } + + let matched = rows.filter((r) => filters.every((f) => f(r))); + + if (mode === 'update') { + matched = matched.map((r) => { + Object.assign(r, payload as Row); + return r; + }); + // when chained .select after update + } + + if (wantSingle || wantMaybe) { + const data = matched[0] ?? null; + if (wantSingle && !data) { + return Promise.resolve({ data: null, error: { message: 'not found' } }); + } + return Promise.resolve({ data, error: null }); + } + + return Promise.resolve({ data: matched, error: null }); + }; + + // bare await on chain (e.g. update without single) + (api as { then?: unknown }).then = (resolve: (v: unknown) => unknown) => + finalize().then(resolve); + + return api; + } + + return { + tables, + activityInserts, + client: { + from: (table: string) => chain(table), + }, + }; +} + +const USER = 'user-1'; +const WALLET = 'GWALLET-PAYER'; +const RESOLVER = 'GWALLET-RESOLVER'; +const AGREEMENT_ID = 'agr-dispute-1'; + +function makeAgreements(db: ReturnType, emitter: EventEmitter2) { + const supabase = { getClient: () => db.client } as never; + const activity = new AgreementActivityService(supabase); + // spy so tests can assert shared path — keep the Mock handle (avoids unbound-method) + const logActivity = jest.spyOn(activity, 'logActivity'); + const svc = new AgreementsService(supabase, emitter, activity); + return { svc, activity, logActivity }; +} + +function makeDisputes( + db: ReturnType, + emitter: EventEmitter2, + agreements: AgreementsService, + activity: AgreementActivityService, +) { + const supabase = { getClient: () => db.client } as never; + return new DisputesService(supabase, agreements, emitter, activity); +} + +describe('dispute-driven Agreement side effects (issue #58)', () => { + let db: ReturnType; + let emitter: EventEmitter2; + let agreements: AgreementsService; + let activity: AgreementActivityService; + let logActivity: jest.SpyInstance; + let disputes: DisputesService; + let emitted: Array<{ event: string; payload: unknown }>; + + beforeEach(() => { + db = buildDb({ + agreements: [ + { + id: AGREEMENT_ID, + status: 'active', + title: 'Escrow job', + amount: '100', + asset: 'USDC', + created_by: WALLET, + milestones: [], + }, + ], + auth_users: [ + { id: USER, wallet_public_key: WALLET }, + { id: 'user-resolver', wallet_public_key: RESOLVER }, + ], + agreement_participants: [ + { agreement_id: AGREEMENT_ID, wallet_address: WALLET, role: 'payer' }, + { agreement_id: AGREEMENT_ID, wallet_address: 'GWALLET-PAYEE', role: 'payee' }, + ], + disputes: [], + dispute_resolutions: [], + agreement_activity: [], + }); + emitter = new EventEmitter2(); + emitted = []; + emitter.onAny((event, payload) => { + emitted.push({ event: String(event), payload }); + }); + ({ svc: agreements, activity, logActivity } = makeAgreements(db, emitter)); + disputes = makeDisputes(db, emitter, agreements, activity); + }); + + it('unit: openDispute uses shared activity + emits DISPUTE_OPENED + status_changed_to_disputed', async () => { + const applySpy = jest.spyOn(agreements, 'applyStatusChange'); + + const result = await disputes.openDispute(USER, { + agreement_id: AGREEMENT_ID, + opened_by: WALLET, + reason: 'Work incomplete', + evidence_urls: [], + }); + + expect(result.error).toBeNull(); + expect(result.dispute).toBeTruthy(); + expect(applySpy).toHaveBeenCalledWith( + AGREEMENT_ID, + WALLET, + 'disputed', + expect.objectContaining({ + activityDetails: expect.objectContaining({ source: 'dispute' }), + }), + ); + + // shared logActivity used (at least status_changed + dispute_opened) + expect(logActivity).toHaveBeenCalled(); + + const actions = logActivity.mock.calls.map((c) => c[2] as string); + expect(actions).toContain('status_changed_to_disputed'); + expect(actions).toContain('dispute_opened'); + + expect(emitted.some((e) => e.event === DISPUTE_OPENED)).toBe(true); + expect(db.tables.agreements[0].status).toBe('disputed'); + }); + + it('unit: resolveDispute uses shared path + emits DISPUTE_RESOLVED + COMPLETED', async () => { + // seed open dispute under review + db.tables.disputes.push({ + id: 'disp-1', + agreement_id: AGREEMENT_ID, + opened_by: WALLET, + reason: 'x', + evidence_urls: [], + status: 'under_review', + resolver_wallet: RESOLVER, + }); + db.tables.agreements[0].status = 'disputed'; + db.tables.auth_users.push({ id: 'user-resolver', wallet_public_key: RESOLVER }); + + // resolver user id must match wallet + const applySpy = jest.spyOn(agreements, 'applyStatusChange'); + const result = await disputes.resolveDispute('user-resolver', 'disp-1', { + resolved_by: RESOLVER, + payer_percentage: 40, + payee_percentage: 60, + resolution_notes: 'Split', + }); + + expect(result.error).toBeNull(); + expect(applySpy).toHaveBeenCalledWith( + AGREEMENT_ID, + RESOLVER, + 'resolved', + expect.objectContaining({ + activityDetails: expect.objectContaining({ source: 'dispute' }), + }), + ); + + const actions = logActivity.mock.calls.map((c) => c[2] as string); + expect(actions).toContain('status_changed_to_resolved'); + expect(actions).toContain('dispute_resolved'); + + expect(emitted.some((e) => e.event === DISPUTE_RESOLVED)).toBe(true); + expect(emitted.some((e) => e.event === AGREEMENT_EVENTS.COMPLETED)).toBe(true); + expect(db.tables.agreements[0].status).toBe('resolved'); + }); + + it('parity: dispute-driven disputed status log matches normal updateStatus shape', async () => { + // normal path + const normalDb = buildDb({ + agreements: [ + { + id: 'agr-normal', + status: 'active', + title: 'N', + amount: '10', + asset: 'USDC', + created_by: WALLET, + milestones: [], + }, + ], + auth_users: [{ id: USER, wallet_public_key: WALLET }], + agreement_participants: [ + { agreement_id: 'agr-normal', wallet_address: WALLET, role: 'payer' }, + ], + }); + const normalEmitter = new EventEmitter2(); + const { svc: normalAgreements, logActivity: normalLog } = makeAgreements( + normalDb, + normalEmitter, + ); + await normalAgreements.updateStatus(USER, 'agr-normal', { + actor_wallet: WALLET, + status: 'disputed', + }); + const normalCall = normalLog.mock.calls.find((c) => c[2] === 'status_changed_to_disputed'); + expect(normalCall).toBeTruthy(); + const normalDetails = normalCall![3] as Record; + expect(normalDetails).toMatchObject({ + status: 'disputed', + from: 'active', + to: 'disputed', + }); + + // dispute path + await disputes.openDispute(USER, { + agreement_id: AGREEMENT_ID, + opened_by: WALLET, + reason: 'parity', + }); + const disputeCall = logActivity.mock.calls.find((c) => c[2] === 'status_changed_to_disputed'); + expect(disputeCall).toBeTruthy(); + const disputeDetails = disputeCall![3] as Record; + // same core shape keys as normal flow + expect(disputeDetails).toMatchObject({ + status: 'disputed', + from: 'active', + to: 'disputed', + }); + }); + + it('regression: listByWallet scoping unchanged (creator + participant union)', async () => { + db.tables.agreements.push({ + id: 'agr-other', + status: 'pending', + title: 'Other', + amount: '1', + asset: 'USDC', + created_by: 'GOTHER', + milestones: [], + }); + // wallet is only on AGREEMENT_ID as creator/participant + const { agreements: listed, error } = await agreements.listByWallet(USER, WALLET); + expect(error).toBeNull(); + const ids = listed.map((a: { id: string }) => a.id); + expect(ids).toContain(AGREEMENT_ID); + expect(ids).not.toContain('agr-other'); + }); +}); diff --git a/src/disputes/disputes.service.ts b/src/disputes/disputes.service.ts index 6548b11..237bc1d 100644 --- a/src/disputes/disputes.service.ts +++ b/src/disputes/disputes.service.ts @@ -7,6 +7,7 @@ import { import { EventEmitter2 } from '@nestjs/event-emitter'; import { SupabaseService } from '../supabase/supabase.service'; import { AgreementsService } from '../agreements/agreements.service'; +import { AgreementActivityService } from '../agreements/agreement-activity.service'; import { DISPUTE_OPENED, DISPUTE_RESOLVED } from '../common/constants/notification-events'; import { OpenDisputeDto, @@ -46,6 +47,7 @@ export class DisputesService { private readonly supabase: SupabaseService, private readonly agreements: AgreementsService, private readonly eventEmitter: EventEmitter2, + private readonly activity: AgreementActivityService, ) {} private async walletForUserId(userId: string): Promise { @@ -130,29 +132,27 @@ export class DisputesService { return { dispute: null, error: error.message }; } - // Capture the agreement's status before moving it to `disputed`, so the - // activity entry records the previous/new state of the transition. - const { data: agreementRow } = await this.supabase - .getClient() - .from('agreements') - .select('status') - .eq('id', dto.agreement_id) - .maybeSingle(); - const previousStatus = (agreementRow?.status as string | undefined) ?? null; - - // Update agreement status to disputed - await this.supabase - .getClient() - .from('agreements') - .update({ status: 'disputed', updated_at: new Date().toISOString() }) - .eq('id', dto.agreement_id); + // Route Agreement status change through the shared side-effect path + // (activity log + domain events) so dispute-driven updates match normal updates. + const statusResult = await this.agreements.applyStatusChange( + dto.agreement_id, + dto.opened_by, + 'disputed', + { + activityDetails: { dispute_id: dispute.id, reason: dto.reason, source: 'dispute' }, + }, + ); + if (!statusResult.success) { + return { dispute: null, error: statusResult.error || 'Failed to mark agreement disputed' }; + } - await this.agreements.logAgreementActivity( + // Dispute-specific activity entry (parity with existing audit vocabulary) + await this.activity.logActivity( dto.agreement_id, dto.opened_by, 'dispute_opened', { dispute_id: dispute.id, reason: dto.reason }, - { previousState: previousStatus, newState: 'disputed' }, + { previousState: statusResult.fromStatus ?? null, newState: 'disputed' }, ); this.eventEmitter.emit(DISPUTE_OPENED, { @@ -197,11 +197,14 @@ export class DisputesService { return { success: false, error: error.message }; } - await this.agreements.logAgreementActivity( + await this.activity.logActivity( dispute.agreement_id, dto.resolver_wallet, 'dispute_resolver_assigned', - { dispute_id: disputeId, resolver_wallet: dto.resolver_wallet }, + { + dispute_id: disputeId, + resolver_wallet: dto.resolver_wallet, + }, ); return { success: true, error: null }; @@ -263,27 +266,25 @@ export class DisputesService { }) .eq('id', disputeId); - // Capture the agreement's status before moving it to `resolved`. - const { data: agreementRow } = await this.supabase - .getClient() - .from('agreements') - .select('status') - .eq('id', dispute.agreement_id) - .maybeSingle(); - const previousStatus = (agreementRow?.status as string | undefined) ?? 'disputed'; - - // Update agreement status - await this.supabase - .getClient() - .from('agreements') - .update({ - status: 'resolved', - completed_at: new Date().toISOString(), - updated_at: new Date().toISOString(), - }) - .eq('id', dispute.agreement_id); + // Shared side-effect path for Agreement status → resolved + const statusResult = await this.agreements.applyStatusChange( + dispute.agreement_id, + dto.resolved_by, + 'resolved', + { + activityDetails: { + dispute_id: disputeId, + payer_percentage: dto.payer_percentage, + payee_percentage: dto.payee_percentage, + source: 'dispute', + }, + }, + ); + if (!statusResult.success) { + return { resolution: null, error: statusResult.error || 'Failed to mark agreement resolved' }; + } - await this.agreements.logAgreementActivity( + await this.activity.logActivity( dispute.agreement_id, dto.resolved_by, 'dispute_resolved', @@ -293,7 +294,7 @@ export class DisputesService { payee_percentage: dto.payee_percentage, resolution_notes: dto.resolution_notes, }, - { previousState: previousStatus, newState: 'resolved' }, + { previousState: statusResult.fromStatus ?? 'disputed', newState: 'resolved' }, ); this.eventEmitter.emit(DISPUTE_RESOLVED, { @@ -343,19 +344,25 @@ export class DisputesService { return { success: false, error: error.message }; } - // Revert agreement status to active - await this.supabase - .getClient() - .from('agreements') - .update({ status: 'active', updated_at: new Date().toISOString() }) - .eq('id', dispute.agreement_id); + // Shared side-effect path: dispute cancelled → agreement back to active + const statusResult = await this.agreements.applyStatusChange( + dispute.agreement_id, + dto.cancelled_by, + 'active', + { + activityDetails: { dispute_id: disputeId, source: 'dispute_cancel' }, + }, + ); + if (!statusResult.success) { + return { success: false, error: statusResult.error || 'Failed to restore agreement status' }; + } - await this.agreements.logAgreementActivity( + await this.activity.logActivity( dispute.agreement_id, dto.cancelled_by, 'dispute_cancelled', { dispute_id: disputeId }, - { previousState: 'disputed', newState: 'active' }, + { previousState: statusResult.fromStatus ?? 'disputed', newState: 'active' }, ); return { success: true, error: null }; diff --git a/src/integration/migrated-flows.integration.spec.ts b/src/integration/migrated-flows.integration.spec.ts index 55289fe..7c8805f 100644 --- a/src/integration/migrated-flows.integration.spec.ts +++ b/src/integration/migrated-flows.integration.spec.ts @@ -12,6 +12,7 @@ import { AuthModule } from '../auth/auth.module'; import { SupabaseService } from '../supabase/supabase.service'; import { ApiClient } from '../common/api/api-client'; import { AgreementsController } from '../agreements/agreements.controller'; +import { AgreementActivityService } from '../agreements/agreement-activity.service'; import { AgreementsService } from '../agreements/agreements.service'; import { DisputesController } from '../disputes/disputes.controller'; import { DisputesService } from '../disputes/disputes.service'; @@ -340,6 +341,7 @@ describe('migrated backend flows (integration)', () => { controllers: [AgreementsController, DisputesController, EscrowsController, WalletsController], providers: [ AgreementsService, + AgreementActivityService, DisputesService, WalletsService, { provide: SupabaseService, useValue: supabase }, diff --git a/src/webhooks/webhooks.module.ts b/src/webhooks/webhooks.module.ts index 37bf7d7..a4f3774 100644 --- a/src/webhooks/webhooks.module.ts +++ b/src/webhooks/webhooks.module.ts @@ -2,9 +2,10 @@ import { Module } from '@nestjs/common'; import { WebhooksController } from './webhooks.controller'; import { WebhooksService } from './webhooks.service'; import { NotificationsModule } from '../notifications/notifications.module'; +import { AgreementsModule } from '../agreements/agreements.module'; @Module({ - imports: [NotificationsModule], + imports: [NotificationsModule, AgreementsModule], controllers: [WebhooksController], providers: [WebhooksService], }) diff --git a/src/webhooks/webhooks.service.ts b/src/webhooks/webhooks.service.ts index 481d244..7d14079 100644 --- a/src/webhooks/webhooks.service.ts +++ b/src/webhooks/webhooks.service.ts @@ -4,6 +4,7 @@ import { EventEmitter2 } from '@nestjs/event-emitter'; import * as crypto from 'crypto'; import { SupabaseService } from '../supabase/supabase.service'; import { NotificationsService } from '../notifications/notifications.service'; +import { AgreementActivityService } from '../agreements/agreement-activity.service'; import { AGREEMENT_EVENTS } from '../common/events/agreement-events.constants'; import type { TrustlessWorkEventDto } from './dto/trustless-work-event.dto'; @@ -40,6 +41,7 @@ export class WebhooksService { private readonly eventEmitter: EventEmitter2, private readonly notifications: NotificationsService, private readonly config: ConfigService, + private readonly activity: AgreementActivityService, ) { this.webhookSecret = this.config.get('TRUSTLESS_WORK_WEBHOOK_SECRET', ''); } @@ -172,7 +174,7 @@ export class WebhooksService { const row = updated; - await this.logActivity( + await this.activity.logActivity( row.id, 'trustless-work-webhook', `webhook_status_changed_to_${targetStatus}`, @@ -237,11 +239,16 @@ export class WebhooksService { throw updateError; } - await this.logActivity(agreement.id, 'trustless-work-webhook', 'webhook_milestone_updated', { - event: payload.event, - contractId: payload.contractId, - milestone_index: milestoneIndex, - }); + await this.activity.logActivity( + agreement.id, + 'trustless-work-webhook', + 'webhook_milestone_updated', + { + event: payload.event, + contractId: payload.contractId, + milestone_index: milestoneIndex, + }, + ); } private async applyInfoUpdate(payload: TrustlessWorkEventDto): Promise { @@ -257,7 +264,7 @@ export class WebhooksService { return; } - await this.logActivity( + await this.activity.logActivity( agreement.id, 'trustless-work-webhook', `webhook_event_${payload.event.replace('.', '_')}`, @@ -302,22 +309,4 @@ export class WebhooksService { this.logger.error('Notification dispatch error', err); } } - - private async logActivity( - agreementId: string, - actorWallet: string, - action: string, - details: Record = {}, - ) { - try { - await this.supabase.getClient().from('agreement_activity').insert({ - agreement_id: agreementId, - actor_wallet: actorWallet, - action, - details, - }); - } catch (e) { - this.logger.error('logActivity', e); - } - } } diff --git a/src/webhooks/webhooks.spec.ts b/src/webhooks/webhooks.spec.ts index 469a7ca..c1a4313 100644 --- a/src/webhooks/webhooks.spec.ts +++ b/src/webhooks/webhooks.spec.ts @@ -40,11 +40,16 @@ function buildService( key === 'TRUSTLESS_WORK_WEBHOOK_SECRET' ? (deps.secret ?? SECRET) : def, }; + const activity = { + logActivity: jest.fn().mockResolvedValue(undefined), + }; + const svc = new (WebhooksService as unknown as new (...args: unknown[]) => WebhooksService)( supabase, eventEmitter, notifications, config, + activity, ) as WebhooksService & { _emit: jest.Mock; _notifyDispute: jest.Mock }; svc._emit = emit;