# Phase 2 — File & DW Adapters + Async Queue

**Status:** ✅ Complete — 2026-04-26
**Duration:** Weeks 6–8
**Goal:** Full async pipeline: SFTP file → Broadway → Data Warehouse insert + DLQ.

---

## Deliverable

1. `adapter_file` polls SFTP, downloads CSV, emits rows to Broadway
2. Broadway processes rows, batches to `adapter_dw` bulk insert
3. Failed rows go to DLQ table; DLQ depth emits to telemetry
4. Async jobs return `202 Accepted` with `job_id`; completion event via PubSub
5. `POST /api/v1/files/upload` triggers file processing pipeline

---

## Tasks

### 1. infra_queue — Broadway Pipelines

**Dependencies to add:** `broadway`, `gen_stage`

```
DB migrations:
  - create table dead_letter_queue (
      id, pipeline, message_payload, error_reason,
      retry_count, source_file, source_row, inserted_at
    )
  - create table async_jobs (
      id, type, status, started_at, completed_at,
      result_summary, error
    )
```

```elixir
defmodule InfraQueue.FilePipeline do
  use Broadway

  def start_link(opts) do
    Broadway.start_link(__MODULE__,
      name: opts[:name] || __MODULE__,
      producer: [
        module: {opts[:producer_module], opts[:producer_opts]},
        concurrency: 1
      ],
      processors: [default: [concurrency: 10]],
      batchers: [
        dw: [batch_size: 500, batch_timeout: 5_000, concurrency: 2]
      ]
    )
  end

  @impl true
  def handle_message(:default, message, _context) do
    case MwTransform.Mapper.row_to_canonical(message.data) do
      {:ok, canonical} -> Broadway.Message.put_batcher(message, :dw) |> Broadway.Message.update_data(fn _ -> canonical end)
      {:error, reason} -> Broadway.Message.failed(message, reason)
    end
  end

  @impl true
  def handle_batch(:dw, messages, _batch_info, _context) do
    payloads = Enum.map(messages, & &1.data)
    case AdapterDw.BatchLoader.bulk_insert(payloads) do
      :ok -> messages
      {:error, reason} -> Enum.map(messages, &Broadway.Message.failed(&1, reason))
    end
  end

  @impl true
  def handle_failed(messages, _context) do
    Enum.each(messages, fn msg ->
      InfraQueue.DeadLetterStore.insert(%{
        pipeline: "file_pipeline",
        message_payload: msg.data,
        error_reason: msg.status |> elem(1) |> inspect()
      })
    end)
    messages
  end
end
```

### 2. adapter_file — SFTP + Parsers

**Dependencies to add:** `nimble_csv`, `sweet_xml`
(SFTP via Erlang `:ssh` — built-in OTP)

```elixir
defmodule AdapterFile.FileWatcher do
  use GenServer

  @poll_interval 60_000

  def init(config) do
    schedule_poll()
    {:ok, %{config: config, seen_files: MapSet.new()}}
  end

  def handle_info(:poll, state) do
    new_files = AdapterFile.SFTPClient.list_new(state.config, state.seen_files)
    Enum.each(new_files, fn file ->
      {:ok, contents} = AdapterFile.SFTPClient.download(state.config, file)
      rows = AdapterFile.CsvParser.parse(contents)
      Enum.each(rows, fn row ->
        InfraQueue.FileProducer.enqueue(%{source_file: file, row: row})
      end)
    end)
    schedule_poll()
    {:noreply, %{state | seen_files: MapSet.union(state.seen_files, MapSet.new(new_files))}}
  end

  defp schedule_poll, do: Process.send_after(self(), :poll, @poll_interval)
end
```

### 3. adapter_dw — Batch Loader

```elixir
defmodule AdapterDw.BatchLoader do
  def bulk_insert(messages) do
    payload = Enum.map(messages, &to_dw_record/1)
    case AdapterDw.Client.post("/bulk/transactions", payload) do
      {:ok, %{status: 201}} -> :ok
      {:error, reason} -> {:error, reason}
    end
  end
end
```

### 4. adapter_http — Generic REST Adapter

Configurable per endpoint in `config/runtime.exs`:

```elixir
config :adapter_http, :endpoints, %{
  "payment_validation" => %{
    url: "https://validation.internal",
    auth: {:bearer, System.get_env("VALIDATION_API_KEY")},
    timeout: 3_000,
    retries: 3
  }
}
```

### 5. Async Job API

`gateway_api` exposes:

```
POST /api/v1/files/upload
  → creates async_job record
  → returns 202 {"job_id": "...", "status": "processing"}

GET /api/v1/jobs/:id
  → returns job status {"status": "completed", "rows_processed": 1200}
```

PubSub completion event:
```elixir
Phoenix.PubSub.broadcast(MwCore.PubSub, "jobs:#{job_id}", %{status: :completed, rows: 1200})
```

---

## Acceptance Criteria

- [ ] FileWatcher detects new SFTP file, downloads, emits to Broadway
- [ ] Broadway processes 10k rows in under 30 seconds (batch size 500, concurrency 10)
- [ ] Malformed rows end up in DLQ table, not silently dropped
- [ ] DLQ depth telemetry counter increments per failed batch
- [ ] `GET /api/v1/jobs/:id` returns `completed` after Broadway finishes
- [ ] `adapter_http` calls internal validation API with configured auth + retry
- [ ] Broadway pipeline drains cleanly on `mix release stop`
