215 lines
9.0 KiB
Go
215 lines
9.0 KiB
Go
package validator
|
|
|
|
import (
|
|
"fmt"
|
|
"strings"
|
|
|
|
"browser.local/platform/domain"
|
|
)
|
|
|
|
const maxJobChannelMessageLength = 256
|
|
|
|
func ValidateRunJobClaim(claim domain.RunJobClaim) error {
|
|
var violations []string
|
|
violations = appendRequired(violations, "runEndpointId", claim.RunEndpointID)
|
|
violations = appendRequired(violations, "sessionToken", claim.SessionToken)
|
|
violations = appendCapacityViolations(violations, claim.Capacity)
|
|
for i, capability := range claim.Capabilities {
|
|
if strings.TrimSpace(capability) == "" {
|
|
violations = append(violations, fmt.Sprintf("capabilities[%d] is required", i))
|
|
}
|
|
}
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateRunJobAck(ack domain.RunJobAck) error {
|
|
var violations []string
|
|
violations = appendLeaseFields(violations, ack.RunEndpointID, ack.SessionToken, ack.JobID, ack.LeaseToken, ack.Attempt)
|
|
violations = appendMessageLength(violations, "message", ack.Message)
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateRunJobProgress(progress domain.RunJobProgress) error {
|
|
var violations []string
|
|
violations = appendLeaseFields(violations, progress.RunEndpointID, progress.SessionToken, progress.JobID, progress.LeaseToken, progress.Attempt)
|
|
violations = appendProgressViolations(violations, progress.Progress)
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateRunJobResult(result domain.RunJobResult) error {
|
|
var violations []string
|
|
violations = appendLeaseFields(violations, result.RunEndpointID, result.SessionToken, result.JobID, result.LeaseToken, result.Attempt)
|
|
if !validTerminalJobState(result.State) {
|
|
violations = append(violations, "state must be succeeded, failed, or cancelled")
|
|
}
|
|
violations = appendProgressViolations(violations, result.Progress)
|
|
violations = appendMessageLength(violations, "message", result.Message)
|
|
violations = appendMessageLength(violations, "errorCode", result.ErrorCode)
|
|
if len([]byte(result.ExecutionResult.Content)) > maxJobChannelMessageLength*256 {
|
|
violations = append(violations, "executionResult.content is too large")
|
|
}
|
|
if result.ExecutionResult.Checksum != "" && !validSHA256Checksum(result.ExecutionResult.Checksum) {
|
|
violations = append(violations, "executionResult.checksum must be sha256:<hex>")
|
|
}
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateRunLifecycleReport(report domain.RunLifecycleReport) error {
|
|
var violations []string
|
|
violations = appendRequired(violations, "runEndpointId", report.RunEndpointID)
|
|
violations = appendRequired(violations, "sessionToken", report.SessionToken)
|
|
violations = appendRequired(violations, "serverInstanceId", report.ServerInstanceID)
|
|
violations = appendRequired(violations, "capability", report.Capability)
|
|
if !validLifecycleReportCapability(report.Capability) {
|
|
violations = append(violations, "capability must be process.install, process.start, process.stop, or process.status")
|
|
}
|
|
if !validTerminalJobState(report.State) {
|
|
violations = append(violations, "state must be succeeded, failed, or cancelled")
|
|
}
|
|
violations = appendProgressViolations(violations, report.Progress)
|
|
violations = appendMessageLength(violations, "message", report.Message)
|
|
violations = appendMessageLength(violations, "errorCode", report.ErrorCode)
|
|
if report.ManagedProcessID != "" && report.ObservationSeq == 0 {
|
|
violations = append(violations, "observationSeq is required with managedProcessId")
|
|
}
|
|
if report.ObservationSeq > 0 && report.ManagedProcessID == "" {
|
|
violations = append(violations, "managedProcessId is required with observationSeq")
|
|
}
|
|
if len([]byte(report.ExecutionResult.Content)) > maxJobChannelMessageLength*256 {
|
|
violations = append(violations, "executionResult.content is too large")
|
|
}
|
|
if report.ExecutionResult.Checksum != "" && !validSHA256Checksum(report.ExecutionResult.Checksum) {
|
|
violations = append(violations, "executionResult.checksum must be sha256:<hex>")
|
|
}
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateDistributionBuildInputRequest(request domain.DistributionBuildInputRequest) error {
|
|
var violations []string
|
|
violations = appendLeaseFields(violations, request.RunEndpointID, request.SessionToken, request.JobID, request.LeaseToken, request.Attempt)
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateDependencyExecutionInputRequest(request domain.DependencyExecutionInputRequest) error {
|
|
return finish(appendLeaseFields(nil, request.RunEndpointID, request.SessionToken, request.JobID, request.LeaseToken, request.Attempt))
|
|
}
|
|
|
|
func ValidateSourceRCONExecutionInputRequest(request domain.SourceRCONExecutionInputRequest) error {
|
|
return finish(appendLeaseFields(nil, request.RunEndpointID, request.SessionToken, request.JobID, request.LeaseToken, request.Attempt))
|
|
}
|
|
|
|
func ValidateRunUpdateInputRequest(request domain.RunUpdateInputRequest) error {
|
|
return finish(appendLeaseFields(nil, request.RunEndpointID, request.SessionToken, request.JobID, request.LeaseToken, request.Attempt))
|
|
}
|
|
|
|
func ValidateRunUpdateChunkRequest(request domain.RunUpdateChunkRequest) error {
|
|
violations := appendLeaseFields(nil, request.RunEndpointID, request.SessionToken, request.JobID, request.LeaseToken, request.Attempt)
|
|
if request.Offset < 0 {
|
|
violations = append(violations, "offset must not be negative")
|
|
}
|
|
if request.Length <= 0 || request.Length > 1024*1024 {
|
|
violations = append(violations, "length must be between 1 and 1048576")
|
|
}
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateRunUpdateHealthReport(report domain.RunUpdateHealthReport) error {
|
|
violations := appendLeaseFields(nil, report.RunEndpointID, report.SessionToken, report.JobID, report.LeaseToken, report.Attempt)
|
|
if report.Outcome != "succeeded" && report.Outcome != "rolled-back" {
|
|
violations = append(violations, "outcome must be succeeded or rolled-back")
|
|
}
|
|
violations = appendRequired(violations, "version", report.Version)
|
|
violations = appendMessageLength(violations, "version", report.Version)
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateRunJobCancelRequest(request domain.RunJobCancelRequest) error {
|
|
var violations []string
|
|
violations = appendRequired(violations, "jobId", request.JobID)
|
|
violations = appendRequired(violations, "reason", request.Reason)
|
|
violations = appendMessageLength(violations, "reason", request.Reason)
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateRunJobCancelPoll(poll domain.RunJobCancelPoll) error {
|
|
var violations []string
|
|
violations = appendLeaseFields(violations, poll.RunEndpointID, poll.SessionToken, poll.JobID, poll.LeaseToken, poll.Attempt)
|
|
return finish(violations)
|
|
}
|
|
|
|
func ValidateRunJobReconcile(reconcile domain.RunJobReconcile) error {
|
|
var violations []string
|
|
violations = appendRequired(violations, "runEndpointId", reconcile.RunEndpointID)
|
|
violations = appendRequired(violations, "sessionToken", reconcile.SessionToken)
|
|
seen := map[string]struct{}{}
|
|
for i, entry := range reconcile.ActiveJobs {
|
|
prefix := fmt.Sprintf("activeJobs[%d]", i)
|
|
violations = appendRequired(violations, prefix+".jobId", entry.JobID)
|
|
violations = appendRequired(violations, prefix+".leaseToken", entry.LeaseToken)
|
|
if entry.Attempt <= 0 {
|
|
violations = append(violations, prefix+".attempt must be positive")
|
|
}
|
|
if _, exists := seen[entry.JobID]; exists {
|
|
violations = append(violations, fmt.Sprintf("%s duplicates %q", prefix, entry.JobID))
|
|
}
|
|
seen[entry.JobID] = struct{}{}
|
|
}
|
|
return finish(violations)
|
|
}
|
|
|
|
func appendLeaseFields(violations []string, runEndpointID string, sessionToken string, jobID string, leaseToken string, attempt int) []string {
|
|
violations = appendRequired(violations, "runEndpointId", runEndpointID)
|
|
violations = appendRequired(violations, "sessionToken", sessionToken)
|
|
violations = appendRequired(violations, "jobId", jobID)
|
|
violations = appendRequired(violations, "leaseToken", leaseToken)
|
|
if attempt <= 0 {
|
|
violations = append(violations, "attempt must be positive")
|
|
}
|
|
return violations
|
|
}
|
|
|
|
func appendProgressViolations(violations []string, progress domain.RunJobProgressReport) []string {
|
|
if progress.Percent < 0 || progress.Percent > 100 {
|
|
violations = append(violations, "progress.percent must be between 0 and 100")
|
|
}
|
|
if progress.Phase != "" && !validDeploymentProgressPhase(progress.Phase) {
|
|
violations = append(violations, "progress.phase is invalid")
|
|
}
|
|
violations = appendMessageLength(violations, "progress.message", progress.Message)
|
|
return violations
|
|
}
|
|
|
|
func validDeploymentProgressPhase(phase string) bool {
|
|
switch phase {
|
|
case "queued", "claimed", "preflight", "install", "configure", "start", "stop", "status", "health":
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
func appendMessageLength(violations []string, field string, message string) []string {
|
|
if len(message) > maxJobChannelMessageLength {
|
|
violations = append(violations, field+" is too long")
|
|
}
|
|
return violations
|
|
}
|
|
|
|
func validTerminalJobState(state domain.JobState) bool {
|
|
switch state {
|
|
case domain.JobStateSucceeded, domain.JobStateFailed, domain.JobStateCancelled:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
func validLifecycleReportCapability(capability string) bool {
|
|
switch capability {
|
|
case domain.LifecycleCapabilityInstall, domain.LifecycleCapabilityStart, domain.LifecycleCapabilityStop, domain.LifecycleCapabilityStatus:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|