Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions packages/client/agents.md
Original file line number Diff line number Diff line change
Expand Up @@ -353,6 +353,23 @@ described, and would briefly empty the store — which, with pruning on, is the
between a reconcile and deleting a customer's skill files. An interrupted transfer therefore
leaves last known good intact, 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 is assumed to be the skill payload.** Delivery
provides one payload per credential and the protocol requires a client to ignore all but the
first payload intent, so `payloads[0]` is both what arrives and what the protocol says to
Expand Down
50 changes: 32 additions & 18 deletions packages/client/src/launchdarkly_ai_server/skills_fdv2.py
Original file line number Diff line number Diff line change
Expand Up @@ -795,6 +795,10 @@ def _payload_transferred(self, data: Any) -> _TransferOutcome:
# whose selector must not become the resume point if it is not the
# payload skills arrive on.
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
Expand Down Expand Up @@ -833,15 +837,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. Adopting the
# selector of a transfer whose contents this layer just threw away
# would ask the next poll or stream to resume from someone else's
# payload, and skill updates could stop arriving while every
# diagnostic still read healthy.
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:
Expand Down Expand Up @@ -2060,10 +2067,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)
Expand All @@ -2085,6 +2096,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:
Expand All @@ -2099,20 +2111,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, so a boot that reconnects with a
# cached basis is not blocked on a transfer the server will not
# send. It is a current answer because the etag that asked for it
# was issued for 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 once the whole body has been applied. A body that
# broke off partway — an ``error`` or ``goodbye`` after an announced
# transfer — left the payload it described unapplied, and keeping
# its etag would let the next 304 report a store that is missing
# that payload as current and healthy.
self._etag = result.etag
self._etag_basis = basis

Expand Down
121 changes: 114 additions & 7 deletions packages/client/tests/test_skills_fdv2.py
Original file line number Diff line number Diff line change
Expand Up @@ -922,18 +922,68 @@ 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")),
("payload-transferred", transferred("basis-2")),
),
)
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()
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand Down
Loading