diff options
| author | srdusr <[email protected]> | 2024-02-14 00:37:00 +0200 |
|---|---|---|
| committer | srdusr <[email protected]> | 2024-02-14 00:37:00 +0200 |
| commit | 8d15c2e0b326933f8fc912e3b13f37e78a9bc0b6 (patch) | |
| tree | 567274fe13d0a8aea8356e9edeba33ef778cf20f /internal/ipc | |
| parent | f2f0a2135a202e3e15d2a8cbfbd791aad9b04f3a (diff) | |
| download | mitmux-8d15c2e0b326933f8fc912e3b13f37e78a9bc0b6.tar.gz mitmux-8d15c2e0b326933f8fc912e3b13f37e78a9bc0b6.zip | |
History view: SQLite storage, daemon/TUI split over Unix socket
Implements build-order step 3. Adds:
- internal/store: SQLite (WAL, single-writer) history table, raw
request/response blobs plus metadata for the list view.
- internal/proxy: request/response capture wired into forward(). HTTP/1.1
legs are captured byte-exact via a teeConn that records wire bytes as
they're read, taken right after the message is fully drained (so no
manual re-reading/replaying is needed - RoundTrip's own streaming does
the draining). HTTP/2 legs (no meaningful "raw bytes" of their own -
multiplexed, HPACK-compressed framing) are reconstructed instead, and
marked as such in storage.
- internal/ipc: JSON-over-Unix-socket protocol between mitmuxd (owns the
proxy and the DB) and any client - list/get for queries, subscribe for
a live push stream of newly captured entries. Keeps the proxy engine
independent of the UI, per the architecture sketch.
- cmd/mitmux: Bubble Tea TUI - a live-updating history table and a
request/response detail view with raw bytes.
Two real bugs surfaced during testing and got fixed before commit:
1. http.Transport's HTTP/2 auto-dispatch does a literal *tls.Conn type
assertion on the dialed connection; wrapping it in a capturing teeConn
broke that silently, and HTTP/2 framing got parsed as HTTP/1.1 text.
Fixed by dropping http.Transport for the upstream leg entirely in
favor of an explicit per-protocol round trip (see PLAN.md stack note).
2. singleConnListener wrapped the client teeConn *inside* a
closeSignalConn, so ConnContext's type assertion for it silently
failed and HTTP/1.1 client-side capture never activated. Fixed the
wrap order; verified via direct SQLite inspection that request_exact
flips back to 1 and the stored bytes are genuinely wire-exact
(preserved chunked-encoding framing, original header casing/order).
Verified live: plain HTTP, HTTPS H1.1, HTTPS H2, and a POST with a body,
checked against the raw stored bytes directly in SQLite; IPC list/get/
subscribe against a throwaway client; and the TUI driven end-to-end in a
tmux session (list, detail view, tab between request/response, live
update on a new request while sitting on the list).
Diffstat (limited to 'internal/ipc')
| -rw-r--r-- | internal/ipc/ipc.go | 126 | ||||
| -rw-r--r-- | internal/ipc/server.go | 136 |
2 files changed, 262 insertions, 0 deletions
diff --git a/internal/ipc/ipc.go b/internal/ipc/ipc.go new file mode 100644 index 0000000..c9d92b4 --- /dev/null +++ b/internal/ipc/ipc.go @@ -0,0 +1,126 @@ +// Package ipc is the protocol between mitmuxd (which owns the proxy and +// the history database) and a client such as the TUI, spoken as +// newline-agnostic JSON messages over a Unix domain socket. This keeps +// the proxy engine running independently of any UI attached to it. +package ipc + +import ( + "encoding/json" + "errors" + "fmt" + "net" + + "mitmux/internal/store" +) + +// Request is sent by a client to the daemon. +type Request struct { + Type string `json:"type"` // "list", "get", or "subscribe" + Limit int `json:"limit,omitempty"` + BeforeID int64 `json:"before_id,omitempty"` + ID int64 `json:"id,omitempty"` +} + +// Response is sent by the daemon to a client. +type Response struct { + Type string `json:"type"` // "list", "get", "new", or "error" + Entries []store.Summary `json:"entries,omitempty"` // for "list" + Detail *EntryDetail `json:"detail,omitempty"` // for "get" + New *store.Summary `json:"new,omitempty"` // for "new" (subscribe push) + Error string `json:"error,omitempty"` +} + +// EntryDetail is a full history entry, raw bytes included. +type EntryDetail struct { + store.Summary + RequestRaw []byte `json:"request_raw"` + ResponseRaw []byte `json:"response_raw"` + RequestExact bool `json:"request_exact"` + ResponseExact bool `json:"response_exact"` +} + +// Client talks to a mitmuxd instance for request/response queries +// (list, get). Use Subscribe separately for the live-update stream. +type Client struct { + conn net.Conn + dec *json.Decoder + enc *json.Encoder +} + +// Dial connects to the daemon's control socket at path. +func Dial(path string) (*Client, error) { + conn, err := net.Dial("unix", path) + if err != nil { + return nil, fmt.Errorf("dial %s: %w", path, err) + } + return &Client{conn: conn, dec: json.NewDecoder(conn), enc: json.NewEncoder(conn)}, nil +} + +// Close closes the connection to the daemon. +func (c *Client) Close() error { + return c.conn.Close() +} + +// List returns up to limit history summaries older than beforeID (0 for +// the most recent), newest first. +func (c *Client) List(limit int, beforeID int64) ([]store.Summary, error) { + if err := c.enc.Encode(Request{Type: "list", Limit: limit, BeforeID: beforeID}); err != nil { + return nil, err + } + var resp Response + if err := c.dec.Decode(&resp); err != nil { + return nil, err + } + if resp.Type == "error" { + return nil, errors.New(resp.Error) + } + return resp.Entries, nil +} + +// Get returns the full entry (raw bytes included) for id. +func (c *Client) Get(id int64) (*EntryDetail, error) { + if err := c.enc.Encode(Request{Type: "get", ID: id}); err != nil { + return nil, err + } + var resp Response + if err := c.dec.Decode(&resp); err != nil { + return nil, err + } + if resp.Type == "error" { + return nil, errors.New(resp.Error) + } + return resp.Detail, nil +} + +// Subscribe opens a dedicated connection that streams newly captured +// history entries as they happen. The returned channel is closed when +// the connection ends; call the returned close func to stop early. +func Subscribe(path string) (<-chan store.Summary, func() error, error) { + conn, err := net.Dial("unix", path) + if err != nil { + return nil, nil, fmt.Errorf("dial %s: %w", path, err) + } + if err := json.NewEncoder(conn).Encode(Request{Type: "subscribe"}); err != nil { + conn.Close() + return nil, nil, err + } + + ch := make(chan store.Summary, 64) + go func() { + defer close(ch) + dec := json.NewDecoder(conn) + for { + var resp Response + if err := dec.Decode(&resp); err != nil { + return + } + if resp.Type == "new" && resp.New != nil { + select { + case ch <- *resp.New: + default: + } + } + } + }() + return ch, conn.Close, nil +} diff --git a/internal/ipc/server.go b/internal/ipc/server.go new file mode 100644 index 0000000..11ba033 --- /dev/null +++ b/internal/ipc/server.go @@ -0,0 +1,136 @@ +package ipc + +import ( + "encoding/json" + "log" + "net" + "sync" + + "mitmux/internal/store" +) + +// Hub fans out newly captured history entries to subscribed clients. +type Hub struct { + mu sync.Mutex + subs map[chan store.Summary]struct{} +} + +// NewHub creates an empty Hub. +func NewHub() *Hub { + return &Hub{subs: make(map[chan store.Summary]struct{})} +} + +// Broadcast notifies all current subscribers of e. Slow subscribers +// drop entries rather than blocking the proxy. +func (h *Hub) Broadcast(e store.Summary) { + h.mu.Lock() + defer h.mu.Unlock() + for ch := range h.subs { + select { + case ch <- e: + default: + } + } +} + +func (h *Hub) subscribe() chan store.Summary { + ch := make(chan store.Summary, 64) + h.mu.Lock() + h.subs[ch] = struct{}{} + h.mu.Unlock() + return ch +} + +func (h *Hub) unsubscribe(ch chan store.Summary) { + h.mu.Lock() + delete(h.subs, ch) + h.mu.Unlock() + close(ch) +} + +// Server serves the daemon side of the mitmux control protocol. +type Server struct { + db *store.Store + hub *Hub +} + +// NewServer creates a control-protocol Server backed by db, broadcasting +// through hub. +func NewServer(db *store.Store, hub *Hub) *Server { + return &Server{db: db, hub: hub} +} + +// Serve accepts connections on ln until it returns an error (e.g. the +// listener is closed). +func (s *Server) Serve(ln net.Listener) error { + for { + conn, err := ln.Accept() + if err != nil { + return err + } + go s.handleConn(conn) + } +} + +func (s *Server) handleConn(conn net.Conn) { + defer conn.Close() + dec := json.NewDecoder(conn) + enc := json.NewEncoder(conn) + + for { + var req Request + if err := dec.Decode(&req); err != nil { + return + } + + switch req.Type { + case "list": + entries, err := s.db.List(req.Limit, req.BeforeID) + if err != nil { + enc.Encode(Response{Type: "error", Error: err.Error()}) + continue + } + enc.Encode(Response{Type: "list", Entries: entries}) + + case "get": + e, err := s.db.Get(req.ID) + if err != nil { + enc.Encode(Response{Type: "error", Error: err.Error()}) + continue + } + enc.Encode(Response{Type: "get", Detail: &EntryDetail{ + Summary: store.Summary{ + ID: e.ID, StartedAt: e.StartedAt, Duration: e.Duration, + Method: e.Method, Scheme: e.Scheme, Host: e.Host, Path: e.Path, + StatusCode: e.StatusCode, ReqSize: len(e.RequestRaw), RespSize: len(e.ResponseRaw), + Error: e.Error, + }, + RequestRaw: e.RequestRaw, ResponseRaw: e.ResponseRaw, + RequestExact: e.RequestExact, ResponseExact: e.ResponseExact, + }}) + + case "subscribe": + sub := s.hub.subscribe() + defer s.hub.unsubscribe(sub) + for e := range sub { + e := e + if err := enc.Encode(Response{Type: "new", New: &e}); err != nil { + return + } + } + return + + default: + enc.Encode(Response{Type: "error", Error: "unknown request type: " + req.Type}) + } + } +} + +// LogAndBroadcast is a convenience OnEntry callback: logs the entry and +// broadcasts it through hub. +func LogAndBroadcast(hub *Hub) func(store.Summary) { + return func(sum store.Summary) { + log.Printf("%s %s%s -> %d (%s)", sum.Method, sum.Host, sum.Path, sum.StatusCode, sum.Duration) + hub.Broadcast(sum) + } +} |