diff --git a/apps/backend/src/app.module.ts b/apps/backend/src/app.module.ts index 9f66edf..a8a636d 100644 --- a/apps/backend/src/app.module.ts +++ b/apps/backend/src/app.module.ts @@ -1,4 +1,5 @@ import { TradeAssistantModule } from './application/trade-assistant/trade-assistant.module'; +import { McpModule } from './application/mcp/mcp.module'; import { Module } from '@nestjs/common'; import { ConfigModule, ConfigService } from '@nestjs/config'; import { TypeOrmModule } from '@nestjs/typeorm'; @@ -193,6 +194,7 @@ import { CustomThrottlerGuard } from './application/guards/throttle.guard'; BlogModule, SubscriptionsModule, TradeAssistantModule, + McpModule, ApiKeysModule, LogsModule, ], diff --git a/apps/backend/src/application/mcp/capabilities/account.capabilities.ts b/apps/backend/src/application/mcp/capabilities/account.capabilities.ts new file mode 100644 index 0000000..463d1b8 --- /dev/null +++ b/apps/backend/src/application/mcp/capabilities/account.capabilities.ts @@ -0,0 +1,48 @@ +import { SubscriptionService } from '../../services/subscription.service'; +import { actorPlan } from '@domain/services/capability-access'; +import { Capability } from '../capability'; + +/** + * Compte et abonnement. + * + * `whoami` n'est pas un gadget : c'est ce qui permet a un agent d'annoncer + * honnetement ce qu'il peut faire, au lieu de proposer une action puis de se + * heurter a un refus. + */ +export function accountCapabilities(subscriptions: SubscriptionService): Capability[] { + return [ + { + policy: { name: 'whoami', scope: 'read' }, + description: + "Identité de l'appelant : identifiant, organisation, rôle et offre effective. À appeler en premier pour savoir ce qui est permis.", + inputSchema: { type: 'object', properties: {}, additionalProperties: false }, + handler: async (_input, actor) => ({ + userId: actor.id, + organizationId: actor.organizationId, + role: actor.role, + plan: actorPlan(actor), + }), + }, + + { + policy: { name: 'get_subscription', scope: 'read', roles: ['ADMIN', 'MANAGER'] }, + description: + "Abonnement de l'organisation : offre, statut, licences utilisées et disponibles. Réservé aux rôles ADMIN et MANAGER.", + inputSchema: { type: 'object', properties: {}, additionalProperties: false }, + handler: async (_input, actor) => { + const overview = await subscriptions.getSubscriptionOverview( + actor.organizationId, + actor.role + ); + return { + plan: overview.plan, + status: overview.status, + usedLicenses: overview.usedLicenses, + maxLicenses: overview.maxLicenses, + availableLicenses: overview.availableLicenses, + currentPeriodEnd: overview.currentPeriodEnd, + }; + }, + }, + ]; +} diff --git a/apps/backend/src/application/mcp/capabilities/admin.capabilities.ts b/apps/backend/src/application/mcp/capabilities/admin.capabilities.ts new file mode 100644 index 0000000..b0cb0a4 --- /dev/null +++ b/apps/backend/src/application/mcp/capabilities/admin.capabilities.ts @@ -0,0 +1,126 @@ +import { UserRepository } from '@domain/ports/out/user.repository'; +import { OrganizationRepository } from '@domain/ports/out/organization.repository'; +import { CsvRateSearchService } from '@domain/services/csv-rate-search.service'; +import { Capability } from '../capability'; + +/** Role d'administration de la plateforme. Les inscriptions creent des MANAGER. */ +const ADMIN_ONLY = ['ADMIN'] as const; + +/** + * Capacites d'administration de la plateforme. + * + * Elles franchissent la frontiere de l'organisation — c'est precisement ce qui + * les distingue du reste du catalogue — et sont donc reservees au role ADMIN, + * verifie dans le processus et non dans un prompt. + * + * Elles sont **en lecture seule**. Modifier un utilisateur, valider un SIRET ou + * remplacer une grille tarifaire touche des comptes clients et de l'argent : + * ces actions restent a la main d'une personne, dans l'espace d'administration, + * tant qu'un mecanisme de confirmation explicite n'existe pas cote agent. + */ +export function adminCapabilities( + users: UserRepository, + organizations: OrganizationRepository, + rateSearch: CsvRateSearchService +): Capability[] { + return [ + { + policy: { name: 'admin_list_users', scope: 'read', roles: ADMIN_ONLY }, + description: + "Liste les comptes de la plateforme, toutes organisations confondues. Réservé à l'administration.", + inputSchema: { + type: 'object', + properties: { + role: { + type: 'string', + description: 'Ne garder que ce rôle.', + enum: ['ADMIN', 'MANAGER', 'USER', 'VIEWER', 'CARRIER'], + }, + search: { + type: 'string', + description: 'Filtre sur l’adresse e-mail ou le nom.', + maxLength: 120, + }, + limit: { + type: 'integer', + description: 'Nombre maximum de comptes.', + minimum: 1, + maximum: 100, + default: 25, + }, + }, + additionalProperties: false, + }, + handler: async input => { + const all = input.role + ? await users.findByRole(input.role as string) + : await users.findAll(); + + const term = (input.search as string | undefined)?.toLowerCase(); + const matching = term + ? all.filter(user => + `${user.email} ${user.firstName} ${user.lastName}`.toLowerCase().includes(term) + ) + : all; + + return { + total: matching.length, + users: matching.slice(0, (input.limit as number) ?? 25).map(user => ({ + id: user.id, + email: user.email, + firstName: user.firstName, + lastName: user.lastName, + role: user.role, + organizationId: user.organizationId, + isActive: user.isActive, + })), + }; + }, + }, + + { + policy: { name: 'admin_list_organizations', scope: 'read', roles: ADMIN_ONLY }, + description: + "Liste les organisations de la plateforme, avec leur nombre de comptes. Réservé à l'administration.", + inputSchema: { + type: 'object', + properties: { + limit: { + type: 'integer', + description: "Nombre maximum d'organisations.", + minimum: 1, + maximum: 100, + default: 25, + }, + }, + additionalProperties: false, + }, + handler: async input => { + const all = await organizations.findAll(); + const page = all.slice(0, (input.limit as number) ?? 25); + + return { + total: all.length, + organizations: await Promise.all( + page.map(async organization => ({ + id: organization.id, + name: organization.name, + userCount: await users.countByOrganization(organization.id), + })) + ), + }; + }, + }, + + { + policy: { name: 'admin_rate_grid_overview', scope: 'read', roles: ADMIN_ONLY }, + description: + "État des grilles tarifaires chargées : transporteurs et types de conteneurs disponibles. Réservé à l'administration.", + inputSchema: { type: 'object', properties: {}, additionalProperties: false }, + handler: async () => ({ + carriers: await rateSearch.getAvailableCompanies(), + containerTypes: await rateSearch.getAvailableContainerTypes(), + }), + }, + ]; +} diff --git a/apps/backend/src/application/mcp/capabilities/bookings.capabilities.ts b/apps/backend/src/application/mcp/capabilities/bookings.capabilities.ts new file mode 100644 index 0000000..fe6d2ec --- /dev/null +++ b/apps/backend/src/application/mcp/capabilities/bookings.capabilities.ts @@ -0,0 +1,113 @@ +import { CsvBookingService } from '../../services/csv-booking.service'; +import { Capability } from '../capability'; + +/** + * Reservations. + * + * Les lectures restent cantonnees a l'appelant, sauf `list_organization_bookings` + * qui demande un role d'encadrement — c'est la meme frontiere que dans + * l'interface, ou seuls ADMIN et MANAGER voient l'onglet organisation. + * + * Les ecritures sont volontairement limitees aux actions reversibles ou + * inoffensives : annuler, et supprimer une reservation impayee. Payer une + * commission, envoyer une demande a un transporteur ou televerser un document + * engagent un tiers ou de l'argent et restent hors de portee d'un agent. + */ +export function bookingsCapabilities(bookings: CsvBookingService): Capability[] { + return [ + { + policy: { name: 'list_my_bookings', scope: 'read' }, + description: + "Liste les réservations de l'utilisateur authentifié, de la plus récente à la plus ancienne.", + inputSchema: { + type: 'object', + properties: { + limit: { + type: 'integer', + description: 'Nombre maximum de réservations.', + minimum: 1, + maximum: 50, + default: 20, + }, + }, + additionalProperties: false, + }, + handler: async (input, actor) => + bookings.getUserBookings(actor.id, 1, (input.limit as number) ?? 20), + }, + + { + policy: { name: 'list_organization_bookings', scope: 'read', roles: ['ADMIN', 'MANAGER'] }, + description: + "Liste les réservations de toute l'organisation. Réservé aux rôles ADMIN et MANAGER.", + inputSchema: { + type: 'object', + properties: { + limit: { + type: 'integer', + description: 'Nombre maximum de réservations.', + minimum: 1, + maximum: 50, + default: 20, + }, + }, + additionalProperties: false, + }, + handler: async (input, actor) => + bookings.getOrganizationBookings(actor.organizationId, 1, (input.limit as number) ?? 20), + }, + + { + policy: { name: 'get_booking', scope: 'read' }, + description: + "Détail d'une réservation : route, marchandise, transporteur, statut, documents.", + inputSchema: { + type: 'object', + properties: { + bookingId: { type: 'string', description: 'Identifiant de la réservation (UUID).' }, + }, + required: ['bookingId'], + additionalProperties: false, + }, + handler: async (input, actor) => bookings.getBookingById(input.bookingId as string, actor.id), + }, + + { + policy: { name: 'booking_statistics', scope: 'read' }, + description: + "Répartition des réservations de l'utilisateur par statut (en attente de paiement, en attente, acceptées, refusées).", + inputSchema: { type: 'object', properties: {}, additionalProperties: false }, + handler: async (_input, actor) => bookings.getUserStats(actor.id), + }, + + { + policy: { name: 'cancel_booking', scope: 'write' }, + description: + "Annule une réservation de l'utilisateur qui n'a pas encore été acceptée par le transporteur. La réservation est conservée avec le statut annulé.", + inputSchema: { + type: 'object', + properties: { + bookingId: { type: 'string', description: 'Identifiant de la réservation (UUID).' }, + }, + required: ['bookingId'], + additionalProperties: false, + }, + handler: async (input, actor) => bookings.cancelBooking(input.bookingId as string, actor.id), + }, + + { + policy: { name: 'delete_unpaid_booking', scope: 'write' }, + description: + "Supprime définitivement une réservation dont la commission n'a pas été payée. Sans effet sur une réservation payée, qui ne peut être qu'annulée.", + inputSchema: { + type: 'object', + properties: { + bookingId: { type: 'string', description: 'Identifiant de la réservation (UUID).' }, + }, + required: ['bookingId'], + additionalProperties: false, + }, + handler: async (input, actor) => bookings.deleteBooking(input.bookingId as string, actor.id), + }, + ]; +} diff --git a/apps/backend/src/application/mcp/capabilities/knowledge.capabilities.ts b/apps/backend/src/application/mcp/capabilities/knowledge.capabilities.ts new file mode 100644 index 0000000..56bfdb4 --- /dev/null +++ b/apps/backend/src/application/mcp/capabilities/knowledge.capabilities.ts @@ -0,0 +1,61 @@ +import { TradeRetrievalPort } from '@domain/ports/out/trade-assistant.port'; +import { Capability } from '../capability'; + +/** + * Documentation du site, exposee comme capacite. + * + * Le meme index que l'assistant integre : un agent externe repond donc a partir + * du wiki Xpeditis, avec les liens vers les pages, plutot que de ses propres + * souvenirs sur le fret maritime. + */ +export function knowledgeCapabilities(retrieval: TradeRetrievalPort): Capability[] { + return [ + { + policy: { name: 'search_documentation', scope: 'read' }, + description: + "Recherche dans le wiki Xpeditis (Incoterms, douanes, conteneurs, IMDG, VGM, calcul du fret, transit times). Renvoie les extraits pertinents et le lien de la page d'origine.", + inputSchema: { + type: 'object', + properties: { + query: { + type: 'string', + description: 'La question ou les mots-clés à rechercher.', + minLength: 2, + maxLength: 500, + }, + language: { + type: 'string', + description: 'Langue de la documentation.', + enum: ['fr', 'en'], + default: 'fr', + }, + limit: { + type: 'integer', + description: "Nombre maximum d'extraits.", + minimum: 1, + maximum: 10, + default: 4, + }, + }, + required: ['query'], + additionalProperties: false, + }, + handler: async input => { + const passages = await retrieval.search( + input.query as string, + (input.language as string) ?? 'fr', + input.limit as number + ); + return { + matches: passages.map(passage => ({ + title: passage.title, + section: passage.section, + url: passage.href, + excerpt: passage.text, + score: passage.score, + })), + }; + }, + }, + ]; +} diff --git a/apps/backend/src/application/mcp/capabilities/rates.capabilities.ts b/apps/backend/src/application/mcp/capabilities/rates.capabilities.ts new file mode 100644 index 0000000..1c20a3b --- /dev/null +++ b/apps/backend/src/application/mcp/capabilities/rates.capabilities.ts @@ -0,0 +1,111 @@ +import { CsvRateSearchService } from '@domain/services/csv-rate-search.service'; +import { RateDirection } from '@domain/entities/csv-rate.entity'; +import { Capability } from '../capability'; + +/** Au-dela, la reponse devient illisible pour un agent et couteuse en contexte. */ +const MAX_RESULTS = 10; + +/** + * Recherche tarifaire LCL — le coeur du produit. + * + * Le resultat est resume : un agent a besoin du transporteur, du delai et du + * total, pas de la structure complete des surcharges. Le detail reste + * accessible dans l'application, dont le lien est renvoye. + */ +export function ratesCapabilities(search: CsvRateSearchService): Capability[] { + return [ + { + policy: { name: 'search_rates', scope: 'read' }, + description: + 'Recherche des tarifs de fret maritime LCL entre deux ports (codes UN/LOCODE, ex. FRLIO, CNSHA). Renvoie les offres disponibles avec transporteur, temps de transit et prix.', + inputSchema: { + type: 'object', + properties: { + origin: { + type: 'string', + description: 'Port de départ, code UN/LOCODE à 5 lettres (ex. FRLIO).', + minLength: 5, + maxLength: 5, + }, + destination: { + type: 'string', + description: "Port d'arrivée, code UN/LOCODE à 5 lettres (ex. CNSHA).", + minLength: 5, + maxLength: 5, + }, + volumeCBM: { + type: 'number', + description: 'Volume de la marchandise en mètres cubes.', + minimum: 0.01, + maximum: 1000, + }, + weightKG: { + type: 'number', + description: 'Poids brut de la marchandise en kilogrammes.', + minimum: 1, + maximum: 1000000, + }, + direction: { + type: 'string', + description: 'Sens de la grille tarifaire.', + enum: ['EXPORT', 'IMPORT'], + }, + hasDangerousGoods: { + type: 'boolean', + description: 'La marchandise relève-t-elle de la réglementation IMDG ?', + default: false, + }, + limit: { + type: 'integer', + description: "Nombre maximum d'offres renvoyées.", + minimum: 1, + maximum: MAX_RESULTS, + default: 5, + }, + }, + required: ['origin', 'destination', 'volumeCBM', 'weightKG'], + additionalProperties: false, + }, + handler: async input => { + const output = await search.execute({ + origin: (input.origin as string).toUpperCase(), + destination: (input.destination as string).toUpperCase(), + volumeCBM: input.volumeCBM as number, + weightKG: input.weightKG as number, + hasDangerousGoods: (input.hasDangerousGoods as boolean) ?? false, + direction: input.direction as RateDirection | undefined, + }); + + const limit = (input.limit as number) ?? 5; + return { + totalResults: output.totalResults, + offers: output.results.slice(0, limit).map(({ rate, priceBreakdown }) => ({ + carrier: rate.companyName, + route: `${rate.originCode.toString()} → ${rate.destinationCode.toString()}`, + routing: rate.routing, + transitDays: rate.transitDays, + frequency: rate.frequency, + freight: { + amount: priceBreakdown.totalFreight, + currency: priceBreakdown.freightCurrency, + }, + destinationCharges: { + amount: priceBreakdown.totalFob, + currency: priceBreakdown.fobCurrency, + }, + dangerousGoods: priceBreakdown.dgSurchargeStatus, + validUntil: rate.validity.getEndDate(), + })), + bookInApp: '/dashboard/search-advanced', + }; + }, + }, + + { + policy: { name: 'list_carriers', scope: 'read' }, + description: 'Liste les transporteurs dont les grilles tarifaires sont chargées.', + inputSchema: { type: 'object', properties: {}, additionalProperties: false }, + handler: async () => ({ carriers: await search.getAvailableCompanies() }), + }, + ]; +} diff --git a/apps/backend/src/application/mcp/capability.registry.spec.ts b/apps/backend/src/application/mcp/capability.registry.spec.ts new file mode 100644 index 0000000..d81f652 --- /dev/null +++ b/apps/backend/src/application/mcp/capability.registry.spec.ts @@ -0,0 +1,237 @@ +import { ForbiddenException, NotFoundException } from '@nestjs/common'; +import { CapabilityActor } from '@domain/services/capability-access'; +import { CapabilityRegistry } from './capability.registry'; +import { CapabilityInputError } from './capability'; + +/** + * Le registre est construit avec les vraies capacites : le test verifie donc le + * catalogue reellement expose, pas un catalogue de laboratoire. + */ +const retrieval = { search: jest.fn().mockResolvedValue([]) }; +const rateSearch = { + execute: jest.fn().mockResolvedValue({ totalResults: 0, results: [] }), + getAvailableCompanies: jest.fn().mockResolvedValue(['CMA CGM']), +}; +const bookings = { + getUserBookings: jest.fn().mockResolvedValue({ bookings: [], total: 0 }), + getOrganizationBookings: jest.fn().mockResolvedValue({ bookings: [], total: 0 }), + getBookingById: jest.fn().mockResolvedValue({ id: 'b1' }), + getUserStats: jest.fn().mockResolvedValue({ pending: 0 }), + cancelBooking: jest.fn().mockResolvedValue({ id: 'b1' }), + deleteBooking: jest.fn().mockResolvedValue({ success: true }), +}; +const subscriptions = { + getSubscriptionOverview: jest.fn().mockResolvedValue({ plan: 'GOLD', status: 'ACTIVE' }), +}; +const usersRepo = { + findAll: jest.fn().mockResolvedValue([]), + findByRole: jest.fn().mockResolvedValue([]), + countByOrganization: jest.fn().mockResolvedValue(0), +}; +const organizationsRepo = { findAll: jest.fn().mockResolvedValue([]) }; +const audit = { log: jest.fn().mockResolvedValue(undefined) }; + +/** Les doubles portent leurs propres types ; seul le passage au registre est force. */ +const build = () => + new CapabilityRegistry( + retrieval as never, + rateSearch as never, + bookings as never, + subscriptions as never, + usersRepo as never, + organizationsRepo as never, + audit as never + ); + +const registry = build(); + +const actor = (overrides: Partial = {}): CapabilityActor => ({ + id: 'u1', + organizationId: 'o1', + role: 'USER', + plan: 'BRONZE', + ...overrides, +}); + +const names = (a: CapabilityActor) => registry.listFor(a).map(c => c.policy.name); + +describe('CapabilityRegistry', () => { + beforeEach(() => jest.clearAllMocks()); + + it('exposes every capability under a unique name', () => { + const all = names(actor({ role: 'ADMIN' })); + expect(new Set(all).size).toBe(all.length); + expect(all).toEqual(expect.arrayContaining(['whoami', 'search_rates', 'list_my_bookings'])); + }); + + it('hides organisation-wide capabilities from a plain user', () => { + expect(names(actor({ role: 'USER' }))).not.toContain('list_organization_bookings'); + expect(names(actor({ role: 'USER' }))).not.toContain('get_subscription'); + expect(names(actor({ role: 'MANAGER' }))).toContain('list_organization_bookings'); + }); + + it('answers "unknown" for a capability the caller has no role for', async () => { + // Ne pas distinguer « interdit » de « inexistant » : sinon la liste filtree + // ne sert a rien, il suffirait de deviner les noms. + await expect( + registry.invoke('list_organization_bookings', {}, actor({ role: 'USER' })) + ).rejects.toThrow(NotFoundException); + }); + + it('scopes every read to the caller, never to a requested identity', async () => { + await registry.invoke('list_my_bookings', { limit: 5 }, actor({ id: 'u42' })); + + expect(bookings.getUserBookings).toHaveBeenCalledWith('u42', 1, 5); + }); + + it('routes organisation reads to the caller organisation', async () => { + await registry.invoke( + 'list_organization_bookings', + {}, + actor({ role: 'ADMIN', organizationId: 'o9' }) + ); + + expect(bookings.getOrganizationBookings).toHaveBeenCalledWith('o9', 1, 20); + }); + + it('validates arguments before touching a service', async () => { + await expect( + registry.invoke('search_rates', { origin: 'FR', destination: 'CNSHA' }, actor()) + ).rejects.toThrow(CapabilityInputError); + + expect(rateSearch.execute).not.toHaveBeenCalled(); + }); + + it('normalises port codes to upper case before searching', async () => { + await registry.invoke( + 'search_rates', + { origin: 'frlio', destination: 'cnsha', volumeCBM: 4, weightKG: 500 }, + actor() + ); + + expect(rateSearch.execute).toHaveBeenCalledWith( + expect.objectContaining({ origin: 'FRLIO', destination: 'CNSHA', volumeCBM: 4 }) + ); + }); + + it('reports the effective plan through whoami', async () => { + await expect(registry.invoke('whoami', {}, actor({ role: 'ADMIN' }))).resolves.toMatchObject({ + role: 'ADMIN', + plan: 'PLATINIUM', + }); + }); + + it('marks read capabilities as read-only and writes as not', () => { + const catalogue = registry.listFor(actor({ role: 'ADMIN' })); + const scopeOf = (name: string) => catalogue.find(c => c.policy.name === name)?.policy.scope; + + expect(scopeOf('search_rates')).toBe('read'); + expect(scopeOf('delete_unpaid_booking')).toBe('write'); + }); + + it('lets a delete reach the service, which enforces the unpaid rule', async () => { + await registry.invoke('delete_unpaid_booking', { bookingId: 'b1' }, actor({ id: 'u1' })); + + // Le registre ne redecide pas la regle metier : il transmet l'appelant et + // laisse le service refuser une reservation payee. + expect(bookings.deleteBooking).toHaveBeenCalledWith('b1', 'u1'); + }); +}); + +describe('journal des appels', () => { + beforeEach(() => jest.clearAllMocks()); + + it('records a successful invocation with its surface and scope', async () => { + await registry.invoke('whoami', {}, actor({ email: 'd@x.com' }), 'assistant'); + + expect(audit.log).toHaveBeenCalledWith( + expect.objectContaining({ + action: 'agent_capability_invoked', + status: 'success', + userId: 'u1', + userEmail: 'd@x.com', + resourceName: 'whoami', + metadata: expect.objectContaining({ surface: 'assistant', scope: 'read' }), + }) + ); + }); + + it('records a refusal too — a repeated attempt is what a journal must reveal', async () => { + await expect( + registry.invoke('admin_list_users', {}, actor({ role: 'USER' }), 'mcp') + ).rejects.toThrow(); + + expect(audit.log).toHaveBeenCalledWith( + expect.objectContaining({ status: 'failure', resourceName: 'admin_list_users' }) + ); + }); + + it('records a handler failure with its message', async () => { + bookings.getUserStats.mockRejectedValueOnce(new Error('database unavailable')); + + await expect(registry.invoke('booking_statistics', {}, actor())).rejects.toThrow(); + + expect(audit.log).toHaveBeenCalledWith( + expect.objectContaining({ status: 'failure', errorMessage: 'database unavailable' }) + ); + }); + + it('defaults the surface to mcp when the caller does not say', async () => { + await registry.invoke('whoami', {}, actor()); + + expect(audit.log.mock.calls[0][0].metadata).toMatchObject({ surface: 'mcp' }); + }); +}); + +describe('capacites d administration', () => { + it('are visible to an ADMIN only', () => { + const adminNames = names(actor({ role: 'ADMIN' })); + expect(adminNames).toEqual( + expect.arrayContaining([ + 'admin_list_users', + 'admin_list_organizations', + 'admin_rate_grid_overview', + ]) + ); + + for (const role of ['MANAGER', 'USER', 'VIEWER']) { + expect(names(actor({ role }))).not.toContain('admin_list_users'); + } + }); + + it('stay read-only until an explicit confirmation mechanism exists', () => { + const adminOnes = registry + .listFor(actor({ role: 'ADMIN' })) + .filter(c => c.policy.name.startsWith('admin_')); + + expect(adminOnes.length).toBeGreaterThan(0); + expect(adminOnes.every(c => c.policy.scope === 'read')).toBe(true); + }); +}); + +describe('plan-gated capabilities', () => { + /** Capacite fictive soumise a une fonctionnalite d'offre. */ + const gated = build(); + beforeAll(() => { + (gated as unknown as { capabilities: unknown[] }).capabilities = [ + { + policy: { name: 'export_everything', scope: 'read', feature: 'api_access' }, + description: '', + inputSchema: { type: 'object', properties: {}, additionalProperties: false }, + handler: async () => ({ ok: true }), + }, + ]; + }); + + it('says plainly that the plan is missing, because the feature can be bought', async () => { + await expect(gated.invoke('export_everything', {}, actor({ plan: 'SILVER' }))).rejects.toThrow( + ForbiddenException + ); + }); + + it('allows it once the plan includes the feature', async () => { + await expect(gated.invoke('export_everything', {}, actor({ plan: 'GOLD' }))).resolves.toEqual({ + ok: true, + }); + }); +}); diff --git a/apps/backend/src/application/mcp/capability.registry.ts b/apps/backend/src/application/mcp/capability.registry.ts new file mode 100644 index 0000000..64bb703 --- /dev/null +++ b/apps/backend/src/application/mcp/capability.registry.ts @@ -0,0 +1,156 @@ +import { ForbiddenException, Inject, Injectable, NotFoundException } from '@nestjs/common'; +import { TRADE_RETRIEVAL, TradeRetrievalPort } from '@domain/ports/out/trade-assistant.port'; +import { CsvRateSearchService } from '@domain/services/csv-rate-search.service'; +import { + CapabilityActor, + denialReason, + grantedCapabilities, +} from '@domain/services/capability-access'; +import { AuditAction, AuditStatus } from '@domain/entities/audit-log.entity'; +import { USER_REPOSITORY, UserRepository } from '@domain/ports/out/user.repository'; +import { + ORGANIZATION_REPOSITORY, + OrganizationRepository, +} from '@domain/ports/out/organization.repository'; +import { AuditService } from '../services/audit.service'; +import { CsvBookingService } from '../services/csv-booking.service'; +import { SubscriptionService } from '../services/subscription.service'; +import { Capability, CapabilityInputError, parseInput } from './capability'; +import { accountCapabilities } from './capabilities/account.capabilities'; +import { adminCapabilities } from './capabilities/admin.capabilities'; +import { bookingsCapabilities } from './capabilities/bookings.capabilities'; +import { knowledgeCapabilities } from './capabilities/knowledge.capabilities'; +import { ratesCapabilities } from './capabilities/rates.capabilities'; + +/** + * Catalogue des capacites du produit. + * + * Un seul endroit declare ce qu'un agent peut faire et sous quelles conditions. + * Les deux consommateurs — le serveur MCP pour les clients externes, l'assistant + * integre pour l'appel de fonctions — lisent ce meme catalogue : une capacite + * ajoutee ici devient disponible des deux cotes, avec les memes droits, sans + * qu'aucune des deux surfaces n'ait a etre modifiee. + * + * Ce registre n'implemente rien : chaque capacite delegue au service applicatif + * qui sert deja l'interface. Le produit n'a pas de seconde logique metier pour + * les agents, donc pas de seconde verite a maintenir. + */ +@Injectable() +export class CapabilityRegistry { + private readonly capabilities: Capability[]; + + constructor( + @Inject(TRADE_RETRIEVAL) retrieval: TradeRetrievalPort, + rateSearch: CsvRateSearchService, + bookings: CsvBookingService, + subscriptions: SubscriptionService, + @Inject(USER_REPOSITORY) users: UserRepository, + @Inject(ORGANIZATION_REPOSITORY) organizations: OrganizationRepository, + private readonly audit: AuditService + ) { + this.capabilities = [ + ...accountCapabilities(subscriptions), + ...knowledgeCapabilities(retrieval), + ...ratesCapabilities(rateSearch), + ...bookingsCapabilities(bookings), + ...adminCapabilities(users, organizations, rateSearch), + ]; + + const duplicate = findDuplicate(this.capabilities.map(c => c.policy.name)); + if (duplicate) { + // Deux capacites homonymes rendraient l'appel ambigu : mieux vaut + // empecher le demarrage que resoudre au hasard. + throw new Error(`Duplicate capability name: ${duplicate}`); + } + } + + /** Capacites visibles par cet appelant, dans l'ordre du catalogue. */ + listFor(actor: CapabilityActor): Capability[] { + return grantedCapabilities(actor, this.capabilities); + } + + /** + * Execute une capacite au nom de l'appelant. + * + * Une capacite hors droits repond « inconnue », comme si elle n'existait pas : + * la liste ne l'expose deja pas, et distinguer les deux cas revelerait + * l'existence de fonctions reservees. + */ + async invoke( + name: string, + rawInput: Record | undefined, + actor: CapabilityActor, + surface: CapabilitySurface = 'mcp' + ): Promise { + const capability = this.capabilities.find(c => c.policy.name === name); + + if (!capability || denialReason(actor, capability.policy) === 'role') { + await this.record(actor, surface, name, 'read', false, 'Unknown or forbidden capability'); + throw new NotFoundException(`Unknown capability "${name}".`); + } + + // Un refus lie a l'offre se dit, lui : la fonction existe, elle s'achete. + if (denialReason(actor, capability.policy) === 'plan') { + const message = `Capability "${name}" requires the "${capability.policy.feature}" feature, not included in your plan.`; + await this.record(actor, surface, name, capability.policy.scope, false, message); + throw new ForbiddenException(message); + } + + try { + const input = parseInput(capability.inputSchema, rawInput); + const result = await capability.handler(input, actor); + await this.record(actor, surface, name, capability.policy.scope, true); + return result; + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + await this.record(actor, surface, name, capability.policy.scope, false, message); + throw error; + } + } + + /** + * Journalise l'appel. + * + * L'audit est pose ici, et non dans chaque adaptateur : le registre est le + * seul passage oblige des deux surfaces, donc le seul endroit ou la trace ne + * peut pas etre oubliee en ajoutant une capacite. Les refus sont journalises + * autant que les succes — c'est ce qui revele une tentative repetee. + * + * `AuditService.log` n'echoue jamais vers l'appelant : une panne du journal + * ne doit pas empecher une action deja autorisee. + */ + private async record( + actor: CapabilityActor, + surface: CapabilitySurface, + name: string, + scope: string, + success: boolean, + errorMessage?: string + ): Promise { + await this.audit.log({ + action: AuditAction.AGENT_CAPABILITY_INVOKED, + status: success ? AuditStatus.SUCCESS : AuditStatus.FAILURE, + userId: actor.id, + userEmail: actor.email ?? 'unknown', + organizationId: actor.organizationId, + resourceType: 'capability', + resourceName: name, + metadata: { surface, scope, role: actor.role }, + ...(errorMessage ? { errorMessage } : {}), + }); + } +} + +/** D'ou vient l'appel : sert a distinguer les usages dans le journal. */ +export type CapabilitySurface = 'mcp' | 'assistant'; + +export { CapabilityInputError }; + +function findDuplicate(names: string[]): string | null { + const seen = new Set(); + for (const name of names) { + if (seen.has(name)) return name; + seen.add(name); + } + return null; +} diff --git a/apps/backend/src/application/mcp/capability.spec.ts b/apps/backend/src/application/mcp/capability.spec.ts new file mode 100644 index 0000000..8a758b1 --- /dev/null +++ b/apps/backend/src/application/mcp/capability.spec.ts @@ -0,0 +1,84 @@ +import { CapabilityInputError, CapabilitySchema, parseInput } from './capability'; + +const schema: CapabilitySchema = { + type: 'object', + properties: { + origin: { type: 'string', description: '', minLength: 5, maxLength: 5 }, + volumeCBM: { type: 'number', description: '', minimum: 0.01, maximum: 1000 }, + limit: { type: 'integer', description: '', minimum: 1, maximum: 10, default: 5 }, + direction: { type: 'string', description: '', enum: ['EXPORT', 'IMPORT'] }, + dangerous: { type: 'boolean', description: '', default: false }, + companies: { type: 'array', description: '', items: { type: 'string' } }, + }, + required: ['origin', 'volumeCBM'], + additionalProperties: false, +}; + +describe('parseInput', () => { + it('accepts a well-formed call and applies defaults', () => { + expect(parseInput(schema, { origin: 'FRLIO', volumeCBM: 4 })).toEqual({ + origin: 'FRLIO', + volumeCBM: 4, + limit: 5, + dangerous: false, + }); + }); + + it('coerces the numeric strings a language model tends to produce', () => { + const parsed = parseInput(schema, { origin: ' FRLIO ', volumeCBM: '4.5', limit: '3' }); + expect(parsed).toMatchObject({ origin: 'FRLIO', volumeCBM: 4.5, limit: 3 }); + }); + + it('treats null and empty string as absent', () => { + expect( + parseInput(schema, { origin: 'FRLIO', volumeCBM: 4, direction: null }) + ).not.toHaveProperty('direction'); + expect(parseInput(schema, { origin: 'FRLIO', volumeCBM: 4, direction: '' })).not.toHaveProperty( + 'direction' + ); + }); + + it.each([ + [{ volumeCBM: 4 }, 'Missing required parameter "origin"'], + [{ origin: 'FRLIO' }, 'Missing required parameter "volumeCBM"'], + [{ origin: 'FR', volumeCBM: 4 }, 'at least 5 characters'], + [{ origin: 'FRLIO', volumeCBM: 'beaucoup' }, 'must be a number'], + [{ origin: 'FRLIO', volumeCBM: 0 }, 'must be >= 0.01'], + [{ origin: 'FRLIO', volumeCBM: 4, limit: 2.5 }, 'whole number'], + [{ origin: 'FRLIO', volumeCBM: 4, limit: 99 }, 'must be <= 10'], + [{ origin: 'FRLIO', volumeCBM: 4, direction: 'BOTH' }, 'must be one of: EXPORT, IMPORT'], + [{ origin: 'FRLIO', volumeCBM: 4, dangerous: 'peut-être' }, 'must be true or false'], + [{ origin: 'FRLIO', volumeCBM: 4, companies: 'CMA' }, 'must be an array'], + ])('rejects %p', (input, message) => { + expect(() => parseInput(schema, input as Record)).toThrow( + CapabilityInputError + ); + expect(() => parseInput(schema, input as Record)).toThrow( + expect.objectContaining({ message: expect.stringContaining(message) }) + ); + }); + + it('refuses an invented parameter instead of passing it through', () => { + // Un modele improvise volontiers un champ : le laisser filer jusqu'au + // service reviendrait a lui laisser choisir la signature de l'appel. + expect(() => parseInput(schema, { origin: 'FRLIO', volumeCBM: 4, orgId: 'autre-org' })).toThrow( + 'Unknown parameter "orgId"' + ); + }); + + it('accepts a call with no arguments at all', () => { + const empty: CapabilitySchema = { type: 'object', properties: {}, additionalProperties: false }; + expect(parseInput(empty, undefined)).toEqual({}); + }); + + it('accepts booleans and arrays in their natural form', () => { + expect( + parseInput(schema, { + origin: 'FRLIO', + volumeCBM: 4, + dangerous: true, + companies: ['CMA CGM', 'MSC'], + }) + ).toMatchObject({ dangerous: true, companies: ['CMA CGM', 'MSC'] }); + }); +}); diff --git a/apps/backend/src/application/mcp/capability.ts b/apps/backend/src/application/mcp/capability.ts new file mode 100644 index 0000000..6a73119 --- /dev/null +++ b/apps/backend/src/application/mcp/capability.ts @@ -0,0 +1,146 @@ +import { CapabilityActor, CapabilityPolicy } from '@domain/services/capability-access'; + +/** + * Sous-ensemble de JSON Schema utilise par les capacites. + * + * Le schema est ecrit une seule fois : il sert a la fois de contrat annonce aux + * clients MCP (`inputSchema`) et de regle de validation a l'entree. Deux + * sources auraient fini par diverger, et c'est la validation qui aurait perdu. + */ +export interface SchemaProperty { + type: 'string' | 'number' | 'integer' | 'boolean' | 'array'; + description: string; + enum?: readonly string[]; + minimum?: number; + maximum?: number; + minLength?: number; + maxLength?: number; + /** Pour `type: 'array'` uniquement. */ + items?: { type: 'string' | 'number' }; + default?: unknown; +} + +export interface CapabilitySchema { + type: 'object'; + properties: Record; + required?: readonly string[]; + additionalProperties: false; +} + +export interface Capability { + policy: CapabilityPolicy; + /** Une phrase : ce que fait l'action, du point de vue de l'utilisateur. */ + description: string; + inputSchema: CapabilitySchema; + handler: (input: Record, actor: CapabilityActor) => Promise; +} + +/** Schema sans aucun parametre, pour les capacites qui n'en prennent pas. */ +export const NO_INPUT: CapabilitySchema = { + type: 'object', + properties: {}, + additionalProperties: false, +}; + +export class CapabilityInputError extends Error {} + +/** + * Valide et normalise une entree contre son schema. + * + * Les entrees viennent d'un modele de langage : elles sont plausibles, pas + * fiables. Un nombre arrive en chaine, un champ facultatif arrive a `null`, un + * champ invente arrive en plus. La validation est donc stricte sur ce qui + * compte (types, valeurs autorisees, bornes) et refuse ce qu'elle ne connait + * pas, plutot que de le transmettre au service. + */ +export function parseInput( + schema: CapabilitySchema, + raw: Record | undefined +): Record { + const input = raw ?? {}; + const parsed: Record = {}; + + for (const key of Object.keys(input)) { + if (!(key in schema.properties)) { + throw new CapabilityInputError(`Unknown parameter "${key}".`); + } + } + + for (const [key, property] of Object.entries(schema.properties)) { + const required = schema.required?.includes(key) ?? false; + const value = input[key]; + + if (value === undefined || value === null || value === '') { + if (required) throw new CapabilityInputError(`Missing required parameter "${key}".`); + if (property.default !== undefined) parsed[key] = property.default; + continue; + } + + parsed[key] = coerce(key, property, value); + } + + return parsed; +} + +function coerce(key: string, property: SchemaProperty, value: unknown): unknown { + switch (property.type) { + case 'string': { + if (typeof value !== 'string') { + throw new CapabilityInputError(`Parameter "${key}" must be a string.`); + } + const text = value.trim(); + if (property.enum && !property.enum.includes(text)) { + throw new CapabilityInputError( + `Parameter "${key}" must be one of: ${property.enum.join(', ')}.` + ); + } + if (property.minLength !== undefined && text.length < property.minLength) { + throw new CapabilityInputError( + `Parameter "${key}" must be at least ${property.minLength} characters.` + ); + } + if (property.maxLength !== undefined && text.length > property.maxLength) { + throw new CapabilityInputError( + `Parameter "${key}" must be at most ${property.maxLength} characters.` + ); + } + return text; + } + + case 'number': + case 'integer': { + // Un modele ecrit volontiers « 12.5 » plutot que 12.5 : la chaine + // numerique est acceptee, le texte non numerique refuse. + const numeric = typeof value === 'number' ? value : Number(String(value).trim()); + if (!Number.isFinite(numeric)) { + throw new CapabilityInputError(`Parameter "${key}" must be a number.`); + } + if (property.type === 'integer' && !Number.isInteger(numeric)) { + throw new CapabilityInputError(`Parameter "${key}" must be a whole number.`); + } + if (property.minimum !== undefined && numeric < property.minimum) { + throw new CapabilityInputError(`Parameter "${key}" must be >= ${property.minimum}.`); + } + if (property.maximum !== undefined && numeric > property.maximum) { + throw new CapabilityInputError(`Parameter "${key}" must be <= ${property.maximum}.`); + } + return numeric; + } + + case 'boolean': { + if (typeof value === 'boolean') return value; + const text = String(value).trim().toLowerCase(); + if (text === 'true') return true; + if (text === 'false') return false; + throw new CapabilityInputError(`Parameter "${key}" must be true or false.`); + } + + case 'array': { + if (!Array.isArray(value)) { + throw new CapabilityInputError(`Parameter "${key}" must be an array.`); + } + const itemType = property.items?.type ?? 'string'; + return value.map(item => coerce(`${key}[]`, { type: itemType, description: '' }, item)); + } + } +} diff --git a/apps/backend/src/application/mcp/mcp.controller.ts b/apps/backend/src/application/mcp/mcp.controller.ts new file mode 100644 index 0000000..68c9962 --- /dev/null +++ b/apps/backend/src/application/mcp/mcp.controller.ts @@ -0,0 +1,190 @@ +import { Body, Controller, HttpCode, Post } from '@nestjs/common'; +import { ApiBearerAuth, ApiOperation, ApiResponse, ApiTags } from '@nestjs/swagger'; +import { CapabilityActor } from '@domain/services/capability-access'; +import { CurrentUser, UserPayload } from '../decorators/current-user.decorator'; +import { SubscriptionService } from '../services/subscription.service'; +import { CapabilityInputError } from './capability'; +import { CapabilityRegistry } from './capability.registry'; + +/** + * Serveur MCP d'Xpeditis. + * + * Expose les capacites du produit au protocole Model Context Protocol, sur une + * unique route HTTP. Le transport est volontairement minimal : un POST + * JSON-RPC, sans session ni flux SSE. Un serveur qui n'expose que des outils + * n'a rien a diffuser au client entre deux appels, et l'absence d'etat rend + * chaque requete authentifiable independamment — ce qui compte ici, puisque + * deux appels consecutifs peuvent venir de deux comptes differents. + * + * L'authentification n'est pas reimplementee : la route passe par le garde + * global `ApiKeyOrJwtGuard`, donc une cle API `X-API-Key` (offres Gold et + * Platinium) ou un jeton JWT. L'identite obtenue porte le role et l'offre, qui + * decident ensuite de ce que le catalogue laisse voir. + * + * Non couvert a ce stade : les ressources et les invites MCP, la negociation + * SSE, et les notifications serveur → client. + */ + +const PROTOCOL_VERSION = '2025-06-18'; +const SERVER_INFO = { name: 'xpeditis', version: '1.0.0' }; + +/** Codes d'erreur JSON-RPC 2.0. */ +const enum RpcError { + InvalidRequest = -32600, + MethodNotFound = -32601, + InvalidParams = -32602, + InternalError = -32603, +} + +interface RpcRequest { + jsonrpc?: string; + id?: string | number | null; + method?: string; + params?: Record; +} + +@ApiTags('MCP') +@ApiBearerAuth() +@Controller('mcp') +export class McpController { + constructor( + private readonly registry: CapabilityRegistry, + private readonly subscriptions: SubscriptionService + ) {} + + @Post() + @HttpCode(200) + @ApiOperation({ + summary: 'Model Context Protocol endpoint', + description: + 'JSON-RPC 2.0 endpoint exposing Xpeditis capabilities as MCP tools. Authenticate with an X-API-Key header (Gold and Platinium plans) or a JWT bearer token. Supported methods: initialize, tools/list, tools/call, ping.', + }) + @ApiResponse({ status: 200, description: 'JSON-RPC response' }) + @ApiResponse({ status: 401, description: 'Unauthorized' }) + async rpc(@CurrentUser() user: UserPayload, @Body() body: RpcRequest | RpcRequest[]) { + // Un lot JSON-RPC est traite element par element, dans l'ordre reçu. + if (Array.isArray(body)) { + const responses = await Promise.all(body.map(entry => this.handle(user, entry))); + return responses.filter(response => response !== null); + } + return this.handle(user, body); + } + + private async handle(user: UserPayload, request: RpcRequest) { + const id = request?.id ?? null; + + // Une notification (sans `id`) n'attend pas de reponse : `notifications/initialized` + // arrive juste apres la poignee de main de tout client MCP. + if (id === null && request?.method?.startsWith('notifications/')) return null; + + if (request?.jsonrpc !== '2.0' || typeof request.method !== 'string') { + return fail(id, RpcError.InvalidRequest, 'Invalid JSON-RPC 2.0 request.'); + } + + try { + switch (request.method) { + case 'initialize': + return ok(id, { + protocolVersion: PROTOCOL_VERSION, + capabilities: { tools: { listChanged: false } }, + serverInfo: SERVER_INFO, + instructions: + "Xpeditis est une plateforme de réservation de fret maritime LCL. Les outils disponibles dépendent du rôle et de l'offre du compte authentifié : appelez `whoami` pour connaître les droits en cours. Pour une question de connaissance métier, préférez `search_documentation`, qui répond à partir du wiki Xpeditis.", + }); + + case 'ping': + return ok(id, {}); + + case 'tools/list': { + const actor = await this.actorOf(user); + return ok(id, { + tools: this.registry.listFor(actor).map(capability => ({ + name: capability.policy.name, + description: capability.description, + inputSchema: capability.inputSchema, + annotations: { readOnlyHint: capability.policy.scope === 'read' }, + })), + }); + } + + case 'tools/call': { + const name = request.params?.name; + if (typeof name !== 'string') { + return fail(id, RpcError.InvalidParams, 'Missing tool name.'); + } + const actor = await this.actorOf(user); + const result = await this.registry.invoke( + name, + request.params?.arguments as Record | undefined, + actor + ); + return ok(id, { + content: [{ type: 'text', text: JSON.stringify(result, null, 2) }], + isError: false, + }); + } + + default: + return fail(id, RpcError.MethodNotFound, `Unknown method "${request.method}".`); + } + } catch (error) { + return this.toRpcError(id, request.method, error); + } + } + + /** + * Une erreur d'outil se rend au modele, pas au transport : MCP demande de + * repondre `isError` dans le resultat pour qu'un agent puisse corriger son + * appel, la ou une erreur JSON-RPC interromprait l'echange. + */ + private toRpcError(id: string | number | null, method: string | undefined, error: unknown) { + const message = error instanceof Error ? error.message : String(error); + + if (method === 'tools/call') { + const invalid = error instanceof CapabilityInputError; + return ok(id, { + content: [{ type: 'text', text: message }], + isError: true, + ...(invalid ? {} : {}), + }); + } + + return fail(id, RpcError.InternalError, message); + } + + /** + * Identite de l'appelant, completee de son offre. + * + * Une cle API porte deja l'offre ; un jeton JWT ne la porte pas, elle est + * alors lue sur l'abonnement. Sans cette resolution, un utilisateur connecte + * a l'application serait traite comme un compte Bronze. + */ + private async actorOf(user: UserPayload & { plan?: string }): Promise { + if (user.plan) { + return { + id: user.id, + organizationId: user.organizationId, + role: user.role, + email: user.email, + plan: user.plan, + }; + } + + const subscription = await this.subscriptions.getOrCreateSubscription(user.organizationId); + return { + id: user.id, + organizationId: user.organizationId, + role: user.role, + email: user.email, + plan: subscription.plan.value, + }; + } +} + +const ok = (id: string | number | null, result: unknown) => ({ jsonrpc: '2.0', id, result }); + +const fail = (id: string | number | null, code: number, message: string) => ({ + jsonrpc: '2.0', + id, + error: { code, message }, +}); diff --git a/apps/backend/src/application/mcp/mcp.module.ts b/apps/backend/src/application/mcp/mcp.module.ts new file mode 100644 index 0000000..f32d462 --- /dev/null +++ b/apps/backend/src/application/mcp/mcp.module.ts @@ -0,0 +1,38 @@ +import { Module } from '@nestjs/common'; +import { TRADE_RETRIEVAL, TRADE_EMBEDDINGS } from '@domain/ports/out/trade-assistant.port'; +import { OpenAiEmbeddingAdapter } from '@infrastructure/ai/openai-embedding.adapter'; +import { WikiRetriever } from '@infrastructure/ai/wiki-retriever'; +import { CsvRateModule } from '@infrastructure/carriers/csv-loader/csv-rate.module'; +import { AuditModule } from '../audit/audit.module'; +import { CsvBookingsModule } from '../csv-bookings/csv-bookings.module'; +import { OrganizationsModule } from '../organizations/organizations.module'; +import { SubscriptionsModule } from '../subscriptions/subscriptions.module'; +import { UsersModule } from '../users/users.module'; +import { CapabilityRegistry } from './capability.registry'; +import { McpController } from './mcp.controller'; + +/** + * Serveur MCP et registre de capacites. + * + * Le module n'apporte aucune logique metier : il assemble des services deja + * exposes ailleurs. C'est le point de la conception — les agents passent par + * les memes chemins que l'interface. + */ +@Module({ + imports: [ + CsvRateModule, + CsvBookingsModule, + SubscriptionsModule, + UsersModule, + OrganizationsModule, + AuditModule, + ], + controllers: [McpController], + providers: [ + CapabilityRegistry, + { provide: TRADE_EMBEDDINGS, useClass: OpenAiEmbeddingAdapter }, + { provide: TRADE_RETRIEVAL, useClass: WikiRetriever }, + ], + exports: [CapabilityRegistry], +}) +export class McpModule {} diff --git a/apps/backend/src/application/trade-assistant/trade-assistant.module.ts b/apps/backend/src/application/trade-assistant/trade-assistant.module.ts index 827929b..e3c5e75 100644 --- a/apps/backend/src/application/trade-assistant/trade-assistant.module.ts +++ b/apps/backend/src/application/trade-assistant/trade-assistant.module.ts @@ -12,12 +12,15 @@ import { OpenAiTradeAdapter } from '@infrastructure/ai/openai-trade.adapter'; import { WikiRetriever } from '@infrastructure/ai/wiki-retriever'; import { TypeOrmTradeConversationRepository } from '@infrastructure/persistence/typeorm/repositories/typeorm-trade-conversation.repository'; import { TypeOrmTradeQuotaRepository } from '@infrastructure/persistence/typeorm/repositories/typeorm-trade-quota.repository'; +import { McpModule } from '../mcp/mcp.module'; import { SubscriptionsModule } from '../subscriptions/subscriptions.module'; import { TradeAssistantController } from './trade-assistant.controller'; import { TradeAssistantService } from './trade-assistant.service'; @Module({ + // `McpModule` fournit le registre de capacites : sans lui, l'assistant // repond mais n'agit jamais. + imports: [ConfigModule, SubscriptionsModule, McpModule], controllers: [TradeAssistantController], providers: [ TradeAssistantService, diff --git a/apps/backend/src/application/trade-assistant/trade-assistant.service.ts b/apps/backend/src/application/trade-assistant/trade-assistant.service.ts index 6550418..2cedc8a 100644 --- a/apps/backend/src/application/trade-assistant/trade-assistant.service.ts +++ b/apps/backend/src/application/trade-assistant/trade-assistant.service.ts @@ -3,6 +3,7 @@ import { Injectable, Logger, NotFoundException, + Optional, ServiceUnavailableException, } from '@nestjs/common'; import { @@ -22,6 +23,8 @@ import { TradeQuotaPort, TradeRetrievalPort, TradeSource, + TradeToolDefinition, + TradeToolInvoker, } from '@domain/ports/out/trade-assistant.port'; import { TRADE_SUPPORT_EMAIL, @@ -29,6 +32,7 @@ import { tradeDailyLimit, } from '@domain/services/trade-assistant-policy'; import { effectivePlan } from '@domain/services/subscription-access'; +import { CapabilityRegistry } from '../mcp/capability.registry'; /** * L'utilisateur qui interroge l'assistant. @@ -59,7 +63,9 @@ export class TradeAssistantService { @Inject(TRADE_QUOTA) private readonly quota: TradeQuotaPort, @Inject(TRADE_AI) private readonly ai: TradeAiPort, @Inject(TRADE_RETRIEVAL) private readonly retrieval: TradeRetrievalPort, - @Inject(TRADE_CONVERSATIONS) private readonly conversations: TradeConversationRepository + @Inject(TRADE_CONVERSATIONS) private readonly conversations: TradeConversationRepository, + // Optionnel : sans registre, l'assistant repond sans jamais agir. + @Optional() private readonly capabilities?: CapabilityRegistry ) {} async status(actor: TradeActor) { @@ -153,10 +159,11 @@ export class TradeAssistantService { : []; const passages = await this.retrieve(question, language); + const { tools, invokeTool } = this.toolsFor({ ...actor, plan: status.plan }); let answer; try { - answer = await this.ai.answer({ question, language, history, passages }); + answer = await this.ai.answer({ question, language, history, passages, tools, invokeTool }); } catch { // Le remboursement vise le jour reserve, meme si la reponse a franchi minuit. await this.quota.release(userId, status.day); @@ -195,6 +202,44 @@ export class TradeAssistantService { }; } + /** + * Outils ouverts a cet utilisateur, et le moyen de les executer. + * + * Le catalogue est filtre par le registre selon le role et l'offre : le + * modele ne voit que ce que la personne a le droit de faire, donc il ne peut + * pas proposer une action interdite — encore moins la declencher. + * + * L'executeur est lie a `actor` : les arguments du modele decrivent *quoi* + * faire, jamais *pour qui*. Une identite ne peut pas etre passee en + * parametre, elle vient de la session. + */ + private toolsFor(actor: TradeActor) { + if (!this.capabilities) return {}; + + const tools: TradeToolDefinition[] = this.capabilities.listFor(actor).map(capability => ({ + name: capability.policy.name, + description: capability.description, + parameters: capability.inputSchema as unknown as Record, + })); + + const invokeTool: TradeToolInvoker = async (name, args) => { + try { + return { + ok: true, + result: await this.capabilities!.invoke(name, args, actor, 'assistant'), + }; + } catch (error) { + // L'echec repart vers le modele comme un resultat : il peut corriger + // son appel ou l'expliquer, au lieu de perdre la reponse en cours. + const message = error instanceof Error ? error.message : String(error); + this.logger.warn(`Assistant tool "${name}" failed: ${message}`); + return { ok: false, result: { error: message } }; + } + }; + + return { tools, invokeTool }; + } + /** * La recherche documentaire ne doit jamais empecher une reponse : sans * extrait, le modele repond sur ses connaissances generales. diff --git a/apps/backend/src/domain/entities/audit-log.entity.ts b/apps/backend/src/domain/entities/audit-log.entity.ts index 792a593..bb5e922 100644 --- a/apps/backend/src/domain/entities/audit-log.entity.ts +++ b/apps/backend/src/domain/entities/audit-log.entity.ts @@ -42,6 +42,10 @@ export enum AuditAction { // Settings actions SETTINGS_UPDATED = 'settings_updated', + + // Agent actions — toute capacite invoquee par un agent, via MCP ou via + // l'assistant integre. Le nom de la capacite est dans `resourceName`. + AGENT_CAPABILITY_INVOKED = 'agent_capability_invoked', } export enum AuditStatus { diff --git a/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788800000000-AddTradeMessageActions.ts b/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788800000000-AddTradeMessageActions.ts new file mode 100644 index 0000000..f017754 --- /dev/null +++ b/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788800000000-AddTradeMessageActions.ts @@ -0,0 +1,16 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class AddTradeMessageActions1788800000000 implements MigrationInterface { + async up(queryRunner: QueryRunner): Promise { + // Les capacites invoquees pour produire la reponse. Conservees avec le + // message : au rechargement de la conversation, l'utilisateur doit toujours + // voir ce que l'assistant a réellement fait, pas seulement ce qu'il a dit. + await queryRunner.query( + `ALTER TABLE trade_messages ADD COLUMN actions jsonb NOT NULL DEFAULT '[]'::jsonb` + ); + } + + async down(queryRunner: QueryRunner): Promise { + await queryRunner.query('ALTER TABLE trade_messages DROP COLUMN actions'); + } +}