diff --git a/lib/codex_pooler/admin/upstream_assignment_workflow.ex b/lib/codex_pooler/admin/upstream_assignment_workflow.ex new file mode 100644 index 00000000..e738bde9 --- /dev/null +++ b/lib/codex_pooler/admin/upstream_assignment_workflow.ex @@ -0,0 +1,131 @@ +defmodule CodexPooler.Admin.UpstreamAssignmentWorkflow do + @moduledoc """ + Coordinates the operator workflow for attaching an existing upstream identity to a Pool. + """ + + require Logger + + alias CodexPooler.Accounts.{Scope, User} + alias CodexPooler.Audit + alias CodexPooler.Events + alias CodexPooler.Jobs + alias CodexPooler.Pools + alias CodexPooler.Pools.Pool + alias CodexPooler.Upstreams + alias CodexPooler.Upstreams.Assignments, as: UpstreamAssignments + alias CodexPooler.Upstreams.Schemas.{PoolUpstreamAssignment, UpstreamIdentity} + + @audit_action "upstream_account.assign_pool" + @catalog_trigger_kind "manual" + + @type workflow_error :: %{required(:code) => atom(), required(:message) => String.t()} + @type workflow_result :: + {:ok, PoolUpstreamAssignment.t()} + | {:error, Ecto.Changeset.t() | workflow_error()} + + @spec assign_to_pool(Scope.t(), Pool.t(), UpstreamIdentity.t() | Ecto.UUID.t()) :: + workflow_result() + def assign_to_pool( + %Scope{user: %User{}} = scope, + %Pool{} = pool, + identity_or_id + ) do + with {:ok, _decision} <- + Pools.require_capability(scope, Pools.capability(:pool_operate), pool_id: pool.id), + {:ok, identity} <- authorize_identity(scope, identity_or_id), + {:ok, assignment} <- UpstreamAssignments.assign_pool_assignment(pool, identity) do + record_audit(scope, pool, identity, assignment) + broadcast_assignment(pool, identity, assignment) + enqueue_catalog_sync(pool) + + {:ok, assignment} + end + end + + def assign_to_pool(_scope, _pool, _identity_or_id), + do: {:error, workflow_error(:invalid_request, "user scope and Pool are required")} + + defp authorize_identity(%Scope{} = scope, identity_or_id) do + case normalize_identity(identity_or_id) do + %UpstreamIdentity{} = identity -> + authorize_loaded_identity(scope, identity) + + nil -> + {:error, workflow_error(:upstream_identity_not_found, "upstream identity was not found")} + end + end + + defp authorize_loaded_identity(scope, identity) do + if Pools.owner?(scope) or visible_identity?(scope, identity.id) do + {:ok, identity} + else + {:error, + workflow_error( + :capability_denied, + "the upstream identity is not available in the operator's Pool scope" + )} + end + end + + defp visible_identity?(scope, identity_id) do + scope + |> Upstreams.list_visible_upstream_identities() + |> Enum.any?(&(&1.id == identity_id)) + end + + defp normalize_identity(%UpstreamIdentity{id: identity_id}), + do: Upstreams.get_upstream_identity(identity_id) + + defp normalize_identity(identity_id) when is_binary(identity_id), + do: Upstreams.get_upstream_identity(identity_id) + + defp normalize_identity(_identity_or_id), do: nil + + defp record_audit(%Scope{user: %User{} = user}, pool, identity, assignment) do + Audit.record_user_event(user, %{ + pool_id: pool.id, + action: @audit_action, + target_type: "upstream_identity", + target_id: identity.id, + details: %{ + upstream_identity_id: identity.id, + pool_assignment_ids: [assignment.id], + assignment_status: assignment.status + } + }) + + :ok + end + + defp broadcast_assignment(pool, identity, assignment) do + Events.broadcast_upstreams(pool, "upstream_assignment_assigned", %{ + assignment_id: assignment.id, + upstream_identity_id: identity.id, + assignment_status: assignment.status + }) + + :ok + end + + defp enqueue_catalog_sync(pool) do + case Jobs.enqueue_catalog_sync(pool, trigger_kind: @catalog_trigger_kind) do + {:ok, _job} -> + :ok + + {:error, reason} -> + Logger.warning(fn -> + "upstream assignment catalog sync enqueue failed pool_id=#{pool.id} " <> + "trigger_kind=#{@catalog_trigger_kind} reason=#{catalog_enqueue_error_code(reason)}" + end) + + :ok + end + end + + defp catalog_enqueue_error_code(%Ecto.Changeset{}), do: "invalid_job" + defp catalog_enqueue_error_code(%{code: code}) when is_atom(code), do: Atom.to_string(code) + defp catalog_enqueue_error_code(reason) when is_atom(reason), do: Atom.to_string(reason) + defp catalog_enqueue_error_code(_reason), do: "unknown" + + defp workflow_error(code, message), do: %{code: code, message: message} +end diff --git a/lib/codex_pooler/audit.ex b/lib/codex_pooler/audit.ex index be2081c0..286ca6db 100644 --- a/lib/codex_pooler/audit.ex +++ b/lib/codex_pooler/audit.ex @@ -36,6 +36,7 @@ defmodule CodexPooler.Audit do {"Pool invite created", "invite.create"}, {"Pool invite revoked", "invite.revoke"}, {"Upstream account imported", "upstream_account.import"}, + {"Upstream account assigned to Pool", "upstream_account.assign_pool"}, {"Upstream account paused", "upstream_account.pause"}, {"Upstream account reactivated", "upstream_account.reactivate"}, {"Upstream account token refresh queued", "upstream_account.refresh_enqueue"}, diff --git a/lib/codex_pooler/upstreams.ex b/lib/codex_pooler/upstreams.ex index 444cacf9..4528da2e 100644 --- a/lib/codex_pooler/upstreams.ex +++ b/lib/codex_pooler/upstreams.ex @@ -49,9 +49,11 @@ defmodule CodexPooler.Upstreams do @spec list_upstream_identities(keyword()) :: [UpstreamIdentity.t()] def list_upstream_identities(opts \\ []) do status = Keyword.get(opts, :status) + excluded_status = Keyword.get(opts, :exclude_status) UpstreamIdentity |> maybe_where_status(status) + |> maybe_exclude_status(excluded_status) |> order_by([identity], asc: identity.account_label) |> Repo.all() end @@ -248,4 +250,9 @@ defmodule CodexPooler.Upstreams do defp maybe_where_status(query, status), do: from(identity in query, where: identity.status == ^status) + + defp maybe_exclude_status(query, nil), do: query + + defp maybe_exclude_status(query, status), + do: from(identity in query, where: identity.status != ^status) end diff --git a/lib/codex_pooler/upstreams/assignments.ex b/lib/codex_pooler/upstreams/assignments.ex index 8bcc8e4d..a8556163 100644 --- a/lib/codex_pooler/upstreams/assignments.ex +++ b/lib/codex_pooler/upstreams/assignments.ex @@ -23,6 +23,9 @@ defmodule CodexPooler.Upstreams.Assignments do defdelegate sync_pool_assignments_for_pool_edit(pool, selected_ids, opts \\ []), to: PoolAssignments + @spec assign_pool_assignment(Pool.t(), identity_ref(), map()) :: assignment_result() + defdelegate assign_pool_assignment(pool, identity_or_id, attrs \\ %{}), to: PoolAssignments + @spec put_assignment_cooldown(assignment_ref(), DateTime.t(), map()) :: assignment_result() defdelegate put_assignment_cooldown(assignment_or_id, cooldown_until, attrs \\ %{}), to: PoolAssignments diff --git a/lib/codex_pooler/upstreams/assignments/pool_assignments.ex b/lib/codex_pooler/upstreams/assignments/pool_assignments.ex index 10905f9f..23983a76 100644 --- a/lib/codex_pooler/upstreams/assignments/pool_assignments.ex +++ b/lib/codex_pooler/upstreams/assignments/pool_assignments.ex @@ -19,7 +19,7 @@ defmodule CodexPooler.Upstreams.Assignments.PoolAssignments do @health_active PoolUpstreamAssignment.active_health_status() @health_cooldown PoolUpstreamAssignment.cooldown_health_status() @health_disabled PoolUpstreamAssignment.disabled_health_status() - @pending UpstreamIdentity.pending_status() + @pending PoolUpstreamAssignment.pending_status() @type lifecycle_error :: %{required(:code) => atom(), required(:message) => String.t()} @type lifecycle_result :: {:ok, map()} | {:error, lifecycle_error()} @@ -60,6 +60,31 @@ defmodule CodexPooler.Upstreams.Assignments.PoolAssignments do def create_pool_assignment(_pool, _identity_or_id, _attrs), do: {:error, lifecycle_error(:pool_not_found, "pool was not found")} + @doc """ + Assigns an already-linked upstream identity to a Pool as an operational + assignment. + + `create_pool_assignment/3` intentionally defaults to `pending` for + onboarding. The admin "Assign to Pool" workflow must instead create an + active assignment and restore an older pending/deleted row when one already + occupies the Pool/identity slot. + """ + @spec assign_pool_assignment(Pool.t(), identity_ref(), map()) :: assignment_result() + def assign_pool_assignment(pool, identity_or_id, attrs \\ %{}) + + def assign_pool_assignment(%Pool{} = pool, identity_or_id, attrs) when is_map(attrs) do + case identity_id(identity_or_id) do + identity_id when is_binary(identity_id) -> + assign_pool_assignment_transaction(pool, identity_id, attrs) + + nil -> + {:error, lifecycle_error(:upstream_identity_not_found, "upstream identity was not found")} + end + end + + def assign_pool_assignment(_pool, _identity_or_id, _attrs), + do: {:error, lifecycle_error(:pool_not_found, "pool was not found")} + @spec sync_pool_assignments_for_pool_edit(Pool.t(), [Ecto.UUID.t()], keyword()) :: :ok | {:error, term()} def sync_pool_assignments_for_pool_edit(pool, selected_ids, opts \\ []) @@ -570,6 +595,73 @@ defmodule CodexPooler.Upstreams.Assignments.PoolAssignments do defp pool_id(id) when is_binary(id), do: id defp pool_id(_id), do: nil + defp assignment_for_pool_identity(pool_id, identity_id) do + Repo.one( + from assignment in PoolUpstreamAssignment, + where: + assignment.pool_id == ^pool_id and + assignment.upstream_identity_id == ^identity_id, + lock: "FOR UPDATE" + ) + end + + defp assign_pool_assignment_transaction(pool, identity_id, attrs) do + Repo.transaction(fn -> + identity_id + |> lock_upstream_identity() + |> assign_locked_identity(pool, attrs) + |> case do + {:ok, assignment} -> assignment + {:error, reason} -> Repo.rollback(reason) + end + end) + end + + defp lock_upstream_identity(identity_id) do + Repo.one( + from identity in UpstreamIdentity, + where: identity.id == ^identity_id, + lock: "FOR UPDATE" + ) + end + + defp assign_locked_identity(nil, _pool, _attrs) do + {:error, lifecycle_error(:upstream_identity_not_found, "upstream identity was not found")} + end + + defp assign_locked_identity(%UpstreamIdentity{status: @deleted}, _pool, _attrs) do + {:error, + lifecycle_error( + :upstream_identity_not_assignable, + "deleted upstream identities cannot be assigned" + )} + end + + defp assign_locked_identity(%UpstreamIdentity{} = identity, pool, attrs) do + assignment_attrs = + attrs + |> atomize_attrs() + |> Map.merge(%{ + status: @assignment_active, + health_status: @health_active, + eligibility_status: @eligible, + cooldown_until: nil, + disabled_at: nil + }) + + case assignment_for_pool_identity(pool.id, identity.id) do + %PoolUpstreamAssignment{status: status} = assignment + when status in [@pending, @assignment_deleted] -> + update_pool_assignment(assignment, assignment_attrs) + + %PoolUpstreamAssignment{} = assignment -> + {:ok, assignment} + + nil -> + create_pool_assignment(pool, identity, assignment_attrs) + end + end + defp atomize_attrs(attrs) when is_map(attrs) do Map.new(attrs, fn {key, value} when is_binary(key) -> {String.to_existing_atom(key), value} diff --git a/lib/codex_pooler/upstreams/lifecycle/account_lifecycle.ex b/lib/codex_pooler/upstreams/lifecycle/account_lifecycle.ex index 8e4d4dc3..77279ebd 100644 --- a/lib/codex_pooler/upstreams/lifecycle/account_lifecycle.ex +++ b/lib/codex_pooler/upstreams/lifecycle/account_lifecycle.ex @@ -20,6 +20,7 @@ defmodule CodexPooler.Upstreams.Lifecycle.AccountLifecycle do @refresh_failed UpstreamIdentity.refresh_failed_status() @reauth_required UpstreamIdentity.reauth_required_status() @deleted UpstreamIdentity.deleted_status() + @assignment_pending PoolUpstreamAssignment.pending_status() @assignment_active PoolUpstreamAssignment.active_status() @assignment_paused PoolUpstreamAssignment.paused_status() @assignment_refresh_due PoolUpstreamAssignment.refresh_due_status() @@ -31,6 +32,7 @@ defmodule CodexPooler.Upstreams.Lifecycle.AccountLifecycle do @health_disabled PoolUpstreamAssignment.disabled_health_status() @reactivatable_statuses [@active, @paused, @refresh_due, @refresh_failed] @reactivatable_assignment_statuses [ + @assignment_pending, @assignment_active, @assignment_paused, @assignment_refresh_due, diff --git a/lib/codex_pooler_web/live/admin/components/pages/upstreams/account_card.ex b/lib/codex_pooler_web/live/admin/components/pages/upstreams/account_card.ex index 2b4dd766..890ce753 100644 --- a/lib/codex_pooler_web/live/admin/components/pages/upstreams/account_card.ex +++ b/lib/codex_pooler_web/live/admin/components/pages/upstreams/account_card.ex @@ -523,6 +523,16 @@ defmodule CodexPoolerWeb.Admin.UpstreamPageComponents.AccountCard do tabindex="0" class="menu dropdown-content z-20 mt-2 w-60 rounded-box border border-base-300 bg-base-100 p-2 text-left shadow-xl" > +
  • + +
  • assign(:pool_available?, Enum.any?(assigns.pool_options, &pool_option_available?/1)) + |> assign(:upstream_actions_docs_url, @upstream_actions_docs_url) + + ~H""" + + + + + """ + end + + defp pool_option_available?({_label, value}) when is_binary(value), do: value != "" + defp pool_option_available?(_option), do: false +end diff --git a/lib/codex_pooler_web/live/admin/components/pages/upstreams/page_components.ex b/lib/codex_pooler_web/live/admin/components/pages/upstreams/page_components.ex index 07affd9b..4fa1bde8 100644 --- a/lib/codex_pooler_web/live/admin/components/pages/upstreams/page_components.ex +++ b/lib/codex_pooler_web/live/admin/components/pages/upstreams/page_components.ex @@ -7,6 +7,7 @@ defmodule CodexPoolerWeb.Admin.UpstreamPageComponents do alias CodexPoolerWeb.Admin.PoolFilterComponents alias CodexPoolerWeb.Admin.UpstreamFilterForm alias CodexPoolerWeb.Admin.UpstreamPageComponents.AccountCard + alias CodexPoolerWeb.Admin.UpstreamPageComponents.AssignPoolDialog alias CodexPoolerWeb.Admin.UpstreamPageComponents.AuthJsonDialog alias CodexPoolerWeb.Admin.UpstreamPageComponents.SavedResetComponents alias Phoenix.HTML.Form @@ -37,6 +38,8 @@ defmodule CodexPoolerWeb.Admin.UpstreamPageComponents do attr :rename_account_form, :any, default: nil attr :deleting_account, :map, default: nil attr :delete_account_form, :any, required: true + attr :assigning_pool_account, :map, default: nil + attr :assign_pool_form, :any, required: true attr :editing_saved_reset_policy, :map, default: nil attr :saved_reset_policy_form, :any, required: true attr :confirming_saved_reset_redemption, :map, default: nil @@ -89,6 +92,11 @@ defmodule CodexPoolerWeb.Admin.UpstreamPageComponents do <.rename_account_dialog account={@renaming_account} form={@rename_account_form} /> <.delete_account_dialog account={@deleting_account} form={@delete_account_form} /> + <.saved_reset_policy_dialog account={@editing_saved_reset_policy} form={@saved_reset_policy_form} diff --git a/lib/codex_pooler_web/live/admin/pages/upstreams_live.ex b/lib/codex_pooler_web/live/admin/pages/upstreams_live.ex index 47c00ce2..976341f6 100644 --- a/lib/codex_pooler_web/live/admin/pages/upstreams_live.ex +++ b/lib/codex_pooler_web/live/admin/pages/upstreams_live.ex @@ -1,8 +1,10 @@ defmodule CodexPoolerWeb.Admin.UpstreamsLive do use CodexPoolerWeb, :admin_live_view + alias CodexPooler.Admin.UpstreamAssignmentWorkflow alias CodexPooler.Events alias CodexPooler.Pools + alias CodexPooler.Pools.Pool alias CodexPooler.Upstreams.Schemas.UpstreamIdentity alias CodexPoolerWeb.Admin.Components, as: AdminComponents alias CodexPoolerWeb.Admin.PoolEventSubscriptions @@ -54,6 +56,8 @@ defmodule CodexPoolerWeb.Admin.UpstreamsLive do rename_account_form: nil, deleting_account: nil, delete_account_form: AccountLifecycleWorkflow.delete_form(nil), + assigning_pool_account: nil, + assign_pool_form: assign_pool_form(), editing_saved_reset_policy: nil, saved_reset_policy_form: saved_reset_policy_form(%{}), confirming_saved_reset_redemption: nil, @@ -273,6 +277,63 @@ defmodule CodexPoolerWeb.Admin.UpstreamsLive do update(socket, :account_panel_views, &toggle_account_panel_view(&1, identity_id, :routing))} end + def handle_event("open_assign_pool", %{"id" => identity_id}, socket) do + case find_account(socket.assigns.upstream_accounts, identity_id) do + %{identity: %UpstreamIdentity{status: status}} = account when status != "deleted" -> + {:noreply, + socket + |> close_account_workflow_dialogs() + |> assign( + assigning_pool_account: account, + assign_pool_form: assign_pool_form(socket.assigns.pools) + )} + + _missing_or_deleted -> + {:noreply, put_flash(socket, :error, "Upstream account is not available to assign")} + end + end + + def handle_event("close_assign_pool", _params, socket) do + {:noreply, close_assign_pool_dialog(socket)} + end + + def handle_event( + "assign_pool_account", + %{"assign_pool" => %{"pool_id" => pool_id}}, + socket + ) do + pool = selected_pool(socket.assigns.pools, pool_id) + + case {socket.assigns.assigning_pool_account, pool} do + {%{identity: %UpstreamIdentity{} = identity} = account, %Pool{} = pool} -> + case UpstreamAssignmentWorkflow.assign_to_pool( + socket.assigns.current_scope, + pool, + identity + ) do + {:ok, _assignment} -> + {:noreply, + socket + |> close_assign_pool_dialog() + |> reload_upstreams() + |> put_flash(:info, "#{account.label} was assigned to #{pool.name}")} + + {:error, reason} -> + {:noreply, put_flash(socket, :error, assignment_error_message(reason))} + end + + {nil, _pool} -> + {:noreply, put_flash(socket, :error, "Upstream account was not found")} + + {_account, nil} -> + {:noreply, put_flash(socket, :error, "Target Pool was not found")} + end + end + + def handle_event("assign_pool_account", _params, socket) do + {:noreply, put_flash(socket, :error, "Select an available Pool")} + end + def handle_event("cancel_saved_reset_redemption", _params, socket) do {:noreply, assign(socket, :confirming_saved_reset_redemption, nil)} end @@ -390,6 +451,8 @@ defmodule CodexPoolerWeb.Admin.UpstreamsLive do rename_account_form={@rename_account_form} deleting_account={@deleting_account} delete_account_form={@delete_account_form} + assigning_pool_account={@assigning_pool_account} + assign_pool_form={@assign_pool_form} editing_saved_reset_policy={@editing_saved_reset_policy} saved_reset_policy_form={@saved_reset_policy_form} confirming_saved_reset_redemption={@confirming_saved_reset_redemption} @@ -533,9 +596,33 @@ defmodule CodexPoolerWeb.Admin.UpstreamsLive do |> close_rename_account_dialog() |> AccountLifecycleWorkflow.close_delete() |> OAuthWorkflow.close() + |> close_assign_pool_dialog() |> close_saved_reset_policy_dialog() end + defp close_assign_pool_dialog(socket) do + assign(socket, + assigning_pool_account: nil, + assign_pool_form: assign_pool_form() + ) + end + + defp assign_pool_form(pools \\ []) do + default_pool_id = + case pools do + [%Pool{id: pool_id}] -> pool_id + _pools -> "" + end + + to_form(%{"pool_id" => default_pool_id}, as: :assign_pool) + end + + defp assignment_error_message(%{message: message}) + when is_binary(message) and message != "", + do: message + + defp assignment_error_message(_reason), do: "Upstream account could not be assigned" + defp close_rename_account_dialog(socket) do assign(socket, renaming_account: nil, diff --git a/lib/codex_pooler_web/live/admin/read_models/upstream_accounts_read_model.ex b/lib/codex_pooler_web/live/admin/read_models/upstream_accounts_read_model.ex index 27321531..11e07184 100644 --- a/lib/codex_pooler_web/live/admin/read_models/upstream_accounts_read_model.ex +++ b/lib/codex_pooler_web/live/admin/read_models/upstream_accounts_read_model.ex @@ -161,9 +161,16 @@ defmodule CodexPoolerWeb.Admin.UpstreamAccountsReadModel do assignments = attach_assignment_model_inventory(assignments, model_inventory) identities = - scope - |> Upstreams.list_visible_upstream_identities() - |> Enum.filter(&Map.has_key?(assignments, &1.id)) + if pool_filter_selected?(filters) do + visible_assigned_identities(scope, assignments) + else + case Upstreams.list_upstream_identities_for_pool_management(scope, + exclude_status: UpstreamIdentity.deleted_status() + ) do + {:ok, identities} -> identities + {:error, _reason} -> visible_assigned_identities(scope, assignments) + end + end token_burns = TokenBurnProjection.summaries(identities) @@ -177,6 +184,18 @@ defmodule CodexPoolerWeb.Admin.UpstreamAccountsReadModel do Enum.filter(pools, &MapSet.member?(visible_pool_ids, &1.id)) end + defp visible_assigned_identities(scope, assignments) do + scope + |> Upstreams.list_visible_upstream_identities() + |> Enum.filter(&Map.has_key?(assignments, &1.id)) + end + + defp pool_filter_selected?(%{"pool_id" => pool_id}) + when is_binary(pool_id) and pool_id != "", + do: true + + defp pool_filter_selected?(_filters), do: false + @spec oauth_flow_state(term(), [term()], DateTimeDisplay.preferences(), keyword()) :: oauth_flow_state() def oauth_flow_state(scope, pools, datetime_preferences, opts \\ []) diff --git a/test/codex_pooler/admin/upstream_assignment_workflow_test.exs b/test/codex_pooler/admin/upstream_assignment_workflow_test.exs new file mode 100644 index 00000000..6b789ff2 --- /dev/null +++ b/test/codex_pooler/admin/upstream_assignment_workflow_test.exs @@ -0,0 +1,90 @@ +defmodule CodexPooler.Admin.UpstreamAssignmentWorkflowTest do + use CodexPooler.DataCase, async: false + + alias CodexPooler.Accounts.Scope + alias CodexPooler.Admin.UpstreamAssignmentWorkflow + alias CodexPooler.Audit.AuditEvent + alias CodexPooler.Events + alias CodexPooler.Jobs.CatalogSyncWorker + alias CodexPooler.Repo + alias CodexPooler.Upstreams + + import CodexPooler.AccountsFixtures + import CodexPooler.PoolerFixtures + + setup do + Repo.delete_all(Oban.Job) + :ok + end + + test "authorized assignment records an audit event, broadcasts, and queues catalog sync" do + %{user: owner} = bootstrap_owner_fixture(%{"email" => unique_user_email()}) + scope = Scope.for_user(owner) + pool = pool_fixture(%{name: "Assignment Workflow Target"}) + + identity = + active_upstream_identity_fixture(%{ + account_label: "Workflow upstream", + chatgpt_account_id: "acct_assignment_workflow" + }) + + assert :ok = Events.subscribe_pool(pool, "upstreams") + + assert {:ok, assignment} = + Task.async(fn -> + UpstreamAssignmentWorkflow.assign_to_pool(scope, pool, identity) + end) + |> Task.await(5_000) + + assert assignment.pool_id == pool.id + assert assignment.upstream_identity_id == identity.id + assert assignment.status == "active" + assert assignment.health_status == "active" + assert assignment.eligibility_status == "eligible" + + assert [job] = all_enqueued(worker: CatalogSyncWorker) + assert job.args == %{"pool_id" => pool.id, "trigger_kind" => "manual"} + + assert_receive {Events, event} + assert event.pool_id == pool.id + assert event.reason == "upstream_assignment_assigned" + assert event.payload["assignment_id"] == assignment.id + assert event.payload["upstream_identity_id"] == identity.id + + assert %AuditEvent{} = + audit_event = + Repo.get_by(AuditEvent, + action: "upstream_account.assign_pool", + target_id: identity.id + ) + + assert audit_event.actor_user_id == owner.id + assert audit_event.pool_id == pool.id + assert audit_event.details["pool_assignment_ids"] == [assignment.id] + assert audit_event.details["assignment_status"] == "active" + end + + test "operator without target Pool access cannot create an assignment or side effects" do + %{user: owner} = bootstrap_owner_fixture(%{"email" => unique_user_email()}) + %{user: admin} = operator_fixture(owner, %{"email" => unique_user_email()}) + scope = Scope.for_user(admin) + pool = pool_fixture(%{name: "Hidden Assignment Target"}) + + identity = + active_upstream_identity_fixture(%{ + account_label: "Unauthorized workflow upstream", + chatgpt_account_id: "acct_assignment_workflow_denied" + }) + + assert {:error, %{code: :capability_denied}} = + UpstreamAssignmentWorkflow.assign_to_pool(scope, pool, identity) + + assert Upstreams.list_pool_assignments_for_identity(identity) == [] + assert [] = all_enqueued(worker: CatalogSyncWorker) + + refute Repo.get_by(AuditEvent, + action: "upstream_account.assign_pool", + target_id: identity.id + ) + end +end diff --git a/test/codex_pooler/upstreams_test.exs b/test/codex_pooler/upstreams_test.exs index df8a3549..61f404a5 100644 --- a/test/codex_pooler/upstreams_test.exs +++ b/test/codex_pooler/upstreams_test.exs @@ -389,6 +389,58 @@ defmodule CodexPooler.UpstreamsTest do assert Upstreams.list_eligible_pool_assignments(pool) == [] end + test "admin assignment creates an active row and restores existing rows" do + pool = pool_fixture() + identity = active_identity_fixture(%{chatgpt_account_id: "acct_admin_assign"}) + + assert Upstreams.list_pool_assignments_for_identity(identity) == [] + + assert {:ok, active_assignment} = + PoolAssignments.assign_pool_assignment(pool, identity) + + assert active_assignment.status == "active" + assert active_assignment.health_status == "active" + assert active_assignment.eligibility_status == "eligible" + + assert {:ok, same_assignment} = + PoolAssignments.assign_pool_assignment(pool, identity) + + assert same_assignment.id == active_assignment.id + + assert {:ok, %{status: :assignment_deleted}} = + PoolAssignments.delete_pool_assignment(pool, active_assignment) + + assert {:ok, restored_assignment} = + PoolAssignments.assign_pool_assignment(pool, identity) + + assert restored_assignment.id == active_assignment.id + assert restored_assignment.status == "active" + assert restored_assignment.health_status == "active" + assert restored_assignment.eligibility_status == "eligible" + assert is_nil(restored_assignment.disabled_at) + + pending_identity = + active_identity_fixture(%{chatgpt_account_id: "acct_admin_assign_pending"}) + + assert {:ok, pending_assignment} = + PoolAssignments.create_pool_assignment(pool, pending_identity) + + assert {:ok, promoted_assignment} = + PoolAssignments.assign_pool_assignment(pool, pending_identity) + + assert promoted_assignment.id == pending_assignment.id + assert promoted_assignment.status == "active" + + deleted_identity = + active_identity_fixture(%{chatgpt_account_id: "acct_admin_assign_deleted"}) + + assert {:ok, deleted_identity} = + IdentityLifecycle.update_upstream_identity(deleted_identity, %{status: "deleted"}) + + assert {:error, %{code: :upstream_identity_not_assignable}} = + PoolAssignments.assign_pool_assignment(pool, deleted_identity) + end + test "assignment list APIs return empty results for invalid pool refs" do assert Upstreams.list_active_pool_assignments(nil) == [] assert Upstreams.list_active_pool_assignments(:invalid_pool) == [] @@ -2512,6 +2564,49 @@ defmodule CodexPooler.UpstreamsTest do assert decrypted == token end + test "reactivation promotes a newly attached pending assignment" do + source_pool = pool_fixture(%{name: "Reactivation Source"}) + target_pool = pool_fixture(%{name: "Reactivation Target"}) + identity = active_identity_fixture(%{chatgpt_account_id: "acct_pending_reactivation"}) + configure_upstream_secret_key!() + + assert {:ok, source_assignment} = + PoolAssignments.create_pool_assignment(source_pool, identity) + + assert {:ok, source_assignment} = + PoolAssignments.activate_pool_assignment(source_assignment) + + assert {:ok, _secret} = + Upstreams.store_encrypted_secret(identity, %{ + secret_kind: "access_token", + plaintext: generated_secret("pending-reactivation") + }) + + scope = fixture_owner_scope() + + assert {:ok, %{status: :paused}} = + Upstreams.pause_account_for_scope(scope, identity, %{}) + + assert {:ok, %{status: :assignment_deleted}} = + PoolAssignments.delete_pool_assignment(source_pool, source_assignment) + + assert {:ok, pending_assignment} = + PoolAssignments.create_pool_assignment(target_pool, identity) + + assert pending_assignment.status == "pending" + + assert {:ok, %{status: :active}} = + Upstreams.reactivate_account_for_scope(scope, identity, %{}) + + assert %PoolUpstreamAssignment{ + status: "active", + health_status: "active", + eligibility_status: "eligible" + } = Repo.get!(PoolUpstreamAssignment, pending_assignment.id) + + assert Repo.get!(PoolUpstreamAssignment, source_assignment.id).status == "deleted" + end + test "reactivation fails when the account has no active routing secret" do pool = pool_fixture() identity = active_identity_fixture(%{chatgpt_account_id: "acct_pause_missing_secret"}) diff --git a/test/codex_pooler_web/live/admin/pages/upstreams_live_test.exs b/test/codex_pooler_web/live/admin/pages/upstreams_live_test.exs index c05a1b2d..ba4211c0 100644 --- a/test/codex_pooler_web/live/admin/pages/upstreams_live_test.exs +++ b/test/codex_pooler_web/live/admin/pages/upstreams_live_test.exs @@ -3195,6 +3195,57 @@ defmodule CodexPoolerWeb.Admin.UpstreamsLiveTest do assert Repo.get!(UpstreamIdentity, identity.id).account_label == "Renamed Codex" end + test "shows an unassigned upstream account and assigns it to a Pool", %{ + conn: conn, + scope: scope + } do + {:ok, pool} = + Pools.create_pool(scope, %{slug: "assign-existing-upstream", name: "Assignment Target"}) + + identity = + active_upstream_identity_fixture(%{ + account_email: "unassigned@example.com", + account_label: "Unassigned Account" + }) + + {:ok, filtered_view, _html} = + live(conn, ~p"/admin/upstreams?pool_id=#{pool.id}") + + refute has_element?(filtered_view, "#upstream-account-#{identity.id}") + + {:ok, view, _html} = live(conn, ~p"/admin/upstreams") + + assert has_element?(view, "#upstream-account-#{identity.id}", "Unassigned Account") + + assert has_element?( + view, + "#upstream-account-#{identity.id} [data-role='upstream-pool-count-cell']", + "No Pools" + ) + + view + |> element("#assign-pool-upstream-account-#{identity.id}") + |> render_click() + + assert has_element?(view, "#assign-pool-dialog[open]", "Assignment Target") + assert_admin_dialog_docs_link(view, "assign-pool-dialog-footer") + + view + |> form("#assign-pool-form", %{"assign_pool" => %{"pool_id" => pool.id}}) + |> render_submit() + + refute has_element?(view, "#assign-pool-dialog") + + assert has_element?( + view, + "#upstream-account-#{identity.id} [data-role='upstream-pool-count-cell']", + "1 Pool" + ) + + assert [assignment] = Upstreams.Assignments.list_active_pool_assignments(pool) + assert assignment.upstream_identity_id == identity.id + end + test "confirms upstream account deletion from the account actions menu", %{ conn: conn, scope: scope @@ -6029,6 +6080,9 @@ defmodule CodexPoolerWeb.Admin.UpstreamsLiveTest do "delete-upstream-account-dialog-footer" -> "https://docs.codex-pooler.com/operators/upstreams/#card-action-menu" + + "assign-pool-dialog-footer" -> + "https://docs.codex-pooler.com/operators/upstreams/#card-action-menu" end assert has_element?(