# How BettaFish Implements Streaming Responses with Flask and SSE

> Learn how BettaFish uses Flask and SSE for real-time streaming responses, pushing progress updates directly to browsers without WebSockets or polling.

- Repository: [BaiFu/bettafish](https://github.com/666ghj/bettafish)
- Tags: how-to-guide
- Published: 2026-02-23

---

**BettaFish uses a thread-safe queue system with Flask's `stream_with_context` to push real-time report generation progress to browsers via Server-Sent Events, eliminating the need for WebSockets or polling.**

The open-source repository **666ghj/bettafish** implements a production-ready streaming API that delivers live updates during long-running report generation tasks. By combining Flask's streaming response capabilities with the SSE protocol, the application pushes stage progress, warnings, and completion events directly to the client as soon as they occur on the backend.

## Architecture Overview

The streaming infrastructure relies on several coordinated components within the `ReportEngine` module:

| Component | Role | Location |
|-----------|------|----------|
| **Flask Blueprint** | Mounts `/api/report/*` endpoints under the main application | [`app.py`](https://github.com/666ghj/bettafish/blob/main/app.py) lines **78‑81** |
| **ReportTask class** | Encapsulates a single job with a bounded event history (`deque`) and thread-locking | [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **11‑15** |
| **publish_event** | Method invoked by the generation pipeline to timestamp and store new events | [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **69‑89** |
| **_broadcast_event** | Multicasts events to all registered subscriber queues | [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **161‑173** |
| **stream_task view** | SSE endpoint that handles client connections, replay, and heartbeats | [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **751‑861** |
| **stream_handler** | Callback bridge connecting the core Report Engine to the SSE layer | `run_report` function, lines **52‑55** |

## End-to-End Event Flow

When a client initiates a report generation request, the system executes the following sequence:

1. **Client initialization** creates a new `EventSource` connection to `/api/report/stream/<task_id>`.
2. **Queue registration** occurs via `_register_stream` (lines **125‑140**), which creates a unique `Queue` object for that specific client connection and stores it in the `stream_subscribers` registry.
3. **Task execution** begins in a background thread via `run_report`, which receives a `stream_handler` callback.
4. **Event generation** happens as the report pipeline calls `task.publish_event(...)` for each stage transition, warning, or progress update.
5. **Broadcast distribution** pushes the event to all listening queues through `_broadcast_event`, using a shallow copy of the subscriber list to prevent race conditions during disconnection.
6. **SSE formatting** converts the JSON event into the wire protocol using `_format_sse` (lines **215‑231**), yielding `id`, `event`, and `data` fields.
7. **Stream termination** occurs when the generator detects a terminal state (`completed`, `error`, or `cancelled`), triggering `_unregister_stream` (lines **43‑58**) to clean up the queue and prevent memory leaks.

## Core Implementation Details

### Event Queue Registration

Each client connection receives an isolated `Queue` object to prevent cross-talk between listeners. The registration logic in [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) uses a module-level lock for thread safety:

```python
def _register_stream(task_id: str) -> Queue:
    """
    Register an event queue for SSE listeners for a specific task.
    """
    queue = Queue()
    with stream_lock:
        stream_subscribers[task_id].append(queue)
    return queue

```

*Source: [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **125‑140**.*

### Publishing and Storing Events

The `ReportTask.publish_event` method assigns monotonically increasing IDs and maintains a bounded history of the last 1000 events for replay capabilities:

```python
def publish_event(self, event_type: str, payload: Dict[str, Any]) -> None:
    timestamp = datetime.utcnow().isoformat() + 'Z'
    event = {
        'id': 0,
        'type': event_type,
        'task_id': self.task_id,
        'timestamp': timestamp,
        'payload': payload,
    }
    with self._event_lock:
        self.last_event_id += 1
        event['id'] = self.last_event_id
        self.event_history.append(event)
    _broadcast_event(self.task_id, event)

```

*Source: [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **69‑89**.*

### Broadcasting to All Subscribers

The `_broadcast_event` function iterates over all queues registered for a specific task ID, using a non-blocking `put` with a short timeout to avoid stalling the generator if a client has disconnected:

```python
def _broadcast_event(task_id: str, event: Dict[str, Any]):
    with stream_lock:
        listeners = list(stream_subscribers.get(task_id, []))
    for queue in listeners:
        try:
            queue.put(event, timeout=0.1)
        except Exception:
            logger.exception("Failed to push stream event, skipping current listener queue")

```

*Source: [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **161‑173**.*

### The SSE Endpoint and Heartbeat Mechanism

The `stream_task` view implements the actual HTTP endpoint that returns `text/event-stream`. It supports reconnection replay via the `Last-Event-ID` header and emits periodic heartbeats to prevent proxy timeouts:

```python
@report_bp.route('/stream/<task_id>', methods=['GET'])
def stream_task(task_id: str):
    task = _get_task(task_id)
    if not task:
        return jsonify({'success': False, 'error': 'Task does not exist'}), 404

    def event_generator():
        queue = _register_stream(task_id)
        # Replay missed events for reconnecting clients

        for event in task.history_since(last_event_id):
            yield _format_sse(event)

        # Main consumption loop

        while True:
            try:
                event = queue.get(timeout=STREAM_HEARTBEAT_INTERVAL)
            except Empty:
                # Emit heartbeat to keep proxies (NGINX, Cloudflare) alive

                event = {
                    'id': f"hb-{int(time.time()*1000)}",
                    'type': 'heartbeat',
                    'task_id': task_id,
                    'timestamp': datetime.utcnow().isoformat() + 'Z',
                    'payload': {'status': task.status}
                }
            yield _format_sse(event)

            if event['type'] in ('completed', 'error', 'cancelled'):
                break
        _unregister_stream(task_id, queue)

    return Response(stream_with_context(event_generator()),
                    mimetype='text/event-stream')

```

*Source: [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **751‑861**.*

### SSE Protocol Formatting

Events conform to the W3C Server-Sent Events specification, with explicit `id`, `event`, and `data` fields separated by newlines:

```python
def _format_sse(event: Dict[str, Any]) -> str:
    payload = json.dumps(event, ensure_ascii=False)
    event_id = event.get('id', 0)
    event_type = event.get('type', 'message')
    return f"id: {event_id}\nevent: {event_type}\ndata: {payload}\n\n"

```

*Source: [`ReportEngine/flask_interface.py`](https://github.com/666ghj/bettafish/blob/main/ReportEngine/flask_interface.py) lines **215‑231**.*

## Practical Implementation Examples

### Server-Side Report Generation

To initiate a streaming report, create a `ReportTask` instance and launch the generation in a background thread:

```python
import uuid
import threading
from flask import jsonify
from ReportEngine.flask_interface import ReportTask, run_report, tasks_registry

def start_report():
    task_id = str(uuid.uuid4())
    task = ReportTask(query=user_query, task_id=task_id)
    
    # Register task in global registry (protected by task_lock)

    with task_lock:
        tasks_registry[task_id] = task
    
    # Launch generation; run_report creates the stream_handler callback

    threading.Thread(target=run_report, args=(task,)).start()
    
    return jsonify({'success': True, 'task_id': task_id})

```

The `run_report` function (lines **52‑55**) creates the `stream_handler` callback that bridges the core engine's output to `task.publish_event()`, converting internal stages into SSE events.

### Client-Side EventSource Consumption

The browser uses the native `EventSource` API to consume the stream without external libraries:

```javascript
const taskId = '123e4567-e89b-12d3-a456-426614174000'; // From /generate response
const source = new EventSource(`/api/report/stream/${taskId}`);

source.addEventListener('stage', e => {
    const data = JSON.parse(e.data);
    console.log('Stage:', data.payload.stage, data.payload.message);
});

source.addEventListener('progress', e => {
    const {progress} = JSON.parse(e.data).payload;
    document.getElementById('progress-bar').style.width = `${progress}%`;
});

source.addEventListener('html_ready', e => {
    const {report_file} = JSON.parse(e.data).payload;
    window.location.href = report_file; // Download generated report
});

source.addEventListener('error', e => {
    console.error('SSE connection error', e);
    source.close();
});

```

The browser automatically handles reconnection and sends the `Last-Event-ID` header, allowing the server to replay missed events from the bounded `deque` history.

## Summary

- **Isolated queues**: Each client receives a dedicated `Queue` object registered via `_register_stream`, ensuring thread-safe event distribution without cross-talk.
- **Bounded history**: The `ReportTask` class maintains a rolling buffer of the last 1000 events, enabling seamless replay for reconnecting clients.
- **Protocol compliance**: The `_format_sse` function generates spec-compliant output with `id`, `event`, and `data` fields, supporting automatic browser reconnection.
- **Heartbeat keep-alive**: The generator emits periodic heartbeat events when the queue is empty, preventing NGINX or Cloudflare from closing idle connections.
- **Automatic cleanup**: The `_unregister_stream` function removes disconnected client queues immediately upon task completion or client disconnect, preventing memory leaks.

## Frequently Asked Questions

### How does BettaFish handle client reconnections without losing events?

The implementation uses the `Last-Event-ID` HTTP header sent automatically by the browser's `EventSource` on reconnection. The `stream_task` view passes this ID to `task.history_since(last_event_id)`, which replays all events from the bounded `deque` that the client missed before entering the live consumption loop.

### Why does BettaFish use Server-Sent Events instead of WebSockets?

SSE provides unidirectional server-to-client streaming over standard HTTP, which simplifies deployment behind corporate proxies and load balancers. Since the report generation only requires server pushes (stages, progress, completion), SSE eliminates the complexity of WebSocket handshake management and fallback handling while supporting automatic reconnection and event IDs natively.

### What prevents memory leaks when clients disconnect unexpectedly?

The `event_generator` function within `stream_task` detects disconnection by monitoring the WSGI input stream via `request.environ.get('wsgi.input')`. When the client closes the connection, this triggers the generator to exit, immediately calling `_unregister_stream` to remove the client's `Queue` from the `stream_subscribers` dictionary.

### How are simultaneous viewers of the same report handled?

The `_broadcast_event` function creates a shallow copy of the subscriber list for a specific `task_id` before iterating, allowing the system to safely multicast events to all connected clients. Each viewer maintains an independent `Queue`, so slow consumers do not block the event publication to other clients.