# pbl_evidence/api/pbl_evidence_collect_from_events.dspy # M5a 主能力:批量从 pbl_runtime_event(M11b 产物)采集学习证据,幂等落 pbl_evidence。 # 铁律:dspy 内不 import、不 new ServerEnv();ServerEnv 注册的契约函数是直接全局;必须显式 return。 debug(f'pbl_evidence_collect_from_events.dspy: START params_kw={dict(params_kw)}') try: tenant = params_kw.get('tenant_id') or (await get_userorgid()) or '0' def _as_list(val): if val is None or val == '': return None if isinstance(val, (list, tuple)): return [str(x).strip() for x in val if str(x).strip()] return [x.strip() for x in str(val).split(',') if x.strip()] res = await pbl_evidence_collect_from_events( tenant_id=tenant, since=params_kw.get('since'), until=params_kw.get('until'), session_id=params_kw.get('session_id'), learner_id=params_kw.get('learner_id'), blueprint_id=params_kw.get('blueprint_id'), event_types=_as_list(params_kw.get('event_types')), evidence_types=_as_list(params_kw.get('evidence_types')), limit=params_kw.get('limit') or 500, dry_run=params_kw.get('dry_run'), update_existing=params_kw.get('update_existing'), ) status = 'OK' if res.get('ok') else 'ERROR' debug(f'pbl_evidence_collect_from_events.dspy: DONE status={status} msg={res.get("message")}') return {'status': status, 'data': {k: v for k, v in res.items() if k != 'items'}, 'items': res.get('items') or [], 'message': res.get('message') or ''} except Exception as e: error(f'pbl_evidence_collect_from_events.dspy: FAIL {format_exc()}') return {'status': 'ERROR', 'data': None, 'message': str(e)}