tarakan/lib/tarakan_web/git_http/service.ex
Maxfield Luke af6077b9c3
All checks were successful
CI and deploy / Test (push) Successful in 4m52s
CI and deploy / Deploy production (push) Successful in 23s
Initial commit on Forgejo
Fresh repository history for elektrine/tarakan hosted at
https://git.elektrine.com/elektrine/tarakan.
2026-07-29 04:43:40 -04:00

398 lines
12 KiB
Elixir

defmodule TarakanWeb.GitHTTP.Service do
@moduledoc """
The smart-HTTP protocol core: ref advertisement and the stateless RPCs.
RPC bodies stream to a `git upload-pack`/`git receive-pack` subprocess
through an Erlang port and the response streams back chunked. Ports cannot
half-close stdin, but that is fine here: stateless-rpc requests are
self-delimiting (an upload-pack request ends with a flush/`done` pkt; a
receive-pack pack stream carries its own object count and trailing hash),
so git writes its response and exits without waiting for EOF. A hard
wall-clock deadline closes the port if a malformed request ever leaves git
waiting on stdin.
"""
import Plug.Conn
alias Tarakan.Git.Concurrency
alias Tarakan.Git.Local
alias Tarakan.HostedRepositories.Storage
require Logger
@read_chunk_bytes 65_536
# An upload-pack request is pure negotiation (wants/haves) and stays tiny;
# only receive-pack carries a pack file.
@max_upload_pack_request_bytes 10 * 1_024 * 1_024
@upload_pack_timeout_ms :timer.seconds(60)
@receive_pack_timeout_ms :timer.seconds(300)
@doc "GET /:owner/:name.git/info/refs?service=…"
def advertise_refs(conn, repository, service) do
case Concurrency.checkout() do
:ok ->
try do
do_advertise_refs(conn, repository, service)
after
Concurrency.checkin()
end
{:error, :busy} ->
conn
|> put_resp_content_type("text/plain")
|> put_resp_header("retry-after", "5")
|> send_resp(503, "git service busy; try again shortly")
|> halt()
end
end
defp do_advertise_refs(conn, repository, service) do
dir = Storage.dir(repository)
case Local.run(dir, [service_subcommand(service), "--stateless-rpc", "--advertise-refs", "."],
extra_env: protocol_env(conn)
) do
{:ok, output} ->
body = [pkt_line("# service=#{service}\n"), "0000", output]
# Exact content type, no charset parameter: git validates this
# header to distinguish smart from dumb HTTP.
conn
|> put_resp_header("content-type", "application/x-#{service}-advertisement")
|> put_resp_header("cache-control", "no-cache")
|> send_resp(200, body)
|> halt()
{:error, reason} ->
Logger.warning(
"git advertise-refs failed for repository #{repository.id}: #{inspect(reason)}"
)
conn
|> put_resp_content_type("text/plain")
|> send_resp(500, "unavailable")
|> halt()
end
end
@doc "POST /:owner/:name.git/git-upload-pack | git-receive-pack"
def rpc(conn, repository, service, scope) do
with :ok <- validate_content_type(conn, service),
{:ok, decoder} <- body_decoder(conn),
:ok <- Concurrency.checkout() do
try do
run_rpc(conn, repository, service, scope, decoder)
after
Concurrency.checkin()
end
else
{:error, :busy} ->
conn
|> put_resp_content_type("text/plain")
|> put_resp_header("retry-after", "5")
|> send_resp(503, "git service busy; try again shortly")
|> halt()
{:error, status, message} ->
conn
|> put_resp_content_type("text/plain")
|> send_resp(status, message)
|> halt()
end
end
defp run_rpc(conn, repository, service, scope, decoder) do
dir = Storage.dir(repository)
deadline = System.monotonic_time(:millisecond) + timeout_ms(service)
port =
Port.open(
{:spawn_executable, git_executable()},
[
:binary,
:exit_status,
:use_stdio,
:hide,
args: [service_subcommand(service), "--stateless-rpc", dir],
env: port_env(conn)
]
)
try do
case forward_request_body(conn, port, decoder, request_cap(service), deadline) do
{:ok, conn} ->
conn =
conn
|> put_resp_header("content-type", "application/x-#{service}-result")
|> put_resp_header("cache-control", "no-cache")
|> send_chunked(200)
case stream_response(conn, port, deadline) do
{:ok, conn, exit_status} ->
if exit_status == 0 and service == "git-receive-pack" do
after_receive_pack(repository, scope)
end
halt(conn)
{:client_closed, conn} ->
halt(conn)
end
{:error, status, message} ->
conn
|> put_resp_content_type("text/plain")
|> send_resp(status, message)
|> halt()
end
after
close_port(port)
end
end
# ── Request streaming ────────────────────────────────────────────────
defp forward_request_body(conn, port, decoder, cap_bytes, deadline) do
if System.monotonic_time(:millisecond) > deadline do
{:error, 408, "request timed out"}
else
case read_body(conn, length: @read_chunk_bytes, read_length: @read_chunk_bytes) do
{tag, data, conn} when tag in [:ok, :more] ->
case decode_and_forward(decoder, data, port, cap_bytes) do
{:ok, decoder} ->
case tag do
:ok -> {:ok, conn}
:more -> forward_request_body(conn, port, decoder, cap_bytes, deadline)
end
# git rejected the request and exited before reading all of it -
# normal for a malformed body. Its diagnostic is already on stdout
# as in-band pkt-lines, so stop forwarding and relay what it wrote
# rather than inventing a transport error.
{:closed, _decoder} ->
{:ok, conn}
{:error, :request_too_large} ->
{:error, 413, "request too large"}
{:error, :invalid_encoding} ->
{:error, 400, "invalid request encoding"}
end
{:error, _reason} ->
{:error, 400, "invalid request body"}
end
end
end
# The decoder tracks inflated bytes so a gzip bomb trips the cap exactly like
# an oversized plain body.
defp body_decoder(conn) do
case get_req_header(conn, "content-encoding") do
[] ->
{:ok, %{zstream: nil, total: 0, finished?: false}}
["gzip"] ->
zstream = :zlib.open()
# 31 = 15 window bits + 16 for the gzip header format.
:ok = :zlib.inflateInit(zstream, 31)
{:ok, %{zstream: zstream, total: 0, finished?: false}}
_unsupported ->
{:error, 415, "unsupported content encoding"}
end
end
defp decode_and_forward(%{zstream: nil} = decoder, data, port, cap_bytes) do
emit(decoder, data, port, cap_bytes)
end
# Compressed bodies inflate in bounded steps, and each step is measured and
# forwarded before the next one runs. Inflating a whole read chunk first
# would let a ~64KB chunk expand past 60MB before anything checked the cap -
# roughly a thousandfold memory amplification, reachable without credentials
# because anonymous clients may clone listed repositories.
defp decode_and_forward(%{finished?: true}, data, _port, _cap_bytes) when data != "" do
{:error, :invalid_encoding}
end
defp decode_and_forward(%{zstream: zstream} = decoder, data, port, cap_bytes) do
inflate_steps(zstream, data, decoder, port, cap_bytes)
end
defp inflate_steps(zstream, input, decoder, port, cap_bytes) do
case safe_inflate(zstream, input) do
{:ok, :continue, output} ->
case emit(decoder, output, port, cap_bytes) do
{:ok, decoder} -> inflate_steps(zstream, [], decoder, port, cap_bytes)
other -> other
end
{:ok, :finished, output} ->
with {:ok, decoder} <- emit(decoder, output, port, cap_bytes) do
{:ok, %{decoder | finished?: true}}
end
{:error, :invalid_encoding} = error ->
error
end
end
# Only the zlib call is guarded: a decompression failure means the client
# sent a bad body, while a failure writing to the port means the subprocess
# went away. Conflating them reports one as the other.
defp safe_inflate(zstream, input) do
case :zlib.safeInflate(zstream, input) do
{:continue, output} -> {:ok, :continue, output}
{:finished, output} -> {:ok, :finished, output}
end
rescue
_error -> {:error, :invalid_encoding}
end
# The cap is checked before the write, so an oversized body is refused
# without the excess ever reaching git.
defp emit(decoder, output, port, cap_bytes) do
binary = IO.iodata_to_binary(output)
total = decoder.total + byte_size(binary)
cond do
total > cap_bytes -> {:error, :request_too_large}
binary == "" -> {:ok, %{decoder | total: total}}
true -> write_to_port(port, binary, %{decoder | total: total})
end
end
defp write_to_port(port, binary, decoder) do
Port.command(port, binary)
{:ok, decoder}
rescue
ArgumentError -> {:closed, decoder}
end
# ── Response streaming ───────────────────────────────────────────────
defp stream_response(conn, port, deadline) do
remaining = max(deadline - System.monotonic_time(:millisecond), 0)
receive do
{^port, {:data, data}} ->
case chunk(conn, data) do
{:ok, conn} -> stream_response(conn, port, deadline)
{:error, :closed} -> {:client_closed, conn}
end
{^port, {:exit_status, status}} ->
if status != 0 do
Logger.warning("git rpc subprocess exited with status #{status}")
end
{:ok, conn, status}
after
remaining ->
Logger.warning("git rpc timed out; closing subprocess")
{:client_closed, conn}
end
end
# ── Post-receive (best-effort, never blocks the protocol response) ──
defp after_receive_pack(repository, scope) do
if synchronous_post_receive?() do
Tarakan.HostedRepositories.PostReceive.run(repository, scope)
else
Task.Supervisor.start_child(Tarakan.TaskSupervisor, fn ->
Tarakan.HostedRepositories.PostReceive.run(repository, scope)
end)
end
:ok
rescue
_error -> :ok
end
defp synchronous_post_receive? do
:tarakan
|> Application.get_env(__MODULE__, [])
|> Keyword.get(:synchronous_post_receive, false)
end
# ── Plumbing ─────────────────────────────────────────────────────────
defp validate_content_type(conn, service) do
expected = "application/x-#{service}-request"
case get_req_header(conn, "content-type") do
[^expected | _rest] ->
:ok
[content_type | _rest] ->
if String.starts_with?(content_type, expected <> ";") do
:ok
else
{:error, 415, "expected #{expected}"}
end
[] ->
{:error, 415, "expected #{expected}"}
end
end
defp request_cap("git-upload-pack"),
do: config(:max_upload_pack_request_bytes, @max_upload_pack_request_bytes)
defp request_cap("git-receive-pack"), do: Storage.max_push_bytes()
defp timeout_ms("git-upload-pack"), do: config(:upload_pack_timeout_ms, @upload_pack_timeout_ms)
defp timeout_ms("git-receive-pack"),
do: config(:receive_pack_timeout_ms, @receive_pack_timeout_ms)
defp config(key, default) do
:tarakan
|> Application.get_env(__MODULE__, [])
|> Keyword.get(key, default)
end
defp service_subcommand("git-upload-pack"), do: "upload-pack"
defp service_subcommand("git-receive-pack"), do: "receive-pack"
defp git_executable do
System.find_executable("git") || raise "git executable not found"
end
# Only well-formed protocol version advertisements; never pass arbitrary
# header bytes into the subprocess environment.
@git_protocol ~r/\Aversion=[0-9]+(?:[:\w.=,-]*)\z/
defp protocol_env(conn) do
case get_req_header(conn, "git-protocol") do
[value | _rest] when byte_size(value) <= 200 ->
if Regex.match?(@git_protocol, value), do: [{"GIT_PROTOCOL", value}], else: []
_other ->
[]
end
end
defp port_env(conn) do
for {key, value} <- Local.env() ++ protocol_env(conn) do
{String.to_charlist(key), String.to_charlist(value)}
end
end
defp pkt_line(data) do
size = byte_size(data) + 4
[String.pad_leading(Integer.to_string(size, 16), 4, "0") |> String.downcase(), data]
end
defp close_port(port) do
if is_port(port) and port in Port.list() do
Port.close(port)
end
:ok
rescue
ArgumentError -> :ok
end
end