Packages

phoenix_kit

1.7.222
1.7.224 1.7.223 1.7.222 1.7.221 1.7.220 1.7.219 1.7.218 1.7.217 1.7.216 1.7.215 1.7.214 1.7.213 1.7.212 1.7.211 1.7.210 1.7.209 1.7.208 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 phoenix_kit integrations events.ex
Raw

lib/phoenix_kit/integrations/events.ex

defmodule PhoenixKit.Integrations.Events do
@moduledoc """
PubSub helpers for broadcasting integration changes in real-time.
Lets the settings pages update live without a refresh. There are two topic
scopes:
* the global system topic `"phoenix_kit:integrations"` — system-wide
connection changes (the website-wide page + consumers subscribe here);
* a per-user topic `"phoenix_kit:integrations:<user_uuid>"` — a user's own
personal connection changes (the personal page subscribes here).
Each broadcaster routes by the connection's owner, so a personal change only
reaches that user's page and a system change only reaches the system page.
## Events
- `{:integration_setup_saved, provider_key, data}` — setup credentials saved
- `{:integration_connected, provider_key, data}` — OAuth flow completed
- `{:integration_disconnected, provider_key}` — disconnected
- `{:integration_validated, provider_key, :ok | {:error, reason}}` — health check
- `{:integration_connection_added | _removed, provider_key, name}`
- `{:integration_connection_renamed, provider_key, old, new}`
The `data` payload on setup_saved/connected is REDACTED — secret fields are
stripped before broadcast so credentials never travel over PubSub (consumers
reload from the owner-scoped context anyway).
"""
alias PhoenixKit.PubSub.Manager
@topic "phoenix_kit:integrations"
@user_topic_prefix "phoenix_kit:integrations:"
@typedoc "Owner scope for routing a broadcast (mirrors `PhoenixKit.Integrations.owner/0`)."
@type owner :: :system | {atom(), String.t()}
# Mirror of PhoenixKit.Integrations.Encryption's sensitive fields — stripped
# from any broadcast payload so decrypted credentials never leave the context.
# `oauth_state` is the CSRF nonce for an in-flight OAuth handshake — not a
# stored credential, but still must never ride a broadcast.
@sensitive_fields ~w(access_token refresh_token client_secret api_key bot_token secret_key password oauth_state)
@doc "The per-user personal-integrations topic."
@spec topic_for_user(String.t()) :: String.t()
def topic_for_user(user_uuid) when is_binary(user_uuid), do: @user_topic_prefix <> user_uuid
@doc "Subscribe to system-wide integration change events (the global topic)."
@spec subscribe() :: :ok | {:error, term()}
def subscribe, do: Manager.subscribe(@topic)
@doc "Subscribe to a user's personal integration change events."
@spec subscribe(String.t()) :: :ok | {:error, term()}
def subscribe(user_uuid) when is_binary(user_uuid),
do: Manager.subscribe(topic_for_user(user_uuid))
@doc "Broadcast that an integration's setup credentials were saved. Routes + redacts by `data`'s owner."
@spec broadcast_setup_saved(String.t(), map()) :: :ok
def broadcast_setup_saved(provider_key, data) do
broadcast(owner_of(data), {:integration_setup_saved, provider_key, redact(data)})
end
@doc "Broadcast that an OAuth integration was connected. Routes + redacts by `data`'s owner."
@spec broadcast_connected(String.t(), map()) :: :ok
def broadcast_connected(provider_key, data) do
broadcast(owner_of(data), {:integration_connected, provider_key, redact(data)})
end
@doc "Broadcast that an integration was disconnected."
@spec broadcast_disconnected(String.t(), owner()) :: :ok
def broadcast_disconnected(provider_key, owner \\ :system) do
broadcast(owner, {:integration_disconnected, provider_key})
end
@doc "Broadcast that an integration's health check completed."
@spec broadcast_validated(String.t(), :ok | {:error, term()}, owner()) ::
:ok
def broadcast_validated(provider_key, status, owner \\ :system) do
broadcast(owner, {:integration_validated, provider_key, status})
end
@doc "Broadcast that a new named connection was added for a provider."
@spec broadcast_connection_added(String.t(), String.t(), owner()) :: :ok
def broadcast_connection_added(provider_key, name, owner \\ :system) do
broadcast(owner, {:integration_connection_added, provider_key, name})
end
@doc "Broadcast that a named connection was removed from a provider."
@spec broadcast_connection_removed(String.t(), String.t(), owner()) :: :ok
def broadcast_connection_removed(provider_key, name, owner \\ :system) do
broadcast(owner, {:integration_connection_removed, provider_key, name})
end
@doc "Broadcast that a named connection was renamed."
@spec broadcast_connection_renamed(
String.t(),
String.t(),
String.t(),
owner()
) :: :ok
def broadcast_connection_renamed(provider_key, old_name, new_name, owner \\ :system) do
broadcast(owner, {:integration_connection_renamed, provider_key, old_name, new_name})
end
# ── Internals ──────────────────────────────────────────────────────────
# A row's owner from its (decrypted) JSONB: a valid-looking uuid ⇒ a typed
# owner (`owner_type`, defaulting to `:user` for pre-typed rows), else system.
# Kept local so Events has no dependency back into the context.
defp owner_of(data) when is_map(data) do
case Map.get(data, "owner_uuid") do
uuid when is_binary(uuid) and uuid != "" -> {owner_type_atom(data), uuid}
_ -> :system
end
end
defp owner_of(_), do: :system
defp owner_type_atom(data) do
case Map.get(data, "owner_type") do
type when is_binary(type) and type != "" ->
try do
String.to_existing_atom(type)
rescue
ArgumentError -> :__unknown_owner__
end
_ ->
:user
end
end
defp redact(data) when is_map(data), do: Map.drop(data, @sensitive_fields)
defp redact(other), do: other
# Route by owner. `:system` → global topic; a typed owner → its own topic
# (`:user` keeps the bare `<prefix><uuid>` topic the personal page subscribes
# to; other types get `<prefix><type>:<id>`). A malformed owner never crashes
# a broadcast — it no-ops.
defp broadcast(:system, message), do: do_broadcast(@topic, message)
defp broadcast({type, id}, message) when is_atom(type) and is_binary(id),
do: do_broadcast(topic_for_owner(type, id), message)
defp broadcast(_owner, _message), do: :ok
defp topic_for_owner(:user, id), do: topic_for_user(id)
defp topic_for_owner(type, id), do: @user_topic_prefix <> Atom.to_string(type) <> ":" <> id
defp do_broadcast(topic, message) do
Manager.broadcast(topic, message)
rescue
_ -> :ok
end
end