""" 处理步骤监控服务 - 记录和监控产品处理流程 """ import os import time import uuid import json import subprocess import threading import logging from datetime import datetime from models.database import db from services.search_service import search_service from services.paramhub_client import paramhub_client logger = logging.getLogger('process_monitor') # 处理步骤定义 PROCESS_STEPS = [ {'num': 1, 'name': '搜索内容库', 'description': '从内容库搜索相关文章'}, {'num': 2, 'name': '搜索互联网', 'description': '从互联网搜索最新数据'}, {'num': 3, 'name': '抓取网页内容', 'description': '抓取搜索结果网页的详细内容'}, {'num': 4, 'name': '提取产品数据(智能体)', 'description': '调用hz4th_editor智能体提取产品相关内容'}, {'num': 5, 'name': '填充字段', 'description': '根据分类字段配置填充数据'}, {'num': 6, 'name': '提交审核', 'description': '提交到ParamHub待审核区'}, ] class ProcessMonitor: """处理步骤监控器""" def __init__(self): self.active_sessions = {} self.step_timers = {} def create_session_id(self): """生成会话ID""" return f"proc_{datetime.now().strftime('%Y%m%d_%H%M%S')}_{uuid.uuid4().hex[:8]}" def start_process(self, product_name, category=None, subcategory=None): """启动产品处理流程""" session_id = self.create_session_id() # 创建会话记录 db.create_process_session(session_id, product_name, category, subcategory) # 初始化控制信息 self.active_sessions[session_id] = { 'paused': False, 'stop': False, 'current_step': 0 } # 启动后台线程处理 thread = threading.Thread( target=self._run_process, args=(session_id, product_name, category, subcategory), daemon=True ) thread.start() logger.info(f"启动处理会话: {session_id}, 产品: {product_name}") return session_id def _run_process(self, session_id, product_name, category, subcategory): """执行处理流程""" try: db.update_session_status(session_id, 'running') result = {'success': False, 'message': '', 'review_id': None} all_data = { 'library_results': [], 'internet_results': [], 'fetched_contents': [], 'extracted_data': None, 'filled_data': None } # 步骤1: 搜索内容库 if not self._check_pause(session_id): self._start_step(session_id, product_name, 1, '搜索内容库') try: articles = db.search_articles(product_name, category) all_data['library_results'] = articles self._complete_step(session_id, 1, {'count': len(articles)}) logger.info(f"[{session_id}] 步骤1完成: 找到 {len(articles)} 篇文章") except Exception as e: self._fail_step(session_id, 1, str(e)) result['message'] = f'搜索内容库失败: {e}' # 步骤2: 搜索互联网 if not self._check_pause(session_id) and not result.get('message'): self._start_step(session_id, product_name, 2, '搜索互联网') try: internet_results = search_service.search_internet(product_name, max_results=10) all_data['internet_results'] = internet_results self._complete_step(session_id, 2, {'count': len(internet_results)}) logger.info(f"[{session_id}] 步骤2完成: 找到 {len(internet_results)} 条结果") except Exception as e: self._complete_step(session_id, 2, {'count': 0, 'error': str(e)}) # 步骤3: 抓取网页内容 if not self._check_pause(session_id) and all_data['internet_results']: self._start_step(session_id, product_name, 3, '抓取网页内容') try: fetched = [] failed_count = 0 urls_to_fetch = [r['url'] for r in all_data['internet_results'][:5]] for i, url in enumerate(urls_to_fetch): if self._check_pause(session_id): break fetch_result = search_service.fetch_url_content(url) if fetch_result.get('success'): title = fetch_result.get('title', '') content = fetch_result.get('content', '') fetched.append({ 'url': url, 'title': title, 'content': content[:500] }) # 保存到内容库 try: existing = db.search_articles(url) if not any(a.get('url') == url for a in existing): db.add_article( product_names=[], category=category or '', keywords=[], summary=content[:200] if content else '', content=content, source=url, url=url, search_title=title ) logger.info(f"[{session_id}] 已保存到内容库: {title[:30]}") except Exception as save_error: logger.warning(f"[{session_id}] 保存内容库失败: {save_error}") else: # 记录失败URL failed_count += 1 error_msg = fetch_result.get('error', '抓取失败') try: db.add_failed_url(url, product_name, error_msg, source='process_monitor') logger.warning(f"[{session_id}] 抓取失败,已记录: {url}") except Exception as e: logger.error(f"[{session_id}] 记录失败URL出错: {e}") time.sleep(0.3) all_data['fetched_contents'] = fetched self._complete_step(session_id, 3, {'count': len(fetched), 'failed': failed_count}) logger.info(f"[{session_id}] 步骤3完成: 抓取 {len(fetched)} 个网页, 失败 {failed_count} 个") except Exception as e: self._fail_step(session_id, 3, str(e)) # 步骤4: 提取产品数据(调用智能体执行) if not self._check_pause(session_id): self._start_step(session_id, product_name, 4, '提取产品数据(智能体)') try: # 构建任务文本 task_text = self._build_agent_task( product_name, category, subcategory, all_data ) # 记录任务文本 self._complete_step(session_id, 4, { 'agent': 'hz4th_editor', 'task_text': task_text, 'status': 'calling_agent' }) # 调用智能体 agent_result = self._call_agent(task_text) if agent_result.get('success'): extracted = self._parse_agent_response(agent_result.get('output', '')) all_data['extracted_data'] = extracted if extracted: self._complete_step(session_id, 4, { 'has_data': True, 'agent': 'hz4th_editor', 'task_text': task_text, 'agent_output': agent_result.get('output', '')[:2000] }) else: self._complete_step(session_id, 4, { 'has_data': False, 'agent': 'hz4th_editor', 'task_text': task_text, 'agent_output': agent_result.get('output', '')[:2000] }, status='skipped') result['message'] = '智能体无法提取有效数据' else: self._fail_step(session_id, 4, f"智能体调用失败: {agent_result.get('error', '未知错误')}") result['message'] = f'智能体调用失败: {agent_result.get("error")}' except Exception as e: self._fail_step(session_id, 4, str(e)) # 步骤5: 填充字段 if not self._check_pause(session_id) and all_data['extracted_data']: self._start_step(session_id, product_name, 5, '填充字段') try: filled = self._fill_fields(all_data['extracted_data'], category, subcategory) all_data['filled_data'] = filled if filled: self._complete_step(session_id, 5, {'filled': True}) else: self._fail_step(session_id, 5, '填充数据失败') except Exception as e: self._fail_step(session_id, 5, str(e)) # 步骤6: 提交审核 if not self._check_pause(session_id) and all_data['filled_data']: self._start_step(session_id, product_name, 6, '提交审核') try: category_type = self._get_category_type(category) success, review_id_or_error = paramhub_client.submit_for_review( category_type, all_data['filled_data'], subcategory ) if success: self._complete_step(session_id, 6, {'review_id': review_id_or_error}) result['success'] = True result['review_id'] = review_id_or_error db.update_session_status(session_id, 'completed', review_id=review_id_or_error, result=json.dumps(result, ensure_ascii=False)) db.add_process_history( product_name=product_name, category=category, subcategory=subcategory, status='submitted', review_id=review_id_or_error, details=all_data ) logger.info(f"[{session_id}] 步骤6完成: 提交成功") else: self._fail_step(session_id, 6, review_id_or_error) db.update_session_status(session_id, 'failed') except Exception as e: self._fail_step(session_id, 6, str(e)) # 清理 if session_id in self.active_sessions: del self.active_sessions[session_id] return result except Exception as e: logger.error(f"处理会话异常: {session_id} - {e}") db.update_session_status(session_id, 'failed') return {'success': False, 'message': str(e)} def _start_step(self, session_id, product_name, step_num, step_name): """开始步骤""" db.update_session_status(session_id, 'running', current_step=step_num) db.add_process_step(session_id, product_name, step_num, step_name) if session_id not in self.step_timers: self.step_timers[session_id] = {} self.step_timers[session_id][step_num] = time.time() def _complete_step(self, session_id, step_num, step_data=None, status='completed'): """完成步骤""" duration_ms = None if session_id in self.step_timers and step_num in self.step_timers[session_id]: duration_ms = int((time.time() - self.step_timers[session_id][step_num]) * 1000) db.update_step_status(session_id, step_num, status, step_data=step_data, duration_ms=duration_ms) def _fail_step(self, session_id, step_num, error_message): """步骤失败""" duration_ms = None if session_id in self.step_timers and step_num in self.step_timers[session_id]: duration_ms = int((time.time() - self.step_timers[session_id][step_num]) * 1000) db.update_step_status(session_id, step_num, 'failed', error_message=error_message, duration_ms=duration_ms) db.update_session_status(session_id, 'failed') def _check_pause(self, session_id): """检查是否暂停""" if session_id not in self.active_sessions: return False session = self.active_sessions[session_id] if session.get('stop'): return True while session.get('paused'): time.sleep(0.5) if session.get('stop'): return True return False def pause_session(self, session_id): """暂停会话""" if session_id in self.active_sessions: self.active_sessions[session_id]['paused'] = True db.pause_session(session_id, '用户暂停') return True return False def resume_session(self, session_id): """继续会话""" if session_id in self.active_sessions: self.active_sessions[session_id]['paused'] = False db.resume_session(session_id) return True return False def stop_session(self, session_id): """停止会话""" if session_id in self.active_sessions: self.active_sessions[session_id]['stop'] = True self.active_sessions[session_id]['paused'] = False db.update_session_status(session_id, 'stopped') return True return False def get_session_status(self, session_id): """获取会话状态""" session = db.get_process_session(session_id) if session: steps = db.get_process_steps(session_id) return {'session': session, 'steps': steps} return None def _build_agent_task(self, product_name, category, subcategory, all_data): """构建智能体任务文本""" # 读取模板 template_file = os.path.join( os.path.dirname(os.path.dirname(__file__)), 'config', 'agent_task_template.txt' ) if os.path.exists(template_file): with open(template_file, 'r', encoding='utf-8') as f: template = f.read() else: # 默认模板 template = ( "请从以下数据中提取产品「{{product_name}}」的相关内容。\n" "类别: {{category}} / {{subcategory}}\n\n" "内容库结果:\n{{library_results}}\n\n" "互联网抓取内容:\n{{internet_results}}\n\n" "要求:只提取与该产品信息直接相关的内容,排除无关产品。以JSON格式输出。" ) # 构建内容库搜索结果 library_lines = [] for i, article in enumerate(all_data.get('library_results', [])[:10], 1): title = article.get('search_title', article.get('title', '无标题')) url = article.get('url', article.get('source', '无URL')) summary = article.get('summary', '')[:200] library_lines.append(f" [{i}] 标题: {title}\n URL: {url}\n 摘要: {summary}") library_text = '\n'.join(library_lines) if library_lines else '(无内容库搜索结果)' # 构建互联网抓取内容 internet_lines = [] for i, item in enumerate(all_data.get('fetched_contents', [])[:10], 1): title = item.get('title', '无标题') url = item.get('url', '无URL') content = item.get('content', '')[:300] internet_lines.append(f" [{i}] 标题: {title}\n URL: {url}\n 内容片段: {content}") internet_text = '\n'.join(internet_lines) if internet_lines else '(无互联网抓取内容)' # 填充模板 task = template.replace('{{product_name}}', product_name or '未知') task = task.replace('{{category}}', category or '未分类') task = task.replace('{{subcategory}}', subcategory or '无') task = task.replace('{{library_results}}', library_text) task = task.replace('{{internet_results}}', internet_text) return task def _call_agent(self, task_text): """调用智能体执行任务""" try: cmd = [ 'openclaw', 'agent', '--agent', 'hz4th_editor', '--message', task_text ] logger.info(f"调用智能体命令: openclaw agent --agent hz4th_editor --message '[任务文本 {len(task_text)} 字符]'") result = subprocess.run( cmd, capture_output=True, text=True, timeout=300 # 5分钟超时 ) if result.returncode == 0: output = result.stdout.strip() logger.info(f"智能体返回: {output[:500]}...") return {'success': True, 'output': output} else: error = result.stderr.strip() or result.stdout.strip() logger.error(f"智能体调用失败: {error}") return {'success': False, 'error': error} except subprocess.TimeoutExpired: return {'success': False, 'error': '智能体执行超时(>5分钟)'} except FileNotFoundError: return {'success': False, 'error': 'openclaw命令未找到'} except Exception as e: return {'success': False, 'error': str(e)} def _parse_agent_response(self, output): """解析智能体返回的结果""" if not output: return None # 尝试从输出中提取JSON import re # 查找JSON块 json_match = re.search(r'```(?:json)?\s*(\{.*?\})\s*```', output, re.DOTALL) if json_match: try: data = json.loads(json_match.group(1)) return { 'name': data.get('name', ''), 'extracted_fields': data.get('extracted_fields', {}), 'sources': data.get('sources', []), 'confidence': data.get('confidence', 'unknown'), 'raw_output': output } except json.JSONDecodeError: pass # 尝试直接解析整个输出为JSON try: data = json.loads(output) return { 'name': data.get('name', ''), 'extracted_fields': data.get('extracted_fields', {}), 'sources': data.get('sources', []), 'confidence': data.get('confidence', 'unknown'), 'raw_output': output } except json.JSONDecodeError: pass # 如果无法解析为JSON,将原始输出作为raw_content保存 return { 'name': '', 'raw_content': output, 'raw_output': output } def _extract_data(self, product_name, all_data): """提取产品数据(备用,已被智能体替代)""" all_content = [] for article in all_data.get('library_results', []): content = article.get('content', '') if content: all_content.append(content) for item in all_data.get('fetched_contents', []): content = item.get('content', '') if content: all_content.append(content) if not all_content: return None return { 'name': product_name, 'raw_content': '\n---\n'.join(all_content[:3]) } def _fill_fields(self, extracted_data, category, subcategory): """填充字段""" if not extracted_data: return None import re filled = { 'name': extracted_data['name'], 'visible': True, 'is_pinned': False } content = extracted_data.get('raw_content', '') params_match = re.search(r'(\d+(?:\.\d+)?)\s*[Bb]', content) if params_match: filled['parameters'] = f"{params_match.group(1)}B" date_match = re.search(r'(\d{4}[-/]\d{1,2}[-/]\d{1,2})', content) if date_match: filled['publish_date'] = date_match.group(1).replace('/', '-') filled['_source'] = 'auto_manager' filled['_extracted_at'] = datetime.now().isoformat() return filled def _get_category_type(self, category): """获取分类类型""" if not category: return 'dynamic' category_lower = category.lower() if 'model' in category_lower or 'ai' in category_lower: return 'model' elif 'gpu' in category_lower: return 'gpu' elif 'cpu' in category_lower: return 'cpu' return 'dynamic' # 全局处理监控实例 process_monitor = ProcessMonitor()