package web import ( "context" "errors" "fmt" "regexp" "sort" "strings" "sync" "time" "gitgud.io/mike/mpv-manager/pkg/constants" "gitgud.io/mike/mpv-manager/pkg/jobhistory" "gitgud.io/mike/mpv-manager/pkg/log" ) // Job represents an installation, uninstallation, or update job type Job struct { ID string // Unique job ID (timestamp-based) Type string // "install", "uninstall", "update" MethodID string // Installation method ID (e.g., "mpv-flatpak") AppName string // Display name (e.g., "MPV") Status string // pending, running, cancelling, handed_off, complete, error, partial, cancelled Phase string // Structured operation phase (used by manager updates) Progress int // 0-100 percentage Message string // Current action message Output []string // Full log output (all lines with timestamps) Error string // Short error message if failed ErrorDetails string // Detailed error reason StartedAt time.Time CompletedAt *time.Time // Context for cancellation support ctx context.Context cancel context.CancelFunc // resources remain leased until the worker performs a terminal transition. // commitStarted prevents cancellation from claiming success once an // irreversible operation has begun. resources []string commitStarted bool outputBytes int outputTrimmed bool } // Context returns the job's context func (j *Job) Context() context.Context { return j.ctx } // Cancel cancels the job's context func (j *Job) Cancel() { if j.cancel != nil { j.cancel() } } // IsCancelled returns true if the job's context has been cancelled func (j *Job) IsCancelled() bool { if j.ctx == nil { return false } select { case <-j.ctx.Done(): return true default: return false } } // JobEvent represents an event broadcast for SSE streaming type JobEvent struct { JobID string // Job identifier Type string // "created", "progress", "output", "status", "complete" Data interface{} // Event-specific data } // JobEventData contains the data for progress/output events type JobProgressData struct { Progress int `json:"progress"` Message string `json:"message"` Phase string `json:"phase,omitempty"` } type JobOutputData struct { Line string `json:"line"` } type JobStatusData struct { Status string `json:"status"` Phase string `json:"phase,omitempty"` Error string `json:"error,omitempty"` ErrorDetails string `json:"errorDetails,omitempty"` AppName string `json:"appName,omitempty"` Type string `json:"type,omitempty"` MethodID string `json:"methodId,omitempty"` } type JobCompleteData struct { Success bool `json:"success"` Progress int `json:"progress"` Message string `json:"message"` AppName string `json:"appName"` Type string `json:"type"` MethodID string `json:"methodId"` Phase string `json:"phase,omitempty"` CompletedAt time.Time `json:"completedAt"` } // JobManager manages all jobs and broadcasts events for SSE type JobManager struct { jobs map[string]*Job // All active jobs (pending/running) recentJobs []*Job // Last 10 completed jobs for history resources map[string]string // normalized resource key -> owning job/request ID mu sync.RWMutex closing bool workers sync.WaitGroup // store persists terminal jobs across restarts (nil = in-memory only). // Appends are serialized by the store after jm.mu is released. store *jobhistory.Store // Fan-out pattern: each SSE client gets its own channel clients map[chan JobEvent]struct{} // Registry of connected SSE clients clientsClosed bool clientsMu sync.RWMutex // Protects client lifecycle and fan-out } // maxRecentJobs is the maximum number of recent jobs to keep in history const maxRecentJobs = 10 // maxJobOutputLines is the maximum number of output lines kept per job. // Oldest lines are dropped beyond this cap to bound memory usage. const maxJobOutputLines = 2000 // maxSSEClients bounds the per-process memory and fan-out work consumed by // persistent browser streams. Each accepted client owns a buffered channel. const maxSSEClients = 64 // progressRegex matches percentages like "(45%)" or "45%" var progressRegex = regexp.MustCompile(`(?:\()?(\d{1,3})%(?:\))?`) var ( ErrJobNotFound = errors.New("job not found") ErrJobTerminal = errors.New("job already terminal") ErrJobCancelling = errors.New("job cancellation already requested") ErrJobCommitStarted = errors.New("job has entered its commit boundary") ) // NewJobManager creates a new job manager with fan-out broadcasting func NewJobManager() *JobManager { return &JobManager{ jobs: make(map[string]*Job), recentJobs: make([]*Job, 0, maxRecentJobs), resources: make(map[string]string), clients: make(map[chan JobEvent]struct{}), } } // RegisterClient registers a new SSE client and returns its dedicated channel. // It returns nil when the bounded client pool is full. The client must call // UnregisterClient when done to prevent memory leaks. func (jm *JobManager) RegisterClient() chan JobEvent { clientChan := make(chan JobEvent, 100) // Buffered to handle bursts jm.clientsMu.Lock() if jm.clientsClosed { close(clientChan) jm.clientsMu.Unlock() return clientChan } if len(jm.clients) >= maxSSEClients { jm.clientsMu.Unlock() return nil } jm.clients[clientChan] = struct{}{} jm.clientsMu.Unlock() return clientChan } // CloseClients closes the broadcaster and every registered SSE channel. It is // idempotent and permanently rejects new clients, allowing graceful server // shutdown to drain persistent stream handlers before its deadline. func (jm *JobManager) CloseClients() { jm.clientsMu.Lock() defer jm.clientsMu.Unlock() if jm.clientsClosed { return } jm.clientsClosed = true for clientChan := range jm.clients { delete(jm.clients, clientChan) close(clientChan) } } // UnregisterClient removes a client and closes its channel // Safe to call multiple times func (jm *JobManager) UnregisterClient(c chan JobEvent) { if c == nil { return } jm.clientsMu.Lock() defer jm.clientsMu.Unlock() if _, exists := jm.clients[c]; exists { delete(jm.clients, c) close(c) } } // broadcastEvent sends an event to ALL registered SSE clients // Uses non-blocking send to prevent slow clients from blocking func (jm *JobManager) broadcastEvent(event JobEvent) { jm.clientsMu.RLock() defer jm.clientsMu.RUnlock() for clientChan := range jm.clients { sendToClient(clientChan, event) } } // isTerminalEventType reports whether an event type carries terminal/status // state that must never be dropped — dropping one leaves the UI stuck forever. func isTerminalEventType(eventType string) bool { return eventType == "status" || eventType == "complete" } // sendToClient delivers an event to a single client channel. // Progress/output events are droppable under backpressure; terminal/status // events are not — when the buffer is full, the oldest buffered events are // drained to make room for them. func sendToClient(clientChan chan JobEvent, event JobEvent) { select { case clientChan <- event: // Event sent successfully default: // Client buffer full — slow clients must not block everyone else if !isTerminalEventType(event.Type) { return } for { select { case clientChan <- event: return default: } select { case <-clientChan: // Dropped the oldest buffered event to make room default: // Buffer drained but send still fails (concurrent senders); // give up rather than spin forever return } } } } // newJob builds a new pending job (does not register it with the manager) func newJob(jobType, methodID, appName string) *Job { jobID := newJobID() // Create a context with cancel for this job ctx, cancel := context.WithCancel(context.Background()) return &Job{ ID: jobID, Type: jobType, MethodID: methodID, AppName: appName, Status: "pending", Progress: 0, Message: "Job created", Output: make([]string, 0), StartedAt: time.Now(), ctx: ctx, cancel: cancel, } } func newJobID() string { return jobhistory.NewID() } // broadcastJobCreated announces a newly created job to all SSE clients func (jm *JobManager) broadcastJobCreated(job *Job) { jm.broadcastEvent(JobEvent{ JobID: job.ID, Type: "created", Data: map[string]interface{}{ "type": job.Type, "methodId": job.MethodID, "appName": job.AppName, }, }) } // CreateJob creates a new pending job and returns it func (jm *JobManager) CreateJob(jobType, methodID, appName string) *Job { job := newJob(jobType, methodID, appName) jm.mu.Lock() jm.jobs[job.ID] = job jm.mu.Unlock() // Broadcast job creation event jm.broadcastJobCreated(job) return job } // CreateJobIfNoneActive atomically checks for an active (pending/running) job // for methodID and creates a new job only if none exists. Holding a single // lock across check-and-create closes the TOCTOU race between HasActiveJob // and CreateJob that allowed duplicate jobs for the same method. // Returns (nil, false) when an active job already exists for the method. func (jm *JobManager) CreateJobIfNoneActive(jobType, methodID, appName string) (*Job, bool) { return jm.CreateJobWithResources(jobType, methodID, appName, "method:"+methodID) } // CreateJobWithResources atomically leases every supplied resource and creates // the job. Resource keys are sorted/deduplicated, so multi-resource acquisition // has no lock-order ambiguity. A cancelling job retains its leases until its // worker acknowledges a terminal result. func (jm *JobManager) CreateJobWithResources(jobType, methodID, appName string, resources ...string) (*Job, bool) { job := newJob(jobType, methodID, appName) job.resources = normalizeResourceKeys(resources) jm.mu.Lock() if jm.closing || !jm.acquireResourcesLocked(job.ID, job.resources) { jm.mu.Unlock() job.cancel() return nil, false } jm.jobs[job.ID] = job jm.mu.Unlock() jm.broadcastJobCreated(job) return job, true } // TryAcquireResources leases resources for a synchronous request. The release // callback is idempotent and releases only keys still owned by this request. func (jm *JobManager) TryAcquireResources(owner string, resources ...string) (func(), bool) { keys := normalizeResourceKeys(resources) jm.mu.Lock() if jm.closing || !jm.acquireResourcesLocked(owner, keys) { jm.mu.Unlock() return nil, false } jm.workers.Add(1) // Synchronous mutating handlers also belong to shutdown. jm.mu.Unlock() var once sync.Once release := func() { once.Do(func() { jm.mu.Lock() jm.releaseResourcesLocked(owner, keys) jm.mu.Unlock() jm.workers.Done() }) } return release, true } func normalizeResourceKeys(resources []string) []string { seen := make(map[string]struct{}, len(resources)) keys := make([]string, 0, len(resources)) for _, resource := range resources { resource = strings.TrimSpace(resource) if resource == "" { continue } if _, exists := seen[resource]; exists { continue } seen[resource] = struct{}{} keys = append(keys, resource) } sort.Strings(keys) return keys } func (jm *JobManager) acquireResourcesLocked(owner string, resources []string) bool { for _, resource := range resources { if existing := jm.resources[resource]; existing != "" && existing != owner { return false } } for _, resource := range resources { jm.resources[resource] = owner } return true } func (jm *JobManager) releaseResourcesLocked(owner string, resources []string) { for _, resource := range resources { if jm.resources[resource] == owner { delete(jm.resources, resource) } } } // GetJob retrieves a job by ID func (jm *JobManager) GetJob(jobID string) (*Job, bool) { jm.mu.RLock() defer jm.mu.RUnlock() job, exists := jm.jobs[jobID] return job, exists } // snapshotJob returns a value copy of the job with a copied Output slice. // Must be called with jm.mu held. func snapshotJob(job *Job) Job { cp := *job cp.Output = append([]string(nil), job.Output...) if job.outputTrimmed { cp.Output = append([]string{constants.OutputTailMarker}, cp.Output...) cp.outputTrimmed = false } return cp } // JobSnapshot returns a value copy of the job with the given ID, including a // copied Output slice. The snapshot is safe to read without holding jm.mu. func (jm *JobManager) JobSnapshot(jobID string) (Job, bool) { jm.mu.RLock() defer jm.mu.RUnlock() job, exists := jm.jobs[jobID] if !exists { return Job{}, false } return snapshotJob(job), true } // UpdateProgress updates the progress and message for a job // Also parses progress from the message if it contains a percentage func (jm *JobManager) UpdateProgress(jobID string, progress int, message string) { jm.updateProgress(jobID, "", progress, message) } // UpdatePhase updates a job's structured lifecycle phase together with its // display progress and message. func (jm *JobManager) UpdatePhase(jobID, phase string, progress int, message string) { jm.updateProgress(jobID, phase, progress, message) } func (jm *JobManager) updateProgress(jobID, phase string, progress int, message string) { jm.mu.Lock() defer jm.mu.Unlock() job, exists := jm.jobs[jobID] if !exists || job.IsComplete() { return } // Parse progress from message if it contains percentage parsedProgress := jm.parseProgress(message) if parsedProgress > progress { progress = parsedProgress } // A completed download is only one part of installation. Reserve 100 for // CompleteJob, after extraction, configuration and tracking have succeeded. if progress < 0 { progress = 0 } else if progress > 99 { progress = 99 } job.Progress = progress job.Message = message if phase != "" { job.Phase = phase } // Update status to running if it was pending if job.Status == "pending" { job.Status = "running" } // Broadcast progress event jm.broadcastEvent(JobEvent{ JobID: jobID, Type: "progress", Data: JobProgressData{ Progress: progress, Message: message, Phase: job.Phase, }, }) } // AddOutput adds a timestamped output line to the job func (jm *JobManager) AddOutput(jobID string, line string) { jm.mu.Lock() defer jm.mu.Unlock() job, exists := jm.jobs[jobID] if !exists { return } // Format output with timestamp if len(line) > constants.MaxOutputLineBytes { const suffix = " [truncated]" line = line[:constants.MaxOutputLineBytes-len(suffix)] + suffix } timestamp := time.Now().Format("15:04:05") timestampedLine := fmt.Sprintf("[%s] %s", timestamp, line) job.Output = append(job.Output, timestampedLine) job.outputBytes += len(timestampedLine) // Cap stored output to bound memory usage (drop oldest lines) for len(job.Output) > maxJobOutputLines-1 || job.outputBytes > constants.MaxJobOutputBytes-len(constants.OutputTailMarker) { job.outputBytes -= len(job.Output[0]) job.Output[0] = "" job.Output = job.Output[1:] job.outputTrimmed = true } // Update status to running if it was pending if job.Status == "pending" { job.Status = "running" } // Try to parse progress from the output line if progress := min(jm.parseProgress(line), 99); !job.IsComplete() && progress > job.Progress { job.Progress = progress } // Broadcast output event jm.broadcastEvent(JobEvent{ JobID: jobID, Type: "output", Data: JobOutputData{ Line: timestampedLine, }, }) } // SetError marks a job as errored with the given error and details func (jm *JobManager) SetError(jobID string, err error, details string) { jm.mu.Lock() job, exists := jm.jobs[jobID] if !exists { jm.mu.Unlock() return } if job.IsComplete() { jm.mu.Unlock() return } now := time.Now() status := "error" errorText := err.Error() if job.Status == "cancelling" || (job.ctx != nil && job.ctx.Err() != nil && !job.commitStarted) { status = "cancelled" errorText = "Job cancelled by user" if details == "" { details = "The worker acknowledged cancellation before entering its commit boundary" } } job.Status = status job.Error = errorText job.ErrorDetails = details job.CompletedAt = &now // Broadcast status event jm.broadcastEvent(JobEvent{ JobID: jobID, Type: "status", Data: JobStatusData{ Status: status, Phase: job.Phase, Error: job.Error, ErrorDetails: job.ErrorDetails, AppName: job.AppName, Type: job.Type, MethodID: job.MethodID, }, }) // Archive the job archived := jm.archiveJobLocked(job) jm.mu.Unlock() jm.persistArchivedJob(archived) } // SetPartialSuccess records that the physical operation committed but its // subsequent setup or authoritative tracking metadata could not be completed. This is terminal // and intentionally distinct from both success and rollback-safe failure. func (jm *JobManager) SetPartialSuccess(jobID string, err error, details string) { jm.mu.Lock() job, exists := jm.jobs[jobID] if !exists || job.IsComplete() { jm.mu.Unlock() return } now := time.Now() job.Status = "partial" job.Error = err.Error() job.ErrorDetails = details job.Message = "Operation needs reconciliation" job.CompletedAt = &now jm.broadcastEvent(JobEvent{ JobID: jobID, Type: "status", Data: JobStatusData{ Status: "partial", Phase: job.Phase, Error: job.Error, ErrorDetails: job.ErrorDetails, AppName: job.AppName, Type: job.Type, MethodID: job.MethodID, }, }) archived := jm.archiveJobLocked(job) jm.mu.Unlock() jm.persistArchivedJob(archived) } // SetHandedOff records that a detached helper owns the remaining update. This // is terminal for the current server process, but deliberately is not final // success; the helper outcome replaces this history record on the next start. func (jm *JobManager) SetHandedOff(jobID, transactionID string) { if transactionID != "" { jm.AddOutput(jobID, "Detached helper transaction: "+transactionID) } jm.mu.Lock() job, exists := jm.jobs[jobID] if !exists || job.IsComplete() { jm.mu.Unlock() return } now := time.Now() job.Status = "handed_off" job.Progress = 99 job.Message = "Verified update handed to the detached helper; final result will be recorded on next start" job.CompletedAt = &now jm.broadcastEvent(JobEvent{ JobID: jobID, Type: "status", Data: JobStatusData{ Status: job.Status, Phase: job.Phase, AppName: job.AppName, Type: job.Type, MethodID: job.MethodID, }, }) archived := jm.archiveJobLocked(job) jm.mu.Unlock() jm.persistArchivedJob(archived) } // CompleteJob marks a job as complete func (jm *JobManager) CompleteJob(jobID string) { jm.mu.Lock() job, exists := jm.jobs[jobID] if !exists { jm.mu.Unlock() return } if job.IsComplete() { jm.mu.Unlock() return } now := time.Now() if !job.commitStarted && (job.Status == "cancelling" || (job.ctx != nil && job.ctx.Err() != nil)) { job.Status = "cancelled" job.Error = "Job cancelled by user" job.ErrorDetails = "The worker acknowledged cancellation before entering its commit boundary" job.CompletedAt = &now jm.broadcastEvent(JobEvent{JobID: jobID, Type: "status", Data: JobStatusData{ Status: job.Status, Error: job.Error, ErrorDetails: job.ErrorDetails, AppName: job.AppName, Type: job.Type, MethodID: job.MethodID, Phase: job.Phase, }}) archived := jm.archiveJobLocked(job) jm.mu.Unlock() jm.persistArchivedJob(archived) return } job.Status = "complete" job.Progress = 100 job.CompletedAt = &now if job.ctx != nil && job.ctx.Err() != nil { job.Message = "Completed; cancellation arrived after the commit boundary" } // Broadcast complete event jm.broadcastEvent(JobEvent{ JobID: jobID, Type: "complete", Data: JobCompleteData{ Success: true, Progress: 100, Message: job.Message, AppName: job.AppName, Type: job.Type, MethodID: job.MethodID, Phase: job.Phase, CompletedAt: now, }, }) // Archive the job archived := jm.archiveJobLocked(job) jm.mu.Unlock() jm.persistArchivedJob(archived) } // BeginCommit atomically declares the point after which cancellation can no // longer truthfully promise that side effects will not commit. func (jm *JobManager) BeginCommit(jobID string) error { jm.mu.Lock() defer jm.mu.Unlock() job, exists := jm.jobs[jobID] if !exists { return fmt.Errorf("%w: %s", ErrJobNotFound, jobID) } if job.IsComplete() { return fmt.Errorf("%w: %s", ErrJobTerminal, jobID) } if job.Status == "cancelling" || (job.ctx != nil && job.ctx.Err() != nil) { return context.Canceled } job.commitStarted = true return nil } // CancelJob requests cancellation but deliberately leaves the job active and // leased. The worker owns the sole terminal transition after it stops. func (jm *JobManager) CancelJob(jobID string) error { jm.mu.Lock() job, exists := jm.jobs[jobID] if !exists { jm.mu.Unlock() return fmt.Errorf("%w: %s", ErrJobNotFound, jobID) } // Check if job can be cancelled if job.IsComplete() { jm.mu.Unlock() return fmt.Errorf("%w: %s", ErrJobTerminal, jobID) } if job.Status == "cancelling" { jm.mu.Unlock() return fmt.Errorf("%w: %s", ErrJobCancelling, jobID) } if job.commitStarted { jm.mu.Unlock() return fmt.Errorf("%w: %s", ErrJobCommitStarted, jobID) } if job.Type == "manager-update" { jm.mu.Unlock() return fmt.Errorf("manager updates cannot be cancelled after staging begins") } // Cancel the context if job.cancel != nil { job.cancel() } job.Status = "cancelling" job.Message = "Cancellation requested; waiting for the worker to stop" // This is non-terminal. Clients keep the job visible until the worker // acknowledges cancellation with SetError/CompleteJob. jm.broadcastEvent(JobEvent{ JobID: jobID, Type: "status", Data: JobStatusData{ Status: "cancelling", Phase: job.Phase, Error: "Job cancelled by user", ErrorDetails: "The job was cancelled by user request", AppName: job.AppName, Type: job.Type, MethodID: job.MethodID, }, }) jm.mu.Unlock() return nil } // jobSummaries takes one coherent active/recent snapshot without copying output. // List and SSE consumers only expose metadata; details retain their output tail. func (jm *JobManager) jobSummaries() (active, recent []Job) { jm.mu.RLock() defer jm.mu.RUnlock() active = make([]Job, 0, len(jm.jobs)) recent = make([]Job, 0, len(jm.recentJobs)) summary := func(job *Job) Job { cp := *job; cp.Output = nil; return cp } for _, job := range jm.jobs { if job.Status == "pending" || job.Status == "running" || job.Status == "cancelling" { active = append(active, summary(job)) } } for _, job := range jm.recentJobs { recent = append(recent, summary(job)) } return active, recent } // GetActiveJobs returns snapshots of all jobs that are pending or running. // Snapshots are value copies (with copied Output slices) that are safe to // read without holding jm.mu — unlike the shared *Job pointers previously // returned, which raced with concurrent job updates. func (jm *JobManager) GetActiveJobs() []Job { jm.mu.RLock() defer jm.mu.RUnlock() jobs := make([]Job, 0, len(jm.jobs)) for _, job := range jm.jobs { if job.Status == "pending" || job.Status == "running" || job.Status == "cancelling" { jobs = append(jobs, snapshotJob(job)) } } return jobs } // GetRecentJobs returns snapshots of the last completed jobs (max 10) func (jm *JobManager) GetRecentJobs() []Job { jm.mu.RLock() defer jm.mu.RUnlock() // Return copies to avoid race conditions result := make([]Job, len(jm.recentJobs)) for i, job := range jm.recentJobs { result[i] = snapshotJob(job) } return result } // parseProgress extracts a percentage from a string // Matches patterns like "(45%)" or "45%" func (jm *JobManager) parseProgress(s string) int { matches := progressRegex.FindStringSubmatch(s) if len(matches) > 1 { var progress int if _, err := fmt.Sscanf(matches[1], "%d", &progress); err == nil { if progress >= 0 && progress <= 100 { return progress } } } return 0 } // archiveJob moves a completed/errored job to recent jobs and removes from active // Must be called with lock already held // archiveJobLocked removes a terminal job from the active map and returns a // detached snapshot for persistence. The caller must hold jm.mu. func (jm *JobManager) archiveJobLocked(job *Job) *Job { jm.releaseResourcesLocked(job.ID, job.resources) // Remove from active jobs delete(jm.jobs, job.ID) // Add to recent jobs (at the beginning) jm.recentJobs = append([]*Job{job}, jm.recentJobs...) // Trim to max size if len(jm.recentJobs) > maxRecentJobs { jm.recentJobs = jm.recentJobs[:maxRecentJobs] } archived := snapshotJob(job) return &archived } // persistArchivedJob performs filesystem I/O without holding the manager's // main state lock, so a slow home directory cannot stall progress or queries. func (jm *JobManager) persistArchivedJob(job *Job) { if jm.store == nil || job == nil { return } if err := jm.store.Append(jobToRecord(job)); err != nil { log.Warn("Failed to persist job history: " + err.Error()) } } // HasActiveJob checks if there's an active job for a specific method func (jm *JobManager) HasActiveJob(methodID string) bool { jm.mu.RLock() defer jm.mu.RUnlock() for _, job := range jm.jobs { if job.MethodID == methodID && (job.Status == "pending" || job.Status == "running" || job.Status == "cancelling") { return true } } return false } // Duration returns the duration of the job func (j *Job) Duration() time.Duration { if j.CompletedAt != nil { return j.CompletedAt.Sub(j.StartedAt) } return time.Since(j.StartedAt) } // IsComplete returns true if the job reached a terminal state // (complete, errored, or cancelled) func (j *Job) IsComplete() bool { return j.Status == "handed_off" || j.Status == "complete" || j.Status == "error" || j.Status == "partial" || j.Status == "cancelled" } // IsError returns true if the job failed with an error func (j *Job) IsError() bool { return j.Status == "error" || j.Status == "partial" }