"""团队沟通引擎 — 通用问题冒泡(多 agent 感知)。 问题类型(pipeline_problem_types)定义冒泡路径(escalation_path,处理方序列)。 处理方 = role + agentid(同一角色可能有多个 agent,必须精确到具体 agent; 人类角色 agentid 为空表示"该角色的任意人")。 问题产生后从路径第一个处理方开始,每个处理方要么解决(resolve)要么未解决 继续冒泡(escalate)到下一个: - 解决了 → status=answered,不再冒泡 - 没解决 → escalation_pos+1,沿路径继续走 - 路径尽头通常是"人"(human 角色)兜底 冒泡路径条目格式: - "role_name" → 该角色任意 agent/人 - {"role":"x", "agentid":"y"} → 指定角色 + 指定 agent """ import json import logging from sqlor.dbpools import DBPools from appPublic.uniqueID import getID DBNAME = "pipeline" logger = logging.getLogger("pipeline.communication") S_PENDING = "pending" # 冒泡中(未解决) S_ANSWERED = "answered" # 已解决,停止冒泡 DEFAULT_PIPELINE_ID = "sdlc_general" def _get_db(): db = DBPools() if not db.databases: from appPublic.jsonConfig import getConfig config = getConfig() if config.databases: db.databases = config.databases return db, DBNAME def _parse_path(raw): """解析冒泡路径为 [{"role":..., "agentid":...}] 列表(统一归一化)。""" if not raw: return [] if isinstance(raw, (list, tuple)): items = raw else: try: items = json.loads(raw) except (json.JSONDecodeError, TypeError): return [] if not isinstance(items, list): return [] result = [] for it in items: if isinstance(it, str): result.append({"role": it, "agentid": ""}) elif isinstance(it, dict): result.append({ "role": str(it.get("role", "") or ""), "agentid": str(it.get("agentid", "") or ""), }) return [x for x in result if x["role"]] def current_handler_of(rec) -> dict: """当前该处理这个问题的处理方 {"role", "agentid"}(由 escalation_path + escalation_pos 决定)。""" path = _parse_path(getattr(rec, 'escalation_path', '')) try: pos = int(getattr(rec, 'escalation_pos', 0) or 0) except (TypeError, ValueError): pos = 0 if not path: # 旧数据/无路径兜底:待主 agent 处理 return {"role": "main_agent", "agentid": ""} if pos < 0 or pos >= len(path): return path[-1] return path[pos] def _handler_matches(handler: dict, role: str, agentid: str = None) -> bool: """处理方是否匹配目标(role, agentid)。agentid 为空=任意 agent。""" if handler.get("role") != role: return False if agentid and handler.get("agentid") and handler.get("agentid") != agentid: return False # 指定了具体 agent,但处理方是别的 agent return True async def _resolve_pipeline_id(sor, project_id): """从项目解析产线 id;查不到用默认 sdlc_general。""" if project_id: recs = await sor.sqlExe( "SELECT pipeline_id FROM sd_projects WHERE id=${p}$", {"p": project_id}) if recs: pid = getattr(recs[0], 'pipeline_id', '') or '' if pid: return pid return DEFAULT_PIPELINE_ID async def get_problem_type(sor, pipeline_id, name): """查问题类型定义(含冒泡路径,已归一化)。找不到返回 None。""" recs = await sor.sqlExe( "SELECT * FROM pipeline_problem_types " "WHERE pipeline_id=${p}$ AND name=${n}$ LIMIT 1", {"p": pipeline_id, "n": name}) if not recs: return None rec = recs[0] return { 'id': getattr(rec, 'id', ''), 'pipeline_id': pipeline_id, 'name': name, 'title': getattr(rec, 'title', '') or name, 'escalation_path': _parse_path(getattr(rec, 'escalation_path', '')), } async def raise_problem(pipeline_id, problem_type, from_role, question, task_id=None, tenant_id=None, context=None, from_agentid=None, first_handler=None, first_handler_agentid=None, suspend_task=True) -> str: """提出问题。按问题类型的冒泡路径路由到第一个处理方。 Args: pipeline_id: 产线 id(为空则从 tenant_id/project 解析) problem_type: 问题类型标识(pipeline_problem_types.name) from_role: 谁提的(角色名) from_agentid: 提问的具体 agent 标识(多 agent 场景必传) first_handler: 可选,指定第一个处理角色(如审核退回指定被退角色) first_handler_agentid: 可选,第一个处理方具体 agent suspend_task: 是否把任务置 waiting(角色 agent 提问时挂起任务等回答) Returns: question_id """ db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: if not pipeline_id: pipeline_id = await _resolve_pipeline_id(sor, tenant_id) pt = await get_problem_type(sor, pipeline_id, problem_type) path = pt['escalation_path'] if pt else [] if first_handler and first_handler not in [p["role"] for p in path]: path = [{"role": first_handler, "agentid": first_handler_agentid or ""}] + path if not path: path = [{"role": from_role or "main_agent", "agentid": from_agentid or ""}] qid = getID() ctx_json = json.dumps(context, ensure_ascii=False, default=str) if context else None await sor.C('pipeline_agent_questions', { 'id': qid, 'tenant_id': tenant_id or '', 'task_id': task_id or '', 'from_role': from_role or '', 'from_agentid': from_agentid or '', 'problem_type': problem_type or '', 'question': question or '', 'context': ctx_json, 'escalation_path': json.dumps(path, ensure_ascii=False), 'escalation_pos': 0, 'status': S_PENDING, }) if suspend_task and task_id: # 挂起任务并清 claimed_by(否则 poller 要求 claimed_by IS NULL 才重新认领,任务僵尸) await sor.sqlExe( "UPDATE pipeline_tasks SET state='waiting', claimed_by=NULL WHERE id=${tid}$", {"tid": task_id}) logger.info("raise_problem: pipeline=%s type=%s from=%s/%s qid=%s path=%s", pipeline_id, problem_type, from_role, from_agentid, qid, path) return qid async def resolve_problem(question_id, answer, answered_by='', answer_source='main_agent', resume_task=True): """当前处理方解决了问题 → 停止冒泡(status=answered),任务恢复 submitted。""" db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: recs = await sor.R('pipeline_agent_questions', {'id': question_id}) if not recs: return None rec = recs[0] task_id = getattr(rec, 'task_id', '') or '' await sor.U('pipeline_agent_questions', { 'id': question_id, 'answer': answer or '', 'answered_by': answered_by or '', 'answer_source': answer_source or 'main_agent', 'status': S_ANSWERED, }) resumed = False if resume_task and task_id: trecs = await sor.R('pipeline_tasks', {'id': task_id}) if trecs: tstate = getattr(trecs[0], 'state', '') if tstate == 'waiting': await sor.sqlExe( "UPDATE pipeline_tasks SET state='submitted', claimed_by=NULL " "WHERE id=${tid}$", {"tid": task_id}) resumed = True logger.info("resolve_problem: qid=%s task=%s by=%s resumed=%s", question_id, task_id, answered_by, resumed) return {'id': question_id, 'status': S_ANSWERED, 'resumed': resumed} async def escalate_problem(question_id, by_role='', by_agentid='', note=''): """当前处理方没解决 → 冒泡到路径下一个处理方(pos+1)。""" db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: recs = await sor.R('pipeline_agent_questions', {'id': question_id}) if not recs: return None rec = recs[0] path = _parse_path(getattr(rec, 'escalation_path', '')) try: pos = int(getattr(rec, 'escalation_pos', 0) or 0) except (TypeError, ValueError): pos = 0 new_pos = pos + 1 if path and new_pos >= len(path): new_pos = len(path) - 1 # 尽头:停在最后一个处理方(通常是"人"兜底) elif not path: new_pos = 0 await sor.U('pipeline_agent_questions', { 'id': question_id, 'escalation_pos': new_pos, }) handler = path[new_pos] if path else {} logger.info("escalate_problem: qid=%s by=%s/%s -> %s (pos %d)", question_id, by_role, by_agentid, handler, new_pos) return {'id': question_id, 'escalation_pos': new_pos, 'current_handler': handler} async def forward_to_customer(question_id, by_role='main_agent', by_agentid=''): """便捷封装:主 agent 答不了,直接冒泡给客户(人)。等价于 escalate。""" return await escalate_problem(question_id, by_role=by_role, by_agentid=by_agentid) async def list_problems_for(role, agentid=None, tenant_id=None, task_id=None, limit=50) -> list: """列出「当前该 (role, agentid) 处理」的 pending 问题。 多 agent 感知:agentid 指定时,只返回"该 agent 处理"或"该角色任意 agent 处理" 的问题;不指定 agentid 时返回该角色所有待处理问题。 替代旧版「按 status 无差别列全部 pending」——把归属不同的问题混给同一个处理者。 """ db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: conditions = ["status='pending'"] params = {} if tenant_id: conditions.append("tenant_id=${tenant_id}$") params["tenant_id"] = tenant_id if task_id: conditions.append("task_id=${task_id}$") params["task_id"] = task_id where = " AND ".join(conditions) try: limit = int(limit) except (TypeError, ValueError): limit = 50 sql = f"SELECT * FROM pipeline_agent_questions WHERE {where} ORDER BY created_at DESC LIMIT {limit}" recs = await sor.sqlExe(sql, params) # 释放 SELECT 元数据锁 await sor.sqlExe("COMMIT", {}) result = [] for rec in (recs or []): h = current_handler_of(rec) if not _handler_matches(h, role, agentid): continue d = _rec_to_dict(rec) d['current_handler'] = h result.append(d) return result def _rec_to_dict(rec): if hasattr(rec, '__dict__'): return {k: getattr(rec, k) for k in dir(rec) if not k.startswith('_') and not callable(getattr(rec, k))} return dict(rec) async def get_task_qa(task_id, include_pending_for=None, include_pending_agentid=None) -> dict: """取任务的历史问答(已回答)+ 可选「待某处理方处理」的 pending 退回意见。 Args: include_pending_for: 角色名,额外返回当前该角色应响应的 pending 问题 include_pending_agentid: 具体 agent 标识(多 agent 场景) Returns: {"answered": [...], "pending": [...]} """ db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: recs = await sor.sqlExe( "SELECT * FROM pipeline_agent_questions " "WHERE task_id=${tid}$ AND status='answered' ORDER BY created_at ASC", {"tid": task_id}) answered = [] for rec in (recs or []): answered.append({ "question": getattr(rec, 'question', '') or '', "answer": getattr(rec, 'answer', '') or '', "answer_source": getattr(rec, 'answer_source', '') or '', "from_role": getattr(rec, 'from_role', '') or '', }) pending = [] if include_pending_for: precs = await sor.sqlExe( "SELECT * FROM pipeline_agent_questions " "WHERE task_id=${tid}$ AND status='pending' ORDER BY created_at ASC", {"tid": task_id}) for rec in (precs or []): h = current_handler_of(rec) if _handler_matches(h, include_pending_for, include_pending_agentid): pending.append({ "question": getattr(rec, 'question', '') or '', "from_role": getattr(rec, 'from_role', '') or '', "problem_type": getattr(rec, 'problem_type', '') or '', "current_handler": h, }) return {"answered": answered, "pending": pending}