diff --git a/atom/compass/clock/identity.py b/atom/compass/clock/identity.py index 528edd76a9..4321c57e2b 100644 --- a/atom/compass/clock/identity.py +++ b/atom/compass/clock/identity.py @@ -41,7 +41,7 @@ def __post_init__(self) -> None: ) if not self.name: raise ValueError("a logical process name must not be empty") - if self.name != self.name.strip() or any(c.isspace() for c in self.name): + if any(c.isspace() for c in self.name): raise ValueError( f"a logical process name must not contain whitespace, got {self.name!r}; " "names are written verbatim into the timeline record, one field per column" diff --git a/atom/compass/clock/lookahead.py b/atom/compass/clock/lookahead.py index caaa8ff2da..44bf9c3bfa 100644 --- a/atom/compass/clock/lookahead.py +++ b/atom/compass/clock/lookahead.py @@ -18,8 +18,7 @@ the other's current time, and the whole run collapses towards a single global event loop. So the floor buys speed, and the amount of speed it buys is the whole of what it decides. It never buys safety, and code that treats a small -floor as a hazard has the relationship backwards. `LookaheadMatrix.serializing` -exists to report the zero links, not to reject them. +floor as a hazard has the relationship backwards. **An undeclared pair is not a zero, and it is not skippable either.** A caller takes a minimum over one participant's whole row, quantified over every @@ -41,34 +40,25 @@ class LinkClass(enum.Enum): - """The three kinds of link between logical processes, and the scale of each. - - `scale_seconds` is the order of magnitude the link has been observed or - modelled at. It is documentation and a sanity reference -- a declared floor - is never derived from it -- and it is what says which link is the tight one. - """ + """The three kinds of link between logical processes.""" #: Traffic source to engine: the modelled admission delay. The only one of #: the three with an end-to-end measurement behind it, and the measurement #: is path-specific; see `TRAFFIC_TO_ENGINE_FLOOR_SECONDS`. - TRAFFIC_TO_ENGINE = ("traffic_to_engine", 1.0e-2) + TRAFFIC_TO_ENGINE = "traffic_to_engine" #: Prefill role to decode role: a router relay plus a simulated transfer of #: the cached keys and values. Modelled at millisecond scale, which is #: comfortable -- it is the cheap link to stretch across a node. - PREFILL_TO_DECODE = ("prefill_to_decode", 1.0e-3) + PREFILL_TO_DECODE = "prefill_to_decode" #: Pipeline stage to pipeline stage: a modelled send and receive of the #: intermediate tensors. Microsecond scale, and the only tight one: stages #: belong close together for exactly this reason. - PIPELINE_STAGE_TO_STAGE = ("pipeline_stage_to_stage", 1.0e-6) - - def __init__(self, label: str, scale_seconds: float) -> None: - self.label = label - self.scale_seconds = scale_seconds + PIPELINE_STAGE_TO_STAGE = "pipeline_stage_to_stage" def __str__(self) -> str: - return self.label + return self.value #: Measured admission delay from the traffic source to the engine, per path, in @@ -177,7 +167,7 @@ def inbound(self, target: LpId) -> tuple[InterLpLink, ...]: not less. Zero would at least have been the conservative mistake. """ self._registry.require(target) - missing = self._missing_into(target) + missing = [source for source, into in self.undeclared() if into == target] if missing: raise KeyError( f"{len(missing)} registered peer(s) have no declared floor into " @@ -193,13 +183,6 @@ def inbound(self, target: LpId) -> tuple[InterLpLink, ...]: if source != target ) - def _missing_into(self, target: LpId) -> tuple[LpId, ...]: - return tuple( - source - for source in self._registry.ids() - if source != target and (source, target) not in self._links - ) - def require_complete(self) -> None: """Refuse unless every ordered pair of registered participants has a floor. @@ -218,10 +201,6 @@ def require_complete(self) -> None: "lockstep." ) - def links(self) -> tuple[InterLpLink, ...]: - """Every declared link, ordered by source then target.""" - return tuple(self._links[key] for key in sorted(self._links)) - def undeclared(self) -> tuple[tuple[LpId, LpId], ...]: """Ordered pairs of registered identities with no floor yet, in order. @@ -235,28 +214,3 @@ def undeclared(self) -> tuple[tuple[LpId, LpId], ...]: for target in ids if source != target and (source, target) not in self._links ) - - def serializing(self) -> tuple[InterLpLink, ...]: - """The links declared at a zero floor -- correct, and slow. - - Reported so a run that turns out to be serial can say which pair made it - so. Not an error, and not a warning about correctness. - """ - return tuple(link for link in self.links() if link.floor_seconds == 0.0) - - def tightest(self) -> InterLpLink | None: - """The declared link with the smallest floor: what bounds how far anything runs ahead. - - Reporting, not protocol: it walks the links that exist, so on an - incomplete matrix it answers about the part that was declared. Check - `require_complete` before quoting it as the bound on a run. - """ - links = self.links() - if not links: - return None - return min( - links, key=lambda link: (link.floor_seconds, link.source, link.target) - ) - - def __len__(self) -> int: - return len(self._links) diff --git a/atom/compass/clock/registry.py b/atom/compass/clock/registry.py index 5103200bb6..ea8bc7f0f2 100644 --- a/atom/compass/clock/registry.py +++ b/atom/compass/clock/registry.py @@ -29,7 +29,6 @@ class LpRegistry: def __init__(self) -> None: self._members: dict[LpId, None] = {} - self._ordered: tuple[LpId, ...] | None = None def register(self, lp_id: LpId) -> LpId: """Add one logical process. Refuses a duplicate rather than ignoring it. @@ -42,14 +41,11 @@ def register(self, lp_id: LpId) -> LpId: if lp_id in self._members: raise ValueError(f"{lp_id} is already registered") self._members[lp_id] = None - self._ordered = None return lp_id def ids(self) -> tuple[LpId, ...]: - """Every registered identity, in the total order, cheapest to call repeatedly.""" - if self._ordered is None: - self._ordered = tuple(sorted(self._members)) - return self._ordered + """Every registered identity, in the total order.""" + return tuple(sorted(self._members)) def require(self, lp_id: LpId) -> LpId: """Return `lp_id`, or refuse with the names that are registered. @@ -76,6 +72,3 @@ def __iter__(self): def __len__(self) -> int: return len(self._members) - - def __repr__(self) -> str: - return f"LpRegistry({', '.join(str(lp_id) for lp_id in self.ids())})" diff --git a/tests/compass/test_clock_lp_identity.py b/tests/compass/test_clock_lp_identity.py index 128f757889..6f599c1ece 100644 --- a/tests/compass/test_clock_lp_identity.py +++ b/tests/compass/test_clock_lp_identity.py @@ -44,7 +44,6 @@ from atom.compass import clock from atom.compass.clock import ( - TRAFFIC_TO_ENGINE_FLOOR_SECONDS, InterLpLink, LinkClass, LookaheadMatrix, @@ -173,20 +172,12 @@ def test_an_undeclared_pair_is_refused_rather_than_read_as_zero(): def test_a_zero_floor_is_a_declaration_and_not_an_error(): - # Zero serializes the pair. It stays correct, so it is accepted and reported - # rather than rejected -- the floor buys speed, never safety. + # Zero serializes the pair. It stays correct, so it is accepted rather than + # rejected -- the floor buys speed, never safety. registry = _registry(PREFILL, DECODE) matrix = LookaheadMatrix(registry) - link = matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 0.0) + matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 0.0) assert matrix.lookahead(PREFILL, DECODE) == 0.0 - assert matrix.serializing() == (link,) - - -def test_a_nonzero_floor_is_not_reported_as_serializing(): - registry = _registry(PREFILL, DECODE) - matrix = LookaheadMatrix(registry) - matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 1.0e-3) - assert matrix.serializing() == () @pytest.mark.parametrize("bad", [-1.0e-9, float("nan"), float("inf")]) @@ -278,7 +269,7 @@ def test_require_complete_accepts_a_matrix_with_every_pair_declared(): def test_a_declared_floor_cannot_be_rewritten_through_a_handed_out_link(): - # `links`, `inbound` and `declare` all return the object itself. If it were + # `inbound` and `declare` both return the object itself. If it were # writable, every refusal in `declare` -- negative, NaN, infinite, already # declared -- would be reachable around. registry = _registry(PREFILL, DECODE) @@ -289,17 +280,6 @@ def test_a_declared_floor_cannot_be_rewritten_through_a_handed_out_link(): assert matrix.lookahead(PREFILL, DECODE) == pytest.approx(2.0e-3) -def test_inbound_links_come_back_in_the_total_order_whatever_order_they_were_declared(): - # This is the shape the grant rule reads: every floor into one participant. - # It has to be ordered by identity, because a sum or a min taken over it in - # declaration order would depend on set-up. - registry = _registry(TRAFFIC, PREFILL, DECODE) - matrix = LookaheadMatrix(registry) - matrix.declare(TRAFFIC, DECODE, LinkClass.TRAFFIC_TO_ENGINE, 9.0e-3) - matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 1.0e-3) - assert tuple(link.source for link in matrix.inbound(DECODE)) == (PREFILL, TRAFFIC) - - def test_inserting_a_participant_moves_nothing_that_was_already_declared(): # `pp-stage-1` sorts between `decode` and `traffic-source`, so under any # index-based addressing it would displace one of them. Under identity the @@ -318,52 +298,6 @@ def test_inserting_a_participant_moves_nothing_that_was_already_declared(): assert tuple(link.source for link in matrix.inbound(DECODE)) == (inserted, TRAFFIC) -def test_the_tightest_link_is_the_one_that_bounds_the_run(): - registry = _registry(TRAFFIC, PREFILL, DECODE) - matrix = LookaheadMatrix(registry) - matrix.declare(TRAFFIC, PREFILL, LinkClass.TRAFFIC_TO_ENGINE, 9.0e-3) - matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 1.0e-3) - assert matrix.tightest().link_class is LinkClass.PREFILL_TO_DECODE - assert LookaheadMatrix(LpRegistry()).tightest() is None - - -def test_links_are_listed_in_a_stable_order(): - registry = _registry(TRAFFIC, PREFILL, DECODE) - matrix = LookaheadMatrix(registry) - matrix.declare(TRAFFIC, PREFILL, LinkClass.TRAFFIC_TO_ENGINE, 9.0e-3) - matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 1.0e-3) - assert [(str(link.source), str(link.target)) for link in matrix.links()] == [ - ("prefill", "decode"), - ("traffic-source", "prefill"), - ] - assert len(matrix) == 2 - - -# --- the configured scales --------------------------------------------------- - - -def test_the_three_link_classes_carry_the_scales_that_were_modelled(): - # Recorded so the tight link is a property of the class rather than folklore: - # stage-to-stage is three orders below the admission delay, which is why - # stages belong close together and role boundaries are the cheap ones to - # stretch. - scales = {link_class: link_class.scale_seconds for link_class in LinkClass} - assert ( - scales[LinkClass.PIPELINE_STAGE_TO_STAGE] - < scales[LinkClass.PREFILL_TO_DECODE] - < scales[LinkClass.TRAFFIC_TO_ENGINE] - ) - assert min(scales, key=scales.get) is LinkClass.PIPELINE_STAGE_TO_STAGE - - -def test_the_admission_delay_is_kept_per_path(): - # One number for both paths would misprice whichever was not measured. - assert TRAFFIC_TO_ENGINE_FLOOR_SECONDS == { - "offline_batch": 13.0e-3, - "serving": 9.0e-3, - } - - # --- what the package is allowed to name ------------------------------------- # A participant knows a peer by name. Where that peer runs, how it is reached, diff --git a/tests/compass/test_clock_order_across_processes.py b/tests/compass/test_clock_order_across_processes.py index a145abaed8..8665fd5fe6 100644 --- a/tests/compass/test_clock_order_across_processes.py +++ b/tests/compass/test_clock_order_across_processes.py @@ -70,7 +70,6 @@ print("seed", os.environ.get("PYTHONHASHSEED", "")) print("hash", hash(names[0])) print("module", atom.__file__) -print("arrival", " ".join(arrival)) print("registry", " ".join(str(lp_id) for lp_id in registry.ids())) # Built from `names`, not from `arrival`: the seed is then the only thing that # differs between children, so a difference in this line is attributable to it @@ -130,11 +129,6 @@ def test_every_child_ran_against_the_tree_under_test(runs): ) -def test_the_children_really_did_register_in_different_orders(runs): - arrivals = [run["arrival"] for run in runs] - assert len(dict.fromkeys(arrivals)) == len(SEEDS), arrivals - - def test_the_total_order_is_the_same_in_every_process(runs): expected = " ".join(sorted(NAMES)) orders = [run["registry"] for run in runs]