forked from retoor/devplacepy
fix: correct typo in progress tracking variable name from Progss to Progress
This commit is contained in:
@@ -0,0 +1,4 @@
|
||||
from devplacepy.services.base import BaseService
|
||||
from devplacepy.services.manager import service_manager
|
||||
|
||||
__all__ = ["BaseService", "service_manager"]
|
||||
@@ -0,0 +1,93 @@
|
||||
import asyncio
|
||||
import logging
|
||||
from abc import ABC, abstractmethod
|
||||
from collections import deque
|
||||
from datetime import datetime
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class BaseService(ABC):
|
||||
def __init__(self, name: str, interval_seconds: int = 3600):
|
||||
self.name = name
|
||||
self.interval_seconds = interval_seconds
|
||||
self.log_buffer = deque(maxlen=20)
|
||||
self._task = None
|
||||
self._running = False
|
||||
self._last_run = None
|
||||
self._next_run = None
|
||||
self._started_at = None
|
||||
|
||||
@property
|
||||
def status(self) -> str:
|
||||
return "running" if self._running else "stopped"
|
||||
|
||||
@property
|
||||
def last_run(self) -> str | None:
|
||||
return self._last_run
|
||||
|
||||
@property
|
||||
def next_run(self) -> str | None:
|
||||
return self._next_run
|
||||
|
||||
@property
|
||||
def uptime(self) -> str | None:
|
||||
if self._started_at is None:
|
||||
return None
|
||||
delta = datetime.utcnow() - self._started_at
|
||||
return str(delta).split(".")[0]
|
||||
|
||||
def log(self, message: str) -> None:
|
||||
stamp = datetime.utcnow().strftime("%H:%M:%S")
|
||||
entry = f"[{stamp}] {message}"
|
||||
self.log_buffer.append(entry)
|
||||
logger.info(f"[{self.name}] {message}")
|
||||
|
||||
@abstractmethod
|
||||
async def run_once(self) -> None:
|
||||
pass
|
||||
|
||||
async def _run_loop(self) -> None:
|
||||
self._running = True
|
||||
self._started_at = datetime.utcnow()
|
||||
self.log(f"Service started (interval={self.interval_seconds}s)")
|
||||
while self._running:
|
||||
try:
|
||||
self._last_run = datetime.utcnow().isoformat()
|
||||
self._next_run = None
|
||||
await self.run_once()
|
||||
except asyncio.CancelledError:
|
||||
self.log("Service cancelled")
|
||||
break
|
||||
except Exception as e:
|
||||
self.log(f"Error in run_once: {e}")
|
||||
if not self._running:
|
||||
break
|
||||
self._next_run = datetime.utcnow().isoformat()
|
||||
self.log(f"Sleeping for {self.interval_seconds}s")
|
||||
try:
|
||||
await asyncio.sleep(self.interval_seconds)
|
||||
except asyncio.CancelledError:
|
||||
self.log("Service cancelled during sleep")
|
||||
break
|
||||
self._running = False
|
||||
self.log("Service stopped")
|
||||
|
||||
async def start(self) -> None:
|
||||
if self._running:
|
||||
self.log("Already running")
|
||||
return
|
||||
self._running = True
|
||||
self._task = asyncio.create_task(self._run_loop())
|
||||
|
||||
async def stop(self) -> None:
|
||||
if not self._running:
|
||||
return
|
||||
self._running = False
|
||||
if self._task is not None:
|
||||
self._task.cancel()
|
||||
try:
|
||||
await asyncio.wait_for(self._task, timeout=10.0)
|
||||
except (asyncio.CancelledError, asyncio.TimeoutError):
|
||||
pass
|
||||
self._task = None
|
||||
@@ -0,0 +1,47 @@
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ServiceManager:
|
||||
def __init__(self):
|
||||
self._services = {}
|
||||
|
||||
def register(self, service) -> None:
|
||||
self._services[service.name] = service
|
||||
logger.info(f"Registered service: {service.name}")
|
||||
|
||||
def get_service(self, name: str):
|
||||
return self._services.get(name)
|
||||
|
||||
def list_services(self) -> list[dict]:
|
||||
return [
|
||||
{
|
||||
"name": svc.name,
|
||||
"status": svc.status,
|
||||
"interval_seconds": svc.interval_seconds,
|
||||
"last_run": svc.last_run,
|
||||
"next_run": svc.next_run,
|
||||
"uptime": svc.uptime,
|
||||
"log_buffer": list(svc.log_buffer),
|
||||
}
|
||||
for svc in self._services.values()
|
||||
]
|
||||
|
||||
async def start_all(self) -> None:
|
||||
for name, svc in self._services.items():
|
||||
await svc.start()
|
||||
logger.info(f"Started service: {name}")
|
||||
|
||||
async def stop_all(self) -> None:
|
||||
tasks = []
|
||||
for name, svc in self._services.items():
|
||||
tasks.append(svc.stop())
|
||||
logger.info(f"Stopping service: {name}")
|
||||
if tasks:
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
logger.info("All services stopped")
|
||||
|
||||
|
||||
service_manager = ServiceManager()
|
||||
@@ -0,0 +1,184 @@
|
||||
import logging
|
||||
import re
|
||||
from datetime import datetime
|
||||
|
||||
import httpx
|
||||
|
||||
from devplacepy.database import get_table
|
||||
from devplacepy.services.base import BaseService
|
||||
from devplacepy.utils import generate_uid
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
NEWS_API_URL_DEFAULT = "https://news.app.molodetz.nl/api"
|
||||
AI_URL_DEFAULT = "https://openai.app.molodetz.nl/v1/chat/completions"
|
||||
AI_MODEL_DEFAULT = "molodetz"
|
||||
GRADE_THRESHOLD_DEFAULT = 7
|
||||
|
||||
|
||||
def _get_setting(key: str, default: str) -> str:
|
||||
table = get_table("site_settings")
|
||||
row = table.find_one(key=key)
|
||||
if row is None:
|
||||
return default
|
||||
return row.get("value", default)
|
||||
|
||||
|
||||
def _extract_grade(text: str) -> int | None:
|
||||
match = re.search(r"\d+", text.strip())
|
||||
if match:
|
||||
val = int(match.group())
|
||||
if 1 <= val <= 10:
|
||||
return val
|
||||
return None
|
||||
|
||||
|
||||
async def _get_article_images(url: str) -> list[dict]:
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10.0) as client:
|
||||
resp = await client.get(url)
|
||||
resp.raise_for_status()
|
||||
html = resp.text
|
||||
pattern = re.compile(r'<img[^>]+src=["\']([^"\']+)["\']', re.IGNORECASE)
|
||||
matches = pattern.findall(html)
|
||||
images = []
|
||||
seen = set()
|
||||
for src in matches:
|
||||
if src in seen:
|
||||
continue
|
||||
seen.add(src)
|
||||
if src.startswith("http") and not any(ext in src.lower() for ext in [".svg", ".ico"]):
|
||||
images.append({"url": src, "alt_text": ""})
|
||||
return images[:10]
|
||||
except Exception:
|
||||
return []
|
||||
|
||||
|
||||
class NewsService(BaseService):
|
||||
def __init__(self):
|
||||
super().__init__(name="news", interval_seconds=3600)
|
||||
|
||||
async def run_once(self) -> None:
|
||||
api_url = _get_setting("news_api_url", NEWS_API_URL_DEFAULT)
|
||||
ai_url = _get_setting("news_ai_url", AI_URL_DEFAULT)
|
||||
ai_model = _get_setting("news_ai_model", AI_MODEL_DEFAULT)
|
||||
threshold = int(_get_setting("news_grade_threshold", str(GRADE_THRESHOLD_DEFAULT)))
|
||||
|
||||
self.log(f"Fetching news from {api_url}")
|
||||
async with httpx.AsyncClient(timeout=30.0) as client:
|
||||
try:
|
||||
resp = await client.get(api_url)
|
||||
resp.raise_for_status()
|
||||
data = resp.json()
|
||||
except Exception as e:
|
||||
self.log(f"Failed to fetch news API: {e}")
|
||||
return
|
||||
|
||||
articles = data.get("articles", [])
|
||||
self.log(f"Received {len(articles)} articles")
|
||||
|
||||
sync_table = get_table("news_sync")
|
||||
news_table = get_table("news")
|
||||
images_table = get_table("news_images")
|
||||
|
||||
synced = 0
|
||||
skipped = 0
|
||||
below = 0
|
||||
failed = 0
|
||||
|
||||
for article in articles:
|
||||
external_id = article.get("guid", "")
|
||||
if not external_id:
|
||||
continue
|
||||
|
||||
existing = sync_table.find_one(external_id=external_id)
|
||||
if existing is not None:
|
||||
skipped += 1
|
||||
continue
|
||||
|
||||
grade = await self._grade_article(article, ai_url, ai_model)
|
||||
now = datetime.utcnow().isoformat()
|
||||
|
||||
if grade is None:
|
||||
failed += 1
|
||||
sync_table.insert({
|
||||
"external_id": external_id,
|
||||
"grade": 0,
|
||||
"status": "grading_failed",
|
||||
"synced_at": now,
|
||||
})
|
||||
continue
|
||||
|
||||
sync_table.insert({
|
||||
"external_id": external_id,
|
||||
"grade": grade,
|
||||
"status": "below_threshold" if grade < threshold else "synced",
|
||||
"synced_at": now,
|
||||
})
|
||||
|
||||
if grade < threshold:
|
||||
below += 1
|
||||
continue
|
||||
|
||||
news_uid = generate_uid()
|
||||
|
||||
news_table.insert({
|
||||
"uid": news_uid,
|
||||
"external_id": external_id,
|
||||
"title": article.get("title", ""),
|
||||
"description": (article.get("description", "") or "")[:5000],
|
||||
"url": article.get("link", ""),
|
||||
"image_url": "",
|
||||
"source_name": article.get("feed_name", ""),
|
||||
"grade": grade,
|
||||
"content": (article.get("content", "") or "")[:10000],
|
||||
"author": article.get("author", ""),
|
||||
"published": article.get("published", ""),
|
||||
"synced_at": now,
|
||||
})
|
||||
|
||||
link = article.get("link", "")
|
||||
if link:
|
||||
images = await _get_article_images(link)
|
||||
for img in images:
|
||||
images_table.insert({
|
||||
"uid": generate_uid(),
|
||||
"news_uid": news_uid,
|
||||
"url": img["url"],
|
||||
"alt_text": img.get("alt_text", ""),
|
||||
})
|
||||
|
||||
synced += 1
|
||||
|
||||
self.log(f"Synced {synced}, skipped {skipped}, below {below}, grading failed {failed}")
|
||||
|
||||
async def _grade_article(self, article: dict, ai_url: str, ai_model: str) -> int | None:
|
||||
title = (article.get("title", "") or "")[:500]
|
||||
description = (article.get("description", "") or "")[:1000]
|
||||
content = (article.get("content", "") or "")[:1500]
|
||||
|
||||
prompt = (
|
||||
"Rate this article's relevance to software developers on a scale of 1-10. "
|
||||
"Return only a single integer between 1 and 10, nothing else.\n\n"
|
||||
f"Title: {title}\n"
|
||||
f"Description: {description}\n"
|
||||
f"Content: {content}"
|
||||
)
|
||||
|
||||
payload = {
|
||||
"model": ai_model,
|
||||
"messages": [{"role": "user", "content": prompt}],
|
||||
"max_tokens": 10,
|
||||
"temperature": 0.0,
|
||||
}
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=15.0) as client:
|
||||
resp = await client.post(ai_url, json=payload)
|
||||
resp.raise_for_status()
|
||||
result = resp.json()
|
||||
text = result.get("choices", [{}])[0].get("message", {}).get("content", "")
|
||||
return _extract_grade(text)
|
||||
except Exception as e:
|
||||
self.log(f"AI grading failed: {e}")
|
||||
return None
|
||||
Reference in New Issue
Block a user