# todo_form_submit.dspy - 待办动态表单提交:校验 + 文件落项目根 + 按 schema.target 移到目标位置 # + 把「回答文本 + 已上传文件路径」写入回答,关闭问题、任务恢复执行。 # # 入参(JSON body 或 params_kw): # question_id 待办问题 id # fields JSON 对象:{字段名: 值}。文件字段值 = base64 字符串(或 data: 前缀); # 文本字段值 = 字符串。 # # 文件处理约定(用户确认): # 统一先落到项目根目录,再按 form_schema 里该字段的 target 移动: # - target 形如 "env/test.json"(带文件名)→ 移到 {root}/env/test.json # - target 形如 "env/"(目录,/ 结尾)→ 移到 {root}/env/{原文件名} # - target 为空 → 留在项目根 # 目标路径相对项目根解析,禁绝对路径与 .. 穿越(服务端二次校验)。 import os import json as _json import base64 as _b64 user_id = await get_user() if not user_id: return _json.dumps({"success": False, "error": "未登录"}, ensure_ascii=False) question_id = ((params_kw or {}).get('question_id') or '').strip() fields_raw = (params_kw or {}).get('fields') or '' if isinstance(fields_raw, str): try: fields = _json.loads(fields_raw) if fields_raw.strip() else {} except Exception: fields = {} else: fields = fields_raw or {} if not isinstance(fields, dict): fields = {} if not question_id: return _json.dumps({"success": False, "error": "缺少 question_id"}, ensure_ascii=False) dbname = get_module_dbname('pipeline-sdlc') # ── 读取待办问题 + form_schema ── async with DBPools().sqlorContext(dbname) as sor: qrecs = await sor.sqlExe( "SELECT id, tenant_id, task_id, context, status FROM pipeline_agent_questions WHERE id=${i}$", {"i": question_id}) await sor.sqlExe("COMMIT", {}) if not qrecs: return _json.dumps({"success": False, "error": "待办不存在或已处理"}, ensure_ascii=False) q = qrecs[0] project_id = str(getattr(q, 'tenant_id', '') or '') q_status = str(getattr(q, 'status', '') or '') if q_status not in ('', 'pending', 'waiting'): return _json.dumps({"success": False, "error": "该待办已处理,无需重复提交"}, ensure_ascii=False) schema = None try: ctx = getattr(q, 'context', '') or '' if ctx: schema = (_json.loads(ctx) or {}).get('form_schema') except Exception: schema = None schema_fields = (schema or {}).get('fields') or [] # ── 项目根目录(与角色 agent 读的位置一致) ── if not project_id: return _json.dumps({"success": False, "error": "待办缺少项目上下文,无法定位工作空间"}, ensure_ascii=False) async with DBPools().sqlorContext(dbname) as sor: project_dir, _base = await get_project_dir_by_id(sor, project_id) await sor.sqlExe("COMMIT", {}) if not project_dir or not os.path.isdir(project_dir): return _json.dumps({"success": False, "error": "项目工作空间不可用"}, ensure_ascii=False) def _safe_rel(target): """校验相对路径:禁绝对路径、禁 .. 穿越。返回规整后的相对路径或 ''。""" t = (target or '').strip() if not t: return '' if t.startswith('/'): return '' parts = [p for p in t.split('/') if p not in ('', '.')] if '..' in parts: return '' return '/'.join(parts) uploaded = [] # [(name, rel_path)] fatal = [] # 致命错误:必填缺失/解码失败/写盘失败 → 整体拒绝提交 warnings = [] # 非致命:文件已落项目根但移动到 target 失败 for sf in schema_fields: name = str(sf.get('name') or '') ftype = str(sf.get('type') or 'text').lower() required = bool(sf.get('required')) val = fields.get(name) if ftype == 'file': if not val: if required: fatal.append(str(sf.get('label') or name) + ':必须上传文件') continue # 解码 base64 try: if isinstance(val, str) and val.startswith('data:'): _, b64 = val.split(',', 1) else: b64 = val data = _b64.b64decode(b64) except Exception as e: fatal.append(str(sf.get('label') or name) + ':文件解码失败 ' + str(e)[:80]) continue # 文件名:优先取随字段附带的 {name}_filename,否则用字段名 fname = str(fields.get(name + '_filename') or '').strip() fname = os.path.basename(fname) if fname else '' if not fname: fname = name # 1) 先落项目根 root_path = os.path.join(project_dir, fname) try: with open(root_path, 'wb') as f: f.write(data) except Exception as e: fatal.append(str(sf.get('label') or name) + ':写入失败 ' + str(e)[:80]) continue # 2) 按 target 移动(失败不致命——文件已在项目根,agent 仍可读到) rel = _safe_rel(sf.get('target')) final_rel = fname try: if rel: if '/' in rel and not rel.endswith('/'): # target 含文件名(如 env/test.json)→ 直接用 dst_rel = rel else: # target 是目录(如 env/)→ 保留原文件名 dst_rel = rel.rstrip('/') + '/' + fname dst_abs = os.path.join(project_dir, dst_rel) os.makedirs(os.path.dirname(dst_abs), exist_ok=True) if os.path.exists(dst_abs): os.remove(dst_abs) # 覆盖(补齐/替换场景) os.rename(root_path, dst_abs) final_rel = dst_rel except Exception as e: warnings.append(str(sf.get('label') or name) + ':已保存到项目根,移动到 ' + rel + ' 失败 ' + str(e)[:80]) uploaded.append((name, final_rel)) else: # text / textarea sv = '' if val is None else str(val).strip() if not sv and required: fatal.append(str(sf.get('label') or name) + ':不能为空') # 致命错误 → 拒绝提交(文件可能已落盘,但回答不写,待办不关闭,让用户补齐重来) if fatal: return _json.dumps({"success": False, "error": ";".join(fatal)}, ensure_ascii=False) # ── 构造回答:文本字段 + 已上传文件路径(回流给 agent 的 _build_qna_section) ── answer_lines = [] for sf in schema_fields: if str(sf.get('type') or 'text').lower() == 'file': continue name = str(sf.get('name') or '') sv = '' if fields.get(name) is None else str(fields.get(name)).strip() if sv: answer_lines.append((sf.get('label') or name) + ':' + sv) if uploaded: answer_lines.append('已上传文件(位于项目工作空间,可直接 read_file):') for name, rel in uploaded: answer_lines.append(' · ' + rel) free_text = str(fields.get('_note') or '').strip() if free_text: answer_lines.append('补充说明:' + free_text) answer = '\n'.join(answer_lines) or '(已按要求上传/填写)' res = await question_answer(question_id, answer, answered_by=user_id, answer_source='owner.superuser') ok = False msg = '' if isinstance(res, tuple) or isinstance(res, list): ok = bool(res[0]) if len(res) > 0 else False msg = str(res[1]) if len(res) > 1 else '' elif isinstance(res, dict): ok = bool(res.get('success', res.get('ok', True))) msg = str(res.get('message') or res.get('error') or '') else: ok = bool(res) if ok: summary = '提交成功。' if uploaded: summary += '已上传 ' + '、'.join(rel for _, rel in uploaded) + ',任务恢复执行。' else: summary += '任务恢复执行。' if warnings: summary += '(部分提示:' + ';'.join(warnings) + ')' return _json.dumps({"success": True, "message": summary}, ensure_ascii=False) return _json.dumps({"success": False, "error": msg or "提交失败"}, ensure_ascii=False)