From 94e5a42fff105accc6c8f82c050f5ea8f72e5fe7 Mon Sep 17 00:00:00 2001 From: Anton Kozyk Date: Wed, 10 Jun 2026 15:20:50 +0300 Subject: [PATCH 1/3] Fix jid issue --- handlers.go | 46 +++++++++++++++++++++++++++++++++++++++++++--- wmiau.go | 20 ++++++++++++++++---- 2 files changed, 59 insertions(+), 7 deletions(-) diff --git a/handlers.go b/handlers.go index ee3c5d3e..9248c041 100644 --- a/handlers.go +++ b/handlers.go @@ -744,6 +744,44 @@ func (s *server) PairPhone() http.HandlerFunc { } } +// resolveSessionJID returns the best-known JID for a user. When logged in, the +func (s *server) resolveSessionJID(txtid string, waClient *whatsmeow.Client, userInfo Values) string { + jid := userInfo.Get("Jid") + + if waClient != nil && waClient.IsLoggedIn() && waClient.Store != nil && waClient.Store.ID != nil { + storeJID := waClient.Store.ID.ToNonAD() + if storeJID.Server == types.HiddenUserServer { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + pn, err := getCachedPNForLID(ctx, waClient, storeJID) + cancel() + if err == nil && !pn.IsEmpty() { + storeJID = pn.ToNonAD() + } + } + resolved := storeJID.String() + if resolved != "" { + jid = resolved + if resolved != userInfo.Get("Jid") { + token := userInfo.Get("Token") + if token != "" { + v := updateUserInfo(userInfo, "Jid", resolved) + userinfocache.Set(token, v, cache.NoExpiration) + } + if _, err := s.db.Exec("UPDATE users SET jid=$1 WHERE id=$2", resolved, txtid); err != nil { + log.Warn().Err(err).Str("user_id", txtid).Msg("Failed to persist resolved JID") + } + } + } + } else if jid == "" { + var dbJid string + if err := s.db.QueryRow("SELECT jid FROM users WHERE id = $1", txtid).Scan(&dbJid); err == nil && dbJid != "" { + jid = dbJid + } + } + + return jid +} + // Gets Connected and LoggedIn Status func (s *server) GetStatus() http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { @@ -763,8 +801,10 @@ func (s *server) GetStatus() http.HandlerFunc { txtid := userInfo.Get("Id") - isConnected := clientManager.GetWhatsmeowClient(txtid).IsConnected() - isLoggedIn := clientManager.GetWhatsmeowClient(txtid).IsLoggedIn() + waClient := clientManager.GetWhatsmeowClient(txtid) + isConnected := waClient != nil && waClient.IsConnected() + isLoggedIn := waClient != nil && waClient.IsLoggedIn() + jid := s.resolveSessionJID(txtid, waClient, userInfo) var proxyURL string s.db.QueryRow("SELECT proxy_url FROM users WHERE id = $1", txtid).Scan(&proxyURL) @@ -821,7 +861,7 @@ func (s *server) GetStatus() http.HandlerFunc { "connected": isConnected, "loggedIn": isLoggedIn, "token": userInfo.Get("Token"), - "jid": userInfo.Get("Jid"), + "jid": jid, "webhook": userInfo.Get("Webhook"), "events": userInfo.Get("Events"), "proxy_url": userInfo.Get("Proxy"), diff --git a/wmiau.go b/wmiau.go index feb69e5a..c1a4604b 100644 --- a/wmiau.go +++ b/wmiau.go @@ -718,11 +718,23 @@ func (mycli *MyClient) myEventHandler(rawEvt interface{}) { log.Error().Err(err).Msg(sqlStmt) return } + if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { + connectedJID := mycli.WAClient.Store.ID.ToNonAD().String() + if _, err := mycli.db.Exec(`UPDATE users SET jid=$1 WHERE id=$2`, connectedJID, mycli.userID); err != nil { + log.Warn().Err(err).Str("user_id", mycli.userID).Msg("Failed to persist JID on connect") + } else if myuserinfo, found := userinfocache.Get(mycli.token); found { + v := updateUserInfo(myuserinfo, "Jid", connectedJID) + userinfocache.Set(mycli.token, v, cache.NoExpiration) + } + } case *events.PairSuccess: log.Info().Str("userid", mycli.userID).Str("token", mycli.token).Str("ID", evt.ID.String()).Str("BusinessName", evt.BusinessName).Str("Platform", evt.Platform).Msg("QR Pair Success") - jid := evt.ID + jidStr := evt.ID.String() + if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { + jidStr = mycli.WAClient.Store.ID.ToNonAD().String() + } sqlStmt := `UPDATE users SET jid=$1 WHERE id=$2` - _, err := mycli.db.Exec(sqlStmt, jid, mycli.userID) + _, err := mycli.db.Exec(sqlStmt, jidStr, mycli.userID) if err != nil { log.Error().Err(err).Msg(sqlStmt) return @@ -737,9 +749,9 @@ func (mycli *MyClient) myEventHandler(rawEvt interface{}) { } else { txtid = myuserinfo.(Values).Get("Id") token := myuserinfo.(Values).Get("Token") - v := updateUserInfo(myuserinfo, "Jid", fmt.Sprintf("%s", jid)) + v := updateUserInfo(myuserinfo, "Jid", jidStr) userinfocache.Set(token, v, cache.NoExpiration) - log.Info().Str("jid", jid.String()).Str("userid", txtid).Str("token", token).Msg("User information set") + log.Info().Str("jid", jidStr).Str("userid", txtid).Str("token", token).Msg("User information set") } // Check if automatic history sync is enabled and trigger it after QR code is scanned From b7741e6f90d21c7a68f4e77f80e9672ca98fad89 Mon Sep 17 00:00:00 2001 From: Anton Kozyk Date: Mon, 10 Aug 2026 15:51:02 +0300 Subject: [PATCH 2/3] Jid fix 2 --- handlers.go | 20 +- jid_resolve_test.go | 72 + wmiau.go | 3284 ++++++++++++++++++++++--------------------- 3 files changed, 1784 insertions(+), 1592 deletions(-) create mode 100644 jid_resolve_test.go diff --git a/handlers.go b/handlers.go index 9248c041..1f8f2d27 100644 --- a/handlers.go +++ b/handlers.go @@ -259,12 +259,24 @@ func (s *server) Connect() http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { - webhook := r.Context().Value("userinfo").(Values).Get("Webhook") - jid := r.Context().Value("userinfo").(Values).Get("Jid") - txtid := r.Context().Value("userinfo").(Values).Get("Id") - token := r.Context().Value("userinfo").(Values).Get("Token") + userInfo := r.Context().Value("userinfo").(Values) + webhook := userInfo.Get("Webhook") + jid := userInfo.Get("Jid") + txtid := userInfo.Get("Id") + token := userInfo.Get("Token") eventstring := "" + // Always prefer the DB jid over a stale in-memory cache entry so reconnect + // can find the whatsmeow device even after a process restart. + var dbJid string + if err := s.db.QueryRow("SELECT jid FROM users WHERE id = $1", txtid).Scan(&dbJid); err == nil && dbJid != "" { + jid = dbJid + if jid != userInfo.Get("Jid") { + v := updateUserInfo(userInfo, "Jid", jid) + userinfocache.Set(token, v, cache.NoExpiration) + } + } + // Decodes request BODY looking for events to subscribe decoder := json.NewDecoder(r.Body) var t connectStruct diff --git a/jid_resolve_test.go b/jid_resolve_test.go new file mode 100644 index 00000000..9fa1173d --- /dev/null +++ b/jid_resolve_test.go @@ -0,0 +1,72 @@ +package main + +import ( + "testing" + + "go.mau.fi/whatsmeow/types" +) + +func TestJidUserKey(t *testing.T) { + tests := []struct { + name string + jid string + want string + }{ + {"empty", "", ""}, + {"plain user", "380996987502@s.whatsapp.net", "380996987502"}, + {"with device suffix", "380996987502:8@s.whatsapp.net", "380996987502"}, + {"lid style", "1234567890:1@lid", "1234567890"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := jidUserKey(tt.jid); got != tt.want { + t.Fatalf("jidUserKey(%q) = %q, want %q", tt.jid, got, tt.want) + } + }) + } +} + +func TestJidLookupCandidates(t *testing.T) { + candidates := jidLookupCandidates("380996987502:8@s.whatsapp.net") + if len(candidates) < 2 { + t.Fatalf("expected multiple candidates, got %d", len(candidates)) + } + + seen := make(map[string]struct{}) + for _, c := range candidates { + seen[c.String()] = struct{}{} + } + + want := []string{ + "380996987502:8@s.whatsapp.net", + "380996987502@s.whatsapp.net", + } + for _, w := range want { + if _, ok := seen[w]; !ok { + t.Fatalf("missing candidate %q in %v", w, candidates) + } + } + + if got := jidUserKey(candidates[0].String()); got != "380996987502" { + t.Fatalf("first candidate key = %q, want 380996987502", got) + } +} + +func TestJidLookupCandidatesInvalid(t *testing.T) { + if got := jidLookupCandidates(""); got != nil { + t.Fatalf("empty jid should return nil, got %v", got) + } + if got := jidLookupCandidates("@s.whatsapp.net"); got != nil { + t.Fatalf("jid without user should return nil, got %v", got) + } +} + +func TestJidLookupCandidatesPhoneOnly(t *testing.T) { + candidates := jidLookupCandidates("380996987502") + if len(candidates) != 1 { + t.Fatalf("phone-only input should yield one candidate, got %d", len(candidates)) + } + if candidates[0].Server != types.DefaultUserServer { + t.Fatalf("expected default user server, got %q", candidates[0].Server) + } +} diff --git a/wmiau.go b/wmiau.go index c1a4604b..deadf3f8 100644 --- a/wmiau.go +++ b/wmiau.go @@ -1,1588 +1,1696 @@ -package main - -import ( - "context" - "crypto/sha256" - "crypto/tls" - "encoding/base64" - "encoding/json" - "errors" - "fmt" - "net/http" - "net/url" - "os" - "path/filepath" - "runtime/debug" - "strconv" - "strings" - "time" - - "github.com/go-resty/resty/v2" - "github.com/jmoiron/sqlx" - "github.com/mdp/qrterminal/v3" - "github.com/patrickmn/go-cache" - "github.com/rs/zerolog/log" - "github.com/skip2/go-qrcode" - "go.mau.fi/whatsmeow" - "go.mau.fi/whatsmeow/appstate" - "go.mau.fi/whatsmeow/proto/waCompanionReg" - "go.mau.fi/whatsmeow/store" - "go.mau.fi/whatsmeow/types" - "go.mau.fi/whatsmeow/types/events" - waLog "go.mau.fi/whatsmeow/util/log" - "golang.org/x/net/proxy" -) - -// db field declaration as *sqlx.DB -type MyClient struct { - WAClient *whatsmeow.Client - eventHandlerID uint32 - userID string - token string - db *sqlx.DB - s *server -} - -// safeGo runs fn in a new goroutine with a defer recover so a panic inside -// fire-and-forget side-effects (webhook delivery, MQ push) cannot crash -// the whole process. Losing one delivery is preferable to taking wuzapi -// down for every connected user. -func safeGo(name string, fn func()) { - go func() { - defer func() { - if r := recover(); r != nil { - log.Error(). - Str("goroutine", name). - Interface("panic", r). - Str("stack", string(debug.Stack())). - Msg("panic recovered in goroutine") - } - }() - fn() - }() -} - -// ensureS3ClientForUser loads S3 config from DB and initializes client if not already present (lazy init for reconnect-after-restart) -func ensureS3ClientForUser(userID string) { - GetS3Manager().EnsureClientFromDB(userID) -} - -func sendToGlobalWebHook(jsonData []byte, token string, userID string) { - jsonDataStr := string(jsonData) - - instance_name := "" - userinfo, found := userinfocache.Get(token) - if found { - instance_name = userinfo.(Values).Get("Name") - } - - if *globalWebhook != "" { - log.Info().Str("url", *globalWebhook).Msg("Calling global webhook") - // Add extra information for the global webhook - globalData := map[string]string{ - "jsonData": jsonDataStr, - "userID": userID, - "instanceName": instance_name, - } - callHookWithHmac(*globalWebhook, globalData, userID, globalHMACKeyEncrypted) - } -} - -func sendToUserWebHook(webhookurl string, path string, jsonData []byte, userID string, token string) { - sendToUserWebHookWithHmac(webhookurl, path, jsonData, userID, token, nil) -} - -func sendToUserWebHookWithHmac(webhookurl string, path string, jsonData []byte, userID string, token string, encryptedHmacKey []byte) { - - instance_name := "" - userinfo, found := userinfocache.Get(token) - if found { - instance_name = userinfo.(Values).Get("Name") - } - data := map[string]string{ - "jsonData": string(jsonData), - "userID": userID, - "instanceName": instance_name, - } - - log.Debug().Interface("webhookData", data).Msg("Data being sent to webhook") - - if webhookurl != "" { - log.Info().Str("url", webhookurl).Msg("Calling user webhook") - - if path == "" { - safeGo("callHookWithHmac", func() { callHookWithHmac(webhookurl, data, userID, encryptedHmacKey) }) - } else { - if err := callHookFileWithHmac(webhookurl, data, userID, path, encryptedHmacKey); err != nil { - log.Error().Err(err).Msg("Error calling hook file") - } - } - } else { - log.Warn().Str("userid", userID).Msg("No webhook set for user") - } -} - -func updateAndGetUserSubscriptions(mycli *MyClient) ([]string, error) { - // Get updated events from cache/database - currentEvents := "" - userinfo2, found2 := userinfocache.Get(mycli.token) - if found2 { - currentEvents = userinfo2.(Values).Get("Events") - } else { - // If not in cache, get from database - if err := mycli.db.Get(¤tEvents, "SELECT events FROM users WHERE id=$1", mycli.userID); err != nil { - log.Warn().Err(err).Str("userID", mycli.userID).Msg("Could not get events from DB") - return nil, err // Propagate the error - } - } - - // Update client subscriptions if changed - eventarray := strings.Split(currentEvents, ",") - var subscribedEvents []string - if len(eventarray) == 1 && eventarray[0] == "" { - subscribedEvents = []string{} - } else { - for _, arg := range eventarray { - arg = strings.TrimSpace(arg) - if arg != "" && Find(supportedEventTypes, arg) { - subscribedEvents = append(subscribedEvents, arg) - } - } - } - - return subscribedEvents, nil -} - -func getUserWebhookUrl(token string) string { - webhookurl := "" - myuserinfo, found := userinfocache.Get(token) - if !found { - log.Warn().Str("token", token).Msg("Could not call webhook as there is no user for this token") - } else { - webhookurl = myuserinfo.(Values).Get("Webhook") - } - return webhookurl -} - -func sendEventWithWebHook(mycli *MyClient, postmap map[string]interface{}, path string) { - webhookurl := getUserWebhookUrl(mycli.token) - - // Get updated events from cache/database - subscribedEvents, err := updateAndGetUserSubscriptions(mycli) - if err != nil { - return - } - - eventType, ok := postmap["type"].(string) - if !ok { - log.Error().Msg("Event type is not a string in postmap") - return - } - - // Log subscription details for debugging - log.Debug(). - Str("userID", mycli.userID). - Str("eventType", eventType). - Strs("subscribedEvents", subscribedEvents). - Msg("Checking event subscription") - - // Check if the current event is in the subscriptions - checkIfSubscribedInEvent := checkIfSubscribedToEvent(subscribedEvents, postmap["type"].(string), mycli.userID) - if !checkIfSubscribedInEvent { - return - } - - // In stdio mode, send as JSON-RPC notification instead of HTTP webhook - if mycli.s != nil && mycli.s.mode == Stdio { - mycli.s.SendNotification(eventType, postmap) - return - } - - // Prepare webhook data - jsonData, err := json.Marshal(postmap) - if err != nil { - log.Error().Err(err).Msg("Failed to marshal postmap to JSON") - return - } - - // Get HMAC key for this user - var encryptedHmacKey []byte - if userinfo, found := userinfocache.Get(mycli.token); found { - encryptedB64 := userinfo.(Values).Get("HmacKeyEncrypted") - if encryptedB64 != "" { - var err error - encryptedHmacKey, err = base64.StdEncoding.DecodeString(encryptedB64) - if err != nil { - log.Error().Err(err).Msg("Failed to decode HMAC key from cache") - } - } - } - - sendToUserWebHookWithHmac(webhookurl, path, jsonData, mycli.userID, mycli.token, encryptedHmacKey) - - // Get global webhook if configured - safeGo("sendToGlobalWebHook", func() { sendToGlobalWebHook(jsonData, mycli.token, mycli.userID) }) - - safeGo("sendToGlobalRabbit", func() { sendToGlobalRabbit(jsonData, mycli.token, mycli.userID) }) -} - -func checkIfSubscribedToEvent(subscribedEvents []string, eventType string, userId string) bool { - if !Find(subscribedEvents, eventType) && !Find(subscribedEvents, "All") { - log.Warn(). - Str("type", eventType). - Strs("subscribedEvents", subscribedEvents). - Str("userID", userId). - Msg("Skipping webhook. Not subscribed for this type") - return false - } - return true -} - -// Connects to Whatsapp Websocket on server startup if last state was connected -func (s *server) connectOnStartup() { - rows, err := s.db.Queryx("SELECT id,name,token,jid,webhook,events,proxy_url,CASE WHEN s3_enabled THEN 'true' ELSE 'false' END AS s3_enabled,media_delivery,COALESCE(history, 0) as history,hmac_key FROM users WHERE connected=1") - if err != nil { - log.Error().Err(err).Msg("DB Problem") - return - } - defer rows.Close() - for rows.Next() { - txtid := "" - token := "" - jid := "" - name := "" - webhook := "" - events := "" - proxy_url := "" - s3_enabled := "" - media_delivery := "" - var history int - var hmac_key []byte - err = rows.Scan(&txtid, &name, &token, &jid, &webhook, &events, &proxy_url, &s3_enabled, &media_delivery, &history, &hmac_key) - if err != nil { - log.Error().Err(err).Msg("DB Problem") - return - } else { - hmacKeyEncrypted := "" - if len(hmac_key) > 0 { - hmacKeyEncrypted = base64.StdEncoding.EncodeToString(hmac_key) - } - - log.Info().Str("token", token).Msg("Connect to Whatsapp on startup") - v := Values{map[string]string{ - "Id": txtid, - "Name": name, - "Jid": jid, - "Webhook": webhook, - "Token": token, - "Proxy": proxy_url, - "Events": events, - "S3Enabled": s3_enabled, - "MediaDelivery": media_delivery, - "History": fmt.Sprintf("%d", history), - "HmacKeyEncrypted": hmacKeyEncrypted, - }} - userinfocache.Set(token, v, cache.NoExpiration) - // Gets and set subscription to webhook events - eventarray := strings.Split(events, ",") - - var subscribedEvents []string - if len(eventarray) == 1 && eventarray[0] == "" { - subscribedEvents = []string{} - } else { - for _, arg := range eventarray { - if !Find(supportedEventTypes, arg) { - log.Warn().Str("Type", arg).Msg("Event type discarded") - continue - } - if !Find(subscribedEvents, arg) { - subscribedEvents = append(subscribedEvents, arg) - } - } - - } - eventstring := strings.Join(subscribedEvents, ",") - log.Info().Str("events", eventstring).Str("jid", jid).Msg("Attempt to connect") - kill := make(chan bool, 1) - setKillChannel(txtid, kill) - go s.startClient(txtid, jid, token, kill) - - // Initialize S3 client if configured - go func(userID string) { - GetS3Manager().EnsureClientFromDB(userID) - }(txtid) - } - } - err = rows.Err() - if err != nil { - log.Error().Err(err).Msg("DB Problem") - } -} - -func parseJID(arg string) (types.JID, bool) { - if arg[0] == '+' { - arg = arg[1:] - } - if !strings.ContainsRune(arg, '@') { - return types.NewJID(arg, types.DefaultUserServer), true - } else { - recipient, err := types.ParseJID(arg) - if err != nil { - log.Error().Err(err).Msg("Invalid JID") - return recipient, false - } else if recipient.User == "" { - log.Error().Err(err).Msg("Invalid JID no server specified") - return recipient, false - } - return recipient, true - } -} - -// getPlatformTypeEnum converts a platform type string to the corresponding DeviceProps enum -// Returns DESKTOP as default if the string doesn't match any known type -func getPlatformTypeEnum(platformType string) *waCompanionReg.DeviceProps_PlatformType { - platformType = strings.ToUpper(strings.TrimSpace(platformType)) - - switch platformType { - case "UNKNOWN": - return waCompanionReg.DeviceProps_UNKNOWN.Enum() - case "CHROME": - return waCompanionReg.DeviceProps_CHROME.Enum() - case "FIREFOX": - return waCompanionReg.DeviceProps_FIREFOX.Enum() - case "IE": - return waCompanionReg.DeviceProps_IE.Enum() - case "OPERA": - return waCompanionReg.DeviceProps_OPERA.Enum() - case "SAFARI": - return waCompanionReg.DeviceProps_SAFARI.Enum() - case "EDGE": - return waCompanionReg.DeviceProps_EDGE.Enum() - case "DESKTOP": - return waCompanionReg.DeviceProps_DESKTOP.Enum() - case "IPAD": - return waCompanionReg.DeviceProps_IPAD.Enum() - case "ANDROID_TABLET": - return waCompanionReg.DeviceProps_ANDROID_TABLET.Enum() - case "OHANA": - return waCompanionReg.DeviceProps_OHANA.Enum() - case "ALOHA": - return waCompanionReg.DeviceProps_ALOHA.Enum() - case "CATALINA": - return waCompanionReg.DeviceProps_CATALINA.Enum() - case "TCL_TV": - return waCompanionReg.DeviceProps_TCL_TV.Enum() - case "IOS_PHONE": - return waCompanionReg.DeviceProps_IOS_PHONE.Enum() - case "IOS_CATALYST": - return waCompanionReg.DeviceProps_IOS_CATALYST.Enum() - case "ANDROID_PHONE": - return waCompanionReg.DeviceProps_ANDROID_PHONE.Enum() - case "ANDROID_AMBIGUOUS": - return waCompanionReg.DeviceProps_ANDROID_AMBIGUOUS.Enum() - case "WEAR_OS": - return waCompanionReg.DeviceProps_WEAR_OS.Enum() - case "AR_WRIST": - return waCompanionReg.DeviceProps_AR_WRIST.Enum() - case "AR_DEVICE": - return waCompanionReg.DeviceProps_AR_DEVICE.Enum() - case "UWP": - return waCompanionReg.DeviceProps_UWP.Enum() - case "VR": - return waCompanionReg.DeviceProps_VR.Enum() - default: - log.Warn().Str("platformType", platformType).Msg("Unknown platform type, defaulting to DESKTOP") - return waCompanionReg.DeviceProps_DESKTOP.Enum() - } -} - -func (s *server) startClient(userID string, textjid string, token string, kill chan bool) { - log.Info().Str("userid", userID).Str("jid", textjid).Msg("Starting websocket connection to Whatsapp") - - // Connection retry constants - const maxConnectionRetries = 3 - const connectionRetryBaseWait = 5 * time.Second - - var deviceStore *store.Device - var err error - - // First handle the device store initialization - if textjid != "" { - jid, _ := parseJID(textjid) - deviceStore, err = container.GetDevice(context.Background(), jid) - if err != nil { - log.Error().Err(err).Msg("Failed to get device") - deviceStore = container.NewDevice() - } - } else { - log.Warn().Msg("No jid found. Creating new device") - deviceStore = container.NewDevice() - } - - if deviceStore == nil { - log.Warn().Msg("No store found. Creating new one") - deviceStore = container.NewDevice() - } - - clientLog := waLog.Stdout("Client", *waDebug, *colorOutput) - - // Create the client with initialized deviceStore - var client *whatsmeow.Client - if *waDebug != "" { - client = whatsmeow.NewClient(deviceStore, clientLog) - } else { - client = whatsmeow.NewClient(deviceStore, nil) - } - - // Now we can use the client with the manager - clientManager.SetWhatsmeowClient(userID, client) - - store.DeviceProps.PlatformType = getPlatformTypeEnum(*platformType) - store.DeviceProps.Os = osName - - mycli := MyClient{ - WAClient: client, - eventHandlerID: 1, - userID: userID, - token: token, - db: s.db, - s: s, - } - mycli.eventHandlerID = mycli.WAClient.AddEventHandler(mycli.myEventHandler) - - // Store the MyClient in clientManager - clientManager.SetMyClient(userID, &mycli) - - httpClient := resty.New() - httpClient.SetRedirectPolicy(resty.FlexibleRedirectPolicy(15)) - if *waDebug == "DEBUG" { - httpClient.SetDebug(true) - } - httpClient.SetTimeout(30 * time.Second) - httpClient.SetTLSClientConfig(&tls.Config{InsecureSkipVerify: true}) - httpClient.OnError(func(req *resty.Request, err error) { - if v, ok := err.(*resty.ResponseError); ok { - // v.Response contains the last response from the server - // v.Err contains the original error - log.Debug().Str("response", v.Response.String()).Msg("resty error") - log.Error().Err(v.Err).Msg("resty error") - } - }) - - // Set proxy if defined in DB (assumes users table contains proxy_url column) - var proxyURL string - err = s.db.Get(&proxyURL, "SELECT proxy_url FROM users WHERE id=$1", userID) - if err == nil && proxyURL != "" { - parsed, perr := url.Parse(proxyURL) - if perr != nil { - log.Warn().Err(perr).Str("proxy", proxyURL).Msg("Invalid proxy URL, skipping proxy setup") - } else { - - log.Info().Str("proxy", proxyURL).Msg("Configuring proxy") - - if parsed.Scheme == "socks5" || parsed.Scheme == "socks5h" { - dialer, derr := proxy.FromURL(parsed, nil) - if derr != nil { - log.Warn().Err(derr).Str("proxy", proxyURL).Msg("Failed to build SOCKS proxy dialer, skipping proxy setup") - } else { - httpClient.SetProxy(proxyURL) - client.SetSOCKSProxy(dialer, whatsmeow.SetProxyOptions{}) - log.Info().Msg("SOCKS proxy configured successfully") - } - } else { - httpClient.SetProxy(proxyURL) - client.SetProxyAddress(parsed.String(), whatsmeow.SetProxyOptions{}) - log.Info().Msg("HTTP/HTTPS proxy configured successfully") - } - } - } - clientManager.SetHTTPClient(userID, httpClient) - - // Initialize S3 client if configured (needed when user reconnects after container restart - connectOnStartup only runs for connected=1) - GetS3Manager().EnsureClientFromDB(userID) - - if client.Store.ID == nil { - // No ID stored, new login - qrChan, err := client.GetQRChannel(context.Background()) - if err != nil { - // This error means that we're already logged in, so ignore it. - if !errors.Is(err, whatsmeow.ErrQRStoreContainsID) { - log.Error().Err(err).Msg("Failed to get QR channel") - return - } - } else { - err = client.Connect() // Must connect to generate QR code - if err != nil { - log.Error().Err(err).Msg("Failed to connect client") - return - } - - myuserinfo, found := userinfocache.Get(token) - - for evt := range qrChan { - if evt.Event == "code" { - // Display QR code in terminal (useful for testing/developing) - // Skip in stdio mode to avoid breaking JSON-RPC - if *logType != "json" && s.mode != Stdio { - qrterminal.GenerateHalfBlock(evt.Code, qrterminal.L, os.Stdout) - fmt.Println("QR code:\n", evt.Code) - } - // Store encoded/embeded base64 QR on database for retrieval with the /qr endpoint - image, _ := qrcode.Encode(evt.Code, qrcode.Medium, 256) - base64qrcode := "data:image/png;base64," + base64.StdEncoding.EncodeToString(image) - sqlStmt := `UPDATE users SET qrcode=$1 WHERE id=$2` - _, err := s.db.Exec(sqlStmt, base64qrcode, userID) - if err != nil { - log.Error().Err(err).Msg(sqlStmt) - } else { - if found { - v := updateUserInfo(myuserinfo, "Qrcode", base64qrcode) - userinfocache.Set(token, v, cache.NoExpiration) - log.Info().Str("qrcode", base64qrcode).Msg("update cache userinfo with qr code") - } - } - - //send QR code with webhook - postmap := make(map[string]interface{}) - postmap["event"] = evt.Event - postmap["qrCodeBase64"] = base64qrcode - postmap["type"] = "QR" - - sendEventWithWebHook(&mycli, postmap, "") - - } else if evt.Event == "timeout" { - // Clear QR code from DB on timeout - // Send webhook notifying QR timeout before cleanup - postmap := make(map[string]interface{}) - postmap["event"] = evt.Event - postmap["type"] = "QRTimeout" - sendEventWithWebHook(&mycli, postmap, "") - - sqlStmt := `UPDATE users SET qrcode='' WHERE id=$1` - _, err := s.db.Exec(sqlStmt, userID) - if err != nil { - log.Error().Err(err).Msg(sqlStmt) - } else { - if found { - v := updateUserInfo(myuserinfo, "Qrcode", "") - userinfocache.Set(token, v, cache.NoExpiration) - } - } - log.Warn().Msg("QR timeout killing channel") - clientManager.DeleteWhatsmeowClient(userID) - clientManager.DeleteMyClient(userID) - clientManager.DeleteHTTPClient(userID) - signalKill(userID) - } else if evt.Event == "success" { - log.Info().Msg("QR pairing ok!") - // Clear QR code after pairing - sqlStmt := `UPDATE users SET qrcode='', connected=1 WHERE id=$1` - _, err := s.db.Exec(sqlStmt, userID) - if err != nil { - log.Error().Err(err).Msg(sqlStmt) - } else { - if found { - v := updateUserInfo(myuserinfo, "Qrcode", "") - userinfocache.Set(token, v, cache.NoExpiration) - } - } - } else { - log.Info().Str("event", evt.Event).Msg("Login event") - } - } - } - - } else { - // Already logged in, just connect - log.Info().Msg("Already logged in, just connect") - - // Retry logic with linear backoff - var lastErr error - - for attempt := 0; attempt < maxConnectionRetries; attempt++ { - if attempt > 0 { - waitTime := time.Duration(attempt) * connectionRetryBaseWait - log.Warn(). - Int("attempt", attempt+1). - Int("max_retries", maxConnectionRetries). - Dur("wait_time", waitTime). - Msg("Retrying connection after delay") - time.Sleep(waitTime) - } - - err = client.Connect() - if err == nil { - log.Info(). - Int("attempt", attempt+1). - Msg("Successfully connected to WhatsApp") - break - } - - lastErr = err - log.Warn(). - Err(err). - Int("attempt", attempt+1). - Int("max_retries", maxConnectionRetries). - Msg("Failed to connect to WhatsApp") - } - - if lastErr != nil { - log.Error(). - Err(lastErr). - Str("userid", userID). - Int("attempts", maxConnectionRetries). - Msg("Failed to connect to WhatsApp after all retry attempts") - - clientManager.DeleteWhatsmeowClient(userID) - clientManager.DeleteMyClient(userID) - clientManager.DeleteHTTPClient(userID) - - sqlStmt := `UPDATE users SET qrcode='', connected=0 WHERE id=$1` - _, dbErr := s.db.Exec(sqlStmt, userID) - if dbErr != nil { - log.Error().Err(dbErr).Msg("Failed to update user status after connection error") - } - - // Use the existing mycli instance from outer scope - postmap := make(map[string]interface{}) - postmap["event"] = "ConnectFailure" - postmap["error"] = lastErr.Error() - postmap["type"] = "ConnectFailure" - postmap["attempts"] = maxConnectionRetries - postmap["reason"] = "Failed to connect after retry attempts" - sendEventWithWebHook(&mycli, postmap, "") - - return - } - } - - // Keep the session goroutine alive until a kill signal arrives. Block on the - // channel (passed in directly, so this goroutine always owns its own channel - // even if a reconnect replaces the map entry) instead of polling — this parks - // the goroutine with zero CPU and no per-second mutex access. - <-kill - log.Info().Str("userid", userID).Msg("Received kill signal") - client.Disconnect() - clientManager.DeleteWhatsmeowClient(userID) - clientManager.DeleteMyClient(userID) - clientManager.DeleteHTTPClient(userID) - if _, err := s.db.Exec(`UPDATE users SET qrcode='', connected=0 WHERE id=$1`, userID); err != nil { - log.Error().Err(err).Msg("failed to mark user disconnected on kill") - } - deleteKillChannel(userID, kill) -} - -func fileToBase64(filepath string) (string, string, error) { - data, err := os.ReadFile(filepath) - if err != nil { - return "", "", err - } - mimeType := http.DetectContentType(data) - return base64.StdEncoding.EncodeToString(data), mimeType, nil -} - -func (mycli *MyClient) myEventHandler(rawEvt interface{}) { - txtid := mycli.userID - postmap := make(map[string]interface{}) - postmap["event"] = rawEvt - dowebhook := 0 - path := "" - - switch evt := rawEvt.(type) { - case *events.AppStateSyncComplete: - if len(mycli.WAClient.Store.PushName) > 0 && evt.Name == appstate.WAPatchCriticalBlock { - err := mycli.WAClient.SendPresence(context.Background(), types.PresenceAvailable) - if err != nil { - log.Warn().Err(err).Msg("Failed to send available presence") - } else { - log.Info().Msg("Marked self as available") - } - } - case *events.Connected, *events.PushNameSetting: - postmap["type"] = "Connected" - dowebhook = 1 - if len(mycli.WAClient.Store.PushName) == 0 { - break - } - // Send presence available when connecting and when the pushname is changed. - // This makes sure that outgoing messages always have the right pushname. - err := mycli.WAClient.SendPresence(context.Background(), types.PresenceAvailable) - if err != nil { - log.Warn().Err(err).Msg("Failed to send available presence") - } else { - log.Info().Msg("Marked self as available") - } - sqlStmt := `UPDATE users SET connected=1 WHERE id=$1` - _, err = mycli.db.Exec(sqlStmt, mycli.userID) - if err != nil { - log.Error().Err(err).Msg(sqlStmt) - return - } - if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { - connectedJID := mycli.WAClient.Store.ID.ToNonAD().String() - if _, err := mycli.db.Exec(`UPDATE users SET jid=$1 WHERE id=$2`, connectedJID, mycli.userID); err != nil { - log.Warn().Err(err).Str("user_id", mycli.userID).Msg("Failed to persist JID on connect") - } else if myuserinfo, found := userinfocache.Get(mycli.token); found { - v := updateUserInfo(myuserinfo, "Jid", connectedJID) - userinfocache.Set(mycli.token, v, cache.NoExpiration) - } - } - case *events.PairSuccess: - log.Info().Str("userid", mycli.userID).Str("token", mycli.token).Str("ID", evt.ID.String()).Str("BusinessName", evt.BusinessName).Str("Platform", evt.Platform).Msg("QR Pair Success") - jidStr := evt.ID.String() - if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { - jidStr = mycli.WAClient.Store.ID.ToNonAD().String() - } - sqlStmt := `UPDATE users SET jid=$1 WHERE id=$2` - _, err := mycli.db.Exec(sqlStmt, jidStr, mycli.userID) - if err != nil { - log.Error().Err(err).Msg(sqlStmt) - return - } - - postmap["type"] = "PairSuccess" - dowebhook = 1 - - myuserinfo, found := userinfocache.Get(mycli.token) - if !found { - log.Warn().Msg("No user info cached on pairing?") - } else { - txtid = myuserinfo.(Values).Get("Id") - token := myuserinfo.(Values).Get("Token") - v := updateUserInfo(myuserinfo, "Jid", jidStr) - userinfocache.Set(token, v, cache.NoExpiration) - log.Info().Str("jid", jidStr).Str("userid", txtid).Str("token", token).Msg("User information set") - } - - // Check if automatic history sync is enabled and trigger it after QR code is scanned - var daysToSyncHistory int - query := "SELECT COALESCE(days_to_sync_history, 0) FROM users WHERE id=$1" - query = mycli.db.Rebind(query) - err = mycli.db.Get(&daysToSyncHistory, query, mycli.userID) - if err != nil { - log.Warn().Err(err).Str("userID", mycli.userID).Msg("Failed to get days_to_sync_history from database") - } else if daysToSyncHistory > 0 { - // Trigger history sync in a goroutine to avoid blocking - // Wait a bit for the connection to be fully established - go func() { - time.Sleep(2 * time.Second) // Give WhatsApp time to fully establish connection - - log.Info(). - Str("userID", mycli.userID). - Int("days", daysToSyncHistory). - Msg("Triggering automatic history sync after QR code scan") - - // Use the SyncWhatsAppHistory logic but for a single user - // Calculate message count based on days (estimate: 15 messages per day) - count := daysToSyncHistory * 15 - if count > 500 { - count = 500 // WhatsApp limit - } - if count < 50 { - count = 50 // Minimum reasonable count - } - - // Get chats from WhatsApp (contacts and groups) - ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) - defer cancel() - - var chatJIDs []string - - // Get all contacts - contacts, err := mycli.WAClient.Store.Contacts.GetAllContacts(ctx) - if err != nil { - log.Error().Err(err).Str("userID", mycli.userID).Msg("Failed to get contacts for history sync") - } else { - for jid := range contacts { - chatJIDs = append(chatJIDs, jid.String()) - } - } - - // Get all groups - groups, err := mycli.WAClient.GetJoinedGroups(ctx) - if err != nil { - log.Error().Err(err).Str("userID", mycli.userID).Msg("Failed to get groups for history sync") - } else { - for _, group := range groups { - chatJIDs = append(chatJIDs, group.JID.String()) - } - } - - // Sync each chat with a small delay between requests - for _, chatJIDStr := range chatJIDs { - chatJID, err := types.ParseJID(chatJIDStr) - if err != nil { - log.Warn().Err(err).Str("chatJID", chatJIDStr).Msg("Failed to parse chat JID, skipping") - continue - } - - // Use the syncHistoryForChat function from handlers.go - err = mycli.s.syncHistoryForChat(context.Background(), mycli.userID, chatJID, count) - if err != nil { - log.Warn().Err(err).Str("chatJID", chatJIDStr).Msg("Failed to sync history for chat") - } else { - log.Info().Str("chatJID", chatJIDStr).Int("count", count).Msg("History sync request sent for chat") - } - - // Small delay between requests to avoid overwhelming WhatsApp - time.Sleep(100 * time.Millisecond) - } - - log.Info(). - Str("userID", mycli.userID). - Int("days", daysToSyncHistory). - Int("chatsSynced", len(chatJIDs)). - Msg("Automatic history sync completed after QR code scan") - }() - } - case *events.StreamReplaced: - log.Info().Msg("Received StreamReplaced event") - return - case *events.Message: - - var s3Config struct { - Enabled string `db:"s3_enabled"` - MediaDelivery string `db:"media_delivery"` - } - - lastMessageCache.Set(mycli.userID, &evt.Info, cache.DefaultExpiration) - myuserinfo, found := userinfocache.Get(mycli.token) - if !found { - err := mycli.db.Get(&s3Config, "SELECT CASE WHEN s3_enabled = 1 THEN 'true' ELSE 'false' END AS s3_enabled, media_delivery FROM users WHERE id = $1", txtid) - if err != nil { - log.Error().Err(err).Msg("onMessage Failed to get S3 config from DB as it was not on cache") - s3Config.Enabled = "false" - s3Config.MediaDelivery = "base64" - } - } else { - s3Config.Enabled = myuserinfo.(Values).Get("S3Enabled") - s3Config.MediaDelivery = myuserinfo.(Values).Get("MediaDelivery") - } - - // Lazy init S3 client if needed (handles reconnect-after-restart when connectOnStartup skipped this user) - if s3Config.Enabled == "true" && (s3Config.MediaDelivery == "s3" || s3Config.MediaDelivery == "both") { - ensureS3ClientForUser(txtid) - } - - postmap["type"] = "Message" - dowebhook = 1 - metaParts := []string{fmt.Sprintf("pushname: %s", evt.Info.PushName), fmt.Sprintf("timestamp: %s", evt.Info.Timestamp)} - if evt.Info.Type != "" { - metaParts = append(metaParts, fmt.Sprintf("type: %s", evt.Info.Type)) - } - if evt.Info.Category != "" { - metaParts = append(metaParts, fmt.Sprintf("category: %s", evt.Info.Category)) - } - if evt.IsViewOnce { - metaParts = append(metaParts, "view once") - } - if evt.IsViewOnce { - metaParts = append(metaParts, "ephemeral") - } - - log.Info().Str("id", evt.Info.ID).Str("source", evt.Info.SourceString()).Str("parts", strings.Join(metaParts, ", ")).Msg("Message Received") - - // If this is a poll vote, decrypt the E2E-encrypted payload so the - // webhook can expose which options were selected. Votes arrive as - // SHA-256 hashes of the option text; we match those back to the - // plaintext options remembered at send time (see SendPoll in - // handlers.go). If the session was restarted between send and vote - // we cannot resolve plaintext; hashes are still emitted so the - // consumer can perform matching itself if it has stored options. - if evt.Message.GetPollUpdateMessage() != nil { - pollMsgID := evt.Message.GetPollUpdateMessage().GetPollCreationMessageKey().GetID() - - pollVote, perr := mycli.WAClient.DecryptPollVote(context.Background(), evt) - if perr != nil { - log.Warn().Err(perr).Str("pollMsgID", pollMsgID).Msg("DecryptPollVote failed") - } - - if perr == nil && pollVote != nil { - hashes := pollVote.GetSelectedOptions() - hashB64 := make([]string, 0, len(hashes)) - for _, h := range hashes { - hashB64 = append(hashB64, base64.StdEncoding.EncodeToString(h)) - } - - selected := make([]string, 0, len(hashes)) - if stored := clientManager.GetPollOptions(mycli.userID, pollMsgID); len(stored) > 0 { - optionsByHash := make(map[string]string, len(stored)) - for _, opt := range stored { - sum := sha256.Sum256([]byte(opt)) - optionsByHash[string(sum[:])] = opt - } - for _, h := range hashes { - if opt, found := optionsByHash[string(h)]; found { - selected = append(selected, opt) - } - } - } - - postmap["pollVote"] = map[string]interface{}{ - "pollCreationMsgID": pollMsgID, - "selectedOptions": selected, - "selectedHashesB64": hashB64, - } - } - } - - if encMessage := evt.Message.GetSecretEncryptedMessage(); encMessage != nil { - decrypted, derr := mycli.WAClient.DecryptSecretEncryptedMessage(context.Background(), evt) - if derr != nil { - log.Warn(). - Err(derr). - Str("messageID", evt.Info.ID). - Str("secretEncType", encMessage.GetSecretEncType().String()). - Msg("DecryptSecretEncryptedMessage failed") - } else if decrypted != nil { - log.Info(). - Str("messageID", evt.Info.ID). - Str("secretEncType", encMessage.GetSecretEncType().String()). - Msg("Decrypted secretEncryptedMessage; swapping evt.Message") - evt.Message = decrypted - } - } - - if !*skipMedia { - - isIncoming := !evt.Info.IsFromMe - chatJID := evt.Info.Sender.String() - if evt.Info.IsGroup { - chatJID = evt.Info.Chat.String() - } - - s3cfg := mediaS3Config{ - Enabled: s3Config.Enabled, - MediaDelivery: s3Config.MediaDelivery, - } - - if img := evt.Message.GetImageMessage(); img != nil { - mycli.processMedia(img, img.GetMimetype(), ".jpg", - downloadTimeoutImage, isIncoming, chatJID, - evt.Info.ID, s3cfg, postmap, nil) - } - - if audio := evt.Message.GetAudioMessage(); audio != nil { - mycli.processMedia(audio, audio.GetMimetype(), ".ogg", - downloadTimeoutAudio, isIncoming, chatJID, - evt.Info.ID, s3cfg, postmap, nil) - } - - if doc := evt.Message.GetDocumentMessage(); doc != nil { - ext := ".bin" - if doc.FileName != nil { - ext = filepath.Ext(*doc.FileName) - } - mycli.processMedia(doc, doc.GetMimetype(), ext, - downloadTimeoutDocument, isIncoming, chatJID, - evt.Info.ID, s3cfg, postmap, nil) - } - - if video := evt.Message.GetVideoMessage(); video != nil { - mycli.processMedia(video, video.GetMimetype(), ".mp4", - downloadTimeoutVideo, isIncoming, chatJID, - evt.Info.ID, s3cfg, postmap, nil) - } - - if sticker := evt.Message.GetStickerMessage(); sticker != nil { - mycli.processMedia(sticker, sticker.GetMimetype(), ".webp", - downloadTimeoutSticker, isIncoming, chatJID, - evt.Info.ID, s3cfg, postmap, map[string]interface{}{ - "isSticker": true, - "stickerAnimated": sticker.GetIsAnimated(), - }) - } - } - - // Save message to history regardless of skipMedia setting - // Get user's history setting from cache - var historyLimit int - userinfo, found := userinfocache.Get(mycli.token) - if found { - historyStr := userinfo.(Values).Get("History") - historyLimit, _ = strconv.Atoi(historyStr) - } else { - log.Warn().Str("userID", mycli.userID).Msg("User info not found in cache, skipping history") - historyLimit = 0 - } - - if historyLimit > 0 { - messageType := "text" - textContent := "" - mediaLink := "" - caption := "" - replyToMessageID := "" - - // Check for delete messages first - if protocolMsg := evt.Message.GetProtocolMessage(); protocolMsg != nil && protocolMsg.GetType() == 0 { - messageType = "delete" - if protocolMsg.GetKey() != nil { - textContent = protocolMsg.GetKey().GetID() // Store the deleted message ID - } - log.Info().Str("deletedMessageID", textContent).Str("messageID", evt.Info.ID).Msg("Delete message detected") - // Check for reactions - } else if reaction := evt.Message.GetReactionMessage(); reaction != nil { - messageType = "reaction" - replyToMessageID = reaction.GetKey().GetID() - textContent = reaction.GetText() // This will be the emoji - } else if img := evt.Message.GetImageMessage(); img != nil { - messageType = "image" - caption = img.GetCaption() - } else if video := evt.Message.GetVideoMessage(); video != nil { - messageType = "video" - caption = video.GetCaption() - } else if audio := evt.Message.GetAudioMessage(); audio != nil { - messageType = "audio" - } else if doc := evt.Message.GetDocumentMessage(); doc != nil { - messageType = "document" - caption = doc.GetCaption() - } else if sticker := evt.Message.GetStickerMessage(); sticker != nil { - messageType = "sticker" - } else if contact := evt.Message.GetContactMessage(); contact != nil { - messageType = "contact" - textContent = contact.GetDisplayName() - } else if location := evt.Message.GetLocationMessage(); location != nil { - messageType = "location" - textContent = location.GetName() - } - - // Extract text content for non-reaction and non-delete messages - if messageType != "reaction" && messageType != "delete" { - if conv := evt.Message.GetConversation(); conv != "" { - textContent = conv - } else if ext := evt.Message.GetExtendedTextMessage(); ext != nil { - textContent = ext.GetText() - // Check if this is a reply to another message - if contextInfo := ext.GetContextInfo(); contextInfo != nil && contextInfo.GetStanzaID() != "" { - replyToMessageID = contextInfo.GetStanzaID() - } - } else { - textContent = caption - } - - // Set default text content for media messages without captions - if textContent == "" { - switch messageType { - case "image": - textContent = ":image:" - case "video": - textContent = ":video:" - case "audio": - textContent = ":audio:" - case "document": - textContent = ":document:" - case "sticker": - textContent = ":sticker:" - case "contact": - if textContent == "" { - textContent = ":contact:" - } - case "location": - if textContent == "" { - textContent = ":location:" - } - } - } - } - - // Check for replies in regular conversation messages too - if messageType == "text" && replyToMessageID == "" { - // For regular text messages, check if there's context info indicating a reply - // This might be available in the message context - if conv := evt.Message.GetConversation(); conv != "" { - // Check if the message has reply context (this depends on WhatsApp message structure) - // For now, we'll rely on ExtendedTextMessage for reply detection - } - } - - // Try to get media link from S3 data if available - if s3Data, ok := postmap["s3"].(map[string]interface{}); ok { - if url, ok := s3Data["url"].(string); ok { - mediaLink = url - } - } - - // Only save if there's meaningful content (including delete messages) - if textContent != "" || mediaLink != "" || (messageType != "text" && messageType != "reaction") || messageType == "delete" { - // Serializar evt para JSON - evtJSON, err := json.Marshal(evt) - if err != nil { - log.Error().Err(err).Msg("Failed to marshal event to JSON") - evtJSON = []byte("{}") - } - - err = mycli.s.saveMessageToHistory( - mycli.userID, - evt.Info.Chat.String(), - evt.Info.Sender.String(), - evt.Info.ID, - messageType, - textContent, - mediaLink, - replyToMessageID, - string(evtJSON), - ) - if err != nil { - log.Error().Err(err).Msg("Failed to save message to history") - } else { - err = mycli.s.trimMessageHistory(mycli.userID, evt.Info.Chat.String(), historyLimit) - if err != nil { - log.Error().Err(err).Msg("Failed to trim message history") - } - } - } else { - log.Debug().Str("messageType", messageType).Str("messageID", evt.Info.ID).Msg("Skipping empty message from history") - } - } - - case *events.Receipt: - postmap["type"] = "ReadReceipt" - dowebhook = 1 - //if evt.Type == events.ReceiptTypeRead || evt.Type == events.ReceiptTypeReadSelf { - if evt.Type == types.ReceiptTypeRead || evt.Type == types.ReceiptTypeReadSelf { - log.Info().Strs("id", evt.MessageIDs).Str("source", evt.SourceString()).Str("timestamp", fmt.Sprintf("%v", evt.Timestamp)).Msg("Message was read") - //if evt.Type == events.ReceiptTypeRead { - if evt.Type == types.ReceiptTypeRead { - postmap["state"] = "Read" - } else { - postmap["state"] = "ReadSelf" - } - //} else if evt.Type == events.ReceiptTypeDelivered { - } else if evt.Type == types.ReceiptTypeDelivered { - postmap["state"] = "Delivered" - log.Info().Str("id", evt.MessageIDs[0]).Str("source", evt.SourceString()).Str("timestamp", fmt.Sprintf("%v", evt.Timestamp)).Msg("Message delivered") - } else { - // Discard webhooks for inactive or other delivery types - return - } - case *events.Presence: - postmap["type"] = "Presence" - dowebhook = 1 - if evt.Unavailable { - postmap["state"] = "offline" - if evt.LastSeen.IsZero() { - log.Info().Str("from", evt.From.String()).Msg("User is now offline") - } else { - log.Info().Str("from", evt.From.String()).Str("lastSeen", fmt.Sprintf("%v", evt.LastSeen)).Msg("User is now offline") - } - } else { - postmap["state"] = "online" - log.Info().Str("from", evt.From.String()).Msg("User is now online") - } - case *events.HistorySync: - postmap["type"] = "HistorySync" - dowebhook = 1 - - // Save HistorySync messages to message_history table - if evt.Data != nil && evt.Data.Conversations != nil { - go func() { - - // Get the account owner's JID for messages sent by the instance - accountOwnerJID := "" - if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { - accountOwnerJID = mycli.WAClient.Store.ID.ToNonAD().String() - } - - savedCount := 0 - for _, conv := range evt.Data.Conversations { - if conv == nil || conv.ID == nil || conv.Messages == nil { - continue - } - - chatJID, err := types.ParseJID(*conv.ID) - if err != nil { - log.Warn().Err(err).Str("convID", *conv.ID).Msg("Failed to parse conversation JID in HistorySync") - continue - } - - for _, msg := range conv.Messages { - if msg == nil || msg.Message == nil { - continue - } - - // Extract message data - messageKey := msg.Message.GetKey() - if messageKey == nil { - continue - } - - messageID := messageKey.GetID() - if messageID == "" { - continue - } - - // Determine sender - never use "me", always use actual JID - // Use GetFromMe() from MessageKey to determine if message is from account owner - // This is more reliable than checking GetParticipant() - isFromMe := messageKey.GetFromMe() - var senderJID string - - if isFromMe { - // Message from account owner - senderJID = accountOwnerJID - if senderJID == "" { - // Fallback: use "me" if account owner JID is not available - senderJID = "me" - log.Warn().Str("messageID", messageID).Msg("accountOwnerJID is not available for a message from me, using 'me' as senderJID") - } - } else { - // Message from someone else - participantJID := messageKey.GetParticipant() - if chatJID.Server == types.GroupServer || chatJID.Server == types.BroadcastServer { - // Group message: use participant JID - senderJID = participantJID - } else { - // Direct message: sender is the chat itself (chat_jid) - senderJID = chatJID.String() - } - } - - // If senderJID is still empty, skip this message - if senderJID == "" { - log.Warn().Str("messageID", messageID).Msg("Cannot determine sender JID, skipping message") - continue - } - - // Get message content - message := msg.Message.GetMessage() - if message == nil { - continue - } - - // Extract message type and content - messageType := "unknown" - textContent := "" - mediaLink := "" - quotedMessageID := "" - - if message.GetConversation() != "" { - messageType = "text" - textContent = message.GetConversation() - } else if ext := message.GetExtendedTextMessage(); ext != nil { - messageType = "text" - textContent = ext.GetText() - if contextInfo := ext.GetContextInfo(); contextInfo != nil { - quotedMessageID = contextInfo.GetStanzaID() - } - } else if img := message.GetImageMessage(); img != nil { - messageType = "image" - textContent = img.GetCaption() - } else if vid := message.GetVideoMessage(); vid != nil { - messageType = "video" - textContent = vid.GetCaption() - } else if audio := message.GetAudioMessage(); audio != nil { - messageType = "audio" - } else if doc := message.GetDocumentMessage(); doc != nil { - messageType = "document" - textContent = doc.GetCaption() - } else if sticker := message.GetStickerMessage(); sticker != nil { - messageType = "sticker" - } else if location := message.GetLocationMessage(); location != nil { - messageType = "location" - textContent = location.GetName() - } else if contact := message.GetContactMessage(); contact != nil { - messageType = "contact" - textContent = contact.GetDisplayName() - } else if buttons := message.GetButtonsResponseMessage(); buttons != nil { - messageType = "buttons_response" - textContent = buttons.GetSelectedButtonID() - } else if list := message.GetListResponseMessage(); list != nil { - messageType = "list_response" - textContent = list.GetSingleSelectReply().GetSelectedRowID() - } else if reaction := message.GetReactionMessage(); reaction != nil { - messageType = "reaction" - textContent = reaction.GetText() - if key := reaction.GetKey(); key != nil { - quotedMessageID = key.GetID() - } - } - - // Set default text for media messages without captions - if textContent == "" && messageType != "text" && messageType != "reaction" && messageType != "delete" { - switch messageType { - case "image": - textContent = ":image:" - case "video": - textContent = ":video:" - case "audio": - textContent = ":audio:" - case "document": - textContent = ":document:" - case "sticker": - textContent = ":sticker:" - case "contact": - textContent = ":contact:" - case "location": - textContent = ":location:" - } - } - - // Get message timestamp - msgTimestamp := time.Now() - if timestamp := msg.Message.GetMessageTimestamp(); timestamp > 0 { - msgTimestamp = time.Unix(int64(timestamp), 0) - } - - // Parse sender JID for MessageInfo - var senderJIDForInfo types.JID - if isFromMe { - if accountOwnerJID != "" { - var pErr error - senderJIDForInfo, pErr = types.ParseJID(accountOwnerJID) - if pErr != nil { - log.Warn().Err(pErr).Str("accountOwnerJID", accountOwnerJID).Msg("Failed to parse account owner JID in HistorySync") - } - } - } else { - if chatJID.Server == types.GroupServer || chatJID.Server == types.BroadcastServer { - // Group: use participant JID - participant := messageKey.GetParticipant() - if participant != "" { - var pErr error - senderJIDForInfo, pErr = types.ParseJID(participant) - if pErr != nil { - log.Warn().Err(pErr).Str("participantJID", participant).Msg("Failed to parse participant JID in HistorySync") - } - } - } else { - // Direct message: sender is the chat - senderJIDForInfo = chatJID - } - } - - // Try to get PushName from store if available - pushName := "" - if !isFromMe && senderJIDForInfo.User != "" { - if mycli.WAClient != nil && mycli.WAClient.Store != nil { - if contact, err := mycli.WAClient.Store.Contacts.GetContact(context.Background(), senderJIDForInfo); err == nil { - pushName = contact.PushName - } - } - } - - // Create MessageInfo structure matching events.Message format - messageInfo := types.MessageInfo{ - MessageSource: types.MessageSource{ - Chat: chatJID, - Sender: senderJIDForInfo, - IsFromMe: isFromMe, - IsGroup: chatJID.Server == types.GroupServer || chatJID.Server == types.BroadcastServer, - }, - ID: messageID, - Timestamp: msgTimestamp, - Type: messageType, - PushName: pushName, - } - - // Create events.Message-like structure for datajson - // This matches the format used in regular message events - // RawMessage should be the full waE2E.Message structure - messageEvent := map[string]interface{}{ - "Info": messageInfo, - "Message": message, - "IsEphemeral": false, - "IsViewOnce": false, - "IsViewOnceV2": false, - "IsViewOnceV2Extension": false, - "IsDocumentWithCaption": false, - "IsLottieSticker": false, - "IsBotInvoke": false, - "IsEdit": false, - "SourceWebMsg": nil, - "UnavailableRequestID": "", - "RetryCount": 0, - "NewsletterMeta": nil, - "RawMessage": msg.Message, - } - - // Serialize to JSON for datajson field - evtJSON, err := json.Marshal(messageEvent) - if err != nil { - log.Error().Err(err).Msg("Failed to marshal HistorySync message event to JSON") - evtJSON = []byte("{}") - } - - // Save message to history - // Only save if there's meaningful content - if textContent != "" || mediaLink != "" || (messageType != "text" && messageType != "reaction") { - err = mycli.s.saveMessageToHistory( - mycli.userID, - chatJID.String(), - senderJID, - messageID, - messageType, - textContent, - mediaLink, - quotedMessageID, - string(evtJSON), - ) - if err != nil { - log.Error().Err(err). - Str("userID", mycli.userID). - Str("chatJID", chatJID.String()). - Str("messageID", messageID). - Msg("Failed to save HistorySync message to history") - } else { - savedCount++ - } - } - } - } - - if savedCount > 0 { - log.Info(). - Str("userID", mycli.userID). - Int("savedCount", savedCount). - Msg("Saved HistorySync messages to message_history") - } - }() - } - - case *events.AppState: - log.Info().Str("index", fmt.Sprintf("%+v", evt.Index)).Str("actionValue", fmt.Sprintf("%+v", evt.SyncActionValue)).Msg("App state event received") - case *events.LoggedOut: - postmap["type"] = "LoggedOut" - dowebhook = 1 - log.Info().Str("reason", evt.Reason.String()).Msg("Logged out") - defer func() { - // Use a non-blocking send to prevent a deadlock if the receiver has already terminated. - signalKill(mycli.userID) - }() - sqlStmt := `UPDATE users SET connected=0 WHERE id=$1` - _, err := mycli.db.Exec(sqlStmt, mycli.userID) - if err != nil { - log.Error().Err(err).Msg(sqlStmt) - return - } - case *events.ChatPresence: - postmap["type"] = "ChatPresence" - dowebhook = 1 - log.Info().Str("state", fmt.Sprintf("%s", evt.State)).Str("media", fmt.Sprintf("%s", evt.Media)).Str("chat", evt.MessageSource.Chat.String()).Str("sender", evt.MessageSource.Sender.String()).Msg("Chat Presence received") - case *events.CallOffer: - postmap["type"] = "CallOffer" - dowebhook = 1 - log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call offer") - case *events.CallAccept: - postmap["type"] = "CallAccept" - dowebhook = 1 - log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call accept") - case *events.CallTerminate: - postmap["type"] = "CallTerminate" - dowebhook = 1 - log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call terminate") - case *events.CallOfferNotice: - postmap["type"] = "CallOfferNotice" - dowebhook = 1 - log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call offer notice") - case *events.CallRelayLatency: - postmap["type"] = "CallRelayLatency" - dowebhook = 1 - log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call relay latency") - case *events.Disconnected: - postmap["type"] = "Disconnected" - dowebhook = 1 - log.Info().Str("reason", fmt.Sprintf("%+v", evt)).Msg("Disconnected from Whatsapp") - case *events.ConnectFailure: - postmap["type"] = "ConnectFailure" - dowebhook = 1 - log.Error().Str("reason", fmt.Sprintf("%+v", evt)).Msg("Failed to connect to Whatsapp") - case *events.UndecryptableMessage: - postmap["type"] = "UndecryptableMessage" - dowebhook = 1 - log.Warn().Str("info", evt.Info.SourceString()).Msg("Undecryptable message received") - case *events.MediaRetry: - postmap["type"] = "MediaRetry" - dowebhook = 1 - log.Info().Str("messageID", evt.MessageID).Msg("Media retry event") - case *events.GroupInfo: - postmap["type"] = "GroupInfo" - dowebhook = 1 - log.Info().Str("jid", evt.JID.String()).Msg("Group info updated") - case *events.JoinedGroup: - postmap["type"] = "JoinedGroup" - dowebhook = 1 - log.Info().Str("jid", evt.JID.String()).Msg("Joined group") - case *events.Picture: - postmap["type"] = "Picture" - dowebhook = 1 - log.Info().Str("jid", evt.JID.String()).Msg("Picture updated") - case *events.BlocklistChange: - postmap["type"] = "BlocklistChange" - dowebhook = 1 - log.Info().Str("jid", evt.JID.String()).Msg("Blocklist changed") - case *events.Blocklist: - postmap["type"] = "Blocklist" - dowebhook = 1 - log.Info().Msg("Blocklist received") - case *events.KeepAliveRestored: - postmap["type"] = "KeepAliveRestored" - dowebhook = 1 - log.Info().Msg("Keep alive restored") - case *events.KeepAliveTimeout: - postmap["type"] = "KeepAliveTimeout" - dowebhook = 1 - log.Warn().Msg("Keep alive timeout") - case *events.ClientOutdated: - postmap["type"] = "ClientOutdated" - dowebhook = 1 - log.Warn().Msg("Client outdated") - case *events.TemporaryBan: - postmap["type"] = "TemporaryBan" - dowebhook = 1 - log.Info().Msg("Temporary ban") - case *events.StreamError: - postmap["type"] = "StreamError" - dowebhook = 1 - log.Error().Str("code", evt.Code).Msg("Stream error") - case *events.PairError: - postmap["type"] = "PairError" - dowebhook = 1 - log.Error().Msg("Pair error") - case *events.PrivacySettings: - postmap["type"] = "PrivacySettings" - dowebhook = 1 - log.Info().Msg("Privacy settings updated") - case *events.UserAbout: - postmap["type"] = "UserAbout" - dowebhook = 1 - log.Info().Str("jid", evt.JID.String()).Msg("User about updated") - case *events.OfflineSyncCompleted: - postmap["type"] = "OfflineSyncCompleted" - dowebhook = 1 - log.Info().Msg("Offline sync completed") - case *events.OfflineSyncPreview: - postmap["type"] = "OfflineSyncPreview" - dowebhook = 1 - log.Info().Msg("Offline sync preview") - case *events.IdentityChange: - postmap["type"] = "IdentityChange" - dowebhook = 1 - log.Info().Str("jid", evt.JID.String()).Msg("Identity changed") - case *events.NewsletterJoin: - postmap["type"] = "NewsletterJoin" - dowebhook = 1 - log.Info().Str("jid", evt.ID.String()).Msg("Newsletter joined") - case *events.NewsletterLeave: - postmap["type"] = "NewsletterLeave" - dowebhook = 1 - log.Info().Str("jid", evt.ID.String()).Msg("Newsletter left") - case *events.NewsletterMuteChange: - postmap["type"] = "NewsletterMuteChange" - dowebhook = 1 - log.Info().Str("jid", evt.ID.String()).Msg("Newsletter mute changed") - case *events.NewsletterLiveUpdate: - postmap["type"] = "NewsletterLiveUpdate" - dowebhook = 1 - log.Info().Msg("Newsletter live update") - case *events.FBMessage: - postmap["type"] = "FBMessage" - dowebhook = 1 - log.Info().Str("info", evt.Info.SourceString()).Msg("Facebook message received") - default: - log.Warn().Str("event", fmt.Sprintf("%+v", evt)).Msg("Unhandled event") - } - - if dowebhook == 1 { - sendEventWithWebHook(mycli, postmap, path) - } -} +package main + +import ( + "context" + "crypto/sha256" + "crypto/tls" + "encoding/base64" + "encoding/json" + "errors" + "fmt" + "net/http" + "net/url" + "os" + "path/filepath" + "runtime/debug" + "strconv" + "strings" + "time" + + "github.com/go-resty/resty/v2" + "github.com/jmoiron/sqlx" + "github.com/mdp/qrterminal/v3" + "github.com/patrickmn/go-cache" + "github.com/rs/zerolog/log" + "github.com/skip2/go-qrcode" + "go.mau.fi/whatsmeow" + "go.mau.fi/whatsmeow/appstate" + "go.mau.fi/whatsmeow/proto/waCompanionReg" + "go.mau.fi/whatsmeow/store" + "go.mau.fi/whatsmeow/types" + "go.mau.fi/whatsmeow/types/events" + waLog "go.mau.fi/whatsmeow/util/log" + "golang.org/x/net/proxy" +) + +// db field declaration as *sqlx.DB +type MyClient struct { + WAClient *whatsmeow.Client + eventHandlerID uint32 + userID string + token string + db *sqlx.DB + s *server +} + +// safeGo runs fn in a new goroutine with a defer recover so a panic inside +// fire-and-forget side-effects (webhook delivery, MQ push) cannot crash +// the whole process. Losing one delivery is preferable to taking wuzapi +// down for every connected user. +func safeGo(name string, fn func()) { + go func() { + defer func() { + if r := recover(); r != nil { + log.Error(). + Str("goroutine", name). + Interface("panic", r). + Str("stack", string(debug.Stack())). + Msg("panic recovered in goroutine") + } + }() + fn() + }() +} + +// ensureS3ClientForUser loads S3 config from DB and initializes client if not already present (lazy init for reconnect-after-restart) +func ensureS3ClientForUser(userID string) { + GetS3Manager().EnsureClientFromDB(userID) +} + +func sendToGlobalWebHook(jsonData []byte, token string, userID string) { + jsonDataStr := string(jsonData) + + instance_name := "" + userinfo, found := userinfocache.Get(token) + if found { + instance_name = userinfo.(Values).Get("Name") + } + + if *globalWebhook != "" { + log.Info().Str("url", *globalWebhook).Msg("Calling global webhook") + // Add extra information for the global webhook + globalData := map[string]string{ + "jsonData": jsonDataStr, + "userID": userID, + "instanceName": instance_name, + } + callHookWithHmac(*globalWebhook, globalData, userID, globalHMACKeyEncrypted) + } +} + +func sendToUserWebHook(webhookurl string, path string, jsonData []byte, userID string, token string) { + sendToUserWebHookWithHmac(webhookurl, path, jsonData, userID, token, nil) +} + +func sendToUserWebHookWithHmac(webhookurl string, path string, jsonData []byte, userID string, token string, encryptedHmacKey []byte) { + + instance_name := "" + userinfo, found := userinfocache.Get(token) + if found { + instance_name = userinfo.(Values).Get("Name") + } + data := map[string]string{ + "jsonData": string(jsonData), + "userID": userID, + "instanceName": instance_name, + } + + log.Debug().Interface("webhookData", data).Msg("Data being sent to webhook") + + if webhookurl != "" { + log.Info().Str("url", webhookurl).Msg("Calling user webhook") + + if path == "" { + safeGo("callHookWithHmac", func() { callHookWithHmac(webhookurl, data, userID, encryptedHmacKey) }) + } else { + if err := callHookFileWithHmac(webhookurl, data, userID, path, encryptedHmacKey); err != nil { + log.Error().Err(err).Msg("Error calling hook file") + } + } + } else { + log.Warn().Str("userid", userID).Msg("No webhook set for user") + } +} + +func updateAndGetUserSubscriptions(mycli *MyClient) ([]string, error) { + // Get updated events from cache/database + currentEvents := "" + userinfo2, found2 := userinfocache.Get(mycli.token) + if found2 { + currentEvents = userinfo2.(Values).Get("Events") + } else { + // If not in cache, get from database + if err := mycli.db.Get(¤tEvents, "SELECT events FROM users WHERE id=$1", mycli.userID); err != nil { + log.Warn().Err(err).Str("userID", mycli.userID).Msg("Could not get events from DB") + return nil, err // Propagate the error + } + } + + // Update client subscriptions if changed + eventarray := strings.Split(currentEvents, ",") + var subscribedEvents []string + if len(eventarray) == 1 && eventarray[0] == "" { + subscribedEvents = []string{} + } else { + for _, arg := range eventarray { + arg = strings.TrimSpace(arg) + if arg != "" && Find(supportedEventTypes, arg) { + subscribedEvents = append(subscribedEvents, arg) + } + } + } + + return subscribedEvents, nil +} + +func getUserWebhookUrl(token string) string { + webhookurl := "" + myuserinfo, found := userinfocache.Get(token) + if !found { + log.Warn().Str("token", token).Msg("Could not call webhook as there is no user for this token") + } else { + webhookurl = myuserinfo.(Values).Get("Webhook") + } + return webhookurl +} + +func sendEventWithWebHook(mycli *MyClient, postmap map[string]interface{}, path string) { + webhookurl := getUserWebhookUrl(mycli.token) + + // Get updated events from cache/database + subscribedEvents, err := updateAndGetUserSubscriptions(mycli) + if err != nil { + return + } + + eventType, ok := postmap["type"].(string) + if !ok { + log.Error().Msg("Event type is not a string in postmap") + return + } + + // Log subscription details for debugging + log.Debug(). + Str("userID", mycli.userID). + Str("eventType", eventType). + Strs("subscribedEvents", subscribedEvents). + Msg("Checking event subscription") + + // Check if the current event is in the subscriptions + checkIfSubscribedInEvent := checkIfSubscribedToEvent(subscribedEvents, postmap["type"].(string), mycli.userID) + if !checkIfSubscribedInEvent { + return + } + + // In stdio mode, send as JSON-RPC notification instead of HTTP webhook + if mycli.s != nil && mycli.s.mode == Stdio { + mycli.s.SendNotification(eventType, postmap) + return + } + + // Prepare webhook data + jsonData, err := json.Marshal(postmap) + if err != nil { + log.Error().Err(err).Msg("Failed to marshal postmap to JSON") + return + } + + // Get HMAC key for this user + var encryptedHmacKey []byte + if userinfo, found := userinfocache.Get(mycli.token); found { + encryptedB64 := userinfo.(Values).Get("HmacKeyEncrypted") + if encryptedB64 != "" { + var err error + encryptedHmacKey, err = base64.StdEncoding.DecodeString(encryptedB64) + if err != nil { + log.Error().Err(err).Msg("Failed to decode HMAC key from cache") + } + } + } + + sendToUserWebHookWithHmac(webhookurl, path, jsonData, mycli.userID, mycli.token, encryptedHmacKey) + + // Get global webhook if configured + safeGo("sendToGlobalWebHook", func() { sendToGlobalWebHook(jsonData, mycli.token, mycli.userID) }) + + safeGo("sendToGlobalRabbit", func() { sendToGlobalRabbit(jsonData, mycli.token, mycli.userID) }) +} + +func checkIfSubscribedToEvent(subscribedEvents []string, eventType string, userId string) bool { + if !Find(subscribedEvents, eventType) && !Find(subscribedEvents, "All") { + log.Warn(). + Str("type", eventType). + Strs("subscribedEvents", subscribedEvents). + Str("userID", userId). + Msg("Skipping webhook. Not subscribed for this type") + return false + } + return true +} + +// Connects to Whatsapp Websocket on server startup if last state was connected +func (s *server) connectOnStartup() { + rows, err := s.db.Queryx("SELECT id,name,token,jid,webhook,events,proxy_url,CASE WHEN s3_enabled THEN 'true' ELSE 'false' END AS s3_enabled,media_delivery,COALESCE(history, 0) as history,hmac_key FROM users WHERE connected=1") + if err != nil { + log.Error().Err(err).Msg("DB Problem") + return + } + defer rows.Close() + for rows.Next() { + txtid := "" + token := "" + jid := "" + name := "" + webhook := "" + events := "" + proxy_url := "" + s3_enabled := "" + media_delivery := "" + var history int + var hmac_key []byte + err = rows.Scan(&txtid, &name, &token, &jid, &webhook, &events, &proxy_url, &s3_enabled, &media_delivery, &history, &hmac_key) + if err != nil { + log.Error().Err(err).Msg("DB Problem") + return + } else { + hmacKeyEncrypted := "" + if len(hmac_key) > 0 { + hmacKeyEncrypted = base64.StdEncoding.EncodeToString(hmac_key) + } + + log.Info().Str("token", token).Msg("Connect to Whatsapp on startup") + v := Values{map[string]string{ + "Id": txtid, + "Name": name, + "Jid": jid, + "Webhook": webhook, + "Token": token, + "Proxy": proxy_url, + "Events": events, + "S3Enabled": s3_enabled, + "MediaDelivery": media_delivery, + "History": fmt.Sprintf("%d", history), + "HmacKeyEncrypted": hmacKeyEncrypted, + }} + userinfocache.Set(token, v, cache.NoExpiration) + // Gets and set subscription to webhook events + eventarray := strings.Split(events, ",") + + var subscribedEvents []string + if len(eventarray) == 1 && eventarray[0] == "" { + subscribedEvents = []string{} + } else { + for _, arg := range eventarray { + if !Find(supportedEventTypes, arg) { + log.Warn().Str("Type", arg).Msg("Event type discarded") + continue + } + if !Find(subscribedEvents, arg) { + subscribedEvents = append(subscribedEvents, arg) + } + } + + } + eventstring := strings.Join(subscribedEvents, ",") + log.Info().Str("events", eventstring).Str("jid", jid).Msg("Attempt to connect") + kill := make(chan bool, 1) + setKillChannel(txtid, kill) + go s.startClient(txtid, jid, token, kill) + + // Initialize S3 client if configured + go func(userID string) { + GetS3Manager().EnsureClientFromDB(userID) + }(txtid) + } + } + err = rows.Err() + if err != nil { + log.Error().Err(err).Msg("DB Problem") + } +} + +func parseJID(arg string) (types.JID, bool) { + if arg == "" { + return types.JID{}, false + } + if arg[0] == '+' { + arg = arg[1:] + } + if !strings.ContainsRune(arg, '@') { + return types.NewJID(arg, types.DefaultUserServer), true + } else { + recipient, err := types.ParseJID(arg) + if err != nil { + log.Error().Err(err).Msg("Invalid JID") + return recipient, false + } else if recipient.User == "" { + log.Error().Err(err).Msg("Invalid JID no server specified") + return recipient, false + } + return recipient, true + } +} + +// jidUserKey returns the phone/account part shared by users.jid and +// whatsmeow_device.jid even when formats differ (380...:8@ vs 380...@). +func jidUserKey(jid string) string { + if jid == "" { + return "" + } + userPart := strings.SplitN(jid, "@", 2)[0] + return strings.SplitN(userPart, ":", 2)[0] +} + +// jidLookupCandidates builds JID variants to try with sqlstore.GetDevice before +// falling back to account-key matching across all stored devices. +func jidLookupCandidates(textjid string) []types.JID { + jid, ok := parseJID(textjid) + if !ok { + return nil + } + + seen := make(map[string]struct{}) + var candidates []types.JID + add := func(candidate types.JID) { + if candidate.IsEmpty() { + return + } + key := candidate.String() + if _, exists := seen[key]; exists { + return + } + seen[key] = struct{}{} + candidates = append(candidates, candidate) + } + + add(jid) + add(jid.ToNonAD()) + + userOnly, _, _ := strings.Cut(jid.User, ":") + if userOnly != "" { + add(types.NewJID(userOnly, jid.Server)) + if jid.Server == types.DefaultUserServer || jid.Server == types.HiddenUserServer { + add(types.NewJID(userOnly, types.DefaultUserServer)) + } + } + + return candidates +} + +func canonicalStoreJID(deviceStore *store.Device) string { + if deviceStore == nil || deviceStore.ID == nil { + return "" + } + return deviceStore.ID.ToNonAD().String() +} + +// resolveDeviceStore loads an existing WhatsApp session from sqlstore. When +// users.jid does not exactly match whatsmeow_device.jid (common after LID/AD +// format changes), it falls back to matching by account key so reconnect does +// not create a fresh device and force QR scan. +func (s *server) resolveDeviceStore(ctx context.Context, textjid string) (*store.Device, string) { + if textjid != "" { + for _, candidate := range jidLookupCandidates(textjid) { + deviceStore, err := container.GetDevice(ctx, candidate) + if err != nil { + log.Error().Err(err).Str("jid", candidate.String()).Msg("Failed to get device") + continue + } + if resolved := canonicalStoreJID(deviceStore); resolved != "" { + if candidate.String() != textjid { + log.Info(). + Str("user_jid", textjid). + Str("store_jid", deviceStore.ID.String()). + Str("resolved_jid", resolved). + Msg("Resolved device by JID variant") + } + return deviceStore, resolved + } + } + + key := jidUserKey(textjid) + if key != "" { + allDevices, err := container.GetAllDevices(ctx) + if err != nil { + log.Error().Err(err).Msg("Failed to list devices from store") + } else { + for _, deviceStore := range allDevices { + if resolved := canonicalStoreJID(deviceStore); resolved != "" && jidUserKey(resolved) == key { + log.Info(). + Str("user_jid", textjid). + Str("store_jid", deviceStore.ID.String()). + Str("resolved_jid", resolved). + Msg("Resolved device by account key") + return deviceStore, resolved + } + } + } + } + + log.Warn().Str("jid", textjid).Msg("No store found for jid. Creating new device") + } else { + log.Warn().Msg("No jid found. Creating new device") + } + + return container.NewDevice(), "" +} + +func (s *server) syncUserJID(userID, token, oldJID, newJID string) { + if newJID == "" || newJID == oldJID { + return + } + if _, err := s.db.Exec(`UPDATE users SET jid=$1 WHERE id=$2`, newJID, userID); err != nil { + log.Warn().Err(err).Str("user_id", userID).Msg("Failed to sync jid from device store") + return + } + if token != "" { + if myuserinfo, found := userinfocache.Get(token); found { + v := updateUserInfo(myuserinfo, "Jid", newJID) + userinfocache.Set(token, v, cache.NoExpiration) + } + } + log.Info().Str("user_id", userID).Str("old_jid", oldJID).Str("resolved_jid", newJID).Msg("Synced user jid from whatsmeow store") +} + +// getPlatformTypeEnum converts a platform type string to the corresponding DeviceProps enum +// Returns DESKTOP as default if the string doesn't match any known type +func getPlatformTypeEnum(platformType string) *waCompanionReg.DeviceProps_PlatformType { + platformType = strings.ToUpper(strings.TrimSpace(platformType)) + + switch platformType { + case "UNKNOWN": + return waCompanionReg.DeviceProps_UNKNOWN.Enum() + case "CHROME": + return waCompanionReg.DeviceProps_CHROME.Enum() + case "FIREFOX": + return waCompanionReg.DeviceProps_FIREFOX.Enum() + case "IE": + return waCompanionReg.DeviceProps_IE.Enum() + case "OPERA": + return waCompanionReg.DeviceProps_OPERA.Enum() + case "SAFARI": + return waCompanionReg.DeviceProps_SAFARI.Enum() + case "EDGE": + return waCompanionReg.DeviceProps_EDGE.Enum() + case "DESKTOP": + return waCompanionReg.DeviceProps_DESKTOP.Enum() + case "IPAD": + return waCompanionReg.DeviceProps_IPAD.Enum() + case "ANDROID_TABLET": + return waCompanionReg.DeviceProps_ANDROID_TABLET.Enum() + case "OHANA": + return waCompanionReg.DeviceProps_OHANA.Enum() + case "ALOHA": + return waCompanionReg.DeviceProps_ALOHA.Enum() + case "CATALINA": + return waCompanionReg.DeviceProps_CATALINA.Enum() + case "TCL_TV": + return waCompanionReg.DeviceProps_TCL_TV.Enum() + case "IOS_PHONE": + return waCompanionReg.DeviceProps_IOS_PHONE.Enum() + case "IOS_CATALYST": + return waCompanionReg.DeviceProps_IOS_CATALYST.Enum() + case "ANDROID_PHONE": + return waCompanionReg.DeviceProps_ANDROID_PHONE.Enum() + case "ANDROID_AMBIGUOUS": + return waCompanionReg.DeviceProps_ANDROID_AMBIGUOUS.Enum() + case "WEAR_OS": + return waCompanionReg.DeviceProps_WEAR_OS.Enum() + case "AR_WRIST": + return waCompanionReg.DeviceProps_AR_WRIST.Enum() + case "AR_DEVICE": + return waCompanionReg.DeviceProps_AR_DEVICE.Enum() + case "UWP": + return waCompanionReg.DeviceProps_UWP.Enum() + case "VR": + return waCompanionReg.DeviceProps_VR.Enum() + default: + log.Warn().Str("platformType", platformType).Msg("Unknown platform type, defaulting to DESKTOP") + return waCompanionReg.DeviceProps_DESKTOP.Enum() + } +} + +func (s *server) startClient(userID string, textjid string, token string, kill chan bool) { + log.Info().Str("userid", userID).Str("jid", textjid).Msg("Starting websocket connection to Whatsapp") + + // Connection retry constants + const maxConnectionRetries = 3 + const connectionRetryBaseWait = 5 * time.Second + + deviceStore, resolvedJID := s.resolveDeviceStore(context.Background(), textjid) + s.syncUserJID(userID, token, textjid, resolvedJID) + + var err error + + clientLog := waLog.Stdout("Client", *waDebug, *colorOutput) + + // Create the client with initialized deviceStore + var client *whatsmeow.Client + if *waDebug != "" { + client = whatsmeow.NewClient(deviceStore, clientLog) + } else { + client = whatsmeow.NewClient(deviceStore, nil) + } + + // Now we can use the client with the manager + clientManager.SetWhatsmeowClient(userID, client) + + store.DeviceProps.PlatformType = getPlatformTypeEnum(*platformType) + store.DeviceProps.Os = osName + + mycli := MyClient{ + WAClient: client, + eventHandlerID: 1, + userID: userID, + token: token, + db: s.db, + s: s, + } + mycli.eventHandlerID = mycli.WAClient.AddEventHandler(mycli.myEventHandler) + + // Store the MyClient in clientManager + clientManager.SetMyClient(userID, &mycli) + + httpClient := resty.New() + httpClient.SetRedirectPolicy(resty.FlexibleRedirectPolicy(15)) + if *waDebug == "DEBUG" { + httpClient.SetDebug(true) + } + httpClient.SetTimeout(30 * time.Second) + httpClient.SetTLSClientConfig(&tls.Config{InsecureSkipVerify: true}) + httpClient.OnError(func(req *resty.Request, err error) { + if v, ok := err.(*resty.ResponseError); ok { + // v.Response contains the last response from the server + // v.Err contains the original error + log.Debug().Str("response", v.Response.String()).Msg("resty error") + log.Error().Err(v.Err).Msg("resty error") + } + }) + + // Set proxy if defined in DB (assumes users table contains proxy_url column) + var proxyURL string + err = s.db.Get(&proxyURL, "SELECT proxy_url FROM users WHERE id=$1", userID) + if err == nil && proxyURL != "" { + parsed, perr := url.Parse(proxyURL) + if perr != nil { + log.Warn().Err(perr).Str("proxy", proxyURL).Msg("Invalid proxy URL, skipping proxy setup") + } else { + + log.Info().Str("proxy", proxyURL).Msg("Configuring proxy") + + if parsed.Scheme == "socks5" || parsed.Scheme == "socks5h" { + dialer, derr := proxy.FromURL(parsed, nil) + if derr != nil { + log.Warn().Err(derr).Str("proxy", proxyURL).Msg("Failed to build SOCKS proxy dialer, skipping proxy setup") + } else { + httpClient.SetProxy(proxyURL) + client.SetSOCKSProxy(dialer, whatsmeow.SetProxyOptions{}) + log.Info().Msg("SOCKS proxy configured successfully") + } + } else { + httpClient.SetProxy(proxyURL) + client.SetProxyAddress(parsed.String(), whatsmeow.SetProxyOptions{}) + log.Info().Msg("HTTP/HTTPS proxy configured successfully") + } + } + } + clientManager.SetHTTPClient(userID, httpClient) + + // Initialize S3 client if configured (needed when user reconnects after container restart - connectOnStartup only runs for connected=1) + GetS3Manager().EnsureClientFromDB(userID) + + if client.Store.ID == nil { + // No ID stored, new login + qrChan, err := client.GetQRChannel(context.Background()) + if err != nil { + // This error means that we're already logged in, so ignore it. + if !errors.Is(err, whatsmeow.ErrQRStoreContainsID) { + log.Error().Err(err).Msg("Failed to get QR channel") + return + } + } else { + err = client.Connect() // Must connect to generate QR code + if err != nil { + log.Error().Err(err).Msg("Failed to connect client") + return + } + + myuserinfo, found := userinfocache.Get(token) + + for evt := range qrChan { + if evt.Event == "code" { + // Display QR code in terminal (useful for testing/developing) + // Skip in stdio mode to avoid breaking JSON-RPC + if *logType != "json" && s.mode != Stdio { + qrterminal.GenerateHalfBlock(evt.Code, qrterminal.L, os.Stdout) + fmt.Println("QR code:\n", evt.Code) + } + // Store encoded/embeded base64 QR on database for retrieval with the /qr endpoint + image, _ := qrcode.Encode(evt.Code, qrcode.Medium, 256) + base64qrcode := "data:image/png;base64," + base64.StdEncoding.EncodeToString(image) + sqlStmt := `UPDATE users SET qrcode=$1 WHERE id=$2` + _, err := s.db.Exec(sqlStmt, base64qrcode, userID) + if err != nil { + log.Error().Err(err).Msg(sqlStmt) + } else { + if found { + v := updateUserInfo(myuserinfo, "Qrcode", base64qrcode) + userinfocache.Set(token, v, cache.NoExpiration) + log.Info().Str("qrcode", base64qrcode).Msg("update cache userinfo with qr code") + } + } + + //send QR code with webhook + postmap := make(map[string]interface{}) + postmap["event"] = evt.Event + postmap["qrCodeBase64"] = base64qrcode + postmap["type"] = "QR" + + sendEventWithWebHook(&mycli, postmap, "") + + } else if evt.Event == "timeout" { + // Clear QR code from DB on timeout + // Send webhook notifying QR timeout before cleanup + postmap := make(map[string]interface{}) + postmap["event"] = evt.Event + postmap["type"] = "QRTimeout" + sendEventWithWebHook(&mycli, postmap, "") + + sqlStmt := `UPDATE users SET qrcode='' WHERE id=$1` + _, err := s.db.Exec(sqlStmt, userID) + if err != nil { + log.Error().Err(err).Msg(sqlStmt) + } else { + if found { + v := updateUserInfo(myuserinfo, "Qrcode", "") + userinfocache.Set(token, v, cache.NoExpiration) + } + } + log.Warn().Msg("QR timeout killing channel") + clientManager.DeleteWhatsmeowClient(userID) + clientManager.DeleteMyClient(userID) + clientManager.DeleteHTTPClient(userID) + signalKill(userID) + } else if evt.Event == "success" { + log.Info().Msg("QR pairing ok!") + // Clear QR code after pairing + sqlStmt := `UPDATE users SET qrcode='', connected=1 WHERE id=$1` + _, err := s.db.Exec(sqlStmt, userID) + if err != nil { + log.Error().Err(err).Msg(sqlStmt) + } else { + if found { + v := updateUserInfo(myuserinfo, "Qrcode", "") + userinfocache.Set(token, v, cache.NoExpiration) + } + } + } else { + log.Info().Str("event", evt.Event).Msg("Login event") + } + } + } + + } else { + // Already logged in, just connect + log.Info().Msg("Already logged in, just connect") + + // Retry logic with linear backoff + var lastErr error + + for attempt := 0; attempt < maxConnectionRetries; attempt++ { + if attempt > 0 { + waitTime := time.Duration(attempt) * connectionRetryBaseWait + log.Warn(). + Int("attempt", attempt+1). + Int("max_retries", maxConnectionRetries). + Dur("wait_time", waitTime). + Msg("Retrying connection after delay") + time.Sleep(waitTime) + } + + err = client.Connect() + if err == nil { + log.Info(). + Int("attempt", attempt+1). + Msg("Successfully connected to WhatsApp") + break + } + + lastErr = err + log.Warn(). + Err(err). + Int("attempt", attempt+1). + Int("max_retries", maxConnectionRetries). + Msg("Failed to connect to WhatsApp") + } + + if lastErr != nil { + log.Error(). + Err(lastErr). + Str("userid", userID). + Int("attempts", maxConnectionRetries). + Msg("Failed to connect to WhatsApp after all retry attempts") + + clientManager.DeleteWhatsmeowClient(userID) + clientManager.DeleteMyClient(userID) + clientManager.DeleteHTTPClient(userID) + + sqlStmt := `UPDATE users SET qrcode='', connected=0 WHERE id=$1` + _, dbErr := s.db.Exec(sqlStmt, userID) + if dbErr != nil { + log.Error().Err(dbErr).Msg("Failed to update user status after connection error") + } + + // Use the existing mycli instance from outer scope + postmap := make(map[string]interface{}) + postmap["event"] = "ConnectFailure" + postmap["error"] = lastErr.Error() + postmap["type"] = "ConnectFailure" + postmap["attempts"] = maxConnectionRetries + postmap["reason"] = "Failed to connect after retry attempts" + sendEventWithWebHook(&mycli, postmap, "") + + return + } + } + + // Keep the session goroutine alive until a kill signal arrives. Block on the + // channel (passed in directly, so this goroutine always owns its own channel + // even if a reconnect replaces the map entry) instead of polling — this parks + // the goroutine with zero CPU and no per-second mutex access. + <-kill + log.Info().Str("userid", userID).Msg("Received kill signal") + client.Disconnect() + clientManager.DeleteWhatsmeowClient(userID) + clientManager.DeleteMyClient(userID) + clientManager.DeleteHTTPClient(userID) + if _, err := s.db.Exec(`UPDATE users SET qrcode='', connected=0 WHERE id=$1`, userID); err != nil { + log.Error().Err(err).Msg("failed to mark user disconnected on kill") + } + deleteKillChannel(userID, kill) +} + +func fileToBase64(filepath string) (string, string, error) { + data, err := os.ReadFile(filepath) + if err != nil { + return "", "", err + } + mimeType := http.DetectContentType(data) + return base64.StdEncoding.EncodeToString(data), mimeType, nil +} + +func (mycli *MyClient) myEventHandler(rawEvt interface{}) { + txtid := mycli.userID + postmap := make(map[string]interface{}) + postmap["event"] = rawEvt + dowebhook := 0 + path := "" + + switch evt := rawEvt.(type) { + case *events.AppStateSyncComplete: + if len(mycli.WAClient.Store.PushName) > 0 && evt.Name == appstate.WAPatchCriticalBlock { + err := mycli.WAClient.SendPresence(context.Background(), types.PresenceAvailable) + if err != nil { + log.Warn().Err(err).Msg("Failed to send available presence") + } else { + log.Info().Msg("Marked self as available") + } + } + case *events.Connected, *events.PushNameSetting: + postmap["type"] = "Connected" + dowebhook = 1 + if len(mycli.WAClient.Store.PushName) == 0 { + break + } + // Send presence available when connecting and when the pushname is changed. + // This makes sure that outgoing messages always have the right pushname. + err := mycli.WAClient.SendPresence(context.Background(), types.PresenceAvailable) + if err != nil { + log.Warn().Err(err).Msg("Failed to send available presence") + } else { + log.Info().Msg("Marked self as available") + } + sqlStmt := `UPDATE users SET connected=1 WHERE id=$1` + _, err = mycli.db.Exec(sqlStmt, mycli.userID) + if err != nil { + log.Error().Err(err).Msg(sqlStmt) + return + } + if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { + connectedJID := mycli.WAClient.Store.ID.ToNonAD().String() + if _, err := mycli.db.Exec(`UPDATE users SET jid=$1 WHERE id=$2`, connectedJID, mycli.userID); err != nil { + log.Warn().Err(err).Str("user_id", mycli.userID).Msg("Failed to persist JID on connect") + } else if myuserinfo, found := userinfocache.Get(mycli.token); found { + v := updateUserInfo(myuserinfo, "Jid", connectedJID) + userinfocache.Set(mycli.token, v, cache.NoExpiration) + } + } + case *events.PairSuccess: + log.Info().Str("userid", mycli.userID).Str("token", mycli.token).Str("ID", evt.ID.String()).Str("BusinessName", evt.BusinessName).Str("Platform", evt.Platform).Msg("QR Pair Success") + jidStr := evt.ID.String() + if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { + jidStr = mycli.WAClient.Store.ID.ToNonAD().String() + } + sqlStmt := `UPDATE users SET jid=$1 WHERE id=$2` + _, err := mycli.db.Exec(sqlStmt, jidStr, mycli.userID) + if err != nil { + log.Error().Err(err).Msg(sqlStmt) + return + } + + postmap["type"] = "PairSuccess" + dowebhook = 1 + + myuserinfo, found := userinfocache.Get(mycli.token) + if !found { + log.Warn().Msg("No user info cached on pairing?") + } else { + txtid = myuserinfo.(Values).Get("Id") + token := myuserinfo.(Values).Get("Token") + v := updateUserInfo(myuserinfo, "Jid", jidStr) + userinfocache.Set(token, v, cache.NoExpiration) + log.Info().Str("jid", jidStr).Str("userid", txtid).Str("token", token).Msg("User information set") + } + + // Check if automatic history sync is enabled and trigger it after QR code is scanned + var daysToSyncHistory int + query := "SELECT COALESCE(days_to_sync_history, 0) FROM users WHERE id=$1" + query = mycli.db.Rebind(query) + err = mycli.db.Get(&daysToSyncHistory, query, mycli.userID) + if err != nil { + log.Warn().Err(err).Str("userID", mycli.userID).Msg("Failed to get days_to_sync_history from database") + } else if daysToSyncHistory > 0 { + // Trigger history sync in a goroutine to avoid blocking + // Wait a bit for the connection to be fully established + go func() { + time.Sleep(2 * time.Second) // Give WhatsApp time to fully establish connection + + log.Info(). + Str("userID", mycli.userID). + Int("days", daysToSyncHistory). + Msg("Triggering automatic history sync after QR code scan") + + // Use the SyncWhatsAppHistory logic but for a single user + // Calculate message count based on days (estimate: 15 messages per day) + count := daysToSyncHistory * 15 + if count > 500 { + count = 500 // WhatsApp limit + } + if count < 50 { + count = 50 // Minimum reasonable count + } + + // Get chats from WhatsApp (contacts and groups) + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + var chatJIDs []string + + // Get all contacts + contacts, err := mycli.WAClient.Store.Contacts.GetAllContacts(ctx) + if err != nil { + log.Error().Err(err).Str("userID", mycli.userID).Msg("Failed to get contacts for history sync") + } else { + for jid := range contacts { + chatJIDs = append(chatJIDs, jid.String()) + } + } + + // Get all groups + groups, err := mycli.WAClient.GetJoinedGroups(ctx) + if err != nil { + log.Error().Err(err).Str("userID", mycli.userID).Msg("Failed to get groups for history sync") + } else { + for _, group := range groups { + chatJIDs = append(chatJIDs, group.JID.String()) + } + } + + // Sync each chat with a small delay between requests + for _, chatJIDStr := range chatJIDs { + chatJID, err := types.ParseJID(chatJIDStr) + if err != nil { + log.Warn().Err(err).Str("chatJID", chatJIDStr).Msg("Failed to parse chat JID, skipping") + continue + } + + // Use the syncHistoryForChat function from handlers.go + err = mycli.s.syncHistoryForChat(context.Background(), mycli.userID, chatJID, count) + if err != nil { + log.Warn().Err(err).Str("chatJID", chatJIDStr).Msg("Failed to sync history for chat") + } else { + log.Info().Str("chatJID", chatJIDStr).Int("count", count).Msg("History sync request sent for chat") + } + + // Small delay between requests to avoid overwhelming WhatsApp + time.Sleep(100 * time.Millisecond) + } + + log.Info(). + Str("userID", mycli.userID). + Int("days", daysToSyncHistory). + Int("chatsSynced", len(chatJIDs)). + Msg("Automatic history sync completed after QR code scan") + }() + } + case *events.StreamReplaced: + log.Info().Msg("Received StreamReplaced event") + return + case *events.Message: + + var s3Config struct { + Enabled string `db:"s3_enabled"` + MediaDelivery string `db:"media_delivery"` + } + + lastMessageCache.Set(mycli.userID, &evt.Info, cache.DefaultExpiration) + myuserinfo, found := userinfocache.Get(mycli.token) + if !found { + err := mycli.db.Get(&s3Config, "SELECT CASE WHEN s3_enabled = 1 THEN 'true' ELSE 'false' END AS s3_enabled, media_delivery FROM users WHERE id = $1", txtid) + if err != nil { + log.Error().Err(err).Msg("onMessage Failed to get S3 config from DB as it was not on cache") + s3Config.Enabled = "false" + s3Config.MediaDelivery = "base64" + } + } else { + s3Config.Enabled = myuserinfo.(Values).Get("S3Enabled") + s3Config.MediaDelivery = myuserinfo.(Values).Get("MediaDelivery") + } + + // Lazy init S3 client if needed (handles reconnect-after-restart when connectOnStartup skipped this user) + if s3Config.Enabled == "true" && (s3Config.MediaDelivery == "s3" || s3Config.MediaDelivery == "both") { + ensureS3ClientForUser(txtid) + } + + postmap["type"] = "Message" + dowebhook = 1 + metaParts := []string{fmt.Sprintf("pushname: %s", evt.Info.PushName), fmt.Sprintf("timestamp: %s", evt.Info.Timestamp)} + if evt.Info.Type != "" { + metaParts = append(metaParts, fmt.Sprintf("type: %s", evt.Info.Type)) + } + if evt.Info.Category != "" { + metaParts = append(metaParts, fmt.Sprintf("category: %s", evt.Info.Category)) + } + if evt.IsViewOnce { + metaParts = append(metaParts, "view once") + } + if evt.IsViewOnce { + metaParts = append(metaParts, "ephemeral") + } + + log.Info().Str("id", evt.Info.ID).Str("source", evt.Info.SourceString()).Str("parts", strings.Join(metaParts, ", ")).Msg("Message Received") + + // If this is a poll vote, decrypt the E2E-encrypted payload so the + // webhook can expose which options were selected. Votes arrive as + // SHA-256 hashes of the option text; we match those back to the + // plaintext options remembered at send time (see SendPoll in + // handlers.go). If the session was restarted between send and vote + // we cannot resolve plaintext; hashes are still emitted so the + // consumer can perform matching itself if it has stored options. + if evt.Message.GetPollUpdateMessage() != nil { + pollMsgID := evt.Message.GetPollUpdateMessage().GetPollCreationMessageKey().GetID() + + pollVote, perr := mycli.WAClient.DecryptPollVote(context.Background(), evt) + if perr != nil { + log.Warn().Err(perr).Str("pollMsgID", pollMsgID).Msg("DecryptPollVote failed") + } + + if perr == nil && pollVote != nil { + hashes := pollVote.GetSelectedOptions() + hashB64 := make([]string, 0, len(hashes)) + for _, h := range hashes { + hashB64 = append(hashB64, base64.StdEncoding.EncodeToString(h)) + } + + selected := make([]string, 0, len(hashes)) + if stored := clientManager.GetPollOptions(mycli.userID, pollMsgID); len(stored) > 0 { + optionsByHash := make(map[string]string, len(stored)) + for _, opt := range stored { + sum := sha256.Sum256([]byte(opt)) + optionsByHash[string(sum[:])] = opt + } + for _, h := range hashes { + if opt, found := optionsByHash[string(h)]; found { + selected = append(selected, opt) + } + } + } + + postmap["pollVote"] = map[string]interface{}{ + "pollCreationMsgID": pollMsgID, + "selectedOptions": selected, + "selectedHashesB64": hashB64, + } + } + } + + if encMessage := evt.Message.GetSecretEncryptedMessage(); encMessage != nil { + decrypted, derr := mycli.WAClient.DecryptSecretEncryptedMessage(context.Background(), evt) + if derr != nil { + log.Warn(). + Err(derr). + Str("messageID", evt.Info.ID). + Str("secretEncType", encMessage.GetSecretEncType().String()). + Msg("DecryptSecretEncryptedMessage failed") + } else if decrypted != nil { + log.Info(). + Str("messageID", evt.Info.ID). + Str("secretEncType", encMessage.GetSecretEncType().String()). + Msg("Decrypted secretEncryptedMessage; swapping evt.Message") + evt.Message = decrypted + } + } + + if !*skipMedia { + + isIncoming := !evt.Info.IsFromMe + chatJID := evt.Info.Sender.String() + if evt.Info.IsGroup { + chatJID = evt.Info.Chat.String() + } + + s3cfg := mediaS3Config{ + Enabled: s3Config.Enabled, + MediaDelivery: s3Config.MediaDelivery, + } + + if img := evt.Message.GetImageMessage(); img != nil { + mycli.processMedia(img, img.GetMimetype(), ".jpg", + downloadTimeoutImage, isIncoming, chatJID, + evt.Info.ID, s3cfg, postmap, nil) + } + + if audio := evt.Message.GetAudioMessage(); audio != nil { + mycli.processMedia(audio, audio.GetMimetype(), ".ogg", + downloadTimeoutAudio, isIncoming, chatJID, + evt.Info.ID, s3cfg, postmap, nil) + } + + if doc := evt.Message.GetDocumentMessage(); doc != nil { + ext := ".bin" + if doc.FileName != nil { + ext = filepath.Ext(*doc.FileName) + } + mycli.processMedia(doc, doc.GetMimetype(), ext, + downloadTimeoutDocument, isIncoming, chatJID, + evt.Info.ID, s3cfg, postmap, nil) + } + + if video := evt.Message.GetVideoMessage(); video != nil { + mycli.processMedia(video, video.GetMimetype(), ".mp4", + downloadTimeoutVideo, isIncoming, chatJID, + evt.Info.ID, s3cfg, postmap, nil) + } + + if sticker := evt.Message.GetStickerMessage(); sticker != nil { + mycli.processMedia(sticker, sticker.GetMimetype(), ".webp", + downloadTimeoutSticker, isIncoming, chatJID, + evt.Info.ID, s3cfg, postmap, map[string]interface{}{ + "isSticker": true, + "stickerAnimated": sticker.GetIsAnimated(), + }) + } + } + + // Save message to history regardless of skipMedia setting + // Get user's history setting from cache + var historyLimit int + userinfo, found := userinfocache.Get(mycli.token) + if found { + historyStr := userinfo.(Values).Get("History") + historyLimit, _ = strconv.Atoi(historyStr) + } else { + log.Warn().Str("userID", mycli.userID).Msg("User info not found in cache, skipping history") + historyLimit = 0 + } + + if historyLimit > 0 { + messageType := "text" + textContent := "" + mediaLink := "" + caption := "" + replyToMessageID := "" + + // Check for delete messages first + if protocolMsg := evt.Message.GetProtocolMessage(); protocolMsg != nil && protocolMsg.GetType() == 0 { + messageType = "delete" + if protocolMsg.GetKey() != nil { + textContent = protocolMsg.GetKey().GetID() // Store the deleted message ID + } + log.Info().Str("deletedMessageID", textContent).Str("messageID", evt.Info.ID).Msg("Delete message detected") + // Check for reactions + } else if reaction := evt.Message.GetReactionMessage(); reaction != nil { + messageType = "reaction" + replyToMessageID = reaction.GetKey().GetID() + textContent = reaction.GetText() // This will be the emoji + } else if img := evt.Message.GetImageMessage(); img != nil { + messageType = "image" + caption = img.GetCaption() + } else if video := evt.Message.GetVideoMessage(); video != nil { + messageType = "video" + caption = video.GetCaption() + } else if audio := evt.Message.GetAudioMessage(); audio != nil { + messageType = "audio" + } else if doc := evt.Message.GetDocumentMessage(); doc != nil { + messageType = "document" + caption = doc.GetCaption() + } else if sticker := evt.Message.GetStickerMessage(); sticker != nil { + messageType = "sticker" + } else if contact := evt.Message.GetContactMessage(); contact != nil { + messageType = "contact" + textContent = contact.GetDisplayName() + } else if location := evt.Message.GetLocationMessage(); location != nil { + messageType = "location" + textContent = location.GetName() + } + + // Extract text content for non-reaction and non-delete messages + if messageType != "reaction" && messageType != "delete" { + if conv := evt.Message.GetConversation(); conv != "" { + textContent = conv + } else if ext := evt.Message.GetExtendedTextMessage(); ext != nil { + textContent = ext.GetText() + // Check if this is a reply to another message + if contextInfo := ext.GetContextInfo(); contextInfo != nil && contextInfo.GetStanzaID() != "" { + replyToMessageID = contextInfo.GetStanzaID() + } + } else { + textContent = caption + } + + // Set default text content for media messages without captions + if textContent == "" { + switch messageType { + case "image": + textContent = ":image:" + case "video": + textContent = ":video:" + case "audio": + textContent = ":audio:" + case "document": + textContent = ":document:" + case "sticker": + textContent = ":sticker:" + case "contact": + if textContent == "" { + textContent = ":contact:" + } + case "location": + if textContent == "" { + textContent = ":location:" + } + } + } + } + + // Check for replies in regular conversation messages too + if messageType == "text" && replyToMessageID == "" { + // For regular text messages, check if there's context info indicating a reply + // This might be available in the message context + if conv := evt.Message.GetConversation(); conv != "" { + // Check if the message has reply context (this depends on WhatsApp message structure) + // For now, we'll rely on ExtendedTextMessage for reply detection + } + } + + // Try to get media link from S3 data if available + if s3Data, ok := postmap["s3"].(map[string]interface{}); ok { + if url, ok := s3Data["url"].(string); ok { + mediaLink = url + } + } + + // Only save if there's meaningful content (including delete messages) + if textContent != "" || mediaLink != "" || (messageType != "text" && messageType != "reaction") || messageType == "delete" { + // Serializar evt para JSON + evtJSON, err := json.Marshal(evt) + if err != nil { + log.Error().Err(err).Msg("Failed to marshal event to JSON") + evtJSON = []byte("{}") + } + + err = mycli.s.saveMessageToHistory( + mycli.userID, + evt.Info.Chat.String(), + evt.Info.Sender.String(), + evt.Info.ID, + messageType, + textContent, + mediaLink, + replyToMessageID, + string(evtJSON), + ) + if err != nil { + log.Error().Err(err).Msg("Failed to save message to history") + } else { + err = mycli.s.trimMessageHistory(mycli.userID, evt.Info.Chat.String(), historyLimit) + if err != nil { + log.Error().Err(err).Msg("Failed to trim message history") + } + } + } else { + log.Debug().Str("messageType", messageType).Str("messageID", evt.Info.ID).Msg("Skipping empty message from history") + } + } + + case *events.Receipt: + postmap["type"] = "ReadReceipt" + dowebhook = 1 + //if evt.Type == events.ReceiptTypeRead || evt.Type == events.ReceiptTypeReadSelf { + if evt.Type == types.ReceiptTypeRead || evt.Type == types.ReceiptTypeReadSelf { + log.Info().Strs("id", evt.MessageIDs).Str("source", evt.SourceString()).Str("timestamp", fmt.Sprintf("%v", evt.Timestamp)).Msg("Message was read") + //if evt.Type == events.ReceiptTypeRead { + if evt.Type == types.ReceiptTypeRead { + postmap["state"] = "Read" + } else { + postmap["state"] = "ReadSelf" + } + //} else if evt.Type == events.ReceiptTypeDelivered { + } else if evt.Type == types.ReceiptTypeDelivered { + postmap["state"] = "Delivered" + log.Info().Str("id", evt.MessageIDs[0]).Str("source", evt.SourceString()).Str("timestamp", fmt.Sprintf("%v", evt.Timestamp)).Msg("Message delivered") + } else { + // Discard webhooks for inactive or other delivery types + return + } + case *events.Presence: + postmap["type"] = "Presence" + dowebhook = 1 + if evt.Unavailable { + postmap["state"] = "offline" + if evt.LastSeen.IsZero() { + log.Info().Str("from", evt.From.String()).Msg("User is now offline") + } else { + log.Info().Str("from", evt.From.String()).Str("lastSeen", fmt.Sprintf("%v", evt.LastSeen)).Msg("User is now offline") + } + } else { + postmap["state"] = "online" + log.Info().Str("from", evt.From.String()).Msg("User is now online") + } + case *events.HistorySync: + postmap["type"] = "HistorySync" + dowebhook = 1 + + // Save HistorySync messages to message_history table + if evt.Data != nil && evt.Data.Conversations != nil { + go func() { + + // Get the account owner's JID for messages sent by the instance + accountOwnerJID := "" + if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { + accountOwnerJID = mycli.WAClient.Store.ID.ToNonAD().String() + } + + savedCount := 0 + for _, conv := range evt.Data.Conversations { + if conv == nil || conv.ID == nil || conv.Messages == nil { + continue + } + + chatJID, err := types.ParseJID(*conv.ID) + if err != nil { + log.Warn().Err(err).Str("convID", *conv.ID).Msg("Failed to parse conversation JID in HistorySync") + continue + } + + for _, msg := range conv.Messages { + if msg == nil || msg.Message == nil { + continue + } + + // Extract message data + messageKey := msg.Message.GetKey() + if messageKey == nil { + continue + } + + messageID := messageKey.GetID() + if messageID == "" { + continue + } + + // Determine sender - never use "me", always use actual JID + // Use GetFromMe() from MessageKey to determine if message is from account owner + // This is more reliable than checking GetParticipant() + isFromMe := messageKey.GetFromMe() + var senderJID string + + if isFromMe { + // Message from account owner + senderJID = accountOwnerJID + if senderJID == "" { + // Fallback: use "me" if account owner JID is not available + senderJID = "me" + log.Warn().Str("messageID", messageID).Msg("accountOwnerJID is not available for a message from me, using 'me' as senderJID") + } + } else { + // Message from someone else + participantJID := messageKey.GetParticipant() + if chatJID.Server == types.GroupServer || chatJID.Server == types.BroadcastServer { + // Group message: use participant JID + senderJID = participantJID + } else { + // Direct message: sender is the chat itself (chat_jid) + senderJID = chatJID.String() + } + } + + // If senderJID is still empty, skip this message + if senderJID == "" { + log.Warn().Str("messageID", messageID).Msg("Cannot determine sender JID, skipping message") + continue + } + + // Get message content + message := msg.Message.GetMessage() + if message == nil { + continue + } + + // Extract message type and content + messageType := "unknown" + textContent := "" + mediaLink := "" + quotedMessageID := "" + + if message.GetConversation() != "" { + messageType = "text" + textContent = message.GetConversation() + } else if ext := message.GetExtendedTextMessage(); ext != nil { + messageType = "text" + textContent = ext.GetText() + if contextInfo := ext.GetContextInfo(); contextInfo != nil { + quotedMessageID = contextInfo.GetStanzaID() + } + } else if img := message.GetImageMessage(); img != nil { + messageType = "image" + textContent = img.GetCaption() + } else if vid := message.GetVideoMessage(); vid != nil { + messageType = "video" + textContent = vid.GetCaption() + } else if audio := message.GetAudioMessage(); audio != nil { + messageType = "audio" + } else if doc := message.GetDocumentMessage(); doc != nil { + messageType = "document" + textContent = doc.GetCaption() + } else if sticker := message.GetStickerMessage(); sticker != nil { + messageType = "sticker" + } else if location := message.GetLocationMessage(); location != nil { + messageType = "location" + textContent = location.GetName() + } else if contact := message.GetContactMessage(); contact != nil { + messageType = "contact" + textContent = contact.GetDisplayName() + } else if buttons := message.GetButtonsResponseMessage(); buttons != nil { + messageType = "buttons_response" + textContent = buttons.GetSelectedButtonID() + } else if list := message.GetListResponseMessage(); list != nil { + messageType = "list_response" + textContent = list.GetSingleSelectReply().GetSelectedRowID() + } else if reaction := message.GetReactionMessage(); reaction != nil { + messageType = "reaction" + textContent = reaction.GetText() + if key := reaction.GetKey(); key != nil { + quotedMessageID = key.GetID() + } + } + + // Set default text for media messages without captions + if textContent == "" && messageType != "text" && messageType != "reaction" && messageType != "delete" { + switch messageType { + case "image": + textContent = ":image:" + case "video": + textContent = ":video:" + case "audio": + textContent = ":audio:" + case "document": + textContent = ":document:" + case "sticker": + textContent = ":sticker:" + case "contact": + textContent = ":contact:" + case "location": + textContent = ":location:" + } + } + + // Get message timestamp + msgTimestamp := time.Now() + if timestamp := msg.Message.GetMessageTimestamp(); timestamp > 0 { + msgTimestamp = time.Unix(int64(timestamp), 0) + } + + // Parse sender JID for MessageInfo + var senderJIDForInfo types.JID + if isFromMe { + if accountOwnerJID != "" { + var pErr error + senderJIDForInfo, pErr = types.ParseJID(accountOwnerJID) + if pErr != nil { + log.Warn().Err(pErr).Str("accountOwnerJID", accountOwnerJID).Msg("Failed to parse account owner JID in HistorySync") + } + } + } else { + if chatJID.Server == types.GroupServer || chatJID.Server == types.BroadcastServer { + // Group: use participant JID + participant := messageKey.GetParticipant() + if participant != "" { + var pErr error + senderJIDForInfo, pErr = types.ParseJID(participant) + if pErr != nil { + log.Warn().Err(pErr).Str("participantJID", participant).Msg("Failed to parse participant JID in HistorySync") + } + } + } else { + // Direct message: sender is the chat + senderJIDForInfo = chatJID + } + } + + // Try to get PushName from store if available + pushName := "" + if !isFromMe && senderJIDForInfo.User != "" { + if mycli.WAClient != nil && mycli.WAClient.Store != nil { + if contact, err := mycli.WAClient.Store.Contacts.GetContact(context.Background(), senderJIDForInfo); err == nil { + pushName = contact.PushName + } + } + } + + // Create MessageInfo structure matching events.Message format + messageInfo := types.MessageInfo{ + MessageSource: types.MessageSource{ + Chat: chatJID, + Sender: senderJIDForInfo, + IsFromMe: isFromMe, + IsGroup: chatJID.Server == types.GroupServer || chatJID.Server == types.BroadcastServer, + }, + ID: messageID, + Timestamp: msgTimestamp, + Type: messageType, + PushName: pushName, + } + + // Create events.Message-like structure for datajson + // This matches the format used in regular message events + // RawMessage should be the full waE2E.Message structure + messageEvent := map[string]interface{}{ + "Info": messageInfo, + "Message": message, + "IsEphemeral": false, + "IsViewOnce": false, + "IsViewOnceV2": false, + "IsViewOnceV2Extension": false, + "IsDocumentWithCaption": false, + "IsLottieSticker": false, + "IsBotInvoke": false, + "IsEdit": false, + "SourceWebMsg": nil, + "UnavailableRequestID": "", + "RetryCount": 0, + "NewsletterMeta": nil, + "RawMessage": msg.Message, + } + + // Serialize to JSON for datajson field + evtJSON, err := json.Marshal(messageEvent) + if err != nil { + log.Error().Err(err).Msg("Failed to marshal HistorySync message event to JSON") + evtJSON = []byte("{}") + } + + // Save message to history + // Only save if there's meaningful content + if textContent != "" || mediaLink != "" || (messageType != "text" && messageType != "reaction") { + err = mycli.s.saveMessageToHistory( + mycli.userID, + chatJID.String(), + senderJID, + messageID, + messageType, + textContent, + mediaLink, + quotedMessageID, + string(evtJSON), + ) + if err != nil { + log.Error().Err(err). + Str("userID", mycli.userID). + Str("chatJID", chatJID.String()). + Str("messageID", messageID). + Msg("Failed to save HistorySync message to history") + } else { + savedCount++ + } + } + } + } + + if savedCount > 0 { + log.Info(). + Str("userID", mycli.userID). + Int("savedCount", savedCount). + Msg("Saved HistorySync messages to message_history") + } + }() + } + + case *events.AppState: + log.Info().Str("index", fmt.Sprintf("%+v", evt.Index)).Str("actionValue", fmt.Sprintf("%+v", evt.SyncActionValue)).Msg("App state event received") + case *events.LoggedOut: + postmap["type"] = "LoggedOut" + dowebhook = 1 + log.Info().Str("reason", evt.Reason.String()).Msg("Logged out") + defer func() { + // Use a non-blocking send to prevent a deadlock if the receiver has already terminated. + signalKill(mycli.userID) + }() + sqlStmt := `UPDATE users SET connected=0 WHERE id=$1` + _, err := mycli.db.Exec(sqlStmt, mycli.userID) + if err != nil { + log.Error().Err(err).Msg(sqlStmt) + return + } + case *events.ChatPresence: + postmap["type"] = "ChatPresence" + dowebhook = 1 + log.Info().Str("state", fmt.Sprintf("%s", evt.State)).Str("media", fmt.Sprintf("%s", evt.Media)).Str("chat", evt.MessageSource.Chat.String()).Str("sender", evt.MessageSource.Sender.String()).Msg("Chat Presence received") + case *events.CallOffer: + postmap["type"] = "CallOffer" + dowebhook = 1 + log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call offer") + case *events.CallAccept: + postmap["type"] = "CallAccept" + dowebhook = 1 + log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call accept") + case *events.CallTerminate: + postmap["type"] = "CallTerminate" + dowebhook = 1 + log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call terminate") + case *events.CallOfferNotice: + postmap["type"] = "CallOfferNotice" + dowebhook = 1 + log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call offer notice") + case *events.CallRelayLatency: + postmap["type"] = "CallRelayLatency" + dowebhook = 1 + log.Info().Str("event", fmt.Sprintf("%+v", evt)).Msg("Got call relay latency") + case *events.Disconnected: + postmap["type"] = "Disconnected" + dowebhook = 1 + log.Info().Str("reason", fmt.Sprintf("%+v", evt)).Msg("Disconnected from Whatsapp") + case *events.ConnectFailure: + postmap["type"] = "ConnectFailure" + dowebhook = 1 + log.Error().Str("reason", fmt.Sprintf("%+v", evt)).Msg("Failed to connect to Whatsapp") + case *events.UndecryptableMessage: + postmap["type"] = "UndecryptableMessage" + dowebhook = 1 + log.Warn().Str("info", evt.Info.SourceString()).Msg("Undecryptable message received") + case *events.MediaRetry: + postmap["type"] = "MediaRetry" + dowebhook = 1 + log.Info().Str("messageID", evt.MessageID).Msg("Media retry event") + case *events.GroupInfo: + postmap["type"] = "GroupInfo" + dowebhook = 1 + log.Info().Str("jid", evt.JID.String()).Msg("Group info updated") + case *events.JoinedGroup: + postmap["type"] = "JoinedGroup" + dowebhook = 1 + log.Info().Str("jid", evt.JID.String()).Msg("Joined group") + case *events.Picture: + postmap["type"] = "Picture" + dowebhook = 1 + log.Info().Str("jid", evt.JID.String()).Msg("Picture updated") + case *events.BlocklistChange: + postmap["type"] = "BlocklistChange" + dowebhook = 1 + log.Info().Str("jid", evt.JID.String()).Msg("Blocklist changed") + case *events.Blocklist: + postmap["type"] = "Blocklist" + dowebhook = 1 + log.Info().Msg("Blocklist received") + case *events.KeepAliveRestored: + postmap["type"] = "KeepAliveRestored" + dowebhook = 1 + log.Info().Msg("Keep alive restored") + case *events.KeepAliveTimeout: + postmap["type"] = "KeepAliveTimeout" + dowebhook = 1 + log.Warn().Msg("Keep alive timeout") + case *events.ClientOutdated: + postmap["type"] = "ClientOutdated" + dowebhook = 1 + log.Warn().Msg("Client outdated") + case *events.TemporaryBan: + postmap["type"] = "TemporaryBan" + dowebhook = 1 + log.Info().Msg("Temporary ban") + case *events.StreamError: + postmap["type"] = "StreamError" + dowebhook = 1 + log.Error().Str("code", evt.Code).Msg("Stream error") + case *events.PairError: + postmap["type"] = "PairError" + dowebhook = 1 + log.Error().Msg("Pair error") + case *events.PrivacySettings: + postmap["type"] = "PrivacySettings" + dowebhook = 1 + log.Info().Msg("Privacy settings updated") + case *events.UserAbout: + postmap["type"] = "UserAbout" + dowebhook = 1 + log.Info().Str("jid", evt.JID.String()).Msg("User about updated") + case *events.OfflineSyncCompleted: + postmap["type"] = "OfflineSyncCompleted" + dowebhook = 1 + log.Info().Msg("Offline sync completed") + case *events.OfflineSyncPreview: + postmap["type"] = "OfflineSyncPreview" + dowebhook = 1 + log.Info().Msg("Offline sync preview") + case *events.IdentityChange: + postmap["type"] = "IdentityChange" + dowebhook = 1 + log.Info().Str("jid", evt.JID.String()).Msg("Identity changed") + case *events.NewsletterJoin: + postmap["type"] = "NewsletterJoin" + dowebhook = 1 + log.Info().Str("jid", evt.ID.String()).Msg("Newsletter joined") + case *events.NewsletterLeave: + postmap["type"] = "NewsletterLeave" + dowebhook = 1 + log.Info().Str("jid", evt.ID.String()).Msg("Newsletter left") + case *events.NewsletterMuteChange: + postmap["type"] = "NewsletterMuteChange" + dowebhook = 1 + log.Info().Str("jid", evt.ID.String()).Msg("Newsletter mute changed") + case *events.NewsletterLiveUpdate: + postmap["type"] = "NewsletterLiveUpdate" + dowebhook = 1 + log.Info().Msg("Newsletter live update") + case *events.FBMessage: + postmap["type"] = "FBMessage" + dowebhook = 1 + log.Info().Str("info", evt.Info.SourceString()).Msg("Facebook message received") + default: + log.Warn().Str("event", fmt.Sprintf("%+v", evt)).Msg("Unhandled event") + } + + if dowebhook == 1 { + sendEventWithWebHook(mycli, postmap, path) + } +} From 9afeb6bed7a6520913c14b23dcd341acb8a0d64d Mon Sep 17 00:00:00 2001 From: AntonKun Date: Mon, 10 Aug 2026 21:52:40 +0300 Subject: [PATCH 3/3] Fix compatibility issues --- handlers.go | 19 +++++++++++-------- wmiau.go | 5 +++-- 2 files changed, 14 insertions(+), 10 deletions(-) diff --git a/handlers.go b/handlers.go index 1f8f2d27..8026410c 100644 --- a/handlers.go +++ b/handlers.go @@ -757,15 +757,16 @@ func (s *server) PairPhone() http.HandlerFunc { } // resolveSessionJID returns the best-known JID for a user. When logged in, the -func (s *server) resolveSessionJID(txtid string, waClient *whatsmeow.Client, userInfo Values) string { +// live whatsmeow store is authoritative — cache/DB often lag after QR pairing. +func (s *server) resolveSessionJID(ctx context.Context, txtid string, waClient *whatsmeow.Client, userInfo Values) string { jid := userInfo.Get("Jid") - if waClient != nil && waClient.IsLoggedIn() && waClient.Store != nil && waClient.Store.ID != nil { + if waClient != nil && waClient.Store != nil && waClient.IsLoggedIn() && waClient.Store.ID != nil { storeJID := waClient.Store.ID.ToNonAD() if storeJID.Server == types.HiddenUserServer { - ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) - pn, err := getCachedPNForLID(ctx, waClient, storeJID) - cancel() + timeoutCtx, cancel := context.WithTimeout(ctx, 5*time.Second) + defer cancel() + pn, err := getCachedPNForLID(timeoutCtx, waClient, storeJID) if err == nil && !pn.IsEmpty() { storeJID = pn.ToNonAD() } @@ -779,14 +780,16 @@ func (s *server) resolveSessionJID(txtid string, waClient *whatsmeow.Client, use v := updateUserInfo(userInfo, "Jid", resolved) userinfocache.Set(token, v, cache.NoExpiration) } - if _, err := s.db.Exec("UPDATE users SET jid=$1 WHERE id=$2", resolved, txtid); err != nil { + query := s.db.Rebind("UPDATE users SET jid=? WHERE id=?") + if _, err := s.db.Exec(query, resolved, txtid); err != nil { log.Warn().Err(err).Str("user_id", txtid).Msg("Failed to persist resolved JID") } } } } else if jid == "" { var dbJid string - if err := s.db.QueryRow("SELECT jid FROM users WHERE id = $1", txtid).Scan(&dbJid); err == nil && dbJid != "" { + query := s.db.Rebind("SELECT jid FROM users WHERE id = ?") + if err := s.db.QueryRow(query, txtid).Scan(&dbJid); err == nil && dbJid != "" { jid = dbJid } } @@ -816,7 +819,7 @@ func (s *server) GetStatus() http.HandlerFunc { waClient := clientManager.GetWhatsmeowClient(txtid) isConnected := waClient != nil && waClient.IsConnected() isLoggedIn := waClient != nil && waClient.IsLoggedIn() - jid := s.resolveSessionJID(txtid, waClient, userInfo) + jid := s.resolveSessionJID(r.Context(), txtid, waClient, userInfo) var proxyURL string s.db.QueryRow("SELECT proxy_url FROM users WHERE id = $1", txtid).Scan(&proxyURL) diff --git a/wmiau.go b/wmiau.go index deadf3f8..28b21e55 100644 --- a/wmiau.go +++ b/wmiau.go @@ -828,7 +828,8 @@ func (mycli *MyClient) myEventHandler(rawEvt interface{}) { } if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { connectedJID := mycli.WAClient.Store.ID.ToNonAD().String() - if _, err := mycli.db.Exec(`UPDATE users SET jid=$1 WHERE id=$2`, connectedJID, mycli.userID); err != nil { + query := mycli.db.Rebind(`UPDATE users SET jid=? WHERE id=?`) + if _, err := mycli.db.Exec(query, connectedJID, mycli.userID); err != nil { log.Warn().Err(err).Str("user_id", mycli.userID).Msg("Failed to persist JID on connect") } else if myuserinfo, found := userinfocache.Get(mycli.token); found { v := updateUserInfo(myuserinfo, "Jid", connectedJID) @@ -841,7 +842,7 @@ func (mycli *MyClient) myEventHandler(rawEvt interface{}) { if mycli.WAClient.Store != nil && mycli.WAClient.Store.ID != nil { jidStr = mycli.WAClient.Store.ID.ToNonAD().String() } - sqlStmt := `UPDATE users SET jid=$1 WHERE id=$2` + sqlStmt := mycli.db.Rebind(`UPDATE users SET jid=? WHERE id=?`) _, err := mycli.db.Exec(sqlStmt, jidStr, mycli.userID) if err != nil { log.Error().Err(err).Msg(sqlStmt)