Automated daily backup: 2026-10-04 03:21
This commit is contained in:
1 parent
417cbacd6b
commit
896be9ea75
1066 files changed
+179069
-79263
No files matched your search
@@ -522,7 +522,7 @@
|
||||
"11175": {
|
||||
"filename": "Dog stuck in woman pussy.mp4",
|
||||
"size": 5693733,
|
||||
"timestamp": "2026-09-22 03:47:04"
|
||||
"timestamp": "2026-10-03 16:58:40"
|
||||
},
|
||||
"11176": {
|
||||
"filename": "Doggy Knot.mp4",
|
||||
@@ -532,7 +532,7 @@
|
||||
"11177": {
|
||||
"filename": "Gay_beast_porn_for_this_dude_getting_nailed_by_his_hung_horse.mp4",
|
||||
"size": 3690584,
|
||||
"timestamp": "2026-09-22 03:47:05"
|
||||
"timestamp": "2026-10-03 16:58:40"
|
||||
},
|
||||
"11178": {
|
||||
"filename": "Milking The Dog - artofzoo.mp4",
|
||||
@@ -547,7 +547,7 @@
|
||||
"11180": {
|
||||
"filename": "Zoo_amateur_strips_naked_and_takes_her_dogs_cock_deep_until_he_unloads.mp4",
|
||||
"size": 29225685,
|
||||
"timestamp": "2026-09-22 03:47:23"
|
||||
"timestamp": "2026-10-03 16:58:43"
|
||||
},
|
||||
"11181": {
|
||||
"filename": "Big_fellow_makes_his_pet_get_down_and_fucks_her.mp4",
|
||||
@@ -652,7 +652,7 @@
|
||||
"11201": {
|
||||
"filename": "Teen girl and her dog alone part 2.mp4",
|
||||
"size": 6061574,
|
||||
"timestamp": "2026-09-22 03:48:08"
|
||||
"timestamp": "2026-10-03 16:58:44"
|
||||
},
|
||||
"11202": {
|
||||
"filename": "Back to back breeding.mp4",
|
||||
@@ -807,12 +807,12 @@
|
||||
"11232": {
|
||||
"filename": "60993d765d8b1.mp4",
|
||||
"size": 4549832,
|
||||
"timestamp": "2026-09-22 03:49:51"
|
||||
"timestamp": "2026-10-03 16:58:44"
|
||||
},
|
||||
"11233": {
|
||||
"filename": "606638a42e3f9.mp4",
|
||||
"size": 27099670,
|
||||
"timestamp": "2026-09-22 03:49:55"
|
||||
"timestamp": "2026-10-03 16:58:47"
|
||||
},
|
||||
"11234": {
|
||||
"filename": "6134d9df5605e.mp4",
|
||||
@@ -827,12 +827,12 @@
|
||||
"11236": {
|
||||
"filename": "61bc5ce80ff80.mp4",
|
||||
"size": 15336036,
|
||||
"timestamp": "2026-09-22 03:50:02"
|
||||
"timestamp": "2026-10-03 16:58:49"
|
||||
},
|
||||
"11237": {
|
||||
"filename": "60e3e4e098a9b.mp4",
|
||||
"size": 182135352,
|
||||
"timestamp": "2026-09-22 03:50:19"
|
||||
"timestamp": "2026-10-03 16:59:07"
|
||||
},
|
||||
"11238": {
|
||||
"filename": "full.mp4",
|
||||
|
||||
File diff suppressed because it is too large.
Load diff
@@ -192,7 +192,7 @@
|
||||
"12208": {
|
||||
"filename": "dogs-pounding-guys-comp-210831.mp4",
|
||||
"size": 43012394,
|
||||
"timestamp": "2026-09-21 23:41:52"
|
||||
"timestamp": "2026-10-03 16:57:59"
|
||||
},
|
||||
"12209": {
|
||||
"filename": "dog-wreaks-owners-ass-180252.mp4",
|
||||
@@ -747,7 +747,7 @@
|
||||
"12319": {
|
||||
"filename": "zoo_amateur_strips_naked_and_takes_her_dogs_cock_deep_until_he_unloads.mp4",
|
||||
"size": 29225685,
|
||||
"timestamp": "2026-09-21 23:52:32"
|
||||
"timestamp": "2026-10-03 16:58:02"
|
||||
},
|
||||
"12320": {
|
||||
"filename": "_zoo_porn_with_pig_giant_pig_unloads_his_huge_balls_inside_this.mp4",
|
||||
@@ -2692,7 +2692,7 @@
|
||||
"12708": {
|
||||
"filename": "toilet-teen-slave-drinks-pee-137377.mp4",
|
||||
"size": 11745200,
|
||||
"timestamp": "2026-09-22 00:20:17"
|
||||
"timestamp": "2026-10-03 16:58:03"
|
||||
},
|
||||
"12709": {
|
||||
"filename": "dog-porn-from-russian-with-teen-273778.mp4",
|
||||
@@ -2787,7 +2787,7 @@
|
||||
"12727": {
|
||||
"filename": "how-to-be-a-good-dog-owner-part-two-191252.mp4",
|
||||
"size": 50196150,
|
||||
"timestamp": "2026-09-22 00:21:55"
|
||||
"timestamp": "2026-10-03 16:58:10"
|
||||
},
|
||||
"12728": {
|
||||
"filename": "2-bitches-are-actively-having-sex-with-a-dog-192948.mp4",
|
||||
@@ -2837,7 +2837,7 @@
|
||||
"12737": {
|
||||
"filename": "wild4dogs-cum-dripping-246981.mp4",
|
||||
"size": 10103185,
|
||||
"timestamp": "2026-09-22 00:23:11"
|
||||
"timestamp": "2026-10-03 16:58:11"
|
||||
},
|
||||
"12738": {
|
||||
"filename": "wild4dogs-outdoor-session-248201.mp4",
|
||||
@@ -2862,27 +2862,27 @@
|
||||
"12742": {
|
||||
"filename": "snapchat-slut-desprate-for-anal-knot-190317.mp4",
|
||||
"size": 4975995,
|
||||
"timestamp": "2026-09-22 00:23:42"
|
||||
"timestamp": "2026-10-03 16:58:11"
|
||||
},
|
||||
"12743": {
|
||||
"filename": "cheating-wife-squirts-to-dog-cock-190310.mp4",
|
||||
"size": 6090261,
|
||||
"timestamp": "2026-09-22 00:23:43"
|
||||
"timestamp": "2026-10-03 16:58:12"
|
||||
},
|
||||
"12744": {
|
||||
"filename": "dog-becomes-addicted-to-human-pussy-191699.mp4",
|
||||
"size": 5283674,
|
||||
"timestamp": "2026-09-22 00:23:44"
|
||||
"timestamp": "2026-10-03 16:58:13"
|
||||
},
|
||||
"12745": {
|
||||
"filename": "how-to-be-a-good-dog-owner-part-one-191255.mp4",
|
||||
"size": 68877843,
|
||||
"timestamp": "2026-09-22 00:23:53"
|
||||
"timestamp": "2026-10-03 16:58:19"
|
||||
},
|
||||
"12746": {
|
||||
"filename": "how-to-be-a-good-dog-owner-part-three-191254.mp4",
|
||||
"size": 68228716,
|
||||
"timestamp": "2026-09-22 00:24:03"
|
||||
"timestamp": "2026-10-03 16:58:27"
|
||||
},
|
||||
"12747": {
|
||||
"filename": "2-women-have-sex-with-a-dog--200120.mp4",
|
||||
@@ -4962,7 +4962,7 @@
|
||||
"13162": {
|
||||
"filename": "lesbians-doing-fisting-on-the-street-206438.mp4",
|
||||
"size": 11180929,
|
||||
"timestamp": "2026-09-22 01:00:36"
|
||||
"timestamp": "2026-10-03 16:58:29"
|
||||
},
|
||||
"13163": {
|
||||
"filename": "lesbians-drink-milk-from-a-goat-195218.mp4",
|
||||
@@ -5537,7 +5537,7 @@
|
||||
"13277": {
|
||||
"filename": "sleeping_sister_gets_dick_inside_of_her_pussy_incest_sex_83645.mp4",
|
||||
"size": 17669095,
|
||||
"timestamp": "2026-09-22 01:12:34"
|
||||
"timestamp": "2026-10-03 16:58:30"
|
||||
},
|
||||
"13278": {
|
||||
"filename": "what-it-takes-to-love-a-dog-62886.mp4",
|
||||
@@ -5632,7 +5632,7 @@
|
||||
"13296": {
|
||||
"filename": "knotting-the-bitch-97087.mp4",
|
||||
"size": 3092672,
|
||||
"timestamp": "2026-09-22 01:14:25"
|
||||
"timestamp": "2026-10-03 16:58:31"
|
||||
},
|
||||
"13297": {
|
||||
"filename": "vixen---boar-to-the-core-83010.mp4",
|
||||
@@ -5847,32 +5847,32 @@
|
||||
"13339": {
|
||||
"filename": "19Yo First Time With Dog.mp4",
|
||||
"size": 23249368,
|
||||
"timestamp": "2026-09-22 01:17:25"
|
||||
"timestamp": "2026-10-03 19:32:21"
|
||||
},
|
||||
"13340": {
|
||||
"filename": "horse-fuck-pussy-girl-with-creampie-cum.mp4",
|
||||
"size": 4058338,
|
||||
"timestamp": "2026-09-22 01:17:25"
|
||||
"timestamp": "2026-10-03 19:32:23"
|
||||
},
|
||||
"13341": {
|
||||
"filename": "lunahotk9 rimming horse.mp4",
|
||||
"size": 4738177,
|
||||
"timestamp": "2026-09-22 01:17:26"
|
||||
"timestamp": "2026-10-03 19:32:24"
|
||||
},
|
||||
"13342": {
|
||||
"filename": "Scarlet Daze.mp4",
|
||||
"size": 1230596885,
|
||||
"timestamp": "2026-09-22 01:18:38"
|
||||
"timestamp": "2026-10-03 19:34:47"
|
||||
},
|
||||
"13343": {
|
||||
"filename": "Stuck in her ass.mp4",
|
||||
"size": 30935116,
|
||||
"timestamp": "2026-09-22 01:18:43"
|
||||
"timestamp": "2026-10-03 19:34:53"
|
||||
},
|
||||
"13344": {
|
||||
"filename": "[VHZone] Veronica Silesto - After Surgery.mp4",
|
||||
"size": 2220398123,
|
||||
"timestamp": "2026-09-22 01:20:59"
|
||||
"timestamp": "2026-10-03 19:39:05"
|
||||
},
|
||||
"13345": {
|
||||
"filename": "VHZone_Veronica_Silesto_Babe_Is_Fucked_Good_By_Her_Dog_Lover.mp4",
|
||||
@@ -5882,16 +5882,16 @@
|
||||
"13346": {
|
||||
"filename": "[VHZone] Veronica Silesto - Bedroom Party.mp4",
|
||||
"size": 3078536403,
|
||||
"timestamp": "2026-09-22 01:24:16"
|
||||
"timestamp": "2026-10-03 19:44:51"
|
||||
},
|
||||
"13347": {
|
||||
"filename": "[VHZone] Veronica Silesto - Blue Party.mp4",
|
||||
"size": 1357851731,
|
||||
"timestamp": "2026-09-22 01:25:34"
|
||||
"timestamp": "2026-10-03 19:47:15"
|
||||
},
|
||||
"13348": {
|
||||
"filename": "[VHZone] Veronica Silesto - Close Up Knot Compilation.mp4",
|
||||
"size": 200015250,
|
||||
"timestamp": "2026-09-22 01:25:54"
|
||||
"timestamp": "2026-10-03 19:47:44"
|
||||
}
|
||||
}
|
||||
@@ -6,22 +6,31 @@ import math
|
||||
import asyncio
|
||||
import inspect
|
||||
import logging
|
||||
|
||||
# Ensure UTF-8 output and enable Windows VT100 / ANSI escape sequences
|
||||
if sys.platform == 'win32':
|
||||
os.system('')
|
||||
try:
|
||||
sys.stdout.reconfigure(encoding='utf-8', errors='replace')
|
||||
sys.stderr.reconfigure(encoding='utf-8', errors='replace')
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
from telethon import TelegramClient, errors, utils
|
||||
from telethon.tl.functions.auth import ExportAuthorizationRequest, ImportAuthorizationRequest
|
||||
from telethon.tl.functions import InvokeWithLayerRequest
|
||||
from telethon.tl.functions.upload import GetFileRequest
|
||||
from telethon.tl.types import MessageMediaPhoto, MessageMediaDocument
|
||||
from telethon.network import MTProtoSender
|
||||
from telethon.tl.alltlobjects import LAYER
|
||||
from FastTelethonhelper.FastTelethon import DownloadSender
|
||||
|
||||
# Suppress Telethon's internal network warnings (e.g., transient server-closed socket resets that auto-reconnect)
|
||||
# Suppress Telethon's internal network warnings
|
||||
logging.basicConfig(level=logging.ERROR)
|
||||
logging.getLogger('telethon').setLevel(logging.ERROR)
|
||||
logging.getLogger('asyncio').setLevel(logging.ERROR)
|
||||
|
||||
# Global cache for persistent downloaders per DC ID to prevent connection churn
|
||||
downloaders = {}
|
||||
downloader_lock = asyncio.Lock()
|
||||
# Maximum chunk size supported by Telegram MTProto upload.getFile is 512 KB
|
||||
CHUNK_SIZE = 512 * 1024
|
||||
|
||||
# Helper function to format sizes
|
||||
def format_size(bytes_count):
|
||||
@@ -55,15 +64,17 @@ class ConcurrentProgressRenderer:
|
||||
if self.slots[i] == "":
|
||||
self.slots[i] = "Initializing..."
|
||||
return i
|
||||
return -1
|
||||
return 0
|
||||
|
||||
async def update_slot(self, slot_idx, filename, received, total, start_time):
|
||||
def update_slot(self, slot_idx, filename, received, total, start_time):
|
||||
if not (0 <= slot_idx < len(self.slots)):
|
||||
return
|
||||
if not total:
|
||||
total = 1
|
||||
percent = (received / total) * 100
|
||||
bar_len = 15
|
||||
filled_len = int(bar_len * received // total)
|
||||
bar = '█' * filled_len + '░' * (bar_len - filled_len)
|
||||
filled_len = int(bar_len * received // total) if total > 0 else 0
|
||||
bar = '█' * min(filled_len, bar_len) + '░' * max(0, bar_len - filled_len)
|
||||
|
||||
now = time.time()
|
||||
elapsed = now - start_time
|
||||
@@ -75,12 +86,12 @@ class ConcurrentProgressRenderer:
|
||||
if len(display_name) > 20:
|
||||
display_name = display_name[:9] + "..." + display_name[-8:]
|
||||
|
||||
async with self.lock:
|
||||
self.slots[slot_idx] = f"Slot {slot_idx+1}: {display_name} [{bar}] {percent:.1f}% ({size_str}) @ {speed_str}"
|
||||
self.slots[slot_idx] = f"Slot {slot_idx+1}: {display_name} [{bar}] {percent:.1f}% ({size_str}) @ {speed_str}"
|
||||
|
||||
async def release_slot(self, slot_idx):
|
||||
async with self.lock:
|
||||
self.slots[slot_idx] = ""
|
||||
if 0 <= slot_idx < len(self.slots):
|
||||
self.slots[slot_idx] = ""
|
||||
|
||||
async def add_bytes(self, size):
|
||||
async with self.lock:
|
||||
@@ -98,8 +109,7 @@ class ConcurrentProgressRenderer:
|
||||
async with self.lock:
|
||||
if self.initialized:
|
||||
num_lines = 2 + len(self.slots)
|
||||
sys.stdout.write(f"\033[{num_lines}A")
|
||||
sys.stdout.write("\033[J")
|
||||
sys.stdout.write(f"\033[{num_lines}A\033[J")
|
||||
sys.stdout.write(message + "\n")
|
||||
sys.stdout.flush()
|
||||
self.initialized = False
|
||||
@@ -112,7 +122,7 @@ class ConcurrentProgressRenderer:
|
||||
percent = (self.completed_files / self.total_files * 100) if self.total_files > 0 else 0.0
|
||||
bar_len = 50
|
||||
filled_len = int(bar_len * self.completed_files // self.total_files) if self.total_files > 0 else 0
|
||||
bar = '█' * filled_len + ' ' * (bar_len - filled_len)
|
||||
bar = '█' * min(filled_len, bar_len) + ' ' * max(0, bar_len - filled_len)
|
||||
|
||||
lines = []
|
||||
lines.append(
|
||||
@@ -135,147 +145,226 @@ class ConcurrentProgressRenderer:
|
||||
sys.stdout.flush()
|
||||
|
||||
async def ui_loop(renderer):
|
||||
if sys.platform == 'win32':
|
||||
import ctypes
|
||||
kernel32 = ctypes.windll.kernel32
|
||||
kernel32.SetConsoleMode(kernel32.GetStdHandle(-11), 7)
|
||||
|
||||
while True:
|
||||
await renderer.render()
|
||||
await asyncio.sleep(0.2)
|
||||
try:
|
||||
await renderer.render()
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception:
|
||||
pass
|
||||
await asyncio.sleep(0.25)
|
||||
|
||||
# Persistent parallel connection downloader to avoid connection churn
|
||||
class PersistentParallelDownloader:
|
||||
def __init__(self, client, dc_id, connection_count):
|
||||
# Thread-safe persistent MTProto sender pool per Telegram Data Center (DC)
|
||||
class SenderPool:
|
||||
def __init__(self, client, dc_id, max_senders=20):
|
||||
self.client = client
|
||||
self.dc_id = dc_id
|
||||
self.connection_count = connection_count
|
||||
self.senders = []
|
||||
self.auth_key = None
|
||||
|
||||
async def initialize(self):
|
||||
self.max_senders = max_senders
|
||||
self.auth_key = (
|
||||
None
|
||||
if self.dc_id and self.client.session.dc_id != self.dc_id
|
||||
else self.client.session.auth_key
|
||||
if dc_id and client.session.dc_id != dc_id
|
||||
else client.session.auth_key
|
||||
)
|
||||
|
||||
# Connect persistent MTProtoSenders sequentially to prevent session ID conflicts
|
||||
for i in range(self.connection_count):
|
||||
dc = await self.client._get_dc(self.dc_id)
|
||||
sender = MTProtoSender(self.auth_key, loggers=self.client._log, retries=10, delay=1, auto_reconnect=True)
|
||||
await sender.connect(
|
||||
self.client._connection(
|
||||
dc.ip_address,
|
||||
dc.port,
|
||||
dc.id,
|
||||
loggers=self.client._log,
|
||||
proxy=self.client._proxy,
|
||||
)
|
||||
self.available = []
|
||||
self.total_created = 0
|
||||
self.condition = asyncio.Condition()
|
||||
|
||||
async def _create_sender(self):
|
||||
dc = await self.client._get_dc(self.dc_id)
|
||||
sender = MTProtoSender(self.auth_key, loggers=self.client._log, retries=10, delay=1, auto_reconnect=True)
|
||||
await sender.connect(
|
||||
self.client._connection(
|
||||
dc.ip_address,
|
||||
dc.port,
|
||||
dc.id,
|
||||
loggers=self.client._log,
|
||||
proxy=self.client._proxy,
|
||||
)
|
||||
if not self.auth_key:
|
||||
auth = await self.client(ExportAuthorizationRequest(self.dc_id))
|
||||
self.client._init_request.query = ImportAuthorizationRequest(
|
||||
id=auth.id, bytes=auth.bytes
|
||||
)
|
||||
req = InvokeWithLayerRequest(LAYER, self.client._init_request)
|
||||
await sender.send(req)
|
||||
self.auth_key = sender.auth_key
|
||||
elif i > 0 and not sender.auth_key:
|
||||
sender.auth_key = self.auth_key
|
||||
|
||||
self.senders.append(sender)
|
||||
|
||||
async def download_file(self, input_file_location, size, out, progress_callback=None):
|
||||
part_size_kb = utils.get_appropriated_part_size(size)
|
||||
part_size = part_size_kb * 1024
|
||||
part_count = math.ceil(size / part_size)
|
||||
|
||||
connections = self.connection_count
|
||||
minimum, remainder = divmod(part_count, connections)
|
||||
|
||||
def get_part_count():
|
||||
nonlocal remainder
|
||||
if remainder > 0:
|
||||
remainder -= 1
|
||||
return minimum + 1
|
||||
return minimum
|
||||
|
||||
download_senders = []
|
||||
for i in range(connections):
|
||||
ds = DownloadSender(
|
||||
self.client,
|
||||
self.senders[i],
|
||||
input_file_location, # Pass the correct InputFileLocation subclass, not the Document TLObject
|
||||
offset=i * part_size,
|
||||
limit=part_size,
|
||||
stride=connections * part_size,
|
||||
count=get_part_count()
|
||||
)
|
||||
if not self.auth_key:
|
||||
auth = await self.client(ExportAuthorizationRequest(self.dc_id))
|
||||
self.client._init_request.query = ImportAuthorizationRequest(
|
||||
id=auth.id, bytes=auth.bytes
|
||||
)
|
||||
download_senders.append(ds)
|
||||
req = InvokeWithLayerRequest(LAYER, self.client._init_request)
|
||||
await sender.send(req)
|
||||
self.auth_key = sender.auth_key
|
||||
elif not sender.auth_key:
|
||||
sender.auth_key = self.auth_key
|
||||
return sender
|
||||
|
||||
part = 0
|
||||
while part < part_count:
|
||||
tasks = []
|
||||
for ds in download_senders:
|
||||
tasks.append(self.client.loop.create_task(ds.next()))
|
||||
try:
|
||||
for task in tasks:
|
||||
data = await task
|
||||
if not data:
|
||||
break
|
||||
out.write(data)
|
||||
part += 1
|
||||
if progress_callback:
|
||||
r = progress_callback(out.tell(), size)
|
||||
if inspect.isawaitable(r):
|
||||
await r
|
||||
except Exception:
|
||||
for task in tasks:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
raise
|
||||
async def acquire_senders(self, count):
|
||||
async with self.condition:
|
||||
senders = []
|
||||
while len(senders) < count:
|
||||
# 1. Drain alive senders from available pool
|
||||
while self.available and len(senders) < count:
|
||||
s = self.available.pop()
|
||||
if s.is_connected():
|
||||
senders.append(s)
|
||||
else:
|
||||
try:
|
||||
await s.disconnect()
|
||||
except Exception:
|
||||
pass
|
||||
self.total_created = max(0, self.total_created - 1)
|
||||
|
||||
if len(senders) == count:
|
||||
break
|
||||
|
||||
# 2. Create new connections up to max_senders
|
||||
if self.total_created < self.max_senders:
|
||||
needed = min(count - len(senders), self.max_senders - self.total_created)
|
||||
for _ in range(needed):
|
||||
try:
|
||||
s = await self._create_sender()
|
||||
self.total_created += 1
|
||||
senders.append(s)
|
||||
except Exception:
|
||||
break
|
||||
|
||||
# If we have at least 1 sender, proceed without blocking
|
||||
if senders:
|
||||
break
|
||||
|
||||
# If 0 senders available, wait for another task to release
|
||||
await self.condition.wait()
|
||||
|
||||
return senders
|
||||
|
||||
async def release_senders(self, senders):
|
||||
async with self.condition:
|
||||
for s in senders:
|
||||
if s and s.is_connected():
|
||||
self.available.append(s)
|
||||
else:
|
||||
if s:
|
||||
try:
|
||||
await s.disconnect()
|
||||
except Exception:
|
||||
pass
|
||||
self.total_created = max(0, self.total_created - 1)
|
||||
self.condition.notify_all()
|
||||
|
||||
async def disconnect_all(self):
|
||||
await asyncio.gather(*[sender.disconnect() for sender in self.senders if sender], return_exceptions=True)
|
||||
self.senders = []
|
||||
async with self.condition:
|
||||
all_senders = list(self.available)
|
||||
self.available.clear()
|
||||
for s in all_senders:
|
||||
try:
|
||||
await s.disconnect()
|
||||
except Exception:
|
||||
pass
|
||||
self.total_created = 0
|
||||
|
||||
async def get_downloader(client, dc_id, connection_count):
|
||||
async with downloader_lock:
|
||||
if dc_id in downloaders:
|
||||
dl = downloaders[dc_id]
|
||||
# Ensure existing downloaders have all active senders matching connection count
|
||||
all_alive = len(dl.senders) == connection_count and all(s.is_connected() for s in dl.senders)
|
||||
if not all_alive:
|
||||
await dl.disconnect_all()
|
||||
del downloaders[dc_id]
|
||||
if dc_id not in downloaders:
|
||||
dl = PersistentParallelDownloader(client, dc_id, connection_count)
|
||||
await dl.initialize()
|
||||
downloaders[dc_id] = dl
|
||||
return downloaders[dc_id]
|
||||
dc_pools = {}
|
||||
pool_lock = asyncio.Lock()
|
||||
|
||||
async def download_media_fast(client, chat, message, filepath, callback, connection_count=8):
|
||||
async def get_dc_pool(client, dc_id):
|
||||
async with pool_lock:
|
||||
if dc_id not in dc_pools:
|
||||
dc_pools[dc_id] = SenderPool(client, dc_id)
|
||||
return dc_pools[dc_id]
|
||||
|
||||
# High-throughput asynchronous pipelined downloader with sequential streaming disk writes
|
||||
async def download_file_parallel(pool, client, input_file_location, size, filepath, callback, desired_connections=4):
|
||||
part_count = math.ceil(size / CHUNK_SIZE) if size > 0 else 0
|
||||
if part_count == 0:
|
||||
with open(filepath, "wb") as f:
|
||||
pass
|
||||
return
|
||||
|
||||
num_senders = min(desired_connections, part_count)
|
||||
senders = await pool.acquire_senders(num_senders)
|
||||
|
||||
try:
|
||||
queue = asyncio.Queue()
|
||||
for i in range(part_count):
|
||||
queue.put_nowait((i, i * CHUNK_SIZE))
|
||||
|
||||
current_part = 0
|
||||
buffer = {}
|
||||
bytes_written = 0
|
||||
write_lock = asyncio.Lock()
|
||||
stop_event = asyncio.Event()
|
||||
|
||||
# Open file in write-binary mode directly (no slow pre-allocation or seek stalls)
|
||||
with open(filepath, "wb") as f:
|
||||
async def handle_chunk(part_idx, data):
|
||||
nonlocal current_part, bytes_written
|
||||
async with write_lock:
|
||||
buffer[part_idx] = data
|
||||
while current_part in buffer:
|
||||
chunk = buffer.pop(current_part)
|
||||
f.write(chunk)
|
||||
bytes_written += len(chunk)
|
||||
current_part += 1
|
||||
if callback:
|
||||
callback(bytes_written, size)
|
||||
|
||||
async def worker(sender):
|
||||
while not queue.empty() and not stop_event.is_set():
|
||||
try:
|
||||
part_idx, offset = queue.get_nowait()
|
||||
except asyncio.QueueEmpty:
|
||||
break
|
||||
|
||||
data = None
|
||||
last_err = None
|
||||
for attempt in range(4):
|
||||
if stop_event.is_set():
|
||||
break
|
||||
try:
|
||||
req = GetFileRequest(input_file_location, offset=offset, limit=CHUNK_SIZE)
|
||||
res = await client._call(sender, req)
|
||||
data = res.bytes
|
||||
break
|
||||
except errors.FloodWaitError as e:
|
||||
await asyncio.sleep(e.seconds)
|
||||
except errors.FileReferenceExpiredError:
|
||||
stop_event.set()
|
||||
raise
|
||||
except Exception as e:
|
||||
last_err = e
|
||||
await asyncio.sleep(0.5 * (attempt + 1))
|
||||
|
||||
if data is None and not stop_event.is_set():
|
||||
stop_event.set()
|
||||
if last_err:
|
||||
raise last_err
|
||||
raise IOError(f"Failed to fetch chunk at offset {offset}")
|
||||
|
||||
if data:
|
||||
await handle_chunk(part_idx, data)
|
||||
|
||||
queue.task_done()
|
||||
|
||||
await asyncio.gather(*[worker(s) for s in senders])
|
||||
f.flush()
|
||||
finally:
|
||||
await pool.release_senders(senders)
|
||||
|
||||
async def download_media_fast(client, chat, message, filepath, callback, connection_count=4):
|
||||
if message.document:
|
||||
dc_id, input_file_location = utils.get_input_location(message.document)
|
||||
dl = await get_downloader(client, dc_id, connection_count)
|
||||
with open(filepath, "wb") as f:
|
||||
await dl.download_file(input_file_location, message.document.size, f, progress_callback=callback)
|
||||
pool = await get_dc_pool(client, dc_id)
|
||||
await download_file_parallel(pool, client, input_file_location, message.document.size, filepath, callback, desired_connections=connection_count)
|
||||
else:
|
||||
# Photos are small enough that native download works perfectly
|
||||
# Photos are small enough that native download works smoothly
|
||||
await client.download_media(message, file=filepath, progress_callback=callback)
|
||||
|
||||
async def download_file_task(client, chat, message, filepath, filename, file_size, index,
|
||||
renderer, semaphore, history_file, downloaded_history, history_lock, connection_count=8):
|
||||
renderer, semaphore, history_file, downloaded_history, history_lock, connection_count=4):
|
||||
msg_id_str = str(message.id)
|
||||
slot_idx = await renderer.get_free_slot()
|
||||
|
||||
|
||||
async with semaphore:
|
||||
slot_idx = await renderer.get_free_slot()
|
||||
start_time = time.time()
|
||||
# Immediately display slot info and starting 0% progress bar
|
||||
renderer.update_slot(slot_idx, filename, 0, file_size, start_time)
|
||||
|
||||
def callback(received, total):
|
||||
asyncio.create_task(renderer.update_slot(slot_idx, filename, received, total or file_size, start_time))
|
||||
renderer.update_slot(slot_idx, filename, received, total or file_size, start_time)
|
||||
|
||||
success = False
|
||||
retries = 0
|
||||
@@ -306,15 +395,6 @@ async def download_file_task(client, chat, message, filepath, filename, file_siz
|
||||
await asyncio.sleep(e.seconds)
|
||||
except (errors.RPCError, asyncio.TimeoutError, ConnectionError, Exception):
|
||||
retries += 1
|
||||
try:
|
||||
if message.document:
|
||||
dc_id, _ = utils.get_input_location(message.document)
|
||||
async with downloader_lock:
|
||||
dl = downloaders.pop(dc_id, None)
|
||||
if dl:
|
||||
await dl.disconnect_all()
|
||||
except Exception:
|
||||
pass
|
||||
if retries >= max_retries:
|
||||
break
|
||||
try:
|
||||
@@ -322,7 +402,7 @@ async def download_file_task(client, chat, message, filepath, filename, file_siz
|
||||
os.truncate(filepath, 0)
|
||||
except Exception:
|
||||
pass
|
||||
await asyncio.sleep(retries * 3)
|
||||
await asyncio.sleep(retries * 2)
|
||||
|
||||
actual_size = os.path.getsize(filepath) if success and os.path.exists(filepath) else 0
|
||||
await renderer.release_slot(slot_idx)
|
||||
@@ -357,8 +437,8 @@ async def main():
|
||||
except ValueError:
|
||||
chat_username = sys.argv[4]
|
||||
output_dir = sys.argv[5] if len(sys.argv) > 5 else "telegram_downloads"
|
||||
concurrency = int(sys.argv[6]) if len(sys.argv) > 6 else 1
|
||||
connections_per_file = int(sys.argv[7]) if len(sys.argv) > 7 else 8
|
||||
concurrency = int(sys.argv[6]) if len(sys.argv) > 6 else 3
|
||||
connections_per_file = int(sys.argv[7]) if len(sys.argv) > 7 else 4
|
||||
|
||||
os.makedirs(output_dir, exist_ok=True)
|
||||
history_file = os.path.join(output_dir, "download_history.json")
|
||||
@@ -370,6 +450,16 @@ async def main():
|
||||
downloaded_history = json.load(f)
|
||||
except Exception as e:
|
||||
print(f"Warning: Failed to load download history: {e}")
|
||||
|
||||
# Fast directory scan to cache local files in memory
|
||||
existing_files = {}
|
||||
try:
|
||||
with os.scandir(output_dir) as it:
|
||||
for entry in it:
|
||||
if entry.is_file():
|
||||
existing_files[entry.name] = entry.stat().st_size
|
||||
except Exception as e:
|
||||
print(f"Warning scanning output directory: {e}")
|
||||
|
||||
client = TelegramClient('session_dumdum', api_id, api_hash)
|
||||
|
||||
@@ -388,7 +478,7 @@ async def main():
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# If entity not found in session cache (common for deleted accounts or raw user IDs), iterate dialogs to find entity and access hash
|
||||
# If entity not found in session cache, search dialogs
|
||||
if not chat:
|
||||
print(f"Direct get_entity failed ({e}). Searching dialogs for ID {chat_username}...")
|
||||
target_id = chat_username if isinstance(chat_username, int) else None
|
||||
@@ -458,9 +548,8 @@ async def main():
|
||||
assigned_filenames[filename] = msg_id_str
|
||||
filepath = os.path.join(output_dir, filename)
|
||||
|
||||
if os.path.exists(filepath):
|
||||
local_size = os.path.getsize(filepath)
|
||||
|
||||
local_size = existing_files.get(filename)
|
||||
if local_size is not None:
|
||||
if file_size > 0 and local_size == file_size:
|
||||
print(f"[SKIP] '{filename}' already exists and is complete ({format_size(local_size)}).")
|
||||
if msg_id_str not in downloaded_history:
|
||||
@@ -483,7 +572,7 @@ async def main():
|
||||
)
|
||||
tasks_to_run.append(task)
|
||||
|
||||
print("\nStarting downloads...")
|
||||
print(f"\nStarting downloads with {concurrency} concurrent files and {connections_per_file} streams per file...")
|
||||
ui_task = asyncio.create_task(ui_loop(renderer))
|
||||
|
||||
if tasks_to_run:
|
||||
@@ -492,9 +581,9 @@ async def main():
|
||||
ui_task.cancel()
|
||||
await renderer.render()
|
||||
|
||||
# Close persistent downloaders
|
||||
for dl in downloaders.values():
|
||||
await dl.disconnect_all()
|
||||
# Close persistent DC pools
|
||||
for pool in dc_pools.values():
|
||||
await pool.disconnect_all()
|
||||
|
||||
try:
|
||||
with open(history_file, 'w', encoding='utf-8') as f:
|
||||
|
||||
Reference in new issue
Block a user