Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/app.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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: [
Expand Down Expand Up @@ -53,6 +54,7 @@ import { AdminModule } from './admin/admin.module';
JobsModule,
CdnModule,
AdminModule,
LiveStreamingModule,
],
})
export class AppModule {}
54 changes: 54 additions & 0 deletions src/live-streaming/README.md
Original file line number Diff line number Diff line change
@@ -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:<id>` 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.
10 changes: 10 additions & 0 deletions src/live-streaming/dto/chat-message.dto.ts
Original file line number Diff line number Diff line change
@@ -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;
}
30 changes: 30 additions & 0 deletions src/live-streaming/dto/create-stream.dto.ts
Original file line number Diff line number Diff line change
@@ -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;
}
4 changes: 4 additions & 0 deletions src/live-streaming/dto/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
export * from './create-stream.dto';
export * from './join-stream.dto';
export * from './chat-message.dto';
export * from './moderate-stream.dto';
35 changes: 35 additions & 0 deletions src/live-streaming/dto/join-stream.dto.ts
Original file line number Diff line number Diff line change
@@ -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;
}
48 changes: 48 additions & 0 deletions src/live-streaming/dto/moderate-stream.dto.ts
Original file line number Diff line number Diff line change
@@ -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;
}
4 changes: 4 additions & 0 deletions src/live-streaming/entities/index.ts
Original file line number Diff line number Diff line change
@@ -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';
86 changes: 86 additions & 0 deletions src/live-streaming/entities/live-stream.entity.ts
Original file line number Diff line number Diff line change
@@ -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;
}
40 changes: 40 additions & 0 deletions src/live-streaming/entities/stream-chat-message.entity.ts
Original file line number Diff line number Diff line change
@@ -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;
}
53 changes: 53 additions & 0 deletions src/live-streaming/entities/stream-moderation-action.entity.ts
Original file line number Diff line number Diff line change
@@ -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;
}
Loading