199 lines
6.8 KiB
Python
199 lines
6.8 KiB
Python
"""商机产线公共层:DB 上下文、爬虫平台 HTTP 客户端、报告/审批状态、人工任务。
|
||
|
||
所有 opp_* capability 模块共用本模块,禁止各自重复实现。
|
||
|
||
爬虫平台接入:内网 HTTP(http_api.py),配置优先级
|
||
1. appbase params 表: tender_api_base / tender_api_token(系统级配置禁硬编码)
|
||
2. 兜底默认: http://192.168.16.2:9085(内网地址,仅默认值)
|
||
"""
|
||
|
||
import json
|
||
import logging
|
||
import urllib.parse
|
||
|
||
from sqlor.dbpools import DBPools
|
||
from appPublic.uniqueID import getID
|
||
|
||
DBNAME = "pipeline"
|
||
PIPELINE_ID = "opportunity_general"
|
||
|
||
logger = logging.getLogger("pipeline.opportunity")
|
||
|
||
# 爬虫平台默认接入点(内网网卡地址,公网不可达);可被 params 表覆盖
|
||
DEFAULT_API_BASE = "http://192.168.16.2:9085"
|
||
|
||
# ── 报告状态 ──
|
||
RP_DRAFT = "draft" # 草稿(agent 编写中/待审阅)
|
||
RP_CONFIRMED = "confirmed" # 人工确认通过(门禁)
|
||
RP_APPROVAL_INITIATED = "approval_initiated" # 已发起研发审批
|
||
RP_APPROVED = "approved" # 审批通过
|
||
RP_REJECTED = "rejected" # 审批驳回
|
||
|
||
# ── 审批状态 ──
|
||
AP_INITIATED = "initiated"
|
||
AP_APPROVED = "approved"
|
||
AP_REJECTED = "rejected"
|
||
|
||
# ── 人工任务类型 ──
|
||
HT_REPORT_CONFIRM = "opp_report_confirm" # 报告人工确认门禁
|
||
HT_DEV_APPROVAL = "opp_dev_approval" # 研发审批人工决策
|
||
|
||
# 报告状态流转(合法迁移)
|
||
REPORT_TRANSITIONS = {
|
||
RP_DRAFT: (RP_CONFIRMED,),
|
||
RP_CONFIRMED: (RP_APPROVAL_INITIATED,),
|
||
RP_APPROVAL_INITIATED: (RP_APPROVED, RP_REJECTED),
|
||
RP_REJECTED: (RP_DRAFT,), # 驳回可改回草稿重来
|
||
}
|
||
|
||
|
||
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 rec_to_dict(rec):
|
||
"""sqlor 行对象 → dict(过滤方法名)。"""
|
||
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 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 表配置 + 默认兜底(系统级配置禁硬编码)。"""
|
||
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 default
|
||
|
||
|
||
async def get_crawler_config(sor):
|
||
"""爬虫平台接入配置 (base, token)。params 表优先,默认内网地址兜底。"""
|
||
base = await get_param(sor, "tender_api_base", DEFAULT_API_BASE)
|
||
token = await get_param(sor, "tender_api_token", "")
|
||
return (base or DEFAULT_API_BASE).rstrip("/"), token
|
||
|
||
|
||
async def crawler_get(sor, path, params=None, timeout=20):
|
||
"""调爬虫平台只读接口。返回 (ok, payload|错误信息)。"""
|
||
base, token = await get_crawler_config(sor)
|
||
qs = ""
|
||
if params:
|
||
clean = {k: v for k, v in params.items() if v not in (None, "")}
|
||
if clean:
|
||
qs = "?" + urllib.parse.urlencode(clean)
|
||
url = base + path + qs
|
||
headers = {"X-API-Token": token} if token else {}
|
||
try:
|
||
import httpx
|
||
except ImportError:
|
||
httpx = None
|
||
try:
|
||
if httpx is not None:
|
||
async with httpx.AsyncClient(timeout=timeout) as cli:
|
||
r = await cli.get(url, headers=headers)
|
||
if r.status_code != 200:
|
||
return False, "爬虫平台返回 %d: %s" % (r.status_code, r.text[:200])
|
||
return True, r.json()
|
||
# 无 httpx 时退回 urllib(同步阻塞,仅在缺依赖时用)
|
||
import asyncio
|
||
import urllib.request
|
||
|
||
def _blocking():
|
||
req = urllib.request.Request(url, headers=headers)
|
||
with urllib.request.urlopen(req, timeout=timeout) as resp:
|
||
return json.loads(resp.read().decode("utf-8"))
|
||
|
||
return True, await asyncio.get_event_loop().run_in_executor(None, _blocking)
|
||
except Exception as e:
|
||
logger.warning("crawler_get %s failed: %s", path, str(e)[:200])
|
||
return False, "爬虫平台调用失败: %s" % str(e)[:200]
|
||
|
||
|
||
async def find_human_task(sor, project_id, task_type, status=None):
|
||
"""查项目下指定类型的人工任务(供门禁判断与幂等创建)。"""
|
||
sql = ("SELECT id, title, status, task_type FROM pipeline_human_tasks "
|
||
"WHERE project_id=${pid}$ AND task_type=${tt}$")
|
||
args = {"pid": project_id, "tt": task_type}
|
||
if status:
|
||
sql += " AND status=${st}$"
|
||
args["st"] = status
|
||
sql += " ORDER BY created_at DESC LIMIT 1"
|
||
recs = await sor.sqlExe(sql, args)
|
||
await sor.sqlExe("COMMIT", {})
|
||
return rec_to_dict(recs[0]) if recs else None
|
||
|
||
|
||
async def create_human_task(sor, project_id, task_type, title, description="",
|
||
assignee_role="owner.superuser", assignee_id=""):
|
||
"""创建项目级人工任务(确认/审批门禁)。幂等由调用方先 find 保证。
|
||
|
||
pipeline_human_tasks NOT NULL 列:id/task_id/step_name/version/task_type/status。
|
||
task_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 "",
|
||
"status": "pending",
|
||
"project_id": project_id or "",
|
||
"title": (title or "")[:480],
|
||
"description": description or "",
|
||
})
|
||
await sor.sqlExe("COMMIT", {})
|
||
logger.info("opp create_human_task type=%s id=%s project=%s",
|
||
task_type, hid, project_id)
|
||
return hid
|
||
|
||
|
||
async def get_report(sor, report_id):
|
||
recs = await sor.sqlExe(
|
||
"SELECT * FROM opp_reports WHERE id=${i}$", {"i": report_id})
|
||
await sor.sqlExe("COMMIT", {})
|
||
return rec_to_dict(recs[0]) if recs else None
|
||
|
||
|
||
async def can_transition(cur_status, new_status):
|
||
return new_status in REPORT_TRANSITIONS.get(cur_status, ())
|