commit 4094192349bfc1942249ac513ada511ed5d7ce5a
parent 8fb73a22ecc40455dc05cc76e6528ab3d6194aae
Author: david cochran <about.trout@gmail.com>
Date: Sun, 11 Feb 2024 19:22:30 -0800
refactor into StreamSampler class
Diffstat:
| D | scrape.py | | | 68 | -------------------------------------------------------------------- |
| A | streamsampler.py | | | 70 | ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ |
2 files changed, 70 insertions(+), 68 deletions(-)
diff --git a/scrape.py b/scrape.py
@@ -1,68 +0,0 @@
-#!/usr/bin/env python3
-
-import aiohttp
-import asyncio
-import datetime
-import logging
-import os
-import random
-import subprocess
-import tempfile
-import time
-
-
-async def get_recent_chunk_url(session, cam_name):
- base_url = f"https://cams.cdn-surfline.com/cdn-wc/wc-{cam_name}"
- async with session.get(f"{base_url}/chunklist.m3u8") as res:
- chunklist = await res.text()
- chunk = next(filter(lambda x: not x.startswith('#'), chunklist.splitlines()))
- return f"{base_url}/{chunk}"
-
-
-def split_keyframes(chunk_url, data_dir, cam_name):
- dt = datetime.datetime.now()
- day, hour, ts = dt.strftime("%m%d%Y"), dt.strftime("%H"), dt.strftime("%s")
-
- keyframes_dir = f"{data_dir}/{cam_name}/{day}/{hour}"
- if not os.path.exists(keyframes_dir):
- os.makedirs(keyframes_dir)
-
- subprocess.run(["ffmpeg",
- "-hide_banner", "-loglevel", "0",
- "-i", chunk_url,
- "-vf", "select=eq(pict_type\,I)",
- "-vsync", "vfr",
- f"{keyframes_dir}/{ts}-thumb-%04d.jpg"])
-
-
-async def save_keyframes_from_chunk(cam_name, data_dir):
- async with aiohttp.ClientSession(raise_for_status=True) as session:
- chunk_url = await get_recent_chunk_url(session, cam_name)
- split_keyframes(chunk_url, data_dir, cam_name)
- logging.info(f"cam_name={cam_name} Extracted keyframes from chunk")
-
-
-async def sample_stream(cam_name, data_dir):
- while True:
- try:
- await save_keyframes_from_chunk(cam_name, data_dir)
- except Exception as ex:
- logging.error(f"cam_name={cam_name} Failed to extract keyframes from chunk: {ex}")
-
- delay = random.randint(30, 90) # seconds.
- logging.info(f"cam_name={cam_name} delay={delay} Sleeping")
- await asyncio.sleep(delay)
-
-
-async def main():
- data_dir = os.environ.get("DATA_DIR", "./data")
- if not os.path.exists(data_dir):
- os.makedirs(data_dir)
- async with asyncio.TaskGroup() as tg:
- tg.create_task(sample_stream("mavericks", data_dir))
- tg.create_task(sample_stream("mavericksov", data_dir))
-
-
-if __name__ == "__main__":
- logging.basicConfig(level=logging.INFO, format='%(asctime)s %(message)s')
- asyncio.run(main())
diff --git a/streamsampler.py b/streamsampler.py
@@ -0,0 +1,70 @@
+#!/usr/bin/env python3
+
+import aiohttp
+import asyncio
+import datetime
+import logging
+import os
+import random
+import subprocess
+import tempfile
+import time
+
+
+class StreamSampler:
+ def __init__(self, cam_name, data_dir="./data"):
+ self.base_url = f"https://cams.cdn-surfline.com/cdn-wc/wc-{cam_name}"
+ self.cam_name = cam_name
+ self.data_dir = data_dir
+
+ async def get_recent_frames(self):
+ async with aiohttp.ClientSession(raise_for_status=True) as session:
+ chunk_url = await self.__get_recent_chunk_url(session)
+ return self.__extract_frames(chunk_url)
+
+ async def __get_recent_chunk_url(self, session):
+ async with session.get(f"{self.base_url}/chunklist.m3u8") as res:
+ chunklist = await res.text()
+ chunk = next(l for l in chunklist.splitlines() if not l.startswith("#"))
+ return f"{self.base_url}/{chunk}"
+
+ def __extract_frames(self, chunk_url):
+ dt = datetime.datetime.now()
+ day, hour, ts = dt.strftime("%m%d%Y"), dt.strftime("%H"), dt.strftime("%s")
+ frames_dir = f"{self.data_dir}/{self.cam_name}/{day}/{hour}"
+ if not os.path.exists(frames_dir):
+ os.makedirs(frames_dir)
+ subprocess.run(["ffmpeg",
+ "-hide_banner", "-loglevel", "0",
+ "-i", chunk_url,
+ "-vf", "select=eq(pict_type\,I)",
+ "-vsync", "vfr",
+ f"{frames_dir}/{ts}-thumb-%04d.jpg"])
+
+ return [f for f in os.scandir(frames_dir) if f.name.startswith(f"{ts}")]
+
+
+async def sample_stream(cam_name):
+ ss = StreamSampler(cam_name)
+ while True:
+ try:
+ frames = await ss.get_recent_frames()
+ logging.info(
+ f"cam_name={cam_name} num_frames={len(frames)} Extracted frames"
+ )
+ except Exception as ex:
+ logging.error(f"cam_name={cam_name} Failed to get_recent_frames: {ex}")
+ delay = random.randint(5, 20) # seconds.
+ logging.info(f"cam_name={cam_name} delay={delay} Sleeping")
+ await asyncio.sleep(delay)
+
+
+async def main():
+ async with asyncio.TaskGroup() as tg:
+ tg.create_task(sample_stream("mavericks"))
+ tg.create_task(sample_stream("mavericksov"))
+
+
+if __name__ == "__main__":
+ logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s")
+ asyncio.run(main())