Add Profile/Channel entities with CRUD and remove legacy parser builder.

Introduce reusable ParserProfile and ParseChannel, pair jobs with enqueue flatten, Events admin CRUD, and drop inline/legacy parser-builder UI and aliases.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
2026-08-16 14:51:03 +03:00
co-authored by Cursor
parent e196de7320
commit 71f8cf169e
30 changed files with 2157 additions and 954 deletions
+3 -2
View File
@@ -13,7 +13,7 @@ for _candidate in (_HERE.parent, *_HERE.parents):
break
from .database import Base, engine, get_db
from .routers import admin, internal, map, objects, parser_builder, v1
from .routers import admin, internal, map, objects, parse_channels, parser_profiles, v1
from .seed import seed_objects, seed_test_consumer
from .services.migrations import migrate_schema
from .services.scheduler import start_scheduler
@@ -50,5 +50,6 @@ app.include_router(objects.router)
app.include_router(map.router)
app.include_router(internal.router)
app.include_router(admin.router)
app.include_router(parser_builder.router)
app.include_router(parser_profiles.router)
app.include_router(parse_channels.router)
app.include_router(v1.router)
+45
View File
@@ -93,12 +93,54 @@ class Event(Base):
)
class ParserProfile(Base):
__tablename__ = "parser_profiles"
id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True)
name: Mapped[str] = mapped_column(String(255), nullable=False)
sample_post: Mapped[str] = mapped_column(Text, default="")
heuristic_profile: Mapped[dict | None] = mapped_column(JSON, nullable=True)
status: Mapped[str] = mapped_column(String(50), default="draft", index=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True),
default=lambda: datetime.now(timezone.utc),
)
jobs: Mapped[list["ParseJob"]] = relationship(back_populates="profile")
class ParseChannel(Base):
__tablename__ = "parse_channels"
id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True)
name: Mapped[str] = mapped_column(String(255), nullable=False)
source_type: Mapped[str] = mapped_column(String(50), nullable=False, default="telegram")
channel: Mapped[str] = mapped_column(String(255), nullable=False)
is_active: Mapped[bool] = mapped_column(Boolean, default=True, nullable=False)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True),
default=lambda: datetime.now(timezone.utc),
)
jobs: Mapped[list["ParseJob"]] = relationship(back_populates="channel")
class ParseJob(Base):
__tablename__ = "parse_jobs"
id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True)
source_type: Mapped[str] = mapped_column(String(50), nullable=False)
source_config: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict)
profile_id: Mapped[int | None] = mapped_column(
ForeignKey("parser_profiles.id", ondelete="SET NULL"),
nullable=True,
index=True,
)
channel_id: Mapped[int | None] = mapped_column(
ForeignKey("parse_channels.id", ondelete="SET NULL"),
nullable=True,
index=True,
)
schedule: Mapped[str | None] = mapped_column(String(100), nullable=True)
interval_seconds: Mapped[int] = mapped_column(Integer, default=3600, nullable=False)
is_active: Mapped[bool] = mapped_column(Boolean, default=True, nullable=False)
@@ -110,6 +152,9 @@ class ParseJob(Base):
default=lambda: datetime.now(timezone.utc),
)
profile: Mapped["ParserProfile | None"] = relationship(back_populates="jobs")
channel: Mapped["ParseChannel | None"] = relationship(back_populates="jobs")
class Consumer(Base):
__tablename__ = "consumers"
+217 -19
View File
@@ -4,14 +4,16 @@ from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy.orm import Session
from ..database import get_db
from ..models import Consumer, Event, ParseJob
from ..models import Consumer, Event, MapObject, ParseChannel, ParseJob, ParserProfile
from ..schemas import (
AnalyticsSummary,
ConsumerCreate,
ConsumerRead,
ConsumerUpdate,
EventCreate,
EventListResponse,
EventRead,
EventUpdate,
ParseJobCreate,
ParseJobRead,
ParseJobUpdate,
@@ -30,9 +32,10 @@ from ..services.filtering import (
create_consumer,
list_consumers,
rotate_consumer_key,
sync_event_to_map_object,
update_consumer,
)
from ..services.jobs import enqueue_job
from ..services.jobs import enqueue_parse_job, flatten_pair_config
router = APIRouter(prefix="/admin", tags=["admin"])
@@ -48,28 +51,85 @@ def _validated_source_config(source_type: str, source_config: dict) -> dict:
raise HTTPException(status_code=400, detail=str(exc)) from exc
def _job_to_read(job: ParseJob) -> ParseJobRead:
return ParseJobRead(
id=job.id,
source_type=job.source_type,
source_config=job.source_config or {},
profile_id=job.profile_id,
channel_id=job.channel_id,
schedule=job.schedule,
interval_seconds=job.interval_seconds,
is_active=job.is_active,
status=job.status,
last_run_at=job.last_run_at,
last_error=job.last_error,
created_at=job.created_at,
profile_name=job.profile.name if job.profile else None,
channel_name=job.channel.name if job.channel else None,
channel_handle=job.channel.channel if job.channel else None,
)
@router.post("/jobs", response_model=ParseJobRead, status_code=201)
def create_parse_job(payload: ParseJobCreate, db: Session = Depends(get_db)):
source_config = _validated_source_config(payload.source_type, payload.source_config)
job = ParseJob(
source_type=payload.source_type,
source_config=source_config,
schedule=payload.schedule,
interval_seconds=payload.interval_seconds,
is_active=payload.is_active,
status="queued",
)
pair_mode = payload.profile_id is not None or payload.channel_id is not None
if pair_mode:
if payload.profile_id is None or payload.channel_id is None:
raise HTTPException(
status_code=400,
detail="Both profile_id and channel_id are required to create a pair",
)
profile = db.query(ParserProfile).filter(ParserProfile.id == payload.profile_id).first()
channel = db.query(ParseChannel).filter(ParseChannel.id == payload.channel_id).first()
if not profile:
raise HTTPException(status_code=404, detail="Profile not found")
if not channel:
raise HTTPException(status_code=404, detail="Channel not found")
if not profile.heuristic_profile:
raise HTTPException(status_code=400, detail="Profile has no heuristic_profile")
if not channel.is_active:
raise HTTPException(status_code=400, detail="Channel is inactive")
try:
source_config = flatten_pair_config(channel, profile, limit=payload.limit)
except Exception as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
job = ParseJob(
source_type=channel.source_type or "telegram",
source_config=source_config,
profile_id=profile.id,
channel_id=channel.id,
schedule=payload.schedule,
interval_seconds=payload.interval_seconds,
is_active=payload.is_active,
status="queued",
)
else:
source_config = _validated_source_config(payload.source_type, payload.source_config)
job = ParseJob(
source_type=payload.source_type,
source_config=source_config,
schedule=payload.schedule,
interval_seconds=payload.interval_seconds,
is_active=payload.is_active,
status="queued",
)
db.add(job)
db.commit()
db.refresh(job)
enqueue_job(job.id, job.source_type, job.source_config)
return job
enqueue_parse_job(db, job)
db.refresh(job)
return _job_to_read(job)
@router.get("/jobs", response_model=list[ParseJobRead])
def list_parse_jobs(db: Session = Depends(get_db)):
return db.query(ParseJob).order_by(ParseJob.id.desc()).all()
jobs = db.query(ParseJob).order_by(ParseJob.id.desc()).all()
return [_job_to_read(job) for job in jobs]
@router.get("/jobs/{job_id}", response_model=ParseJobRead)
@@ -77,7 +137,7 @@ def get_parse_job(job_id: int, db: Session = Depends(get_db)):
job = db.query(ParseJob).filter(ParseJob.id == job_id).first()
if not job:
raise HTTPException(status_code=404, detail="Job not found")
return job
return _job_to_read(job)
@router.post("/jobs/{job_id}/retry", response_model=ParseJobRead)
@@ -93,8 +153,12 @@ def retry_parse_job(job_id: int, db: Session = Depends(get_db)):
db.commit()
db.refresh(job)
enqueue_job(job.id, job.source_type, job.source_config)
return job
try:
enqueue_parse_job(db, job)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
db.refresh(job)
return _job_to_read(job)
@router.patch("/jobs/{job_id}", response_model=ParseJobRead)
@@ -107,8 +171,36 @@ def update_parse_job(
if not job:
raise HTTPException(status_code=404, detail="Job not found")
if payload.source_config is not None:
if payload.profile_id is not None:
profile = db.query(ParserProfile).filter(ParserProfile.id == payload.profile_id).first()
if not profile:
raise HTTPException(status_code=404, detail="Profile not found")
job.profile_id = profile.id
if payload.channel_id is not None:
channel = db.query(ParseChannel).filter(ParseChannel.id == payload.channel_id).first()
if not channel:
raise HTTPException(status_code=404, detail="Channel not found")
job.channel_id = channel.id
if job.profile_id and job.channel_id:
profile = db.query(ParserProfile).filter(ParserProfile.id == job.profile_id).first()
channel = db.query(ParseChannel).filter(ParseChannel.id == job.channel_id).first()
if profile and channel:
existing_limit = (job.source_config or {}).get("limit", 100)
limit = payload.limit if payload.limit is not None else existing_limit
if not isinstance(limit, int):
limit = 100
try:
job.source_config = flatten_pair_config(channel, profile, limit=limit)
except Exception as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
elif payload.source_config is not None:
job.source_config = _validated_source_config(job.source_type, payload.source_config)
elif payload.limit is not None and isinstance(job.source_config, dict):
cfg = dict(job.source_config)
cfg["limit"] = payload.limit
job.source_config = _validated_source_config(job.source_type, cfg)
if payload.interval_seconds is not None:
job.interval_seconds = payload.interval_seconds
if payload.is_active is not None:
@@ -116,7 +208,7 @@ def update_parse_job(
db.commit()
db.refresh(job)
return job
return _job_to_read(job)
@router.delete("/jobs/{job_id}", status_code=204)
@@ -159,6 +251,112 @@ def list_events(
return EventListResponse(items=items, total=total)
def _ensure_unique_source_url(db: Session, source_url: str, *, exclude_id: int | None = None) -> None:
query = db.query(Event).filter(Event.source_url == source_url)
if exclude_id is not None:
query = query.filter(Event.id != exclude_id)
if query.first():
raise HTTPException(status_code=409, detail="Event with this source_url already exists")
def _sync_or_clear_map_object(db: Session, event: Event) -> None:
if event.latitude is not None and event.longitude is not None:
sync_event_to_map_object(db, event)
return
linked = db.query(MapObject).filter(MapObject.event_id == event.id).first()
if linked:
db.delete(linked)
@router.post("/events", response_model=EventRead, status_code=201)
def create_event(payload: EventCreate, db: Session = Depends(get_db)):
source_url = payload.source_url.strip()
if not source_url:
raise HTTPException(status_code=400, detail="source_url is required")
_ensure_unique_source_url(db, source_url)
if (payload.latitude is None) != (payload.longitude is None):
raise HTTPException(status_code=400, detail="latitude and longitude must be set together")
event = Event(
source_type=payload.source_type.strip() or "manual",
source_url=source_url,
raw_text=payload.raw_text or "",
title=payload.title or "",
description=payload.description or "",
locality=payload.locality or "",
latitude=payload.latitude,
longitude=payload.longitude,
event_date=payload.event_date,
region=payload.region,
topic=payload.topic,
tags=payload.tags,
metadata_=payload.metadata,
)
db.add(event)
db.flush()
_sync_or_clear_map_object(db, event)
db.commit()
db.refresh(event)
return event
@router.get("/events/{event_id}", response_model=EventRead)
def get_event(event_id: int, db: Session = Depends(get_db)):
event = db.query(Event).filter(Event.id == event_id).first()
if not event:
raise HTTPException(status_code=404, detail="Event not found")
return event
@router.patch("/events/{event_id}", response_model=EventRead)
def update_event(event_id: int, payload: EventUpdate, db: Session = Depends(get_db)):
event = db.query(Event).filter(Event.id == event_id).first()
if not event:
raise HTTPException(status_code=404, detail="Event not found")
data = payload.model_dump(exclude_unset=True)
if "source_url" in data and data["source_url"] is not None:
source_url = data["source_url"].strip()
if not source_url:
raise HTTPException(status_code=400, detail="source_url cannot be empty")
_ensure_unique_source_url(db, source_url, exclude_id=event.id)
event.source_url = source_url
data.pop("source_url")
if "metadata" in data:
event.metadata_ = data.pop("metadata")
if "source_type" in data and data["source_type"] is not None:
event.source_type = data.pop("source_type").strip() or event.source_type
for field, value in data.items():
setattr(event, field, value)
lat = event.latitude
lon = event.longitude
if (lat is None) != (lon is None):
raise HTTPException(status_code=400, detail="latitude and longitude must be set together")
_sync_or_clear_map_object(db, event)
db.commit()
db.refresh(event)
return event
@router.delete("/events/{event_id}", status_code=204)
def delete_event(event_id: int, db: Session = Depends(get_db)):
event = db.query(Event).filter(Event.id == event_id).first()
if not event:
raise HTTPException(status_code=404, detail="Event not found")
linked = db.query(MapObject).filter(MapObject.event_id == event.id).all()
for obj in linked:
db.delete(obj)
db.delete(event)
db.commit()
@router.get("/analytics/summary", response_model=AnalyticsSummary)
def analytics_summary(db: Session = Depends(get_db)):
return get_analytics_summary(db)
@@ -8,6 +8,7 @@ from ..deps import verify_internal_token
from ..models import ParseJob
from ..schemas import IngestRequest, IngestResponse, ListenerSubscription
from ..services.ingest import ingest_events
from ..services.jobs import resolve_job_source_config
router = APIRouter(prefix="/internal", tags=["internal"])
@@ -68,14 +69,19 @@ def listener_subscriptions(
)
result: list[ListenerSubscription] = []
for job in jobs:
channel = job.source_config.get("channel") if job.source_config else None
try:
source_config = resolve_job_source_config(db, job)
except ValueError:
continue
channel = source_config.get("channel")
if not channel:
continue
result.append(
ListenerSubscription(
job_id=job.id,
channel=str(channel),
source_config=dict(job.source_config or {}),
source_config=dict(source_config),
)
)
db.commit()
return result
@@ -0,0 +1,98 @@
"""Admin CRUD for parse channels (Telegram handles in MVP)."""
from __future__ import annotations
from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy.orm import Session
from ..database import get_db
from ..models import ParseChannel, ParseJob
from ..schemas import ParseChannelCreate, ParseChannelRead, ParseChannelUpdate
router = APIRouter(prefix="/admin/parse-channels", tags=["parse-channels"])
ALLOWED_SOURCE_TYPES = {"telegram"}
@router.get("", response_model=list[ParseChannelRead])
def list_channels(db: Session = Depends(get_db)):
return db.query(ParseChannel).order_by(ParseChannel.id.desc()).all()
@router.post("", response_model=ParseChannelRead, status_code=201)
def create_channel(payload: ParseChannelCreate, db: Session = Depends(get_db)):
source_type = (payload.source_type or "telegram").strip()
if source_type not in ALLOWED_SOURCE_TYPES:
raise HTTPException(
status_code=400,
detail=f"source_type must be one of: {', '.join(sorted(ALLOWED_SOURCE_TYPES))}",
)
channel = ParseChannel(
name=payload.name.strip(),
source_type=source_type,
channel=payload.channel.strip().lstrip("@"),
is_active=payload.is_active,
)
db.add(channel)
db.commit()
db.refresh(channel)
return channel
@router.get("/{channel_id}", response_model=ParseChannelRead)
def get_channel(channel_id: int, db: Session = Depends(get_db)):
channel = db.query(ParseChannel).filter(ParseChannel.id == channel_id).first()
if not channel:
raise HTTPException(status_code=404, detail="Channel not found")
return channel
@router.patch("/{channel_id}", response_model=ParseChannelRead)
def update_channel(
channel_id: int,
payload: ParseChannelUpdate,
db: Session = Depends(get_db),
):
channel = db.query(ParseChannel).filter(ParseChannel.id == channel_id).first()
if not channel:
raise HTTPException(status_code=404, detail="Channel not found")
if payload.name is not None:
channel.name = payload.name.strip()
if payload.source_type is not None:
source_type = payload.source_type.strip()
if source_type not in ALLOWED_SOURCE_TYPES:
raise HTTPException(
status_code=400,
detail=f"source_type must be one of: {', '.join(sorted(ALLOWED_SOURCE_TYPES))}",
)
channel.source_type = source_type
if payload.channel is not None:
channel.channel = payload.channel.strip().lstrip("@")
if payload.is_active is not None:
channel.is_active = payload.is_active
db.commit()
db.refresh(channel)
return channel
@router.delete("/{channel_id}", status_code=204)
def delete_channel(channel_id: int, db: Session = Depends(get_db)):
channel = db.query(ParseChannel).filter(ParseChannel.id == channel_id).first()
if not channel:
raise HTTPException(status_code=404, detail="Channel not found")
linked = (
db.query(ParseJob)
.filter(ParseJob.channel_id == channel_id, ParseJob.status.in_(("queued", "running")))
.count()
)
if linked:
raise HTTPException(
status_code=409,
detail="Cannot delete channel used by queued/running jobs",
)
db.delete(channel)
db.commit()
@@ -1,108 +0,0 @@
"""Admin API: construct static telegram parsers from a sample post."""
from __future__ import annotations
from typing import Any
from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel, Field
from sqlalchemy.orm import Session
from contracts.heuristic_profile import HeuristicProfile
from contracts.sources import TelegramSourceConfig
from ..database import get_db
from ..models import ParseJob
from ..schemas import ParseJobRead
from ..services.jobs import enqueue_job
from ..services import parser_builder as builder
router = APIRouter(prefix="/admin/parser-builder", tags=["parser-builder"])
class GenerateRequest(BaseModel):
sample_post: str = Field(min_length=1)
class GenerateResponse(BaseModel):
profile: dict[str, Any]
class PreviewRequest(BaseModel):
sample_post: str = Field(min_length=1)
profile: dict[str, Any]
class PreviewResponse(BaseModel):
fields: dict[str, str]
class CreateProfileJobRequest(BaseModel):
channel: str = Field(min_length=1)
profile: dict[str, Any]
sample_post: str | None = None
limit: int = Field(default=100, ge=1, le=1000)
interval_seconds: int = Field(default=3600, ge=60, le=604800)
is_active: bool = True
@router.get("/target-fields")
def target_fields():
return {"fields": builder.get_target_fields()}
@router.post("/generate", response_model=GenerateResponse)
async def generate_profile(payload: GenerateRequest):
if not builder.deepseek_enabled():
raise HTTPException(
status_code=503,
detail="DEEPSEEK_API_KEY is not set on ca-api. Add it to .env for parser generation.",
)
try:
profile = await builder.generate_profile(payload.sample_post)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
except Exception as exc:
raise HTTPException(
status_code=502,
detail=f"DeepSeek generate failed: {exc}",
) from exc
return GenerateResponse(profile=profile.model_dump())
@router.post("/preview", response_model=PreviewResponse)
def preview_profile(payload: PreviewRequest):
try:
HeuristicProfile.model_validate(payload.profile)
fields = builder.preview_with_profile(payload.sample_post, payload.profile)
except Exception as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
return PreviewResponse(fields=fields)
@router.post("/jobs", response_model=ParseJobRead, status_code=201)
def create_profile_job(payload: CreateProfileJobRequest, db: Session = Depends(get_db)):
try:
cfg = TelegramSourceConfig(
channel=payload.channel.strip(),
limit=payload.limit,
extract_mode="profile",
heuristic_profile=payload.profile,
sample_post=payload.sample_post,
)
except Exception as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
job = ParseJob(
source_type="telegram",
source_config=cfg.model_dump(),
interval_seconds=payload.interval_seconds,
is_active=payload.is_active,
status="queued",
)
db.add(job)
db.commit()
db.refresh(job)
enqueue_job(job.id, job.source_type, job.source_config)
return job
@@ -0,0 +1,164 @@
"""Admin CRUD for reusable parser profiles + generate/preview."""
from __future__ import annotations
from typing import Any
from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel, Field
from sqlalchemy.orm import Session
from contracts.heuristic_profile import HeuristicProfile
from ..database import get_db
from ..models import ParseJob, ParserProfile
from ..schemas import ParserProfileCreate, ParserProfileRead, ParserProfileUpdate
from ..services import parser_builder as builder
router = APIRouter(prefix="/admin/parser-profiles", tags=["parser-profiles"])
class GenerateRequest(BaseModel):
sample_post: str = Field(min_length=1)
class GenerateResponse(BaseModel):
profile: dict[str, Any]
class PreviewRequest(BaseModel):
sample_post: str = Field(min_length=1)
heuristic_profile: dict[str, Any] | None = None
# Accept alias used by older builder UI
profile: dict[str, Any] | None = None
class PreviewResponse(BaseModel):
fields: dict[str, str]
def _validate_heuristic_profile(raw: dict[str, Any] | None) -> dict[str, Any] | None:
if raw is None:
return None
try:
return HeuristicProfile.model_validate(raw).model_dump()
except Exception as exc:
raise HTTPException(status_code=400, detail=f"Invalid heuristic_profile: {exc}") from exc
def _profile_status(heuristic_profile: dict | None, explicit: str | None = None) -> str:
if explicit:
return explicit
return "ready" if heuristic_profile else "draft"
@router.get("/target-fields")
def target_fields():
return {"fields": builder.get_target_fields()}
@router.post("/generate", response_model=GenerateResponse)
async def generate_profile(payload: GenerateRequest):
if not builder.deepseek_enabled():
raise HTTPException(
status_code=503,
detail="DEEPSEEK_API_KEY is not set on ca-api. Add it to .env for parser generation.",
)
try:
profile = await builder.generate_profile(payload.sample_post)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
except Exception as exc:
raise HTTPException(
status_code=502,
detail=f"DeepSeek generate failed: {exc}",
) from exc
return GenerateResponse(profile=profile.model_dump())
@router.post("/preview", response_model=PreviewResponse)
def preview_profile(payload: PreviewRequest):
raw = payload.heuristic_profile if payload.heuristic_profile is not None else payload.profile
if raw is None:
raise HTTPException(status_code=400, detail="heuristic_profile is required")
try:
HeuristicProfile.model_validate(raw)
fields = builder.preview_with_profile(payload.sample_post, raw)
except Exception as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
return PreviewResponse(fields=fields)
@router.get("", response_model=list[ParserProfileRead])
def list_profiles(db: Session = Depends(get_db)):
return db.query(ParserProfile).order_by(ParserProfile.id.desc()).all()
@router.post("", response_model=ParserProfileRead, status_code=201)
def create_profile(payload: ParserProfileCreate, db: Session = Depends(get_db)):
heuristic = _validate_heuristic_profile(payload.heuristic_profile)
profile = ParserProfile(
name=payload.name.strip(),
sample_post=payload.sample_post or "",
heuristic_profile=heuristic,
status=_profile_status(heuristic, payload.status),
)
db.add(profile)
db.commit()
db.refresh(profile)
return profile
@router.get("/{profile_id}", response_model=ParserProfileRead)
def get_profile(profile_id: int, db: Session = Depends(get_db)):
profile = db.query(ParserProfile).filter(ParserProfile.id == profile_id).first()
if not profile:
raise HTTPException(status_code=404, detail="Profile not found")
return profile
@router.patch("/{profile_id}", response_model=ParserProfileRead)
def update_profile(
profile_id: int,
payload: ParserProfileUpdate,
db: Session = Depends(get_db),
):
profile = db.query(ParserProfile).filter(ParserProfile.id == profile_id).first()
if not profile:
raise HTTPException(status_code=404, detail="Profile not found")
if payload.name is not None:
profile.name = payload.name.strip()
if payload.sample_post is not None:
profile.sample_post = payload.sample_post
if "heuristic_profile" in payload.model_fields_set:
profile.heuristic_profile = _validate_heuristic_profile(payload.heuristic_profile)
if payload.status is not None:
profile.status = payload.status
elif "heuristic_profile" in payload.model_fields_set:
profile.status = _profile_status(profile.heuristic_profile)
db.commit()
db.refresh(profile)
return profile
@router.delete("/{profile_id}", status_code=204)
def delete_profile(profile_id: int, db: Session = Depends(get_db)):
profile = db.query(ParserProfile).filter(ParserProfile.id == profile_id).first()
if not profile:
raise HTTPException(status_code=404, detail="Profile not found")
linked = (
db.query(ParseJob)
.filter(ParseJob.profile_id == profile_id, ParseJob.status.in_(("queued", "running")))
.count()
)
if linked:
raise HTTPException(
status_code=409,
detail="Cannot delete profile used by queued/running jobs",
)
db.delete(profile)
db.commit()
+93
View File
@@ -92,6 +92,38 @@ class EventRead(BaseModel):
metadata: dict[str, Any] | None = Field(validation_alias="metadata_")
class EventCreate(BaseModel):
source_type: str = Field(default="manual", min_length=1, max_length=50)
source_url: str = Field(min_length=1, max_length=512)
raw_text: str = ""
title: str = ""
description: str = ""
locality: str = ""
latitude: float | None = Field(default=None, ge=-90, le=90)
longitude: float | None = Field(default=None, ge=-180, le=180)
event_date: datetime | None = None
region: str | None = None
topic: str | None = None
tags: list[str] | None = None
metadata: dict[str, Any] | None = None
class EventUpdate(BaseModel):
source_type: str | None = Field(default=None, min_length=1, max_length=50)
source_url: str | None = Field(default=None, min_length=1, max_length=512)
raw_text: str | None = None
title: str | None = None
description: str | None = None
locality: str | None = None
latitude: float | None = Field(default=None, ge=-90, le=90)
longitude: float | None = Field(default=None, ge=-180, le=180)
event_date: datetime | None = None
region: str | None = None
topic: str | None = None
tags: list[str] | None = None
metadata: dict[str, Any] | None = None
class IngestEventItem(BaseModel):
source_type: str = "telegram"
source_url: str
@@ -127,9 +159,62 @@ class IngestResponse(BaseModel):
map_objects_synced: int
class ParserProfileCreate(BaseModel):
name: str = Field(min_length=1, max_length=255)
sample_post: str = ""
heuristic_profile: dict[str, Any] | None = None
status: str | None = None
class ParserProfileUpdate(BaseModel):
name: str | None = Field(default=None, min_length=1, max_length=255)
sample_post: str | None = None
heuristic_profile: dict[str, Any] | None = None
status: str | None = None
class ParserProfileRead(BaseModel):
model_config = ConfigDict(from_attributes=True)
id: int
name: str
sample_post: str
heuristic_profile: dict[str, Any] | None
status: str
created_at: datetime
class ParseChannelCreate(BaseModel):
name: str = Field(min_length=1, max_length=255)
source_type: str = "telegram"
channel: str = Field(min_length=1, max_length=255)
is_active: bool = True
class ParseChannelUpdate(BaseModel):
name: str | None = Field(default=None, min_length=1, max_length=255)
source_type: str | None = None
channel: str | None = Field(default=None, min_length=1, max_length=255)
is_active: bool | None = None
class ParseChannelRead(BaseModel):
model_config = ConfigDict(from_attributes=True)
id: int
name: str
source_type: str
channel: str
is_active: bool
created_at: datetime
class ParseJobCreate(BaseModel):
source_type: str = "telegram"
source_config: dict[str, Any] = Field(default_factory=dict)
profile_id: int | None = None
channel_id: int | None = None
limit: int = Field(default=100, ge=1, le=1000)
schedule: str | None = None
interval_seconds: int = Field(default=3600, ge=60, le=604800)
is_active: bool = True
@@ -137,6 +222,9 @@ class ParseJobCreate(BaseModel):
class ParseJobUpdate(BaseModel):
source_config: dict[str, Any] | None = None
profile_id: int | None = None
channel_id: int | None = None
limit: int | None = Field(default=None, ge=1, le=1000)
interval_seconds: int | None = Field(default=None, ge=60, le=604800)
is_active: bool | None = None
@@ -147,6 +235,8 @@ class ParseJobRead(BaseModel):
id: int
source_type: str
source_config: dict[str, Any]
profile_id: int | None = None
channel_id: int | None = None
schedule: str | None
interval_seconds: int
is_active: bool
@@ -154,6 +244,9 @@ class ParseJobRead(BaseModel):
last_run_at: datetime | None
last_error: str | None
created_at: datetime
profile_name: str | None = None
channel_name: str | None = None
channel_handle: str | None = None
class ConsumerCreate(BaseModel):
+52 -1
View File
@@ -4,6 +4,7 @@ import sys
from pathlib import Path
import redis
from sqlalchemy.orm import Session
# Allow importing shared contracts from monorepo root in local runs
_HERE = Path(__file__).resolve()
@@ -14,7 +15,9 @@ for _candidate in (_HERE.parent, *_HERE.parents):
break
from contracts.queues import queue_key_for_source # noqa: E402
from contracts.sources import parse_source_config # noqa: E402
from contracts.sources import TelegramSourceConfig, parse_source_config # noqa: E402
from ..models import ParseChannel, ParseJob, ParserProfile # noqa: E402
REDIS_URL = os.getenv("REDIS_URL", "redis://redis:6379/0")
@@ -23,6 +26,47 @@ def get_redis() -> redis.Redis:
return redis.from_url(REDIS_URL, decode_responses=True)
def flatten_pair_config(
channel: ParseChannel,
profile: ParserProfile,
*,
limit: int = 100,
) -> dict:
"""Expand Profile + Channel into Redis/CP source_config."""
if not profile.heuristic_profile:
raise ValueError("ParserProfile.heuristic_profile is empty")
cfg = TelegramSourceConfig(
channel=channel.channel.strip(),
limit=limit,
extract_mode="profile",
heuristic_profile=profile.heuristic_profile,
sample_post=profile.sample_post or None,
)
return cfg.model_dump()
def resolve_job_source_config(db: Session, job: ParseJob) -> dict:
"""
For pair jobs, rebuild flat source_config from current Profile/Channel.
Legacy jobs keep stored source_config.
"""
if job.profile_id and job.channel_id:
profile = db.query(ParserProfile).filter(ParserProfile.id == job.profile_id).first()
channel = db.query(ParseChannel).filter(ParseChannel.id == job.channel_id).first()
if not profile:
raise ValueError(f"ParserProfile {job.profile_id} not found")
if not channel:
raise ValueError(f"ParseChannel {job.channel_id} not found")
existing = job.source_config or {}
limit = existing.get("limit", 100)
if not isinstance(limit, int) or limit < 1:
limit = 100
flattened = flatten_pair_config(channel, profile, limit=limit)
job.source_config = flattened
return flattened
return dict(job.source_config or {})
def enqueue_job(job_id: int, source_type: str, source_config: dict) -> None:
# Validate known configs early; unknown types still raise from queue_key_for_source
try:
@@ -38,3 +82,10 @@ def enqueue_job(job_id: int, source_type: str, source_config: dict) -> None:
}
key = queue_key_for_source(source_type)
get_redis().rpush(key, json.dumps(payload))
def enqueue_parse_job(db: Session, job: ParseJob) -> None:
"""Resolve (flatten if pair) then push to Redis."""
source_config = resolve_job_source_config(db, job)
db.commit()
enqueue_job(job.id, job.source_type, source_config)
@@ -5,20 +5,29 @@ from sqlalchemy.engine import Engine
def migrate_schema(engine: Engine) -> None:
"""Apply lightweight schema updates for existing deployments."""
inspector = inspect(engine)
if "parse_jobs" not in inspector.get_table_names():
return
columns = {col["name"] for col in inspector.get_columns("parse_jobs")}
tables = set(inspector.get_table_names())
statements: list[str] = []
if "interval_seconds" not in columns:
statements.append(
"ALTER TABLE parse_jobs ADD COLUMN interval_seconds INTEGER NOT NULL DEFAULT 3600"
)
if "is_active" not in columns:
statements.append(
"ALTER TABLE parse_jobs ADD COLUMN is_active BOOLEAN NOT NULL DEFAULT TRUE"
)
if "parse_jobs" in tables:
columns = {col["name"] for col in inspector.get_columns("parse_jobs")}
if "interval_seconds" not in columns:
statements.append(
"ALTER TABLE parse_jobs ADD COLUMN interval_seconds INTEGER NOT NULL DEFAULT 3600"
)
if "is_active" not in columns:
statements.append(
"ALTER TABLE parse_jobs ADD COLUMN is_active BOOLEAN NOT NULL DEFAULT TRUE"
)
if "profile_id" not in columns:
statements.append("ALTER TABLE parse_jobs ADD COLUMN profile_id INTEGER")
if "channel_id" not in columns:
statements.append("ALTER TABLE parse_jobs ADD COLUMN channel_id INTEGER")
# create_all handles new tables; FKs on existing DBs may need indexes
if "parse_jobs" in tables:
# Re-inspect after potential adds is not needed for FK constraints here —
# create_all + nullable FKs are enough for MVP; optional constraints below.
pass
if not statements:
return
@@ -5,7 +5,7 @@ from datetime import datetime, timezone
from ..database import SessionLocal
from ..models import ParseJob
from .jobs import enqueue_job
from .jobs import enqueue_parse_job
logger = logging.getLogger(__name__)
@@ -36,7 +36,14 @@ def run_scheduler_tick() -> None:
job.status = "queued"
job.last_error = None
db.commit()
enqueue_job(job.id, job.source_type, job.source_config)
try:
enqueue_parse_job(db, job)
except ValueError as exc:
job.status = "failed"
job.last_error = str(exc)
db.commit()
logger.warning("Skip re-queue job %s: %s", job.id, exc)
continue
logger.info("Re-queued recurring job %s (interval %ss)", job.id, job.interval_seconds)
except Exception:
logger.exception("Scheduler tick failed")