feat: add account roles audit and task authorization
This commit is contained in:
1 parent
337aaf7c88
commit
191c1a1aad
15 files changed
+4961
-457
No files matched your search
+338
-28
@@ -21,12 +21,15 @@ import {
|
||||
import {
|
||||
TaskError,
|
||||
TaskService,
|
||||
canAccessTask,
|
||||
canViewOperationsDashboard,
|
||||
type ParseDecisionInput,
|
||||
type ParseTaskClaim,
|
||||
type TaskContext,
|
||||
type TaskEvent
|
||||
} from './task-service.js';
|
||||
import { ParserOrchestrator, type AiParser as ExternalParser } from './parser-orchestrator.js';
|
||||
import { BUSINESS_ROUTES, businessRouteById } from './business-routes.js';
|
||||
import {
|
||||
decodeInlineInputAttachment,
|
||||
InputAttachmentError,
|
||||
@@ -51,6 +54,42 @@ const loginSchema = z.object({
|
||||
password: z.string().min(1).max(512)
|
||||
});
|
||||
|
||||
const changePasswordSchema = z.object({
|
||||
current_password: z.string().min(1).max(512),
|
||||
new_password: z.string().min(12).max(512)
|
||||
});
|
||||
|
||||
const accountCreateSchema = z.object({
|
||||
username: z.string().min(1).max(160),
|
||||
password: z.string().min(12).max(512),
|
||||
role: z.enum(['admin', 'team_lead', 'user']).default('user'),
|
||||
must_change_password: z.boolean().default(true),
|
||||
business_route_ids: z.array(
|
||||
z.string().trim().refine((routeId) => Boolean(businessRouteById(routeId)), '业务类型不存在。')
|
||||
).max(BUSINESS_ROUTES.length).default([])
|
||||
.refine((routeIds) => new Set(routeIds).size === routeIds.length, '业务类型不能重复。')
|
||||
});
|
||||
|
||||
const accountUpdateSchema = z.object({
|
||||
role: z.enum(['admin', 'team_lead', 'user']).optional(),
|
||||
is_active: z.boolean().optional()
|
||||
}).refine((body) => body.role !== undefined || body.is_active !== undefined, {
|
||||
message: '至少提供一个账号更新字段。'
|
||||
});
|
||||
|
||||
const accountPasswordResetSchema = z.object({
|
||||
password: z.string().min(12).max(512),
|
||||
must_change_password: z.boolean().default(true)
|
||||
});
|
||||
|
||||
const accountBusinessAuthorizationsSchema = z.object({
|
||||
business_route_ids: z.array(
|
||||
z.string().trim().refine((routeId) => Boolean(businessRouteById(routeId)), '业务类型不存在。')
|
||||
).max(BUSINESS_ROUTES.length)
|
||||
.refine((routeIds) => new Set(routeIds).size === routeIds.length, '业务类型不能重复。'),
|
||||
expected_revision: z.number().int().min(0)
|
||||
});
|
||||
|
||||
const encodedInputAttachmentSchema = z.object({
|
||||
name: z.string().min(1).max(200),
|
||||
content_type: z.string().max(200).optional(),
|
||||
@@ -118,18 +157,41 @@ const listTasksQuerySchema = z.object({
|
||||
search: z.string().max(200).optional(),
|
||||
limit: z.coerce.number().int().min(1).max(200).default(200),
|
||||
offset: z.coerce.number().int().min(0).max(1_000_000).default(0),
|
||||
archive: z.enum(['active', 'archived', 'all']).default('active'),
|
||||
include_total: z.preprocess(
|
||||
(value) => value === undefined ? true : String(value).toLowerCase() !== 'false',
|
||||
z.boolean()
|
||||
).default(true)
|
||||
});
|
||||
const taskBulkDeleteSchema = z.object({
|
||||
task_ids: z.array(z.string().trim().min(1).max(200)).min(1).max(100)
|
||||
task_ids: z.array(z.string().trim().min(1).max(200)).min(1).max(100),
|
||||
reason: z.string().trim().max(500).optional()
|
||||
}).refine(
|
||||
(body) => new Set(body.task_ids).size === body.task_ids.length,
|
||||
{ message: '任务编号不能重复。', path: ['task_ids'] }
|
||||
);
|
||||
|
||||
const taskArchiveSchema = z.object({ reason: z.string().trim().max(500).optional() });
|
||||
|
||||
const auditQuerySchema = z.object({
|
||||
event_type: z.string().trim().max(200).optional(),
|
||||
actor_user_id: z.string().uuid().optional(),
|
||||
entity_type: z.string().trim().max(120).optional(),
|
||||
limit: z.coerce.number().int().min(1).max(200).default(100),
|
||||
offset: z.coerce.number().int().min(0).max(1_000_000).default(0)
|
||||
});
|
||||
|
||||
const operationsDashboardQuerySchema = z.object({
|
||||
from: z.string().datetime({ offset: true }).optional(),
|
||||
to: z.string().datetime({ offset: true }).optional(),
|
||||
actor_user_id: z.string().uuid().optional(),
|
||||
business_route_id: z.string().trim().max(120).optional(),
|
||||
status: z.enum(['all', 'active', 'completed', 'attention', 'failed', 'cancelled', 'archived']).default('all'),
|
||||
search: z.string().trim().max(200).optional(),
|
||||
limit: z.coerce.number().int().min(1).max(100).default(50),
|
||||
offset: z.coerce.number().int().min(0).max(1_000_000).default(0)
|
||||
});
|
||||
|
||||
const LOG_REDACTION_PATHS = [
|
||||
'req.headers.authorization',
|
||||
'req.headers.cookie',
|
||||
@@ -302,7 +364,7 @@ function publicUser(session: ActiveSession) {
|
||||
id: session.user.id,
|
||||
username: session.user.username,
|
||||
role: session.user.role,
|
||||
organization_id: session.user.organizationId
|
||||
must_change_password: session.user.mustChangePassword
|
||||
};
|
||||
}
|
||||
|
||||
@@ -402,6 +464,21 @@ export async function buildServer({
|
||||
return reply.sendFile('index.html');
|
||||
});
|
||||
|
||||
app.get('/accounts', async (_request, reply) => {
|
||||
reply.header('Cache-Control', 'no-store');
|
||||
return reply.sendFile('index.html');
|
||||
});
|
||||
|
||||
app.get('/audit', async (_request, reply) => {
|
||||
reply.header('Cache-Control', 'no-store');
|
||||
return reply.sendFile('index.html');
|
||||
});
|
||||
|
||||
app.get('/operations-dashboard', async (_request, reply) => {
|
||||
reply.header('Cache-Control', 'no-store');
|
||||
return reply.sendFile('index.html');
|
||||
});
|
||||
|
||||
app.addHook('onRequest', async (request, reply) => {
|
||||
if (request.url.startsWith('/api/') || request.url.startsWith('/health/')) return;
|
||||
const redirectUrl = canonicalStaticRedirect(config, request);
|
||||
@@ -439,14 +516,35 @@ export async function buildServer({
|
||||
return session;
|
||||
};
|
||||
|
||||
const getReadySession = async (request: FastifyRequest): Promise<ActiveSession> => {
|
||||
const session = await getSession(request);
|
||||
if (session.user.mustChangePassword) {
|
||||
throw new AuthError('password_change_required', '首次登录或密码重置后必须先修改密码。', 403);
|
||||
}
|
||||
return session;
|
||||
};
|
||||
|
||||
const requireAdmin = (session: ActiveSession): ActiveSession => {
|
||||
if (session.user.role !== 'admin') throw new AuthError('admin_required', '需要管理员权限。', 403);
|
||||
return session;
|
||||
};
|
||||
|
||||
const requireLeadership = (session: ActiveSession): ActiveSession => {
|
||||
if (!canViewOperationsDashboard(session.user.role)) {
|
||||
throw new AuthError('leadership_required', '需要组长或管理员权限。', 403);
|
||||
}
|
||||
return session;
|
||||
};
|
||||
|
||||
const contextFor = (session: ActiveSession, request: FastifyRequest): TaskContext => ({
|
||||
organizationId: session.user.organizationId,
|
||||
userId: session.user.id,
|
||||
requestId: requestId(request),
|
||||
role: session.user.role,
|
||||
source: 'manual'
|
||||
});
|
||||
|
||||
const requireMutationSession = async (request: FastifyRequest): Promise<ActiveSession> => {
|
||||
const requireAuthenticatedMutationSession = async (request: FastifyRequest): Promise<ActiveSession> => {
|
||||
requireSameOrigin(config, request);
|
||||
const session = await getSession(request);
|
||||
const csrf = String(request.headers['x-csrf-token'] || '');
|
||||
@@ -456,6 +554,26 @@ export async function buildServer({
|
||||
return session;
|
||||
};
|
||||
|
||||
const requireMutationSession = async (request: FastifyRequest): Promise<ActiveSession> => {
|
||||
const session = await requireAuthenticatedMutationSession(request);
|
||||
if (session.user.mustChangePassword) {
|
||||
throw new AuthError('password_change_required', '首次登录或密码重置后必须先修改密码。', 403);
|
||||
}
|
||||
return session;
|
||||
};
|
||||
|
||||
const requireAdminSession = async (request: FastifyRequest): Promise<ActiveSession> => (
|
||||
requireAdmin(await getReadySession(request))
|
||||
);
|
||||
|
||||
const requireAdminMutationSession = async (request: FastifyRequest): Promise<ActiveSession> => (
|
||||
requireAdmin(await requireMutationSession(request))
|
||||
);
|
||||
|
||||
const requireLeadershipSession = async (request: FastifyRequest): Promise<ActiveSession> => (
|
||||
requireLeadership(await getReadySession(request))
|
||||
);
|
||||
|
||||
async function persistParseOutcome(
|
||||
claim: ParseTaskClaim,
|
||||
result: unknown,
|
||||
@@ -834,13 +952,97 @@ export async function buildServer({
|
||||
return { ok: true, csrf_token: await auth.rotateCsrf(session.id) };
|
||||
});
|
||||
|
||||
app.put('/api/auth/password', async (request) => {
|
||||
const session = await requireAuthenticatedMutationSession(request);
|
||||
const body = changePasswordSchema.parse(request.body);
|
||||
await auth.changeOwnPassword(session, body.current_password, body.new_password, requestId(request));
|
||||
return { ok: true, password_changed: true };
|
||||
});
|
||||
|
||||
app.get('/api/accounts', async (request) => {
|
||||
const session = await requireAdminSession(request);
|
||||
return {
|
||||
ok: true,
|
||||
accounts: await auth.listAccounts(session.user),
|
||||
task_types: BUSINESS_ROUTES.map((route, displayOrder) => ({
|
||||
route_id: route.routeId,
|
||||
directive: route.directive,
|
||||
action: route.action,
|
||||
display_order: displayOrder + 1
|
||||
}))
|
||||
};
|
||||
});
|
||||
|
||||
app.post('/api/accounts', { config: { rateLimit: { max: 20, timeWindow: '1 minute' } } }, async (request) => {
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const body = accountCreateSchema.parse(request.body);
|
||||
const account = await auth.createAccount(session.user, {
|
||||
username: body.username,
|
||||
password: body.password,
|
||||
role: body.role,
|
||||
mustChangePassword: body.must_change_password,
|
||||
businessRouteIds: body.business_route_ids
|
||||
}, requestId(request));
|
||||
return { ok: true, account };
|
||||
});
|
||||
|
||||
app.patch('/api/accounts/:userId', async (request) => {
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const params = request.params as { userId: string };
|
||||
const userId = z.string().uuid().parse(params.userId);
|
||||
const body = accountUpdateSchema.parse(request.body);
|
||||
const account = await auth.updateAccount(session.user, userId, {
|
||||
role: body.role,
|
||||
isActive: body.is_active
|
||||
}, requestId(request));
|
||||
return { ok: true, account };
|
||||
});
|
||||
|
||||
app.put('/api/accounts/:userId/business-authorizations', async (request) => {
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const params = request.params as { userId: string };
|
||||
const userId = z.string().uuid().parse(params.userId);
|
||||
const body = accountBusinessAuthorizationsSchema.parse(request.body);
|
||||
const account = await auth.setBusinessRouteAuthorizations(
|
||||
session.user,
|
||||
userId,
|
||||
body.business_route_ids,
|
||||
body.expected_revision,
|
||||
requestId(request)
|
||||
);
|
||||
return { ok: true, account };
|
||||
});
|
||||
|
||||
app.post('/api/accounts/:userId/reset-password', { config: { rateLimit: { max: 20, timeWindow: '1 minute' } } }, async (request) => {
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const params = request.params as { userId: string };
|
||||
const userId = z.string().uuid().parse(params.userId);
|
||||
const body = accountPasswordResetSchema.parse(request.body);
|
||||
await auth.resetAccountPassword(
|
||||
session.user,
|
||||
userId,
|
||||
body.password,
|
||||
body.must_change_password,
|
||||
requestId(request)
|
||||
);
|
||||
return { ok: true, password_reset: true, sessions_revoked: true };
|
||||
});
|
||||
|
||||
app.post('/api/accounts/:userId/revoke-sessions', async (request) => {
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const params = request.params as { userId: string };
|
||||
const userId = z.string().uuid().parse(params.userId);
|
||||
const revoked = await auth.revokeAccountSessions(session.user, userId, requestId(request));
|
||||
return { ok: true, sessions_revoked: revoked };
|
||||
});
|
||||
|
||||
app.get('/api/settings/automation', async (request) => {
|
||||
const session = await getSession(request);
|
||||
const session = await requireAdminSession(request);
|
||||
return { ok: true, settings: await tasks.getAutomationSettings(session.user.organizationId) };
|
||||
});
|
||||
|
||||
app.put('/api/settings/automation', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const body = automationSettingsSchema.parse(request.body);
|
||||
return {
|
||||
ok: true,
|
||||
@@ -849,12 +1051,12 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.get('/api/settings/parser-routing', async (request) => {
|
||||
const session = await getSession(request);
|
||||
const session = await requireAdminSession(request);
|
||||
return { ok: true, routes: await tasks.getParserRoutingSettings(session.user.organizationId) };
|
||||
});
|
||||
|
||||
app.put('/api/settings/parser-routing/:routeId', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const routeId = String((request.params as { routeId?: string }).routeId || '');
|
||||
const body = parserRoutingUpdateSchema.parse(request.body);
|
||||
return {
|
||||
@@ -869,7 +1071,7 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.post('/api/settings/parser-routing/emergency-ai', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const body = parserEmergencyAiSchema.parse(request.body);
|
||||
return {
|
||||
ok: true,
|
||||
@@ -878,7 +1080,7 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.post('/api/tasks/:taskId/reparse', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const taskId = String((request.params as { taskId?: string }).taskId || '');
|
||||
const body = parserReparseSchema.parse(request.body);
|
||||
const task = await tasks.reparseTaskWithAi(contextFor(session, request), taskId, body.reason);
|
||||
@@ -890,7 +1092,7 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.put('/api/parser-decisions/:decisionId/review', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const decisionId = String((request.params as { decisionId?: string }).decisionId || '');
|
||||
const body = parserDecisionReviewSchema.parse(request.body);
|
||||
return {
|
||||
@@ -902,7 +1104,7 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.get('/api/channels', async (request) => {
|
||||
const session = await getSession(request);
|
||||
const session = await requireAdminSession(request);
|
||||
const channels = await channelService.list(session.user.organizationId);
|
||||
return {
|
||||
ok: true,
|
||||
@@ -911,7 +1113,7 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.post('/api/channels', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const body = channelCreateSchema.parse(request.body);
|
||||
const channel = await channelService.create(contextFor(session, request), {
|
||||
displayName: body.display_name,
|
||||
@@ -925,7 +1127,7 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.patch('/api/channels/:channelId', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const params = request.params as { channelId: string };
|
||||
const body = channelUpdateSchema.parse(request.body);
|
||||
const channel = await channelService.update(contextFor(session, request), params.channelId, {
|
||||
@@ -939,7 +1141,7 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.post('/api/channels/:channelId/rotate-key', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const params = request.params as { channelId: string };
|
||||
const body = channelRotateKeySchema.parse(request.body);
|
||||
const channel = await channelService.rotateKey(
|
||||
@@ -952,7 +1154,7 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.delete('/api/channels/:channelId', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const params = request.params as { channelId: string };
|
||||
const result = await channelService.delete(contextFor(session, request), params.channelId);
|
||||
await agentBus?.reload();
|
||||
@@ -961,23 +1163,25 @@ export async function buildServer({
|
||||
|
||||
app.post('/api/auth/logout', async (request, reply) => {
|
||||
setAuthNoStore(reply);
|
||||
const session = await requireMutationSession(request);
|
||||
const session = await requireAuthenticatedMutationSession(request);
|
||||
await auth.revokeSession(request.cookies[config.SESSION_COOKIE_NAME]);
|
||||
reply.clearCookie(config.SESSION_COOKIE_NAME, { path: '/' });
|
||||
await auth.recordAudit(session.user.organizationId, session.user.id, 'logout', requestId(request));
|
||||
app.log.info({ user_id: session.user.id, request_id: requestId(request) }, 'administrator logged out');
|
||||
app.log.info({ user_id: session.user.id, request_id: requestId(request) }, 'user logged out');
|
||||
return { ok: true };
|
||||
});
|
||||
|
||||
app.get('/api/tasks', async (request) => {
|
||||
const session = await getSession(request);
|
||||
const session = await getReadySession(request);
|
||||
const query = listTasksQuerySchema.parse(request.query || {});
|
||||
const page = await tasks.listTasksPage(session.user.organizationId, {
|
||||
status: query.status || undefined,
|
||||
search: query.search || undefined,
|
||||
limit: query.limit,
|
||||
offset: query.offset,
|
||||
includeTotal: query.include_total
|
||||
includeTotal: query.include_total,
|
||||
archive: query.archive,
|
||||
access: contextFor(session, request)
|
||||
});
|
||||
return {
|
||||
ok: true,
|
||||
@@ -1020,17 +1224,28 @@ export async function buildServer({
|
||||
});
|
||||
|
||||
app.get('/api/tasks/:taskId', async (request) => {
|
||||
const session = await getSession(request);
|
||||
const session = await getReadySession(request);
|
||||
const params = request.params as { taskId: string };
|
||||
return { ok: true, task: await tasks.getTask(session.user.organizationId, params.taskId) };
|
||||
return { ok: true, task: await tasks.getTask(session.user.organizationId, params.taskId, contextFor(session, request)) };
|
||||
});
|
||||
|
||||
app.get('/api/tasks/:taskId/input-history', async (request) => {
|
||||
const session = await getReadySession(request);
|
||||
const params = request.params as { taskId: string };
|
||||
return { ok: true, ...(await tasks.getTaskInputHistory(contextFor(session, request), params.taskId)) };
|
||||
});
|
||||
|
||||
app.get('/api/tasks/:taskId/artifacts/:artifactId', async (request, reply) => {
|
||||
const session = await getSession(request);
|
||||
const session = await getReadySession(request);
|
||||
const params = request.params as { taskId: string; artifactId: string };
|
||||
const artifactId = z.string().uuid().safeParse(params.artifactId);
|
||||
if (!artifactId.success) throw new TaskError('artifact_not_found', '附件不存在或无权访问。', 404);
|
||||
const artifact = await tasks.getTaskArtifact(session.user.organizationId, params.taskId, artifactId.data);
|
||||
const artifact = await tasks.getTaskArtifact(
|
||||
session.user.organizationId,
|
||||
params.taskId,
|
||||
artifactId.data,
|
||||
contextFor(session, request)
|
||||
);
|
||||
if (artifact.storage_backend === 'oss' && artifact.public_url) {
|
||||
reply.header('Cache-Control', 'no-store');
|
||||
return reply.redirect(artifact.public_url, 302);
|
||||
@@ -1090,13 +1305,49 @@ export async function buildServer({
|
||||
app.post('/api/tasks/bulk-delete', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const body = taskBulkDeleteSchema.parse(request.body);
|
||||
return { ok: true, ...(await tasks.hardDeleteTasks(contextFor(session, request), body.task_ids)) };
|
||||
return {
|
||||
ok: true,
|
||||
archived: true,
|
||||
...(await tasks.archiveTasks(contextFor(session, request), body.task_ids, body.reason))
|
||||
};
|
||||
});
|
||||
|
||||
app.post('/api/tasks/bulk-archive', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const body = taskBulkDeleteSchema.parse(request.body);
|
||||
return {
|
||||
ok: true,
|
||||
archived: true,
|
||||
...(await tasks.archiveTasks(contextFor(session, request), body.task_ids, body.reason))
|
||||
};
|
||||
});
|
||||
|
||||
app.delete('/api/tasks/:taskId', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const params = request.params as { taskId: string };
|
||||
return { ok: true, ...(await tasks.hardDeleteTask(contextFor(session, request), params.taskId)) };
|
||||
const body = taskArchiveSchema.parse(request.body || {});
|
||||
return {
|
||||
ok: true,
|
||||
archived: true,
|
||||
task: await tasks.archiveTask(contextFor(session, request), params.taskId, body.reason)
|
||||
};
|
||||
});
|
||||
|
||||
app.post('/api/tasks/:taskId/archive', 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)
|
||||
};
|
||||
});
|
||||
|
||||
app.post('/api/tasks/:taskId/restore', async (request) => {
|
||||
const session = await requireMutationSession(request);
|
||||
const params = request.params as { taskId: string };
|
||||
return { ok: true, restored: true, task: await tasks.restoreTask(contextFor(session, request), params.taskId) };
|
||||
});
|
||||
|
||||
app.post('/api/connections/heartbeat', async (request) => {
|
||||
@@ -1106,8 +1357,53 @@ export async function buildServer({
|
||||
return { ok: true, connected: true };
|
||||
});
|
||||
|
||||
app.get('/api/audit', async (request) => {
|
||||
const session = await requireAdminSession(request);
|
||||
const query = auditQuerySchema.parse(request.query || {});
|
||||
return {
|
||||
ok: true,
|
||||
...(await tasks.listAuditEvents(contextFor(session, request), {
|
||||
eventType: query.event_type,
|
||||
actorUserId: query.actor_user_id,
|
||||
entityType: query.entity_type,
|
||||
limit: query.limit,
|
||||
offset: query.offset
|
||||
}))
|
||||
};
|
||||
});
|
||||
|
||||
app.get('/api/operations-dashboard', async (request, reply) => {
|
||||
setAuthNoStore(reply);
|
||||
const session = await requireLeadershipSession(request);
|
||||
const query = operationsDashboardQuerySchema.parse(request.query || {});
|
||||
return {
|
||||
ok: true,
|
||||
...(await tasks.listOperationsDashboard(contextFor(session, request), {
|
||||
from: query.from,
|
||||
to: query.to,
|
||||
actorUserId: query.actor_user_id,
|
||||
businessRouteId: query.business_route_id,
|
||||
status: query.status,
|
||||
search: query.search,
|
||||
limit: query.limit,
|
||||
offset: query.offset
|
||||
}))
|
||||
};
|
||||
});
|
||||
|
||||
app.get('/api/operations-dashboard/tasks/:taskId', async (request, reply) => {
|
||||
setAuthNoStore(reply);
|
||||
const session = await requireLeadershipSession(request);
|
||||
const params = request.params as { taskId: string };
|
||||
return {
|
||||
ok: true,
|
||||
read_only: true,
|
||||
...(await tasks.getOperationsDashboardTask(contextFor(session, request), params.taskId))
|
||||
};
|
||||
});
|
||||
|
||||
app.get('/api/events', async (request, reply) => {
|
||||
const session = await getSession(request);
|
||||
const session = await getReadySession(request);
|
||||
const query = (request.query || {}) as Record<string, unknown>;
|
||||
const querySince = Number(query.since || 0);
|
||||
const reconnectSince = Number(request.headers['last-event-id'] || 0);
|
||||
@@ -1125,9 +1421,23 @@ export async function buildServer({
|
||||
});
|
||||
const send = (event: TaskEvent) => {
|
||||
if (event.organization_id !== session.user.organizationId) return;
|
||||
response.write(`id: ${event.id}\nevent: task\ndata: ${JSON.stringify(event)}\n\n`);
|
||||
if (!canAccessTask(contextFor(session, request), {
|
||||
createdBy: event.owner_user_id,
|
||||
source: event.task_source || 'agentbus'
|
||||
})) return;
|
||||
const publicEvent = {
|
||||
id: event.id,
|
||||
organization_id: event.organization_id,
|
||||
task_id: event.task_id,
|
||||
status: event.status,
|
||||
stage: event.stage,
|
||||
message: event.message,
|
||||
payload: event.payload,
|
||||
created_at: event.created_at
|
||||
};
|
||||
response.write(`id: ${event.id}\nevent: task\ndata: ${JSON.stringify(publicEvent)}\n\n`);
|
||||
};
|
||||
for (const event of await tasks.eventsSince(session.user.organizationId, since)) send(event);
|
||||
for (const event of await tasks.eventsSince(session.user.organizationId, since, contextFor(session, request))) send(event);
|
||||
const heartbeat = setInterval(() => response.write(': heartbeat\n\n'), 20_000);
|
||||
const onTask = (event: TaskEvent) => send(event);
|
||||
tasks.events.on('task', onTask);
|
||||
|
||||
Reference in new issue
Block a user