#include "StorageSendThread.h" #include "Logger.h" #include "ProcessConfig.h" #include #include #include #include #include #define SQLITE_DB_LOCK_TIMEOUT 1100 // 세션 중에 SQLITE_BUSY 상황 발생 문제 해결을 위해 DB Lock 관련 timeout 설정값. (1.1 sec) // 생성자.. CStorageSendThread::CStorageSendThread() { // Thread Handle 초기화 m_threadHandle = 0; } // 소멸자.. CStorageSendThread::~CStorageSendThread() { // Thread 가 동작 중인 경우 Thread 동작 정지 처리. if( m_threadHandle != 0 ) pthread_cancel(m_threadHandle); } // Thread 를 생성하여 작업을 시작한다... bool CStorageSendThread::Start() { int nResult = pthread_create( &m_threadHandle, NULL, CStorageSendThread::threadFunc, this ); if( nResult != 0 ) { // Thread 생성 실패시... int errorNum = errno; LOG( LERR, "StorageSendThread: Thread create failed.[%d][%s]", errorNum, strerror(errorNum) ); return false; } return true; } // Thread 함수... void * CStorageSendThread::threadFunc( void * arg ) { CStorageSendThread * pObject = reinterpret_cast(arg); pthread_detach( pthread_self() ); // 초기 기동시...통계 정보가 바로 DB 에 저장되지 않은 상태이므로... // 잠시 sleep 했다가.. 작업을 시작한다. sleep( 60 ); // 본 Thread 에서는 Local DB 에 저장된 서비스별 Storage 통계 정보를 // 주기적으로 select 하여... GTS cc_statd 모듈로 전송처리한다. // 실제 본 thread function 에서 loop 를 동작하는 것이 아니라.... // 객체 execute 함수 내에서 loop 로 동작하도록 한다. pthread_testcancel(); pObject->Execute(); pthread_testcancel(); // Thread 종료시 Handle 초기화 pObject->m_threadHandle = 0; return NULL; } // Local DB 에 저장된 서비스 별 Storage 통계 정보를 계산하여... 주기적으로 GTS cc_statd 에 전송 작업을 반복 수행한다... void CStorageSendThread::Execute() { // 변수 선언 및 초기화.. time_t currentTime; // 현재 시간 정보를 저장하기 위한 변수.. unsigned long nowTimestamp; unsigned long lastSendTimestamp = 0; // cc_statd 로 전송 시도한 마지막 timestamp 값 (5분단위) // Local DB 관련 변수 초기화. sqlite3 * db = NULL; sqlite3_stmt * stmt = NULL; int db_result; int nRetryCount; // 기준 timestamp 값을 조건으로.. cc_statd 로 전송되지 않은 서비스별 통계 정보를 모두 조회하여... 전송 처리한다. // cc_statd 로 전송 완료된 통계 정보는 send_flag 값이 1 로 설정, 전송 무시 처리인 경우 2로 설정됨. // CHG 2014-02-17 huibong // - 한번에 많은 양의 통계 정보를 GTS 로 전송할 경우.. GTS cc_statd 모듈의 메모리 사용량이 증가할 수 있으므로... // - 한번에 1000 개 통계정보만 전송되도록 수정한다. const char * select_sql = "SELECT timestamp, user_seq, svc_seq, quota, used " "FROM storage_stat " "WHERE timestamp <= ? AND send_flag = 0 " "ORDER BY timestamp ASC, user_seq ASC, svc_seq ASC " "LIMIT 1000 "; // Local DB 와 연결을 수행한다. db_result = sqlite3_open( CProcessConfig::GetInstance()->GetLocalDbFileName(), &db ); if( db_result != SQLITE_OK ) { LOG( LERR, "StorageSendThread: sqlite3 open failed.[%s][%d][%s]", CProcessConfig::GetInstance()->GetLocalDbFileName(), db_result, sqlite3_errmsg( db ) ); sqlite3_close(db); db = NULL; } else { // DB File 이 open 된 경우... nRetryCount = 0; // bind 관련 Query 를 수행할 stmt 객체 생성. // - 만약 sqlite DB File 이 사전에 생성되어 있지 않은 경우... 위의 open 함수에서 0 파일을 생성되어 정상적으로 open 처리되지만.. // sqlite3_prepare 함수에서 table 이 생성되어 있지 않으므로 오류가 발생함. // 따라서 loop 중에 재시도하도록 한다. db_result = sqlite3_prepare(db, select_sql, strlen(select_sql), &stmt, NULL ); // sqlite3_prepare() 함수 실행시 SQLITE_BUSY 오류가 발생할 수 있으므로.. 재시도 로직을 추가한다. while( db_result == SQLITE_BUSY && nRetryCount < 30 ) { // sqlite 에서 SQLITE_BUSY 반환 관련 lock 대기 timeout 을 설정. sqlite3_busy_timeout( db, SQLITE_DB_LOCK_TIMEOUT ); ++nRetryCount; LOG( LWAR, "StorageSendThread: sqlite3 prepare SQLITE_BUSY, Retry[%d]", nRetryCount ); db_result = sqlite3_prepare(db, select_sql, strlen(select_sql), &stmt, NULL ); } if( db_result != SQLITE_OK ) { LOG( LERR, "StorageSendThread: sqlite3 prepare failed.[%d][%s]", db_result, sqlite3_errmsg(db)); sqlite3_close(db); db = NULL; stmt = NULL; } } // Local DB 조회 정보를 저장하기 위한 구조체... struct gts_send_stat_storage stStorageStat; // Logging _LOG( LINF, "StorageSendThread: start..."); // 루프를 돌면서... Local DB 의 Storage 통계 정보를 조회하여... 주기적으로 cc_statd 에 전송한다. while( 1 ) { // cc_statd 로 전송을 수행할 것인지.. timestatmp 값을 계산하여... 확인한다. // 1. 현재 시간 정보 추출.. currentTime = time(0); // 2. 현재 timestamp 값을 계산한다... // 계산시.. 스토리지 사용량 계산 관련 시간이 오래 걸리수 있으므로.... // 240 sec 보정 처리하여 5분 단위 timestamp 값을 추출한다. nowTimestamp = (currentTime - 240 + 299)/300; // 3. 최초 기동 또는 시간 동기화 이상으로 reset 처리된 경우... lastSendTimestamp = 0 // 트래픽 통계 정보를 합산하여 바로 전송하면... timestamp 관련 중복 등의 문제가 발생할 수 있으므로... // 5분 후에 이전 통계 정보까지 전달될 수 있도록 한다. if( lastSendTimestamp == 0 ) { _LOG( LDBG, "StorageSendThread: first try... send skip..now[%lu] last[%lu]", nowTimestamp, lastSendTimestamp ); lastSendTimestamp = nowTimestamp; sleep(30); continue; } // 4. 마지막 timestamp 값과 현재 timestamp 계산 값이 같은 경우... // 이미 전송한 것으로 판단하고.. 쉰다. if( nowTimestamp == lastSendTimestamp ) { //_LOG( LDBG, "NetworkSendThread: send check.. skip.. now[%lu] last[%lu]", nowTimestamp, lastSendTimestamp ); // 10 sec 대기 후 timestamp 변경이 되었는지 다시 계산하도록 처리... // - 원래 코드는 1 마다 검사였으나.. // - 이럴 경우.. 순간적으로 여러 rc_sscd 에서 cc_statd 로 통계 정보가 일괄 전송 되어 // - cc_statd 에 부하가 발생될 확률이 높으므로.. 다음과 같이 10 sec 단위로 체크토록 변경한다. sleep(10); continue; } else { // 만약 lastSendTimestamp > nowTimestamp 인 경우... // - 시간 동기화 이상으로 이전에 계산된 timestamp 값이 잘못된 경우.... 보정처리.. if( nowTimestamp < lastSendTimestamp ) { LOG( LWAR, "StorageSendThread: Last timestamp not valid. skip and reset. now[%lu] last[%lu]", nowTimestamp, lastSendTimestamp ); lastSendTimestamp = 0; sleep(30); continue; } } _LOG( LDBG, "StorageSendThread: stats info select & send start... now[%lu] last[%lu]", nowTimestamp, lastSendTimestamp ); // 이제 전송작업을 수행할 것이지만.... 실제 전송 실패가 발생할 수도 있다... // 따라서... lastSendTimestamp 값을 변경하지 않을 경우.... while 루프에 의해 금방 재시도를 하므로.... // 아예 5분 지난 timestamp 변경 후.. 다시 시도하도록 lastSendTimestamp 값을 현재 값으로 변경 처리한다. lastSendTimestamp = nowTimestamp; // Local DB 조회 정보를 저장할 vector 초기화. m_vecStorageStat.clear(); // 이제 실제 전송 관련 처리 작업을 수행한다. // Local DB 연결이 되지 않은 경우... if( db == NULL ) { _LOG( LINF, "StorageSendThread: sqlite3 open retry.[%s]", CProcessConfig::GetInstance()->GetLocalDbFileName() ); db_result = sqlite3_open( CProcessConfig::GetInstance()->GetLocalDbFileName(), &db ); if( db_result != SQLITE_OK ) { LOG( LERR, "StorageSendThread: sqlite3 open failed2.[%s][%d][%s]" , CProcessConfig::GetInstance()->GetLocalDbFileName(), db_result, sqlite3_errmsg( db ) ); sqlite3_close(db); db = NULL; } else { nRetryCount = 0; // bind 관련 Query 를 수행할 stmt 객체 생성. db_result = sqlite3_prepare(db, select_sql, strlen(select_sql), &stmt, NULL ); // sqlite3_prepare() 함수 실행시 SQLITE_BUSY 오류가 발생할 수 있으므로.. 재시도 로직을 추가한다. while( db_result == SQLITE_BUSY && nRetryCount < 30 ) { // sqlite 에서 SQLITE_BUSY 반환 관련 lock 대기 timeout 을 설정. sqlite3_busy_timeout( db, SQLITE_DB_LOCK_TIMEOUT ); ++nRetryCount; LOG( LWAR, "StorageSendThread: sqlite3 prepare SQLITE_BUSY, Retry[%d]", nRetryCount ); db_result = sqlite3_prepare(db, select_sql, strlen(select_sql), &stmt, NULL ); } if( db_result != SQLITE_OK ) { LOG( LERR, "StorageSendThread: sqlite3 prepare failed2.[%d][%s]", db_result, sqlite3_errmsg(db)); sqlite3_close(db); db = NULL; stmt = NULL; } } } // sqlite 접속이 정상적으로 된 경우... if( stmt != NULL ) { // 매개변수 bind 처리. sqlite3_bind_int64( stmt, 1, nowTimestamp ); // 64bit 에서 unsigned long 변수는 8Byte 이므로.. // select query 수행. db_result = sqlite3_step(stmt); nRetryCount = 0; // 만약 BUSY 상태로 인해 오류가 발생하면..30번 재시도 while( db_result == SQLITE_BUSY && nRetryCount < 30 ) { sqlite3_busy_timeout( db, SQLITE_DB_LOCK_TIMEOUT ); ++nRetryCount; LOG( LWAR, "StorageSendThread: sqlite3 step SQLITE_BUSY, Retry[%d]",nRetryCount ); // query 재시도 db_result = sqlite3_step(stmt); } // select 조회 결과가 존재하는 경우... SQLITE_ROW (100 ) 을 반환... // select 조회 결과가 없는 경우.. SQLITE_DONE (101) 을 반환. // select query 수행시 오류가 발생한 경우... if( db_result != SQLITE_ROW && db_result != SQLITE_DONE ) { LOG( LERR, "StorageSendThread: sqlite3 step error.[%d][%s]", db_result, sqlite3_errmsg(db)); // Local DB 오류 발생시.. DB 에 문제가 있는 것으로 보고.. DB 객체를 초기화한다. sqlite3_finalize(stmt); sqlite3_close(db); db = NULL; stmt = NULL; } else { // select Query 가 정상적으로 수행된 경우.... while( db_result == SQLITE_ROW ) { // 임시 저장변수 초기화. memset( &stStorageStat, 0x00, sizeof(struct gts_send_stat_storage) ); // 조회 결과를 임시 변수에 저장.. stStorageStat.timestamp = sqlite3_column_int (stmt, 0); stStorageStat.user_seq = sqlite3_column_int (stmt, 1); stStorageStat.svc_seq = sqlite3_column_int (stmt, 2); stStorageStat.quota = sqlite3_column_int64(stmt, 3); stStorageStat.used = sqlite3_column_int64(stmt, 4); stStorageStat.send_flag = 0; _LOG( LDBG, "StorageSendThread: select result [%u][%d][%d] [%lu][%lu]" , stStorageStat.timestamp, stStorageStat.user_seq, stStorageStat.svc_seq , stStorageStat.quota, stStorageStat.used ); // vector 객체에 query 결과를 저장한다. m_vecStorageStat.push_back( stStorageStat ); // Next 결과 처리.. db_result = sqlite3_step( stmt ); nRetryCount = 0; // 만약 BUSY 상태로 인해 오류가 발생하면..60번 재시도 while( db_result == SQLITE_BUSY && nRetryCount < 60 ) { sqlite3_busy_timeout( db, SQLITE_DB_LOCK_TIMEOUT ); ++nRetryCount; LOG( LWAR, "StorageSendThread: sqlite3 step SQLITE_BUSY, Retry[%d]",nRetryCount ); // query 재시도 db_result = sqlite3_step(stmt); } } // Query 수행된 경우.. prepared 문을 재사용하기 위해 reset 처리한다. sqlite3_reset(stmt); } } // cc_statd 로 전송할 Data 가 존재하는 경우... if( m_vecStorageStat.empty() == false ) { // GTS 로 Data 전송을 담당하는 Socket 객체에 Data 를 전달하여... // 해당 Socket 객체에서 Data 를 전달하도록 처리한다. m_socketGts.SendStatStorage( m_vecStorageStat ); // 전송 처리 작업 호출 후.... // 전송 결과를 확인하여... 전송이 완료된 정보에 대해 Local DB 의 send_flag, send_time 정보를 변경처리한다. // - local DB 에서 조회 결과가 있어야만... m_vecStorageStat 의 값이 존재하므로. local db 연결 확인은 따로 하지 않아도 됨. char update_sql[1024]; char * db_err; std::vector::iterator it; for( it = m_vecStorageStat.begin(); it != m_vecStorageStat.end(); ++it ) { // GTS 로 전송 처리 또는 전송 무시 처리된 경우.. ( it->send_flag 가 1 or 2 인 경우.. ) if( it->send_flag != 0 ) { // GTS 로 전송 완료된 Local DB 항목에 대해 전송 완료 flag 값을 변경하기 위한 Query. snprintf( update_sql, 1023, "UPDATE storage_stat " "SET send_flag = %d, send_time = datetime('now', 'localtime') " "WHERE timestamp = %u AND user_seq = %d AND svc_seq = %d AND send_flag = 0 " , it->send_flag, it->timestamp, it->user_seq, it->svc_seq ); // Query 수행... db_result = sqlite3_exec(db, update_sql, NULL, NULL, &db_err); nRetryCount = 0; // 만약 BUSY 상태로 인해 오류가 발생하면.. 40번 재시도 while( db_result == SQLITE_BUSY && nRetryCount < 40 ) { sqlite3_free(db_err); sqlite3_busy_timeout( db, SQLITE_DB_LOCK_TIMEOUT ); ++nRetryCount; LOG( LWAR, "StorageSendThread: local db update sqlite3_exec SQLITE_BUSY. Retry [%d]",nRetryCount ); // Query 재수행 db_result = sqlite3_exec(db, update_sql, NULL, NULL, &db_err); } if( db_result != SQLITE_OK ) { LOG( LERR, "StorageSendThread: local db send_flag update fail.[%s][%s]", update_sql, db_err ); sqlite3_free(db_err); } _LOG( LINF, "Storage stats send to GTS OK [%u] [%d][%d] [%lu / %lu]" , it->timestamp, it->user_seq, it->svc_seq, it->used, it->quota ); } } } else { // 전송할 Data 가 없는 경우.. 로깅 처리.. LOG( LINF, "StorageSendThread: Storage stats data is empty. Not send to GTS .."); } // GTS cc_statd 로 통계 정보 전송 처리 후.. 불필요하게 loop 를 동작시킬 필요가 없으므로.... // 잠깐 sleep 한다. sleep(120); } // 종료 처리 코드 sqlite3_finalize(stmt); sqlite3_close(db); db = NULL; stmt = NULL; m_vecStorageStat.clear(); return; }