已合并
chore: 清理旧 app 框架及消费方 #407
chore: 清理旧 app 框架及消费方 #407
已合并
王明琦创建于 7月30日
共 100 个文件变更+1-14638
@@ -1,4 +0,0 @@
1-[submodule "agent-studio"]
2- path = agent-studio
3- url = https://gitcode.com/openJiuwen/agent-studio.git
4- branch = develop
@@ -1 +0,0 @@
1-Subproject commit 3d86cdeed137e5ebe19ba7f8bd03273768e480b3
@@ -1,7 +1,5 @@
1# ── 配置项加解密────────────────1# ── 配置项加解密────────────────
2# 未设置时以下变量均为明文;设置后可将敏感值换为 script/crypto 生成的密文。2# 未设置时以下变量均为明文;设置后可将敏感值换为 script/crypto 生成的密文。
3-# 注意: ir_execution_service 使用的 SERVER_AES_MASTER_KEY 用于「内存加密」,与
4-# 本处的 AES_MASTER_KEY(配置字段加解密)用途不同,请勿混用。
5AES_MASTER_KEY=3AES_MASTER_KEY=
6 4 
7# ── App ──────────────────────────────────────────────────────────────────────5# ── App ──────────────────────────────────────────────────────────────────────
@@ -1,357 +0,0 @@
1-# =============================================================================
2-# ir_execution_service 本地环境变量说明
3-# 复制本文件为 .env 后按需填写;含密钥请勿提交到公共仓库。
4-# =============================================================================
5- 
6-# 本服务版本号(用于 DFX 日志 version 字段,以及 app version)
7-LOWCODE_IR_EXECUTION_SERVICE_VERSION=0.1.1
8-# ----------------------------------------------------------------------------------------------------
9-# openjiuwen 工作流检查点(建议先看)
10-# FORCE_DEL_WORKFLOW_STATE:对应 SDK 内配置键 _force_del_workflow_state;设为 true 时,若同一
11-# conversation_id 下已有持久化 workflow 状态,且本次请求为「普通 dict 入参」(非交互续跑),则会在
12-# 执行前删除旧图/工作流检查点,相当于从头跑。对「__interactive_reply 续跑」仍走恢复逻辑,一般不受影响。
13-# 需避免:在依赖异常已落盘状态后、仍希望用普通输入直接顶掉续跑场景区分业务时,请结合产品评估是否开启。
14-# ----------------------------------------------------------------------------------------------------
15-FORCE_DEL_WORKFLOW_STATE=false
16- 
17-# ----------------------------------------------------------------------------------------------------
18-# 一、应用基础与代码沙箱
19-# DEBUG:是否开启调试类日志或行为,一般生产环境为 false。
20-# ----------------------------------------------------------------------------------------------------
21-DEBUG=false
22- 
23-# ----------------------------------------------------------------------------------------------------
24-# 运行时依赖(openjiuwen_runtime.foundation Settings 必填)
25-# 说明:部分代码会 import openjiuwen_runtime.foundation,
26-# 其 Settings 在 import 阶段就会读取一批字段;因此即使“只跑 IR/工作流、不启 Docker”,
27-# 也建议给齐下面两个变量(可用占位值),避免 import 时因缺字段直接失败:
28-# - IP:给部署/对外 URL 拼接用;本机调试用 127.0.0.1 即可
29-# - LOWCODE_IMAGE:给 Docker/镜像选择用;不启相关能力时填占位即可
30-# ----------------------------------------------------------------------------------------------------
31-IP=127.0.0.1
32-LOWCODE_IMAGE=dummy/lowcode:local
33- 
34-# CODE_SANDBOX_URL:代码类工作流节点(以及同类“远程执行代码”能力)的 HTTP 执行端点,启动时必填;
35-# 不跑代码节点也建议填一个可达占位地址,或确保工作流中不存在代码节点。
36-CODE_SANDBOX_URL=http://127.0.0.1:8188/run
37- 
38- 
39-# ----------------------------------------------------------------------------------------------------
40-# 二、按模型服务地址区分的 API Key(LLM_KEY__ 前缀)
41-# 说明:当工作流/画布节点里某 LLM 的 base_url 命中时,可改为从环境变量取 key:
42-# 环境变量名 = "LLM_KEY__" + 由 base_url 派生的后缀
43-# 后缀由主机名+路径等规则生成(大写、非字母数字变下划线)。
44-# 这样可以用“按网关分 key”的方式管理,而不把 key 写死在导出的 JSON 里。
45-# 注意:仅使用 DEFAULT_LLM_* 走默认大模型时,一般仍需要 DEFAULT_LLM_API_KEY(或能推导到 LLM_KEY__*)至少一处有效。
46-# ----------------------------------------------------------------------------------------------------
47-LLM_KEY__DASHSCOPE_ALIYUNCS_COM_COMPATIBLE_MODE_V1= # 需加密
48- 
49-# ----------------------------------------------------------------------------------------------------
50-# 三、默认大模型(记忆引擎 LongTermMemory 与部分补齐逻辑会使用)
51-# DEFAULT_LLM_API_BASE:OpenAI 兼容接口的 Base URL,末尾一般带 /v1。
52-# DEFAULT_LLM_MODEL_PROVIDER:客户端类型标识,常见为 OpenAI(兼容多种网关)。
53-# DEFAULT_LLM_MODEL_NAME:实际调用的模型名,需与网关支持的名称一致。
54-# DEFAULT_LLM_API_KEY:上述地址对应的密钥;若留空,部分场景会尝试用 LLM_KEY__* 按 base_url 推导。
55-# 简单理解:DEFAULT_LLM_API_KEY 是“默认大模型”的主 key;LLM_KEY__* 是“按不同 base_url 分 key”的补充机制。
56-# ----------------------------------------------------------------------------------------------------
57-DEFAULT_LLM_API_BASE="https://dashscope.aliyuncs.com/compatible-mode/v1"
58-DEFAULT_LLM_MODEL_PROVIDER="OpenAI"
59-DEFAULT_LLM_MODEL_NAME="qwen3-max"
60-DEFAULT_LLM_API_KEY= # 需加密
61- 
62- 
63-# ----------------------------------------------------------------------------------------------------
64-# 四、访问大模型 HTTPS 是否校验证书
65-# LLM_SSL_VERIFY:false 时关闭 SSL 证书校验,内网或自签证书时可开;公网生产建议 true。
66-# RESTFUL_SSL_VERIFY / RESTFUL_SSL_CERT:工作流里 HTTP 服务插件(openjiuwen RestfulApi)使用,与 LLM_SSL_* 无关。
67-# RESTFUL_SSL_VERIFY=true 表示校验 HTTPS 证书;若目标站点证书链在运行环境不可信,可能失败(常见为自签/内网 CA)。
68-# 与 LLM 侧类似:公网/生产建议 true;本地调试可临时 false。
69-# RESTFUL_SSL_CERT 可填额外信任证书(PEM 等,具体以运行时代码/客户端要求为准),用于解决“系统不信任该 CA”的问题。
70-# ----------------------------------------------------------------------------------------------------
71-# LLM_SSL_VERIFY:false 时关闭 SSL 证书校验,内网或自签证书时可开;公网生产建议 true。
72-LLM_SSL_VERIFY=false
73- 
74-RESTFUL_SSL_VERIFY=false
75-RESTFUL_SSL_CERT=
76- 
77- 
78-# ----------------------------------------------------------------------------------------------------
79-# 五、记忆引擎(LongTermMemory)行为与作用域
80-# IR_ENABLE_AGENT_MEMORY:全局记忆开关(默认 false)。设为 false 时:
81-# - 不初始化记忆引擎
82-# - 不进行记忆读写(Agent MemoryRail 也会跳过)
83-# - 启动时跳过记忆相关环境变量校验(LLM/Embedding/DB/KV/Vector 等)
84-# KV_STORE_TYPE=redis 时,记忆引擎会把 KV 元数据/索引写进 MEMORY_REDIS_URL 指向的 Redis;
85-# key 由 openjiuwen 记忆模块实现,当前可见前缀包括(以 SDK 实现为准,通常不可通过本服务单变量整体换前缀):
86-# - UMD/... :用户记忆相关
87-# - user_var/... / session_var/...:变量型记忆
88-# - MEMORY_MIGRATION_KV_SCHEMA_VERSION:迁移版本号
89-# 若与别的业务/别的服务共用一个 Redis 逻辑库,建议至少使用不同的 DB(URL 里 /db),并配合 MEMORY_SCOPE_ID/向量库隔离。
90-# MEMORY_SCOPE_ID:本应用记忆数据的逻辑隔离 ID,多实例或多人共用同一向量库时勿重复。
91-# MEMORY_INPUT_MSG_MAX_LEN:写入记忆的单条用户消息最大字符长度,过大可能被截断。
92-# MEMORY_SINGLE_TURN_HISTORY_SUMMARY_MAX_TOKEN:单轮摘要等场景允许的最大 token 上限(具体截断策略以实现为准)。
93-# ----------------------------------------------------------------------------------------------------
94-IR_ENABLE_AGENT_MEMORY=false
95-MEMORY_SCOPE_ID=ir_agent_runner_memory
96-MEMORY_INPUT_MSG_MAX_LEN=8192
97-MEMORY_SINGLE_TURN_HISTORY_SUMMARY_MAX_TOKEN=128
98- 
99- 
100-# ----------------------------------------------------------------------------------------------------
101-# 六、向量 Embedding(记忆检索、向量化写入)
102-# EMBED_MODEL_NAME:嵌入模型名称,需与 EMBED_BASE_URL 所在服务支持的列表一致。
103-# EMBED_BASE_URL:嵌入接口完整 URL,一般为某网关的 embeddings 路径。
104-# EMBED_API_KEY:调用嵌入服务的密钥。
105-# EMBEDDING_SSL_VERIFY:访问嵌入服务时是否校验 HTTPS 证书。
106-# ----------------------------------------------------------------------------------------------------
107-EMBED_MODEL_NAME=text-embedding-v3
108-EMBED_BASE_URL=https://dashscope.aliyuncs.com/compatible-mode/v1/embeddings
109-EMBED_API_KEY= # 需加密
110-EMBEDDING_SSL_VERIFY=false
111- 
112- 
113-# ----------------------------------------------------------------------------------------------------
114-# 七、向量数据库类型(记忆向量存哪里)
115-# INDEX_MANAGER_TYPE:填 milvus 或 chroma。只选其一;与下面 Milvus 或 Chroma 小节配套填写。
116-# ----------------------------------------------------------------------------------------------------
117-INDEX_MANAGER_TYPE=milvus
118- 
119- 
120-# ----------------------------------------------------------------------------------------------------
121-# 八、Chroma 向量库(仅当 INDEX_MANAGER_TYPE=chroma)
122-# MEMORY_DATA_PATH:Chroma 持久化目录,相对路径相对进程当前工作目录,不存在时部分环境会自动创建。
123-# ----------------------------------------------------------------------------------------------------
124-MEMORY_DATA_PATH=openjiuwen_runtime/memory-data
125- 
126- 
127-# ----------------------------------------------------------------------------------------------------
128-# 九、Milvus 向量库(仅当 INDEX_MANAGER_TYPE=milvus)
129-# MILVUS_HOST / MILVUS_PORT:Milvus 服务地址与 gRPC 端口,默认 19530。
130-# MILVUS_USER / MILVUS_PASSWORD:若 Milvus 开启鉴权则填写,否则可保留默认。
131-# ----------------------------------------------------------------------------------------------------
132-MILVUS_HOST=localhost
133-MILVUS_PORT=19530
134-MILVUS_USER=root
135-MILVUS_PASSWORD=123456 # 需加密
136- 
137- 
138-# ----------------------------------------------------------------------------------------------------
139-# 十、业务关系型数据库(与 openjiuwen_studio.ops.config.Settings 同名)
140-# DB_TYPE:选择数据库类型;不同取值只会启用对应小节中的配置,其它小节的值可忽略。
141-# 可选:sqlite、mysql、gaussdb、opengauss
142-# ----------------------------------------------------------------------------------------------------
143-DB_TYPE=mysql
144- 
145- 
146-# ----------------------------------------------------------------------------------------------------
147-# 十一、SQLite(仅当 DB_TYPE=sqlite)
148-# SQLITE_DB_PATH:存放 sqlite 文件的目录。
149-# OPS_SQLITE_DB / AGENT_SQLITE_DB:Studio 运维库与 Agent 库文件名(本应用 memory 等用 AGENT_SQLITE_DB 对应文件)。
150-# ----------------------------------------------------------------------------------------------------
151-SQLITE_DB_PATH=openjiuwen_runtime/data/databases
152-OPS_SQLITE_DB=ops.db
153-AGENT_SQLITE_DB=openjiuwen_runtime.db
154- 
155- 
156-# ----------------------------------------------------------------------------------------------------
157-# 十二、MySQL(仅当 DB_TYPE=mysql)
158-# DB_HOST / DB_PORT / DB_USER / DB_PASSWORD:MySQL 连接四元组(与 Studio 一致)。
159-# OPS_DB_NAME / AGENT_DB_NAME:两个“逻辑库名/Schema”用途不同:
160-# - OPS_DB_NAME:偏运维/管理面数据
161-# - AGENT_DB_NAME:偏 Agent/运行面数据(记忆等常落在这套库上)
162-# 它们都是“库名(database/schema)”,不是表名,也不是 URL 路径片段。
163-# ----------------------------------------------------------------------------------------------------
164-DB_HOST=localhost
165-DB_PORT=3306
166-DB_USER=root
167-DB_PASSWORD=123456 # 需加密
168-AGENT_DB_NAME=openjiuwen_runtime
169- 
170-# DB_NAME:MySQL 场景下建议与 AGENT_DB_NAME 保持一致(部分 foundation 校验/默认连接会依赖 DB_NAME 字段存在)。
171-DB_NAME=openjiuwen_runtime
172-OPS_DB_NAME=openjiuwen_ops
173- 
174- 
175-# ----------------------------------------------------------------------------------------------------
176-# 十三、KV 存储(记忆引擎中间状态等)
177-# KV_STORE_TYPE:
178-# - redis:多进程/多实例共享(生产常用)
179-# - inmemory:仅当前进程内存(适合本地、单测;重启丢数据)
180-# - db:走关系型数据库的异步连接(与 DB_TYPE 等配置相关;需确保 DB 已正确配置且驱动可用)
181-# ----------------------------------------------------------------------------------------------------
182-KV_STORE_TYPE=redis
183- 
184- 
185-# ----------------------------------------------------------------------------------------------------
186-# 十四、Redis URL(按业务场景拆分,互不替补/回退)
187-# 说明:为避免不同业务/不同服务的 key 冲突,这里不再支持 REDIS_HOST/PORT/... 拼装,也不再允许场景间互相回退。
188-# 你需要为每个场景提供各自的 URL(可指向不同 Redis 实例或不同逻辑库)。
189-#
190-# LOWCODE_DEFAULT_REDIS_URL:默认/基础 Redis,用于“可配 key 前缀”的场景(IR 二级缓存、Workflow 对话上下文)。
191-# - 这两类场景会用各自的前缀变量保证不冲突(见下方 LOWCODE_CONTEXT_REDIS_KEY_PREFIX、LOWCODE_IR_REDIS_KEY_PREFIX)。
192-# LOWCODE_DEFAULT_REDIS_URL=redis://default:123456@localhost:6379/0
193-# ----------------------------------------------------------------------------------------------------
194-LOWCODE_DEFAULT_REDIS_URL=
195- 
196-# MEMORY_REDIS_URL:记忆引擎(LongTermMemory)KV 专用 Redis URL(KV_STORE_TYPE=redis 时必填)。
197-# MEMORY_REDIS_URL=redis://default:123456@localhost:6379/1
198-MEMORY_REDIS_URL=
199- 
200- 
201-# ----------------------------------------------------------------------------------------------------
202-# 十五、工作流 / Agent 会话 Checkpointer(Redis 持久化图状态,可选)
203-# 用于把 工作流运行态在 Redis 中落地(多机续跑/断点等能力,取决于你启用的图执行路径)。
204-# 连接信息:只使用 CHECKPOINTER_REDIS_URL。
205-# CHECKPOINTER_REDIS_URL=redis://default:123456@localhost:6379/2
206-# key 结构由 openjiuwen.extensions.checkpointer.redis 实现(本服务未暴露单独“前缀”环境变量;需隔离时优先独立 Redis/独立 DB(URL 里 /db),并避免 session_id 与其它业务撞车):
207-# <session_id>:<namespace>:<entity_id>:<suffix...>
208-# namespace 常见为 agent、agent-team、workflow、workflow-graph;GraphStore/WorkflowStorage/Agent*Storage 会组合不同 suffix。
209-# session_exists/release 会按 "<session_id>:" 做前缀扫描/清理。
210-# CHECKPOINTER_DISABLED:设为 1 或 true 则关闭自定义 Redis checkpointer,使用 SDK 默认(多为内存)。
211-# CHECKPOINTER_DEFAULT_TTL_MINUTES:状态过期分钟数,0 表示不启用 TTL/不过期(以实现为准)。
212-# CHECKPOINTER_REFRESH_TTL_ON_READ:读时是否续期/刷新过期时间(以实现为准)。
213-# ----------------------------------------------------------------------------------------------------
214-CHECKPOINTER_REDIS_URL=redis://root:123456@localhost:6379/2
215-CHECKPOINTER_DISABLED=0
216-CHECKPOINTER_DEFAULT_TTL_MINUTES=0
217-CHECKPOINTER_REFRESH_TTL_ON_READ=0
218- 
219- 
220-# -----------------------------------------------------------------------------
221-# Workflow 对话上下文(SessionModelContext)持久化到 Redis(多机/跨请求续问)
222-# -----------------------------------------------------------------------------
223-# 说明:用于把 SessionModelContext(工作流/提问器等对话上下文)持久化到 Redis,支持跨请求续问。
224-# - 使用 LOWCODE_DEFAULT_REDIS_URL;
225-# - key 增加业务前缀,避免与其他业务/其他服务冲突;
226-# - 何时落盘/何时清理,以实现为准(一般续问场景会更依赖该持久化)。
227-#
228-# LOWCODE_CONTEXT_REDIS_KEY_PREFIX:Redis key 前缀,默认明确表示 workflow 对话上下文。
229-# LOWCODE_CONTEXT_REDIS_TTL_SECONDS:对话上下文状态 TTL(秒),每次保存会刷新 TTL。
230-# LOWCODE_CONTEXT_REDIS_LOCK_TTL_SECONDS:分布式锁 TTL(秒),避免同一会话并发请求互相覆盖。
231-# 实际 key 形态(实现见 runtime_support/context_persistence.py):
232-# <LOWCODE_CONTEXT_REDIS_KEY_PREFIX>:<conversation_id>:state:<context_id> # JSON 状态
233-# <LOWCODE_CONTEXT_REDIS_KEY_PREFIX>:<conversation_id>:lock # 分布式锁
234-LOWCODE_CONTEXT_REDIS_KEY_PREFIX=lowcode:workflow_context
235-LOWCODE_CONTEXT_REDIS_TTL_SECONDS=3600
236-LOWCODE_CONTEXT_REDIS_LOCK_TTL_SECONDS=120
237- 
238- 
239-# ----------------------------------------------------------------------------------------------------
240-# 十六、敏感配置加密(openjiuwen_studio SecurityUtils,与 Studio 一致)
241-#
242-# 根密钥配好后,下列“业务敏感项”可写成密文:DB_PASSWORD、MILVUS_TOKEN、MILVUS_PASSWORD、
243-# DEFAULT_LLM_API_KEY、LLM_KEY__*、EMBED_API_KEY、OBS_ACCESS_KEY_ID、OBS_SECRET_ACCESS_KEY 等;IR/DSL 内密文 api_key 也同理。
244-# 根密钥未配置时,上述值按普通明文处理。
245-# 更细的加解密链路与密文格式以 Studio/SecurityUtils 实现为准(此处不展开)。
246-#
247-# 运行模式(用空行分隔子模块)
248-# - SERVICE_MODE:develop(本地/联调)或 product(生产注入)
249-# - HUAWEICLOUD_KMS_ENABLED:true=走华为云 KMS 解密链;false=走本地/进程内 AES 根密钥(SERVER_AES_*)
250-#
251-# 非 KMS 根密钥(HUAWEICLOUD_KMS_ENABLED=false)
252-# - SERVER_AES_MASTER_KEY:32 字节随机根密钥的 Base64(develop 可用;生产不建议明文落盘)
253-# - SERVER_AES_MASTER_KEY_ENV:生产建议只注入此项;develop 下可与 SERVER_AES_MASTER_KEY 二选一或并存
254-#
255-# KMS 相关(HUAWEICLOUD_KMS_ENABLED=true;关闭 KMS 时可留空)
256-# - HUAWEICLOUD_USERNAME / HUAWEICLOUD_PASSWORD:IAM 用户名与密码(必填)
257-# - HUAWEICLOUD_DOMAIN_NAME:租户域名(可选,常用 Default)
258-# - HUAWEICLOUD_PROJECT_NAME / HUAWEICLOUD_PROJECT_ID:至少填其一(仅填 name 时会请求 IAM 解析 project_id)
259-# - HUAWEICLOUD_IAM_ENDPOINT:可选,默认 https://iam.myhuaweicloud.com
260-# - HUAWEICLOUD_REGION:KMS 区域,SecurityUtils 默认 cn-north-4
261-# - HUAWEICLOUD_KMS_ENDPOINT:可选,默认 https://kms.{HUAWEICLOUD_REGION}.myhuaweicloud.com
262-# - HUAWEICLOUD_KMS_KEY_ID:KMS 数据密钥 ID(解密根密钥时必填)
263-# - HUAWEICLOUD_KMS_ENCRYPTION_ALGORITHM:可选,HuaweiCloudKMS 默认 RSAES_OAEP_SHA_256
264-# ----------------------------------------------------------------------------------------------------
265-SERVICE_MODE=develop
266-HUAWEICLOUD_KMS_ENABLED=false
267- 
268-SERVER_AES_MASTER_KEY=
269-SERVER_AES_MASTER_KEY_ENV=
270- 
271-HUAWEICLOUD_USERNAME=
272-HUAWEICLOUD_PASSWORD=
273-HUAWEICLOUD_DOMAIN_NAME=
274-HUAWEICLOUD_PROJECT_NAME=
275-HUAWEICLOUD_PROJECT_ID=
276-HUAWEICLOUD_IAM_ENDPOINT=
277- 
278-HUAWEICLOUD_REGION=cn-north-4
279-HUAWEICLOUD_KMS_ENDPOINT=
280-HUAWEICLOUD_KMS_KEY_ID=
281-HUAWEICLOUD_KMS_ENCRYPTION_ALGORITHM=
282- 
283- 
284-# ----------------------------------------------------------------------------------------------------
285-# 十七、华为云 OBS 与 IR 内容二级缓存(LRU + Redis)
286-#
287-# 请求体 ir_path 为桶内对象键;IR 正文不落盘:先查进程内 LRU(可选,带 TTL),再查 Redis(二级缓存,可选,带 TTL),最后读 OBS。
288-# OBS 五项与 LOWCODE_IR_OBS_BUCKET 启动校验必填。
289-#
290-# OBS 连接(用空行分隔子模块)
291-# - OBS_ACCESS_KEY_ID:访问密钥 AK(可加密)
292-# - OBS_SECRET_ACCESS_KEY:访问密钥 SK(可加密)
293-# - OBS_SERVER:终端节点 URL,例如 https://obs.cn-north-4.myhuaweicloud.com(与桶区域一致)
294-# - OBS_REGION:区域代号,例如 cn-north-4
295-# - LOWCODE_IR_OBS_BUCKET:桶名
296-#
297-# 进程内缓存(LRU)
298-# - LOWCODE_IR_MEMORY_CACHE_ENABLED:是否启用进程内 LRU(读写都控制)
299-# - LOWCODE_IR_MEMORY_LRU_MAX:LRU 最大条数;0 表示关闭
300-# - LOWCODE_IR_MEMORY_TTL_SECONDS:LRU TTL(秒,滑动过期:命中会续期)
301-#
302-# Redis 二级缓存
303-# - LOWCODE_IR_REDIS_CACHE_ENABLED:是否启用 Redis 二级缓存(读写都控制)
304-# - LOWCODE_IR_REDIS_KEY_PREFIX:IR 缓存与锁 key 的前缀根(避免与其他业务冲突)
305-# - LOWCODE_IR_REDIS_TTL_SECONDS:Redis 中 IR 正文缓存 TTL(秒)
306-# - LOWCODE_IR_REDIS_LOCK_TTL_SECONDS:跨进程下载互斥锁 key TTL(秒)
307-# - LOWCODE_IR_FETCH_LOCK_TIMEOUT_S:等待其他实例写入完成的最大秒数(避免长时间卡死)
308-# - LOWCODE_IR_LOCAL_LOCK_TTL_SECONDS:进程内去重锁对象的清理 TTL(秒,仅影响本地内存字典大小,不影响 Redis/OBS 行为)
309-#
310-# Redis 连接:使用 LOWCODE_DEFAULT_REDIS_URL(不允许回退/替补)。
311-# LOWCODE_IR_REDIS_KEY_PREFIX:IR 内容缓存与互斥锁在 Redis 中的命名空间根(仅影响本服务 IR 二级缓存,不含 OBS 对象键);
312-# 正文:<LOWCODE_IR_REDIS_KEY_PREFIX>:data:<sha256>
313-# 锁: <LOWCODE_IR_REDIS_KEY_PREFIX>:lock:<sha256>
314-# 其中 sha256=SHA256(f"{LOWCODE_IR_OBS_BUCKET}\\n{object_key}")(见 ir_cache_fetch._dedup_token),用于避免 key 过长。
315-# ----------------------------------------------------------------------------------------------------
316-OBS_ACCESS_KEY_ID= # 需加密
317-OBS_SECRET_ACCESS_KEY= # 需加密
318-OBS_SERVER=
319-OBS_REGION=
320-LOWCODE_IR_OBS_BUCKET=
321- 
322-LOWCODE_IR_MEMORY_CACHE_ENABLED=true
323-LOWCODE_IR_MEMORY_LRU_MAX=256
324-LOWCODE_IR_MEMORY_TTL_SECONDS=300
325- 
326-LOWCODE_IR_REDIS_CACHE_ENABLED=true
327-LOWCODE_IR_REDIS_KEY_PREFIX=ir_exec
328-LOWCODE_IR_REDIS_TTL_SECONDS=86400
329-LOWCODE_IR_REDIS_LOCK_TTL_SECONDS=45
330-LOWCODE_IR_FETCH_LOCK_TIMEOUT_S=120
331-LOWCODE_IR_LOCAL_LOCK_TTL_SECONDS=300
332- 
333- 
334-# ----------------------------------------------------------------------------------------------------
335-# 十八、接口结构化日志(Server/Client 统一落盘)
336-#
337-# 说明:本服务可按 dfx.md 约定把 server(execute_invoke/execute_stream)与 client(obs/redis/llm)日志
338-# 统一写入一个 JSON Lines 文件。该目录由环境变量控制;为空则不启用。
339-# - LOWCODE_INTERFACE_LOG_PATH:日志文件完整路径(包含文件名),例如 /data/logs/interface.log;为空则不启用。
340-# - LOWCODE_INTERFACE_LOG_MAX_BYTES:轮转单文件最大字节数,默认 20971520(20MB)。
341-# - LOWCODE_INTERFACE_LOG_BACKUP_COUNT:轮转保留文件数,默认 20。
342-# ----------------------------------------------------------------------------------------------------
343-LOWCODE_INTERFACE_LOG_PATH=
344-LOWCODE_INTERFACE_LOG_MAX_BYTES=20971520
345-LOWCODE_INTERFACE_LOG_BACKUP_COUNT=20
346- 
347-# ----------------------------------------------------------------------------------------------------
348-# 十九、告警日志(alarm.log)
349-#
350-# 说明:仅记录告警类事件(warning/error/critical),用于依赖服务异常与服务自身异常告警。
351-# 需填写完整文件路径(含文件名),例如 /data/logs/alarm.log;为空则不启用。
352-# - LOWCODE_ALARM_LOG_MAX_BYTES:轮转单文件最大字节数,默认 20971520(20MB)。
353-# - LOWCODE_ALARM_LOG_BACKUP_COUNT:轮转保留文件数,默认 20。
354-# ----------------------------------------------------------------------------------------------------
355-LOWCODE_ALARM_LOG_PATH=
356-LOWCODE_ALARM_LOG_MAX_BYTES=20971520
357-LOWCODE_ALARM_LOG_BACKUP_COUNT=20
@@ -1,5 +0,0 @@
1-logs/
2-openjiuwen_runtime/
3-ir_execution_service/logs/
4-ir_execution_service/openjiuwen_runtime/
5-__pycache__/
@@ -1,11 +0,0 @@
1-exclude logs/*
2-recursive-exclude logs *
3- 
4-exclude openjiuwen_runtime/*
5-recursive-exclude openjiuwen_runtime *
6- 
7-exclude ir_execution_service/logs/*
8-recursive-exclude ir_execution_service/logs *
9- 
10-exclude ir_execution_service/openjiuwen_runtime/*
11-recursive-exclude ir_execution_service/openjiuwen_runtime *
@@ -1,226 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved
3- 
4-"""
5-DSL 工作流依赖解析:从导出 JSON 的 dependencies.workflows 解析子工作流,实现 IWorkflowLoader,
6-使 ExecutorWorkflow.compile 无需查库。
7- 
8-不修改 openjiuwen 与 openjiuwen_studio 源码;与 workflow_ir_builder 模块配合使用。
9-"""
10- 
11-from __future__ import annotations
12- 
13-from typing import Any, Dict, Tuple
14- 
15-from openjiuwen.core.workflow.workflow import Workflow as InvokableWorkflow
16-from openjiuwen_studio.core.common import dsl as studio_dsl
17-from openjiuwen_studio.core.executor.workflow.context import Context
18-from openjiuwen_studio.core.executor.workflow.workflow import IWorkflowLoader, Workflow as ExecutorWorkflow
19- 
20-from .runtime_support.http_response_contract import LowcodeApiResponseCode
21-from .runtime_support.runtime_env import llm_api_key_env_var_name, resolve_llm_api_key_from_env
22-from .runtime_support.studio_secrets import decrypt_optional_secret
23- 
24- 
25-WorkflowKey = Tuple[str, str]
26- 
27- 
28-class WorkflowLlmApiKeyMissingError(Exception):
29- """DSL LLM/意图/提问器节点:可解析的 LLM_KEY__* 与 JSON 内 api_key 均未配置。"""
30- 
31- 
32-def strip_dependencies(wf: Dict[str, Any]) -> Dict[str, Any]:
33- return {k: v for k, v in wf.items() if k != "dependencies"}
34- 
35- 
36-def collect_workflow_registry(root: Dict[str, Any]) -> Dict[WorkflowKey, Dict[str, Any]]:
37- """扁平化收集所有嵌套 dependencies.workflows,字典键为二元组 (id, version)。"""
38- reg: Dict[WorkflowKey, Dict[str, Any]] = {}
39- 
40- def walk(deps: Any) -> None:
41- if not isinstance(deps, dict):
42- return
43- for wf in deps.get("workflows") or []:
44- if not isinstance(wf, dict):
45- continue
46- wid = str(wf.get("id") or "").strip()
47- if not wid:
48- continue
49- wver = str(wf.get("version") or "draft").strip() or "draft"
50- key = (wid, wver)
51- if key not in reg:
52- reg[key] = wf
53- walk(wf.get("dependencies"))
54- 
55- walk(root.get("dependencies"))
56- return reg
57- 
58- 
59-def _inject_llm_into_component(comp: Dict[str, Any]) -> None:
60- if not isinstance(comp, dict):
61- return
62- t = comp.get("type")
63- try:
64- ti = int(t) if t is not None else -1
65- except (TypeError, ValueError):
66- ti = -1
67- if ti == int(studio_dsl.ComponentType.COMPONENT_TYPE_LOOP):
68- cfg = comp.get("configs") or {}
69- lb = cfg.get("loop_body") or {}
70- for c in lb.get("components") or []:
71- _inject_llm_into_component(c)
72- return
73- if ti not in (
74- int(studio_dsl.ComponentType.COMPONENT_TYPE_LLM),
75- int(studio_dsl.ComponentType.COMPONENT_TYPE_INTENT),
76- int(studio_dsl.ComponentType.COMPONENT_TYPE_QUESTION),
77- ):
78- return
79- 
80- cfg = comp.get("configs") or {}
81- model = cfg.get("model")
82- if not isinstance(model, dict):
83- return
84- mcc = model.get("model_client_config") or {}
85- if not isinstance(mcc, dict):
86- return
87- base_url = str(mcc.get("api_base") or "").strip()
88- envn = llm_api_key_env_var_name(base_url)
89- env_val = resolve_llm_api_key_from_env(base_url)
90- if "<SLUG_FROM_BASE_URL>" in envn:
91- return
92- 
93- json_val = str(mcc.get("api_key") or "").strip()
94- if env_val:
95- mcc["api_key"] = env_val
96- elif json_val:
97- mcc["api_key"] = json_val
98- else:
99- cid = str(comp.get("id") or comp.get("component_id") or "").strip() or "?"
100- raise WorkflowLlmApiKeyMissingError(
101- LowcodeApiResponseCode.LLM_API_KEY_MISSING.format_message(env_var=envn)
102- + f" (component id={cid})"
103- )
104- 
105- mcc["api_key"] = decrypt_optional_secret(str(mcc.get("api_key") or "").strip())
106- 
107- model["model_client_config"] = mcc
108- cfg["model"] = model
109- comp["configs"] = cfg
110- 
111- 
112-def inject_llm_api_keys_into_workflow_tree(wf: Dict[str, Any]) -> None:
113- """按 api_base 解析 LLM_KEY__*:仅当对应环境变量非空时覆盖 api_key,否则保留 DSL 内配置;二者皆空则抛 WorkflowLlmApiKeyMissingError。"""
114- for comp in wf.get("components") or []:
115- if isinstance(comp, dict):
116- _inject_llm_into_component(comp)
117- 
118- 
119-def _scalar_endpoint_from_config(value: Any) -> Any:
120- """若 connection 的 source 或 target 被写成 list,取第一个元素;仅本应用入口兜底。"""
121- if isinstance(value, list):
122- return value[0] if value else ""
123- return value
124- 
125- 
126-def _normalize_connection_endpoints_in_workflow_dict(wf: Dict[str, Any]) -> None:
127- conns = wf.get("connections")
128- if isinstance(conns, list):
129- for c in conns:
130- if not isinstance(c, dict):
131- continue
132- c["source"] = _scalar_endpoint_from_config(c.get("source"))
133- c["target"] = _scalar_endpoint_from_config(c.get("target"))
134- for comp in wf.get("components") or []:
135- if not isinstance(comp, dict):
136- continue
137- try:
138- ti = int(comp.get("type")) if comp.get("type") is not None else -1
139- except (TypeError, ValueError):
140- ti = -1
141- if ti != int(studio_dsl.ComponentType.COMPONENT_TYPE_LOOP):
142- continue
143- cfg = comp.get("configs") or {}
144- lb = cfg.get("loop_body")
145- if isinstance(lb, dict):
146- _normalize_connection_endpoints_in_workflow_dict(lb)
147- 
148- 
149-def workflow_dict_to_dl_workflow(wf: Dict[str, Any]) -> studio_dsl.Workflow:
150- stripped = strip_dependencies(wf)
151- _normalize_connection_endpoints_in_workflow_dict(stripped)
152- inject_llm_api_keys_into_workflow_tree(stripped)
153- return studio_dsl.Workflow.model_validate(stripped)
154- 
155- 
156-class DependencyWorkflowLoader(IWorkflowLoader):
157- """按 id/version 在 dependencies 扁平表中查找子工作流并递归 compile。"""
158- 
159- def __init__(
160- self,
161- registry: Dict[WorkflowKey, Dict[str, Any]],
162- space_id: str,
163- current_user: Dict[str, Any],
164- ) -> None:
165- self._registry = registry
166- self._space_id = space_id
167- self._current_user = current_user
168- self._cache: Dict[WorkflowKey, InvokableWorkflow] = {}
169- self._compiling: set[WorkflowKey] = set()
170- 
171- def _resolve(self, wid: str, version: str) -> tuple[Dict[str, Any], WorkflowKey]:
172- vid = str(wid or "").strip()
173- ver = str(version or "").strip() or "draft"
174- key: WorkflowKey = (vid, ver)
175- d = self._registry.get(key)
176- if d is None and ver != "draft":
177- key = (vid, "draft")
178- d = self._registry.get(key)
179- if d is None:
180- raise ValueError(
181- f"dependencies.workflows 中未找到子工作流 id={wid!r} version={version!r},"
182- f"已注册 id 列表: {[k[0] for k in self._registry]}"
183- )
184- return d, key
185- 
186- async def get_compiled_workflow(
187- self,
188- context: Context,
189- workflow_id: str,
190- version: str,
191- space_id: str,
192- current_user: Dict[str, Any],
193- ) -> InvokableWorkflow:
194- wf_dict, cache_key = self._resolve(workflow_id, version)
195- if cache_key in self._cache:
196- return self._cache[cache_key]
197- if cache_key in self._compiling:
198- raise ValueError(f"子工作流循环依赖: id={cache_key[0]!r} version={cache_key[1]!r}")
199- self._compiling.add(cache_key)
200- try:
201- dl = workflow_dict_to_dl_workflow(wf_dict)
202- user = current_user if current_user is not None else self._current_user
203- executor = ExecutorWorkflow(dl, space_id, user)
204- compiled = await executor.compile(context, loader=self)
205- self._cache[cache_key] = compiled
206- return compiled
207- finally:
208- self._compiling.discard(cache_key)
209- 
210- 
211-def unwrap_workflow_document(ir: Dict[str, Any]) -> Dict[str, Any]:
212- """若导出为外层 workflow 键包裹的内层 DSL,则合并 dependencies 后返回内层字典。"""
213- if isinstance(ir.get("components"), list) and isinstance(ir.get("connections"), list):
214- return ir
215- w = ir.get("workflow")
216- if isinstance(w, dict) and isinstance(w.get("components"), list):
217- merged = dict(w)
218- if isinstance(ir.get("dependencies"), dict) and "dependencies" not in w:
219- merged["dependencies"] = ir["dependencies"]
220- return merged
221- return ir
222- 
223- 
224-def looks_like_dsl_workflow_export(ir: Dict[str, Any]) -> bool:
225- ir = unwrap_workflow_document(ir)
226- return isinstance(ir.get("components"), list) and isinstance(ir.get("connections"), list)
@@ -1,275 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved
3- 
4-"""execute_invoke 的业务逻辑与返回体转换。"""
5- 
6-from __future__ import annotations
7- 
8-import asyncio
9-from typing import Any
10- 
11-from fastapi.responses import JSONResponse
12- 
13-from openjiuwen.core.runner import Runner
14-from openjiuwen.core.context_engine.schema.config import ContextEngineConfig
15-from openjiuwen_runtime.foundation.log import get_logger
16-from openjiuwen_studio.schemas import ResponseModel
17- 
18-from .dsl_workflow_dependency_loader import WorkflowLlmApiKeyMissingError
19-from .react_agent_builder import build_react_agent_from_ir_dict
20-from .runtime_support.context_persistence import RedisContextPersistence
21-from .runtime_support.http_response_contract import (
22- LowcodeApiResponseCode,
23- ResponseDataType,
24- build_error_response_model,
25- to_jsonable,
26-)
27-from .runtime_support.execution_request import ExecutionPrepareError, prepare_execution_request
28-from .runtime_support.runtime_bootstrap import ensure_runtime_ready
29-from .runtime_support.workflow_context_helpers import append_user_input_message_if_needed, stable_workflow_context_id
30- 
31-JSON_MEDIA_TYPE = "application/json; charset=utf-8"
32- 
33-_log = get_logger(__name__)
34- 
35-_ctx_store = RedisContextPersistence()
36- 
37- 
38-def _json_response(model: ResponseModel) -> JSONResponse:
39- # HTTP 始终返回 200,由 body.code 表达业务状态。
40- return JSONResponse(model.model_dump(), media_type=JSON_MEDIA_TYPE)
41- 
42- 
43-def _invoke_exception_to_model(exc: Exception) -> ResponseModel:
44- """将 invoke 路径异常转换为 error 响应。"""
45- from openjiuwen.core.common.exception.errors import BaseError
46- 
47- if isinstance(exc, WorkflowLlmApiKeyMissingError):
48- c = LowcodeApiResponseCode.LLM_API_KEY_MISSING
49- return build_error_response_model(c, message=str(exc))
50- if isinstance(exc, asyncio.TimeoutError):
51- code = LowcodeApiResponseCode.EXECUTION_TIMEOUT
52- return build_error_response_model(code, message=code.format_message())
53- if isinstance(exc, BaseError):
54- detail_code = int(getattr(exc, "code", LowcodeApiResponseCode.INTERNAL_ERROR))
55- msg = str(getattr(exc, "message", "") or exc)
56- return build_error_response_model(
57- LowcodeApiResponseCode.EXECUTION_FAILED,
58- message=msg,
59- payload={"detail_code": detail_code},
60- )
61- if isinstance(exc, ValueError):
62- msg = str(exc)
63- return build_error_response_model(LowcodeApiResponseCode.INVALID_PARAM, message=msg)
64- msg = str(exc)
65- return build_error_response_model(LowcodeApiResponseCode.INTERNAL_ERROR, message=msg)
66- 
67- 
68-def _agent_invoke_result_to_model(result: Any) -> ResponseModel:
69- """将 agent invoke 返回值映射为 ResponseModel。"""
70- ok = LowcodeApiResponseCode.SUCCESS
71- 
72- if isinstance(result, dict):
73- result_type = str(result.get("result_type") or "").strip()
74- 
75- if result_type == "answer":
76- payload = {"output": str(result.get("output", "") or "")}
77- return ResponseModel(code=int(ok), message=ok.default_message, data={"type": "result", "payload": payload})
78- 
79- if result_type == "error":
80- raw_msg = result.get("message")
81- if raw_msg is None:
82- raw_msg = result.get("output")
83- msg = str(raw_msg or "")
84- code = LowcodeApiResponseCode.EXECUTION_FAILED
85- return ResponseModel(
86- code=int(code),
87- message=msg.strip() or code.default_message,
88- data={"type": "error", "payload": {"message": msg}},
89- )
90- 
91- if result_type == "interrupt":
92- payload = {
93- "workflow_execution_state": to_jsonable(result.get("workflow_execution_state")),
94- "component_ids": to_jsonable(result.get("component_ids", [])),
95- }
96- return ResponseModel(
97- code=int(ok),
98- message=ok.default_message,
99- data={"type": "interaction", "payload": payload},
100- )
101- 
102- if result_type:
103- return ResponseModel(
104- code=int(ok),
105- message=ok.default_message,
106- data={
107- "type": ResponseDataType.FORCE_FINISH.value,
108- "payload": to_jsonable(result),
109- },
110- )
111- 
112- return ResponseModel(
113- code=int(ok),
114- message=ok.default_message,
115- data={"type": ResponseDataType.UNKNOWN.value, "payload": to_jsonable(result)},
116- )
117- 
118- return ResponseModel(
119- code=int(ok),
120- message=ok.default_message,
121- data={"type": "result", "payload": to_jsonable(result)},
122- )
123- 
124- 
125-def _workflow_invoke_result_to_model(result: Any) -> ResponseModel:
126- """将 workflow invoke 返回值映射为 ResponseModel。"""
127- from openjiuwen.core.common.constants.constant import INTERACTION
128- from openjiuwen.core.session.stream import OutputSchema
129- from openjiuwen.core.workflow import WorkflowExecutionState, WorkflowOutput
130- 
131- ok = LowcodeApiResponseCode.SUCCESS
132- 
133- if isinstance(result, WorkflowOutput):
134- state = getattr(result, "state", None)
135- inner = getattr(result, "result", None)
136- if state == WorkflowExecutionState.INPUT_REQUIRED:
137- result = inner if inner is not None else []
138- elif state == WorkflowExecutionState.COMPLETED:
139- result = inner
140- else:
141- result = inner
142- 
143- if isinstance(result, dict):
144- return ResponseModel(
145- code=int(ok),
146- message=ok.default_message,
147- data={"type": "result", "payload": {"data": to_jsonable(result)}},
148- )
149- 
150- if isinstance(result, list):
151- for item in result:
152- output_type = None
153- payload_obj: Any = None
154- if isinstance(item, OutputSchema):
155- output_type = item.type
156- payload_obj = item.payload
157- elif isinstance(item, dict):
158- output_type = item.get("type")
159- payload_obj = item.get("payload")
160- 
161- if output_type == INTERACTION:
162- payload_json = to_jsonable(payload_obj)
163- interaction_id = payload_json.get("id") if isinstance(payload_json, dict) else None
164- interaction_value = payload_json.get("value") if isinstance(payload_json, dict) else None
165- return ResponseModel(
166- code=int(ok),
167- message=ok.default_message,
168- data={"type": "interaction", "payload": {"id": interaction_id, "value": interaction_value}},
169- )
170- 
171- c = LowcodeApiResponseCode.INVOKE_NOT_SUPPORTED
172- err_payload = {"error_code": int(c), "error_message": c.default_message}
173- return ResponseModel(
174- code=int(c),
175- message=c.default_message,
176- data={"type": "error", "payload": err_payload},
177- )
178- 
179- return ResponseModel(
180- code=int(ok),
181- message=ok.default_message,
182- data={"type": "result", "payload": {"data": to_jsonable(result)}},
183- )
184- 
185- 
186-async def handle_execute_invoke(body: Any) -> JSONResponse:
187- """FastAPI 路由层入口。"""
188- try:
189- await ensure_runtime_ready()
190- except Exception as e:
191- _log.exception("service unavailable during startup: %s", e)
192- return _json_response(build_error_response_model(LowcodeApiResponseCode.SERVICE_UNAVAILABLE, message=str(e)))
193- 
194- try:
195- prepared = await prepare_execution_request(body)
196- except ExecutionPrepareError as exc:
197- return _json_response(build_error_response_model(exc.code, message=exc.message))
198- 
199- if prepared.executable_kind == "workflow":
200- session_id = str(prepared.session_id or "").strip()
201- if not session_id:
202- c = LowcodeApiResponseCode.MISSING_PARAM
203- return _json_response(
204- build_error_response_model(c, message=c.format_message(field="conversation_id"))
205- )
206- 
207- from .workflow_ir_builder import build_core_workflow_from_ir_dict
208- from openjiuwen.core.workflow import WorkflowExecutionState, WorkflowOutput
209- 
210- try:
211- workflow = await build_core_workflow_from_ir_dict(
212- prepared.ir_root,
213- space_id=prepared.space_id,
214- current_user=prepared.current_user,
215- )
216- except WorkflowLlmApiKeyMissingError as e:
217- c = LowcodeApiResponseCode.LLM_API_KEY_MISSING
218- return _json_response(build_error_response_model(c, message=str(e)))
219- except Exception as e:
220- return _json_response(build_error_response_model(LowcodeApiResponseCode.IR_LOAD_FAILED, message=str(e)))
221- 
222- conversation_id = session_id
223- context_id = stable_workflow_context_id(workflow)
224- 
225- try:
226- async with _ctx_store.conversation_lock(conversation_id=conversation_id):
227- context = await _ctx_store.load_context(
228- conversation_id=conversation_id,
229- context_id=context_id,
230- config=ContextEngineConfig(),
231- )
232- await append_user_input_message_if_needed(context, prepared.inputs_obj)
233- 
234- wf_output = await asyncio.wait_for(
235- Runner.run_workflow(
236- workflow=workflow,
237- inputs=prepared.inputs_obj,
238- session=prepared.session_id,
239- context=context,
240- ),
241- timeout=prepared.timeout_seconds,
242- )
243- 
244- if isinstance(wf_output, WorkflowOutput) and wf_output.state == WorkflowExecutionState.INPUT_REQUIRED:
245- await _ctx_store.save_on_interaction(
246- conversation_id=conversation_id,
247- context_id=context_id,
248- context=context,
249- )
250- else:
251- await _ctx_store.delete(conversation_id=conversation_id, context_id=context_id)
252- except Exception as e:
253- _log.exception("workflow invoke failed: %s", e)
254- if conversation_id:
255- await _ctx_store.delete(conversation_id=conversation_id, context_id=context_id)
256- return _json_response(_invoke_exception_to_model(e))
257- 
258- return _json_response(_workflow_invoke_result_to_model(wf_output))
259- 
260- try:
261- react_agent = await build_react_agent_from_ir_dict(prepared.ir_root, prepared.current_user)
262- except Exception as e:
263- return _json_response(build_error_response_model(LowcodeApiResponseCode.IR_LOAD_FAILED, message=str(e)))
264- 
265- try:
266- agent_output = await asyncio.wait_for(
267- Runner.run_agent(agent=react_agent, inputs=prepared.inputs_obj, session=prepared.session_id),
268- timeout=prepared.timeout_seconds,
269- )
270- except Exception as e:
271- _log.exception("agent invoke failed: %s", e)
272- return _json_response(_invoke_exception_to_model(e))
273- 
274- return _json_response(_agent_invoke_result_to_model(agent_output))
275- 
@@ -1,246 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved.
3- 
4-"""IR 执行服务 HTTP 入口。"""
5- 
6-from __future__ import annotations
7- 
8-import json
9-import logging
10-import os
11-import time
12-from pathlib import Path
13- 
14-from fastapi.exceptions import RequestValidationError
15-from fastapi.responses import JSONResponse
16-from pydantic import BaseModel, Field
17-from sse_starlette import EventSourceResponse
18-from starlette.middleware.base import BaseHTTPMiddleware
19-from starlette.requests import Request
20-from starlette.responses import Response
21- 
22-_APP_DIR = Path(__file__).resolve().parent
23-_APP_ROOT = _APP_DIR.parent
24- 
25-try:
26- from dotenv import load_dotenv
27- 
28- load_dotenv(_APP_ROOT / ".env", override=False)
29-except ImportError:
30- pass
31- 
32-from openjiuwen.core.runner import Runner
33-from openjiuwen.core.common.logging import set_session_id
34-from openjiuwen_runtime.service.app.base_app import BaseApp
35- 
36-from .runtime_support.runtime_env_prepare import prepare_runtime_environment
37-from .runtime_support.error_logging import setup_error_file_logging
38-from .runtime_support.alarm_logger import (
39- init_alarm_logger_from_env,
40- install_core_alarm_sink,
41- install_runner_tool_alarm_callbacks,
42-)
43-from .runtime_support.core_log_bridge import install_core_log_bridge
44-from .runtime_support.interface_logger import (
45- init_interface_logger_from_env,
46- install_core_interface_sink,
47- install_runner_llm_stream_interface_callbacks,
48- install_runner_tool_interface_callbacks,
49- log_server,
50- set_request_context,
51-)
52- 
53-prepare_runtime_environment()
54-setup_error_file_logging()
55-init_alarm_logger_from_env()
56-install_core_log_bridge()
57-install_core_alarm_sink()
58-install_runner_tool_alarm_callbacks()
59-init_interface_logger_from_env()
60-install_core_interface_sink()
61-install_runner_tool_interface_callbacks()
62-install_runner_llm_stream_interface_callbacks()
63- 
64-# Workflow 默认超时较短,复杂 DSL 续跑时容易误判超时。
65-os.environ.setdefault("WORKFLOW_EXECUTE_TIMEOUT", "300")
66- 
67-from .runtime_support.http_response_contract import (
68- LowcodeApiResponseCode,
69- build_error_response_model,
70-)
71-from .runtime_support.runtime_bootstrap import ensure_runtime_ready
72- 
73-_JSON_MEDIA_TYPE = "application/json; charset=utf-8"
74-_PY_LOG = logging.getLogger(__name__)
75- 
76- 
77-async def _response_body_bytes(resp: Response) -> bytes | None:
78- """取响应正文。注意 BaseHTTPMiddleware.call_next 返回的是流式包装体,通常没有物化的 .body。"""
79- body_iter = getattr(resp, "body_iterator", None)
80- if body_iter is not None:
81- parts: list[bytes] = []
82- async for chunk in body_iter:
83- if not chunk:
84- continue
85- if isinstance(chunk, memoryview):
86- parts.append(chunk.tobytes())
87- elif isinstance(chunk, (bytes, bytearray)):
88- parts.append(bytes(chunk))
89- else:
90- parts.append(str(chunk).encode(resp.charset))
91- return b"".join(parts)
92- raw = getattr(resp, "body", None)
93- if isinstance(raw, memoryview) and raw:
94- return raw.tobytes()
95- if isinstance(raw, (bytes, bytearray)) and raw:
96- return bytes(raw)
97- return None
98- 
99- 
100-def _invalid_request_json_response(exc: RequestValidationError) -> JSONResponse:
101- body = build_error_response_model(
102- LowcodeApiResponseCode.INVALID_REQUEST,
103- message=LowcodeApiResponseCode.INVALID_REQUEST.default_message,
104- payload={"errors": exc.errors()},
105- )
106- return JSONResponse(body.model_dump(), media_type=_JSON_MEDIA_TYPE)
107- 
108- 
109-class IrQueryBody(BaseModel):
110- user_id: str
111- conversation_id: str
112- ir_path: str = Field(
113- ...,
114- description="OBS 桶内对象键(Object Key);正文经进程内缓存与可选 Redis 缓存,加速读取,不落盘。",
115- )
116- inputs: str
117- timeout_ms: int = Field(120_000, ge=1)
118- 
119- 
120-class IrExecutionServiceApp(BaseApp):
121- """BaseApp 上挂载自定义 POST 路由,不继承 AgentApp。"""
122- 
123- def __init__(self) -> None:
124- super().__init__(
125- app_name="IrExecutionService",
126- app_description="面向低代码的工作流 IR 执行 HTTP 服务",
127- version=(os.environ.get("LOWCODE_IR_EXECUTION_SERVICE_VERSION") or "").strip(),
128- )
129- 
130- class _DfxMiddleware(BaseHTTPMiddleware):
131- async def dispatch(self, request: Request, call_next):
132- # Only record DFX logs for the two public endpoints.
133- path = (request.url.path or "").rstrip("/")
134- if not (path.endswith("/execute_invoke") or path.endswith("/execute_stream")):
135- return await call_next(request)
136- 
137- import uuid
138- 
139- request_id = uuid.uuid4().hex
140- source_ip = getattr(getattr(request, "client", None), "host", "") or ""
141- set_request_context(request_id=request_id, source_ip=source_ip)
142- set_session_id(request_id)
143- 
144- interface_name = "execute_invoke" if path.endswith("/execute_invoke") else "execute_stream"
145- t0 = time.perf_counter()
146- try:
147- resp = await call_next(request)
148- except Exception as e:
149- log_server(
150- interface_name=interface_name,
151- cost_ms=(time.perf_counter() - t0) * 1000.0,
152- ok=False,
153- return_code=int(LowcodeApiResponseCode.INTERNAL_ERROR),
154- return_info=str(e),
155- source_ip=source_ip,
156- add_info={"path": path},
157- )
158- raise
159- 
160- # For invoke: BaseHTTPMiddleware 下 call_next 得到的是流式包装响应,无物化 .body,需先读全再解析。
161- if interface_name == "execute_invoke":
162- had_stream_wrapper = getattr(resp, "body_iterator", None) is not None
163- code = int(LowcodeApiResponseCode.SUCCESS)
164- msg = str(LowcodeApiResponseCode.SUCCESS.default_message)
165- ok = True
166- parsed_ok = False
167- raw_bytes: bytes | None = None
168- try:
169- raw_bytes = await _response_body_bytes(resp)
170- if raw_bytes:
171- parsed = json.loads(raw_bytes.decode("utf-8"))
172- if isinstance(parsed, dict):
173- code = int(parsed.get("code", 0))
174- msg = str(parsed.get("message", "") or "")
175- ok = code == 0
176- parsed_ok = True
177- except Exception as exc:
178- _PY_LOG.warning("failed to parse execute_invoke response body: %s", exc)
179- if not parsed_ok:
180- ok = False
181- code = int(LowcodeApiResponseCode.INTERNAL_ERROR)
182- msg = "parse failed"
183- log_server(
184- interface_name=interface_name,
185- cost_ms=(time.perf_counter() - t0) * 1000.0,
186- ok=ok,
187- return_code=code,
188- return_info=msg,
189- source_ip=source_ip,
190- add_info={"path": path},
191- )
192- if had_stream_wrapper:
193- return Response(
194- content=raw_bytes or b"",
195- status_code=resp.status_code,
196- headers=resp.headers,
197- media_type=resp.media_type,
198- )
199- # Stream path end-state is logged inside stream generator (see stream_api).
200- return resp
201- 
202- self.app.add_middleware(_DfxMiddleware)
203- 
204- @self.app.exception_handler(RequestValidationError)
205- async def _validation_on_stream_routes(request: Request, exc: RequestValidationError):
206- path = (request.url.path or "").rstrip("/")
207- if path.endswith("/execute_stream"):
208- from .stream_api import validation_error_stream_events
209- 
210- return EventSourceResponse(validation_error_stream_events(exc))
211- if path.endswith("/execute_invoke"):
212- return _invalid_request_json_response(exc)
213- return JSONResponse(status_code=422, content={"detail": exc.errors()})
214- 
215- @self.app.post("/execute_stream")
216- async def execute_stream(body: IrQueryBody):
217- from .stream_api import execute_stream_event_source
218- 
219- return EventSourceResponse(execute_stream_event_source(body))
220- 
221- @self.app.post("/execute_invoke")
222- async def execute_invoke(body: IrQueryBody):
223- from .invoke_api import handle_execute_invoke
224- 
225- return await handle_execute_invoke(body)
226- 
227- 
228-runner = IrExecutionServiceApp()
229- 
230- 
231-@runner.init
232-async def _startup() -> None:
233- await ensure_runtime_ready()
234- await Runner.start()
235- 
236- 
237-@runner.shutdown
238-async def _shutdown() -> None:
239- await Runner.stop()
240- 
241- 
242-app = runner.app
243- 
244- 
245-if __name__ == "__main__":
246- runner.run()
@@ -1,352 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved.
3- 
4-"""从 Studio 导出 IR 构建 ReActAgent。"""
5-from __future__ import annotations
6- 
7-import json
8-import os
9-from pathlib import Path
10-from typing import Any
11- 
12-from openjiuwen.core.common.schema.param import Param
13-from openjiuwen.core.context_engine.schema.config import ContextEngineConfig
14-from openjiuwen.core.foundation.llm.schema.config import ModelClientConfig, ModelRequestConfig
15-from openjiuwen.core.memory.config.config import AgentMemoryConfig
16-from openjiuwen.core.single_agent.agents.react_agent import ReActAgent, ReActAgentConfig as NewReActAgentConfig
17-from openjiuwen.core.single_agent.legacy.config import LegacyReActAgentConfig
18-from openjiuwen.core.single_agent.schema.agent_card import AgentCard
19-from openjiuwen_studio.lowcode.compiler import AgentCompiler
20-from openjiuwen_studio.lowcode.config_adapter import ConfigAdapter
21-from openjiuwen_studio.lowcode.schemas import ModelOverride
22- 
23-from .runtime_support.runtime_env import (
24- clean_env_value,
25- get_bool_env,
26- get_env,
27- resolve_llm_api_key_from_env,
28- resolve_memory_scope_id,
29-)
30-from .runtime_support.studio_secrets import decrypt_optional_secret, resolve_secret_env
31- 
32- 
33-def build_model_overrides_from_default_llm_env(export_data: dict[str, Any]) -> dict[str, ModelOverride]:
34- """按 DEFAULT_LLM_* 与 LLM_KEY__ 推导规则生成 ModelOverride,仅含键 str(agent.model_id)。"""
35- agent = export_data.get("agent") if isinstance(export_data.get("agent"), dict) else {}
36- mid = agent.get("model_id")
37- if mid is None or not str(mid).strip():
38- return {}
39- 
40- model_name = clean_env_value("DEFAULT_LLM_MODEL_NAME")
41- base_url = clean_env_value("DEFAULT_LLM_API_BASE")
42- api_key = resolve_secret_env("DEFAULT_LLM_API_KEY", "")
43- provider = clean_env_value("DEFAULT_LLM_MODEL_PROVIDER", "")
44- 
45- if not api_key and base_url:
46- api_key = resolve_llm_api_key_from_env(base_url)
47- 
48- override_kwargs: dict[str, Any] = {}
49- if model_name:
50- override_kwargs["name"] = model_name
51- if base_url:
52- override_kwargs["base_url"] = base_url
53- if api_key:
54- override_kwargs["api_key"] = api_key
55- if provider:
56- override_kwargs["provider"] = provider
57- 
58- if not override_kwargs:
59- return {}
60- 
61- return {str(mid): ModelOverride(**override_kwargs)}
62- 
63- 
64-def normalize_runtime_config_for_react_agent(
65- config: LegacyReActAgentConfig | NewReActAgentConfig,
66-) -> NewReActAgentConfig:
67- """兼容 legacy 与新版 ReActAgentConfig,统一为 core 新版配置。"""
68- if isinstance(config, NewReActAgentConfig):
69- return config
70- 
71- model_name = getattr(config, "model_name", "") or ""
72- m = getattr(config, "model_config", None)
73- info = getattr(m, "model_info", None) if m is not None else None
74- model_provider = str(getattr(m, "model_provider", "") or "")
75- raw_api_key = str(getattr(info, "api_key", "") or "").strip()
76- api_key_plain = decrypt_optional_secret(raw_api_key)
77- mcc = ModelClientConfig(
78- model_provider=model_provider,
79- api_key=api_key_plain,
80- api_base=str(getattr(info, "api_base", "") or ""),
81- verify_ssl=get_bool_env("LLM_SSL_VERIFY", True),
82- )
83- mrc = ModelRequestConfig(
84- temperature=getattr(info, "temperature", None),
85- max_tokens=getattr(info, "max_tokens", None),
86- timeout=float(getattr(info, "timeout", 60) or 60),
87- )
88- ctx_cfg = ContextEngineConfig(
89- max_context_message_num=200,
90- default_window_round_num=config.constrain.reserved_max_chat_rounds,
91- )
92- return NewReActAgentConfig(
93- mem_scope_id=config.memory_scope_id or "",
94- model_name=str(model_name),
95- model_provider=model_provider,
96- api_key=api_key_plain,
97- api_base=str(getattr(info, "api_base", "") or ""),
98- prompt_template_name=config.prompt_template_name or "",
99- prompt_template=list(config.prompt_template or []),
100- max_iterations=config.constrain.max_iteration,
101- model_client_config=mcc,
102- model_config_obj=mrc,
103- context_engine_config=ctx_cfg,
104- )
105- 
106- 
107-def _agent_memory_config_from_export_memory(memory: Any) -> AgentMemoryConfig:
108- """从导出 JSON 的 agent.memory 构建 AgentMemoryConfig;缺省为 false 或空列表。"""
109- if not isinstance(memory, dict):
110- memory = {}
111- 
112- raw_vars = memory.get("variable_config")
113- if not isinstance(raw_vars, list):
114- raw_vars = []
115- 
116- mem_variables: list[Any] = []
117- for var in raw_vars:
118- if not isinstance(var, dict):
119- continue
120- if not var.get("enabled", False):
121- continue
122- name = str(var.get("name") or "").strip()
123- if not name:
124- continue
125- desc = str(var.get("description") or "")
126- mem_variables.append(Param.string(name, description=desc, required=False))
127- 
128- return AgentMemoryConfig(
129- mem_variables=mem_variables,
130- # 注意:AgentMemoryConfig 在 core 里默认都是 True;这里必须以导出 IR 为准,
131- # 且缺省按 False 处理,避免“开关没开也加载/写入记忆”。
132- enable_long_term_mem=bool(memory.get("longterm_memory_config", False)),
133- enable_user_profile=bool(memory.get("user_profile_config", False)),
134- enable_semantic_memory=bool(memory.get("semantic_memory_config", False)),
135- enable_episodic_memory=bool(memory.get("episodic_memory_config", False)),
136- enable_summary_memory=bool(memory.get("summary_memory_config", False)),
137- )
138- 
139- 
140-def _memory_switch_enabled() -> bool:
141- """全局记忆开关:默认开启;设置 IR_ENABLE_AGENT_MEMORY=false 可关闭所有记忆加载/写入。"""
142- v = (os.environ.get("IR_ENABLE_AGENT_MEMORY") or "true").strip().lower()
143- return v not in {"0", "false", "no", "off"}
144- 
145- 
146-def _is_agent_memory_cfg_enabled(cfg: AgentMemoryConfig) -> bool:
147- """只要任一记忆能力开启,就认为需要挂载 MemoryRail。"""
148- if not isinstance(cfg, AgentMemoryConfig):
149- return False
150- return bool(
151- cfg.mem_variables
152- or cfg.enable_long_term_mem
153- or cfg.enable_user_profile
154- or cfg.enable_semantic_memory
155- or cfg.enable_episodic_memory
156- or cfg.enable_summary_memory
157- )
158- 
159- 
160-def _ensure_memory_placeholders_in_system_prompt(agent: Any, agent_memory_cfg: AgentMemoryConfig) -> None:
161- """
162- 仅当导出 JSON 的开关开启时,才注入对应占位符。
163- MemoryRail 使用 PromptTemplate(占位符前后缀为 {{ 与 }})渲染 system message 中的记忆变量。
164- """
165- enable_vars = bool(getattr(agent_memory_cfg, "mem_variables", None))
166- enable_long_term = bool(getattr(agent_memory_cfg, "enable_long_term_mem", False))
167- if not (enable_vars or enable_long_term):
168- return
169- 
170- cfg = getattr(agent, "_config", None)
171- if cfg is None:
172- return
173- prompt_template = getattr(cfg, "prompt_template", None)
174- if not isinstance(prompt_template, list):
175- prompt_template = []
176- try:
177- cfg.prompt_template = prompt_template
178- except (AttributeError, TypeError):
179- return
180- # 空列表时也必须继续:否则 MemoryRail 写入 ctx.extra 的占位符永远不会出现在任何 system 消息里。
181- 
182- def _has_placeholder(key: str) -> bool:
183- token = "{{" + key + "}}"
184- for m in prompt_template:
185- if not isinstance(m, dict):
186- continue
187- if m.get("role") != "system":
188- continue
189- c = m.get("content")
190- if isinstance(c, str) and token in c:
191- return True
192- return False
193- 
194- ok_long_term = (not enable_long_term) or _has_placeholder("sys_long_term_memory")
195- ok_vars = (not enable_vars) or _has_placeholder("sys_memory_variables")
196- if ok_long_term and ok_vars:
197- return
198- 
199- lines = ["【系统记忆注入(由服务端自动追加)】"]
200- lines.append("你可能会获得以下记忆信息(均为 JSON 字符串),用于辅助回答:")
201- if enable_long_term:
202- lines.append("- 长期记忆(列表,可能为空):{{sys_long_term_memory}}")
203- if enable_vars:
204- lines.append("- 用户记忆变量(字典,可能为空):{{sys_memory_variables}}")
205- lines.append("规则:若相关字段为空,不要编造用户信息。")
206- 
207- prompt_template.append({"role": "system", "content": "\n".join(lines)})
208- 
209- 
210-async def _register_memory_rail_from_export(agent: Any, export_agent: dict[str, Any]) -> None:
211- """根据导出 IR 的 memory 配置挂载 MemoryRail(如未开启则跳过)。"""
212- if not _memory_switch_enabled():
213- return
214- 
215- memory = export_agent.get("memory") if isinstance(export_agent, dict) else None
216- if not isinstance(memory, dict) or not memory:
217- return
218- 
219- agent_memory_cfg = _agent_memory_config_from_export_memory(memory)
220- if not _is_agent_memory_cfg_enabled(agent_memory_cfg):
221- return
222- 
223- scope_id = resolve_memory_scope_id(
224- raw_memory_scope_id=str(getattr(agent, "_config", None).mem_scope_id or ""),
225- default_memory_scope_id=get_env("DEFAULT_MEMORY_SCOPE_ID", ""),
226- )
227- 
228- from openjiuwen.core.application.llm_agent.rails.memory_rail import MemoryRail
229- 
230- await agent.register_rail(MemoryRail(scope_id, agent_memory_cfg))
231- _ensure_memory_placeholders_in_system_prompt(agent, agent_memory_cfg)
232- 
233- 
234-def _adapt_runtime_config(agent_config_dict: dict[str, Any]) -> Any:
235- adapt_to_runtime = getattr(ConfigAdapter, "adapt_to_runtime_config", None)
236- if callable(adapt_to_runtime):
237- return adapt_to_runtime(agent_config_dict)
238- return ConfigAdapter.adapt(agent_config_dict)
239- 
240- 
241-def _agent_card_from_export_agent(export_agent: dict[str, Any]) -> AgentCard:
242- return AgentCard(
243- id=export_agent.get("agent_id", ""),
244- name=export_agent.get("agent_name", "Agent"),
245- description=export_agent.get("description", ""),
246- version=export_agent.get("agent_version", "draft"),
247- )
248- 
249- 
250-def _prepend_configs_system_prompt_first(
251- export_agent: dict[str, Any],
252- runtime_config: NewReActAgentConfig,
253-) -> None:
254- """将导出 IR 中 agent.configs.system_prompt 作为第一条 system 与 prompt_template 合并(置前)。"""
255- configs = export_agent.get("configs") if isinstance(export_agent.get("configs"), dict) else None
256- if not configs:
257- return
258- raw = configs.get("system_prompt")
259- if not isinstance(raw, str):
260- return
261- text = raw.strip()
262- if not text:
263- return
264- existing = list(runtime_config.prompt_template or [])
265- runtime_config.prompt_template = [{"role": "system", "content": text}, *existing]
266- 
267- 
268-async def _compile_runtime_config_from_export_data(
269- export_data: dict[str, Any],
270- current_user: dict[str, Any],
271- model_overrides: dict[str, Any] | None,
272-) -> tuple[AgentCard, NewReActAgentConfig]:
273- compiler = AgentCompiler()
274- export_agent = export_data.get("agent") if isinstance(export_data.get("agent"), dict) else {}
275- compile_for_runtime = getattr(compiler, "compile_for_runtime", None)
276- if callable(compile_for_runtime):
277- compile_result = await compile_for_runtime(
278- config=export_data,
279- model_overrides=model_overrides or None,
280- current_user=current_user,
281- )
282- runtime_config = normalize_runtime_config_for_react_agent(compile_result["runtime_config"])
283- _prepend_configs_system_prompt_first(export_agent, runtime_config)
284- return (compile_result["agent_card"], runtime_config)
285- 
286- compiled = await compiler.compile_with_overrides_config(
287- config=export_data,
288- model_overrides=model_overrides,
289- current_user=current_user,
290- )
291- agent_config_dict = compiled["agent_config"]
292- runtime_config = normalize_runtime_config_for_react_agent(_adapt_runtime_config(agent_config_dict))
293- _prepend_configs_system_prompt_first(export_agent, runtime_config)
294- return (_agent_card_from_export_agent(export_agent), runtime_config)
295- 
296- 
297-async def build_react_agent_from_export_data(
298- export_data: dict[str, Any],
299- current_user: dict[str, Any],
300- *,
301- model_overrides: dict[str, Any] | None = None,
302-) -> ReActAgent:
303- """由已解析的导出数据构建 ReActAgent。"""
304- export_agent = export_data.get("agent") if isinstance(export_data.get("agent"), dict) else {}
305- agent_card, runtime_config = await _compile_runtime_config_from_export_data(
306- export_data,
307- current_user,
308- model_overrides or None,
309- )
310- agent = ReActAgent(card=agent_card)
311- agent.configure(runtime_config)
312- await _register_memory_rail_from_export(agent, export_agent)
313- return agent
314- 
315- 
316-async def build_react_agent(
317- ir_path: Path,
318- current_user: dict[str, Any],
319- *,
320- model_overrides: dict[str, Any] | None = None,
321-) -> ReActAgent:
322- """由 IR 文件构建 ReActAgent。"""
323- export_data = json.loads(ir_path.read_text(encoding="utf-8"))
324- return await build_react_agent_from_export_data(
325- export_data,
326- current_user,
327- model_overrides=model_overrides,
328- )
329- 
330- 
331-async def build_react_agent_from_ir(ir_path: Path, current_user: dict[str, Any]) -> ReActAgent:
332- """读取 IR 文件,按进程环境补齐模型覆盖后构建 ReActAgent。"""
333- export_data = json.loads(ir_path.read_text(encoding="utf-8"))
334- model_overrides = build_model_overrides_from_default_llm_env(export_data)
335- return await build_react_agent_from_export_data(
336- export_data,
337- current_user,
338- model_overrides=model_overrides or None,
339- )
340- 
341- 
342-async def build_react_agent_from_ir_dict(
343- ir_root: dict[str, Any],
344- current_user: dict[str, Any],
345-) -> ReActAgent:
346- """由已解析的 IR 根对象构建 ReActAgent,模型覆盖规则与按文件读取路径一致。"""
347- model_overrides = build_model_overrides_from_default_llm_env(ir_root)
348- return await build_react_agent_from_export_data(
349- ir_root,
350- current_user,
351- model_overrides=model_overrides or None,
352- )
@@ -1,6 +0,0 @@
1-from __future__ import annotations
2- 
3-from .runtime_bootstrap import ensure_runtime_ready
4- 
5-__all__ = ["ensure_runtime_ready"]
6- 
@@ -1,274 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved.
3- 
4-from __future__ import annotations
5- 
6-import json
7-import logging
8-import os
9-import socket
10-from dataclasses import dataclass
11-from datetime import datetime, timezone
12-from enum import Enum
13-from logging.handlers import RotatingFileHandler
14-from pathlib import Path
15-from typing import Any
16- 
17-from openjiuwen_runtime.foundation.log import get_logger
18-from .core_log_bridge import CoreLogEvent, register_core_log_sink
19- 
20-_LOG = get_logger(__name__)
21- 
22- 
23-class AlarmServerName(str, Enum):
24- IR_EXECUTION_SERVICE = "ir_execution_service"
25- OBS = "obs"
26- REDIS = "redis"
27- LLM = "llm"
28- TOOL = "tool"
29- 
30- 
31-class AlarmSeverity(str, Enum):
32- CRITICAL = "critical"
33- MAJOR = "major"
34- MINOR = "minor"
35- WARNING = "warning"
36- INDETERMINATE = "indeterminate"
37- CLEARED = "cleared"
38- 
39- 
40-def _now_alarm_time_str() -> str:
41- # Follow sample in dfx.md: "20221108 10:23:13"
42- return datetime.now(timezone.utc).strftime("%Y%m%d %H:%M:%S")
43- 
44- 
45-def _is_loopback(ip: str) -> bool:
46- ip = (ip or "").strip()
47- return ip in {"127.0.0.1", "::1"} or ip.startswith("127.")
48- 
49- 
50-def resolve_local_ip() -> str:
51- """Resolve a non-loopback local ip, best-effort."""
52- env_ip = (os.environ.get("IP") or "").strip()
53- if env_ip and not _is_loopback(env_ip):
54- return env_ip
55- try:
56- s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
57- try:
58- s.connect(("8.8.8.8", 80))
59- ip = s.getsockname()[0]
60- return ip if ip and not _is_loopback(ip) else ""
61- finally:
62- s.close()
63- except Exception:
64- try:
65- ip = socket.gethostbyname(socket.gethostname())
66- return ip if ip and not _is_loopback(ip) else ""
67- except Exception:
68- return ""
69- 
70- 
71-def map_level_from_python(levelno: int) -> AlarmSeverity:
72- # Required mapping:
73- # CRITICAL -> critical
74- # ERROR -> major
75- # WARNING -> minor
76- if levelno >= logging.CRITICAL:
77- return AlarmSeverity.CRITICAL
78- if levelno >= logging.ERROR:
79- return AlarmSeverity.MAJOR
80- return AlarmSeverity.MINOR
81- 
82- 
83-@dataclass(frozen=True, slots=True)
84-class AlarmLogRecord:
85- timestamp: str
86- server_name: str
87- ip: str
88- level: str
89- module: str
90- message: str
91- 
92- def to_json_line(self) -> str:
93- return json.dumps(
94- {
95- "timestamp": self.timestamp,
96- "server_name": self.server_name,
97- "ip": self.ip,
98- "level": self.level,
99- "module": self.module,
100- "message": self.message,
101- },
102- ensure_ascii=False,
103- default=str,
104- )
105- 
106- 
107-class AlarmLogger:
108- def __init__(self, *, log_file: Path) -> None:
109- self._logger = logging.getLogger("ir_execution_service.alarm")
110- self._logger.setLevel(logging.INFO)
111- self._logger.propagate = False
112- 
113- for h in list(self._logger.handlers):
114- self._logger.removeHandler(h)
115- try:
116- h.close()
117- except Exception as exc:
118- _LOG.warning("failed to close logging handler: %s", exc)
119- 
120- log_file.parent.mkdir(parents=True, exist_ok=True)
121- max_bytes = _env_int("LOWCODE_ALARM_LOG_MAX_BYTES", 20 * 1024 * 1024)
122- backup_count = _env_int("LOWCODE_ALARM_LOG_BACKUP_COUNT", 20)
123- handler = RotatingFileHandler(
124- filename=str(log_file),
125- maxBytes=max_bytes,
126- backupCount=backup_count,
127- encoding="utf-8",
128- )
129- handler.setLevel(logging.INFO)
130- handler.setFormatter(logging.Formatter("%(message)s"))
131- self._logger.addHandler(handler)
132- 
133- _LOG.info("Alarm logger ready: %s", log_file)
134- 
135- def write(self, record: AlarmLogRecord) -> None:
136- self._logger.info(record.to_json_line())
137- 
138- 
139-_ALARM: AlarmLogger | None = None
140-_ALARM_BRIDGE_INSTALLED = False
141-_ALARM_TOOL_CALLBACK_INSTALLED = False
142- 
143- 
144-def _env_int(name: str, default: int) -> int:
145- raw = (os.environ.get(name) or "").strip()
146- if not raw:
147- return default
148- try:
149- value = int(raw)
150- return value if value > 0 else default
151- except Exception:
152- return default
153- 
154- 
155-def init_alarm_logger_from_env() -> AlarmLogger | None:
156- """Initialize alarm logger if path env is set. Env must include filename."""
157- global _ALARM
158- if _ALARM is not None:
159- return _ALARM
160- 
161- raw = (os.environ.get("LOWCODE_ALARM_LOG_PATH") or "").strip()
162- if not raw:
163- _LOG.info("Alarm logger disabled (LOWCODE_ALARM_LOG_PATH not set).")
164- return None
165- p = Path(raw).expanduser().resolve()
166- _ALARM = AlarmLogger(log_file=p)
167- return _ALARM
168- 
169- 
170-def log_alarm(
171- *,
172- server_name: AlarmServerName,
173- level: AlarmSeverity,
174- module: str,
175- message: str,
176- ip: str = "",
177-) -> None:
178- if _ALARM is None:
179- return
180- _ALARM.write(
181- AlarmLogRecord(
182- timestamp=_now_alarm_time_str(),
183- server_name=server_name.value,
184- ip=str(ip or "").strip(),
185- level=level.value,
186- module=str(module or "").strip(),
187- message=str(message or "").strip(),
188- )
189- )
190- 
191- 
192-def install_alarm_log_bridge() -> None:
193- install_core_alarm_sink()
194- 
195- 
196-def install_core_alarm_sink() -> None:
197- global _ALARM_BRIDGE_INSTALLED
198- if _ALARM is None or _ALARM_BRIDGE_INSTALLED:
199- return
200- 
201- def _sink(event: CoreLogEvent) -> None:
202- if event.kind == "llm_upstream":
203- if event.ok:
204- return
205- payload = event.payload or {}
206- msg = str(payload.get("error_message") or payload.get("exception") or event.return_info)
207- log_alarm(
208- server_name=AlarmServerName.LLM,
209- level=AlarmSeverity.MAJOR,
210- module="llm.call",
211- message=msg,
212- ip="",
213- )
214- return
215- 
216- if event.kind != "core_warning":
217- return
218- 
219- if event.logger_name == "llm":
220- server = AlarmServerName.LLM
221- elif event.logger_name == "tool":
222- server = AlarmServerName.TOOL
223- else:
224- return
225- 
226- payload = event.payload or {}
227- msg = event.return_info
228- if server == AlarmServerName.LLM:
229- msg = str(payload.get("error_message") or payload.get("exception") or payload.get("message") or msg)
230- elif server == AlarmServerName.TOOL:
231- msg = str(payload.get("error_message") or payload.get("message") or msg)
232- 
233- log_alarm(
234- server_name=server,
235- level=map_level_from_python(event.levelno),
236- module=event.logger_name,
237- message=msg,
238- ip="",
239- )
240- 
241- register_core_log_sink(_sink)
242- _ALARM_BRIDGE_INSTALLED = True
243- 
244- 
245-def install_runner_tool_alarm_callbacks() -> None:
246- global _ALARM_TOOL_CALLBACK_INSTALLED
247- if _ALARM is None or _ALARM_TOOL_CALLBACK_INSTALLED:
248- return
249- 
250- from openjiuwen.core.runner import Runner
251- from openjiuwen.core.runner.callback.events import ToolCallEvents
252- 
253- fw = Runner.callback_framework
254- 
255- async def _on_tool_error(*, tool_name: str = "", error: BaseException | None = None, **_: Any) -> None:
256- err_text = str(error) if error is not None else ""
257- label = str(tool_name or "").strip()
258- if label and err_text:
259- msg = f"{label} failed: {err_text}"
260- elif label:
261- msg = f"{label} failed"
262- else:
263- msg = err_text or "tool failed"
264- log_alarm(
265- server_name=AlarmServerName.TOOL,
266- level=AlarmSeverity.MAJOR,
267- module="tool.call",
268- message=msg,
269- ip="",
270- )
271- 
272- fw.register_sync(ToolCallEvents.TOOL_CALL_ERROR, _on_tool_error, priority=-1100)
273- _ALARM_TOOL_CALLBACK_INSTALLED = True
274- 
@@ -1,340 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved.
3- 
4-"""
5-SessionModelContext persistence for ir_execution_service.
6- 
7-- Persist ONLY user input messages (role=user) into Redis as JSON.
8-- Key prefix / TTL / lock TTL are controlled by environment variables.
9-- A simple Redis distributed lock is used to prevent concurrent execution
10- or stale overwrites for the same conversation_id.
11- 
12-This module intentionally does NOT re-implement SessionModelContext. It only
13-creates / restores it by using the SDK implementation.
14-"""
15- 
16-from __future__ import annotations
17- 
18-import json
19-import os
20-import socket
21-import time
22-from contextlib import asynccontextmanager
23-from dataclasses import dataclass
24-from typing import Any, AsyncIterator, Iterable, Optional
25-from urllib.parse import urlparse
26- 
27-from openjiuwen_runtime.foundation.log import get_logger
28- 
29-from .alarm_logger import AlarmServerName, AlarmSeverity, log_alarm
30-from .interface_logger import log_client
31-from .runtime_env import clean_env_value, get_int_env
32- 
33-_log = get_logger(__name__)
34- 
35- 
36-def _env_prefix() -> str:
37- # Use a fixed prefix to avoid cross-business key conflicts.
38- # Default is explicit to indicate "workflow dialogue context".
39- raw = clean_env_value("LOWCODE_CONTEXT_REDIS_KEY_PREFIX", "lowcode:workflow_context")
40- raw = raw.strip(":").strip()
41- return raw or "lowcode:workflow_context"
42- 
43- 
44-def _env_ttl_seconds() -> int:
45- # TTL for persisted context state; refreshed on each save.
46- return max(1, get_int_env("LOWCODE_CONTEXT_REDIS_TTL_SECONDS", 3600))
47- 
48- 
49-def _env_lock_ttl_seconds() -> int:
50- # Lock TTL to avoid deadlock when a worker crashes.
51- return max(1, get_int_env("LOWCODE_CONTEXT_REDIS_LOCK_TTL_SECONDS", 120))
52- 
53- 
54-def _redis_url() -> str:
55- # 每个业务场景只用自己的 Redis URL;对话上下文使用默认/基础 Redis(可配前缀避免冲突)。
56- url = clean_env_value("LOWCODE_DEFAULT_REDIS_URL")
57- if not url:
58- raise RuntimeError("Context persistence requires LOWCODE_DEFAULT_REDIS_URL.")
59- return url
60- 
61- 
62-def _redis_dest_ip() -> str:
63- url = _redis_url()
64- try:
65- host = urlparse(url).hostname or ""
66- return socket.gethostbyname(host) if host else ""
67- except Exception:
68- return ""
69- 
70- 
71-def _context_state_key(conversation_id: str, context_id: str) -> str:
72- # ctx:{prefix}:{conversation}:{context}
73- return f"{_env_prefix()}:{conversation_id}:state:{context_id}"
74- 
75- 
76-def _lock_key(conversation_id: str) -> str:
77- return f"{_env_prefix()}:{conversation_id}:lock"
78- 
79- 
80-# Only delete the lock if the value still matches our token (avoids deleting a successor lock after TTL expiry).
81-_UNLOCK_IF_TOKEN_MATCHES_LUA = """
82-if redis.call("get", KEYS[1]) == ARGV[1] then
83- return redis.call("del", KEYS[1])
84-else
85- return 0
86-end
87-"""
88- 
89- 
90-def _now_ms() -> int:
91- return int(time.time() * 1000)
92- 
93- 
94-@dataclass(frozen=True, slots=True)
95-class PersistedContextState:
96- """JSON-friendly payload saved in Redis."""
97- 
98- # Outer mapping expected by SessionModelContext.load_state:
99- # {context_id: {"messages": [...], "offload_messages": {...}}}
100- states: dict[str, Any]
101- 
102- 
103-def _filter_user_messages(messages: Iterable[Any]) -> list[dict[str, Any]]:
104- """Keep ONLY role=user messages; serialize to JSON-friendly dict."""
105- out: list[dict[str, Any]] = []
106- for m in messages or []:
107- role = getattr(m, "role", None)
108- if role != "user":
109- continue
110- dump = getattr(m, "model_dump", None)
111- if callable(dump):
112- d = dump()
113- if isinstance(d, dict):
114- # Ensure role/content exist for reconstruction.
115- d.setdefault("role", "user")
116- out.append(d)
117- continue
118- content = getattr(m, "content", None)
119- out.append({"role": "user", "content": content})
120- return out
121- 
122- 
123-def _restore_user_messages(message_dicts: list[dict[str, Any]]) -> list[Any]:
124- """Rebuild user messages for SessionModelContext history."""
125- from openjiuwen.core.foundation.llm import UserMessage
126- 
127- restored: list[Any] = []
128- for d in message_dicts or []:
129- if not isinstance(d, dict):
130- continue
131- content = d.get("content")
132- # UserMessage accepts role/content; extra fields are ignored by pydantic if not declared.
133- restored.append(UserMessage(role="user", content=content))
134- return restored
135- 
136- 
137-class RedisContextPersistence:
138- """Persist/restore SessionModelContext state by (conversation_id, context_id)."""
139- 
140- def __init__(self) -> None:
141- self._redis = None
142- 
143- def _get_redis(self):
144- if self._redis is not None:
145- return self._redis
146- try:
147- from redis.asyncio import Redis # type: ignore
148- except Exception as e: # pragma: no cover
149- raise RuntimeError("redis-py (redis.asyncio) is required for context persistence") from e
150- url = _redis_url()
151- self._redis = Redis.from_url(url, decode_responses=True)
152- safe = url.split("@")[-1] if "@" in url else url
153- _log.info("Context persistence Redis client ready (%s)", safe)
154- return self._redis
155- 
156- async def load_context(self, *, conversation_id: str, context_id: str, config: Any) -> Any:
157- """
158- Create SessionModelContext and restore persisted user messages (if any).
159- 
160- Returns:
161- SessionModelContext instance (SDK type).
162- """
163- from openjiuwen.core.context_engine.context.context import SessionModelContext
164- 
165- ctx = SessionModelContext(
166- context_id=context_id,
167- session_id=conversation_id,
168- config=config,
169- history_messages=[],
170- processors=[],
171- )
172- 
173- key = _context_state_key(conversation_id, context_id)
174- t0 = time.perf_counter()
175- raw = await self._get_redis().get(key)
176- log_client(
177- interface_name="redis.get",
178- cost_ms=(time.perf_counter() - t0) * 1000.0,
179- ok=True,
180- return_code=0,
181- return_info="hit" if raw else "miss",
182- dest_ip=_redis_dest_ip(),
183- add_info={"key": key, "scene": "workflow_context"},
184- )
185- if not raw:
186- return ctx
187- 
188- try:
189- payload = json.loads(raw)
190- except Exception:
191- _log.warning("Context state JSON decode failed, key=%s", key, exc_info=True)
192- log_alarm(
193- server_name=AlarmServerName.REDIS,
194- level=AlarmSeverity.MINOR,
195- module="redis.get",
196- message=f"Context state JSON decode failed, key={key}",
197- ip=_redis_dest_ip(),
198- )
199- return ctx
200- 
201- # Expected format: {"context_id": {"messages":[...], "offload_messages":{...}}}
202- if not isinstance(payload, dict):
203- return ctx
204- ctx_bucket = payload.get(context_id)
205- if not isinstance(ctx_bucket, dict):
206- return ctx
207- msg_dicts = ctx_bucket.get("messages", [])
208- if not isinstance(msg_dicts, list):
209- msg_dicts = []
210- restored_user_messages = _restore_user_messages(msg_dicts) # only user messages
211- if restored_user_messages:
212- # Seed as history; we intentionally do not restore offload cache.
213- ctx = SessionModelContext(
214- context_id=context_id,
215- session_id=conversation_id,
216- config=config,
217- history_messages=restored_user_messages,
218- processors=[],
219- )
220- return ctx
221- 
222- async def save_on_interaction(self, *, conversation_id: str, context_id: str, context: Any) -> None:
223- """Persist current user-message history when workflow yields interaction."""
224- try:
225- saved = context.save_state()
226- except Exception:
227- _log.warning("Context save_state failed; skip persistence", exc_info=True)
228- log_alarm(
229- server_name=AlarmServerName.IR_EXECUTION_SERVICE,
230- level=AlarmSeverity.MINOR,
231- module="context.save_state",
232- message="Context save_state failed; skip persistence",
233- ip="",
234- )
235- return
236- 
237- # Keep ONLY user messages.
238- msgs = []
239- if isinstance(saved, dict):
240- msgs = saved.get("messages", []) or []
241- user_msgs = _filter_user_messages(msgs)
242- 
243- payload = {context_id: {"messages": user_msgs, "offload_messages": {}}}
244- key = _context_state_key(conversation_id, context_id)
245- t0 = time.perf_counter()
246- await self._get_redis().set(key, json.dumps(payload, ensure_ascii=False), ex=_env_ttl_seconds())
247- log_client(
248- interface_name="redis.set",
249- cost_ms=(time.perf_counter() - t0) * 1000.0,
250- ok=True,
251- return_code=0,
252- return_info="saved",
253- dest_ip=_redis_dest_ip(),
254- add_info={"key": key, "ttl": _env_ttl_seconds(), "scene": "workflow_context"},
255- )
256- 
257- async def delete(self, *, conversation_id: str, context_id: str) -> None:
258- key = _context_state_key(conversation_id, context_id)
259- try:
260- t0 = time.perf_counter()
261- await self._get_redis().delete(key)
262- log_client(
263- interface_name="redis.delete",
264- cost_ms=(time.perf_counter() - t0) * 1000.0,
265- ok=True,
266- return_code=0,
267- return_info="deleted",
268- dest_ip=_redis_dest_ip(),
269- add_info={"key": key, "scene": "workflow_context"},
270- )
271- except Exception:
272- _log.warning("Context delete failed, key=%s", key, exc_info=True)
273- log_alarm(
274- server_name=AlarmServerName.REDIS,
275- level=AlarmSeverity.MAJOR,
276- module="redis.delete",
277- message=f"Context delete failed, key={key}",
278- ip=_redis_dest_ip(),
279- )
280- 
281- @asynccontextmanager
282- async def conversation_lock(self, *, conversation_id: str) -> AsyncIterator[None]:
283- """
284- Acquire a simple distributed lock for a conversation.
285- 
286- Strategy: SET lock_key value NX EX ttl
287- """
288- redis = self._get_redis()
289- key = _lock_key(conversation_id)
290- ttl = _env_lock_ttl_seconds()
291- token = f"{os.getpid()}-{_now_ms()}"
292- 
293- acquired = False
294- try:
295- # redis-py returns True/False for set(..., nx=True)
296- t0 = time.perf_counter()
297- acquired = bool(await redis.set(key, token, ex=ttl, nx=True))
298- log_client(
299- interface_name="redis.setnx",
300- cost_ms=(time.perf_counter() - t0) * 1000.0,
301- ok=acquired,
302- return_code=0 if acquired else 1,
303- return_info="acquired" if acquired else "busy",
304- dest_ip=_redis_dest_ip(),
305- add_info={"key": key, "ttl": ttl, "scene": "workflow_context_lock"},
306- )
307- if not acquired:
308- log_alarm(
309- server_name=AlarmServerName.REDIS,
310- level=AlarmSeverity.MINOR,
311- module="redis.setnx",
312- message=f"conversation lock busy: {conversation_id}",
313- ip=_redis_dest_ip(),
314- )
315- raise RuntimeError(f"conversation lock busy: {conversation_id}")
316- yield
317- finally:
318- if acquired:
319- try:
320- t1 = time.perf_counter()
321- await redis.eval(_UNLOCK_IF_TOKEN_MATCHES_LUA, 1, key, token)
322- log_client(
323- interface_name="redis.eval_unlock",
324- cost_ms=(time.perf_counter() - t1) * 1000.0,
325- ok=True,
326- return_code=0,
327- return_info="released",
328- dest_ip=_redis_dest_ip(),
329- add_info={"key": key, "scene": "workflow_context_lock"},
330- )
331- except Exception:
332- _log.warning("Lock release failed, key=%s", key, exc_info=True)
333- log_alarm(
334- server_name=AlarmServerName.REDIS,
335- level=AlarmSeverity.MAJOR,
336- module="redis.eval_unlock",
337- message=f"Lock release failed, key={key}",
338- ip=_redis_dest_ip(),
339- )
340- 
@@ -1,166 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved.
3- 
4-from __future__ import annotations
5- 
6-import json
7-import logging
8-from dataclasses import dataclass
9-from typing import Any, Callable
10- 
11-_MOD_LOG = logging.getLogger(__name__)
12- 
13- 
14-@dataclass(frozen=True, slots=True)
15-class CoreLogEvent:
16- kind: str
17- logger_name: str
18- levelno: int
19- request_id: str
20- interface_name: str
21- ok: bool | None
22- return_info: str
23- payload: dict[str, Any] | None
24- raw_message: str
25- 
26- 
27-_SINKS: list[Callable[[CoreLogEvent], None]] = []
28-_BRIDGE_INSTALLED = False
29- 
30- 
31-def register_core_log_sink(sink: Callable[[CoreLogEvent], None]) -> None:
32- for cb in _SINKS:
33- if cb is sink:
34- return
35- _SINKS.append(sink)
36- 
37- 
38-def _emit_to_sinks(event: CoreLogEvent) -> None:
39- for cb in list(_SINKS):
40- try:
41- cb(event)
42- except Exception:
43- _MOD_LOG.warning("core log sink callback raised", exc_info=True)
44- continue
45- 
46- 
47-def _is_upstream_llm_http_completion_end(payload: dict[str, Any]) -> bool:
48- msg = str(payload.get("message") or "")
49- if "API response received." in msg:
50- return True
51- md = payload.get("metadata")
52- if not isinstance(md, dict):
53- return False
54- resp = md.get("response")
55- if isinstance(resp, dict) and isinstance(resp.get("choices"), list):
56- return True
57- if isinstance(resp, str) and "choices=" in resp and "ChatCompletion(" in resp:
58- return True
59- return False
60- 
61- 
62-def _is_upstream_llm_hard_failure(payload: dict[str, Any]) -> bool:
63- msg = str(payload.get("message") or payload.get("error_message") or "")
64- lower = msg.lower()
65- if "failed to decode json from llm output" in lower:
66- return False
67- if "unsupported llm_output type for parse" in lower:
68- return False
69- if "stream parser attempt error" in lower:
70- return False
71- if "api async invoke error" in lower:
72- return True
73- if "api async stream error" in lower:
74- return True
75- if "api invoke error" in lower:
76- return True
77- if "invoke error" in lower and "parser" not in lower:
78- return True
79- if "kv cache release failed" in lower:
80- return True
81- if "kv cache release error" in lower:
82- return True
83- return False
84- 
85- 
86-class CoreLogBridge(logging.Handler):
87- def emit(self, record: logging.LogRecord) -> None:
88- logger_name = str(record.name or "")
89- if logger_name not in {"llm", "tool"}:
90- return
91- try:
92- msg = record.getMessage()
93- except Exception:
94- _MOD_LOG.warning("failed to format core bridge log record", exc_info=True)
95- return
96- 
97- payload: dict[str, Any] | None = None
98- if msg and msg[0] == "{":
99- try:
100- parsed = json.loads(msg)
101- payload = parsed if isinstance(parsed, dict) else None
102- except Exception:
103- payload = None
104- 
105- if logger_name == "llm" and payload is not None:
106- event_type = payload.get("event_type")
107- if event_type == "llm_call_end" and _is_upstream_llm_http_completion_end(payload):
108- rid = str(payload.get("trace_id") or payload.get("session_id") or "")
109- _emit_to_sinks(
110- CoreLogEvent(
111- kind="llm_upstream",
112- logger_name=logger_name,
113- levelno=record.levelno,
114- request_id=rid,
115- interface_name="llm.call",
116- ok=True,
117- return_info=str(payload.get("message") or ""),
118- payload=payload,
119- raw_message=msg,
120- )
121- )
122- return
123- if event_type == "llm_call_error" and _is_upstream_llm_hard_failure(payload):
124- rid = str(payload.get("trace_id") or payload.get("session_id") or "")
125- _emit_to_sinks(
126- CoreLogEvent(
127- kind="llm_upstream",
128- logger_name=logger_name,
129- levelno=record.levelno,
130- request_id=rid,
131- interface_name="llm.call",
132- ok=False,
133- return_info=str(payload.get("message") or payload.get("error_message") or ""),
134- payload=payload,
135- raw_message=msg,
136- )
137- )
138- return
139- 
140- if record.levelno >= logging.WARNING:
141- _emit_to_sinks(
142- CoreLogEvent(
143- kind="core_warning",
144- logger_name=logger_name,
145- levelno=record.levelno,
146- request_id="",
147- interface_name="",
148- ok=None,
149- return_info=msg,
150- payload=payload,
151- raw_message=msg,
152- )
153- )
154- 
155- 
156-def install_core_log_bridge() -> None:
157- global _BRIDGE_INSTALLED
158- if _BRIDGE_INSTALLED:
159- return
160- bridge = CoreLogBridge()
161- # 只挂 root:子 logger(llm/tool)默认 propagate=True,记录会冒泡一次,避免与再挂子 logger 导致 emit 双份。
162- # CoreLogBridge.emit 已按 record.name 过滤非 llm/tool。
163- root = logging.getLogger()
164- root.addHandler(bridge)
165- _BRIDGE_INSTALLED = True
166- 
@@ -1,51 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved
3- 
4-from __future__ import annotations
5- 
6-import logging
7-import os
8-from pathlib import Path
9- 
10- 
11-def setup_error_file_logging() -> Path:
12- """为本应用单独落盘错误日志(ERROR+),避免被大量 INFO 淹没。
13- 
14- 默认路径:<app_root>/logs/error.log
15- 可通过环境变量 IR_ERROR_LOG_PATH 覆盖。
16- """
17- 
18- # Service root: .../applications/ir_execution_service
19- app_root = Path(__file__).resolve().parent.parent.parent
20- default_path = app_root / "logs" / "error.log"
21- raw = (os.environ.get("IR_ERROR_LOG_PATH") or "").strip()
22- log_path = Path(raw).expanduser().resolve() if raw else default_path.resolve()
23- log_path.parent.mkdir(parents=True, exist_ok=True)
24- 
25- handler = logging.FileHandler(str(log_path), encoding="utf-8")
26- handler.setLevel(logging.ERROR)
27- handler.setFormatter(
28- logging.Formatter(
29- fmt="%(asctime)s | %(levelname)s | %(name)s | %(message)s",
30- datefmt="%Y-%m-%d %H:%M:%S",
31- )
32- )
33- 
34- def _already_added(logger_obj: logging.Logger) -> bool:
35- for h in logger_obj.handlers:
36- if isinstance(h, logging.FileHandler):
37- h_base = getattr(h, "baseFilename", None)
38- hand_base = getattr(handler, "baseFilename", None)
39- if h_base == hand_base:
40- return True
41- return False
42- 
43- # 尽量覆盖:根 logger、uvicorn、以及 openjiuwen 的 logger 层级
44- for name in ("", "uvicorn", "uvicorn.error", "uvicorn.access", "openjiuwen"):
45- lg = logging.getLogger(name)
46- if not _already_added(lg):
47- lg.addHandler(handler)
48- if lg.level > logging.ERROR:
49- lg.setLevel(logging.ERROR)
50- 
51- return log_path
@@ -1,164 +0,0 @@
1-#!/usr/bin/env python
2-# coding: utf-8
3-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved
4- 
5-from __future__ import annotations
6- 
7-import json
8-import os
9-from dataclasses import dataclass
10-from typing import Any
11- 
12-from fastapi import HTTPException
13- 
14-from .http_response_contract import LowcodeApiResponseCode
15-from .ir_resolver import (
16- detect_executable_kind,
17- ensure_ir_root,
18- lowcode_code_from_http_exception,
19-)
20- 
21- 
22-class ExecutionPrepareError(Exception):
23- def __init__(self, code: LowcodeApiResponseCode, message: str):
24- super().__init__(message)
25- self.code = code
26- self.message = message
27- 
28- 
29-@dataclass(slots=True)
30-class PreparedExecutionRequest:
31- ir_root: dict[str, Any]
32- executable_kind: str
33- inputs_obj: Any
34- space_id: str
35- current_user: dict[str, Any]
36- session_id: str
37- timeout_seconds: float
38- 
39- 
40-def _detect_executable_kind(ir_root: dict[str, Any]) -> str:
41- try:
42- return detect_executable_kind(ir_root)
43- except HTTPException as exc:
44- code, message = lowcode_code_from_http_exception(exc)
45- if exc.status_code == 400 and "neither workflow" in (message or "").lower():
46- code = LowcodeApiResponseCode.IR_INVALID
47- raise ExecutionPrepareError(code, message) from exc
48- 
49- 
50-def _decode_inputs(inputs_raw: str) -> dict[str, Any]:
51- try:
52- inputs_obj = json.loads(inputs_raw)
53- except (TypeError, json.JSONDecodeError) as exc:
54- c = LowcodeApiResponseCode.INVALID_INPUTS
55- raise ExecutionPrepareError(c, f"{c.default_message}: {exc}") from exc
56- if not isinstance(inputs_obj, dict):
57- c = LowcodeApiResponseCode.INVALID_INPUTS
58- raise ExecutionPrepareError(c, "inputs must decode to a JSON object")
59- return inputs_obj
60- 
61- 
62-# Lowcode workflow IR: input component uses numeric type 8 (COMPONENT_TYPE_INPUT).
63-_IR_INPUT_COMPONENT_TYPE = 8
64- 
65- 
66-def _unwrap_interactive_reply_value(value: Any) -> Any:
67- """客户端若把嵌套对象二次编码成字符串(如 PowerShell ConvertTo-Json 默认 Depth=2),此处尽量还原为 dict。"""
68- while isinstance(value, str):
69- s = value.strip()
70- if len(s) < 2 or s[0] != "{":
71- break
72- try:
73- parsed = json.loads(s)
74- except json.JSONDecodeError:
75- break
76- if not isinstance(parsed, dict):
77- break
78- value = parsed
79- return value
80- 
81- 
82-def _find_ir_component(ir_root: dict[str, Any], comp_id: str) -> dict[str, Any] | None:
83- components = ir_root.get("components")
84- if not isinstance(components, list):
85- return None
86- for item in components:
87- if isinstance(item, dict) and item.get("id") == comp_id:
88- return item
89- return None
90- 
91- 
92-def _coerce_workflow_interactive_reply_value(ir_root: dict[str, Any], node_id: str, value: Any) -> Any:
93- """Input 节点要求 dict(字段名 -> 值);单字段时可把纯字符串包成 dict,与提问器式续跑兼容。"""
94- if isinstance(value, dict) or value is None:
95- return value
96- comp = _find_ir_component(ir_root, node_id)
97- if comp is None or comp.get("type") != _IR_INPUT_COMPONENT_TYPE:
98- return value
99- configs = comp.get("configs") if isinstance(comp.get("configs"), dict) else {}
100- fields = configs.get("inputs")
101- if not isinstance(fields, list) or len(fields) != 1:
102- return value
103- field0 = fields[0]
104- if not isinstance(field0, dict):
105- return value
106- name = field0.get("input_name")
107- if not isinstance(name, str) or not name:
108- return value
109- return {name: value}
110- 
111- 
112-def _normalize_workflow_resume_inputs(inputs_obj: dict[str, Any], ir_root: dict[str, Any]) -> Any:
113- if set(inputs_obj.keys()) != {"__interactive_reply"}:
114- return inputs_obj
115- 
116- from openjiuwen.core.session import InteractiveInput
117- 
118- reply = inputs_obj["__interactive_reply"]
119- if isinstance(reply, dict) and (reply.get("id") is not None):
120- node_id = str(reply.get("id"))
121- value = _unwrap_interactive_reply_value(reply.get("value"))
122- value = _coerce_workflow_interactive_reply_value(ir_root, node_id, value)
123- interactive_input = InteractiveInput()
124- interactive_input.update(node_id, value)
125- return interactive_input
126- return InteractiveInput(reply)
127- 
128- 
129-def _prepare_inputs_for_kind(
130- inputs_obj: dict[str, Any], executable_kind: str, user_id: str, ir_root: dict[str, Any]
131-) -> Any:
132- if executable_kind == "workflow":
133- return _normalize_workflow_resume_inputs(inputs_obj, ir_root)
134- if "user_id" in inputs_obj:
135- return inputs_obj
136- enriched = dict(inputs_obj)
137- enriched["user_id"] = user_id
138- return enriched
139- 
140- 
141-async def prepare_execution_request(body: Any) -> PreparedExecutionRequest:
142- try:
143- ir_root = await ensure_ir_root(getattr(body, "ir_path"))
144- except HTTPException as exc:
145- code, message = lowcode_code_from_http_exception(exc)
146- raise ExecutionPrepareError(code, message) from exc
147- 
148- executable_kind = _detect_executable_kind(ir_root)
149- user_id = str(getattr(body, "user_id"))
150- inputs_obj = _prepare_inputs_for_kind(
151- _decode_inputs(getattr(body, "inputs")), executable_kind, user_id, ir_root
152- )
153- space_id = os.environ.get("WORKFLOW_SPACE_ID", "default")
154- 
155- return PreparedExecutionRequest(
156- ir_root=ir_root,
157- executable_kind=executable_kind,
158- inputs_obj=inputs_obj,
159- space_id=space_id,
160- current_user={"user_id": user_id, "space_id": space_id},
161- session_id=getattr(body, "conversation_id"),
162- timeout_seconds=getattr(body, "timeout_ms") / 1000.0,
163- )
164- 
@@ -1,26 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved
3- 
4-"""IR 内部兼容别名:实际实现位于 foundation.db.dialects.gaussdb_asyncgaussdb。
5- 
6-保留本模块是为了 IR 作为独立部署单元时的包内 import 路径稳定
7-(`from .gaussdb_sqlalchemy_dialect import ensure_gaussdb_dialect_registered`),
8-同时避免与 foundation 维护两份几乎相同的方言实现造成代码漂移。
9-"""
10-from __future__ import annotations
11- 
12-from openjiuwen_runtime.foundation.db.dialects.gaussdb_asyncgaussdb import (
13- AsyncAdapt_async_gaussdb_dbapi,
14- PGDialect_async_gaussdb,
15- dialect,
16- ensure_async_gaussdb_installed,
17- ensure_gaussdb_dialect_registered,
18-)
19- 
20-__all__ = [
21- "AsyncAdapt_async_gaussdb_dbapi",
22- "PGDialect_async_gaussdb",
23- "dialect",
24- "ensure_async_gaussdb_installed",
25- "ensure_gaussdb_dialect_registered",
26-]
@@ -1,115 +0,0 @@
1-# coding: utf-8
2-# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved
3- 
4-"""低码 Runner HTTP 响应体约定。"""
5- 
6-from __future__ import annotations
7- 
8-from enum import Enum, IntEnum
9-from typing import Any
10- 
11-from openjiuwen_studio.schemas import ResponseModel
12- 
13- 
14-class LowcodeApiResponseCode(IntEnum):
15- """低码 Runner HTTP 接口使用的数值 code 及默认英文 message。"""
16- 
17- def __new__(cls, value: int, default_message: str):
18- obj = int.__new__(cls, value)
19- obj._value_ = value
20- obj.default_message = default_message
21- return obj
22- 
23- # 成功
24- SUCCESS = (0, "success")
25- 
26- # 1xxx - 参数错误(客户端问题)
27- INVALID_REQUEST = (1001, "invalid request body")
28- MISSING_PARAM = (1002, "missing required param: {field}")
29- INVALID_PARAM = (1003, "invalid param: {field}")
30- INVALID_IR_PATH = (1004, "invalid ir_path format")
31- INVALID_INPUTS = (1005, "inputs is not valid json string")
32- INVALID_TIMEOUT = (1006, "timeout_ms must be positive integer")
33- 
34- # 2xxx - 资源错误(加载、找不到)
35- IR_NOT_FOUND = (2001, "ir not found: {ir_path}")
36- IR_DOWNLOAD_FAILED = (2002, "failed to download ir")
37- IR_INVALID = (2003, "invalid ir format")
38- IR_LOAD_FAILED = (2004, "failed to load ir")
39- SESSION_LOAD_FAILED = (2005, "failed to load session")
40- LLM_API_KEY_MISSING = (2006, "LLM api_key missing: set env {env_var} or api_key in DSL")
41- 
42- # 3xxx - 执行错误(运行时)
43- EXECUTION_TIMEOUT = (3001, "execution timeout")
44- EXECUTION_FAILED = (3002, "agent execution failed")
45- EXECUTION_CANCELLED = (3003, "execution cancelled")
46- OUTPUT_INVALID = (3004, "agent output invalid")
47- INVOKE_NOT_SUPPORTED = (3005, "invoke not supported")
48- 
49- # 4xxx - 系统限制(预留,当前不可用)
50- RATE_LIMITED = (4001, "rate limit exceeded")
51- CONCURRENCY_LIMITED = (4002, "too many concurrent executions")
52- QUEUE_FULL = (4003, "execution queue full")
53- RESOURCE_EXHAUSTED = (4004, "server resource exhausted")
54- 
55- # 5xxx - 内部错误(平台问题)
56- INTERNAL_ERROR = (5001, "internal server error")
57- SERVICE_UNAVAILABLE = (5002, "service temporarily unavailable")
58- DEPENDENCY_ERROR = (5003, "dependency service error")
59- 
60- def format_message(self, **kwargs: str) -> str:
61- """使用 default_message 做 str.format;无占位符时忽略 kwargs。"""
62- if not kwargs:
63- return self.default_message
64- return self.default_message.format(**kwargs)
65- 
66- 
67-class ResponseDataType(str, Enum):
68- """ResponseModel.data.type 取值(字符串枚举)。"""
69- 
70- ERROR = "error"
71- STREAM = "stream"
72- RESULT = "result"
73- TRACE = "trace"
74- NODE_OUTPUT = "node_output"
75- INPUT_REQUIRED = "input_required"
76- INTERACTION = "interaction"
77- FORCE_FINISH = "force_finish"
78- UNKNOWN = "unknown"
79- 
80- 
81-def build_error_response_model(
82- code: LowcodeApiResponseCode,
83- *,
84- message: str | None = None,
85- payload: dict[str, Any] | None = None,
86-) -> ResponseModel:
87- msg = message if message is not None else code.default_message
88- body: dict[str, Any] = {"message": msg}
89- if payload:
90- body.update(payload)
91- return ResponseModel(
92- code=int(code),
93- message=msg,
94- data={"type": ResponseDataType.ERROR.value, "payload": body},
95- )
96- 
97- 
98-def to_jsonable(obj: Any) -> Any:
99- """将 core 或 studio 对象转为可 JSON 序列化的 dict、list 或标量。"""
100- if obj is None or isinstance(obj, (str, int, float, bool)):
101- return obj
102- if isinstance(obj, dict):
103- return {str(k): to_jsonable(v) for k, v in obj.items()}
104- if isinstance(obj, (list, tuple)):
105- return [to_jsonable(x) for x in obj]
106- model_dump = getattr(obj, "model_dump", None)
107- if callable(model_dump):
108- try:
109- return model_dump(mode="json")
110- except TypeError:
111- return model_dump()
112- if hasattr(obj, "__dict__"):
113- return {k: to_jsonable(v) for k, v in vars(obj).items() if not k.startswith("_")}
114- return str(obj)
115-