diff --git a/packages/client/README.md b/packages/client/README.md index 4ab0a15..1337c8f 100644 --- a/packages/client/README.md +++ b/packages/client/README.md @@ -628,7 +628,7 @@ HTTP 422 when it will not serve one. - **Not the empty case:** an environment with zero skills is served an empty payload that commits normally. - **Recovery:** once the cause is fixed, call `start()` on the same store. Only `close()` is - final: a restart clears `failed`, resets the retry budget, and held content stays readable + final: a restart clears `failed`, resets the backoff, and held content stays readable throughout. Restarting the process also works. **Nothing above the store changes.** The accessors, verification, and `write_skills` see raw @@ -651,15 +651,23 @@ never followed, so a 3xx stops delivery instead of forwarding the key to the `Lo Construct a new store to resume. **Reads are memory-bounded.** A poll body or streamed event larger than `MAX_RESPONSE_BYTES` -(64 MiB) is dropped without being applied; the store keeps serving what it last held, and -delivery retries on its normal backoff. +(64 MiB) is dropped without being applied, and delivery stops: the payload's size belongs to the +environment, so a retry would download it again and be refused the same way. `failed` carries +the reason, the store keeps serving what it last held, and `start()` resumes delivery once the +payload is back under the bound. **Streaming is the default, and it is what makes revocation fast.** A `delete-object` reaches a live stream in seconds; with `mode="poll"` it arrives within one `poll_interval`. With -`watch_skills`, a revoked skill's `SKILL.md` leaves the disk without a restart. During an -outage the store keeps serving its last content, and `write_skills`' default -`on_unavailable="keep"` leaves managed files alone, so an outage does not read as "everything -was revoked". +`watch_skills("*", ...)`, a revoked skill's `SKILL.md` leaves the disk without a restart. During +an outage the store keeps serving its last content and retries for as long as it runs, and +`write_skills`' default `on_unavailable="keep"` leaves managed files alone, so an outage does +not read as "everything was revoked". + +**With an explicit skill list, revocation does not reach the disk.** Given a list such as +`skill_refs(config)`, a skill deleted in LaunchDarkly is reported as an `error` action under +its key, and its `SKILL.md` is kept rather than pruned. The watcher also listens only to the +skill store, not to flag changes, so unpinning a skill from a variation is not seen either. Use +`"*"` when revocation must reach the disk. **Without the watcher, the revocation bound is process lifetime.** If you call `write_skills` once at boot and never run `watch_skills`, a skill revoked after boot stays on disk, and in the @@ -697,7 +705,7 @@ that skips verification. | `write_skills(skills, root, *, prune=True, timeout=10.0, on_unavailable="keep")` | Materialize skills under `root`, returning a `ReconcileReport`. `prune` removes formerly-managed skills no longer requested. `on_unavailable="raise"` raises instead of reporting when content cannot be retrieved. Raises `ValueError` for an unusable root, a negative or non-finite `timeout`, or an unrecognised `on_unavailable`. `timeout=0` is valid and makes every skill report an error. **Performs synchronous filesystem I/O — see the note below.** | | `SkillStore` | The structural interface content arrives through: `get_object(kind, key, version=None)`, `all_objects(kind)`, optional `is_initialized()`, `add_listener(kind, fn)` / `remove_listener(kind, fn)`. A store without `is_initialized()` is treated as initialized. Both shipped stores deliver only the skill kind, so `add_listener` on any other kind raises. | | `InMemorySkillStore(objects=None)` | A dict-backed store with `put(raw)`, for local development and testing. Holds several versions of a key. | -| `FDv2SkillStore(sdk_key, *, base_uri=…, stream_uri=…, mode="stream", …)` | The delivery transport: a store fed by LaunchDarkly over the SDK-facing FDv2 channel. `start()`, `wait_for_skills(timeout)`, `is_initialized()`, `close()`, `diagnostics`, `failed`; also a context manager. `close()` is **final** — `start()` afterwards raises. `poll_interval` and `read_timeout` must be positive and finite. **Server-side only.** See *Receiving skills from LaunchDarkly* above. | +| `FDv2SkillStore(sdk_key, *, base_uri=…, stream_uri=…, mode="stream", …)` | The delivery transport: a store fed by LaunchDarkly over the SDK-facing FDv2 channel. `start()`, `wait_for_skills(timeout)`, `is_initialized()`, `close()`, `diagnostics`, `failed`; also a context manager. `close()` is **final** — `start()` afterwards raises. `poll_interval`, `read_timeout`, `initial_backoff` and `max_backoff` must be positive and finite, and `initial_backoff` may not exceed `max_backoff`. **Server-side only.** See *Receiving skills from LaunchDarkly* above. | | `watch_skills(skills, root, *, debounce=0.5, on_reconcile=None, …)` | `write_skills` plus a re-reconcile on every delivery change, so revocation takes effect within `debounce` rather than at the next restart. Returns `(initial report, SkillWatcher)`; close the watcher when done. `debounce` is in **seconds**, non-negative and finite. `on_reconcile` receives each *subsequent* report. One watcher per root. | | `StoreDiagnostics` | What the transport has seen: `payloads_transferred`, `skill_objects_received`, `objects_ignored`, `objects_revoked`, `payloads_ignored`, `hashless_objects`, `connection_failures`, `last_error`. | diff --git a/packages/client/agents.md b/packages/client/agents.md index b867387..1c56990 100644 --- a/packages/client/agents.md +++ b/packages/client/agents.md @@ -273,7 +273,7 @@ asserts both the cause and the absences. **A fatal stops the run, not the store, so every surface says `start()`, not "restart the process".** `_give_up` does not close; only `close` sets `_closed`, the one thing `start` -refuses. A store that gave up — on a 401, 403, 404, 422, or an exhausted retry budget — +refuses. A store that gave up — on a 401, 403, 404, 422, or another fatal status — resumes in place once the cause is fixed, clearing the terminal reason through `_rearm_waiters`. Asserted by `test_the_give_up_line_points_at_start_not_a_process_restart`, `test_a_store_that_gave_up_on_a_422_resumes_on_start`, and @@ -290,10 +290,12 @@ resumes in place once the cause is fixed, clearing the terminal reason through **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 `_RecoverableTransportError`: nothing from that body or event -is applied, the in-flight payload is abandoned, the failure is recorded in -`connection_failures`/`last_error`, and delivery retries on the usual backoff while the -committed set stays served. This bound is independent of +(64 MiB). Crossing it raises `_ResponseTooLargeError`, a fatal error: nothing from that +body or event is applied, the in-flight payload is abandoned, `failed` and `last_error` are +set, and the committed set stays served. It is fatal because the size belongs to the +environment, not the connection: retried, it would re-download up to 64 MiB on every backoff +step forever. The requester wrappers re-raise fatal errors unchanged; do not let a generic +`except Exception` turn one back into a recoverable error. This bound is independent of `skills_core.MAX_SKILL_CONTENT_BYTES` (one skill's content, at verification); do not derive one from the other. @@ -720,7 +722,7 @@ create false confidence. The operator's verification steps are in the README. `timeout` is a monotonic deadline, checked before each retrieval, each write, and each prune; only the final manifest rewrite runs past it, so files already written are never orphaned. Bounded retries are **not** implemented at this layer; retry policy belongs to the delivery -transport (`FDv2SkillStore`'s backoff and retry budget). Why: +transport (`FDv2SkillStore`'s capped backoff). Why: 1. **There is nothing transient to retry.** `SkillStore.get_object` is a synchronous in-process read of already-delivered data, modelled on the LaunchDarkly data-store API. A @@ -882,6 +884,12 @@ true. Tampered content must never trigger deletion. ### 6. Expecting revocation to reach a boot-only `write_skills` deployment +Even with `watch_skills`, only the `"*"` form removes a revoked skill from disk. With an +explicit list such as `skill_refs(config)`, the list is fixed: a skill the store answers +`absent` for is reported as an `error` and its files are kept, and the watcher listens only to +the skill store, so unpinning a skill or moving it to a new version is not seen until the refs +are read again. + Without `watch_skills`, the revocation bound is process lifetime: a skill revoked after boot stays on disk until the process reconciles again, so a restart (or an explicit re-run of `write_skills`) is the incident-response action — and content an agent has already read into diff --git a/packages/client/src/launchdarkly_ai_server/skills_fdv2.py b/packages/client/src/launchdarkly_ai_server/skills_fdv2.py index b9bf96b..d26856c 100644 --- a/packages/client/src/launchdarkly_ai_server/skills_fdv2.py +++ b/packages/client/src/launchdarkly_ai_server/skills_fdv2.py @@ -25,6 +25,7 @@ import re import socket import threading +import time import urllib.error import urllib.parse import urllib.request @@ -86,11 +87,19 @@ """The most the transport holds in memory from one poll body or one streamed event. A memory backstop set far above any real payload, separate from the per-skill -content limit verification enforces. Crossing it is a recoverable failure: -nothing is applied, the store keeps its current content, and delivery retries.""" +content limit verification enforces. Crossing it is fatal: nothing is applied, +the store keeps its current content, and ``failed`` is set. The payload's size +belongs to the environment, so a retry would be refused the same way, after +downloading up to this much again.""" _READ_CHUNK_BYTES = 64 * 1024 +_BACKOFF_RESET_INTERVAL = 60.0 +""" +Seconds a stream must stay open before its end resets the backoff delay, as in +the base server-side SDKs. Overridden in tests. +""" + _EVENT_SERVER_INTENT = "server-intent" _EVENT_PUT_OBJECT = "put-object" _EVENT_DELETE_OBJECT = "delete-object" @@ -478,6 +487,9 @@ class _TransferOutcome: up_to_date: bool = False """A ``none`` intent: the content held is current. Counts as a healthy exchange though it commits nothing, like a 304 to a poll.""" + recycled: bool = False + """The disconnect is a ``goodbye``: routine if the connection has already + completed an exchange.""" def _identity_of(raw: dict[str, Any]) -> tuple[str, Any]: @@ -589,12 +601,17 @@ def _server_intent(self, data: Any) -> _TransferOutcome: self._pending = _SkillObjectSet() elif intent == _INTENT_TRANSFER_CHANGES: self._pending = self._committed.copy() + elif intent == _INTENT_TRANSFER_NONE: + # Current, but not finished: later edits arrive on this connection + # with no second intent, so expect changes, as the base SDK's + # ``ChangeSetBuilder.expect_changes()`` does. The pending set is + # copied when the first object arrives. + self._intent = _INTENT_TRANSFER_CHANGES + self._pending = None + return _TransferOutcome(up_to_date=True) else: - if intent != _INTENT_TRANSFER_NONE: - logger.debug("Ignoring FDv2 server-intent with intentCode %r", intent) + logger.debug("Ignoring FDv2 server-intent with intentCode %r", intent) self._pending = None - # Only ``none`` is an answer; an unrecognised intent is not. - return _TransferOutcome(up_to_date=intent == _INTENT_TRANSFER_NONE) return _TransferOutcome() def _target_for(self, data: Any) -> _SkillObjectSet | None: @@ -712,12 +729,16 @@ def _goodbye(self, data: Any) -> _TransferOutcome: silent = bool(data.get("silent")) if isinstance(data, dict) else False self._abandon_in_flight() if not silent: - logger.info("FDv2 connection closing: %s", reason) + # 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}") + return _TransferOutcome( + disconnect=f"server said goodbye: {reason}", recycled=True + ) # -- payload identity ---------------------------------------------------- @@ -824,12 +845,25 @@ class _FatalTransportError(Exception): """A failure retrying cannot fix: bad credential, forbidden, wrong URI.""" +class _ResponseTooLargeError(_FatalTransportError): + """A poll body, stream line or stream event over ``MAX_RESPONSE_BYTES``.""" + + class _RecoverableTransportError(Exception): """A failure worth retrying. Carries a server-requested delay when given one.""" - def __init__(self, message: str, retry_after: float | None = None) -> None: + def __init__( + self, + message: str, + retry_after: float | None = None, + *, + recycled: bool = False, + ) -> None: super().__init__(message) self.retry_after = retry_after + self.recycled = recycled + """A ``goodbye`` after a completed exchange: the server recycling the + stream. Reconnected, but not counted or reported as a failure.""" class _StaleRequestStateError(_RecoverableTransportError): @@ -1081,7 +1115,7 @@ def poll(self, basis: str | None, etag: str | None) -> _PollResult: # urllib raises on 304; here it means "unchanged". return _PollResult(not_modified=True, events=[], etag=etag) raise _classify_status(exc.code, exc.headers) from exc - except _RecoverableTransportError: + except (_RecoverableTransportError, _FatalTransportError): raise except Exception as exc: raise _RecoverableTransportError( @@ -1124,7 +1158,7 @@ def _read_bounded(response: Any, limit: int) -> bytes: return b"".join(chunks) seen += len(chunk) if seen > limit: - raise _RecoverableTransportError( + raise _ResponseTooLargeError( f"polling response exceeded the {limit}-byte transport bound " f"(at least {seen} bytes received); nothing from it was applied" ) @@ -1169,13 +1203,13 @@ def _iter_stream_lines(response: Any, limit: int) -> Any: if not line: return if len(line) > limit: - raise _RecoverableTransportError( + raise _ResponseTooLargeError( f"an FDv2 stream line exceeded the {limit}-byte transport " "bound; the connection was dropped and nothing from the " "in-flight payload was applied" ) yield line - except _RecoverableTransportError: + except (_RecoverableTransportError, _FatalTransportError): raise except Exception as exc: raise _RecoverableTransportError( @@ -1189,8 +1223,8 @@ def _iter_sse(response: Any) -> Any: Minimal: ``event:``/``data:`` fields, multi-line ``data`` joined with newlines, blank line dispatches, ``:`` comments skipped. An event over - ``MAX_RESPONSE_BYTES`` raises a recoverable error and the in-flight payload - is abandoned. + ``MAX_RESPONSE_BYTES`` raises a fatal error and the in-flight payload is + abandoned. """ limit = MAX_RESPONSE_BYTES try: @@ -1219,7 +1253,7 @@ def _iter_sse(response: Any) -> Any: continue event_bytes += len(raw_line) if event_bytes > limit: - raise _RecoverableTransportError( + raise _ResponseTooLargeError( f"an FDv2 stream event exceeded the {limit}-byte transport " f"bound (at least {event_bytes} bytes received); the " "connection was dropped and nothing from the in-flight " @@ -1251,8 +1285,11 @@ def _backoff_delay( Jitter is subtracted, never added, so *maximum* is a true ceiling. """ - # float(2 ** n): the integer power is untyped to mypy. - ceiling: float = min(maximum, base * float(2 ** max(0, attempt - 1))) + # Clamped: retries are unbounded, and ``float(2 ** n)`` overflows past + # about 1024. ``2 ** 62`` exceeds any real cap. ``float(...)``: the integer + # power is untyped to mypy. + exponent = min(max(0, attempt - 1), 62) + ceiling: float = min(maximum, base * float(2**exponent)) return ceiling * (1.0 - jitter * random.random()) @@ -1307,7 +1344,6 @@ def __init__( read_timeout: float | None = None, initial_backoff: float = 1.0, max_backoff: float = 30.0, - max_consecutive_failures: int = 10, _requester: Any = None, ) -> None: """ @@ -1326,10 +1362,12 @@ def __init__( ``"poll"`` mode it bounds the whole request (``DEFAULT_POLL_TIMEOUT``); in ``"stream"`` mode, each wait for more bytes (``DEFAULT_STREAM_READ_TIMEOUT``). - - *max_backoff*: caps every retry delay, including ``Retry-After``. - - *max_consecutive_failures*: after this many failures in a row, - delivery stops, logs an error, and ``failed`` is set; the store keeps - serving last known good. A completed exchange resets the count. + - *initial_backoff*: the first retry delay; positive and finite. + - *max_backoff*: caps every retry delay, including ``Retry-After``; + positive, finite, and at least *initial_backoff*. + + Recoverable failures are retried for as long as the store runs; only a + fatal status stops delivery and sets ``failed``. """ _require_server_side_credential(sdk_key) # A lone ``base_uri`` serves both endpoints. @@ -1353,12 +1391,24 @@ def __init__( ) elif not (math.isfinite(read_timeout) and read_timeout > 0): raise ValueError(f"read_timeout must be positive, got {read_timeout!r}") + # With no failure bound these are the only limit on the retry loop: a + # zero or negative value reconnects as fast as the network allows. + for option, value in ( + ("initial_backoff", initial_backoff), + ("max_backoff", max_backoff), + ): + if not (math.isfinite(value) and value > 0): + raise ValueError(f"{option} must be positive and finite, got {value!r}") + if initial_backoff > max_backoff: + raise ValueError( + f"initial_backoff ({initial_backoff!r}) must not exceed " + f"max_backoff ({max_backoff!r})" + ) self._mode: Mode = mode self._poll_interval = poll_interval self._initial_backoff = initial_backoff self._max_backoff = max_backoff - self._max_consecutive_failures = max_consecutive_failures self._objects = _SkillObjectSet() self._reader = _ProtocolReader(self._objects) @@ -1393,6 +1443,12 @@ def __init__( # Recoverable failures in a row; reset by a completed exchange, not by a # connection ending (a stream only ends by being dropped). self._failures = 0 + # The backoff's attempt number, kept apart from ``_failures``: reset + # only by a stream that stayed open ``_BACKOFF_RESET_INTERVAL``, or by + # a completed poll, so a server that answers and drops is backed off. + self._backoff_attempts = 0 + # When the current stream connected, for the reset above. + self._connected_at: float | None = None # Whether the current attempt got a complete answer, which separates a # recycled healthy stream from a failed one. self._attempt_answered = False @@ -1407,7 +1463,7 @@ def start(self) -> FDv2SkillStore: Raises ``RuntimeError`` if the store has been closed. A store whose delivery stopped on its own (``failed`` is set) can be started again, - with a fresh retry budget. + with its backoff reset. """ with self._lock: if self._closed: @@ -1437,12 +1493,13 @@ def start(self) -> FDv2SkillStore: def _rearm_waiters(self) -> None: """ Resets per-run state for a store being started again after it gave up: - the ended-delivery flag, ``failed``, and the retry budget. A payload + the ended-delivery flag, ``failed``, and the failure count. A payload already held still answers ``wait_for_skills``. Call with the lock held. """ self._delivery_ended.clear() self._failed_reason = None self._failures = 0 + self._backoff_attempts = 0 if not self._first_payload.is_set(): self._released.clear() @@ -1609,6 +1666,8 @@ def _run(self) -> None: # A returned poll is a current answer, even a 304. Stream # successes are recorded in ``_apply``. self._record_success() + with self._lock: + self._backoff_attempts = 0 except _FatalTransportError as exc: self._give_up(str(exc)) return @@ -1616,7 +1675,6 @@ def _run(self) -> None: if self._stop.is_set(): # ``close`` interrupted the request; not a failure. return - repairing_state = False if isinstance(exc, _StaleRequestStateError): # With no basis or etag to drop, the 400 is fatal. Otherwise # drop them and request a full transfer once. @@ -1626,30 +1684,32 @@ def _run(self) -> None: self._basis = None self._etag = None self._etag_basis = None - repairing_state = True if exhausted: self._give_up(str(exc)) return with self._lock: # Discard any partial payload from the dropped connection. self._reader._abandon_in_flight() - self._failures += 1 - failures = self._failures answered = self._attempt_answered - self._reader.diagnostics.connection_failures = failures - self._reader.diagnostics.last_error = str(exc) - if failures > self._max_consecutive_failures and not repairing_state: - # The one stateless retry after a 400 is exempt, so an - # exhausted budget cannot block that repair. - self._give_up( - f"gave up after {failures} consecutive failures; " - f"last error: {exc}" - ) - return + # A goodbye after a completed exchange is a routine recycle. + recycled = exc.recycled and answered + if not recycled: + self._failures += 1 + self._reader.diagnostics.connection_failures = self._failures + self._reader.diagnostics.last_error = str(exc) + connected_at = self._connected_at + self._connected_at = None + if ( + connected_at is not None + and time.monotonic() - connected_at >= _BACKOFF_RESET_INTERVAL + ): + self._backoff_attempts = 0 + self._backoff_attempts += 1 + attempt = self._backoff_attempts delay = exc.retry_after if delay is None or not math.isfinite(delay): delay = _backoff_delay( - failures, base=self._initial_backoff, maximum=self._max_backoff + attempt, base=self._initial_backoff, maximum=self._max_backoff ) else: # Floor at ``initial_backoff`` so ``Retry-After: 0`` cannot @@ -1713,7 +1773,7 @@ def _apply(self, name: str, data: Any) -> None: self._basis = outcome.basis if outcome.committed or outcome.up_to_date: # Both are completed exchanges. Counting ``up_to_date`` keeps a - # stream for an unchanging environment from exhausting its budget. + # stream for an unchanging environment from reading as failing. self._record_success() if outcome.committed: self._publish_first_payload() @@ -1722,7 +1782,9 @@ def _apply(self, name: str, data: Any) -> None: if outcome.fatal: raise _FatalTransportError(outcome.fatal) if outcome.disconnect: - raise _RecoverableTransportError(outcome.disconnect) + raise _RecoverableTransportError( + outcome.disconnect, recycled=outcome.recycled + ) def _poll_once(self) -> None: with self._lock: @@ -1751,6 +1813,7 @@ def _stream_once(self) -> None: connection = self._requester.stream(basis) with self._lock: self._connection = connection + self._connected_at = time.monotonic() try: # ``close`` may have run during the connect, before there was a # connection to interrupt. diff --git a/packages/client/src/launchdarkly_ai_server/skills_fs.py b/packages/client/src/launchdarkly_ai_server/skills_fs.py index e5323ae..2d7bf29 100644 --- a/packages/client/src/launchdarkly_ai_server/skills_fs.py +++ b/packages/client/src/launchdarkly_ai_server/skills_fs.py @@ -377,9 +377,10 @@ def _resolve_requests( """ Turns the caller's input into one request per skill. - Also returns whether any retrieval was incomplete (absent, uninitialized or + Also returns whether any retrieval was incomplete (an uninitialized or raising store, or an exhausted timeout). That flag suppresses pruning, so an - outage never deletes managed files. + outage never deletes managed files. An ``absent`` reference does not set it: + the key stays requested, carrying an error, so prune keeps its files. """ if isinstance(skills, str): if skills != "*": diff --git a/packages/client/src/launchdarkly_ai_server/skills_watch.py b/packages/client/src/launchdarkly_ai_server/skills_watch.py index 0475f1c..452c3ee 100644 --- a/packages/client/src/launchdarkly_ai_server/skills_watch.py +++ b/packages/client/src/launchdarkly_ai_server/skills_watch.py @@ -2,9 +2,11 @@ Agent Skills — keep skills on disk in sync as delivery changes. ``write_skills`` is a one-shot reconcile of what the store holds now. -``watch_skills`` re-runs it whenever the store reports a change, so a skill -revoked in LaunchDarkly is removed from disk within a debounce interval rather -than at the next restart. +``watch_skills`` re-runs it whenever the store reports a change. With ``"*"``, +a skill revoked in LaunchDarkly is removed from disk within a debounce interval +rather than at the next restart. An explicit list is fixed: a skill the store +answers ``absent`` for keeps its files and reports an error, and a config change +that unpins a skill or moves it to a new version is not seen. ``on_unavailable="keep"`` remains the default, so an outage never deletes the application's skill files. diff --git a/packages/client/tests/test_skills_fdv2.py b/packages/client/tests/test_skills_fdv2.py index b9caf25..3e31531 100644 --- a/packages/client/tests/test_skills_fdv2.py +++ b/packages/client/tests/test_skills_fdv2.py @@ -22,6 +22,7 @@ import hashlib import inspect import json +import math import socket import threading import time @@ -65,6 +66,7 @@ _ProtocolReader, _RecoverableTransportError, _Requester, + _ResponseTooLargeError, _retry_after_seconds, _SkillObjectSet, _StaleRequestStateError, @@ -935,6 +937,46 @@ def test_transfer_none_holds_everything_and_commits(self) -> None: ) assert len(held) == 1 + def test_a_put_after_a_none_intent_is_applied(self) -> None: + """``none`` means current, not finished: later edits follow it on the + same connection with no second intent, and must land.""" + held = _SkillObjectSet() + reader = _ProtocolReader(held) + drive( + reader, + events( + ("server-intent", server_intent("none")), + ("put-object", put_skill(object_version=4)), + ("payload-transferred", transferred("basis-2")), + ), + ) + assert held.get("pdf-extraction", 4) is not None + assert reader.diagnostics.objects_ignored == 0 + + def test_a_delete_after_a_none_intent_revokes(self) -> None: + """The case that matters: a skill revoked after a routine reconnect. + + Dropping the delete while adopting the selector would keep serving the + revoked skill, and a reconnect would not recover it, because the basis + has already moved past the change. + """ + held = _SkillObjectSet() + reader = _ProtocolReader(held) + drive(reader, full_payload(("put-object", put_skill()))) + outcomes = drive( + reader, + events( + ("server-intent", server_intent("none")), + ("delete-object", delete_skill()), + ("payload-transferred", transferred("basis-2")), + ), + ) + assert len(held) == 0 + assert reader.diagnostics.objects_revoked == 1 + assert reader.diagnostics.objects_ignored == 0 + assert outcomes[-1].changes == [{"key": "pdf-extraction", "version": 3}] + assert outcomes[-1].basis == "basis-2" + def test_an_object_arriving_with_no_intent_is_treated_as_a_delta(self) -> None: held = _SkillObjectSet() reader = _ProtocolReader(held) @@ -1538,6 +1580,43 @@ def test_the_stream_request_declares_the_payload_kind(self, endpoint: Any) -> No store.close() assert endpoint.requests[0]["query"]["kinds"] == "agent-skill" + def test_a_revocation_after_an_up_to_date_answer_arrives( + self, endpoint: Any + ) -> None: + """A reconnect answered ``none``, then a revocation on the same stream. + + The first connection ends after its payload, as a recycled stream does. + """ + endpoint.queue_stream(full_payload(("put-object", put_skill()))) + endpoint.queue_stream( + events( + ("server-intent", server_intent("none")), + ("delete-object", delete_skill()), + ("payload-transferred", transferred("basis-2")), + ) + ) + store = FDv2SkillStore( + SDK_KEY, + base_uri=endpoint.base_uri, + mode="stream", + initial_backoff=0.01, + max_backoff=0.02, + ) + notified: list[dict[str, Any]] = [] + store.add_listener(SKILL_OBJECT_KIND, notified.append) + try: + store.start() + assert store.wait_for_skills(timeout=5) is True + assert wait_until( + lambda: store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") is None + ) + assert store.diagnostics.objects_revoked == 1 + assert {"key": "pdf-extraction", "version": 3} in notified + finally: + store.close() + # The premise: the revocation came on a reconnect carrying a basis. + assert endpoint.requests[1]["query"].get("basis") == "basis-1" + def test_a_streamed_revocation_arrives_without_a_restart( self, endpoint: Any ) -> None: @@ -1813,6 +1892,19 @@ def stream(self, basis: str | None) -> Any: return _BlockingConnection() +def _record_backoff_attempts(monkeypatch: Any) -> list[int]: + """Records the attempt number of every backoff the store computes.""" + attempts: list[int] = [] + real = skills_fdv2._backoff_delay + + def recording(attempt: int, **kwargs: Any) -> float: + attempts.append(attempt) + return real(attempt, **kwargs) + + monkeypatch.setattr(skills_fdv2, "_backoff_delay", recording) + return attempts + + def stream_store(**kwargs: Any) -> FDv2SkillStore: return FDv2SkillStore( SDK_KEY, @@ -1843,8 +1935,7 @@ def test_the_give_up_line_points_at_start_not_a_process_restart( ``_give_up`` ends the run and not the store, and the test below is what proves a restarted store delivers. Telling an operator to restart their process is therefore an overstatement wherever it appears, and this - line appears on all of them — a 401, a 403, a 404, a 422, and an - exhausted retry budget alike. + line appears on all of them — a 401, a 403, a 404, and a 422 alike. """ endpoint.queue_poll(status=401) with caplog.at_level("ERROR"): @@ -1922,25 +2013,37 @@ def blocking_error(msg: Any, *args: Any, **kwargs: Any) -> None: assert store.wait_for_skills(timeout=5) is True assert store.failed is None - def test_a_restart_returns_the_retry_budget(self, endpoint: Any) -> None: - """The retry budget belongs to the run that spent it. + def test_a_restart_resets_the_failure_count(self) -> None: + """The failure count belongs to the run that accumulated it. - A store that gave up at its failure limit would otherwise carry the - spent count into the restarted run and give up again on its first - recoverable failure, without retrying once. + Carried into a restarted run, it would start that run's backoff part + way up and report failures the new run never had. """ - store = poll_store(endpoint, max_consecutive_failures=1) - # One over the limit, so the first run retries once and then gives up. - endpoint.queue_poll(status=500) - endpoint.queue_poll(status=500) + requester = _ScriptedRequester( + _RecoverableTransportError("x"), + _RecoverableTransportError("x"), + _FatalTransportError("401"), + ) + store = FDv2SkillStore( + SDK_KEY, + mode="poll", + initial_backoff=0.001, + max_backoff=0.002, + _requester=requester, + ) with store: + store.start() assert wait_until(lambda: store.failed is not None) + assert store.diagnostics.connection_failures == 2 - endpoint.queue_poll(status=500) - endpoint.queue_poll(full_payload(("put-object", put_skill()))) + requester.outcomes = [ + _RecoverableTransportError("x"), + _FatalTransportError("401"), + ] store.start() - assert store.wait_for_skills(timeout=5) is True - assert store.failed is None + assert wait_until(lambda: store.failed is not None) + # One failure in the new run, not three. + assert store.diagnostics.connection_failures == 1 def test_a_404_stops_delivery_immediately(self, endpoint: Any) -> None: """A 404 means the endpoint does not exist for this credential. @@ -1968,7 +2071,7 @@ def test_a_400_reconnects_once_from_scratch_and_is_then_fatal( The selector and the etag are the only client state the request carries, so a fresh connection built from nothing is the one repair available. - It gets exactly one: the retry bound is what keeps this from being + It gets exactly one: the one-retry limit is what keeps this from being "400 is recoverable". """ endpoint.queue_poll(full_payload(("put-object", put_skill()))) @@ -1992,49 +2095,6 @@ def test_a_400_reconnects_once_from_scratch_and_is_then_fatal( # Last known good survives both. assert store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") is not None - def test_a_400_still_reconnects_from_scratch_on_a_spent_budget( - self, endpoint: Any - ) -> None: - """The one repair available does not compete with the retry bound. - - A 400 arriving on a budget an outage has already spent would otherwise - give up while holding the one request known to fix it, and delivery - would stop for the process lifetime over state the store was about to - drop. Exempting that request cannot unbound the loop: it carries no - state, so a second 400 is fatal on its own. - """ - store = poll_store(endpoint, max_consecutive_failures=1) - endpoint.queue_poll(full_payload(("put-object", put_skill()))) - endpoint.queue_poll(status=500) - endpoint.queue_poll(status=400) - endpoint.queue_poll(full_payload(("put-object", put_skill()))) - with store: - assert store.wait_for_skills(timeout=5) is True - # The premise: the budget is spent by the time the 400 arrives. - assert wait_until(lambda: len(endpoint.requests) == 4) - assert store.failed is None - # The repair went out from scratch rather than never going out at all. - repair = endpoint.requests[3] - assert repair["query"] == {"kinds": FDV2_PAYLOAD_KIND} - assert repair["if_none_match"] is None - - def test_a_non_400_after_the_repair_meets_the_spent_budget( - self, endpoint: Any - ) -> None: - """The exemption is for the repair, not for the run that follows it.""" - store = poll_store(endpoint, max_consecutive_failures=1) - endpoint.queue_poll(full_payload(("put-object", put_skill()))) - endpoint.queue_poll(status=500) - endpoint.queue_poll(status=400) - endpoint.queue_poll(status=500) - endpoint.queue_poll(full_payload(("put-object", put_skill()))) - with store: - assert store.wait_for_skills(timeout=5) is True - assert wait_until(lambda: store.failed is not None) - assert "gave up after 3 consecutive failures" in store.failed - # The fifth queued payload is never asked for. - assert len(endpoint.requests) == 4 - def test_a_400_carrying_no_client_state_is_fatal_at_once( self, endpoint: Any ) -> None: @@ -2121,8 +2181,7 @@ def test_a_422_on_the_first_response_stops_delivery(self, endpoint: Any) -> None """ endpoint.queue_poll(status=422) endpoint.queue_poll(full_payload(("put-object", put_skill()))) - # Well above one, so the bound is not what stopped it. - with poll_store(endpoint, max_consecutive_failures=5) as store: + with poll_store(endpoint) as store: assert wait_until(lambda: store.failed is not None) assert "422" in store.failed # The queued payload is never asked for. @@ -2292,32 +2351,111 @@ def test_a_retry_resets_the_failure_count_on_success(self, endpoint: Any) -> Non assert store.wait_for_skills(timeout=5) assert wait_until(lambda: store.diagnostics.connection_failures == 0) - def test_retries_are_bounded(self) -> None: + def test_there_is_no_consecutive_failure_option(self) -> None: + """A count bound would freeze a give-up contract into the public API. + + Asserted by name, so a reintroduction is caught here rather than in a + process that stopped receiving revocations after a short outage. + """ + assert ( + "max_consecutive_failures" + not in inspect.signature(FDv2SkillStore.__init__).parameters + ) + with pytest.raises(TypeError): + FDv2SkillStore(SDK_KEY, max_consecutive_failures=3) # type: ignore[call-arg] + + @pytest.mark.parametrize("option", ["initial_backoff", "max_backoff"]) + @pytest.mark.parametrize("value", [0, 0.0, -5, math.nan, math.inf]) + def test_a_backoff_option_must_be_positive_and_finite( + self, option: str, value: float + ) -> None: + """With no failure bound, these two numbers are the only limit on the + retry loop: zero or less reconnects as fast as the network allows.""" + kwargs = {"initial_backoff": 1.0, "max_backoff": 30.0, option: value} + with pytest.raises(ValueError, match=option): + FDv2SkillStore(SDK_KEY, **kwargs) + + def test_the_initial_backoff_may_not_exceed_the_cap(self) -> None: + with pytest.raises(ValueError, match="must not exceed"): + FDv2SkillStore(SDK_KEY, initial_backoff=5.0, max_backoff=1.0) + # Equal is a fixed delay, and fine. + FDv2SkillStore(SDK_KEY, initial_backoff=2.0, max_backoff=2.0) + + def test_wait_for_skills_runs_to_its_timeout_during_an_outage(self) -> None: + """Recoverable failures no longer end delivery, so the store does not + know the answer yet, and the wait is not cut short.""" + store = stream_store(_requester=_ScriptedRequester()) + try: + store.start() + started = time.monotonic() + assert store.wait_for_skills(timeout=0.5) is False + assert time.monotonic() - started >= 0.4 + assert store.failed is None + assert store.diagnostics.connection_failures > 0 + finally: + store.close() + + def test_a_400_after_a_recoverable_failure_still_repairs( + self, endpoint: Any + ) -> None: + """The repair still goes out from scratch when an outage came first.""" + endpoint.queue_poll(full_payload(("put-object", put_skill()))) + endpoint.queue_poll(status=500) + endpoint.queue_poll(status=400) + endpoint.queue_poll(full_payload(("put-object", put_skill()))) + with poll_store(endpoint) as store: + assert store.wait_for_skills(timeout=5) is True + assert wait_until(lambda: len(endpoint.requests) >= 4) + assert store.failed is None + # The 500 left the basis in place; the 400 is what dropped it. + assert "basis" in endpoint.requests[2]["query"] + repair = endpoint.requests[3] + assert repair["query"] == {"kinds": FDV2_PAYLOAD_KIND} + assert repair["if_none_match"] is None + + def test_the_second_400_still_stops_with_a_recoverable_failure_between( + self, endpoint: Any + ) -> None: + """A recoverable failure between them does not earn the 400 another + repair: the request after it still carries no client state.""" + endpoint.queue_poll(full_payload(("put-object", put_skill()))) + endpoint.queue_poll(status=400) + endpoint.queue_poll(status=500) + endpoint.queue_poll(status=400) + endpoint.queue_poll(full_payload(("put-object", put_skill()))) + with poll_store(endpoint) as store: + assert store.wait_for_skills(timeout=5) is True + assert wait_until(lambda: store.failed is not None) + assert "400" in store.failed + # The fifth queued payload is never asked for. + assert len(endpoint.requests) == 4 + + def test_recoverable_failures_are_retried_indefinitely(self) -> None: + """Well past ten in a row, the bound this replaces, and still retrying.""" + requester = _ScriptedRequester() store = FDv2SkillStore( SDK_KEY, mode="poll", poll_interval=0.01, initial_backoff=0.001, max_backoff=0.002, - max_consecutive_failures=3, - _requester=_ScriptedRequester(), + _requester=requester, ) try: store.start() - assert wait_until(lambda: store.failed is not None) - # Four, not three: the bound is the number of failures *tolerated*, - # so the run that exceeds it is the one that gives up. - assert "gave up after 4 consecutive failures" in store.failed + assert wait_until(lambda: store.diagnostics.connection_failures >= 25) + assert store.failed is None + assert store.diagnostics.last_error is not None finally: store.close() + assert len(requester.calls) >= 25 def test_recycled_stream_connections_are_not_failures(self) -> None: # A streaming connection only ever ends by being dropped, so a loop - # that counted every drop as a failure would give up on a healthy - # server after max_consecutive_failures + 1 recycles, and delivery - # (including revocation) would silently stop for the process lifetime. + # that counted every drop as a failure would back a healthy server off + # to ``max_backoff`` and report it as failing. requester = _RecyclingRequester() - store = stream_store(max_consecutive_failures=3, _requester=requester) + store = stream_store(_requester=requester) try: store.start() assert wait_until(lambda: requester.connections >= 8) @@ -2337,21 +2475,45 @@ def test_an_up_to_date_recycled_stream_is_not_a_failure( # Resetting at a commit covers only a connection that carried new # content. An environment whose skills are not changing answers every # reconnect with ``intentCode: "none"`` and transfers nothing, so a loop - # that counted those drops would give up on a *healthy* idle stream - # after max_consecutive_failures + 1 recycles — and revocation, the one - # thing streaming exists to deliver promptly, would never arrive again. + # that counted those drops would back a *healthy* idle stream off to + # ``max_backoff`` and report it as failing. requester = _UpToDateRecyclingRequester(farewell=farewell) - store = stream_store(max_consecutive_failures=3, _requester=requester) + store = stream_store(_requester=requester) try: store.start() assert wait_until(lambda: requester.connections >= 8) assert store.failed is None - # As with a payload-carrying recycle, the count may read 1 - # mid-reconnect. What it must never do is climb. - assert store.diagnostics.connection_failures <= 1 + if farewell: + # A goodbye after an answer is the server recycling the stream, + # not a failure: it is neither counted nor reported. + assert store.diagnostics.connection_failures == 0 + assert store.diagnostics.last_error is None + else: + # A plain drop counts until the next answer clears it, so the + # count may read 1 mid-reconnect. What it must never do is climb. + assert store.diagnostics.connection_failures <= 1 finally: store.close() + def test_a_goodbye_before_any_answer_is_a_failure(self, caplog: Any) -> None: + """Only a completed exchange makes a goodbye routine. A server that only + ever says goodbye has delivered nothing, and must stay visible.""" + + class _OnlyGoodbye(_FakeRequester): + def stream(self, basis: str | None) -> Any: + return _ScriptedConnection([("goodbye", {"reason": "go away"})]) + + store = stream_store(_requester=_OnlyGoodbye()) + with caplog.at_level("DEBUG", logger="launchdarkly_ai_server.skills_fdv2"): + try: + store.start() + assert wait_until(lambda: store.diagnostics.connection_failures >= 3) + finally: + store.close() + assert "goodbye" in (store.diagnostics.last_error or "") + 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_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 @@ -2371,9 +2533,7 @@ def test_a_recycled_connection_reconnects_quietly(self, caplog: Any) -> None: 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. - store = stream_store( - max_consecutive_failures=10, _requester=_ScriptedRequester() - ) + store = stream_store(_requester=_ScriptedRequester()) with caplog.at_level("DEBUG", logger="launchdarkly_ai_server.skills_fdv2"): try: store.start() @@ -2399,7 +2559,7 @@ def test_a_stream_that_dies_mid_read_reconnects(self, exc: BaseException) -> Non # such a failure as unexpected would stop delivery — including # revocation — for the process lifetime the first time a socket died. requester = _DyingStreamRequester(exc) - store = stream_store(max_consecutive_failures=3, _requester=requester) + store = stream_store(_requester=requester) try: store.start() assert store.wait_for_skills(timeout=5) is True @@ -2417,31 +2577,116 @@ def test_a_stream_commit_resets_the_failure_count(self) -> None: _RecoverableTransportError("x"), _RecoverableTransportError("x"), payload, + _FatalTransportError("401"), ) - store = stream_store(max_consecutive_failures=3, _requester=requester) + store = stream_store(_requester=requester) try: store.start() assert store.wait_for_skills(timeout=5) - # Three failures reach the bound, then a commit, then the exhausted - # requester fails on every reconnect. The count must start again at - # the commit: the stream's own drop is failure one, and three more - # connects are owed before giving up. Carrying the three over would - # give up on the drop itself, with no further connect at all. assert wait_until(lambda: store.failed is not None) - assert "gave up after 4 consecutive failures" in store.failed - assert "last error: x" in store.failed - assert len(requester.calls) == 7 + # Three failures, then a commit, then the stream's own drop. The + # count starts again at the commit, so the drop is failure one, not + # four. + assert store.diagnostics.connection_failures == 1 + assert len(requester.calls) == 5 finally: store.close() - def test_stream_retries_are_bounded(self) -> None: - store = stream_store( - max_consecutive_failures=3, _requester=_ScriptedRequester() - ) + def test_an_answer_alone_does_not_reset_the_backoff_delay( + self, monkeypatch: Any + ) -> None: + """A server that answers ``none`` and drops at once is backed off. + + Resetting the delay on every answer would reconnect it about once a + second, from every process, for as long as it stays degraded. The + failure count still clears on each answer. + """ + attempts = _record_backoff_attempts(monkeypatch) + requester = _UpToDateRecyclingRequester() + store = stream_store(_requester=requester) try: store.start() - assert wait_until(lambda: store.failed is not None) - assert "gave up after 4 consecutive failures" in store.failed + assert wait_until(lambda: len(attempts) >= 6) + assert store.diagnostics.connection_failures <= 1 + finally: + store.close() + assert attempts[:6] == [1, 2, 3, 4, 5, 6] + + def test_a_stream_held_past_the_threshold_resets_the_backoff_delay( + self, monkeypatch: Any + ) -> None: + monkeypatch.setattr(skills_fdv2, "_BACKOFF_RESET_INTERVAL", 0.05) + attempts = _record_backoff_attempts(monkeypatch) + + class _HeldThenDropped(_FakeRequester): + def stream(self, basis: str | None) -> Any: + def held() -> Any: + yield ("server-intent", server_intent("none")) + time.sleep(0.1) + + return _ScriptedConnection(held()) + + store = stream_store(_requester=_HeldThenDropped()) + try: + store.start() + assert wait_until(lambda: len(attempts) >= 3) + finally: + store.close() + assert attempts[:3] == [1, 1, 1] + + def test_a_completed_poll_resets_the_backoff_delay( + self, endpoint: Any, monkeypatch: Any + ) -> None: + """Fail, succeed, fail: the second failure retries at the first step, + since ``poll_interval`` already spaces the requests.""" + attempts = _record_backoff_attempts(monkeypatch) + endpoint.queue_poll(status=500) + endpoint.queue_poll(full_payload(("put-object", put_skill()))) + endpoint.queue_poll(status=500) + with poll_store(endpoint) as store: + assert store.wait_for_skills(timeout=5) + assert wait_until(lambda: len(attempts) >= 2) + assert store.diagnostics.connection_failures <= 1 + assert attempts[:2] == [1, 1] + + def test_stream_failures_are_retried_indefinitely(self) -> None: + """Last known good is served throughout, however long the outage.""" + payload = [ + (e["event"], e["data"]) for e in full_payload(("put-object", put_skill())) + ] + store = stream_store(_requester=_ScriptedRequester(payload)) + try: + store.start() + assert store.wait_for_skills(timeout=5) + assert wait_until(lambda: store.diagnostics.connection_failures >= 25) + assert store.failed is None + assert store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") is not None + finally: + store.close() + + def test_an_announced_transfer_that_never_completes_is_a_failure(self) -> None: + """An intent is a promise, not a delivery, so it does not reset the count. + + A server that announces a transfer and drops before + ``payload-transferred``, every time, has delivered nothing. Counted as + health, it would be retried forever at the initial backoff. + """ + + class _AnnouncesThenDrops(_FakeRequester): + def stream(self, basis: str | None) -> Any: + return _ScriptedConnection( + [ + ("server-intent", server_intent("xfer-full")), + ("put-object", put_skill()), + ] + ) + + store = stream_store(_requester=_AnnouncesThenDrops()) + try: + store.start() + assert wait_until(lambda: store.diagnostics.connection_failures >= 5) + assert store.failed is None + assert store.is_initialized() is False finally: store.close() @@ -2561,6 +2806,17 @@ def test_backoff_is_exponential_and_capped(self) -> None: assert _backoff_delay(3, base=1.0, maximum=30.0, jitter=0.0) == 4.0 assert _backoff_delay(20, base=1.0, maximum=30.0, jitter=0.0) == 30.0 + def test_backoff_stays_finite_at_any_attempt_number(self) -> None: + """Retries are unbounded, so the attempt number is too. + + Unclamped, ``float(2 ** n)`` raises ``OverflowError`` past about 1024, + which a long enough outage reaches. + """ + assert _backoff_delay(10_000, base=1.0, maximum=30.0, jitter=0.0) == 30.0 + delay = _backoff_delay(10_000, base=1.0, maximum=30.0) + assert math.isfinite(delay) + assert 0.0 <= delay <= 30.0 + def test_jitter_never_exceeds_the_cap(self) -> None: for attempt in range(1, 12): for _ in range(50): @@ -2643,25 +2899,33 @@ class TestTransportMemoryBound: def test_the_bound_is_far_above_any_legitimate_payload(self) -> None: assert MAX_RESPONSE_BYTES == 64 * 1024 * 1024 - def test_an_over_cap_poll_body_is_not_applied_and_is_retried( + def test_crossing_the_bound_is_fatal(self) -> None: + assert issubclass(_ResponseTooLargeError, _FatalTransportError) + assert not issubclass(_ResponseTooLargeError, _RecoverableTransportError) + + def test_an_over_cap_poll_body_stops_delivery_without_applying_it( self, endpoint: Any, monkeypatch: Any ) -> None: + """Fatal, like a 422: the size belongs to the environment, so a retry + would download it again and be refused the same way.""" monkeypatch.setattr(skills_fdv2, "MAX_RESPONSE_BYTES", 2048) endpoint.queue_poll(full_payload(("put-object", put_skill(content="x" * 8192)))) - # 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) + with poll_store(endpoint) as store: + assert wait_until(lambda: store.failed is not None) + assert "2048-byte transport bound" in store.failed assert "2048-byte transport bound" in (store.diagnostics.last_error or "") assert store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") is None 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. - assert wait_until(lambda: len(endpoint.requests) >= 2) + assert store.diagnostics.connection_failures == 0 + # Fatal means one request, not a retry loop. + assert len(endpoint.requests) == 1 + + # Once the payload is back under the bound, start() resumes. + endpoint.queue_poll(full_payload(("put-object", put_skill()))) + store.start() assert store.wait_for_skills(timeout=5) is True - assert wait_until(lambda: store.diagnostics.connection_failures == 0) - assert "2048-byte transport bound" in (store.diagnostics.last_error or "") - assert all(r["path"] == "/sdk/poll" for r in endpoint.requests) + assert store.failed is None def test_a_poll_body_exactly_at_the_cap_is_accepted( self, endpoint: Any, monkeypatch: Any @@ -2683,7 +2947,7 @@ def test_a_poll_body_one_byte_over_the_cap_is_refused( body = json.dumps({"events": payload}).encode("utf-8") monkeypatch.setattr(skills_fdv2, "MAX_RESPONSE_BYTES", len(body) - 1) endpoint.queue_poll(payload) - with poll_store(endpoint, max_consecutive_failures=0) as store: + with poll_store(endpoint) as store: assert wait_until(lambda: store.failed is not None) assert store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") is None assert f"{len(body) - 1}-byte transport bound" in store.failed @@ -2704,8 +2968,8 @@ def test_an_over_cap_stream_event_abandons_the_payload_in_flight( ) -> None: """ The first payload commits. The second starts, then carries an event - over the cap: that connection is dropped, the half-received payload is - never committed, and the reconnect finds the committed set intact. + over the cap: delivery stops, the half-received payload is never + committed, and the committed set is still served. """ monkeypatch.setattr(skills_fdv2, "MAX_RESPONSE_BYTES", 2048) endpoint.queue_stream( @@ -2716,8 +2980,6 @@ def test_an_over_cap_stream_event_abandons_the_payload_in_flight( ("payload-transferred", transferred("basis-2")), ) ) - endpoint.hold_stream_open = True - endpoint.queue_stream(events(("server-intent", server_intent("none")))) store = FDv2SkillStore( SDK_KEY, base_uri=endpoint.base_uri, @@ -2728,16 +2990,13 @@ def test_an_over_cap_stream_event_abandons_the_payload_in_flight( try: store.start() assert store.wait_for_skills(timeout=5) is True - assert wait_until( - lambda: "transport bound" in (store.diagnostics.last_error or "") - ) - assert wait_until(lambda: len(endpoint.requests) >= 2) + assert wait_until(lambda: store.failed is not None) + assert "transport bound" in store.failed assert store.get_object(SKILL_OBJECT_KIND, "pdf-extraction") is not None assert store.get_object(SKILL_OBJECT_KIND, "oversized") is None assert store.diagnostics.payloads_transferred == 1 assert store.diagnostics.skill_objects_received == 1 - assert store.failed is None - assert endpoint.requests[1]["query"].get("basis") == "basis-1" + assert len(endpoint.requests) == 1 finally: store.close() @@ -2746,7 +3005,7 @@ def test_a_stream_line_that_never_ends_is_refused_and_the_body_closed( ) -> None: monkeypatch.setattr(skills_fdv2, "MAX_RESPONSE_BYTES", 1024) source = _LineSource(b"data: " + b"x" * 4096) - with pytest.raises(_RecoverableTransportError, match="1024-byte"): + with pytest.raises(_ResponseTooLargeError, match="1024-byte"): list(_iter_sse(source)) assert source.closed @@ -2754,7 +3013,7 @@ def test_an_event_is_measured_across_its_data_lines(self, monkeypatch: Any) -> N monkeypatch.setattr(skills_fdv2, "MAX_RESPONSE_BYTES", 1024) lines = b"".join(b"data: " + b"x" * 500 + b"\n" for _ in range(3)) source = _LineSource(b"event: put-object\n" + lines + b"\n") - with pytest.raises(_RecoverableTransportError, match="1024-byte"): + with pytest.raises(_ResponseTooLargeError, match="1024-byte"): list(_iter_sse(source)) assert source.closed @@ -3594,6 +3853,7 @@ def test_it_adds_no_dependency(self) -> None: "re", "socket", "threading", + "time", "typing", "urllib", } @@ -3681,9 +3941,7 @@ def test_a_failed_delivery_records_nothing_either( success path alone would not cover it. """ skills_module._set_emitter_for_testing(recording_emitter) - store = stream_store( - max_consecutive_failures=1, _requester=_ScriptedRequester() - ) + store = stream_store(_requester=_ScriptedRequester(_FatalTransportError("401"))) try: store.start() assert wait_until(lambda: store.failed is not None) @@ -3974,7 +4232,7 @@ def test_a_store_restarted_after_giving_up_waits_again(self) -> None: # Restarting after a *close* is not available — see # ``TestCloseIsFinal`` — so the give-up path is what exercises this. class _FailsThenGoesQuiet(_FakeRequester): - """One failure, enough to give up; silent on every run after.""" + """One fatal failure; silent on every run after.""" def __init__(self) -> None: self.attempts = 0 @@ -3982,12 +4240,10 @@ def __init__(self) -> None: def stream(self, basis: str | None) -> Any: self.attempts += 1 if self.attempts == 1: - raise _RecoverableTransportError("x") + raise _FatalTransportError("x") return _BlockingConnection() - store = stream_store( - max_consecutive_failures=0, _requester=_FailsThenGoesQuiet() - ) + store = stream_store(_requester=_FailsThenGoesQuiet()) store.start() assert wait_until(lambda: store.failed is not None) assert store.wait_for_skills(timeout=0.1) is False