Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
32 commits
Select commit Hold shift + click to select a range
c69f8ca
fix(api): protocol command 支持透传 stream_modes / stream_subgraphs
LiangMuYuan Aug 5, 2026
e0959c1
fix(streaming): 默认改用官方 astream 输出,修复 Send 并行 updates 重复 (#48)
LiangMuYuan Aug 5, 2026
85305c8
fix(streaming): updates 内多条无 id 消息分配递增 message_index,避免协议消息 id 冲突
LiangMuYuan Aug 5, 2026
53d44d3
test(run_executor): 适配官方 astream 执行路径,覆盖 updates 不重复与协议事件
LiangMuYuan Aug 5, 2026
a430afb
fix(streaming): 通过原生 tools stream mode 恢复 tool 事件,tool_name 按 tool_ca…
LiangMuYuan Aug 5, 2026
d1a9a2c
fix(api): 默认 /runs/{id}/stream 重放线程级 protocol 事件(含 tool/message),end 收尾
LiangMuYuan Aug 5, 2026
9d05e61
test(streaming): 集成测试适配 protocol 事件(tools 通道/重放/子图 namespace)
LiangMuYuan Aug 5, 2026
be61bdd
fix(api): 默认重放线程事件带 seq 并按 after_seq 过滤,支持 Last-Event-ID 续传
LiangMuYuan Aug 5, 2026
a05dfab
test(e2e): live provider 流式断言适配(从 values 快照取最终 ai 文本)
LiangMuYuan Aug 5, 2026
cd0261b
fix(streaming): 中断改写为 values.__interrupt__ 事件并去重(对齐官方)
LiangMuYuan Aug 5, 2026
d8c6a78
test(e2e): live provider 适配(store PUT 204 断言 + hitl 落点到 /stream/events)
LiangMuYuan Aug 6, 2026
df80fbf
test(e2e): live provider 实时流 messages/partial 增量累积验证
LiangMuYuan Aug 6, 2026
f1802d5
test(run_executor): 补 events 模式 / subgraphs 三元组 / 中断合并覆盖,覆盖率回到 90%
LiangMuYuan Aug 6, 2026
d42138c
fix(ci): verify_docker_api 适配重放端点协议事件格式(tool-started/end.get)
LiangMuYuan Aug 6, 2026
9b53faf
fix(api): default run stream uses single monotonic SSE cursor for Las…
LiangMuYuan Aug 8, 2026
244e3a0
fix(streaming): publish protocol events into run stream for single mo…
LiangMuYuan Aug 10, 2026
32e744a
test(live-provider): align default-stream assertions with protocol re…
LiangMuYuan Aug 10, 2026
552d3ad
fix(api): return 400 for invalid stream modes; raise langgraph min to…
LiangMuYuan Aug 10, 2026
d6e7851
fix(streaming): publish raw astream_events items onto the events channel
LiangMuYuan Aug 10, 2026
18a0998
fix(streaming): normalize tuple subgraph namespaces to list for live …
LiangMuYuan Aug 10, 2026
90d0ed1
test(streaming): add HTTP-level regression for events mode raw astrea…
LiangMuYuan Aug 11, 2026
e62c8ec
test(streaming): add HTTP-level mid-run reconnect regressions for inl…
LiangMuYuan Aug 11, 2026
4087fca
test(api): protocol run.start invalid stream_mode returns 400
LiangMuYuan Aug 11, 2026
0298d79
test(ci): assert run-stream mid-run reconnect exactly-once in docker …
LiangMuYuan Aug 11, 2026
c08fe10
test(protocol): guard live namespace filter against tuple subgraph na…
LiangMuYuan Aug 11, 2026
0cb21d6
fix(streaming): keep run-stream seq monotonic across resume and cold …
LiangMuYuan Aug 11, 2026
d9b0f38
test(live-provider): restore token-level messages/partial assertion i…
LiangMuYuan Aug 11, 2026
9bd73bc
fix(streaming): keep thread-level protocol seq monotonic across cold …
LiangMuYuan Aug 12, 2026
e8a11e5
fix(streaming): atomic durable append for run/thread stream seq; defa…
LiangMuYuan Aug 12, 2026
4311a86
fix(streaming): derive id-less message identity from subgraph namespa…
LiangMuYuan Aug 12, 2026
a389d04
test(live-provider): assert messages/metadata identity aligns with st…
LiangMuYuan Aug 12, 2026
582e7b1
docs(AGENTS): align live-provider proof target with official messages…
LiangMuYuan Aug 12, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,4 +22,4 @@
- `workflow_dispatch` inputs may override the model and base URL for a single run.
- API keys stay in GitHub secrets only; do not add dispatch inputs for secrets.
- The workflow runs `tests/integration/test_live_provider_streaming.py` against manifest-registered graphs in `examples/live_provider_graphs/manifest.json`.
- The proof target is SSE `message_chunk` events from real provider-backed graphs, not just a final successful response.
- The proof target is incremental SSE `messages/partial` frames (aligned to the official LangGraph v1 messages wire contract: `messages/metadata` + `messages/partial`/`messages/complete`) from real provider-backed graphs, not just a final successful response. The legacy `message_chunk` event no longer exists in the official SDK wire format.
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ dependencies = [
"aiomysql>=0.2.0",
"aiosqlite>=0.20.0",
"redis>=5.0.0",
"langgraph>=1.0.3",
"langgraph>=1.2.0",
"langgraph-sdk>=0.3.5",
"langchain-core>=1.0.0",
"langchain-openai>=1.0.0",
Expand Down
120 changes: 120 additions & 0 deletions scripts/check_atomic_append_concurrency.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
"""Cross-dialect concurrency check for the atomic stream append API.

Runs N concurrent appends against the same run and thread on the given
metadata backend and asserts every seq is unique and gapless, and that the
event table contains exactly N rows (no frame lost, no duplicate). This is the
sqlite-independent proof of the counter-row locking claim: MySQL/PostgreSQL
honour SELECT ... FOR UPDATE (sqlite does not), so the concurrent publishers
must serialize on the counter row rather than rely on uniqueness retries.

Usage (repo root, from WSL2):
UV_PROJECT_ENVIRONMENT=... uv run python scripts/check_atomic_append_concurrency.py mysql|postgresql|sqlite
"""

import asyncio
import sys
import time
from pathlib import Path

sys.path.insert(0, str(Path(__file__).resolve().parents[1]))

from sqlalchemy import func, select # noqa: E402

from agentseek_api.core.database import db_manager # noqa: E402
from agentseek_api.core.orm import RunStreamEvent, ThreadStreamEvent # noqa: E402
from agentseek_api.services import stream_persistence as stream_module # noqa: E402
from agentseek_api.settings import settings # noqa: E402

_N = 12


async def _count_events(scope: str, scope_id: str) -> int:
model = RunStreamEvent if scope == "run" else ThreadStreamEvent
id_column = model.run_id if scope == "run" else model.thread_id
async with db_manager.get_session_factory()() as session:
return int(await session.scalar(select(func.count()).select_from(model).where(id_column == scope_id)))


async def _exercise(backend: str) -> None:
# Unique scope ids per run so a pre-existing counter row (from an earlier
# invocation against the same database) cannot shift the expected range.
stamp = int(time.time() * 1000)
run_id, thread_id = f"run-conc-{backend}-{stamp}", f"thread-conc-{backend}-{stamp}"
for scope, scope_id in (("run", run_id), ("thread", thread_id)):
if scope == "run":
append = stream_module.append_run_stream_event_atomic
results = await asyncio.gather(*(append(scope_id, {"event": "message", "data": i}) for i in range(_N)))
rows = await stream_module.load_run_stream_events(scope_id)
seqs = sorted(seq for seq, _ in results)
persisted = [seq for seq, _ in rows]
else:
append = stream_module.append_thread_stream_event_atomic
payload = lambda i: {"method": "values", "params": {"namespace": [], "timestamp": 1, "data": i}} # noqa: E731
results = await asyncio.gather(*(append(scope_id, payload(i)) for i in range(_N)))
rows = await stream_module.load_thread_stream_events(scope_id, channels=["values"], namespaces=None, depth=None)
seqs = sorted(seq for seq, _ in results)
persisted = [event["seq"] for event in rows]

assert seqs == list(range(1, _N + 1)), f"[{backend}/{scope}] non-gapless seqs: {seqs}"
assert persisted == list(range(1, _N + 1)), f"[{backend}/{scope}] persisted mismatch: {persisted}"
assert len(results) == _N, f"[{backend}/{scope}] frame dropped: got {len(results)}"
count = await _count_events(scope, scope_id)
assert count == _N, f"[{backend}/{scope}] event rows != {_N}: {count}"
print(f"[{backend}/{scope}] OK: {_N} concurrent appends -> unique gapless 1..{_N}, {count} rows")

print(f"[{backend}] PASS")


def main() -> None:
backend = sys.argv[1] if len(sys.argv) > 1 else "sqlite"
if backend == "sqlite":
settings.SEEKDB_EMBED = False
settings.SEEKDB_URL = "sqlite+aiosqlite:////tmp/atomic-conc.db"
settings.METADATA_DB_BACKEND = "sqlite"
elif backend == "mysql":
settings.SEEKDB_EMBED = False
settings.SEEKDB_URL = "mysql+aiomysql://root:root@127.0.0.1:33306/agentseek"
settings.METADATA_DB_BACKEND = "mysql"
elif backend == "postgresql":
settings.SEEKDB_EMBED = False
settings.SEEKDB_URL = "postgresql://postgres:postgres@127.0.0.1:35432/agentseek"
settings.METADATA_DB_BACKEND = "postgresql"
else:
raise SystemExit(f"unknown backend: {backend}")
asyncio.run(_run(backend))


async def _run(backend: str) -> None:
# The concurrency check only exercises the metadata-DB stream events; the
# langgraph checkpointer/store backends (OceanBase/MySQL via pymysql)
# connect eagerly in their constructors and would fail under a postgres
# metadata URL, so they are replaced with inert stand-ins.
import agentseek_api.core.database as database_module

class _FakeCheckpointer:
def __init__(self, *args: object, **kwargs: object) -> None:
pass

def setup(self) -> None:
return None

def save_checkpoint(self, **kwargs: object) -> None:
return None

class _FakeStore:
def __init__(self, *args: object, **kwargs: object) -> None:
pass

database_module.OceanBaseCheckpointSaver = _FakeCheckpointer # type: ignore[assignment]
database_module.LangGraphOceanBaseCheckpointSaver = _FakeCheckpointer # type: ignore[assignment]
database_module.OceanBaseStore = _FakeStore # type: ignore[assignment]
await db_manager.initialize()
try:
await _exercise(backend)
finally:
await db_manager.close()
print(f"[{backend}] db closed")


if __name__ == "__main__":
main()
182 changes: 176 additions & 6 deletions scripts/verify_docker_api.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
from __future__ import annotations

import argparse
import http.client
import json
from urllib import error as urllib_error
from urllib import parse as urllib_parse
from urllib import request as urllib_request


Expand Down Expand Up @@ -51,6 +53,67 @@ def _stream_payloads(stream_text: str) -> list[dict[str, object]]:
]


def _sse_frames(stream_text: str) -> list[tuple[int | None, str, dict[str, object]]]:
"""Parse SSE text into (id, event, data) frames. ``id`` may be absent."""
frames: list[tuple[int | None, str, dict[str, object]]] = []
current_id: int | None = None
current_event = ""
current_data: list[str] = []
for line in stream_text.splitlines():
if line.startswith("id: "):
current_id = int(line[len("id: "):].strip())
elif line.startswith("event: "):
current_event = line[len("event: "):].strip()
elif line.startswith("data: "):
current_data.append(line[len("data: "):].strip())
elif line == "" and current_data:
frames.append((current_id, current_event, json.loads("".join(current_data))))
current_id, current_event, current_data = None, "", []
return frames


def _read_sse_prefix_then_disconnect(
*,
base_url: str,
path: str,
headers: dict[str, str],
max_frames: int,
timeout_seconds: float = 15.0,
) -> list[tuple[int | None, str, dict[str, object]]]:
"""Connect to an SSE endpoint, read up to ``max_frames`` frames, then close
the connection to simulate a client disconnecting mid-run.

``urllib`` (used by ``_request``) blocks until the whole body is read, so a
mid-stream disconnect has to be driven with a lower-level ``http.client``
connection that we can tear down after a few frames.
"""
parsed = urllib_parse.urlsplit(base_url)
conn = http.client.HTTPConnection(parsed.hostname, parsed.port, timeout=timeout_seconds)
frames: list[tuple[int | None, str, dict[str, object]]] = []
try:
conn.request("GET", path, headers=headers)
response = conn.getresponse()
current_id: int | None = None
current_event = ""
current_data: list[str] = []
for raw_line in response:
line = raw_line.decode("utf-8", errors="replace").rstrip("\r\n")
if line.startswith("id: "):
current_id = int(line[len("id: "):].strip())
elif line.startswith("event: "):
current_event = line[len("event: "):].strip()
elif line.startswith("data: "):
current_data.append(line[len("data: "):].strip())
elif line == "" and current_data:
frames.append((current_id, current_event, json.loads("".join(current_data))))
current_id, current_event, current_data = None, "", []
if len(frames) >= max_frames:
break
finally:
conn.close()
return frames



def _assert_sample_run(
*,
Expand Down Expand Up @@ -234,8 +297,13 @@ def _assert_common_flow(base_url: str) -> None:
assert isinstance(stream_body, str)
assert "text/event-stream" in stream_content_type
payloads = _stream_payloads(stream_body)
assert any(payload["event"] == "start" for payload in payloads)
assert any(payload["event"] == "end" and payload.get("status") == "success" for payload in payloads)
assert any(isinstance(payload, dict) and payload.get("event") == "start" for payload in payloads)
assert any(
isinstance(payload, dict)
and payload.get("event") == "end"
and payload.get("status") == "success"
for payload in payloads
)

_, stateless_run, _ = _request(
base_url=base_url,
Expand Down Expand Up @@ -326,7 +394,7 @@ def _assert_common_flow(base_url: str) -> None:
end_statuses = {
payload["status"]
for payload in _stream_payloads(resumed_stream)
if payload["event"] == "end"
if isinstance(payload, dict) and payload.get("event") == "end"
}
assert "success" in end_statuses
stress_waited, _ = _assert_sample_run(
Expand Down Expand Up @@ -358,7 +426,12 @@ def _assert_common_flow(base_url: str) -> None:
react_output = react_waited["output"]
assert isinstance(react_output, dict)
assert "42" in str(react_output["final_text"])
assert any(payload["event"] == "tool_start" and payload["name"] == "lookup" for payload in react_payloads)
assert any(
isinstance(payload, dict)
and payload.get("event") == "tool-started"
and payload.get("tool_name") == "lookup"
for payload in react_payloads
)

stress_tool_waited, stress_tool_payloads = _assert_sample_run(
base_url=base_url,
Expand All @@ -372,10 +445,107 @@ def _assert_common_flow(base_url: str) -> None:
tool_messages = [message for message in stress_tool_output["transcript"] if message["type"] == "ToolMessage"]
assert len(tool_messages) == 3
tool_starts = [
payload for payload in stress_tool_payloads if payload["event"] == "tool_start" and payload["name"] == "slow_process"
payload
for payload in stress_tool_payloads
if isinstance(payload, dict)
and payload.get("event") == "tool-started"
and payload.get("tool_name") == "slow_process"
]
assert len(tool_starts) == 3

_assert_run_stream_reconnect_exactly_once(base_url=base_url, user_headers=alice)


def _assert_run_stream_reconnect_exactly_once(*, base_url: str, user_headers: dict[str, str]) -> None:
"""Regression for the run-stream reconnect contract.

A client that connects to ``GET /runs/{id}/stream`` while the run is still
executing, disconnects, and reconnects with ``Last-Event-ID`` must not
replay frames already delivered and must not lose frames produced while it
was disconnected (exactly-once). This runs against the real HTTP boundary
and, in the redis-durable job, against the real Redis run-stream store —
the path the in-process pytest regression covers with a faked store.
"""
_, assistant, _ = _request(
base_url=base_url,
path="/assistants",
method="POST",
payload={"name": "docker-reconnect", "graph_id": "stress_tool_agent"},
)
assert isinstance(assistant, dict)
assistant_id = str(assistant["assistant_id"])

_, thread, _ = _request(
base_url=base_url,
path="/threads",
method="POST",
payload={"metadata": {"suite": "docker-reconnect"}},
headers=user_headers,
)
assert isinstance(thread, dict)
thread_id = str(thread["thread_id"])

# ~4.5s total runtime (3 steps * 1.5s) guarantees a mid-run window.
_, run, _ = _request(
base_url=base_url,
path=f"/threads/{thread_id}/runs",
method="POST",
payload={"assistant_id": assistant_id, "input": {"delay": 1.5, "steps": 3}},
headers=user_headers,
)
assert isinstance(run, dict)
run_id = str(run["run_id"])

stream_path = f"/threads/{thread_id}/runs/{run_id}/stream"

# Phase 1: connect mid-run, read a few frames, then disconnect.
phase1 = _read_sse_prefix_then_disconnect(
base_url=base_url,
path=stream_path,
headers=user_headers,
max_frames=3,
)
if not phase1:
phase1 = _read_sse_prefix_then_disconnect(
base_url=base_url,
path=stream_path,
headers=user_headers,
max_frames=3,
)
assert phase1, "expected to observe the run stream while the run is still active"
phase1_ids = [frame_id for frame_id, _, _ in phase1 if frame_id is not None]
last_id = int(phase1_ids[-1])

_, waited, _ = _request(
base_url=base_url,
path=f"/threads/{thread_id}/runs/{run_id}/wait",
headers=user_headers,
)
assert isinstance(waited, dict)
assert waited["status"] == "success"

# Phase 2: reconnect with Last-Event-ID and assert exactly-once.
_, phase2_body, _ = _request(
base_url=base_url,
path=stream_path,
headers={**user_headers, "Last-Event-ID": str(last_id)},
)
assert isinstance(phase2_body, str)
phase2_ids = [frame_id for frame_id, _, _ in _sse_frames(phase2_body) if frame_id is not None]

assert all(frame_id > last_id for frame_id in phase2_ids), (
f"reconnect replayed an already-delivered id: phase1={phase1_ids} phase2={phase2_ids}"
)
all_ids = phase1_ids + phase2_ids
assert len(all_ids) == len(set(all_ids)), f"duplicate ids delivered across reconnect: {all_ids}"

_, full_body, _ = _request(base_url=base_url, path=stream_path, headers=user_headers)
assert isinstance(full_body, str)
full_ids = [frame_id for frame_id, _, _ in _sse_frames(full_body) if frame_id is not None]
assert [frame_id for frame_id in full_ids if frame_id > last_id] == phase2_ids, (
f"reconnect content mismatch: full={full_ids} phase2={phase2_ids}"
)


def _assert_smoke_flow(base_url: str) -> None:
headers = {"x-user-id": "autobuild"}
Expand Down Expand Up @@ -501,7 +671,7 @@ def _assert_resume_check(base_url: str, *, thread_id: str, run_id: str, resume:
end_statuses = [
payload["status"]
for payload in _stream_payloads(run_stream)
if payload["event"] == "end"
if isinstance(payload, dict) and payload.get("event") == "end"
]
assert "interrupted" in end_statuses
assert "success" in end_statuses
Expand Down
Loading
Loading