Add Home Assistant Stelloauth app #1

Merged
dennis merged 17 commits from feat/home-assistant-stelloauth-addon into main 2026-09-24 20:19:49 +02:00
24 changed files with 3060 additions and 0 deletions
+4
View File
@@ -0,0 +1,4 @@
.git
.venv
.pytest_cache
**/__pycache__
+1
View File
@@ -0,0 +1 @@
*.patch -whitespace
+27
View File
@@ -0,0 +1,27 @@
name: CI
"on":
push:
branches: [main]
pull_request:
branches: [main]
permissions:
contents: read
jobs:
validate:
runs-on: ubuntu-latest
timeout-minutes: 45
steps:
- uses: actions/checkout@11bd71901bbe5b1630ceea73d27597364c9af683
- name: Validate runner architecture
run: test "$(uname -m)" = "x86_64"
- name: Build validation image
run: docker build --file tests/Dockerfile.ci --tag homeassistant-stelloauth-tests:test .
- name: Validate metadata, patches, and process manager
run: docker run --rm homeassistant-stelloauth-tests:test -q
- name: Build amd64 image
run: docker build --build-arg TARGETARCH=amd64 --build-arg BUILD_ARCH=amd64 --tag homeassistant-stelloauth-addon:test stelloauth
- name: Exercise runtime
run: SKIP_BUILD=1 tests/test_runtime.sh
+6
View File
@@ -0,0 +1,6 @@
.venv/
.pytest_cache/
__pycache__/
*.pyc
.coverage
artifacts/
+21
View File
@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2026 Dennis / Radix ApS
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+59
View File
@@ -0,0 +1,59 @@
# Stelloauth for Home Assistant
Et Home Assistant custom add-on, som kører Stelloauth og CloakBrowser lokalt til
OAuth-opsætning af integrationen Stellantis Vehicles. Repositoryet understøtter
`amd64` og `aarch64`.
```text
Home Assistant Core
|
| POST http://0031621f-stelloauth:8080/worker
v
+--------------------------------------------------+
| Ét add-on / én container |
| |
| Stelloauth 0.0.0.0:8080 |
| | |
| | CDP http://127.0.0.1:9222 |
| v |
| CloakBrowser 127.0.0.1:9222 |
+--------------------------------------------------+
```
Én add-on giver de to processer samme Supervisor-livscyklus og holder
CloakBrowsers CDP-port på containerens loopback-interface. Add-onen bygges
lokalt fra kilde, fordi CloakBrowser Binary License ikke tillader, at dette
repository genudgiver en afledt image med den proprietære CloakBrowser-binær.
Der publiceres derfor ingen prebuilt images.
## Installation
1. Tilføj
`https://git.radixadm.dk/dennis/homeassistant-stelloauth-addon.git` som
repository i Home Assistants Tilføjelsesbutik.
2. Installér **Stelloauth**, aktivér **Start ved opstart** og **Watchdog**, og
start add-onen.
3. Følg den fulde vejledning i [stelloauth/DOCS.md](stelloauth/DOCS.md).
HAOS-installation er ikke gennemført eller påstået som valideret. Målmaskinen
havde ved inspektionen approximately 4 GB free, mens det lokalt byggede image
fyldte 2.571 GB; frigør om nødvendigt mere diskplads før installation.
## Fastlåste upstream-kilder
| Komponent | Version | Commit / OCI index digest |
| --- | --- | --- |
| Stelloauth | `v0.6.0` | `367d4f8c02a3b072c59142c49dffc129edc8548b` |
| Stelloauth image contract | `v0.6.0` | `sha256:51b2194ec9b80cc484d11c016ec5436a12161277a493afdab7078959966d5aa9` |
| CloakBrowser | `0.5.10` | `f04c23da285b3b3d3cf10c8f9d282e7adc1d52ce` / `sha256:2ed5b2d047cbdde22cde7ef1a796526c716aadaa5bccbe1db5ade49282b64a76` |
| Go-builder | `1.27.1-bookworm` | `sha256:69a7b9788769bec032d238959b61854e9ae87f57be9029ec04e9885fabf99195` |
| Stellantis Vehicles | `2026.9.4` | `9e0ef96f8fe478af291da4c38c923ada78d0ebf6` |
Dockerfile og patches er build-opskriften. Supervisor henter det officielle
CloakBrowser-image og bygger alene et lokalt image til intern brug.
## Licens
Repositoryets eget arbejde er MIT-licenseret; se [LICENSE](LICENSE). Dette
omfatter ikke CloakBrowsers proprietære binær. Den er fortsat omfattet af den
separate **CloakBrowser Binary License** og redistribueres ikke af repositoryet.
+4
View File
@@ -0,0 +1,4 @@
[pytest]
markers =
upstream: fetches and validates pinned upstream source
docker: builds or runs the add-on image
+3
View File
@@ -0,0 +1,3 @@
name: Stelloauth for Home Assistant
url: https://git.radixadm.dk/dennis/homeassistant-stelloauth-addon.git
maintainer: Dennis / Radix ApS
+2
View File
@@ -0,0 +1,2 @@
pytest==8.4.2
PyYAML==6.0.2
+11
View File
@@ -0,0 +1,11 @@
# Changelog
## 0.1.0
- Fastlåser Stelloauth `v0.6.0` og CloakBrowser `0.5.10` til verificerede
commits og OCI-digests.
- Tilføjer fælles procesovervågning, readiness, watchdog og begrænset shutdown.
- Dokumenterer intern Login service URL:
`http://0031621f-stelloauth:8080/worker`.
- Begrænser CDP til loopback og hardener URL-validering, request-størrelse,
rate limiting og logredigering.
+89
View File
@@ -0,0 +1,89 @@
# Installation og drift
1. Gå til **Indstillinger → Tilføjelser → Tilføjelsesbutik → ⋮ →
Repositorier**, og tilføj præcis
`https://git.radixadm.dk/dennis/homeassistant-stelloauth-addon.git`.
2. Installér **Stelloauth**, aktivér **Start ved opstart** og **Watchdog**, og
start derefter add-onen. Installationen bygger et lokalt image fra kilde til
den valgte `amd64`- eller `aarch64`-arkitektur. Repositoryet publicerer ikke
et prebuilt image.
3. Vent, til loggen i denne rækkefølge viser de fem faste readiness-beskeder:
```text
Cleaning CloakBrowser profiles
Starting CloakBrowser
CloakBrowser ready
Starting Stelloauth
Stelloauth listening on 0.0.0.0:8080
```
4. Behold host-porten deaktiveret i normal drift. Ved kortvarig fejlfinding kan
`8080/tcp` tilknyttes host-port `8080`. Kontrollér derefter
`http://192.168.1.20:8080/` eller worker-endpointet
`http://192.168.1.20:8080/worker`, og **deaktivér porttilknytningen igen**,
når kontrollen er færdig. Worker-endpointet modtager MyOpel-oplysninger og
har ingen egen autentificering.
5. Åbn konfigurationen af **Stellantis Vehicles**. Angiv præcis
`http://0031621f-stelloauth:8080/worker` som **Login service URL**.
Integrationen tilføjer ikke `/worker`; hele stien skal derfor stå i feltet.
6. Vælg **Brand: Opel** og **Country: DK**, og gennemfør derefter integrationens
OAuth-opsætning med dine MyOpel-oplysninger.
7. Add-onens tre muligheder er:
- `queue_timeout`: hvor længe et loginforsøg må vente på den ene session.
- `rate_limit_count`: højeste antal loginforsøg i hver periode.
- `rate_limit_duration`: længden af rate limit-perioden.
`CLOAK_MAX_SESSIONS` er fastlåst til én session, fordi CloakBrowsers gratis
niveau tillader ét samtidigt login. Samtidige forsøg bliver derfor køet.
8. Hvert OAuth-forsøg får en midlertidig profil under `/tmp/cloakserve`.
CloakBrowser rydder inaktive browserprocesser efter 30 sekunder, og
process manageren rydder gamle profiler ved opstart. Credentials, cookies,
tokens og OAuth-koder gemmes ikke i `/data`. CDP lytter kun på loopback
`127.0.0.1:9222`, og logs bruger faste, redigerede hændelser uden email,
passwords, URLs, koder eller tokens.
9. De målte resultater fra den reelle `linux/amd64`-kørsel under Rosetta var:
- Image: 2,571,693,650 bytes (2.571 GB decimal / 2452.56 MiB).
- Dokumenteret Task 4-måling: 121,5 MiB.
- Stop: cirka 9.3 sekunder.
- `amd64` runtime bestod under Rosetta; `aarch64` build bestod.
- Mål-HAOS havde ved inspektionen approximately 4 GB free. Den knappe
plads sammenholdt med image- og build-lag kan forhindre installationen;
frigør plads først. Der er ikke verificeret en vellykket HAOS-installation.
Der er ikke gennemført et live MyOpel-login.
Login-flow RAM: not measured without real MyOpel credentials.
10. Fejlfinding og fjernelse:
- Mangler en readiness-besked, så se efter timeout: CloakBrowser har 60
sekunder og Stelloauth 30 sekunder. Ret årsagen og genstart add-onen.
- Et ugyldigt eller ikke-tilladt authorize-URL giver HTTP `400`.
- For mange loginforsøg giver HTTP `429`; vent den konfigurerede periode.
- Hvis repository-URL eller hostname ændres, ændres Supervisor-repository-ID
og dermed `0031621f-stelloauth`. Beregn og brug den nye interne URL.
- Ved disk pressure: kontrollér fri plads og fjern unødvendige images eller
backups via de normale Supervisor-funktioner før et nyt build.
- Hvis den interne URL ikke kan nås, brug kun den midlertidige portkontrol
fra trin 4 og deaktivér porttilknytningen bagefter.
- Fjernelse sker i **Indstillinger → Tilføjelser → Stelloauth → Afinstallér**.
Supervisor stopper containeren og fjerner add-onens lokale data; fjern
også repositoryet fra Tilføjelsesbutikken, hvis det ikke længere bruges.
## Kilder og licenser
Add-on-version `0.1.0` bygger Stelloauth `v0.6.0` fra commit
`367d4f8c02a3b072c59142c49dffc129edc8548b` og bruger det officielle
CloakBrowser `0.5.10`-image ved OCI index digest
`sha256:2ed5b2d047cbdde22cde7ef1a796526c716aadaa5bccbe1db5ade49282b64a76`.
Repositoryets egne filer og patches er MIT-licenserede. CloakBrowsers
proprietære binær er fortsat under den separate **CloakBrowser Binary License**;
den er ikke MIT-licenseret eller redistribueret af dette repository.
+55
View File
@@ -0,0 +1,55 @@
# syntax=docker/dockerfile:1.7
FROM golang:1.27.1-bookworm@sha256:69a7b9788769bec032d238959b61854e9ae87f57be9029ec04e9885fabf99195 AS stelloauth-builder
ARG TARGETARCH
RUN apt-get update \
&& apt-get install -y --no-install-recommends ca-certificates git patch \
&& rm -rf /var/lib/apt/lists/*
WORKDIR /src
RUN git clone --filter=blob:none https://github.com/tamcore/stelloauth.git . \
&& git checkout --detach 367d4f8c02a3b072c59142c49dffc129edc8548b \
&& test "$(git rev-parse HEAD)" = "367d4f8c02a3b072c59142c49dffc129edc8548b"
COPY patches/stelloauth-security.patch /tmp/stelloauth-security.patch
RUN git apply --check /tmp/stelloauth-security.patch \
&& git apply /tmp/stelloauth-security.patch \
&& go test ./... \
&& CGO_ENABLED=0 GOOS=linux GOARCH="$TARGETARCH" go build -trimpath -ldflags="-s -w" -o /out/stelloauth ./cmd/stelloauth
FROM cloakhq/cloakbrowser:0.5.10@sha256:2ed5b2d047cbdde22cde7ef1a796526c716aadaa5bccbe1db5ade49282b64a76
ARG BUILD_ARCH
ARG BUILD_DATE
ARG BUILD_DESCRIPTION="Local OAuth worker for Stellantis Vehicles using CloakBrowser"
ARG BUILD_NAME="Stelloauth"
ARG BUILD_REF
ARG BUILD_REPOSITORY="https://git.radixadm.dk/dennis/homeassistant-stelloauth-addon"
ARG BUILD_VERSION="0.1.0"
LABEL io.hass.name="$BUILD_NAME" \
io.hass.description="$BUILD_DESCRIPTION" \
io.hass.arch="$BUILD_ARCH" \
io.hass.type="app" \
io.hass.version="$BUILD_VERSION" \
org.opencontainers.image.created="$BUILD_DATE" \
org.opencontainers.image.revision="$BUILD_REF" \
org.opencontainers.image.source="$BUILD_REPOSITORY" \
org.opencontainers.image.version="$BUILD_VERSION"
USER root
RUN apt-get update \
&& apt-get install -y --no-install-recommends patch \
&& rm -rf /var/lib/apt/lists/*
COPY patches/cloakserve-loopback.patch /tmp/cloakserve-loopback.patch
RUN patch --dry-run -p2 -d /usr/local/bin < /tmp/cloakserve-loopback.patch \
&& patch -p2 -d /usr/local/bin < /tmp/cloakserve-loopback.patch \
&& rm /tmp/cloakserve-loopback.patch \
&& apt-get purge -y --auto-remove patch \
&& rm -rf /var/lib/apt/lists/*
COPY --from=stelloauth-builder /out/stelloauth /usr/local/bin/stelloauth
COPY rootfs/ /
RUN chmod 0755 /usr/local/bin/stelloauth /usr/local/bin/addon-supervisor /usr/local/bin/cloakserve \
&& mkdir -p /data /tmp/cloakserve \
&& chmod 0700 /tmp/cloakserve
EXPOSE 8080
ENTRYPOINT []
CMD ["/usr/local/bin/addon-supervisor"]
+10
View File
@@ -0,0 +1,10 @@
# Stelloauth
Lokal OAuth-worker til integrationen Stellantis Vehicles. Add-onen bygger
Stelloauth `v0.6.0` og CloakBrowser `0.5.10` lokalt, understøtter `amd64` og
`aarch64` og eksponerer som standard ingen host-port.
Se [den fulde installations- og fejlfindingsvejledning](DOCS.md).
CloakBrowsers binær er under den separate CloakBrowser Binary License og er
ikke omfattet af repositoryets MIT-licens.
+24
View File
@@ -0,0 +1,24 @@
name: Stelloauth
version: "0.1.0"
slug: stelloauth
description: Local OAuth worker for Stellantis Vehicles using CloakBrowser
url: https://git.radixadm.dk/dennis/homeassistant-stelloauth-addon
arch:
- amd64
- aarch64
startup: application
boot: auto
init: true
watchdog: http://[HOST]:[PORT:8080]/
ports:
8080/tcp: null
ports_description:
8080/tcp: Temporary LAN access for troubleshooting only
options:
queue_timeout: 60s
rate_limit_count: 5
rate_limit_duration: 1h
schema:
queue_timeout: match(^[1-9][0-9]*(ms|s|m|h)$)
rate_limit_count: int(1,20)
rate_limit_duration: match(^[1-9][0-9]*(ms|s|m|h)$)
@@ -0,0 +1,14 @@
diff --git a/bin/cloakserve b/bin/cloakserve
index 0e545b2..04354aa 100755
--- a/bin/cloakserve
+++ b/bin/cloakserve
@@ -903,8 +903,7 @@ def main() -> None:
port,
)
- in_container = os.path.exists("/.dockerenv") or os.path.exists("/run/.containerenv")
- host = "0.0.0.0" if in_container else "127.0.0.1"
+ host = "127.0.0.1"
web.run_app(app, host=host, port=port, print=None)
@@ -0,0 +1,552 @@
diff --git a/internal/app/oauth.go b/internal/app/oauth.go
index f0b9548..0c59a11 100644
--- a/internal/app/oauth.go
+++ b/internal/app/oauth.go
@@ -162,12 +162,11 @@ func performChromedpOAuth(
reqURL := e.Request.URL
// Capture OAuth redirect
if strings.HasPrefix(reqURL, redirectPrefix) {
- log.Printf("[%s] Redirect URL: %s", requestID, reqURL)
+ logCapturedOAuthRedirect(requestID, reqURL)
parsed, err := url.Parse(reqURL)
if err == nil {
if code := parsed.Query().Get("code"); code != "" {
oauthCode = code
- log.Printf("[%s] Captured OAuth code from redirect request", requestID)
}
}
} else if strings.Contains(reqURL, "OPErrorPage.php") {
@@ -175,7 +174,7 @@ func performChromedpOAuth(
// contextId when the login took too long).
if parsed, perr := url.Parse(reqURL); perr == nil {
flowError = friendlyOPError(parsed.Query().Get("code"), parsed.Query().Get("message"))
- log.Printf("[%s] Stellantis error page: %s", requestID, reqURL)
+ logStellantisErrorPage(requestID, reqURL)
}
} else if debug != nil && isRelevantURL(reqURL) {
// Only show relevant OAuth flow URLs in debug output
@@ -452,7 +451,7 @@ func friendlyOPError(code, message string) string {
func codeFromLocation(browserCtx context.Context, redirectPrefix string) string {
var currentURL string
_ = chromedp.Run(browserCtx, chromedp.Location(&currentURL))
- log.Printf("Current URL: %s", currentURL)
+ logCurrentBrowserLocation(currentURL)
if !strings.HasPrefix(currentURL, redirectPrefix) {
return ""
}
@@ -462,3 +461,15 @@ func codeFromLocation(browserCtx context.Context, redirectPrefix string) string
}
return parsed.Query().Get("code")
}
+
+func logCapturedOAuthRedirect(requestID, _ string) {
+ log.Printf("[%s] Captured OAuth redirect", requestID)
+}
+
+func logStellantisErrorPage(requestID, _ string) {
+ log.Printf("[%s] Stellantis error page received", requestID)
+}
+
+func logCurrentBrowserLocation(_ string) {
+ log.Printf("Current browser location checked")
+}
diff --git a/internal/app/oauth_test.go b/internal/app/oauth_test.go
index 3589572..b9ad90d 100644
--- a/internal/app/oauth_test.go
+++ b/internal/app/oauth_test.go
@@ -6,6 +6,32 @@ import (
"testing"
)
+func TestOAuthEventLogsDoNotContainURLsCodesOrTokens(t *testing.T) {
+ cookie := "event-cookie"
+ accessToken := "event-access-token"
+ refreshToken := "event-refresh-token"
+ redirectURL := "mypeugeot://oauth2redirect/dk?code=event-oauth-code&access_token=" + accessToken
+ errorURL := "https://idpcvs.peugeot.com/am/oauth2/OPErrorPage.php?cookie=" + cookie
+ currentURL := redirectURL + "&refresh_token=" + refreshToken
+ oauthCode := "event-oauth-code"
+
+ logs := captureLogs(t, func() {
+ logCapturedOAuthRedirect("redirect-request-id", redirectURL)
+ logStellantisErrorPage("error-request-id", errorURL)
+ logCurrentBrowserLocation(currentURL)
+ })
+ want := "[redirect-request-id] Captured OAuth redirect\n" +
+ "[error-request-id] Stellantis error page received\n" +
+ "Current browser location checked\n"
+ if logs != want {
+ t.Fatalf("logs = %q, want %q", logs, want)
+ }
+ assertLogRedacted(t, logs,
+ cookie, accessToken, refreshToken,
+ redirectURL, errorURL, currentURL, oauthCode,
+ )
+}
+
func TestFriendlyOPError(t *testing.T) {
cases := []struct {
name string
diff --git a/internal/app/server.go b/internal/app/server.go
index 48a487e..946d8cf 100644
--- a/internal/app/server.go
+++ b/internal/app/server.go
@@ -121,6 +121,13 @@ func getClientIP(r *http.Request) string {
return r.RemoteAddr
}
+func remoteClientIP(r *http.Request) string {
+ if host, _, err := net.SplitHostPort(r.RemoteAddr); err == nil {
+ return host
+ }
+ return strings.TrimSpace(r.RemoteAddr)
+}
+
// refundIfExpired gives back the rate-limit charge when the OAuth attempt
// failed with a transient "session expired" error (bounded by the limiter).
func refundIfExpired(clientIP, requestID string, err error) {
@@ -129,6 +136,14 @@ func refundIfExpired(clientIP, requestID string, err error) {
}
}
+func logOAuthRequest(requestID, clientIP string, req OAuthRequest) {
+ log.Printf("[%s] OAuth request from %s (%s/%s)", requestID, clientIP, req.Brand, req.Country)
+}
+
+func logOAuthFailure(requestID string, _ error) {
+ log.Printf("[%s] OAuth failed", requestID)
+}
+
func handleOAuth(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
sendError(w, "Method not allowed", http.StatusMethodNotAllowed)
@@ -136,7 +151,7 @@ func handleOAuth(w http.ResponseWriter, r *http.Request) {
}
// Get client IP early for rate limiting
- clientIP := getClientIP(r)
+ clientIP := remoteClientIP(r)
// Check rate limit
if !rateLimiter.isAllowed(clientIP) {
@@ -160,7 +175,7 @@ func handleOAuth(w http.ResponseWriter, r *http.Request) {
// Generate request ID
requestID := uuid.New().String()
- log.Printf("[%s] OAuth request from %s for user %s (%s/%s)", requestID, clientIP, req.Email, req.Brand, req.Country)
+ logOAuthRequest(requestID, clientIP, req)
// Check if client accepts SSE
if r.Header.Get("Accept") == "text/event-stream" {
@@ -171,7 +186,7 @@ func handleOAuth(w http.ResponseWriter, r *http.Request) {
code, err := performOAuth(req, requestID, nil, nil)
if err != nil {
refundIfExpired(clientIP, requestID, err)
- log.Printf("[%s] OAuth failed: %s", requestID, err.Error())
+ logOAuthFailure(requestID, err)
sendError(w, err.Error(), http.StatusBadRequest)
return
}
@@ -206,7 +221,7 @@ func handleOAuthSSE(w http.ResponseWriter, req OAuthRequest, requestID, clientIP
code, err := performOAuth(req, requestID, progress, debug)
if err != nil {
refundIfExpired(clientIP, requestID, err)
- log.Printf("[%s] OAuth failed: %s", requestID, err.Error())
+ logOAuthFailure(requestID, err)
_, _ = fmt.Fprintf(w, "data: {\"type\":\"error\",\"message\":\"%s\"}\n\n", err.Error())
flusher.Flush()
return
diff --git a/internal/app/server_test.go b/internal/app/server_test.go
index b56bb27..730e658 100644
--- a/internal/app/server_test.go
+++ b/internal/app/server_test.go
@@ -1,12 +1,16 @@
package app
import (
+ "bytes"
"encoding/json"
+ "errors"
+ "log"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing"
+ "time"
)
func TestMain(m *testing.M) {
@@ -101,6 +105,47 @@ func TestHandleOAuth_InvalidBody(t *testing.T) {
}
}
+func useSingleRequestRateLimiter(t *testing.T) {
+ t.Helper()
+ previous := rateLimiter
+ rateLimiter = &RateLimiter{
+ requests: make(map[string][]time.Time),
+ refunds: make(map[string][]time.Time),
+ limit: 1,
+ window: time.Hour,
+ enabled: true,
+ }
+ t.Cleanup(func() { rateLimiter = previous })
+}
+
+func setSpoofedClientHeaders(req *http.Request, suffix string) {
+ req.Header.Set("Forwarded", "for=203.0.113."+suffix)
+ req.Header.Set("X-Forwarded-For", "198.51.100."+suffix)
+ req.Header.Set("X-Real-IP", "192.0.2."+suffix)
+}
+
+func TestHandleOAuthRateLimitUsesDirectPeer(t *testing.T) {
+ useSingleRequestRateLimiter(t)
+
+ first := httptest.NewRequest(http.MethodPost, "/oauth", strings.NewReader("invalid json"))
+ first.RemoteAddr = "10.0.0.8:41001"
+ setSpoofedClientHeaders(first, "11")
+ firstResponse := httptest.NewRecorder()
+ handleOAuth(firstResponse, first)
+ if firstResponse.Code != http.StatusBadRequest {
+ t.Fatalf("first status = %d, want %d", firstResponse.Code, http.StatusBadRequest)
+ }
+
+ second := httptest.NewRequest(http.MethodPost, "/oauth", strings.NewReader("invalid json"))
+ second.RemoteAddr = "10.0.0.8:41001"
+ setSpoofedClientHeaders(second, "22")
+ secondResponse := httptest.NewRecorder()
+ handleOAuth(secondResponse, second)
+ if secondResponse.Code != http.StatusTooManyRequests {
+ t.Fatalf("second status = %d, want %d", secondResponse.Code, http.StatusTooManyRequests)
+ }
+}
+
func TestHandleOAuth_MissingFields(t *testing.T) {
body := `{"brand":"MyPeugeot","country":"","email":"","password":""}`
req := httptest.NewRequest(http.MethodPost, "/oauth", strings.NewReader(body))
@@ -188,3 +233,74 @@ func TestParseClientIP(t *testing.T) {
}
}
}
+
+func TestRemoteClientIPIgnoresForwardedHeaders(t *testing.T) {
+ req := httptest.NewRequest(http.MethodPost, "/oauth", nil)
+ req.RemoteAddr = "192.0.2.10:1234"
+ req.Header.Set("Forwarded", "for=203.0.113.1")
+ req.Header.Set("X-Forwarded-For", "203.0.113.2")
+ req.Header.Set("X-Real-IP", "203.0.113.3")
+
+ if got := remoteClientIP(req); got != "192.0.2.10" {
+ t.Errorf("remoteClientIP() = %q, want %q", got, "192.0.2.10")
+ }
+}
+
+func captureLogs(t *testing.T, fn func()) string {
+ t.Helper()
+ var logs bytes.Buffer
+ originalWriter := log.Writer()
+ originalFlags := log.Flags()
+ log.SetOutput(&logs)
+ log.SetFlags(0)
+ defer func() {
+ log.SetOutput(originalWriter)
+ log.SetFlags(originalFlags)
+ }()
+
+ fn()
+ return logs.String()
+}
+
+func assertLogRedacted(t *testing.T, logs string, sentinels ...string) {
+ t.Helper()
+ for _, sentinel := range sentinels {
+ if strings.Contains(logs, sentinel) {
+ t.Errorf("logs contain sensitive value %q: %s", sentinel, logs)
+ }
+ }
+}
+
+func TestOAuthRequestAndFailureLogsDoNotContainSensitiveValues(t *testing.T) {
+ email := "sentinel-email@example.invalid"
+ password := "sentinel-password"
+ cookie := "sentinel-cookie"
+ accessToken := "sentinel-access-token"
+ refreshToken := "sentinel-refresh-token"
+ authorizeURL := "https://idpcvs.peugeot.com/am/oauth2/authorize?sentinel=authorize-url"
+ redirectURL := "mypeugeot://oauth2redirect/dk?sentinel=redirect-url"
+ oauthCode := "sentinel-oauth-code"
+ req := OAuthRequest{
+ Brand: "MyPeugeot",
+ Country: "DK",
+ Email: email,
+ Password: password,
+ }
+ failure := errors.New(strings.Join([]string{
+ cookie, accessToken, refreshToken, authorizeURL, redirectURL, oauthCode,
+ }, "|"))
+
+ logs := captureLogs(t, func() {
+ logOAuthRequest("request-id", "192.0.2.10", req)
+ logOAuthFailure("request-id", failure)
+ })
+ want := "[request-id] OAuth request from 192.0.2.10 (MyPeugeot/DK)\n" +
+ "[request-id] OAuth failed\n"
+ if logs != want {
+ t.Fatalf("logs = %q, want %q", logs, want)
+ }
+ assertLogRedacted(t, logs,
+ email, password, cookie, accessToken, refreshToken,
+ authorizeURL, redirectURL, oauthCode,
+ )
+}
diff --git a/internal/app/worker.go b/internal/app/worker.go
index 00b9cf8..3dbe54a 100644
--- a/internal/app/worker.go
+++ b/internal/app/worker.go
@@ -7,9 +7,20 @@ import (
"log"
"net/http"
"net/url"
+ "strings"
"uuid"
)
+const maxWorkerRequestBytes = 64 << 10
+
+var approvedWorkerHosts = map[string]struct{}{
+ "idpcvs.citroen.com": {},
+ "idpcvs.driveds.com": {},
+ "idpcvs.opel.com": {},
+ "idpcvs.peugeot.com": {},
+ "idpcvs.vauxhall.co.uk": {},
+}
+
// workerRequest is the worker-v2-compatible OAuth request: the caller supplies a
// fully-built authorize URL instead of brand/country. This mirrors the API of
// github.com/andreadegiovine/homeassistant-stellantis-vehicles-worker-v2 so the
@@ -40,15 +51,21 @@ func handleWorker(w http.ResponseWriter, r *http.Request) {
return
}
- clientIP := getClientIP(r)
+ clientIP := remoteClientIP(r)
if !rateLimiter.isAllowed(clientIP) {
log.Printf("Rate limit exceeded for %s (worker)", clientIP)
sendWorkerError(w, "Rate limit exceeded. Try again later.", http.StatusTooManyRequests)
return
}
+ r.Body = http.MaxBytesReader(w, r.Body, maxWorkerRequestBytes)
var req workerRequest
if err := json.UnmarshalRead(r.Body, &req); err != nil {
+ var maxBytesErr *http.MaxBytesError
+ if errors.As(err, &maxBytesErr) {
+ sendWorkerError(w, "Request body too large", http.StatusRequestEntityTooLarge)
+ return
+ }
sendWorkerError(w, "Invalid request body", http.StatusBadRequest)
return
}
@@ -58,6 +75,11 @@ func handleWorker(w http.ResponseWriter, r *http.Request) {
return
}
+ if err := validateWorkerURL(req.URL); err != nil {
+ sendWorkerError(w, err.Error(), http.StatusBadRequest)
+ return
+ }
+
scheme, err := redirectScheme(req.URL)
if err != nil {
sendWorkerError(w, err.Error(), http.StatusBadRequest)
@@ -65,12 +87,12 @@ func handleWorker(w http.ResponseWriter, r *http.Request) {
}
requestID := uuid.New().String()
- log.Printf("[%s] worker OAuth request from %s for user %s", requestID, clientIP, req.Email)
+ logWorkerOAuthRequest(requestID, clientIP, req)
code, err := performChromedpOAuth(req.URL, req.Email, req.Password, scheme, requestID, nil, nil)
if err != nil {
refundIfExpired(clientIP, requestID, err)
- log.Printf("[%s] worker OAuth failed: %s", requestID, err.Error())
+ logWorkerOAuthFailure(requestID, err)
sendWorkerError(w, err.Error(), http.StatusBadRequest)
return
}
@@ -80,6 +102,31 @@ func handleWorker(w http.ResponseWriter, r *http.Request) {
_ = json.MarshalWrite(w, workerCode{Code: code})
}
+func validateWorkerURL(rawURL string) error {
+ parsed, err := url.Parse(rawURL)
+ if err != nil {
+ return errors.New("invalid url")
+ }
+ if parsed.Scheme != "https" || parsed.User != nil || parsed.Port() != "" || parsed.Fragment != "" {
+ return errors.New("invalid url")
+ }
+ if _, ok := approvedWorkerHosts[strings.ToLower(parsed.Hostname())]; !ok {
+ return errors.New("invalid url")
+ }
+ if parsed.EscapedPath() != "/am/oauth2/authorize" {
+ return errors.New("invalid url")
+ }
+ return nil
+}
+
+func logWorkerOAuthRequest(requestID, clientIP string, _ workerRequest) {
+ log.Printf("[%s] worker OAuth request from %s", requestID, clientIP)
+}
+
+func logWorkerOAuthFailure(requestID string, _ error) {
+ log.Printf("[%s] worker OAuth failed", requestID)
+}
+
// redirectScheme extracts the custom redirect scheme (e.g. "mymap") from the
// authorize URL's redirect_uri query parameter; the code-capture listener keys on
// "<scheme>://".
diff --git a/internal/app/worker_test.go b/internal/app/worker_test.go
index e64964d..7141640 100644
--- a/internal/app/worker_test.go
+++ b/internal/app/worker_test.go
@@ -2,12 +2,81 @@ package app
import (
"encoding/json"
+ "errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
)
+func TestWorkerRequestAndFailureLogsDoNotContainSensitiveValues(t *testing.T) {
+ email := "worker-email@example.invalid"
+ password := "worker-password"
+ cookie := "worker-cookie"
+ accessToken := "worker-access-token"
+ refreshToken := "worker-refresh-token"
+ authorizeURL := "https://idpcvs.peugeot.com/am/oauth2/authorize?sentinel=worker-authorize-url"
+ redirectURL := "mypeugeot://oauth2redirect/dk?sentinel=worker-redirect-url"
+ oauthCode := "worker-oauth-code"
+ req := workerRequest{
+ URL: authorizeURL,
+ Email: email,
+ Password: password,
+ }
+ failure := errors.New(strings.Join([]string{
+ cookie, accessToken, refreshToken, redirectURL, oauthCode,
+ }, "|"))
+
+ logs := captureLogs(t, func() {
+ logWorkerOAuthRequest("worker-request-id", "192.0.2.20", req)
+ logWorkerOAuthFailure("worker-request-id", failure)
+ })
+ want := "[worker-request-id] worker OAuth request from 192.0.2.20\n" +
+ "[worker-request-id] worker OAuth failed\n"
+ if logs != want {
+ t.Fatalf("logs = %q, want %q", logs, want)
+ }
+ assertLogRedacted(t, logs,
+ email, password, cookie, accessToken, refreshToken,
+ authorizeURL, redirectURL, oauthCode,
+ )
+}
+
+func TestValidateWorkerURL(t *testing.T) {
+ allowed := []string{
+ "https://idpcvs.citroen.com/am/oauth2/authorize?redirect_uri=mycitroen%3A%2F%2Foauth2redirect%2Fdk",
+ "https://idpcvs.driveds.com/am/oauth2/authorize?redirect_uri=mymap%3A%2F%2Foauth2redirect%2Fdk",
+ "https://idpcvs.opel.com/am/oauth2/authorize?redirect_uri=myopel%3A%2F%2Foauth2redirect%2Fdk",
+ "https://idpcvs.peugeot.com/am/oauth2/authorize?redirect_uri=mypeugeot%3A%2F%2Foauth2redirect%2Fdk",
+ "https://idpcvs.vauxhall.co.uk/am/oauth2/authorize?redirect_uri=myvauxhall%3A%2F%2Foauth2redirect%2Fgb",
+ "https://IDPCVS.OPEL.COM/am/oauth2/authorize?redirect_uri=myopel%3A%2F%2Foauth2redirect%2Fdk",
+ }
+ for _, rawURL := range allowed {
+ t.Run(rawURL, func(t *testing.T) {
+ if err := validateWorkerURL(rawURL); err != nil {
+ t.Fatalf("validateWorkerURL() error = %v", err)
+ }
+ })
+ }
+
+ rejected := map[string]string{
+ "http": "http://idpcvs.peugeot.com/am/oauth2/authorize?redirect_uri=mypeugeot%3A%2F%2Foauth2redirect%2Fdk",
+ "userinfo": "https://user@idpcvs.peugeot.com/am/oauth2/authorize?redirect_uri=mypeugeot%3A%2F%2Foauth2redirect%2Fdk",
+ "explicit port": "https://idpcvs.peugeot.com:443/am/oauth2/authorize?redirect_uri=mypeugeot%3A%2F%2Foauth2redirect%2Fdk",
+ "evil suffix": "https://idpcvs.peugeot.com.evil.example/am/oauth2/authorize?redirect_uri=mypeugeot%3A%2F%2Foauth2redirect%2Fdk",
+ "extra path component": "https://idpcvs.peugeot.com/am/oauth2/authorize/extra?redirect_uri=mypeugeot%3A%2F%2Foauth2redirect%2Fdk",
+ "escaped path": "https://idpcvs.peugeot.com/am/oauth2%2Fauthorize?redirect_uri=mypeugeot%3A%2F%2Foauth2redirect%2Fdk",
+ "fragment": "https://idpcvs.peugeot.com/am/oauth2/authorize?redirect_uri=mypeugeot%3A%2F%2Foauth2redirect%2Fdk#fragment",
+ }
+ for name, rawURL := range rejected {
+ t.Run(name, func(t *testing.T) {
+ if err := validateWorkerURL(rawURL); err == nil {
+ t.Fatal("validateWorkerURL() error = nil, want rejection")
+ }
+ })
+ }
+}
+
func TestHandleWorker_MethodNotAllowed(t *testing.T) {
req := httptest.NewRequest(http.MethodGet, "/worker", nil)
w := httptest.NewRecorder()
@@ -30,6 +99,47 @@ func TestHandleWorker_InvalidBody(t *testing.T) {
}
}
+func TestHandleWorkerRateLimitUsesDirectPeer(t *testing.T) {
+ useSingleRequestRateLimiter(t)
+
+ first := httptest.NewRequest(http.MethodPost, "/worker", strings.NewReader("invalid json"))
+ first.RemoteAddr = "10.0.0.9:41002"
+ setSpoofedClientHeaders(first, "33")
+ firstResponse := httptest.NewRecorder()
+ handleWorker(firstResponse, first)
+ if firstResponse.Code != http.StatusBadRequest {
+ t.Fatalf("first status = %d, want %d", firstResponse.Code, http.StatusBadRequest)
+ }
+
+ second := httptest.NewRequest(http.MethodPost, "/worker", strings.NewReader("invalid json"))
+ second.RemoteAddr = "10.0.0.9:41002"
+ setSpoofedClientHeaders(second, "44")
+ secondResponse := httptest.NewRecorder()
+ handleWorker(secondResponse, second)
+ if secondResponse.Code != http.StatusTooManyRequests {
+ t.Fatalf("second status = %d, want %d", secondResponse.Code, http.StatusTooManyRequests)
+ }
+}
+
+func TestHandleWorker_RequestTooLarge(t *testing.T) {
+ body := `{"url":"` + strings.Repeat("x", 65<<10) + `"}`
+ req := httptest.NewRequest(http.MethodPost, "/worker", strings.NewReader(body))
+ w := httptest.NewRecorder()
+
+ handleWorker(w, req)
+
+ if w.Code != http.StatusRequestEntityTooLarge {
+ t.Fatalf("status = %d, want %d", w.Code, http.StatusRequestEntityTooLarge)
+ }
+ var resp workerError
+ if err := json.NewDecoder(w.Body).Decode(&resp); err != nil {
+ t.Fatalf("decode response: %v", err)
+ }
+ if resp.Message != "Request body too large" {
+ t.Errorf("message = %q, want %q", resp.Message, "Request body too large")
+ }
+}
+
func TestHandleWorker_MissingParams(t *testing.T) {
body := `{"url":"","email":"","password":""}`
req := httptest.NewRequest(http.MethodPost, "/worker", strings.NewReader(body))
@@ -83,7 +193,7 @@ func TestRedirectScheme(t *testing.T) {
},
{
name: "opel custom scheme",
- authURL: "https://example.com/authorize?redirect_uri=myopel%3A%2F%2Foauth2redirect%2Fde",
+ authURL: "https://idpcvs.opel.com/am/oauth2/authorize?redirect_uri=myopel%3A%2F%2Foauth2redirect%2Fde",
want: "myopel",
},
{
+438
View File
@@ -0,0 +1,438 @@
#!/usr/bin/env python3
from __future__ import annotations
import dataclasses
import json
import logging
import os
import re
import shutil
import signal
import subprocess
import sys
import time
import urllib.parse
import urllib.request
from pathlib import Path
from typing import Callable, NamedTuple
OPTIONS_PATH = Path("/data/options.json")
PROFILE_PATH = Path("/tmp/cloakserve")
CLOAK_ROOT = "http://127.0.0.1:9222/"
CLOAK_VERSION = "http://127.0.0.1:9222/json/version?fingerprint=addon-readiness"
CLOAK_CLOSE = "http://127.0.0.1:9222/fingerprint/addon-readiness/close"
STELLOAUTH_ROOT = "http://127.0.0.1:8080/"
DURATION = re.compile(r"^[1-9][0-9]*(?:ms|s|m|h)$")
CLOAK_COMMAND = [
"/usr/local/bin/cloakserve",
"--headless=true",
"--idle-timeout=30",
"--data-dir=/tmp/cloakserve",
]
STELLOAUTH_COMMAND = ["/usr/local/bin/stelloauth"]
SHUTDOWN_GRACE_SECONDS = 9.0
class ConfigError(RuntimeError):
pass
@dataclasses.dataclass(frozen=True)
class Options:
queue_timeout: str
rate_limit_count: int
rate_limit_duration: str
class HttpResponse(NamedTuple):
status: int
body: bytes
class ReadinessError(RuntimeError):
pass
def load_options(path: Path) -> Options:
try:
value = json.loads(path.read_text(encoding="utf-8"))
if not isinstance(value, dict) or set(value) != {
"queue_timeout",
"rate_limit_count",
"rate_limit_duration",
}:
raise ValueError
queue_timeout = value["queue_timeout"]
rate_limit_count = value["rate_limit_count"]
rate_limit_duration = value["rate_limit_duration"]
if not isinstance(queue_timeout, str) or DURATION.fullmatch(queue_timeout) is None:
raise ValueError
if (
isinstance(rate_limit_count, bool)
or not isinstance(rate_limit_count, int)
or not 1 <= rate_limit_count <= 20
):
raise ValueError
if (
not isinstance(rate_limit_duration, str)
or DURATION.fullmatch(rate_limit_duration) is None
):
raise ValueError
return Options(queue_timeout, rate_limit_count, rate_limit_duration)
except (OSError, UnicodeError, json.JSONDecodeError, KeyError, TypeError, ValueError):
raise ConfigError("Invalid add-on configuration") from None
def build_environment(options: Options) -> dict[str, str]:
environment = os.environ.copy()
environment.update(
{
"CLOAK_CDP_URL": "http://127.0.0.1:9222",
"CLOAK_MAX_SESSIONS": "1",
"CLOAK_QUEUE_TIMEOUT": options.queue_timeout,
"RATE_LIMIT_COUNT": str(options.rate_limit_count),
"RATE_LIMIT_DURATION": options.rate_limit_duration,
"HTTP_ADDRESS": "0.0.0.0",
"PORT": "8080",
"METRICS_ADDRESS": "127.0.0.1",
"METRICS_PORT": "9090",
}
)
return environment
def http_request(method: str, url: str, timeout: float) -> HttpResponse:
opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
data = b"" if method == "POST" else None
request = urllib.request.Request(url, data=data, method=method)
with opener.open(request, timeout=timeout) as response:
return HttpResponse(response.status, response.read())
def probe_cloak(
request: Callable[[str, str, float], HttpResponse] = http_request,
) -> None:
if request("GET", CLOAK_ROOT, 2.0).status != 200:
raise ReadinessError("CloakBrowser root probe failed")
version = request("GET", CLOAK_VERSION, 2.0)
if version.status != 200:
raise ReadinessError("CloakBrowser CDP probe failed")
try:
document = json.loads(version.body)
if not isinstance(document, dict):
raise ValueError
websocket_url = document["webSocketDebuggerUrl"]
if (
not isinstance(websocket_url, str)
or not websocket_url.startswith("ws://127.0.0.1:")
):
raise ValueError
parsed = urllib.parse.urlsplit(websocket_url)
if (
parsed.scheme != "ws"
or parsed.hostname != "127.0.0.1"
or parsed.username is not None
or parsed.password is not None
or parsed.port is None
):
raise ValueError
except (json.JSONDecodeError, KeyError, TypeError, UnicodeError, ValueError):
raise ReadinessError("CloakBrowser CDP response invalid") from None
if request("POST", CLOAK_CLOSE, 2.0).status != 200:
raise ReadinessError("CloakBrowser readiness profile close failed")
def probe_stelloauth(
request: Callable[[str, str, float], HttpResponse] = http_request,
) -> None:
if request("GET", STELLOAUTH_ROOT, 2.0).status != 200:
raise ReadinessError("Stelloauth root probe failed")
def wait_until_ready(
name: str,
probe: Callable[[], None],
timeout: float,
stopping: Callable[[], bool],
monotonic: Callable[[], float] = time.monotonic,
sleep: Callable[[float], None] = time.sleep,
) -> None:
deadline = monotonic() + timeout
delay = 0.25
while monotonic() < deadline and not stopping():
try:
probe()
return
except (OSError, ValueError, ReadinessError):
sleep(delay)
delay = min(delay * 2, 2.0)
raise ReadinessError(f"{name} did not become ready")
def process_group_alive(
pgid: int,
killpg: Callable[[int, int], None] = os.killpg,
) -> bool:
try:
killpg(pgid, 0)
except ProcessLookupError:
return False
except PermissionError:
return True
except OSError:
return True
return True
class ProcessManager:
def __init__(
self,
environment: dict[str, str],
*,
popen: Callable[..., subprocess.Popen] = subprocess.Popen,
killpg: Callable[[int, int], None] = os.killpg,
group_alive: Callable[[int], bool] = process_group_alive,
cloak_probe: Callable[[], None] = probe_cloak,
stelloauth_probe: Callable[[], None] = probe_stelloauth,
monotonic: Callable[[], float] = time.monotonic,
sleep: Callable[[float], None] = time.sleep,
profile_path: Path = PROFILE_PATH,
) -> None:
self.environment = environment
self._popen = popen
self._killpg = killpg
self._group_alive = group_alive
self._cloak_probe = cloak_probe
self._stelloauth_probe = stelloauth_probe
self._monotonic = monotonic
self._sleep = sleep
self._profile_path = profile_path
self._children: list[tuple[str, subprocess.Popen]] = []
self._reaped_children: set[int] = set()
self._groups: list[int] = []
self._dead_groups: set[int] = set()
self._stopping = False
self._term_sent = False
self._shutdown_deadline: float | None = None
def _handle_signal(self, _signum: int, _frame: object) -> None:
if self._stopping:
return
logging.info("Shutdown requested")
self._stopping = True
self._begin_shutdown()
def _clean_profiles(self) -> None:
logging.info("Cleaning CloakBrowser profiles")
if self._profile_path.is_symlink():
self._profile_path.unlink()
elif self._profile_path.exists():
shutil.rmtree(self._profile_path)
self._profile_path.mkdir(parents=True, mode=0o700)
self._profile_path.chmod(0o700)
def _spawn(self, name: str, command: list[str]) -> subprocess.Popen:
process = self._popen(
command,
env=self.environment,
start_new_session=True,
)
self._children.append((name, process))
self._groups.append(process.pid)
if self._stopping:
self._signal_groups(signal.SIGTERM, [process.pid])
return process
def _first_exited_child(self) -> tuple[str, int] | None:
for name, process in self._children:
returncode = process.poll()
if returncode is not None:
return name, returncode
return None
def _startup_stopping(self) -> bool:
return self._stopping or self._first_exited_child() is not None
def _reap_exited_children(self) -> None:
for _name, process in self._children:
if process.pid in self._reaped_children or process.poll() is None:
continue
try:
process.wait(timeout=None)
except OSError:
logging.error("Child process reap failed")
else:
self._reaped_children.add(process.pid)
def _live_groups(self) -> list[int]:
self._reap_exited_children()
live_groups = []
for pgid in self._groups:
if pgid in self._dead_groups:
continue
if self._group_alive(pgid):
live_groups.append(pgid)
else:
self._dead_groups.add(pgid)
return live_groups
def _signal_groups(
self, sent_signal: int, groups: list[int] | None = None
) -> None:
for pgid in self._live_groups() if groups is None else groups:
if pgid in self._dead_groups:
continue
if groups is not None and not self._group_alive(pgid):
self._dead_groups.add(pgid)
continue
try:
self._killpg(pgid, sent_signal)
except ProcessLookupError:
self._dead_groups.add(pgid)
except OSError:
pass
def _begin_shutdown(self) -> None:
if self._shutdown_deadline is None:
self._shutdown_deadline = self._monotonic() + SHUTDOWN_GRACE_SECONDS
if not self._term_sent:
self._signal_groups(signal.SIGTERM)
self._term_sent = True
def _reap_all(self) -> None:
for _name, process in self._children:
if process.pid in self._reaped_children:
continue
try:
process.wait(timeout=None)
except OSError:
logging.error("Child process reap failed")
else:
self._reaped_children.add(process.pid)
def _shutdown(self) -> None:
if not self._children:
return
logging.info("Stopping child processes")
self._begin_shutdown()
assert self._shutdown_deadline is not None
live_groups = self._live_groups()
while live_groups and self._monotonic() < self._shutdown_deadline:
remaining = self._shutdown_deadline - self._monotonic()
self._sleep(min(0.25, max(0.0, remaining)))
live_groups = self._live_groups()
if live_groups:
logging.info("Forcing child processes to stop")
self._signal_groups(signal.SIGKILL, live_groups)
self._reap_all()
logging.info("Child processes stopped")
@staticmethod
def _failure_status(returncode: int) -> int:
return returncode if returncode != 0 else 1
def _startup_failure(self, service: str) -> int:
exited = self._first_exited_child()
if exited is not None:
name, returncode = exited
logging.error(f"{name} exited unexpectedly")
status = self._failure_status(returncode)
else:
logging.error(f"{service} readiness failed")
status = 1
self._shutdown()
return status
def _run(self) -> int:
try:
self._clean_profiles()
if self._stopping:
self._shutdown()
return 0
logging.info("Starting CloakBrowser")
self._spawn("CloakBrowser", CLOAK_COMMAND)
wait_until_ready(
"CloakBrowser",
self._cloak_probe,
60.0,
self._startup_stopping,
self._monotonic,
self._sleep,
)
if self._stopping:
self._shutdown()
return 0
if self._first_exited_child() is not None:
return self._startup_failure("CloakBrowser")
logging.info("CloakBrowser ready")
logging.info("Starting Stelloauth")
self._spawn("Stelloauth", STELLOAUTH_COMMAND)
wait_until_ready(
"Stelloauth",
self._stelloauth_probe,
30.0,
self._startup_stopping,
self._monotonic,
self._sleep,
)
if self._stopping:
self._shutdown()
return 0
if self._first_exited_child() is not None:
return self._startup_failure("Stelloauth")
logging.info("Stelloauth listening on 0.0.0.0:8080")
except ReadinessError:
if self._stopping:
self._shutdown()
return 0
service = "Stelloauth" if len(self._children) > 1 else "CloakBrowser"
return self._startup_failure(service)
except OSError:
logging.error("Process startup failed")
self._shutdown()
return 1
while not self._stopping:
exited = self._first_exited_child()
if exited is not None:
name, returncode = exited
logging.error(f"{name} exited unexpectedly")
status = self._failure_status(returncode)
self._shutdown()
return status
self._sleep(0.25)
self._shutdown()
return 0
def run(self) -> int:
previous_handlers = {
signal.SIGTERM: signal.signal(signal.SIGTERM, self._handle_signal),
signal.SIGINT: signal.signal(signal.SIGINT, self._handle_signal),
}
try:
return self._run()
finally:
for handled_signal, previous_handler in previous_handlers.items():
signal.signal(handled_signal, previous_handler)
def main() -> int:
logging.basicConfig(level=logging.INFO, format="%(message)s")
try:
options = load_options(OPTIONS_PATH)
except ConfigError:
logging.error("Invalid add-on configuration")
return 2
return ProcessManager(build_environment(options)).run()
if __name__ == "__main__":
sys.exit(main())
+10
View File
@@ -0,0 +1,10 @@
configuration:
queue_timeout:
name: Køventetid
description: Angiver hvor længe et loginforsøg må vente i køen.
rate_limit_count:
name: Loginforsøg
description: Angiver det maksimale antal loginforsøg i hver periode.
rate_limit_duration:
name: Rate limit-periode
description: Angiver periodens længde for begrænsning af loginforsøg.
+10
View File
@@ -0,0 +1,10 @@
configuration:
queue_timeout:
name: Queue timeout
description: Sets how long a login attempt may wait in the queue.
rate_limit_count:
name: Login attempts
description: Sets the maximum number of login attempts in each period.
rate_limit_duration:
name: Rate limit period
description: Sets the length of the login attempt rate limit period.
+15
View File
@@ -0,0 +1,15 @@
FROM golang:1.27.1-bookworm@sha256:69a7b9788769bec032d238959b61854e9ae87f57be9029ec04e9885fabf99195
RUN apt-get update \
&& apt-get install --yes --no-install-recommends python3 python3-venv
WORKDIR /workspace
COPY requirements-dev.txt ./
RUN python3 -m venv /opt/venv \
&& /opt/venv/bin/pip install --no-cache-dir --requirement requirements-dev.txt
COPY . ./
RUN git init
ENTRYPOINT ["/opt/venv/bin/pytest"]
+330
View File
@@ -0,0 +1,330 @@
from __future__ import annotations
import hashlib
import os
import re
import subprocess
from pathlib import Path
import pytest
import yaml
ROOT = Path(__file__).parents[1]
REPOSITORY_URL = "https://git.radixadm.dk/dennis/homeassistant-stelloauth-addon.git"
def load_yaml(path: str) -> dict:
with (ROOT / path).open(encoding="utf-8") as handle:
value = yaml.safe_load(handle)
assert isinstance(value, dict)
return value
def test_repository_metadata() -> None:
metadata = load_yaml("repository.yaml")
assert metadata["url"] == REPOSITORY_URL
assert metadata["name"] == "Stelloauth for Home Assistant"
assert metadata["maintainer"] == "Dennis / Radix ApS"
def test_addon_contract() -> None:
config = load_yaml("stelloauth/config.yaml")
assert config["slug"] == "stelloauth"
assert config["version"] == "0.1.0"
assert config["arch"] == ["amd64", "aarch64"]
assert config["startup"] == "application"
assert config["boot"] == "auto"
assert config["init"] is True
assert config["watchdog"] == "http://[HOST]:[PORT:8080]/"
assert config["ports"] == {"8080/tcp": None}
for key in ("ingress", "host_network", "privileged", "full_access", "docker_api", "devices", "map"):
assert key not in config
def test_defaults_and_schema_are_aligned() -> None:
config = load_yaml("stelloauth/config.yaml")
assert config["options"] == {
"queue_timeout": "60s",
"rate_limit_count": 5,
"rate_limit_duration": "1h",
}
assert set(config["schema"]) == set(config["options"])
assert config["schema"]["rate_limit_count"] == "int(1,20)"
def test_repository_hostname_derivation() -> None:
repository_id = hashlib.sha1(REPOSITORY_URL.lower().encode()).hexdigest()[:8]
assert repository_id == "0031621f"
assert f"{repository_id}-stelloauth" == "0031621f-stelloauth"
def test_user_guide_documentation_contract() -> None:
documentation = (ROOT / "stelloauth/DOCS.md").read_text(encoding="utf-8")
for required_text in (
REPOSITORY_URL,
"Installér **Stelloauth**",
"aktivér **Start ved opstart** og **Watchdog**",
"http://0031621f-stelloauth:8080/worker",
"http://192.168.1.20:8080/worker",
"Brand: Opel",
"Country: DK",
"Image: 2,571,693,650 bytes (2.571 GB decimal / 2452.56 MiB).",
"Dokumenteret Task 4-måling: 121,5 MiB.",
"Stop: cirka 9.3 sekunder.",
"approximately 4 GB free",
"Der er ikke gennemført et live MyOpel-login.",
"Login-flow RAM: not measured without real MyOpel credentials.",
):
assert required_text in documentation
assert "deaktivér porttilknytningen igen" in documentation
assert "Seneste idle RAM" not in documentation
def test_root_readme_repository_source_and_license_facts() -> None:
readme = (ROOT / "README.md").read_text(encoding="utf-8")
for required_text in (
REPOSITORY_URL,
"v0.6.0",
"367d4f8c02a3b072c59142c49dffc129edc8548b",
"0.5.10",
"f04c23da285b3b3d3cf10c8f9d282e7adc1d52ce",
"CloakBrowser Binary License",
"MIT-licenseret",
):
assert required_text in readme
def test_translations_cover_every_option() -> None:
keys = set(load_yaml("stelloauth/config.yaml")["options"])
for language in ("da", "en"):
translation = load_yaml(f"stelloauth/translations/{language}.yaml")
assert set(translation["configuration"]) == keys
for entry in translation["configuration"].values():
assert set(entry) == {"name", "description"}
assert all(isinstance(value, str) and value.strip() for value in entry.values())
def test_dockerfile_uses_approved_pins_and_builds_patched_stelloauth() -> None:
dockerfile = (ROOT / "stelloauth/Dockerfile").read_text(encoding="utf-8")
assert (
"golang:1.27.1-bookworm@sha256:"
"69a7b9788769bec032d238959b61854e9ae87f57be9029ec04e9885fabf99195"
) in dockerfile
assert (
"cloakhq/cloakbrowser:0.5.10@sha256:"
"2ed5b2d047cbdde22cde7ef1a796526c716aadaa5bccbe1db5ade49282b64a76"
) in dockerfile
assert "367d4f8c02a3b072c59142c49dffc129edc8548b" in dockerfile
assert "go test ./..." in dockerfile
assert not re.search(r"^FROM\s+\S+:latest(?:\s|$)", dockerfile, re.MULTILINE)
def test_dockerfile_declares_home_assistant_runtime_contract() -> None:
dockerfile = (ROOT / "stelloauth/Dockerfile").read_text(encoding="utf-8")
for label in (
"io.hass.name",
"io.hass.description",
"io.hass.arch",
"io.hass.type",
"io.hass.version",
):
assert label in dockerfile
assert 'io.hass.type="app"' in dockerfile
assert 'io.hass.type="addon"' not in dockerfile
assert re.search(r"^EXPOSE 8080$", dockerfile, re.MULTILINE)
assert not re.search(r"^EXPOSE .*\b9222\b", dockerfile, re.MULTILINE)
assert "ENTRYPOINT []" in dockerfile
assert 'CMD ["/usr/local/bin/addon-supervisor"]' in dockerfile
def test_dockerfile_patches_parent_cloakserve_instead_of_copying_a_binary() -> None:
dockerfile = (ROOT / "stelloauth/Dockerfile").read_text(encoding="utf-8")
assert "COPY patches/cloakserve-loopback.patch" in dockerfile
for line in dockerfile.splitlines():
if line.lstrip().startswith("COPY "):
source = line.split()[1]
assert Path(source).name != "cloakserve"
def test_runtime_accepts_docker_port_unpublished_status() -> None:
runtime_test = (ROOT / "tests/test_runtime.sh").read_text(encoding="utf-8")
assert 'docker port "$container" 9222/tcp 2>/dev/null || true' in runtime_test
def test_runtime_copies_options_without_a_host_bind_mount() -> None:
runtime_test = (ROOT / "tests/test_runtime.sh").read_text(encoding="utf-8")
assert "docker create --name" in runtime_test
assert 'docker cp "$options_file" "$container:/data/options.json"' in runtime_test
assert 'docker start "$container"' in runtime_test
assert "--mount" not in runtime_test
def test_runtime_probes_service_inside_container_network_namespace() -> None:
runtime_test = (ROOT / "tests/test_runtime.sh").read_text(encoding="utf-8")
assert 'probe_root "$container"' in runtime_test
assert 'post_invalid_worker "$first_container"' in runtime_test
assert 'assert_8080_loopback_mapping "$first_container"' in runtime_test
assert 'assert_8080_loopback_mapping "$second_container"' in runtime_test
assert 'f"http://127.0.0.1:{sys.argv[1]}/"' not in runtime_test
def test_ci_builds_and_loads_only_the_amd64_test_image() -> None:
workflow = (ROOT / ".gitea/workflows/ci.yml").read_text(encoding="utf-8")
parsed = yaml.safe_load(workflow)
steps = parsed["jobs"]["validate"]["steps"]
architecture_command = next(
step["run"] for step in steps if step.get("name") == "Validate runner architecture"
)
assert architecture_command == 'test "$(uname -m)" = "x86_64"'
build_command = next(
step["run"] for step in steps if step.get("name") == "Build amd64 image"
)
assert build_command.split() == [
"docker",
"build",
"--build-arg",
"TARGETARCH=amd64",
"--build-arg",
"BUILD_ARCH=amd64",
"--tag",
"homeassistant-stelloauth-addon:test",
"stelloauth",
]
for publication_primitive in (
"--push",
"docker push",
"docker/login-action",
"docker/build-push-action",
"packages: write",
):
assert publication_primitive not in workflow
def test_ci_runs_validation_in_pinned_container() -> None:
workflow = (ROOT / ".gitea/workflows/ci.yml").read_text(encoding="utf-8")
parsed = yaml.safe_load(workflow)
steps = parsed["jobs"]["validate"]["steps"]
assert "actions/setup-python" not in workflow
assert "actions/setup-go" not in workflow
build_command = next(
step["run"] for step in steps if step.get("name") == "Build validation image"
)
assert build_command.split() == [
"docker",
"build",
"--file",
"tests/Dockerfile.ci",
"--tag",
"homeassistant-stelloauth-tests:test",
".",
]
validation_command = next(
step["run"]
for step in steps
if step.get("name") == "Validate metadata, patches, and process manager"
)
assert validation_command == "docker run --rm homeassistant-stelloauth-tests:test -q"
dockerfile = (ROOT / "tests/Dockerfile.ci").read_text(encoding="utf-8")
assert (
"golang:1.27.1-bookworm@sha256:"
"69a7b9788769bec032d238959b61854e9ae87f57be9029ec04e9885fabf99195"
) in dockerfile
assert "python3 python3-venv" in dockerfile
assert 'ENTRYPOINT ["/opt/venv/bin/pytest"]' in dockerfile
def test_patch_payloads_disable_git_whitespace_errors() -> None:
result = subprocess.run(
[
"git",
"check-attr",
"whitespace",
"--",
"stelloauth/patches/stelloauth-security.patch",
"stelloauth/patches/cloakserve-loopback.patch",
],
cwd=ROOT,
text=True,
capture_output=True,
check=True,
)
assert result.stdout.splitlines() == [
"stelloauth/patches/stelloauth-security.patch: whitespace: unset",
"stelloauth/patches/cloakserve-loopback.patch: whitespace: unset",
]
def test_runtime_rejects_every_ipv6_cdp_listener() -> None:
runtime_test = (ROOT / "tests/test_runtime.sh").read_text(encoding="utf-8")
assert 'cat /proc/net/tcp6 > "$tcp6_artifact"' in runtime_test
assert 'for table in ("/proc/net/tcp", "/proc/net/tcp6"):' in runtime_test
assert 'if table == "/proc/net/tcp6":' in runtime_test
def test_runtime_requires_process_baseline_after_cdp_close_and_zero_stopped_pid() -> None:
runtime_test = (ROOT / "tests/test_runtime.sh").read_text(encoding="utf-8")
assert 'docker top "$container" -eo pid,args' in runtime_test
for prefix in ("first", "second"):
container = f"${prefix}_container"
baseline = (
f'{prefix}_baseline="$(capture_process_baseline "{container}")"'
)
mapping_index = runtime_test.index(
f'assert_8080_loopback_mapping "{container}"'
)
ready_index = runtime_test.index(f'wait_ready "{container}"')
baseline_index = runtime_test.index(baseline)
close_index = runtime_test.index(
f'probe_and_close_cdp "{container}"', baseline_index
)
return_index = runtime_test.index(
f'assert_processes_return_to_baseline "{container}" "${prefix}_baseline"',
close_index,
)
assert mapping_index < ready_index < baseline_index < close_index < return_index
assert "{{.State.Pid}}" in runtime_test
assert '[ "$state" = "exited 0 0" ]' in runtime_test
def test_runtime_process_normalization_retains_every_unknown_wrapped_child(
tmp_path: Path,
) -> None:
runtime_test = ROOT / "tests/test_runtime.sh"
top_file = tmp_path / "docker-top.txt"
environment = os.environ.copy()
environment.update(
{
"DOCKER_HOST": "unix:///nonexistent-runtime-normalization.sock",
"SKIP_BUILD": "1",
}
)
service_lines = [
"101 /run/rosetta/rosetta /usr/local/bin/python3 python3 /usr/local/bin/addon-supervisor",
"102 /run/rosetta/rosetta /usr/local/bin/python3 python3 /usr/local/bin/cloakserve --headless=true --idle-timeout=30 --data-dir=/tmp/cloakserve",
"103 /usr/bin/qemu-x86_64-static /usr/local/bin/stelloauth",
]
def normalize(extra_line: str | None = None) -> list[str]:
lines = ["PID COMMAND", *service_lines]
if extra_line is not None:
lines.append(extra_line)
top_file.write_text("\n".join(lines) + "\n", encoding="utf-8")
completed = subprocess.run(
[str(runtime_test), "--normalize-processes", str(top_file)],
cwd=ROOT,
env=environment,
text=True,
capture_output=True,
check=False,
)
assert completed.returncode == 0, completed.stderr
return completed.stdout.splitlines()
baseline = ["addon-supervisor", "cloakserve", "stelloauth"]
assert normalize() == baseline
wrapped_children = [
"/run/rosetta/rosetta /opt/vendor/headless-shell --user-data-dir=/tmp/profile",
"/run/rosetta/rosetta /opt/vendor/browser --profile runtime-readiness",
"/run/rosetta/rosetta /opt/vendor/crashpad_handler --database=/tmp/profile",
"/run/rosetta/rosetta /opt/vendor/opaque-child --flag",
"/run/rosetta/rosetta /opt/vendor/opaque-child --parent=/usr/local/bin/cloakserve",
"/usr/bin/qemu-x86_64-static /opt/vendor/opaque-qemu-child",
]
for offset, child in enumerate(wrapped_children, start=104):
normalized = normalize(f"{offset} {child}")
assert normalized == sorted([*baseline, f"unexpected:{child}"])
assert normalized != baseline
+412
View File
@@ -0,0 +1,412 @@
#!/usr/bin/env bash
set -Eeuo pipefail
repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
image="homeassistant-stelloauth-addon:test"
run_id="$(date +%s)-$$"
first_container="stelloauth-runtime-${run_id}-first"
second_container="stelloauth-runtime-${run_id}-second"
tmp_dir="$(mktemp -d)"
artifacts_dir="${repo_root}/artifacts"
options_file="${tmp_dir}/options.json"
sentinels=(
"sentinel-email@example.invalid"
"SENTINEL_PASSWORD_9a34"
"SENTINEL_COOKIE_7b21"
"SENTINEL_OAUTH_CODE_5c88"
"SENTINEL_ACCESS_TOKEN_1d62"
"SENTINEL_REFRESH_TOKEN_4e73"
)
cleanup() {
set +e
for container in "$first_container" "$second_container"; do
if docker container inspect "$container" >/dev/null 2>&1; then
if [ "$(docker inspect --format '{{.State.Running}}' "$container" 2>/dev/null)" = "true" ]; then
docker stop --time 10 "$container" >/dev/null 2>&1
fi
docker rm "$container" >/dev/null 2>&1
fi
done
rm -r "$tmp_dir"
}
trap cleanup EXIT
fail() {
printf 'runtime test failed: %s\n' "$*" >&2
exit 1
}
write_options() {
python3 - "$options_file" <<'PY'
import json
import pathlib
import sys
options = {
"queue_timeout": "60s",
"rate_limit_count": 5,
"rate_limit_duration": "1h",
}
pathlib.Path(sys.argv[1]).write_text(
json.dumps(options, separators=(",", ":")) + "\n",
encoding="utf-8",
)
PY
}
start_container() {
local container="$1"
docker create --name "$container" \
--platform linux/amd64 \
--publish 127.0.0.1::8080 \
"$image" >/dev/null
docker cp "$options_file" "$container:/data/options.json"
docker start "$container" >/dev/null
}
assert_8080_loopback_mapping() {
local container="$1"
local mapping
mapping="$(docker port "$container" 8080/tcp)"
case "$mapping" in
127.0.0.1:[0-9]*) ;;
*) fail "$container 8080 mapping is '$mapping', want 127.0.0.1:<port>" ;;
esac
}
probe_root() {
local container="$1"
docker exec -i "$container" python3 - <<'PY'
import urllib.request
opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
with opener.open("http://127.0.0.1:8080/", timeout=2) as response:
if response.status != 200:
raise SystemExit(f"root status {response.status}")
response.read()
PY
}
wait_ready() {
local container="$1"
local deadline=$((SECONDS + 90))
while (( SECONDS < deadline )); do
if [ "$(docker inspect --format '{{.State.Running}}' "$container")" != "true" ]; then
docker logs "$container" >&2
fail "$container exited during readiness"
fi
if docker logs "$container" 2>&1 | grep -Fq "Stelloauth listening on 0.0.0.0:8080"; then
if probe_root "$container" >/dev/null 2>&1; then
return
fi
fi
sleep 1
done
docker logs "$container" >&2
fail "$container did not become ready within 90 seconds"
}
normalize_process_file() {
local top_file="$1"
python3 - "$top_file" <<'PY'
import pathlib
import sys
commands = []
for raw_line in pathlib.Path(sys.argv[1]).read_text(encoding="utf-8").splitlines()[1:]:
fields = raw_line.split(maxsplit=1)
command = " ".join(fields[1].split()) if len(fields) == 2 else ""
tokens = command.split()
while tokens and (
tokens[0] == "/run/rosetta/rosetta"
or pathlib.Path(tokens[0]).name.startswith("qemu-")
):
tokens = tokens[1:]
python_prefixes = (
[],
["python3"],
["/usr/local/bin/python3"],
["/usr/local/bin/python3", "python3"],
)
supervisor_commands = [
[*prefix, "/usr/local/bin/addon-supervisor"]
for prefix in python_prefixes
]
cloak_commands = [
[
*prefix,
"/usr/local/bin/cloakserve",
"--headless=true",
"--idle-timeout=30",
"--data-dir=/tmp/cloakserve",
]
for prefix in python_prefixes
]
stelloauth_commands = (
["/usr/local/bin/stelloauth"],
["/usr/local/bin/stelloauth", "/usr/local/bin/stelloauth"],
)
if tokens in supervisor_commands:
commands.append("addon-supervisor")
elif tokens in cloak_commands:
commands.append("cloakserve")
elif tokens in stelloauth_commands:
commands.append("stelloauth")
elif command:
commands.append(f"unexpected:{command}")
print("\n".join(sorted(commands)))
PY
}
normalized_process_commands() {
local container="$1"
local top_file="${tmp_dir}/${container}-processes.txt"
docker top "$container" -eo pid,args > "$top_file"
normalize_process_file "$top_file"
}
capture_process_baseline() {
local container="$1"
local expected current
local deadline=$((SECONDS + 10))
expected=$'addon-supervisor\ncloakserve\nstelloauth'
while (( SECONDS < deadline )); do
current="$(normalized_process_commands "$container")"
if [ "$current" = "$expected" ]; then
printf '%s\n' "$current"
return
fi
sleep 0.25
done
printf 'expected startup process baseline:\n%s\ncurrent process commands:\n%s\n' \
"$expected" "$current" >&2
docker top "$container" >&2
fail "$container did not reach the expected startup process baseline"
}
assert_processes_return_to_baseline() {
local container="$1"
local baseline="$2"
local current
local deadline=$((SECONDS + 10))
while (( SECONDS < deadline )); do
current="$(normalized_process_commands "$container")"
if [ "$current" = "$baseline" ]; then
return
fi
sleep 0.25
done
printf 'expected process baseline after CDP close:\n%s\ncurrent process commands:\n%s\n' \
"$baseline" "$current" >&2
docker top "$container" >&2
fail "$container retained browser or profile processes after CDP close"
}
assert_loopback_cdp_listener() {
local container="$1"
local tcp_artifact="$2"
local tcp6_artifact="$3"
docker exec "$container" cat /proc/net/tcp > "$tcp_artifact"
docker exec "$container" cat /proc/net/tcp6 > "$tcp6_artifact"
docker exec -i "$container" python3 - <<'PY'
expected = f"0100007F:{9222:04X}"
if expected != "0100007F:2406":
raise SystemExit(f"unexpected 9222 hexadecimal encoding: {expected}")
listeners = {"/proc/net/tcp": set(), "/proc/net/tcp6": set()}
for table in ("/proc/net/tcp", "/proc/net/tcp6"):
with open(table, encoding="ascii") as handle:
next(handle)
for line in handle:
fields = line.split()
if len(fields) < 4 or fields[3] != "0A":
continue
local_address = fields[1].upper()
if local_address.rsplit(":", 1)[-1] != "2406":
continue
listeners[table].add(local_address)
if table == "/proc/net/tcp6":
raise SystemExit(f"IPv6 CDP listener present: {local_address}")
if listeners["/proc/net/tcp"] != {expected}:
raise SystemExit(
f"IPv4 CDP listeners = {sorted(listeners['/proc/net/tcp'])}, want [{expected}]"
)
PY
}
probe_and_close_cdp() {
local container="$1"
docker exec -i "$container" python3 - <<'PY'
import json
import urllib.request
opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
version_url = "http://127.0.0.1:9222/json/version?fingerprint=runtime-readiness"
close_url = "http://127.0.0.1:9222/fingerprint/runtime-readiness/close"
with opener.open(version_url, timeout=10) as response:
if response.status != 200:
raise SystemExit(f"CDP version status {response.status}")
document = json.load(response)
websocket_url = document.get("webSocketDebuggerUrl")
if not isinstance(websocket_url, str) or not websocket_url:
raise SystemExit("CDP response lacks webSocketDebuggerUrl")
request = urllib.request.Request(close_url, data=b"", method="POST")
with opener.open(request, timeout=10) as response:
if response.status != 200:
raise SystemExit(f"CDP close status {response.status}")
response.read()
PY
}
post_invalid_worker() {
local container="$1"
local response_artifact="$2"
docker exec -i "$container" python3 - > "$response_artifact" <<'PY'
import json
import sys
import urllib.error
import urllib.request
body = {
"url": (
"https://example.invalid/am/oauth2/authorize"
"?redirect_uri=sentinel%3A%2F%2Fcallback"
"&code=SENTINEL_OAUTH_CODE_5c88"
"&access_token=SENTINEL_ACCESS_TOKEN_1d62"
"&refresh_token=SENTINEL_REFRESH_TOKEN_4e73"
"&cookie=SENTINEL_COOKIE_7b21"
),
"email": "sentinel-email@example.invalid",
"password": "SENTINEL_PASSWORD_9a34",
"cookie": "SENTINEL_COOKIE_7b21",
"oauth_code": "SENTINEL_OAUTH_CODE_5c88",
"access_token": "SENTINEL_ACCESS_TOKEN_1d62",
"refresh_token": "SENTINEL_REFRESH_TOKEN_4e73",
}
request = urllib.request.Request(
"http://127.0.0.1:8080/worker",
data=json.dumps(body, separators=(",", ":")).encode(),
headers={"Content-Type": "application/json"},
method="POST",
)
opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
try:
with opener.open(request, timeout=10) as response:
status = response.status
response_body = response.read()
except urllib.error.HTTPError as error:
status = error.code
response_body = error.read()
if status != 400:
raise SystemExit(f"invalid worker status {status}, want 400")
sys.stdout.buffer.write(response_body)
PY
}
assert_no_9222_mapping() {
local container="$1"
local mapping
mapping="$(docker port "$container" 9222/tcp 2>/dev/null || true)"
[ -z "$mapping" ] || fail "container 9222 is mapped: $mapping"
}
scan_logs() {
local container="$1"
local log_file="${tmp_dir}/${container}.log"
docker logs "$container" > "$log_file" 2>&1
for sentinel in "${sentinels[@]}"; do
if grep -Fq "$sentinel" "$log_file"; then
fail "$container logs contain sentinel $sentinel"
fi
done
if grep -Fq "worker OAuth request" "$log_file"; then
fail "$container began an OAuth flow for the rejected worker body"
fi
}
stop_and_assert() {
local container="$1"
local timing_artifact="$2"
local started_ns ended_ns elapsed state
started_ns="$(python3 -c 'import time; print(time.monotonic_ns())')"
docker stop --time 10 "$container" >/dev/null
ended_ns="$(python3 -c 'import time; print(time.monotonic_ns())')"
elapsed="$(python3 - "$started_ns" "$ended_ns" <<'PY'
import sys
print((int(sys.argv[2]) - int(sys.argv[1])) / 1_000_000_000)
PY
)"
printf 'seconds=%s\n' "$elapsed" > "$timing_artifact"
python3 - "$elapsed" <<'PY'
import sys
if float(sys.argv[1]) > 10.0:
raise SystemExit(f"container stop exceeded 10 seconds: {sys.argv[1]}")
PY
state="$(docker inspect --format '{{.State.Status}} {{.State.ExitCode}} {{.State.Pid}}' "$container")"
[ "$state" = "exited 0 0" ] || fail "$container state is $state, want exited 0 with PID 0"
}
if [ "${1:-}" = "--normalize-processes" ]; then
[ "$#" -eq 2 ] || fail "--normalize-processes requires one docker top file"
normalize_process_file "$2"
exit
fi
[ "$#" -eq 0 ] || fail "unexpected runtime test arguments"
mkdir -p "$artifacts_dir"
write_options
if [ "${SKIP_BUILD:-0}" != "1" ]; then
docker buildx build \
--platform linux/amd64 \
--build-arg BUILD_ARCH=amd64 \
--load \
--tag "$image" \
"$repo_root/stelloauth"
fi
start_container "$first_container"
assert_8080_loopback_mapping "$first_container"
wait_ready "$first_container"
first_baseline="$(capture_process_baseline "$first_container")"
assert_no_9222_mapping "$first_container"
assert_loopback_cdp_listener \
"$first_container" \
"${artifacts_dir}/runtime-proc-net-tcp.txt" \
"${artifacts_dir}/runtime-proc-net-tcp6.txt"
probe_and_close_cdp "$first_container"
assert_processes_return_to_baseline "$first_container" "$first_baseline"
post_invalid_worker "$first_container" "${artifacts_dir}/runtime-invalid-worker-response.json"
scan_logs "$first_container"
stop_and_assert "$first_container" "${artifacts_dir}/runtime-first-stop.txt"
scan_logs "$first_container"
start_container "$second_container"
assert_8080_loopback_mapping "$second_container"
wait_ready "$second_container"
second_baseline="$(capture_process_baseline "$second_container")"
assert_no_9222_mapping "$second_container"
assert_loopback_cdp_listener \
"$second_container" \
"${artifacts_dir}/runtime-restart-proc-net-tcp.txt" \
"${artifacts_dir}/runtime-restart-proc-net-tcp6.txt"
probe_and_close_cdp "$second_container"
assert_processes_return_to_baseline "$second_container" "$second_baseline"
sleep 5
docker stats --no-stream "$second_container" > "${artifacts_dir}/runtime-docker-stats.txt"
docker top "$second_container" > "${artifacts_dir}/runtime-docker-top.txt"
docker image inspect "$image" --format '{{.Size}}' > "${artifacts_dir}/runtime-image-size-bytes.txt"
docker image inspect "$image" --format '{{json .Config.ExposedPorts}}' > "${artifacts_dir}/runtime-image-exposed-ports.json"
docker inspect "$second_container" --format '{{json .HostConfig.PortBindings}}' > "${artifacts_dir}/runtime-host-port-bindings.json"
stop_and_assert "$second_container" "${artifacts_dir}/runtime-second-stop.txt"
scan_logs "$second_container"
printf 'runtime acceptance PASS: root, CDP, loopback bind, invalid worker, redaction, stop, restart\n'
+56
View File
@@ -0,0 +1,56 @@
from __future__ import annotations
import subprocess
from pathlib import Path
import pytest
ROOT = Path(__file__).parents[1]
STELLOAUTH_COMMIT = "367d4f8c02a3b072c59142c49dffc129edc8548b"
CLOAK_COMMIT = "f04c23da285b3b3d3cf10c8f9d282e7adc1d52ce"
def run(*args: str, cwd: Path):
return subprocess.run(args, cwd=cwd, text=True, capture_output=True, check=True)
def fetch_exact(tmp_path: Path, name: str, url: str, commit: str) -> Path:
target = tmp_path / name
run("git", "init", str(target), cwd=tmp_path)
run("git", "remote", "add", "origin", url, cwd=target)
run("git", "fetch", "--depth=1", "origin", commit, cwd=target)
run("git", "checkout", "--detach", "FETCH_HEAD", cwd=target)
assert run("git", "rev-parse", "HEAD", cwd=target).stdout.strip() == commit
return target
@pytest.mark.upstream
def test_stelloauth_patch_applies_and_tests_pass(tmp_path: Path) -> None:
source = fetch_exact(
tmp_path,
"stelloauth",
"https://github.com/tamcore/stelloauth.git",
STELLOAUTH_COMMIT,
)
patch = ROOT / "stelloauth/patches/stelloauth-security.patch"
run("git", "apply", "--check", str(patch), cwd=source)
run("git", "apply", str(patch), cwd=source)
run("go", "test", "./...", cwd=source)
@pytest.mark.upstream
def test_cloak_patch_binds_only_loopback(tmp_path: Path) -> None:
source = fetch_exact(
tmp_path,
"cloakbrowser",
"https://github.com/CloakHQ/CloakBrowser.git",
CLOAK_COMMIT,
)
patch = ROOT / "stelloauth/patches/cloakserve-loopback.patch"
run("git", "apply", "--check", str(patch), cwd=source)
run("git", "apply", str(patch), cwd=source)
wrapper = (source / "bin/cloakserve").read_text(encoding="utf-8")
assert 'host = "127.0.0.1"' in wrapper
assert 'host = "0.0.0.0" if in_container' not in wrapper
run("python3", "-m", "py_compile", "bin/cloakserve", cwd=source)
+907
View File
@@ -0,0 +1,907 @@
from __future__ import annotations
import importlib.machinery
import importlib.util
import json
import logging
import signal
import sys
from pathlib import Path
from types import SimpleNamespace
import pytest
ROOT = Path(__file__).parents[1]
SUPERVISOR_PATH = ROOT / "stelloauth/rootfs/usr/local/bin/addon-supervisor"
def load_supervisor():
loader = importlib.machinery.SourceFileLoader("addon_supervisor", str(SUPERVISOR_PATH))
spec = importlib.util.spec_from_loader(loader.name, loader)
assert spec is not None
module = importlib.util.module_from_spec(spec)
sys.modules[loader.name] = module
loader.exec_module(module)
return module
@pytest.fixture
def supervisor():
return load_supervisor()
def test_options_load_valid_defaults(supervisor, tmp_path: Path) -> None:
path = tmp_path / "options.json"
path.write_text(
json.dumps(
{
"queue_timeout": "60s",
"rate_limit_count": 5,
"rate_limit_duration": "1h",
}
),
encoding="utf-8",
)
assert supervisor.load_options(path) == supervisor.Options("60s", 5, "1h")
def test_environment_maps_options_and_preserves_parent(supervisor, monkeypatch) -> None:
monkeypatch.setenv("PARENT_SENTINEL", "preserved")
environment = supervisor.build_environment(supervisor.Options("60s", 5, "1h"))
assert environment["PARENT_SENTINEL"] == "preserved"
assert {key: environment[key] for key in (
"CLOAK_CDP_URL",
"CLOAK_MAX_SESSIONS",
"CLOAK_QUEUE_TIMEOUT",
"RATE_LIMIT_COUNT",
"RATE_LIMIT_DURATION",
"HTTP_ADDRESS",
"PORT",
"METRICS_ADDRESS",
"METRICS_PORT",
)} == {
"CLOAK_CDP_URL": "http://127.0.0.1:9222",
"CLOAK_MAX_SESSIONS": "1",
"CLOAK_QUEUE_TIMEOUT": "60s",
"RATE_LIMIT_COUNT": "5",
"RATE_LIMIT_DURATION": "1h",
"HTTP_ADDRESS": "0.0.0.0",
"PORT": "8080",
"METRICS_ADDRESS": "127.0.0.1",
"METRICS_PORT": "9090",
}
@pytest.mark.parametrize(
"content",
[
"{}",
json.dumps(
{
"queue_timeout": "60s",
"rate_limit_count": 5,
"rate_limit_duration": "1h",
"unknown-SENTINEL": "secret-SENTINEL",
}
),
json.dumps({"queue_timeout": "60s", "rate_limit_count": 5}),
json.dumps(
{
"queue_timeout": "0s-SENTINEL",
"rate_limit_count": 5,
"rate_limit_duration": "1h",
}
),
json.dumps(
{
"queue_timeout": "60-SENTINEL",
"rate_limit_count": 5,
"rate_limit_duration": "1h",
}
),
json.dumps(
{
"queue_timeout": "60s",
"rate_limit_count": True,
"rate_limit_duration": "1h",
}
),
json.dumps(
{
"queue_timeout": "60s",
"rate_limit_count": 0,
"rate_limit_duration": "1h",
}
),
json.dumps(
{
"queue_timeout": "60s",
"rate_limit_count": 21,
"rate_limit_duration": "1h",
}
),
json.dumps(
{
"queue_timeout": "60s",
"rate_limit_count": 5,
"rate_limit_duration": "1-SENTINEL",
}
),
'{"queue_timeout":"malformed-SENTINEL"',
],
)
def test_invalid_options_raise_fixed_non_secret_error(
supervisor, tmp_path: Path, caplog, content: str
) -> None:
path = tmp_path / "options.json"
path.write_text(content, encoding="utf-8")
with caplog.at_level(logging.INFO), pytest.raises(
supervisor.ConfigError, match="^Invalid add-on configuration$"
):
supervisor.load_options(path)
assert "SENTINEL" not in caplog.text
assert "SENTINEL" not in str(sys.exc_info())
def test_unreadable_options_raise_fixed_non_secret_error(
supervisor, tmp_path: Path, caplog
) -> None:
missing = tmp_path / "missing-SENTINEL.json"
with caplog.at_level(logging.INFO), pytest.raises(
supervisor.ConfigError, match="^Invalid add-on configuration$"
):
supervisor.load_options(missing)
assert "SENTINEL" not in caplog.text
def test_invalid_options_main_logs_only_fixed_error(
supervisor, tmp_path: Path, caplog, monkeypatch
) -> None:
path = tmp_path / "options.json"
path.write_text('{"credential":"secret-SENTINEL"}', encoding="utf-8")
monkeypatch.setattr(supervisor, "OPTIONS_PATH", path)
with caplog.at_level(logging.INFO):
assert supervisor.main() == 2
assert caplog.messages == ["Invalid add-on configuration"]
assert "SENTINEL" not in caplog.text
def test_probe_cloak_uses_real_cdp_websocket_and_closes_profile(supervisor) -> None:
calls: list[tuple[str, str, float]] = []
def request(method: str, url: str, timeout: float):
calls.append((method, url, timeout))
if url == supervisor.CLOAK_VERSION:
return supervisor.HttpResponse(
200,
json.dumps(
{
"webSocketDebuggerUrl": (
"ws://127.0.0.1:9222/devtools/browser/readiness"
)
}
).encode(),
)
return supervisor.HttpResponse(200, b"ok")
supervisor.probe_cloak(request)
assert [(method, url) for method, url, _ in calls] == [
("GET", supervisor.CLOAK_ROOT),
("GET", supervisor.CLOAK_VERSION),
("POST", supervisor.CLOAK_CLOSE),
]
assert all(timeout == 2.0 for _, _, timeout in calls)
@pytest.mark.parametrize(
"root_status,version_status,version_body,close_status",
[
(200, 200, b"{}", 200),
(200, 200, b'{"webSocketDebuggerUrl":""}', 200),
(200, 200, b"not-json-SENTINEL", 200),
(
200,
200,
b'{"webSocketDebuggerUrl":"ws://192.0.2.1:9222/devtools/browser/x"}',
200,
),
(
200,
200,
b'{"webSocketDebuggerUrl":"WS://127.0.0.1:9222/devtools/browser/x"}',
200,
),
(
503,
200,
b'{"webSocketDebuggerUrl":"ws://127.0.0.1:9222/devtools/browser/x"}',
200,
),
(
200,
503,
b'{"webSocketDebuggerUrl":"ws://127.0.0.1:9222/devtools/browser/x"}',
200,
),
(
200,
200,
b'{"webSocketDebuggerUrl":"ws://127.0.0.1:9222/devtools/browser/x"}',
503,
),
],
)
def test_probe_cloak_rejects_invalid_readiness(
supervisor, root_status, version_status, version_body, close_status
) -> None:
responses = {
supervisor.CLOAK_ROOT: supervisor.HttpResponse(root_status, b"root"),
supervisor.CLOAK_VERSION: supervisor.HttpResponse(version_status, version_body),
supervisor.CLOAK_CLOSE: supervisor.HttpResponse(close_status, b"close"),
}
with pytest.raises(supervisor.ReadinessError):
supervisor.probe_cloak(lambda _method, url, _timeout: responses[url])
def test_probe_stelloauth_requires_http_200(supervisor) -> None:
calls = []
def request(method: str, url: str, timeout: float):
calls.append((method, url, timeout))
return supervisor.HttpResponse(200, b"SENTINEL-body-is-ignored")
supervisor.probe_stelloauth(request)
assert calls == [("GET", supervisor.STELLOAUTH_ROOT, 2.0)]
with pytest.raises(supervisor.ReadinessError):
supervisor.probe_stelloauth(
lambda _method, _url, _timeout: supervisor.HttpResponse(503, b"")
)
class FakeClock:
def __init__(self) -> None:
self.now = 0.0
self.sleeps: list[float] = []
def monotonic(self) -> float:
return self.now
def sleep(self, delay: float) -> None:
self.sleeps.append(delay)
self.now += delay
def test_readiness_backoff_starts_at_quarter_second_and_caps_at_two(supervisor) -> None:
clock = FakeClock()
attempts = 0
def probe() -> None:
nonlocal attempts
attempts += 1
if attempts <= 5:
raise OSError("transient-SENTINEL")
supervisor.wait_until_ready(
"Service", probe, 60.0, lambda: False, clock.monotonic, clock.sleep
)
assert clock.sleeps == [0.25, 0.5, 1.0, 2.0, 2.0]
@pytest.mark.parametrize("name,timeout", [("CloakBrowser", 60.0), ("Stelloauth", 30.0)])
def test_readiness_expires_at_service_timeout(supervisor, name: str, timeout: float) -> None:
clock = FakeClock()
with pytest.raises(
supervisor.ReadinessError, match=f"^{name} did not become ready$"
):
supervisor.wait_until_ready(
name,
lambda: (_ for _ in ()).throw(ValueError("transient-SENTINEL")),
timeout,
lambda: False,
clock.monotonic,
clock.sleep,
)
assert timeout <= clock.now <= timeout + 2.0
assert max(clock.sleeps) == 2.0
def test_readiness_stop_flag_aborts_immediately(supervisor) -> None:
clock = FakeClock()
called = False
def probe() -> None:
nonlocal called
called = True
with pytest.raises(
supervisor.ReadinessError, match="^CloakBrowser did not become ready$"
):
supervisor.wait_until_ready(
"CloakBrowser", probe, 60.0, lambda: True, clock.monotonic, clock.sleep
)
assert called is False
assert clock.sleeps == []
def test_probe_http_request_disables_proxies(supervisor, monkeypatch) -> None:
observed = {}
class Response:
status = 200
def read(self) -> bytes:
return b"response"
def __enter__(self):
return self
def __exit__(self, *_args):
return None
class Opener:
def open(self, request, timeout):
observed["request"] = request
observed["timeout"] = timeout
return Response()
def build_opener(handler):
observed["handler"] = handler
return Opener()
monkeypatch.setattr(supervisor.urllib.request, "build_opener", build_opener)
assert supervisor.http_request("POST", "http://127.0.0.1/", 3.0) == (
200,
b"response",
)
assert observed["handler"].proxies == {}
assert observed["request"].get_method() == "POST"
assert observed["timeout"] == 3.0
class FakeProcess:
def __init__(self, pid: int, *, ignores_term: bool = False) -> None:
self.pid = pid
self.returncode: int | None = None
self.ignores_term = ignores_term
self.wait_calls: list[float | None] = []
self.on_wait = None
def poll(self) -> int | None:
return self.returncode
def wait(self, timeout: float | None) -> int:
self.wait_calls.append(timeout)
if self.returncode is None:
self.returncode = -9
if self.on_wait is not None:
self.on_wait()
return self.returncode
def manager_harness(
supervisor,
tmp_path: Path,
*,
cloak_probe=None,
stelloauth_probe=None,
ignores_term: tuple[bool, bool] = (False, False),
live_descendants: tuple[bool, bool] = (False, False),
):
clock = FakeClock()
events: list[object] = []
processes = [
FakeProcess(1001, ignores_term=ignores_term[0]),
FakeProcess(1002, ignores_term=ignores_term[1]),
]
groups = {process.pid: True for process in processes}
leaders_reaped = {process.pid: False for process in processes}
term_delivered = {process.pid: False for process in processes}
for index, process in enumerate(processes):
def on_wait(index=index, process=process) -> None:
leaders_reaped[process.pid] = True
if not live_descendants[index] or (
term_delivered[process.pid] and not process.ignores_term
):
groups[process.pid] = False
process.on_wait = on_wait
group_checks: list[int] = []
spawned: list[FakeProcess] = []
def popen(command, *, env, start_new_session):
process = processes[len(spawned)]
spawned.append(process)
events.append(("start", list(command), env, start_new_session))
return process
def killpg(pid: int, sent_signal: int) -> None:
events.append(("signal", pid, sent_signal))
process = next(item for item in processes if item.pid == pid)
if sent_signal == signal.SIGKILL:
groups[pid] = False
if process.returncode is None:
process.returncode = -sent_signal
else:
term_delivered[pid] = True
if not process.ignores_term:
if process.returncode is None:
process.returncode = -sent_signal
if leaders_reaped[pid]:
groups[pid] = False
def group_alive(pid: int) -> bool:
group_checks.append(pid)
return groups[pid]
def default_cloak_probe() -> None:
events.append("probe cloak")
def default_stelloauth_probe() -> None:
events.append("probe stelloauth")
profile_path = tmp_path / "cloakserve"
environment = {"SENSITIVE_SENTINEL": "credential-code-cookie-token-SENTINEL"}
manager = supervisor.ProcessManager(
environment,
popen=popen,
killpg=killpg,
group_alive=group_alive,
cloak_probe=cloak_probe or default_cloak_probe,
stelloauth_probe=stelloauth_probe or default_stelloauth_probe,
monotonic=clock.monotonic,
sleep=clock.sleep,
profile_path=profile_path,
)
return SimpleNamespace(
manager=manager,
clock=clock,
events=events,
processes=processes,
groups=groups,
group_checks=group_checks,
spawned=spawned,
profile_path=profile_path,
environment=environment,
)
@pytest.mark.parametrize(
"exception,expected",
[(None, True), (ProcessLookupError(), False), (PermissionError(), True)],
)
def test_process_group_liveness_uses_signal_zero(
supervisor, exception: OSError | None, expected: bool
) -> None:
calls = []
def killpg(pgid: int, sent_signal: int) -> None:
calls.append((pgid, sent_signal))
if exception is not None:
raise exception
assert supervisor.process_group_alive(4242, killpg) is expected
assert calls == [(4242, 0)]
def test_lifecycle_startup_order_and_profiles(supervisor, tmp_path: Path, monkeypatch) -> None:
harness = manager_harness(supervisor, tmp_path)
harness.profile_path.mkdir()
stale = harness.profile_path / "stale-profile"
stale.write_text("stale", encoding="utf-8")
original_cloak = harness.manager._cloak_probe
original_stelloauth = harness.manager._stelloauth_probe
def cloak_probe() -> None:
assert not stale.exists()
assert harness.profile_path.stat().st_mode & 0o777 == 0o700
original_cloak()
def stelloauth_probe() -> None:
original_stelloauth()
harness.manager._handle_signal(signal.SIGTERM, None)
harness.manager._cloak_probe = cloak_probe
harness.manager._stelloauth_probe = stelloauth_probe
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 0
starts_and_probes = [event for event in harness.events if event == "probe cloak" or event == "probe stelloauth" or (isinstance(event, tuple) and event[0] == "start")]
assert [(event[0], event[1]) if isinstance(event, tuple) else event for event in starts_and_probes] == [
("start", supervisor.CLOAK_COMMAND),
"probe cloak",
("start", supervisor.STELLOAUTH_COMMAND),
"probe stelloauth",
]
for event in starts_and_probes:
if isinstance(event, tuple):
assert event[2] is harness.environment
assert event[3] is True
def test_cloak_readiness_failure_never_starts_stelloauth(
supervisor, tmp_path: Path, monkeypatch
) -> None:
def failing_probe() -> None:
raise supervisor.ReadinessError("secret-SENTINEL")
harness = manager_harness(supervisor, tmp_path, cloak_probe=failing_probe)
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 1
assert len(harness.spawned) == 1
assert harness.clock.now >= 60.0
assert harness.processes[0].wait_calls == [None]
def test_stelloauth_readiness_failure_stops_both_children(
supervisor, tmp_path: Path, monkeypatch
) -> None:
def failing_probe() -> None:
raise OSError("credential-SENTINEL")
harness = manager_harness(supervisor, tmp_path, stelloauth_probe=failing_probe)
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 1
assert len(harness.spawned) == 2
assert {(event[1], event[2]) for event in harness.events if event[0] == "signal"} == {
(1001, signal.SIGTERM),
(1002, signal.SIGTERM),
}
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
@pytest.mark.parametrize(
"child_index,child_status,expected_status",
[(0, 0, 1), (0, 7, 7), (1, 0, 1), (1, 9, 9)],
)
def test_child_exit_stops_sibling_and_returns_failure(
supervisor,
tmp_path: Path,
monkeypatch,
child_index: int,
child_status: int,
expected_status: int,
) -> None:
harness = manager_harness(supervisor, tmp_path)
slept = False
def sleep(delay: float) -> None:
nonlocal slept
harness.clock.sleep(delay)
if not slept:
slept = True
harness.processes[child_index].returncode = child_status
harness.manager._sleep = sleep
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == expected_status
sibling = harness.processes[1 - child_index]
assert ("signal", sibling.pid, signal.SIGTERM) in harness.events
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_exited_leader_with_live_descendants_still_gets_group_sigterm(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(
supervisor, tmp_path, live_descendants=(True, False)
)
triggered = False
def sleep(delay: float) -> None:
nonlocal triggered
harness.clock.sleep(delay)
if not triggered:
triggered = True
harness.processes[0].returncode = 7
assert harness.groups[1001] is True
harness.manager._sleep = sleep
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 7
assert ("signal", 1001, signal.SIGTERM) in harness.events
assert ("signal", 1002, signal.SIGTERM) in harness.events
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_leader_exit_on_sigterm_without_descendants_finishes_without_waiting(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(supervisor, tmp_path)
def stelloauth_probe() -> None:
harness.manager._handle_signal(signal.SIGTERM, None)
harness.manager._stelloauth_probe = stelloauth_probe
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 0
assert harness.clock.now == 0.0
assert not any(
event[0] == "signal" and event[2] == signal.SIGKILL
for event in harness.events
)
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_exited_leader_with_live_descendants_still_uses_deadline_and_sigkill(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(
supervisor,
tmp_path,
ignores_term=(True, False),
live_descendants=(True, False),
)
triggered = False
def sleep(delay: float) -> None:
nonlocal triggered
harness.clock.sleep(delay)
if not triggered:
triggered = True
harness.processes[0].returncode = 7
harness.manager._sleep = sleep
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 7
assert harness.clock.now == pytest.approx(9.25)
assert ("signal", 1001, signal.SIGTERM) in harness.events
assert ("signal", 1001, signal.SIGKILL) in harness.events
assert ("signal", 1002, signal.SIGTERM) in harness.events
assert ("signal", 1002, signal.SIGKILL) not in harness.events
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_group_disappearance_ends_shared_wait_without_sigkill(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(supervisor, tmp_path, ignores_term=(True, True))
triggered = False
def sleep(delay: float) -> None:
nonlocal triggered
if not triggered:
triggered = True
harness.manager._handle_signal(signal.SIGTERM, None)
harness.clock.sleep(delay)
if harness.clock.now >= 0.5:
harness.groups[1001] = False
harness.groups[1002] = False
harness.manager._sleep = sleep
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 0
assert harness.clock.now == pytest.approx(0.5)
assert all(event[2] != signal.SIGKILL for event in harness.events if event[0] == "signal")
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_dead_group_is_not_rechecked_or_signalled_after_possible_pgid_reuse(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(supervisor, tmp_path)
calls = []
def group_alive(pgid: int) -> bool:
calls.append(pgid)
if pgid == 1001:
return len([item for item in calls if item == 1001]) > 1
return harness.groups[pgid]
harness.manager._group_alive = group_alive
triggered = False
def sleep(delay: float) -> None:
nonlocal triggered
harness.clock.sleep(delay)
if not triggered:
triggered = True
harness.processes[0].returncode = 7
harness.manager._sleep = sleep
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 7
assert calls.count(1001) == 1
assert not any(
event[0] == "signal" and event[1] == 1001 for event in harness.events
)
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_group_disappearing_at_deadline_is_rechecked_before_sigkill(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(supervisor, tmp_path, ignores_term=(True, True))
deadline_checks = {1001: 0, 1002: 0}
def group_alive(pgid: int) -> bool:
if harness.clock.now < supervisor.SHUTDOWN_GRACE_SECONDS:
return True
deadline_checks[pgid] += 1
return deadline_checks[pgid] == 1
harness.manager._group_alive = group_alive
triggered = False
def sleep(delay: float) -> None:
nonlocal triggered
if not triggered:
triggered = True
harness.manager._handle_signal(signal.SIGTERM, None)
harness.clock.sleep(delay)
harness.manager._sleep = sleep
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 0
assert deadline_checks == {1001: 2, 1002: 2}
assert not any(
event[0] == "signal" and event[2] == signal.SIGKILL
for event in harness.events
)
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
@pytest.mark.parametrize("incoming_signal", [signal.SIGTERM, signal.SIGINT])
def test_signal_shutdown_forwards_sigterm_and_returns_zero(
supervisor, tmp_path: Path, monkeypatch, incoming_signal: int
) -> None:
harness = manager_harness(supervisor, tmp_path)
triggered = False
def sleep(delay: float) -> None:
nonlocal triggered
if not triggered:
triggered = True
harness.manager._handle_signal(incoming_signal, None)
harness.clock.sleep(delay)
harness.manager._sleep = sleep
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 0
assert {(event[1], event[2]) for event in harness.events if event[0] == "signal"} == {
(1001, signal.SIGTERM),
(1002, signal.SIGTERM),
}
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_signal_during_cloak_readiness_never_starts_stelloauth(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(supervisor, tmp_path)
def cloak_probe() -> None:
harness.manager._handle_signal(signal.SIGTERM, None)
harness.manager._cloak_probe = cloak_probe
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 0
assert len(harness.spawned) == 1
assert ("signal", 1001, signal.SIGTERM) in harness.events
assert harness.processes[0].wait_calls == [None]
def test_child_exit_during_cloak_readiness_never_starts_stelloauth(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(supervisor, tmp_path)
def cloak_probe() -> None:
harness.processes[0].returncode = 7
harness.manager._cloak_probe = cloak_probe
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 7
assert len(harness.spawned) == 1
assert harness.processes[0].wait_calls == [None]
def test_signal_while_starting_child_still_terminates_new_process_group(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(supervisor, tmp_path)
original_popen = harness.manager._popen
starts = 0
def popen(command, *, env, start_new_session):
nonlocal starts
starts += 1
if starts == 2:
harness.manager._handle_signal(signal.SIGTERM, None)
return original_popen(command, env=env, start_new_session=start_new_session)
harness.manager._popen = popen
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 0
assert ("signal", 1002, signal.SIGTERM) in harness.events
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_shutdown_reserves_time_before_outer_ten_second_stop_deadline(
supervisor, tmp_path: Path, monkeypatch
) -> None:
harness = manager_harness(supervisor, tmp_path, ignores_term=(True, True))
triggered = False
def sleep(delay: float) -> None:
nonlocal triggered
if not triggered:
triggered = True
harness.manager._handle_signal(signal.SIGTERM, None)
harness.clock.sleep(delay)
harness.manager._sleep = sleep
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
assert harness.manager.run() == 0
assert harness.clock.now < 10.0
assert harness.clock.now == pytest.approx(supervisor.SHUTDOWN_GRACE_SECONDS)
assert [
(event[1], event[2]) for event in harness.events if event[0] == "signal"
] == [
(1001, signal.SIGTERM),
(1002, signal.SIGTERM),
(1001, signal.SIGKILL),
(1002, signal.SIGKILL),
]
assert [process.wait_calls for process in harness.processes] == [[None], [None]]
def test_lifecycle_logs_are_fixed_and_contain_no_sensitive_values(
supervisor, tmp_path: Path, monkeypatch, caplog
) -> None:
harness = manager_harness(supervisor, tmp_path)
def stelloauth_probe() -> None:
harness.manager._handle_signal(signal.SIGTERM, None)
harness.manager._stelloauth_probe = stelloauth_probe
monkeypatch.setattr(supervisor.signal, "signal", lambda *_args: None)
with caplog.at_level(logging.INFO):
assert harness.manager.run() == 0
assert caplog.messages == [
"Cleaning CloakBrowser profiles",
"Starting CloakBrowser",
"CloakBrowser ready",
"Starting Stelloauth",
"Shutdown requested",
"Stopping child processes",
"Child processes stopped",
]
lowered = caplog.text.lower()
for sensitive in ("sentinel", "credential", "code", "cookie", "token"):
assert sensitive not in lowered