From 34c985a910810a232eb7abf5da5de2d34bea382a Mon Sep 17 00:00:00 2001 From: Shevchuk Date: Wed, 23 Jul 2025 21:21:33 +0700 Subject: [PATCH 1/3] [WIP] --- cmd/copytrader/server/config.go | 55 +++++++++ cmd/copytrader/server/main.go | 99 ++++++++++++++++ cmd/copytrader/worker/config.go | 66 +++++++++++ cmd/copytrader/worker/main.go | 128 +++++++++++++++++++++ go.mod | 1 + plugin/copytrader/copytrader.go | 193 ++++++++++++++++++++++++++++++++ plugin/copytrader/policy.go | 54 +++++++++ plugin/payroll/payroll.go | 5 +- plugin/payroll/transaction.go | 18 +-- 9 files changed, 608 insertions(+), 11 deletions(-) create mode 100644 cmd/copytrader/server/config.go create mode 100644 cmd/copytrader/server/main.go create mode 100644 cmd/copytrader/worker/config.go create mode 100644 cmd/copytrader/worker/main.go create mode 100644 plugin/copytrader/copytrader.go create mode 100644 plugin/copytrader/policy.go diff --git a/cmd/copytrader/server/config.go b/cmd/copytrader/server/config.go new file mode 100644 index 0000000..bb705ca --- /dev/null +++ b/cmd/copytrader/server/config.go @@ -0,0 +1,55 @@ +package main + +import ( + "fmt" + "os" + "strings" + + "github.com/spf13/viper" + "github.com/vultisig/verifier/vault_config" + + "github.com/vultisig/plugin/api" + "github.com/vultisig/plugin/storage" +) + +type CopytraderServerConfig struct { + Server api.ServerConfig `mapstructure:"server" json:"server"` + Database struct { + DSN string `mapstructure:"dsn" json:"dsn,omitempty"` + } `mapstructure:"database" json:"database,omitempty"` + BaseConfigPath string `mapstructure:"base_config_path" json:"base_config_path,omitempty"` + Redis storage.RedisConfig `mapstructure:"redis" json:"redis,omitempty"` + BlockStorage vault_config.BlockStorage `mapstructure:"block_storage" json:"block_storage,omitempty"` + Datadog struct { + Host string `mapstructure:"host" json:"host,omitempty"` + Port string `mapstructure:"port" json:"port,omitempty"` + } `mapstructure:"datadog" json:"datadog"` +} + +func GetConfigure() (*CopytraderServerConfig, error) { + configName := os.Getenv("VS_CONFIG_NAME") + if configName == "" { + configName = "config" + } + + return ReadConfig(configName) +} + +func ReadConfig(configName string) (*CopytraderServerConfig, error) { + viper.SetConfigName(configName) + viper.AddConfigPath(".") + viper.SetEnvKeyReplacer(strings.NewReplacer(".", "_")) + viper.AutomaticEnv() + + viper.SetDefault("Server.VaultsFilePath", "vaults") + + if err := viper.ReadInConfig(); err != nil { + return nil, fmt.Errorf("failed to read config file, %w", err) + } + var cfg CopytraderServerConfig + err := viper.Unmarshal(&cfg) + if err != nil { + return nil, fmt.Errorf("unable to decode into struct, %w", err) + } + return &cfg, nil +} diff --git a/cmd/copytrader/server/main.go b/cmd/copytrader/server/main.go new file mode 100644 index 0000000..73faeff --- /dev/null +++ b/cmd/copytrader/server/main.go @@ -0,0 +1,99 @@ +package main + +import ( + "fmt" + "net" + + "github.com/DataDog/datadog-go/statsd" + "github.com/hibiken/asynq" + "github.com/sirupsen/logrus" + "github.com/vultisig/verifier/vault" + + "github.com/vultisig/plugin/api" + "github.com/vultisig/plugin/storage" + "github.com/vultisig/plugin/storage/postgres" +) + +func main() { + //ctx := context.Background() + + cfg, err := GetConfigure() + if err != nil { + panic(err) + } + logger := logrus.New() + + sdClient, err := statsd.New(net.JoinHostPort(cfg.Datadog.Host, cfg.Datadog.Port)) + if err != nil { + panic(err) + } + redisStorage, err := storage.NewRedisStorage(cfg.Redis) + if err != nil { + panic(err) + } + redisOptions := asynq.RedisClientOpt{ + Addr: net.JoinHostPort(cfg.Redis.Host, cfg.Redis.Port), + Username: cfg.Redis.User, + Password: cfg.Redis.Password, + DB: cfg.Redis.DB, + } + + client := asynq.NewClient(redisOptions) + defer func() { + if err := client.Close(); err != nil { + fmt.Println("fail to close asynq client,", err) + } + }() + + inspector := asynq.NewInspector(redisOptions) + + vaultStorage, err := vault.NewBlockStorageImp(cfg.BlockStorage) + if err != nil { + panic(err) + } + + db, err := postgres.NewPostgresBackend(cfg.Database.DSN, nil) + if err != nil { + logger.Fatalf("Failed to connect to database: %v", err) + } + + //txIndexerStore, err := tx_indexer_storage.NewPostgresTxIndexStore(ctx, cfg.Database.DSN) + //if err != nil { + // panic(fmt.Errorf("tx_indexer_storage.NewPostgresTxIndexStore: %w", err)) + //} + + //txIndexerService := tx_indexer.NewService( + // logger, + // txIndexerStore, + // tx_indexer.Chains(), + //) + + //p, err := payroll.NewPlugin( + // db, + // nil, // not used by server + // vaultStorage, + // nil, + // txIndexerService, + // client, + // cfg.Server.EncryptionSecret, + //) + //if err != nil { + // logger.Fatalf("failed to create payroll plugin,err: %s", err) + //} + + server := api.NewServer( + cfg.Server, + db, + redisStorage, + vaultStorage, + redisOptions, + client, + inspector, + sdClient, + nil, + //p, + ) + if err := server.StartServer(); err != nil { + panic(err) + } +} diff --git a/cmd/copytrader/worker/config.go b/cmd/copytrader/worker/config.go new file mode 100644 index 0000000..7f6a6b7 --- /dev/null +++ b/cmd/copytrader/worker/config.go @@ -0,0 +1,66 @@ +package main + +import ( + "fmt" + "os" + "strings" + + "github.com/spf13/viper" + "github.com/vultisig/verifier/vault_config" + + "github.com/vultisig/plugin/storage" +) + +type CopytraderWorkerConfig struct { + Redis storage.RedisConfig `mapstructure:"redis" json:"redis,omitempty"` + Rpc Rpc `mapstructure:"Rpc" json:"Rpc,omitempty"` + Verifier verifier `mapstructure:"verifier" json:"verifier,omitempty"` + BlockStorage vault_config.BlockStorage `mapstructure:"block_storage" json:"block_storage,omitempty"` + VaultServiceConfig vault_config.Config `mapstructure:"vault_service" json:"vault_service,omitempty"` + Datadog struct { + Host string `mapstructure:"host" json:"host,omitempty"` + Port string `mapstructure:"port" json:"port,omitempty"` + } `mapstructure:"datadog" json:"datadog"` + Database struct { + DSN string `mapstructure:"dsn" json:"dsn,omitempty"` + } `mapstructure:"database" json:"database,omitempty"` +} + +type Rpc struct { + Ethereum rpcItem `mapstructure:"ethereum" json:"ethereum,omitempty"` +} + +type rpcItem struct { + URL string `mapstructure:"url" json:"url,omitempty"` +} + +type verifier struct { + URL string `mapstructure:"url"` + Token string `mapstructure:"token"` + PartyPrefix string `mapstructure:"party_prefix"` +} + +func GetConfigure() (*CopytraderWorkerConfig, error) { + configName := os.Getenv("VS_CONFIG_NAME") + if configName == "" { + configName = "config" + } + return ReadConfig(configName) +} + +func ReadConfig(configName string) (*CopytraderWorkerConfig, error) { + viper.SetConfigName(configName) + viper.AddConfigPath(".") + viper.SetEnvKeyReplacer(strings.NewReplacer(".", "_")) + viper.AutomaticEnv() + + if err := viper.ReadInConfig(); err != nil { + return nil, fmt.Errorf("fail to reading config file, %w", err) + } + var cfg CopytraderWorkerConfig + err := viper.Unmarshal(&cfg) + if err != nil { + return nil, fmt.Errorf("unable to decode into struct, %w", err) + } + return &cfg, nil +} diff --git a/cmd/copytrader/worker/main.go b/cmd/copytrader/worker/main.go new file mode 100644 index 0000000..1099161 --- /dev/null +++ b/cmd/copytrader/worker/main.go @@ -0,0 +1,128 @@ +package main + +import ( + "context" + "fmt" + + "github.com/DataDog/datadog-go/statsd" + "github.com/ethereum/go-ethereum/ethclient" + "github.com/hibiken/asynq" + "github.com/sirupsen/logrus" + "github.com/vultisig/plugin/internal/keysign" + "github.com/vultisig/plugin/plugin/copytrader" + "github.com/vultisig/plugin/storage/postgres" + "github.com/vultisig/verifier/tx_indexer" + "github.com/vultisig/verifier/tx_indexer/pkg/storage" + "github.com/vultisig/verifier/vault" + "github.com/vultisig/vultiserver/relay" + + "github.com/vultisig/plugin/internal/tasks" +) + +// Don't scale payroll.worker, it has scheduler which must be single instance running +// Consider moving scheduler to separate worker +func main() { + ctx := context.Background() + + cfg, err := GetConfigure() + if err != nil { + panic(err) + } + + sdClient, err := statsd.New(cfg.Datadog.Host + ":" + cfg.Datadog.Port) + if err != nil { + panic(err) + } + vaultStorage, err := vault.NewBlockStorageImp(cfg.BlockStorage) + if err != nil { + panic(fmt.Sprintf("failed to initialize vault storage: %v", err)) + } + + redisOptions := asynq.RedisClientOpt{ + Addr: cfg.Redis.Host + ":" + cfg.Redis.Port, + Username: cfg.Redis.User, + Password: cfg.Redis.Password, + DB: cfg.Redis.DB, + } + + logger := logrus.StandardLogger() + client := asynq.NewClient(redisOptions) + srv := asynq.NewServer( + redisOptions, + asynq.Config{ + Logger: logger, + Concurrency: 10, + Queues: map[string]int{ + tasks.QUEUE_NAME: 10, + "scheduled_plugin_queue": 10, // new queue + }, + }, + ) + + txIndexerStore, err := storage.NewPostgresTxIndexStore(ctx, cfg.Database.DSN) + if err != nil { + panic(fmt.Errorf("storage.NewPostgresTxIndexStore: %w", err)) + } + + txIndexerService := tx_indexer.NewService( + logger, + txIndexerStore, + tx_indexer.Chains(), + ) + + vaultService, err := vault.NewManagementService( + cfg.VaultServiceConfig, + client, + sdClient, + vaultStorage, + txIndexerService, + ) + if err != nil { + panic(fmt.Errorf("failed to create vault service: %w", err)) + } + + postgressDB, err := postgres.NewPostgresBackend(cfg.Database.DSN, nil) + if err != nil { + panic(fmt.Errorf("failed to create postgres backend: %w", err)) + } + + rpcClient, err := ethclient.Dial(cfg.Rpc.Ethereum.URL) + if err != nil { + panic(fmt.Errorf("failed to create eth client: %w", err)) + } + + ct, err := copytrader.NewPlugin( + postgressDB, + keysign.NewSigner( + logger.WithField("pkg", "keysign.Signer").Logger, + relay.NewRelayClient(cfg.VaultServiceConfig.Relay.Server), + []keysign.Emitter{ + keysign.NewVerifierEmitter(cfg.Verifier.URL, cfg.Verifier.Token), + keysign.NewPluginEmitter(client, tasks.TypeKeySignDKLS, tasks.QUEUE_NAME), + }, + []string{ + cfg.Verifier.PartyPrefix, + cfg.VaultServiceConfig.LocalPartyPrefix, + }, + ), + vaultStorage, + rpcClient, + txIndexerService, + client, + cfg.VaultServiceConfig.EncryptionSecret, + ) + if err != nil { + panic(fmt.Errorf("failed to create copytrader plugin: %w", err)) + } + + _ = ct + + mux := asynq.NewServeMux() + // TRIGGER??? + //mux.HandleFunc(tasks.TypePluginTransaction, p.HandleSchedulerTrigger) + mux.HandleFunc(tasks.TypeKeySignDKLS, vaultService.HandleKeySignDKLS) + mux.HandleFunc(tasks.TypeReshareDKLS, vaultService.HandleReshareDKLS) + if err := srv.Run(mux); err != nil { + panic(fmt.Errorf("could not run server: %w", err)) + } +} diff --git a/go.mod b/go.mod index bcf2c5c..2377402 100644 --- a/go.mod +++ b/go.mod @@ -196,4 +196,5 @@ replace ( github.com/cwespare/xxhash/v2 => github.com/cespare/xxhash/v2 v2.1.1 github.com/gogo/protobuf => github.com/gogo/protobuf v1.3.2 nhooyr.io/websocket => github.com/coder/websocket v1.8.6 + github.com/vultisig/verifier => ../verifier ) diff --git a/plugin/copytrader/copytrader.go b/plugin/copytrader/copytrader.go new file mode 100644 index 0000000..e6384ba --- /dev/null +++ b/plugin/copytrader/copytrader.go @@ -0,0 +1,193 @@ +package copytrader + +import ( + "context" + "fmt" + + gcommon "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/ethclient" + "github.com/hibiken/asynq" + "github.com/sirupsen/logrus" + "github.com/vultisig/mobile-tss-lib/tss" + "github.com/vultisig/recipes/sdk/evm" + rtypes "github.com/vultisig/recipes/types" + "github.com/vultisig/verifier/common" + vcommon "github.com/vultisig/verifier/common" + "github.com/vultisig/verifier/plugin" + "github.com/vultisig/verifier/tx_indexer" + vtypes "github.com/vultisig/verifier/types" + "github.com/vultisig/verifier/vault" + + "github.com/vultisig/plugin/internal/keysign" + "github.com/vultisig/plugin/storage" +) + +var _ plugin.Plugin = (*Plugin)(nil) + +type Plugin struct { + db storage.DatabaseStorage + signer *keysign.Signer + eth *evm.SDK + logger logrus.FieldLogger + txIndexerService *tx_indexer.Service + client *asynq.Client + vaultStorage vault.Storage + vaultEncryptionSecret string +} + +func NewPlugin( + db storage.DatabaseStorage, + signer *keysign.Signer, + vaultStorage vault.Storage, + ethRpc *ethclient.Client, + txIndexerService *tx_indexer.Service, + client *asynq.Client, + vaultEncryptionSecret string, +) (*Plugin, error) { + if db == nil { + return nil, fmt.Errorf("database storage cannot be nil") + } + + var eth *evm.SDK + if ethRpc != nil { + ethEvmChainID, err := common.Ethereum.EvmID() + if err != nil { + return nil, fmt.Errorf("failed to get Ethereum EVM ID: %w", err) + } + eth = evm.NewSDK(ethEvmChainID, ethRpc, ethRpc.Client()) + } + + return &Plugin{ + db: db, + signer: signer, + eth: eth, + logger: logrus.WithField("plugin", "payroll"), + txIndexerService: txIndexerService, + client: client, + vaultStorage: vaultStorage, + vaultEncryptionSecret: vaultEncryptionSecret, + }, nil +} + +func (p *Plugin) GetRecipeSpecification() *rtypes.RecipeSchema { + return &rtypes.RecipeSchema{ + Version: 1, // Schema version + ScheduleVersion: 1, // Schedule specification version + // TODO: configure + PluginId: string(vtypes.PluginVultisigCopytrader_0000), + PluginName: "Copy trading plugin", + PluginVersion: 1, // Convert from "0.1.0" to int32 + SupportedResources: []*rtypes.ResourcePattern{ + { + ResourcePath: &rtypes.ResourcePath{ + ChainId: "ethereum", + ProtocolId: "uniswapv2_router", + FunctionId: "swapExactTokensForTokens", + Full: "ethereum.uniswapv2_router.swapExactTokensForTokens", + }, + ParameterCapabilities: []*rtypes.ParameterConstraintCapability{ + { + ParameterName: "aim", + SupportedTypes: []rtypes.ConstraintType{ + rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, + rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, + }, + Required: true, + }, + { + ParameterName: "source_token", + SupportedTypes: []rtypes.ConstraintType{ + rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, + rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, + }, + Required: true, + }, + { + ParameterName: "destination_token", + SupportedTypes: []rtypes.ConstraintType{ + rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, + rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, + }, + Required: true, + }, + { + ParameterName: "amount", + SupportedTypes: []rtypes.ConstraintType{ + rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, + rtypes.ConstraintType_CONSTRAINT_TYPE_MAX, + rtypes.ConstraintType_CONSTRAINT_TYPE_RANGE, + }, + Required: true, + }, + }, + Required: true, + }, + }, + Requirements: &rtypes.PluginRequirements{ + MinVultisigVersion: 1, + SupportedChains: []string{"ethereum"}, + }, + } +} + +func (p *Plugin) ProposeTransactions(policy vtypes.PluginPolicy) ([]vtypes.PluginKeysignRequest, error) { + //TODO implement me + panic("implement me") +} + +func (p *Plugin) initSign( + ctx context.Context, + req vtypes.PluginKeysignRequest, + pluginPolicy vtypes.PluginPolicy, +) error { + sigs, err := p.signer.Sign(ctx, req) + if err != nil { + p.logger.WithError(err).Error("Keysign failed") + return fmt.Errorf("failed to sign transaction: %w", err) + } + + if len(sigs) != 1 { + p.logger. + WithField("sigs_count", len(sigs)). + Error("expected only 1 message+sig per request for evm") + return fmt.Errorf("failed to sign transaction: invalid signature count: %d", len(sigs)) + } + var sig tss.KeysignResponse + for _, s := range sigs { + sig = s + } + + err = p.SigningComplete(ctx, sig, req, pluginPolicy) + if err != nil { + p.logger.WithError(err).Error("failed to complete signing process (broadcast tx)") + return fmt.Errorf("failed to complete signing process: %w", err) + } + return nil +} + +func (p *Plugin) SigningComplete( + ctx context.Context, + signature tss.KeysignResponse, + signRequest vtypes.PluginKeysignRequest, + _ vtypes.PluginPolicy, +) error { + tx, err := p.eth.Send( + ctx, + gcommon.FromHex(signRequest.Transaction), + gcommon.Hex2Bytes(signature.R), + gcommon.Hex2Bytes(signature.S), + gcommon.Hex2Bytes(signature.RecoveryID), + ) + if err != nil { + p.logger.WithError(err).WithField("tx_hex", signRequest.Transaction).Error("p.eth.Send") + return fmt.Errorf("p.eth.Send(tx_hex=%s): %w", signRequest.Transaction, err) + } + + p.logger.WithFields(logrus.Fields{ + "from_public_key": signRequest.PublicKey, + "to_address": tx.To().Hex(), + "hash": tx.Hash().Hex(), + "chain": vcommon.Ethereum.String(), + }).Info("tx successfully signed and broadcasted") + return nil +} diff --git a/plugin/copytrader/policy.go b/plugin/copytrader/policy.go new file mode 100644 index 0000000..6b7cf05 --- /dev/null +++ b/plugin/copytrader/policy.go @@ -0,0 +1,54 @@ +package copytrader + +import ( + "fmt" + "strings" + + "github.com/vultisig/plugin/internal/plugin" + "github.com/vultisig/recipes/chain" + "github.com/vultisig/recipes/engine" + vtypes "github.com/vultisig/verifier/types" +) + +func (p *Plugin) ValidateProposedTransactions(policy vtypes.PluginPolicy, txs []vtypes.PluginKeysignRequest) error { + err := p.ValidatePluginPolicy(policy) + if err != nil { + return fmt.Errorf("failed to validate plugin policy: %v", err) + } + + recipe, err := policy.GetRecipe() + if err != nil { + return fmt.Errorf("failed to get recipe from policy: %v", err) + } + + eng := engine.NewEngine() + + for _, tx := range txs { + for _, keysignMessage := range tx.Messages { + messageChain, err := chain.GetChain(strings.ToLower(keysignMessage.Chain.String())) + if err != nil { + return fmt.Errorf("failed to get chain: %w", err) + } + + decodedTx, err := messageChain.ParseTransaction(keysignMessage.Message) + if err != nil { + return fmt.Errorf("failed to parse transaction: %w", err) + } + + transactionAllowed, _, err := eng.Evaluate(recipe, messageChain, decodedTx) + if err != nil { + return fmt.Errorf("failed to evaluate transaction: %w", err) + } + + if !transactionAllowed { + return fmt.Errorf("transaction %s on %s not allowed by policy", keysignMessage.Hash, keysignMessage.Chain) + } + } + } + + return nil +} + +func (p *Plugin) ValidatePluginPolicy(policyDoc vtypes.PluginPolicy) error { + return plugin.ValidatePluginPolicy(policyDoc, p.GetRecipeSpecification()) +} diff --git a/plugin/payroll/payroll.go b/plugin/payroll/payroll.go index d04cfaf..ba3f894 100644 --- a/plugin/payroll/payroll.go +++ b/plugin/payroll/payroll.go @@ -6,13 +6,14 @@ import ( "github.com/ethereum/go-ethereum/ethclient" "github.com/hibiken/asynq" "github.com/sirupsen/logrus" - "github.com/vultisig/plugin/internal/keysign" - "github.com/vultisig/plugin/storage" "github.com/vultisig/recipes/sdk/evm" "github.com/vultisig/verifier/common" "github.com/vultisig/verifier/plugin" "github.com/vultisig/verifier/tx_indexer" "github.com/vultisig/verifier/vault" + + "github.com/vultisig/plugin/internal/keysign" + "github.com/vultisig/plugin/storage" ) var _ plugin.Plugin = (*Plugin)(nil) diff --git a/plugin/payroll/transaction.go b/plugin/payroll/transaction.go index cb3e7e3..58533d8 100644 --- a/plugin/payroll/transaction.go +++ b/plugin/payroll/transaction.go @@ -11,25 +11,25 @@ import ( "sync" "time" + gcommon "github.com/ethereum/go-ethereum/common" etypes "github.com/ethereum/go-ethereum/core/types" "github.com/google/uuid" - "github.com/vultisig/recipes/ethereum" - "github.com/vultisig/recipes/sdk/evm" - "github.com/vultisig/verifier/tx_indexer/pkg/storage" - "github.com/vultisig/vultiserver/contexthelper" - "golang.org/x/sync/errgroup" - - gcommon "github.com/ethereum/go-ethereum/common" "github.com/hibiken/asynq" "github.com/sirupsen/logrus" + "golang.org/x/sync/errgroup" + "github.com/vultisig/mobile-tss-lib/tss" - "github.com/vultisig/plugin/common" - "github.com/vultisig/plugin/internal/scheduler" + "github.com/vultisig/recipes/ethereum" + "github.com/vultisig/recipes/sdk/evm" rtypes "github.com/vultisig/recipes/types" "github.com/vultisig/verifier/address" vcommon "github.com/vultisig/verifier/common" + "github.com/vultisig/verifier/tx_indexer/pkg/storage" vtypes "github.com/vultisig/verifier/types" + "github.com/vultisig/vultiserver/contexthelper" + "github.com/vultisig/plugin/common" + "github.com/vultisig/plugin/internal/scheduler" "github.com/vultisig/plugin/internal/types" ) From ae270cd0e29b146072d3fb63153a831154003a61 Mon Sep 17 00:00:00 2001 From: Shevchuk Date: Thu, 24 Jul 2025 22:59:05 +0700 Subject: [PATCH 2/3] [WIP] Swap watcher added --- plugin/copytrader/copytrader.go | 20 +++++++++-- plugin/copytrader/watcher.go | 64 +++++++++++++++++++++++++++++++++ 2 files changed, 82 insertions(+), 2 deletions(-) create mode 100644 plugin/copytrader/watcher.go diff --git a/plugin/copytrader/copytrader.go b/plugin/copytrader/copytrader.go index e6384ba..22f0695 100644 --- a/plugin/copytrader/copytrader.go +++ b/plugin/copytrader/copytrader.go @@ -22,17 +22,23 @@ import ( "github.com/vultisig/plugin/storage" ) +const UniswapV2RouterAddress = "0x7a250d5630B4cF539739dF2C5dAcb4c659F2488D" + +var UniswapSwapTopic = gcommon.HexToHash("0xd78ad95fa46c994b6551d0da85fc275fe613ce37657fb8d5e3d130840159d822") + var _ plugin.Plugin = (*Plugin)(nil) type Plugin struct { db storage.DatabaseStorage signer *keysign.Signer eth *evm.SDK + ethRpc *ethclient.Client logger logrus.FieldLogger txIndexerService *tx_indexer.Service client *asynq.Client vaultStorage vault.Storage vaultEncryptionSecret string + blockID uint64 } func NewPlugin( @@ -48,24 +54,34 @@ func NewPlugin( return nil, fmt.Errorf("database storage cannot be nil") } - var eth *evm.SDK + var ( + eth *evm.SDK + currentBlock uint64 + ) if ethRpc != nil { ethEvmChainID, err := common.Ethereum.EvmID() if err != nil { return nil, fmt.Errorf("failed to get Ethereum EVM ID: %w", err) } eth = evm.NewSDK(ethEvmChainID, ethRpc, ethRpc.Client()) + + currentBlock, err = ethRpc.BlockNumber(context.Background()) + if err != nil { + return nil, fmt.Errorf("failed to get block: %w", err) + } } return &Plugin{ db: db, signer: signer, eth: eth, - logger: logrus.WithField("plugin", "payroll"), + ethRpc: ethRpc, + logger: logrus.WithField("plugin", "copytrader"), txIndexerService: txIndexerService, client: client, vaultStorage: vaultStorage, vaultEncryptionSecret: vaultEncryptionSecret, + blockID: currentBlock, }, nil } diff --git a/plugin/copytrader/watcher.go b/plugin/copytrader/watcher.go new file mode 100644 index 0000000..f1e1fc0 --- /dev/null +++ b/plugin/copytrader/watcher.go @@ -0,0 +1,64 @@ +package copytrader + +import ( + "context" + "math/big" + "time" + + "github.com/ethereum/go-ethereum/core/types" + "github.com/sirupsen/logrus" +) + +func (p *Plugin) WatchSwap(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + case <-time.After(5 * time.Second): + currentBlock, err := p.ethRpc.BlockNumber(ctx) + if err != nil { + p.logger.WithError(err).Error("failed to get block") + continue + } + + for p.blockID < currentBlock { + p.blockID++ + p.logger.Info("Processing block: ", p.blockID) + + block, err := p.ethRpc.BlockByNumber(ctx, big.NewInt(0).SetUint64(p.blockID)) + if err != nil { + p.logger.WithError(err).Error("failed to get block") + continue + } + + for _, tx := range block.Transactions() { + // is Uniswap tx check + if tx.To().String() == UniswapV2RouterAddress { + txReceipt, err := p.ethRpc.TransactionReceipt(ctx, tx.Hash()) + if err != nil { + p.logger.WithError(err).Error("failed to get block") + continue + } + + for _, log := range txReceipt.Logs { + if log.Topics[0] == UniswapSwapTopic { + signer := types.LatestSignerForChainID(tx.ChainId()) + sender, err := signer.Sender(tx) + if err != nil { + p.logger.WithError(err).Error("failed to get signer") + continue + } + p.logger.WithFields(logrus.Fields{ + "sender": sender.String(), + "txHash": tx.Hash().String(), + "pair": log.Address.String(), + }) + //TODO: Trigger swaps there + } + } + } + } + } + } + } +} From 7f531bb819df0ac2f99dc3197295cd2e6fa28041 Mon Sep 17 00:00:00 2001 From: Shevchuk Date: Fri, 25 Jul 2025 19:35:37 +0700 Subject: [PATCH 3/3] Added swap trigger and handler --- cmd/copytrader/server/main.go | 51 ++++++------ cmd/copytrader/worker/main.go | 8 +- plugin/copytrader/const.go | 6 ++ plugin/copytrader/copytrader.go | 132 ------------------------------- plugin/copytrader/policy.go | 65 ++++++++++++++- plugin/copytrader/transaction.go | 97 +++++++++++++++++++++++ plugin/copytrader/types.go | 13 +++ plugin/copytrader/watcher.go | 69 +++++++++++----- 8 files changed, 259 insertions(+), 182 deletions(-) create mode 100644 plugin/copytrader/const.go create mode 100644 plugin/copytrader/transaction.go create mode 100644 plugin/copytrader/types.go diff --git a/cmd/copytrader/server/main.go b/cmd/copytrader/server/main.go index 73faeff..ebd56fe 100644 --- a/cmd/copytrader/server/main.go +++ b/cmd/copytrader/server/main.go @@ -1,12 +1,16 @@ package main import ( + "context" "fmt" "net" "github.com/DataDog/datadog-go/statsd" "github.com/hibiken/asynq" "github.com/sirupsen/logrus" + "github.com/vultisig/plugin/plugin/copytrader" + "github.com/vultisig/verifier/tx_indexer" + tx_indexer_storage "github.com/vultisig/verifier/tx_indexer/pkg/storage" "github.com/vultisig/verifier/vault" "github.com/vultisig/plugin/api" @@ -15,7 +19,7 @@ import ( ) func main() { - //ctx := context.Background() + ctx := context.Background() cfg, err := GetConfigure() if err != nil { @@ -57,29 +61,29 @@ func main() { logger.Fatalf("Failed to connect to database: %v", err) } - //txIndexerStore, err := tx_indexer_storage.NewPostgresTxIndexStore(ctx, cfg.Database.DSN) - //if err != nil { - // panic(fmt.Errorf("tx_indexer_storage.NewPostgresTxIndexStore: %w", err)) - //} + txIndexerStore, err := tx_indexer_storage.NewPostgresTxIndexStore(ctx, cfg.Database.DSN) + if err != nil { + panic(fmt.Errorf("tx_indexer_storage.NewPostgresTxIndexStore: %w", err)) + } - //txIndexerService := tx_indexer.NewService( - // logger, - // txIndexerStore, - // tx_indexer.Chains(), - //) + txIndexerService := tx_indexer.NewService( + logger, + txIndexerStore, + tx_indexer.Chains(), + ) - //p, err := payroll.NewPlugin( - // db, - // nil, // not used by server - // vaultStorage, - // nil, - // txIndexerService, - // client, - // cfg.Server.EncryptionSecret, - //) - //if err != nil { - // logger.Fatalf("failed to create payroll plugin,err: %s", err) - //} + ct, err := copytrader.NewPlugin( + db, + nil, // not used by server + vaultStorage, + nil, + txIndexerService, + client, + cfg.Server.EncryptionSecret, + ) + if err != nil { + logger.Fatalf("failed to create payroll plugin,err: %s", err) + } server := api.NewServer( cfg.Server, @@ -90,8 +94,7 @@ func main() { client, inspector, sdClient, - nil, - //p, + ct, ) if err := server.StartServer(); err != nil { panic(err) diff --git a/cmd/copytrader/worker/main.go b/cmd/copytrader/worker/main.go index 1099161..00482a0 100644 --- a/cmd/copytrader/worker/main.go +++ b/cmd/copytrader/worker/main.go @@ -53,8 +53,7 @@ func main() { Logger: logger, Concurrency: 10, Queues: map[string]int{ - tasks.QUEUE_NAME: 10, - "scheduled_plugin_queue": 10, // new queue + tasks.QUEUE_NAME: 10, }, }, ) @@ -115,11 +114,8 @@ func main() { panic(fmt.Errorf("failed to create copytrader plugin: %w", err)) } - _ = ct - mux := asynq.NewServeMux() - // TRIGGER??? - //mux.HandleFunc(tasks.TypePluginTransaction, p.HandleSchedulerTrigger) + mux.HandleFunc(tasks.TypePluginTransaction, ct.HandleSwapTask) mux.HandleFunc(tasks.TypeKeySignDKLS, vaultService.HandleKeySignDKLS) mux.HandleFunc(tasks.TypeReshareDKLS, vaultService.HandleReshareDKLS) if err := srv.Run(mux); err != nil { diff --git a/plugin/copytrader/const.go b/plugin/copytrader/const.go new file mode 100644 index 0000000..093aa86 --- /dev/null +++ b/plugin/copytrader/const.go @@ -0,0 +1,6 @@ +package copytrader + +const ( + UniswapV2RouterAddress = "0x7a250d5630B4cF539739dF2C5dAcb4c659F2488D" + SwapExactTokensForTokens = "38ed1739" +) diff --git a/plugin/copytrader/copytrader.go b/plugin/copytrader/copytrader.go index 22f0695..c767cc2 100644 --- a/plugin/copytrader/copytrader.go +++ b/plugin/copytrader/copytrader.go @@ -4,28 +4,19 @@ import ( "context" "fmt" - gcommon "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/ethclient" "github.com/hibiken/asynq" "github.com/sirupsen/logrus" - "github.com/vultisig/mobile-tss-lib/tss" "github.com/vultisig/recipes/sdk/evm" - rtypes "github.com/vultisig/recipes/types" "github.com/vultisig/verifier/common" - vcommon "github.com/vultisig/verifier/common" "github.com/vultisig/verifier/plugin" "github.com/vultisig/verifier/tx_indexer" - vtypes "github.com/vultisig/verifier/types" "github.com/vultisig/verifier/vault" "github.com/vultisig/plugin/internal/keysign" "github.com/vultisig/plugin/storage" ) -const UniswapV2RouterAddress = "0x7a250d5630B4cF539739dF2C5dAcb4c659F2488D" - -var UniswapSwapTopic = gcommon.HexToHash("0xd78ad95fa46c994b6551d0da85fc275fe613ce37657fb8d5e3d130840159d822") - var _ plugin.Plugin = (*Plugin)(nil) type Plugin struct { @@ -84,126 +75,3 @@ func NewPlugin( blockID: currentBlock, }, nil } - -func (p *Plugin) GetRecipeSpecification() *rtypes.RecipeSchema { - return &rtypes.RecipeSchema{ - Version: 1, // Schema version - ScheduleVersion: 1, // Schedule specification version - // TODO: configure - PluginId: string(vtypes.PluginVultisigCopytrader_0000), - PluginName: "Copy trading plugin", - PluginVersion: 1, // Convert from "0.1.0" to int32 - SupportedResources: []*rtypes.ResourcePattern{ - { - ResourcePath: &rtypes.ResourcePath{ - ChainId: "ethereum", - ProtocolId: "uniswapv2_router", - FunctionId: "swapExactTokensForTokens", - Full: "ethereum.uniswapv2_router.swapExactTokensForTokens", - }, - ParameterCapabilities: []*rtypes.ParameterConstraintCapability{ - { - ParameterName: "aim", - SupportedTypes: []rtypes.ConstraintType{ - rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, - rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, - }, - Required: true, - }, - { - ParameterName: "source_token", - SupportedTypes: []rtypes.ConstraintType{ - rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, - rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, - }, - Required: true, - }, - { - ParameterName: "destination_token", - SupportedTypes: []rtypes.ConstraintType{ - rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, - rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, - }, - Required: true, - }, - { - ParameterName: "amount", - SupportedTypes: []rtypes.ConstraintType{ - rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, - rtypes.ConstraintType_CONSTRAINT_TYPE_MAX, - rtypes.ConstraintType_CONSTRAINT_TYPE_RANGE, - }, - Required: true, - }, - }, - Required: true, - }, - }, - Requirements: &rtypes.PluginRequirements{ - MinVultisigVersion: 1, - SupportedChains: []string{"ethereum"}, - }, - } -} - -func (p *Plugin) ProposeTransactions(policy vtypes.PluginPolicy) ([]vtypes.PluginKeysignRequest, error) { - //TODO implement me - panic("implement me") -} - -func (p *Plugin) initSign( - ctx context.Context, - req vtypes.PluginKeysignRequest, - pluginPolicy vtypes.PluginPolicy, -) error { - sigs, err := p.signer.Sign(ctx, req) - if err != nil { - p.logger.WithError(err).Error("Keysign failed") - return fmt.Errorf("failed to sign transaction: %w", err) - } - - if len(sigs) != 1 { - p.logger. - WithField("sigs_count", len(sigs)). - Error("expected only 1 message+sig per request for evm") - return fmt.Errorf("failed to sign transaction: invalid signature count: %d", len(sigs)) - } - var sig tss.KeysignResponse - for _, s := range sigs { - sig = s - } - - err = p.SigningComplete(ctx, sig, req, pluginPolicy) - if err != nil { - p.logger.WithError(err).Error("failed to complete signing process (broadcast tx)") - return fmt.Errorf("failed to complete signing process: %w", err) - } - return nil -} - -func (p *Plugin) SigningComplete( - ctx context.Context, - signature tss.KeysignResponse, - signRequest vtypes.PluginKeysignRequest, - _ vtypes.PluginPolicy, -) error { - tx, err := p.eth.Send( - ctx, - gcommon.FromHex(signRequest.Transaction), - gcommon.Hex2Bytes(signature.R), - gcommon.Hex2Bytes(signature.S), - gcommon.Hex2Bytes(signature.RecoveryID), - ) - if err != nil { - p.logger.WithError(err).WithField("tx_hex", signRequest.Transaction).Error("p.eth.Send") - return fmt.Errorf("p.eth.Send(tx_hex=%s): %w", signRequest.Transaction, err) - } - - p.logger.WithFields(logrus.Fields{ - "from_public_key": signRequest.PublicKey, - "to_address": tx.To().Hex(), - "hash": tx.Hash().Hex(), - "chain": vcommon.Ethereum.String(), - }).Info("tx successfully signed and broadcasted") - return nil -} diff --git a/plugin/copytrader/policy.go b/plugin/copytrader/policy.go index 6b7cf05..3158da7 100644 --- a/plugin/copytrader/policy.go +++ b/plugin/copytrader/policy.go @@ -4,10 +4,12 @@ import ( "fmt" "strings" - "github.com/vultisig/plugin/internal/plugin" "github.com/vultisig/recipes/chain" "github.com/vultisig/recipes/engine" + rtypes "github.com/vultisig/recipes/types" vtypes "github.com/vultisig/verifier/types" + + "github.com/vultisig/plugin/internal/plugin" ) func (p *Plugin) ValidateProposedTransactions(policy vtypes.PluginPolicy, txs []vtypes.PluginKeysignRequest) error { @@ -52,3 +54,64 @@ func (p *Plugin) ValidateProposedTransactions(policy vtypes.PluginPolicy, txs [] func (p *Plugin) ValidatePluginPolicy(policyDoc vtypes.PluginPolicy) error { return plugin.ValidatePluginPolicy(policyDoc, p.GetRecipeSpecification()) } + +func (p *Plugin) GetRecipeSpecification() *rtypes.RecipeSchema { + return &rtypes.RecipeSchema{ + Version: 1, // Schema version + ScheduleVersion: 1, // Schedule specification version + // TODO: configure + PluginId: string(vtypes.PluginVultisigCopytrader_0000), + PluginName: "Copy trading plugin", + PluginVersion: 1, // Convert from "0.1.0" to int32 + SupportedResources: []*rtypes.ResourcePattern{ + { + ResourcePath: &rtypes.ResourcePath{ + ChainId: "ethereum", + ProtocolId: "uniswapv2_router", + FunctionId: "swapExactTokensForTokens", + Full: "ethereum.uniswapv2_router.swapExactTokensForTokens", + }, + ParameterCapabilities: []*rtypes.ParameterConstraintCapability{ + { + ParameterName: "aim", + SupportedTypes: []rtypes.ConstraintType{ + rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, + rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, + }, + Required: true, + }, + { + ParameterName: "source_token", + SupportedTypes: []rtypes.ConstraintType{ + rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, + rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, + }, + Required: true, + }, + { + ParameterName: "destination_token", + SupportedTypes: []rtypes.ConstraintType{ + rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, + rtypes.ConstraintType_CONSTRAINT_TYPE_WHITELIST, + }, + Required: true, + }, + { + ParameterName: "amount", + SupportedTypes: []rtypes.ConstraintType{ + rtypes.ConstraintType_CONSTRAINT_TYPE_FIXED, + rtypes.ConstraintType_CONSTRAINT_TYPE_MAX, + rtypes.ConstraintType_CONSTRAINT_TYPE_RANGE, + }, + Required: true, + }, + }, + Required: true, + }, + }, + Requirements: &rtypes.PluginRequirements{ + MinVultisigVersion: 1, + SupportedChains: []string{"ethereum"}, + }, + } +} diff --git a/plugin/copytrader/transaction.go b/plugin/copytrader/transaction.go new file mode 100644 index 0000000..94c8802 --- /dev/null +++ b/plugin/copytrader/transaction.go @@ -0,0 +1,97 @@ +package copytrader + +import ( + "context" + "encoding/json" + "fmt" + "time" + + gcommon "github.com/ethereum/go-ethereum/common" + "github.com/hibiken/asynq" + "github.com/sirupsen/logrus" + "github.com/vultisig/mobile-tss-lib/tss" + vcommon "github.com/vultisig/verifier/common" + vtypes "github.com/vultisig/verifier/types" + "github.com/vultisig/vultiserver/contexthelper" +) + +func (p *Plugin) HandleSwapTask(c context.Context, t *asynq.Task) error { + ctx, cancel := context.WithTimeout(c, 5*time.Minute) + defer cancel() + + if err := contexthelper.CheckCancellation(ctx); err != nil { + p.logger.WithError(err).Warn("Context cancelled, skipping trigger") + return err + } + var swapTask SwapTask + if err := json.Unmarshal(t.Payload(), &swapTask); err != nil { + p.logger.WithError(err).Error("Failed to unmarshal swapTask payload") + return fmt.Errorf("failed to unmarshal swapTask payload: %s, %w", err, asynq.SkipRetry) + } + + //TODO: implement aim <-> policy db + //TODO: trigger swaps + return nil +} + +func (p *Plugin) ProposeTransactions(policy vtypes.PluginPolicy) ([]vtypes.PluginKeysignRequest, error) { + //TODO implement me + panic("implement me") +} + +func (p *Plugin) initSign( + ctx context.Context, + req vtypes.PluginKeysignRequest, + pluginPolicy vtypes.PluginPolicy, +) error { + sigs, err := p.signer.Sign(ctx, req) + if err != nil { + p.logger.WithError(err).Error("Keysign failed") + return fmt.Errorf("failed to sign transaction: %w", err) + } + + if len(sigs) != 1 { + p.logger. + WithField("sigs_count", len(sigs)). + Error("expected only 1 message+sig per request for evm") + return fmt.Errorf("failed to sign transaction: invalid signature count: %d", len(sigs)) + } + var sig tss.KeysignResponse + for _, s := range sigs { + sig = s + } + + err = p.SigningComplete(ctx, sig, req, pluginPolicy) + if err != nil { + p.logger.WithError(err).Error("failed to complete signing process (broadcast tx)") + return fmt.Errorf("failed to complete signing process: %w", err) + } + return nil +} + +func (p *Plugin) SigningComplete( + ctx context.Context, + signature tss.KeysignResponse, + signRequest vtypes.PluginKeysignRequest, + _ vtypes.PluginPolicy, +) error { + tx, err := p.eth.Send( + ctx, + gcommon.FromHex(signRequest.Transaction), + gcommon.Hex2Bytes(signature.R), + gcommon.Hex2Bytes(signature.S), + gcommon.Hex2Bytes(signature.RecoveryID), + ) + if err != nil { + p.logger.WithError(err).WithField("tx_hex", signRequest.Transaction).Error("p.eth.Send") + return fmt.Errorf("p.eth.Send(tx_hex=%s): %w", signRequest.Transaction, err) + } + + p.logger.WithFields(logrus.Fields{ + "from_public_key": signRequest.PublicKey, + "to_address": tx.To().Hex(), + "hash": tx.Hash().Hex(), + "chain": vcommon.Ethereum.String(), + }).Info("tx successfully signed and broadcasted") + return nil +} diff --git a/plugin/copytrader/types.go b/plugin/copytrader/types.go new file mode 100644 index 0000000..a97e5fd --- /dev/null +++ b/plugin/copytrader/types.go @@ -0,0 +1,13 @@ +package copytrader + +import ( + "math/big" + + common "github.com/ethereum/go-ethereum/common" +) + +type SwapTask struct { + Sender common.Address + Path []common.Address + Amount *big.Int +} diff --git a/plugin/copytrader/watcher.go b/plugin/copytrader/watcher.go index f1e1fc0..4c1dcd7 100644 --- a/plugin/copytrader/watcher.go +++ b/plugin/copytrader/watcher.go @@ -2,14 +2,24 @@ package copytrader import ( "context" + "encoding/hex" + "encoding/json" + "fmt" "math/big" "time" - "github.com/ethereum/go-ethereum/core/types" - "github.com/sirupsen/logrus" + "github.com/ethereum/go-ethereum/accounts/abi" + "github.com/ethereum/go-ethereum/common" + "github.com/vultisig/recipes/sdk/evm/codegen/uniswapv2_router" ) func (p *Plugin) WatchSwap(ctx context.Context) { + var uniswapABI abi.ABI + err := json.Unmarshal([]byte(uniswapv2_router.Uniswapv2RouterMetaData.ABI), &uniswapABI) + if err != nil { + panic(err) + } + for { select { case <-ctx.Done(): @@ -31,31 +41,52 @@ func (p *Plugin) WatchSwap(ctx context.Context) { continue } + // Process txs to find UniswapV2Router interactions for _, tx := range block.Transactions() { + if tx.To() == nil { + continue + } // is Uniswap tx check if tx.To().String() == UniswapV2RouterAddress { - txReceipt, err := p.ethRpc.TransactionReceipt(ctx, tx.Hash()) + inputBytes := tx.Data() + signature, data := inputBytes[:4], inputBytes[4:] + if hex.EncodeToString(signature) != SwapExactTokensForTokens { + continue + } + + method, err := uniswapABI.MethodById(signature) if err != nil { - p.logger.WithError(err).Error("failed to get block") + p.logger.WithError(err).Error("unknown method") continue } - for _, log := range txReceipt.Logs { - if log.Topics[0] == UniswapSwapTopic { - signer := types.LatestSignerForChainID(tx.ChainId()) - sender, err := signer.Sender(tx) - if err != nil { - p.logger.WithError(err).Error("failed to get signer") - continue - } - p.logger.WithFields(logrus.Fields{ - "sender": sender.String(), - "txHash": tx.Hash().String(), - "pair": log.Address.String(), - }) - //TODO: Trigger swaps there - } + // Getting args from tx to find necessary info + var args = make(map[string]interface{}) + err = method.Inputs.UnpackIntoMap(args, data) + if err != nil { + p.logger.WithError(err).Error("failed to unpack data") + continue } + + path := args["path"] + tokens, valid := path.([]common.Address) + if !valid { + p.logger.Error("invalid path", path) + continue + } + + amountIn, _ := new(big.Int).SetString(fmt.Sprint(args["amountIn"]), 10) + to := args["to"] + sender, valid := to.(common.Address) + if !valid { + p.logger.Error("invalid sender", to) + continue + } + + //Triggering swaps + fmt.Println("sender", sender.String()) + fmt.Println("amount", amountIn.String()) + fmt.Println("path: ", tokens) } } }