Skip to content
Merged
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
3 changes: 3 additions & 0 deletions functions/compose-gke-cluster/function/fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,9 @@ def resolve_project(self) -> str | None:
)
return project
response.normal(self.rsp, f"Waiting for GCP {cred_kind} {cred_name}")
# Nothing is composed while waiting for the project, and an XR with
# no composed resources would otherwise be trivially ready.
self.rsp.desired.composite.ready = fnv1.READY_FALSE
return None
if cred_kind == "ClusterProviderConfig":
pc = gcpcpcv1beta1.ClusterProviderConfig.model_validate(d)
Expand Down
28 changes: 28 additions & 0 deletions functions/compose-gke-cluster/tests/test_fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -580,8 +580,36 @@ async def test_compose(self) -> None:
)
want2.requirements.resources["gcp-provider-config"].CopyFrom(_GCP_PROVIDER_CONFIG_SELECTOR)

# The ProviderConfig resolved to nothing and no cluster is observed to
# take the project from, so nothing can be composed. The XR is marked
# not ready rather than left to aggregate to trivially ready.
req3 = fnv1.RunFunctionRequest(
observed=fnv1.State(
composite=fnv1.Resource(
resource=resource.dict_to_struct(
_gke_xr().model_dump(exclude_none=True, mode="json"),
),
),
),
)
req3.required_resources["gcp-provider-config"].SetInParent()

want3 = fnv1.RunFunctionResponse(
meta=fnv1.ResponseMeta(ttl=durationpb.Duration(seconds=60)),
desired=fnv1.State(composite=fnv1.Resource(ready=fnv1.READY_FALSE)),
results=[
fnv1.Result(
severity=fnv1.SEVERITY_NORMAL,
message="Waiting for GCP ClusterProviderConfig default",
),
],
context=structpb.Struct(),
)
want3.requirements.resources["gcp-provider-config"].CopyFrom(_GCP_PROVIDER_CONFIG_SELECTOR)

cases = [
Case(name="first pass composes infra resources; IAM binding gated", req=req1, want=want1),
Case(name="a missing ProviderConfig composes nothing and isn't ready", req=req3, want=want3),
Case(
name="second pass with observed SA email composes IAM binding and marks ready resources",
req=req2,
Expand Down
3 changes: 3 additions & 0 deletions functions/compose-inference-cluster/function/fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -519,6 +519,9 @@ def resolve_classes(self) -> bool:
),
)
response.normal(self.rsp, f"Waiting for InferenceClasses: {', '.join(missing)}")
# Only the guard and namespaces, both marked ready, are composed
# while waiting, so the XR would otherwise be ready.
self.rsp.desired.composite.ready = fnv1.READY_FALSE
return False

return True
Expand Down
5 changes: 4 additions & 1 deletion functions/compose-inference-cluster/tests/test_fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -311,10 +311,13 @@ def _early_return_guard_case() -> tuple[fnv1.RunFunctionRequest, fnv1.RunFunctio
want = fnv1.RunFunctionResponse(
meta=fnv1.ResponseMeta(ttl=durationpb.Duration(seconds=60)),
desired=fnv1.State(
# The guard and namespace are marked ready, so the XR is marked not
# ready while it waits for its classes.
composite=fnv1.Resource(ready=fnv1.READY_FALSE),
resources={
"usage-replicas": _guard_clusterusage(),
"namespace-team-a": _namespace_object("team-a", "mp-team-a-bd964"),
}
},
),
context=structpb.Struct(),
)
Expand Down
6 changes: 6 additions & 0 deletions functions/compose-model-cache/function/fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -550,6 +550,9 @@ def derive_conditions(
reason=CONDITION_REASON_NO_CLUSTERS,
),
)
# Nothing is composed with no cluster to stage onto, and an XR with
# no composed resources would otherwise be trivially ready.
self.rsp.desired.composite.ready = fnv1.READY_FALSE
return
response.set_conditions(
self.rsp,
Expand Down Expand Up @@ -583,6 +586,9 @@ def derive_conditions(
self.rsp,
f"authSecret {_namespace(self.xr.metadata)}/{auth.name} is missing or has no key {key!r}",
)
# Only the PVCs are composed while the token is missing, and once
# they bind the XR would otherwise be ready.
self.rsp.desired.composite.ready = fnv1.READY_FALSE
elif any(p == PHASE_FAILED for _, p in per_cluster_phase):
response.set_conditions(
self.rsp,
Expand Down
6 changes: 5 additions & 1 deletion functions/compose-model-cache/tests/test_fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -643,7 +643,9 @@ async def test_compose(self) -> None: # noqa: PLR0915
# doesn't depend on the token, so a cache isn't pruned for a missing one
# - but the hydration Job and token Secret are held back. ArtifactReady
# is False with reason AuthSecretMissing, and a warning names the Secret
# and key so the user can fix it instead of seeing the XR stall. ---
# and key so the user can fix it instead of seeing the XR stall. The XR
# is marked not ready, since the PVC alone would make it ready once it
# binds. ---
want10 = fnv1.RunFunctionResponse(
meta=fnv1.ResponseMeta(ttl=durationpb.Duration(seconds=60)),
desired=fnv1.State(
Expand All @@ -656,6 +658,7 @@ async def test_compose(self) -> None: # noqa: PLR0915
},
},
),
ready=fnv1.READY_FALSE,
),
resources={
"pvc-cluster-a": fnv1.Resource(resource=resource.dict_to_struct(_pvc_object("cluster-a-pc"))),
Expand Down Expand Up @@ -782,6 +785,7 @@ async def test_compose(self) -> None: # noqa: PLR0915
desired=fnv1.State(
composite=fnv1.Resource(
resource=resource.dict_to_struct({"status": {"summary": {"ready": "0/0"}, "clusters": []}}),
ready=fnv1.READY_FALSE,
),
),
conditions=[
Expand Down
9 changes: 6 additions & 3 deletions functions/compose-model-deployment/function/fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -292,6 +292,9 @@ def resolve_inputs(self) -> bool:
),
)
response.warning(self.rsp, "No InferenceClusters found")
# Nothing is composed with no cluster to schedule onto, and an XR
# with no composed resources would otherwise be trivially ready.
self.rsp.desired.composite.ready = fnv1.READY_FALSE
return False

self.clusters = [_inference_cluster(icv1alpha1.InferenceCluster.model_validate(c)) for c in cluster_dicts]
Expand Down Expand Up @@ -524,9 +527,9 @@ def compose_endpoints(self, matched: list[scheduling.Candidate]) -> None:

A replica whose cluster has no hostname gets no endpoint, and so does
one whose ModelReplica isn't Ready. The replica's Ready tracks the
engine workloads serving and the remote Service and HTTPRoute that front
them, which is the whole path this endpoint advertises. Composing it any
earlier routes traffic at pods still pulling images or loading weights,
engine workloads and endpoint picker being available and the rest of the
routing objects in front of them being applied. Composing it any earlier
routes traffic at pods still pulling images or loading weights,
returning 503s during deployment and scale-up (#102). The endpoint
appears on the reconcile that first observes the replica Ready, and is
withdrawn again if the replica stops being Ready, pulling a dead backend
Expand Down
2 changes: 1 addition & 1 deletion functions/compose-model-deployment/tests/test_fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -588,7 +588,7 @@ async def test_compose(self) -> None:
want=_want(
fnv1.RunFunctionResponse(
meta=fnv1.ResponseMeta(ttl=durationpb.Duration(seconds=60)),
desired=fnv1.State(),
desired=fnv1.State(composite=fnv1.Resource(ready=fnv1.READY_FALSE)),
conditions=[
fnv1.Condition(
type="ReplicasScheduled",
Expand Down
32 changes: 14 additions & 18 deletions functions/compose-model-replica/function/fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,8 @@
engines. An engine's member roles select its backend: a Standalone member composes
to a Deployment (native), a Leader plus Worker to the cluster's chosen multi-node
backend - a LeaderWorkerSet (llm-d) or a PodCliqueSet (Grove), per the
InferenceCluster's stack. One shared Service and HTTPRoute front all of
a replica's engines.
InferenceCluster's stack. One HTTPRoute, InferencePool and endpoint picker
front all of a replica's engines.

Each member's template is a curated subset of PodTemplateSpec. The container
named "engine" is the inference engine; its image, command, and args are passed
Expand Down Expand Up @@ -116,6 +116,9 @@ def resolve_inputs(self) -> bool:
),
)
response.normal(self.rsp, "Waiting for cluster to be resolved")
# Nothing is composed while waiting, and an XR with no composed
# resources would otherwise be trivially ready.
self.rsp.desired.composite.ready = fnv1.READY_FALSE
return False

self.ic = icv1alpha1.InferenceCluster.model_validate(ic_dict)
Expand All @@ -134,6 +137,7 @@ def resolve_inputs(self) -> bool:
),
)
response.normal(self.rsp, "Waiting for cluster providerConfigRef")
self.rsp.desired.composite.ready = fnv1.READY_FALSE
return False

return True
Expand Down Expand Up @@ -218,23 +222,15 @@ def derive_conditions(self) -> None:
)

# Per-resource readiness. Crossplane gates the XR's Ready on every
# composed resource being ready, so the function must mark each one - a
# composed resource isn't ready just because provider-kubernetes set its
# Object's own Ready condition. Marking a resource ready asserts the
# function observed it ready, so we only ever mark a resource we can see
# in observed state. A workload additionally gates on the model actually
# serving; the Service, HTTPRoute, and ResourceClaimTemplates have no
# runtime readiness to wait on (existing is being ready), so observing
# them is enough. A freshly composed resource isn't in observed yet, so
# it stays unready until the next reconcile sees it applied.
workloads = set(workload_keys)
# composed resource being ready, but doesn't read a composed resource's
# own Ready condition, so the function relays it. Each Object's
# readiness policy decides that condition: a workload and the endpoint
# picker are Ready once their CEL query sees them available, anything
# else once applied. Envoy AI Gateway fails closed without a picker,
# whatever the pool's failureMode says, so the replica can't serve
# until its picker is available.
for key in self.rsp.desired.resources:
if key not in self.req.observed.resources:
continue
if key in workloads:
if resource.get_condition(self.req.observed.resources.get(key), "Ready").status == "True":
self.rsp.desired.resources[key].ready = fnv1.READY_TRUE
else:
if resource.get_condition(self.req.observed.resources.get(key), "Ready").status == "True":
self.rsp.desired.resources[key].ready = fnv1.READY_TRUE

def _workload_accepted(self, key: str) -> bool:
Expand Down
111 changes: 85 additions & 26 deletions functions/compose-model-replica/tests/test_fn.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,29 @@ def setUpModule() -> None:
logging.configure(level=logging.Level.DISABLED)


def _observed_object(*, ready: bool) -> fnv1.Resource:
"""A composed provider-kubernetes Object as observed back, with the Ready
condition its readiness policy derives."""
return fnv1.Resource(
resource=resource.dict_to_struct(
{
"apiVersion": "kubernetes.m.crossplane.io/v1alpha1",
"kind": "Object",
"status": {
"conditions": [
{
"type": "Ready",
"status": "True" if ready else "False",
"reason": "Available" if ready else "Unavailable",
"lastTransitionTime": "2025-01-01T00:00:00Z",
},
],
},
}
),
)


class TestFunctionRunner(unittest.IsolatedAsyncioTestCase):
"""Tests for FunctionRunner.RunFunction."""

Expand Down Expand Up @@ -369,7 +392,9 @@ async def test_compose(self) -> None:

want2 = fnv1.RunFunctionResponse(
meta=fnv1.ResponseMeta(ttl=durationpb.Duration(seconds=60)),
desired=fnv1.State(),
# Nothing is composed while waiting, so the XR is marked not ready
# rather than left to aggregate to trivially ready.
desired=fnv1.State(composite=fnv1.Resource(ready=fnv1.READY_FALSE)),
conditions=[
fnv1.Condition(
type="ModelAccepted",
Expand Down Expand Up @@ -415,7 +440,7 @@ async def test_compose(self) -> None:

want3 = fnv1.RunFunctionResponse(
meta=fnv1.ResponseMeta(ttl=durationpb.Duration(seconds=60)),
desired=fnv1.State(),
desired=fnv1.State(composite=fnv1.Resource(ready=fnv1.READY_FALSE)),
conditions=[
fnv1.Condition(
type="ModelAccepted",
Expand All @@ -438,12 +463,29 @@ async def test_compose(self) -> None:
)
want3.requirements.resources["cluster"].CopyFrom(cluster_requirement)

# Unified routing fronts the serving pods with an InferencePool + endpoint
# picker; their manifests are asserted in detail in test_backends. Here we
# only check the function wired the whole set in (and dropped the plain
# Service), then drop their manifests so the golden covers the dispatch,
# wiring and readiness the function itself owns.
routing_keys = {
"inference-pool",
"epp",
"epp-config",
"epp-role",
"epp-rolebinding",
"epp-serviceaccount",
"epp-service",
}
for key in routing_keys:
want1.desired.resources[key].CopyFrom(fnv1.Resource())

# Case 4: the resources from case 1 now exist in observed, and the
# workload Object reports Available (so its derived Ready is True). The
# function marks every composed resource ready (it can now observe them),
# the workload because it's serving and the rest because existing is
# being ready for them. Built from case 1, mutating only what the
# observed-ready transition changes: the four ready flags, the
# function marks each observed resource ready once its Object reports
# Ready: the workload because it's serving and the rest because existing
# is being ready for them. Built from case 1, mutating only what the
# observed-ready transition changes: the three ready flags, the
# acceptance/readiness conditions, and the dropped first-reconcile event.
req4 = fnv1.RunFunctionRequest()
req4.CopyFrom(req1)
Expand Down Expand Up @@ -471,12 +513,12 @@ async def test_compose(self) -> None:
),
)
)
# The other two are observed simply by being present; their content
# doesn't matter, only that the function can see them. (The InferencePool +
# endpoint picker resources aren't observed here, so they stay unready and
# are excluded from the golden below - test_backends covers their manifests.)
# The other two have no runtime readiness to wait on, so under their
# SuccessfulCreate policy provider-kubernetes reports them Ready once
# applied. (The InferencePool + endpoint picker resources aren't
# observed here, so they stay unready.)
for key in ("model-route", "resource-claim-main-standalone"):
req4.observed.resources[key].CopyFrom(fnv1.Resource(resource=structpb.Struct()))
req4.observed.resources[key].CopyFrom(_observed_object(ready=True))

want4 = fnv1.RunFunctionResponse()
want4.CopyFrom(want1)
Expand All @@ -493,27 +535,44 @@ async def test_compose(self) -> None:
# not yet observed), so it's gone now.
del want4.results[:]

# Case 5: everything from case 4 plus the routing objects is observed,
# but the endpoint picker's Service failed to apply, say because its
# name was invalid, so its Object isn't Ready. Being observed isn't
# being applied, so it stays unready and holds the replica unready with
# it. Built from case 4, mutating only the routing objects' ready flags.
req5 = fnv1.RunFunctionRequest()
req5.CopyFrom(req4)
for key in routing_keys:
req5.observed.resources[key].CopyFrom(_observed_object(ready=key != "epp-service"))

want5 = fnv1.RunFunctionResponse()
want5.CopyFrom(want4)
for key in routing_keys - {"epp-service"}:
want5.desired.resources[key].ready = fnv1.READY_TRUE

# Case 6: as case 5, but everything applied and the endpoint picker's
# Deployment isn't Available yet, so its Object's CEL-derived Ready is
# False. The gateway fails closed without a picker, so the replica stays
# unready with it.
req6 = fnv1.RunFunctionRequest()
req6.CopyFrom(req4)
for key in routing_keys:
req6.observed.resources[key].CopyFrom(_observed_object(ready=key != "epp"))

want6 = fnv1.RunFunctionResponse()
want6.CopyFrom(want4)
for key in routing_keys - {"epp"}:
want6.desired.resources[key].ready = fnv1.READY_TRUE

cases = [
Case(name="cluster ready composes native Deployment", req=req1, want=want1),
Case(name="cluster not resolved returns waiting conditions", req=req2, want=want2),
Case(name="cluster without providerConfigRef returns waiting conditions", req=req3, want=want3),
Case(name="observed resources are marked ready", req=req4, want=want4),
Case(name="an object that failed to apply stays unready", req=req5, want=want5),
Case(name="an unavailable endpoint picker stays unready", req=req6, want=want6),
]

# Unified routing fronts the serving pods with an InferencePool + endpoint
# picker; their manifests are asserted in detail in test_backends. Here we
# only check the function wired the whole set in (and dropped the plain
# Service), then drop them so the golden covers the dispatch and wiring the
# function itself owns.
routing_keys = {
"inference-pool",
"epp",
"epp-config",
"epp-role",
"epp-rolebinding",
"epp-serviceaccount",
"epp-service",
}
for case in cases:
with self.subTest(case.name):
got = await self.runner.RunFunction(case.req, None)
Expand All @@ -527,7 +586,7 @@ async def test_compose(self) -> None:
for key in routing_keys:
manifest = resources[key]["resource"]["spec"]["forProvider"]["manifest"]
self.assertEqual(manifest["metadata"]["namespace"], "mp-ml-team-51733", key)
resources.pop(key, None)
del resources[key]["resource"]
self.assertEqual(
json_format.MessageToDict(case.want),
got_dict,
Expand Down
Loading