diff --git a/packages/client/agents.md b/packages/client/agents.md index b867387..a00d912 100644 --- a/packages/client/agents.md +++ b/packages/client/agents.md @@ -331,6 +331,23 @@ unit of consistency. A half-applied full transfer would publish a state the serv described and briefly empty the store, which with pruning on deletes skill files. An interrupted transfer keeps last known good, and listeners fire once per commit. +**A commit is the only thing that publishes a first payload, and a 304 is not one.** +`is_initialized()` — the fact `write_skills("*")` authorizes a prune on — goes true when a +payload *commits*, so every other answer has to stop short of claiming one. A +`payload-transferred` that applied nothing reports neither a commit nor an up-to-date +answer, and does not adopt the selector of a payload it never applied: a `none` intent +builds no pending set, nor does an `intentCode` this SDK does not recognise, and a foreign +payload's contents are declined. A poll adopts the response `ETag` only from a body that +completed an exchange — a commit, or a `none` intent, which is the server saying the +content held is what the etag describes. And a 304 *confirms* the payload held rather than +establishing one, because the exchange it stands in for cannot establish one either. +Loosen any of the three and the other two carry a store that received nothing into a prune +of every managed `SKILL.md` on disk: an empty committed set reads as an environment that +revoked every skill, and a 304 carries nothing to notice it on. There is no cached basis +to make it safe — `_basis` and `_etag` both start as `None` with no injection point, so a +304 reaching a store that holds nothing takes a server answering a request that carried no +etag at all. + **The first payload intent is read, and assumed to be the skill payload.** Delivery sends one payload per credential and the protocol says to ignore all but the first intent, so `payloads[0]` is read. The risk: an `xfer-full` for *another* payload would start an empty diff --git a/packages/client/src/launchdarkly_ai_server/skills_fdv2.py b/packages/client/src/launchdarkly_ai_server/skills_fdv2.py index b9bf96b..f4f7ade 100644 --- a/packages/client/src/launchdarkly_ai_server/skills_fdv2.py +++ b/packages/client/src/launchdarkly_ai_server/skills_fdv2.py @@ -655,6 +655,10 @@ def _payload_transferred(self, data: Any) -> _TransferOutcome: # Checked even with no pending set (a ``none`` intent), so a foreign # payload's selector never becomes the resume point. foreign = self._is_foreign_payload(payload_id) + # A ``none`` intent builds no pending set, nor does an intent code this + # SDK does not recognise, and a foreign payload's contents are declined + # below. + applied = not foreign and self._pending is not None if foreign: self._warn_foreign_payload(payload_id) self.diagnostics.payloads_ignored += 1 @@ -685,12 +689,18 @@ def _payload_transferred(self, data: Any) -> _TransferOutcome: version, len(self._committed), ) + if not applied: + # Nothing applied, so nothing to report — and in particular not a + # commit, because a commit publishes the first payload, and + # ``is_initialized`` (what ``write_skills("*")`` authorises a prune + # on) must not go true over a store that received nothing. Not up to + # date either: only the ``none`` intent says that, on its own event. + # The selector is withheld too — it names a payload never applied. + return _TransferOutcome() return _TransferOutcome( committed=True, changes=changes, - # A declined payload must not move the resume point, or skill - # updates could silently stop arriving. - basis=state if not foreign and isinstance(state, str) and state else None, + basis=state if isinstance(state, str) and state else None, ) def _abandon_in_flight(self) -> None: @@ -1485,9 +1495,9 @@ def wait_for_skills(self, timeout: float = 10.0) -> bool: """ Blocks until the first payload arrives, or *timeout* seconds elapse. - Returns ``True`` once a payload has committed or a 304 confirmed the one - held is current. That does not mean any skill verified, or that the - environment has skills; see ``diagnostics``. + Returns ``True`` once a payload has committed; a 304 alone does not + count. That does not mean any skill verified, or that the environment + has skills; see ``diagnostics``. Returns ``False`` on timeout, or early if delivery ends first (``close``, or a failure that will not be retried). @@ -1702,10 +1712,14 @@ def _give_up(self, reason: str) -> None: reason, ) - def _apply(self, name: str, data: Any) -> None: + def _apply(self, name: str, data: Any) -> bool: """ Feeds one event to the reader, publishes a commit, and raises the transport error the event calls for, if any. + + Returns whether the event completed an exchange — a commit, or the + server confirming the content held is current — which is what + ``_poll_once`` adopts an etag on. """ with self._lock: outcome = self._reader.handle(name, data) @@ -1723,6 +1737,7 @@ def _apply(self, name: str, data: Any) -> None: raise _FatalTransportError(outcome.fatal) if outcome.disconnect: raise _RecoverableTransportError(outcome.disconnect) + return outcome.committed or outcome.up_to_date def _poll_once(self) -> None: with self._lock: @@ -1733,15 +1748,22 @@ def _poll_once(self) -> None: result = self._requester.poll(basis, etag) if result.not_modified: logger.debug("Skill payload unchanged (HTTP 304)") - # A 304 counts as a first payload: the etag belongs to a body this - # store applied in full. - self._publish_first_payload() + # A 304 confirms the payload this store holds; it cannot establish + # one. ``is_initialized`` stays false until something commits, so a + # store that has received nothing never authorises a prune of the + # files on disk. ``_run`` counts the poll as a healthy answer. return + completed = False for name, data in result.events: - self._apply(name, data) + completed = self._apply(name, data) or completed + if not completed: + # An etag describes the body it came with, so it is adopted only when + # that body is also what the store now holds: a commit, or a ``none`` + # intent. A transfer this SDK could not apply is neither, and keeping + # its etag would let the next 304 confirm content never applied. Any + # etag already held stays valid, so it is left alone, not cleared. + return with self._lock: - # Adopted only after the whole body applied, so a body that broke - # off partway cannot earn a later 304. self._etag = result.etag self._etag_basis = basis diff --git a/packages/client/tests/test_skills_fdv2.py b/packages/client/tests/test_skills_fdv2.py index b9caf25..bfe5a7c 100644 --- a/packages/client/tests/test_skills_fdv2.py +++ b/packages/client/tests/test_skills_fdv2.py @@ -922,11 +922,19 @@ def test_a_catastrophic_goodbye_is_fatal(self) -> None: ) assert outcome.fatal is not None - def test_transfer_none_holds_everything_and_commits(self) -> None: + def test_transfer_none_holds_everything_and_commits_nothing(self) -> None: + """ + ``none`` is the server saying the payload held is current, so nothing is + applied and nothing is committed. It is a *completed exchange* — which + is what breaks a row of failures, and what lets the next request offer + the body's etag — but it is not a payload. Publishing one would make + ``is_initialized`` true over whatever the store happens to hold, and + that is the fact ``write_skills("*")`` prunes on. + """ held = _SkillObjectSet() reader = _ProtocolReader(held) drive(reader, full_payload(("put-object", put_skill()))) - drive( + intent, outcome = drive( reader, events( ("server-intent", server_intent("none")), @@ -934,6 +942,48 @@ def test_transfer_none_holds_everything_and_commits(self) -> None: ), ) assert len(held) == 1 + assert intent.up_to_date is True + assert outcome.committed is False + # The intent event above already carried the up-to-date answer; the + # transfer completing it adds nothing to report. + assert outcome.up_to_date is False + # Nor does it move the resume point. ``basis-2`` names a payload this + # store was never sent, and resuming from it would ask every later + # connection for changes since a payload it never applied. + assert outcome.basis is None + + def test_a_transfer_that_applied_nothing_is_not_a_commit(self) -> None: + """ + Three shapes reach ``payload-transferred`` with no pending set to apply: + an intent code this SDK does not recognise, a ``none`` intent, and a + lone transfer under no intent at all. None of them applied anything, so + none of them claims anything — neither a commit, which is what publishes + the first payload ``write_skills("*")`` prunes on, nor an up-to-date + answer, which only the server can give and only the ``none`` intent + does, on its own event. + + The transfer is still a wire fact, counted either way. + """ + for payload_events in ( + events( + ("server-intent", server_intent("xfer-future")), + ("put-object", put_skill()), + ("payload-transferred", transferred("basis-1")), + ), + events( + ("server-intent", server_intent("none")), + ("payload-transferred", transferred("basis-1")), + ), + events(("payload-transferred", transferred("basis-1"))), + ): + held = _SkillObjectSet() + reader = _ProtocolReader(held) + outcome = drive(reader, payload_events)[-1] + assert outcome.committed is False + assert outcome.up_to_date is False + assert outcome.basis is None + assert len(held) == 0 + assert reader.diagnostics.payloads_transferred == 1 def test_an_object_arriving_with_no_intent_is_treated_as_a_delta(self) -> None: held = _SkillObjectSet() @@ -1445,15 +1495,68 @@ def test_a_304_keeps_the_held_content(self, endpoint: Any) -> None: assert store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") is not None assert store.diagnostics.payloads_transferred == 1 assert store.failed is None + # A payload arrived, and a 304 does not take that back. + assert store.is_initialized() is True - def test_a_304_before_any_payload_still_releases_wait_for_skills( + def test_a_304_confirms_a_payload_but_cannot_establish_one( self, endpoint: Any ) -> None: - """A reconnect with a cached basis has nothing to transfer; boot must not - block on a payload the server has no reason to send.""" + """ + A 304 answers for content this store already holds, and the etag that + asked for it is only ever adopted from a body that completed an + exchange. Reaching one with nothing held therefore takes a server + answering a request that carried no etag at all, and that 304 says + nothing about a payload this store never received. + + Releasing ``wait_for_skills`` on it would make ``is_initialized`` true + over an empty committed set — the fact ``write_skills("*")`` prunes on — + so a reconcile that raced delivery would delete every managed skill on + disk instead of reporting the retrieval unavailable (§3.21). Failing + closed costs a boot that is genuinely waiting nothing it was not already + waiting for. + """ endpoint.queue_poll(status=304) with poll_store(endpoint) as store: - assert store.wait_for_skills(timeout=5) is True + assert store.wait_for_skills(timeout=0.5) is False + assert store.is_initialized() is False + # Not a failure either: the poll was answered, and the store is + # still asking. + assert store.failed is None + assert endpoint.requests[0]["if_none_match"] is None + + def test_an_intent_it_cannot_apply_does_not_lend_its_etag_to_a_304( + self, endpoint: Any + ) -> None: + """ + The chain this closes: a body under a future intent code announces and + transfers a payload this SDK cannot apply, its etag is adopted as though + the body had been applied in full, and the next 304 reports the empty + store it left behind as current. That store is initialized, healthy, and + authorises a prune of every managed skill on disk, with nothing in the + 304 to notice it on. + + An etag is adopted only from a body that completed an exchange, so the + second request carries none and the endpoint's standing 304 cannot + answer for content that never arrived. + """ + endpoint.queue_poll( + events( + ("server-intent", server_intent("xfer-future")), + ("put-object", put_skill()), + ("payload-transferred", transferred("basis-1")), + ), + etag='W/"v1"', + ) + with poll_store(endpoint) as store: + assert wait_until(lambda: len(endpoint.requests) >= 2) + assert store.is_initialized() is False + assert store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") is None + assert endpoint.requests[1]["if_none_match"] is None + # Nor is the selector of a payload it could not apply a resume point. + assert [r["query"].get("basis") for r in endpoint.requests[:2]] == [ + None, + None, + ] def test_a_mixed_payload_over_the_wire_yields_only_the_skill( self, endpoint: Any @@ -2648,6 +2751,7 @@ def test_an_over_cap_poll_body_is_not_applied_and_is_retried( ) -> None: monkeypatch.setattr(skills_fdv2, "MAX_RESPONSE_BYTES", 2048) endpoint.queue_poll(full_payload(("put-object", put_skill(content="x" * 8192)))) + endpoint.queue_poll(full_payload(("put-object", put_skill()))) # A long enough backoff to observe the failure before the retry lands. with poll_store(endpoint, initial_backoff=0.3, max_backoff=0.3) as store: assert wait_until(lambda: store.diagnostics.connection_failures == 1) @@ -2656,7 +2760,10 @@ def test_an_over_cap_poll_body_is_not_applied_and_is_retried( assert store.diagnostics.payloads_transferred == 0 assert store.diagnostics.skill_objects_received == 0 assert store.failed is None - # The retry is an ordinary poll; the endpoint answers it 304. + # The retry is an ordinary poll, and the payload it is answered + # with is what releases the waiter. A 304 could not: nothing has + # committed, and a 304 confirms a payload rather than establishing + # one. assert wait_until(lambda: len(endpoint.requests) >= 2) assert store.wait_for_skills(timeout=5) is True assert wait_until(lambda: store.diagnostics.connection_failures == 0)