Skip to content

Commit 9fb9b55

Browse files
committed
docs(examples): add conversation/identity_config/custom_tools/streaming_deltas to match Go parity
Task: 1789906941
1 parent 6bd48b7 commit 9fb9b55

8 files changed

Lines changed: 521 additions & 0 deletions

File tree

‎examples/forward/__main__.py‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,19 +4,25 @@
44
from qca import Forward
55

66
from .batch import run as batch
7+
from .conversation import run as conversation
8+
from .identity_config import run as identity_config
79
from .memory import run as memory
810
from .models import run as models
911
from .resources import run as resources
1012
from .schedule import run as schedule
1113
from .session import run as session
14+
from .streaming_deltas import run as streaming_deltas
1215

1316
SCENARIOS = {
1417
"models": models,
1518
"session": session,
19+
"conversation": conversation,
1620
"resources": resources,
21+
"identity_config": identity_config,
1722
"memory": memory,
1823
"schedule": schedule,
1924
"batch": batch,
25+
"streaming_deltas": streaming_deltas,
2026
}
2127

2228

‎examples/forward/conversation.py‎

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,70 @@
1+
"""复用同一个 Session 进行多轮对话,并分页读取会话历史。
2+
3+
运行:python -m examples.forward.conversation
4+
"""
5+
6+
from __future__ import annotations
7+
8+
from examples.common.live import Run, choose_model, marker, name, run_cli
9+
from qca import Forward
10+
11+
from ._cleanup import finish_session
12+
13+
14+
def run(client: Forward, context: Run) -> None:
15+
environment = client.environments.create(name=name("env"), config={"type": "cloud"})
16+
environment_id = context.track("environment", environment.id, lambda: client.environments.archive(environment.id))
17+
18+
identity = client.identities.create(external_id=name("identity"), name="SDK 示例用户")
19+
identity_id = context.track("identity", identity.id, lambda: client.identities.delete(identity.id))
20+
21+
model = choose_model(client.models.list(), context.config.model)
22+
context.output("selected_model", model)
23+
template = client.templates.create(
24+
name=name("template"),
25+
environment_id=environment_id,
26+
model=model,
27+
system="你是一个 SDK 示例助手。必要时调用工具,只使用可实际读取的数据回答问题。",
28+
tools=[{"type": "agent_toolset_20260401"}],
29+
)
30+
template_id = context.track("template", template.id, lambda: client.templates.archive(template.id))
31+
32+
session = client.sessions.create(identity_id=identity_id, template_id=template_id)
33+
session_id = context.track("session", session.id, lambda: finish_session(client, context, session.id))
34+
35+
def ask(prompt: str) -> None:
36+
context.output("user", prompt)
37+
sent = client.sessions.events.send(
38+
session_id,
39+
events=[{"type": "user.message", "content": [{"type": "text", "text": prompt}]}],
40+
extra_headers={"Idempotency-Key": name("event")},
41+
)
42+
if not sent.data or not sent.data[0].id:
43+
raise RuntimeError("Send returned no user event")
44+
with client.sessions.events.stream(
45+
session_id,
46+
extra_headers={"Last-Event-ID": sent.data[0].id},
47+
timeout=context.remaining(),
48+
) as stream:
49+
for event in stream:
50+
context.remaining()
51+
if event.type == "agent.message":
52+
context.output("assistant", event.to_json())
53+
elif event.type in ("session.error", "session.status_terminated"):
54+
raise RuntimeError(f"Session stopped: {event.type}")
55+
elif event.type == "session.status_idle":
56+
break
57+
58+
# 同一个 Session 支持多轮:服务端在 session_id 下保留完整历史,无需客户端携带上文。
59+
code = "project-" + marker()
60+
ask(f"这次项目代号是 {code}。请在本次对话中记住它,不要使用工具或写入记忆库。现在只回复:已记住。")
61+
ask("只根据本次会话上文,告诉我刚才约定的项目代号。只回复代号,不要使用工具。")
62+
63+
# 分页读取已有的用户消息和助手回复,重建对话文字记录。
64+
for event in client.sessions.events.list(session_id, order="asc", limit=100):
65+
if event.type in ("user.message", "agent.message"):
66+
context.output(f"history.{event.type}", event.to_json())
67+
68+
69+
if __name__ == "__main__":
70+
run_cli("forward", Forward, {"conversation": run})
Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
"""两个 Identity 共用一个 Template,写入各自的个性化配置,再读取生效配置并在会话中验证。
2+
3+
运行:python -m examples.forward.identity_config
4+
"""
5+
6+
from __future__ import annotations
7+
8+
from examples.common.live import Run, choose_model, marker, name, run_cli
9+
from qca import Forward
10+
11+
from ._cleanup import finish_session
12+
13+
14+
def run(client: Forward, context: Run) -> None:
15+
environment = client.environments.create(name=name("env"), config={"type": "cloud"})
16+
environment_id = context.track("environment", environment.id, lambda: client.environments.archive(environment.id))
17+
18+
model = choose_model(client.models.list(), context.config.model)
19+
context.output("selected_model", model)
20+
shared, baseline = marker(), marker()
21+
template = client.templates.create(
22+
name=name("template"),
23+
environment_id=environment_id,
24+
model=model,
25+
system="你是一个 SDK 示例助手。必要时调用工具,只使用可实际读取的数据回答问题。",
26+
tools=[{"type": "agent_toolset_20260401"}],
27+
environment_variables={"SDK_SHARED_VALUE": shared, "SDK_PERSONAL_VALUE": baseline},
28+
)
29+
template_id = context.track("template", template.id, lambda: client.templates.archive(template.id))
30+
31+
# 先创建两个 Identity 并写入个性化配置:SDK_SHARED_VALUE 继承模板默认,SDK_PERSONAL_VALUE 被各自覆盖。
32+
identities: list[str] = []
33+
personal: list[str] = []
34+
for index in range(2):
35+
identity = client.identities.create(external_id=name("identity"), name=f"SDK 示例用户 {index + 1}")
36+
identity_id = context.track("identity", identity.id, lambda ref=identity.id: client.identities.delete(ref))
37+
identities.append(identity_id)
38+
value = marker()
39+
personal.append(value)
40+
client.identities.configs.upsert(
41+
template_id,
42+
identity_id=identity_id,
43+
identity_config={"environment_variables": {"SDK_PERSONAL_VALUE": {"op": "set", "value": value}}},
44+
)
45+
effective = client.identities.configs.retrieve_effective(template_id, identity_id=identity_id)
46+
context.output(
47+
f"identity_{index + 1}_effective_env",
48+
effective.session.environment_variables if effective.session else None,
49+
)
50+
51+
# 每个 Identity 各起一个会话,读取自己实际生效的环境变量。
52+
for index, identity_id in enumerate(identities):
53+
session = client.sessions.create(identity_id=identity_id, template_id=template_id)
54+
session_id = context.track("session", session.id, lambda ref=session.id: finish_session(client, context, ref))
55+
prompt = "请使用工具读取 SDK_SHARED_VALUE 和 SDK_PERSONAL_VALUE 两个环境变量,只返回这两个变量的实际值。"
56+
context.output(f"identity_{index + 1}_user", prompt)
57+
sent = client.sessions.events.send(
58+
session_id,
59+
events=[{"type": "user.message", "content": [{"type": "text", "text": prompt}]}],
60+
extra_headers={"Idempotency-Key": name("event")},
61+
)
62+
if not sent.data or not sent.data[0].id:
63+
raise RuntimeError("Send returned no user event")
64+
with client.sessions.events.stream(
65+
session_id,
66+
extra_headers={"Last-Event-ID": sent.data[0].id},
67+
timeout=context.remaining(),
68+
) as stream:
69+
for event in stream:
70+
context.remaining()
71+
if event.type == "agent.message":
72+
context.output(f"identity_{index + 1}_assistant", event.to_json())
73+
elif event.type in ("session.error", "session.status_terminated"):
74+
raise RuntimeError(f"Session stopped: {event.type}")
75+
elif event.type == "session.status_idle":
76+
break
77+
78+
79+
if __name__ == "__main__":
80+
run_cli("forward", Forward, {"identity_config": run})
Lines changed: 84 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,84 @@
1+
"""订阅 agent.message 文本增量,逐步更新消息预览,再用最终完整消息替换预览。
2+
3+
运行:python -m examples.forward.streaming_deltas
4+
"""
5+
6+
from __future__ import annotations
7+
8+
from examples.common.live import Run, choose_model, marker, name, run_cli
9+
from qca import Forward
10+
11+
from ._cleanup import finish_session
12+
13+
14+
def _delta_text(event: object) -> str:
15+
"""从一条文本增量事件中取出片段文本;不是文本增量时返回空串。"""
16+
delta = event.to_dict().get("delta") or {}
17+
if not isinstance(delta, dict) or delta.get("type") != "content_delta":
18+
return ""
19+
content = delta.get("content") or {}
20+
if not isinstance(content, dict) or content.get("type") != "text":
21+
return ""
22+
return content.get("text") or ""
23+
24+
25+
def run(client: Forward, context: Run) -> None:
26+
environment = client.environments.create(name=name("env"), config={"type": "cloud"})
27+
environment_id = context.track("environment", environment.id, lambda: client.environments.archive(environment.id))
28+
29+
identity = client.identities.create(external_id=name("identity"), name="SDK 示例用户")
30+
identity_id = context.track("identity", identity.id, lambda: client.identities.delete(identity.id))
31+
32+
model = choose_model(client.models.list(), context.config.model)
33+
context.output("selected_model", model)
34+
template = client.templates.create(
35+
name=name("template"),
36+
environment_id=environment_id,
37+
model=model,
38+
system="你是一个 SDK 示例助手。必要时调用工具,只使用可实际读取的数据回答问题。",
39+
tools=[{"type": "agent_toolset_20260401"}],
40+
)
41+
template_id = context.track("template", template.id, lambda: client.templates.archive(template.id))
42+
43+
session = client.sessions.create(identity_id=identity_id, template_id=template_id)
44+
session_id = context.track("session", session.id, lambda: finish_session(client, context, session.id))
45+
46+
marker_value = marker()
47+
prompt = "请分三句话解释为什么多轮对话要复用 Session ID,最后原样附上:" + marker_value
48+
# 预览增量不写入历史:先订阅、开启 agent.message 增量,再发送消息。
49+
with client.sessions.events.stream(
50+
session_id,
51+
event_deltas=["agent.message"],
52+
timeout=context.remaining(),
53+
) as stream:
54+
context.output("user", prompt)
55+
sent = client.sessions.events.send(
56+
session_id,
57+
events=[{"type": "user.message", "content": [{"type": "text", "text": prompt}]}],
58+
extra_headers={"Idempotency-Key": name("event")},
59+
)
60+
if not sent.data or not sent.data[0].id:
61+
raise RuntimeError("Send returned no user event")
62+
63+
previews: dict[str, str] = {}
64+
deltas = 0
65+
for event in stream:
66+
context.remaining()
67+
if event.type == "event_delta":
68+
text = _delta_text(event)
69+
if text and event.event_id:
70+
deltas += 1
71+
# 同一条消息的多个片段按 event_id 累加,逐步刷新该消息的预览。
72+
previews[event.event_id] = previews.get(event.event_id, "") + text
73+
context.output("preview", f"{event.event_id}: {previews[event.event_id]}")
74+
elif event.type == "agent.message":
75+
context.output("assistant", event.to_json())
76+
elif event.type in ("session.error", "session.status_terminated"):
77+
raise RuntimeError(f"Session stopped: {event.type}")
78+
elif event.type == "session.status_idle":
79+
break
80+
context.output("deltas", deltas)
81+
82+
83+
if __name__ == "__main__":
84+
run_cli("forward", Forward, {"streaming_deltas": run})

‎examples/managed/__main__.py‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,20 +3,26 @@
33
from examples.common.live import run_cli
44
from qca import Managed
55

6+
from .conversation import run as conversation
7+
from .custom_tools import run as custom_tools
68
from .deployment import run as deployment
79
from .dream import run as dream
810
from .memory import run as memory
911
from .models import run as models
1012
from .resources import run as resources
1113
from .session import run as session
14+
from .streaming_deltas import run as streaming_deltas
1215

1316
SCENARIOS = {
1417
"models": models,
1518
"session": session,
19+
"conversation": conversation,
1620
"resources": resources,
21+
"custom_tools": custom_tools,
1722
"memory": memory,
1823
"deployment": deployment,
1924
"dream": dream,
25+
"streaming_deltas": streaming_deltas,
2026
}
2127

2228

‎examples/managed/conversation.py‎

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,66 @@
1+
"""复用同一个 Session 进行多轮对话,并分页读取会话历史。
2+
3+
运行:python -m examples.managed.conversation
4+
"""
5+
6+
from __future__ import annotations
7+
8+
from examples.common.live import Run, choose_model, marker, name, run_cli
9+
from qca import Managed
10+
11+
from ._cleanup import finish_session
12+
13+
14+
def run(client: Managed, context: Run) -> None:
15+
environment = client.environments.create(name=name("env"), config={"type": "cloud"})
16+
environment_id = context.track("environment", environment.id, lambda: client.environments.archive(environment.id))
17+
18+
model = choose_model(client.models.list(), context.config.model)
19+
context.output("selected_model", model)
20+
agent = client.agents.create(
21+
name=name("agent"),
22+
model={"id": model},
23+
system="你是一个 SDK 示例助手。必要时调用工具,只使用可实际读取的数据回答问题。",
24+
tools=[{"type": "agent_toolset_20260401"}],
25+
)
26+
agent_id = context.track("agent", agent.id, lambda: client.agents.archive(agent.id))
27+
28+
session = client.sessions.create(environment_id=environment_id, agent=agent_id)
29+
session_id = context.track("session", session.id, lambda: finish_session(client, context, session.id))
30+
31+
def ask(prompt: str) -> None:
32+
context.output("user", prompt)
33+
sent = client.sessions.events.send(
34+
session_id,
35+
events=[{"type": "user.message", "content": [{"type": "text", "text": prompt}]}],
36+
extra_headers={"Idempotency-Key": name("event")},
37+
)
38+
if not sent.data or not sent.data[0].id:
39+
raise RuntimeError("Send returned no user event")
40+
with client.sessions.events.stream(
41+
session_id,
42+
extra_headers={"Last-Event-ID": sent.data[0].id},
43+
timeout=context.remaining(),
44+
) as stream:
45+
for event in stream:
46+
context.remaining()
47+
if event.type == "agent.message":
48+
context.output("assistant", event.to_json())
49+
elif event.type in ("session.error", "session.status_terminated"):
50+
raise RuntimeError(f"Session stopped: {event.type}")
51+
elif event.type == "session.status_idle":
52+
break
53+
54+
# 同一个 Session 支持多轮:服务端在 session_id 下保留完整历史,无需客户端携带上文。
55+
code = "project-" + marker()
56+
ask(f"这次项目代号是 {code}。请在本次对话中记住它,不要使用工具或写入记忆库。现在只回复:已记住。")
57+
ask("只根据本次会话上文,告诉我刚才约定的项目代号。只回复代号,不要使用工具。")
58+
59+
# 分页读取已有的用户消息和助手回复,重建对话文字记录。
60+
for event in client.sessions.events.list(session_id, order="asc", limit=100):
61+
if event.type in ("user.message", "agent.message"):
62+
context.output(f"history.{event.type}", event.to_json())
63+
64+
65+
if __name__ == "__main__":
66+
run_cli("managed", Managed, {"conversation": run})

0 commit comments

Comments
 (0)