Skip to content

run_luxorASAP.py

Daemon que monitora a planilha historico_trades_ativos.xlsx, gera as tabelas Parquet e atualiza posições & caixa em tempo real.

  • Inicialização e carregamento de dados do LuxorDB.
  • Processamento de trades e atualização de posições de fundos.
  • Exportação de dados processados de volta para o LuxorDB.
  • Gerenciamento de logs e controle de execução.

O programa opera continuamente durante o horário de mercado, verificando periodicamente por novas transações e atualizando as posições dos fundos.

LuxorASAPManager

Source code in LuxorASAP/run_luxorASAP.py
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
class LuxorASAPManager:

    @logger.catch
    def __init__(self, reset=False, is_develop_mode=False):

        self.today = dt.date.today()

        self.__initialize(reset, is_develop_mode)

        self.source_bases_exec_count = 0


    def __reset(self):
        self._register_assets()
        self.save_funds()


    def __initialize(self, reset, is_develop_mode):

        self.is_develop_mode = is_develop_mode
        # Criando db manager, inicia carregando todos os dados
        self.lq = LuxorQuery(is_develop_mode=self.is_develop_mode)
        self.dl = DataLoader(Path().absolute()/"LuxorDB"/"tables")
        logger.info("LuxorDB: Dados Inicializados.")

        self.dl.add_file_tracker(Path().absolute()/"source_bases"/"historico_trades_ativos.xlsx", 
                                sheet_names={"Boletas": "trades", "Ativos": "assets", "VCs_prices":"hist_vc_px_last",
                                             "asset_field_map": "asset_field_map", "field_map" : "field_map"}, index=False,
                                             normalize_columns=True)
        self.dl.load_file_if_modified(Path().absolute()/"source_bases"/"historico_trades_ativos.xlsx", export_to_blob=True)

        # Carregando/criando fundos
        self.funds = {}

        if reset:
            self.__reset()
        else:
            self.funds = self._load_funds()


    def get_unique_funds(self):
        """ Retorna uma lista de ocorrencia unicas de fundos na base de boletas """
        unq_funds = list(self.lq.get_table("trades")["Fund"].unique()) 
        unq_funds.remove("lipizzaner(eduardo)")


        return unq_funds


    def get_unique_assets(self, fund_filter=None):
        """ Retorna lista de ocorrencias unicas de ativos. Se 'fund_filter' for informado, retorna as ocorrencias unicas nesse fundo."""
        trades = self.lq.get_table("trades").sort_values(by="Date", kind="stable")

        trades = self.lq.text_to_lowercase(trades)
        trades["FRGN_KEY"] = trades["Asset"] + "_" + trades["Ticker"]

        if fund_filter is None:
            return trades["FRGN_KEY"].unique()

        return trades.loc[ trades["Fund"] == fund_filter.lower(), "FRGN_KEY"].unique()


    def get_asset_properties(self, asset_key):
        """ Encontra as propriedades do ativo dado por 'asset_key'. Se nao for um ativo valido, retorna 'None'."""

        valid_assets = self.lq.get_table("assets") 
        properties = valid_assets.loc[ valid_assets["Key"] == asset_key ]

        if properties.empty:
            return properties

        return properties.iloc[0]


    def _load_funds(self):
        # Seleciona o arquivo mais recente e carrega os dados dos fundos

        files = os.listdir("funds_data")
        files.sort(reverse=True)
        savefile = os.path.join("funds_data", files[0])

        with open(savefile, "rb") as f:
            funds = pickle.load(f)

        return funds


    def save_funds(self):
        # salva os dados de cada fundo num pickle
        attempts = 5
        file_name = os.path.join("funds_data", "FundsData_"+str(self.today)+".pickle")
        while attempts > 0:
            try:
                with open(file_name, "wb") as f:
                    pickle.dump(self.funds, f, pickle.HIGHEST_PROTOCOL)
                    attempts = 0
            except PermissionError:
                attempts -=1
                logger.error(f"Erro ao tentar salvar FundsData_{str(self.today)}.pickle\nSerão feitas 5 tentativas...")
                time.sleep(5)

    def redirect_trades(self):
        # Filtra o df de trades e redireciona para processamento individual em cada fundo.
        trades = self.lq.get_table("trades")  
        for fund_name, fund in self.funds.items():
            trades_df = trades.loc[ trades["Fund"] == fund_name ]

            fund.set_trades_to_process(trades_df)


    # TODO: Nao precisa usar o db_manager. Pode simplesmente salvar dentro da luxorDB.
    def export_data_to_db(self, date_reference=None):
        """
        Insere as tabelas processadas (posicoes, posicoes_dia) para luxor_db.
        """
        if date_reference is None: date_reference=self.today
        last_positions = []
        hist_positions = []
        cash_movements = [] # vamos agrupar as movimentacoes de caixa aqui para exportar
        positions_by_bank = [] # vamos agrupar as posicoes por banco aqui para exportar
        last_positions_by_bank = [] # vamos agrupar as posicoes por banco aqui para exportar
        # Recuperamos a base de posicoes.. Mais recentes e historicas
        for fund_name, fund_data in self.funds.items():
            last_positions.append(fund_data.get_positions(all=False, date_reference=date_reference))
            hist_positions.append(fund_data.get_positions(all=True))
            cash_movements.append(fund_data.get_cash_movements())
            positions_by_bank.append(fund_data.get_positions_by_bank())
            last_positions_by_bank.append(fund_data.get_positions_by_bank(all=False, date_reference=date_reference))

        last_update = dt.datetime.now().timestamp() # forcar atualizacao

        # Concatenamos tudo numa tabela fato so
        last_positions = pd.concat(last_positions)
        # Exportamos as tabelas para luxorDB
        # sobrescrevendo as movimentacoes na base (modificacoes ou nao, vamos sobrescrever pois base pode estar desatualizada)
        logger.info(f"Exportando dados para a base de dados.")
        self.dl.load_table_if_modified("last_positions", last_positions, last_update=last_update, normalize_columns=True,
                                       export_to_blob=True)

        #O mesmo eh feito com o restante
        hist_positions = pd.concat(hist_positions)

        self.dl.load_table_if_modified("hist_positions", hist_positions, last_update=last_update, normalize_columns=True,
                                       export_to_blob=True)
        # Exportamos as tabelas para luxorDB
        cash_movements = pd.concat(cash_movements)
        self.dl.load_table_if_modified("cash_movements", cash_movements, last_update=last_update, normalize_columns=True,
                                       export_to_blob=True)

        positions_by_bank = pd.concat(positions_by_bank)
        self.dl.load_table_if_modified("hist_positions_by_bank", positions_by_bank, last_update=last_update, normalize_columns=True,
                                       export_to_blob=True)

        last_positions_by_bank = pd.concat(last_positions_by_bank)
        self.dl.load_table_if_modified("last_positions_by_bank", last_positions_by_bank, last_update=last_update, normalize_columns=True,
                                       export_to_blob=True)


    def _register_assets(self):
        # TODO PRECISAMOS SEPARAR POR CONTA BANCÁRIA!
        # Restagando nomes unicos dos fundos
        fund_names = self.get_unique_funds()
        for name in fund_names:     

            self.funds[name] = Fund(name)
            # Resgatando ativos unicos de cada fundo
            fund_assets = self.get_unique_assets(fund_filter=name)

            # Inicializamos cada ativo valido com suas respectivas props
            for asset_key in fund_assets:
                asset_props = self.get_asset_properties(asset_key)

                if asset_props.empty:
                    continue    # Ativo nao eh valido

                props = {
                "asset_type" : asset_props["Type"], "group" : asset_props["Group"], "ticker_bbg" : asset_props["Ticker_BBG"],
                "curncy_exp" : asset_props["Currency Exposure"], "location" : asset_props["Location"]
                }

                self.funds[name].register_asset(asset_props["Asset"], asset_props["Ticker"], props, self.lq)


    def _register_new_assets(self, fund_name):

        assets_list = self.get_unique_assets(fund_filter=fund_name)
        assets_already_registered = self.funds[fund_name].get_assets()

        for asset_key in assets_list:
            if asset_key not in assets_already_registered:
                # Entao eh um ativo novo que precisamos registrar no fundo
                asset_props = self.get_asset_properties(asset_key)

                if asset_props.empty:
                    continue    # Ativo nao eh valido

                props = {
                "asset_type" : asset_props["Type"], "group" : asset_props["Group"], "ticker_bbg" : asset_props["Ticker_BBG"],
                "curncy_exp" : asset_props["Currency Exposure"], "location" : asset_props["Location"]
                }

                self.funds[fund_name].register_asset(asset_props["Asset"], asset_props["Ticker"], props, self.lq)


    def process_trades(self, fund_names=None, date_reference=None):

        if date_reference is None: date_reference = self.today

        flag_modified = False
        flag_reset = False

        if fund_names is None: #Faz para todos os fundos cadastrados
            fund_names = list(self.funds.keys())
            # Nesse caso processa todos os fundos

        for f in fund_names:
            # Processa apenas para os fundos informados, se existirem
            if f in self.funds:

                # Primeiro criamos ativos novos (caso onde foram feitos trades com ativos novos)
                self._register_new_assets(fund_name=f)

                # Agora podemos processas os trades#
                _flag_modified, _flag_reset = self.funds[f].process_trades(self.lq)
                flag_modified = flag_modified | _flag_modified
                flag_reset = flag_reset | _flag_reset

        return flag_modified, flag_reset


    def stop(self):
        self.save_funds()
        #self.export_data_to_db()


    def update(self):
        # Faz diferenca olhar se eh fim de semana ou nao? Acho que nao mais...
        #if self.today != dt.date.today():
        #    tday= dt.date.today()
        #    if tday.weekday() <= 6:
        #        if not self.lq.is_holiday(tday, location="any"):  
        #            self.source_bases_exec_count = 0
        #            self.today = tday
        #        else: return
        #    else: return #final de semana, nao vamos rodar

        # Foi criado um script separado para atualizacao das source_bases
        #if self.source_bases_exec_count == 0:
        #    now = dt.datetime.now()
        #    if now.hour*100 + now.minute >= 50:
        #        source_bases_updater.update(is_develop_mode=self.is_develop_mode)
        #        self.source_bases_exec_count += 1

        # Processamos todos os trades
        self.redirect_trades()
        positions_changed, flag_reset = self.process_trades()

        if not flag_reset:
            if positions_changed:
                # Enviamos as tabelas resultantes para a base de dados 
                self.export_data_to_db()

            # Salvar pickle com informacao de processamento dos fundos
            self.save_funds()

            # Buscando alteracoes na historico_trades_ativos. Carregando caso haja.
            self.dl.scan_files(export_to_blob=True)

            # Atualizando tabelas carregadas.
            self.lq.update()

        else:
            logger.warning("RESET automático agendado para próxima exec.")
            self.__initialize(reset=flag_reset, is_develop_mode=self.is_develop_mode)

export_data_to_db(date_reference=None)

Insere as tabelas processadas (posicoes, posicoes_dia) para luxor_db.

Source code in LuxorASAP/run_luxorASAP.py
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
def export_data_to_db(self, date_reference=None):
    """
    Insere as tabelas processadas (posicoes, posicoes_dia) para luxor_db.
    """
    if date_reference is None: date_reference=self.today
    last_positions = []
    hist_positions = []
    cash_movements = [] # vamos agrupar as movimentacoes de caixa aqui para exportar
    positions_by_bank = [] # vamos agrupar as posicoes por banco aqui para exportar
    last_positions_by_bank = [] # vamos agrupar as posicoes por banco aqui para exportar
    # Recuperamos a base de posicoes.. Mais recentes e historicas
    for fund_name, fund_data in self.funds.items():
        last_positions.append(fund_data.get_positions(all=False, date_reference=date_reference))
        hist_positions.append(fund_data.get_positions(all=True))
        cash_movements.append(fund_data.get_cash_movements())
        positions_by_bank.append(fund_data.get_positions_by_bank())
        last_positions_by_bank.append(fund_data.get_positions_by_bank(all=False, date_reference=date_reference))

    last_update = dt.datetime.now().timestamp() # forcar atualizacao

    # Concatenamos tudo numa tabela fato so
    last_positions = pd.concat(last_positions)
    # Exportamos as tabelas para luxorDB
    # sobrescrevendo as movimentacoes na base (modificacoes ou nao, vamos sobrescrever pois base pode estar desatualizada)
    logger.info(f"Exportando dados para a base de dados.")
    self.dl.load_table_if_modified("last_positions", last_positions, last_update=last_update, normalize_columns=True,
                                   export_to_blob=True)

    #O mesmo eh feito com o restante
    hist_positions = pd.concat(hist_positions)

    self.dl.load_table_if_modified("hist_positions", hist_positions, last_update=last_update, normalize_columns=True,
                                   export_to_blob=True)
    # Exportamos as tabelas para luxorDB
    cash_movements = pd.concat(cash_movements)
    self.dl.load_table_if_modified("cash_movements", cash_movements, last_update=last_update, normalize_columns=True,
                                   export_to_blob=True)

    positions_by_bank = pd.concat(positions_by_bank)
    self.dl.load_table_if_modified("hist_positions_by_bank", positions_by_bank, last_update=last_update, normalize_columns=True,
                                   export_to_blob=True)

    last_positions_by_bank = pd.concat(last_positions_by_bank)
    self.dl.load_table_if_modified("last_positions_by_bank", last_positions_by_bank, last_update=last_update, normalize_columns=True,
                                   export_to_blob=True)

get_asset_properties(asset_key)

Encontra as propriedades do ativo dado por 'asset_key'. Se nao for um ativo valido, retorna 'None'.

Source code in LuxorASAP/run_luxorASAP.py
86
87
88
89
90
91
92
93
94
95
def get_asset_properties(self, asset_key):
    """ Encontra as propriedades do ativo dado por 'asset_key'. Se nao for um ativo valido, retorna 'None'."""

    valid_assets = self.lq.get_table("assets") 
    properties = valid_assets.loc[ valid_assets["Key"] == asset_key ]

    if properties.empty:
        return properties

    return properties.iloc[0]

get_unique_assets(fund_filter=None)

Retorna lista de ocorrencias unicas de ativos. Se 'fund_filter' for informado, retorna as ocorrencias unicas nesse fundo.

Source code in LuxorASAP/run_luxorASAP.py
73
74
75
76
77
78
79
80
81
82
83
def get_unique_assets(self, fund_filter=None):
    """ Retorna lista de ocorrencias unicas de ativos. Se 'fund_filter' for informado, retorna as ocorrencias unicas nesse fundo."""
    trades = self.lq.get_table("trades").sort_values(by="Date", kind="stable")

    trades = self.lq.text_to_lowercase(trades)
    trades["FRGN_KEY"] = trades["Asset"] + "_" + trades["Ticker"]

    if fund_filter is None:
        return trades["FRGN_KEY"].unique()

    return trades.loc[ trades["Fund"] == fund_filter.lower(), "FRGN_KEY"].unique()

get_unique_funds()

Retorna uma lista de ocorrencia unicas de fundos na base de boletas

Source code in LuxorASAP/run_luxorASAP.py
64
65
66
67
68
69
70
def get_unique_funds(self):
    """ Retorna uma lista de ocorrencia unicas de fundos na base de boletas """
    unq_funds = list(self.lq.get_table("trades")["Fund"].unique()) 
    unq_funds.remove("lipizzaner(eduardo)")


    return unq_funds