Next alpha stage commit
This commit is contained in:
131
app/jobs.py
131
app/jobs.py
@@ -23,6 +23,7 @@ from app.data_management import (
|
||||
from app.db import SessionLocal, engine, init_db
|
||||
from app.db_lock import DatabaseWriteBusy, database_write_lock
|
||||
from app.gtfs_storage import missing_sidecar_paths as gtfs_missing_sidecar_paths
|
||||
from app.harmonization import create_gtfs_harmonized_snapshot
|
||||
from app.models import Dataset, Job, JobEvent, Source, SourceUpdateCheck
|
||||
from app.osm_storage import missing_sidecar_paths as osm_missing_sidecar_paths
|
||||
from app.pipeline.gtfs import backfill_gtfs_shapes
|
||||
@@ -34,6 +35,7 @@ from app.pipeline.route_layer import rebuild_route_layer
|
||||
from app.pipeline.run import run_source
|
||||
from app.pipeline.sample_data import clear_project_data, load_sample_project
|
||||
from app.source_catalog import import_ingestable_sources, import_source_catalog, source_catalog_summary
|
||||
from app.workbench import generate_map_gtfs_review_items
|
||||
|
||||
|
||||
ROUTE_MATCHING_JOB_KIND = "route_matching"
|
||||
@@ -44,6 +46,8 @@ SOURCE_IMPORT_JOB_KIND = "source_import"
|
||||
SOURCE_DELETE_JOB_KIND = "source_delete"
|
||||
DATASET_DELETE_JOB_KIND = "dataset_delete"
|
||||
MAINTENANCE_JOB_KIND = "maintenance"
|
||||
MAP_GTFS_REVIEW_JOB_KIND = "map_gtfs_review"
|
||||
GTFS_HARMONIZED_SNAPSHOT_JOB_KIND = "gtfs_harmonized_snapshot"
|
||||
TERMINAL_JOB_STATUSES = {"completed", "failed", "cancelled"}
|
||||
ACTIVE_JOB_STATUSES = {"queued", "running", "paused"}
|
||||
LEASE_SECONDS = max(300, int(settings.queue_job_lease_seconds))
|
||||
@@ -124,6 +128,70 @@ def create_route_matching_job(session: Session, *, priority: int = 0) -> Job:
|
||||
return job
|
||||
|
||||
|
||||
def create_map_gtfs_review_job(
|
||||
session: Session,
|
||||
*,
|
||||
source_id: int | None = None,
|
||||
dataset_id: int | None = None,
|
||||
limit: int = 1000,
|
||||
priority: int = 0,
|
||||
) -> Job:
|
||||
result = {
|
||||
"source_id": source_id,
|
||||
"dataset_id": dataset_id,
|
||||
"limit": max(1, min(int(limit), 50_000)),
|
||||
}
|
||||
job = Job(
|
||||
kind=MAP_GTFS_REVIEW_JOB_KIND,
|
||||
status="queued",
|
||||
description="Generate Map/GTFS workbench review queue",
|
||||
progress_current=0,
|
||||
progress_total=2,
|
||||
priority=int(priority),
|
||||
result_json=json.dumps(result, separators=(",", ":")),
|
||||
)
|
||||
session.add(job)
|
||||
session.flush()
|
||||
add_job_event(
|
||||
session,
|
||||
job,
|
||||
event_type="queued",
|
||||
message="Map/GTFS review generation queued.",
|
||||
progress_current=0,
|
||||
progress_total=job.progress_total,
|
||||
metadata=result,
|
||||
)
|
||||
return job
|
||||
|
||||
|
||||
def create_gtfs_harmonized_snapshot_job(session: Session, *, activate: bool = True, priority: int = 0) -> Job:
|
||||
result = {
|
||||
"activate": bool(activate),
|
||||
"queued_pid": os.getpid(),
|
||||
}
|
||||
job = Job(
|
||||
kind=GTFS_HARMONIZED_SNAPSHOT_JOB_KIND,
|
||||
status="queued",
|
||||
description="Build harmonized GTFS snapshot",
|
||||
progress_current=0,
|
||||
progress_total=2,
|
||||
priority=int(priority),
|
||||
result_json=json.dumps(result, separators=(",", ":")),
|
||||
)
|
||||
session.add(job)
|
||||
session.flush()
|
||||
add_job_event(
|
||||
session,
|
||||
job,
|
||||
event_type="queued",
|
||||
message="Harmonized GTFS snapshot build queued.",
|
||||
progress_current=0,
|
||||
progress_total=job.progress_total,
|
||||
metadata=result,
|
||||
)
|
||||
return job
|
||||
|
||||
|
||||
def create_osm_relabel_job(
|
||||
session: Session,
|
||||
*,
|
||||
@@ -855,12 +923,14 @@ def job_queue_revision(session: Session, *, include_dismissed: bool = False) ->
|
||||
|
||||
|
||||
def job_events(session: Session, job_id: int, *, limit: int = 100) -> list[JobEvent]:
|
||||
return session.scalars(
|
||||
event_limit = max(1, min(limit, 500))
|
||||
events = session.scalars(
|
||||
select(JobEvent)
|
||||
.where(JobEvent.job_id == job_id)
|
||||
.order_by(JobEvent.created_at, JobEvent.id)
|
||||
.limit(max(1, min(limit, 500)))
|
||||
.order_by(JobEvent.created_at.desc(), JobEvent.id.desc())
|
||||
.limit(event_limit)
|
||||
).all()
|
||||
return list(reversed(events))
|
||||
|
||||
|
||||
def request_job_control(job_id: int, action: str) -> dict[str, Any]:
|
||||
@@ -990,6 +1060,10 @@ def run_claimed_job(job_id: int, worker_id: str) -> None:
|
||||
_run_dataset_delete_job(job_id, worker_id)
|
||||
elif job.kind == MAINTENANCE_JOB_KIND:
|
||||
_run_maintenance_job(job_id, worker_id)
|
||||
elif job.kind == MAP_GTFS_REVIEW_JOB_KIND:
|
||||
_run_map_gtfs_review_job(job_id, worker_id)
|
||||
elif job.kind == GTFS_HARMONIZED_SNAPSHOT_JOB_KIND:
|
||||
_run_gtfs_harmonized_snapshot_job(job_id, worker_id)
|
||||
else:
|
||||
raise ValueError(f"unsupported job kind: {job.kind}")
|
||||
except JobPaused:
|
||||
@@ -1292,6 +1366,46 @@ def _run_osm_relabel_job(job_id: int, worker_id: str) -> None:
|
||||
session.commit()
|
||||
|
||||
|
||||
def _run_map_gtfs_review_job(job_id: int, worker_id: str) -> None:
|
||||
init_db()
|
||||
with SessionLocal() as session:
|
||||
job = _job_for_worker(session, job_id, worker_id)
|
||||
options = _json_object(job.result_json)
|
||||
_job_running(session, job, worker_id, "started", "Generating Map/GTFS workbench review queue.", 1)
|
||||
_check_job_control(session, job)
|
||||
progress_callback = _job_progress_callback(session, job, worker_id, update_job_progress=False)
|
||||
result = generate_map_gtfs_review_items(
|
||||
session,
|
||||
source_id=_optional_int(options.get("source_id")),
|
||||
dataset_id=_optional_int(options.get("dataset_id")),
|
||||
limit=int(options.get("limit") or 1000),
|
||||
progress_callback=progress_callback,
|
||||
)
|
||||
job = _job_for_worker(session, job_id, worker_id)
|
||||
_complete_job(session, job, "Map/GTFS workbench review queue generated.", {**options, "review_result": result})
|
||||
session.commit()
|
||||
|
||||
|
||||
def _run_gtfs_harmonized_snapshot_job(job_id: int, worker_id: str) -> None:
|
||||
init_db()
|
||||
with SessionLocal() as session:
|
||||
job = _job_for_worker(session, job_id, worker_id)
|
||||
options = _json_object(job.result_json)
|
||||
options["worker_pid"] = os.getpid()
|
||||
options["worker_id"] = worker_id
|
||||
job.result_json = json.dumps(options, separators=(",", ":"))
|
||||
_job_running(session, job, worker_id, "started", "Building harmonized GTFS snapshot.", 1)
|
||||
_check_job_control(session, job)
|
||||
snapshot = create_gtfs_harmonized_snapshot(
|
||||
session,
|
||||
activate=bool(options.get("activate", True)),
|
||||
note=f"job #{job.id}",
|
||||
)
|
||||
job = _job_for_worker(session, job_id, worker_id)
|
||||
_complete_job(session, job, "Harmonized GTFS snapshot built.", {**options, "snapshot": snapshot})
|
||||
session.commit()
|
||||
|
||||
|
||||
def _run_source_import_job(job_id: int, worker_id: str) -> None:
|
||||
init_db()
|
||||
with SessionLocal() as session:
|
||||
@@ -1335,7 +1449,16 @@ def _run_source_import_job(job_id: int, worker_id: str) -> None:
|
||||
if options.get("run_match"):
|
||||
_job_running(session, job, worker_id, "matching", "Running route matcher after import.", 3)
|
||||
progress_callback = _job_progress_callback(session, job, worker_id)
|
||||
match_result = run_route_matching(session, progress_callback=progress_callback)
|
||||
imported_gtfs_dataset_id = (
|
||||
int(options["dataset_id"])
|
||||
if options.get("dataset_kind") == "gtfs" and options.get("dataset_id")
|
||||
else None
|
||||
)
|
||||
match_result = run_route_matching(
|
||||
session,
|
||||
progress_callback=progress_callback,
|
||||
gtfs_dataset_ids=[imported_gtfs_dataset_id] if imported_gtfs_dataset_id is not None else None,
|
||||
)
|
||||
job = _job_for_worker(session, job_id, worker_id)
|
||||
options = _json_object(job.result_json)
|
||||
options["match_result"] = match_result
|
||||
|
||||
Reference in New Issue
Block a user