Commit 4f95f686 by heheer Committed by GitHub

app delete queue (#6122)

* app delete queue

* test

* perf: del app queue

* perf: log

* perf: query

* perf: retry del s3

* fix: ts

* perf: add job

* redis retry

* perf: mq check

* update log

* perf: mq concurrency

* perf: error check

* perf: mq

* perf: init model

---------

Co-authored-by: archer <545436317@qq.com>
parent 36821600
......@@ -11,6 +11,7 @@ description: 'FastGPT V4.14.5 更新说明'
## ⚙️ 优化
1. 优化获取 redis 所有 key 的逻辑,避免大量获取时导致阻塞。
2. Redis 和 MQ 的重连逻辑优化。
## 🐛 修复
......
......@@ -120,7 +120,7 @@
"document/content/docs/upgrading/4-14/4142.mdx": "2025-11-18T19:27:14+08:00",
"document/content/docs/upgrading/4-14/4143.mdx": "2025-11-26T20:52:05+08:00",
"document/content/docs/upgrading/4-14/4144.mdx": "2025-12-16T14:56:04+08:00",
"document/content/docs/upgrading/4-14/4145.mdx": "2025-12-18T23:25:48+08:00",
"document/content/docs/upgrading/4-14/4145.mdx": "2025-12-19T00:08:30+08:00",
"document/content/docs/upgrading/4-8/40.mdx": "2025-08-02T19:38:37+08:00",
"document/content/docs/upgrading/4-8/41.mdx": "2025-08-02T19:38:37+08:00",
"document/content/docs/upgrading/4-8/42.mdx": "2025-08-02T19:38:37+08:00",
......
......@@ -59,6 +59,9 @@ export type AppSchema = {
inited?: boolean;
/** @deprecated */
teamTags: string[];
// 软删除字段
deleteTime?: Date | null;
};
export type AppListItemType = {
......
......@@ -8,6 +8,7 @@ import {
} from 'bullmq';
import { addLog } from '../system/log';
import { newQueueRedisConnection, newWorkerRedisConnection } from '../redis';
import { delay } from '@fastgpt/global/common/system/utils';
const defaultWorkerOpts: Omit<ConnectionOptions, 'connection'> = {
removeOnComplete: {
......@@ -25,6 +26,7 @@ export enum QueueNames {
// Delete Queue
datasetDelete = 'datasetDelete',
appDelete = 'appDelete',
// @deprecated
websiteSync = 'websiteSync'
}
......@@ -77,15 +79,41 @@ export function getWorker<DataType, ReturnType = void>(
const newWorker = new Worker<DataType, ReturnType>(name.toString(), processor, {
connection: newWorkerRedisConnection(),
...defaultWorkerOpts,
// 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)
...opts
});
// default error handler, to avoid unhandled exceptions
newWorker.on('error', (error) => {
addLog.error(`MQ Worker [${name}]: ${error.message}`, error);
newWorker.on('error', async (error) => {
addLog.error(`MQ Worker error`, {
message: error.message,
data: { name }
});
await newWorker.close();
});
newWorker.on('failed', (jobId, error) => {
addLog.error(`MQ Worker [${name}]: ${error.message}`, error);
// Critical: Worker has been closed - remove from pool
newWorker.on('closed', async () => {
addLog.error(`MQ Worker [${name}] closed unexpectedly`, {
data: {
name,
message: 'Worker will need to be manually restarted'
}
});
try {
await delay(1000);
workers.delete(name);
getWorker(name, processor, opts);
} catch (error) {}
});
newWorker.on('paused', async () => {
addLog.warn(`MQ Worker [${name}] paused`);
await delay(1000);
newWorker.resume();
});
workers.set(name, newWorker);
return newWorker;
}
......
......@@ -3,26 +3,54 @@ import Redis from 'ioredis';
const REDIS_URL = process.env.REDIS_URL ?? 'redis://localhost:6379';
// Base Redis options for connection reliability
const REDIS_BASE_OPTION = {
// Retry strategy: exponential backoff with unlimited retries for stability
retryStrategy: (times: number) => {
// Never give up retrying to ensure worker keeps running
const delay = Math.min(times * 50, 2000); // Max 2s between retries
if (times > 10) {
addLog.error(`[Redis connection failed] attempt ${times}, will keep retrying...`);
} else {
addLog.warn(`Redis reconnecting... attempt ${times}, delay ${delay}ms`);
}
return delay; // Always return a delay to keep retrying
},
// Reconnect on specific errors (Redis master-slave switch, network issues)
reconnectOnError: (err: any) => {
const reconnectErrors = ['READONLY', 'ECONNREFUSED', 'ETIMEDOUT', 'ECONNRESET'];
const shouldReconnect = reconnectErrors.some((errType) => err.message.includes(errType));
if (shouldReconnect) {
addLog.warn(`Redis reconnecting due to error: ${err.message}`);
}
return shouldReconnect;
},
// Connection timeout
connectTimeout: 10000, // 10 seconds
// Enable offline queue to buffer commands when disconnected
enableOfflineQueue: true
};
export const newQueueRedisConnection = () => {
const redis = new Redis(REDIS_URL);
redis.on('connect', () => {
console.log('Redis connected');
const redis = new Redis(REDIS_URL, {
...REDIS_BASE_OPTION,
// Limit retries for queue operations
maxRetriesPerRequest: 3
});
redis.on('error', (error) => {
console.error('Redis connection error', error);
addLog.error('[Redis Queue connection error]', error);
});
return redis;
};
export const newWorkerRedisConnection = () => {
const redis = new Redis(REDIS_URL, {
...REDIS_BASE_OPTION,
// BullMQ requires maxRetriesPerRequest: null for blocking operations
maxRetriesPerRequest: null
});
redis.on('connect', () => {
console.log('Redis connected');
});
redis.on('error', (error) => {
console.error('Redis connection error', error);
addLog.error('[Redis Worker connection error]', error);
});
return redis;
};
......@@ -31,13 +59,17 @@ export const FASTGPT_REDIS_PREFIX = 'fastgpt:';
export const getGlobalRedisConnection = () => {
if (global.redisClient) return global.redisClient;
global.redisClient = new Redis(REDIS_URL, { keyPrefix: FASTGPT_REDIS_PREFIX });
global.redisClient.on('connect', () => {
addLog.info('Redis connected');
global.redisClient = new Redis(REDIS_URL, {
...REDIS_BASE_OPTION,
keyPrefix: FASTGPT_REDIS_PREFIX,
maxRetriesPerRequest: 3
});
global.redisClient.on('error', (error) => {
addLog.error('Redis connection error', error);
addLog.error('[Redis Global connection error]', error);
});
global.redisClient.on('close', () => {
addLog.warn('[Redis Global connection closed]');
});
return global.redisClient;
......
import { Client, type RemoveOptions, type CopyConditions, S3Error } from 'minio';
import {
Client,
type RemoveOptions,
type CopyConditions,
S3Error,
InvalidObjectNameError,
InvalidXMLError
} from 'minio';
import {
type CreatePostPresignedUrlOptions,
type CreatePostPresignedUrlParams,
......@@ -17,6 +24,25 @@ import { type Readable } from 'node:stream';
import { type UploadFileByBufferParams, UploadFileByBufferSchema } from '../type';
import { parseFileExtensionFromUrl } from '@fastgpt/global/common/string/tools';
// Check if the error is a "file not found" type error, which should be treated as success
export const isFileNotFoundError = (error: any): boolean => {
if (error instanceof S3Error) {
// Handle various "not found" error codes
return (
error.code === 'NoSuchKey' ||
error.code === 'InvalidObjectName' ||
error.message === 'Not Found' ||
error.message ===
'The request signature we calculated does not match the signature you provided. Check your key and signing method.' ||
error.message.includes('Object name contains unsupported characters.')
);
}
if (error instanceof InvalidObjectNameError || error instanceof InvalidXMLError) {
return true;
}
return false;
};
export class S3BaseBucket {
private _client: Client;
private _externalClient: Client | undefined;
......@@ -94,7 +120,7 @@ export class S3BaseBucket {
temporary: false
}
});
await this.delete(from);
await this.removeObject(from);
}
async copy({
......@@ -120,24 +146,17 @@ export class S3BaseBucket {
return this.client.copyObject(bucket, to, `${bucket}/${from}`, options?.copyConditions);
}
async delete(objectKey: string, options?: RemoveOptions): Promise<void> {
try {
if (!objectKey) return Promise.resolve();
// 把连带的 parsed 数据一起删除
const fileParsedPrefix = `${path.dirname(objectKey)}/${path.basename(objectKey, path.extname(objectKey))}-parsed`;
await this.addDeleteJob({ prefix: fileParsedPrefix });
return await this.client.removeObject(this.bucketName, objectKey, options);
} catch (error) {
if (error instanceof S3Error) {
if (error.code === 'InvalidObjectName') {
addLog.warn(`${this.bucketName} delete object not found: ${objectKey}`, error);
return Promise.resolve();
}
async removeObject(objectKey: string, options?: RemoveOptions): Promise<void> {
return this.client.removeObject(this.bucketName, objectKey, options).catch((err) => {
if (isFileNotFoundError(err)) {
return Promise.resolve();
}
return Promise.reject(error);
}
addLog.error(`[S3 delete error]`, {
message: err.message,
data: { code: err.code, key: objectKey }
});
throw err;
});
}
// 列出文件
......
......@@ -27,26 +27,7 @@ export async function clearExpiredMinioFiles() {
const bucket = global.s3BucketMap[bucketName];
if (bucket) {
await bucket.delete(file.minioKey);
if (!file.minioKey.includes('-parsed/')) {
try {
const dir = path.dirname(file.minioKey);
const basename = path.basename(file.minioKey);
const ext = path.extname(basename);
if (ext) {
const nameWithoutExt = path.basename(basename, ext);
const parsedPrefix = `${dir}/${nameWithoutExt}-parsed`;
await bucket.addDeleteJob({ prefix: parsedPrefix });
addLog.info(`Scheduled deletion of parsed images: ${parsedPrefix}`);
}
} catch (error) {
addLog.debug(`Failed to schedule parsed images deletion for ${file.minioKey}`);
}
}
await bucket.addDeleteJob({ key: file.minioKey });
await MongoS3TTL.deleteOne({ _id: file._id });
success++;
......@@ -57,12 +38,6 @@ export async function clearExpiredMinioFiles() {
addLog.warn(`Bucket not found: ${file.bucketName}`);
}
} catch (error) {
if (
error instanceof S3Error &&
error.message.includes('Object name contains unsupported characters.')
) {
await MongoS3TTL.deleteOne({ _id: file._id });
}
fail++;
addLog.error(`Failed to delete minio file: ${file.minioKey}`, error);
}
......
import { getQueue, getWorker, QueueNames } from '../bullmq';
import pLimit from 'p-limit';
import { retryFn } from '@fastgpt/global/common/system/utils';
import { addLog } from '../system/log';
import path from 'path';
import { batchRun } from '@fastgpt/global/common/system/utils';
import { isFileNotFoundError, type S3BaseBucket } from './buckets/base';
export type S3MQJobData = {
key?: string;
......@@ -10,89 +11,99 @@ export type S3MQJobData = {
bucketName: string;
};
const jobOption = {
attempts: 10,
removeOnFail: {
count: 10000, // 保留10000个失败任务
age: 14 * 24 * 60 * 60 // 14 days
},
removeOnComplete: true,
backoff: {
delay: 2000,
type: 'exponential'
}
};
export const addS3DelJob = async (data: S3MQJobData): Promise<void> => {
const queue = getQueue<S3MQJobData>(QueueNames.s3FileDelete);
const jobId = (() => {
if (data.key) {
return data.key;
}
if (data.keys) {
return undefined;
}
if (data.prefix) {
return data.prefix;
}
throw new Error('Invalid s3 delete job data');
})();
await queue.add('delete-s3-files', data, { jobId, ...jobOption });
};
await queue.add(
'delete-s3-files',
{ ...data },
{
attempts: 3,
removeOnFail: false,
removeOnComplete: true,
backoff: {
delay: 2000,
type: 'exponential'
const prefixDel = async (bucket: S3BaseBucket, prefix: string) => {
addLog.debug(`[S3 delete] delete prefix: ${prefix}`);
let tasks: Promise<any>[] = [];
return new Promise<void>(async (resolve, reject) => {
const stream = bucket.listObjectsV2(prefix, true);
stream.on('data', (file) => {
if (!file.name) return;
tasks.push(bucket.removeObject(file.name));
});
stream.on('end', async () => {
if (tasks.length === 0) {
return resolve();
}
}
);
const results = await Promise.allSettled(tasks);
const failed = results.some((r) => r.status === 'rejected');
if (failed) {
addLog.error(`[S3 delete] delete prefix failed: ${prefix}`);
reject('Some deletes failed');
}
resolve();
});
stream.on('error', (err) => {
if (isFileNotFoundError(err)) {
return resolve();
}
addLog.error(`[S3 delete] delete prefix: ${prefix} error`, err);
reject(err);
});
});
};
export const startS3DelWorker = async () => {
return getWorker<S3MQJobData>(
QueueNames.s3FileDelete,
async (job) => {
const { prefix, bucketName, key, keys } = job.data;
const limit = pLimit(10);
const bucket = s3BucketMap[bucketName];
let { prefix, bucketName, key, keys } = job.data;
const bucket = global.s3BucketMap[bucketName];
if (!bucket) {
return Promise.reject(`Bucket not found: ${bucketName}`);
addLog.error(`Bucket not found: ${bucketName}`);
return;
}
if (key) {
addLog.info(`[S3 delete] delete key: ${key}`);
await bucket.delete(key);
addLog.info(`[S3 delete] delete key: ${key} success`);
keys = [key];
}
if (keys) {
addLog.info(`[S3 delete] delete keys: ${keys.length}`);
const tasks: Promise<void>[] = [];
for (const key of keys) {
const p = limit(() => retryFn(() => bucket.delete(key)));
tasks.push(p);
}
await Promise.all(tasks);
addLog.info(`[S3 delete] delete keys: ${keys.length} success`);
addLog.debug(`[S3 delete] delete keys: ${keys.length}`);
await batchRun(keys, async (key) => {
await bucket.removeObject(key);
// Delete parsed
if (!key.includes('-parsed/')) {
const fileParsedPrefix = `${path.dirname(key)}/${path.basename(key, path.extname(key))}-parsed`;
await prefixDel(bucket, fileParsedPrefix);
}
});
}
if (prefix) {
addLog.info(`[S3 delete] delete prefix: ${prefix}`);
const tasks: Promise<void>[] = [];
return new Promise<void>(async (resolve, reject) => {
const stream = bucket.listObjectsV2(prefix, true);
stream.on('data', async (file) => {
if (!file.name) return;
const p = limit(() =>
// 因为封装的 delete 方法里,包含前缀删除,这里不能再使用,避免循环。
retryFn(() => bucket.client.removeObject(bucket.bucketName, file.name))
);
tasks.push(p);
});
stream.on('end', async () => {
try {
const results = await Promise.allSettled(tasks);
const failed = results.filter((r) => r.status === 'rejected');
if (failed.length > 0) {
addLog.error(`[S3 delete] delete prefix: ${prefix} failed`);
reject('Some deletes failed');
}
addLog.info(`[S3 delete] delete prefix: ${prefix} success`);
resolve();
} catch (err) {
addLog.error(`[S3 delete] delete prefix: ${prefix} error`, err);
reject(err);
}
});
stream.on('error', (err) => {
addLog.error(`[S3 delete] delete prefix: ${prefix} error`, err);
reject(err);
});
});
await prefixDel(bucket, prefix);
}
},
{
concurrency: 1
concurrency: 3
}
);
};
......@@ -42,7 +42,7 @@ class S3AvatarSource extends S3PublicBucket {
async deleteAvatar(avatar: string, session?: ClientSession): Promise<void> {
const key = avatar.slice(this.prefix.length);
await MongoS3TTL.deleteOne({ minioKey: key, bucketName: this.bucketName }, session);
await this.delete(key);
await this.removeObject(key);
}
async refreshAvatar(newAvatar?: string, oldAvatar?: string, session?: ClientSession) {
......
......@@ -120,42 +120,40 @@ export const loadSystemModels = async (init = false, language = 'en') => {
]);
// Load system model from local
await Promise.all(
systemModels.map(async (model) => {
const mergeObject = (obj1: any, obj2: any) => {
if (!obj1 && !obj2) return undefined;
const formatObj1 = typeof obj1 === 'object' ? obj1 : {};
const formatObj2 = typeof obj2 === 'object' ? obj2 : {};
return { ...formatObj1, ...formatObj2 };
};
systemModels.forEach((model) => {
const mergeObject = (obj1: any, obj2: any) => {
if (!obj1 && !obj2) return undefined;
const formatObj1 = typeof obj1 === 'object' ? obj1 : {};
const formatObj2 = typeof obj2 === 'object' ? obj2 : {};
return { ...formatObj1, ...formatObj2 };
};
const dbModel = dbModels.find((item) => item.model === model.model);
const provider = getModelProvider(dbModel?.metadata?.provider || model.provider, language);
const dbModel = dbModels.find((item) => item.model === model.model);
const provider = getModelProvider(dbModel?.metadata?.provider || model.provider, language);
const modelData: any = {
...model,
...dbModel?.metadata,
provider: provider.id,
avatar: provider.avatar,
type: dbModel?.metadata?.type || model.type,
isCustom: false,
const modelData: any = {
...model,
...dbModel?.metadata,
provider: provider.id,
avatar: provider.avatar,
type: dbModel?.metadata?.type || model.type,
isCustom: false,
...(model.type === ModelTypeEnum.llm && {
maxResponse: model.maxTokens || 4000
}),
...(model.type === ModelTypeEnum.llm && {
maxResponse: model.maxTokens || 4000
}),
...(model.type === ModelTypeEnum.llm && dbModel?.metadata?.type === ModelTypeEnum.llm
? {
maxResponse: dbModel?.metadata?.maxResponse ?? model.maxTokens ?? 4000,
defaultConfig: mergeObject(model.defaultConfig, dbModel?.metadata?.defaultConfig),
fieldMap: mergeObject(model.fieldMap, dbModel?.metadata?.fieldMap),
maxTokens: undefined
}
: {})
};
pushModel(modelData);
})
);
...(model.type === ModelTypeEnum.llm && dbModel?.metadata?.type === ModelTypeEnum.llm
? {
maxResponse: dbModel?.metadata?.maxResponse ?? model.maxTokens ?? 4000,
defaultConfig: mergeObject(model.defaultConfig, dbModel?.metadata?.defaultConfig),
fieldMap: mergeObject(model.fieldMap, dbModel?.metadata?.fieldMap),
maxTokens: undefined
}
: {})
};
pushModel(modelData);
});
// Custom model(Not in system config)
dbModels.forEach((dbModel) => {
......@@ -240,8 +238,7 @@ export const loadSystemModels = async (init = false, language = 'en') => {
);
} catch (error) {
console.error('Load models error', error);
// @ts-ignore
global.systemModelList = undefined;
return Promise.reject(error);
}
};
......
......@@ -4,12 +4,10 @@ import {
FlowNodeInputTypeEnum,
FlowNodeTypeEnum
} from '@fastgpt/global/core/workflow/node/constant';
import { AppFolderTypeList } from '@fastgpt/global/core/app/constants';
import { MongoApp } from './schema';
import type { StoreNodeItemType } from '@fastgpt/global/core/workflow/type/node';
import { encryptSecretValue, storeSecretValue } from '../../common/secret/utils';
import { SystemToolSecretInputTypeEnum } from '@fastgpt/global/core/app/tool/systemTool/constants';
import { type ClientSession } from '../../common/mongo';
import { MongoEvaluation } from './evaluation/evalSchema';
import { removeEvaluationJob } from './evaluation/mq';
import { MongoChatItem } from '../chat/chatItemSchema';
......@@ -23,10 +21,12 @@ import { MongoChatSetting } from '../chat/setting/schema';
import { MongoResourcePermission } from '../../support/permission/schema';
import { PerResourceTypeEnum } from '@fastgpt/global/support/permission/constant';
import { removeImageByPath } from '../../common/file/image/controller';
import { mongoSessionRun } from '../../common/mongo/sessionRun';
import { MongoAppLogKeys } from './logs/logkeysSchema';
import { MongoChatItemResponse } from '../chat/chatItemResponseSchema';
import { getS3ChatSource } from '../../common/s3/sources/chat';
import { MongoAppChatLog } from './logs/chatLogsSchema';
import { MongoAppRegistration } from '../../support/appRegistration/schema';
import { MongoMcpKey } from '../../support/mcp/schema';
export const beforeUpdateAppFormat = ({ nodes }: { nodes?: StoreNodeItemType[] }) => {
if (!nodes) return;
......@@ -136,112 +136,71 @@ export const getAppBasicInfoByIds = async ({ teamId, ids }: { teamId: string; id
}));
};
export const onDelOneApp = async ({
teamId,
appId,
session
export const deleteAppDataProcessor = async ({
app,
teamId
}: {
app: AppSchema;
teamId: string;
appId: string;
session?: ClientSession;
}) => {
const apps = await findAppAndAllChildren({
teamId,
appId,
fields: '_id avatar'
});
const deletedAppIds = apps
.filter((app) => !AppFolderTypeList.includes(app.type))
.map((app) => String(app._id));
// Remove eval job
const evalJobs = await MongoEvaluation.find(
{
appId: { $in: apps.map((app) => app._id) }
},
'_id'
).lean();
await Promise.all(evalJobs.map((evalJob) => removeEvaluationJob(evalJob._id)));
const del = async (app: AppSchema, session: ClientSession) => {
const appId = String(app._id);
const appId = String(app._id);
// 1. 删除应用头像
await removeImageByPath(app.avatar);
// 2. 删除聊天记录和S3文件
await getS3ChatSource().deleteChatFilesByPrefix({ appId });
await MongoAppChatLog.deleteMany({ teamId, appId });
await MongoChatItemResponse.deleteMany({ appId });
await MongoChatItem.deleteMany({ appId });
await MongoChat.deleteMany({ appId });
// 3. 删除应用相关数据(使用事务)
{
// 删除分享链接
await MongoOutLink.deleteMany({
appId
}).session(session);
// Openapi
await MongoOpenApi.deleteMany({
appId
}).session(session);
// delete version
await MongoAppVersion.deleteMany({
appId
}).session(session);
await MongoChatInputGuide.deleteMany({
appId
}).session(session);
await MongoOutLink.deleteMany({ appId });
// 删除 OpenAPI 配置
await MongoOpenApi.deleteMany({ appId });
// 删除应用版本
await MongoAppVersion.deleteMany({ appId });
// 删除聊天输入引导
await MongoChatInputGuide.deleteMany({ appId });
// 删除精选应用记录
await MongoChatFavouriteApp.deleteMany({
teamId,
appId
}).session(session);
await MongoChatFavouriteApp.deleteMany({ teamId, appId });
// 从快捷应用中移除对应应用
await MongoChatSetting.updateMany(
{ teamId },
{ $pull: { quickAppIds: { id: String(appId) } } }
).session(session);
// Del permission
await MongoChatSetting.updateMany({ teamId }, { $pull: { quickAppIds: { $in: [appId] } } });
// 删除权限记录
await MongoResourcePermission.deleteMany({
resourceType: PerResourceTypeEnum.app,
teamId,
resourceId: appId
}).session(session);
await MongoAppLogKeys.deleteMany({
appId
}).session(session);
// delete app
await MongoApp.deleteOne(
{
_id: appId
},
{ session }
);
// Delete avatar
await removeImageByPath(app.avatar, session);
};
// Delete chats
for await (const app of apps) {
const appId = String(app._id);
await getS3ChatSource().deleteChatFilesByPrefix({ appId });
await MongoChatItemResponse.deleteMany({
appId
});
await MongoChatItem.deleteMany({
appId
});
await MongoChat.deleteMany({
appId
});
}
// 删除日志密钥
await MongoAppLogKeys.deleteMany({ appId });
for await (const app of apps) {
if (session) {
await del(app, session);
}
// 删除应用注册记录
await MongoAppRegistration.deleteMany({ appId });
// 删除应用从MCP key apps数组中移除
await MongoMcpKey.updateMany({ teamId, 'apps.appId': appId }, { $pull: { apps: { appId } } });
await mongoSessionRun((session) => del(app, session));
// 删除应用本身
await MongoApp.deleteOne({ _id: appId });
}
};
return deletedAppIds;
export const deleteAppsImmediate = async ({
teamId,
appIds
}: {
teamId: string;
appIds: string[];
}) => {
// Remove eval job
const evalJobs = await MongoEvaluation.find(
{
teamId,
appId: { $in: appIds }
},
'_id'
).lean();
await Promise.all(evalJobs.map((evalJob) => removeEvaluationJob(evalJob._id)));
};
import { getQueue, getWorker, QueueNames } from '../../../common/bullmq';
import { appDeleteProcessor } from './processor';
export type AppDeleteJobData = {
teamId: string;
appId: string;
};
// 创建工作进程
export const initAppDeleteWorker = () => {
return getWorker<AppDeleteJobData>(QueueNames.appDelete, appDeleteProcessor, {
concurrency: 1, // 确保同时只有1个删除任务
removeOnFail: {
age: 30 * 24 * 60 * 60 // 保留30天失败记录
}
});
};
// 添加删除任务
export const addAppDeleteJob = (data: AppDeleteJobData) => {
// 创建删除队列
const appDeleteQueue = getQueue<AppDeleteJobData>(QueueNames.appDelete, {
defaultJobOptions: {
attempts: 10,
backoff: {
type: 'exponential',
delay: 5000
},
removeOnComplete: true,
removeOnFail: { age: 30 * 24 * 60 * 60 } // 保留30天失败记录
}
});
const jobId = `${data.teamId}:${data.appId}`;
// Use jobId to automatically prevent duplicate deletion tasks (BullMQ feature)
return appDeleteQueue.add('deleteapp', data, {
jobId,
delay: 1000 // Delay 1 second to ensure API response completes
});
};
import type { Processor } from 'bullmq';
import type { AppDeleteJobData } from './index';
import { findAppAndAllChildren, deleteAppDataProcessor } from '../controller';
import { addLog } from '../../../common/system/log';
import { batchRun } from '@fastgpt/global/common/system/utils';
import type { AppSchema } from '@fastgpt/global/core/app/type';
import { MongoApp } from '../schema';
const deleteApps = async ({ teamId, apps }: { teamId: string; apps: AppSchema[] }) => {
const results = await batchRun(
apps,
async (app) => {
await deleteAppDataProcessor({ app, teamId });
},
3
);
return results.flat();
};
export const appDeleteProcessor: Processor<AppDeleteJobData> = async (job) => {
const { teamId, appId } = job.data;
const startTime = Date.now();
addLog.info(`[App Delete] Start deleting app: ${appId} for team: ${teamId}`);
try {
// 1. 查找应用及其所有子应用
const apps = await findAppAndAllChildren({
teamId,
appId
});
if (!apps || apps.length === 0) {
addLog.warn(`[App Delete] App not found: ${appId}`);
return;
}
// 2. 安全检查:确保所有要删除的应用都已标记为 deleteTime
const markedForDelete = await MongoApp.find(
{
_id: { $in: apps.map((app) => app._id) },
teamId,
deleteTime: { $ne: null }
},
{ _id: 1 }
).lean();
if (markedForDelete.length !== apps.length) {
addLog.warn(
`[App Delete] Safety check: ${markedForDelete.length}/${apps.length} apps marked for deletion`,
{
markedAppIds: markedForDelete.map((app) => app._id),
totalAppIds: apps.map((app) => app._id)
}
);
}
const childrenLen = apps.length - 1;
const appIds = apps.map((app) => app._id);
// 3. 执行真正的删除操作(只删除已经标记为 deleteTime 的数据)
await deleteApps({
teamId,
apps
});
addLog.info(`[App Delete] Successfully deleted app: ${appId} and ${childrenLen} children`, {
duration: Date.now() - startTime,
totalApps: appIds.length,
appIds
});
} catch (error: any) {
addLog.error(`[App Delete] Failed to delete app: ${appId}`, error);
throw error;
}
};
......@@ -119,6 +119,12 @@ const AppSchema = new Schema(
inited: Boolean,
teamTags: {
type: [String]
},
// 软删除标记字段
deleteTime: {
type: Date,
default: null // null表示未删除,有值表示删除时间
}
},
{
......@@ -138,5 +144,6 @@ AppSchema.index(
);
// Admin count
AppSchema.index({ type: 1 });
AppSchema.index({ deleteTime: 1 });
export const MongoApp = getMongoModel<AppType>(AppCollectionName, AppSchema);
import { NextAPI } from '@/service/middleware/entry';
import { batchRun } from '@fastgpt/global/common/system/utils';
import { getQueue, QueueNames } from '@fastgpt/service/common/bullmq';
import type { S3MQJobData } from '@fastgpt/service/common/s3/mq';
import { addLog } from '@fastgpt/service/common/system/log';
import type { ApiRequestProps, ApiResponseType } from '@fastgpt/service/type/next';
import { authCert } from '@fastgpt/service/support/permission/auth/common';
export type ResponseType = {
message: string;
retriedCount: number;
failedCount: number;
};
async function handler(
req: ApiRequestProps,
res: ApiResponseType<ResponseType>
): Promise<ResponseType> {
await authCert({ req, authRoot: true });
const queue = getQueue<S3MQJobData>(QueueNames.s3FileDelete);
// Get all failed jobs and retry them
const failedJobs = await queue.getFailed();
console.log(`Found ${failedJobs.length} failed jobs`);
let retriedCount = 0;
await batchRun(
failedJobs,
async (job) => {
addLog.debug(`Retrying job with 3 new attempts`, { retriedCount });
try {
// Remove old job and recreate with new attempts
const jobData = job.data;
await job.remove();
// Add new job with 3 more attempts
await queue.add('delete-s3-files', jobData, {
attempts: 10,
removeOnFail: {
count: 10000, // 保留10000个失败任务
age: 14 * 24 * 60 * 60 // 14 days
},
removeOnComplete: true,
backoff: {
delay: 2000,
type: 'exponential'
}
});
retriedCount++;
console.log(`Retried job ${job.id} with 3 new attempts`);
} catch (error) {
console.error(`Failed to retry job ${job.id}:`, error);
}
},
100
);
return {
message: 'Successfully retried all failed S3 delete jobs with 3 new attempts',
retriedCount,
failedCount: failedJobs.length
};
}
export default NextAPI(handler);
......@@ -2,7 +2,11 @@ import type { NextApiRequest, NextApiResponse } from 'next';
import { authApp } from '@fastgpt/service/support/permission/app/auth';
import { NextAPI } from '@/service/middleware/entry';
import { OwnerPermissionVal } from '@fastgpt/global/support/permission/constant';
import { onDelOneApp } from '@fastgpt/service/core/app/controller';
import { findAppAndAllChildren } from '@fastgpt/service/core/app/controller';
import { MongoApp } from '@fastgpt/service/core/app/schema';
import { mongoSessionRun } from '@fastgpt/service/common/mongo/sessionRun';
import { addAppDeleteJob } from '@fastgpt/service/core/app/delete';
import { deleteAppsImmediate } from '@fastgpt/service/core/app/controller';
import { pushTrack } from '@fastgpt/service/common/middle/tracks/utils';
import { addAuditLog } from '@fastgpt/service/support/user/audit/util';
import { AuditEventEnum } from '@fastgpt/global/support/user/audit/constants';
......@@ -23,7 +27,31 @@ async function handler(req: NextApiRequest, res: NextApiResponse<string[]>) {
per: OwnerPermissionVal
});
const deletedAppIds = await onDelOneApp({ teamId, appId });
const deleteAppsList = await findAppAndAllChildren({
teamId,
appId
});
await mongoSessionRun(async (session) => {
// Mark app as deleted
await MongoApp.updateMany(
{ _id: deleteAppsList.map((app) => app._id), teamId },
{ deleteTime: new Date() },
{ session }
);
// Stop background tasks immediately
await deleteAppsImmediate({
teamId,
appIds: deleteAppsList.map((app) => app._id)
});
// Add to delete queue for async cleanup
await addAppDeleteJob({
teamId,
appId
});
});
(async () => {
addAuditLog({
......@@ -40,7 +68,9 @@ async function handler(req: NextApiRequest, res: NextApiResponse<string[]>) {
// Tracks
pushTrack.countAppNodes({ teamId, tmbId, uid: userId, appId });
return deletedAppIds;
return deleteAppsList
.filter((app) => !['folder'].includes(app.type))
.map((app) => String(app._id));
}
export default NextAPI(handler);
......@@ -11,7 +11,6 @@ import { type ApiRequestProps } from '@fastgpt/service/type/next';
import { type ParentIdType } from '@fastgpt/global/common/parentFolder/type';
import { parseParentIdInMongo } from '@fastgpt/global/common/parentFolder/utils';
import { AppFolderTypeList, AppTypeEnum } from '@fastgpt/global/core/app/constants';
import { AppDefaultRoleVal } from '@fastgpt/global/support/permission/app/constant';
import { authApp } from '@fastgpt/service/support/permission/app/auth';
import { authUserPer } from '@fastgpt/service/support/permission/user/auth';
import { replaceRegChars } from '@fastgpt/global/common/string/tools';
......@@ -97,7 +96,7 @@ async function handler(req: ApiRequestProps<ListAppBody>): Promise<AppListItemTy
const findAppsQuery = (() => {
if (getRecentlyChat) {
return {
// get all chat app, excluding hidden apps
// get all chat app, excluding hidden apps and deleted apps
teamId,
type: { $in: [AppTypeEnum.workflow, AppTypeEnum.simple, AppTypeEnum.workflowTool] }
};
......@@ -160,7 +159,7 @@ async function handler(req: ApiRequestProps<ListAppBody>): Promise<AppListItemTy
})();
const myApps = await MongoApp.find(
findAppsQuery,
{ ...findAppsQuery, deleteTime: null },
'_id parentId avatar type name intro tmbId updateTime pluginData inheritPermission modules',
{
limit: limit
......
import { addLog } from '@fastgpt/service/common/system/log';
import { initS3MQWorker } from '@fastgpt/service/common/s3';
import { initDatasetDeleteWorker } from '@fastgpt/service/core/dataset/delete';
import { initAppDeleteWorker } from '@fastgpt/service/core/app/delete';
export const initBullMQWorkers = () => {
addLog.info('Init BullMQ Workers...');
initS3MQWorker();
initDatasetDeleteWorker();
initAppDeleteWorker();
};
......@@ -13,7 +13,7 @@ import { getUser } from '@test/datas/users';
import { Call } from '@test/utils/request';
import { describe, expect, it, beforeEach } from 'vitest';
describe('closeCustom api test', () => {
describe.sequential('closeCustom api test', () => {
let testUser: Awaited<ReturnType<typeof getUser>>;
let appId: string;
let chatId: string;
......
......@@ -13,7 +13,7 @@ import { getUser } from '@test/datas/users';
import { Call } from '@test/utils/request';
import { describe, expect, it, beforeEach } from 'vitest';
describe('updateFeedbackReadStatus api test', () => {
describe.sequential('updateFeedbackReadStatus api test', () => {
let testUser: Awaited<ReturnType<typeof getUser>>;
let appId: string;
let chatId: string;
......
......@@ -14,7 +14,7 @@ import { getUser } from '@test/datas/users';
import { Call } from '@test/utils/request';
import { describe, expect, it, beforeEach } from 'vitest';
describe('updateUserFeedback api test', () => {
describe.sequential('updateUserFeedback api test', () => {
let testUser: Awaited<ReturnType<typeof getUser>>;
let appId: string;
let chatId: string;
......
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