Загрузка данных


import builtins
import json
import os
import sys
import logging
import traceback
import requests
from datetime import datetime, timedelta
from calendar import monthrange
from argparse import ArgumentParser
from io import StringIO
import csv
import time

from umdap_config import LOG_FOLDER, ENV_FOLDER, EROUTER_HOST, PRIMARY_MADAB, UMDAP_ETL_URL_LIST

APPLICATION_NAME = 'externalGTT'
LOG_FILE = os.path.join(LOG_FOLDER, "externalGTT.log")
CONFIG_FILE = os.path.join(ENV_FOLDER, "externalGtt.json")
CREDENTIALS_FILE = os.path.join(os.path.dirname(os.path.abspath(__file__)), "config.json")

BASE_URL = 'https://www.globaltradetracker.com/api/rest'
builtins.BASE_URL = BASE_URL
builtins.XML_DATE_FMT = '%d.%m.%Y %H:%M:%S'
builtins.DELIMITER = ','

results = list()
last_update_file = os.path.join(os.path.dirname(os.path.abspath(__file__)), "last_update.json")


def log(msg=''):
    """Логирование с временной меткой"""
    str_with_time = f'{datetime.now().strftime("%Y/%m/%d %H:%M:%S")} {msg}'
    print(str_with_time)
    with open(LOG_FILE, "a", encoding="utf-8") as file:
        file.write(f'{str_with_time}\n')

builtins.log = log


def get_token():
    """Получение токена для API"""
    try:
        with open(CREDENTIALS_FILE, 'r') as f:
            credentials = json.load(f)
        
        link = f"{BASE_URL}/gettoken?userid={credentials['user_id']}&password={credentials['password']}"
        log(f"Getting token from GTT API")
        response = requests.get(link, timeout=30)
        
        if response.status_code == 200:
            token = response.content.decode().strip()
            if token:
                log("Token received successfully")
                return token
            else:
                log("Empty token received")
                return None
        else:
            log(f"Failed to get token. Status code: {response.status_code}")
            log(f"Response: {response.text}")
            return None
    except FileNotFoundError:
        log(f"ERROR: Credentials file not found: {CREDENTIALS_FILE}")
        return None
    except Exception as e:
        log(f"Error getting token: {str(e)}")
        return None


def get_subscription_countries(token):
    """Получение списка стран из подписки (только default reporters)"""
    try:
        link = f"{BASE_URL}/countries?token={token}"
        log(f"Getting available countries from subscription")
        response = requests.get(link, timeout=30)
        
        if response.status_code == 200:
            data = response.json()
            countries = []
            
            if isinstance(data, list):
                for item in data:
                    if isinstance(item, dict):
                        country_data = item.get('country', item)
                        if country_data.get('reportingcountry') == 'true' and country_data.get('defaultReporter') == True:
                            code = country_data.get('reportercode', '')
                            if code:
                                countries.append(code)
            
            countries = list(dict.fromkeys(countries))
            
            # ВРЕМЕННО: только 9 стран для теста
            test_countries = [
                'US_district',
                'BR_national',
                'IN_national',
                'CN_national',
                'DE_national',
                'RU_national',
                'JP_national',
                'KR_national',
                'TR_national'
            ]
            
            countries = [c for c in test_countries if c in countries]
            
            log(f"Test countries ({len(countries)}): {countries}")
            return countries
        else:
            log(f"ERROR: Failed to get countries. Status: {response.status_code}")
            return []
    except Exception as e:
        log(f"ERROR: Failed to get countries: {str(e)}")
        return []


def get_subscription_hs_codes(token):
    """Получение списка HS кодов из подписки (динамически)"""
    try:
        link = f"{BASE_URL}/subscriptions?token={token}"
        log(f"Getting available HS codes from subscription")
        response = requests.get(link, timeout=30)
        
        if response.status_code == 200:
            data = response.json()
            hs_codes = []
            
            if isinstance(data, list):
                for item in data:
                    if isinstance(item, dict):
                        codes = item.get('hsCodes', [])
                        if codes:
                            hs_codes.extend([str(code) for code in codes])
            
            # ВРЕМЕННО: только 1 HS-код для теста
            hs_codes = ['31']
            
            log(f"Test HS codes ({len(hs_codes)}): {hs_codes}")
            return hs_codes
        else:
            log(f"ERROR: Failed to get subscriptions. Status: {response.status_code}")
            return []
    except Exception as e:
        log(f"ERROR: Failed to get subscriptions: {str(e)}")
        return []


def get_last_update_time():
    """Чтение времени последнего обновления"""
    try:
        if os.path.exists(last_update_file):
            with open(last_update_file, 'r') as f:
                data = json.load(f)
                return data.get('last_update', None)
        return None
    except:
        return None


def save_last_update_time(update_time):
    """Сохранение времени последнего обновления"""
    try:
        with open(last_update_file, 'w') as f:
            json.dump({'last_update': update_time}, f)
        log(f"Last update time saved: {update_time}")
    except Exception as e:
        log(f"Error saving last update time: {str(e)}")


def check_data_updates(token):
    """Проверка наличия обновлений данных"""
    try:
        last_update = get_last_update_time()
        if not last_update:
            return True
        
        link = f"{BASE_URL}/dataupdates?token={token}&updatedAfter={last_update}"
        log(f"Checking for data updates since {last_update}")
        
        response = requests.get(link, timeout=30)
        if response.status_code == 200:
            data = response.json()
            if data and len(data) > 0:
                log(f"Found {len(data)} updates available")
                return True
            else:
                log("No updates found")
                return False
        else:
            log(f"Error checking updates. Status code: {response.status_code}")
            return False
    except Exception as e:
        log(f"Error checking updates: {str(e)}")
        return False


def parse_gtt_data(data, reporter_code="", trade_flow=""):
    """Парсинг данных из JSON ответа GTT API"""
    global results
    
    if not data:
        return
    
    if isinstance(data, list):
        for record in data:
            parse_single_record(record, reporter_code, trade_flow)
    elif isinstance(data, dict):
        if 'data' in data:
            parse_gtt_data(data['data'], reporter_code, trade_flow)
        elif 'records' in data:
            parse_gtt_data(data['records'], reporter_code, trade_flow)
        elif 'result' in data:
            parse_gtt_data(data['result'], reporter_code, trade_flow)
        elif 'items' in data:
            parse_gtt_data(data['items'], reporter_code, trade_flow)
        else:
            for key, value in data.items():
                if isinstance(value, list):
                    parse_gtt_data(value, reporter_code, trade_flow)
                    return
            parse_single_record(data, reporter_code, trade_flow)


def parse_single_record(record, reporter_code="", trade_flow=""):
    """Парсинг одной записи данных"""
    try:
        if not isinstance(record, dict):
            return
            
        def safe_get(obj, key, default=''):
            if isinstance(obj, dict):
                return obj.get(key, default)
            return default
        
        reporter = record.get('reporter', {})
        partner = record.get('partner', {})
        commodity = record.get('commodity', {})
        
        period_data = record.get('period', None)
        period_str = parse_period(period_data)
        
        value = record.get('value', record.get('monetaryValue', 0))
        if isinstance(value, dict):
            value = value.get('number', value.get('value', 0))
        
        quantity1 = record.get('quantity1', record.get('primaryQuantity', 0))
        if isinstance(quantity1, dict):
            quantity1 = quantity1.get('number', quantity1.get('value', 0))
        
        quantity2 = record.get('quantity2', record.get('secondaryQuantity', 0))
        if isinstance(quantity2, dict):
            quantity2 = quantity2.get('number', quantity2.get('value', 0))
        
        new_result = {
            'Reporter Code': safe_get(reporter, 'code', reporter_code) if isinstance(reporter, dict) else reporter_code,
            'Trade Flow': record.get('tradeFlow', record.get('flow', trade_flow)),
            'Is Mirror Data': record.get('isMirrorData', ''),
            'Period': period_str,
            'Reporter Name': safe_get(reporter, 'name', '') if isinstance(reporter, dict) else '',
            'Reporter Description': safe_get(reporter, 'description', '') if isinstance(reporter, dict) else '',
            'Reporter Source': safe_get(reporter, 'source', '') if isinstance(reporter, dict) else '',
            'Incoterm': record.get('incoterm', ''),
            'Partner Code': safe_get(partner, 'code', '') if isinstance(partner, dict) else '',
            'Partner Name': safe_get(partner, 'name', '') if isinstance(partner, dict) else '',
            'Commodity HS Code': safe_get(commodity, 'hsCode', record.get('hsCode', '')),
            'Commodity Description': safe_get(commodity, 'description', ''),
            'HS6 Code': record.get('hs6Code', safe_get(commodity, 'hs6Code', '')),
            'HS6 Code Description': record.get('hs6Description', ''),
            'Subdivision': record.get('subdivision', ''),
            'Port': record.get('port', ''),
            'Transport': record.get('transport', ''),
            'Foreign Port': record.get('foreignPort', ''),
            'US State': record.get('usState', ''),
            'Customs Regime': record.get('customsRegime', ''),
            'Suppression': record.get('suppression', ''),
            'Monetary Value': value,
            'Currency': record.get('currency', 'USD'),
            'Secondary Monetary Value': record.get('secondaryValue', record.get('secondaryMonetaryValue', 0)),
            'Secondary Currency': record.get('secondaryCurrency', ''),
            'Primary Quantity': quantity1,
            'Primary Quantity Unit': record.get('quantity1Unit', record.get('primaryQuantityUnit', '')),
            'Secondary Quantity': quantity2,
            'Secondary Quantity Unit': record.get('quantity2Unit', record.get('secondaryQuantityUnit', '')),
            'Primary Quantity Price': record.get('price1', record.get('primaryQuantityPrice', 0)),
            'Primary Quantity Price Unit': record.get('price1Unit', record.get('primaryQuantityPriceUnit', '')),
            'Secondary Quantity Price': record.get('price2', record.get('secondaryQuantityPrice', 0)),
            'Secondary Quantity Price Unit': record.get('price2Unit', record.get('secondaryQuantityPriceUnit', '')),
        }
        results.append(new_result)
    except Exception as e:
        log(f"Error parsing record: {str(e)}")
        log(f"Record data: {str(record)[:500]}")


def parse_period(period_data):
    """Парсинг периода в формат даты"""
    try:
        if isinstance(period_data, dict):
            year = period_data.get('year', period_data.get(0, 2000))
            month = period_data.get('month', period_data.get(1, 1))
        elif isinstance(period_data, list) and len(period_data) >= 2:
            year, month = period_data[0], period_data[1]
        elif isinstance(period_data, str):
            parts = period_data.split('-')
            if len(parts) >= 2:
                year, month = int(parts[0]), int(parts[1])
                if len(parts) == 3:
                    day = int(parts[2])
                    return datetime(year, month, day).strftime('%Y-%m-%d')
            else:
                return datetime(2000, 1, 1).strftime('%Y-%m-%d')
        else:
            return datetime(2000, 1, 1).strftime('%Y-%m-%d')
        
        last_day = monthrange(int(year), int(month))[1]
        return datetime(int(year), int(month), last_day).strftime('%Y-%m-%d')
    except Exception as e:
        log(f"Error parsing period {period_data}: {str(e)}")
        return datetime(2000, 1, 1).strftime('%Y-%m-%d')


def get_historical_data(token, countries, hs_codes):
    """БЛОК 1: Получение исторических данных только за 2026 год"""
    log("=" * 60)
    log("Starting TEST data collection for 2026")
    log(f"Countries ({len(countries)}): {countries}")
    log(f"HS Codes ({len(hs_codes)}): {hs_codes}")
    log("=" * 60)
    
    year = 2026
    current_month = datetime.now().month
    
    total_requests = 0
    successful_requests = 0
    empty_requests = 0
    error_requests = 0
    
    log(f"\n{'='*60}")
    log(f"Processing year: {year}")
    log(f"{'='*60}")
    
    for country in countries:
        for hscode in hs_codes:
            for impexp in ['E', 'I']:
                try:
                    from_date = f"{year}-01"
                    to_date = f"{year}-{current_month:02d}"
                    
                    link = f"{BASE_URL}/getreport"
                    params = {
                        'token': token,
                        'hscode': hscode,
                        'reporter': country,
                        'impexp': impexp,
                        'from': from_date,
                        'to': to_date,
                        'currency': 'USD',
                        'format': 'json'
                    }
                    
                    total_requests += 1
                    
                    response = requests.get(link, params=params, timeout=300)
                    
                    if response.status_code == 200:
                        successful_requests += 1
                        data = response.json()
                        
                        if data and isinstance(data, list) and len(data) > 0:
                            parse_gtt_data(data, country, impexp)
                            log(f"✓ [{total_requests}] Country={country}, HS={hscode}, Flow={impexp}, Records={len(data)}, Total={len(results)}")
                        else:
                            empty_requests += 1
                    else:
                        error_requests += 1
                        log(f"✗ [{total_requests}] Country={country}, HS={hscode}, Flow={impexp}, Error={response.status_code}")
                    
                except Exception as e:
                    log(f"✗ Exception: {str(e)}")
                    continue
    
    log(f"\n{'='*60}")
    log(f"COMPLETED")
    log(f"Total requests: {total_requests}")
    log(f"Successful: {successful_requests}")
    log(f"Empty: {empty_requests}")
    log(f"Errors: {error_requests}")
    log(f"Records collected: {len(results)}")
    log(f"{'='*60}")


def get_updated_data(token, countries, hs_codes):
    """БЛОК 2: Получение обновленных данных"""
    log("Checking for updated data")
    
    last_update = get_last_update_time()
    current_date = datetime.now().strftime('%Y-%m')
    
    if not check_data_updates(token):
        log("No updates available")
        return
    
    for country in countries:
        for hscode in hs_codes:
            for impexp in ['E', 'I']:
                try:
                    link = f"{BASE_URL}/getreport"
                    params = {
                        'token': token,
                        'hscode': hscode,
                        'reporter': country,
                        'impexp': impexp,
                        'updatedAfter': last_update if last_update else '2000-01',
                        'latestavailablemonths': 3,
                        'currency': 'USD',
                        'format': 'json'
                    }
                    
                    response = requests.get(link, params=params, timeout=300)
                    if response.status_code == 200:
                        data = response.json()
                        parse_gtt_data(data, country, impexp)
                    
                except Exception as e:
                    log(f"Error: {str(e)}")
                    continue
    
    save_last_update_time(current_date)
    log(f"Updated data collection completed. Total records: {len(results)}")


def gtt_request(url):
    """Формирование CSV для отправки в DWH"""
    global results
    
    log(f'{len(results)} import(s) generated')
    if results:
        with StringIO() as f:
            writer = csv.DictWriter(f, results[0].keys())
            writer.writeheader()
            writer.writerows(results)
            f.seek(0)
            return f.read()
    return ""


def main():
    """Основная функция"""
    global results
    
    log('--------------------- Script started ------------------------------')
    
    try:
        token = get_token()
        if not token:
            log("ERROR: Failed to get token")
            sys.exit(1)
        
        countries = get_subscription_countries(token)
        hs_codes = get_subscription_hs_codes(token)
        
        if not countries or not hs_codes:
            log("ERROR: Failed to get subscription parameters")
            sys.exit(1)
        
        last_update = get_last_update_time()
        
        if last_update is None:
            log("First run detected - starting TEST data collection")
            get_historical_data(token, countries, hs_codes)
            save_last_update_time(datetime.now().strftime('%Y-%m'))
        else:
            log(f"Regular update - last update was {last_update}")
            get_updated_data(token, countries, hs_codes)
        
        if results:
            log(f"Preparing to send {len(results)} records to DWH")
            from etl_common import process_etls
            process_etls(
                application_name=APPLICATION_NAME,
                etl_desc_url=f'{EROUTER_HOST}/procedure/{PRIMARY_MADAB}/umdap.etl.description/call.json?application={APPLICATION_NAME}',
                data_type=None,
                api_requestor=gtt_request,
                erouter_url_list=UMDAP_ETL_URL_LIST
            )
        else:
            log("No data to send")
        
    except Exception as e:
        log(f'Main thread unexpected error: {traceback.format_exception(*sys.exc_info())}')
        sys.exit(1)
    finally:
        log('--------------------- Script finished -----------------------------')


if __name__ == "__main__":
    main()