- 12张 bid_* 表(models+DDL+CRUD):招标信息/邮箱配置/人员/招标文件/评分项/资质/ 投标文件要求/章节/评审/评分/知识库/标书 - 7个角色能力模块(bid_*_capability):采集摘要/解析/知识库资质/编写/评审/合成/评分 - bid_flow 状态对账器:角色任务唯一创建者+阻塞门禁(源头守门,引擎零特判) - bid_mail:阿里云企业邮 IMAP 采集(凭据存DB AES,客户端专用密码兼容) - 能力包注册(bid_ability)+ 角色工具自注册(role_tool_schemas) - 技能:6 common 状态机/规范 + 9 角色技能(pipelines/bidding_general,随 pipeline-core 部署) - i18n(zh/en)、init/data.json(appcodes+pipelines记录)、build.sh、RBAC load_path
334 lines
12 KiB
Python
334 lines
12 KiB
Python
"""投标产线公共层:DB 上下文、角色规范化、门限参数、LLM 调用、项目/招标信息取数。
|
||
|
||
所有 bid_* capability 模块共用本模块,禁止各自重复实现(dspy/capability 都只调这里)。
|
||
"""
|
||
|
||
import json
|
||
import logging
|
||
|
||
from sqlor.dbpools import DBPools
|
||
from appPublic.uniqueID import getID
|
||
|
||
DBNAME = "pipeline"
|
||
PIPELINE_ID = "bidding_general"
|
||
|
||
logger = logging.getLogger("pipeline.bidding")
|
||
|
||
# ── 章节 / 标书状态 ──
|
||
CH_PENDING = "pending"
|
||
CH_WRITING = "writing"
|
||
CH_WRITTEN = "written"
|
||
CH_REVIEWING = "reviewing"
|
||
CH_APPROVED = "approved"
|
||
CH_REJECTED = "rejected"
|
||
|
||
# ── 招标信息状态 ──
|
||
TD_NEW = "new"
|
||
TD_SUMMARIZED = "summarized"
|
||
TD_PENDING_APPROVAL = "pending_approval"
|
||
TD_APPROVED = "approved"
|
||
TD_REJECTED = "rejected"
|
||
TD_ARCHIVED = "archived"
|
||
|
||
# ── 标书状态 ──
|
||
DOC_DRAFT = "draft"
|
||
DOC_REVIEWING = "reviewing"
|
||
DOC_PASSED = "passed"
|
||
DOC_REJECTED = "rejected"
|
||
DOC_DELIVERED = "delivered"
|
||
|
||
# ── 人类任务类型(阻塞门禁)──
|
||
HT_TENDER_APPROVAL = "tender_approval"
|
||
HT_TENDER_FILE = "tender_file_upload"
|
||
HT_HUMAN_DOCS = "human_docs_request"
|
||
HT_DELIVERY_CONFIRM = "bid_delivery_confirm"
|
||
|
||
BLOCKING_HT_TYPES = (HT_TENDER_FILE, HT_HUMAN_DOCS)
|
||
|
||
# 任务未终结状态(判断「是否已有在办任务」用)
|
||
OPEN_TASK_STATES = ("submitted", "running", "review", "qc_review")
|
||
|
||
# ── 门限默认值(可用 appbase params 表覆盖)──
|
||
DEFAULT_PARAMS = {
|
||
"bid_chapter_pass_ratio": "0.8", # 章节评审通过门限
|
||
"bid_pass_ratio": "0.85", # 整书评分通过门限
|
||
"bid_max_revise": "3", # 单章节最大重做次数
|
||
"bid_max_round": "3", # 整书最大评分轮次
|
||
}
|
||
|
||
|
||
def get_db():
|
||
"""取 DBPools(databases 为空时从 config 兜底注入,与引擎其他模块一致)。"""
|
||
db = DBPools()
|
||
if not db.databases:
|
||
from appPublic.jsonConfig import getConfig
|
||
config = getConfig()
|
||
if config and config.databases:
|
||
db.databases = config.databases
|
||
return db, DBNAME
|
||
|
||
|
||
def new_id():
|
||
return getID()
|
||
|
||
|
||
def normalize_role(role):
|
||
"""agent 角色补 agent. 前缀;人角色 {orgtype}.{role} 原样返回。"""
|
||
role = (role or "").strip()
|
||
if not role:
|
||
return ""
|
||
if "." in role:
|
||
return role
|
||
return "agent." + role
|
||
|
||
|
||
def rec_to_dict(rec):
|
||
"""sqlor 行对象 → dict(过滤方法名,避免 DictObject 序列化泄漏内置方法)。"""
|
||
if rec is None:
|
||
return {}
|
||
if isinstance(rec, dict):
|
||
return {k: v for k, v in rec.items() if not callable(v)}
|
||
d = {}
|
||
for k, v in vars(rec).items():
|
||
if k.startswith("_") or callable(v):
|
||
continue
|
||
d[k] = v
|
||
return d
|
||
|
||
|
||
def rows_to_dicts(recs, limit=200):
|
||
return [rec_to_dict(r) for r in (recs or [])[:limit]]
|
||
|
||
|
||
def to_float(v, default=0.0):
|
||
try:
|
||
if v is None or v == "":
|
||
return default
|
||
return float(v)
|
||
except (TypeError, ValueError):
|
||
return default
|
||
|
||
|
||
def to_int(v, default=0):
|
||
try:
|
||
if v is None or v == "":
|
||
return default
|
||
return int(float(v))
|
||
except (TypeError, ValueError):
|
||
return default
|
||
|
||
|
||
def json_loads(s, default=None):
|
||
if not s:
|
||
return default if default is not None else []
|
||
if isinstance(s, (list, dict)):
|
||
return s
|
||
try:
|
||
return json.loads(s)
|
||
except (json.JSONDecodeError, TypeError, ValueError):
|
||
return default if default is not None else []
|
||
|
||
|
||
async def get_param(sor, key, default=""):
|
||
"""读 appbase params 表配置(复用引擎 workspace.get_param,不重复实现)+ 默认兜底。"""
|
||
try:
|
||
from pipeline_service.workspace import get_param as _gp
|
||
v = await _gp(sor, key, "")
|
||
if v not in (None, ""):
|
||
return str(v)
|
||
except Exception:
|
||
pass
|
||
return str(DEFAULT_PARAMS.get(key, default))
|
||
|
||
|
||
async def resolve_project_id(project_id="", **entity_ids):
|
||
"""补齐 project_id:调用方(会话上下文)已给则原样返回;否则从实体反查项目。
|
||
|
||
角色工具经引擎分发时默认注入会话当前项目;为空时(无会话上下文)按实体反查,
|
||
避免「项目不明」导致跨项目读写。全部反查失败抛 ValueError。
|
||
"""
|
||
if project_id:
|
||
return project_id
|
||
db, dbname = get_db()
|
||
async with db.sqlorContext(dbname) as sor:
|
||
for tbl, col, val in (("bid_chapters", "chapter_id", entity_ids.get("chapter_id")),
|
||
("bid_reviews", "review_id", entity_ids.get("review_id")),
|
||
("bid_documents", "document_id", entity_ids.get("document_id")),
|
||
("bid_tender_files", "file_id", entity_ids.get("file_id"))):
|
||
if not val:
|
||
continue
|
||
recs = await sor.sqlExe(
|
||
"SELECT project_id FROM " + tbl + " WHERE id=${v}$ LIMIT 1", {"v": val})
|
||
await sor.sqlExe("COMMIT", {})
|
||
if recs and getattr(recs[0], "project_id", ""):
|
||
return getattr(recs[0], "project_id", "")
|
||
raise ValueError("缺少 project_id(请指定项目或切换到项目)")
|
||
|
||
|
||
async def get_thresholds(sor):
|
||
"""取四个门限参数(章节门限/整书门限/章节重做上限/整书轮次上限)。"""
|
||
return {
|
||
"chapter_pass_ratio": to_float(await get_param(sor, "bid_chapter_pass_ratio"), 0.8),
|
||
"bid_pass_ratio": to_float(await get_param(sor, "bid_pass_ratio"), 0.85),
|
||
"max_revise": to_int(await get_param(sor, "bid_max_revise"), 3),
|
||
"max_round": to_int(await get_param(sor, "bid_max_round"), 3),
|
||
}
|
||
|
||
|
||
async def get_project(sor, project_id):
|
||
"""取项目记录(dict),不存在返回 {}。"""
|
||
if not project_id:
|
||
return {}
|
||
recs = await sor.sqlExe(
|
||
"SELECT id, name, org_id, status, pipeline_id, created_by, default_model "
|
||
"FROM sd_projects WHERE id=${pid}$", {"pid": project_id})
|
||
await sor.sqlExe("COMMIT", {})
|
||
return rec_to_dict(recs[0]) if recs else {}
|
||
|
||
|
||
async def get_tender_by_project(sor, project_id):
|
||
"""取项目关联的招标信息(dict)。"""
|
||
if not project_id:
|
||
return {}
|
||
recs = await sor.sqlExe(
|
||
"SELECT * FROM bid_tenders WHERE project_id=${pid}$ ORDER BY created_at DESC LIMIT 1",
|
||
{"pid": project_id})
|
||
await sor.sqlExe("COMMIT", {})
|
||
return rec_to_dict(recs[0]) if recs else {}
|
||
|
||
|
||
async def create_role_task(sor, project_id, role, title, params=None, owner_id="bid_flow"):
|
||
"""创建角色任务(pipeline_tasks,tenant_id=project_id,pipeline_id='role_task')。
|
||
|
||
与开发产线共用同一套任务表和 poller —— 投标产线不另建任务机制。
|
||
"""
|
||
tid = new_id()
|
||
await sor.C("pipeline_tasks", {
|
||
"id": tid,
|
||
"tenant_id": project_id,
|
||
"pipeline_id": "role_task",
|
||
"owner_id": owner_id,
|
||
"title": (title or "")[:250],
|
||
"params": json.dumps(params or {}, ensure_ascii=False),
|
||
"role": normalize_role(role),
|
||
"state": "submitted",
|
||
"claimed_by": None,
|
||
})
|
||
logger.info("bid create_role_task role=%s task=%s project=%s", role, tid, project_id)
|
||
return tid
|
||
|
||
|
||
async def count_open_tasks(sor, project_id, role=None, chapter_id=None):
|
||
"""统计项目下未终结任务数(可按角色/章节过滤)。用于「已有在办任务就不重复派」。"""
|
||
states = "','".join(OPEN_TASK_STATES)
|
||
sql = ("SELECT id, params FROM pipeline_tasks WHERE tenant_id=${pid}$ "
|
||
"AND pipeline_id='role_task' AND state IN ('" + states + "')")
|
||
p = {"pid": project_id}
|
||
if role:
|
||
sql += " AND role=${role}$"
|
||
p["role"] = normalize_role(role)
|
||
recs = await sor.sqlExe(sql, p)
|
||
await sor.sqlExe("COMMIT", {})
|
||
if chapter_id is None:
|
||
return len(recs or [])
|
||
n = 0
|
||
for r in (recs or []):
|
||
prm = json_loads(getattr(r, "params", "") or "{}", {})
|
||
if isinstance(prm, dict) and prm.get("chapter_id") == chapter_id:
|
||
n += 1
|
||
return n
|
||
|
||
|
||
async def has_done_task(sor, project_id, role):
|
||
"""项目下该角色是否已有验收通过(approved/completed)的任务。"""
|
||
recs = await sor.sqlExe(
|
||
"SELECT id FROM pipeline_tasks WHERE tenant_id=${pid}$ AND pipeline_id='role_task' "
|
||
"AND role=${role}$ AND state IN ('approved','completed') LIMIT 1",
|
||
{"pid": project_id, "role": normalize_role(role)})
|
||
await sor.sqlExe("COMMIT", {})
|
||
return bool(recs)
|
||
|
||
|
||
async def has_pending_blocking_human_task(sor, project_id):
|
||
"""项目是否有阻塞中的人类任务(投标产线无迭代,按 project_id 判定)。"""
|
||
types = "','".join(BLOCKING_HT_TYPES + (HT_TENDER_APPROVAL,))
|
||
recs = await sor.sqlExe(
|
||
"SELECT COUNT(*) AS c FROM pipeline_human_tasks "
|
||
"WHERE project_id=${pid}$ AND status='pending' "
|
||
"AND task_type IN ('" + types + "')", {"pid": project_id})
|
||
await sor.sqlExe("COMMIT", {})
|
||
return to_int(getattr(recs[0], "c", 0) if recs else 0) > 0
|
||
|
||
|
||
async def create_human_task(sor, project_id, task_type, title, description="",
|
||
assignee_role="owner.superuser", assignee_id="",
|
||
form_schema=None):
|
||
"""创建人类任务(pipeline_human_tasks)。iteration_id 留空 → 走项目级门禁。"""
|
||
hid = new_id()
|
||
await sor.C("pipeline_human_tasks", {
|
||
"id": hid,
|
||
"task_id": "",
|
||
"step_name": task_type,
|
||
"version": 1,
|
||
"task_type": task_type,
|
||
"assignee_role": assignee_role or "owner.superuser",
|
||
"assignee_id": assignee_id or "",
|
||
"form_schema": json.dumps(form_schema or {}, ensure_ascii=False),
|
||
"status": "pending",
|
||
"project_id": project_id or "",
|
||
"title": (title or "")[:480],
|
||
"description": description or "",
|
||
})
|
||
logger.info("bid create_human_task type=%s id=%s project=%s", task_type, hid, project_id)
|
||
return hid
|
||
|
||
|
||
async def find_human_task(sor, project_id, task_type, status=None):
|
||
"""查项目下指定类型人类任务(最新一条)。"""
|
||
sql = ("SELECT * FROM pipeline_human_tasks WHERE project_id=${pid}$ "
|
||
"AND task_type=${tt}$")
|
||
p = {"pid": project_id, "tt": task_type}
|
||
if status:
|
||
sql += " AND status=${st}$"
|
||
p["st"] = status
|
||
sql += " ORDER BY created_at DESC LIMIT 1"
|
||
recs = await sor.sqlExe(sql, p)
|
||
await sor.sqlExe("COMMIT", {})
|
||
return rec_to_dict(recs[0]) if recs else {}
|
||
|
||
|
||
async def llm_json(prompt, model_name="", retries=1):
|
||
"""调 LLM 并解析 JSON(剥离 markdown code fence)。失败返回 (None, err)。"""
|
||
from pipeline_service.llm_bridge import llm_call
|
||
last_err = ""
|
||
for _ in range(max(1, retries + 1)):
|
||
try:
|
||
raw = await llm_call(prompt, model=model_name or None)
|
||
txt = (raw or "").strip()
|
||
if txt.startswith("```"):
|
||
parts = txt.split("\n", 1)
|
||
txt = parts[1] if len(parts) > 1 else txt
|
||
if "```" in txt:
|
||
txt = txt.rsplit("```", 1)[0]
|
||
txt = txt.strip()
|
||
# 容错:截取第一个 { 到最后一个 }
|
||
if not txt.startswith("{") and "{" in txt:
|
||
txt = txt[txt.find("{"):txt.rfind("}") + 1]
|
||
return json.loads(txt), ""
|
||
except Exception as e:
|
||
last_err = str(e)[:200]
|
||
return None, last_err or "LLM 返回无法解析为 JSON"
|
||
|
||
|
||
async def record(sor, project_id, entity, entity_id, action,
|
||
from_state=None, to_state=None, who=None, agent_id=None, detail=None):
|
||
"""写审计(复用引擎 audit,失败不影响主流程)。"""
|
||
try:
|
||
from pipeline_service.audit import record_audit
|
||
await record_audit(project_id, entity, entity_id, action,
|
||
from_state=from_state, to_state=to_state,
|
||
who=normalize_role(who), agent_id=agent_id,
|
||
detail=detail, sor=sor)
|
||
except Exception as e:
|
||
logger.warning("bid audit failed: %s", str(e)[:120])
|