// 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" "sync" "time" "mitmux/internal/clientcert" "mitmux/internal/proxy" "mitmux/internal/rules" "mitmux/internal/scope" "mitmux/internal/store" ) // Request is sent by a client to the daemon. type Request struct { Type string `json:"type"` // "list", "get", "subscribe", "repeat", "intrude", "rules_list", "rules_save", "rules_delete", "rules_toggle", "delete_entry", or "clear_history" Limit int `json:"limit,omitempty"` BeforeID int64 `json:"before_id,omitempty"` ID int64 `json:"id,omitempty"` // For "list": a non-empty Query switches from most-recent-first to // an FTS5 search (see store.Store.Search for syntax), ranked by // relevance. Query string `json:"query,omitempty"` // For "repeat": send Raw to scheme://host exactly as given. // For "intrude": Raw is the §marked§ template - see proxy.Intrude. Scheme string `json:"scheme,omitempty"` Host string `json:"host,omitempty"` Raw []byte `json:"raw,omitempty"` // For "intrude": the attack mode (proxy.Sniper and so on - empty // defaults to Sniper) and its payload sets. Sniper and BatteringRam // only ever use PayloadSets[0] (one shared set); Pitchfork and // ClusterBomb require exactly one set per §marked§ position, in // order - see proxy.Intrude. Mode proxy.AttackMode `json:"mode,omitempty"` PayloadSets [][]string `json:"payload_sets,omitempty"` // For "intrude": optional Go regexps evaluated against each result's // response bytes. GrepMatch flags whether it matched at all; // GrepExtract additionally captures text (first submatch if the // pattern has a capturing group, else the whole match) into the // result. Either or both may be empty to skip that check. Compiled // and validated once, server-side, before the attack starts - a bad // pattern fails the same way a bad marker or empty payload set does. GrepMatch string `json:"grep_match,omitempty"` GrepExtract string `json:"grep_extract,omitempty"` // For "rules_save": add (Rule.ID == 0) or update (Rule.ID != 0) a // match-and-replace rule. For "rules_delete"/"rules_toggle": RuleID // (and RuleEnabled for toggle) identify the target. Rule *rules.Rule `json:"rule,omitempty"` RuleID int64 `json:"rule_id,omitempty"` RuleEnabled bool `json:"rule_enabled,omitempty"` // For "scope_add": the new rule (always an add - scope rules are // simple enough that edit-in-place isn't worth a separate update // path; delete and re-add covers it). For "scope_delete"/ // "scope_toggle": ScopeRuleID (and RuleEnabled for toggle) identify // the target. ScopeRule *scope.Rule `json:"scope_rule,omitempty"` ScopeRuleID int64 `json:"scope_rule_id,omitempty"` // For "clientcert_add": the new certificate (always an add, same // reasoning as scope rules above). For "clientcert_delete"/ // "clientcert_toggle": ClientCertID (and RuleEnabled for toggle) // identify the target. ClientCert *clientcert.Cert `json:"client_cert,omitempty"` ClientCertID int64 `json:"client_cert_id,omitempty"` // For "tag_entry": ID identifies the history entry (same field // "get"/"set_flagged"/"delete_entry" use). Plugin names who's // tagging it - informational only, not an identity or auth // mechanism, since anything that can reach the socket can claim any // name. Tag is the short marker itself (e.g. "jwt", "authz-bypass"). // Data is an opaque, plugin-defined JSON blob a panel view renders // later without needing this plugin still connected - empty is // fine for a plugin that only needs the tag itself, no extra detail. TagPlugin string `json:"tag_plugin,omitempty"` Tag string `json:"tag,omitempty"` TagData string `json:"tag_data,omitempty"` // For "set_flagged" and "delete_entry": ID identifies the history // entry. "clear_history" needs no fields at all. Flagged bool `json:"flagged,omitempty"` // For "import": entries to insert directly into history - the // client has already parsed/converted them (e.g. from a HAR file, // see cmd/mitmux/har.go); the daemon just stores them. ImportEntries []ImportEntry `json:"import_entries,omitempty"` } // ImportEntry is one entry to insert directly into history via // "import". RequestExact/ResponseExact are always false once stored: // an imported entry is reconstructed from whatever structured format it // came from (HAR, say), never the literal bytes that were on the wire // for the original request - the same situation an HTTP/2 capture is // already in. type ImportEntry struct { Method string `json:"method"` Scheme string `json:"scheme"` Host string `json:"host"` Path string `json:"path"` StatusCode int `json:"status_code"` RequestRaw []byte `json:"request_raw"` ResponseRaw []byte `json:"response_raw"` StartedAt time.Time `json:"started_at"` Duration time.Duration `json:"duration"` } // Response is sent by the daemon to a client. type Response struct { Type string `json:"type"` // "list", "get", "new", "repeat", "rules", "scope_rules", "intrude_result", "intrude_done", "import_done", "status", "flagged", "deleted", "cleared", or "error" Entries []store.Summary `json:"entries,omitempty"` // for "list" Detail *EntryDetail `json:"detail,omitempty"` // for "get" and "repeat" New *store.Summary `json:"new,omitempty"` // for "new" (subscribe push) Rules []rules.Rule `json:"rules,omitempty"` // for "rules" ScopeRules []scope.Rule `json:"scope_rules,omitempty"` // for "scope_rules" ClientCerts []clientcert.Cert `json:"client_certs,omitempty"` // for "client_certs" WSMessages []store.WSMessage `json:"ws_messages,omitempty"` // for "ws_messages" Status *StatusMsg `json:"status,omitempty"` // for "status" // For "import_done": how many entries were actually inserted (a // per-entry insert failure is skipped, not fatal to the batch). Imported int `json:"imported,omitempty"` // For "tag_entry": the new tag row's assigned ID. TagID int64 `json:"tag_id,omitempty"` // For "intrude_result": one completed attack request. IntrudeResult *IntrudeResultMsg `json:"intrude_result,omitempty"` Error string `json:"error,omitempty"` } // StatusMsg is basic daemon info for a TUI status bar. type StatusMsg struct { ProxyAddr string `json:"proxy_addr"` HistoryCount int64 `json:"history_count"` } // IntrudeResultMsg is one completed Intruder attack request. Values holds // what was substituted into each §marked§ position for this request, in // position order - for Sniper, every entry but the one fuzzed position // equals that position's base value; for the other three modes every // entry is an actual payload. type IntrudeResultMsg struct { Iteration int `json:"iteration"` Values []string `json:"values"` EntryID int64 `json:"entry_id"` StatusCode int `json:"status_code"` RespSize int `json:"resp_size"` Duration time.Duration `json:"duration"` Error string `json:"error,omitempty"` GrepMatch bool `json:"grep_match,omitempty"` GrepExtract string `json:"grep_extract,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"` // *Truncated is true when the matching *Exact is false specifically // because capture hit its size cap and dropped bytes off the end - // as opposed to false because there was never a wire-exact // representation to begin with (HTTP/2). Only meaningful alongside // a false *Exact. RequestTruncated bool `json:"request_truncated,omitempty"` ResponseTruncated bool `json:"response_truncated,omitempty"` // Tags is every plugin-contributed marker on this entry - see the // "tag_entry" request. Summary.Tags (from List/Search) is just the // comma-joined names for a compact list-view badge; this is the // full record, including each tag's plugin and opaque Data blob, for // a panel view to render. Tags []store.EntryTag `json:"tags,omitempty"` } // Client talks to a mitmuxd instance for request/response queries // (list, get). Use Subscribe separately for the live-update stream. // // One request/response round trip is in flight on the connection at a // time, guarded by mu - a caller like the mitmux TUI dispatches each // request as its own goroutine (a Bubble Tea tea.Cmd), and without this // two overlapping calls (e.g. opening two entries in quick succession) // would interleave their JSON on the wire or hand one call the other's // response. type Client struct { mu sync.Mutex 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() } // SetFlagged sets the flagged marker on a history entry - a simple // "mark this, revisit later" bit, filterable via flagged:true/false in // Search. // TagEntry marks history entry id with tag, attributed to plugin (any // non-empty name a plugin chooses to identify itself by - informational // only), with an optional opaque data blob a panel view can render // later. This is the core plugin-integration primitive: any process // that can reach the control socket - the reference Go client here, or // a plugin in any other language following the same JSON wire protocol // (see PLAN.md) - can tag entries it finds interesting without mitmux // needing to know anything about it in advance. Returns the new tag's // assigned ID. func (c *Client) TagEntry(entryID int64, plugin, tag, data string) (int64, error) { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(Request{Type: "tag_entry", ID: entryID, TagPlugin: plugin, Tag: tag, TagData: data}); err != nil { return 0, err } var resp Response if err := c.dec.Decode(&resp); err != nil { return 0, err } if resp.Type == "error" { return 0, errors.New(resp.Error) } return resp.TagID, nil } func (c *Client) SetFlagged(id int64, flagged bool) error { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(Request{Type: "set_flagged", ID: id, Flagged: flagged}); err != nil { return err } var resp Response if err := c.dec.Decode(&resp); err != nil { return err } if resp.Type == "error" { return errors.New(resp.Error) } return nil } // DeleteEntry removes a single history entry. func (c *Client) DeleteEntry(id int64) error { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(Request{Type: "delete_entry", ID: id}); err != nil { return err } var resp Response if err := c.dec.Decode(&resp); err != nil { return err } if resp.Type == "error" { return errors.New(resp.Error) } return nil } // ClearHistory removes every history entry. Rules are untouched. func (c *Client) ClearHistory() error { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(Request{Type: "clear_history"}); err != nil { return err } var resp Response if err := c.dec.Decode(&resp); err != nil { return err } if resp.Type == "error" { return errors.New(resp.Error) } return nil } // Import inserts entries directly into history - used to bring in // traffic from an external source (a HAR file, say) rather than // something mitmux itself captured. Returns how many were actually // inserted; a per-entry insert failure is skipped rather than aborting // the whole batch. func (c *Client) Import(entries []ImportEntry) (int, error) { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(Request{Type: "import", ImportEntries: entries}); err != nil { return 0, err } var resp Response if err := c.dec.Decode(&resp); err != nil { return 0, err } if resp.Type == "error" { return 0, errors.New(resp.Error) } return resp.Imported, nil } // Status returns basic daemon info for a status bar. func (c *Client) Status() (*StatusMsg, error) { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(Request{Type: "status"}); 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.Status, nil } // 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) { return c.list(Request{Type: "list", Limit: limit, BeforeID: beforeID}) } // Search returns up to limit history summaries matching an FTS5 query // (see store.Store.Search for syntax), ranked by relevance. func (c *Client) Search(query string, limit int, beforeID int64) ([]store.Summary, error) { return c.list(Request{Type: "list", Query: query, Limit: limit, BeforeID: beforeID}) } func (c *Client) list(req Request) ([]store.Summary, error) { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(req); 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) { c.mu.Lock() defer c.mu.Unlock() 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 } // ListWSMessages returns every WebSocket frame captured for entryID's // connection, in the order they were sent - empty (not an error) if the // entry wasn't a WebSocket upgrade or nothing was captured. func (c *Client) ListWSMessages(entryID int64) ([]store.WSMessage, error) { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(Request{Type: "ws_messages", ID: entryID}); 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.WSMessages, nil } // Repeat sends raw to scheme://host exactly as given (no re-serialization, // no header injection) and returns the resulting entry, including the raw // response bytes. The exchange is also recorded to history. func (c *Client) Repeat(scheme, host string, raw []byte) (*EntryDetail, error) { c.mu.Lock() defer c.mu.Unlock() if err := c.enc.Encode(Request{Type: "repeat", Scheme: scheme, Host: host, Raw: raw}); 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 } // ListRules returns every match-and-replace rule. func (c *Client) ListRules() ([]rules.Rule, error) { c.mu.Lock() defer c.mu.Unlock() return c.rulesRoundTrip(Request{Type: "rules_list"}) } // SaveRule adds r (if r.ID == 0) or updates the existing rule with that // ID, and returns its ID. func (c *Client) SaveRule(r rules.Rule) (int64, error) { c.mu.Lock() defer c.mu.Unlock() saved, err := c.rulesRoundTrip(Request{Type: "rules_save", Rule: &r}) if err != nil { return 0, err } if len(saved) == 0 { return 0, errors.New("rules_save: daemon returned no rule") } return saved[0].ID, nil } // DeleteRule removes a rule. func (c *Client) DeleteRule(id int64) error { c.mu.Lock() defer c.mu.Unlock() _, err := c.rulesRoundTrip(Request{Type: "rules_delete", RuleID: id}) return err } // SetRuleEnabled toggles a rule without touching its other fields. func (c *Client) SetRuleEnabled(id int64, enabled bool) error { c.mu.Lock() defer c.mu.Unlock() _, err := c.rulesRoundTrip(Request{Type: "rules_toggle", RuleID: id, RuleEnabled: enabled}) return err } func (c *Client) rulesRoundTrip(req Request) ([]rules.Rule, error) { if err := c.enc.Encode(req); 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.Rules, nil } // ListScopeRules returns every scope rule, enabled or not. func (c *Client) ListScopeRules() ([]scope.Rule, error) { c.mu.Lock() defer c.mu.Unlock() return c.scopeRoundTrip(Request{Type: "scope_list"}) } // AddScopeRule adds r and returns its assigned ID. func (c *Client) AddScopeRule(r scope.Rule) (int64, error) { c.mu.Lock() defer c.mu.Unlock() saved, err := c.scopeRoundTrip(Request{Type: "scope_add", ScopeRule: &r}) if err != nil { return 0, err } if len(saved) == 0 { return 0, errors.New("scope_add: daemon returned no rule") } return saved[0].ID, nil } // DeleteScopeRule removes a scope rule. func (c *Client) DeleteScopeRule(id int64) error { c.mu.Lock() defer c.mu.Unlock() _, err := c.scopeRoundTrip(Request{Type: "scope_delete", ScopeRuleID: id}) return err } // SetScopeRuleEnabled toggles a scope rule without touching its pattern. func (c *Client) SetScopeRuleEnabled(id int64, enabled bool) error { c.mu.Lock() defer c.mu.Unlock() _, err := c.scopeRoundTrip(Request{Type: "scope_toggle", ScopeRuleID: id, RuleEnabled: enabled}) return err } func (c *Client) scopeRoundTrip(req Request) ([]scope.Rule, error) { if err := c.enc.Encode(req); 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.ScopeRules, nil } // ListClientCerts returns every client certificate, enabled or not. func (c *Client) ListClientCerts() ([]clientcert.Cert, error) { c.mu.Lock() defer c.mu.Unlock() return c.clientCertRoundTrip(Request{Type: "clientcert_list"}) } // AddClientCert adds cert and returns its assigned ID. func (c *Client) AddClientCert(cert clientcert.Cert) (int64, error) { c.mu.Lock() defer c.mu.Unlock() saved, err := c.clientCertRoundTrip(Request{Type: "clientcert_add", ClientCert: &cert}) if err != nil { return 0, err } if len(saved) == 0 { return 0, errors.New("clientcert_add: daemon returned no certificate") } return saved[0].ID, nil } // DeleteClientCert removes a client certificate. func (c *Client) DeleteClientCert(id int64) error { c.mu.Lock() defer c.mu.Unlock() _, err := c.clientCertRoundTrip(Request{Type: "clientcert_delete", ClientCertID: id}) return err } // SetClientCertEnabled toggles a client cert without touching its // content. func (c *Client) SetClientCertEnabled(id int64, enabled bool) error { c.mu.Lock() defer c.mu.Unlock() _, err := c.clientCertRoundTrip(Request{Type: "clientcert_toggle", ClientCertID: id, RuleEnabled: enabled}) return err } func (c *Client) clientCertRoundTrip(req Request) ([]clientcert.Cert, error) { if err := c.enc.Encode(req); 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.ClientCerts, 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 } // Intrude starts an attack (see proxy.Intrude and proxy.AttackMode): // template must contain at least one §marked§ position. mode selects how // payloadSets combine across positions; "" defaults to Sniper. // grepMatch/grepExtract are optional Go regexps evaluated server-side // against each result's response bytes (empty string disables either // check) - see IntrudeResultMsg. Unlike Subscribe's live feed, no result // is ever dropped for a slow consumer - each one is the attack's actual // data, not a notification with the real thing recoverable elsewhere. A // setup error (bad markers, empty payload set, too many requests, an // unparseable grep regexp) is returned directly rather than through the // channel. The returned channel closes when the attack finishes or the // connection is closed early. func Intrude(path, scheme, host string, template []byte, mode proxy.AttackMode, payloadSets [][]string, grepMatch, grepExtract string) (<-chan IntrudeResultMsg, func() error, error) { conn, err := net.Dial("unix", path) if err != nil { return nil, nil, fmt.Errorf("dial %s: %w", path, err) } req := Request{Type: "intrude", Scheme: scheme, Host: host, Raw: template, Mode: mode, PayloadSets: payloadSets, GrepMatch: grepMatch, GrepExtract: grepExtract} if err := json.NewEncoder(conn).Encode(req); err != nil { conn.Close() return nil, nil, err } dec := json.NewDecoder(conn) var first Response if err := dec.Decode(&first); err != nil { conn.Close() return nil, nil, err } if first.Type == "error" { conn.Close() return nil, nil, errors.New(first.Error) } ch := make(chan IntrudeResultMsg) go func() { defer close(ch) deliver := func(resp Response) bool { switch resp.Type { case "intrude_result": if resp.IntrudeResult != nil { ch <- *resp.IntrudeResult } return true default: // "intrude_done", or anything else - stop return false } } if !deliver(first) { return } for { var resp Response if err := dec.Decode(&resp); err != nil { return } if !deliver(resp) { return } } }() return ch, conn.Close, nil }