// Package proxy is the mitmux proxy engine: a forward HTTP proxy. // Plain HTTP requests pass through unmodified. CONNECT requests (HTTPS) // are intercepted: mitmux terminates TLS with the client using a leaf // certificate signed by its own CA, and separately terminates TLS with // the real server, forwarding requests between the two. ALPN is // negotiated independently on each side (see handleConnect) so HTTP/2 // stays HTTP/2 end to end without one side being forced to match the // other. Every request/response pair is captured to the history store - // exactly, byte for byte, on HTTP/1.1 legs; reconstructed on HTTP/2 legs, // which have no meaningful "raw bytes" of their own (see capture.go). // // Upstream requests are round-tripped manually (write the request, // read the response off the same connection) rather than through // http.Transport: Transport's automatic HTTP/2 dispatch keys off a // literal *tls.Conn type assertion on the connection it dials, which a // capturing wrapper around that connection defeats - the request would // silently be parsed as HTTP/1.1 over what is actually HTTP/2 framing. // Handling both protocols explicitly here, per request, avoids that and // also removes any ambiguity about which connection served which // request, since each request gets its own connection either way. package proxy import ( "bufio" "bytes" "context" "crypto/tls" "errors" "fmt" "io" "log" "net" "net/http" "net/url" "strconv" "strings" "sync" "time" "golang.org/x/net/http2" xproxy "golang.org/x/net/proxy" "mitmux/internal/ca" "mitmux/internal/clientcert" "mitmux/internal/rules" "mitmux/internal/scope" "mitmux/internal/store" ) // upstreamTimeout bounds the write-request/read-response phase of an // upstream exchange, once dialing has already succeeded. const upstreamTimeout = 60 * time.Second // clientHeaderTimeout bounds how long a client connection can sit // sending request headers (or a TLS ClientHello, on the CONNECT-tunnel // leg) before mitmux gives up on it - a slow-loris style connection // that opens and then trickles bytes (or never sends a ClientHello at // all) would otherwise hold a connection and its goroutine open // indefinitely, with nothing else in the codebase bounding it. Narrow // on purpose: this only covers the pre-body phase, not overall request // duration - a legitimately slow multi-minute upload/download must // still work, so this is deliberately not a blanket ReadTimeout/ // WriteTimeout on the whole connection. const clientHeaderTimeout = 30 * time.Second // clientIdleTimeout bounds how long a keep-alive client connection can // sit idle between requests before mitmux closes it - cleans up // abandoned idle connections without affecting any connection that's // actively mid-transfer. const clientIdleTimeout = 120 * time.Second // hopByHopHeaders are stripped before forwarding a request or response, // per RFC 7230 6.1 - they are meaningful only between a client and its // immediate next hop, not end-to-end. var hopByHopHeaders = []string{ "Connection", "Proxy-Connection", "Keep-Alive", "Proxy-Authenticate", "Proxy-Authorization", "TE", "Trailers", "Transfer-Encoding", "Upgrade", } // Server is a forward proxy listener. It can bind more than one address // at once (Addrs) - all sharing the same handler, history store, CA and // rules, so a client on any of them sees identical behavior; this is for // cases like wanting a separate port per client/network segment, not // for running logically different proxies in one process. type Server struct { Addrs []string // UpstreamProxy, if set, chains every outbound connection through // another proxy instead of dialing origins directly - e.g. routing // mitmux's own traffic through Burp, a corporate proxy, a network- // access proxy, or Tor. Bare "host:port" (or an "http://" prefix, // stripped before it gets here) means an HTTP CONNECT proxy; a // "socks5://[user:pass@]host:port" prefix means SOCKS5 - see // parseSOCKS5. UpstreamProxy string // OnEntry, if set, is called after each request/response pair is // stored, so a daemon can broadcast it to live TUI subscribers. OnEntry func(store.Summary) ca *ca.CA store *store.Store server *http.Server } // New creates a proxy Server bound to addrs (e.g. ["127.0.0.1:8080"]), // signing intercepted TLS connections with root and recording history to // db. upstreamProxy chains outbound connections through another HTTP // CONNECT proxy (host:port, no scheme) instead of dialing origins // directly; empty disables chaining. func New(addrs []string, root *ca.CA, db *store.Store, upstreamProxy string) *Server { s := &Server{Addrs: addrs, ca: root, store: db, UpstreamProxy: upstreamProxy} s.server = &http.Server{ Handler: http.HandlerFunc(s.handle), ConnContext: withClientTee, ReadHeaderTimeout: clientHeaderTimeout, IdleTimeout: clientIdleTimeout, } return s } // ListenAndServe binds every address in Addrs and blocks until one of // them stops (including on Shutdown, which closes all of them - every // Serve call below then returns http.ErrServerClosed). Addresses are all // bound up front before any of them start serving, so a bad address // (already in use, unparseable, ...) fails startup immediately rather // than leaving the daemon partially listening. func (s *Server) ListenAndServe() error { if len(s.Addrs) == 0 { return errors.New("no listen addresses configured") } lns := make([]net.Listener, len(s.Addrs)) for i, addr := range s.Addrs { ln, err := net.Listen("tcp", addr) if err != nil { for _, opened := range lns[:i] { opened.Close() } return fmt.Errorf("listen on %s: %w", addr, err) } lns[i] = ln } errCh := make(chan error, len(lns)) for i, ln := range lns { log.Printf("proxy listening on %s", s.Addrs[i]) go func(ln net.Listener) { errCh <- s.server.Serve(&teeListener{Listener: ln}) }(ln) } return <-errCh } // Shutdown gracefully stops the proxy. func (s *Server) Shutdown(ctx context.Context) error { return s.server.Shutdown(ctx) } func (s *Server) handle(w http.ResponseWriter, r *http.Request) { if r.Method == http.MethodConnect { s.handleConnect(w, r) return } s.handleHTTP(w, r) } // dialer resolves a fresh upstream connection for one request, along // with the ALPN protocol negotiated for it ("http/1.1", "h2", or "" if // not applicable/negotiated). type dialer func(ctx context.Context) (conn net.Conn, negotiated string, err error) // handleConnect intercepts a CONNECT request: it terminates TLS with the // client using a leaf certificate signed by mitmux's CA, then forwards // each request upstream over its own independently negotiated TLS // connection. Client-side and upstream-side ALPN are negotiated // separately (each offering both HTTP/2 and HTTP/1.1) rather than one // being forced to match the other, so e.g. an HTTP/1.1-only client // reaching an HTTP/2-preferring server doesn't fail to connect. func (s *Server) handleConnect(w http.ResponseWriter, r *http.Request) { hostPort := r.Host hostname, _, err := net.SplitHostPort(hostPort) if err != nil { hostname = hostPort hostPort = net.JoinHostPort(hostPort, "443") } hijacker, ok := w.(http.Hijacker) if !ok { http.Error(w, "hijacking not supported", http.StatusInternalServerError) return } client, _, err := hijacker.Hijack() if err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) return } if _, err := client.Write([]byte("HTTP/1.1 200 Connection Established\r\n\r\n")); err != nil { client.Close() return } clientTLS := tls.Server(client, &tls.Config{ GetCertificate: func(hello *tls.ClientHelloInfo) (*tls.Certificate, error) { name := hello.ServerName if name == "" { name = hostname } return s.ca.LeafFor(name) }, NextProtos: []string{http2.NextProtoTLS, "http/1.1"}, MinVersion: tls.VersionTLS12, }) // Bounded the same way the upstream leg already is (see forward's // conn.SetDeadline): without this, a client that completes CONNECT // and then never sends a ClientHello at all holds the connection and // its goroutine open indefinitely. Cleared after a successful // handshake - the request/response phase that follows has no // business inheriting a short handshake-only deadline. client.SetDeadline(time.Now().Add(clientHeaderTimeout)) if err := clientTLS.Handshake(); err != nil { log.Printf("mitm handshake with client for %s: %v", hostname, err) client.Close() return } client.SetDeadline(time.Time{}) dial := func(ctx context.Context) (net.Conn, string, error) { return dialUpstreamTLS(ctx, hostPort, hostname, s.UpstreamProxy, s.clientCertFor(hostname)) } handler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { s.forward(dial, "https", hostname, w, r) }) if clientTLS.ConnectionState().NegotiatedProtocol == http2.NextProtoTLS { (&http2.Server{}).ServeConn(clientTLS, &http2.ServeConnOpts{Handler: handler}) return } h1 := &http.Server{ Handler: handler, ConnContext: withClientTee, ReadHeaderTimeout: clientHeaderTimeout, IdleTimeout: clientIdleTimeout, } err = h1.Serve(newSingleConnListener(clientTLS)) if err != nil && !errors.Is(err, io.EOF) { log.Printf("h1 serve for %s: %v", hostname, err) } } // dialUpstreamTLS connects to the real server (directly, or tunneled // through upstreamProxy if set - see dialViaProxy), offering both // HTTP/2 and HTTP/1.1 over ALPN and letting the server pick. Chaining // through another proxy is transparent to everything from here on: once // the CONNECT tunnel is up, TLS and the request/response on top of it // look identical to a direct connection. cert, if non-nil, is presented // during the handshake for servers that require mutual TLS - see // clientCertFor. func dialUpstreamTLS(ctx context.Context, hostPort, sni, upstreamProxy string, cert *tls.Certificate) (net.Conn, string, error) { raw, err := dialViaProxy(ctx, hostPort, upstreamProxy) if err != nil { return nil, "", err } cfg := &tls.Config{ ServerName: sni, NextProtos: []string{http2.NextProtoTLS, "http/1.1"}, } if cert != nil { cfg.Certificates = []tls.Certificate{*cert} } conn := tls.Client(raw, cfg) if err := conn.HandshakeContext(ctx); err != nil { raw.Close() return nil, "", err } return conn, conn.ConnectionState().NegotiatedProtocol, nil } // clientCertFor returns the client certificate configured for host, if // any - see clientcert.FindFor. Errors (a bad DB read, an unparseable // PEM pair) are logged and treated as "no certificate configured" rather // than failing the connection outright: a broken client-cert config // shouldn't take down otherwise-working proxying for that host. func (s *Server) clientCertFor(host string) *tls.Certificate { if s.store == nil { return nil } certs, err := s.store.ListClientCerts() if err != nil { log.Printf("list client certs: %v", err) return nil } c := clientcert.FindFor(certs, host) if c == nil { return nil } tc, err := c.TLSCertificate() if err != nil { log.Printf("client cert %q: %v", c.Name, err) return nil } return &tc } // dialUpstreamPlain connects to a plain (non-TLS) upstream for the // non-CONNECT proxy path, which is always HTTP/1.1. Unlike the TLS/ // CONNECT path, chaining here means dialing the upstream proxy's own // address directly and writing it an absolute-form request (what a // proxy expects) rather than tunneling - see forward()'s viaProxyForm. func dialUpstreamPlain(ctx context.Context, host, upstreamProxy string) (net.Conn, string, error) { if _, _, err := net.SplitHostPort(host); err != nil { host = net.JoinHostPort(host, "80") } // SOCKS5 tunnels straight to host, same as dialViaProxy's TLS path - // see dialSOCKS5's doc comment for why that needs no absolute-form // adjustment the way chaining through an HTTP proxy does below. if addr, auth, err := parseSOCKS5(upstreamProxy); err != nil { return nil, "", err } else if addr != "" { conn, err := dialSOCKS5(ctx, addr, auth, host) return conn, "http/1.1", err } target := host if upstreamProxy != "" { target = upstreamProxy } nd := &net.Dialer{Timeout: 10 * time.Second} conn, err := nd.DialContext(ctx, "tcp", target) return conn, "http/1.1", err } // parseSOCKS5 returns the proxy's bare "host:port" and optional // credentials if upstreamProxy has a "socks5://" prefix - the marker // this tool uses to distinguish a SOCKS5 upstream from the default HTTP // CONNECT proxy chaining every other non-empty value means. addr == "" // means upstreamProxy isn't a SOCKS5 proxy, which includes the "no // upstream proxy configured at all" empty-string case - callers branch // on that the same way they'd branch on upstreamProxy == "". func parseSOCKS5(upstreamProxy string) (addr string, auth *xproxy.Auth, err error) { if !strings.HasPrefix(upstreamProxy, "socks5://") { return "", nil, nil } u, err := url.Parse(upstreamProxy) if err != nil { return "", nil, fmt.Errorf("invalid socks5 upstream proxy %q: %w", upstreamProxy, err) } if u.User != nil { pass, _ := u.User.Password() auth = &xproxy.Auth{User: u.User.Username(), Password: pass} } return u.Host, auth, nil } // dialSOCKS5 tunnels to target through the SOCKS5 proxy at proxyAddr. // Unlike an HTTP CONNECT proxy, SOCKS5 is transport-level and protocol- // agnostic: the resulting connection behaves exactly like one dialed // directly to target, with no "absolute-form request" adjustment needed // on top (see forward's proxyForm). func dialSOCKS5(ctx context.Context, proxyAddr string, auth *xproxy.Auth, target string) (net.Conn, error) { d, err := xproxy.SOCKS5("tcp", proxyAddr, auth, xproxy.Direct) if err != nil { return nil, fmt.Errorf("configure SOCKS5 proxy %s: %w", proxyAddr, err) } // xproxy.Direct (the forward dialer passed above) always yields a // ContextDialer-capable SOCKS5 client, per the library's own // implementation - this fallback exists so a future forward-dialer // change can't silently drop context cancellation rather than fail // to compile against a changed interface. cd, ok := d.(xproxy.ContextDialer) if !ok { return d.Dial("tcp", target) } conn, err := cd.DialContext(ctx, "tcp", target) if err != nil { return nil, fmt.Errorf("dial %s via SOCKS5 proxy %s: %w", target, proxyAddr, err) } return conn, nil } // dialViaProxy returns a raw TCP connection ready to speak TLS to // hostPort - dialed directly if upstreamProxy is empty, tunneled through // a SOCKS5 proxy if upstreamProxy has a "socks5://" prefix, or tunneled // through an HTTP CONNECT proxy otherwise. func dialViaProxy(ctx context.Context, hostPort, upstreamProxy string) (net.Conn, error) { if addr, auth, err := parseSOCKS5(upstreamProxy); err != nil { return nil, err } else if addr != "" { return dialSOCKS5(ctx, addr, auth, hostPort) } nd := &net.Dialer{Timeout: 10 * time.Second} if upstreamProxy == "" { return nd.DialContext(ctx, "tcp", hostPort) } conn, err := nd.DialContext(ctx, "tcp", upstreamProxy) if err != nil { return nil, fmt.Errorf("dial upstream proxy %s: %w", upstreamProxy, err) } connectReq := &http.Request{ Method: http.MethodConnect, URL: &url.URL{Opaque: hostPort}, Host: hostPort, Header: make(http.Header), } if err := connectReq.Write(conn); err != nil { conn.Close() return nil, fmt.Errorf("write CONNECT to upstream proxy %s: %w", upstreamProxy, err) } br := bufio.NewReader(conn) resp, err := http.ReadResponse(br, connectReq) if err != nil { conn.Close() return nil, fmt.Errorf("read CONNECT response from upstream proxy %s: %w", upstreamProxy, err) } if resp.StatusCode != http.StatusOK { conn.Close() return nil, fmt.Errorf("upstream proxy %s refused CONNECT to %s: %s", upstreamProxy, hostPort, resp.Status) } if br.Buffered() > 0 { // The upstream proxy shouldn't send anything past the CONNECT // response before the tunnel starts, but if it did, those bytes // are sitting in br's buffer, not on conn - replay them first // rather than silently dropping the start of the TLS handshake. return &prefixedConn{Conn: conn, r: br}, nil } return conn, nil } // prefixedConn serves buffered bytes from r before falling through to // reading directly off the underlying connection. type prefixedConn struct { net.Conn r *bufio.Reader } func (c *prefixedConn) Read(p []byte) (int, error) { if c.r.Buffered() > 0 { return c.r.Read(p) } return c.Conn.Read(p) } // roundTripH1 writes outReq directly to conn and reads the response back // off the same connection, wrapping conn in a teeConn so the exact wire // bytes of both can be captured. proxyForm selects an absolute-form // request line ("GET http://host/path HTTP/1.1") instead of origin-form // - needed when conn is a connection to another proxy, which expects // that form, rather than to the origin server itself. func roundTripH1(conn net.Conn, outReq *http.Request, proxyForm bool) (*http.Response, *teeConn, error) { tee := newTeeConn(conn) var writeErr error if proxyForm { writeErr = outReq.WriteProxy(tee) } else { writeErr = outReq.Write(tee) } if writeErr != nil { return nil, nil, writeErr } resp, err := http.ReadResponse(bufio.NewReader(tee), outReq) if err != nil { return nil, nil, err } return resp, tee, nil } // roundTripH2 sends outReq over a new single-connection HTTP/2 client. func roundTripH2(conn net.Conn, outReq *http.Request) (*http.Response, error) { cc, err := (&http2.Transport{}).NewClientConn(conn) if err != nil { return nil, err } return cc.RoundTrip(outReq) } // forward dials upstream, sends r, copies the response back to w, and // records the exchange to history. r's URL is rewritten from // origin-form (as read off the terminated connection) to absolute-form // for the round trip. func (s *Server) forward(dial dialer, scheme, hostname string, w http.ResponseWriter, r *http.Request) { clientTee := teeConnFromContext(r.Context()) outReq := r.Clone(r.Context()) outReq.URL.Scheme = scheme outReq.URL.Host = hostname outReq.RequestURI = "" if isWebSocketUpgradeRequest(r) { stripHopByHopKeepingUpgrade(outReq.Header) } else { stripHopByHop(outReq.Header) } // Header rules are applied to outReq only, after cloning and header // stripping - history's request_raw keeps showing what the client // actually sent (clientTee/reqBodyCap already capture from r, not // outReq), while what actually reaches the upstream server reflects // the rules. That split is deliberate: match-and-replace is a wire // transform, not a rewrite of the audit trail. reqRules, err := s.enabledRules("request") if err != nil { log.Printf("load request rules: %v", err) reqRules = nil } if len(reqRules) > 0 { outReq.Header = rules.ApplyHeaders(outReq.Header, reqRules) } // A body rule needs the body materialized in memory to rewrite it - // the exact opposite of the normal streamed-straight-through path, // which is what makes exact capture of an arbitrarily large body // possible without ever buffering it. Only paid when a body rule is // actually configured and enabled; everyone else keeps streaming. var reqBodyCap *cappedTee if rules.HasBodyRules(reqRules) && outReq.Body != nil { if clientTee == nil { // H2: nothing has captured this body's original bytes yet - // wrap it first so draining it below (to apply the rule) // captures them as a side effect, same as the unconditional // wrap further down does when no body rule is in play. reqBodyCap = newCappedTee(outReq.Body) outReq.Body = io.NopCloser(reqBodyCap) } // H1: clientTee already captures every byte read off the client // connection regardless of who's doing the reading, so draining // outReq.Body here (which for H1 is the same underlying reader // r.Body was, per Clone's documented behavior of not deep- // copying Body) is captured exactly as if roundTripH1 had read // it directly during the actual upstream write. newBody, newLen, applied, bodyErr := applyBodyRules(outReq.Body, reqRules) if bodyErr != nil { log.Printf("apply request body rules: %v", bodyErr) } else { outReq.Body = newBody if applied { outReq.ContentLength = newLen } } } // Only needed when the client leg isn't tee-captured (HTTP/2) and a // body rule hasn't already wrapped/captured it above. if clientTee == nil && outReq.Body != nil && reqBodyCap == nil { reqBodyCap = newCappedTee(outReq.Body) outReq.Body = io.NopCloser(reqBodyCap) } started := time.Now() conn, negotiated, dialErr := dial(r.Context()) if dialErr != nil { reqRaw, reqExact, reqTrunc := captureRequest(r, clientTee, reqBodyCap) s.record(started, time.Since(started), scheme, hostname, r, reqRaw, reqExact, reqTrunc, nil, false, false, 0, dialErr.Error()) http.Error(w, dialErr.Error(), http.StatusBadGateway) return } defer conn.Close() // The dial itself is bounded (net.Dialer.Timeout / HandshakeContext); // without this, a server that accepts the connection and then never // writes or never finishes writing would hang the request forever - // there's no other timeout covering the write-request/read-response // phase. Bounds the whole exchange, so a legitimately slow multi- // minute transfer would also get cut off; a fixed default is enough // for now, not worth a config surface yet. conn.SetDeadline(time.Now().Add(upstreamTimeout)) var resp *http.Response var upstreamTee *teeConn if negotiated == http2.NextProtoTLS { resp, err = roundTripH2(conn, outReq) } else { // Only the plain-HTTP path needs absolute-form, and only when // chained through an HTTP proxy specifically - a CONNECT tunnel // (chained or not) is transparent from here on, so it always uses // origin-form like a direct connection, and so does a SOCKS5 // upstream: SOCKS5 tunnels straight to the origin, invisible to // the HTTP layer, same as dialSOCKS5's doc comment explains. proxyForm := scheme == "http" && s.UpstreamProxy != "" && !strings.HasPrefix(s.UpstreamProxy, "socks5://") resp, upstreamTee, err = roundTripH1(conn, outReq, proxyForm) } duration := time.Since(started) reqRaw, reqExact, reqTrunc := captureRequest(r, clientTee, reqBodyCap) if err != nil { s.record(started, duration, scheme, hostname, r, reqRaw, reqExact, reqTrunc, nil, false, false, 0, err.Error()) http.Error(w, err.Error(), http.StatusBadGateway) return } defer resp.Body.Close() // A WebSocket upgrade stops being one-shot request/response the // instant the 101 lands - match-and-replace rules, body capture, and // the normal write-response-then-record flow below all assume a // bounded response with a body, none of which applies here. HTTP/2 // client legs are excluded: they can't be hijacked for raw access // the way an HTTP/1.1 connection can (see handleWebSocketUpgrade). if negotiated != http2.NextProtoTLS && isWebSocketUpgradeResponse(resp) { s.handleWebSocketUpgrade(w, r, scheme, hostname, started, duration, reqRaw, reqExact, reqTrunc, resp, upstreamTee) return } respRules, err := s.enabledRules("response") if err != nil { log.Printf("load response rules: %v", err) respRules = nil } var respBodyCap *cappedTee if upstreamTee == nil { respBodyCap = newCappedTee(resp.Body) resp.Body = io.NopCloser(respBodyCap) } // Same split as the request side: response_raw keeps reflecting what // the origin server actually sent (captured below, from upstreamTee // or respBodyCap, both already wired to resp.Body independent of // resp.Header), while the client actually receives the rule-modified // headers/body. if len(respRules) > 0 { resp.Header = rules.ApplyHeaders(resp.Header, respRules) } if rules.HasBodyRules(respRules) && resp.Body != nil { newBody, newLen, applied, bodyErr := applyBodyRules(resp.Body, respRules) if bodyErr != nil { log.Printf("apply response body rules: %v", bodyErr) } else { resp.Body = newBody if applied { // Unlike http.Request.Write (which derives the wire // Content-Length from req.ContentLength regardless of any // stale header), http.ResponseWriter does not: the loop // below just forwards whatever's in resp.Header verbatim. // A body rule that changes length would otherwise leave a // stale Content-Length on the wire and corrupt response // framing for the client. resp.ContentLength = newLen resp.Header.Set("Content-Length", strconv.FormatInt(newLen, 10)) } } } stripHopByHop(resp.Header) for k, vv := range resp.Header { for _, v := range vv { w.Header().Add(k, v) } } w.WriteHeader(resp.StatusCode) io.Copy(w, resp.Body) var respRaw []byte var respExact, respTrunc bool if upstreamTee != nil { respRaw, respTrunc = upstreamTee.Take() respExact = !respTrunc } else { respRaw, respExact = captureResponse(resp, respBodyCap) } s.record(started, duration, scheme, hostname, r, reqRaw, reqExact, reqTrunc, respRaw, respExact, respTrunc, resp.StatusCode, "") } // handleWebSocketUpgrade takes over the connection after resp (a 101 // matching isWebSocketUpgradeResponse) comes back from the origin. It // records the upgrade request/response pair to history exactly like a // normal exchange, then relays WebSocket frames bidirectionally, byte- // for-byte unmodified, until either side closes - decoding each frame's // payload along the way for capture into the ws_messages table, tagged // to this exchange's own history entry. func (s *Server) handleWebSocketUpgrade(w http.ResponseWriter, r *http.Request, scheme, hostname string, started time.Time, duration time.Duration, reqRaw []byte, reqExact, reqTrunc bool, resp *http.Response, upstreamTee *teeConn) { hijacker, ok := w.(http.Hijacker) if !ok { s.record(started, duration, scheme, hostname, r, reqRaw, reqExact, reqTrunc, nil, false, false, resp.StatusCode, "websocket: client connection doesn't support hijacking") return } clientConn, brw, err := hijacker.Hijack() if err != nil { s.record(started, duration, scheme, hostname, r, reqRaw, reqExact, reqTrunc, nil, false, false, resp.StatusCode, err.Error()) return } defer clientConn.Close() // A WebSocket connection is expected to live far longer than one // request/response - unlike the bounded upstreamTimeout the rest of // forward() uses, there's no natural cutoff here. clientConn.SetDeadline(time.Time{}) upstreamTee.SetDeadline(time.Time{}) // Everything upstreamTee has captured so far is exactly the raw 101 // response bytes, possibly with some already-arrived WebSocket frame // bytes tacked on the end (bufio's own read-ahead inside // roundTripH1) - split at the header/body boundary so the header // portion can be relayed and recorded as this exchange's // response_raw, and any leftover treated as the start of the frame // stream rather than lost. respRaw, _ := upstreamTee.Take() headerEnd := len(respRaw) if idx := bytes.Index(respRaw, []byte("\r\n\r\n")); idx >= 0 { headerEnd = idx + 4 } headerBytes, upstreamLeftover := respRaw[:headerEnd], respRaw[headerEnd:] if _, err := clientConn.Write(headerBytes); err != nil { s.record(started, duration, scheme, hostname, r, reqRaw, reqExact, reqTrunc, headerBytes, true, false, resp.StatusCode, err.Error()) return } entryID := s.record(started, duration, scheme, hostname, r, reqRaw, reqExact, reqTrunc, headerBytes, true, false, resp.StatusCode, "") var clientLeftover []byte if brw.Reader.Buffered() > 0 { clientLeftover = make([]byte, brw.Reader.Buffered()) io.ReadFull(brw.Reader, clientLeftover) } clientReader := io.MultiReader(bytes.NewReader(clientLeftover), clientConn) upstreamReader := io.MultiReader(bytes.NewReader(upstreamLeftover), upstreamTee) // Each direction pumps independently. done := make(chan struct{}, 2) go func() { pumpWS(clientReader, upstreamTee, func(opcode byte, payload []byte) { s.recordWSMessage(entryID, "client_to_server", opcode, payload) }) done <- struct{}{} }() go func() { pumpWS(upstreamReader, clientConn, func(opcode byte, payload []byte) { s.recordWSMessage(entryID, "server_to_client", opcode, payload) }) done <- struct{}{} }() // Wait for the first direction to stop, then give the other one a // bounded window to finish its own close sequence too - typically // relaying the peer's own close-frame reply - rather than tearing // the connection down the instant either side sees a close frame // pass through. Without this, a client that closes gracefully would // see its own close frame answered with an abrupt EOF instead of // the origin's actual close reply. If the other direction doesn't // finish in time (a slow or non-compliant peer), the deadlines below // force it to unblock rather than leak the goroutine indefinitely. <-done deadline := time.Now().Add(5 * time.Second) clientConn.SetDeadline(deadline) upstreamTee.SetDeadline(deadline) select { case <-done: case <-time.After(5 * time.Second): } } func (s *Server) recordWSMessage(entryID int64, direction string, opcode byte, payload []byte) { if s.store == nil || entryID == 0 { return } if _, err := s.store.AddWSMessage(store.WSMessage{ EntryID: entryID, StartedAt: time.Now(), Direction: direction, Opcode: int(opcode), Payload: payload, }); err != nil { log.Printf("store websocket message: %v", err) } } // enabledRules fetches the current enabled match-and-replace rules for // scope ("request" or "response") fresh from the store on every call - // simple and always current, and cheap enough (a local, in-process // SQLite query) not to bother caching for how this is actually used. func (s *Server) enabledRules(scope string) ([]rules.Rule, error) { if s.store == nil { return nil, nil } return s.store.EnabledRules(scope) } // record stores one history entry and notifies OnEntry, returning the // entry's assigned ID (0 if it wasn't stored at all - no store attached, // scope excluded it, or the insert itself failed) so a caller that needs // to attach more data to this specific entry afterward (see // handleWebSocketUpgrade's ws_messages rows) can do so. func (s *Server) record(started time.Time, duration time.Duration, scheme, host string, r *http.Request, reqRaw []byte, reqExact, reqTruncated bool, respRaw []byte, respExact, respTruncated bool, status int, errMsg string) int64 { if s.store == nil { return 0 } // Scope only filters what gets recorded here - the request has // already been forwarded and its response already written to the // client by the time record() runs (see forward()), so an // out-of-scope host still proxies completely normally, it just // doesn't clutter history. Repeat/Intrude (recordRaw, a different // function) deliberately don't go through this check: a user // explicitly resending or fuzzing a specific request wants to see // the result regardless of scope, which exists to cut passive- // capture noise, not to second-guess a deliberate action. if scopeRules, err := s.store.ListScopeRules(); err == nil && !scope.InScope(scopeRules, host) { return 0 } e := &store.Entry{ StartedAt: started, Duration: duration, Method: r.Method, Scheme: scheme, Host: host, Path: r.URL.Path, StatusCode: status, RequestRaw: reqRaw, ResponseRaw: respRaw, RequestExact: reqExact, ResponseExact: respExact, RequestTruncated: reqTruncated, ResponseTruncated: respTruncated, Error: errMsg, } id, err := s.store.Insert(e) if err != nil { log.Printf("store history entry: %v", err) return 0 } if s.OnEntry != nil { s.OnEntry(store.Summary{ ID: id, StartedAt: e.StartedAt, Duration: e.Duration, Method: e.Method, Scheme: e.Scheme, Host: e.Host, Path: e.Path, StatusCode: e.StatusCode, ReqSize: len(reqRaw), RespSize: len(respRaw), Error: errMsg, Source: "proxy", }) } return id } // singleConnListener adapts one already-accepted net.Conn into a // net.Listener so http.Serve can drive it, returning io.EOF from the // second Accept once the connection closes. type singleConnListener struct { ch chan net.Conn addr net.Addr } // newSingleConnListener wraps c for one Accept, teeConn on the outside // so a *teeConn is what ConnContext sees (see withClientTee) - wrapping // it the other way around lets closeSignalConn's concrete type mask the // teeConn from that type assertion, silently disabling capture. func newSingleConnListener(c net.Conn) *singleConnListener { ch := make(chan net.Conn, 1) signaled := &closeSignalConn{Conn: c, onClose: sync.OnceFunc(func() { close(ch) })} ch <- newTeeConn(signaled) return &singleConnListener{ch: ch, addr: c.LocalAddr()} } func (l *singleConnListener) Accept() (net.Conn, error) { c, ok := <-l.ch if !ok { return nil, io.EOF } return c, nil } func (l *singleConnListener) Close() error { return nil } func (l *singleConnListener) Addr() net.Addr { return l.addr } type closeSignalConn struct { net.Conn onClose func() } func (c *closeSignalConn) Close() error { err := c.Conn.Close() c.onClose() return err } // handleHTTP forwards a plain (non-CONNECT) proxy request, copies the // response back, and records it to history. func (s *Server) handleHTTP(w http.ResponseWriter, r *http.Request) { if !r.URL.IsAbs() { http.Error(w, "mitmux: request must use absolute-form URI (configure as a proxy, not a target)", http.StatusBadRequest) return } if isCertDownloadHost(r.URL.Host) { s.serveCACert(w) return } host := r.URL.Host dial := func(ctx context.Context) (net.Conn, string, error) { return dialUpstreamPlain(ctx, host, s.UpstreamProxy) } s.forward(dial, r.URL.Scheme, r.URL.Host, w, r) } // certDownloadHost is a magic hostname mitmux intercepts and answers // itself, serving its own CA certificate - reachable over plain HTTP // from any client configured to use mitmux as its proxy, including a // mobile browser, which otherwise has no easy way to get a file onto // the device to trust as a CA at all. Deliberately not a real, // resolvable domain (".cert" isn't a registered TLD) so it can never // collide with an actual site someone meant to visit - the same idea // as mitmproxy's own http://mitm.it/, arrived at independently rather // than reusing their domain. HTTP only, on purpose: fetching this over // HTTPS would require the client to already trust mitmux's CA to MITM // that very connection - exactly the chicken-and-egg problem this page // exists to solve, so intercepting it on the CONNECT/TLS path wouldn't // make sense and isn't attempted. const certDownloadHost = "mitmux.cert" func isCertDownloadHost(hostPort string) bool { host := hostPort if h, _, err := net.SplitHostPort(hostPort); err == nil { host = h } return strings.EqualFold(host, certDownloadHost) } // serveCACert answers with the CA certificate as a download. The // content type (application/x-x509-ca-cert) is what makes iOS and // Android offer to install it as a trusted certificate directly from // the browser's download prompt, rather than just saving a plain file. func (s *Server) serveCACert(w http.ResponseWriter) { w.Header().Set("Content-Type", "application/x-x509-ca-cert") w.Header().Set("Content-Disposition", `attachment; filename="mitmux-ca.pem"`) w.WriteHeader(http.StatusOK) w.Write(s.ca.CertPEM) } func stripHopByHop(h http.Header) { for _, k := range hopByHopHeaders { h.Del(k) } } // stripHopByHopKeepingUpgrade is stripHopByHop for a request that's // asking to upgrade the connection (see isWebSocketUpgradeRequest): // every other hop-by-hop header is still stripped, but Connection and // Upgrade are left alone since they're the upgrade request itself, not // leftover framing from the client's hop to mitmux. func stripHopByHopKeepingUpgrade(h http.Header) { for _, k := range hopByHopHeaders { if k == "Connection" || k == "Upgrade" { continue } h.Del(k) } }