add support for websocket api for real time updates.

This commit is contained in:
Fizer Khan
2025-12-05 16:22:59 +05:30
parent 5e5a20c9a7
commit 66569fe2ca
24 changed files with 1713 additions and 175 deletions
+5
View File
@@ -9,6 +9,7 @@ import { connectDatabase, disconnectDatabase } from './lib/prisma';
import { redis } from './lib/redis';
import { connectInternalNats, disconnectInternalNats, disconnectAllClusters } from './lib/nats';
import { closeClickHouseClient } from './lib/clickhouse';
import { setupWebSocket } from './lib/websocket';
// Import routes
import { authRoutes } from './modules/auth/auth.routes';
@@ -137,6 +138,10 @@ async function start() {
// Connect to internal NATS
await connectInternalNats();
// Setup WebSocket server
setupWebSocket(app);
app.log.info('WebSocket server ready at /ws');
// Start HTTP server
await app.listen({
port: config.PORT,
+194 -28
View File
@@ -92,11 +92,11 @@ export async function queryStreamMetrics(
format: 'JSONEachRow',
});
const rows = await result.json<any[]>();
const rows = await result.json() as Record<string, unknown>[];
return rows.map((row) => ({
clusterId: row.cluster_id,
streamName: row.stream_name,
timestamp: new Date(row.timestamp),
clusterId: row.cluster_id as string,
streamName: row.stream_name as string,
timestamp: new Date(row.timestamp as string),
messagesTotal: Number(row.messages_total),
bytesTotal: Number(row.bytes_total),
messagesRate: Number(row.messages_rate),
@@ -104,7 +104,7 @@ export async function queryStreamMetrics(
consumerCount: Number(row.consumer_count),
firstSeq: Number(row.first_seq),
lastSeq: Number(row.last_seq),
subjects: row.subjects || [],
subjects: (row.subjects as string[]) || [],
}));
}
@@ -175,12 +175,12 @@ export async function queryConsumerMetrics(
format: 'JSONEachRow',
});
const rows = await result.json<any[]>();
const rows = await result.json() as Record<string, unknown>[];
return rows.map((row) => ({
clusterId: row.cluster_id,
streamName: row.stream_name,
consumerName: row.consumer_name,
timestamp: new Date(row.timestamp),
clusterId: row.cluster_id as string,
streamName: row.stream_name as string,
consumerName: row.consumer_name as string,
timestamp: new Date(row.timestamp as string),
pendingCount: Number(row.pending_count),
ackPending: Number(row.ack_pending),
redelivered: Number(row.redelivered),
@@ -293,7 +293,7 @@ export async function queryAuditLogs(
query_params: params,
format: 'JSONEachRow',
});
const countRows = await countResult.json<{ total: string }[]>();
const countRows = await countResult.json() as { total: string }[];
const total = parseInt(countRows[0]?.total || '0');
// Get logs
@@ -309,29 +309,195 @@ export async function queryAuditLogs(
format: 'JSONEachRow',
});
const rows = await result.json<any[]>();
const rows = await result.json() as Record<string, unknown>[];
const logs = rows.map((row) => ({
id: row.id,
orgId: row.org_id,
userId: row.user_id,
userEmail: row.user_email,
timestamp: new Date(row.timestamp),
action: row.action,
resourceType: row.resource_type,
resourceId: row.resource_id,
resourceName: row.resource_name,
clusterId: row.cluster_id,
ipAddress: row.ip_address,
userAgent: row.user_agent,
requestId: row.request_id,
changes: row.changes,
status: row.status,
errorMessage: row.error_message,
id: row.id as string,
orgId: row.org_id as string,
userId: row.user_id as string,
userEmail: row.user_email as string,
timestamp: new Date(row.timestamp as string),
action: row.action as string,
resourceType: row.resource_type as string,
resourceId: row.resource_id as string,
resourceName: row.resource_name as string,
clusterId: (row.cluster_id as string) || null,
ipAddress: row.ip_address as string,
userAgent: row.user_agent as string,
requestId: row.request_id as string,
changes: row.changes as Record<string, unknown>,
status: row.status as string,
errorMessage: row.error_message as string | undefined,
}));
return { logs, total };
}
// ==================== Cluster Metrics Queries ====================
export async function queryClusterOverview(
clusterId?: string
): Promise<{
totalStreams: number;
totalConsumers: number;
totalMessages: number;
messageRate: number;
activeAlerts: number;
}> {
const ch = getClickHouseClient();
// Get latest stream metrics for counts
const streamQuery = clusterId
? `
SELECT
count(DISTINCT stream_name) as total_streams,
sum(consumer_count) as total_consumers,
sum(messages_total) as total_messages,
avg(messages_rate) as message_rate
FROM stream_metrics
WHERE cluster_id = {clusterId:UUID}
AND timestamp >= now() - INTERVAL 5 MINUTE
`
: `
SELECT
count(DISTINCT stream_name) as total_streams,
count(DISTINCT cluster_id, consumer_name) as total_consumers,
sum(messages_total) as total_messages,
avg(messages_rate) as message_rate
FROM stream_metrics
WHERE timestamp >= now() - INTERVAL 5 MINUTE
`;
try {
const result = await ch.query({
query: streamQuery,
query_params: clusterId ? { clusterId } : {},
format: 'JSONEachRow',
});
const rows = await result.json() as Record<string, unknown>[];
const row = rows[0] || {};
return {
totalStreams: Number(row.total_streams) || 0,
totalConsumers: Number(row.total_consumers) || 0,
totalMessages: Number(row.total_messages) || 0,
messageRate: Number(row.message_rate) || 0,
activeAlerts: 0, // Alerts are managed separately
};
} catch {
return {
totalStreams: 0,
totalConsumers: 0,
totalMessages: 0,
messageRate: 0,
activeAlerts: 0,
};
}
}
export async function queryOverviewMetrics(
clusterId?: string,
timeRange: string = '1h'
): Promise<{
totalMessages: number;
totalBytes: number;
avgThroughput: number;
avgLatency: number;
messagesTrend: number;
bytesTrend: number;
throughputTrend: number;
latencyTrend: number;
}> {
const ch = getClickHouseClient();
// Parse time range
const rangeHours: Record<string, number> = {
'1h': 1,
'6h': 6,
'24h': 24,
'7d': 168,
};
const hours = rangeHours[timeRange] || 1;
const clusterFilter = clusterId ? 'AND cluster_id = {clusterId:UUID}' : '';
try {
// Current period metrics
const currentResult = await ch.query({
query: `
SELECT
sum(messages_total) as total_messages,
sum(bytes_total) as total_bytes,
avg(messages_rate) as avg_throughput
FROM stream_metrics
WHERE timestamp >= now() - INTERVAL ${hours} HOUR
${clusterFilter}
`,
query_params: clusterId ? { clusterId } : {},
format: 'JSONEachRow',
});
const currentRows = await currentResult.json() as Record<string, unknown>[];
const current = currentRows[0] || {};
// Previous period for trend calculation
const previousResult = await ch.query({
query: `
SELECT
sum(messages_total) as total_messages,
sum(bytes_total) as total_bytes,
avg(messages_rate) as avg_throughput
FROM stream_metrics
WHERE timestamp >= now() - INTERVAL ${hours * 2} HOUR
AND timestamp < now() - INTERVAL ${hours} HOUR
${clusterFilter}
`,
query_params: clusterId ? { clusterId } : {},
format: 'JSONEachRow',
});
const previousRows = await previousResult.json() as Record<string, unknown>[];
const previous = previousRows[0] || {};
// Calculate trends (percentage change)
const calcTrend = (curr: number, prev: number) => {
if (prev === 0) return curr > 0 ? 100 : 0;
return ((curr - prev) / prev) * 100;
};
return {
totalMessages: Number(current.total_messages) || 0,
totalBytes: Number(current.total_bytes) || 0,
avgThroughput: Number(current.avg_throughput) || 0,
avgLatency: 0, // Latency requires message-level tracking
messagesTrend: calcTrend(
Number(current.total_messages) || 0,
Number(previous.total_messages) || 0
),
bytesTrend: calcTrend(
Number(current.total_bytes) || 0,
Number(previous.total_bytes) || 0
),
throughputTrend: calcTrend(
Number(current.avg_throughput) || 0,
Number(previous.avg_throughput) || 0
),
latencyTrend: 0,
};
} catch {
return {
totalMessages: 0,
totalBytes: 0,
avgThroughput: 0,
avgLatency: 0,
messagesTrend: 0,
bytesTrend: 0,
throughputTrend: 0,
latencyTrend: 0,
};
}
}
// ==================== Helpers ====================
function parseInterval(interval: string): number {
+34
View File
@@ -298,6 +298,40 @@ export async function deleteConsumer(clusterId: string, streamName: string, cons
return jsm.consumers.delete(streamName, consumerName);
}
export async function pauseConsumer(
clusterId: string,
streamName: string,
consumerName: string,
pauseUntil?: Date
): Promise<ConsumerInfo> {
const jsm = getClusterJetStreamManager(clusterId);
// Get current consumer info
const consumer = await jsm.consumers.info(streamName, consumerName);
// Update with pause_until set to a future time
const until = pauseUntil || new Date(Date.now() + 100 * 365 * 24 * 60 * 60 * 1000); // 100 years if no time specified
return jsm.consumers.update(streamName, consumerName, {
...consumer.config,
pause_until: until,
});
}
export async function resumeConsumer(
clusterId: string,
streamName: string,
consumerName: string
): Promise<ConsumerInfo> {
const jsm = getClusterJetStreamManager(clusterId);
// Get current consumer info
const consumer = await jsm.consumers.info(streamName, consumerName);
// Update with pause_until removed (set to epoch/past time)
return jsm.consumers.update(streamName, consumerName, {
...consumer.config,
pause_until: undefined,
});
}
// ==================== Message Operations ====================
export async function publishMessage(
+64
View File
@@ -145,3 +145,67 @@ export async function getUserPermissions(userId: string, orgId: string): Promise
export async function invalidateUserPermissions(userId: string, orgId: string): Promise<void> {
await redis.del(`${PERMISSIONS_PREFIX}${userId}:${orgId}`);
}
// Pub/Sub for real-time metrics
const METRICS_CHANNEL = 'metrics';
const ALERTS_CHANNEL = 'alerts';
// Subscriber client for pub/sub (separate from main client)
let subscriber: Redis | null = null;
const messageHandlers = new Map<string, Set<(message: any) => void>>();
export function getSubscriber(): Redis {
if (!subscriber) {
subscriber = new Redis(config.REDIS_URL, {
maxRetriesPerRequest: 3,
retryStrategy(times) {
const delay = Math.min(times * 50, 2000);
return delay;
},
});
subscriber.on('message', (channel, message) => {
try {
const data = JSON.parse(message);
const handlers = messageHandlers.get(channel);
handlers?.forEach((handler) => handler(data));
} catch (err) {
console.error('Failed to parse Redis message:', err);
}
});
}
return subscriber;
}
export function subscribeToChannel(channel: string, handler: (data: any) => void): () => void {
const sub = getSubscriber();
if (!messageHandlers.has(channel)) {
messageHandlers.set(channel, new Set());
sub.subscribe(channel);
}
messageHandlers.get(channel)!.add(handler);
// Return unsubscribe function
return () => {
const handlers = messageHandlers.get(channel);
if (handlers) {
handlers.delete(handler);
if (handlers.size === 0) {
sub.unsubscribe(channel);
messageHandlers.delete(channel);
}
}
};
}
export async function publishMetrics(data: any): Promise<void> {
await redis.publish(METRICS_CHANNEL, JSON.stringify(data));
}
export async function publishAlert(data: any): Promise<void> {
await redis.publish(ALERTS_CHANNEL, JSON.stringify(data));
}
export { METRICS_CHANNEL, ALERTS_CHANNEL };
+21
View File
@@ -2,6 +2,7 @@ import { FastifyInstance } from 'fastify';
import { WebSocket, WebSocketServer } from 'ws';
import { verifyToken } from '../common/middleware/auth';
import { logger } from './logger';
import { subscribeToChannel, METRICS_CHANNEL, ALERTS_CHANNEL } from './redis';
interface Client {
ws: WebSocket;
@@ -76,9 +77,29 @@ export function setupWebSocket(fastify: FastifyInstance): WebSocketServer {
}
});
// Set up Redis subscription bridge for real-time updates
setupRedisBridge();
return wss;
}
function setupRedisBridge() {
// Subscribe to metrics channel and broadcast to WebSocket clients
subscribeToChannel(METRICS_CHANNEL, (data) => {
const channel = data.channel || 'metrics';
broadcast(channel, data);
logger.debug({ channel }, 'Broadcasting metrics to WebSocket clients');
});
// Subscribe to alerts channel and broadcast to WebSocket clients
subscribeToChannel(ALERTS_CHANNEL, (data) => {
broadcast('alerts', data);
logger.debug('Broadcasting alert to WebSocket clients');
});
logger.info('Redis to WebSocket bridge established');
}
function handleMessage(clientId: string, client: Client, message: any) {
switch (message.type) {
case 'subscribe':
@@ -4,6 +4,8 @@ import {
queryStreamMetrics,
queryConsumerMetrics,
queryAuditLogs,
queryClusterOverview,
queryOverviewMetrics,
} from '../../lib/clickhouse';
import { authenticate } from '../../common/middleware/auth';
@@ -106,14 +108,8 @@ export const analyticsRoutes: FastifyPluginAsync = async (fastify) => {
fastify.get<{ Querystring: { clusterId?: string } }>(
'/cluster/overview',
async (request) => {
// TODO: Implement cluster overview metrics
return {
totalStreams: 0,
totalConsumers: 0,
totalMessages: 0,
messageRate: 0,
activeAlerts: 0,
};
const { clusterId } = request.query;
return queryClusterOverview(clusterId);
}
);
@@ -121,17 +117,8 @@ export const analyticsRoutes: FastifyPluginAsync = async (fastify) => {
fastify.get<{ Querystring: { clusterId?: string; timeRange?: string } }>(
'/overview',
async (request) => {
// TODO: Implement real metrics from ClickHouse
return {
totalMessages: 0,
totalBytes: 0,
avgThroughput: 0,
avgLatency: 0,
messagesTrend: 0,
bytesTrend: 0,
throughputTrend: 0,
latencyTrend: 0,
};
const { clusterId, timeRange } = request.query;
return queryOverviewMetrics(clusterId, timeRange || '1h');
}
);
};
@@ -93,4 +93,39 @@ export const consumerRoutes: FastifyPluginAsync = async (fastify) => {
return { consumer };
}
);
// POST /clusters/:cid/streams/:sid/consumers/:name/pause - Pause consumer
fastify.post<{
Params: { cid: string; sid: string; name: string };
Body: { pauseUntil?: string };
}>(
'/:cid/streams/:sid/consumers/:name/pause',
async (request) => {
const pauseUntil = request.body?.pauseUntil
? new Date(request.body.pauseUntil)
: undefined;
const consumer = await consumerService.pauseConsumer(
request.user!.orgId,
request.params.cid,
request.params.sid,
request.params.name,
pauseUntil
);
return { consumer, paused: true };
}
);
// POST /clusters/:cid/streams/:sid/consumers/:name/resume - Resume consumer
fastify.post<{ Params: { cid: string; sid: string; name: string } }>(
'/:cid/streams/:sid/consumers/:name/resume',
async (request) => {
const consumer = await consumerService.resumeConsumer(
request.user!.orgId,
request.params.cid,
request.params.sid,
request.params.name
);
return { consumer, resumed: true };
}
);
};
@@ -5,6 +5,8 @@ import {
createConsumer as natsCreateConsumer,
updateConsumer as natsUpdateConsumer,
deleteConsumer as natsDeleteConsumer,
pauseConsumer as natsPauseConsumer,
resumeConsumer as natsResumeConsumer,
} from '../../lib/nats';
import { NotFoundError } from '../../../../shared/src/index';
import type { ConsumerInfo, CreateConsumerInput, UpdateConsumerInput } from '../../../../shared/src/index';
@@ -228,3 +230,42 @@ function mapReplayPolicy(policy: string): any {
};
return map[policy] || 'instant';
}
// ==================== Pause/Resume Operations ====================
export async function pauseConsumer(
orgId: string,
clusterId: string,
streamName: string,
consumerName: string,
pauseUntil?: Date
): Promise<ConsumerInfo> {
// Verify cluster belongs to org
const cluster = await prisma.natsCluster.findFirst({
where: { id: clusterId, orgId },
});
if (!cluster) {
throw new NotFoundError('Cluster', clusterId);
}
return natsPauseConsumer(clusterId, streamName, consumerName, pauseUntil);
}
export async function resumeConsumer(
orgId: string,
clusterId: string,
streamName: string,
consumerName: string
): Promise<ConsumerInfo> {
// Verify cluster belongs to org
const cluster = await prisma.natsCluster.findFirst({
where: { id: clusterId, orgId },
});
if (!cluster) {
throw new NotFoundError('Cluster', clusterId);
}
return natsResumeConsumer(clusterId, streamName, consumerName);
}