// Package simplex runs the SimpleX Chat command-line client and talks // to it over its WebSocket API. // // The protocol, documented in the simplex-chat repository under bots/: // a command goes out as {"corrId": "...", "cmd": "..."}, and the client // answers it with {"corrId": "...", "resp": {...}} carrying the same id. // Everything it sends without a corrId is an event. The API has no // authentication; the client binds it to localhost only. package simplex import ( "context" "encoding/json" "errors" "fmt" "log/slog" "strconv" "strings" "sync" "github.com/gorilla/websocket" ) // maxMessageSize bounds one message from the chat client. The largest // thing it sends is a record carrying a contact's profile picture, well // under this. const maxMessageSize = 16 << 20 var ( // ErrClosed is returned by Command once the connection has ended. ErrClosed = errors.New("connection to the chat client closed") errUnexpected = errors.New("unexpected response") errCommand = errors.New("command failed") ) // EventHandler receives each event the chat client sends. Events are // delivered one at a time, on the goroutine that also delivers command // responses: a handler may Send, but must never wait on Command, whose // response could then never arrive. type EventHandler func(c *Client, ev Event) // Client is a connection to the chat client's WebSocket API. type Client struct { conn *websocket.Conn log *slog.Logger onEvent EventHandler // writeMu serialises writes: the connection allows one writer at // a time, and Command and Send are called from different // goroutines. writeMu sync.Mutex mu sync.Mutex lastID uint64 waiting map[string]chan Event done chan struct{} err error // why the read loop ended; valid once done is closed } // Dial connects to the chat client's API at url and starts reading from // it, passing every event to onEvent. func Dial( ctx context.Context, url string, log *slog.Logger, onEvent EventHandler, ) (*Client, error) { conn, resp, err := websocket.DefaultDialer.DialContext(ctx, url, nil) if resp != nil { _ = resp.Body.Close() } if err != nil { return nil, fmt.Errorf("connecting to %s: %w", url, err) } conn.SetReadLimit(maxMessageSize) c := &Client{ conn: conn, log: log, onEvent: onEvent, waiting: make(map[string]chan Event), done: make(chan struct{}), } go c.read() return c, nil } // Done is closed when the connection ends; Err then says why. func (c *Client) Done() <-chan struct{} { return c.done } // Err returns why the connection ended. Call it only after Done is // closed. func (c *Client) Err() error { return c.err } // Close ends the connection. func (c *Client) Close() error { err := c.conn.Close() if err != nil { return fmt.Errorf("closing connection: %w", err) } return nil } // ActiveUser returns the chat client's active user profile. func (c *Client) ActiveUser(ctx context.Context) (User, error) { var r struct { User User `json:"user"` } err := c.command(ctx, cmdShowActiveUser, TypeActiveUser, &r) return r.User, err } // Address returns the user's long-term contact address, and false if // the user has none. func (c *Client) Address(ctx context.Context, userID int64) (ConnLink, bool, error) { //nolint:tagliatelle // the chat client's wire format. var r struct { ContactLink struct { ConnLinkContact ConnLink `json:"connLinkContact"` } `json:"contactLink"` } err := c.command(ctx, cmdShowAddress(userID), TypeUserContactLink, &r) var cerr *CommandError if errors.As(err, &cerr) && cerr.Detail == "userContactLinkNotFound" { return ConnLink{}, false, nil } if err != nil { return ConnLink{}, false, err } return r.ContactLink.ConnLinkContact, true, nil } // CreateAddress creates the user's long-term contact address. func (c *Client) CreateAddress(ctx context.Context, userID int64) (ConnLink, error) { //nolint:tagliatelle // the chat client's wire format. var r struct { ConnLinkContact ConnLink `json:"connLinkContact"` } err := c.command(ctx, cmdCreateAddress(userID), TypeUserContactLinkCreated, &r) return r.ConnLinkContact, err } // SetAddressSettings replaces the settings of the user's address. func (c *Client) SetAddressSettings( ctx context.Context, userID int64, s AddressSettings, ) error { cmd, err := cmdSetAddressSettings(userID, s) if err != nil { return err } return c.command(ctx, cmd, TypeUserContactLinkUpdated, nil) } // SendText sends a text message to a contact, as a reply to the message // quotedItemID (0 for none). It does not wait for the chat client to // accept it; a failure is logged when the client's answer arrives. func (c *Client) SendText(contactID, quotedItemID int64, text string) error { cmd, err := cmdSendText(contactID, quotedItemID, text) if err != nil { return err } c.mu.Lock() id := c.nextID() c.mu.Unlock() return c.write(id, cmd) } // CommandError is a command the chat client refused. Type and Detail // are the discriminators of its chatError record, such as "errorStore" // and "userContactLinkNotFound". type CommandError struct { Type string Detail string } func (e *CommandError) Error() string { return fmt.Sprintf("%s: %s/%s", errCommand, e.Type, e.Detail) } func (e *CommandError) Unwrap() error { return errCommand } // command sends cmd, waits for its response, and, if the response has // type want, decodes it into out (unless out is nil). func (c *Client) command(ctx context.Context, cmd, want string, out any) error { ch := make(chan Event, 1) c.mu.Lock() id := c.nextID() c.waiting[id] = ch c.mu.Unlock() defer func() { c.mu.Lock() delete(c.waiting, id) c.mu.Unlock() }() err := c.write(id, cmd) if err != nil { return err } var ev Event select { case ev = <-ch: case <-c.done: return fmt.Errorf("%w: %w", ErrClosed, c.err) case <-ctx.Done(): return fmt.Errorf("waiting for a response to %q: %w", cmdName(cmd), ctx.Err()) } switch ev.Type { case want: if out == nil { return nil } return ev.Decode(out) case TypeChatCmdError: return commandError(ev) default: return fmt.Errorf("%w to %q: %s", errUnexpected, cmdName(cmd), ev.Type) } } // nextID returns a fresh correlation id. The caller holds c.mu. func (c *Client) nextID() string { c.lastID++ return strconv.FormatUint(c.lastID, 10) } func (c *Client) write(id, cmd string) error { c.writeMu.Lock() defer c.writeMu.Unlock() err := c.conn.WriteJSON(command{CorrID: id, Cmd: cmd}) if err != nil { return fmt.Errorf("sending %q: %w", cmdName(cmd), err) } return nil } // read is the only reader of the connection. It runs until the // connection fails or is closed. func (c *Client) read() { defer close(c.done) for { _, data, err := c.conn.ReadMessage() if err != nil { c.err = err return } c.dispatch(data) } } // dispatch routes one message: a response to whoever waits for it, an // event to the handler. A message that does not parse is logged and // skipped rather than ending the connection; the API documentation // warns that records change between releases. func (c *Client) dispatch(data []byte) { var env envelope err := json.Unmarshal(data, &env) if err != nil { c.log.Warn("undecodable message from the chat client", "error", err) return } var head tagged err = json.Unmarshal(env.Resp, &head) if err != nil { c.log.Warn("undecodable record from the chat client", "error", err) return } ev := Event{Type: head.Type, raw: env.Resp} if env.CorrID == "" { c.onEvent(c, ev) return } c.mu.Lock() ch, ok := c.waiting[env.CorrID] c.mu.Unlock() if ok { ch <- ev return } // The response to a SendText: nobody waits for it, so a failure // is reported here or nowhere. if ev.Type == TypeChatCmdError { c.log.Warn("sending a message failed", "error", commandError(ev)) } } func commandError(ev Event) error { var r cmdError err := ev.Decode(&r) if err != nil { return err } e := &CommandError{Type: r.ChatError.Type} for _, detail := range []*tagged{ r.ChatError.ErrorType, r.ChatError.StoreError, r.ChatError.AgentError, } { if detail != nil { e.Detail = detail.Type } } return e } // cmdName is a command without its arguments, for error messages: the // arguments of /_send are a message someone wrote. func cmdName(cmd string) string { name, _, _ := strings.Cut(cmd, " ") return name }