Start LangGraph backend migration foundation

This commit is contained in:
Xin Wang
2026-07-27 17:21:29 +08:00
parent 5c719ed2ea
commit 1c8e9da486
33 changed files with 1859 additions and 86 deletions

View File

@@ -0,0 +1,121 @@
from copy import deepcopy
import pytest
from src.agent.state import AccidentGraphState, GeneratedTurn
from src.backends.chat import ChatInput, FormUpdate, TextDelta
from src.backends.fastgpt import FastGPTBackend
from src.backends.langgraph import LangGraphBackend
from src.core.config import Settings
from src.core.fastgpt_client import create_chat_backend
class FakeResponseGenerator:
def __init__(self):
self.states = []
self.closed = False
async def generate(self, state: AccidentGraphState) -> GeneratedTurn:
self.states.append(deepcopy(state))
turn_number = state.get("turn_count", 0) + 1
patch = {"turn": turn_number} if state["need_form_update"] else {}
return GeneratedTurn(
content=f"<state>1002</state>第{turn_number}",
form_update=patch,
)
async def aclose(self) -> None:
self.closed = True
def langgraph_settings():
return Settings(
_env_file=None,
environment="test",
agent_backend="langgraph",
langgraph_checkpointer="memory",
llm_api_key="test-key",
llm_model="test-model",
)
@pytest.mark.asyncio
async def test_langgraph_backend_preserves_thread_scoped_turn_state():
generator = FakeResponseGenerator()
backend = create_chat_backend(
langgraph_settings(),
response_generator=generator,
)
assert isinstance(backend, LangGraphBackend)
first = await backend.complete(
ChatInput("session-1", "第一轮", need_form_update=True)
)
second = await backend.complete(
ChatInput("session-1", "第二轮", need_form_update=True)
)
other_session = await backend.complete(
ChatInput("session-2", "独立会话", need_form_update=True)
)
assert first.content == "<state>1002</state>第1轮"
assert first.form_update == {"turn": 1}
assert second.content == "<state>1002</state>第2轮"
assert second.form_update == {"turn": 2}
assert other_session.content == "<state>1002</state>第1轮"
assert generator.states[0].get("turn_count", 0) == 0
assert generator.states[1]["turn_count"] == 1
assert generator.states[2].get("turn_count", 0) == 0
@pytest.mark.asyncio
async def test_langgraph_stream_bridges_form_update_before_text():
generator = FakeResponseGenerator()
backend = create_chat_backend(
langgraph_settings(),
response_generator=generator,
)
events = [
event
async for event in backend.stream(
ChatInput("session-stream", "开始", need_form_update=True)
)
]
assert events == [
FormUpdate({"turn": 1}),
TextDelta("<state>1002</state>第1轮"),
]
@pytest.mark.asyncio
async def test_langgraph_backend_closes_owned_generator():
generator = FakeResponseGenerator()
backend = create_chat_backend(
langgraph_settings(),
response_generator=generator,
)
await backend.aclose()
assert generator.closed
def test_factory_keeps_fastgpt_as_default_compatible_backend():
settings = Settings(
_env_file=None,
environment="test",
agent_backend="fastgpt",
fastgpt_api_key="test-key",
fastgpt_base_url="http://fastgpt.test",
fastgpt_app_id="test-app",
)
fake_client = object()
backend = create_chat_backend(
settings,
fastgpt_client=fake_client,
)
assert isinstance(backend, FastGPTBackend)

View File

@@ -1,8 +1,11 @@
GET http://101.89.151.141:3000
@fastgptBaseUrl = http://127.0.0.1:3000
@fastgptApiKey = replace-with-local-api-key
GET {{fastgptBaseUrl}}
###
POST http://101.89.151.141:3000/api/v1/chat/completions
POST {{fastgptBaseUrl}}/api/v1/chat/completions
content-type: application/json
Authorization: Bearer fastgpt-xCH4CaEoNEyVtq7fkBEI5UP3O6sABKdpGszTtSYk4R2TVW5VgrPp1YPfuLX1iH
Authorization: Bearer {{fastgptApiKey}}
{
@@ -23,4 +26,4 @@ content-type: application/json; charset=utf-8
etag: "s14v22uu1g5f"
content-length: 219
date: Fri, 20 Jun 2025 02:37:16 GMT
connection: close
connection: close

View File

@@ -0,0 +1,116 @@
import json
import pytest
from src.api.endpoints import chat
from src.backends.chat import ChatInput, ChatResult, FormUpdate, TextDelta
from src.schemas.models import ProcessRequest_chat
def make_request(**overrides):
payload = {
"sessionId": "session-001",
"timeStamp": "20260726120000",
"text": "发生了交通事故",
"needFormUpdate": True,
}
payload.update(overrides)
return ProcessRequest_chat(**payload)
async def response_text(response):
chunks = []
async for chunk in response.body_iterator:
chunks.append(chunk.decode() if isinstance(chunk, bytes) else chunk)
return "".join(chunks)
def parse_sse(body):
events = []
for block in body.strip().split("\n\n"):
lines = block.splitlines()
event = lines[0].removeprefix("event: ")
data = json.loads(lines[1].removeprefix("data: "))
events.append((event, data))
return events
class OrderedBackend:
async def stream(self, chat_input: ChatInput):
yield TextDelta("<sta")
yield TextDelta("te>1002</state>")
yield FormUpdate({"jdcsl": 2})
yield TextDelta("第一句。")
yield TextDelta("第二句。")
async def complete(self, chat_input: ChatInput):
return ChatResult("<state>1002</state>第一句。第二句。", "1002", {"jdcsl": 2})
class MissingPrefixBackend:
async def stream(self, chat_input: ChatInput):
yield TextDelta("没有状态前缀")
async def complete(self, chat_input: ChatInput):
return ChatResult("没有状态前缀")
@pytest.mark.asyncio
async def test_sse_success_event_order_and_cardinality():
response = await chat(make_request(), stream=True, backend=OrderedBackend())
events = parse_sse(await response_text(response))
names = [name for name, _ in events]
assert names == [
"stage_code",
"formUpdate",
"text_delta",
"text_delta",
"done",
]
assert names.count("stage_code") == 1
assert names.count("formUpdate") == 1
assert names.count("done") == 1
assert "error" not in names
assert "".join(data["text"] for name, data in events if name == "text_delta") == (
"第一句。第二句。"
)
@pytest.mark.asyncio
async def test_use_text_chunk_only_changes_delta_boundaries():
request = make_request(useTextChunk=True)
response = await chat(request, stream=True, backend=OrderedBackend())
events = parse_sse(await response_text(response))
assert "".join(data["text"] for name, data in events if name == "text_delta") == (
"第一句。第二句。"
)
assert [name for name, _ in events].count("done") == 1
@pytest.mark.asyncio
async def test_missing_stream_prefix_characterizes_current_legacy_behavior():
response = await chat(
make_request(needFormUpdate=False),
stream=True,
backend=MissingPrefixBackend(),
)
events = parse_sse(await response_text(response))
assert [name for name, _ in events] == ["text_delta", "done"]
assert events[0][1]["text"] == "没有状态前缀"
@pytest.mark.asyncio
async def test_missing_non_stream_prefix_returns_compatible_business_error():
response = await chat(
make_request(needFormUpdate=False),
stream=False,
backend=MissingPrefixBackend(),
)
assert response.code == "500"
assert response.outputText == ""
assert response.nextStageCode == ""
assert response.msg == "大模型服务返回消息不完整"

View File

@@ -0,0 +1,145 @@
import json
import pytest
from src.api import endpoints
from src.api.endpoints import get_info, set_info
from src.schemas.models import ProcessRequest_get, ProcessRequest_set
class FakeResponse:
def __init__(self, payload):
self._payload = payload
def raise_for_status(self):
return None
def json(self):
return self._payload
class FakeInfoClient:
def __init__(self, state):
self.state = state
self.completion_calls = []
async def create_chat_completion(self, **kwargs):
self.completion_calls.append(kwargs)
if "variables" in kwargs:
self.state = kwargs["variables"]["state"]
return FakeResponse({"newVariables": {"state": self.state}})
@pytest.fixture
def skip_helper_record_deletion(monkeypatch):
calls = []
async def fake_delete(client, session_id):
calls.append(session_id)
monkeypatch.setattr(endpoints, "delete_last_two_chat_records", fake_delete)
return calls
def make_set_request(**overrides):
payload = {
"sessionId": "session-001",
"timeStamp": "20260726120000",
"key": "hphm1",
"value": "<PLATE_1>",
}
payload.update(overrides)
return ProcessRequest_set(**payload)
def make_get_request(**overrides):
payload = {
"sessionId": "session-001",
"timeStamp": "20260726120000",
"key": "all",
}
payload.update(overrides)
return ProcessRequest_get(**payload)
@pytest.mark.asyncio
async def test_set_info_reads_then_writes_fastgpt_state(
skip_helper_record_deletion,
):
client = FakeInfoClient({"hphm1": "<PLATE_OLD>", "jdcsl": 1})
response = await set_info(make_set_request(), client=client)
assert response.code == "200"
assert client.state == {"hphm1": "<PLATE_1>", "jdcsl": 1}
assert len(client.completion_calls) == 2
assert client.completion_calls[1]["variables"]["state"] == client.state
assert skip_helper_record_deletion == ["session-001", "session-001"]
@pytest.mark.asyncio
async def test_set_info_include_input_info_keeps_legacy_magic_payload(
skip_helper_record_deletion,
):
client = FakeInfoClient({})
await set_info(make_set_request(includeInputInfo=True), client=client)
message = client.completion_calls[0]["messages"][0]["content"]
assert message == '<setInfo>{"key": "hphm1", "value": "<PLATE_1>"}</setInfo>'
@pytest.mark.asyncio
async def test_get_info_all_keeps_json_string_and_boolean_encoding(
skip_helper_record_deletion,
):
client = FakeInfoClient(
{
"ywrysw": False,
"ywfjdc": True,
"jdcsl": 2,
"xm1": "<PERSON_1>",
"hphm1": "<PLATE_1>",
"xm2": "<PERSON_2>",
}
)
response = await get_info(make_get_request(), client=client)
value = json.loads(response.value)
assert response.code == "200"
assert isinstance(response.value, str)
assert value["acdinfo"]["ywrysw"] == "0"
assert value["acdinfo"]["ywfjdc"] == "1"
assert value["acdinfo"]["jdcsl"] == 2
assert value["acdhuman1"]["xm1"] == "<PERSON_1>"
assert value["acdhuman2"]["xm2"] == "<PERSON_2>"
assert skip_helper_record_deletion == ["session-001"]
@pytest.mark.asyncio
async def test_get_info_unknown_key_returns_json_encoded_empty_string(
skip_helper_record_deletion,
):
client = FakeInfoClient({})
response = await get_info(
make_get_request(key="unknown_legacy_key"),
client=client,
)
assert response.code == "200"
assert response.value == '""'
@pytest.mark.asyncio
async def test_get_info_fastgpt_shape_error_keeps_legacy_business_error():
class InvalidClient:
async def create_chat_completion(self, **kwargs):
return FakeResponse({})
response = await get_info(make_get_request(), client=InvalidClient())
assert response.code == "500"
assert response.value == ""
assert response.msg == "大模型服务器无响应"

78
test/core/test_config.py Normal file
View File

@@ -0,0 +1,78 @@
import pytest
from pydantic import ValidationError
from src.core.config import Settings
def test_fastgpt_backend_accepts_legacy_environment_names():
settings = Settings(
_env_file=None,
ZNJJ_ENVIRONMENT="test",
AGENT_BACKEND="fastgpt",
ANALYSIS_AUTH_TOKEN="test-fastgpt-key",
ANALYSIS_SERVICE_URL="http://fastgpt.test",
APP_ID="test-app",
)
assert settings.environment == "test"
assert settings.agent_backend == "fastgpt"
assert settings.has_fastgpt_config
assert settings.fastgpt_api_key.get_secret_value() == "test-fastgpt-key"
def test_fastgpt_backend_rejects_incomplete_configuration():
with pytest.raises(ValidationError, match="FastGPT backend requires"):
Settings(
_env_file=None,
environment="test",
agent_backend="fastgpt",
fastgpt_api_key="test-key",
)
def test_langgraph_backend_requires_only_minimal_llm_configuration():
settings = Settings(
_env_file=None,
environment="test",
agent_backend="langgraph",
langgraph_checkpointer="memory",
llm_api_key="test-llm-key",
llm_model="test-model",
)
assert settings.agent_backend == "langgraph"
assert settings.langgraph_checkpointer == "memory"
assert not settings.has_fastgpt_config
def test_langgraph_backend_rejects_missing_llm_configuration():
with pytest.raises(ValidationError, match="LLM_API_KEY, LLM_MODEL"):
Settings(
_env_file=None,
environment="test",
agent_backend="langgraph",
)
def test_production_cannot_use_memory_checkpointer():
with pytest.raises(ValidationError, match="cannot use the memory checkpointer"):
Settings(
_env_file=None,
environment="production",
agent_backend="langgraph",
langgraph_checkpointer="memory",
llm_api_key="test-llm-key",
llm_model="test-model",
)
def test_secret_values_are_masked_in_settings_repr():
settings = Settings(
_env_file=None,
environment="test",
agent_backend="langgraph",
llm_api_key="never-print-this-key",
llm_model="test-model",
)
assert "never-print-this-key" not in repr(settings)

View File

@@ -0,0 +1,162 @@
{
"version": "2026-07-26",
"description": "LangGraph 迁移的脱敏黄金场景。占位符不是可用个人信息。",
"scenarios": [
{
"case_id": "single_vehicle_happy_path",
"category": "end_to_end",
"initial_stage": "1001",
"turns": [
{"event": "session_started", "input": "【继续办理】", "expected_stage": "1002"},
{"event": "user_message", "input": "车辆倒车时碰到了固定物体", "expected_stage": "1002", "expected_patch": {"sgyy": "倒车碰到固定物体"}},
{"event": "user_message", "input": "没有人员受伤", "expected_stage": "1002", "expected_patch": {"ywrysw": false}},
{"event": "user_message", "input": "没有非机动车", "expected_stage": "1002", "expected_patch": {"ywfjdc": false}},
{"event": "user_message", "input": "事故时间是十分钟前", "expected_stage": "1002"},
{"event": "user_message", "input": "我还在现场", "expected_stage": "1002", "expected_patch": {"sfsgxc": true}},
{"event": "user_message", "input": "只有一辆机动车", "expected_stage": "2000", "expected_patch": {"jdcsl": 1}},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2001"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2002"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2003"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2004"},
{"event": "user_message", "input": "车牌正确", "expected_stage": "2005"},
{"event": "user_message", "input": "车损在车辆前方", "expected_stage": "3001", "expected_patch": {"csbw1": "前方"}}
],
"expected_handoff_reason": null
},
{
"case_id": "double_vehicle_happy_path",
"category": "end_to_end",
"initial_stage": "1002",
"turns": [
{"event": "user_message", "input": "两辆机动车发生追尾,没有人受伤,也没有非机动车", "expected_stage": "1002", "expected_patch": {"jdcsl": 2, "ywrysw": false, "ywfjdc": false, "sgyy": "追尾"}},
{"event": "user_message", "input": "时间正确,我还在现场", "expected_stage": "2010", "expected_patch": {"sfsgxc": true}},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2011"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2012"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2013"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2014"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2015"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2016"},
{"event": "user_message", "input": "车牌正确", "expected_stage": "3002"}
],
"expected_handoff_reason": null
},
{
"case_id": "explicit_handoff_global",
"category": "handoff",
"initial_stage": "2002",
"turns": [
{"event": "user_message", "input": "请帮我转人工", "expected_stage": "0001"}
],
"expected_handoff_reason": "user_requested"
},
{
"case_id": "injury_handoff",
"category": "safety",
"initial_stage": "1002",
"turns": [
{"event": "user_message", "input": "有人倒地并且不舒服", "expected_stage": "0003", "expected_patch": {"ywrysw": true}}
],
"expected_handoff_reason": "injury_or_complex"
},
{
"case_id": "injury_negation_does_not_handoff",
"category": "safety",
"initial_stage": "1002",
"turns": [
{"event": "user_message", "input": "不是人受伤,是车受损,人没事", "expected_stage": "1002", "expected_patch": {"ywrysw": false}}
],
"expected_handoff_reason": null
},
{
"case_id": "three_vehicle_complex_handoff",
"category": "safety",
"initial_stage": "1002",
"turns": [
{"event": "user_message", "input": "一共涉及三辆机动车", "expected_stage": "0003", "expected_patch": {"jdcsl": 3}}
],
"expected_handoff_reason": "complex_accident"
},
{
"case_id": "photo_failure_single",
"category": "deterministic_event",
"initial_stage": "2001",
"turns": [
{"event": "photo_recognition_failed", "input": "【客户端连续3次拍摄识别失败图片过于模糊】", "expected_stage": "0005"}
],
"expected_handoff_reason": "photo_recognition_failed"
},
{
"case_id": "photo_failure_double_confirmation",
"category": "deterministic_event",
"initial_stage": "2016",
"turns": [
{"event": "photo_recognition_failed", "input": "【客户端连续3次拍摄识别失败未识别到完整车牌】", "expected_stage": "0005"}
],
"expected_handoff_reason": "photo_recognition_failed"
},
{
"case_id": "two_consecutive_no_responses",
"category": "deterministic_event",
"initial_stage": "1002",
"turns": [
{"event": "no_response", "input": "【用户无回复】", "expected_stage": "1002", "expected_no_response_count": 1},
{"event": "no_response", "input": "【用户无回复】", "expected_stage": "0004", "expected_no_response_count": 2}
],
"expected_handoff_reason": "no_response"
},
{
"case_id": "photo_step_cannot_skip",
"category": "transition_guard",
"initial_stage": "2011",
"turns": [
{"event": "user_message", "input": "后面的照片我都拍好了", "expected_stage": "2011"},
{"event": "photo_completed", "input": "【拍摄完成】", "expected_stage": "2012"}
],
"expected_handoff_reason": null
},
{
"case_id": "external_update_then_chat",
"category": "state_consistency",
"initial_stage": "1002",
"initial_form": {"hphm1": "<PLATE_1>"},
"turns": [
{"event": "set_info", "input": {"key": "hphm1", "value": "<PLATE_CORRECTED>"}, "expected_stage": "1002", "expected_patch": {"hphm1": "<PLATE_CORRECTED>"}},
{"event": "user_message", "input": "请继续办理", "expected_stage": "1002", "expected_state_contains": {"hphm1": "<PLATE_CORRECTED>"}}
],
"expected_handoff_reason": null
},
{
"case_id": "invalid_external_internal_field",
"category": "state_consistency",
"initial_stage": "1002",
"turns": [
{"event": "set_info", "input": {"key": "stage_code", "value": "0000"}, "expected_error": "INVALID_FIELD", "expected_stage": "1002"}
],
"expected_handoff_reason": null
},
{
"case_id": "prefix_split_across_chunks",
"category": "prefix_parser",
"initial_stage": "1002",
"stream_chunks": ["<sta", "te>1002", "</sta", "te>请描述事故经过"],
"expected_stage": "1002",
"expected_text": "请描述事故经过"
},
{
"case_id": "prefix_unknown_code",
"category": "prefix_parser",
"initial_stage": "1002",
"model_output": "<state>9999</state>继续处理",
"expected_error": "UNKNOWN_STAGE_CODE",
"expected_stage": "1002"
},
{
"case_id": "prefix_illegal_photo_jump",
"category": "prefix_parser",
"initial_stage": "2000",
"model_output": "<state>2004</state>请确认车牌",
"expected_error": "ILLEGAL_STAGE_TRANSITION",
"expected_stage": "2000"
}
]
}

View File

@@ -0,0 +1,119 @@
import json
import re
from pathlib import Path
ROOT = Path(__file__).resolve().parents[1]
DOMAIN = ROOT / "docs" / "domain"
GOLDEN = ROOT / "test" / "fixtures" / "golden" / "accident-scenarios.json"
def load_json(path):
return json.loads(path.read_text(encoding="utf-8"))
def test_stage_registry_and_transition_matrix_are_closed():
stage_registry = load_json(DOMAIN / "stage-codes.json")
transitions = load_json(DOMAIN / "stage-transitions.json")
codes = {entry["code"] for entry in stage_registry["codes"]}
assert len(codes) == len(stage_registry["codes"])
assert {"0000", "0004", "0005", "1001", "1002", "3001", "3002"} <= codes
assert set(transitions["allowed"]) == codes
for source, targets in transitions["allowed"].items():
assert set(targets) <= codes, source
for terminal in {"0000", "0001", "0002", "0003", "0004", "0005"}:
assert transitions["allowed"][terminal] == []
def test_photo_sequences_cannot_skip_steps():
transitions = load_json(DOMAIN / "stage-transitions.json")
for sequence in transitions["photo_sequences"].values():
for current, following in zip(sequence, sequence[1:]):
assert following in transitions["allowed"][current]
later_steps = set(sequence[sequence.index(following) + 1 :])
assert later_steps.isdisjoint(transitions["allowed"][current])
def test_field_registry_groups_are_complete_and_unique():
registry = load_json(DOMAIN / "field-registry.json")
fields = registry["fields"]
keys = [field["key"] for field in fields]
assert len(keys) == len(set(keys))
assert set(keys) == {
key for group_keys in registry["groups"].values() for key in group_keys
}
assert {"sfzmhm1", "sjhm1", "hphm1", "sfzmhm2", "sjhm2", "hphm2"} <= {
field["key"] for field in fields if field["sensitive"]
}
assert {"sfzmwh1", "sjwh1", "sfzmwh2", "sjwh2"} <= {
field["key"] for field in fields if not field["external_write"]
}
def test_golden_scenarios_are_unique_and_use_known_stages():
stages = {
entry["code"]
for entry in load_json(DOMAIN / "stage-codes.json")["codes"]
}
scenarios = load_json(GOLDEN)["scenarios"]
case_ids = [case["case_id"] for case in scenarios]
assert len(scenarios) >= 12
assert len(case_ids) == len(set(case_ids))
for case in scenarios:
assert case["initial_stage"] in stages
if "expected_stage" in case:
assert case["expected_stage"] in stages
for turn in case.get("turns", []):
if "expected_stage" in turn:
assert turn["expected_stage"] in stages
def test_golden_turn_sequences_follow_the_transition_matrix():
allowed = load_json(DOMAIN / "stage-transitions.json")["allowed"]
scenarios = load_json(GOLDEN)["scenarios"]
for case in scenarios:
current = case["initial_stage"]
for turn in case.get("turns", []):
expected = turn.get("expected_stage")
if expected is None:
continue
if "expected_error" not in turn:
assert expected in allowed[current], (
case["case_id"],
current,
expected,
)
current = expected
def test_golden_fixture_contains_no_realistic_phone_or_national_id():
raw = GOLDEN.read_text(encoding="utf-8")
assert not re.search(r"(?<!\d)1[3-9]\d{9}(?!\d)", raw)
assert not re.search(r"(?<!\d)\d{17}[\dXx](?!\d)", raw)
def test_repository_sources_and_examples_do_not_embed_fastgpt_tokens():
token_pattern = re.compile(r"fastgpt-[A-Za-z0-9]{20,}")
candidates = [ROOT / ".env.example"]
for directory in ("src", "test", "docs"):
candidates.extend(
path
for path in (ROOT / directory).rglob("*")
if path.is_file() and "__pycache__" not in path.parts
)
for path in candidates:
try:
text = path.read_text(encoding="utf-8")
except UnicodeDecodeError:
continue
assert not token_pattern.search(text), path