Spaces:
Running
Running
File size: 28,573 Bytes
83d6851 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 | package instance_service
import (
"bufio"
"context"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"slices"
"sort"
"strings"
"time"
"agentdeck-whatsapp-service/pkg/config"
instance_model "agentdeck-whatsapp-service/pkg/instance/model"
instance_repository "agentdeck-whatsapp-service/pkg/instance/repository"
event_types "agentdeck-whatsapp-service/pkg/internal/event_types"
logger_wrapper "agentdeck-whatsapp-service/pkg/logger"
"agentdeck-whatsapp-service/pkg/utils"
whatsmeow_service "agentdeck-whatsapp-service/pkg/whatsmeow/service"
"go.mau.fi/whatsmeow"
"go.mau.fi/whatsmeow/types"
)
type InstanceService interface {
Create(data *CreateStruct) (*instance_model.Instance, error)
Connect(data *ConnectStruct, instance *instance_model.Instance) (*instance_model.Instance, string, string, error)
Reconnect(instance *instance_model.Instance) error
Disconnect(instance *instance_model.Instance) (*instance_model.Instance, error)
Logout(instance *instance_model.Instance) (*instance_model.Instance, error)
Status(instance *instance_model.Instance) (*StatusStruct, error)
GetQr(instance *instance_model.Instance) (*QrcodeStruct, error)
Pair(data *PairStruct, instance *instance_model.Instance) (*PairReturnStruct, error)
GetAll() ([]*instance_model.Instance, error)
Info(instanceId string) (*instance_model.Instance, error)
Delete(id string) error
SetProxy(id string, proxyConfig *ProxyConfig) error
SetProxyFromStruct(id string, data *SetProxyStruct) error
RemoveProxy(id string) error
ForceReconnect(instanceId string, number string) error
GetInstanceByToken(token string) (*instance_model.Instance, error)
GetLogs(instanceId string, startDate, endDate time.Time, level string, limit int) ([]logger_wrapper.LogEntry, error)
GetAdvancedSettings(instanceId string) (*instance_model.AdvancedSettings, error)
UpdateAdvancedSettings(instanceId string, settings *instance_model.AdvancedSettings) error
}
type instances struct {
instanceRepository instance_repository.InstanceRepository
config *config.Config
killChannel map[string](chan bool)
clientPointer map[string]*whatsmeow.Client
whatsmeowService whatsmeow_service.WhatsmeowService
loggerWrapper *logger_wrapper.LoggerManager
}
type ProxyConfig struct {
Protocol string `json:"protocol,omitempty"`
Port string `json:"port"`
Password string `json:"password"`
Username string `json:"username"`
Host string `json:"host"`
}
type CreateStruct struct {
InstanceId string `json:"instanceId"`
Name string `json:"name"`
Token string `json:"token"`
Proxy *ProxyConfig `json:"proxy"`
AdvancedSettings *instance_model.AdvancedSettings `json:"advancedSettings"`
}
type ConnectStruct struct {
WebhookUrl string `json:"webhookUrl"`
Subscribe []string `json:"subscribe"`
Immediate bool `json:"immediate"`
Phone string `json:"phone"`
RabbitmqEnable string `json:"rabbitmqEnable"`
WebSocketEnable string `json:"websocketEnable"`
NatsEnable string `json:"natsEnable"`
}
type StatusStruct struct {
Connected bool
LoggedIn bool
myJid *types.JID
Name string
}
type QrcodeStruct struct {
Qrcode string `json:"qrcode"`
Code string `json:"code"`
// Passkey ceremony fields. Populated when the account requires a WebAuthn
// passkey to finish linking (no QR to scan at that point). The manager uses
// PasskeyStage to switch its UI and PasskeyOpenUrl for the
// "Abrir WhatsApp Web" button that launches the passkey ceremony.
PasskeyStage string `json:"passkeyStage,omitempty"`
PasskeyOpenURL string `json:"passkeyOpenUrl,omitempty"`
PasskeyCode string `json:"passkeyCode,omitempty"`
}
type PairStruct struct {
Subscribe []string `json:"subscribe"`
Phone string `json:"phone"`
}
type PairReturnStruct struct {
PairingCode string
}
type SetProxyStruct struct {
Protocol string `json:"protocol,omitempty"`
Host string `json:"host" validate:"required"`
Port string `json:"port" validate:"required"`
Username string `json:"username"`
Password string `json:"password"`
}
type ForceReconnectStruct struct {
Number string `json:"number"`
}
func (i *instances) ensureClientConnected(instanceId string) (*whatsmeow.Client, error) {
logger := i.loggerWrapper.GetLogger(instanceId)
client := i.clientPointer[instanceId]
logger.LogInfo("[%s] Checking client connection status - Client exists: %v", instanceId, client != nil)
if client == nil {
logger.LogInfo("[%s] No client found, attempting to start new instance", instanceId)
err := i.whatsmeowService.StartInstance(instanceId)
if err != nil {
logger.LogError("[%s] Failed to start instance: %v", instanceId, err)
return nil, errors.New("no active session found")
}
logger.LogInfo("[%s] Instance started, waiting 2 seconds...", instanceId)
time.Sleep(2 * time.Second)
client = i.clientPointer[instanceId]
logger.LogInfo("[%s] Checking new client - Exists: %v, Connected: %v",
instanceId,
client != nil,
client != nil && client.IsConnected())
if client == nil || !client.IsConnected() {
logger.LogError("[%s] New client validation failed - Exists: %v, Connected: %v",
instanceId,
client != nil,
client != nil && client.IsConnected())
return nil, errors.New("no active session found")
}
} else if !client.IsConnected() {
logger.LogError("[%s] Existing client is disconnected - Connected status: %v",
instanceId,
client.IsConnected())
return nil, errors.New("client disconnected")
}
logger.LogInfo("[%s] Client successfully validated - Connected: %v", instanceId, client.IsConnected())
return client, nil
}
func (i instances) Create(data *CreateStruct) (*instance_model.Instance, error) {
if data.Proxy != nil {
data.Proxy.Protocol = utils.NormalizeProxyProtocol(data.Proxy.Protocol, data.Proxy.Port)
}
proxyJson, err := json.Marshal(data.Proxy)
if err != nil {
return nil, err
}
findInstance, _ := i.instanceRepository.GetInstanceByName(data.Name)
if findInstance != nil {
return nil, fmt.Errorf("instance already exists")
}
instance := instance_model.Instance{
Id: data.InstanceId,
Name: data.Name,
Token: data.Token,
OsName: i.config.OsName,
Proxy: string(proxyJson),
Connected: false,
ClientName: i.config.ClientName,
}
// Set advanced settings if provided
if data.AdvancedSettings != nil {
instance.AlwaysOnline = data.AdvancedSettings.AlwaysOnline
instance.RejectCall = data.AdvancedSettings.RejectCall
instance.MsgRejectCall = data.AdvancedSettings.MsgRejectCall
instance.ReadMessages = data.AdvancedSettings.ReadMessages
instance.IgnoreGroups = data.AdvancedSettings.IgnoreGroups
instance.IgnoreStatus = data.AdvancedSettings.IgnoreStatus
}
createdInstance, err := i.instanceRepository.Create(instance)
if err != nil {
return nil, err
}
return createdInstance, nil
}
func (i instances) Connect(data *ConnectStruct, instance *instance_model.Instance) (*instance_model.Instance, string, string, error) {
var subscribedEvents []string
i.loggerWrapper.GetLogger(instance.Id).LogInfo("[%s] Processing subscribe events: %v", instance.Id, data.Subscribe)
if len(data.Subscribe) == 0 {
subscribedEvents = append(subscribedEvents, event_types.MESSAGE)
} else if len(data.Subscribe) > 0 && data.Subscribe[0] == "ALL" {
for _, event := range event_types.AllEventTypes {
subscribedEvents = append(subscribedEvents, event)
}
} else {
for _, arg := range data.Subscribe {
if !event_types.IsEventType(arg) {
i.loggerWrapper.GetLogger(instance.Id).LogWarn("[%s] Message type discarded '%s'", instance.Id, arg)
continue
}
subscribedEvents = append(subscribedEvents, arg)
}
}
eventString := strings.Join(subscribedEvents, ",")
instance.Events = eventString
instance.Webhook = data.WebhookUrl
instance.RabbitmqEnable = data.RabbitmqEnable
instance.NatsEnable = data.NatsEnable
instance.WebSocketEnable = data.WebSocketEnable
err := i.instanceRepository.Update(instance)
if err != nil {
i.loggerWrapper.GetLogger(instance.Id).LogError("[%s] Error updating instance: %s", instance.Id, err)
return nil, "", "", err
}
// Verifica se a instância já está rodando
isInstanceRunning := i.clientPointer[instance.Id] != nil
// Sincroniza as configurações na instância em execução (se já estiver conectada)
err = i.whatsmeowService.UpdateInstanceSettings(instance.Id)
if err != nil {
i.loggerWrapper.GetLogger(instance.Id).LogInfo("[%s] Instance not in runtime yet, will be updated when connected", instance.Id)
isInstanceRunning = false
} else {
i.loggerWrapper.GetLogger(instance.Id).LogInfo("[%s] Instance settings updated successfully in runtime", instance.Id)
isInstanceRunning = true
}
// Se a instância não estiver rodando, inicia uma nova
if !isInstanceRunning {
i.loggerWrapper.GetLogger(instance.Id).LogInfo("[%s] Starting new client instance", instance.Id)
i.killChannel[instance.Id] = make(chan bool)
clientData := &whatsmeow_service.ClientData{
Instance: instance,
Subscriptions: subscribedEvents,
Phone: data.Phone,
IsProxy: false,
}
if instance.Proxy != "" || i.config.ProxyHost != "" {
var proxyConfig ProxyConfig
err := json.Unmarshal([]byte(instance.Proxy), &proxyConfig)
if err != nil {
i.loggerWrapper.GetLogger(instance.Id).LogError("[%s] error unmarshalling proxy config: %v", instance.Id, err)
return nil, "", "", err
}
if proxyConfig.Host != "" || i.config.ProxyHost != "" {
clientData.IsProxy = true
}
}
go i.whatsmeowService.StartClient(clientData)
} else {
i.loggerWrapper.GetLogger(instance.Id).LogInfo("[%s] Instance already running, settings updated without restarting client", instance.Id)
}
// logger.LogInfo("Waiting 1 seconds")
// time.Sleep(1000 * time.Millisecond)
// if i.clientPointer[instance.Id] != nil {
// if !i.clientPointer[instance.Id].IsConnected() {
// return instance, "", "", fmt.Errorf("failed to connect")
// }
// } else {
// return instance, "", "", fmt.Errorf("failed to connect")
// }
return instance, instance.Jid, eventString, nil
}
func (i instances) Reconnect(instance *instance_model.Instance) error {
_, err := i.ensureClientConnected(instance.Id)
if err != nil {
return err
}
return i.whatsmeowService.ReconnectClient(instance.Id)
}
func (i instances) Disconnect(instance *instance_model.Instance) (*instance_model.Instance, error) {
client, err := i.ensureClientConnected(instance.Id)
if err != nil {
return instance, err
}
if client.IsConnected() {
if client.IsLoggedIn() {
i.loggerWrapper.GetLogger(instance.Id).LogInfo("[%s] Disconnection successful", instance.Id)
i.killChannel[instance.Id] <- true
instance.Events = ""
err := i.instanceRepository.Update(instance)
if err != nil {
return instance, err
}
return instance, nil
}
}
i.loggerWrapper.GetLogger(instance.Id).LogWarn("[%s] Ignoring disconnect as it was not connected", instance.Id)
return instance, nil
}
func (i instances) Logout(instance *instance_model.Instance) (*instance_model.Instance, error) {
client, err := i.ensureClientConnected(instance.Id)
if err != nil {
return instance, err
}
if client.IsLoggedIn() && client.IsConnected() {
err := client.Logout(context.Background())
if err != nil {
return instance, err
}
instance.Connected = false
err = i.instanceRepository.Update(instance)
if err != nil {
return instance, err
}
select {
case i.killChannel[instance.Id] <- true:
case <-time.After(5 * time.Second):
}
delete(i.clientPointer, instance.Id)
delete(i.killChannel, instance.Id)
i.loggerWrapper.GetLogger(instance.Id).LogInfo("[%s] Logout successful", instance.Id)
return instance, nil
}
if client.IsConnected() {
client.Disconnect()
select {
case i.killChannel[instance.Id] <- true:
case <-time.After(5 * time.Second):
}
delete(i.clientPointer, instance.Id)
delete(i.killChannel, instance.Id)
i.loggerWrapper.GetLogger(instance.Id).LogInfo("[%s] Disconnection successful", instance.Id)
return instance, nil
}
i.loggerWrapper.GetLogger(instance.Id).LogWarn("[%s] Ignoring logout as it was not connected", instance.Id)
return instance, fmt.Errorf("ignoring logout as it was not connected")
}
func (i instances) Status(instance *instance_model.Instance) (*StatusStruct, error) {
client := i.clientPointer[instance.Id]
if client == nil {
return &StatusStruct{
Connected: false,
LoggedIn: false,
}, nil
}
isConnected := client.IsConnected()
isLoggedIn := client.IsLoggedIn()
var myJid *types.JID
var name string
if isLoggedIn {
myJid = client.Store.ID
name = client.Store.PushName
}
return &StatusStruct{
Connected: isConnected,
LoggedIn: isLoggedIn,
myJid: myJid,
Name: name,
}, nil
}
func (i instances) GetQr(instance *instance_model.Instance) (*QrcodeStruct, error) {
logger := i.loggerWrapper.GetLogger(instance.Id)
client := i.clientPointer[instance.Id]
// Se não há cliente ou o cliente está logado, precisamos iniciar um novo cliente
if client == nil || client.IsLoggedIn() {
if client != nil && client.IsLoggedIn() {
logger.LogInfo("[%s] Client is logged in, starting new instance for QR code", instance.Id)
} else {
logger.LogInfo("[%s] No client found, starting new instance for QR code", instance.Id)
}
// Iniciar nova instância para gerar QR code
err := i.whatsmeowService.StartInstance(instance.Id)
if err != nil {
logger.LogError("[%s] Failed to start instance: %v", instance.Id, err)
return nil, fmt.Errorf("failed to start instance: %w", err)
}
// Aguardar um pouco para o cliente iniciar e gerar QR code
logger.LogInfo("[%s] Waiting for QR code generation...", instance.Id)
time.Sleep(3 * time.Second)
// Verificar novamente se há cliente
client = i.clientPointer[instance.Id]
if client != nil && client.IsLoggedIn() {
return nil, fmt.Errorf("session already logged in")
}
} else if !client.IsConnected() {
// Se o cliente existe mas não está conectado, pode estar aguardando QR code
logger.LogInfo("[%s] Client exists but not connected, checking for existing QR code", instance.Id)
}
// Buscar instância atualizada do banco para pegar o QR code mais recente
instance, err := i.instanceRepository.GetInstanceByID(instance.Id)
if err != nil {
return nil, err
}
// If a passkey ceremony is in progress, there is no QR to scan — return the
// passkey stage + the #wapk openUrl so the manager can render the
// "Abrir WhatsApp Web" button. Checked before the empty-QR branch because
// during a passkey ceremony instance.Qrcode is empty.
if store := i.whatsmeowService.PasskeyCeremonyStore(); store != nil {
if token, state, ok := store.StateByInstance(instance.Id); ok {
logger.LogInfo("[%s] Passkey ceremony active (stage=%s) — returning passkey info instead of QR", instance.Id, state.Stage)
return &QrcodeStruct{
PasskeyStage: state.Stage,
PasskeyCode: state.Code,
PasskeyOpenURL: buildPasskeyOpenURL(token),
}, nil
}
}
code := instance.Qrcode
if code == "" {
// Se não há QR code ainda, aguardar um pouco mais e tentar novamente
logger.LogInfo("[%s] No QR code available yet, waiting a bit more...", instance.Id)
time.Sleep(2 * time.Second)
instance, err = i.instanceRepository.GetInstanceByID(instance.Id)
if err != nil {
return nil, err
}
code = instance.Qrcode
if code == "" {
return nil, fmt.Errorf("no QR code available. Please wait a moment and try again")
}
}
parts := strings.Split(code, "|")
if len(parts) < 2 {
return nil, fmt.Errorf("invalid QR code format")
}
qr := &QrcodeStruct{
Qrcode: parts[0],
Code: parts[1],
}
return qr, nil
}
// buildPasskeyOpenURL builds the URL the manager opens to start the passkey
// ceremony: https://web.whatsapp.com/#wapk=<base64url({t:token,b:publicBase})>.
// publicBase must be the PUBLICLY reachable API base the browser can hit; set it
// via PASSKEY_PUBLIC_URL. Kept in sync with the event handler in whatsmeow.go.
func buildPasskeyOpenURL(token string) string {
publicBase := os.Getenv("PASSKEY_PUBLIC_URL")
if publicBase == "" {
publicBase = "<SET_PASSKEY_PUBLIC_URL>"
}
payload := fmt.Sprintf(`{"t":%q,"b":%q}`, token, publicBase)
wapk := base64.RawURLEncoding.EncodeToString([]byte(payload))
return "https://web.whatsapp.com/#wapk=" + wapk
}
func (i instances) Pair(data *PairStruct, instance *instance_model.Instance) (*PairReturnStruct, error) {
logger := i.loggerWrapper.GetLogger(instance.Id)
client := i.clientPointer[instance.Id]
if client == nil || !client.IsConnected() {
if client != nil && client.IsLoggedIn() {
return nil, fmt.Errorf("instance is already authenticated")
}
logger.LogInfo("[%s] No active connection, starting instance for phone pairing", instance.Id)
if err := i.whatsmeowService.StartInstance(instance.Id); err != nil {
logger.LogError("[%s] Failed to start instance for pairing: %v", instance.Id, err)
return nil, fmt.Errorf("failed to start instance: %w", err)
}
// Wait for the WA websocket connection and initial QR generation to establish.
// PairPhone must be called after the QR event is received per whatsmeow docs.
time.Sleep(3 * time.Second)
client = i.clientPointer[instance.Id]
if client == nil {
return nil, fmt.Errorf("failed to initialize client for pairing")
}
}
if client.IsLoggedIn() {
return nil, fmt.Errorf("instance is already authenticated")
}
code, err := client.PairPhone(context.Background(), data.Phone, true, whatsmeow.PairClientChrome, "Chrome (Linux)")
if err != nil {
logger.LogError("[%s] PairPhone failed: %v", instance.Id, err)
return nil, fmt.Errorf("pairing failed: %w", err)
}
return &PairReturnStruct{PairingCode: code}, nil
}
func (i instances) GetAll() ([]*instance_model.Instance, error) {
instances, err := i.instanceRepository.GetAll(i.config.ClientName)
if err != nil {
return nil, err
}
for _, instance := range instances {
if client := i.clientPointer[instance.Id]; client != nil {
instance.Connected = client.IsLoggedIn()
} else {
instance.Connected = false
}
instance.Proxy = ""
}
return instances, nil
}
func (i instances) Info(instanceId string) (*instance_model.Instance, error) {
instance, err := i.instanceRepository.GetInstanceByID(instanceId)
if err != nil {
return nil, err
}
// Atualiza o status connected com base no estado real do cliente
if client := i.clientPointer[instance.Id]; client != nil {
instance.Connected = client.IsLoggedIn()
} else {
instance.Connected = false
}
instance.Proxy = ""
return instance, nil
}
func (i instances) Delete(id string) error {
instance, err := i.instanceRepository.GetInstanceByID(id)
if err != nil {
return err
}
if i.clientPointer[instance.Id] != nil && i.clientPointer[instance.Id].IsConnected() {
if i.clientPointer[instance.Id].IsLoggedIn() {
i.clientPointer[instance.Id].Logout(context.Background())
}
i.clientPointer[instance.Id].Disconnect()
}
// Limpar todos os recursos da instância antes de deletar
delete(i.clientPointer, instance.Id)
if i.killChannel[instance.Id] != nil {
close(i.killChannel[instance.Id])
delete(i.killChannel, instance.Id)
}
// Limpar cache via whatsmeow service
err = i.whatsmeowService.ClearInstanceCache(instance.Id, instance.Token)
if err != nil {
i.loggerWrapper.GetLogger(instance.Id).LogWarn("[%s] Failed to clear instance cache: %v", instance.Id, err)
}
err = i.instanceRepository.Delete(id)
if err != nil {
return err
}
return nil
}
func (i instances) SetProxy(id string, proxyConfig *ProxyConfig) error {
instance, err := i.instanceRepository.GetInstanceByID(id)
if err != nil {
return err
}
// Validate proxy configuration
if proxyConfig == nil {
return fmt.Errorf("proxy configuration cannot be nil")
}
if proxyConfig.Host == "" {
return fmt.Errorf("proxy host is required")
}
if proxyConfig.Port == "" {
return fmt.Errorf("proxy port is required")
}
proxyConfig.Protocol = utils.NormalizeProxyProtocol(proxyConfig.Protocol, proxyConfig.Port)
// Convert proxy config to JSON
proxyJSON, err := json.Marshal(proxyConfig)
if err != nil {
i.loggerWrapper.GetLogger(id).LogError("[%s] Failed to marshal proxy config: %v", id, err)
return fmt.Errorf("failed to marshal proxy configuration: %v", err)
}
instance.Proxy = string(proxyJSON)
// Update instance in database
err = i.instanceRepository.Update(instance)
if err != nil {
i.loggerWrapper.GetLogger(id).LogError("[%s] Failed to update instance with proxy: %v", id, err)
return err
}
i.loggerWrapper.GetLogger(id).LogInfo("[%s] Proxy configuration updated: %s://%s:%s", id, proxyConfig.Protocol, proxyConfig.Host, proxyConfig.Port)
// Reconnect to apply proxy changes
go i.Reconnect(instance)
return nil
}
func (i instances) SetProxyFromStruct(id string, data *SetProxyStruct) error {
if data == nil {
return fmt.Errorf("proxy data cannot be nil")
}
proxyConfig := &ProxyConfig{
Protocol: data.Protocol,
Host: data.Host,
Port: data.Port,
Username: data.Username,
Password: data.Password,
}
return i.SetProxy(id, proxyConfig)
}
func (i instances) RemoveProxy(id string) error {
instance, err := i.instanceRepository.GetInstanceByID(id)
if err != nil {
return err
}
instance.Proxy = ""
err = i.instanceRepository.Update(instance)
if err != nil {
return err
}
i.loggerWrapper.GetLogger(id).LogInfo("[%s] Proxy configuration removed", id)
go i.Reconnect(instance)
return nil
}
func (i instances) ForceReconnect(instanceId string, number string) error {
if i.clientPointer[instanceId].IsConnected() && i.clientPointer[instanceId].IsLoggedIn() {
return fmt.Errorf("client already connected")
}
err := i.whatsmeowService.ForceUpdateJid(instanceId, number)
if err != nil {
return err
}
instance, err := i.instanceRepository.GetInstanceByID(instanceId)
if err != nil {
return err
}
subscribedEvents := strings.Split(instance.Events, ",")
i.killChannel[instance.Id] = make(chan bool)
clientData := &whatsmeow_service.ClientData{
Instance: instance,
Subscriptions: subscribedEvents,
Phone: "",
IsProxy: false,
}
if instance.Proxy != "" || i.config.ProxyHost != "" {
var proxyConfig ProxyConfig
err := json.Unmarshal([]byte(instance.Proxy), &proxyConfig)
if err != nil {
i.loggerWrapper.GetLogger(instance.Id).LogError("[%s] error unmarshalling proxy config: %v", instance.Id, err)
return err
}
if proxyConfig.Host != "" || i.config.ProxyHost != "" {
clientData.IsProxy = true
}
}
if i.clientPointer[instance.Id] != nil {
client := i.clientPointer[instance.Id]
client.Disconnect()
select {
case i.killChannel[instance.Id] <- true:
case <-time.After(5 * time.Second):
}
delete(i.clientPointer, instance.Id)
delete(i.killChannel, instance.Id)
}
go i.whatsmeowService.StartClient(clientData)
time.Sleep(2 * time.Second)
if i.clientPointer[instance.Id] != nil {
if !i.clientPointer[instance.Id].IsConnected() {
return fmt.Errorf("failed to connect")
}
if !i.clientPointer[instance.Id].IsLoggedIn() {
return fmt.Errorf("failed to login")
}
} else {
return fmt.Errorf("failed to connect")
}
return nil
}
func (i instances) GetInstanceByToken(token string) (*instance_model.Instance, error) {
return i.instanceRepository.GetInstanceByToken(token)
}
func (i instances) GetLogs(instanceId string, startDate, endDate time.Time, level string, limit int) ([]logger_wrapper.LogEntry, error) {
// Inicializa o slice vazio para garantir que nunca retorne null
logs := make([]logger_wrapper.LogEntry, 0)
// Define valores padrão
if limit <= 0 {
limit = 100 // Limite padrão de 100 registros
}
// Se não foi fornecida data inicial, usa 7 dias atrás
if startDate.IsZero() {
startDate = time.Now().AddDate(0, 0, -7)
}
// Se não foi fornecida data final, usa data atual
if endDate.IsZero() {
endDate = time.Now()
}
// Ajusta as datas para início e fim do dia
startDate = time.Date(startDate.Year(), startDate.Month(), startDate.Day(), 0, 0, 0, 0, time.UTC)
endDate = time.Date(endDate.Year(), endDate.Month(), endDate.Day(), 23, 59, 59, 999999999, time.UTC)
// Garante que a data inicial não seja posterior à data final
if startDate.After(endDate) {
return logs, fmt.Errorf("data inicial não pode ser posterior à data final")
}
// Níveis de log válidos
validLevels := map[string]bool{
"INFO": true,
"ERROR": true,
"WARN": true,
"DEBUG": true,
}
var levelArray []string
if level == "" {
// Se nenhum nível foi especificado, usa todos
levelArray = []string{"INFO", "ERROR", "WARN", "DEBUG"}
} else {
// Divide e normaliza os níveis fornecidos
for _, l := range strings.Split(level, ",") {
l = strings.TrimSpace(strings.ToUpper(l))
if !validLevels[l] {
return logs, fmt.Errorf("nível de log inválido: %s", l)
}
levelArray = append(levelArray, l)
}
}
// Lê os logs do arquivo
logPath := filepath.Join(i.config.LogDirectory, instanceId, "instance.log")
file, err := os.Open(logPath)
if err != nil {
if os.IsNotExist(err) {
return logs, nil // Retorna array vazio se arquivo não existir
}
return logs, fmt.Errorf("erro ao abrir arquivo de log: %v", err)
}
defer file.Close()
scanner := bufio.NewScanner(file)
// Aumenta o buffer do scanner para lidar com linhas grandes
const maxCapacity = 1024 * 1024 // 1MB
buf := make([]byte, maxCapacity)
scanner.Buffer(buf, maxCapacity)
for scanner.Scan() {
var entry logger_wrapper.LogEntry
if err := json.Unmarshal(scanner.Bytes(), &entry); err != nil {
continue // Ignora linhas inválidas
}
// Ajusta o timestamp da entrada para UTC para comparação correta
entry.Timestamp = entry.Timestamp.UTC()
// Aplica os filtros
if entry.Timestamp.Before(startDate) || entry.Timestamp.After(endDate) {
continue
}
if !slices.Contains(levelArray, entry.Level) {
continue
}
logs = append(logs, entry)
// Verifica o limite
if len(logs) >= limit {
break
}
}
if err := scanner.Err(); err != nil {
return logs, fmt.Errorf("erro ao ler arquivo de log: %v", err)
}
// Ordena os logs por timestamp em ordem decrescente
sort.Slice(logs, func(i, j int) bool {
return logs[i].Timestamp.After(logs[j].Timestamp)
})
return logs, nil
}
func (i instances) GetAdvancedSettings(instanceId string) (*instance_model.AdvancedSettings, error) {
i.loggerWrapper.GetLogger(instanceId).LogInfo("[%s] Getting advanced settings", instanceId)
settings, err := i.instanceRepository.GetAdvancedSettings(instanceId)
if err != nil {
i.loggerWrapper.GetLogger(instanceId).LogError("[%s] Error getting advanced settings: %v", instanceId, err)
return nil, err
}
return settings, nil
}
func (i instances) UpdateAdvancedSettings(instanceId string, settings *instance_model.AdvancedSettings) error {
i.loggerWrapper.GetLogger(instanceId).LogInfo("[%s] Updating advanced settings", instanceId)
err := i.instanceRepository.UpdateAdvancedSettings(instanceId, settings)
if err != nil {
i.loggerWrapper.GetLogger(instanceId).LogError("[%s] Error updating advanced settings: %v", instanceId, err)
return err
}
// Sincroniza as configurações na instância em execução
err = i.whatsmeowService.UpdateInstanceAdvancedSettings(instanceId)
if err != nil {
i.loggerWrapper.GetLogger(instanceId).LogWarn("[%s] Error syncing advanced settings to runtime: %v", instanceId, err)
// Não falha a operação, apenas loga o warning
}
i.loggerWrapper.GetLogger(instanceId).LogInfo("[%s] Advanced settings updated successfully", instanceId)
return nil
}
func NewInstanceService(
instanceRepository instance_repository.InstanceRepository,
killChannel map[string](chan bool),
clientPointer map[string]*whatsmeow.Client,
whatsmeowService whatsmeow_service.WhatsmeowService,
config *config.Config,
loggerWrapper *logger_wrapper.LoggerManager,
) InstanceService {
return &instances{
instanceRepository: instanceRepository,
killChannel: killChannel,
clientPointer: clientPointer,
whatsmeowService: whatsmeowService,
config: config,
loggerWrapper: loggerWrapper,
}
}
|