Files
mustafa.ulukaya ef26860df9 feat(logging): unified request/response log with configurable retention (v1.11.0)
Until now the only record of what happened was `user_activity_logs`, which
stores non-GET 2xx operations with no bodies. When something failed you could
see that a counter went up, never what was sent or what came back.

This adds one queryable timeline covering both directions:

- inbound: every API call, including GETs and including 4xx/5xx, with the
  user, client IP, status, duration and — redacted, size-capped — the request
  and response bodies.
- outbound: every HTTP call the backend makes, tagged with who it went to
  (ACME/Let's Encrypt, Cloudflare, GoDaddy, HAProxy stats, agents, the ACME
  diagnostics probe).

Outbound rows inherit the inbound request's id, so one operator action and the
CA/DNS calls it triggered read as a single trace: opening a failed "Request
Certificate" shows the exact POST /acme/new-order and the CA's 429 underneath.

Implementation notes:

- Capture is a pure-ASGI middleware that TEES the request and response streams
  rather than draining them. `await request.body()` inside a BaseHTTPMiddleware
  would consume the receive channel and break the raw-body agent heartbeat
  handler. Registered last so it is outermost: it then sees the final
  client-visible response and seeds correlation_id_context before the error
  handler reads it.
- Rows are written by a batching background writer with a bounded queue, so the
  request path never awaits the database and a saturated logger drops rows
  visibly (surfaced on the page) instead of blocking. Redaction runs on the
  writer, off the request coroutine.
- Secrets never land: headers are an allowlist with Authorization/Cookie kept
  only as a presence marker; body keys and value shapes are redacted
  (passwords, tokens, api_token, API keys, private-key PEMs, JWTs); the ACME
  JWS request body is never stored, because a stored protected+signature pair
  is a replayable credential — a summary is logged instead; DNS-provider errors
  record only the exception type; the ACME HTTP-01 challenge endpoint is
  excluded so key_authorization is never captured.
- Retention is operator-configurable in Settings -> Request Log: separate day
  counts for successful and failed rows (7 / 30) plus a hard row cap (500k),
  whichever is reached first. Pruned in batches under a Postgres advisory lock,
  with the day counts bound as parameters, never interpolated.
- New permissions requestlog.read / requestlog.manage. super_admin and
  security_admin get both, operator gets read, viewer gets neither.

Schema: one new table (request_logs) plus its settings seed, SCHEMA_VERSION
10 -> 11, auto-migrated. No existing table altered, no agent or rendered-config
change. Kill switches: REQUEST_LOG_ENABLED=false (middleware never registered)
or the `enabled` toggle in Settings.

Tests: 245 new (7 backend files + 1 frontend), full suite 1655 backend +
17 frontend passing.
2026-08-11 02:36:03 +03:00

869 lines
41 KiB
Python

"""
ACME v2 (RFC 8555) client service for automated certificate management.
Supports Let's Encrypt, ZeroSSL, Google Trust Services, and any ACME-compatible CA.
"""
import aiohttp
import json
import logging
import hashlib
import base64
import time
from datetime import datetime, timedelta
from typing import Optional, Dict, List, Tuple, Any
from cryptography.hazmat.primitives.asymmetric import rsa, ec, padding
from cryptography.hazmat.primitives import hashes, serialization
from cryptography.hazmat.backends import default_backend
from cryptography.x509.oid import NameOID
from cryptography import x509
import josepy as jose
from database.connection import get_database_connection, close_database_connection
logger = logging.getLogger(__name__)
def _b64url(data: bytes) -> str:
return base64.urlsafe_b64encode(data).rstrip(b'=').decode('ascii')
def _b64url_decode(s: str) -> bytes:
s += '=' * (-len(s) % 4) # pad to a multiple of 4 (0 pad when already aligned)
return base64.urlsafe_b64decode(s)
class ACMEService:
"""Async ACME v2 client with multi-provider support."""
def __init__(self):
self._directory_cache: Dict[str, dict] = {}
# Anti-replay nonces are scoped PER CA (directory_url). A Replay-Nonce issued by one ACME
# server must never be sent in a JWS to another, or the second server rejects it (e.g. ZeroSSL
# "malformed: The Replay Nonce could not be base64url-decoded"). This client is a process-wide
# singleton shared across CAs, so a single shared nonce was leaking across them.
self._nonce_by_dir: Dict[str, str] = {}
async def _get_settings(self) -> dict:
conn = await get_database_connection()
try:
rows = await conn.fetch(
"SELECT key, value FROM system_settings WHERE category = 'acme'"
)
settings = {}
for row in rows:
key = row['key'].split('.', 1)[1] if '.' in row['key'] else row['key']
val = row['value']
if isinstance(val, str):
try:
val = json.loads(val)
except (json.JSONDecodeError, TypeError):
pass
settings[key] = val
return settings
finally:
await close_database_connection(conn)
async def get_directory(self, directory_url: str) -> dict:
if directory_url in self._directory_cache:
cached = self._directory_cache[directory_url]
if cached.get('_fetched_at', 0) > time.time() - 3600:
return cached
# SECURITY (GHSA-3vh4-gvxx-wm2p): directory_url can come from a stored
# account row; validate it (https + public IP, no redirects) before the
# server-side fetch so it cannot be pointed at internal/metadata targets.
from utils.ssrf_guard import assert_public_url, safe_connector
await assert_public_url(directory_url)
# v1.11.0: recorded in request_logs as an outbound call so an operator
# can see exactly which CA was contacted and what it answered.
from utils.http_instrumentation import outbound_span, TARGET_ACME
async with outbound_span(target=TARGET_ACME, method="GET", url=directory_url) as span:
async with aiohttp.ClientSession(connector=safe_connector()) as session:
async with session.get(directory_url, timeout=aiohttp.ClientTimeout(total=15), allow_redirects=False) as resp:
if resp.status != 200:
span.set_response(resp.status, dict(resp.headers))
raise Exception(f"Failed to fetch ACME directory: HTTP {resp.status}")
data = await resp.json()
span.set_response(resp.status, dict(resp.headers), data)
if 'Replay-Nonce' in resp.headers:
self._nonce_by_dir[directory_url] = resp.headers['Replay-Nonce']
data['_fetched_at'] = time.time()
self._directory_cache[directory_url] = data
return data
async def _get_nonce(self, directory_url: str) -> str:
# Use a cached nonce for THIS CA only; otherwise fetch a fresh one from THIS CA's newNonce.
cached = self._nonce_by_dir.pop(directory_url, None)
if cached:
return cached
directory = await self.get_directory(directory_url)
# get_directory may have just captured a nonce for this CA from the directory response.
cached = self._nonce_by_dir.pop(directory_url, None)
if cached:
return cached
# SECURITY (GHSA-3vh4-gvxx-wm2p): newNonce is taken from the (attacker-
# influenceable) directory JSON and is fetched here BEFORE the guarded
# _signed_request POST, so it must be guarded too — otherwise a directory
# that returns an internal newNonce (and omits Replay-Nonce) is a live SSRF.
# https + public IP only, IPv4-pinned connector, no redirects, bounded timeout.
from utils.ssrf_guard import assert_public_url, safe_connector
nonce_url = directory['newNonce']
await assert_public_url(nonce_url)
# v1.11.0: a HEAD with no body and no status check — capture the status
# and the allowlisted headers only. `Replay-Nonce` itself is redacted by
# the header rules: it is a single-use credential.
from utils.http_instrumentation import outbound_span, TARGET_ACME
async with outbound_span(
target=TARGET_ACME, method="HEAD", url=nonce_url, capture_body=False
) as span:
async with aiohttp.ClientSession(connector=safe_connector()) as session:
async with session.head(nonce_url, timeout=aiohttp.ClientTimeout(total=15), allow_redirects=False) as resp:
span.set_response(resp.status, dict(resp.headers))
return resp.headers['Replay-Nonce']
def _generate_account_key(self) -> Tuple[str, dict]:
private_key = rsa.generate_private_key(
public_exponent=65537,
key_size=2048,
backend=default_backend()
)
pem = private_key.private_bytes(
encoding=serialization.Encoding.PEM,
format=serialization.PrivateFormat.PKCS8,
encryption_algorithm=serialization.NoEncryption()
).decode('utf-8')
pub = private_key.public_key()
pub_numbers = pub.public_numbers()
jwk = {
"kty": "RSA",
"n": _b64url(pub_numbers.n.to_bytes((pub_numbers.n.bit_length() + 7) // 8, 'big')),
"e": _b64url(pub_numbers.e.to_bytes((pub_numbers.e.bit_length() + 7) // 8, 'big')),
}
return pem, jwk
def _load_private_key(self, pem_str: str):
return serialization.load_pem_private_key(
pem_str.encode('utf-8'),
password=None,
backend=default_backend()
)
def _get_jwk(self, private_key) -> dict:
pub = private_key.public_key()
pub_numbers = pub.public_numbers()
return {
"kty": "RSA",
"n": _b64url(pub_numbers.n.to_bytes((pub_numbers.n.bit_length() + 7) // 8, 'big')),
"e": _b64url(pub_numbers.e.to_bytes((pub_numbers.e.bit_length() + 7) // 8, 'big')),
}
def _jwk_thumbprint(self, jwk: dict) -> str:
ordered = json.dumps({"e": jwk["e"], "kty": jwk["kty"], "n": jwk["n"]}, separators=(',', ':'))
digest = hashlib.sha256(ordered.encode('utf-8')).digest()
return _b64url(digest)
@staticmethod
def _dns_txt_value(key_authorization: str) -> str:
"""RFC 8555 §8.4: the DNS-01 TXT value is base64url(SHA256(key_authorization)) over the
RAW 32-byte digest (NOT the hexdigest)."""
return _b64url(hashlib.sha256(key_authorization.encode('utf-8')).digest())
@staticmethod
def _challenge_dns_name(identifier: str) -> str:
"""The `_acme-challenge.<base>` record name for an ACME identifier. A leading wildcard
`*.` is stripped, so both `*.example.com` and bare `example.com` map to the SAME name
`_acme-challenge.example.com` (which is why apex+wildcard need two coexisting TXT values)."""
base = identifier[2:] if identifier.startswith('*.') else identifier
return f"_acme-challenge.{base}"
def _sign_jws(self, private_key, protected: dict, payload: Any) -> dict:
protected_b64 = _b64url(json.dumps(protected).encode('utf-8'))
if payload == "":
payload_b64 = ""
else:
payload_b64 = _b64url(json.dumps(payload).encode('utf-8'))
sign_input = f"{protected_b64}.{payload_b64}".encode('ascii')
signature = private_key.sign(sign_input, padding.PKCS1v15(), hashes.SHA256())
return {
"protected": protected_b64,
"payload": payload_b64,
"signature": _b64url(signature),
}
async def _signed_request(
self,
url: str,
directory_url: str,
private_key,
payload: Any,
account_url: Optional[str] = None,
jwk: Optional[dict] = None,
) -> Tuple[int, dict, dict]:
nonce = await self._get_nonce(directory_url)
protected = {"alg": "RS256", "nonce": nonce, "url": url}
if account_url:
protected["kid"] = account_url
elif jwk:
protected["jwk"] = jwk
else:
protected["jwk"] = self._get_jwk(private_key)
body = self._sign_jws(private_key, protected, payload)
# SECURITY (GHSA-3vh4-gvxx-wm2p): `url` is taken from the CA directory /
# order responses. The directory is already fetched from a validated
# public CA, but guard the follow-up POST target too (defence in depth)
# so a tampered/malicious directory cannot steer the request internally.
from utils.ssrf_guard import assert_public_url, safe_connector
await assert_public_url(url)
# v1.11.0: instrument each ATTEMPT separately (the span goes inside the
# retry loop, the session stays outside it) so a badNonce retry shows up
# as its own row instead of being folded into the successful one.
#
# capture_body=False is mandatory here. The JWS body is
# {protected, payload, signature}: `protected` carries the nonce and the
# account kid/jwk, and `signature` is made with the account private key.
# The key itself never crosses the wire, but a stored (protected,
# signature) pair is a REPLAYABLE ACME credential for the lifetime of the
# nonce. We log a description of the request instead of the request.
from utils.http_instrumentation import outbound_span, TARGET_ACME
async with aiohttp.ClientSession(connector=safe_connector()) as session:
for attempt in range(3):
jws_summary = {
"jws": True,
"acme_url": protected.get("url"),
"kid_present": bool(protected.get("kid")),
"jwk_present": bool(protected.get("jwk")),
"payload_empty": payload == "",
"attempt": attempt + 1,
}
async with outbound_span(
target=TARGET_ACME,
method="POST",
url=url,
request_body=jws_summary,
capture_body=False,
) as span:
async with session.post(
url,
json=body,
headers={"Content-Type": "application/jose+json"},
timeout=aiohttp.ClientTimeout(total=30),
allow_redirects=False,
) as resp:
if 'Replay-Nonce' in resp.headers:
self._nonce_by_dir[directory_url] = resp.headers['Replay-Nonce']
if resp.status == 400 and attempt < 2:
err = await resp.json()
etype = (err.get('type') or '')
edetail = (err.get('detail') or '').lower()
# Retry on badNonce, and on any nonce-related malformed rejection (e.g.
# "The Replay Nonce could not be base64url-decoded") — refetch a FRESH nonce
# from the target CA and resign. With per-CA scoping the cross-CA cause is gone;
# this is defense-in-depth so a stale/rejected nonce always self-heals.
if etype.endswith('badNonce') or 'nonce' in edetail:
span.set_response(resp.status, dict(resp.headers), err)
nonce = resp.headers.get('Replay-Nonce') or await self._get_nonce(directory_url)
protected['nonce'] = nonce
body = self._sign_jws(private_key, protected, payload)
continue
resp_data = {}
content_type = resp.headers.get('Content-Type', '')
if 'json' in content_type:
resp_data = await resp.json()
elif resp.status < 300:
text = await resp.text()
if text:
try:
resp_data = json.loads(text)
except json.JSONDecodeError:
resp_data = {"raw": text}
headers = dict(resp.headers)
span.set_response(resp.status, headers, resp_data)
return resp.status, resp_data, headers
raise Exception(f"ACME request to {url} failed after retries")
async def register_account(
self,
email: str,
directory_url: str,
tos_agreed: bool = True,
eab_kid: Optional[str] = None,
eab_hmac_key: Optional[str] = None,
challenge_type: str = 'http-01',
dns_provider: Optional[str] = None,
) -> dict:
directory = await self.get_directory(directory_url)
pem, jwk = self._generate_account_key()
private_key = self._load_private_key(pem)
payload: Dict[str, Any] = {
"termsOfServiceAgreed": tos_agreed,
"contact": [f"mailto:{email}"],
}
if eab_kid and eab_hmac_key:
eab_key_bytes = _b64url_decode(eab_hmac_key)
eab_protected = {
"alg": "HS256",
"kid": eab_kid,
"url": directory['newAccount'],
}
eab_protected_b64 = _b64url(json.dumps(eab_protected).encode('utf-8'))
eab_payload_b64 = _b64url(json.dumps(jwk).encode('utf-8'))
import hmac as hmac_mod
eab_sign_input = f"{eab_protected_b64}.{eab_payload_b64}".encode('ascii')
eab_signature = hmac_mod.new(eab_key_bytes, eab_sign_input, hashlib.sha256).digest()
payload["externalAccountBinding"] = {
"protected": eab_protected_b64,
"payload": eab_payload_b64,
"signature": _b64url(eab_signature),
}
status, data, headers = await self._signed_request(
directory['newAccount'], directory_url, private_key, payload, jwk=jwk
)
if status not in (200, 201):
raise Exception(f"Account registration failed: {data}")
account_url = headers.get('Location', '')
conn = await get_database_connection()
try:
row = await conn.fetchrow("""
INSERT INTO letsencrypt_accounts (email, directory_url, account_url, jwk_private_key, status, tos_agreed, eab_kid, challenge_type, dns_provider)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
ON CONFLICT (email, directory_url) DO UPDATE SET
account_url = $3, jwk_private_key = $4, status = $5, tos_agreed = $6,
challenge_type = $8, dns_provider = $9, updated_at = NOW()
RETURNING id, email, directory_url, account_url, status, tos_agreed, created_at, challenge_type, dns_provider
""", email, directory_url, account_url, pem,
data.get('status') or 'valid', tos_agreed, eab_kid, challenge_type, dns_provider)
return dict(row)
finally:
await close_database_connection(conn)
async def deactivate_account(self, account_id: int) -> dict:
"""Deactivate an ACME account at the CA and mark it locally."""
conn = await get_database_connection()
try:
account = await conn.fetchrow(
"SELECT id, email, directory_url, account_url, jwk_private_key, status FROM letsencrypt_accounts WHERE id = $1",
account_id
)
if not account:
raise Exception("Account not found")
if account['status'] == 'deactivated':
raise Exception("Account is already deactivated")
private_key = self._load_private_key(account['jwk_private_key'])
payload = {"status": "deactivated"}
status, data, _headers = await self._signed_request(
account['account_url'],
account['directory_url'],
private_key,
payload,
account_url=account['account_url'],
)
if status not in (200, 201):
logger.warning(f"ACME account deactivation returned {status}: {data}")
raise Exception(f"CA rejected deactivation (HTTP {status}): {data.get('detail') or 'Unknown error'}")
await conn.execute(
"UPDATE letsencrypt_accounts SET status = 'deactivated', updated_at = NOW() WHERE id = $1",
account_id
)
return {"id": account_id, "email": account['email'], "status": "deactivated"}
finally:
await close_database_connection(conn)
async def create_order(
self,
account_id: int,
domains: List[str],
cluster_ids: Optional[List[int]] = None,
challenge_type: str = 'http-01',
created_by: Optional[int] = None,
) -> dict:
logger.info(f"ACME: Creating order for domains={domains}, account_id={account_id}, challenge_type={challenge_type}")
conn = await get_database_connection()
try:
account = await conn.fetchrow(
"SELECT * FROM letsencrypt_accounts WHERE id = $1", account_id
)
if not account:
raise Exception(f"Account {account_id} not found")
private_key = self._load_private_key(account['jwk_private_key'])
directory = await self.get_directory(account['directory_url'])
identifiers = [{"type": "dns", "value": d} for d in domains]
payload = {"identifiers": identifiers}
status, data, headers = await self._signed_request(
directory['newOrder'],
account['directory_url'],
private_key,
payload,
account_url=account['account_url'],
)
if status not in (200, 201):
logger.error(f"ACME: Order creation failed: HTTP {status}, response={data}")
raise Exception(f"Order creation failed: {data}")
order_url = headers.get('Location', '')
logger.info(f"ACME: Order created, order_url={order_url}, status={data.get('status')}, authorizations={len(data.get('authorizations', []))}")
expires_at = None
if data.get('expires'):
try:
expires_at = datetime.fromisoformat(data['expires'].replace('Z', '+00:00'))
except (ValueError, TypeError):
pass
order_row = await conn.fetchrow("""
INSERT INTO letsencrypt_orders
(account_id, order_url, status, domains, finalize_url, expires_at, cluster_ids, challenge_type, created_by)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
RETURNING id
""", account_id, order_url, data.get('status') or 'pending',
json.dumps(domains), data.get('finalize') or '', expires_at,
json.dumps(cluster_ids or []), challenge_type, created_by)
order_id = order_row['id']
# Commit 5e: track authorization fetch outcomes per-domain.
# Previously a 'continue' on auth fetch failure silently dropped HTTP-01
# challenges; the order proceeded but had no challenges to respond to,
# leaving it stuck in 'pending' forever.
auth_fetch_failures = [] # [{auth_url, http_status, error}]
domains_with_http01 = set()
for auth_url in data.get('authorizations', []):
try:
auth_status, auth_data, _ = await self._signed_request(
auth_url, account['directory_url'], private_key, "",
account_url=account['account_url'],
)
except Exception as auth_err:
logger.warning(f"ACME: authorization fetch raised: {auth_url} - {auth_err}")
auth_fetch_failures.append({
"auth_url": auth_url, "http_status": None, "error": str(auth_err)
})
continue
if auth_status != 200:
logger.warning(
f"ACME: Failed to fetch authorization {auth_url}: HTTP {auth_status}, body={str(auth_data)[:200]}"
)
auth_fetch_failures.append({
"auth_url": auth_url, "http_status": auth_status,
"error": str(auth_data)[:500] if auth_data else ""
})
continue
domain = (auth_data.get('identifier') or {}).get('value', '')
http01_for_domain = False
for challenge in (auth_data.get('challenges') or []):
# Store only the challenge of the CHOSEN method (default 'http-01' keeps the
# existing behaviour byte-identical; 'dns-01' selects the TXT challenge instead).
if challenge.get('type') == challenge_type:
token = challenge['token']
jwk = self._get_jwk(private_key)
thumbprint = self._jwk_thumbprint(jwk)
key_auth = f"{token}.{thumbprint}"
dns_txt = self._dns_txt_value(key_auth) if challenge_type == 'dns-01' else None
await conn.execute("""
INSERT INTO acme_challenges
(order_id, domain, token, key_authorization, challenge_url, status,
challenge_type, dns_txt_value)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
""", order_id, domain, token, key_auth,
challenge.get('url') or '', challenge.get('status') or 'pending',
challenge_type, dns_txt)
logger.info(f"ACME: {challenge_type} challenge stored for domain={domain}, token={token[:20]}..., challenge_url={(challenge.get('url') or '')[:60]}")
http01_for_domain = True
if http01_for_domain and domain:
domains_with_http01.add(domain)
# If NO http-01 challenges were registered at all, the order cannot
# proceed — fail-fast and persist diagnostic detail.
if not domains_with_http01:
error_payload = json.dumps({
"stage": "create_order_authorizations",
"auth_fetch_failures": auth_fetch_failures,
"domains": domains,
"timestamp": datetime.utcnow().isoformat(),
})
await conn.execute(
"UPDATE letsencrypt_orders SET status = 'invalid', error_detail = $1, updated_at = NOW() WHERE id = $2",
error_payload, order_id
)
raise Exception(
f"ACME order {order_id} created but no {challenge_type} challenges available "
f"(auth fetch failures: {len(auth_fetch_failures)}). See order.error_detail for diagnostics."
)
elif auth_fetch_failures:
# Partial failure: some domains have challenges, others don't. Record warning.
error_payload = json.dumps({
"stage": "create_order_authorizations_partial",
"auth_fetch_failures": auth_fetch_failures,
"domains_with_challenges": sorted(domains_with_http01),
"timestamp": datetime.utcnow().isoformat(),
})
await conn.execute(
"UPDATE letsencrypt_orders SET error_detail = $1, updated_at = NOW() WHERE id = $2",
error_payload, order_id
)
logger.warning(
f"ACME order {order_id}: partial authorization fetch failure — "
f"{len(auth_fetch_failures)} failed, {len(domains_with_http01)} succeeded"
)
return {
"order_id": order_id,
"order_url": order_url,
"status": data.get('status') or 'pending',
"domains": domains,
"authorizations": data.get('authorizations') or [],
"finalize": data.get('finalize') or '',
}
finally:
await close_database_connection(conn)
async def respond_to_challenges(self, order_id: int) -> List[dict]:
conn = await get_database_connection()
try:
# Issue #12 / Commit 5a: include 'failed' challenges so they can be retried,
# but rate-limit per challenge: max 5 attempts in last 5 minutes.
# DNS-01 skip-gate (the single safe choke point): never POST a challenge response for a
# dns-01 row whose TXT record has not been published yet — that would make the CA validate
# against a missing record and burn the order. http-01 rows (challenge_type 'http-01'/NULL)
# are never excluded, so the existing flow is byte-identical.
challenges = await conn.fetch(
"""SELECT * FROM acme_challenges
WHERE order_id = $1
AND (status IN ('pending', 'failed') OR status IS NULL)
AND NOT (COALESCE(challenge_type, 'http-01') = 'dns-01' AND COALESCE(dns_record_published, FALSE) = FALSE)""",
order_id
)
order = await conn.fetchrow(
"SELECT o.*, a.jwk_private_key, a.account_url, a.directory_url FROM letsencrypt_orders o JOIN letsencrypt_accounts a ON o.account_id = a.id WHERE o.id = $1",
order_id
)
if not order:
raise Exception(f"Order {order_id} not found")
private_key = self._load_private_key(order['jwk_private_key'])
results = []
logger.info(f"ACME: Responding to {len(challenges)} challenge(s) for order_id={order_id}")
for ch in challenges:
if not ch['challenge_url']:
logger.warning(f"ACME: Skipping challenge id={ch['id']} domain={ch['domain']} - no challenge_url")
continue
# Rate-limit retry: skip if attempted >=5 times in last 5 minutes
attempts = ch.get('attempts') or 0
last_attempt = ch.get('last_attempt_at')
if attempts >= 5 and last_attempt:
age_seconds = (datetime.utcnow().replace(tzinfo=None) -
(last_attempt.replace(tzinfo=None) if last_attempt.tzinfo else last_attempt)).total_seconds()
if age_seconds < 300:
logger.warning(
f"ACME: Rate-limit: skipping challenge id={ch['id']} domain={ch['domain']} "
f"(attempts={attempts}, age={int(age_seconds)}s < 300s)"
)
results.append({"domain": ch['domain'], "token": ch['token'],
"status": ch['status'], "skipped": "rate-limit"})
continue
status, data, _ = await self._signed_request(
ch['challenge_url'],
order['directory_url'],
private_key,
{},
account_url=order['account_url'],
)
new_status = (data.get('status') or 'processing') if status == 200 else 'failed'
logger.info(f"ACME: Challenge response for domain={ch['domain']}, token={ch['token'][:20]}..., CA_HTTP={status}, CA_status_raw={data.get('status')!r}, stored_status={new_status}")
await conn.execute(
"""UPDATE acme_challenges
SET status = $1,
attempts = COALESCE(attempts, 0) + 1,
last_attempt_at = NOW()
WHERE id = $2""",
new_status, ch['id']
)
results.append({"domain": ch['domain'], "token": ch['token'], "status": new_status})
return results
finally:
await close_database_connection(conn)
async def finalize_order(self, order_id: int) -> dict:
conn = await get_database_connection()
try:
order = await conn.fetchrow("""
SELECT o.*, a.jwk_private_key, a.account_url, a.directory_url
FROM letsencrypt_orders o
JOIN letsencrypt_accounts a ON o.account_id = a.id
WHERE o.id = $1
""", order_id)
if not order:
raise Exception(f"Order {order_id} not found")
domains = json.loads(order['domains']) if isinstance(order['domains'], str) else order['domains']
logger.info(f"ACME: Finalizing order_id={order_id}, domains={domains}, finalize_url={(order['finalize_url'] or 'N/A')[:60]}")
private_key = self._load_private_key(order['jwk_private_key'])
cert_key = rsa.generate_private_key(
public_exponent=65537, key_size=2048, backend=default_backend()
)
cert_key_pem = cert_key.private_bytes(
serialization.Encoding.PEM,
serialization.PrivateFormat.PKCS8,
serialization.NoEncryption()
).decode('utf-8')
builder = x509.CertificateSigningRequestBuilder()
builder = builder.subject_name(x509.Name([
x509.NameAttribute(NameOID.COMMON_NAME, domains[0]),
]))
san_names = [x509.DNSName(d) for d in domains]
builder = builder.add_extension(
x509.SubjectAlternativeName(san_names), critical=False
)
csr = builder.sign(cert_key, hashes.SHA256(), default_backend())
csr_der = csr.public_bytes(serialization.Encoding.DER)
payload = {"csr": _b64url(csr_der)}
status, data, _ = await self._signed_request(
order['finalize_url'],
order['directory_url'],
private_key,
payload,
account_url=order['account_url'],
)
if status not in (200, 201):
# Commit 5h: structured JSON-as-TEXT error_detail for consistency
# with check_order_status (5c) and download_certificate (5i).
error_msg = data.get('detail') if isinstance(data, dict) else str(data)
error_payload = json.dumps({
"stage": "finalize_order",
"http_status": status,
"ca_response": data if isinstance(data, dict) else str(data)[:1000],
"timestamp": datetime.utcnow().isoformat(),
})
logger.error(f"ACME: Finalize failed for order_id={order_id}: HTTP {status}, error={error_msg}")
await conn.execute(
"UPDATE letsencrypt_orders SET status = 'invalid', error_detail = $1, updated_at = NOW() WHERE id = $2",
error_payload, order_id
)
raise Exception(f"Finalize failed: {error_msg}")
order_status = data.get('status') or 'processing'
certificate_url = data.get('certificate') or ''
logger.info(f"ACME: Finalize success for order_id={order_id}, status={order_status}, certificate_url={certificate_url[:60] if certificate_url else 'N/A'}")
await conn.execute(
"UPDATE letsencrypt_orders SET status = $1, certificate_url = $2, cert_private_key = $3, updated_at = NOW() WHERE id = $4",
order_status, certificate_url, cert_key_pem, order_id
)
return {
"order_id": order_id,
"status": order_status,
"certificate_url": certificate_url,
"private_key_pem": cert_key_pem,
}
finally:
await close_database_connection(conn)
async def download_certificate(self, order_id: int) -> dict:
conn = await get_database_connection()
try:
order = await conn.fetchrow("""
SELECT o.*, a.jwk_private_key, a.account_url, a.directory_url
FROM letsencrypt_orders o
JOIN letsencrypt_accounts a ON o.account_id = a.id
WHERE o.id = $1
""", order_id)
if not order or not order['certificate_url']:
raise Exception("Certificate not ready for download")
private_key = self._load_private_key(order['jwk_private_key'])
status, data, headers = await self._signed_request(
order['certificate_url'],
order['directory_url'],
private_key,
"",
account_url=order['account_url'],
)
if status != 200:
# Commit 5i: persist structured error_detail to TEXT column.
# Previously download failures only raised an exception, leaving the
# order in 'valid' state with no DB diagnostic — operators were
# blind to root cause.
try:
error_payload = json.dumps({
"stage": "download_certificate",
"http_status": status,
"ca_response": data if isinstance(data, dict) else str(data)[:1000],
"timestamp": datetime.utcnow().isoformat(),
})
await conn.execute(
"UPDATE letsencrypt_orders SET error_detail = $1, updated_at = NOW() WHERE id = $2",
error_payload, order_id
)
except Exception as persist_err:
logger.warning(f"ACME: failed to persist download error for order {order_id}: {persist_err}")
raise Exception(f"Certificate download failed: HTTP {status}")
# Commit 5d: harden response parsing. ACME RFC 8555 §7.4.2 specifies
# `application/pem-certificate-chain` as the Content-Type. Some test
# CAs (Pebble) return raw PEM, others return JSON-wrapped data. We
# accept either format and extract the PEM body robustly.
cert_pem = ''
if isinstance(data, dict):
cert_pem = data.get('raw', '') or data.get('certificate', '') or ''
elif isinstance(data, (str, bytes)):
cert_pem = data.decode('utf-8') if isinstance(data, bytes) else data
elif data is not None:
cert_pem = str(data)
if not cert_pem or '-----BEGIN CERTIFICATE-----' not in cert_pem:
raise Exception(
f"Certificate download succeeded (HTTP 200) but response body "
f"does not contain a PEM certificate. Content-Type={headers.get('Content-Type', 'unknown') if headers else 'unknown'}, "
f"body_len={len(cert_pem) if cert_pem else 0}"
)
parts = cert_pem.strip().split('-----END CERTIFICATE-----')
certificate = (parts[0] + '-----END CERTIFICATE-----').strip() if parts else cert_pem
chain = '-----END CERTIFICATE-----'.join(parts[1:]).strip() if len(parts) > 1 else ''
if chain and not chain.startswith('-----'):
chain = chain.lstrip('\n')
await conn.execute(
"UPDATE letsencrypt_orders SET status = 'valid', updated_at = NOW() WHERE id = $1",
order_id
)
return {
"certificate_pem": certificate,
"chain_pem": chain,
"full_chain_pem": cert_pem,
}
finally:
await close_database_connection(conn)
async def check_order_status(self, order_id: int) -> dict:
conn = await get_database_connection()
try:
order = await conn.fetchrow("""
SELECT o.*, a.jwk_private_key, a.account_url, a.directory_url
FROM letsencrypt_orders o
JOIN letsencrypt_accounts a ON o.account_id = a.id
WHERE o.id = $1
""", order_id)
if not order:
raise Exception(f"Order {order_id} not found")
if not order['order_url']:
return {"order_id": order_id, "status": order['status']}
private_key = self._load_private_key(order['jwk_private_key'])
status, data, _ = await self._signed_request(
order['order_url'],
order['directory_url'],
private_key,
"",
account_url=order['account_url'],
)
if status == 200:
new_status = data.get('status') or order['status']
certificate_url = data.get('certificate') or order['certificate_url']
await conn.execute(
"UPDATE letsencrypt_orders SET status = $1, certificate_url = $2, updated_at = NOW() WHERE id = $3",
new_status, certificate_url, order_id
)
return {"order_id": order_id, "status": new_status, "certificate_url": certificate_url}
# Commit 5c: persist CA error to error_detail (TEXT) as structured JSON.
# Previously a non-200 was silently swallowed (no log, no DB record),
# leaving operators without diagnostic info for stuck orders.
try:
error_payload = json.dumps({
"stage": "check_order_status",
"http_status": status,
"ca_response": data if isinstance(data, dict) else str(data)[:1000],
"timestamp": datetime.utcnow().isoformat(),
})
await conn.execute(
"UPDATE letsencrypt_orders SET error_detail = $1, updated_at = NOW() WHERE id = $2",
error_payload, order_id
)
except Exception as persist_err:
logger.warning(f"ACME: failed to persist error_detail for order {order_id}: {persist_err}")
logger.warning(
f"ACME: check_order_status non-200 for order {order_id}: HTTP {status}, response={data}"
)
return {"order_id": order_id, "status": order['status']}
finally:
await close_database_connection(conn)
async def revoke_certificate(self, certificate_pem: str, account_id: int, reason: int = 0) -> bool:
conn = await get_database_connection()
try:
account = await conn.fetchrow(
"SELECT * FROM letsencrypt_accounts WHERE id = $1", account_id
)
if not account:
raise Exception("Account not found")
private_key = self._load_private_key(account['jwk_private_key'])
directory = await self.get_directory(account['directory_url'])
cert = x509.load_pem_x509_certificate(certificate_pem.encode('utf-8'), default_backend())
cert_der = cert.public_bytes(serialization.Encoding.DER)
payload = {"certificate": _b64url(cert_der), "reason": reason}
status, data, _ = await self._signed_request(
directory['revokeCert'],
account['directory_url'],
private_key,
payload,
account_url=account['account_url'],
)
return status == 200
finally:
await close_database_connection(conn)
acme_service = ACMEService()