Files
param-auto-manager/services/process_monitor.py
T
hz4th_coder eea7269eda feat: 步骤5填充字段改用智能体执行
- 新增步骤5任务模板(config/agent_fill_fields_template.txt)
- 步骤5调用智能体hz4th_editor执行:
  1. 获取API文档了解对应类别字段定义
  2. 从内容库获取相关内容数据
  3. 整理产品参数
  4. 通过API提交到ParamHub审核系统
- /process页面新增步骤5模板编辑面板
- 步骤6改为确认提交结果(备用本地提交)
- 新增API: GET/POST /api/process/fill-fields-template
2026-07-15 16:31:00 +08:00

715 lines
32 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
处理步骤监控服务 - 记录和监控产品处理流程
"""
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': '调用hz4th_editor智能体整理产品数据并提交审核'},
{'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', '')
article_id = None
# 保存到内容库
try:
existing = db.search_articles(url)
if existing and len(existing) > 0:
# 已存在,使用现有ID
article_id = existing[0].get('id')
logger.info(f"[{session_id}] 内容库已存在: {title[:30]}, ID={article_id}")
else:
# 新增,获取返回的ID
article_id = 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]}, ID={article_id}")
except Exception as save_error:
logger.warning(f"[{session_id}] 保存内容库失败: {save_error}")
fetched.append({
'id': article_id,
'url': url,
'title': title,
'content': content[:500]
})
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
)
# 调用智能体
agent_result = self._call_agent(task_text)
if agent_result.get('success'):
parsed = self._parse_agent_response(agent_result.get('output', ''))
if parsed and parsed.get('relevant_ids'):
# 根据ID从内容库获取实际内容
relevant_contents = []
for aid in parsed['relevant_ids']:
article = db.get_article_by_id(aid)
if article:
relevant_contents.append({
'id': aid,
'title': article.get('search_title', ''),
'url': article.get('url', ''),
'content': article.get('content', ''),
'summary': article.get('summary', ''),
'analysis': parsed.get('analysis', {}).get(str(aid), '')
})
all_data['extracted_data'] = {
'name': product_name,
'relevant_ids': parsed['relevant_ids'],
'relevant_contents': relevant_contents,
'confidence': parsed.get('confidence', 'unknown'),
'raw_output': agent_result.get('output', '')
}
self._complete_step(session_id, 4, {
'has_data': True,
'agent': 'hz4th_editor',
'task_text': task_text,
'relevant_ids': parsed['relevant_ids'],
'relevant_count': len(relevant_contents),
'confidence': parsed.get('confidence', 'unknown'),
'agent_output': agent_result.get('output', '')[:2000]
})
logger.info(f"[{session_id}] 步骤4完成: 智能体返回 {len(parsed['relevant_ids'])} 个相关ID")
else:
all_data['extracted_data'] = None
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'] = '智能体未找到相关数据ID'
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:
# 构建任务文本
fill_task_text = self._build_fill_fields_task(
product_name, category, subcategory, all_data['extracted_data']
)
# 调用智能体
fill_agent_result = self._call_agent(fill_task_text)
if fill_agent_result.get('success'):
fill_parsed = self._parse_fill_agent_response(fill_agent_result.get('output', ''))
if fill_parsed and fill_parsed.get('success'):
review_id = fill_parsed.get('review_id')
submitted_data = fill_parsed.get('submitted_data', {})
all_data['filled_data'] = submitted_data
self._complete_step(session_id, 5, {
'filled': True,
'agent': 'hz4th_editor',
'task_text': fill_task_text,
'review_id': review_id,
'submitted_data': submitted_data,
'agent_output': fill_agent_result.get('output', '')[:2000]
})
# 如果智能体返回了review_id,直接标记步骤6完成
if review_id:
result['success'] = True
result['review_id'] = review_id
db.update_session_status(session_id, 'completed',
review_id=review_id,
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,
details=all_data
)
logger.info(f"[{session_id}] 步骤5完成: 智能体提交成功, review_id={review_id}")
else:
self._complete_step(session_id, 5, {
'filled': True,
'agent': 'hz4th_editor',
'task_text': fill_task_text,
'submitted_data': submitted_data,
'message': '数据已整理但未获取到review_id',
'agent_output': fill_agent_result.get('output', '')[:2000]
})
result['message'] = '数据已整理但提交状态未知'
else:
self._fail_step(session_id, 5, f"智能体执行失败: {fill_parsed.get('message', '未知错误') if fill_parsed else '解析失败'}")
result['message'] = f'智能体执行失败'
else:
self._fail_step(session_id, 5, f"智能体调用失败: {fill_agent_result.get('error', '未知错误')}")
result['message'] = f'智能体调用失败: {fill_agent_result.get("error")}'
except Exception as e:
self._fail_step(session_id, 5, str(e))
# 步骤6: 确认提交结果(如果步骤5未获取到review_id,尝试本地提交)
if not self._check_pause(session_id) and all_data['filled_data'] and not result.get('review_id'):
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 = (
"请分析以下数据ID是否与产品「{{product_name}}」相关且对提取参数有用。\n"
"类别: {{category}} / {{subcategory}}\n\n"
"内容库结果ID: {{library_results}}\n\n"
"互联网已入库ID: {{internet_results}}\n\n"
"要求:输出相关且有用的ID列表,以JSON格式输出。"
)
# 构建内容库搜索结果ID列表
library_ids = []
for article in all_data.get('library_results', []):
aid = article.get('id')
if aid:
title = article.get('search_title', article.get('title', ''))
library_ids.append(f"ID {aid}: {title}")
library_text = '\n'.join(library_ids) if library_ids else '(无内容库搜索结果)'
# 构建互联网已入库数据ID列表
internet_ids = []
for item in all_data.get('fetched_contents', []):
aid = item.get('id')
if aid:
title = item.get('title', '')
internet_ids.append(f"ID {aid}: {title}")
internet_text = '\n'.join(internet_ids) if internet_ids 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):
"""解析智能体返回的结果,提取relevant_ids"""
if not output:
return None
# 尝试从输出中提取JSON
import re
parsed_data = None
# 查找JSON块
json_match = re.search(r'```(?:json)?\s*(\{.*?\})\s*```', output, re.DOTALL)
if json_match:
try:
parsed_data = json.loads(json_match.group(1))
except json.JSONDecodeError:
pass
# 尝试直接解析整个输出为JSON
if not parsed_data:
try:
parsed_data = json.loads(output)
except json.JSONDecodeError:
pass
if parsed_data:
relevant_ids = parsed_data.get('relevant_ids', [])
# 确保都是整数
relevant_ids = [int(x) for x in relevant_ids if str(x).isdigit()]
return {
'relevant_ids': relevant_ids,
'analysis': parsed_data.get('analysis', {}),
'excluded_ids': parsed_data.get('excluded_ids', []),
'exclusion_reasons': parsed_data.get('exclusion_reasons', {}),
'confidence': parsed_data.get('confidence', 'unknown'),
'raw_output': output
}
# 无法解析为JSON,尝试从文本中提取ID
id_matches = re.findall(r'(?:ID|id)[\s:]*(\d+)', output)
if id_matches:
return {
'relevant_ids': [int(x) for x in id_matches],
'analysis': {},
'confidence': 'low',
'raw_output': output
}
return None
def _build_fill_fields_task(self, product_name, category, subcategory, extracted_data):
"""构建步骤5填充字段的智能体任务文本"""
# 读取模板
template_file = os.path.join(
os.path.dirname(os.path.dirname(__file__)),
'config', 'agent_fill_fields_template.txt'
)
if os.path.exists(template_file):
with open(template_file, 'r', encoding='utf-8') as f:
template = f.read()
else:
# 默认模板
template = (
"请根据内容库数据ID {{relevant_content_ids}} 整理产品「{{product_name}}」的参数并提交审核。\n"
"类别: {{category}} / {{subcategory}}\n"
"参考API文档: http://192.168.2.8:12007/hz4th_coder/param-hub-python/src/branch/master/API.md"
)
# 构建相关内容ID列表
relevant_ids = extracted_data.get('relevant_ids', [])
relevant_contents = extracted_data.get('relevant_contents', [])
if relevant_contents:
content_lines = []
for item in relevant_contents:
aid = item.get('id', '')
title = item.get('title', '')
content_lines.append(f"ID {aid}: {title}")
relevant_text = '\n'.join(content_lines)
elif relevant_ids:
relevant_text = '\n'.join([f"ID {aid}" for aid in relevant_ids])
else:
relevant_text = '(无相关内容ID'
# 填充模板
task = template.replace('{{product_name}}', product_name or '未知')
task = task.replace('{{category}}', category or '未分类')
task = task.replace('{{subcategory}}', subcategory or '无')
task = task.replace('{{relevant_content_ids}}', relevant_text)
return task
def _parse_fill_agent_response(self, output):
"""解析步骤5智能体返回的结果"""
if not output:
return None
import re
parsed_data = None
# 查找JSON块
json_match = re.search(r'```(?:json)?\s*(\{.*?\})\s*```', output, re.DOTALL)
if json_match:
try:
parsed_data = json.loads(json_match.group(1))
except json.JSONDecodeError:
pass
# 尝试直接解析整个输出为JSON
if not parsed_data:
try:
parsed_data = json.loads(output)
except json.JSONDecodeError:
pass
if parsed_data:
return {
'success': parsed_data.get('success', False),
'review_id': parsed_data.get('review_id'),
'submitted_data': parsed_data.get('submitted_data', {}),
'data_sources': parsed_data.get('data_sources', []),
'message': parsed_data.get('message', ''),
'raw_output': output
}
# 尝试从文本中提取review_id
review_match = re.search(r'review[_-]?id[\s:]*([\w-]+)', output, re.I)
if review_match:
return {
'success': True,
'review_id': review_match.group(1),
'submitted_data': {},
'message': '从输出中提取到review_id',
'raw_output': output
}
return None
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.get('name', ''),
'visible': True,
'is_pinned': False
}
# 从relevant_contents中拼接所有内容
relevant_contents = extracted_data.get('relevant_contents', [])
all_content = '\n---\n'.join([
c.get('content', '') or c.get('summary', '')
for c in relevant_contents
if c.get('content') or c.get('summary')
])
# 兼容旧格式
if not all_content:
all_content = extracted_data.get('raw_content', '')
params_match = re.search(r'(\d+(?:\.\d+)?)\s*[Bb]', all_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})', all_content)
if date_match:
filled['publish_date'] = date_match.group(1).replace('/', '-')
filled['_source'] = 'auto_manager'
filled['_extracted_at'] = datetime.now().isoformat()
filled['_relevant_ids'] = extracted_data.get('relevant_ids', [])
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()