장시간 무응답 attempt를 안전하게 fence하고 provider health와 분리 관측해야 중복 출력 없이 기존 recovery budget으로 재실행할 수 있다.
147 lines
4.8 KiB
Go
147 lines
4.8 KiB
Go
package node
|
|
|
|
import (
|
|
"fmt"
|
|
|
|
"iop/apps/node/internal/transport"
|
|
runtime "iop/packages/go/execution"
|
|
iop "iop/proto/gen/iop"
|
|
)
|
|
|
|
// runRequestFromProto is the Edge-Node wire boundary. Common runtime packages
|
|
// remain independent of protobuf and Node transport details. ResponseStallTimeoutMS
|
|
// is deliberately left unset here: the handler validates the raw wire value and
|
|
// assigns the effective timeout via ValidateStallTimeoutOnWire before routing.
|
|
func runRequestFromProto(req *iop.RunRequest) runtime.RunRequest {
|
|
return runtime.RunRequest{
|
|
RunID: req.GetRunId(),
|
|
Adapter: req.GetAdapter(),
|
|
Target: req.GetTarget(),
|
|
SessionID: req.GetSessionId(),
|
|
Background: req.GetBackground(),
|
|
Policy: structAsMap(req.GetPolicy()),
|
|
Input: structAsMap(req.GetInput()),
|
|
TimeoutSec: int(req.GetTimeoutSec()),
|
|
Metadata: req.GetMetadata(),
|
|
}
|
|
}
|
|
|
|
var allowlistedLivenessMetadataKeys = map[string]bool{
|
|
"failure_code": true,
|
|
"provider_health": true,
|
|
"liveness_classification": true,
|
|
"idle_duration_ms": true,
|
|
"run_id": true,
|
|
"attempt_id": true,
|
|
"attempt_fence": true,
|
|
"adapter": true,
|
|
"target": true,
|
|
"health_observation_seq": true,
|
|
}
|
|
|
|
func allowlistedLivenessMetadata(metadata map[string]string) map[string]string {
|
|
if len(metadata) == 0 {
|
|
return nil
|
|
}
|
|
var filtered map[string]string
|
|
for k, v := range metadata {
|
|
if allowlistedLivenessMetadataKeys[k] {
|
|
if filtered == nil {
|
|
filtered = make(map[string]string)
|
|
}
|
|
filtered[k] = v
|
|
}
|
|
}
|
|
return filtered
|
|
}
|
|
|
|
func executionFailureToProto(failure *runtime.Failure) *iop.ExecutionFailure {
|
|
if failure == nil || failure.Code != runtime.FailureCodeResponseStalled {
|
|
return nil
|
|
}
|
|
msg := failure.Message
|
|
if msg == "" {
|
|
msg = failure.Error()
|
|
}
|
|
return &iop.ExecutionFailure{
|
|
Code: string(failure.Code),
|
|
Message: msg,
|
|
Retryable: failure.Retryable,
|
|
Metadata: allowlistedLivenessMetadata(failure.Metadata),
|
|
}
|
|
}
|
|
|
|
// runEventToProto preserves the existing Edge-Node event values while
|
|
// translating the host-neutral common event into the Node wire response.
|
|
func runEventToProto(event runtime.RuntimeEvent, nodeID, sessionID string, background bool) *iop.RunEvent {
|
|
errorMessage := event.Error
|
|
if errorMessage == "" && event.Failure != nil {
|
|
errorMessage = event.Failure.Error()
|
|
}
|
|
wireEvent := &iop.RunEvent{
|
|
RunId: event.RunID,
|
|
Type: string(event.Type),
|
|
Delta: event.Delta,
|
|
Message: event.Message,
|
|
Error: errorMessage,
|
|
Failure: executionFailureToProto(event.Failure),
|
|
Metadata: event.Metadata,
|
|
Timestamp: event.Timestamp.UnixNano(),
|
|
SessionId: sessionID,
|
|
Background: background,
|
|
NodeId: nodeID,
|
|
}
|
|
if event.Usage != nil {
|
|
wireEvent.Usage = &iop.Usage{
|
|
InputTokens: int32(event.Usage.InputTokens),
|
|
OutputTokens: int32(event.Usage.OutputTokens),
|
|
ReasoningTokens: int32(event.Usage.ReasoningTokens),
|
|
CachedInputTokens: int32(event.Usage.CachedInputTokens),
|
|
}
|
|
}
|
|
return wireEvent
|
|
}
|
|
|
|
// ValidateStallTimeoutOnWire validates a raw wire value and returns the
|
|
// effective timeout before the request reaches the router or provider. Zero
|
|
// resolves to the documented default; safe positive values pass through;
|
|
// negative and overflow values are rejected instead of being silently defaulted.
|
|
func ValidateStallTimeoutOnWire(ms int64) (int64, error) {
|
|
effective, err := runtime.ResolveStallTimeoutMS(ms)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("response_stall_timeout_ms: %w", err)
|
|
}
|
|
return effective, nil
|
|
}
|
|
|
|
func applyValidatedRunStallTimeout(req *iop.RunRequest, runReq *runtime.RunRequest) error {
|
|
effective, err := ValidateStallTimeoutOnWire(req.GetResponseStallTimeoutMs())
|
|
if err != nil {
|
|
return err
|
|
}
|
|
runReq.ResponseStallTimeoutMS = effective
|
|
return nil
|
|
}
|
|
|
|
func (n *Node) validateRunStallTimeout(sess *transport.Session, req *iop.RunRequest, runReq *runtime.RunRequest) error {
|
|
if err := applyValidatedRunStallTimeout(req, runReq); err != nil {
|
|
n.sendPreExecuteError(sess, req.GetRunId(), req.GetSessionId(), req.GetBackground(), n.nodeID, err.Error())
|
|
return fmt.Errorf("node: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func providerTunnelRequestFromProto(req *iop.ProviderTunnelRequest) (runtime.ProviderTunnelRequest, error) {
|
|
tr := runtime.ProviderTunnelRequest{
|
|
RunID: req.GetRunId(), TunnelID: req.GetTunnelId(), Adapter: req.GetAdapter(), Target: req.GetTarget(),
|
|
Method: req.GetMethod(), Path: req.GetPath(), Operation: req.GetOperation(), Headers: req.GetHeaders(),
|
|
Body: req.GetBody(), Stream: req.GetStream(), TimeoutSec: int(req.GetTimeoutSec()), Metadata: req.GetMetadata(),
|
|
SessionID: req.GetSessionId(),
|
|
}
|
|
effective, err := ValidateStallTimeoutOnWire(req.GetResponseStallTimeoutMs())
|
|
if err != nil {
|
|
return tr, err
|
|
}
|
|
tr.ResponseStallTimeoutMS = effective
|
|
return tr, nil
|
|
}
|