Files
run/runtime/job_journal.go
2026-08-26 09:56:43 +08:00

294 lines
9.3 KiB
Go

package runtime
import (
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"sync"
"browser.local/run/protocol"
)
const jobJournalVersion = 1
type JobJournal struct {
mu sync.Mutex
path string
active map[string]protocol.RunJobAssignment
pendingResults map[string]protocol.RunJobResultRequest
pendingActivations map[string]string
}
type jobJournalFile struct {
Version int `json:"version"`
Active []protocol.RunJobAssignment `json:"active"`
PendingResults []protocol.RunJobResultRequest `json:"pendingResults,omitempty"`
PendingActivations map[string]string `json:"pendingActivations,omitempty"`
}
func NewJobJournal() *JobJournal {
return &JobJournal{active: map[string]protocol.RunJobAssignment{}, pendingResults: map[string]protocol.RunJobResultRequest{}, pendingActivations: map[string]string{}}
}
func NewPersistentJobJournal(workspaceRoot string) (*JobJournal, error) {
if workspaceRoot == "" {
workspaceRoot = filepath.Join(".", ".run-workspace")
}
dir := filepath.Join(workspaceRoot, "state")
if err := os.MkdirAll(dir, 0o700); err != nil {
return nil, fmt.Errorf("create job journal directory: %w", err)
}
if err := os.Chmod(dir, 0o700); err != nil {
return nil, fmt.Errorf("secure job journal directory: %w", err)
}
journal := &JobJournal{path: filepath.Join(dir, "jobs.json"), active: map[string]protocol.RunJobAssignment{}, pendingResults: map[string]protocol.RunJobResultRequest{}, pendingActivations: map[string]string{}}
if err := journal.load(); err != nil {
return nil, err
}
return journal, nil
}
func (journal *JobJournal) Store(job protocol.RunJobAssignment) error {
journal.mu.Lock()
defer journal.mu.Unlock()
if err := validateJournalAssignment(job); err != nil {
return err
}
previous, existed := journal.active[job.JobID]
journal.active[job.JobID] = job
if err := journal.persistLocked(); err != nil {
if existed {
journal.active[job.JobID] = previous
} else {
delete(journal.active, job.JobID)
}
return err
}
return nil
}
func (journal *JobJournal) Delete(jobID string) error {
journal.mu.Lock()
defer journal.mu.Unlock()
previous, existed := journal.active[jobID]
previousResult, hadResult := journal.pendingResults[jobID]
previousActivation, hadActivation := journal.pendingActivations[jobID]
delete(journal.active, jobID)
delete(journal.pendingResults, jobID)
delete(journal.pendingActivations, jobID)
if err := journal.persistLocked(); err != nil {
if existed {
journal.active[jobID] = previous
}
if hadResult {
journal.pendingResults[jobID] = previousResult
}
if hadActivation {
journal.pendingActivations[jobID] = previousActivation
}
return err
}
return nil
}
func (journal *JobJournal) StorePendingResult(result protocol.RunJobResultRequest, activationManifest string) error {
journal.mu.Lock()
defer journal.mu.Unlock()
assignment, exists := journal.active[result.JobID]
if !exists || assignment.Attempt != result.Attempt || assignment.LeaseToken != result.LeaseToken || assignment.RunEndpointID != result.RunEndpointID {
return fmt.Errorf("pending result does not match active journal attempt")
}
if result.State != "succeeded" && result.State != "failed" && result.State != "cancelled" {
return fmt.Errorf("pending result state is not terminal")
}
result.SessionToken = ""
previous, hadPrevious := journal.pendingResults[result.JobID]
previousActivation, hadActivation := journal.pendingActivations[result.JobID]
journal.pendingResults[result.JobID] = result
if activationManifest != "" {
journal.pendingActivations[result.JobID] = activationManifest
} else {
delete(journal.pendingActivations, result.JobID)
}
if err := journal.persistLocked(); err != nil {
if hadPrevious {
journal.pendingResults[result.JobID] = previous
} else {
delete(journal.pendingResults, result.JobID)
}
if hadActivation {
journal.pendingActivations[result.JobID] = previousActivation
} else {
delete(journal.pendingActivations, result.JobID)
}
return err
}
return nil
}
func (journal *JobJournal) PendingActivation(jobID string) string {
journal.mu.Lock()
defer journal.mu.Unlock()
return journal.pendingActivations[jobID]
}
func (journal *JobJournal) PendingResult(jobID string) (protocol.RunJobResultRequest, bool) {
journal.mu.Lock()
defer journal.mu.Unlock()
result, exists := journal.pendingResults[jobID]
return result, exists
}
func (journal *JobJournal) MarkActive(job protocol.RunJobAssignment) {
_ = journal.Store(job)
}
func (journal *JobJournal) MarkTerminal(jobID string) {
_ = journal.Delete(jobID)
}
func (journal *JobJournal) ActiveJobs() []protocol.RunJobAssignment {
journal.mu.Lock()
defer journal.mu.Unlock()
jobs := make([]protocol.RunJobAssignment, 0, len(journal.active))
for _, job := range journal.active {
jobs = append(jobs, job)
}
sort.Slice(jobs, func(i, j int) bool { return jobs[i].JobID < jobs[j].JobID })
return jobs
}
func (journal *JobJournal) ReconcileEntries() []protocol.RunJobReconcileEntry {
jobs := journal.ActiveJobs()
entries := make([]protocol.RunJobReconcileEntry, len(jobs))
for i, job := range jobs {
entries[i] = protocol.RunJobReconcileEntry{JobID: job.JobID, LeaseToken: job.LeaseToken, Attempt: job.Attempt}
}
return entries
}
func (journal *JobJournal) ActiveJobIDs() []string {
jobs := journal.ActiveJobs()
ids := make([]string, len(jobs))
for i, job := range jobs {
ids[i] = job.JobID
}
return ids
}
func (journal *JobJournal) ActiveCount() int {
journal.mu.Lock()
defer journal.mu.Unlock()
return len(journal.active)
}
func (journal *JobJournal) load() error {
payload, err := os.ReadFile(journal.path)
if err != nil {
if errors.Is(err, os.ErrNotExist) {
return nil
}
return fmt.Errorf("read job journal: %w", err)
}
var snapshot jobJournalFile
if err := json.Unmarshal(payload, &snapshot); err != nil {
return fmt.Errorf("decode job journal: %w", err)
}
if snapshot.Version != jobJournalVersion {
return fmt.Errorf("unsupported job journal version %d", snapshot.Version)
}
for _, job := range snapshot.Active {
if err := validateJournalAssignment(job); err != nil {
return fmt.Errorf("invalid job journal entry: %w", err)
}
if _, exists := journal.active[job.JobID]; exists {
return fmt.Errorf("duplicate job journal entry %q", job.JobID)
}
journal.active[job.JobID] = job
}
for _, result := range snapshot.PendingResults {
assignment, exists := journal.active[result.JobID]
if !exists || result.SessionToken != "" || assignment.Attempt != result.Attempt || assignment.LeaseToken != result.LeaseToken || assignment.RunEndpointID != result.RunEndpointID {
return fmt.Errorf("invalid pending result journal entry for %q", result.JobID)
}
if _, duplicate := journal.pendingResults[result.JobID]; duplicate {
return fmt.Errorf("duplicate pending result journal entry %q", result.JobID)
}
journal.pendingResults[result.JobID] = result
}
for jobID, manifest := range snapshot.PendingActivations {
if _, exists := journal.pendingResults[jobID]; !exists || manifest == "" {
return fmt.Errorf("invalid pending activation journal entry for %q", jobID)
}
journal.pendingActivations[jobID] = manifest
}
return nil
}
func (journal *JobJournal) persistLocked() error {
if journal.path == "" {
return nil
}
jobs := make([]protocol.RunJobAssignment, 0, len(journal.active))
for _, job := range journal.active {
jobs = append(jobs, job)
}
sort.Slice(jobs, func(i, j int) bool { return jobs[i].JobID < jobs[j].JobID })
results := make([]protocol.RunJobResultRequest, 0, len(journal.pendingResults))
for _, result := range journal.pendingResults {
result.SessionToken = ""
results = append(results, result)
}
sort.Slice(results, func(i, j int) bool { return results[i].JobID < results[j].JobID })
activations := make(map[string]string, len(journal.pendingActivations))
for jobID, manifest := range journal.pendingActivations {
activations[jobID] = manifest
}
payload, err := json.MarshalIndent(jobJournalFile{Version: jobJournalVersion, Active: jobs, PendingResults: results, PendingActivations: activations}, "", " ")
if err != nil {
return fmt.Errorf("encode job journal: %w", err)
}
temporary := journal.path + ".tmp"
file, err := os.OpenFile(temporary, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o600)
if err != nil {
return fmt.Errorf("open temporary job journal: %w", err)
}
removeTemporary := true
defer func() {
_ = file.Close()
if removeTemporary {
_ = os.Remove(temporary)
}
}()
if _, err := file.Write(payload); err != nil {
return fmt.Errorf("write job journal: %w", err)
}
if err := file.Sync(); err != nil {
return fmt.Errorf("sync job journal: %w", err)
}
if err := file.Close(); err != nil {
return fmt.Errorf("close job journal: %w", err)
}
if err := os.Rename(temporary, journal.path); err != nil {
return fmt.Errorf("replace job journal: %w", err)
}
removeTemporary = false
if err := os.Chmod(journal.path, 0o600); err != nil {
return fmt.Errorf("secure job journal: %w", err)
}
return nil
}
func validateJournalAssignment(job protocol.RunJobAssignment) error {
if err := protocol.ValidateRunJobAssignment(job); err != nil {
return err
}
if job.Attempt <= 0 || job.LeaseToken == "" {
return fmt.Errorf("job attempt and lease token are required")
}
return nil
}