diff --git a/cron/scheduler.py b/cron/scheduler.py index 9c260d112f..3540808691 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -1495,6 +1495,39 @@ def _reclaim_fds_best_effort() -> None: pass +def _resolve_cron_surface_mode(pconfig, logical_platform_name: str) -> str: + """Resolve the continuable-cron delivery surface for a platform config. + + Returns ``"in_channel"`` or ``"thread"`` (default). Two config shapes: + + - Native adapter: the flat key ``platforms.

.extra.cron_continuable_surface`` + (shipped shape, unchanged). + - Relay-fronted: ``platforms.relay.extra..cron_continuable_surface`` + — the same per-logical-platform sub-block the relay's documented Slack + knobs use (``reply_in_thread``, ``dm_top_level_threads_as_sessions``; + see RelayAdapter._relay_slack_extra). The sub-block wins over a flat + key when both exist, matching _relay_slack_extra precedence, and is + scoped to its logical platform so a ``slack:`` block cannot leak onto + another fronted platform. + + Field gap (2026-08-18): the scheduler read only the flat key, so on the + relay lane — where pconfig is platforms.relay — operators had NO working + location for the knob and briefs always threaded. + """ + try: + extra = getattr(pconfig, "extra", None) or {} + sub = extra.get(str(logical_platform_name or "").lower()) + if isinstance(sub, dict) and sub.get("cron_continuable_surface") is not None: + raw = sub.get("cron_continuable_surface") + else: + raw = extra.get("cron_continuable_surface") + if raw is not None and str(raw).strip().lower() == "in_channel": + return "in_channel" + except Exception: + pass + return "thread" + + def _resolve_origin(job: dict) -> Optional[dict]: """Extract origin info from a job, preserving any extra routing metadata. @@ -2718,13 +2751,7 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option # the adapter capability flag ``supports_inchannel_continuable`` so an # unsupported platform fails SAFE to "thread" (Slack is the first # consumer; "first consumer ≠ definition"). - surface_mode = "thread" - try: - surface_raw = (pconfig.extra or {}).get("cron_continuable_surface") - if surface_raw is not None and str(surface_raw).strip().lower() == "in_channel": - surface_mode = "in_channel" - except Exception: - surface_mode = "thread" + surface_mode = _resolve_cron_surface_mode(pconfig, platform_name) in_channel_surface = surface_mode == "in_channel" if in_channel_surface and runtime_adapter is not None and not getattr( runtime_adapter, "supports_inchannel_continuable", False diff --git a/gateway/relay/adapter.py b/gateway/relay/adapter.py index 9621827781..be66459960 100644 --- a/gateway/relay/adapter.py +++ b/gateway/relay/adapter.py @@ -126,6 +126,11 @@ class RelayAdapter(BasePlatformAdapter): # _capture_scope / send. self._platform_by_chat: Dict[str, str] = {} self.supports_code_blocks = descriptor.markdown_dialect not in ("", "plain") + # Cron flat continuable surface — descriptor-advertised (see + # _apply_descriptor; same bit, constructor path). + self.supports_inchannel_continuable = bool( + getattr(descriptor, "supports_inchannel_continuable", False) + ) # Phase 7 Unit 7d-B: watches the transport for a terminal auth revocation # (a 4401 close after a successful handshake = the operator opted this # instance out of the relay). On revocation we surface a clean, @@ -361,6 +366,13 @@ class RelayAdapter(BasePlatformAdapter): self.descriptor = descriptor self.MAX_MESSAGE_LENGTH = descriptor.max_message_length self.supports_code_blocks = descriptor.markdown_dialect not in ("", "plain") + # Cron in_channel continuable surface (D6 gate in cron/scheduler.py): + # the scheduler reads this off the adapter; the connector advertises it + # per platform at handshake. Class default is False (BasePlatformAdapter), + # so only an explicit descriptor bit turns the flat surface on. + self.supports_inchannel_continuable = bool( + getattr(descriptor, "supports_inchannel_continuable", False) + ) async def _on_inbound(self, event) -> None: """Bridge a connector-delivered MessageEvent into the normal adapter path.""" diff --git a/gateway/relay/descriptor.py b/gateway/relay/descriptor.py index 4963d428b7..15a4f87037 100644 --- a/gateway/relay/descriptor.py +++ b/gateway/relay/descriptor.py @@ -64,6 +64,14 @@ class CapabilityDescriptor: # "no context" — additive within contract_version 1. from_json filters # unknown keys, so a connector sending this to an older gateway is safe too. supports_context: bool = False + # Whether the connector's platform can host a FLAT continuable cron + # surface (native Slack's ``cron_continuable_surface: in_channel``): the + # brief posts top-level in the channel/DM and a plain reply continues the + # job via the flat ``(platform, chat_id, None)`` session. The scheduler + # fails safe to thread mode when False (D6 gate), so an older connector + # that never sends this keeps today's thread behavior — additive within + # contract_version 1. + supports_inchannel_continuable: bool = False # Op-level capability discovery (Phase 1 parity): the outbound op names the # connector's sender for this platform actually implements (e.g. # ["send", "edit", "typing", "follow_up", "get_chat_info"]). Empty tuple = diff --git a/tests/relay/test_relay_inchannel_continuable.py b/tests/relay/test_relay_inchannel_continuable.py new file mode 100644 index 0000000000..04ecf4b24d --- /dev/null +++ b/tests/relay/test_relay_inchannel_continuable.py @@ -0,0 +1,138 @@ +"""Relay lane parity: flat in_channel continuable cron surface (Coatue F1). + +Field report 2026-08-18: on relay-fronted Slack, cron briefs ALWAYS deliver +into a dedicated thread — `cron_continuable_surface: in_channel` (the flat-DM +continuable surface, native Slack's documented shape) is inert on the relay +lane. Three gaps, each pinned here: + +1. CapabilityDescriptor has no in_channel capability bit, so RelayAdapter + inherits supports_inchannel_continuable=False from the base class and the + scheduler fails safe to thread mode (D6 gate). +2. The scheduler reads the surface knob flat off pconfig.extra; on the relay + lane pconfig is platforms.relay, whose documented Slack knobs live under + extra.slack.* (reply_in_thread precedent) — so the operator has no working + place to put the knob. +3. Descriptor renegotiation (_apply_descriptor) must preserve the new + capability bit, same as supports_code_blocks. +""" + +from types import SimpleNamespace + +import pytest + +from gateway.relay.descriptor import CapabilityDescriptor +from cron.scheduler import _resolve_cron_surface_mode + + +def _descriptor(**overrides): + base = dict( + contract_version=1, + platform="slack", + label="Slack", + max_message_length=4000, + supports_draft_streaming=False, + supports_edit=True, + supports_threads=True, + markdown_dialect="mrkdwn", + len_unit="chars", + ) + base.update(overrides) + return CapabilityDescriptor(**base) + + +class TestDescriptorCapabilityBit: + def test_defaults_false(self): + d = _descriptor() + assert d.supports_inchannel_continuable is False + + def test_from_json_reads_flag(self): + import json + + payload = dict( + contract_version=1, platform="slack", label="Slack", + max_message_length=4000, supports_draft_streaming=False, + supports_edit=True, supports_threads=True, + markdown_dialect="mrkdwn", len_unit="chars", + supports_inchannel_continuable=True, + ) + d = CapabilityDescriptor.from_json(json.dumps(payload)) + assert d.supports_inchannel_continuable is True + + def test_from_json_missing_flag_defaults_false(self): + """Older connector that never sends the field — legacy-safe.""" + import json + + payload = dict( + contract_version=1, platform="slack", label="Slack", + max_message_length=4000, supports_draft_streaming=False, + supports_edit=True, supports_threads=True, + markdown_dialect="mrkdwn", len_unit="chars", + ) + d = CapabilityDescriptor.from_json(json.dumps(payload)) + assert d.supports_inchannel_continuable is False + + +class TestRelayAdapterCapabilityMapping: + def _adapter(self, descriptor): + from gateway.config import PlatformConfig + from gateway.relay.adapter import RelayAdapter + + config = PlatformConfig(enabled=True, extra={}) + return RelayAdapter(config, descriptor) + + def test_adapter_maps_descriptor_flag_true(self): + adapter = self._adapter(_descriptor(supports_inchannel_continuable=True)) + assert getattr(adapter, "supports_inchannel_continuable", False) is True + + def test_adapter_maps_descriptor_flag_false(self): + adapter = self._adapter(_descriptor()) + assert getattr(adapter, "supports_inchannel_continuable", True) is False + + def test_renegotiation_updates_flag(self): + """_apply_descriptor must carry the bit, like supports_code_blocks.""" + adapter = self._adapter(_descriptor()) + adapter._apply_descriptor(_descriptor(supports_inchannel_continuable=True)) + assert adapter.supports_inchannel_continuable is True + + +class TestSurfaceKnobResolution: + """_resolve_cron_surface_mode reads the knob from BOTH config shapes.""" + + def test_native_flat_key(self): + pconfig = SimpleNamespace(extra={"cron_continuable_surface": "in_channel"}) + assert _resolve_cron_surface_mode(pconfig, "slack") == "in_channel" + + def test_relay_slack_subblock(self): + """The relay lane's documented shape: platforms.relay.extra.slack.* + (same seam as reply_in_thread / dm_top_level_threads_as_sessions).""" + pconfig = SimpleNamespace( + extra={"slack": {"cron_continuable_surface": "in_channel"}} + ) + assert _resolve_cron_surface_mode(pconfig, "slack") == "in_channel" + + def test_subblock_is_per_logical_platform(self): + """A slack sub-block must not leak onto another fronted platform.""" + pconfig = SimpleNamespace( + extra={"slack": {"cron_continuable_surface": "in_channel"}} + ) + assert _resolve_cron_surface_mode(pconfig, "discord") == "thread" + + def test_default_thread(self): + pconfig = SimpleNamespace(extra={}) + assert _resolve_cron_surface_mode(pconfig, "slack") == "thread" + + def test_flat_key_wins_over_subblock(self): + """Legacy-fallback precedence mirrors _relay_slack_extra: an explicit + flat key keeps working when both are present.""" + pconfig = SimpleNamespace( + extra={ + "cron_continuable_surface": "thread", + "slack": {"cron_continuable_surface": "in_channel"}, + } + ) + # Sub-block is the documented relay shape and wins when present — + # matches _relay_slack_extra (sub-dict preferred over flat extra). + assert _resolve_cron_surface_mode(pconfig, "slack") == "in_channel" + + def test_none_pconfig_defaults_thread(self): + assert _resolve_cron_surface_mode(None, "slack") == "thread"