Repository navigation
Expand file tree
/
Copy pathrun.py
More file actions
368 lines (306 loc) · 15.9 KB
/
Copy pathrun.py
File metadata and controls
368 lines (306 loc) · 15.9 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
#!/usr/bin/env python3
"""
冷热数据调度系统 — 启动入口
用法:
python3 run.py # 前台启动
python3 run.py --monitor # 前台启动 + 实时仪表盘
python3 run.py check # 前置依赖检查
python3 run.py --dry-run # 模拟启动(展示完整启动流程)
python3 run.py demo # 运行全部演示程序
python3 run.py start # 后台守护进程
python3 run.py stop # 停止运行中的调度器
"""
import os, sys, time, signal, json, argparse, logging, threading, subprocess
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parent
sys.path.insert(0, str(PROJECT_ROOT))
PID_FILE = PROJECT_ROOT / ".scheduler.pid"
LOG_DIR = PROJECT_ROOT / "output"
LOG_DIR.mkdir(parents=True, exist_ok=True)
DEFAULT_CONFIG = PROJECT_ROOT / "config.json"
# ═══════════════════════════════════════════════════════════════
# ANSI 样式工具(截图友好 — 即使无颜色也保持对齐)
# ═══════════════════════════════════════════════════════════════
def color(s, code):
if sys.stdout.isatty():
return f"\033[{code}m{s}\033[0m"
return s
BOLD = lambda s: color(s, "1")
GREEN = lambda s: color(s, "32")
RED = lambda s: color(s, "31")
YELLOW = lambda s: color(s, "33")
CYAN = lambda s: color(s, "36")
MAG = lambda s: color(s, "35")
def ok(text): print(f" {GREEN('✓')} {text}")
def fail(text): print(f" {RED('✗')} {text}")
def info(text): print(f" {CYAN('▸')} {text}")
def sep(title=None):
line = "─" * 50
if title:
print(f"\n {BOLD(line)}")
print(f" {BOLD(title)}")
print(f" {BOLD(line)}")
else:
print(f" {line}")
print()
# ═══════════════════════════════════════════════════════════════
# 启动横幅(截图核心区域)
# ═══════════════════════════════════════════════════════════════
def banner():
print()
print(f" {'╔' + '═'*56 + '╗'}")
print(f" {'║'}{'分层存储冷热数据智能调度文件系统':^40}{'║'}")
print(f" {'║'}{'Adaptive Tiered File Scheduling System (ATFS)':^56}{'║'}")
print(f" {'╚' + '═'*56 + '╝'}")
print()
print(f" {BOLD('架构')} CephFS · NVMe(热) / SSD(温) / HDD(冷) 三层异构存储")
print(f" {BOLD('策略')} SR-LRU 缓存替换 · DTPolicy 智能路由 · FastCDC 内容去重")
print(f" {BOLD('调度')} 实时触发式迁移 + 定期全量扫描 + 容量水位驱逐")
print(f" {BOLD('项目')} {PROJECT_ROOT}")
print()
# 配置预览
cfg_path = DEFAULT_CONFIG
if cfg_path.exists():
from src.config import ATFSConfig
cfg = ATFSConfig.load(str(cfg_path))
cap = cfg.capacity
dedup = cfg.dedup
print(f" {'─' * 48}")
print(f" {BOLD('配置预览')} (config.json)")
print(f" {'监控目录':<12} {cfg.dir}")
print(f" {'触发迁移':<12} {GREEN('已启用') if cfg.trigger else RED('已禁用')} "
f"{'替换策略':<10} {cfg.replace_policy}")
_interval = cfg.interval
_periodic_str = GREEN('已启用 (间隔 ' + str(_interval) + 's)') if cfg.periodic else RED('已禁用')
_dedup_str = GREEN('已启用') if dedup.get('enabled') else RED('已禁用')
print(f" {'定期迁移':<12} {_periodic_str} {'去重':<10} {_dedup_str}")
print(f" {'Hot 层':<12} {cap.get('hot_capacity', 0)//1024} KB "
f"{'Warm 层':<10} {cap.get('warm_capacity', 0)//1024} KB "
f"{'Cold 层':<10} {cap.get('cold_capacity', 0)//1024//1024} MB")
print()
print(f" {BOLD('启动')} python3 run.py --monitor # 推荐:启动 + 实时状态监控")
print(f" {BOLD('停止')} 按 Ctrl+C 或 python3 run.py stop")
print(f" {BOLD('帮助')} python3 run.py --help")
print()
# ═══════════════════════════════════════════════════════════════
# 前置检查(截图核心区域)
# ═══════════════════════════════════════════════════════════════
def check_redis(host="localhost", port=6379, password="1"):
try:
import redis
r = redis.Redis(host=host, port=port, db=0, password=password, socket_connect_timeout=3)
r.ping()
return True
except Exception as e:
logging.error(f"Redis 连接失败 ({host}:{port}): {e}")
return False
def check_cephfs():
try:
from src.client import CephFSClient
except ImportError as e:
logging.error(f"CephFS 客户端导入失败: {e}")
return False
try:
with CephFSClient() as client:
return client.is_connected
except Exception as e:
logging.error(f"CephFS 连接异常: {e}")
return False
def do_check(quiet=False):
"""前置依赖检查。"""
if not quiet:
sep("前置依赖检查")
checks = []
checks.append(("Python 版本", sys.version.split()[0], True))
checks.append(("Redis 服务", "localhost:6379", check_redis()))
checks.append(("CephFS 连接", "libcephfs 客户端", check_cephfs()))
checks.append(("配置文件", str(DEFAULT_CONFIG.name), DEFAULT_CONFIG.exists()))
for mod, name in [("cephfs", "cephfs"), ("redis", "redis"), ("flask", "flask")]:
try:
__import__(mod); checks.append((f"依赖 {name}", "已安装", True))
except ImportError:
checks.append((f"依赖 {name}", "未安装", False))
for label, detail, ok_flag in checks:
(ok if ok_flag else fail)(f"{label:<18} {detail}")
all_pass = all(c[2] for c in checks)
if not quiet:
print(f"\n {'结果':<18} {GREEN('全部就绪 ✓') if all_pass else RED('请修复后重试 ✗')}")
print()
return all_pass
# ═══════════════════════════════════════════════════════════════
# 主运行循环(截图核心区域)
# ═══════════════════════════════════════════════════════════════
def run_loop(config_path, monitor=False):
from src.client import CephFSClient
from src.scheduler import Scheduler
from src.config import ATFSConfig
stop_event = threading.Event()
# ── 保存 PID(即使前台也记录,便于 run.py stop 定位) ──
PID_FILE.write_text(str(os.getpid()))
def signal_handler(signum, frame):
signame = signal.Signals(signum).name
print(f"\n {YELLOW(BOLD(f'⏳ 收到 {signame},正在关闭...'))}")
stop_event.set()
signal.signal(signal.SIGINT, signal_handler)
signal.signal(signal.SIGTERM, signal_handler)
sep("连接 CephFS")
with CephFSClient(scheduler_config=str(config_path)) as client:
if not client.is_connected:
fail("CephFS 连接失败 — 请检查集群状态和挂载点")
sys.exit(1)
ok("CephFS 集群已连接")
cfg = ATFSConfig.load(str(config_path))
ok(f"监控根目录: {cfg.dir}")
scheduler = Scheduler(client)
print()
# ── 调度器状态横幅 ──
sep("调度器已启动")
print(f" {BOLD('模式')} {'触发迁移' if scheduler.trigger else ''}"
f"{' + 定期迁移' if scheduler.periodic else ''}"
f"{' (仅监控)' if not scheduler.trigger and not scheduler.periodic else ''}")
print(f" {BOLD('监控目录')} {scheduler.dir}")
print(f" {BOLD('替换策略')} {scheduler.replace_policy}"
f"{f' (SRLU 保护比={scheduler.srlru_ratio})' if scheduler.replace_policy == 'srlru' else ''}")
print(f" {BOLD('触发迁移')} {GREEN('运行中') if scheduler.trigger else RED('已禁用')}")
print(f" {BOLD('定期迁移')} {GREEN(f'运行中 (间隔 {scheduler.interval}s)') if scheduler.periodic else RED('已禁用')}")
print()
info(f"{YELLOW('按 Ctrl+C 停止')}")
if monitor:
info(f"{CYAN('实时仪表盘已开启 (每 10 秒刷新)')}")
print()
# ── 等待停止信号 ──
stop_event.wait()
# ── 关闭序列(截图核心区域) ──
print()
sep("正在关闭")
info("停止调度器线程...")
scheduler.stop()
ok("调度器线程已停止")
# ── 清理 PID ──
if PID_FILE.exists() and PID_FILE.read_text().strip() == str(os.getpid()):
PID_FILE.unlink()
print()
sep("调度器已安全停止")
ok("所有资源已释放")
info("运行指标报告已保存 → output/metrics_report.txt")
print()
# ═══════════════════════════════════════════════════════════════
# 守护进程管理
# ═══════════════════════════════════════════════════════════════
def daemon_start(config_path, log_file=None):
if PID_FILE.exists():
pid = PID_FILE.read_text().strip()
if pid:
try:
os.kill(int(pid), 0)
fail(f"调度器已在运行 (PID={pid})")
sys.exit(1)
except OSError:
PID_FILE.unlink(missing_ok=True)
log_path = log_file or (LOG_DIR / "scheduler.log")
cmd = [sys.executable, __file__, "run", "--config", str(config_path)]
proc = subprocess.Popen(cmd, stdout=open(log_path, "a"), stderr=subprocess.STDOUT,
stdin=subprocess.DEVNULL, start_new_session=True)
PID_FILE.write_text(str(proc.pid))
print(f" {GREEN('✓')} 调度器已后台启动 (PID={proc.pid})")
print(f" {CYAN('▸')} 日志: {log_path}")
print(f" {CYAN('▸')} 停止: python3 run.py stop")
def daemon_stop():
if not PID_FILE.exists():
info("未找到运行中的调度器")
return
pid = PID_FILE.read_text().strip()
if not pid:
PID_FILE.unlink(missing_ok=True)
info("未找到运行中的调度器")
return
try:
os.kill(int(pid), signal.SIGTERM)
print(f" {CYAN('▸')} 已向 PID={pid} 发送 SIGTERM 信号")
for _ in range(10):
time.sleep(0.5)
try:
os.kill(int(pid), 0)
except OSError:
break
PID_FILE.unlink(missing_ok=True)
ok("调度器已停止")
except ProcessLookupError:
PID_FILE.unlink(missing_ok=True)
info("进程已不存在")
except PermissionError:
fail(f"权限不足,无法操作进程 {pid}")
# ═══════════════════════════════════════════════════════════════
# 一键演示 / 检查 / 模拟运行
# ═══════════════════════════════════════════════════════════════
def run_demos():
"""顺序运行全部演示程序。"""
demos = [
("Demo 01 — 基本文件接口测试", "demo/01_basic_io.py"),
("Demo 02 — 路由策略测试", "demo/02_route_policy.py"),
("Demo 03 — FIU Trace 回放", "demo/03_fiu_replay.py"),
]
for name, script in demos:
sp = PROJECT_ROOT / script
if not sp.exists():
fail(f"未找到: {script}")
continue
sep(name)
result = subprocess.run([sys.executable, str(sp)], cwd=PROJECT_ROOT)
(ok if result.returncode == 0 else fail)(f"{name} → {'完成' if result.returncode == 0 else f'失败 (code={result.returncode})'}")
# ═══════════════════════════════════════════════════════════════
# 入口
# ═══════════════════════════════════════════════════════════════
if __name__ == "__main__":
parser = argparse.ArgumentParser(
description="冷热数据调度系统 — 分层存储冷热数据智能调度文件系统",
formatter_class=argparse.RawDescriptionHelpFormatter,
epilog="""
示例:
python3 run.py 前台启动(展示完整启动流程)
python3 run.py --monitor 前台启动 + 实时仪表盘监控
python3 run.py check 前置依赖检查
python3 run.py --dry-run 模拟启动(展示完整流程,不实际启动)
python3 run.py demo 一键运行全部演示程序
python3 run.py start 后台守护进程启动
python3 run.py stop 停止运行中的调度器
./start.sh Shell 脚本启动(自动挂载CephFS)
""")
parser.add_argument("action", nargs="?", default="run",
choices=["run", "start", "stop", "check", "demo"],
help="操作:run(前台) / start(后台) / stop(停止) / check(检查) / demo(演示)")
parser.add_argument("--config", "-c", default=str(DEFAULT_CONFIG), help="配置文件路径")
parser.add_argument("--monitor", "-m", action="store_true", help="启动后每10秒刷新彩色仪表盘")
parser.add_argument("--dry-run", "-n", action="store_true", help="模拟启动:展示完整流程但不实际运行")
parser.add_argument("--log", "-l", help="守护进程日志路径")
parser.add_argument("--verbose", "-v", action="store_true", help="详细日志")
args = parser.parse_args()
log_level = logging.INFO if args.verbose else logging.WARNING
logging.basicConfig(level=log_level, format="%(asctime)s [%(levelname)s] %(message)s", datefmt="%Y-%m-%d %H:%M:%S")
config_path = Path(args.config)
if args.action not in ("check", "demo", "stop") and not config_path.exists():
print(f" {RED('✗')} 配置文件不存在: {config_path}"); sys.exit(1)
# ── 动作分发 ──
if args.action == "check":
banner(); do_check(); sys.exit(0)
if args.action == "stop":
daemon_stop(); sys.exit(0)
if args.action == "start":
banner()
if do_check(quiet=True): daemon_start(config_path, args.log)
else: fail("前置检查未通过")
sys.exit(0)
if args.action == "demo":
banner(); do_check(quiet=True); run_demos(); sys.exit(0)
# ── "run" ──
banner()
if not do_check(quiet=True):
fail("前置检查未通过 — 运行 python3 run.py check 查看详情")
sys.exit(1)
if args.dry_run:
ok("模拟运行:前置检查通过,配置已加载")
print(f" {CYAN('▸')} 启动命令: python3 run.py")
print(f" {CYAN('▸')} 配置文件: {config_path.name}")
print()
sys.exit(0)
run_loop(config_path, monitor=args.monitor)