Files

689 lines
19 KiB
Go
Raw Permalink Normal View History

package node
import (
"context"
2021-06-30 13:40:54 +02:00
"database/sql"
"errors"
"fmt"
"os"
"path/filepath"
"reflect"
"sync"
"sync/atomic"
"time"
errorspkg "github.com/pkg/errors"
"go.uber.org/zap"
2020-01-02 10:10:19 +01:00
"github.com/ethereum/go-ethereum/event"
gethrpc "github.com/ethereum/go-ethereum/rpc"
2020-01-02 10:10:19 +01:00
accsmanagement "github.com/status-im/status-go/internal/accounts-management"
"github.com/status-im/status-go/internal/connection"
"github.com/status-im/status-go/internal/crypto"
"github.com/status-im/status-go/internal/db/multiaccounts"
"github.com/status-im/status-go/internal/db/multiaccounts/accounts"
2025-12-16 12:54:23 +00:00
"github.com/status-im/status-go/internal/ipfs"
2026-08-19 20:49:00 +01:00
"github.com/status-im/status-go/internal/panics"
"github.com/status-im/status-go/internal/pausable"
"github.com/status-im/status-go/internal/rpc"
2025-10-16 17:36:49 +01:00
"github.com/status-im/status-go/internal/timesource"
"github.com/status-im/status-go/internal/transactions"
"github.com/status-im/status-go/params"
noderpc "github.com/status-im/status-go/pkg/backend/node/rpc"
2025-06-17 01:02:34 -03:00
"github.com/status-im/status-go/pkg/pubsub"
2026-08-20 21:44:51 +01:00
accountssvc "github.com/status-im/status-go/pkg/services/accounts"
appgeneral "github.com/status-im/status-go/pkg/services/app-general"
"github.com/status-im/status-go/pkg/services/backup"
"github.com/status-im/status-go/pkg/services/browsers"
"github.com/status-im/status-go/pkg/services/chat"
"github.com/status-im/status-go/pkg/services/communitytokens"
"github.com/status-im/status-go/pkg/services/connector"
"github.com/status-im/status-go/pkg/services/ens"
"github.com/status-im/status-go/pkg/services/eth"
"github.com/status-im/status-go/pkg/services/gif"
"github.com/status-im/status-go/pkg/services/linkpreview"
localnotifications "github.com/status-im/status-go/pkg/services/local-notifications"
"github.com/status-im/status-go/pkg/services/media"
"github.com/status-im/status-go/pkg/services/networks"
2026-08-20 21:44:51 +01:00
"github.com/status-im/status-go/pkg/services/newsfeed"
"github.com/status-im/status-go/pkg/services/permissions"
"github.com/status-im/status-go/pkg/services/personal"
"github.com/status-im/status-go/pkg/services/preferences"
"github.com/status-im/status-go/pkg/services/rpcstats"
"github.com/status-im/status-go/pkg/services/sharedurls"
"github.com/status-im/status-go/pkg/services/stickers"
"github.com/status-im/status-go/pkg/services/storagestats"
2026-08-20 21:44:51 +01:00
"github.com/status-im/status-go/pkg/services/updates"
"github.com/status-im/status-go/pkg/services/utils"
"github.com/status-im/status-go/pkg/services/wakuv2ext"
"github.com/status-im/status-go/pkg/services/wallet"
"github.com/status-im/status-go/pkg/services/wallet/community"
"github.com/status-im/status-go/pkg/services/wallet/pendingtxtracker"
"github.com/status-im/status-go/pkg/services/wallet/token"
"github.com/status-im/status-go/pkg/services/wallet/tokenbalances"
)
// errors
var (
ErrNodeRunning = errors.New("node is already running")
ErrNoRunningNode = errors.New("there is no running node")
)
// StatusNode abstracts contained geth node and provides helper methods to
// interact with it.
type StatusNode struct {
mu sync.RWMutex
2021-06-30 13:40:54 +02:00
appDB *sql.DB
multiaccountsDB *multiaccounts.Database
walletDB *sql.DB
2021-06-30 13:40:54 +02:00
running atomic.Bool
config *params.NodeConfig // Status node configuration
rpcClient *rpc.Client // reference to an RPC client
networkManager *networks.Manager
2026-08-19 20:49:00 +01:00
services []StatusService
rpcServer *gethrpc.Server
downloader *ipfs.Downloader
2025-09-10 10:19:04 +01:00
mediaServerAddress *string
mediaServerAdvertizeHost string
mediaServerAdvertizePort int
mediaServerEnableTLS *bool
2026-08-19 20:50:03 +01:00
mediaServer *media.Server
tokenManager *token.Manager
logger *zap.Logger
2021-06-30 13:40:54 +02:00
gethAccountsManager *accsmanagement.AccountsManager
transactor *transactions.Transactor
2021-06-30 13:40:54 +02:00
publicMethods map[string]bool
2021-07-07 08:11:09 +02:00
// we explicitly list every service, we could use interfaces
// and store them in a nicer way and user reflection, but for now stupid is good
2021-06-30 13:40:54 +02:00
rpcStatsSrvc *rpcstats.Service
accountsSrvc *accountssvc.Service
browsersSrvc *browsers.Service
permissionsSrvc *permissions.Service
walletSrvc *wallet.Service
networksSrvc *networks.Service
2021-06-30 13:40:54 +02:00
localNotificationsSrvc *localnotifications.Service
personalSrvc *personal.Service
2025-10-16 17:36:49 +01:00
timeSourceSrvc timesource.Service
2021-06-30 13:40:54 +02:00
wakuV2ExtSrvc *wakuv2ext.Service
2021-12-21 16:05:09 +01:00
ensSrvc *ens.Service
communityTokensSrvc *communitytokens.Service
2022-01-31 13:58:03 +01:00
gifSrvc *gif.Service
2022-02-02 18:50:55 -04:00
stickersSrvc *stickers.Service
2022-02-10 11:15:27 -04:00
chatSrvc *chat.Service
updatesSrvc *updates.Service
pendingTracker *pendingtxtracker.PendingTxTracker
connectorSrvc *connector.Service
appGeneralSrvc *appgeneral.Service
ethSrvc *eth.Service
newsfeedSrvc *newsfeed.Service
preferencesSrvc *preferences.Service
sharedUrlsSrvc *sharedurls.Service
storageStatsSrvc *storagestats.Service
linkPreviewSrvc *linkpreview.Service
2025-06-17 01:02:34 -03:00
walletFeed event.Feed
accountsPublisher *pubsub.Publisher
localBackup *backup.Controller
2026-07-21 18:24:32 -04:00
tokenManagerStartDone chan struct{}
serviceRegistry *ServiceRegistry
}
// ServiceRegistry returns the node's service registry for granular pause/resume control.
func (n *StatusNode) ServiceRegistry() *ServiceRegistry {
return n.serviceRegistry
}
// New makes new instance of StatusNode.
func New(transactor *transactions.Transactor, gethAccountsManager *accsmanagement.AccountsManager, logger *zap.Logger) *StatusNode {
logger = logger.Named("StatusNode")
return &StatusNode{
transactor: transactor,
gethAccountsManager: gethAccountsManager,
logger: logger,
publicMethods: make(map[string]bool),
accountsPublisher: pubsub.NewPublisher(),
rpcServer: gethrpc.NewServer(),
serviceRegistry: newServiceRegistry(),
}
}
2018-04-16 14:36:09 +02:00
// Config exposes reference to running node's configuration
func (n *StatusNode) Config() *params.NodeConfig {
n.mu.RLock()
defer n.mu.RUnlock()
return n.config
}
2026-08-19 20:50:03 +01:00
func (n *StatusNode) MediaServer() *media.Server {
n.mu.RLock()
defer n.mu.RUnlock()
2025-09-10 10:19:04 +01:00
return n.mediaServer
}
2025-09-10 10:19:04 +01:00
func (n *StatusNode) startMediaServer() error {
if n.mediaServer != nil {
if err := n.mediaServer.Stop(); err != nil {
return err
}
}
2026-08-19 20:50:03 +01:00
var opts []media.MediaServerOption
if n.mediaServerEnableTLS != nil {
2026-08-19 20:50:03 +01:00
opts = append(opts, media.WithMediaServerDisableTLS(!*n.mediaServerEnableTLS))
}
2025-09-10 10:19:04 +01:00
if n.mediaServerAddress != nil {
2026-08-19 20:50:03 +01:00
opts = append(opts, media.WithMediaServerAddress(*n.mediaServerAddress))
2025-09-10 10:19:04 +01:00
}
2026-08-19 20:50:03 +01:00
opts = append(opts, media.WithMediaServerAdvertizeAddress(n.mediaServerAdvertizeHost, n.mediaServerAdvertizePort))
mediaServer, err := media.NewServer(nil, nil, n.multiaccountsDB, nil, opts...)
if err != nil {
return err
}
2025-09-10 10:19:04 +01:00
n.mediaServer = mediaServer
2025-09-10 10:19:04 +01:00
if err := n.mediaServer.Start(); err != nil {
return err
}
return nil
}
2025-09-10 10:19:04 +01:00
// StartMediaServerWithoutDB starts media server without starting the node
// The server can only handle requests that don't require appdb or IPFS downloader
func (n *StatusNode) StartMediaServerWithoutDB() error {
if n.IsRunning() {
n.logger.Debug("node is already running, no need to StartMediaServerWithoutDB")
return nil
}
return n.startMediaServer()
}
// StartWithOptions starts current StatusNode, failing if it's already started.
// It takes some options that allows to further configure starting process.
func (n *StatusNode) Start(config *params.NodeConfig) error {
n.mu.Lock()
defer n.mu.Unlock()
if !n.running.CompareAndSwap(false, true) {
n.logger.Debug("node is already running")
return ErrNodeRunning
}
n.logger.Debug("starting with options", zap.Stringer("ClusterConfig", &config.ClusterConfig))
return n.startWithDB(config)
}
func (n *StatusNode) StartLocalBackup() error {
if n.localBackup != nil {
return errors.New("local backup already started")
}
backupPath, err := n.accountsSrvc.GetBackupPath()
if err != nil {
return err
}
if backupPath == "" {
// No path set yet, set it to the user's config directory
dir, err := os.UserConfigDir()
// We do not return the error as it's not a major issue
if err != nil {
n.logger.Error("failed to get user config dir", zap.Error(err))
} else {
err = n.accountsSrvc.SetBackupPath(filepath.Join(dir, "Status", "backups"))
if err != nil {
n.logger.Error("failed to set backup path", zap.Error(err))
}
}
}
chatAccount, err := n.gethAccountsManager.SelectedChatAccount()
if err != nil {
return err
}
privateKey := chatAccount.PrivateKey()
n.localBackup, err = backup.NewController(backup.Config{
2025-12-05 12:41:17 +00:00
PrivateKey: crypto.Keccak256(crypto.FromECDSA(privateKey)),
FileNameProvider: n,
BackupEnabled: true,
Interval: time.Minute * 30,
}, n.logger.Named("LocalBackup"))
if err != nil {
return err
}
if n.accountsSrvc != nil {
n.localBackup.Register("settings", n.accountsSrvc)
}
if n.walletSrvc != nil {
n.localBackup.Register("wallet", n.walletSrvc)
}
2025-11-26 20:40:01 +01:00
if n.wakuV2ExtSrvc != nil {
n.localBackup.Register("messenger", n.wakuV2ExtSrvc.Messenger())
}
n.serviceRegistry.Register(n.localBackup)
n.localBackup.Start()
return nil
}
2025-12-05 12:41:17 +00:00
func (n *StatusNode) GetBackupFilename() (string, error) {
chatAccount, err := n.gethAccountsManager.SelectedChatAccount()
if err != nil {
return "", err
}
privateKey := chatAccount.PrivateKey()
backupPath, err := n.accountsSrvc.GetBackupPath()
if err != nil {
return "", err
}
compressedPubKey, err := utils.SerializePublicKey(crypto.CompressPubkey(&privateKey.PublicKey))
if err != nil {
return "", err
}
if backupPath == "" {
return "", errors.New("backup path is not set")
}
fullPath := filepath.Join(backupPath, fmt.Sprintf("%s_user_data.bkp", compressedPubKey[len(compressedPubKey)-6:]))
return fullPath, nil
}
func (n *StatusNode) PerformLocalBackup() (string, error) {
return n.localBackup.PerformBackup()
}
func (n *StatusNode) LoadLocalBackup(filePath string) error {
return n.localBackup.LoadBackup(filePath)
}
2025-09-10 10:19:04 +01:00
func (n *StatusNode) SetMediaServerOptions(address *string, enableTLS *bool, advertizeHost string, advertizePort int) {
2026-07-07 17:08:05 +04:00
n.mu.Lock()
defer n.mu.Unlock()
2025-09-10 10:19:04 +01:00
n.mediaServerAddress = address
n.mediaServerEnableTLS = enableTLS
2025-09-10 10:19:04 +01:00
n.mediaServerAdvertizeHost = advertizeHost
n.mediaServerAdvertizePort = advertizePort
if n.mediaServer != nil {
if err := n.startMediaServer(); err != nil {
n.logger.Error("failed to restart media server with updated options", zap.Error(err))
}
}
}
func (n *StatusNode) StopMediaServer() error {
2026-07-07 17:08:05 +04:00
n.mu.Lock()
defer n.mu.Unlock()
if n.mediaServer == nil {
return nil
}
err := n.mediaServer.Stop()
n.mediaServer = nil
return err
}
func (n *StatusNode) startWithDB(config *params.NodeConfig) error {
n.config = config
if err := n.setupRPCClient(); err != nil {
return err
}
n.downloader = ipfs.NewDownloader(config.RootDataDir)
2025-09-10 10:19:04 +01:00
if n.mediaServer == nil {
if err := n.startMediaServer(); err != nil {
return err
}
}
2025-09-10 10:19:04 +01:00
n.mediaServer.SetDataProviders(n.appDB, n.walletDB, n.downloader)
2026-07-21 18:24:32 -04:00
if err := n.createTokenManager(); err != nil {
return err
}
2026-08-03 14:57:04 -04:00
if err := n.tokenManager.Start(context.Background()); err != nil {
return errorspkg.Wrap(err, "failed to start token manager")
}
n.tokenManagerStartDone = make(chan struct{})
close(n.tokenManagerStartDone)
2025-09-10 10:19:04 +01:00
if err := n.initServices(config, n.mediaServer); err != nil {
2021-06-30 13:40:54 +02:00
return err
}
// Run migrations
err := n.runServicesMigrations()
if err != nil {
return errorspkg.Wrap(err, "failed to run services migrations")
}
// Register services
for _, service := range n.services {
2025-09-10 10:19:04 +01:00
err := n.registerService(service)
if err != nil {
name := reflect.TypeOf(service).Name()
text := fmt.Sprintf("failed to register service '%s'", name)
return errorspkg.Wrap(err, text)
}
}
// Start services
err = n.timeSourceSrvc.Start(context.Background())
if err != nil {
return errorspkg.Wrap(err, "failed to start time source")
}
for _, service := range n.services {
err := service.Start()
if err != nil {
name := reflect.TypeOf(service).Name()
text := fmt.Sprintf("failed to start service '%s'", name)
return errorspkg.Wrap(err, text)
}
}
n.populateServiceRegistry()
return nil
2018-04-16 14:36:09 +02:00
}
// populateServiceRegistry registers all services that implement Pausable and wraps the
// media server as a pausable. Called once after all services are started.
func (n *StatusNode) populateServiceRegistry() {
// Register services implementing Pausable
for _, service := range n.services {
2026-08-19 20:49:00 +01:00
if p, ok := service.(pausable.Pausable); ok {
n.serviceRegistry.Register(p)
}
}
// Wrap and register the media server
if n.mediaServer != nil {
n.serviceRegistry.Register(newPausableMediaServer(n.mediaServer))
}
// Wrap and register the messenger so that PauseServices/ResumeServices
// gate filter health-check pings and mailserver syncs in background.
if n.wakuV2ExtSrvc != nil {
if m := n.wakuV2ExtSrvc.Messenger(); m != nil {
n.serviceRegistry.Register(newPausableMessenger(m))
}
}
// Register infrastructure components that are not go-ethereum services
if n.downloader != nil {
n.serviceRegistry.Register(n.downloader)
}
2026-08-19 20:49:00 +01:00
if p, ok := n.timeSourceSrvc.(pausable.Pausable); ok {
n.serviceRegistry.Register(p)
}
}
2026-07-21 18:24:32 -04:00
func (n *StatusNode) createTokenManager() error {
const (
defaultAutoRefreshInterval = 30 * time.Minute // interval after which we should fetch the token lists from the remote source (or use the default one if remote source is not set)
defaultAutoRefreshCheckInterval = 3 * time.Minute // interval after which we should check if we should trigger the auto-refresh
)
autoRefreshInterval := defaultAutoRefreshInterval
autoRefreshCheckInterval := defaultAutoRefreshCheckInterval
if n.config.WalletConfig.TokensListsAutoRefreshInterval > 0 &&
n.config.WalletConfig.TokensListsAutoRefreshCheckInterval > 0 &&
n.config.WalletConfig.TokensListsAutoRefreshInterval > n.config.WalletConfig.TokensListsAutoRefreshCheckInterval {
autoRefreshInterval = time.Duration(n.config.WalletConfig.TokensListsAutoRefreshInterval) * time.Second
autoRefreshCheckInterval = time.Duration(n.config.WalletConfig.TokensListsAutoRefreshCheckInterval) * time.Second
}
accDB, err := accounts.NewDB(n.appDB)
if err != nil {
return err
}
n.tokenManager, err = token.NewTokenManager(n.walletDB, n.rpcClient, community.NewManager(n.appDB, n.mediaServer, nil),
n.rpcClient.GetNetworkManager(), n.appDB, n.mediaServer, &n.walletFeed, n.accountsPublisher, accDB,
autoRefreshInterval, autoRefreshCheckInterval)
if err != nil {
return err
}
2025-12-15 18:14:10 +08:00
// check for possible custom tokens in the config
if len(n.config.WalletConfig.CustomTokens) > 0 {
for _, token := range n.config.WalletConfig.CustomTokens {
err := n.tokenManager.UpsertCustom(*token)
if err != nil {
return err
}
}
}
2026-07-21 18:24:32 -04:00
return nil
}
func (n *StatusNode) StartTokenManager() {
n.mu.Lock()
if !n.running.Load() || n.tokenManager == nil || n.tokenManagerStartDone != nil {
n.mu.Unlock()
return
}
tokenManager := n.tokenManager
done := make(chan struct{})
n.tokenManagerStartDone = done
n.mu.Unlock()
go func() {
2026-08-19 20:49:00 +01:00
defer panics.LogOnPanic()
2026-07-21 18:24:32 -04:00
defer close(done)
if err := tokenManager.Start(context.Background()); err != nil {
n.logger.Error("failed to start token manager", zap.Error(err))
}
}()
}
// NetworkManager returns the network manager owned by the node.
func (n *StatusNode) NetworkManager() *networks.Manager {
return n.networkManager
}
2018-04-16 10:01:37 +02:00
func (n *StatusNode) setupRPCClient() (err error) {
networkManager := networks.NewManager(n.appDB, n.accountsPublisher)
if networkManager == nil {
return errorspkg.New("failed to create network manager")
}
if err = networkManager.InitEmbeddedNetworks(n.config.Networks); err != nil {
return errorspkg.Wrap(err, "network manager failed to initialize")
}
n.networkManager = networkManager
config := rpc.ClientConfig{
NetworkManager: networkManager,
2025-06-17 01:02:34 -03:00
AccountsPublisher: n.accountsPublisher,
}
n.rpcClient, err = rpc.NewClient(config)
if err != nil {
return
}
n.rpcClient.Start(context.Background())
2018-04-16 10:01:37 +02:00
return
}
// Stop will stop current StatusNode. A stopped node cannot be resumed.
func (n *StatusNode) Stop() error {
n.mu.Lock()
defer n.mu.Unlock()
2018-04-16 14:36:09 +02:00
n.logger.Debug("stopping")
if !n.running.CompareAndSwap(true, false) {
2018-04-16 14:36:09 +02:00
return ErrNoRunningNode
}
var errs []error
2026-07-21 18:24:32 -04:00
if n.tokenManagerStartDone != nil {
<-n.tokenManagerStartDone
n.tokenManagerStartDone = nil
2026-08-03 14:57:04 -04:00
n.tokenManager.Stop()
2026-07-21 18:24:32 -04:00
}
n.timeSourceSrvc.Stop()
for _, service := range n.services {
err := service.Stop()
errs = append(errs, err)
}
if n.localBackup != nil {
n.localBackup.Stop()
n.localBackup = nil
}
2025-06-17 01:02:34 -03:00
n.accountsPublisher.Close()
n.rpcClient.Stop()
n.rpcClient = nil
2018-04-16 14:36:09 +02:00
n.config = nil
if n.mediaServer != nil {
n.mediaServer.SetDataProviders(nil, nil, nil)
}
n.downloader.Stop()
n.downloader = nil
2021-07-07 08:11:09 +02:00
n.rpcStatsSrvc = nil
n.accountsSrvc = nil
n.browsersSrvc = nil
n.permissionsSrvc = nil
n.walletSrvc = nil
n.networksSrvc = nil
n.networkManager = nil
2021-07-07 08:11:09 +02:00
n.localNotificationsSrvc = nil
n.personalSrvc = nil
n.timeSourceSrvc = nil
n.wakuV2ExtSrvc = nil
2021-12-21 16:05:09 +01:00
n.ensSrvc = nil
n.communityTokensSrvc = nil
2022-02-02 18:50:55 -04:00
n.stickersSrvc = nil
n.connectorSrvc = nil
n.publicMethods = make(map[string]bool)
n.pendingTracker = nil
n.appGeneralSrvc = nil
n.newsfeedSrvc = nil
n.preferencesSrvc = nil
n.storageStatsSrvc = nil
n.logger.Debug("status node stopped")
return errors.Join(errs...)
}
// IsRunning confirm that node is running.
func (n *StatusNode) IsRunning() bool {
return n.running.Load()
}
// Pause reduces non-essential background work while preserving core messaging/Waku runtime.
// Delegates to the ServiceRegistry so callers can also pause individual services via PauseService.
func (n *StatusNode) Pause() error {
if n.serviceRegistry == nil {
return nil
}
return n.serviceRegistry.PauseAll()
}
// Resume restores regular service work after paused background mode.
// Delegates to the ServiceRegistry so callers can also resume individual services via ResumeService.
func (n *StatusNode) Resume() error {
if n.serviceRegistry == nil {
return nil
}
return n.serviceRegistry.ResumeAll()
}
func (n *StatusNode) CallInProcessRPC(inputJSON string) string {
codec := noderpc.NewSingleRequestCodec(inputJSON)
n.rpcServer.ServeCodec(codec.GethCodec(), 0)
return codec.Output()
}
// RPCClient exposes reference to RPC client connected to the running node.
func (n *StatusNode) RPCClient() *rpc.Client {
n.mu.RLock()
defer n.mu.RUnlock()
return n.rpcClient
}
2017-09-01 20:44:50 +02:00
2021-06-30 13:40:54 +02:00
func (n *StatusNode) SetAppDB(db *sql.DB) {
n.appDB = db
}
func (n *StatusNode) GetAppDB() *sql.DB {
return n.appDB
}
2021-06-30 13:40:54 +02:00
func (n *StatusNode) SetMultiaccountsDB(db *multiaccounts.Database) {
n.multiaccountsDB = db
}
func (n *StatusNode) SetWalletDB(db *sql.DB) {
n.walletDB = db
}
func (n *StatusNode) GetWalletDB() *sql.DB {
return n.walletDB
}
func (n *StatusNode) TokenManager() *token.Manager {
return n.tokenManager
}
func (n *StatusNode) TokenBalancesFetcher() tokenbalances.FetcherIface {
if n.walletSrvc != nil {
return n.walletSrvc.GetTokenBalancesFetcher()
}
return nil
}
func (n *StatusNode) TokenBalancesStorage() tokenbalances.Storage {
if n.walletSrvc != nil {
return n.walletSrvc.GetTokenBalancesStorage()
}
return nil
}
func (n *StatusNode) ConnectionChanged(state connection.State) {
if n.wakuV2ExtSrvc != nil {
n.wakuV2ExtSrvc.ConnectionChanged(state)
}
}