feat(relay): flat in_channel continuable cron surface on the relay lane

Field report (enterprise side-by-side, 2026-08-18, finding 1 — the relay-only
blocker): on relay-fronted Slack, cron briefs always deliver into a dedicated
thread; the flat continuable surface (cron_continuable_surface: in_channel)
that native Slack supports is inert, so plain DM replies never continue the
job and the main conversation never sees the brief.

Three gaps closed:
- CapabilityDescriptor gains supports_inchannel_continuable (default False,
  additive within contract_version 1; from_json ignores it from old
  connectors, old gateways filter it as unknown). The connector advertises
  it per platform at handshake.
- RelayAdapter maps the bit onto the adapter capability surface in both the
  constructor and _apply_descriptor (renegotiation), so the scheduler's D6
  fail-safe gate sees it exactly like native Slack's class attribute.
- _resolve_cron_surface_mode replaces the scheduler's inline flat-key read:
  native keeps the shipped flat shape; the relay lane reads the same
  per-logical-platform sub-block as the documented relay Slack knobs
  (platforms.relay.extra.slack.cron_continuable_surface), sub-block wins,
  scoped so a slack block cannot leak onto other fronted platforms.

The seed path needs no changes: RelayAdapter inherits set_session_store
(wired by the generic adapter boot loop) and _seed_cron_channel_session
keys the flat session off the logical platform_name.

12 new tests: descriptor default/from_json/legacy-absence, adapter mapping
constructor + renegotiation, and the surface-knob matrix (native flat key,
relay sub-block, per-platform scoping, precedence, defaults).
This commit is contained in:
Victor Kyriazakos
2026-08-19 13:44:27 +00:00
parent b5f95d0890
commit 85b89451f6
4 changed files with 192 additions and 7 deletions
+34 -7
View File
@@ -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.<p>.extra.cron_continuable_surface``
(shipped shape, unchanged).
- Relay-fronted: ``platforms.relay.extra.<logical>.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
+12
View File
@@ -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."""
+8
View File
@@ -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 =
@@ -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"