srdusr
aboutsummaryrefslogtreecommitdiffstats
path: root/internal/ipc
diff options
context:
space:
mode:
Diffstat (limited to 'internal/ipc')
-rw-r--r--internal/ipc/ipc.go126
-rw-r--r--internal/ipc/server.go136
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)
+ }
+}