AsyncStream: convertir callbacks en secuencias
El adaptador universal entre el mundo de las notificaciones y el de la iteración asíncrona. Anatomía de la continuación, semántica de `yield` y su resultado, el contrato de terminación con `finish`, la variante `AsyncThrowingStream` y el patrón `makeStream` que separa producción de consumo.
Casi ninguna fuente real de eventos nace hablando el idioma de AsyncSequence. Un delegado llama a un método, un observador recibe una notificación, una biblioteca en C invoca un puntero a función, un temporizador dispara un bloque. Todas esas APIs son de empuje: deciden cuándo hablar y no aceptan que nadie les pida el siguiente valor. AsyncStream es la pieza que traduce ese mundo al de tracción sin obligar a reescribir la fuente: crea una secuencia asíncrona vacía y te entrega un mando a distancia —la continuación— que puedes guardar donde haga falta y usar desde cualquier hilo para empujar elementos dentro. Entender bien ese objeto, su semántica de encolado y su contrato de terminación es la diferencia entre un adaptador correcto y una fuga de memoria que nadie encuentra.
- Construir una
AsyncStreama partir de una API basada en callbacks o delegados. - Explicar el ciclo de vida de la continuación y por qué debe conservarse fuera del inicializador.
- Interpretar el resultado de
yieldy aplicar correctamentefinishyfinish(throwing:). - Elegir entre el inicializador con clausura, la variante por desdoblamiento y el patrón
makeStream.
El adaptador y su mando a distancia
La forma canónica envuelve una API de callbacks en el inicializador que recibe una clausura de construcción.
func posiciones(del sensor: Sensor) -> AsyncStream<Posicion> {
AsyncStream { continuation in
sensor.alRecibir = { punto in
continuation.yield(punto)
}
sensor.alTerminar = {
continuation.finish()
}
sensor.empezar()
}
}
La clausura se ejecuta una sola vez, cuando alguien crea el iterador, y su único cometido es entregar la continuación a quien vaya a producir. La continuación es Sendable: puede viajar a otro hilo, guardarse en una propiedad, capturarse en un bloque de C. yield no suspende nunca y puede llamarse desde cualquier contexto, síncrono o asíncrono, aislado o no; esa es exactamente la propiedad que hace posible el puente, porque el código antiguo que la llama no sabe nada de tareas ni de actores.
En el otro extremo, el consumidor no percibe nada de esta maquinaria:
for await p in posiciones(del: sensor) {
dibujar(p)
}
Hay una regla que separa el uso correcto del incorrecto y conviene enunciarla pronto: si la continuación no sobrevive a la clausura de construcción, el flujo está roto. Registrarla en un objeto, en una propiedad o en un manejador es lo que la mantiene viva; dejarla morir al terminar la clausura produce un flujo que no emite jamás y que además nunca termina, porque nadie podrá llamar a finish.
Un flujo cuya continuación se pierde no falla, no avisa y no se cierra: el bucle consumidor simplemente se queda suspendido para siempre. El síntoma llega tarde, disfrazado de tarea colgada o de pantalla que no carga. Cuando un for await no produce nada, la primera hipótesis debe ser siempre que la continuación se soltó demasiado pronto.
Semántica de yield y contrato de terminación
yield devuelve un valor que casi todo el mundo descarta y que contiene información importante.
func yield(_ value: Element) -> YieldResult // el resultado es descartable
// .enqueued(remaining: Int) el elemento entro en el bufer
// .dropped(Element) la politica de bufer lo descarto
// .terminated el flujo ya habia terminado
Ese resultado es el único canal por el que el productor se entera de lo que ocurre aguas abajo: si el búfer está lleno y la política descarta, si el consumidor ya se marchó. Ignorarlo es aceptable en un adaptador simple; en un productor caro —uno que lee un fichero o descodifica imágenes— es la señal que permite dejar de trabajar en vano.
La terminación tiene tres puertas, y las tres conducen al mismo estado final:
continuation.finish() // fin normal
continuation.finish(throwing: ErrorRed.caida) // solo en AsyncThrowingStream
// y la tercera: que la continuacion se destruya sin que nadie llame a finish
finish es idempotente: la segunda llamada no hace nada, igual que un yield posterior, que devolverá .terminated. Los elementos ya encolados se entregan antes de que el bucle vea el final, de modo que terminar no equivale a truncar. Y la tercera puerta importa más de lo que parece: si la continuación se libera sin haber llamado a finish, el flujo se cierra igualmente. Es una red de seguridad contra tareas eternamente colgadas, no una invitación a olvidarse de cerrar.
La variante que lanza cambia solo la firma del error:
func lineas(de socket: Socket) -> AsyncThrowingStream<String, Error> {
AsyncThrowingStream { continuation in
socket.alRecibir = { continuation.yield($0) }
socket.alFallar = { continuation.finish(throwing: $0) }
socket.alCerrar = { continuation.finish() }
}
}
El consumidor pasa entonces a escribir for try await, y el error interrumpe el bucle como cualquier otro lanzamiento. Un detalle que se olvida: tras un error el flujo está terminado, no pausado. No hay reanudación posible ni reintento implícito; reintentar es responsabilidad de quien consume, envolviendo el bucle.
stateDiagram-v2 [*] --> Inactivo Inactivo --> Activo: se crea el iterador y corre la clausura Activo --> Activo: yield encola un elemento Activo --> Terminado: finish normal Activo --> Terminado: finish con error Activo --> Terminado: la continuacion se libera Terminado --> [*]: onTermination y limpieza note right of Terminado yield posterior devuelve terminated finish posterior no hace nada end note
Tres formas de construir, tres situaciones
El inicializador con clausura no es la única vía, y elegir mal complica código que podría ser trivial.
El desdoblamiento sirve cuando la fuente ya es de tracción y solo hay que repetir una operación asíncrona hasta agotarla. No hay continuación en juego.
let paginas = AsyncStream(unfolding: {
await api.siguientePagina() // devolver nil termina el flujo
})
El patrón makeStream, incorporado en Swift 5.9, devuelve la pareja de golpe y resuelve de raíz el problema de la continuación fugitiva, además de permitir que productor y consumidor se creen en momentos distintos.
let (flujo, cont) = AsyncStream.makeStream(of: Posicion.self, bufferingPolicy: .bufferingNewest(32))
final class Rastreador {
let continuacion: AsyncStream<Posicion>.Continuation
init(_ c: AsyncStream<Posicion>.Continuation) { continuacion = c }
func recibir(_ p: Posicion) { continuacion.yield(p) }
deinit { continuacion.finish() }
}
Ese deinit es el idioma más limpio para atar la vida del flujo a la vida del productor. Y el inicializador con clausura sigue siendo la mejor opción cuando la suscripción a la fuente debe ocurrir justo al empezar a consumir y cancelarse al dejar de hacerlo, porque su clausura se ejecuta de forma perezosa y onTermination le da un lugar natural donde deshacer el registro.
Adaptador, no motor
Un flujo no produce nada por sí mismo. Es una tubería con un extremo de empuje y otro de tracción; el trabajo real sigue estando en la fuente que envuelve.
yield no espera
La llamada es no bloqueante y siempre retorna al instante. Lo que ocurra con el elemento lo decide la política de búfer, no el productor.
Cerrar siempre
Todo camino de salida de la fuente debe acabar en finish. Un flujo sin cierre es una tarea consumidora suspendida de forma indefinida.
Los tropiezos habituales
Cuatro fallos concentran la mayoría de los problemas reales.
Un solo consumidor. AsyncStream está diseñada para un único bucle. Dos consumidores sobre el mismo flujo no reciben copias: se reparten los elementos de forma impredecible. Si necesitas difusión, construye un tipo que la implemente por encima.
Consumir dos veces. Iterar de nuevo un flujo ya consumido no lo reinicia. Una AsyncStream es un objeto de un solo uso; la receta reutilizable es la función que la fabrica.
Producir sin límite. Si el productor emite más rápido de lo que el consumidor procesa, la política de búfer por defecto —ilimitada— hace crecer la memoria hasta donde llegue. Es el tema entero de la lección siguiente y el motivo de que makeStream pida la política de forma explícita.
Confundir el aislamiento. Que yield sea seguro desde cualquier hilo no convierte en seguro lo que hagas alrededor. Mutar estado compartido dentro del callback de la fuente sigue siendo tu problema, y el compilador en modo estricto lo señalará como tal.
Debajo de este adaptador hay una idea que atraviesa la concurrencia moderna entera: la reificación de la continuación, es decir, convertir en un objeto que puedes guardar y pasar aquello que normalmente es invisible, el resto del programa que queda por ejecutar. Cuando llamas a withCheckedContinuation para adaptar un callback a un await, el compilador toma el punto exacto en que la función se suspendió y te lo entrega como un valor con un método resume: has cogido un instante del flujo de control y lo has metido en una variable. AsyncStream.Continuation es la misma jugada llevada de un solo valor a una serie indefinida de ellos, y por eso su parecido no es casual sino estructural; la diferencia es que aquella se reanuda exactamente una vez —reanudarla dos veces es un fallo grave— mientras que esta admite tantos yield como haga falta antes del cierre. De ahí se explica su comportamiento entero sin necesidad de memorizar reglas: yield no puede suspender porque quien lo llama es código ajeno que quizá no vive dentro de ninguna tarea, así que el desajuste de ritmos tiene que absorberse en un búfer y no en una espera del productor; la continuación es Sendable porque su razón de existir es cruzar la frontera entre un mundo sin concurrencia estructurada y otro que sí la tiene; y finish debe ser idempotente porque las APIs de empuje son notoriamente descuidadas al notificar su final, a veces dos veces, a veces ninguna. Merece la pena ver también qué se pierde en la traducción: la fuente original no tiene forma de saber si alguien la escucha ni de frenar cuando el consumidor se retrasa, y ninguna cantidad de ingeniería en el adaptador puede inventar una contrapresión que la fuente no ofrece. El flujo puede descartar, puede acumular o puede avisar mediante el resultado de yield, pero no puede pedirle a un sensor que emita más despacio. Esa asimetría —tracción aguas abajo, empuje aguas arriba, y un búfer en medio como única moneda de cambio— es el precio exacto de tener un adaptador universal, y explica por qué la decisión importante al escribir un AsyncStream no es cómo se construye, sino qué se hace cuando los dos ritmos dejan de cuadrar.
- Adapta una API de delegado real a
AsyncStreamy comprueba qué ocurre si no guardas la continuación fuera de la clausura. - Registra el resultado de cada
yieldy provoca a propósito un.droppedcon un consumidor lento y una política acotada. - Escribe la misma fuente con
AsyncThrowingStreamy verifica que trasfinish(throwing:)unyieldposterior devuelve.terminated. - Ata la continuación al
deinitde un objeto productor y observa el cierre del bucle al liberarlo. - Consume el mismo flujo desde dos bucles concurrentes y documenta cómo se reparten los elementos.