Primeros pasos¶
Instalación¶
pip install lapinbeam
Para desarrollar contra el propio código fuente en su lugar:
uv sync # crea .venv, compila la extensión Rust, instala deps
uv run maturin develop # recompila la extensión tras tocar código Rust
uv run pytest # suite de tests en Python
cargo test # suite de tests en Rust
No se instala nada a nivel de sistema: todo vive en .venv.
Conceptos clave¶
| Concepto | Qué es |
|---|---|
Node |
El extremo local del clúster: nombre@host:puerto. Posee el runtime de Tokio en segundo plano, el listener y las conexiones a peers. |
@actor |
Marca una clase como actor. Supervisor.spawn lee estos metadatos; no es un envoltorio de runtime. |
Supervisor |
Crea actores y los reinicia ante excepciones no controladas (estrategia one_for_one, con reinicios limitados y backoff). |
ActorRef / RemoteRef |
Un handle para enviar mensajes a un actor local o remoto. Ambos exponen el mismo await ref.send(msg). |
Un único actor¶
import asyncio
from lapinbeam import ActorRef, Node, Supervisor, actor
@actor(name="echo")
class Echo:
async def receive(self, msg):
print("recibido:", msg)
async def main():
node = Node("app@127.0.0.1:0") # el puerto 0 elige uno efímero
await node.start()
sup = Supervisor(strategy="one_for_one", node=node)
echo: ActorRef = sup.spawn(Echo)
await echo.send({"hello": "world"})
await asyncio.sleep(0.1) # deja que la mailbox se vacíe antes de parar
await node.stop()
asyncio.run(main())
Node también funciona como gestor de contexto asíncrono, la forma más
idiomática para cualquier cosa de vida más larga:
async with Node("app@127.0.0.1:0") as node:
sup = Supervisor(node=node)
ref = sup.spawn(Echo)
await ref.send({"hello": "world"})
await asyncio.sleep(0.1)
# node.stop() se ejecuta automáticamente al salir, incluso si el bloque lanza una excepción.
node.stop() también cancela cualquier tarea de actor lanzada por
cualquier Supervisor en ese nodo — ninguna se queda corriendo para
siempre, bloqueada en un mailbox que nadie va a rellenar nunca. Para tirar
abajo solo los actores de un Supervisor en concreto en vez de todo el
nodo (p.ej. varios supervisores compartiendo un mismo nodo), llama
directamente a await sup.shutdown().
Concurrencia: un actor procesa un mensaje cada vez¶
Cada actor tiene exactamente un mailbox y exactamente una tarea leyendo de
él, en un bucle: coge un mensaje, ejecuta el handler, espera a que
termine, coge el siguiente. send() no cambia eso — solo encola el
mensaje y devuelve el control al instante, sin importar cuántos envíes
seguidos:
import asyncio
from lapinbeam import ActorRef, Node, Supervisor, actor
@actor(name="processor")
class Processor:
async def receive(self, msg):
await asyncio.sleep(1) # p.ej. una llamada lenta a otro servicio
print("terminado con", msg["id"])
async def main():
async with Node("app@127.0.0.1:0") as node:
sup = Supervisor(node=node)
ref: ActorRef = sup.spawn(Processor)
for i in range(5):
await ref.send({"id": i}) # cada uno vuelve al instante...
await asyncio.sleep(6) # ...pero este actor tarda ~5s en total igualmente
Los cinco send() devuelven el control en una fracción de segundo, pero
Processor los termina uno a uno — el quinto llega sobre el segundo 5,
no el primero. Esto es deliberado, no una limitación: como dentro de un
mismo actor nunca hay más de un mensaje "en vuelo" a la vez, el código
del handler puede leer y escribir self.lo_que_sea libremente, sin
locks — la misma garantía que dan los procesos de Erlang.
Si quieres que varias de esas llamadas corran de verdad a la vez, crea un
pool de actores en vez de esperar que un solo actor se paralelice a
sí mismo. Supervisor.spawn_pool() hace justo eso — n_workers actores
creados una vez, compartiendo una cola interna, y el que esté libre coge
el siguiente mensaje:
from lapinbeam import PoolRef
async def process(msg):
await asyncio.sleep(1)
print("terminado con", msg["id"])
async def main():
async with Node("app@127.0.0.1:0") as node:
sup = Supervisor(node=node)
pool: PoolRef = await sup.spawn_pool(process, 5, name="processors")
for i in range(5):
await pool.send({"id": i})
await asyncio.sleep(2) # ahora ~1s en total, no ~5s
Cada actor del pool tiene su propio mailbox y su propia tarea, así que sus
cinco asyncio.sleep(1) se solapan de verdad — los cinco terminan sobre
el segundo 1 en vez del segundo 5. current_message()/.reply() y
ask()/ask_stream() enviados a pool funcionan exactamente igual que
contra un actor spawn()eado normal, sin importar qué worker acabe
respondiendo. Un process que lanza una excepción no tira abajo a su
worker — se captura internamente y se reporta vía
on_event(kind="pool_worker_error"), y ese worker sigue con el siguiente
mensaje de la cola.
process de arriba es una función suelta, así que los workers no guardan
estado entre llamadas. Si cada worker debe mantener su propio estado
entre los mensajes que le vayan tocando (una caché, un contador, una
conexión), pasa una clase @actor en su lugar — spawn_pool() crea una
instancia por worker (args/kwargs van a su constructor, igual que en
spawn()) y despacha cada mensaje a través de los handlers @on de esa
instancia, o a receive si no define ningún @on:
from lapinbeam import PoolRef, actor
@actor(name="processor")
class Processor:
def __init__(self):
self.procesados = 0
async def receive(self, msg):
await asyncio.sleep(1)
self.procesados += 1
print("terminado con", msg["id"], "— este worker lleva", self.procesados)
async def main():
async with Node("app@127.0.0.1:0") as node:
sup = Supervisor(node=node)
pool: PoolRef = await sup.spawn_pool(Processor, 5, name="processors")
for i in range(5):
await pool.send({"id": i})
await asyncio.sleep(2)
Misma cola, misma concurrencia, mismo salvavidas
on_event(kind="pool_worker_error") — pero ahora cada una de las 5
instancias de Processor mantiene procesados a lo largo de todos los
mensajes que esa instancia llegue a procesar. Ese salvavidas importa más
aquí que con una función: si un handler lanza una excepción, la misma
instancia sigue funcionando después, así que cualquier estado de self a
medio actualizar se queda a medio actualizar para el siguiente mensaje que
le toque a ese worker (a diferencia de un actor spawn()eado, cuyo
reinicio siempre parte de una instancia nueva).
Cuándo usar spawn_pool(): llegan más elementos de trabajo de los
que un solo actor podría procesar, pero cada uno es lo bastante barato (o
I/O-bound — ver el aviso de abajo) como para que unos pocos workers den
abasto, y no te importa cuál de ellos procese cada elemento — usa la
forma función para trabajo sin estado, y la forma clase cuando los workers
deban mantener estado entre mensajes. Cuándo no: si necesitas esperar
varias respuestas independientes a la vez en vez de que un pool
responda a un único ask(), la herramienta es asyncio.gather() sobre
varios ask() (con pool o sin él) — ver el patrón de mixture-of-experts
en Agentes de IA y MCP.
Este paralelismo es para trabajo I/O-bound, no CPU-bound — salvo que pidas un executor
Un pool ayuda porque await asyncio.sleep(...) (o una llamada de red,
o cualquier otro await que ceda el control de verdad) permite a
asyncio intercalar la espera de cada actor en el mismo hilo. No
ayuda a un handler que hace trabajo de CPU síncrono de verdad, sin
ningún await dentro — el bucle de eventos de asyncio es de un solo
hilo, así que N actores machacando números siguen corriendo uno detrás
de otro, igual de lento que N llamadas secuenciales dentro de un solo
actor.
Para trabajo genuinamente CPU-bound, pasa executor="process" (o
"thread", para una llamada bloqueante que no es CPU-bound pero no
tiene equivalente async, p.ej. una librería C o un driver de BBDD
síncronos) — handler debe ser entonces una función síncrona
normal, ejecutada fuera del event loop en un
ProcessPoolExecutor/ThreadPoolExecutor real, dimensionado a
n_workers. Su valor de retorno se envía de vuelta automáticamente si
el mensaje llegó por ask()/ask_stream():
def crunch(msg): # def normal, no async def
return sum(i * i for i in range(msg["n"]))
async def main():
async with Node("app@127.0.0.1:0") as node:
sup = Supervisor(node=node)
pool = await sup.spawn_pool(crunch, 4, name="crunchers", executor="process")
result = await pool.ask({"n": 10_000_000})
if __name__ == "__main__": # necesario con executor="process" — ver abajo
asyncio.run(main())
executor="process" no admite una clase @actor como handler — no
hay forma de mantener estado de Python en self a través de una
frontera de proceso — y exige que handler y cada mensaje/args/
kwargs con el que se llame sean serializables con pickle (una
función a nivel de módulo, no un closure ni una lambda). También
hereda el requisito habitual de multiprocessing: el script de
entrada debe proteger su código de nivel de módulo con
if __name__ == "__main__": — sin eso, un proceso worker reimporta el
script como __main__ y vuelve a ejecutar todo lo de nivel de módulo
(incluido el propio asyncio.run(main())), lo que en el mejor caso
duplica trabajo y en el peor se queda colgado. executor="thread" no
tiene esta restricción — comparte el proceso padre en vez de lanzar
procesos nuevos. Una alternativa que evita ambas restricciones:
reparte el trabajo entre procesos de sistema operativo separados —
p.ej. varios Node de lapinbeam, posiblemente en máquinas distintas,
hablando por red igual que el ejemplo de dos nodos de abajo.
Backpressure: acotar la cola del pool¶
Por defecto la cola interna del pool no tiene límite — si pool.send() se
llama más rápido de lo que los workers pueden vaciarla, la cola sigue
creciendo. Pasa queue_capacity para acotarla, misma convención que
Node(mailbox_capacity=...): una vez llena, un send() de más se
descarta y se reporta vía on_event(kind="pool_queue_full") en vez de
consumir memoria sin límite:
pool = await sup.spawn_pool(process, 5, name="processors", queue_capacity=1000)
Orden por clave: pools particionados¶
El reparto por defecto es "el worker que esté libre" — bueno para
throughput, pero no da ninguna garantía de orden entre dos mensajes de la
misma entidad lógica (p.ej. dos actualizaciones del mismo order_id
podrían acabar procesadas fuera de orden por dos workers distintos). Pasa
key para particionar el pool: cada worker tiene su propia cola, y todo
mensaje con la misma clave siempre cae en el mismo worker, en el orden de
llegada, mientras que claves distintas siguen corriendo en paralelo:
pool: PoolRef = await sup.spawn_pool(
process, 5, name="processors", key=lambda msg: msg["order_id"]
)
Ahora todo mensaje de order_id=42 lo procesa un worker concreto, en
orden, mientras order_id=43 corre a la vez en otro distinto. El
compromiso: una distribución de claves desigual (unas pocas claves muy
"calientes") puede dejar workers ociosos mientras otros acumulan cola —
esta no es la herramienta para repartir throughput de forma uniforme, solo
para los casos en los que el orden importa más que la carga equilibrada.
Recoger varias respuestas: pool.map()¶
asyncio.gather() sobre varios ask() ya es la herramienta para
"esperar N respuestas independientes a la vez" (mencionado arriba) —
pool.map() es ese mismo patrón resumido en una sola llamada, cuando
todos los elementos van al mismo pool:
results = await pool.map([{"id": i} for i in range(5)])
Los resultados llegan en el mismo orden que los elementos de entrada, no
en orden de finalización — el tercer resultado siempre es la respuesta
del tercer elemento, aunque haya sido el primero en terminar. Pasa
return_exceptions=True para recibir en la lista cada resultado o
excepción en vez de dejar que el primer fallo (p.ej. un TimeoutError de
un elemento) aborte todo el lote.
Parar un solo pool: pool.stop()¶
Un pool creado con spawn_pool() normalmente vive lo mismo que su
Supervisor. Si un servidor crea y destruye pools a lo largo de su vida
en vez de mantener uno fijo desde el arranque, pool.stop() para solo el
dispatcher y los workers de ese pool (y apaga su
ThreadPoolExecutor/ProcessPoolExecutor, si tiene uno) sin tocar nada
más de ese mismo Supervisor — y libera name para que una llamada
posterior a spawn_pool() lo reutilice:
pool = await sup.spawn_pool(process, 5, name="processors")
...
await pool.stop()
O acotado automáticamente a un bloque — lo que devuelve spawn_pool() es
a la vez awaitable y gestor de contexto async:
async with sup.spawn_pool(process, 5, name="processors") as pool:
for i in range(5):
await pool.send({"id": i})
await asyncio.sleep(2)
# pool.stop() ya se ha llamado aquí
Dos nodos hablando entre sí¶
Esta es la demo de dos nodos de examples/, y lo único que merece la pena
dejar extra claro: un Node es un servidor/proceso, no un concepto que
viva solo en memoria. Abajo hay dos scripts separados —
app_node_a.py y app_node_b.py —, cada uno con solo el actor que ese
servidor necesita. No son dos ramas del mismo fichero: son dos ficheros,
pensados para correr como dos procesos python separados, potencialmente
en dos máquinas distintas.
sequenceDiagram
box Servidor A (node_a@host:9001)
participant Ingestor as Actor Ingestor
end
box Servidor B (node_b@host:9002)
participant Processor as Actor Processor
end
Note over Ingestor,Processor: node.connect_peer() abre una única conexión TCP,<br/>compartida por todos los actores de ambos servidores
Ingestor->>Processor: send(TASK, reply_to="ingestor")
Note right of Processor: Processor.receive(msg) se ejecuta aquí, en el servidor B
Processor->>Ingestor: send(ACK) — enrutado a "ingestor" por nombre
Note left of Ingestor: Ingestor.receive(msg) se ejecuta aquí, en el servidor A
El servidor A envía mensajes TASK al actor processor del servidor B; B
responde con ACK al actor que A haya indicado como reply_to — A no deja
fijo "responder a Ingestor", simplemente le dice a B a quién contestar, que
es justo lo que hace que el mismo código de Processor se pueda reutilizar
sin importar quién lo llame.
import asyncio
import os
from lapinbeam import Node, RemoteRef, Supervisor, actor
# Este actor existe solo en el servidor A. Su trabajo es recibir el
# ACK que el servidor B envía de vuelta una vez procesada una tarea.
@actor(name="ingestor")
class Ingestor:
async def receive(self, msg):
if msg.get("type") == "ACK":
print("ack recibido para", msg["payload_id"])
async def main():
node_name = os.environ["NODE_NAME"] # p.ej. node_a@127.0.0.1:9001 (este servidor)
peer = os.environ["PEER"] # p.ej. node_b@127.0.0.1:9002 (el otro servidor)
node = Node(node_name)
await node.start() # vincula el socket de escucha de ESTE servidor
sup = Supervisor(node=node)
sup.spawn(Ingestor) # registra el actor que recibirá los ACK
await node.connect_peer(peer) # marca al servidor B y espera el handshake TCP
remote: RemoteRef = node.get_remote_actor(peer, "processor") # una referencia al actor "processor" de B
for i in range(100):
# Cada send() de aquí cruza de verdad la red hasta el servidor B.
await remote.send({"type": "TASK", "payload_id": i, "reply_to": "ingestor"})
await asyncio.sleep(0.01)
asyncio.run(main())
import asyncio
import os
from lapinbeam import Node, RemoteRef, Supervisor, actor
# Este actor existe solo en el servidor B. Recibe los mensajes TASK
# que envía el servidor A, y responde con un ACK — no a una dirección
# fija, sino a cualquier nombre de actor que A haya puesto en
# `reply_to`.
@actor(name="processor")
class Processor:
def __init__(self, node_ref, peer_id):
# `node_ref` es el propio Node de ESTE proceso (el del
# servidor B) — se usa para enviar las respuestas. `peer_id`
# es el id del OTRO servidor (el de A); ambos servidores ya
# conocen la dirección del otro de antemano, gracias a las
# variables de entorno NODE_NAME/PEER de abajo.
self.node = node_ref
self.peer_id = peer_id
async def receive(self, msg):
if msg.get("type") == "TASK":
# get_remote_actor() NO abre una conexión nueva — solo
# construye una referencia que reutiliza la única
# conexión TCP ya establecida entre los dos servidores.
remote: RemoteRef = self.node.get_remote_actor(self.peer_id, msg["reply_to"])
await remote.send({"type": "ACK", "payload_id": msg["payload_id"]})
async def main():
node_name = os.environ["NODE_NAME"] # p.ej. node_b@127.0.0.1:9002 (este servidor)
peer = os.environ["PEER"] # p.ej. node_a@127.0.0.1:9001 (el otro servidor)
node = Node(node_name)
await node.start() # vincula el socket de escucha de ESTE servidor
sup = Supervisor(node=node)
sup.spawn(Processor, node, peer) # registra el actor que responde a los TASK
await node.wait_until_stopped() # el servidor B solo reacciona a mensajes entrantes; nunca marca hacia fuera
asyncio.run(main())
Fíjate en qué es igual y qué no: los dos scripts leen NODE_NAME/PEER del
entorno de la misma forma, pero no hay ninguna rama condicional en ningún
sitio — cada fichero solo desempeña un papel, exactamente igual que
examples/app_node_a.py y examples/app_node_b.py.
Processor recibe peer_id inyectado por su constructor porque es la
única forma que tiene de saber a qué nodo responder. Si un handler
prefiere averiguarlo a partir del propio mensaje en vez de depender de un
argumento del constructor, lapinbeam.current_message() lo devuelve
directamente:
from lapinbeam import MessageMeta, RemoteRef, current_message
async def receive(self, msg):
if msg.get("type") == "TASK":
meta: MessageMeta | None = current_message() # ¿quién envió esto, y a quién responder?
remote: RemoteRef = self.node.get_remote_actor(meta.src, meta.reply_to)
await remote.send({"type": "ACK", "payload_id": msg["payload_id"]})
current_message() devuelve un MessageMeta(src, reply_to, correlation_id,
msg_id, node) — poblado con lo que sea que el emisor haya pasado a
send() — mientras el propio handler que recibió msg sigue en ejecución,
y None fuera de uno (p.ej. desde una tarea en segundo plano que el propio
actor haya lanzado). Para un mensaje enviado por un actor local, src es
el id de este mismo nodo, y msg_id siempre es None (es un id por
conexión que el transporte solo asigna a mensajes remotos). reply_to y
correlation_id son None a menos que el emisor los indique:
await ref.send(msg, reply_to="ingestor", correlation_id=7), tanto en
ActorRef como en RemoteRef.
Como responder a quien envió un mensaje — a reply_to, con el mismo
correlation_id, sea local o remoto — es lo bastante común como para
tener su propio atajo, el fragmento de arriba se puede escribir así:
async def receive(self, msg):
if msg.get("type") == "TASK":
await current_message().reply({"type": "ACK", "payload_id": msg["payload_id"]})
meta.reply(msg) lanza RuntimeError si meta.reply_to es None —
nadie le dio una dirección de respuesta.
Request/response con ask()¶
send() siempre es fire-and-forget — nada conecta una respuesta con el
envío que la provocó a menos que lo construyas tú mismo. ask() hace
justo eso: etiqueta el envío con un correlation_id nuevo, espera una
única respuesta, y funciona igual en ActorRef y en RemoteRef:
respuesta: dict = await remote_processor.ask({"type": "TASK", "payload_id": 1})
El handler que recibe el mensaje sigue teniendo que responder de verdad —
ask() no cambia lo que hace un handler, solo cambia cómo espera quien
llama:
@actor(name="processor")
class Processor:
async def receive(self, msg):
result = msg["payload_id"] * 2
await current_message().reply({"type": "ACK", "result": result})
Si nada responde en timeout segundos (5 por defecto; None espera
indefinidamente), ask() lanza TimeoutError. Por debajo registra un
mailbox oculto de un solo uso como dirección de respuesta y lo limpia
después — no queda ningún actor ni recurso adicional persistente. Ver
Agentes de IA y MCP para un ejemplo trabajado: despachar
tool calls de MCP a un nodo worker, y repartir una pregunta entre varios
actores expertos concurrentemente.
Respuestas en streaming con ask_stream()¶
ask() es para un handler que calcula una única respuesta. Cuando el
handler necesita reportar progreso por el camino — un trabajo largo con
varios pasos, cada uno interesante de mostrar antes del resultado final —
ask_stream() es la misma idea, repetida: el handler llama a
current_message().reply_stream() tantas veces como quiera, y luego a
reply_final() exactamente una vez, y quien pregunta los va leyendo a
medida que llegan:
@actor(name="importer")
class Importer:
async def receive(self, msg):
for row in msg["rows"]:
await do_slow_import(row)
await current_message().reply_stream({"imported": row["id"]})
await current_message().reply_final({"status": "done"})
async def watch_import(ref, rows):
async for update in ref.ask_stream({"rows": rows}, timeout=None):
print(update) # {"imported": ...} unas cuantas veces, luego {"status": "done"}
timeout (5s por defecto, igual que ask()) aplica por elemento aquí,
no al stream entero — el reloj se reinicia tras cada reply_stream()/
reply_final(), así que un handler que sigue trabajando activamente
nunca expira solo porque el trabajo completo tarde mucho; solo expira si
se queda callado durante timeout segundos. Funciona igual en
ActorRef, RemoteRef, y en un PoolRef de
Supervisor.spawn_pool() — el worker que acabe procesando el mensaje es
el que verás respondiendo.
Cuándo usarlo: siempre que quien pregunta de verdad quiera observar
el progreso, no solo esperar una respuesta única — alimentar una barra de
progreso, un log, o (el caso común) una respuesta Server-Sent Events en
un handler web. Cuándo no: si solo te importa el resultado final,
ask() a secas es más simple y no exige que el handler se acuerde de
llamar a reply_final(). Y ask_stream() solo le llega a quien lo
llamó — si varios observadores independientes necesitan las mismas
actualizaciones en vivo (p.ej. más de una pestaña abierta sobre el mismo
trabajo en marcha), repartirlo entre ellos es cosa tuya, no de
ask_stream(): que una sola tarea llame a ask_stream() y reenvíe cada
actualización a un pub/sub local pequeño al que se suscriban los demás
observadores, en vez de que cada uno llame a ask_stream() por su
cuenta. examples/order_stream/ es exactamente esto: una tarea de relay
por pedido, alimentando cuantas conexiones SSE lo estén observando.
Ejecútalos como dos procesos separados (ver Ejemplos para
correr esto entre contenedores u hosts reales en vez de 127.0.0.1):
# terminal 1
NODE_NAME=node_a@127.0.0.1:9001 PEER=node_b@127.0.0.1:9002 python app_node_a.py
# terminal 2
NODE_NAME=node_b@127.0.0.1:9002 PEER=node_a@127.0.0.1:9001 python app_node_b.py
Observar el clúster¶
connect_peer ya espera a que el handshake se complete antes de devolver el
control, pero rara vez querrás volar a ciegas sobre el estado de la conexión
o los errores de entrega en un proceso de larga duración — suscríbete a los
eventos de sistema:
def on_event(event: dict) -> None:
if event["kind"] == "peer_disconnected":
print("peer perdido:", event["peer"])
elif event["kind"] == "error":
print("error de entrega desde", event["peer"], ":", event["detail"],
"correlation_id:", event["correlation_id"])
elif event["kind"] == "decode_error":
print("mensaje inválido para", event["actor"], ":", event["detail"])
elif event["kind"] == "reconnect_gave_up":
print("dejando de reintentar con:", event["peer"])
elif event["kind"] == "supervisor_gave_up":
print("actor detenido definitivamente:", event["actor"], ":", event["detail"])
elif event["kind"] == "mailbox_full":
print("mensaje descartado para:", event["actor"])
node.on_event(on_event)
event["kind"] es uno de "peer_connected", "peer_disconnected",
"error" (un peer reportó un fallo de entrega, p.ej. un mensaje enviado a
un nombre de actor que no existe en el nodo remoto — event["correlation_id"]
repite lo que sea que llevara el send() fallido, o None),
"decode_error" (un mensaje para un actor local no se pudo decodificar —
p.ej. un ValidationError de Pydantic sobre un payload mal formado — y se
descartó antes de llegar a la mailbox de ese actor, en vez de perderse en
una línea de log de asyncio sin relación aparente), "reconnect_gave_up"
(la reconexión automática a event["peer"] se abandonó tras
reconnect_max_attempts intentos fallidos — ya no se reintenta ni se
sigue rastreando, así que no queda ninguna fuga; llama a connect_peer()
de nuevo si quieres reintentarlo), o "supervisor_gave_up" (un
Supervisor dejó de reiniciar event["actor"] tras demasiados fallos
dentro de su ventana de reinicios — incluyendo un fallo en el propio
__init__ del actor, no solo en sus handlers receive/@on — y ya no
sigue en ejecución), o "mailbox_full" (se descartó un mensaje para
event["actor"] porque su mailbox estaba lleno — solo posible si el
Node de ese actor se creó con mailbox_capacity; sin límite por
defecto, ver Limitaciones). Si ya sabes que no
necesitas más un peer, llama a node.forget_peer(peer_id) en vez de
esperar a que esto pase solo.
Ajustar la detección de fallos y el backpressure¶
Node(...) acepta algunos parámetros más además de los ya vistos, todos
opcionales y con valores por defecto que mantienen el comportamiento
actual si no se indican:
node = Node(
"app@127.0.0.1:0",
heartbeat_interval=1.0, # cada cuánto hacer ping a cada peer
peer_timeout=3.0, # abandonar un peer que no ha enviado nada en este tiempo
peer_queue_capacity=256, # tramas salientes en cola por peer
mailbox_capacity=None, # límite del mailbox por actor; None = sin límite
)
heartbeat_interval/peer_timeout controlan cuán rápido se detecta un
peer silenciosamente caído — acortar peer_timeout sin acortar también
heartbeat_interval en ambos lados dará falsos positivos ante el jitter
normal de la red. peer_queue_capacity acota cuántas tramas salientes se
pueden encolar para un peer cuya escritura TCP está congestionada.
mailbox_capacity es el único que cambia de verdad el comportamiento por
defecto una vez que se fija — ver "mailbox_full" arriba.
Siguiente: Mensajes tipados para enviar tipos reales de Python en vez de dicts, Patrones inspirados en OTP para árboles de supervisión, links, monitors, grupos y registro de nombres a nivel de clúster, o Benchmarks para saber cuánto cuesta esto en latencia y throughput.