monarch.actor#
The monarch.actor module provides the actor-based programming model for distributed computation. See Getting Started for an overview.
Creating Actors#
Actors are created on multidimensional meshes of processes that are launched across hosts. HostMesh represents a mesh of hosts. ProcMesh is a mesh of processes.
- class monarch.actor.HostMesh(hy_host_mesh, region, stream_logs, is_fake_in_process, code_sync_proc_mesh)[source]#
Bases:
MeshTraitHostMesh represents a collection of compute hosts that can be used to spawn processes and actors.
Can be used as an async context manager, which will shut down the hosts on exit: ``` host_mesh = job.state().hosts async with host_mesh:
# spawn proc meshes, actor meshes, etc.
# shutdown() is called automatically on exit. ```
If you don’t want to shutdown the hosts, you don’t need to use it as a context manager.
- spawn_procs(per_host=None, bootstrap=None, name=None, proc_bind=None, bootstrap_command=None)[source]#
Spawn a ProcMesh onto this host mesh.
- Parameters:
per_host (Dict[str, int] | None) – shape of procs per host, e.g.
{"gpus": 4}.bootstrap (Callable[[], None] | Callable[[], Awaitable[None]] | None) – optional setup callable run on each proc.
name (str | None) – optional name for the proc mesh.
proc_bind (list[dict[str, str]] | None) – optional per-process CPU/NUMA binding config. Length must equal
math.prod(per_host.values()). Each dict maps binding keys (cpunodebind,membind,physcpubind,cpus) to values.bootstrap_command (BootstrapCommand | Callable[[Point], BootstrapCommand] | None) – optional BootstrapCommand or callable that returns a BootstrapCommand for each coordinate. The callable receives a
Point(combined coordinate across host and per_host dimensions). This allows full customization of the bootstrap command per coordinate.
- property region: Region#
- with_python_executable(python_executable)[source]#
Return a new HostMesh that will use the given Python executable when spawning procs. Procs spawned from this mesh will also inherit this Python executable when they call
this_host().spawn_procs(...).
- shutdown()[source]#
Shutdown the host mesh and all of its processes. It will throw an exception if this host mesh is a reference rather than owned, which can happen if this HostMesh object was received from a remote actor or if it was produced by slicing. After shutting down, the hosts in this mesh will be unusable, and no new HostMeshes will be able to connect to them. If you want to stop everything on the host but keep them available for new clients, use stop() instead.
This is run automatically on __aexit__ when used as an async context manager.
- Returns:
A future that completes when the host mesh has been shut down.
- Return type:
Future[None]
- stop()[source]#
Stop the host mesh, releasing all resources but keeping worker processes alive for reconnection. A new HostMesh can be created that points to the same hosts.
Like shutdown, this throws if the host mesh is a reference rather than owned.
- Returns:
A future that completes when the host mesh has been stopped.
- Return type:
Future[None]
- async sync_workspace(workspace, conda=False, auto_reload=False)[source]#
Sync local code changes to the remote hosts.
- property initialized: Future[Literal[True]]#
Future completes with ‘True’ when the HostMesh has initialized. Because HostMesh are remote objects, there is no guarantee that the HostMesh is still usable after this completes, only that at some point in the past it was usable.
- flatten(name)#
Returns a new device mesh with all dimensions flattened into a single dimension with the given name.
Currently this supports only dense meshes: that is, all ranks must be contiguous in the mesh.
- rename(**kwargs)#
Returns a new device mesh with some of dimensions renamed. Dimensions not mentioned are retained:
new_mesh = mesh.rename(host=’dp’, gpu=’tp’)
- size(dim=None)#
Returns the number of elements (total) of the subset of mesh asked for. If dims is None, returns the total number of devices in the mesh.
- slice(**kwargs)#
Select along named dimensions. Integer values remove dimensions, slice objects keep dimensions but restrict them.
Examples: mesh.slice(batch=3, gpu=slice(2, 6))
- split(**kwargs)#
Returns a new device mesh with some dimensions of this mesh split. For instance, this call splits the host dimension into dp and pp dimensions, The size of ‘pp’ is specified and the dimension size is derived from it:
new_mesh = mesh.split(host=(‘dp’, ‘pp’), gpu=(‘tp’,’cp’), pp=16, cp=2)
Dimensions not specified will remain unchanged.
- class monarch.actor.ProcMesh(hy_proc_mesh, host_mesh, region, root_region, _device_mesh=None)[source]#
Bases:
MeshTraitA distributed mesh of processes for actor computation.
ProcMesh represents a collection of processes that can spawn and manage actors. It provides the foundation for distributed actor systems by managing process allocation, lifecycle, and communication across multiple hosts and devices.
The ProcMesh supports spawning actors, monitoring process health, logging configuration, and code synchronization across distributed processes.
- property initialized: Future[Literal[True]]#
Future completes with ‘True’ when the ProcMesh has initialized. Because ProcMesh are remote objects, there is no guarantee that the ProcMesh is still usable after this completes, only that at some point in the past it was usable.
- spawn(name, Class, *args, **kwargs)[source]#
Spawn a T-typed actor mesh on the process mesh.
Args: - name: The name of the actor. - Class: The class of the actor to spawn. - args: Positional arguments to pass to the actor’s constructor. - kwargs: Keyword arguments to pass to the actor’s constructor.
Returns: - The actor mesh reference typed as T.
Note:
The method returns immediately, initializing the underlying actor instances asynchronously. Thus, return of this method does not guarantee the actor’s __init__ has be executed; but rather that __init__ will be executed before the first call to the actor’s endpoints.
Nonblocking enhances composition, permitting the user to easily pipeline mesh creation, for exmaple to construct complex mesh object graphs without introducing additional latency.
If __init__ fails, the actor will be stopped and a supervision event will be raised.
- classmethod from_host_mesh(host_mesh, hy_proc_mesh, region, setup=None, _attach_controller_controller=True)[source]#
- activate()[source]#
Activate the device mesh. Operations done from insided this context manager will be distributed tensor operations. Each operation will be excuted on each device in the mesh.
See https://meta-pytorch.org/monarch/generated/examples/distributed_tensors.html for more information
- with mesh.activate():
t = torch.rand(3, 4, device=”cuda”)
- stop(reason='stopped by client')[source]#
This will stop all processes (and actors) in the mesh and release any resources associated with the mesh.
- flatten(name)#
Returns a new device mesh with all dimensions flattened into a single dimension with the given name.
Currently this supports only dense meshes: that is, all ranks must be contiguous in the mesh.
- rename(**kwargs)#
Returns a new device mesh with some of dimensions renamed. Dimensions not mentioned are retained:
new_mesh = mesh.rename(host=’dp’, gpu=’tp’)
- size(dim=None)#
Returns the number of elements (total) of the subset of mesh asked for. If dims is None, returns the total number of devices in the mesh.
- slice(**kwargs)#
Select along named dimensions. Integer values remove dimensions, slice objects keep dimensions but restrict them.
Examples: mesh.slice(batch=3, gpu=slice(2, 6))
- split(**kwargs)#
Returns a new device mesh with some dimensions of this mesh split. For instance, this call splits the host dimension into dp and pp dimensions, The size of ‘pp’ is specified and the dimension size is derived from it:
new_mesh = mesh.split(host=(‘dp’, ‘pp’), gpu=(‘tp’,’cp’), pp=16, cp=2)
Dimensions not specified will remain unchanged.
- monarch.actor.get_or_spawn_controller(name, Class, *args, **kwargs)[source]#
Creates a singleton actor (controller) indexed by name, or if it already exists, returns the existing actor.
- Parameters:
name (str) – The unique name of the actor, used as a key for retrieval.
Class (Type) – The class of the actor to spawn. Must be a subclass of Actor.
*args (Any) – Positional arguments to pass to the actor constructor.
**kwargs (Any) – Keyword arguments to pass to the actor constructor.
- Returns:
A Future that resolves to a reference to the actor.
- Return type:
Future[TActor]
- monarch.actor.this_host()[source]#
The current machine.
This is just shorthand for looking it up via the context
- monarch.actor.this_proc()[source]#
The current singleton process that this specific actor is running on
- monarch.actor.default_bootstrap_cmd()[source]#
Get the default bootstrap command for the current environment.
Returns a BootstrapCommand configured with the current Python executable and environment. This can be used as a base for customization with
with_env()or by modifying its attributes directly.- Returns:
The default bootstrap command.
- Return type:
BootstrapCommand
- monarch.actor.hosts_from_config(name)[source]#
Get the host mesh ‘name’ from the monarch configuration for the project.
This config can be modified so that the same code can create meshes from scheduler sources, and different sizes etc.
WARNING: This function is a standin so that our getting_started example code works. The real implementation needs an RFC design.
- monarch.actor.enable_transport(transport)[source]#
Allow monarch to communicate with transport type ‘transport’ This must be called before any other calls in the monarch API. If it isn’t called, we will implicitly call monarch.enable_transport(ChannelTransport.Unix) on the first monarch call.
Currently only one transport type may be enabled at one time. In the future we may allow multiple to be enabled.
- Supported transport values:
ChannelTransport enum: ChannelTransport.Unix, ChannelTransport.TcpWithHostname, etc.
- string short cuts for the ChannelTransport enum:
“tcp”: ChannelTransport.TcpWithHostname
“ipc”: ChannelTransport.Unix
“metatls”: ChannelTransport.MetaTlsWithIpV6
“metatls-hostname”: ChannelTransport.MetaTlsWithHostname
“tls”: ChannelTransport.Tls (uses configurable TLS certs)
- ZMQ-style URL format string for explicit address, e.g.:
For Meta usage, use metatls-hostname
Defining Actors#
All actor classes subclass the Actor base object, which provides them mesh slicing API. Each publicly exposed function of the actor is annotated with @endpoint:
- class monarch.actor.Actor[source]#
Bases:
MeshTraitBase class for actors.
Subclass
Actorto define an actor, and decorate the methods that form its public API with@endpoint. Actors are spawned onto aProcMesh(for example withthis_proc().spawn(...)), after which their endpoints are invoked remotely through the messaging adverbs. Each actor processes its messages sequentially and participates in the supervision tree.- flatten(name)#
Returns a new device mesh with all dimensions flattened into a single dimension with the given name.
Currently this supports only dense meshes: that is, all ranks must be contiguous in the mesh.
- rename(**kwargs)#
Returns a new device mesh with some of dimensions renamed. Dimensions not mentioned are retained:
new_mesh = mesh.rename(host=’dp’, gpu=’tp’)
- size(dim=None)#
Returns the number of elements (total) of the subset of mesh asked for. If dims is None, returns the total number of devices in the mesh.
- slice(**kwargs)#
Select along named dimensions. Integer values remove dimensions, slice objects keep dimensions but restrict them.
Examples: mesh.slice(batch=3, gpu=slice(2, 6))
- split(**kwargs)#
Returns a new device mesh with some dimensions of this mesh split. For instance, this call splits the host dimension into dp and pp dimensions, The size of ‘pp’ is specified and the dimension size is derived from it:
new_mesh = mesh.split(host=(‘dp’, ‘pp’), gpu=(‘tp’,’cp’), pp=16, cp=2)
Dimensions not specified will remain unchanged.
- monarch.actor.endpoint(method: Callable[[Concatenate[Any, P]], Awaitable[R]], *, propagate: None | Literal['cached', 'inspect', 'mocked'] | Callable[[...], Any] = None, explicit_response_port: Literal[False] = False, instrument: bool = True) EndpointProperty[P, R][source]#
- monarch.actor.endpoint(method: Callable[[Concatenate[Any, P]], R], *, propagate: None | Literal['cached', 'inspect', 'mocked'] | Callable[[...], Any] = None, explicit_response_port: Literal[False] = False, instrument: bool = True) EndpointProperty[P, R]
- monarch.actor.endpoint(*, propagate: None | Literal['cached', 'inspect', 'mocked'] | Callable[[...], Any] = None, explicit_response_port: Literal[False] = False, instrument: bool = True) EndpointIfy
- monarch.actor.endpoint(method: Callable[Concatenate[Any, 'Port[R]', P], Awaitable[None]], *, propagate: None | Literal['cached', 'inspect', 'mocked'] | Callable[[...], Any] = None, explicit_response_port: Literal[True], instrument: bool = True) EndpointProperty[P, R]
- monarch.actor.endpoint(method: Callable[Concatenate[Any, 'Port[R]', P], None], *, propagate: None | Literal['cached', 'inspect', 'mocked'] | Callable[[...], Any] = None, explicit_response_port: Literal[True], instrument: bool = True) EndpointProperty[P, R]
- monarch.actor.endpoint(*, propagate: None | Literal['cached', 'inspect', 'mocked'] | Callable[[...], Any] = None, explicit_response_port: Literal[True], instrument: bool = True) PortedEndpointIfy
Mark an
Actormethod as an endpoint callable from other actors.An endpoint defines part of an actor’s public API. Once the actor is spawned onto a mesh, you invoke its endpoints through messaging adverbs (
call,call_one,choose,stream,broadcast, andrref) rather than calling the method directly. Each adverb controls how the message is delivered and how the response is returned. Endpoint methods may be synchronous orasync.Apply it bare or with options:
class Counter(Actor): @endpoint def increment(self) -> int: ... @endpoint(explicit_response_port=True) async def stream_results(self, port: Port[int]) -> None: ...
- Parameters:
method – The actor method to wrap. Supplied automatically when the decorator is applied without parentheses; leave it unset when passing options.
propagate – Controls how the tensor engine infers output tensor shapes for
rrefwithout running the endpoint. Pass a callable that takes the endpoint’s arguments (excludingself) and returns tensors of the shapes the endpoint would produce, or one of the strings"cached","inspect", or"mocked". Defaults toNone, which matters only for endpoints used with distributed tensors.explicit_response_port – When
True, the endpoint receives aPortas its first argument (afterself) and is responsible for sending its result through that port instead of returning a value. This supports custom response protocols, such as sending several responses or deferring one. The method’s return annotation should beNone.instrument – When
True(the default), wrap each invocation in a tracing span named for the method.
- Returns:
An
EndpointPropertydescriptor. Accessing it on a spawned actor yields anEndpointexposing the messaging adverbs.
- monarch.actor.concurrent_endpoint(method=None, *, propagate=None, explicit_response_port=False, instrument=True)[source]#
Run one async actor endpoint as an
asynciobackground task.The decorator accepts the same endpoint options as
@endpoint. It rewrites the endpoint to use an explicit response port, schedules the original async body as a task, and returns immediately so that the actor can handle another queued message.If two messages from the same source actor call
@concurrent_endpointmethods on the same target actor, the first endpoint body starts before the second. This is only a start-order guarantee: the first endpoint runs until its firstawait, not to completion, before the second starts.If you mix
@concurrent_endpointwith normal@endpointmethods, a normal endpoint that follows a concurrent endpoint may run before the concurrent endpoint body has started.Outstanding tasks created by this decorator are cancelled and awaited before user
__cleanup__runs. This is a cleanup-time guarantee only: meshes owned by the actor may already have been stopped by the core actor lifecycle before these tasks are cancelled. General pending tasks on the actor’s asyncio loop are cancelled later, after user__cleanup__completes.If the endpoint body raises an exception that it does not forward through its response port, the actor fails with a supervision error, just as an exception escaping an
@endpoint(explicit_response_port=True)method does. This follows the principle of no silent errors. To report an error to the caller without failing the actor, forward it through the port (for example,port.exception(e)); to recover, catch it inside the endpoint.More complex protocols can still use
@endpoint(explicit_response_port=True)and manage response ports manually.
Messaging Actor#
Messaging is done through the “adverbs” defined for each endpoint
- class monarch.actor.Endpoint(*args, **kwargs)[source]#
Bases:
Protocol[P,R]A callable endpoint on a spawned actor or actor mesh.
Accessing an
@endpointmethod on a spawned actor yields anEndpointrather than calling the method directly. Its methods are the messaging adverbs that send a message to the actor(s) and control how responses are returned:call,call_one,choose,stream,broadcast, andrref.- stream(*args, **kwargs)[source]#
Broadcasts to all actors and yields their responses as a stream.
This enables processing results from multiple actors incrementally as they become available. Returns an iterator of Future values.
- class monarch.actor.Future(*, coro)[source]#
Bases:
Generic[R]The Future class wraps a PythonTask, which is a handle to a asyncio coroutine running on the Tokio event loop. These coroutines do not use asyncio or asyncio.Future; instead, they are executed directly on the Tokio runtime. The Future class provides both synchronous (.get()) and asynchronous APIs (await) for interacting with these tasks.
- Parameters:
coro (Coroutine[Any, Any, R] | PythonTask[R]) – The coroutine or PythonTask representing the asynchronous computation.
- get(timeout=None)[source]#
Get the result of the Future.
Caveats:
This method is designed to be used in places where event loops are not available. Besides that, you should avoid using this method if possible. Instead, use as_asyncio() (or await). This is because when Future.get() is called from within an active event loop, it blocks synchronously and does not yield control. That may degrade performance by preventing other tasks from running, and can potentially cause deadlocks if this future depends on them.
A timeout never consumes the Future: on TimeoutError the underlying task keeps running, so a later get()/await still observes its result.
examples:
This is not recommended because fut.get() blocks the event loop and might lead to issues explained above. ``` def inner_func(fut):
result = fut.get() # …
- async def out_func(fut):
inner_func(fut)
This is okay because everything is running synchronously. ``` def inner_func(fut):
result = fut.get() # …
- def main():
# … inner_func(fut)
- as_asyncio()[source]#
Return a standard
asyncio.Futurethat resolves when this Future does.Requires a running event loop; off a loop it raises
RuntimeErrorwithout consuming the underlying task (the Future stays unawaited, so a laterget()still drives it). Observation is non-consuming: repeatedas_asyncio()/awaiteach return a fresh loop-local future observing the same result.
- class monarch.actor.ValueMesh(shape, values)[source]#
Bases:
object- get(rank)#
Get value by linear rank (0..num_ranks-1).
- values()#
Return the values in region/iteration order as a Python list.
- flatten(name)#
Returns a new device mesh with all dimensions flattened into a single dimension with the given name.
Currently this supports only dense meshes: that is, all ranks must be contiguous in the mesh.
- static from_indexed(shape, pairs)#
Build from (rank, value) pairs with last-write-wins semantics.
- rename(**kwargs)#
Returns a new device mesh with some of dimensions renamed. Dimensions not mentioned are retained:
new_mesh = mesh.rename(host=’dp’, gpu=’tp’)
- size(dim=None)#
Returns the number of elements (total) of the subset of mesh asked for. If dims is None, returns the total number of devices in the mesh.
- slice(**kwargs)#
Select along named dimensions. Integer values remove dimensions, slice objects keep dimensions but restrict them.
Examples: mesh.slice(batch=3, gpu=slice(2, 6))
- split(**kwargs)#
Returns a new device mesh with some dimensions of this mesh split. For instance, this call splits the host dimension into dp and pp dimensions, The size of ‘pp’ is specified and the dimension size is derived from it:
new_mesh = mesh.split(host=(‘dp’, ‘pp’), gpu=(‘tp’,’cp’), pp=16, cp=2)
Dimensions not specified will remain unchanged.
- class monarch.actor.ActorError(exception, message='A remote actor call has failed.')[source]#
Bases:
ExceptionDeterministic problem with the user’s code. For example, an OOM resulting in trying to allocate too much GPU memory, or violating some invariant enforced by the various APIs.
- add_note()#
Exception.add_note(note) – add a note to the exception
- args#
- with_traceback()#
Exception.with_traceback(tb) – set self.__traceback__ to tb and return self.
- class monarch.actor.Accumulator(endpoint, identity, combine)[source]#
Bases:
Generic[P,R,A]Accumulate the result of a broadcast invocation of an endpoint across a sliced mesh.
- Usage:
>>> counter = Accumulator(Actor.increment, 0, lambda x, y: x + y)
- monarch.actor.send(endpoint, args, kwargs, port=None, selection='all')[source]#
Fire-and-forget broadcast invocation of the endpoint across a given selection of the mesh.
This sends the message to all actors but does not wait for any result. Use the port provided to send the response back to the caller.
- class monarch.actor.Channel[source]#
Bases:
Generic[R]An advanced low level API for a communication channel used for message passing between actors.
Provides static methods to create communication channels with port pairs for sending and receiving messages of type R.
- class monarch.actor.Port(port_ref, instance, rank)[source]#
Bases:
objectA port that sends messages to a remote receiver. Wraps an EitherPortRef with the actor instance needed for sending.
- send(obj)#
- send_message(message)#
- exception(e)#
- return_undeliverable#
- class monarch.actor.PortReceiver(mailbox, receiver)[source]#
Bases:
Generic[R]Receiver for messages sent through a communication channel.
Handles receiving R-typed objects sent from a corresponding Port.
- monarch.actor.as_endpoint(not_an_endpoint: Callable[[P], R], *, propagate: None | Literal['cached', 'inspect', 'mocked'] | Callable[[...], Any] = None, explicit_response_port: Literal[False] = False) Endpoint[P, R][source]#
- monarch.actor.as_endpoint(not_an_endpoint: Callable[Concatenate['PortProtocol[R]', P], None], *, propagate: None | Literal['cached', 'inspect', 'mocked'] | Callable[[...], Any] = None, explicit_response_port: Literal[True]) Endpoint[P, R]
Treat an actor method that is not an
@endpointas one.Use this to call a plain method of a spawned actor through the messaging adverbs when the method was not decorated with
@endpoint, for exampleas_endpoint(actor.method).call(...). The options match those ofendpoint.- Parameters:
not_an_endpoint – An unannotated method of a spawned actor.
propagate – Tensor-shape propagation for
rref; seeendpoint.explicit_response_port – When
True, the method receives aPortas its first argument and sends its result through it; seeendpoint.
- Returns:
An
Endpointexposing the messaging adverbs.- Raises:
ValueError – If
not_an_endpointis not a method of a spawned actor.
Context API#
Use these functions to look up what actor is running the currently executing code.
- monarch.actor.current_actor_name()[source]#
Return the actor id of the currently executing actor as a string.
- monarch.actor.current_rank()[source]#
Return the current message’s position within its mesh as a
Point.
- monarch.actor.current_size()[source]#
Return the size of each mesh dimension for the current message.
The result maps each dimension label to the number of actors along it.
- monarch.actor.context()[source]#
Return the
Contextfor the currently executing actor.Call this from within an endpoint to inspect the running actor and the current message’s position in the mesh. Outside an actor (on the client) it returns the root client context.
- monarch.actor.shutdown_context()[source]#
Shutdown global actor context resources.
Idempotent: subsequent calls return an immediately-resolved future. This is safe to call both explicitly and from atexit.
- Returns:
- A future that completes when shutdown is
finished. Call with .get() to wait for completion.
- Return type:
Future[None]
- class monarch.actor.Point(rank, extent)[source]#
Bases:
Point,MappingA coordinate within a mesh.
A
Pointmaps each mesh dimension label (for example"hosts"or"gpus") to its integer index along that dimension. It behaves as a read-only mapping, sopoint["gpus"]is the index along the"gpus"dimension.current_rankandContext.message_rankreturn thePointthat locates the current actor within the mesh that received the message.- extent#
- get(k[, d]) D[k] if k in D, else d. d defaults to None.#
- items() a set-like object providing a view on D's items#
- keys() a set-like object providing a view on D's keys#
- rank#
- size(label)#
- values() an object providing a view on D's values#
Supervision#
Types used for error handling and supervision in actor meshes.
- monarch.actor.unhandled_fault_hook(failure)[source]#
When a supervision event is unhandled and is propagated back to the client, this hook is called. The default implementation is to exit the process with error code 1 after logging the event. If this function raises any exception (including BaseException classes such as SystemExit from sys.exit), this fault is considered unhandled. Any normal return value will cause the fault to be dropped. Logs will be written containing the failure message in either case.
To customize this behavior, overwrite this function in your client code like so: ``` import monarch.actor
- def my_unhandled_fault_hook(failure: MeshFailure) -> None:
# log it, add metrics, etc. print(f”Mesh failure was not handled: {failure}”) # To ignore this error, return any value (including None) without an exception.
monarch.actor.unhandled_fault_hook = my_unhandled_fault_hook ```
If the fault is unhandled, it exits the main thread by delivering a KeyboardInterrupt. This is done because the Python Interpreter can only be finalized to run destructors and atexit hooks from the main thread. So if you see a “KeyboardInterrupt” happening that you didn’t send, it’s because there was an unhandled fault.
Telemetry#
Utilities for tracing actor execution.
- monarch.actor.traced(fn: Callable[[_P], _R]) Callable[[_P], _R][source]#
- monarch.actor.traced(*, name: str) Callable[[Callable[[_P], _R]], Callable[[_P], _R]]
Decorator that wraps a function in a telemetry span.
Works with both sync and async functions. The span is automatically associated with the current actor context, if any. When no name is provided, the function’s
__name__is used as the span name.Usage:
@traced async def do_work(): ... @traced(name="custom_name") def compute(): ...