diff --git a/CHANGELOG.md b/CHANGELOG.md index 5e5b29b0..a53c2586 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -42,6 +42,11 @@ Sync wave from `chat@4.31.0` to `chat@4.41.1` (tracking #184). `UPSTREAM_PARITY` - **Twilio: authenticated media downloads restricted to the configured API origin** (#235; **security**, **breaking (Twilio, custom `api_url` only)**). Ports upstream `d8103a10` (vercel/chat#831, chat@4.38.1). `fetch_twilio_media` gains keyword-only `api_url` / `api_base_url` and raises `TwilioApiError("Twilio media URL must match the configured Twilio API origin", status=0)` for any URL whose scheme, host or effective port differs from `api_url` → `api_base_url` → `https://api.twilio.com`. The check runs before credentials are resolved or a request is made. The adapter passes its `api_url` to every attachment download, freshly received webhook media (`MediaUrlN`) as well as rehydrated attachments, so the existing Python-only Twilio host allowlist stays in front as defence in depth (documented in `docs/UPSTREAM_SYNC.md`). - **Consumer-visible (custom `api_url` only):** with a non-default `api_url`, media hosted on any other origin, including inbound media on `https://api.twilio.com`, is now refused (upstream behaves the same). With a non-Twilio or `http` `api_url` (a proxy or local mock), no media URL passes both layers: `api.twilio.com` fails the origin check and the proxy origin fails the host allowlist, so attachment downloads always raise. Before this change such configs downloaded `api.twilio.com` media. The default config (`api_url` unset) is unaffected. - **Teams: cap `microsoft-teams-{apps,api,cards}` at `<2.1`.** `uv.lock` is not committed and the extras were unbounded, so fresh installs resolved `microsoft-teams-apps` 2.1.0 (released 2026-09-16). 2.1.0 removed `App.activity_sender`, which the adapter uses to create native DM streams (`teams/adapter.py:943`), and changed the activities-client `update` signature that `edit_message`'s service-URL retargeting relies on. Native streaming and edits could fail on a fresh install, and CI turned red. The cap resolves to 2.0.16 until the adapter supports 2.1. +- **Teams: support `microsoft-teams-{apps,api,cards,common}` 2.1; cap lifted to `<2.2`** (#250). Supersedes the `<2.1` cap above. The adapter feature-detects the SDK line instead of sniffing versions. The Teams tests were run locally on 2.0.16 and 2.1.0; CI (`uv sync --group dev`, no lockfile) installs 2.1.x, so only 2.1.x is covered in CI. `microsoft-teams-common` is now declared (with the same `<2.2` cap) because the adapter imports it directly and `microsoft-teams-apps` leaves it unbounded. + - **Native DM streaming:** uses `app.activity_sender.create_stream(ref)` when the App has it (2.0.x). Otherwise it builds `HttpStream` on `app.api.from_service_url(ref.service_url)`, as 2.1's own `ctx.stream` does. On both lines the stream keeps its own client on the inbound service URL, so outbound calls that retarget `App.api` cannot redirect it. + - **Edit/delete:** no runtime change. `edit_message` retargeting already worked on 2.1; only the test double broke: 2.1 passes new `service_url=` / `agentic_identity=` keywords to the activities client's `update`, and the test's fake `update` did not accept them. The tests now check the request URL at the HTTP boundary. + - **Inbound auth (security):** 2.1's validator also accepts Entra ID ("Agent ID") tokens for the app id from any tenant, and picks that branch from the token's unverified issuer, fetching the JWKS of whatever tenant the token names. The adapter does not support Agent ID activities, so the bridge now reads the Bearer token's unverified `iss` before the SDK validator runs and answers 401 unless it is the cloud's Bot Framework issuer (`App.cloud.token_issuer`, so sovereign clouds keep working). An unauthenticated request therefore cannot make the bot fetch an attacker-chosen tenant's JWKS, and only Bot Framework tokens are accepted, as on 2.0.x. Tokens that pass still get the SDK's full validation, and the same issuer check runs again on the validated token in `_dispatch_activity`. On 2.0.x the only change is that such tokens are refused without a JWKS fetch. `dangerously_allow_unauthenticated_requests` mode is untouched. See `docs/UPSTREAM_SYNC.md`. + - **Consumer-visible:** fresh installs of `chat-sdk[teams]` now resolve `microsoft-teams-apps` 2.1.x. To stay on 2.0.x, pin all four SDK packages together (`microsoft-teams-apps<2.1 microsoft-teams-api<2.1 microsoft-teams-cards<2.1 microsoft-teams-common<2.1`); pinning only `microsoft-teams-apps` resolves apps 2.0.16 next to api/cards/common 2.1.0, a mix that has not been reviewed. A live Teams check of streaming and edit on 2.1 has not been done yet. - **Postgres state: expired `set_if_not_exists` claims are reclaimed; migration-owned schemas** (#240). Ports upstream `d88789c9` (vercel/chat#636, chat@4.35.0) and `ea025af7` (vercel/chat#913, chat@4.41.0). - **Fix (consumer-visible):** `PostgresStateAdapter.set_if_not_exists` used `ON CONFLICT DO NOTHING`, so an *expired* row blocked every new claim until a `get()` of that exact key happened to delete it. Dedupe keys and any lease built on `set_if_not_exists` (e.g. Telegram `update_id` claims) refused work they should accept. The query is now upstream's conditional upsert (`DO UPDATE ... WHERE expires_at IS NOT NULL AND expires_at <= now() RETURNING cache_key`), which reclaims an expired row atomically and never overwrites a live or permanent (no-TTL) one. Postgres-backed dedupe and leases now recover without cleanup. No schema change. - **New, opt-in:** keyword-only `auto_create_schema: bool = True` on `PostgresStateAdapter` and `create_postgres_state`. With `False`, `connect()` runs no DDL. It runs `SELECT 1`, then one read-only query that checks every table, the column privileges the adapter uses and the list/queue sequences. It raises `chat_sdk.StateSchemaError` naming the problem (the first table PostgreSQL cannot resolve, or every object whose grants are missing), so a wrong `search_path` or a forgotten grant fails at startup. The default is unchanged, and `auto_create_schema=None` also means `True` (upstream `autoCreateSchema ?? true`). diff --git a/docs/UPSTREAM_SYNC.md b/docs/UPSTREAM_SYNC.md index d90db929..043aef19 100644 --- a/docs/UPSTREAM_SYNC.md +++ b/docs/UPSTREAM_SYNC.md @@ -1004,7 +1004,8 @@ stay explicit instead of being rediscovered in code review. | Slack `stream()` on Enterprise Grid (`chat.startStream` `team_not_found`) | Threads the workspace `team_id` into `client.chat_stream(...)` (= `options.recipient_team_id`, the `team.id` extracted on the inbound path), which slack_sdk forwards into `_stream_args` → `chat.startStream`. `chat.appendStream`/`chat.stopStream` don't receive `_stream_args` and don't need `team_id`. Harmless on non-Grid workspaces — a correct `team_id` is always valid. (chat-sdk-python#95) | Builds the `chat.startStream` args from `channel`/`threadTs`/`recipientUserId`/`recipientTeamId`/`taskDisplayMode` only (`adapter-slack/src/index.ts` `stream()`); never passes a workspace `team_id`. On Grid orgs `chat.startStream` then fails with `team_not_found` (the per-workspace bot token alone isn't sufficient to disambiguate the team), even though `chat.postMessage` on the same workspace succeeds without it. | Upstream has the same gap — its `stream()` never threads `team_id`, so streaming is broken on Grid while non-streaming posts work. `chat.startStream` requires `team_id` for Grid disambiguation; `chat.postMessage` does not, which is why only streaming regresses. We source `team_id` from the already-plumbed `recipient_team_id` (the workspace where the interaction happened = the streaming target workspace). Live verification needs a real Grid workspace; the unit regression (`tests/test_slack_api.py::TestStream::test_stream_threads_team_id_to_chat_stream_for_grid` + the `team_not_found` mutation guard) simulates Grid by raising `team_not_found` from the streamer's lazy `chat.startStream` when `team_id` is absent. Tracked for contribution upstream. | | Fallback streaming final SentMessage content (non-Teams adapters) | SentMessage + final edit carry `final_content` (remend'd — inline markers auto-closed) | SentMessage + final edit carry raw `accumulated` | Narrow UX refinement. If a stream ends with an unclosed `*`/`~~`/etc., upstream ships the unclosed marker; we run `_remend` so the user sees a clean final message. Not observable in the common case where streams close their own markers. Teams DMs stream through the SDK `IStreamer` and the Teams accumulate-and-post path ships raw `accumulated` via `post_message`, matching upstream; this divergence applies only to the remaining adapters that still route through `_fallback_stream`. | | Teams group-chat / channel streaming via accumulate-and-post | `TeamsAdapter.stream` accumulates the full text and issues a single `post_message` (SDK-backed) instead of post+edit, even for group chats and channel threads | Same (`@chat-adapter/teams@4.30.0`: `if (activeStream && !activeStream.canceled) … else { accumulate; postMessage }`) — no divergence at the adapter level | Documented for clarity: the Python port matches upstream's behavior of avoiding the post+edit flicker where Teams doesn't support native streaming. The buffered fallback routes through the same SDK `App.send` path as a normal `post_message`. | -| Teams native streaming via the SDK `IStreamer` (DMs) | `TeamsAdapter._handle_message_activity` captures a Teams SDK `IStreamer` (`microsoft_teams.apps.StreamerProtocol` / `HttpStream`) for DMs via `app.activity_sender.create_stream(ref)`, registers it in `_active_streams`, and `await`s a `processing_done` gate (a wrapped `wait_until` shim) so the streamer stays alive while the handler streams. `stream()` → `_stream_via_emit` calls `stream.emit(text)` per chunk and NEVER calls `close()`; the adapter's `_handle_message_activity` `finally` calls `stream.close()` once (the lifecycle-owner role the SDK App's `process_activity` plays upstream). | `@chat-adapter/teams@4.30.0` `index.ts` does exactly this: `this.activeStreams.set(threadId, ctx.stream)`, build `processingDone` + wrapped `waitUntil`, `await processingDone`, `streamViaEmit` calls `stream.emit(text)` and never `close()` (the SDK App auto-closes after the handler returns). | **No adapter-level divergence.** The only mechanical difference is the close call site: upstream lets the SDK `App` auto-close `ctx.stream` because the SDK owns dispatch; our bridge overrides `server.on_request`, so we own dispatch and reproduce the close in `_handle_message_activity`'s `finally`. The SDK `HttpStream.close` no-ops when the stream was canceled or had no content, so closing in both success and cancel paths is safe (matching the SDK App, which closes in both its success and `StreamCancelledError` branches). Cancellation is detected via `stream.canceled` (checked before each emit) and by catching `StreamCancelledError` (other exceptions re-raise). The first chunk id is captured via `on_chunk` and awaited only when text was emitted and the stream was not canceled. Replaces the prior hand-rolled wire format, the 1500ms emit throttle, and the `RawMessage.text` / `update_interval_ms` divergences (all unwound in #93 PR 3). | +| Teams native streaming via the SDK `IStreamer` (DMs) | `TeamsAdapter._handle_message_activity` captures a Teams SDK `IStreamer` (`microsoft_teams.apps.StreamerProtocol` / `HttpStream`) for DMs via `app.activity_sender.create_stream(ref)` on `microsoft-teams-apps` 2.0.x, or `HttpStream(app.api.from_service_url(ref.service_url), ref)` on 2.1.x, which removed `ActivitySender` (#250), registers it in `_active_streams`, and `await`s a `processing_done` gate (a wrapped `wait_until` shim) so the streamer stays alive while the handler streams. `stream()` → `_stream_via_emit` calls `stream.emit(text)` per chunk and NEVER calls `close()`; the adapter's `_handle_message_activity` `finally` calls `stream.close()` once (the lifecycle-owner role the SDK App's `process_activity` plays upstream). | `@chat-adapter/teams@4.30.0` `index.ts` does exactly this: `this.activeStreams.set(threadId, ctx.stream)`, build `processingDone` + wrapped `waitUntil`, `await processingDone`, `streamViaEmit` calls `stream.emit(text)` and never `close()` (the SDK App auto-closes after the handler returns). | **No adapter-level divergence.** The only mechanical difference is the close call site: upstream lets the SDK `App` auto-close `ctx.stream` because the SDK owns dispatch; our bridge overrides `server.on_request`, so we own dispatch and reproduce the close in `_handle_message_activity`'s `finally`. The SDK `HttpStream.close` no-ops when the stream was canceled or had no content, so closing in both success and cancel paths is safe (matching the SDK App, which closes in both its success and `StreamCancelledError` branches). Cancellation is detected via `stream.canceled` (checked before each emit) and by catching `StreamCancelledError` (other exceptions re-raise). The first chunk id is captured via `on_chunk` and awaited only when text was emitted and the stream was not canceled. Replaces the prior hand-rolled wire format, the 1500ms emit throttle, and the `RawMessage.text` / `update_interval_ms` divergences (all unwound in #93 PR 3). | +| Teams: `microsoft-teams-apps` 2.0.x and 2.1.x both supported (#250) | The `[teams]` extra allows `>=2.0.13,<2.2` for `microsoft-teams-{apps,api,cards,common}` (common is imported directly and apps leaves it unbounded, so it is declared and capped too). SDK differences are feature-detected, never version-sniffed. **Streaming:** `_create_streamer` uses `app.activity_sender.create_stream(ref)` when the App has one (2.0.x) and otherwise builds `HttpStream` on `app.api.from_service_url(ref.service_url)` (2.1.x, mirroring 2.1's `ActivityContext.stream`). Either way the stream owns a client pinned to the inbound service URL, so `_point_app_api_at` retargeting the shared `App.api` for an outbound call cannot redirect an in-flight stream. An SDK with neither entry point falls back to buffered posting. **Edit/delete:** unchanged. `app.api.conversations.activities(id).update/delete` exists on every supported version; on 2.1 it routes through `conversations.update_activity(..., service_url=None)`, which falls back to the retargeted client URL. The flattened `update_activity`/`delete_activity` are not used because 2.0.13.4 (the floor) lacks them. **Inbound auth:** 2.1 replaced the Bot Framework-only `TokenValidator.for_service` with `InboundActivityTokenValidator`, which also accepts Entra ID (Agent 365 "Agent ID") tokens whose audience is the app id, from any tenant (issuer taken from the token's `tid`) and without the `serviceurl` claim check. It also picks the Entra branch from the token's *unverified* `iss`, building a per-`tid` validator and fetching that tenant's JWKS (blocking) before the signature is checked. So `BridgeHttpAdapter.dispatch` calls the adapter's `_rejects_before_auth` hook before the SDK route handler: it decodes the Bearer token without verification and answers 401 unless `iss == app.cloud.token_issuer` (an unreadable token is refused too; the SDK would refuse it). A token that passes still gets the SDK's full validation. `_dispatch_activity` repeats the issuer check on the validated `JsonWebToken` as defence in depth. On 2.0.x the SDK already enforces the issuer; the pre-check only spares the Bot Framework JWKS fetch for tokens it would reject. Requests with no `Bearer` header, and all requests in `dangerously_allow_unauthenticated_requests` / `skip_auth` mode (the SDK ignores the header there), go to the SDK unchanged | `@chat-adapter/teams@4.41.1` still depends on `@microsoft/teams.*` `^2.0.14` and calls `ctx.stream` / `app.api.conversations.activities(...)` | Python-only SDK compatibility, not a behavior divergence: on either SDK line the adapter streams, edits and authenticates as it did on 2.0.x. Agent ID activities are not supported by this adapter, so accepting their tokens would only widen who can deliver activities. Pinned by `TestCreateStreamer` (both SDK shapes plus an unstubbed real-SDK test that retargets `App.api` mid-stream), `TestOutboundServiceUrlRouting` (edit/delete checked at the HTTP boundary) and `TestInboundTokenIssuer` (real bridge → SDK validator path with RS256 test-key tokens, only JWKS key resolution stubbed: Entra tokens refused with no JWKS fetched, Bot Framework tokens accepted and still audience/`serviceurl`-checked, sovereign-cloud issuer via `CLOUD=USGov`, the post-validation guard, and unauthenticated mode) and `TestSdkDependencyDeclarations` (every imported `microsoft_teams.*` package declared and capped). CI installs 2.1.x only; 2.0.16 was run locally. **A live Teams check of native DM streaming, edit and delete on 2.1 is still owed before release.** | | Teams streaming throttle / Bot Framework wire format ownership | The SDK `HttpStream` owns the entire Bot Framework streaming wire format (`streamType`/`streamSequence`/`streamId`), the inter-flush throttle, and 429 retry. We hand it text via `emit()` and read back the assigned id via `on_chunk`. | Same — `@microsoft/teams.apps`'s `IStreamer` owns all of this in the JS SDK. | **THROTTLE PARITY (verified against the installed SDK source — `microsoft-teams-apps==2.0.13.4`):** the SDK throttles and is 429-safe, so we don't regress to rate-limit errors: (1) `http_stream.py:266` — after a flush, if more is queued, the next flush is scheduled via `call_later(0.5, …)`, i.e. a 500ms inter-flush delay (the module docstring at `http_stream.py:39-41` states this is "to ensure we dont hit API rate limits with Microsoft Teams"); (2) `http_stream.py:283,290` — `add_stream_update(self._index)` stamps the Bot Framework `streamSequence` and `self._index` increments per stream activity; (3) `http_stream.py:285-288` — each chunk send goes through `retry(..., RetryOptions(max_delay=4.0, max_attempts=8))`, so transient 429s are retried with backoff; (4) `http_stream.py:180-201` — `close()` waits for the queue to drain (`_wait_for_id_and_queue`) and the final `add_stream_final()` send also goes through `retry()`. **A LIVE Teams check (streaming a real long response without a 429) is out of scope for this build and is flagged for the reviewers/maintainer.** | | Teams divider rendering | `card_to_adaptive_card`'s `_hoist_dividers` post-processing pass (`teams/cards.py`) hoists `separator: True` onto the next sibling (or emits a non-empty Container for a trailing divider) | `convertDividerToElement` emits an empty `Container` with `separator: True` | Upstream shares the same bug: Microsoft Teams renders an empty Container at zero height, so the separator line is effectively invisible. Python port fixes locally (issue #45) via the `_hoist_dividers` pass rather than blocking on upstream. | | `SlackAdapter.current_token` / `current_token_async` / `current_client` | Public accessors that return the request-context-bound token and a preconfigured `AsyncWebClient`. `current_token` (sync `@property`) reads the cache; `current_token_async` (async method) invokes the resolver on demand for callable `bot_token` configs used outside `handle_webhook`. | Not exposed (`getToken()` is private on the TS `SlackAdapter`) | Python-only addition (issue #47). Downstream code that calls Slack Web APIs from inside a handler — email resolution, user profile fetches, reaction bookkeeping — otherwise depends on underscore-prefixed helpers. The async variant is required because the sync `current_token` cannot drive an async resolver (see `bot_token` resolver invocation site row). | diff --git a/pyproject.toml b/pyproject.toml index 284301f6..061a44a8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -56,9 +56,13 @@ teams = [ "aiohttp>=3.9", # Official Microsoft Teams Apps Python SDK (issue #93 migration; mirrors upstream # adapter-teams@4.30.0's @microsoft/teams.* deps). Graph stays hand-rolled (no [graph] extra). - "microsoft-teams-apps>=2.0.13,<2.1", - "microsoft-teams-api>=2.0.13,<2.1", - "microsoft-teams-cards>=2.0.13,<2.1", + # Both 2.0.x and 2.1.x are supported (feature-detected, #250). The upper bound stays one + # minor ahead because uv.lock is not committed and 2.1.0 shipped breaking changes (#251). + # common is imported directly (ClientOptions) and apps leaves it unbounded, so bound it too. + "microsoft-teams-apps>=2.0.13,<2.2", + "microsoft-teams-api>=2.0.13,<2.2", + "microsoft-teams-cards>=2.0.13,<2.2", + "microsoft-teams-common>=2.0.13,<2.2", ] telegram = ["aiohttp>=3.9"] whatsapp = ["aiohttp>=3.9"] @@ -75,9 +79,10 @@ all = [ "pynacl>=1.5", "aiohttp>=3.9", "google-auth>=2.0", - "microsoft-teams-apps>=2.0.13,<2.1", - "microsoft-teams-api>=2.0.13,<2.1", - "microsoft-teams-cards>=2.0.13,<2.1", + "microsoft-teams-apps>=2.0.13,<2.2", + "microsoft-teams-api>=2.0.13,<2.2", + "microsoft-teams-cards>=2.0.13,<2.2", + "microsoft-teams-common>=2.0.13,<2.2", ] [build-system] @@ -136,9 +141,10 @@ dev = [ "pyrefly==0.61.1", # Teams adapter SDK (issue #93) — in dev so CI (`uv sync --group dev`) installs it # and the Teams test suite runs; matches how the other optional adapter deps are wired. - "microsoft-teams-apps>=2.0.13,<2.1", - "microsoft-teams-api>=2.0.13,<2.1", - "microsoft-teams-cards>=2.0.13,<2.1", + "microsoft-teams-apps>=2.0.13,<2.2", + "microsoft-teams-api>=2.0.13,<2.2", + "microsoft-teams-cards>=2.0.13,<2.2", + "microsoft-teams-common>=2.0.13,<2.2", ] # --------------------------------------------------------------------------- diff --git a/src/chat_sdk/adapters/teams/adapter.py b/src/chat_sdk/adapters/teams/adapter.py index 2e588290..52f2914c 100644 --- a/src/chat_sdk/adapters/teams/adapter.py +++ b/src/chat_sdk/adapters/teams/adapter.py @@ -19,6 +19,7 @@ from urllib.parse import quote, urlparse if TYPE_CHECKING: + from microsoft_teams.api import ConversationReference from microsoft_teams.apps import StreamerProtocol from chat_sdk.adapters.teams.bridge import BridgeHttpAdapter @@ -327,7 +328,7 @@ def __init__(self, config: TeamsAdapterConfig | None = None) -> None: # ``handle_webhook`` can dispatch serverless webhooks through it. The # SDK is an optional ([teams] extra) dependency, so it is imported # lazily here rather than at module scope. - self._bridge = BridgeHttpAdapter(self._logger) + self._bridge = BridgeHttpAdapter(self._logger, reject_before_auth=self._rejects_before_auth) self._app = self._build_app(config) self._app_initialized = False @@ -539,6 +540,9 @@ async def _dispatch_activity(self, event: Any) -> Any: the HTTP response. Card actions (``invoke``) return the Bot Framework invoke acknowledgement; everything else returns ``200`` with no body. """ + if not self._is_bot_framework_token(getattr(event, "token", None)): + return {"status": 401, "body": {"error": "Unauthorized"}} + activity = self._activity_to_dict(event) activity_type = activity.get("type", "") self._logger.debug("Teams activity received", {"type": activity_type}) @@ -569,6 +573,89 @@ async def _dispatch_activity(self, event: Any) -> Any: return {"status": 200, "body": None} + def _rejects_before_auth(self, headers: dict[str, str]) -> bool: + """Return ``True`` to answer 401 before the SDK's JWT validator runs. + + ``microsoft-teams-apps`` 2.0.x validates inbound activities against the + Bot Framework issuer only (``TokenValidator.for_service``). 2.1.x + switched to ``InboundActivityTokenValidator``, which picks its branch + from the token's *unverified* ``iss``. An Entra-looking issuer + (``login.microsoftonline.com/...`` or ``sts.windows.net/...``) sends it + down the Agent 365 "Agent ID" branch: it builds a per-``tid`` Entra + validator and fetches that tenant's JWKS (a blocking HTTP call), and it + accepts any tenant's Entra token whose audience is our app id, without + the Bot Framework ``serviceurl`` claim check. This adapter does not + support Agent ID activities, so lifting the ``<2.1`` cap (#250) must + not widen who can deliver activities, nor let an unauthenticated + request pick which JWKS the bot fetches. + + We therefore read the Bearer token's unverified ``iss`` here and reject + anything but the cloud's Bot Framework issuer. A token that passes still + goes through the SDK's full validation (signature, audience, expiry, + ``serviceurl``), so this only narrows. On 2.0.x the SDK would reject + the same tokens after fetching the fixed Bot Framework JWKS. + + Requests without a ``Bearer`` Authorization header, and every request + in ``dangerously_allow_unauthenticated_requests`` mode (where the SDK + ignores the header), are left to the SDK unchanged. + """ + if self._auth_disabled(): + return False + authorization = headers.get("authorization") or headers.get("Authorization") or "" + if not authorization.startswith("Bearer "): + return False + import jwt + + try: + claims = jwt.decode(authorization.removeprefix("Bearer "), options={"verify_signature": False}) + except jwt.InvalidTokenError: + # Unreadable token: the SDK would reject it too. Its issuer is not + # the Bot Framework's, so reject here without touching the SDK. + claims = {} + return not self._is_bot_framework_issuer(claims.get("iss")) + + def _auth_disabled(self) -> bool: + """Whether the SDK App skips inbound JWT validation. + + ``dangerously_allow_unauthenticated_requests`` on 2.0.14+ (set from the + option or the SDK's env var), ``skip_auth`` on 2.0.13. + """ + options = self._app.options + allow = getattr(options, "dangerously_allow_unauthenticated_requests", None) + if allow is None: + allow = getattr(options, "skip_auth", False) + return bool(allow) + + def _is_bot_framework_token(self, token: Any) -> bool: + """Return ``False`` for a validated inbound token not issued by the Bot Framework. + + Defence in depth behind :meth:`_rejects_before_auth`, which already + turns such tokens away before validation: the same issuer rule, applied + to the token the SDK validated, in case a request ever reaches + :meth:`_dispatch_activity` without passing the bridge's pre-check. + + Only a ``JsonWebToken`` is checked: the SDK wraps every token it + validated in one. In ``dangerously_allow_unauthenticated_requests`` + mode the SDK passes a placeholder token instead, and that mode is left + as the SDK defines it. + """ + from microsoft_teams.api import JsonWebToken + + if not isinstance(token, JsonWebToken): + return True + return self._is_bot_framework_issuer(token.issuer) + + def _is_bot_framework_issuer(self, issuer: Any) -> bool: + """Whether ``issuer`` is this cloud's Bot Framework token issuer; logs a rejection.""" + expected = self._app.cloud.token_issuer + if issuer == expected: + return True + self._logger.warn( + "Teams activity rejected: inbound token was not issued by the Bot Framework", + {"issuer": issuer, "expectedIssuer": expected}, + ) + return False + @staticmethod def _activity_to_dict(event: Any) -> dict[str, Any]: """Extract the camelCase activity dict from a Teams SDK activity event. @@ -921,13 +1008,25 @@ def _create_streamer(self, activity: dict[str, Any], thread_id: str) -> Streamer Builds the :class:`ConversationReference` the streamer needs from the inbound activity (``recipient`` → bot account, ``conversation``, - ``channelId``, ``serviceUrl``) and hands it to the SDK's - ``ActivitySender.create_stream`` — the same call the SDK's own - ``ActivityContext`` makes to expose ``ctx.stream``. The returned - ``HttpStream`` owns the Bot Framework streaming wire format + ``channelId``, ``serviceUrl``) and builds the stream the same way the + installed SDK's own ``ActivityContext`` builds ``ctx.stream``. The + returned ``HttpStream`` owns the Bot Framework streaming wire format (``streamType``/``streamSequence``/``streamId``), the per-flush throttle, and 429 retry/backoff. We never poke its internals. + Supports both SDK lines by feature detection (not version sniffing): + + * ``microsoft-teams-apps`` 2.0.x exposes ``App.activity_sender``; + ``ActivitySender.create_stream(ref)`` builds a fresh ``ApiClient`` + on ``ref.service_url``. + * 2.1.x removed ``ActivitySender``. ``ActivityContext.stream`` now + calls ``HttpStream(app.api, ref)`` directly, and ``HttpStream`` + scopes its own client with ``api.from_service_url(ref.service_url)``. + + Either way the stream gets its own client pinned to the inbound + activity's service URL, so a later :meth:`_point_app_api_at` (which + retargets the shared ``App.api`` in place) cannot move it. + Returns ``None`` when the activity lacks the fields the SDK ref requires (``serviceUrl``), so the caller can fall back to fire-and-forget processing + accumulate-and-post. @@ -940,7 +1039,12 @@ def _create_streamer(self, activity: dict[str, Any], thread_id: str) -> Streamer try: _validate_service_url(service_url) ref = build_conversation_reference(activity, bot_app_id=self._app_id) - return self._app.activity_sender.create_stream(ref) + activity_sender = getattr(self._app, "activity_sender", None) + if activity_sender is not None: + # microsoft-teams-apps 2.0.x + return activity_sender.create_stream(ref) + # microsoft-teams-apps 2.1+: no ActivitySender. + return self._new_http_stream(ref) except Exception as exc: self._logger.warn( "Failed to create Teams streamer; falling back to buffered post", @@ -948,6 +1052,27 @@ def _create_streamer(self, activity: dict[str, Any], thread_id: str) -> Streamer ) return None + def _new_http_stream(self, ref: ConversationReference) -> StreamerProtocol: + """Build an ``HttpStream`` for ``ref`` on SDKs without ``ActivitySender`` (2.1+). + + Mirrors 2.1's ``ActivityContext.stream`` (``HttpStream(self.api, ref)``), + except that we hand ``HttpStream`` a client already scoped with + ``App.api.from_service_url(ref.service_url)`` rather than the shared + ``App.api``. 2.1's ``HttpStream`` scopes the client again itself, so + this costs one extra lightweight clone (the HTTP connection is shared). + It keeps the stream off the shared client even if a future + ``HttpStream`` stops scoping, because ``_point_app_api_at`` retargets + ``App.api`` in place for outbound calls. + + An SDK whose ``ApiClient`` has no ``from_service_url`` raises + ``AttributeError`` here, which :meth:`_create_streamer` turns into the + buffered-post fallback. + """ + from microsoft_teams.apps import HttpStream + + scoped_api = self._app.api.from_service_url(ref.service_url) + return HttpStream(scoped_api, ref) + # Keys injected by the SDK's card renderer or Teams transport — not user input. _ACTION_TRANSPORT_KEYS = frozenset({"actionId", "msteams"}) @@ -1422,6 +1547,12 @@ def _point_app_api_at(self, service_url: str) -> None: that replace ``self._app.api`` with a mock lacking those attributes — an ``AttributeError`` there is harmless because the mock ignores the service URL anyway. + + Works on both ``microsoft-teams-apps`` 2.0.x and 2.1.x. On 2.1 the + activity methods gained optional ``service_url=`` / ``agentic_identity=`` + keywords; we pass neither, so they fall back to the client's own (just + retargeted) ``service_url``, and ``app.send`` sees + ``ref.service_url == api.service_url`` and sends on ``App.api`` itself. """ _validate_service_url(service_url) normalized = service_url.rstrip("/") diff --git a/src/chat_sdk/adapters/teams/bridge.py b/src/chat_sdk/adapters/teams/bridge.py index 19b5e3e8..91eaee5d 100644 --- a/src/chat_sdk/adapters/teams/bridge.py +++ b/src/chat_sdk/adapters/teams/bridge.py @@ -24,7 +24,7 @@ import inspect import json -from collections.abc import Iterable +from collections.abc import Callable, Iterable from typing import TYPE_CHECKING, Any, Literal, cast from chat_sdk.logger import Logger @@ -49,10 +49,23 @@ class BridgeHttpAdapter: ``packages/adapter-teams/src/bridge-adapter.ts``. """ - def __init__(self, logger: Logger) -> None: + def __init__( + self, + logger: Logger, + reject_before_auth: Callable[[dict[str, str]], bool] | None = None, + ) -> None: + """Create the bridge. + + ``reject_before_auth`` (Python-only, #250) is called with the request + headers before the SDK route handler runs. When it returns ``True`` the + request is answered ``401`` without reaching the SDK's JWT validator, + so an unauthenticated request cannot steer which JWKS the validator + fetches. ``None`` keeps upstream's behavior (everything goes to the SDK). + """ self._handler: HttpRouteHandler | None = None self._webhook_options: dict[str, WebhookOptions] = {} self._logger = logger + self._reject_before_auth = reject_before_auth # ------------------------------------------------------------------ # HttpServerAdapter protocol @@ -122,6 +135,13 @@ async def dispatch(self, request: Any, options: WebhookOptions | None = None) -> headers = self._read_headers(request) + if self._reject_before_auth is not None and self._reject_before_auth(headers): + return _make_response( + json.dumps({"error": "Unauthorized"}), + 401, + content_type="application/json", + ) + activity_id = parsed_body.get("id") if activity_id and options is not None: self._webhook_options[activity_id] = options diff --git a/src/chat_sdk/adapters/teams/streamer.py b/src/chat_sdk/adapters/teams/streamer.py index f1f6251e..ecef8a58 100644 --- a/src/chat_sdk/adapters/teams/streamer.py +++ b/src/chat_sdk/adapters/teams/streamer.py @@ -7,11 +7,12 @@ isolated from ``adapter.py`` and importable lazily (Port Rule: optional/SDK deps imported inside functions, not at module top). -Mirrors what the Teams SDK's own ``ActivityContext`` does in -``microsoft_teams/apps/app_process.py`` ``_build_context`` (it builds a -``ConversationReference`` from the activity and calls -``ActivitySender.create_stream(ref)`` to expose ``ctx.stream``). Our bridge -owns dispatch, so we reproduce just the reference-building step. +Mirrors what the Teams SDK's own ``ActivityContext`` does: it builds a +``ConversationReference`` from the activity and exposes ``ctx.stream`` from it +(``ActivitySender.create_stream(ref)`` on ``microsoft-teams-apps`` 2.0.x, +``HttpStream(api, ref)`` on 2.1.x). Our bridge owns dispatch, so we reproduce +just the reference-building step; ``TeamsAdapter._create_streamer`` picks the +stream constructor for the installed SDK. """ from __future__ import annotations diff --git a/tests/test_fixture_replay.py b/tests/test_fixture_replay.py index e4792c61..8e3f7d6e 100644 --- a/tests/test_fixture_replay.py +++ b/tests/test_fixture_replay.py @@ -15,7 +15,6 @@ import hashlib import hmac -import inspect import json import time from typing import Any @@ -125,35 +124,32 @@ def _teams_request(body: str) -> _FakeRequest: def _teams_skip_auth(): - """Force the Microsoft Teams SDK to skip inbound JWT validation. + """Configure the Microsoft Teams SDK App to skip inbound JWT validation. Inbound auth moved into the SDK ``App`` (issue #93 PR 1). Fixture replays - carry no signed Bot Framework token, so this patches ``HttpServer.initialize`` - to enable ``skip_auth`` — exercising the real bridge → SDK → handler path - without signature checks. Returns a context manager covering both - ``adapter.initialize()`` (where the route + validator are set up) and the - subsequent ``handle_webhook`` dispatch. + carry no signed Bot Framework token, so this turns on the App's own + unauthenticated-requests option just before ``App.initialize`` hands it to + the ``HttpServer`` -- the same state a consumer gets from setting the + option, so the adapter's pre-validation issuer check (#250) sees the mode + too -- exercising the real bridge -> SDK -> handler path without signature + checks. Returns a context manager covering both ``adapter.initialize()`` + (where the route + validator are set up) and the subsequent + ``handle_webhook`` dispatch. """ - from microsoft_teams.apps.http.http_server import HttpServer - - real_initialize = HttpServer.initialize - - # microsoft-teams-apps 2.0.14+ renamed the SDK's ``skip_auth`` flag to - # ``dangerously_allow_unauthenticated_requests``; force whichever flag this - # version has and forward everything else untouched (the SDK calls - # ``initialize`` with keywords only). - skip_flag = ( - "dangerously_allow_unauthenticated_requests" - if "dangerously_allow_unauthenticated_requests" in inspect.signature(real_initialize).parameters - else "skip_auth" - ) + from microsoft_teams.apps import App + + real_initialize = App.initialize - def _initialize_skip_auth(self, *args, **kwargs): - kwargs.pop("skip_auth", None) - kwargs.pop("dangerously_allow_unauthenticated_requests", None) - return real_initialize(self, *args, **{**kwargs, skip_flag: True}) + async def _initialize_skip_auth(self, *args, **kwargs): + # microsoft-teams-apps 2.0.14+ renamed the option ``skip_auth`` to + # ``dangerously_allow_unauthenticated_requests``; set whichever exists. + if hasattr(self.options, "dangerously_allow_unauthenticated_requests"): + self.options.dangerously_allow_unauthenticated_requests = True + else: + self.options.skip_auth = True + return await real_initialize(self, *args, **kwargs) - return patch.object(HttpServer, "initialize", _initialize_skip_auth) + return patch.object(App, "initialize", _initialize_skip_auth) def _gchat_request(body: str) -> _FakeRequest: diff --git a/tests/test_teams_adapter.py b/tests/test_teams_adapter.py index 9f33190b..34da2498 100644 --- a/tests/test_teams_adapter.py +++ b/tests/test_teams_adapter.py @@ -5,6 +5,7 @@ from __future__ import annotations +import functools import inspect import re from datetime import datetime, timezone @@ -847,6 +848,193 @@ async def test_400_for_invalid_json(self): assert response["status"] == 400 +@functools.cache +def _test_rsa_key(): + """One RSA key per test run for signing test JWTs (generation is slow).""" + from cryptography.hazmat.primitives.asymmetric import rsa + + return rsa.generate_private_key(public_exponent=65537, key_size=2048) + + +class TestInboundTokenIssuer: + """Only Bot Framework-issued inbound tokens reach the handlers (#250). + + ``microsoft-teams-apps`` 2.1 also accepts Entra ID ("Agent ID") tokens + for our app id from any tenant, and picks that branch (with a per-tenant + JWKS fetch) from the token's unverified issuer. These tests run the real + bridge -> SDK ``HttpServer`` -> SDK validator -> ``_dispatch_activity`` path + with RS256 tokens signed by a test key. Only JWKS key *resolution* is + stubbed (``PyJWKClient.get_signing_key_from_jwt`` returns the test public + key and records which JWKS URI was asked), so the SDK's own signature, + issuer, audience, expiry and ``serviceurl`` checks all run. + """ + + TENANT = "11111111-2222-3333-4444-555555555555" + SERVICE_URL = "https://smba.trafficmanager.net/teams/" + ACTIVITY = { + "type": "message", + "id": "msg-1", + "text": "hello", + "from": {"id": "user-1", "name": "Alice"}, + "recipient": {"id": "28:test-app-id", "name": "bot"}, + "conversation": {"id": "19:abc@thread.tacv2", "conversationType": "channel"}, + "channelId": "msteams", + "serviceUrl": SERVICE_URL, + } + + @pytest.fixture + def signing_key(self): + return _test_rsa_key() + + @pytest.fixture(autouse=True) + def jwks_uris(self, monkeypatch, signing_key): + """Resolve every JWKS lookup to the test key; return the URIs asked for.""" + from types import SimpleNamespace + + import jwt + + asked: list[str] = [] + + def get_signing_key_from_jwt(self, token): + asked.append(self.uri) + return SimpleNamespace(key=signing_key.public_key()) + + monkeypatch.setattr(jwt.PyJWKClient, "get_signing_key_from_jwt", get_signing_key_from_jwt) + monkeypatch.delenv("CLOUD", raising=False) + monkeypatch.delenv("DANGEROUSLY_ALLOW_UNAUTHENTICATED_REQUESTS", raising=False) + return asked + + def _token(self, signing_key, issuer: str, **overrides) -> str: + import time + + import jwt + + now = int(time.time()) + claims = { + "iss": issuer, + "aud": "test-app-id", + "tid": self.TENANT, + "serviceurl": self.SERVICE_URL, + "iat": now, + "nbf": now, + "exp": now + 600, + **overrides, + } + return jwt.encode(claims, signing_key, algorithm="RS256", headers={"kid": "test-kid"}) + + async def _post(self, token: str): + import json + + logger = _make_logger() + adapter = _make_adapter(logger=logger) + chat = MagicMock() + chat.get_state = MagicMock(return_value=MagicMock(set=AsyncMock(), get=AsyncMock(return_value=None))) + chat.process_message = MagicMock() + await adapter.initialize(chat) + request = _FakeRequest( + json.dumps(self.ACTIVITY), + {"content-type": "application/json", "authorization": f"Bearer {token}"}, + ) + response = await adapter.handle_webhook(request) + return response, chat, logger + + @staticmethod + def _issuer_warnings(logger): + return [c.args[1] for c in logger.warn.call_args_list if "not issued by the Bot Framework" in c.args[0]] + + @pytest.mark.asyncio + async def test_entra_issued_token_is_rejected_before_any_jwks_fetch(self, signing_key, jwks_uris): + # A correctly signed Entra token for our app id: 2.1's validator would + # accept it on its Entra branch after fetching the tenant's JWKS. + issuer = f"https://login.microsoftonline.com/{self.TENANT}/v2.0" + response, chat, logger = await self._post(self._token(signing_key, issuer)) + assert response["status"] == 401 + chat.process_message.assert_not_called() + assert jwks_uris == [] + assert self._issuer_warnings(logger) == [{"issuer": issuer, "expectedIssuer": "https://api.botframework.com"}] + + @pytest.mark.asyncio + async def test_attacker_chosen_tenant_jwks_is_never_fetched(self, jwks_uris): + # Unsigned tokens naming arbitrary tenants must not make the bot fetch + # those tenants' JWKS (v2 and v1 Entra issuer shapes). + import jwt + + for issuer in ( + "https://login.microsoftonline.com/attacker-tenant/v2.0", + "https://sts.windows.net/attacker-tenant/", + ): + token = jwt.encode({"iss": issuer, "aud": "test-app-id", "tid": "attacker-tenant"}, "k" * 40) + response, _chat, _logger = await self._post(token) + assert response["status"] == 401 + assert jwks_uris == [] + + @pytest.mark.asyncio + async def test_bot_framework_issued_token_is_dispatched(self, signing_key, jwks_uris): + response, chat, logger = await self._post(self._token(signing_key, "https://api.botframework.com")) + assert response["status"] == 200 + chat.process_message.assert_called_once() + assert jwks_uris == ["https://login.botframework.com/v1/.well-known/keys"] + assert self._issuer_warnings(logger) == [] + + @pytest.mark.asyncio + async def test_bot_framework_token_still_gets_sdk_validation(self, signing_key): + # The pre-check only narrows: a Bot Framework-issued token for another + # app id (audience) or another service URL is still refused by the SDK. + for overrides in ({"aud": "someone-elses-app"}, {"serviceurl": "https://smba.trafficmanager.net/other/"}): + token = self._token(signing_key, "https://api.botframework.com", **overrides) + response, chat, _logger = await self._post(token) + assert response["status"] == 401, overrides + chat.process_message.assert_not_called() + + @pytest.mark.asyncio + async def test_sovereign_cloud_uses_its_own_issuer(self, monkeypatch, signing_key): + # App.cloud comes from the CLOUD env var; the issuer rule must follow it. + monkeypatch.setenv("CLOUD", "USGov") + response, chat, _logger = await self._post(self._token(signing_key, "https://api.botframework.us")) + assert response["status"] == 200 + chat.process_message.assert_called_once() + + response, chat, logger = await self._post(self._token(signing_key, "https://api.botframework.com")) + assert response["status"] == 401 + chat.process_message.assert_not_called() + assert self._issuer_warnings(logger) == [ + {"issuer": "https://api.botframework.com", "expectedIssuer": "https://api.botframework.us"} + ] + + @pytest.mark.asyncio + async def test_dispatch_rejects_validated_non_bot_framework_token(self, signing_key): + # Defence in depth: even if a validated Entra token reached the SDK's + # on_request callback, _dispatch_activity refuses it before routing. + from types import SimpleNamespace + + from microsoft_teams.api import JsonWebToken + + logger = _make_logger() + adapter = _make_adapter(logger=logger) + adapter._handle_message_activity = AsyncMock() + issuer = f"https://login.microsoftonline.com/{self.TENANT}/v2.0" + event = SimpleNamespace(body=MagicMock(), token=JsonWebToken(value=self._token(signing_key, issuer))) + response = await adapter._dispatch_activity(event) + assert response == {"status": 401, "body": {"error": "Unauthorized"}} + adapter._handle_message_activity.assert_not_awaited() + assert self._issuer_warnings(logger) == [{"issuer": issuer, "expectedIssuer": "https://api.botframework.com"}] + + @pytest.mark.asyncio + async def test_unauthenticated_mode_ignores_the_header(self, monkeypatch, jwks_uris): + # In the SDK's dangerously_allow_unauthenticated_requests mode the SDK + # ignores Authorization entirely; the pre-check must not add auth back. + import jwt + + logger = _make_logger() + adapter = _make_adapter(logger=logger) + adapter._app.options.dangerously_allow_unauthenticated_requests = True + token = jwt.encode({"iss": "https://login.microsoftonline.com/t/v2.0"}, "k" * 40) + assert adapter._rejects_before_auth({"authorization": f"Bearer {token}"}) is False + adapter._app.options.dangerously_allow_unauthenticated_requests = False + assert adapter._rejects_before_auth({"authorization": f"Bearer {token}"}) is True + assert jwks_uris == [] + + # --------------------------------------------------------------------------- # initialize # --------------------------------------------------------------------------- @@ -899,6 +1087,35 @@ def test_app_id_mapped_to_sdk_client_id(self): assert adapter._app.id == "my-bot-id" +class TestSdkDependencyDeclarations: + def test_every_imported_sdk_package_is_declared_with_an_upper_bound(self): + # uv.lock is not committed, so an SDK package that only arrives + # transitively (microsoft-teams-apps leaves -common unbounded) can + # drift to a new minor on a fresh install (#250/#251). Every + # microsoft_teams. the adapter imports must be declared, capped, + # in each dependency list that installs the Teams SDK. + import pathlib + import tomllib + + root = pathlib.Path(__file__).resolve().parent.parent + source = "\n".join(p.read_text() for p in (root / "src/chat_sdk/adapters/teams").glob("*.py")) + imported = set(re.findall(r"\bmicrosoft_teams\.(\w+)", source)) + assert {"api", "apps", "common"} <= imported + + project = tomllib.loads((root / "pyproject.toml").read_text()) + lists = { + "[teams]": project["project"]["optional-dependencies"]["teams"], + "[all]": project["project"]["optional-dependencies"]["all"], + "dev": project["dependency-groups"]["dev"], + } + for list_name, deps in lists.items(): + specs = {d.split(">")[0].split("<")[0].split("=")[0]: d for d in deps} + for pkg in sorted(imported): + dist = f"microsoft-teams-{pkg}" + assert dist in specs, f"{dist} missing from {list_name}" + assert "<" in specs[dist], f"{dist} has no upper bound in {list_name}" + + # --------------------------------------------------------------------------- # renderFormatted # --------------------------------------------------------------------------- @@ -1031,23 +1248,51 @@ async def fake_send(conversation_id, activity): assert seen["conversations"] == self.SOVEREIGN_URL.rstrip("/") assert seen["activities"] == self.SOVEREIGN_URL.rstrip("/") + @staticmethod + def _capture_wire(adapter: TeamsAdapter, method: str) -> list[str]: + """Stub the REAL activities client's HTTP ``method`` and record each URL. + + Patching at the HTTP boundary (not ``activities_client.update``) keeps + the SDK's own call chain in the test, so it holds on every supported + ``microsoft-teams-apps`` line: 2.0.x calls ``http.put(url, json=...)``; + 2.1.x routes ``activities(...).update`` through + ``conversations.update_activity(..., service_url=None, agentic_identity=None)`` + and adds a ``_metadata=`` kwarg (#250). + """ + urls: list[str] = [] + + class _Response: + def json(self) -> dict[str, str]: + return {"id": "edit-1"} + + async def fake_http(url, **_kwargs): + urls.append(url) + return _Response() + + http = adapter._app.api.conversations.activities_client.http + setattr(http, method, fake_http) + return urls + @pytest.mark.asyncio async def test_edit_message_retargets_real_activities_client(self): adapter = _make_adapter(app_id="test-app-id", logger=_make_logger()) - seen: dict[str, str] = {} - - # Patch the real activities_client.update so the routing target is read - # off the REAL client chain (not a wholesale api mock). - async def fake_update(conversation_id, activity_id, activity): - seen["url"] = adapter._app.api.conversations.activities_client.service_url - return _SentActivity(activity_id) + urls = self._capture_wire(adapter, "put") + tid = adapter.encode_thread_id( + TeamsThreadId(conversation_id="19:abc@thread.tacv2", service_url=self.SOVEREIGN_URL) + ) + result = await adapter.edit_message(tid, "edit-1", {"markdown": "x"}) + assert urls == [f"{self.SOVEREIGN_URL.rstrip('/')}/v3/conversations/19:abc@thread.tacv2/activities/edit-1"] + assert result.id == "edit-1" - adapter._app.api.conversations.activities_client.update = fake_update # type: ignore[method-assign] + @pytest.mark.asyncio + async def test_delete_message_retargets_real_activities_client(self): + adapter = _make_adapter(app_id="test-app-id", logger=_make_logger()) + urls = self._capture_wire(adapter, "delete") tid = adapter.encode_thread_id( TeamsThreadId(conversation_id="19:abc@thread.tacv2", service_url=self.SOVEREIGN_URL) ) - await adapter.edit_message(tid, "edit-1", {"markdown": "x"}) - assert seen["url"] == self.SOVEREIGN_URL.rstrip("/") + await adapter.delete_message(tid, "gone-1") + assert urls == [f"{self.SOVEREIGN_URL.rstrip('/')}/v3/conversations/19:abc@thread.tacv2/activities/gone-1"] class TestFileAttachments: diff --git a/tests/test_teams_native_streaming.py b/tests/test_teams_native_streaming.py index 0742345a..dccf3f5e 100644 --- a/tests/test_teams_native_streaming.py +++ b/tests/test_teams_native_streaming.py @@ -730,24 +730,95 @@ def process_message(adapter_arg, thread_id, message, options): class TestCreateStreamer: - def test_creates_streamer_for_valid_dm_activity(self): + """``_create_streamer`` builds the SDK stream on both SDK lines (#250). + + ``microsoft-teams-apps`` 2.0.x exposes ``App.activity_sender.create_stream``; + 2.1.x removed ``ActivitySender`` and builds ``HttpStream(api, ref)``. The + branch tests force each shape by stubbing the App, so both run on + whichever SDK is installed. ``test_real_sdk_stream_is_pinned_to_inbound_service_url`` + runs unstubbed against the installed SDK. + """ + + SOVEREIGN_URL = "https://smba.infra.gov.teams.microsoft.us/teams/" + + @staticmethod + def _drop_activity_sender(adapter: TeamsAdapter) -> None: + """Make the App look like 2.1.x (no ``activity_sender``).""" + if hasattr(adapter._app, "activity_sender"): + del adapter._app.activity_sender + + def test_real_sdk_stream_is_pinned_to_inbound_service_url(self): + """No stubs: the installed SDK yields an ``HttpStream`` on its own client. + + The stream's client must target the inbound activity's service URL and + must not be the shared ``App.api``, which ``_point_app_api_at`` retargets + in place for every outbound call. If it were shared, an outbound call to + another region mid-stream would redirect the stream's chunks. + """ + from microsoft_teams.apps import HttpStream + adapter = _make_adapter() tid = _dm_thread_id(adapter) - created = {} + streamer = adapter._create_streamer(_dm_activity(), tid) + + assert isinstance(streamer, HttpStream) + assert streamer._ref.conversation.id == "a:1Abc-DM-conversation-id" + assert streamer._ref.bot.id == "28:test-app-id" + assert streamer._client is not adapter._app.api + assert streamer._client.service_url == "https://smba.trafficmanager.net/teams" + + adapter._point_app_api_at(self.SOVEREIGN_URL) + assert adapter._app.api.service_url == self.SOVEREIGN_URL.rstrip("/") + assert streamer._client.service_url == "https://smba.trafficmanager.net/teams" + + def test_uses_activity_sender_when_present(self): + """2.0.x shape: ``activity_sender.create_stream(ref)`` builds the stream.""" + adapter = _make_adapter() + tid = _dm_thread_id(adapter) + created: dict[str, Any] = {} + fake = FakeStreamer() def create_stream(ref): created["ref"] = ref - return FakeStreamer() + return fake - adapter._app.activity_sender.create_stream = create_stream # type: ignore[method-assign] + adapter._app.activity_sender = MagicMock(create_stream=create_stream) - streamer = adapter._create_streamer(_dm_activity(), tid) - assert streamer is not None + assert adapter._create_streamer(_dm_activity(), tid) is fake ref = created["ref"] assert ref.conversation.id == "a:1Abc-DM-conversation-id" assert ref.service_url == "https://smba.trafficmanager.net/teams/" assert ref.bot.id == "28:test-app-id" + def test_builds_http_stream_on_scoped_client_without_activity_sender(self, monkeypatch): + """2.1.x shape: ``HttpStream`` gets ``App.api`` scoped to the inbound service URL.""" + import microsoft_teams.apps as teams_apps + + adapter = _make_adapter() + tid = _dm_thread_id(adapter) + self._drop_activity_sender(adapter) + + scoped_api = object() + api = MagicMock() + api.from_service_url = MagicMock(return_value=scoped_api) + adapter._app.api = api + + built: dict[str, Any] = {} + fake = FakeStreamer() + + def fake_http_stream(client, ref): + built["client"] = client + built["ref"] = ref + return fake + + monkeypatch.setattr(teams_apps, "HttpStream", fake_http_stream) + + assert adapter._create_streamer(_dm_activity(), tid) is fake + api.from_service_url.assert_called_once_with("https://smba.trafficmanager.net/teams/") + assert built["client"] is scoped_api + assert built["ref"].conversation.id == "a:1Abc-DM-conversation-id" + assert built["ref"].bot.id == "28:test-app-id" + def test_returns_none_without_service_url(self): adapter = _make_adapter() tid = _dm_thread_id(adapter) @@ -770,8 +841,24 @@ def test_returns_none_when_create_stream_raises(self): def boom(ref): raise RuntimeError("sdk exploded") - adapter._app.activity_sender.create_stream = boom # type: ignore[method-assign] + adapter._app.activity_sender = MagicMock(create_stream=boom) + assert adapter._create_streamer(_dm_activity(), tid) is None + adapter._logger.warn.assert_called_once() + assert adapter._logger.warn.call_args.args[1]["error"] == "sdk exploded" + + def test_returns_none_when_sdk_has_neither_stream_entry_point(self): + """No ``activity_sender`` and no ``ApiClient.from_service_url`` → buffered fallback. + + Guards an SDK shape we do not support: falling through to + ``HttpStream(App.api, ref)`` there could share the retargetable client. + """ + adapter = _make_adapter() + tid = _dm_thread_id(adapter) + self._drop_activity_sender(adapter) + adapter._app.api = MagicMock(spec=["service_url"]) + assert adapter._create_streamer(_dm_activity(), tid) is None + adapter._logger.warn.assert_called_once() # ---------------------------------------------------------------------------