Files
browser/platform/service/job_scheduler_test.go

242 lines
11 KiB
Go

package service
import (
"strings"
"testing"
"time"
"browser.local/platform/domain"
"browser.local/platform/repo"
)
func TestDurableJobPlatformRestartPreservesLeaseFencing(t *testing.T) {
store := repo.NewMemoryStore()
stamp := fixedTime
now := func() time.Time { return stamp }
svc, sessionToken := newMutableRunJobService(t, store, now)
createQueuedRunJob(t, svc, "job-restart", "idem-restart")
claim := mustClaimJob(t, svc, sessionToken, "run-local")
restarted := newCoreService(store, now)
ack, err := restarted.AckRunJob(domain.RunJobAck{
RunEndpointID: "run-local", SessionToken: sessionToken, JobID: claim.Job.JobID,
LeaseToken: claim.Job.LeaseToken, Attempt: claim.Job.Attempt, Message: "recovered after restart",
})
if err != nil || ack.Job.State != domain.JobStateRunning {
t.Fatalf("restart ack failed: ack=%+v err=%v", ack, err)
}
stored, err := store.Jobs().Get("job-restart")
if err != nil {
t.Fatalf("get stored job: %v", err)
}
if stored.LeaseTokenHash == "" || stored.LeaseTokenHash == claim.Job.LeaseToken || stored.Attempt != 1 || stored.LeaseSessionGen != 1 {
t.Fatalf("expected hashed durable lease metadata, got %+v", stored)
}
}
func TestDurableJobAckTimeoutBackoffAndAttemptFencing(t *testing.T) {
store := repo.NewMemoryStore()
stamp := fixedTime
now := func() time.Time { return stamp }
svc, sessionToken := newMutableRunJobService(t, store, now)
createQueuedRunJob(t, svc, "job-timeout", "idem-timeout")
first := mustClaimJob(t, svc, sessionToken, "run-local")
stamp = stamp.Add(defaultJobAckTimeout)
_, err := svc.AckRunJob(domain.RunJobAck{RunEndpointID: "run-local", SessionToken: sessionToken, JobID: first.Job.JobID, LeaseToken: first.Job.LeaseToken, Attempt: first.Job.Attempt})
if err == nil || !strings.Contains(err.Error(), "ack deadline") {
t.Fatalf("expected late ack rejection, got %v", err)
}
retrying, _ := svc.GetJob(first.Job.JobID)
if retrying.State != domain.JobStateRetrying || !retrying.NextAttemptAt.Equal(stamp.Add(2*time.Second)) {
t.Fatalf("expected persisted retry wait, got %+v", retrying)
}
empty, err := svc.ClaimRunJob(runJobClaim(sessionToken, "run-local"))
if err != nil || empty.HasJob {
t.Fatalf("job claimed before backoff elapsed: %+v err=%v", empty, err)
}
stamp = retrying.NextAttemptAt
second := mustClaimJob(t, svc, sessionToken, "run-local")
if second.Job.Attempt != first.Job.Attempt+1 || second.Job.LeaseToken == first.Job.LeaseToken {
t.Fatalf("expected fenced second attempt, first=%+v second=%+v", first.Job, second.Job)
}
_, err = svc.CompleteRunJob(domain.RunJobResult{
RunEndpointID: "run-local", SessionToken: sessionToken, JobID: first.Job.JobID,
LeaseToken: first.Job.LeaseToken, Attempt: first.Job.Attempt, State: domain.JobStateSucceeded,
Progress: domain.RunJobProgressReport{Percent: 100}, Message: "late result",
})
if err == nil || !strings.Contains(err.Error(), "attempt or leaseToken") {
t.Fatalf("expected old attempt result rejection, got %v", err)
}
}
func TestDurableJobLeaseExpiryAndRetryBudget(t *testing.T) {
store := repo.NewMemoryStore()
stamp := fixedTime
now := func() time.Time { return stamp }
svc, sessionToken := newMutableRunJobService(t, store, now)
_, err := svc.CreateJob(domain.Job{
ID: "job-retry", RunEndpointID: "run-local", Capability: "process.start", IdempotencyKey: "idem-retry",
RetryPolicy: domain.JobRetryPolicy{MaxAttempts: 2, InitialBackoffSeconds: 3, MaxBackoffSeconds: 3},
})
if err != nil {
t.Fatalf("create retry job: %v", err)
}
first := mustClaimJob(t, svc, sessionToken, "run-local")
mustAckJob(t, svc, sessionToken, first.Job)
stamp = stamp.Add(defaultJobLeaseDuration)
_, err = svc.UpdateRunJobProgress(domain.RunJobProgress{
RunEndpointID: "run-local", SessionToken: sessionToken, JobID: first.Job.JobID,
LeaseToken: first.Job.LeaseToken, Attempt: first.Job.Attempt, Sequence: 1,
Progress: domain.RunJobProgressReport{Percent: 20, Message: "late progress"},
})
if err == nil || !strings.Contains(err.Error(), "lease expired") {
t.Fatalf("expected expired lease rejection, got %v", err)
}
retrying, _ := svc.GetJob(first.Job.JobID)
if retrying.State != domain.JobStateRetrying {
t.Fatalf("expected retrying after lease expiry, got %+v", retrying)
}
stamp = retrying.NextAttemptAt
second := mustClaimJob(t, svc, sessionToken, "run-local")
mustAckJob(t, svc, sessionToken, second.Job)
result, err := svc.CompleteRunJob(domain.RunJobResult{
RunEndpointID: "run-local", SessionToken: sessionToken, JobID: second.Job.JobID,
LeaseToken: second.Job.LeaseToken, Attempt: second.Job.Attempt, State: domain.JobStateFailed,
Progress: domain.RunJobProgressReport{Percent: 100, Message: "still failing"}, Message: "still failing", Retryable: true,
})
if err != nil || result.Job.State != domain.JobStateFailed {
t.Fatalf("expected terminal failure after retry budget, result=%+v err=%v", result, err)
}
}
func TestDurableJobCancellationBeforeAndAfterClaimIsIdempotent(t *testing.T) {
store := repo.NewMemoryStore()
stamp := fixedTime
now := func() time.Time { return stamp }
svc, sessionToken := newMutableRunJobService(t, store, now)
createQueuedRunJob(t, svc, "job-cancel-queued", "idem-cancel-queued")
firstCancel, err := svc.RequestRunJobCancel(domain.RunJobCancelRequest{JobID: "job-cancel-queued", Reason: "operator cancelled queue"})
if err != nil || firstCancel.State != domain.JobStateCancelled || firstCancel.CompletedAt.IsZero() {
t.Fatalf("cancel queued job: result=%+v err=%v", firstCancel, err)
}
repeated, err := svc.RequestRunJobCancel(domain.RunJobCancelRequest{JobID: "job-cancel-queued", Reason: "operator cancelled queue"})
if err != nil || !repeated.CompletedAt.Equal(firstCancel.CompletedAt) {
t.Fatalf("repeat cancel was not idempotent: result=%+v err=%v", repeated, err)
}
createQueuedRunJob(t, svc, "job-cancel-running", "idem-cancel-running")
claim := mustClaimJob(t, svc, sessionToken, "run-local")
mustAckJob(t, svc, sessionToken, claim.Job)
intent, err := svc.RequestRunJobCancel(domain.RunJobCancelRequest{JobID: claim.Job.JobID, Reason: "operator stop"})
if err != nil || intent.State != domain.JobStateRunning || !intent.CompletedAt.IsZero() {
t.Fatalf("cancel active intent: result=%+v err=%v", intent, err)
}
poll, err := svc.PollRunJobCancel(domain.RunJobCancelPoll{
RunEndpointID: "run-local", SessionToken: sessionToken, JobID: claim.Job.JobID,
LeaseToken: claim.Job.LeaseToken, Attempt: claim.Job.Attempt,
})
if err != nil || !poll.HasCancel || poll.Reason != "operator stop" {
t.Fatalf("poll cancel: result=%+v err=%v", poll, err)
}
terminalRequest := domain.RunJobResult{
RunEndpointID: "run-local", SessionToken: sessionToken, JobID: claim.Job.JobID,
LeaseToken: claim.Job.LeaseToken, Attempt: claim.Job.Attempt, State: domain.JobStateCancelled,
Progress: domain.RunJobProgressReport{Percent: 100, Message: "cancelled"}, Message: "cancelled",
}
terminal, err := svc.CompleteRunJob(terminalRequest)
if err != nil || terminal.Job.State != domain.JobStateCancelled {
t.Fatalf("complete cancellation: result=%+v err=%v", terminal, err)
}
if _, err := svc.CompleteRunJob(terminalRequest); err != nil {
t.Fatalf("duplicate cancelled result should be idempotent: %v", err)
}
}
func TestDurableJobReconcileRotatedSessionAndMissingAttempt(t *testing.T) {
store := repo.NewMemoryStore()
stamp := fixedTime
now := func() time.Time { return stamp }
svc, oldSession := newMutableRunJobService(t, store, now)
createQueuedRunJob(t, svc, "job-confirmed", "idem-confirmed")
confirmedClaim := mustClaimJob(t, svc, oldSession, "run-local")
mustAckJob(t, svc, oldSession, confirmedClaim.Job)
createQueuedRunJob(t, svc, "job-missing", "idem-missing")
missingClaim := mustClaimJob(t, svc, oldSession, "run-local")
mustAckJob(t, svc, oldSession, missingClaim.Job)
hello := validRunControlHello()
hello.CapabilityReport.Capabilities = append(hello.CapabilityReport.Capabilities, "process.start")
hello.CapabilityReport.Fingerprint = "cap-jobs-rotated"
rotated, err := svc.RegisterRunHello(hello)
if err != nil {
t.Fatalf("rotate Run session: %v", err)
}
_, err = svc.UpdateRunJobProgress(domain.RunJobProgress{
RunEndpointID: "run-local", SessionToken: oldSession, JobID: confirmedClaim.Job.JobID,
LeaseToken: confirmedClaim.Job.LeaseToken, Attempt: confirmedClaim.Job.Attempt,
Progress: domain.RunJobProgressReport{Percent: 20}, Sequence: 1,
})
if err == nil {
t.Fatal("expected rotated Run session to reject progress")
}
reconciled, err := svc.ReconcileRunJobs(domain.RunJobReconcile{
RunEndpointID: "run-local", SessionToken: rotated.SessionToken,
ActiveJobs: []domain.RunJobReconcileEntry{{JobID: confirmedClaim.Job.JobID, LeaseToken: confirmedClaim.Job.LeaseToken, Attempt: confirmedClaim.Job.Attempt}},
})
if err != nil || len(reconciled.ConfirmedJobs) != 1 || len(reconciled.DiscardJobIDs) != 0 {
t.Fatalf("reconcile rotated session: result=%+v err=%v", reconciled, err)
}
progress, err := svc.UpdateRunJobProgress(domain.RunJobProgress{
RunEndpointID: "run-local", SessionToken: rotated.SessionToken, JobID: confirmedClaim.Job.JobID,
LeaseToken: confirmedClaim.Job.LeaseToken, Attempt: confirmedClaim.Job.Attempt,
Progress: domain.RunJobProgressReport{Percent: 30, Message: "reconciled"}, Sequence: 1,
})
if err != nil || progress.Job.Progress.Percent != 30 {
t.Fatalf("progress after reconcile: result=%+v err=%v", progress, err)
}
missing, _ := svc.GetJob(missingClaim.Job.JobID)
if missing.State != domain.JobStateRetrying || missing.ReconcileOutcome != "missing from Run journal" {
t.Fatalf("expected missing active attempt to retry, got %+v", missing)
}
}
func newMutableRunJobService(t *testing.T, store repo.Store, now func() time.Time) (*CoreService, string) {
t.Helper()
svc := newCoreService(store, now)
hello := validRunControlHello()
hello.CapabilityReport.Capabilities = append(hello.CapabilityReport.Capabilities, "process.start")
hello.CapabilityReport.Fingerprint = "cap-jobs"
result, err := svc.RegisterRunHello(hello)
if err != nil {
t.Fatalf("register Run: %v", err)
}
return svc, result.SessionToken
}
func runJobClaim(sessionToken string, endpointID string) domain.RunJobClaim {
return domain.RunJobClaim{RunEndpointID: endpointID, SessionToken: sessionToken, Capabilities: []string{"process.start"}, Capacity: domain.RunCapacity{MaxJobs: 4}}
}
func mustClaimJob(t *testing.T, svc *CoreService, sessionToken string, endpointID string) domain.RunJobClaimResult {
t.Helper()
claim, err := svc.ClaimRunJob(runJobClaim(sessionToken, endpointID))
if err != nil || !claim.HasJob || claim.Job == nil {
t.Fatalf("claim job: result=%+v err=%v", claim, err)
}
return claim
}
func mustAckJob(t *testing.T, svc *CoreService, sessionToken string, job *domain.RunJobAssignment) {
t.Helper()
if _, err := svc.AckRunJob(domain.RunJobAck{
RunEndpointID: job.RunEndpointID, SessionToken: sessionToken, JobID: job.JobID,
LeaseToken: job.LeaseToken, Attempt: job.Attempt,
}); err != nil {
t.Fatalf("ack job %s: %v", job.JobID, err)
}
}