Initialer Import des AssetManagers
This commit is contained in:
@@ -0,0 +1,187 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from .models import AssetJobState, SoftwareJob
|
||||
|
||||
|
||||
TERMINAL_JOB_STATES = {"success", "failed", "partial", "timeout", "cancelled"}
|
||||
|
||||
|
||||
def _int_or_none(value: Any) -> int | None:
|
||||
try:
|
||||
parsed = int(value)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
return parsed if parsed > 0 else None
|
||||
|
||||
|
||||
def job_definition_id(job: SoftwareJob) -> int | None:
|
||||
if str(job.job_type or "") != "job_definition":
|
||||
return None
|
||||
parameters = job.parameters if isinstance(job.parameters, dict) else {}
|
||||
return _int_or_none(parameters.get("definition_id"))
|
||||
|
||||
|
||||
def job_definition_revision(job: SoftwareJob) -> int | None:
|
||||
if str(job.job_type or "") != "job_definition":
|
||||
return None
|
||||
parameters = job.parameters if isinstance(job.parameters, dict) else {}
|
||||
return _int_or_none(parameters.get("definition_revision"))
|
||||
|
||||
|
||||
def job_state_key(job: SoftwareJob) -> str:
|
||||
job_type = str(job.job_type or "software_inventory").strip() or "software_inventory"
|
||||
if job_type == "job_definition":
|
||||
definition_id = job_definition_id(job)
|
||||
return f"job_definition:{definition_id}" if definition_id else f"job_definition:job:{job.id}"
|
||||
return job_type[:120]
|
||||
|
||||
|
||||
def filter_state_key(job_type: str, definition_id: int | None = None) -> str:
|
||||
normalized = str(job_type or "").strip()
|
||||
if normalized == "job_definition" and definition_id:
|
||||
return f"job_definition:{int(definition_id)}"
|
||||
return normalized[:120]
|
||||
|
||||
|
||||
def sync_asset_job_state(
|
||||
db: Session,
|
||||
job: SoftwareJob,
|
||||
*,
|
||||
increment_execution: bool = False,
|
||||
status_at: datetime | None = None,
|
||||
) -> AssetJobState:
|
||||
"""Update the compact per-asset state for one job.
|
||||
|
||||
This table is an intentionally small projection of ``software_jobs``. The
|
||||
full job history remains authoritative; this projection exists only for
|
||||
fast filters, dashboards and the later software-package desired-state
|
||||
module.
|
||||
"""
|
||||
state_key = job_state_key(job)
|
||||
state = (
|
||||
db.query(AssetJobState)
|
||||
.filter(AssetJobState.asset_id == job.asset_id, AssetJobState.state_key == state_key)
|
||||
.with_for_update()
|
||||
.first()
|
||||
)
|
||||
now = status_at or datetime.utcnow()
|
||||
if state is None:
|
||||
state = AssetJobState(
|
||||
asset_id=job.asset_id,
|
||||
state_key=state_key,
|
||||
job_type=str(job.job_type or "software_inventory")[:50],
|
||||
job_definition_id=job_definition_id(job),
|
||||
execution_count=0,
|
||||
)
|
||||
db.add(state)
|
||||
db.flush()
|
||||
|
||||
revision = job_definition_revision(job)
|
||||
state.job_type = str(job.job_type or "software_inventory")[:50]
|
||||
state.job_definition_id = job_definition_id(job)
|
||||
state.definition_revision = revision
|
||||
state.last_job_id = job.id
|
||||
state.last_status = str(job.status or "created")[:30]
|
||||
state.last_status_at = now
|
||||
state.last_created_at = job.created_at
|
||||
state.last_sent_at = job.sent_at
|
||||
state.last_started_at = job.started_at
|
||||
state.last_completed_at = job.finished_at or job.callback_completed_at
|
||||
state.updated_at = now
|
||||
if increment_execution:
|
||||
state.execution_count = int(state.execution_count or 0) + 1
|
||||
elif int(state.execution_count or 0) <= 0:
|
||||
state.execution_count = max(1, int(job.attempt_count or 0))
|
||||
|
||||
if state.last_status == "success":
|
||||
state.last_success_at = job.finished_at or job.callback_completed_at or now
|
||||
state.last_success_revision = revision
|
||||
return state
|
||||
|
||||
|
||||
def backfill_asset_job_states(db: Session) -> int:
|
||||
"""Create compact states from historic jobs once after the table is added."""
|
||||
if db.query(AssetJobState.id).limit(1).first() is not None:
|
||||
return 0
|
||||
|
||||
rows = (
|
||||
db.query(
|
||||
SoftwareJob.id,
|
||||
SoftwareJob.asset_id,
|
||||
SoftwareJob.job_type,
|
||||
SoftwareJob.parameters,
|
||||
SoftwareJob.status,
|
||||
SoftwareJob.attempt_count,
|
||||
SoftwareJob.created_at,
|
||||
SoftwareJob.sent_at,
|
||||
SoftwareJob.started_at,
|
||||
SoftwareJob.finished_at,
|
||||
SoftwareJob.callback_completed_at,
|
||||
)
|
||||
.order_by(SoftwareJob.created_at.asc(), SoftwareJob.id.asc())
|
||||
.all()
|
||||
)
|
||||
if not rows:
|
||||
return 0
|
||||
|
||||
states: dict[tuple[int, str], dict[str, Any]] = {}
|
||||
for row in rows:
|
||||
parameters = row.parameters if isinstance(row.parameters, dict) else {}
|
||||
definition_id = _int_or_none(parameters.get("definition_id")) if row.job_type == "job_definition" else None
|
||||
revision = _int_or_none(parameters.get("definition_revision")) if row.job_type == "job_definition" else None
|
||||
state_key = (
|
||||
f"job_definition:{definition_id}"
|
||||
if row.job_type == "job_definition" and definition_id
|
||||
else (str(row.job_type or "software_inventory")[:120])
|
||||
)
|
||||
key = (int(row.asset_id), state_key)
|
||||
current = states.get(key)
|
||||
if current is None:
|
||||
current = {
|
||||
"asset_id": int(row.asset_id),
|
||||
"state_key": state_key,
|
||||
"job_type": str(row.job_type or "software_inventory")[:50],
|
||||
"job_definition_id": definition_id,
|
||||
"definition_revision": revision,
|
||||
"last_success_revision": None,
|
||||
"last_job_id": int(row.id),
|
||||
"last_status": str(row.status or "created")[:30],
|
||||
"last_status_at": row.finished_at or row.callback_completed_at or row.started_at or row.sent_at or row.created_at or datetime.utcnow(),
|
||||
"last_created_at": row.created_at,
|
||||
"last_sent_at": row.sent_at,
|
||||
"last_started_at": row.started_at,
|
||||
"last_completed_at": row.finished_at or row.callback_completed_at,
|
||||
"last_success_at": None,
|
||||
"execution_count": 0,
|
||||
"updated_at": datetime.utcnow(),
|
||||
}
|
||||
states[key] = current
|
||||
|
||||
current.update({
|
||||
"job_type": str(row.job_type or "software_inventory")[:50],
|
||||
"job_definition_id": definition_id,
|
||||
"definition_revision": revision,
|
||||
"last_job_id": int(row.id),
|
||||
"last_status": str(row.status or "created")[:30],
|
||||
"last_status_at": row.finished_at or row.callback_completed_at or row.started_at or row.sent_at or row.created_at or datetime.utcnow(),
|
||||
"last_created_at": row.created_at,
|
||||
"last_sent_at": row.sent_at,
|
||||
"last_started_at": row.started_at,
|
||||
"last_completed_at": row.finished_at or row.callback_completed_at,
|
||||
"updated_at": datetime.utcnow(),
|
||||
})
|
||||
current["execution_count"] = int(current.get("execution_count") or 0) + max(1, int(row.attempt_count or 0))
|
||||
if str(row.status or "") == "success":
|
||||
success_at = row.finished_at or row.callback_completed_at or row.started_at or row.sent_at or row.created_at
|
||||
if success_at and (current.get("last_success_at") is None or success_at >= current["last_success_at"]):
|
||||
current["last_success_at"] = success_at
|
||||
current["last_success_revision"] = revision
|
||||
|
||||
db.bulk_insert_mappings(AssetJobState, list(states.values()))
|
||||
db.commit()
|
||||
return len(states)
|
||||
Reference in New Issue
Block a user