diff --git a/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788600000000-CreateTradeAssistantUsage.ts b/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788600000000-CreateTradeAssistantUsage.ts new file mode 100644 index 0000000..a2b8492 --- /dev/null +++ b/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788600000000-CreateTradeAssistantUsage.ts @@ -0,0 +1,17 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class CreateTradeAssistantUsage1788600000000 implements MigrationInterface { + async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(`CREATE TABLE trade_assistant_usage ( + user_id uuid NOT NULL REFERENCES users(id) ON DELETE CASCADE, + day date NOT NULL, + used integer NOT NULL DEFAULT 0 CHECK (used >= 0), + input_tokens bigint NOT NULL DEFAULT 0, + output_tokens bigint NOT NULL DEFAULT 0, + PRIMARY KEY (user_id, day) + )`); + } + async down(queryRunner: QueryRunner): Promise { + await queryRunner.query('DROP TABLE trade_assistant_usage'); + } +} diff --git a/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788700000000-CreateTradeConversations.ts b/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788700000000-CreateTradeConversations.ts new file mode 100644 index 0000000..e030174 --- /dev/null +++ b/apps/backend/src/infrastructure/persistence/typeorm/migrations/1788700000000-CreateTradeConversations.ts @@ -0,0 +1,37 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class CreateTradeConversations1788700000000 implements MigrationInterface { + async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(`CREATE TABLE trade_conversations ( + id uuid PRIMARY KEY DEFAULT uuid_generate_v4(), + user_id uuid NOT NULL REFERENCES users(id) ON DELETE CASCADE, + title text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now() + )`); + + // La liste laterale n'affiche que les conversations d'un utilisateur, de la + // plus recemment active a la plus ancienne : l'index sert exactement cela. + await queryRunner.query( + 'CREATE INDEX idx_trade_conversations_user ON trade_conversations (user_id, updated_at DESC)' + ); + + await queryRunner.query(`CREATE TABLE trade_messages ( + id uuid PRIMARY KEY DEFAULT uuid_generate_v4(), + conversation_id uuid NOT NULL REFERENCES trade_conversations(id) ON DELETE CASCADE, + role text NOT NULL CHECK (role IN ('user', 'assistant')), + content text NOT NULL, + sources jsonb NOT NULL DEFAULT '[]'::jsonb, + created_at timestamptz NOT NULL DEFAULT now() + )`); + + await queryRunner.query( + 'CREATE INDEX idx_trade_messages_conversation ON trade_messages (conversation_id, created_at)' + ); + } + + async down(queryRunner: QueryRunner): Promise { + await queryRunner.query('DROP TABLE trade_messages'); + await queryRunner.query('DROP TABLE trade_conversations'); + } +} diff --git a/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-conversation.repository.ts b/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-conversation.repository.ts new file mode 100644 index 0000000..ef680e4 --- /dev/null +++ b/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-conversation.repository.ts @@ -0,0 +1,142 @@ +import { Injectable } from '@nestjs/common'; +import { DataSource } from 'typeorm'; +import { + TradeAction, + TradeConversationRepository, + TradeConversationSummary, + TradeMessage, + TradeSource, +} from '@domain/ports/out/trade-assistant.port'; + +/** + * Conversations de l'assistant. + * + * Comme le reste de la feature (voir `typeorm-trade-quota.repository.ts`), les + * acces passent par du SQL parametre plutot que par des entites TypeORM : les + * requetes utiles ici sont des agregats et des mises a jour conditionnelles que + * l'ORM rendrait plus longs a lire, pas plus surs. + * + * Chaque requete porte `user_id` : une conversation ne peut etre lue, renommee + * ou supprimee que par son proprietaire, sans controle d'acces separe a oublier. + */ +@Injectable() +export class TypeOrmTradeConversationRepository implements TradeConversationRepository { + constructor(private readonly db: DataSource) {} + + async list(userId: string): Promise { + const rows: RawSummary[] = await this.db.query( + `SELECT c.id, c.title, c.created_at, c.updated_at, + (SELECT COUNT(*) FROM trade_messages m WHERE m.conversation_id = c.id) AS message_count + FROM trade_conversations c + WHERE c.user_id = $1 + ORDER BY c.updated_at DESC`, + [userId] + ); + return rows.map(toSummary); + } + + async create(userId: string, title: string): Promise { + const rows: RawSummary[] = await this.db.query( + `INSERT INTO trade_conversations (user_id, title) VALUES ($1, $2) + RETURNING id, title, created_at, updated_at, 0 AS message_count`, + [userId, title] + ); + return toSummary(rows[0]); + } + + async find(userId: string, conversationId: string): Promise { + const rows: RawSummary[] = await this.db.query( + `SELECT c.id, c.title, c.created_at, c.updated_at, + (SELECT COUNT(*) FROM trade_messages m WHERE m.conversation_id = c.id) AS message_count + FROM trade_conversations c + WHERE c.id = $1 AND c.user_id = $2`, + [conversationId, userId] + ); + return rows.length ? toSummary(rows[0]) : null; + } + + async messages(userId: string, conversationId: string): Promise { + const rows: RawMessage[] = await this.db.query( + `SELECT m.id, m.role, m.content, m.sources, m.actions, m.created_at + FROM trade_messages m + JOIN trade_conversations c ON c.id = m.conversation_id AND c.user_id = $2 + WHERE m.conversation_id = $1 + ORDER BY m.created_at, m.id`, + [conversationId, userId] + ); + return rows.map(toMessage); + } + + async addMessage( + conversationId: string, + role: 'user' | 'assistant', + content: string, + sources: TradeSource[] = [], + actions: TradeAction[] = [] + ): Promise { + const rows: RawMessage[] = await this.db.query( + `INSERT INTO trade_messages (conversation_id, role, content, sources, actions) + VALUES ($1, $2, $3, $4::jsonb, $5::jsonb) + RETURNING id, role, content, sources, actions, created_at`, + [conversationId, role, content, JSON.stringify(sources), JSON.stringify(actions)] + ); + + // La date de mise a jour classe la liste laterale : elle suit le dernier + // message, pas la creation. + await this.db.query('UPDATE trade_conversations SET updated_at = now() WHERE id = $1', [ + conversationId, + ]); + + return toMessage(rows[0]); + } + + async rename(userId: string, conversationId: string, title: string): Promise { + await this.db.query( + 'UPDATE trade_conversations SET title = $3 WHERE id = $1 AND user_id = $2', + [conversationId, userId, title] + ); + } + + async remove(userId: string, conversationId: string): Promise { + await this.db.query('DELETE FROM trade_conversations WHERE id = $1 AND user_id = $2', [ + conversationId, + userId, + ]); + } +} + +/* -------------------------------------------------------------------------- */ + +interface RawSummary { + id: string; + title: string; + created_at: Date; + updated_at: Date; + message_count: string | number; +} + +interface RawMessage { + id: string; + role: 'user' | 'assistant'; + content: string; + sources: TradeSource[] | null; + actions: TradeAction[] | null; + created_at: Date; +} + +const toSummary = (row: RawSummary): TradeConversationSummary => ({ + id: row.id, + title: row.title, + createdAt: row.created_at.toISOString(), + updatedAt: row.updated_at.toISOString(), + messageCount: Number(row.message_count), +}); + +const toMessage = (row: RawMessage): TradeMessage => ({ + id: row.id, + role: row.role, + content: row.content, + sources: row.sources ?? [], + actions: row.actions ?? [], + createdAt: row.created_at.toISOString(), +}); diff --git a/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-quota.repository.spec.ts b/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-quota.repository.spec.ts new file mode 100644 index 0000000..a7d8661 --- /dev/null +++ b/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-quota.repository.spec.ts @@ -0,0 +1,102 @@ +import { DataSource } from 'typeorm'; +import { randomUUID } from 'crypto'; +import { TypeOrmTradeQuotaRepository } from './typeorm-trade-quota.repository'; +import { CreateTradeAssistantUsage1788600000000 } from '../migrations/1788600000000-CreateTradeAssistantUsage'; + +// Opt in only against the disposable PostgreSQL documented in docs/features/trade-assistant.md. +const run = process.env.TRADE_TEST_DATABASE_URL ? describe : describe.skip; +run('Trade quota PostgreSQL integration', () => { + let db: DataSource; + let quota: TypeOrmTradeQuotaRepository; + const firstUser = randomUUID(); + const secondUser = randomUUID(); + const schema = 'trade_test_' + randomUUID().replace(/-/g, ''); + beforeAll(async () => { + db = new DataSource({ + type: 'postgres', + url: process.env.TRADE_TEST_DATABASE_URL, + extra: { options: `-c search_path=${schema}` }, + }); + await db.initialize(); + await db.query(`CREATE SCHEMA "${schema}"`); + await db.query('CREATE TABLE users (id uuid PRIMARY KEY)'); + const runner = db.createQueryRunner(); + try { + await new CreateTradeAssistantUsage1788600000000().up(runner); + } finally { + await runner.release(); + } + await db.query('INSERT INTO users VALUES ($1), ($2)', [firstUser, secondUser]); + quota = new TypeOrmTradeQuotaRepository(db); + }); + afterAll(async () => { + if (db?.isInitialized) { + await db.query(`DROP SCHEMA "${schema}" CASCADE`); + await db.destroy(); + } + }); + it('accepts exactly three of twenty concurrent Bronze requests', async () => { + const initial = await quota.get(firstUser); + expect(initial.used).toBe(0); + expect(new Date(initial.resetsAt).getTime()).toBeGreaterThan(Date.now()); + const results = await Promise.all( + Array.from({ length: 20 }, () => quota.reserve(firstUser, initial.day, 3)) + ); + expect(results.filter(Boolean)).toHaveLength(3); + expect((await quota.get(firstUser)).used).toBe(3); + expect((await quota.get(secondUser)).used).toBe(0); + await quota.release(firstUser, initial.day); + expect(await quota.reserve(firstUser, initial.day, 3)).toBe(true); + expect(await quota.reserve(firstUser, initial.day, 3)).toBe(false); + }); + it('never blocks an unlimited plan, and keeps counting it', async () => { + const unlimitedUser = randomUUID(); + await db.query('INSERT INTO users VALUES ($1)', [unlimitedUser]); + const { day } = await quota.get(unlimitedUser); + + // Avec `-1`, la condition `used < -1` etait toujours fausse : la premiere + // question passait par l'INSERT, toutes les suivantes etaient refusees. + const results = await Promise.all( + Array.from({ length: 25 }, () => quota.reserve(unlimitedUser, day, -1)) + ); + + expect(results.filter(Boolean)).toHaveLength(25); + expect((await quota.get(unlimitedUser)).used).toBe(25); + }); + + it('ignores previous-day usage and never reserves an expired window', async () => { + await db.query( + "INSERT INTO trade_assistant_usage (user_id, day, used) VALUES ($1, DATE '2000-01-01', 15)", + [secondUser] + ); + expect((await quota.get(secondUser)).used).toBe(0); + expect(await quota.reserve(secondUser, '2000-01-01', 15)).toBe(false); + await quota.release(secondUser, '2000-01-01'); + expect((await quota.get(secondUser)).used).toBe(0); + }); + it('records tokens and removes usage when its user is deleted', async () => { + const { day } = await quota.get(secondUser); + await quota.reserve(secondUser, day, 10); + await quota.recordTokens(secondUser, day, { + text: 'unused', + inputTokens: 100, + outputTokens: 50, + }); + const rows = await db.query( + 'SELECT input_tokens, output_tokens FROM trade_assistant_usage WHERE user_id=$1 AND day=$2', + [secondUser, day] + ); + expect(rows[0]).toEqual({ input_tokens: '100', output_tokens: '50' }); + await db.query('DELETE FROM users WHERE id=$1', [secondUser]); + expect( + await db.query('SELECT * FROM trade_assistant_usage WHERE user_id=$1', [secondUser]) + ).toEqual([]); + }); + it('computes Paris midnight correctly across daylight saving changes', async () => { + const rows = await db.query(`SELECT + ((DATE '2026-03-29' + 1)::timestamp AT TIME ZONE 'Europe/Paris') AS spring, + ((DATE '2026-10-25' + 1)::timestamp AT TIME ZONE 'Europe/Paris') AS autumn`); + expect(rows[0].spring.toISOString()).toBe('2026-03-29T22:00:00.000Z'); + expect(rows[0].autumn.toISOString()).toBe('2026-10-25T23:00:00.000Z'); + }); +}); diff --git a/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-quota.repository.ts b/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-quota.repository.ts new file mode 100644 index 0000000..c98de2b --- /dev/null +++ b/apps/backend/src/infrastructure/persistence/typeorm/repositories/typeorm-trade-quota.repository.ts @@ -0,0 +1,60 @@ +import { Injectable } from '@nestjs/common'; +import { DataSource } from 'typeorm'; +import { TradeQuotaPort, TradeUsage, TradeAnswer } from '@domain/ports/out/trade-assistant.port'; + +@Injectable() +export class TypeOrmTradeQuotaRepository implements TradeQuotaPort { + constructor(private readonly db: DataSource) {} + + async get(userId: string): Promise { + const rows: Array<{ day: string; resetsAt: Date; used: number }> = await this.db.query( + ` + SELECT to_char(w.day, 'YYYY-MM-DD') AS day, + ((w.day + 1)::timestamp AT TIME ZONE 'Europe/Paris') AS "resetsAt", + COALESCE(q.used, 0)::integer AS used + FROM (SELECT (CURRENT_TIMESTAMP AT TIME ZONE 'Europe/Paris')::date AS day) w + LEFT JOIN trade_assistant_usage q ON q.user_id = $1 AND q.day = w.day`, + [userId] + ); + return { ...rows[0], resetsAt: rows[0].resetsAt.toISOString() }; + } + + /** + * Reserve une question pour la journee. + * + * `limit` negatif signifie illimite (offre Platinium) : la consommation est + * toujours comptee — c'est la base du suivi de cout — mais la mise a jour + * n'est plus conditionnee au plafond. Sans cette branche, `used < -1` etait + * toujours faux et l'offre illimitee etait en realite bloquee des la + * deuxieme question de la journee. + */ + async reserve(userId: string, day: string, limit: number): Promise { + const cap = limit < 0 ? 'TRUE' : 'trade_assistant_usage.used < $3'; + const parameters = limit < 0 ? [userId, day] : [userId, day, limit]; + + const rows: Array<{ used: number }> = await this.db.query( + ` + INSERT INTO trade_assistant_usage (user_id, day, used) + SELECT $1, $2::date, 1 WHERE $2::date = (CURRENT_TIMESTAMP AT TIME ZONE 'Europe/Paris')::date + ON CONFLICT (user_id, day) DO UPDATE SET used = trade_assistant_usage.used + 1 + WHERE ${cap} RETURNING used`, + parameters + ); + return rows.length > 0; + } + + async release(userId: string, day: string): Promise { + await this.db.query( + 'UPDATE trade_assistant_usage SET used = GREATEST(0, used - 1) WHERE user_id = $1 AND day = $2', + [userId, day] + ); + } + + async recordTokens(userId: string, day: string, answer: TradeAnswer): Promise { + await this.db.query( + `UPDATE trade_assistant_usage SET input_tokens = input_tokens + $3, + output_tokens = output_tokens + $4 WHERE user_id = $1 AND day = $2`, + [userId, day, answer.inputTokens, answer.outputTokens] + ); + } +}