tps/1/scripts/publisher.py (53 lines)
1 import paho.mqtt.client as mqtt 2 import random 3 import time 4 5 # Configuración 6 broker_address = "broker.hivemq.com" 7 # broker_address = "mqtt-dashboard.com" 8 9 topic = "tp1/aguilar_klockner" 10 min_size = 50 # Tamaño mínimo del fragmento 11 max_size = 70 # Tamaño máximo del fragmento 12 file_to_publish = 'input.txt' 13 14 def on_connect(client, userdata, flags, rc): 15 # Al conectarse, configuramos la opción TCP_NODELAY 16 client_socket = client._socket().socket # Accede al socket subyacente 17 client_socket.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) # Desactiva Nagle 18 19 def on_subscribe(self, mqttc, obj, mid, granted_qos): 20 print("Subscribed: "+str(mid)+" "+str(granted_qos)) 21 22 def publish_file(client, filename, min_size, max_size): 23 with open(filename, 'r') as file: 24 content = file.read() 25 26 index = 0 27 fragment_number = 0 28 while index < len(content): 29 fragment_size = random.randint(min_size, max_size) 30 fragment = content[index:index+fragment_size] 31 32 # Metadatos: número de fragmento, tamaño, y bandera de último fragmento 33 is_last = 1 if index + fragment_size >= len(content) else 0 34 payload = f'{fragment_number}|{fragment_size}|{is_last}|{fragment}' 35 36 if fragment_number != 6: 37 client.publish(topic, payload, qos=2, retain=False) 38 print(f"Fragmento publicado {fragment_number} (size: {fragment_size})") 39 40 fragment_number += 1 41 index += fragment_size 42 time.sleep(1) 43 44 # Configuración del cliente MQTT 45 client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2) 46 client.on_connect = on_connect # Añade el manejador de eventos para cuando se conecte 47 client.on_subscribe = on_subscribe # Añade el manejador de suscripción 48 client.connect(broker_address, 1883, 60) 49 50 # Publicar el archivo fragmentado 51 publish_file(client, file_to_publish, min_size, max_size) 52 53 client.disconnect()
