init
This commit is contained in:
parent
df2eade8c8
commit
fb07e0b010
82
main.go
Normal file
82
main.go
Normal file
@ -0,0 +1,82 @@
|
||||
//go:build !Develop
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"AngkorWalletScanning/model"
|
||||
"AngkorWalletScanning/src"
|
||||
"AngkorWalletScanning/umlog"
|
||||
"AngkorWalletScanning/umsql"
|
||||
"fmt"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
var (
|
||||
MODE string = "" // Release, Info, Debug
|
||||
DB_USER string = ""
|
||||
DB_PASS string = ""
|
||||
DB_HOST string = ""
|
||||
DB_PORT string = ""
|
||||
DB_NAME string = ""
|
||||
KEY string = ""
|
||||
CODE_STRC model.CodeStrc
|
||||
)
|
||||
|
||||
func main() {
|
||||
var header model.Header
|
||||
if len(os.Args) != 1 {
|
||||
DB_USER = os.Args[1]
|
||||
DB_PASS = os.Args[2]
|
||||
DB_HOST = os.Args[3]
|
||||
DB_PORT = os.Args[4]
|
||||
DB_NAME = os.Args[5]
|
||||
KEY = os.Args[6]
|
||||
MODE = os.Args[7]
|
||||
header.AppKey = os.Args[8]
|
||||
header.SecretKey = os.Args[9]
|
||||
}
|
||||
|
||||
err := umlog.Init("/var/tmp/AnkWalletScanning.log", "7208765670:AAEN35cnQPeoDho33QvrchhVzuwnUcyDZdk", CODE_STRC.MESSAGE_ID)
|
||||
if err != nil {
|
||||
fmt.Println("umlog.Init Error : " + err.Error())
|
||||
}
|
||||
|
||||
if MODE == "Release" {
|
||||
umlog.SetLevel("E")
|
||||
} else if MODE == "Info" {
|
||||
umlog.SetLevel("I")
|
||||
} else {
|
||||
umlog.SetLevel("D")
|
||||
}
|
||||
|
||||
walletDB := umsql.DbConnect(DB_HOST, DB_PORT, DB_USER, DB_PASS, DB_NAME)
|
||||
defer umsql.DbDisconnect(walletDB)
|
||||
_, CODE_STRC = src.GetCodeInfo(walletDB, KEY)
|
||||
umlog.Info("main\nCode Info Loaded: %+v", CODE_STRC)
|
||||
|
||||
// 트랜잭션 모니터링 시작 (백그라운드 고루틴)
|
||||
// CYBX 토큰 컨트랙트 주소는 환경에 맞게 설정 필요
|
||||
var contractAddress string
|
||||
contractAddress = CODE_STRC.CYBX_MAINNET
|
||||
|
||||
// 백그라운드에서 트랜잭션 모니터링 시작
|
||||
// go src.StartTransactionMonitor(
|
||||
// src.StartTransactionMonitor(
|
||||
// walletDB, // DB 연결
|
||||
// contractAddress, // CYBX 토큰 컨트랙트 주소
|
||||
// 10*time.Second, // 10초마다 체크
|
||||
// 40, // 한 번에 40명 체크
|
||||
// NETTYPE, // MAINNET 또는 TESTNET
|
||||
// CODE_STRC, // 코드 정보
|
||||
// )
|
||||
|
||||
src.StartTransactionMonitorAlchemy(
|
||||
walletDB,
|
||||
header,
|
||||
contractAddress,
|
||||
3*time.Second,
|
||||
15,
|
||||
CODE_STRC)
|
||||
|
||||
}
|
||||
16
model/code.go
Normal file
16
model/code.go
Normal file
@ -0,0 +1,16 @@
|
||||
package model
|
||||
|
||||
type CodeStrc struct {
|
||||
CIPHERKEY string
|
||||
CIPHERIVKEY string
|
||||
MESSAGE_TOKEN string
|
||||
MESSAGE_ID string
|
||||
SENDBIRD_KEY string
|
||||
SENDBIRD_VERSION string
|
||||
SENDBIRD_TOKEN string
|
||||
ETHERSCAN_TOKEN string
|
||||
CYBX_MAINNET string
|
||||
CYBX_TESTNET string
|
||||
ALCHEMY_MAINNET string
|
||||
ALCHEMY_TESTNET string
|
||||
}
|
||||
34
model/common.go
Normal file
34
model/common.go
Normal file
@ -0,0 +1,34 @@
|
||||
package model
|
||||
|
||||
type Header struct {
|
||||
AppKey string `form:"appKey" json:"appKey" example:"adpsofimerpifgarpo8rt320"`
|
||||
SecretKey string `form:"secretKey" json:"secretKey" example:"sdfjlsdkfjsldkfjlsdkjflsdkjflsdkjflsdkj"`
|
||||
}
|
||||
|
||||
type ServerUrlStrc struct {
|
||||
AAuthUrl string `form:"aauthUrl" json:"aauthUrl" example:"http://aauth.angkorlife.com"`
|
||||
LifeUrl string `form:"lifeUrl" json:"lifeUrl" example:"http://life.angkorlife.com"`
|
||||
WattUrl string `form:"wattUrl" json:"wattUrl" example:"http://watt.angkorlife.com"`
|
||||
}
|
||||
|
||||
type UserInfoStrc struct {
|
||||
AngkorId string `form:"angkorId" json:"angkorId" example:"akDEvQ1693297580"`
|
||||
UserId string `form:"userId" json:"userId" example:"user01"`
|
||||
UserAngkorId string `form:"userAngkorId" json:"userAngkorId" example:"alice"`
|
||||
PhoneNumber string `form:"phoneNumber" json:"phoneNumber" example:"8551012345678"`
|
||||
Email string `form:"email" json:"email" example:"user01@example.com"`
|
||||
}
|
||||
|
||||
type AuthInitStrc struct {
|
||||
GrantType string `form:"grantType" json:"grantType" example:"session"`
|
||||
UserAngkorId string `form:"userAngkorId" json:"userAngkorId" example:"alice"`
|
||||
}
|
||||
|
||||
type AuthKeyStrc struct {
|
||||
TokenType string `json:"tokenType" example:"bearer"`
|
||||
AuthKey string `json:"authKey" example:"adpsofimerpifgarpo8rt320"`
|
||||
}
|
||||
|
||||
type WattInfoStrc struct {
|
||||
WattAmount string `form:"wattAmount" json:"wattAmount" example:"0.0"`
|
||||
}
|
||||
28
model/default.go
Normal file
28
model/default.go
Normal file
@ -0,0 +1,28 @@
|
||||
package model
|
||||
|
||||
type RequestDefaultStrc struct {
|
||||
GrantType string `form:"grantType" json:"grantType" example:"session"`
|
||||
AngkorId string `form:"angkorid" json:"angkorid" example:"ak123214"`
|
||||
}
|
||||
|
||||
type ResponseDefaultModel struct {
|
||||
Code int `json:"code" example:"200"`
|
||||
Message string `json:"message" example:"success"`
|
||||
}
|
||||
|
||||
type DefaultErrorModel struct {
|
||||
Code int `json:"code" example:"400"`
|
||||
Message string `json:"message" example:"message"`
|
||||
}
|
||||
|
||||
type ApplicationInfo struct {
|
||||
AppId string `json:"appId" example:"11"`
|
||||
ApplicationName string `json:"appName" example:"angkorlife"`
|
||||
ApplicationCompany string `json:"appCompany" example:"digital angkor"`
|
||||
ApplicationIcon string `json:"appIcon" example:"https://icon.angkorlifes.com/images/angkorlife.jpg"`
|
||||
ApplicationConsent string `json:"appConsent" example:"127"`
|
||||
ApplicationCallbackUrl string `json:"appCallback" example:"https://callback.angkorlifes.com/receive"`
|
||||
ServiceCategory string `json:"serviceCategory" example:"game"`
|
||||
AppleBundleIDs string `json:"appleBundleIDs" example:"com.angkorlifes"`
|
||||
AndroidPackageName string `json:"androidPackageName" example:"com.angkorlifes"`
|
||||
}
|
||||
693
src/alchemy.go
Normal file
693
src/alchemy.go
Normal file
@ -0,0 +1,693 @@
|
||||
package src
|
||||
|
||||
import (
|
||||
"AngkorWalletScanning/model"
|
||||
"AngkorWalletScanning/umlog"
|
||||
"AngkorWalletScanning/umsql"
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"math/big"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"golang.org/x/time/rate"
|
||||
)
|
||||
|
||||
// Alchemy API Rate Limiter (무료 티어: 500 CU/s, eth_getLogs=75 CU → 약 6.67 req/s)
|
||||
var alchemyApiLimiter = rate.NewLimiter(rate.Limit(6.5), 7)
|
||||
|
||||
// HTTP 클라이언트 초기화 함수
|
||||
func initAlchemyHTTPClient() *http.Client {
|
||||
transport := &http.Transport{
|
||||
Dial: (&net.Dialer{
|
||||
Timeout: 30 * time.Second, // 10초 → 30초로 증가
|
||||
KeepAlive: 30 * time.Second,
|
||||
Resolver: &net.Resolver{
|
||||
PreferGo: false, // 시스템 DNS resolver 사용
|
||||
},
|
||||
}).Dial,
|
||||
TLSHandshakeTimeout: 15 * time.Second, // 10초 → 15초로 증가
|
||||
ResponseHeaderTimeout: 15 * time.Second, // 10초 → 15초로 증가
|
||||
ExpectContinueTimeout: 1 * time.Second,
|
||||
MaxIdleConns: 100,
|
||||
MaxIdleConnsPerHost: 10,
|
||||
IdleConnTimeout: 90 * time.Second,
|
||||
TLSClientConfig: &tls.Config{
|
||||
InsecureSkipVerify: false,
|
||||
},
|
||||
DisableKeepAlives: false,
|
||||
// DNS 캐시 시간 연장
|
||||
MaxConnsPerHost: 10,
|
||||
}
|
||||
|
||||
// 환경 변수 또는 시스템 프록시 사용
|
||||
if proxyURL := os.Getenv("HTTPS_PROXY"); proxyURL != "" {
|
||||
umlog.Debug("[Alchemy] Using proxy from HTTPS_PROXY: %s\n", proxyURL)
|
||||
if proxy, err := url.Parse(proxyURL); err == nil {
|
||||
transport.Proxy = http.ProxyURL(proxy)
|
||||
}
|
||||
} else if proxyURL := os.Getenv("HTTP_PROXY"); proxyURL != "" {
|
||||
umlog.Debug("[Alchemy] Using proxy from HTTP_PROXY: %s\n", proxyURL)
|
||||
if proxy, err := url.Parse(proxyURL); err == nil {
|
||||
transport.Proxy = http.ProxyURL(proxy)
|
||||
}
|
||||
} else {
|
||||
// 시스템 기본 프록시 사용
|
||||
transport.Proxy = http.ProxyFromEnvironment
|
||||
}
|
||||
|
||||
return &http.Client{
|
||||
Timeout: 60 * time.Second, // 30초 → 60초로 증가 (전체 요청 타임아웃)
|
||||
Transport: transport,
|
||||
}
|
||||
}
|
||||
|
||||
// HTTP 클라이언트 (타임아웃 및 연결 설정)
|
||||
var alchemyHTTPClient = initAlchemyHTTPClient()
|
||||
|
||||
// ERC20 Transfer Event Topic (Transfer 이벤트 시그니처 해시)
|
||||
const transferEventTopic = "0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef"
|
||||
|
||||
// Alchemy JSON-RPC Request/Response Structures
|
||||
type AlchemyRequest struct {
|
||||
JsonRPC string `json:"jsonrpc"`
|
||||
Method string `json:"method"`
|
||||
Params []interface{} `json:"params"`
|
||||
ID int `json:"id"`
|
||||
}
|
||||
|
||||
type AlchemyResponse struct {
|
||||
JsonRPC string `json:"jsonrpc"`
|
||||
ID int `json:"id"`
|
||||
Result map[string]interface{} `json:"result"`
|
||||
Error *AlchemyError `json:"error,omitempty"`
|
||||
}
|
||||
|
||||
type AlchemyError struct {
|
||||
Code int `json:"code"`
|
||||
Message string `json:"message"`
|
||||
}
|
||||
|
||||
type AlchemyLog struct {
|
||||
Address string `json:"address"`
|
||||
Topics []string `json:"topics"`
|
||||
Data string `json:"data"`
|
||||
BlockNumber string `json:"blockNumber"`
|
||||
TransactionHash string `json:"transactionHash"`
|
||||
TransactionIndex string `json:"transactionIndex"`
|
||||
BlockHash string `json:"blockHash"`
|
||||
LogIndex string `json:"logIndex"`
|
||||
Removed bool `json:"removed"`
|
||||
}
|
||||
|
||||
type AlchemyTransferInfo struct {
|
||||
BlockNumber int64
|
||||
TransactionHash string
|
||||
From string
|
||||
To string
|
||||
Value string
|
||||
TokenAddress string
|
||||
Timestamp int64
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : Alchemy API를 통해 특정 주소의 ERC20 토큰 Transfer 이벤트 조회
|
||||
parameter
|
||||
- address : 조회할 지갑 주소
|
||||
- contractAddress : 토큰 컨트랙트 주소 (CYBX)
|
||||
- fromBlock : 시작 블록 (hex 형식)
|
||||
- toBlock : 종료 블록 (hex 형식, "latest" 가능)
|
||||
- netType : 네트워크 타입 ("MAINNET" 또는 "TESTNET")
|
||||
- CODE_STRC : 설정 정보
|
||||
return
|
||||
- []AlchemyTransferInfo : Transfer 이벤트 목록
|
||||
- error : 에러 정보
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func getAlchemyTokenTransfers(address string, contractAddress string, fromBlock string, toBlock string, CODE_STRC model.CodeStrc) ([]AlchemyTransferInfo, error) {
|
||||
// 입력 값 검증
|
||||
if address == "" || len(address) < 40 {
|
||||
return nil, fmt.Errorf("invalid wallet address: %s", address)
|
||||
}
|
||||
if contractAddress == "" || len(contractAddress) < 40 {
|
||||
return nil, fmt.Errorf("invalid contract address: %s", contractAddress)
|
||||
}
|
||||
|
||||
// Rate limiter 체크
|
||||
ctx := context.Background()
|
||||
err := alchemyApiLimiter.Wait(ctx)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("rate limiter error: %v", err)
|
||||
}
|
||||
|
||||
// eth_getLogs 요청 파라미터 구성
|
||||
// Transfer(address indexed from, address indexed to, uint256 value)
|
||||
// topic[0] = Transfer 이벤트 시그니처
|
||||
// topic[1] = from 주소 (32바이트 패딩)
|
||||
// topic[2] = to 주소 (32바이트 패딩)
|
||||
|
||||
// 입금 트랜잭션 필터: to = 내 주소
|
||||
paddedAddress := "0x" + strings.Repeat("0", 24) + strings.TrimPrefix(strings.ToLower(address), "0x")
|
||||
|
||||
filterParams := map[string]interface{}{
|
||||
"fromBlock": fromBlock,
|
||||
"toBlock": toBlock,
|
||||
"address": contractAddress,
|
||||
"topics": []interface{}{
|
||||
transferEventTopic, // Transfer 이벤트
|
||||
nil, // from (any)
|
||||
paddedAddress, // to (내 주소)
|
||||
},
|
||||
}
|
||||
|
||||
request := AlchemyRequest{
|
||||
JsonRPC: "2.0",
|
||||
Method: "eth_getLogs",
|
||||
Params: []interface{}{filterParams},
|
||||
ID: 1,
|
||||
}
|
||||
|
||||
// JSON 요청 생성
|
||||
requestBody, err := json.Marshal(request)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to marshal request: %v", err)
|
||||
}
|
||||
|
||||
// HTTP POST 요청
|
||||
umlog.Debug("Request - Address: %s, From: %s, To: %s", address, fromBlock, toBlock)
|
||||
// umlog.Debug("[getAlchemyTokenTransfers] Padded address (topic[2]): %s\n", paddedAddress)
|
||||
// umlog.Debug("[getAlchemyTokenTransfers] Request body: %s", string(requestBody))
|
||||
req, err := http.NewRequest("POST", CODE_STRC.ALCHEMY_MAINNET, bytes.NewBuffer(requestBody))
|
||||
if err != nil {
|
||||
umlog.Warn("Failed to create request: %v\n", err)
|
||||
return nil, fmt.Errorf("failed to create request: %v", err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Accept", "application/json")
|
||||
|
||||
// 재시도 로직 (DNS 에러 대응)
|
||||
var resp *http.Response
|
||||
maxRetries := 3
|
||||
for attempt := 1; attempt <= maxRetries; attempt++ {
|
||||
resp, err = alchemyHTTPClient.Do(req)
|
||||
if err == nil {
|
||||
break // 성공
|
||||
}
|
||||
|
||||
// DNS 에러인 경우 재시도
|
||||
if attempt < maxRetries {
|
||||
umlog.Debug("Attempt %d failed: %v, retrying in 2s...", attempt, err)
|
||||
time.Sleep(2 * time.Second)
|
||||
// 요청 바디 재설정 (한 번 읽히면 소진되므로)
|
||||
req.Body = io.NopCloser(bytes.NewBuffer(requestBody))
|
||||
continue
|
||||
}
|
||||
|
||||
// 최종 실패
|
||||
umlog.Warn("HTTP request failed after %d attempts: %v\n", maxRetries, err)
|
||||
return nil, fmt.Errorf("http request failed: %v", err)
|
||||
}
|
||||
|
||||
defer resp.Body.Close()
|
||||
// umlog.Debug("Response status: %d", resp.StatusCode)
|
||||
|
||||
// 응답 읽기
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to read response: %v", err)
|
||||
}
|
||||
|
||||
// umlog.Debug("Response body: %s\n", string(body))
|
||||
|
||||
// JSON 파싱
|
||||
var alchemyResp struct {
|
||||
JsonRPC string `json:"jsonrpc"`
|
||||
ID int `json:"id"`
|
||||
Result []AlchemyLog `json:"result"`
|
||||
Error *AlchemyError `json:"error,omitempty"`
|
||||
}
|
||||
|
||||
err = json.Unmarshal(body, &alchemyResp)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to parse json: %v", err)
|
||||
}
|
||||
|
||||
// API 에러 체크
|
||||
if alchemyResp.Error != nil {
|
||||
return nil, fmt.Errorf("alchemy api error: %s (code: %d)", alchemyResp.Error.Message, alchemyResp.Error.Code)
|
||||
}
|
||||
|
||||
// Transfer 이벤트 파싱
|
||||
var transfers []AlchemyTransferInfo
|
||||
for _, log := range alchemyResp.Result {
|
||||
if len(log.Topics) < 3 {
|
||||
continue // Invalid Transfer event
|
||||
}
|
||||
umlog.Debug("Log entry: %+v\n", log)
|
||||
|
||||
// from 주소 추출 (topic[1])
|
||||
fromAddr := "0x" + log.Topics[1][26:] // 24개의 0 제거
|
||||
|
||||
// to 주소 추출 (topic[2])
|
||||
toAddr := "0x" + log.Topics[2][26:] // 24개의 0 제거
|
||||
|
||||
// value 추출 (data)
|
||||
value := strings.TrimPrefix(log.Data, "0x")
|
||||
|
||||
// 블록 번호 파싱
|
||||
blockNum := new(big.Int)
|
||||
blockNum.SetString(strings.TrimPrefix(log.BlockNumber, "0x"), 16)
|
||||
|
||||
// 타임스탬프는 별도 조회 필요 (블록 정보에서)
|
||||
timestamp := time.Now().Unix() // 기본값, 필요시 eth_getBlockByNumber 호출
|
||||
|
||||
transfers = append(transfers, AlchemyTransferInfo{
|
||||
BlockNumber: blockNum.Int64(),
|
||||
TransactionHash: log.TransactionHash,
|
||||
From: fromAddr,
|
||||
To: toAddr,
|
||||
Value: value,
|
||||
TokenAddress: log.Address,
|
||||
Timestamp: timestamp,
|
||||
})
|
||||
}
|
||||
|
||||
return transfers, nil
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 현재 블록 번호 조회 (Alchemy API)
|
||||
parameter
|
||||
- netType : 네트워크 타입 ("MAINNET" 또는 "TESTNET")
|
||||
- CODE_STRC : 설정 정보
|
||||
return
|
||||
- int64 : 현재 블록 번호
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func getCurrentBlockNumberAlchemy(CODE_STRC model.CodeStrc) int64 {
|
||||
// Alchemy API URL
|
||||
alchemyURL := CODE_STRC.ALCHEMY_MAINNET
|
||||
|
||||
if alchemyURL == "" {
|
||||
umlog.Error("Alchemy API URL not configured")
|
||||
return 0
|
||||
}
|
||||
|
||||
request := AlchemyRequest{
|
||||
JsonRPC: "2.0",
|
||||
Method: "eth_blockNumber",
|
||||
Params: []interface{}{},
|
||||
ID: 1,
|
||||
}
|
||||
|
||||
requestBody, err := json.Marshal(request)
|
||||
if err != nil {
|
||||
umlog.Warn("Failed to marshal request: %v\n", err)
|
||||
return 0
|
||||
}
|
||||
|
||||
umlog.Info("Requesting: %s", alchemyURL)
|
||||
req, err := http.NewRequest("POST", alchemyURL, bytes.NewBuffer(requestBody))
|
||||
if err != nil {
|
||||
umlog.Warn("Failed to create request: %v\n", err)
|
||||
return 0
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Accept", "application/json")
|
||||
|
||||
resp, err := alchemyHTTPClient.Do(req)
|
||||
if err != nil {
|
||||
umlog.Warn("HTTP request failed: %v\n", err)
|
||||
return 0
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
body, _ := io.ReadAll(resp.Body)
|
||||
|
||||
var result struct {
|
||||
JsonRPC string `json:"jsonrpc"`
|
||||
ID int `json:"id"`
|
||||
Result string `json:"result"`
|
||||
}
|
||||
|
||||
json.Unmarshal(body, &result)
|
||||
|
||||
blockNum := new(big.Int)
|
||||
blockNum.SetString(strings.TrimPrefix(result.Result, "0x"), 16)
|
||||
|
||||
return blockNum.Int64()
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 해당 트랜잭션에 대해 이미 알림을 보냈는지 확인
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- txHash : 트랜잭션 해시
|
||||
return
|
||||
- bool : 이미 알림 보냈으면 true
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func isTransactionNotified(walletDBConn *sql.DB, txHash string) bool {
|
||||
query := "SELECT COUNT(*) FROM ank_wallet_notification WHERE tx_hash = ?"
|
||||
|
||||
var count int
|
||||
res, _, err := umsql.SqlSelect(walletDBConn, query, txHash)
|
||||
if err != nil {
|
||||
umlog.Error("SELECT Error(%s): %s", err.Error(), query)
|
||||
return false
|
||||
}
|
||||
defer res.Close()
|
||||
|
||||
if res.Next() {
|
||||
if err := res.Scan(&count); err != nil {
|
||||
umlog.Error("SCAN Error(%s): %s", err.Error(), query)
|
||||
return false
|
||||
}
|
||||
}
|
||||
return count > 0
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 특정 지갑의 입금 트랜잭션 체크 (Alchemy API 사용)
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- userId : 사용자 ID
|
||||
- walletAddress : 지갑 주소
|
||||
- contractAddress : CYBX 토큰 컨트랙트 주소
|
||||
- lastCheckTime : 마지막 체크 시각
|
||||
- netType : 네트워크 타입
|
||||
- CODE_STRC : 설정 정보
|
||||
return
|
||||
- int : 새로 발견된 입금 건수
|
||||
- error : 에러 정보
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func checkWalletWithAlchemy(walletDBConn *sql.DB, header model.Header, userId int, walletAddress string, contractAddress string, lastCheckTime time.Time, currentBlock int64, CODE_STRC model.CodeStrc) (int, error) {
|
||||
// 현재 블록은 호출자로부터 받음 (배치 시작 시 한 번만 조회)
|
||||
|
||||
// 입력 값 검증
|
||||
if walletAddress == "" {
|
||||
umlog.Warn("User %d has empty wallet address, skipping\n", userId)
|
||||
return 0, nil
|
||||
}
|
||||
if contractAddress == "" {
|
||||
umlog.Warn("Empty contract address, skipping\n")
|
||||
return 0, fmt.Errorf("contract address is empty")
|
||||
}
|
||||
|
||||
// Alchemy 무료 티어 제한: 한 번에 최대 10 블록만 조회 가능
|
||||
const maxBlockRange = 9 // 10블록이 아니라 9블록 차이 (0부터 9까지 = 10개)
|
||||
|
||||
// 마지막 체크 시각을 블록 번호로 변환 (BSC는 약 3초당 1블록)
|
||||
timeDiff := time.Since(lastCheckTime)
|
||||
blocksToCheck := int64(timeDiff.Seconds() / 3)
|
||||
|
||||
// umlog.Debug("Time diff: %v, Calculated blocks: %d\n", timeDiff, blocksToCheck)
|
||||
|
||||
// Alchemy 무료 티어 제한 적용 (최소 1블록은 체크)
|
||||
if blocksToCheck < 1 {
|
||||
blocksToCheck = 1
|
||||
}
|
||||
if blocksToCheck > maxBlockRange {
|
||||
blocksToCheck = maxBlockRange
|
||||
umlog.Debug("Limiting block range to %d due to Alchemy free tier\n", maxBlockRange)
|
||||
}
|
||||
|
||||
startBlock := currentBlock - blocksToCheck
|
||||
if startBlock < 0 {
|
||||
startBlock = 0
|
||||
}
|
||||
|
||||
// 블록 번호를 hex로 변환
|
||||
fromBlockHex := fmt.Sprintf("0x%x", startBlock)
|
||||
toBlockHex := fmt.Sprintf("0x%x", currentBlock)
|
||||
|
||||
umlog.Debug("Checking blocks %d to %d (range: %d)\n", startBlock, currentBlock, blocksToCheck)
|
||||
|
||||
// Alchemy API로 Transfer 이벤트 조회
|
||||
transfers, err := getAlchemyTokenTransfers(walletAddress, contractAddress, fromBlockHex, toBlockHex, CODE_STRC)
|
||||
if err != nil {
|
||||
umlog.Warn("Failed to get transfers for user %d: %v\n", userId, err)
|
||||
return 0, err
|
||||
}
|
||||
|
||||
newTransactionCount := 0
|
||||
|
||||
for _, transfer := range transfers {
|
||||
umlog.Debug("transfer info: %v\n", transfer)
|
||||
|
||||
// 트랜잭션 시각 체크 (Timestamp 사용)
|
||||
txTime := time.Unix(transfer.Timestamp, 0)
|
||||
|
||||
// 마지막 체크 이후의 트랜잭션만 처리
|
||||
if txTime.Before(lastCheckTime) {
|
||||
continue
|
||||
}
|
||||
|
||||
// 이미 알림 보낸 트랜잭션인지 확인
|
||||
if isTransactionNotified(walletDBConn, transfer.TransactionHash) {
|
||||
continue
|
||||
}
|
||||
|
||||
// 금액 계산 (18 decimals for most tokens, CYBX는 확인 필요)
|
||||
amount := calculateTokenAmountFromHex(transfer.Value, "18")
|
||||
|
||||
// 알림 전송
|
||||
sendIncomingTransactionNotification(walletDBConn, header, userId, transfer.From, transfer.To, amount, transfer.TransactionHash)
|
||||
|
||||
// 알림 기록 저장
|
||||
saveNotificationRecord(walletDBConn, userId, transfer.TransactionHash, transfer.From, transfer.To, amount, "CYBX")
|
||||
|
||||
newTransactionCount++
|
||||
|
||||
umlog.Info("checkWalletWithAlchemy\nNew incoming tx detected - User: %d, From: %s, Amount: %s CYBX, Hash: %s",
|
||||
userId, transfer.From, amount, transfer.TransactionHash)
|
||||
}
|
||||
|
||||
return newTransactionCount, nil
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : Hex 값으로 토큰 금액 계산 (decimal 적용)
|
||||
parameter
|
||||
- hexValue : 토큰 raw value (hex 문자열)
|
||||
- decimal : 토큰 decimal
|
||||
return
|
||||
- string : 사람이 읽을 수 있는 금액
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func calculateTokenAmountFromHex(hexValue string, decimal string) string {
|
||||
// hex를 big.Int로 변환
|
||||
valueInt := new(big.Int)
|
||||
if strings.HasPrefix(hexValue, "0x") {
|
||||
hexValue = strings.TrimPrefix(hexValue, "0x")
|
||||
}
|
||||
valueInt.SetString(hexValue, 16)
|
||||
|
||||
// decimal을 int로 변환
|
||||
decimalInt, _ := strconv.Atoi(decimal)
|
||||
|
||||
// 10^decimal 계산
|
||||
divisor := new(big.Int).Exp(big.NewInt(10), big.NewInt(int64(decimalInt)), nil)
|
||||
|
||||
// 나눗셈
|
||||
result := new(big.Float).Quo(
|
||||
new(big.Float).SetInt(valueInt),
|
||||
new(big.Float).SetInt(divisor),
|
||||
)
|
||||
|
||||
// 소수점 4자리까지 표시
|
||||
return result.Text('f', 4)
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 모든 활성 사용자의 지갑 체크 (Alchemy API 사용)
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- contractAddress : CYBX 토큰 컨트랙트 주소
|
||||
- batchSize : 한 번에 처리할 사용자 수
|
||||
- netType : 네트워크 타입
|
||||
- CODE_STRC : 설정 정보
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func checkAllWalletsWithAlchemy(walletDBConn *sql.DB, header model.Header, contractAddress string, batchSize int, CODE_STRC model.CodeStrc) {
|
||||
startTime := time.Now()
|
||||
|
||||
// 배치 시작 시 현재 블록 번호 한 번만 조회
|
||||
currentBlock := getCurrentBlockNumberAlchemy(CODE_STRC)
|
||||
if currentBlock == 0 {
|
||||
umlog.Error("Failed to get current block number")
|
||||
return
|
||||
}
|
||||
umlog.Debug("Current block: %d - will be used for all wallets in this batch", currentBlock)
|
||||
|
||||
// 활성 사용자 목록 조회 (최근 활동 순)
|
||||
query := `
|
||||
SELECT id, address, checkedAt
|
||||
FROM ank_wallet_user
|
||||
WHERE address > ''
|
||||
ORDER BY checkedAt ASC
|
||||
LIMIT ?
|
||||
`
|
||||
|
||||
rows, err := walletDBConn.Query(query, batchSize)
|
||||
if err != nil {
|
||||
umlog.Error("checkAllWalletsWithAlchemy\nDB query failed: %v", err)
|
||||
return
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
totalChecked := 0
|
||||
totalNewTransactions := 0
|
||||
|
||||
// 요청 분산 간격 계산 (batchSize개를 3초에 균등 분산)
|
||||
// 예: 12개 → 3000ms ÷ 12 = 250ms 간격
|
||||
delayBetweenRequests := time.Duration(3000/batchSize) * time.Millisecond
|
||||
|
||||
for rows.Next() {
|
||||
var userId int
|
||||
var walletAddress string
|
||||
var lastCheckStr string
|
||||
|
||||
err := rows.Scan(&userId, &walletAddress, &lastCheckStr)
|
||||
if err != nil {
|
||||
umlog.Error("checkAllWalletsWithAlchemy\nRow scan error: %v", err)
|
||||
continue
|
||||
}
|
||||
|
||||
// 마지막 체크 시각 파싱
|
||||
lastCheckTime, _ := time.Parse("2006-01-02 15:04:05", lastCheckStr)
|
||||
|
||||
// 개별 지갑 체크 (현재 블록 번호를 파라미터로 전달)
|
||||
newTxCount, err := checkWalletWithAlchemy(walletDBConn, header, userId, walletAddress, contractAddress, lastCheckTime, currentBlock, CODE_STRC)
|
||||
if err != nil {
|
||||
// 에러 로그는 checkWalletWithAlchemy 내부에서 이미 출력됨
|
||||
continue
|
||||
}
|
||||
|
||||
// 마지막 체크 시각 업데이트
|
||||
updateLastCheckTime(walletDBConn, userId)
|
||||
|
||||
totalChecked++
|
||||
totalNewTransactions += newTxCount
|
||||
|
||||
// 다음 요청 전 대기 (마지막 요청 후에는 대기 불필요)
|
||||
if totalChecked < batchSize {
|
||||
time.Sleep(delayBetweenRequests)
|
||||
}
|
||||
}
|
||||
|
||||
elapsed := time.Since(startTime)
|
||||
|
||||
umlog.Info("checkAllWalletsWithAlchemy\nBatch completed - Checked: %d wallets, New transactions: %d, Time: %.2f seconds",
|
||||
totalChecked, totalNewTransactions, elapsed.Seconds())
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 트랜잭션 모니터링 시작 (Alchemy API 사용)
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- contractAddress : CYBX 토큰 컨트랙트 주소
|
||||
- checkInterval : 체크 주기
|
||||
- batchSize : 한 번에 처리할 사용자 수
|
||||
- netType : 네트워크 타입
|
||||
- CODE_STRC : 설정 정보
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func StartTransactionMonitorAlchemy(walletDBConn *sql.DB, header model.Header, contractAddress string, checkInterval time.Duration, batchSize int, CODE_STRC model.CodeStrc) {
|
||||
umlog.Debug("[StartTransactionMonitorAlchemy] Starting monitor - Interval: %v, Batch size: %d\n", checkInterval, batchSize)
|
||||
umlog.Debug("[StartTransactionMonitorAlchemy] Contract Address: %s\n", contractAddress)
|
||||
umlog.Debug("[StartTransactionMonitorAlchemy] Alchemy URL: %s\n", CODE_STRC.ALCHEMY_MAINNET)
|
||||
|
||||
ticker := time.NewTicker(checkInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
// 즉시 한 번 실행
|
||||
checkAllWalletsWithAlchemy(walletDBConn, header, contractAddress, batchSize, CODE_STRC)
|
||||
|
||||
// 주기적 실행
|
||||
for range ticker.C {
|
||||
checkAllWalletsWithAlchemy(walletDBConn, header, contractAddress, batchSize, CODE_STRC)
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 입금 알림 전송 (FCM Push + Telegram)
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- userId : 사용자 ID
|
||||
- fromAddress : 발신 주소
|
||||
- toAddress : 수신 주소
|
||||
- amount : 금액
|
||||
- txHash : 트랜잭션 해시
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func sendIncomingTransactionNotification(walletDBConn *sql.DB, header model.Header, userId int, fromAddress string, toAddress string, amount string, txHash string) {
|
||||
// FCM 토큰 조회
|
||||
var fcmToken, angkorId string
|
||||
query := "SELECT fcm_token, angkorid FROM ank_wallet_user WHERE id = ?"
|
||||
umlog.Debug("authCode count query: %s", query)
|
||||
rows, _, err := umsql.SqlSelect(walletDBConn, query, userId)
|
||||
if err != nil {
|
||||
umlog.Error("SELECT Error(%s): %s", err.Error(), query)
|
||||
return
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
if rows.Next() {
|
||||
if scanErr := rows.Scan(&fcmToken, &angkorId); scanErr != nil {
|
||||
umlog.Error("SELECT Scan Error(%s): %s", scanErr.Error(), query)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// FCM Push 알림 (FCM 토큰이 있는 경우)
|
||||
// if fcmToken.Valid && fcmToken.String != "" {
|
||||
// sendFCMPushNotification(fcmToken.String, amount, fromAddress, txHash)
|
||||
// }
|
||||
|
||||
// Telegram 알림 (관리자용)
|
||||
messageboxy := fmt.Sprintf(
|
||||
"💰 Deposit notification(from Alchemy)\nFrom: %s\nTo: %s\nAmount: %s CYBX\nTxHash: %s",
|
||||
fromAddress, toAddress, amount, txHash)
|
||||
message := map[string]string{"message": messageboxy}
|
||||
jsonData, err := json.Marshal(message)
|
||||
if err != nil {
|
||||
umlog.Warn("SendMessage Error(%s): %v", err.Error(), message)
|
||||
}
|
||||
SendMessage(header, angkorId, string(jsonData))
|
||||
}
|
||||
149
src/common.go
Normal file
149
src/common.go
Normal file
@ -0,0 +1,149 @@
|
||||
package src
|
||||
|
||||
import (
|
||||
"AngkorWalletScanning/model"
|
||||
"AngkorWalletScanning/umlog"
|
||||
"AngkorWalletScanning/umsql"
|
||||
"bytes"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"net/http"
|
||||
"strconv"
|
||||
)
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : DB에서 코드값을 읽어온다.
|
||||
parameter
|
||||
- dbConn : DB 핸들러
|
||||
- key : AES 복호화 키
|
||||
return
|
||||
- int : 상태 코드
|
||||
- model.CodeStrc: 코드 구조체
|
||||
error code
|
||||
- 4001202 : DB 연결이 nil
|
||||
- 4001213 : DB 조회 오류
|
||||
- 4001215 : DB 스캔 오류
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func GetCodeInfo(dbConn *sql.DB, key string) (int, model.CodeStrc) {
|
||||
var status int = 200
|
||||
var code_strc model.CodeStrc
|
||||
|
||||
// DB 연결 검증
|
||||
if dbConn == nil {
|
||||
umlog.Error("DB connection is nil")
|
||||
return 4001202, code_strc
|
||||
}
|
||||
|
||||
// 매개변수 검증
|
||||
if len(key) == 0 {
|
||||
umlog.Error("Decryption key is empty")
|
||||
return 4001213, code_strc
|
||||
}
|
||||
|
||||
// SQL 쿼리 (매개변수 바인딩 사용)
|
||||
query := `SELECT code_name, AES_DECRYPT(UNHEX(code_value), ?) AS code_value
|
||||
FROM angkor_code
|
||||
WHERE code_number IN ('001001', '001002', '001010', '001011', '001003', '001004', '001005', '100001', '100002', '100003', '100004', '100005')`
|
||||
|
||||
res, _, err := umsql.SqlSelect(dbConn, query, key)
|
||||
if err != nil {
|
||||
umlog.Error("SELECT Error(%s): %s", err.Error(), query)
|
||||
return 4001213, code_strc
|
||||
}
|
||||
defer res.Close()
|
||||
|
||||
for res.Next() {
|
||||
var code_name, code_value string
|
||||
if err := res.Scan(&code_name, &code_value); err != nil {
|
||||
umlog.Error("SCAN Error(%s): %s", err.Error(), query)
|
||||
return 4001215, code_strc
|
||||
}
|
||||
|
||||
// 코드명에 따라 구조체 필드 설정
|
||||
switch code_name {
|
||||
case "CIPHERKEY":
|
||||
code_strc.CIPHERKEY = code_value
|
||||
case "CIPHERIVKEY":
|
||||
code_strc.CIPHERIVKEY = code_value
|
||||
case "MESSAGE_TOKEN":
|
||||
code_strc.MESSAGE_TOKEN = code_value
|
||||
case "MESSAGE_ID":
|
||||
code_strc.MESSAGE_ID = code_value
|
||||
case "SENDBIRD_KEY":
|
||||
code_strc.SENDBIRD_KEY = code_value
|
||||
case "SENDBIRD_TOKEN":
|
||||
code_strc.SENDBIRD_TOKEN = code_value
|
||||
case "SENDBIRD_VERSION":
|
||||
code_strc.SENDBIRD_VERSION = code_value
|
||||
case "ETHERSCAN":
|
||||
code_strc.ETHERSCAN_TOKEN = code_value
|
||||
case "CYBX_MAINNET":
|
||||
code_strc.CYBX_MAINNET = code_value
|
||||
case "CYBX_TESTNET":
|
||||
code_strc.CYBX_TESTNET = code_value
|
||||
case "ALCHEMY_MAINNET":
|
||||
code_strc.ALCHEMY_MAINNET = code_value
|
||||
case "ALCHEMY_TESTNET":
|
||||
code_strc.ALCHEMY_TESTNET = code_value
|
||||
default:
|
||||
umlog.Warn("Unknown code_name: %s", code_name)
|
||||
}
|
||||
}
|
||||
|
||||
return status, code_strc
|
||||
}
|
||||
|
||||
func SendMessage(header model.Header, angkorId string, message string) error {
|
||||
// Media 서버에서 사용자 정보 조회
|
||||
fullUrl := "https://aauth.angkorlifes.com/message/v1/chat"
|
||||
var sendBody = map[string]any{
|
||||
"grantType": "session",
|
||||
"angkorId": angkorId,
|
||||
"messageType": "WAL",
|
||||
"notiMessage": message,
|
||||
}
|
||||
sendBodyByte, err := json.Marshal(sendBody)
|
||||
if err != nil {
|
||||
umlog.Warn("JSON marshal error: %s", err.Error())
|
||||
return err
|
||||
}
|
||||
req, err := http.NewRequest("POST", fullUrl, bytes.NewBuffer(sendBodyByte))
|
||||
if err != nil {
|
||||
umlog.Warn("Error creating request: => %s : %s", fullUrl, err.Error())
|
||||
return err
|
||||
}
|
||||
//Content-Type 헤더 추가
|
||||
req.Header.Add("Content-Type", "application/json")
|
||||
req.Header.Add("AppKey", header.AppKey)
|
||||
req.Header.Add("Authorization", header.SecretKey)
|
||||
|
||||
client := &http.Client{}
|
||||
resp, err := client.Do(req)
|
||||
if err != nil {
|
||||
umlog.Warn("Error: => %s", err.Error())
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
umlog.Debug("StatusCode: => %d", resp.StatusCode)
|
||||
|
||||
// 응답 본문 읽기
|
||||
resBody, err := ioutil.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
umlog.Warn("Error reading response body: %s", err.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
if resp.StatusCode != 200 {
|
||||
// 실패 응답
|
||||
umlog.Warn("User info retrieval failed - Status: %d, Response: %s", resp.StatusCode, string(resBody))
|
||||
return fmt.Errorf("User info retrieval failed - Status: %s, Response: %s", strconv.Itoa(resp.StatusCode), string(resBody))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
464
src/etherscan.go
Normal file
464
src/etherscan.go
Normal file
@ -0,0 +1,464 @@
|
||||
package src
|
||||
|
||||
import (
|
||||
"AngkorWalletScanning/model"
|
||||
"AngkorWalletScanning/umlog"
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"math/big"
|
||||
"net/http"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"golang.org/x/time/rate"
|
||||
)
|
||||
|
||||
// BSCScan API Rate Limiter (5 requests per second for free tier)
|
||||
var bscApiLimiter = rate.NewLimiter(rate.Limit(5), 5)
|
||||
|
||||
// BSCScan API Response Structures
|
||||
type BSCScanTransaction struct {
|
||||
BlockNumber string `json:"blockNumber"`
|
||||
TimeStamp string `json:"timeStamp"`
|
||||
Hash string `json:"hash"`
|
||||
From string `json:"from"`
|
||||
To string `json:"to"`
|
||||
Value string `json:"value"`
|
||||
ContractAddress string `json:"contractAddress"`
|
||||
TokenName string `json:"tokenName"`
|
||||
TokenSymbol string `json:"tokenSymbol"`
|
||||
TokenDecimal string `json:"tokenDecimal"`
|
||||
TransactionIndex string `json:"transactionIndex"`
|
||||
Gas string `json:"gas"`
|
||||
GasPrice string `json:"gasPrice"`
|
||||
GasUsed string `json:"gasUsed"`
|
||||
CumulativeGasUsed string `json:"cumulativeGasUsed"`
|
||||
Input string `json:"input"`
|
||||
Confirmations string `json:"confirmations"`
|
||||
}
|
||||
|
||||
type BSCScanResponse struct {
|
||||
Status string `json:"status"`
|
||||
Message string `json:"message"`
|
||||
Result []BSCScanTransaction `json:"result"`
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : BSCScan API를 통해 특정 주소의 토큰 트랜잭션 조회
|
||||
parameter
|
||||
- address : 조회할 지갑 주소
|
||||
- contractAddress : 토큰 컨트랙트 주소 (CYBX)
|
||||
- startBlock : 시작 블록 번호 (0 = 처음부터)
|
||||
- netType : 네트워크 타입 ("MAINNET" 또는 "TESTNET")
|
||||
return
|
||||
- []BSCScanTransaction : 트랜잭션 목록
|
||||
- error : 에러 정보
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func getBSCScanTokenTransactions(address string, contractAddress string, startBlock int64, netType string, CODE_STRC model.CodeStrc) ([]BSCScanTransaction, error) {
|
||||
// Rate limiter 체크 (5 requests per second)
|
||||
ctx := context.Background()
|
||||
err := bscApiLimiter.Wait(ctx)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("rate limiter error: %v", err)
|
||||
}
|
||||
|
||||
// chainid 결정 (MAINNET: 56, TESTNET: 97)
|
||||
chainId := "56" // BSC Mainnet
|
||||
if netType == "TESTNET" {
|
||||
chainId = "97" // BSC Testnet
|
||||
}
|
||||
|
||||
// Etherscan API V2 URL 구성
|
||||
url := fmt.Sprintf(
|
||||
"https://api.etherscan.io/v2/api?chainid=%s&module=account&action=tokentx&contractaddress=%s&address=%s&startblock=%d&endblock=999999999&sort=desc&apikey=%s",
|
||||
chainId,
|
||||
contractAddress,
|
||||
address,
|
||||
startBlock,
|
||||
CODE_STRC.ETHERSCAN_TOKEN,
|
||||
)
|
||||
|
||||
// HTTP GET 요청
|
||||
resp, err := http.Get(url)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("http request failed: %v", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
// 응답 읽기
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to read response: %v", err)
|
||||
}
|
||||
|
||||
// JSON 파싱
|
||||
var bscResp BSCScanResponse
|
||||
err = json.Unmarshal(body, &bscResp)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to parse json: %v", err)
|
||||
}
|
||||
|
||||
// API 응답 상태 체크
|
||||
if bscResp.Status != "1" {
|
||||
if bscResp.Message == "No transactions found" {
|
||||
return []BSCScanTransaction{}, nil
|
||||
}
|
||||
return nil, fmt.Errorf("bscscan api error: %s", bscResp.Message)
|
||||
}
|
||||
|
||||
return bscResp.Result, nil
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 특정 지갑의 입금 트랜잭션 체크 및 알림 전송
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- userId : 사용자 ID
|
||||
- walletAddress : 지갑 주소
|
||||
- contractAddress : CYBX 토큰 컨트랙트 주소
|
||||
- lastCheckTime : 마지막 체크 시각
|
||||
- netType : 네트워크 타입 ("MAINNET" 또는 "TESTNET")
|
||||
return
|
||||
- int : 새로 발견된 입금 건수
|
||||
- error : 에러 정보
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func checkWalletWithAPI(walletDBConn *sql.DB, header model.Header, userId int, walletAddress string, contractAddress string, lastCheckTime time.Time, netType string, CODE_STRC model.CodeStrc) (int, error) {
|
||||
// 마지막 체크 시각을 블록 번호로 변환 (BSC는 약 3초당 1블록)
|
||||
timeDiff := time.Since(lastCheckTime)
|
||||
blocksToCheck := int64(timeDiff.Seconds() / 3) // 3초당 1블록
|
||||
|
||||
// 최소 100블록은 체크 (약 5분)
|
||||
if blocksToCheck < 100 {
|
||||
blocksToCheck = 100
|
||||
}
|
||||
|
||||
// 현재 블록에서 역산
|
||||
currentBlock := getCurrentBlockNumber(netType)
|
||||
startBlock := currentBlock - blocksToCheck
|
||||
if startBlock < 0 {
|
||||
startBlock = 0
|
||||
}
|
||||
|
||||
// BSCScan API로 트랜잭션 조회
|
||||
transactions, err := getBSCScanTokenTransactions(walletAddress, contractAddress, startBlock, netType, CODE_STRC)
|
||||
if err != nil {
|
||||
umlog.Error("checkWalletWithAPI\nFailed to get transactions for user %d: %v", userId, err)
|
||||
return 0, err
|
||||
}
|
||||
|
||||
newTransactionCount := 0
|
||||
|
||||
// 입금 트랜잭션 필터링 (To == 내 지갑 주소)
|
||||
for _, tx := range transactions {
|
||||
// 소문자로 변환하여 비교
|
||||
if strings.ToLower(tx.To) != strings.ToLower(walletAddress) {
|
||||
continue // 출금 트랜잭션은 무시
|
||||
}
|
||||
|
||||
// 트랜잭션 시각 체크
|
||||
txTimestamp, _ := strconv.ParseInt(tx.TimeStamp, 10, 64)
|
||||
txTime := time.Unix(txTimestamp, 0)
|
||||
|
||||
// 마지막 체크 이후의 트랜잭션만 처리
|
||||
if txTime.Before(lastCheckTime) {
|
||||
continue
|
||||
}
|
||||
|
||||
// 이미 알림 보낸 트랜잭션인지 확인
|
||||
if isTransactionNotified(walletDBConn, tx.Hash) {
|
||||
continue
|
||||
}
|
||||
|
||||
// 금액 계산 (토큰 decimal 적용)
|
||||
amount := calculateTokenAmount(tx.Value, tx.TokenDecimal)
|
||||
|
||||
// 알림 전송
|
||||
sendIncomingTransactionNotification(walletDBConn, header, userId, tx.From, walletAddress, amount, tx.Hash)
|
||||
|
||||
// 알림 기록 저장
|
||||
saveNotificationRecord(walletDBConn, userId, tx.Hash, tx.From, tx.To, amount, tx.TokenSymbol)
|
||||
|
||||
newTransactionCount++
|
||||
|
||||
umlog.Info("checkWalletWithAPI\nNew incoming tx detected - User: %d, From: %s, Amount: %s %s, Hash: %s",
|
||||
userId, tx.From, amount, tx.TokenSymbol, tx.Hash)
|
||||
}
|
||||
|
||||
return newTransactionCount, nil
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 현재 BSC 블록 번호 조회
|
||||
parameter
|
||||
- netType : 네트워크 타입 ("MAINNET" 또는 "TESTNET")
|
||||
return
|
||||
- int64 : 현재 블록 번호
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func getCurrentBlockNumber(netType string) int64 {
|
||||
apiKey := os.Getenv("BSCSCAN_API_KEY")
|
||||
if apiKey == "" {
|
||||
return 0
|
||||
}
|
||||
|
||||
// chainid 결정 (MAINNET: 56, TESTNET: 97)
|
||||
chainId := "56" // BSC Mainnet
|
||||
if netType == "TESTNET" {
|
||||
chainId = "97" // BSC Testnet
|
||||
}
|
||||
|
||||
url := fmt.Sprintf("https://api.etherscan.io/v2/api?chainid=%s&module=proxy&action=eth_blockNumber&apikey=%s", chainId, apiKey)
|
||||
|
||||
resp, err := http.Get(url)
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
|
||||
var result map[string]interface{}
|
||||
json.Unmarshal(body, &result)
|
||||
|
||||
if blockHex, ok := result["result"].(string); ok {
|
||||
blockNum := new(big.Int)
|
||||
blockNum.SetString(blockHex[2:], 16) // "0x" 제거 후 16진수 파싱
|
||||
return blockNum.Int64()
|
||||
}
|
||||
|
||||
return 0
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 토큰 금액 계산 (decimal 적용)
|
||||
parameter
|
||||
- value : 토큰 raw value (wei 단위)
|
||||
- decimal : 토큰 decimal
|
||||
return
|
||||
- string : 사람이 읽을 수 있는 금액
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func calculateTokenAmount(value string, decimal string) string {
|
||||
// value를 big.Int로 변환
|
||||
valueInt := new(big.Int)
|
||||
valueInt.SetString(value, 10)
|
||||
|
||||
// decimal을 int로 변환
|
||||
decimalInt, _ := strconv.Atoi(decimal)
|
||||
|
||||
// 10^decimal 계산
|
||||
divisor := new(big.Int).Exp(big.NewInt(10), big.NewInt(int64(decimalInt)), nil)
|
||||
|
||||
// 나눗셈
|
||||
result := new(big.Float).Quo(
|
||||
new(big.Float).SetInt(valueInt),
|
||||
new(big.Float).SetInt(divisor),
|
||||
)
|
||||
|
||||
// 소수점 4자리까지 표시
|
||||
return result.Text('f', 4)
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 알림 기록을 DB에 저장
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- userId : 사용자 ID
|
||||
- txHash : 트랜잭션 해시
|
||||
- txFrom : 발신 주소
|
||||
- txTo : 수신 주소
|
||||
- amount : 금액
|
||||
- token : 토큰 심볼
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func saveNotificationRecord(walletDBConn *sql.DB, userId int, txHash string, txFrom string, txTo string, amount string, token string) {
|
||||
query := `
|
||||
INSERT INTO ank_wallet_notification
|
||||
(user_id, tx_hash, tx_from, tx_to, amount, token, notified_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, NOW())
|
||||
`
|
||||
|
||||
_, err := walletDBConn.Exec(query, userId, txHash, txFrom, txTo, amount, token)
|
||||
if err != nil {
|
||||
umlog.Error("saveNotificationRecord", fmt.Sprintf("Failed to save notification: %v", err))
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 사용자의 마지막 체크 시각 업데이트
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- userId : 사용자 ID
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func updateLastCheckTime(walletDBConn *sql.DB, userId int) {
|
||||
query := "UPDATE ank_wallet_user SET checkedAt = NOW() WHERE id = ?"
|
||||
|
||||
_, err := walletDBConn.Exec(query, userId)
|
||||
if err != nil {
|
||||
umlog.Error("updateLastCheckTime\nFailed to update last_check_time for user %d: %v", userId, err)
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : FCM Push 알림 전송
|
||||
parameter
|
||||
- fcmToken : FCM 토큰
|
||||
- amount : 입금 금액
|
||||
- fromAddress : 발신 주소
|
||||
- txHash : 트랜잭션 해시
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func sendFCMPushNotification(fcmToken string, amount string, fromAddress string, txHash string) {
|
||||
// TODO: Firebase Admin SDK를 사용하여 FCM Push 알림 구현
|
||||
// 여기서는 구조만 제공하고, 실제 구현은 프로젝트 환경에 맞게 작성 필요
|
||||
|
||||
umlog.Info("sendFCMPushNotification\nSending FCM push to token: %s, Amount: %s CYBX", fcmToken[:20]+"...", amount)
|
||||
|
||||
// Example:
|
||||
// message := &messaging.Message{
|
||||
// Token: fcmToken,
|
||||
// Notification: &messaging.Notification{
|
||||
// Title: "입금 알림",
|
||||
// Body: fmt.Sprintf("%s CYBX가 입금되었습니다.", amount),
|
||||
// },
|
||||
// Data: map[string]string{
|
||||
// "type": "incoming_transaction",
|
||||
// "amount": amount,
|
||||
// "from": fromAddress,
|
||||
// "txHash": txHash,
|
||||
// },
|
||||
// }
|
||||
//
|
||||
// _, err := fcmClient.Send(ctx, message)
|
||||
// if err != nil {
|
||||
// umlog.Error("sendFCMPushNotification", fmt.Sprintf("FCM send failed: %v", err))
|
||||
// }
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 모든 활성 사용자의 지갑 체크 (배치 처리)
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- contractAddress : CYBX 토큰 컨트랙트 주소
|
||||
- batchSize : 한 번에 처리할 사용자 수 (기본 300)
|
||||
- netType : 네트워크 타입 ("MAINNET" 또는 "TESTNET")
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func checkAllWalletsForIncomingTransactions(walletDBConn *sql.DB, header model.Header, contractAddress string, batchSize int, netType string, CODE_STRC model.CodeStrc) {
|
||||
startTime := time.Now()
|
||||
|
||||
// 활성 사용자 목록 조회 (최근 활동 순)
|
||||
query := `
|
||||
SELECT id, address, checkedAt
|
||||
FROM ank_wallet_user
|
||||
WHERE address > ''
|
||||
ORDER BY checkedAt ASC
|
||||
LIMIT ?
|
||||
`
|
||||
|
||||
rows, err := walletDBConn.Query(query, batchSize)
|
||||
if err != nil {
|
||||
umlog.Error("checkAllWalletsForIncomingTransactions\nDB query failed: %v", err)
|
||||
return
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
totalChecked := 0
|
||||
totalNewTransactions := 0
|
||||
|
||||
for rows.Next() {
|
||||
var userId int
|
||||
var walletAddress string
|
||||
var lastCheckStr string
|
||||
|
||||
err := rows.Scan(&userId, &walletAddress, &lastCheckStr)
|
||||
if err != nil {
|
||||
umlog.Error("checkAllWalletsForIncomingTransactions\nRow scan error: %v", err)
|
||||
continue
|
||||
}
|
||||
|
||||
// 마지막 체크 시각 파싱
|
||||
lastCheckTime, _ := time.Parse("2006-01-02 15:04:05", lastCheckStr)
|
||||
|
||||
// 개별 지갑 체크
|
||||
newTxCount, err := checkWalletWithAPI(walletDBConn, header, userId, walletAddress, contractAddress, lastCheckTime, netType, CODE_STRC)
|
||||
if err != nil {
|
||||
// 에러 로그는 checkWalletWithAPI 내부에서 이미 출력됨
|
||||
continue
|
||||
}
|
||||
|
||||
// 마지막 체크 시각 업데이트
|
||||
updateLastCheckTime(walletDBConn, userId)
|
||||
|
||||
totalChecked++
|
||||
totalNewTransactions += newTxCount
|
||||
}
|
||||
|
||||
elapsed := time.Since(startTime)
|
||||
|
||||
umlog.Info("checkAllWalletsForIncomingTransactions\nBatch completed - Checked: %d wallets, New transactions: %d, Time: %.2f seconds",
|
||||
totalChecked, totalNewTransactions, elapsed.Seconds())
|
||||
}
|
||||
|
||||
/*
|
||||
**************************************************************
|
||||
|
||||
desc : 트랜잭션 모니터링 시작 (백그라운드 고루틴)
|
||||
parameter
|
||||
- walletDBConn : Wallet DB 핸들러
|
||||
- contractAddress : CYBX 토큰 컨트랙트 주소
|
||||
- checkInterval : 체크 주기 (예: 5분)
|
||||
- batchSize : 한 번에 처리할 사용자 수
|
||||
- netType : 네트워크 타입 ("MAINNET" 또는 "TESTNET")
|
||||
|
||||
**************************************************************
|
||||
*/
|
||||
func StartTransactionMonitor(walletDBConn *sql.DB, header model.Header, contractAddress string, checkInterval time.Duration, batchSize int, netType string, CODE_STRC model.CodeStrc) {
|
||||
umlog.Info("StartTransactionMonitor: Transaction monitor started - Network: %s, Interval: %v, Batch size: %d", netType, checkInterval, batchSize)
|
||||
|
||||
ticker := time.NewTicker(checkInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
// 즉시 한 번 실행
|
||||
checkAllWalletsForIncomingTransactions(walletDBConn, header, contractAddress, batchSize, netType, CODE_STRC)
|
||||
|
||||
// // 주기적 실행
|
||||
// for range ticker.C {
|
||||
// checkAllWalletsForIncomingTransactions(walletDBConn, header, contractAddress, batchSize, netType, CODE_STRC)
|
||||
// }
|
||||
}
|
||||
103
umlog/umlog.go
Normal file
103
umlog/umlog.go
Normal file
@ -0,0 +1,103 @@
|
||||
package umlog
|
||||
|
||||
import (
|
||||
"AngkorWalletScanning/ummessage"
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"runtime"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
type Logger struct {
|
||||
Debug *log.Logger
|
||||
Info *log.Logger
|
||||
Warn *log.Logger
|
||||
Error *log.Logger
|
||||
}
|
||||
|
||||
var (
|
||||
logFile *os.File
|
||||
umLogger Logger
|
||||
log_level string = ""
|
||||
)
|
||||
|
||||
// 패키지 외부에서 사용하기 위해서는 첫문자는 대문자로...
|
||||
func Init(logfile string, messageToken string, messageId string) error {
|
||||
var err error
|
||||
logFile, err = os.OpenFile(logfile, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o666)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
umLogger.Debug = log.New(logFile, "[DEBUG] ", log.Ldate|log.Ltime)
|
||||
umLogger.Info = log.New(logFile, "[INFO] ", log.Ldate|log.Ltime)
|
||||
umLogger.Warn = log.New(logFile, "[WARN] ", log.Ldate|log.Ltime)
|
||||
umLogger.Error = log.New(logFile, "[ERROR] ", log.Ldate|log.Ltime)
|
||||
|
||||
ummessage.InitTelegramBot(messageToken, messageId)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func Close() {
|
||||
logFile.Close()
|
||||
}
|
||||
|
||||
func SetLevel(level string) {
|
||||
switch level {
|
||||
case "I":
|
||||
log_level = "info"
|
||||
case "W":
|
||||
log_level = "warn"
|
||||
case "E":
|
||||
log_level = "error"
|
||||
default:
|
||||
log_level = "debug"
|
||||
}
|
||||
}
|
||||
|
||||
func debugInfo() (string, string, string) {
|
||||
funcname, filename, line, _ := runtime.Caller(2)
|
||||
functionname := runtime.FuncForPC(funcname).Name()
|
||||
filenames := strings.Split(filename, "/")
|
||||
return filenames[len(filenames)-1], functionname, strconv.Itoa(line)
|
||||
}
|
||||
|
||||
func Message(title string, body string) {
|
||||
chatId := int64(-5029236432)
|
||||
token := "7108014537:AAHU7299mCzHsBJp1KmPMp7tZEGfGsQ-ffM"
|
||||
ummessage.SendMessageDirect(chatId, token, title, body)
|
||||
}
|
||||
func Error(format string, args ...interface{}) {
|
||||
if log_level == "error" || log_level == "warn" || log_level == "info" || log_level == "debug" {
|
||||
filename, function, line := debugInfo()
|
||||
message := fmt.Sprintf(format, args...)
|
||||
umLogger.Error.Println("[" + filename + "][" + function + "][" + line + "] " + message)
|
||||
ummessage.SendMessage("["+filename+"]["+function+"]["+line+"]", message)
|
||||
}
|
||||
}
|
||||
|
||||
func Warn(format string, args ...interface{}) {
|
||||
if log_level == "warn" || log_level == "info" || log_level == "debug" {
|
||||
filename, _, line := debugInfo()
|
||||
message := fmt.Sprintf(format, args...)
|
||||
umLogger.Warn.Println("[" + filename + "][" + line + "] " + message)
|
||||
}
|
||||
}
|
||||
|
||||
func Info(format string, args ...interface{}) {
|
||||
if log_level == "info" || log_level == "debug" {
|
||||
filename, _, line := debugInfo()
|
||||
message := fmt.Sprintf(format, args...)
|
||||
umLogger.Info.Println("[" + filename + "][" + line + "] " + message)
|
||||
}
|
||||
}
|
||||
|
||||
func Debug(format string, args ...interface{}) {
|
||||
if log_level == "debug" {
|
||||
filename, function, line := debugInfo()
|
||||
message := fmt.Sprintf(format, args...)
|
||||
umLogger.Debug.Println("[" + filename + "][" + function + "][" + line + "] " + message)
|
||||
}
|
||||
}
|
||||
73
ummessage/ummessage.go
Normal file
73
ummessage/ummessage.go
Normal file
@ -0,0 +1,73 @@
|
||||
package ummessage
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"net/http"
|
||||
"time"
|
||||
)
|
||||
|
||||
var telegramToken string
|
||||
var telegramChatId string
|
||||
|
||||
func InitTelegramBot(token string, chatId string) {
|
||||
telegramToken = token
|
||||
telegramChatId = chatId
|
||||
}
|
||||
|
||||
type SendMessageModel struct {
|
||||
Parse_mode string
|
||||
Chat_id string
|
||||
Disable_web_page_preview bool
|
||||
Text string
|
||||
}
|
||||
|
||||
func SendMessage(title string, body string) {
|
||||
message := `<b>` + title + `</b>
|
||||
<b>발생시간(UTC)</b> : <code>` + time.Now().Format("2006-01-02 15:04:05") + `</code>
|
||||
<b>내용</b> :
|
||||
<pre>` + body + `</pre>`
|
||||
messageBody, _ := json.Marshal(map[string]any{
|
||||
"chat_id": telegramChatId,
|
||||
"text": message,
|
||||
"parse_mode": "html",
|
||||
})
|
||||
sendUrl := "https://api.telegram.org/bot" + telegramToken + "/sendMessage"
|
||||
resp, err := http.Post(sendUrl, "application/json", bytes.NewBuffer(messageBody))
|
||||
if err != nil {
|
||||
fmt.Println(err.Error())
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
// 결과 출력
|
||||
_, err = ioutil.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
fmt.Println(err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func SendMessageDirect(chatId int64, telegramToken string, title string, body string) {
|
||||
message := `<b>` + title + `</b>
|
||||
<b>발생시간(UTC)</b> : <code>` + time.Now().Format("2006-01-02 15:04:05") + `</code>
|
||||
<b>내용</b> :
|
||||
` + body
|
||||
messageBody, _ := json.Marshal(map[string]any{
|
||||
"chat_id": chatId,
|
||||
"text": message,
|
||||
"parse_mode": "html",
|
||||
})
|
||||
sendUrl := "https://api.telegram.org/bot" + telegramToken + "/sendMessage"
|
||||
resp, err := http.Post(sendUrl, "application/json", bytes.NewBuffer(messageBody))
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
// 결과 출력
|
||||
_, err = ioutil.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
fmt.Println(err.Error())
|
||||
}
|
||||
}
|
||||
100
umredis/umredis-cluster.go
Normal file
100
umredis/umredis-cluster.go
Normal file
@ -0,0 +1,100 @@
|
||||
package umredis
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"AngkorWalletScanning/umlog"
|
||||
|
||||
"github.com/redis/go-redis/v9"
|
||||
)
|
||||
|
||||
// version : 1.0.0
|
||||
// date : 2022-10-19
|
||||
// writer : lee.beomhee
|
||||
// modify : lee.beomhee (2022-10-19)
|
||||
|
||||
func RedisClusterConnect(host []string, password string) *redis.ClusterClient {
|
||||
var rdb *redis.ClusterClient
|
||||
if len(password) > 0 {
|
||||
rdb = redis.NewClusterClient(&redis.ClusterOptions{
|
||||
Addrs: host,
|
||||
Password: password,
|
||||
})
|
||||
} else {
|
||||
rdb = redis.NewClusterClient(&redis.ClusterOptions{
|
||||
Addrs: host,
|
||||
})
|
||||
}
|
||||
return rdb
|
||||
}
|
||||
|
||||
func RedisClusterDisConnect(cluster *redis.ClusterClient) {
|
||||
umlog.Debug("REDIS DISCONNECTION")
|
||||
cluster.Close()
|
||||
}
|
||||
|
||||
func ClusterGetKeys(cluster *redis.ClusterClient, keys string) ([]string, error) {
|
||||
return cluster.Keys(context.Background(), keys+"*").Result()
|
||||
}
|
||||
|
||||
func ClusterGetValue(cluster *redis.ClusterClient, key string) (string, error) {
|
||||
|
||||
return cluster.Get(context.Background(), key).Result()
|
||||
}
|
||||
|
||||
func ClusterHGetValue(cluster *redis.ClusterClient, key string, field string) (string, error) {
|
||||
return cluster.HGet(context.Background(), key, field).Result()
|
||||
}
|
||||
|
||||
func ClusterHGetAllValue(cluster *redis.ClusterClient, key string) (map[string]string, error) {
|
||||
return cluster.HGetAll(context.Background(), key).Result()
|
||||
}
|
||||
|
||||
func ClusterSetValue(cluster *redis.ClusterClient, key string, value any) error {
|
||||
_, err := cluster.Set(context.Background(), key, value, 0).Result()
|
||||
return err
|
||||
}
|
||||
|
||||
func ClusterHMSetValue(cluster *redis.ClusterClient, key string, value any) error {
|
||||
_, err := cluster.HMSet(context.Background(), key, value).Result()
|
||||
return err
|
||||
}
|
||||
|
||||
func ClusterHSetValue(cluster *redis.ClusterClient, key string, field string, value any) error {
|
||||
_, err := cluster.HSet(context.Background(), key, field, value).Result()
|
||||
return err
|
||||
}
|
||||
|
||||
func ClusterHINCRBY(cluster *redis.ClusterClient, key string, field string, value int64) error {
|
||||
_, err := cluster.HIncrBy(context.Background(), key, field, value).Result()
|
||||
return err
|
||||
}
|
||||
|
||||
func ClusterRemoveValue(cluster *redis.ClusterClient, key string, field string) error {
|
||||
_, err := cluster.HDel(context.Background(), key, field).Result()
|
||||
return err
|
||||
}
|
||||
|
||||
func ClusterRemoveKey(cluster *redis.ClusterClient, key string) error {
|
||||
_, err := cluster.Del(context.Background(), key).Result()
|
||||
return err
|
||||
}
|
||||
|
||||
func ClusterRemoveKeys(cluster *redis.ClusterClient, keys string) error {
|
||||
findKeys, err := cluster.Keys(context.Background(), keys+"*").Result()
|
||||
if err == nil {
|
||||
for _, key := range findKeys {
|
||||
_, err = cluster.Del(context.Background(), key).Result()
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func ClusterExpireKey(cluster *redis.ClusterClient, key string, sec time.Duration) error {
|
||||
_, err := cluster.Expire(context.Background(), key, sec).Result()
|
||||
return err
|
||||
}
|
||||
209
umsql/umsql.go
Normal file
209
umsql/umsql.go
Normal file
@ -0,0 +1,209 @@
|
||||
package umsql
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"AngkorWalletScanning/model"
|
||||
"AngkorWalletScanning/umlog"
|
||||
|
||||
_ "github.com/go-sql-driver/mysql"
|
||||
)
|
||||
|
||||
// type DB_handler struct {
|
||||
// conn *sql.DB
|
||||
// }
|
||||
|
||||
func DbConnect(host string, port string, user string, password string, dbname string) *sql.DB {
|
||||
conn, connect_err := sql.Open("mysql", user+":"+password+"@tcp("+host+":"+port+")/"+dbname)
|
||||
if connect_err != nil {
|
||||
umlog.Error(connect_err.Error())
|
||||
panic(connect_err.Error())
|
||||
}
|
||||
conn.SetConnMaxLifetime(time.Minute * 3)
|
||||
conn.SetMaxOpenConns(1000)
|
||||
conn.SetMaxIdleConns(1000)
|
||||
return conn
|
||||
}
|
||||
|
||||
func DbDisconnect(connA *sql.DB) {
|
||||
umlog.Debug("DataBase DISCONNECTION")
|
||||
err := connA.Close()
|
||||
if err != nil {
|
||||
umlog.Warn("DataBase DISCONNECTION Error A: " + err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func DbBegin(conn *sql.DB) (*sql.Tx, error) {
|
||||
if conn != nil {
|
||||
return conn.Begin()
|
||||
} else {
|
||||
return nil, errors.New("db connecter is nil")
|
||||
}
|
||||
}
|
||||
|
||||
func DbRollback(tx *sql.Tx) error {
|
||||
return tx.Rollback()
|
||||
}
|
||||
|
||||
func DbCommit(tx *sql.Tx) error {
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
func SqlSelect(conn *sql.DB, query string, args ...any) (*sql.Rows, []string, error) {
|
||||
var rows *sql.Rows
|
||||
var err error
|
||||
|
||||
if conn != nil {
|
||||
if len(args) == 0 {
|
||||
rows, err = conn.Query(query)
|
||||
} else {
|
||||
rows, err = conn.Query(query, args...)
|
||||
}
|
||||
if err != nil {
|
||||
umlog.Error(err.Error())
|
||||
return nil, nil, err
|
||||
}
|
||||
cols, err := rows.Columns()
|
||||
if err != nil {
|
||||
umlog.Error(err.Error())
|
||||
return nil, nil, err
|
||||
}
|
||||
return rows, cols, nil
|
||||
}
|
||||
umlog.Error("db connecter is nil")
|
||||
return nil, nil, errors.New("db connecter is nil")
|
||||
}
|
||||
|
||||
func SqlTxExecute(tx *sql.Tx, query string, args ...any) (int64, error) {
|
||||
var result sql.Result
|
||||
var err error
|
||||
|
||||
if tx != nil {
|
||||
if len(args) == 0 {
|
||||
result, err = tx.Exec(query)
|
||||
} else {
|
||||
result, err = tx.Exec(query, args...)
|
||||
}
|
||||
if err != nil {
|
||||
umlog.Error(err.Error())
|
||||
return -1, err
|
||||
}
|
||||
n, err := result.RowsAffected()
|
||||
if err != nil {
|
||||
umlog.Error(err.Error())
|
||||
return -1, err
|
||||
}
|
||||
return n, nil
|
||||
} else {
|
||||
umlog.Error("tx is nil")
|
||||
return -1, errors.New("tx is nil")
|
||||
}
|
||||
}
|
||||
|
||||
func SqlExecute(conn *sql.DB, query string, args ...any) (int64, error) {
|
||||
var result sql.Result
|
||||
var err error
|
||||
|
||||
if conn != nil {
|
||||
if len(args) == 0 {
|
||||
result, err = conn.Exec(query)
|
||||
} else {
|
||||
result, err = conn.Exec(query, args...)
|
||||
}
|
||||
if err != nil {
|
||||
umlog.Error(err.Error())
|
||||
return -1, err
|
||||
}
|
||||
n, err := result.RowsAffected()
|
||||
if err != nil {
|
||||
umlog.Error(err.Error())
|
||||
return -1, err
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
umlog.Error("db connecter is nil")
|
||||
return -1, errors.New("db connecter is nil")
|
||||
}
|
||||
|
||||
func SqlStatement(conn *sql.DB, sql string, param [][]string) error {
|
||||
stmt, err := conn.Prepare(sql)
|
||||
if err != nil {
|
||||
umlog.Error(err.Error())
|
||||
return err
|
||||
}
|
||||
defer stmt.Close()
|
||||
|
||||
for i, sub_val := range param {
|
||||
fmt.Println(sub_val)
|
||||
err = sqlStatementExec(stmt, param[i], len(param[i]))
|
||||
if err != nil {
|
||||
umlog.Error(err.Error())
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func sqlStatementExec(stmt *sql.Stmt, param []string, param_count int) error {
|
||||
args := make([]interface{}, param_count)
|
||||
for i, v := range param {
|
||||
args[i] = v
|
||||
}
|
||||
_, err := stmt.Exec(args...)
|
||||
return err
|
||||
}
|
||||
|
||||
func SqlTransaction(conn *sql.DB, maxRetries int, txFunc func(*sql.Tx) error) model.DefaultErrorModel {
|
||||
var err error
|
||||
var value string
|
||||
for i := 0; i <= maxRetries; i++ {
|
||||
tx, beginErr := conn.Begin()
|
||||
if beginErr != nil {
|
||||
return model.DefaultErrorModel{Code: 4001201, Message: fmt.Sprintf("failed to begin transaction: %v", beginErr)}
|
||||
}
|
||||
|
||||
err = txFunc(tx)
|
||||
|
||||
if err != nil {
|
||||
// 트랜잭션 롤백
|
||||
_ = tx.Rollback()
|
||||
|
||||
// 1205 오류인지 확인
|
||||
if isLockWaitTimeout(err) {
|
||||
umlog.Info(value + ": Transaction lock timeout (1205), retrying... (" + strconv.Itoa(i+1) + "/" + strconv.Itoa(maxRetries) + ")")
|
||||
time.Sleep(time.Duration(100*(i+1)) * time.Millisecond) // simple backoff
|
||||
continue
|
||||
}
|
||||
return model.DefaultErrorModel{Code: 4001201, Message: fmt.Sprintf("transaction failed: %v", err)}
|
||||
}
|
||||
|
||||
// 커밋 시도
|
||||
commitErr := tx.Commit()
|
||||
if commitErr != nil {
|
||||
// 커밋 중 에러가 1205일 수도 있음
|
||||
if isLockWaitTimeout(commitErr) {
|
||||
umlog.Info(value + ": Commit failed due to lock timeout, retrying... (" + strconv.Itoa(i+1) + "/" + strconv.Itoa(maxRetries) + ")")
|
||||
time.Sleep(time.Duration(100*(i+1)) * time.Millisecond)
|
||||
continue
|
||||
}
|
||||
umlog.Error(value + ": commit error: " + commitErr.Error())
|
||||
return model.DefaultErrorModel{Code: 4001201, Message: fmt.Sprintf("commit error: %v", commitErr)}
|
||||
}
|
||||
// 성공
|
||||
return model.DefaultErrorModel{Code: 200, Message: ""}
|
||||
}
|
||||
|
||||
umlog.Error(value + ": Transaction failed after " + strconv.Itoa(maxRetries) + " retries: " + err.Error())
|
||||
return model.DefaultErrorModel{Code: 4001201, Message: fmt.Sprintf("%s: Transaction failed after %d retries: %w", value, maxRetries, err)}
|
||||
}
|
||||
|
||||
func isLockWaitTimeout(err error) bool {
|
||||
// mysql 드라이버의 오류 메시지에서 1205 감지
|
||||
return err != nil && (strings.Contains(err.Error(), "Error 1205") || strings.Contains(err.Error(), "Lock wait timeout"))
|
||||
}
|
||||
Loading…
x
Reference in New Issue
Block a user