feat: 服务端集成自动打包流程
This commit is contained in:
@@ -15,12 +15,17 @@ class LogStreamer:
|
||||
self._queues: Dict[str, asyncio.Queue] = {}
|
||||
self._subscribers: Dict[str, list] = {}
|
||||
self._log_lines: Dict[str, list] = {}
|
||||
self._completed: set[str] = set()
|
||||
|
||||
def create_queue(self, task_id: str) -> asyncio.Queue:
|
||||
"""创建任务的日志队列"""
|
||||
"""创建任务的日志队列(重复调用时保留已有日志)"""
|
||||
if task_id in self._queues:
|
||||
return self._queues[task_id]
|
||||
|
||||
queue = asyncio.Queue(maxsize=LOG_QUEUE_MAX_SIZE)
|
||||
self._queues[task_id] = queue
|
||||
self._log_lines[task_id] = []
|
||||
self._log_lines.setdefault(task_id, [])
|
||||
self._completed.discard(task_id)
|
||||
return queue
|
||||
|
||||
async def emit(self, task_id: str, message: str, level: str = "info"):
|
||||
@@ -66,11 +71,19 @@ class LogStreamer:
|
||||
if not queue:
|
||||
return
|
||||
|
||||
# 回放已缓存的历史日志
|
||||
existing = self._log_lines.get(task_id, [])
|
||||
# Messages are stored both in history and the live queue. Replay the
|
||||
# history once, then discard its queue copies before waiting for new logs.
|
||||
existing = list(self._log_lines.get(task_id, []))
|
||||
completed = task_id in self._completed
|
||||
while not queue.empty():
|
||||
queue.get_nowait()
|
||||
|
||||
for entry in existing:
|
||||
yield entry
|
||||
|
||||
if completed:
|
||||
return
|
||||
|
||||
# 流式推送新消息
|
||||
while True:
|
||||
try:
|
||||
@@ -86,6 +99,7 @@ class LogStreamer:
|
||||
|
||||
def complete(self, task_id: str):
|
||||
"""标记任务日志结束"""
|
||||
self._completed.add(task_id)
|
||||
if task_id in self._queues:
|
||||
try:
|
||||
self._queues[task_id].put_nowait(None)
|
||||
@@ -95,6 +109,7 @@ class LogStreamer:
|
||||
def cleanup(self, task_id: str):
|
||||
"""清理任务日志队列"""
|
||||
self._queues.pop(task_id, None)
|
||||
self._completed.discard(task_id)
|
||||
|
||||
def save_log(self, task_id: str, build_dir: Path):
|
||||
"""将收集的日志保存到独立日志目录(不会随构建目录删除)"""
|
||||
|
||||
Reference in New Issue
Block a user