mirror of
https://github.com/langbot-app/LangBot.git
synced 2026-09-16 14:57:15 +00:00
fix(telemetry): require acknowledged adapter evidence and prepare beta.3
This commit is contained in:
@@ -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
|
||||
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
|
||||
+1
-1
@@ -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"]
|
||||
|
||||
@@ -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) ----
|
||||
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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')
|
||||
|
||||
@@ -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 }}'
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user