fix: isolate ERP queues and force deletion by account

This commit is contained in:
inman committed 2026-09-03 09:45:44 +08:00
1 parent 69ea6d2517
commit d09b3032c0
9 files changed
+677 -97

No files matched your search

+37 -14
View File
@@ -21,10 +21,10 @@ import {
import {
TaskError,
TaskService,
canAccessTask,
canViewOperationsDashboard,
type ParseDecisionInput,
type ParseTaskClaim,
type TaskBrowserCommand,
type TaskContext,
type TaskEvent
} from './task-service.js';
@@ -168,8 +168,19 @@ const listTasksQuerySchema = z.object({
z.boolean()
).default(true)
});
const taskEventsQuerySchema = z.object({
since: z.coerce.number().int().min(0).default(0),
executable_by: z.literal('me').default('me')
});
const taskIdListSchema = z.array(z.string().trim().min(1).max(200)).min(1).max(100);
const taskBulkDeleteSchema = z.object({
task_ids: z.array(z.string().trim().min(1).max(200)).min(1).max(100),
task_ids: taskIdListSchema
}).refine(
(body) => new Set(body.task_ids).size === body.task_ids.length,
{ message: '任务编号不能重复。', path: ['task_ids'] }
);
const taskBulkArchiveSchema = z.object({
task_ids: taskIdListSchema,
reason: z.string().trim().max(500).optional()
}).refine(
(body) => new Set(body.task_ids).size === body.task_ids.length,
@@ -1302,14 +1313,14 @@ export async function buildServer({
const body = taskBulkDeleteSchema.parse(request.body);
return {
ok: true,
archived: true,
...(await tasks.archiveTasks(contextFor(session, request), body.task_ids, body.reason))
deleted: true,
...(await tasks.hardDeleteTasks(contextFor(session, request), body.task_ids))
};
});
app.post('/api/tasks/bulk-archive', async (request) => {
const session = await requireMutationSession(request);
const body = taskBulkDeleteSchema.parse(request.body);
const body = taskBulkArchiveSchema.parse(request.body);
return {
ok: true,
archived: true,
@@ -1320,11 +1331,9 @@ export async function buildServer({
app.delete('/api/tasks/:taskId', async (request) => {
const session = await requireMutationSession(request);
const params = request.params as { taskId: string };
const body = taskArchiveSchema.parse(request.body || {});
return {
ok: true,
archived: true,
task: await tasks.archiveTask(contextFor(session, request), params.taskId, body.reason)
...(await tasks.hardDeleteTask(contextFor(session, request), params.taskId))
};
});
@@ -1422,8 +1431,8 @@ export async function buildServer({
app.get('/api/events', async (request, reply) => {
const session = await getSession(request);
const query = (request.query || {}) as Record<string, unknown>;
const querySince = Number(query.since || 0);
const query = taskEventsQuerySchema.parse(request.query || {});
const querySince = query.since;
const reconnectSince = Number(request.headers['last-event-id'] || 0);
const since = Math.max(
Number.isFinite(querySince) ? querySince : 0,
@@ -1439,9 +1448,10 @@ export async function buildServer({
});
const send = (event: TaskEvent) => {
if (event.organization_id !== session.user.organizationId) return;
if (!canAccessTask(contextFor(session, request), {
assignedUserId: event.owner_user_id
})) return;
// This is the executable wake-up feed, not the administrator's read
// model. Every role, including admin, receives only its own assigned
// task events so visibility can never turn into plugin dispatch.
if (event.owner_user_id !== session.user.id) return;
const publicEvent = {
id: event.id,
organization_id: event.organization_id,
@@ -1454,13 +1464,26 @@ export async function buildServer({
};
response.write(`id: ${event.id}\nevent: task\ndata: ${JSON.stringify(publicEvent)}\n\n`);
};
for (const event of await tasks.eventsSince(session.user.organizationId, since, contextFor(session, request))) send(event);
const sendBrowserCommand = (command: TaskBrowserCommand) => {
if (command.organization_id !== session.user.organizationId) return;
if (command.assigned_user_id !== session.user.id) return;
response.write(`event: browser-command\ndata: ${JSON.stringify({
action: command.action,
task_id: command.task_id,
target_user_id: command.assigned_user_id,
created_at: command.created_at
})}\n\n`);
};
for (const event of await tasks.eventsSince(session.user.organizationId, session.user.id, since)) send(event);
const heartbeat = setInterval(() => response.write(': heartbeat\n\n'), 20_000);
const onTask = (event: TaskEvent) => send(event);
const onBrowserCommand = (command: TaskBrowserCommand) => sendBrowserCommand(command);
tasks.events.on('task', onTask);
tasks.events.on('browser-command', onBrowserCommand);
request.raw.on('close', () => {
clearInterval(heartbeat);
tasks.events.off('task', onTask);
tasks.events.off('browser-command', onBrowserCommand);
});
});
+160 -7
View File
@@ -63,6 +63,14 @@ export interface TaskEvent {
task_source?: TaskSource;
}
export interface TaskBrowserCommand {
organization_id: string;
task_id: string;
assigned_user_id: string;
action: 'hard_delete';
created_at: string;
}
export interface PublicTaskEvent {
id: number;
status: string;
@@ -2925,6 +2933,17 @@ export class TaskService {
this.events.emit('task', event);
}
private notifyBrowserCommand(command: TaskBrowserCommand): void {
this.log('info', {
diagnostic_event: 'task.browser_command.emitted',
diagnostic_stage: 'browser_routing',
task_id: command.task_id,
command_action: command.action,
assigned_user_id: command.assigned_user_id
}, 'task browser command emitted to assigned account');
this.events.emit('browser-command', command);
}
private async audit(
client: import('pg').PoolClient,
context: TaskContext,
@@ -6173,12 +6192,23 @@ export class TaskService {
const leaseOwner = `browser:${connectionId}`;
const outcome = await withTransaction(this.config, async (client) => {
const organization = await client.query(
`SELECT id FROM organizations WHERE id = $1 FOR UPDATE`,
`SELECT id FROM organizations WHERE id = $1`,
[context.organizationId]
);
if (!organization.rowCount) throw new TaskError('organization_not_found', '组织不存在。', 404);
if (!context.userId) {
throw new TaskError('task_execution_assignee_mismatch', '任务不属于当前账号的云电脑,禁止领取执行。', 403);
}
// Serialize claims only for this execution account. Other accounts use
// different advisory locks and may claim their own FIFO concurrently.
await client.query(
`SELECT pg_advisory_xact_lock(
hashtextextended($1::text || ':erp-account-queue:' || $2::text, 0)
)`,
[context.organizationId, context.userId]
);
const row = await this.lockTaskForAccess(client, context, taskId);
if (!context.userId || text(row.assigned_user_id) !== context.userId) {
if (text(row.assigned_user_id) !== context.userId) {
throw new TaskError('task_execution_assignee_mismatch', '任务不属于当前账号的云电脑,禁止领取执行。', 403);
}
await this.assertTaskCreatorBusinessAuthorizationInTransaction(client, row, context.requestId);
@@ -6226,12 +6256,13 @@ export class TaskService {
FROM tasks t
JOIN task_attempts a ON a.task_id = t.id AND a.phase = 'erp'
WHERE t.organization_id = $1
AND t.assigned_user_id = $2
AND a.status IN ('accepted', 'running')
AND t.lease_expires_at IS NOT NULL
AND t.lease_expires_at > now()
ORDER BY t.created_at ASC, t.id ASC
FOR UPDATE OF t, a`,
[context.organizationId]
[context.organizationId, context.userId]
);
const activeOther = (activeExecutions.rows as Record<string, unknown>[])
.find((active) => text(active.id) !== text(row.id));
@@ -6239,6 +6270,7 @@ export class TaskService {
`SELECT id, task_id
FROM tasks
WHERE organization_id = $1
AND assigned_user_id = $2
AND status = 'confirmed'
AND handoff_status = 'awaiting_handoff'
AND (
@@ -6265,7 +6297,7 @@ export class TaskService {
)
ORDER BY created_at ASC, id ASC
FOR UPDATE`,
[context.organizationId]
[context.organizationId, context.userId]
);
const queueRows = queue.rows as Record<string, unknown>[];
const queueIndex = queueRows.findIndex((candidate) => text(candidate.id) === text(row.id));
@@ -7067,6 +7099,127 @@ export class TaskService {
return this.getTask(context.organizationId, taskId, context);
}
async hardDeleteTask(
context: TaskContext,
taskId: string
): Promise<{ task_id: string; deleted: boolean }> {
const outcome = await this.hardDeleteTasks(context, [taskId]);
return { task_id: taskId, deleted: outcome.deleted_count === 1 };
}
async hardDeleteTasks(
context: TaskContext,
taskIds: string[]
): Promise<{ task_ids: string[]; deleted_count: number }> {
const normalizedTaskIds = [...new Set(taskIds.map((taskId) => text(taskId).trim()).filter(Boolean))];
if (!normalizedTaskIds.length) throw new TaskError('invalid_task_ids', '请至少选择一个待删除任务。', 400);
const outcome = await withTransaction(this.config, async (client) => {
const lookup = await client.query(
`SELECT * FROM tasks
WHERE organization_id = $1
AND task_id = ANY($2::text[])
AND ($3::boolean = false OR assigned_user_id = $4)
ORDER BY task_id
FOR UPDATE`,
[context.organizationId, normalizedTaskIds, this.isOwnerRestrictedUser(context), context.userId || null]
);
const rows = lookup.rows as Record<string, unknown>[];
const foundTaskIds = new Set(rows.map((row) => text(row.task_id)));
const missingTaskIds = normalizedTaskIds.filter((taskId) => !foundTaskIds.has(taskId));
if (missingTaskIds.length) throw new TaskError('task_not_found', '任务不存在。', 404);
// Hard delete intentionally has no status or handoff-state gate. The row
// locks settle concurrent transitions; existing task foreign keys then
// remove every task-owned record through ON DELETE CASCADE.
const internalTaskIds = rows.map((row) => text(row.id));
const artifactLookup = await client.query(
`SELECT artifact.id, artifact.organization_id, task.task_id AS public_task_id,
artifact.execution_id, artifact.artifact_index, artifact.artifact_type,
artifact.file_name, artifact.content_type, artifact.byte_size, artifact.sha256,
artifact.storage_backend, artifact.storage_key, artifact.created_at
FROM task_artifacts artifact
JOIN tasks task ON task.id = artifact.task_id
WHERE artifact.task_id = ANY($1::uuid[])
ORDER BY artifact.created_at, artifact.id`,
[internalTaskIds]
);
const artifacts: StoredTaskArtifact[] = (artifactLookup.rows as Record<string, unknown>[]).map((row) => ({
id: text(row.id),
organization_id: text(row.organization_id),
task_id: text(row.public_task_id),
execution_id: text(row.execution_id),
artifact_index: Number(row.artifact_index),
type: text(row.artifact_type),
file_name: text(row.file_name),
content_type: text(row.content_type),
byte_size: Number(row.byte_size),
sha256: text(row.sha256),
storage_backend: text(row.storage_backend),
storage_key: text(row.storage_key) || null,
public_url: null,
created_at: new Date(String(row.created_at)).toISOString()
}));
// Retain only a minimal audit marker for the destructive action itself;
// it has no foreign key to the task and contains no task input/output.
for (const row of rows) {
await this.audit(client, context, 'task.hard_deleted', text(row.task_id), {
forced: true,
previous_status: text(row.status),
previous_handoff_status: text(row.handoff_status) || null,
was_archived: Boolean(row.archived_at)
});
}
await client.query(
`DELETE FROM outbox_events
WHERE organization_id = $1
AND aggregate_type = 'task'
AND aggregate_id = ANY($2::text[])`,
[context.organizationId, normalizedTaskIds]
);
const deleted = await client.query(
`DELETE FROM tasks
WHERE id = ANY($1::uuid[])
RETURNING task_id`,
[internalTaskIds]
);
const deletedCount = Number(deleted.rowCount || 0);
if (deletedCount !== rows.length) {
throw new TaskError('task_delete_incomplete', '任务未能完整删除,请重试。', 500);
}
const browserCommands: TaskBrowserCommand[] = rows.flatMap((row) => {
const assignedUserId = text(row.assigned_user_id);
if (!assignedUserId) return [];
return [{
organization_id: context.organizationId,
task_id: text(row.task_id),
assigned_user_id: assignedUserId,
action: 'hard_delete' as const,
created_at: new Date().toISOString()
}];
});
return { artifacts, browserCommands, deletedCount };
});
for (const command of outcome.browserCommands) this.notifyBrowserCommand(command);
if (outcome.artifacts.length && this.artifactStore.cleanup) {
try {
await this.artifactStore.cleanup(outcome.artifacts);
} catch (error) {
this.log('error', {
diagnostic_event: 'task.hard_delete_artifact_cleanup_failed',
diagnostic_stage: 'artifact_cleanup',
request_id: context.requestId,
task_count: normalizedTaskIds.length,
artifact_count: outcome.artifacts.length,
error_type: error instanceof Error ? error.name : 'unknown_error'
}, 'hard-deleted task artifact cleanup failed');
}
}
return { task_ids: normalizedTaskIds, deleted_count: outcome.deletedCount };
}
async cancelTask(context: TaskContext, taskId: string): Promise<PublicTask> {
const outcome = await withTransaction(this.config, async (client) => {
const row = await this.lockTaskForAccess(client, context, taskId);
@@ -7213,15 +7366,15 @@ export class TaskService {
});
}
async eventsSince(organizationId: string, since = 0, access?: TaskAccessScope): Promise<TaskEvent[]> {
async eventsSince(organizationId: string, assignedUserId: string, since = 0): Promise<TaskEvent[]> {
const result = await getPool(this.config).query(
`SELECT e.id, e.organization_id, t.task_id, t.assigned_user_id, t.source,
e.status, e.stage, e.message, e.payload, e.created_at
FROM task_events e JOIN tasks t ON t.id = e.task_id
WHERE e.organization_id = $1 AND e.id > $2
AND ($3::boolean = false OR t.assigned_user_id = $4)
AND t.assigned_user_id = $3
ORDER BY e.id ASC LIMIT 500`,
[organizationId, since, isTaskOwnerRestricted(access?.role), access?.userId || null]
[organizationId, since, assignedUserId]
);
return result.rows.map((row) => ({
id: Number(row.id),