diff --git a/packages/client/README.md b/packages/client/README.md index 6d18ab65..2836d07f 100644 --- a/packages/client/README.md +++ b/packages/client/README.md @@ -459,8 +459,9 @@ async def main(): # 1. Which skills does this config reference? Pure projection — no I/O. info = await inspect_config("doc-agent", {"kind": "user", "key": "user-123"}) if info["config"] is None: - # The config could not be resolved. Stop here: an empty reference list - # passed to write_skills would prune every skill it manages. + # The config could not be resolved. skill_refs raises ValueError on + # None rather than return [], which write_skills would read as "prune + # every skill it manages". Stop here and keep what is on disk. return refs = skill_refs(info["config"]) # [SkillReference(key='pdf-extraction', version=2)] @@ -488,7 +489,9 @@ as above, writes only what the resolved variation asked for. but is not a list of `{key, version}` objects (key matching `^[a-z0-9][a-z0-9-]*$`, version an integer ≥ 1), including `skills: null`. One bad entry rejects the whole field, so `write_skills` never receives a partial list that would prune skills the config still -references. Config parsing does not check `skills`, so a malformed field never fails +references. It also raises when the config itself is not a dict, including the `None` a failed +`inspect_config` returns, so an outage cannot reach `write_skills` as an empty list and prune +every managed skill. A config with no `skills` field returns `[]`. Config parsing does not check `skills`, so a malformed field never fails `config().invoke()` or other core calls. **Integrity is not optional.** Content is returned only when its sha256 (lowercase hex, over @@ -566,9 +569,11 @@ any handler. The same mapping is attached as `extra["ld_skills"]` for structured ERROR ld.skills.integrity_failure {"action":"withheld","event":"ld.skills.integrity_failure","expected_hash":"0000…0000","language":"python","observed_hash":"5fc8…6ec0","reason":"content hash mismatch","reason_code":"hash_mismatch","skill_key":"pdf-extraction","version":2} ``` -**`ld.skills.integrity_failure` is a stability commitment.** Match on it; it will not be -renamed. JSON keys are sorted, so the line is byte-identical across LaunchDarkly's AI SDKs for -the same input. +**Match on `ld.skills.integrity_failure`.** It is the name LaunchDarkly's AI SDKs share for this +record. While Agent Skills is experimental, the event name, its fields, and its `reason_code` +values may change in a minor release; any such change gets a changelog entry under +**Experimental**, so review it before upgrading. JSON keys are sorted, so the line is byte-identical across LaunchDarkly's +AI SDKs for the same input. | Field | Description | |---|---| @@ -782,7 +787,7 @@ that skips verification. | Export | Description | |---|---| -| `skill_refs(config)` | Project a config's `skills` array into `list[SkillReference]`. Pure — no client, store, or network needed. Returns `[]` when the field is absent or the config is not a dict. Raises `ValueError` when the field is present but malformed (including `null`), so an unreadable field never reaches a pruning reconcile as "no skills". | +| `skill_refs(config)` | Project a config's `skills` array into `list[SkillReference]`. Pure — no client, store, or network needed. Returns `[]` when the field is absent. Raises `ValueError` when the config is not a dict (including the `None` a failed `inspect_config` returns) or the field is present but malformed (including `null`), so neither reaches a pruning reconcile as "no skills". | | `get_skill(key, *, version=None)` | One verified skill, or `None`. `version=None` means newest available; a specific `version` matches exactly. Raises only when no store is configured. | | `get_skill_result(key, *, version=None)` | The same retrieval, reporting **why**: a frozen `SkillOutcome` with `.skill`, `.reason` (`ok` / `absent` / `integrity_failure` / `store_unavailable` / `wrong_version`), and `.detail`. See *Failing closed on tampering* above. Raises only when no store is configured. | | `get_skills(refs)` | Batch form. Accepts `SkillReference` values and bare key strings (string = latest). Results follow input order; missing or unverifiable entries are omitted. | diff --git a/packages/client/agents.md b/packages/client/agents.md index 9888dcb9..1ca5ae3a 100644 --- a/packages/client/agents.md +++ b/packages/client/agents.md @@ -207,7 +207,8 @@ Three layers, in increasing order of blast radius: typed `SkillReference` values. Pure: no network, no client, no store, no telemetry. It also validates the array, and **fails closed**: a present but malformed field (including `null`, or one bad entry) raises `ValueError` rather than returning a partial - list that would authorize a prune. `parse_ai_config` deliberately does not check `skills`, + list that would authorize a prune. So does a config that is not a dict, including the + `None` a failed `inspect_config` returns; only a dict with no `skills` key yields `[]`. `parse_ai_config` deliberately does not check `skills`, so an experimental field cannot fail a core config call (TESTING.md §0.3). 2. **Content accessors** — `get_skill`, `get_skill_result`, `get_skills`, `all_skills` read through the `SkillStore` seam. Configure a store with `set_skill_store(store)`; with @@ -299,6 +300,38 @@ resumes in place once the cause is fixed, clearing the terminal reason through answer) surfaces as an `HTTPError` that `_classify_status` maps to a fatal, non-retried failure. +**A known stream event whose data is not JSON drops the connection; it is never skipped.** +`_iter_sse` logs a WARNING and raises `_RecoverableTransportError`, as the base SDK's FDv2 +stream interrupts on a `JSONDecodeError`, so the in-flight payload is abandoned and the +reconnect sends the last committed basis. Skipping the event would let the +`payload-transferred` after it commit without it and advance the basis: a lost +`delete-object` would then never be retransmitted, and a lost `put-object` in an `xfer-full` +would revoke that skill by omission. The warning is logged in `_iter_sse` because `_run` logs +a disconnect after a completed exchange at debug. Both cases are asserted by +`test_a_malformed_event_abandons_its_transfer_and_keeps_the_basis`. Only the events in +`_EVENTS_WITH_DATA` (the ones `handle` reads) are parsed: a `heart-beat` or an unknown event +is yielded with `None` data, so its data cannot drop the connection, as `ldclient` parses only +inside its known-event branches. Empty data on a known event is not JSON either, and drops the +connection the same way: read as `None`, a `delete-object` would be ignored and a +`payload-transferred` would commit with no selector. + +**A put or delete that is valid JSON but has no usable `key` is warned and ignored, for the +experimental stage.** Its transfer still commits, so a keyless `delete-object` loses its +revocation and a keyless `put-object` in an `xfer-full` revokes that skill by omission (and +`write_skills("*")` prunes it). `ldclient` interrupts the stream instead. Interrupting is not +strictly safer, because a server that keeps resending the object would then block every later +change, so the choice is to be revisited before 1.0 (TESTING.md §3.25). Pinned by +`test_a_delete_with_no_usable_key_is_ignored_and_its_transfer_commits` and +`test_a_put_with_no_usable_key_in_an_xfer_full_revokes_by_omission`; changing the rule should fail +them. + +**A `goodbye` with `catastrophe: true` is a recoverable, counted disconnect, not a fatal.** +The Python base SDK does not read the flag and the Go SDK only logs it, so stopping delivery +on it would leave this store the only LaunchDarkly SDK that needs `start()` after a server +incident. `_goodbye` logs it at ERROR and returns a disconnect without `recycled`, so the +delivery loop counts it even after a completed exchange. Asserted by +`test_a_catastrophic_goodbye_reconnects_and_is_counted`. + **Reads are memory-bounded.** `_read_bounded` (poll bodies) and `_iter_stream_lines`/`_iter_sse` (each line and each event) enforce `MAX_RESPONSE_BYTES` (64 MiB). Crossing it raises `_ResponseTooLargeError`, a fatal error: nothing from that @@ -588,7 +621,9 @@ Do not undo any of these as a simplification: discriminate (`resolve_from_store` and `list_raw_objects` also log ERROR for a raising store), and the stdlib's default formatter drops `extra`, so an `extra`-only record is invisible under a plain `logging.basicConfig()`. -- **`ld.skills.integrity_failure` is documented for customers to match on.** Never rename it. +- **`ld.skills.integrity_failure` is documented for customers to match on.** Do not rename it + casually. The experimental stage allows a rename in a minor release, but only in every SDK + at once and with a changelog entry, because a rename silently breaks customers' alerts. - **`sort_keys=True` makes the line byte-identical across SDKs** (modulo `language`), since the other implementations build the object in alphabetical key order. - **Optional fields are omitted, never nulled**, so a SIEM field-existence check means @@ -940,6 +975,6 @@ a conversation is out of reach at this layer either way. - `Skill.content` is opaque `bytes`. Do not add anything that parses or interprets it — no YAML library in this package's dependencies at any tier, and no accessor that decodes content. - Do not route skills telemetry through `client.track()`, and do not introduce an LD context anywhere in the skills path. Signals go through the `skills_core.py` emitter seam, whose default is a no-op, and only via its `record_*` functions. - Do not add a signal name outside the three in the Agent Skills table above — the list is an allowlist. `AgentControl Skill SDK Reference Returned` and `AgentControl Skill Content Retrieved` were considered and deliberately excluded from SDK emission. -- Do not rename `ld.skills.integrity_failure`, and do not add an eleventh `reason_code` in one language only — both are documented compatibility surfaces. See "The integrity-failure log record" above. +- Do not rename `ld.skills.integrity_failure` in one language only or without a changelog entry, and do not add an eleventh `reason_code` in one language only — both are documented compatibility surfaces, though the experimental stage lets either change in a minor release. See "The integrity-failure log record" above. - Do not relax any of the `write_skills` filesystem defenses (local key re-validation, symlink refusal, manifest-authorized destruction, corrupt-manifest fail-closed, atomic `0644` writes). Each is a deliberate security property with abuse-case tests attached. - Do not make `SkillStore` lookups key-only. Version is part of the lookup identity because a payload holds several versions of one key; a key-only seam cannot express a version-pinned reference. diff --git a/packages/client/src/launchdarkly_ai_server/skills.py b/packages/client/src/launchdarkly_ai_server/skills.py index c02e1a6c..c44b9ff7 100644 --- a/packages/client/src/launchdarkly_ai_server/skills.py +++ b/packages/client/src/launchdarkly_ai_server/skills.py @@ -204,21 +204,30 @@ def skill_refs(config: AiConfigRep | None) -> list[SkillReference]: Returns the skill references attached to a resolved AI Config. Pure: no network, store, or telemetry. Returns ``[]`` when the config has no - ``skills`` field, or when *config* is not a dict (for example ``None`` from a - failed ``inspect_config``). Typical use: ``await get_skills(skill_refs(config))``. + ``skills`` field. Typical use: ``await get_skills(skill_refs(config))``. The config parser does not validate ``skills``, so a malformed field does not fail core config calls. It is validated here instead, and rejected whole: ``write_skills`` with ``prune=True`` would delete the files of any skill missing from the list, so a partial or empty list is never returned - for a field that is present. + for a field that is present. A *config* that is not a dict, such as the + ``None`` a failed ``inspect_config`` returns, is rejected for the same + reason: read as "no skills", it would prune every managed skill. Raises: - ValueError: If ``skills`` is present but is not a list of ``{key, - version}`` objects with a valid key and an integer version >= 1. - This includes ``skills: null``. + ValueError: If *config* is not a dict (including ``None``), or if + ``skills`` is present but is not a list of ``{key, version}`` + objects with a valid key and an integer version >= 1. This + includes ``skills: null``. """ - if not isinstance(config, dict) or "skills" not in config: + if not isinstance(config, dict): + raise ValueError( + f"skill_refs was given {type(config).__name__}, not an AI Config. " + "A failed inspect_config returns None as its config; handle that " + "before deriving skill references, because an empty list passed to " + "write_skills would prune every skill it manages." + ) + if "skills" not in config: return [] raw = config["skills"] diff --git a/packages/client/src/launchdarkly_ai_server/skills_fdv2.py b/packages/client/src/launchdarkly_ai_server/skills_fdv2.py index d3108f1e..90aee139 100644 --- a/packages/client/src/launchdarkly_ai_server/skills_fdv2.py +++ b/packages/client/src/launchdarkly_ai_server/skills_fdv2.py @@ -108,6 +108,20 @@ _EVENT_GOODBYE = "goodbye" _EVENT_ERROR = "error" +# The events whose data ``_FDv2Reader.handle`` reads. Only these are parsed, so +# a ``heart-beat`` or an unknown event with data that is not JSON is ignored, +# as ``ldclient``'s FDv2 stream ignores it. +_EVENTS_WITH_DATA = frozenset( + ( + _EVENT_SERVER_INTENT, + _EVENT_PUT_OBJECT, + _EVENT_DELETE_OBJECT, + _EVENT_PAYLOAD_TRANSFERRED, + _EVENT_GOODBYE, + _EVENT_ERROR, + ) +) + _INTENT_TRANSFER_FULL = "xfer-full" _INTENT_TRANSFER_CHANGES = "xfer-changes" _INTENT_TRANSFER_NONE = "none" @@ -482,7 +496,6 @@ class _TransferOutcome: committed: bool = False changes: list[dict[str, Any]] = field(default_factory=list) basis: str | None = None - fatal: str | None = None disconnect: str | None = None up_to_date: bool = False """A ``none`` intent: the content held is current. Counts as a healthy @@ -738,14 +751,23 @@ def _goodbye(self, data: Any) -> _TransferOutcome: catastrophe = bool(data.get("catastrophe")) if isinstance(data, dict) else False silent = bool(data.get("silent")) if isinstance(data, dict) else False self._abandon_in_flight() + if catastrophe: + # Recoverable, as in the base SDKs: Python's does not read the flag + # and Go's only logs it. Not ``recycled``, so it is counted even + # after a completed exchange, and logged here at ERROR because the + # delivery loop logs a disconnect after one at debug. + logger.error( + "The FDv2 server reported a catastrophic failure (%s); " + "reconnecting with backoff", + reason, + ) + return _TransferOutcome( + disconnect=f"server sent a catastrophic goodbye: {reason}" + ) if not silent: # Debug only: a goodbye after a completed exchange is a routine # recycle, and the delivery loop warns when one is not. logger.debug("FDv2 connection closing: %s", reason) - if catastrophe: - return _TransferOutcome( - fatal=f"server sent a catastrophic goodbye: {reason}" - ) return _TransferOutcome( disconnect=f"server said goodbye: {reason}", recycled=True ) @@ -1235,6 +1257,17 @@ def _iter_sse(response: Any) -> Any: newlines, blank line dispatches, ``:`` comments skipped. An event over ``MAX_RESPONSE_BYTES`` raises a fatal error and the in-flight payload is abandoned. + + A known event whose data is not JSON, empty data included, is logged at + WARNING and raises ``_RecoverableTransportError``, as the base SDK's FDv2 + stream does (``json.loads("")`` raises there too): the + in-flight payload is abandoned and the reconnect resumes from the last + committed basis. Skipping the event instead would let the + ``payload-transferred`` after it commit the transfer without it and advance + the basis past it, so a lost ``delete-object`` would never be retransmitted. + The warning is logged here because the delivery loop logs a disconnect + after a completed exchange at debug. Any other event is yielded with + ``None`` data, unparsed. """ limit = MAX_RESPONSE_BYTES try: @@ -1246,15 +1279,19 @@ def _iter_sse(response: Any) -> Any: if line == "": if name is not None: payload = "\n".join(data_lines) - try: - parsed = json.loads(payload) if payload else None - except json.JSONDecodeError: - logger.warning( - "Discarding FDv2 '%s' event whose data was not JSON", name - ) - parsed = None - else: - yield name, parsed + parsed = None + if name in _EVENTS_WITH_DATA: + try: + parsed = json.loads(payload) + except json.JSONDecodeError as exc: + message = ( + f"an FDv2 '{name}' event's data was not JSON " + f"({exc}); the connection was dropped and " + "nothing from the in-flight payload was applied" + ) + logger.warning("%s; reconnecting", message) + raise _RecoverableTransportError(message) from exc + yield name, parsed name = None data_lines = [] event_bytes = 0 @@ -1796,8 +1833,6 @@ def _apply(self, name: str, data: Any) -> bool: self._publish_first_payload() if outcome.changes: self._notify(outcome.changes) - if outcome.fatal: - raise _FatalTransportError(outcome.fatal) if outcome.disconnect: raise _RecoverableTransportError( outcome.disconnect, recycled=outcome.recycled diff --git a/packages/client/tests/test_skills.py b/packages/client/tests/test_skills.py index bbfcb548..ebfe133a 100644 --- a/packages/client/tests/test_skills.py +++ b/packages/client/tests/test_skills.py @@ -287,8 +287,23 @@ def test_emits_no_telemetry(self, recording_emitter: Any) -> None: skill_refs(self._config(skills=[{"key": "a", "version": 1}])) assert recording_emitter.records == [] - def test_a_non_dict_config_returns_empty_list(self) -> None: - assert skill_refs(None) == [] + @pytest.mark.parametrize( + "config", + [ + pytest.param(None, id="none"), + pytest.param([], id="list"), + pytest.param("doc-agent", id="string"), + ], + ) + def test_a_config_that_is_not_a_dict_raises(self, config: Any) -> None: + """``None`` is what a failed ``inspect_config`` returns as its config. + + Read as "no skills", it would let ``write_skills(skill_refs(config), + root)``, which prunes by default, delete every managed skill during an + outage. + """ + with pytest.raises(ValueError, match="not an AI Config"): + skill_refs(config) @pytest.mark.parametrize( "malformed", diff --git a/packages/client/tests/test_skills_fdv2.py b/packages/client/tests/test_skills_fdv2.py index 6e3770db..ce5c0e69 100644 --- a/packages/client/tests/test_skills_fdv2.py +++ b/packages/client/tests/test_skills_fdv2.py @@ -200,6 +200,22 @@ def full_payload( ) +def sse_body( + payload_events: list[dict[str, Any]], *, truncate: int = -1, blank: bool = False +) -> bytes: + """Serialises *payload_events* as SSE, corrupting the data at *truncate*. + + The data there is cut short, or emptied when *blank* is set. + """ + lines = [] + for i, event in enumerate(payload_events): + data = json.dumps(event["data"]) + if i == truncate: + data = "" if blank else data[:-5] + lines.append(f"event: {event['event']}\ndata: {data}\n\n") + return "".join(lines).encode() + + # --------------------------------------------------------------------------- # The fake endpoint # --------------------------------------------------------------------------- @@ -868,7 +884,7 @@ def test_flag_and_segment_objects_are_skipped_cleanly(self) -> None: assert held.get("pdf-extraction", None) is not None assert reader.diagnostics.objects_ignored == 4 assert reader.diagnostics.skill_objects_received == 1 - assert all(o.fatal is None and o.disconnect is None for o in outcomes) + assert all(o.disconnect is None for o in outcomes) def test_an_unknown_kind_is_ignored_rather_than_fatal(self) -> None: """ @@ -886,15 +902,74 @@ def test_an_unknown_kind_is_ignored_rather_than_fatal(self) -> None: } outcomes = drive(reader, full_payload(("put-object", exotic))) assert len(held) == 0 - assert all(o.fatal is None and o.disconnect is None for o in outcomes) + assert all(o.disconnect is None for o in outcomes) def test_an_unknown_event_name_is_ignored(self) -> None: held = _SkillObjectSet() reader = _ProtocolReader(held) outcome = reader.handle("some-future-event", {"anything": True}) - assert outcome.fatal is None assert outcome.disconnect is None + def test_a_delete_with_no_usable_key_is_ignored_and_its_transfer_commits( + self, + ) -> None: + """Pins TESTING.md §3.25's experimental-stage rule: ignored, not fatal. + + The revocation it carried is lost, so ``revoked`` stays served. Base + ``ldclient`` interrupts the stream instead; the choice is to be + revisited before 1.0, and changing it should fail this test. + """ + held = _SkillObjectSet() + reader = _ProtocolReader(held) + drive( + reader, + full_payload( + ("put-object", put_skill("a", object_version=1)), + ("put-object", put_skill("revoked", object_version=1)), + ), + ) + keyless = delete_skill("revoked", object_version=1) + del keyless["key"] + outcomes = drive( + reader, + events( + ("server-intent", server_intent("xfer-changes")), + ("put-object", put_skill("fresh", object_version=1)), + ("delete-object", keyless), + ("payload-transferred", transferred("basis-2")), + ), + ) + assert all(o.disconnect is None for o in outcomes) + assert outcomes[-1].committed + assert held.get("revoked", None) is not None + assert held.get("fresh", None) is not None + + def test_a_put_with_no_usable_key_in_an_xfer_full_revokes_by_omission( + self, + ) -> None: + """Pins the other half of the same rule. + + The skill the put carried is missing from the full transfer, so it is + revoked, and ``write_skills("*")`` would prune it. + """ + held = _SkillObjectSet() + reader = _ProtocolReader(held) + drive(reader, full_payload(("put-object", put_skill("a", object_version=1)))) + keyless = put_skill("a", object_version=1) + del keyless["key"] + outcomes = drive( + reader, + full_payload( + ("put-object", put_skill("b", object_version=1)), + ("put-object", keyless), + state="basis-2", + ), + ) + assert all(o.disconnect is None for o in outcomes) + assert outcomes[-1].committed + assert held.get("a", None) is None + assert held.get("b", None) is not None + def test_a_heartbeat_does_nothing(self) -> None: reader = _ProtocolReader(_SkillObjectSet()) outcome = reader.handle("heart-beat", None) @@ -922,14 +997,21 @@ def test_a_goodbye_asks_for_a_reconnect(self) -> None: reader = _ProtocolReader(_SkillObjectSet()) outcome = reader.handle("goodbye", {"reason": "rebalancing", "silent": False}) assert outcome.disconnect is not None - assert outcome.fatal is None + assert outcome.recycled is True + + def test_a_catastrophic_goodbye_is_a_counted_disconnect(self) -> None: + """Not fatal: neither the Python nor the Go base SDK stops on it. - def test_a_catastrophic_goodbye_is_fatal(self) -> None: + Not ``recycled`` either, so the delivery loop counts it as a failure + even after a completed exchange. + """ reader = _ProtocolReader(_SkillObjectSet()) outcome = reader.handle( "goodbye", {"reason": "no", "silent": False, "catastrophe": True} ) - assert outcome.fatal is not None + assert outcome.disconnect is not None + assert "catastroph" in outcome.disconnect + assert outcome.recycled is False def test_transfer_none_holds_everything_and_commits_nothing(self) -> None: """ @@ -2691,6 +2773,61 @@ def stream(self, basis: str | None) -> Any: warnings = [r for r in caplog.records if r.levelname == "WARNING"] assert any("server said goodbye" in r.getMessage() for r in warnings) + def test_a_catastrophic_goodbye_reconnects_and_is_counted( + self, caplog: Any + ) -> None: + """Delivery keeps going, from the basis reached, and says so at ERROR. + + It follows a commit here, which would make an ordinary goodbye a quiet + recycle; a catastrophe is still counted and still logged. It is also + ``silent``, which a catastrophe's ERROR ignores (unlike Go's). + """ + + class _CatastropheThenQuiet(_FakeRequester): + def __init__(self) -> None: + self.bases: list[str | None] = [] + + def stream(self, basis: str | None) -> Any: + self.bases.append(basis) + if len(self.bases) > 1: + return _BlockingConnection() + # Through the SSE parser, so a goodbye whose data it stopped + # parsing would arrive as None and never be read as catastrophic. + return _StreamConnection( + _LineSource( + sse_body( + full_payload(("put-object", put_skill())) + + events( + ( + "goodbye", + { + "reason": "meltdown", + "catastrophe": True, + "silent": True, + }, + ) + ) + ) + ) + ) + + requester = _CatastropheThenQuiet() + store = stream_store(_requester=requester) + with caplog.at_level("DEBUG", logger="launchdarkly_ai_server.skills_fdv2"): + try: + store.start() + assert wait_until(lambda: len(requester.bases) >= 2) + assert store.failed is None + assert store.diagnostics.connection_failures == 1 + assert "catastroph" in (store.diagnostics.last_error or "") + assert requester.bases[1] == "basis-1" + assert store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") + finally: + store.close() + errors = [r for r in caplog.records if r.levelname == "ERROR"] + assert any("meltdown" in r.getMessage() for r in errors) + assert not any("will not retry" in r.getMessage() for r in errors) + def test_a_recycled_connection_reconnects_quietly(self, caplog: Any) -> None: # A healthy idle stream reconnects for as long as the process runs, so # warning on each one would fill a customer's logs with a fault they do @@ -2707,6 +2844,137 @@ def test_a_recycled_connection_reconnects_quietly(self, caplog: Any) -> None: assert not [r for r in caplog.records if r.levelname == "WARNING"] assert [r for r in caplog.records if "reconnecting in" in r.getMessage()] + @pytest.mark.parametrize( + ( + "initial_keys", + "retransmission", + "held_at_reconnect", + "after", + "revoked", + "blank", + ), + [ + pytest.param( + ("a", "revoked"), + events( + ("server-intent", server_intent("xfer-changes")), + ("put-object", put_skill("fresh", object_version=1)), + ("delete-object", delete_skill("revoked", object_version=1)), + ("payload-transferred", transferred("basis-2")), + ), + {"a": True, "revoked": True, "fresh": False}, + {"a": True, "revoked": False, "fresh": True}, + 1, + False, + id="delete-object-in-xfer-changes", + ), + pytest.param( + ("a",), + events( + ("server-intent", server_intent("xfer-full")), + ("put-object", put_skill("b", object_version=1)), + ("put-object", put_skill("a", object_version=1)), + ("payload-transferred", transferred("basis-2")), + ), + {"a": True, "b": False}, + {"a": True, "b": True}, + 0, + False, + id="put-object-in-xfer-full", + ), + pytest.param( + ("a", "revoked"), + events( + ("server-intent", server_intent("xfer-changes")), + ("put-object", put_skill("fresh", object_version=1)), + ("delete-object", delete_skill("revoked", object_version=1)), + ("payload-transferred", transferred("basis-2")), + ), + {"a": True, "revoked": True, "fresh": False}, + {"a": True, "revoked": False, "fresh": True}, + 1, + True, + id="empty-delete-object-in-xfer-changes", + ), + ], + ) + def test_a_malformed_event_abandons_its_transfer_and_keeps_the_basis( + self, + caplog: Any, + initial_keys: tuple[str, ...], + retransmission: list[dict[str, Any]], + held_at_reconnect: dict[str, bool], + after: dict[str, bool], + revoked: int, + blank: bool, + ) -> None: + """A corrupted event fails the connection; it is not skipped. + + The event at index 2 of the second transfer is cut short, or emptied. Skipping it + would let the ``payload-transferred`` after it commit the rest of the + transfer and advance the basis past it: a lost ``delete-object`` would + leave a revoked skill served indefinitely, because the reconnect asks + only for changes since that basis, and a lost ``put-object`` in an + ``xfer-full`` would revoke that skill by omission. The drop follows a + commit on the same connection, so the delivery loop logs it at debug; + the reader warns. + """ + initial = full_payload( + *(("put-object", put_skill(key, object_version=1)) for key in initial_keys), + state="basis-1", + ) + keys = held_at_reconnect.keys() + + class _CorruptedThenClean(_FakeRequester): + def __init__(self) -> None: + self.bases: list[str | None] = [] + self.held_at_reconnect: dict[str, bool] = {} + self.store: FDv2SkillStore | None = None + + def stream(self, basis: str | None) -> Any: + self.bases.append(basis) + if len(self.bases) == 1: + body = sse_body( + initial + retransmission, truncate=len(initial) + 2, blank=blank + ) + return _StreamConnection(_LineSource(body)) + if len(self.bases) == 2: + assert self.store is not None + self.held_at_reconnect = { + key: self.store.get_object(SKILL_OBJECT_KIND, key) is not None + for key in keys + } + return _StreamConnection(_LineSource(sse_body(retransmission))) + return _BlockingConnection() + + requester = _CorruptedThenClean() + store = stream_store(_requester=requester) + requester.store = store + with caplog.at_level("DEBUG", logger="launchdarkly_ai_server.skills_fdv2"): + try: + store.start() + assert wait_until(lambda: len(requester.bases) >= 3) + # Nothing from the corrupted transfer committed: not the event + # it lost, and not the put that arrived intact beside it. + assert requester.held_at_reconnect == held_at_reconnect + # The reconnect resumed from the last committed basis. + assert requester.bases[:2] == [None, "basis-1"] + # The clean retransmission applied in full. + assert { + key: store.get_object(SKILL_OBJECT_KIND, key) is not None + for key in keys + } == after + assert store.diagnostics.objects_revoked == revoked + assert requester.bases[2] == "basis-2" + assert store.failed is None + finally: + store.close() + warnings = [r.getMessage() for r in caplog.records if r.levelname == "WARNING"] + corrupted = retransmission[2]["event"] + assert any(f"'{corrupted}' event's data was not JSON" in w for w in warnings), ( + warnings + ) + def test_a_connection_that_never_answered_still_warns(self, caplog: Any) -> None: # The quiet path is earned by answering. A connection that failed before # it told us anything is the case the warning exists for. @@ -3199,6 +3467,62 @@ def test_multi_line_data_under_the_cap_still_decodes(self) -> None: assert list(_iter_sse(source)) == [("put-object", {"a": 1})] assert source.closed + def test_an_event_whose_data_is_not_json_drops_the_connection(self) -> None: + source = _LineSource(b'event: delete-object\ndata: {"key": "a:1\n\n') + with pytest.raises(_RecoverableTransportError, match="not JSON"): + list(_iter_sse(source)) + assert source.closed + + @pytest.mark.parametrize( + "name", + [ + "server-intent", + "put-object", + "delete-object", + "payload-transferred", + "goodbye", + "error", + ], + ) + def test_an_event_handle_reads_is_parsed(self, name: str) -> None: + # Pins the parse set: an event dropped from it would reach ``handle`` + # with None data, and a catastrophic goodbye would read as a recycle. + source = _LineSource( + f'event: {name}\ndata: {{"reason": "x", "catastrophe": true}}\n\n'.encode() + ) + assert list(_iter_sse(source)) == [(name, {"reason": "x", "catastrophe": True})] + + @pytest.mark.parametrize("data_line", [b"", b"data:\n"], ids=["absent", "blank"]) + @pytest.mark.parametrize( + "name", + [ + "server-intent", + "put-object", + "delete-object", + "payload-transferred", + "goodbye", + "error", + ], + ) + def test_an_event_handle_reads_with_no_data_drops_the_connection( + self, name: str, data_line: bytes + ) -> None: + # Empty data is not JSON. Read as None, a delete would be ignored and a + # payload-transferred would commit with no selector. + source = _LineSource(f"event: {name}\n".encode() + data_line + b"\n") + with pytest.raises(_RecoverableTransportError, match="not JSON"): + list(_iter_sse(source)) + + @pytest.mark.parametrize("name", ["heart-beat", "x-future-event"]) + def test_an_event_handle_does_not_read_is_not_parsed(self, name: str) -> None: + # Unknown events are ignored by contract, so data that is not JSON on + # one, or on a heart-beat, must not drop the connection. + source = _LineSource( + f"event: {name}\ndata: not json\n\n".encode() + + b'event: put-object\ndata: {"a": 1}\n\n' + ) + assert list(_iter_sse(source)) == [(name, None), ("put-object", {"a": 1})] + @pytest.fixture def second_endpoint() -> Any: