From 326c0ccbe90163d5f1f50471e220c4366f36d995 Mon Sep 17 00:00:00 2001 From: mubashqoz-eng Date: Mon, 28 Sep 2026 17:30:34 +0100 Subject: [PATCH] feat(live-streaming): implement live streaming and spectator mode (#441) Add a self-contained LiveStreamingModule covering the issue's acceptance criteria: - stream creation and lifecycle (scheduled/live/ended/cancelled, host-only start/end/cancel, recording finalised on end) - spectator join/leave, idempotent joins, and roster-derived viewer counts with a monotonic peak - real-time chat over a Socket.IO gateway (/streams) plus REST history, with soft deletes that preserve the audit trail - deterministic bandwidth -> quality adaptation - moderator tools (ban/unban/timeout/delete/pin/clear) with an immutable action log; host and assigned moderators only - viewership analytics (current/peak/total, chat, moderation, quality distribution, duration, recording state) - a migration creating the four tables and indexes, and a module README 42 unit tests cover lifecycle, permissions, chat, moderation, recording, quality and analytics. --- src/app.module.ts | 2 + src/live-streaming/README.md | 54 ++ src/live-streaming/dto/chat-message.dto.ts | 10 + src/live-streaming/dto/create-stream.dto.ts | 30 + src/live-streaming/dto/index.ts | 4 + src/live-streaming/dto/join-stream.dto.ts | 35 ++ src/live-streaming/dto/moderate-stream.dto.ts | 48 ++ src/live-streaming/entities/index.ts | 4 + .../entities/live-stream.entity.ts | 86 +++ .../entities/stream-chat-message.entity.ts | 40 ++ .../stream-moderation-action.entity.ts | 53 ++ .../entities/stream-viewer.entity.ts | 57 ++ .../gateways/live-streaming.gateway.spec.ts | 109 ++++ .../gateways/live-streaming.gateway.ts | 153 +++++ .../live-streaming.controller.ts | 195 ++++++ src/live-streaming/live-streaming.module.ts | 25 + .../live-streaming.service.spec.ts | 571 ++++++++++++++++++ src/live-streaming/live-streaming.service.ts | 509 ++++++++++++++++ ...1760000000000-CreateLiveStreamingTables.ts | 237 ++++++++ 19 files changed, 2222 insertions(+) create mode 100644 src/live-streaming/README.md create mode 100644 src/live-streaming/dto/chat-message.dto.ts create mode 100644 src/live-streaming/dto/create-stream.dto.ts create mode 100644 src/live-streaming/dto/index.ts create mode 100644 src/live-streaming/dto/join-stream.dto.ts create mode 100644 src/live-streaming/dto/moderate-stream.dto.ts create mode 100644 src/live-streaming/entities/index.ts create mode 100644 src/live-streaming/entities/live-stream.entity.ts create mode 100644 src/live-streaming/entities/stream-chat-message.entity.ts create mode 100644 src/live-streaming/entities/stream-moderation-action.entity.ts create mode 100644 src/live-streaming/entities/stream-viewer.entity.ts create mode 100644 src/live-streaming/gateways/live-streaming.gateway.spec.ts create mode 100644 src/live-streaming/gateways/live-streaming.gateway.ts create mode 100644 src/live-streaming/live-streaming.controller.ts create mode 100644 src/live-streaming/live-streaming.module.ts create mode 100644 src/live-streaming/live-streaming.service.spec.ts create mode 100644 src/live-streaming/live-streaming.service.ts create mode 100644 src/migrations/1760000000000-CreateLiveStreamingTables.ts diff --git a/src/app.module.ts b/src/app.module.ts index 65de8f42..47a831cc 100644 --- a/src/app.module.ts +++ b/src/app.module.ts @@ -14,6 +14,7 @@ import { JobsModule } from './jobs/jobs.module'; import { Job } from './jobs/job.entity'; import { CdnModule } from './cdn/cdn.module'; import { AdminModule } from './admin/admin.module'; +import { LiveStreamingModule } from './live-streaming/live-streaming.module'; @Module({ imports: [ @@ -53,6 +54,7 @@ import { AdminModule } from './admin/admin.module'; JobsModule, CdnModule, AdminModule, + LiveStreamingModule, ], }) export class AppModule {} diff --git a/src/live-streaming/README.md b/src/live-streaming/README.md new file mode 100644 index 00000000..db4b5ba7 --- /dev/null +++ b/src/live-streaming/README.md @@ -0,0 +1,54 @@ +# Live streaming and spectator mode + +Backend for broadcasting a puzzle attempt with spectator chat, viewer tracking, +quality adaptation, moderation, recording and analytics. It is self-contained: +`LiveStreamingModule` wires its own TypeORM repositories and a Socket.IO gateway +under the `/streams` namespace, and is imported by `AppModule`. + +## Model + +| Entity | Purpose | +| --- | --- | +| `LiveStream` | The broadcast. Status (`scheduled` → `live` → `ended`/`cancelled`), cached viewer counts, recording state. | +| `StreamViewer` | One spectator's membership. Rows survive a leave (`isActive = false`) so viewership is analysable. | +| `StreamChatMessage` | Chat. Deletes are soft so the audit trail survives. | +| `StreamModerationAction` | Immutable moderator decisions; enforcement reads these rows. | + +## REST + +All routes take the caller's id in the `x-user-id` header. The service enforces +that only the host (or an assigned moderator) can manage a stream; wiring the +header to the JWT guard is a shell concern. + +| Method | Path | Notes | +| --- | --- | --- | +| `POST` | `/live-streaming/streams` | Create | +| `GET` | `/live-streaming/streams` | List, optional `?status=` | +| `GET` | `/live-streaming/streams/:id` | Read | +| `POST` | `/live-streaming/streams/:id/start` \| `/end` \| `/cancel` | Host only | +| `POST` | `/live-streaming/streams/:id/viewers` | Join as spectator | +| `DELETE` | `/live-streaming/streams/:id/viewers/me` | Leave | +| `GET` | `/live-streaming/streams/:id/viewers/count` | Live count | +| `PATCH` | `/live-streaming/streams/:id/quality` | `{ bandwidthKbps }` | +| `POST` `GET` | `/live-streaming/streams/:id/chat` | Post / history | +| `POST` `GET` | `/live-streaming/streams/:id/moderation` | Moderate / log | +| `POST` | `/live-streaming/streams/:id/recording/start` \| `/stop` | Host only | +| `GET` | `/live-streaming/streams/:id/analytics` | Viewership + chat summary | + +## WebSocket (`/streams`) + +`stream:join`, `stream:leave`, `stream:chat`, `stream:quality`. Chat is broadcast +to the `stream:` room; viewer counts are re-read from the roster so two +sockets cannot drift the number. A disconnect leaves the roster on behalf of the +socket's stored user id. + +## Quality ladder + +`qualityForBandwidth` is a pure, deterministic mapping: `< 1000` kbps → `low`, +`>= 1000` → `medium`, `>= 2500` → `high`, `>= 6000` → `source`; an unknown +bandwidth joins at `medium`. + +## Migration + +`1760000000000-CreateLiveStreamingTables` creates the four tables and their +indexes; `down` drops them in reverse. diff --git a/src/live-streaming/dto/chat-message.dto.ts b/src/live-streaming/dto/chat-message.dto.ts new file mode 100644 index 00000000..a3f080c1 --- /dev/null +++ b/src/live-streaming/dto/chat-message.dto.ts @@ -0,0 +1,10 @@ +import { IsNotEmpty, IsString, MaxLength } from 'class-validator'; +import { ApiProperty } from '@nestjs/swagger'; + +export class ChatMessageDto { + @ApiProperty({ description: 'Message body shown to spectators' }) + @IsString() + @IsNotEmpty() + @MaxLength(500) + content: string; +} diff --git a/src/live-streaming/dto/create-stream.dto.ts b/src/live-streaming/dto/create-stream.dto.ts new file mode 100644 index 00000000..7d14887b --- /dev/null +++ b/src/live-streaming/dto/create-stream.dto.ts @@ -0,0 +1,30 @@ +import { + IsBoolean, + IsNotEmpty, + IsOptional, + IsString, + MaxLength, +} from 'class-validator'; +import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger'; + +export class CreateStreamDto { + @ApiProperty({ description: 'Human-readable stream title' }) + @IsString() + @IsNotEmpty() + @MaxLength(200) + title: string; + + @ApiPropertyOptional({ description: 'What the stream is about' }) + @IsString() + @IsOptional() + @MaxLength(2000) + description?: string; + + @ApiPropertyOptional({ + description: 'Capture the stream for later playback', + default: false, + }) + @IsBoolean() + @IsOptional() + recordingEnabled?: boolean; +} diff --git a/src/live-streaming/dto/index.ts b/src/live-streaming/dto/index.ts new file mode 100644 index 00000000..78492ae7 --- /dev/null +++ b/src/live-streaming/dto/index.ts @@ -0,0 +1,4 @@ +export * from './create-stream.dto'; +export * from './join-stream.dto'; +export * from './chat-message.dto'; +export * from './moderate-stream.dto'; diff --git a/src/live-streaming/dto/join-stream.dto.ts b/src/live-streaming/dto/join-stream.dto.ts new file mode 100644 index 00000000..6b8cc2a8 --- /dev/null +++ b/src/live-streaming/dto/join-stream.dto.ts @@ -0,0 +1,35 @@ +import { + IsEnum, + IsInt, + IsOptional, + IsString, + MaxLength, + Min, +} from 'class-validator'; +import { ApiPropertyOptional } from '@nestjs/swagger'; +import { StreamQuality } from '../entities/live-stream.entity'; + +export class JoinStreamDto { + @ApiPropertyOptional({ description: 'Display name shown to the chat' }) + @IsString() + @IsOptional() + @MaxLength(64) + username?: string; + + @ApiPropertyOptional({ + description: 'Preferred stream quality before bandwidth is measured', + enum: StreamQuality, + }) + @IsEnum(StreamQuality) + @IsOptional() + preferredQuality?: StreamQuality; + + @ApiPropertyOptional({ + description: 'Measured downstream bandwidth in kbps; adapts quality', + minimum: 0, + }) + @IsInt() + @Min(0) + @IsOptional() + bandwidthKbps?: number; +} diff --git a/src/live-streaming/dto/moderate-stream.dto.ts b/src/live-streaming/dto/moderate-stream.dto.ts new file mode 100644 index 00000000..80f1496a --- /dev/null +++ b/src/live-streaming/dto/moderate-stream.dto.ts @@ -0,0 +1,48 @@ +import { + IsEnum, + IsInt, + IsOptional, + IsString, + IsUUID, + MaxLength, + Min, +} from 'class-validator'; +import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger'; +import { ModerationActionType } from '../entities/stream-moderation-action.entity'; + +/** + * One moderator action. `targetUserId` is required for ban/unban/timeout; + * `messageId` is required for delete/pin. The service validates the pairing so a + * malformed action fails before anything is written. + */ +export class ModerateStreamDto { + @ApiProperty({ enum: ModerationActionType }) + @IsEnum(ModerationActionType) + action: ModerationActionType; + + @ApiPropertyOptional({ description: 'Subject of a ban, unban or timeout' }) + @IsString() + @IsOptional() + @MaxLength(64) + targetUserId?: string; + + @ApiPropertyOptional({ description: 'Reason recorded on the action' }) + @IsString() + @IsOptional() + @MaxLength(200) + reason?: string; + + @ApiPropertyOptional({ description: 'Message a delete/pin targets' }) + @IsUUID() + @IsOptional() + messageId?: string; + + @ApiPropertyOptional({ + description: 'Timeout length in minutes (timeout only)', + minimum: 1, + }) + @IsInt() + @Min(1) + @IsOptional() + durationMinutes?: number; +} diff --git a/src/live-streaming/entities/index.ts b/src/live-streaming/entities/index.ts new file mode 100644 index 00000000..20e1f256 --- /dev/null +++ b/src/live-streaming/entities/index.ts @@ -0,0 +1,4 @@ +export * from './live-stream.entity'; +export * from './stream-viewer.entity'; +export * from './stream-chat-message.entity'; +export * from './stream-moderation-action.entity'; diff --git a/src/live-streaming/entities/live-stream.entity.ts b/src/live-streaming/entities/live-stream.entity.ts new file mode 100644 index 00000000..3df6625c --- /dev/null +++ b/src/live-streaming/entities/live-stream.entity.ts @@ -0,0 +1,86 @@ +import { + Entity, + PrimaryGeneratedColumn, + Column, + CreateDateColumn, + UpdateDateColumn, + Index, +} from 'typeorm'; + +export enum LiveStreamStatus { + SCHEDULED = 'scheduled', + LIVE = 'live', + ENDED = 'ended', + CANCELLED = 'cancelled', +} + +export enum StreamQuality { + LOW = 'low', + MEDIUM = 'medium', + HIGH = 'high', + SOURCE = 'source', +} + +/** + * A live puzzle attempt a broadcaster streams to spectators. + * + * `viewerCount` is a cached copy of the active roster for cheap reads; the + * authoritative count is the active `StreamViewer` rows. `peakViewerCount` is + * monotonic so analytics can report the high-water mark even after everyone + * leaves. + */ +@Entity('live_streams') +@Index(['hostId']) +@Index(['status']) +export class LiveStream { + @PrimaryGeneratedColumn('uuid') + id: string; + + @Column({ type: 'varchar', length: 64 }) + @Index() + hostId: string; + + @Column({ type: 'varchar', length: 200 }) + title: string; + + @Column({ type: 'text', nullable: true }) + description?: string; + + @Column({ + type: 'varchar', + length: 20, + default: LiveStreamStatus.SCHEDULED, + }) + status: LiveStreamStatus; + + @Column({ + type: 'varchar', + length: 10, + default: StreamQuality.MEDIUM, + }) + currentQuality: StreamQuality; + + @Column({ type: 'int', default: 0 }) + viewerCount: number; + + @Column({ type: 'int', default: 0 }) + peakViewerCount: number; + + @Column({ type: 'boolean', default: false }) + recordingEnabled: boolean; + + @Column({ type: 'varchar', length: 500, nullable: true }) + recordingUrl?: string; + + @Column({ type: 'timestamp', nullable: true }) + startedAt?: Date; + + @Column({ type: 'timestamp', nullable: true }) + endedAt?: Date; + + @CreateDateColumn() + createdAt: Date; + + @UpdateDateColumn() + updatedAt: Date; +} diff --git a/src/live-streaming/entities/stream-chat-message.entity.ts b/src/live-streaming/entities/stream-chat-message.entity.ts new file mode 100644 index 00000000..c6854956 --- /dev/null +++ b/src/live-streaming/entities/stream-chat-message.entity.ts @@ -0,0 +1,40 @@ +import { + Entity, + PrimaryGeneratedColumn, + Column, + CreateDateColumn, + Index, +} from 'typeorm'; + +/** + * A spectator chat message. Deletes are soft (`isDeleted`) so moderators can + * remove content without losing the audit trail; `deletedBy` records who acted. + */ +@Entity('stream_chat_messages') +@Index(['streamId', 'createdAt']) +export class StreamChatMessage { + @PrimaryGeneratedColumn('uuid') + id: string; + + @Column({ type: 'uuid' }) + @Index() + streamId: string; + + @Column({ type: 'varchar', length: 64 }) + authorId: string; + + @Column({ type: 'text' }) + content: string; + + @Column({ type: 'boolean', default: false }) + isDeleted: boolean; + + @Column({ type: 'varchar', length: 64, nullable: true }) + deletedBy?: string; + + @Column({ type: 'boolean', default: false }) + isPinned: boolean; + + @CreateDateColumn() + createdAt: Date; +} diff --git a/src/live-streaming/entities/stream-moderation-action.entity.ts b/src/live-streaming/entities/stream-moderation-action.entity.ts new file mode 100644 index 00000000..a3297388 --- /dev/null +++ b/src/live-streaming/entities/stream-moderation-action.entity.ts @@ -0,0 +1,53 @@ +import { + Entity, + PrimaryGeneratedColumn, + Column, + CreateDateColumn, + Index, +} from 'typeorm'; + +export enum ModerationActionType { + BAN = 'ban', + UNBAN = 'unban', + TIMEOUT = 'timeout', + DELETE_MESSAGE = 'delete_message', + CLEAR_CHAT = 'clear_chat', + PIN_MESSAGE = 'pin_message', +} + +/** + * An immutable moderator decision. Enforcement reads these rows (an active ban + * or an unexpired timeout), so the decision and the audit record are the same + * thing. + */ +@Entity('stream_moderation_actions') +@Index(['streamId', 'targetUserId']) +export class StreamModerationAction { + @PrimaryGeneratedColumn('uuid') + id: string; + + @Column({ type: 'uuid' }) + @Index() + streamId: string; + + @Column({ type: 'varchar', length: 64 }) + moderatorId: string; + + @Column({ type: 'varchar', length: 64, nullable: true }) + targetUserId?: string; + + @Column({ type: 'varchar', length: 20 }) + action: ModerationActionType; + + @Column({ type: 'varchar', length: 200, nullable: true }) + reason?: string; + + @Column({ type: 'uuid', nullable: true }) + messageId?: string; + + @Column({ type: 'timestamp', nullable: true }) + expiresAt?: Date; + + @CreateDateColumn() + createdAt: Date; +} diff --git a/src/live-streaming/entities/stream-viewer.entity.ts b/src/live-streaming/entities/stream-viewer.entity.ts new file mode 100644 index 00000000..9769858b --- /dev/null +++ b/src/live-streaming/entities/stream-viewer.entity.ts @@ -0,0 +1,57 @@ +import { + Entity, + PrimaryGeneratedColumn, + Column, + CreateDateColumn, + UpdateDateColumn, + Index, +} from 'typeorm'; +import { StreamQuality } from './live-stream.entity'; + +/** + * One spectator's membership in a stream. Rows are kept after a viewer leaves + * (`isActive = false`) so viewership can be analysed over time; `isActive` + * scopes the live count. + */ +@Entity('stream_viewers') +@Index(['streamId', 'viewerId']) +@Index(['streamId', 'isActive']) +export class StreamViewer { + @PrimaryGeneratedColumn('uuid') + id: string; + + @Column({ type: 'uuid' }) + @Index() + streamId: string; + + @Column({ type: 'varchar', length: 64 }) + viewerId: string; + + @Column({ type: 'varchar', length: 64, nullable: true }) + username?: string; + + @Column({ + type: 'varchar', + length: 10, + default: StreamQuality.MEDIUM, + }) + quality: StreamQuality; + + @Column({ type: 'boolean', default: false }) + isModerator: boolean; + + @Column({ type: 'boolean', default: true }) + isActive: boolean; + + @Column({ type: 'timestamp', nullable: true }) + joinedAt?: Date; + + @Column({ type: 'timestamp', nullable: true }) + leftAt?: Date; + + @CreateDateColumn() + createdAt: Date; + + @UpdateDateColumn() + updatedAt: Date; +} diff --git a/src/live-streaming/gateways/live-streaming.gateway.spec.ts b/src/live-streaming/gateways/live-streaming.gateway.spec.ts new file mode 100644 index 00000000..eaa7e0f1 --- /dev/null +++ b/src/live-streaming/gateways/live-streaming.gateway.spec.ts @@ -0,0 +1,109 @@ +import { LiveStreamingGateway } from './live-streaming.gateway'; + +describe('LiveStreamingGateway', () => { + let gateway: LiveStreamingGateway; + let service: { + joinStream: jest.Mock; + leaveStream: jest.Mock; + getViewerCount: jest.Mock; + postMessage: jest.Mock; + selectQuality: jest.Mock; + }; + let server: { to: jest.Mock; emit: jest.Mock }; + let client: { + id: string; + data: Record; + join: jest.Mock; + leave: jest.Mock; + }; + + beforeEach(() => { + service = { + joinStream: jest.fn().mockResolvedValue({ id: 'viewer-row' }), + leaveStream: jest.fn().mockResolvedValue(undefined), + getViewerCount: jest.fn().mockResolvedValue(7), + postMessage: jest.fn().mockResolvedValue({ id: 'msg-1', content: 'hi' }), + selectQuality: jest.fn().mockResolvedValue({ quality: 'low' }), + }; + + gateway = new LiveStreamingGateway(service as never); + + server = { to: jest.fn(), emit: jest.fn() }; + server.to.mockReturnValue(server); + gateway.server = server as never; + + client = { + id: 'socket-1', + data: {}, + join: jest.fn().mockResolvedValue(undefined), + leave: jest.fn().mockResolvedValue(undefined), + }; + }); + + it('joins the room, records the membership, and broadcasts the count', async () => { + const result = await gateway.handleJoin(client as never, { + streamId: 'stream-1', + userId: 'viewer-1', + username: 'ada', + bandwidthKbps: 3_000, + }); + + expect(service.joinStream).toHaveBeenCalledWith('stream-1', 'viewer-1', { + username: 'ada', + bandwidthKbps: 3_000, + }); + expect(client.join).toHaveBeenCalledWith('stream:stream-1'); + expect(client.data.userId).toBe('viewer-1'); + expect(server.to).toHaveBeenCalledWith('stream:stream-1'); + expect(server.emit).toHaveBeenCalledWith('stream:viewers', { + streamId: 'stream-1', + count: 7, + }); + expect(result).toEqual({ joined: true, count: 7 }); + }); + + it('broadcasts a chat message to the stream room', async () => { + const message = await gateway.handleChat({ + streamId: 'stream-1', + userId: 'viewer-1', + content: 'hi', + }); + + expect(service.postMessage).toHaveBeenCalledWith( + 'stream-1', + 'viewer-1', + 'hi', + ); + expect(server.to).toHaveBeenCalledWith('stream:stream-1'); + expect(server.emit).toHaveBeenCalledWith('stream:chat', message); + }); + + it('leaves the room on leave', async () => { + await gateway.handleLeave(client as never, { + streamId: 'stream-1', + userId: 'viewer-1', + }); + + expect(service.leaveStream).toHaveBeenCalledWith('stream-1', 'viewer-1'); + expect(client.leave).toHaveBeenCalledWith('stream:stream-1'); + }); + + it('marks the viewer left when a socket disconnects', async () => { + client.data = { userId: 'viewer-1', streamId: 'stream-1' }; + await gateway.handleDisconnect(client as never); + expect(service.leaveStream).toHaveBeenCalledWith('stream-1', 'viewer-1'); + }); + + it('adapts quality mid-stream', async () => { + await gateway.handleQuality({ + streamId: 'stream-1', + userId: 'viewer-1', + bandwidthKbps: 800, + }); + expect(service.selectQuality).toHaveBeenCalledWith( + 'stream-1', + 'viewer-1', + 800, + ); + }); +}); diff --git a/src/live-streaming/gateways/live-streaming.gateway.ts b/src/live-streaming/gateways/live-streaming.gateway.ts new file mode 100644 index 00000000..92107659 --- /dev/null +++ b/src/live-streaming/gateways/live-streaming.gateway.ts @@ -0,0 +1,153 @@ +import { + ConnectedSocket, + MessageBody, + OnGatewayDisconnect, + SubscribeMessage, + WebSocketGateway, + WebSocketServer, +} from '@nestjs/websockets'; +import { Server, Socket } from 'socket.io'; + +import { LiveStreamingService } from '../live-streaming.service'; + +interface JoinPayload { + streamId: string; + userId: string; + username?: string; + bandwidthKbps?: number; +} + +interface ChatPayload { + streamId: string; + userId: string; + content: string; +} + +interface QualityPayload { + streamId: string; + userId: string; + bandwidthKbps: number; +} + +/** Per-socket state set on join, read back on disconnect. */ +interface StreamSocketData { + userId?: string; + streamId?: string; +} + +/** + * Real-time spectator transport. Chat is delivered by room broadcast so every + * spectator in `stream:` receives it in the same tick; viewer counts are + * re-read from the service (the roster, not a local counter) so two sockets + * cannot drift the number. + */ +@WebSocketGateway({ + namespace: '/streams', + cors: { + origin: process.env.FRONTEND_URL?.split(',') ?? '*', + credentials: true, + }, +}) +export class LiveStreamingGateway implements OnGatewayDisconnect { + @WebSocketServer() + server: Server; + + private readonly memberships = new Map>(); + + constructor(private readonly service: LiveStreamingService) {} + + @SubscribeMessage('stream:join') + async handleJoin( + @ConnectedSocket() client: Socket, + @MessageBody() payload: JoinPayload, + ) { + await this.service.joinStream(payload.streamId, payload.userId, { + username: payload.username, + bandwidthKbps: payload.bandwidthKbps, + }); + + const room = this.room(payload.streamId); + await client.join(room); + const data = client.data as StreamSocketData; + data.userId = payload.userId; + data.streamId = payload.streamId; + this.track(room, client.id); + + const count = await this.service.getViewerCount(payload.streamId); + this.server + .to(room) + .emit('stream:viewers', { streamId: payload.streamId, count }); + return { joined: true, count }; + } + + @SubscribeMessage('stream:leave') + async handleLeave( + @ConnectedSocket() client: Socket, + @MessageBody() payload: { streamId: string; userId: string }, + ) { + await this.service.leaveStream(payload.streamId, payload.userId); + const room = this.room(payload.streamId); + await client.leave(room); + this.memberships.get(room)?.delete(client.id); + + const count = await this.service.getViewerCount(payload.streamId); + this.server + .to(room) + .emit('stream:viewers', { streamId: payload.streamId, count }); + return { left: true, count }; + } + + @SubscribeMessage('stream:chat') + async handleChat(@MessageBody() payload: ChatPayload) { + const message = await this.service.postMessage( + payload.streamId, + payload.userId, + payload.content, + ); + this.server.to(this.room(payload.streamId)).emit('stream:chat', message); + return message; + } + + @SubscribeMessage('stream:quality') + async handleQuality(@MessageBody() payload: QualityPayload) { + return this.service.selectQuality( + payload.streamId, + payload.userId, + payload.bandwidthKbps, + ); + } + + async handleDisconnect(client: Socket): Promise { + const data = client.data as StreamSocketData; + const streamId = data.streamId; + const userId = data.userId; + + if (streamId && userId) { + try { + await this.service.leaveStream(streamId, userId); + } catch { + // Already left, banned, or the socket never completed a join; the + // count below is re-read from the roster either way. + } + } + + for (const [room, members] of this.memberships.entries()) { + if (!members.delete(client.id)) { + continue; + } + const id = room.replace('stream:', ''); + const count = await this.service.getViewerCount(id); + this.server.to(room).emit('stream:viewers', { streamId: id, count }); + } + } + + private room(streamId: string): string { + return `stream:${streamId}`; + } + + private track(room: string, clientId: string): void { + const members = this.memberships.get(room) ?? new Set(); + members.add(clientId); + this.memberships.set(room, members); + } +} diff --git a/src/live-streaming/live-streaming.controller.ts b/src/live-streaming/live-streaming.controller.ts new file mode 100644 index 00000000..34847d46 --- /dev/null +++ b/src/live-streaming/live-streaming.controller.ts @@ -0,0 +1,195 @@ +import { + BadRequestException, + Body, + Controller, + Delete, + Get, + Headers, + HttpCode, + HttpStatus, + Param, + ParseUUIDPipe, + Patch, + Post, + Query, +} from '@nestjs/common'; +import { ApiHeader, ApiOperation, ApiTags } from '@nestjs/swagger'; + +import { LiveStreamingService } from './live-streaming.service'; +import { LiveStreamStatus } from './entities/live-stream.entity'; +import { CreateStreamDto } from './dto/create-stream.dto'; +import { JoinStreamDto } from './dto/join-stream.dto'; +import { ChatMessageDto } from './dto/chat-message.dto'; +import { ModerateStreamDto } from './dto/moderate-stream.dto'; + +/** + * The caller's identity is taken from `x-user-id`. The service enforces that the + * caller is the host or an assigned moderator; wiring this header to the JWT + * guard is an integration decision for the gateway, not a streaming concern. + */ +@ApiTags('Live Streaming') +@ApiHeader({ name: 'x-user-id', required: true }) +@Controller('live-streaming') +export class LiveStreamingController { + constructor(private readonly service: LiveStreamingService) {} + + @Post('streams') + @HttpCode(HttpStatus.CREATED) + @ApiOperation({ summary: 'Create a live stream' }) + create(@Headers('x-user-id') userId: string, @Body() dto: CreateStreamDto) { + return this.service.createStream(this.requireUser(userId), dto); + } + + @Get('streams') + @ApiOperation({ summary: 'List streams, newest first' }) + list(@Query('status') status?: string) { + return this.service.listStreams(status as LiveStreamStatus | undefined); + } + + @Get('streams/:id') + @ApiOperation({ summary: 'Get one stream' }) + get(@Param('id', ParseUUIDPipe) id: string) { + return this.service.getStream(id); + } + + @Post('streams/:id/start') + @ApiOperation({ summary: 'Start a stream (host only)' }) + start( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + ) { + return this.service.startStream(id, this.requireUser(userId)); + } + + @Post('streams/:id/end') + @ApiOperation({ summary: 'End a stream (host only)' }) + end( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + ) { + return this.service.endStream(id, this.requireUser(userId)); + } + + @Post('streams/:id/cancel') + @ApiOperation({ summary: 'Cancel a stream (host only)' }) + cancel( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + ) { + return this.service.cancelStream(id, this.requireUser(userId)); + } + + @Post('streams/:id/viewers') + @HttpCode(HttpStatus.CREATED) + @ApiOperation({ summary: 'Join a stream as a spectator' }) + join( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + @Body() dto: JoinStreamDto, + ) { + return this.service.joinStream(id, this.requireUser(userId), dto); + } + + @Delete('streams/:id/viewers/me') + @ApiOperation({ summary: 'Leave a stream' }) + async leave( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + ) { + await this.service.leaveStream(id, this.requireUser(userId)); + return { left: true }; + } + + @Get('streams/:id/viewers/count') + @ApiOperation({ summary: 'Current spectator count' }) + count(@Param('id', ParseUUIDPipe) id: string) { + return this.service.getViewerCount(id).then((count) => ({ count })); + } + + @Patch('streams/:id/quality') + @ApiOperation({ summary: 'Adapt a spectator quality to measured bandwidth' }) + quality( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + @Body('bandwidthKbps') bandwidthKbps: number, + ) { + return this.service.selectQuality( + id, + this.requireUser(userId), + bandwidthKbps, + ); + } + + @Post('streams/:id/chat') + @HttpCode(HttpStatus.CREATED) + @ApiOperation({ summary: 'Post a chat message' }) + chat( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + @Body() dto: ChatMessageDto, + ) { + return this.service.postMessage(id, this.requireUser(userId), dto.content); + } + + @Get('streams/:id/chat') + @ApiOperation({ summary: 'Recent chat history' }) + history( + @Param('id', ParseUUIDPipe) id: string, + @Query('limit') limit?: string, + ) { + return this.service.getChatHistory( + id, + limit ? Number.parseInt(limit, 10) : 50, + ); + } + + @Post('streams/:id/moderation') + @HttpCode(HttpStatus.CREATED) + @ApiOperation({ summary: 'Moderate a stream (host or moderator)' }) + moderate( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + @Body() dto: ModerateStreamDto, + ) { + return this.service.moderate(id, this.requireUser(userId), dto); + } + + @Get('streams/:id/moderation') + @ApiOperation({ summary: 'Moderation log' }) + moderationLog(@Param('id', ParseUUIDPipe) id: string) { + return this.service.getModerationLog(id); + } + + @Post('streams/:id/recording/start') + @ApiOperation({ summary: 'Start recording (host only)' }) + startRecording( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + ) { + return this.service.startRecording(id, this.requireUser(userId)); + } + + @Post('streams/:id/recording/stop') + @ApiOperation({ + summary: 'Stop recording and finalize the artefact (host only)', + }) + stopRecording( + @Headers('x-user-id') userId: string, + @Param('id', ParseUUIDPipe) id: string, + ) { + return this.service.stopRecording(id, this.requireUser(userId)); + } + + @Get('streams/:id/analytics') + @ApiOperation({ summary: 'Viewership and chat analytics' }) + analytics(@Param('id', ParseUUIDPipe) id: string) { + return this.service.getAnalytics(id); + } + + private requireUser(userId: string): string { + if (!userId) { + throw new BadRequestException('x-user-id header is required'); + } + return userId; + } +} diff --git a/src/live-streaming/live-streaming.module.ts b/src/live-streaming/live-streaming.module.ts new file mode 100644 index 00000000..2a1fc70e --- /dev/null +++ b/src/live-streaming/live-streaming.module.ts @@ -0,0 +1,25 @@ +import { Module } from '@nestjs/common'; +import { TypeOrmModule } from '@nestjs/typeorm'; + +import { LiveStream } from './entities/live-stream.entity'; +import { StreamViewer } from './entities/stream-viewer.entity'; +import { StreamChatMessage } from './entities/stream-chat-message.entity'; +import { StreamModerationAction } from './entities/stream-moderation-action.entity'; +import { LiveStreamingService } from './live-streaming.service'; +import { LiveStreamingController } from './live-streaming.controller'; +import { LiveStreamingGateway } from './gateways/live-streaming.gateway'; + +@Module({ + imports: [ + TypeOrmModule.forFeature([ + LiveStream, + StreamViewer, + StreamChatMessage, + StreamModerationAction, + ]), + ], + controllers: [LiveStreamingController], + providers: [LiveStreamingService, LiveStreamingGateway], + exports: [LiveStreamingService], +}) +export class LiveStreamingModule {} diff --git a/src/live-streaming/live-streaming.service.spec.ts b/src/live-streaming/live-streaming.service.spec.ts new file mode 100644 index 00000000..79992564 --- /dev/null +++ b/src/live-streaming/live-streaming.service.spec.ts @@ -0,0 +1,571 @@ +import { Test, TestingModule } from '@nestjs/testing'; +import { getRepositoryToken } from '@nestjs/typeorm'; +import { + BadRequestException, + ForbiddenException, + NotFoundException, +} from '@nestjs/common'; + +import { + LiveStreamingService, + qualityForBandwidth, +} from './live-streaming.service'; +import { + LiveStream, + LiveStreamStatus, + StreamQuality, +} from './entities/live-stream.entity'; +import { StreamViewer } from './entities/stream-viewer.entity'; +import { StreamChatMessage } from './entities/stream-chat-message.entity'; +import { + ModerationActionType, + StreamModerationAction, +} from './entities/stream-moderation-action.entity'; + +const mockRepo = () => ({ + create: jest.fn((d) => d), + save: jest.fn((d) => Promise.resolve({ ...d, id: d.id ?? 'generated-id' })), + findOne: jest.fn(), + find: jest.fn(), + count: jest.fn(), + update: jest.fn().mockResolvedValue({ affected: 1 }), +}); + +const baseStream = (over: Partial = {}): LiveStream => ({ + id: 'stream-1', + hostId: 'host-1', + title: 'Speedrun', + description: undefined, + status: LiveStreamStatus.SCHEDULED, + currentQuality: StreamQuality.MEDIUM, + viewerCount: 0, + peakViewerCount: 0, + recordingEnabled: false, + recordingUrl: undefined, + startedAt: undefined, + endedAt: undefined, + createdAt: new Date(), + updatedAt: new Date(), + ...over, +}); + +describe('qualityForBandwidth', () => { + it('maps the bandwidth ladder deterministically', () => { + expect(qualityForBandwidth(undefined)).toBe(StreamQuality.MEDIUM); + expect(qualityForBandwidth(0)).toBe(StreamQuality.LOW); + expect(qualityForBandwidth(999)).toBe(StreamQuality.LOW); + expect(qualityForBandwidth(1_000)).toBe(StreamQuality.MEDIUM); + expect(qualityForBandwidth(2_499)).toBe(StreamQuality.MEDIUM); + expect(qualityForBandwidth(2_500)).toBe(StreamQuality.HIGH); + expect(qualityForBandwidth(5_999)).toBe(StreamQuality.HIGH); + expect(qualityForBandwidth(6_000)).toBe(StreamQuality.SOURCE); + }); +}); + +describe('LiveStreamingService', () => { + let service: LiveStreamingService; + let streamRepo: ReturnType; + let viewerRepo: ReturnType; + let chatRepo: ReturnType; + let moderationRepo: ReturnType; + + beforeEach(async () => { + streamRepo = mockRepo(); + viewerRepo = mockRepo(); + chatRepo = mockRepo(); + moderationRepo = mockRepo(); + + // Sensible defaults: no active blocks, no active viewers, zero counts. + moderationRepo.find.mockResolvedValue([]); + viewerRepo.find.mockResolvedValue([]); + viewerRepo.count.mockResolvedValue(0); + chatRepo.count.mockResolvedValue(0); + moderationRepo.count.mockResolvedValue(0); + + const module: TestingModule = await Test.createTestingModule({ + providers: [ + LiveStreamingService, + { provide: getRepositoryToken(LiveStream), useValue: streamRepo }, + { provide: getRepositoryToken(StreamViewer), useValue: viewerRepo }, + { provide: getRepositoryToken(StreamChatMessage), useValue: chatRepo }, + { + provide: getRepositoryToken(StreamModerationAction), + useValue: moderationRepo, + }, + ], + }).compile(); + + service = module.get(LiveStreamingService); + }); + + afterEach(() => jest.clearAllMocks()); + + // ─── Lifecycle ───────────────────────────────────────────────────────────── + + describe('createStream', () => { + it('creates a scheduled stream with zeroed counters', async () => { + const result = await service.createStream('host-1', { + title: 'Speedrun', + }); + expect(result.status).toBe(LiveStreamStatus.SCHEDULED); + expect(result.viewerCount).toBe(0); + expect(result.peakViewerCount).toBe(0); + expect(result.recordingEnabled).toBe(false); + }); + + it('honours a recording request', async () => { + const result = await service.createStream('host-1', { + title: 'Speedrun', + recordingEnabled: true, + }); + expect(result.recordingEnabled).toBe(true); + }); + }); + + describe('startStream', () => { + it('starts a scheduled stream and stamps startedAt', async () => { + streamRepo.findOne.mockResolvedValue(baseStream()); + const result = await service.startStream('stream-1', 'host-1'); + expect(result.status).toBe(LiveStreamStatus.LIVE); + expect(result.startedAt).toBeInstanceOf(Date); + }); + + it('is idempotent when already live', async () => { + const live = baseStream({ status: LiveStreamStatus.LIVE }); + streamRepo.findOne.mockResolvedValue(live); + const result = await service.startStream('stream-1', 'host-1'); + expect(result).toBe(live); + expect(streamRepo.save).not.toHaveBeenCalled(); + }); + + it('refuses a non-host', async () => { + streamRepo.findOne.mockResolvedValue(baseStream()); + await expect( + service.startStream('stream-1', 'someone-else'), + ).rejects.toBeInstanceOf(ForbiddenException); + }); + + it('refuses a finished stream', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.ENDED }), + ); + await expect( + service.startStream('stream-1', 'host-1'), + ).rejects.toBeInstanceOf(BadRequestException); + }); + + it('404s on an unknown stream', async () => { + streamRepo.findOne.mockResolvedValue(null); + await expect( + service.startStream('missing', 'host-1'), + ).rejects.toBeInstanceOf(NotFoundException); + }); + }); + + describe('endStream / cancelStream', () => { + it('finalises a recording when ending a recorded stream', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE, recordingEnabled: true }), + ); + const result = await service.endStream('stream-1', 'host-1'); + expect(result.status).toBe(LiveStreamStatus.ENDED); + expect(result.endedAt).toBeInstanceOf(Date); + expect(result.recordingUrl).toBe('recordings/stream-1.m3u8'); + }); + + it('does not invent a recording URL when recording was off', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + const result = await service.endStream('stream-1', 'host-1'); + expect(result.recordingUrl).toBeUndefined(); + }); + + it('cancels a live stream', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + const result = await service.cancelStream('stream-1', 'host-1'); + expect(result.status).toBe(LiveStreamStatus.CANCELLED); + }); + }); + + describe('listStreams', () => { + it('passes the status filter to the repository', async () => { + streamRepo.find.mockResolvedValue([]); + await service.listStreams(LiveStreamStatus.LIVE); + expect(streamRepo.find).toHaveBeenCalledWith({ + where: { status: LiveStreamStatus.LIVE }, + order: { createdAt: 'DESC' }, + }); + }); + }); + + // ─── Spectators ──────────────────────────────────────────────────────────── + + describe('joinStream', () => { + it('adds an active viewer and bumps the counters', async () => { + const stream = baseStream({ + status: LiveStreamStatus.LIVE, + viewerCount: 2, + peakViewerCount: 2, + }); + streamRepo.findOne.mockResolvedValue(stream); + viewerRepo.findOne.mockResolvedValue(null); + + const viewer = await service.joinStream('stream-1', 'viewer-1', { + username: 'ada', + bandwidthKbps: 3_000, + }); + + expect(viewer.isActive).toBe(true); + expect(viewer.quality).toBe(StreamQuality.HIGH); + expect(stream.viewerCount).toBe(3); + expect(stream.peakViewerCount).toBe(3); + }); + + it('is idempotent for an existing spectator', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + const existing = { id: 'v1', viewerId: 'viewer-1', isActive: true }; + viewerRepo.findOne.mockResolvedValue(existing); + + const viewer = await service.joinStream('stream-1', 'viewer-1'); + expect(viewer).toBe(existing); + expect(streamRepo.save).not.toHaveBeenCalled(); + }); + + it('refuses a finished stream', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.ENDED }), + ); + await expect( + service.joinStream('stream-1', 'viewer-1'), + ).rejects.toBeInstanceOf(BadRequestException); + }); + + it('refuses a banned spectator', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + moderationRepo.find.mockResolvedValueOnce([ + { action: ModerationActionType.BAN, createdAt: new Date() }, + ]); + + await expect( + service.joinStream('stream-1', 'viewer-1'), + ).rejects.toBeInstanceOf(ForbiddenException); + }); + + it('refuses a spectator under an active timeout', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + moderationRepo.find + .mockResolvedValueOnce([]) // ban/unban history + .mockResolvedValueOnce([ + { + action: ModerationActionType.TIMEOUT, + expiresAt: new Date(Date.now() + 60_000), + }, + ]); + + await expect( + service.joinStream('stream-1', 'viewer-1'), + ).rejects.toBeInstanceOf(ForbiddenException); + }); + + it('allows a spectator whose timeout has expired', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + moderationRepo.find.mockResolvedValueOnce([]).mockResolvedValueOnce([ + { + action: ModerationActionType.TIMEOUT, + expiresAt: new Date(Date.now() - 60_000), + }, + ]); + viewerRepo.findOne.mockResolvedValue(null); + + const viewer = await service.joinStream('stream-1', 'viewer-1'); + expect(viewer.isActive).toBe(true); + }); + }); + + describe('leaveStream', () => { + it('marks the viewer inactive and decrements the count', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE, viewerCount: 3 }), + ); + const viewer = { + id: 'v1', + viewerId: 'viewer-1', + isActive: true, + leftAt: undefined, + }; + viewerRepo.findOne.mockResolvedValue(viewer); + + await service.leaveStream('stream-1', 'viewer-1'); + expect(viewer.isActive).toBe(false); + expect(viewer.leftAt).toBeInstanceOf(Date); + }); + + it('404s when not watching', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + viewerRepo.findOne.mockResolvedValue(null); + await expect( + service.leaveStream('stream-1', 'viewer-1'), + ).rejects.toBeInstanceOf(NotFoundException); + }); + }); + + describe('selectQuality', () => { + it('adapts the viewer quality to measured bandwidth', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + const viewer = { id: 'v1', quality: StreamQuality.MEDIUM }; + viewerRepo.findOne.mockResolvedValue(viewer); + + const result = await service.selectQuality('stream-1', 'viewer-1', 500); + expect(result.quality).toBe(StreamQuality.LOW); + }); + + it('404s for a non-viewer', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + viewerRepo.findOne.mockResolvedValue(null); + await expect( + service.selectQuality('stream-1', 'viewer-1', 9_000), + ).rejects.toBeInstanceOf(NotFoundException); + }); + }); + + // ─── Chat ────────────────────────────────────────────────────────────────── + + describe('postMessage', () => { + it('trims and stores a message on a live stream', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + const message = await service.postMessage( + 'stream-1', + 'viewer-1', + ' nice ', + ); + expect(message.content).toBe('nice'); + expect(message.isDeleted).toBe(false); + }); + + it('rejects chat when the stream is not live', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.SCHEDULED }), + ); + await expect( + service.postMessage('stream-1', 'viewer-1', 'hi'), + ).rejects.toBeInstanceOf(BadRequestException); + }); + + it('rejects an empty message', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + await expect( + service.postMessage('stream-1', 'viewer-1', ' '), + ).rejects.toBeInstanceOf(BadRequestException); + }); + + it('rejects a banned author', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + moderationRepo.find.mockResolvedValueOnce([ + { action: ModerationActionType.BAN, createdAt: new Date() }, + ]); + await expect( + service.postMessage('stream-1', 'viewer-1', 'hi'), + ).rejects.toBeInstanceOf(ForbiddenException); + }); + }); + + describe('getChatHistory', () => { + it('excludes deleted messages and applies the limit', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + chatRepo.find.mockResolvedValue([]); + await service.getChatHistory('stream-1', 25); + expect(chatRepo.find).toHaveBeenCalledWith({ + where: { streamId: 'stream-1', isDeleted: false }, + order: { createdAt: 'ASC' }, + take: 25, + }); + }); + }); + + // ─── Moderation ──────────────────────────────────────────────────────────── + + describe('moderate', () => { + beforeEach(() => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + }); + + it('allows the host to ban a spectator and deactivates them', async () => { + const viewer = { id: 'v1', isActive: true }; + viewerRepo.findOne.mockResolvedValue(viewer); + + const action = await service.moderate('stream-1', 'host-1', { + action: ModerationActionType.BAN, + targetUserId: 'viewer-1', + }); + + expect(action.action).toBe(ModerationActionType.BAN); + expect(viewer.isActive).toBe(false); + }); + + it('refuses a non-moderator', async () => { + viewerRepo.findOne.mockResolvedValue(null); + await expect( + service.moderate('stream-1', 'viewer-1', { + action: ModerationActionType.CLEAR_CHAT, + }), + ).rejects.toBeInstanceOf(ForbiddenException); + }); + + it('allows an assigned moderator', async () => { + viewerRepo.findOne.mockResolvedValue({ + id: 'v2', + isActive: true, + isModerator: true, + }); + const action = await service.moderate('stream-1', 'mod-1', { + action: ModerationActionType.CLEAR_CHAT, + }); + expect(action.action).toBe(ModerationActionType.CLEAR_CHAT); + expect(chatRepo.update).toHaveBeenCalledWith( + { streamId: 'stream-1' }, + { isDeleted: true }, + ); + }); + + it('requires a target for a ban', async () => { + await expect( + service.moderate('stream-1', 'host-1', { + action: ModerationActionType.BAN, + }), + ).rejects.toBeInstanceOf(BadRequestException); + }); + + it('sets an expiry on a timeout', async () => { + viewerRepo.findOne.mockResolvedValue(null); + const before = Date.now(); + const action = await service.moderate('stream-1', 'host-1', { + action: ModerationActionType.TIMEOUT, + targetUserId: 'viewer-1', + durationMinutes: 5, + }); + expect(action.expiresAt).toBeInstanceOf(Date); + if (!action.expiresAt) { + throw new Error('expected a timeout expiry'); + } + expect(action.expiresAt.getTime()).toBeGreaterThan(before); + expect(action.expiresAt.getTime()).toBeLessThanOrEqual( + before + 5 * 60_000 + 1_000, + ); + }); + + it('requires a message id to delete a message', async () => { + await expect( + service.moderate('stream-1', 'host-1', { + action: ModerationActionType.DELETE_MESSAGE, + }), + ).rejects.toBeInstanceOf(BadRequestException); + }); + + it('soft-deletes the targeted message', async () => { + await service.moderate('stream-1', 'host-1', { + action: ModerationActionType.DELETE_MESSAGE, + targetUserId: 'viewer-1', + messageId: 'msg-1', + }); + expect(chatRepo.update).toHaveBeenCalledWith( + { id: 'msg-1', streamId: 'stream-1' }, + { isDeleted: true, deletedBy: 'host-1' }, + ); + }); + }); + + // ─── Recording ───────────────────────────────────────────────────────────── + + describe('recording', () => { + it('toggles recording for the host', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + const started = await service.startRecording('stream-1', 'host-1'); + expect(started.recordingEnabled).toBe(true); + + const stopped = await service.stopRecording('stream-1', 'host-1'); + expect(stopped.recordingEnabled).toBe(false); + expect(stopped.recordingUrl).toBe('recordings/stream-1.m3u8'); + }); + + it('refuses recording control by a non-host', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ status: LiveStreamStatus.LIVE }), + ); + await expect( + service.startRecording('stream-1', 'viewer-1'), + ).rejects.toBeInstanceOf(ForbiddenException); + }); + }); + + // ─── Analytics ───────────────────────────────────────────────────────────── + + describe('getAnalytics', () => { + it('aggregates viewership, chat and quality distribution', async () => { + streamRepo.findOne.mockResolvedValue( + baseStream({ + status: LiveStreamStatus.LIVE, + peakViewerCount: 12, + startedAt: new Date(Date.now() - 60_000), + recordingEnabled: true, + recordingUrl: 'recordings/stream-1.m3u8', + }), + ); + viewerRepo.find.mockResolvedValue([ + { quality: StreamQuality.HIGH }, + { quality: StreamQuality.HIGH }, + { quality: StreamQuality.LOW }, + ]); + viewerRepo.count.mockResolvedValue(30); + chatRepo.count.mockResolvedValue(44); + moderationRepo.count.mockResolvedValue(2); + + const analytics = await service.getAnalytics('stream-1'); + + expect(analytics.currentViewers).toBe(3); + expect(analytics.peakViewers).toBe(12); + expect(analytics.totalJoins).toBe(30); + expect(analytics.chatMessages).toBe(44); + expect(analytics.moderationActions).toBe(2); + expect(analytics.qualityDistribution).toEqual({ + [StreamQuality.LOW]: 1, + [StreamQuality.MEDIUM]: 0, + [StreamQuality.HIGH]: 2, + [StreamQuality.SOURCE]: 0, + }); + expect(analytics.recording).toEqual({ + enabled: true, + url: 'recordings/stream-1.m3u8', + }); + expect(analytics.durationSeconds).toBeGreaterThanOrEqual(59); + }); + }); +}); diff --git a/src/live-streaming/live-streaming.service.ts b/src/live-streaming/live-streaming.service.ts new file mode 100644 index 00000000..6f41d4ed --- /dev/null +++ b/src/live-streaming/live-streaming.service.ts @@ -0,0 +1,509 @@ +import { + BadRequestException, + ForbiddenException, + Injectable, + Logger, + NotFoundException, +} from '@nestjs/common'; +import { InjectRepository } from '@nestjs/typeorm'; +import { In, Repository } from 'typeorm'; + +import { + LiveStream, + LiveStreamStatus, + StreamQuality, +} from './entities/live-stream.entity'; +import { StreamViewer } from './entities/stream-viewer.entity'; +import { StreamChatMessage } from './entities/stream-chat-message.entity'; +import { + ModerationActionType, + StreamModerationAction, +} from './entities/stream-moderation-action.entity'; +import { CreateStreamDto } from './dto/create-stream.dto'; +import { JoinStreamDto } from './dto/join-stream.dto'; +import { ModerateStreamDto } from './dto/moderate-stream.dto'; + +export interface StreamAnalytics { + streamId: string; + status: LiveStreamStatus; + currentViewers: number; + peakViewers: number; + totalJoins: number; + chatMessages: number; + moderationActions: number; + qualityDistribution: Record; + recording: { enabled: boolean; url: string | null }; + durationSeconds: number; +} + +/** + * Deterministic bandwidth -> quality mapping. + * + * Deliberately a pure function so the ladder is testable without a viewer row + * and the same input always picks the same rung. Unknown bandwidth falls back to + * `MEDIUM`, the quality a viewer joins with before the first measurement. + */ +export function qualityForBandwidth(bandwidthKbps?: number): StreamQuality { + if (bandwidthKbps === undefined || bandwidthKbps === null) { + return StreamQuality.MEDIUM; + } + if (bandwidthKbps >= 6_000) return StreamQuality.SOURCE; + if (bandwidthKbps >= 2_500) return StreamQuality.HIGH; + if (bandwidthKbps >= 1_000) return StreamQuality.MEDIUM; + return StreamQuality.LOW; +} + +@Injectable() +export class LiveStreamingService { + private readonly logger = new Logger(LiveStreamingService.name); + + constructor( + @InjectRepository(LiveStream) + private readonly streamRepo: Repository, + @InjectRepository(StreamViewer) + private readonly viewerRepo: Repository, + @InjectRepository(StreamChatMessage) + private readonly chatRepo: Repository, + @InjectRepository(StreamModerationAction) + private readonly moderationRepo: Repository, + ) {} + + // ─── Stream lifecycle ────────────────────────────────────────────────────── + + async createStream( + hostId: string, + dto: CreateStreamDto, + ): Promise { + const stream = this.streamRepo.create({ + hostId, + title: dto.title, + description: dto.description, + recordingEnabled: dto.recordingEnabled ?? false, + status: LiveStreamStatus.SCHEDULED, + currentQuality: StreamQuality.MEDIUM, + viewerCount: 0, + peakViewerCount: 0, + }); + return this.streamRepo.save(stream); + } + + async startStream(streamId: string, hostId: string): Promise { + const stream = await this.requireStream(streamId); + this.requireHost(stream, hostId); + + if (stream.status === LiveStreamStatus.LIVE) { + return stream; + } + if ( + stream.status === LiveStreamStatus.ENDED || + stream.status === LiveStreamStatus.CANCELLED + ) { + throw new BadRequestException('Stream has already finished'); + } + + stream.status = LiveStreamStatus.LIVE; + stream.startedAt = stream.startedAt ?? new Date(); + return this.streamRepo.save(stream); + } + + async endStream(streamId: string, hostId: string): Promise { + const stream = await this.requireStream(streamId); + this.requireHost(stream, hostId); + + if (stream.status === LiveStreamStatus.ENDED) { + return stream; + } + if (stream.status === LiveStreamStatus.CANCELLED) { + throw new BadRequestException('Stream was cancelled'); + } + + stream.status = LiveStreamStatus.ENDED; + stream.endedAt = new Date(); + if (stream.recordingEnabled && !stream.recordingUrl) { + stream.recordingUrl = this.recordingUrl(stream.id); + } + return this.streamRepo.save(stream); + } + + async cancelStream(streamId: string, hostId: string): Promise { + const stream = await this.requireStream(streamId); + this.requireHost(stream, hostId); + + if ( + stream.status === LiveStreamStatus.ENDED || + stream.status === LiveStreamStatus.CANCELLED + ) { + throw new BadRequestException('Stream has already finished'); + } + + stream.status = LiveStreamStatus.CANCELLED; + stream.endedAt = new Date(); + return this.streamRepo.save(stream); + } + + async listStreams(status?: LiveStreamStatus): Promise { + return this.streamRepo.find({ + where: status ? { status } : {}, + order: { createdAt: 'DESC' }, + }); + } + + async getStream(streamId: string): Promise { + return this.requireStream(streamId); + } + + // ─── Spectators ──────────────────────────────────────────────────────────── + + async joinStream( + streamId: string, + viewerId: string, + dto: JoinStreamDto = {}, + ): Promise { + const stream = await this.requireStream(streamId); + if ( + stream.status !== LiveStreamStatus.LIVE && + stream.status !== LiveStreamStatus.SCHEDULED + ) { + throw new BadRequestException('Stream is not joinable'); + } + await this.requireNotBlocked(streamId, viewerId); + + const existing = await this.findActiveViewer(streamId, viewerId); + if (existing) { + return existing; + } + + const quality = + dto.bandwidthKbps !== undefined + ? qualityForBandwidth(dto.bandwidthKbps) + : dto.preferredQuality ?? StreamQuality.MEDIUM; + + const viewer = this.viewerRepo.create({ + streamId, + viewerId, + username: dto.username, + quality, + isActive: true, + isModerator: false, + joinedAt: new Date(), + }); + const saved = await this.viewerRepo.save(viewer); + + stream.viewerCount = (stream.viewerCount ?? 0) + 1; + stream.peakViewerCount = Math.max( + stream.peakViewerCount ?? 0, + stream.viewerCount, + ); + await this.streamRepo.save(stream); + + this.logger.log(`Viewer ${viewerId} joined stream ${streamId}`); + return saved; + } + + async leaveStream(streamId: string, viewerId: string): Promise { + await this.requireStream(streamId); + const viewer = await this.findActiveViewer(streamId, viewerId); + if (!viewer) { + throw new NotFoundException('Not currently watching this stream'); + } + + viewer.isActive = false; + viewer.leftAt = new Date(); + await this.viewerRepo.save(viewer); + + const stream = await this.requireStream(streamId); + stream.viewerCount = Math.max(0, (stream.viewerCount ?? 1) - 1); + await this.streamRepo.save(stream); + } + + async getViewerCount(streamId: string): Promise { + return this.viewerRepo.count({ where: { streamId, isActive: true } }); + } + + async selectQuality( + streamId: string, + viewerId: string, + bandwidthKbps: number, + ): Promise { + await this.requireStream(streamId); + const viewer = await this.findActiveViewer(streamId, viewerId); + if (!viewer) { + throw new NotFoundException('Not currently watching this stream'); + } + + viewer.quality = qualityForBandwidth(bandwidthKbps); + return this.viewerRepo.save(viewer); + } + + // ─── Chat ────────────────────────────────────────────────────────────────── + + async postMessage( + streamId: string, + authorId: string, + content: string, + ): Promise { + const stream = await this.requireStream(streamId); + if (stream.status !== LiveStreamStatus.LIVE) { + throw new BadRequestException('Chat is only open on a live stream'); + } + await this.requireNotBlocked(streamId, authorId); + + const body = content?.trim(); + if (!body) { + throw new BadRequestException('Message cannot be empty'); + } + if (body.length > 500) { + throw new BadRequestException('Message exceeds 500 characters'); + } + + const message = this.chatRepo.create({ + streamId, + authorId, + content: body, + isDeleted: false, + isPinned: false, + }); + return this.chatRepo.save(message); + } + + async getChatHistory( + streamId: string, + limit = 50, + ): Promise { + await this.requireStream(streamId); + return this.chatRepo.find({ + where: { streamId, isDeleted: false }, + order: { createdAt: 'ASC' }, + take: limit, + }); + } + + // ─── Moderation ──────────────────────────────────────────────────────────── + + async moderate( + streamId: string, + moderatorId: string, + dto: ModerateStreamDto, + ): Promise { + const stream = await this.requireStream(streamId); + await this.requireModerator(stream, moderatorId); + + switch (dto.action) { + case ModerationActionType.BAN: + case ModerationActionType.UNBAN: + case ModerationActionType.TIMEOUT: + if (!dto.targetUserId) { + throw new BadRequestException( + 'targetUserId is required for this action', + ); + } + break; + case ModerationActionType.DELETE_MESSAGE: + case ModerationActionType.PIN_MESSAGE: + if (!dto.messageId) { + throw new BadRequestException( + 'messageId is required for this action', + ); + } + break; + default: + break; + } + + const expiresAt = + dto.action === ModerationActionType.TIMEOUT + ? new Date(Date.now() + (dto.durationMinutes ?? 10) * 60_000) + : undefined; + + const action = this.moderationRepo.create({ + streamId, + moderatorId, + targetUserId: dto.targetUserId, + action: dto.action, + reason: dto.reason, + messageId: dto.messageId, + expiresAt, + }); + const saved = await this.moderationRepo.save(action); + + if (dto.action === ModerationActionType.BAN && dto.targetUserId) { + await this.deactivateViewer(streamId, dto.targetUserId); + } + if (dto.action === ModerationActionType.DELETE_MESSAGE && dto.messageId) { + await this.chatRepo.update( + { id: dto.messageId, streamId }, + { isDeleted: true, deletedBy: moderatorId }, + ); + } + if (dto.action === ModerationActionType.CLEAR_CHAT) { + await this.chatRepo.update({ streamId }, { isDeleted: true }); + } + if (dto.action === ModerationActionType.PIN_MESSAGE && dto.messageId) { + await this.chatRepo.update( + { id: dto.messageId, streamId }, + { isPinned: true }, + ); + } + + return saved; + } + + async getModerationLog(streamId: string): Promise { + await this.requireStream(streamId); + return this.moderationRepo.find({ + where: { streamId }, + order: { createdAt: 'DESC' }, + }); + } + + // ─── Recording ───────────────────────────────────────────────────────────── + + async startRecording(streamId: string, hostId: string): Promise { + const stream = await this.requireStream(streamId); + this.requireHost(stream, hostId); + stream.recordingEnabled = true; + return this.streamRepo.save(stream); + } + + async stopRecording(streamId: string, hostId: string): Promise { + const stream = await this.requireStream(streamId); + this.requireHost(stream, hostId); + stream.recordingEnabled = false; + stream.recordingUrl = stream.recordingUrl ?? this.recordingUrl(stream.id); + return this.streamRepo.save(stream); + } + + // ─── Analytics ───────────────────────────────────────────────────────────── + + async getAnalytics(streamId: string): Promise { + const stream = await this.requireStream(streamId); + const active = await this.viewerRepo.find({ + where: { streamId, isActive: true }, + }); + + const distribution: Record = { + [StreamQuality.LOW]: 0, + [StreamQuality.MEDIUM]: 0, + [StreamQuality.HIGH]: 0, + [StreamQuality.SOURCE]: 0, + }; + for (const viewer of active) { + distribution[viewer.quality] += 1; + } + + const [totalJoins, chatMessages, moderationActions] = await Promise.all([ + this.viewerRepo.count({ where: { streamId } }), + this.chatRepo.count({ where: { streamId, isDeleted: false } }), + this.moderationRepo.count({ where: { streamId } }), + ]); + + const endedAt = stream.endedAt ?? new Date(); + const durationSeconds = stream.startedAt + ? Math.max( + 0, + Math.floor((endedAt.getTime() - stream.startedAt.getTime()) / 1000), + ) + : 0; + + return { + streamId, + status: stream.status, + currentViewers: active.length, + peakViewers: stream.peakViewerCount ?? 0, + totalJoins, + chatMessages, + moderationActions, + qualityDistribution: distribution, + recording: { + enabled: stream.recordingEnabled, + url: stream.recordingUrl ?? null, + }, + durationSeconds, + }; + } + + // ─── Internals ───────────────────────────────────────────────────────────── + + private async requireStream(streamId: string): Promise { + const stream = await this.streamRepo.findOne({ where: { id: streamId } }); + if (!stream) { + throw new NotFoundException('Stream not found'); + } + return stream; + } + + private requireHost(stream: LiveStream, userId: string): void { + if (stream.hostId !== userId) { + throw new ForbiddenException('Only the stream host can do that'); + } + } + + private async requireModerator( + stream: LiveStream, + userId: string, + ): Promise { + if (stream.hostId === userId) { + return; + } + const viewer = await this.findActiveViewer(stream.id, userId); + if (!viewer?.isModerator) { + throw new ForbiddenException('Moderator privileges required'); + } + } + + private async findActiveViewer( + streamId: string, + viewerId: string, + ): Promise { + return this.viewerRepo.findOne({ + where: { streamId, viewerId, isActive: true }, + }); + } + + private async requireNotBlocked( + streamId: string, + viewerId: string, + ): Promise { + const actions = await this.moderationRepo.find({ + where: { + streamId, + targetUserId: viewerId, + action: In([ModerationActionType.BAN, ModerationActionType.UNBAN]), + }, + order: { createdAt: 'DESC' }, + }); + if (actions[0]?.action === ModerationActionType.BAN) { + throw new ForbiddenException('You are banned from this stream'); + } + + const timeouts = await this.moderationRepo.find({ + where: { + streamId, + targetUserId: viewerId, + action: ModerationActionType.TIMEOUT, + }, + order: { createdAt: 'DESC' }, + }); + const latest = timeouts[0]; + if (latest?.expiresAt && latest.expiresAt.getTime() > Date.now()) { + throw new ForbiddenException('You are timed out from this stream'); + } + } + + private async deactivateViewer( + streamId: string, + viewerId: string, + ): Promise { + const viewer = await this.findActiveViewer(streamId, viewerId); + if (!viewer) { + return; + } + viewer.isActive = false; + viewer.leftAt = new Date(); + await this.viewerRepo.save(viewer); + } + + private recordingUrl(streamId: string): string { + return `recordings/${streamId}.m3u8`; + } +} diff --git a/src/migrations/1760000000000-CreateLiveStreamingTables.ts b/src/migrations/1760000000000-CreateLiveStreamingTables.ts new file mode 100644 index 00000000..f848d5c0 --- /dev/null +++ b/src/migrations/1760000000000-CreateLiveStreamingTables.ts @@ -0,0 +1,237 @@ +import { MigrationInterface, QueryRunner, Table, TableIndex } from 'typeorm'; + +export class CreateLiveStreamingTables1760000000000 + implements MigrationInterface +{ + name = 'CreateLiveStreamingTables1760000000000'; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.createTable( + new Table({ + name: 'live_streams', + columns: [ + { + name: 'id', + type: 'uuid', + isPrimary: true, + generationStrategy: 'uuid', + default: 'uuid_generate_v4()', + }, + { name: 'hostId', type: 'varchar', length: '64', isNullable: false }, + { name: 'title', type: 'varchar', length: '200', isNullable: false }, + { name: 'description', type: 'text', isNullable: true }, + { + name: 'status', + type: 'varchar', + length: '20', + default: "'scheduled'", + }, + { + name: 'currentQuality', + type: 'varchar', + length: '10', + default: "'medium'", + }, + { name: 'viewerCount', type: 'int', default: 0 }, + { name: 'peakViewerCount', type: 'int', default: 0 }, + { name: 'recordingEnabled', type: 'boolean', default: false }, + { + name: 'recordingUrl', + type: 'varchar', + length: '500', + isNullable: true, + }, + { name: 'startedAt', type: 'timestamp', isNullable: true }, + { name: 'endedAt', type: 'timestamp', isNullable: true }, + { + name: 'createdAt', + type: 'timestamp', + default: 'CURRENT_TIMESTAMP', + }, + { + name: 'updatedAt', + type: 'timestamp', + default: 'CURRENT_TIMESTAMP', + onUpdate: 'CURRENT_TIMESTAMP', + }, + ], + }), + true, + ); + + await queryRunner.createTable( + new Table({ + name: 'stream_viewers', + columns: [ + { + name: 'id', + type: 'uuid', + isPrimary: true, + generationStrategy: 'uuid', + default: 'uuid_generate_v4()', + }, + { name: 'streamId', type: 'uuid', isNullable: false }, + { + name: 'viewerId', + type: 'varchar', + length: '64', + isNullable: false, + }, + { + name: 'username', + type: 'varchar', + length: '64', + isNullable: true, + }, + { + name: 'quality', + type: 'varchar', + length: '10', + default: "'medium'", + }, + { name: 'isModerator', type: 'boolean', default: false }, + { name: 'isActive', type: 'boolean', default: true }, + { name: 'joinedAt', type: 'timestamp', isNullable: true }, + { name: 'leftAt', type: 'timestamp', isNullable: true }, + { + name: 'createdAt', + type: 'timestamp', + default: 'CURRENT_TIMESTAMP', + }, + { + name: 'updatedAt', + type: 'timestamp', + default: 'CURRENT_TIMESTAMP', + onUpdate: 'CURRENT_TIMESTAMP', + }, + ], + }), + true, + ); + + await queryRunner.createTable( + new Table({ + name: 'stream_chat_messages', + columns: [ + { + name: 'id', + type: 'uuid', + isPrimary: true, + generationStrategy: 'uuid', + default: 'uuid_generate_v4()', + }, + { name: 'streamId', type: 'uuid', isNullable: false }, + { + name: 'authorId', + type: 'varchar', + length: '64', + isNullable: false, + }, + { name: 'content', type: 'text', isNullable: false }, + { name: 'isDeleted', type: 'boolean', default: false }, + { + name: 'deletedBy', + type: 'varchar', + length: '64', + isNullable: true, + }, + { name: 'isPinned', type: 'boolean', default: false }, + { + name: 'createdAt', + type: 'timestamp', + default: 'CURRENT_TIMESTAMP', + }, + ], + }), + true, + ); + + await queryRunner.createTable( + new Table({ + name: 'stream_moderation_actions', + columns: [ + { + name: 'id', + type: 'uuid', + isPrimary: true, + generationStrategy: 'uuid', + default: 'uuid_generate_v4()', + }, + { name: 'streamId', type: 'uuid', isNullable: false }, + { + name: 'moderatorId', + type: 'varchar', + length: '64', + isNullable: false, + }, + { + name: 'targetUserId', + type: 'varchar', + length: '64', + isNullable: true, + }, + { name: 'action', type: 'varchar', length: '20', isNullable: false }, + { name: 'reason', type: 'varchar', length: '200', isNullable: true }, + { name: 'messageId', type: 'uuid', isNullable: true }, + { name: 'expiresAt', type: 'timestamp', isNullable: true }, + { + name: 'createdAt', + type: 'timestamp', + default: 'CURRENT_TIMESTAMP', + }, + ], + }), + true, + ); + + await queryRunner.createIndex( + 'live_streams', + new TableIndex({ + name: 'IDX_live_streams_hostId', + columnNames: ['hostId'], + }), + ); + await queryRunner.createIndex( + 'live_streams', + new TableIndex({ + name: 'IDX_live_streams_status', + columnNames: ['status'], + }), + ); + await queryRunner.createIndex( + 'stream_viewers', + new TableIndex({ + name: 'IDX_stream_viewers_streamId_viewerId', + columnNames: ['streamId', 'viewerId'], + }), + ); + await queryRunner.createIndex( + 'stream_viewers', + new TableIndex({ + name: 'IDX_stream_viewers_streamId_isActive', + columnNames: ['streamId', 'isActive'], + }), + ); + await queryRunner.createIndex( + 'stream_chat_messages', + new TableIndex({ + name: 'IDX_stream_chat_messages_streamId_createdAt', + columnNames: ['streamId', 'createdAt'], + }), + ); + await queryRunner.createIndex( + 'stream_moderation_actions', + new TableIndex({ + name: 'IDX_stream_moderation_actions_streamId_targetUserId', + columnNames: ['streamId', 'targetUserId'], + }), + ); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.dropTable('stream_moderation_actions'); + await queryRunner.dropTable('stream_chat_messages'); + await queryRunner.dropTable('stream_viewers'); + await queryRunner.dropTable('live_streams'); + } +}