merge: integrate AgentBus account workers

This commit is contained in:
inman committed 2026-09-02 15:15:01 +08:00
commit b5f58477d9
30 files changed
+1054 -257

No files matched your search

+30 -8
View File
@@ -63,6 +63,7 @@ const accountCreateSchema = z.object({
username: z.string().min(1).max(160),
password: z.string().min(1),
role: z.enum(['admin', 'team_lead', 'user']).default('user'),
erp_account: z.string().trim().max(200).optional(),
business_route_ids: z.array(
z.string().trim().refine((routeId) => Boolean(businessRouteById(routeId)), '业务类型不存在。')
).max(BUSINESS_ROUTES.length).default([])
@@ -71,8 +72,9 @@ const accountCreateSchema = z.object({
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, {
is_active: z.boolean().optional(),
erp_account: z.string().trim().max(200).nullable().optional()
}).refine((body) => body.role !== undefined || body.is_active !== undefined || body.erp_account !== undefined, {
message: '至少提供一个账号更新字段。'
});
@@ -123,6 +125,8 @@ const resultSchema = z.object({
const heartbeatSchema = z.object({
connection_id: z.string().min(1).max(200),
extension_version: z.string().max(80).optional(),
erp_account: z.string().trim().max(200).optional(),
erp_account_matched: z.boolean().optional(),
metadata: z.record(z.unknown()).optional()
});
const automationSettingsSchema = z.object({ enabled: z.boolean() });
@@ -138,6 +142,7 @@ const parserDecisionReviewSchema = z.object({
});
const channelCreateSchema = z.object({
display_name: z.string().min(1).max(120),
owner_user_id: z.string().uuid(),
external_user_ref: z.string().max(200).optional(),
agentbus_key: z.string().min(1).max(4_000),
bot_address: z.string().max(200).optional(),
@@ -145,6 +150,7 @@ const channelCreateSchema = z.object({
});
const channelUpdateSchema = z.object({
display_name: z.string().min(1).max(120).optional(),
owner_user_id: z.string().uuid().optional(),
external_user_ref: z.string().max(200).optional(),
bot_address: z.string().max(200).optional(),
enabled: z.boolean().optional()
@@ -156,6 +162,7 @@ const listTasksQuerySchema = z.object({
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'),
executable_by: z.enum(['me']).optional(),
include_total: z.preprocess(
(value) => value === undefined ? true : String(value).toLowerCase() !== 'false',
z.boolean()
@@ -361,7 +368,8 @@ function publicUser(session: ActiveSession) {
return {
id: session.user.id,
username: session.user.username,
role: session.user.role
role: session.user.role,
erp_account: session.user.erpAccount
};
}
@@ -963,6 +971,7 @@ export async function buildServer({
username: body.username,
password: body.password,
role: body.role,
erpAccount: body.erp_account,
businessRouteIds: body.business_route_ids
}, requestId(request));
return { ok: true, account };
@@ -975,8 +984,10 @@ export async function buildServer({
const body = accountUpdateSchema.parse(request.body);
const account = await auth.updateAccount(session.user, userId, {
role: body.role,
isActive: body.is_active
isActive: body.is_active,
erpAccount: body.erp_account
}, requestId(request));
await agentBus?.reload();
return { ok: true, account };
});
@@ -1098,6 +1109,7 @@ export async function buildServer({
const body = channelCreateSchema.parse(request.body);
const channel = await channelService.create(contextFor(session, request), {
displayName: body.display_name,
ownerUserId: body.owner_user_id,
externalUserRef: body.external_user_ref,
agentbusKey: body.agentbus_key,
botAddress: body.bot_address,
@@ -1113,6 +1125,7 @@ export async function buildServer({
const body = channelUpdateSchema.parse(request.body);
const channel = await channelService.update(contextFor(session, request), params.channelId, {
displayName: body.display_name,
ownerUserId: body.owner_user_id,
externalUserRef: body.external_user_ref,
botAddress: body.bot_address,
enabled: body.enabled
@@ -1162,6 +1175,7 @@ export async function buildServer({
offset: query.offset,
includeTotal: query.include_total,
archive: query.archive,
assignedToUserId: query.executable_by === 'me' ? session.user.id : undefined,
access: contextFor(session, request)
});
return {
@@ -1334,8 +1348,17 @@ export async function buildServer({
app.post('/api/connections/heartbeat', async (request) => {
const session = await requireMutationSession(request);
const body = heartbeatSchema.parse(request.body);
await tasks.heartbeat(contextFor(session, request), body.connection_id, body.extension_version || '', body.metadata || {});
return { ok: true, connected: true };
const worker = await tasks.heartbeat(
contextFor(session, request),
body.connection_id,
body.extension_version || '',
body.metadata || {},
{
erpAccount: body.erp_account || '',
erpAccountMatched: body.erp_account_matched === true
}
);
return { ok: true, connected: true, ...worker };
});
app.get('/api/audit', async (request) => {
@@ -1417,8 +1440,7 @@ export async function buildServer({
const send = (event: TaskEvent) => {
if (event.organization_id !== session.user.organizationId) return;
if (!canAccessTask(contextFor(session, request), {
createdBy: event.owner_user_id,
source: event.task_source || 'agentbus'
assignedUserId: event.owner_user_id
})) return;
const publicEvent = {
id: event.id,