Packages

Clustering strategy for connecting nodes running on Google App Engine.

Current section

Files

Jump to
libcluster_gae lib google_app_engine.ex
Raw

lib/google_app_engine.ex

defmodule Cluster.Strategy.GoogleAppEngine do
@moduledoc """
Clustering strategy for Google App Engine.
This strategy only connect nodes that are able to receive HTTP traffic.
"""
use GenServer
use Cluster.Strategy
require Logger
alias Cluster.Strategy.State
@default_polling_interval 5_000
@access_token_path 'http://metadata.google.internal/computeMetadata/v1/instance/service-accounts/default/token'
def start_link(args) do
GenServer.start_link(__MODULE__, args)
end
@impl true
def init([%State{} = state]) do
{:ok, load(state)}
end
@impl true
def handle_info(:timeout, state) do
handle_info(:load, state)
end
def handle_info(:load, %State{} = state) do
{:noreply, load(state)}
end
def handle_info(_, state) do
{:noreply, state}
end
defp load(%State{} = state) do
connect = state.connect
list_nodes = state.list_nodes
topology = state.topology
nodes = get_nodes(state)
Cluster.Strategy.connect_nodes(topology, connect, list_nodes, nodes)
Process.send_after(self(), :load, polling_interval(state))
state
end
defp polling_interval(%State{config: config}) do
Keyword.get(config, :polling_interval, @default_polling_interval)
end
defp get_nodes(%State{}) do
instances = get_running_instances()
release_name = System.get_env("REL_NAME")
Enum.map(instances, & :"#{release_name}@#{&1}")
end
defp get_running_instances do
project_id = System.get_env("GOOGLE_CLOUD_PROJECT")
service_id = System.get_env("GAE_SERVICE")
versions = get_running_versions(project_id, service_id)
Enum.flat_map(versions, &get_instances_for_version(project_id, service_id, &1))
end
defp get_running_versions(project_id, service_id) do
access_token = access_token()
headers = [{'Authorization', 'Bearer #{access_token}'}]
api_url = 'https://appengine.googleapis.com/v1/apps/#{project_id}/services/#{service_id}'
case :httpc.request(:get, {api_url, headers}, [], []) do
{:ok, {{_, 200, _}, _headers, body}} ->
%{"split" => %{"allocations" => allocations}} = Jason.decode!(body)
Map.keys(allocations)
end
end
defp get_instances_for_version(project_id, service_id, version) do
access_token = access_token()
headers = [{'Authorization', 'Bearer #{access_token}'}]
api_url = 'https://appengine.googleapis.com/v1/apps/#{project_id}/services/#{service_id}/versions/#{version}/instances'
case :httpc.request(:get, {api_url, headers}, [], []) do
{:ok, {{_, 200, _}, _headers, body}} ->
handle_instances(Jason.decode!(body))
end
end
defp handle_instances(%{"instances" => instances}) do
instances
|> Enum.filter(& &1["vmStatus"] == "RUNNING")
|> Enum.map(& &1["id"])
end
defp handle_instances(_), do: []
defp access_token do
headers = [{'Metadata-Flavor', 'Google'}]
http_options = [ssl: [verify: :verify_none], timeout: 15000]
case :httpc.request(:get, {@access_token_path, headers}, http_options, []) do
{:ok, {{_, 200, _}, _headers, body}} ->
%{"access_token" => token} = Jason.decode!(body)
token
end
end
end