feat: http-post-stream for SSE, pure-Scheme socks5 rewrite
ober
71a5ff38dcf183d9efb8ff8ba3e72cfa18ba7f8e
--- a/lib/std/net/request.sls +++ b/lib/std/net/request.sls @@ -8,6 +8,7 @@ (library (std net request) (export http-get http-post http-put http-delete http-head + http-post-stream request-status request-text request-content request-headers request-header request-close parse-url url-parts-scheme url-parts-host url-parts-port url-parts-path @@ -286,7 +287,7 @@ (let* ([max-body (*http-max-body-size*)] [cl (assoc "content-length" headers)] [chunked? (let ([te (assoc "transfer-encoding" headers)]) - (and te (string-contains (cdr te) "chunked")))]) + (and te (string-contains (string-downcase (cdr te)) "chunked")))]) (cond ;; Chunked transfer encoding [chunked? @@ -404,4 +405,184 @@ (if (null? rest) acc (loop (cdr rest) (string-append acc sep (car rest)))))])) + ;; ========== Streaming POST ========== + ;; + ;; http-post-stream calls event-cb with each complete SSE event string, + ;; e.g. "data: {...}" (without the trailing \n\n). + ;; Called with #f when the stream ends. + ;; Returns the HTTP status code. + + (define (http-post-stream url headers body event-cb) + (let* ([parsed (parse-url url)] + [scheme (url-parts-scheme parsed)] + [host (url-parts-host parsed)] + [port (url-parts-port parsed)] + [path (url-parts-path parsed)]) + (if (string=? scheme "https") + (stream-post-https host port path headers body event-cb) + (stream-post-http host port path headers body event-cb)))) + + ;; ─── HTTPS streaming ─────────────────────────────────────────── + + (define (stream-post-https host port path headers data event-cb) + (let ([handle #f]) + (dynamic-wind + (lambda () (void)) + (lambda () + (set! handle (rustls-connect host port)) + (send-https-request handle "POST" host path headers data) + (let* ([hdr+lo (stream-read-headers handle)] + [hdr-str (car hdr+lo)] + [leftover (cdr hdr+lo)] + [hdr-port (open-input-string hdr-str)] + [status (parse-status-code (read-line-crlf hdr-port))] + [hdrs (read-headers hdr-port)] + [chunked? (let ([te (assoc "transfer-encoding" hdrs)]) + (and te (string-contains (string-downcase (cdr te)) "chunked")))]) + (stream-sse handle leftover chunked? event-cb) + (event-cb #f) + status)) + (lambda () + (when handle + (guard (e [#t (void)]) (rustls-close handle))))))) + + (define (stream-read-headers handle) + ;; Read bytes until \r\n\r\n, return (cons header-string leftover-string) + (let ([buf (make-bytevector 4096)] + [acc ""]) + (let loop () + (let ([n (rustls-read handle buf 4096)]) + (if (<= n 0) + (cons acc "") + (let* ([chunk (utf8->string (bv-slice buf 0 n))] + [combined (string-append acc chunk)] + [sep (string-contains combined "\r\n\r\n")]) + (if sep + (cons (substring combined 0 (+ sep 4)) + (substring combined (+ sep 4) (string-length combined))) + (begin (set! acc combined) (loop))))))))) + + (define (stream-sse handle leftover chunked? event-cb) + ;; Dispatch to chunked or plain SSE streaming + (if chunked? + (stream-sse-chunked handle leftover event-cb) + (stream-sse-plain handle leftover event-cb))) + + (define (stream-sse-plain handle leftover event-cb) + ;; Stream plain (non-chunked) SSE body + (let ([buf (make-bytevector 4096)] + [acc leftover]) + (let loop () + ;; Flush complete events already in acc + (let event-loop () + (let ([sep (string-contains acc "\n\n")]) + (when sep + (let ([event (string-trim (substring acc 0 sep))]) + (when (> (string-length event) 0) + (event-cb event))) + (set! acc (substring acc (+ sep 2) (string-length acc))) + (event-loop)))) + ;; Read more bytes + (let ([n (rustls-read handle buf 4096)]) + (when (> n 0) + (set! acc (string-append acc (utf8->string (bv-slice buf 0 n)))) + (loop)))))) + + (define (stream-sse-chunked handle leftover event-cb) + ;; Stream chunked-TE SSE body: decode chunks then extract SSE events + (let ([buf (make-bytevector 4096)] + [raw-acc leftover] ; raw bytes including chunk headers + [sse-acc ""] ; decoded data bytes + [done? #f]) + (define (decode-chunks!) + ;; Decode as many complete chunks as possible from raw-acc + (let loop () + (when (not done?) + (let ([crlf (string-contains raw-acc "\r\n")]) + (when crlf + (let* ([size-str (substring raw-acc 0 crlf)] + [semi (string-contains size-str ";")] + [hex (string-trim (if semi + (substring size-str 0 semi) + size-str))] + [chunk-len (string->number hex 16)]) + (when chunk-len + (if (= chunk-len 0) + (set! done? #t) + (let ([data-start (+ crlf 2)] + [data-end (+ crlf 2 chunk-len)]) + ;; Need data + trailing \r\n in raw-acc + (when (>= (string-length raw-acc) (+ data-end 2)) + (set! sse-acc + (string-append sse-acc + (substring raw-acc data-start data-end))) + (set! raw-acc + (substring raw-acc (+ data-end 2) + (string-length raw-acc))) + (loop))))))))))) + (define (flush-sse-events!) + (let loop () + (let ([sep (string-contains sse-acc "\n\n")]) + (when sep + (let ([event (string-trim (substring sse-acc 0 sep))]) + (when (> (string-length event) 0) + (event-cb event))) + (set! sse-acc (substring sse-acc (+ sep 2) (string-length sse-acc))) + (loop))))) + (let outer-loop () + (decode-chunks!) + (flush-sse-events!) + (when (not done?) + (let ([n (rustls-read handle buf 4096)]) + (when (> n 0) + (set! raw-acc + (string-append raw-acc (utf8->string (bv-slice buf 0 n)))) + (outer-loop))))))) + + ;; ─── HTTP (plain) streaming ───────────────────────────────────── + + (define (stream-post-http host port path headers data event-cb) + (let-values ([(in out) (tcp-connect host port)]) + (dynamic-wind + (lambda () (void)) + (lambda () + (send-http-request out "POST" host path headers data) + (let* ([status-code (parse-status-code (read-line-crlf in))] + [resp-headers (read-headers in)] + [chunked? (let ([te (assoc "transfer-encoding" resp-headers)]) + (and te (string-contains (string-downcase (cdr te)) "chunked")))]) + ;; For plain TCP, read body as string then process as SSE + (let ([body (if chunked? + (read-chunked-body in (*http-max-body-size*)) + (let ([out-str (open-output-string)]) + (let loop () + (let ([c (read-char in)]) + (unless (eof-object? c) + (write-char c out-str) + (loop)))) + (get-output-string out-str)))]) + ;; Process SSE events from complete body + (let ([acc body]) + (let event-loop () + (let ([sep (string-contains acc "\n\n")]) + (when sep + (let ([event (string-trim (substring acc 0 sep))]) + (when (> (string-length event) 0) + (event-cb event))) + (set! acc (substring acc (+ sep 2) (string-length acc))) + (event-loop)))))) + (event-cb #f) + status-code)) + (lambda () + (close-port in) + (close-port out))))) + + ;; ─── Helpers ──────────────────────────────────────────────────── + + (define (bv-slice bv start end) + (let ([n (- end start)] + [out (make-bytevector (- end start))]) + (bytevector-copy! bv start out 0 n) + out)) + ) ;; end library --- a/lib/std/net/socks5-server.sls +++ b/lib/std/net/socks5-server.sls @@ -1,162 +1,283 @@ #!chezscheme -;;; (std net socks5-server) — SOCKS5 proxy server via Rust FFI +;;; :std/net/socks5-server -- SOCKS5 proxy server ;;; -;;; Starts a local SOCKS5 proxy server (RFC 1928) backed by the Rust -;;; implementation in libjerboa_native. Supports no-auth and -;;; username/password authentication (RFC 1929). +;;; Implements RFC 1928 SOCKS5 proxy with optional username/password auth (RFC 1929). +;;; Runs a listener thread that accepts connections and spawns relay threads. ;;; -;;; Usage: -;;; (import (std net socks5-server)) +;;; Uses ffi-shim functions exclusively — no raw libc foreign-procedures, no +;;; hardcoded platform constants. Works on Linux, FreeBSD, macOS, Android. ;;; -;;; ;; Start a proxy on a random port, no auth -;;; (define proxy (socks5-start)) -;;; (socks5-port proxy) ;; → actual port number -;;; -;;; ;; Start on specific port with auth -;;; (define proxy (socks5-start 1080 "user" "pass")) -;;; -;;; ;; Set *_PROXY env vars for child processes -;;; (socks5-set-proxy-env! proxy) -;;; -;;; ;; Unset proxy env vars -;;; (socks5-unset-proxy-env!) -;;; -;;; ;; Get stats -;;; (socks5-stats proxy) ;; → "active:0 total:5" -;;; -;;; ;; Stop -;;; (socks5-stop proxy) +;;; API: +;;; (socks5-start port) → handle (no auth, random port if 0) +;;; (socks5-start port user pass bind-addr) → handle (with auth) +;;; (socks5-stop handle) → void +;;; (socks5-port handle) → integer +;;; (socks5-stats handle) → string +;;; (socks5-set-proxy-env! handle) → void (sets ALL_PROXY etc.) +;;; (socks5-unset-proxy-env!) → void (library (std net socks5-server) - (export - socks5-start - socks5-stop - socks5-port - socks5-stats - socks5-set-proxy-env! - socks5-unset-proxy-env!) + (export socks5-start socks5-stop socks5-port socks5-stats + socks5-set-proxy-env! socks5-unset-proxy-env!) (import (chezscheme)) - ;; Load native library (dynamic builds). - ;; In static builds, symbols are pre-registered via Sforeign_symbol. - (define _native-loaded - (or (guard (e [#t #f]) (load-shared-object "libjerboa_native.so") #t) - (guard (e [#t #f]) (load-shared-object "lib/libjerboa_native.so") #t) - (guard (e [#t #f]) (load-shared-object "libjerboa_native.dylib") #t) - (guard (e [#t #f]) (load-shared-object "lib/libjerboa_native.dylib") #t) - #t)) - - ;; ========== FFI declarations ========== - - (define c-socks5-start - (foreign-procedure "jerboa_socks5_server_start" - (u8* size_t unsigned-16 u8* size_t u8* size_t) unsigned-64)) - - (define c-socks5-stop - (foreign-procedure "jerboa_socks5_server_stop" - (unsigned-64) int)) - - (define c-socks5-port - (foreign-procedure "jerboa_socks5_server_port" - (unsigned-64) unsigned-16)) - - (define c-socks5-stats - (foreign-procedure "jerboa_socks5_server_stats" - (unsigned-64 u8* size_t) int)) - - (define c-last-error - (foreign-procedure "jerboa_last_error" (u8* size_t) size_t)) - - ;; ========== Error helper ========== - - (define (subbytevector bv start end) - (let* ([n (- end start)] - [result (make-bytevector n)]) - (bytevector-copy! bv start result 0 n) - result)) - - (define (get-last-error) - (let ([buf (make-bytevector 512)]) - (let ([len (c-last-error buf 512)]) - (if (> len 0) - (utf8->string (subbytevector buf 0 (min len 511))) - "unknown error")))) - - ;; ========== Public API ========== - - ;; (socks5-start) → handle ; random port, no auth - ;; (socks5-start port) → handle ; specific port, no auth - ;; (socks5-start port user pw) → handle ; specific port, user/pass auth - ;; (socks5-start port user pw bind-addr) → handle + ;; ========== FFI (all through ffi-shim — portable) ========== + + (define c-listen-tcp-addr + (foreign-procedure "ffi_stream_listen_tcp_addr" (string int int) int)) + (define c-connect-tcp + (foreign-procedure "ffi_stream_connect_tcp" (string int) int)) + (define c-accept + (foreign-procedure "ffi_stream_accept" (int) int)) + (define c-accept-nb + (foreign-procedure "ffi_stream_accept_nonblock" (int) int)) + (define c-set-nonblock + (foreign-procedure "ffi_set_nonblock" (int) int)) + (define c-clear-nonblock + (foreign-procedure "ffi_clear_nonblock" (int) int)) + (define c-close + (foreign-procedure "ffi_stream_close" (int) void)) + (define c-poll1 + (foreign-procedure "ffi_poll_readable_ms" (int int) int)) + (define c-poll2 + (foreign-procedure "ffi_poll2" (int int int) int)) + (define c-getsockname-port + (foreign-procedure "ffi_getsockname_port" (int) int)) + (define c-bv-read-exact + (foreign-procedure "ffi_bv_read_exact" (int u8* int int) int)) + (define c-bv-write-all + (foreign-procedure "ffi_bv_write_all" (int u8* int int) int)) + (define c-bv-recv + (foreign-procedure "ffi_bv_recv" (int u8* int int) int)) + (define c-bv-write + (foreign-procedure "ffi_bv_write" (int u8* int int) int)) + + ;; ========== Utilities ========== + + (define (read-bytes fd n) + ;; Read exactly n bytes, returning bytevector or #f + (let ([buf (make-bytevector n 0)]) + (if (>= (c-bv-read-exact fd buf 0 n) 0) + buf + #f))) + + (define (write-bytes fd bv) + ;; Write all bytes, return #t on success + (>= (c-bv-write-all fd bv 0 (bytevector-length bv)) 0)) + + (define (make-bv . bytes) + (u8-list->bytevector bytes)) + + ;; ========== SOCKS5 Protocol ========== + + (define (bytevector-contains bv val) + (let loop ([i 0]) + (cond + [(= i (bytevector-length bv)) #f] + [(= (bytevector-u8-ref bv i) val) #t] + [else (loop (+ i 1))]))) + + (define (handle-greeting fd auth?) + (let ([hdr (read-bytes fd 2)]) + (and hdr + (= (bytevector-u8-ref hdr 0) 5) ;; SOCKS5 + (let* ([nmethods (bytevector-u8-ref hdr 1)] + [methods (read-bytes fd nmethods)]) + (and methods + (if auth? + (if (bytevector-contains methods 2) + (write-bytes fd (make-bv 5 2)) + (begin (write-bytes fd (make-bv 5 #xff)) #f)) + (if (bytevector-contains methods 0) + (write-bytes fd (make-bv 5 0)) + (begin (write-bytes fd (make-bv 5 #xff)) #f)))))))) + + (define (handle-auth fd username password) + ;; RFC 1929: VER=1 ULEN USER PLEN PASS → VER STATUS (0=success) + (let ([ver (read-bytes fd 1)]) + (and ver (= (bytevector-u8-ref ver 0) 1) + (let ([ulen-bv (read-bytes fd 1)]) + (and ulen-bv + (let* ([ulen (bytevector-u8-ref ulen-bv 0)] + [user (read-bytes fd ulen)]) + (and user + (let ([plen-bv (read-bytes fd 1)]) + (and plen-bv + (let* ([plen (bytevector-u8-ref plen-bv 0)] + [pass (read-bytes fd plen)]) + (and pass + (let ([u (utf8->string user)] + [p (utf8->string pass)]) + (if (and (string=? u username) + (string=? p password)) + (write-bytes fd (u8-list->bytevector '(1 0))) + (begin + (write-bytes fd (u8-list->bytevector '(1 1))) + #f)))))))))))))) + + (define (handle-request fd) + ;; Parse CONNECT request, return (values host port) or (values #f #f) + (let ([hdr (read-bytes fd 4)]) + (if (not hdr) + (values #f #f) + (let ([ver (bytevector-u8-ref hdr 0)] + [cmd (bytevector-u8-ref hdr 1)] + [atyp (bytevector-u8-ref hdr 3)]) + (if (not (and (= ver 5) (= cmd 1))) ;; Only CONNECT + (begin + (write-bytes fd (u8-list->bytevector '(5 7 0 1 0 0 0 0 0 0))) + (values #f #f)) + (case atyp + [(1) ;; IPv4 + (let ([addr-bv (read-bytes fd 4)] + [port-bv (read-bytes fd 2)]) + (if (or (not addr-bv) (not port-bv)) + (values #f #f) + (let ([a (bytevector-u8-ref addr-bv 0)] + [b (bytevector-u8-ref addr-bv 1)] + [c (bytevector-u8-ref addr-bv 2)] + [d (bytevector-u8-ref addr-bv 3)] + [p (bytevector-u16-ref port-bv 0 (endianness big))]) + (values (format "~a.~a.~a.~a" a b c d) p))))] + [(3) ;; Domain name + (let ([len-bv (read-bytes fd 1)]) + (and len-bv + (let* ([len (bytevector-u8-ref len-bv 0)] + [name-bv (read-bytes fd len)] + [port-bv (read-bytes fd 2)]) + (if (or (not name-bv) (not port-bv)) + (values #f #f) + (values (utf8->string name-bv) + (bytevector-u16-ref port-bv 0 (endianness big)))))))] + [else + (write-bytes fd (u8-list->bytevector '(5 8 0 1 0 0 0 0 0 0))) + (values #f #f)])))))) + + (define (send-reply fd rep) + (write-bytes fd (u8-list->bytevector (list 5 rep 0 1 0 0 0 0 0 0)))) + + ;; ========== Relay ========== + + (define (relay-data client-fd target-fd) + ;; Bidirectional relay using ffi_poll2 + ffi_bv_recv/ffi_bv_write + (let ([buf (make-bytevector 8192 0)]) + (let loop () + (let ([mask (c-poll2 client-fd target-fd 60000)]) + (when (> mask 0) + (let ([ok #t]) + ;; client → target (bit 0) + (when (> (fxlogand mask 1) 0) + (let ([r (c-bv-recv client-fd buf 0 8192)]) + (if (<= r 0) + (set! ok #f) + (when (< (c-bv-write-all target-fd buf 0 r) 0) + (set! ok #f))))) + ;; target → client (bit 1) + (when (and ok (> (fxlogand mask 2) 0)) + (let ([r (c-bv-recv target-fd buf 0 8192)]) + (if (<= r 0) + (set! ok #f) + (when (< (c-bv-write-all client-fd buf 0 r) 0) + (set! ok #f))))) + ;; HUP/ERR bits: 4 = fd1, 8 = fd2 + (when (and ok + (= (fxlogand mask 4) 0) + (= (fxlogand mask 8) 0)) + (loop)))))))) + + ;; ========== Handle a single client ========== + + (define (handle-client client-fd auth-info stats-box) + (guard (e [#t (c-close client-fd)]) + (let ([auth? (and auth-info #t)]) + (when (handle-greeting client-fd auth?) + (when (or (not auth?) + (handle-auth client-fd (car auth-info) (cdr auth-info))) + (let-values ([(host port) (handle-request client-fd)]) + (when (and host port) + (let ([target-fd (c-connect-tcp host port)]) + (if (< target-fd 0) + (send-reply client-fd 5) ;; connection refused + (begin + (send-reply client-fd 0) ;; success + (let ([cur (unbox stats-box)]) + (set-box! stats-box + (cons (+ (car cur) 1) (cdr cur)))) + (relay-data client-fd target-fd) + (c-close target-fd))))))) + (c-close client-fd))))) + + ;; ========== Server ========== + + (define (make-handle fd stop stats port thd) + (vector fd stop stats port thd)) + (define (handle-fd h) (vector-ref h 0)) + (define (handle-stop h) (vector-ref h 1)) + (define (handle-stats h) (vector-ref h 2)) + (define (handle-port* h) (vector-ref h 3)) + (define (handle-thd h) (vector-ref h 4)) + + (define (server-loop listener-fd stop-box auth-info stats-box) + (c-set-nonblock listener-fd) + (let loop () + (unless (unbox stop-box) + (let ([ready (c-poll1 listener-fd 500)]) + (when (> ready 0) + (let ([client-fd (c-accept listener-fd)]) + (when (>= client-fd 0) + (fork-thread + (lambda () + (handle-client client-fd auth-info stats-box))))))) + (loop)))) + (define socks5-start (case-lambda - [() - (socks5-start* "127.0.0.1" 0 #f #f)] [(port) - (socks5-start* "127.0.0.1" port #f #f)] - [(port username password) - (socks5-start* "127.0.0.1" port username password)] + (start-server port #f "127.0.0.1")] [(port username password bind-addr) - (socks5-start* bind-addr port username password)])) - - (define (socks5-start* bind-addr port username password) - (let* ([addr-bv (string->utf8 bind-addr)] - [user-bv (if username (string->utf8 username) (make-bytevector 0))] - [pass-bv (if password (string->utf8 password) (make-bytevector 0))] - [handle (c-socks5-start - addr-bv (bytevector-length addr-bv) - port - user-bv (if username (bytevector-length user-bv) 0) - pass-bv (if password (bytevector-length pass-bv) 0))]) - (when (= handle 0) - (error 'socks5-start (get-last-error))) - handle)) - - ;; Stop a running SOCKS5 proxy server. + (start-server port (cons username password) bind-addr)])) + + (define (start-server port auth-info bind-addr) + (let ([fd (c-listen-tcp-addr bind-addr port 128)]) + (when (< fd 0) + (error 'socks5-start (format "failed to listen on ~a:~a (errno ~a)" + bind-addr port (- fd)))) + (let* ([actual-port (c-getsockname-port fd)] + [stop-box (box #f)] + [stats-box (box (cons 0 0))] + [thd (fork-thread + (lambda () + (server-loop fd stop-box auth-info stats-box)))]) + (make-handle fd stop-box stats-box actual-port thd)))) + (define (socks5-stop handle) - (let ([rc (c-socks5-stop handle)]) - (when (< rc 0) - (error 'socks5-stop (get-last-error))))) + (set-box! (handle-stop handle) #t) + (sleep (make-time 'time-duration 0 1)) + (c-close (handle-fd handle))) - ;; Get the actual bound port. (define (socks5-port handle) - (let ([p (c-socks5-port handle)]) - (when (= p 0) - (error 'socks5-port (get-last-error))) - p)) + (handle-port* handle)) - ;; Get stats string: "active:N total:N" (define (socks5-stats handle) - (let* ([buf (make-bytevector 256)] - [n (c-socks5-stats handle buf 256)]) - (if (> n 0) - (utf8->string (subbytevector buf 0 n)) - ""))) - - ;; ========== Proxy environment variables ========== - - ;; All the environment variable names that programs check for proxy config. - (define *proxy-env-vars* - '("HTTP_PROXY" "http_proxy" - "HTTPS_PROXY" "https_proxy" - "ALL_PROXY" "all_proxy" - "SOCKS_PROXY" "socks_proxy" - "SOCKS5_PROXY" "socks5_proxy")) - - ;; Set all proxy env vars to point to the running SOCKS5 proxy. - ;; Uses socks5h:// scheme (h = proxy resolves DNS, standard for SOCKS5). + (let ([s (unbox (handle-stats handle))]) + (format "connections: ~a" (car s)))) + (define (socks5-set-proxy-env! handle) - (let* ([port (socks5-port handle)] - [url (string-append "socks5h://127.0.0.1:" (number->string port))]) - (for-each - (lambda (var) (putenv var url)) - *proxy-env-vars*))) + (let ([port (socks5-port handle)]) + (putenv "ALL_PROXY" (format "socks5h://127.0.0.1:~a" port)) + (putenv "all_proxy" (format "socks5h://127.0.0.1:~a" port)) + (putenv "http_proxy" (format "socks5h://127.0.0.1:~a" port)) + (putenv "https_proxy" (format "socks5h://127.0.0.1:~a" port)) + (putenv "HTTP_PROXY" (format "socks5h://127.0.0.1:~a" port)) + (putenv "HTTPS_PROXY" (format "socks5h://127.0.0.1:~a" port)))) - ;; Unset all proxy env vars. (define (socks5-unset-proxy-env!) - (for-each - (lambda (var) (putenv var "")) - *proxy-env-vars*)) + (putenv "ALL_PROXY" "") + (putenv "all_proxy" "") + (putenv "http_proxy" "") + (putenv "https_proxy" "") + (putenv "HTTP_PROXY" "") + (putenv "HTTPS_PROXY" "")) ) ;; end library