# deploy_account.dspy - 应用部署主机逻辑账号系统 API # action: ensure / list / remove / run / file_read / file_write / file_list # 零 root:逻辑隔离目录 + bwrap 沙箱,不建真实 Linux 账号 user_id = await get_user() if not user_id: return json.dumps({'ok': False, 'error': '未登录'}, ensure_ascii=False) action = (params_kw or {}).get('action', 'ensure') from pipeline_service.deploy_account import ( ensure_account, remove_account, list_accounts, run_in_sandbox, write_file, read_file, list_dir, sanitize_username, ) from pipeline_service.work_env import run_in_work_env from app_audit.audit_service import audit_log client_ip = '' try: client_ip = request.get('client_ip', '') or '' except Exception: pass # 查当前用户 username(账号名 ag_) username = user_id try: async with get_sor_context(request._run_ns, 'pipeline') as _sor: _u = await _sor.sqlExe("SELECT username FROM users WHERE id=${u}$", {'u': user_id}) if _u: username = getattr(_u[0], 'username', user_id) or user_id except Exception: pass if action == 'ensure': # 确保账号存在(幂等,首次自动创建) try: async with get_sor_context(request._run_ns, 'pipeline') as sor: r = await ensure_account(sor, user_id, username) await audit_log(sor, user_id, username, 'account_create', target='ag_' + sanitize_username(username), detail='created=' + str(r.get('created', '?')), result='ok' if r.get('created') is not None else 'fail', client_ip=client_ip) return json.dumps({'ok': True, **r}, ensure_ascii=False) except Exception as e: return json.dumps({'ok': False, 'error': str(e)}, ensure_ascii=False) elif action == 'list': try: async with get_sor_context(request._run_ns, 'pipeline') as sor: rows = await list_accounts(sor, user_id) return json.dumps({'ok': True, 'accounts': rows}, ensure_ascii=False) except Exception as e: return json.dumps({'ok': False, 'error': str(e)}, ensure_ascii=False) elif action == 'remove': account_name = (params_kw or {}).get('account_name', '') try: async with get_sor_context(request._run_ns, 'pipeline') as sor: r = await remove_account(sor, account_name) await audit_log(sor, user_id, username, 'account_remove', target=account_name, detail=str(r.get('error', '')), result='ok' if r.get('ok') else 'fail', client_ip=client_ip) return json.dumps(r, ensure_ascii=False) except Exception as e: return json.dumps({'ok': False, 'error': str(e)}, ensure_ascii=False) elif action == 'run': # 沙箱内执行命令 account_name = (params_kw or {}).get('account_name', '') command = (params_kw or {}).get('command', '') workdir = (params_kw or {}).get('workdir', '') timeout = int((params_kw or {}).get('timeout', 120) or 120) network = (params_kw or {}).get('network', 'shared') if not account_name: account_name = 'ag_' + sanitize_username(username) try: async with get_sor_context(request._run_ns, 'pipeline') as sor: _u = await sor.sqlExe("SELECT orgid FROM users WHERE id=${u}$", {'u': user_id}) _org = getattr(_u[0], 'orgid', '') if _u else '' r = await run_in_work_env(sor, account_name, user_id, _org, command, workdir, timeout, network) await audit_log(sor, user_id, username, 'sandbox_run', target=account_name, detail=str(command)[:500], result='ok' if r.get('rc') == 0 else 'fail', client_ip=client_ip) return json.dumps({'ok': True, **r}, ensure_ascii=False) except Exception as e: return json.dumps({'ok': False, 'error': str(e)}, ensure_ascii=False) elif action == 'file_write': account_name = (params_kw or {}).get('account_name', '') or ('ag_' + sanitize_username(username)) relpath = (params_kw or {}).get('path', '') content = (params_kw or {}).get('content', '') try: async with get_sor_context(request._run_ns, 'pipeline') as sor: r = await write_file(sor, account_name, relpath, content) return json.dumps(r, ensure_ascii=False) except Exception as e: return json.dumps({'ok': False, 'error': str(e)}, ensure_ascii=False) elif action == 'file_read': account_name = (params_kw or {}).get('account_name', '') or ('ag_' + sanitize_username(username)) relpath = (params_kw or {}).get('path', '') try: async with get_sor_context(request._run_ns, 'pipeline') as sor: r = await read_file(sor, account_name, relpath) return json.dumps(r, ensure_ascii=False) except Exception as e: return json.dumps({'ok': False, 'error': str(e)}, ensure_ascii=False) elif action == 'file_list': account_name = (params_kw or {}).get('account_name', '') or ('ag_' + sanitize_username(username)) relpath = (params_kw or {}).get('path', '') try: async with get_sor_context(request._run_ns, 'pipeline') as sor: r = await list_dir(sor, account_name, relpath) return json.dumps(r, ensure_ascii=False) except Exception as e: return json.dumps({'ok': False, 'error': str(e)}, ensure_ascii=False) else: return json.dumps({'ok': False, 'error': '未知 action: ' + str(action)}, ensure_ascii=False)