feat(api): registre de capacites et serveur MCP
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018BAUeCFpDkRD6tU5wGsc1C
This commit is contained in:
parent
8b45d2f1a6
commit
177b70ea96
@ -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,
|
||||
],
|
||||
|
||||
@ -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,
|
||||
};
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
@ -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(),
|
||||
}),
|
||||
},
|
||||
];
|
||||
}
|
||||
@ -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),
|
||||
},
|
||||
];
|
||||
}
|
||||
@ -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,
|
||||
})),
|
||||
};
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
@ -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() }),
|
||||
},
|
||||
];
|
||||
}
|
||||
237
apps/backend/src/application/mcp/capability.registry.spec.ts
Normal file
237
apps/backend/src/application/mcp/capability.registry.spec.ts
Normal file
@ -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> = {}): 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,
|
||||
});
|
||||
});
|
||||
});
|
||||
156
apps/backend/src/application/mcp/capability.registry.ts
Normal file
156
apps/backend/src/application/mcp/capability.registry.ts
Normal file
@ -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<string, unknown> | undefined,
|
||||
actor: CapabilityActor,
|
||||
surface: CapabilitySurface = 'mcp'
|
||||
): Promise<unknown> {
|
||||
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<void> {
|
||||
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<string>();
|
||||
for (const name of names) {
|
||||
if (seen.has(name)) return name;
|
||||
seen.add(name);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
84
apps/backend/src/application/mcp/capability.spec.ts
Normal file
84
apps/backend/src/application/mcp/capability.spec.ts
Normal file
@ -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<string, unknown>)).toThrow(
|
||||
CapabilityInputError
|
||||
);
|
||||
expect(() => parseInput(schema, input as Record<string, unknown>)).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'] });
|
||||
});
|
||||
});
|
||||
146
apps/backend/src/application/mcp/capability.ts
Normal file
146
apps/backend/src/application/mcp/capability.ts
Normal file
@ -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<string, SchemaProperty>;
|
||||
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<string, unknown>, actor: CapabilityActor) => Promise<unknown>;
|
||||
}
|
||||
|
||||
/** 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<string, unknown> | undefined
|
||||
): Record<string, unknown> {
|
||||
const input = raw ?? {};
|
||||
const parsed: Record<string, unknown> = {};
|
||||
|
||||
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));
|
||||
}
|
||||
}
|
||||
}
|
||||
190
apps/backend/src/application/mcp/mcp.controller.ts
Normal file
190
apps/backend/src/application/mcp/mcp.controller.ts
Normal file
@ -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<string, unknown>;
|
||||
}
|
||||
|
||||
@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<string, unknown> | 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<CapabilityActor> {
|
||||
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 },
|
||||
});
|
||||
38
apps/backend/src/application/mcp/mcp.module.ts
Normal file
38
apps/backend/src/application/mcp/mcp.module.ts
Normal file
@ -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 {}
|
||||
@ -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,
|
||||
|
||||
@ -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<string, unknown>,
|
||||
}));
|
||||
|
||||
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.
|
||||
|
||||
@ -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 {
|
||||
|
||||
@ -0,0 +1,16 @@
|
||||
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||
|
||||
export class AddTradeMessageActions1788800000000 implements MigrationInterface {
|
||||
async up(queryRunner: QueryRunner): Promise<void> {
|
||||
// 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<void> {
|
||||
await queryRunner.query('ALTER TABLE trade_messages DROP COLUMN actions');
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue
Block a user