# Implementation Plan: Streaming híbrido en chat de conversación

**Branch**: `003-chat-streaming` | **Date**: 2026-05-24 | **Spec**: [spec.md](./spec.md)

**Input**: Feature specification from `/specs/003-chat-streaming/spec.md`

## Summary

Añadir **streaming híbrido** al chat de conversación (`/assistant`): eventos de **estado** durante clasificación, consulta de datos y preparación de respuesta, seguidos de **fragmentos de texto** del asistente en tiempo real. Solo mensajes de **texto** usan el nuevo canal; adjuntos PDF/JPG mantienen `POST /api/v1/chat/message` síncrono. Backend expone SSE (`text/event-stream`); frontend Angular consume con `fetch` + `ReadableStream` (JWT en header). Se extiende `LlmChatCompletionContract` con streaming; el orquestador emite estados vía `AgentStreamEmitter`.

## Technical Context

**Language/Version**: PHP 8.2+, Laravel 12; TypeScript / Angular 19  
**Primary Dependencies**: `LlmChatCompletionContract`, `AgentOrchestratorService`, `ChatConversationService`, Laravel `StreamedResponse`, Angular `fetch` API  
**Storage**: Tabla existente `chat_messages` (sin migración v1)  
**Testing**: PHPUnit (`ChatConversationStreamServiceTest`, `AgentOrchestratorStreamTest`); Karma/Jasmine para parser SSE en `ChatApiService`  
**Target Platform**: API Laravel + SPA `ChatFront/` (`conversation` feature)  
**Project Type**: Extensión full-stack brownfield (backend + frontend)  
**Performance Goals**: Primer evento `status` visible en UI < 2 s (SC-001); time-to-first-token mejorado vs. baseline síncrono  
**Constraints**: JWT Bearer obligatorio; no persistir mensajes assistant incompletos en error; throttle existente; nginx/PHP buffer deshabilitado para SSE  
**Scale/Scope**: 1 endpoint SSE nuevo, 1 contrato LLM ampliado, refactor orquestador + 2 handlers principales (`GeneralChat`, `AgentDatabase`), UI `conversation.component`  

## Constitution Check

*GATE: Must pass before Phase 0 research. Re-check after Phase 1 design.*

| Principle | Status | Notes |
|-----------|--------|-------|
| **I. Spec-First** | ✅ | Trazado a FR-001–FR-010, US1–US3 |
| **II. Skinny Controllers** | ✅ | `ChatController::messageStream` delega en `ChatConversationStreamService` |
| **III. Contratos IA** | ✅ | `LlmChatCompletionContract::chatStream()`; OpenAI/Gemini adaptadores |
| **IV. API v1** | ✅ | `POST /api/v1/chat/message/stream`; errores vía evento `error` + códigos runtime |
| **V. Contrato API** | ✅ | `contracts/chat-stream-sse.md`; tipos en `api.types.ts` |
| **VI. Frontend desacoplado** | ✅ | Lógica SSE en `ChatApiService`; componente solo UI/estado |
| **VII. Fidelidad PDF** | ✅ N/A | Adjuntos fuera del stream; flujo PDF sin cambios |
| **VIII. Simplicidad** | ✅ | SSE unidireccional; sin WebSockets; endpoint síncrono conservado |

**Post-design re-check**: ✅ Todos los gates pasan. Sin entradas en Complexity Tracking.

## Project Structure

### Documentation (this feature)

```text
specs/003-chat-streaming/
├── spec.md
├── plan.md                          # Este archivo
├── research.md                      # Phase 0
├── data-model.md                    # Phase 1
├── quickstart.md                    # Phase 1
├── contracts/
│   └── chat-stream-sse.md           # Contrato SSE Backend ↔ Frontend
├── checklists/
│   └── requirements.md
└── tasks.md                         # (/speckit-tasks — pendiente)
```

### Source Code (cambios previstos)

```text
app/
├── Contracts/
│   ├── LlmChatCompletionContract.php       # + chatStream()
│   └── Agent/
│       └── AgentStreamEmitterInterface.php # NUEVO — status + chunk + done + error
├── DTOs/
│   └── ChatStreamEvent.php                 # NUEVO — serialización SSE
├── Http/Controllers/Api/V1/
│   └── ChatController.php                  # + messageStream()
├── Services/
│   ├── ChatConversationService.php         # sendMessage sin cambio (adjuntos)
│   ├── ChatConversationStreamService.php   # NUEVO — orquesta stream + persistencia
│   ├── Agent/
│   │   ├── AgentOrchestratorService.php    # + handleStreaming()
│   │   └── Handlers/
│   │       ├── GeneralChatCapabilityHandler.php      # stream final LLM
│   │       └── AgentDatabaseQueryCapabilityHandler.php # status + stream interpret
│   └── Llm/
│       ├── OpenAiLlmChatClient.php         # stream OpenAI
│       └── GeminiLlmChatClient.php         # stream Gemini
config/
└── agent.php                               # + streaming.enabled, status messages (opcional)
routes/
└── api.php                                 # POST chat/message/stream

ChatFront/src/app/
├── core/
│   ├── models/api.types.ts                 # StreamEvent types
│   └── services/chat-api.service.ts        # streamConversationMessage()
└── features/conversation/
    ├── conversation.component.ts           # consume stream, bubble dinámico
    └── conversation.component.html         # quitar spinner estático; bubble streaming

tests/
├── Feature/Chat/
│   └── ChatConversationStreamTest.php
└── Unit/Services/
    └── ChatStreamEventTest.php
```

**Structure Decision**: Feature full-stack. Contrato HTTP documentado en `contracts/`. Endpoint síncrono `POST /chat/message` se mantiene para adjuntos y fallback.

## Architecture

### Flujo objetivo (texto)

```mermaid
sequenceDiagram
    participant UI as ConversationComponent
    participant API as ChatApiService
    participant C as ChatController
    participant S as ChatConversationStreamService
    participant O as AgentOrchestrator
    participant L as LLM

    UI->>API: fetch POST message/stream + JWT
    API->>C: StreamedResponse SSE
    C->>S: streamMessage(user, text)
    S->>S: persist user ChatMessage
    S-->>UI: event status analyzing
    S->>O: handleStreaming(emitter)
    O->>O: classify intent
    O-->>UI: event status querying
    O->>L: chatStream (interpret / reply)
    loop tokens
        L-->>O: delta
        O-->>UI: event chunk
    end
    S->>S: persist assistant ChatMessage
    S-->>UI: event done assistant_id
```

### Fases de estado (mensajes UI)

| Fase | Evento `status` | Cuándo |
|------|-----------------|--------|
| `analyzing` | Analizando tu pregunta… | Tras persistir user; inicio classify/handler |
| `querying` | Consultando datos… | Handler `agent_database` ejecutando herramienta |
| `generating` | Preparando respuesta… | Antes de `chatStream` / redacción final |
| `chunk` | — | Fragmentos de texto assistant |
| `done` | — | `{ assistant_id, content }` |
| `error` | — | Mensaje amigable ES |

## Backend

### Endpoint

- **Nuevo**: `POST /api/v1/chat/message/stream`
- **Auth**: `auth:api` + `throttle:30,1` (igual que message)
- **Body**: `{ "message": "texto" }` (JSON) — sin adjunto
- **Response**: `Content-Type: text/event-stream`; ver [contracts/chat-stream-sse.md](./contracts/chat-stream-sse.md)

### Servicios

1. **`ChatConversationStreamService`**
   - Valida texto no vacío
   - Persiste mensaje **user** al inicio (FR-004, SC-005: no duplicar user en retry si se implementa idempotency key en v2)
   - Invoca `AgentOrchestratorService::handleStreaming($sessionId, $text, $referer, $emitter)`
   - Acumula chunks; en `done` persiste **assistant** completo
   - En `error`: no persiste assistant parcial (edge case spec)

2. **`AgentOrchestratorService::handleStreaming`**
   - Emite `status: analyzing` inmediato
   - Clasifica intención (reutiliza `IntentClassifierService`)
   - Delega handler con `$emitter`
   - Handlers emiten `querying` / `generating` según fase

3. **`LlmChatCompletionContract`**
   - Añadir: `chatStream(array $messages, callable $onChunk, array $options = []): string`
   - Retorna texto completo al final (para persistencia)

4. **Handlers**
   - `GeneralChatCapabilityHandler`: stream directo del LLM
   - `AgentDatabaseQueryCapabilityHandler`: status durante SQL; stream en `interpretResultsAndRespond` (turno 2)
   - Respuestas determinísticas (conteos, rankings formateados en PHP en el futuro): emitir `generating` + un `chunk` con contenido completo o chunks simulados por párrafo

### Controller (skinny)

```php
public function messageStream(ChatConversationRequest $request): StreamedResponse
{
    return $this->chatConversationStreamService->streamResponse(
        $this->resolveUser($request),
        trim((string) $request->validated('message')),
        (string) $request->headers->get('referer', ''),
    );
}
```

### Config (opcional v1)

```php
// config/agent.php
'streaming' => [
    'enabled' => env('AGENT_STREAMING_ENABLED', true),
    'status_messages' => [
        'analyzing' => 'Analizando tu pregunta…',
        'querying' => 'Consultando datos…',
        'generating' => 'Preparando respuesta…',
    ],
],
```

## Frontend

### `ChatApiService`

- Nuevo método `streamConversationMessage(message: string, handlers: StreamHandlers): AbortController`
- Usa `fetch(url, { method: 'POST', headers: { Authorization, Accept: 'text/event-stream', Content-Type: 'application/json' }, body })`
- Parser línea a línea `data: {...}\n\n`
- **No** usar `HttpClient` (no soporta SSE POST nativo)
- **No** usar `EventSource` (no soporta POST ni Bearer custom sin query hack)

### `ConversationComponent`

- `send()`: si hay adjunto → flujo actual `sendConversationMessage`; si solo texto → `streamConversationMessage`
- Añadir mensaje user optimista + bubble assistant vacío con `streamingContent` signal
- Actualizar bubble en cada `status` (subtítulo) y `chunk` (append)
- Al `done`: fijar `id` del assistant, limpiar estado streaming, opcional `loadConversation({ silent: true })` para sincronizar IDs
- Eliminar bloque `@if (sending())` spinner estático; reemplazar por bubble assistant en modo stream
- v1: texto plano durante stream; al `done` opcional render Markdown (ngx-markdown o pipe simple) — documentado en research R7

### Tipos (`api.types.ts`)

```typescript
export type ChatStreamEventType = 'status' | 'chunk' | 'done' | 'error';

export interface ChatStreamEvent {
  type: ChatStreamEventType;
  phase?: 'analyzing' | 'querying' | 'generating';
  message?: string;
  content?: string;
  assistant_id?: number;
}
```

## Testing Strategy

| Layer | Test |
|-------|------|
| Backend | Feature test: POST stream devuelve `Content-Type: text/event-stream`, eventos ordenados, assistant persistido |
| Backend | Unit: `ChatStreamEvent` serialización SSE |
| Backend | Mock LLM: `chatStream` invoca callback con deltas |
| Frontend | Unit: parser SSE en `ChatApiService` |
| Manual | UAT SC-001/SC-002 según [quickstart.md](./quickstart.md) |

## Risks & Mitigations

| Risk | Mitigation |
|------|------------|
| Buffering nginx/PHP impide SSE | `X-Accel-Buffering: no`, `ob_end_flush`, documentar en quickstart |
| 3 LLM calls en agent_database; solo último streamea | Eventos `status` cubren latencia previa (SC-002) |
| Timeout proxy 60s | Alinear `OPENAI_REQUEST_TIMEOUT`; status keep-alive |
| JWT expira mid-stream | Evento `error` + interceptor 401 existente |

## Phase 0 / Phase 1 Artifacts

- [research.md](./research.md) — decisiones SSE, fetch, persistencia, Markdown v1
- [data-model.md](./data-model.md) — StreamEvent, flujo persistencia
- [contracts/chat-stream-sse.md](./contracts/chat-stream-sse.md) — contrato HTTP
- [quickstart.md](./quickstart.md) — prueba manual y env

## Next Step

Ejecutar **`/speckit-tasks`** para generar `tasks.md` e implementación.
