Files

81 lines
2.4 KiB
Python

# -*- coding: utf-8 -*-
"""检测模块:并发检测 IPTV 源可用性"""
import requests
from concurrent.futures import ThreadPoolExecutor, as_completed
from config import CHECK_TIMEOUT, MAX_WORKERS, HEADERS
def check_channel(channel):
"""检测单个频道是否可用
返回 (channel, is_ok) 元组
"""
url = channel.get("url", "")
if not url:
return channel, False
try:
# 使用流式请求,只读取少量数据
resp = requests.get(
url,
headers=HEADERS,
timeout=CHECK_TIMEOUT,
stream=True,
allow_redirects=True,
)
# 状态码检查
if resp.status_code != 200:
return channel, False
# 读取前 1KB 数据验证内容
content = b""
for chunk in resp.iter_content(chunk_size=1024):
content += chunk
if len(content) >= 1024:
break
# 检查是否为有效的流内容(m3u8 或 ts 流)
if content:
text = content[:200].decode("utf-8", errors="ignore")
# m3u8 文件或以 #EXTM3U 开头
if "#EXTM3U" in text or "#EXTINF" in text:
return channel, True
# ts 流通常以 0x47 开头
if content[:1] == b"\x47":
return channel, True
# 其他情况认为可用(有些流返回的是二进制数据)
if len(content) > 100:
return channel, True
return channel, False
except Exception:
return channel, False
def check_all(channels):
"""并发检测所有频道"""
print(f"[检测] 开始检测 {len(channels)} 个频道,并发数 {MAX_WORKERS}")
valid = []
total = len(channels)
done = 0
with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor:
futures = {executor.submit(check_channel, ch): ch for ch in channels}
for future in as_completed(futures):
done += 1
channel, ok = future.result()
if ok:
valid.append(channel)
if done % 50 == 0 or done == total:
print(f"[检测进度] {done}/{total},可用 {len(valid)}")
print(f"[检测完成] 可用频道 {len(valid)}/{total}")
return valid
if __name__ == "__main__":
from collector import collect_all
channels = collect_all()
valid = check_all(channels)
print(f"最终可用:{len(valid)}")