From 73896fe796e0ec2141d8b74d013cf50fa5c1098e Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Fri, 20 Mar 2026 23:23:23 +0200 Subject: [PATCH 1/9] connector & twittermeow: use proper key version to determine latest conversation key --- pkg/connector/backfill.go | 7 +- pkg/connector/keystore.go | 19 ++- pkg/twittermeow/client.go | 1 + pkg/twittermeow/crypto/message.go | 9 +- pkg/twittermeow/data/endpoints/endpoints.go | 2 +- pkg/twittermeow/data/payload/form.go | 1 + pkg/twittermeow/data/payload/json.go | 4 +- pkg/twittermeow/messaging.go | 139 +++++++++----------- pkg/twittermeow/xchat_processor.go | 8 +- pkg/twittermeow/xchat_send.go | 106 +++++---------- 10 files changed, 133 insertions(+), 163 deletions(-) diff --git a/pkg/connector/backfill.go b/pkg/connector/backfill.go index 9022f2b..1b2f250 100644 --- a/pkg/connector/backfill.go +++ b/pkg/connector/backfill.go @@ -388,11 +388,16 @@ func (tc *TwitterClient) storeConversationKeyFromChangeEvent(ctx context.Context return err } + keyCreatedAt := methods.ParseMsecTimestamp(ptr.Val(evt.CreatedAtMsec)) + if keyCreatedAt.IsZero() { + return fmt.Errorf("missing valid XChat key timestamp for conversation %s key %s", conversationID, newKeyVersion) + } + return tc.client.GetKeyManager().PutConversationKey(ctx, &crypto.ConversationKey{ ConversationID: conversationID, KeyVersion: newKeyVersion, Key: convKeyBytes, - CreatedAt: time.Now(), + CreatedAt: keyCreatedAt, }) } diff --git a/pkg/connector/keystore.go b/pkg/connector/keystore.go index 52fbd85..1da11d0 100644 --- a/pkg/connector/keystore.go +++ b/pkg/connector/keystore.go @@ -10,6 +10,7 @@ import ( "maunium.net/go/mautrix/bridgev2" "go.mau.fi/mautrix-twitter/pkg/twittermeow/crypto" + "go.mau.fi/mautrix-twitter/pkg/twittermeow/methods" ) // userLoginKeyStore stores cryptographic keys. @@ -43,6 +44,8 @@ func (ks *userLoginKeyStore) getPortalMetadata(ctx context.Context, conversation func (ks *userLoginKeyStore) GetConversationKey(ctx context.Context, conversationID, keyVersion string) (*crypto.ConversationKey, error) { log := zerolog.Ctx(ctx) + ks.mu.RLock() + defer ks.mu.RUnlock() meta, _, err := ks.getPortalMetadata(ctx, conversationID) if err != nil { @@ -94,6 +97,9 @@ func (ks *userLoginKeyStore) PutConversationKey(ctx context.Context, key *crypto return fmt.Errorf("conversation key cannot be nil") } + ks.mu.Lock() + defer ks.mu.Unlock() + meta, portal, err := ks.getPortalMetadata(ctx, key.ConversationID) if err != nil { log.Err(err). @@ -139,6 +145,9 @@ func (ks *userLoginKeyStore) PutConversationKey(ctx context.Context, key *crypto } func (ks *userLoginKeyStore) DeleteConversationKey(ctx context.Context, conversationID, keyVersion string) error { + ks.mu.Lock() + defer ks.mu.Unlock() + meta, portal, err := ks.getPortalMetadata(ctx, conversationID) if err != nil { return err @@ -152,6 +161,8 @@ func (ks *userLoginKeyStore) DeleteConversationKey(ctx context.Context, conversa func (ks *userLoginKeyStore) GetLatestConversationKey(ctx context.Context, conversationID string) (*crypto.ConversationKey, error) { log := zerolog.Ctx(ctx) + ks.mu.RLock() + defer ks.mu.RUnlock() meta, _, err := ks.getPortalMetadata(ctx, conversationID) if err != nil { @@ -168,10 +179,9 @@ func (ks *userLoginKeyStore) GetLatestConversationKey(ctx context.Context, conve return nil, crypto.ErrKeyNotFound } - // Find key with latest CreatedAt var latest *ConversationKeyData for _, data := range meta.ConversationKeys { - if latest == nil || data.CreatedAt.After(latest.CreatedAt) { + if latest == nil || methods.CompareSnowflake(data.KeyVersion, latest.KeyVersion) > 0 { latest = data } } @@ -250,6 +260,9 @@ func (ks *userLoginKeyStore) PutPublicKey(_ context.Context, _ *crypto.PublicKey } func (ks *userLoginKeyStore) GetConversationToken(ctx context.Context, conversationID string) (string, error) { + ks.mu.RLock() + defer ks.mu.RUnlock() + meta, _, err := ks.getPortalMetadata(ctx, conversationID) if err != nil { return "", err @@ -263,6 +276,8 @@ func (ks *userLoginKeyStore) GetConversationToken(ctx context.Context, conversat func (ks *userLoginKeyStore) PutConversationToken(ctx context.Context, conversationID, token string) error { log := zerolog.Ctx(ctx) + ks.mu.Lock() + defer ks.mu.Unlock() meta, portal, err := ks.getPortalMetadata(ctx, conversationID) if err != nil { diff --git a/pkg/twittermeow/client.go b/pkg/twittermeow/client.go index 3645808..139889d 100644 --- a/pkg/twittermeow/client.go +++ b/pkg/twittermeow/client.go @@ -216,6 +216,7 @@ func (c *Client) LoadMessagesPage(ctx context.Context) (*response.AccountSetting IncludeNSFWAdminFlag: true, IncludeRankedTimeline: true, IncludeAltTextCompose: true, + IncludeExtDMAVCallSettings: true, Ext: "ssoConnections", IncludeCountryCode: true, IncludeExtDMNSFWMediaFilter: true, diff --git a/pkg/twittermeow/crypto/message.go b/pkg/twittermeow/crypto/message.go index 1686ddd..d6ef594 100644 --- a/pkg/twittermeow/crypto/message.go +++ b/pkg/twittermeow/crypto/message.go @@ -6,7 +6,6 @@ import ( "encoding/base64" "encoding/hex" "encoding/json" - "errors" "fmt" "github.com/rs/zerolog" @@ -482,13 +481,9 @@ func (b *MessageBuilder) Build(ctx context.Context) (*payload.MessageEvent, erro } else if b.km != nil { key, err := b.km.GetLatestConversationKey(ctx, b.conversationID) if err != nil { - if errors.Is(err, ErrKeyNotFound) { - unencrypted = true - } else { - return nil, fmt.Errorf("get conversation key: %w", err) - } + return nil, fmt.Errorf("get conversation key: %w", err) } else if key == nil || len(key.Key) == 0 { - unencrypted = true + return nil, ErrKeyNotFound } else { convKey = key.Key keyVersion = key.KeyVersion diff --git a/pkg/twittermeow/data/endpoints/endpoints.go b/pkg/twittermeow/data/endpoints/endpoints.go index c05b07c..6f6ee37 100644 --- a/pkg/twittermeow/data/endpoints/endpoints.go +++ b/pkg/twittermeow/data/endpoints/endpoints.go @@ -16,7 +16,7 @@ const ( API_BASE_HOST = "api.x.com" API_BASE_URL = "https://" + API_BASE_HOST - ACCOUNT_SETTINGS_URL = API_BASE_URL + "/1.1/account/settings.json" + ACCOUNT_SETTINGS_URL = BASE_URL + "/i/api/1.1/account/settings.json" INBOX_INITIAL_STATE_URL = BASE_URL + "/i/api/1.1/dm/inbox_initial_state.json" DM_USER_UPDATES_URL = BASE_URL + "/i/api/1.1/dm/user_updates.json" TRUSTED_INBOX_TIMELINE_URL = BASE_URL + "/i/api/1.1/dm/inbox_timeline/trusted.json" diff --git a/pkg/twittermeow/data/payload/form.go b/pkg/twittermeow/data/payload/form.go index e5731b3..fa1bcce 100644 --- a/pkg/twittermeow/data/payload/form.go +++ b/pkg/twittermeow/data/payload/form.go @@ -33,6 +33,7 @@ type AccountSettingsQuery struct { IncludeNSFWAdminFlag bool `url:"include_nsfw_admin_flag"` IncludeRankedTimeline bool `url:"include_ranked_timeline"` IncludeAltTextCompose bool `url:"include_alt_text_compose"` + IncludeExtDMAVCallSettings bool `url:"include_ext_dm_av_call_settings"` Ext string `url:"ext"` IncludeCountryCode bool `url:"include_country_code"` IncludeExtDMNSFWMediaFilter bool `url:"include_ext_dm_nsfw_media_filter"` diff --git a/pkg/twittermeow/data/payload/json.go b/pkg/twittermeow/data/payload/json.go index 5ab138d..0052aee 100644 --- a/pkg/twittermeow/data/payload/json.go +++ b/pkg/twittermeow/data/payload/json.go @@ -305,7 +305,7 @@ func NewSendMessageMutationPayload(vars SendMessageMutationVariables) *SendMessa } type GenerateXChatTokenMutationPayload struct { - Variables map[string]any `json:"variables"` + Variables string `json:"variables"` Features map[string]any `json:"features,omitempty"` Extensions struct { PersistedQuery struct { @@ -316,7 +316,7 @@ type GenerateXChatTokenMutationPayload struct { } func (p *GenerateXChatTokenMutationPayload) Default() *GenerateXChatTokenMutationPayload { - p.Variables = map[string]any{} + p.Variables = "{}" p.Extensions.PersistedQuery.Version = 1 return p } diff --git a/pkg/twittermeow/messaging.go b/pkg/twittermeow/messaging.go index 26726fd..878e4cd 100644 --- a/pkg/twittermeow/messaging.go +++ b/pkg/twittermeow/messaging.go @@ -414,10 +414,65 @@ type SendEncryptedEditOpts struct { Entities []*payload.RichTextEntity } +func (c *Client) sendMessageMutationOnce(ctx context.Context, pl *payload.SendMessageMutationPayload) (*response.SendMessageMutationResponse, error) { + jsonBody, err := json.Marshal(pl) + if err != nil { + return nil, err + } + + c.Logger.Debug(). + RawJSON("payload", jsonBody). + Msg("SendMessageMutation payload") + + _, respBody, err := c.makeAPIRequest(ctx, apiRequestOpts{ + URL: endpoints.SEND_MESSAGE_MUTATION_URL, + Method: http.MethodPost, + WithClientUUID: true, + Origin: endpoints.BASE_URL, + ContentType: types.ContentTypeJSON, + Body: jsonBody, + }) + if err != nil { + return nil, err + } + + var resp response.SendMessageMutationResponse + if err := json.Unmarshal(respBody, &resp); err != nil { + return nil, err + } + if len(resp.Errors) > 0 { + return nil, fmt.Errorf("send message mutation error: %s", resp.Errors[0].Message) + } + if resp.Data.XChatSendCreateMessageEvent.EncodedMessageEvent == "" { + return nil, fmt.Errorf("send message mutation returned no encoded message event") + } + return &resp, nil +} + +func (c *Client) sendMessageMutation(ctx context.Context, pl *payload.SendMessageMutationPayload) (*response.SendMessageMutationResponse, error) { + resp, err := c.sendMessageMutationOnce(ctx, pl) + if err == nil { + return resp, nil + } + + if err := c.refreshConversationToken(ctx, pl.Variables.ConversationID); err != nil { + return nil, err + } + + nextToken, tokenErr := c.keyManager.GetConversationToken(ctx, pl.Variables.ConversationID) + if tokenErr != nil || nextToken == pl.Variables.ConversationToken { + return nil, err + } + + retryPayload := *pl + retryPayload.Variables.ConversationToken = nextToken + return c.sendMessageMutationOnce(ctx, &retryPayload) +} + // SendEncryptedReaction sends a reaction add/remove via the XChat protocol. // targetMessageSequenceID must be the XChat message sequence ID of the message being reacted to. func (c *Client) SendEncryptedReaction(ctx context.Context, conversationID, targetMessageSequenceID, emoji string, remove bool) (*response.SendMessageMutationResponse, error) { - token, err := c.keyManager.GetConversationToken(ctx, conversationID) + token, err := c.ensureConversationToken(ctx, conversationID) if err != nil { return nil, fmt.Errorf("get conversation token: %w", err) } @@ -452,29 +507,7 @@ func (c *Client) SendEncryptedReaction(ctx context.Context, conversationID, targ EncodedMessageEventSignature: sigPtr, }) - jsonBody, err := json.Marshal(pl) - if err != nil { - return nil, err - } - - c.Logger.Debug(). - RawJSON("payload", jsonBody). - Msg("SendMessageMutation reaction payload") - - _, respBody, err := c.makeAPIRequest(ctx, apiRequestOpts{ - URL: endpoints.SEND_MESSAGE_MUTATION_URL, - Method: http.MethodPost, - WithClientUUID: true, - Origin: endpoints.BASE_URL, - ContentType: types.ContentTypeJSON, - Body: jsonBody, - }) - if err != nil { - return nil, err - } - - var resp response.SendMessageMutationResponse - return &resp, json.Unmarshal(respBody, &resp) + return c.sendMessageMutation(ctx, pl) } // SendEncryptedEdit sends a message edit via the XChat protocol. @@ -487,7 +520,7 @@ func (c *Client) SendEncryptedEdit(ctx context.Context, opts SendEncryptedEditOp return nil, fmt.Errorf("target message sequence ID is required") } - token, err := c.keyManager.GetConversationToken(ctx, opts.ConversationID) + token, err := c.ensureConversationToken(ctx, opts.ConversationID) if err != nil { return nil, fmt.Errorf("get conversation token: %w", err) } @@ -520,35 +553,13 @@ func (c *Client) SendEncryptedEdit(ctx context.Context, opts SendEncryptedEditOp EncodedMessageEventSignature: sigPtr, }) - jsonBody, err := json.Marshal(pl) - if err != nil { - return nil, err - } - - c.Logger.Debug(). - RawJSON("payload", jsonBody). - Msg("SendMessageMutation edit payload") - - _, respBody, err := c.makeAPIRequest(ctx, apiRequestOpts{ - URL: endpoints.SEND_MESSAGE_MUTATION_URL, - Method: http.MethodPost, - WithClientUUID: true, - Origin: endpoints.BASE_URL, - ContentType: types.ContentTypeJSON, - Body: jsonBody, - }) - if err != nil { - return nil, err - } - - var resp response.SendMessageMutationResponse - return &resp, json.Unmarshal(respBody, &resp) + return c.sendMessageMutation(ctx, pl) } // SendEncryptedMessage sends an encrypted message via the XChat protocol. func (c *Client) SendEncryptedMessage(ctx context.Context, opts SendEncryptedMessageOpts) (*response.SendMessageMutationResponse, error) { // Get the server-provided conversation token for this conversation - token, err := c.keyManager.GetConversationToken(ctx, opts.ConversationID) + token, err := c.ensureConversationToken(ctx, opts.ConversationID) if err != nil { return nil, fmt.Errorf("get conversation token: %w", err) } @@ -605,29 +616,7 @@ func (c *Client) SendEncryptedMessage(ctx context.Context, opts SendEncryptedMes EncodedMessageEventSignature: sigPtr, }) - jsonBody, err := json.Marshal(pl) - if err != nil { - return nil, err - } - - c.Logger.Debug(). - RawJSON("payload", jsonBody). - Msg("SendMessageMutation payload") - - _, respBody, err := c.makeAPIRequest(ctx, apiRequestOpts{ - URL: endpoints.SEND_MESSAGE_MUTATION_URL, - Method: http.MethodPost, - WithClientUUID: true, - Origin: endpoints.BASE_URL, - ContentType: types.ContentTypeJSON, - Body: jsonBody, - }) - if err != nil { - return nil, err - } - - var resp response.SendMessageMutationResponse - return &resp, json.Unmarshal(respBody, &resp) + return c.sendMessageMutation(ctx, pl) } func (c *Client) GetInitialXChatPage(ctx context.Context, variables *payload.GetInitialXChatPageQueryVariables) (*response.GetInitialXChatPageQueryResponse, error) { @@ -671,7 +660,7 @@ func (c *Client) DeleteXChatMessage(ctx context.Context, opts DeleteXChatMessage return fmt.Errorf("sender ID is required") } - token, err := c.keyManager.GetConversationToken(ctx, opts.ConversationID) + token, err := c.ensureConversationToken(ctx, opts.ConversationID) if err != nil { return fmt.Errorf("get conversation token: %w", err) } @@ -798,7 +787,7 @@ func (c *Client) MuteConversation(ctx context.Context, conversationID string) er return fmt.Errorf("sender ID is required") } - token, err := c.keyManager.GetConversationToken(ctx, conversationID) + token, err := c.ensureConversationToken(ctx, conversationID) if err != nil { return fmt.Errorf("get conversation token: %w", err) } @@ -902,7 +891,7 @@ func (c *Client) UnmuteConversation(ctx context.Context, conversationID string) return fmt.Errorf("sender ID is required") } - token, err := c.keyManager.GetConversationToken(ctx, conversationID) + token, err := c.ensureConversationToken(ctx, conversationID) if err != nil { return fmt.Errorf("get conversation token: %w", err) } diff --git a/pkg/twittermeow/xchat_processor.go b/pkg/twittermeow/xchat_processor.go index 7b622ca..5014370 100644 --- a/pkg/twittermeow/xchat_processor.go +++ b/pkg/twittermeow/xchat_processor.go @@ -17,6 +17,7 @@ import ( "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/payload" "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/response" "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/types" + "go.mau.fi/mautrix-twitter/pkg/twittermeow/methods" ) // XChatEventHandler processes XChat events. @@ -419,12 +420,17 @@ func (p *XChatEventProcessor) processConversationKeyChange(ctx context.Context, return err } + keyCreatedAt := methods.ParseMsecTimestamp(ptr.Val(evt.CreatedAtMsec)) + if keyCreatedAt.IsZero() { + return fmt.Errorf("missing valid XChat key timestamp for conversation %s key %s", conversationID, newKeyVersion) + } + // Store the new key if err := p.client.keyManager.PutConversationKey(ctx, &crypto.ConversationKey{ ConversationID: conversationID, KeyVersion: newKeyVersion, Key: convKeyBytes, - CreatedAt: time.Now(), + CreatedAt: keyCreatedAt, }); err != nil { p.log.Err(err). Str("conversation_id", conversationID). diff --git a/pkg/twittermeow/xchat_send.go b/pkg/twittermeow/xchat_send.go index 08827a6..4506ff6 100644 --- a/pkg/twittermeow/xchat_send.go +++ b/pkg/twittermeow/xchat_send.go @@ -3,10 +3,8 @@ package twittermeow import ( "context" "encoding/base64" - "encoding/json" "errors" "fmt" - "net/http" "strconv" "time" @@ -14,10 +12,9 @@ import ( "go.mau.fi/util/ptr" "go.mau.fi/mautrix-twitter/pkg/twittermeow/crypto" - "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/endpoints" "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/payload" "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/response" - "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/types" + "go.mau.fi/mautrix-twitter/pkg/twittermeow/methods" ) func (c *Client) SendXChatReadReceipt(ctx context.Context, conversationID, lastReadEventID string, readAt time.Time) error { @@ -87,7 +84,7 @@ func (c *Client) SendXChatReadReceipt(ctx context.Context, conversationID, lastR func (c *Client) ensureConversationToken(ctx context.Context, conversationID string) (string, error) { token, err := c.keyManager.GetConversationToken(ctx, conversationID) - if err == nil && token != "" { + if err == nil { return token, nil } if err != nil && !errors.Is(err, crypto.ErrKeyNotFound) { @@ -99,13 +96,9 @@ func (c *Client) ensureConversationToken(ctx context.Context, conversationID str } token, err = c.keyManager.GetConversationToken(ctx, conversationID) - if err != nil || token == "" { - if err == nil { - err = crypto.ErrKeyNotFound - } + if err != nil { return "", fmt.Errorf("get conversation token: %w", err) } - return token, nil } @@ -117,21 +110,6 @@ func (c *Client) refreshConversationToken(ctx context.Context, conversationID st } item := resp.Data.GetInboxPageConversationData.Data - if err := c.storeConversationTokenFromInboxItem(ctx, conversationID, &item); err != nil { - if errors.Is(err, crypto.ErrKeyNotFound) { - return err - } - return fmt.Errorf("store conversation token: %w", err) - } - - return nil -} - -func (c *Client) storeConversationTokenFromInboxItem(ctx context.Context, conversationID string, item *response.XChatInboxItem) error { - if item == nil { - return crypto.ErrKeyNotFound - } - encoded := make([]string, 0, len(item.LatestMessageEvents)+len(item.EncodedMessageEvents)+len(item.LatestConversationKeyChangeEvents)+2) encoded = append(encoded, item.LatestMessageEvents...) encoded = append(encoded, item.EncodedMessageEvents...) @@ -142,51 +120,38 @@ func (c *Client) storeConversationTokenFromInboxItem(ctx context.Context, conver if item.ConversationDetail.LatestGroupTitleChangeMessageEvent != "" { encoded = append(encoded, item.ConversationDetail.LatestGroupTitleChangeMessageEvent) } + for _, readEvt := range item.LatestReadEventsPerParticipant { + encoded = append(encoded, readEvt.LatestMarkConversationReadEvent) + } for _, encodedEvt := range encoded { - ok, err := c.storeConversationTokenFromEncodedEvent(ctx, conversationID, encodedEvt) - if err != nil { - return err - } - if ok { + err := c.putConversationTokenFromEncodedEvent(ctx, conversationID, encodedEvt) + if err == nil { return nil } - } - - for _, readEvt := range item.LatestReadEventsPerParticipant { - ok, err := c.storeConversationTokenFromEncodedEvent(ctx, conversationID, readEvt.LatestMarkConversationReadEvent) - if err != nil { + if !errors.Is(err, crypto.ErrKeyNotFound) { return err } - if ok { - return nil - } } return crypto.ErrKeyNotFound } -func (c *Client) storeConversationTokenFromEncodedEvent(ctx context.Context, fallbackConversationID, encoded string) (bool, error) { +func (c *Client) putConversationTokenFromEncodedEvent(ctx context.Context, conversationID, encoded string) error { if encoded == "" { - return false, nil + return crypto.ErrKeyNotFound } evt, err := DecodeMessageEvent(encoded) if err != nil { - return false, nil + return crypto.ErrKeyNotFound } if evt == nil || evt.ConversationToken == nil || *evt.ConversationToken == "" { - return false, nil + return crypto.ErrKeyNotFound } - - conversationID := fallbackConversationID if evt.ConversationId != nil && *evt.ConversationId != "" { conversationID = *evt.ConversationId } - - if err := c.keyManager.PutConversationToken(ctx, conversationID, *evt.ConversationToken); err != nil { - return false, err - } - return true, nil + return c.keyManager.PutConversationToken(ctx, conversationID, *evt.ConversationToken) } // getSelfConversationID returns the user's self-conversation ID (user_id:user_id format). @@ -233,7 +198,8 @@ func (c *Client) SendXChatPinConversation(ctx context.Context, targetConversatio EncodedMessageEventSignature: sigPtr, }) - return c.sendMessageMutation(ctx, pl) + _, err = c.sendMessageMutation(ctx, pl) + return err } // SendXChatUnpinConversation unpins a conversation via XChat. @@ -274,28 +240,7 @@ func (c *Client) SendXChatUnpinConversation(ctx context.Context, targetConversat EncodedMessageEventSignature: sigPtr, }) - return c.sendMessageMutation(ctx, pl) -} - -// sendMessageMutation sends a SendMessageMutation request. -func (c *Client) sendMessageMutation(ctx context.Context, pl *payload.SendMessageMutationPayload) error { - jsonBody, err := json.Marshal(pl) - if err != nil { - return err - } - - c.Logger.Debug(). - RawJSON("payload", jsonBody). - Msg("SendMessageMutation payload") - - _, _, err = c.makeAPIRequest(ctx, apiRequestOpts{ - URL: endpoints.SEND_MESSAGE_MUTATION_URL, - Method: http.MethodPost, - WithClientUUID: true, - Origin: endpoints.BASE_URL, - ContentType: types.ContentTypeJSON, - Body: jsonBody, - }) + _, err = c.sendMessageMutation(ctx, pl) return err } @@ -402,12 +347,25 @@ func (c *Client) processKeyChangeEventsFromItem(ctx context.Context, conversatio if err != nil { continue } - c.keyManager.PutConversationKey(ctx, &crypto.ConversationKey{ + keyCreatedAt := methods.ParseMsecTimestamp(ptr.Val(evt.CreatedAtMsec)) + if keyCreatedAt.IsZero() { + c.Logger.Warn(). + Str("conversation_id", conversationID). + Str("key_version", ptr.Val(ckce.ConversationKeyVersion)). + Str("created_at_msec", ptr.Val(evt.CreatedAtMsec)). + Msg("Skipping conversation key update without valid XChat timestamp") + continue + } + + err = c.keyManager.PutConversationKey(ctx, &crypto.ConversationKey{ ConversationID: conversationID, KeyVersion: ptr.Val(ckce.ConversationKeyVersion), Key: convKeyBytes, - CreatedAt: time.Now(), + CreatedAt: keyCreatedAt, }) + if err != nil { + return err + } } return nil } From c5f12fa388451bc63758d410d7206e3a9324157a Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Fri, 20 Mar 2026 23:29:27 +0200 Subject: [PATCH 2/9] twittermeow: update sendMessageMutation to use ensureConversationToken function --- pkg/twittermeow/messaging.go | 32 ++++++++++---------------------- 1 file changed, 10 insertions(+), 22 deletions(-) diff --git a/pkg/twittermeow/messaging.go b/pkg/twittermeow/messaging.go index 878e4cd..82bc699 100644 --- a/pkg/twittermeow/messaging.go +++ b/pkg/twittermeow/messaging.go @@ -414,8 +414,16 @@ type SendEncryptedEditOpts struct { Entities []*payload.RichTextEntity } -func (c *Client) sendMessageMutationOnce(ctx context.Context, pl *payload.SendMessageMutationPayload) (*response.SendMessageMutationResponse, error) { - jsonBody, err := json.Marshal(pl) +func (c *Client) sendMessageMutation(ctx context.Context, pl *payload.SendMessageMutationPayload) (*response.SendMessageMutationResponse, error) { + token, err := c.ensureConversationToken(ctx, pl.Variables.ConversationID) + if err != nil { + return nil, fmt.Errorf("get conversation token: %w", err) + } + + requestPayload := *pl + requestPayload.Variables.ConversationToken = token + + jsonBody, err := json.Marshal(&requestPayload) if err != nil { return nil, err } @@ -449,26 +457,6 @@ func (c *Client) sendMessageMutationOnce(ctx context.Context, pl *payload.SendMe return &resp, nil } -func (c *Client) sendMessageMutation(ctx context.Context, pl *payload.SendMessageMutationPayload) (*response.SendMessageMutationResponse, error) { - resp, err := c.sendMessageMutationOnce(ctx, pl) - if err == nil { - return resp, nil - } - - if err := c.refreshConversationToken(ctx, pl.Variables.ConversationID); err != nil { - return nil, err - } - - nextToken, tokenErr := c.keyManager.GetConversationToken(ctx, pl.Variables.ConversationID) - if tokenErr != nil || nextToken == pl.Variables.ConversationToken { - return nil, err - } - - retryPayload := *pl - retryPayload.Variables.ConversationToken = nextToken - return c.sendMessageMutationOnce(ctx, &retryPayload) -} - // SendEncryptedReaction sends a reaction add/remove via the XChat protocol. // targetMessageSequenceID must be the XChat message sequence ID of the message being reacted to. func (c *Client) SendEncryptedReaction(ctx context.Context, conversationID, targetMessageSequenceID, emoji string, remove bool) (*response.SendMessageMutationResponse, error) { From 83734da3d7bcd5491df57e258e6c2b5abf237d4d Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Sat, 21 Mar 2026 02:05:21 +0200 Subject: [PATCH 3/9] connector & twittermeow: replace settings.json api with proper user lookup and remote profile filling --- pkg/connector/login.go | 37 ++++++++++++++++++++------- pkg/twittermeow/account.go | 52 ++++++++++++++++++++++++++++++++++++++ pkg/twittermeow/client.go | 37 ++++++++++----------------- 3 files changed, 94 insertions(+), 32 deletions(-) diff --git a/pkg/connector/login.go b/pkg/connector/login.go index 0900d95..d91164a 100644 --- a/pkg/connector/login.go +++ b/pkg/connector/login.go @@ -35,6 +35,7 @@ import ( "go.mau.fi/mautrix-twitter/pkg/twittermeow/crypto" "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/payload" "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/response" + "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/types" ) type TwitterLogin struct { @@ -47,8 +48,8 @@ type TwitterLogin struct { isMigration bool // True if upgrading from main branch (had cookies but no encryption keys) needsPINSetup bool - client *twittermeow.Client - settings *response.AccountSettingsResponse + client *twittermeow.Client + profile twittermeow.CurrentUserProfile } var ( @@ -209,12 +210,12 @@ func (t *TwitterLogin) SubmitCookies(ctx context.Context, cookies map[string]str client := twittermeow.NewClient(cookieStruct, nil, t.User.Log.With().Str("component", "login_twitter_client").Logger()) - settings, err := client.LoadMessagesPage(ctx) + profile, err := client.LoadMessagesPage(ctx) if err != nil { return nil, fmt.Errorf("failed to load messages page after submitting cookies: %w", err) } t.client = client - t.settings = settings + t.profile = profile t.persistClientCookiesAndUserID() t.refreshPINSetupState(ctx, "Failed to determine PIN setup state, using recovery prompt") @@ -244,11 +245,11 @@ func (t *TwitterLogin) ensureClientForPIN(ctx context.Context) error { } cookieStruct := twitCookies.NewCookiesFromString(t.Cookies) t.client = twittermeow.NewClient(cookieStruct, nil, t.User.Log.With().Str("component", "login_twitter_client").Logger()) - settings, err := t.client.LoadMessagesPage(ctx) + profile, err := t.client.LoadMessagesPage(ctx) if err != nil { return fmt.Errorf("failed to load messages page: %w", err) } - t.settings = settings + t.profile = profile return nil } @@ -528,13 +529,15 @@ func (t *TwitterLogin) SubmitUserInput(ctx context.Context, input map[string]str t.User.Log.Info().Msg("Migration: flagged for full encrypted room backfill") } - remoteProfile := &status.RemoteProfile{ - Username: t.settings.ScreenName, - } currentUserID := strings.TrimSpace(t.client.GetCurrentUserID()) if currentUserID == "" { return nil, ErrMissingUserID } + + remoteProfile := &status.RemoteProfile{ + Username: strings.TrimSpace(t.profile.ScreenName), + Name: strings.TrimSpace(t.profile.Name), + } id := MakeUserLoginID(currentUserID) ul, err := t.User.NewLogin( ctx, @@ -561,6 +564,22 @@ func (t *TwitterLogin) SubmitUserInput(ctx context.Context, input map[string]str return nil, err } + if profile := t.profile; profile.AvatarURL != "" { + updatedProfile := ul.Client.(*TwitterClient).makeXChatRemoteProfile(ctx, &types.User{ + IDStr: currentUserID, + ScreenName: profile.ScreenName, + Name: profile.Name, + ProfileImageURLHTTPS: profile.AvatarURL, + }) + if ul.UserLogin.RemoteName != updatedProfile.Username || ul.UserLogin.RemoteProfile != *updatedProfile { + ul.UserLogin.RemoteName = updatedProfile.Username + ul.UserLogin.RemoteProfile = *updatedProfile + if err := ul.Save(ctx); err != nil { + t.User.Log.Warn().Err(err).Msg("Failed to save login profile after syncing avatar") + } + } + } + go func(ctx context.Context, client *TwitterClient) { client.DoConnect(ctx) }(context.WithoutCancel(ctx), ul.Client.(*TwitterClient)) diff --git a/pkg/twittermeow/account.go b/pkg/twittermeow/account.go index 85ae938..34f5601 100644 --- a/pkg/twittermeow/account.go +++ b/pkg/twittermeow/account.go @@ -15,6 +15,13 @@ import ( "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/types" ) +type CurrentUserProfile struct { + ID string + ScreenName string + Name string + AvatarURL string +} + func (c *Client) Login(ctx context.Context) error { err := c.loadPage(ctx, endpoints.BASE_LOGIN_URL) if err != nil { @@ -42,6 +49,51 @@ func (c *Client) GetAccountSettings(ctx context.Context, params payload.AccountS return &data, json.Unmarshal(respBody, &data) } +func (c *Client) GetCurrentUserProfile(ctx context.Context) (CurrentUserProfile, error) { + currentUserID := strings.TrimSpace(c.GetCurrentUserID()) + if currentUserID == "" { + return CurrentUserProfile{}, fmt.Errorf("current user ID is empty") + } + + resp, err := c.GetUsersByIdsForXChat(ctx, payload.NewGetUsersByIdsForXChatVariables([]string{currentUserID})) + if err != nil { + return CurrentUserProfile{}, err + } + if len(resp.Errors) > 0 && resp.Errors[0].Message != "" { + return CurrentUserProfile{}, fmt.Errorf("GetUsersByIdsForXChat error: %s", resp.Errors[0].Message) + } + if len(resp.Data.GetMemberResults.Results) != 1 { + return CurrentUserProfile{}, fmt.Errorf("expected 1 user result for %s, got %d", currentUserID, len(resp.Data.GetMemberResults.Results)) + } + + result := resp.Data.GetMemberResults.Results[0] + if result.MemberResults == nil || result.MemberResults.Result == nil || result.MemberResults.Result.Core == nil { + return CurrentUserProfile{}, fmt.Errorf("GetUsersByIdsForXChat returned no user for %s", currentUserID) + } + + resultUserID := result.MemberResults.RestID + if resultUserID == "" { + resultUserID = result.MemberResults.Result.RestID + } + if resultUserID != "" && resultUserID != currentUserID { + return CurrentUserProfile{}, fmt.Errorf("GetUsersByIdsForXChat returned user %s for %s", resultUserID, currentUserID) + } + + profile := CurrentUserProfile{ + ID: currentUserID, + ScreenName: strings.TrimSpace(result.MemberResults.Result.Core.ScreenName), + Name: strings.TrimSpace(result.MemberResults.Result.Core.Name), + } + if resultUserID != "" { + profile.ID = resultUserID + } + if result.MemberResults.Result.Avatar != nil { + profile.AvatarURL = strings.TrimSpace(result.MemberResults.Result.Avatar.ImageURL) + } + + return profile, nil +} + func (c *Client) GetDMPermissions(ctx context.Context, params payload.GetDMPermissionsQuery) (*response.GetDMPermissionsResponse, error) { encodedQuery, err := params.Encode() if err != nil { diff --git a/pkg/twittermeow/client.go b/pkg/twittermeow/client.go index 139889d..f42dfb0 100644 --- a/pkg/twittermeow/client.go +++ b/pkg/twittermeow/client.go @@ -203,40 +203,31 @@ func (c *Client) Logout(ctx context.Context) error { return c.loadPage(ctx, endpoints.BASE_LOGOUT_URL) } -func (c *Client) LoadMessagesPage(ctx context.Context) (*response.AccountSettingsResponse, error) { +func (c *Client) LoadMessagesPage(ctx context.Context) (CurrentUserProfile, error) { err := c.loadPage(ctx, endpoints.BASE_MESSAGES_URL) if err != nil { - return nil, fmt.Errorf("failed to load messages page: %w", err) - } - - data, err := c.GetAccountSettings(ctx, payload.AccountSettingsQuery{ - IncludeExtSharingAudiospacesListeningDataWithFollowers: true, - IncludeMentionFilter: true, - IncludeNSFWUserFlag: true, - IncludeNSFWAdminFlag: true, - IncludeRankedTimeline: true, - IncludeAltTextCompose: true, - IncludeExtDMAVCallSettings: true, - Ext: "ssoConnections", - IncludeCountryCode: true, - IncludeExtDMNSFWMediaFilter: true, - }) + return CurrentUserProfile{}, fmt.Errorf("failed to load messages page: %w", err) + } + + profile, err := c.GetCurrentUserProfile(ctx) if err != nil { if IsAuthError(err) { - return nil, err + return CurrentUserProfile{}, err } - c.Logger.Warn().Err(err).Msg("Failed to get account settings") - data = &response.AccountSettingsResponse{} + c.Logger.Warn().Err(err).Msg("Failed to fetch current user profile after loading messages page") + profile = CurrentUserProfile{ID: c.GetCurrentUserID()} } c.session.InitializedAt = time.Now() c.session.CacheVersion = CurrentCacheVersion - c.Logger.Info(). - Str("screen_name", data.ScreenName). - Msg("Successfully loaded and authenticated as user") + logEvt := c.Logger.Info().Str("user_id", profile.ID) + if profile.ScreenName != "" { + logEvt = logEvt.Str("screen_name", profile.ScreenName) + } + logEvt.Msg("Successfully loaded and authenticated as user") - return data, nil + return profile, nil } func (c *Client) GetCurrentUserID() string { From 9ee79212d660d5b44b0285b908d85a73622a0e5e Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Sat, 21 Mar 2026 02:07:50 +0200 Subject: [PATCH 4/9] twittermeow: remove settings.json dead code --- pkg/twittermeow/account.go | 19 ------ pkg/twittermeow/data/endpoints/endpoints.go | 1 - pkg/twittermeow/data/payload/form.go | 21 ------- pkg/twittermeow/data/response/account.go | 65 --------------------- 4 files changed, 106 deletions(-) diff --git a/pkg/twittermeow/account.go b/pkg/twittermeow/account.go index 34f5601..9041920 100644 --- a/pkg/twittermeow/account.go +++ b/pkg/twittermeow/account.go @@ -30,25 +30,6 @@ func (c *Client) Login(ctx context.Context) error { return nil } -func (c *Client) GetAccountSettings(ctx context.Context, params payload.AccountSettingsQuery) (*response.AccountSettingsResponse, error) { - encodedQuery, err := params.Encode() - if err != nil { - return nil, err - } - url := fmt.Sprintf("%s?%s", endpoints.ACCOUNT_SETTINGS_URL, string(encodedQuery)) - apiRequestOpts := apiRequestOpts{ - URL: url, - Method: http.MethodGet, - } - _, respBody, err := c.makeAPIRequest(ctx, apiRequestOpts) - if err != nil { - return nil, err - } - - data := response.AccountSettingsResponse{} - return &data, json.Unmarshal(respBody, &data) -} - func (c *Client) GetCurrentUserProfile(ctx context.Context) (CurrentUserProfile, error) { currentUserID := strings.TrimSpace(c.GetCurrentUserID()) if currentUserID == "" { diff --git a/pkg/twittermeow/data/endpoints/endpoints.go b/pkg/twittermeow/data/endpoints/endpoints.go index 6f6ee37..87f3345 100644 --- a/pkg/twittermeow/data/endpoints/endpoints.go +++ b/pkg/twittermeow/data/endpoints/endpoints.go @@ -16,7 +16,6 @@ const ( API_BASE_HOST = "api.x.com" API_BASE_URL = "https://" + API_BASE_HOST - ACCOUNT_SETTINGS_URL = BASE_URL + "/i/api/1.1/account/settings.json" INBOX_INITIAL_STATE_URL = BASE_URL + "/i/api/1.1/dm/inbox_initial_state.json" DM_USER_UPDATES_URL = BASE_URL + "/i/api/1.1/dm/user_updates.json" TRUSTED_INBOX_TIMELINE_URL = BASE_URL + "/i/api/1.1/dm/inbox_timeline/trusted.json" diff --git a/pkg/twittermeow/data/payload/form.go b/pkg/twittermeow/data/payload/form.go index fa1bcce..82be6cd 100644 --- a/pkg/twittermeow/data/payload/form.go +++ b/pkg/twittermeow/data/payload/form.go @@ -26,27 +26,6 @@ func (p *JotClientEventPayload) Encode() ([]byte, error) { return []byte(values.Encode()), nil } -type AccountSettingsQuery struct { - IncludeExtSharingAudiospacesListeningDataWithFollowers bool `url:"include_ext_sharing_audiospaces_listening_data_with_followers"` - IncludeMentionFilter bool `url:"include_mention_filter"` - IncludeNSFWUserFlag bool `url:"include_nsfw_user_flag"` - IncludeNSFWAdminFlag bool `url:"include_nsfw_admin_flag"` - IncludeRankedTimeline bool `url:"include_ranked_timeline"` - IncludeAltTextCompose bool `url:"include_alt_text_compose"` - IncludeExtDMAVCallSettings bool `url:"include_ext_dm_av_call_settings"` - Ext string `url:"ext"` - IncludeCountryCode bool `url:"include_country_code"` - IncludeExtDMNSFWMediaFilter bool `url:"include_ext_dm_nsfw_media_filter"` -} - -func (p *AccountSettingsQuery) Encode() ([]byte, error) { - values, err := query.Values(p) - if err != nil { - return nil, err - } - return []byte(values.Encode()), nil -} - type ContextInfo string const ( diff --git a/pkg/twittermeow/data/response/account.go b/pkg/twittermeow/data/response/account.go index da1f710..f1a7f4f 100644 --- a/pkg/twittermeow/data/response/account.go +++ b/pkg/twittermeow/data/response/account.go @@ -2,71 +2,6 @@ package response import "go.mau.fi/mautrix-twitter/pkg/twittermeow/data/types" -type AccountSettingsResponse struct { - Protected bool `json:"protected,omitempty"` - ScreenName string `json:"screen_name,omitempty"` - AlwaysUseHTTPS bool `json:"always_use_https,omitempty"` - UseCookiePersonalization bool `json:"use_cookie_personalization,omitempty"` - SleepTime SleepTime `json:"sleep_time,omitempty"` - GeoEnabled bool `json:"geo_enabled,omitempty"` - Language string `json:"language,omitempty"` - DiscoverableByEmail bool `json:"discoverable_by_email,omitempty"` - DiscoverableByMobilePhone bool `json:"discoverable_by_mobile_phone,omitempty"` - DisplaySensitiveMedia bool `json:"display_sensitive_media,omitempty"` - PersonalizedTrends bool `json:"personalized_trends,omitempty"` - AllowMediaTagging string `json:"allow_media_tagging,omitempty"` - AllowContributorRequest string `json:"allow_contributor_request,omitempty"` - AllowAdsPersonalization bool `json:"allow_ads_personalization,omitempty"` - AllowLoggedOutDevicePersonalization bool `json:"allow_logged_out_device_personalization,omitempty"` - AllowLocationHistoryPersonalization bool `json:"allow_location_history_personalization,omitempty"` - AllowSharingDataForThirdPartyPersonalization bool `json:"allow_sharing_data_for_third_party_personalization,omitempty"` - AllowDmsFrom string `json:"allow_dms_from,omitempty"` - AlwaysAllowDmsFromSubscribers any `json:"always_allow_dms_from_subscribers,omitempty"` - AllowDmGroupsFrom string `json:"allow_dm_groups_from,omitempty"` - TranslatorType string `json:"translator_type,omitempty"` - CountryCode string `json:"country_code,omitempty"` - NSFWUser bool `json:"nsfw_user,omitempty"` - NSFWAdmin bool `json:"nsfw_admin,omitempty"` - RankedTimelineSetting any `json:"ranked_timeline_setting,omitempty"` - RankedTimelineEligible any `json:"ranked_timeline_eligible,omitempty"` - AddressBookLiveSyncEnabled bool `json:"address_book_live_sync_enabled,omitempty"` - UniversalQualityFilteringEnabled string `json:"universal_quality_filtering_enabled,omitempty"` - DMReceiptSetting string `json:"dm_receipt_setting,omitempty"` - AltTextComposeEnabled any `json:"alt_text_compose_enabled,omitempty"` - MentionFilter string `json:"mention_filter,omitempty"` - AllowAuthenticatedPeriscopeRequests bool `json:"allow_authenticated_periscope_requests,omitempty"` - ProtectPasswordReset bool `json:"protect_password_reset,omitempty"` - RequirePasswordLogin bool `json:"require_password_login,omitempty"` - RequiresLoginVerification bool `json:"requires_login_verification,omitempty"` - ExtSharingAudiospacesListeningDataWithFollowers bool `json:"ext_sharing_audiospaces_listening_data_with_followers,omitempty"` - Ext Ext `json:"ext,omitempty"` - DmQualityFilter string `json:"dm_quality_filter,omitempty"` - AutoplayDisabled bool `json:"autoplay_disabled,omitempty"` - SettingsMetadata SettingsMetadata `json:"settings_metadata,omitempty"` -} -type SleepTime struct { - Enabled bool `json:"enabled,omitempty"` - EndTime any `json:"end_time,omitempty"` - StartTime any `json:"start_time,omitempty"` -} -type Ok struct { - SSOIDHash string `json:"ssoIdHash,omitempty"` - SSOProvider string `json:"ssoProvider,omitempty"` -} -type R struct { - Ok []Ok `json:"ok,omitempty"` -} -type SsoConnections struct { - R R `json:"r,omitempty"` - TTL int `json:"ttl,omitempty"` -} -type Ext struct { - SsoConnections SsoConnections `json:"ssoConnections,omitempty"` -} -type SettingsMetadata struct { - IsEU string `json:"is_eu,omitempty"` -} - type GetDMPermissionsResponse struct { Permissions Permissions `json:"permissions,omitempty"` Users map[string]types.User `json:"users,omitempty"` From dcbd212c3564e64d323ed942f5e0799501cc01f8 Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Sun, 22 Mar 2026 19:13:20 +0200 Subject: [PATCH 5/9] twittermeow: fix ondemand script parsing --- pkg/twittermeow/client.go | 17 ++++---- pkg/twittermeow/methods/html.go | 20 +++++++--- pkg/twittermeow/methods/html_test.go | 60 ++++++++++++++++++++++++++++ 3 files changed, 84 insertions(+), 13 deletions(-) create mode 100644 pkg/twittermeow/methods/html_test.go diff --git a/pkg/twittermeow/client.go b/pkg/twittermeow/client.go index f42dfb0..54c13aa 100644 --- a/pkg/twittermeow/client.go +++ b/pkg/twittermeow/client.go @@ -277,11 +277,11 @@ func (c *Client) fetchScript(ctx context.Context, url string) ([]byte, error) { return scriptRespBody, err } -func (c *Client) fetchAndParseMainScript(ctx context.Context, scriptURL string) { +func (c *Client) fetchAndParseMainScript(ctx context.Context, scriptURL string) string { scriptRespBody, err := c.fetchScript(ctx, scriptURL) if err != nil { zerolog.Ctx(ctx).Warn().Err(err).Msg("Failed to fetch main script") - return + return "" } authTokenBytes := methods.ParseBearerToken(scriptRespBody) authTokens := exslices.CastFunc(authTokenBytes, func(from []byte) string { @@ -299,6 +299,7 @@ func (c *Client) fetchAndParseMainScript(ctx context.Context, scriptURL string) Msg("Hardcoded token doesn't match fetched one") c.session.bearerToken = authTokens[0] } + return methods.ParseOndemandSURLFromScript(scriptRespBody) } func (c *Client) fetchAndParseSScript(ctx context.Context, scriptURL string) (*[4]int, error) { @@ -395,16 +396,16 @@ func (c *Client) parseMainPageHTML(ctx context.Context, mainPageResp *http.Respo } mainScriptURL := methods.ParseMainScriptURL(mainPageHTML) + ondemandSURL := methods.ParseOndemandSURLFromScript([]byte(mainPageHTML)) if mainScriptURL == "" { zerolog.Ctx(ctx).Warn().Int("status_code", mainPageResp.StatusCode).Msg("Main script URL not found in main page HTML") - } else { - c.fetchAndParseMainScript(ctx, mainScriptURL) + } else if ondemandSURL == "" { + ondemandSURL = c.fetchAndParseMainScript(ctx, mainScriptURL) } - ondemandS := methods.ParseOndemandS(mainPageHTML) - if ondemandS == "" { - c.Logger.Warn().Msg("ondemand.s not found in main page HTML") - } else if indexes, err := c.fetchAndParseSScript(ctx, fmt.Sprintf("https://abs.twimg.com/responsive-web/client-web/ondemand.s.%sa.js", ondemandS)); err != nil { + if ondemandSURL == "" { + c.Logger.Warn().Msg("ondemand.s URL not found in bootstrap sources") + } else if indexes, err := c.fetchAndParseSScript(ctx, ondemandSURL); err != nil { c.Logger.Warn().Err(err).Msg("Failed to fetch and parse s script") } else { c.session.variableIndexes = indexes diff --git a/pkg/twittermeow/methods/html.go b/pkg/twittermeow/methods/html.go index d275b8b..d6e88b2 100644 --- a/pkg/twittermeow/methods/html.go +++ b/pkg/twittermeow/methods/html.go @@ -14,7 +14,7 @@ var ( guestTokenRegex = regexp.MustCompile(`gt=([0-9]+)`) verificationTokenRegex = regexp.MustCompile(`meta name="twitter-site-verification" content="([^"]+)"`) countryCodeRegex = regexp.MustCompile(`"country":\s*"([A-Z]{2})"`) - ondemandSRegex = regexp.MustCompile(`"ondemand.s":"([a-f0-9]+)"`) + ondemandSChunkIDRegex = regexp.MustCompile(`(\d+):"ondemand\.s"`) variableIndexesRegex = regexp.MustCompile(`\[.+?\(\w{1,2}\[(\d{1,2})],16\).+?\(\w{1,2}\[(\d{1,2})],16\).+?\(\w{1,2}\[(\d{1,2})],16\).+?\(\w{1,2}\[(\d{1,2})],16\)`) ) @@ -75,10 +75,20 @@ func ParseCountry(html string) string { return match[1] } -func ParseOndemandS(html string) string { - match := ondemandSRegex.FindStringSubmatch(html) - if len(match) < 2 { +func ParseOndemandSURLFromScript(js []byte) string { + chunkIDMatch := ondemandSChunkIDRegex.FindSubmatchIndex(js) + if len(chunkIDMatch) < 4 { return "" } - return match[1] + + chunkID := string(js[chunkIDMatch[2]:chunkIDMatch[3]]) + hashRegex := regexp.MustCompile(`(?:^|[,{])` + regexp.QuoteMeta(chunkID) + `:"([0-9a-f]+)"`) + jsAfterNameMap := js[chunkIDMatch[1]:] + hashMatch := hashRegex.FindSubmatchIndex(jsAfterNameMap) + if len(hashMatch) < 4 { + return "" + } + + hash := string(jsAfterNameMap[hashMatch[2]:hashMatch[3]]) + return "https://abs.twimg.com/responsive-web/client-web/ondemand.s." + hash + "a.js" } diff --git a/pkg/twittermeow/methods/html_test.go b/pkg/twittermeow/methods/html_test.go new file mode 100644 index 0000000..8385d34 --- /dev/null +++ b/pkg/twittermeow/methods/html_test.go @@ -0,0 +1,60 @@ +package methods + +import ( + "io" + "net/http" + "testing" + "time" +) + +func TestParseOndemandSURLFromScript(t *testing.T) { + client := &http.Client{Timeout: 20 * time.Second} + req, err := http.NewRequest(http.MethodGet, "https://x.com/", nil) + if err != nil { + t.Fatalf("failed to create request: %v", err) + } + req.Header.Set("User-Agent", "Mozilla/5.0") + req.Header.Set("Accept-Language", "en-US,en;q=0.9") + + resp, err := client.Do(req) + if err != nil { + t.Skipf("failed to fetch x.com: %v", err) + } + defer resp.Body.Close() + + html, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatalf("failed to read x.com response: %v", err) + } + + ondemandURL := ParseOndemandSURLFromScript(html) + if ondemandURL == "" { + mainScriptURL := ParseMainScriptURL(string(html)) + if mainScriptURL == "" { + t.Fatalf("failed to locate main script URL from x.com response") + } + + req, err = http.NewRequest(http.MethodGet, mainScriptURL, nil) + if err != nil { + t.Fatalf("failed to create main script request: %v", err) + } + req.Header.Set("User-Agent", "Mozilla/5.0") + req.Header.Set("Accept-Language", "en-US,en;q=0.9") + + resp, err = client.Do(req) + if err != nil { + t.Fatalf("failed to fetch main script: %v", err) + } + defer resp.Body.Close() + + script, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatalf("failed to read main script response: %v", err) + } + ondemandURL = ParseOndemandSURLFromScript(script) + } + + if ondemandURL == "" { + t.Fatalf("failed to resolve ondemand.s URL from live x.com bootstrap") + } +} From 169f2fb703e93f998d2493032c36e8608226299d Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Mon, 30 Mar 2026 13:16:45 +0300 Subject: [PATCH 6/9] connector & twittermeow: make room registration and backfill go through ChatResync rather than direct calls --- pkg/connector/chatsync.go | 96 +++++---------- pkg/connector/conversationdata.go | 2 +- pkg/connector/handlematrix.go | 8 +- pkg/connector/handletwit.go | 167 +++++++++++++++------------ pkg/twittermeow/account.go | 29 ++--- pkg/twittermeow/client.go | 12 +- pkg/twittermeow/messaging.go | 37 +++--- pkg/twittermeow/methods/html_test.go | 71 ++++-------- pkg/twittermeow/polling.go | 8 +- pkg/twittermeow/stream_client.go | 4 +- pkg/twittermeow/xchat_send.go | 44 +++---- 11 files changed, 204 insertions(+), 274 deletions(-) diff --git a/pkg/connector/chatsync.go b/pkg/connector/chatsync.go index d6ac570..e578677 100644 --- a/pkg/connector/chatsync.go +++ b/pkg/connector/chatsync.go @@ -20,13 +20,11 @@ import ( "context" "encoding/base64" "strings" - "time" "github.com/rs/zerolog" "go.mau.fi/util/ptr" "maunium.net/go/mautrix/bridgev2" "maunium.net/go/mautrix/bridgev2/database" - "maunium.net/go/mautrix/bridgev2/simplevent" "go.mau.fi/mautrix-twitter/pkg/twittermeow" "go.mau.fi/mautrix-twitter/pkg/twittermeow/crypto" @@ -61,7 +59,7 @@ func shouldEmitChatInfoUpdate(chatInfo *bridgev2.ChatInfo, portalRoomType databa } // syncXChatChannel syncs a single conversation from XChat inbox data. -// Creates the portal synchronously if it doesn't exist. +// It queues a portal resync and lets bridgev2 create the room if needed. func (tc *TwitterClient) syncXChatChannel(ctx context.Context, item *response.XChatInboxItem, users map[string]*types.User) { log := zerolog.Ctx(ctx) @@ -94,51 +92,23 @@ func (tc *TwitterClient) syncXChatChannel(ctx context.Context, item *response.XC } } - // Ensure a backfill task exists even if we don't end up emitting a ChatInfoChange. - // Beeper scrollback relies on the backfill task existing for the portal. - if portal.MXID != "" { - if chatInfo.CanBackfill { - if err := tc.connector.br.DB.BackfillTask.EnsureExists(ctx, portal.PortalKey, tc.userLogin.ID); err != nil { - log.Warn().Err(err). - Str("conversation_id", conv.ConversationID). - Msg("Failed to ensure backfill task exists") - } else { - tc.connector.br.WakeupBackfillQueue() - } - } - } - - // Create Matrix room if it doesn't exist + // Queue a ChatResync so bridgev2 owns room creation and backfill task registration. if portal.MXID == "" { - err = portal.CreateMatrixRoom(ctx, tc.userLogin, chatInfo) - if err != nil { - log.Warn().Err(err). + resync := tc.queueChatResyncNow(chatResyncCreate, portal.PortalKey, chatInfo) + if !resync.Success { + log.Warn(). Str("conversation_id", conv.ConversationID). - Msg("Failed to create Matrix room") + Err(resync.Error). + Msg("Failed to queue ChatResync for XChat conversation") return } - // Register backfill task for the newly created room - if chatInfo.CanBackfill { - if err := tc.connector.br.DB.BackfillTask.EnsureExists(ctx, portal.PortalKey, tc.userLogin.ID); err != nil { - log.Warn().Err(err). - Str("conversation_id", conv.ConversationID). - Msg("Failed to ensure backfill task exists for new room") - } else { - tc.connector.br.WakeupBackfillQueue() - } - } - } else { - if shouldEmitChatInfoUpdate(chatInfo, portal.RoomType) { - tc.userLogin.QueueRemoteEvent(&simplevent.ChatInfoChange{ - EventMeta: simplevent.EventMeta{ - Type: bridgev2.RemoteEventChatInfoChange, - PortalKey: portal.PortalKey, - Timestamp: time.Now(), - }, - ChatInfoChange: &bridgev2.ChatInfoChange{ - ChatInfo: chatInfo, - }, - }) + } else if chatInfo.CanBackfill || shouldEmitChatInfoUpdate(chatInfo, portal.RoomType) { + resync := tc.queueChatResyncNow(chatResyncUpdate, portal.PortalKey, chatInfo) + if !resync.Success { + log.Warn(). + Str("conversation_id", conv.ConversationID). + Err(resync.Error). + Msg("Failed to queue ChatResync for existing XChat conversation") } } @@ -516,37 +486,23 @@ func (tc *TwitterClient) syncUntrustedConversation(ctx context.Context, conv *ty chatInfo := tc.conversationToChatInfo(ctx, conv, inbox) - // Create Matrix room if it doesn't exist + // Queue a ChatResync so bridgev2 owns room creation and backfill task registration. if portal.MXID == "" { - err = portal.CreateMatrixRoom(ctx, tc.userLogin, chatInfo) - if err != nil { - log.Warn().Err(err). + resync := tc.queueChatResyncNow(chatResyncCreate, portal.PortalKey, chatInfo) + if !resync.Success { + log.Warn(). Str("conversation_id", conv.ConversationID). - Msg("Failed to create Matrix room for untrusted conversation") + Err(resync.Error). + Msg("Failed to queue ChatResync for untrusted conversation") return } } else { - // Room already exists - update MessageRequest status via ChatInfoChange - tc.userLogin.QueueRemoteEvent(&simplevent.ChatInfoChange{ - EventMeta: simplevent.EventMeta{ - Type: bridgev2.RemoteEventChatInfoChange, - PortalKey: portal.PortalKey, - Timestamp: time.Now(), - }, - ChatInfoChange: &bridgev2.ChatInfoChange{ - ChatInfo: chatInfo, - }, - }) - } - - // Ensure untrusted conversations also have a queue backfill task once a room exists. - if portal.MXID != "" && chatInfo.CanBackfill { - if err := tc.connector.br.DB.BackfillTask.EnsureExists(ctx, portal.PortalKey, tc.userLogin.ID); err != nil { - log.Warn().Err(err). + resync := tc.queueChatResyncNow(chatResyncUpdate, portal.PortalKey, chatInfo) + if !resync.Success { + log.Warn(). Str("conversation_id", conv.ConversationID). - Msg("Failed to ensure backfill task exists for untrusted conversation") - } else { - tc.connector.br.WakeupBackfillQueue() + Err(resync.Error). + Msg("Failed to queue ChatResync for existing untrusted conversation") } } @@ -581,7 +537,7 @@ func (tc *TwitterClient) processUntrustedMessages(ctx context.Context, conversat } // Queue the message event - tc.HandlePollingEvent(msg, inbox) + tc.HandlePollingEvent(ctx, msg, inbox) } } diff --git a/pkg/connector/conversationdata.go b/pkg/connector/conversationdata.go index a811c52..693d038 100644 --- a/pkg/connector/conversationdata.go +++ b/pkg/connector/conversationdata.go @@ -112,7 +112,7 @@ func (tc *TwitterClient) ensurePortalForConversation(ctx context.Context, conver log.Warn().Err(err).Msg("Failed to process key change events for fetched conversation data") } - // Sync channel (creates portal if needed) + // Sync channel and queue a resync if the portal still needs a room. tc.syncXChatChannel(ctx, item, users) // Process messages/read events to backfill and register any keys embedded there diff --git a/pkg/connector/handlematrix.go b/pkg/connector/handlematrix.go index 6682cc8..1435e0e 100644 --- a/pkg/connector/handlematrix.go +++ b/pkg/connector/handlematrix.go @@ -590,7 +590,11 @@ func (tc *TwitterClient) doHandleMatrixReaction(ctx context.Context, remove bool // XChat reactions are sent as encrypted MessageCreateEvents (reaction_add/reaction_remove). xchatConvID := NormalizeConversationID(conversationID) - _, err := tc.client.SendEncryptedReaction(ctx, xchatConvID, messageID, emoji, remove) + action := twittermeow.SendEncryptedReactionAdd + if remove { + action = twittermeow.SendEncryptedReactionRemove + } + _, err := tc.client.SendEncryptedReaction(ctx, xchatConvID, messageID, emoji, action) return err } @@ -781,7 +785,7 @@ func (tc *TwitterClient) HandleMatrixViewingChat(ctx context.Context, chat *brid if chat.Portal != nil { conversationID = ParsePortalID(chat.Portal.ID) } - tc.client.SetActiveConversation(ConvertConversationIDToREST(conversationID)) + tc.client.SetActiveConversation(context.WithoutCancel(ctx), ConvertConversationIDToREST(conversationID)) return nil } diff --git a/pkg/connector/handletwit.go b/pkg/connector/handletwit.go index 87493de..f753144 100644 --- a/pkg/connector/handletwit.go +++ b/pkg/connector/handletwit.go @@ -88,6 +88,51 @@ func (tc *TwitterClient) HandleStreamEvent(evt response.StreamEvent) { } } +type chatResyncMode uint8 + +const ( + chatResyncUpdate chatResyncMode = iota + chatResyncCreate + chatResyncCreateWithBackfill +) + +func (tc *TwitterClient) queueChatResync( + mode chatResyncMode, + portalKey networkid.PortalKey, + chatInfo *bridgev2.ChatInfo, + timestamp time.Time, + streamOrder int64, +) bridgev2.EventHandlingResult { + evt := &simplevent.ChatResync{ + EventMeta: simplevent.EventMeta{ + Type: bridgev2.RemoteEventChatResync, + PortalKey: portalKey, + Timestamp: timestamp, + StreamOrder: streamOrder, + }, + ChatInfo: chatInfo, + } + + switch mode { + case chatResyncUpdate: + case chatResyncCreate: + evt.CreatePortal = true + case chatResyncCreateWithBackfill: + evt.CreatePortal = true + evt.CheckNeedsBackfillFunc = func(context.Context, *database.Message) (bool, error) { + return true, nil + } + default: + panic("unknown chat resync mode") + } + + return tc.userLogin.QueueRemoteEvent(evt) +} + +func (tc *TwitterClient) queueChatResyncNow(mode chatResyncMode, portalKey networkid.PortalKey, chatInfo *bridgev2.ChatInfo) bridgev2.EventHandlingResult { + return tc.queueChatResync(mode, portalKey, chatInfo, time.Now(), 0) +} + // buildMemberChangeEvent creates a ChatInfoChange event for participant joins/leaves. func (tc *TwitterClient) buildMemberChangeEvent( conversationID, eventID, eventTime string, @@ -152,7 +197,8 @@ func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.Twit targetMessageID = eventID } - if ctx == nil || ctx.Value(ensurePortalContextKey{}) == nil { + bootstrapCreatePortal := ctx.Value(ensurePortalContextKey{}) != nil + if !bootstrapCreatePortal { if _, err := tc.ensurePortalForConversation(ctx, evt.ConversationID, requiredKeyVersion); err != nil { log.Warn(). Err(err). @@ -178,7 +224,7 @@ func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.Twit Bool("is_from_me", isFromMe) }, PortalKey: portalKey, - CreatePortal: false, + CreatePortal: bootstrapCreatePortal, Sender: tc.MakeEventSender(evt.MessageData.SenderID), StreamOrder: streamOrder, Timestamp: methods.ParseMsecTimestamp(evt.Time), @@ -206,7 +252,8 @@ func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.Twit streamOrder = methods.ParseInt64(evt.Time) } - if ctx == nil || ctx.Value(ensurePortalContextKey{}) == nil { + bootstrapCreatePortal := ctx.Value(ensurePortalContextKey{}) != nil + if !bootstrapCreatePortal { if _, err := tc.ensurePortalForConversation(ctx, evt.ConversationID, requiredKeyVersion); err != nil { log.Warn(). Err(err). @@ -233,7 +280,7 @@ func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.Twit Bool("is_from_me", isFromMe) }, PortalKey: portalKey, - CreatePortal: false, // Portal should already exist from initial sync + CreatePortal: bootstrapCreatePortal, Sender: tc.MakeEventSender(evt.MessageData.SenderID), StreamOrder: streamOrder, Timestamp: methods.ParseMsecTimestamp(evt.Time), @@ -361,10 +408,6 @@ func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.Twit return tc.userLogin.QueueRemoteEvent(portalDeleteRemoteEvent).Success case *types.ConversationNameUpdate: - if ctx == nil { - ctx = context.Background() - } - // XChat group titles are encrypted. Decrypt before forwarding to Matrix so // we don't set the room name to ciphertext. newName := evt.ConversationName @@ -479,19 +522,13 @@ func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.Twit return false } - return tc.userLogin.QueueRemoteEvent(&simplevent.ChatResync{ - EventMeta: simplevent.EventMeta{ - Type: bridgev2.RemoteEventChatResync, - PortalKey: portalKey, - CreatePortal: true, - Timestamp: methods.ParseMsecTimestamp(evt.Time), - StreamOrder: methods.ParseInt64(evt.ID), - }, - ChatInfo: chatInfo, - CheckNeedsBackfillFunc: func(ctx context.Context, latestMessage *database.Message) (bool, error) { - return true, nil - }, - }).Success + return tc.queueChatResync( + chatResyncCreateWithBackfill, + portalKey, + chatInfo, + methods.ParseMsecTimestamp(evt.Time), + methods.ParseInt64(evt.ID), + ).Success } return true @@ -530,16 +567,13 @@ func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.Twit return false } - return tc.userLogin.QueueRemoteEvent(&simplevent.ChatResync{ - EventMeta: simplevent.EventMeta{ - Type: bridgev2.RemoteEventChatResync, - PortalKey: tc.MakePortalKeyFromID(evt.ConversationID), - CreatePortal: true, - Timestamp: methods.ParseMsecTimestamp(evt.Time), - StreamOrder: methods.ParseInt64(evt.ID), - }, - ChatInfo: chatInfo, - }).Success + return tc.queueChatResync( + chatResyncCreate, + tc.MakePortalKeyFromID(evt.ConversationID), + chatInfo, + methods.ParseMsecTimestamp(evt.Time), + methods.ParseInt64(evt.ID), + ).Success default: log.Debug(). @@ -557,10 +591,10 @@ var _ = payload.FailureType(0) // This is used for untrusted (message request) conversations that don't // receive real-time updates via XChat WebSocket. // Returns true to continue polling, false to stop. -func (tc *TwitterClient) HandlePollingEvent(evt types.TwitterEvent, inbox *response.TwitterInboxData) bool { +func (tc *TwitterClient) HandlePollingEvent(ctx context.Context, evt types.TwitterEvent, inbox *response.TwitterInboxData) bool { // Always cache users from inbox when available - needed for portal creation if inbox != nil { - tc.updateTwitterUserInfo(context.Background(), inbox) + tc.updateTwitterUserInfo(ctx, inbox) tc.userCacheLock.Lock() for userID, user := range inbox.Users { tc.userCache[userID] = user @@ -615,7 +649,6 @@ func (tc *TwitterClient) HandlePollingEvent(evt types.TwitterEvent, inbox *respo // Skip if conversation can use XChat (has encryption keys) - XChat WebSocket handles those portalKey := tc.MakePortalKeyFromID(conversationID) - ctx := context.Background() if portal, err := tc.connector.br.GetPortalByKey(ctx, portalKey); err == nil && portal != nil { meta := portal.Metadata.(*PortalMetadata) if meta.CanUseXChat() { @@ -629,7 +662,7 @@ func (tc *TwitterClient) HandlePollingEvent(evt types.TwitterEvent, inbox *respo // Dispatch to the appropriate handler based on event type switch e := evt.(type) { case *types.Message: - return tc.handlePollingMessage(e, inbox) + return tc.handlePollingMessage(ctx, e, inbox) case *types.MessageReactionCreate: reaction := (*types.MessageReaction)(e) portalKey := tc.MakePortalKeyFromID(conversationID) @@ -681,7 +714,6 @@ func (tc *TwitterClient) HandlePollingEvent(evt types.TwitterEvent, inbox *respo Msg("Conversation became trusted via polling") // Update portal metadata to mark as trusted (same as XChat path) - ctx := context.Background() portalKey := tc.MakePortalKeyFromID(conversationID) portal, err := tc.connector.br.GetPortalByKey(ctx, portalKey) if err != nil { @@ -705,16 +737,13 @@ func (tc *TwitterClient) HandlePollingEvent(evt types.TwitterEvent, inbox *respo return false } - return tc.userLogin.QueueRemoteEvent(&simplevent.ChatResync{ - EventMeta: simplevent.EventMeta{ - Type: bridgev2.RemoteEventChatResync, - PortalKey: tc.MakePortalKeyFromID(conversationID), - CreatePortal: true, - Timestamp: methods.ParseMsecTimestamp(e.Time), - StreamOrder: methods.ParseInt64(e.ID), - }, - ChatInfo: chatInfo, - }).Success + return tc.queueChatResync( + chatResyncCreate, + tc.MakePortalKeyFromID(conversationID), + chatInfo, + methods.ParseMsecTimestamp(e.Time), + methods.ParseInt64(e.ID), + ).Success } return true @@ -752,7 +781,7 @@ func (tc *TwitterClient) markPollingChatResyncSuccess(conversationID string, now } // handlePollingMessage handles a message event from REST API polling. -func (tc *TwitterClient) handlePollingMessage(evt *types.Message, inbox *response.TwitterInboxData) bool { +func (tc *TwitterClient) handlePollingMessage(ctx context.Context, evt *types.Message, inbox *response.TwitterInboxData) bool { isFromMe := MakeUserLoginID(evt.MessageData.SenderID) == tc.userLogin.ID portalKey := tc.MakePortalKeyFromID(evt.ConversationID) msgID := evt.ID @@ -763,7 +792,6 @@ func (tc *TwitterClient) handlePollingMessage(evt *types.Message, inbox *respons Logger() // For polling messages, ensure the portal exists - ctx := context.Background() portal, err := tc.connector.br.GetPortalByKey(ctx, portalKey) if err != nil { log.Warn(). @@ -772,7 +800,7 @@ func (tc *TwitterClient) handlePollingMessage(evt *types.Message, inbox *respons return false } - // Create portal if it doesn't exist + // Queue a portal resync if the room doesn't exist yet. if portal.MXID == "" { chatInfo := tc.getOrFetchChatInfoForPolling(ctx, evt.ConversationID, inbox) if chatInfo == nil { @@ -780,19 +808,19 @@ func (tc *TwitterClient) handlePollingMessage(evt *types.Message, inbox *respons Msg("Failed to get chat info for polling message") return false } - if err := portal.CreateMatrixRoom(ctx, tc.userLogin, chatInfo); err != nil { + resync := tc.queueChatResync( + chatResyncCreate, + portal.PortalKey, + chatInfo, + methods.ParseMsecTimestamp(evt.Time), + methods.ParseInt64(msgID), + ) + if !resync.Success { log.Warn(). - Err(err). - Msg("Failed to create Matrix room for polling message") + Err(resync.Error). + Msg("Failed to queue ChatResync for polling message") return false } - // Register backfill task for the newly created room - if err := tc.connector.br.DB.BackfillTask.EnsureExists(ctx, portal.PortalKey, tc.userLogin.ID); err != nil { - log.Warn().Err(err). - Msg("Failed to ensure backfill task exists for new polling room") - } else { - tc.connector.br.WakeupBackfillQueue() - } } else if tc.shouldAttemptPollingChatResync(evt.ConversationID, now) { chatInfo := tc.getOrFetchChatInfoForPolling(ctx, evt.ConversationID, inbox) if !tc.isDMChatInfoComplete(chatInfo) { @@ -804,21 +832,18 @@ func (tc *TwitterClient) handlePollingMessage(evt *types.Message, inbox *respons Int("missing_userinfo", missingUserInfo). Msg("Skipping polling ChatResync: incomplete chat info") } else { - ok := tc.userLogin.QueueRemoteEvent(&simplevent.ChatResync{ - EventMeta: simplevent.EventMeta{ - Type: bridgev2.RemoteEventChatResync, - PortalKey: portal.PortalKey, - CreatePortal: true, - Timestamp: methods.ParseMsecTimestamp(evt.Time), - StreamOrder: methods.ParseInt64(msgID), - }, - ChatInfo: chatInfo, - }).Success - if ok { + resync := tc.queueChatResync( + chatResyncUpdate, + portal.PortalKey, + chatInfo, + methods.ParseMsecTimestamp(evt.Time), + methods.ParseInt64(msgID), + ) + if resync.Success { log.Debug().Msg("Queued polling ChatResync") tc.markPollingChatResyncSuccess(evt.ConversationID, now) } else { - log.Debug().Msg("Failed to queue polling ChatResync") + log.Debug().Err(resync.Error).Msg("Failed to queue polling ChatResync") } } } diff --git a/pkg/twittermeow/account.go b/pkg/twittermeow/account.go index 9041920..62e506f 100644 --- a/pkg/twittermeow/account.go +++ b/pkg/twittermeow/account.go @@ -40,36 +40,37 @@ func (c *Client) GetCurrentUserProfile(ctx context.Context) (CurrentUserProfile, if err != nil { return CurrentUserProfile{}, err } - if len(resp.Errors) > 0 && resp.Errors[0].Message != "" { - return CurrentUserProfile{}, fmt.Errorf("GetUsersByIdsForXChat error: %s", resp.Errors[0].Message) + if len(resp.Errors) > 0 { + return CurrentUserProfile{}, fmt.Errorf("get user profile: %s", resp.Errors[0].Message) } if len(resp.Data.GetMemberResults.Results) != 1 { return CurrentUserProfile{}, fmt.Errorf("expected 1 user result for %s, got %d", currentUserID, len(resp.Data.GetMemberResults.Results)) } result := resp.Data.GetMemberResults.Results[0] - if result.MemberResults == nil || result.MemberResults.Result == nil || result.MemberResults.Result.Core == nil { + member := result.MemberResults + if member == nil || member.Result == nil || member.Result.Core == nil { return CurrentUserProfile{}, fmt.Errorf("GetUsersByIdsForXChat returned no user for %s", currentUserID) } - resultUserID := result.MemberResults.RestID + resultUserID := member.RestID if resultUserID == "" { - resultUserID = result.MemberResults.Result.RestID + resultUserID = member.Result.RestID } - if resultUserID != "" && resultUserID != currentUserID { + if resultUserID == "" { + resultUserID = currentUserID + } + if resultUserID != currentUserID { return CurrentUserProfile{}, fmt.Errorf("GetUsersByIdsForXChat returned user %s for %s", resultUserID, currentUserID) } profile := CurrentUserProfile{ - ID: currentUserID, - ScreenName: strings.TrimSpace(result.MemberResults.Result.Core.ScreenName), - Name: strings.TrimSpace(result.MemberResults.Result.Core.Name), - } - if resultUserID != "" { - profile.ID = resultUserID + ID: resultUserID, + ScreenName: strings.TrimSpace(member.Result.Core.ScreenName), + Name: strings.TrimSpace(member.Result.Core.Name), } - if result.MemberResults.Result.Avatar != nil { - profile.AvatarURL = strings.TrimSpace(result.MemberResults.Result.Avatar.ImageURL) + if member.Result.Avatar != nil { + profile.AvatarURL = strings.TrimSpace(member.Result.Avatar.ImageURL) } return profile, nil diff --git a/pkg/twittermeow/client.go b/pkg/twittermeow/client.go index 54c13aa..263b236 100644 --- a/pkg/twittermeow/client.go +++ b/pkg/twittermeow/client.go @@ -27,7 +27,7 @@ import ( "go.mau.fi/mautrix-twitter/pkg/twittermeow/methods" ) -type EventHandler func(evt types.TwitterEvent, inbox *response.TwitterInboxData) bool +type EventHandler func(ctx context.Context, evt types.TwitterEvent, inbox *response.TwitterInboxData) bool type StreamEventHandler func(evt response.StreamEvent) // ConversationDataCallback is called when conversation data is refreshed (e.g., during key refresh). @@ -211,11 +211,7 @@ func (c *Client) LoadMessagesPage(ctx context.Context) (CurrentUserProfile, erro profile, err := c.GetCurrentUserProfile(ctx) if err != nil { - if IsAuthError(err) { - return CurrentUserProfile{}, err - } - c.Logger.Warn().Err(err).Msg("Failed to fetch current user profile after loading messages page") - profile = CurrentUserProfile{ID: c.GetCurrentUserID()} + return CurrentUserProfile{}, fmt.Errorf("failed to fetch current user profile after loading messages page: %w", err) } c.session.InitializedAt = time.Now() @@ -480,8 +476,8 @@ func (c *Client) makeAPIRequest(ctx context.Context, apiRequestOpts apiRequestOp return c.MakeRequest(ctx, apiRequestOpts.URL, apiRequestOpts.Method, headers, apiRequestOpts.Body, apiRequestOpts.ContentType) } -func (c *Client) SetActiveConversation(conversationID string) { - c.stream.startOrUpdateEventStream(conversationID) +func (c *Client) SetActiveConversation(ctx context.Context, conversationID string) { + c.stream.startOrUpdateEventStream(ctx, conversationID) } // FetchRaw performs an authenticated request to the given URL and returns the response and body. diff --git a/pkg/twittermeow/messaging.go b/pkg/twittermeow/messaging.go index 82bc699..5eda773 100644 --- a/pkg/twittermeow/messaging.go +++ b/pkg/twittermeow/messaging.go @@ -414,6 +414,13 @@ type SendEncryptedEditOpts struct { Entities []*payload.RichTextEntity } +type SendEncryptedReactionAction uint8 + +const ( + SendEncryptedReactionAdd SendEncryptedReactionAction = iota + SendEncryptedReactionRemove +) + func (c *Client) sendMessageMutation(ctx context.Context, pl *payload.SendMessageMutationPayload) (*response.SendMessageMutationResponse, error) { token, err := c.ensureConversationToken(ctx, pl.Variables.ConversationID) if err != nil { @@ -459,22 +466,20 @@ func (c *Client) sendMessageMutation(ctx context.Context, pl *payload.SendMessag // SendEncryptedReaction sends a reaction add/remove via the XChat protocol. // targetMessageSequenceID must be the XChat message sequence ID of the message being reacted to. -func (c *Client) SendEncryptedReaction(ctx context.Context, conversationID, targetMessageSequenceID, emoji string, remove bool) (*response.SendMessageMutationResponse, error) { - token, err := c.ensureConversationToken(ctx, conversationID) - if err != nil { - return nil, fmt.Errorf("get conversation token: %w", err) - } - +func (c *Client) SendEncryptedReaction(ctx context.Context, conversationID, targetMessageSequenceID, emoji string, action SendEncryptedReactionAction) (*response.SendMessageMutationResponse, error) { messageID := uuid.NewString() builder := crypto.NewMessageBuilder(c.keyManager, c.GetCurrentUserID()). SetMessageID(messageID). SetConversationID(conversationID) - if remove { - builder.SetReactionRemove(targetMessageSequenceID, emoji) - } else { + switch action { + case SendEncryptedReactionAdd: builder.SetReactionAdd(targetMessageSequenceID, emoji) + case SendEncryptedReactionRemove: + builder.SetReactionRemove(targetMessageSequenceID, emoji) + default: + panic("unknown encrypted reaction action") } encodedMCE, encodedSig, err := builder.BuildForSend(ctx) @@ -490,7 +495,6 @@ func (c *Client) SendEncryptedReaction(ctx context.Context, conversationID, targ pl := payload.NewSendMessageMutationPayload(payload.SendMessageMutationVariables{ ConversationID: conversationID, MessageID: messageID, - ConversationToken: token, EncodedMessageCreateEvent: encodedMCE, EncodedMessageEventSignature: sigPtr, }) @@ -508,11 +512,6 @@ func (c *Client) SendEncryptedEdit(ctx context.Context, opts SendEncryptedEditOp return nil, fmt.Errorf("target message sequence ID is required") } - token, err := c.ensureConversationToken(ctx, opts.ConversationID) - if err != nil { - return nil, fmt.Errorf("get conversation token: %w", err) - } - messageID := opts.MessageID if messageID == "" { messageID = uuid.NewString() @@ -536,7 +535,6 @@ func (c *Client) SendEncryptedEdit(ctx context.Context, opts SendEncryptedEditOp pl := payload.NewSendMessageMutationPayload(payload.SendMessageMutationVariables{ ConversationID: opts.ConversationID, MessageID: messageID, - ConversationToken: token, EncodedMessageCreateEvent: encodedMCE, EncodedMessageEventSignature: sigPtr, }) @@ -546,12 +544,6 @@ func (c *Client) SendEncryptedEdit(ctx context.Context, opts SendEncryptedEditOp // SendEncryptedMessage sends an encrypted message via the XChat protocol. func (c *Client) SendEncryptedMessage(ctx context.Context, opts SendEncryptedMessageOpts) (*response.SendMessageMutationResponse, error) { - // Get the server-provided conversation token for this conversation - token, err := c.ensureConversationToken(ctx, opts.ConversationID) - if err != nil { - return nil, fmt.Errorf("get conversation token: %w", err) - } - messageID := opts.MessageID if messageID == "" { messageID = uuid.NewString() @@ -599,7 +591,6 @@ func (c *Client) SendEncryptedMessage(ctx context.Context, opts SendEncryptedMes pl := payload.NewSendMessageMutationPayload(payload.SendMessageMutationVariables{ ConversationID: opts.ConversationID, MessageID: messageID, - ConversationToken: token, EncodedMessageCreateEvent: encodedMCE, EncodedMessageEventSignature: sigPtr, }) diff --git a/pkg/twittermeow/methods/html_test.go b/pkg/twittermeow/methods/html_test.go index 8385d34..1818bfc 100644 --- a/pkg/twittermeow/methods/html_test.go +++ b/pkg/twittermeow/methods/html_test.go @@ -1,60 +1,33 @@ package methods import ( - "io" - "net/http" "testing" - "time" ) func TestParseOndemandSURLFromScript(t *testing.T) { - client := &http.Client{Timeout: 20 * time.Second} - req, err := http.NewRequest(http.MethodGet, "https://x.com/", nil) - if err != nil { - t.Fatalf("failed to create request: %v", err) + tests := []struct { + name string + js string + want string + }{ + { + name: "find url", + js: `123:"ondemand.s",{123:"deadbeef"}`, + want: "https://abs.twimg.com/responsive-web/client-web/ondemand.s.deadbeefa.js", + }, + { + name: "missing chunk", + js: `123:"main",{123:"deadbeef"}`, + want: "", + }, } - req.Header.Set("User-Agent", "Mozilla/5.0") - req.Header.Set("Accept-Language", "en-US,en;q=0.9") - resp, err := client.Do(req) - if err != nil { - t.Skipf("failed to fetch x.com: %v", err) - } - defer resp.Body.Close() - - html, err := io.ReadAll(resp.Body) - if err != nil { - t.Fatalf("failed to read x.com response: %v", err) - } - - ondemandURL := ParseOndemandSURLFromScript(html) - if ondemandURL == "" { - mainScriptURL := ParseMainScriptURL(string(html)) - if mainScriptURL == "" { - t.Fatalf("failed to locate main script URL from x.com response") - } - - req, err = http.NewRequest(http.MethodGet, mainScriptURL, nil) - if err != nil { - t.Fatalf("failed to create main script request: %v", err) - } - req.Header.Set("User-Agent", "Mozilla/5.0") - req.Header.Set("Accept-Language", "en-US,en;q=0.9") - - resp, err = client.Do(req) - if err != nil { - t.Fatalf("failed to fetch main script: %v", err) - } - defer resp.Body.Close() - - script, err := io.ReadAll(resp.Body) - if err != nil { - t.Fatalf("failed to read main script response: %v", err) - } - ondemandURL = ParseOndemandSURLFromScript(script) - } - - if ondemandURL == "" { - t.Fatalf("failed to resolve ondemand.s URL from live x.com bootstrap") + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + got := ParseOndemandSURLFromScript([]byte(test.js)) + if got != test.want { + t.Fatalf("unexpected ondemand url: got %q want %q", got, test.want) + } + }) } } diff --git a/pkg/twittermeow/polling.go b/pkg/twittermeow/polling.go index e920b05..0effbc5 100644 --- a/pkg/twittermeow/polling.go +++ b/pkg/twittermeow/polling.go @@ -76,7 +76,7 @@ func (pc *PollingClient) doPoll(ctx context.Context) { } else if err != nil { log.Err(err).Msg("Failed to poll for updates") authError := IsAuthError(err) - pc.client.eventHandler(&types.PollingError{Error: err, IsAuth: authError}, nil) + pc.client.eventHandler(ctx, &types.PollingError{Error: err, IsAuth: authError}, nil) if authError { return } @@ -88,7 +88,7 @@ func (pc *PollingClient) doPoll(ctx context.Context) { tick.Reset(backoffInterval) } else if failing { failing = false - pc.client.eventHandler(&types.PollingError{}, nil) + pc.client.eventHandler(ctx, &types.PollingError{}, nil) backoffInterval = defaultPollingInterval / 2 tick.Reset(defaultPollingInterval) } @@ -119,7 +119,7 @@ func (pc *PollingClient) poll(ctx context.Context) error { return nil } - if !pc.client.eventHandler(nil, userUpdatesResponse.UserEvents) { + if !pc.client.eventHandler(ctx, nil, userUpdatesResponse.UserEvents) { return errEventHandlerFailed } for _, entry := range userUpdatesResponse.UserEvents.Entries { @@ -128,7 +128,7 @@ func (pc *PollingClient) poll(ctx context.Context) error { } parsed := entry.ParseWithErrorLog(&pc.client.Logger) if parsed != nil { - if !pc.client.eventHandler(parsed, userUpdatesResponse.UserEvents) { + if !pc.client.eventHandler(ctx, parsed, userUpdatesResponse.UserEvents) { return errEventHandlerFailed } } diff --git a/pkg/twittermeow/stream_client.go b/pkg/twittermeow/stream_client.go index f88c0e7..61c37a5 100644 --- a/pkg/twittermeow/stream_client.go +++ b/pkg/twittermeow/stream_client.go @@ -42,8 +42,8 @@ func (c *Client) newStreamClient() *StreamClient { return sc } -func (sc *StreamClient) startOrUpdateEventStream(conversationID string) { - ctx := sc.client.Logger.With().Str("action", "event stream").Logger().WithContext(context.Background()) +func (sc *StreamClient) startOrUpdateEventStream(ctx context.Context, conversationID string) { + ctx = sc.client.Logger.With().Str("action", "event stream").Logger().WithContext(ctx) if sc.conversationID == "" { sc.conversationID = conversationID go sc.start(ctx) diff --git a/pkg/twittermeow/xchat_send.go b/pkg/twittermeow/xchat_send.go index 4506ff6..ccf411b 100644 --- a/pkg/twittermeow/xchat_send.go +++ b/pkg/twittermeow/xchat_send.go @@ -87,19 +87,16 @@ func (c *Client) ensureConversationToken(ctx context.Context, conversationID str if err == nil { return token, nil } - if err != nil && !errors.Is(err, crypto.ErrKeyNotFound) { + if !errors.Is(err, crypto.ErrKeyNotFound) { return "", fmt.Errorf("get conversation token: %w", err) } - if err := c.refreshConversationToken(ctx, conversationID); err != nil && !errors.Is(err, crypto.ErrKeyNotFound) { + err = c.refreshConversationToken(ctx, conversationID) + if err != nil { return "", err } - token, err = c.keyManager.GetConversationToken(ctx, conversationID) - if err != nil { - return "", fmt.Errorf("get conversation token: %w", err) - } - return token, nil + return c.keyManager.GetConversationToken(ctx, conversationID) } func (c *Client) refreshConversationToken(ctx context.Context, conversationID string) error { @@ -125,33 +122,32 @@ func (c *Client) refreshConversationToken(ctx context.Context, conversationID st } for _, encodedEvt := range encoded { - err := c.putConversationTokenFromEncodedEvent(ctx, conversationID, encodedEvt) - if err == nil { - return nil - } - if !errors.Is(err, crypto.ErrKeyNotFound) { - return err + tokenConversationID, token := conversationTokenFromEncodedEvent(conversationID, encodedEvt) + if token == "" { + continue } + return c.keyManager.PutConversationToken(ctx, tokenConversationID, token) } return crypto.ErrKeyNotFound } -func (c *Client) putConversationTokenFromEncodedEvent(ctx context.Context, conversationID, encoded string) error { +func conversationTokenFromEncodedEvent(fallbackConversationID, encoded string) (conversationID, token string) { if encoded == "" { - return crypto.ErrKeyNotFound + return "", "" } evt, err := DecodeMessageEvent(encoded) if err != nil { - return crypto.ErrKeyNotFound + return "", "" } if evt == nil || evt.ConversationToken == nil || *evt.ConversationToken == "" { - return crypto.ErrKeyNotFound + return "", "" } + conversationID = fallbackConversationID if evt.ConversationId != nil && *evt.ConversationId != "" { conversationID = *evt.ConversationId } - return c.keyManager.PutConversationToken(ctx, conversationID, *evt.ConversationToken) + return conversationID, *evt.ConversationToken } // getSelfConversationID returns the user's self-conversation ID (user_id:user_id format). @@ -168,11 +164,6 @@ func (c *Client) SendXChatPinConversation(ctx context.Context, targetConversatio selfConvID := c.getSelfConversationID() - token, err := c.ensureConversationToken(ctx, selfConvID) - if err != nil { - return fmt.Errorf("get self conversation token: %w", err) - } - messageID := uuid.NewString() builder := crypto.NewMessageBuilder(c.keyManager, c.GetCurrentUserID()). @@ -193,7 +184,6 @@ func (c *Client) SendXChatPinConversation(ctx context.Context, targetConversatio pl := payload.NewSendMessageMutationPayload(payload.SendMessageMutationVariables{ ConversationID: selfConvID, MessageID: messageID, - ConversationToken: token, EncodedMessageCreateEvent: encodedMCE, EncodedMessageEventSignature: sigPtr, }) @@ -210,11 +200,6 @@ func (c *Client) SendXChatUnpinConversation(ctx context.Context, targetConversat selfConvID := c.getSelfConversationID() - token, err := c.ensureConversationToken(ctx, selfConvID) - if err != nil { - return fmt.Errorf("get self conversation token: %w", err) - } - messageID := uuid.NewString() builder := crypto.NewMessageBuilder(c.keyManager, c.GetCurrentUserID()). @@ -235,7 +220,6 @@ func (c *Client) SendXChatUnpinConversation(ctx context.Context, targetConversat pl := payload.NewSendMessageMutationPayload(payload.SendMessageMutationVariables{ ConversationID: selfConvID, MessageID: messageID, - ConversationToken: token, EncodedMessageCreateEvent: encodedMCE, EncodedMessageEventSignature: sigPtr, }) From 06cd44ffd8a03f6cfff1695126f90d022a50a7f3 Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Mon, 30 Mar 2026 13:47:05 +0300 Subject: [PATCH 7/9] connector,twittermeow: fix merge conflicts --- pkg/connector/handletwit.go | 11 ++++++ pkg/twittermeow/account.go | 6 ++- pkg/twittermeow/client.go | 6 ++- pkg/twittermeow/methods/html_test.go | 55 ++++++++++++++++++++++++++++ pkg/twittermeow/stream_client.go | 3 ++ pkg/twittermeow/xchat_send.go | 6 ++- 6 files changed, 84 insertions(+), 3 deletions(-) diff --git a/pkg/connector/handletwit.go b/pkg/connector/handletwit.go index f753144..95cff46 100644 --- a/pkg/connector/handletwit.go +++ b/pkg/connector/handletwit.go @@ -170,6 +170,9 @@ func (tc *TwitterClient) buildMemberChangeEvent( // HandleXChatEvent handles events from the XChat websocket processor. func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.TwitterEvent) bool { + if ctx == nil { + ctx = context.TODO() + } if rawEvt == nil { return true } @@ -592,6 +595,10 @@ var _ = payload.FailureType(0) // receive real-time updates via XChat WebSocket. // Returns true to continue polling, false to stop. func (tc *TwitterClient) HandlePollingEvent(ctx context.Context, evt types.TwitterEvent, inbox *response.TwitterInboxData) bool { + if ctx == nil { + ctx = context.TODO() + } + // Always cache users from inbox when available - needed for portal creation if inbox != nil { tc.updateTwitterUserInfo(ctx, inbox) @@ -782,6 +789,10 @@ func (tc *TwitterClient) markPollingChatResyncSuccess(conversationID string, now // handlePollingMessage handles a message event from REST API polling. func (tc *TwitterClient) handlePollingMessage(ctx context.Context, evt *types.Message, inbox *response.TwitterInboxData) bool { + if ctx == nil { + ctx = context.TODO() + } + isFromMe := MakeUserLoginID(evt.MessageData.SenderID) == tc.userLogin.ID portalKey := tc.MakePortalKeyFromID(evt.ConversationID) msgID := evt.ID diff --git a/pkg/twittermeow/account.go b/pkg/twittermeow/account.go index 62e506f..d464ea1 100644 --- a/pkg/twittermeow/account.go +++ b/pkg/twittermeow/account.go @@ -41,7 +41,11 @@ func (c *Client) GetCurrentUserProfile(ctx context.Context) (CurrentUserProfile, return CurrentUserProfile{}, err } if len(resp.Errors) > 0 { - return CurrentUserProfile{}, fmt.Errorf("get user profile: %s", resp.Errors[0].Message) + msg := strings.TrimSpace(resp.Errors[0].Message) + if msg == "" { + msg = "unknown error" + } + return CurrentUserProfile{}, fmt.Errorf("get user profile: %s", msg) } if len(resp.Data.GetMemberResults.Results) != 1 { return CurrentUserProfile{}, fmt.Errorf("expected 1 user result for %s, got %d", currentUserID, len(resp.Data.GetMemberResults.Results)) diff --git a/pkg/twittermeow/client.go b/pkg/twittermeow/client.go index 263b236..6144a8c 100644 --- a/pkg/twittermeow/client.go +++ b/pkg/twittermeow/client.go @@ -211,7 +211,11 @@ func (c *Client) LoadMessagesPage(ctx context.Context) (CurrentUserProfile, erro profile, err := c.GetCurrentUserProfile(ctx) if err != nil { - return CurrentUserProfile{}, fmt.Errorf("failed to fetch current user profile after loading messages page: %w", err) + if IsAuthError(err) { + return CurrentUserProfile{}, err + } + c.Logger.Warn().Err(err).Msg("Failed to fetch current user profile after loading messages page") + profile = CurrentUserProfile{ID: c.GetCurrentUserID()} } c.session.InitializedAt = time.Now() diff --git a/pkg/twittermeow/methods/html_test.go b/pkg/twittermeow/methods/html_test.go index 1818bfc..6cf0d26 100644 --- a/pkg/twittermeow/methods/html_test.go +++ b/pkg/twittermeow/methods/html_test.go @@ -1,7 +1,10 @@ package methods import ( + "io" + "net/http" "testing" + "time" ) func TestParseOndemandSURLFromScript(t *testing.T) { @@ -31,3 +34,55 @@ func TestParseOndemandSURLFromScript(t *testing.T) { }) } } + +func TestParseOndemandSURLFromLiveXBootstrap(t *testing.T) { + client := &http.Client{Timeout: 20 * time.Second} + req, err := http.NewRequest(http.MethodGet, "https://x.com/", nil) + if err != nil { + t.Fatalf("failed to create request: %v", err) + } + req.Header.Set("User-Agent", "Mozilla/5.0") + req.Header.Set("Accept-Language", "en-US,en;q=0.9") + + resp, err := client.Do(req) + if err != nil { + t.Skipf("failed to fetch x.com: %v", err) + } + defer resp.Body.Close() + + html, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatalf("failed to read x.com response: %v", err) + } + + ondemandURL := ParseOndemandSURLFromScript(html) + if ondemandURL == "" { + mainScriptURL := ParseMainScriptURL(string(html)) + if mainScriptURL == "" { + t.Fatalf("failed to locate main script URL from x.com response") + } + + req, err = http.NewRequest(http.MethodGet, mainScriptURL, nil) + if err != nil { + t.Fatalf("failed to create main script request: %v", err) + } + req.Header.Set("User-Agent", "Mozilla/5.0") + req.Header.Set("Accept-Language", "en-US,en;q=0.9") + + resp, err = client.Do(req) + if err != nil { + t.Fatalf("failed to fetch main script: %v", err) + } + defer resp.Body.Close() + + script, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatalf("failed to read main script response: %v", err) + } + ondemandURL = ParseOndemandSURLFromScript(script) + } + + if ondemandURL == "" { + t.Fatalf("failed to resolve ondemand.s URL from live x.com bootstrap") + } +} diff --git a/pkg/twittermeow/stream_client.go b/pkg/twittermeow/stream_client.go index 61c37a5..a977eda 100644 --- a/pkg/twittermeow/stream_client.go +++ b/pkg/twittermeow/stream_client.go @@ -43,6 +43,9 @@ func (c *Client) newStreamClient() *StreamClient { } func (sc *StreamClient) startOrUpdateEventStream(ctx context.Context, conversationID string) { + if ctx == nil { + ctx = context.TODO() + } ctx = sc.client.Logger.With().Str("action", "event stream").Logger().WithContext(ctx) if sc.conversationID == "" { sc.conversationID = conversationID diff --git a/pkg/twittermeow/xchat_send.go b/pkg/twittermeow/xchat_send.go index ccf411b..8cf0726 100644 --- a/pkg/twittermeow/xchat_send.go +++ b/pkg/twittermeow/xchat_send.go @@ -96,7 +96,11 @@ func (c *Client) ensureConversationToken(ctx context.Context, conversationID str return "", err } - return c.keyManager.GetConversationToken(ctx, conversationID) + token, err = c.keyManager.GetConversationToken(ctx, conversationID) + if err != nil { + return "", fmt.Errorf("get conversation token: %w", err) + } + return token, nil } func (c *Client) refreshConversationToken(ctx context.Context, conversationID string) error { From 8fe6414821f60d54de495e7fedc7b9dc0c627944 Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Mon, 30 Mar 2026 13:59:00 +0300 Subject: [PATCH 8/9] connector,twittermeow: fix context.TODO calls --- pkg/connector/client.go | 5 ++++- pkg/connector/handletwit.go | 11 ----------- pkg/twittermeow/stream_client.go | 3 --- 3 files changed, 4 insertions(+), 15 deletions(-) diff --git a/pkg/connector/client.go b/pkg/connector/client.go index f187309..e6e2094 100644 --- a/pkg/connector/client.go +++ b/pkg/connector/client.go @@ -81,7 +81,10 @@ func NewTwitterClient(login *bridgev2.UserLogin, connector *TwitterConnector, cl if !ok { return displayname } - ghost, err := tc.connector.br.GetGhostByID(context.TODO(), userID) + if ctx.Ctx == nil { + return displayname + } + ghost, err := tc.connector.br.GetGhostByID(ctx.Ctx, userID) if err != nil || len(ghost.Identifiers) < 1 { return displayname } diff --git a/pkg/connector/handletwit.go b/pkg/connector/handletwit.go index 95cff46..f753144 100644 --- a/pkg/connector/handletwit.go +++ b/pkg/connector/handletwit.go @@ -170,9 +170,6 @@ func (tc *TwitterClient) buildMemberChangeEvent( // HandleXChatEvent handles events from the XChat websocket processor. func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.TwitterEvent) bool { - if ctx == nil { - ctx = context.TODO() - } if rawEvt == nil { return true } @@ -595,10 +592,6 @@ var _ = payload.FailureType(0) // receive real-time updates via XChat WebSocket. // Returns true to continue polling, false to stop. func (tc *TwitterClient) HandlePollingEvent(ctx context.Context, evt types.TwitterEvent, inbox *response.TwitterInboxData) bool { - if ctx == nil { - ctx = context.TODO() - } - // Always cache users from inbox when available - needed for portal creation if inbox != nil { tc.updateTwitterUserInfo(ctx, inbox) @@ -789,10 +782,6 @@ func (tc *TwitterClient) markPollingChatResyncSuccess(conversationID string, now // handlePollingMessage handles a message event from REST API polling. func (tc *TwitterClient) handlePollingMessage(ctx context.Context, evt *types.Message, inbox *response.TwitterInboxData) bool { - if ctx == nil { - ctx = context.TODO() - } - isFromMe := MakeUserLoginID(evt.MessageData.SenderID) == tc.userLogin.ID portalKey := tc.MakePortalKeyFromID(evt.ConversationID) msgID := evt.ID diff --git a/pkg/twittermeow/stream_client.go b/pkg/twittermeow/stream_client.go index a977eda..61c37a5 100644 --- a/pkg/twittermeow/stream_client.go +++ b/pkg/twittermeow/stream_client.go @@ -43,9 +43,6 @@ func (c *Client) newStreamClient() *StreamClient { } func (sc *StreamClient) startOrUpdateEventStream(ctx context.Context, conversationID string) { - if ctx == nil { - ctx = context.TODO() - } ctx = sc.client.Logger.With().Str("action", "event stream").Logger().WithContext(ctx) if sc.conversationID == "" { sc.conversationID = conversationID From caa23823a20b0bdb0ab1c362639d41339a839404 Mon Sep 17 00:00:00 2001 From: Selyatin Ismet <50295732+Selyatin@users.noreply.github.com> Date: Tue, 31 Mar 2026 22:05:58 +0300 Subject: [PATCH 9/9] connector,twittermeow: fix context.TODO regression --- pkg/connector/handletwit.go | 9 --------- pkg/twittermeow/stream_client.go | 3 --- 2 files changed, 12 deletions(-) diff --git a/pkg/connector/handletwit.go b/pkg/connector/handletwit.go index ba5f367..f753144 100644 --- a/pkg/connector/handletwit.go +++ b/pkg/connector/handletwit.go @@ -408,9 +408,6 @@ func (tc *TwitterClient) HandleXChatEvent(ctx context.Context, rawEvt types.Twit return tc.userLogin.QueueRemoteEvent(portalDeleteRemoteEvent).Success case *types.ConversationNameUpdate: - if ctx == nil { - ctx = context.TODO() - } // XChat group titles are encrypted. Decrypt before forwarding to Matrix so // we don't set the room name to ciphertext. newName := evt.ConversationName @@ -595,9 +592,6 @@ var _ = payload.FailureType(0) // receive real-time updates via XChat WebSocket. // Returns true to continue polling, false to stop. func (tc *TwitterClient) HandlePollingEvent(ctx context.Context, evt types.TwitterEvent, inbox *response.TwitterInboxData) bool { - if ctx == nil { - ctx = context.TODO() - } // Always cache users from inbox when available - needed for portal creation if inbox != nil { tc.updateTwitterUserInfo(ctx, inbox) @@ -788,9 +782,6 @@ func (tc *TwitterClient) markPollingChatResyncSuccess(conversationID string, now // handlePollingMessage handles a message event from REST API polling. func (tc *TwitterClient) handlePollingMessage(ctx context.Context, evt *types.Message, inbox *response.TwitterInboxData) bool { - if ctx == nil { - ctx = context.TODO() - } isFromMe := MakeUserLoginID(evt.MessageData.SenderID) == tc.userLogin.ID portalKey := tc.MakePortalKeyFromID(evt.ConversationID) msgID := evt.ID diff --git a/pkg/twittermeow/stream_client.go b/pkg/twittermeow/stream_client.go index a977eda..61c37a5 100644 --- a/pkg/twittermeow/stream_client.go +++ b/pkg/twittermeow/stream_client.go @@ -43,9 +43,6 @@ func (c *Client) newStreamClient() *StreamClient { } func (sc *StreamClient) startOrUpdateEventStream(ctx context.Context, conversationID string) { - if ctx == nil { - ctx = context.TODO() - } ctx = sc.client.Logger.With().Str("action", "event stream").Logger().WithContext(ctx) if sc.conversationID == "" { sc.conversationID = conversationID