feat(db): tables de quota et de conversations de l assistant
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018BAUeCFpDkRD6tU5wGsc1C
This commit is contained in:
parent
b94ceee73f
commit
d32eecd0bf
@ -0,0 +1,17 @@
|
|||||||
|
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||||
|
|
||||||
|
export class CreateTradeAssistantUsage1788600000000 implements MigrationInterface {
|
||||||
|
async up(queryRunner: QueryRunner): Promise<void> {
|
||||||
|
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<void> {
|
||||||
|
await queryRunner.query('DROP TABLE trade_assistant_usage');
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -0,0 +1,37 @@
|
|||||||
|
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||||
|
|
||||||
|
export class CreateTradeConversations1788700000000 implements MigrationInterface {
|
||||||
|
async up(queryRunner: QueryRunner): Promise<void> {
|
||||||
|
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<void> {
|
||||||
|
await queryRunner.query('DROP TABLE trade_messages');
|
||||||
|
await queryRunner.query('DROP TABLE trade_conversations');
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -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<TradeConversationSummary[]> {
|
||||||
|
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<TradeConversationSummary> {
|
||||||
|
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<TradeConversationSummary | null> {
|
||||||
|
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<TradeMessage[]> {
|
||||||
|
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<TradeMessage> {
|
||||||
|
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<void> {
|
||||||
|
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<void> {
|
||||||
|
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(),
|
||||||
|
});
|
||||||
@ -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');
|
||||||
|
});
|
||||||
|
});
|
||||||
@ -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<TradeUsage> {
|
||||||
|
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<boolean> {
|
||||||
|
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<void> {
|
||||||
|
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<void> {
|
||||||
|
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]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Reference in New Issue
Block a user