Skip to content
Merged
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
13 changes: 13 additions & 0 deletions pkg/connector/backfill.go
Original file line number Diff line number Diff line change
Expand Up @@ -616,6 +616,7 @@ func ensureRESTFallbackCursor(resp *bridgev2.FetchMessagesResponse) *bridgev2.Fe

const (
xchatSkipReasonEmptyContents = "empty_contents"
xchatSkipReasonUserMetadata = "user_metadata"
xchatSkipReasonKeyMissing = "key_missing"
xchatSkipReasonDecryptFailed = "decrypt_failed"
xchatSkipReasonParseFailed = "parse_failed"
Expand Down Expand Up @@ -644,8 +645,12 @@ func (tc *TwitterClient) decodeXChatMessageCreateForBackfill(ctx context.Context
mce := evt.Detail.MessageCreateEvent
keyVersion := ptr.Val(mce.ConversationKeyVersion)
contentsBytes := mce.Contents
isUserMetadata := isEncryptedXChatBackfillMessage(mce) && mce.IsUserMetadata()

if len(contentsBytes) == 0 {
if isUserMetadata {
return nil, time.Time{}, 0, xchatSkipReasonDecryptFailed
}
return nil, time.Time{}, 0, xchatSkipReasonEmptyContents
}

Expand All @@ -664,6 +669,14 @@ func (tc *TwitterClient) decodeXChatMessageCreateForBackfill(ctx context.Context
Str("key_version", keyVersion).
Logger()

if isUserMetadata {
err = crypto.DecryptUserMetadataBytes(contentsBytes, convKey.Key)
if err != nil {
return nil, time.Time{}, 0, xchatSkipReasonDecryptFailed
}
return nil, time.Time{}, 0, xchatSkipReasonUserMetadata
}

decrypted, err := crypto.DecryptMessageEntryContentsBytesDebug(contentsBytes, convKey.Key, &debugLog)
if err != nil {
return nil, time.Time{}, 0, xchatSkipReasonDecryptFailed
Expand Down
14 changes: 14 additions & 0 deletions pkg/twittermeow/crypto/message.go
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,20 @@ func ParseMessageEntryContentsBytes(data []byte) (*payload.MessageEntryContents,
return decodeMessageEntryHolder(data, nil)
}

func DecryptUserMetadataBytes(ciphertext, conversationKey []byte) (err error) {
defer func() {
if recovered := recover(); recovered != nil {
err = fmt.Errorf("thrift decode panic: %v", recovered)
}
}()

plaintext, err := SecretboxDecrypt(ciphertext, conversationKey)
if err != nil {
return fmt.Errorf("secretbox decrypt: %w", err)
}
return payload.Decode(plaintext, new(payload.UserMetadata))
}

func decodeMessageEntryHolder(data []byte, log *zerolog.Logger) (_ *payload.MessageEntryContents, err error) {
defer func() {
if r := recover(); r != nil {
Expand Down
38 changes: 30 additions & 8 deletions pkg/twittermeow/data/payload/thrift.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

31 changes: 31 additions & 0 deletions pkg/twittermeow/data/payload/thrift_codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -349,6 +349,37 @@ func readValue(ctx context.Context, proto thrift.TProtocol, v reflect.Value, goT
}
v.Set(slice)

case thrift.MAP:
keyType, valueType, size, err := proto.ReadMapBegin(ctx)
if err != nil {
return err
}
if goType.Kind() != reflect.Map {
return fmt.Errorf("expected map type, got %s", goType.Kind())
}
if expected := goTypeToThriftType(goType.Key(), goType.Key()); keyType != expected {
return fmt.Errorf("unexpected map key type %d, expected %d", keyType, expected)
}
if expected := goTypeToThriftType(goType.Elem(), goType.Elem()); valueType != expected {
return fmt.Errorf("unexpected map value type %d, expected %d", valueType, expected)
}
decodedMap := reflect.MakeMapWithSize(goType, size)
for i := 0; i < size; i++ {
key := reflect.New(goType.Key()).Elem()
if err := readValue(ctx, proto, key, goType.Key(), keyType); err != nil {
return err
}
value := reflect.New(goType.Elem()).Elem()
if err := readValue(ctx, proto, value, goType.Elem(), valueType); err != nil {
return err
}
decodedMap.SetMapIndex(key, value)
}
if err := proto.ReadMapEnd(ctx); err != nil {
return err
}
v.Set(decodedMap)

case thrift.STRUCT:
return readStruct(ctx, proto, v)

Expand Down
12 changes: 12 additions & 0 deletions pkg/twittermeow/xchat_processor.go
Original file line number Diff line number Diff line change
Expand Up @@ -516,8 +516,12 @@ func (p *XChatEventProcessor) processMessageCreateEvent(ctx context.Context, evt
conversationID := ptr.Val(evt.ConversationId)
contentsBytes := mce.Contents
keyVersion := ptr.Val(mce.ConversationKeyVersion)
isUserMetadata := isEncryptedMessageCreateEvent(mce) && mce.IsUserMetadata()

if len(contentsBytes) == 0 {
if isUserMetadata {
return errors.New("XChat user metadata event has no contents")
}
p.log.Debug().
Str("sequence_id", ptr.Val(evt.SequenceId)).
Str("conversation_id", conversationID).
Expand Down Expand Up @@ -568,6 +572,14 @@ func (p *XChatEventProcessor) processMessageCreateEvent(ctx context.Context, evt
Str("conversation_id", conversationID).
Str("key_version", keyVersion).
Logger()
if isUserMetadata {
err = crypto.DecryptUserMetadataBytes(contentsBytes, convKey.Key)
if err != nil {
return fmt.Errorf("decrypt XChat user metadata: %w", err)
}
return nil
}

contents, err = crypto.DecryptMessageEntryContentsBytesDebug(contentsBytes, convKey.Key, &debugLog)
if err != nil {
p.log.Warn().
Expand Down
Loading