knight-server
v1.3.9
Published
Knight server architecture library
Maintainers
Readme
Knight-Server 服务插件使用手册
概述
Knight-Server 是一个基于 Express + WebSocket 的 Node.js 服务端框架,集成了 MySQL、Neo4j、Redis 三大一等模块,通过装饰器 @OnRegister 实现路由自动注册,提供开箱即用的 HTTP/WebSocket/MySQL/Neo4j/Redis 全栈开发体验。
- 包名:
knight-server - 入口:
dist/index.js - 类型声明:
dist/index.d.ts
一、安装与启动
安装
npm install knight-server最小启动示例
import { Server } from "knight-server";
const server = Server.Instance;
await server.OnStart({
mysql: { /* ... */ },
neo4j: { /* ... */ },
redis: { /* ... */ },
http: { port: 3000 },
websocket: {
port: 6666,
heartbeat_timeout_interval: 30000,
router_timeout_interval: 60000,
},
serverHost: "127.0.0.1",
});Server 是单例类,构造函数私有,通过 Server.Instance 访问。OnStart 按 mysql → neo4j → redis → http → websocket 顺序启动各模块,重复调用会被忽略。
二、服务配置 (ServerConfig)
export interface ServerConfig {
mysql?: MySQLConfig; // 可选,关系型数据库
neo4j?: Neo4jConfig; // 可选,图数据库
redis?: RedisConfig; // 可选,缓存
http?: HttpConfig; // 可选,HTTP 服务
websocket?: WebSocketConfig; // 可选,WebSocket 服务
token?: TokenConfig; // 可选,JWT 配置
log?: LogConfig; // 可选,日志落盘配置,详见 9.1
serverHost: string; // 服务器对外地址,如 "127.0.0.1" 或域名
}2.1 MySQLConfig
interface MySQLConfig {
host: string;
port: number;
user: string;
password: string;
database: string;
connectionLimit: number; // 连接池大小
tableRelativePath: string; // .sql 表定义文件目录(相对于项目根目录)
}模块启动时会自动校验必填字段,并调用 SmartTable 根据 .sql 文件同步数据库表结构。
2.2 Neo4jConfig
interface Neo4jConfig {
uri: string; // 连接 URI,如 bolt://localhost:7687 或 neo4j://localhost:7687
username: string; // 用户名
password: string; // 密码
database?: string; // 数据库名(可选,默认 neo4j)
}模块启动时通过 verifyConnectivity 校验连接,连接成功后才注册 Neo4j 路由,并预加载 Embedding 模型(基于 transformers.js,失败仅告警不阻塞启动)。
2.3 RedisConfig
interface RedisConfig {
host: string;
port: number;
password?: string;
db?: number;
url?: string; // 连接 URL,优先级高于 host/port
connectTimeout?: number; // 连接超时(ms)
heartbeatInterval?: number; // 心跳间隔(ms),默认 240000(4分钟)
reconnectMaxRetries?: number; // 默认重连策略的最大重试次数,默认 10
reconnectStrategy?: RedisReconnectStrategy; // 自定义重连策略,配置后忽略 reconnectMaxRetries
recoveryInterval?: number; // 连接自动恢复检查间隔(ms),默认 10000
flushOnReconnect?: boolean; // 恢复就绪后是否清空当前 db,默认 false(除冷启动一次连上外都会清;前提:该 db 由本应用独占)
}模块启动后维护两个客户端:主客户端执行命令,订阅客户端处理 Pub/Sub。内置心跳机制,每隔 heartbeatInterval 毫秒向主客户端和订阅客户端发送 PING 命令,防止服务端因空闲超时断开连接(如 Redis timeout 配置)。重连时自动恢复所有频道和模式订阅。
断线重连由底层 node-redis 驱动自动完成(指数退避);启动时连接失败不中断模块启动,路由注册与后续逻辑照常执行;驱动重试耗尽后由框架的连接看门狗定期自动恢复。也可调用 server.redis.reconnect() 手动恢复。详见 8.2 心跳保活与断线重连。
2.4 HttpConfig
interface HttpConfig {
port: number;
}2.5 WebSocketConfig
interface WebSocketConfig {
port: number;
heartbeat_timeout_interval: number; // 心跳超时间隔(ms),超时客户端会被关闭(code 1001)
router_timeout_interval: number; // 非持久路由闲置检测间隔(ms),无客户端时自动注销
}2.6 TokenConfig (可选)
interface TokenConfig {
secret: string;
accessExpiresIn?: string; // 默认 "15m"
refreshExpiresIn?: string; // 默认 "7d"
}此配置存放在 server.config.token 中供业务代码使用,框架本身不自动接入鉴权逻辑——你需要在 Router 的 OnCheck 中自行调用 TokenUtility.verify()。
三、路由注册机制
3.1 @OnRegister 装饰器
import { OnRegister } from "knight-server";
@OnRegister("/api/users")
export class UserRouter extends HttpRouter { /* ... */ }行为:
- 将路由元数据
{ path, type, permanent: true }推入全局RouteRegistry数组 path自动规范化为小写且以/开头- 各 Module 在
OnStart时通过instanceof过滤出归属自己的 Router 并实例化 - 装饰器标记的路由
permanent默认为true,不会被自动清理
3.2 手动注册 / 注销
// 动态注册
server.http.OnRegistRouter({
path: "/api/dynamic",
type: DynamicRouter,
permanent: false,
});
// 动态注销
server.http.OnDegistRouter("/api/dynamic");四、HTTP 模块
4.1 HttpRouter 抽象类
继承 HttpRouter 后需实现以下方法:
| 方法 | 签名 | 说明 |
|------|------|------|
| OnCheck | (server, request: Request, response: Response): Promise<CheckResult> | 鉴权/校验 |
| OnGetHandler | (server, checkResult, request, response): Promise<void> | GET 请求 |
| OnPostHandler | 同上 | POST 请求 |
| OnPutHandler | 同上 | PUT 请求 |
| OnDeleteHandler | 同上 | DELETE 请求 |
| OnOptionsHandler | 同上 | OPTIONS 请求 |
| OnHeadHandler | 同上 | HEAD 请求 |
| OnPatchHandler | 同上 | PATCH 请求 |
| OnTraceHandler | 同上 | TRACE 请求 |
4.2 调度流程
HTTP 请求 → express.all(path) 代理层 → 查找 router → OnCheck
├─ valid=false → 返回 401
└─ valid=true → 按 method 分发到对应 OnXxxHandler4.3 示例
import { HttpRouter, OnRegister, CheckResult, TokenUtility } from "knight-server";
@OnRegister("/api/user/info")
export class UserInfoRouter extends HttpRouter {
async OnCheck(server, request, response): Promise<CheckResult> {
const token = request.headers.authorization?.replace("Bearer ", "");
if (!token) return { valid: false, error: "请先登录" };
try {
const payload = TokenUtility.verify(token, server.config.token.secret);
return { valid: true, payload };
} catch {
return { valid: false, error: "令牌无效" };
}
}
async OnGetHandler(server, checkResult, request, response) {
const users = await server.mysql.find("users", "id = ?", [checkResult.payload.uid]);
response.json({ success: true, data: users.data });
}
async OnPostHandler(server, checkResult, request, response) {
const result = await server.mysql.insert("users", request.body);
response.json(result);
}
// 未实现的 method handler 会导致请求无响应,建议至少实现 OnOptionsHandler
async OnOptionsHandler(server, checkResult, request, response) {
response.status(204).end();
}
}五、WebSocket 模块
5.1 WebSocketRouter 抽象类
abstract class WebSocketRouter {
// 鉴权检查
abstract OnCheck(server: Server, url: URL): Promise<CheckResult>;
// 客户端连接
abstract OnConnect(server: Server, checkResult: CheckResult, client: WebSocketClient): Promise<void>;
// 客户端断开
abstract OnDisconnect(server: Server, checkResult: CheckResult, client: WebSocketClient): Promise<void>;
// 收到消息
abstract OnMessage(server: Server, checkResult: CheckResult, client: WebSocketClient, data: string): Promise<void>;
}5.2 WebSocketClient 关键方法
| 方法 | 说明 |
|------|------|
| OnSend(data) | 向该客户端发送消息 |
| OnBroadcast(data, excludeSelf?) | 向同路由下所有客户端广播 |
| OnUpdateHeartbeat() | 刷新心跳时间戳(收到消息时调用;协议层 Ping 帧会自动刷新) |
| OnChangeRouter(router) | 将客户端迁移到另一个路由 |
| OnGetSearchParams(key) | 获取连接 URL 上的查询参数 |
| OnGetCheckResultPayload<T>() | 获取鉴权时写入的 payload |
5.3 关键属性
| 属性 | 说明 |
|------|------|
| client.key | "ip:port" 格式的唯一标识 |
| client.checkResult | 鉴权结果 |
5.4 WebSocketRouter 客户端管理
| 方法 | 说明 |
|------|------|
| OnGetClient(key) | 按 key 获取客户端 |
| OnGetClients() | 获取所有客户端 |
| OnGetClientBySearchParams(key, value) | 按 URL 参数查找 |
| OnGetClientByCheckResultPayload(key, value) | 按鉴权 payload 字段查找 |
5.5 连接流程
WebSocket 连接 → 按 pathname 找到 router → router.OnHandler
→ 去重检查(ip:port)
→ OnCheck(鉴权)
→ 创建 WebSocketClient → 注册到 client pool
→ OnConnect
→ 启动心跳检测定时器客户端发送 Ping 帧时,服务端自动回复 Pong 并刷新心跳时间戳;此外在 OnMessage 中调用 client.OnUpdateHeartbeat() 也可维持心跳。
5.6 示例
import { WebSocketRouter, WebSocketClient, OnRegister, CheckResult, Server } from "knight-server";
@OnRegister("/ws/chat")
export class ChatRouter extends WebSocketRouter {
async OnCheck(server: Server, url: URL): Promise<CheckResult> {
const token = url.searchParams.get("token");
if (!token) return { valid: false, error: "缺少令牌" };
// ... 验证 token
return { valid: true, payload: { userId: "xxx" } };
}
async OnConnect(server, checkResult, client) {
client.OnBroadcast(JSON.stringify({
type: "join",
userId: (checkResult.payload as any).userId,
}));
}
async OnDisconnect(server, checkResult, client) {
client.OnBroadcast(JSON.stringify({
type: "leave",
userId: (checkResult.payload as any).userId,
}));
}
async OnMessage(server, checkResult, client, data) {
client.OnUpdateHeartbeat(); // 刷新心跳
client.OnBroadcast(data, true); // 广播给同路由其他人
}
}六、MySQL 模块
6.1 访问方式
// 直接通过 server.mysql 调用
const result = await server.mysql.find("users", "age > ?", [18]);6.2 MySQLRouter 钩子
@OnRegister("/users")
export class UsersRouter extends MySQLRouter {
// 访问控制
async OnCheck(server): Promise<CheckResult> {
if (!server.mysql.isConnected()) return { valid: false, error: "数据库未连接" };
return { valid: true };
}
// 插入前钩子
async OnBeforeInsert(data: any): Promise<HookResult> {
return { proceed: true, data };
}
// 更新前钩子
async OnBeforeUpdate(data: any, where: string, params: any[]): Promise<HookResult> {
return { proceed: true, data };
}
// 删除前钩子
async OnBeforeDelete(where: string, params: any[]): Promise<HookResult> {
return { proceed: true };
}
// 查询后钩子(脱敏、转换等)
async OnAfterSelect(rows: any[]): Promise<any[]> {
return rows.map(({ password, ...rest }) => rest);
}
}6.3 HookResult
interface HookResult<T = any> {
proceed: boolean; // true=继续, false=阻止
data?: T; // 修改后的数据
error?: string; // 阻止时的错误信息
}6.4 CRUD 方法一览
所有方法返回 Promise<QueryResult<T>>:
interface QueryResult<T = any> {
success: boolean;
data?: T[]; // 多行结果
item?: T; // 单行结果 (findOne, findById)
affectedRows?: number; // insert/update/delete
insertId?: number; // insert
error?: string; // 失败时的错误信息
}查询:
| 方法 | 说明 |
|------|------|
| find<T>(table, where?, params?) | 查询多条 |
| findOne<T>(table, where, params?) | 查询单条 (含 .item) |
| findById<T>(table, id, idField?) | 按 ID 查询 |
| paginate<T>(table, page?, pageSize?, where?, params?, orderBy?) | 分页查询 |
写入:
| 方法 | 说明 |
|------|------|
| insert<T>(table, data) | 插入单条 |
| insertBatch<T>(table, dataArray) | 批量插入 |
| update<T>(table, data, where, whereParams?) | 条件更新 |
| updateById<T>(table, id, data, idField?) | 按 ID 更新 |
| delete(table, where, params?) | 条件删除 |
| deleteById(table, id, idField?) | 按 ID 删除 |
高级:
| 方法 | 说明 |
|------|------|
| query<T>(sql, params?) | 执行自定义 SQL(绕过钩子与访问控制) |
| transaction<T>(callback) | 事务执行,自动 begin/commit/rollback/release |
6.5 事务示例
const result = await server.mysql.transaction(async (connection) => {
await connection.execute("UPDATE accounts SET balance = balance - 100 WHERE id = ?", [1]);
await connection.execute("UPDATE accounts SET balance = balance + 100 WHERE id = ?", [2]);
return { msg: "转账成功" };
});
if (!result.success) {
console.error(result.error); // 事务已自动回滚
}七、Neo4j 模块
7.1 访问方式
// 直接通过 server.neo4j 调用(默认剥离每条结果的 embedding 字段)
const result = await server.neo4j.query("MATCH (n:User) RETURN n");
// 需要保留向量:传第三参 { keepEmbedding: true };自定义向量字段名用 { embeddingField }
const full = await server.neo4j.query("MATCH (n:Tool) RETURN n", {}, { keepEmbedding: true });7.2 Neo4jRouter 钩子
@OnRegister("/users")
export class UsersRouter extends Neo4jRouter {
// 访问控制
async OnCheck(server): Promise<CheckResult> {
if (!server.neo4j.isConnected()) return { valid: false, error: "Neo4j未连接" };
return { valid: true };
}
// 创建节点前拦截
async OnBeforeCreateNode(label, props): Promise<GraphHookResult> {
props.created_at = Date.now();
return { proceed: true, data: props };
}
// 更新节点前拦截
async OnBeforeUpdateNode(label, where, params, props): Promise<GraphHookResult> {
props.updated_at = Date.now();
return { proceed: true, data: props };
}
// 删除节点前拦截
async OnBeforeDeleteNode(label, where, params): Promise<GraphHookResult> {
return { proceed: true };
}
// 查询后拦截(脱敏、转换等)
async OnAfterSelect(records: any[]): Promise<any[]> {
return records.map(({ password, ...rest }) => rest);
}
}访问控制说明:
OnCheck会被所有节点与关系操作(createNode/findNodes/updateNode/createRelationship/findNeighbors等)自动调用;返回valid: false时该操作被拦截并返回{ success: false, error }。基类默认实现已包含isConnected()检查,子类覆盖时若需保留可自行调用。关系操作会分别校验from/to两个标签对应的路由。
7.3 GraphHookResult
interface GraphHookResult<T = any> {
proceed: boolean; // true=继续, false=阻止
data?: T; // 修改后的数据(仅 create/update)
error?: string; // 阻止时的错误信息
}7.4 图查询结果 GraphQueryResult
interface GraphQueryResult<T = any> {
success: boolean;
data?: T[]; // 映射后的数据 (query)
item?: T; // 首条记录 (query)
records?: Neo4jRecord[]; // 原始记录
summary?: ResultSummary; // 执行摘要
error?: string; // 失败时的错误信息
total?: number; // 分页总数 (paginateNodes)
page?: number; // 当前页码
pageSize?: number; // 每页数量
}7.5 节点操作
查询:
| 方法 | 说明 |
|------|------|
| findNodes<T>(label, where?, params?, options?) | 按标签查找节点 |
| findNode<T>(label, where, params?, options?) | 查询单个节点(含 .item) |
| findNodeById<T>(label, id, idField?, options?) | 按 ID 查询 |
| paginateNodes<T>(label, page?, pageSize?, where?, params?, orderBy?, options?) | 分页查询(含 .total/.page/.pageSize) |
| countNodes(label, where?, params?) | 统计节点数量,返回 number |
写入:
| 方法 | 说明 |
|------|------|
| createNode(label, props?, options?) | 创建节点(触发 OnBeforeCreateNode;options 见下方 Neo4jOptions) |
| createNodes(label, propsArray, options?) | 批量创建节点(UNWIND,逐元素触发 OnBeforeCreateNode;options 见下方 Neo4jOptions) |
| mergeNode(label, matchProps, setProps?) | 合并节点(不存在则创建,存在则跳过) |
| updateNode(label, props, where, params?, options?) | 条件更新节点(触发 OnBeforeUpdateNode;options 见下方 Neo4jOptions) |
| deleteNode(label, where, params?) | 条件删除节点(DETACH DELETE,触发 OnBeforeDeleteNode) |
写操作选项 Neo4jOptions:
interface Neo4jOptions {
embedding?: Neo4jEmbeddingOptions; // 将检索文本向量化写入节点
returnFields?: string[]; // RETURN 仅返回这些字段,省略则返回整个节点
timestamps?: boolean; // 自动填充 created_at/updated_at(毫秒)
}
await server.neo4j.createNode("Doc", { category: "docs" }, {
timestamps: true, // 自动填 created_at / updated_at
returnFields: ["id", "category"], // 仅返回这两个字段
});
timestamps批量创建时所有节点共享同一时间戳;embedding批量创建时所有节点使用同一检索文本,需逐条不同向量请循环createNode。mergeNode的钩子修改只作用于非匹配字段(matchProps匹配键保持不变)。
所有节点方法会先触发 OnCheck 访问校验;查询类方法(findNodes/findNode/findNodeById/paginateNodes)还会触发 OnAfterSelect 钩子。embedding 剥离统一在底层 query() 执行:所有经由 query() 的方法(含写方法返回的节点、findNeighbors / findRelationships / findShortestPath)默认剥离 embedding 字段;查询类方法传入 options.includeEmbedding = true 可保留(vectorSearch 自身也会默认剥离),其余方法用 query(cypher, params, { keepEmbedding: true }) 保留;向量字段非 embedding 时通过 options.embeddingField 指定;若只需部分字段,通过 options.returnFields 在 Cypher 层投影(见下)。
查询字段投影
returnFields:findNodes/findNode/findNodeById/paginateNodes支持options.returnFields,在 Cypher 层只RETURN白名单字段(如["id", "name"]),embedding 向量不进入传输链路——比事后stripEmbedding剥离省掉每条 ~3KB 的序列化+传输+拷贝。列表/全量拉取等非向量场景应优先使用。const list = await server.neo4j.findNodes("Tool", undefined, undefined, { returnFields: ["id", "name", "status"], // 只拉这三个字段,不带 embedding });
- 默认(不传
returnFields)行为不变:返回整节点、剥离 embedding、含__labels。- 投影后返回对象只含白名单字段,不含
__labels(标量投影不走toPlain(node))。- 单字段时返回标量而非对象:
returnFields: ["name"]得到string[],需对象形态请传 ≥2 个字段。
数值类型注意:neo4j-driver v6 会把 JS
number序列化成 Float(20→20.0)。Neo4j 对SKIP/LIMIT、向量检索k这类「必须为整数」的参数会拒绝 Float。框架已在内部用neo4j.int()处理——paginateNodes的 SKIP/LIMIT、vectorSearch的 k、findNodeById的整数 id 均自动转换,直接传 JS number 即可。但写入节点的整数属性不会被自动转换——
createNode("User", { id: 123 })会以 Float123.0存储。需要 Integer 属性(整数比较、索引、排序)时请显式neo4j.int():import neo4j from "neo4j-driver"; await server.neo4j.createNode("User", { id: neo4j.int(123), name: "alice" }); await server.neo4j.findNodeById("User", 123); // 读侧自动 int(),命中上面写入的 Integer同理,
findNodes/findNode自定义where里若用$param匹配整数属性,参数也应传neo4j.int(...)。
7.6 关系操作
// NodeRef:节点引用(关系查询中,from 侧变量为 a,to 侧变量为 b)
interface NodeRef {
label: string; // 节点标签
where?: string; // 匹配条件,如 "a.id = $id"(from)或 "b.id = $id"(to)
params?: Record<string, any>; // 匹配参数
}| 方法 | 说明 |
|------|------|
| createRelationship(from, type, to, props?) | 创建关系 (a)-[r:type]->(b) |
| findRelationships<T>(type, from?, to?) | 按类型查询关系 |
| findNeighbors<T>(label, where, options?, params?) | 查询邻居节点(支持方向/深度/类型过滤) |
| findShortestPath(from, to, relType?) | 查询两点间最短路径,返回 { nodes, relationships, length } |
| deleteRelationship(type, from?, to?) | 删除关系 |
关系操作会先对涉及的
from/to标签(findNeighbors还会校验targetLabel)触发OnCheck访问校验。
deleteRelationship 的 from/to 语义:
from/to均为可选,决定删除范围——
- 都传:只删
a→b这一条关系;- 只传
from:删a的所有该类型出边;- 只传
to:删指向b的所有该类型入边;- 都不传:删全图该类型所有关系(破坏性,慎用)。
const from = { label: "wiki", where: "a.name = $fromName", params: { fromName: "A" } }; const to = { label: "wiki", where: "b.name = $toName", params: { toName: "B" } }; await server.neo4j.deleteRelationship("DEPEND", from, to); // 只删 A→B注意
from/to的where分别以变量a/b书写,且两边的params键名需互不相同(同名会被后者覆盖)。
const from = { label: "User", where: "a.id = $fromId", params: { fromId: 1 } };
const to = { label: "Group", where: "b.id = $toId", params: { toId: 100 } };
await server.neo4j.createRelationship(from, "MEMBER_OF", to, { since: "2026-01-01" });
// 查询 User 的 MEMBER_OF 邻居(二度、双向、目标为 Group)
await server.neo4j.findNeighbors("User", "n.id = $id", {
relType: "MEMBER_OF",
direction: "BOTH",
depth: 2,
targetLabel: "Group",
}, { id: 1 });7.7 事务
// 写事务(自动提交/回滚)
const result = await server.neo4j.transaction(async (tx) => {
await tx.run("CREATE (n:User $props)", { props: { id: 1, name: "alice" } });
return { msg: "创建成功" };
});
// 读事务
const data = await server.neo4j.readTransaction(async (tx) => {
const res = await tx.run("MATCH (n:User) RETURN n");
return res.records;
});7.8 图结构管理
// 创建索引
await server.neo4j.createIndex("User", "email");
// 创建唯一约束
await server.neo4j.createUniqueConstraint("User", "email");索引/约束语法要求 Neo4j 5.x 及以上。
7.9 向量检索(Neo4j 5.13+)
Neo4j 只负责「存储 + 相似度检索」。向量可用外部模型(OpenAI text-embedding、本地 bge 等)生成,也可用框架内置的 EmbeddingUtility(基于 transformers.js 本地 ONNX 推理,默认模型 Xenova/jina-embeddings-v2-base-zh)。
文本向量化:
// 纯向量化(模型在 Neo4j 模块启动时已预加载,此处直接调用)
const vec = await server.neo4j.vectorizeText("Knight-Server 使用手册");
// vec: number[](float 数组)向量化 + 创建节点(通过统一的 embedding 选项一步到位):
await server.neo4j.createNode("Doc", { category: "docs", text: "Knight-Server 使用手册" }, {
embedding: {
text: "Knight-Server 使用手册", // 调用方组合好的检索文本(仅用于向量化,不写入节点)
embeddingField: "embedding", // 向量字段名,默认 "embedding"
},
});
// 等价于创建节点 { category: "docs", text: "...", embedding: [...] }更新节点重算向量:
await server.neo4j.updateNode("Doc", { text: "新文本" }, "n.id = $id", { id: 1 }, {
embedding: { text: "新文本", embeddingField: "embedding" },
});检索文本只来自
options.embedding.text,不落库;真正的业务字段由调用方放进props照常存储。
建向量索引:
await server.neo4j.createVectorIndex("Doc", "embedding", 768, "cosine");
// 参数:label, 字段名, 向量维度, 相似度函数(cosine/euclidean,默认 cosine)向量维度必须与模型输出一致(
jina-embeddings-v2-base-zh为 768 维)。
写入向量(手动指定向量时,向量就是普通属性):
await server.neo4j.createNode("Doc", {
text: "Knight-Server 使用手册",
embedding: [0.0123, -0.0456 /* ... 共 768 维 */],
});相似度检索:
const result = await server.neo4j.vectorSearch("Doc", "embedding", queryVector, 10);
// 结果节点附带 __score 相似度分数
for (const doc of result.data) {
console.log(doc.__score, doc.text);
}
// 带过滤条件(在相似结果上再 WHERE,过滤不影响检索的 k 值)
const filtered = await server.neo4j.vectorSearch(
"Doc", "embedding", queryVector, 10,
"node.status = $status",
{ status: "published" },
);
// 相似度底线 minScore:score 低于该值的命中被剔除,实际返回数可能少于 k
const qualified = await server.neo4j.vectorSearch(
"Doc", "embedding", queryVector, 10, undefined, {},
{ minScore: 0.6 },
);
// where 与 minScore 可同时传,取交集
const strict = await server.neo4j.vectorSearch(
"Doc", "embedding", queryVector, 10,
"node.status = $status", { status: "published" },
{ minScore: 0.6 },
);
- 向量以 float 数组属性存储;
queryVector需为 float 数组,索引名自动为vector_{label}_{field}。- 结果默认剥离
embedding字段,仅保留__score相似度分数;传入options.includeEmbedding = true可保留完整 embedding。k是候选池(ANN 检索量),minScore是相似度底线,实际返回数 = 两者共同过滤后的数量,可少于k。- 若使用 Neo4j 6.x 原生
Vector类型,可通过run()传入vector(new Float32Array(...))自行处理。
图扩散组合(关系感知召回由业务层组合两个核心原语完成,框架不做业务组装):
// 业务层:关系感知召回 = 向量检索 + 图扩散(findNeighbors 即图扩散原语)
const hits = await server.neo4j.vectorSearch("Tool", "embedding", queryVector, 20, undefined, {}, { minScore: 0.6 });
const ids = hits.data.map(h => h.id); // 命中节点 id(默认 idField = "id")
const related = await server.neo4j.findNeighbors(
"Tool", "id(n) IN $ids",
{ relType: "DEPEND", direction: "OUTGOING" },
{ ids }
);
// 组装逻辑(合并/去重/加权)由业务层自行处理7.10 原生 Cypher 与底层驱动
// 执行原生 Cypher(返回 records/summary)
const raw = await server.neo4j.run("MATCH (n) RETURN n LIMIT 10");
// 获取底层 neo4j-driver 驱动
const driver = server.neo4j.getNativeDriver();安全提示:所有方法的
label、field、where等均为原始 Cypher 拼接(where本身就是 Cypher 片段)。这些值必须来自可信常量,禁止直接拼接用户输入,否则存在 Cypher 注入风险。用户数据一律通过params参数化传入。
八、Redis 模块
8.1 访问方式
await server.redis.set("key", "value");
const val = await server.redis.get("key");8.2 心跳保活与断线重连
模块启动后自动对主客户端和订阅客户端定时发送 PING 保活,防止 Redis 服务端因空闲超时(timeout)断开连接。心跳间隔由 RedisConfig.heartbeatInterval 控制(默认 4 分钟)。连接意外断开重连后,自动恢复所有频道和模式订阅,心跳也会自动恢复。
断线重连由底层 node-redis 驱动自动完成,框架不重复实现重连,只负责重连后的状态同步(订阅确认、心跳恢复)。isConnected() 直接取驱动两个客户端的实时就绪状态,不维护本地标记,因此不会出现状态漂移。
默认重连策略:指数退避(50ms 起、2000ms 封顶,带随机抖动),达到 reconnectMaxRetries(默认 10,累计约 10 秒)后放弃。
连接看门狗(自动恢复):驱动放弃重连后不会自行恢复,框架每隔 recoveryInterval(默认 10 秒)检查一次连接状态,发现掉线就自动调用 reconnect()(重连两个客户端并补订阅)。因此下面两种场景都能自愈,无需人工介入:
| 场景 | 过程 | |---|---| | 启动时 Redis 未就绪 | 驱动重试约 10s 后放弃 → 服务照常启动(路由已注册)→ 看门狗接手,Redis 恢复后自动连上并补订阅 | | 运行中 Redis 当机 > 10s | 驱动重试耗尽后放弃 → 看门狗接手定期重试 → Redis 恢复后自动恢复连接与订阅 |
连接正常时看门狗只是空转检查,开销可忽略(可用
recoveryInterval调整频率)。若希望完全不依赖看门狗,可把reconnectMaxRetries调大或传入永不放弃的reconnectStrategy:
// 方式一:加大重试次数(约 1 分钟)
redis: { host: "127.0.0.1", port: 6379, reconnectMaxRetries: 60 }
// 方式二:自定义策略 —— 永不放弃(返回延迟毫秒;返回 false 表示放弃)
redis: {
host: "127.0.0.1", port: 6379,
reconnectStrategy: (retries) => Math.min(2 ** retries * 50, 5000),
}
// 方式三:不重连(失败即放弃,由调用方自行处理)
redis: { host: "127.0.0.1", port: 6379, reconnectStrategy: false }启动行为:OnStart 会等待首次连接(阻塞至多约 10 秒);连接失败时只记录日志、不中断启动,路由注册与服务其余模块照常启动,随后由看门狗自动恢复(也可手动 reconnect())。
await server.OnStart(config); // Redis 未启动也不会卡死,约 10s 后继续
if (!server.redis.isConnected()) { /* 降级处理:此刻尚未连上,看门狗会在后台自动恢复 */ }
await server.redis.reconnect(); // 也可立即手动恢复(成功与否返回 boolean)订阅恢复是幂等的:驱动在重连时已自动重放订阅,框架使用稳定的分发器引用再次确认,因此不会出现重复投递(同一条消息不会被处理多次)。
重连后清空缓存(flushOnReconnect,默认关):断连期间「写库成功、但删缓存没送达」的条目,应用侧无论怎么记账都只能等连接回来才处理。打开这个开关后,框架在每次恢复到就绪时执行一次 flushDb(),把整库缓存一次性作废,读路径按需重建即可 —— 应用侧不必再维护自己的「待校准」状态。
- 触发点是恢复沿,不是断连沿:断连那一刻命令发不出去,FLUSHDB 同样发不出去。
- 挂在主客户端的
ready上:驱动自行重连成功(最常见)与看门狗reconnect()成功都经过这一个事件,所以短到驱动自己就恢复的抖动也会清,不会漏。 - 全新启动且一次就绪时不触发:标记只在「连接中断」或「连接尝试失败」之后置位。注意「启动时 Redis 就不可用、等它恢复」也算恢复沿,会清 —— 这段时间里的写同样没能送达删缓存,本就该清。
- 失败会重试:清空失败会记为「待清空」,看门狗按
recoveryInterval(默认 10s)再试直到成功;reconnect()的返回值仍只表示「连接是否就绪」,与清空成败无关。 - 不开这个开关时的行为:框架只看连接状态,缓存一致性由应用侧自己负责(例如自行记录「哪次删除没送达」并在连接恢复后清理)。
- ⚠️ 前提是该
db由本应用独占:flushDb删的是整个 db(含其他应用写入的键)。框架自身不往 Redis 写键,但同一个 Redis 实例的同一个db上若还有别的应用,请为它们指定独立的db。
8.3 库级操作
| 方法 | Redis 命令 | 返回值 |
|------|------------|--------|
| flushDb(mode?) | FLUSHDB | Promise<boolean> |
mode 可选值为 "ASYNC"(异步清空,不阻塞服务端,Redis 4.0+)或 "SYNC"(同步清空);省略时使用服务端默认行为。返回 reply === "OK",即是否清空成功。
await server.redis.flushDb(); // 清空当前 db 的全部键
await server.redis.flushDb("ASYNC"); // 异步清空⚠️ 不可逆操作:该命令会删除当前
db内的所有键(不仅是本框架写入的),请勿在生产环境随意调用。多库场景务必在配置中通过db指定专用库(如测试库),避免误清业务数据。
8.4 KV 操作
| 方法 | Redis 命令 | 返回值 |
|------|------------|--------|
| get(key) | GET | Promise<string \| {}> |
| set(key, value) | SET | Promise<string \| {}> |
| deleteKeys(...keys) | DEL | Promise<number> |
| expire(key, seconds) | EXPIRE | Promise<number> |
| exists(...keys) | EXISTS | Promise<number> |
| increment(key) | INCR | Promise<number> |
| incrementBy(key, n) | INCRBY | Promise<number> |
| decrement(key) | DECR | Promise<number> |
| decrementBy(key, n) | DECRBY | Promise<number> |
| timeToLive(key) | TTL | Promise<number> |
| keys(pattern) | KEYS | Promise<string[]> |
| setIfNotExists(key, value) | SETNX | Promise<number> |
| getAndSet(key, value) | GETSET | Promise<string \| {}> |
| multiGet(...keys) | MGET | Promise<(string \| {})[]> |
| multiSet(keyValues) | MSET | Promise<string \| null> |
8.5 Hash 操作
| 方法 | Redis 命令 |
|------|------------|
| hashSet(key, field, value) | HSET |
| hashGet(key, field) | HGET |
| hashGetAll(key) | HGETALL |
| hashDelete(key, ...fields) | HDEL |
| hashExists(key, field) | HEXISTS |
| hashKeys(key) | HKEYS |
| hashValues(key) | HVALS |
| hashLength(key) | HLEN |
| hashMultiSet(key, fieldValues) | HSET (批量) |
8.6 List 操作
| 方法 | Redis 命令 |
|------|------------|
| listLeftPush(key, ...elements) | LPUSH |
| listRightPush(key, ...elements) | RPUSH |
| listLeftPop(key) | LPOP |
| listRightPop(key) | RPOP |
| listRange(key, start, stop) | LRANGE |
| listLength(key) | LLEN |
| listRemove(key, count, element) | LREM |
8.7 Set 操作
| 方法 | Redis 命令 |
|------|------------|
| setAdd(key, ...members) | SADD |
| setRemove(key, ...members) | SREM |
| setMembers(key) | SMEMBERS |
| setIsMember(key, member) | SISMEMBER |
| setCardinality(key) | SCARD |
| setIntersect(...keys) | SINTER |
| setUnion(...keys) | SUNION |
| setDifference(...keys) | SDIFF |
8.8 SortedSet 操作
| 方法 | Redis 命令 |
|------|------------|
| sortedSetAdd(key, ...{score, value}[]) | ZADD |
| sortedSetRange(key, start, stop, withScores?) | ZRANGE |
| sortedSetRangeByScore(key, min, max, withScores?) | ZRANGEBYSCORE |
| sortedSetRemove(key, ...members) | ZREM |
| sortedSetCardinality(key) | ZCARD |
| sortedSetScore(key, member) | ZSCORE |
| sortedSetRank(key, member) | ZRANK |
| sortedSetReverseRank(key, member) | ZREVRANK |
8.9 Pub/Sub 操作
// 订阅频道
await server.redis.subscribe("game:events", (message, channel) => {
console.log(`[${channel}] ${message}`);
});
// 发布消息
await server.redis.publish("game:events", JSON.stringify({ type: "start" }));
// 退订
await server.redis.unsubscribe("game:events");
// 模式订阅(通配符)
await server.redis.patternSubscribe("game:*", (message, channel) => { /* ... */ });
await server.redis.patternUnsubscribe("game:*");8.10 RedisRouter 自动订阅
@OnRegister 的路径即频道名,RedisModule 启动后自动订阅,消息到达时调用 OnHandler:
@OnRegister("/channel/notifications")
export class NotificationRouter extends RedisRouter {
// 有人 publish 到 "/channel/notifications" 时自动触发
async OnHandler(channel: string, message: string): Promise<void> {
const data = JSON.parse(message);
console.log(`收到频道 ${channel} 的消息:`, data);
}
}关于
OnCheck:Redis 为 Pub/Sub 订阅模式,没有请求入口,因此RedisRouter的OnCheck基类默认直接放行(valid: true)。连接失败时模块启动会跳过该路由注册,业务侧无需在此做连接检查。
8.11 Pipeline 操作
const pipeline = server.redis.pipeline();
pipeline.set("key1", "val1")
.get("key1")
.deleteKeys("key2")
.listLeftPush("queue", "item1", "item2");
const results = await pipeline.exec(server.redis.getNativeClient());8.12 Lua 脚本
// 执行脚本
const result = await server.redis.evalScript("return redis.call('GET', KEYS[1])", ["mykey"]);
// 加载 + 缓存执行
const sha = await server.redis.scriptLoad("return redis.call('INCR', KEYS[1])");
const count = await server.redis.evalSha(sha, ["counter"]);8.13 底层客户端访问
const native = server.redis.getNativeClient(); // 主客户端
const sub = server.redis.getNativeSubscriber(); // 订阅客户端九、工具集 (Utility)
9.1 LogUtility
import { LogUtility } from "knight-server";
LogUtility.LogTip("提示信息");
LogUtility.LogInfo("信息");
LogUtility.LogWarning("警告");
LogUtility.LogError("错误");
LogUtility.LogSystem("系统信息");日志格式:[系统名称][标签][yyyy-MM-dd HH:mm:ss.SSS]🔻 消息...
类名不自动获取,由调用方自己写进消息里(发布包经过混淆,从 Error.stack 取出的方法名会退化成 e/t 这类无意义短名)。
日志落盘(可选)
默认只输出到控制台。 Server.OnStart 总会初始化一次日志:配置了 log.dir 才写文件,否则等价于走全默认值(控制台照常输出,不写文件、也不会创建任何目录)。初始化在模块启动之前完成,各模块的启动日志同样会进文件:
const config: ServerConfig = {
serverHost: "127.0.0.1",
log: { dir: "logs", timezone: "Asia/Shanghai" }, // dir 相对 process.cwd()
// log: { dir: "/var/log/knight", filePrefix: "api", dailyRotation: true, console: false },
};
await Server.OnStart(config);| 字段 | 默认 | 说明 |
|------|------|------|
| dir | 不配置 | 日志目录,绝对路径,或以进程工作目录 process.cwd() 为基准的相对路径;多级路径自动逐级创建(等价 mkdir -p)。不配置则不开文件日志 |
| filePrefix | "⚡KNIGHT" | 系统名称 / 文件名前缀,同时是每行开头的 [...],产出 ⚡KNIGHT-2026-09-20.log |
| timezone | 不配置 | 日志时区(IANA 名称,如 "Asia/Shanghai")。时间戳与按天切分的文件名都按其计算;不配置则跟随进程本地时区,详见下方「日志时区」 |
| dailyRotation | true | 按天切分文件,跨天自动换新文件;置 false 则固定写 ⚡KNIGHT.log |
| console | true | 是否同时输出到控制台。生产由 pm2/systemd 接管 stdout 时可置 false 避免重复 |
文件行用纯文本级别而非 emoji(便于 grep 与日志采集),控制台格式不变;Error 参数落盘时取完整堆栈,普通对象序列化为 JSON:
[⚡KNIGHT][ERROR][2026-09-20 16:13:53.530] Redis连接失败: Error: connect ECONNREFUSED 127.0.0.1:6379
at ...
[⚡KNIGHT][INFO][2026-09-20 16:13:53.531] 服务已启动: {"port":3000}日志时区(服务端必读):timezone 不配置时用的是「进程本地时区」,不是 UTC —— 但容器(Docker / 宝塔)默认就是 UTC,此时日志时间会比北京时间少 8 小时,而且按天切分的日界会落在北京时间早上 08:00,一个自然日的日志被劈成两个文件。
在服务器上用这条命令确认进程时区:
node -e "const d=new Date();console.log('时钟:',d.toString());console.log('TZ:',process.env.TZ??'(未设置)');console.log('时区:',Intl.DateTimeFormat().resolvedOptions().timeZone);console.log('getHours:',d.getHours())"两种解法,任选其一:
- 配置
timezone(推荐,不受部署环境影响):log: { dir: "logs", timezone: "Asia/Shanghai" }。生效后启动时会打印一条[LogUtility]🔻日志时区: Asia/Shanghai (GMT+8),可直接核对。 - 改环境让进程时区跟随机器:Docker 加
-e TZ=Asia/Shanghai、pm2 在ecosystem.config.js里设env: { TZ: 'Asia/Shanghai' }、systemd 在 unit 里写Environment=TZ=Asia/Shanghai。
时区名写错(如 Asia/Shanghia)不会中断启动:框架直出 时区配置无效,已回退为进程本地时区 并降级为本地时间。
路径为什么以 cwd 为基准? 本库的发布形态是 rollup 打出的单文件 bundle(生产构建再经混淆),运行时 __dirname 指向的是使用方放置该文件的位置 —— node_modules/knight-server/dist/(npm ci 会被清空、也可能只读),并不是「项目根目录」。process.cwd() 才是由运维侧掌控的运行时锚点:
| 部署方式 | 建议写法 |
|---|---|
| systemd | WorkingDirectory=/opt/knight + log: { dir: "logs" },或直接 dir: "/var/log/knight" |
| pm2 | cwd 字段决定 process.cwd();应用目录归 root 而服务以普通用户运行时务必用绝对路径(如 /var/log/knight),否则 mkdir 会 EACCES |
| Docker | 用绝对路径指向挂载卷,如 dir: "/app/logs",并把该目录挂载出来 |
启动时会打印一条解析后的绝对路径,便于确认日志实际落在哪里:
[⚡KNIGHT][⚙️ 系统][2026-09-20 16:13:53.530]🔻 [LogUtility]🔻日志文件: /opt/knight/logs/⚡KNIGHT-2026-09-20.log初始化失败不中断启动:目录不可写(权限不足、路径上已存在同名文件等)时,框架在控制台直出 日志目录不可用,文件日志已关闭 并降级为仅控制台,OnStart 照常继续。运行期写入失败同样只降级一次并关闭文件日志,不刷屏、不抛出。
也可以不经过 ServerConfig,直接调用 LogUtility.OnInit(config?):config 省略或未给 dir 时只输出控制台并返回 false,给了 dir 则返回是否启用成功。
LogUtility.OnInit({ dir: "/var/log/knight" }); // 返回是否启用成功注意:
Server.OnStart每次启动都会调用一次OnInit。若在OnStart之前自行初始化过,随后ServerConfig里又没有log,这次调用会以全默认值覆盖掉先前的目录配置(表现为文件日志被关闭)。要保留请把配置写进ServerConfig.log。注意:日志落盘只覆盖
LogUtility的五个方法。AbstractObject上的this.log/logInfo/logWarn/logError以及各类中直接调用的console.*不走LogUtility,不会进文件。
9.2 TokenUtility (JWT)
import { TokenUtility } from "knight-server";
// 签发
const accessToken = TokenUtility.sign({ uid: 1 }, server.config.token.secret, "2h");
const refreshToken = TokenUtility.signRefresh({ uid: 1 }, server.config.token.secret, "7d");
// 验证(抛出异常时视为无效)
const payload = TokenUtility.verify<{ uid: number }>(accessToken, server.config.token.secret);
// 解码(不验证签名)
const decoded = TokenUtility.decode<{ uid: number }>(accessToken);9.3 StringUtility
| 方法 | 说明 |
|------|------|
| empty(str) | null/undefined/"" → true |
| isBlank(str) | 空或纯空白 → true |
| trimToEmpty(str) | 安全 trim,null → "" |
| uuid() | 生成 UUID |
| random(length?) | 随机字符串,默认 16 位 |
| randomNumber(length?) | 随机数字串,默认 6 位 |
| format(template, params) | "Hello {name}" 模板替换 |
| camelToSnake(str) | 驼峰转下划线 |
| snakeToCamel(str) | 下划线转驼峰 |
| joinUrl(...parts) | URL 路径拼接 |
9.4 TimeUtility
| 方法 | 说明 |
|------|------|
| nowTimestamp() | Date.now() |
| nowUnix() | 秒级时间戳 |
| formatDate(date?) | yyyy-MM-dd |
| formatDateTime(date?) | yyyy-MM-dd HH:mm:ss |
| sleep(ms) | 异步等待 |
| startOfDay(date?) | 当天 00:00:00 |
| endOfDay(date?) | 当天 23:59:59 |
| addDays(date?, n) | 加/减天数 |
| diffInDays(d1, d2) | 日期差 |
| isExpired(date) | 是否已过期 |
| timer(fn) | 测量异步执行耗时 |
9.5 NetworkUtility
import { NetworkUtility } from "knight-server";
// HTTP 请求
const res = await NetworkUtility.HTTP.OnRequest<{ data: any }>("https://api.example.com", {
method: "post",
data: { key: "value" },
params: { page: 1 },
timeout: 5000,
});
// SSE 流式请求
for await (const chunk of NetworkUtility.HTTP.OnStreamRequest("https://api.example.com/stream")) {
console.log(chunk);
}
// 客户端 WebSocket
const ws = new NetworkUtility.Websocket();
ws.On("message", (event) => console.log(event.data));
ws.OnConnect("ws://localhost:6666/ws/chat?token=xxx");
ws.OnSend("hello");
ws.OnDisConnect();常用枚举:
NetworkUtility.WebSocketCloseCode.CONFLICT // 4000
NetworkUtility.WebSocketCloseCode.AUTH_FAILED // 4001
NetworkUtility.ResponseCode.UNAUTHORIZED // 401
NetworkUtility.ResponseCode.TOO_MANY_REQUESTS // 429
NetworkUtility.ResponseCode.BUSINESS_ERROR // 10009.6 EmbeddingUtility (文本向量化)
基于 @huggingface/transformers(本地 ONNX 推理,无需外部 API),提供文本向量化能力,供 Neo4j 向量检索使用。
import { EmbeddingUtility } from "knight-server";
// 加载模型(幂等;默认模型 Xenova/jina-embeddings-v2-base-zh,默认镜像 https://hf-mirror.com)
await EmbeddingUtility.OnLoad();
// 自定义模型与镜像
await EmbeddingUtility.OnLoad("Xenova/bge-small-zh-v1.5", "https://hf-mirror.com");
// 文本向量化,返回 Float32Array
const vector = await EmbeddingUtility.OnEmbed("你好,世界");
- 首次加载会下载模型权重,体积较大;
OnLoad幂等,重复调用不会重复加载。@huggingface/transformers及其 onnxruntime 等依赖较重,Neo4j 模块启动时会预加载一次。- 更多用法见「7.9 向量检索」中的
vectorizeText与createNode/updateNode的embedding选项。
9.7 CryptoUtility (加密/解密)
基于 Node 原生 crypto 模块,提供密码哈希、对称加密、消息认证与摘要。
import { CryptoUtility } from "knight-server";
// 密码哈希与验证
const hashed = CryptoUtility.OnHashPassword("123456");
const ok = CryptoUtility.OnVerifyPassword("123456", hashed); // true
// AES-256-GCM 对称加密(key 经 SHA-256 派生为 32 字节)
const encrypted = CryptoUtility.OnEncrypt("敏感数据", "secret-key");
const decrypted = CryptoUtility.OnDecrypt(encrypted, "secret-key");
// HMAC-SHA256 签名与验签
const sign = CryptoUtility.OnSign("payload", "secret");
const valid = CryptoUtility.OnVerifySign("payload", "secret", sign); // true
// 摘要
const sha256 = CryptoUtility.OnSha256("hello");
const md5 = CryptoUtility.OnMd5("hello"); // 仅用于非安全场景| 方法 | 说明 |
|------|------|
| OnHashPassword(password) | scrypt 密码哈希,返回 salt:hash |
| OnVerifyPassword(password, stored) | 验证密码,匹配返回 true |
| OnEncrypt(plainText, key) | AES-256-GCM 加密,返回 iv:authTag:cipherText(hex) |
| OnDecrypt(cipherText, key) | AES-256-GCM 解密,key 错误/密文被篡改时抛异常 |
| OnSign(data, secret) | HMAC-SHA256 签名(hex) |
| OnVerifySign(data, secret, signature) | HMAC-SHA256 验签 |
| OnSha256(data) | SHA-256 摘要(hex) |
| OnMd5(data) | MD5 摘要(hex,仅非安全场景) |
OnHashPassword使用 scrypt(Node 原生,抗 GPU 暴力破解),每次生成随机盐。- 验签/验密均使用
timingSafeEqual防时序攻击。
十、CheckResult 与 CheckPayload
interface CheckResult<T extends CheckPayload = CheckPayload> {
valid: boolean;
error?: string; // valid=false 时的错误信息
payload?: T; // 携带的业务数据(用户信息等)
}
interface CheckPayload { }在 OnCheck 中返回 payload,后续可通过 checkResult.payload 获取,避免重复查询。
扩展 Payload 类型:
interface ChatPayload extends CheckPayload {
userId: string;
roomId: string;
}
// 在 OnCheck 中
return { valid: true, payload: { userId: "xxx", roomId: "room1" } as ChatPayload };十一、完整示例
import {
Server,
HttpRouter, WebSocketRouter, WebSocketClient,
MySQLRouter, Neo4jRouter, RedisRouter,
OnRegister, CheckResult, CheckPayload,
TokenUtility, TimeUtility, LogUtility,
} from "knight-server";
// ────────── HTTP 路由 ──────────
@OnRegister("/api/user/info")
class UserInfoRouter extends HttpRouter {
async OnCheck(server, request, response): Promise<CheckResult> {
const token = request.headers.authorization?.replace("Bearer ", "");
if (!token) return { valid: false, error: "请先登录" };
try {
const payload = TokenUtility.verify(token, server.config.token.secret);
return { valid: true, payload };
} catch {
return { valid: false, error: "令牌无效" };
}
}
async OnGetHandler(server, checkResult, request, response) {
const uid = (checkResult.payload as any).uid;
const result = await server.mysql.findOne("users", "id = ?", [uid]);
if (!result.item) {
response.status(404).json({ success: false, error: "用户不存在" });
return;
}
const { password, ...safe } = result.item;
response.json({ success: true, data: safe });
}
}
// ────────── WebSocket 路由 ──────────
@OnRegister("/ws/chat")
class ChatRouter extends WebSocketRouter {
async OnCheck(server, url): Promise<CheckResult> {
const token = url.searchParams.get("token");
if (!token) return { valid: false, error: "缺少令牌" };
try {
const payload = TokenUtility.verify(token, server.config.token.secret);
return { valid: true, payload };
} catch {
return { valid: false, error: "令牌无效" };
}
}
async OnConnect(server, checkResult, client) {
const uid = (checkResult.payload as any).uid;
await server.mysql.update("users", { online: 1 }, "id = ?", [uid]);
client.OnBroadcast(JSON.stringify({ type: "join", uid }));
}
async OnDisconnect(server, checkResult, client) {
const uid = (checkResult.payload as any).uid;
await server.mysql.update("users", { online: 0, last_seen: TimeUtility.formatDateTime() }, "id = ?", [uid]);
client.OnBroadcast(JSON.stringify({ type: "leave", uid }));
}
async OnMessage(server, checkResult, client, data) {
client.OnUpdateHeartbeat();
client.OnBroadcast(data, true);
}
}
// ────────── DB 路由 (钩子) ──────────
@OnRegister("/users")
class UsersRouter extends MySQLRouter {
async OnBeforeInsert(data) {
data.created_at = TimeUtility.formatDateTime();
return { proceed: true, data };
}
async OnAfterSelect(rows) {
return rows.map(({ password, ...rest }) => rest);
}
}
// ────────── Neo4j 路由 (钩子) ──────────
@OnRegister("/user")
class UserGraphRouter extends Neo4jRouter {
async OnBeforeCreateNode(label, props) {
props.created_at = TimeUtility.formatDateTime();
return { proceed: true, data: props };
}
}
// ────────── Redis 路由 (Pub/Sub) ──────────
@OnRegister("/channel/notifications")
class NotificationRouter extends RedisRouter {
async OnHandler(channel, message) {
LogUtility.LogInfo(`[${channel}] 收到消息:`, message);
}
}
// ────────── 启动 ──────────
(async () => {
const server = Server.Instance;
await server.OnStart({
mysql: {
host: "localhost",
port: 3306,
user: "root",
password: "your_password",
database: "myapp",
connectionLimit: 10,
tableRelativePath: "./src/database/tables",
},
neo4j: {
uri: "bolt://localhost:7687",
username: "neo4j",
password: "your_password",
database: "neo4j",
},
redis: {
host: "localhost",
port: 6379,
},
http: {
port: 3000,
},
websocket: {
port: 6666,
heartbeat_timeout_interval: 30000,
router_timeout_interval: 60000,
},
token: {
secret: "your-jwt-secret-key",
accessExpiresIn: "2h",
},
log: {
dir: "logs", // 相对 process.cwd();部署到 Linux 建议用绝对路径,如 "/var/log/knight"
},
serverHost: "127.0.0.1",
});
console.log("[KNIGHT] 服务启动完成");
})();十二、设计要点
- 所有业务 Router 必须用
@OnRegister装饰,否则 Module 无法发现和注册它 - Router 归属由
instanceof判断,继承正确的基类(HttpRouter / WebSocketRouter / MySQLRouter / Neo4jRouter / RedisRouter)即可 - WebSocket 心跳由协议层 Ping/Pong + 业务层共同维持:客户端发送 Ping 帧时服务端自动回复 Pong 并刷新心跳;在
OnMessage中调用client.OnUpdateHeartbeat()也是备选方案。超时客户端会被自动断开 - MySQL / Neo4j 的钩子可以阻止或改写数据,返回
{ proceed: false }可拦截操作 server.config在所有 Router 中都可通过this.server访问(继承自AbstractObject)- Express 中间件已内置:CORS、JSON body 解析(50MB)、URL-encoded 解析(50MB)
