feat(system): expand runtime queue task inspection and safe actions

Replace failed-only task views with state-aware runtime task listing,
operator-safe cancel/run-now/delete actions, orphan detection, and
multimodal queue cleanup for deleted knowledge rows.
This commit is contained in:
wizardchen
2026-07-14 15:17:45 +08:00
committed by lyingbug
parent 1eae475a1f
commit 4ae48d6acf
22 changed files with 2337 additions and 481 deletions
+261 -2
View File
@@ -11316,6 +11316,113 @@ const docTemplate = `{
}
}
},
"/system/admin/runtime/queues/{queue}/tasks": {
"get": {
"produces": [
"application/json"
],
"tags": [
"System Admin"
],
"summary": "List runtime queue tasks by state",
"parameters": [
{
"type": "string",
"description": "Queue name",
"name": "queue",
"in": "path",
"required": true
},
{
"enum": [
"pending",
"active",
"scheduled",
"retry",
"archived",
"completed"
],
"type": "string",
"description": "Task state",
"name": "state",
"in": "query",
"required": true
},
{
"type": "integer",
"default": 1,
"description": "Page",
"name": "page",
"in": "query"
},
{
"type": "integer",
"default": 20,
"description": "Page size",
"name": "page_size",
"in": "query"
}
],
"responses": {
"200": {
"description": "OK",
"schema": {
"$ref": "#/definitions/internal_handler.RuntimeTasksResponse"
}
}
}
}
},
"/system/admin/runtime/queues/{queue}/tasks/{task_id}/actions/{action}": {
"post": {
"produces": [
"application/json"
],
"tags": [
"System Admin"
],
"summary": "Run a safe runtime task action",
"parameters": [
{
"type": "string",
"description": "Queue name",
"name": "queue",
"in": "path",
"required": true
},
{
"type": "string",
"description": "Task ID",
"name": "task_id",
"in": "path",
"required": true
},
{
"enum": [
"cancel",
"run_now",
"delete"
],
"type": "string",
"description": "Action",
"name": "action",
"in": "path",
"required": true
}
],
"responses": {
"200": {
"description": "OK",
"schema": {
"type": "object",
"additionalProperties": {
"type": "boolean"
}
}
}
}
}
},
"/system/admin/settings/{key}": {
"get": {
"description": "Returns the row matching :key. 404 when the key is unknown\nto the registry; 200 with the row when known.",
@@ -14170,7 +14277,11 @@ const docTemplate = `{
"system.setting_changed",
"system.admin_promoted",
"system.admin_revoked",
"system.user_password_reset"
"system.user_password_reset",
"system.queue_task_retried",
"system.queue_task_deleted",
"system.queue_task_run_now",
"system.queue_task_cancelled"
],
"x-enum-varnames": [
"AuditActionMemberAdded",
@@ -14192,7 +14303,11 @@ const docTemplate = `{
"AuditActionSystemSettingChanged",
"AuditActionSystemAdminPromoted",
"AuditActionSystemAdminRevoked",
"AuditActionSystemUserPasswordReset"
"AuditActionSystemUserPasswordReset",
"AuditActionSystemQueueTaskRetried",
"AuditActionSystemQueueTaskDeleted",
"AuditActionSystemQueueTaskRunNow",
"AuditActionSystemQueueTaskCancelled"
]
},
"github_com_Tencent_WeKnora_internal_types.AuditLog": {
@@ -17324,6 +17439,127 @@ const docTemplate = `{
}
}
},
"github_com_Tencent_WeKnora_internal_types.RuntimeTaskAction": {
"type": "string",
"enum": [
"cancel",
"run_now",
"delete"
],
"x-enum-varnames": [
"RuntimeTaskActionCancel",
"RuntimeTaskActionRunNow",
"RuntimeTaskActionDelete"
]
},
"github_com_Tencent_WeKnora_internal_types.RuntimeTaskInfo": {
"type": "object",
"properties": {
"allowed_actions": {
"type": "array",
"items": {
"$ref": "#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskAction"
}
},
"completed_at": {
"type": "string"
},
"data_source_id": {
"type": "string"
},
"deadline": {
"type": "string"
},
"enqueued_at": {
"type": "string"
},
"id": {
"type": "string"
},
"is_orphaned": {
"type": "boolean"
},
"knowledge_base_id": {
"type": "string"
},
"knowledge_count": {
"type": "integer"
},
"knowledge_id": {
"type": "string"
},
"last_error": {
"type": "string"
},
"last_failed_at": {
"type": "string"
},
"max_retry": {
"type": "integer"
},
"next_process_at": {
"type": "string"
},
"queue": {
"type": "string"
},
"retried": {
"type": "integer"
},
"source_id": {
"type": "string"
},
"source_kb_id": {
"type": "string"
},
"started_at": {
"type": "string"
},
"state": {
"$ref": "#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskState"
},
"sync_log_id": {
"type": "string"
},
"target_id": {
"type": "string"
},
"target_kb_id": {
"type": "string"
},
"task_id": {
"type": "string"
},
"tenant_id": {
"type": "integer"
},
"type": {
"type": "string"
},
"worker": {
"type": "string"
}
}
},
"github_com_Tencent_WeKnora_internal_types.RuntimeTaskState": {
"type": "string",
"enum": [
"pending",
"active",
"scheduled",
"retry",
"archived",
"completed"
],
"x-enum-varnames": [
"RuntimeTaskPending",
"RuntimeTaskActive",
"RuntimeTaskScheduled",
"RuntimeTaskRetry",
"RuntimeTaskArchived",
"RuntimeTaskCompleted"
]
},
"github_com_Tencent_WeKnora_internal_types.S3EngineConfig": {
"type": "object",
"properties": {
@@ -20140,6 +20376,29 @@ const docTemplate = `{
}
}
},
"internal_handler.RuntimeTasksResponse": {
"type": "object",
"properties": {
"available": {
"type": "boolean"
},
"has_more": {
"type": "boolean"
},
"page": {
"type": "integer"
},
"page_size": {
"type": "integer"
},
"tasks": {
"type": "array",
"items": {
"$ref": "#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskInfo"
}
}
}
},
"internal_handler.RuntimeWorkerPool": {
"description": "Return every row in the system_settings table (system-scope, not tenant-scope). SystemAdmin only.",
"type": "object",
+261 -2
View File
@@ -11309,6 +11309,113 @@
}
}
},
"/system/admin/runtime/queues/{queue}/tasks": {
"get": {
"produces": [
"application/json"
],
"tags": [
"System Admin"
],
"summary": "List runtime queue tasks by state",
"parameters": [
{
"type": "string",
"description": "Queue name",
"name": "queue",
"in": "path",
"required": true
},
{
"enum": [
"pending",
"active",
"scheduled",
"retry",
"archived",
"completed"
],
"type": "string",
"description": "Task state",
"name": "state",
"in": "query",
"required": true
},
{
"type": "integer",
"default": 1,
"description": "Page",
"name": "page",
"in": "query"
},
{
"type": "integer",
"default": 20,
"description": "Page size",
"name": "page_size",
"in": "query"
}
],
"responses": {
"200": {
"description": "OK",
"schema": {
"$ref": "#/definitions/internal_handler.RuntimeTasksResponse"
}
}
}
}
},
"/system/admin/runtime/queues/{queue}/tasks/{task_id}/actions/{action}": {
"post": {
"produces": [
"application/json"
],
"tags": [
"System Admin"
],
"summary": "Run a safe runtime task action",
"parameters": [
{
"type": "string",
"description": "Queue name",
"name": "queue",
"in": "path",
"required": true
},
{
"type": "string",
"description": "Task ID",
"name": "task_id",
"in": "path",
"required": true
},
{
"enum": [
"cancel",
"run_now",
"delete"
],
"type": "string",
"description": "Action",
"name": "action",
"in": "path",
"required": true
}
],
"responses": {
"200": {
"description": "OK",
"schema": {
"type": "object",
"additionalProperties": {
"type": "boolean"
}
}
}
}
}
},
"/system/admin/settings/{key}": {
"get": {
"description": "Returns the row matching :key. 404 when the key is unknown\nto the registry; 200 with the row when known.",
@@ -14163,7 +14270,11 @@
"system.setting_changed",
"system.admin_promoted",
"system.admin_revoked",
"system.user_password_reset"
"system.user_password_reset",
"system.queue_task_retried",
"system.queue_task_deleted",
"system.queue_task_run_now",
"system.queue_task_cancelled"
],
"x-enum-varnames": [
"AuditActionMemberAdded",
@@ -14185,7 +14296,11 @@
"AuditActionSystemSettingChanged",
"AuditActionSystemAdminPromoted",
"AuditActionSystemAdminRevoked",
"AuditActionSystemUserPasswordReset"
"AuditActionSystemUserPasswordReset",
"AuditActionSystemQueueTaskRetried",
"AuditActionSystemQueueTaskDeleted",
"AuditActionSystemQueueTaskRunNow",
"AuditActionSystemQueueTaskCancelled"
]
},
"github_com_Tencent_WeKnora_internal_types.AuditLog": {
@@ -17317,6 +17432,127 @@
}
}
},
"github_com_Tencent_WeKnora_internal_types.RuntimeTaskAction": {
"type": "string",
"enum": [
"cancel",
"run_now",
"delete"
],
"x-enum-varnames": [
"RuntimeTaskActionCancel",
"RuntimeTaskActionRunNow",
"RuntimeTaskActionDelete"
]
},
"github_com_Tencent_WeKnora_internal_types.RuntimeTaskInfo": {
"type": "object",
"properties": {
"allowed_actions": {
"type": "array",
"items": {
"$ref": "#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskAction"
}
},
"completed_at": {
"type": "string"
},
"data_source_id": {
"type": "string"
},
"deadline": {
"type": "string"
},
"enqueued_at": {
"type": "string"
},
"id": {
"type": "string"
},
"is_orphaned": {
"type": "boolean"
},
"knowledge_base_id": {
"type": "string"
},
"knowledge_count": {
"type": "integer"
},
"knowledge_id": {
"type": "string"
},
"last_error": {
"type": "string"
},
"last_failed_at": {
"type": "string"
},
"max_retry": {
"type": "integer"
},
"next_process_at": {
"type": "string"
},
"queue": {
"type": "string"
},
"retried": {
"type": "integer"
},
"source_id": {
"type": "string"
},
"source_kb_id": {
"type": "string"
},
"started_at": {
"type": "string"
},
"state": {
"$ref": "#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskState"
},
"sync_log_id": {
"type": "string"
},
"target_id": {
"type": "string"
},
"target_kb_id": {
"type": "string"
},
"task_id": {
"type": "string"
},
"tenant_id": {
"type": "integer"
},
"type": {
"type": "string"
},
"worker": {
"type": "string"
}
}
},
"github_com_Tencent_WeKnora_internal_types.RuntimeTaskState": {
"type": "string",
"enum": [
"pending",
"active",
"scheduled",
"retry",
"archived",
"completed"
],
"x-enum-varnames": [
"RuntimeTaskPending",
"RuntimeTaskActive",
"RuntimeTaskScheduled",
"RuntimeTaskRetry",
"RuntimeTaskArchived",
"RuntimeTaskCompleted"
]
},
"github_com_Tencent_WeKnora_internal_types.S3EngineConfig": {
"type": "object",
"properties": {
@@ -20133,6 +20369,29 @@
}
}
},
"internal_handler.RuntimeTasksResponse": {
"type": "object",
"properties": {
"available": {
"type": "boolean"
},
"has_more": {
"type": "boolean"
},
"page": {
"type": "integer"
},
"page_size": {
"type": "integer"
},
"tasks": {
"type": "array",
"items": {
"$ref": "#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskInfo"
}
}
}
},
"internal_handler.RuntimeWorkerPool": {
"description": "Return every row in the system_settings table (system-scope, not tenant-scope). SystemAdmin only.",
"type": "object",
+182
View File
@@ -294,6 +294,10 @@ definitions:
- system.admin_promoted
- system.admin_revoked
- system.user_password_reset
- system.queue_task_retried
- system.queue_task_deleted
- system.queue_task_run_now
- system.queue_task_cancelled
type: string
x-enum-varnames:
- AuditActionMemberAdded
@@ -316,6 +320,10 @@ definitions:
- AuditActionSystemAdminPromoted
- AuditActionSystemAdminRevoked
- AuditActionSystemUserPasswordReset
- AuditActionSystemQueueTaskRetried
- AuditActionSystemQueueTaskDeleted
- AuditActionSystemQueueTaskRunNow
- AuditActionSystemQueueTaskCancelled
github_com_Tencent_WeKnora_internal_types.AuditLog:
properties:
action:
@@ -2739,6 +2747,91 @@ definitions:
description: 'Optional: role to assign when approving; overrides applicant''s
requested role'
type: object
github_com_Tencent_WeKnora_internal_types.RuntimeTaskAction:
enum:
- cancel
- run_now
- delete
type: string
x-enum-varnames:
- RuntimeTaskActionCancel
- RuntimeTaskActionRunNow
- RuntimeTaskActionDelete
github_com_Tencent_WeKnora_internal_types.RuntimeTaskInfo:
properties:
allowed_actions:
items:
$ref: '#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskAction'
type: array
completed_at:
type: string
data_source_id:
type: string
deadline:
type: string
enqueued_at:
type: string
id:
type: string
is_orphaned:
type: boolean
knowledge_base_id:
type: string
knowledge_count:
type: integer
knowledge_id:
type: string
last_error:
type: string
last_failed_at:
type: string
max_retry:
type: integer
next_process_at:
type: string
queue:
type: string
retried:
type: integer
source_id:
type: string
source_kb_id:
type: string
started_at:
type: string
state:
$ref: '#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskState'
sync_log_id:
type: string
target_id:
type: string
target_kb_id:
type: string
task_id:
type: string
tenant_id:
type: integer
type:
type: string
worker:
type: string
type: object
github_com_Tencent_WeKnora_internal_types.RuntimeTaskState:
enum:
- pending
- active
- scheduled
- retry
- archived
- completed
type: string
x-enum-varnames:
- RuntimeTaskPending
- RuntimeTaskActive
- RuntimeTaskScheduled
- RuntimeTaskRetry
- RuntimeTaskArchived
- RuntimeTaskCompleted
github_com_Tencent_WeKnora_internal_types.S3EngineConfig:
properties:
access_key:
@@ -4832,6 +4925,21 @@ definitions:
description: compatibility field
type: integer
type: object
internal_handler.RuntimeTasksResponse:
properties:
available:
type: boolean
has_more:
type: boolean
page:
type: integer
page_size:
type: integer
tasks:
items:
$ref: '#/definitions/github_com_Tencent_WeKnora_internal_types.RuntimeTaskInfo'
type: array
type: object
internal_handler.RuntimeWorkerPool:
description: Return every row in the system_settings table (system-scope, not
tenant-scope). SystemAdmin only.
@@ -12614,6 +12722,80 @@ paths:
summary: 获取解析任务队列运行时状态
tags:
- 系统管理
/system/admin/runtime/queues/{queue}/tasks:
get:
parameters:
- description: Queue name
in: path
name: queue
required: true
type: string
- description: Task state
enum:
- pending
- active
- scheduled
- retry
- archived
- completed
in: query
name: state
required: true
type: string
- default: 1
description: Page
in: query
name: page
type: integer
- default: 20
description: Page size
in: query
name: page_size
type: integer
produces:
- application/json
responses:
"200":
description: OK
schema:
$ref: '#/definitions/internal_handler.RuntimeTasksResponse'
summary: List runtime queue tasks by state
tags:
- System Admin
/system/admin/runtime/queues/{queue}/tasks/{task_id}/actions/{action}:
post:
parameters:
- description: Queue name
in: path
name: queue
required: true
type: string
- description: Task ID
in: path
name: task_id
required: true
type: string
- description: Action
enum:
- cancel
- run_now
- delete
in: path
name: action
required: true
type: string
produces:
- application/json
responses:
"200":
description: OK
schema:
additionalProperties:
type: boolean
type: object
summary: Run a safe runtime task action
tags:
- System Admin
/system/admin/settings/{key}:
delete:
description: |-
+35 -17
View File
@@ -555,23 +555,42 @@ export interface RuntimeQueuesResponse {
timestamp: number
}
export interface RuntimeFailedTask {
export type RuntimeTaskState = 'pending' | 'active' | 'scheduled' | 'retry' | 'archived' | 'completed'
export type RuntimeTaskAction = 'cancel' | 'run_now' | 'delete'
export interface RuntimeTask {
id: string
queue: string
type: string
last_error: string
last_failed_at: string
state: RuntimeTaskState
allowed_actions: RuntimeTaskAction[]
last_error?: string
last_failed_at?: string
next_process_at?: string
started_at?: string
completed_at?: string
deadline?: string
enqueued_at?: string
retried: number
max_retry: number
is_orphaned?: boolean
worker?: string
tenant_id?: number
knowledge_base_id?: string
knowledge_id?: string
task_id?: string
source_id?: string
target_id?: string
source_kb_id?: string
target_kb_id?: string
data_source_id?: string
sync_log_id?: string
knowledge_count?: number
}
export interface RuntimeFailedTasksResponse {
export interface RuntimeTasksResponse {
available: boolean
tasks: RuntimeFailedTask[]
tasks: RuntimeTask[]
page: number
page_size: number
has_more: boolean
@@ -588,24 +607,23 @@ export async function getRuntimeQueues(): Promise<RuntimeQueuesResponse> {
return response as unknown as RuntimeQueuesResponse
}
export async function getRuntimeFailedTasks(
export async function getRuntimeTasks(
queue: string,
state: RuntimeTaskState,
page = 1,
pageSize = 20,
): Promise<RuntimeFailedTasksResponse> {
return get(`/api/v1/system/admin/runtime/queues/${encodeURIComponent(queue)}/failed-tasks`, {
params: { page, page_size: pageSize },
): Promise<RuntimeTasksResponse> {
return get(`/api/v1/system/admin/runtime/queues/${encodeURIComponent(queue)}/tasks`, {
params: { state, page, page_size: pageSize },
})
}
export async function retryRuntimeFailedTask(queue: string, taskID: string): Promise<void> {
export async function mutateRuntimeTask(
queue: string,
taskID: string,
action: RuntimeTaskAction,
): Promise<void> {
await post(
`/api/v1/system/admin/runtime/queues/${encodeURIComponent(queue)}/failed-tasks/${encodeURIComponent(taskID)}/retry`,
)
}
export async function deleteRuntimeFailedTask(queue: string, taskID: string): Promise<void> {
await del(
`/api/v1/system/admin/runtime/queues/${encodeURIComponent(queue)}/failed-tasks/${encodeURIComponent(taskID)}`,
`/api/v1/system/admin/runtime/queues/${encodeURIComponent(queue)}/tasks/${encodeURIComponent(taskID)}/actions/${encodeURIComponent(action)}`,
)
}
+51 -26
View File
@@ -3545,7 +3545,7 @@ export default {
},
runtime: {
title: "Task Queue Runtime",
description: "Live background-queue load and per-process capacity for each isolated worker pool. Read-only; auto-refreshes every 5s.",
description: "Live background-queue load and per-process capacity for each isolated worker pool. Includes task details and safe controls; auto-refreshes every 5s.",
refresh: "Refresh",
autoRefresh: "Auto-refresh (every 5s)",
loading: "Loading...",
@@ -3584,6 +3584,7 @@ export default {
scheduled: "Scheduled",
retry: "Retry",
archived: "Failed",
completed: "Completed",
latency: "Oldest wait",
status: "Status",
},
@@ -3600,43 +3601,65 @@ export default {
title: "{count} failed tasks need attention",
description: "Open a red Failed count in the table below to inspect causes, then retry after fixing the root issue.",
},
failedTasks: {
title: "Failed tasks · {queue}",
description: "Tasks that stopped after exceeding their automatic retry limit, newest failure first.",
guideTitle: "Fix the cause before running again",
guideDescription: "Run once performs one immediate attempt without resetting the retry counter. Clear record only removes it from this queue; it does not complete the original task, and historical logs remain.",
listTitle: "Failed records",
openAria: "View {count} failed tasks in {queue}",
unavailable: "Failed-task details are unavailable in this deployment",
empty: "This queue has no stopped failed tasks",
tasks: {
title: "Task details · {queue}",
description: "Inspect every task state and use only the actions the backend marks safe for the current task.",
listTitle: "{state} tasks",
openAria: "View {count} {state} tasks in {queue}",
unavailable: "Task details are unavailable in this deployment",
empty: "This queue has no {state} tasks",
loadError: "Failed to load task details",
loadMore: "Load more",
loadedSummary: "{count} loaded — scroll or tap to load more",
loadedAll: "All {count} loaded",
loadingMore: "Loading more…",
attempts: "Retry count {count}",
attemptsShort: "{count}×",
showError: "View error",
hideError: "Hide error",
lastError: "Last error",
noError: "No error details recorded",
attempts: "Attempt {current}/{max}",
unknownTarget: "No related object identified",
knowledgeBaseLabel: "Knowledge base ID",
knowledgeLabel: "Document ID",
taskIDLabel: "Business task ID",
tenantLabel: "Tenant ID",
tenant: "Workspace {id}",
knowledgeBase: "Knowledge base {id}",
knowledge: "Document {id}",
taskID: "Business task {id}",
retryOnce: "Run once",
retryConfirm: "Confirm that the cause has been fixed. This task will run once immediately.",
retrySuccess: "Task returned to the queue",
retryError: "Failed to run the task",
sourceLabel: "Source ID",
targetLabel: "Target ID",
sourceKBLabel: "Source knowledge base",
targetKBLabel: "Target knowledge base",
dataSourceLabel: "Data source ID",
syncLogLabel: "Sync log ID",
knowledgeCountLabel: "Documents",
enqueuedAt: "Enqueued",
startedAt: "Started",
nextProcessAt: "Next run",
lastFailedAt: "Last failed",
completedAt: "Completed",
deadline: "Deadline",
worker: "Worker",
health: "Runtime health",
orphaned: "Worker heartbeat lost; waiting for recovery",
cancel: "Cancel task",
cancelConfirm: "This uses the business cancellation flow and also stops related tasks for the same document. Cancel it?",
runNow: "Run now",
runNowConfirm: "This task will move to pending immediately without resetting its retry count. Continue?",
deleteRecord: "Clear record",
deleteConfirm: "This only removes the record from this queue. It will not run or complete the original task, and historical logs remain. Clear it?",
deleteSuccess: "Failure record cleared",
deleteError: "Failed to clear the record",
stateFilter: "Filter by task state",
states: {
active: "Active",
pending: "Pending",
scheduled: "Scheduled",
retry: "Retrying",
archived: "Failed",
completed: "Completed",
},
guides: {
active: "Inspect the worker, start time, deadline, and orphan status. Cancel is offered only when the business state can be updated safely.",
pending: "Pending tasks have not been claimed. Cancellation updates business state instead of only removing Redis data.",
scheduled: "Scheduled tasks can run early; document pipeline tasks can also use safe business cancellation.",
retry: "Use the last error, attempt count, and next run time to decide whether to run now or cancel.",
archived: "Fix the root cause before running again. Clearing a record does not complete the original business task.",
completed: "Only recently completed tasks with result retention are shown. No management actions are available.",
},
actionSuccess: { cancel: "Task cancelled", run_now: "Task moved to pending", delete: "Failure record cleared" },
actionError: { cancel: "Failed to cancel task", run_now: "Failed to run task", delete: "Failed to clear record" },
taskTypes: {
documentProcess: "Document parsing",
manualProcess: "Manual reprocessing",
@@ -3903,6 +3926,8 @@ export default {
'system.user_password_reset': 'User password reset',
'system.queue_task_retried': 'Failed task run again',
'system.queue_task_deleted': 'Failed task record cleared',
'system.queue_task_run_now': 'Queue task run now',
'system.queue_task_cancelled': 'Queue task cancelled',
},
outcome: {
success: 'Success',
+42 -7
View File
@@ -2569,6 +2569,7 @@ export default {
scheduled: "예약",
retry: "재시도",
archived: "실패",
completed: "완료",
latency: "최장 대기",
status: "상태",
},
@@ -2585,21 +2586,21 @@ export default {
title: "처리 대기 중인 실패 작업 {count}개",
description: "아래 표의 빨간색 '최종 실패' 숫자를 눌러 원인을 확인한 뒤, 수정 후 수동으로 재시도하세요.",
},
failedTasks: {
title: "최종 실패 · {queue}",
description: "자동 재시도 한도를 초과하여 중지된 작업이며 최근 실패 순으로 표시됩니다.",
tasks: {
title: "작업 상세 · {queue}",
description: "모든 작업 상태를 확인하고 백엔드가 현재 상태에 안전하다고 표시한 작업만 실행합니다.",
guideTitle: "원인을 해결한 뒤 다시 실행하세요",
guideDescription: "‘한 번 다시 실행’은 재시도 횟수를 초기화하지 않고 즉시 한 번 실행합니다. ‘기록 지우기’는 현재 큐에서만 제거하며 원래 작업을 완료하지 않고 기록 로그는 유지됩니다.",
listTitle: "실패 기록",
openAria: "{queue}의 최종 실패 작업 {count}개 보기",
listTitle: "{state} 작업",
openAria: "{queue}의 {state} 작업 {count}개 보기",
unavailable: "현재 배포에서는 실패 작업 상세 정보를 볼 수 없습니다",
empty: "이 큐에는 최종 실패 작업이 없습니다",
empty: "이 큐에는 {state} 작업이 없습니다",
loadError: "실패 작업 상세 정보를 불러오지 못했습니다",
loadMore: "더 보기",
loadedSummary: "{count}개 로드됨 · 스크롤하거나 눌러 더 보기",
loadedAll: "총 {count}개 모두 로드됨",
loadingMore: "더 불러오는 중…",
attempts: "재시도 횟수 {count}",
attempts: "실행 {current}/{max}",
attemptsShort: "{count}회",
showError: "오류 보기",
hideError: "오류 접기",
@@ -2610,6 +2611,38 @@ export default {
knowledgeLabel: "문서 ID",
taskIDLabel: "비즈니스 작업 ID",
tenantLabel: "테넌트 ID",
sourceLabel: "소스 ID",
targetLabel: "대상 ID",
sourceKBLabel: "소스 지식 베이스",
targetKBLabel: "대상 지식 베이스",
dataSourceLabel: "데이터 소스 ID",
syncLogLabel: "동기화 기록 ID",
knowledgeCountLabel: "문서 수",
enqueuedAt: "큐 등록 시간",
startedAt: "시작 시간",
nextProcessAt: "다음 실행",
lastFailedAt: "마지막 실패",
completedAt: "완료 시간",
deadline: "마감 시간",
worker: "실행 인스턴스",
health: "실행 상태",
orphaned: "실행 인스턴스 연결 끊김; 복구 대기 중",
cancel: "작업 취소",
cancelConfirm: "비즈니스 취소 흐름을 사용하여 같은 문서의 관련 작업도 중지합니다. 취소할까요?",
runNow: "지금 실행",
runNowConfirm: "재시도 횟수를 초기화하지 않고 즉시 대기 상태로 이동합니다. 계속할까요?",
stateFilter: "작업 상태별 필터",
states: { active: "실행 중", pending: "대기 중", scheduled: "예약", retry: "재시도 중", archived: "최종 실패", completed: "완료" },
guides: {
active: "실행 인스턴스와 시작 시간, 마감 및 고립 상태를 확인합니다. 안전한 비즈니스 취소가 가능한 작업만 취소할 수 있습니다.",
pending: "아직 worker가 가져가지 않은 작업입니다. 취소 시 비즈니스 상태도 함께 갱신됩니다.",
scheduled: "예약 작업을 즉시 실행하거나 지원되는 문서 작업을 안전하게 취소할 수 있습니다.",
retry: "마지막 오류, 시도 횟수 및 다음 실행 시간을 확인하세요.",
archived: "원인을 해결한 후 다시 실행하세요. 기록 삭제는 원래 작업을 완료하지 않습니다.",
completed: "보존 기간이 설정된 최근 완료 작업만 표시되며 관리 작업은 없습니다.",
},
actionSuccess: { cancel: "작업을 취소했습니다", run_now: "작업이 대기열로 이동했습니다", delete: "실패 기록을 지웠습니다" },
actionError: { cancel: "작업을 취소하지 못했습니다", run_now: "작업을 실행하지 못했습니다", delete: "실패 기록을 지우지 못했습니다" },
tenant: "공간 {id}",
knowledgeBase: "지식 베이스 {id}",
knowledge: "문서 {id}",
@@ -2888,6 +2921,8 @@ export default {
"system.user_password_reset": "사용자 비밀번호 재설정",
"system.queue_task_retried": "실패 작업 다시 실행",
"system.queue_task_deleted": "실패 작업 기록 삭제",
"system.queue_task_run_now": "큐 작업 즉시 실행",
"system.queue_task_cancelled": "큐 작업 취소",
},
outcome: {
success: "성공",
+43 -8
View File
@@ -2266,6 +2266,7 @@ export default {
scheduled: "Запланир.",
retry: "Повтор",
archived: "Сбой",
completed: "Завершено",
latency: "Макс. ожидание",
status: "Статус",
},
@@ -2282,21 +2283,21 @@ export default {
title: "{count} сбойных задач требуют внимания",
description: "Нажмите красное число в столбце «Окончательные сбои», чтобы увидеть причину, затем повторите после исправления.",
},
failedTasks: {
title: "Окончательные сбои · {queue}",
description: "Задачи, остановленные после исчерпания автоматических повторов; новые сбои показаны первыми.",
tasks: {
title: "Сведения о задачах · {queue}",
description: "Просмотр всех состояний; доступны только действия, которые сервер считает безопасными для текущей задачи.",
guideTitle: "Сначала устраните причину",
guideDescription: "«Запустить один раз» выполняет одну попытку без сброса счётчика повторов. «Очистить запись» удаляет её только из этой очереди, не завершает исходную задачу и сохраняет исторические журналы.",
listTitle: "Сбойные записи",
openAria: "Показать сбоев в очереди {queue}: {count}",
listTitle: "Задачи: {state}",
openAria: "Показать задач {state} в очереди {queue}: {count}",
unavailable: "Подробности сбоев недоступны в этой конфигурации",
empty: "В этой очереди нет остановленных задач",
empty: "В этой очереди нет задач {state}",
loadError: "Не удалось загрузить сведения о сбоях",
loadMore: "Показать ещё",
loadedSummary: "Загружено {count} · прокрутите или нажмите, чтобы загрузить ещё",
loadedAll: "Загружены все {count}",
loadingMore: "Загрузка…",
attempts: "Число повторов {count}",
attempts: "Попытка {current}/{max}",
attemptsShort: "×{count}",
showError: "Показать ошибку",
hideError: "Скрыть ошибку",
@@ -2307,6 +2308,38 @@ export default {
knowledgeLabel: "ID документа",
taskIDLabel: "ID бизнес-задачи",
tenantLabel: "ID тенанта",
sourceLabel: "ID источника",
targetLabel: "ID цели",
sourceKBLabel: "Исходная база знаний",
targetKBLabel: "Целевая база знаний",
dataSourceLabel: "ID источника данных",
syncLogLabel: "ID журнала синхронизации",
knowledgeCountLabel: "Документов",
enqueuedAt: "Поставлено в очередь",
startedAt: "Запущено",
nextProcessAt: "Следующий запуск",
lastFailedAt: "Последний сбой",
completedAt: "Завершено",
deadline: "Срок",
worker: "Исполнитель",
health: "Состояние выполнения",
orphaned: "Связь с исполнителем потеряна; ожидается восстановление",
cancel: "Отменить задачу",
cancelConfirm: "Будет использована бизнес-отмена и остановлены связанные задачи документа. Отменить?",
runNow: "Запустить сейчас",
runNowConfirm: "Задача сразу перейдёт в ожидание без сброса счётчика повторов. Продолжить?",
stateFilter: "Фильтр по состоянию задачи",
states: { active: "Активные", pending: "Ожидают", scheduled: "Запланированы", retry: "Повтор", archived: "Окончательный сбой", completed: "Завершены" },
guides: {
active: "Проверьте исполнителя, время запуска, срок и признак потери связи. Отмена доступна только при безопасном обновлении бизнес-состояния.",
pending: "Задачи ещё не взяты исполнителем. Отмена обновляет бизнес-состояние, а не только удаляет данные Redis.",
scheduled: "Запланированные задачи можно запустить раньше; задачи документов также поддерживают безопасную отмену.",
retry: "Оцените последнюю ошибку, число попыток и время следующего запуска.",
archived: "Сначала устраните причину. Удаление записи не завершает исходную бизнес-задачу.",
completed: "Показаны только недавние завершённые задачи с хранением результата. Действий нет.",
},
actionSuccess: { cancel: "Задача отменена", run_now: "Задача переведена в ожидание", delete: "Запись о сбое удалена" },
actionError: { cancel: "Не удалось отменить задачу", run_now: "Не удалось запустить задачу", delete: "Не удалось удалить запись" },
tenant: "Пространство {id}",
knowledgeBase: "База знаний {id}",
knowledge: "Документ {id}",
@@ -2584,7 +2617,9 @@ export default {
'system.admin_revoked': 'Отозван системный администратор',
'system.user_password_reset': 'Сброшен пароль пользователя',
'system.queue_task_retried': 'Повторно запущена сбойная задача',
'system.queue_task_deleted': 'Удалена запись о сбойной задаче'
'system.queue_task_deleted': 'Удалена запись о сбойной задаче',
'system.queue_task_run_now': 'Задача очереди запущена сейчас',
'system.queue_task_cancelled': 'Задача очереди отменена',
},
outcome: {
success: 'Успешно',
+60 -27
View File
@@ -2533,7 +2533,7 @@ export default {
},
runtime: {
title: "任务队列运行时",
description: "后台任务队列的实时负载,以及各独立 worker 池的每实例并发配置。仅供观察,每 5 秒自动刷新。",
description: "后台任务队列的实时负载,以及各独立 worker 池的每实例并发配置。支持查看任务明细和安全管理,每 5 秒自动刷新。",
refresh: "刷新",
autoRefresh: "自动刷新(每 5 秒)",
loading: "加载中...",
@@ -2572,6 +2572,7 @@ export default {
scheduled: "定时",
retry: "重试",
archived: "最终失败",
completed: "已完成",
latency: "最早等待",
status: "状态",
},
@@ -2588,43 +2589,73 @@ export default {
title: "{count} 个任务待处理",
description: "点击下方表中红色「最终失败」数字查看原因,修复后可手动重试。",
},
failedTasks: {
title: "最终失败 · {queue}",
description: "这里展示超过自动重试上限后停止的任务,按最近失败时间排序。",
guideTitle: "先修复原因,再重新执行",
guideDescription: "“重新执行一次”只会立即再运行一次,不会重置自动重试次数;“清除记录”只从当前队列移除记录,不会完成原任务,历史日志仍会保留。",
listTitle: "失败记录",
openAria: "查看{queue}的 {count} 个最终失败任务",
unavailable: "当前部署不支持查看失败任务明细",
empty: "这个队列当前没有最终失败任务",
loadError: "获取失败任务明细失败",
tasks: {
title: "任务明细 · {queue}",
description: "查看各状态任务及其安全管理动作。可用动作由后端根据任务当前状态和业务语义返回。",
listTitle: "{state}任务",
openAria: "查看{queue}中 {count} 个{state}任务",
unavailable: "当前部署不支持查看任务明细",
empty: "这个队列当前没有{state}任务",
loadError: "获取任务明细失败",
loadMore: "加载更多",
loadedSummary: "已加载 {count} 条,继续下滑或点击加载",
loadedAll: "已全部加载,共 {count} 条",
loadingMore: "正在加载更多…",
attempts: "重试次数 {count}",
attemptsShort: "{count} 次",
showError: "查看错误",
hideError: "收起错误",
lastError: "最后错误",
noError: "未记录错误详情",
attempts: "执行 {current}/{max}",
unknownTarget: "未识别到关联对象",
knowledgeBaseLabel: "知识库 ID",
knowledgeLabel: "文档 ID",
taskIDLabel: "业务任务 ID",
tenantLabel: "租户 ID",
tenant: "空间 {id}",
knowledgeBase: "知识库 {id}",
knowledge: "文档 {id}",
taskID: "业务任务 {id}",
retryOnce: "重新执行一次",
retryConfirm: "请确认已经修复失败原因。该任务将立即重新执行一次。",
retrySuccess: "任务已重新进入队列",
retryError: "重新执行任务失败",
sourceLabel: "来源 ID",
targetLabel: "目标 ID",
sourceKBLabel: "来源知识库",
targetKBLabel: "目标知识库",
dataSourceLabel: "数据源 ID",
syncLogLabel: "同步记录 ID",
knowledgeCountLabel: "文档数量",
enqueuedAt: "入队时间",
startedAt: "开始时间",
nextProcessAt: "下次执行",
lastFailedAt: "最后失败",
completedAt: "完成时间",
deadline: "执行截止",
worker: "执行实例",
health: "运行健康",
orphaned: "执行实例已失联,等待恢复",
cancel: "终止任务",
cancelConfirm: "将通过业务取消流程终止该任务及同一文档的相关队列任务。确认终止?",
runNow: "立即执行",
runNowConfirm: "该任务将立即进入待执行队列,重试次数不会重置。确认继续?",
deleteRecord: "清除记录",
deleteConfirm: "仅从当前队列移除这条失败记录,不会执行或完成原任务,历史日志仍会保留。确认清除?",
deleteSuccess: "失败记录已清除",
deleteError: "清除失败记录失败",
stateFilter: "按任务状态筛选",
states: {
active: "运行中",
pending: "排队中",
scheduled: "定时执行",
retry: "重试中",
archived: "最终失败",
completed: "已完成",
},
guides: {
active: "运行中任务可查看执行实例、开始时间和截止时间。只有具备完整业务取消语义的任务才允许终止。",
pending: "排队任务尚未被 worker 领取。终止操作会同步更新业务状态,而不是只删除 Redis 记录。",
scheduled: "定时任务可提前立即执行;支持业务取消的文档任务也可以安全终止。",
retry: "请结合最后错误、重试次数和下次执行时间判断是否立即执行或终止。",
archived: "请先修复失败原因再立即执行。清除记录不会完成原业务任务。",
completed: "这里只展示设置了结果保留时间的近期完成任务,不提供管理动作。",
},
actionSuccess: {
cancel: "任务已终止",
run_now: "任务已进入待执行队列",
delete: "失败记录已清除",
},
actionError: {
cancel: "终止任务失败",
run_now: "立即执行任务失败",
delete: "清除失败记录失败",
},
taskTypes: {
documentProcess: "文档解析",
manualProcess: "手工重新处理",
@@ -2902,6 +2933,8 @@ export default {
"system.user_password_reset": "重置用户密码",
"system.queue_task_retried": "重新执行失败任务",
"system.queue_task_deleted": "清除失败任务记录",
"system.queue_task_run_now": "立即执行队列任务",
"system.queue_task_cancelled": "终止队列任务",
},
outcome: {
success: "成功",
+450 -204
View File
@@ -171,18 +171,46 @@
</div>
</template>
<template #active="{ row }">
<span class="rq-number" :class="{ 'rq-number--active': row.active > 0 }">{{ row.active }}</span>
<t-button
v-if="row.active > 0"
variant="text"
size="small"
class="rq-task-count rq-task-count--active"
@click="openRuntimeTasks(row, 'active')"
>
{{ row.active }}<t-icon name="chevron-right" />
</t-button>
<span v-else class="rq-number">0</span>
</template>
<template #pending="{ row }">
<div class="rq-backlog">
<span class="rq-number" :class="{ 'rq-number--active': row.pending > 0 }">{{ row.pending }}</span>
<small v-if="row.scheduled > 0">
+{{ row.scheduled }} {{ t('system.globalSettings.runtime.columns.scheduled') }}
</small>
<t-button
v-if="row.pending > 0"
variant="text"
size="small"
class="rq-task-count"
@click="openRuntimeTasks(row, 'pending')"
>{{ row.pending }}<t-icon name="chevron-right" /></t-button>
<span v-else class="rq-number">0</span>
<t-button
v-if="row.scheduled > 0"
variant="text"
size="small"
class="rq-scheduled-count"
@click="openRuntimeTasks(row, 'scheduled')"
>+{{ row.scheduled }} {{ t('system.globalSettings.runtime.columns.scheduled') }}</t-button>
</div>
</template>
<template #retry="{ row }">
<span class="rq-number" :class="{ 'rq-number--warning': row.retry > 0 }">{{ row.retry }}</span>
<t-button
v-if="row.retry > 0"
variant="text"
theme="warning"
size="small"
class="rq-task-count"
@click="openRuntimeTasks(row, 'retry')"
>{{ row.retry }}<t-icon name="chevron-right" /></t-button>
<span v-else class="rq-number">0</span>
</template>
<template #archived="{ row }">
<t-button
@@ -190,14 +218,24 @@
variant="text"
theme="danger"
size="small"
class="rq-failed-count"
:aria-label="t('system.globalSettings.runtime.failedTasks.openAria', { queue: queueLabel(row.name), count: row.archived })"
@click="openFailedTasks(row)"
class="rq-task-count rq-failed-count"
:aria-label="t('system.globalSettings.runtime.tasks.openAria', { state: taskStateLabel('archived'), queue: queueLabel(row.name), count: row.archived })"
@click="openRuntimeTasks(row, 'archived')"
>
{{ row.archived }}<t-icon name="chevron-right" />
</t-button>
<span v-else class="rq-number">0</span>
</template>
<template #completed="{ row }">
<t-button
v-if="row.completed > 0"
variant="text"
size="small"
class="rq-task-count rq-task-count--completed"
@click="openRuntimeTasks(row, 'completed')"
>{{ row.completed }}<t-icon name="chevron-right" /></t-button>
<span v-else class="rq-number">0</span>
</template>
<template #latency_ms="{ row }">
<span class="rq-latency">{{ formatLatency(row.latency_ms) }}</span>
</template>
@@ -254,109 +292,113 @@
</template>
<SettingDrawer
v-model:visible="failedTasksVisible"
v-model:visible="taskDrawerVisible"
class="rq-failed-drawer"
:title="t('system.globalSettings.runtime.failedTasks.title', { queue: failedTaskQueueLabel })"
:description="t('system.globalSettings.runtime.failedTasks.description')"
icon="error-circle"
width="680px"
:title="t('system.globalSettings.runtime.tasks.title', { queue: taskQueueLabel })"
:description="t('system.globalSettings.runtime.tasks.description')"
icon="queue"
width="720px"
:min-width="520"
:max-width="960"
storage-key="setting-drawer:width:runtime-failed-tasks"
:max-width="1040"
storage-key="setting-drawer:width:runtime-tasks"
hide-footer
>
<section class="setting-drawer__section">
<h4 class="setting-drawer__section-title">
{{ t('system.globalSettings.runtime.failedTasks.guideTitle') }}
</h4>
<p class="rq-failed-guide-desc">
{{ t('system.globalSettings.runtime.failedTasks.guideDescription') }}
</p>
<div
class="rq-task-state-filter"
role="tablist"
:aria-label="t('system.globalSettings.runtime.tasks.stateFilter')"
>
<button
v-for="state in taskStates"
:key="state"
type="button"
role="tab"
class="rq-task-state-option"
:class="{ 'is-active': taskState === state }"
:aria-selected="taskState === state"
@click="selectTaskState(state)"
>
<span class="rq-task-state-option__label">{{ taskStateLabel(state) }}</span>
<span
class="rq-task-state-option__count"
:class="{ 'has-value': taskStateCount(taskQueue, state) > 0 }"
>
{{ taskStateCount(taskQueue, state) }}
</span>
</button>
</div>
<p class="rq-failed-guide-desc">{{ taskStateGuide }}</p>
</section>
<section class="setting-drawer__section">
<div class="rq-failed-section-head">
<h4 class="setting-drawer__section-title">
{{ t('system.globalSettings.runtime.failedTasks.listTitle') }}
{{ t('system.globalSettings.runtime.tasks.listTitle', { state: taskStateLabel(taskState) }) }}
</h4>
<t-button
variant="text"
size="small"
:loading="failedTasksLoading && !failedTasksLoadingMore"
@click="reloadFailedTasks"
:loading="tasksLoading && !tasksLoadingMore"
@click="reloadRuntimeTasks"
>
<template #icon><t-icon name="refresh" /></template>
{{ t('system.globalSettings.runtime.refresh') }}
</t-button>
</div>
<div v-if="failedTasksLoading && failedTasks.length === 0" class="rq-failed-loading">
<div v-if="tasksLoading && tasks.length === 0" class="rq-failed-loading">
<t-loading size="small" />
<span>{{ t('system.globalSettings.runtime.loading') }}</span>
</div>
<div v-else-if="failedTasksError" class="rq-failed-error-state">
<span>{{ failedTasksError }}</span>
<t-button size="small" variant="outline" @click="reloadFailedTasks">
<div v-else-if="tasksError" class="rq-failed-error-state">
<span>{{ tasksError }}</span>
<t-button size="small" variant="outline" @click="reloadRuntimeTasks">
{{ t('system.globalSettings.runtime.retry') }}
</t-button>
</div>
<t-empty
v-else-if="failedTasks.length === 0"
:description="t('system.globalSettings.runtime.failedTasks.empty')"
v-else-if="tasks.length === 0"
:description="t('system.globalSettings.runtime.tasks.empty', { state: taskStateLabel(taskState) })"
/>
<div v-else class="rq-failed-list-panel">
<article
v-for="task in failedTasks"
v-for="task in tasks"
:key="task.id"
class="rq-failed-row"
>
<div class="rq-failed-row-content">
<div class="rq-failed-row-summary">
<span class="rq-failed-row-type">{{ failedTaskTypeLabel(task.type) }}</span>
<span class="rq-failed-row-type">{{ runtimeTaskTypeLabel(task.type) }}</span>
<span class="rq-failed-row-sep" aria-hidden="true">·</span>
<span class="rq-failed-row-stat">
{{ t('system.globalSettings.runtime.failedTasks.attempts', { count: task.retried + 1 }) }}
<span class="rq-task-state-pill" :class="`rq-task-state-pill--${task.state}`">
{{ taskStateLabel(task.state) }}
</span>
<span class="rq-failed-row-sep" aria-hidden="true">·</span>
<time :datetime="task.last_failed_at">{{ formatFailedAt(task.last_failed_at) }}</time>
<span class="rq-failed-row-stat">
{{ t('system.globalSettings.runtime.tasks.attempts', { current: task.retried + 1, max: task.max_retry + 1 }) }}
</span>
</div>
<dl v-if="failedTaskRefs(task).length > 0" class="rq-failed-row-refs">
<div v-for="ref in failedTaskRefs(task)" :key="ref.key" class="rq-failed-ref">
<dl v-if="runtimeTaskMeta(task).length > 0" class="rq-failed-row-refs">
<div v-for="ref in runtimeTaskMeta(task)" :key="ref.key" class="rq-failed-ref">
<dt>{{ ref.label }}</dt>
<dd :title="ref.value">{{ ref.value }}</dd>
</div>
</dl>
<p v-else class="rq-failed-row-unknown">
{{ t('system.globalSettings.runtime.failedTasks.unknownTarget') }}
{{ t('system.globalSettings.runtime.tasks.unknownTarget') }}
</p>
<p class="rq-failed-row-error">
{{ task.last_error || t('system.globalSettings.runtime.failedTasks.noError') }}
<p v-if="task.last_error" class="rq-failed-row-error">
{{ task.last_error }}
</p>
</div>
<div class="rq-failed-row-actions">
<t-popconfirm
theme="warning"
:content="t('system.globalSettings.runtime.failedTasks.retryConfirm')"
@confirm="retryFailedTask(task)"
>
<t-button
shape="square"
variant="text"
size="small"
class="rq-failed-icon-btn"
:title="t('system.globalSettings.runtime.failedTasks.retryOnce')"
:aria-label="t('system.globalSettings.runtime.failedTasks.retryOnce')"
:loading="failedTaskActionID === task.id && failedTaskAction === 'retry'"
:disabled="Boolean(failedTaskActionID)"
>
<t-icon name="refresh" />
</t-button>
</t-popconfirm>
<div v-if="task.allowed_actions.length > 0" class="rq-failed-row-actions">
<t-popconfirm
v-if="task.allowed_actions.includes('cancel')"
theme="danger"
:content="t('system.globalSettings.runtime.failedTasks.deleteConfirm')"
@confirm="removeFailedTask(task)"
:content="t('system.globalSettings.runtime.tasks.cancelConfirm')"
@confirm="runTaskAction(task, 'cancel')"
>
<t-button
shape="square"
@@ -364,10 +406,47 @@
size="small"
theme="danger"
class="rq-failed-icon-btn"
:title="t('system.globalSettings.runtime.failedTasks.deleteRecord')"
:aria-label="t('system.globalSettings.runtime.failedTasks.deleteRecord')"
:loading="failedTaskActionID === task.id && failedTaskAction === 'delete'"
:disabled="Boolean(failedTaskActionID)"
:title="t('system.globalSettings.runtime.tasks.cancel')"
:aria-label="t('system.globalSettings.runtime.tasks.cancel')"
:loading="taskActionID === task.id && taskAction === 'cancel'"
:disabled="Boolean(taskActionID)"
><t-icon name="close-circle" /></t-button>
</t-popconfirm>
<t-popconfirm
v-if="task.allowed_actions.includes('run_now')"
theme="warning"
:content="t('system.globalSettings.runtime.tasks.runNowConfirm')"
@confirm="runTaskAction(task, 'run_now')"
>
<t-button
shape="square"
variant="text"
size="small"
class="rq-failed-icon-btn"
:title="t('system.globalSettings.runtime.tasks.runNow')"
:aria-label="t('system.globalSettings.runtime.tasks.runNow')"
:loading="taskActionID === task.id && taskAction === 'run_now'"
:disabled="Boolean(taskActionID)"
>
<t-icon name="refresh" />
</t-button>
</t-popconfirm>
<t-popconfirm
v-if="task.allowed_actions.includes('delete')"
theme="danger"
:content="t('system.globalSettings.runtime.tasks.deleteConfirm')"
@confirm="runTaskAction(task, 'delete')"
>
<t-button
shape="square"
variant="text"
size="small"
theme="danger"
class="rq-failed-icon-btn"
:title="t('system.globalSettings.runtime.tasks.deleteRecord')"
:aria-label="t('system.globalSettings.runtime.tasks.deleteRecord')"
:loading="taskActionID === task.id && taskAction === 'delete'"
:disabled="Boolean(taskActionID)"
>
<t-icon name="delete" />
</t-button>
@@ -375,28 +454,28 @@
</div>
</article>
<div ref="failedTasksSentinelRef" class="rq-failed-load-sentinel" aria-hidden="true" />
<div ref="tasksSentinelRef" class="rq-failed-load-sentinel" aria-hidden="true" />
<div class="rq-failed-list-footer">
<span class="rq-failed-list-status">
<template v-if="failedTasksLoadingMore">
{{ t('system.globalSettings.runtime.failedTasks.loadingMore') }}
<template v-if="tasksLoadingMore">
{{ t('system.globalSettings.runtime.tasks.loadingMore') }}
</template>
<template v-else-if="!failedTasksHasMore">
{{ t('system.globalSettings.runtime.failedTasks.loadedAll', { count: failedTasks.length }) }}
<template v-else-if="!tasksHasMore">
{{ t('system.globalSettings.runtime.tasks.loadedAll', { count: tasks.length }) }}
</template>
<template v-else>
{{ t('system.globalSettings.runtime.failedTasks.loadedSummary', { count: failedTasks.length }) }}
{{ t('system.globalSettings.runtime.tasks.loadedSummary', { count: tasks.length }) }}
</template>
</span>
<t-button
v-if="failedTasksHasMore"
v-if="tasksHasMore"
variant="outline"
block
:loading="failedTasksLoadingMore"
@click="loadMoreFailedTasks"
:loading="tasksLoadingMore"
@click="loadMoreRuntimeTasks"
>
{{ t('system.globalSettings.runtime.failedTasks.loadMore') }}
{{ t('system.globalSettings.runtime.tasks.loadMore') }}
</t-button>
</div>
</div>
@@ -411,13 +490,14 @@ import { useI18n } from 'vue-i18n'
import { MessagePlugin } from 'tdesign-vue-next'
import SettingDrawer from '@/components/settings/SettingDrawer.vue'
import {
deleteRuntimeFailedTask,
getRuntimeFailedTasks,
getRuntimeTasks,
getRuntimeQueues,
retryRuntimeFailedTask,
mutateRuntimeTask,
type ModelRuntimeStat,
type QueueStat,
type RuntimeFailedTask,
type RuntimeTask,
type RuntimeTaskAction,
type RuntimeTaskState,
type RuntimeWorkerPool,
} from '@/api/system'
@@ -435,20 +515,22 @@ const loadedOnce = ref(false)
const error = ref('')
const autoRefresh = ref(true)
const updatedAt = ref('')
const failedTasksVisible = ref(false)
const failedTaskQueue = ref<QueueStat | null>(null)
const failedTasks = ref<RuntimeFailedTask[]>([])
const failedTasksLoading = ref(false)
const failedTasksLoadingMore = ref(false)
const failedTasksError = ref('')
const failedTasksPage = ref(1)
const failedTasksHasMore = ref(false)
const failedTasksSentinelRef = ref<HTMLElement | null>(null)
const failedTaskActionID = ref('')
const failedTaskAction = ref<'retry' | 'delete' | ''>('')
const taskDrawerVisible = ref(false)
const taskQueue = ref<QueueStat | null>(null)
const taskState = ref<RuntimeTaskState>('archived')
const tasks = ref<RuntimeTask[]>([])
const tasksLoading = ref(false)
const tasksLoadingMore = ref(false)
const tasksError = ref('')
const tasksPage = ref(1)
const tasksHasMore = ref(false)
const tasksSentinelRef = ref<HTMLElement | null>(null)
const taskActionID = ref('')
const taskAction = ref<RuntimeTaskAction | ''>('')
const FAILED_TASK_PAGE_SIZE = 20
const failedTaskTypeKeys: Record<string, string> = {
const TASK_PAGE_SIZE = 20
const taskStates: RuntimeTaskState[] = ['active', 'pending', 'scheduled', 'retry', 'archived', 'completed']
const runtimeTaskTypeKeys: Record<string, string> = {
'document:process': 'documentProcess',
'manual:process': 'manualProcess',
'knowledge:post_process': 'postProcess',
@@ -470,7 +552,8 @@ const failedTaskTypeKeys: Record<string, string> = {
}
let pollTimer: ReturnType<typeof setInterval> | null = null
let failedTasksScrollObserver: IntersectionObserver | null = null
let tasksScrollObserver: IntersectionObserver | null = null
let tasksRequestID = 0
const columns = computed(() => [
{ colKey: 'name', title: t('system.globalSettings.runtime.columns.queue'), minWidth: 188 },
@@ -478,6 +561,7 @@ const columns = computed(() => [
{ colKey: 'pending', title: t('system.globalSettings.runtime.columns.pending'), width: 84, align: 'center' as const },
{ colKey: 'retry', title: t('system.globalSettings.runtime.columns.retry'), width: 68, align: 'center' as const },
{ colKey: 'archived', title: t('system.globalSettings.runtime.columns.archived'), width: 96, align: 'center' as const },
{ colKey: 'completed', title: t('system.globalSettings.runtime.columns.completed'), width: 84, align: 'center' as const },
{ colKey: 'latency_ms', title: t('system.globalSettings.runtime.columns.latency'), width: 104, align: 'center' as const },
{ colKey: 'status', title: t('system.globalSettings.runtime.columns.status'), width: 96 },
])
@@ -504,7 +588,8 @@ const totalActive = computed(() => queues.value.reduce((s, q) => s + q.active, 0
const totalPending = computed(() => queues.value.reduce((s, q) => s + q.pending, 0))
const totalRetry = computed(() => queues.value.reduce((s, q) => s + q.retry, 0))
const totalArchived = computed(() => queues.value.reduce((s, q) => s + q.archived, 0))
const failedTaskQueueLabel = computed(() => failedTaskQueue.value ? queueLabel(failedTaskQueue.value.name) : '')
const taskQueueLabel = computed(() => taskQueue.value ? queueLabel(taskQueue.value.name) : '')
const taskStateGuide = computed(() => t(`system.globalSettings.runtime.tasks.guides.${taskState.value}`))
// Friendly per-queue label lives in i18n; falls back to the raw queue
// name so a queue added on the backend still renders before translations
@@ -527,53 +612,30 @@ function queueMeta(row: QueueStat): string {
return scope
}
function failedTaskTypeLabel(type: string): string {
const key = failedTaskTypeKeys[type]
function runtimeTaskTypeLabel(type: string): string {
const key = runtimeTaskTypeKeys[type]
if (!key) return type
const path = `system.globalSettings.runtime.failedTasks.taskTypes.${key}`
const path = `system.globalSettings.runtime.tasks.taskTypes.${key}`
return te(path) ? (t(path) as string) : type
}
interface FailedTaskRef {
interface RuntimeTaskMeta {
key: string
label: string
value: string
}
function failedTaskRefs(task: RuntimeFailedTask): FailedTaskRef[] {
const refs: FailedTaskRef[] = []
if (task.knowledge_base_id) {
refs.push({
key: 'kb',
label: t('system.globalSettings.runtime.failedTasks.knowledgeBaseLabel'),
value: task.knowledge_base_id,
})
}
if (task.knowledge_id) {
refs.push({
key: 'knowledge',
label: t('system.globalSettings.runtime.failedTasks.knowledgeLabel'),
value: task.knowledge_id,
})
}
if (task.task_id) {
refs.push({
key: 'task',
label: t('system.globalSettings.runtime.failedTasks.taskIDLabel'),
value: task.task_id,
})
}
if (refs.length === 0 && task.tenant_id) {
refs.push({
key: 'tenant',
label: t('system.globalSettings.runtime.failedTasks.tenantLabel'),
value: String(task.tenant_id),
})
}
return refs
function taskStateLabel(state: RuntimeTaskState): string {
return t(`system.globalSettings.runtime.tasks.states.${state}`)
}
function formatFailedAt(value: string): string {
function taskStateCount(row: QueueStat | null, state: RuntimeTaskState): number {
if (!row) return 0
return row[state] ?? 0
}
function formatTaskTime(value?: string): string {
if (!value) return '—'
const date = new Date(value)
if (Number.isNaN(date.getTime())) return '—'
return date.toLocaleString(locale.value, {
@@ -587,6 +649,56 @@ function formatFailedAt(value: string): string {
})
}
function runtimeTaskMeta(task: RuntimeTask): RuntimeTaskMeta[] {
const refs: RuntimeTaskMeta[] = []
if (task.knowledge_base_id) {
refs.push({
key: 'kb',
label: t('system.globalSettings.runtime.tasks.knowledgeBaseLabel'),
value: task.knowledge_base_id,
})
}
if (task.knowledge_id) {
refs.push({
key: 'knowledge',
label: t('system.globalSettings.runtime.tasks.knowledgeLabel'),
value: task.knowledge_id,
})
}
if (task.task_id) {
refs.push({
key: 'task',
label: t('system.globalSettings.runtime.tasks.taskIDLabel'),
value: task.task_id,
})
}
if (task.source_id) refs.push({ key: 'source', label: t('system.globalSettings.runtime.tasks.sourceLabel'), value: task.source_id })
if (task.target_id) refs.push({ key: 'target', label: t('system.globalSettings.runtime.tasks.targetLabel'), value: task.target_id })
if (task.source_kb_id) refs.push({ key: 'source-kb', label: t('system.globalSettings.runtime.tasks.sourceKBLabel'), value: task.source_kb_id })
if (task.target_kb_id) refs.push({ key: 'target-kb', label: t('system.globalSettings.runtime.tasks.targetKBLabel'), value: task.target_kb_id })
if (task.data_source_id) refs.push({ key: 'datasource', label: t('system.globalSettings.runtime.tasks.dataSourceLabel'), value: task.data_source_id })
if (task.sync_log_id) refs.push({ key: 'sync-log', label: t('system.globalSettings.runtime.tasks.syncLogLabel'), value: task.sync_log_id })
if (task.knowledge_count) {
refs.push({ key: 'knowledge-count', label: t('system.globalSettings.runtime.tasks.knowledgeCountLabel'), value: String(task.knowledge_count) })
}
if (task.tenant_id) {
refs.push({
key: 'tenant',
label: t('system.globalSettings.runtime.tasks.tenantLabel'),
value: String(task.tenant_id),
})
}
if (task.enqueued_at) refs.push({ key: 'enqueued', label: t('system.globalSettings.runtime.tasks.enqueuedAt'), value: formatTaskTime(task.enqueued_at) })
if (task.started_at) refs.push({ key: 'started', label: t('system.globalSettings.runtime.tasks.startedAt'), value: formatTaskTime(task.started_at) })
if (task.next_process_at) refs.push({ key: 'next', label: t('system.globalSettings.runtime.tasks.nextProcessAt'), value: formatTaskTime(task.next_process_at) })
if (task.last_failed_at) refs.push({ key: 'failed', label: t('system.globalSettings.runtime.tasks.lastFailedAt'), value: formatTaskTime(task.last_failed_at) })
if (task.completed_at) refs.push({ key: 'completed', label: t('system.globalSettings.runtime.tasks.completedAt'), value: formatTaskTime(task.completed_at) })
if (task.deadline) refs.push({ key: 'deadline', label: t('system.globalSettings.runtime.tasks.deadline'), value: formatTaskTime(task.deadline) })
if (task.worker) refs.push({ key: 'worker', label: t('system.globalSettings.runtime.tasks.worker'), value: task.worker })
if (task.is_orphaned) refs.push({ key: 'orphaned', label: t('system.globalSettings.runtime.tasks.health'), value: t('system.globalSettings.runtime.tasks.orphaned') })
return refs
}
function poolLabel(pool: string): string {
const path = `system.globalSettings.runtime.pools.${pool}`
return te(path) ? (t(path) as string) : pool
@@ -634,110 +746,106 @@ function queueState(row: QueueStat): { label: string; tone: string } {
return { label: t('system.globalSettings.runtime.status.idle'), tone: 'idle' }
}
async function fetchFailedTasks(reset: boolean) {
const queue = failedTaskQueue.value?.name
async function fetchRuntimeTasks(reset: boolean) {
const queue = taskQueue.value?.name
if (!queue) return
if (!reset && (failedTasksLoadingMore.value || !failedTasksHasMore.value)) return
if (!reset && (tasksLoadingMore.value || !tasksHasMore.value)) return
const page = reset ? 1 : failedTasksPage.value + 1
const requestedState = taskState.value
const requestID = ++tasksRequestID
const page = reset ? 1 : tasksPage.value + 1
if (reset) {
failedTasksPage.value = 1
failedTasks.value = []
failedTasksLoading.value = true
tasksPage.value = 1
tasks.value = []
tasksLoading.value = true
} else {
failedTasksLoadingMore.value = true
tasksLoadingMore.value = true
}
failedTasksError.value = ''
tasksError.value = ''
try {
const response = await getRuntimeFailedTasks(queue, page, FAILED_TASK_PAGE_SIZE)
const response = await getRuntimeTasks(queue, requestedState, page, TASK_PAGE_SIZE)
if (requestID !== tasksRequestID || taskQueue.value?.name !== queue || taskState.value !== requestedState) return
if (!response.available) {
failedTasksError.value = t('system.globalSettings.runtime.failedTasks.unavailable')
tasksError.value = t('system.globalSettings.runtime.tasks.unavailable')
return
}
failedTasks.value = reset ? response.tasks : [...failedTasks.value, ...response.tasks]
failedTasksPage.value = page
failedTasksHasMore.value = response.has_more
tasks.value = reset ? response.tasks : [...tasks.value, ...response.tasks]
tasksPage.value = page
tasksHasMore.value = response.has_more
} catch (err: any) {
failedTasksError.value = err?.message || t('system.globalSettings.runtime.failedTasks.loadError')
if (requestID !== tasksRequestID) return
tasksError.value = err?.message || t('system.globalSettings.runtime.tasks.loadError')
} finally {
if (requestID !== tasksRequestID) return
if (reset) {
failedTasksLoading.value = false
tasksLoading.value = false
} else {
failedTasksLoadingMore.value = false
tasksLoadingMore.value = false
}
await nextTick()
attachFailedTasksScrollObserver()
attachTasksScrollObserver()
}
}
function detachFailedTasksScrollObserver() {
failedTasksScrollObserver?.disconnect()
failedTasksScrollObserver = null
function detachTasksScrollObserver() {
tasksScrollObserver?.disconnect()
tasksScrollObserver = null
}
function attachFailedTasksScrollObserver() {
detachFailedTasksScrollObserver()
const sentinel = failedTasksSentinelRef.value
if (!sentinel || !failedTasksVisible.value || !failedTasksHasMore.value) return
function attachTasksScrollObserver() {
detachTasksScrollObserver()
const sentinel = tasksSentinelRef.value
if (!sentinel || !taskDrawerVisible.value || !tasksHasMore.value) return
const root = sentinel.closest('.t-drawer__body') as HTMLElement | null
if (!root) return
failedTasksScrollObserver = new IntersectionObserver(
tasksScrollObserver = new IntersectionObserver(
(entries) => {
if (entries.some((entry) => entry.isIntersecting)) {
void loadMoreFailedTasks()
void loadMoreRuntimeTasks()
}
},
{ root, rootMargin: '96px 0px', threshold: 0 },
)
failedTasksScrollObserver.observe(sentinel)
tasksScrollObserver.observe(sentinel)
}
function openFailedTasks(row: QueueStat) {
failedTaskQueue.value = row
failedTasksVisible.value = true
fetchFailedTasks(true)
function openRuntimeTasks(row: QueueStat, state: RuntimeTaskState) {
taskQueue.value = row
taskState.value = state
taskDrawerVisible.value = true
void fetchRuntimeTasks(true)
}
function reloadFailedTasks() {
return fetchFailedTasks(true)
function selectTaskState(state: RuntimeTaskState) {
if (taskState.value === state) return
taskState.value = state
void fetchRuntimeTasks(true)
}
function loadMoreFailedTasks() {
if (failedTasksLoading.value || failedTasksLoadingMore.value || !failedTasksHasMore.value) return
return fetchFailedTasks(false)
function reloadRuntimeTasks() {
return fetchRuntimeTasks(true)
}
async function retryFailedTask(task: RuntimeFailedTask) {
const queue = failedTaskQueue.value?.name
function loadMoreRuntimeTasks() {
if (tasksLoading.value || tasksLoadingMore.value || !tasksHasMore.value) return
return fetchRuntimeTasks(false)
}
async function runTaskAction(task: RuntimeTask, action: RuntimeTaskAction) {
const queue = taskQueue.value?.name
if (!queue) return
failedTaskActionID.value = task.id
failedTaskAction.value = 'retry'
taskActionID.value = task.id
taskAction.value = action
try {
await retryRuntimeFailedTask(queue, task.id)
MessagePlugin.success(t('system.globalSettings.runtime.failedTasks.retrySuccess'))
await Promise.all([reloadFailedTasks(), load(false)])
await mutateRuntimeTask(queue, task.id, action)
MessagePlugin.success(t(`system.globalSettings.runtime.tasks.actionSuccess.${action}`))
await Promise.all([reloadRuntimeTasks(), load(false)])
taskQueue.value = queues.value.find((item) => item.name === queue) ?? taskQueue.value
} catch (err: any) {
MessagePlugin.error(err?.message || t('system.globalSettings.runtime.failedTasks.retryError'))
MessagePlugin.error(err?.message || t(`system.globalSettings.runtime.tasks.actionError.${action}`))
} finally {
failedTaskActionID.value = ''
failedTaskAction.value = ''
}
}
async function removeFailedTask(task: RuntimeFailedTask) {
const queue = failedTaskQueue.value?.name
if (!queue) return
failedTaskActionID.value = task.id
failedTaskAction.value = 'delete'
try {
await deleteRuntimeFailedTask(queue, task.id)
MessagePlugin.success(t('system.globalSettings.runtime.failedTasks.deleteSuccess'))
await Promise.all([reloadFailedTasks(), load(false)])
} catch (err: any) {
MessagePlugin.error(err?.message || t('system.globalSettings.runtime.failedTasks.deleteError'))
} finally {
failedTaskActionID.value = ''
failedTaskAction.value = ''
taskActionID.value = ''
taskAction.value = ''
}
}
@@ -748,6 +856,9 @@ async function load(showSpinner: boolean) {
available.value = resp.available
pools.value = resp.pools || []
queues.value = resp.queues || []
if (taskQueue.value) {
taskQueue.value = queues.value.find((item) => item.name === taskQueue.value?.name) ?? taskQueue.value
}
models.value = resp.models || []
modelLimiterAvailable.value = Boolean(resp.model_limiter_available)
updatedAt.value = new Date((resp.timestamp || Date.now() / 1000) * 1000)
@@ -786,19 +897,19 @@ watch(autoRefresh, (on) => {
else stopPolling()
})
watch(failedTasksVisible, async (open) => {
watch(taskDrawerVisible, async (open) => {
if (!open) {
detachFailedTasksScrollObserver()
detachTasksScrollObserver()
return
}
await nextTick()
attachFailedTasksScrollObserver()
attachTasksScrollObserver()
}, { flush: 'post' })
watch(failedTasksHasMore, async () => {
if (!failedTasksVisible.value) return
watch(tasksHasMore, async () => {
if (!taskDrawerVisible.value) return
await nextTick()
attachFailedTasksScrollObserver()
attachTasksScrollObserver()
})
onMounted(() => {
@@ -808,7 +919,7 @@ onMounted(() => {
onUnmounted(() => {
stopPolling()
detachFailedTasksScrollObserver()
detachTasksScrollObserver()
})
</script>
@@ -1220,6 +1331,81 @@ onUnmounted(() => {
line-height: 1.6;
}
.rq-task-state-filter {
display: flex;
margin-bottom: 12px;
border-bottom: 1px solid var(--td-component-stroke);
}
.rq-task-state-option {
position: relative;
display: inline-flex;
flex: 1;
align-items: center;
justify-content: center;
gap: 4px;
min-width: 0;
padding: 10px 6px;
border: 0;
border-radius: 0;
background: transparent;
color: var(--td-text-color-secondary);
cursor: pointer;
font: inherit;
font-size: 13px;
line-height: 1.2;
white-space: nowrap;
transition: color 0.15s ease;
&:hover:not(.is-active) {
color: var(--td-text-color-primary);
}
&:focus-visible {
outline: 2px solid var(--td-brand-color);
outline-offset: -2px;
}
&.is-active {
color: var(--td-brand-color);
font-weight: 600;
&::after {
content: '';
position: absolute;
left: 8px;
right: 8px;
bottom: -1px;
height: 2px;
border-radius: 2px 2px 0 0;
background: var(--td-brand-color);
}
}
}
.rq-task-state-option__label {
overflow: hidden;
text-overflow: ellipsis;
}
.rq-task-state-option__count {
flex-shrink: 0;
color: var(--td-text-color-placeholder);
font-size: 11px;
font-weight: 500;
line-height: 1;
font-variant-numeric: tabular-nums;
&.has-value {
color: var(--td-text-color-secondary);
}
.rq-task-state-option.is-active & {
color: var(--td-brand-color);
font-weight: 600;
}
}
.rq-failed-section-head {
display: flex;
align-items: center;
@@ -1317,7 +1503,7 @@ onUnmounted(() => {
font-weight: 600;
}
.rq-failed-count {
.rq-task-count {
min-width: 0;
height: 28px;
padding: 0 2px;
@@ -1331,6 +1517,22 @@ onUnmounted(() => {
}
}
.rq-task-count--active {
color: var(--td-brand-color);
}
.rq-task-count--completed {
color: var(--td-success-color);
}
.rq-scheduled-count {
min-width: 0;
height: 18px;
padding: 0;
color: var(--td-text-color-placeholder);
font-size: 11px;
}
.rq-backlog {
display: flex;
align-items: center;
@@ -1412,7 +1614,7 @@ onUnmounted(() => {
text-align: center;
}
&:deep(td.t-align-center .rq-failed-count) {
&:deep(td.t-align-center .rq-task-count) {
margin-inline: auto;
}
@@ -1472,6 +1674,35 @@ onUnmounted(() => {
font-weight: 600;
}
.rq-task-state-pill {
padding: 1px 6px;
border-radius: 999px;
color: var(--td-text-color-secondary);
font-size: 11px;
background: var(--td-bg-color-secondarycontainer);
&--active {
color: var(--td-brand-color);
background: var(--td-brand-color-light);
}
&--retry,
&--scheduled {
color: var(--td-warning-color);
background: var(--td-warning-color-1);
}
&--archived {
color: var(--td-error-color);
background: var(--td-error-color-1);
}
&--completed {
color: var(--td-success-color);
background: var(--td-success-color-1);
}
}
.rq-failed-row-sep {
color: var(--td-text-color-placeholder);
}
@@ -1602,6 +1833,21 @@ onUnmounted(() => {
}
@media (max-width: 620px) {
.rq-task-state-filter {
overflow-x: auto;
scrollbar-width: none;
&::-webkit-scrollbar {
display: none;
}
}
.rq-task-state-option {
flex: 0 0 auto;
min-width: 72px;
padding: 10px 10px;
}
.rq-loading-metrics {
grid-template-columns: 1fr;
}
+8 -1
View File
@@ -1552,9 +1552,11 @@ function auditActionTheme(
case 'system.admin_revoked':
case 'system.setting_changed':
case 'system.queue_task_retried':
case 'system.queue_task_run_now':
return 'warning'
case 'system.user_password_reset':
case 'system.queue_task_deleted':
case 'system.queue_task_cancelled':
return 'danger'
case 'rbac.access_denied':
return 'danger'
@@ -1633,7 +1635,12 @@ function auditTargetKey(row: AuditLog): string {
if (name && mail) return `${name} (${mail})`
return name || mail || (row.target_user_id ? row.target_user_id.slice(0, 8) : '')
}
if (row.action === 'system.queue_task_retried' || row.action === 'system.queue_task_deleted') {
if (
row.action === 'system.queue_task_retried'
|| row.action === 'system.queue_task_run_now'
|| row.action === 'system.queue_task_cancelled'
|| row.action === 'system.queue_task_deleted'
) {
const queue = details && typeof details.queue === 'string' ? details.queue : ''
const taskID = details && typeof details.task_id === 'string' ? details.task_id : row.target_id
return queue && taskID ? `${queue}:${taskID}` : taskID || queue
@@ -3,12 +3,14 @@ package service
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"strings"
"time"
"github.com/Tencent/WeKnora/internal/application/repository"
filesvc "github.com/Tencent/WeKnora/internal/application/service/file"
"github.com/Tencent/WeKnora/internal/application/service/retriever"
"github.com/Tencent/WeKnora/internal/logger"
@@ -138,16 +140,18 @@ func (s *ImageMultimodalService) Handle(ctx context.Context, task *asynq.Task) e
ctx = context.WithValue(ctx, types.LanguageContextKey, payload.Language)
}
// Short-circuit when the parent knowledge has been cancelled by the user
// or marked for deletion. Skip the VLM call entirely so we don't burn
// model quota on already-aborted work.
if k, kerr := s.knowledgeRepo.GetKnowledgeByIDOnly(ctx, payload.KnowledgeID); kerr == nil && k != nil {
switch k.ParseStatus {
case types.ParseStatusCancelled, types.ParseStatusDeleting:
logger.Infof(ctx, "[ImageMultimodal] Knowledge %s aborted (%s), skipping image %s",
payload.KnowledgeID, k.ParseStatus, payload.ImageURL)
return nil
}
// Drop orphaned or user-aborted work before touching VLM. Missing
// knowledge/KB rows are permanent failures — retrying only burns queue
// capacity (asynq default MaxRetry=25 on legacy tasks).
drop, dropErr := s.shouldDropOrphanedMultimodal(ctx, &payload)
if dropErr != nil {
return dropErr
}
if drop {
logger.Infof(ctx,
"[ImageMultimodal] Dropping task chunk=%s knowledge=%s kb=%s image=%s",
payload.ChunkID, payload.KnowledgeID, payload.KnowledgeBaseID, payload.ImageURL)
return nil
}
// Open a per-image subspan under the parent attempt's multimodal
@@ -350,6 +354,41 @@ func (s *ImageMultimodalService) Handle(ctx context.Context, task *asynq.Task) e
return nil
}
// shouldDropOrphanedMultimodal reports whether the task should exit without
// retrying. True for user-cancelled/deleting knowledge, or when the parent
// knowledge / knowledge-base row no longer exists (deleted while queue entries
// survived).
func (s *ImageMultimodalService) shouldDropOrphanedMultimodal(
ctx context.Context, payload *types.ImageMultimodalPayload,
) (bool, error) {
if payload.KnowledgeID != "" && s.knowledgeRepo != nil {
k, err := s.knowledgeRepo.GetKnowledgeByIDOnly(ctx, payload.KnowledgeID)
if errors.Is(err, repository.ErrKnowledgeNotFound) {
return true, nil
}
if err != nil {
return false, err
}
switch k.ParseStatus {
case types.ParseStatusCancelled, types.ParseStatusDeleting:
return true, nil
}
}
if payload.KnowledgeBaseID != "" && s.kbService != nil {
kb, err := s.kbService.GetKnowledgeBaseByIDOnly(ctx, payload.KnowledgeBaseID)
if errors.Is(err, repository.ErrKnowledgeBaseNotFound) {
return true, nil
}
if err != nil {
return false, err
}
if kb == nil {
return true, nil
}
}
return false, nil
}
// isFinalAsynqAttempt reports whether the current task context belongs to the
// last retry attempt before asynq archives the task as a dead-letter. We use
// this to flip multimodal finalize semantics: during normal retries we skip
@@ -0,0 +1,117 @@
package service
import (
"context"
"encoding/json"
"errors"
"testing"
"github.com/Tencent/WeKnora/internal/application/repository"
"github.com/Tencent/WeKnora/internal/types"
"github.com/Tencent/WeKnora/internal/types/interfaces"
"github.com/hibiken/asynq"
)
type orphanKnowledgeRepo struct {
interfaces.KnowledgeRepository
knowledge *types.Knowledge
err error
}
func (r *orphanKnowledgeRepo) GetKnowledgeByIDOnly(_ context.Context, _ string) (*types.Knowledge, error) {
if r.err != nil {
return nil, r.err
}
return r.knowledge, nil
}
type orphanKBService struct {
interfaces.KnowledgeBaseService
kb *types.KnowledgeBase
err error
}
func (s *orphanKBService) GetKnowledgeBaseByIDOnly(_ context.Context, _ string) (*types.KnowledgeBase, error) {
if s.err != nil {
return nil, s.err
}
return s.kb, nil
}
func TestShouldDropOrphanedMultimodal(t *testing.T) {
t.Parallel()
svc := &ImageMultimodalService{}
drop, err := svc.shouldDropOrphanedMultimodal(context.Background(), &types.ImageMultimodalPayload{
KnowledgeID: "missing",
})
if err != nil || drop {
t.Fatalf("nil repo should not drop: drop=%v err=%v", drop, err)
}
svc.knowledgeRepo = &orphanKnowledgeRepo{err: repository.ErrKnowledgeNotFound}
drop, err = svc.shouldDropOrphanedMultimodal(context.Background(), &types.ImageMultimodalPayload{
KnowledgeID: "missing",
})
if err != nil || !drop {
t.Fatalf("missing knowledge should drop: drop=%v err=%v", drop, err)
}
svc.knowledgeRepo = &orphanKnowledgeRepo{knowledge: &types.Knowledge{ParseStatus: types.ParseStatusCancelled}}
drop, err = svc.shouldDropOrphanedMultimodal(context.Background(), &types.ImageMultimodalPayload{
KnowledgeID: "cancelled",
})
if err != nil || !drop {
t.Fatalf("cancelled knowledge should drop: drop=%v err=%v", drop, err)
}
svc.knowledgeRepo = &orphanKnowledgeRepo{knowledge: &types.Knowledge{ParseStatus: types.ParseStatusProcessing}}
svc.kbService = &orphanKBService{err: repository.ErrKnowledgeBaseNotFound}
drop, err = svc.shouldDropOrphanedMultimodal(context.Background(), &types.ImageMultimodalPayload{
KnowledgeID: "live",
KnowledgeBaseID: "missing-kb",
})
if err != nil || !drop {
t.Fatalf("missing kb should drop: drop=%v err=%v", drop, err)
}
}
func TestImageMultimodalHandleDropsMissingKnowledge(t *testing.T) {
t.Parallel()
svc := &ImageMultimodalService{
knowledgeRepo: &orphanKnowledgeRepo{err: repository.ErrKnowledgeNotFound},
kbService: &orphanKBService{kb: &types.KnowledgeBase{ID: "kb-1"}},
}
payload, err := json.Marshal(types.ImageMultimodalPayload{
TenantID: 1,
KnowledgeID: "missing",
KnowledgeBaseID: "kb-1",
ImageURL: "minio://bucket/img.png",
})
if err != nil {
t.Fatal(err)
}
if err := svc.Handle(context.Background(), asynq.NewTask(types.TypeImageMultimodal, payload)); err != nil {
t.Fatalf("orphan task should succeed without retry: %v", err)
}
}
func TestImageMultimodalHandlePropagatesTransientKnowledgeError(t *testing.T) {
t.Parallel()
dbErr := errors.New("db unavailable")
svc := &ImageMultimodalService{
knowledgeRepo: &orphanKnowledgeRepo{err: dbErr},
kbService: &orphanKBService{kb: &types.KnowledgeBase{ID: "kb-1"}},
}
payload, err := json.Marshal(types.ImageMultimodalPayload{
TenantID: 1,
KnowledgeID: "k-1",
KnowledgeBaseID: "kb-1",
})
if err != nil {
t.Fatal(err)
}
if err := svc.Handle(context.Background(), asynq.NewTask(types.TypeImageMultimodal, payload)); !errors.Is(err, dbErr) {
t.Fatalf("transient error should propagate: %v", err)
}
}
+174 -62
View File
@@ -19,6 +19,7 @@ import (
"github.com/Tencent/WeKnora/internal/application/service/file"
"github.com/Tencent/WeKnora/internal/config"
"github.com/Tencent/WeKnora/internal/database"
apperrors "github.com/Tencent/WeKnora/internal/errors"
"github.com/Tencent/WeKnora/internal/infrastructure/docparser"
"github.com/Tencent/WeKnora/internal/logger"
modellimiter "github.com/Tencent/WeKnora/internal/models/limiter"
@@ -32,6 +33,10 @@ import (
"github.com/neo4j/neo4j-go-driver/v6/neo4j"
)
type runtimeKnowledgeCanceller interface {
CancelKnowledgeParse(ctx context.Context, knowledgeID string) (*types.Knowledge, error)
}
// SystemHandler handles system-related requests
type SystemHandler struct {
cfg *config.Config
@@ -49,6 +54,10 @@ type SystemHandler struct {
// Lite mode), so GetRuntimeQueues can distinguish "no queues in this
// deployment" from "queues are empty".
taskInspector interfaces.TaskInspector
// knowledgeSvc supplies the domain-level cancellation path used by the
// runtime task console. It updates business state and tracing before queue
// records are removed, unlike a raw Redis deletion.
knowledgeSvc runtimeKnowledgeCanceller
}
// NewSystemHandler creates a new system handler
@@ -60,6 +69,7 @@ func NewSystemHandler(cfg *config.Config,
systemSettingSvc interfaces.SystemSettingService,
auditSvc interfaces.AuditLogService,
taskInspector interfaces.TaskInspector,
knowledgeSvc interfaces.KnowledgeService,
) *SystemHandler {
return &SystemHandler{
cfg: cfg,
@@ -70,6 +80,7 @@ func NewSystemHandler(cfg *config.Config,
systemSettingSvc: systemSettingSvc,
auditSvc: auditSvc,
taskInspector: taskInspector,
knowledgeSvc: knowledgeSvc,
}
}
@@ -1593,15 +1604,12 @@ func (h *SystemHandler) GetRuntimeQueues(c *gin.Context) {
c.JSON(http.StatusOK, resp)
}
// RuntimeFailedTasksResponse is a page of tasks that exhausted automatic
// retries. Available=false is used by Lite mode, where there is no durable
// queue backend to inspect or mutate.
type RuntimeFailedTasksResponse struct {
Available bool `json:"available"`
Tasks []types.FailedTaskInfo `json:"tasks"`
Page int `json:"page"`
PageSize int `json:"page_size"`
HasMore bool `json:"has_more"`
type RuntimeTasksResponse struct {
Available bool `json:"available"`
Tasks []types.RuntimeTaskInfo `json:"tasks"`
Page int `json:"page"`
PageSize int `json:"page_size"`
HasMore bool `json:"has_more"`
}
func isKnownRuntimeQueue(name string) bool {
@@ -1613,7 +1621,7 @@ func isKnownRuntimeQueue(name string) bool {
return false
}
func runtimeFailedTaskPage(c *gin.Context) (int, int) {
func runtimeTaskPage(c *gin.Context) (int, int) {
page, _ := strconv.Atoi(c.DefaultQuery("page", "1"))
pageSize, _ := strconv.Atoi(c.DefaultQuery("page_size", "20"))
if page < 1 {
@@ -1628,40 +1636,37 @@ func runtimeFailedTaskPage(c *gin.Context) (int, int) {
return page, pageSize
}
// ListRuntimeFailedTasks godoc
// @Summary List tasks that stopped after exhausting automatic retries
// @Tags System Admin
// @Produce json
// @Param queue path string true "Queue name"
// @Success 200 {object} RuntimeFailedTasksResponse
// @Router /system/admin/runtime/queues/{queue}/failed-tasks [get]
func (h *SystemHandler) ListRuntimeFailedTasks(c *gin.Context) {
func (h *SystemHandler) listRuntimeTasks(c *gin.Context, state types.RuntimeTaskState) {
queue := c.Param("queue")
if !isKnownRuntimeQueue(queue) {
c.JSON(http.StatusBadRequest, gin.H{"error": "Unknown task queue"})
return
}
page, pageSize := runtimeFailedTaskPage(c)
inspector, ok := h.taskInspector.(interfaces.FailedTaskInspector)
if !state.Valid() {
c.JSON(http.StatusBadRequest, gin.H{"error": "Unknown task state"})
return
}
page, pageSize := runtimeTaskPage(c)
inspector, ok := h.taskInspector.(interfaces.RuntimeTaskInspector)
if !ok {
c.JSON(http.StatusOK, RuntimeFailedTasksResponse{
c.JSON(http.StatusOK, RuntimeTasksResponse{
Available: false,
Tasks: []types.FailedTaskInfo{},
Tasks: []types.RuntimeTaskInfo{},
Page: page,
PageSize: pageSize,
})
return
}
tasks, supported, err := inspector.ListFailedTasks(c.Request.Context(), queue, page, pageSize)
tasks, supported, err := inspector.ListRuntimeTasks(c.Request.Context(), queue, state, page, pageSize)
if err != nil {
logger.Errorf(c.Request.Context(), "list failed queue tasks queue=%s: %v", queue, err)
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to list stopped tasks"})
logger.Errorf(c.Request.Context(), "list runtime queue tasks queue=%s state=%s: %v", queue, state, err)
c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to list queue tasks"})
return
}
if tasks == nil {
tasks = []types.FailedTaskInfo{}
tasks = []types.RuntimeTaskInfo{}
}
c.JSON(http.StatusOK, RuntimeFailedTasksResponse{
c.JSON(http.StatusOK, RuntimeTasksResponse{
Available: supported,
Tasks: tasks,
Page: page,
@@ -1670,29 +1675,55 @@ func (h *SystemHandler) ListRuntimeFailedTasks(c *gin.Context) {
})
}
// ListRuntimeTasks returns one live task-state page for the unified operator
// drawer. Raw task payloads are never included.
// @Summary List runtime queue tasks by state
// @Tags System Admin
// @Produce json
// @Param queue path string true "Queue name"
// @Param state query string true "Task state" Enums(pending,active,scheduled,retry,archived,completed)
// @Param page query int false "Page" default(1)
// @Param page_size query int false "Page size" default(20)
// @Success 200 {object} RuntimeTasksResponse
// @Router /system/admin/runtime/queues/{queue}/tasks [get]
func (h *SystemHandler) ListRuntimeTasks(c *gin.Context) {
h.listRuntimeTasks(c, types.RuntimeTaskState(c.Query("state")))
}
func (h *SystemHandler) emitQueueTaskAudit(
ctx context.Context,
action types.AuditAction,
queue, taskID string,
extra map[string]string,
) {
if h.auditSvc == nil {
return
}
actorID, _ := types.UserIDFromContext(ctx)
details, _ := json.Marshal(map[string]string{"queue": queue, "task_id": taskID})
detailMap := map[string]string{"queue": queue, "task_id": taskID}
for key, value := range extra {
detailMap[key] = value
}
details, _ := json.Marshal(detailMap)
targetType := "queue_task"
targetID := taskID
if targetID == "" {
targetType = "task_queue"
targetID = queue
}
_ = h.auditSvc.Log(ctx, &types.AuditLog{
TenantID: 0,
ActorUserID: actorID,
ActorRole: "system_admin",
Action: action,
TargetType: "queue_task",
TargetID: taskID,
TargetType: targetType,
TargetID: targetID,
Outcome: types.AuditOutcomeSuccess,
Details: types.JSON(details),
})
}
func runtimeFailedTaskParams(c *gin.Context) (string, string, bool) {
func runtimeTaskParams(c *gin.Context) (string, string, bool) {
queue := c.Param("queue")
taskID := c.Param("task_id")
if !isKnownRuntimeQueue(queue) || taskID == "" {
@@ -1702,56 +1733,137 @@ func runtimeFailedTaskParams(c *gin.Context) (string, string, bool) {
return queue, taskID, true
}
// RetryRuntimeFailedTask schedules one stopped task for one immediate manual
// run. It does not reset the automatic retry counter.
func (h *SystemHandler) RetryRuntimeFailedTask(c *gin.Context) {
queue, taskID, ok := runtimeFailedTaskParams(c)
func (h *SystemHandler) mutateRuntimeTask(c *gin.Context, action types.RuntimeTaskAction) {
queue, taskID, ok := runtimeTaskParams(c)
if !ok {
return
}
inspector, supported := h.taskInspector.(interfaces.FailedTaskInspector)
inspector, supported := h.taskInspector.(interfaces.RuntimeTaskInspector)
if !supported {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Task queue is unavailable"})
return
}
available, err := inspector.RetryFailedTask(c.Request.Context(), queue, taskID)
task, available, err := inspector.GetRuntimeTask(c.Request.Context(), queue, taskID)
if !available {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Task queue is unavailable"})
return
}
if err != nil {
logger.Errorf(c.Request.Context(), "retry failed queue task queue=%s task=%s: %v", queue, taskID, err)
c.JSON(http.StatusConflict, gin.H{"error": "Task is no longer available for retry"})
if err != nil || task == nil {
c.JSON(http.StatusConflict, gin.H{"error": "Task is no longer available"})
return
}
h.emitQueueTaskAudit(c.Request.Context(), types.AuditActionSystemQueueTaskRetried, queue, taskID)
if !task.Allows(action) {
c.JSON(http.StatusConflict, gin.H{"error": "Action is not allowed for the current task state"})
return
}
var auditAction types.AuditAction
auditDetails := map[string]string{
"task_type": task.Type,
"from_state": string(task.State),
"action": string(action),
}
switch action {
case types.RuntimeTaskActionCancel:
if h.knowledgeSvc == nil || task.TenantID == 0 || task.KnowledgeID == "" {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Business cancellation is unavailable"})
return
}
cancelCtx := context.WithValue(c.Request.Context(), types.TenantIDContextKey, task.TenantID)
if _, err = h.knowledgeSvc.CancelKnowledgeParse(cancelCtx, task.KnowledgeID); err != nil {
if isKnowledgeGone(err) {
if h.purgeOrphanRuntimeTask(c.Request.Context(), inspector, queue, taskID, task.KnowledgeID) {
logger.Warnf(c.Request.Context(),
"cancel runtime queue task queue=%s task=%s: knowledge gone, purged queue record",
queue, taskID)
auditDetails["orphan_purge"] = "true"
auditAction = types.AuditActionSystemQueueTaskCancelled
break
}
}
logger.Errorf(c.Request.Context(), "cancel runtime queue task queue=%s task=%s: %v", queue, taskID, err)
c.JSON(http.StatusConflict, gin.H{"error": "Task can no longer be cancelled"})
return
}
auditAction = types.AuditActionSystemQueueTaskCancelled
case types.RuntimeTaskActionRunNow:
available, err = inspector.RunRuntimeTask(c.Request.Context(), queue, taskID)
if !available {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Task queue is unavailable"})
return
}
if err != nil {
logger.Errorf(c.Request.Context(), "run queue task now queue=%s task=%s: %v", queue, taskID, err)
c.JSON(http.StatusConflict, gin.H{"error": "Task is no longer available to run"})
return
}
auditAction = types.AuditActionSystemQueueTaskRunNow
case types.RuntimeTaskActionDelete:
available, err = inspector.DeleteRuntimeTask(c.Request.Context(), queue, taskID)
if !available {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Task queue is unavailable"})
return
}
if err != nil {
logger.Errorf(c.Request.Context(), "delete queue task queue=%s task=%s: %v", queue, taskID, err)
c.JSON(http.StatusConflict, gin.H{"error": "Task is no longer available for deletion"})
return
}
auditAction = types.AuditActionSystemQueueTaskDeleted
default:
c.JSON(http.StatusBadRequest, gin.H{"error": "Unknown task action"})
return
}
h.emitQueueTaskAudit(c.Request.Context(), auditAction, queue, taskID, auditDetails)
c.JSON(http.StatusOK, gin.H{"success": true})
}
// DeleteRuntimeFailedTask removes one stopped task from the queue archive. It
// does not rerun the task or change its related business object.
func (h *SystemHandler) DeleteRuntimeFailedTask(c *gin.Context) {
queue, taskID, ok := runtimeFailedTaskParams(c)
if !ok {
return
func isKnowledgeGone(err error) bool {
if errors.Is(err, repository.ErrKnowledgeNotFound) {
return true
}
inspector, supported := h.taskInspector.(interfaces.FailedTaskInspector)
if !supported {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Task queue is unavailable"})
return
var appErr *apperrors.AppError
return errors.As(err, &appErr) && appErr.Code == apperrors.ErrNotFound
}
// purgeOrphanRuntimeTask removes queue records for a task whose knowledge row
// is already gone. Best-effort: CancelTasksForKnowledge sweeps every queue
// state; ForceDeleteRuntimeTask is the fallback for the specific task ID.
func (h *SystemHandler) purgeOrphanRuntimeTask(
ctx context.Context,
inspector interfaces.RuntimeTaskInspector,
queue, taskID, knowledgeID string,
) bool {
if h.taskInspector != nil && knowledgeID != "" {
deleted, _, err := h.taskInspector.CancelTasksForKnowledge(ctx, knowledgeID)
if err != nil {
logger.Warnf(ctx, "purge orphan tasks for knowledge %s: %v", knowledgeID, err)
} else if deleted > 0 {
return true
}
}
available, err := inspector.DeleteFailedTask(c.Request.Context(), queue, taskID)
if !available {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "Task queue is unavailable"})
return
supported, err := inspector.ForceDeleteRuntimeTask(ctx, queue, taskID)
if !supported || err != nil {
if err != nil {
logger.Warnf(ctx, "force delete orphan task queue=%s task=%s: %v", queue, taskID, err)
}
return false
}
if err != nil {
logger.Errorf(c.Request.Context(), "delete failed queue task queue=%s task=%s: %v", queue, taskID, err)
c.JSON(http.StatusConflict, gin.H{"error": "Task is no longer available for deletion"})
return
}
h.emitQueueTaskAudit(c.Request.Context(), types.AuditActionSystemQueueTaskDeleted, queue, taskID)
c.JSON(http.StatusOK, gin.H{"success": true})
return true
}
// MutateRuntimeTask executes a backend-advertised action after re-reading the
// current task state to close click/race windows.
// @Summary Run a safe runtime task action
// @Tags System Admin
// @Produce json
// @Param queue path string true "Queue name"
// @Param task_id path string true "Task ID"
// @Param action path string true "Action" Enums(cancel,run_now,delete)
// @Success 200 {object} map[string]bool
// @Router /system/admin/runtime/queues/{queue}/tasks/{task_id}/actions/{action} [post]
func (h *SystemHandler) MutateRuntimeTask(c *gin.Context) {
h.mutateRuntimeTask(c, types.RuntimeTaskAction(c.Param("action")))
}
func (h *SystemHandler) ListSystemSettings(c *gin.Context) {
+161 -23
View File
@@ -7,6 +7,7 @@ import (
"net/http/httptest"
"testing"
"github.com/Tencent/WeKnora/internal/application/repository"
"github.com/Tencent/WeKnora/internal/types"
"github.com/gin-gonic/gin"
)
@@ -75,21 +76,79 @@ func (runtimeTestInspector) WorkerServerStats(context.Context) ([]types.WorkerSe
}, true, nil
}
type runtimeFailedTestInspector struct {
type runtimeTaskTestInspector struct {
runtimeTestInspector
tasks []types.FailedTaskInfo
retriedTask string
deletedTask string
mutatedQueue string
tasks []types.RuntimeTaskInfo
retriedTask string
deletedTask string
forceDeleted string
cancelKnowledge string
cancelDeleted int
mutatedQueue string
}
func (r *runtimeFailedTestInspector) ListFailedTasks(
context.Context, string, int, int,
) ([]types.FailedTaskInfo, bool, error) {
func (r *runtimeTaskTestInspector) CancelTasksForKnowledge(
_ context.Context, knowledgeID string,
) (int, int, error) {
r.cancelKnowledge = knowledgeID
if r.cancelDeleted > 0 {
return r.cancelDeleted, 0, nil
}
return 0, 0, nil
}
type runtimeKnowledgeCancelTest struct {
tenantID uint64
knowledgeID string
err error
}
func (r *runtimeKnowledgeCancelTest) CancelKnowledgeParse(
ctx context.Context, knowledgeID string,
) (*types.Knowledge, error) {
r.tenantID, _ = ctx.Value(types.TenantIDContextKey).(uint64)
r.knowledgeID = knowledgeID
if r.err != nil {
return nil, r.err
}
return &types.Knowledge{ID: knowledgeID, TenantID: r.tenantID}, nil
}
func (r *runtimeTaskTestInspector) ListRuntimeTasks(
_ context.Context, _ string, state types.RuntimeTaskState, _, _ int,
) ([]types.RuntimeTaskInfo, bool, error) {
for i := range r.tasks {
if r.tasks[i].State == "" {
r.tasks[i].State = state
}
if r.tasks[i].AllowedActions == nil && state == types.RuntimeTaskArchived {
r.tasks[i].AllowedActions = []types.RuntimeTaskAction{
types.RuntimeTaskActionRunNow,
types.RuntimeTaskActionDelete,
}
}
}
return r.tasks, true, nil
}
func (r *runtimeFailedTestInspector) RetryFailedTask(
func (r *runtimeTaskTestInspector) GetRuntimeTask(
_ context.Context, queue, taskID string,
) (*types.RuntimeTaskInfo, bool, error) {
for i := range r.tasks {
if r.tasks[i].ID == taskID {
return &r.tasks[i], true, nil
}
}
return &types.RuntimeTaskInfo{
ID: taskID, Queue: queue, State: types.RuntimeTaskArchived,
AllowedActions: []types.RuntimeTaskAction{
types.RuntimeTaskActionRunNow,
types.RuntimeTaskActionDelete,
},
}, true, nil
}
func (r *runtimeTaskTestInspector) RunRuntimeTask(
_ context.Context, queue, taskID string,
) (bool, error) {
r.mutatedQueue = queue
@@ -97,7 +156,7 @@ func (r *runtimeFailedTestInspector) RetryFailedTask(
return true, nil
}
func (r *runtimeFailedTestInspector) DeleteFailedTask(
func (r *runtimeTaskTestInspector) DeleteRuntimeTask(
_ context.Context, queue, taskID string,
) (bool, error) {
r.mutatedQueue = queue
@@ -105,6 +164,14 @@ func (r *runtimeFailedTestInspector) DeleteFailedTask(
return true, nil
}
func (r *runtimeTaskTestInspector) ForceDeleteRuntimeTask(
_ context.Context, queue, taskID string,
) (bool, error) {
r.mutatedQueue = queue
r.forceDeleted = taskID
return true, nil
}
func TestGetRuntimeQueuesReportsIsolatedPoolCapacity(t *testing.T) {
gin.SetMode(gin.TestMode)
handler := &SystemHandler{
@@ -186,9 +253,9 @@ func TestGetRuntimeQueuesFallsBackFromInvalidHistoricalConcurrency(t *testing.T)
}
}
func TestListRuntimeFailedTasksReturnsSafeTaskDetails(t *testing.T) {
func TestListRuntimeTasksReturnsSafeTaskDetails(t *testing.T) {
gin.SetMode(gin.TestMode)
inspector := &runtimeFailedTestInspector{tasks: []types.FailedTaskInfo{{
inspector := &runtimeTaskTestInspector{tasks: []types.RuntimeTaskInfo{{
ID: "task-1",
Queue: types.QueueDefault,
Type: types.TypeDocumentProcess,
@@ -202,13 +269,13 @@ func TestListRuntimeFailedTasksReturnsSafeTaskDetails(t *testing.T) {
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Params = gin.Params{{Key: "queue", Value: types.QueueDefault}}
ctx.Request = httptest.NewRequest(http.MethodGet, "/api/v1/system/admin/runtime/queues/default/failed-tasks", nil)
ctx.Request = httptest.NewRequest(http.MethodGet, "/api/v1/system/admin/runtime/queues/default/tasks?state=archived", nil)
handler.ListRuntimeFailedTasks(ctx)
handler.ListRuntimeTasks(ctx)
if recorder.Code != http.StatusOK {
t.Fatalf("status = %d, body=%s", recorder.Code, recorder.Body.String())
}
var response RuntimeFailedTasksResponse
var response RuntimeTasksResponse
if err := json.Unmarshal(recorder.Body.Bytes(), &response); err != nil {
t.Fatalf("decode response: %v", err)
}
@@ -220,9 +287,9 @@ func TestListRuntimeFailedTasksReturnsSafeTaskDetails(t *testing.T) {
}
}
func TestRuntimeFailedTaskMutationsDelegateToInspector(t *testing.T) {
func TestRuntimeTaskMutationsDelegateToInspector(t *testing.T) {
gin.SetMode(gin.TestMode)
inspector := &runtimeFailedTestInspector{}
inspector := &runtimeTaskTestInspector{}
handler := &SystemHandler{taskInspector: inspector}
retryRecorder := httptest.NewRecorder()
@@ -230,9 +297,10 @@ func TestRuntimeFailedTaskMutationsDelegateToInspector(t *testing.T) {
retryCtx.Params = gin.Params{
{Key: "queue", Value: types.QueueDefault},
{Key: "task_id", Value: "task-1"},
{Key: "action", Value: string(types.RuntimeTaskActionRunNow)},
}
retryCtx.Request = httptest.NewRequest(http.MethodPost, "/retry", nil)
handler.RetryRuntimeFailedTask(retryCtx)
handler.MutateRuntimeTask(retryCtx)
if retryRecorder.Code != http.StatusOK || inspector.retriedTask != "task-1" {
t.Fatalf("retry failed: status=%d inspector=%+v", retryRecorder.Code, inspector)
}
@@ -242,24 +310,94 @@ func TestRuntimeFailedTaskMutationsDelegateToInspector(t *testing.T) {
deleteCtx.Params = gin.Params{
{Key: "queue", Value: types.QueueDefault},
{Key: "task_id", Value: "task-2"},
{Key: "action", Value: string(types.RuntimeTaskActionDelete)},
}
deleteCtx.Request = httptest.NewRequest(http.MethodDelete, "/task-2", nil)
handler.DeleteRuntimeFailedTask(deleteCtx)
handler.MutateRuntimeTask(deleteCtx)
if deleteRecorder.Code != http.StatusOK || inspector.deletedTask != "task-2" {
t.Fatalf("delete failed: status=%d inspector=%+v", deleteRecorder.Code, inspector)
}
}
func TestListRuntimeFailedTasksRejectsUnknownQueue(t *testing.T) {
func TestListRuntimeTasksRejectsUnknownQueue(t *testing.T) {
gin.SetMode(gin.TestMode)
handler := &SystemHandler{taskInspector: &runtimeFailedTestInspector{}}
handler := &SystemHandler{taskInspector: &runtimeTaskTestInspector{}}
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Params = gin.Params{{Key: "queue", Value: "unknown"}}
ctx.Request = httptest.NewRequest(http.MethodGet, "/failed-tasks", nil)
ctx.Request = httptest.NewRequest(http.MethodGet, "/tasks?state=archived", nil)
handler.ListRuntimeFailedTasks(ctx)
handler.ListRuntimeTasks(ctx)
if recorder.Code != http.StatusBadRequest {
t.Fatalf("status = %d, want %d", recorder.Code, http.StatusBadRequest)
}
}
func TestListRuntimeTasksRejectsUnknownState(t *testing.T) {
gin.SetMode(gin.TestMode)
handler := &SystemHandler{taskInspector: &runtimeTaskTestInspector{}}
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Params = gin.Params{{Key: "queue", Value: types.QueueDefault}}
ctx.Request = httptest.NewRequest(http.MethodGet, "/tasks?state=unknown", nil)
handler.ListRuntimeTasks(ctx)
if recorder.Code != http.StatusBadRequest {
t.Fatalf("status = %d, want %d", recorder.Code, http.StatusBadRequest)
}
}
func TestRuntimeTaskCancelUsesDomainCancellationWithTaskTenant(t *testing.T) {
gin.SetMode(gin.TestMode)
inspector := &runtimeTaskTestInspector{tasks: []types.RuntimeTaskInfo{{
ID: "task-cancel", Queue: types.QueueDefault, Type: types.TypeDocumentProcess,
State: types.RuntimeTaskActive, TenantID: 42, KnowledgeID: "knowledge-42",
AllowedActions: []types.RuntimeTaskAction{types.RuntimeTaskActionCancel},
}}}
canceller := &runtimeKnowledgeCancelTest{}
handler := &SystemHandler{taskInspector: inspector, knowledgeSvc: canceller}
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Params = gin.Params{
{Key: "queue", Value: types.QueueDefault},
{Key: "task_id", Value: "task-cancel"},
{Key: "action", Value: string(types.RuntimeTaskActionCancel)},
}
ctx.Request = httptest.NewRequest(http.MethodPost, "/cancel", nil)
handler.MutateRuntimeTask(ctx)
if recorder.Code != http.StatusOK {
t.Fatalf("status = %d, body=%s", recorder.Code, recorder.Body.String())
}
if canceller.tenantID != 42 || canceller.knowledgeID != "knowledge-42" {
t.Fatalf("domain cancellation context mismatch: %+v", canceller)
}
}
func TestRuntimeTaskCancelPurgesOrphanWhenKnowledgeGone(t *testing.T) {
gin.SetMode(gin.TestMode)
inspector := &runtimeTaskTestInspector{tasks: []types.RuntimeTaskInfo{{
ID: "task-orphan", Queue: types.QueueMultimodal, Type: types.TypeImageMultimodal,
State: types.RuntimeTaskRetry, TenantID: 42, KnowledgeID: "knowledge-gone",
AllowedActions: []types.RuntimeTaskAction{types.RuntimeTaskActionCancel},
}}}
canceller := &runtimeKnowledgeCancelTest{err: repository.ErrKnowledgeNotFound}
handler := &SystemHandler{taskInspector: inspector, knowledgeSvc: canceller}
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Params = gin.Params{
{Key: "queue", Value: types.QueueMultimodal},
{Key: "task_id", Value: "task-orphan"},
{Key: "action", Value: string(types.RuntimeTaskActionCancel)},
}
ctx.Request = httptest.NewRequest(http.MethodPost, "/cancel", nil)
handler.MutateRuntimeTask(ctx)
if recorder.Code != http.StatusOK {
t.Fatalf("status = %d, body=%s", recorder.Code, recorder.Body.String())
}
if inspector.forceDeleted != "task-orphan" {
t.Fatalf("expected force delete, got deleted=%q force=%q cancel=%q",
inspector.deletedTask, inspector.forceDeleted, inspector.cancelKnowledge)
}
}
+5 -7
View File
@@ -899,14 +899,12 @@ func RegisterSystemAdminRoutes(
adminRoutes.PUT("/settings/:key", handler.UpdateSystemSetting)
adminRoutes.DELETE("/settings/:key", handler.ResetSystemSetting)
// Runtime observability: live asynq queue depths + worker pool
// concurrency for the parse/wiki pools. Read-only snapshot for the
// SystemAdmin runtime dashboard. Returns available=false in Lite
// mode (no Redis) so the UI degrades gracefully.
// Runtime operations: live asynq queue depths, safe task projections,
// and state-checked task actions for the SystemAdmin dashboard. Lite
// mode returns available=false.
adminRoutes.GET("/runtime/queues", handler.GetRuntimeQueues)
adminRoutes.GET("/runtime/queues/:queue/failed-tasks", handler.ListRuntimeFailedTasks)
adminRoutes.POST("/runtime/queues/:queue/failed-tasks/:task_id/retry", handler.RetryRuntimeFailedTask)
adminRoutes.DELETE("/runtime/queues/:queue/failed-tasks/:task_id", handler.DeleteRuntimeFailedTask)
adminRoutes.GET("/runtime/queues/:queue/tasks", handler.ListRuntimeTasks)
adminRoutes.POST("/runtime/queues/:queue/tasks/:task_id/actions/:action", handler.MutateRuntimeTask)
// Bulk action — write the current default-quota setting onto
// every existing tenant. Lives under /tenants instead of
+236 -53
View File
@@ -5,6 +5,7 @@ import (
"encoding/json"
"errors"
"fmt"
"time"
"github.com/Tencent/WeKnora/internal/logger"
"github.com/Tencent/WeKnora/internal/types"
@@ -46,12 +47,21 @@ type knowledgeIDProbe struct {
KnowledgeID string `json:"knowledge_id,omitempty"`
}
type failedTaskPayloadProbe struct {
TenantID uint64 `json:"tenant_id,omitempty"`
KnowledgeBaseID string `json:"knowledge_base_id,omitempty"`
KBID string `json:"kb_id,omitempty"`
KnowledgeID string `json:"knowledge_id,omitempty"`
TaskID string `json:"task_id,omitempty"`
type runtimeTaskPayloadProbe struct {
TenantID uint64 `json:"tenant_id,omitempty"`
KnowledgeBaseID string `json:"knowledge_base_id,omitempty"`
KBID string `json:"kb_id,omitempty"`
KnowledgeID string `json:"knowledge_id,omitempty"`
TaskID string `json:"task_id,omitempty"`
SourceID string `json:"source_id,omitempty"`
TargetID string `json:"target_id,omitempty"`
SourceKBID string `json:"source_kb_id,omitempty"`
TargetKBID string `json:"target_kb_id,omitempty"`
DataSourceID string `json:"data_source_id,omitempty"`
SyncLogID string `json:"sync_log_id,omitempty"`
KnowledgeIDs []string `json:"knowledge_ids,omitempty"`
EnqueuedAt int64 `json:"enqueued_at,omitempty"`
CreatedAt int64 `json:"created_at,omitempty"`
}
// queuesScanned is the fixed set of queue names this codebase enqueues
@@ -192,17 +202,165 @@ func (a *asynqTaskInspector) QueueStats(
return stats, true, nil
}
// ListFailedTasks returns archived tasks newest-first. Only routing metadata is
// projected from the payload so the SystemAdmin dashboard never exposes raw
// document content or connector secrets.
func (a *asynqTaskInspector) ListFailedTasks(
type runtimeWorkerMetadata struct {
started time.Time
worker string
}
func runtimeTaskState(state asynq.TaskState) (types.RuntimeTaskState, error) {
switch state {
case asynq.TaskStatePending:
return types.RuntimeTaskPending, nil
case asynq.TaskStateActive:
return types.RuntimeTaskActive, nil
case asynq.TaskStateScheduled:
return types.RuntimeTaskScheduled, nil
case asynq.TaskStateRetry:
return types.RuntimeTaskRetry, nil
case asynq.TaskStateArchived:
return types.RuntimeTaskArchived, nil
case asynq.TaskStateCompleted:
return types.RuntimeTaskCompleted, nil
default:
return "", fmt.Errorf("unsupported runtime task state %v", state)
}
}
func runtimeTaskTime(value time.Time) *time.Time {
if value.IsZero() {
return nil
}
copy := value
return &copy
}
func runtimePayloadTime(value int64) *time.Time {
if value <= 0 {
return nil
}
// Payload timestamps in this repository are seconds today. Accept the
// common higher-precision Unix forms as well for connector-originated jobs.
var parsed time.Time
if value > 100_000_000_000_000_000 {
parsed = time.Unix(0, value)
} else if value > 100_000_000_000_000 {
parsed = time.UnixMicro(value)
} else if value > 10_000_000_000 {
parsed = time.UnixMilli(value)
} else {
parsed = time.Unix(value, 0)
}
return &parsed
}
func runtimeTaskActions(info types.RuntimeTaskInfo) []types.RuntimeTaskAction {
actions := make([]types.RuntimeTaskAction, 0, 3)
if _, cancellable := taskTypesForKnowledgeCancel[info.Type]; cancellable &&
info.TenantID > 0 && info.KnowledgeID != "" {
switch info.State {
case types.RuntimeTaskPending, types.RuntimeTaskActive,
types.RuntimeTaskScheduled, types.RuntimeTaskRetry:
actions = append(actions, types.RuntimeTaskActionCancel)
}
}
switch info.State {
case types.RuntimeTaskScheduled, types.RuntimeTaskRetry:
actions = append(actions, types.RuntimeTaskActionRunNow)
case types.RuntimeTaskArchived:
actions = append(actions, types.RuntimeTaskActionRunNow, types.RuntimeTaskActionDelete)
}
return actions
}
func projectRuntimeTask(task *asynq.TaskInfo, worker runtimeWorkerMetadata) (types.RuntimeTaskInfo, error) {
state, err := runtimeTaskState(task.State)
if err != nil {
return types.RuntimeTaskInfo{}, err
}
probe := runtimeTaskPayloadProbe{}
_ = json.Unmarshal(task.Payload, &probe)
kbID := probe.KnowledgeBaseID
if kbID == "" {
kbID = probe.KBID
}
enqueuedAt := probe.EnqueuedAt
if enqueuedAt == 0 {
enqueuedAt = probe.CreatedAt
}
info := types.RuntimeTaskInfo{
ID: task.ID,
Queue: task.Queue,
Type: task.Type,
State: state,
LastError: task.LastErr,
LastFailedAt: runtimeTaskTime(task.LastFailedAt),
NextProcessAt: runtimeTaskTime(task.NextProcessAt),
StartedAt: runtimeTaskTime(worker.started),
CompletedAt: runtimeTaskTime(task.CompletedAt),
Deadline: runtimeTaskTime(task.Deadline),
EnqueuedAt: runtimePayloadTime(enqueuedAt),
Retried: task.Retried,
MaxRetry: task.MaxRetry,
IsOrphaned: task.IsOrphaned,
Worker: worker.worker,
TenantID: probe.TenantID,
KnowledgeBaseID: kbID,
KnowledgeID: probe.KnowledgeID,
TaskID: probe.TaskID,
SourceID: probe.SourceID,
TargetID: probe.TargetID,
SourceKBID: probe.SourceKBID,
TargetKBID: probe.TargetKBID,
DataSourceID: probe.DataSourceID,
SyncLogID: probe.SyncLogID,
KnowledgeCount: len(probe.KnowledgeIDs),
}
info.AllowedActions = runtimeTaskActions(info)
return info, nil
}
func (a *asynqTaskInspector) activeWorkerMetadata() map[string]runtimeWorkerMetadata {
result := make(map[string]runtimeWorkerMetadata)
servers, err := a.inspector.Servers()
if err != nil {
return result
}
for _, server := range servers {
if server == nil {
continue
}
workerName := server.Host
if server.PID > 0 {
workerName = fmt.Sprintf("%s:%d", server.Host, server.PID)
}
for _, worker := range server.ActiveWorkers {
if worker == nil {
continue
}
result[worker.Queue+"\x00"+worker.TaskID] = runtimeWorkerMetadata{
started: worker.Started,
worker: workerName,
}
}
}
return result
}
// ListRuntimeTasks returns one task state page. Only allow-listed routing
// metadata is projected from the payload so the dashboard never exposes raw
// document content, signed URLs, or connector secrets.
func (a *asynqTaskInspector) ListRuntimeTasks(
ctx context.Context,
queue string,
state types.RuntimeTaskState,
page, pageSize int,
) ([]types.FailedTaskInfo, bool, error) {
) ([]types.RuntimeTaskInfo, bool, error) {
if a == nil || a.inspector == nil {
return nil, false, nil
}
if !state.Valid() {
return nil, true, fmt.Errorf("unsupported runtime task state %q", state)
}
if page < 1 {
page = 1
}
@@ -212,79 +370,104 @@ func (a *asynqTaskInspector) ListFailedTasks(
if pageSize > 100 {
pageSize = 100
}
tasks, err := a.inspector.ListArchivedTasks(queue, asynq.Page(page), asynq.PageSize(pageSize))
opts := []asynq.ListOption{asynq.Page(page), asynq.PageSize(pageSize)}
var tasks []*asynq.TaskInfo
var err error
switch state {
case types.RuntimeTaskPending:
tasks, err = a.inspector.ListPendingTasks(queue, opts...)
case types.RuntimeTaskActive:
tasks, err = a.inspector.ListActiveTasks(queue, opts...)
case types.RuntimeTaskScheduled:
tasks, err = a.inspector.ListScheduledTasks(queue, opts...)
case types.RuntimeTaskRetry:
tasks, err = a.inspector.ListRetryTasks(queue, opts...)
case types.RuntimeTaskArchived:
tasks, err = a.inspector.ListArchivedTasks(queue, opts...)
case types.RuntimeTaskCompleted:
tasks, err = a.inspector.ListCompletedTasks(queue, opts...)
}
if errors.Is(err, asynq.ErrQueueNotFound) {
return []types.FailedTaskInfo{}, true, nil
return []types.RuntimeTaskInfo{}, true, nil
}
if err != nil {
return nil, true, err
}
result := make([]types.FailedTaskInfo, 0, len(tasks))
workers := map[string]runtimeWorkerMetadata{}
if state == types.RuntimeTaskActive {
workers = a.activeWorkerMetadata()
}
result := make([]types.RuntimeTaskInfo, 0, len(tasks))
for _, task := range tasks {
if task == nil {
continue
}
probe := failedTaskPayloadProbe{}
_ = json.Unmarshal(task.Payload, &probe)
kbID := probe.KnowledgeBaseID
if kbID == "" {
kbID = probe.KBID
info, projectErr := projectRuntimeTask(task, workers[task.Queue+"\x00"+task.ID])
if projectErr != nil {
logger.Warnf(ctx, "[TaskInspector] project runtime task queue=%s id=%s: %v", queue, task.ID, projectErr)
continue
}
result = append(result, types.FailedTaskInfo{
ID: task.ID,
Queue: task.Queue,
Type: task.Type,
LastError: task.LastErr,
LastFailedAt: task.LastFailedAt,
Retried: task.Retried,
MaxRetry: task.MaxRetry,
TenantID: probe.TenantID,
KnowledgeBaseID: kbID,
KnowledgeID: probe.KnowledgeID,
TaskID: probe.TaskID,
})
result = append(result, info)
}
return result, true, nil
}
// RetryFailedTask schedules one archived task for one immediate manual run.
// Asynq deliberately preserves the retry counter, so a repeated failure is
// archived again instead of silently starting a fresh automatic retry cycle.
func (a *asynqTaskInspector) RetryFailedTask(
ctx context.Context,
queue, taskID string,
) (bool, error) {
func (a *asynqTaskInspector) GetRuntimeTask(
ctx context.Context, queue, taskID string,
) (*types.RuntimeTaskInfo, bool, error) {
if a == nil || a.inspector == nil {
return nil, false, nil
}
task, err := a.inspector.GetTaskInfo(queue, taskID)
if err != nil {
return nil, true, err
}
workers := map[string]runtimeWorkerMetadata{}
if task.State == asynq.TaskStateActive {
workers = a.activeWorkerMetadata()
}
info, err := projectRuntimeTask(task, workers[task.Queue+"\x00"+task.ID])
if err != nil {
return nil, true, err
}
return &info, true, nil
}
// RunRuntimeTask moves a scheduled, retry, or archived task to pending. Asynq
// deliberately preserves the retry counter.
func (a *asynqTaskInspector) RunRuntimeTask(ctx context.Context, queue, taskID string) (bool, error) {
if a == nil || a.inspector == nil {
return false, nil
}
if err := a.ensureTaskArchived(queue, taskID); err != nil {
task, _, err := a.GetRuntimeTask(ctx, queue, taskID)
if err != nil {
return true, err
}
if task == nil || !task.Allows(types.RuntimeTaskActionRunNow) {
return true, fmt.Errorf("task %s in queue %s cannot run now", taskID, queue)
}
return true, a.inspector.RunTask(queue, taskID)
}
func (a *asynqTaskInspector) DeleteFailedTask(
ctx context.Context,
queue, taskID string,
) (bool, error) {
func (a *asynqTaskInspector) DeleteRuntimeTask(ctx context.Context, queue, taskID string) (bool, error) {
if a == nil || a.inspector == nil {
return false, nil
}
if err := a.ensureTaskArchived(queue, taskID); err != nil {
task, _, err := a.GetRuntimeTask(ctx, queue, taskID)
if err != nil {
return true, err
}
if task == nil || !task.Allows(types.RuntimeTaskActionDelete) {
return true, fmt.Errorf("task %s in queue %s cannot be deleted", taskID, queue)
}
return true, a.inspector.DeleteTask(queue, taskID)
}
func (a *asynqTaskInspector) ensureTaskArchived(queue, taskID string) error {
task, err := a.inspector.GetTaskInfo(queue, taskID)
if err != nil {
return err
func (a *asynqTaskInspector) ForceDeleteRuntimeTask(ctx context.Context, queue, taskID string) (bool, error) {
if a == nil || a.inspector == nil {
return false, nil
}
if task.State != asynq.TaskStateArchived {
return fmt.Errorf("task %s in queue %s is no longer archived", taskID, queue)
}
return nil
return true, a.inspector.DeleteTask(queue, taskID)
}
func (a *asynqTaskInspector) WorkerServerStats(
@@ -0,0 +1,96 @@
package router
import (
"encoding/json"
"testing"
"time"
"github.com/Tencent/WeKnora/internal/types"
"github.com/hibiken/asynq"
)
func TestProjectRuntimeTaskRedactsPayloadAndBuildsSafeActions(t *testing.T) {
payload, err := json.Marshal(map[string]any{
"tenant_id": 42,
"knowledge_base_id": "kb-1",
"knowledge_id": "knowledge-1",
"file_url": "secret://signed-document-url",
})
if err != nil {
t.Fatal(err)
}
started := time.Unix(1_700_000_000, 0)
info, err := projectRuntimeTask(&asynq.TaskInfo{
ID: "task-1", Queue: types.QueueDefault, Type: types.TypeDocumentProcess,
Payload: payload, State: asynq.TaskStateActive, MaxRetry: 3, Retried: 1,
}, runtimeWorkerMetadata{started: started, worker: "worker-a:123"})
if err != nil {
t.Fatalf("project task: %v", err)
}
if info.State != types.RuntimeTaskActive || info.TenantID != 42 ||
info.KnowledgeBaseID != "kb-1" || info.KnowledgeID != "knowledge-1" {
t.Fatalf("safe routing metadata missing: %+v", info)
}
if info.StartedAt == nil || !info.StartedAt.Equal(started) || info.Worker != "worker-a:123" {
t.Fatalf("worker metadata missing: %+v", info)
}
if len(info.AllowedActions) != 1 || info.AllowedActions[0] != types.RuntimeTaskActionCancel {
t.Fatalf("active document actions = %v", info.AllowedActions)
}
}
func TestProjectRuntimeTaskActionsFollowCurrentState(t *testing.T) {
payload := []byte(`{"tenant_id":7,"knowledge_id":"knowledge-7"}`)
cases := []struct {
state asynq.TaskState
want []types.RuntimeTaskAction
}{
{asynq.TaskStateScheduled, []types.RuntimeTaskAction{types.RuntimeTaskActionCancel, types.RuntimeTaskActionRunNow}},
{asynq.TaskStateRetry, []types.RuntimeTaskAction{types.RuntimeTaskActionCancel, types.RuntimeTaskActionRunNow}},
{asynq.TaskStateArchived, []types.RuntimeTaskAction{types.RuntimeTaskActionRunNow, types.RuntimeTaskActionDelete}},
{asynq.TaskStateCompleted, []types.RuntimeTaskAction{}},
}
for _, tc := range cases {
info, err := projectRuntimeTask(&asynq.TaskInfo{
ID: "task", Queue: types.QueueDefault, Type: types.TypeDocumentProcess,
Payload: payload, State: tc.state,
}, runtimeWorkerMetadata{})
if err != nil {
t.Fatalf("state %v: %v", tc.state, err)
}
if len(info.AllowedActions) != len(tc.want) {
t.Fatalf("state %v actions = %v, want %v", tc.state, info.AllowedActions, tc.want)
}
for i := range tc.want {
if info.AllowedActions[i] != tc.want[i] {
t.Fatalf("state %v actions = %v, want %v", tc.state, info.AllowedActions, tc.want)
}
}
}
}
func TestProjectRuntimeTaskUsesAllowListedBatchMetadata(t *testing.T) {
payload := []byte(`{
"tenant_id":9,
"task_id":"move-1",
"source_kb_id":"source-kb",
"target_kb_id":"target-kb",
"knowledge_ids":["a","b"],
"created_at":1700000000,
"content":"must-not-be-projected"
}`)
info, err := projectRuntimeTask(&asynq.TaskInfo{
ID: "task-move", Queue: types.QueueMaintenance, Type: types.TypeKnowledgeMove,
Payload: payload, State: asynq.TaskStatePending,
}, runtimeWorkerMetadata{})
if err != nil {
t.Fatal(err)
}
if info.TaskID != "move-1" || info.SourceKBID != "source-kb" ||
info.TargetKBID != "target-kb" || info.KnowledgeCount != 2 || info.EnqueuedAt == nil {
t.Fatalf("batch projection mismatch: %+v", info)
}
if len(info.AllowedActions) != 0 {
t.Fatalf("generic maintenance task must not expose raw deletion: %v", info.AllowedActions)
}
}
+4 -2
View File
@@ -125,8 +125,10 @@ const (
// Runtime queue mutations are privileged SystemAdmin actions. Retrying an
// archived task can repeat its original side effects; deleting one removes
// the Redis failure record. Both must leave a platform audit trail.
AuditActionSystemQueueTaskRetried AuditAction = "system.queue_task_retried"
AuditActionSystemQueueTaskDeleted AuditAction = "system.queue_task_deleted"
AuditActionSystemQueueTaskRetried AuditAction = "system.queue_task_retried"
AuditActionSystemQueueTaskDeleted AuditAction = "system.queue_task_deleted"
AuditActionSystemQueueTaskRunNow AuditAction = "system.queue_task_run_now"
AuditActionSystemQueueTaskCancelled AuditAction = "system.queue_task_cancelled"
)
// AuditOutcome distinguishes successful mutations from middleware-level
+8
View File
@@ -40,6 +40,8 @@ func TestAuditAction_DotNamespaceConvention(t *testing.T) {
AuditActionSystemUserPasswordReset,
AuditActionSystemQueueTaskRetried,
AuditActionSystemQueueTaskDeleted,
AuditActionSystemQueueTaskRunNow,
AuditActionSystemQueueTaskCancelled,
}
for _, a := range all {
s := string(a)
@@ -124,6 +126,8 @@ func TestAuditAction_NoCollisionsAcrossNamespaces(t *testing.T) {
register("AuditActionSystemUserPasswordReset", AuditActionSystemUserPasswordReset)
register("AuditActionSystemQueueTaskRetried", AuditActionSystemQueueTaskRetried)
register("AuditActionSystemQueueTaskDeleted", AuditActionSystemQueueTaskDeleted)
register("AuditActionSystemQueueTaskRunNow", AuditActionSystemQueueTaskRunNow)
register("AuditActionSystemQueueTaskCancelled", AuditActionSystemQueueTaskCancelled)
}
// TestAuditAction_SystemNamespacePrefix pins the system.* actions
@@ -140,6 +144,8 @@ func TestAuditAction_SystemNamespacePrefix(t *testing.T) {
AuditActionSystemUserPasswordReset,
AuditActionSystemQueueTaskRetried,
AuditActionSystemQueueTaskDeleted,
AuditActionSystemQueueTaskRunNow,
AuditActionSystemQueueTaskCancelled,
}
for _, a := range cases {
assert.True(t,
@@ -164,6 +170,8 @@ func TestAuditAction_SystemWireValues(t *testing.T) {
{AuditActionSystemUserPasswordReset, "system.user_password_reset"},
{AuditActionSystemQueueTaskRetried, "system.queue_task_retried"},
{AuditActionSystemQueueTaskDeleted, "system.queue_task_deleted"},
{AuditActionSystemQueueTaskRunNow, "system.queue_task_run_now"},
{AuditActionSystemQueueTaskCancelled, "system.queue_task_cancelled"},
}
for _, c := range cases {
assert.Equal(t, c.wire, string(c.constant))
-21
View File
@@ -1,21 +0,0 @@
package types
import "time"
// FailedTaskInfo is the operator-facing projection of an asynq task that
// exhausted its automatic retry budget. Payloads are intentionally not
// exposed by the SystemAdmin API: they may contain document content or
// connector credentials. Only stable routing identifiers are copied out.
type FailedTaskInfo struct {
ID string `json:"id"`
Queue string `json:"queue"`
Type string `json:"type"`
LastError string `json:"last_error"`
LastFailedAt time.Time `json:"last_failed_at"`
Retried int `json:"retried"`
MaxRetry int `json:"max_retry"`
TenantID uint64 `json:"tenant_id,omitempty"`
KnowledgeBaseID string `json:"knowledge_base_id,omitempty"`
KnowledgeID string `json:"knowledge_id,omitempty"`
TaskID string `json:"task_id,omitempty"`
}
+15 -9
View File
@@ -64,16 +64,22 @@ type TaskInspector interface {
WorkerServerStats(ctx context.Context) (stats []types.WorkerServerStat, supported bool, err error)
}
// FailedTaskInspector is the optional operator surface implemented by queue
// backends that retain tasks after their automatic retry budget is exhausted.
// It is separate from TaskInspector so Lite-mode and test implementations that
// only support cancellation/metrics do not need to pretend failed tasks exist.
type FailedTaskInspector interface {
ListFailedTasks(
// RuntimeTaskInspector is the optional operator surface implemented by queue
// backends that retain inspectable task state. It is separate from
// TaskInspector so Lite mode and light-weight tests do not need to implement
// queue mutations.
type RuntimeTaskInspector interface {
ListRuntimeTasks(
ctx context.Context,
queue string,
state types.RuntimeTaskState,
page, pageSize int,
) (tasks []types.FailedTaskInfo, supported bool, err error)
RetryFailedTask(ctx context.Context, queue, taskID string) (supported bool, err error)
DeleteFailedTask(ctx context.Context, queue, taskID string) (supported bool, err error)
) (tasks []types.RuntimeTaskInfo, supported bool, err error)
GetRuntimeTask(ctx context.Context, queue, taskID string) (task *types.RuntimeTaskInfo, supported bool, err error)
RunRuntimeTask(ctx context.Context, queue, taskID string) (supported bool, err error)
DeleteRuntimeTask(ctx context.Context, queue, taskID string) (supported bool, err error)
// ForceDeleteRuntimeTask removes a queue record without checking the
// operator-facing AllowedActions. Used when the business row is already
// gone but a retry/pending task survived (orphan cleanup).
ForceDeleteRuntimeTask(ctx context.Context, queue, taskID string) (supported bool, err error)
}
+79
View File
@@ -0,0 +1,79 @@
package types
import "time"
// RuntimeTaskState is the stable operator-facing task lifecycle. It mirrors
// the durable states exposed by asynq while keeping the HTTP API independent
// from the queue library's Go enum.
type RuntimeTaskState string
const (
RuntimeTaskPending RuntimeTaskState = "pending"
RuntimeTaskActive RuntimeTaskState = "active"
RuntimeTaskScheduled RuntimeTaskState = "scheduled"
RuntimeTaskRetry RuntimeTaskState = "retry"
RuntimeTaskArchived RuntimeTaskState = "archived"
RuntimeTaskCompleted RuntimeTaskState = "completed"
)
func (s RuntimeTaskState) Valid() bool {
switch s {
case RuntimeTaskPending, RuntimeTaskActive, RuntimeTaskScheduled,
RuntimeTaskRetry, RuntimeTaskArchived, RuntimeTaskCompleted:
return true
default:
return false
}
}
// RuntimeTaskAction is returned by the backend for every task. The frontend
// renders only these actions instead of inferring safety from the state.
type RuntimeTaskAction string
const (
RuntimeTaskActionCancel RuntimeTaskAction = "cancel"
RuntimeTaskActionRunNow RuntimeTaskAction = "run_now"
RuntimeTaskActionDelete RuntimeTaskAction = "delete"
)
// RuntimeTaskInfo is the safe SystemAdmin projection of one queue task.
// Raw payloads and results are deliberately excluded because they may contain
// document content, signed object URLs, or connector credentials.
type RuntimeTaskInfo struct {
ID string `json:"id"`
Queue string `json:"queue"`
Type string `json:"type"`
State RuntimeTaskState `json:"state"`
AllowedActions []RuntimeTaskAction `json:"allowed_actions"`
LastError string `json:"last_error,omitempty"`
LastFailedAt *time.Time `json:"last_failed_at,omitempty"`
NextProcessAt *time.Time `json:"next_process_at,omitempty"`
StartedAt *time.Time `json:"started_at,omitempty"`
CompletedAt *time.Time `json:"completed_at,omitempty"`
Deadline *time.Time `json:"deadline,omitempty"`
EnqueuedAt *time.Time `json:"enqueued_at,omitempty"`
Retried int `json:"retried"`
MaxRetry int `json:"max_retry"`
IsOrphaned bool `json:"is_orphaned,omitempty"`
Worker string `json:"worker,omitempty"`
TenantID uint64 `json:"tenant_id,omitempty"`
KnowledgeBaseID string `json:"knowledge_base_id,omitempty"`
KnowledgeID string `json:"knowledge_id,omitempty"`
TaskID string `json:"task_id,omitempty"`
SourceID string `json:"source_id,omitempty"`
TargetID string `json:"target_id,omitempty"`
SourceKBID string `json:"source_kb_id,omitempty"`
TargetKBID string `json:"target_kb_id,omitempty"`
DataSourceID string `json:"data_source_id,omitempty"`
SyncLogID string `json:"sync_log_id,omitempty"`
KnowledgeCount int `json:"knowledge_count,omitempty"`
}
func (t RuntimeTaskInfo) Allows(action RuntimeTaskAction) bool {
for _, allowed := range t.AllowedActions {
if allowed == action {
return true
}
}
return false
}