From 55300459765633c28dd2e969637b693943cb0acc Mon Sep 17 00:00:00 2001 From: highesttt Date: Wed, 16 Sep 2026 18:45:40 -0400 Subject: [PATCH] xchat: fixed user metadata blocking catch-up --- pkg/connector/backfill.go | 13 +++++++ pkg/twittermeow/crypto/message.go | 14 ++++++++ pkg/twittermeow/data/payload/thrift.go | 38 +++++++++++++++----- pkg/twittermeow/data/payload/thrift_codec.go | 31 ++++++++++++++++ pkg/twittermeow/xchat_processor.go | 12 +++++++ 5 files changed, 100 insertions(+), 8 deletions(-) diff --git a/pkg/connector/backfill.go b/pkg/connector/backfill.go index 9a9b54e..98e2129 100644 --- a/pkg/connector/backfill.go +++ b/pkg/connector/backfill.go @@ -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" @@ -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 } @@ -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 diff --git a/pkg/twittermeow/crypto/message.go b/pkg/twittermeow/crypto/message.go index 9db3b09..225091c 100644 --- a/pkg/twittermeow/crypto/message.go +++ b/pkg/twittermeow/crypto/message.go @@ -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 { diff --git a/pkg/twittermeow/data/payload/thrift.go b/pkg/twittermeow/data/payload/thrift.go index 51c1139..ae51bec 100644 --- a/pkg/twittermeow/data/payload/thrift.go +++ b/pkg/twittermeow/data/payload/thrift.go @@ -348,14 +348,36 @@ type MessageContents struct { } type MessageCreateEvent struct { - Contents []byte `thrift:"contents,100" json:"contents,omitempty"` - ConversationKeyVersion *string `thrift:"conversation_key_version,101" json:"conversation_key_version,omitempty"` - ShouldNotify *bool `thrift:"should_notify,102" json:"should_notify,omitempty"` - TtlMsec *int64 `thrift:"ttl_msec,103" json:"ttl_msec,omitempty"` - DeliveredAtMsec *int64 `thrift:"delivered_at_msec,104" json:"delivered_at_msec,omitempty"` - IsPendingPublicKey *bool `thrift:"is_pending_public_key,105" json:"is_pending_public_key,omitempty"` - Priority *int32 `thrift:"priority,106" json:"priority,omitempty"` - AdditionalActionList []int32 `thrift:"additional_action_list,107" json:"additional_action_list,omitempty"` + Contents []byte `thrift:"contents,100" json:"contents,omitempty"` + ConversationKeyVersion *string `thrift:"conversation_key_version,101" json:"conversation_key_version,omitempty"` + ShouldNotify *bool `thrift:"should_notify,102" json:"should_notify,omitempty"` + TtlMsec *int64 `thrift:"ttl_msec,103" json:"ttl_msec,omitempty"` + DeliveredAtMsec *int64 `thrift:"delivered_at_msec,104" json:"delivered_at_msec,omitempty"` + IsPendingPublicKey *bool `thrift:"is_pending_public_key,105" json:"is_pending_public_key,omitempty"` + Priority *int32 `thrift:"priority,106" json:"priority,omitempty"` + AdditionalActionList []int32 `thrift:"additional_action_list,107" json:"additional_action_list,omitempty"` + SideEffects []*MessageCreateSideEffect `thrift:"side_effects,111" json:"side_effects,omitempty"` +} + +func (mce *MessageCreateEvent) IsUserMetadata() bool { + for _, sideEffect := range mce.SideEffects { + if sideEffect != nil && (sideEffect.PinConversations != nil || sideEffect.SetNicknames != nil || sideEffect.SetVerifiedUsers != nil) { + return true + } + } + return false +} + +type MessageCreateSideEffect struct { + PinConversations *struct{} `thrift:"pin_conversations,3" json:"pin_conversations,omitempty"` + SetNicknames *struct{} `thrift:"set_nicknames,4" json:"set_nicknames,omitempty"` + SetVerifiedUsers *struct{} `thrift:"set_verified_users,5" json:"set_verified_users,omitempty"` +} + +type UserMetadata struct { + PinnedConversationIds []string `thrift:"pinned_conversation_ids,1" json:"pinned_conversation_ids,omitempty"` + UserIdToNickname map[int64]string `thrift:"user_id_to_nickname,2" json:"user_id_to_nickname,omitempty"` + SafetyNumberVerifiedUserIds []int64 `thrift:"safety_number_verified_user_ids,3" json:"safety_number_verified_user_ids,omitempty"` } type MessageDeleteEvent struct { diff --git a/pkg/twittermeow/data/payload/thrift_codec.go b/pkg/twittermeow/data/payload/thrift_codec.go index e676f8b..06aa7de 100644 --- a/pkg/twittermeow/data/payload/thrift_codec.go +++ b/pkg/twittermeow/data/payload/thrift_codec.go @@ -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) diff --git a/pkg/twittermeow/xchat_processor.go b/pkg/twittermeow/xchat_processor.go index 039006d..df223a4 100644 --- a/pkg/twittermeow/xchat_processor.go +++ b/pkg/twittermeow/xchat_processor.go @@ -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). @@ -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().