22# Licensed under the MIT License.
33
44from dataclasses import dataclass , field
5+ from decimal import Decimal , InvalidOperation
56from typing import Callable , Iterable , Optional , Sequence
67
78import grpc
@@ -134,7 +135,7 @@ def build_serverless_activity_declaration(
134135 environment_variables : Optional [dict [str , str ]] = None ,
135136 max_concurrent_activities : int = DEFAULT_MAX_CONCURRENT_ACTIVITIES ,
136137 entrypoint : Optional [Iterable [str ]] = None ,
137- cmd : Optional [Iterable [str ]] = None ) -> pb .ServerlessActivityDeclaration :
138+ cmd : Optional [Iterable [str ]] = None ) -> pb .OnDemandSandboxActivityDeclaration :
138139 resolved_activity_names = resolve_activity_names (activity_names )
139140 if not resolved_activity_names :
140141 raise ValueError ("Serverless activity declaration requires at least one activity name." )
@@ -154,19 +155,16 @@ def build_serverless_activity_declaration(
154155 if not image_ref :
155156 raise ValueError ("Serverless activity image metadata requires a container image reference." )
156157
157- if not cpu or not cpu . strip ():
158- raise ValueError ( "Serverless activity declaration requires CPU resources." )
158+ resolved_cpu = _normalize_cpu ( cpu )
159+ resolved_memory = _normalize_memory ( memory )
159160
160- if not memory or not memory .strip ():
161- raise ValueError ("Serverless activity declaration requires memory resources." )
162-
163- declaration = pb .ServerlessActivityDeclaration (
161+ declaration = pb .OnDemandSandboxActivityDeclaration (
164162 worker_profile_id = worker_profile_id .strip (),
165- image = pb .ServerlessActivityImage (
163+ image = pb .OnDemandSandboxActivityImage (
166164 image_ref = image_ref ),
167- resources = pb .ServerlessActivityResources (
168- cpu = cpu . strip () ,
169- memory = memory . strip () ),
165+ resources = pb .OnDemandSandboxActivityResources (
166+ cpu = resolved_cpu ,
167+ memory = resolved_memory ),
170168 max_concurrent_activities = max_concurrent_activities )
171169 declaration .activity_names .extend (resolved_activity_names )
172170 declaration .environment_variables .update (environment_variables or {})
@@ -175,9 +173,9 @@ def build_serverless_activity_declaration(
175173 return declaration
176174
177175
178- def build_profile_serverless_activity_declarations () -> list [pb .ServerlessActivityDeclaration ]:
176+ def build_profile_serverless_activity_declarations () -> list [pb .OnDemandSandboxActivityDeclaration ]:
179177 """Build serverless declarations from worker profile configuration."""
180- declarations : list [pb .ServerlessActivityDeclaration ] = []
178+ declarations : list [pb .OnDemandSandboxActivityDeclaration ] = []
181179 activity_owners : dict [str , str ] = {}
182180 for profile in _worker_profiles .values ():
183181 activity_names = resolve_activity_names (profile .activity_names )
@@ -217,7 +215,7 @@ def build_serverless_worker_start(
217215 max_activities_count : int ,
218216 activity_names : Iterable [str ],
219217 substrate : Optional [str ] = None ,
220- dts_sandbox_identifier : Optional [str ] = None ) -> pb .ServerlessActivityWorkerMessage :
218+ dts_sandbox_identifier : Optional [str ] = None ) -> pb .OnDemandSandboxActivityWorkerMessage :
221219 if not taskhub or not taskhub .strip ():
222220 raise ValueError ("Serverless activity worker registration requires a task hub name." )
223221
@@ -231,8 +229,8 @@ def build_serverless_worker_start(
231229 if not resolved_activity_names :
232230 raise ValueError ("Serverless activity worker registration requires at least one registered activity." )
233231
234- message = pb .ServerlessActivityWorkerMessage (
235- start = pb .ServerlessActivityWorkerStart (
232+ message = pb .OnDemandSandboxActivityWorkerMessage (
233+ start = pb .OnDemandSandboxActivityWorkerStart (
236234 task_hub = taskhub .strip (),
237235 worker_profile_id = worker_profile_id .strip (),
238236 max_activities_count = max_activities_count ,
@@ -242,12 +240,12 @@ def build_serverless_worker_start(
242240 return message
243241
244242
245- def build_serverless_worker_heartbeat (active_activities_count : int ) -> pb .ServerlessActivityWorkerMessage :
243+ def build_serverless_worker_heartbeat (active_activities_count : int ) -> pb .OnDemandSandboxActivityWorkerMessage :
246244 if active_activities_count < 0 :
247245 raise ValueError ("Serverless activity worker active activity count cannot be negative." )
248246
249- return pb .ServerlessActivityWorkerMessage (
250- heartbeat = pb .ServerlessActivityWorkerHeartbeat (
247+ return pb .OnDemandSandboxActivityWorkerMessage (
248+ heartbeat = pb .OnDemandSandboxActivityWorkerHeartbeat (
251249 active_activities_count = active_activities_count ))
252250
253251
@@ -278,7 +276,7 @@ def __init__(
278276 interceptors = resolved_interceptors ,
279277 channel_options = channel_options )
280278 self ._channel = channel
281- self ._stub = stubs .ServerlessActivitiesStub (channel )
279+ self ._stub = stubs .OnDemandSandboxActivitiesStub (channel )
282280
283281 def close (self ) -> None :
284282 if self ._owns_channel :
@@ -291,17 +289,18 @@ def enable_serverless_activities(self) -> None:
291289 raise ValueError ("No configured serverless activities were found." )
292290
293291 for declaration in declarations :
294- self ._stub .DeclareServerlessActivities (declaration )
292+ self ._stub .DeclareOnDemandSandboxActivities (declaration )
295293
296294 def remove_serverless_activity_declaration (self , worker_profile_id : str ) -> None :
297295 worker_profile_id = _normalize_required (worker_profile_id , "Worker profile ID is required." )
298- self ._stub .RemoveServerlessActivityDeclaration (
299- pb .RemoveServerlessActivityDeclarationRequest (worker_profile_id = worker_profile_id ))
296+ self ._stub .RemoveOnDemandSandboxActivityDeclaration (
297+ pb .RemoveOnDemandSandboxActivityDeclarationRequest (worker_profile_id = worker_profile_id ))
300298
301299 def connect_serverless_activity_worker (
302300 self ,
303- messages : Iterable [pb .ServerlessActivityWorkerMessage ]) -> pb .ServerlessActivityWorkerSessionResult :
304- return self ._stub .ConnectServerlessActivityWorker (messages )
301+ messages : Iterable [pb .OnDemandSandboxActivityWorkerMessage ]
302+ ) -> pb .OnDemandSandboxActivityWorkerSessionResult :
303+ return self ._stub .ConnectOnDemandSandboxActivityWorker (messages )
305304
306305
307306def _normalize_optional_strings (values : Iterable [str ]) -> list [str ]:
@@ -314,6 +313,46 @@ def _normalize_required(value: str, message: str) -> str:
314313 return value .strip ()
315314
316315
316+ def _normalize_cpu (value : str ) -> str :
317+ normalized = _normalize_required (value , "Serverless activity declaration requires CPU resources." )
318+ milli_cpu = _try_parse_cpu_millicores (normalized )
319+ if milli_cpu is None or milli_cpu <= 0 :
320+ raise ValueError (
321+ "Serverless activity CPU resources must be a positive Kubernetes-style CPU quantity. "
322+ "Use formats like '500m', '2', or '0.5'." )
323+ return normalized
324+
325+
326+ def _normalize_memory (value : str ) -> str :
327+ normalized = _normalize_required (value , "Serverless activity declaration requires memory resources." )
328+ memory_mib = _try_parse_memory_mib (normalized )
329+ if memory_mib is None or memory_mib <= 0 :
330+ raise ValueError (
331+ "Serverless activity memory resources must be a positive Kubernetes-style memory quantity. "
332+ "Use formats like '256Mi', '1Gi', or '2048'." )
333+ return normalized
334+
335+
336+ def _try_parse_cpu_millicores (value : str ) -> Optional [int ]:
337+ try :
338+ if value [- 1 :].lower () == "m" :
339+ return int (Decimal (value [:- 1 ]))
340+ return int (Decimal (value ) * 1000 )
341+ except (InvalidOperation , ValueError ):
342+ return None
343+
344+
345+ def _try_parse_memory_mib (value : str ) -> Optional [int ]:
346+ try :
347+ if value [- 2 :].lower () == "gi" :
348+ return int (Decimal (value [:- 2 ]) * 1024 )
349+ if value [- 2 :].lower () == "mi" :
350+ return int (Decimal (value [:- 2 ]))
351+ return int (Decimal (value ))
352+ except (InvalidOperation , ValueError ):
353+ return None
354+
355+
317356def _parse_substrate (substrate : Optional [str ]) -> "pb.SubstrateKind" :
318357 if not substrate :
319358 return pb .SUBSTRATE_KIND_UNSPECIFIED
0 commit comments