Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -160,3 +160,4 @@ There are a bunch of environmental variables that can be set inside the docker c
* `JSON_RPC_IGNORE_AVATARS`: When set to `true`, avatars are not automatically downloaded in json-rpc mode (default: `false`)
* `JSON_RPC_IGNORE_STICKERS`: When set to `true`, sticker packs are not automatically downloaded in json-rpc mode (default: `false`)
* `JSON_RPC_TRUST_NEW_IDENTITIES`: Choose how to trust new identities in json-rpc mode. Supported values: `on-first-use`, `always`, `never`. (default: `on-first-use`)
* `JSON_RPC_RECEIVE_MODE`: Controls when signal-cli pulls inbound messages from the Signal servers in json-rpc mode. Supported values: `on-start` (default), `manual`. In the default `on-start` mode, signal-cli auto-receives messages and pushes them as JSON-RPC notifications regardless of whether any websocket subscriber is currently attached to `/v1/receive/{number}` — messages arriving while no subscriber is attached (e.g. during a brief subscriber redeploy) are silently dropped. In `manual` mode signal-cli only fetches messages while at least one websocket subscriber is attached; while no subscriber is attached, messages remain queued server-side under Signal's normal retention rules and are delivered on the next subscriber re-attach. Recommended for any deployment where the websocket consumer is restarted or scaled.
2 changes: 1 addition & 1 deletion src/api/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -572,7 +572,7 @@ func (a *Api) SendV2(c *gin.Context) {
}

func (a *Api) handleSignalReceive(ws *websocket.Conn, number string, stop chan struct{}) {
receiveChannel, channelUuid, err := a.signalClient.GetReceiveChannel()
receiveChannel, channelUuid, err := a.signalClient.GetReceiveChannel(number)
if err != nil {
log.Error("Couldn't get receive channel: ", err.Error())
return
Expand Down
4 changes: 2 additions & 2 deletions src/client/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -1059,12 +1059,12 @@ func (s *SignalClient) Receive(number string, timeout int64, ignoreAttachments b
}
}

func (s *SignalClient) GetReceiveChannel() (chan JsonRpc2ReceivedMessage, string, error) {
func (s *SignalClient) GetReceiveChannel(number string) (chan JsonRpc2ReceivedMessage, string, error) {
jsonRpc2Client, err := s.getJsonRpc2Client()
if err != nil {
return nil, "", err
}
return jsonRpc2Client.GetReceiveChannel()
return jsonRpc2Client.GetReceiveChannel(number)
}

func (s *SignalClient) RemoveReceiveChannel(channelUuid string) {
Expand Down
140 changes: 129 additions & 11 deletions src/client/jsonrpc2.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"bytes"
"encoding/json"
"errors"
"fmt"
"net"
"net/http"
"strconv"
Expand Down Expand Up @@ -56,15 +57,26 @@ func (r *RateLimitErrorType) Error() string {
return r.Err.Error()
}

// receiveSubscription tracks the state of a manual-mode signal-cli
// subscribeReceive call: the subscription id assigned by signal-cli and
// a refcount of websocket subscribers attached to that account.
type receiveSubscription struct {
id int64
refcount int
}

type JsonRpc2Client struct {
conn net.Conn
receivedResponsesById map[string]chan JsonRpc2MessageResponse
receivedMessagesChannels map[string]chan JsonRpc2ReceivedMessage
signalCliApiConfig *utils.SignalCliApiConfig
number string
receivedMessagesMutex sync.Mutex
receivedResponsesMutex sync.Mutex
address string
conn net.Conn
receivedResponsesById map[string]chan JsonRpc2MessageResponse
receivedMessagesChannels map[string]chan JsonRpc2ReceivedMessage
receiveSubscriptions map[string]*receiveSubscription // account -> sub state
channelAccountByUuid map[string]string // channelUuid -> account
signalCliApiConfig *utils.SignalCliApiConfig
number string
receivedMessagesMutex sync.Mutex
receivedResponsesMutex sync.Mutex
receiveSubscriptionsMutex sync.Mutex
address string
}

func NewJsonRpc2Client(signalCliApiConfig *utils.SignalCliApiConfig, number string) *JsonRpc2Client {
Expand All @@ -73,6 +85,8 @@ func NewJsonRpc2Client(signalCliApiConfig *utils.SignalCliApiConfig, number stri
number: number,
receivedResponsesById: make(map[string]chan JsonRpc2MessageResponse),
receivedMessagesChannels: make(map[string]chan JsonRpc2ReceivedMessage),
receiveSubscriptions: make(map[string]*receiveSubscription),
channelAccountByUuid: make(map[string]string),
}
}

Expand Down Expand Up @@ -236,6 +250,19 @@ func (r *JsonRpc2Client) ReceiveData(number string, receiveWebhookUrl string) {
var resp1 JsonRpc2ReceivedMessage
json.Unmarshal([]byte(str), &resp1)
if resp1.Method == "receive" {
// In manual receive-mode signal-cli wraps the envelope in
// {"subscription":N,"result":{...}}; in auto mode it sends
// the envelope directly. Unwrap so the broadcast format is
// the same in both modes and downstream consumers (e.g. the
// websocket handler) don't have to know which mode is in use.
var manualWrapper struct {
Subscription int64 `json:"subscription"`
Result json.RawMessage `json:"result"`
}
if err := json.Unmarshal(resp1.Params, &manualWrapper); err == nil && len(manualWrapper.Result) > 0 {
resp1.Params = manualWrapper.Result
}

r.receivedMessagesMutex.Lock()
for _, c := range r.receivedMessagesChannels {
select {
Expand All @@ -244,7 +271,6 @@ func (r *JsonRpc2Client) ReceiveData(number string, receiveWebhookUrl string) {
default:
log.Debug("Couldn't send message to golang channel, as there's no receiver")
}
continue
}
r.receivedMessagesMutex.Unlock()

Expand All @@ -270,16 +296,102 @@ func (r *JsonRpc2Client) ReceiveData(number string, receiveWebhookUrl string) {
}
}

func (r *JsonRpc2Client) GetReceiveChannel() (chan JsonRpc2ReceivedMessage, string, error) {
c := make(chan JsonRpc2ReceivedMessage)
// subscribeReceive starts receiving messages for an account by calling
// signal-cli's subscribeReceive JSON-RPC method. Only relevant when
// signal-cli was launched with --receive-mode=manual; in auto mode the
// daemon pushes notifications without an explicit subscribe call.
// Returns the subscription id assigned by signal-cli.
func (r *JsonRpc2Client) subscribeReceive(account string) (int64, error) {
type subscribeReceiveArgs struct{}
resultStr, err := r.getRaw("subscribeReceive", &account, subscribeReceiveArgs{})
if err != nil {
return 0, err
}
var subscriptionId int64
if err := json.Unmarshal([]byte(resultStr), &subscriptionId); err != nil {
return 0, fmt.Errorf("subscribeReceive: couldn't parse subscription id from %q: %w", resultStr, err)
}
return subscriptionId, nil
}

// unsubscribeReceive cancels a manual-mode subscription previously
// returned by subscribeReceive.
func (r *JsonRpc2Client) unsubscribeReceive(account string, subscriptionId int64) error {
type unsubscribeReceiveArgs struct {
Subscription int64 `json:"subscription"`
}
_, err := r.getRaw("unsubscribeReceive", &account, unsubscribeReceiveArgs{Subscription: subscriptionId})
return err
}

// acquireReceiveSubscription ensures an active manual-mode subscription
// exists for the given account, refcounting concurrent websocket
// subscribers. The first caller for an account triggers a real
// subscribeReceive RPC; subsequent callers just bump the refcount.
func (r *JsonRpc2Client) acquireReceiveSubscription(account string) error {
r.receiveSubscriptionsMutex.Lock()
defer r.receiveSubscriptionsMutex.Unlock()

if sub, ok := r.receiveSubscriptions[account]; ok {
sub.refcount++
return nil
}
id, err := r.subscribeReceive(account)
if err != nil {
return err
}
r.receiveSubscriptions[account] = &receiveSubscription{id: id, refcount: 1}
log.Infof("Subscribed to receive notifications for account %s (subscription=%d)", account, id)
return nil
}

// releaseReceiveSubscription decrements the per-account refcount and
// cancels the subscription with signal-cli if it drops to zero.
func (r *JsonRpc2Client) releaseReceiveSubscription(account string) {
r.receiveSubscriptionsMutex.Lock()
defer r.receiveSubscriptionsMutex.Unlock()

sub, ok := r.receiveSubscriptions[account]
if !ok {
return
}
sub.refcount--
if sub.refcount > 0 {
return
}
if err := r.unsubscribeReceive(account, sub.id); err != nil {
log.Warnf("unsubscribeReceive failed for account %s (subscription=%d): %s", account, sub.id, err.Error())
} else {
log.Infof("Unsubscribed from receive notifications for account %s (subscription=%d)", account, sub.id)
}
delete(r.receiveSubscriptions, account)
}

// GetReceiveChannel returns a channel that will receive messages for the
// given account. If signal-cli is in manual receive-mode, it also acquires
// a subscription so that signal-cli starts forwarding messages for the
// account; in auto mode, account is unused (notifications flow regardless).
//
// account may be empty when the caller does not need a subscription
// (e.g. legacy callers in auto mode); in that case no subscribeReceive
// RPC is issued.
func (r *JsonRpc2Client) GetReceiveChannel(account string) (chan JsonRpc2ReceivedMessage, string, error) {
c := make(chan JsonRpc2ReceivedMessage, 64)

channelUuid, err := uuid.NewV4()
if err != nil {
return c, "", err
}

if account != "" {
if err := r.acquireReceiveSubscription(account); err != nil {

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think I don't understand that code part here. Isn't account always non-empty? So, I guess we would end up always in this code branch, no matter which mode is selected via JSON_RPC_RECEIVE_MODE.

@MykolaBalakin MykolaBalakin Aug 11, 2026

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Looking further into the code, I'd say that even if gin, for any reason, calls this method with an empty string, this method shouldn't be responsible for deciding whether the value of account is a valid account. If signal-cli does not return an error for its subscribeReceive call, then account is a valid account.

So if @strichter has no objections and there are no other remarks from @bbernhard, I’d push a fix to this comment to get the PR merged and released.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks!

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@strichter, I've created a PR to your fork just so that the code you've implemented stays as your PR (otherwise I'll need to push a commit into my fork and create a new PR to the bbernard's repo, which will look like I've implemented the actual receive-mode fix).

return c, "", fmt.Errorf("subscribeReceive failed for account %s: %w", account, err)
}
}

r.receivedMessagesMutex.Lock()
r.receivedMessagesChannels[channelUuid.String()] = c
r.channelAccountByUuid[channelUuid.String()] = account
r.receivedMessagesMutex.Unlock()

return c, channelUuid.String(), nil
Expand All @@ -288,5 +400,11 @@ func (r *JsonRpc2Client) GetReceiveChannel() (chan JsonRpc2ReceivedMessage, stri
func (r *JsonRpc2Client) RemoveReceiveChannel(channelUuid string) {
r.receivedMessagesMutex.Lock()
delete(r.receivedMessagesChannels, channelUuid)
account := r.channelAccountByUuid[channelUuid]
delete(r.channelAccountByUuid, channelUuid)
r.receivedMessagesMutex.Unlock()

if account != "" {
r.releaseReceiveSubscription(account)
}
}
20 changes: 18 additions & 2 deletions src/scripts/jsonrpc2-helper.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ import (
const supervisorctlConfigTemplate = `
[program:%s]
process_name=%s
command=%s --output=json --config %s%s daemon %s%s%s%s --tcp 127.0.0.1:%d
command=%s --output=json --config %s%s daemon%s %s%s%s%s --tcp 127.0.0.1:%d
autostart=true
autorestart=true
startretries=10
Expand Down Expand Up @@ -75,6 +75,22 @@ func main() {
signalCliIgnoreStickers = " --ignore-stickers"
}

// Receive mode: by default signal-cli auto-receives messages on the
// daemon and pushes them as JSON-RPC notifications, regardless of
// whether any websocket subscriber is connected. In a deploy where
// the subscribing client briefly disconnects, those notifications
// have nowhere to go and are silently dropped. With manual receive
// mode, signal-cli only fetches messages when explicitly requested
// via subscribeReceive — see jsonrpc2.go — so messages stay buffered
// on the Signal servers until a subscriber re-attaches.
signalCliReceiveMode := ""
receiveMode := utils.GetEnv("JSON_RPC_RECEIVE_MODE", "")
if receiveMode == "manual" {
signalCliReceiveMode = " --receive-mode=manual"
} else if receiveMode != "" && receiveMode != "on-start" {
log.Fatal("Invalid JSON_RPC_RECEIVE_MODE environment variable set! Must be 'manual' or 'on-start'.")
}

supervisorctlProgramName := "signal-cli-json-rpc-1"
supervisorctlLogFolder := "/var/log/" + supervisorctlProgramName
_, err := exec.Command("mkdir", "-p", supervisorctlLogFolder).Output()
Expand All @@ -100,7 +116,7 @@ func main() {
supervisorctlConfigFilename := "/etc/supervisor/conf.d/" + "signal-cli-json-rpc-1.conf"

supervisorctlConfig := fmt.Sprintf(supervisorctlConfigTemplate, supervisorctlProgramName, supervisorctlProgramName, signalCliBinary,
signalCliConfigDir, trustNewIdentities, signalCliIgnoreAttachments, signalCliIgnoreStories,
signalCliConfigDir, trustNewIdentities, signalCliReceiveMode, signalCliIgnoreAttachments, signalCliIgnoreStories,
signalCliIgnoreAvatars, signalCliIgnoreStickers, tcpPort,
supervisorctlProgramName, supervisorctlProgramName)

Expand Down