feat: add ProcessDock process and port manager
This commit is contained in:
499
processdock/services.py
Normal file
499
processdock/services.py
Normal file
@@ -0,0 +1,499 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import json
|
||||
import socket
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from collections import defaultdict
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import psutil
|
||||
|
||||
|
||||
STATUS_LABELS = {
|
||||
# psutil exposes different status constants on each operating system, while
|
||||
# the returned values are stable strings. Literal keys keep imports portable.
|
||||
"running": "运行中",
|
||||
"sleeping": "休眠",
|
||||
"disk-sleep": "等待磁盘",
|
||||
"stopped": "已停止",
|
||||
"tracing-stop": "调试暂停",
|
||||
"zombie": "僵尸进程",
|
||||
"dead": "已结束",
|
||||
"wake-kill": "唤醒终止",
|
||||
"waking": "唤醒中",
|
||||
"parked": "已驻留",
|
||||
"idle": "空闲",
|
||||
"locked": "已锁定",
|
||||
"waiting": "等待中",
|
||||
"suspended": "已挂起",
|
||||
}
|
||||
|
||||
|
||||
class ProcessDockError(RuntimeError):
|
||||
"""Base error with a user-facing message."""
|
||||
|
||||
|
||||
class ProcessGoneError(ProcessDockError):
|
||||
pass
|
||||
|
||||
|
||||
class ProcessAccessError(ProcessDockError):
|
||||
pass
|
||||
|
||||
|
||||
class ProtectedProcessError(ProcessDockError):
|
||||
pass
|
||||
|
||||
|
||||
class ProcessChangedError(ProcessDockError):
|
||||
pass
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ActionResult:
|
||||
success: bool
|
||||
message: str
|
||||
still_running: bool = False
|
||||
|
||||
|
||||
def format_command(arguments: list[str] | None) -> str:
|
||||
if not arguments:
|
||||
return ""
|
||||
if os.name == "nt":
|
||||
return subprocess.list2cmdline(arguments)
|
||||
return " ".join(_shell_quote(argument) for argument in arguments)
|
||||
|
||||
|
||||
def _shell_quote(value: str) -> str:
|
||||
if not value:
|
||||
return "''"
|
||||
if all(character.isalnum() or character in "@%_-+=:,./" for character in value):
|
||||
return value
|
||||
return "'" + value.replace("'", "'\"'\"'") + "'"
|
||||
|
||||
|
||||
def humanize_bytes(size: int | float) -> str:
|
||||
value = float(max(size, 0))
|
||||
for unit in ("B", "KB", "MB", "GB", "TB"):
|
||||
if value < 1024 or unit == "TB":
|
||||
if unit == "B":
|
||||
return f"{value:.0f} {unit}"
|
||||
return f"{value:.1f} {unit}"
|
||||
value /= 1024
|
||||
return f"{value:.1f} TB"
|
||||
|
||||
|
||||
def matches_process(process: dict[str, Any], query: str) -> bool:
|
||||
normalized = query.strip().casefold()
|
||||
if not normalized:
|
||||
return True
|
||||
values = [
|
||||
str(process.get("pid", "")),
|
||||
str(process.get("name", "")),
|
||||
str(process.get("username", "")),
|
||||
str(process.get("exe", "")),
|
||||
str(process.get("command", "")),
|
||||
",".join(str(port) for port in process.get("ports", [])),
|
||||
]
|
||||
return any(normalized in value.casefold() for value in values)
|
||||
|
||||
|
||||
class ProcessService:
|
||||
"""Reads process/port data and performs guarded local process actions."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.current_pid = os.getpid()
|
||||
self._cpu_count = max(psutil.cpu_count(logical=True) or 1, 1)
|
||||
self._cpu_samples: dict[int, tuple[float, float, float]] = {}
|
||||
self._last_cpu: dict[int, float] = {}
|
||||
self._lock = threading.RLock()
|
||||
|
||||
def snapshot(self) -> dict[str, Any]:
|
||||
port_map, port_warning = self._collect_ports()
|
||||
warnings = [port_warning] if port_warning else []
|
||||
if os.name == "nt":
|
||||
try:
|
||||
processes = self._collect_windows_processes(port_map)
|
||||
except (OSError, subprocess.SubprocessError, json.JSONDecodeError) as error:
|
||||
warnings.append(f"快速采集不可用,已切换兼容模式:{error}")
|
||||
processes = self._collect_psutil_processes(port_map)
|
||||
else:
|
||||
processes = self._collect_psutil_processes(port_map)
|
||||
|
||||
processes.sort(key=lambda item: (item["name"].casefold(), item["pid"]))
|
||||
unique_ports = {
|
||||
(detail["protocol"], detail["address"], detail["port"])
|
||||
for process in processes
|
||||
for detail in process["port_details"]
|
||||
}
|
||||
memory = psutil.virtual_memory()
|
||||
summary = {
|
||||
"process_count": len(processes),
|
||||
"port_count": len(unique_ports),
|
||||
"cpu_percent": round(float(psutil.cpu_percent(interval=None)), 1),
|
||||
"memory_percent": round(float(memory.percent), 1),
|
||||
"memory_used": humanize_bytes(memory.used),
|
||||
"memory_total": humanize_bytes(memory.total),
|
||||
}
|
||||
return {
|
||||
"processes": processes,
|
||||
"summary": summary,
|
||||
"warnings": warnings,
|
||||
"refreshed_at": datetime.now().strftime("%H:%M:%S"),
|
||||
}
|
||||
|
||||
def _collect_psutil_processes(
|
||||
self, port_map: dict[int, list[dict[str, Any]]]
|
||||
) -> list[dict[str, Any]]:
|
||||
now = time.monotonic()
|
||||
processes: list[dict[str, Any]] = []
|
||||
live_pids: set[int] = set()
|
||||
for process in psutil.process_iter():
|
||||
try:
|
||||
row = self._read_process(process, port_map.get(process.pid, []), now)
|
||||
except (psutil.NoSuchProcess, psutil.ZombieProcess):
|
||||
continue
|
||||
except Exception:
|
||||
# A single unusual system process should never break the full refresh.
|
||||
continue
|
||||
processes.append(row)
|
||||
live_pids.add(row["pid"])
|
||||
|
||||
with self._lock:
|
||||
self._cpu_samples = {
|
||||
pid: sample for pid, sample in self._cpu_samples.items() if pid in live_pids
|
||||
}
|
||||
self._last_cpu = {
|
||||
pid: value for pid, value in self._last_cpu.items() if pid in live_pids
|
||||
}
|
||||
return processes
|
||||
|
||||
def _collect_windows_processes(
|
||||
self, port_map: dict[int, list[dict[str, Any]]]
|
||||
) -> list[dict[str, Any]]:
|
||||
"""Use Windows CIM's bulk providers to avoid hundreds of slow per-PID calls."""
|
||||
script = r"""
|
||||
[Console]::OutputEncoding = [System.Text.Encoding]::UTF8
|
||||
$processes = Get-CimInstance Win32_Process -ErrorAction Stop | ForEach-Object {
|
||||
$created = 0
|
||||
if ($_.CreationDate) {
|
||||
$created = [DateTimeOffset]$_.CreationDate
|
||||
$created = $created.ToUnixTimeMilliseconds() / 1000
|
||||
}
|
||||
[PSCustomObject]@{
|
||||
pid = [int]$_.ProcessId
|
||||
ppid = [int]$_.ParentProcessId
|
||||
name = [string]$_.Name
|
||||
exe = [string]$_.ExecutablePath
|
||||
command = [string]$_.CommandLine
|
||||
created = [double]$created
|
||||
}
|
||||
}
|
||||
$performance = Get-CimInstance Win32_PerfFormattedData_PerfProc_Process -ErrorAction Stop |
|
||||
Where-Object { $_.Name -ne '_Total' } |
|
||||
ForEach-Object {
|
||||
[PSCustomObject]@{
|
||||
pid = [int]$_.IDProcess
|
||||
cpu = [double]$_.PercentProcessorTime
|
||||
memory = [int64]$_.WorkingSetPrivate
|
||||
}
|
||||
}
|
||||
[PSCustomObject]@{ processes = @($processes); performance = @($performance) } |
|
||||
ConvertTo-Json -Depth 4 -Compress
|
||||
"""
|
||||
completed = subprocess.run(
|
||||
["powershell.exe", "-NoProfile", "-NonInteractive", "-Command", script],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
encoding="utf-8",
|
||||
errors="replace",
|
||||
timeout=8,
|
||||
creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0),
|
||||
)
|
||||
if completed.returncode != 0:
|
||||
detail = (completed.stderr or completed.stdout).strip()
|
||||
raise OSError(detail or "PowerShell CIM 查询失败")
|
||||
payload = json.loads(completed.stdout.lstrip("\ufeff"))
|
||||
static_by_pid = {
|
||||
int(item["pid"]): item
|
||||
for item in self._ensure_list(payload.get("processes"))
|
||||
if item.get("pid") is not None
|
||||
}
|
||||
performance_by_pid = {
|
||||
int(item["pid"]): item
|
||||
for item in self._ensure_list(payload.get("performance"))
|
||||
if item.get("pid") is not None
|
||||
}
|
||||
total_memory = max(int(psutil.virtual_memory().total), 1)
|
||||
processes: list[dict[str, Any]] = []
|
||||
|
||||
for pid, values in static_by_pid.items():
|
||||
if pid == 0:
|
||||
continue
|
||||
performance = performance_by_pid.get(pid, {})
|
||||
memory_bytes = max(int(performance.get("memory") or 0), 0)
|
||||
raw_cpu = max(float(performance.get("cpu") or 0.0), 0.0)
|
||||
cpu_percent = round(min(raw_cpu / self._cpu_count, 100.0), 1)
|
||||
try:
|
||||
username = psutil.Process(pid).username() or "—"
|
||||
except (psutil.Error, OSError):
|
||||
username = "—"
|
||||
port_details = port_map.get(pid, [])
|
||||
ports = sorted({detail["port"] for detail in port_details})
|
||||
executable = values.get("exe") or ""
|
||||
command = values.get("command") or ""
|
||||
create_time = float(values.get("created") or 0.0)
|
||||
processes.append(
|
||||
{
|
||||
"pid": pid,
|
||||
"name": values.get("name") or f"PID {pid}",
|
||||
"ports": ports,
|
||||
"port_details": port_details,
|
||||
"cpu_percent": cpu_percent,
|
||||
"memory_bytes": memory_bytes,
|
||||
"memory_text": humanize_bytes(memory_bytes),
|
||||
"memory_percent": round(memory_bytes / total_memory * 100, 1),
|
||||
"username": username,
|
||||
"status": "running",
|
||||
"status_label": "运行中",
|
||||
"exe": executable,
|
||||
"command": command,
|
||||
"command_parts": [],
|
||||
"create_time": create_time,
|
||||
"create_time_text": self._format_timestamp(create_time),
|
||||
"ppid": int(values.get("ppid") or 0),
|
||||
"access_limited": not executable and not command and username == "—",
|
||||
"is_current": pid == self.current_pid,
|
||||
"can_terminate": pid != self.current_pid,
|
||||
"can_open_location": bool(executable),
|
||||
}
|
||||
)
|
||||
return processes
|
||||
|
||||
@staticmethod
|
||||
def _ensure_list(value: Any) -> list[Any]:
|
||||
if value is None:
|
||||
return []
|
||||
return value if isinstance(value, list) else [value]
|
||||
|
||||
def _collect_ports(self) -> tuple[dict[int, list[dict[str, Any]]], str | None]:
|
||||
ports: dict[int, set[tuple[str, str, int]]] = defaultdict(set)
|
||||
try:
|
||||
connections = psutil.net_connections(kind="inet")
|
||||
except psutil.AccessDenied:
|
||||
return {}, "权限不足:无法读取完整端口列表,建议以管理员身份运行。"
|
||||
except OSError as error:
|
||||
return {}, f"端口信息暂时不可用:{error}"
|
||||
|
||||
for connection in connections:
|
||||
if connection.pid is None or not connection.laddr:
|
||||
continue
|
||||
protocol = "TCP" if connection.type == socket.SOCK_STREAM else "UDP"
|
||||
if protocol == "TCP" and connection.status != psutil.CONN_LISTEN:
|
||||
continue
|
||||
try:
|
||||
address = connection.laddr.ip
|
||||
port = int(connection.laddr.port)
|
||||
except AttributeError:
|
||||
address, port = connection.laddr[:2]
|
||||
port = int(port)
|
||||
if port > 0:
|
||||
ports[connection.pid].add((protocol, str(address), port))
|
||||
|
||||
result = {
|
||||
pid: [
|
||||
{"protocol": protocol, "address": address, "port": port}
|
||||
for protocol, address, port in sorted(
|
||||
items, key=lambda item: (item[2], item[0], item[1])
|
||||
)
|
||||
]
|
||||
for pid, items in ports.items()
|
||||
}
|
||||
return result, None
|
||||
|
||||
def _read_process(
|
||||
self,
|
||||
process: psutil.Process,
|
||||
port_details: list[dict[str, Any]],
|
||||
sample_time: float,
|
||||
) -> dict[str, Any]:
|
||||
values = process.as_dict(
|
||||
attrs=[
|
||||
"pid",
|
||||
"name",
|
||||
"username",
|
||||
"status",
|
||||
"exe",
|
||||
"cmdline",
|
||||
"create_time",
|
||||
"ppid",
|
||||
"memory_info",
|
||||
"memory_percent",
|
||||
"cpu_times",
|
||||
],
|
||||
ad_value=None,
|
||||
)
|
||||
pid = int(values["pid"])
|
||||
create_time = float(values.get("create_time") or 0.0)
|
||||
cpu_percent = self._calculate_cpu(
|
||||
pid, values.get("cpu_times"), create_time, sample_time
|
||||
)
|
||||
memory_info = values.get("memory_info")
|
||||
memory_bytes = int(getattr(memory_info, "rss", 0) or 0)
|
||||
command_parts = values.get("cmdline") or []
|
||||
status = str(values.get("status") or "unknown")
|
||||
ports = sorted({detail["port"] for detail in port_details})
|
||||
exe = values.get("exe") or ""
|
||||
username = values.get("username") or "—"
|
||||
limited = not exe and not command_parts and username == "—"
|
||||
|
||||
return {
|
||||
"pid": pid,
|
||||
"name": values.get("name") or f"PID {pid}",
|
||||
"ports": ports,
|
||||
"port_details": port_details,
|
||||
"cpu_percent": cpu_percent,
|
||||
"memory_bytes": memory_bytes,
|
||||
"memory_text": humanize_bytes(memory_bytes),
|
||||
"memory_percent": round(float(values.get("memory_percent") or 0.0), 1),
|
||||
"username": username,
|
||||
"status": status,
|
||||
"status_label": STATUS_LABELS.get(status, "未知"),
|
||||
"exe": exe,
|
||||
"command": format_command(command_parts),
|
||||
"command_parts": command_parts,
|
||||
"create_time": create_time,
|
||||
"create_time_text": self._format_timestamp(create_time),
|
||||
"ppid": int(values.get("ppid") or 0),
|
||||
"access_limited": limited,
|
||||
"is_current": pid == self.current_pid,
|
||||
"can_terminate": pid != self.current_pid,
|
||||
"can_open_location": bool(exe),
|
||||
}
|
||||
|
||||
def _calculate_cpu(
|
||||
self,
|
||||
pid: int,
|
||||
cpu_times: Any,
|
||||
create_time: float,
|
||||
sample_time: float,
|
||||
) -> float:
|
||||
if cpu_times is None:
|
||||
return 0.0
|
||||
total = float(getattr(cpu_times, "user", 0.0)) + float(
|
||||
getattr(cpu_times, "system", 0.0)
|
||||
)
|
||||
with self._lock:
|
||||
previous = self._cpu_samples.get(pid)
|
||||
result = self._last_cpu.get(pid, 0.0)
|
||||
if previous is not None:
|
||||
previous_time, previous_total, previous_created = previous
|
||||
elapsed = sample_time - previous_time
|
||||
if elapsed >= 0.2 and abs(previous_created - create_time) < 0.01:
|
||||
used = max(total - previous_total, 0.0)
|
||||
result = min(max(used / elapsed / self._cpu_count * 100, 0.0), 100.0)
|
||||
result = round(result, 1)
|
||||
self._cpu_samples[pid] = (sample_time, total, create_time)
|
||||
self._last_cpu[pid] = result
|
||||
return result
|
||||
|
||||
@staticmethod
|
||||
def _format_timestamp(timestamp: float) -> str:
|
||||
if timestamp <= 0:
|
||||
return "—"
|
||||
try:
|
||||
return datetime.fromtimestamp(timestamp).strftime("%Y-%m-%d %H:%M:%S")
|
||||
except (OSError, OverflowError, ValueError):
|
||||
return "—"
|
||||
|
||||
def terminate(
|
||||
self, pid: int, *, force: bool, expected_create_time: float | None = None
|
||||
) -> ActionResult:
|
||||
if pid == self.current_pid:
|
||||
raise ProtectedProcessError("ProcessDock 正在使用该进程,不能结束自身。")
|
||||
try:
|
||||
process = psutil.Process(pid)
|
||||
actual_create_time = process.create_time()
|
||||
if (
|
||||
expected_create_time is not None
|
||||
and expected_create_time > 0
|
||||
and abs(actual_create_time - expected_create_time) >= 0.01
|
||||
):
|
||||
raise ProcessChangedError("PID 已被新的进程复用,请刷新后重试。")
|
||||
name = process.name()
|
||||
|
||||
if force:
|
||||
process.kill()
|
||||
elif os.name == "nt":
|
||||
completed = subprocess.run(
|
||||
["taskkill", "/PID", str(pid)],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=5,
|
||||
creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0),
|
||||
)
|
||||
# Console/background processes commonly reject taskkill without /F.
|
||||
# Treat that as "still running" below so the GUI can offer force,
|
||||
# rather than presenting a misleading hard failure.
|
||||
else:
|
||||
process.terminate()
|
||||
|
||||
try:
|
||||
process.wait(timeout=3)
|
||||
except psutil.TimeoutExpired:
|
||||
return ActionResult(
|
||||
False,
|
||||
f"{name} 尚未退出,可以尝试强制结束。",
|
||||
still_running=True,
|
||||
)
|
||||
action = "强制结束" if force else "正常结束"
|
||||
return ActionResult(True, f"已{action} {name}(PID {pid})。")
|
||||
except ProtectedProcessError:
|
||||
raise
|
||||
except ProcessChangedError:
|
||||
raise
|
||||
except psutil.NoSuchProcess as error:
|
||||
raise ProcessGoneError("进程已经结束。") from error
|
||||
except psutil.AccessDenied as error:
|
||||
raise ProcessAccessError("权限不足,请尝试以管理员身份运行 ProcessDock。") from error
|
||||
except subprocess.TimeoutExpired as error:
|
||||
raise ProcessDockError("系统结束命令超时,请稍后重试。") from error
|
||||
|
||||
def open_file_location(
|
||||
self, pid: int, *, expected_create_time: float | None = None
|
||||
) -> ActionResult:
|
||||
try:
|
||||
process = psutil.Process(pid)
|
||||
actual_create_time = process.create_time()
|
||||
if (
|
||||
expected_create_time is not None
|
||||
and expected_create_time > 0
|
||||
and abs(actual_create_time - expected_create_time) >= 0.01
|
||||
):
|
||||
raise ProcessChangedError("PID 已被新的进程复用,请刷新后重试。")
|
||||
executable = process.exe()
|
||||
if not executable:
|
||||
raise ProcessDockError("该进程没有可访问的程序路径。")
|
||||
path = Path(executable)
|
||||
if sys.platform == "win32":
|
||||
subprocess.Popen(["explorer.exe", f"/select,{path}"])
|
||||
elif sys.platform == "darwin":
|
||||
subprocess.Popen(["open", "-R", str(path)])
|
||||
else:
|
||||
subprocess.Popen(["xdg-open", str(path.parent)])
|
||||
return ActionResult(True, "已打开程序所在位置。")
|
||||
except ProcessChangedError:
|
||||
raise
|
||||
except psutil.NoSuchProcess as error:
|
||||
raise ProcessGoneError("进程已经结束。") from error
|
||||
except psutil.AccessDenied as error:
|
||||
raise ProcessAccessError("权限不足,无法读取程序路径。") from error
|
||||
except OSError as error:
|
||||
raise ProcessDockError(f"无法打开文件位置:{error}") from error
|
||||
Reference in New Issue
Block a user