GET /api/v1/cards/stream (Bearer ou ?access_token=). Events card.created|updated|deleted, response.added + heartbeat. Co-authored-by: Cursor <cursoragent@cursor.com>
112 lines
2.9 KiB
TypeScript
112 lines
2.9 KiB
TypeScript
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;
|
||
}
|
||
}
|