Revert "merge: integrate extension auto-update"

This reverts commit 322475860a, reversing
changes made to f52d9d7413.
This commit is contained in:
inman
2026-09-03 16:45:14 +08:00
parent b2e33e2e5d
commit 81a0cdac8e
47 changed files with 169 additions and 3108 deletions

View File

@@ -28,8 +28,6 @@ export interface PublicAccount {
username: string;
role: AuthRole;
erp_account: string | null;
extension_ecs_region_id: string | null;
extension_ecs_instance_id: string | null;
is_active: boolean;
authorized_business_route_ids: BusinessRouteId[];
business_authorization_revision: number;
@@ -103,42 +101,6 @@ function validateAccountRouting(role: AuthRole, value: unknown): string | null {
return erpAccount;
}
function normalizeEcsRegionId(value: unknown): string | null {
const normalized = String(value ?? '').trim().toLocaleLowerCase('en-US');
if (!normalized) return null;
if (!/^[a-z0-9][a-z0-9-]{0,63}$/u.test(normalized)) {
throw new AuthError('extension_ecs_region_invalid', 'ECS 地域 ID 格式无效。', 400);
}
return normalized;
}
function normalizeEcsInstanceId(value: unknown): string | null {
const normalized = String(value ?? '').trim();
if (!normalized) return null;
if (!/^i-[A-Za-z0-9]{6,64}$/u.test(normalized)) {
throw new AuthError('extension_ecs_instance_invalid', 'ECS 实例 ID 格式无效。', 400);
}
return normalized;
}
function validateExtensionHostBinding(
role: AuthRole,
regionValue: unknown,
instanceValue: unknown
): { regionId: string | null; instanceId: string | null } {
const regionId = normalizeEcsRegionId(regionValue);
const instanceId = normalizeEcsInstanceId(instanceValue);
if (role === 'admin' && (regionId || instanceId)) {
throw new AuthError('admin_extension_host_forbidden', '管理员账号不能绑定 ERP 插件云主机。', 409);
}
if (Boolean(regionId) !== Boolean(instanceId)) {
throw new AuthError('extension_ecs_binding_incomplete', 'ECS 地域 ID 与实例 ID 必须同时填写或同时清空。', 400);
}
return role === 'admin'
? { regionId: null, instanceId: null }
: { regionId, instanceId };
}
function normalizeRole(value: unknown): AuthRole {
if (value === 'admin' || value === 'team_lead') return value;
return 'user';
@@ -185,8 +147,6 @@ function mapAccount(row: Record<string, unknown>): PublicAccount {
username: String(row.username),
role,
erp_account: normalizeErpAccount(row.erp_account),
extension_ecs_region_id: normalizeEcsRegionId(row.extension_ecs_region_id),
extension_ecs_instance_id: normalizeEcsInstanceId(row.extension_ecs_instance_id),
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)),
@@ -202,8 +162,7 @@ async function loadPublicAccount(
userId: string
): Promise<PublicAccount | null> {
const result = await client.query(
`SELECT u.id, u.username, u.role, u.erp_account,
u.extension_ecs_region_id, u.extension_ecs_instance_id, 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(
@@ -409,8 +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.erp_account,
u.extension_ecs_region_id, u.extension_ecs_instance_id, 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(
@@ -434,8 +392,6 @@ export class AuthService {
password: string;
role: AuthRole;
erpAccount?: string;
extensionEcsRegionId?: string;
extensionEcsInstanceId?: string;
businessRouteIds?: readonly string[];
},
requestId: string
@@ -445,11 +401,6 @@ export class AuthService {
const passwordHash = await argon2.hash(validatePassword(input.password), { type: argon2.argon2id });
const role = normalizeRole(input.role);
const erpAccount = validateAccountRouting(role, input.erpAccount);
const extensionHost = validateExtensionHostBinding(
role,
input.extensionEcsRegionId,
input.extensionEcsInstanceId
);
const businessRouteIds = role === 'admin' ? [] : normalizeBusinessRouteIds(input.businessRouteIds);
try {
return await withTransaction(this.config, async (client) => {
@@ -464,19 +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, erp_account,
extension_ecs_region_id, extension_ecs_instance_id, password_changed_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, 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,
erpAccount,
extensionHost.regionId,
extensionHost.instanceId
]
[actor.organizationId, username, passwordHash, role, erpAccount]
);
const accountId = String(created.rows[0].id);
if (businessRouteIds.length) {
@@ -493,7 +435,6 @@ export class AuthService {
await this.accountAudit(client, actor, 'account.created', account.id, requestId, {
role,
erp_account_configured: Boolean(erpAccount),
extension_host_configured: Boolean(extensionHost.instanceId),
authorized_business_route_ids: account.authorized_business_route_ids
});
return account;
@@ -512,30 +453,17 @@ export class AuthService {
async updateAccount(
actor: AuthUser,
targetUserId: string,
input: {
role?: AuthRole;
isActive?: boolean;
erpAccount?: string | null;
extensionEcsRegionId?: string | null;
extensionEcsInstanceId?: string | null;
},
input: { role?: AuthRole; isActive?: boolean; erpAccount?: string | null },
requestId: string
): Promise<PublicAccount> {
this.requireAdmin(actor);
if (
input.role === undefined
&& input.isActive === undefined
&& input.erpAccount === undefined
&& input.extensionEcsRegionId === undefined
&& input.extensionEcsInstanceId === undefined
) {
if (input.role === undefined && input.isActive === undefined && input.erpAccount === undefined) {
throw new AuthError('account_update_empty', '没有需要更新的账号字段。', 400);
}
try {
return await withTransaction(this.config, async (client) => {
const target = await client.query(
`SELECT id, username, role, erp_account,
extension_ecs_region_id, extension_ecs_instance_id, 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
@@ -552,19 +480,6 @@ export class AuthService {
? null
: input.erpAccount === undefined ? before.erp_account : input.erpAccount
);
const extensionHost = validateExtensionHostBinding(
role,
role === 'admin'
? null
: input.extensionEcsRegionId === undefined
? before.extension_ecs_region_id
: input.extensionEcsRegionId,
role === 'admin'
? null
: input.extensionEcsInstanceId === undefined
? before.extension_ecs_instance_id
: input.extensionEcsInstanceId
);
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);
@@ -582,24 +497,12 @@ export class AuthService {
}
const updated = await client.query(
`UPDATE users
SET role = $1, is_active = $2, erp_account = $3,
extension_ecs_region_id = $4, extension_ecs_instance_id = $5,
updated_at = now()
WHERE organization_id = $6 AND id = $7
SET role = $1, is_active = $2, erp_account = $3, updated_at = now()
WHERE organization_id = $4 AND id = $5
RETURNING id`,
[
role,
isActive,
erpAccount,
extensionHost.regionId,
extensionHost.instanceId,
actor.organizationId,
before.id
]
[role, isActive, erpAccount, actor.organizationId, before.id]
);
const routingIdentityChanged = before.erp_account !== erpAccount;
const extensionHostChanged = before.extension_ecs_region_id !== extensionHost.regionId
|| before.extension_ecs_instance_id !== extensionHost.instanceId;
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',
@@ -615,7 +518,7 @@ export class AuthService {
[role === 'admin' ? '绑定账号已变更为管理员,渠道已解除绑定。' : '绑定账号已停用,渠道已解除绑定。', actor.organizationId, before.id]
);
}
if (!isActive || routingIdentityChanged || extensionHostChanged || before.role !== role) {
if (!isActive || routingIdentityChanged || before.role !== role) {
await client.query(
`UPDATE browser_connections
SET status = 'superseded', erp_account_verified = false
@@ -629,7 +532,6 @@ export class AuthService {
previous_active: before.is_active,
active: isActive,
erp_account_changed: routingIdentityChanged,
extension_host_changed: extensionHostChanged,
sessions_revoked: before.role !== role || before.is_active !== isActive || routingIdentityChanged
});
const account = await loadPublicAccount(client, actor.organizationId, String(updated.rows[0].id));

View File

@@ -32,20 +32,6 @@ const optionalOssEndpoint = z.preprocess(
z.string().regex(/^[a-z0-9.-]+$/i).optional()
);
const optionalWindowsExtensionPath = z.preprocess(
(value) => {
const normalized = String(value ?? '').trim().replace(/\/+$/u, '');
return normalized || undefined;
},
z.string()
.max(240)
.refine(
(value) => /^[A-Za-z]:\\ProgramData\\LTJT\\[A-Za-z0-9._\\-]+$/u.test(value) && !value.includes('..'),
'EXTENSION_WINDOWS_INSTALL_PATH must be an absolute path below <drive>:\\ProgramData\\LTJT.'
)
.optional()
);
function parseDurationMs(value: string): number {
const normalized = String(value || '').trim().toLowerCase();
const match = /^(\d+(?:\.\d+)?)(ms|s|m|h)?$/.exec(normalized);
@@ -111,18 +97,6 @@ const envSchema = z.object({
OSS_BUCKET_NAME: optionalString,
OSS_REGION: optionalString,
OSS_KEY_PREFIX: z.string().trim().min(1).max(200).default('liansyn-platform/attachments'),
EXTENSION_AUTO_UPDATE_ENABLED: z.enum(['true', 'false']).default('false').transform((value) => value === 'true'),
EXTENSION_UPDATE_OSS_KEY_PREFIX: z.string().trim().min(1).max(200)
.default('liansyn-platform/chrome-extension')
.transform((value) => value.replace(/^\/+|\/+$/gu, '')),
EXTENSION_WINDOWS_INSTALL_PATH: optionalWindowsExtensionPath
.default('C:\\ProgramData\\LTJT\\chrome-extension\\ltjt-order-assistant'),
EXTENSION_UPDATE_MAX_PACKAGE_BYTES: z.coerce.number().int().positive().max(50_000_000).default(15_000_000),
EXTENSION_UPDATE_COMMAND_TIMEOUT_SECONDS: z.coerce.number().int().min(60).max(3_600).default(600),
EXTENSION_UPDATE_POLL_INTERVAL_MS: z.coerce.number().int().min(1_000).max(60_000).default(5_000),
ALIBABA_CLOUD_ACCESS_KEY_ID: optionalString,
ALIBABA_CLOUD_ACCESS_KEY_SECRET: optionalString,
ALIBABA_CLOUD_SECURITY_TOKEN: optionalString,
DATA_RETENTION_ENABLED: z.enum(['true', 'false']).default('false').transform((value) => value === 'true'),
DATA_RETENTION_DAYS: z.coerce.number().int().positive().default(180),
LOG_LEVEL: z.enum(['trace', 'debug', 'info', 'warn', 'error', 'fatal', 'silent']).default('info')
@@ -186,23 +160,6 @@ export function loadConfig(env: NodeJS.ProcessEnv = process.env): AppConfig {
if (parsed.ARTIFACT_STORAGE_BACKEND === 'oss' && !ossRegion) {
throw new Error('OSS_REGION is required when OSS_ENDPOINT does not use the standard oss-<region> endpoint format.');
}
if (parsed.EXTENSION_AUTO_UPDATE_ENABLED) {
const missing = [
['OSS_ACCESS_KEY_ID', parsed.OSS_ACCESS_KEY_ID],
['OSS_ACCESS_KEY_SECRET', parsed.OSS_ACCESS_KEY_SECRET],
['OSS_ENDPOINT', parsed.OSS_ENDPOINT],
['OSS_BUCKET_NAME', parsed.OSS_BUCKET_NAME],
['OSS_REGION', ossRegion],
['ALIBABA_CLOUD_ACCESS_KEY_ID', parsed.ALIBABA_CLOUD_ACCESS_KEY_ID],
['ALIBABA_CLOUD_ACCESS_KEY_SECRET', parsed.ALIBABA_CLOUD_ACCESS_KEY_SECRET]
].filter(([, value]) => !value).map(([name]) => name);
if (missing.length) {
throw new Error(`Extension auto-update is enabled but missing configuration: ${missing.join(', ')}`);
}
if (parsed.NODE_ENV === 'production' && new URL(parsed.APP_ORIGIN).protocol !== 'https:') {
throw new Error('APP_ORIGIN must use HTTPS when extension auto-update is enabled in production.');
}
}
return {
...parsed,
fieldEncryptionKey: decodeEncryptionKey(parsed.FIELD_ENCRYPTION_KEY, parsed.NODE_ENV),

View File

@@ -5,7 +5,7 @@ import { writeEmergencyDiagnostic } from './diagnostics.js';
const { Pool } = pg;
let pool: pg.Pool | null = null;
export const REQUIRED_SCHEMA_VERSION = '019_extension_host_updates';
export const REQUIRED_SCHEMA_VERSION = '018_agentbus_account_workers';
export interface DatabaseReadiness {
ready: boolean;

View File

@@ -97,17 +97,9 @@ export function diagnosticRequestPath(value: unknown): string {
const raw = text(value);
if (!raw) return '/';
try {
const pathname = new URL(raw, 'http://diagnostic.invalid').pathname;
if (pathname.startsWith('/api/extension-updates/package/')) {
return '/api/extension-updates/package/:token';
}
return pathname.slice(0, 500) || '/';
return new URL(raw, 'http://diagnostic.invalid').pathname.slice(0, 500) || '/';
} catch {
const pathname = raw.split(/[?#]/u, 1)[0];
if (pathname.startsWith('/api/extension-updates/package/')) {
return '/api/extension-updates/package/:token';
}
return pathname.slice(0, 500) || '/';
return raw.split(/[?#]/u, 1)[0].slice(0, 500) || '/';
}
}

File diff suppressed because it is too large Load Diff

View File

@@ -9,7 +9,6 @@ export interface OssPutObjectInput {
content: Buffer;
contentType: string;
contentDisposition?: string;
acl?: 'private' | 'public-read';
}
export interface OssObjectClient {
@@ -169,7 +168,7 @@ export class AliyunOssClient implements OssObjectClient {
const headers: Record<string, string> = {
'content-type': input.contentType,
'x-oss-content-sha256': UNSIGNED_PAYLOAD,
'x-oss-object-acl': input.acl || 'public-read'
'x-oss-object-acl': 'public-read'
};
if (input.contentDisposition) headers['content-disposition'] = input.contentDisposition;
const additionalHeaders = input.contentDisposition ? ['content-disposition'] : [];

View File

@@ -44,11 +44,6 @@ import {
normalizeRequestId,
writeEmergencyDiagnostic
} from './diagnostics.js';
import {
ExtensionUpdateError,
createExtensionUpdateService,
type ExtensionUpdateController
} from './extension-updates.js';
export function aiServiceConnected(databaseIsReady: boolean, probe: Record<string, unknown>): boolean {
return databaseIsReady && probe.configured === true && probe.reachable === true;
@@ -69,8 +64,6 @@ const accountCreateSchema = z.object({
password: z.string().min(1),
role: z.enum(['admin', 'team_lead', 'user']).default('user'),
erp_account: z.string().trim().max(200).optional(),
extension_ecs_region_id: z.string().trim().max(64).optional(),
extension_ecs_instance_id: z.string().trim().max(80).optional(),
business_route_ids: z.array(
z.string().trim().refine((routeId) => Boolean(businessRouteById(routeId)), '业务类型不存在。')
).max(BUSINESS_ROUTES.length).default([])
@@ -80,16 +73,8 @@ const accountCreateSchema = z.object({
const accountUpdateSchema = z.object({
role: z.enum(['admin', 'team_lead', 'user']).optional(),
is_active: z.boolean().optional(),
erp_account: z.string().trim().max(200).nullable().optional(),
extension_ecs_region_id: z.string().trim().max(64).nullable().optional(),
extension_ecs_instance_id: z.string().trim().max(80).nullable().optional()
}).refine((body) => (
body.role !== undefined
|| body.is_active !== undefined
|| body.erp_account !== undefined
|| body.extension_ecs_region_id !== undefined
|| body.extension_ecs_instance_id !== undefined
), {
erp_account: z.string().trim().max(200).nullable().optional()
}).refine((body) => body.role !== undefined || body.is_active !== undefined || body.erp_account !== undefined, {
message: '至少提供一个账号更新字段。'
});
@@ -142,12 +127,8 @@ const heartbeatSchema = z.object({
extension_version: z.string().max(80).optional(),
erp_account: z.string().trim().max(200).optional(),
erp_account_matched: z.boolean().optional(),
extension_update_safe: z.boolean().optional(),
metadata: z.record(z.unknown()).optional()
});
const extensionReleasePublishSchema = z.object({
package_base64: z.string().min(4).max(70_000_000)
});
const automationSettingsSchema = z.object({ enabled: z.boolean() });
const parserRoutingUpdateSchema = z.object({
mode: z.enum(['ai', 'shadow', 'auto', 'program']),
@@ -415,14 +396,12 @@ export async function buildServer({
config = loadConfig(),
parser,
startParserLoop = true,
loggerDestination,
extensionUpdates: extensionUpdatesOverride
loggerDestination
}: {
config?: AppConfig;
parser?: ExternalParser;
startParserLoop?: boolean;
loggerDestination?: DestinationStream;
extensionUpdates?: ExtensionUpdateController;
} = {}) {
const requestStartedAt = new WeakMap<FastifyRequest, bigint>();
const app = Fastify({
@@ -430,11 +409,7 @@ export async function buildServer({
logController: new LogController({ disableRequestLogging: true }),
genReqId: (rawRequest) => normalizeRequestId(rawRequest.headers['x-request-id']),
trustProxy: true,
bodyLimit: Math.min(75_000_000, Math.max(
2_000_000,
Math.ceil(config.ARTIFACT_MAX_BYTES * 1.4) + 1_000_000,
Math.ceil(config.EXTENSION_UPDATE_MAX_PACKAGE_BYTES * 1.4) + 1_000_000
))
bodyLimit: Math.min(75_000_000, Math.max(2_000_000, Math.ceil(config.ARTIFACT_MAX_BYTES * 1.4) + 1_000_000))
});
await app.register(cookie);
await app.register(helmet, { contentSecurityPolicy: false });
@@ -536,11 +511,6 @@ export async function buildServer({
warn: (metadata, message) => app.log.warn(metadata, message),
error: (metadata, message) => app.log.error(metadata, message)
});
const extensionUpdates = extensionUpdatesOverride || createExtensionUpdateService(config, {
info: (metadata, message) => app.log.info(metadata, message),
warn: (metadata, message) => app.log.warn(metadata, message),
error: (metadata, message) => app.log.error(metadata, message)
});
const externalParser = parser || await loadExternalParser();
const parserOrchestrator = new ParserOrchestrator(externalParser);
const activeParseWorkers = new Map<string, string>();
@@ -1013,8 +983,6 @@ export async function buildServer({
password: body.password,
role: body.role,
erpAccount: body.erp_account,
extensionEcsRegionId: body.extension_ecs_region_id,
extensionEcsInstanceId: body.extension_ecs_instance_id,
businessRouteIds: body.business_route_ids
}, requestId(request));
return { ok: true, account };
@@ -1028,9 +996,7 @@ export async function buildServer({
const account = await auth.updateAccount(session.user, userId, {
role: body.role,
isActive: body.is_active,
erpAccount: body.erp_account,
extensionEcsRegionId: body.extension_ecs_region_id,
extensionEcsInstanceId: body.extension_ecs_instance_id
erpAccount: body.erp_account
}, requestId(request));
await agentBus?.reload();
return { ok: true, account };
@@ -1395,77 +1361,13 @@ export async function buildServer({
contextFor(session, request),
body.connection_id,
body.extension_version || '',
{
...(body.metadata || {}),
extension_update_safe: body.extension_update_safe === true
},
body.metadata || {},
{
erpAccount: body.erp_account || '',
erpAccountMatched: body.erp_account_matched === true
}
);
const { extension_update_target: updateTarget, ...publicWorker } = worker;
const extensionUpdate = await extensionUpdates.observeHeartbeat({
organizationId: session.user.organizationId,
userId: session.user.id,
connectionId: body.connection_id,
currentVersion: body.extension_version || '',
extensionUpdateSafe: body.extension_update_safe === true,
target: {
ecsRegionId: updateTarget.ecs_region_id,
ecsInstanceId: updateTarget.ecs_instance_id,
serverUpdateSafe: updateTarget.server_update_safe
}
});
return { ok: true, connected: true, ...publicWorker, extension_update: extensionUpdate };
});
app.get('/api/extension-updates/releases', async (request) => {
const session = await requireAdminSession(request);
return {
ok: true,
enabled: config.EXTENSION_AUTO_UPDATE_ENABLED,
releases: await extensionUpdates.listReleases(session.user.organizationId)
};
});
app.post('/api/extension-updates/releases', {
config: { rateLimit: { max: 5, timeWindow: '1 minute' } }
}, async (request) => {
const session = await requireAdminMutationSession(request);
const body = extensionReleasePublishSchema.parse(request.body);
const normalized = body.package_base64.replace(/\s+/gu, '');
if (!/^[A-Za-z0-9+/]+={0,2}$/u.test(normalized)) {
throw new ExtensionUpdateError('extension_package_base64_invalid', '插件包编码无效。');
}
const content = Buffer.from(normalized, 'base64');
const canonical = content.toString('base64').replace(/=+$/u, '');
if (canonical !== normalized.replace(/=+$/u, '')) {
throw new ExtensionUpdateError('extension_package_base64_invalid', '插件包编码无效。');
}
return {
ok: true,
release: await extensionUpdates.publishRelease({
organizationId: session.user.organizationId,
actorUserId: session.user.id,
requestId: requestId(request),
content
})
};
});
app.get('/api/extension-updates/package/:token', {
// A fleet may share one outbound NAT address and start together after a
// release. The HMAC token is already release/organization/host/expiry
// scoped, so keep only a generous abuse ceiling here.
config: { rateLimit: { max: 300, timeWindow: '1 minute' } }
}, async (request, reply) => {
const params = request.params as { token: string };
const download = await extensionUpdates.downloadPackage(params.token);
reply.header('Cache-Control', 'private, no-store, max-age=0');
reply.header('Content-Type', 'application/zip');
reply.header('Content-Disposition', attachmentContentDisposition(download.fileName));
return reply.send(download.content);
return { ok: true, connected: true, ...worker };
});
app.get('/api/audit', async (request) => {
@@ -1586,7 +1488,7 @@ export async function buildServer({
});
app.setErrorHandler((error, request, reply) => {
if (error instanceof AuthError || error instanceof TaskError || error instanceof ExtensionUpdateError) {
if (error instanceof AuthError || error instanceof TaskError) {
request.log.warn({
diagnostic_event: 'http.request.rejected',
diagnostic_stage: 'http',
@@ -1633,14 +1535,13 @@ export async function buildServer({
}
app.addHook('onClose', async () => {
extensionUpdates.close();
app.log.info({
diagnostic_event: 'service.closing',
diagnostic_stage: 'shutdown'
}, 'control plane closing');
await closePool();
});
return { app, auth, tasks, agentBus, channelService, extensionUpdates };
return { app, auth, tasks, agentBus, channelService };
}
function installProcessDiagnostics(app: Awaited<ReturnType<typeof buildServer>>['app']): void {

View File

@@ -49,7 +49,6 @@ import {
createAgentBusAcceptedDeliveryPayload,
type AgentBusAcceptedDeliveryOptions
} from './agentbus-delivery.js';
import { compareExtensionVersions } from './extension-updates.js';
export interface TaskEvent {
id: number;
@@ -2721,84 +2720,6 @@ export class TaskService {
}
}
private async requireExtensionHostAvailable(
client: import('pg').PoolClient,
context: TaskContext,
connectionId: string
): Promise<void> {
const update = await client.query(
`SELECT host_update.status
FROM users account
JOIN extension_host_updates host_update
ON host_update.organization_id = account.organization_id
AND host_update.ecs_region_id = account.extension_ecs_region_id
AND host_update.ecs_instance_id = account.extension_ecs_instance_id
WHERE account.organization_id = $1
AND account.id = $2
AND host_update.status IN ('waiting_for_idle', 'running', 'deployed')
FOR UPDATE OF host_update`,
[context.organizationId, context.userId]
);
if (update.rowCount) {
throw new TaskError(
'extension_host_update_in_progress',
'当前云主机正在等待或执行插件更新;新 ERP 任务会在更新并验版后继续。',
409,
{ update_status: text(update.rows[0].status) }
);
}
if (!this.config.EXTENSION_AUTO_UPDATE_ENABLED) return;
const requiredVersion = await client.query(
`SELECT release.version AS target_version,
connection.extension_version AS current_version,
account.extension_ecs_region_id,
account.extension_ecs_instance_id,
host_update.status AS update_status
FROM extension_releases release
JOIN browser_connections connection
ON connection.organization_id = release.organization_id
AND connection.user_id = $2
AND connection.connection_id = $3
JOIN users account
ON account.organization_id = connection.organization_id
AND account.id = connection.user_id
LEFT JOIN extension_host_updates host_update
ON host_update.organization_id = account.organization_id
AND host_update.ecs_region_id = account.extension_ecs_region_id
AND host_update.ecs_instance_id = account.extension_ecs_instance_id
AND host_update.release_id = release.id
WHERE release.organization_id = $1
AND release.is_active = true`,
[context.organizationId, context.userId, connectionId]
);
if (!requiredVersion.rowCount) return;
const versionRow = requiredVersion.rows[0] as Record<string, unknown>;
const currentVersion = text(versionRow.current_version);
const targetVersion = text(versionRow.target_version);
let isCurrent = false;
try {
isCurrent = compareExtensionVersions(currentVersion, targetVersion) >= 0;
} catch {
isCurrent = false;
}
if (!isCurrent) {
throw new TaskError(
'extension_update_required',
`当前插件版本 ${currentVersion || '未知'} 尚未更新到 ${targetVersion};新 ERP 任务会在更新并验版后继续。`,
409,
{
current_version: currentVersion || null,
target_version: targetVersion,
update_status: text(versionRow.update_status) || (
versionRow.extension_ecs_region_id && versionRow.extension_ecs_instance_id
? 'not_started'
: 'unconfigured'
)
}
);
}
}
private log(
level: 'info' | 'warn' | 'error',
metadata: Record<string, unknown>,
@@ -6292,7 +6213,6 @@ export class TaskService {
}
await this.assertTaskCreatorBusinessAuthorizationInTransaction(client, row, context.requestId);
await this.requireBrowserConnection(client, context, connectionId);
await this.requireExtensionHostAvailable(client, context, connectionId);
const existingAttempt = await client.query(
`SELECT id, status, response_hash
FROM task_attempts
@@ -7349,20 +7269,10 @@ export class TaskService {
extensionVersion: string,
metadata: Record<string, unknown> = {},
routing: { erpAccount?: string; erpAccountMatched?: boolean } = {}
): Promise<{
execution_ready: boolean;
erp_account_matched: boolean;
worker_connection_id: string;
extension_update_target: {
ecs_region_id: string | null;
ecs_instance_id: string | null;
server_update_safe: boolean;
};
}> {
): Promise<{ execution_ready: boolean; erp_account_matched: boolean; worker_connection_id: string }> {
return withTransaction(this.config, async (client) => {
const account = await client.query(
`SELECT id, role, is_active, erp_account,
extension_ecs_region_id, extension_ecs_instance_id
`SELECT id, role, is_active, erp_account
FROM users
WHERE organization_id = $1 AND id = $2
FOR UPDATE`,
@@ -7448,59 +7358,10 @@ export class TaskService {
if (!result.rowCount) {
throw new TaskError('browser_connection_not_owned', '浏览器连接已绑定其他账号,请重新生成连接标识。', 403);
}
const ecsRegionId = text(accountRow.extension_ecs_region_id) || null;
const ecsInstanceId = text(accountRow.extension_ecs_instance_id) || null;
let serverUpdateSafe = false;
if (ecsRegionId && ecsInstanceId) {
const unsafe = await client.query(
`SELECT (
EXISTS (
SELECT 1
FROM users host_account
JOIN tasks task
ON task.organization_id = host_account.organization_id
AND task.assigned_user_id = host_account.id
WHERE host_account.organization_id = $1
AND host_account.extension_ecs_region_id = $2
AND host_account.extension_ecs_instance_id = $3
AND (
task.status = 'reconciliation_pending'
OR (task.lease_expires_at IS NOT NULL AND task.lease_expires_at > now())
OR EXISTS (
SELECT 1 FROM task_attempts attempt
WHERE attempt.task_id = task.id
AND attempt.phase = 'erp'
AND attempt.status IN ('accepted', 'running', 'reconciliation_pending')
)
)
)
OR EXISTS (
SELECT 1
FROM users host_account
JOIN browser_connections connection
ON connection.organization_id = host_account.organization_id
AND connection.user_id = host_account.id
WHERE host_account.organization_id = $1
AND host_account.extension_ecs_region_id = $2
AND host_account.extension_ecs_instance_id = $3
AND connection.status = 'connected'
AND connection.last_seen_at >= now() - interval '90 seconds'
AND lower(COALESCE(connection.metadata->>'extension_update_safe', 'false')) <> 'true'
)
) AS unsafe`,
[context.organizationId, ecsRegionId, ecsInstanceId]
);
serverUpdateSafe = !databaseBoolean(unsafe.rows[0]?.unsafe);
}
return {
execution_ready: executionReady,
erp_account_matched: erpAccountMatched,
worker_connection_id: connectionId,
extension_update_target: {
ecs_region_id: ecsRegionId,
ecs_instance_id: ecsInstanceId,
server_update_safe: serverUpdateSafe
}
worker_connection_id: connectionId
};
});
}