Commit bd140f71 by Archer Committed by GitHub

feat(code-sandbox): add queueId concurrency limit (#7008)

* feat(code-sandbox): add queue id concurrency limit

* doc

* test(code-sandbox): avoid duplicate app logger setup
parent 4af1ef77
# 代码沙盒 queueId 并发排队设计
## 目标
在代码沙盒运行接口中新增可选 `queueId`,并通过环境变量控制同一个 `queueId` 同时可进入执行流程的请求数。
目标行为:
- 默认不启用,完全兼容现有调用;
- 启用后仅对带 `queueId` 的请求生效;
- 同一 `queueId` 超出并发上限时 FIFO 等待;
- 进程池原有 worker 等待队列继续负责真实 worker 分配。
## 接口设计
### 请求体
`POST /sandbox/js``POST /sandbox/python` 增加字段:
```json
{
"code": "async function main() { return {} }",
"variables": {},
"queueId": "team-xxx"
}
```
字段约束:
- `queueId` 可选;
- 空字符串会按未传处理;
- 非字符串返回 400;
- 最大长度限制为 128,避免异常请求造成队列 key 膨胀。
### 环境变量
新增:
| 变量 | 说明 | 默认值 |
| --- | --- | --- |
| `SANDBOX_QUEUE_ID_CONCURRENCY` | 同一 `queueId` 同时可进入执行流程的请求数;为空时不启用 queueId 排队 | 空 |
## 实现方案
新增 `QueueIdLimiter`
- 内部维护 `Map<string, QueueState>`
- 每个 `queueId` 有独立 FIFO 等待队列、运行计数和上限;
- `run(queueId, task)` 在未启用或无 queueId 时直接执行 task;
- 同一 queueId 运行数达到上限时,把请求放入对应等待队列;
- task 完成后释放许可,唤醒同 queueId 的下一个等待请求;
- 队列空且运行数为 0 时删除对应 Map entry。
与进程池组合:
```ts
await queueIdLimiter.run(queueId, () => pool.execute(options));
```
这形成两层队列:
- queueId 队列:限制同一个业务 id 的并发进入执行流程;
- process pool 队列:限制真实 worker 数量,并复用已有 worker 生命周期管理。
## 测试设计
- 单元测试 `QueueIdLimiter`
- 未启用时不排队;
- 同 queueId 按上限限制并发;
- 不同 queueId 互不阻塞;
- 空 queueId 不排队;
- 队列空后清理状态。
- API 集成测试:
- 接口接受 `queueId` 并正常执行;
- `queueId` 非字符串返回 400。
- 真实 HTTP 集成测试:
- 启动本地 Hono server;
- 通过 `fetch` 并发请求 `/sandbox/js`
- 验证同一 `queueId` 串行、不同 `queueId` 并行、未传 `queueId` 不受限。
## TODO
- [x] 梳理现有接口、env、进程池和测试结构。
- [x] 新增 `SANDBOX_QUEUE_ID_CONCURRENCY` 环境变量。
- [x] 新增 `QueueIdLimiter` 并补单元测试。
- [x] 在 JS/Python API 执行入口接入 queueId limiter。
- [x] 更新 `ExecuteOptions``CodeSandbox.runCode()` 类型和 README。
- [x] 运行代码沙盒局部测试。
- [x] 运行真实 HTTP 集成测试。
- [x] 运行代码沙盒全量测试。
# 代码沙盒 queueId 排队能力问题分析
## 背景
代码沙盒当前通过 `ProcessPool` / `PythonProcessPool` 维护每种语言的 worker 池。请求进入 `/sandbox/js``/sandbox/python` 后,会直接调用对应进程池的 `execute()`
- 有空闲 worker 时立即执行;
- 没有空闲 worker 时进入进程池内部 `waitQueue`,等待 worker 释放;
- 所有请求共享同一个语言池队列,无法按业务维度限制某一类请求的并发。
当某个业务方在短时间内提交大量代码执行请求时,会占用同语言 worker 和池内等待队列,影响其他业务方的请求延迟。
## 需求
为代码沙盒运行接口增加 `queueId`
- `POST /sandbox/js`
- `POST /sandbox/python`
新增环境变量控制同一个 `queueId` 同时允许多少个请求进入执行流程。环境变量为空时认为不启用排队能力,保持现有行为。
## 现状分析
### 可复用能力
- `projects/code-sandbox/src/pool/base-process-pool.ts` 已有 worker 维度等待队列 `waitQueue`,负责“待运行/待分配 worker”的队列。
- `projects/code-sandbox/src/utils/semaphore.ts` 已提供简单 FIFO 信号量,语义可用于单个 `queueId` 的并发控制。
### 插入位置
排队控制应放在 HTTP API 边界和进程池之间:
```
HTTP request -> queueId limiter -> process pool waitQueue -> worker execution
```
这样可以保持进程池只关心 worker 生命周期,不把业务 queueId 概念扩散到 worker 管理层。
## 边界语义
- 环境变量为空:不创建 queueId 队列,所有请求走现有进程池逻辑。
- 环境变量有值但请求未传 `queueId`:不按 queueId 排队,避免把所有未标识请求挤到同一个匿名队列。
- 同一个 `queueId` 内按 FIFO 唤醒。
- 不同 `queueId` 之间不做额外公平调度,仍由进程池 worker 队列决定实际执行顺序。
- `queueId` 队列在无运行请求且无等待请求后清理,避免高基数 id 导致内存长期增长。
## 实现约定
- 环境变量名称采用 `SANDBOX_QUEUE_ID_CONCURRENCY`
- FastGPT 主应用如果要对工作流代码节点启用业务排队,需要在调用 `codeSandbox.runCode()` 时显式传入 queueId。本次先实现代码沙盒运行接口和 SDK 方法的可选参数,不替业务侧猜测默认 queueId。
......@@ -3,24 +3,24 @@ title: 服务协议
description: ' FastGPT 服务协议'
---
最后更新时间:2024年3月3
最后更新时间:2026 年 5 月 28
FastGPT 服务协议是您与珠海环界云计算有限公司(以下简称“我们”或“本公司”)之间就FastGPT云服务(以下简称“本服务”)的使用等相关事项所订立的协议。请您仔细阅读并充分理解本协议各条款,特别是免除或者限制我们责任的条款、对您权益的限制条款、争议解决和法律适用条款等。如您不同意本协议任一内容,请勿注册或使用本服务。
FastGPT 服务协议是您与**广州环际云计算有限公司**(以下简称“我们”或“本公司”)之间就 FastGPT 云服务(以下简称“本服务”)的使用等相关事项所订立的协议。请您仔细阅读并充分理解本协议各条款,特别是免除或者限制我们责任的条款、对您权益的限制条款、争议解决和法律适用条款等。如您不同意本协议任一内容,请勿注册或使用本服务。
**第1条 服务内容**
**第 1 条服务内容**
1. 我们将向您提供存储、计算、网络传输等基于互联网的信息技术服务。
2. 我们将不定期向您通过站内信、电子邮件或短信等形式向您推送最新的动态。
3. 我们将为您提供相关技术支持和客户服务,帮助您更好地使用本服务。
4. 我们将为您提供稳定的在线服务,保证每月服务可用性不低于99%。
4. 我们将为您提供稳定的在线服务,保证每月服务可用性不低于 99%。
**第2条 用户注册与账户管理**
**第 2 条用户注册与账户管理**
1. 您在使用本服务前需要注册一个账户。您保证在注册时提供的信息真实、准确、完整,并及时更新。
2. 您应妥善保管账户名和密码,对由此产生的全部行为负责。如发现他人使用您的账户,请及时修改账号密码或与我们进行联系。
3. 我们有权对您的账户进行审查,如发现您的账户存在异常或违法情况,我们有权暂停或终止向您提供服务。
**第3条 使用规则**
**第 3 条使用规则**
1. 您不得利用本服务从事任何违法活动或侵犯他人合法权益的行为,包括但不限于侵犯知识产权、泄露他人商业机密等。
2. 您不得通过任何手段恶意注册账户,包括但不限于以牟利、炒作、套现等目的。
......@@ -45,26 +45,26 @@ FastGPT 服务协议是您与珠海环界云计算有限公司(以下简称“
- 禁止故意欺骗并可能对公共利益产生不利影响的内容,包括与健康、安全、选举诚信或公民参与相关的欺骗性或不真实内容。
- 直接支持非法主动攻击或造成技术危害的恶意软件活动的内容,例如提供恶意可执行文件、组织拒绝服务攻击或管理命令和控制服务器。
**第4条 费用及支付**
**第 4 条费用及支付**
1. 您同意支付与本服务相关的费用,具体费用标准以我们公布的价格为准。
2. 我们可能会根据运营成本和市场情况调整费用标准。最新价格以您付款时刻的价格为准。
**第5条 服务免责与责任限制**
**第 5 条服务免责与责任限制**
1. 本服务按照现有技术和条件所能达到的水平提供。我们不能保证本服务完全无故障或满足您的所有需求。
2. 对于因您自身误操作导致的数据丢失、损坏等情况,我们不承担责任。
3. 由于生成式 AI 的特性,其在不同国家的管控措施也会有所不同,请所有使用者务必遵守所在地的相关法律。如果您以任何违反 FastGPT 可接受使用政策的方式使用,包括但不限于法律、法规、政府命令或法令禁止的任何用途,或任何侵犯他人权利的使用;由使用者自行承担。我们对由客户使用产生的问题概不负责。下面是各国对生成式AI的管控条例的链接:
3. 由于生成式 AI 的特性,其在不同国家的管控措施也会有所不同,请所有使用者务必遵守所在地的相关法律。如果您以任何违反 FastGPT 可接受使用政策的方式使用,包括但不限于法律、法规、政府命令或法令禁止的任何用途,或任何侵犯他人权利的使用;由使用者自行承担。我们对由客户使用产生的问题概不负责。下面是各国对生成式 AI 的管控条例的链接:
[中国生成式人工智能服务管理办法(征求意见稿)](http://www.cac.gov.cn/2023-04/11/c_1682854275475410.htm)
**第6条 知识产权**
**第 6 条知识产权**
1. 我们对本服务及相关软件、技术、文档等拥有全部知识产权,除非经我们明确许可,您不得进行复制、分发、出租、反向工程等行为。
2. 您在使用本服务过程中产生的所有数据和内容(包括但不限于文件、图片等)的知识产权归您所有。我们不会对您的数据和内容进行使用、复制、修改等行为。
3. 在线服务中其他用户的数据和内容的知识产权归原用户所有,未经原用户许可,您不得进行使用、复制、修改等行为。
**第7条 其他条款**
**第 7 条其他条款**
1. 如本协议中部分条款因违反法律法规而被视为无效,不影响其他条款的效力。
2. 本公司保留对本协议及隐私政策的最终解释权。如您对本协议或隐私政策有任何疑问,请联系我们:archer@fastgpt.io。
......@@ -10,20 +10,23 @@ description: 'FastGPT V4.15.0-beta3 Release Notes'
## 🔧 Environment Variable Changes
Code Sandbox adds security-related environment variables such as `SANDBOX_API_MAX_BODY_MB` and `SANDBOX_MAX_OUTPUT_MB`. The full defaults are listed below:
| Variable | Default | Description |
| --------------------------------- | ------- | ------------------------------------------------------------------------------------------ |
| `SANDBOX_API_MAX_BODY_MB` | `8` | Maximum `/sandbox` API JSON body size, including `variables`, in MB. |
| `SANDBOX_MAX_OUTPUT_MB` | `10` | Maximum output JSON size for one code execution, including return values and logs, in MB. |
| `CHECK_INTERNAL_IP` | `true` | Enables internal IP checks for sandbox network requests by default to reduce SSRF risk. |
| `SANDBOX_MAX_TIMEOUT` | `60000` | Timeout for one code execution, in milliseconds. |
| `SANDBOX_MAX_MEMORY_MB` | `256` | Memory limit for one sandbox, in MB. The runtime reserves an extra `50` MB for overhead. |
| `SANDBOX_POOL_SIZE` | `20` | Number of pre-warmed JS/Python workers. |
| `SANDBOX_REQUEST_MAX_COUNT` | `30` | Maximum number of network requests allowed during one code execution. |
| `SANDBOX_REQUEST_TIMEOUT` | `60000` | Timeout for one network request from inside the sandbox, in milliseconds. |
| `SANDBOX_REQUEST_MAX_RESPONSE_MB` | `10` | Maximum response body size for one sandbox network request, in MB. |
| `SANDBOX_REQUEST_MAX_BODY_MB` | `5` | Maximum request body size for one sandbox network request, in MB. |
Code Sandbox adds security-related environment variables such as `SANDBOX_API_MAX_BODY_MB` and `SANDBOX_MAX_OUTPUT_MB`, and now supports grouped request queuing for run APIs through `queueId`. The full defaults are listed below:
| Variable | Default | Description |
| --------------------------------- | ------- | ----------------------------------------------------------------------------------------------------- |
| `SANDBOX_API_MAX_BODY_MB` | `8` | Maximum `/sandbox` API JSON body size, including `variables`, in MB. |
| `SANDBOX_MAX_OUTPUT_MB` | `10` | Maximum output JSON size for one code execution, including return values and logs, in MB. |
| `CHECK_INTERNAL_IP` | `true` | Enables internal IP checks for sandbox network requests by default to reduce SSRF risk. |
| `SANDBOX_MAX_TIMEOUT` | `60000` | Timeout for one code execution, in milliseconds. |
| `SANDBOX_MAX_MEMORY_MB` | `256` | Memory limit for one sandbox, in MB. The runtime reserves an extra `50` MB for overhead. |
| `SANDBOX_POOL_SIZE` | `20` | Number of pre-warmed JS/Python workers. |
| `SANDBOX_REQUEST_MAX_COUNT` | `30` | Maximum number of network requests allowed during one code execution. |
| `SANDBOX_REQUEST_TIMEOUT` | `60000` | Timeout for one network request from inside the sandbox, in milliseconds. |
| `SANDBOX_REQUEST_MAX_RESPONSE_MB` | `10` | Maximum response body size for one sandbox network request, in MB. |
| `SANDBOX_REQUEST_MAX_BODY_MB` | `5` | Maximum request body size for one sandbox network request, in MB. |
| `SANDBOX_QUEUE_ID_CONCURRENCY` | Empty | Number of requests with the same `queueId` that may enter execution at once. Empty disables queueing. |
`/sandbox/js` and `/sandbox/python` request bodies now support an optional `queueId` field. When `SANDBOX_QUEUE_ID_CONCURRENCY` is configured and a valid `queueId` is provided, requests with the same `queueId` are queued in FIFO order. Requests with different `queueId` values, or requests without `queueId`, are not limited by this setting and continue to be limited only by the worker pool concurrency.
## ⚙️ Improvements
......
......@@ -7,7 +7,7 @@ description: 'FastGPT V4.15.0-beta3 更新说明'
### 🔧 环境变量变更
Code Sandbox 新增 `SANDBOX_API_MAX_BODY_MB`、`SANDBOX_MAX_OUTPUT_MB` 等安全相关环境变量;完整默认值如下:
Code Sandbox 新增 `SANDBOX_API_MAX_BODY_MB`、`SANDBOX_MAX_OUTPUT_MB` 等安全相关环境变量,并支持通过 `queueId` 对运行接口做分组排队;完整默认值如下:
| 变量 | 默认值 | 说明 |
| --------------------------------- | ------- | ----------------------------------------------------------------- |
......@@ -21,6 +21,9 @@ Code Sandbox 新增 `SANDBOX_API_MAX_BODY_MB`、`SANDBOX_MAX_OUTPUT_MB` 等安
| `SANDBOX_REQUEST_TIMEOUT` | `60000` | 沙箱内单次网络请求超时时间,单位毫秒。 |
| `SANDBOX_REQUEST_MAX_RESPONSE_MB` | `10` | 沙箱内单次网络响应体最大大小,单位 MB。 |
| `SANDBOX_REQUEST_MAX_BODY_MB` | `5` | 沙箱内单次网络请求体最大大小,单位 MB。 |
| `SANDBOX_QUEUE_ID_CONCURRENCY` | 空 | 同一个 `queueId` 同时可进入执行流程的请求数;为空时不启用排队。 |
`/sandbox/js` 和 `/sandbox/python` 请求体新增可选字段 `queueId`。当配置 `SANDBOX_QUEUE_ID_CONCURRENCY` 且请求传入有效 `queueId` 时,同一 `queueId` 的请求会按 FIFO 排队;不同 `queueId` 或未传 `queueId` 的请求不受该限制,仍只受 worker 池并发限制影响。
## 🚀 新增内容
......
......@@ -280,8 +280,8 @@
"content/self-host/upgrading/4-15/4150.mdx": "2026-05-20T17:52:26+08:00",
"content/self-host/upgrading/4-15/41502.en.mdx": "2026-05-25T11:21:30+08:00",
"content/self-host/upgrading/4-15/41502.mdx": "2026-05-25T11:21:30+08:00",
"content/self-host/upgrading/4-15/41503.en.mdx": "2026-05-27T12:17:46+08:00",
"content/self-host/upgrading/4-15/41503.mdx": "2026-05-27T14:50:36+08:00",
"content/self-host/upgrading/4-15/41503.en.mdx": "2026-05-28T12:45:18+08:00",
"content/self-host/upgrading/4-15/41503.mdx": "2026-05-28T12:45:18+08:00",
"content/self-host/upgrading/outdated/40.en.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/outdated/40.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/outdated/41.en.mdx": "2026-04-26T21:08:47+08:00",
......
......@@ -57,11 +57,13 @@ export class CodeSandbox {
async runCode({
codeType,
code,
variables
variables,
queueId
}: {
codeType: string;
code: string;
variables: Record<string, any>;
queueId?: string;
}) {
const url = (() => {
if (codeType == SandboxCodeTypeEnum.py) {
......@@ -74,7 +76,7 @@ export class CodeSandbox {
const { data } = await this.client.post<{
codeReturn: Record<string, any>;
log: string;
}>(url, { code, variables });
}>(url, { code, variables, queueId });
return data;
}
......
......@@ -75,10 +75,13 @@ docker run -p 3000:3000 \
```json
{
"code": "async function main(variables) {\n return { result: variables.a + variables.b }\n}",
"variables": { "a": 1, "b": 2 }
"variables": { "a": 1, "b": 2 },
"queueId": "team-xxx"
}
```
`queueId` 可选;仅当配置 `SANDBOX_QUEUE_ID_CONCURRENCY` 时,同一 `queueId` 会按该并发数排队执行。
### `POST /sandbox/python`
执行 Python 代码。
......@@ -86,7 +89,8 @@ docker run -p 3000:3000 \
```json
{
"code": "def main(variables):\n return {'result': variables['a'] + variables['b']}",
"variables": { "a": 1, "b": 2 }
"variables": { "a": 1, "b": 2 },
"queueId": "team-xxx"
}
```
......@@ -140,6 +144,7 @@ docker run -p 3000:3000 \
| 变量 | 说明 | 默认值 |
|------|------|--------|
| `SANDBOX_POOL_SIZE` | 每种语言的 worker 进程数 | `20` |
| `SANDBOX_QUEUE_ID_CONCURRENCY` | 同一 `queueId` 同时可进入执行流程的请求数,空值表示不按 `queueId` 排队 | 空 |
### 资源限制
......
......@@ -47,6 +47,8 @@ export const env = createEnv({
// ===== 进程池 =====
/** 进程池大小(预热 worker 数量) */
SANDBOX_POOL_SIZE: IntSchema.min(1).max(100).default(20),
/** 同一 queueId 同时可进入执行流程的请求数;为空时不启用 queueId 排队 */
SANDBOX_QUEUE_ID_CONCURRENCY: IntSchema.min(1).max(100).optional(),
// ===== 资源限制 =====
SANDBOX_API_MAX_BODY_MB: IntSchema.min(1).max(100).default(8),
......
......@@ -8,6 +8,7 @@ import { PythonProcessPool } from './pool/python-process-pool';
import type { ExecuteOptions } from './types';
import { getErrText } from './utils';
import { configureLogger, getLogger, LogCategories } from './utils/logger';
import { QueueIdLimiter } from './utils/queue-id-limiter';
await configureLogger();
......@@ -71,12 +72,19 @@ async function readLimitedJsonBody(c: Context): Promise<unknown> {
}
/** 请求体校验 schema */
const queueIdSchema = z.preprocess((value) => {
if (typeof value !== 'string') return value;
const queueId = value.trim();
return queueId || undefined;
}, z.string().max(128).optional());
const executeSchema = z.object({
code: z
.string()
.min(1)
.max(5 * 1024 * 1024), // 最大 5MB 代码
variables: z.record(z.string(), z.any()).default({})
variables: z.record(z.string(), z.any()).default({}),
queueId: queueIdSchema
});
const app = new Hono();
......@@ -84,6 +92,13 @@ const app = new Hono();
/** 进程池 */
const jsPool = new ProcessPool(env.SANDBOX_POOL_SIZE);
const pythonPool = new PythonProcessPool(env.SANDBOX_POOL_SIZE);
const queueIdLimiter = new QueueIdLimiter(env.SANDBOX_QUEUE_ID_CONCURRENCY);
if (queueIdLimiter.enabled) {
serverLogger.info(
`QueueId limiter enabled: max ${env.SANDBOX_QUEUE_ID_CONCURRENCY} concurrent requests per queueId`
);
}
const poolReady = Promise.all([jsPool.init(), pythonPool.init()])
.then(() => {
......@@ -173,7 +188,9 @@ app.post('/sandbox/js', async (c) => {
400
);
}
const result = await jsPool.execute(parsed.data as ExecuteOptions);
const result = await queueIdLimiter.run(parsed.data.queueId, () =>
jsPool.execute(parsed.data as ExecuteOptions)
);
return c.json(result);
} catch (err: any) {
const status = err instanceof ApiBodyError ? err.status : 200;
......@@ -201,7 +218,9 @@ app.post('/sandbox/python', async (c) => {
400
);
}
const result = await pythonPool.execute(parsed.data as ExecuteOptions);
const result = await queueIdLimiter.run(parsed.data.queueId, () =>
pythonPool.execute(parsed.data as ExecuteOptions)
);
return c.json(result);
} catch (err: any) {
const status = err instanceof ApiBodyError ? err.status : 200;
......
......@@ -2,6 +2,7 @@
export type ExecuteOptions = {
code: string;
variables: Record<string, any>;
queueId?: string;
};
/** 执行结果 */
......
type QueueState = {
running: number;
waiters: Array<() => void>;
};
export type QueueIdLimiterStats = {
enabled: boolean;
maxConcurrency?: number;
queueCount: number;
queues: Array<{
queueId: string;
running: number;
queued: number;
}>;
};
/**
* 按 queueId 控制进入执行流程的请求并发。
*
* limiter 只负责同一个 queueId 内的 FIFO 限流;真实 worker 分配仍交给
* ProcessPool 的等待队列处理,避免把业务排队规则耦合到 worker 生命周期管理。
*/
export class QueueIdLimiter {
private readonly queues = new Map<string, QueueState>();
constructor(private readonly maxConcurrency?: number) {
if (maxConcurrency !== undefined && (!Number.isInteger(maxConcurrency) || maxConcurrency < 1)) {
throw new Error('QueueIdLimiter maxConcurrency must be a positive integer');
}
}
get enabled(): boolean {
return this.maxConcurrency !== undefined;
}
/**
* 在指定 queueId 的并发限制下执行任务。
*
* 未启用 limiter 或 queueId 为空时直接执行,保持历史接口行为。
*/
async run<T>(queueId: string | undefined, task: () => Promise<T>): Promise<T> {
if (!this.enabled || !queueId) {
return task();
}
await this.acquire(queueId);
try {
return await task();
} finally {
this.release(queueId);
}
}
get stats(): QueueIdLimiterStats {
return {
enabled: this.enabled,
maxConcurrency: this.maxConcurrency,
queueCount: this.queues.size,
queues: Array.from(this.queues.entries()).map(([queueId, state]) => ({
queueId,
running: state.running,
queued: state.waiters.length
}))
};
}
private acquire(queueId: string): Promise<void> {
const state = this.getOrCreateQueue(queueId);
if (state.running < this.maxConcurrency!) {
state.running++;
return Promise.resolve();
}
return new Promise<void>((resolve) => {
state.waiters.push(resolve);
});
}
private release(queueId: string): void {
const state = this.queues.get(queueId);
if (!state) return;
const next = state.waiters.shift();
if (next) {
// 直接把当前运行名额转交给等待队列头部,running 保持不变。
next();
return;
}
state.running = Math.max(0, state.running - 1);
if (state.running === 0) {
this.queues.delete(queueId);
}
}
private getOrCreateQueue(queueId: string): QueueState {
const existing = this.queues.get(queueId);
if (existing) return existing;
const state: QueueState = { running: 0, waiters: [] };
this.queues.set(queueId, state);
return state;
}
}
......@@ -2,10 +2,54 @@
* API 测试 - 使用 app.request() 直接测试 Hono 路由
* 无需启动服务或配置 CODE_SANDBOX_URL
*/
import { describe, it, expect, beforeAll } from 'vitest';
import { serve, type ServerType } from '@hono/node-server';
import { describe, it, expect, beforeAll, afterAll } from 'vitest';
import { app, poolReady } from '../../src/index';
import { env } from '../../src/env';
type RunWindow = {
label: string;
startedAt: number;
finishedAt: number;
};
type SandboxResponse = {
success: boolean;
data?: {
codeReturn: RunWindow;
log: string;
};
message?: string;
};
const delayCode = `
async function main(v) {
const startedAt = Date.now();
await new Promise((resolve) => setTimeout(resolve, v.delayMs));
return {
label: v.label,
startedAt,
finishedAt: Date.now()
};
}
`;
function hasOverlap(a: RunWindow, b: RunWindow) {
return a.startedAt < b.finishedAt && b.startedAt < a.finishedAt;
}
function closeServer(server: ServerType) {
return new Promise<void>((resolve, reject) => {
server.close((err) => {
if (err) {
reject(err);
return;
}
resolve();
});
});
}
/** 构造请求 headers,自动带上 auth(如果配置了 token) */
function headers(extra: Record<string, string> = {}): Record<string, string> {
const h: Record<string, string> = { ...extra };
......@@ -15,11 +59,11 @@ function headers(extra: Record<string, string> = {}): Record<string, string> {
return h;
}
async function executeJs(code: string, variables: Record<string, any> = {}) {
async function executeJs(code: string, variables: Record<string, any> = {}, queueId?: string) {
const res = await app.request('/sandbox/js', {
method: 'POST',
headers: headers({ 'Content-Type': 'application/json' }),
body: JSON.stringify({ code, variables })
body: JSON.stringify({ code, variables, queueId })
});
return res.json();
......@@ -66,6 +110,17 @@ describe('API Routes', () => {
expect(data.success).toBe(true);
});
it('POST /sandbox/js 接受 queueId 正常执行', async () => {
const data = await executeJs(
'async function main(v) { return { ok: true, name: v.name } }',
{ name: 'queue' },
'team-queue-test'
);
expect(data.success).toBe(true);
expect(data.data.codeReturn).toEqual({ ok: true, name: 'queue' });
});
it('POST /sandbox/js 安全拦截', async () => {
const res = await app.request('/sandbox/js', {
method: 'POST',
......@@ -247,6 +302,85 @@ describe('API Routes', () => {
});
});
describe('HTTP queueId integration', () => {
let server: ServerType | undefined;
let baseUrl = '';
beforeAll(async () => {
await poolReady;
const info = await new Promise<{ port: number }>((resolve) => {
server = serve({ fetch: app.fetch, hostname: '127.0.0.1', port: 0 }, (address) => {
resolve({ port: address.port });
});
});
baseUrl = `http://127.0.0.1:${info.port}`;
}, 30000);
afterAll(async () => {
if (server) {
await closeServer(server);
}
});
async function runJs({
label,
queueId,
delayMs = 300
}: {
label: string;
queueId?: string;
delayMs?: number;
}) {
const res = await fetch(`${baseUrl}/sandbox/js`, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
...(env.SANDBOX_TOKEN ? { Authorization: `Bearer ${env.SANDBOX_TOKEN}` } : {})
},
body: JSON.stringify({
code: delayCode,
variables: { label, delayMs },
queueId
})
});
expect(res.status).toBe(200);
const body = (await res.json()) as SandboxResponse;
expect(body.success).toBe(true);
expect(body.data?.codeReturn.label).toBe(label);
return body.data!.codeReturn;
}
it('同一 queueId 的真实 HTTP 请求会串行进入执行流程', async () => {
const queueId = `same-${Date.now()}`;
const [first, second] = await Promise.all([
runJs({ label: 'first', queueId }),
runJs({ label: 'second', queueId })
]);
expect(hasOverlap(first, second)).toBe(false);
});
it('不同 queueId 的真实 HTTP 请求可以并行执行', async () => {
const [first, second] = await Promise.all([
runJs({ label: 'queue-a', queueId: `queue-a-${Date.now()}` }),
runJs({ label: 'queue-b', queueId: `queue-b-${Date.now()}` })
]);
expect(hasOverlap(first, second)).toBe(true);
});
it('未传 queueId 的真实 HTTP 请求不受 queueId 并发限制', async () => {
const [first, second] = await Promise.all([
runJs({ label: 'no-queue-a' }),
runJs({ label: 'no-queue-b' })
]);
expect(hasOverlap(first, second)).toBe(true);
});
});
// ===== 错误处理安全 =====
describe('API 错误处理安全', () => {
beforeAll(async () => {
......@@ -349,6 +483,22 @@ describe('API Zod 校验失败', () => {
expect(data.success).toBe(false);
});
it('JS: queueId 非字符串返回 400', async () => {
const res = await app.request('/sandbox/js', {
method: 'POST',
headers: headers({ 'Content-Type': 'application/json' }),
body: JSON.stringify({
code: 'async function main() { return { ok: true } }',
variables: {},
queueId: 123
})
});
expect(res.status).toBe(400);
const data = await res.json();
expect(data.success).toBe(false);
expect(data.message).toMatch(/Invalid request/i);
});
it('Python: code 为数字返回 400', async () => {
const res = await app.request('/sandbox/python', {
method: 'POST',
......
import { describe, expect, it } from 'vitest';
import { QueueIdLimiter } from '../../src/utils/queue-id-limiter';
const delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));
const flush = () => delay(0);
function createDeferred() {
let resolve!: () => void;
const promise = new Promise<void>((r) => {
resolve = r;
});
return { promise, resolve };
}
describe('QueueIdLimiter', () => {
it('未启用时不按 queueId 排队', async () => {
const limiter = new QueueIdLimiter();
let running = 0;
let maxRunning = 0;
await Promise.all(
Array.from({ length: 3 }, () =>
limiter.run('same-queue', async () => {
running++;
maxRunning = Math.max(maxRunning, running);
await delay(10);
running--;
})
)
);
expect(maxRunning).toBe(3);
expect(limiter.stats).toMatchObject({
enabled: false,
queueCount: 0,
queues: []
});
});
it('同一个 queueId 超出并发后按 FIFO 排队', async () => {
const limiter = new QueueIdLimiter(1);
const firstGate = createDeferred();
const secondGate = createDeferred();
const order: string[] = [];
const first = limiter.run('team-a', async () => {
order.push('first-start');
await firstGate.promise;
order.push('first-end');
return 'first';
});
await flush();
const second = limiter.run('team-a', async () => {
order.push('second-start');
await secondGate.promise;
order.push('second-end');
return 'second';
});
const third = limiter.run('team-a', async () => {
order.push('third-start');
return 'third';
});
await flush();
expect(order).toEqual(['first-start']);
expect(limiter.stats.queues).toEqual([{ queueId: 'team-a', running: 1, queued: 2 }]);
firstGate.resolve();
await first;
await flush();
expect(order).toEqual(['first-start', 'first-end', 'second-start']);
expect(limiter.stats.queues).toEqual([{ queueId: 'team-a', running: 1, queued: 1 }]);
secondGate.resolve();
await Promise.all([second, third]);
expect(order).toEqual([
'first-start',
'first-end',
'second-start',
'second-end',
'third-start'
]);
expect(limiter.stats.queueCount).toBe(0);
});
it('不同 queueId 之间互不阻塞', async () => {
const limiter = new QueueIdLimiter(1);
const gate = createDeferred();
const order: string[] = [];
const first = limiter.run('team-a', async () => {
order.push('team-a-start');
await gate.promise;
});
await flush();
await limiter.run('team-b', async () => {
order.push('team-b-start');
});
expect(order).toEqual(['team-a-start', 'team-b-start']);
gate.resolve();
await first;
expect(limiter.stats.queueCount).toBe(0);
});
it('queueId 为空时不排队', async () => {
const limiter = new QueueIdLimiter(1);
let running = 0;
let maxRunning = 0;
await Promise.all(
Array.from({ length: 3 }, () =>
limiter.run(undefined, async () => {
running++;
maxRunning = Math.max(maxRunning, running);
await delay(10);
running--;
})
)
);
expect(maxRunning).toBe(3);
expect(limiter.stats.queueCount).toBe(0);
});
});
......@@ -19,6 +19,7 @@ export default defineConfig({
CHECK_INTERNAL_IP: 'true',
SANDBOX_API_MAX_BODY_MB: '1',
SANDBOX_MAX_TIMEOUT: '5000',
SANDBOX_QUEUE_ID_CONCURRENCY: '1',
SANDBOX_TOKEN: 'test'
}
}
......
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