Commit 6fe7cedf by Xianquan Committed by GitHub

refactor(redis): consolidate Redis runtime and BullMQ services in DAL (#7393)

* refactor(redis): establish runtime and keyspace

* refactor(redis): add capabilities and script registry

* refactor(redis): centralize phase 3a stores

* refactor(redis): migrate wechat qr login store

* refactor(redis): move Redis access into DAL

* refactor(redis): migrate lease repository

* refactor(redis): migrate stop signal and stream resume repositories

* refactor(redis): complete DAL migration

* test(redis): complete DAL Redis test coverage

* chore(redis): align pro lockfile importer

* test(redis): update service and app redis mocks

* chore(redis): align pro migration submodule

* refactor(redis): consolidate cache migration

* refactor(redis): convert runtime to class

* chore(deps): update redis and tooling dependencies

* refactor(redis): move BullMQ services into DAL

* fix(redis): harden DAL checks and lifecycle

* chore: fix bullmq to 5.52.2

* chore: bump deps

* test(redis): mock captured BullMQ binding

* lock

* docs: consolidate DAL and WeChat designs

* chore: remove DAL review artifacts

---------

Co-authored-by: Archer <545436317@qq.com>
parent b7b78be7
# SigNoz Redis / BullMQ 告警手册
> 本文面向 SigNoz 管理员,记录可直接配置的 Redis/BullMQ 日志告警、阈值和降噪规则。
> DAL 架构、失败合同和发布约束见 [FastGPT Data Access Layer 设计](../dal/data-access-layer.md)。
## 1. 使用前提
### 1.1 Service name
示例统一使用:
```text
service.name = 'fastgpt-cn-client'
```
实际值由 `OTEL_SERVICE_NAME` 决定。App 与 Pro 使用不同 service name 时分别创建规则,不要为了复用查询合并两个进程。
### 1.2 OTel 字段
| SigNoz 字段 | 用途 |
| --- | --- |
| `service.name` | 应用/进程边界 |
| `severity_text` | `error``warn``info``debug` |
| `body.__log_message` | 稳定日志消息,告警的主要匹配字段 |
| `body.*` | `role``name` 等低基数 metadata |
| `category` | LogTape category,目前不同 Cache/Runtime 尚未统一 |
`error`、完整 Redis URL、teamId、userId、sessionId、shareId、jobId、key 等字段只用于排查,不得用于通知分组或模板。
Runtime 与 BullMQ logger 当前通常继承 `['system']`,业务 Cache 使用各自领域 category。category 统一前,以下规则默认不增加 category 条件,否则会漏掉 Runtime、BullMQ 和多数 Cache 日志。
## 2. 通用配置
| 参数 | 默认值 |
| --- | --- |
| Evaluation window | Rolling 5 分钟 |
| Evaluation frequency | 1 分钟 |
| No-data | No alert / OK |
| Match type | `at least once` |
| Repeat notification | P1 15 分钟;P2 60 分钟 |
| 推荐 Group By | `body.__log_message``body.role``body.name``service.name` |
发布、Redis 维护和故障演练使用 SigNoz maintenance window 或通知路由静默,不修改 Filter 绕过一次发布。
## 3. 告警规则
### R1. Redis 连接错误
- 严重度:P1
- Filter:
```text
service.name = 'fastgpt-cn-client'
AND severity_text = 'error'
AND body.__log_message = 'Redis connection error'
```
- Aggregate:`count()`;Group By:`body.role``service.name`
- 阈值:同一 role 5 分钟内 `>= 3`
- 处理:检查 Redis 实例、网络、认证和连接数,并关联 R3/R4。
### R2. Redis 进程关闭失败
- 严重度:P1
- Filter:
```text
service.name = 'fastgpt-cn-client'
AND severity_text = 'error'
AND body.__log_message = 'Redis runtime close failed after process shutdown signal'
```
- Aggregate:`count()`;Group By:`service.name`
- 阈值:5 分钟内 `>= 1`
- 处理:关联部署信号、进程退出码、Next.js drain 和 BullMQ close 日志。
### R3. Redis 重连抖动
- 严重度:P2
- Filter:
```text
service.name = 'fastgpt-cn-client'
AND severity_text = 'warn'
AND (
body.__log_message = 'Redis connection closed'
OR body.__log_message = 'Redis reconnect scheduled'
OR body.__log_message = 'Redis reconnect requested by command error'
)
```
- Aggregate:`count()`;Group By:`body.role``service.name`
- 阈值:同一 role 5 分钟内 `>= 5`;单条 close 不告警。
- 处理:R1 同时触发时按连接故障处理,否则检查 failover、网络抖动和连接压力。
### R4. BullMQ Queue/Worker 错误
- 严重度:P1
- Filter:
```text
service.name = 'fastgpt-cn-client'
AND severity_text = 'error'
AND (
body.__log_message = 'BullMQ queue error'
OR body.__log_message = 'BullMQ worker error'
OR body.__log_message = 'BullMQ worker restart failed, will retry'
)
```
- Aggregate:`count()`;Group By:`body.name``service.name`
- 阈值:同一 queue/worker 5 分钟内 `>= 3`
- 处理:先排除 R1;Redis 正常时检查 processor、job payload 和 BullMQ 版本。
### R5. BullMQ 生命周期异常
- 严重度:P2;非部署窗口可提升到 P1。
- Filter 消息:
```text
BullMQ worker closed, attempting restart
BullMQ worker resume failed
BullMQ queue close failed
BullMQ worker close failed
BullMQ queue connection forced disconnect failed
BullMQ worker connection forced disconnect failed
Failed to release Redis connection after BullMQ queue creation error
Failed to release Redis connection after BullMQ worker creation error
```
- Filter 叠加 `severity_text = 'warn'`,用 OR 匹配上述 `body.__log_message`
- Aggregate:`count()`;Group By:`body.name``body.__log_message`
- 阈值:5 分钟内 `>= 3`;close/forced-disconnect failure 可单独设置 `>= 1`
- 降噪:部署期间静默;单条 worker closed 不代表故障。
### R6. Redis Cache 持续降级
- 严重度:P2
- 关注消息:
```text
Daily active dedupe failed open
DingTalk accessToken cache read failed
DingTalk accessToken cache write failed
Invalid Redis session record
Failed to delete invalid Redis session record
Redis lease renew failed
Redis lease acquire failed
Redis lease release failed
Failed to get team point cache
Failed to set team point cache
Failed to increment team point cache
Failed to clear team point cache
Workflow stop signal read failed open
Workflow stop signal clear failed
Failed to get team vector count cache
Failed to set team vector count cache
Failed to invalidate team vector count cache
```
- Filter 叠加 `severity_text IN ('error', 'warn')`,用 OR 匹配上述消息。
- Aggregate:`count()`;Group By:`body.__log_message``service.name`
- 阈值:5 分钟内 `>= 10`;session 损坏或 lease acquire failure 可设 `>= 3`
- 处理:R1 未触发时,检查单一 Cache 的数据格式、TTL 和降级路径。
`Skipped invalid team point cache value` 属于数据质量信号,不并入 Redis 可用性告警。
### R7. Stream Resume 镜像失败
- 严重度:P1;内存压力状态变化为 P2。
- P1 消息:
```text
Failed to clear stream resume redis keys before mirror
Failed to mirror stream response to redis
Failed to shrink stream resume redis ttl
```
- 阈值:合计 5 分钟内 `>= 5`;恢复能力属于核心 SLA 时可低频通知单次失败。
- P2 消息:
```text
Disabling new stream resume mirrors due to Redis memory pressure
Failed to inspect Redis memory pressure for stream resume mirror
Failed to persist stream resume unavailable state
cleanStaleGeneratingChats: failed to inspect stream resume activity
```
- 阈值:内存压力状态变化 `>= 1/10 分钟`;检查或状态写入失败 `>= 3/5 分钟`
### R8. 限流服务 Fail-closed
- 严重度:P1,因为错误会直接转换为 429。
- Filter:
```text
service.name = 'fastgpt-cn-client'
AND severity_text = 'error'
AND (
body.__log_message = 'Fixed window rate limit failed closed'
OR body.__log_message = 'Team QPM configuration lookup failed closed'
OR body.__log_message = 'Team QPM rate limit failed closed'
)
```
- Aggregate:`count()`;Group By:`body.__log_message``body.type`
- 阈值:5 分钟内 `>= 1`;高流量环境可使用 `>= 3`
- 禁止按 teamId 或 key 分组。
### R9. Pro 企微订单 Cache 故障
- 严重度:P1,由支付错误率决定是否降为 P2。
- Filter:
```text
service.name = 'fastgpt-cn-client'
AND severity_text = 'error'
AND (
body.__log_message = 'Failed to get WeChat Work pending order cache'
OR body.__log_message = 'Failed to set WeChat Work pending order cache'
OR body.__log_message = 'Failed to clean WeChat Work pending order cache'
)
```
- Aggregate:`count()`;Group By:`body.__log_message``service.name`
- 阈值:5 分钟内 `>= 1`,只关注持续故障时使用 `>= 3`
- 禁止按 teamId 或 orderId 分组。
### R10. BullMQ 业务任务失败
这类错误不能单独证明 Redis 不可用,应与 R1/R4 分开配置:
```text
Failed to push collection update job
Failed to update collection
Failed to check evaluation job status
Failed to remove evaluation job
Failed to enqueue skill creation job, marked skill creation failed
Failed to resume pending skill creation jobs
Failed to resume marked skill delete jobs
Reply job failed
Wechat getUpdates request failed
getUpdates API error
Schedule next poll (completed) failed
Schedule next poll (failed) failed
```
- 严重度和阈值按业务 SLA,默认 `>= 3/5 分钟`
- Group By 使用稳定消息或业务 category;shareId 只能做去重计数,不能形成通知组。
- R1/R4 未触发时,按 processor、第三方 API 或 job payload 排查。
## 4. Dashboard 与告警边界
以下信息只用于 dashboard 或趋势,不直接通知:
```text
Redis connection established
Redis connection ready
Redis runtime closed after process shutdown signal
BullMQ worker ready
BullMQ worker restarted successfully
BullMQ worker paused
Collection update job pushed
Dataset sync scheduler reconcile finished
Redis memory pressure recovered; stream resume mirror creation resumed
```
幂等补偿、资源已不存在和 active job 无法删除等 warning 也属于业务 dashboard。只有它们与 R1/R4 或明确业务 SLA 同时满足时才升级。
## 5. Category 统一约束
基础设施 category 统一后,Redis Runtime/health/shutdown 使用 `LogCategories.INFRA.REDIS`,BullMQ Runtime 使用 `LogCategories.INFRA.QUEUE`;业务 Cache 保留领域 category,并额外提供低基数 `component` 字段。
在 SigNoz 中确认 category 的实际存储类型前,不批量修改现有规则。日志消息重命名必须同步更新对应规则;本文不维护一份与源码重复的全量消息清单。
......@@ -45,9 +45,12 @@
"typescript-eslint": "catalog:lint",
"vitest": "catalog:"
},
"engines": {
"node": ">=20.19.0",
"pnpm": "10.x"
"devEngines": {
"runtime": {
"name": "node",
"version": "24.18.1",
"onFail": "download"
}
},
"packageManager": "pnpm@10.33.4"
}
{
"name": "@fastgpt/dal",
"version": "1.0.0",
"type": "module",
"exports": {
"./redis": "./redis/index.ts",
"./redis/adapter": "./redis/adapter.ts",
"./redis/bullmq": "./redis/bullmq/index.ts",
"./redis/bullmq/services/*": "./redis/bullmq/services/*.ts",
"./redis/caches": "./redis/caches/index.ts",
"./redis/types": "./redis/types.ts",
"./redis/runtime": "./redis/runtime/index.ts"
},
"scripts": {
"test": "vitest run -c vitest.config.ts",
"test:watch": "vitest -c vitest.config.ts",
"test:integration:redis": "vitest run -c vitest.config.ts test/integrations/redis",
"typecheck": "tsc --noEmit --pretty"
},
"engines": {
"node": ">=20.19.0",
"pnpm": "10.x"
},
"dependencies": {
"@fastgpt/global": "workspace:*",
"bullmq": "5.52.2",
"ioredis": "^5.11.1",
"zod": "catalog:"
},
"devDependencies": {
"@types/node": "catalog:",
"typescript": "catalog:",
"vitest": "catalog:"
}
}
import { getRedisRuntime } from '../runtime';
import { getRedisBullMQRuntime } from './context';
import type { QueueNames } from './names';
import type { Processor, Queue, QueueOptions, Worker, WorkerOptions } from './types';
const defaultWorkerOpts: Omit<WorkerOptions, 'connection'> = {
removeOnComplete: {
count: 0 // Delete jobs immediately on completion
},
removeOnFail: {
count: 0 // Delete jobs immediately on failure
},
// BullMQ Worker important settings
lockDuration: 600000, // 10 minutes for large file operations
stalledInterval: 30000, // Check for stalled jobs every 30s
maxStalledCount: 3 // Move job to failed after 1 stall (default behavior)
};
/**
* DAL BullMQ Runtime 的业务绑定。
*
* 这个类只放默认 Worker 配置,不缓存 Runtime。Runtime 本身由 DAL 通过进程级 context
* 复用;这样 Redis Runtime 关闭后重新配置时,binding 不会继续持有旧实例。
*/
export class BullMQBinding {
private getRuntime() {
const redisRuntime = getRedisRuntime();
return getRedisBullMQRuntime({
redisRuntime,
logger: redisRuntime.getLogger(),
workerLifecycle: {
restartOnClose: true,
resumeOnPause: true
}
});
}
/** 返回当前绑定的通用日志 port,业务合同只依赖这个最小能力。 */
getLogger() {
return this.getRuntime().getLogger();
}
/** 获取或创建业务队列;连接和 Queue 生命周期由 DAL 管理。 */
getQueue<DataType, ReturnType = void>(
name: QueueNames,
opts?: Omit<QueueOptions, 'connection'>
): Queue<DataType, ReturnType> {
return this.getRuntime().getQueue<DataType, ReturnType>(name, opts);
}
/** 获取或创建业务 Worker;默认项和生命周期均由 DAL 管理。 */
getWorker<DataType, ReturnType = void>(
name: QueueNames,
processor: Processor<DataType, ReturnType>,
opts?: Omit<WorkerOptions, 'connection'>
): Worker<DataType, ReturnType> {
return this.getRuntime().getWorker<DataType, ReturnType>(name, processor, {
...defaultWorkerOpts,
...opts
});
}
}
/** 进程级 BullMQ 绑定。 */
export const bullMQ = new BullMQBinding();
import type { RedisRuntimeLogger } from '../types';
export const closeWithTimeout = ({
operation,
resource,
timeoutMs
}: {
operation: () => Promise<void>;
resource: string;
timeoutMs: number;
}) =>
new Promise<void>((resolve, reject) => {
const timeout = setTimeout(() => reject(new Error(`${resource} close timed out`)), timeoutMs);
Promise.resolve()
.then(operation)
.then(resolve, reject)
.finally(() => clearTimeout(timeout));
});
/** BullMQ close 超时后断开指定连接;强制断开失败只记录日志,不阻塞其他资源关闭。 */
export const forceDisconnect = ({
name,
resource,
disconnect,
logger
}: {
name: string;
resource: string;
disconnect: () => Promise<void> | void;
logger: RedisRuntimeLogger;
}) => {
try {
void Promise.resolve(disconnect()).catch((error) => {
logger.warn(`BullMQ ${resource} forced disconnect failed`, { name, error });
});
} catch (error) {
logger.warn(`BullMQ ${resource} forced disconnect failed`, { name, error });
}
};
import type { RedisRuntimeLogger } from '../types';
export const DEFAULT_CLOSE_TIMEOUT_MS = 5_000;
export const DEFAULT_RESTART_DELAY_MS = 1_000;
export const DEFAULT_BULLMQ_RESOURCE_ID = 'redis:default';
export const BULLMQ_RUNTIME_CONTEXT_SYMBOL = Symbol.for('@fastgpt/dal/redis/bullmq-context');
export const silentBullMQLogger: RedisRuntimeLogger = {
info: () => undefined,
warn: () => undefined,
error: () => undefined
};
export const delay = (milliseconds: number) =>
new Promise<void>((resolve) => {
setTimeout(resolve, milliseconds);
});
import { BULLMQ_RUNTIME_CONTEXT_SYMBOL, DEFAULT_BULLMQ_RESOURCE_ID } from './constants';
import { RedisBullMQRuntime } from './runtime';
import type { RedisBullMQRuntimeOptions } from './types';
type BullMQRuntimeContext = {
resources: Map<string, RedisBullMQRuntime>;
};
const getBullMQRuntimeContext = (): BullMQRuntimeContext => {
const existing = Reflect.get(globalThis, BULLMQ_RUNTIME_CONTEXT_SYMBOL) as
| BullMQRuntimeContext
| undefined;
if (existing) return existing;
const context: BullMQRuntimeContext = { resources: new Map() };
Reflect.set(globalThis, BULLMQ_RUNTIME_CONTEXT_SYMBOL, context);
return context;
};
/** 获取或复用进程级 BullMQ Runtime,避免 Next.js 热重载重复创建 Queue/Worker。 */
export const getRedisBullMQRuntime = (options: RedisBullMQRuntimeOptions) => {
const context = getBullMQRuntimeContext();
const existing = context.resources.get(DEFAULT_BULLMQ_RESOURCE_ID);
if (existing?.getState() === 'closed') {
context.resources.delete(DEFAULT_BULLMQ_RESOURCE_ID);
} else if (existing) {
if (existing.redisRuntime !== options.redisRuntime) {
throw new Error('BullMQ runtime is already bound to a different Redis runtime');
}
return existing;
}
const runtime = new RedisBullMQRuntime(options);
context.resources.set(DEFAULT_BULLMQ_RESOURCE_ID, runtime);
return runtime;
};
/** 返回已配置的进程级 Runtime,不会因读取而创建 Redis 连接。 */
export const getConfiguredRedisBullMQRuntime = () =>
getBullMQRuntimeContext().resources.get(DEFAULT_BULLMQ_RESOURCE_ID);
export { bullMQ, BullMQBinding } from './binding';
export { getConfiguredRedisBullMQRuntime, getRedisBullMQRuntime } from './context';
export { QueueNames } from './names';
export { RedisBullMQRuntime } from './runtime';
export { UnrecoverableError } from 'bullmq';
export * from './services';
export type {
BullMQRuntimeState,
BullMQWorkerLifecycleOptions,
ConnectionOptions,
Job,
JobSchedulerJson,
Processor,
Queue,
QueueOptions,
RedisBullMQRuntimeOptions,
Worker,
WorkerOptions
} from './types';
import type { Queue, Worker } from 'bullmq';
import type { BullMQEventListener, WorkerListenerSnapshot } from './types';
/** 管理 DAL 自己挂载的 listener,避免关闭或重启时删除业务 listener。 */
export class BullMQLifecycleListeners {
private readonly queueErrorHandlers = new WeakMap<Queue, BullMQEventListener>();
private readonly workerLifecycleHandlers = new WeakMap<Worker, Set<BullMQEventListener>>();
registerQueueErrorHandler(queue: Queue, handler: BullMQEventListener) {
this.queueErrorHandlers.set(queue, handler);
}
registerWorkerLifecycleHandlers(worker: Worker, handlers: Set<BullMQEventListener>) {
this.workerLifecycleHandlers.set(worker, handlers);
}
captureWorkerBusinessListeners(
worker: Worker,
lifecycleHandlers: Set<BullMQEventListener>
): WorkerListenerSnapshot[] {
const snapshots: WorkerListenerSnapshot[] = [];
for (const eventName of worker.eventNames()) {
for (const rawListener of worker.rawListeners(eventName)) {
const listenerWithOriginal = rawListener as BullMQEventListener & {
listener?: BullMQEventListener;
};
const listener = listenerWithOriginal.listener ?? listenerWithOriginal;
if (
lifecycleHandlers.has(rawListener as BullMQEventListener) ||
lifecycleHandlers.has(listener)
) {
continue;
}
snapshots.push({
eventName,
listener,
once: listenerWithOriginal.listener !== undefined
});
}
}
return snapshots;
}
restoreWorkerBusinessListeners(
worker: Worker,
listeners: readonly WorkerListenerSnapshot[] | undefined
) {
for (const { eventName, listener, once } of listeners ?? []) {
if (once) {
worker.once(eventName as any, listener as any);
} else {
worker.on(eventName as any, listener as any);
}
}
}
removeWorkerLifecycleListeners(worker: Worker) {
const lifecycleHandlers = this.workerLifecycleHandlers.get(worker);
if (!lifecycleHandlers) return;
for (const eventName of worker.eventNames()) {
for (const rawListener of worker.rawListeners(eventName)) {
const listenerWithOriginal = rawListener as BullMQEventListener & {
listener?: BullMQEventListener;
};
const listener = listenerWithOriginal.listener ?? listenerWithOriginal;
if (
lifecycleHandlers.has(rawListener as BullMQEventListener) ||
lifecycleHandlers.has(listener)
) {
worker.removeListener(eventName as any, rawListener as any);
}
}
}
this.workerLifecycleHandlers.delete(worker);
}
removeQueueLifecycleListener(queue: Queue) {
const errorHandler = this.queueErrorHandlers.get(queue);
if (!errorHandler) return;
queue.removeListener('error', errorHandler as any);
this.queueErrorHandlers.delete(queue);
}
}
/** Redis DAL 管理的业务队列名称和稳定 job namespace。 */
export enum QueueNames {
datasetSync = 'datasetSync',
evaluation = 'evaluation',
s3FileDelete = 's3FileDelete',
collectionUpdate = 'collectionUpdate',
agentSkillCreate = 'agentSkillCreate',
// Delete Queue
datasetDelete = 'datasetDelete',
appDelete = 'appDelete',
agentSkillDelete = 'agentSkillDelete',
teamDelete = 'teamDelete',
// Publish
wechatPoll = 'wechatPoll',
wechatReply = 'wechatReply',
/** @deprecated */
websiteSync = 'websiteSync'
}
import { Queue, type QueueOptions } from 'bullmq';
import type { RedisRuntime } from '../runtime/connection';
import type { RedisRuntimeLogger } from '../types';
import { closeWithTimeout, forceDisconnect } from './close';
import type { BullMQDisconnectable } from './types';
import { BullMQLifecycleListeners } from './listeners';
/** 管理 Queue 的创建、复用、错误 listener 和有序关闭。 */
export class BullMQQueueManager {
private readonly queues = new Map<string, Queue>();
private readonly listeners = new BullMQLifecycleListeners();
constructor(
private readonly options: {
redisRuntime: RedisRuntime;
logger: RedisRuntimeLogger;
closeTimeoutMs: number;
}
) {}
getQueue<DataType, ReturnType = void>(
name: string,
opts?: Omit<QueueOptions, 'connection'>
): Queue<DataType, ReturnType> {
const existing = this.queues.get(name);
if (existing) return existing as Queue<DataType, ReturnType>;
const connection = this.options.redisRuntime.createQueueConnection();
try {
const queue = new Queue<DataType, ReturnType>(name, {
...opts,
connection
});
const errorHandler = (error: Error) => {
this.options.logger.error('BullMQ queue error', { name, error });
};
this.listeners.registerQueueErrorHandler(queue, errorHandler);
queue.on('error', errorHandler);
this.queues.set(name, queue);
return queue;
} catch (error) {
this.releaseConnection(connection);
throw error;
}
}
async close() {
const activeQueues = Array.from(this.queues.entries());
this.queues.clear();
await Promise.all(activeQueues.map(([name, queue]) => this.closeQueue({ name, queue })));
}
private releaseConnection(connection: ReturnType<RedisRuntime['createQueueConnection']>) {
void this.options.redisRuntime.releaseConnection(connection).catch((error) => {
this.options.logger.warn(
'Failed to release Redis connection after BullMQ queue creation error',
{
error
}
);
});
}
private async closeQueue({ name, queue }: { name: string; queue: Queue }) {
await closeWithTimeout({
operation: () => queue.close(),
resource: `BullMQ queue ${name}`,
timeoutMs: this.options.closeTimeoutMs
}).catch((error) => {
this.options.logger.warn('BullMQ queue close failed', { name, error });
const queueConnection = (queue as unknown as { connection?: BullMQDisconnectable })
.connection;
forceDisconnect({
name,
resource: 'queue connection',
disconnect: queueConnection
? () => queueConnection.disconnect(false)
: () => queue.disconnect(),
logger: this.options.logger
});
});
this.listeners.removeQueueLifecycleListener(queue);
}
}
import type { Queue, Processor, QueueOptions, Worker, WorkerOptions } from 'bullmq';
import type { RedisRuntime } from '../runtime/connection';
import {
DEFAULT_CLOSE_TIMEOUT_MS,
DEFAULT_RESTART_DELAY_MS,
silentBullMQLogger
} from './constants';
import { BullMQQueueManager } from './queue-manager';
import type {
BullMQRuntimeState,
BullMQWorkerLifecycleOptions,
RedisBullMQRuntimeOptions
} from './types';
import { BullMQWorkerManager } from './worker-manager';
import type { RedisRuntimeLogger } from '../types';
/** 管理 DAL Redis Runtime 所拥有的 BullMQ Queue/Worker 生命周期。 */
export class RedisBullMQRuntime {
readonly redisRuntime: RedisRuntime;
private readonly logger: RedisRuntimeLogger;
private readonly queueManager: BullMQQueueManager;
private readonly workerManager: BullMQWorkerManager;
private readonly unregisterBeforeCloseHook: () => void;
private state: BullMQRuntimeState = 'running';
private closePromise: Promise<void> | undefined;
constructor({
redisRuntime,
logger = silentBullMQLogger,
closeTimeoutMs = DEFAULT_CLOSE_TIMEOUT_MS,
workerLifecycle = {},
hookName = 'bullmq'
}: RedisBullMQRuntimeOptions) {
this.redisRuntime = redisRuntime;
this.logger = logger;
const normalizedLifecycle: Required<BullMQWorkerLifecycleOptions> = {
restartOnClose: workerLifecycle.restartOnClose ?? false,
resumeOnPause: workerLifecycle.resumeOnPause ?? false,
restartDelayMs: workerLifecycle.restartDelayMs ?? DEFAULT_RESTART_DELAY_MS
};
this.queueManager = new BullMQQueueManager({
redisRuntime,
logger,
closeTimeoutMs
});
this.workerManager = new BullMQWorkerManager({
redisRuntime,
logger,
closeTimeoutMs,
workerLifecycle: normalizedLifecycle,
getState: () => this.state
});
this.unregisterBeforeCloseHook = redisRuntime.registerBeforeCloseHook({
name: hookName,
close: () => this.close()
});
}
getState() {
return this.state;
}
/** 返回当前 BullMQ Runtime 使用的通用日志 port。 */
getLogger() {
return this.logger;
}
getQueue<DataType, ReturnType = void>(
name: string,
opts?: Omit<QueueOptions, 'connection'>
): Queue<DataType, ReturnType> {
this.assertRunning();
return this.queueManager.getQueue<DataType, ReturnType>(name, opts);
}
getWorker<DataType, ReturnType = void>(
name: string,
processor: Processor<DataType, ReturnType>,
opts?: Omit<WorkerOptions, 'connection'>
): Worker<DataType, ReturnType> {
this.assertRunning();
return this.workerManager.getWorker<DataType, ReturnType>(name, processor, opts);
}
close() {
if (this.closePromise) return this.closePromise;
this.state = 'shutting-down';
this.closePromise = Promise.resolve()
.then(async () => {
// Worker 内部拥有 blocking duplicate,必须先于 Queue 和 Redis Runtime 连接关闭。
let firstError: unknown;
let hasError = false;
try {
await this.workerManager.close();
} catch (error) {
firstError = error;
hasError = true;
}
// Worker 关闭失败也不能跳过 Queue,否则队列连接会被 Redis Runtime 强制回收。
try {
await this.queueManager.close();
} catch (error) {
if (!hasError) {
firstError = error;
hasError = true;
}
}
if (hasError) throw firstError;
})
.finally(() => {
this.state = 'closed';
this.unregisterBeforeCloseHook();
});
return this.closePromise;
}
private assertRunning() {
if (this.state !== 'running') {
throw new Error(`BullMQ runtime is ${this.state}`);
}
}
}
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker } from '../types';
export type AppDeleteJobData = {
teamId: string;
appId: string;
};
const appDeleteQueueOptions = {
defaultJobOptions: {
attempts: 10,
backoff: {
type: 'exponential' as const,
delay: 5000
},
removeOnComplete: true,
removeOnFail: { age: 30 * 24 * 60 * 60 }
}
};
/** App 删除队列的业务合同和生命周期入口。 */
export class AppDeleteMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** 获取 App 删除队列;队列对象由 binding 按名称懒加载并复用。 */
getQueue(): Queue<AppDeleteJobData> {
return this.binding.getQueue<AppDeleteJobData>(QueueNames.appDelete, appDeleteQueueOptions);
}
/** 使用 service 统一的删除任务保留策略创建 App 删除 Worker。 */
getWorker(processor: Processor<AppDeleteJobData>): Worker<AppDeleteJobData> {
return this.binding.getWorker<AppDeleteJobData>(QueueNames.appDelete, processor, {
concurrency: 1,
removeOnFail: {
age: 90 * 24 * 60 * 60,
count: 10000
}
});
}
/** 投递幂等的 App 删除任务,并延迟一秒让请求先完成。 */
addJob(data: AppDeleteJobData) {
const jobId = `${String(data.teamId)}-${String(data.appId)}`;
return this.getQueue().add('delete_app', data, {
jobId,
delay: 1000
});
}
}
export const appDeleteMQService = new AppDeleteMQService();
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker, WorkerOptions } from '../types';
export type CollectionUpdateJobData = {
teamId: string;
datasetId: string;
collectionId: string;
};
const collectionUpdateQueueOptions = {
defaultJobOptions: {
attempts: 3,
backoff: {
type: 'exponential' as const,
delay: 1000
},
removeOnFail: true
}
};
/** Collection update 队列的业务合同和生命周期入口。 */
export class CollectionUpdateMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** 获取 collection update 队列;队列本身不持有 Mongo processor。 */
getQueue(): Queue<CollectionUpdateJobData> {
return this.binding.getQueue<CollectionUpdateJobData>(
QueueNames.collectionUpdate,
collectionUpdateQueueOptions
);
}
/** 注入领域 processor 创建 collection update Worker。 */
getWorker(
processor: Processor<CollectionUpdateJobData>,
opts?: Omit<WorkerOptions, 'connection'>
): Worker<CollectionUpdateJobData> {
return this.binding.getWorker<CollectionUpdateJobData>(QueueNames.collectionUpdate, processor, {
concurrency: 3,
removeOnComplete: {
count: 0
},
...opts,
// 固定 jobId 需要在最终失败后释放,否则后续 collection 更新会永久撞上旧 job。
removeOnFail: {
count: 0
}
});
}
/** 以 collectionId 去重并延迟五秒投递 update 任务。 */
async pushJob(data: CollectionUpdateJobData) {
const jobId = `collection-update-${data.collectionId}`;
try {
const queue = this.getQueue();
const addJob = () =>
queue.add('updateCollection', data, {
jobId,
delay: 5000
});
// 清理旧版本可能留下的终态 job,避免迁移前的保留策略继续阻塞固定 jobId。
const existingJob = await queue.getJob(jobId);
if (existingJob) {
const state = await existingJob.getState();
if (state === 'completed' || state === 'failed') {
await existingJob.remove();
}
}
try {
await addJob();
} catch (error) {
const isDuplicateJobError =
error instanceof Error && /already exists|duplicate/i.test(error.message);
if (!isDuplicateJobError) throw error;
// 并发调用可能在终态清理后争抢同一个 jobId;活动中的那一个已经代表本次更新。
const duplicateJob = await queue.getJob(jobId);
if (!duplicateJob) throw error;
const state = await duplicateJob.getState();
if (state === 'completed' || state === 'failed') {
await duplicateJob.remove();
await addJob();
} else {
this.binding.getLogger().info('Collection update job already queued', {
collectionId: data.collectionId,
state
});
return;
}
}
this.binding.getLogger().info('Collection update job pushed', {
collectionId: data.collectionId
});
} catch (error) {
this.binding.getLogger().error('Failed to push collection update job', {
collectionId: data.collectionId,
error
});
throw error;
}
}
}
export const collectionUpdateMQService = new CollectionUpdateMQService();
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker } from '../types';
export type DatasetDeleteJobData = {
teamId: string;
datasetId: string;
};
const datasetDeleteQueueOptions = {
defaultJobOptions: {
attempts: 10,
backoff: {
type: 'exponential' as const,
delay: 5000
},
removeOnComplete: true,
removeOnFail: { age: 30 * 24 * 60 * 60 }
}
};
/** Dataset 删除队列的业务合同和生命周期入口。 */
export class DatasetDeleteMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** 获取 Dataset 删除队列;队列对象由 binding 按名称懒加载并复用。 */
getQueue(): Queue<DatasetDeleteJobData> {
return this.binding.getQueue<DatasetDeleteJobData>(
QueueNames.datasetDelete,
datasetDeleteQueueOptions
);
}
/** 使用 service 统一的删除任务保留策略创建 Dataset 删除 Worker。 */
getWorker(processor: Processor<DatasetDeleteJobData>): Worker<DatasetDeleteJobData> {
return this.binding.getWorker<DatasetDeleteJobData>(QueueNames.datasetDelete, processor, {
concurrency: 1,
removeOnFail: {
age: 90 * 24 * 60 * 60,
count: 10000
}
});
}
/** 投递幂等的 Dataset 删除任务,并延迟一秒让请求先完成。 */
addJob(data: DatasetDeleteJobData) {
const jobId = `${String(data.teamId)}-${String(data.datasetId)}`;
return this.getQueue().add('delete_dataset', data, {
jobId,
delay: 1000
});
}
}
export const datasetDeleteMQService = new DatasetDeleteMQService();
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker } from '../types';
import { DatasetStatusEnum } from '@fastgpt/global/core/dataset/constants';
export type DatasetSyncJobData = {
datasetId: string;
};
const repeatDuration = 24 * 60 * 60 * 1000;
/** Dataset sync 队列、scheduler 和状态转换的业务服务。 */
export class DatasetSyncMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** 获取 dataset sync 队列;队列配置在首次使用时才交给 binding。 */
getQueue(): Queue<DatasetSyncJobData> {
return this.binding.getQueue<DatasetSyncJobData>(QueueNames.datasetSync, {
defaultJobOptions: {
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000
}
}
});
}
/** 创建 dataset sync Worker;实际同步 processor 由 app/pro 注入。 */
getWorker(processor: Processor<DatasetSyncJobData>): Worker<DatasetSyncJobData> {
return this.binding.getWorker<DatasetSyncJobData>(QueueNames.datasetSync, processor, {
removeOnFail: {
age: 15 * 24 * 60 * 60,
count: 1000
},
concurrency: 1
});
}
/** 投递以 datasetId 去重的同步任务。 */
addJob(data: DatasetSyncJobData) {
const datasetId = String(data.datasetId);
return this.getQueue().add(datasetId, data, { deduplication: { id: datasetId } });
}
/** 将 BullMQ 状态转换为业务侧 dataset sync 状态。 */
async getDatasetStatus(datasetId: string) {
const queue = this.getQueue();
const jobId = await queue.getDeduplicationJobId(datasetId);
if (!jobId) {
return { status: DatasetStatusEnum.active, errorMsg: undefined };
}
const job = await queue.getJob(jobId);
if (!job) {
return { status: DatasetStatusEnum.active, errorMsg: undefined };
}
const jobState = await job.getState();
if (jobState === 'failed' || jobState === 'unknown') {
return { status: DatasetStatusEnum.error, errorMsg: job.failedReason };
}
if (['waiting-children', 'waiting'].includes(jobState)) {
return { status: DatasetStatusEnum.waiting, errorMsg: undefined };
}
if (jobState === 'active') {
return { status: DatasetStatusEnum.syncing, errorMsg: undefined };
}
return { status: DatasetStatusEnum.active, errorMsg: undefined };
}
/** 创建或更新 dataset sync 的每日 scheduler。 */
upsertScheduler(data: DatasetSyncJobData, startDate?: number) {
const datasetId = String(data.datasetId);
return this.getQueue().upsertJobScheduler(
datasetId,
{
every: repeatDuration,
startDate: startDate ?? new Date().getTime() + repeatDuration
},
{
name: datasetId,
data
}
);
}
/** 读取 dataset sync scheduler。 */
getScheduler(datasetId: string) {
return this.getQueue().getJobScheduler(String(datasetId));
}
/** 删除 dataset sync scheduler。 */
removeScheduler(datasetId: string) {
return this.getQueue().removeJobScheduler(String(datasetId));
}
}
export const datasetSyncMQService = new DatasetSyncMQService();
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker, WorkerOptions } from '../types';
export type EvaluationJobData = {
evalId: string;
};
/** Evaluation 队列和状态操作的业务服务。 */
export class EvaluationMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** 获取评测队列;队列连接在首次调用时才创建。 */
getQueue(): Queue<EvaluationJobData> {
return this.binding.getQueue<EvaluationJobData>(QueueNames.evaluation, {
defaultJobOptions: {
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000
}
}
});
}
/** 获取评测 Worker;队列状态操作和重试策略集中在 BullMQ service。 */
getWorker(
processor: Processor<EvaluationJobData>,
opts?: Omit<WorkerOptions, 'connection'>
): Worker<EvaluationJobData> {
return this.binding.getWorker<EvaluationJobData>(QueueNames.evaluation, processor, {
removeOnFail: {
count: 1000
},
...opts
});
}
/** 投递以 evalId 去重的评测任务。 */
addJob(data: EvaluationJobData) {
const evalId = String(data.evalId);
return this.getQueue().add(evalId, data, { deduplication: { id: evalId } });
}
/** 查询评测任务是否仍处于可执行状态。 */
async isJobActive(evalId: string): Promise<boolean> {
try {
const queue = this.getQueue();
const jobId = await queue.getDeduplicationJobId(String(evalId));
if (!jobId) return false;
const job = await queue.getJob(jobId);
if (!job) return false;
const jobState = await job.getState();
return ['waiting', 'delayed', 'prioritized', 'active'].includes(jobState);
} catch (error) {
this.binding.getLogger().error('Failed to check evaluation job status', { evalId, error });
return false;
}
}
/** 删除尚未开始执行的评测任务,active/completed 任务保持原状态。 */
async removeJob(evalId: string): Promise<boolean> {
const formatEvalId = String(evalId);
try {
const queue = this.getQueue();
const jobId = await queue.getDeduplicationJobId(formatEvalId);
if (!jobId) {
this.binding.getLogger().warn('No evaluation job found to remove', { evalId });
return false;
}
const job = await queue.getJob(jobId);
if (!job) {
this.binding.getLogger().warn('Evaluation job not found in queue', { evalId, jobId });
return false;
}
const jobState = await job.getState();
if (['waiting', 'delayed', 'prioritized'].includes(jobState)) {
await job.remove();
this.binding.getLogger().info('Evaluation job removed successfully', {
evalId,
jobId,
jobState
});
return true;
}
this.binding.getLogger().warn('Cannot remove active or completed evaluation job', {
evalId,
jobId,
jobState
});
return false;
} catch (error) {
this.binding.getLogger().error('Failed to remove evaluation job', { evalId, error });
return false;
}
}
}
export const evaluationMQService = new EvaluationMQService();
/**
* BullMQ 业务队列合同的集中入口。
*
* 领域代码通过 DAL 入口使用这里的队列 data type、enqueue、scheduler 和 worker binding。
*/
export * from './appDelete';
export * from './collectionUpdate';
export * from './datasetDelete';
export * from './datasetSync';
export * from './evaluation';
export * from './s3FileDelete';
export * from './skillCreate';
export * from './skillDelete';
export * from './teamDelete';
export * from './wechat';
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker } from '../types';
export type S3MQJobData = {
key?: string;
keys?: string[];
prefix?: string;
bucketName: string;
};
const s3DeleteJobOptions = {
attempts: 10,
removeOnFail: {
count: 10000,
age: 14 * 24 * 60 * 60
},
removeOnComplete: true,
backoff: {
delay: 2000,
type: 'exponential' as const
}
};
const encodeJobIdPart = (value: string) => encodeURIComponent(value);
/** S3 文件删除队列的业务合同和生命周期入口。 */
export class S3FileDeleteMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** 获取 S3 文件删除队列;对象存储删除 processor 由 common/s3 注入。 */
getQueue(): Queue<S3MQJobData> {
return this.binding.getQueue<S3MQJobData>(QueueNames.s3FileDelete);
}
/** 创建 S3 文件删除 Worker,统一保留策略仍由队列 service 管理。 */
getWorker(processor: Processor<S3MQJobData>): Worker<S3MQJobData> {
return this.binding.getWorker<S3MQJobData>(QueueNames.s3FileDelete, processor, {
concurrency: 6
});
}
/** 根据对象 key/prefix 生成幂等任务并投递到 S3 删除队列。 */
async addJob(data: S3MQJobData): Promise<void> {
const jobId = (() => {
if (data.key) {
return `s3-key-${encodeJobIdPart(data.bucketName)}|${encodeJobIdPart(data.key)}`;
}
if (data.keys) return undefined;
if (data.prefix) {
return `s3-prefix-${encodeJobIdPart(data.bucketName)}|${encodeJobIdPart(data.prefix)}`;
}
throw new Error('Invalid s3 delete job data');
})();
await this.getQueue().add('delete-s3-files', data, { jobId, ...s3DeleteJobOptions });
}
}
export const s3FileDeleteMQService = new S3FileDeleteMQService();
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker } from '../types';
export type AgentSkillCreateJobData = {
skillId: string;
teamId: string;
tmbId: string;
};
/** Skill 创建队列和幂等任务操作的业务服务。 */
export class SkillCreateMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** Skill 创建任务的默认队列配置;只在首次调用时创建 Queue。 */
getQueue(): Queue<AgentSkillCreateJobData> {
return this.binding.getQueue<AgentSkillCreateJobData>(QueueNames.agentSkillCreate, {
defaultJobOptions: {
attempts: 1,
removeOnComplete: true,
removeOnFail: {
age: 30 * 24 * 60 * 60,
count: 1000
}
}
});
}
/** 创建 Skill 创建 Worker;workspace、Mongo 和对象存储逻辑由 processor 所属领域维护。 */
getWorker(processor: Processor<AgentSkillCreateJobData>): Worker<AgentSkillCreateJobData> {
return this.binding.getWorker<AgentSkillCreateJobData>(QueueNames.agentSkillCreate, processor, {
concurrency: 2
});
}
/** 投递以 skillId 去重的 Skill 创建任务,并清理已经完成的旧任务。 */
async addJob(data: AgentSkillCreateJobData) {
const skillId = String(data.skillId);
const queue = this.getQueue();
const existingJob = await queue.getJob(skillId);
if (existingJob) {
const state = await existingJob.getState();
if (state !== 'completed' && state !== 'failed') {
return existingJob;
}
await existingJob.remove();
}
return queue.add(skillId, data, {
jobId: skillId
});
}
}
export const skillCreateMQService = new SkillCreateMQService();
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker } from '../types';
export type AgentSkillDeleteJobData = {
teamId: string;
skillId: string;
};
const agentSkillDeleteQueueOptions = {
defaultJobOptions: {
attempts: 10,
backoff: {
type: 'exponential' as const,
delay: 5000
},
removeOnComplete: true,
removeOnFail: { age: 30 * 24 * 60 * 60 }
}
};
/** Skill 删除队列的业务合同和生命周期入口。 */
export class SkillDeleteMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** 获取 Skill 删除队列;队列对象由 binding 按名称懒加载并复用。 */
getQueue(): Queue<AgentSkillDeleteJobData> {
return this.binding.getQueue<AgentSkillDeleteJobData>(
QueueNames.agentSkillDelete,
agentSkillDeleteQueueOptions
);
}
/** 创建 Skill 删除 Worker;具体清理由调用方注入 processor。 */
getWorker(processor: Processor<AgentSkillDeleteJobData>): Worker<AgentSkillDeleteJobData> {
return this.binding.getWorker<AgentSkillDeleteJobData>(QueueNames.agentSkillDelete, processor, {
concurrency: 1,
removeOnFail: {
age: 90 * 24 * 60 * 60,
count: 10000
}
});
}
/** 投递以 teamId-skillId 去重的 Skill 删除任务。 */
addJob(data: AgentSkillDeleteJobData) {
const jobId = `${String(data.teamId)}-${String(data.skillId)}`;
return this.getQueue().add('delete_agent_skill', data, {
jobId,
delay: 1000
});
}
}
export const skillDeleteMQService = new SkillDeleteMQService();
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker } from '../types';
export type TeamDeleteJobData = {
teamId: string;
};
const teamDeleteQueueOptions = {
defaultJobOptions: {
attempts: 10,
backoff: {
type: 'exponential' as const,
delay: 5000
},
removeOnComplete: true,
removeOnFail: { age: 30 * 24 * 60 * 60 }
}
};
/** Team 删除队列的业务合同和生命周期入口。 */
export class TeamDeleteMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
/** 获取 Team 删除队列;队列对象由 binding 按名称懒加载并复用。 */
getQueue(): Queue<TeamDeleteJobData> {
return this.binding.getQueue<TeamDeleteJobData>(QueueNames.teamDelete, teamDeleteQueueOptions);
}
/** 使用 service 统一的删除任务保留策略创建 Team 删除 Worker。 */
getWorker(processor: Processor<TeamDeleteJobData>): Worker<TeamDeleteJobData> {
return this.binding.getWorker<TeamDeleteJobData>(QueueNames.teamDelete, processor, {
concurrency: 1,
removeOnFail: {
age: 90 * 24 * 60 * 60,
count: 10000
}
});
}
/** 投递幂等的 Team 删除任务,并延迟一秒让请求先完成。 */
addJob(data: TeamDeleteJobData) {
return this.getQueue().add('delete_team', data, {
jobId: String(data.teamId),
delay: 1000
});
}
}
export const teamDeleteMQService = new TeamDeleteMQService();
import { bullMQ, type BullMQBinding } from '../binding';
import { QueueNames } from '../names';
import type { Processor, Queue, Worker, WorkerOptions } from '../types';
export type WechatPollJobData = {
shareId: string;
};
export type WechatReplyJobData = {
shareId: string;
userId: string;
text: string;
contextToken: string;
lastMsgId: string;
};
export const WECHAT_POLL_JOB_NAME = 'wechatPublishPoll';
export const WECHAT_REPLY_JOB_NAME = 'wechatPublishReply';
/** 微信轮询和回复队列的业务服务。 */
export class WechatMQService {
constructor(private readonly binding: BullMQBinding = bullMQ) {}
getPollQueue(): Queue<WechatPollJobData> {
return this.binding.getQueue<WechatPollJobData>(QueueNames.wechatPoll);
}
getReplyQueue(): Queue<WechatReplyJobData> {
return this.binding.getQueue<WechatReplyJobData>(QueueNames.wechatReply);
}
/** 创建微信轮询 Worker,轮询 processor 的返回值用于 completed 续链策略。 */
getPollWorker<ReturnType = boolean>(
processor: Processor<WechatPollJobData, ReturnType>,
opts: Omit<WorkerOptions, 'connection'>
): Worker<WechatPollJobData, ReturnType> {
return this.binding.getWorker<WechatPollJobData, ReturnType>(
QueueNames.wechatPoll,
processor,
opts
);
}
/** 创建微信回复 Worker。 */
getReplyWorker(
processor: Processor<WechatReplyJobData>,
opts: Omit<WorkerOptions, 'connection'>
): Worker<WechatReplyJobData> {
return this.binding.getWorker<WechatReplyJobData>(QueueNames.wechatReply, processor, opts);
}
/** 投递微信回复任务,具体幂等键由调用方根据消息 id 生成。 */
addReplyJob(
data: WechatReplyJobData,
opts: {
jobId: string;
backoff?: { type: 'fixed'; delay: number };
}
) {
return this.getReplyQueue().add(WECHAT_REPLY_JOB_NAME, data, opts);
}
/** 投递下一轮微信轮询任务。 */
addPollJob(
data: WechatPollJobData,
opts: {
jobId: string;
delay?: number;
removeOnComplete?: boolean;
removeOnFail?: boolean;
}
) {
return this.getPollQueue().add(WECHAT_POLL_JOB_NAME, data, opts);
}
/** 删除指定渠道的轮询任务。active 任务由 BullMQ 自身状态决定是否可移除。 */
removePollJob(jobId: string) {
return this.getPollQueue().remove(jobId);
}
}
export const wechatMQService = new WechatMQService();
import type {
ConnectionOptions,
Job,
JobSchedulerJson,
Processor,
Queue,
QueueOptions,
Worker,
WorkerOptions
} from 'bullmq';
import type { RedisRuntime } from '../runtime/connection';
import type { RedisRuntimeLogger } from '../types';
export type BullMQRuntimeState = 'running' | 'shutting-down' | 'closed';
/** BullMQ worker 在 Redis 连接异常关闭时的可选恢复策略。 */
export type BullMQWorkerLifecycleOptions = {
restartOnClose?: boolean;
resumeOnPause?: boolean;
restartDelayMs?: number;
};
export type RedisBullMQRuntimeOptions = {
redisRuntime: RedisRuntime;
logger?: RedisRuntimeLogger;
closeTimeoutMs?: number;
workerLifecycle?: BullMQWorkerLifecycleOptions;
hookName?: string;
};
export type BullMQDisconnectable = {
disconnect: (wait?: boolean) => Promise<void> | void;
};
export type BullMQEventListener = (...args: any[]) => unknown;
export type WorkerListenerSnapshot = {
eventName: string | symbol;
listener: BullMQEventListener;
once: boolean;
};
export type {
ConnectionOptions,
Job,
JobSchedulerJson,
Processor,
Queue,
QueueOptions,
Worker,
WorkerOptions
};
import { Worker, type Processor, type WorkerOptions } from 'bullmq';
import type { RedisRuntime } from '../runtime/connection';
import type { RedisRuntimeLogger } from '../types';
import { closeWithTimeout, forceDisconnect } from './close';
import { delay } from './constants';
import { BullMQLifecycleListeners } from './listeners';
import type {
BullMQDisconnectable,
BullMQEventListener,
BullMQRuntimeState,
BullMQWorkerLifecycleOptions,
WorkerListenerSnapshot
} from './types';
/** 管理 Worker 的创建、异常恢复、业务 listener 迁移和有序关闭。 */
export class BullMQWorkerManager {
private readonly workers = new Map<string, Worker>();
private readonly restartingWorkers = new Set<string>();
private readonly listeners = new BullMQLifecycleListeners();
private readonly lifecycle: Required<BullMQWorkerLifecycleOptions>;
constructor(
private readonly options: {
redisRuntime: RedisRuntime;
logger: RedisRuntimeLogger;
closeTimeoutMs: number;
workerLifecycle: Required<BullMQWorkerLifecycleOptions>;
getState: () => BullMQRuntimeState;
}
) {
this.lifecycle = options.workerLifecycle;
}
getWorker<DataType, ReturnType = void>(
name: string,
processor: Processor<DataType, ReturnType>,
opts?: Omit<WorkerOptions, 'connection'>
): Worker<DataType, ReturnType> {
const existing = this.workers.get(name);
if (existing) return existing as Worker<DataType, ReturnType>;
if (this.restartingWorkers.has(name)) {
throw new Error(`BullMQ worker ${name} is restarting`);
}
const worker = this.createWorker({ name, processor, opts });
this.workers.set(name, worker);
return worker;
}
async close() {
const activeWorkers = Array.from(this.workers.entries());
this.workers.clear();
this.restartingWorkers.clear();
await Promise.all(activeWorkers.map(([name, worker]) => this.closeWorker({ name, worker })));
}
private createWorker<DataType, ReturnType>({
name,
processor,
opts,
listeners
}: {
name: string;
processor: Processor<DataType, ReturnType>;
opts?: Omit<WorkerOptions, 'connection'>;
listeners?: readonly WorkerListenerSnapshot[];
}): Worker<DataType, ReturnType> {
const connection = this.options.redisRuntime.createWorkerConnection();
try {
const worker = new Worker<DataType, ReturnType>(name, processor, {
...opts,
connection
});
const lifecycleHandlers = new Set<BullMQEventListener>();
const readyHandler: BullMQEventListener = () => {
this.options.logger.info('BullMQ worker ready', { name });
};
const errorHandler: BullMQEventListener = (error) => {
this.options.logger.error('BullMQ worker error', { name, error });
};
const closedHandler: BullMQEventListener = () => {
if (this.workers.get(name) !== worker) return;
this.workers.delete(name);
const shouldRestart =
this.options.getState() === 'running' && this.lifecycle.restartOnClose;
const businessListeners = shouldRestart
? this.listeners.captureWorkerBusinessListeners(worker, lifecycleHandlers)
: undefined;
this.listeners.removeWorkerLifecycleListeners(worker);
if (!shouldRestart) return;
this.restartingWorkers.add(name);
this.options.logger.warn('BullMQ worker closed, attempting restart', { name });
void this.restartWorker({ name, processor, opts, listeners: businessListeners }).finally(
() => {
this.restartingWorkers.delete(name);
}
);
};
const pausedHandler: BullMQEventListener = () => {
if (
this.options.getState() !== 'running' ||
!this.lifecycle.resumeOnPause ||
this.workers.get(name) !== worker
) {
return;
}
this.options.logger.warn('BullMQ worker paused', { name });
void delay(this.lifecycle.restartDelayMs)
.then(() => {
if (this.options.getState() === 'running' && this.workers.get(name) === worker) {
return worker.resume();
}
return undefined;
})
.catch((error) => {
this.options.logger.warn('BullMQ worker resume failed', { name, error });
});
};
lifecycleHandlers.add(readyHandler);
lifecycleHandlers.add(errorHandler);
lifecycleHandlers.add(closedHandler);
lifecycleHandlers.add(pausedHandler);
this.listeners.registerWorkerLifecycleHandlers(worker, lifecycleHandlers);
worker.on('ready', readyHandler);
worker.on('error', errorHandler);
worker.on('closed', closedHandler);
worker.on('paused', pausedHandler);
this.listeners.restoreWorkerBusinessListeners(worker, listeners);
return worker;
} catch (error) {
this.releaseConnection(connection);
throw error;
}
}
private async restartWorker<DataType, ReturnType>({
name,
processor,
opts,
listeners
}: {
name: string;
processor: Processor<DataType, ReturnType>;
opts?: Omit<WorkerOptions, 'connection'>;
listeners?: readonly WorkerListenerSnapshot[];
}) {
while (this.options.getState() === 'running') {
try {
const worker = this.createWorker({ name, processor, opts, listeners });
if (this.options.getState() !== 'running') {
await this.closeWorker({ name, worker });
return;
}
this.workers.set(name, worker);
this.options.logger.info('BullMQ worker restarted successfully', { name });
return;
} catch (error) {
this.options.logger.error('BullMQ worker restart failed, will retry', { name, error });
await delay(this.lifecycle.restartDelayMs);
}
}
}
private releaseConnection(connection: ReturnType<RedisRuntime['createWorkerConnection']>) {
void this.options.redisRuntime.releaseConnection(connection).catch((error) => {
this.options.logger.warn(
'Failed to release Redis connection after BullMQ worker creation error',
{ error }
);
});
}
private async closeWorker({ name, worker }: { name: string; worker: Worker }) {
await closeWithTimeout({
operation: () => worker.close(true),
resource: `BullMQ worker ${name}`,
timeoutMs: this.options.closeTimeoutMs
}).catch((error) => {
this.options.logger.warn('BullMQ worker close failed', { name, error });
const workerConnection = (worker as unknown as { connection?: BullMQDisconnectable })
.connection;
forceDisconnect({
name,
resource: 'worker connection',
disconnect: workerConnection
? () => workerConnection.disconnect(false)
: () => worker.disconnect(),
logger: this.options.logger
});
const blockingConnection = (
worker as unknown as { blockingConnection?: BullMQDisconnectable }
).blockingConnection;
if (blockingConnection) {
forceDisconnect({
name,
resource: 'worker blocking connection',
disconnect: () => blockingConnection.disconnect(false),
logger: this.options.logger
});
}
});
this.listeners.removeWorkerLifecycleListeners(worker);
}
}
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import type { RedisCacheLogger } from '../types';
const DAILY_ACTIVE_DEDUPE_TTL_SECONDS = 24 * 60 * 60;
export type DailyActiveDedupeCacheOptions = {
redis?: RedisCacheAdapter;
logger: RedisCacheLogger<'warn'>;
};
/**
* 每日活跃用户去重 Cache。
*
* 同一 UTC 日期内只有第一个请求能原子声明历史 key;Redis 故障时 fail-open,允许本次
* tracking 继续写入事实存储,避免缓存故障造成活跃事件丢失。
*/
export class DailyActiveDedupeCache {
private readonly redis: RedisCacheAdapter;
private readonly logger: RedisCacheLogger<'warn'>;
constructor({ redis = redisCacheAdapter, logger }: DailyActiveDedupeCacheOptions) {
this.redis = redis;
this.logger = logger;
}
/** 返回本次请求是否应记录 daily active;Redis 故障时降级为 true。 */
async shouldRecord({ uid, date }: { uid: string; date: string }) {
try {
return await this.redis.setIfAbsent({
key: asRedisLogicalKey(`cache:dailyUserActive:${uid}_${date}`),
value: '1',
ttlSeconds: DAILY_ACTIVE_DEDUPE_TTL_SECONDS
});
} catch (error) {
this.logger.warn('Daily active dedupe failed open', { error });
return true;
}
}
}
import { createHash } from 'node:crypto';
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import type { RedisCacheLogger } from '../types';
const TOKEN_SAFE_WINDOW_SECONDS = 5 * 60;
type AccessTokenResponse = {
accessToken: string;
expireIn: number;
};
type DingtalkAccessTokenServer = {
appKey: string;
appSecret?: string;
};
export type DingtalkAccessTokenCacheOptions = {
redis?: RedisCacheAdapter;
logger: RedisCacheLogger<'warn'>;
};
/**
* Dingtalk access token Cache。
*
* Cache 保持历史物理 key、动态 TTL 和进程内 single-flight。Redis 读写失败按缓存
* miss 降级,不阻断上游 token 获取;上游错误仍原样抛给调用方。
*/
export class DingtalkAccessTokenCache {
private readonly redis: RedisCacheAdapter;
private readonly logger: RedisCacheLogger<'warn'>;
private readonly refreshingTokenMap = new Map<string, Promise<string>>();
constructor({ redis = redisCacheAdapter, logger }: DingtalkAccessTokenCacheOptions) {
this.redis = redis;
this.logger = logger;
}
private getCacheKey = ({ appKey, appSecret }: DingtalkAccessTokenServer) => {
const secretHash = createHash('sha256')
.update(appSecret ?? '')
.digest('hex')
.slice(0, 12);
return asRedisLogicalKey(`cache:dataset:dingtalk:accessToken:${appKey}:${secretHash}`);
};
/** 读取缓存或合并并发刷新;Redis 只承担加速,不作为 token 事实来源。 */
async getOrRefresh({
server,
fetchToken
}: {
server: DingtalkAccessTokenServer;
fetchToken: () => Promise<AccessTokenResponse>;
}) {
const cacheKey = this.getCacheKey(server);
try {
const cachedToken = await this.redis.get(cacheKey);
if (cachedToken) return cachedToken;
} catch (error) {
this.logger.warn('DingTalk accessToken cache read failed', {
provider: 'dingtalk',
appKey: server.appKey,
error
});
}
const refreshingToken = this.refreshingTokenMap.get(cacheKey);
if (refreshingToken) return refreshingToken;
const refreshPromise = (async () => {
try {
const { accessToken, expireIn } = await fetchToken();
const ttlSeconds = Math.max(expireIn - TOKEN_SAFE_WINDOW_SECONDS, 60);
try {
await this.redis.set({
key: cacheKey,
value: accessToken,
ttlMs: ttlSeconds * 1000
});
} catch (error) {
this.logger.warn('DingTalk accessToken cache write failed', {
provider: 'dingtalk',
appKey: server.appKey,
ttl: ttlSeconds,
error
});
}
return accessToken;
} catch (error) {
await this.redis.delete(cacheKey).catch(() => undefined);
throw error;
} finally {
this.refreshingTokenMap.delete(cacheKey);
}
})();
this.refreshingTokenMap.set(cacheKey, refreshPromise);
return refreshPromise;
}
}
import {
asRedisLogicalKey,
redisCacheAdapter,
RedisInvalidArgumentError,
type RedisCacheAdapter
} from '../adapter';
import { PositiveSafeIntegerSchema } from '../runtime/schema';
export type FixedWindowRateLimitResult = {
allowed: boolean;
currentCount: number;
remaining: number;
ttlSeconds: number;
resetAt: number;
};
export type FixedWindowRateLimitCacheOptions = {
redis?: RedisCacheAdapter;
now?: () => number;
};
/**
* 固定窗口限流 Cache。
*
* 计数与 TTL 的原子性由 adapter 保证;Cache 只负责限制值校验和业务决策结果。
* Redis 执行错误向上抛出,由认证或 API service 统一映射为 fail-closed。
*/
export class FixedWindowRateLimitCache {
private readonly redis: RedisCacheAdapter;
private readonly now: () => number;
constructor({
redis = redisCacheAdapter,
now = Date.now
}: FixedWindowRateLimitCacheOptions = {}) {
this.redis = redis;
this.now = now;
}
async consume({
key,
limit,
windowSeconds = 60
}: {
key: string;
limit: number;
windowSeconds?: number;
}): Promise<FixedWindowRateLimitResult> {
const parsedLimit = PositiveSafeIntegerSchema.safeParse(limit);
if (!parsedLimit.success) {
throw new RedisInvalidArgumentError({
operation: 'fixedWindow.consume',
message: 'limit must be a positive safe integer'
});
}
const { currentCount, ttlSeconds } = await this.redis.consumeFixedWindow({
key: asRedisLogicalKey(key),
windowSeconds
});
return {
allowed: currentCount <= parsedLimit.data,
currentCount,
remaining: Math.max(0, parsedLimit.data - currentCount),
ttlSeconds,
resetAt: this.now() + ttlSeconds * 1000
};
}
}
export const fixedWindowRateLimitCache = new FixedWindowRateLimitCache();
export { DailyActiveDedupeCache } from './dailyActiveDedupe';
export type { DailyActiveDedupeCacheOptions } from './dailyActiveDedupe';
export { DingtalkAccessTokenCache } from './dingtalkAccessToken';
export type { DingtalkAccessTokenCacheOptions } from './dingtalkAccessToken';
export { SystemVersionCache, systemVersionCache } from './systemVersion';
export type { SystemVersionCacheOptions } from './systemVersion';
export { FixedWindowRateLimitCache, fixedWindowRateLimitCache } from './fixedWindowRateLimit';
export type {
FixedWindowRateLimitCacheOptions,
FixedWindowRateLimitResult
} from './fixedWindowRateLimit';
export { TeamQpmCache, teamQpmCache } from './teamQpm';
export type { TeamQpmCacheOptions } from './teamQpm';
export { TeamPointCache, teamPointCache } from './teamPoint';
export type { TeamPointCacheOptions, TeamPointSnapshot } from './teamPoint';
export { TeamVectorCountCache } from './teamVectorCount';
export type { TeamVectorCountCacheOptions } from './teamVectorCount';
export {
WECHAT_QR_LOGIN_TTL_SECONDS,
WechatQrLoginCache,
wechatQrLoginCache
} from './wechatQrLogin';
export type { WechatQrLoginCacheOptions, WechatQrLoginData } from './wechatQrLogin';
export { SESSION_TTL_SECONDS, SessionCache, SessionDataSchema } from './session';
export type { SessionCacheOptions, SessionData, SessionRecord } from './session';
export {
LeaseCache,
isRedisLeaseError,
RedisLeaseAcquireError,
RedisLeaseLostError,
RedisLeaseUnavailableError
} from './lease';
export type { LeaseCacheOptions, RedisLeaseContext, WithLeaseOptions } from './lease';
export {
WORKFLOW_STOP_SIGNAL_TTL_SECONDS,
WorkflowStopSignalCache,
WorkflowStopSignalParamsSchema,
getWorkflowStopSignalKey
} from './workflowStopSignal';
export type {
WorkflowStopSignalCacheOptions,
WorkflowStopSignalParams
} from './workflowStopSignal';
export {
StreamResumeActiveStateSchema,
StreamResumeCache,
StreamResumeParamsSchema,
StreamResumeUnavailableStateSchema
} from './streamResume';
export type {
StreamResumeActiveState,
StreamResumeCacheOptions,
StreamResumeKeys,
StreamResumeParams,
StreamResumeUnavailableState
} from './streamResume';
export {
OUTLINK_STREAM_CONTENT_TTL_SECONDS,
OUTLINK_STREAM_END_FLAG,
OUTLINK_STREAM_INITIAL_TTL_SECONDS,
OutLinkStreamCache,
getOutLinkStreamKey,
outLinkStreamCache
} from './outLinkStream';
export type { OutLinkStreamCacheOptions } from './outLinkStream';
export {
WECHAT_POLLING_FAILURE_TTL_SECONDS,
WechatPollingFailureCache,
getWechatPollingFailureKey,
wechatPollingFailureCache
} from './wechatPollingFailure';
export type { WechatPollingFailureCacheOptions } from './wechatPollingFailure';
export type { RedisCacheLogger } from '../types';
import { randomUUID } from 'node:crypto';
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import { RedisInvalidArgumentError } from '../runtime/errors';
import { parsePositiveInteger } from '../runtime/validation';
import type { RedisCacheLogger } from '../types';
const LEASE_KEY_PREFIX = 'lock:';
const MAX_TIMER_DELAY_MS = 2_147_483_647;
export class RedisLeaseUnavailableError extends Error {
constructor({ key, label }: { key: string; label: string }) {
super(`Redis lease is already held for ${label}: ${key}`);
this.name = 'RedisLeaseUnavailableError';
}
}
export class RedisLeaseLostError extends Error {
constructor({ key, label }: { key: string; label: string }) {
super(`Redis lease was lost while running ${label}: ${key}`);
this.name = 'RedisLeaseLostError';
}
}
export class RedisLeaseAcquireError extends Error {
cause: unknown;
constructor({ key, label, cause }: { key: string; label: string; cause: unknown }) {
super(`Failed to acquire Redis lease for ${label}: ${key}`);
this.name = 'RedisLeaseAcquireError';
this.cause = cause;
}
}
export const isRedisLeaseError = (error: unknown) =>
error instanceof RedisLeaseUnavailableError ||
error instanceof RedisLeaseLostError ||
error instanceof RedisLeaseAcquireError;
export type LeaseCacheOptions = {
redis?: RedisCacheAdapter;
logger: RedisCacheLogger<'warn'>;
};
export type RedisLeaseContext = {
/** lease 丢失后触发,支持传递给可取消的 provider 请求。 */
signal: AbortSignal;
/** 在进入下一步不可逆副作用前确认当前执行者仍持有 lease。 */
assertValid: () => void;
};
export type WithLeaseOptions<T> = {
key: string;
label: string;
ttlMs: number;
renewIntervalMs?: number;
fn: (context: RedisLeaseContext) => Promise<T>;
};
/**
* Redis Lease Cache。
*
* Lease 只负责 token 生命周期和续租/释放竞态;key 的物理前缀与 Lua 原子操作由 adapter
* 负责。获取异常和租约丢失会阻断临界区,释放失败只记录 warning,避免误报业务结果。
*/
export class LeaseCache {
private readonly redis: RedisCacheAdapter;
private readonly logger: RedisCacheLogger<'warn'>;
constructor({ redis = redisCacheAdapter, logger }: LeaseCacheOptions) {
this.redis = redis;
this.logger = logger;
}
private getLeaseKey = (key: string) => {
if (typeof key !== 'string' || key.length === 0) {
throw new RedisInvalidArgumentError({
operation: 'lease.with',
message: 'key must be a non-empty string'
});
}
return asRedisLogicalKey(`${LEASE_KEY_PREFIX}${key}`);
};
private parseLeaseOptions = ({
ttlMs,
renewIntervalMs
}: Pick<WithLeaseOptions<unknown>, 'ttlMs' | 'renewIntervalMs'>) => {
const parsedTtlMs = parsePositiveInteger({
value: ttlMs,
operation: 'lease.with',
field: 'ttlMs'
});
const parsedRenewIntervalMs = parsePositiveInteger({
value: renewIntervalMs ?? Math.max(1, Math.floor(parsedTtlMs / 6)),
operation: 'lease.with',
field: 'renewIntervalMs'
});
if (parsedRenewIntervalMs >= parsedTtlMs) {
throw new RedisInvalidArgumentError({
operation: 'lease.with',
message: 'renewIntervalMs must be smaller than ttlMs'
});
}
return { parsedTtlMs, parsedRenewIntervalMs };
};
/** 在临界区执行带自动续租的 Redis Lease。 */
async withLease<T>({ key, label, ttlMs, renewIntervalMs, fn }: WithLeaseOptions<T>): Promise<T> {
const leaseKey = this.getLeaseKey(key);
const { parsedTtlMs, parsedRenewIntervalMs } = this.parseLeaseOptions({
ttlMs,
renewIntervalMs
});
const token = randomUUID();
let leaseLostError: RedisLeaseLostError | undefined;
let leaseExpiresAt = Date.now() + parsedTtlMs;
let active = true;
let renewalInFlight = false;
let expiryTimer: ReturnType<typeof setTimeout> | undefined;
let renewTimer: ReturnType<typeof setTimeout> | undefined;
const abortController = new AbortController();
const clearExpiryTimer = () => {
if (expiryTimer !== undefined) {
clearTimeout(expiryTimer);
expiryTimer = undefined;
}
};
const clearRenewTimer = () => {
if (renewTimer !== undefined) {
clearTimeout(renewTimer);
renewTimer = undefined;
}
};
const markLeaseLost = () => {
leaseLostError ??= new RedisLeaseLostError({ key: leaseKey, label });
clearExpiryTimer();
if (!abortController.signal.aborted) {
abortController.abort(leaseLostError);
}
return leaseLostError;
};
const assertValid = () => {
if (!leaseLostError && Date.now() >= leaseExpiresAt) {
markLeaseLost();
}
if (leaseLostError) throw leaseLostError;
};
const scheduleExpiryCheck = () => {
clearExpiryTimer();
if (!active || leaseLostError) return;
const remainingMs = leaseExpiresAt - Date.now();
if (remainingMs <= 0) {
markLeaseLost();
return;
}
expiryTimer = setTimeout(
() => {
expiryTimer = undefined;
if (!active || leaseLostError) return;
if (Date.now() >= leaseExpiresAt) {
markLeaseLost();
} else {
scheduleExpiryCheck();
}
},
Math.max(0, Math.min(remainingMs, MAX_TIMER_DELAY_MS))
);
expiryTimer.unref?.();
};
const renew = async () => {
if (!active || leaseLostError || renewalInFlight) return;
renewalInFlight = true;
const renewStartedAt = Date.now();
try {
const renewed = await this.redis.renewLease({
key: leaseKey,
token,
ttlMs: parsedTtlMs
});
if (!active || leaseLostError) return;
if (renewed) {
leaseExpiresAt = renewStartedAt + parsedTtlMs;
scheduleExpiryCheck();
return;
}
markLeaseLost();
this.logger.warn('Redis lease renew failed because token no longer matches', {
key: leaseKey,
label
});
} catch (error) {
this.logger.warn('Redis lease renew failed', { key: leaseKey, label, error });
if (Date.now() >= leaseExpiresAt) {
markLeaseLost();
}
} finally {
renewalInFlight = false;
}
};
const scheduleRenewal = (remainingMs = parsedRenewIntervalMs) => {
clearRenewTimer();
if (!active || leaseLostError) return;
// Node 会截断超长 timer;分段等待避免大 TTL 被误调度成高频续租。
const delayMs = Math.min(remainingMs, MAX_TIMER_DELAY_MS);
renewTimer = setTimeout(() => {
renewTimer = undefined;
if (!active || leaseLostError) return;
if (remainingMs > MAX_TIMER_DELAY_MS) {
scheduleRenewal(remainingMs - MAX_TIMER_DELAY_MS);
return;
}
void renew()
.finally(() => scheduleRenewal())
.catch((error) => {
this.logger.warn('Redis lease renewal scheduler failed', {
key: leaseKey,
label,
error
});
});
}, delayMs);
renewTimer.unref?.();
};
let acquired: boolean;
const acquireStartedAt = Date.now();
try {
acquired = await this.redis.acquireLease({
key: leaseKey,
token,
ttlMs: parsedTtlMs
});
} catch (error) {
this.logger.warn('Redis lease acquire failed', { key: leaseKey, label, error });
throw new RedisLeaseAcquireError({ key: leaseKey, label, cause: error });
}
if (!acquired) {
throw new RedisLeaseUnavailableError({ key: leaseKey, label });
}
leaseExpiresAt = acquireStartedAt + parsedTtlMs;
scheduleExpiryCheck();
scheduleRenewal();
try {
assertValid();
const result = await fn({ signal: abortController.signal, assertValid });
assertValid();
return result;
} catch (error) {
assertValid();
throw error;
} finally {
active = false;
clearRenewTimer();
clearExpiryTimer();
try {
await this.redis.releaseLease({ key: leaseKey, token });
} catch (error) {
this.logger.warn('Redis lease release failed', { key: leaseKey, label, error });
}
}
}
}
import { z } from 'zod';
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import { PositiveSafeIntegerSchema } from '../runtime/schema';
export const OUTLINK_STREAM_INITIAL_TTL_SECONDS = 120;
export const OUTLINK_STREAM_CONTENT_TTL_SECONDS = 60;
export const OUTLINK_STREAM_END_FLAG = '[DONE]';
const OUTLINK_STREAM_NAMESPACE = 'cache:streamResponse';
const StreamIdSchema = z.string().min(1);
export type OutLinkStreamCacheOptions = {
redis?: RedisCacheAdapter;
};
/** 构造 OutLink 字符串响应缓存的逻辑 key,禁止业务层直接拼接 cache 前缀。 */
export const getOutLinkStreamKey = (streamId: string) =>
asRedisLogicalKey(`${OUTLINK_STREAM_NAMESPACE}:${StreamIdSchema.parse(streamId)}`);
/**
* OutLink Stream Cache。
*
* 该 Cache 保持 Wechat/Wecom 的字符串拼接协议;追加和 TTL 在同一 Redis 事务中执行,
* 读取 miss 返回 undefined,删除结果由 adapter 透传。通道响应和加密编排留在调用方。
*/
export class OutLinkStreamCache {
private readonly redis: RedisCacheAdapter;
constructor({ redis = redisCacheAdapter }: OutLinkStreamCacheOptions = {}) {
this.redis = redis;
}
getKey = getOutLinkStreamKey;
private parseTtlSeconds = (ttlSeconds: number) => PositiveSafeIntegerSchema.parse(ttlSeconds);
/** 初始化或追加响应片段,并原子刷新对应 TTL。 */
append = ({
streamId,
value,
ttlSeconds
}: {
streamId: string;
value: string;
ttlSeconds: number;
}) =>
this.redis.appendStringWithTtl({
key: getOutLinkStreamKey(streamId),
value,
ttlSeconds: this.parseTtlSeconds(ttlSeconds)
});
/** 读取当前已拼接响应;Redis miss 映射为 undefined。 */
async get(streamId: string) {
const value = await this.redis.get(getOutLinkStreamKey(streamId));
return value ?? undefined;
}
/** 删除已完成响应,保留 adapter 的删除结果以便调用方按需记录。 */
delete = (streamId: string) => this.redis.delete(getOutLinkStreamKey(streamId));
}
export const outLinkStreamCache = new OutLinkStreamCache();
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import { RedisInvalidArgumentError } from '../runtime/errors';
import { z } from 'zod';
import { NonNegativeSafeIntegerSchema } from '../runtime/schema';
import type { RedisCacheLogger } from '../types';
const SESSION_KEY_PREFIX = 'session:';
export const SESSION_TTL_SECONDS = 7 * 24 * 60 * 60;
/** Session 业务数据的类型和写入校验合同。 */
export const SessionDataSchema = z.object({
userId: z.string().min(1),
teamId: z.string().min(1),
tmbId: z.string().min(1),
isRoot: z.boolean(),
createdAt: NonNegativeSafeIntegerSchema,
ip: z.string().nullable().optional()
});
export type SessionData = z.infer<typeof SessionDataSchema>;
const SessionHashSchema = z.object({
userId: z.string().min(1),
teamId: z.string().min(1),
tmbId: z.string().min(1),
isRoot: z.enum(['0', '1']).transform((value) => value === '1'),
createdAt: z.string().regex(/^\d+$/).transform(Number).pipe(NonNegativeSafeIntegerSchema),
ip: z.string().optional()
});
export type SessionCacheOptions = {
redis?: RedisCacheAdapter;
logger: RedisCacheLogger;
};
export type SessionRecord = {
sessionId: string;
data: SessionData;
};
/**
* 用户 Session Cache。
*
* Cache 保持历史 session hash key、字段编码和 7 天 TTL;Redis 只存认证态,读取
* 错误向上抛出保持 fail-closed。损坏 hash 会被记录并尽力删除,避免后续请求重复解析。
*/
export class SessionCache {
private readonly redis: RedisCacheAdapter;
private readonly logger: RedisCacheLogger;
constructor({ redis = redisCacheAdapter, logger }: SessionCacheOptions) {
this.redis = redis;
this.logger = logger;
}
private getKey = (sessionId: string) => asRedisLogicalKey(`${SESSION_KEY_PREFIX}${sessionId}`);
private getUserPrefix = (userId: string) => asRedisLogicalKey(`${SESSION_KEY_PREFIX}${userId}`);
private decodeHash = async ({
logicalKey,
sessionId,
hash
}: {
logicalKey: ReturnType<typeof asRedisLogicalKey>;
sessionId: string;
hash: Record<string, string>;
}): Promise<SessionData | undefined> => {
if (Object.keys(hash).length === 0) return;
const parsed = SessionHashSchema.safeParse(hash);
if (parsed.success) return parsed.data;
this.logger.error('Invalid Redis session record', {
sessionId,
issues: parsed.error.issues
});
await this.redis.delete(logicalKey).catch((error) => {
this.logger.warn('Failed to delete invalid Redis session record', { sessionId, error });
});
};
private readByLogicalKey = async ({
logicalKey,
sessionId
}: {
logicalKey: ReturnType<typeof asRedisLogicalKey>;
sessionId: string;
}) => {
const hash = await this.redis.getHashAll(logicalKey);
return this.decodeHash({ logicalKey, sessionId, hash });
};
/** 读取一个 session ID;miss 或损坏记录返回 undefined,Redis 错误继续向上抛出。 */
get = (sessionId: string) =>
this.readByLogicalKey({ logicalKey: this.getKey(sessionId), sessionId });
/** 原子写入历史 session hash 并设置 7 天 TTL。 */
async set({ sessionId, data }: { sessionId: string; data: SessionData }) {
const parsedData = SessionDataSchema.safeParse(data);
if (!parsedData.success) {
throw new RedisInvalidArgumentError({
operation: 'session.set',
message: 'session data is invalid'
});
}
const { userId, teamId, tmbId, isRoot, createdAt, ip } = parsedData.data;
await this.redis.setHashWithTtl({
key: this.getKey(sessionId),
fields: {
userId,
teamId,
tmbId,
isRoot: isRoot ? '1' : '0',
createdAt: String(createdAt),
...(ip === undefined || ip === null ? {} : { ip })
},
ttlSeconds: SESSION_TTL_SECONDS
});
}
/** 删除一个 session ID;调用方不需要知道物理 key。 */
async delete(sessionId: string) {
await this.redis.delete(this.getKey(sessionId));
}
/** 分页扫描某个用户的全部 typed session,损坏记录会被尽力清理。 */
async listByUser(userId: string): Promise<SessionRecord[]> {
const records: SessionRecord[] = [];
for await (const logicalKeys of this.redis.iterateByPrefix({
prefix: this.getUserPrefix(userId)
})) {
const batch = await Promise.all(
logicalKeys.map(async (logicalKey) => {
const sessionId = logicalKey.slice(SESSION_KEY_PREFIX.length);
const data = await this.readByLogicalKey({ logicalKey, sessionId });
return data ? { sessionId, data } : undefined;
})
);
records.push(...batch.filter((record): record is SessionRecord => record !== undefined));
}
return records;
}
/** 批量删除 session ID;空集合不获取 Redis connection。 */
deleteMany = (sessionIds: readonly string[]) => {
if (sessionIds.length === 0) return Promise.resolve();
return this.redis.deleteMany(sessionIds.map(this.getKey));
};
}
import {
asRedisLogicalKey,
redisCacheAdapter,
type RedisLogicalKey,
type RedisCacheAdapter
} from '../adapter';
import { PositiveSafeIntegerSchema } from '../runtime/schema';
import type { RedisCacheLogger, RedisStreamEntry } from '../types';
import { z } from 'zod';
const STREAM_RESUME_NAMESPACE = 'stream:resume';
export const StreamResumeParamsSchema = z.object({
teamId: z.string().min(1),
sourceType: z.string().min(1),
sourceId: z.string().min(1),
chatId: z.string().min(1)
});
export type StreamResumeParams = z.infer<typeof StreamResumeParamsSchema>;
export const StreamResumeUnavailableStateSchema = z.object({
reason: z.string().min(1)
});
export type StreamResumeUnavailableState = z.infer<typeof StreamResumeUnavailableStateSchema>;
export const StreamResumeActiveStateSchema = z.object({
updatedAt: PositiveSafeIntegerSchema
});
export type StreamResumeActiveState = z.infer<typeof StreamResumeActiveStateSchema>;
export type StreamResumeKeys = {
keyOfStream: RedisLogicalKey;
keyOfUnavailable: RedisLogicalKey;
keyOfActive: RedisLogicalKey;
};
export type StreamResumeCacheOptions = {
redis?: RedisCacheAdapter;
logger: RedisCacheLogger<'error'>;
streamTtlSeconds: number;
postCompleteTtlSeconds: number;
ttlTouchIntervalMs: number;
};
type StreamResumeBlockingReader = ReturnType<RedisCacheAdapter['createBlockingStreamReader']>;
const parsePositiveConfig = ({
operation,
field,
value
}: {
operation: string;
field: string;
value: number;
}) => {
const parsed = PositiveSafeIntegerSchema.safeParse(value);
if (!parsed.success) {
throw new Error(`${operation}.${field} must be a positive safe integer`);
}
return parsed.data;
};
/**
* Stream Resume Cache。
*
* Cache 固定历史 stream/state key 和 TTL,负责镜像写入的顺序、Stream 返回解析以及
* blocking reader 生命周期;HTTP/SSE response 和终止事件由 service 层继续编排。
*/
export class StreamResumeCache {
private readonly redis: RedisCacheAdapter;
private readonly logger: RedisCacheLogger<'error'>;
private readonly parsedStreamTtlSeconds: number;
private readonly parsedPostCompleteTtlSeconds: number;
private readonly parsedTtlTouchIntervalMs: number;
constructor({
redis = redisCacheAdapter,
logger,
streamTtlSeconds,
postCompleteTtlSeconds,
ttlTouchIntervalMs
}: StreamResumeCacheOptions) {
this.redis = redis;
this.logger = logger;
this.parsedStreamTtlSeconds = parsePositiveConfig({
operation: 'streamResume',
field: 'streamTtlSeconds',
value: streamTtlSeconds
});
this.parsedPostCompleteTtlSeconds = parsePositiveConfig({
operation: 'streamResume',
field: 'postCompleteTtlSeconds',
value: postCompleteTtlSeconds
});
this.parsedTtlTouchIntervalMs = parsePositiveConfig({
operation: 'streamResume',
field: 'ttlTouchIntervalMs',
value: ttlTouchIntervalMs
});
}
private parseParams = (params: StreamResumeParams): StreamResumeParams =>
StreamResumeParamsSchema.parse(params);
getKeys = (params: StreamResumeParams): StreamResumeKeys => {
const parsed = this.parseParams(params);
const { teamId, sourceType, sourceId, chatId } = parsed;
return {
keyOfStream: asRedisLogicalKey(
`${STREAM_RESUME_NAMESPACE}:data:${teamId}:${sourceType}:${sourceId}:${chatId}`
),
keyOfUnavailable: asRedisLogicalKey(
`${STREAM_RESUME_NAMESPACE}:unavailable:${teamId}:${sourceType}:${sourceId}:${chatId}`
),
keyOfActive: asRedisLogicalKey(
`${STREAM_RESUME_NAMESPACE}:active:${teamId}:${sourceType}:${sourceId}:${chatId}`
)
};
};
private touchState = async (keys: StreamResumeKeys) => {
await Promise.all([
this.redis.expireStream({ key: keys.keyOfStream, ttlSeconds: this.parsedStreamTtlSeconds }),
this.redis.set({
key: keys.keyOfActive,
value: JSON.stringify({ updatedAt: Date.now() } satisfies StreamResumeActiveState),
ttlMs: this.parsedStreamTtlSeconds * 1000
})
]);
};
private clearMirror = async (keys: StreamResumeKeys) => {
await Promise.all([
this.redis.delete(keys.keyOfUnavailable),
this.redis.delete(keys.keyOfStream),
this.redis.delete(keys.keyOfActive)
]);
};
/** 持久化当前请求无法创建镜像的原因。 */
async setUnavailable(params: StreamResumeParams, state: StreamResumeUnavailableState) {
const parsedState = StreamResumeUnavailableStateSchema.parse(state);
await this.redis.set({
key: this.getKeys(params).keyOfUnavailable,
value: JSON.stringify(parsedState),
ttlMs: this.parsedStreamTtlSeconds * 1000
});
}
/** 读取 unavailable 状态;miss 由调用方解释为可继续读取 Stream。 */
async getUnavailable(params: StreamResumeParams) {
const value = await this.redis.get(this.getKeys(params).keyOfUnavailable);
if (!value) return;
try {
const parsed = StreamResumeUnavailableStateSchema.safeParse(JSON.parse(value));
return parsed.success ? parsed.data : undefined;
} catch {
return;
}
}
/** 读取 active 状态;损坏 JSON 按 miss 处理。 */
async getActive(params: StreamResumeParams) {
const value = await this.redis.get(this.getKeys(params).keyOfActive);
if (!value) return;
try {
const parsed = StreamResumeActiveStateSchema.safeParse(JSON.parse(value));
return parsed.success ? parsed.data : undefined;
} catch {
return;
}
}
/** 读取 Redis 内存水位;是否阻止创建镜像由 service 的运行策略决定。 */
getMemoryInfo = () => this.redis.getMemoryInfo();
/** 创建一个顺序写入镜像;Redis 写入失败只记录日志并保持后续 flush 可完成。 */
createMirror(params: StreamResumeParams) {
const parsedParams = this.parseParams(params);
const keys = this.getKeys(parsedParams);
let queue: Promise<void> = this.clearMirror(keys).catch((error) => {
this.logger.error('Failed to clear stream resume redis keys before mirror', {
params: parsedParams,
error
});
});
let lastTouchedAt = 0;
const enqueueRaw = (raw: string) => {
queue = queue
.then(async () => {
await this.redis.appendStreamEntry({
key: keys.keyOfStream,
fields: { raw }
});
const now = Date.now();
if (lastTouchedAt === 0 || now - lastTouchedAt >= this.parsedTtlTouchIntervalMs) {
await this.touchState(keys);
lastTouchedAt = now;
}
})
.catch((error) => {
this.logger.error('Failed to mirror stream response to redis', {
params: parsedParams,
error
});
});
return queue;
};
return {
...keys,
enqueueRaw,
flush: async () => {
await queue;
},
shrinkTTLAfterComplete: async () => {
try {
await Promise.all([
this.redis.expireStream({
key: keys.keyOfStream,
ttlSeconds: this.parsedPostCompleteTtlSeconds
}),
this.redis.expireStream({
key: keys.keyOfActive,
ttlSeconds: this.parsedPostCompleteTtlSeconds
})
]);
} catch (error) {
this.logger.error('Failed to shrink stream resume redis ttl', {
params: parsedParams,
error
});
}
}
};
}
/** 读取 history;返回值已脱离 Redis 的交替数组协议。 */
async range({
params,
start,
end,
count
}: {
params: StreamResumeParams;
start: string;
end: string;
count: number;
}): Promise<RedisStreamEntry[]> {
return this.redis.rangeStream({ key: this.getKeys(params).keyOfStream, start, end, count });
}
/**
* 在 DAL 内运行请求级 blocking reader,并保证无论读取循环如何结束都会释放连接。
* callback 只接收 typed reader,不会获得 raw ioredis client。
*/
async withBlockingReader<T>({
params,
blockMs,
count,
callback
}: {
params: StreamResumeParams;
blockMs: number;
count?: number;
callback: (reader: StreamResumeBlockingReader) => Promise<T> | T;
}): Promise<T> {
const reader = this.redis.createBlockingStreamReader({
key: this.getKeys(params).keyOfStream,
blockMs,
count
});
try {
return await callback(reader);
} finally {
await reader.close();
}
}
}
import { randomUUID } from 'node:crypto';
import {
asRedisLogicalKey,
redisCacheAdapter,
type RedisLogicalKey,
type RedisCacheAdapter
} from '../adapter';
const SYSTEM_VERSION_PREFIX = 'VERSION_KEY:';
const SYSTEM_VERSION_SCAN_BATCH_SIZE = 100;
export type SystemVersionCacheOptions = {
redis?: RedisCacheAdapter;
createVersion?: () => string;
};
/**
* System Version Cache。
*
* Cache 保持历史永久 key 和 UUID value;首次读取使用单条 SET NX GET 原子初始化。
* wildcard refresh 会扫描并删除指定 base key 下的全部子 key,但不会删除 base key 本身。
* Redis 是版本一致性的事实来源,所有错误均向上传播。
*/
export class SystemVersionCache {
private readonly redis: RedisCacheAdapter;
private readonly createVersion: () => string;
constructor({
redis = redisCacheAdapter,
createVersion = randomUUID
}: SystemVersionCacheOptions = {}) {
this.redis = redis;
this.createVersion = createVersion;
}
private getBaseKey = (key: string) => asRedisLogicalKey(`${SYSTEM_VERSION_PREFIX}${key}`);
private getKey = ({ key, id }: { key: string; id?: string }) =>
id ? asRedisLogicalKey(`${this.getBaseKey(key)}:${id}`) : this.getBaseKey(key);
/** 返回已有版本;key 不存在时原子写入并返回新的永久版本。 */
getOrInitialize = ({ key, id }: { key: string; id?: string }) =>
this.redis.getOrSet({
key: this.getKey({ key, id }),
value: this.createVersion()
});
/** 刷新单个版本,或在 id='*' 时只删除该 base key 下的全部子版本。 */
async refresh({ key, id }: { key: string; id?: string | '*' }) {
if (id !== '*') {
await this.redis.set({
key: this.getKey({ key, id }),
value: this.createVersion()
});
return;
}
// 先完成遍历再删除,避免修改 keyspace 导致 SCAN 游标漏过尚未返回的子 key。
const childKeys: RedisLogicalKey[] = [];
for await (const keys of this.redis.iterateByPrefix({
prefix: this.getBaseKey(key),
batchSize: SYSTEM_VERSION_SCAN_BATCH_SIZE
})) {
childKeys.push(...keys);
}
if (childKeys.length > 0) {
await this.redis.deleteMany(childKeys);
}
}
}
export const systemVersionCache = new SystemVersionCache();
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import { FiniteNumberSchema } from '../runtime/schema';
import type { RedisCacheLogger } from '../types';
import { z } from 'zod';
const TEAM_POINT_CACHE_TTL_MS = 60 * 1000;
const TEAM_POINT_CACHE_TTL_SECONDS = 60;
const TEAM_POINT_SURPLUS_KEY_PREFIX = 'cache:team_point_surplus:';
const TEAM_POINT_TOTAL_KEY_PREFIX = 'cache:team_point_total:';
const TeamPointCacheValueSchema = z
.string()
.refine((value) => value.trim() !== '', { error: 'cache value must not be empty' })
.transform(Number)
.pipe(FiniteNumberSchema);
export type TeamPointCacheOptions = {
redis?: RedisCacheAdapter;
logger?: RedisCacheLogger<'warn'>;
};
export type TeamPointSnapshot = {
totalPoints: number;
surplusPoints: number;
};
const noopLogger: RedisCacheLogger<'warn'> = {
warn: () => undefined
};
/**
* 团队积分双 key Cache。
*
* 两个积分值必须作为同一版本成对读取和刷新;Redis 只承担加速,读取异常或 partial hit
* 返回 miss 让 wallet service 回源 Mongo。写入、增量和清理失败只记录 warning,不覆盖钱包主流程。
*/
export class TeamPointCache {
private readonly redis: RedisCacheAdapter;
private readonly logger: RedisCacheLogger<'warn'>;
constructor({ redis = redisCacheAdapter, logger = noopLogger }: TeamPointCacheOptions = {}) {
this.redis = redis;
this.logger = logger;
}
private getKeys = (teamId: string) => ({
surplus: asRedisLogicalKey(`${TEAM_POINT_SURPLUS_KEY_PREFIX}${teamId}`),
total: asRedisLogicalKey(`${TEAM_POINT_TOTAL_KEY_PREFIX}${teamId}`)
});
private parsePoint = (value: string | null) => {
if (value === null) return;
const parsed = TeamPointCacheValueSchema.safeParse(value);
return parsed.success ? parsed.data : undefined;
};
private parsePointInput = (value: number, field: string) => {
const parsed = FiniteNumberSchema.safeParse(value);
if (parsed.success) return parsed.data;
this.logger.warn('Skipped invalid team point cache value', { field, value });
};
/** 成对读取积分缓存;partial hit、损坏值和 Redis 故障均降级为 miss。 */
async get(teamId: string): Promise<TeamPointSnapshot | undefined> {
try {
const keys = this.getKeys(teamId);
const [surplusValue, totalValue] = await this.redis.getPair({
first: keys.surplus,
second: keys.total
});
const surplusPoints = this.parsePoint(surplusValue);
const totalPoints = this.parsePoint(totalValue);
if (surplusPoints === undefined || totalPoints === undefined) {
return undefined;
}
return { totalPoints, surplusPoints };
} catch (error) {
this.logger.warn('Failed to get team point cache', { teamId, error });
return undefined;
}
}
/** 原子刷新两个积分 key;Redis 故障只记录降级日志。 */
async set({
teamId,
totalPoints,
surplusPoints
}: {
teamId: string;
totalPoints: number;
surplusPoints: number;
}) {
const parsedTotalPoints = this.parsePointInput(totalPoints, 'totalPoints');
const parsedSurplusPoints = this.parsePointInput(surplusPoints, 'surplusPoints');
if (parsedTotalPoints === undefined || parsedSurplusPoints === undefined) return;
try {
const keys = this.getKeys(teamId);
await this.redis.setPair({
first: { key: keys.surplus, value: String(parsedSurplusPoints) },
second: { key: keys.total, value: String(parsedTotalPoints) },
ttlMs: TEAM_POINT_CACHE_TTL_MS
});
} catch (error) {
this.logger.warn('Failed to set team point cache', { teamId, error });
}
}
/** 原子增加 surplus;0 增量不会创建或刷新 cache key。 */
async incrementSurplus({ teamId, value }: { teamId: string; value: number }) {
const parsedValue = this.parsePointInput(value, 'value');
if (parsedValue === undefined) return;
if (parsedValue === 0) return;
try {
await this.redis.incrementWithTtl({
key: this.getKeys(teamId).surplus,
increment: parsedValue,
ttlSeconds: TEAM_POINT_CACHE_TTL_SECONDS
});
} catch (error) {
this.logger.warn('Failed to increment team point cache', {
teamId,
value: parsedValue,
error
});
}
}
/** 单条 multi-key DEL 清理两个积分 key;缓存故障不覆盖钱包主流程。 */
async clear(teamId: string) {
try {
const keys = this.getKeys(teamId);
await this.redis.deleteMany([keys.surplus, keys.total]);
} catch (error) {
this.logger.warn('Failed to clear team point cache', { teamId, error });
}
}
}
export const teamPointCache = new TeamPointCache();
import {
asRedisLogicalKey,
redisCacheAdapter,
RedisInvalidArgumentError,
type RedisCacheAdapter
} from '../adapter';
import { PositiveSafeIntegerSchema } from '../runtime/schema';
import { z } from 'zod';
const TEAM_QPM_CACHE_TTL_MS = 60 * 60 * 1000;
const TEAM_QPM_KEY_PREFIX = 'cache:team_qpm_limit:';
const TeamQpmCacheValueSchema = z
.string()
.regex(/^[1-9]\d*$/)
.transform(Number)
.pipe(PositiveSafeIntegerSchema);
export type TeamQpmCacheOptions = {
redis?: RedisCacheAdapter;
};
/**
* Team QPM 配置 Cache。
*
* 只拥有 Redis cache key、字符串 codec、TTL 和清理语义;套餐回源与默认值由 wallet service
* 决定。损坏缓存按 miss 处理,避免把 NaN 当成“不限流”。
*/
export class TeamQpmCache {
private readonly redis: RedisCacheAdapter;
constructor({ redis = redisCacheAdapter }: TeamQpmCacheOptions = {}) {
this.redis = redis;
}
private getKey = (teamId: string) => asRedisLogicalKey(`${TEAM_QPM_KEY_PREFIX}${teamId}`);
async getCachedLimit(teamId: string): Promise<number | null> {
const cached = await this.redis.get(this.getKey(teamId));
if (cached === null) return null;
const parsed = TeamQpmCacheValueSchema.safeParse(cached);
return parsed.success ? parsed.data : null;
}
async setCachedLimit({ teamId, limit }: { teamId: string; limit: number }) {
const parsedLimit = PositiveSafeIntegerSchema.safeParse(limit);
if (!parsedLimit.success) {
throw new RedisInvalidArgumentError({
operation: 'teamQpm.set',
message: 'limit must be a positive safe integer'
});
}
await this.redis.set({
key: this.getKey(teamId),
value: String(parsedLimit.data),
ttlMs: TEAM_QPM_CACHE_TTL_MS
});
}
async clearCachedLimit(teamId: string) {
await this.redis.delete(this.getKey(teamId));
}
}
export const teamQpmCache = new TeamQpmCache();
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import type { RedisCacheLogger } from '../types';
const TEAM_VECTOR_COUNT_CACHE_TTL_MS = 30 * 60 * 1000;
const TEAM_VECTOR_CACHE_OPERATION_TIMEOUT_MS = 3000;
export type TeamVectorCountCacheOptions = {
redis?: RedisCacheAdapter;
logger: RedisCacheLogger<'warn'>;
};
/**
* 团队向量数量 Cache。
*
* 所有 Redis 操作都有独立 3 秒 deadline;读取失败按 miss 回源,写入和失效失败只记录
* 日志。Cache 不在业务层叠加 legacy cache helper 的通用重试。
*/
export class TeamVectorCountCache {
private readonly redis: RedisCacheAdapter;
private readonly logger: RedisCacheLogger<'warn'>;
constructor({ redis = redisCacheAdapter, logger }: TeamVectorCountCacheOptions) {
this.redis = redis;
this.logger = logger;
}
private getKey = (teamId: string) => asRedisLogicalKey(`cache:team_vector_count:${teamId}`);
private runWithTimeout = async <T>({
promise,
timeoutMessage
}: {
promise: Promise<T>;
timeoutMessage: string;
}): Promise<T> => {
let timer: ReturnType<typeof setTimeout> | undefined;
try {
return await Promise.race([
promise,
new Promise<never>((_, reject) => {
timer = setTimeout(
() => reject(new Error(timeoutMessage)),
TEAM_VECTOR_CACHE_OPERATION_TIMEOUT_MS
);
})
]);
} finally {
if (timer) clearTimeout(timer);
}
};
private runOperation = async <T>({
teamId,
operation,
warnMessage,
action
}: {
teamId: string;
operation: string;
warnMessage: string;
action: () => Promise<T>;
}) => {
try {
return await this.runWithTimeout({
promise: action(),
timeoutMessage: `${operation} timed out after ${TEAM_VECTOR_CACHE_OPERATION_TIMEOUT_MS}ms`
});
} catch (error) {
this.logger.warn(warnMessage, { teamId, error });
return undefined;
}
};
/** 读取团队向量数量;miss、错误和超时统一返回 undefined 触发 VectorDB 回源。 */
async get(teamId: string) {
const count = await this.runOperation({
teamId,
operation: 'Get team vector count cache',
warnMessage: 'Failed to get team vector count cache',
action: () => this.redis.get(this.getKey(teamId))
});
if (count === null || count === undefined || count.trim().length === 0) return undefined;
if (!/^\d+$/.test(count)) return undefined;
const parsedCount = Number(count);
return Number.isSafeInteger(parsedCount) && parsedCount >= 0 ? parsedCount : undefined;
}
/** best-effort 写入缓存;调用方无需等待该结果才能返回 VectorDB 主结果。 */
async set({ teamId, count }: { teamId: string; count: number }) {
await this.runOperation({
teamId,
operation: 'Set team vector count cache',
warnMessage: 'Failed to set team vector count cache',
action: () =>
this.redis.set({
key: this.getKey(teamId),
value: String(count),
ttlMs: TEAM_VECTOR_COUNT_CACHE_TTL_MS
})
});
}
/** best-effort 失效缓存;Redis 故障不得覆盖 VectorDB 写入或删除结果。 */
async invalidate(teamId: string) {
await this.runOperation({
teamId,
operation: 'Invalidate team vector count cache',
warnMessage: 'Failed to invalidate team vector count cache',
action: () => this.redis.delete(this.getKey(teamId))
});
}
}
import { z } from 'zod';
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
export const WECHAT_POLLING_FAILURE_TTL_SECONDS = 300;
const WECHAT_POLLING_FAILURE_NAMESPACE = 'cache:wechat:publish:failures';
const ShareIdSchema = z.string().min(1);
export type WechatPollingFailureCacheOptions = {
redis?: RedisCacheAdapter;
};
/** 构造 Wechat polling failure counter 的逻辑 key,物理前缀由 Redis adapter 统一添加。 */
export const getWechatPollingFailureKey = (shareId: string) =>
asRedisLogicalKey(`${WECHAT_POLLING_FAILURE_NAMESPACE}:${ShareIdSchema.parse(shareId)}`);
/**
* Wechat polling failure counter Cache。
*
* 计数递增由 adapter 以 INCRBY + EXPIRE NX 事务完成,避免多个 worker 读取后写回时丢失更新。
* Redis 错误不在 Cache 内吞掉,worker 继续沿用 failed/退避语义。
*/
export class WechatPollingFailureCache {
private readonly redis: RedisCacheAdapter;
constructor({ redis = redisCacheAdapter }: WechatPollingFailureCacheOptions = {}) {
this.redis = redis;
}
getKey = getWechatPollingFailureKey;
/** 原子递增连续失败次数;首次递增建立 300 秒 TTL,后续递增保留原 TTL。 */
increment = (shareId: string) =>
this.redis.incrementIntegerWithTtl({
key: getWechatPollingFailureKey(shareId),
increment: 1,
ttlSeconds: WECHAT_POLLING_FAILURE_TTL_SECONDS
});
/** 成功轮询后将计数归零并刷新 300 秒 TTL,保持历史 value 合同。 */
reset = (shareId: string) =>
this.redis.set({
key: getWechatPollingFailureKey(shareId),
value: '0',
ttlMs: WECHAT_POLLING_FAILURE_TTL_SECONDS * 1000
});
/** 达到阈值后删除计数 key;删除结果由调用方按需处理。 */
clear = (shareId: string) => this.redis.delete(getWechatPollingFailureKey(shareId));
}
export const wechatPollingFailureCache = new WechatPollingFailureCache();
import { z } from 'zod';
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import { RedisInvalidArgumentError, RedisInvalidResponseError } from '../runtime/errors';
export const WECHAT_QR_LOGIN_TTL_SECONDS = 8 * 60;
const WechatQrLoginDataSchema = z.looseObject({
qrcode: z.string().min(1),
qrcode_img_content: z.string().min(1)
});
export type WechatQrLoginData = z.infer<typeof WechatQrLoginDataSchema>;
export type WechatQrLoginCacheOptions = {
redis?: RedisCacheAdapter;
};
type WechatQrLoginKey = {
outLinkId: string;
tmbId: string;
};
/**
* 微信二维码登录 Cache。
*
* Cache 保持历史物理 key、JSON value 和 480 秒 TTL。Redis miss 返回 undefined;
* Redis 操作错误和损坏的缓存数据均 fail-closed,由应用边界处理。
*/
export class WechatQrLoginCache {
private readonly redis: RedisCacheAdapter;
constructor({ redis = redisCacheAdapter }: WechatQrLoginCacheOptions = {}) {
this.redis = redis;
}
private getKey = ({ outLinkId, tmbId }: WechatQrLoginKey) =>
asRedisLogicalKey(`cache:publish:wechat:qrcode:${outLinkId}:${tmbId}`);
/** 校验并保存 iLink 二维码响应;上游响应不合法时不会写入损坏缓存。 */
async set({
outLinkId,
tmbId,
data
}: WechatQrLoginKey & {
data: WechatQrLoginData;
}) {
const parsed = WechatQrLoginDataSchema.safeParse(data);
if (!parsed.success) {
throw new RedisInvalidArgumentError({
operation: 'wechatQrLogin.set',
message: 'Wechat QR login data is invalid'
});
}
await this.redis.set({
key: this.getKey({ outLinkId, tmbId }),
value: JSON.stringify(parsed.data),
ttlMs: WECHAT_QR_LOGIN_TTL_SECONDS * 1000
});
}
/** 读取并校验缓存;只有 Redis 中不存在该 key 时才返回 undefined。 */
async get({ outLinkId, tmbId }: WechatQrLoginKey) {
const raw = await this.redis.get(this.getKey({ outLinkId, tmbId }));
if (raw === null) return undefined;
const data = (() => {
try {
return JSON.parse(raw) as unknown;
} catch {
throw new RedisInvalidResponseError({
operation: 'wechatQrLogin.get',
message: 'Wechat QR login cache contains invalid JSON'
});
}
})();
const parsed = WechatQrLoginDataSchema.safeParse(data);
if (!parsed.success) {
throw new RedisInvalidResponseError({
operation: 'wechatQrLogin.get',
message: 'Wechat QR login cache contains invalid data'
});
}
return parsed.data;
}
/** 删除已确认的二维码状态;删除失败继续向上抛出,保持 fail-closed 语义。 */
delete = ({ outLinkId, tmbId }: WechatQrLoginKey) =>
this.redis.delete(this.getKey({ outLinkId, tmbId }));
}
export const wechatQrLoginCache = new WechatQrLoginCache();
import { z } from 'zod';
import { asRedisLogicalKey, redisCacheAdapter, type RedisCacheAdapter } from '../adapter';
import { RedisInvalidArgumentError } from '../runtime/errors';
import type { RedisCacheLogger } from '../types';
const WORKFLOW_STOP_SIGNAL_PREFIX = 'agent_runtime_stopping';
export const WORKFLOW_STOP_SIGNAL_TTL_SECONDS = 60;
export const WorkflowStopSignalParamsSchema = z.object({
sourceType: z.string().min(1),
sourceId: z.string().min(1),
chatId: z.string().min(1)
});
export type WorkflowStopSignalParams = z.infer<typeof WorkflowStopSignalParamsSchema>;
export type WorkflowStopSignalCacheOptions = {
redis?: RedisCacheAdapter;
logger: RedisCacheLogger<'warn'>;
};
/** 构造运行态停止标记 logical key;物理前缀由 Redis adapter 统一添加。 */
export const getWorkflowStopSignalKey = ({
sourceType,
sourceId,
chatId
}: WorkflowStopSignalParams) => {
const parsed = WorkflowStopSignalParamsSchema.safeParse({ sourceType, sourceId, chatId });
if (!parsed.success) {
throw new RedisInvalidArgumentError({
operation: 'workflowStopSignal.key',
message: 'sourceType, sourceId and chatId must be non-empty strings'
});
}
return asRedisLogicalKey(
`${WORKFLOW_STOP_SIGNAL_PREFIX}:${parsed.data.sourceType}:${parsed.data.sourceId}:${parsed.data.chatId}`
);
};
/**
* workflow 与 auxiliary generation 共用的停止信号 Cache。
*
* 写入故障由调用方处理为 fail-closed;读取故障按 false 降级;清理为 best-effort,避免
* Redis 故障覆盖工作流或辅助生成的最终结果。
*/
export class WorkflowStopSignalCache {
private readonly redis: RedisCacheAdapter;
private readonly logger: RedisCacheLogger<'warn'>;
constructor({ redis = redisCacheAdapter, logger }: WorkflowStopSignalCacheOptions) {
this.redis = redis;
this.logger = logger;
}
/** 设置 60 秒停止标记;Redis 错误继续向上抛出。 */
async set(params: WorkflowStopSignalParams) {
await this.redis.set({
key: getWorkflowStopSignalKey(params),
value: '1',
ttlMs: WORKFLOW_STOP_SIGNAL_TTL_SECONDS * 1000
});
}
/** 检查停止标记;Redis 读取故障按 false 降级。 */
async isStopping(params: WorkflowStopSignalParams) {
const key = getWorkflowStopSignalKey(params);
try {
const result = await this.redis.get(key);
return result !== null;
} catch (error) {
this.logger.warn('Workflow stop signal read failed open', { key, error });
return false;
}
}
/** 删除停止标记;清理失败只记录 warning,不覆盖业务结果。 */
async clear(params: WorkflowStopSignalParams) {
const key = getWorkflowStopSignalKey(params);
try {
await this.redis.delete(key);
} catch (error) {
this.logger.warn('Workflow stop signal clear failed', { key, error });
}
}
}
import { getRedisRuntime } from './runtime/connection';
export { closeRedisRuntime, configureRedisRuntime } from './runtime/connection';
export type { RedisRuntimeLogger, RedisRuntimeOptions } from './runtime/connection';
export type {
RedisCacheLogger,
RedisLogMetadata,
RedisLogMethod,
RedisRuntimeHealthMetric,
RedisRuntimeMetrics,
RedisRuntimeShutdownMetric
} from './types';
export { RedisConfigurationError, parseRedisConnectionConfig } from './runtime/config';
export type { RedisConnectionConfig } from './runtime/config';
export {
isRedisOperationError,
RedisInvalidArgumentError,
RedisInvalidResponseError,
RedisOperationError,
RedisOperationExecutionError,
RedisOperationTimeoutError
} from './runtime/errors';
export type {
RedisOperationErrorCode,
RedisOperationOutcome,
RedisOperationRole
} from './runtime/errors';
export { asRedisLogicalKey, createRedisLogicalKey } from './runtime/keyspace';
export type { RedisLogicalKey } from './runtime/keyspace';
/** 检查已经由应用配置的默认 Redis Runtime,不在 DAL 内隐式读取环境变量。 */
export const checkRedisHealth = () => getRedisRuntime().checkHealth();
import type { RedisOptions } from 'ioredis';
import { z } from 'zod';
export type RedisEndpoint = {
transport: 'tcp' | 'unix';
host?: string;
port?: number;
path?: string;
db?: number;
tls: boolean;
hasUsername: boolean;
hasPassword: boolean;
};
export type RedisConnectionConfig = {
options: RedisOptions;
endpoint: RedisEndpoint;
};
export class RedisConfigurationError extends Error {
constructor(message: string) {
super(message);
this.name = 'RedisConfigurationError';
}
}
const RedisUrlInputSchema = z
.string({ error: 'REDIS_URL must be a string' })
.trim()
.min(1, { error: 'REDIS_URL must not be empty' })
.refine((value) => !value.includes('?') && !value.includes('#'), {
error: 'REDIS_URL query parameters and fragments are not supported'
});
const RedisDbPathSchema = z
.string()
.regex(/^\d+$/, { error: 'REDIS_URL database must be a non-negative integer' })
.transform(Number)
.pipe(
z
.number({ error: 'REDIS_URL database is outside the supported integer range' })
.max(Number.MAX_SAFE_INTEGER, {
error: 'REDIS_URL database is outside the supported integer range'
})
);
const RedisTcpUrlSchema = z.object({
protocol: z.enum(['redis:', 'rediss:'], {
error: 'REDIS_URL protocol must be redis or rediss'
}),
hostname: z.string().min(1, { error: 'REDIS_URL host must not be empty' }),
port: z
.number({ error: 'REDIS_URL port must be between 1 and 65535' })
.int({ error: 'REDIS_URL port must be between 1 and 65535' })
.min(1, { error: 'REDIS_URL port must be between 1 and 65535' })
.max(65535, { error: 'REDIS_URL port must be between 1 and 65535' })
});
const parseConfigSchema = <T>(schema: z.ZodType<T>, value: unknown): T => {
const result = schema.safeParse(value);
if (!result.success) {
throw new RedisConfigurationError(result.error.issues[0]!.message);
}
return result.data;
};
const parseRedisDb = (pathname: string) => {
const dbPath = pathname.replace(/^\//, '');
if (!dbPath) return;
return parseConfigSchema(RedisDbPathSchema, dbPath);
};
const decodeCredential = (value: string, field: 'username' | 'password') => {
try {
return decodeURIComponent(value);
} catch {
throw new RedisConfigurationError(`REDIS_URL ${field} is not valid percent-encoded text`);
}
};
/**
* 解析 standalone Redis 连接配置。
*
* 兼容现有的无协议地址和 Unix socket,但拒绝未知协议、query/hash 以及非法 db。
* 抛出的错误不会包含原始 URL,避免账号密码进入启动日志。
*/
export const parseRedisConnectionConfig = (input: string): RedisConnectionConfig => {
const redisUrl = parseConfigSchema(RedisUrlInputSchema, input);
if (redisUrl.startsWith('/')) {
return {
options: { path: redisUrl },
endpoint: {
transport: 'unix',
path: redisUrl,
tls: false,
hasUsername: false,
hasPassword: false
}
};
}
const normalizedRedisUrl = redisUrl.includes('://') ? redisUrl : `redis://${redisUrl}`;
let parsedUrl: URL;
try {
parsedUrl = new URL(normalizedRedisUrl);
} catch {
throw new RedisConfigurationError('REDIS_URL is not a valid Redis connection URL');
}
const db = parseRedisDb(parsedUrl.pathname);
const { protocol, hostname, port } = parseConfigSchema(RedisTcpUrlSchema, {
protocol: parsedUrl.protocol.toLowerCase(),
hostname: parsedUrl.hostname,
port: parsedUrl.port ? Number(parsedUrl.port) : 6379
});
const tls = protocol === 'rediss:';
const host = hostname.replace(/^\[|\]$/g, '');
const options: RedisOptions = {
host,
port
};
if (parsedUrl.username) {
options.username = decodeCredential(parsedUrl.username, 'username');
}
if (parsedUrl.password) {
options.password = decodeCredential(parsedUrl.password, 'password');
}
if (db !== undefined) {
options.db = db;
}
if (tls) {
options.tls = {};
}
return {
options,
endpoint: {
transport: 'tcp',
host,
port,
db,
tls,
hasUsername: parsedUrl.username.length > 0,
hasPassword: parsedUrl.password.length > 0
}
};
};
export type RedisOperationRole = 'command';
export type RedisOperationErrorCode =
| 'REDIS_INVALID_ARGUMENT'
| 'REDIS_INVALID_RESPONSE'
| 'REDIS_OPERATION_FAILED'
| 'REDIS_OPERATION_TIMEOUT';
export type RedisOperationOutcome = 'not-started' | 'failed' | 'unknown';
/** Redis adapter 对外暴露的稳定 operation 错误基类,不包含完整 key 或业务数据。 */
export class RedisOperationError extends Error {
readonly code: RedisOperationErrorCode;
readonly operation: string;
readonly role: RedisOperationRole;
readonly outcome: RedisOperationOutcome;
override readonly cause?: unknown;
constructor({
// name,
message,
code,
operation,
role,
outcome,
cause
}: {
message: string;
code: RedisOperationErrorCode;
operation: string;
role: RedisOperationRole;
outcome: RedisOperationOutcome;
cause?: unknown;
}) {
super(message);
this.name = new.target.name;
this.code = code;
this.operation = operation;
this.role = role;
this.outcome = outcome;
this.cause = cause;
}
}
export class RedisInvalidArgumentError extends RedisOperationError {
constructor({ operation, message }: { operation: string; message: string }) {
super({
message,
code: 'REDIS_INVALID_ARGUMENT',
operation,
role: 'command',
outcome: 'not-started'
});
}
}
export class RedisInvalidResponseError extends RedisOperationError {
constructor({ operation, message }: { operation: string; message: string }) {
super({
message,
code: 'REDIS_INVALID_RESPONSE',
operation,
role: 'command',
outcome: 'failed'
});
}
}
export class RedisOperationExecutionError extends RedisOperationError {
readonly attempt: number;
constructor({
operation,
attempt,
outcome,
cause
}: {
operation: string;
attempt: number;
outcome: Exclude<RedisOperationOutcome, 'not-started'>;
cause: unknown;
}) {
super({
message: `Redis operation ${operation} failed`,
code: 'REDIS_OPERATION_FAILED',
operation,
role: 'command',
outcome,
cause
});
this.attempt = attempt;
}
}
export class RedisOperationTimeoutError extends RedisOperationError {
readonly timeoutMs: number;
readonly attempt: number;
constructor({
operation,
timeoutMs,
attempt,
outcome
}: {
operation: string;
timeoutMs: number;
attempt: number;
outcome: Exclude<RedisOperationOutcome, 'not-started'>;
}) {
super({
message: `Redis operation ${operation} timed out`,
code: 'REDIS_OPERATION_TIMEOUT',
operation,
role: 'command',
outcome
});
this.timeoutMs = timeoutMs;
this.attempt = attempt;
}
}
export const isRedisOperationError = (error: unknown): error is RedisOperationError => {
return error instanceof RedisOperationError;
};
export {
closeRedisRuntime,
configureRedisRuntime,
RedisRuntime,
getConfiguredRedisRuntime,
getRedisRuntime
} from './connection';
export type {
RedisBeforeCloseHook,
RedisClient,
RedisClientFactory,
RedisConnectionRole,
RedisConnectionSnapshot,
RedisConnectionState,
RedisEndpoint,
RedisRuntimeLogger,
RedisRuntimeOptions
} from './connection';
export type {
RedisRuntimeHealthMetric,
RedisRuntimeMetrics,
RedisRuntimeShutdownMetric
} from '../types';
export { registerRedisRuntimeShutdown } from './shutdown';
export type { RedisRuntimeShutdownOptions } from './shutdown';
export { RedisConfigurationError, parseRedisConnectionConfig } from './config';
export type { RedisConnectionConfig } from './config';
export {
FiniteNumberSchema,
NonNegativeSafeIntegerSchema,
PositiveSafeIntegerSchema
} from './schema';
export {
isRedisOperationError,
RedisInvalidArgumentError,
RedisInvalidResponseError,
RedisOperationError,
RedisOperationExecutionError,
RedisOperationTimeoutError
} from './errors';
export type { RedisOperationErrorCode, RedisOperationOutcome, RedisOperationRole } from './errors';
export {
asRedisLogicalKey,
createChildRedisScanPattern,
createRedisLogicalKey,
FASTGPT_REDIS_PREFIX,
toLogicalRedisKey,
toPhysicalRedisKey
} from './keyspace';
export type { RedisLogicalKey, RedisPhysicalKey } from './keyspace';
import type { RedisLogicalKey, RedisPhysicalKey } from '../types';
export const FASTGPT_REDIS_PREFIX = 'fastgpt:';
export type { RedisLogicalKey, RedisPhysicalKey } from '../types';
type RedisKeySegment = string | number;
const namespacePattern = /^[A-Za-z0-9_-]+(?::[A-Za-z0-9_-]+)*$/;
const encodeRedisKeySegment = (value: string) =>
encodeURIComponent(value).replace(
/[!'()*]/g,
(character) => `%${character.charCodeAt(0).toString(16).toUpperCase()}`
);
const escapeRedisGlob = (value: string) => value.replace(/[*?[\]\\]/g, '\\$&');
const ensureLogicalKey = (key: string) => {
if (!key) {
throw new Error('Redis logical key must not be empty');
}
if (key.startsWith(FASTGPT_REDIS_PREFIX)) {
throw new Error('Redis logical key must not include the physical prefix');
}
return key as RedisLogicalKey;
};
/** 构造新业务使用的逻辑 key;segment 使用 RFC3986 编码,禁止 Redis glob 字符泄漏。 */
export const createRedisLogicalKey = ({
namespace,
version,
segments = []
}: {
namespace: string;
version?: number;
segments?: readonly RedisKeySegment[];
}): RedisLogicalKey => {
if (!namespacePattern.test(namespace)) {
throw new Error('Redis key namespace contains unsupported characters');
}
if (version !== undefined && (!Number.isSafeInteger(version) || version < 1)) {
throw new Error('Redis key version must be a positive integer');
}
const encodedSegments = segments.map((segment) => {
const value = String(segment);
if (!value) {
throw new Error('Redis key segment must not be empty');
}
return encodeRedisKeySegment(value);
});
const keyParts = [
namespace,
...(version === undefined ? [] : [`v${version}`]),
...encodedSegments
];
return ensureLogicalKey(keyParts.join(':'));
};
/** 将历史调用方传入的逻辑 key 收窄为 RedisLogicalKey,不改变现有 key 格式。 */
export const asRedisLogicalKey = (key: string): RedisLogicalKey => ensureLogicalKey(key);
/** 显式添加 FastGPT 物理前缀,替代业务代码依赖 ioredis keyPrefix。 */
export const toPhysicalRedisKey = (key: string): RedisPhysicalKey =>
`${FASTGPT_REDIS_PREFIX}${ensureLogicalKey(key)}` as RedisPhysicalKey;
/** 只移除 key 开头的 FastGPT 前缀,拒绝读取其他应用的物理 key。 */
export const toLogicalRedisKey = (key: string): RedisLogicalKey => {
if (!key.startsWith(FASTGPT_REDIS_PREFIX)) {
throw new Error('Redis physical key does not belong to the FastGPT keyspace');
}
return ensureLogicalKey(key.slice(FASTGPT_REDIS_PREFIX.length));
};
/** 为历史“删除某前缀下全部子 key”语义生成物理 SCAN pattern。 */
export const createChildRedisScanPattern = (logicalPrefix: string): string =>
`${escapeRedisGlob(toPhysicalRedisKey(logicalPrefix))}:*`;
import {
isRedisOperationError,
RedisInvalidArgumentError,
RedisOperationExecutionError,
RedisOperationTimeoutError,
type RedisOperationOutcome
} from './errors';
/** Redis command 的执行语义;operation 名称本身只用于错误和观测标签。 */
export type RedisOperationMode = 'read' | 'idempotent-write' | 'uncertain-write';
type RedisOperationPolicy = {
maxAttempts: 1 | 2;
timeoutOutcome: Exclude<RedisOperationOutcome, 'not-started'>;
};
const DEFAULT_OPERATION_TIMEOUT_MS = 3_000;
const operationPolicies: Record<RedisOperationMode, RedisOperationPolicy> = {
read: {
maxAttempts: 2,
timeoutOutcome: 'failed'
},
'idempotent-write': {
maxAttempts: 2,
timeoutOutcome: 'unknown'
},
'uncertain-write': {
maxAttempts: 1,
timeoutOutcome: 'unknown'
}
};
const transientErrorMessages = [
'ECONNREFUSED',
'ECONNRESET',
'EPIPE',
'ETIMEDOUT',
'EAI_AGAIN',
'READONLY',
'Connection is closed',
'Reached the max retries per request limit'
];
class RedisAttemptTimeoutError extends Error {}
const isTransientRedisError = (error: unknown) => {
if (error instanceof RedisAttemptTimeoutError) return true;
const message = error instanceof Error ? error.message : String(error);
return transientErrorMessages.some((item) => message.includes(item));
};
const executeAttempt = <T>({
execute,
timeoutMs
}: {
execute: () => Promise<T>;
timeoutMs: number;
}) =>
new Promise<T>((resolve, reject) => {
const timeout = setTimeout(() => reject(new RedisAttemptTimeoutError()), timeoutMs);
Promise.resolve()
.then(execute)
.then(resolve, reject)
.finally(() => clearTimeout(timeout));
});
export type RedisOperationInput<T> = {
operation: string;
execute: () => Promise<T>;
timeoutMs?: number;
};
/**
* 集中执行 Redis operation,并按写入是否可能重复应用选择 retry 语义。
*
* operation 只作为错误和观测标签,不再需要维护全量 operation allowlist。调用方只能选择
* read、幂等写入或结果未知写入三种固定语义,不能自行声明 retry 次数。timeout 仅终止等待,
* 不能取消已经发往 Redis 的命令,因此写操作超时会标记 outcome=unknown。
*/
export class RedisOperationExecutor {
constructor(private readonly defaultTimeoutMs = DEFAULT_OPERATION_TIMEOUT_MS) {}
readonly read = <T>(input: RedisOperationInput<T>): Promise<T> =>
this.execute({ ...input, mode: 'read' });
readonly idempotentWrite = <T>(input: RedisOperationInput<T>): Promise<T> =>
this.execute({ ...input, mode: 'idempotent-write' });
readonly uncertainWrite = <T>(input: RedisOperationInput<T>): Promise<T> =>
this.execute({ ...input, mode: 'uncertain-write' });
/** 按调用方选择的执行模式运行 operation,并统一处理 timeout、retry 和错误结果。 */
private async execute<T>({
operation,
mode,
execute,
timeoutMs
}: RedisOperationInput<T> & { mode: RedisOperationMode }): Promise<T> {
const policy = operationPolicies[mode];
const effectiveTimeoutMs = timeoutMs ?? this.defaultTimeoutMs;
if (!Number.isSafeInteger(effectiveTimeoutMs) || effectiveTimeoutMs <= 0) {
throw new RedisInvalidArgumentError({
operation,
message: 'timeoutMs must be a positive safe integer'
});
}
for (let attempt = 1; ; attempt += 1) {
try {
return await executeAttempt({ execute, timeoutMs: effectiveTimeoutMs });
} catch (error) {
if (isRedisOperationError(error)) throw error;
const canRetry = attempt < policy.maxAttempts && isTransientRedisError(error);
if (canRetry) continue;
if (error instanceof RedisAttemptTimeoutError) {
throw new RedisOperationTimeoutError({
operation,
timeoutMs: effectiveTimeoutMs,
attempt,
outcome: policy.timeoutOutcome
});
}
throw new RedisOperationExecutionError({
operation,
attempt,
outcome: policy.timeoutOutcome,
cause: error
});
}
}
}
}
const defaultRedisOperationExecutor = new RedisOperationExecutor();
export const executeRedisRead = <T>(input: RedisOperationInput<T>) =>
defaultRedisOperationExecutor.read(input);
export const executeRedisIdempotentWrite = <T>(input: RedisOperationInput<T>) =>
defaultRedisOperationExecutor.idempotentWrite(input);
export const executeRedisUncertainWrite = <T>(input: RedisOperationInput<T>) =>
defaultRedisOperationExecutor.uncertainWrite(input);
import { RedisInvalidResponseError } from './errors';
import type { RedisStreamEntry } from '../types';
/** 从 Redis INFO 文本中读取非负整数;字段缺失或格式不匹配时返回 undefined。 */
export const parseRedisInfoNumber = (info: string, key: string) => {
const match = info.match(new RegExp(`(?:^|\\r?\\n)${key}:(\\d+)`));
if (!match) return undefined;
const value = Number(match[1]);
return Number.isFinite(value) ? value : undefined;
};
/** 校验并解析 Redis Stream 返回的交替 field/value 数组。 */
export const parseStreamFields = ({
operation,
rawFields
}: {
operation: string;
rawFields: unknown;
}): Record<string, string> => {
if (
!Array.isArray(rawFields) ||
rawFields.length % 2 !== 0 ||
rawFields.some((field) => typeof field !== 'string')
) {
throw new RedisInvalidResponseError({
operation,
message: 'Redis Stream fields returned an unsupported response'
});
}
const fields: Record<string, string> = {};
for (let index = 0; index < rawFields.length; index += 2) {
fields[rawFields[index] as string] = rawFields[index + 1] as string;
}
return fields;
};
/** 校验并解析 XRANGE/XREAD 返回的 Stream entry 列表。 */
export const parseStreamEntries = ({
operation,
rawEntries
}: {
operation: string;
rawEntries: unknown;
}): RedisStreamEntry[] => {
if (!Array.isArray(rawEntries)) {
throw new RedisInvalidResponseError({
operation,
message: 'Redis Stream entries returned an unsupported response'
});
}
return rawEntries.map((entry) => {
if (!Array.isArray(entry) || entry.length !== 2 || typeof entry[0] !== 'string') {
throw new RedisInvalidResponseError({
operation,
message: 'Redis Stream entry returned an unsupported response'
});
}
return {
id: entry[0],
fields: parseStreamFields({ operation, rawFields: entry[1] })
};
});
};
import type { RedisOptions } from 'ioredis';
import type { RedisClient, RedisConnectionRole, RedisConnectionState } from './connection';
import type { RedisRuntimeLogger } from '../types';
const DEFAULT_COMMAND_TIMEOUT_MS = 5_000;
const roleOptions: Record<
RedisConnectionRole,
Partial<
Pick<
RedisOptions,
| 'autoResendUnfulfilledCommands'
| 'commandTimeout'
| 'enableOfflineQueue'
| 'maxRetriesPerRequest'
>
>
> = {
command: {
enableOfflineQueue: true,
maxRetriesPerRequest: 1,
commandTimeout: DEFAULT_COMMAND_TIMEOUT_MS,
autoResendUnfulfilledCommands: false
},
blocking: {
enableOfflineQueue: true,
maxRetriesPerRequest: null,
autoResendUnfulfilledCommands: false
},
queue: {
enableOfflineQueue: true,
maxRetriesPerRequest: 3
},
worker: {
enableOfflineQueue: true,
maxRetriesPerRequest: null
}
};
const reconnectErrorMessages = ['READONLY', 'ECONNREFUSED', 'ETIMEDOUT', 'ECONNRESET'];
const getErrorMessage = (error: unknown) => {
return error instanceof Error ? error.message : String(error ?? 'Unknown Redis error');
};
/** 将 ioredis 当前 status 映射为 Runtime 对外暴露的连接状态。 */
export const getInitialConnectionState = (client: RedisClient): RedisConnectionState => {
if (client.status === 'ready') return 'ready';
if (client.status === 'connect') return 'connected';
if (client.status === 'reconnecting') return 'reconnecting';
if (client.status === 'close') return 'closed';
return 'connecting';
};
/** 为不同 Redis 连接角色合并 endpoint、重连和队列策略。 */
export const getConnectionOptions = ({
endpointOptions,
role,
logger
}: {
endpointOptions: RedisOptions;
role: RedisConnectionRole;
logger: RedisRuntimeLogger;
}): RedisOptions => ({
...endpointOptions,
retryStrategy: (times: number) => {
const delayMs = Math.min(times * 50, 2000);
if (times === 1 || times % 30 === 0) {
logger.warn('Redis reconnect scheduled', { role, attempt: times, delayMs });
}
return delayMs;
},
reconnectOnError: (error: Error) => {
const message = getErrorMessage(error);
const shouldReconnect = reconnectErrorMessages.some((errorType) => message.includes(errorType));
if (shouldReconnect) {
logger.warn('Redis reconnect requested by command error', { role, message });
}
return shouldReconnect;
},
connectTimeout: 10_000,
...roleOptions[role]
});
import { z } from 'zod';
/** Redis 参数和 Cache 配置共用的正安全整数 schema。 */
export const PositiveSafeIntegerSchema = z.int().positive();
/** Redis 返回值和计数器共用的非负安全整数 schema。 */
export const NonNegativeSafeIntegerSchema = z.int().nonnegative();
/** Redis cache 中允许正负小数,但拒绝 NaN 和 Infinity。 */
export const FiniteNumberSchema = z.number().finite();
import { getConfiguredRedisRuntime, type RedisRuntime } from './connection';
import type { RedisRuntimeLogger } from '../types';
const REDIS_SHUTDOWN_CONTEXT_SYMBOL = Symbol.for('@fastgpt/dal/redis/shutdown-context');
const DEFAULT_SHUTDOWN_SIGNALS = ['SIGTERM', 'SIGINT'] as const;
const silentLogger: RedisRuntimeLogger = {
info: () => undefined,
warn: () => undefined,
error: () => undefined
};
type RedisProcessLike = {
on: (signal: NodeJS.Signals, listener: () => void) => unknown;
removeListener: (signal: NodeJS.Signals, listener: () => void) => unknown;
exit: (code?: number) => void;
};
type RedisShutdownRegistration = {
runtime: RedisRuntime;
processRef: RedisProcessLike;
unregister: () => void;
};
type RedisShutdownContext = {
registration?: RedisShutdownRegistration;
};
export type RedisRuntimeShutdownOptions = {
runtime?: RedisRuntime;
logger?: RedisRuntimeLogger;
processRef?: RedisProcessLike;
signals?: readonly NodeJS.Signals[];
exitProcess?: boolean;
};
const getRedisShutdownContext = (): RedisShutdownContext => {
const existing = Reflect.get(globalThis, REDIS_SHUTDOWN_CONTEXT_SYMBOL) as
| RedisShutdownContext
| undefined;
if (existing) return existing;
const context: RedisShutdownContext = {};
Reflect.set(globalThis, REDIS_SHUTDOWN_CONTEXT_SYMBOL, context);
return context;
};
/**
* 注册进程级 Redis 优雅关闭处理器。
*
* 该函数只在应用 instrumentation 显式调用时安装 SIGTERM/SIGINT listener;library import
* 不产生进程副作用。重复注册同一 Runtime 是幂等的,信号触发后只执行一次 close。
*/
export const registerRedisRuntimeShutdown = ({
runtime = getConfiguredRedisRuntime(),
logger = silentLogger,
processRef = process,
signals = DEFAULT_SHUTDOWN_SIGNALS,
exitProcess = true
}: RedisRuntimeShutdownOptions = {}) => {
if (!runtime) {
throw new Error('Redis runtime has not been configured');
}
const context = getRedisShutdownContext();
const existing = context.registration;
if (existing?.runtime === runtime && existing.processRef === processRef) {
return existing.unregister;
}
existing?.unregister();
let closePromise: Promise<void> | undefined;
let registered = true;
const handlers = new Map<NodeJS.Signals, () => void>();
const unregister = () => {
if (!registered) return;
registered = false;
for (const [signal, handler] of handlers) {
processRef.removeListener(signal, handler);
}
handlers.clear();
if (context.registration?.unregister === unregister) {
context.registration = undefined;
}
};
const handleSignal = () => {
if (closePromise) return;
unregister();
// 通过 Promise 链统一接住 Runtime.close 的同步异常和异步拒绝,避免信号 listener
// 把异常抛回 Node 的 EventEmitter 调用栈。
closePromise = Promise.resolve().then(() => runtime.close());
void closePromise.then(
() => {
logger.info('Redis runtime closed after process shutdown signal');
if (exitProcess) processRef.exit(0);
},
(error) => {
logger.error('Redis runtime close failed after process shutdown signal', { error });
if (exitProcess) processRef.exit(1);
}
);
};
for (const signal of signals) {
const handler = () => handleSignal();
handlers.set(signal, handler);
processRef.on(signal, handler);
}
const registration = { runtime, processRef, unregister } satisfies RedisShutdownRegistration;
context.registration = registration;
return unregister;
};
/** 在指定 deadline 内等待 operation;超时只终止等待,不取消底层 Redis 操作。 */
export const runWithTimeout = <T>({
operation,
timeoutMs,
timeoutMessage
}: {
operation: Promise<T>;
timeoutMs: number;
timeoutMessage: string;
}): Promise<T> => {
return new Promise<T>((resolve, reject) => {
const timeout = setTimeout(() => reject(new Error(timeoutMessage)), timeoutMs);
operation.then(resolve, reject).finally(() => clearTimeout(timeout));
});
};
import { RedisInvalidArgumentError } from './errors';
import { PositiveSafeIntegerSchema } from './schema';
/** 严格解析正安全整数,并将 Zod issue 映射为稳定、脱敏的 Redis 参数错误。 */
export const parsePositiveInteger = ({
value,
operation,
field,
maximum
}: {
value: unknown;
operation: string;
field: string;
maximum?: number;
}): number => {
const result = PositiveSafeIntegerSchema.safeParse(value);
if (!result.success || (maximum !== undefined && result.data > maximum)) {
throw new RedisInvalidArgumentError({
operation,
message: `${field} must be a positive safe integer${
maximum === undefined ? '' : ` no greater than ${maximum}`
}`
});
}
return result.data;
};
/** 严格解析可选毫秒 TTL;undefined 保持为未设置,不做隐式类型转换。 */
export const parseOptionalTtlMs = ({
ttlMs,
operation
}: {
ttlMs: unknown;
operation: string;
}): number | undefined => {
if (ttlMs === undefined) return undefined;
return parsePositiveInteger({ value: ttlMs, operation, field: 'ttlMs' });
};
/** Redis 结构化日志 metadata 的最小约定。 */
export type RedisLogMetadata = Record<string, unknown>;
/** Redis Runtime/Cache 共用的日志方法签名。 */
export type RedisLogMethod = (message: string, metadata?: RedisLogMetadata) => void;
/** Runtime 使用完整日志能力;Cache 通过泛型只声明实际使用的级别。 */
export type RedisRuntimeLogger = {
info: RedisLogMethod;
warn: RedisLogMethod;
error: RedisLogMethod;
};
/** Redis Runtime 健康检查的 metrics 结果。metrics 不参与业务错误处理。 */
export type RedisRuntimeHealthMetric = {
success: boolean;
latencyMs: number;
};
/** Redis Runtime 关闭耗时的 metrics 结果。 */
export type RedisRuntimeShutdownMetric = {
durationMs: number;
};
/** Redis Runtime 的可选观测 port;实现方可以接入 OpenTelemetry 或测试 recorder。 */
export type RedisRuntimeMetrics = {
connectionCreated?: (role: string) => void;
connectionClosed?: (role: string) => void;
connectionError?: (role: string) => void;
healthCheck?: (result: RedisRuntimeHealthMetric) => void;
shutdownCompleted?: (result: RedisRuntimeShutdownMetric) => void;
};
export type RedisCacheLoggerLevel = 'warn' | 'error';
export type RedisCacheLogger<Level extends RedisCacheLoggerLevel = RedisCacheLoggerLevel> = Pick<
RedisRuntimeLogger,
Level
>;
declare const redisLogicalKeyBrand: unique symbol;
declare const redisPhysicalKeyBrand: unique symbol;
export type RedisLogicalKey = string & { readonly [redisLogicalKeyBrand]: true };
export type RedisPhysicalKey = string & { readonly [redisPhysicalKeyBrand]: true };
/** Redis INFO MEMORY operation 的 typed 结果。 */
export type RedisMemoryInfo = {
usedMemory?: number;
maxMemory?: number;
};
/** Redis Stream entry 的 normalized 结果,不向 Cache 暴露 raw ioredis response。 */
export type RedisStreamEntry = {
id: string;
fields: Record<string, string>;
};
import fs from 'node:fs';
/** 按测试环境优先级定位 DAL 专属环境文件。 */
const getEnvFilePath = (): URL | undefined => {
const packageRoot = new URL('../', import.meta.url);
const envFiles = ['.env.test.local', '.env.test', '.env.local', '.env'];
return envFiles
.map((fileName) => new URL(fileName, packageRoot))
.find((filePath) => fs.existsSync(filePath));
};
/** 在 Vitest worker 启动前加载 DAL integration 所需的环境变量。 */
export const setup = () => {
const envFilePath = getEnvFilePath();
if (envFilePath) {
process.loadEnvFile(envFilePath);
}
};
export const teardown = () => undefined;
import Redis from 'ioredis';
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest';
import { RedisCacheAdapter } from '@fastgpt/dal/redis/adapter';
import { DailyActiveDedupeCache } from '@fastgpt/dal/redis/caches';
const redisUrl = process.env.REDIS_INTEGRATION_URL;
const describeWithRedis = redisUrl ? describe : describe.skip;
describeWithRedis('DailyActiveDedupeCache Redis 7.2 integration', () => {
const uid = `integration-${process.pid}-${Date.now()}`;
const date = '2026-07-24';
const physicalKey = `fastgpt:cache:dailyUserActive:${uid}_${date}`;
const logger = { warn: vi.fn() };
let client: Redis;
beforeAll(async () => {
client = new Redis(redisUrl!, {
enableOfflineQueue: false,
lazyConnect: true,
maxRetriesPerRequest: 1
});
await client.connect();
});
afterAll(async () => {
await client.del(physicalKey);
await client.quit();
});
it('allows exactly one winner across concurrent claims and keeps the historical TTL', async () => {
const adapter = new RedisCacheAdapter({ getCommandClient: () => client });
const cache = new DailyActiveDedupeCache({ redis: adapter, logger });
const results = await Promise.all(
Array.from({ length: 64 }, () => cache.shouldRecord({ uid, date }))
);
expect(results.filter(Boolean)).toHaveLength(1);
expect(await client.get(physicalKey)).toBe('1');
expect(await client.ttl(physicalKey)).toBeGreaterThan(86_390);
expect(await client.ttl(physicalKey)).toBeLessThanOrEqual(86_400);
expect(logger.warn).not.toHaveBeenCalled();
});
});
import Redis from 'ioredis';
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
import { RedisCacheAdapter } from '@fastgpt/dal/redis/adapter';
import { FixedWindowRateLimitCache } from '@fastgpt/dal/redis/caches';
const redisUrl = process.env.REDIS_INTEGRATION_URL;
const describeWithRedis = redisUrl ? describe : describe.skip;
describeWithRedis('FixedWindowRateLimitCache Redis 7.2 integration', () => {
const key = `integration-fixed-window-${process.pid}-${Date.now()}`;
const physicalKey = `fastgpt:${key}`;
let client: Redis;
beforeAll(async () => {
client = new Redis(redisUrl!, {
enableOfflineQueue: false,
lazyConnect: true,
maxRetriesPerRequest: 1
});
await client.connect();
});
afterAll(async () => {
await client.del(physicalKey);
await client.quit();
});
it('assigns unique counts under concurrency and keeps one fixed TTL', async () => {
const adapter = new RedisCacheAdapter({ getCommandClient: () => client });
const cache = new FixedWindowRateLimitCache({ redis: adapter });
const results = await Promise.all(
Array.from({ length: 64 }, () => cache.consume({ key, limit: 32, windowSeconds: 60 }))
);
expect(new Set(results.map((result) => result.currentCount))).toHaveLength(64);
expect(results.filter((result) => result.allowed)).toHaveLength(32);
expect(await client.get(physicalKey)).toBe('64');
expect(await client.ttl(physicalKey)).toBeGreaterThan(0);
expect(await client.ttl(physicalKey)).toBeLessThanOrEqual(60);
});
});
import Redis from 'ioredis';
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
import { asRedisLogicalKey, RedisCacheAdapter } from '@fastgpt/dal/redis/adapter';
const redisUrl = process.env.REDIS_INTEGRATION_URL;
const describeWithRedis = redisUrl ? describe : describe.skip;
describeWithRedis('Redis 7.2 kernel integration', () => {
const namespace = `integration:redis-kernel:${process.pid}:${Date.now()}`;
const valueKey = asRedisLogicalKey(`${namespace}:value`);
const getOrSetKey = asRedisLogicalKey(`${namespace}:get-or-set`);
const scanPrefix = asRedisLogicalKey(`${namespace}:scan:a*b`);
const scanChildKeys = Array.from({ length: 256 }, (_, index) => `${scanPrefix}:child-${index}`);
const unrelatedScanKey = `${namespace}:scan:aXb:child`;
const leaseKey = asRedisLogicalKey(`${namespace}:lease`);
const streamKey = asRedisLogicalKey(`${namespace}:stream`);
const physical = (key: string) => `fastgpt:${key}`;
const cleanupKeys = [
physical(valueKey),
physical(getOrSetKey),
...scanChildKeys.map(physical),
physical(unrelatedScanKey),
physical(leaseKey),
physical(streamKey)
];
let client: Redis;
beforeAll(async () => {
client = new Redis(redisUrl!, {
enableOfflineQueue: false,
lazyConnect: true,
maxRetriesPerRequest: 1
});
await client.connect();
const serverInfo = await client.info('server');
const version = serverInfo.match(/redis_version:(\d+)\.(\d+)\./);
expect(version).not.toBeNull();
const majorVersion = Number(version?.[1]);
const minorVersion = Number(version?.[2]);
expect(majorVersion).toBeGreaterThanOrEqual(7);
expect(majorVersion > 7 || minorVersion >= 2).toBe(true);
});
afterAll(async () => {
if (!client) return;
await client.del(...cleanupKeys);
await client.quit();
});
const createAdapter = () =>
new RedisCacheAdapter({
getCommandClient: () => client,
createBlockingConnection: () =>
client.duplicate({ enableOfflineQueue: true, maxRetriesPerRequest: null }),
releaseConnection: async (blockingClient) => {
const redisClient = blockingClient as Redis;
if (redisClient.status === 'end' || redisClient.status === 'close') {
redisClient.disconnect();
return;
}
await redisClient.quit().catch(() => redisClient.disconnect());
}
});
it('uses one explicit physical keyspace and safely paginates SCAN patterns', async () => {
const adapter = createAdapter();
await adapter.set({ key: valueKey, value: 'value' });
expect(await client.get(physical(valueKey))).toBe('value');
expect(await client.get(`fastgpt:${physical(valueKey)}`)).toBeNull();
const pipeline = client.pipeline();
scanChildKeys.forEach((key) => pipeline.set(physical(key), 'child'));
pipeline.set(physical(unrelatedScanKey), 'unrelated');
await pipeline.exec();
const scannedKeys: string[] = [];
for await (const keys of adapter.iterateByPrefix({ prefix: scanPrefix, batchSize: 16 })) {
scannedKeys.push(...keys);
}
expect(new Set(scannedKeys)).toEqual(new Set(scanChildKeys));
expect(await client.get(physical(unrelatedScanKey))).toBe('unrelated');
});
it('keeps SET NX GET atomic under concurrency', async () => {
const adapter = createAdapter();
const values = await Promise.all(
Array.from({ length: 128 }, (_, index) =>
adapter.getOrSet({ key: getOrSetKey, value: `candidate-${index}` })
)
);
expect(new Set(values)).toHaveLength(1);
expect(await client.get(physical(getOrSetKey))).toBe(values[0]);
});
it('keeps Lua lease ownership and token checks atomic under concurrency', async () => {
const adapter = createAdapter();
const tokens = Array.from({ length: 64 }, (_, index) => `token-${index}`);
const acquired = await Promise.all(
tokens.map((token) => adapter.acquireLease({ key: leaseKey, token, ttlMs: 5_000 }))
);
expect(acquired.filter(Boolean)).toHaveLength(1);
const winner = tokens[acquired.findIndex(Boolean)]!;
await expect(
adapter.renewLease({ key: leaseKey, token: 'wrong-token', ttlMs: 5_000 })
).resolves.toBe(false);
await expect(adapter.renewLease({ key: leaseKey, token: winner, ttlMs: 5_000 })).resolves.toBe(
true
);
await expect(adapter.releaseLease({ key: leaseKey, token: 'wrong-token' })).resolves.toBe(
false
);
await expect(adapter.releaseLease({ key: leaseKey, token: winner })).resolves.toBe(true);
await expect(client.get(physical(leaseKey))).resolves.toBeNull();
});
it('writes, ranges, and blocks on a Redis Stream with an isolated reader connection', async () => {
const adapter = createAdapter();
const reader = adapter.createBlockingStreamReader({ key: streamKey, blockMs: 1_000 });
try {
const readPromise = reader.read('$');
const appendPromise = new Promise<string>((resolve, reject) => {
setTimeout(() => {
void adapter
.appendStreamEntry({ key: streamKey, fields: { raw: 'hello' } })
.then(resolve, reject);
}, 25);
});
const [entries, streamId] = await Promise.all([readPromise, appendPromise]);
expect(entries).toEqual([{ id: streamId, fields: { raw: 'hello' } }]);
await expect(
adapter.rangeStream({ key: streamKey, start: '-', end: '+', count: 10 })
).resolves.toEqual([{ id: streamId, fields: { raw: 'hello' } }]);
} finally {
await reader.close();
}
});
});
import Redis from 'ioredis';
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest';
import { RedisCacheAdapter } from '@fastgpt/dal/redis/adapter';
import { SystemVersionCache } from '@fastgpt/dal/redis/caches';
const redisUrl = process.env.REDIS_INTEGRATION_URL;
const describeWithRedis = redisUrl ? describe : describe.skip;
describeWithRedis('SystemVersionCache Redis 7.2 integration', () => {
const key = `integration-system-version-${process.pid}-${Date.now()}`;
const basePhysicalKey = `fastgpt:VERSION_KEY:${key}`;
const childPhysicalKeys = Array.from(
{ length: 512 },
(_, index) => `${basePhysicalKey}:team-${index}`
);
const unrelatedPhysicalKey = `${basePhysicalKey}-other:team-1`;
let client: Redis;
beforeAll(async () => {
client = new Redis(redisUrl!, {
enableOfflineQueue: false,
lazyConnect: true,
maxRetriesPerRequest: 1
});
await client.connect();
});
afterAll(async () => {
await client.del(basePhysicalKey, unrelatedPhysicalKey, ...childPhysicalKeys);
await client.quit();
});
it('returns one permanent UUID across concurrent first initialization', async () => {
const adapter = new RedisCacheAdapter({ getCommandClient: () => client });
const cache = new SystemVersionCache({ redis: adapter });
const versions = await Promise.all(
Array.from({ length: 64 }, () => cache.getOrInitialize({ key }))
);
expect(new Set(versions)).toHaveLength(1);
expect(await client.get(basePhysicalKey)).toBe(versions[0]);
expect(await client.ttl(basePhysicalKey)).toBe(-1);
});
it('uses paginated SCAN to delete only child keys and preserves the base key', async () => {
const pipeline = client.pipeline();
childPhysicalKeys.forEach((childKey) => pipeline.set(childKey, 'child-version'));
pipeline.set(unrelatedPhysicalKey, 'unrelated-version');
await pipeline.exec();
const scanSpy = vi.spyOn(client, 'scan');
const adapter = new RedisCacheAdapter({ getCommandClient: () => client });
const cache = new SystemVersionCache({ redis: adapter });
await cache.refresh({ key, id: '*' });
expect(scanSpy.mock.calls.length).toBeGreaterThan(1);
expect(await client.get(basePhysicalKey)).toBeTruthy();
expect((await client.mget(...childPhysicalKeys)).every((value) => value === null)).toBe(true);
expect(await client.get(unrelatedPhysicalKey)).toBe('unrelated-version');
});
});
import Redis from 'ioredis';
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
import { RedisCacheAdapter } from '@fastgpt/dal/redis/adapter';
import { TeamPointCache } from '@fastgpt/dal/redis/caches';
const redisUrl = process.env.REDIS_INTEGRATION_URL;
const describeWithRedis = redisUrl ? describe : describe.skip;
describeWithRedis('TeamPointCache Redis 7.2 integration', () => {
const teamId = `integration-team-point-${process.pid}-${Date.now()}`;
const physicalKeys = [
`fastgpt:cache:team_point_surplus:${teamId}`,
`fastgpt:cache:team_point_total:${teamId}`
];
let client: Redis;
beforeAll(async () => {
client = new Redis(redisUrl!, {
enableOfflineQueue: false,
lazyConnect: true,
maxRetriesPerRequest: 1
});
await client.connect();
});
afterAll(async () => {
await client.del(...physicalKeys);
await client.quit();
});
it('reads and refreshes both keys with the same TTL', async () => {
const adapter = new RedisCacheAdapter({ getCommandClient: () => client });
const cache = new TeamPointCache({ redis: adapter });
await cache.set({ teamId, totalPoints: 2_000, surplusPoints: 1_500 });
await expect(cache.get(teamId)).resolves.toEqual({
totalPoints: 2_000,
surplusPoints: 1_500
});
expect(await client.ttl(physicalKeys[0])).toBeGreaterThan(55);
expect(await client.ttl(physicalKeys[1])).toBeGreaterThan(55);
});
it('keeps pair reads coherent while concurrent pair writes occur', async () => {
const adapter = new RedisCacheAdapter({ getCommandClient: () => client });
const cache = new TeamPointCache({ redis: adapter });
await cache.set({ teamId, totalPoints: 0, surplusPoints: 0 });
const reads = await Promise.all(
Array.from({ length: 64 }, (_, index) =>
Promise.all([
cache.set({ teamId, totalPoints: index, surplusPoints: -index }),
cache.get(teamId)
]).then(([, value]) => value)
)
);
reads.forEach((value) => {
if (!value) return;
expect(value.surplusPoints + value.totalPoints).toBe(0);
});
});
it('increments surplus with TTL and clears both keys together', async () => {
const adapter = new RedisCacheAdapter({ getCommandClient: () => client });
const cache = new TeamPointCache({ redis: adapter });
await cache.set({ teamId, totalPoints: 100, surplusPoints: 40 });
await cache.incrementSurplus({ teamId, value: -5 });
await expect(cache.get(teamId)).resolves.toEqual({
totalPoints: 100,
surplusPoints: 35
});
await cache.clear(teamId);
await expect(cache.get(teamId)).resolves.toBeUndefined();
});
});
This source diff could not be displayed because it is too large. You can view the blob instead.
This diff is collapsed. Click to expand it.
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or sign in to comment