feat: bind AgentBus work to account workers

This commit is contained in:
inman committed 2026-09-02 15:09:07 +08:00
1 parent d034f649c4
commit 6f9fd0f0bd
31 files changed
+1119 -257

No files matched your search

+113 -26
View File
@@ -20,12 +20,14 @@ export interface AuthUser {
organizationId: string;
username: string;
role: AuthRole;
erpAccount: string | null;
}
export interface PublicAccount {
id: string;
username: string;
role: AuthRole;
erp_account: string | null;
is_active: boolean;
authorized_business_route_ids: BusinessRouteId[];
business_authorization_revision: number;
@@ -78,6 +80,27 @@ function validatePassword(value: string): string {
return password;
}
function normalizeErpAccount(value: unknown): string | null {
const normalized = String(value ?? '').trim();
if (!normalized) return null;
if (normalized.length > 200) {
throw new AuthError('erp_account_invalid', 'ERP 账号必须为 1—200 个字符。', 400);
}
return normalized;
}
function validateAccountRouting(role: AuthRole, value: unknown): string | null {
const erpAccount = normalizeErpAccount(value);
if (role === 'admin') {
if (erpAccount) throw new AuthError('admin_erp_account_forbidden', '管理员账号不能绑定员工 ERP 账号。', 409);
return null;
}
if (!erpAccount) {
throw new AuthError('erp_account_required', '普通用户或组长必须绑定 ERP 账号。', 400);
}
return erpAccount;
}
function normalizeRole(value: unknown): AuthRole {
if (value === 'admin' || value === 'team_lead') return value;
return 'user';
@@ -109,7 +132,8 @@ function mapUser(row: Record<string, unknown>): AuthUser {
id: String(row.id),
organizationId: String(row.organization_id),
username: String(row.username),
role: normalizeRole(row.role)
role: normalizeRole(row.role),
erpAccount: normalizeErpAccount(row.erp_account)
};
}
@@ -122,6 +146,7 @@ function mapAccount(row: Record<string, unknown>): PublicAccount {
id: String(row.id),
username: String(row.username),
role,
erp_account: normalizeErpAccount(row.erp_account),
is_active: row.is_active === true || String(row.is_active) === 'true',
authorized_business_route_ids: role === 'admin' ? [...ALL_BUSINESS_ROUTE_IDS] : storedRouteIds,
business_authorization_revision: Math.max(0, Number(row.business_authorization_revision || 0)),
@@ -137,7 +162,7 @@ async function loadPublicAccount(
userId: string
): Promise<PublicAccount | null> {
const result = await client.query(
`SELECT u.id, u.username, u.role, u.is_active,
`SELECT u.id, u.username, u.role, u.erp_account, u.is_active,
u.business_authorization_revision,
u.last_login_at, u.created_at, u.updated_at,
COALESCE(ARRAY(
@@ -209,21 +234,34 @@ export class AuthService {
const result = existing.rowCount
? await client.query(
`UPDATE users
SET password_hash = $1, role = 'admin', is_active = true,
SET password_hash = $1, role = 'admin', erp_account = NULL, is_active = true,
must_change_password = false, password_changed_at = now(),
failed_login_count = 0, locked_until = NULL, updated_at = now()
WHERE id = $2
RETURNING id, organization_id, username, role`,
RETURNING id, organization_id, username, role, erp_account`,
[passwordHash, existing.rows[0].id]
)
: await client.query(
`INSERT INTO users (organization_id, username, password_hash, role)
VALUES ($1, $2, $3, 'admin')
RETURNING id, organization_id, username, role`,
RETURNING id, organization_id, username, role, erp_account`,
[organization.id, normalized, passwordHash]
);
const user = mapUser(result.rows[0]);
if (existing.rowCount && force) {
await client.query(
`UPDATE user_channels
SET owner_user_id = NULL, enabled = false, status = 'disabled',
last_error = '绑定账号已重置为管理员,渠道已解除绑定。', updated_at = now()
WHERE organization_id = $1 AND owner_user_id = $2`,
[organization.id, user.id]
);
await client.query(
`UPDATE browser_connections
SET status = 'superseded', erp_account_verified = false
WHERE organization_id = $1 AND user_id = $2 AND status = 'connected'`,
[organization.id, user.id]
);
await client.query(
'UPDATE sessions SET revoked_at = now() WHERE user_id = $1 AND revoked_at IS NULL',
[user.id]
@@ -237,7 +275,7 @@ export class AuthService {
const normalized = normalizeUsername(username);
const pool = getPool(this.config);
const lookup = await pool.query(
`SELECT id, organization_id, username, password_hash, role, is_active,
`SELECT id, organization_id, username, password_hash, role, erp_account, is_active,
failed_login_count, locked_until
FROM users
WHERE organization_id = (SELECT id FROM organizations WHERE slug = $1)
@@ -330,7 +368,7 @@ export class AuthService {
async listAccounts(actor: AuthUser): Promise<PublicAccount[]> {
this.requireAdmin(actor);
const result = await getPool(this.config).query(
`SELECT u.id, u.username, u.role, u.is_active,
`SELECT u.id, u.username, u.role, u.erp_account, u.is_active,
u.business_authorization_revision,
u.last_login_at, u.created_at, u.updated_at,
COALESCE(ARRAY(
@@ -353,6 +391,7 @@ export class AuthService {
username: string;
password: string;
role: AuthRole;
erpAccount?: string;
businessRouteIds?: readonly string[];
},
requestId: string
@@ -361,8 +400,10 @@ export class AuthService {
const username = validateUsername(input.username);
const passwordHash = await argon2.hash(validatePassword(input.password), { type: argon2.argon2id });
const role = normalizeRole(input.role);
const erpAccount = validateAccountRouting(role, input.erpAccount);
const businessRouteIds = role === 'admin' ? [] : normalizeBusinessRouteIds(input.businessRouteIds);
return withTransaction(this.config, async (client) => {
try {
return await withTransaction(this.config, async (client) => {
await client.query(
`SELECT pg_advisory_xact_lock(hashtextextended($1::text || ':' || $2::text, 0))`,
[actor.organizationId, username]
@@ -374,10 +415,10 @@ export class AuthService {
if (existing.rowCount) throw new AuthError('account_exists', '该账号已存在。', 409);
const created = await client.query(
`INSERT INTO users
(organization_id, username, password_hash, role, password_changed_at)
VALUES ($1, $2, $3, $4, now())
(organization_id, username, password_hash, role, erp_account, password_changed_at)
VALUES ($1, $2, $3, $4, $5, now())
RETURNING id`,
[actor.organizationId, username, passwordHash, role]
[actor.organizationId, username, passwordHash, role, erpAccount]
);
const accountId = String(created.rows[0].id);
if (businessRouteIds.length) {
@@ -393,25 +434,36 @@ export class AuthService {
if (!account) throw new AuthError('account_not_found', '账号创建后未能读取。', 500);
await this.accountAudit(client, actor, 'account.created', account.id, requestId, {
role,
erp_account_configured: Boolean(erpAccount),
authorized_business_route_ids: account.authorized_business_route_ids
});
return account;
});
return account;
});
} catch (error) {
if (error && typeof error === 'object' && String((error as { code?: unknown }).code || '') === '23505') {
const constraint = String((error as { constraint?: unknown }).constraint || '');
if (constraint.includes('erp_account')) {
throw new AuthError('erp_account_conflict', '该 ERP 账号已经绑定另一个平台账号。', 409);
}
}
throw error;
}
}
async updateAccount(
actor: AuthUser,
targetUserId: string,
input: { role?: AuthRole; isActive?: boolean },
input: { role?: AuthRole; isActive?: boolean; erpAccount?: string | null },
requestId: string
): Promise<PublicAccount> {
this.requireAdmin(actor);
if (input.role === undefined && input.isActive === undefined) {
if (input.role === undefined && input.isActive === undefined && input.erpAccount === undefined) {
throw new AuthError('account_update_empty', '没有需要更新的账号字段。', 400);
}
return withTransaction(this.config, async (client) => {
try {
return await withTransaction(this.config, async (client) => {
const target = await client.query(
`SELECT id, username, role, is_active,
`SELECT id, username, role, erp_account, is_active,
last_login_at, created_at, updated_at
FROM users
WHERE organization_id = $1 AND id = $2
@@ -422,6 +474,12 @@ export class AuthService {
const before = mapAccount(target.rows[0] as Record<string, unknown>);
const role = input.role === undefined ? before.role : normalizeRole(input.role);
const isActive = input.isActive === undefined ? before.is_active : input.isActive;
const erpAccount = validateAccountRouting(
role,
role === 'admin'
? null
: input.erpAccount === undefined ? before.erp_account : input.erpAccount
);
const removesActiveAdmin = before.role === 'admin' && before.is_active && (role !== 'admin' || !isActive);
if (actor.id === before.id && (role !== 'admin' || !isActive)) {
throw new AuthError('self_lockout_forbidden', '不能停用或降级当前登录的管理员账号。', 409);
@@ -439,28 +497,56 @@ export class AuthService {
}
const updated = await client.query(
`UPDATE users
SET role = $1, is_active = $2, updated_at = now()
WHERE organization_id = $3 AND id = $4
SET role = $1, is_active = $2, erp_account = $3, updated_at = now()
WHERE organization_id = $4 AND id = $5
RETURNING id`,
[role, isActive, actor.organizationId, before.id]
[role, isActive, erpAccount, actor.organizationId, before.id]
);
if (before.role !== role || before.is_active !== isActive) {
const routingIdentityChanged = before.erp_account !== erpAccount;
if (before.role !== role || before.is_active !== isActive || routingIdentityChanged) {
await client.query(
'UPDATE sessions SET revoked_at = now() WHERE user_id = $1 AND revoked_at IS NULL',
[before.id]
);
}
if (!isActive || role === 'admin') {
await client.query(
`UPDATE user_channels
SET owner_user_id = NULL, enabled = false, status = 'disabled',
last_error = $1, updated_at = now()
WHERE organization_id = $2 AND owner_user_id = $3`,
[role === 'admin' ? '绑定账号已变更为管理员,渠道已解除绑定。' : '绑定账号已停用,渠道已解除绑定。', actor.organizationId, before.id]
);
}
if (!isActive || routingIdentityChanged || before.role !== role) {
await client.query(
`UPDATE browser_connections
SET status = 'superseded', erp_account_verified = false
WHERE organization_id = $1 AND user_id = $2 AND status = 'connected'`,
[actor.organizationId, before.id]
);
}
await this.accountAudit(client, actor, 'account.updated', before.id, requestId, {
previous_role: before.role,
role,
previous_active: before.is_active,
active: isActive,
sessions_revoked: before.role !== role || before.is_active !== isActive
erp_account_changed: routingIdentityChanged,
sessions_revoked: before.role !== role || before.is_active !== isActive || routingIdentityChanged
});
const account = await loadPublicAccount(client, actor.organizationId, String(updated.rows[0].id));
if (!account) throw new AuthError('account_not_found', '账号更新后未能读取。', 500);
return account;
});
return account;
});
} catch (error) {
if (error && typeof error === 'object' && String((error as { code?: unknown }).code || '') === '23505') {
const constraint = String((error as { constraint?: unknown }).constraint || '');
if (constraint.includes('erp_account')) {
throw new AuthError('erp_account_conflict', '该 ERP 账号已经绑定另一个平台账号。', 409);
}
}
throw error;
}
}
async setBusinessRouteAuthorizations(
@@ -648,7 +734,7 @@ export class AuthService {
async getActiveSession(token: string | undefined): Promise<ActiveSession | null> {
if (!token) return null;
const result = await getPool(this.config).query(
`SELECT s.id AS session_id, s.csrf_token_hash, u.id, u.organization_id, u.username, u.role
`SELECT s.id AS session_id, s.csrf_token_hash, u.id, u.organization_id, u.username, u.role, u.erp_account
FROM sessions s
JOIN users u ON u.id = s.user_id
WHERE s.token_hash = $1
@@ -671,7 +757,8 @@ export class AuthService {
id: row.id,
organization_id: row.organization_id,
username: row.username,
role: row.role
role: row.role,
erp_account: row.erp_account
})
};
}