Files
run/runtime/worker.go
T

1688 lines
69 KiB
Go

package runtime
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"log"
"net/http"
"os"
"path/filepath"
"runtime"
"strings"
"sync"
"time"
"browser.local/run/config"
"browser.local/run/protocol"
"browser.local/run/spool"
)
type WorkerClient interface {
Hello(context.Context, protocol.RunHelloRequest) (protocol.RunHelloResponse, error)
Heartbeat(context.Context, protocol.RunHeartbeatRequest) (protocol.RunHeartbeatResponse, error)
StreamControlEvents(context.Context, protocol.RunControlStreamRequest, func(protocol.RunControlEvent) error) error
ReportLifecycle(context.Context, protocol.RunLifecycleReportRequest) (protocol.RunLifecycleReportResponse, error)
GetRunLogStreamProgress(context.Context, protocol.RunLogStreamProgressRequest) (protocol.RunLogStreamProgressResponse, error)
ClaimJob(context.Context, protocol.RunJobClaimRequest) (protocol.RunJobClaimResponse, error)
AckJob(context.Context, protocol.RunJobAckRequest) (protocol.RunJobAckResponse, error)
UpdateJobProgress(context.Context, protocol.RunJobProgressRequest) (protocol.RunJobProgressResponse, error)
CompleteJob(context.Context, protocol.RunJobResultRequest) (protocol.RunJobResultResponse, error)
GetDistributionBuildInput(context.Context, protocol.DistributionBuildInputRequest) (protocol.DistributionBuildInputResponse, error)
GetDependencyExecutionInput(context.Context, protocol.DependencyExecutionInputRequest) (protocol.DependencyExecutionInputResponse, error)
GetSourceRCONExecutionInput(context.Context, protocol.SourceRCONExecutionInputRequest) (protocol.SourceRCONExecutionInputResponse, error)
GetProtectedRequestExecutionInput(context.Context, protocol.ProtectedRequestExecutionInputRequest) (protocol.ProtectedRequestExecutionInputResponse, error)
GetRunUpdateInput(context.Context, protocol.RunUpdateInputRequest) (protocol.RunUpdateInputResponse, error)
ReadRunUpdateChunk(context.Context, protocol.RunUpdateChunkRequest) (protocol.RunUpdateChunkResponse, error)
ReportRunUpdateHealth(context.Context, protocol.RunUpdateHealthRequest) (protocol.RunUpdateHealthResponse, error)
IngestMetricBatch(context.Context, protocol.MetricBatchIngestRequest) (protocol.MetricBatchIngestResponse, error)
OpenArtifactTransfer(context.Context, protocol.ArtifactTransferOpenRequest) (protocol.ArtifactTransferOpenResponse, error)
UploadArtifactChunk(context.Context, protocol.ArtifactChunkUploadRequest) (protocol.ArtifactChunkUploadResponse, error)
CompleteArtifactTransfer(context.Context, protocol.ArtifactTransferCompleteRequest) (protocol.ArtifactTransferCompleteResponse, error)
PollJobCancel(context.Context, protocol.RunJobCancelPollRequest) (protocol.RunJobCancelPollResponse, error)
ReconcileJobs(context.Context, protocol.RunJobReconcileRequest) (protocol.RunJobReconcileResponse, error)
}
const (
jobActivePollInterval = 10 * time.Second
jobClaimWaitSeconds = 25
durableUploaderFlushTimeout = 5 * time.Second
fileArtifactChunkSize = 1024 * 1024
maxFileArtifactBytes = int64(512 * 1024 * 1024)
)
type Worker struct {
cfg config.Config
client WorkerClient
executor LifecycleExecutor
state WorkerState
journal *JobJournal
stateMu sync.RWMutex
sessionRefreshMu sync.Mutex
sequenceMu sync.Mutex
restartMu sync.Mutex
restartRequested bool
observationMu sync.Mutex
reportedObservations map[string]uint64
metricCollector MetricCollector
}
type WorkerState struct {
RunEndpointID string
SessionToken string
SessionExpiresAt time.Time
Capabilities []string
Capacity protocol.RunCapacityReport
LastHeartbeat time.Time
Sequence uint64
}
func NewWorker(cfg config.Config, client WorkerClient, options ...LifecycleExecutorOption) (*Worker, error) {
if client == nil {
return nil, fmt.Errorf("worker client is required")
}
if cfg.RunEndpointID == "" {
cfg.RunEndpointID = config.DefaultEndpointID
}
if cfg.DisplayName == "" {
cfg.DisplayName = config.DefaultDisplayName
}
if cfg.Version == "" {
cfg.Version = config.DefaultVersion
}
if cfg.MaxJobs <= 0 {
cfg.MaxJobs = 1
}
executorOptions := append([]LifecycleExecutorOption{
WithLifecycleWorkspaceRoot(cfg.WorkspaceRoot),
WithManagedProcessStateRoot(managedProcessStateRoot(cfg)),
WithManagedProcessOutputRoot(cfg.WorkspaceRoot),
WithLocalStartupDiagnostics(cfg.LocalStartupDiagnostics),
}, options...)
if err := migrateManagedProcessState(cfg, managedProcessStateRoot(cfg)); err != nil {
return nil, err
}
journal, err := NewPersistentJobJournal(cfg.WorkspaceRoot)
if err != nil {
return nil, err
}
worker := &Worker{
cfg: cfg,
client: client,
executor: NewLifecycleExecutor(executorOptions...),
state: WorkerState{
RunEndpointID: cfg.RunEndpointID,
Capabilities: SupportedRunCapabilitiesForComponent(cfg.ComponentKind),
Capacity: protocol.RunCapacityReport{MaxJobs: cfg.MaxJobs},
},
journal: journal,
reportedObservations: map[string]uint64{},
metricCollector: defaultMetricCollector{},
}
if worker.executor.metricCollector != nil {
worker.metricCollector = worker.executor.metricCollector
}
return worker, nil
}
func migrateLegacyManagedProcessState(cfg config.Config, stateRoot string) error {
workspaceRoot := cfg.WorkspaceRoot
if strings.TrimSpace(workspaceRoot) == "" {
workspaceRoot = filepath.Join(".", ".run-workspace")
}
legacyPath := filepath.Join(workspaceRoot, "state", "processes.json")
newPath := filepath.Join(stateRoot, "state", "processes.json")
if filepath.Clean(legacyPath) == filepath.Clean(newPath) {
return nil
}
if _, err := os.Stat(newPath); err == nil {
return nil
} else if !os.IsNotExist(err) {
return err
}
body, err := os.ReadFile(legacyPath)
if os.IsNotExist(err) {
return nil
}
if err != nil {
return fmt.Errorf("read legacy managed process state: %w", err)
}
var legacy processJournal
if err := json.Unmarshal(body, &legacy); err != nil {
return fmt.Errorf("decode legacy managed process state: %w", err)
}
filtered := processJournal{Version: managedProcessJournalVersion, Items: map[string]ProcessIdentity{}, Retired: map[string]ProcessIdentity{}}
for key, item := range legacy.Items {
if managedProcessBelongsToRun(item, cfg) {
filtered.Items[key] = item
}
}
for key, item := range legacy.Retired {
if managedProcessBelongsToRun(item, cfg) {
filtered.Retired[key] = item
}
}
if len(filtered.Items) == 0 && len(filtered.Retired) == 0 {
return nil
}
if err := ensureDirectory(filepath.Dir(newPath)); err != nil {
return fmt.Errorf("create isolated managed process state: %w", err)
}
encoded, err := json.Marshal(filtered)
if err != nil {
return fmt.Errorf("encode isolated managed process state: %w", err)
}
if err := os.WriteFile(newPath, encoded, 0o600); err != nil {
return fmt.Errorf("write isolated managed process state: %w", err)
}
log.Printf("RUN phase=process.managed status=legacy_state_migrated items=%d retired=%d", len(filtered.Items), len(filtered.Retired))
return nil
}
// migrateManagedProcessState also imports matching journals from prior Run
// state namespaces. A generated package can change its component metadata
// while retaining the same endpoint/server/profile; that must not create a
// second autonomous process for the same server during a Run update.
func migrateManagedProcessState(cfg config.Config, stateRoot string) error {
if err := migrateLegacyManagedProcessState(cfg, stateRoot); err != nil {
return err
}
workspaceRoot := cfg.WorkspaceRoot
if strings.TrimSpace(workspaceRoot) == "" {
workspaceRoot = filepath.Join(".", ".run-workspace")
}
targetPath := filepath.Join(stateRoot, "state", "processes.json")
target := processJournal{Version: managedProcessJournalVersion, Items: map[string]ProcessIdentity{}, Retired: map[string]ProcessIdentity{}}
if body, err := os.ReadFile(targetPath); err == nil {
if err := json.Unmarshal(body, &target); err != nil {
return fmt.Errorf("decode managed process state: %w", err)
}
if target.Items == nil {
target.Items = map[string]ProcessIdentity{}
}
if target.Retired == nil {
target.Retired = map[string]ProcessIdentity{}
}
} else if !os.IsNotExist(err) {
return fmt.Errorf("read managed process state: %w", err)
}
changed := false
err := filepath.WalkDir(workspaceRoot, func(path string, entry os.DirEntry, walkErr error) error {
if walkErr != nil {
return walkErr
}
if entry.IsDir() || entry.Name() != "processes.json" {
return nil
}
if filepath.Clean(path) == filepath.Clean(targetPath) {
return nil
}
body, err := os.ReadFile(path)
if err != nil {
return err
}
var source processJournal
if err := json.Unmarshal(body, &source); err != nil {
return fmt.Errorf("decode managed process state candidate: %w", err)
}
for key, item := range source.Items {
if !managedProcessBelongsToRun(item, cfg) {
continue
}
existing, exists := target.Items[key]
if !exists || (!processAlivePID(existing.PID) && processAlivePID(item.PID)) {
target.Items[key] = item
changed = true
}
}
for key, item := range source.Retired {
if managedProcessBelongsToRun(item, cfg) {
if _, exists := target.Retired[key]; !exists {
target.Retired[key] = item
changed = true
}
}
}
return nil
})
if err != nil {
return fmt.Errorf("scan managed process state: %w", err)
}
if !changed {
return nil
}
if err := ensureDirectory(filepath.Dir(targetPath)); err != nil {
return fmt.Errorf("create managed process state directory: %w", err)
}
body, err := json.Marshal(target)
if err != nil {
return fmt.Errorf("encode managed process state: %w", err)
}
temporary := targetPath + ".tmp"
if err := os.WriteFile(temporary, body, 0o600); err != nil {
return fmt.Errorf("write managed process state: %w", err)
}
if err := os.Rename(temporary, targetPath); err != nil {
_ = os.Remove(temporary)
return fmt.Errorf("replace managed process state: %w", err)
}
log.Printf("RUN phase=process.managed status=state_namespaces_merged items=%d retired=%d", len(target.Items), len(target.Retired))
return nil
}
func managedProcessBelongsToRun(item ProcessIdentity, cfg config.Config) bool {
if item.RunEndpointID == "" || item.ServerInstanceID == "" || item.RunEndpointID != cfg.RunEndpointID || item.ServerInstanceID != cfg.ServerInstanceID {
return false
}
return cfg.ComponentKey == "" || item.ProfileKey == "" || item.ProfileKey == cfg.ComponentKey
}
func managedProcessStateRoot(cfg config.Config) string {
workspaceRoot := cfg.WorkspaceRoot
if strings.TrimSpace(workspaceRoot) == "" {
workspaceRoot = filepath.Join(".", ".run-workspace")
}
identity := strings.Join([]string{
"run-process-state-v1",
cfg.RunEndpointID,
cfg.ServerInstanceID,
cfg.PluginID,
cfg.ComponentKind,
cfg.ComponentKey,
}, "\x00")
digest := sha256.Sum256([]byte(identity))
return filepath.Join(workspaceRoot, "run-services", hex.EncodeToString(digest[:]))
}
func (worker *Worker) Register(ctx context.Context) error {
worker.sessionRefreshMu.Lock()
defer worker.sessionRefreshMu.Unlock()
return worker.registerUnlocked(ctx)
}
func (worker *Worker) registerUnlocked(ctx context.Context) error {
state := worker.State()
log.Printf("RUN phase=register status=starting endpoint=%s version=%s server=%s plugin=%s component=%s componentKey=%s capabilities=%d maxJobs=%d", worker.cfg.RunEndpointID, worker.cfg.Version, worker.cfg.ServerInstanceID, worker.cfg.PluginID, worker.cfg.ComponentKind, safeOptional(worker.cfg.ComponentKey), len(state.Capabilities), worker.cfg.MaxJobs)
response, err := worker.client.Hello(ctx, protocol.RunHelloRequest{
RegistrationToken: worker.cfg.RegistrationToken,
RunEndpointID: worker.cfg.RunEndpointID,
ServerInstanceID: worker.cfg.ServerInstanceID,
PluginID: worker.cfg.PluginID,
ComponentKind: worker.cfg.ComponentKind,
ComponentKey: worker.cfg.ComponentKey,
KeyGeneration: worker.cfg.KeyGeneration,
DisplayName: worker.cfg.DisplayName,
Version: worker.cfg.Version,
Status: "online",
Platform: runtime.GOOS,
Architecture: runtime.GOARCH,
CapabilityReport: protocol.RunCapabilityReport{
Capabilities: state.Capabilities,
Fingerprint: capabilityFingerprint(state.Capabilities),
},
Capacity: worker.capacityReportFor(state),
})
if err != nil {
log.Printf("RUN phase=register status=failed endpoint=%s error=%s", worker.cfg.RunEndpointID, err.Error())
return err
}
if !response.Accepted || response.SessionToken == "" {
log.Printf("RUN phase=register status=rejected endpoint=%s accepted=%t sessionTokenPresent=%t", worker.cfg.RunEndpointID, response.Accepted, response.SessionToken != "")
return fmt.Errorf("run hello was not accepted")
}
worker.stateMu.Lock()
worker.state.SessionToken = response.SessionToken
worker.state.SessionExpiresAt = response.SessionExpiresAt
state = worker.state
state.Capabilities = append([]string(nil), state.Capabilities...)
worker.stateMu.Unlock()
if sink, ok := worker.executor.logSink.(*SpoolLogSink); ok {
sink.mu.Lock()
sink.RunEndpointID = state.RunEndpointID
sink.SessionToken = state.SessionToken
sink.Progress = func(ctx context.Context, serverInstanceID string, streamID string) (uint64, error) {
current, err := worker.registeredState()
if err != nil {
return 0, err
}
response, err := worker.client.GetRunLogStreamProgress(ctx, protocol.RunLogStreamProgressRequest{RunEndpointID: current.RunEndpointID, SessionToken: current.SessionToken, ServerInstanceID: serverInstanceID, LogStreamID: streamID})
if err != nil {
return 0, err
}
if !response.Accepted || response.LogStreamID != streamID {
return 0, fmt.Errorf("run log stream progress was not accepted")
}
return response.LatestSeq, nil
}
sink.mu.Unlock()
}
if sink, ok := worker.executor.logSink.(*LiveLogSink); ok {
sink.SetSession(state.RunEndpointID, state.SessionToken)
}
if hook, ok := worker.executor.artifactHook.(*QueueArtifactHook); ok {
hook.RunEndpointID = state.RunEndpointID
hook.SessionToken = state.SessionToken
}
worker.executor.ResumeManagedProcessLogs(ctx)
log.Printf("RUN phase=register status=accepted endpoint=%s sessionExpiresAt=%s heartbeatSeconds=%d", state.RunEndpointID, response.SessionExpiresAt.Format(time.RFC3339), response.HeartbeatIntervalSeconds)
return nil
}
func (worker *Worker) registeredState() (WorkerState, error) {
state := worker.State()
if state.SessionToken == "" {
return state, fmt.Errorf("worker is not registered")
}
return state, nil
}
func (worker *Worker) HeartbeatOnce(ctx context.Context) error {
state, err := worker.registeredState()
if err != nil {
return err
}
if !state.SessionExpiresAt.IsZero() && !time.Now().UTC().Add(time.Minute).Before(state.SessionExpiresAt) {
log.Printf("RUN phase=heartbeat status=session_expiring endpoint=%s sessionExpiresAt=%s", state.RunEndpointID, state.SessionExpiresAt.Format(time.RFC3339))
return worker.reregisterAndReconcile(ctx, "heartbeat_session_expiring", state.SessionToken)
}
log.Printf("RUN phase=heartbeat status=starting endpoint=%s activeJobs=%d", state.RunEndpointID, worker.journal.ActiveCount())
response, err := worker.client.Heartbeat(ctx, protocol.RunHeartbeatRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
Version: worker.cfg.Version,
Status: "online",
CapabilityFingerprint: capabilityFingerprint(state.Capabilities),
Capacity: worker.capacityReportFor(state),
})
if err != nil {
if sessionInvalidError(err) {
log.Printf("RUN phase=heartbeat status=session_invalid endpoint=%s error=%s", state.RunEndpointID, err.Error())
return worker.reregisterAndReconcile(ctx, "heartbeat_session_invalid", state.SessionToken)
}
log.Printf("RUN phase=heartbeat status=failed endpoint=%s error=%s", state.RunEndpointID, err.Error())
return err
}
if !response.Accepted {
log.Printf("RUN phase=heartbeat status=rejected endpoint=%s", state.RunEndpointID)
return fmt.Errorf("heartbeat was not accepted")
}
worker.stateMu.Lock()
if worker.state.SessionToken == state.SessionToken {
worker.state.LastHeartbeat = response.ServerTime
}
worker.stateMu.Unlock()
log.Printf("RUN phase=heartbeat status=accepted endpoint=%s serverTime=%s nextSeconds=%d", state.RunEndpointID, response.ServerTime.Format(time.RFC3339), response.NextHeartbeatSeconds)
return nil
}
func (worker *Worker) reregisterAndReconcile(ctx context.Context, reason string, observedToken string) error {
worker.sessionRefreshMu.Lock()
defer worker.sessionRefreshMu.Unlock()
state := worker.State()
if observedToken != "" && state.SessionToken != "" && state.SessionToken != observedToken {
log.Printf("RUN phase=register status=already_refreshed reason=%s endpoint=%s", safeOptional(reason), state.RunEndpointID)
return worker.ReconcileOnce(ctx)
}
log.Printf("RUN phase=register status=refreshing reason=%s endpoint=%s", safeOptional(reason), state.RunEndpointID)
if err := worker.registerUnlocked(ctx); err != nil {
return err
}
return worker.ReconcileOnce(ctx)
}
func (worker *Worker) ClaimAndRunOnce(ctx context.Context) (bool, error) {
return worker.claimAndRunOnce(ctx, jobClaimWaitSeconds)
}
func (worker *Worker) claimAndRunOnce(ctx context.Context, waitSeconds int) (bool, error) {
state, err := worker.registeredState()
if err != nil {
return false, err
}
log.Printf("RUN phase=claim status=polling endpoint=%s activeJobs=%d", state.RunEndpointID, worker.journal.ActiveCount())
claim, err := worker.client.ClaimJob(ctx, protocol.RunJobClaimRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
Capabilities: state.Capabilities,
Capacity: worker.capacityReportFor(state),
WaitSeconds: waitSeconds,
})
if err != nil {
if sessionInvalidError(err) {
log.Printf("RUN phase=claim status=session_invalid endpoint=%s error=%s", state.RunEndpointID, err.Error())
return false, worker.reregisterAndReconcile(ctx, "claim_session_invalid", state.SessionToken)
}
log.Printf("RUN phase=claim status=failed endpoint=%s error=%s", state.RunEndpointID, err.Error())
return false, err
}
if !claim.Accepted || !claim.HasJob || claim.Job == nil {
log.Printf("RUN phase=claim status=idle endpoint=%s accepted=%t hasJob=%t", state.RunEndpointID, claim.Accepted, claim.HasJob)
return false, nil
}
assignment := *claim.Job
log.Printf("RUN phase=claim status=assigned job=%s capability=%s target=%s attempt=%d server=%s workspaceScope=%s", assignment.JobID, assignment.Capability, safeOptional(assignment.TargetKey), assignment.Attempt, assignment.ServerInstanceID, safeOptional(assignment.ExecutionInput.WorkspaceScope))
return true, worker.runAssignment(ctx, assignment)
}
func sessionInvalidError(err error) bool {
var sessionError interface{ SessionInvalid() bool }
return errors.As(err, &sessionError) && sessionError.SessionInvalid()
}
func (worker *Worker) runAssignment(ctx context.Context, assignment protocol.RunJobAssignment) error {
state, err := worker.registeredState()
if err != nil {
return err
}
if assignment.RunEndpointID != state.RunEndpointID {
return fmt.Errorf("job assignment endpoint does not match registered Run endpoint")
}
log.Printf("RUN phase=job status=journal_store job=%s capability=%s attempt=%d", assignment.JobID, assignment.Capability, assignment.Attempt)
if err := worker.journal.Store(assignment); err != nil {
log.Printf("RUN phase=job status=journal_store_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
log.Printf("RUN phase=job status=ack_start job=%s capability=%s", assignment.JobID, assignment.Capability)
ack, err := worker.client.AckJob(ctx, protocol.RunJobAckRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
JobID: assignment.JobID,
LeaseToken: assignment.LeaseToken,
Attempt: assignment.Attempt,
Message: "job accepted by run worker",
})
if err != nil {
log.Printf("RUN phase=job status=ack_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
assignment = ack.Job
if !ack.Accepted {
log.Printf("RUN phase=job status=ack_rejected job=%s", assignment.JobID)
return fmt.Errorf("job acknowledgement was not accepted")
}
log.Printf("RUN phase=job status=ack_accepted job=%s attempt=%d", assignment.JobID, assignment.Attempt)
if err := worker.journal.Store(assignment); err != nil {
log.Printf("RUN phase=job status=journal_store_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
progressSequence := worker.nextProgressSequence(assignment.ProgressSequence)
state, err = worker.registeredState()
if err != nil {
return err
}
log.Printf("RUN phase=job status=progress_start job=%s sequence=%d", assignment.JobID, progressSequence)
progress, err := worker.client.UpdateJobProgress(ctx, protocol.RunJobProgressRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
JobID: assignment.JobID,
LeaseToken: assignment.LeaseToken,
Attempt: assignment.Attempt,
Progress: protocol.RunJobProgressReport{Percent: 10, Message: "lifecycle execution started"},
Sequence: progressSequence,
})
if err != nil {
log.Printf("RUN phase=job status=progress_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
if !progress.Accepted {
log.Printf("RUN phase=job status=progress_rejected job=%s", assignment.JobID)
return fmt.Errorf("job progress was not accepted")
}
assignment = progress.Job
log.Printf("RUN phase=job status=progress_accepted job=%s percent=%d", assignment.JobID, assignment.Progress.Percent)
if err := worker.journal.Store(assignment); err != nil {
log.Printf("RUN phase=job status=journal_store_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
log.Printf("RUN phase=job status=execute_start job=%s capability=%s", assignment.JobID, assignment.Capability)
execution, assignment, cancelledByPlatform, err := worker.executeWithJobPolling(ctx, assignment)
if err != nil {
log.Printf("RUN phase=job status=execute_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
log.Printf("RUN phase=job status=execute_done job=%s state=%s errorCode=%s message=%s cancelledByPlatform=%t", assignment.JobID, execution.State, safeOptional(execution.ErrorCode), safeOptional(execution.Message), cancelledByPlatform)
if cancelledByPlatform && execution.State == lifecycleResultStateSucceeded {
execution = LifecycleExecutionResult{
State: lifecycleResultStateCancelled,
Progress: protocol.RunJobProgressReport{Percent: 100, Message: "cancelled by platform"},
Message: "cancelled by platform",
ErrorCode: "lifecycle_cancelled",
}
}
state, err = worker.registeredState()
if err != nil {
return err
}
resultRequest := LifecycleResultRequest(assignment, state.SessionToken, execution)
log.Printf("RUN phase=job status=result_store job=%s state=%s", assignment.JobID, resultRequest.State)
if err := worker.journal.StorePendingResult(resultRequest, execution.ActivationManifest); err != nil {
log.Printf("RUN phase=job status=result_store_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
log.Printf("RUN phase=job status=result_submit job=%s state=%s", assignment.JobID, resultRequest.State)
result, err := worker.client.CompleteJob(ctx, resultRequest)
if err != nil {
log.Printf("RUN phase=job status=result_submit_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
if !result.Accepted {
log.Printf("RUN phase=job status=result_rejected job=%s", assignment.JobID)
return fmt.Errorf("job result was not accepted")
}
log.Printf("RUN phase=job status=result_accepted job=%s state=%s", assignment.JobID, resultRequest.State)
if err := worker.journal.Delete(assignment.JobID); err != nil {
log.Printf("RUN phase=job status=journal_delete_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
log.Printf("RUN phase=job status=complete job=%s state=%s", assignment.JobID, resultRequest.State)
if execution.ActivationManifest != "" {
log.Printf("RUN phase=self_update status=activate_start job=%s", assignment.JobID)
if err := worker.executor.selfUpdateActivator.Activate(execution.ActivationManifest); err != nil {
log.Printf("RUN phase=self_update status=activate_failed job=%s error=%s", assignment.JobID, err.Error())
return fmt.Errorf("launch self-update helper: %w", err)
}
worker.restartMu.Lock()
worker.restartRequested = true
worker.restartMu.Unlock()
}
return nil
}
func (worker *Worker) executeWithJobPolling(ctx context.Context, assignment protocol.RunJobAssignment) (LifecycleExecutionResult, protocol.RunJobAssignment, bool, error) {
jobCtx, cancel := context.WithCancel(ctx)
defer cancel()
cancelledByPlatform := false
pollCancel := func() {
state, err := worker.registeredState()
if err != nil {
log.Printf("RUN phase=job.cancel_poll status=skipped_unregistered job=%s error=%s", assignment.JobID, err.Error())
return
}
log.Printf("RUN phase=job.cancel_poll status=starting job=%s", assignment.JobID)
response, err := worker.client.PollJobCancel(ctx, protocol.RunJobCancelPollRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
JobID: assignment.JobID,
LeaseToken: assignment.LeaseToken,
Attempt: assignment.Attempt,
})
if err != nil {
log.Printf("RUN phase=job.cancel_poll status=failed job=%s error=%s", assignment.JobID, err.Error())
return
}
if err == nil && response.HasCancel {
log.Printf("RUN phase=job.cancel_poll status=cancel_requested job=%s", assignment.JobID)
cancelledByPlatform = true
cancel()
return
}
log.Printf("RUN phase=job.cancel_poll status=clear job=%s", assignment.JobID)
}
pollCancel()
executionCh := make(chan LifecycleExecutionResult, 1)
executionAssignment := assignment
go func() {
executionCh <- worker.executeAssignment(jobCtx, executionAssignment)
}()
log.Printf("RUN phase=job.execute status=worker_started job=%s pollSeconds=%d", assignment.JobID, int(jobActivePollInterval/time.Second))
ticker := time.NewTicker(jobActivePollInterval)
defer ticker.Stop()
for {
select {
case execution := <-executionCh:
log.Printf("RUN phase=job.execute status=worker_finished job=%s state=%s", assignment.JobID, execution.State)
return execution, assignment, cancelledByPlatform, nil
case <-ctx.Done():
log.Printf("RUN phase=job.execute status=context_done job=%s error=%s", assignment.JobID, ctx.Err().Error())
cancel()
execution := <-executionCh
return execution, assignment, cancelledByPlatform, ctx.Err()
case <-ticker.C:
log.Printf("RUN phase=job.execute status=active job=%s percent=%d", assignment.JobID, assignment.Progress.Percent)
if !cancelledByPlatform {
pollCancel()
}
if cancelledByPlatform {
continue
}
state, err := worker.registeredState()
if err != nil {
cancel()
log.Printf("RUN phase=job.execute status=lease_renew_failed job=%s error=%s", assignment.JobID, err.Error())
return LifecycleExecutionResult{}, assignment, cancelledByPlatform, err
}
progress, err := worker.client.UpdateJobProgress(ctx, protocol.RunJobProgressRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
JobID: assignment.JobID,
LeaseToken: assignment.LeaseToken,
Attempt: assignment.Attempt,
Progress: protocol.RunJobProgressReport{Percent: assignment.Progress.Percent, Message: "lifecycle execution active"},
Sequence: worker.nextProgressSequence(assignment.ProgressSequence),
})
if err != nil || !progress.Accepted {
cancel()
if err == nil {
err = fmt.Errorf("job lease renewal was not accepted")
}
log.Printf("RUN phase=job.execute status=lease_renew_failed job=%s error=%s", assignment.JobID, err.Error())
return LifecycleExecutionResult{}, assignment, cancelledByPlatform, err
}
assignment = progress.Job
log.Printf("RUN phase=job.execute status=lease_renewed job=%s percent=%d", assignment.JobID, assignment.Progress.Percent)
if err := worker.journal.Store(assignment); err != nil {
cancel()
log.Printf("RUN phase=job.execute status=journal_store_failed job=%s error=%s", assignment.JobID, err.Error())
return LifecycleExecutionResult{}, assignment, cancelledByPlatform, err
}
}
}
}
func (worker *Worker) nextProgressSequence(minimum uint64) uint64 {
worker.sequenceMu.Lock()
defer worker.sequenceMu.Unlock()
worker.stateMu.Lock()
defer worker.stateMu.Unlock()
if worker.state.Sequence < minimum {
worker.state.Sequence = minimum
}
worker.state.Sequence++
return worker.state.Sequence
}
func (worker *Worker) executeAssignment(ctx context.Context, assignment protocol.RunJobAssignment) LifecycleExecutionResult {
log.Printf("RUN phase=job.dispatch status=select job=%s capability=%s inputRef=%s target=%s", assignment.JobID, assignment.Capability, safeOptional(assignment.InputRef), safeOptional(assignment.TargetKey))
if protocol.IsProtectedRequestCapability(assignment.Capability) {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=protected_request", assignment.JobID)
return worker.executeProtectedRequestJob(ctx, assignment)
}
if assignment.ExecutionInput.SourceRCON != nil {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=source_rcon", assignment.JobID)
return worker.executeSourceRCONJob(ctx, assignment)
}
if assignment.Capability == protocol.RunCapabilityDistributionBuild {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=distribution_build", assignment.JobID)
return worker.executeDistributionBuild(ctx, assignment)
}
if assignment.Capability == protocol.RunCapabilityDependenciesCheck || assignment.Capability == protocol.RunCapabilityDependenciesInstall {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=dependencies", assignment.JobID)
return worker.executeDependencyJob(ctx, assignment)
}
if assignment.Capability == protocol.RunCapabilityRunSelfUpdate {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=self_update", assignment.JobID)
return worker.executeRunSelfUpdate(ctx, assignment)
}
if assignment.Capability == protocol.RunCapabilityLogsBackfill {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=logs_backfill", assignment.JobID)
return worker.executor.ExecuteLogBackfill(ctx, assignment)
}
if assignment.Capability == protocol.RunCapabilityRemoteRunDBSQLiteProbe {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=sqlite_schema_probe", assignment.JobID)
if worker.executor.sqliteSchemaProbe == nil {
return lifecycleFailure("sqlite_probe_unavailable", "SQLite schema probe executor is unavailable")
}
if assignment.ExecutionInput.SQLiteSchemaProbe != nil {
targetKey, err := worker.materializeSQLiteProbeDataTarget(ctx, assignment)
if err != nil {
return sqliteProbeFailureForDataTarget(assignment, err)
}
assignment.TargetKey = targetKey
}
return worker.executor.sqliteSchemaProbe.Execute(ctx, assignment)
}
if assignment.Capability == protocol.RunCapabilityRemoteRunDBSQLiteQuery {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=sqlite_query", assignment.JobID)
if worker.executor.sqliteQuery == nil {
return lifecycleFailure("sqlite_query_unavailable", "SQLite query executor is unavailable")
}
targetKey, err := worker.materializeSQLiteQueryDataTarget(ctx, assignment)
if err != nil {
return sqliteQueryFailureForDataTarget(err)
}
assignment.TargetKey = targetKey
return worker.executor.sqliteQuery.Execute(ctx, assignment)
}
if isSupportedLifecycleCapability(assignment.Capability) {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=lifecycle", assignment.JobID)
return worker.executor.ExecuteContext(ctx, assignment)
}
if supportedCapability(SupportedFileCapabilities(), assignment.Capability) {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=file", assignment.JobID)
return worker.executeFileJob(ctx, assignment)
}
if isSupportedDistributionCapability(assignment.Capability) {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=distribution", assignment.JobID)
return ExecuteDistributionJob(ctx, assignment)
}
if isSupportedRemoteCapability(assignment.Capability) {
log.Printf("RUN phase=job.dispatch status=selected job=%s executor=remote", assignment.JobID)
return ExecuteRemoteAccessJob(ctx, assignment)
}
log.Printf("RUN phase=job.dispatch status=unsupported job=%s capability=%s", assignment.JobID, assignment.Capability)
return lifecycleFailure("unsupported_run_capability", "unsupported run capability")
}
func (worker *Worker) executeFileJob(ctx context.Context, assignment protocol.RunJobAssignment) LifecycleExecutionResult {
if assignment.Capability != protocol.RunCapabilityFilesRead {
return worker.executor.ExecuteContext(ctx, assignment)
}
fileExecutor := worker.executor.fileExecutor
if fileExecutor == nil {
return lifecycleExecutionFailure("file_executor_unavailable", "file executor is unavailable", false)
}
scope, targetPath, deploymentRoot, err := fileExecutor.existingReadTargetForAssignment(assignment)
if err != nil {
return lifecycleExecutionFailure("file_read_failed", err.Error(), false)
}
info, err := os.Stat(targetPath)
if err != nil {
return lifecycleExecutionFailure("file_read_failed", err.Error(), false)
}
limit := assignment.ExecutionInput.MaxReadBytes
if limit <= 0 || limit > maxExecutionContentBytes {
limit = maxExecutionContentBytes
}
if info.Size() <= int64(limit) {
return fileExecutor.read(ctx, scope, deploymentRoot, assignment)
}
if info.Size() > maxFileArtifactBytes {
return lifecycleExecutionFailure("file_read_too_large", "file exceeds artifact transfer limit", false)
}
checksum, err := checksumServerFileForArtifact(ctx, targetPath)
if err != nil {
if errors.Is(err, context.Canceled) {
return lifecycleExecutionFailure("file_cancelled", "file read cancelled", false)
}
return lifecycleExecutionFailure("file_read_failed", err.Error(), false)
}
artifactID := "artifact-" + safeWorkspaceName(assignment.JobID) + "-file-read"
if err := worker.uploadFileArtifact(ctx, assignment, artifactID, targetPath, info.Size(), checksum); err != nil {
if errors.Is(err, context.Canceled) {
return lifecycleExecutionFailure("file_cancelled", "file read cancelled", false)
}
return lifecycleExecutionFailure("file_artifact_upload_failed", "file artifact upload failed", false)
}
metadata := fileExecutor.metadata(scope, assignment.TargetKey, checksum, info.Size())
return LifecycleExecutionResult{State: lifecycleResultStateSucceeded, Progress: protocol.RunJobProgressReport{Percent: 100, Message: "file read completed"}, ResultRef: "artifact://" + artifactID, Message: "file read completed", ExecutionResult: protocol.RunJobExecutionResult{Kind: "file.read", Version: metadata.Version, Checksum: checksum, SizeBytes: info.Size(), Summary: "large file transferred as artifact"}}
}
func (worker *Worker) executeProtectedRequestJob(ctx context.Context, assignment protocol.RunJobAssignment) LifecycleExecutionResult {
if err := protocol.ValidateRunJobAssignment(assignment); err != nil {
return protectedRequestFailure(assignment.Capability, ProtectedRequestStatusFailed, "protected_request_assignment_invalid")
}
state, err := worker.registeredState()
if err != nil {
return protectedRequestFailure(assignment.Capability, ProtectedRequestStatusFailed, "protected_request_unregistered")
}
input, err := worker.client.GetProtectedRequestExecutionInput(ctx, protocol.ProtectedRequestExecutionInputRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
JobID: assignment.JobID,
LeaseToken: assignment.LeaseToken,
Attempt: assignment.Attempt,
FencingToken: assignment.FencingToken,
})
if err != nil || !protectedRequestInputMatchesAssignment(input, assignment, state.RunEndpointID) {
return protectedRequestFailure(assignment.Capability, ProtectedRequestStatusFailed, "protected_request_input_unavailable")
}
return worker.executor.ExecuteProtectedRequest(ctx, assignment, input)
}
func protectedRequestInputMatchesAssignment(input protocol.ProtectedRequestExecutionInputResponse, assignment protocol.RunJobAssignment, endpointID string) bool {
return protocol.ValidProtectedRequestExecutionInput(input) && input.JobID == assignment.JobID && input.ServerInstanceID == assignment.ServerInstanceID && input.RunEndpointID == endpointID && input.FencingToken == assignment.FencingToken && input.TargetKey == assignment.TargetKey && input.TransportKey == assignment.ExecutionInput.RemoteAdapterKey && input.Kind == protectedRequestKindForCapability(assignment.Capability)
}
func (worker *Worker) executeSourceRCONJob(ctx context.Context, assignment protocol.RunJobAssignment) LifecycleExecutionResult {
if err := protocol.ValidateRunJobAssignment(assignment); err != nil {
return lifecycleFailure("unsafe_source_rcon_plan", "Source RCON plan is invalid")
}
state, err := worker.registeredState()
if err != nil {
return lifecycleFailure("source_rcon_unregistered", "Run worker is not registered")
}
input, err := worker.client.GetSourceRCONExecutionInput(ctx, protocol.SourceRCONExecutionInputRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
JobID: assignment.JobID,
LeaseToken: assignment.LeaseToken,
Attempt: assignment.Attempt,
})
if err != nil || input.JobID != assignment.JobID || input.ServerInstanceID != assignment.ServerInstanceID || input.RunEndpointID != state.RunEndpointID {
return lifecycleFailure("source_rcon_input_unavailable", "Source RCON command input is unavailable")
}
return worker.executor.ExecuteSourceRCON(ctx, assignment, input.Command)
}
func supportedCapability(capabilities []string, target string) bool {
for _, capability := range capabilities {
if capability == target {
return true
}
}
return false
}
func (worker *Worker) ReconcileOnce(ctx context.Context) error {
state, err := worker.registeredState()
if err != nil {
return err
}
activeBefore := worker.journal.ActiveCount()
log.Printf("RUN phase=reconcile status=starting endpoint=%s activeJobs=%d", state.RunEndpointID, activeBefore)
response, err := worker.client.ReconcileJobs(ctx, protocol.RunJobReconcileRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
ActiveJobs: worker.journal.ReconcileEntries(),
})
if err != nil {
log.Printf("RUN phase=reconcile status=failed endpoint=%s error=%s", state.RunEndpointID, err.Error())
return err
}
if !response.Accepted {
log.Printf("RUN phase=reconcile status=rejected endpoint=%s", state.RunEndpointID)
return fmt.Errorf("job reconciliation was not accepted")
}
confirmed := map[string]struct{}{}
for _, job := range response.ConfirmedJobs {
if job.RunEndpointID != state.RunEndpointID {
return fmt.Errorf("reconciled job endpoint does not match registered Run endpoint")
}
if err := worker.journal.Store(job); err != nil {
return err
}
confirmed[job.JobID] = struct{}{}
}
for _, jobID := range response.DiscardJobIDs {
if err := worker.journal.Delete(jobID); err != nil {
return err
}
confirmed[jobID] = struct{}{}
}
for _, job := range worker.journal.ActiveJobs() {
if _, accounted := confirmed[job.JobID]; !accounted {
if err := worker.journal.Delete(job.JobID); err != nil {
return err
}
}
}
log.Printf("RUN phase=reconcile status=accepted endpoint=%s confirmed=%d discarded=%d activeBefore=%d activeAfter=%d", state.RunEndpointID, len(response.ConfirmedJobs), len(response.DiscardJobIDs), activeBefore, worker.journal.ActiveCount())
return nil
}
func (worker *Worker) RecoverActiveJobs(ctx context.Context) error {
activeJobs := worker.journal.ActiveJobs()
log.Printf("RUN phase=recover status=starting activeJobs=%d", len(activeJobs))
for _, assignment := range activeJobs {
if pending, ok := worker.journal.PendingResult(assignment.JobID); ok {
log.Printf("RUN phase=recover status=pending_result job=%s", assignment.JobID)
activationManifest := worker.journal.PendingActivation(assignment.JobID)
state, err := worker.registeredState()
if err != nil {
return err
}
pending.SessionToken = state.SessionToken
result, err := worker.client.CompleteJob(ctx, pending)
if err != nil {
log.Printf("RUN phase=recover status=result_submit_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
if !result.Accepted {
log.Printf("RUN phase=recover status=result_rejected job=%s", assignment.JobID)
return fmt.Errorf("recovered job result was not accepted")
}
if err := worker.journal.Delete(assignment.JobID); err != nil {
return err
}
if activationManifest != "" {
if err := worker.executor.selfUpdateActivator.Activate(activationManifest); err != nil {
return fmt.Errorf("launch recovered self-update helper: %w", err)
}
worker.restartMu.Lock()
worker.restartRequested = true
worker.restartMu.Unlock()
}
continue
}
log.Printf("RUN phase=recover status=rerun_active_job job=%s capability=%s", assignment.JobID, assignment.Capability)
if err := worker.runAssignment(ctx, assignment); err != nil {
log.Printf("RUN phase=recover status=rerun_failed job=%s error=%s", assignment.JobID, err.Error())
return err
}
}
log.Printf("RUN phase=recover status=complete activeJobs=%d", worker.journal.ActiveCount())
return nil
}
func (worker *Worker) Run(ctx context.Context) error {
log.Printf("RUN phase=run status=starting endpoint=%s", worker.cfg.RunEndpointID)
if err := worker.Register(ctx); err != nil {
return err
}
if err := worker.ReconcileOnce(ctx); err != nil {
return err
}
if err := worker.RecoverActiveJobs(ctx); err != nil {
return err
}
if err := worker.RunAutonomousLifecycleOnce(ctx); err != nil {
return err
}
if err := worker.reportAutonomousProcessObservations(ctx); err != nil {
log.Printf("RUN phase=autonomous_lifecycle.observation status=degraded error=%s", err.Error())
}
worker.reportMetricsDegraded(ctx, "startup")
if err := MarkSelfUpdateHealthy(worker.cfg.UpdateHealthFile); err != nil {
log.Printf("RUN phase=self_update_health status=mark_failed error=%s", err.Error())
return err
}
if err := worker.reportRunUpdateHealth(ctx); err != nil {
return err
}
heartbeatInterval := durationOrDefault(worker.cfg.HeartbeatInterval, 15*time.Second)
jobInterval := durationOrDefault(worker.cfg.PollInterval, 2*time.Second)
state := worker.State()
log.Printf("RUN phase=run status=ready endpoint=%s heartbeatSeconds=%d jobPollSeconds=%d", state.RunEndpointID, int(heartbeatInterval/time.Second), int(jobInterval/time.Second))
heartbeatTicker := time.NewTicker(heartbeatInterval)
defer heartbeatTicker.Stop()
workerCtx, cancelWorker := context.WithCancel(ctx)
uploaderDone := make(chan struct{})
go worker.runDurableUploaders(workerCtx, uploaderDone)
jobDone := make(chan error, 1)
controlWake := make(chan struct{}, 1)
controlDone := make(chan error, 1)
go func() { controlDone <- worker.runControlStreamLoop(workerCtx, controlWake) }()
go func() { jobDone <- worker.runJobLoop(workerCtx, jobInterval, controlWake) }()
defer func() {
cancelWorker()
<-uploaderDone
}()
for {
select {
case <-ctx.Done():
log.Printf("RUN phase=run status=context_done error=%s", ctx.Err().Error())
return ctx.Err()
case err := <-jobDone:
log.Printf("RUN phase=run status=job_loop_done error=%s", errorSummary(err))
return err
case err := <-controlDone:
if err != nil && err != context.Canceled {
log.Printf("RUN phase=run status=control_stream_done error=%s", errorSummary(err))
return err
}
return err
case <-heartbeatTicker.C:
if err := worker.HeartbeatOnce(ctx); err != nil {
log.Printf("RUN phase=heartbeat status=retry_scheduled backoffMs=%d", boundedRetryBackoff(worker.cfg.RetryBackoff).Milliseconds())
heartbeatTicker.Reset(boundedRetryBackoff(worker.cfg.RetryBackoff))
continue
}
if err := worker.reportAutonomousProcessObservations(ctx); err != nil {
log.Printf("RUN phase=autonomous_lifecycle.observation status=degraded error=%s", err.Error())
}
worker.reportMetricsDegraded(ctx, "heartbeat")
heartbeatTicker.Reset(heartbeatInterval)
}
}
}
func (worker *Worker) reportRunUpdateHealth(ctx context.Context) error {
if worker.cfg.UpdateJobID == "" && worker.cfg.UpdateOutcome == "" && worker.cfg.UpdateAttempt == 0 && worker.cfg.UpdateLeaseToken == "" {
log.Printf("RUN phase=self_update_health status=skipped")
return nil
}
if worker.cfg.UpdateJobID == "" || (worker.cfg.UpdateOutcome != "succeeded" && worker.cfg.UpdateOutcome != "rolled-back") || worker.cfg.UpdateAttempt <= 0 || worker.cfg.UpdateLeaseToken == "" {
log.Printf("RUN phase=self_update_health status=invalid_config")
return fmt.Errorf("self-update health report configuration is incomplete")
}
log.Printf("RUN phase=self_update_health status=reporting job=%s outcome=%s attempt=%d", worker.cfg.UpdateJobID, worker.cfg.UpdateOutcome, worker.cfg.UpdateAttempt)
state, err := worker.registeredState()
if err != nil {
return err
}
response, err := worker.client.ReportRunUpdateHealth(ctx, protocol.RunUpdateHealthRequest{
RunEndpointID: state.RunEndpointID,
SessionToken: state.SessionToken,
JobID: worker.cfg.UpdateJobID,
LeaseToken: worker.cfg.UpdateLeaseToken,
Attempt: worker.cfg.UpdateAttempt,
Outcome: worker.cfg.UpdateOutcome,
Version: worker.cfg.Version,
})
if err != nil {
log.Printf("RUN phase=self_update_health status=failed job=%s error=%s", worker.cfg.UpdateJobID, err.Error())
return err
}
if !response.Accepted || response.JobID != worker.cfg.UpdateJobID {
log.Printf("RUN phase=self_update_health status=rejected job=%s", worker.cfg.UpdateJobID)
return fmt.Errorf("self-update health report was not accepted")
}
log.Printf("RUN phase=self_update_health status=accepted job=%s", worker.cfg.UpdateJobID)
return nil
}
func (worker *Worker) runControlStreamLoop(ctx context.Context, wake chan<- struct{}) error {
log.Printf("RUN phase=control_stream status=starting endpoint=%s", worker.cfg.RunEndpointID)
var lastSeq uint64
for {
state, err := worker.registeredState()
if err != nil {
if waitErr := waitWorkerLoop(ctx, boundedRetryBackoff(worker.cfg.RetryBackoff)); waitErr != nil {
return waitErr
}
continue
}
err = worker.client.StreamControlEvents(ctx, protocol.RunControlStreamRequest{RunEndpointID: state.RunEndpointID, SessionToken: state.SessionToken, LastEventSeq: lastSeq}, func(event protocol.RunControlEvent) error {
if event.RunEndpointID != "" && event.RunEndpointID != state.RunEndpointID {
return fmt.Errorf("control event endpoint mismatch")
}
if event.Sequence > lastSeq {
lastSeq = event.Sequence
}
log.Printf("RUN phase=control_stream status=event endpoint=%s type=%s seq=%d", state.RunEndpointID, safeOptional(event.Type), event.Sequence)
signalControlWake(wake)
return nil
})
if ctx.Err() != nil {
log.Printf("RUN phase=control_stream status=context_done error=%s", ctx.Err().Error())
return ctx.Err()
}
if err != nil {
if sessionInvalidError(err) {
log.Printf("RUN phase=control_stream status=session_invalid endpoint=%s error=%s", state.RunEndpointID, err.Error())
if refreshErr := worker.reregisterAndReconcile(ctx, "control_stream_session_invalid", state.SessionToken); refreshErr != nil {
log.Printf("RUN phase=control_stream status=reregister_failed endpoint=%s error=%s", state.RunEndpointID, refreshErr.Error())
}
lastSeq = 0
signalControlWake(wake)
} else {
log.Printf("RUN phase=control_stream status=failed endpoint=%s error=%s", state.RunEndpointID, err.Error())
}
}
if waitErr := waitWorkerLoop(ctx, boundedRetryBackoff(worker.cfg.RetryBackoff)); waitErr != nil {
return waitErr
}
}
}
func signalControlWake(wake chan<- struct{}) {
select {
case wake <- struct{}{}:
default:
}
}
func (worker *Worker) runJobLoop(ctx context.Context, interval time.Duration, wake <-chan struct{}) error {
log.Printf("RUN phase=job_loop status=starting fallbackPollMs=%d controlStream=true", interval.Milliseconds())
for {
select {
case <-ctx.Done():
log.Printf("RUN phase=job_loop status=context_done error=%s", ctx.Err().Error())
return ctx.Err()
default:
}
if worker.journal.ActiveCount() > 0 {
log.Printf("RUN phase=job_loop status=active_jobs activeJobs=%d", worker.journal.ActiveCount())
if err := worker.ReconcileOnce(ctx); err != nil {
log.Printf("RUN phase=job_loop status=reconcile_failed error=%s", err.Error())
if err := waitWorkerLoop(ctx, boundedRetryBackoff(worker.cfg.RetryBackoff)); err != nil {
return err
}
continue
}
if err := worker.RecoverActiveJobs(ctx); err != nil {
log.Printf("RUN phase=job_loop status=recover_failed error=%s", err.Error())
if err := waitWorkerLoop(ctx, boundedRetryBackoff(worker.cfg.RetryBackoff)); err != nil {
return err
}
continue
}
}
claimStartedAt := time.Now()
handled, err := worker.claimAndRunOnce(ctx, 0)
if err != nil {
log.Printf("RUN phase=job_loop status=claim_failed error=%s", err.Error())
if err := waitWorkerLoop(ctx, boundedRetryBackoff(worker.cfg.RetryBackoff)); err != nil {
return err
}
continue
}
worker.restartMu.Lock()
restartRequested := worker.restartRequested
worker.restartMu.Unlock()
if restartRequested {
log.Printf("RUN phase=job_loop status=restart_requested")
return ErrSelfUpdateRestartRequested
}
if !handled && time.Since(claimStartedAt) < interval {
if err := waitWorkerJobWake(ctx, interval-time.Since(claimStartedAt), wake); err != nil {
return err
}
}
}
}
func waitWorkerJobWake(ctx context.Context, delay time.Duration, wake <-chan struct{}) error {
if delay <= 0 {
return nil
}
timer := time.NewTimer(delay)
defer timer.Stop()
select {
case <-ctx.Done():
return ctx.Err()
case <-wake:
return nil
case <-timer.C:
return nil
}
}
func waitWorkerLoop(ctx context.Context, delay time.Duration) error {
if delay <= 0 {
return nil
}
timer := time.NewTimer(delay)
defer timer.Stop()
select {
case <-ctx.Done():
return ctx.Err()
case <-timer.C:
return nil
}
}
type durableLogClient interface {
IngestLogBatch(context.Context, protocol.LogBatchIngestRequest) (protocol.LogBatchIngestResponse, error)
}
type durableLogProgressClient interface {
GetRunLogStreamProgress(context.Context, protocol.RunLogStreamProgressRequest) (protocol.RunLogStreamProgressResponse, error)
}
type durableArtifactClient interface {
UploadArtifactChunk(context.Context, protocol.ArtifactChunkUploadRequest) (protocol.ArtifactChunkUploadResponse, error)
}
type sessionLogBatchClient struct {
client durableLogClient
progressClient durableLogProgressClient
runEndpointID string
sessionToken string
}
type sessionLogStreamProgressClient struct {
client interface {
GetRunLogStreamProgress(context.Context, protocol.RunLogStreamProgressRequest) (protocol.RunLogStreamProgressResponse, error)
}
runEndpointID string
sessionToken string
serverID string
}
func (client sessionLogStreamProgressClient) GetRunLogStreamProgress(ctx context.Context, streamID string) (uint64, error) {
response, err := client.client.GetRunLogStreamProgress(ctx, protocol.RunLogStreamProgressRequest{
RunEndpointID: client.runEndpointID,
SessionToken: client.sessionToken,
ServerInstanceID: client.serverID,
LogStreamID: streamID,
})
if err != nil {
return 0, err
}
if !response.Accepted || response.LogStreamID != streamID {
return 0, fmt.Errorf("platform log stream progress response is invalid")
}
return response.LatestSeq, nil
}
func (client sessionLogBatchClient) IngestLogBatch(ctx context.Context, batch protocol.LogBatchIngestRequest) (protocol.LogBatchIngestResponse, error) {
batch.RunEndpointID = client.runEndpointID
batch.SessionToken = client.sessionToken
if checksum, err := checksumForLogEntries(batch.Entries); err == nil {
batch.Checksum = checksum
}
response, err := client.client.IngestLogBatch(ctx, batch)
if err != nil && logBatchNotFoundError(err) {
log.Printf("RUN phase=durable_uploaders.logs status=spool_quarantine stream=%s firstSeq=%d lastSeq=%d reason=platform_not_found error=%s", safeOptional(batch.LogStreamID), batch.FirstSeq, batch.LastSeq, err.Error())
return protocol.LogBatchIngestResponse{}, spool.PermanentLogBatchRejection("platform_not_found", err)
}
if err != nil && (logBatchSequenceGapError(err) || logBatchAcknowledgedRangeConflict(err)) {
reason := "platform_sequence_gap"
if logBatchAcknowledgedRangeConflict(err) {
reason = "platform_acknowledged_range_conflict"
}
log.Printf("RUN phase=durable_uploaders.logs status=spool_quarantine stream=%s firstSeq=%d lastSeq=%d reason=%s error=%s", safeOptional(batch.LogStreamID), batch.FirstSeq, batch.LastSeq, reason, err.Error())
return protocol.LogBatchIngestResponse{}, spool.PermanentLogBatchRejection(reason, err)
}
if err != nil && logBatchLegacySessionMetadataError(err) && strings.TrimSpace(batch.LogSessionID) == "" && !batch.SessionStartedAt.IsZero() {
log.Printf("RUN phase=durable_uploaders.logs status=spool_quarantine stream=%s firstSeq=%d lastSeq=%d reason=legacy_session_metadata error=%s", safeOptional(batch.LogStreamID), batch.FirstSeq, batch.LastSeq, err.Error())
return protocol.LogBatchIngestResponse{}, spool.PermanentLogBatchRejection("legacy_session_metadata", err)
}
if err != nil && logBatchSessionMetadataMismatchError(err) {
log.Printf("RUN phase=durable_uploaders.logs status=spool_quarantine stream=%s firstSeq=%d lastSeq=%d reason=session_metadata_mismatch error=%s", safeOptional(batch.LogStreamID), batch.FirstSeq, batch.LastSeq, err.Error())
return protocol.LogBatchIngestResponse{}, spool.PermanentLogBatchRejection("session_metadata_mismatch", err)
}
return response, err
}
func (client sessionLogBatchClient) LogStreamLatestSeq(ctx context.Context, batch protocol.LogBatchIngestRequest) (uint64, error) {
if client.progressClient == nil {
return 0, fmt.Errorf("platform log progress client is unavailable")
}
response, err := client.progressClient.GetRunLogStreamProgress(ctx, protocol.RunLogStreamProgressRequest{
RunEndpointID: client.runEndpointID,
SessionToken: client.sessionToken,
ServerInstanceID: batch.ServerInstanceID,
LogStreamID: batch.LogStreamID,
})
if err != nil {
return 0, err
}
if !response.Accepted || response.LogStreamID != batch.LogStreamID {
return 0, fmt.Errorf("platform log stream progress response is invalid")
}
log.Printf("RUN phase=durable_uploaders.logs status=sequence_recover stream=%s localFirstSeq=%d platformLatestSeq=%d", safeOptional(batch.LogStreamID), batch.FirstSeq, response.LatestSeq)
return response.LatestSeq, nil
}
func logBatchNotFoundError(err error) bool {
var httpErr interface{ HTTPStatus() int }
return errors.As(err, &httpErr) && httpErr.HTTPStatus() == http.StatusNotFound
}
func logBatchSequenceGapError(err error) bool {
var gapError interface{ LogBatchSequenceGap() bool }
return errors.As(err, &gapError) && gapError.LogBatchSequenceGap()
}
func logBatchAcknowledgedRangeConflict(err error) bool {
var conflict interface{ LogBatchAcknowledgedRangeConflict() bool }
return errors.As(err, &conflict) && conflict.LogBatchAcknowledgedRangeConflict()
}
func logBatchLegacySessionMetadataError(err error) bool {
var legacyMetadata interface{ LogBatchLegacySessionMetadata() bool }
return errors.As(err, &legacyMetadata) && legacyMetadata.LogBatchLegacySessionMetadata()
}
func logBatchSessionMetadataMismatchError(err error) bool {
var mismatch interface{ LogBatchSessionMetadataMismatch() bool }
return errors.As(err, &mismatch) && mismatch.LogBatchSessionMetadataMismatch()
}
func artifactChunkMissingOnPlatform(err error) bool {
var httpErr interface{ HTTPStatus() int }
return errors.As(err, &httpErr) && httpErr.HTTPStatus() == http.StatusNotFound
}
type sessionArtifactChunkClient struct {
client durableArtifactClient
runEndpointID string
sessionToken string
}
func (client sessionArtifactChunkClient) UploadArtifactChunk(ctx context.Context, chunk protocol.ArtifactChunkUploadRequest) (protocol.ArtifactChunkUploadResponse, error) {
chunk.RunEndpointID = client.runEndpointID
chunk.SessionToken = client.sessionToken
response, err := client.client.UploadArtifactChunk(ctx, chunk)
if err != nil {
if artifactChunkMissingOnPlatform(err) {
log.Printf("RUN phase=durable_uploaders.artifacts status=drop_stale transfer=%s artifact=%s chunk=%d reason=platform_transfer_missing", safeOptional(chunk.TransferID), safeOptional(chunk.ArtifactID), chunk.ChunkIndex)
return protocol.ArtifactChunkUploadResponse{Accepted: true, TransferID: chunk.TransferID, ArtifactID: chunk.ArtifactID, ChunkIndex: chunk.ChunkIndex}, nil
}
return response, err
}
return response, nil
}
func (worker *Worker) runDurableUploaders(ctx context.Context, done chan<- struct{}) {
defer close(done)
logSink, hasLogSink := worker.executor.logSink.(*SpoolLogSink)
artifactHook, hasArtifactHook := worker.executor.artifactHook.(*QueueArtifactHook)
logClient, hasLogClient := worker.client.(durableLogClient)
progressClient, _ := worker.client.(durableLogProgressClient)
artifactClient, hasArtifactClient := worker.client.(durableArtifactClient)
if (!hasLogSink || !hasLogClient) && (!hasArtifactHook || !hasArtifactClient) {
log.Printf("RUN phase=durable_uploaders status=disabled logs=%t artifacts=%t", hasLogSink && hasLogClient, hasArtifactHook && hasArtifactClient)
return
}
log.Printf("RUN phase=durable_uploaders status=starting logs=%t artifacts=%t", hasLogSink && hasLogClient, hasArtifactHook && hasArtifactClient)
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
log.Printf("RUN phase=durable_uploaders status=stopping")
return
case <-ticker.C:
state, err := worker.registeredState()
if err != nil {
log.Printf("RUN phase=durable_uploaders status=skipped_unregistered error=%s", err.Error())
continue
}
if hasLogSink && hasLogClient {
startedAt := time.Now()
log.Printf("RUN phase=durable_uploaders.logs status=flush_start timeoutMs=%d", durableUploaderFlushTimeout.Milliseconds())
flushCtx, cancel := context.WithTimeout(ctx, durableUploaderFlushTimeout)
client := sessionLogBatchClient{client: logClient, progressClient: progressClient, runEndpointID: state.RunEndpointID, sessionToken: state.SessionToken}
flushed, err := logSink.Spool.Flush(flushCtx, client)
if err != nil {
log.Printf("RUN phase=durable_uploaders.logs status=flush_failed flushed=%d durationMs=%d error=%s", flushed, time.Since(startedAt).Milliseconds(), err.Error())
} else if flushed > 0 {
log.Printf("RUN phase=durable_uploaders.logs status=flushed count=%d durationMs=%d", flushed, time.Since(startedAt).Milliseconds())
}
cancel()
}
if hasArtifactHook && hasArtifactClient {
startedAt := time.Now()
log.Printf("RUN phase=durable_uploaders.artifacts status=flush_start timeoutMs=%d", durableUploaderFlushTimeout.Milliseconds())
flushCtx, cancel := context.WithTimeout(ctx, durableUploaderFlushTimeout)
client := sessionArtifactChunkClient{client: artifactClient, runEndpointID: state.RunEndpointID, sessionToken: state.SessionToken}
flushed, err := artifactHook.Queue.Flush(flushCtx, client)
if err != nil {
log.Printf("RUN phase=durable_uploaders.artifacts status=flush_failed flushed=%d durationMs=%d error=%s", flushed, time.Since(startedAt).Milliseconds(), err.Error())
} else if flushed > 0 {
log.Printf("RUN phase=durable_uploaders.artifacts status=flushed count=%d durationMs=%d", flushed, time.Since(startedAt).Milliseconds())
}
cancel()
}
}
}
}
func checksumForLogEntries(entries []protocol.LogEntry) (string, error) {
stable := make([]logEntryChecksumBody, len(entries))
for i, entry := range entries {
stable[i] = logEntryChecksumBody{
Seq: entry.Seq,
Timestamp: entry.Timestamp.UTC().Format("2006-01-02T15:04:05.000000000Z07:00"),
Level: entry.Level,
Line: entry.Line,
Fields: entry.Fields,
}
}
encoded, err := json.Marshal(stable)
if err != nil {
return "", err
}
sum := sha256.Sum256(encoded)
return "sha256:" + hex.EncodeToString(sum[:]), nil
}
type logEntryChecksumBody struct {
Seq uint64 `json:"seq"`
Timestamp string `json:"timestamp"`
Level string `json:"level,omitempty"`
Line string `json:"line"`
Fields map[string]string `json:"fields,omitempty"`
}
func (worker *Worker) capacityReport() protocol.RunCapacityReport {
return worker.capacityReportFor(worker.State())
}
func (worker *Worker) capacityReportFor(state WorkerState) protocol.RunCapacityReport {
return protocol.RunCapacityReport{
MaxJobs: state.Capacity.MaxJobs,
RunningJobs: worker.journal.ActiveCount(),
QueuedJobs: 0,
Summary: "worker control active; job capacity reported separately",
}
}
func (worker *Worker) State() WorkerState {
worker.stateMu.RLock()
defer worker.stateMu.RUnlock()
state := worker.state
state.Capabilities = append([]string(nil), state.Capabilities...)
return state
}
type SpoolLogSink struct {
RunEndpointID string
SessionToken string
Spool spool.LogSpool
mu sync.Mutex
Progress func(context.Context, string, string) (uint64, error)
}
type LiveLogClient interface {
RelayLiveLogBatch(context.Context, protocol.LogBatchIngestRequest) (protocol.LogBatchIngestResponse, error)
}
// LiveLogSink is process-local and best-effort. It has no disk spool, no
// resend backlog, and does not make lifecycle execution wait for the platform.
type LiveLogSink struct {
Client LiveLogClient
RunEndpointID string
SessionToken string
mu sync.Mutex
sequences map[string]uint64
queue chan protocol.LogBatchIngestRequest
closed bool
closeOnce sync.Once
}
func NewLiveLogSink(client LiveLogClient) *LiveLogSink {
sink := &LiveLogSink{Client: client, sequences: map[string]uint64{}, queue: make(chan protocol.LogBatchIngestRequest, 512)}
go sink.dispatch()
return sink
}
func (sink *LiveLogSink) SetSession(endpointID string, sessionToken string) {
sink.mu.Lock()
sink.RunEndpointID = endpointID
sink.SessionToken = sessionToken
sink.mu.Unlock()
}
func (sink *LiveLogSink) Append(ctx context.Context, assignment protocol.RunJobAssignment, stream string, line string) error {
return sink.append(assignment, stream, line)
}
func (sink *LiveLogSink) AppendWithCursor(ctx context.Context, assignment protocol.RunJobAssignment, stream string, line string, _ ProcessLogCursor) error {
return sink.append(assignment, stream, line)
}
func (sink *LiveLogSink) append(assignment protocol.RunJobAssignment, stream string, line string) error {
if sink == nil {
return nil
}
streamKey := declaredProcessStreamKey(assignment, stream)
logStreamID := logStreamIDForAssignment(assignment, streamKey)
sink.mu.Lock()
defer sink.mu.Unlock()
if sink.closed || sink.queue == nil {
return nil
}
endpointID := sink.RunEndpointID
if endpointID == "" {
endpointID = assignment.RunEndpointID
}
if endpointID == "" || sink.SessionToken == "" || assignment.ServerInstanceID == "" {
return nil
}
sink.sequences[logStreamID]++
sequence := sink.sequences[logStreamID]
batch := protocol.LogBatchIngestRequest{
RunEndpointID: endpointID,
SessionToken: sink.SessionToken,
LogStreamID: logStreamID,
ServerInstanceID: assignment.ServerInstanceID,
StreamKey: streamKey,
Source: "process",
LogSessionID: assignment.LogSessionID,
SessionStartedAt: assignment.SessionStartedAt,
FirstSeq: sequence,
LastSeq: sequence,
Compression: "none",
Entries: []protocol.LogEntry{{Seq: sequence, Timestamp: time.Now().UTC(), Level: "info", Line: line}},
}
batch.Checksum, _ = checksumForLogEntries(batch.Entries)
select {
case sink.queue <- batch:
default:
log.Printf("RUN phase=live_log_relay status=dropped stream=%s sequence=%d reason=transient_queue_full", safeOptional(logStreamID), sequence)
}
return nil
}
func (sink *LiveLogSink) dispatch() {
for batch := range sink.queue {
if sink.Client == nil {
continue
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
_, err := sink.Client.RelayLiveLogBatch(ctx, batch)
cancel()
if err != nil {
log.Printf("RUN phase=live_log_relay status=dropped stream=%s sequence=%d error=%s", safeOptional(batch.LogStreamID), batch.FirstSeq, err.Error())
}
}
}
func (sink *LiveLogSink) Close() {
if sink == nil {
return
}
sink.closeOnce.Do(func() {
sink.mu.Lock()
sink.closed = true
if sink.queue != nil {
close(sink.queue)
}
sink.mu.Unlock()
})
}
func (sink *SpoolLogSink) Append(ctx context.Context, assignment protocol.RunJobAssignment, stream string, line string) error {
return sink.append(ctx, assignment, stream, line, nil)
}
func (sink *SpoolLogSink) AppendWithCursor(ctx context.Context, assignment protocol.RunJobAssignment, stream string, line string, cursor ProcessLogCursor) error {
return sink.append(ctx, assignment, stream, line, &spool.LogSourceCursor{StartOffset: cursor.StartOffset, EndOffset: cursor.EndOffset})
}
func (sink *SpoolLogSink) append(ctx context.Context, assignment protocol.RunJobAssignment, stream string, line string, cursor *spool.LogSourceCursor) error {
sink.mu.Lock()
defer sink.mu.Unlock()
streamKey := declaredProcessStreamKey(assignment, stream)
logStreamID := logStreamIDForAssignment(assignment, streamKey)
entry := protocol.LogEntry{Timestamp: time.Now().UTC(), Level: "info", Line: line}
source := "process"
if strings.HasPrefix(stream, "management-program.") {
source = "management-program"
} else if assignment.Capability == protocol.RunCapabilityLogsBackfill {
source = "file"
}
var recoverProgress func(context.Context, string) (uint64, error)
runStreamPrefix := "run." + assignment.RunEndpointID + "." + assignment.ServerInstanceID + "."
if sink.Progress != nil && strings.HasPrefix(logStreamID, runStreamPrefix) {
recoverProgress = func(ctx context.Context, streamID string) (uint64, error) {
return sink.Progress(ctx, assignment.ServerInstanceID, streamID)
}
}
if strings.TrimSpace(assignment.LogSessionID) != "" {
// A generation-scoped stream is globally fresh, so its first durable
// append must not depend on platform availability or the worker ctx.
recoverProgress = nil
}
_, _, err := sink.Spool.EnqueueNextAggregated(ctx, protocol.LogBatchIngestRequest{
RunEndpointID: sink.RunEndpointID,
SessionToken: sink.SessionToken,
LogStreamID: logStreamID,
ServerInstanceID: assignment.ServerInstanceID,
StreamKey: streamKey,
Source: source,
LogSessionID: assignment.LogSessionID,
SessionStartedAt: assignment.SessionStartedAt,
Entries: []protocol.LogEntry{entry},
}, cursor, recoverProgress, checksumForLogEntries)
return err
}
func logStreamIDForAssignment(assignment protocol.RunJobAssignment, streamKey string) string {
if strings.TrimSpace(assignment.LogSessionID) != "" {
return fmt.Sprintf("run.%s.%s.%s.%s", assignment.RunEndpointID, assignment.ServerInstanceID, assignment.LogSessionID, streamKey)
}
if autonomousLifecycleLogAssignment(assignment) {
return fmt.Sprintf("run.%s.%s.%s", assignment.RunEndpointID, assignment.ServerInstanceID, streamKey)
}
return fmt.Sprintf("job.%s.%s", assignment.JobID, streamKey)
}
func autonomousLifecycleLogAssignment(assignment protocol.RunJobAssignment) bool {
return strings.HasPrefix(assignment.IdempotencyKey, "autonomous:") || strings.HasPrefix(assignment.LeaseToken, "local-autonomous-") || strings.HasPrefix(assignment.JobID, "autonomous-")
}
func declaredProcessStreamKey(assignment protocol.RunJobAssignment, stream string) string {
kind := ""
if stream == "stdout" {
kind = "process.stdout"
} else if stream == "stderr" {
kind = "process.stderr"
}
if kind != "" {
for _, source := range assignment.ExecutionInput.LogSources {
if source.Kind == kind && strings.TrimSpace(source.StreamKey) != "" {
return source.StreamKey
}
}
}
return stream
}
type QueueArtifactHook struct {
RunEndpointID string
SessionToken string
Queue spool.ArtifactQueue
}
func (hook QueueArtifactHook) QueueLifecycleResult(_ context.Context, assignment protocol.RunJobAssignment, result ProcessResult) (string, error) {
ref := fmt.Sprintf("artifact://jobs/%s/lifecycle-result", assignment.JobID)
payload := []byte(result.Stdout + result.Stderr)
if len(payload) == 0 {
payload = []byte("lifecycle result metadata")
}
artifactID := "artifact-" + assignment.JobID + "-lifecycle"
if err := hook.Queue.Enqueue(protocol.ArtifactChunkUploadRequest{
RunEndpointID: hook.RunEndpointID,
SessionToken: hook.SessionToken,
TransferID: "transfer-" + assignment.JobID,
ArtifactID: artifactID,
ChunkIndex: 0,
Offset: 0,
SizeBytes: len(payload),
Checksum: checksumForText(string(payload)),
Payload: payload,
}); err != nil {
return "", err
}
return ref, nil
}
func capabilityFingerprint(capabilities []string) string {
return checksumForText(strings.Join(capabilities, ","))
}
func durationOrDefault(value time.Duration, fallback time.Duration) time.Duration {
if value <= 0 {
return fallback
}
return value
}
func boundedRetryBackoff(value time.Duration) time.Duration {
value = durationOrDefault(value, time.Second)
if value > 30*time.Second {
return 30 * time.Second
}
return value
}