更新依赖项,优化 WebSocket 处理,添加文件流管道功能,改进用户认证逻辑

This commit is contained in:
2025-12-21 02:06:38 +08:00
parent d22be3a840
commit 15fcfdad18
16 changed files with 150 additions and 262 deletions

View File

@@ -38,6 +38,21 @@ export const checkAuth = async (req: http.IncomingMessage, res: http.ServerRespo
return { tokenUser, token };
};
export const getLoginUserByToken = async (token: string) => {
if (token) {
token = token.replace('Bearer ', '');
}
if (!token) {
return null;
}
let tokenUser;
try {
tokenUser = await User.verifyToken(token);
return { tokenUser, token };
} catch (e) {
return null;
}
}
export const getLoginUser = async (req: http.IncomingMessage) => {
let token = (req.headers?.['authorization'] as string) || (req.headers?.['Authorization'] as string) || '';
const url = new URL(req.url || '', 'http://localhost');

View File

@@ -3,7 +3,7 @@ import { useFileStore } from '@kevisual/use-config/file-store';
import { minioResources } from './minio.ts';
export const config = useConfig() as any;
export const port = config.PORT || 4005;
export const port = config.PORT ? Number(config.PORT) : 4005;
export const fileStore = useFileStore('pages');
type ConfigType = {
api: {

View File

@@ -7,4 +7,6 @@ export * from './get-router.ts'
export * from './get-content-type.ts'
export * from './utils.ts'
export * from './utils.ts'
export { pipeFileStream, pipeStream } from './pipe.ts'

View File

@@ -0,0 +1,19 @@
import * as http from 'http';
import * as fs from 'fs';
import { isBun } from '../../utils/get-engine.ts';
export const pipeFileStream = (filePath: string, res: http.ServerResponse) => {
const readStream = fs.createReadStream(filePath);
if (isBun) {
res.pipe(readStream as any);
} else {
readStream.pipe(res, { end: true });
}
}
export const pipeStream = (readStream: fs.ReadStream, res: http.ServerResponse) => {
if (isBun) {
res.pipe(readStream as any);
} else {
readStream.pipe(res, { end: true });
}
}

View File

@@ -10,7 +10,7 @@ import { addStat } from '@/modules/html/stat/index.ts';
import path from 'path';
import { getTextContentType } from '@/modules/fm-manager/index.ts';
import { logger } from '@/modules/logger.ts';
import { pipeStream } from '../pipe.ts';
const pipelineAsync = promisify(pipeline);
export async function downloadFileFromMinio(fileUrl: string, destFile: string) {
@@ -74,7 +74,7 @@ export async function minioProxy(
res.writeHead(200, {
...headers,
});
objectStream.pipe(res, { end: true });
pipeStream(objectStream as any, res);
}
return true;
} catch (error) {
@@ -154,7 +154,7 @@ export const httpProxy = async (
res.writeHead(proxyRes.statusCode, {
...headers,
});
proxyRes.pipe(res, { end: true });
pipeStream(proxyRes as any, res);
}
});
proxyReq.on('error', (err) => {

View File

@@ -1,4 +1,4 @@
import http from 'http';
import http from 'node:http';
import { minioClient } from '@/modules/minio.ts';
type ProxyInfo = {
path?: string;

View File

@@ -1,5 +1,6 @@
import { IncomingMessage } from 'node:http';
import http from 'node:http';
import { logger } from '../logger.ts';
export const getUserFromRequest = (req: IncomingMessage) => {
const url = new URL(req.url, `http://${req.headers.host}`);
@@ -14,8 +15,8 @@ export const getUserFromRequest = (req: IncomingMessage) => {
export const getDNS = (req: http.IncomingMessage) => {
const hostName = req.headers.host;
const ip = req.socket.remoteAddress;
const hostName = req.headers?.host;
const ip = req?.socket?.remoteAddress || '';
return { hostName, ip };
};

View File

@@ -1,136 +1,29 @@
import { WebSocketServer } from 'ws';
import { nanoid } from 'nanoid';
import { WsProxyManager } from './manager.ts';
import { getLoginUser } from '@/modules/auth.ts';
import { getLoginUserByToken } from '@/modules/auth.ts';
import { logger } from '../logger.ts';
export const wsProxyManager = new WsProxyManager();
export const upgrade = (request: any, socket: any, head: any) => {
const req = request as any;
const url = new URL(req.url, 'http://localhost');
const id = url.searchParams.get('id');
if (url.pathname === '/ws/proxy') {
console.log('upgrade', request.url, id);
wss.handleUpgrade(req, socket, head, (ws) => {
// 这里手动触发 connection 事件
console.log('emitting connection event');
// @ts-ignore
wss.emit('connection', ws, req);
});
return true;
}
return false;
};
export const wss = new WebSocketServer({
noServer: true,
path: '/ws/proxy',
});
wss.on('connection', async (ws, req) => {
console.log('connected', req.url);
const url = new URL(req.url, 'http://localhost');
const _id = url?.searchParams?.get('id');
const id = _id || nanoid();
const loginUser = await getLoginUser(req);
if (!loginUser) {
console.log('未登录,断开连接');
ws.send(JSON.stringify({ code: 401, message: '未登录' }));
setTimeout(() => {
ws.close();
}, 1000);
return;
}
const user = loginUser.tokenUser.username;
const userApp = user + '-' + id;
console.log('注册 ws 连接', userApp);
wsProxyManager.register(userApp, { user, ws });
ws.send(
JSON.stringify({
type: 'connected',
user: user,
id,
}),
);
ws.on('message', async (event: Buffer) => {
const eventData = event.toString();
if (!eventData) {
return;
}
const data = JSON.parse(eventData);
logger.debug('message', data);
});
ws.on('close', () => {
logger.debug('ws closed');
wsProxyManager.unregister(userApp);
});
});
export class WssApp {
wss: WebSocketServer;
bunWSS = websocket;
constructor() {
this.wss = wss;
}
upgrade(request: any, socket: any, head: any) {
// return upgrade(request, socket, head);
return bunUpgrade(request);
}
}
export const bunUpgrade = (request: Request) => {
const url = new URL(request.url, 'http://localhost');
const isUpgrade = url.pathname === '/ws/proxy';
if (isUpgrade) {
console.log('upgrade', request.url);
// 使用 Bun 原生 WebSocket
new Response(null, {
status: 101,
headers: {
'Upgrade': 'websocket',
},
});
return true;
}
return false;
};
// Bun WebSocket 处理器
export const websocket = {
async open(ws: any) {
console.log('WebSocket opened');
const { url, token } = ws.data;
const urlObj = new URL(url, 'http://localhost');
const _id = urlObj.searchParams.get('id');
const id = _id || nanoid();
// 创建一个模拟的 request 对象用于认证
const mockReq: any = {
url: url,
headers: {
authorization: token ? `Bearer ${token}` : undefined,
},
};
const loginUser = await getLoginUser(mockReq);
if (!loginUser) {
console.log('未登录,断开连接');
import { WebScoketListenerFun } from '@kevisual/router/src/server/server-type.ts'
export const wssFun: WebScoketListenerFun = async (req, res) => {
// do nothing, just to enable ws upgrade event
const { id, ws, token, data, emitter } = req;
logger.debug('ws proxy connected, id=', id, ' token=', token, ' data=', data);
// console.log('req', req)
const { type } = data || {};
if (type === 'registryClient') {
const loginUser = await getLoginUserByToken(token);
if (!loginUser?.tokenUser) {
logger.debug('未登录,断开连接');
ws.send(JSON.stringify({ code: 401, message: '未登录' }));
ws.close();
setTimeout(() => {
ws.close(401, 'Unauthorized');
}, 1000);
return;
}
const user = loginUser.tokenUser.username;
const user = loginUser?.tokenUser?.username;
const userApp = user + '-' + id;
console.log('注册 ws 连接', userApp);
ws.data.userApp = userApp;
ws.data.user = user;
logger.debug('注册 ws 连接', userApp);
// @ts-ignore
wsProxyManager.register(userApp, { user, ws });
ws.send(
JSON.stringify({
type: 'connected',
@@ -138,26 +31,22 @@ export const websocket = {
id,
}),
);
},
async message(ws: any, message: string) {
try {
const data = JSON.parse(message);
logger.debug('message', data);
} catch (error) {
logger.error('Failed to parse message', error);
}
},
close(ws: any) {
const { userApp } = ws.data;
logger.debug('ws closed', userApp);
if (userApp) {
emitter.once('close--' + id, () => {
logger.debug('ws emitter closed');
wsProxyManager.unregister(userApp);
}
},
error(ws: any, error: Error) {
console.error('WebSocket error:', error);
},
};
});
// @ts-ignore
ws.data.userApp = userApp;
return;
}
// @ts-ignore
const userApp = ws.data.userApp;
logger.debug('message', data, ' userApp=', userApp);
const wsMessage = wsProxyManager.get(userApp);
if (wsMessage) {
wsMessage.sendResponse(data);
} else {
// @ts-ignore
logger.debug('账号应用未注册,无法处理消息。未授权?', ws.data);
}
}

View File

@@ -10,18 +10,11 @@ class WsMessage {
this.ws = ws;
this.user = user;
this.emitter = new EventEmitter();
this.listenMessage();
}
async listenMessage() {
this.ws.on('message', (event: Buffer) => {
const eventData = event.toString();
if (!eventData) {
return;
}
const data = JSON.parse(eventData);
logger.debug('ws-proxy listenMessage', data);
this.emitter.emit(data.id, data.data);
});
async sendResponse(data: any) {
if (data.id) {
this.emitter.emit(data.id, data?.data);
}
}
async sendData(data: any, opts?: { timeout?: number }) {
if (this.ws.readyState !== WebSocket.OPEN) {
@@ -65,6 +58,7 @@ export class WsProxyManager {
value.ws.close();
}
}
console.log('WsProxyManager register', id);
const value = new WsMessage({ ws: opts?.ws, user: opts?.user });
this.wssMap.set(id, value);
}

View File

@@ -30,7 +30,7 @@ export const UserV1Proxy = async (req: IncomingMessage, res: ServerResponse, opt
const client = wsProxyManager.get(userAppKey);
const ids = wsProxyManager.getIds();
if (!client) {
opts?.createNotFoundPage?.(`未找到应用, ${userAppKey}, ${ids.join(',')}`);
opts?.createNotFoundPage?.(`未找到应用 [${userAppKey}], 当前应用列表: ${ids.join(',')}`);
return false;
}
const value = await client.sendData(data);