defmodule OpenAgents.Repositories do
@moduledoc "Canonical repository identity and membership authorization."
import Ecto.Query, warn: false
alias OpenAgents.Accounts.User
alias OpenAgents.Agents.Agent
alias OpenAgents.Forge.{Repos, WAL}
alias OpenAgents.{Analytics, Audit, Repo}
alias OpenAgents.Machines.Machine
alias OpenAgents.Repositories.{
IdempotencyRequest,
MachineGrant,
Membership,
Namespace,
NamespaceAlias,
ProvisioningOutbox,
Repository,
RepositoryImport
}
@writable_roles ~w(owner maintainer contributor)
@all_roles ~w(owner maintainer contributor viewer)
# A computer grant hands a long-lived `smct_` token standing access to a
# repository, so it takes the two roles that already administer one.
@machine_grant_roles ~w(owner maintainer)
@repository_namespace_limit 100
# GitHub's default label set. Every created or imported repository starts
# with this vocabulary so triage has something to attach on day one; a
# repository that does not want a name can delete it.
@default_labels [
{"bug", "d73a4a", "Something isn't working"},
{"documentation", "0075ca", "Improvements or additions to docs"},
{"duplicate", "cfd3d7", "This issue or pull request already exists"},
{"enhancement", "a2eeef", "New feature or request"},
{"good first issue", "7057ff", "Good for newcomers"},
{"help wanted", "008672", "Extra attention is needed"},
{"invalid", "e4e669", "This doesn't seem right"},
{"question", "d876e3", "Further information is requested"},
{"wontfix", "ffffff", "This will not be worked on"}
]
# The two durable receipts that say where a repository is in provisioning:
# the outbox row is the work, the import row is the GitHub snapshot. Both are
# `has_one`, so they are preloaded together wherever a surface renders
# progress or provenance.
@provisioning_assocs [:repository_import, :provisioning_outbox]
def get_by_path!(owner, name) when is_binary(owner) and is_binary(name) do
owner_key = String.downcase(owner)
name_key = String.downcase(name)
Repo.one!(repository_path_query(owner_key, name_key))
end
@doc """
The repository at `owner/name` an anonymous reader may read, or a raise.
Composes `readable_by/2` with a `nil` reader rather than restating the public
half. It used to carry its own copy of the predicate, which is how the copy
in `list_visible_repositories/1` came to lose the `ready` half and disagree
with the paged read about the same repository (REPOSITORY-001).
"""
def get_public_by_path!(owner, name) when is_binary(owner) and is_binary(name) do
get_visible_by_path!(owner, name, nil)
end
@doc """
The repository at `owner/name` this reader may read, or a raise.
This and `visible_by_path/3` are the two ways a caller-supplied owner and
name become a repository row a reader may see. Both compose `readable_by/2`,
so the surfaces owing the predicate are the callers of one of them and can be
read from compiled import tables
(`OpenAgents.Repositories.VisibilityJoinTest`).
"""
def get_visible_by_path!(owner, name, user) do
owner
|> visible_path_query(name)
|> readable_by(user)
|> Repo.one!()
end
@doc """
The repository at `owner/name` this reader may read, or `nil`.
The same decision as `get_visible_by_path!/3` for a caller that answers a
missing repository and an unreadable one the same way, so it does not need
to restate the join to avoid the raise.
"""
def visible_by_path(owner, name, user) do
owner
|> visible_path_query(name)
|> readable_by(user)
|> Repo.one()
end
@doc """
Narrows a repository query to the repositories `user` may read.
One predicate, composed by every surface that lists or resolves a repository,
so a new surface cannot arrive at a looser rule by rewriting the join from
memory. Reading is public on a repository that is both public and
provisioned; anything else needs a membership in a reading role. `nil` is an
anonymous visitor, who sees only the public half.
It composes into a query that already carries joins and preloads: the
membership join is appended, so the caller's own bindings keep their
positions.
"""
def readable_by(query, user)
def readable_by(query, nil) do
from repository in query,
where: repository.visibility == "public" and repository.lifecycle_state == "ready"
end
def readable_by(query, %User{id: user_id}) do
from repository in query,
left_join: reader in Membership,
on: reader.repository_id == repository.id and reader.user_id == ^user_id,
where:
(repository.visibility == "public" and repository.lifecycle_state == "ready") or
(not is_nil(reader.user_id) and reader.role in ^@all_roles)
end
def get_writable_by_path!(owner, name, %User{id: user_id}) do
owner_key = String.downcase(owner)
name_key = String.downcase(name)
Repo.one!(
from repository in repository_path_query(owner_key, name_key),
join: membership in Membership,
on:
membership.repository_id == repository.id and membership.user_id == ^user_id and
membership.role in ^@writable_roles,
join: user in User,
on: user.id == membership.user_id and user.status == "active",
where: repository.lifecycle_state == "ready"
)
end
@doc "Updates whether an owner allows new pull requests for a repository."
def update_pull_request_setting(owner, name, %User{} = actor, enabled)
when is_boolean(enabled) do
repository = get_visible_by_path!(owner, name, actor)
if owner?(repository, actor) do
repository
|> Repository.changeset(%{pull_requests_enabled: enabled})
|> Repo.update()
else
{:error, :forbidden}
end
rescue
Ecto.NoResultsError -> {:error, :not_found}
end
def update_pull_request_setting(_owner, _name, _actor, _enabled),
do: {:error, :invalid_pull_request_setting}
def create_repository(attrs) do
owner = fetch_attr!(attrs, :owner)
Repo.transaction(fn ->
namespace = ensure_legacy_namespace!(owner)
now = DateTime.utc_now()
attrs =
attrs
|> Map.new()
|> Map.put(:namespace_id, namespace.id)
|> Map.put_new(:storage_key, Ecto.UUID.generate())
|> Map.put_new(:lifecycle_state, "ready")
|> Map.put_new(:provisioning_kind, "empty")
|> Map.put_new(:ready_at, now)
%Repository{}
|> Repository.changeset(attrs)
|> Repo.insert!()
|> Repo.preload(:namespace)
end)
|> case do
{:ok, repository} -> {:ok, repository}
{:error, reason} -> {:error, reason}
end
rescue
error in Ecto.InvalidChangesetError -> {:error, error.changeset}
end
def ensure_user_namespace(%User{} = user) do
upsert_github_namespace(%{
provider_account_id: user.github_id,
slug: user.github_login,
kind: "user",
owner_user_id: user.id
})
end
def upsert_github_namespace(attrs) when is_map(attrs) do
attrs = Map.put_new(attrs, :provider_refreshed_at, DateTime.utc_now())
Repo.transaction(fn ->
provider_account_id = fetch_attr!(attrs, :provider_account_id)
kind = fetch_attr!(attrs, :kind)
existing =
Repo.one(
from namespace in Namespace,
where:
namespace.provider == "github" and
namespace.provider_account_id == ^provider_account_id and
namespace.kind == ^kind,
lock: "FOR UPDATE"
)
case existing do
nil ->
%Namespace{}
|> Namespace.changeset(attrs)
|> Repo.insert!()
%Namespace{} = namespace ->
maybe_update_namespace!(namespace, attrs)
end
end)
|> unwrap_transaction()
rescue
error in Ecto.InvalidChangesetError -> {:error, error.changeset}
end
def get_namespace_by_slug!(slug) when is_binary(slug) do
slug_key = String.downcase(slug)
case Repo.get_by(Namespace, slug_key: slug_key, state: "active") do
%Namespace{} = namespace ->
namespace
nil ->
Repo.one!(
from namespace_alias in NamespaceAlias,
join: namespace in assoc(namespace_alias, :namespace),
where: namespace_alias.slug_key == ^slug_key and namespace.state == "active",
select: namespace
)
end
end
@doc "One active namespace by slug or alias, or `nil` when no such namespace exists."
def get_namespace_by_slug(slug) when is_binary(slug) do
get_namespace_by_slug!(slug)
rescue
Ecto.NoResultsError -> nil
end
@doc """
Creates a repository in an already-resolved namespace of either kind.
`create_user_repository/3` and `create_organization_repository/4` each name
one kind because the GitHub-compatible route they serve names one kind. This
one takes whichever namespace resolving an owner produced, so the caller does
not have to know the kind before it asks.
A user namespace that belongs to somebody else is refused here as well as at
resolution. The check costs one comparison and closes the gap between the
two callers this function will eventually have.
"""
def create_namespace_repository(
%User{} = user,
%Namespace{} = namespace,
attrs,
idempotency_key
)
when is_map(attrs) and is_binary(idempotency_key) do
if namespace.kind == "user" and namespace.owner_user_id != user.id do
{:error, :namespace_not_allowed}
else
create_repository_transaction(
user,
namespace,
attrs,
nil,
"create",
"empty",
idempotency_key
)
end
end
def create_user_repository(%User{} = user, attrs, idempotency_key)
when is_map(attrs) and is_binary(idempotency_key) do
with {:ok, namespace} <- ensure_user_namespace(user) do
create_repository_transaction(
user,
namespace,
attrs,
nil,
"create",
"empty",
idempotency_key
)
end
end
def create_organization_repository(
%User{} = user,
%Namespace{kind: "organization"} = namespace,
attrs,
idempotency_key
)
when is_map(attrs) and is_binary(idempotency_key) do
create_repository_transaction(
user,
namespace,
attrs,
nil,
"create",
"empty",
idempotency_key
)
end
def create_user_import(%User{} = user, source, attrs, idempotency_key)
when is_map(source) and is_map(attrs) and is_binary(idempotency_key) do
with {:ok, namespace} <- ensure_user_namespace(user),
true <-
fetch_attr!(source, :source_owner_id) == user.github_id or
{:error, :source_namespace_mismatch} do
create_repository_transaction(
user,
namespace,
attrs,
source,
"github_import",
"github_import",
idempotency_key
)
end
end
def create_organization_import(
%User{} = user,
%Namespace{kind: "organization"} = namespace,
source,
attrs,
idempotency_key
)
when is_map(source) and is_map(attrs) and is_binary(idempotency_key) do
if fetch_attr!(source, :source_owner_id) == namespace.provider_account_id do
create_repository_transaction(
user,
namespace,
attrs,
source,
"github_import",
"github_import",
idempotency_key
)
else
{:error, :source_namespace_mismatch}
end
end
@doc """
Bring in a public repository this account does not own, as an upstream
mirror in the account's own namespace.
This is not the import gate relaxed. `create_user_import/4` still refuses a
source owned by anyone else, because an import claims the snapshot as this
account's own repository and says nothing about where it came from. A mirror
makes the opposite claim: it names the upstream on every surface that shows
it, carries the upstream's license or records that there is none, and
refuses every push, because there is no path back to the upstream and a
copy that silently diverges from the source it names is worse than one that
will not move.
The source must be public. A private repository this account can read
through its own GitHub grant is not something the forge may republish under
a mirror's provenance, so `:source_repository_not_public` refuses it before
any row exists.
"""
def create_user_mirror(%User{} = user, source, attrs, idempotency_key)
when is_map(source) and is_map(attrs) and is_binary(idempotency_key) do
with {:ok, namespace} <- ensure_user_namespace(user),
{:ok, mirror} <- mirror_provenance(source) do
create_repository_transaction(
user,
namespace,
attrs,
source,
"github_mirror",
"github_import",
idempotency_key,
mirror
)
end
end
@doc "Bring in a public repository as an upstream mirror in an organization namespace."
def create_organization_mirror(
%User{} = user,
%Namespace{kind: "organization"} = namespace,
source,
attrs,
idempotency_key
)
when is_map(source) and is_map(attrs) and is_binary(idempotency_key) do
with {:ok, mirror} <- mirror_provenance(source) do
create_repository_transaction(
user,
namespace,
attrs,
source,
"github_mirror",
"github_import",
idempotency_key,
mirror
)
end
end
@doc """
Whether this repository is an upstream mirror.
One predicate, so the Git plane, the API projection, and the page all decide
mirror-ness the same way and a fifth reader cannot invent a sixth rule.
"""
def mirror?(%Repository{} = repository), do: Repository.mirror?(repository)
defp mirror_provenance(source) do
full_name = fetch_attr!(source, :source_full_name)
license = fetch_attr(source, :source_license)
cond do
fetch_attr(source, :source_public) != true ->
{:error, :source_repository_not_public}
not is_binary(full_name) ->
{:error, :invalid_import}
true ->
{:ok, {"https://github.com/" <> full_name, normalize_license(license)}}
end
end
# An upstream with no license is not an upstream whose license is unknown.
# GitHub answers `null` for a repository with no license file and
# `NOASSERTION` for one whose license it cannot identify; both become the
# literal "none", which the surfaces render as a statement rather than a
# blank.
defp normalize_license(license)
when is_binary(license) and license != "" and license != "NOASSERTION",
do: String.slice(license, 0, 60)
defp normalize_license(_license), do: "none"
defp repository_creation_changeset(repository, attrs, namespace, user_id, kind, nil),
do: Repository.creation_changeset(repository, attrs, namespace, user_id, kind)
defp repository_creation_changeset(
repository,
attrs,
namespace,
user_id,
_kind,
{upstream_url, license}
),
do:
Repository.mirror_creation_changeset(
repository,
attrs,
namespace,
user_id,
upstream_url,
license
)
def list_visible_repositories(%User{} = user) do
Repo.all(
from repository in readable_by(Repository, user),
join: namespace in assoc(repository, :namespace),
order_by: [asc: namespace.slug_key, asc: repository.name_key, asc: repository.id],
preload: [namespace: namespace]
)
end
@doc """
Whether `user` can read any repository at all.
What a workspace-wide list asks to tell an empty page from an empty
workspace: "no open issues" and "no repositories yet" want different words
and different next steps, and only this distinguishes them.
"""
def any_visible_repository?(user), do: Repo.exists?(readable_by(Repository, user))
@doc "Delete a repository owned by `user`, including its durable and node-local Git data."
def delete_owned_repository(owner, name, %User{} = user, options \\ [])
when is_binary(owner) and is_binary(name) do
with %Repository{} = repository <- owned_repository(owner, name, user) do
result =
:global.trans({{:forge_push, repository.storage_key}, self()}, fn ->
delete_owned_repository_locked(repository, user)
end)
case result do
{:ok, deleted} = success ->
broadcast_repository_change(deleted.id)
Analytics.capture("repository_deleted", Analytics.distinct_id(user), %{
"repository_id" => deleted.id,
"provisioning_kind" => deleted.provisioning_kind,
"surface" => Keyword.get(options, :surface, "api")
})
success
error ->
error
end
else
nil -> {:error, :not_found}
end
end
defp owned_repository(owner, name, %User{id: user_id}) do
owner_key = String.downcase(owner)
name_key = String.downcase(name)
Repo.one(
from repository in repository_path_query(owner_key, name_key),
join: membership in Membership,
on:
membership.repository_id == repository.id and membership.user_id == ^user_id and
membership.role == "owner"
)
end
defp delete_owned_repository_locked(repository, user) do
Repo.transaction(fn ->
locked_repository =
Repo.one(
from candidate in Repository,
join: membership in Membership,
on:
membership.repository_id == candidate.id and membership.user_id == ^user.id and
membership.role == "owner",
where: candidate.id == ^repository.id,
lock: "FOR UPDATE",
select: candidate
)
if is_nil(locked_repository), do: Repo.rollback(:not_found)
provisioning =
Repo.one(
from outbox in ProvisioningOutbox,
where: outbox.repository_id == ^locked_repository.id,
lock: "FOR UPDATE"
)
if provisioning && provisioning.state == "running", do: Repo.rollback(:repository_busy)
with :ok <- WAL.delete_repo(locked_repository.storage_key),
:ok <- delete_local_caches(locked_repository.storage_key) do
Audit.record!(
"repository.deleted",
{:user, user.id},
"repository",
locked_repository.id,
repository_id: locked_repository.id,
metadata: %{
"owner" => locked_repository.owner,
"name" => locked_repository.name,
"provisioning_kind" => locked_repository.provisioning_kind
}
)
Repo.delete!(locked_repository)
else
{:error, reason} -> Repo.rollback({:storage_cleanup_failed, reason})
end
end)
end
defp delete_local_caches(storage_key) do
[node() | Node.list()]
|> Enum.uniq()
|> Task.async_stream(
fn target ->
if target == node() do
Repos.delete_repo(storage_key)
else
:erpc.call(target, Repos, :delete_repo, [storage_key], 30_000)
end
end,
ordered: false,
timeout: 31_000,
on_timeout: :kill_task,
max_concurrency: max(1, 1 + length(Node.list()))
)
|> Enum.reduce_while(:ok, fn
{:ok, :ok}, :ok -> {:cont, :ok}
{:ok, {:error, reason}}, :ok -> {:halt, {:error, reason}}
{:exit, reason}, :ok -> {:halt, {:error, reason}}
end)
end
def list_visible_repositories_page(
%User{} = user,
per_page,
after_cursor,
namespace_key \\ nil
)
when per_page in 1..100 do
query =
from repository in readable_by(Repository, user),
join: namespace in assoc(repository, :namespace),
as: :namespace,
order_by: [asc: namespace.slug_key, asc: repository.name_key, asc: repository.id],
preload: [namespace: namespace]
query =
query
|> apply_namespace_filter(namespace_key)
|> apply_repository_cursor(after_cursor)
rows =
from(row in query, limit: ^(per_page + 1))
|> Repo.all()
|> Repo.preload(@provisioning_assocs)
{Enum.take(rows, per_page), length(rows) > per_page}
end
@doc """
One repository the user may see, by id, with its provisioning receipts.
The list page's per-row counterpart: a surface that has already rendered a
row and then hears the repository changed reloads exactly that row rather
than the whole page. Returns `nil` rather than raising, because a repository
can stop being visible between the broadcast and the read.
"""
def get_visible_repository(id, user) when is_binary(id) do
visible_repository(
from repository in readable_by(Repository, user),
join: namespace in assoc(repository, :namespace),
where: repository.id == ^id,
preload: [namespace: namespace]
)
end
defp visible_repository(query) do
case Repo.one(query) do
nil -> nil
%Repository{} = repository -> Repo.preload(repository, @provisioning_assocs)
end
end
@doc """
Subscribes the caller to one repository's provisioning transitions.
DATA-001: PostgreSQL stays authoritative. The message carries the repository
id and nothing else, so a subscriber re-reads through its own visibility
predicate and can never be handed a row the database would not have given it.
"""
def subscribe_provisioning(repository_id) when is_binary(repository_id),
do: Phoenix.PubSub.subscribe(OpenAgents.PubSub, provisioning_topic(repository_id))
@doc "Subscribes a repository index to repository creation and deletion events."
def subscribe_repository_changes,
do: Phoenix.PubSub.subscribe(OpenAgents.PubSub, repository_changes_topic())
@doc "Stops the caller hearing about one repository, once it has settled."
def unsubscribe_provisioning(repository_id) when is_binary(repository_id),
do: Phoenix.PubSub.unsubscribe(OpenAgents.PubSub, provisioning_topic(repository_id))
@doc """
Announces that one repository's provisioning or import state moved.
Called after the owning transaction commits, never inside it: a subscriber
re-reads immediately, and a message sent from inside the transaction races
the commit and hands it the old row.
"""
def broadcast_provisioning(repository_id) when is_binary(repository_id) do
Phoenix.PubSub.broadcast(
OpenAgents.PubSub,
provisioning_topic(repository_id),
{:repository_provisioning, repository_id}
)
end
@doc "Announces that a repository entered or left the visible repository collection."
def broadcast_repository_change(repository_id) when is_binary(repository_id) do
Phoenix.PubSub.broadcast(
OpenAgents.PubSub,
repository_changes_topic(),
{:repository_changed, repository_id}
)
end
@doc "Records accepted Git activity and refreshes repository-list subscribers."
def record_push_activity(storage_key, occurred_at \\ DateTime.utc_now())
when is_binary(storage_key) do
case Repo.update_all(
from(repository in Repository,
where: repository.storage_key == ^storage_key,
select: repository.id
),
set: [updated_at: occurred_at]
) do
{1, [repository_id]} ->
broadcast_repository_change(repository_id)
:ok
{0, []} ->
{:error, :repository_not_found}
end
end
defp provisioning_topic(repository_id), do: "repository:" <> repository_id
defp repository_changes_topic, do: "repositories:changes"
defp apply_namespace_filter(query, nil), do: query
defp apply_namespace_filter(query, namespace_key) when is_binary(namespace_key) do
where(query, [namespace: namespace], namespace.slug_key == ^namespace_key)
end
def get_import_for_user!(id, %User{id: user_id}) do
Repo.one!(
from repository_import in RepositoryImport,
join: repository in assoc(repository_import, :repository),
join: membership in Membership,
on: membership.repository_id == repository.id and membership.user_id == ^user_id,
where: repository_import.id == ^id,
preload: [repository: {repository, [:namespace]}]
)
end
def add_member(%Repository{} = repository, %User{} = user, role \\ "contributor") do
Repo.transaction(fn ->
membership =
%Membership{}
|> Membership.changeset(%{repository_id: repository.id, user_id: user.id, role: role})
|> Repo.insert!(
on_conflict: {:replace, [:role, :updated_at]},
conflict_target: [:repository_id, :user_id],
returning: true
)
Audit.record!(
"repository.membership.updated",
{:user, user.id},
"membership",
membership_subject_id(membership),
repository_id: repository.id,
metadata: %{"role" => membership.role}
)
membership
end)
end
@doc """
Lists one repository's members with their users, owners first.
The order is role rank and then login, so the page reads the same way every
time it renders.
"""
def list_members(%Repository{id: repository_id}) do
from(membership in Membership,
join: user in assoc(membership, :user),
where: membership.repository_id == ^repository_id,
order_by: [
asc:
fragment(
"array_position(array['owner','maintainer','contributor','viewer'], ?)",
membership.role
),
asc: fragment("lower(?)", user.github_login)
],
preload: [user: user]
)
|> Repo.all()
end
@doc "The acting owner's view of one member row: add by GitHub login."
def add_member_by_login(%Repository{} = repository, %User{} = actor, login, role)
when is_binary(login) and role in @all_roles do
with_owner_memberships(repository, actor, fn _memberships ->
case active_user_by_login(login) do
%User{} = user ->
membership = upsert_membership!(repository, user, role)
audit_membership(repository, actor, user, "added", role)
membership
nil ->
Repo.rollback(:unknown_user)
end
end)
end
def change_member_role(%Repository{} = repository, %User{} = actor, user_id, role)
when role in @all_roles do
with_owner_memberships(repository, actor, fn memberships ->
user = Repo.get(User, user_id)
membership = Enum.find(memberships, &(&1.user_id == user_id))
if is_nil(user) or is_nil(membership), do: Repo.rollback(:unknown_member)
guard_last_owner!(memberships, membership, role)
updated = upsert_membership!(repository, user, role)
audit_membership(repository, actor, user, "role changed", role)
updated
end)
end
def remove_member(%Repository{} = repository, %User{} = actor, user_id) do
case with_owner_memberships(repository, actor, fn memberships ->
user = Repo.get(User, user_id)
membership = Enum.find(memberships, &(&1.user_id == user_id))
if is_nil(user) or is_nil(membership), do: Repo.rollback(:unknown_member)
guard_last_owner!(memberships, membership, nil)
Repo.delete!(membership)
Audit.record!(
"repository.membership.removed",
{:user, actor.id},
"membership",
membership_subject_id(membership),
repository_id: repository.id,
metadata: %{"login" => user.github_login}
)
:ok
end) do
{:ok, :ok} -> :ok
{:error, reason} -> {:error, reason}
end
end
# Serialize membership administration per repository. This makes both the
# owner check and the last-owner rule true at the instant of the write, even
# when an owner keeps an old LiveView open or two owners act concurrently.
defp with_owner_memberships(%Repository{id: repository_id}, %User{id: actor_id}, operation) do
Repo.transaction(fn ->
memberships =
Repo.all(
from membership in Membership,
where: membership.repository_id == ^repository_id,
order_by: [asc: membership.user_id],
lock: "FOR UPDATE"
)
active_actor? =
Repo.exists?(
from user in User,
where: user.id == ^actor_id and user.status == "active"
)
actor_owns? =
Enum.any?(memberships, &(&1.user_id == actor_id and &1.role == "owner"))
if active_actor? and actor_owns? do
operation.(memberships)
else
Repo.rollback(:forbidden)
end
end)
end
# A repository with no owner cannot be administered any more, so the last
# owner cannot be demoted or removed, including by themselves. The caller
# holds row locks over every membership in this repository.
defp guard_last_owner!(memberships, %Membership{} = membership, new_role) do
leaving_owner? = membership.role == "owner" and new_role != "owner"
owner_count = Enum.count(memberships, &(&1.role == "owner"))
if leaving_owner? and owner_count <= 1, do: Repo.rollback(:last_owner)
end
defp upsert_membership!(repository, user, role) do
%Membership{}
|> Membership.changeset(%{
repository_id: repository.id,
user_id: user.id,
role: role
})
|> Repo.insert!(
on_conflict: {:replace, [:role, :updated_at]},
conflict_target: [:repository_id, :user_id],
returning: true
)
end
defp audit_membership(repository, actor, subject, action, role) do
Audit.record!(
"repository.membership.updated",
{:user, actor.id},
"membership",
"#{repository.id}:#{subject.id}",
repository_id: repository.id,
metadata: %{"action" => action, "role" => role, "login" => subject.github_login}
)
end
defp active_user_by_login(login) do
Repo.one(
from user in User,
where:
user.status == "active" and
fragment("lower(?)", user.github_login) == ^String.downcase(String.trim(login))
)
end
@doc """
Repositories this account administers, and can therefore hand to one of its
computers.
This is a membership question, not a visibility one: it starts from the
actor's own `repository_memberships` rows in the two roles that may
administer a repository, so it neither composes nor restates
`readable_by/2` (REPOSITORY-001).
"""
@spec list_grantable_repositories(User.t() | nil) :: [Repository.t()]
def list_grantable_repositories(%User{id: user_id}) do
Repo.all(
from repository in Repository,
join: membership in Membership,
on: membership.repository_id == repository.id,
where: membership.user_id == ^user_id and membership.role in @machine_grant_roles,
order_by: [asc: repository.owner, asc: repository.name]
)
end
def list_grantable_repositories(_actor), do: []
@doc """
The repository grants one of this account's computers holds.
A computer this account does not own is indistinguishable from one that holds
nothing.
"""
@spec list_machine_grants(User.t(), String.t()) :: [MachineGrant.t()]
def list_machine_grants(%User{} = actor, machine_id) when is_binary(machine_id) do
case owned_machine(actor, machine_id) do
{:ok, machine} ->
Repo.all(
from grant in MachineGrant,
join: repository in Repository,
on: repository.id == grant.repository_id,
where: grant.machine_id == ^machine.id,
order_by: [asc: repository.owner, asc: repository.name],
preload: [repository: repository]
)
{:error, _reason} ->
[]
end
end
def list_machine_grants(_actor, _machine_id), do: []
@doc """
Grant one of this account's computers access to one repository it
administers.
This is the entry point the `/computers` surface calls, and the only one:
`grant_machine/4` below is private, so a caller cannot reach the write
without passing an acting account for both halves. The grant was previously
unreachable — the table was read by `OpenAgents.Forge.GitHTTP` on every
request a paired computer made and written by nothing outside tests, so every
such request answered `404 unknown repository`.
Both halves are resolved from the acting account rather than from the
caller's identifiers: a computer another account owns, a repository this
account does not administer, and an identifier that names nothing are all the
same refusal.
"""
@spec grant_machine_access(User.t(), String.t(), String.t(), [String.t()]) ::
{:ok, MachineGrant.t()} | {:error, atom() | Ecto.Changeset.t()}
def grant_machine_access(%User{} = actor, machine_id, repository_id, operations)
when is_binary(machine_id) and is_binary(repository_id) and is_list(operations) do
with {:ok, machine} <- owned_machine(actor, machine_id),
{:ok, repository} <- administered_repository(actor, repository_id) do
grant_machine(repository, actor, machine, operations)
end
end
def grant_machine_access(_actor, _machine_id, _repository_id, _operations),
do: {:error, :machine_not_owned}
@doc """
Withdraw a computer's access to one repository.
A grant outlives the delegation that needed it, so the surface that creates
one owes a way to take it back. Revoking the computer itself already ends the
grant's effect — `machine_access?/3` joins the computer's status and token
expiry — but that is all or nothing, and this is not.
The audit record carries the account that created the grant, which is the
only place that survives the row.
"""
@spec revoke_machine_access(User.t(), String.t(), String.t()) ::
{:ok, MachineGrant.t()} | {:error, atom()}
def revoke_machine_access(%User{} = actor, machine_id, repository_id)
when is_binary(machine_id) and is_binary(repository_id) do
with {:ok, machine} <- owned_machine(actor, machine_id),
{:ok, repository} <- administered_repository(actor, repository_id),
%MachineGrant{} = grant <-
Repo.get_by(MachineGrant, repository_id: repository.id, machine_id: machine.id) do
Repo.transaction(fn ->
Repo.delete!(grant)
Audit.record!(
"repository.machine_grant.revoked",
{:user, actor.id},
"machine_grant",
grant.id,
repository_id: repository.id,
metadata: %{
"machine_id" => machine.id,
"operations" => grant.operations,
"granted_by_user_id" => grant.created_by_user_id
}
)
grant
end)
else
nil -> {:error, :grant_not_found}
{:error, reason} -> {:error, reason}
end
end
def revoke_machine_access(_actor, _machine_id, _repository_id),
do: {:error, :machine_not_owned}
defp owned_machine(%User{id: user_id}, machine_id) do
case OpenAgents.Machines.get_machine(user_id, machine_id) do
{:ok, machine} -> {:ok, machine}
{:error, _reason} -> {:error, :machine_not_owned}
end
end
defp administered_repository(%User{} = actor, repository_id) do
with {:ok, id} <- Ecto.UUID.cast(repository_id),
%Repository{} = repository <- Repo.get(Repository, id),
role when role in @machine_grant_roles <- membership_role(repository, actor) do
{:ok, repository}
else
_absent_or_foreign -> {:error, :repository_not_allowed}
end
end
defp grant_machine(
%Repository{} = repository,
%User{} = actor,
%Machine{} = machine,
operations
)
when is_list(operations) do
with true <- machine.user_id == actor.id or {:error, :machine_not_owned},
role when role in @machine_grant_roles <- membership_role(repository, actor) do
Repo.transaction(fn ->
grant =
%MachineGrant{}
|> MachineGrant.changeset(%{
repository_id: repository.id,
machine_id: machine.id,
created_by_user_id: actor.id,
operations: operations |> Enum.uniq() |> Enum.sort()
})
|> Repo.insert!(
on_conflict: {:replace, [:operations, :created_by_user_id, :updated_at]},
conflict_target: [:repository_id, :machine_id],
returning: true
)
Audit.record!(
"repository.machine_grant.updated",
{:user, actor.id},
"machine_grant",
grant.id,
repository_id: repository.id,
metadata: %{"machine_id" => machine.id, "operations" => grant.operations}
)
grant
end)
else
nil -> {:error, :repository_not_allowed}
false -> {:error, :machine_not_owned}
{:error, reason} -> {:error, reason}
_role -> {:error, :repository_not_allowed}
end
rescue
error in Ecto.InvalidChangesetError -> {:error, error.changeset}
end
def machine_access?(%Repository{id: repository_id}, machine_id, operation)
when operation in ~w(read write) and is_binary(machine_id) do
now = DateTime.utc_now()
Repo.exists?(
from grant in MachineGrant,
join: machine in Machine,
on:
machine.id == grant.machine_id and machine.status == "active" and
machine.token_expires_at > ^now,
where:
grant.repository_id == ^repository_id and grant.machine_id == ^machine_id and
fragment("? = ANY(?)", ^operation, grant.operations)
)
end
def machine_access?(%Repository{}, _machine_id, _operation), do: false
def writable?(%Repository{id: repository_id}, %User{id: user_id}) do
Repo.exists?(
from membership in Membership,
join: user in User,
on: user.id == membership.user_id and user.status == "active",
where:
membership.repository_id == ^repository_id and membership.user_id == ^user_id and
membership.role in ^@writable_roles
)
end
def writable?(%Repository{}, nil), do: false
def membership_role(%Repository{id: repository_id}, %User{id: user_id}) do
Repo.one(
from membership in Membership,
where: membership.repository_id == ^repository_id and membership.user_id == ^user_id,
select: membership.role
)
end
def membership_role(%Repository{}, nil), do: nil
@doc "Whether the repository is publicly readable."
def public?(%Repository{visibility: "public"}), do: true
def public?(%Repository{}), do: false
@doc """
Whether the user holds any membership role, including read-only `viewer`.
"""
def member?(%Repository{id: repository_id}, %User{id: user_id}) do
Repo.exists?(
from membership in Membership,
join: user in User,
on: user.id == membership.user_id and user.status == "active",
where: membership.repository_id == ^repository_id and membership.user_id == ^user_id
)
end
def member?(%Repository{}, nil), do: false
@doc """
Whether the user may take part in issue conversations: open issues and
comment.
GitHub's model, which is ours: an active signed-in person can join the
conversation on any public repository without membership; a private
repository admits its own members. Triage writes (label, assign, close,
edit) stay behind writability and are not governed by this predicate.
"""
def issue_participant?(%Repository{}, nil), do: false
def issue_participant?(%Repository{} = repository, %User{} = user) do
active_user?(user) and (public?(repository) or member?(repository, user))
end
def issue_participant?(%Repository{} = repository, %Agent{} = agent) do
active_agent?(agent) and public?(repository)
end
@doc "Whether the user holds the repository's `owner` role."
def owner?(%Repository{} = repository, %User{} = user) do
active_user?(user) and membership_role(repository, user) == "owner"
end
def owner?(%Repository{}, nil), do: false
defp active_user?(%User{id: user_id}) do
Repo.exists?(from user in User, where: user.id == ^user_id and user.status == "active")
end
defp active_agent?(%Agent{id: agent_id}) do
Repo.exists?(from agent in Agent, where: agent.id == ^agent_id and agent.status == "active")
end
@doc "Subscribes the caller to one repository's issue activity."
@all_issues_topic "issues:all"
def subscribe_issues(repository_id),
do: Phoenix.PubSub.subscribe(OpenAgents.PubSub, issues_topic(repository_id))
def unsubscribe_issues(repository_id),
do: Phoenix.PubSub.unsubscribe(OpenAgents.PubSub, issues_topic(repository_id))
@doc """
Announces that one repository's issues moved.
Called after the owning transaction commits. The message carries the
repository id and nothing else, so every subscriber re-reads through its own
visibility and authorization predicates.
"""
def broadcast_issues(repository_id) do
Phoenix.PubSub.broadcast(
OpenAgents.PubSub,
issues_topic(repository_id),
{:issues_changed, repository_id}
)
Phoenix.PubSub.broadcast(
OpenAgents.PubSub,
@all_issues_topic,
{:issues_changed, repository_id}
)
end
defp issues_topic(repository_id), do: "issues:" <> repository_id
@doc "Receives `{:issues_changed, repository_id}` for every repository at once."
def subscribe_all_issues,
do: Phoenix.PubSub.subscribe(OpenAgents.PubSub, @all_issues_topic)
@doc "Subscribes the caller to one repository's project activity."
@all_projects_topic "projects:all"
def subscribe_projects(repository_id),
do: Phoenix.PubSub.subscribe(OpenAgents.PubSub, projects_topic(repository_id))
def unsubscribe_projects(repository_id),
do: Phoenix.PubSub.unsubscribe(OpenAgents.PubSub, projects_topic(repository_id))
@doc """
Announces that one repository's projects moved.
Called after the owning transaction commits. The message carries the
repository id and nothing else, so every subscriber re-reads through its own
visibility and authorization predicates.
"""
def broadcast_projects(repository_id) do
Phoenix.PubSub.broadcast(
OpenAgents.PubSub,
projects_topic(repository_id),
{:projects_changed, repository_id}
)
Phoenix.PubSub.broadcast(
OpenAgents.PubSub,
@all_projects_topic,
{:projects_changed, repository_id}
)
end
defp projects_topic(repository_id), do: "projects:" <> repository_id
@doc "Receives `{:projects_changed, repository_id}` for every repository at once."
def subscribe_all_projects,
do: Phoenix.PubSub.subscribe(OpenAgents.PubSub, @all_projects_topic)
@doc "Seeds GitHub's default label vocabulary onto a new or imported repository."
def seed_default_labels!(%Repository{} = repository) do
Enum.each(@default_labels, fn {name, color, description} ->
Repo.insert!(
%OpenAgents.Labels.Label{}
|> OpenAgents.Labels.Label.changeset(%{
name: name,
color: color,
description: description,
repository_id: repository.id
}),
on_conflict: :nothing,
conflict_target: [:repository_id, :name]
)
end)
:ok
end
def list_assignable_users(%Repository{id: repository_id}) do
Repo.all(
from user in User,
join: membership in Membership,
on: membership.user_id == user.id,
where:
membership.repository_id == ^repository_id and
membership.role in ^@writable_roles and user.status == "active",
order_by: [asc: fragment("lower(?)", user.github_login)]
)
end
def get_assignable_user_by_login!(%Repository{id: repository_id}, login)
when is_binary(login) do
Repo.one!(
from user in User,
join: membership in Membership,
on: membership.user_id == user.id,
where:
membership.repository_id == ^repository_id and
membership.role in ^@writable_roles and user.status == "active" and
fragment("lower(?)", user.github_login) == ^String.downcase(URI.decode(login))
)
end
defp create_repository_transaction(
user,
namespace,
attrs,
source,
operation,
provisioning_kind,
idempotency_key,
mirror \\ nil
) do
normalized_request = %{
namespace_id: namespace.id,
repository: normalize_repository_attrs(attrs),
source: normalize_source(source),
mirror: mirror
}
request_digest = digest(normalized_request)
result =
Repo.transaction(fn ->
case get_idempotency_request(user.id, operation, idempotency_key) do
%IdempotencyRequest{request_digest: ^request_digest} = request ->
replay_result(request, operation)
%IdempotencyRequest{} ->
Repo.rollback(:idempotency_conflict)
nil ->
create_repository_rows!(
user,
namespace,
normalized_request.repository,
normalized_request.source,
operation,
provisioning_kind,
idempotency_key,
request_digest,
normalized_request.mirror
)
end
end)
case result do
{:ok, {:ok, %Repository{} = repository, :created} = value} ->
broadcast_repository_change(repository.id)
value
{:ok, {:ok, %Repository{} = repository, %RepositoryImport{}, :created} = value} ->
broadcast_repository_change(repository.id)
value
{:ok, value} ->
value
{:error, reason} ->
{:error, reason}
end
rescue
error in Ecto.InvalidChangesetError -> {:error, error.changeset}
end
defp create_repository_rows!(
user,
namespace,
attrs,
source,
operation,
provisioning_kind,
idempotency_key,
request_digest,
mirror
) do
lock_and_validate_quota!(namespace.id)
repository =
%Repository{}
|> repository_creation_changeset(attrs, namespace, user.id, provisioning_kind, mirror)
|> Repo.insert!()
Audit.record!("repository.created", {:user, user.id}, "repository", repository.id,
repository_id: repository.id,
metadata: %{
"namespace_id" => namespace.id,
"provisioning_kind" => provisioning_kind,
"visibility" => repository.visibility
}
)
membership =
%Membership{}
|> Membership.changeset(%{repository_id: repository.id, user_id: user.id, role: "owner"})
|> Repo.insert!()
Audit.record!(
"repository.membership.created",
{:user, user.id},
"membership",
membership_subject_id(membership),
repository_id: repository.id,
metadata: %{"role" => "owner"}
)
seed_default_labels!(repository)
repository_import =
if source do
created_import =
%RepositoryImport{}
|> RepositoryImport.changeset(repository.id, source)
|> Repo.insert!()
Audit.record!(
"repository.import.created",
{:user, user.id},
"repository_import",
created_import.id,
repository_id: repository.id,
metadata: %{
"provider" => created_import.provider,
"source_repository_id" => created_import.source_repository_id
}
)
created_import
end
# The outbox names the executor, not the caller's intent. A mirror and an
# import are copied by the same provisioning path, so both queue
# `github_import` here while the idempotency record keeps them apart: one
# key must never replay an import as a mirror or the reverse.
outbox =
%ProvisioningOutbox{}
|> ProvisioningOutbox.changeset(
repository.id,
repository_import && repository_import.id,
outbox_operation(provisioning_kind)
)
|> Repo.insert!()
Audit.record!(
"repository.provisioning.pending",
{:user, user.id},
"provisioning_outbox",
outbox.id,
repository_id: repository.id,
metadata: %{"operation" => operation}
)
%IdempotencyRequest{}
|> IdempotencyRequest.changeset(
user.id,
operation,
idempotency_key,
request_digest,
%{
repository_id: repository.id,
repository_import_id: repository_import && repository_import.id
}
)
|> Repo.insert!()
repository = Repo.preload(repository, [:namespace, :memberships, :repository_import])
if repository_import do
{:ok, repository, repository.repository_import, :created}
else
{:ok, repository, :created}
end
end
defp outbox_operation("empty"), do: "create"
defp outbox_operation("github_import"), do: "github_import"
defp replay_result(request, operation) when operation in ["github_import", "github_mirror"] do
repository =
Repository
|> Repo.get!(request.repository_id)
|> Repo.preload([:namespace, :memberships, :repository_import])
{:ok, repository, repository.repository_import, :replayed}
end
defp replay_result(request, _operation) do
repository =
Repository
|> Repo.get!(request.repository_id)
|> Repo.preload([:namespace, :memberships, :repository_import])
{:ok, repository, :replayed}
end
defp get_idempotency_request(user_id, operation, idempotency_key) do
Repo.one(
from request in IdempotencyRequest,
where:
request.user_id == ^user_id and request.operation == ^operation and
request.idempotency_key == ^idempotency_key,
lock: "FOR UPDATE"
)
end
defp maybe_update_namespace!(namespace, attrs) do
next_slug = fetch_attr!(attrs, :slug)
if String.downcase(next_slug) != namespace.slug_key do
assert_namespace_slug_available!(next_slug, namespace.id)
%NamespaceAlias{}
|> NamespaceAlias.changeset(namespace.id, namespace.slug)
|> Repo.insert!(on_conflict: :nothing, conflict_target: [:slug_key])
end
namespace
|> Namespace.changeset(attrs)
|> Repo.update!()
end
# Legacy domain tests create already-ready repositories without a GitHub
# principal. Keep that fixture seam separate from the public creation API,
# which always projects an immutable GitHub account ID.
defp ensure_legacy_namespace!(owner) do
slug_key = String.downcase(owner)
case Repo.get_by(Namespace, slug_key: slug_key, state: "active") do
%Namespace{} = namespace ->
namespace
nil ->
provider_account_id = 8_000_000_000 + :erlang.phash2(slug_key, 1_000_000_000)
%Namespace{}
|> Namespace.changeset(%{
provider_account_id: provider_account_id,
slug: owner,
kind: "organization",
provider_refreshed_at: DateTime.utc_now()
})
|> Repo.insert!()
end
end
defp assert_namespace_slug_available!(slug, namespace_id) do
slug_key = String.downcase(slug)
collision? =
Repo.exists?(
from namespace in Namespace,
where:
namespace.slug_key == ^slug_key and namespace.state == "active" and
namespace.id != ^namespace_id
) or
Repo.exists?(
from namespace_alias in NamespaceAlias,
where:
namespace_alias.slug_key == ^slug_key and
namespace_alias.namespace_id != ^namespace_id
)
if collision?, do: Repo.rollback(:namespace_slug_conflict), else: :ok
end
defp lock_and_validate_quota!(namespace_id) do
_namespace =
Repo.one!(
from namespace in Namespace, where: namespace.id == ^namespace_id, lock: "FOR UPDATE"
)
repository_count =
Repo.aggregate(
from(repository in Repository, where: repository.namespace_id == ^namespace_id),
:count
)
if repository_count >= repository_namespace_limit(),
do: Repo.rollback(:repository_quota_exceeded)
end
defp repository_namespace_limit do
case Application.get_env(
:openagents,
:repository_namespace_limit,
@repository_namespace_limit
) do
limit when is_integer(limit) and limit > 0 -> limit
_invalid -> @repository_namespace_limit
end
end
# Resolution only: which row the caller-supplied path names, never who may
# read it. Every caller composes `readable_by/2` on the result.
defp visible_path_query(owner, name) when is_binary(owner) and is_binary(name) do
repository_path_query(String.downcase(owner), String.downcase(name))
end
defp repository_path_query(owner_key, name_key) do
from repository in Repository,
join: namespace in assoc(repository, :namespace),
left_join: namespace_alias in NamespaceAlias,
on: namespace_alias.namespace_id == namespace.id and namespace_alias.slug_key == ^owner_key,
where:
repository.name_key == ^name_key and namespace.state == "active" and
(namespace.slug_key == ^owner_key or not is_nil(namespace_alias.id)),
distinct: repository.id,
preload: [namespace: namespace]
end
defp apply_repository_cursor(query, nil), do: query
defp apply_repository_cursor(query, {owner_key, name_key, id}) do
where(
query,
[repository, namespace: namespace],
namespace.slug_key > ^owner_key or
(namespace.slug_key == ^owner_key and repository.name_key > ^name_key) or
(namespace.slug_key == ^owner_key and repository.name_key == ^name_key and
repository.id > ^id)
)
end
defp normalize_repository_attrs(attrs) do
%{
name: fetch_attr!(attrs, :name),
description: fetch_attr(attrs, :description),
visibility: fetch_attr(attrs, :visibility) || "private",
default_branch: fetch_attr(attrs, :default_branch) || "main"
}
end
defp normalize_source(nil), do: nil
defp normalize_source(source) do
%{
provider: fetch_attr(source, :provider) || "github",
source_repository_id: fetch_attr!(source, :source_repository_id),
source_owner_id: fetch_attr!(source, :source_owner_id),
source_full_name: fetch_attr!(source, :source_full_name),
source_default_branch: fetch_attr!(source, :source_default_branch),
source_ref_digest: fetch_attr!(source, :source_ref_digest),
source_head_sha: fetch_attr(source, :source_head_sha),
source_refs: fetch_attr!(source, :source_refs),
source_uses_lfs: fetch_attr(source, :source_uses_lfs) || false
}
end
defp digest(value) do
value
|> :erlang.term_to_binary([:deterministic])
|> then(&:crypto.hash(:sha256, &1))
|> Base.encode16(case: :lower)
end
defp fetch_attr(attrs, key) do
Map.get(attrs, key, Map.get(attrs, Atom.to_string(key)))
end
defp fetch_attr!(attrs, key) do
case fetch_attr(attrs, key) do
nil -> raise ArgumentError, "missing #{key}"
value -> value
end
end
defp membership_subject_id(membership),
do: "#{membership.repository_id}:#{membership.user_id}"
defp unwrap_transaction({:ok, value}), do: {:ok, value}
defp unwrap_transaction({:error, reason}), do: {:error, reason}
end