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

230 lines
9.8 KiB
Go

package runtime
import (
"context"
"os"
"path/filepath"
"reflect"
"strings"
"testing"
"browser.local/run/protocol"
)
func TestPersistentJobJournalReloadsAtomicallyWithOwnerOnlyPermissions(t *testing.T) {
root := t.TempDir()
journal, err := NewPersistentJobJournal(root)
if err != nil {
t.Fatalf("new journal: %v", err)
}
assignment := workerJobAssignment(protocol.RunCapabilityProcessStart)
if err := journal.Store(assignment); err != nil {
t.Fatalf("store assignment: %v", err)
}
journalPath := filepath.Join(root, "state", "jobs.json")
info, err := os.Stat(journalPath)
if err != nil {
t.Fatalf("stat journal: %v", err)
}
if info.Mode().Perm() != 0o600 {
t.Fatalf("expected journal mode 0600, got %o", info.Mode().Perm())
}
dirInfo, err := os.Stat(filepath.Dir(journalPath))
if err != nil || dirInfo.Mode().Perm() != 0o700 {
t.Fatalf("expected state directory mode 0700, info=%+v err=%v", dirInfo, err)
}
reloaded, err := NewPersistentJobJournal(root)
if err != nil {
t.Fatalf("reload journal: %v", err)
}
if jobs := reloaded.ActiveJobs(); len(jobs) != 1 || jobs[0].JobID != assignment.JobID || jobs[0].LeaseToken != assignment.LeaseToken || jobs[0].Attempt != assignment.Attempt {
t.Fatalf("unexpected reloaded assignments: %+v", jobs)
}
if err := reloaded.Delete(assignment.JobID); err != nil {
t.Fatalf("delete assignment: %v", err)
}
third, err := NewPersistentJobJournal(root)
if err != nil || third.ActiveCount() != 0 {
t.Fatalf("terminal delete did not persist: count=%d err=%v", third.ActiveCount(), err)
}
}
func TestPersistentJobJournalRejectsCorruptState(t *testing.T) {
root := t.TempDir()
dir := filepath.Join(root, "state")
if err := os.MkdirAll(dir, 0o700); err != nil {
t.Fatalf("create state dir: %v", err)
}
if err := os.WriteFile(filepath.Join(dir, "jobs.json"), []byte(`{"version":1,"active":[`), 0o600); err != nil {
t.Fatalf("write corrupt journal: %v", err)
}
if _, err := NewPersistentJobJournal(root); err == nil {
t.Fatal("expected corrupt journal to fail closed")
}
}
func TestWorkerRestartReconcilesAndRecoversConfirmedAttempt(t *testing.T) {
cfg := workerTestConfig(t)
assignment := workerJobAssignment(protocol.RunCapabilityProcessStart)
assignment.ProgressSequence = 5
journal, err := NewPersistentJobJournal(cfg.WorkspaceRoot)
if err != nil {
t.Fatalf("new journal: %v", err)
}
if err := journal.Store(assignment); err != nil {
t.Fatalf("store interrupted assignment: %v", err)
}
client := newFakeWorkerClient()
client.claimJob = assignment
client.reconcileResponse = protocol.RunJobReconcileResponse{
Accepted: true, RunEndpointID: cfg.RunEndpointID, ConfirmedJobs: []protocol.RunJobAssignment{assignment}, ServerTime: workerTestTime(),
}
restarted, err := NewWorker(cfg, client, WithProcessSupervisor(staticSupervisor{stdout: "recovered\n"}))
if err != nil {
t.Fatalf("restart worker: %v", err)
}
if err := restarted.Register(context.Background()); err != nil {
t.Fatalf("register restarted worker: %v", err)
}
if err := restarted.ReconcileOnce(context.Background()); err != nil {
t.Fatalf("reconcile restarted worker: %v", err)
}
if len(client.reconcileRequests) != 1 || !reflect.DeepEqual(client.reconcileRequests[0].ActiveJobs, []protocol.RunJobReconcileEntry{{JobID: assignment.JobID, LeaseToken: assignment.LeaseToken, Attempt: assignment.Attempt}}) {
t.Fatalf("reconcile omitted attempt evidence: %+v", client.reconcileRequests)
}
if err := restarted.RecoverActiveJobs(context.Background()); err != nil {
t.Fatalf("recover confirmed assignment: %v", err)
}
if restarted.journal.ActiveCount() != 0 || len(client.ackRequests) != 1 || len(client.progressRequests) != 1 || client.progressRequests[0].Sequence != 6 || len(client.resultRequests) != 1 {
t.Fatalf("confirmed attempt was not recovered: journal=%d ack=%d result=%d", restarted.journal.ActiveCount(), len(client.ackRequests), len(client.resultRequests))
}
}
func TestWorkerReconcileDiscardsStaleAttemptWithoutExecuting(t *testing.T) {
cfg := workerTestConfig(t)
assignment := workerJobAssignment(protocol.RunCapabilityProcessStart)
journal, err := NewPersistentJobJournal(cfg.WorkspaceRoot)
if err != nil {
t.Fatalf("new journal: %v", err)
}
if err := journal.Store(assignment); err != nil {
t.Fatalf("store assignment: %v", err)
}
client := newFakeWorkerClient()
client.reconcileResponse = protocol.RunJobReconcileResponse{Accepted: true, RunEndpointID: cfg.RunEndpointID, DiscardJobIDs: []string{assignment.JobID}, ServerTime: workerTestTime()}
worker, err := NewWorker(cfg, client)
if err != nil {
t.Fatalf("new worker: %v", err)
}
if err := worker.Register(context.Background()); err != nil {
t.Fatalf("register: %v", err)
}
if err := worker.ReconcileOnce(context.Background()); err != nil {
t.Fatalf("reconcile: %v", err)
}
if worker.journal.ActiveCount() != 0 || len(client.ackRequests) != 0 {
t.Fatalf("stale assignment was not discarded safely: journal=%d ack=%d", worker.journal.ActiveCount(), len(client.ackRequests))
}
}
func TestWorkerRetainsJournalWhenPlatformRejectsResultTransport(t *testing.T) {
cfg := workerTestConfig(t)
client := newFakeWorkerClient()
client.claimJob = workerJobAssignment(protocol.RunCapabilityProcessStart)
client.resultErr = context.DeadlineExceeded
worker, err := NewWorker(cfg, client, WithProcessSupervisor(staticSupervisor{stdout: "completed locally\n"}))
if err != nil {
t.Fatalf("new worker: %v", err)
}
if err := worker.Register(context.Background()); err != nil {
t.Fatalf("register: %v", err)
}
if handled, err := worker.ClaimAndRunOnce(context.Background()); !handled || err == nil {
t.Fatalf("expected retained failed result transport, handled=%v err=%v", handled, err)
}
if worker.journal.ActiveCount() != 1 {
t.Fatalf("result transport failure removed journal entry")
}
reloaded, err := NewPersistentJobJournal(cfg.WorkspaceRoot)
if err != nil {
t.Fatalf("reload retained journal: %v", err)
}
if reloaded.ActiveCount() != 1 {
t.Fatalf("retained journal did not survive restart: count=%d", reloaded.ActiveCount())
}
if pending, ok := reloaded.PendingResult(client.claimJob.JobID); !ok || pending.SessionToken != "" || pending.State != "succeeded" {
t.Fatalf("pending terminal result was not retained safely: result=%+v ok=%v", pending, ok)
}
payload, err := os.ReadFile(filepath.Join(cfg.WorkspaceRoot, "state", "jobs.json"))
if err != nil {
t.Fatalf("read retained journal: %v", err)
}
if strings.Contains(string(payload), "session-token") {
t.Fatalf("journal persisted raw Run session token: %s", payload)
}
client.resultErr = nil
client.reconcileResponse = protocol.RunJobReconcileResponse{Accepted: true, RunEndpointID: cfg.RunEndpointID, ConfirmedJobs: []protocol.RunJobAssignment{client.claimJob}, ServerTime: workerTestTime()}
restarted, err := NewWorker(cfg, client)
if err != nil {
t.Fatalf("restart result worker: %v", err)
}
if err := restarted.Register(context.Background()); err != nil {
t.Fatalf("register result worker: %v", err)
}
if err := restarted.ReconcileOnce(context.Background()); err != nil {
t.Fatalf("reconcile result worker: %v", err)
}
ackCount := len(client.ackRequests)
if err := restarted.RecoverActiveJobs(context.Background()); err != nil {
t.Fatalf("replay pending result: %v", err)
}
if len(client.ackRequests) != ackCount || restarted.journal.ActiveCount() != 0 {
t.Fatalf("pending result replay re-executed work: ack before=%d after=%d journal=%d", ackCount, len(client.ackRequests), restarted.journal.ActiveCount())
}
}
func TestWorkerRecoversAcceptedSelfUpdateResultAndActivatesOnce(t *testing.T) {
cfg := workerTestConfig(t)
client := newFakeWorkerClient()
assignment := workerJobAssignment(protocol.RunCapabilityRunSelfUpdate)
assignment.TargetKey = "run/update"
assignment.InputRef = "artifact://artifact-run-recovery"
client.claimJob = assignment
client.reconcileResponse = protocol.RunJobReconcileResponse{Accepted: true, RunEndpointID: cfg.RunEndpointID, ConfirmedJobs: []protocol.RunJobAssignment{assignment}, ServerTime: workerTestTime()}
journal, err := NewPersistentJobJournal(cfg.WorkspaceRoot)
if err != nil {
t.Fatalf("create self-update journal: %v", err)
}
if err := journal.Store(assignment); err != nil {
t.Fatalf("store self-update assignment: %v", err)
}
pending := protocol.RunJobResultRequest{RunEndpointID: assignment.RunEndpointID, JobID: assignment.JobID, LeaseToken: assignment.LeaseToken, Attempt: assignment.Attempt, State: "succeeded", Progress: protocol.RunJobProgressReport{Percent: 100, Message: "Run update staged"}, ResultRef: "artifact://jobs/run-update/staged", ExecutionResult: protocol.RunJobExecutionResult{Kind: "run.update.staged", Checksum: bytesChecksum([]byte("archive")), Summary: "verified update staged"}}
manifestPath := filepath.Join(cfg.WorkspaceRoot, "self-updates", assignment.JobID, "manifest.json")
if err := journal.StorePendingResult(pending, manifestPath); err != nil {
t.Fatalf("store pending self-update result: %v", err)
}
activator := &recordingSelfUpdateActivator{}
restarted, err := NewWorker(cfg, client, WithSelfUpdateActivator(activator))
if err != nil {
t.Fatalf("restart worker: %v", err)
}
if err := restarted.Register(context.Background()); err != nil {
t.Fatalf("register restarted worker: %v", err)
}
if err := restarted.ReconcileOnce(context.Background()); err != nil {
t.Fatalf("reconcile restarted worker: %v", err)
}
if err := restarted.RecoverActiveJobs(context.Background()); err != nil {
t.Fatalf("recover accepted self-update result: %v", err)
}
if activator.manifestPath != manifestPath || restarted.journal.ActiveCount() != 0 || restarted.journal.PendingActivation(assignment.JobID) != "" {
t.Fatalf("self-update activation was not recovered exactly once: path=%q active=%d pending=%q", activator.manifestPath, restarted.journal.ActiveCount(), restarted.journal.PendingActivation(assignment.JobID))
}
}