Packages

phoenix_kit

1.7.41
1.7.207 1.7.206 1.7.205 1.7.204 1.7.203 1.7.202 1.7.201 1.7.200 1.7.199 1.7.198 1.7.197 1.7.196 1.7.194 1.7.193 1.7.192 1.7.191 1.7.190 1.7.189 1.7.187 1.7.186 1.7.185 1.7.184 1.7.183 1.7.182 1.7.181 1.7.180 1.7.179 1.7.178 1.7.177 1.7.176 1.7.175 1.7.174 1.7.173 1.7.172 1.7.171 1.7.170 1.7.169 1.7.168 1.7.167 1.7.166 1.7.165 1.7.164 1.7.162 1.7.161 1.7.160 1.7.159 1.7.157 1.7.156 1.7.155 1.7.154 1.7.153 1.7.152 1.7.151 1.7.150 1.7.149 1.7.146 1.7.145 1.7.144 1.7.143 1.7.138 1.7.133 1.7.132 1.7.131 1.7.130 1.7.128 1.7.126 1.7.125 1.7.121 1.7.120 1.7.119 1.7.118 1.7.117 1.7.116 1.7.115 1.7.114 1.7.113 1.7.112 1.7.111 1.7.110 1.7.109 1.7.108 1.7.107 1.7.106 1.7.105 1.7.104 1.7.103 1.7.102 1.7.101 1.7.100 1.7.99 1.7.98 1.7.97 1.7.96 1.7.95 1.7.94 1.7.93 1.7.92 1.7.91 1.7.90 1.7.89 1.7.88 1.7.87 1.7.86 1.7.85 1.7.84 1.7.83 1.7.82 1.7.81 1.7.80 1.7.79 1.7.78 1.7.77 1.7.76 1.7.75 1.7.74 1.7.71 1.7.70 1.7.69 1.7.66 1.7.65 1.7.64 1.7.63 1.7.62 1.7.61 1.7.59 1.7.58 1.7.57 1.7.56 1.7.55 1.7.54 1.7.53 1.7.52 1.7.51 1.7.49 1.7.44 1.7.43 1.7.42 1.7.41 1.7.39 1.7.38 1.7.37 1.7.36 1.7.34 1.7.33 1.7.31 1.7.30 1.7.29 1.7.28 1.7.27 1.7.26 1.7.25 1.7.24 1.7.23 1.7.22 1.7.21 1.7.20 1.7.19 1.7.18 1.7.17 1.7.16 1.7.15 1.7.14 1.7.13 1.7.12 1.7.11 1.7.10 1.7.9 1.7.8 1.7.7 1.7.6 1.7.5 1.7.4 1.7.3 1.7.2 1.7.1 1.7.0 1.6.20 1.6.19 1.6.18 1.6.17 1.6.16 1.6.15 1.6.14 1.6.13 1.6.12 1.6.11 1.6.10 1.6.9 1.6.8 1.6.7 1.6.6 1.6.5 1.6.4 1.6.3 1.5.2 1.5.1 1.5.0 1.4.9 1.4.8 1.4.7 1.4.6 1.4.5 1.4.4 1.4.3 1.4.2 1.4.1 1.4.0 1.3.2 1.3.1 1.3.0 1.2.10 1.2.9 1.2.8 1.2.7 1.2.5 1.2.4 1.2.2 1.2.1 1.2.0 1.1.0 1.0.0

A foundation for building Elixir Phoenix apps — SaaS, social networks, ERP systems, marketplaces, and more

Current section

Files

Jump to
phoenix_kit lib modules emails sqs_processor.ex
Raw

lib/modules/emails/sqs_processor.ex

defmodule PhoenixKit.Modules.Emails.SQSProcessor do
@moduledoc """
Processor for handling email events from AWS SQS messages.
This module is responsible for:
- Parsing SNS messages from SQS
- Processing different types of SES events
- Updating email statuses in the database
- Creating event records for tracking
## Supported Event Types
- **Send** - Email send confirmation through SES
- **Delivery** - Successful email delivery to recipient
- **Bounce** - Email bounce (hard/soft bounce)
- **Complaint** - Spam complaint
- **Open** - Email open (AWS SES tracking)
- **Click** - Link click in email
## Processing Architecture
```
SQS Message → SNS Parsing → Event Processing → Database Update
```
## Security
- Message structure validation
- Event type checking
- Protection against event duplication
- Graceful handling of invalid data
## Examples
# Parse SNS message
{:ok, event_data} = SQSProcessor.parse_sns_message(sqs_message)
# Process event
{:ok, result} = SQSProcessor.process_email_event(event_data)
"""
require Logger
alias PhoenixKit.Modules.Emails
alias PhoenixKit.Modules.Emails.Event
alias PhoenixKit.Modules.Emails.Log
## --- Public API ---
@doc """
Parses SNS message from SQS into event data structure.
## Parameters
- `sqs_message` - message from SQS queue
## Returns
- `{:ok, event_data}` - successfully parsed event data
- `{:error, reason}` - parsing error
## Examples
iex> SQSProcessor.parse_sns_message(sqs_message)
{:ok, %{
"eventType" => "delivery",
"mail" => %{"messageId" => "abc123"},
"delivery" => %{"timestamp" => "2025-09-20T15:30:45.000Z"}
}}
"""
def parse_sns_message(%{"Body" => body}) do
parse_sns_body(body)
end
def parse_sns_message(%{"body" => body}) do
parse_sns_body(body)
end
def parse_sns_message(%{body: body}) do
parse_sns_body(body)
end
def parse_sns_message(_), do: {:error, :invalid_message_format}
# Helper function to parse SNS body content
defp parse_sns_body(body) when is_binary(body) do
# Validate body is not empty
if String.trim(body) == "" do
Logger.error("Received empty SNS message body")
{:error, :empty_message_body}
else
with {:ok, sns_data} when is_map(sns_data) <- Jason.decode(body),
{:ok, event_data} <- extract_ses_event(sns_data) do
{:ok, event_data}
else
{:ok, invalid_data} ->
Logger.error("SNS body decoded but not a map", %{
data_type: type_of(invalid_data),
data_preview: inspect(invalid_data) |> String.slice(0, 200)
})
{:error, :invalid_sns_format}
{:error, %Jason.DecodeError{} = error} ->
Logger.error("Invalid JSON in SNS message body", %{
error: inspect(error),
position: error.position,
body_preview: String.slice(body, 0, 500)
})
{:error, :invalid_json}
{:error, reason} ->
{:error, reason}
end
end
end
defp parse_sns_body(_) do
Logger.error("SNS body is not a binary string")
{:error, :invalid_body_type}
end
# Helper to get type name for better error logging
defp type_of(data) when is_list(data), do: :list
defp type_of(data) when is_map(data), do: :map
defp type_of(data) when is_binary(data), do: :binary
defp type_of(data) when is_integer(data), do: :integer
defp type_of(data) when is_float(data), do: :float
defp type_of(data) when is_atom(data), do: :atom
defp type_of(_), do: :unknown
@doc """
Processes email event and updates corresponding database records.
## Parameters
- `event_data` - event data from SNS
## Returns
- `{:ok, result}` - successful processing
- `{:error, reason}` - processing error
## Examples
iex> SQSProcessor.process_email_event(event_data)
{:ok, %{type: "delivery", log_id: 123, updated: true}}
"""
def process_email_event(event_data) when is_map(event_data) do
case determine_event_type(event_data) do
"send" ->
process_send_event(event_data)
"delivery" ->
process_delivery_event(event_data)
"bounce" ->
process_bounce_event(event_data)
"complaint" ->
process_complaint_event(event_data)
"open" ->
process_open_event(event_data)
"click" ->
process_click_event(event_data)
"reject" ->
process_reject_event(event_data)
"delivery_delay" ->
process_delivery_delay_event(event_data)
"subscription" ->
process_subscription_event(event_data)
"rendering_failure" ->
process_rendering_failure_event(event_data)
unknown_type ->
Logger.warning("Unknown email event type", %{type: unknown_type})
{:error, {:unknown_event_type, unknown_type}}
end
end
def process_email_event(_), do: {:error, :invalid_event_data}
## --- Private Helper Functions ---
# Helper function to handle placeholder log creation with configuration check
defp handle_placeholder_creation(event_data, message_id, event_type, status, callback_fn) do
if Emails.placeholder_logs_enabled?() do
Logger.warning(
"[SYNC ISSUE] #{String.capitalize(event_type)} event for unknown email - creating placeholder log",
%{
message_id: message_id,
event_type: event_type,
recommendation: "Check EmailInterceptor synchronization"
}
)
case create_placeholder_log_from_event(event_data, status) do
{:ok, log} ->
case callback_fn.(log) do
{:ok, result} ->
{:ok, Map.put(result, :created_placeholder, true)}
error ->
error
end
{:error, reason} ->
Logger.error("Failed to create placeholder log for #{event_type} event", %{
message_id: message_id,
reason: inspect(reason)
})
{:error, :email_log_not_found}
end
else
Logger.error(
"[SYNC ISSUE] #{String.capitalize(event_type)} event for unknown email - placeholder log creation disabled",
%{
message_id: message_id,
event_type: event_type,
action: "Event dropped - no email log found",
recommendation:
"Enable placeholder logs with Emails.set_placeholder_logs(true) or investigate EmailInterceptor synchronization"
}
)
{:error, :email_log_not_found}
end
end
# Extracts SES event from SNS message
defp extract_ses_event(%{"Type" => "Notification", "Message" => message_json}) do
with {:ok, :not_empty} <- validate_message_not_empty(message_json),
{:ok, :not_validation} <- validate_not_sns_validation(message_json),
{:ok, ses_event} <- decode_ses_message(message_json),
{:ok, validated_event} <- validate_ses_event_fields(ses_event) do
{:ok, validated_event}
else
error -> error
end
end
defp extract_ses_event(%{"Type" => "SubscriptionConfirmation"}) do
# SNS subscription confirmation - ignore
{:error, :subscription_confirmation}
end
defp extract_ses_event(%{"Type" => "UnsubscribeConfirmation"}) do
# SNS unsubscribe confirmation - ignore
{:error, :unsubscribe_confirmation}
end
defp extract_ses_event(data) do
Logger.error("Unknown SNS event format", %{
data_keys: Map.keys(data),
data_preview: inspect(data) |> String.slice(0, 500)
})
{:error, :unknown_sns_format}
end
# Validates message is not empty
defp validate_message_not_empty(message_json) do
if String.trim(message_json) == "" do
Logger.error("Received empty SES message JSON")
{:error, :empty_ses_message}
else
{:ok, :not_empty}
end
end
# Validates this is not an SNS topic validation message
defp validate_not_sns_validation(message_json) do
if String.contains?(message_json, "Successfully validated SNS topic") do
{:error, :sns_validation_message}
else
{:ok, :not_validation}
end
end
# Decodes the JSON message
defp decode_ses_message(message_json) do
case Jason.decode(message_json) do
{:ok, ses_event} when is_map(ses_event) ->
{:ok, ses_event}
{:ok, invalid_data} ->
Logger.error("SES message decoded but not a map", %{
data_type: type_of(invalid_data),
data_preview: inspect(invalid_data) |> String.slice(0, 500)
})
{:error, :invalid_ses_format}
{:error, %Jason.DecodeError{} = error} ->
Logger.error("Failed to decode SES message JSON - invalid JSON format", %{
error: inspect(error),
position: error.position,
message_preview: String.slice(message_json, 0, 500),
message_length: String.length(message_json)
})
{:error, :invalid_ses_message}
end
end
# Validates required SES event fields
defp validate_ses_event_fields(ses_event) do
event_type = ses_event["eventType"]
message_id = get_in(ses_event, ["mail", "messageId"])
if event_type && message_id do
{:ok, ses_event}
else
Logger.error("SES event missing required fields", %{
event_type: event_type,
message_id: message_id,
available_keys: Map.keys(ses_event),
raw_event: inspect(ses_event) |> String.slice(0, 1000)
})
{:error, :missing_required_fields}
end
end
# Determines event type based on eventType field
defp determine_event_type(event_data) do
event_data
|> Map.get("eventType", "unknown")
|> String.downcase()
end
## --- Event Processing Functions ---
# Processes send event
defp process_send_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
case find_email_log_by_message_id(message_id) do
{:ok, log} ->
# Update headers if empty
update_log_headers_if_empty(log, mail_data)
{:ok, %{type: "send", log_id: log.id, updated: false}}
{:error, :not_found} ->
# Rare case - received send event without preliminary logging
handle_placeholder_creation(event_data, message_id, "send", "sent", fn log ->
Logger.info("Created placeholder log for send event", %{
log_id: log.id,
message_id: message_id
})
{:ok, %{type: "send", log_id: log.id, updated: true}}
end)
end
end
# Processes delivery event
defp process_delivery_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
delivery_data = event_data["delivery"] || %{}
delivery_timestamp = get_in(delivery_data, ["timestamp"])
case find_email_log_by_message_id(message_id) do
{:ok, log} ->
# Update headers if empty
update_log_headers_if_empty(log, mail_data)
# Update status to delivered
update_attrs = %{
status: "delivered",
delivered_at: parse_timestamp(delivery_timestamp)
}
case Log.update_log(log, update_attrs) do
{:ok, updated_log} ->
# Create event record
create_delivery_event(updated_log, delivery_data)
Logger.info("Email delivered", %{
log_id: updated_log.id,
message_id: message_id,
delivered_at: updated_log.delivered_at
})
{:ok, %{type: "delivery", log_id: updated_log.id, updated: true}}
{:error, reason} ->
Logger.error("Failed to update delivery status", %{
log_id: log.id,
reason: inspect(reason)
})
{:error, reason}
end
{:error, :not_found} ->
handle_placeholder_creation(event_data, message_id, "delivery", "delivered", fn log ->
# Update status to delivered and add timestamp
update_attrs = %{
status: "delivered",
delivered_at: parse_timestamp(delivery_timestamp)
}
case Log.update_log(log, update_attrs) do
{:ok, updated_log} ->
# Create event record
create_delivery_event(updated_log, delivery_data)
Logger.info("Created placeholder log for delivery event", %{
log_id: updated_log.id,
message_id: message_id,
delivered_at: updated_log.delivered_at
})
{:ok, %{type: "delivery", log_id: updated_log.id, updated: true}}
{:error, reason} ->
Logger.error("Failed to update placeholder log for delivery", %{
log_id: log.id,
reason: inspect(reason)
})
{:error, reason}
end
end)
end
end
# Processes bounce event
defp process_bounce_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
bounce_data = event_data["bounce"]
bounce_type = get_in(bounce_data, ["bounceType"])
bounce_subtype = get_in(bounce_data, ["bounceSubType"])
status = determine_bounce_status(bounce_type)
update_attrs = %{
status: status,
bounced_at: DateTime.utc_now(),
error_message: build_bounce_error_message(bounce_data)
}
extra_log_data = %{
bounce_type: bounce_type,
bounce_subtype: bounce_subtype
}
process_ses_event(
message_id,
mail_data,
update_attrs,
bounce_data,
"bounce",
&create_bounce_event/2,
extra_log_data,
nil
)
end
defp determine_bounce_status(bounce_type) do
case String.downcase(bounce_type || "") do
"permanent" -> "hard_bounced"
"temporary" -> "soft_bounced"
_ -> "bounced"
end
end
# Processes complaint event
defp process_complaint_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
complaint_data = event_data["complaint"]
complaint_type = get_in(complaint_data, ["complaintFeedbackType"])
update_attrs = %{
status: "complaint",
complained_at: DateTime.utc_now(),
error_message: "Spam complaint: #{complaint_type || "unknown"}"
}
extra_log_data = %{complaint_type: complaint_type}
placeholder_opts = %{event_data: event_data, status: "complaint"}
process_ses_event(
message_id,
mail_data,
update_attrs,
complaint_data,
"complaint",
&create_complaint_event/2,
extra_log_data,
placeholder_opts
)
end
# Processes email open event
defp process_open_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
open_data = event_data["open"]
open_timestamp = get_in(open_data, ["timestamp"])
case find_email_log_by_message_id(message_id) do
{:ok, log} ->
# Update headers if empty
update_log_headers_if_empty(log, mail_data)
# Update status only if current status is not "clicked"
# (click is more important than open)
status_update =
case log.status do
# Do not change
"clicked" -> %{}
_ -> %{status: "opened"}
end
case Log.update_log(log, status_update) do
{:ok, updated_log} ->
# Create event record
create_open_event(updated_log, open_data, open_timestamp)
{:ok, %{type: "open", log_id: updated_log.id, updated: true}}
{:error, reason} ->
Logger.error("Failed to update open status", %{
log_id: log.id,
reason: inspect(reason)
})
{:error, reason}
end
{:error, :not_found} ->
handle_placeholder_creation(event_data, message_id, "open", "opened", fn log ->
# Create event record for created log
create_open_event(log, open_data, open_timestamp)
Logger.info("Created placeholder log for open event", %{
log_id: log.id,
message_id: message_id
})
{:ok, %{type: "open", log_id: log.id, updated: true}}
end)
end
end
# Processes click event
defp process_click_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
click_data = event_data["click"]
click_timestamp = get_in(click_data, ["timestamp"])
case find_email_log_by_message_id(message_id) do
{:ok, log} ->
# Update headers if empty
update_log_headers_if_empty(log, mail_data)
# Click - highest engagement level
update_attrs = %{status: "clicked"}
case Log.update_log(log, update_attrs) do
{:ok, updated_log} ->
# Create event record
create_click_event(updated_log, click_data, click_timestamp)
Logger.info("Email link clicked", %{
log_id: updated_log.id,
message_id: message_id,
link_url: get_in(click_data, ["link"]),
ip_address: get_in(click_data, ["ipAddress"])
})
{:ok, %{type: "click", log_id: updated_log.id, updated: true}}
{:error, reason} ->
Logger.error("Failed to update click status", %{
log_id: log.id,
reason: inspect(reason)
})
{:error, reason}
end
{:error, :not_found} ->
handle_placeholder_creation(event_data, message_id, "click", "clicked", fn log ->
# Click - highest engagement level
update_attrs = %{status: "clicked"}
case Log.update_log(log, update_attrs) do
{:ok, updated_log} ->
# Create event record
create_click_event(updated_log, click_data, click_timestamp)
Logger.info("Created placeholder log for click event", %{
log_id: updated_log.id,
message_id: message_id,
link_url: get_in(click_data, ["link"]),
ip_address: get_in(click_data, ["ipAddress"])
})
{:ok, %{type: "click", log_id: updated_log.id, updated: true}}
{:error, reason} ->
Logger.error("Failed to update placeholder log for click", %{
log_id: log.id,
reason: inspect(reason)
})
{:error, reason}
end
end)
end
end
# Processes reject event
defp process_reject_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
reject_data = event_data["reject"]
reject_reason = get_in(reject_data, ["reason"])
update_attrs = %{
status: "rejected",
rejected_at: DateTime.utc_now(),
error_message: build_reject_error_message(reject_data)
}
extra_log_data = %{reject_reason: reject_reason}
placeholder_opts = %{event_data: event_data, status: "rejected"}
process_ses_event(
message_id,
mail_data,
update_attrs,
reject_data,
"reject",
&create_reject_event/2,
extra_log_data,
placeholder_opts
)
end
# Processes delivery delay event
defp process_delivery_delay_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
delay_data = event_data["deliveryDelay"]
delay_type = get_in(delay_data, ["delayType"])
expiration_time = get_in(delay_data, ["expirationTime"])
case find_email_log_by_message_id(message_id) do
{:ok, log} ->
# Update headers if empty
update_log_headers_if_empty(log, mail_data)
# Only update if current status is not more advanced
status_update =
case log.status do
s
when s in [
"delivered",
"bounced",
"hard_bounced",
"soft_bounced",
"clicked",
"opened"
] ->
%{delayed_at: DateTime.utc_now()}
_ ->
%{status: "delayed", delayed_at: DateTime.utc_now()}
end
case Log.update_log(log, status_update) do
{:ok, updated_log} ->
# Create event record
create_delivery_delay_event(updated_log, delay_data)
Logger.info("Email delivery delayed", %{
log_id: updated_log.id,
message_id: message_id,
delay_type: delay_type,
expiration_time: expiration_time
})
{:ok, %{type: "delivery_delay", log_id: updated_log.id, updated: true}}
{:error, reason} ->
Logger.error("Failed to update delay status", %{
log_id: log.id,
reason: inspect(reason)
})
{:error, reason}
end
{:error, :not_found} ->
Logger.warning(
"Delivery delay event for unknown email - attempting to create placeholder log",
%{
message_id: message_id
}
)
case create_placeholder_log_from_event(event_data, "delayed") do
{:ok, log} ->
# Create event record for created log
create_delivery_delay_event(log, delay_data)
Logger.info("Created placeholder log for delay event", %{
log_id: log.id,
message_id: message_id,
delay_type: delay_type
})
{:ok,
%{type: "delivery_delay", log_id: log.id, updated: true, created_placeholder: true}}
{:error, reason} ->
Logger.error("Failed to create placeholder log for delay event", %{
message_id: message_id,
reason: inspect(reason)
})
{:error, :email_log_not_found}
end
end
end
# Processes subscription event
defp process_subscription_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
subscription_data = event_data["subscription"]
subscription_type = get_in(subscription_data, ["subscriptionType"])
case find_email_log_by_message_id(message_id) do
{:ok, log} ->
# Update headers if empty
update_log_headers_if_empty(log, mail_data)
# Create event record
create_subscription_event(log, subscription_data)
Logger.info("Email subscription event", %{
log_id: log.id,
message_id: message_id,
subscription_type: subscription_type
})
{:ok, %{type: "subscription", log_id: log.id, updated: false}}
{:error, :not_found} ->
Logger.warning(
"Subscription event for unknown email - attempting to create placeholder log",
%{
message_id: message_id
}
)
case create_placeholder_log_from_event(event_data, "sent") do
{:ok, log} ->
# Create event record for created log
create_subscription_event(log, subscription_data)
Logger.info("Created placeholder log for subscription event", %{
log_id: log.id,
message_id: message_id,
subscription_type: subscription_type
})
{:ok,
%{type: "subscription", log_id: log.id, updated: true, created_placeholder: true}}
{:error, reason} ->
Logger.error("Failed to create placeholder log for subscription event", %{
message_id: message_id,
reason: inspect(reason)
})
{:error, :email_log_not_found}
end
end
end
# Processes rendering failure event
defp process_rendering_failure_event(event_data) do
message_id = get_in(event_data, ["mail", "messageId"])
mail_data = event_data["mail"] || %{}
failure_data = event_data["failure"]
error_message = get_in(failure_data, ["errorMessage"])
template_name = get_in(failure_data, ["templateName"])
update_attrs = %{
status: "failed",
failed_at: DateTime.utc_now(),
error_message: build_rendering_failure_message(failure_data)
}
extra_log_data = %{template_name: template_name, error_message: error_message}
placeholder_opts = %{event_data: event_data, status: "failed"}
process_ses_event(
message_id,
mail_data,
update_attrs,
failure_data,
"rendering_failure",
&create_rendering_failure_event/2,
extra_log_data,
placeholder_opts
)
end
# Generic SES event processor to reduce duplication
defp process_ses_event(
message_id,
mail_data,
update_attrs,
event_specific_data,
event_type,
create_event_fn,
extra_log_data,
placeholder_opts
) do
case find_email_log_by_message_id(message_id) do
{:ok, log} ->
update_log_headers_if_empty(log, mail_data)
update_log_and_create_event(
log,
update_attrs,
event_specific_data,
event_type,
create_event_fn,
message_id,
extra_log_data
)
{:error, :not_found} when not is_nil(placeholder_opts) ->
handle_missing_log_with_placeholder(
message_id,
event_specific_data,
update_attrs,
event_type,
create_event_fn,
extra_log_data,
placeholder_opts
)
{:error, :not_found} ->
Logger.warning("#{event_type} event for unknown email", %{message_id: message_id})
{:error, :email_log_not_found}
end
end
defp handle_missing_log_with_placeholder(
message_id,
full_event_data,
update_attrs,
event_type,
create_event_fn,
extra_log_data,
placeholder_opts
) do
Logger.warning(
"#{event_type} event for unknown email - attempting to create placeholder log",
%{message_id: message_id}
)
case create_placeholder_log_from_event(
placeholder_opts[:event_data],
placeholder_opts[:status]
) do
{:ok, log} ->
case Log.update_log(log, update_attrs) do
{:ok, updated_log} ->
create_event_fn.(updated_log, full_event_data)
Logger.info(
"Created placeholder log for #{event_type} event",
Map.merge(%{log_id: updated_log.id, message_id: message_id}, extra_log_data)
)
{:ok,
%{type: event_type, log_id: updated_log.id, updated: true, created_placeholder: true}}
{:error, reason} ->
Logger.error("Failed to update placeholder log for #{event_type}", %{
log_id: log.id,
reason: inspect(reason)
})
{:error, reason}
end
{:error, reason} ->
Logger.error("Failed to create placeholder log for #{event_type} event", %{
message_id: message_id,
reason: inspect(reason)
})
{:error, :email_log_not_found}
end
end
defp update_log_and_create_event(
log,
update_attrs,
event_data,
event_type,
create_event_fn,
message_id,
extra_log_data
) do
case Log.update_log(log, update_attrs) do
{:ok, updated_log} ->
handle_event_creation(
updated_log,
event_data,
event_type,
create_event_fn,
message_id,
extra_log_data
)
{:ok, %{type: event_type, log_id: updated_log.id, updated: true}}
{:error, reason} ->
Logger.error("Failed to update #{event_type} status", %{
log_id: log.id,
reason: inspect(reason)
})
{:error, reason}
end
end
defp handle_event_creation(
log,
event_data,
event_type,
create_event_fn,
message_id,
extra_log_data
) do
case create_event_fn.(log, event_data) do
{:ok, :duplicate_event} ->
log_level = if event_type in ["bounce", "complaint", "reject"], do: :warning, else: :info
Logger.log(
log_level,
"#{event_type} event (duplicate skipped)",
Map.merge(%{log_id: log.id, message_id: message_id}, extra_log_data)
)
{:ok, _event} ->
log_level = if event_type in ["rendering_failure"], do: :error, else: :info
Logger.log(
log_level,
"#{event_type} event created",
Map.merge(%{log_id: log.id, message_id: message_id}, extra_log_data)
)
{:error, reason} ->
Logger.error("Failed to create #{event_type} event", %{
log_id: log.id,
reason: inspect(reason)
})
end
end
## --- Helper Functions ---
# Finds email log by message_id with extended search
defp find_email_log_by_message_id(message_id) when is_binary(message_id) do
# First search - direct search by message_id
case Emails.get_log_by_message_id(message_id) do
{:ok, log} ->
{:ok, log}
{:error, :not_found} ->
# Second search - search by AWS message ID
case Log.find_by_aws_message_id(message_id) do
{:ok, log} ->
{:ok, log}
{:error, :not_found} ->
Logger.warning("No email log found for message_id", %{
message_id: message_id,
searched_strategies: ["direct", "aws_field", "metadata"]
})
{:error, :not_found}
end
{:error, reason} ->
Logger.error("Error during email log search", %{
message_id: message_id,
reason: inspect(reason)
})
{:error, reason}
end
end
defp find_email_log_by_message_id(message_id) do
Logger.error("Invalid message_id format", %{
message_id: inspect(message_id),
message_id_type: type_of(message_id)
})
{:error, :invalid_message_id}
end
# Creates event record for delivery
defp create_delivery_event(log, delivery_data) do
# Check if delivery event already exists to prevent duplicates
if Event.event_exists?(log.id, "delivery") do
{:ok, :duplicate_event}
else
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "delivery",
event_data: delivery_data,
occurred_at: parse_timestamp(get_in(delivery_data, ["timestamp"]))
}
Emails.create_event(event_attrs)
end
end
# Creates event record for bounce
defp create_bounce_event(log, bounce_data) do
# Check if bounce event already exists to prevent duplicates
if Event.event_exists?(log.id, "bounce") do
{:ok, :duplicate_event}
else
# Convert AWS bounce types to our internal types
aws_bounce_type = get_in(bounce_data, ["bounceType"])
bounce_type = normalize_bounce_type(aws_bounce_type)
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "bounce",
event_data: bounce_data,
occurred_at: parse_timestamp(get_in(bounce_data, ["timestamp"])),
bounce_type: bounce_type,
bounce_subtype: get_in(bounce_data, ["bounceSubType"])
}
Emails.create_event(event_attrs)
end
end
# Converts AWS SES bounce types to internal bounce types
defp normalize_bounce_type("Permanent"), do: "hard"
defp normalize_bounce_type("Transient"), do: "soft"
defp normalize_bounce_type(_), do: "hard"
# Creates event record for complaint
defp create_complaint_event(log, complaint_data) do
# Check if complaint event already exists to prevent duplicates
if Event.event_exists?(log.id, "complaint") do
{:ok, :duplicate_event}
else
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "complaint",
event_data: complaint_data,
occurred_at: parse_timestamp(get_in(complaint_data, ["timestamp"])),
complaint_type: get_in(complaint_data, ["complaintFeedbackType"])
}
Emails.create_event(event_attrs)
end
end
# Creates event record for open
defp create_open_event(log, open_data, timestamp) do
# Check if open event already exists to prevent duplicates
if Event.event_exists?(log.id, "open") do
{:ok, :duplicate_event}
else
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "open",
event_data: open_data,
occurred_at: parse_timestamp(timestamp),
ip_address: get_in(open_data, ["ipAddress"]),
user_agent: get_in(open_data, ["userAgent"])
}
Emails.create_event(event_attrs)
end
end
# Creates event record for click
defp create_click_event(log, click_data, timestamp) do
# For clicks, we might want to allow multiple click events (different links)
# but for now, let's prevent duplicate click events too
if Event.event_exists?(log.id, "click") do
{:ok, :duplicate_event}
else
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "click",
event_data: click_data,
occurred_at: parse_timestamp(timestamp),
link_url: get_in(click_data, ["link"]),
ip_address: get_in(click_data, ["ipAddress"]),
user_agent: get_in(click_data, ["userAgent"])
}
Emails.create_event(event_attrs)
end
end
# Creates event record for reject
defp create_reject_event(log, reject_data) do
# Check if reject event already exists to prevent duplicates
if Event.event_exists?(log.id, "reject") do
{:ok, :duplicate_event}
else
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "reject",
event_data: reject_data,
occurred_at: parse_timestamp(get_in(reject_data, ["timestamp"])),
reject_reason: get_in(reject_data, ["reason"])
}
Emails.create_event(event_attrs)
end
end
# Creates event record for delivery delay
defp create_delivery_delay_event(log, delay_data) do
# Check if delivery_delay event already exists to prevent duplicates
if Event.event_exists?(log.id, "delivery_delay") do
{:ok, :duplicate_event}
else
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "delivery_delay",
event_data: delay_data,
occurred_at: parse_timestamp(get_in(delay_data, ["timestamp"])),
delay_type: get_in(delay_data, ["delayType"])
}
Emails.create_event(event_attrs)
end
end
# Creates event record for subscription
defp create_subscription_event(log, subscription_data) do
# Check if subscription event already exists to prevent duplicates
if Event.event_exists?(log.id, "subscription") do
{:ok, :duplicate_event}
else
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "subscription",
event_data: subscription_data,
occurred_at: parse_timestamp(get_in(subscription_data, ["timestamp"])),
subscription_type: get_in(subscription_data, ["subscriptionType"])
}
Emails.create_event(event_attrs)
end
end
# Creates event record for rendering failure
defp create_rendering_failure_event(log, failure_data) do
# Check if rendering_failure event already exists to prevent duplicates
if Event.event_exists?(log.id, "rendering_failure") do
{:ok, :duplicate_event}
else
event_attrs = %{
email_log_id: log.id,
email_log_uuid: log.uuid,
event_type: "rendering_failure",
event_data: failure_data,
occurred_at: parse_timestamp(get_in(failure_data, ["timestamp"])),
failure_reason: get_in(failure_data, ["errorMessage"])
}
Emails.create_event(event_attrs)
end
end
# Parses timestamp string to DateTime
defp parse_timestamp(timestamp_string) when is_binary(timestamp_string) do
case DateTime.from_iso8601(timestamp_string) do
{:ok, datetime, _} -> datetime
{:error, _} -> DateTime.utc_now()
end
end
defp parse_timestamp(_), do: DateTime.utc_now()
# Creates error message for bounce
defp build_bounce_error_message(bounce_data) do
bounce_type = get_in(bounce_data, ["bounceType"])
bounce_subtype = get_in(bounce_data, ["bounceSubType"])
recipients = get_in(bounce_data, ["bouncedRecipients"]) || []
recipient_details =
Enum.map(recipients, fn recipient ->
email = recipient["emailAddress"]
status = recipient["status"]
diagnostic = recipient["diagnosticCode"]
parts = [email, status, diagnostic] |> Enum.filter(& &1) |> Enum.join(" - ")
parts
end)
base_message = "#{bounce_type} bounce"
base_message =
if bounce_subtype, do: "#{base_message} (#{bounce_subtype})", else: base_message
if Enum.empty?(recipient_details) do
base_message
else
"#{base_message}: #{Enum.join(recipient_details, "; ")}"
end
end
# Creates placeholder email log from event data for cases
# when we receive events without a pre-created log
defp create_placeholder_log_from_event(event_data, initial_status) do
mail_data = event_data["mail"] || %{}
message_id = get_in(mail_data, ["messageId"])
# Extract main data from event
destination = get_in(mail_data, ["destination"]) || []
source = get_in(mail_data, ["source"])
# Determine recipient (first in destination list)
to_email =
case destination do
[first | _] when is_binary(first) -> first
_ -> "unknown@example.com"
end
# Determine sender
from_email =
case source do
email when is_binary(email) -> email
_ -> "unknown@example.com"
end
# Get general information from mail object
subject = get_in(mail_data, ["commonHeaders", "subject"]) || "(no subject)"
timestamp = get_in(mail_data, ["timestamp"])
log_attrs = %{
message_id: message_id,
# Store AWS message ID in dedicated field
aws_message_id: message_id,
to: to_email,
from: from_email,
subject: subject,
status: initial_status,
sent_at: parse_timestamp(timestamp),
headers: %{
"x-placeholder-log" => "true",
"x-created-from-event" => event_data["eventType"] || "unknown"
},
body_preview: "(email body not available - created from event)",
provider: "aws_ses",
template_name: "placeholder",
campaign_id: "recovered_from_event"
}
Emails.create_log(log_attrs)
end
# Builds error message for reject events
defp build_reject_error_message(reject_data) do
reason = get_in(reject_data, ["reason"]) || "unknown"
"Email rejected by SES: #{reason}"
end
# Builds error message for rendering failure events
defp build_rendering_failure_message(failure_data) do
error_message = get_in(failure_data, ["errorMessage"]) || "unknown error"
template_name = get_in(failure_data, ["templateName"])
base_message = "Template rendering failed: #{error_message}"
if template_name do
"#{base_message} (template: #{template_name})"
else
base_message
end
end
# Extract headers from AWS SES mail object
defp extract_headers_from_mail(mail_data) do
# Get headers array from mail object
headers_array = get_in(mail_data, ["headers"]) || []
common_headers = get_in(mail_data, ["commonHeaders"]) || %{}
# Parse headers array into map
parsed_headers =
headers_array
|> Enum.map(fn
%{"name" => name, "value" => value} -> {name, value}
_ -> nil
end)
|> Enum.reject(&is_nil/1)
|> Enum.into(%{})
# Normalize commonHeaders to simple map
normalized_common = normalize_common_headers(common_headers)
# Merge with priority to parsed_headers (they are more complete)
Map.merge(normalized_common, parsed_headers)
end
# Normalize commonHeaders to simple string map
defp normalize_common_headers(common_headers) when is_map(common_headers) do
common_headers
|> Enum.map(fn
{"from", [first | _]} -> {"From", first}
{"from", value} when is_binary(value) -> {"From", value}
{"to", [first | _]} -> {"To", first}
{"to", value} when is_binary(value) -> {"To", value}
{"subject", value} -> {"Subject", value}
{"messageId", value} -> {"Message-ID", value}
{"date", value} -> {"Date", value}
{"returnPath", value} -> {"Return-Path", value}
{"replyTo", [first | _]} -> {"Reply-To", first}
{"replyTo", value} when is_binary(value) -> {"Reply-To", value}
{_key, _value} -> nil
end)
|> Enum.reject(&is_nil/1)
|> Enum.into(%{})
end
defp normalize_common_headers(_), do: %{}
# Update email log headers if they are empty
defp update_log_headers_if_empty(log, mail_data) do
cond do
not Emails.save_headers_enabled?() ->
{:ok, log}
not is_nil(log.headers) and map_size(log.headers) > 0 ->
{:ok, log}
true ->
do_update_log_headers(log, mail_data)
end
end
defp do_update_log_headers(log, mail_data) do
headers = extract_headers_from_mail(mail_data)
if map_size(headers) > 0 do
case Log.update_log(log, %{headers: headers}) do
{:ok, updated_log} ->
Logger.info("Updated email log headers from SES event")
{:ok, updated_log}
{:error, changeset} ->
Logger.error("Failed to update email log headers: #{inspect(changeset.errors)}")
{:error, changeset}
end
else
{:ok, log}
end
end
end