]> git.ekhem.eu.org Git - guix.git/commitdiff
[proxy] Stream responses and bound waits with an idle timeout. main
authorJakub Czajka <jakub@ekhem.eu.org>
Thu, 3 Sep 2026 22:54:31 +0000 (22:54 +0000)
committerJakub Czajka <jakub@ekhem.eu.org>
Thu, 3 Sep 2026 22:54:31 +0000 (22:54 +0000)
Rewrite the proxy to relay upstream responses incrementally:
each chunked SSE frame reaches the client as it arrives instead
of being buffered, and every network wait — connect, response
headers, and gaps mid-stream — is bounded by an idle timeout
(CLAUDE_PROXY_TIMEOUT, default 60 s).  A mid-stream expiry emits
a typed SSE error and aborts the relay.

DeepSeek reasoning models default to thinking on, which stalls
non-streaming requests, so such requests now get thinking:disabled
injected unless they set a thinking field themselves; streaming
and explicit-thinking requests pass through untouched.  This
replaces the old safety-classifier sniffing.

Log one line per event with a per-request id, so a single request
can be traced with grep 'req=<id>'.  Handle portless https
upstreams (scheme-default port and Host header) and give the
Shepherd service a SIGTERM stop with a paced respawn that
survives the port overlap during guix home reconfigure.

Rework the tests around a deterministic fake upstream that can
stream with gaps and stall: assert incremental relaying,
thinking injection, 504 on stalled headers, the mid-stream typed
error, log traceability, and error passthrough.  Add a TLS twin
pair — an HTTPS fake on 127.0.0.1:443 behind a self-signed test
certificate and a proxy on portless https://127.0.0.1 — to guard
scheme-default-port and TLS record handling.  vps-base.scm now
builds its test services from claude-proxy-test-os-services
instead of duplicating the wrapper.

Co-Authored-By: Claude <noreply@anthropic.com>
conf/common/claude-proxy.scm
tests/claude-proxy-test-cert.pem [new file with mode: 0644]
tests/claude-proxy-test-key.pem [new file with mode: 0644]
tests/claude-proxy.scm
tests/vps-base.scm

index f2346db0d52a594033cafdb390e23eeb6f1208fa..e9ad008c39418011f530e413980cba622a5bb86c 100644 (file)
@@ -1,9 +1,14 @@
 ;; Copyright (c) 2026 Jakub Czajka <jakub@ekhem.eu.org>
 ;; License: GPL-3.0 or later.
 ;;
-;; claude-proxy.scm — Minimal Guile HTTP proxy for Claude Code.
-;; Detects safety-classifier requests, injects thinking:disabled.
-;; Provides both the proxy program-file and a home Shepherd service.
+;; claude-proxy.scm — Guile HTTP proxy for Claude Code.
+;; Forwards Anthropic-format requests to DeepSeek: injects
+;; thinking:disabled on non-streaming requests that do not set a
+;; thinking field, streams the upstream response back incrementally,
+;; and bounds every network wait with an idle timeout.  Logs one
+;; line per event with a per-request id so a single request can be
+;; traced with grep 'req=<id>'.  Provides both the proxy
+;; program-file and a home Shepherd service.
 
 (define-module (conf common claude-proxy)
   #:use-module (gnu packages guile)
   #:export (claude-proxy-script claude-proxy-daemon-service))
 
 (define claude-proxy-code
+  ;; The whole script is one top-level (begin ...), so use-modules
+  ;; here resolves its modules at expansion time.  The launching
+  ;; wrapper therefore sets GUILE_LOAD_PATH and GUILE_LOAD_COMPILED_PATH
+  ;; to the store directories of the off-path (gnutls) and (json)
+  ;; modules before exec'ing guile.
   #~(begin
-      (set! %load-path
-            (cons* #$(file-append guile-gnutls "/share/guile/site/3.0")
-                   #$(file-append guile-json-4 "/share/guile/site/3.0")
-                   %load-path))
-      (set! %load-compiled-path
-            (cons* #$(file-append guile-gnutls "/lib/guile/3.0/site-ccache")
-                   #$(file-append guile-json-4 "/lib/guile/3.0/site-ccache")
-                   %load-compiled-path))
-      (use-modules (web server)
+      (use-modules (gnutls)
                    (web client)
                    (web request)
                    (web response)
                    (web uri)
+                   (web http)
                    (json)
                    (rnrs bytevectors)
-                   (srfi srfi-1)
+                   (rnrs io ports)
+                   (ice-9 binary-ports)
+                   (ice-9 threads)
                    (ice-9 match)
-                   (ice-9 receive))
+                   (ice-9 receive)
+                   (ice-9 rdelim)
+                   (srfi srfi-1)
+                   (srfi srfi-13))
 
+      ;; Make the GnuTLS bindings visible to (web client) so its
+      ;; internal TLS helpers (certificate verification) resolve.
       (catch #t
              (lambda ()
                (module-use! (resolve-module '(web client))
@@ -44,6 +54,8 @@
                (format (current-error-port)
                        "[claude-proxy] gnutls NOT available~%")))
 
+      ;; ── Configuration ───────────────────────────────
+      
       (define port
         (or (and=> (getenv "PROXY_PORT") string->number) 16890))
 
         (string->uri (or (getenv "CLAUDE_PROXY_UPSTREAM")
                          "https://api.deepseek.com/anthropic")))
 
+      ;; Idle timeout in seconds: bounds connect, waiting for
+      ;; upstream response headers, and every gap in a streamed
+      ;; response.  Long generations survive as long as bytes flow.
+      (define idle-timeout
+        (or (and=> (getenv "CLAUDE_PROXY_TIMEOUT") string->number) 60))
+
+      ;; ── Logging and request ids ─────────────────────
+      
       (define (log fmt . args)
         (apply format
                (current-error-port)
                (string-append "[claude-proxy] " fmt "~%") args)
         (force-output (current-error-port)))
 
-      (define (classifier? json)
-        (define blocks
-          (or (and=> (assoc-ref json "system") vector->list)
-              '()))
-        (define (has? prefix)
-          (any (lambda (b)
-                 (and (equal? (assoc-ref b "type") "text")
-                      (let ((t (assoc-ref b "text")))
-                        (and (string? t)
-                             (string-prefix? prefix t))))) blocks))
-        (and (has? "x-anthropic-billing-header:")
-             (has? "You are a security monitor")))
-
-      (define (disable-thinking json)
-        (alist-delete "output_config"
-                      (alist-delete "reasoning_effort"
-                                    (assoc-set! json "thinking"
-                                                '(("type" . "disabled"))))))
-
-      (define (handler req body)
-        (define (ok code str)
-          (values (build-response #:code code
-                                  #:headers '((content-type application/json)))
-                  (string->utf8 str)))
-        (match (request-method req)
-          ('GET (ok 200 "{\"status\":\"ok\"}"))
-          ('POST (catch #t
+      (define (req-log id fmt . args)
+        (apply format
+               (current-error-port)
+               (string-append "[claude-proxy] req=" id " " fmt "~%") args)
+        (force-output (current-error-port)))
+
+      (define (ms-now)
+        (let ((t (gettimeofday)))
+          (+ (* (car t) 1000)
+             (quotient (cdr t) 1000))))
+
+      (define req-lock
+        (make-mutex))
+      (define req-counter
+        0)
+
+      (define (next-req-id!)
+        (lock-mutex req-lock)
+        (set! req-counter
+              (+ req-counter 1))
+        (let ((id (number->string req-counter 16)))
+          (unlock-mutex req-lock) id))
+
+      (define (peer->string addr)
+        (catch #t
+               (lambda ()
+                 (inet-ntop AF_INET
+                            (sockaddr:addr addr)))
+               (lambda _
+                 (number->string (sockaddr:addr addr)))))
+
+      ;; ── Deadline port ───────────────────────────────
+      ;;
+      ;; Every network read and write goes through this custom
+      ;; port.  Before pulling from the underlying socket the port
+      ;; waits on select(2) with an idle deadline; each pull drains
+      ;; everything that is available (a whole TLS record for
+      ;; https upstreams) so select never hides data buffered at a
+      ;; lower layer.
+      
+      (define stage-size
+        65536)
+
+      (define (make-deadline-port sock fdes timeout-s session)
+        "Wrap SOCK (an unbuffered socket port) in a custom port
+whose reads and writes wait for FDES via select, with an idle
+deadline of TIMEOUT-S seconds.  SESSION is #f for plain TCP or a
+GnuTLS session over SOCK.  Throws 'claude-proxy-timeout when the
+deadline passes before any data arrives."
+        (define stage
+          (make-bytevector stage-size))
+        (define stage-start
+          0)
+        (define stage-len
+          0)
+        (define (wait! kind deadline-ms)
+          (let ((left (- deadline-ms
+                         (ms-now))))
+            (when (<= left 0)
+              (throw 'claude-proxy-timeout kind))
+            (let* ((ready (if (eq? kind
+                                   'read)
+                              (select (list fdes)
+                                      '()
+                                      '()
+                                      (/ left 1000.0))
+                              (select '()
+                                      (list fdes)
+                                      '()
+                                      (/ left 1000.0))))
+                   ;; SELECT returns lists of fds; the single fd shows
+                   ;; up in the first list when readable and in the
+                   ;; second when writable.
+                   (got (if (eq? kind
+                                 'read)
+                            (car ready)
+                            (cadr ready))))
+              (if (pair? got) #t
+                  (wait! kind deadline-ms)))))
+        (define (tls-eof? err)
+          ;; Abrupt close after the response body: treat as EOF.
+          (and session
+               (eq? err error/premature-termination)))
+        (define (tls-again? err)
+          ;; GnuTLS asks to be called again: it consumed bytes that do
+          ;; not yet form a complete record (TLS 1.3 post-handshake
+          ;; messages, a partial record).  SELECT cannot see the bytes
+          ;; GnuTLS already buffered, so on error/again wait for more
+          ;; and retry instead of reporting an error.
+          (and session
+               (eq? err error/again)))
+        (define (fill! deadline-ms)
+          (let loop
+            ()
+            (wait! 'read deadline-ms)
+            (let ((n (catch 'gnutls-error
+                            (lambda ()
+                              (if session
+                                  (record-receive! session stage)
+                                  (recv! sock stage)))
+                            (lambda (key err proc . rest)
+                              (cond
+                                ((tls-eof? err)
+                                 0)
+                                ((tls-again? err)
+                                 'again)
+                                (else (apply throw key err proc rest)))))))
+              (cond
+                ((eq? n
+                      'again)
+                 (loop))
+                ((zero? n)
+                 'eof)
+                (else (set! stage-start 0)
+                      (set! stage-len n)
+                      'ok)))))
+        (define (read! bv start count)
+          (cond
+            ((< stage-start stage-len)
+             (let ((n (min count
+                           (- stage-len stage-start))))
+               (bytevector-copy! stage stage-start bv start n)
+               (set! stage-start
+                     (+ stage-start n)) n))
+            (else (let ((deadline (+ (ms-now)
+                                     (* timeout-s 1000))))
+                    (if (eq? (fill! deadline)
+                             'eof) 0
+                        (read! bv start count))))))
+        (define (send-slice! slice)
+          "Write SLICE, retrying when GnuTLS asks for more time."
+          (let attempt
+            ()
+            (wait! 'write
+                   (+ (ms-now)
+                      (* timeout-s 1000)))
+            (catch 'gnutls-error
+                   (lambda ()
+                     (if session
+                         (record-send session slice)
+                         (send sock slice)))
+                   (lambda (key err proc . rest)
+                     (if (tls-again? err)
+                         (attempt)
+                         (apply throw key err proc rest))))))
+        (define (write! bv start count)
+          (let loop
+            ((off 0))
+            (if (= off count) count
+                (let* ((left (- (+ start count)
+                                (+ start off)))
+                       (slice (make-bytevector left))
+                       (n (begin
+                            (bytevector-copy! bv
+                                              (+ start off) slice 0 left)
+                            (send-slice! slice))))
+                  (if (zero? n)
+                      (throw 'claude-proxy-timeout
+                             'write)
+                      (loop (+ off n)))))))
+        (define (close*)
+          (false-if-exception (close-port sock)))
+        (let ((wrapped (make-custom-binary-input/output-port
+                        "claude-proxy-port"
+                        read!
+                        write!
+                        #f
+                        #f
+                        close*)))
+          (set-port-encoding! wrapped "UTF-8") wrapped))
+
+      ;; ── Upstream connection ─────────────────────────
+      
+      (define (wrap-tls sock host)
+        "Perform the (web client) tls-wrap handshake over SOCK for
+HOST, reusing its certificate-verification helpers, but return the
+GnuTLS session so reads can pull whole TLS records via
+record-receive!."
+        (define session
+          (make-session connection-end/client))
+        (if (module-defined? (resolve-interface '(gnutls))
+                             'set-session-server-name!)
+            (set-session-server-name! session server-name-type/dns host))
+        (set-session-transport-fd! session
+                                   (fileno sock))
+        (set-session-default-priority! session)
+        (set-session-priorities! session "NORMAL:%COMPAT:-VERS-SSL3.0")
+        (set-session-credentials! session
+                                  ((@@ (web client)
+                                       make-credendials-with-ca-trust-files)
+                                   ((@@ (web client)
+                                        x509-certificate-directory))))
+        (let loop
+          ((retries 5))
+          (catch 'gnutls-error
+                 (lambda ()
+                   (handshake session))
+                 (lambda (key err proc . rest)
+                   (cond
+                     ((eq? err error/warning-alert-received)
+                      (format (current-error-port)
+                              "warning: TLS warning alert received: ~a~%"
+                              (alert-description->string (alert-get session)))
+                      (handshake session))
+                     (else (if (or (fatal-error? err)
+                                   (zero? retries))
+                               (apply throw key err proc rest)
+                               (begin
+                                 (format (current-error-port)
+                                         "warning: TLS non-fatal error: ~a~%"
+                                         (error->string err))
+                                 (loop (- retries 1)))))))))
+        (catch 'tls-certificate-error
+               (lambda ()
+                 ((@@ (web client) assert-valid-server-certificate)
+                  session host))
+               (lambda args
+                 (false-if-exception (close-port sock))
+                 (apply throw args)))
+        session)
+
+      (define (uri-effective-port url)
+        "The port URL connects to: its explicit port, or the scheme
+        default when none is written (uri-port returns #f then, and
+        guile's own (web client) falls back the same way)."
+        (or (uri-port url)
+            (match (uri-scheme url)
+              ((or 'http
+                   'ws)
+               80)
+              ((or 'https
+                   'wss)
+               443)
+              (_ #f))))
+
+      (define (host-header url)
+        "Return the Host header value for URL as a (host . port)
+pair — guile validates Host values in that shape — with the port
+when it is not the scheme default."
+        (let ((p (uri-effective-port url)))
+          (if (or (and (eq? (uri-scheme url)
+                            'https)
+                       (eqv? p 443))
+                  (and (eq? (uri-scheme url)
+                            'http)
+                       (eqv? p 80)))
+              (cons (uri-host url) #f)
+              (cons (uri-host url) p))))
+
+      (define (open-upstream url)
+        "Connect to URL (with TLS when https) in a worker thread,
+bounded by IDLE-TIMEOUT seconds, and return the deadline-wrapped
+port.  Throws 'claude-proxy-timeout on expiry or rethrows the
+worker's error."
+        (define host
+          (uri-host url))
+        (define svc
+          (number->string (uri-effective-port url)))
+        (define result-mutex
+          (make-mutex))
+        (define result-cv
+          (make-condition-variable))
+        (define result
+          #f)
+        (define sock-cell
+          (list #f))
+        (define (worker)
+          (catch #t
+                 (lambda ()
+                   (let ((sock (socket PF_INET SOCK_STREAM 0)))
+                     (set-car! sock-cell sock)
+                     ;; Backstop for GnuTLS's own C-level reads and
+                     ;; writes; the select-based deadline port is the
+                     ;; primary bound.
+                     (setsockopt sock SOL_SOCKET SO_RCVTIMEO
+                                 (cons idle-timeout 0))
+                     (setsockopt sock SOL_SOCKET SO_SNDTIMEO
+                                 (cons idle-timeout 0))
+                     (connect sock
+                              (vector-ref (car (getaddrinfo host svc)) 4))
+                     (setvbuf sock
+                              'none)
+                     (let ((session (if (eq? (uri-scheme url)
+                                             'https)
+                                        (wrap-tls sock host) #f)))
+                       (lock-mutex result-mutex)
+                       (set! result
+                             (cons 'done
+                                   (make-deadline-port sock
+                                                       (fileno sock)
+                                                       idle-timeout session)))
+                       (signal-condition-variable result-cv)
+                       (unlock-mutex result-mutex))))
+                 (lambda (key . args)
+                   (lock-mutex result-mutex)
+                   (set! result
+                         (cons 'error
+                               (cons key args)))
+                   (signal-condition-variable result-cv)
+                   (unlock-mutex result-mutex))))
+        (define worker-thread
+          (call-with-new-thread worker))
+        ;; Re-check RESULT before waiting: the worker may publish and
+        ;; signal before the wait starts, and a signal with no waiter
+        ;; is lost.  WAIT-CONDITION-VARIABLE returns #t when awakened
+        ;; and #f once the absolute DEADLINE has passed.
+        (lock-mutex result-mutex)
+        (let* ((t0 (gettimeofday))
+               (deadline (cons (+ (car t0) idle-timeout)
+                               (cdr t0)))
+               (done (let loop
+                       ()
+                       (or result
+                           (and (wait-condition-variable result-cv
+                                                         result-mutex deadline)
+                                (loop))))))
+          (unlock-mutex result-mutex)
+          (if done
+              ;; RESULT is (done . port) or (error key args...).
+              (cond
+                ((eq? (car result)
+                      'done)
+                 (cdr result))
+                (else (apply throw
+                             (cadr result)
+                             (cddr result))))
+              (begin
+                ;; The worker may still be blocked in connect(2) or
+                ;; the TLS handshake; closing the socket bounds the
+                ;; damage.  cancel-thread cannot interrupt C-level
+                ;; blocking, so the thread may linger until the
+                ;; kernel gives up — it only publishes a result
+                ;; nobody reads.
+                (cancel-thread worker-thread)
+                (let ((sock (car sock-cell)))
+                  (when (port? sock)
+                    (false-if-exception (close-port sock))))
+                (throw 'claude-proxy-timeout
+                       'connect)))))
+
+      ;; ── Request rewriting ───────────────────────────
+      
+      (define (prepare-payload json id)
+        "Return JSON with thinking disabled on non-streaming
+requests that do not set a thinking field themselves.  DeepSeek
+reasoning models default to thinking on, which stalls non-streaming
+calls behind full reasoning time."
+        (let ((stream? (eq? (assoc-ref json "stream") #t))
+              (thinking (assoc-ref json "thinking")))
+          (cond
+            (thinking (req-log id
+                       "event=thinking action=omit reason=explicit-thinking")
+                      json)
+            (stream? (req-log id "event=thinking action=omit reason=streaming")
+                     json)
+            (else (req-log id
+                   "event=thinking action=inject reason=non-streaming")
+                  (assoc-set! (alist-delete "output_config"
+                                            (alist-delete "reasoning_effort"
+                                                          json)) "thinking"
+                              '(("type" . "disabled")))))))
+
+      ;; ── Client responses ────────────────────────────
+      
+      (define (send-simple id
+                           cp
+                           t0
+                           code
+                           reason
+                           content-type
+                           body-str)
+        (let ((body (string->utf8 body-str)))
+          (put-string cp
+                      (string-append (format #f "HTTP/1.1 ~d ~a\r\n" code
+                                             reason)
+                                     (format #f "Content-Type: ~a\r\n"
+                                             content-type)
+                                     (format #f "Content-Length: ~d\r\n"
+                                             (bytevector-length body))
+                                     "Connection: close\r\n\r\n"))
+          (put-bytevector cp body)
+          (force-output cp)
+          (req-log id "event=response code=~d bytes=~d ms=~d" code
+                   (bytevector-length body)
+                   (- (ms-now) t0))))
+
+      (define (relay id
+                     cp
+                     version
+                     rsp
+                     up
+                     t0)
+        "Stream the upstream response body to the client, chunked
+for HTTP/1.1 clients and raw-then-close for HTTP/1.0.  The body is
+read from the raw upstream port UP and decoded from its own
+framing, so each upstream chunk reaches the client as it arrives."
+        (define status
+          (response-code rsp))
+        (define reason
+          (response-reason-phrase rsp))
+        (define ct
+          (assoc-ref (response-headers rsp)
+                     'content-type))
+        ;; The parsed header value is a list whose car is the
+        ;; media-type symbol, e.g. (text/event-stream).
+        (define sse?
+          (and (pair? ct)
+               (eq? (car ct)
+                    'text/event-stream)))
+        (define buf
+          (make-bytevector 65536))
+        (define total
+          0)
+        (define first?
+          #t)
+        (define (parse-chunk-size line)
+          ;; Chunk size is 1*HEXDIG with an optional extension.
+          (define (hex? c)
+            (or (char-numeric? c)
+                (and (char>=? c #\a)
+                     (char<=? c #\f))
+                (and (char>=? c #\A)
+                     (char<=? c #\F))))
+          (let ((end (or (string-index line
+                                       (lambda (c)
+                                         (not (hex? c))))
+                         (string-length line))))
+            (if (zero? end) #f
+                (string->number (substring line 0 end) 16))))
+        (put-string cp
+                    (string-append "HTTP/1.1 "
+                                   (number->string status) " " reason "\r\n"))
+        ;; Content-type arrives parsed (e.g. (text/event-stream));
+        ;; write-header knows how to serialize that form.
+        (when ct
+          (write-header 'content-type ct cp))
+        (put-string cp
+                    (string-append "Connection: close\r\n"
+                                   (if (equal? version
+                                               '(1 . 1))
+                                       "Transfer-Encoding: chunked\r\n\r\n"
+                                       "\r\n")))
+        (force-output cp)
+        (let ((chunked? (equal? version
+                                '(1 . 1))))
+          (define (emit bv count)
+            "Write COUNT bytes of BV to the client as one chunk
+(HTTP/1.1) or raw (HTTP/1.0), then flush CP itself.  Chunk framing
+is written by hand: make-chunked-output-port cannot relay
+incrementally — its flush only moves bytes into CP's conversion
+buffer, which is transmitted only when CP is flushed or closed."
+            (when chunked?
+              (put-string cp
+                          (string-append (number->string count 16) "\r\n")))
+            (put-bytevector cp bv 0 count)
+            (when chunked?
+              (put-string cp "\r\n"))
+            (force-output cp)
+            (set! total
+                  (+ total count))
+            (when first?
+              (set! first? #f)
+              (req-log id "event=first-byte ms=~d bytes=~d"
+                       (- (ms-now) t0) total)))
+          (define (attempt thunk)
+            "Run THUNK; a deadline expiry logs a mid-stream timeout,
+any other error a mid-stream error.  Both abandon the relay with
+the reason as the throw argument."
+            (catch #t thunk
+                   (lambda (key . args)
+                     (cond
+                       ((eq? key
+                             'claude-proxy-timeout)
+                        (req-log id "event=timeout phase=mid-stream ms=~d"
+                                 (- (ms-now) t0))
+                        (throw 'relay-abort
+                               'claude-proxy-timeout))
+                       (else (req-log id
+                              "event=error phase=mid-stream key=~s ms=~d" key
+                              (- (ms-now) t0))
+                             (throw 'relay-abort key))))))
+          (define (premature!)
+            (req-log id "event=error phase=mid-stream key=premature-eof ms=~d"
+                     (- (ms-now) t0))
+            (throw 'relay-abort
+                   'premature-eof))
+          (define (read-exact! left)
+            "Read exactly LEFT bytes from UP in buffer-sized slices,
+emitting each; an EOF beforehand aborts the relay."
+            (let loop
+              ((left left))
+              (if (zero? left)
+                  'done
+                  (let* ((k (min left
+                                 (bytevector-length buf)))
+                         (n (attempt (lambda ()
+                                       (get-bytevector-n! up buf 0 k)))))
+                    (if (eof-object? n)
+                        (premature!)
+                        (begin
+                          (emit buf n)
+                          (loop (- left n))))))))
+          (define (decode-chunked)
+            "Relay a chunked body one chunk at a time.  Exact-count
+reads keep chunk data from bleeding into the next chunk header."
+            (let loop
+              ()
+              (let* ((line (attempt (lambda ()
+                                      (read-line up))))
+                     (size (and (not (eof-object? line))
+                                (parse-chunk-size line))))
+                (cond
+                  ((eof-object? line)
+                   (premature!))
+                  ((not size)
+                   (attempt (lambda ()
+                              (throw 'bad-chunk-size))))
+                  ((zero? size)
+                   ;; Final chunk: swallow the closing CRLF; any
+                   ;; trailers are ignored, the connection closes.
+                   (attempt (lambda ()
+                              (get-bytevector-n! up buf 0 2)))
+                   'done)
+                  (else (read-exact! size)
+                        ;; The CRLF that terminates each chunk.
+                        (attempt (lambda ()
+                                   (get-bytevector-n! up buf 0 2)))
+                        (loop))))))
+          (define (decode-framed)
+            (cond
+              ((response-must-not-include-body? rsp)
+               'done)
+              ;; Response-transfer-encoding holds a list of codings,
+              ;; e.g. ((chunked)); match guile's response-body-port.
+              ((member '(chunked)
+                       (response-transfer-encoding rsp))
+               (decode-chunked))
+              ((response-content-length rsp)
+               =>
+               (lambda (len)
+                 (read-exact! len)))
+              (else
+               ;; No framing at all: relay bytes until the upstream
+               ;; closes the connection.
+               (let loop
+                 ()
+                 (let ((chunk (attempt (lambda ()
+                                         (get-bytevector-some up)))))
+                   (if (eof-object? chunk)
+                       'done
+                       (begin
+                         (emit chunk
+                               (bytevector-length chunk))
+                         (loop))))))))
+          (define (abort-message reason)
+            (case reason
+              ((claude-proxy-timeout)
+               "upstream idle timeout")
+              ((premature-eof)
+               "upstream connection closed")
+              ((bad-chunk-size)
+               "upstream sent invalid chunked framing")
+              (else "upstream error")))
+          (define (emit-error-frame reason)
+            "Terminate an SSE body with the Anthropic-style error
+event: a leading newline keeps the frame line-aligned even when the
+upstream stalled mid-frame.  The body then ends normally, so the
+client sees a typed error instead of a truncated stream."
+            (let* ((frame (string-append
+                           "\nevent: error\ndata: {\"type\":\"error\","
+                           "\"error\":{\"type\":\"api_error\",\"message\":\""
+                           (abort-message reason) "\"}}\n\n"))
+                   (b (string->utf8 frame)))
+              (emit b
+                    (bytevector-length b))))
+          (define (try-error-frame reason)
+            "Emit the typed SSE error frame for REASON, surviving a
+client that vanished mid-abort.  Returns #t on success."
+            (catch #t
+                   (lambda ()
+                     (emit-error-frame reason) #t)
+                   (lambda (key . args)
+                     (req-log id "event=error phase=client-write key=~s" key)
+                     #f)))
+          (let ((clean? (catch 'relay-abort
+                               (lambda ()
+                                 (decode-framed)
+                                 (req-log id
+                                  "event=done status=~d bytes=~d ms=~d" status
+                                  total
+                                  (- (ms-now) t0)) #t)
+                               ;; The throw argument is the abort
+                               ;; reason; see ATTEMPT and PREMATURE!.
+                               (lambda (key reason)
+                                 (req-log id
+                                  "event=aborted status=~d bytes=~d ms=~d"
+                                  status total
+                                  (- (ms-now) t0))
+                                 ;; SSE clients get a typed error event
+                                 ;; and a well-formed body; everything
+                                 ;; else keeps the hard truncation.
+                                 (if sse?
+                                     (try-error-frame reason) #f)))))
+            (if clean?
+                ;; Clean end: HTTP/1.1 bodies close with the final
+                ;; zero chunk; HTTP/1.0 clients end at the close.
+                (if chunked?
+                    (begin
+                      (false-if-exception (put-string cp "0\r\n\r\n"))
+                      (force-output cp)
+                      (false-if-exception (close-port cp)))
+                    (false-if-exception (close-port cp)))
+                ;; Non-SSE abort: close without the terminating chunk
+                ;; so the client sees an incomplete body.
+                (false-if-exception (close-port cp))))))
+
+      ;; ── Request handling ────────────────────────────
+      
+      (define (call-upstream id
+                             cp
+                             version
+                             t0
+                             url
+                             hdrs
+                             payload)
+        ;; Only the response is kept; the body port http-request
+        ;; returns is discarded.  UP is read raw: RELAY decodes the
+        ;; framing itself, since that decoder port would swallow
+        ;; whole small bodies before yielding a single byte.
+        (let* ((up (open-upstream url))
+               (request (lambda ()
+                          (http-request url
+                                        #:method 'POST
+                                        #:headers hdrs
+                                        #:body payload
+                                        #:port up
+                                        #:streaming? #t
+                                        #:decode-body? #f
+                                        #:keep-alive? #f)))
+               (relay-response (lambda (rsp body)
+                                 (req-log id
+                                  "event=upstream-headers status=~d ms=~d"
+                                  (response-code rsp)
+                                  (- (ms-now) t0))
+                                 (relay id
+                                        cp
+                                        version
+                                        rsp
+                                        up
+                                        t0))))
+          (dynamic-wind (lambda ()
+                          #f)
+                        (lambda ()
+                          (req-log id "event=upstream-connect ms=~d"
+                                   (- (ms-now) t0))
+                          (call-with-values request relay-response))
                         (lambda ()
-                          (let* ((json (json-string->scm (utf8->string body)))
-                                 (is? (classifier? json))
-                                 (to-send (if is?
-                                              (disable-thinking json) json))
-                                 (payload (string->utf8 (scm->json-string
-                                                         to-send
-                                                         #:unicode #t)))
-                                 (path (uri-path (request-uri req)))
-                                 (base-path (uri-path upstream))
-                                 (url (build-uri (uri-scheme upstream)
-                                                 #:host (uri-host upstream)
-                                                 #:port (uri-port upstream)
-                                                 #:path (string-append
-                                                         base-path path)))
-                                 (auth (assoc-ref (request-headers req)
-                                                  'x-api-key))
-                                 (hdrs (list (list 'content-type
-                                                   'application/json)
-                                             (cons 'x-api-key auth))))
-                            (log "~a ~a [~a~a]"
-                                 'POST path
-                                 (if is? "classifier" "pass")
-                                 (if is? ", thinking-disabled" ""))
-                            (receive (rsp body-port)
-                                     (http-request url
-                                                   #:method 'POST
-                                                   #:headers hdrs
-                                                   #:body payload
-                                                   #:streaming? #f
-                                                   #:decode-body? #f)
-                                     (let* ((code (response-code rsp))
-                                            (ct (assoc-ref (response-headers
-                                                            rsp)
-                                                           'content-type))
-                                            (rhdrs (if ct
-                                                       (list (cons 'content-type
-                                                                   ct))
-                                                       '())))
-                                       (log " -> ~d" code)
-                                       (values (build-response #:code code
-                                                               #:headers rhdrs)
-                                               body-port)))))
-                        (lambda (k . a)
-                          (log "~a ~a error: ~a"
-                               'POST
-                               (uri-path (request-uri req)) k)
-                          (ok 502 "{\"error\":\"upstream unavailable\"}"))))
-          (_ (ok 405 "{\"error\":\"method not allowed\"}"))))
-
-      (log "listening on 127.0.0.1:~d" port)
-      (run-server handler
-                  'http
-                  `(#:port ,port
-                    #:addr ,INADDR_ANY))))
+                          (false-if-exception (close-port up))))))
+
+      (define (forward id
+                       request
+                       body
+                       cp
+                       version
+                       t0)
+        (define timeout-body
+          "{\"error\":\"upstream timeout\"}")
+        (define unavailable-body
+          "{\"error\":\"upstream unavailable\"}")
+        (define (fail code reason msg)
+          (send-simple id
+                       cp
+                       t0
+                       code
+                       reason
+                       "application/json"
+                       msg))
+        (define (try-upstream json)
+          (let* ((prepared (prepare-payload json id))
+                 (payload (string->utf8 (scm->json-string prepared
+                                                          #:unicode #t)))
+                 (rpath (uri-path (request-uri request)))
+                 (rquery (uri-query (request-uri request)))
+                 (base-path (uri-path upstream))
+                 (url (build-uri (uri-scheme upstream)
+                                 #:host (uri-host upstream)
+                                 #:port (uri-port upstream)
+                                 #:path (string-append base-path rpath)
+                                 #:query rquery))
+                 (auth (assoc-ref (request-headers request)
+                                  'x-api-key))
+                 (av (assoc-ref (request-headers request)
+                                'anthropic-version))
+                 (hdrs (append (list (cons 'host
+                                           (host-header upstream))
+                                     (cons 'content-type
+                                           (list 'application/json)))
+                               (if auth
+                                   (list (cons 'x-api-key auth))
+                                   '())
+                               (if av
+                                   (list (cons 'anthropic-version av))
+                                   '()))))
+            (req-log id "event=upstream url=~a ms=~d"
+                     (uri->string url)
+                     (- (ms-now) t0))
+            (catch #t
+                   (lambda ()
+                     (call-upstream id
+                                    cp
+                                    version
+                                    t0
+                                    url
+                                    hdrs
+                                    payload))
+                   (lambda (key . args)
+                     (cond
+                       ((eq? key
+                             'claude-proxy-timeout)
+                        (req-log id "event=timeout phase=headers ms=~d"
+                                 (- (ms-now) t0))
+                        (fail 504 "Gateway Timeout" timeout-body))
+                       (else (req-log id
+                              "event=error phase=headers key=~s ms=~d" key
+                              (- (ms-now) t0))
+                             (req-log id "event=error phase=headers msg=~s"
+                                      args)
+                             (fail 502 "Bad Gateway" unavailable-body)))))))
+        (let ((json (catch #t
+                           (lambda ()
+                             (json-string->scm (utf8->string body)))
+                           (lambda _
+                             #f))))
+          (if (not json)
+              (fail 400 "Bad Request" "{\"error\":\"invalid json\"}")
+              (try-upstream json))))
+
+      (define max-body
+        (expt 2 24))
+      ;; 16 MiB request-body cap
+      
+      (define (body-allowed? request)
+        ;; read-headers decodes content-length to an integer; clients
+        ;; that omit it leave no header at all.
+        (let ((cl (assoc-ref (request-headers request)
+                             'content-length)))
+          (or (not cl)
+              (<= (if (number? cl) cl
+                      (or (string->number cl) 0)) max-body))))
+
+      (define (read-request* id cp t0)
+        (catch #t
+               (lambda ()
+                 (read-request cp))
+               (lambda (key . args)
+                 (cond
+                   ((eq? key
+                         'claude-proxy-timeout)
+                    (req-log id "event=timeout phase=client-read ms=~d"
+                             (- (ms-now) t0)) #f)
+                   ((any eof-object? args)
+                    'eof)
+                   (else (send-simple id
+                                      cp
+                                      t0
+                                      400
+                                      "Bad Request"
+                                      "application/json"
+                                      "{\"error\":\"bad request\"}") #f)))))
+
+      (define (read-body* id cp t0 request)
+        (catch #t
+               (lambda ()
+                 (read-request-body request))
+               (lambda (key . args)
+                 (cond
+                   ((eq? key
+                         'claude-proxy-timeout)
+                    (req-log id "event=timeout phase=client-read ms=~d"
+                             (- (ms-now) t0)) #f)
+                   (else (throw key args))))))
+
+      (define (serve-one id cp t0)
+        (let ((request (read-request* id cp t0)))
+          (when request
+            (let ((method (request-method request))
+                  (version (request-version request)))
+              (req-log id "event=request method=~s path=~a version=~s" method
+                       (uri-path (request-uri request)) version)
+              (cond
+                ((eq? method
+                      'GET)
+                 (send-simple id
+                              cp
+                              t0
+                              200
+                              "OK"
+                              "application/json"
+                              "{\"status\":\"ok\"}"))
+                ((not (eq? method
+                           'POST))
+                 (send-simple id
+                              cp
+                              t0
+                              405
+                              "Method Not Allowed"
+                              "application/json"
+                              "{\"error\":\"method not allowed\"}"))
+                ((not (body-allowed? request))
+                 (send-simple id
+                              cp
+                              t0
+                              413
+                              "Payload Too Large"
+                              "application/json"
+                              "{\"error\":\"payload too large\"}"))
+                (else (let ((body (read-body* id cp t0 request)))
+                        (when body
+                          (forward id
+                                   request
+                                   body
+                                   cp
+                                   version
+                                   t0)))))))))
+
+      ;; ── Server ──────────────────────────────────────
+      
+      (define (handle-connection client-sock peer)
+        (let ((id (next-req-id!))
+              (t0 (ms-now)))
+          (req-log id "event=accept peer=~a"
+                   (peer->string peer))
+          (catch #t
+                 (lambda ()
+                   (setvbuf client-sock
+                            'none)
+                   (let ((cp (make-deadline-port client-sock
+                                                 (fileno client-sock)
+                                                 idle-timeout #f)))
+                     (dynamic-wind (lambda ()
+                                     #f)
+                                   (lambda ()
+                                     (serve-one id cp t0))
+                                   (lambda ()
+                                     (false-if-exception (close-port cp))))))
+                 (lambda (key . args)
+                   (req-log id "event=error phase=connection key=~s" key)
+                   (req-log id "event=error phase=connection msg=~s"
+                            (and (pair? args)
+                                 (car args)))
+                   (req-log id "event=error phase=connection ms=~d"
+                            (- (ms-now) t0))))))
+
+      ;; Orphan watchdog cadence in seconds.
+      (define watch-interval
+        (or (and=> (getenv "CLAUDE_PROXY_WATCHDOG_SECONDS") string->number) 15))
+
+      (define (run-proxy)
+        (sigaction SIGPIPE SIG_IGN)
+        (let ((listener (socket PF_INET SOCK_STREAM 0))
+              (start-ppid (getppid)))
+          ;; One handler for TERM and INT: log, close the listener,
+          ;; die cleanly.  PRIMITIVE-EXIT, not EXIT: a handler can run
+          ;; on any thread, and EXIT only unwinds its own thread (it
+          ;; is a THROW to 'quit).
+          (for-each (lambda (sig)
+                      (sigaction sig
+                                 (lambda (sig)
+                                   (log "shutdown signal=~d ppid=~d" sig
+                                        (getppid))
+                                   (false-if-exception (close-port listener))
+                                   (primitive-exit 0))))
+                    (list SIGTERM SIGINT))
+          (setsockopt listener SOL_SOCKET SO_REUSEADDR 1)
+          (catch 'system-error
+                 (lambda ()
+                   (bind listener
+                         (make-socket-address AF_INET INADDR_LOOPBACK port))
+                   (listen listener 64))
+                 (lambda args
+                   ;; One line, no backtrace; shepherd respawns and
+                   ;; retries until the port is free again.
+                   (log "bind failed port=~d: ~a (errno ~d); exiting" port
+                        (strerror (system-error-errno args))
+                        (system-error-errno args))
+                   (primitive-exit 1)))
+          (log "startup listen=~d ppid=~d upstream=~a timeout=~d" port
+               start-ppid
+               (uri->string upstream) idle-timeout)
+          ;; Orphan watchdog: when the launching shepherd dies without
+          ;; stopping us (guix home reconfigure kills the daemon), we
+          ;; are reparented; free the port so the next daemon's
+          ;; respawn can bind.  Dormant when the parent is PID 1.
+          (when (> start-ppid 1)
+            (call-with-new-thread (lambda ()
+                                    (let loop
+                                      ()
+                                      (sleep watch-interval)
+                                      (if (= (getppid) start-ppid)
+                                          (loop)
+                                          (begin
+                                            (log
+                                             "orphaned parent ~d -> ~d; exiting"
+                                             start-ppid
+                                             (getppid))
+                                            (primitive-exit 0)))))))
+          (let loop
+            ()
+            ;; SELECT, not ACCEPT: a blocked select runs pending
+            ;; signal asyncs on EINTR, so the TERM/INT handlers
+            ;; fire even while the server idles in the kernel
+            ;; wait (a blocked accept swallows the interrupt).
+            (when (pair? (car (select (list listener)
+                                      '()
+                                      '() 60.0)))
+              (let ((client (accept listener)))
+                ;; One thread per connection; handle-connection takes
+                ;; care of its own errors and cleanup.
+                (call-with-new-thread (lambda ()
+                                        (handle-connection (car client)
+                                                           (cdr client))))))
+            (loop))))
+
+      (run-proxy)))
 
 (define claude-proxy-script
   (program-file "claude-proxy" claude-proxy-code))
                     (setenv "PROXY_PORT" "16890")
                     (setenv "CLAUDE_PROXY_UPSTREAM"
                             "https://api.deepseek.com/anthropic")
+                    (setenv "CLAUDE_PROXY_TIMEOUT" "60")
+                    ;; Off-path modules: (gnutls) and (json) live in
+                    ;; these store directories, which shepherd does
+                    ;; not put on GUILE_LOAD_PATH.
+                    (setenv "GUILE_LOAD_PATH"
+                            (string-append #$(file-append guile-gnutls
+                                              "/share/guile/site/3.0") ":"
+                                           #$(file-append guile-json-4
+                                              "/share/guile/site/3.0")))
+                    (setenv "GUILE_LOAD_COMPILED_PATH"
+                            (string-append #$(file-append guile-gnutls
+                                              "/lib/guile/3.0/site-ccache")
+                                           ":"
+                                           #$(file-append guile-json-4
+                                              "/lib/guile/3.0/site-ccache")))
                     (apply execl
                            #$claude-proxy-script "claude-proxy"
                            '()))))
 (define claude-proxy-service
   (shepherd-service (provision '(claude-proxy))
                     (start claude-proxy-start)
-                    (stop #~(make-kill-destructor))
+                    ;; Explicit TERM with a short SIGKILL escalation:
+                    ;; the proxy's signal handler exits promptly, and
+                    ;; the grace period covers any handler latency.
+                    (stop #~(make-kill-destructor SIGTERM
+                                                  #:grace-period 3))
+                    ;; Respawn with a 5-second pace: a guix home
+                    ;; reconfigure overlaps two shepherd daemons for
+                    ;; tens of seconds while the old proxy still holds
+                    ;; the port, and slow retries keep the default
+                    ;; 5-in-7s respawn limit from disabling the
+                    ;; service before the port frees.  (respawn-limit
+                    ;; cannot express that here: guix emits the value
+                    ;; unquoted into shepherd.conf, so a pair is a
+                    ;; syntax error and an integer crashes shepherd's
+                    ;; respawn bookkeeping.)
+                    (respawn? #t)
+                    (respawn-delay 5)
                     (auto-start? #t)
                     (documentation "Shepherd service for claude-proxy.")))
 
diff --git a/tests/claude-proxy-test-cert.pem b/tests/claude-proxy-test-cert.pem
new file mode 100644 (file)
index 0000000..5b72072
--- /dev/null
@@ -0,0 +1,19 @@
+-----BEGIN CERTIFICATE-----
+MIIDHDCCAgSgAwIBAgIUOX0m6IigoZ9XlXsv4tFfa8LcbaQwDQYJKoZIhvcNAQEL
+BQAwFDESMBAGA1UEAwwJMTI3LjAuMC4xMCAXDTI2MDkwMzE3NTQyOFoYDzIxMjYw
+ODEwMTc1NDI4WjAUMRIwEAYDVQQDDAkxMjcuMC4wLjEwggEiMA0GCSqGSIb3DQEB
+AQUAA4IBDwAwggEKAoIBAQC7OSYfhTQwzUSYcY2QPm0PoMiuysvlEpd+18b/b8pO
+T7vBA+w/0QPPtExrX4TbO9nZox7/SW9+WOigcsy2gaTRs5vCsKCpPQlH5atao4Pv
+cLJ3uESt+XqwMpowIp/GXFp3NL3j0q0rcrvisa2pcOOpLiVi2lmAVqmhKXVL+ITG
+u8HtDnw4duTErXFVQBRslkwrbcwW1itiiqdJXZ6oz18u5aQ/e5rRV5Rk7Vojxzus
+lN4vqpDTmUiPVk5aCQrOdgiTcIYEKI8Sy4rtGyBgKmAUEBzqIu8iphshlKWJ5B/N
+i9PQsQV+Zyx8rY72tlpj3NLcg6KhW8RnJ2qkV9QhFTpLAgMBAAGjZDBiMB0GA1Ud
+DgQWBBRKBG6nhcv95RLD+OfDnNFg+Y5pNjAfBgNVHSMEGDAWgBRKBG6nhcv95RLD
++OfDnNFg+Y5pNjAPBgNVHREECDAGhwR/AAABMA8GA1UdEwEB/wQFMAMBAf8wDQYJ
+KoZIhvcNAQELBQADggEBACrqbJ7HzFIBecusXnh+MpiyZ9y1VLwvMNe7MqdBY5Qt
+w/d2UxS+jgyyt/TR4V7y4ZsfWZdtLNCK414+3eiQqHXjD2O0Yc2upoC0lziahYNU
+O68ehOFUyrf1HTpCbBZDzxxKdvCekCEbE5nj6a4Gv+OTqiEt4+kHcBkS8fmiZwdc
+RQuuMdel3BpBBwb0o9a+SiKV+3MOr2BD8BnA/9u0lEo1xONBY6K4yikWp9lztmqw
+vuuyCM+VZg5tsm1CQISIdthoAr9bVFsO1Zc2ER6HAWy+yJBXdF12Juj7bBHFS7a6
+9l8PDR6eSWU9ebKvCAcfN46Vrm4CIhr2UPH3DVHRIWw=
+-----END CERTIFICATE-----
diff --git a/tests/claude-proxy-test-key.pem b/tests/claude-proxy-test-key.pem
new file mode 100644 (file)
index 0000000..d0c0183
--- /dev/null
@@ -0,0 +1,28 @@
+-----BEGIN PRIVATE KEY-----
+MIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQC7OSYfhTQwzUSY
+cY2QPm0PoMiuysvlEpd+18b/b8pOT7vBA+w/0QPPtExrX4TbO9nZox7/SW9+WOig
+csy2gaTRs5vCsKCpPQlH5atao4PvcLJ3uESt+XqwMpowIp/GXFp3NL3j0q0rcrvi
+sa2pcOOpLiVi2lmAVqmhKXVL+ITGu8HtDnw4duTErXFVQBRslkwrbcwW1itiiqdJ
+XZ6oz18u5aQ/e5rRV5Rk7VojxzuslN4vqpDTmUiPVk5aCQrOdgiTcIYEKI8Sy4rt
+GyBgKmAUEBzqIu8iphshlKWJ5B/Ni9PQsQV+Zyx8rY72tlpj3NLcg6KhW8RnJ2qk
+V9QhFTpLAgMBAAECggEAMlFGLzOAFtucI3JlTw6QBiK6vMtxKBQqliMM7wrO7uQb
+/GB/Bpm4sgJthXALB1bbElR2WLsWVXk0sCaaYTSPpPJmGtmYnFN0opeVyYrnwxrh
+RC7Ieo9xy1gWO3iaw1d/0sjgfhHZi7NOzrrdBwx5izcFQj+XzVe9SDyZszzMBpLp
+ma80qvyNw34RHQEQpUTvl1fXuIh1VH42SUwIRQtMrUYpyoEQOJfc7Cwy0/dzMnt+
+XCF9mXu9wyzrfXvnGJv2TjobzqhUx1JkqHiW5TGI/3cKpTecrT6WMUGzt0v3xTS4
+TxGhmnCy4qE5xT9po29k94vulBhVD9PT4/1OaSke9QKBgQD/GkaxIKYqqtPGnqf8
+Li4GKJspMuhbF7gM3LCqIqM6wgUlIazBPVedk9BTXlMiYQAQYBlp5Rh6xxxh8+Dg
+UFVMPFq2LAIb41IEV51efkZNgGnVC0zeyLYmI0tVBJEp4A8itqgerizD4Lbg9sBP
+7q5cClgu12LSaFMZuXmlRaIcDwKBgQC74b8Pc8y70LsKoaIQJLALwOONpJuPCLDj
+B89mh7RtYRH0slB8QiJYRNOVp9mw9FYZFQkwBBbCCiTagjtBcTLiDc/OkZdln2pa
+KWzRztDN97Ijnj3zCac8SyrJHoUu12yUjzQ2GrC2kQSEjj/RkpsQ3vYd6aEkKmAw
+6mKqKNRyBQKBgQCDxtEJorHzjHcFPOPN0xUXPVaZe6CnnaMHkeD4ohfrvFnoCnMx
+Bz0BO1/8ENelBLKBwwKdvyhcFArHVrGtbhIB5ZN+U1FrkovFjnTOYTBzzIfe841r
+8AaXwNejPU63cPSgm/ZQkuyw6p3Nq+k/4S3Ugct3tu9nfVigCz2ZcFUDZQKBgEec
+JGlsVqVjSlckAhQrF5pzO4gaLFxZEKqHqIpIwQFLlT9x03F494QzP330CuoCRuqq
+dOUDOfVdTmymZJVt4tn8L69pGI5YM34H+f0B2d4XQaOHxc7jaAV4FOexJUwUOcNp
+zZmtlJsRLOqlGTf0q/vDL4V5Lb0OFbmvLEn04/xNAoGBAIfcx+E8sdVvs0gIrOv5
+wsmyx0Mcr10GF01QlePxgBgZNfWIVuZomGc82gT606K5R6L4R129Jzvs087gBLVK
+C/j9JqbFZ146gKD91lQMinGffCGa9XcXLzLph2MCeAyO46z5jNqSuWk0TL3tP8o7
+XCrVbiudrikX2WtdtF3aNz02
+-----END PRIVATE KEY-----
index 79285ba6af2c973d88fba9ba57555773d01dce10..6ce2f10235382fcf49318171071bdfd4b6bb3be8 100644 (file)
 ;;; Copyright (c) 2026 Jakub Czajka <jakub@ekhem.eu.org>
 ;;; License: GPL-3.0 or later.
 ;;;
-;;; claude-proxy end-to-end test — validates HTTP responses from the
-;;; Shepherd-started proxy (health endpoint and error handling).
+;;; claude-proxy end-to-end tests.  The proxy runs against a fake
+;;; upstream (127.0.0.1:19998) that can stall and stream with gaps,
+;;; so incremental relaying, thinking injection, timeouts and
+;;; log traceability can all be asserted deterministically.  A TLS
+;;; twin pair — an HTTPS fake on 127.0.0.1:443 with a self-signed
+;;; test certificate, and a second proxy instance pointed at the
+;;; portless URL https://127.0.0.1 — guards the scheme-default-port
+;;; and TLS record handling regressions hermetically.
 
 (define-module (tests claude-proxy)
   #:use-module (conf common claude-proxy)
+  #:use-module (gnu packages guile)
+  #:use-module (gnu packages tls)
   #:use-module (gnu services)
   #:use-module (gnu services shepherd)
   #:use-module (gnu system)
   #:use-module (guix gexp)
   #:use-module (tests common)
-  #:export (claude-proxy-test-cases claude-proxy-test-os-services))
+  #:export (fake-upstream-service claude-proxy-test-service
+                                  fake-upstream-tls-service
+                                  claude-proxy-tls-test-service
+                                  claude-proxy-test-os-services
+                                  claude-proxy-test-cases))
+
+;;;
+;;; Fake upstream — a raw-socket HTTP/1.1 server on 127.0.0.1:19998.
+;;; Routes:
+;;;   POST /v1/messages with "stream":true    — chunked SSE: one frame,
+;;;     sleep 3 s, two more frames + [DONE] (the 3 s gap proves the
+;;;     proxy relays incrementally instead of buffering).
+;;;   POST /v1/messages otherwise             — 200 JSON echoing the
+;;;     received body (asserts thinking injection).
+;;;   POST /slow-headers                      — sleep 15 s, then reply.
+;;;   POST /mid-stall                         — chunked headers + one
+;;;     frame, then sleep 15 s.
+;;;   anything else                           — 401 authentication_error.
+;;;
+
+(define fake-upstream-code
+  #~(begin
+      (use-modules (ice-9 match)
+                   (ice-9 rdelim)
+                   (ice-9 threads)
+                   (rnrs bytevectors)
+                   (rnrs io ports)
+                   (srfi srfi-13))
+      (define (read-headers sock)
+        "Read header lines until the blank line; return an alist of
+lowercased names to values.  Right-trim each line: guile read-line
+keeps the \\r of CRLF endings, which would otherwise make the blank
+line read as \"\\r\" and the header scan run on into the body."
+        (let loop
+          ((acc '()))
+          (let* ((raw (read-line sock))
+                 (line (if (eof-object? raw) raw
+                           (string-trim-right raw))))
+            (if (or (eof-object? line)
+                    (string-null? line)) acc
+                (let ((i (string-index line #\:)))
+                  (if i
+                      (loop (acons (string-downcase (substring line 0 i))
+                                   (string-trim (substring line
+                                                           (+ i 1))) acc))
+                      (loop acc)))))))
+      (define (respond sock code reason content-type body)
+        (display (string-append "HTTP/1.1 "
+                                code
+                                " "
+                                reason
+                                "\r\n"
+                                "Content-Type: "
+                                content-type
+                                "\r\n"
+                                "Content-Length: "
+                                (number->string (string-length body))
+                                "\r\n"
+                                "Connection: close\r\n\r\n"
+                                body) sock)
+        (force-output sock))
+      (define (chunk sock str)
+        "Write STR as one HTTP chunk."
+        (let ((b (string->utf8 str)))
+          (display (number->string (bytevector-length b) 16) sock)
+          (display "\r\n" sock)
+          (put-bytevector sock b)
+          (display "\r\n" sock)
+          (force-output sock)))
+      (define (respond-stream sock)
+        "Stream three SSE frames with a 3-second gap after the first,
+then terminate the chunked body."
+        (display "HTTP/1.1 200 OK\r\n" sock)
+        (display "Content-Type: text/event-stream\r\n" sock)
+        (display "Transfer-Encoding: chunked\r\n" sock)
+        (display "Connection: close\r\n\r\n" sock)
+        (force-output sock)
+        (chunk sock
+               (string-append "event: message_start\n"
+                              "data: {\"type\":\"message_start\","
+                              "\"message\":{\"id\":\"msg_01\","
+                              "\"content\":[]}}\n\n"))
+        (sleep 3)
+        (chunk sock
+               (string-append "event: content_block_delta\n"
+                              "data: {\"type\":\"content_block_delta\","
+                              "\"index\":0,"
+                              "\"delta\":{\"type\":\"text_delta\","
+                              "\"text\":\"Hello from \"}}\n\n"))
+        (chunk sock
+               (string-append "event: content_block_delta\n"
+                              "data: {\"type\":\"content_block_delta\","
+                              "\"index\":0,"
+                              "\"delta\":{\"type\":\"text_delta\","
+                              "\"text\":\"the fake\"}}\n\n"
+                              "event: message_delta\n"
+                              "data: {\"type\":\"message_delta\","
+                              "\"delta\":{\"stop_reason\":\"end_turn\"}}\n\n"
+                              "event: message_stop\n"
+                              "data: {\"type\":\"message_stop\"}\n\n"
+                              "data: [DONE]\n\n"))
+        (display "0\r\n\r\n" sock)
+        (force-output sock))
+      (define (respond-mid-stall sock)
+        "Reply with headers and one frame, then stall forever (until
+the proxy gives up and closes)."
+        (display "HTTP/1.1 200 OK\r\n" sock)
+        (display "Content-Type: text/event-stream\r\n" sock)
+        (display "Transfer-Encoding: chunked\r\n" sock)
+        (display "Connection: close\r\n\r\n" sock)
+        (force-output sock)
+        (chunk sock "event: ping\ndata: {\"type\":\"ping\"}\n\n")
+        (sleep 15)
+        (false-if-exception (chunk sock "event: message_stop
+data: {\"type\":\"message_stop\"}
+
+"))
+        (false-if-exception (begin
+                              (display "0\r\n\r\n" sock)
+                              (force-output sock))))
+      (define (handle sock)
+        (catch #t
+               (lambda ()
+                 (let* ((line (read-line sock))
+                        (parts (and (not (eof-object? line))
+                                    (string-split line #\space)))
+                        (path (and (>= (length parts) 2)
+                                   (list-ref parts 1)))
+                        (hdrs (read-headers sock))
+                        (cl (assoc-ref hdrs "content-length"))
+                        (body (if cl
+                                  (get-bytevector-n sock
+                                                    (string->number cl))
+                                  #vu8()))
+                        (body-str (utf8->string body)))
+                   (cond
+                     ((and (equal? path "/v1/messages")
+                           (string-contains body-str "\"stream\":true"))
+                      (respond-stream sock))
+                     ((equal? path "/v1/messages")
+                      (respond sock "200" "OK" "application/json" body-str))
+                     ((equal? path "/slow-headers")
+                      (sleep 15)
+                      (respond sock "200" "OK" "application/json"
+                               "{\"done\":true}"))
+                     ((equal? path "/mid-stall")
+                      (respond-mid-stall sock))
+                     (else (respond sock "401" "Unauthorized"
+                                    "application/json"
+                                    (string-append
+                                     "{\"type\":\"error\",\"error\":"
+                                     "{\"type\":\"authentication_error\","
+                                     "\"message\":\"invalid x-api-key\"}}"))))))
+               (lambda _
+                 #f))
+        (false-if-exception (close-port sock)))
+      ;; Clients (the proxy under test) close connections mid-stream;
+      ;; a write to a closed socket must not kill the fake via SIGPIPE.
+      (sigaction SIGPIPE SIG_IGN)
+      (let ((listener (socket PF_INET SOCK_STREAM 0)))
+        (setsockopt listener SOL_SOCKET SO_REUSEADDR 1)
+        (bind listener
+              (make-socket-address AF_INET INADDR_LOOPBACK 19998))
+        (listen listener 16)
+        (let loop
+          ()
+          (match (accept listener)
+            ((conn . peer) (setvbuf conn
+                                    'none)
+             (call-with-new-thread (lambda ()
+                                     (handle conn)))
+             (loop)))))))
+
+(define fake-upstream-script
+  (program-file "fake-claude-upstream" fake-upstream-code))
+
+(define fake-upstream-start
+  #~(make-forkexec-constructor (list #$fake-upstream-script)
+                               #:log-file "/var/log/fake-upstream.log"))
+
+(define fake-upstream-service
+  (shepherd-service (provision '(fake-upstream))
+                    (start fake-upstream-start)
+                    (stop #~(make-kill-destructor))
+                    (auto-start? #t)))
+
+;;;
+;;; Test proxy — env vars point it at the fake upstream and set a
+;;; 5-second idle timeout, so every stall test completes quickly.
+;;;
+
+(define claude-proxy-test-wrapper
+  (program-file "test-claude-proxy-start"
+                #~(begin
+                    (setenv "PROXY_PORT" "16890")
+                    (setenv "CLAUDE_PROXY_UPSTREAM" "http://127.0.0.1:19998")
+                    (setenv "CLAUDE_PROXY_TIMEOUT" "5")
+                    ;; Off-path modules: (gnutls) and (json) live in
+                    ;; these store directories, which shepherd does
+                    ;; not put on GUILE_LOAD_PATH.
+                    (setenv "GUILE_LOAD_PATH"
+                            (string-append #$(file-append guile-gnutls
+                                              "/share/guile/site/3.0") ":"
+                                           #$(file-append guile-json-4
+                                              "/share/guile/site/3.0")))
+                    (setenv "GUILE_LOAD_COMPILED_PATH"
+                            (string-append #$(file-append guile-gnutls
+                                              "/lib/guile/3.0/site-ccache")
+                                           ":"
+                                           #$(file-append guile-json-4
+                                              "/lib/guile/3.0/site-ccache")))
+                    (apply execl
+                           #$claude-proxy-script "claude-proxy"
+                           '()))))
+
+(define claude-proxy-test-start
+  #~(make-forkexec-constructor (list #$claude-proxy-test-wrapper)
+                               #:log-file "/var/log/claude-proxy.log"))
+
+(define claude-proxy-test-service
+  (shepherd-service (provision '(claude-proxy))
+                    (start claude-proxy-test-start)
+                    (stop #~(make-kill-destructor))
+                    (auto-start? #t)))
+
+;;;
+;;; TLS fake upstream — the HTTPS twin of the plain fake: a GnuTLS
+;;; server session on 127.0.0.1:443, serving the same routes over
+;;; the session record port.  Its certificate and key (self-signed,
+;;; CN=127.0.0.1, SAN IP:127.0.0.1, CA:TRUE) are passed as
+;;; arguments; a third, optional argument overrides the port (used
+;;; for host-side rehearsal, where binding 443 needs root).
+;;; Responding in TLS 1.3 sends a post-handshake NewSessionTicket,
+;;; which is what makes the proxy's record reads hit their
+;;; transient error/again retry path.
+;;;
+
+(define fake-tls-code
+  #~(begin
+      (use-modules (gnutls)
+                   (ice-9 match)
+                   (ice-9 rdelim)
+                   (ice-9 threads)
+                   (rnrs bytevectors)
+                   (rnrs io ports)
+                   (srfi srfi-13))
+      (define (read-headers sock)
+        "Read header lines until the blank line; return an alist of
+lowercased names to values."
+        (let loop
+          ((acc '()))
+          (let* ((raw (read-line sock))
+                 (line (if (eof-object? raw) raw
+                           (string-trim-right raw))))
+            (if (or (eof-object? line)
+                    (string-null? line)) acc
+                (let ((i (string-index line #\:)))
+                  (if i
+                      (loop (acons (string-downcase (substring line 0 i))
+                                   (string-trim (substring line
+                                                           (+ i 1))) acc))
+                      (loop acc)))))))
+      (define (respond sock code reason content-type body)
+        (display (string-append "HTTP/1.1 "
+                                code
+                                " "
+                                reason
+                                "\r\n"
+                                "Content-Type: "
+                                content-type
+                                "\r\n"
+                                "Content-Length: "
+                                (number->string (string-length body))
+                                "\r\n"
+                                "Connection: close\r\n\r\n"
+                                body) sock)
+        (force-output sock))
+      (define (chunk sock str)
+        "Write STR as one HTTP chunk."
+        (let ((b (string->utf8 str)))
+          (display (number->string (bytevector-length b) 16) sock)
+          (display "\r\n" sock)
+          (put-bytevector sock b)
+          (display "\r\n" sock)
+          (force-output sock)))
+      (define (respond-stream sock)
+        "Stream three SSE frames with a 3-second gap after the first,
+then terminate the chunked body."
+        (display "HTTP/1.1 200 OK\r\n" sock)
+        (display "Content-Type: text/event-stream\r\n" sock)
+        (display "Transfer-Encoding: chunked\r\n" sock)
+        (display "Connection: close\r\n\r\n" sock)
+        (force-output sock)
+        (chunk sock
+               (string-append "event: message_start\n"
+                              "data: {\"type\":\"message_start\","
+                              "\"message\":{\"id\":\"msg_01\","
+                              "\"content\":[]}}\n\n"))
+        (sleep 3)
+        (chunk sock
+               (string-append "event: content_block_delta\n"
+                              "data: {\"type\":\"content_block_delta\","
+                              "\"index\":0,"
+                              "\"delta\":{\"type\":\"text_delta\","
+                              "\"text\":\"Hello from \"}}\n\n"))
+        (chunk sock
+               (string-append "event: content_block_delta\n"
+                              "data: {\"type\":\"content_block_delta\","
+                              "\"index\":0,"
+                              "\"delta\":{\"type\":\"text_delta\","
+                              "\"text\":\"the fake\"}}\n\n"
+                              "event: message_delta\n"
+                              "data: {\"type\":\"message_delta\","
+                              "\"delta\":{\"stop_reason\":\"end_turn\"}}\n\n"
+                              "event: message_stop\n"
+                              "data: {\"type\":\"message_stop\"}\n\n"
+                              "data: [DONE]\n\n"))
+        (display "0\r\n\r\n" sock)
+        (force-output sock))
+      (define (handle io)
+        "Dispatch one HTTP request read from IO, mirroring the plain
+fake's routes."
+        (catch #t
+               (lambda ()
+                 (let* ((line (read-line io))
+                        (parts (and (not (eof-object? line))
+                                    (string-split line #\space)))
+                        (path (and (>= (length parts) 2)
+                                   (list-ref parts 1)))
+                        (hdrs (read-headers io))
+                        (cl (assoc-ref hdrs "content-length"))
+                        (body (if cl
+                                  (get-bytevector-n io
+                                                    (string->number cl))
+                                  #vu8()))
+                        (body-str (utf8->string body)))
+                   (cond
+                     ((and (equal? path "/v1/messages")
+                           (string-contains body-str "\"stream\":true"))
+                      (respond-stream io))
+                     ((equal? path "/v1/messages")
+                      (respond io "200" "OK" "application/json" body-str))
+                     (else (respond io "401" "Unauthorized" "application/json"
+                                    (string-append
+                                     "{\"type\":\"error\",\"error\":"
+                                     "{\"type\":\"authentication_error\","
+                                     "\"message\":\"invalid x-api-key\"}}"))))))
+               (lambda _
+                 #f))
+        (false-if-exception (close-port io)))
+      (define (serve conn cert key)
+        "Complete a TLS handshake with CONN, then serve HTTP over
+the session record port."
+        (catch #t
+               (lambda ()
+                 (let* ((session (make-session connection-end/server))
+                        (cred (make-certificate-credentials)))
+                   (set-session-default-priority! session)
+                   (set-session-priorities! session
+                                            "NORMAL:%COMPAT:-VERS-SSL3.0")
+                   (set-session-transport-fd! session
+                                              (fileno conn))
+                   (set-certificate-credentials-x509-key-files! cred cert key
+                    x509-certificate-format/pem)
+                   (set-session-credentials! session cred)
+                   (handshake session)
+                   (handle (session-record-port session))))
+               (lambda _
+                 #f))
+        (false-if-exception (close-port conn)))
+      ;; Clients (the proxy under test) close connections mid-stream;
+      ;; a write to a closed socket must not kill the fake via SIGPIPE.
+      (sigaction SIGPIPE SIG_IGN)
+      (let* ((argv (command-line))
+             (cert (list-ref argv 1))
+             (key (list-ref argv 2))
+             (pnum (if (> (length argv) 3)
+                       (string->number (list-ref argv 3)) 443)))
+        (let ((listener (socket PF_INET SOCK_STREAM 0)))
+          (setsockopt listener SOL_SOCKET SO_REUSEADDR 1)
+          (bind listener
+                (make-socket-address AF_INET INADDR_LOOPBACK pnum))
+          (listen listener 16)
+          (let loop
+            ()
+            (match (accept listener)
+              ((conn . peer) (call-with-new-thread (lambda ()
+                                                     (serve conn cert key)))
+               (loop))))))))
+
+(define fake-upstream-tls-script
+  (program-file "fake-claude-upstream-tls" fake-tls-code))
+
+(define fake-upstream-tls-wrapper
+  ;; Sets the load path for (gnutls) — off the shepherd path — then
+  ;; execs the fake with the certificate and key passed as arguments.
+  (program-file "fake-claude-upstream-tls-start"
+                #~(begin
+                    (setenv "GUILE_LOAD_PATH"
+                            (string-append #$(file-append guile-gnutls
+                                              "/share/guile/site/3.0")))
+                    (setenv "GUILE_LOAD_COMPILED_PATH"
+                            (string-append #$(file-append guile-gnutls
+                                              "/lib/guile/3.0/site-ccache")))
+                    (apply execl
+                           #$fake-upstream-tls-script
+                           "fake-claude-upstream-tls"
+                           (cdr (command-line))))))
+
+(define fake-tls-cert
+  ;; Literal relative names: local-file resolves them against the
+  ;; directory of the module's source file — the worktree's tests/
+  ;; during guix builds.  Nothing touches the file system until the
+  ;; object is lowered, so hook compiles (which load the module from
+  ;; the main checkout, where these files do not exist) stay clean.
+  (local-file "claude-proxy-test-cert.pem" "claude-proxy-test-cert.pem"))
+
+(define fake-tls-key
+  (local-file "claude-proxy-test-key.pem" "claude-proxy-test-key.pem"))
+
+(define claude-proxy-test-ca-dir
+  ;; Directory the proxy scans for CA certificates (see
+  ;; GUILE_TLS_CERTIFICATE_DIRECTORY in web/client.scm).
+  (file-union "claude-proxy-test-ca"
+              `(("ca.pem" ,fake-tls-cert))))
+
+(define fake-upstream-tls-start
+  #~(make-forkexec-constructor (list #$fake-upstream-tls-wrapper
+                                     #$fake-tls-cert
+                                     #$fake-tls-key)
+                               #:log-file "/var/log/fake-upstream-tls.log"))
+
+(define fake-upstream-tls-service
+  (shepherd-service (provision '(fake-upstream-tls))
+                    (start fake-upstream-tls-start)
+                    (stop #~(make-kill-destructor))
+                    (auto-start? #t)))
+
+;;;
+;;; TLS test proxy — same test env as the plain proxy, but pointed
+;;; at the portless URL https://127.0.0.1.  The upstream port is
+;;; then the scheme default (443), which is where the TLS fake
+;;; listens; the proxy must also trust the test certificate via
+;;; GUILE_TLS_CERTIFICATE_DIRECTORY.
+;;;
+
+(define claude-proxy-tls-test-wrapper
+  (program-file "test-claude-proxy-tls-start"
+                #~(begin
+                    (setenv "PROXY_PORT" "16892")
+                    (setenv "CLAUDE_PROXY_UPSTREAM" "https://127.0.0.1")
+                    (setenv "CLAUDE_PROXY_TIMEOUT" "5")
+                    (setenv "GUILE_TLS_CERTIFICATE_DIRECTORY"
+                            #$claude-proxy-test-ca-dir)
+                    (setenv "GUILE_LOAD_PATH"
+                            (string-append #$(file-append guile-gnutls
+                                              "/share/guile/site/3.0") ":"
+                                           #$(file-append guile-json-4
+                                              "/share/guile/site/3.0")))
+                    (setenv "GUILE_LOAD_COMPILED_PATH"
+                            (string-append #$(file-append guile-gnutls
+                                              "/lib/guile/3.0/site-ccache")
+                                           ":"
+                                           #$(file-append guile-json-4
+                                              "/lib/guile/3.0/site-ccache")))
+                    (apply execl
+                           #$claude-proxy-script "claude-proxy-tls"
+                           '()))))
+
+(define claude-proxy-tls-test-start
+  #~(make-forkexec-constructor (list #$claude-proxy-tls-test-wrapper)
+                               #:log-file "/var/log/claude-proxy-tls.log"))
+
+(define claude-proxy-tls-test-service
+  (shepherd-service (provision '(claude-proxy-tls))
+                    (start claude-proxy-tls-test-start)
+                    (stop #~(make-kill-destructor))
+                    (auto-start? #t)))
+
+;;;
+;;; Test bodies.
+;;;
 
 (define %script-exists-body
-  (let ((path (file-append claude-proxy-script "/bin/claude-proxy")))
-    `(catch #t
-            (lambda ()
-              (let ((st (stat ,path)))
-                (not (zero? (logand (stat:mode st) #x40)))))
-            (lambda (key . args)
-              #f))))
+  (let ((path claude-proxy-script))
+    `(begin
+       (catch #t
+              (lambda ()
+                (let ((st (stat ,path)))
+                  (not (zero? (logand (stat:mode st) #x40)))))
+              (lambda (key . args)
+                #f)))))
 
 (define %health-test-body
   '(begin
             (lambda (key . args)
               #f))))
 
-(define %post-502-test-body
+(define %fetch-post-code
+  ;; Evaluates to a closure of (path body rcvtimeo [port-num]) -> list
+  ;; (first-ms total-ms text) or #f.  FIRST-MS is the time until the
+  ;; first byte of the response, TOTAL-MS until the connection ends.
+  ;; PORT-NUM defaults to 16890 (the plain proxy); the TLS tests
+  ;; pass 16892.  EVAL compiles its argument as one expression,
+  ;; where top-level defines are illegal — ms-now is therefore an
+  ;; internal define of the returned lambda.
+  '(begin
+     (use-modules (ice-9 rdelim)
+                  (ice-9 rw))
+     (lambda (path body rcvtimeo . port-arg)
+       (define (ms-now)
+         (let ((t (gettimeofday)))
+           (+ (* 1000
+                 (car t))
+              (quotient (cdr t) 1000))))
+       (catch #t
+              (lambda ()
+                (let* ((pnum (if (pair? port-arg)
+                                 (car port-arg) 16890))
+                       (clen (string-length body))
+                       (req (string-append "POST "
+                                           path
+                                           " HTTP/1.0\r\n"
+                                           "Content-Type: application/json"
+                                           "\r\n"
+                                           "x-api-key: sk-test\r\n"
+                                           "Content-Length: "
+                                           (number->string clen)
+                                           "\r\n\r\n"
+                                           body))
+                       (s (socket PF_INET SOCK_STREAM 0)))
+                  (connect s
+                           (make-socket-address AF_INET INADDR_LOOPBACK pnum))
+                  (setsockopt s SOL_SOCKET SO_RCVTIMEO
+                              (cons rcvtimeo 0))
+                  (let ((t0 (ms-now)))
+                    (display req s)
+                    (force-output s)
+                    (let loop
+                      ((acc '())
+                       (first-ms #f))
+                      (let* ((buf (make-string 8192))
+                             (n (catch #t
+                                       (lambda ()
+                                         (read-string!/partial buf s 0 8192))
+                                       (lambda _
+                                         0))))
+                        ;; READ-STRING!/PARTIAL returns #f at end of
+                        ;; file, not 0.
+                        (if (or (not n)
+                                (zero? n))
+                            (let ((text (string-concatenate-reverse acc)))
+                              (close-port s)
+                              (list first-ms
+                                    (- (ms-now) t0) text))
+                            (loop (cons (substring buf 0 n) acc)
+                                  (or first-ms
+                                      (- (ms-now) t0)))))))))
+              (lambda (key . args)
+                #f)))))
+
+(define %stream-test-body
+  ;; The fake sends its first SSE frame immediately and the rest only
+  ;; after a 3-second sleep.  An incremental relay delivers the first
+  ;; bytes in well under 2.5 s while the fake is still sleeping; a
+  ;; buffering relay would deliver everything at once, 3+ seconds in.
+  `(begin
+     (let* ((fetch (eval ,%fetch-post-code
+                         (interaction-environment)))
+            (r (fetch "/v1/messages"
+                      (string-append
+                       "{\"stream\":true,\"model\":\"deepseek-chat\","
+                       "\"max_tokens\":16,"
+                       "\"messages\":[{\"role\":\"user\","
+                       "\"content\":\"hello\"}]}") 4)))
+       (and r
+            (list? r)
+            (< (car r) 2500)
+            (>= (cadr r) 2500)
+            (string-contains (caddr r) "data: [DONE]")
+            (string-contains (caddr r) "message_start")))))
+
+(define %injection-test-body
+  ;; The fake echoes the body it received.  A non-streaming request
+  ;; without an explicit thinking field must arrive with
+  ;; thinking disabled and reasoning_effort/output_config stripped;
+  ;; an explicit thinking field must pass through untouched.
+  `(begin
+     (use-modules (srfi srfi-13))
+     (let* ((fetch (eval ,%fetch-post-code
+                         (interaction-environment)))
+            (plain (fetch "/v1/messages"
+                          (string-append
+                           "{\"stream\":false,\"model\":\"deepseek-chat\","
+                           "\"max_tokens\":16,\"reasoning_effort\":\"high\","
+                           "\"output_config\":{\"budget\":\"auto\"},"
+                           "\"messages\":[{\"role\":\"user\","
+                           "\"content\":\"hello\"}]}") 4))
+            (explicit (fetch "/v1/messages"
+                             (string-append
+                              "{\"stream\":false,\"model\":\"deepseek-chat\","
+                              "\"max_tokens\":16,"
+                              "\"thinking\":{\"type\":\"enabled\","
+                              "\"budget_tokens\":1000},"
+                              "\"messages\":[{\"role\":\"user\","
+                              "\"content\":\"hello\"}]}") 4)))
+       (and plain
+            (list? plain)
+            (string-contains (caddr plain) "200")
+            (string-contains (caddr plain) "\"type\":\"disabled\"")
+            (not (string-contains (caddr plain) "reasoning_effort"))
+            (not (string-contains (caddr plain) "output_config"))
+            explicit
+            (list? explicit)
+            (string-contains (caddr explicit) "\"type\":\"enabled\"")
+            (string-contains (caddr explicit) "budget_tokens")))))
+
+(define %slow-headers-test-body
+  ;; The fake holds the response headers for 15 s; the proxy's
+  ;; 5-second idle timeout must produce a 504 in about 5 s.
+  `(begin
+     (use-modules (srfi srfi-13))
+     (let* ((fetch (eval ,%fetch-post-code
+                         (interaction-environment)))
+            (r (fetch "/slow-headers"
+                      (string-append
+                       "{\"stream\":true,\"model\":\"deepseek-chat\","
+                       "\"max_tokens\":16,"
+                       "\"messages\":[{\"role\":\"user\","
+                       "\"content\":\"hello\"}]}") 10)))
+       (and r
+            (list? r)
+            (string-contains (caddr r) "504")
+            (string-contains (caddr r) "upstream timeout")
+            (>= (cadr r) 4000)
+            (<= (cadr r) 9000)))))
+
+(define %mid-stall-test-body
+  ;; The fake streams one frame then stalls; the proxy must cut the
+  ;; connection at the 5-second timeout and deliver a typed SSE error
+  ;; event instead of a bare truncation (no [DONE]).
+  `(begin
+     (use-modules (srfi srfi-13))
+     (let* ((fetch (eval ,%fetch-post-code
+                         (interaction-environment)))
+            (r (fetch "/mid-stall"
+                      (string-append
+                       "{\"stream\":true,\"model\":\"deepseek-chat\","
+                       "\"max_tokens\":16,"
+                       "\"messages\":[{\"role\":\"user\","
+                       "\"content\":\"hello\"}]}") 8)))
+       (and r
+            (list? r)
+            (string-contains (caddr r) "200")
+            (string-contains (caddr r) "\"type\":\"ping\"")
+            (string-contains (caddr r) "event: error")
+            (string-contains (caddr r) "\"type\":\"error\"")
+            (string-contains (caddr r) "api_error")
+            (string-contains (caddr r) "upstream idle timeout")
+            (not (string-contains (caddr r) "[DONE]"))
+            (>= (cadr r) 4000)
+            (<= (cadr r) 9000)))))
+
+(define %still-healthy-body
   '(begin
      (use-modules (ice-9 rdelim))
      (catch #t
             (lambda ()
-              (let* ((json-body (string-append
-                                 "{\"model\":\"claude-sonnet-4-6\","
-                                 "\"max_tokens\":1,"
-                                 "\"messages\":[{\"role\":\"user\","
-                                 "\"content\":\"hello\"}]}"))
-                     (clen (string-length json-body))
-                     (req (string-append "POST /v1/messages HTTP/1.0\r\n"
-                                         "Content-Type: application/json\r\n"
-                                         "x-api-key: sk-test\r\n"
-                                         "Content-Length: "
-                                         (number->string clen)
-                                         "\r\n\r\n"
-                                         json-body))
-                     (s (socket PF_INET SOCK_STREAM 0)))
+              (let* ((s (socket PF_INET SOCK_STREAM 0)))
                 (connect s
                          (make-socket-address AF_INET INADDR_LOOPBACK 16890))
-                (display req s)
+                (display "GET /health HTTP/1.0\r\n\r\n" s)
                 (force-output s)
                 (let ((rsp (read-string s)))
                   (close-port s)
                   (and (string? rsp)
-                       (string-contains rsp "502")
-                       (string-contains rsp "application/json")
-                       (string-contains rsp "upstream")))))
+                       (string-contains rsp "200")))))
+            (lambda (key . args)
+              #f))))
+
+(define %proxy-restart-test-body
+  ;; Stopping the service must free port 16890 (stop-service blocks
+  ;; until the destructor reaped the process), and a restart must
+  ;; serve again.  Poll in-VM: wait-for-tcp-port is host-side only.
+  '(begin
+     (use-modules (gnu services herd)
+                  (ice-9 rdelim))
+     (define (health)
+       (catch #t
+              (lambda ()
+                (let* ((s (socket PF_INET SOCK_STREAM 0)))
+                  (connect s
+                           (make-socket-address AF_INET INADDR_LOOPBACK 16890))
+                  (display "GET /health HTTP/1.0\r\n\r\n" s)
+                  (force-output s)
+                  (let ((rsp (read-string s)))
+                    (close-port s)
+                    (and (string? rsp)
+                         (string-contains rsp "200")))))
+              (lambda (key . args)
+                #f)))
+     (and (health)
+          (stop-service 'claude-proxy)
+          (not (health))
+          (start-service 'claude-proxy)
+          (let loop
+            ((n 0))
+            (cond
+              ((health)
+               #t)
+              ((>= n 20)
+               #f)
+              (else (sleep 1)
+                    (loop (+ n 1))))))))
+
+(define %auth-error-test-body
+  ;; The fake answers unknown routes with 401; the proxy must pass
+  ;; the status and body through instead of masking it.
+  `(begin
+     (use-modules (srfi srfi-13))
+     (let* ((fetch (eval ,%fetch-post-code
+                         (interaction-environment)))
+            (r (fetch "/other"
+                      "{\"stream\":false,\"model\":\"x\",\"max_tokens\":1}" 4)))
+       (and r
+            (list? r)
+            (string-contains (caddr r) "401")
+            (string-contains (caddr r) "authentication_error")
+            (< (cadr r) 2000)))))
+
+(define %connect-refused-test-body
+  ;; With the fake upstream stopped, a POST must fail fast with 502.
+  `(begin
+     (use-modules (gnu services herd)
+                  (srfi srfi-13))
+     (and (stop-service 'fake-upstream)
+          (let* ((fetch (eval ,%fetch-post-code
+                              (interaction-environment)))
+                 (r (fetch "/v1/messages"
+                           (string-append
+                            "{\"stream\":true,\"model\":\"deepseek-chat\","
+                            "\"max_tokens\":16,"
+                            "\"messages\":[{\"role\":\"user\","
+                            "\"content\":\"hello\"}]}") 4)))
+            (start-service 'fake-upstream)
+            (and r
+                 (list? r)
+                 (string-contains (caddr r) "502")
+                 (string-contains (caddr r) "upstream unavailable")
+                 (< (cadr r) 2000))))))
+
+(define %log-traceability-body
+  ;; Every proxy event carries a per-request req=<hex> id; at least
+  ;; one request must have accept, first-byte and done lines sharing
+  ;; that id, and the timeouts must be logged too — so one grep
+  ;; traces a whole request.
+  '(begin
+     (use-modules (ice-9 rdelim)
+                  (ice-9 regex))
+     (define (has-event? log id event)
+       (string-contains log
+                        (string-append id " event=" event)))
+     (define (traceable? log ids)
+       (and (pair? ids)
+            (or (and (has-event? log
+                                 (car ids) "accept")
+                     (has-event? log
+                                 (car ids) "first-byte")
+                     (has-event? log
+                                 (car ids) "done"))
+                (traceable? log
+                            (cdr ids)))))
+     (catch #t
+            (lambda ()
+              (let* ((log (call-with-input-file "/var/log/claude-proxy.log"
+                            read-string))
+                     (rx (make-regexp "req=[0-9a-f]+"))
+                     (ids (let loop
+                            ((m (regexp-exec rx log))
+                             (acc '()))
+                            (if m
+                                (loop (regexp-exec rx log
+                                                   (match:end m))
+                                      (cons (match:substring m) acc)) acc))))
+                (and (> (length ids) 2)
+                     (string-contains log "event=accept")
+                     (string-contains log "event=first-byte")
+                     (string-contains log "event=done")
+                     (string-contains log "event=timeout")
+                     (traceable? log ids))))
             (lambda (key . args)
               #f))))
 
+(define %tls-echo-test-body
+  ;; The TLS proxy (16892) reaches the TLS fake through the portless
+  ;; URL https://127.0.0.1.  A non-streaming request must come back
+  ;; with thinking injected — proving the upstream leg resolved the
+  ;; scheme-default port 443, trusted the test certificate, and
+  ;; completed a TLS request/response cycle.
+  `(begin
+     (use-modules (srfi srfi-13))
+     (let* ((fetch (eval ,%fetch-post-code
+                         (interaction-environment)))
+            (r (fetch "/v1/messages"
+                      (string-append
+                       "{\"stream\":false,\"model\":\"deepseek-chat\","
+                       "\"max_tokens\":16,"
+                       "\"messages\":[{\"role\":\"user\","
+                       "\"content\":\"hello\"}]}") 4 16892)))
+       (and r
+            (list? r)
+            (string-contains (caddr r) "200")
+            (string-contains (caddr r) "\"type\":\"disabled\"")))))
+
+(define %tls-stream-test-body
+  ;; Over TLS the fake behaves like the plain one: the first SSE
+  ;; frame is sent immediately, the rest only after a 3-second
+  ;; sleep.  An incremental relay delivers the first bytes in well
+  ;; under 2.5 s while the fake is still sleeping, and the full
+  ;; stream ends with [DONE].
+  `(begin
+     (use-modules (srfi srfi-13))
+     (let* ((fetch (eval ,%fetch-post-code
+                         (interaction-environment)))
+            (r (fetch "/v1/messages"
+                      (string-append
+                       "{\"stream\":true,\"model\":\"deepseek-chat\","
+                       "\"max_tokens\":16,"
+                       "\"messages\":[{\"role\":\"user\","
+                       "\"content\":\"hello\"}]}") 4 16892)))
+       (and r
+            (list? r)
+            (< (car r) 2500)
+            (>= (cadr r) 2500)
+            (string-contains (caddr r) "data: [DONE]")
+            (string-contains (caddr r) "message_start")))))
+
 (define (claude-proxy-test-cases marionette)
   "Return a gexp with claude-proxy test assertions."
   #~(begin
 
       #$(assert-service-running "claude-proxy: service running"
                                 'claude-proxy marionette)
+      #$(assert-service-running "fake-upstream: service running"
+                                'fake-upstream marionette)
 
       (test-assert "claude-proxy: port 16890 TCP"
                    (wait-for-tcp-port 16890
                                       #$marionette))
+      (test-assert "fake-upstream: port 19998 TCP"
+                   (wait-for-tcp-port 19998
+                                      #$marionette))
+
+      #$(assert-service-running "claude-proxy-tls: service running"
+                                'claude-proxy-tls marionette)
+      #$(assert-service-running "fake-upstream-tls: service running"
+                                'fake-upstream-tls marionette)
+      (test-assert "claude-proxy-tls: port 16892 TCP"
+                   (wait-for-tcp-port 16892
+                                      #$marionette))
+      (test-assert "fake-upstream-tls: port 443 TCP"
+                   (wait-for-tcp-port 443
+                                      #$marionette))
 
       (test-assert "claude-proxy: health endpoint returns 200"
                    (marionette-eval '#$%health-test-body
                                     #$marionette))
 
-      (test-assert "claude-proxy: POST forwarding returns 502"
-                   (marionette-eval '#$%post-502-test-body
-                                    #$marionette))))
+      (test-assert "claude-proxy: SSE first frame relays in < 2.5 s"
+                   (marionette-eval '#$%stream-test-body
+                                    #$marionette))
 
-(define test-proxy-wrapper
-  (program-file "test-claude-proxy-start"
-                #~(begin
-                    (setenv "PROXY_PORT" "16890")
-                    (setenv "CLAUDE_PROXY_UPSTREAM" "http://127.0.0.1:19998")
-                    (apply execl
-                           #$claude-proxy-script "claude-proxy"
-                           '()))))
+      (test-assert "claude-proxy: thinking disabled injected"
+                   (marionette-eval '#$%injection-test-body
+                                    #$marionette))
 
-(define test-proxy-start
-  #~(make-forkexec-constructor (list #$test-proxy-wrapper)
-                               #:log-file "/var/log/claude-proxy.log"))
+      (test-assert "claude-proxy: 504 when upstream stalls headers"
+                   (marionette-eval '#$%slow-headers-test-body
+                                    #$marionette))
+
+      (test-assert "claude-proxy: mid-stream stall emits typed SSE error"
+                   (marionette-eval '#$%mid-stall-test-body
+                                    #$marionette))
+
+      (test-assert "claude-proxy: healthy after mid-stream stall"
+                   (marionette-eval '#$%still-healthy-body
+                                    #$marionette))
+
+      (test-assert "claude-proxy: stop releases port; restart serves again"
+                   (marionette-eval '#$%proxy-restart-test-body
+                                    #$marionette))
+
+      (test-assert "claude-proxy: upstream 401 passed through"
+                   (marionette-eval '#$%auth-error-test-body
+                                    #$marionette))
+
+      (test-assert "claude-proxy: 502 when upstream is down"
+                   (marionette-eval '#$%connect-refused-test-body
+                                    #$marionette))
+
+      (test-assert "claude-proxy: one grep traces a request"
+                   (marionette-eval '#$%log-traceability-body
+                                    #$marionette))
+
+      (test-assert "claude-proxy-tls: echo via portless https://127.0.0.1"
+                   (marionette-eval '#$%tls-echo-test-body
+                                    #$marionette))
+
+      (test-assert "claude-proxy-tls: SSE over TLS relays to [DONE]"
+                   (marionette-eval '#$%tls-stream-test-body
+                                    #$marionette))))
 
 (define (claude-proxy-test-os-services base-os)
-  "Return a test OS service list that includes the proxy
-pointed at a non-existent upstream (19998), so POST returns 502."
-  (let ((proxy-svc (shepherd-service (provision '(claude-proxy))
-                                     (start test-proxy-start)
-                                     (stop #~(make-kill-destructor))
-                                     (auto-start? #t))))
-    (cons (simple-service 'claude-proxy-test shepherd-root-service-type
-                          (list proxy-svc))
-          (operating-system-user-services base-os))))
+  "Return a test OS service list with the fake upstreams (plain and
+TLS) and the proxy instances pointed at them."
+  (cons* (simple-service 'claude-proxy-test shepherd-root-service-type
+                         (list fake-upstream-service claude-proxy-test-service
+                               fake-upstream-tls-service
+                               claude-proxy-tls-test-service))
+         (operating-system-user-services base-os)))
index 49f675b9a2897c1b69e74677b62100e435cae5e0..0cbf709978ef762505fea5ccfbf3167d887c0cb7 100644 (file)
@@ -91,24 +91,10 @@ ed25519 key to the production SSH config."
       (openssh-configuration (inherit config)
                              (authorized-keys `(("dak" ,key))))))
   (cons* %static-networking
-         (let* ((wrapper (program-file "test-claude-proxy-start"
-                                       #~(begin
-                                           (setenv "PROXY_PORT" "16890")
-                                           (setenv "CLAUDE_PROXY_UPSTREAM"
-                                                   "http://127.0.0.1:19998")
-                                           (apply execl
-                                                  #$claude-proxy-script
-                                                  "claude-proxy"
-                                                  '()))))
-                (start #~(make-forkexec-constructor (list #$wrapper)
-                          #:log-file "/var/log/claude-proxy.log"))
-                (svc (shepherd-service (provision '(claude-proxy))
-                                       (start start)
-                                       (stop #~(make-kill-destructor))
-                                       (auto-start? #t))))
-           (simple-service 'claude-proxy-test shepherd-root-service-type
-                           (list svc)))
-         (modify-services (operating-system-user-services base-os)
+         ;; The claude-proxy test pair (plain and TLS fakes plus the
+         ;; two proxy instances) is added by claude-proxy-test-os-services,
+         ;; which both the main checkout and this worktree export.
+         (modify-services (claude-proxy-test-os-services base-os)
            (delete dhcpcd-service-type)
            (openssh-service-type config =>
                                  (add-test-key config)))))