| from __future__ import annotations |
|
|
| import os |
|
|
| from services.agent_firewall import server_management_refusal |
|
|
| |
| |
| |
|
|
| from services.context_compressor import ( |
| build_request_spec_from_messages, |
| extract_request_spec, |
| merge_request_spec, |
| render_request_spec, |
| ) |
|
|
| def _inject(context): |
| globals().update(context) |
|
|
| async def stream_chat(tid,prompt): |
| global last_llm_error |
| usage_user_id=thread_user_ids.get(tid,"") |
|
|
| pending_uploads=thread_upload_ids.pop(tid,[]) |
| upload_context="" |
|
|
| if usage_user_id and pending_uploads: |
| upload_context=await build_user_upload_context( |
| usage_user_id, |
| pending_uploads, |
| ) |
|
|
| if upload_context: |
| asyncio.create_task( |
| safe_memory_event( |
| usage_user_id, |
| "attachment_used", |
| { |
| "thread_id":tid, |
| "upload_ids":pending_uploads[:10], |
| "count":len(pending_uploads[:10]), |
| }, |
| ) |
| ) |
| processing_context="" |
|
|
| if ( |
| usage_user_id |
| and pending_uploads |
| and _quality_check_requested(prompt) |
| ): |
| task_id="proc_"+secrets.token_hex(8) |
|
|
| await safe_memory_event( |
| usage_user_id, |
| "processing_started", |
| { |
| "task_id":task_id, |
| "thread_id":tid, |
| "operation":"quality_check", |
| "upload_ids": |
| pending_uploads[:10], |
| "status":"running", |
| }, |
| ) |
|
|
| yield out( |
| "status", |
| { |
| "text": |
| "正在使用本地 Python 检查上传数据…" |
| }, |
| ) |
|
|
| try: |
| processing_result=await asyncio.to_thread( |
| _run_upload_quality_checks, |
| usage_user_id, |
| pending_uploads, |
| ) |
|
|
| record=_compact_processing_record( |
| processing_result |
| ) |
|
|
| await safe_memory_event( |
| usage_user_id, |
| "processing_completed", |
| { |
| "task_id":task_id, |
| "thread_id":tid, |
| "operation": |
| "quality_check", |
| "upload_ids": |
| pending_uploads[:10], |
| "status":"completed", |
| "result":record, |
| }, |
| ) |
|
|
| processing_context=( |
| "[USER_DATA_PROCESSING_RESULT]\n" |
| "The following result was computed " |
| "locally with Python from the user's " |
| "uploaded file. Treat these computed " |
| "values as authoritative for this " |
| "answer. Do not estimate them from " |
| "the raw file. Explain the result " |
| "clearly in Chinese.\n" |
| + json.dumps( |
| record, |
| ensure_ascii=False, |
| default=str, |
| ) |
| + "\n[/USER_DATA_PROCESSING_RESULT]" |
| ) |
|
|
| yield out( |
| "status", |
| { |
| "text": |
| "数据质检完成,正在整理结果…" |
| }, |
| ) |
|
|
| except Exception as exc: |
| await safe_memory_event( |
| usage_user_id, |
| "processing_failed", |
| { |
| "task_id":task_id, |
| "thread_id":tid, |
| "operation": |
| "quality_check", |
| "upload_ids": |
| pending_uploads[:10], |
| "status":"failed", |
| "error": |
| str(exc)[:500], |
| }, |
| ) |
|
|
| raise |
|
|
| use_fisheries = _is_fisheries_prompt(prompt) |
| use_ocean = _needs_ocean_mcp(prompt) |
| mcp_error = None |
| turn_error = "" |
| export_tool_completed=False |
| export_tool_results=[] |
| try: |
| hf_task = None |
| if use_fisheries: |
| yield out("status",{"text":"正在读取 Hugging Face 渔业数据…","stage":"source","step":2,"total_steps":5}) |
| hf_task = asyncio.create_task(build_hf_fisheries_context(prompt)) |
|
|
| |
| |
| |
| |
| |
| |
| |
| if use_ocean: |
| yield out("status",{"text":"Ocean 数据工具已就绪,正在处理请求…","stage":"source","step":2,"total_steps":5}) |
|
|
| grounded_prompt = prompt |
|
|
| if processing_context: |
| grounded_prompt += ( |
| "\n\n" + processing_context |
| ) |
| elif upload_context: |
| grounded_prompt += ( |
| "\n\n" + upload_context |
| ) |
| if hf_task is not None: |
| try: |
| hf_context = await hf_task |
| grounded_prompt = grounded_prompt + "\n\n" + hf_context |
|
|
| if usage_user_id: |
| asyncio.create_task( |
| safe_memory_event( |
| usage_user_id, |
| "fisheries_query", |
| { |
| "thread_id":tid, |
| "source":"huggingface-fisheries", |
| "status":"completed", |
| "prompt_chars":len(prompt.strip()), |
| }, |
| ) |
| ) |
| except Exception as exc: |
| grounded_prompt = ( |
| prompt |
| + "\n\n[HF_FISHERIES_LIVE_CONTEXT_ERROR]\n" |
| + str(exc) |
| + "\n[/HF_FISHERIES_LIVE_CONTEXT_ERROR]" |
| ) |
|
|
| try: |
| det=await rjson(f"/v1/threads/{tid}") |
| since=int(det.get("latest_seq") or 0) |
| tr=await rjson(f"/v1/threads/{tid}/turns",method="POST",body={ |
| "prompt":grounded_prompt, |
| "input_summary":prompt[:200], |
| "model":MODEL, |
| "mode":"agent", |
| "allow_shell":False, |
| "trust_mode":False, |
| "auto_approve":False, |
| }) |
| except Exception as exc: |
| if not _is_codewhale_missing_thread_error(exc): |
| raise |
| history_text=_format_recent_thread_history(tid) |
| if history_text: |
| grounded_prompt=( |
| "[PRIOR_CONVERSATION_HISTORY]\n" |
| + history_text |
| + "\n[/PRIOR_CONVERSATION_HISTORY]\n\n" |
| + "[CONVERSATION_RECOVERY_INSTRUCTION]\n" |
| + "用户在本次对话中已经描述过下载/查询要求。上面是同一用户的完整对话记录;" |
| + "请把最新用户消息当作该对话的后续补充或修改继续处理," |
| + "不要要求用户重新提交完整需求。\n" |
| + "[/CONVERSATION_RECOVERY_INSTRUCTION]\n\n" |
| + grounded_prompt |
| ) |
| old_tid=tid |
| tid=await recover_codewhale_thread(old_tid) |
| log.warning( |
| "CodeWhale thread expired and was rebuilt: " |
| "old=%s new=%s", |
| old_tid, |
| tid, |
| ) |
| yield out( |
| "status", |
| { |
| "text": |
| "检测到原对话会话已过期,正在自动保留上下文并恢复…", |
| "stage":"source", |
| "step":1, |
| "total_steps":5, |
| }, |
| ) |
| yield out( |
| "thread_recovered", |
| { |
| "thread_id":tid, |
| "old_thread_id":old_tid, |
| }, |
| ) |
| det=await rjson(f"/v1/threads/{tid}") |
| since=int(det.get("latest_seq") or 0) |
| tr=await rjson(f"/v1/threads/{tid}/turns",method="POST",body={ |
| "prompt":grounded_prompt, |
| "input_summary":prompt[:200], |
| "model":MODEL, |
| "mode":"agent", |
| "allow_shell":False, |
| "trust_mode":False, |
| "auto_approve":False, |
| }) |
| turn=((tr or {}).get("turn") or {}).get("id") |
| answer="" |
| yield out("status",{"text":"Codex 正在处理…","stage":"compose","step":4,"total_steps":5}) |
|
|
| async for rec in events(tid,since): |
| if turn and rec.get("turn_id") and rec["turn_id"]!=turn: |
| continue |
|
|
| e=rec.get("event") |
| p=pl(rec) |
|
|
| if e=="item.started": |
| tool=_tool_name_from_payload(p) |
|
|
| if tool.startswith("mcp_marine_"): |
| yield out("status",{"text":"正在查询学校 Ocean 数据服务器…","stage":"query","step":3,"total_steps":5}) |
|
|
| if usage_user_id: |
| asyncio.create_task( |
| safe_memory_event( |
| usage_user_id, |
| _marine_event_type(tool), |
| { |
| "thread_id":tid, |
| "tool":tool, |
| "operation":_marine_event_type(tool), |
| "source":_data_source_from_prompt(prompt), |
| "status":"started", |
| "prompt_chars":len(prompt.strip()), |
| }, |
| ) |
| ) |
|
|
| if e=="item.completed": |
| item=p.get("item") or {} |
| summary=str(item.get("summary") or "") |
| completed_tool=_tool_name_from_payload(p) |
| try: |
| completed_raw=json.dumps(p,ensure_ascii=False,default=str) |
| except Exception: |
| completed_raw=str(p) |
|
|
| if ( |
| completed_tool.endswith("_export") |
| or "mcp_marine_marine_export" in completed_raw |
| ): |
| export_tool_completed=True |
| export_tool_results.append(completed_raw) |
|
|
| if "MCP server 'marine' connected" in summary or "mcp_marine_" in summary: |
| marine_threads.add(tid) |
| |
| |
| recovered=_agent_text_from_payload(p) |
| if recovered and not answer: |
| answer=recovered |
|
|
| if e=="item.delta" and p.get("kind")=="agent_message": |
| d=p.get("delta") or "" |
| if d: |
| answer+=d |
|
|
| if e in {"item.failed","item.interrupted"}: |
| err=_error_from_payload(p) |
| if err: |
| turn_error=err |
|
|
| if e=="approval.required": |
| aid=p.get("approval_id") or p.get("id") |
| tool=( |
| p.get("tool_name") |
| or ((p.get("tool") or {}).get("name") if isinstance(p.get("tool"),dict) else p.get("tool")) |
| or ((p.get("item") or {}).get("tool_name") if isinstance(p.get("item"),dict) else "") |
| or "" |
| ) |
| if aid: |
| if tool.startswith("mcp_marine_"): |
| await approve(aid,"allow") |
| elif tool=="start_mcp_server": |
| |
| |
| await approve(aid,"deny") |
| yield out("status",{"text":"已阻止重复启动 Ocean MCP"}) |
| else: |
| await approve(aid,"deny") |
| yield out("status",{"text":"已保持安全数据访问模式"}) |
|
|
| if e=="turn.lifecycle": |
| st=_turn_status(p) |
| if st in {"failed","canceled","interrupted"}: |
| detail=_error_from_payload(p) or turn_error or f"Turn {st}" |
| raise RuntimeError(detail) |
|
|
| if e=="turn.completed": |
| st=_turn_status(p) |
| if st in {"failed","canceled","interrupted"}: |
| detail=_error_from_payload(p) or turn_error or f"Turn {st}" |
| raise RuntimeError(detail) |
| if not answer.strip(): |
| detail=_error_from_payload(p) or turn_error |
| if detail: |
| raise RuntimeError(detail) |
| raise RuntimeError( |
| "Codex 回合已结束,但 CodeWhale 没有产生 assistant 文本。" |
| ) |
| final_answer=_sanitize_final_answer(answer) |
| if not final_answer: |
| raise RuntimeError("模型返回内容在输出清理后为空。") |
|
|
| export_error=_ocean_export_execution_error( |
| prompt, |
| export_tool_completed=export_tool_completed, |
| tool_result_text="\n".join(export_tool_results), |
| final_answer=final_answer, |
| ) |
| if export_error: |
| raise RuntimeError(export_error) |
|
|
| if usage_user_id: |
| asyncio.create_task( |
| record_generated_download_assets( |
| usage_user_id, |
| tid, |
| prompt, |
| final_answer, |
| ) |
| ) |
|
|
| last_llm_error=None |
| yield out("token",{"text":final_answer}) |
| yield out("done",{"text":final_answer}) |
| return |
|
|
| raise RuntimeError(turn_error or "Runtime stream ended early") |
| except Exception as exc: |
| last_llm_error=str(exc) |
| log.exception( |
| "chat failed: thread=%s model=%s ocean=%s fisheries=%s", |
| tid, MODEL, use_ocean, use_fisheries, |
| ) |
| yield out("error",{"text":_public_error(str(exc)),"stage":"chat"}) |
|
|
|
|
| async def harness_stream_chat(tid,prompt): |
| global last_llm_error |
|
|
| usage_user_id=thread_user_ids.get(tid,"") |
| pending_uploads=thread_upload_ids.pop(tid,[]) |
|
|
| upload_context="" |
| processing_context="" |
|
|
| if usage_user_id and pending_uploads: |
| upload_context=await build_user_upload_context( |
| usage_user_id, |
| pending_uploads, |
| ) |
|
|
| asyncio.create_task( |
| safe_memory_event( |
| usage_user_id, |
| "attachment_used", |
| { |
| "thread_id":tid, |
| "upload_ids":pending_uploads[:10], |
| "count":len(pending_uploads[:10]), |
| "runtime":"codex-harness", |
| }, |
| ) |
| ) |
|
|
| if ( |
| usage_user_id |
| and pending_uploads |
| and _quality_check_requested(prompt) |
| ): |
| task_id="proc_"+secrets.token_hex(8) |
|
|
| await safe_memory_event( |
| usage_user_id, |
| "processing_started", |
| { |
| "task_id":task_id, |
| "thread_id":tid, |
| "operation":"quality_check", |
| "upload_ids":pending_uploads[:10], |
| "status":"running", |
| "runtime":"codex-harness", |
| }, |
| ) |
|
|
| yield out( |
| "status", |
| { |
| "text": |
| "正在使用本地 Python 检查数据…" |
| }, |
| ) |
|
|
| try: |
| processing_result=await asyncio.to_thread( |
| _run_upload_quality_checks, |
| usage_user_id, |
| pending_uploads, |
| ) |
|
|
| record=_compact_processing_record( |
| processing_result |
| ) |
|
|
| await safe_memory_event( |
| usage_user_id, |
| "processing_completed", |
| { |
| "task_id":task_id, |
| "thread_id":tid, |
| "operation":"quality_check", |
| "upload_ids":pending_uploads[:10], |
| "status":"completed", |
| "runtime":"codex-harness", |
| "result":record, |
| }, |
| ) |
|
|
| processing_context=( |
| "[USER_DATA_PROCESSING_RESULT]\n" |
| "These values were computed locally " |
| "with Python from the current user's " |
| "uploaded file. Use them as the " |
| "authoritative result. Do not guess " |
| "or recompute them mentally. Explain " |
| "the result clearly in Chinese.\n" |
| + json.dumps( |
| record, |
| ensure_ascii=False, |
| default=str, |
| ) |
| + "\n[/USER_DATA_PROCESSING_RESULT]" |
| ) |
|
|
| yield out( |
| "status", |
| { |
| "text": |
| "数据计算完成,Codex Harness 正在整理结果…" |
| }, |
| ) |
|
|
| except Exception as exc: |
| await safe_memory_event( |
| usage_user_id, |
| "processing_failed", |
| { |
| "task_id":task_id, |
| "thread_id":tid, |
| "operation":"quality_check", |
| "upload_ids":pending_uploads[:10], |
| "status":"failed", |
| "runtime":"codex-harness", |
| "error":str(exc)[:500], |
| }, |
| ) |
| raise |
|
|
| app_context=thread_system_prompts.get( |
| tid, |
| USER_SYSTEM, |
| ) |
|
|
| harness_prompt=( |
| "[APPLICATION_CONTEXT]\n" |
| + app_context |
| + "\n[/APPLICATION_CONTEXT]\n\n" |
| + "[CURRENT_USER_MESSAGE]\n" |
| + prompt |
| + "\n[/CURRENT_USER_MESSAGE]" |
| ) |
|
|
| if processing_context: |
| harness_prompt += ( |
| "\n\n" + processing_context |
| ) |
| elif upload_context: |
| harness_prompt += ( |
| "\n\n" + upload_context |
| ) |
|
|
| yield out( |
| "status", |
| { |
| "text": |
| f"Codex Harness · {HARNESS_MODEL} 正在处理…" |
| }, |
| ) |
|
|
| try: |
| task=asyncio.create_task( |
| asyncio.to_thread( |
| dsh.run, |
| harness_prompt, |
| session_id=tid, |
| ) |
| ) |
|
|
| while not task.done(): |
| try: |
| await asyncio.wait_for( |
| asyncio.shield(task), |
| timeout=8, |
| ) |
| except asyncio.TimeoutError: |
| yield ": keepalive\n\n" |
|
|
| result=await task |
|
|
| final_answer=_sanitize_final_answer( |
| result.final_response |
| ) |
|
|
| if not final_answer: |
| raise RuntimeError( |
| "Codex Harness 没有返回有效文本。" |
| ) |
|
|
| log.info( |
| "Codex Harness completed: " |
| "thread=%s model=%s reason=%s uploads=%s", |
| tid, |
| HARNESS_MODEL, |
| result.finish_reason, |
| len(pending_uploads), |
| ) |
|
|
| last_llm_error=None |
|
|
| yield out( |
| "token", |
| {"text":final_answer}, |
| ) |
|
|
| yield out( |
| "done", |
| { |
| "text":final_answer, |
| "runtime":"codex-harness", |
| "model":HARNESS_MODEL, |
| "finish_reason":result.finish_reason, |
| }, |
| ) |
|
|
| except Exception as exc: |
| last_llm_error=str(exc) |
|
|
| log.exception( |
| "Codex Harness failed: " |
| "thread=%s model=%s uploads=%s", |
| tid, |
| HARNESS_MODEL, |
| len(pending_uploads), |
| ) |
|
|
| yield out( |
| "error", |
| { |
| "text": |
| "Codex Harness 调用失败:" |
| + str(exc)[:500], |
| "stage":"harness", |
| }, |
| ) |
|
|
|
|
| DATA_REQUEST_TTL_SECONDS = 6 * 3600 |
|
|
|
|
| def _is_runtime_info_prompt(prompt: str) -> bool: |
| """Recognize the read-only deployment-info question locally.""" |
| text = str(prompt or "").strip().lower() |
| return ( |
| ("供应商" in text or "provider" in text) |
| and ("模型" in text or "model" in text) |
| and ("运行时" in text or "runtime" in text) |
| and not _needs_ocean_mcp(text) |
| and not _is_fisheries_prompt(text) |
| ) |
|
|
|
|
| def _is_marine_health_prompt(prompt: str) -> bool: |
| """Recognize a live Marine API status question without invoking the model.""" |
| text = str(prompt or "").strip().lower() |
| health_terms = ( |
| "海洋数据服务器连接", "检查海洋数据服务器", "数据服务器连接", |
| "三个数据域", "三个域是否在线", "ocean、tuna、squid", |
| "ocean, tuna, squid", "marine health", "marine server health", |
| ) |
| return any(term in text for term in health_terms) |
|
|
|
|
| def _is_fisheries_inventory_prompt(prompt: str) -> bool: |
| """Recognize file-list requests that do not need an agent turn. |
| |
| A repository inventory is a read-only server operation. Sending this |
| request through the Codex CLI is both slow and, in headless deployments, |
| vulnerable to the interactive-MCP approval limitation. Keep actual file |
| reads, analysis and exports on their normal guarded paths. |
| """ |
| text = str(prompt or "").strip().lower() |
| |
| |
| |
| explicit_inventory = ( |
| any(term in text for term in ("完整文件数", "总大小", "本地镜像", "回退读取")) |
| and any(term in text for term in ("squid", "鱿鱼", "柔鱼", "tuna", "金枪鱼")) |
| ) |
| if not _is_fisheries_prompt(text) and not explicit_inventory: |
| return False |
| inventory_terms = ( |
| "可用数据集", "当前可用", "文件来源", "文件清单", "列出", |
| "有哪些文件", "入库", "目录", "inventory", "catalog", |
| ) |
| if not any(term in text for term in inventory_terms) and not explicit_inventory: |
| return False |
| |
| |
| |
| metadata_text = text |
| for phrase in ( |
| "不读取文件内容", "不读文件内容", "不读取内容", "不读内容", |
| "不读取文件", "不读文件", "无需读取", "不要读取", |
| ): |
| metadata_text = metadata_text.replace(phrase, "") |
| content_terms = ( |
| "字段", "记录数", "缺失", "重复", "筛选", "汇总", "聚合", |
| "导出", "下载", "实际读取", "读取csv", "读取 csv", |
| ) |
| return not any(term in metadata_text for term in content_terms) |
|
|
|
|
| async def _direct_fisheries_inventory_answer(prompt: str) -> str: |
| """Return a live fisheries file inventory without invoking the LLM/MCP.""" |
| text = str(prompt or "").lower() |
| wants_tuna = any(term in text for term in ("tuna", "金枪鱼")) |
| wants_squid = any(term in text for term in ("squid", "鱿鱼", "柔鱼")) |
| domains = [] |
| if wants_squid or not wants_tuna: |
| domains.append("squid") |
| if wants_tuna: |
| domains.append("tuna") |
|
|
| repo_map = globals().get("HF_DATASET_REPOS") or {} |
| tree_reader = globals().get("hf_live_tree") |
| files_only = globals().get("_hf_live_files") |
| human_bytes = globals().get("_human_bytes") |
| local_tree = globals().get("local_dataset_tree") |
| if not callable(tree_reader) or not callable(files_only) or not callable(human_bytes): |
| return "渔业数据清单服务尚未就绪。" |
|
|
| source_terms = { |
| "GFW": ("gfw", "global_fishing_watch", "global-fishing-watch"), |
| "SPRFMO": ("sprfmo",), |
| "NPFC": ("npfc",), |
| "WCPFC": ("wcpfc",), |
| "IATTC": ("iattc",), |
| "ICCAT": ("iccat",), |
| "IOTC": ("iotc",), |
| "CCSBT": ("ccsbt",), |
| "FAO": ("fao",), |
| "Sea Around Us": ("sea_around", "sea-around", "sea around"), |
| "RAM Legacy": ("ram",), |
| "VIIRS": ("viirs",), |
| } |
| lines = [] |
| for domain in domains: |
| repo = str(repo_map.get(domain) or "").strip() |
| if not repo: |
| lines.append(f"- {domain.title()}:未配置数据仓库。") |
| continue |
| try: |
| raw_items = await tree_reader(repo) |
| files = files_only(raw_items) |
| except Exception as exc: |
| lines.append(f"- {domain.title()}:清单读取失败({str(exc)[:220]})。") |
| continue |
|
|
| total = sum(int(item.get("size_bytes") or 0) for item in files) |
| is_local = False |
| try: |
| is_local = bool(callable(local_tree) and local_tree(repo)) |
| except Exception: |
| pass |
| origin = "学校服务器本地镜像" if is_local else "Hugging Face main 分支" |
| lines.extend([ |
| f"## {domain.title()} 数据清单", |
| f"- 仓库:`{repo}`", |
| f"- 来源:{origin}", |
| f"- 可用文件:{len(files)} 个,约 {human_bytes(total)}", |
| ]) |
|
|
| groups = [] |
| for name, terms in source_terms.items(): |
| matched = [ |
| item for item in files |
| if any(term in str(item.get("path") or "").lower() for term in terms) |
| ] |
| if matched: |
| groups.append( |
| f"- {name}:{len(matched)} 个文件,约 " |
| f"{human_bytes(sum(int(item.get('size_bytes') or 0) for item in matched))}" |
| ) |
| if groups: |
| lines.append("- 按来源:") |
| lines.extend(groups) |
| else: |
| lines.append("- 按来源:当前文件路径未匹配到可识别来源标签。") |
|
|
| preview = files[:12] |
| if preview: |
| lines.append("- 文件示例:") |
| lines.extend( |
| f" - `{item.get('path')}`({human_bytes(item.get('size_bytes') or 0)})" |
| for item in preview |
| ) |
| lines.append("- 本次仅读取文件清单元数据,未打开或分析文件内容。") |
|
|
| return "\n".join(lines) if lines else "未找到可查询的渔业数据域。" |
|
|
|
|
| async def _direct_marine_health_answer(): |
| """Query the configured Marine API directly for a deterministic status answer.""" |
| httpx_module = globals().get("httpx") |
| if httpx_module is None: |
| try: |
| import httpx as httpx_module |
| except Exception as exc: |
| return f"海洋数据服务器连接检查失败:HTTP 客户端不可用({exc})" |
| base = str(globals().get("MARINE_API_URL") or "").rstrip("/") |
| if not base: |
| return "海洋数据服务器未配置 MARINE_API_URL。" |
|
|
| candidates = [f"{base}/health"] |
| |
| if not base.endswith("/marine"): |
| candidates.append(f"{base}/marine/health") |
|
|
| last_error = "" |
| async with httpx_module.AsyncClient(timeout=12.0, follow_redirects=True) as client: |
| for url in candidates: |
| try: |
| response = await client.get(url) |
| if response.is_success: |
| payload = response.json() |
| domains = payload.get("domains") if isinstance(payload, dict) else None |
| labels = { |
| "ocean": "在线", |
| "tuna": "在线", |
| "squid": "在线", |
| } |
| if isinstance(domains, dict): |
| for key in labels: |
| value = domains.get(key) |
| if isinstance(value, dict): |
| value = value.get("status") or value.get("online") |
| if value in (False, "offline", "未连接", "不可用"): |
| labels[key] = "离线/不可用" |
| elif isinstance(domains, list): |
| for key in labels: |
| if key not in {str(item).lower() for item in domains}: |
| labels[key] = "未在健康响应中列出" |
| return ( |
| "海洋数据服务器连接正常。\n" |
| f"- Ocean:{labels['ocean']}\n" |
| f"- Tuna:{labels['tuna']}\n" |
| f"- Squid:{labels['squid']}" |
| ) |
| last_error = f"HTTP {response.status_code}" |
| except Exception as exc: |
| last_error = str(exc) |
| return f"海洋数据服务器连接检查失败:{last_error[:240]}" |
|
|
|
|
| def _context_spec_has_data(spec) -> bool: |
| if not isinstance(spec, dict): |
| return False |
| return any( |
| spec.get(key) |
| for key in ( |
| "domain", "dataset", "variables", "region", |
| "depth", "date", "format", |
| ) |
| ) |
|
|
|
|
| async def dispatch_chat_stream(tid,prompt): |
| blocked = server_management_refusal(prompt) |
| if blocked: |
| yield out("token", {"text": blocked}) |
| yield out( |
| "done", |
| { |
| "text": blocked, |
| "runtime": "server-management-firewall", |
| "model": "none", |
| "finish_reason": "blocked", |
| }, |
| ) |
| return |
|
|
| |
| |
| |
| |
| if _is_runtime_info_prompt(prompt): |
| provider = str(globals().get("CODEX_PROVIDER_NAME") or os.environ.get("CODEX_PROVIDER", "ZAI")).upper() |
| model = str(globals().get("HARNESS_MODEL") or os.environ.get("CODEX_MODEL", "glm-5.3")) |
| runtime = "官方 Codex CLI Harness(native)" |
| answer = f"- 供应商:{provider}\n- 模型名称:{model}\n- 运行时:{runtime}" |
| yield out("token", {"text": answer}) |
| yield out("done", {"text": answer, "runtime": "local-fastpath", "model": model, "finish_reason": "completed"}) |
| return |
|
|
| if _is_marine_health_prompt(prompt): |
| answer = await _direct_marine_health_answer() |
| yield out("token", {"text": answer}) |
| yield out("done", {"text": answer, "runtime": "direct-marine-health", "model": "none", "finish_reason": "completed"}) |
| return |
|
|
| if _is_fisheries_inventory_prompt(prompt): |
| |
| |
| |
| log.info("direct fisheries inventory fast path: prompt_chars=%s", len(str(prompt or ""))) |
| try: |
| answer = await asyncio.wait_for( |
| _direct_fisheries_inventory_answer(prompt), |
| timeout=float(os.environ.get("FISHERIES_INVENTORY_TIMEOUT", "45")), |
| ) |
| except asyncio.TimeoutError: |
| answer = ( |
| "学校服务器的 Squid 清单查询超过 45 秒,已停止等待模型调用。" |
| "请稍后重试;本次没有读取文件内容,也没有生成文件。" |
| ) |
| yield out("token", {"text": answer}) |
| yield out( |
| "done", |
| { |
| "text": answer, |
| "runtime": "direct-fisheries-inventory", |
| "model": "none", |
| "finish_reason": "completed", |
| }, |
| ) |
| return |
|
|
| greeting = str(prompt or "").strip().lower() |
| |
| |
| greeting = greeting.replace("!", "!").replace("?", "?") |
| greeting_replies = { |
| "你好": "你好!有什么可以帮你的吗?", |
| "你好!": "你好!有什么可以帮你的吗?", |
| "你好?": "你好!有什么可以帮你的吗?", |
| "你好不好": "你好!我很好,谢谢关心。有什么可以帮你的吗?", |
| "你好吗": "你好!我很好,谢谢关心。有什么可以帮你的吗?", |
| "你好啊": "你好!有什么可以帮你的吗?", |
| "你好呀": "你好!有什么可以帮你的吗?", |
| "您好": "您好!有什么可以帮你的吗?", |
| "您好!": "您好!有什么可以帮你的吗?", |
| "您好?": "您好!有什么可以帮你的吗?", |
| "嗨": "你好!有什么可以帮你的吗?", |
| "hello": "Hello!有什么可以帮你的吗?", |
| "hi": "你好!有什么可以帮你的吗?", |
| "在吗": "在的,请告诉我你的海洋数据需求。", |
| } |
| if greeting in greeting_replies: |
| answer = greeting_replies[greeting] |
| yield out("token", {"text": answer}) |
| yield out( |
| "done", |
| { |
| "text": answer, |
| "runtime": "local-fastpath", |
| "model": "none", |
| "finish_reason": "completed", |
| }, |
| ) |
| return |
|
|
| now=time.time() |
| pending=thread_last_data_requests.get(tid) or {} |
| if pending and now-float(pending.get("ts") or 0)>DATA_REQUEST_TTL_SECONDS: |
| thread_last_data_requests.pop(tid, None) |
| pending={} |
|
|
| history_spec=build_request_spec_from_messages( |
| thread_recent_history.get(tid) or [] |
| ) |
| context_spec=thread_compressed_contexts.get(tid) or {} |
| if not _context_spec_has_data(context_spec): |
| context_spec=history_spec |
| if not _context_spec_has_data(context_spec) and pending.get("spec"): |
| context_spec=pending.get("spec") |
|
|
| explicit_data=( |
| _needs_ocean_mcp(prompt) |
| or _is_fisheries_prompt(prompt) |
| ) |
| is_confirmation=_is_confirmation_prompt(prompt) |
| is_amendment=bool( |
| (not is_confirmation) |
| and pending.get("prompt") |
| and _is_data_request_followup(prompt) |
| ) |
| resumed=bool( |
| pending.get("prompt") |
| and (is_confirmation or is_amendment) |
| ) |
|
|
| |
| |
| |
| |
| if is_amendment: |
| merged_spec=merge_request_spec( |
| context_spec, |
| prompt, |
| revision=True, |
| ) |
| thread_compressed_contexts[tid]=merged_spec |
| thread_last_data_requests[tid]={ |
| "prompt":render_request_spec(merged_spec), |
| "ts":now, |
| "spec":merged_spec, |
| } |
| elif explicit_data and not resumed: |
| parsed_spec=extract_request_spec(prompt) |
| if _context_spec_has_data(parsed_spec): |
| merged_spec=merge_request_spec( |
| None, |
| prompt, |
| revision=False, |
| ) |
| thread_compressed_contexts[tid]=merged_spec |
| else: |
| merged_spec={"raw_request":str(prompt or "")[:1600]} |
| thread_compressed_contexts[tid]=merged_spec |
| thread_last_data_requests[tid]={ |
| "prompt":render_request_spec(merged_spec) or prompt, |
| "ts":now, |
| "spec":merged_spec, |
| } |
|
|
| |
| |
| batch_spec=parse_ocean_batch_request(prompt) |
| if batch_spec: |
| if not resumed: |
| batch_spec_context=thread_compressed_contexts.get(tid) or {} |
| thread_last_data_requests[tid]={ |
| "prompt":( |
| render_request_spec(batch_spec_context) |
| or prompt |
| ), |
| "ts":now, |
| "spec":batch_spec_context, |
| } |
| uid=thread_user_ids.get(tid,"") |
| job=await OCEAN_BATCH_MANAGER.create(uid,prompt,batch_spec,tid) |
| yield out("status",{ |
| "text":f"已拆分为 {job['total']} 个单日子任务,后台并行执行中…", |
| "stage":"query","step":3,"total_steps":5, |
| }) |
| yield out("batch_job",job) |
| variables="、".join(batch_spec.get("variables") or []) |
| bbox=batch_spec.get("bbox") or {} |
| region_note=( |
| f"区域:{batch_spec.get('region_name')} " |
| f"({bbox.get('lon_min')}°E–{bbox.get('lon_max')}°E, " |
| f"{bbox.get('lat_min')}°N–{bbox.get('lat_max')}°N)。" |
| if batch_spec.get("region_name") else "" |
| ) |
| text=( |
| f"已创建 Ocean 批量导出任务,共 {job['total']} 个单日子任务。" |
| f"范围 {batch_spec['start_date']} 至 {batch_spec['end_date']},变量:{variables}。" |
| + region_note + |
| "学校 Ocean 接口只接收单日 YYYY-MM-DD;平台已自动把月份/日期范围拆成单日并行执行," |
| "不会再要求您补充某一天。任务卡会实时显示完成、失败、执行中和等待数量。" |
| ) |
| yield out("done",{"text":text,"batch_job_id":job["job_id"]}) |
| return |
|
|
| has_upload = bool(thread_upload_ids.get(tid)) |
| compact_spec_text=render_request_spec( |
| thread_compressed_contexts.get(tid) or context_spec |
| ) |
| routed_prompt=prompt |
| if is_confirmation and resumed: |
| thread_last_data_requests[tid]={ |
| "prompt":compact_spec_text or str(pending["prompt"]), |
| "ts":now, |
| "spec":thread_compressed_contexts.get(tid) or context_spec, |
| } |
| routed_prompt=( |
| "请下载或继续执行以下已确认的数据请求,保留全部参数:\n\n" |
| + (compact_spec_text or str(pending["prompt"])) |
| + "\n\n[USER_CONFIRMATION]\n" |
| + "用户刚刚回复确认。请立即继续执行上一项数据查询或导出," |
| + "沿用已经给出的日期、区域、变量和数据源,不要再次询问确认。\n" |
| + "[/USER_CONFIRMATION]" |
| ) |
| elif is_amendment: |
| routed_prompt=( |
| "请下载或继续执行以下数据请求,保留全部参数:\n\n" |
| + str( |
| thread_last_data_requests[tid].get("prompt") |
| or pending["prompt"] |
| or "" |
| ) |
| + "\n\n[AMENDMENT_DIRECTIVE]\n" |
| + "用户正在补充或修改上一项数据下载/查询请求。" |
| + "请把“最近补充/修改”作为本次唯一变更,其余条件继续沿用卡片中的参数;" |
| + "按最新要求立即执行,不要要求用户重新提交完整需求。\n" |
| + "[/AMENDMENT_DIRECTIVE]" |
| ) |
| routed_prompt=_apply_ocean_export_defaults(routed_prompt) |
| if _is_parse_only_request(prompt): |
| routed_prompt += ( |
| "\n\n[READ_ONLY_PARAMETER_PARSE]\n" |
| "这是只解析参数的请求。禁止调用任何 Marine/MCP 数据工具,禁止查询、" |
| "导出、下载或创建任务。只返回从用户原话提取出的参数、缺失参数和校验结果。" |
| "[/READ_ONLY_PARAMETER_PARSE]" |
| ) |
|
|
| |
| |
| |
| |
| keep_codewhale = dsh is None |
|
|
| if keep_codewhale: |
| async for chunk in stream_chat(tid, routed_prompt): |
| yield chunk |
| return |
|
|
| async for chunk in harness_stream_chat(tid, routed_prompt): |
| yield chunk |
|
|
|
|
| def init_chat_runtime(context): |
| _inject(context) |
| return stream_chat, harness_stream_chat, dispatch_chat_stream |
|
|