diff --git a/src/jobu/recovery.cpp b/src/jobu/recovery.cpp index 381b16e..243057e 100644 --- a/src/jobu/recovery.cpp +++ b/src/jobu/recovery.cpp @@ -1,5 +1,6 @@ #include "recovery.hpp" +#include "job_lifecycle_priv.hpp" #include "recovery_priv.hpp" #include "recovery_repository_priv.hpp" #include "run_repository_priv.hpp" @@ -46,6 +47,7 @@ class Recovery final { std::function const& should_stop) : _database{database} , _repository{database, attributes} + , _lifecycle{database, attributes} , _runs{database, attributes} , _scheduler{database, attributes} , _cron{cron} @@ -184,6 +186,10 @@ class Recovery final { [this](auto after) { return _repository.list_jobs(_batch_size, after); }, [](auto const& row) { return row.id; }, [this, final](JobDefinition const& job) -> RecoveryResult<> { + auto lifecycle = _lifecycle.validate_job_lifecycle(job); + if (!lifecycle) { + return RecoveryResult<>::failure(std::move(lifecycle).error()); + } if (!final) { return RecoveryResult<>::success(); } @@ -255,6 +261,7 @@ class Recovery final { constexpr auto counters = std::array{&RecoveryReport::interrupted_attempts, &RecoveryReport::retrying_runs, &RecoveryReport::terminal_runs, + &RecoveryReport::finished_jobs, &RecoveryReport::inserted_successors, &RecoveryReport::suspended_jobs, &RecoveryReport::suspended_queues}; @@ -298,6 +305,14 @@ class Recovery final { .retrying_runs = decision->retry ? 1U : 0U, .terminal_runs = decision->retry ? 0U : 1U}; if (!decision->retry) { + // Reconcile the final one-time definition in the same unit as its terminal run. + // A retry remains outstanding work and cannot finish the definition. + auto finished = _lifecycle.finish_after_terminal_run(run.id, _recovery_time); + if (!finished) { + return RecoveryResult::failure(std::move(finished).error()); + } + delta.finished_jobs = *finished ? 1U : 0U; + auto successor = _repository.insert_interrupted_successor(run.id, lower_bound, _cron, _generator); if (!successor) { return RecoveryResult::failure(std::move(successor).error()); @@ -341,6 +356,7 @@ class Recovery final { jb::db::Database& _database; RecoveryRepository _repository; + JobLifecycleRepository _lifecycle; RunRepository _runs; SchedulerRepository _scheduler; CronEngine const& _cron; diff --git a/src/jobu/recovery.hpp b/src/jobu/recovery.hpp index 5aacd86..dc29135 100644 --- a/src/jobu/recovery.hpp +++ b/src/jobu/recovery.hpp @@ -36,6 +36,8 @@ struct RecoveryReport { std::uint64_t retrying_runs{}; /// Interrupted runs made terminal instead of retried. std::uint64_t terminal_runs{}; + /// One-time definitions made terminal by interrupted-run processing in this invocation. + std::uint64_t finished_jobs{}; /// Recurring runs inserted during interruption or independent missing-successor repair. std::uint64_t inserted_successors{}; /// Jobs changed from Suspending to Suspended after their running work was gone. @@ -55,7 +57,8 @@ struct RecoveryReport { /// and reopen after storage failure and rerun recovery before serving. A commit error can also /// mean the current unit committed but its acknowledgement was lost; inspect reopened state. /// Repeated recovery is safe and does not duplicate retries or recurring successors. Success proves -/// no Running work remains, required recurring work exists, barriers are valid, and suspensions drained. +/// no Running work remains, current job lifecycles and barriers are valid, required recurring +/// work exists, and suspensions are drained. /// /// @return Committed change counts, `jobu.recovery.invalid_options` for an out-of-range batch /// size, `jobu.recovery.invariant` for inconsistent durable state or counter overflow, or a diff --git a/test/jobu-recovery-public-header-test.cpp b/test/jobu-recovery-public-header-test.cpp index e5206c2..e5a0656 100644 --- a/test/jobu-recovery-public-header-test.cpp +++ b/test/jobu-recovery-public-header-test.cpp @@ -6,6 +6,7 @@ static_assert(std::is_default_constructible_v); static_assert(std::is_copy_constructible_v); static_assert(jb::jobu::RecoveryOptions{}.scan_batch_size == 256); static_assert(jb::jobu::RecoveryReport{}.interrupted_attempts == 0); +static_assert(jb::jobu::RecoveryReport{}.finished_jobs == 0); using RecoveryFunction = jb::core::Result (*)(jb::db::Database&, diff --git a/test/jobud-recovery-integration-test.cpp b/test/jobud-recovery-integration-test.cpp index 220903a..e60514b 100644 --- a/test/jobud-recovery-integration-test.cpp +++ b/test/jobud-recovery-integration-test.cpp @@ -176,14 +176,14 @@ class CrashFixture final { REQUIRE(sqlite3_busy_timeout(raw, 100) == SQLITE_OK); } - void require_schema_startup_failure() + void require_startup_failure(std::string_view expected_code) { launch(false); until([&] { return exit.has_value(); }, true); INFO(log); REQUIRE(exit->kind == ProcessExitKind::Exited); CHECK(exit->exit_code == EXIT_FAILURE); - CHECK(log.find("jobu.schema.") != std::string::npos); + CHECK(log.find(expected_code) != std::string::npos); CHECK_FALSE(std::filesystem::exists(socket_path)); daemon.reset(); REQUIRE(storage.database.open()); @@ -422,8 +422,20 @@ TEST_CASE("daemon crash recovers each runner under both policies and exhausted r fixture.start(type == JobType::Cli); REQUIRE(fixture.count("SELECT count(*) FROM jobu_attempts WHERE state='running'") == 0); REQUIRE(fixture.count("SELECT count(*) FROM jobu_runs") == 1); + auto const retry = policy == RecoveryPolicy::RetryInterrupted && !exhausted; + CHECK(fixture.count(retry ? "SELECT count(*) FROM jobu_jobs WHERE state='active'" + : "SELECT count(*) FROM jobu_jobs WHERE state='failed'") == 1); fixture.crash(); - require_interrupted(fixture, running, policy == RecoveryPolicy::RetryInterrupted && !exhausted); + require_interrupted(fixture, running, retry); + JobRepository jobs{fixture.storage.database, fixture.storage.registry}; + auto current = jobs.find_by_id(job.id, false); + REQUIRE(current); + REQUIRE(*current); + CHECK((*current)->state == (retry ? JobState::Active : JobState::Failed)); + CHECK((*current)->revision == job.revision + (retry ? 0U : 1U)); + if (!retry) { + CHECK((*current)->updated_at == fixture.read_run(seed.run.id).run.completed_at); + } fixture.unchanged_restart(type == JobType::Cli); } @@ -663,6 +675,24 @@ TEST_CASE("daemon validates the current schema before recovery and serving", "[j fixture.unchanged_restart(); } +TEST_CASE("daemon rejects an unfinished one-time owner with no work before serving", "[jobud][recovery][integration]") +{ + CrashFixture fixture; + auto queue = recovery_queue(recovery_id(1)); + fixture.storage.insert_queue(queue); + auto broken = fixture.storage.make_job(recovery_id(2), queue.id, JobType::Http); + fixture.storage.insert_job(broken); + auto healthy = fixture.storage.make_job(recovery_id(4), queue.id, JobType::Http); + fixture.storage.insert_job(healthy); + auto running = fixture.storage.make_run(recovery_id(5), healthy, RunState::Running); + fixture.storage.insert_run(running); + + auto before = storage_snapshot(fixture.storage.database); + fixture.require_startup_failure("jobu.storage.invariant"); + CHECK(storage_snapshot(fixture.storage.database) == before); + fixture.storage.require_run(running); +} + TEST_CASE("daemon schema rejection leaves recovery rows untouched and never listens", "[jobud][recovery][schema][integration]") { @@ -683,7 +713,7 @@ TEST_CASE("daemon schema rejection leaves recovery rows untouched and never list REQUIRE(query.exec(corrupt)); } auto const before = storage_snapshot(fixture.storage.database); - fixture.require_schema_startup_failure(); + fixture.require_startup_failure("jobu.schema."); CHECK(storage_snapshot(fixture.storage.database) == before); fixture.storage.require_run(running); } diff --git a/test/jobud-runtime-test.cpp b/test/jobud-runtime-test.cpp index cf7dcc3..8295d59 100644 --- a/test/jobud-runtime-test.cpp +++ b/test/jobud-runtime-test.cpp @@ -181,6 +181,10 @@ struct RuntimeFixture { -> RecoveryRunFixture { auto job = storage.make_job(recovery_id(suffix + 10), recovery_id(1), type); + if (state == RunState::Succeeded) { + // Startup requires a finished one-time definition when its final run has completed. + job.state = JobState::Succeeded; + } auto run = storage.make_run(recovery_id(suffix + 100), job, state); storage.insert_job(job); storage.insert_run(run); @@ -1014,6 +1018,7 @@ TEST_CASE("Daemon history and statistics RPC read retained data without exposing auto other_queue = recovery_queue(recovery_id(2)); fixture.storage.insert_queue(other_queue); auto other_job = fixture.storage.make_job(recovery_id(22), other_queue.id); + other_job.state = JobState::Succeeded; auto other_run = fixture.storage.make_run(recovery_id(202), other_job, RunState::Succeeded); auto captured = ByteBuffer(65'536U, std::byte{0xff}); other_run.attempts.back().output = jb::jobu::detail::AttemptOutput{.stdout_bytes = captured}; diff --git a/test/recovery-service-test.cpp b/test/recovery-service-test.cpp index f6afa21..1c72e9d 100644 --- a/test/recovery-service-test.cpp +++ b/test/recovery-service-test.cpp @@ -29,7 +29,9 @@ #include #include #include +#include #include +#include #include #include #include @@ -127,6 +129,7 @@ void require_report(RecoveryReport const& actual, RecoveryReport const& expected CHECK(actual.interrupted_attempts == expected.interrupted_attempts); CHECK(actual.retrying_runs == expected.retrying_runs); CHECK(actual.terminal_runs == expected.terminal_runs); + CHECK(actual.finished_jobs == expected.finished_jobs); CHECK(actual.inserted_successors == expected.inserted_successors); CHECK(actual.suspended_jobs == expected.suspended_jobs); CHECK(actual.suspended_queues == expected.suspended_queues); @@ -205,11 +208,17 @@ TEST_CASE("Recovery pages interruption units and preserves complete snapshots", auto result = fixture.recover(batch); REQUIRE(result); auto retries = policy == RecoveryPolicy::RetryInterrupted ? 2U : 0U; - require_report(*result, {.interrupted_attempts = 3, .retrying_runs = retries, .terminal_runs = 3 - retries}); + require_report(*result, + {.interrupted_attempts = 3, + .retrying_runs = retries, + .terminal_runs = 3 - retries, + .finished_jobs = 3 - retries}); for (auto const& original : originals) { - require_interrupted(fixture, - original, - policy == RecoveryPolicy::RetryInterrupted && original.attempts.size() < 3); + auto const retry = policy == RecoveryPolicy::RetryInterrupted && original.attempts.size() < 3; + require_interrupted(fixture, original, retry); + auto current = fixture.job(original.run.job_id); + CHECK(current.state == (retry ? JobState::Active : JobState::Failed)); + CHECK(current.revision == original.run.job_revision + (retry ? 0U : 1U)); } fixture.storage.reopen(); auto again = fixture.recover(batch); @@ -222,6 +231,177 @@ TEST_CASE("Recovery pages interruption units and preserves complete snapshots", } } +TEST_CASE("Recovery finishes the last accepted one-time manual run after original cancellation", + "[jobu][recovery][sqlite]") +{ + auto const retry = GENERATE(false, true); + ServiceFixture fixture; + auto queue = recovery_queue(recovery_id(1), + QueueState::Active, + retry ? RecoveryPolicy::RetryInterrupted : RecoveryPolicy::FailInterrupted); + fixture.storage.insert_queue(queue); + auto job = fixture.storage.make_job(recovery_id(2), queue.id); + job.attributes.at("retry.initial_delay").data = Duration{5s}; + fixture.storage.insert_job(job); + auto original = fixture.storage.make_run(recovery_id(3), job, RunState::Cancelled); + auto manual = fixture.storage.make_run(recovery_id(4), job, RunState::Running, 0, RunOrigin::Manual); + fixture.storage.insert_run(original); + fixture.storage.insert_run(manual); + + auto report = fixture.recover(); + REQUIRE(report); + require_report(*report, + {.interrupted_attempts = 1, + .retrying_runs = retry ? 1U : 0U, + .terminal_runs = retry ? 0U : 1U, + .finished_jobs = retry ? 0U : 1U}); + fixture.storage.require_run(original); + require_interrupted(fixture, manual, retry); + auto current = fixture.job(job.id); + CHECK(current.state == (retry ? JobState::Active : JobState::Failed)); + CHECK(current.revision == job.revision + (retry ? 0U : 1U)); + RunRepository runs{fixture.storage.database, fixture.storage.registry}; + auto successor = runs.find_schedule_owned(job.id); + REQUIRE(successor); + CHECK_FALSE(*successor); + CHECK(fixture.cron.next_calls().empty()); + + fixture.storage.reopen(); + auto again = fixture.recover(); + REQUIRE(again); + require_report(*again); + auto unchanged = fixture.job(job.id); + CHECK(unchanged.state == current.state); + CHECK(unchanged.revision == current.revision); + CHECK(unchanged.updated_at == current.updated_at); +} + +TEST_CASE("Recovery keeps a one-time definition unfinished while its original occurrence remains", + "[jobu][recovery][sqlite]") +{ + ServiceFixture fixture; + auto queue = recovery_queue(recovery_id(1)); + fixture.storage.insert_queue(queue); + auto job = fixture.storage.make_job(recovery_id(2), queue.id); + fixture.storage.insert_job(job); + auto original = fixture.storage.make_run(recovery_id(3), job); + auto manual = fixture.storage.make_run(recovery_id(4), job, RunState::Running, 0, RunOrigin::Manual); + fixture.storage.insert_run(original); + fixture.storage.insert_run(manual); + + auto report = fixture.recover(); + REQUIRE(report); + require_report(*report, {.interrupted_attempts = 1, .terminal_runs = 1}); + fixture.storage.require_run(original); + require_interrupted(fixture, manual, false); + auto current = fixture.job(job.id); + CHECK(current.state == JobState::Active); + CHECK(current.revision == job.revision); + CHECK(current.updated_at == job.updated_at); + CHECK(fixture.cron.next_calls().empty()); +} + +TEST_CASE("Recovery preserves terminal one-time definitions without retained history", "[jobu][recovery][sqlite]") +{ + auto const state = GENERATE(JobState::Succeeded, JobState::Failed, JobState::Cancelled); + ServiceFixture fixture; + auto queue = recovery_queue(recovery_id(1)); + fixture.storage.insert_queue(queue); + auto job = fixture.storage.make_job(recovery_id(2), queue.id); + job.state = state; + fixture.storage.insert_job(job); + + auto report = fixture.recover(); + REQUIRE(report); + require_report(*report); + auto current = fixture.job(job.id); + CHECK(current.state == state); + CHECK(current.revision == job.revision); + CHECK(current.updated_at == job.updated_at); + CHECK(fixture.cron.next_calls().empty()); + + fixture.storage.reopen(); + auto again = fixture.recover(); + REQUIRE(again); + require_report(*again); +} + +TEST_CASE("Recovery rejects unfinished one-time definitions with no outstanding work before repair", + "[jobu][recovery][sqlite]") +{ + auto const history = GENERATE(false, true); + ServiceFixture fixture; + auto queue = recovery_queue(recovery_id(1)); + fixture.storage.insert_queue(queue); + auto broken = fixture.storage.make_job(recovery_id(2), queue.id); + fixture.storage.insert_job(broken); + if (history) { + fixture.storage.insert_run(fixture.storage.make_run(recovery_id(3), broken, RunState::Failed, 1)); + } + auto healthy = fixture.storage.make_job(recovery_id(4), queue.id); + fixture.storage.insert_job(healthy); + auto running = fixture.storage.make_run(recovery_id(5), healthy, RunState::Running); + fixture.storage.insert_run(running); + + auto rejected = fixture.recover(); + REQUIRE_FALSE(rejected); + CHECK(rejected.error().code == "jobu.storage.invariant"); + fixture.storage.require_run(running); + CHECK(fixture.job(broken.id).state == JobState::Active); +} + +TEST_CASE("Recovery rejects terminal one-time definitions with outstanding work", "[jobu][recovery][sqlite]") +{ + ServiceFixture fixture; + auto queue = recovery_queue(recovery_id(1)); + fixture.storage.insert_queue(queue); + auto job = fixture.storage.make_job(recovery_id(2), queue.id); + job.state = JobState::Failed; + fixture.storage.insert_job(job); + auto work = fixture.storage.make_run(recovery_id(3), job); + fixture.storage.insert_run(work); + + auto rejected = fixture.recover(); + REQUIRE_FALSE(rejected); + fixture.storage.require_run(work); + CHECK(fixture.job(job.id).state == JobState::Failed); +} + +TEST_CASE("Recovery final validation rejects a lifecycle broken after an interruption commit", + "[jobu][recovery][sqlite]") +{ + ServiceFixture fixture; + auto queue = recovery_queue(recovery_id(1)); + fixture.storage.insert_queue(queue); + auto job = fixture.storage.make_job(recovery_id(2), queue.id); + fixture.storage.insert_job(job); + auto running = fixture.storage.make_run(recovery_id(3), job, RunState::Running); + fixture.storage.insert_run(running); + + bool saw_precommit = false; + bool corrupted = false; + auto rejected = fixture.recover(1, [&] { + if (fixture.run(running.run.id).state == RunState::Interrupted && !corrupted) { + if (saw_precommit) { + // Insert a current-format inconsistency after the interruption unit commits. + execute(fixture.storage.database, "UPDATE jobu_jobs SET state = 'active'"); + corrupted = true; + } + else { + saw_precommit = true; + } + } + return false; + }); + REQUIRE(corrupted); + REQUIRE_FALSE(rejected); + CHECK(rejected.error().code == "jobu.storage.invariant"); + require_interrupted(fixture, running, false); + CHECK(fixture.job(job.id).state == JobState::Active); + fixture.storage.reopen(); + CHECK_FALSE(fixture.recover()); +} + TEST_CASE("Recovery repairs recurring work and all drained owner shapes", "[jobu][recovery][sqlite]") { ServiceFixture fixture; @@ -235,6 +415,8 @@ TEST_CASE("Recovery repairs recurring work and all drained owner shapes", "[jobu auto once = fixture.storage.make_job(recovery_id(13), queue.id); once.state = JobState::Suspending; fixture.storage.insert_job(once); + auto waiting = fixture.storage.make_run(recovery_id(21), once, RunState::RetryWait, 1); + fixture.storage.insert_run(waiting); auto original = fixture.storage.make_run(recovery_id(20), running_job, RunState::Running); fixture.storage.insert_run(original); @@ -261,7 +443,10 @@ TEST_CASE("Recovery repairs recurring work and all drained owner shapes", "[jobu CHECK(fixture.job(suspended.id).revision == suspended.revision); auto no_once = runs.find_schedule_owned(once.id); REQUIRE(no_once); - CHECK_FALSE(*no_once); + REQUIRE(*no_once); + CHECK((*no_once)->id == waiting.run.id); + fixture.storage.require_run(waiting); + CHECK(fixture.job(once.id).state == JobState::Suspended); CHECK(fixture.queue(queue.id).state == QueueState::Suspended); CHECK(fixture.queue(empty.id).state == QueueState::Suspended); auto again = fixture.recover(); @@ -308,7 +493,8 @@ TEST_CASE("Recovery checks every family before committing any repair", "[jobu][r fixture.storage.insert_job(job); auto original = fixture.storage.make_run(recovery_id(3), job, RunState::Running); fixture.storage.insert_run(original); - auto historical_job = fixture.storage.make_job(recovery_id(4), queue.id); + auto historical_job = fixture.storage.make_job(recovery_id(4), queue.id); + historical_job.state = JobState::Failed; fixture.storage.insert_job(historical_job); auto historical = fixture.storage.make_run(recovery_id(5), historical_job, RunState::Failed, 1); fixture.storage.insert_run(historical); @@ -379,6 +565,66 @@ TEST_CASE("Recovery cancellation rolls back the complete current repair unit", " .suspended_queues = 1}); } +TEST_CASE("Recovery cancellation rolls back the terminal run and one-time definition together", + "[jobu][recovery][sqlite]") +{ + ServiceFixture fixture; + auto queue = recovery_queue(recovery_id(1), QueueState::Suspending); + fixture.storage.insert_queue(queue); + auto job = fixture.storage.make_job(recovery_id(2), queue.id); + job.state = JobState::Suspending; + fixture.storage.insert_job(job); + auto running = fixture.storage.make_run(recovery_id(3), job, RunState::Running); + fixture.storage.insert_run(running); + + bool saw_terminal = false; + auto cancelled = fixture.recover(1, [&] { + saw_terminal = fixture.run(running.run.id).state == RunState::Interrupted; + if (saw_terminal) { + CHECK(fixture.job(job.id).state == JobState::Failed); + CHECK(fixture.queue(queue.id).state == QueueState::Suspended); + } + return saw_terminal; + }); + REQUIRE_FALSE(cancelled); + CHECK(cancelled.error().code == "jobu.recovery.cancelled"); + REQUIRE(saw_terminal); + fixture.storage.reopen(); + fixture.storage.require_run(running); + CHECK(fixture.job(job.id).state == JobState::Suspending); + CHECK(fixture.job(job.id).revision == job.revision); + CHECK(fixture.queue(queue.id).state == QueueState::Suspending); + + auto recovered = fixture.recover(); + REQUIRE(recovered); + require_report(*recovered, + {.interrupted_attempts = 1, .terminal_runs = 1, .finished_jobs = 1, .suspended_queues = 1}); + CHECK(fixture.job(job.id).state == JobState::Failed); + CHECK(fixture.job(job.id).revision == job.revision + 1); + CHECK(fixture.job(job.id).updated_at == UtcTimePoint{120s}); +} + +TEST_CASE("Recovery rolls back interruption when the final job revision is exhausted", "[jobu][recovery][sqlite]") +{ + ServiceFixture fixture; + auto queue = recovery_queue(recovery_id(1)); + fixture.storage.insert_queue(queue); + auto job = fixture.storage.make_job(recovery_id(2), queue.id); + fixture.storage.insert_job(job); + auto running = fixture.storage.make_run(recovery_id(3), job, RunState::Running); + fixture.storage.insert_run(running); + auto const maximum = std::numeric_limits::max(); + execute(fixture.storage.database, "UPDATE jobu_jobs SET revision = " + std::to_string(maximum)); + + auto rejected = fixture.recover(); + REQUIRE_FALSE(rejected); + CHECK(rejected.error().code == "jobu.job.revision_exhausted"); + fixture.storage.reopen(); + fixture.storage.require_run(running); + CHECK(fixture.job(job.id).state == JobState::Active); + CHECK(fixture.job(job.id).revision == static_cast(maximum)); +} + TEST_CASE("Recovery restart retains committed units and retries only unfinished units", "[jobu][recovery][sqlite]") { ServiceFixture fixture; @@ -403,6 +649,8 @@ TEST_CASE("Recovery restart retains committed units and retries only unfinished CHECK(failed.error().detail.find("private") == std::string::npos); fixture.storage.reopen(); require_interrupted(fixture, first, false); + CHECK(fixture.job(first_job.id).state == JobState::Failed); + CHECK(fixture.job(first_job.id).revision == first_job.revision + 1); fixture.storage.require_run(second); fixture.cron.set_next_error({}); auto resumed = fixture.recover(); diff --git a/test/recovery-transaction-fault-test.cpp b/test/recovery-transaction-fault-test.cpp index 9e74aaf..86c51f2 100644 --- a/test/recovery-transaction-fault-test.cpp +++ b/test/recovery-transaction-fault-test.cpp @@ -55,7 +55,7 @@ auto has_run(Scenario scenario) -> bool return scenario <= Scenario::RecurringSuspension; } -auto boundary(std::string_view sql) -> std::string +auto boundary(std::string_view sql, Scenario scenario) -> std::string { if (sql.starts_with("SELECT id FROM jobu_runs WHERE 1 = 1")) { return sql.find("AND state = :state") == std::string_view::npos ? "scan.runs" : "scan.running"; @@ -88,6 +88,11 @@ auto boundary(std::string_view sql) -> std::string return "repair.successor"; } if (sql.starts_with("UPDATE jobu_jobs SET state = :next_state")) { + // Both transitions share JobRepository::set_state SQL; the seeded scenario identifies + // whether this write finishes a definition or drains its suspension. + if (scenario == Scenario::Terminal || scenario == Scenario::Exhausted) { + return "repair.job_terminal"; + } return "repair.job_suspension"; } if (sql.starts_with("UPDATE jobu_queues SET state = :next_state")) { @@ -101,6 +106,7 @@ void check_report(RecoveryReport const& actual, RecoveryReport const& expected = CHECK(actual.interrupted_attempts == expected.interrupted_attempts); CHECK(actual.retrying_runs == expected.retrying_runs); CHECK(actual.terminal_runs == expected.terminal_runs); + CHECK(actual.finished_jobs == expected.finished_jobs); CHECK(actual.inserted_successors == expected.inserted_successors); CHECK(actual.suspended_jobs == expected.suspended_jobs); CHECK(actual.suspended_queues == expected.suspended_queues); @@ -122,7 +128,7 @@ struct Fixture { explicit Fixture(Scenario scenario = Scenario::Terminal) { - faults->classify = boundary; + faults->classify = [scenario](std::string_view sql) { return boundary(sql, scenario); }; time.set_utc(UtcTimePoint{120s}); if (scenario == Scenario::Retry || scenario == Scenario::Exhausted) { queue.recovery_policy = RecoveryPolicy::RetryInterrupted; @@ -136,6 +142,10 @@ struct Fixture { if (scenario == Scenario::RecurringSuspension || scenario == Scenario::JobSuspension) { job.state = JobState::Suspending; } + if (scenario == Scenario::QueueSuspension) { + // A completed one-time definition needs no retained run history. + job.state = JobState::Failed; + } if (scenario == Scenario::RecurringSuspension || scenario == Scenario::MissingSuccessor) { job.schedule = CronSchedule{.expression = "* * * * *", .timezone = "UTC"}; cron.set_occurrences(std::get(job.schedule), {UtcTimePoint{180s}, UtcTimePoint{300s}}); @@ -146,6 +156,10 @@ struct Fixture { storage.make_run(recovery_id(3), job, RunState::Running, scenario == Scenario::Exhausted ? 2 : 0); storage.insert_run(*original); } + if (scenario == Scenario::JobSuspension) { + // RetryWait is unfinished work, but does not keep suspension draining. + storage.insert_run(storage.make_run(recovery_id(3), job, RunState::RetryWait, 1)); + } faults->calls.clear(); } @@ -216,9 +230,18 @@ struct Fixture { REQUIRE(*owner); REQUIRE(parent); REQUIRE(*parent); - bool const job_drained = job.state == JobState::Suspending; - CHECK((*owner)->state == (job_drained ? JobState::Suspended : job.state)); - CHECK((*owner)->revision == job.revision + (job_drained ? 1 : 0)); + auto expected_state = job.state; + auto expected_revision = job.revision; + if (scenario == Scenario::Terminal || scenario == Scenario::Exhausted) { + expected_state = JobState::Failed; + ++expected_revision; + } + else if (job.state == JobState::Suspending) { + expected_state = JobState::Suspended; + ++expected_revision; + } + CHECK((*owner)->state == expected_state); + CHECK((*owner)->revision == expected_revision); CHECK((*parent)->state == (queue.state == QueueState::Suspending ? QueueState::Suspended : queue.state)); } @@ -242,6 +265,9 @@ auto repair_faults(Scenario scenario) -> std::vector if (has_run(scenario)) { writes = {"repair.attempt", "repair.output", scenario == Scenario::Retry ? "repair.retry" : "repair.terminal"}; } + if (scenario == Scenario::Terminal || scenario == Scenario::Exhausted) { + writes.emplace_back("repair.job_terminal"); + } if (scenario == Scenario::RecurringSuspension || scenario == Scenario::MissingSuccessor) { writes.emplace_back("repair.successor"); } @@ -355,7 +381,8 @@ TEST_CASE("Recovery failures after a committed unit retain progress under both p check_report(*resumed, {.interrupted_attempts = 1, .retrying_runs = scenario == Scenario::Retry ? 1U : 0U, - .terminal_runs = scenario == Scenario::Terminal ? 1U : 0U}); + .terminal_runs = scenario == Scenario::Terminal ? 1U : 0U, + .finished_jobs = scenario == Scenario::Terminal ? 1U : 0U}); fixture.check_interrupted(*fixture.original, scenario == Scenario::Retry); fixture.check_interrupted(second, scenario == Scenario::Retry); fixture.check_idempotent();