task_contract.py 2.9 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182
  1. """Unified task state vocabulary + adaptive-batch field contract (P3-M4).
  2. Reconciles the two scheduling vocabularies that existed in parallel:
  3. - TaskManager / Task ORM:
  4. pending -> dispatched -> running -> completed / failed / cancelled
  5. - BatchScheduler:
  6. queued -> running -> completed / failed / cancelled
  7. This module is the single source of truth for status normalization and the
  8. adaptive-batch fields shared across TaskManager, BatchScheduler and the
  9. adaptive orchestrator. It must stay dependency-free so any of those services
  10. can import it without cycles.
  11. All source is ASCII only.
  12. """
  13. from typing import Any, Dict, Optional
  14. # Canonical statuses (TaskManager vocabulary).
  15. TASK_STATUS_PENDING = "pending"
  16. TASK_STATUS_DISPATCHED = "dispatched"
  17. TASK_STATUS_RUNNING = "running"
  18. TASK_STATUS_COMPLETED = "completed"
  19. TASK_STATUS_COMPLETED_WITH_ERRORS = "completed_with_errors"
  20. TASK_STATUS_FAILED = "failed"
  21. TASK_STATUS_CANCELLED = "cancelled"
  22. # BatchScheduler-only vocabulary.
  23. TASK_STATUS_QUEUED = "queued"
  24. TERMINAL_STATUSES = {
  25. TASK_STATUS_COMPLETED,
  26. TASK_STATUS_COMPLETED_WITH_ERRORS,
  27. TASK_STATUS_FAILED,
  28. TASK_STATUS_CANCELLED,
  29. }
  30. # scheduler / legacy word -> canonical word
  31. STATUS_ALIASES = {
  32. TASK_STATUS_QUEUED: TASK_STATUS_PENDING,
  33. TASK_STATUS_PENDING: TASK_STATUS_PENDING,
  34. TASK_STATUS_DISPATCHED: TASK_STATUS_DISPATCHED,
  35. TASK_STATUS_RUNNING: TASK_STATUS_RUNNING,
  36. TASK_STATUS_COMPLETED: TASK_STATUS_COMPLETED,
  37. TASK_STATUS_COMPLETED_WITH_ERRORS: TASK_STATUS_COMPLETED,
  38. TASK_STATUS_FAILED: TASK_STATUS_FAILED,
  39. TASK_STATUS_CANCELLED: TASK_STATUS_CANCELLED,
  40. "canceled": TASK_STATUS_CANCELLED,
  41. }
  42. # Adaptive-batch fields introduced in P3-M2, shared by Task ORM,
  43. # BatchScheduler and AdaptiveOrchestrator.
  44. ADAPTIVE_BATCH_FIELDS = ("task_type", "loop_id", "batch_id", "point_ids", "dynamic")
  45. def normalize_status(status: Optional[str]) -> Optional[str]:
  46. """Map any scheduler/legacy status word to the canonical status.
  47. Unknown words pass through unchanged so validation can report them.
  48. """
  49. if status is None:
  50. return None
  51. return STATUS_ALIASES.get(str(status), str(status))
  52. def is_terminal(status: Optional[str]) -> bool:
  53. return normalize_status(status) in TERMINAL_STATUSES
  54. def merge_adaptive_fields(task_info: Dict[str, Any], **kwargs) -> Dict[str, Any]:
  55. """Merge adaptive-batch fields into a task dict (no-ops stay default).
  56. Accepts kwargs that may include task_type/loop_id/batch_id/point_ids/
  57. dynamic and sets the canonical defaults when absent.
  58. """
  59. out = dict(task_info)
  60. out.setdefault("task_type", kwargs.get("task_type", "scan"))
  61. out.setdefault("loop_id", kwargs.get("loop_id"))
  62. out.setdefault("batch_id", kwargs.get("batch_id"))
  63. out.setdefault("point_ids", kwargs.get("point_ids") or [])
  64. out.setdefault("dynamic", bool(kwargs.get("dynamic", False)))
  65. return out