Compare commits

...
Author SHA1 Message Date
chengyongru 42ebe3a6b0 fix(whatsapp): inspect remote media before dispatch 2026-08-05 15:37:34 +08:00
chengyongru 47029757e5 fix(whatsapp): normalize detected audio MIME aliases 2026-08-01 12:31:53 +08:00
chengyongru 10ce839e3d test(whatsapp): reuse media send client mock 2026-08-01 12:14:17 +08:00
chengyongru 353dfed502 fix(whatsapp): route media by detected MIME type 2026-08-01 12:09:54 +08:00
chengyongruandGitHub cdb75f8e7d feat(providers): support DeepSeek Responses API (#5197) 2026-08-01 11:53:51 +08:00
chengyongruandGitHub 971b977a84 fix(weixin): recover refreshed state after session expiry (#5196) 2026-08-01 00:28:21 +08:00
54650332fb fix(slack): scope channel thread openers to their own session
A top-level channel message that opens a thread fell back to the
channel-wide session, because the session key required `raw_thread_ts` —
which Slack only sets on messages that already arrived inside a thread.
Every new thread therefore began life in one shared channel session and
only became thread-scoped from its first reply onward, so unrelated
threads saw each other's opening turns.

Key off `thread_ts` instead. It is set both for messages arriving inside
a thread and for channel messages that `reply_in_thread` opens a thread
for. DM roots never get a `thread_ts`, so they keep the default per-chat
session and the DM routing from 82c5083 is preserved; with
`reply_in_thread` disabled no thread exists and the channel session is
still used.

This restores the per-thread isolation introduced in #1048.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-31 23:47:26 +08:00
chengyongruandGitHub 172fe4f991 fix(webui): preserve user scroll ownership near tail (#5193) 2026-07-31 23:37:26 +08:00
shixi-liandchengyongru dda9b61b1e fix(config): install timezone data on all platforms 2026-07-31 19:55:22 +08:00
22 changed files with 952 additions and 71 deletions
+1 -2
View File
@@ -356,8 +356,7 @@ Providers that use the Responses API can keep reasoning context across a
conversation, which helps with multi-step tasks. Supported providers can also conversation, which helps with multi-step tasks. Supported providers can also
compact long conversations automatically. compact long conversations automatically.
nanobot preserves Responses conversation state automatically for OpenAI nanobot preserves Responses conversation state automatically for OpenAI Responses, OpenAI Codex, Azure OpenAI, DeepSeek V4 Flash, and compatible GitHub Copilot models.
Responses, OpenAI Codex, Azure OpenAI, and compatible GitHub Copilot models.
Native compaction is also automatic when the provider supports it. The Native compaction is also automatic when the provider supports it. The
threshold is derived from the active model's context window and reserved output threshold is derived from the active model's context window and reserved output
headroom; no provider configuration is required. headroom; no provider configuration is required.
+2
View File
@@ -231,6 +231,8 @@ Arbitrary custom provider names are OpenAI-compatible only; they do not use the
`providers.openai.apiType` may be set when you need to force a specific OpenAI API surface. Other providers reject `apiType`; leave it unset outside `providers.openai`. Replace the model with a model ID available to your OpenAI account. Direct OpenAI Responses, OpenAI Codex, Azure OpenAI Responses, and eligible GitHub Copilot models share [opaque Responses state retention](./configuration.md#responses-state-and-compaction); native compaction is enabled only where the backend supports it. `providers.openai.apiType` may be set when you need to force a specific OpenAI API surface. Other providers reject `apiType`; leave it unset outside `providers.openai`. Replace the model with a model ID available to your OpenAI account. Direct OpenAI Responses, OpenAI Codex, Azure OpenAI Responses, and eligible GitHub Copilot models share [opaque Responses state retention](./configuration.md#responses-state-and-compaction); native compaction is enabled only where the backend supports it.
DeepSeek is the model-level exception in the OpenAI-compatible provider: `deepseek-v4-flash` automatically uses DeepSeek's native Responses API, while `deepseek-v4-pro` remains on Chat Completions.
### Custom OpenAI-Compatible Endpoint ### Custom OpenAI-Compatible Endpoint
The `custom` provider fits one OpenAI-compatible endpoint that is not represented by a named provider. The `custom` provider fits one OpenAI-compatible endpoint that is not represented by a named provider.
+5 -6
View File
@@ -493,12 +493,11 @@ class SlackChannel(BaseChannel):
except Exception as e: except Exception as e:
self.logger.debug("reactions_add failed: {}", e) self.logger.debug("reactions_add failed: {}", e)
# Thread-scoped session key whenever the user is in a real thread # Thread-scoped session key whenever the turn lives in a thread: either the
# (raw_thread_ts is set). DM threads get their own session, separate # message arrived inside one (raw_thread_ts) or reply_in_thread opens a new
# from the DM root, so context doesn't bleed across thread boundaries. # thread for this channel message. DM roots have no thread_ts and keep the
session_key = ( # default per-chat session, so context doesn't bleed across thread boundaries.
f"slack:{chat_id}:{thread_ts}" if thread_ts and raw_thread_ts else None session_key = f"slack:{chat_id}:{thread_ts}" if thread_ts else None
)
media_paths: list[str] = [] media_paths: list[str] = []
file_markers: list[str] = [] file_markers: list[str] = []
for file_info in _as_json_list(event.get("files")) or []: for file_info in _as_json_list(event.get("files")) or []:
@@ -555,6 +555,113 @@ async def test_dm_thread_message_keeps_thread_ts_and_threaded_session() -> None:
assert kwargs["metadata"]["slack"]["thread_ts"] == "1700000000.000100" assert kwargs["metadata"]["slack"]["thread_ts"] == "1700000000.000100"
def _channel_mention_request(envelope_id: str, ts: str) -> SimpleNamespace:
return SimpleNamespace(
type="events_api",
envelope_id=envelope_id,
payload={
"event": {
"type": "app_mention",
"user": "U1",
"channel": "C123",
"text": "<@UBOT> hello",
"ts": ts,
}
},
)
@pytest.mark.asyncio
async def test_channel_root_message_uses_thread_scoped_session() -> None:
"""A channel mention that opens a thread belongs to that thread's session."""
channel = SlackChannel(SlackConfig(enabled=True), MessageBus())
channel._bot_user_id = "UBOT"
channel._web_client = _FakeAsyncWebClient()
channel._handle_message = AsyncMock() # type: ignore[method-assign]
client = SimpleNamespace(send_socket_mode_response=AsyncMock())
req = _channel_mention_request("env-c1", "1700000000.000100")
await channel._on_socket_request(client, req)
channel._handle_message.assert_awaited_once()
kwargs = channel._handle_message.await_args.kwargs
assert kwargs["session_key"] == "slack:C123:1700000000.000100"
assert kwargs["metadata"]["slack"]["thread_ts"] == "1700000000.000100"
@pytest.mark.asyncio
async def test_channel_root_messages_do_not_share_one_session() -> None:
"""Two threads opened in the same channel must not collapse into one session."""
channel = SlackChannel(SlackConfig(enabled=True), MessageBus())
channel._bot_user_id = "UBOT"
channel._web_client = _FakeAsyncWebClient()
channel._handle_message = AsyncMock() # type: ignore[method-assign]
client = SimpleNamespace(send_socket_mode_response=AsyncMock())
first = _channel_mention_request("env-c1", "1700000000.000100")
second = _channel_mention_request("env-c2", "1700000000.000200")
await channel._on_socket_request(client, first)
await channel._on_socket_request(client, second)
session_keys = [call.kwargs["session_key"] for call in channel._handle_message.await_args_list]
assert session_keys == [
"slack:C123:1700000000.000100",
"slack:C123:1700000000.000200",
]
@pytest.mark.asyncio
async def test_channel_root_message_without_reply_in_thread_uses_channel_session() -> None:
"""With reply_in_thread disabled no thread is opened, so the channel session is used."""
channel = SlackChannel(SlackConfig(enabled=True, reply_in_thread=False), MessageBus())
channel._bot_user_id = "UBOT"
channel._web_client = _FakeAsyncWebClient()
channel._handle_message = AsyncMock() # type: ignore[method-assign]
client = SimpleNamespace(send_socket_mode_response=AsyncMock())
req = _channel_mention_request("env-c3", "1700000000.000300")
await channel._on_socket_request(client, req)
channel._handle_message.assert_awaited_once()
kwargs = channel._handle_message.await_args.kwargs
assert kwargs["session_key"] is None
assert kwargs["metadata"]["slack"]["thread_ts"] is None
@pytest.mark.asyncio
async def test_channel_thread_reply_keeps_thread_session() -> None:
"""A reply inside a channel thread stays in the session opened by the root message."""
channel = SlackChannel(SlackConfig(enabled=True), MessageBus())
channel._bot_user_id = "UBOT"
channel._web_client = _FakeAsyncWebClient()
channel._handle_message = AsyncMock() # type: ignore[method-assign]
channel._with_thread_context = AsyncMock(return_value="hello") # type: ignore[method-assign]
client = SimpleNamespace(send_socket_mode_response=AsyncMock())
req = SimpleNamespace(
type="events_api",
envelope_id="env-c4",
payload={
"event": {
"type": "app_mention",
"user": "U1",
"channel": "C123",
"text": "<@UBOT> follow up",
"ts": "1700000000.000400",
"thread_ts": "1700000000.000100",
}
},
)
await channel._on_socket_request(client, req)
channel._handle_message.assert_awaited_once()
kwargs = channel._handle_message.await_args.kwargs
assert kwargs["session_key"] == "slack:C123:1700000000.000100"
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_slack_slash_command_skips_thread_context() -> None: async def test_slack_slash_command_skips_thread_context() -> None:
channel = SlackChannel(SlackConfig(enabled=True, allow_from=[]), MessageBus()) channel = SlackChannel(SlackConfig(enabled=True, allow_from=[]), MessageBus())
+25 -2
View File
@@ -230,9 +230,30 @@ class WeixinChannel(BaseChannel):
self.logger.error("Failed to load Weixin account state", exc_info=True) self.logger.error("Failed to load Weixin account state", exc_info=True)
return False return False
def _save_state(self) -> None: def _save_state(self, *, force: bool = False) -> None:
state_file = self._get_state_dir() / "account.json" state_file = self._get_state_dir() / "account.json"
with suppress(Exception): with suppress(Exception):
if not force and state_file.exists():
persisted: object = None
try:
persisted = json.loads(state_file.read_text())
except Exception:
persisted = None
persisted_token = ""
if isinstance(persisted, dict):
persisted_mapping = cast(dict[str, object], persisted)
persisted_token = str(persisted_mapping.get("token", "") or "")
configured_token_is_authoritative: bool = bool(self.config.token) and (
self._token == self.config.token
)
if (
persisted_token
and persisted_token != self._token
and not configured_token_is_authoritative
):
# A concurrent QR login may have committed a newer token.
# Never let an older runtime snapshot overwrite it.
return
data = { data = {
"token": self._token, "token": self._token,
"get_updates_buf": self._get_updates_buf, "get_updates_buf": self._get_updates_buf,
@@ -489,7 +510,7 @@ class WeixinChannel(BaseChannel):
self._token = token self._token = token
if base_url: if base_url:
self.config.base_url = base_url self.config.base_url = base_url
self._save_state() self._save_state(force=True)
async def connect_close_client(self) -> None: async def connect_close_client(self) -> None:
self._running = False self._running = False
@@ -613,6 +634,8 @@ class WeixinChannel(BaseChannel):
remaining = self._session_pause_remaining_s() remaining = self._session_pause_remaining_s()
if remaining > 0: if remaining > 0:
await asyncio.sleep(remaining) await asyncio.sleep(remaining)
if not self.config.token:
self._load_state()
return return
body: dict[str, Any] = { body: dict[str, Any] = {
@@ -98,6 +98,80 @@ def test_save_and_load_state_persists_context_tokens(tmp_path) -> None:
assert restored._context_tokens == {"wx-user": "ctx-1"} assert restored._context_tokens == {"wx-user": "ctx-1"}
def test_save_state_preserves_token_committed_by_another_instance(tmp_path) -> None:
channel = WeixinChannel(
WeixinConfig(enabled=True, allow_from=["*"], state_dir=str(tmp_path)),
MessageBus(),
)
channel._token = "old-token"
channel._save_state()
replacement = {
"token": "new-token",
"base_url": "https://new.example",
"get_updates_buf": "",
"context_tokens": {},
"typing_tickets": {},
}
(tmp_path / "account.json").write_text(json.dumps(replacement), encoding="utf-8")
channel._get_updates_buf = "stale-cursor"
channel._save_state()
assert json.loads((tmp_path / "account.json").read_text()) == replacement
def test_save_state_force_overwrites_replaced_token(tmp_path) -> None:
channel = WeixinChannel(
WeixinConfig(enabled=True, allow_from=["*"], state_dir=str(tmp_path)),
MessageBus(),
)
(tmp_path / "account.json").write_text(json.dumps({"token": "old-token"}), encoding="utf-8")
channel.connect_commit_account(token="new-token", base_url="https://new.example")
saved = json.loads((tmp_path / "account.json").read_text())
assert saved["token"] == "new-token"
assert saved["base_url"] == "https://new.example"
def test_save_state_persists_explicit_config_token_over_stale_state(tmp_path) -> None:
channel = WeixinChannel(
WeixinConfig(
enabled=True,
allow_from=["*"],
token="configured-token",
state_dir=str(tmp_path),
),
MessageBus(),
)
channel._token = "configured-token"
channel._get_updates_buf = "current-cursor"
(tmp_path / "account.json").write_text(
json.dumps({"token": "stale-token", "get_updates_buf": "stale-cursor"}),
encoding="utf-8",
)
channel._save_state()
saved = json.loads((tmp_path / "account.json").read_text())
assert saved["token"] == "configured-token"
assert saved["get_updates_buf"] == "current-cursor"
def test_save_state_with_empty_runtime_token_preserves_persisted_account(tmp_path) -> None:
channel = WeixinChannel(
WeixinConfig(enabled=True, allow_from=["*"], state_dir=str(tmp_path)),
MessageBus(),
)
persisted = {"token": "persisted-token", "get_updates_buf": "persisted-cursor"}
(tmp_path / "account.json").write_text(json.dumps(persisted), encoding="utf-8")
channel._save_state()
assert json.loads((tmp_path / "account.json").read_text()) == persisted
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_process_message_deduplicates_inbound_ids() -> None: async def test_process_message_deduplicates_inbound_ids() -> None:
channel, bus = _make_channel() channel, bus = _make_channel()
@@ -462,6 +536,56 @@ async def test_poll_once_pauses_session_on_expired_errcode() -> None:
assert channel._session_pause_remaining_s() > 0 assert channel._session_pause_remaining_s() > 0
@pytest.mark.asyncio
async def test_poll_once_reloads_refreshed_state_after_session_pause(
tmp_path, monkeypatch: pytest.MonkeyPatch
) -> None:
channel = WeixinChannel(
WeixinConfig(enabled=True, allow_from=["*"], state_dir=str(tmp_path)),
MessageBus(),
)
channel._token = "old-token"
channel._save_state()
(tmp_path / "account.json").write_text(
json.dumps({"token": "new-token", "base_url": "https://new.example"}),
encoding="utf-8",
)
channel._session_pause_until = time.time() + 10
monkeypatch.setattr(weixin_mod.asyncio, "sleep", AsyncMock())
await channel._poll_once()
assert channel._token == "new-token"
assert channel.config.base_url == "https://new.example"
@pytest.mark.asyncio
async def test_poll_once_keeps_explicit_token_after_session_pause(
tmp_path, monkeypatch: pytest.MonkeyPatch
) -> None:
channel = WeixinChannel(
WeixinConfig(
enabled=True,
allow_from=["*"],
token="configured-token",
state_dir=str(tmp_path),
),
MessageBus(),
)
channel._token = "configured-token"
(tmp_path / "account.json").write_text(
json.dumps({"token": "stale-token", "base_url": "https://stale.example"}),
encoding="utf-8",
)
channel._session_pause_until = time.time() + 10
monkeypatch.setattr(weixin_mod.asyncio, "sleep", AsyncMock())
await channel._poll_once()
assert channel._token == "configured-token"
assert channel.config.base_url == "https://ilinkai.weixin.qq.com"
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_qr_login_refreshes_expired_qr_and_then_succeeds( async def test_qr_login_refreshes_expired_qr_and_then_succeeds(
no_qr_poll_delay, no_qr_poll_delay,
+92 -9
View File
@@ -12,7 +12,9 @@ from collections import OrderedDict
from contextlib import suppress from contextlib import suppress
from pathlib import Path from pathlib import Path
from typing import Any, Literal, NamedTuple, cast from typing import Any, Literal, NamedTuple, cast
from urllib.parse import urlparse
import httpx
from pydantic import Field from pydantic import Field
from nanobot.bus.events import OutboundMessage from nanobot.bus.events import OutboundMessage
@@ -20,6 +22,7 @@ from nanobot.bus.queue import MessageBus
from nanobot.channels.base import BaseChannel from nanobot.channels.base import BaseChannel
from nanobot.config.paths import get_media_dir, get_runtime_subdir from nanobot.config.paths import get_media_dir, get_runtime_subdir
from nanobot.config.schema import Base from nanobot.config.schema import Base
from nanobot.security.network import PinnedDNSAsyncTransport
class WhatsAppConfig(Base): class WhatsAppConfig(Base):
@@ -39,6 +42,8 @@ class _NeonizeAPI(NamedTuple):
MessageEv: Any MessageEv: Any
PairStatusEv: Any PairStatusEv: Any
build_jid: Any build_jid: Any
detect_mime: Any
detect_buffer: Any
class _MediaInfo(NamedTuple): class _MediaInfo(NamedTuple):
@@ -52,6 +57,15 @@ class _MediaInfo(NamedTuple):
_NEONIZE_API: _NeonizeAPI | None = None _NEONIZE_API: _NeonizeAPI | None = None
_JID_RE = re.compile(r"^(?P<user>[^@]+)@(?P<server>[^@]+)$") _JID_RE = re.compile(r"^(?P<user>[^@]+)@(?P<server>[^@]+)$")
_LEGACY_BRIDGE_CONFIG_FIELDS = ("bridgeUrl", "bridgeToken", "bridge_url", "bridge_token") _LEGACY_BRIDGE_CONFIG_FIELDS = ("bridgeUrl", "bridgeToken", "bridge_url", "bridge_token")
_REMOTE_MEDIA_MAX_BYTES = 32 * 1024 * 1024
_REMOTE_MEDIA_MAX_REDIRECTS = 5
_REMOTE_MEDIA_TIMEOUT_SECONDS = 120.0
# OGG is intentionally excluded: WhatsApp accepts only mono Opus, which MIME sniffing cannot prove.
_DIRECT_AUDIO_MIMETYPES = {"audio/aac", "audio/amr", "audio/mp4", "audio/mpeg"}
_MIMETYPE_ALIASES = {
"audio/x-hx-aac-adts": "audio/aac",
"audio/x-m4a": "audio/mp4",
}
def _default_database_path() -> Path: def _default_database_path() -> Path:
@@ -68,9 +82,15 @@ def _load_neonize() -> _NeonizeAPI:
return _NEONIZE_API return _NEONIZE_API
try: try:
import magic
from neonize.aioze.client import NewAClient from neonize.aioze.client import NewAClient
from neonize.aioze.events import ConnectedEv, DisconnectedEv, MessageEv, PairStatusEv from neonize.aioze.events import ConnectedEv, DisconnectedEv, MessageEv, PairStatusEv
from neonize.utils.jid import build_jid from neonize.utils.jid import build_jid
detect_mime = getattr(magic, "from_file", None)
detect_buffer = getattr(magic, "from_buffer", None)
if not callable(detect_mime) or not callable(detect_buffer):
raise ImportError("python-magic does not expose from_file/from_buffer")
except ImportError as exc: except ImportError as exc:
raise RuntimeError( raise RuntimeError(
"WhatsApp dependencies not installed. Run: nanobot plugins enable whatsapp" "WhatsApp dependencies not installed. Run: nanobot plugins enable whatsapp"
@@ -83,6 +103,8 @@ def _load_neonize() -> _NeonizeAPI:
MessageEv=MessageEv, MessageEv=MessageEv,
PairStatusEv=PairStatusEv, PairStatusEv=PairStatusEv,
build_jid=build_jid, build_jid=build_jid,
detect_mime=detect_mime,
detect_buffer=detect_buffer,
) )
return _NEONIZE_API return _NEONIZE_API
@@ -417,23 +439,84 @@ class WhatsAppChannel(BaseChannel):
return api.build_jid(user, server) return api.build_jid(user, server)
async def _send_media(self, client: Any, to: Any, media_path: str) -> None: async def _send_media(self, client: Any, to: Any, media_path: str) -> None:
path = str(Path(media_path).expanduser()) source: str | bytes
mime, _ = mimetypes.guess_type(path) if media_path.startswith(("http://", "https://")):
mimetype = mime or "application/octet-stream" source = await self._fetch_remote_media(media_path)
filename = Path(urlparse(media_path).path).name or "attachment"
else:
source = str(Path(media_path).expanduser())
filename = Path(source).name
mimetype = self._detect_mimetype(source)
if mimetype.startswith("image/"): if mimetype.startswith("image/"):
await client.send_image(to, path) await client.send_image(to, source)
elif mimetype.startswith("video/"): elif mimetype.startswith("video/"):
await client.send_video(to, path) await client.send_video(to, source)
elif mimetype.startswith("audio/"): elif mimetype in _DIRECT_AUDIO_MIMETYPES:
await client.send_audio(to, path) await client.send_audio(to, source)
else: else:
await client.send_document( await client.send_document(
to, to,
path, source,
filename=Path(path).name, filename=filename,
mimetype=mimetype, mimetype=mimetype,
) )
async def _fetch_remote_media(self, url: str) -> bytes:
timeout = httpx.Timeout(_REMOTE_MEDIA_TIMEOUT_SECONDS, connect=10.0)
async with httpx.AsyncClient(
transport=PinnedDNSAsyncTransport(),
follow_redirects=True,
max_redirects=_REMOTE_MEDIA_MAX_REDIRECTS,
timeout=timeout,
trust_env=False,
) as http:
async with http.stream("GET", url) as response:
response.raise_for_status()
declared_size = response.headers.get("content-length")
if (
declared_size
and declared_size.isdigit()
and int(declared_size) > _REMOTE_MEDIA_MAX_BYTES
):
raise ValueError(
f"Remote WhatsApp media exceeds the {_REMOTE_MEDIA_MAX_BYTES}-byte limit"
)
chunks: list[bytes] = []
total = 0
async for chunk in response.aiter_bytes():
total += len(chunk)
if total > _REMOTE_MEDIA_MAX_BYTES:
raise ValueError(
f"Remote WhatsApp media exceeds the {_REMOTE_MEDIA_MAX_BYTES}-byte limit"
)
chunks.append(chunk)
return b"".join(chunks)
def _detect_mimetype(self, source: str | bytes) -> str:
try:
api = _load_neonize()
detected = (
api.detect_buffer(source, mime=True)
if isinstance(source, bytes)
else api.detect_mime(source, mime=True)
)
except Exception as exc:
label = f"{len(source)} downloaded bytes" if isinstance(source, bytes) else source
self.logger.debug("Failed to inspect WhatsApp media {}: {}", label, exc)
detected = None
if isinstance(detected, str) and "/" in detected:
mimetype = detected.partition(";")[0].strip().lower()
return _MIMETYPE_ALIASES.get(mimetype, mimetype)
if isinstance(source, bytes):
return "application/octet-stream"
guessed, _ = mimetypes.guess_type(source)
return guessed or "application/octet-stream"
def _register_handlers( def _register_handlers(
self, self,
client: Any, client: Any,
@@ -1,11 +1,13 @@
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import mimetypes
import sys import sys
import types import types
from types import SimpleNamespace from types import SimpleNamespace
from unittest.mock import AsyncMock, MagicMock from unittest.mock import AsyncMock, MagicMock
import httpx
import pytest import pytest
import nanobot.channels.whatsapp.runtime as whatsapp_module import nanobot.channels.whatsapp.runtime as whatsapp_module
@@ -78,7 +80,21 @@ def _make_channel(config: dict | None = None) -> WhatsAppChannel:
return ch return ch
def _patch_neonize_api(monkeypatch) -> None: def _make_send_client() -> SimpleNamespace:
return SimpleNamespace(
send_message=AsyncMock(),
send_image=AsyncMock(),
send_video=AsyncMock(),
send_audio=AsyncMock(),
send_document=AsyncMock(),
)
def _patch_neonize_api(monkeypatch, detect_mime=None, detect_buffer=None) -> None:
detect_mime = detect_mime or (
lambda path, *, mime: mimetypes.guess_type(path)[0] or "application/octet-stream"
)
detect_buffer = detect_buffer or (lambda data, *, mime: "application/octet-stream")
monkeypatch.setattr( monkeypatch.setattr(
whatsapp_module, whatsapp_module,
"_NEONIZE_API", "_NEONIZE_API",
@@ -89,6 +105,8 @@ def _patch_neonize_api(monkeypatch) -> None:
MessageEv=object(), MessageEv=object(),
PairStatusEv=object(), PairStatusEv=object(),
build_jid=lambda user, server="s.whatsapp.net": (user, server), build_jid=lambda user, server="s.whatsapp.net": (user, server),
detect_mime=detect_mime,
detect_buffer=detect_buffer,
), ),
) )
@@ -178,13 +196,7 @@ async def test_login_fails_when_connect_task_fails(monkeypatch) -> None:
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_send_text_uses_neonize_send_message(monkeypatch) -> None: async def test_send_text_uses_neonize_send_message(monkeypatch) -> None:
_patch_neonize_api(monkeypatch) _patch_neonize_api(monkeypatch)
client = SimpleNamespace( client = _make_send_client()
send_message=AsyncMock(),
send_image=AsyncMock(),
send_video=AsyncMock(),
send_audio=AsyncMock(),
send_document=AsyncMock(),
)
ch = _make_channel() ch = _make_channel()
ch._client = client ch._client = client
ch._connected = True ch._connected = True
@@ -197,13 +209,7 @@ async def test_send_text_uses_neonize_send_message(monkeypatch) -> None:
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_send_media_dispatches_by_mimetype(monkeypatch) -> None: async def test_send_media_dispatches_by_mimetype(monkeypatch) -> None:
_patch_neonize_api(monkeypatch) _patch_neonize_api(monkeypatch)
client = SimpleNamespace( client = _make_send_client()
send_message=AsyncMock(),
send_image=AsyncMock(),
send_video=AsyncMock(),
send_audio=AsyncMock(),
send_document=AsyncMock(),
)
ch = _make_channel() ch = _make_channel()
ch._client = client ch._client = client
ch._connected = True ch._connected = True
@@ -213,14 +219,14 @@ async def test_send_media_dispatches_by_mimetype(monkeypatch) -> None:
channel="whatsapp", channel="whatsapp",
chat_id="12345@s.whatsapp.net", chat_id="12345@s.whatsapp.net",
content="", content="",
media=["photo.jpg", "clip.mp4", "voice.ogg", "report.pdf"], media=["photo.jpg", "clip.mp4", "voice.mp3", "report.pdf"],
) )
) )
jid = ("12345", "s.whatsapp.net") jid = ("12345", "s.whatsapp.net")
client.send_image.assert_awaited_once_with(jid, "photo.jpg") client.send_image.assert_awaited_once_with(jid, "photo.jpg")
client.send_video.assert_awaited_once_with(jid, "clip.mp4") client.send_video.assert_awaited_once_with(jid, "clip.mp4")
client.send_audio.assert_awaited_once_with(jid, "voice.ogg") client.send_audio.assert_awaited_once_with(jid, "voice.mp3")
client.send_document.assert_awaited_once_with( client.send_document.assert_awaited_once_with(
jid, jid,
"report.pdf", "report.pdf",
@@ -229,6 +235,191 @@ async def test_send_media_dispatches_by_mimetype(monkeypatch) -> None:
) )
@pytest.mark.asyncio
async def test_send_mislabeled_audio_as_document(monkeypatch) -> None:
_patch_neonize_api(monkeypatch, detect_mime=lambda path, *, mime: "audio/x-wav")
client = _make_send_client()
ch = _make_channel()
ch._client = client
ch._connected = True
await ch.send(
OutboundMessage(
channel="whatsapp",
chat_id="12345@s.whatsapp.net",
content="",
media=["recording.mpeg"],
)
)
jid = ("12345", "s.whatsapp.net")
client.send_document.assert_awaited_once_with(
jid,
"recording.mpeg",
filename="recording.mpeg",
mimetype="audio/x-wav",
)
client.send_video.assert_not_awaited()
@pytest.mark.asyncio
async def test_send_remote_mislabeled_audio_as_document(monkeypatch) -> None:
payload = b"remote wav payload"
media_url = "https://cdn.example/recording.mpeg?token=secret"
def handle_request(request: httpx.Request) -> httpx.Response:
assert str(request.url) == media_url
return httpx.Response(200, content=payload)
monkeypatch.setattr(
whatsapp_module,
"PinnedDNSAsyncTransport",
lambda: httpx.MockTransport(handle_request),
)
def detect_buffer(data: bytes, *, mime: bool) -> str:
assert data == payload
assert mime is True
return "audio/x-wav"
_patch_neonize_api(
monkeypatch,
detect_buffer=detect_buffer,
)
client = _make_send_client()
ch = _make_channel()
ch._client = client
ch._connected = True
await ch.send(
OutboundMessage(
channel="whatsapp",
chat_id="12345@s.whatsapp.net",
content="",
media=[media_url],
)
)
jid = ("12345", "s.whatsapp.net")
client.send_document.assert_awaited_once_with(
jid,
payload,
filename="recording.mpeg",
mimetype="audio/x-wav",
)
client.send_video.assert_not_awaited()
@pytest.mark.asyncio
async def test_send_remote_media_blocks_private_url(monkeypatch) -> None:
_patch_neonize_api(monkeypatch)
client = _make_send_client()
ch = _make_channel()
ch._client = client
ch._connected = True
with pytest.raises(httpx.RequestError, match="private/internal"):
await ch.send(
OutboundMessage(
channel="whatsapp",
chat_id="12345@s.whatsapp.net",
content="",
media=["http://127.0.0.1/recording.mpeg"],
)
)
client.send_video.assert_not_awaited()
client.send_document.assert_not_awaited()
@pytest.mark.asyncio
async def test_send_remote_media_enforces_download_limit(monkeypatch) -> None:
monkeypatch.setattr(whatsapp_module, "_REMOTE_MEDIA_MAX_BYTES", 3)
monkeypatch.setattr(
whatsapp_module,
"PinnedDNSAsyncTransport",
lambda: httpx.MockTransport(lambda request: httpx.Response(200, content=b"1234")),
)
_patch_neonize_api(monkeypatch)
client = _make_send_client()
ch = _make_channel()
ch._client = client
ch._connected = True
with pytest.raises(ValueError, match="exceeds the 3-byte limit"):
await ch.send(
OutboundMessage(
channel="whatsapp",
chat_id="12345@s.whatsapp.net",
content="",
media=["https://cdn.example/recording.mpeg"],
)
)
client.send_video.assert_not_awaited()
client.send_document.assert_not_awaited()
@pytest.mark.asyncio
async def test_send_unsupported_ogg_audio_as_document(monkeypatch) -> None:
_patch_neonize_api(monkeypatch, detect_mime=lambda path, *, mime: "audio/ogg")
client = _make_send_client()
ch = _make_channel()
ch._client = client
ch._connected = True
await ch.send(
OutboundMessage(
channel="whatsapp",
chat_id="12345@s.whatsapp.net",
content="",
media=["voice.ogg"],
)
)
jid = ("12345", "s.whatsapp.net")
client.send_document.assert_awaited_once_with(
jid,
"voice.ogg",
filename="voice.ogg",
mimetype="audio/ogg",
)
client.send_audio.assert_not_awaited()
@pytest.mark.parametrize(
("detected_mimetype", "filename"),
[
("audio/x-m4a", "recording.m4a"),
("audio/x-hx-aac-adts", "recording.aac"),
],
)
@pytest.mark.asyncio
async def test_send_supported_audio_magic_aliases_inline(
monkeypatch, detected_mimetype: str, filename: str
) -> None:
_patch_neonize_api(
monkeypatch,
detect_mime=lambda path, *, mime: detected_mimetype,
)
client = _make_send_client()
ch = _make_channel()
ch._client = client
ch._connected = True
await ch.send(
OutboundMessage(
channel="whatsapp",
chat_id="12345@s.whatsapp.net",
content="",
media=[filename],
)
)
client.send_audio.assert_awaited_once_with(("12345", "s.whatsapp.net"), filename)
client.send_document.assert_not_awaited()
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_send_when_disconnected_raises() -> None: async def test_send_when_disconnected_raises() -> None:
ch = _make_channel() ch = _make_channel()
+21 -6
View File
@@ -958,22 +958,34 @@ class OpenAICompatProvider(LLMProvider):
model: str | None, model: str | None,
reasoning_effort: str | None, reasoning_effort: str | None,
) -> bool: ) -> bool:
"""Use Responses API only for direct OpenAI requests that benefit from it.""" """Choose Responses for providers/models that explicitly support it."""
if self._api_type == "chat_completions": if self._api_type == "chat_completions":
return False return False
if self._spec and self._spec.name not in ("openai", "github_copilot"): spec_name = self._spec.name if self._spec is not None else None
model_name = self._request_model_name(model or self.default_model).lower()
supported_models = {
supported.lower()
for supported in getattr(self._spec, "responses_models", ())
}
model_responses = any(
model_name == supported or model_name.endswith(f"/{supported}")
for supported in supported_models
)
provider_responses = spec_name in ("openai", "github_copilot")
if not provider_responses and not model_responses:
return False return False
if self._api_type == "responses": if self._api_type == "responses":
# Explicit configuration means Responses is mandatory; do not # Explicit configuration means Responses is mandatory; do not
# consult the circuit breaker or fall back to Chat Completions. # consult the circuit breaker or fall back to Chat Completions.
return True return True
if self._spec is None or self._spec.name != "github_copilot": if provider_responses and (self._spec is None or self._spec.name != "github_copilot"):
if not _is_direct_openai_base(self._effective_base): if not _is_direct_openai_base(self._effective_base):
return False return False
model_name = (model or self.default_model).lower()
wants = False wants = False
if reasoning_effort and reasoning_effort.lower() != "none": if model_responses:
wants = True
elif reasoning_effort and reasoning_effort.lower() != "none":
wants = True wants = True
elif any(token in model_name for token in ("gpt-5", "o1", "o3", "o4")): elif any(token in model_name for token in ("gpt-5", "o1", "o3", "o4")):
wants = True wants = True
@@ -1099,11 +1111,13 @@ class OpenAICompatProvider(LLMProvider):
self._sanitize_empty_content(sanitized_state.pending_messages) self._sanitize_empty_content(sanitized_state.pending_messages)
) )
) )
preserve_reasoning = bool(self._spec and self._spec.name == "deepseek")
instructions, input_items, replayed = prepare_responses_input( instructions, input_items, replayed = prepare_responses_input(
sanitized_messages, sanitized_messages,
state=sanitized_state, state=sanitized_state,
provider=self._responses_state_provider(), provider=self._responses_state_provider(),
model=model_name, model=model_name,
preserve_reasoning=preserve_reasoning,
) )
body: dict[str, Any] = { body: dict[str, Any] = {
@@ -1131,7 +1145,7 @@ class OpenAICompatProvider(LLMProvider):
if self._supports_temperature(model_name, reasoning_effort): if self._supports_temperature(model_name, reasoning_effort):
body["temperature"] = temperature body["temperature"] = temperature
if not self._supports_temperature(model_name, reasoning_effort): if not self._supports_temperature(model_name, reasoning_effort) and not preserve_reasoning:
body["include"] = ["reasoning.encrypted_content"] body["include"] = ["reasoning.encrypted_content"]
if reasoning_effort and reasoning_effort.lower() != "none": if reasoning_effort and reasoning_effort.lower() != "none":
body["reasoning"] = {"effort": reasoning_effort} body["reasoning"] = {"effort": reasoning_effort}
@@ -1827,6 +1841,7 @@ class OpenAICompatProvider(LLMProvider):
_timed_stream(), _timed_stream(),
on_content_delta, on_content_delta,
on_tool_call_delta=on_tool_call_delta, on_tool_call_delta=on_tool_call_delta,
on_reasoning_delta=on_thinking_delta,
capture=capture, capture=capture,
) )
self._record_responses_success(model, reasoning_effort) self._record_responses_success(model, reasoning_effort)
@@ -12,7 +12,11 @@ def _as_json_object(value: object) -> dict[str, Any] | None:
return cast(dict[str, Any], value) if isinstance(value, dict) else None return cast(dict[str, Any], value) if isinstance(value, dict) else None
def convert_messages(messages: list[dict[str, Any]]) -> tuple[str, list[dict[str, Any]]]: def convert_messages(
messages: list[dict[str, Any]],
*,
preserve_reasoning: bool = False,
) -> tuple[str, list[dict[str, Any]]]:
"""Convert Chat Completions messages to Responses API input items. """Convert Chat Completions messages to Responses API input items.
Returns ``(system_prompt, input_items)`` where *system_prompt* is extracted Returns ``(system_prompt, input_items)`` where *system_prompt* is extracted
@@ -36,6 +40,13 @@ def convert_messages(messages: list[dict[str, Any]]) -> tuple[str, list[dict[str
continue continue
if role == "assistant": if role == "assistant":
if preserve_reasoning:
reasoning = msg.get("reasoning_content")
if isinstance(reasoning, str) and reasoning:
input_items.append({
"type": "reasoning",
"content": reasoning,
})
if isinstance(content, str) and content: if isinstance(content, str) and content:
message_id = _unique_item_id(f"msg_{idx}", used_item_ids) message_id = _unique_item_id(f"msg_{idx}", used_item_ids)
input_items.append({ input_items.append({
+35 -13
View File
@@ -69,7 +69,9 @@ def _response_object(value: object) -> dict[str, Any] | None:
return object_value return object_value
dump = getattr(value, "model_dump", None) dump = getattr(value, "model_dump", None)
if callable(dump): if callable(dump):
return _as_json_object(dump()) dumped = _as_json_object(dump())
if dumped is not None:
return dumped
try: try:
return _as_json_object(vars(value)) return _as_json_object(vars(value))
except TypeError: except TypeError:
@@ -444,6 +446,14 @@ def _extract_reasoning_summary_from_output(output: object) -> str | None:
for item in _response_object_list(output): for item in _response_object_list(output):
if item.get("type") != "reasoning": if item.get("type") != "reasoning":
continue continue
content = item.get("content")
if isinstance(content, str) and content:
parts.append(content)
elif isinstance(content, list):
for block in _response_object_list(cast(list[object], content)):
text = block.get("text")
if isinstance(text, str) and text:
parts.append(text)
for summary in _response_object_list(item.get("summary")): for summary in _response_object_list(item.get("summary")):
if summary.get("type") == "summary_text" and summary.get("text"): if summary.get("type") == "summary_text" and summary.get("text"):
text = summary.get("text") text = summary.get("text")
@@ -483,11 +493,9 @@ def parse_response_output(
if isinstance(refusal, str): if isinstance(refusal, str):
content_parts.append(refusal) content_parts.append(refusal)
elif item_type == "reasoning": elif item_type == "reasoning":
for s in _response_object_list(item.get("summary")): text = _extract_reasoning_summary_from_output([item])
if s.get("type") == "summary_text" and s.get("text"): if text:
text = s.get("text") reasoning_content = (reasoning_content or "") + text
if isinstance(text, str):
reasoning_content = (reasoning_content or "") + text
elif item_type == "function_call": elif item_type == "function_call":
call_id = item.get("call_id") or "" call_id = item.get("call_id") or ""
item_id = item.get("id") or "fc_0" item_id = item.get("id") or "fc_0"
@@ -532,6 +540,7 @@ async def consume_sdk_stream(
stream: Any, stream: Any,
on_content_delta: Callable[[str], Awaitable[None]] | None = None, on_content_delta: Callable[[str], Awaitable[None]] | None = None,
on_tool_call_delta: Callable[[dict[str, Any]], Awaitable[None]] | None = None, on_tool_call_delta: Callable[[dict[str, Any]], Awaitable[None]] | None = None,
on_reasoning_delta: Callable[[str], Awaitable[None]] | None = None,
capture: ResponsesStreamCapture | None = None, capture: ResponsesStreamCapture | None = None,
) -> tuple[str, list[ToolCallRequest], str, dict[str, int], str | None]: ) -> tuple[str, list[ToolCallRequest], str, dict[str, int], str | None]:
"""Consume an SDK async stream from ``client.responses.create(stream=True)``.""" """Consume an SDK async stream from ``client.responses.create(stream=True)``."""
@@ -542,6 +551,7 @@ async def consume_sdk_stream(
finish_reason = "stop" finish_reason = "stop"
usage: dict[str, int] = {} usage: dict[str, int] = {}
reasoning_content: str | None = None reasoning_content: str | None = None
streamed_reasoning = False
refusal_seen = False refusal_seen = False
refusal_deltas: dict[tuple[str | None, int | None], str] = {} refusal_deltas: dict[tuple[str | None, int | None], str] = {}
emitted_refusal_text = "" emitted_refusal_text = ""
@@ -572,6 +582,19 @@ async def consume_sdk_stream(
content += delta_text content += delta_text
if on_content_delta and delta_text: if on_content_delta and delta_text:
await on_content_delta(delta_text) await on_content_delta(delta_text)
elif event_type == "response.reasoning_text.delta":
delta_text = getattr(event, "delta", "") or ""
if delta_text:
reasoning_content = (reasoning_content or "") + delta_text
streamed_reasoning = True
if on_reasoning_delta:
await on_reasoning_delta(delta_text)
elif event_type == "response.reasoning_text.done":
text = getattr(event, "text", "") or ""
if text and not streamed_reasoning and not reasoning_content:
reasoning_content = text
if on_reasoning_delta:
await on_reasoning_delta(text)
elif event_type == "response.refusal.delta": elif event_type == "response.refusal.delta":
refusal_seen = True refusal_seen = True
delta_text = getattr(event, "delta", None) delta_text = getattr(event, "delta", None)
@@ -689,13 +712,12 @@ async def consume_sdk_stream(
"completion_tokens": int(getattr(usage_obj, "output_tokens", 0) or 0), "completion_tokens": int(getattr(usage_obj, "output_tokens", 0) or 0),
"total_tokens": int(getattr(usage_obj, "total_tokens", 0) or 0), "total_tokens": int(getattr(usage_obj, "total_tokens", 0) or 0),
} }
for out_item in cast(list[Any], getattr(resp, "output", None) or []): if not reasoning_content:
if getattr(out_item, "type", None) == "reasoning": reasoning_content = _extract_reasoning_summary_from_output(
for s in cast(list[Any], getattr(out_item, "summary", None) or []): getattr(resp, "output", None)
if getattr(s, "type", None) == "summary_text": )
text = getattr(s, "text", None) if reasoning_content and on_reasoning_delta:
if text: await on_reasoning_delta(reasoning_content)
reasoning_content = (reasoning_content or "") + text
elif event_type in {"error", "response.failed"}: elif event_type in {"error", "response.failed"}:
detail = getattr(event, "error", None) or getattr(event, "message", None) or event detail = getattr(event, "error", None) or getattr(event, "message", None) or event
raise RuntimeError(f"Response failed: {str(detail)[:500]}") raise RuntimeError(f"Response failed: {str(detail)[:500]}")
+9 -2
View File
@@ -43,6 +43,7 @@ def prepare_responses_input(
state: ProviderConversationState | None, state: ProviderConversationState | None,
provider: str, provider: str,
model: str, model: str,
preserve_reasoning: bool = False,
) -> tuple[str, list[dict[str, Any]], bool]: ) -> tuple[str, list[dict[str, Any]], bool]:
"""Build a request from exact prior items plus only newly appended messages. """Build a request from exact prior items plus only newly appended messages.
@@ -50,7 +51,10 @@ def prepare_responses_input(
When no compatible state exists, it is converted normally as a safe When no compatible state exists, it is converted normally as a safe
fallback. fallback.
""" """
instructions, fallback_items = convert_messages(messages) instructions, fallback_items = convert_messages(
messages,
preserve_reasoning=preserve_reasoning,
)
if state is None or not responses_state_matches( if state is None or not responses_state_matches(
state, state,
provider=provider, provider=provider,
@@ -62,7 +66,10 @@ def prepare_responses_input(
if prior_items is None: if prior_items is None:
return instructions, fallback_items, False return instructions, fallback_items, False
_, delta_items = convert_messages(state.pending_messages) _, delta_items = convert_messages(
state.pending_messages,
preserve_reasoning=preserve_reasoning,
)
logger.debug( logger.debug(
"Replaying Responses state: prior_items={} pending_messages={}", "Replaying Responses state: prior_items={} pending_messages={}",
len(prior_items), len(prior_items),
+6
View File
@@ -111,6 +111,11 @@ class ProviderSpec:
# Substring match against the wire model name (lowercased). # Substring match against the wire model name (lowercased).
implicit_reasoning_models: tuple[str, ...] = () implicit_reasoning_models: tuple[str, ...] = ()
# Models that expose the OpenAI Responses wire format. This is model-level
# because providers may add Responses support incrementally (DeepSeek V4
# Flash is supported before V4 Pro).
responses_models: tuple[str, ...] = ()
# When the model returns content as a list of {"type":"thinking",...} + # When the model returns content as a list of {"type":"thinking",...} +
# {"type":"text",...} blocks, extract the thinking text into # {"type":"text",...} blocks, extract the thinking text into
# reasoning_content. Mistral's Magistral / reasoning-enabled responses use # reasoning_content. Mistral's Magistral / reasoning-enabled responses use
@@ -461,6 +466,7 @@ PROVIDERS: tuple[ProviderSpec, ...] = (
backend="openai_compat", backend="openai_compat",
default_api_base="https://api.deepseek.com", default_api_base="https://api.deepseek.com",
thinking_style="thinking_type", thinking_style="thinking_type",
responses_models=("deepseek-v4-flash",),
), ),
# Gemini: Google's OpenAI-compatible endpoint # Gemini: Google's OpenAI-compatible endpoint
ProviderSpec( ProviderSpec(
+1 -1
View File
@@ -51,7 +51,7 @@ dependencies = [
"filelock>=3.25.2", "filelock>=3.25.2",
"watchfiles>=1.1.1,<2.0.0", "watchfiles>=1.1.1,<2.0.0",
"packaging>=24.0", "packaging>=24.0",
"tzdata>=2025.2; sys_platform == 'win32'", "tzdata>=2025.2",
"defusedxml>=0.7.1,<1.0.0", "defusedxml>=0.7.1,<1.0.0",
"pypdf>=5.0.0,<6.0.0", "pypdf>=5.0.0,<6.0.0",
"python-docx>=1.1.0,<2.0.0", "python-docx>=1.1.0,<2.0.0",
+1 -1
View File
@@ -2558,7 +2558,7 @@ def test_optional_dependency_metadata_for_enable():
): ):
assert not any(dep.startswith(dep_name) for dep in required) assert not any(dep.startswith(dep_name) for dep in required)
for dependency in ( for dependency in (
"tzdata>=2025.2; sys_platform == 'win32'", "tzdata>=2025.2",
"defusedxml>=0.7.1,<1.0.0", "defusedxml>=0.7.1,<1.0.0",
"pypdf>=5.0.0,<6.0.0", "pypdf>=5.0.0,<6.0.0",
"python-docx>=1.1.0,<2.0.0", "python-docx>=1.1.0,<2.0.0",
+30
View File
@@ -1,4 +1,8 @@
import json import json
import os
import subprocess
import sys
import textwrap
import warnings import warnings
import pytest import pytest
@@ -42,6 +46,32 @@ def test_agent_timezone_rejects_unknown_iana_name() -> None:
Config.model_validate({"agents": {"defaults": {"timezone": "Not/AZone"}}}) Config.model_validate({"agents": {"defaults": {"timezone": "Not/AZone"}}})
def test_agent_timezones_use_packaged_data_without_system_database() -> None:
script = textwrap.dedent(
"""\
from zoneinfo import TZPATH
from nanobot.config.schema import Config
assert not TZPATH
for name in ("UTC", "Asia/Shanghai"):
config = Config.model_validate({"agents": {"defaults": {"timezone": name}}})
serialized = config.model_dump(mode="json", by_alias=True)
restored = Config.model_validate(serialized)
assert restored.agents.defaults.timezone == name
"""
)
result = subprocess.run(
[sys.executable, "-c", script],
env=os.environ | {"PYTHONTZPATH": ""},
capture_output=True,
text=True,
check=False,
)
assert result.returncode == 0, result.stderr
def test_provider_api_type_accepts_exact_values_only() -> None: def test_provider_api_type_accepts_exact_values_only() -> None:
config = Config.model_validate({ config = Config.model_validate({
"providers": { "providers": {
+56
View File
@@ -150,6 +150,22 @@ class TestConvertMessages:
assert items[0]["content"][0]["type"] == "output_text" assert items[0]["content"][0]["type"] == "output_text"
assert items[0]["content"][0]["text"] == "I'll help" assert items[0]["content"][0]["text"] == "I'll help"
def test_preserves_deepseek_reasoning_content(self):
_, items = convert_messages([
{"role": "assistant", "reasoning_content": "think first", "content": "answer"},
], preserve_reasoning=True)
assert items == [
{"type": "reasoning", "content": "think first"},
{
"type": "message",
"role": "assistant",
"content": [{"type": "output_text", "text": "answer"}],
"status": "completed",
"id": "msg_0",
},
]
def test_assistant_empty_content_skipped(self): def test_assistant_empty_content_skipped(self):
_, items = convert_messages([{"role": "assistant", "content": ""}]) _, items = convert_messages([{"role": "assistant", "content": ""}])
assert len(items) == 0 assert len(items) == 0
@@ -539,6 +555,22 @@ class TestParseResponseOutput:
assert result.content == "42" assert result.content == "42"
assert result.reasoning_content == "I think therefore I am." assert result.reasoning_content == "I think therefore I am."
def test_deepseek_reasoning_content_extracted(self):
resp = {
"output": [
{"type": "reasoning", "content": "think first"},
{"type": "message", "content": [
{"type": "output_text", "text": "answer"},
]},
],
"status": "completed", "usage": {},
}
result = parse_response_output(resp)
assert result.content == "answer"
assert result.reasoning_content == "think first"
def test_empty_output(self): def test_empty_output(self):
resp = {"output": [], "status": "completed", "usage": {}} resp = {"output": [], "status": "completed", "usage": {}}
result = parse_response_output(resp) result = parse_response_output(resp)
@@ -1633,6 +1665,30 @@ class TestConsumeSdkStream:
_, _, _, _, reasoning = await consume_sdk_stream(stream()) _, _, _, _, reasoning = await consume_sdk_stream(stream())
assert reasoning == "thinking..." assert reasoning == "thinking..."
@pytest.mark.asyncio
async def test_deepseek_reasoning_text_streamed(self):
events = [
MagicMock(type="response.reasoning_text.delta", delta="step 1 "),
MagicMock(type="response.reasoning_text.delta", delta="step 2"),
MagicMock(type="response.reasoning_text.done", text="step 1 step 2"),
]
emitted: list[str] = []
async def stream():
for event in events:
yield event
async def on_reasoning_delta(delta: str) -> None:
emitted.append(delta)
_, _, _, _, reasoning = await consume_sdk_stream(
stream(),
on_reasoning_delta=on_reasoning_delta,
)
assert reasoning == "step 1 step 2"
assert emitted == ["step 1 ", "step 2"]
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_error_event_raises(self): async def test_error_event_raises(self):
ev = MagicMock(type="error", error="rate_limit_exceeded") ev = MagicMock(type="error", error="rate_limit_exceeded")
@@ -29,6 +29,32 @@ def test_responses_api_available_by_default(provider):
assert provider._should_use_responses_api("gpt-5", None) is True assert provider._should_use_responses_api("gpt-5", None) is True
def test_deepseek_v4_flash_uses_responses_by_model(provider):
provider._spec = type("Spec", (), {
"name": "deepseek",
"responses_models": ("deepseek-v4-flash",),
"strip_model_prefix": False,
"strip_model_prefixes": (),
})()
provider._effective_base = "https://api.deepseek.com"
provider.default_model = "deepseek-v4-flash"
assert provider._should_use_responses_api("deepseek-v4-flash", None) is True
assert provider._should_use_responses_api("deepseek-v4-pro", None) is False
def test_deepseek_v4_flash_matches_provider_prefixed_model(provider):
provider._spec = type("Spec", (), {
"name": "deepseek",
"responses_models": ("deepseek-v4-flash",),
"strip_model_prefix": False,
"strip_model_prefixes": (),
})()
provider._effective_base = "https://api.deepseek.com"
assert provider._should_use_responses_api("deepseek/deepseek-v4-flash", None) is True
def test_direct_openai_enables_server_compaction(provider): def test_direct_openai_enables_server_compaction(provider):
provider._extra_body = {} provider._extra_body = {}
@@ -542,7 +542,7 @@ export const ThreadViewport = forwardRef<ThreadViewportHandle, ThreadViewportPro
const distance = el.scrollHeight - el.scrollTop - el.clientHeight; const distance = el.scrollHeight - el.scrollTop - el.clientHeight;
const near = distance < NEAR_BOTTOM_PX; const near = distance < NEAR_BOTTOM_PX;
const owner = threadMotionRef.current?.observeScroll(near) ?? "automatic"; const owner = threadMotionRef.current?.observeScroll(near) ?? "automatic";
const logicallyAtBottom = owner === "automatic" || near; const logicallyAtBottom = owner === "automatic" || (owner === "navigation" && near);
setAtBottom((current) => setAtBottom((current) =>
current === logicallyAtBottom ? current : logicallyAtBottom, current === logicallyAtBottom ? current : logicallyAtBottom,
); );
@@ -557,6 +557,7 @@ export const ThreadViewport = forwardRef<ThreadViewportHandle, ThreadViewportPro
if (!direction) return; if (!direction) return;
threadMotionRef.current?.handleUserScrollIntent( threadMotionRef.current?.handleUserScrollIntent(
canScrollInDirection(el, direction), canScrollInDirection(el, direction),
direction === "forward",
); );
}; };
const handleWheel = (event: WheelEvent) => { const handleWheel = (event: WheelEvent) => {
@@ -572,20 +573,21 @@ export const ThreadViewport = forwardRef<ThreadViewportHandle, ThreadViewportPro
const handlePointerDown = (event: PointerEvent) => { const handlePointerDown = (event: PointerEvent) => {
if (event.button === 0 && event.target === el) yieldCameraToUser(); if (event.button === 0 && event.target === el) yieldCameraToUser();
}; };
let touchStartY: number | null = null; let lastTouchY: number | null = null;
const handleTouchStart = (event: TouchEvent) => { const handleTouchStart = (event: TouchEvent) => {
touchStartY = event.touches[0]?.clientY ?? null; lastTouchY = event.touches[0]?.clientY ?? null;
}; };
const handleTouchMove = (event: TouchEvent) => { const handleTouchMove = (event: TouchEvent) => {
const currentY = event.touches[0]?.clientY; const currentY = event.touches[0]?.clientY;
const scrollDeltaY = const scrollDeltaY =
touchStartY !== null && currentY !== undefined lastTouchY !== null && currentY !== undefined
? touchStartY - currentY ? lastTouchY - currentY
: 0; : 0;
lastTouchY = currentY ?? null;
handleDirectionalInput(directionFromDelta(scrollDeltaY)); handleDirectionalInput(directionFromDelta(scrollDeltaY));
}; };
const handleTouchEnd = () => { const handleTouchEnd = () => {
touchStartY = null; lastTouchY = null;
}; };
const handleKeyDown = (event: KeyboardEvent) => { const handleKeyDown = (event: KeyboardEvent) => {
if ( if (
+34 -5
View File
@@ -168,6 +168,9 @@ export class ThreadMotionCoordinator {
private measurementFrameId: number | null = null; private measurementFrameId: number | null = null;
private geometryDirty = false; private geometryDirty = false;
private composerInputDuringTurn = false; private composerInputDuringTurn = false;
// A user leaving the live tail must first move beyond the near-bottom
// boundary, or explicitly reverse toward latest, before follow can resume.
private resumeFollowArmed = false;
constructor(options: ThreadMotionCoordinatorOptions) { constructor(options: ThreadMotionCoordinatorOptions) {
this.camera = options.camera; this.camera = options.camera;
@@ -198,6 +201,7 @@ export class ThreadMotionCoordinator {
if (isNewTurn) { if (isNewTurn) {
this.camera.cancel(); this.camera.cancel();
this.composerInputDuringTurn = false; this.composerInputDuringTurn = false;
this.resumeFollowArmed = false;
this.promptPositioned = turn.entry === "restored"; this.promptPositioned = turn.entry === "restored";
this.mode = this.promptPositioned && turn.hasOutput this.mode = this.promptPositioned && turn.hasOutput
? "follow-output" ? "follow-output"
@@ -249,15 +253,31 @@ export class ThreadMotionCoordinator {
this.handleUserScrollIntent(true); this.handleUserScrollIntent(true);
} }
handleUserScrollIntent(canScroll: boolean): void { handleUserScrollIntent(canScroll: boolean, towardLatest = false): void {
if (this.mode === "browsing-history" && towardLatest && !canScroll) {
this.transitionToAutoFollow(false);
return;
}
const event = canScroll ? "user-scroll" : "boundary-scroll"; const event = canScroll ? "user-scroll" : "boundary-scroll";
if (!this.transition(event)) return; const transitioned = this.transition(event);
if (this.mode === "browsing-history" && canScroll) {
this.resumeFollowArmed = towardLatest;
} else if (transitioned && this.mode === "browsing-history") {
this.resumeFollowArmed = false;
}
if (!transitioned) return;
this.camera.cancel(); this.camera.cancel();
} }
resumeAutoFollow(): void { resumeAutoFollow(): void {
this.transitionToAutoFollow(true);
}
private transitionToAutoFollow(cancelCamera: boolean): void {
if (!this.transition("resume-follow")) return; if (!this.transition("resume-follow")) return;
this.camera.cancel(); this.resumeFollowArmed = false;
if (cancelCamera) this.camera.cancel();
this.onAutoFollow?.();
this.invalidateGeometry(); this.invalidateGeometry();
} }
@@ -317,11 +337,19 @@ export class ThreadMotionCoordinator {
case "navigating-history": case "navigating-history":
if (!this.camera.isFollowing()) { if (!this.camera.isFollowing()) {
this.transition("navigation-settled"); this.transition("navigation-settled");
if (nearBottom) this.resumeAutoFollow(); if (nearBottom) {
this.resumeAutoFollow();
} else {
this.resumeFollowArmed = true;
}
} }
return "navigation"; return "navigation";
case "browsing-history": case "browsing-history":
if (!nearBottom) return "user"; if (!nearBottom) {
this.resumeFollowArmed = true;
return "user";
}
if (!this.resumeFollowArmed) return "user";
this.resumeAutoFollow(); this.resumeAutoFollow();
return "automatic"; return "automatic";
default: default:
@@ -339,6 +367,7 @@ export class ThreadMotionCoordinator {
this.camera.cancel(); this.camera.cancel();
this.turn = { id: null, promptId: null, hasOutput: false }; this.turn = { id: null, promptId: null, hasOutput: false };
this.composerInputDuringTurn = false; this.composerInputDuringTurn = false;
this.resumeFollowArmed = false;
this.mode = "idle"; this.mode = "idle";
this.promptPositioned = false; this.promptPositioned = false;
} }
+54
View File
@@ -410,6 +410,9 @@ describe("ThreadMotionCoordinator", () => {
expect(camera.jumpTo).toHaveBeenCalledWith(780); expect(camera.jumpTo).toHaveBeenCalledWith(780);
coordinator.takeUserControl(); coordinator.takeUserControl();
expect(coordinator.observeScroll(true)).toBe("user");
expect(coordinator.snapshot().mode).toBe("browsing-history");
expect(coordinator.observeScroll(false)).toBe("user"); expect(coordinator.observeScroll(false)).toBe("user");
expect(coordinator.snapshot().mode).toBe("browsing-history"); expect(coordinator.snapshot().mode).toBe("browsing-history");
@@ -417,6 +420,57 @@ describe("ThreadMotionCoordinator", () => {
expect(coordinator.snapshot().mode).toBe("anchor-prompt"); expect(coordinator.snapshot().mode).toBe("anchor-prompt");
}); });
it("resumes shallow history browsing when user intent turns toward latest", () => {
const {
camera,
coordinator,
advanceFrame,
} = motionHarness({
scrollTop: 1_400,
});
coordinator.updateTurn({
id: "turn-1",
promptId: "prompt-1",
hasOutput: true,
});
advanceFrame();
camera.followTo.mockClear();
coordinator.handleUserScrollIntent(true);
expect(coordinator.observeScroll(true)).toBe("user");
advanceFrame();
expect(camera.followTo).not.toHaveBeenCalled();
coordinator.handleUserScrollIntent(true, true);
expect(coordinator.observeScroll(true)).toBe("automatic");
expect(coordinator.snapshot().mode).toBe("follow-output");
advanceFrame();
expect(camera.followTo).toHaveBeenCalledWith(1_400);
});
it("resumes shallow history browsing from forward intent at the boundary", () => {
const {
advanceFrame,
coordinator,
onAutoFollow,
} = motionHarness({
scrollTop: 1_400,
});
coordinator.updateTurn({
id: "turn-1",
promptId: "prompt-1",
hasOutput: true,
});
advanceFrame();
coordinator.handleUserScrollIntent(true);
expect(coordinator.observeScroll(true)).toBe("user");
coordinator.handleUserScrollIntent(false, true);
expect(coordinator.snapshot().mode).toBe("follow-output");
expect(onAutoFollow).toHaveBeenCalledTimes(1);
});
it("preserves history browsing when an active turn is cleared", () => { it("preserves history browsing when an active turn is cleared", () => {
const { const {
camera, camera,
+95
View File
@@ -763,6 +763,101 @@ describe("ThreadViewport", () => {
} }
}); });
it("keeps shallow wheel and touch scrolling user-owned until intent reverses", async () => {
const followTo = vi.spyOn(ThreadCameraController.prototype, "followTo");
const threaded: UIMessage[] = [
{ id: "u1", role: "user", content: "old question", turnId: "turn-1", createdAt: 1 },
{ id: "a1", role: "assistant", content: "old answer", turnId: "turn-1", createdAt: 2 },
{ id: "u2", role: "user", content: "new question", turnId: "turn-2", createdAt: 3 },
];
const answer: UIMessage = {
id: "a2",
role: "assistant",
content: "streaming answer",
turnId: "turn-2",
isStreaming: true,
createdAt: 4,
};
const { container, rerender } = render(
<ThreadViewport
messages={threaded}
isStreaming
composer={<div>composer</div>}
/>,
);
const scroller = getScroller(container);
Object.defineProperties(scroller, {
scrollHeight: { configurable: true, value: 1_904 },
clientHeight: { configurable: true, value: 500 },
scrollTop: { configurable: true, writable: true, value: 1_404 },
});
const prompt = container.querySelector<HTMLElement>('[data-user-prompt-id="u2"]');
expect(prompt).not.toBeNull();
Object.defineProperty(prompt, "offsetTop", {
configurable: true,
value: 1_420,
});
rerender(
<ThreadViewport
messages={[...threaded, answer]}
isStreaming
composer={<div>composer</div>}
activeTurnId="turn-2"
activeTurnStartedHere
/>,
);
await flushAnimationFrame();
followTo.mockClear();
act(() => {
fireEvent.wheel(scroller, { deltaY: -24 });
scroller.scrollTop = 1_380;
scroller.dispatchEvent(new Event("scroll"));
});
await flushAnimationFrame();
expect(followTo).not.toHaveBeenCalled();
expect(scroller.scrollTop).toBe(1_380);
expect(screen.getByRole("button", { name: "Scroll to bottom" })).toBeInTheDocument();
act(() => {
scroller.scrollTop = 1_404;
scroller.dispatchEvent(new Event("scroll"));
fireEvent.wheel(scroller, { deltaY: 24 });
});
await flushAnimationFrame();
expect(followTo).toHaveBeenCalledWith(1_404);
expect(scroller.scrollTop).toBe(1_404);
expect(screen.queryByRole("button", { name: "Scroll to bottom" }))
.not.toBeInTheDocument();
followTo.mockClear();
act(() => {
fireEvent.touchStart(scroller, { touches: [{ clientY: 300 }] });
fireEvent.touchMove(scroller, { touches: [{ clientY: 324 }] });
scroller.scrollTop = 1_380;
scroller.dispatchEvent(new Event("scroll"));
});
await flushAnimationFrame();
expect(followTo).not.toHaveBeenCalled();
expect(screen.getByRole("button", { name: "Scroll to bottom" })).toBeInTheDocument();
act(() => {
fireEvent.touchMove(scroller, { touches: [{ clientY: 300 }] });
scroller.scrollTop = 1_404;
scroller.dispatchEvent(new Event("scroll"));
fireEvent.touchEnd(scroller);
});
await flushAnimationFrame();
expect(followTo).toHaveBeenCalledWith(1_404);
expect(screen.queryByRole("button", { name: "Scroll to bottom" }))
.not.toBeInTheDocument();
});
it("keeps the scroll-to-bottom button above a growing composer", async () => { it("keeps the scroll-to-bottom button above a growing composer", async () => {
const resizeObserver = stubResizeObserver(); const resizeObserver = stubResizeObserver();