tps/1/scripts/subscriber-data-lost.py (59 lines)
1 import paho.mqtt.client as mqtt 2 3 # Configuración 4 broker = "broker.hivemq.com" 5 #broker = "mqtt-dashboard.com" 6 topic = "tp1/aguilar_klockner" # <<<<<<<<<<====== Completar con el nombre del grupo 7 output_file = 'output.txt' 8 received_fragments = {} 9 last_fragment = False 10 11 def on_subscribe(self, mqttc, obj, mid, granted_qos): 12 print("Subscribed: "+str(mid)+" "+str(granted_qos)) 13 14 def on_message(client, userdata, msg): 15 global last_fragment 16 17 18 # Decodificar mensaje: número de fragmento, tamaño, bandera de último fragmento, y contenido 19 payload = msg.payload.decode('utf-8') 20 fragment_info, fragment = payload.rsplit('|', 1) 21 fragment_number, fragment_size, is_last = map(int, fragment_info.split('|')[:3]) 22 23 received_fragments[fragment_number] = (fragment_size, fragment) 24 print(f"Fragmento recibido {fragment_number} (size: {fragment_size})") 25 if is_last == 1: 26 last_fragment = True 27 28 # Reensamblar si es el último fragmento 29 if last_fragment: 30 reassemble_file(output_file) 31 quit() 32 33 def reassemble_file(filename): 34 total_size = 0 #Creo la variable que guarda el peso total 35 control_size = 0 #Creo la variable que guarda el peso recibido 36 with open(filename, 'w') as file: 37 for fragment_number in sorted(received_fragments): 38 fragment_size, fragment = received_fragments[fragment_number] 39 if fragment_number != 0: #Excluyo de la reconstruccion al peso del archivo 40 file.write(fragment[:fragment_size]) # Reescribimos usando el largo correcto 41 total_size += fragment_size 42 else: # Guardo cuanto pesa el archivo como numero 43 control_size = int(fragment) 44 print(f"File reassembled as {filename}") 45 if control_size > total_size: #Contrasto el peso recibido con el esperado para definir si hubo error 46 print(f"Error en la transmision: Perdida de datos detectada") 47 else: 48 print(f"Transmision recibida exitosamente") 49 50 # Configuración del cliente MQTT 51 client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2) 52 client.on_message = on_message 53 client.on_subscribe = on_subscribe 54 55 client.connect(broker, 1883, 60) 56 client.subscribe(topic, qos=2) 57 58 # Mantener el cliente en funcionamiento 59 client.loop_forever()
