753 lines
29 KiB
Python
753 lines
29 KiB
Python
import json
|
||
import shlex
|
||
import time
|
||
import hashlib
|
||
import hmac
|
||
import docker
|
||
from flask import Flask, jsonify, request
|
||
from flask_socketio import SocketIO
|
||
|
||
# 全局签名密钥,前后端统一,生产环境使用环境变量注入
|
||
SIGN_SECRET = b"DockerAdminPlatform2026SecretKey123456789"
|
||
# 签名有效期 300秒(5分钟)
|
||
SIGN_EXPIRE = 300
|
||
|
||
# 兼容新旧docker-py客户端
|
||
try:
|
||
cli = docker.DockerClient(base_url='unix:///var/run/docker.sock')
|
||
api_cli = docker.APIClient(base_url='unix:///var/run/docker.sock')
|
||
except AttributeError:
|
||
cli = docker.Client(base_url='unix:///var/run/docker.sock')
|
||
api_cli = docker.APIClient(base_url='unix:///var/run/docker.sock')
|
||
|
||
app = Flask(__name__)
|
||
app.config['SECRET_KEY'] = 'docker-terminal-secret-2026'
|
||
|
||
# WebSocket心跳保活,防止Nginx/反向代理断连
|
||
socketio = SocketIO(
|
||
app,
|
||
cors_allowed_origins="*",
|
||
async_mode="threading",
|
||
ping_timeout=60,
|
||
ping_interval=25
|
||
)
|
||
|
||
# 全局会话存储
|
||
session_exec_map = {} # sid => TTY终端流信息
|
||
log_task_map = {} # sid => {"cid": "", "running": True}
|
||
ws_user_map = {} # 新增:sid 映射 WebSocket鉴权用户信息
|
||
registry_config = {
|
||
"repo_name": "",
|
||
"repo_addr": "",
|
||
"repo_user": "",
|
||
"repo_pwd": ""
|
||
}
|
||
# 多系统shell兼容列表(Ubuntu/Alpine/标准Linux)
|
||
SHELL_CANDIDATES = [
|
||
"/bin/bash",
|
||
"/usr/bin/bash",
|
||
"/bin/sh",
|
||
"/bin/ash"
|
||
]
|
||
|
||
# ===================== 工具函数1:docker socket兼容recv =====================
|
||
def docker_recv(sock, size=1024):
|
||
try:
|
||
return sock.recv(size)
|
||
except Exception:
|
||
try:
|
||
return sock._sock.recv(size)
|
||
except Exception:
|
||
return b""
|
||
|
||
# ===================== 工具函数2:操作审计日志(终端输入脱敏,不存明文密码) =====================
|
||
def write_audit_log(user, role, action, target_name, target_id, result, client_ip):
|
||
"""企业审计记录,可对接数据库docker_audit_log表"""
|
||
# 终端输入仅标记事件,不记录明文命令,规避密码泄露风险
|
||
if action == "terminal_input":
|
||
result = "command_entered"
|
||
log_item = {
|
||
"timestamp": time.strftime("%Y-%m-%d %H:%M:%S"),
|
||
"user": user,
|
||
"role": role,
|
||
"action": action,
|
||
"container_name": target_name,
|
||
"container_id": target_id,
|
||
"result": result,
|
||
"client_ip": client_ip
|
||
}
|
||
# 生产环境替换为mysql insert语句
|
||
print(json.dumps(log_item, ensure_ascii=False))
|
||
|
||
# ===================== 工具函数3:HTTP Header HMAC签名鉴权(原有,HTTP接口专用) =====================
|
||
def verify_auth_sign():
|
||
"""
|
||
校验前端PHP生成的HMAC签名,防止伪造X-Role:admin
|
||
请求Header要求:
|
||
X-User: 用户名
|
||
X-Role: admin/user/guest
|
||
X-Time: 当前时间戳(秒)
|
||
X-Sign: hmac-sha256签名
|
||
返回 (user, role, client_ip, True/False)
|
||
"""
|
||
user = request.headers.get("X-User", "anonymous")
|
||
role = request.headers.get("X-Role", "guest")
|
||
ts_str = request.headers.get("X-Time", "0")
|
||
sign = request.headers.get("X-Sign", "")
|
||
client_ip = request.headers.get("X-Client-Ip", request.remote_addr)
|
||
try:
|
||
ts = int(ts_str)
|
||
except ValueError:
|
||
return user, role, client_ip, False
|
||
# 校验签名有效期
|
||
now = int(time.time())
|
||
if abs(now - ts) > SIGN_EXPIRE:
|
||
return user, role, client_ip, False
|
||
# 校验HMAC签名
|
||
raw_data = f"{user}|{role}|{ts}".encode("utf-8")
|
||
calc_sign = hmac.new(SIGN_SECRET, raw_data, hashlib.sha256).hexdigest()
|
||
if calc_sign != sign:
|
||
return user, role, client_ip, False
|
||
return user, role, client_ip, True
|
||
|
||
# ===================== 新增工具函数4:WebSocket连接专用鉴权函数 =====================
|
||
def verify_socket_auth(auth):
|
||
user = auth.get("user", "")
|
||
role = auth.get("role", "")
|
||
ts_str = auth.get("ts", "0")
|
||
sign = auth.get("sign", "")
|
||
try:
|
||
ts = int(ts_str)
|
||
except:
|
||
return user, role, "", False
|
||
if abs(int(time.time()) - ts) > SIGN_EXPIRE:
|
||
return user, role, "", False
|
||
raw_data = f"{user}|{role}|{ts}".encode("utf-8")
|
||
calc_sign = hmac.new(
|
||
SIGN_SECRET,
|
||
raw_data,
|
||
hashlib.sha256
|
||
).hexdigest()
|
||
if calc_sign != sign:
|
||
return user, role, "", False
|
||
return user, role, "", True
|
||
|
||
# ===================== 工具函数5:环境变量脱敏,屏蔽密码/密钥 =====================
|
||
def safe_env_filter(env_list):
|
||
safe_env = []
|
||
sensitive_keywords = ["PASS", "PASSWORD", "TOKEN", "SECRET", "KEY", "CREDENTIAL"]
|
||
for item in env_list:
|
||
equal_pos = item.find("=")
|
||
if equal_pos == -1:
|
||
safe_env.append(item)
|
||
continue
|
||
key = item[:equal_pos]
|
||
val = item[equal_pos+1:]
|
||
if any(word in key.upper() for word in sensitive_keywords):
|
||
safe_env.append(f"{key}=******")
|
||
else:
|
||
safe_env.append(item)
|
||
return safe_env
|
||
|
||
# ===================== 全部HTTP业务接口 =====================
|
||
# 1. 容器列表(端口预拼接字符串,前端直接渲染,无需formatPorts)
|
||
@app.route("/list")
|
||
def api_list():
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "list_container", "", "", "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
arr = []
|
||
for c in cli.containers.list(all=True):
|
||
img_name = c.image.tags[0] if c.image.tags else ""
|
||
port_str_arr = []
|
||
ports_raw = c.ports
|
||
for c_port, host_list in ports_raw.items():
|
||
if not host_list:
|
||
continue
|
||
for host_item in host_list:
|
||
host_port = host_item.get("HostPort", "")
|
||
port_str_arr.append(f"{host_port}:{c_port.split('/')[0]}")
|
||
port_display = ",".join(port_str_arr) if port_str_arr else "-"
|
||
arr.append({
|
||
"id": c.short_id,
|
||
"full_id": c.id,
|
||
"name": c.name,
|
||
"image": img_name,
|
||
"status": c.status,
|
||
"ports": port_display
|
||
})
|
||
write_audit_log(user, role, "list_container", "", "", "success", ip)
|
||
return jsonify({"code":0,"msg":"success","list":arr})
|
||
|
||
# 2. 镜像列表
|
||
@app.route("/image_list")
|
||
def api_image_list():
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "list_image", "", "", "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
arr = []
|
||
for img in cli.images.list():
|
||
if not img.tags:
|
||
continue
|
||
full_tag = img.tags[0]
|
||
repo, tag = full_tag.split(":", 1)
|
||
size_mb = round(img.attrs["Size"] / 1024 / 1024, 1)
|
||
arr.append({
|
||
"repo": repo,
|
||
"tag": tag,
|
||
"id": img.short_id,
|
||
"full_id": img.id,
|
||
"size": f"{size_mb}MB",
|
||
"created": img.attrs["Created"]
|
||
})
|
||
write_audit_log(user, role, "list_image", "", "", "success", ip)
|
||
return jsonify({"code":0,"msg":"success","list":arr})
|
||
|
||
# 3. 删除镜像
|
||
@app.route("/remove_image/<image_id>")
|
||
def api_remove_image(image_id):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "remove_image", "", image_id, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
if role != "admin":
|
||
write_audit_log(user, role, "remove_image", "", image_id, "denied", ip)
|
||
return jsonify({"code":403,"msg":"权限不足,仅管理员可删除镜像"})
|
||
try:
|
||
img = cli.images.get(image_id)
|
||
cli.images.remove(image_id, force=True)
|
||
write_audit_log(user, role, "remove_image", img.tags[0] if img.tags else "", image_id, "success", ip)
|
||
return jsonify({"code":0,"msg":"镜像已强制删除"})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "remove_image", "", image_id, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e)})
|
||
|
||
# 4. 重启容器
|
||
@app.route("/restart/<cid>")
|
||
def api_restart(cid):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "restart", "", cid, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
if role != "admin":
|
||
write_audit_log(user, role, "restart", "", cid, "denied", ip)
|
||
return jsonify({"code":403,"msg":"权限不足,仅管理员可操作容器"})
|
||
try:
|
||
container = cli.containers.get(cid)
|
||
container.restart()
|
||
write_audit_log(user, role, "restart", container.name, cid, "success", ip)
|
||
return jsonify({"code":0,"msg":"容器已重启"})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "restart", "", cid, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e)})
|
||
|
||
# 5. 停止容器
|
||
@app.route("/stop/<cid>")
|
||
def api_stop(cid):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "stop", "", cid, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
if role != "admin":
|
||
write_audit_log(user, role, "stop", "", cid, "denied", ip)
|
||
return jsonify({"code":403,"msg":"权限不足,仅管理员可操作容器"})
|
||
try:
|
||
container = cli.containers.get(cid)
|
||
container.stop()
|
||
write_audit_log(user, role, "stop", container.name, cid, "success", ip)
|
||
return jsonify({"code":0,"msg":"容器已停止"})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "stop", "", cid, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e)})
|
||
|
||
# 6. 删除容器
|
||
@app.route("/rm/<cid>")
|
||
def api_rm(cid):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "rm", "", cid, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
if role != "admin":
|
||
write_audit_log(user, role, "rm", "", cid, "denied", ip)
|
||
return jsonify({"code":403,"msg":"权限不足,仅管理员可删除容器"})
|
||
try:
|
||
container = cli.containers.get(cid)
|
||
container.remove(force=True)
|
||
write_audit_log(user, role, "rm", container.name, cid, "success", ip)
|
||
return jsonify({"code":0,"msg":"容器已强制删除"})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "rm", "", cid, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e)})
|
||
|
||
# 7. 一次性拉取日志(兼容旧页面单次查看)
|
||
@app.route("/logs/<cid>/<int:tail_lines>")
|
||
def api_logs(cid, tail_lines):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "read_log", "", cid, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
try:
|
||
container = cli.containers.get(cid)
|
||
log_bytes = container.logs(tail=tail_lines)
|
||
log_str = log_bytes.decode("utf-8", errors="replace")
|
||
write_audit_log(user, role, "read_log", container.name, cid, "success", ip)
|
||
return jsonify({"code":0,"msg":"success","data":log_str})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "read_log", "", cid, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e),"data":""})
|
||
|
||
# 8. 容器网络详情
|
||
@app.route("/network/<cid>")
|
||
def api_network(cid):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "read_network", "", cid, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
try:
|
||
container = cli.containers.get(cid)
|
||
info = container.attrs["NetworkSettings"]["Networks"]
|
||
write_audit_log(user, role, "read_network", container.name, cid, "success", ip)
|
||
return jsonify({"code":0,"msg":"success","data":json.dumps(info,indent=2)})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "read_network", "", cid, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e),"data":""})
|
||
|
||
# 9. 短连接单次执行命令(旧弹窗兼容,非实时终端)
|
||
@app.route("/exec/<cid>", methods=["POST"])
|
||
def api_exec(cid):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "single_exec", "", cid, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
if role != "admin":
|
||
write_audit_log(user, role, "single_exec", "", cid, "denied", ip)
|
||
return jsonify({"code":403,"msg":"权限不足,仅管理员可执行命令"})
|
||
try:
|
||
data = request.get_json()
|
||
cmd = data.get("cmd", "")
|
||
if not cmd:
|
||
return jsonify({"code":400,"msg":"命令不能为空","data":""})
|
||
cmd_list = shlex.split(cmd)
|
||
container = cli.containers.get(cid)
|
||
exec_obj = container.exec_run(cmd_list, stdout=True, stderr=True)
|
||
output = exec_obj.output.decode("utf-8", errors="replace")
|
||
write_audit_log(user, role, "single_exec", container.name, cid, "command_executed", ip)
|
||
return jsonify({"code":0,"msg":"success","data":output})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "single_exec", "", cid, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e),"data":""})
|
||
|
||
# 10. 创建容器(端口标准格式 + 前端可配置CPU/内存)
|
||
@app.route("/create", methods=["POST"])
|
||
def api_create():
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "create_container", "", "", "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
if role != "admin":
|
||
write_audit_log(user, role, "create_container", "", "", "denied", ip)
|
||
return jsonify({"code":403,"msg":"权限不足,仅管理员可创建容器"})
|
||
try:
|
||
data = request.get_json()
|
||
name = data.get("name", "")
|
||
image = data.get("image")
|
||
cmd = data.get("cmd")
|
||
ports_input = data.get("ports", {})
|
||
# 修复:docker标准端口字典格式 {内部端口:宿主机端口}
|
||
port_bind = {}
|
||
for inner_port, host_port in ports_input.items():
|
||
port_bind[f"{inner_port}/tcp"] = int(host_port)
|
||
# 前端可配置资源限制,默认1核512M
|
||
cpu_core = int(data.get("cpu", 1))
|
||
memory_limit = data.get("memory", "512m")
|
||
nano_cpus = cpu_core * 1000000000
|
||
container = cli.containers.run(
|
||
image=image,
|
||
name=name if name else None,
|
||
command=cmd if cmd else None,
|
||
ports=port_bind,
|
||
detach=True,
|
||
restart_policy={"Name": "always"},
|
||
mem_limit=memory_limit,
|
||
nano_cpus=nano_cpus
|
||
)
|
||
write_audit_log(user, role, "create_container", container.name, container.id, "success", ip)
|
||
return jsonify({"code":0,"msg":f"容器创建成功,ID:{container.short_id}"})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "create_container", "", "", f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e)})
|
||
|
||
# 11. 登录私有镜像仓库
|
||
@app.route("/bind_repo", methods=["POST"])
|
||
def api_bind_repo():
|
||
global registry_config
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "repo_login", "", "", "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
if role != "admin":
|
||
write_audit_log(user, role, "repo_login", "", "", "denied", ip)
|
||
return jsonify({"code":403,"msg":"权限不足,仅管理员可登录镜像仓库"})
|
||
try:
|
||
data = request.get_json()
|
||
repo_name = data.get("repo_name","")
|
||
repo_addr = data.get("repo_addr")
|
||
repo_user = data.get("repo_user","")
|
||
repo_pwd = data.get("repo_pwd","")
|
||
api_cli.login(username=repo_user, password=repo_pwd, registry=repo_addr)
|
||
registry_config["repo_name"] = repo_name
|
||
registry_config["repo_addr"] = repo_addr
|
||
registry_config["repo_user"] = repo_user
|
||
registry_config["repo_pwd"] = repo_pwd
|
||
write_audit_log(user, role, "repo_login", repo_addr, "", "success", ip)
|
||
return jsonify({"code":0,"msg":"镜像仓库登录绑定成功"})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "repo_login", repo_addr, "", f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e)})
|
||
|
||
# 修复:stats接口online_cpus兜底兼容,无KeyError
|
||
@app.route("/stats/<cid>")
|
||
def api_stats(cid):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "read_stats", "", cid, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
try:
|
||
container = cli.containers.get(cid)
|
||
raw = container.stats(stream=False)
|
||
# CPU兼容兜底,不存在online_cpus自动计算
|
||
cpu_stats = raw["cpu_stats"]
|
||
online_cpus = cpu_stats.get("online_cpus", len(cpu_stats["cpu_usage"].get("percpu_usage", [])) or 1)
|
||
cpu_delta = cpu_stats["cpu_usage"]["total_usage"] - raw["precpu_stats"]["cpu_usage"]["total_usage"]
|
||
system_delta = cpu_stats["system_cpu_usage"] - raw["precpu_stats"]["system_cpu_usage"]
|
||
cpu_percent = 0.0
|
||
if system_delta > 0 and cpu_delta > 0:
|
||
cpu_percent = round((cpu_delta / system_delta) * online_cpus * 100, 1)
|
||
# 内存计算
|
||
mem_usage = raw["memory_stats"]["usage"] / 1024 / 1024
|
||
mem_limit = raw["memory_stats"]["limit"] / 1024 / 1024
|
||
mem_percent = round((mem_usage / mem_limit) * 100, 1) if mem_limit > 0 else 0
|
||
# 网络流量
|
||
net_rx = 0
|
||
net_tx = 0
|
||
for net in raw["networks"].values():
|
||
net_rx += net["rx_bytes"] / 1024 / 1024
|
||
net_tx += net["tx_bytes"] / 1024 / 1024
|
||
res = {
|
||
"cpu": cpu_percent,
|
||
"mem_mb": round(mem_usage, 1),
|
||
"mem_limit_mb": round(mem_limit, 1),
|
||
"mem_percent": mem_percent,
|
||
"net_rx_mb": round(net_rx, 1),
|
||
"net_tx_mb": round(net_tx, 1)
|
||
}
|
||
write_audit_log(user, role, "read_stats", container.name, cid, "success", ip)
|
||
return jsonify({"code":0,"msg":"success","data":res})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "read_stats", "", cid, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e),"data":{}})
|
||
|
||
# 修复:detail接口ENV脱敏,多网卡IP兼容
|
||
@app.route("/detail/<cid>")
|
||
def api_detail(cid):
|
||
user, role, ip, auth_ok = verify_auth_sign()
|
||
if not auth_ok:
|
||
write_audit_log(user, role, "read_detail", "", cid, "auth_failed", ip)
|
||
return jsonify({"code":403,"msg":"身份校验失败,非法访问"})
|
||
try:
|
||
container = cli.containers.get(cid)
|
||
attrs = container.attrs
|
||
nets = attrs["NetworkSettings"]["Networks"]
|
||
ip_map = {}
|
||
for net_name, net_info in nets.items():
|
||
ip_map[net_name] = net_info.get("IPAddress", "")
|
||
# 敏感环境变量脱敏
|
||
raw_env = attrs["Config"].get("Env", [])
|
||
safe_env = safe_env_filter(raw_env)
|
||
res = {
|
||
"network_ips": ip_map,
|
||
"mounts": attrs["Mounts"],
|
||
"env": safe_env,
|
||
"created": attrs["Created"],
|
||
"entrypoint": attrs["Config"]["Entrypoint"],
|
||
"cmd": attrs["Config"]["Cmd"],
|
||
"restart_policy": attrs["HostConfig"]["RestartPolicy"]
|
||
}
|
||
write_audit_log(user, role, "read_detail", container.name, cid, "success", ip)
|
||
return jsonify({"code":0,"msg":"success","data":res})
|
||
except Exception as e:
|
||
write_audit_log(user, role, "read_detail", "", cid, f"fail:{str(e)}", ip)
|
||
return jsonify({"code":500,"msg":str(e),"data":{}})
|
||
|
||
# 根路由
|
||
@app.route("/")
|
||
def index():
|
||
return jsonify({"code":400,"msg":"无效操作","list":[],"data":""})
|
||
|
||
# ===================== WebSocket 连接鉴权事件(新增) =====================
|
||
@socketio.on("connect")
|
||
def handle_connect(auth):
|
||
user, role, _, auth_ok = verify_socket_auth(auth)
|
||
if not auth_ok:
|
||
return False
|
||
ws_user_map[request.sid] = {
|
||
"user": user,
|
||
"role": role
|
||
}
|
||
print("WS Connected:", user, role)
|
||
|
||
# ===================== WebSocket 交互式TTY实时终端 =====================
|
||
def stream_terminal_output(sid, sock):
|
||
"""后台异步读取docker终端输出,兼容两种socket接收方式"""
|
||
try:
|
||
while True:
|
||
chunk = docker_recv(sock, 1024)
|
||
if not chunk:
|
||
break
|
||
text = chunk.decode("utf-8", errors="replace")
|
||
socketio.emit("terminal_output", {"data": text}, room=sid)
|
||
except Exception:
|
||
pass
|
||
# 流结束自动清理会话
|
||
if sid in session_exec_map:
|
||
try:
|
||
session_exec_map[sid]["socket"].close()
|
||
except Exception:
|
||
pass
|
||
del session_exec_map[sid]
|
||
|
||
@socketio.on("terminal_start")
|
||
def handle_terminal_start(json_data):
|
||
sid = request.sid
|
||
# 新增调试输出
|
||
socketio.emit(
|
||
"terminal_output",
|
||
{
|
||
"data":"[DEBUG] terminal_start reached\r\n"
|
||
},
|
||
room=sid
|
||
)
|
||
# 替换原verify_auth_sign(),从ws_user_map读取身份
|
||
info = ws_user_map.get(request.sid)
|
||
if not info:
|
||
socketio.emit(
|
||
"terminal_output",
|
||
{"data":"身份认证失效\n"},
|
||
room=request.sid
|
||
)
|
||
return
|
||
user = info["user"]
|
||
role = info["role"]
|
||
client_ip = request.remote_addr
|
||
|
||
# 权限拦截:非管理员禁止打开交互式终端
|
||
if role != "admin":
|
||
socketio.emit("terminal_output", {"data": "Permission denied: 仅管理员可访问交互式终端\n"}, room=sid)
|
||
write_audit_log(user, role, "terminal_start", "", json_data.get("cid"), "denied", client_ip)
|
||
return
|
||
|
||
cid = json_data.get("cid")
|
||
if not cid:
|
||
socketio.emit("terminal_output", {"data": "错误:容器ID不能为空\n"}, room=sid)
|
||
write_audit_log(user, role, "terminal_start", "", "", "empty_cid", client_ip)
|
||
return
|
||
|
||
# 自动探测容器可用shell,兼容全版本docker-py(exec_create返回dict取Id)
|
||
target_shell = None
|
||
for shell in SHELL_CANDIDATES:
|
||
try:
|
||
ret = api_cli.exec_create(
|
||
container=cid,
|
||
cmd=[shell, "-c", "echo ok"]
|
||
)
|
||
test_exec_id = ret["Id"]
|
||
test_stream = api_cli.exec_start(test_exec_id)
|
||
test_stream.read()
|
||
target_shell = shell
|
||
break
|
||
except Exception:
|
||
continue
|
||
if not target_shell:
|
||
socketio.emit("terminal_output", {"data": "错误:容器内无可用shell(bash/sh/ash)\n"}, room=sid)
|
||
write_audit_log(user, role, "terminal_start", "", cid, "no_shell", client_ip)
|
||
return
|
||
|
||
try:
|
||
# 创建交互式TTY会话,统一提取Id兼容所有SDK版本
|
||
exec_ret = api_cli.exec_create(
|
||
container=cid,
|
||
cmd=target_shell,
|
||
stdin=True,
|
||
stdout=True,
|
||
stderr=True,
|
||
tty=True
|
||
)
|
||
exec_id = exec_ret["Id"]
|
||
sock = api_cli.exec_start(exec_id=exec_id, tty=True, socket=True, detach=False)
|
||
session_exec_map[sid] = {
|
||
"cid": cid,
|
||
"exec_id": exec_id,
|
||
"socket": sock
|
||
}
|
||
socketio.start_background_task(stream_terminal_output, sid, sock)
|
||
container_name = cli.containers.get(cid).name
|
||
write_audit_log(user, role, "terminal_start", container_name, cid, "success", client_ip)
|
||
except Exception as e:
|
||
socketio.emit("terminal_output", {"data": f"启动终端失败:{str(e)}\n"}, room=sid)
|
||
write_audit_log(user, role, "terminal_start", "", cid, f"fail:{str(e)}", client_ip)
|
||
|
||
@socketio.on("terminal_input")
|
||
def handle_terminal_input(json_data):
|
||
"""前端输入字符写入容器标准输入,审计脱敏不记录明文"""
|
||
sid = request.sid
|
||
# 替换原verify_auth_sign()
|
||
info = ws_user_map.get(request.sid)
|
||
if not info:
|
||
return
|
||
user = info["user"]
|
||
role = info["role"]
|
||
client_ip = request.remote_addr
|
||
|
||
if sid not in session_exec_map:
|
||
socketio.emit("terminal_output", {"data": "终端会话已断开,请重新打开\n"}, room=sid)
|
||
return
|
||
input_text = json_data.get("data", "")
|
||
sock = session_exec_map[sid]["socket"]
|
||
sock.send(input_text.encode("utf-8"))
|
||
cid = session_exec_map[sid]["cid"]
|
||
# 审计仅标记事件,不再记录明文命令,规避密码泄露
|
||
write_audit_log(user, role, "terminal_input", "", cid, "command_entered", client_ip)
|
||
|
||
@socketio.on("terminal_resize")
|
||
def handle_terminal_resize(json_data):
|
||
"""前端xterm窗口大小自适应"""
|
||
sid = request.sid
|
||
if sid not in session_exec_map:
|
||
return
|
||
exec_id = session_exec_map[sid]["exec_id"]
|
||
height = int(json_data.get("height", 24))
|
||
width = int(json_data.get("width", 80))
|
||
try:
|
||
api_cli.exec_resize(exec_id=exec_id, height=height, width=width)
|
||
except Exception:
|
||
pass
|
||
|
||
# 新增主动关闭终端事件(解决切换容器残留TTY)
|
||
@socketio.on("terminal_close")
|
||
def handle_terminal_close():
|
||
sid = request.sid
|
||
info = ws_user_map.get(request.sid)
|
||
if not info:
|
||
return
|
||
user = info["user"]
|
||
role = info["role"]
|
||
client_ip = request.remote_addr
|
||
|
||
if sid in session_exec_map:
|
||
cid = session_exec_map[sid]["cid"]
|
||
try:
|
||
session_exec_map[sid]["socket"].close()
|
||
except Exception:
|
||
pass
|
||
del session_exec_map[sid]
|
||
write_audit_log(user, role, "terminal_close", "", cid, "manual_close", client_ip)
|
||
|
||
@socketio.on("disconnect")
|
||
def handle_disconnect():
|
||
"""前端断开连接,释放docker socket、清理会话,防止句柄泄漏"""
|
||
sid = request.sid
|
||
# 清理ws_user_map会话
|
||
if sid in ws_user_map:
|
||
del ws_user_map[sid]
|
||
info = ws_user_map.get(sid)
|
||
if not info:
|
||
user = ""
|
||
role = ""
|
||
client_ip = request.remote_addr
|
||
else:
|
||
user = info["user"]
|
||
role = info["role"]
|
||
client_ip = request.remote_addr
|
||
|
||
# 清理TTY终端
|
||
if sid in session_exec_map:
|
||
cid = session_exec_map[sid]["cid"]
|
||
try:
|
||
session_exec_map[sid]["socket"].close()
|
||
except Exception:
|
||
pass
|
||
del session_exec_map[sid]
|
||
write_audit_log(user, role, "terminal_close", "", cid, "disconnect", client_ip)
|
||
# 清理实时日志标记,终止日志循环
|
||
if sid in log_task_map:
|
||
del log_task_map[sid]
|
||
|
||
# ===================== WebSocket 实时流式日志 docker logs -f(生成器主动释放) =====================
|
||
def stream_follow_log(sid, cid):
|
||
generator = None
|
||
try:
|
||
container = cli.containers.get(cid)
|
||
generator = container.logs(stream=True, follow=True)
|
||
log_task_map[sid] = {"cid": cid, "running": True}
|
||
for line in generator:
|
||
task_info = log_task_map.get(sid, {})
|
||
if not task_info.get("running"):
|
||
break
|
||
txt = line.decode("utf-8", errors="replace")
|
||
socketio.emit("log_output", {"data": txt}, room=sid)
|
||
except Exception:
|
||
pass
|
||
finally:
|
||
# 主动关闭日志生成器,释放docker daemon长连接
|
||
if generator is not None:
|
||
generator.close()
|
||
if sid in log_task_map:
|
||
del log_task_map[sid]
|
||
|
||
@socketio.on("log_start")
|
||
def handle_log_start(json_data):
|
||
sid = request.sid
|
||
# 替换verify_auth_sign,读取ws_user_map
|
||
info = ws_user_map.get(request.sid)
|
||
if not info:
|
||
socketio.emit("log_output", {"data": "身份校验失败"}, room=sid)
|
||
write_audit_log("", "", "log_start", "", "", "auth_failed", request.remote_addr)
|
||
return
|
||
user = info["user"]
|
||
role = info["role"]
|
||
client_ip = request.remote_addr
|
||
|
||
cid = json_data.get("cid")
|
||
if not cid:
|
||
socketio.emit("log_output", {"data": "容器ID不能为空"}, room=sid)
|
||
write_audit_log(user, role, "log_start", "", "", "empty_cid", client_ip)
|
||
return
|
||
if role != "admin":
|
||
socketio.emit("log_output", {"data": "权限不足,仅管理员可查看实时日志"}, room=sid)
|
||
write_audit_log(user, role, "log_start", "", cid, "denied", client_ip)
|
||
return
|
||
socketio.start_background_task(stream_follow_log, sid, cid)
|
||
container_name = cli.containers.get(cid).name
|
||
write_audit_log(user, role, "log_start", container_name, cid, "success", client_ip)
|
||
|
||
@socketio.on("log_stop")
|
||
def handle_log_stop():
|
||
sid = request.sid
|
||
info = ws_user_map.get(request.sid)
|
||
if not info:
|
||
return
|
||
user = info["user"]
|
||
role = info["role"]
|
||
client_ip = request.remote_addr
|
||
|
||
if sid in log_task_map:
|
||
log_task_map[sid]["running"] = False
|
||
cid = log_task_map[sid]["cid"]
|
||
write_audit_log(user, role, "log_stop", "", cid, "manual_close", client_ip)
|
||
|
||
# 启动入口(eventlet异步,必须socketio.run)
|
||
if __name__ == "__main__":
|
||
socketio.run(app, host="0.0.0.0", port=8090, debug=False, allow_unsafe_werkzeug=True)
|