Skip to content

Commit 2e60cce

Browse files
Merge pull request #577 from backend-developers-ltd/job-error-propagation
Job error propagation
2 parents f574dff + f48c319 commit 2e60cce

47 files changed

Lines changed: 1964 additions & 1305 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

compute_horde/compute_horde/fv_protocol/validator_requests.py

Lines changed: 53 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,18 @@
1-
from enum import StrEnum
21
from typing import Any, Literal, Self
32

43
import bittensor
54
from pydantic import BaseModel
65

6+
from compute_horde.protocol_consts import (
7+
HordeFailureReason,
8+
JobFailureReason,
9+
JobParticipantType,
10+
JobRejectionReason,
11+
JobStage,
12+
JobStatus,
13+
)
14+
from compute_horde.protocol_messages import FailureContext
15+
716

817
class V0Heartbeat(BaseModel, extra="forbid"):
918
"""Message sent from validator to facilitator to keep connection alive"""
@@ -39,46 +48,71 @@ def ss58_address(self) -> str:
3948
return address
4049

4150

42-
class MinerResponse(BaseModel, extra="allow"):
43-
job_uuid: str
44-
message_type: str | None
51+
class JobResultDetails(BaseModel, extra="allow"):
4552
docker_process_stderr: str
4653
docker_process_stdout: str
4754
artifacts: dict[str, str] | None = None
4855
upload_results: dict[str, str] | None = None
4956

5057

58+
class JobRejectionDetails(BaseModel):
59+
rejected_by: JobParticipantType
60+
reason: JobRejectionReason
61+
message: str
62+
context: FailureContext | None = None
63+
64+
65+
class JobFailureDetails(BaseModel):
66+
reason: JobFailureReason
67+
stage: JobStage
68+
message: str
69+
context: FailureContext | None = None
70+
docker_process_exit_status: int | None = None
71+
docker_process_stdout: str | None = None
72+
docker_process_stderr: str | None = None
73+
74+
75+
class HordeFailureDetails(BaseModel):
76+
reported_by: JobParticipantType
77+
reason: HordeFailureReason
78+
message: str
79+
context: FailureContext | None = None
80+
81+
5182
class StreamingServerDetails(BaseModel, extra="forbid"):
5283
streaming_server_cert: str | None = None
5384
streaming_server_address: str | None = None
5485
streaming_server_port: int | None = None
5586

5687

5788
class JobStatusMetadata(BaseModel, extra="allow"):
58-
comment: str
59-
miner_response: MinerResponse | None = None
89+
miner_response: JobResultDetails | None = None
90+
job_rejection_details: JobRejectionDetails | None = None
91+
job_failure_details: JobFailureDetails | None = None
92+
horde_failure_details: HordeFailureDetails | None = None
6093
streaming_details: StreamingServerDetails | None = None
6194

95+
@classmethod
96+
def from_uncaught_exception(cls, reported_by: JobParticipantType, exception: Exception) -> Self:
97+
return cls(
98+
horde_failure_details=HordeFailureDetails(
99+
reported_by=reported_by,
100+
reason=HordeFailureReason.UNHANDLED_EXCEPTION,
101+
message="Uncaught exception",
102+
context={"exception_type": type(exception).__qualname__},
103+
),
104+
)
105+
62106

63107
class JobStatusUpdate(BaseModel, extra="forbid"):
108+
# TODO(post error propagation): remove "extra"
64109
"""
65-
Message sent from validator to facilitator in response to NewJobRequest.
110+
Message sent from validator to facilitator when the job's state changes.
66111
"""
67112

68-
class Status(StrEnum):
69-
RECEIVED = "received"
70-
ACCEPTED = "accepted"
71-
EXECUTOR_READY = "executor_ready"
72-
VOLUMES_READY = "volumes_ready"
73-
EXECUTION_DONE = "execution_done"
74-
COMPLETED = "completed"
75-
REJECTED = "rejected"
76-
FAILED = "failed"
77-
STREAMING_READY = "streaming_ready"
78-
79113
message_type: Literal["V0JobStatusUpdate"] = "V0JobStatusUpdate"
80114
uuid: str
81-
status: Status
115+
status: JobStatus
82116
metadata: JobStatusMetadata | None = None
83117

84118

Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,86 @@
1+
import logging
2+
from typing import Generic, TypeVar
3+
4+
from compute_horde.protocol_consts import HordeFailureReason, JobFailureReason, JobRejectionReason
5+
from compute_horde.protocol_messages import FailureContext
6+
7+
ReasonType = TypeVar("ReasonType")
8+
9+
logger = logging.getLogger(__name__)
10+
11+
12+
class BaseJobError(Exception, Generic[ReasonType]):
13+
"""
14+
Common structure for job-related errors.
15+
"""
16+
17+
def __init__(self, message: str, reason: ReasonType, context: FailureContext | None = None):
18+
"""
19+
Parameters:
20+
message: str
21+
Short message for a human seeing the error somewhere.
22+
reason: ReasonType
23+
Structured error reason.
24+
context: FailureContext | None
25+
Optional arbitrary key-value context to be sent along the error.
26+
Avoid
27+
"""
28+
super().__init__(message)
29+
self.message = message
30+
self.reason = reason
31+
self.context = context
32+
33+
def add_context(self, context: FailureContext):
34+
if self.context is None:
35+
self.context = {**context}
36+
else:
37+
self.context.update(context)
38+
39+
def __str__(self):
40+
return f"{self.message}, {self.reason=}"
41+
42+
def __repr__(self):
43+
return f"{type(self).__name__}: {str(self)}"
44+
45+
46+
class JobRejection(BaseJobError[JobRejectionReason]):
47+
def __str__(self):
48+
return f"Job rejected ({self.reason}): {self.message} "
49+
50+
51+
class JobError(BaseJobError[JobFailureReason]):
52+
def __str__(self):
53+
return f"Job failed ({self.reason}): {self.message} "
54+
55+
56+
class HordeError(BaseJobError[HordeFailureReason]):
57+
def __init__(
58+
self,
59+
message: str,
60+
reason: HordeFailureReason = HordeFailureReason.GENERIC_ERROR,
61+
context: FailureContext | None = None,
62+
):
63+
super().__init__(message, reason, context)
64+
65+
def __str__(self):
66+
return f"Horde failed ({self.reason}): {self.message}"
67+
68+
@classmethod
69+
def wrap_unhandled(cls, e: Exception, context: FailureContext | None = None) -> "HordeError":
70+
"""
71+
Intended for "catch-all" except blocks - consistently wraps the exception in a HordeError.
72+
For catching errors somewhere in the job code, do instead:
73+
`raise HordeError("Spline reticulation failed") from e`
74+
"""
75+
if isinstance(e, HordeError):
76+
# Don't wrap another horde failure
77+
return e
78+
79+
if isinstance(e, BaseJobError):
80+
# If another job error is to be wrapped by accident - don't blow up, but scream for help
81+
logger.error(f"Wrapping a {type(e).__qualname__} in a HordeError", exc_info=False)
82+
83+
failure = cls("Unhandled exception", HordeFailureReason.UNHANDLED_EXCEPTION, context)
84+
failure.__cause__ = e
85+
failure.add_context({"exception_type": type(e).__qualname__})
86+
return failure

0 commit comments

Comments
 (0)