feat(#195): SSE realtime Cartes (stream + emit audience)
GET /api/v1/cards/stream (Bearer ou ?access_token=). Events card.created|updated|deleted, response.added + heartbeat. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -0,0 +1,111 @@
|
||||
import { Injectable, MessageEvent, OnModuleDestroy } from '@nestjs/common';
|
||||
import { Observable, Subject, interval, merge, takeUntil } from 'rxjs';
|
||||
import { map } from 'rxjs/operators';
|
||||
|
||||
export type CardRealtimeEventType =
|
||||
| 'card.created'
|
||||
| 'card.updated'
|
||||
| 'card.deleted'
|
||||
| 'response.added'
|
||||
| 'heartbeat';
|
||||
|
||||
export interface CardRealtimePayload {
|
||||
event: CardRealtimeEventType;
|
||||
card_id?: string;
|
||||
data?: unknown;
|
||||
at: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Bus SSE in-memory par utilisateur (V1 mono-instance).
|
||||
* Rooms = userId (audience carte).
|
||||
*/
|
||||
@Injectable()
|
||||
export class CardsRealtimeService implements OnModuleDestroy {
|
||||
private readonly byUser = new Map<string, Set<Subject<CardRealtimePayload>>>();
|
||||
private readonly destroy$ = new Subject<void>();
|
||||
|
||||
onModuleDestroy(): void {
|
||||
this.destroy$.next();
|
||||
this.destroy$.complete();
|
||||
for (const set of this.byUser.values()) {
|
||||
for (const s of set) s.complete();
|
||||
}
|
||||
this.byUser.clear();
|
||||
}
|
||||
|
||||
/** Flux SSE pour un utilisateur authentifié. */
|
||||
streamFor(userId: string): Observable<MessageEvent> {
|
||||
const subject = new Subject<CardRealtimePayload>();
|
||||
let set = this.byUser.get(userId);
|
||||
if (!set) {
|
||||
set = new Set();
|
||||
this.byUser.set(userId, set);
|
||||
}
|
||||
set.add(subject);
|
||||
|
||||
const heartbeat$ = interval(25_000).pipe(
|
||||
map(
|
||||
(): CardRealtimePayload => ({
|
||||
event: 'heartbeat',
|
||||
at: new Date().toISOString(),
|
||||
}),
|
||||
),
|
||||
);
|
||||
|
||||
return new Observable<MessageEvent>((observer) => {
|
||||
const sub = merge(subject, heartbeat$)
|
||||
.pipe(takeUntil(this.destroy$))
|
||||
.subscribe({
|
||||
next: (payload) =>
|
||||
observer.next({
|
||||
type: payload.event,
|
||||
data: payload,
|
||||
} as MessageEvent),
|
||||
error: (err) => observer.error(err),
|
||||
complete: () => observer.complete(),
|
||||
});
|
||||
|
||||
// ping initial
|
||||
subject.next({
|
||||
event: 'heartbeat',
|
||||
at: new Date().toISOString(),
|
||||
});
|
||||
|
||||
return () => {
|
||||
sub.unsubscribe();
|
||||
set!.delete(subject);
|
||||
subject.complete();
|
||||
if (set!.size === 0) this.byUser.delete(userId);
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
emitToUsers(
|
||||
userIds: string[],
|
||||
event: Exclude<CardRealtimeEventType, 'heartbeat'>,
|
||||
cardId: string | undefined,
|
||||
data?: unknown,
|
||||
): void {
|
||||
const unique = [...new Set(userIds.filter(Boolean))];
|
||||
const payload: CardRealtimePayload = {
|
||||
event,
|
||||
card_id: cardId,
|
||||
data,
|
||||
at: new Date().toISOString(),
|
||||
};
|
||||
for (const uid of unique) {
|
||||
const set = this.byUser.get(uid);
|
||||
if (!set) continue;
|
||||
for (const s of set) s.next(payload);
|
||||
}
|
||||
}
|
||||
|
||||
/** Test / debug : nombre d’abonnés actifs. */
|
||||
subscriberCount(userId?: string): number {
|
||||
if (userId) return this.byUser.get(userId)?.size ?? 0;
|
||||
let n = 0;
|
||||
for (const s of this.byUser.values()) n += s.size;
|
||||
return n;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user