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
10 changes: 9 additions & 1 deletion src/app.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,8 +29,12 @@ import { AlertsModule } from "./alerts/alerts.module";
import { MetricsModule } from "./metrics/metrics.module";
import { AnalyticsModule } from "./analytics/analytics.module";
import { RateLimitModule } from "./quota/rate-limit.module";
import { MessagingModule } from "./messaging/messaging.module";

// Auth entities
import { Conversation } from "./messaging/entities/conversation.entity";
import { Message } from "./messaging/entities/message.entity";
import { UserPresence } from "./messaging/entities/user-presence.entity";
import { User } from "./user/entities/user.entity";
import { EmailVerification } from "./auth/entities/email-verification.entity";
import { Wallet } from "./auth/entities/wallet.entity";
Expand Down Expand Up @@ -118,6 +122,9 @@ import { QuotaGuard } from "./common/guard/quota.guard";
AlertDeliveryLog,
AnalyticsEvent,
DailyMetric,
Conversation,
Message,
UserPresence,
],
synchronize: !isProduction,
logging: isProduction ? ["error"] : ["error", "warn", "schema"],
Expand Down Expand Up @@ -154,6 +161,7 @@ import { QuotaGuard } from "./common/guard/quota.guard";
MetricsModule,
AnalyticsModule,
RateLimitModule,
MessagingModule,
],

controllers: [AppController],
Expand Down Expand Up @@ -203,4 +211,4 @@ export class AppModule implements NestModule, OnModuleInit {
onModuleInit() {
this.verifier.start();
}
}
}
82 changes: 82 additions & 0 deletions src/messaging/MESSAGING_SCALING.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
# Messaging Module: Scaling Strategy & Infrastructure

## Overview
The StellAIverse messaging module uses Socket.IO for real-time WebSocket communication, combined with PostgreSQL for persistence and Redis for horizontal scaling. This document outlines the infrastructure requirements and scaling strategy to support a growing user base.

## Required Infrastructure

### 1. Redis (Required for Horizontal Scaling)
Socket.IO requires a Redis adapter to enable broadcasting across multiple backend instances. Redis acts as a pub/sub layer for message distribution.

#### Redis Setup:
```typescript
// In messaging.module.ts add the Redis adapter
import { RedisAdapter } from "@socket.io/redis-adapter";
import { createClient } from "redis";

// Then in the MessagingGateway after server initialization:
const pubClient = createClient({ url: process.env.REDIS_URL });
const subClient = pubClient.duplicate();
await Promise.all([pubClient.connect(), subClient.connect()]);
this.server.adapter(createAdapter(pubClient, subClient));
```

### 2. PostgreSQL (Current Persistence Layer)
We currently use PostgreSQL for storing all conversations, messages, and presence data. This is sufficient for up to millions of messages, but for larger scale we can implement the following optimizations:

### 3. Additional Recommended Infrastructure
- **S3/Object Storage**: For archiving old messages (older than 1 year)
- **Full-text Search Service**: For implementing message search functionality (Elasticsearch or Typesense)

## Message Lifecycle & Cleanup Strategy

### Persistence Rules:
1. **Active Messages**: Last 12 months stored in PostgreSQL
2. **Archived Messages**: Older than 12 months moved to S3/Glacier
3. **Cron Job**: Run `messagingService.archiveOldMessages()` monthly

### Implementation:
```typescript
// Add a cron service to run monthly archiving
import { Cron } from "@nestjs/schedule";

@Cron("0 0 1 * *") // Run on the 1st of every month
async handleArchiving() {
const archivedCount = await this.messagingService.archiveOldMessages();
logger.log(`Archived ${archivedCount} old messages`);
}
```

## WebSocket Connection Flow
1. Client connects to `wss://api.example.com/messaging`
2. Client sends JWT token in handshake
3. Server validates token, attaches user to socket
4. Server adds user to their personal room: `user:{userId}`
5. When joining a conversation, server adds socket to `conversation:{conversationId}`

## Reconnection & Message Guarantee
To handle dropped connections and ensure no message loss:
1. **Client-side buffering**: Messages are queued when disconnected
2. **Message IDs**: Every message gets a UUID client-side to prevent duplicates
3. **Last seen synchronization**: On reconnection, client fetches all messages since last seen
4. **At-least-once delivery**: Server persists messages before broadcasting

## Scaling to Multiple Instances
1. Deploy behind a load balancer that supports sticky sessions (or use Redis adapter which removes this requirement)
2. Use the Redis adapter to enable cross-instance communication
3. Implement connection draining during deployments to avoid abrupt disconnections
4. Monitor connection counts per instance, scale out when approaching 10k connections per instance

## Monitoring & Metrics
Key metrics to track:
- Active connections per instance
- Messages sent/sec
- Delivery receipts latency
- Archive job success rate
- WebSocket error rates

## Testing Strategy
1. **Unit tests**: Test all service methods
2. **Integration tests**: Test WebSocket connection, message flow
3. **Load tests**: Simulate thousands of concurrent users
4. **Chaos tests**: Simulate network splits, instance failures
19 changes: 19 additions & 0 deletions src/messaging/dto/create-conversation.dto.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
import { IsArray, IsEnum, IsOptional, IsString, IsUUID } from "class-validator";
import { ConversationType } from "../entities/message.enum";

export class CreateConversationDto {
@IsArray()
@IsUUID("4", { each: true })
participantIds: string[];

@IsEnum(ConversationType)
type: ConversationType;

@IsOptional()
@IsString()
name?: string;

@IsOptional()
@IsString()
avatar?: string;
}
14 changes: 14 additions & 0 deletions src/messaging/dto/get-messages.dto.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
import { IsOptional, IsUUID, IsInt, Min } from "class-validator";
import { Type } from "class-transformer";

export class GetMessagesDto {
@IsOptional()
@IsUUID()
before?: string;

@IsOptional()
@IsInt()
@Min(1)
@Type(() => Number)
limit?: number = 50;
}
13 changes: 13 additions & 0 deletions src/messaging/dto/send-message.dto.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
import { IsString, IsUUID, IsOptional, IsObject } from "class-validator";

export class SendMessageDto {
@IsUUID()
conversationId: string;

@IsString()
content: string;

@IsOptional()
@IsObject()
metadata?: Record<string, any>;
}
48 changes: 48 additions & 0 deletions src/messaging/entities/conversation.entity.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
import {
Entity,
PrimaryGeneratedColumn,
Column,
ManyToMany,
JoinTable,
OneToMany,
CreateDateColumn,
UpdateDateColumn,
} from "typeorm";
import { User } from "../../user/entities/user.entity";
import { Message } from "./message.entity";
import { ConversationType } from "./message.enum";

@Entity()
export class Conversation {
@PrimaryGeneratedColumn("uuid")
id: string;

@Column({
type: "enum",
enum: ConversationType,
default: ConversationType.PRIVATE,
})
type: ConversationType;

@Column({ nullable: true })
name?: string;

@Column({ nullable: true })
avatar?: string;

@ManyToMany(() => User)
@JoinTable()
participants: User[];

@OneToMany(() => Message, (message) => message.conversation)
messages: Message[];

@CreateDateColumn()
createdAt: Date;

@UpdateDateColumn()
updatedAt: Date;

@Column({ nullable: true })
lastMessageAt?: Date;
}
48 changes: 48 additions & 0 deletions src/messaging/entities/message.entity.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
import {
Entity,
PrimaryGeneratedColumn,
Column,
ManyToOne,
CreateDateColumn,
UpdateDateColumn,
} from "typeorm";
import { User } from "../../user/entities/user.entity";
import { Conversation } from "./conversation.entity";
import { MessageStatus } from "./message.enum";

@Entity()
export class Message {
@PrimaryGeneratedColumn("uuid")
id: string;

@Column("text")
content: string;

@ManyToOne(() => User)
sender: User;

@ManyToOne(() => Conversation, (conversation) => conversation.messages)
conversation: Conversation;

@Column({
type: "enum",
enum: MessageStatus,
default: MessageStatus.SENT,
})
status: MessageStatus;

@Column({ type: "jsonb", nullable: true })
metadata?: Record<string, any>;

@CreateDateColumn()
createdAt: Date;

@UpdateDateColumn()
updatedAt: Date;

@Column({ nullable: true })
deliveredAt?: Date;

@Column({ nullable: true })
readAt?: Date;
}
10 changes: 10 additions & 0 deletions src/messaging/entities/message.enum.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
export enum MessageStatus {
SENT = "sent",
DELIVERED = "delivered",
READ = "read",
}

export enum ConversationType {
PRIVATE = "private",
GROUP = "group",
}
35 changes: 35 additions & 0 deletions src/messaging/entities/user-presence.entity.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
import {
Entity,
PrimaryGeneratedColumn,
Column,
OneToOne,
JoinColumn,
CreateDateColumn,
UpdateDateColumn,
} from "typeorm";
import { User } from "../../user/entities/user.entity";

@Entity()
export class UserPresence {
@PrimaryGeneratedColumn("uuid")
id: string;

@OneToOne(() => User)
@JoinColumn()
user: User;

@Column({ default: false })
isOnline: boolean;

@Column({ nullable: true })
lastSeenAt?: Date;

@Column({ nullable: true })
currentSocketId?: string;

@CreateDateColumn()
createdAt: Date;

@UpdateDateColumn()
updatedAt: Date;
}
31 changes: 31 additions & 0 deletions src/messaging/events/message.events.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
export class MessageSentEvent {
constructor(
public readonly message: any,
public readonly conversationId: string,
public readonly recipientIds: string[],
) {}
}

export class MessageDeliveredEvent {
constructor(
public readonly messageId: string,
public readonly conversationId: string,
public readonly userId: string,
) {}
}

export class MessageReadEvent {
constructor(
public readonly messageId: string,
public readonly conversationId: string,
public readonly userId: string,
) {}
}

export class UserPresenceChangedEvent {
constructor(
public readonly userId: string,
public readonly isOnline: boolean,
public readonly lastSeenAt?: Date,
) {}
}
Loading
Loading