From 03854b5d33d8b66fec4c86a77714e0e7512a31bd Mon Sep 17 00:00:00 2001 From: dadachann <185672915+dadachann@users.noreply.github.com> Date: Tue, 15 Sep 2026 13:23:12 +0000 Subject: [PATCH] fix(telemetry): require acknowledged adapter evidence and prepare beta.3 --- .github/workflows/build-docker-image.yml | 9 +- pyproject.toml | 2 +- .../pkg/platform/adapters/telegram/adapter.py | 13 +- .../platform/adapters/telegram/api_impl.py | 21 +- .../pkg/telemetry/adapter_diagnostics.py | 182 ++++++++++ src/langbot/pkg/telemetry/diagnostics.py | 25 +- .../core/test_release_docker_tags.py | 52 +++ .../telemetry/test_adapter_confirmation.py | 318 ++++++++++++++++++ uv.lock | 2 +- 9 files changed, 599 insertions(+), 25 deletions(-) create mode 100644 tests/unit_tests/core/test_release_docker_tags.py create mode 100644 tests/unit_tests/telemetry/test_adapter_confirmation.py diff --git a/.github/workflows/build-docker-image.yml b/.github/workflows/build-docker-image.yml index 9ef4506b8..215783735 100644 --- a/.github/workflows/build-docker-image.yml +++ b/.github/workflows/build-docker-image.yml @@ -44,4 +44,11 @@ jobs: run: docker buildx build --build-arg LANGBOT_BUILD_REVISION=${{ github.sha }} --platform linux/arm64,linux/amd64 -t rockchin/langbot:${{ steps.check_version.outputs.version }} -t rockchin/langbot:latest . --push - name: Build for Pre-release # no update for latest tag if: ${{ github.event.release.prerelease == true }} - run: docker buildx build --build-arg LANGBOT_BUILD_REVISION=${{ github.sha }} --platform linux/arm64,linux/amd64 -t rockchin/langbot:${{ steps.check_version.outputs.version }} . --push \ No newline at end of file + env: + RELEASE_VERSION: ${{ steps.check_version.outputs.version }} + run: | + tags=(-t "rockchin/langbot:$RELEASE_VERSION") + if [[ "$RELEASE_VERSION" =~ ^v[0-9]+\.[0-9]+\.[0-9]+-beta\.[0-9]+$ ]]; then + tags+=(-t rockchin/langbot:beta) + fi + docker buildx build --build-arg LANGBOT_BUILD_REVISION=${{ github.sha }} --platform linux/arm64,linux/amd64 "${tags[@]}" . --push \ No newline at end of file diff --git a/pyproject.toml b/pyproject.toml index 73c15cae4..6e99165cf 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "langbot" -version = "4.11.0-beta.2" +version = "4.11.0-beta.3" description = "Production-grade platform for building agentic IM bots" readme = "README.md" license-files = ["LICENSE"] diff --git a/src/langbot/pkg/platform/adapters/telegram/adapter.py b/src/langbot/pkg/platform/adapters/telegram/adapter.py index 8d1d8541b..0e4ae9c08 100644 --- a/src/langbot/pkg/platform/adapters/telegram/adapter.py +++ b/src/langbot/pkg/platform/adapters/telegram/adapter.py @@ -7,6 +7,7 @@ Preserves all existing functionality (messaging, streaming output, markdown card from __future__ import annotations from langbot.pkg.telemetry import diagnostics +from langbot.pkg.telemetry.adapter_diagnostics import record_api_result import typing import traceback @@ -238,20 +239,20 @@ class TelegramAdapter(TelegramAPIMixin, abstract_platform_adapter.AbstractPlatfo text = telegramify_markdown.markdownify(content=text) args['parse_mode'] = 'MarkdownV2' args['text'] = text - await self.bot.send_message(**args) + record_api_result(await self.bot.send_message(**args)) elif component_type == 'photo': photo = component.get('photo') if photo is None: continue args['photo'] = telegram.InputFile(photo) - await self.bot.send_photo(**args) + record_api_result(await self.bot.send_photo(**args)) elif component_type == 'document': doc = component.get('document') if doc is None: continue filename = component.get('filename', 'file') args['document'] = telegram.InputFile(doc, filename=filename) - await self.bot.send_document(**args) + record_api_result(await self.bot.send_document(**args)) @diagnostics.observe('api', 'reply_message', source='platform', stage='accepted') async def reply_message( @@ -285,20 +286,20 @@ class TelegramAdapter(TelegramAPIMixin, abstract_platform_adapter.AbstractPlatfo if self.config['markdown_card'] is True: args['parse_mode'] = 'MarkdownV2' args['text'] = content - await self.bot.send_message(**args) + record_api_result(await self.bot.send_message(**args)) elif component_type == 'photo': photo = component.get('photo') if photo is None: continue args['photo'] = telegram.InputFile(photo) - await self.bot.send_photo(**args) + record_api_result(await self.bot.send_photo(**args)) elif component_type == 'document': doc = component.get('document') if doc is None: continue filename = component.get('filename', 'file') args['document'] = telegram.InputFile(doc, filename=filename) - await self.bot.send_document(**args) + record_api_result(await self.bot.send_document(**args)) # ---- Streaming Output (preserving original logic) ---- diff --git a/src/langbot/pkg/platform/adapters/telegram/api_impl.py b/src/langbot/pkg/platform/adapters/telegram/api_impl.py index f9aa7de40..e6e053485 100644 --- a/src/langbot/pkg/platform/adapters/telegram/api_impl.py +++ b/src/langbot/pkg/platform/adapters/telegram/api_impl.py @@ -6,6 +6,7 @@ Implements optional API methods defined in AbstractPlatformAdapter. from __future__ import annotations from langbot.pkg.telemetry import diagnostics +from langbot.pkg.telemetry.adapter_diagnostics import record_api_result import typing @@ -52,7 +53,7 @@ class TelegramAPIMixin: } if self.config.get('markdown_card', False): args['parse_mode'] = 'MarkdownV2' - await self.bot.edit_message_text(**args) + record_api_result(await self.bot.edit_message_text(**args), edited_content=new_content) return @diagnostics.observe('api', 'delete_message', source='platform', stage='accepted') @@ -63,7 +64,7 @@ class TelegramAPIMixin: message_id: typing.Union[int, str], ) -> None: """Delete / recall a message.""" - await self.bot.delete_message(chat_id=chat_id, message_id=message_id) + record_api_result(await self.bot.delete_message(chat_id=chat_id, message_id=message_id)) @diagnostics.observe('api', 'forward_message', source='platform', stage='accepted') async def forward_message( @@ -223,7 +224,7 @@ class TelegramAPIMixin: } if duration > 0: kwargs['until_date'] = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(seconds=duration) - await self.bot.restrict_chat_member(**kwargs) + record_api_result(await self.bot.restrict_chat_member(**kwargs)) @diagnostics.observe('api', 'unmute_member', source='platform', stage='accepted') async def unmute_member( @@ -243,10 +244,12 @@ class TelegramAPIMixin: can_send_video_notes=True, can_send_voice_notes=True, ) - await self.bot.restrict_chat_member( - chat_id=group_id, - user_id=user_id, - permissions=permissions, + record_api_result( + await self.bot.restrict_chat_member( + chat_id=group_id, + user_id=user_id, + permissions=permissions, + ) ) @diagnostics.observe('api', 'kick_member', source='platform', stage='accepted') @@ -256,7 +259,7 @@ class TelegramAPIMixin: user_id: typing.Union[int, str], ) -> None: """Kick a member from the group.""" - await self.bot.ban_chat_member(chat_id=group_id, user_id=user_id) + record_api_result(await self.bot.ban_chat_member(chat_id=group_id, user_id=user_id)) @diagnostics.observe('api', 'leave_group', source='platform', stage='accepted') async def leave_group( @@ -264,4 +267,4 @@ class TelegramAPIMixin: group_id: typing.Union[int, str], ) -> None: """Make the bot leave a group.""" - await self.bot.leave_chat(chat_id=group_id) + record_api_result(await self.bot.leave_chat(chat_id=group_id)) diff --git a/src/langbot/pkg/telemetry/adapter_diagnostics.py b/src/langbot/pkg/telemetry/adapter_diagnostics.py index a62ab4e72..301dd9113 100644 --- a/src/langbot/pkg/telemetry/adapter_diagnostics.py +++ b/src/langbot/pkg/telemetry/adapter_diagnostics.py @@ -23,6 +23,18 @@ def boundary_fields(module, kind, operation, bound, parent): else: # Keep the generic boundary for investigation, not a false named test. evidence = False + if ( + entry['adapter'] == 'lark-omni' + and resolved + in { + 'platform_api.check_tenant_access_token', + 'platform_api.refresh_app_access_token', + 'platform_api.refresh_tenant_access_token', + } + and getattr(bound.get('self'), 'config', {}).get('app_type', 'self') != 'isv' + ): + # Self-built apps return {'ok': True} without checking/refreshing a token. + evidence = False # A forwarding API may call another decorated method. Count the outer call. if parent and parent.adapter_api_active: evidence = False @@ -30,6 +42,8 @@ def boundary_fields(module, kind, operation, bound, parent): return {'_adapter_api_active': kind == 'api'} return { '_adapter_api_active': kind == 'api', + '_adapter_void_ack': completed_void_api(module, resolved), + '_adapter_read_ack': completed_read_api(module, resolved), 'operation': resolved, 'attributes': {'adapter_evidence': True, **message_scenario(bound)}, } @@ -68,3 +82,171 @@ def message_scenario(bound): ) result['content_type'] = next(iter(types)) if len(types) == 1 else 'mixed' if types else 'unknown' return result + + +# These exact API implementations always await an exception-raising SDK call, +# including SDKs whose successful acknowledgement has no return payload. +# Do not extend this to methods with conditional no-op or queued work. +VOID_ACKNOWLEDGED_APIS = { + 'aiocqhttp': frozenset( + { + 'delete_message', + 'set_group_name', + 'mute_member', + 'unmute_member', + 'kick_member', + 'leave_group', + 'approve_friend_request', + 'approve_group_invite', + } + ), + 'discord': frozenset( + { + 'edit_message', + 'delete_message', + 'mute_member', + 'unmute_member', + 'kick_member', + 'leave_group', + } + ), + 'kook': frozenset({'delete_message'}), +} + + +def completed_void_api(module, operation): + return module.endswith('.api_impl') and any( + directory in module.split('.') and operation in operations + for directory, operations in VOID_ACKNOWLEDGED_APIS.items() + ) + + +# These read methods return actual fetched/cached data or raise. Exclude the +# known identity-only placeholders and file-ID passthroughs: a typed result alone +# is not evidence that a lookup worked. Unsupported methods still raise normally. +_READ_APIS = frozenset( + { + 'get_message', + 'get_group_info', + 'get_group_list', + 'get_group_member_list', + 'get_group_member_info', + 'get_user_info', + 'get_friend_list', + 'get_file_url', + } +) +READ_ACKNOWLEDGED_APIS = { + 'aiocqhttp': _READ_APIS, + 'telegram': _READ_APIS, + 'discord': _READ_APIS - {'get_file_url'}, + 'dingtalk': _READ_APIS - {'get_group_info', 'get_user_info'}, + 'kook': _READ_APIS - {'get_file_url'}, + 'lark': _READ_APIS - {'get_user_info', 'get_group_member_info', 'get_file_url'}, + 'officialaccount': _READ_APIS, + 'qqofficial': _READ_APIS, + 'slack': _READ_APIS, + 'wecom': _READ_APIS - {'get_user_info'}, + 'wecombot': _READ_APIS, + 'wecomcs': _READ_APIS, +} + + +def completed_read_api(module, operation): + return module.endswith('.api_impl') and any( + directory in module.split('.') and operation in operations + for directory, operations in READ_ACKNOWLEDGED_APIS.items() + ) + + +def response_outcome(value, adapter=None): + """Finite acknowledgement contracts; never inspect or export message content. + + A MessageResult may contain an *inbound* message ID even when nothing was + sent. Only its actual raw acknowledgement is eligible. Empty/queued/unknown + returns are unconfirmed, not failures. Unknown SDK shapes need an explicit + confirmation at their operation boundary, not a truthiness fallback. + """ + from langbot_plugin.api.entities.builtin.platform.events import MessageResult + + message_result = isinstance(value, MessageResult) + if message_result: + value = value.raw + if not isinstance(value, dict): + return 'skipped' + if ( + value.get('ok') is False + or value.get('status') == 'failed' + or any(type(value.get(key)) is int and value[key] != 0 for key in ('retcode', 'errcode', 'code')) + ): + return 'failed' + if value.get('queued') is True or value.get('status') == 'async': + return 'skipped' + # The wrappers that batch sends retain the real responses in these fields. + if 'results' in value: + results = value['results'] + if not isinstance(results, list) or not results: + return 'skipped' + outcomes = [response_outcome(item, adapter) for item in results] + return 'failed' if 'failed' in outcomes else 'skipped' if 'skipped' in outcomes else 'succeeded' + for key in ('result', 'raw'): + if key in value: + return response_outcome(value[key], adapter) + if ( + value.get('ok') is True + or value.get('status') == 'ok' + or any(type(value.get(key)) is int and value[key] == 0 for key in ('retcode', 'errcode', 'code')) + ): + return 'succeeded' + if message_result: + # These wrappers put only SDK-returned IDs in raw (never the source ID). + if adapter in {'aiocqhttp-omni', 'discord-omni', 'telegram-omni'}: + message_id = value.get('message_id') + if type(message_id) in (str, int) and message_id: + return 'succeeded' + if adapter == 'lark-omni': + ids = value.get('message_ids') + if isinstance(ids, list) and ids and all(isinstance(item, str) and item for item in ids): + return 'succeeded' + return 'skipped' + + +def record_api_result(value, *, edited_content=None): + """Record Telegram's Message/True acknowledgement and return it unchanged. + + This runs after the SDK await, not after conversion or local stream setup. + No payload is retained; exceptions and the adapter's public return stay intact. + """ + try: + from .diagnostics import current_span + + span = current_span() + if span is not None and span.kind == 'api' and span.fields.get('attributes', {}).get('adapter_evidence'): + from telegram import Message + + outcome = ( + 'succeeded' + if value is True or isinstance(value, Message) + else 'failed' + if value is False + else 'skipped' + ) + if outcome == 'succeeded' and edited_content is not None: + from langbot_plugin.api.entities.builtin.platform.message import Plain + + if any(not isinstance(item, Plain) for item in edited_content): + outcome = 'partial' + # One unacknowledged component must not be hidden by a later success. + prior = span.adapter_api_result + span.adapter_api_result = ( + 'failed' + if 'failed' in (prior, outcome) + else 'skipped' + if 'skipped' in (prior, outcome) + else 'partial' + if 'partial' in (prior, outcome) + else outcome + ) + except Exception: + pass + return value diff --git a/src/langbot/pkg/telemetry/diagnostics.py b/src/langbot/pkg/telemetry/diagnostics.py index 2636528ec..871efc2f8 100644 --- a/src/langbot/pkg/telemetry/diagnostics.py +++ b/src/langbot/pkg/telemetry/diagnostics.py @@ -186,6 +186,9 @@ class Span: self.operation = operation self.fields = dict(fields) self.adapter_api_active = bool(self.fields.pop('_adapter_api_active', False)) + self.adapter_void_ack = bool(self.fields.pop('_adapter_void_ack', False)) + self.adapter_read_ack = bool(self.fields.pop('_adapter_read_ack', False)) + self.adapter_api_result = None if ( parent and parent.fields.get('workspace_uuid') @@ -264,13 +267,21 @@ def result_outcome(value): if not isinstance(value, EBAEvent) and span.fields.get('attributes', {}).get('adapter_evidence'): span.fields['attributes']['adapter_evidence'] = False if span is not None and span.kind == 'api' and span.fields.get('attributes', {}).get('adapter_evidence'): - # Common adapter response contracts expose status without inspecting content. - if isinstance(value, dict) and ( - value.get('ok') is False - or value.get('status') == 'failed' - or (type(value.get('retcode')) is int and value['retcode'] != 0) - ): - set_outcome('failed', reason_code='response_error') + from .adapter_diagnostics import response_outcome + + # A normal return is not an acknowledgement. Fail closed for acceptance + # (but never for the business call) even if result inspection raises. + if span.outcome is None: + set_outcome('skipped') + outcome = span.adapter_api_result + if outcome is None: + if (value is None and span.adapter_void_ack) or ( + span.adapter_read_ack and value is not None and value != '' + ): + outcome = 'succeeded' + else: + outcome = response_outcome(value, span.fields.get('adapter')) + set_outcome(outcome, **({'reason_code': 'response_error'} if outcome == 'failed' else {})) if isinstance(value, ActionResponse): if value.code != 0: set_outcome('failed', reason_code='response_error') diff --git a/tests/unit_tests/core/test_release_docker_tags.py b/tests/unit_tests/core/test_release_docker_tags.py new file mode 100644 index 000000000..72180b56a --- /dev/null +++ b/tests/unit_tests/core/test_release_docker_tags.py @@ -0,0 +1,52 @@ +"""Exercise release tag selection without building or pushing an image.""" + +import os +from pathlib import Path +import subprocess + +import pytest +import yaml + + +ROOT = Path(__file__).resolve().parents[3] + + +def _build_args(tmp_path, step_name, version): + workflow = yaml.safe_load((ROOT / '.github/workflows/build-docker-image.yml').read_text()) + steps = workflow['jobs']['publish-docker-image']['steps'] + step = next(step for step in steps if step['name'] == step_name) + command = step['run'].replace('${{ github.sha }}', 'a' * 40) + command = command.replace('${{ steps.check_version.outputs.version }}', version) + capture = tmp_path / 'docker-args' + script = 'docker() { printf "%s\\n" "$@" > "$CAPTURE"; };\n' + command + if 'env' in step: + assert step['env']['RELEASE_VERSION'] == '${{ steps.check_version.outputs.version }}' + subprocess.run( + ['bash', '-euo', 'pipefail', '-c', script], + check=True, + env={**os.environ, 'CAPTURE': str(capture), 'RELEASE_VERSION': version}, + ) + args = capture.read_text().splitlines() + assert '--push' in args + assert args[args.index('--platform') + 1] == 'linux/arm64,linux/amd64' + assert args[args.index('--build-arg') + 1] == 'LANGBOT_BUILD_REVISION=' + 'a' * 40 + return [args[i + 1] for i, value in enumerate(args) if value == '-t'], step['if'] + + +@pytest.mark.parametrize('version', ['v4.11.0-beta.3', 'v4.12.0-beta.1']) +def test_beta_release_updates_beta_but_not_latest(tmp_path, version): + tags, condition = _build_args(tmp_path, 'Build for Pre-release', version) + assert tags == [f'rockchin/langbot:{version}', 'rockchin/langbot:beta'] + assert condition == '${{ github.event.release.prerelease == true }}' + + +@pytest.mark.parametrize('version', ['v4.11.0-rc.1', 'v4.12.0-alpha.1']) +def test_other_prereleases_do_not_update_channels(tmp_path, version): + tags, _ = _build_args(tmp_path, 'Build for Pre-release', version) + assert tags == [f'rockchin/langbot:{version}'] + + +def test_stable_release_keeps_latest_without_updating_beta(tmp_path): + tags, condition = _build_args(tmp_path, 'Build for Release', 'v4.11.0') + assert tags == ['rockchin/langbot:v4.11.0', 'rockchin/langbot:latest'] + assert condition == '${{ github.event.release.prerelease == false }}' diff --git a/tests/unit_tests/telemetry/test_adapter_confirmation.py b/tests/unit_tests/telemetry/test_adapter_confirmation.py new file mode 100644 index 000000000..48307db67 --- /dev/null +++ b/tests/unit_tests/telemetry/test_adapter_confirmation.py @@ -0,0 +1,318 @@ +"""Real adapter acceptance must reflect confirmed operations, not normal returns.""" + +import base64 +import json +from types import SimpleNamespace +from unittest.mock import AsyncMock +from uuid import uuid4 + +import httpx +import pytest +from langbot_plugin.api.entities.builtin.platform.message import File, MessageChain, Plain + +from langbot.pkg.api.http.context import ExecutionContext +from langbot.pkg.platform.adapters.telegram.adapter import TelegramAdapter +from langbot.pkg.telemetry import diagnostics as d + + +def make_adapter(version='4.11.0b3', **config): + ap = SimpleNamespace(instance_config=SimpleNamespace(data={'space': {'url': 'https://example.invalid', **config}})) + manager = d.DiagnosticsManager(ap, version=version, instance_id='instance-test') + ap.diagnostics = manager + context = ExecutionContext(instance_uuid=manager.instance_id, workspace_uuid=str(uuid4()), placement_generation=1) + sdk = SimpleNamespace(edit_message_text=AsyncMock(return_value=True), send_message=AsyncMock(return_value=True)) + adapter = TelegramAdapter.model_construct( + bot=sdk, config={'markdown_card': False}, logger=SimpleNamespace(ap=ap, execution_context=context), listeners={} + ) + return adapter, manager, sdk + + +async def wire_events(manager): + batches = [] + + async def sender(request): + batch = json.loads(request.content) + batches.append(batch) + return httpx.Response( + 200, + json={ + 'code': 200, + 'data': {'accepted_event_ids': [e['event_id'] for e in batch['events']], 'rejected': []}, + }, + ) + + manager.client = httpx.AsyncClient(transport=httpx.MockTransport(sender)) + await manager.flush_once() + assert not manager.pending + await manager.shutdown(drain_timeout=0) + assert 'PRIVATE_' not in json.dumps(batches) + return [event for batch in batches for event in batch['events']] + + +def successes(events): + return [e for e in events if e['attributes'].get('adapter_evidence') and e['outcome'] == 'succeeded'] + + +@pytest.mark.asyncio +async def test_real_telegram_file_only_edit_is_not_acceptance_through_sender(): + adapter, manager, sdk = make_adapter() + for _ in range(2): + result = await adapter.edit_message( + 'person', + 'PRIVATE_CHAT', + 'PRIVATE_MESSAGE', + MessageChain([File(name='PRIVATE_FILE', base64=base64.b64encode(b'PRIVATE_BYTES').decode())]), + ) + assert result is None + sdk.edit_message_text.assert_not_awaited() + events = await wire_events(manager) + assert len([e for e in events if e['operation'] == 'edit_message']) == 4 + assert not successes(events) + + +@pytest.mark.asyncio +async def test_real_telegram_acknowledged_void_edit_is_acceptance_through_sender(): + adapter, manager, sdk = make_adapter() + assert ( + await adapter.edit_message( + 'person', 'PRIVATE_CHAT', 'PRIVATE_MESSAGE', MessageChain([Plain(text='PRIVATE_TEXT')]) + ) + is None + ) + sdk.edit_message_text.assert_awaited_once_with( + chat_id='PRIVATE_CHAT', message_id='PRIVATE_MESSAGE', text='PRIVATE_TEXT' + ) + events = await wire_events(manager) + assert len(successes(events)) == 1 + assert successes(events)[0]['attributes']['content_type'] == 'text' + + +@pytest.mark.asyncio +async def test_real_telegram_empty_send_is_not_acceptance(): + adapter, manager, sdk = make_adapter() + assert await adapter.send_message('person', 'PRIVATE_CHAT', MessageChain([])) is None + sdk.send_message.assert_not_awaited() + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +async def test_real_telegram_edit_exception_is_unchanged(): + adapter, manager, sdk = make_adapter() + error = ValueError('PRIVATE_ERROR') + sdk.edit_message_text.side_effect = error + with pytest.raises(ValueError) as raised: + await adapter.edit_message( + 'person', 'PRIVATE_CHAT', 'PRIVATE_MESSAGE', MessageChain([Plain(text='PRIVATE_TEXT')]) + ) + assert raised.value is error + events = await wire_events(manager) + assert not successes(events) + assert any(e['outcome'] == 'failed' and e['operation'] == 'edit_message' for e in events) + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + 'version,config', + [('4.11.0', {}), ('4.11.0b3', {'disable_telemetry': True}), ('4.11.0b3', {'disable_beta_diagnostics': True})], +) +async def test_real_adapter_disabled_gates_keep_business_behavior(version, config): + adapter, manager, sdk = make_adapter(version, **config) + assert ( + await adapter.edit_message( + 'person', 'PRIVATE_CHAT', 'PRIVATE_MESSAGE', MessageChain([Plain(text='PRIVATE_TEXT')]) + ) + is None + ) + sdk.edit_message_text.assert_awaited_once() + assert await wire_events(manager) == [] + + +@pytest.mark.asyncio +@pytest.mark.parametrize('result', [None, False, {}, {'queued': True}, {'stream': False}]) +async def test_unconfirmed_normal_return_is_not_acceptance(result): + adapter, manager, _ = make_adapter() + + async def call(self): + return result + + call.__module__ = 'langbot.pkg.platform.adapters.telegram.adapter' + observed = d.observe('api', 'send_message', source='platform', stage='accepted')(call) + assert await observed(adapter) is result + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +@pytest.mark.parametrize('raw', [{}, {'results': []}, {'result': None}, {'queued': True}]) +async def test_message_result_source_id_is_not_confirmation(raw): + from langbot_plugin.api.entities.builtin.platform.events import MessageResult + + adapter, manager, _ = make_adapter() + result = MessageResult(message_id='PRIVATE_SOURCE_ID', raw=raw) + + async def call(self): + return result + + call.__module__ = 'langbot.pkg.platform.adapters.wecom.adapter' + observed = d.observe('api', 'reply_message', source='platform', stage='accepted')(call) + assert await observed(adapter) is result + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +async def test_real_discord_void_delete_is_confirmed(): + from langbot.pkg.platform.adapters.discord.adapter import DiscordAdapter + + adapter, manager, _ = make_adapter() + message = SimpleNamespace(delete=AsyncMock(return_value=None)) + channel = SimpleNamespace(fetch_message=AsyncMock(return_value=message)) + discord = DiscordAdapter.model_construct( + bot=SimpleNamespace(get_channel=lambda _: channel), logger=adapter.logger, config={}, listeners={} + ) + assert await discord.delete_message('group', '123', '456') is None + message.delete.assert_awaited_once() + assert len(successes(await wire_events(manager))) == 1 + + +@pytest.mark.asyncio +@pytest.mark.parametrize('result', [False, None]) +async def test_real_telegram_unacknowledged_edit_keeps_void_return(result): + adapter, manager, sdk = make_adapter() + sdk.edit_message_text.return_value = result + assert ( + await adapter.edit_message( + 'person', 'PRIVATE_CHAT', 'PRIVATE_MESSAGE', MessageChain([Plain(text='PRIVATE_TEXT')]) + ) + is None + ) + sdk.edit_message_text.assert_awaited_once() + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +async def test_real_telegram_acknowledged_void_delete_is_acceptance(): + adapter, manager, sdk = make_adapter() + sdk.delete_message = AsyncMock(return_value=True) + assert await adapter.delete_message('person', 'PRIVATE_CHAT', 'PRIVATE_MESSAGE') is None + sdk.delete_message.assert_awaited_once() + assert len(successes(await wire_events(manager))) == 1 + + +@pytest.mark.asyncio +async def test_real_telegram_sdk_message_send_is_acceptance(): + import datetime + import telegram + + adapter, manager, sdk = make_adapter() + sdk.send_message.return_value = telegram.Message( + message_id=123, date=datetime.datetime.now(datetime.timezone.utc), chat=telegram.Chat(id=456, type='private') + ) + assert await adapter.send_message('person', 'PRIVATE_CHAT', MessageChain([Plain(text='PRIVATE_TEXT')])) is None + sdk.send_message.assert_awaited_once() + assert len(successes(await wire_events(manager))) == 1 + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + 'raw', [{'ok': True, 'raw': None}, {'ok': True, 'raw': {'errcode': 42}}, {'results': [{'ok': True}, None]}] +) +async def test_wrapped_ack_does_not_hide_missing_or_failed_response(raw): + adapter, manager, _ = make_adapter() + + async def call(self): + return raw + + call.__module__ = 'langbot.pkg.platform.adapters.wecombot.adapter' + observed = d.observe('api', 'reply_message', source='platform', stage='accepted')(call) + assert await observed(adapter) is raw + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + 'action', ['check_tenant_access_token', 'refresh_app_access_token', 'refresh_tenant_access_token'] +) +async def test_real_lark_self_app_token_noop_is_not_acceptance(action): + from langbot.pkg.platform.adapters.lark.adapter import LarkAdapter + + adapter, manager, _ = make_adapter() + lark = LarkAdapter.model_construct(config={'app_type': 'self'}, logger=adapter.logger, listeners={}) + assert await lark.call_platform_api(action, {}) == {'ok': True} + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +async def test_real_telegram_mixed_edit_cannot_confirm_ignored_file(): + adapter, manager, sdk = make_adapter() + message = MessageChain([Plain(text='PRIVATE_TEXT'), File(name='PRIVATE_FILE', base64='eA==')]) + assert await adapter.edit_message('person', 'PRIVATE_CHAT', 'PRIVATE_MESSAGE', message) is None + sdk.edit_message_text.assert_awaited_once() + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +async def test_real_discord_sent_message_result_is_confirmed(): + from langbot.pkg.platform.adapters.discord.adapter import DiscordAdapter + + adapter, manager, _ = make_adapter() + channel = SimpleNamespace(send=AsyncMock(return_value=SimpleNamespace(id=123))) + discord = DiscordAdapter.model_construct( + bot=SimpleNamespace(get_channel=lambda _: channel), logger=adapter.logger, config={}, listeners={} + ) + result = await discord.send_message('group', '456', MessageChain([Plain(text='PRIVATE_TEXT')])) + assert result.message_id == 123 + channel.send.assert_awaited_once_with(content='PRIVATE_TEXT') + assert len(successes(await wire_events(manager))) == 1 + + +@pytest.mark.asyncio +async def test_real_aiocqhttp_empty_forward_is_not_acceptance(): + from langbot.pkg.platform.adapters.aiocqhttp.adapter import AiocqhttpAdapter + from langbot_plugin.api.entities.builtin.platform.message import Forward + + adapter, manager, _ = make_adapter() + sdk = SimpleNamespace(call_action=AsyncMock()) + onebot = AiocqhttpAdapter.model_construct(bot=sdk, logger=adapter.logger, config={}, listeners={}) + result = await onebot.send_message('group', '123', MessageChain([Forward(node_list=[])])) + assert result.message_id is None and result.raw == {} + sdk.call_action.assert_not_awaited() + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +async def test_real_telegram_unsupported_exception_is_unchanged(): + from langbot_plugin.api.entities.builtin.platform.errors import NotSupportedError + + adapter, manager, _ = make_adapter() + with pytest.raises(NotSupportedError): + await adapter.upload_file(b'PRIVATE_BYTES', 'PRIVATE_FILE') + assert not successes(await wire_events(manager)) + + +@pytest.mark.asyncio +async def test_real_lark_sent_message_result_is_confirmed(): + from langbot.pkg.platform.adapters.lark.adapter import LarkAdapter + + adapter, manager, _ = make_adapter() + create = AsyncMock( + return_value=SimpleNamespace(success=lambda: True, data=SimpleNamespace(message_id='PRIVATE_SENT_ID')) + ) + lark = LarkAdapter.model_construct( + config={'app_type': 'self'}, + logger=adapter.logger, + listeners={}, + api_client=SimpleNamespace(im=SimpleNamespace(v1=SimpleNamespace(message=SimpleNamespace(acreate=create)))), + ) + result = await lark.send_message('group', 'PRIVATE_CHAT', MessageChain([Plain(text='PRIVATE_TEXT')])) + assert result.message_id == 'PRIVATE_SENT_ID' + create.assert_awaited_once() + assert len(successes(await wire_events(manager))) == 1 + + +@pytest.mark.asyncio +async def test_real_telegram_returned_user_info_is_confirmed(): + adapter, manager, sdk = make_adapter() + sdk.get_chat = AsyncMock(return_value=SimpleNamespace(id=123, first_name='PRIVATE_NAME', username='PRIVATE_USER')) + result = await adapter.get_user_info('PRIVATE_USER') + assert result.id == 123 + sdk.get_chat.assert_awaited_once_with(chat_id='PRIVATE_USER') + assert len(successes(await wire_events(manager))) == 1 diff --git a/uv.lock b/uv.lock index b18ae9598..1cbfd8425 100644 --- a/uv.lock +++ b/uv.lock @@ -1999,7 +1999,7 @@ wheels = [ [[package]] name = "langbot" -version = "4.11.0b2" +version = "4.11.0b3" source = { editable = "." } dependencies = [ { name = "aiocqhttp" },