Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions src/jobu/recovery.cpp
Original file line number Diff line number Diff line change
@@ -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"
Expand Down Expand Up @@ -46,6 +47,7 @@ class Recovery final {
std::function<bool()> const& should_stop)
: _database{database}
, _repository{database, attributes}
, _lifecycle{database, attributes}
, _runs{database, attributes}
, _scheduler{database, attributes}
, _cron{cron}
Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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};
Expand Down Expand Up @@ -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<RecoveryReport>::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<RecoveryReport>::failure(std::move(successor).error());
Expand Down Expand Up @@ -341,6 +356,7 @@ class Recovery final {

jb::db::Database& _database;
RecoveryRepository _repository;
JobLifecycleRepository _lifecycle;
RunRepository _runs;
SchedulerRepository _scheduler;
CronEngine const& _cron;
Expand Down
5 changes: 4 additions & 1 deletion src/jobu/recovery.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
Expand Down
1 change: 1 addition & 0 deletions test/jobu-recovery-public-header-test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ static_assert(std::is_default_constructible_v<jb::jobu::RecoveryOptions>);
static_assert(std::is_copy_constructible_v<jb::jobu::RecoveryReport>);
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::jobu::RecoveryReport, jb::core::Error> (*)(jb::db::Database&,
Expand Down
38 changes: 34 additions & 4 deletions test/jobud-recovery-integration-test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand Down Expand Up @@ -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);
}

Expand Down Expand Up @@ -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]")
{
Expand All @@ -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);
}
5 changes: 5 additions & 0 deletions test/jobud-runtime-test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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};
Expand Down
Loading