one move: make an operational decision
Simulate a bounded response stream: deliver chunks to a slow consumer, record drops, and stop cleanly on cancellation. In the Async Ingestion Service system, implement this level as a deterministic contract before composing it with the next service.
def stream_with_backpressure(chunks, capacity, consume_per_tick, cancel_after=None):
return {'delivered': [...], 'dropped': [...], 'peak_buffer': ..., 'cancelled': ...}