# Database Schema

STA Connector uses PostgreSQL with Ecto for all deployments, providing robust JSONB support for flexible metadata storage and strong data integrity through Ecto changesets.

## Schema Overview

```
+-----------------------------------------------------------------------+
|                         Database Schema                                |
+-----------------------------------------------------------------------+
|                                                                        |
|  +-------------------------+      +-------------------------+          |
|  |  files                  |      |  file_events            |          |
|  +-------------------------+      +-------------------------+          |
|  | id (PK, UUID)           |<-----| id (PK, UUID)           |          |
|  | protocol                |      | file_id (FK)            |          |
|  | system                  |      | event_type              |          |
|  | direction               |      | details (JSONB)         |          |
|  | status                  |      | inserted_at             |          |
|  | metadata (JSONB)        |      +-------------------------+          |
|  | inserted_at             |                                           |
|  | updated_at              |      +-------------------------+          |
|  +-------------------------+      |  audit_logs             |          |
|                                   +-------------------------+          |
|  +-------------------------+      | id (PK, UUID)           |          |
|  |  configs                |      | action                  |          |
|  +-------------------------+      | actor                   |          |
|  | id (PK, UUID)           |      | details (JSONB)         |          |
|  | key (unique)            |      | inserted_at             |          |
|  | value (JSONB)           |      +-------------------------+          |
|  | updated_at              |                                           |
|  +-------------------------+                                           |
|                                                                        |
+-----------------------------------------------------------------------+
```

## Ecto Schemas

### File Schema

Tracks all files processed through the STA connector (both inbound and outbound).

```elixir
# lib/sta_connector/files/file.ex
defmodule StaConnector.Files.File do
  use Ecto.Schema
  import Ecto.Changeset

  @primary_key {:id, :binary_id, autogenerate: true}
  @foreign_key_type :binary_id

  @directions ~w(inbound outbound)
  @statuses ~w(pending uploading uploaded processing completed failed retrying)
  @systems ~w(CCS CIR CMP STR SPI CAM LDL DICT)

  schema "files" do
    field :protocol, :string
    field :system, :string
    field :direction, :string
    field :status, :string, default: "pending"
    field :metadata, :map, default: %{}

    has_many :events, StaConnector.Files.FileEvent

    timestamps(type: :utc_datetime_usec)
  end

  @required_fields ~w(protocol system direction)a
  @optional_fields ~w(status metadata)a

  @doc false
  def changeset(file, attrs) do
    file
    |> cast(attrs, @required_fields ++ @optional_fields)
    |> validate_required(@required_fields)
    |> validate_inclusion(:direction, @directions)
    |> validate_inclusion(:system, @systems)
    |> validate_inclusion(:status, @statuses)
    |> validate_protocol_format()
    |> unique_constraint(:protocol)
  end

  @doc """
  Changeset for status transitions with validation.
  """
  def status_changeset(file, new_status) do
    file
    |> cast(%{status: new_status}, [:status])
    |> validate_inclusion(:status, @statuses)
    |> validate_status_transition()
  end

  defp validate_protocol_format(changeset) do
    validate_format(changeset, :protocol, ~r/^[A-Z0-9]{10,20}$/,
      message: "must be 10-20 alphanumeric characters"
    )
  end

  defp validate_status_transition(changeset) do
    case {get_field(changeset, :status), get_change(changeset, :status)} do
      {_, nil} -> changeset
      {"completed", _} -> add_error(changeset, :status, "cannot change status of completed file")
      {"failed", status} when status not in ["retrying", "pending"] ->
        add_error(changeset, :status, "failed files can only be retried")
      _ -> changeset
    end
  end
end
```

### FileEvent Schema

Tracks all events/state changes for a file.

```elixir
# lib/sta_connector/files/file_event.ex
defmodule StaConnector.Files.FileEvent do
  use Ecto.Schema
  import Ecto.Changeset

  @primary_key {:id, :binary_id, autogenerate: true}
  @foreign_key_type :binary_id

  @event_types ~w(
    created status_changed uploaded downloaded
    processing_started processing_completed processing_failed
    retry_scheduled callback_sent error
  )

  schema "file_events" do
    field :event_type, :string
    field :details, :map, default: %{}

    belongs_to :file, StaConnector.Files.File

    timestamps(type: :utc_datetime_usec, updated_at: false)
  end

  @required_fields ~w(event_type file_id)a
  @optional_fields ~w(details)a

  @doc false
  def changeset(event, attrs) do
    event
    |> cast(attrs, @required_fields ++ @optional_fields)
    |> validate_required(@required_fields)
    |> validate_inclusion(:event_type, @event_types)
    |> foreign_key_constraint(:file_id)
  end
end
```

### Config Schema

Stores runtime configuration with JSONB values.

```elixir
# lib/sta_connector/settings/config.ex
defmodule StaConnector.Settings.Config do
  use Ecto.Schema
  import Ecto.Changeset

  @primary_key {:id, :binary_id, autogenerate: true}

  schema "configs" do
    field :key, :string
    field :value, :map

    timestamps(type: :utc_datetime_usec, inserted_at: false)
  end

  @required_fields ~w(key value)a

  @doc false
  def changeset(config, attrs) do
    config
    |> cast(attrs, @required_fields)
    |> validate_required(@required_fields)
    |> validate_key_format()
    |> unique_constraint(:key)
  end

  defp validate_key_format(changeset) do
    validate_format(changeset, :key, ~r/^[a-z][a-z0-9_.]+$/,
      message: "must start with lowercase letter and contain only lowercase letters, numbers, dots, and underscores"
    )
  end
end
```

### AuditLog Schema

Tracks all administrative actions for compliance and debugging.

```elixir
# lib/sta_connector/audit/audit_log.ex
defmodule StaConnector.Audit.AuditLog do
  use Ecto.Schema
  import Ecto.Changeset

  @primary_key {:id, :binary_id, autogenerate: true}

  @actions ~w(
    config_created config_updated config_deleted
    file_retry_requested file_status_changed
    route_enabled route_disabled
    system_started system_stopped
    user_login user_logout
  )

  schema "audit_logs" do
    field :action, :string
    field :actor, :string
    field :details, :map, default: %{}

    timestamps(type: :utc_datetime_usec, updated_at: false)
  end

  @required_fields ~w(action actor)a
  @optional_fields ~w(details)a

  @doc false
  def changeset(log, attrs) do
    log
    |> cast(attrs, @required_fields ++ @optional_fields)
    |> validate_required(@required_fields)
    |> validate_inclusion(:action, @actions)
  end
end
```

## Ecto Migrations

### Create Files Table

```elixir
# priv/repo/migrations/20260131000001_create_files.exs
defmodule StaConnector.Repo.Migrations.CreateFiles do
  use Ecto.Migration

  def change do
    create table(:files, primary_key: false) do
      add :id, :binary_id, primary_key: true
      add :protocol, :string, null: false
      add :system, :string, null: false
      add :direction, :string, null: false
      add :status, :string, null: false, default: "pending"
      add :metadata, :jsonb, null: false, default: "{}"

      timestamps(type: :utc_datetime_usec)
    end

    create unique_index(:files, [:protocol])
    create index(:files, [:system])
    create index(:files, [:direction])
    create index(:files, [:status])
    create index(:files, [:inserted_at])
    create index(:files, [:system, :direction, :status])

    # GIN index for JSONB queries
    execute(
      "CREATE INDEX files_metadata_gin ON files USING GIN (metadata jsonb_path_ops)",
      "DROP INDEX files_metadata_gin"
    )
  end
end
```

### Create File Events Table

```elixir
# priv/repo/migrations/20260131000002_create_file_events.exs
defmodule StaConnector.Repo.Migrations.CreateFileEvents do
  use Ecto.Migration

  def change do
    create table(:file_events, primary_key: false) do
      add :id, :binary_id, primary_key: true
      add :file_id, references(:files, type: :binary_id, on_delete: :delete_all), null: false
      add :event_type, :string, null: false
      add :details, :jsonb, null: false, default: "{}"

      timestamps(type: :utc_datetime_usec, updated_at: false)
    end

    create index(:file_events, [:file_id])
    create index(:file_events, [:event_type])
    create index(:file_events, [:inserted_at])
    create index(:file_events, [:file_id, :event_type])

    # GIN index for JSONB queries
    execute(
      "CREATE INDEX file_events_details_gin ON file_events USING GIN (details jsonb_path_ops)",
      "DROP INDEX file_events_details_gin"
    )
  end
end
```

### Create Configs Table

```elixir
# priv/repo/migrations/20260131000003_create_configs.exs
defmodule StaConnector.Repo.Migrations.CreateConfigs do
  use Ecto.Migration

  def change do
    create table(:configs, primary_key: false) do
      add :id, :binary_id, primary_key: true
      add :key, :string, null: false
      add :value, :jsonb, null: false

      timestamps(type: :utc_datetime_usec, inserted_at: false)
    end

    create unique_index(:configs, [:key])

    # GIN index for JSONB value queries
    execute(
      "CREATE INDEX configs_value_gin ON configs USING GIN (value jsonb_path_ops)",
      "DROP INDEX configs_value_gin"
    )
  end
end
```

### Create Audit Logs Table

```elixir
# priv/repo/migrations/20260131000004_create_audit_logs.exs
defmodule StaConnector.Repo.Migrations.CreateAuditLogs do
  use Ecto.Migration

  def change do
    create table(:audit_logs, primary_key: false) do
      add :id, :binary_id, primary_key: true
      add :action, :string, null: false
      add :actor, :string, null: false
      add :details, :jsonb, null: false, default: "{}"

      timestamps(type: :utc_datetime_usec, updated_at: false)
    end

    create index(:audit_logs, [:action])
    create index(:audit_logs, [:actor])
    create index(:audit_logs, [:inserted_at])
    create index(:audit_logs, [:action, :inserted_at])

    # GIN index for JSONB queries
    execute(
      "CREATE INDEX audit_logs_details_gin ON audit_logs USING GIN (details jsonb_path_ops)",
      "DROP INDEX audit_logs_details_gin"
    )

    # Partial index for recent logs (commonly queried)
    execute(
      "CREATE INDEX audit_logs_recent ON audit_logs (inserted_at DESC) WHERE inserted_at > NOW() - INTERVAL '30 days'",
      "DROP INDEX audit_logs_recent"
    )
  end
end
```

## Database Constraints

### Check Constraints Migration

```elixir
# priv/repo/migrations/20260131000005_add_check_constraints.exs
defmodule StaConnector.Repo.Migrations.AddCheckConstraints do
  use Ecto.Migration

  def change do
    # Files table constraints
    create constraint(:files, :valid_direction,
      check: "direction IN ('inbound', 'outbound')"
    )

    create constraint(:files, :valid_status,
      check: "status IN ('pending', 'uploading', 'uploaded', 'processing', 'completed', 'failed', 'retrying')"
    )

    create constraint(:files, :valid_system,
      check: "system IN ('CCS', 'CIR', 'CMP', 'STR', 'SPI', 'CAM', 'LDL', 'DICT')"
    )

    create constraint(:files, :protocol_format,
      check: "protocol ~ '^[A-Z0-9]{10,20}$'"
    )

    # File events constraint
    create constraint(:file_events, :valid_event_type,
      check: "event_type IN ('created', 'status_changed', 'uploaded', 'downloaded', 'processing_started', 'processing_completed', 'processing_failed', 'retry_scheduled', 'callback_sent', 'error')"
    )

    # Configs constraint
    create constraint(:configs, :key_format,
      check: "key ~ '^[a-z][a-z0-9_.]+$'"
    )
  end
end
```

## Ecto Changesets and Validation

### File Metadata Validation

```elixir
# lib/sta_connector/files/file_metadata.ex
defmodule StaConnector.Files.FileMetadata do
  @moduledoc """
  Validates and structures file metadata stored in JSONB.
  """

  import Ecto.Changeset

  @outbound_required ~w(file_name file_type file_size checksum source_system)
  @inbound_required ~w(file_name file_type sta_protocol)

  @doc """
  Validates metadata based on file direction.
  """
  def validate_metadata(changeset) do
    direction = get_field(changeset, :direction)
    metadata = get_field(changeset, :metadata) || %{}

    case direction do
      "outbound" -> validate_outbound_metadata(changeset, metadata)
      "inbound" -> validate_inbound_metadata(changeset, metadata)
      _ -> changeset
    end
  end

  defp validate_outbound_metadata(changeset, metadata) do
    missing = @outbound_required -- Map.keys(metadata)

    if Enum.empty?(missing) do
      changeset
      |> validate_file_size(metadata)
      |> validate_checksum(metadata)
    else
      add_error(changeset, :metadata, "missing required fields: #{Enum.join(missing, ", ")}")
    end
  end

  defp validate_inbound_metadata(changeset, metadata) do
    missing = @inbound_required -- Map.keys(metadata)

    if Enum.empty?(missing) do
      changeset
    else
      add_error(changeset, :metadata, "missing required fields: #{Enum.join(missing, ", ")}")
    end
  end

  defp validate_file_size(changeset, %{"file_size" => size}) when is_integer(size) and size > 0 do
    changeset
  end
  defp validate_file_size(changeset, _) do
    add_error(changeset, :metadata, "file_size must be a positive integer")
  end

  defp validate_checksum(changeset, %{"checksum" => checksum}) when is_binary(checksum) do
    if String.match?(checksum, ~r/^[a-f0-9]{64}$/i) do
      changeset
    else
      add_error(changeset, :metadata, "checksum must be a valid SHA-256 hash")
    end
  end
  defp validate_checksum(changeset, _) do
    add_error(changeset, :metadata, "checksum is required")
  end
end
```

### Embedded Schema for STA State

```elixir
# lib/sta_connector/files/sta_state.ex
defmodule StaConnector.Files.STAState do
  @moduledoc """
  Embedded schema for STA-specific state stored in file metadata.
  """

  use Ecto.Schema
  import Ecto.Changeset

  @primary_key false
  embedded_schema do
    field :transfer_id, :string
    field :external_id, :string
    field :sta_status, :string
    field :sta_message, :string
    field :sta_timestamp, :utc_datetime_usec
    field :acknowledged_at, :utc_datetime_usec
    field :delivered_at, :utc_datetime_usec
  end

  def changeset(state, attrs) do
    state
    |> cast(attrs, [:transfer_id, :external_id, :sta_status, :sta_message,
                    :sta_timestamp, :acknowledged_at, :delivered_at])
  end
end
```

## Context Modules

### Files Context

```elixir
# lib/sta_connector/files.ex
defmodule StaConnector.Files do
  @moduledoc """
  Context for file operations.
  """

  import Ecto.Query
  alias StaConnector.Repo
  alias StaConnector.Files.{File, FileEvent}

  @doc """
  Creates a new file record.
  """
  def create_file(attrs) do
    %File{}
    |> File.changeset(attrs)
    |> Repo.insert()
    |> case do
      {:ok, file} ->
        create_event(file, "created", %{initial_status: file.status})
        {:ok, file}
      error ->
        error
    end
  end

  @doc """
  Updates file status with event tracking.
  """
  def update_status(%File{} = file, new_status, details \\ %{}) do
    old_status = file.status

    file
    |> File.status_changeset(new_status)
    |> Repo.update()
    |> case do
      {:ok, updated_file} ->
        create_event(updated_file, "status_changed", %{
          from: old_status,
          to: new_status,
          details: details
        })
        {:ok, updated_file}
      error ->
        error
    end
  end

  @doc """
  Gets a file by ID with preloaded events.
  """
  def get_file(id) do
    File
    |> Repo.get(id)
    |> Repo.preload(:events)
  end

  @doc """
  Gets a file by protocol number.
  """
  def get_file_by_protocol(protocol) do
    Repo.get_by(File, protocol: protocol)
  end

  @doc """
  Checks if a protocol has been processed.
  """
  def protocol_processed?(protocol) do
    query = from f in File,
      where: f.protocol == ^protocol and f.status == "completed",
      select: count(f.id)

    Repo.one(query) > 0
  end

  @doc """
  Lists files with filters.
  """
  def list_files(filters \\ %{}, opts \\ []) do
    limit = Keyword.get(opts, :limit, 50)
    offset = Keyword.get(opts, :offset, 0)

    File
    |> apply_filters(filters)
    |> order_by([f], desc: f.inserted_at)
    |> limit(^limit)
    |> offset(^offset)
    |> Repo.all()
  end

  defp apply_filters(query, filters) do
    Enum.reduce(filters, query, fn
      {:system, system}, q -> where(q, [f], f.system == ^system)
      {:direction, dir}, q -> where(q, [f], f.direction == ^dir)
      {:status, status}, q -> where(q, [f], f.status == ^status)
      {:since, since}, q -> where(q, [f], f.inserted_at >= ^since)
      _, q -> q
    end)
  end

  @doc """
  Creates a file event.
  """
  def create_event(%File{id: file_id}, event_type, details \\ %{}) do
    %FileEvent{}
    |> FileEvent.changeset(%{
      file_id: file_id,
      event_type: event_type,
      details: details
    })
    |> Repo.insert()
  end

  @doc """
  Gets file statistics.
  """
  def get_stats(since \\ nil) do
    since = since || DateTime.add(DateTime.utc_now(), -24, :hour)

    query = from f in File,
      where: f.inserted_at >= ^since,
      group_by: [f.status, f.system, f.direction],
      select: %{
        status: f.status,
        system: f.system,
        direction: f.direction,
        count: count(f.id)
      }

    Repo.all(query)
  end
end
```

### Settings Context

```elixir
# lib/sta_connector/settings.ex
defmodule StaConnector.Settings do
  @moduledoc """
  Context for configuration management.
  """

  import Ecto.Query
  alias StaConnector.Repo
  alias StaConnector.Settings.Config
  alias StaConnector.Audit

  @doc """
  Gets a configuration value by key.
  """
  def get_config(key) do
    case Repo.get_by(Config, key: key) do
      nil -> {:error, :not_found}
      config -> {:ok, config.value}
    end
  end

  @doc """
  Sets a configuration value (upsert).
  """
  def set_config(key, value, actor \\ "system") do
    result =
      case Repo.get_by(Config, key: key) do
        nil ->
          %Config{}
          |> Config.changeset(%{key: key, value: value})
          |> Repo.insert()

        existing ->
          existing
          |> Config.changeset(%{value: value})
          |> Repo.update()
      end

    case result do
      {:ok, config} ->
        Audit.log("config_updated", actor, %{key: key, value: value})
        {:ok, config}
      error ->
        error
    end
  end

  @doc """
  Deletes a configuration key.
  """
  def delete_config(key, actor \\ "system") do
    case Repo.get_by(Config, key: key) do
      nil ->
        {:error, :not_found}

      config ->
        Audit.log("config_deleted", actor, %{key: key})
        Repo.delete(config)
    end
  end

  @doc """
  Lists all configurations.
  """
  def list_configs do
    Config
    |> order_by([c], asc: c.key)
    |> Repo.all()
  end
end
```

## PostgreSQL Configuration

### Database Configuration

```elixir
# config/config.exs
config :sta_connector, StaConnector.Repo,
  adapter: Ecto.Adapters.Postgres

# config/dev.exs
config :sta_connector, StaConnector.Repo,
  username: "postgres",
  password: "postgres",
  hostname: "localhost",
  database: "sta_connector_dev",
  stacktrace: true,
  show_sensitive_data_on_connection_error: true,
  pool_size: 10

# config/prod.exs
config :sta_connector, StaConnector.Repo,
  url: {:system, "DATABASE_URL"},
  pool_size: String.to_integer(System.get_env("POOL_SIZE") || "25"),
  ssl: true,
  ssl_opts: [
    verify: :verify_peer,
    cacertfile: "/etc/ssl/certs/ca-certificates.crt"
  ]
```

### Repo Configuration

```elixir
# lib/sta_connector/repo.ex
defmodule StaConnector.Repo do
  use Ecto.Repo,
    otp_app: :sta_connector,
    adapter: Ecto.Adapters.Postgres

  @doc """
  Dynamically configures the repo from runtime environment.
  """
  def init(_type, config) do
    {:ok, Keyword.put(config, :url, System.get_env("DATABASE_URL"))}
  end
end
```

## Database Indexes Summary

| Table | Index | Type | Purpose |
|-------|-------|------|---------|
| `files` | `protocol` | Unique B-tree | Fast protocol lookups, prevent duplicates |
| `files` | `system` | B-tree | Filter by BCB system |
| `files` | `direction` | B-tree | Filter inbound/outbound |
| `files` | `status` | B-tree | Filter by processing status |
| `files` | `inserted_at` | B-tree | Time-based queries |
| `files` | `system, direction, status` | Composite | Dashboard queries |
| `files` | `metadata` | GIN | JSONB queries |
| `file_events` | `file_id` | B-tree | Event lookups by file |
| `file_events` | `event_type` | B-tree | Filter by event type |
| `file_events` | `inserted_at` | B-tree | Time-based queries |
| `file_events` | `details` | GIN | JSONB queries |
| `configs` | `key` | Unique B-tree | Fast key lookups |
| `configs` | `value` | GIN | JSONB queries |
| `audit_logs` | `action` | B-tree | Filter by action type |
| `audit_logs` | `actor` | B-tree | Filter by actor |
| `audit_logs` | `inserted_at` | B-tree | Time-based queries |
| `audit_logs` | `details` | GIN | JSONB queries |
| `audit_logs` | Recent partial | B-tree | Last 30 days optimization |

## Query Examples

### JSONB Queries

```elixir
# Find files by metadata field
def files_by_source_system(source) do
  from f in File,
    where: fragment("metadata->>'source_system' = ?", ^source)
end

# Find files with specific STA status in metadata
def files_with_sta_status(sta_status) do
  from f in File,
    where: fragment("metadata->'sta_state'->>'sta_status' = ?", ^sta_status)
end

# Search audit logs by details
def audit_logs_for_file(file_id) do
  from a in AuditLog,
    where: fragment("details->>'file_id' = ?", ^file_id),
    order_by: [desc: a.inserted_at]
end
```

### Statistics Queries

```elixir
# Count files by status for dashboard
def status_counts do
  from f in File,
    group_by: f.status,
    select: {f.status, count(f.id)}
end

# Processing metrics over time
def hourly_processing_stats(hours \\ 24) do
  since = DateTime.add(DateTime.utc_now(), -hours, :hour)

  from f in File,
    where: f.inserted_at >= ^since,
    group_by: fragment("date_trunc('hour', ?)", f.inserted_at),
    select: %{
      hour: fragment("date_trunc('hour', ?)", f.inserted_at),
      total: count(f.id),
      completed: count(fragment("CASE WHEN status = 'completed' THEN 1 END")),
      failed: count(fragment("CASE WHEN status = 'failed' THEN 1 END"))
    },
    order_by: fragment("date_trunc('hour', ?)", f.inserted_at)
end
```

## Best Practices

### 1. Use Ecto Changesets for Validation

```elixir
# Always validate through changesets
def update_file(%File{} = file, attrs) do
  file
  |> File.changeset(attrs)
  |> Repo.update()
end
```

### 2. Use Transactions for Related Operations

```elixir
def create_file_with_event(attrs) do
  Repo.transaction(fn ->
    with {:ok, file} <- create_file(attrs),
         {:ok, _event} <- create_event(file, "created") do
      file
    else
      {:error, reason} -> Repo.rollback(reason)
    end
  end)
end
```

### 3. Preload Associations When Needed

```elixir
def get_file_with_events(id) do
  File
  |> Repo.get(id)
  |> Repo.preload(events: from(e in FileEvent, order_by: [desc: e.inserted_at]))
end
```

### 4. Use Streaming for Large Datasets

```elixir
def export_all_files do
  File
  |> order_by([f], asc: f.inserted_at)
  |> Repo.stream()
  |> Stream.each(&process_file/1)
  |> Stream.run()
end
```

### 5. Configure Connection Pool Properly

```elixir
# Production pool settings
config :sta_connector, StaConnector.Repo,
  pool_size: 25,
  queue_target: 5000,
  queue_interval: 1000
```

## Next Steps

- [Concurrency Design](/architecture/concurrency) - Elixir concurrency patterns
- [Architecture Overview](/architecture/) - System components
- [API Reference](/api/) - REST API documentation
