mirror of
https://github.com/hansjone/netx.git
synced 2026-10-09 03:10:46 +08:00
feat(collect): draft jobs, merged create UI, and unified start
Create collection jobs without auto-running; add /start for first run and re-run. Merge create form with filterable NE picker and selected list. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
fed7f56ad0
commit
f1ed4fe753
9 changed files with 343 additions and 147 deletions
|
|
@ -67,8 +67,8 @@ def reconcile_stale_collection_job(db: Session, job_id: str) -> bool:
|
|||
if st in _TERMINAL:
|
||||
continue
|
||||
if st == "pending":
|
||||
# Pending while job is running means queued in the worker pool, not stuck.
|
||||
if str(job.status or "") == "running":
|
||||
# Pending while job is running/queued, or draft job not started yet.
|
||||
if str(job.status or "") in ("running", "pending"):
|
||||
continue
|
||||
anchor = job.started_at or job.created_at
|
||||
limit = pending_stale_sec
|
||||
|
|
|
|||
|
|
@ -6,7 +6,7 @@ from sqlalchemy.orm import Session
|
|||
|
||||
from .collection_service import (
|
||||
build_collection_job_zip,
|
||||
create_and_start_collection,
|
||||
create_collection,
|
||||
delete_collection_job,
|
||||
get_collection_job,
|
||||
list_collection_jobs,
|
||||
|
|
@ -15,6 +15,7 @@ from .collection_service import (
|
|||
pause_collection_job,
|
||||
resolve_run_output_path,
|
||||
restart_collection_job,
|
||||
start_collection_job,
|
||||
retry_failed_collection_job,
|
||||
)
|
||||
from .collection_schemas import CollectionJobCreate
|
||||
|
|
@ -29,25 +30,15 @@ router = APIRouter(prefix="/v1/ne-collections", tags=["ne-collections"])
|
|||
def api_eligible_ne(
|
||||
page: int = Query(default=1, ge=1),
|
||||
page_size: int = Query(default=200, ge=1, le=500),
|
||||
keyword: str = Query(default=""),
|
||||
db: Session = Depends(get_db),
|
||||
):
|
||||
return list_eligible_ne(db, page=page, page_size=page_size)
|
||||
return list_eligible_ne(db, page=page, page_size=page_size, keyword=keyword)
|
||||
|
||||
|
||||
@router.post("")
|
||||
def api_create_collection(
|
||||
body: CollectionJobCreate,
|
||||
background_tasks: BackgroundTasks,
|
||||
db: Session = Depends(get_db),
|
||||
):
|
||||
out, payload = create_and_start_collection(db, body)
|
||||
background_tasks.add_task(
|
||||
dispatch_collection_runs,
|
||||
payload["job_id"],
|
||||
payload["run_ids"],
|
||||
payload["commands"],
|
||||
)
|
||||
return out.model_dump()
|
||||
def api_create_collection(body: CollectionJobCreate, db: Session = Depends(get_db)):
|
||||
return create_collection(db, body).model_dump()
|
||||
|
||||
|
||||
@router.get("")
|
||||
|
|
@ -103,12 +94,29 @@ def api_pause_collection(job_id: str, db: Session = Depends(get_db)):
|
|||
return pause_collection_job(db, job_id).model_dump()
|
||||
|
||||
|
||||
@router.post("/{job_id}/start")
|
||||
def api_start_collection(
|
||||
job_id: str,
|
||||
background_tasks: BackgroundTasks,
|
||||
db: Session = Depends(get_db),
|
||||
):
|
||||
out, payload = start_collection_job(db, job_id)
|
||||
background_tasks.add_task(
|
||||
dispatch_collection_runs,
|
||||
payload["job_id"],
|
||||
payload["run_ids"],
|
||||
payload["commands"],
|
||||
)
|
||||
return out.model_dump()
|
||||
|
||||
|
||||
@router.post("/{job_id}/restart")
|
||||
def api_restart_collection(
|
||||
job_id: str,
|
||||
background_tasks: BackgroundTasks,
|
||||
db: Session = Depends(get_db),
|
||||
):
|
||||
"""Alias of /start for backward compatibility."""
|
||||
out, payload = restart_collection_job(db, job_id)
|
||||
background_tasks.add_task(
|
||||
dispatch_collection_runs,
|
||||
|
|
|
|||
|
|
@ -111,8 +111,25 @@ def run_to_out(row: NeCollectionRun) -> CollectionRunOut:
|
|||
)
|
||||
|
||||
|
||||
def list_eligible_ne(db: Session, *, page: int = 1, page_size: int = 200) -> dict[str, Any]:
|
||||
def list_eligible_ne(
|
||||
db: Session,
|
||||
*,
|
||||
page: int = 1,
|
||||
page_size: int = 200,
|
||||
keyword: str = "",
|
||||
) -> dict[str, Any]:
|
||||
stmt = db.query(ManagedNE).filter(ManagedNE.connect_status == "pass")
|
||||
kw = str(keyword or "").strip()
|
||||
if kw:
|
||||
like = f"%{kw}%"
|
||||
stmt = stmt.filter(
|
||||
or_(
|
||||
ManagedNE.name.ilike(like),
|
||||
ManagedNE.ip_address.ilike(like),
|
||||
ManagedNE.vendor.ilike(like),
|
||||
ManagedNE.device_type.ilike(like),
|
||||
)
|
||||
)
|
||||
total = int(stmt.count())
|
||||
rows = (
|
||||
stmt.order_by(ManagedNE.name.asc())
|
||||
|
|
@ -135,9 +152,7 @@ def list_eligible_ne(db: Session, *, page: int = 1, page_size: int = 200) -> dic
|
|||
return {"total": total, "page": page, "page_size": page_size, "items": items}
|
||||
|
||||
|
||||
def create_and_start_collection(
|
||||
db: Session, body: CollectionJobCreate
|
||||
) -> tuple[CollectionJobOut, CollectionSchedulePayload]:
|
||||
def create_collection(db: Session, body: CollectionJobCreate) -> CollectionJobOut:
|
||||
commands = _parse_commands(body.commands)
|
||||
if not commands:
|
||||
raise HTTPException(status_code=400, detail="commands_empty")
|
||||
|
|
@ -168,16 +183,15 @@ def create_and_start_collection(
|
|||
job = NeCollectionJob(
|
||||
title=str(body.title or "").strip() or f"collect-{now.strftime('%Y%m%d-%H%M%S')}",
|
||||
commands="\n".join(commands),
|
||||
status="running",
|
||||
status="pending",
|
||||
ne_count=len(ne_rows),
|
||||
created_at=now,
|
||||
started_at=now,
|
||||
last_run_at=now,
|
||||
started_at=None,
|
||||
last_run_at=None,
|
||||
)
|
||||
db.add(job)
|
||||
db.flush()
|
||||
|
||||
run_ids: list[str] = []
|
||||
for ne in ne_rows:
|
||||
run = NeCollectionRun(
|
||||
job_id=str(job.id),
|
||||
|
|
@ -187,16 +201,9 @@ def create_and_start_collection(
|
|||
status="pending",
|
||||
)
|
||||
db.add(run)
|
||||
run_ids.append(str(run.id))
|
||||
db.commit()
|
||||
db.refresh(job)
|
||||
|
||||
payload: CollectionSchedulePayload = {
|
||||
"job_id": str(job.id),
|
||||
"run_ids": run_ids,
|
||||
"commands": commands,
|
||||
}
|
||||
return job_to_out(job, output_count=0), payload
|
||||
return job_to_out(job, output_count=0)
|
||||
|
||||
|
||||
def list_collection_jobs(db: Session, *, page: int = 1, page_size: int = 20) -> dict[str, Any]:
|
||||
|
|
@ -355,7 +362,8 @@ def _start_job_retry(
|
|||
return job_to_out(job, output_count=_output_count_for_job(db, job_id)), payload
|
||||
|
||||
|
||||
def restart_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, CollectionSchedulePayload]:
|
||||
def start_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, CollectionSchedulePayload]:
|
||||
"""Start a draft job or re-run all NEs after pause/completion."""
|
||||
job = db.get(NeCollectionJob, job_id)
|
||||
if not job:
|
||||
raise HTTPException(status_code=404, detail="collection_job_not_found")
|
||||
|
|
@ -367,6 +375,13 @@ def restart_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut,
|
|||
commands = _parse_commands(str(job.commands or ""))
|
||||
if not commands:
|
||||
raise HTTPException(status_code=400, detail="commands_empty")
|
||||
|
||||
if str(job.status or "") == "pending":
|
||||
run_ids = [str(r.id) for r in runs if str(r.status or "") == "pending"]
|
||||
if not run_ids:
|
||||
raise HTTPException(status_code=400, detail="collection_nothing_to_start")
|
||||
return _start_job_retry(db, job, job_id, run_ids, commands, reset_all_counts=False)
|
||||
|
||||
retry_ids = _reset_runs_for_retry(
|
||||
db,
|
||||
job_id,
|
||||
|
|
@ -376,6 +391,10 @@ def restart_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut,
|
|||
return _start_job_retry(db, job, job_id, retry_ids, commands, reset_all_counts=True)
|
||||
|
||||
|
||||
def restart_collection_job(db: Session, job_id: str) -> tuple[CollectionJobOut, CollectionSchedulePayload]:
|
||||
return start_collection_job(db, job_id)
|
||||
|
||||
|
||||
def retry_failed_collection_job(
|
||||
db: Session, job_id: str
|
||||
) -> tuple[CollectionJobOut, CollectionSchedulePayload]:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue