2025-11-18 19:59:37 +01:00
|
|
|
package processor
|
|
|
|
|
|
|
|
|
|
import (
|
2025-12-08 21:31:32 +01:00
|
|
|
"context"
|
2025-11-18 19:59:37 +01:00
|
|
|
"crypto/ecdsa"
|
|
|
|
|
"encoding/hex"
|
|
|
|
|
|
|
|
|
|
"github.com/golang/protobuf/proto"
|
|
|
|
|
"github.com/pkg/errors"
|
2025-12-08 21:31:32 +01:00
|
|
|
otelattribute "go.opentelemetry.io/otel/attribute"
|
|
|
|
|
oteltrace "go.opentelemetry.io/otel/trace"
|
2025-11-18 19:59:37 +01:00
|
|
|
"go.uber.org/zap"
|
|
|
|
|
|
2025-12-18 12:24:40 +00:00
|
|
|
"github.com/status-im/status-go/internal/crypto"
|
|
|
|
|
"github.com/status-im/status-go/internal/crypto/types"
|
2025-12-08 21:31:32 +01:00
|
|
|
"github.com/status-im/status-go/internal/instrumentation/trace"
|
2026-04-16 13:22:54 +01:00
|
|
|
adapters "github.com/status-im/status-go/pkg/messaging/adapters"
|
|
|
|
|
common "github.com/status-im/status-go/pkg/messaging/common"
|
2025-12-17 19:40:40 +00:00
|
|
|
"github.com/status-im/status-go/pkg/messaging/controller/utils"
|
2026-04-16 13:22:54 +01:00
|
|
|
encryption "github.com/status-im/status-go/pkg/messaging/layers/encryption"
|
2025-12-17 19:40:40 +00:00
|
|
|
"github.com/status-im/status-go/pkg/messaging/layers/encryption/sharedsecret"
|
|
|
|
|
"github.com/status-im/status-go/pkg/messaging/layers/segmentation"
|
2025-12-22 19:57:36 +00:00
|
|
|
messagingtypes "github.com/status-im/status-go/pkg/messaging/types"
|
2025-11-18 19:59:37 +01:00
|
|
|
"github.com/status-im/status-go/pkg/pubsub"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
const (
|
|
|
|
|
maxNumOfEphemeralKeys = 3
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
var errReliabilityNotStarted = errors.New("reliability not started")
|
|
|
|
|
|
|
|
|
|
type Processor struct {
|
|
|
|
|
identity *ecdsa.PrivateKey
|
2026-04-16 13:22:54 +01:00
|
|
|
stack *common.MessagingStack
|
2025-11-18 19:59:37 +01:00
|
|
|
|
2026-04-16 13:22:54 +01:00
|
|
|
messageConfirmationStorage common.MessageConfirmationPersistence
|
|
|
|
|
hashRatchetStorage common.HashRatchetPersistence
|
2025-11-18 19:59:37 +01:00
|
|
|
|
|
|
|
|
ephemeralKeysManager *EphemeralKeysManager
|
|
|
|
|
|
|
|
|
|
publisher *pubsub.Publisher
|
|
|
|
|
logger *zap.Logger
|
2025-12-08 21:31:32 +01:00
|
|
|
|
|
|
|
|
tracer trace.Tracer
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func NewProcessor(
|
|
|
|
|
identity *ecdsa.PrivateKey,
|
2026-04-16 13:22:54 +01:00
|
|
|
stack *common.MessagingStack,
|
|
|
|
|
messageConfirmationStorage common.MessageConfirmationPersistence,
|
|
|
|
|
hashRatchetStorage common.HashRatchetPersistence,
|
2025-11-18 19:59:37 +01:00
|
|
|
logger *zap.Logger,
|
2025-12-08 21:31:32 +01:00
|
|
|
tracer trace.Tracer,
|
2025-11-18 19:59:37 +01:00
|
|
|
) *Processor {
|
|
|
|
|
return &Processor{
|
|
|
|
|
identity: identity,
|
|
|
|
|
stack: stack,
|
|
|
|
|
messageConfirmationStorage: messageConfirmationStorage,
|
|
|
|
|
hashRatchetStorage: hashRatchetStorage,
|
|
|
|
|
ephemeralKeysManager: NewEphemeralKeysManager(maxNumOfEphemeralKeys),
|
|
|
|
|
publisher: pubsub.NewPublisher(),
|
|
|
|
|
logger: logger.Named("processor"),
|
2025-12-08 21:31:32 +01:00
|
|
|
tracer: tracer,
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (r *Processor) Publisher() *pubsub.Publisher {
|
|
|
|
|
return r.publisher
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (r *Processor) GetEphemeralKey() (*ecdsa.PrivateKey, error) {
|
|
|
|
|
key, err := r.ephemeralKeysManager.GetRandom()
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
_, err = r.stack.Transport.LoadKeyFilters(key)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
return key, nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
func (r *Processor) ProcessMessage(msg *messagingtypes.ReceivedMessage) (*messagingtypes.HandleMessageResponse, error) {
|
2025-11-18 19:59:37 +01:00
|
|
|
response, err := r.processMessage(msg)
|
|
|
|
|
if response == nil || err != nil {
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Process queued hash ratchet messages
|
|
|
|
|
for _, message := range response.messages {
|
|
|
|
|
queuedMessagesResponse, err := r.processQueuedHashRatchetMessages(message.EncryptionLayer.HashRatchetInfo)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
response.messages = append(response.messages, queuedMessagesResponse.messages...)
|
|
|
|
|
response.ackedMessageIDs = append(response.ackedMessageIDs, queuedMessagesResponse.ackedMessageIDs...)
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
return &messagingtypes.HandleMessageResponse{
|
2025-11-18 19:59:37 +01:00
|
|
|
Messages: response.messages,
|
|
|
|
|
AckedMessageIDs: response.ackedMessageIDs,
|
|
|
|
|
}, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type processMessageResponse struct {
|
2025-12-22 19:57:36 +00:00
|
|
|
messages []*messagingtypes.Message
|
2025-12-18 12:24:40 +00:00
|
|
|
ackedMessageIDs []types.HexBytes
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
func (r *Processor) processMessage(m *messagingtypes.ReceivedMessage) (*processMessageResponse, error) {
|
2025-12-18 12:24:40 +00:00
|
|
|
logger := r.logger.With(zap.Stringer("hash", types.HexBytes(m.Hash)))
|
2025-11-18 19:59:37 +01:00
|
|
|
logger.Debug("processing received message")
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
responseMessage := &messagingtypes.Message{}
|
2025-11-18 19:59:37 +01:00
|
|
|
|
|
|
|
|
response := &processMessageResponse{
|
2025-12-22 19:57:36 +00:00
|
|
|
messages: []*messagingtypes.Message{responseMessage},
|
2025-12-18 12:24:40 +00:00
|
|
|
ackedMessageIDs: []types.HexBytes{},
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
err := processTransportLayer(responseMessage, m)
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Error("failed to process transport layer", zap.Error(err))
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-08 21:31:32 +01:00
|
|
|
err = r.processSegmentationLayer(responseMessage)
|
2025-11-18 19:59:37 +01:00
|
|
|
if err != nil {
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-08 21:31:32 +01:00
|
|
|
hashes := [][]byte{m.Hash}
|
|
|
|
|
if responseMessage.SegmentationLayer.Segmented {
|
|
|
|
|
// Segments not completed yet, stop processing
|
|
|
|
|
if !responseMessage.SegmentationLayer.Completed {
|
|
|
|
|
return nil, nil
|
|
|
|
|
}
|
|
|
|
|
hashes = responseMessage.SegmentationLayer.Hashes
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
|
|
|
|
|
2025-12-08 21:31:32 +01:00
|
|
|
ctx, span := r.tracer.Start(trace.DeriveRemoteContext(utils.MergeByteSlices(hashes)), "Processor.processMessage",
|
|
|
|
|
oteltrace.WithAttributes(
|
2025-12-18 12:24:40 +00:00
|
|
|
otelattribute.String("hash", types.EncodeHex(m.Hash)),
|
|
|
|
|
otelattribute.StringSlice("hashes", types.EncodeHexes(hashes)),
|
2025-12-08 21:31:32 +01:00
|
|
|
),
|
|
|
|
|
)
|
|
|
|
|
defer span.End()
|
2025-11-18 19:59:37 +01:00
|
|
|
|
2025-12-08 21:31:32 +01:00
|
|
|
err = r.processEncryptionLayer(ctx, responseMessage, logger)
|
|
|
|
|
if err == nil {
|
|
|
|
|
span.AddEvent("encryption layer processed")
|
|
|
|
|
} else {
|
2025-11-18 19:59:37 +01:00
|
|
|
// Hash ratchet with a group id not found yet, save the message for future processing
|
2026-04-16 13:22:54 +01:00
|
|
|
if err == encryption.ErrHashRatchetGroupIDNotFound && len(responseMessage.EncryptionLayer.HashRatchetInfo) == 1 {
|
2025-11-18 19:59:37 +01:00
|
|
|
info := responseMessage.EncryptionLayer.HashRatchetInfo[0]
|
2025-12-08 21:31:32 +01:00
|
|
|
span.AddEvent("hash ratchet with group id not found yet", oteltrace.WithAttributes(
|
2025-12-18 12:24:40 +00:00
|
|
|
otelattribute.String("groupID", types.ToHex(info.GroupID)),
|
2025-12-08 21:31:32 +01:00
|
|
|
))
|
2025-11-18 19:59:37 +01:00
|
|
|
return nil, r.hashRatchetStorage.SaveMessage(info.GroupID, info.KeyID, m)
|
2025-12-08 21:31:32 +01:00
|
|
|
} else {
|
|
|
|
|
span.AddEvent("encryption layer not processed", oteltrace.WithAttributes(
|
|
|
|
|
otelattribute.String("error", err.Error()),
|
|
|
|
|
))
|
|
|
|
|
logger.Debug("failed to process encryption layer", zap.Error(err))
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-06 14:01:49 -04:00
|
|
|
// A broken SDS layer yields a payload that silently fails to decode further up,
|
|
|
|
|
// so fail the whole envelope instead: it must be retried, not marked processed.
|
2025-12-22 19:57:36 +00:00
|
|
|
err = r.processSDSLayer(responseMessage)
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Error("failed to unwrap payload for SDS", zap.Error(err))
|
2026-08-06 14:01:49 -04:00
|
|
|
return nil, err
|
2025-12-22 19:57:36 +00:00
|
|
|
}
|
|
|
|
|
|
2025-11-18 19:59:37 +01:00
|
|
|
messages, ackedMessageIDs, err := r.processReliabilityLayer(responseMessage, logger)
|
|
|
|
|
if err == nil {
|
2025-12-08 21:31:32 +01:00
|
|
|
span.AddEvent("reliability layer processed")
|
2025-11-18 19:59:37 +01:00
|
|
|
response.messages = messages
|
|
|
|
|
response.ackedMessageIDs = ackedMessageIDs
|
|
|
|
|
} else {
|
2025-12-08 21:31:32 +01:00
|
|
|
span.AddEvent("reliability layer not processed")
|
2025-11-18 19:59:37 +01:00
|
|
|
logger.Debug("failed to process reliability layer", zap.Error(err))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return response, nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
func (r *Processor) processQueuedHashRatchetMessages(hashRatchetInfos []*messagingtypes.HashRatchetInfo) (*processMessageResponse, error) {
|
2025-11-18 19:59:37 +01:00
|
|
|
response := &processMessageResponse{
|
2025-12-22 19:57:36 +00:00
|
|
|
messages: []*messagingtypes.Message{},
|
2025-12-18 12:24:40 +00:00
|
|
|
ackedMessageIDs: []types.HexBytes{},
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for _, hashRatchetInfo := range hashRatchetInfos {
|
|
|
|
|
messages, err := r.hashRatchetStorage.GetMessages(hashRatchetInfo.KeyID)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
var processedIds [][]byte
|
|
|
|
|
for _, message := range messages {
|
2025-12-18 12:24:40 +00:00
|
|
|
logger := r.logger.With(zap.String("hash", types.EncodeHex(message.Hash)))
|
2025-11-18 19:59:37 +01:00
|
|
|
logger.Debug("processing queued hash ratchet message")
|
|
|
|
|
|
|
|
|
|
r, err := r.processMessage(message)
|
|
|
|
|
if err != nil {
|
|
|
|
|
continue
|
|
|
|
|
}
|
2026-08-13 15:49:36 +04:00
|
|
|
if r == nil {
|
|
|
|
|
// Not processed: the message is still incomplete (segmentation) or was
|
|
|
|
|
// re-queued waiting for its hash ratchet key — leave it in the queue.
|
|
|
|
|
continue
|
|
|
|
|
}
|
2025-11-18 19:59:37 +01:00
|
|
|
|
|
|
|
|
processedIds = append(processedIds, message.Hash)
|
|
|
|
|
|
|
|
|
|
response.messages = append(response.messages, r.messages...)
|
|
|
|
|
response.ackedMessageIDs = append(response.ackedMessageIDs, r.ackedMessageIDs...)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
err = r.hashRatchetStorage.DeleteMessages(processedIds)
|
|
|
|
|
if err != nil {
|
|
|
|
|
r.logger.Warn("failed to delete hash ratchet messages", zap.Error(err))
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return response, nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
func processTransportLayer(m *messagingtypes.Message, receivedMessage *messagingtypes.ReceivedMessage) error {
|
2025-11-18 19:59:37 +01:00
|
|
|
publicKey, err := crypto.UnmarshalPubkey(receivedMessage.Sig)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return errors.Wrap(err, "failed to get signature")
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
m.TransportLayer.Message = receivedMessage
|
|
|
|
|
m.TransportLayer.Hash = receivedMessage.Hash
|
|
|
|
|
m.TransportLayer.SigPubKey = publicKey
|
|
|
|
|
m.TransportLayer.Payload = receivedMessage.Payload
|
|
|
|
|
|
|
|
|
|
if receivedMessage.Dst != nil {
|
|
|
|
|
publicKey, err := crypto.UnmarshalPubkey(receivedMessage.Dst)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
m.TransportLayer.Dst = publicKey
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
func (r *Processor) processSegmentationLayer(m *messagingtypes.Message) error {
|
2025-12-08 21:31:32 +01:00
|
|
|
reconstructedPayload, transportIDs, err := r.stack.Segmentation.Reconstruct(
|
|
|
|
|
m.TransportLayer.Payload,
|
|
|
|
|
m.TransportLayer.SigPubKey,
|
|
|
|
|
m.TransportLayer.Hash)
|
2025-11-18 19:59:37 +01:00
|
|
|
|
|
|
|
|
switch err {
|
|
|
|
|
case nil:
|
|
|
|
|
m.TransportLayer.Payload = reconstructedPayload
|
2025-12-08 21:31:32 +01:00
|
|
|
m.SegmentationLayer.Segmented = true
|
|
|
|
|
m.SegmentationLayer.Completed = true
|
|
|
|
|
m.SegmentationLayer.Hashes = transportIDs
|
2025-11-18 19:59:37 +01:00
|
|
|
case segmentation.ErrIncomplete:
|
2025-12-08 21:31:32 +01:00
|
|
|
m.SegmentationLayer.Segmented = true
|
|
|
|
|
m.SegmentationLayer.Completed = false
|
2025-11-18 19:59:37 +01:00
|
|
|
err = nil
|
2026-03-18 10:49:47 -04:00
|
|
|
case segmentation.ErrAlreadyCompleted:
|
|
|
|
|
// A duplicate segment for an already reconstructed message should be ignored.
|
|
|
|
|
m.SegmentationLayer.Segmented = true
|
|
|
|
|
m.SegmentationLayer.Completed = false
|
|
|
|
|
err = nil
|
2025-11-18 19:59:37 +01:00
|
|
|
case segmentation.ErrInvalidPayload:
|
2025-12-08 21:31:32 +01:00
|
|
|
m.SegmentationLayer.Segmented = false
|
|
|
|
|
m.SegmentationLayer.Completed = false
|
2025-11-18 19:59:37 +01:00
|
|
|
err = nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-08 21:31:32 +01:00
|
|
|
return err
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
func (r *Processor) processEncryptionLayer(ctx context.Context, m *messagingtypes.Message, logger *zap.Logger) error {
|
2025-11-18 19:59:37 +01:00
|
|
|
logger = logger.Named("processEncryptionLayer")
|
|
|
|
|
|
2025-12-08 21:31:32 +01:00
|
|
|
ctx, span := r.tracer.Start(ctx, "Processor.processEncryptionLayer")
|
|
|
|
|
defer span.End()
|
|
|
|
|
|
2025-11-18 19:59:37 +01:00
|
|
|
// As we handle non-encrypted messages, we make sure that DecryptPayload
|
|
|
|
|
// is set regardless of whether this step is successful
|
|
|
|
|
m.EncryptionLayer.Payload = m.TransportLayer.Payload
|
|
|
|
|
|
|
|
|
|
// if it's an ephemeral key, we don't negotiate a topic
|
|
|
|
|
ephemeralKey := r.ephemeralKeysManager.GetPrivateKeyFor(m.TransportLayer.Dst)
|
|
|
|
|
if ephemeralKey != nil {
|
2025-12-08 21:31:32 +01:00
|
|
|
span.AddEvent("targeted ephemeral key")
|
2025-11-18 19:59:37 +01:00
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2026-04-16 13:22:54 +01:00
|
|
|
var protocolMessage encryption.ProtocolMessage
|
2025-11-18 19:59:37 +01:00
|
|
|
err := proto.Unmarshal(m.TransportLayer.Payload, &protocolMessage)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return errors.Wrap(err, "failed to unmarshal ProtocolMessage")
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
response, err := r.stack.Encryption.HandleMessage(
|
2025-12-08 21:31:32 +01:00
|
|
|
ctx,
|
2025-11-18 19:59:37 +01:00
|
|
|
r.identity,
|
|
|
|
|
m.SigPubKey(),
|
|
|
|
|
&protocolMessage,
|
|
|
|
|
m.TransportLayer.Hash,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
switch err {
|
|
|
|
|
case nil:
|
|
|
|
|
m.EncryptionLayer.Payload = response.DecryptedMessage
|
2026-04-16 13:22:54 +01:00
|
|
|
m.EncryptionLayer.Installations = adapters.FromEncryptionInstallations(response.Installations)
|
|
|
|
|
m.EncryptionLayer.HashRatchetInfo = adapters.FromEncryptionHashRatchets(response.HashRatchetInfo)
|
2025-11-18 19:59:37 +01:00
|
|
|
|
|
|
|
|
err := r.ProcessSharedSecrets(response.SharedSecrets)
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Error("failed to process shared secrets", zap.Error(err))
|
|
|
|
|
}
|
2026-04-16 13:22:54 +01:00
|
|
|
case encryption.ErrHashRatchetGroupIDNotFound:
|
2025-11-18 19:59:37 +01:00
|
|
|
if response != nil {
|
2026-04-16 13:22:54 +01:00
|
|
|
m.EncryptionLayer.HashRatchetInfo = adapters.FromEncryptionHashRatchets(response.HashRatchetInfo)
|
2025-11-18 19:59:37 +01:00
|
|
|
}
|
2026-04-16 13:22:54 +01:00
|
|
|
case encryption.ErrDeviceNotFound:
|
2025-11-18 19:59:37 +01:00
|
|
|
pubsub.Publish(r.publisher, SenderUnawareOfInstallation{
|
|
|
|
|
PublicKey: m.SigPubKey(),
|
|
|
|
|
})
|
|
|
|
|
default:
|
|
|
|
|
logger.Error("failed to decrypt message", zap.Error(err))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
func (r *Processor) processSDSLayer(msg *messagingtypes.Message) error {
|
|
|
|
|
if len(msg.EncryptionLayer.Payload) <= 0 {
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
unwrappedPayload, err := r.stack.Reliability.UnwrapPayloadFromSDS(msg.EncryptionLayer.Payload)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
msg.EncryptionLayer.Payload = unwrappedPayload
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-11-18 19:59:37 +01:00
|
|
|
func (r *Processor) ProcessSharedSecrets(secrets []*sharedsecret.Secret) error {
|
|
|
|
|
for _, secret := range secrets {
|
2025-12-22 19:57:36 +00:00
|
|
|
_, err := r.stack.Transport.ProcessNegotiatedSecret(messagingtypes.NegotiatedSecret{
|
2025-11-18 19:59:37 +01:00
|
|
|
PublicKey: secret.Identity,
|
|
|
|
|
Key: secret.Key,
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
func (r *Processor) processReliabilityLayer(m *messagingtypes.Message, logger *zap.Logger) ([]*messagingtypes.Message, []types.HexBytes, error) {
|
2025-11-18 19:59:37 +01:00
|
|
|
if !r.stack.Reliability.Started() {
|
|
|
|
|
return nil, nil, errReliabilityNotStarted
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
datasyncMessage, err := r.stack.Reliability.UnwrapAndAcknowledgeMessage(
|
|
|
|
|
m.SigPubKey(),
|
|
|
|
|
m.EncryptionLayer.Payload,
|
|
|
|
|
)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, nil, err
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-22 19:57:36 +00:00
|
|
|
var statusMessages []*messagingtypes.Message
|
2025-11-18 19:59:37 +01:00
|
|
|
for _, ds := range datasyncMessage.Messages {
|
|
|
|
|
message, err := m.Clone()
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, nil, err
|
|
|
|
|
}
|
|
|
|
|
message.EncryptionLayer.Payload = ds.Body
|
|
|
|
|
statusMessages = append(statusMessages, message)
|
|
|
|
|
}
|
|
|
|
|
|
2025-12-18 12:24:40 +00:00
|
|
|
ackedMessageIDs := make([]types.HexBytes, 0, len(datasyncMessage.Acks))
|
2025-11-18 19:59:37 +01:00
|
|
|
for _, ack := range datasyncMessage.Acks {
|
|
|
|
|
messageID, err := r.messageConfirmationStorage.MarkAsConfirmed(ack, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Info("got datasync acknowledge for message we don't have in db", zap.String("ack", hex.EncodeToString(ack)))
|
|
|
|
|
continue
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
r.stack.Transport.ConfirmMessageDelivered(messageID.String())
|
|
|
|
|
|
|
|
|
|
ackedMessageIDs = append(ackedMessageIDs, messageID)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return statusMessages, ackedMessageIDs, nil
|
|
|
|
|
}
|