release: bump version to 0.69.0
Some checks failed
ci / backend (push) Has been cancelled
ci / frontend (push) Has been cancelled
ci / delivery (push) Has been cancelled
release / images (push) Has been cancelled

This commit is contained in:
linkong
2026-06-03 17:27:00 +08:00
parent 06aca980d0
commit acbbfdf9e2
57 changed files with 5907 additions and 316 deletions

View File

@@ -21,6 +21,12 @@ from app.models.datasource_config import DataSourceConfig
from app.models.system_setting import SystemSetting
from app.models.user import User
from app.services.tv_streams import get_tv_settings_payload
from app.services.earth_news import (
get_earth_news_sources_payload,
reset_earth_news_sources_payload,
save_earth_news_sources_payload,
test_news_source_config,
)
from app.services.earth_boundaries import (
EarthBoundaryBuildError,
get_boundary_build_status,
@@ -100,6 +106,19 @@ class EarthAboutPayload(BaseModel):
meta: list[EarthAboutMetaItem] = Field(default_factory=list)
class EarthNewsSourcesPayload(BaseModel):
cache_version: int | None = None
source_tags: list[dict[str, Any]] = Field(default_factory=list)
categories: list[dict[str, Any]] = Field(default_factory=list)
item_tag_rules: list[dict[str, Any]] = Field(default_factory=list)
sources: list[dict[str, Any]] = Field(default_factory=list)
health: dict[str, Any] = Field(default_factory=dict)
class EarthNewsSourceTestPayload(BaseModel):
source: dict[str, Any] = Field(default_factory=dict)
def _normalize_earth_brand_payload(payload: dict[str, Any] | None) -> dict[str, str]:
merged = DEFAULT_EARTH_BRAND.copy()
if payload:
@@ -324,6 +343,38 @@ async def reset_earth_about(
return {"status": "reset", "about": _normalize_earth_about_payload(None), "is_default": True}
@router.get("/news-sources")
async def get_earth_news_sources(db: AsyncSession = Depends(get_db)):
return await get_earth_news_sources_payload(db)
@router.put("/news-sources")
async def update_earth_news_sources(
payload: EarthNewsSourcesPayload,
_current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
return await save_earth_news_sources_payload(db, payload.model_dump())
@router.delete("/news-sources")
@router.post("/news-sources/reset")
async def reset_earth_news_sources(
_current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
return await reset_earth_news_sources_payload(db)
@router.post("/news-sources/test")
async def test_earth_news_source(
payload: EarthNewsSourceTestPayload,
_current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
return await test_news_source_config(payload.source, db=db)
@router.get("/oobe-status")
async def get_earth_oobe_status(
current_user: User | None = Depends(_get_optional_current_user),

View File

@@ -90,7 +90,7 @@ async def get_interactables_geojson(
return interactables_to_geojson(items)
payload = await get_or_build_layer_payload(
key=earth_layer_cache.key("interactables", layer=layer or "all"),
key=earth_layer_cache.key("interactables", interactable_layer=layer or "all"),
policy=INTERACTABLE_CACHE_POLICY,
builder=build_payload,
response=response,

View File

@@ -1,16 +1,92 @@
from fastapi import APIRouter, Depends, Query
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy.ext.asyncio import AsyncSession
from app.db.session import get_db
from app.services.earth_news import get_earth_news_payload
from app.services.earth_news import (
ALLOWED_NEWS_CATEGORY_KEYS,
SUPPORTED_NEWS_LOCALES,
REGION_ANCHORS,
get_earth_news_payload,
)
router = APIRouter()
def _parse_categories(raw: str | None) -> set[str] | None:
if raw is None or not raw.strip():
return None
requested = {item.strip().lower() for item in raw.split(",") if item.strip()}
invalid = sorted(requested - set(ALLOWED_NEWS_CATEGORY_KEYS))
if invalid:
raise HTTPException(
status_code=422,
detail={
"message": "Unsupported news categories.",
"invalid_categories": invalid,
"allowed_categories": list(ALLOWED_NEWS_CATEGORY_KEYS),
},
)
return requested or None
def _parse_source_ids(raw: str | None) -> set[str] | None:
if raw is None or not raw.strip():
return None
return {item.strip() for item in raw.split(",") if item.strip()} or None
def _parse_limit(raw: int | None) -> int:
if raw is None:
return 12
if raw < 1:
raise HTTPException(status_code=422, detail={"message": "News limit must be greater than 0."})
return min(raw, 100)
def _parse_locale(raw: str | None) -> str:
if raw is None or not raw.strip():
return "zh-CN"
requested = raw.strip()
if requested not in SUPPORTED_NEWS_LOCALES:
raise HTTPException(
status_code=422,
detail={
"message": "Unsupported news locale.",
"invalid_locale": requested,
"allowed_locales": sorted(SUPPORTED_NEWS_LOCALES),
},
)
return requested
@router.get("/earth-feed")
async def get_earth_feed(
lat: float | None = Query(None, description="Current Earth view center latitude"),
lon: float | None = Query(None, description="Current Earth view center longitude"),
region: str | None = Query(None, description="Explicit Earth news region for UE/client integrations"),
categories: str | None = Query(None, description="Comma-separated news category keys"),
sources: str | None = Query(None, description="Comma-separated news source ids"),
limit: int | None = Query(None, description="Maximum news items to return, capped at 100"),
locale: str | None = Query(None, description="Display locale, zh-CN or en-US"),
db: AsyncSession = Depends(get_db),
):
return await get_earth_news_payload(lat=lat, lon=lon, db=db)
normalized_region = region.strip().lower() if isinstance(region, str) and region.strip() else None
if normalized_region is not None and normalized_region not in REGION_ANCHORS:
raise HTTPException(
status_code=422,
detail={
"message": "Unsupported news region.",
"invalid_region": normalized_region,
"allowed_regions": list(REGION_ANCHORS.keys()),
},
)
return await get_earth_news_payload(
lat=lat,
lon=lon,
region=normalized_region,
categories=_parse_categories(categories),
source_ids=_parse_source_ids(sources),
limit=_parse_limit(limit),
locale=_parse_locale(locale),
db=db,
)

View File

@@ -1,16 +1,17 @@
from __future__ import annotations
import os
import secrets
import subprocess
import sys
from datetime import datetime
from fastapi import APIRouter, Depends, HTTPException, Query, Request, status
from fastapi import APIRouter, Depends, Header, HTTPException, Query, Request, status
from pydantic import BaseModel
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.config import ROOT_DIR
from app.core.config import ROOT_DIR, settings
from app.core.security import get_current_user
from app.db.session import get_db
from app.models.user import User
@@ -37,6 +38,9 @@ from app.services.system_logs import (
normalize_log_level,
read_database_log_snapshot,
read_log_snapshot,
read_observability_group_events,
read_observability_groups,
read_observability_raw_events,
)
from app.services.earth_layer_cache import earth_layer_cache
@@ -108,12 +112,34 @@ class EarthClientLogEventCreate(BaseModel):
url: str | None = None
module: str | None = None
detail: str | None = None
fingerprint: str | None = None
occurrence_count: int = 1
metadata: dict[str, object] | None = None
class EarthClientLogEventResponse(BaseModel):
accepted: bool
source_id: str
level: str
fingerprint: str | None = None
class ServiceLogEventCreate(BaseModel):
source: str = "ai-provider"
service: str = "ai-provider"
module: str | None = None
category: str | None = None
event: str = "service.runtime_log"
level: str = "error"
message: str
fingerprint: str | None = None
occurrence_count: int = 1
request_id: str | None = None
trace_id: str | None = None
task_id: str | None = None
source_id: int | str | None = None
provider: str | None = None
context: dict[str, object] | None = None
async def ingest_client_log_event(
@@ -136,6 +162,9 @@ async def ingest_client_log_event(
"url": payload.url or "",
"module": payload.module or "",
"detail": payload.detail or "",
"fingerprint": payload.fingerprint or "",
"occurrence_count": max(1, int(payload.occurrence_count or 1)),
"metadata": payload.metadata or {},
},
)
await record_system_log(
@@ -151,9 +180,36 @@ async def ingest_client_log_event(
"detail": payload.detail or "",
"module": payload.module or "",
"client_ip": request.client.host if request.client else "",
"metadata": payload.metadata or {},
},
fingerprint=payload.fingerprint,
occurrence_count=max(1, int(payload.occurrence_count or 1)),
)
return EarthClientLogEventResponse(accepted=True, source_id=source_id, level=normalized_level)
return EarthClientLogEventResponse(accepted=True, source_id=source_id, level=normalized_level, fingerprint=payload.fingerprint)
def require_observability_ingest_token(
authorization: str | None,
ingest_token: str | None,
) -> None:
expected_token = settings.OBSERVABILITY_INGEST_TOKEN.strip()
if not expected_token:
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="Observability service ingestion is not configured",
)
provided = ""
if ingest_token:
provided = ingest_token.strip()
elif authorization:
scheme, _, token = authorization.partition(" ")
if scheme.lower() == "bearer":
provided = token.strip()
if not provided or not secrets.compare_digest(provided, expected_token):
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN,
detail="Invalid observability ingestion token",
)
class EarthLayerCacheStatusResponse(BaseModel):
@@ -378,6 +434,118 @@ async def get_system_log_sources(
}
@router.get("/logs/observability/groups")
async def get_observability_log_groups(
limit: int = DEFAULT_LOG_LINE_LIMIT,
level: str = "all",
levels: str | None = Query(None, description="Comma-separated log levels"),
start_date: str | None = Query(None, description="Filter logs from this date (YYYY-MM-DD)"),
end_date: str | None = Query(None, description="Filter logs until this date (YYYY-MM-DD)"),
search: str | None = Query(None, description="Case-insensitive substring search"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
ensure_super_admin(current_user)
if limit < 1 or limit > MAX_LOG_LINE_LIMIT:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=f"limit must be between 1 and {MAX_LOG_LINE_LIMIT}")
normalized_start_date = validate_log_date(start_date, "start_date")
normalized_end_date = validate_log_date(end_date, "end_date")
if normalized_start_date and normalized_end_date and normalized_start_date > normalized_end_date:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="start_date must be earlier than or equal to end_date")
return await read_observability_groups(
limit=limit,
level=level,
levels=levels,
start_date=normalized_start_date,
end_date=normalized_end_date,
search=search,
db=db,
)
@router.get("/logs/observability/groups/{fingerprint}/events")
async def get_observability_group_events(
fingerprint: str,
limit: int = DEFAULT_LOG_LINE_LIMIT,
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
ensure_super_admin(current_user)
if limit < 1 or limit > MAX_LOG_LINE_LIMIT:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=f"limit must be between 1 and {MAX_LOG_LINE_LIMIT}")
payload = await read_observability_group_events(fingerprint, limit=limit, db=db)
if payload is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Observability group not found")
return payload
@router.get("/logs/observability/raw")
async def get_observability_raw_events(
limit: int = DEFAULT_LOG_LINE_LIMIT,
level: str = "all",
levels: str | None = Query(None, description="Comma-separated log levels"),
start_date: str | None = Query(None, description="Filter logs from this date (YYYY-MM-DD)"),
end_date: str | None = Query(None, description="Filter logs until this date (YYYY-MM-DD)"),
search: str | None = Query(None, description="Case-insensitive substring search"),
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
ensure_super_admin(current_user)
if limit < 1 or limit > MAX_LOG_LINE_LIMIT:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=f"limit must be between 1 and {MAX_LOG_LINE_LIMIT}")
normalized_start_date = validate_log_date(start_date, "start_date")
normalized_end_date = validate_log_date(end_date, "end_date")
return await read_observability_raw_events(
limit=limit,
level=level,
levels=levels,
start_date=normalized_start_date,
end_date=normalized_end_date,
search=search,
db=db,
)
@router.post("/logs/service", response_model=EarthClientLogEventResponse)
async def ingest_service_log(
payload: ServiceLogEventCreate,
authorization: str | None = Header(default=None),
ingest_token: str | None = Header(default=None, alias="X-Planet-Observability-Token"),
):
require_observability_ingest_token(authorization, ingest_token)
normalized_level = normalize_log_level(payload.level)
source = (payload.source or "ai-provider").strip() or "ai-provider"
context = dict(payload.context or {})
if payload.request_id:
context["request_id"] = payload.request_id
if payload.trace_id:
context["trace_id"] = payload.trace_id
if payload.task_id:
context["task_id"] = payload.task_id
if payload.source_id is not None:
context["source_id"] = payload.source_id
if payload.provider:
context["provider"] = payload.provider
await record_system_log(
source=source,
service=(payload.service or source).strip() or source,
module=payload.module or source,
event=(payload.event or "service.runtime_log").strip() or "service.runtime_log",
level=normalized_level,
message=payload.message,
category=payload.category or "service-runtime",
context=context,
fingerprint=payload.fingerprint,
occurrence_count=max(1, int(payload.occurrence_count or 1)),
)
return EarthClientLogEventResponse(
accepted=True,
source_id=source,
level=normalized_level,
fingerprint=payload.fingerprint,
)
@router.get("/logs/{source_id}", response_model=SystemLogSnapshotResponse)
async def get_system_log_snapshot(
source_id: str,

View File

@@ -1,3 +1,4 @@
import re
from urllib.parse import quote, urljoin
import httpx
@@ -10,6 +11,26 @@ from app.services.tv_streams import get_public_tv_payload, is_allowed_tv_proxy_u
router = APIRouter()
_HLS_URI_ATTRIBUTE_RE = re.compile(r'URI="([^"]+)"')
def _proxied_tv_url(url: str) -> str:
return f"/api/v1/tv/proxy?url={quote(url, safe='')}"
def _rewrite_hls_uri_attributes(line: str, *, base_url: str) -> str:
def replace(match: re.Match[str]) -> str:
uri = match.group(1)
absolute_url = urljoin(base_url, uri)
return f'URI="{_proxied_tv_url(absolute_url)}"'
return _HLS_URI_ATTRIBUTE_RE.sub(replace, line)
def _should_strip_hls_metadata_line(line: str) -> bool:
normalized = line.strip().upper()
return normalized.startswith("#EXT-X-MEDIA:") and "TYPE=SUBTITLES" in normalized
@router.get("/streams")
async def list_public_tv_streams(
@@ -56,11 +77,16 @@ async def proxy_tv_stream(
rewritten_lines: list[str] = []
for line in manifest_text.splitlines():
stripped = line.strip()
if not stripped or stripped.startswith("#"):
if not stripped:
rewritten_lines.append(line)
continue
if stripped.startswith("#"):
if _should_strip_hls_metadata_line(line):
continue
rewritten_lines.append(_rewrite_hls_uri_attributes(line, base_url=response_url))
continue
absolute_url = urljoin(response_url, stripped)
rewritten_lines.append(f"/api/v1/tv/proxy?url={quote(absolute_url, safe='')}")
rewritten_lines.append(_proxied_tv_url(absolute_url))
return Response(
content="\n".join(rewritten_lines),
media_type="application/vnd.apple.mpegurl",

View File

@@ -41,6 +41,7 @@ class Settings(BaseSettings):
AI_PROVIDER_SERVICE_TOKEN: str = ""
AI_PROVIDER_TIMEOUT_SECONDS: int = 60
AI_PROVIDER_RETRY_ATTEMPTS: int = 2
OBSERVABILITY_INGEST_TOKEN: str = ""
@property
def REDIS_URL(self) -> str:

View File

@@ -13,7 +13,7 @@ from app.models.compute_center_location import ComputeCenterLocationRecord
from app.models.system_setting import SystemSetting
from app.models.playground_session import PlaygroundSession
from app.models.playground_message import PlaygroundMessage
from app.models.system_log import SystemLog, AuditLog
from app.models.system_log import AuditLog, ObservabilityEvent, ObservabilityEventGroup, SystemLog
from app.models.vessel import AISConflictRecord, AISRawObservation, AISSourceHealth, VesselPosition, VesselStatic
from app.models.datasource_mapping import DataSourceMappingTemplate
from app.models.earth_news import EarthNewsItem
@@ -37,6 +37,8 @@ __all__ = [
"ComputeCenterLocationRecord",
"SystemLog",
"AuditLog",
"ObservabilityEvent",
"ObservabilityEventGroup",
"PlaygroundSession",
"PlaygroundMessage",
"VesselPosition",

View File

@@ -38,3 +38,46 @@ class AuditLog(Base):
ip = Column(String(64), nullable=True)
details = Column(JSON, nullable=False, default=dict)
created_at = Column(DateTime(timezone=True), server_default=func.now())
class ObservabilityEvent(Base):
__tablename__ = "observability_events"
id = Column(Integer, primary_key=True, autoincrement=True)
occurred_at = Column(DateTime(timezone=True), server_default=func.now(), index=True)
source = Column(String(50), nullable=False, index=True)
service = Column(String(50), nullable=True, index=True)
module = Column(String(120), nullable=True, index=True)
category = Column(String(80), nullable=True, index=True)
event = Column(String(160), nullable=True, index=True)
level = Column(String(20), nullable=False, index=True)
message = Column(Text, nullable=False)
fingerprint = Column(String(80), nullable=False, index=True)
request_id = Column(String(64), nullable=True, index=True)
trace_id = Column(String(64), nullable=True, index=True)
task_id = Column(String(120), nullable=True, index=True)
source_ref_id = Column(String(120), nullable=True, index=True)
provider = Column(String(120), nullable=True, index=True)
user_id = Column(Integer, nullable=True, index=True)
context = Column(JSON, nullable=False, default=dict)
occurrence_count = Column(Integer, nullable=False, default=1)
created_at = Column(DateTime(timezone=True), server_default=func.now())
class ObservabilityEventGroup(Base):
__tablename__ = "observability_event_groups"
fingerprint = Column(String(80), primary_key=True)
source = Column(String(50), nullable=False, index=True)
service = Column(String(50), nullable=True, index=True)
module = Column(String(120), nullable=True, index=True)
category = Column(String(80), nullable=True, index=True)
event = Column(String(160), nullable=True, index=True)
last_level = Column(String(20), nullable=False, index=True)
sample_message = Column(Text, nullable=False)
sample_detail = Column(Text, nullable=True)
affected_sources = Column(JSON, nullable=False, default=list)
count = Column(Integer, nullable=False, default=0)
first_seen_at = Column(DateTime(timezone=True), nullable=False, index=True)
last_seen_at = Column(DateTime(timezone=True), nullable=False, index=True)
updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now())

View File

@@ -62,11 +62,11 @@ def interactables_to_geojson(items: list[EarthInteractable]) -> dict[str, Any]:
def invalidate_interactable_cache(layer: str | None = None) -> int:
layer_key = str(layer or "*").strip() or "*"
deleted = earth_layer_cache.delete_pattern(
f"{EARTH_LAYER_CACHE_PREFIX}:interactables:layer:{layer_key}*"
f"{EARTH_LAYER_CACHE_PREFIX}:interactables:interactable_layer:{layer_key}*"
)
if layer_key != "all":
deleted += earth_layer_cache.delete_pattern(
f"{EARTH_LAYER_CACHE_PREFIX}:interactables:layer:all*"
f"{EARTH_LAYER_CACHE_PREFIX}:interactables:interactable_layer:all*"
)
return deleted

File diff suppressed because it is too large Load Diff

View File

@@ -11,6 +11,7 @@ from app.services.earth_news import (
ParsedNewsItem,
apply_enrichment_patch_to_item,
build_anchor_location_patch,
_news_meta_patch,
)
@@ -34,6 +35,8 @@ def _location_patch_from_record(record: EarthNewsItem) -> dict[str, Any]:
def record_to_parsed_news_item(record: EarthNewsItem) -> ParsedNewsItem:
location_meta = dict(record.location_meta or {})
news_meta = location_meta.get("news_meta") if isinstance(location_meta.get("news_meta"), dict) else {}
item = ParsedNewsItem(
id=record.id,
title=record.title,
@@ -49,11 +52,29 @@ def record_to_parsed_news_item(record: EarthNewsItem) -> ParsedNewsItem:
enrichment_status=record.enrichment_status or "pending",
enrichment_error=record.enrichment_error,
enriched_at=_coerce_datetime(record.enriched_at),
source_tags=list(news_meta.get("source_tags") or []),
feed_id=str(news_meta.get("feed_id") or ""),
feed_type=str(news_meta.get("feed_type") or "rss"),
feed_default_category=str(news_meta.get("feed_default_category") or "other"),
category=str(news_meta.get("category") or "other"),
item_tags=list(news_meta.get("item_tags") or []),
tagging_source=str(news_meta.get("tagging_source") or "rules"),
tagging_confidence=float(news_meta.get("tagging_confidence") or 0),
importance_score=int(news_meta.get("importance_score") or 0),
importance_level=str(news_meta.get("importance_level") or "low"),
importance_reasons=list(news_meta.get("importance_reasons") or []),
market_impact=str(news_meta.get("market_impact") or "none"),
)
return apply_enrichment_patch_to_item(item, _location_patch_from_record(record))
def _query_sort_key(active_region: str):
if active_region == "global":
return (
EarthNewsItem.published_at.is_(None),
EarthNewsItem.published_at.desc().nullslast(),
EarthNewsItem.feed_name.asc(),
)
return (
EarthNewsItem.region != active_region,
EarthNewsItem.published_at.is_(None),
@@ -62,28 +83,93 @@ def _query_sort_key(active_region: str):
)
def _category_filter_clause(categories: set[str] | None):
if not categories:
return None
return EarthNewsItem.location_meta.op("->")("news_meta").op("->>")("category").in_(sorted(categories))
def _source_filter_clause(source_ids: set[str] | None):
if not source_ids:
return None
return EarthNewsItem.location_meta.op("->")("news_meta").op("->>")("source_id").in_(sorted(source_ids))
def _record_source_id(record: EarthNewsItem) -> str:
location_meta = dict(record.location_meta or {})
news_meta = location_meta.get("news_meta") if isinstance(location_meta.get("news_meta"), dict) else {}
source_id = str(news_meta.get("source_id") or "").strip()
if source_id:
return source_id
if isinstance(record.id, str) and ":" in record.id:
return record.id.split(":", 1)[0]
return record.feed_name or record.source or record.id
def _diversify_records_by_source(records: list[EarthNewsItem], *, limit: int) -> list[EarthNewsItem]:
if limit <= 0 or len(records) <= limit:
return records[:limit]
buckets: dict[str, list[EarthNewsItem]] = {}
order: list[str] = []
for record in records:
source_id = _record_source_id(record)
if source_id not in buckets:
buckets[source_id] = []
order.append(source_id)
buckets[source_id].append(record)
diversified: list[EarthNewsItem] = []
while len(diversified) < limit and order:
next_order: list[str] = []
for source_id in order:
bucket = buckets.get(source_id) or []
if bucket and len(diversified) < limit:
diversified.append(bucket.pop(0))
if bucket:
next_order.append(source_id)
order = next_order
return diversified
async def list_earth_news_items(
db: AsyncSession,
*,
active_region: str,
limit: int,
categories: set[str] | None = None,
source_ids: set[str] | None = None,
) -> list[ParsedNewsItem]:
regions = {"global", active_region}
result = await db.execute(
query_limit = limit if source_ids else min(max(limit * 4, limit), 100)
query = (
select(EarthNewsItem)
.where(EarthNewsItem.region.in_(regions))
.order_by(*_query_sort_key(active_region))
.limit(limit)
.limit(query_limit)
)
return [record_to_parsed_news_item(record) for record in result.scalars().all()]
if active_region != "global":
query = query.where(EarthNewsItem.region.in_({"global", active_region}))
category_clause = _category_filter_clause(categories)
if category_clause is not None:
query = query.where(category_clause)
source_clause = _source_filter_clause(source_ids)
if source_clause is not None:
query = query.where(source_clause)
result = await db.execute(query)
records = list(result.scalars().all())
if not source_ids:
records = _diversify_records_by_source(records, limit=limit)
else:
records = records[:limit]
return [record_to_parsed_news_item(record) for record in records]
async def list_earth_news_cruise_items(
db: AsyncSession,
*,
limit: int,
categories: set[str] | None = None,
source_ids: set[str] | None = None,
) -> list[ParsedNewsItem]:
result = await db.execute(
query = (
select(EarthNewsItem)
.order_by(
EarthNewsItem.region.asc(),
@@ -93,6 +179,13 @@ async def list_earth_news_cruise_items(
)
.limit(limit)
)
category_clause = _category_filter_clause(categories)
if category_clause is not None:
query = query.where(category_clause)
source_clause = _source_filter_clause(source_ids)
if source_clause is not None:
query = query.where(source_clause)
result = await db.execute(query)
return [record_to_parsed_news_item(record) for record in result.scalars().all()]
@@ -101,13 +194,13 @@ async def get_earth_news_freshness(
*,
active_region: str,
) -> tuple[int, datetime | None]:
regions = {"global", active_region}
result = await db.execute(
select(
func.count(EarthNewsItem.id),
func.max(func.coalesce(EarthNewsItem.published_at, EarthNewsItem.last_seen_at)),
).where(EarthNewsItem.region.in_(regions))
query = select(
func.count(EarthNewsItem.id),
func.max(func.coalesce(EarthNewsItem.published_at, EarthNewsItem.last_seen_at)),
)
if active_region != "global":
query = query.where(EarthNewsItem.region.in_({"global", active_region}))
result = await db.execute(query)
count, newest = result.one()
item_count = int(count or 0)
if item_count == 0:
@@ -115,6 +208,33 @@ async def get_earth_news_freshness(
return item_count, _coerce_datetime(newest)
async def get_earth_news_feed_coverage(
db: AsyncSession,
*,
active_region: str,
recent_after: datetime | None = None,
) -> set[tuple[str, str]]:
query = select(
EarthNewsItem.id,
EarthNewsItem.location_meta.op("->")("news_meta").op("->>")("source_id"),
EarthNewsItem.location_meta.op("->")("news_meta").op("->>")("feed_id"),
)
if active_region != "global":
query = query.where(EarthNewsItem.region.in_({"global", active_region}))
if recent_after is not None:
query = query.where(func.coalesce(EarthNewsItem.published_at, EarthNewsItem.last_seen_at) >= recent_after)
result = await db.execute(query)
coverage: set[tuple[str, str]] = set()
for item_id, source_id, feed_id in result.all():
normalized_source_id = str(source_id or "").strip()
normalized_feed_id = str(feed_id or "").strip()
if not normalized_source_id and isinstance(item_id, str) and ":" in item_id:
normalized_source_id = item_id.split(":", 1)[0]
if normalized_source_id and normalized_feed_id:
coverage.add((normalized_source_id, normalized_feed_id))
return coverage
async def upsert_earth_news_items(db: AsyncSession, items: list[ParsedNewsItem]) -> int:
if not items:
return 0
@@ -165,12 +285,20 @@ async def upsert_earth_news_items(db: AsyncSession, items: list[ParsedNewsItem])
record.homepage_url = item.homepage_url
record.published_at = item.published_at
record.last_seen_at = now
location_meta = dict(record.location_meta or {})
location_meta["news_meta"] = _news_meta_patch(item)
record.location_meta = location_meta
if item.localizations:
merged_localizations = {
**dict(record.localizations or {}),
**dict(item.localizations or {}),
}
record.content_language = item.content_language
record.localizations = dict(item.localizations or {})
record.enrichment_status = item.enrichment_status
record.enrichment_error = item.enrichment_error
record.enriched_at = item.enriched_at
record.localizations = merged_localizations
if item.enrichment_status != "pending" or item.enrichment_error or item.enriched_at:
record.enrichment_status = item.enrichment_status
record.enrichment_error = item.enrichment_error
record.enriched_at = item.enriched_at
changed += 1
await db.flush()
return changed

View File

@@ -1,14 +1,185 @@
from __future__ import annotations
import hashlib
import re
from datetime import UTC, datetime
from typing import Any
from app.core.logging import get_logger, sanitize_log_value
from app.core.request_context import get_request_id
from app.db.session import async_session_factory
from app.models.system_log import AuditLog, SystemLog
from app.models.system_log import AuditLog, ObservabilityEvent, ObservabilityEventGroup, SystemLog
logger = get_logger(__name__)
HLS_TRANSIENT_RE = re.compile(r"(index|chunk|segment)[_-]?\d+(?:_\d+)?\.(?:ts|m4s|vtt)", re.IGNORECASE)
QUERY_RE = re.compile(r"([?&](?:m|t|token|expires|signature|X-Amz-[^=]+)=[^&\\s]+)", re.IGNORECASE)
UUID_RE = re.compile(r"\b[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\b", re.IGNORECASE)
CONNECTION_RE = re.compile(r"\bconn_[A-Za-z0-9:._-]+\b")
NUMBER_RE = re.compile(r"\b\d{5,}\b")
def normalize_observability_text(value: Any) -> str:
text = str(sanitize_log_value(value or "")).strip()
text = QUERY_RE.sub("", text)
text = HLS_TRANSIENT_RE.sub("<hls-fragment>", text)
text = UUID_RE.sub("<uuid>", text)
text = CONNECTION_RE.sub("<connection>", text)
text = NUMBER_RE.sub("<number>", text)
return re.sub(r"\s+", " ", text).strip()
def build_observability_fingerprint(
*,
source: str,
service: str | None = None,
module: str | None = None,
category: str | None = None,
event: str | None = None,
message: str,
context: dict[str, Any] | None = None,
) -> str:
context = context or {}
stable_context = {
key: context.get(key)
for key in (
"task_type",
"source_id",
"source",
"provider",
"status_code",
"error_type",
"details",
)
if context.get(key) not in (None, "")
}
raw = "|".join(
[
normalize_observability_text(source),
normalize_observability_text(service),
normalize_observability_text(module),
normalize_observability_text(category),
normalize_observability_text(event),
normalize_observability_text(message),
normalize_observability_text(stable_context),
]
)
return hashlib.sha1(raw.encode("utf-8", errors="replace")).hexdigest()
def _context_text(context: dict[str, Any] | None, key: str) -> str | None:
value = (context or {}).get(key)
if value in (None, ""):
return None
return str(value)
async def record_observability_event(
*,
source: str,
level: str,
message: str,
service: str | None = None,
module: str | None = None,
event: str | None = None,
request_id: str | None = None,
trace_id: str | None = None,
user_id: int | None = None,
category: str | None = None,
context: dict[str, Any] | None = None,
fingerprint: str | None = None,
occurred_at: datetime | None = None,
occurrence_count: int = 1,
) -> None:
normalized_context = sanitize_log_value(context or {})
if not isinstance(normalized_context, dict):
normalized_context = {"value": normalized_context}
safe_message = str(sanitize_log_value(message))
normalized_level = str(level or "info").lower()
count = max(1, int(occurrence_count or 1))
event_time = occurred_at or datetime.now(UTC)
event_fingerprint = fingerprint or build_observability_fingerprint(
source=source,
service=service,
module=module,
category=category,
event=event,
message=safe_message,
context=normalized_context,
)
detail = _context_text(normalized_context, "detail") or _context_text(normalized_context, "error")
affected_sources = sorted(
{
item
for item in (
source,
service,
module,
_context_text(normalized_context, "source_id"),
_context_text(normalized_context, "source"),
)
if item
}
)
try:
async with async_session_factory() as session:
session.add(
ObservabilityEvent(
source=source,
service=service,
module=module,
category=category,
event=event,
level=normalized_level,
message=safe_message,
fingerprint=event_fingerprint,
occurred_at=event_time,
request_id=request_id or get_request_id(),
trace_id=trace_id,
user_id=user_id,
task_id=_context_text(normalized_context, "task_id"),
source_ref_id=_context_text(normalized_context, "source_id") or _context_text(normalized_context, "source"),
provider=_context_text(normalized_context, "provider"),
context=normalized_context,
occurrence_count=count,
)
)
group = await session.get(ObservabilityEventGroup, event_fingerprint)
if group is None:
session.add(
ObservabilityEventGroup(
fingerprint=event_fingerprint,
source=source,
service=service,
module=module,
category=category,
event=event,
last_level=normalized_level,
sample_message=safe_message,
sample_detail=detail,
affected_sources=affected_sources,
count=count,
first_seen_at=event_time,
last_seen_at=event_time,
)
)
else:
group.count = int(group.count or 0) + count
group.last_seen_at = event_time
group.last_level = normalized_level
group.sample_message = safe_message
group.sample_detail = detail
merged_sources = sorted(set(group.affected_sources or []) | set(affected_sources))
group.affected_sources = merged_sources
await session.commit()
except Exception:
logger.exception_event(
"Failed to persist observability event",
event="observability_event.persist.failed",
context={"event_name": event, "source": source},
)
async def record_system_log(
*,
@@ -23,6 +194,8 @@ async def record_system_log(
user_id: int | None = None,
category: str | None = None,
context: dict[str, Any] | None = None,
fingerprint: str | None = None,
occurrence_count: int = 1,
) -> None:
try:
async with async_session_factory() as session:
@@ -48,6 +221,21 @@ async def record_system_log(
event="system_log.persist.failed",
context={"event_name": event, "source": source},
)
await record_observability_event(
source=source,
service=service,
module=module,
event=event,
level=level,
message=message,
request_id=request_id,
trace_id=trace_id,
user_id=user_id,
category=category,
context=context,
fingerprint=fingerprint,
occurrence_count=occurrence_count,
)
async def record_audit_log(

View File

@@ -14,7 +14,7 @@ from pathlib import Path
from typing import Any
from app.core.security import redis_client
from app.models.system_log import AuditLog, SystemLog
from app.models.system_log import AuditLog, ObservabilityEvent, ObservabilityEventGroup, SystemLog
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
@@ -120,6 +120,10 @@ class DailyLogMarker:
dominant_level: str
def _normalize_search_query(search: str | None) -> str:
return (search or "").strip().lower()
def _planet_state_dir() -> Path:
configured = os.getenv("PLANET_STATE_DIR")
if configured:
@@ -660,6 +664,228 @@ async def read_database_log_snapshot(
}
def _observability_group_matches(
group: ObservabilityEventGroup,
*,
selected_levels: tuple[str, ...],
start_date: str | None,
end_date: str | None,
search: str | None,
) -> bool:
if selected_levels and group.last_level not in selected_levels:
return False
if start_date or end_date:
if group.last_seen_at is None:
return False
date_token = group.last_seen_at.astimezone(UTC).date().isoformat()
if start_date and date_token < start_date:
return False
if end_date and date_token > end_date:
return False
query = _normalize_search_query(search)
if not query:
return True
haystack = " ".join(
[
group.fingerprint or "",
group.source or "",
group.service or "",
group.module or "",
group.category or "",
group.event or "",
group.last_level or "",
group.sample_message or "",
group.sample_detail or "",
json.dumps(group.affected_sources or [], ensure_ascii=False, sort_keys=True),
]
).lower()
return query in haystack
def _serialize_observability_group(group: ObservabilityEventGroup) -> dict[str, Any]:
return {
"fingerprint": group.fingerprint,
"source": group.source,
"service": group.service,
"module": group.module,
"category": group.category,
"event": group.event,
"level": group.last_level,
"message": group.sample_message,
"detail": group.sample_detail,
"affected_sources": group.affected_sources or [],
"count": group.count or 0,
"first_seen_at": group.first_seen_at.isoformat() if group.first_seen_at else None,
"last_seen_at": group.last_seen_at.isoformat() if group.last_seen_at else None,
}
def _serialize_observability_event(record: ObservabilityEvent) -> dict[str, Any]:
return {
"id": record.id,
"source": record.source,
"service": record.service,
"module": record.module,
"category": record.category,
"event": record.event,
"level": record.level,
"message": record.message,
"fingerprint": record.fingerprint,
"occurred_at": record.occurred_at.isoformat() if record.occurred_at else None,
"request_id": record.request_id,
"trace_id": record.trace_id,
"task_id": record.task_id,
"source_id": record.source_ref_id,
"provider": record.provider,
"user_id": record.user_id,
"context": record.context or {},
"occurrence_count": record.occurrence_count or 1,
}
async def read_observability_groups(
*,
limit: int,
level: str = LOG_LEVEL_ALL,
levels: str | None = None,
start_date: str | None = None,
end_date: str | None = None,
search: str | None = None,
db: AsyncSession,
) -> dict[str, Any]:
selected_levels = normalize_log_levels(level, levels)
scan_limit = max(limit * 5, limit, DEFAULT_LOG_LINE_LIMIT)
result = await db.execute(
select(ObservabilityEventGroup)
.order_by(ObservabilityEventGroup.last_seen_at.desc().nullslast())
.limit(scan_limit)
)
groups = [
group
for group in result.scalars().all()
if _observability_group_matches(
group,
selected_levels=selected_levels,
start_date=start_date,
end_date=end_date,
search=search,
)
][:limit]
return {
"mode": "grouped",
"line_limit": limit,
"line_count": len(groups),
"groups": [_serialize_observability_group(group) for group in groups],
"filters": {
"level": level,
"levels": list(selected_levels),
"start_date": start_date,
"end_date": end_date,
"search": search or "",
},
}
async def read_observability_group_events(
fingerprint: str,
*,
limit: int,
db: AsyncSession,
) -> dict[str, Any] | None:
group = await db.get(ObservabilityEventGroup, fingerprint)
if group is None:
return None
result = await db.execute(
select(ObservabilityEvent)
.where(ObservabilityEvent.fingerprint == fingerprint)
.order_by(ObservabilityEvent.occurred_at.desc().nullslast(), ObservabilityEvent.id.desc())
.limit(limit)
)
events = list(reversed(result.scalars().all()))
return {
"fingerprint": fingerprint,
"group": _serialize_observability_group(group),
"line_limit": limit,
"line_count": len(events),
"events": [_serialize_observability_event(record) for record in events],
}
async def read_observability_raw_events(
*,
limit: int,
level: str = LOG_LEVEL_ALL,
levels: str | None = None,
start_date: str | None = None,
end_date: str | None = None,
search: str | None = None,
db: AsyncSession,
) -> dict[str, Any]:
selected_levels = normalize_log_levels(level, levels)
query = select(ObservabilityEvent).order_by(ObservabilityEvent.occurred_at.desc().nullslast(), ObservabilityEvent.id.desc())
if selected_levels:
query = query.where(ObservabilityEvent.level.in_(selected_levels))
result = await db.execute(query.limit(max(limit * 5, limit)))
records = result.scalars().all()
search_query = _normalize_search_query(search)
visible: list[ObservabilityEvent] = []
for record in records:
if start_date or end_date:
if record.occurred_at is None:
continue
date_token = record.occurred_at.astimezone(UTC).date().isoformat()
if start_date and date_token < start_date:
continue
if end_date and date_token > end_date:
continue
if search_query:
haystack = " ".join(
[
record.source or "",
record.service or "",
record.module or "",
record.category or "",
record.event or "",
record.message or "",
record.fingerprint or "",
record.request_id or "",
record.trace_id or "",
record.task_id or "",
record.source_ref_id or "",
record.provider or "",
json.dumps(record.context or {}, ensure_ascii=False, sort_keys=True),
]
).lower()
if search_query not in haystack:
continue
visible.append(record)
if len(visible) >= limit:
break
visible = list(reversed(visible))
return {
"mode": "raw",
"line_limit": limit,
"line_count": len(visible),
"events": [_serialize_observability_event(record) for record in visible],
"lines": [
" ".join(
part
for part in [
record.occurred_at.isoformat() if record.occurred_at else "",
record.level.upper(),
record.source,
record.category or "",
record.event or "",
f"fingerprint={record.fingerprint}",
record.message,
]
if part
)
for record in visible
],
}
def _stable_hash(value: str) -> str:
return hashlib.sha1(value.encode("utf-8", errors="replace")).hexdigest()[:16]

View File

@@ -655,6 +655,8 @@ async def test_ingest_earth_client_log_accepts_public_events():
"message": "登陆点加载失败: 登陆点接口返回 HTTP 500",
"category": "startup-load",
"module": "layer-startup",
"fingerprint": "client-test",
"occurrence_count": 3,
},
)
assert response.status_code == 200
@@ -668,6 +670,8 @@ async def test_ingest_earth_client_log_accepts_public_events():
assert persisted_kwargs["event"] == "earth.client.runtime_log"
assert persisted_kwargs["category"] == "startup-load"
assert persisted_kwargs["level"] == "error"
assert persisted_kwargs["fingerprint"] == "client-test"
assert persisted_kwargs["occurrence_count"] == 3
finally:
app.dependency_overrides.clear()
@@ -705,6 +709,59 @@ async def test_ingest_admin_client_log_accepts_public_events():
app.dependency_overrides.clear()
@pytest.mark.asyncio
async def test_ingest_service_log_requires_configured_token(monkeypatch):
monkeypatch.setattr(settings, "OBSERVABILITY_INGEST_TOKEN", "")
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
response = await client.post(
"/api/v1/system/logs/service",
json={"message": "AI provider failed"},
headers={"X-Planet-Observability-Token": "secret"},
)
assert response.status_code == 503
@pytest.mark.asyncio
async def test_ingest_service_log_accepts_internal_token(monkeypatch):
monkeypatch.setattr(settings, "OBSERVABILITY_INGEST_TOKEN", "service-secret")
transport = ASGITransport(app=app)
with patch("app.api.v1.system_control.record_system_log", new_callable=AsyncMock) as mock_record_system_log:
async with AsyncClient(transport=transport, base_url="http://test") as client:
response = await client.post(
"/api/v1/system/logs/service",
json={
"source": "ai-provider",
"service": "ai-provider",
"module": "provider",
"category": "connectivity",
"event": "ai.provider.test.failed",
"level": "error",
"message": "Provider connectivity failed",
"fingerprint": "ai-provider-test",
"occurrence_count": 4,
"provider": "minimax",
"trace_id": "trace-123",
"context": {"status_code": 502},
},
headers={"Authorization": "Bearer service-secret"},
)
assert response.status_code == 200
data = response.json()
assert data["accepted"] is True
assert data["source_id"] == "ai-provider"
mock_record_system_log.assert_awaited_once()
persisted_kwargs = mock_record_system_log.await_args.kwargs
assert persisted_kwargs["event"] == "ai.provider.test.failed"
assert persisted_kwargs["fingerprint"] == "ai-provider-test"
assert persisted_kwargs["occurrence_count"] == 4
assert persisted_kwargs["context"]["provider"] == "minimax"
assert persisted_kwargs["context"]["trace_id"] == "trace-123"
assert persisted_kwargs["context"]["status_code"] == 502
@pytest.mark.asyncio
async def test_earth_layer_cache_status_requires_super_admin(auth_headers, monkeypatch):
def override_get_current_user():

View File

@@ -74,6 +74,6 @@ def test_interactable_cache_invalidation_clears_layer_and_all(monkeypatch):
assert deleted == 2
assert patterns == [
"earth:layer:v1:interactables:layer:places*",
"earth:layer:v1:interactables:layer:all*",
"earth:layer:v1:interactables:interactable_layer:places*",
"earth:layer:v1:interactables:interactable_layer:all*",
]

View File

@@ -4,14 +4,20 @@ from types import SimpleNamespace
import pytest
from app.services.earth_news import (
NewsFeedEndpoint,
NewsFeedSource,
NewsTargetLocation,
ParsedNewsItem,
apply_news_classification,
default_earth_news_sources_payload,
normalize_earth_news_sources_payload,
_fetch_source,
_enrich_items_with_target_locations,
_extract_target_location_from_text,
_parse_feed_entries,
_serialize_item,
get_earth_news_payload,
test_news_source_config as run_news_source_config_test,
)
from app.services.earth_news_queue import NewsTargetLocationMessage
from app.services.earth_news_worker import process_target_location_message
@@ -164,6 +170,479 @@ def test_parse_aggregated_rss_splits_publisher_from_title():
assert items[0].source == "Reuters"
def test_parse_chinese_rss_marks_source_language_and_keeps_zh_localization():
source = NewsFeedSource(
id="36kr",
name="36氪",
region="asia-pacific",
feed_url="https://36kr.com/feed",
homepage_url="https://www.36kr.com/",
source_tags=("china", "business_news"),
default_category="business",
)
xml = """
<rss>
<channel>
<item>
<title>中国电商平台发布季度增长数据</title>
<description>平台表示,跨境电商订单量同比增长。</description>
<link>https://36kr.com/p/example</link>
</item>
</channel>
</rss>
"""
items = _parse_feed_entries(xml, source)
payload_zh = _serialize_item(items[0], active_region="global", locale="zh-CN")
payload_en = _serialize_item(items[0], active_region="global", locale="en-US")
assert items[0].content_language == "zh-CN"
assert items[0].localizations["zh-CN"]["title"] == "中国电商平台发布季度增长数据"
assert payload_zh["display_title"] == "中国电商平台发布季度增长数据"
assert payload_en["display_title"] == "中国电商平台发布季度增长数据"
def test_default_news_sources_include_business_and_ecommerce_sources():
payload = default_earth_news_sources_payload()
sources_by_id = {source["id"]: source for source in payload["sources"]}
source_ids = {source["id"] for source in payload["sources"]}
category_keys = {category["key"] for category in payload["categories"]}
tag_keys = {tag["key"] for tag in payload["source_tags"]}
assert "cnbc-business" in source_ids
assert "36kr" in source_ids
assert "techcrunch" in source_ids
assert "retaildive" in source_ids
assert "prnewswire-retail" in source_ids
assert "google-news" in source_ids
assert "global-scan" not in source_ids
assert "google-americas" not in source_ids
assert "google-europe" not in source_ids
assert "google-mea" not in source_ids
assert "google-apac" not in source_ids
assert "businesswire-ecommerce" in source_ids
assert "us-census-ecommerce" in source_ids
assert "mofcom-data" in source_ids
assert "stats-china-online-retail" in source_ids
assert "ebrun" in source_ids
assert sources_by_id["36kr"]["source_type"] == "rss"
assert sources_by_id["36kr"]["homepage_url"] == "https://www.36kr.com/"
assert sources_by_id["36kr"]["feed_directory_url"] == "https://www.36kr.com/rss-center"
kr_feeds = {feed["id"]: feed for feed in sources_by_id["36kr"]["feeds"]}
assert set(kr_feeds) == {"feed", "article", "newsflash", "moment"}
assert kr_feeds["feed"]["url"] == "https://36kr.com/feed"
assert kr_feeds["article"]["url"] == "https://36kr.com/feed-article"
assert kr_feeds["newsflash"]["url"] == "https://36kr.com/feed-newsflash"
assert kr_feeds["moment"]["url"] == "https://36kr.com/feed-moment"
assert all(feed["enabled"] is True for feed in kr_feeds.values())
assert all(feed["default_category"] == "business" for feed in kr_feeds.values())
assert "https://36kr.com/feed-article" in sources_by_id["36kr"]["feed_urls"]
assert "https://36kr.com/feed-newsflash" in sources_by_id["36kr"]["feed_urls"]
assert "https://36kr.com/feed-moment" in sources_by_id["36kr"]["feed_urls"]
assert sources_by_id["ebrun"]["source_type"] == "rss"
assert sources_by_id["ebrun"]["homepage_url"] == "https://www.ebrun.com/"
assert sources_by_id["ebrun"]["feed_directory_url"] == "https://www.ebrun.com/rss/"
ebrun_feeds = {feed["id"]: feed for feed in sources_by_id["ebrun"]["feeds"]}
assert {"b2c", "b2b", "retail", "o2o", "service", "data", "policy"}.issubset(ebrun_feeds)
assert all(feed["enabled"] is True for feed in ebrun_feeds.values())
assert all(feed["default_category"] == "ecommerce" for feed in ebrun_feeds.values())
assert "https://www.ebrun.com/rss/news_b2c.xml" in sources_by_id["ebrun"]["feed_urls"]
assert "https://www.ebrun.com/rss/news_retail.xml" in sources_by_id["ebrun"]["feed_urls"]
assert sources_by_id["businesswire-ecommerce"]["source_type"] == "reference"
assert sources_by_id["businesswire-ecommerce"]["enabled"] is False
assert sources_by_id["google-news"]["source_type"] == "aggregated"
assert sources_by_id["google-news"]["homepage_url"] == "https://news.google.com/"
assert sources_by_id["google-news"]["feed_directory_url"] == "https://news.google.com/rss"
google_feeds = {feed["id"]: feed for feed in sources_by_id["google-news"]["feeds"]}
assert set(google_feeds) == {"world", "americas", "europe", "middle-east-africa", "asia-pacific"}
assert all(feed["type"] == "aggregated" for feed in google_feeds.values())
assert all(feed["enabled"] is True for feed in google_feeds.values())
assert google_feeds["world"]["region"] == "global"
assert google_feeds["europe"]["region"] == "europe"
assert sources_by_id["stats-china-online-retail"]["source_type"] == "rss"
assert sources_by_id["stats-china-online-retail"]["enabled"] is True
assert "https://www.stats.gov.cn/sj/zxfb/rss.xml" in sources_by_id["stats-china-online-retail"]["feed_urls"]
assert {"business", "ecommerce", "finance"}.issubset(category_keys)
assert {"official_data", "business_news", "ecommerce", "press_release", "finance", "logistics"}.issubset(tag_keys)
def test_default_enabled_fetchable_sources_have_explicit_types_and_urls():
payload = default_earth_news_sources_payload()
for source in payload["sources"]:
source_type = source["source_type"]
assert source_type in {"rss", "atom", "aggregated", "reference"}
if source_type == "reference":
assert source["enabled"] is False
assert source["feeds"] == []
continue
if source["enabled"]:
assert source["feed_url"]
assert source["feed_urls"]
assert source["feeds"]
assert any(feed["enabled"] for feed in source["feeds"])
for feed in source["feeds"]:
assert feed["url"] != source["homepage_url"]
assert feed["url"] != source.get("feed_directory_url", "")
def test_legacy_news_source_urls_migrate_to_feed_children():
payload = normalize_earth_news_sources_payload(
{
"sources": [
{
"id": "legacy-source",
"name": "Legacy Source",
"region": "global",
"source_type": "rss",
"feed_urls": ["https://example.com/a.xml", "https://example.com/b.xml"],
"default_category": "business",
}
]
}
)
source = payload["sources"][0]
assert source["feed_urls"] == ["https://example.com/a.xml", "https://example.com/b.xml"]
assert [feed["url"] for feed in source["feeds"]] == ["https://example.com/a.xml", "https://example.com/b.xml"]
assert [feed["id"] for feed in source["feeds"]] == ["feed-1", "feed-2"]
assert all(feed["default_category"] == "business" for feed in source["feeds"])
def test_builtin_news_source_legacy_directory_url_is_repaired():
payload = normalize_earth_news_sources_payload(
{
"sources": [
{
"id": "36kr",
"name": "36氪",
"region": "asia-pacific",
"source_type": "rss",
"homepage_url": "https://www.36kr.com/",
"feed_url": "https://www.36kr.com/rss-center",
"feed_urls": ["https://www.36kr.com/rss-center"],
"feeds": [
{
"id": "feed-1",
"name": "36氪",
"url": "https://www.36kr.com/rss-center",
"type": "rss",
"enabled": True,
"default_category": "business",
}
],
"default_category": "business",
}
]
}
)
source = payload["sources"][0]
feed_urls = {feed["url"] for feed in source["feeds"]}
assert source["homepage_url"] == "https://www.36kr.com/"
assert source["feed_directory_url"] == "https://www.36kr.com/rss-center"
assert "https://www.36kr.com/rss-center" not in feed_urls
assert {
"https://36kr.com/feed",
"https://36kr.com/feed-article",
"https://36kr.com/feed-newsflash",
"https://36kr.com/feed-moment",
}.issubset(feed_urls)
def test_builtin_news_source_without_feed_children_gets_explicit_defaults():
payload = normalize_earth_news_sources_payload(
{
"sources": [
{
"id": "ebrun",
"name": "亿邦动力",
"region": "asia-pacific",
"source_type": "rss",
"homepage_url": "https://www.ebrun.com/",
"feed_url": "https://www.ebrun.com/rss/news_b2c.xml",
"feed_urls": ["https://www.ebrun.com/rss/news_b2c.xml"],
"default_category": "ecommerce",
}
]
}
)
source = payload["sources"][0]
feed_urls = {feed["url"] for feed in source["feeds"]}
assert source["feed_directory_url"] == "https://www.ebrun.com/rss/"
assert "https://www.ebrun.com/rss/" not in feed_urls
assert {
"https://www.ebrun.com/rss/news_b2c.xml",
"https://www.ebrun.com/rss/news_b2b.xml",
"https://www.ebrun.com/rss/news_retail.xml",
"https://www.ebrun.com/rss/news_o2o.xml",
"https://www.ebrun.com/rss/news_service.xml",
"https://www.ebrun.com/rss/news_data.xml",
"https://www.ebrun.com/rss/news_policy.xml",
}.issubset(feed_urls)
def test_builtin_fetchable_source_saved_as_reference_is_repaired():
payload = normalize_earth_news_sources_payload(
{
"sources": [
{
"id": "stats-china-online-retail",
"name": "国家统计局数据发布",
"region": "asia-pacific",
"source_type": "reference",
"enabled": False,
"homepage_url": "https://www.stats.gov.cn/sj/zxfb/",
"feed_url": "https://www.stats.gov.cn/sj/zxfb/",
"default_category": "ecommerce",
}
]
}
)
source = payload["sources"][0]
assert source["source_type"] == "rss"
assert source["enabled"] is True
assert source["priority"] == 19
assert source["source_tags"] == ["official_data", "ecommerce", "retail", "china"]
assert source["default_category"] == "ecommerce"
assert source["importance_weight"] == 36
assert source["feed_directory_url"] == ""
assert source["feeds"] == [
{
"id": "release",
"name": "数据发布",
"url": "https://www.stats.gov.cn/sj/zxfb/rss.xml",
"type": "rss",
"region": "asia-pacific",
"enabled": True,
"default_category": "ecommerce",
"tags": [],
"priority": 1,
}
]
def test_legacy_google_sources_merge_into_google_news_source():
payload = normalize_earth_news_sources_payload(
{
"sources": [
{
"id": "global-scan",
"name": "Global Monitor / World",
"region": "global",
"source_type": "aggregated",
"feed_url": "https://news.google.com/rss/search?q=world",
"homepage_url": "https://news.google.com/",
},
{
"id": "google-europe",
"name": "Global Monitor / Europe",
"region": "europe",
"source_type": "aggregated",
"feed_url": "https://news.google.com/rss/search?q=europe",
"homepage_url": "https://news.google.com/",
},
]
}
)
sources_by_id = {source["id"]: source for source in payload["sources"]}
assert "global-scan" not in sources_by_id
assert "google-europe" not in sources_by_id
assert "google-news" in sources_by_id
assert {feed["id"] for feed in sources_by_id["google-news"]["feeds"]} == {
"world",
"americas",
"europe",
"middle-east-africa",
"asia-pacific",
}
def test_feed_child_default_category_overrides_source_default():
source = NewsFeedSource(
id="multi-feed",
name="Multi Feed",
region="global",
feed_url="https://example.com/source.xml",
homepage_url="https://example.com",
default_category="business",
)
feed = NewsFeedEndpoint(
id="ecommerce-feed",
name="Ecommerce Feed",
url="https://example.com/ecommerce.xml",
default_category="ecommerce",
)
xml = """
<rss>
<channel>
<item>
<title>Quarterly results released</title>
<description>Company update.</description>
<link>https://example.com/results</link>
</item>
</channel>
</rss>
"""
items = _parse_feed_entries(xml, source, feed=feed)
assert items[0].feed_id == "ecommerce-feed"
assert items[0].feed_name == "Ecommerce Feed"
assert items[0].feed_default_category == "ecommerce"
assert items[0].category == "ecommerce"
@pytest.mark.asyncio
async def test_fetch_source_only_requests_enabled_feed_children(monkeypatch):
source = NewsFeedSource(
id="multi-feed",
name="Multi Feed",
region="global",
feed_url="https://example.com/source.xml",
homepage_url="https://example.com",
feeds=(
NewsFeedEndpoint(id="enabled", name="Enabled", url="https://example.com/enabled.xml", enabled=True),
NewsFeedEndpoint(id="disabled", name="Disabled", url="https://example.com/disabled.xml", enabled=False),
),
)
calls = []
async def fake_fetch_single(_client, feed_source, feed, *, config_payload=None):
calls.append(feed.id)
item = ParsedNewsItem(
id=f"{feed_source.id}:{feed.id}:1",
title="Fetched story",
summary="Fetched summary",
url=f"https://example.com/{feed.id}",
source="Example",
feed_name=feed.name,
feed_region="global",
homepage_url="https://example.com",
published_at=None,
feed_id=feed.id,
)
return feed_source, [item], None, {"source_id": feed_source.id, "feed_id": feed.id, "ok": True, "status": "ok", "item_count": 1, "count": 1}
monkeypatch.setattr("app.services.earth_news._fetch_single_feed_url", fake_fetch_single)
source_result, items, error, health = await _fetch_source(object(), source)
assert source_result.id == "multi-feed"
assert calls == ["enabled"]
assert error is None
assert [item.feed_id for item in items] == ["enabled"]
assert health["ok"] is True
assert [result["feed_id"] for result in health["feed_results"]] == ["enabled"]
@pytest.mark.asyncio
async def test_fetch_source_filters_google_feed_children_by_active_region(monkeypatch):
source = NewsFeedSource(
id="google-news",
name="Google News",
region="global",
feed_url="https://news.google.com/rss",
homepage_url="https://news.google.com/",
source_type="aggregated",
feeds=(
NewsFeedEndpoint(id="world", name="全球", url="https://example.com/world.xml", type="aggregated", region="global"),
NewsFeedEndpoint(id="europe", name="欧洲", url="https://example.com/europe.xml", type="aggregated", region="europe"),
NewsFeedEndpoint(id="americas", name="美洲", url="https://example.com/americas.xml", type="aggregated", region="americas"),
),
)
calls = []
async def fake_fetch_single(_client, feed_source, feed, *, config_payload=None):
calls.append(feed.id)
item = ParsedNewsItem(
id=f"{feed_source.id}:{feed.id}:1",
title=f"{feed.name} headline",
summary="Fetched summary",
url=f"https://example.com/{feed.id}",
source="Example",
feed_name=feed.name,
feed_region=feed.region,
homepage_url="https://example.com",
published_at=None,
feed_id=feed.id,
)
return feed_source, [item], None, {"source_id": feed_source.id, "feed_id": feed.id, "ok": True, "status": "ok", "item_count": 1, "count": 1}
monkeypatch.setattr("app.services.earth_news._fetch_single_feed_url", fake_fetch_single)
_source_result, items, error, health = await _fetch_source(object(), source, active_region="europe")
assert error is None
assert calls == ["world", "europe"]
assert [item.feed_region for item in items] == ["global", "europe"]
assert [result["feed_id"] for result in health["feed_results"]] == ["world", "europe"]
def test_parse_rdf_rss_items_with_namespaces():
source = NewsFeedSource(
id="dw-top",
name="DW Top Stories",
region="europe",
feed_url="https://rss.dw.com/rdf/rss-en-top",
homepage_url="https://www.dw.com/en/top-stories/s-9097",
)
xml = """
<rdf:RDF xmlns:rdf="http://www.w3.org/1999/02/22-rdf-syntax-ns#"
xmlns="http://purl.org/rss/1.0/">
<item rdf:about="https://example.com/dw">
<title>German retail sales rise</title>
<link>https://example.com/dw</link>
<description>Retail summary</description>
</item>
</rdf:RDF>
"""
items = _parse_feed_entries(xml, source)
assert len(items) == 1
assert items[0].title == "German retail sales rise"
def test_news_classification_marks_ecommerce_and_importance():
source = NewsFeedSource(
id="ebrun",
name="亿邦动力",
region="asia-pacific",
feed_url="https://www.ebrun.com/rss/",
homepage_url="https://www.ebrun.com/",
source_tags=("business_news", "ecommerce", "china"),
default_category="ecommerce",
importance_weight=14,
)
item = ParsedNewsItem(
id="ebrun:test",
title="跨境电商平台 GMV 同比增长,物流履约效率提升",
summary="订单量和网上零售额继续增长。",
url="https://example.com/ecommerce",
source="亿邦动力",
feed_name="亿邦动力",
feed_region="asia-pacific",
homepage_url="https://www.ebrun.com/",
published_at=None,
)
apply_news_classification(item, source)
assert item.category == "ecommerce"
assert "cross_border_ecommerce" in item.item_tags
assert "logistics_fulfillment" in item.item_tags
assert item.importance_level in {"high", "critical"}
assert "命中电商数据指标" in item.importance_reasons
@pytest.mark.asyncio
async def test_enrich_items_with_target_locations_uses_ai_and_geocode(monkeypatch):
item = ParsedNewsItem(
@@ -338,8 +817,8 @@ async def test_earth_news_payload_returns_anchor_items_and_enqueues_location_job
published_at=datetime(2026, 5, 15, 3, 0, tzinfo=UTC),
)
async def fake_fetch_source(_client, feed_source):
return feed_source, [item], None
async def fake_fetch_source(_client, feed_source, **_kwargs):
return feed_source, [item], None, {"source_id": feed_source.id, "ok": True, "status": "ok", "count": 1}
async def fake_get_cached_target_location_patch(_item_id):
return None
@@ -402,7 +881,7 @@ async def test_earth_news_payload_uses_fresh_database_items_without_rss(monkeypa
async def fake_get_earth_news_freshness(_db, *, active_region):
return 12, datetime.now(UTC)
async def fake_list_earth_news_items(_db, *, active_region, limit):
async def fake_list_earth_news_items(_db, *, active_region, limit, categories=None):
assert limit == 12
return [item]
@@ -452,10 +931,10 @@ async def test_earth_news_payload_keeps_current_items_and_all_cruise_items(monke
async def fake_get_earth_news_freshness(_db, *, active_region):
return 12, datetime.now(UTC)
async def fake_list_earth_news_items(_db, *, active_region, limit):
async def fake_list_earth_news_items(_db, *, active_region, limit, categories=None):
return [current_item]
async def fake_list_earth_news_cruise_items(_db, *, limit):
async def fake_list_earth_news_cruise_items(_db, *, limit, categories=None):
return [current_item, cruise_item]
async def fake_enqueue_target_location_job(_payload, **_kwargs):
@@ -477,6 +956,89 @@ async def test_earth_news_payload_keeps_current_items_and_all_cruise_items(monke
assert payload["cruise_items"][1]["region"] == "asia-pacific"
@pytest.mark.asyncio
async def test_earth_news_payload_passes_region_and_category_filters_to_store(monkeypatch):
class FakeDb:
execute = object()
captured = {}
item = ParsedNewsItem(
id="db:business",
title="Business story",
summary="Business summary",
url="https://example.com/business",
source="Stored Source",
feed_name="Stored Feed",
feed_region="europe",
homepage_url="https://example.com",
published_at=datetime(2026, 5, 15, 3, 0, tzinfo=UTC),
category="business",
)
async def fake_get_earth_news_freshness(_db, *, active_region):
captured["freshness_region"] = active_region
return 12, datetime.now(UTC)
async def fake_list_earth_news_items(_db, *, active_region, limit, categories=None, source_ids=None):
captured["items_region"] = active_region
captured["items_categories"] = categories
captured["items_source_ids"] = source_ids
return [item]
async def fake_list_earth_news_cruise_items(_db, *, limit, categories=None, source_ids=None):
captured["cruise_categories"] = categories
captured["cruise_source_ids"] = source_ids
return [item]
async def fake_enqueue_target_location_job(_payload, **_kwargs):
return True
monkeypatch.setattr("app.services.earth_news_store.get_earth_news_freshness", fake_get_earth_news_freshness)
monkeypatch.setattr("app.services.earth_news_store.list_earth_news_items", fake_list_earth_news_items)
monkeypatch.setattr("app.services.earth_news_store.list_earth_news_cruise_items", fake_list_earth_news_cruise_items)
monkeypatch.setattr("app.services.earth_news_queue.enqueue_target_location_job", fake_enqueue_target_location_job)
monkeypatch.setattr("app.services.earth_news._fetch_rss_items_for_sources", lambda _sources: (_ for _ in ()).throw(AssertionError("fresh database items should not fetch RSS")))
payload = await get_earth_news_payload(
lat=35.0,
lon=-100.0,
region="europe",
categories={"business", "ecommerce"},
db=FakeDb(),
)
assert captured["freshness_region"] == "europe"
assert captured["items_region"] == "europe"
assert captured["items_categories"] == {"business", "ecommerce"}
assert captured["items_source_ids"] is None
assert captured["cruise_categories"] == {"business", "ecommerce"}
assert captured["cruise_source_ids"] is None
assert payload["filters"] == {
"region": "europe",
"categories": ["business", "ecommerce"],
"sources": [],
"limit": 12,
"locale": "zh-CN",
}
assert payload["items"][0]["category"] == "business"
@pytest.mark.asyncio
async def test_news_source_test_treats_type_reference_as_non_fetching():
result = await run_news_source_config_test(
{
"id": "reference-only",
"name": "Reference Only",
"type": "reference",
"feed_url": "https://example.com",
}
)
assert result["ok"] is False
assert result["health"]["status"] == "reference"
assert "不参与 RSS/Atom 抓取" in result["error"]
@pytest.mark.asyncio
async def test_earth_news_payload_initializes_empty_database_from_rss(monkeypatch):
db = object()
@@ -512,14 +1074,14 @@ async def test_earth_news_payload_initializes_empty_database_from_rss(monkeypatc
async def fake_get_earth_news_freshness(_db, *, active_region):
return 0, None
async def fake_fetch_rss_items_for_sources(_sources):
return [item], []
async def fake_fetch_rss_items_for_sources(_sources, **_kwargs):
return [item], [], {"test-feed": {"source_id": "test-feed", "ok": True, "status": "ok", "count": 1}}
async def fake_upsert_earth_news_items(_db, items):
upserted.extend(items)
return len(items)
async def fake_list_earth_news_items(_db, *, active_region, limit):
async def fake_list_earth_news_items(_db, *, active_region, limit, categories=None):
return [item]
async def fake_enqueue_target_location_job(payload, **_kwargs):
@@ -568,14 +1130,14 @@ async def test_earth_news_payload_supplements_stale_database_items(monkeypatch):
async def fake_get_earth_news_freshness(_db, *, active_region):
return 12, datetime(2026, 5, 14, 3, 0, tzinfo=UTC)
async def fake_fetch_rss_items_for_sources(_sources):
async def fake_fetch_rss_items_for_sources(_sources, **_kwargs):
fetched.append(True)
return [old_item], []
return [old_item], [], {"stored": {"source_id": "stored", "ok": True, "status": "ok", "count": 1}}
async def fake_upsert_earth_news_items(_db, items):
return len(items)
async def fake_list_earth_news_items(_db, *, active_region, limit):
async def fake_list_earth_news_items(_db, *, active_region, limit, categories=None):
return [old_item]
async def fake_enqueue_target_location_job(_payload, **_kwargs):
@@ -630,8 +1192,8 @@ async def test_earth_news_payload_merges_cached_location_patch(monkeypatch):
},
}
async def fake_fetch_source(_client, feed_source):
return feed_source, [item], None
async def fake_fetch_source(_client, feed_source, **_kwargs):
return feed_source, [item], None, {"source_id": feed_source.id, "ok": True, "status": "ok", "count": 1}
async def fake_get_cached_target_location_patch(_item_id):
return cached_patch
@@ -697,8 +1259,8 @@ async def test_earth_news_payload_requeues_cached_failed_localization(monkeypatc
}
enqueued = []
async def fake_fetch_source(_client, feed_source):
return feed_source, [item], None
async def fake_fetch_source(_client, feed_source, **_kwargs):
return feed_source, [item], None, {"source_id": feed_source.id, "ok": True, "status": "ok", "count": 1}
async def fake_get_cached_target_location_patch(_item_id):
return cached_patch

View File

@@ -9,6 +9,8 @@ import pytest
from app.core.logging import PlanetContextFilter, PlanetFormatter, get_logger
from app.core.request_context import set_request_id
from app.services import business_logs
from app.services import persistent_logs
from app.models.system_log import ObservabilityEvent, ObservabilityEventGroup
def _capture_output(callback):
@@ -98,6 +100,86 @@ def test_business_context_redacts_nested_sensitive_values():
assert context["nested"]["safe"] == "visible"
def test_observability_fingerprint_normalizes_hls_fragments():
first = persistent_logs.build_observability_fingerprint(
source="earth-client",
service="earth",
module="tv",
category="hls-proxy",
event="hls.fragment.failed",
message="HLS 分片加载失败: index_5_9086220.ts?m=1725933270",
context={"status_code": 502},
)
second = persistent_logs.build_observability_fingerprint(
source="earth-client",
service="earth",
module="tv",
category="hls-proxy",
event="hls.fragment.failed",
message="HLS 分片加载失败: index_5_9086361.ts?m=1725934270",
context={"status_code": 502},
)
assert first == second
@pytest.mark.asyncio
async def test_record_observability_event_updates_group_count(monkeypatch):
events: list[ObservabilityEvent] = []
groups: dict[str, ObservabilityEventGroup] = {}
class FakeSession:
async def __aenter__(self):
return self
async def __aexit__(self, exc_type, exc, tb):
return False
def add(self, item):
if isinstance(item, ObservabilityEvent):
events.append(item)
elif isinstance(item, ObservabilityEventGroup):
groups[item.fingerprint] = item
async def get(self, model, key):
if model is ObservabilityEventGroup:
return groups.get(key)
return None
async def commit(self):
return None
monkeypatch.setattr(persistent_logs, "async_session_factory", lambda: FakeSession())
await persistent_logs.record_observability_event(
source="earth-client",
level="error",
service="earth",
module="tv",
category="hls-proxy",
event="hls.fragment.failed",
message="HLS 分片加载失败: index_5_9086220.ts?m=1725933270",
context={"status_code": 502},
occurrence_count=2,
)
await persistent_logs.record_observability_event(
source="earth-client",
level="error",
service="earth",
module="tv",
category="hls-proxy",
event="hls.fragment.failed",
message="HLS 分片加载失败: index_5_9086361.ts?m=1725934270",
context={"status_code": 502},
occurrence_count=1,
)
assert len(events) == 2
assert len(groups) == 1
group = next(iter(groups.values()))
assert group.count == 3
@pytest.mark.asyncio
async def test_emit_business_log_persists_sanitized_system_event(monkeypatch):
events = []

View File

@@ -0,0 +1,35 @@
from app.api.v1.tv import _rewrite_hls_uri_attributes, _should_strip_hls_metadata_line
def test_rewrite_hls_uri_attributes_rewrites_subtitle_manifest_url():
line = '#EXT-X-MEDIA:TYPE=SUBTITLES,GROUP-ID="subs",NAME="English",URI="index_3_0.m3u8"'
rewritten = _rewrite_hls_uri_attributes(
line,
base_url="https://example.com/live/master.m3u8",
)
assert 'URI="/api/v1/tv/proxy?url=https%3A%2F%2Fexample.com%2Flive%2Findex_3_0.m3u8"' in rewritten
def test_rewrite_hls_uri_attributes_rewrites_absolute_uri():
line = '#EXT-X-I-FRAME-STREAM-INF:BANDWIDTH=1234,URI="https://cdn.example.com/live/iframe.m3u8"'
rewritten = _rewrite_hls_uri_attributes(
line,
base_url="https://example.com/live/master.m3u8",
)
assert 'URI="/api/v1/tv/proxy?url=https%3A%2F%2Fcdn.example.com%2Flive%2Fiframe.m3u8"' in rewritten
def test_strip_hls_subtitle_media_metadata():
line = '#EXT-X-MEDIA:TYPE=SUBTITLES,GROUP-ID="subs",NAME="English",URI="index_3_0.m3u8"'
assert _should_strip_hls_metadata_line(line) is True
def test_keep_hls_audio_media_metadata():
line = '#EXT-X-MEDIA:TYPE=AUDIO,GROUP-ID="audio",NAME="English",URI="audio.m3u8"'
assert _should_strip_hls_metadata_line(line) is False