From 8a78baf1eca20daaa07821b8b346d31ae5972976 Mon Sep 17 00:00:00 2001 From: Cole Robertson Date: Wed, 9 Sep 2026 19:27:06 +0000 Subject: [PATCH 1/2] Run Litestream commands without a shell and surface failures Commands.run joined argv into one shell string, ignored the exit status and discarded stderr, so a failing command returned "" (or [] after table parsing) and looked like success. Arguments with spaces broke, and argv reached the shell unescaped. The runner now uses Open3.popen3 with an argv array, raises CommandFailedException with the exit status and stderr on failure, and takes an explicit output mode: :table (the existing header/rows parsing), :raw, or :json when a caller passes json: true (Litestream >= 0.5). In JSON mode the two opt-in restore skips, which print one logfmt line on stdout with exit 0, come back as {"skipped" => true, "message" => ...} so callers can tell "did nothing" from data. timeout: runs the command in its own process group and TERMs then KILLs it on expiry, raising CommandTimeoutException with the child reaped. The LITESTREAM_INSTALL_DIR note prints once per process. --- CHANGELOG.md | 3 + README.md | 2 + lib/litestream/commands.rb | 92 ++++++++++++--- test/litestream/test_commands.rb | 195 ++++++++++++++++++++++++++++++- 4 files changed, 275 insertions(+), 17 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5b61a0f..6a07d51 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,8 @@ ## [Unreleased] +- Run Litestream commands without a shell and raise on command failures and timeouts. +- Support parsed JSON output from Litestream 0.5 commands with `json: true`. + ## [0.14.0] - 2025-06-14 - Change async behaviour of replicate and other commands ([@hschne](https://github.com/fractaledmind/litestream-ruby/pull/62)) diff --git a/README.md b/README.md index e5c3056..c9173e1 100644 --- a/README.md +++ b/README.md @@ -434,6 +434,8 @@ s3 a295b16a796689f3 1 0 2036 2024-04-17T00:01:19Z In addition to the provided rake tasks, you can also run Litestream commands directly from Ruby. The gem provides a `Litestream::Commands` module that wraps the Litestream CLI commands. This is particularly useful for the introspection commands, as you can use the output in your Ruby code. +Commands raise `Litestream::Commands::CommandFailedException` when the Litestream binary exits non-zero, and the exception message includes stderr. Pass `json: true` to receive parsed JSON when using Litestream 0.5 or newer (Litestream 0.3 does not support the `-json` flag), and use `timeout:` to limit how many seconds a command may run. + The `Litestream::Commands.databases` method returns an array of hashes with the "path" and "replicas" keys for each database: ```ruby diff --git a/lib/litestream/commands.rb b/lib/litestream/commands.rb index 0d86aad..35ebddf 100644 --- a/lib/litestream/commands.rb +++ b/lib/litestream/commands.rb @@ -1,3 +1,5 @@ +require "json" +require "open3" require_relative "upstream" module Litestream @@ -20,6 +22,9 @@ module Commands # raised when a litestream command fails CommandFailedException = Class.new(StandardError) + # raised when a litestream command times out + CommandTimeoutException = Class.new(CommandFailedException) + module Output class << self def format(data) @@ -47,7 +52,10 @@ def executable(exe_path: DEFAULT_DIR) litestream_install_dir = ENV["LITESTREAM_INSTALL_DIR"] if litestream_install_dir if File.directory?(litestream_install_dir) - warn "NOTE: using LITESTREAM_INSTALL_DIR to find litestream executable: #{litestream_install_dir}" + unless @litestream_install_dir_noted + warn "NOTE: using LITESTREAM_INSTALL_DIR to find litestream executable: #{litestream_install_dir}" + @litestream_install_dir_noted = true + end exe_path = litestream_install_dir exe_file = File.expand_path(File.join(litestream_install_dir, "litestream")) else @@ -135,14 +143,19 @@ def wal(database, **argv) private def execute(command, argv = {}, database = nil, tabled_output: true) - cmd = prepare(command, argv, database) - results = run(cmd, tabled_output: tabled_output) - - if Array === results && results.one? && results[0]["level"] == "ERROR" - raise CommandFailedException, "Failed to execute `#{cmd.join(" ")}`; Reason: #{results[0]["error"]}" + argv = argv.stringify_keys + timeout = argv.delete("timeout") + output = if argv.delete("json") + argv["-json"] = nil + :json + elsif tabled_output + :table else - results + :raw end + + cmd = prepare(command, argv, database) + run(cmd, output: output, timeout: timeout) end def prepare(command, argv = {}, database = nil) @@ -154,18 +167,71 @@ def prepare(command, argv = {}, database = nil) args = { "--config" => Litestream.config_path.to_s - }.merge(argv.stringify_keys).to_a.flatten.compact - cmd = [executable, command, *args, database].compact + }.merge(argv.stringify_keys).to_a.flatten.compact.map(&:to_s) + cmd = [executable, command, *args, database].compact.map(&:to_s) puts cmd.inspect if ENV["DEBUG"] cmd end - def run(cmd, tabled_output:) - stdout = `#{cmd.join(" ")}`.chomp - return stdout unless tabled_output + # Runs the command without a shell and returns its parsed stdout. A non-zero + # exit raises with stderr. With a timeout, the command runs in its own + # process group and is killed (TERM, then KILL) when the deadline passes. + def run(cmd, output:, timeout: nil) + stdin, stdout, stderr, wait_thread = Open3.popen3(*cmd, pgroup: true) + stdin.close + stdout_reader = Thread.new { stdout.read } + stderr_reader = Thread.new { stderr.read } + + # The readers finish when the last process holding the pipes exits, so + # waiting on them covers descendants the direct child may have left behind. + unless wait_thread.join(timeout) && stdout_reader.join(timeout) && stderr_reader.join(timeout) + kill_process_group("TERM", wait_thread.pid) + wait_thread.join(1) + kill_process_group("KILL", wait_thread.pid) + wait_thread.join + [stdout_reader, stderr_reader].each(&:join) + raise CommandTimeoutException, "Failed to execute `#{cmd[1]}`: timed out after #{timeout} seconds" + end + + status = wait_thread.value + unless status.success? + raise CommandFailedException, "Failed to execute `#{cmd[1]}` (exit status #{status.exitstatus}): #{stderr_reader.value.strip[0, 500]}" + end + + case output + when :json then parse_json(cmd, stdout_reader.value) + when :table then parse_table(stdout_reader.value) + else stdout_reader.value + end + ensure + [stdout_reader, stderr_reader].each { |reader| reader&.join } + [stdin, stdout, stderr].each { |io| io&.close unless io&.closed? } + end + + def kill_process_group(signal, pid) + Process.kill(signal, -pid) + rescue Errno::ESRCH + end + + # Two opt-in restore skips (-if-db-not-exists when the output exists, + # -if-replica-exists with no backups) exit 0 and print one logfmt line on + # stdout instead of JSON. They come back as {"skipped" => true, "message" => ...}. + def parse_json(cmd, stdout) + stdout = stdout.strip + return if stdout.empty? + return JSON.parse(stdout) if stdout.start_with?("{", "[") + + skipped = stdout.match(/\Atime=\S+ level=\S+ msg=(?:"([^"]*)"|(\S+))/) + return {"skipped" => true, "message" => skipped[1] || skipped[2]} if skipped + + raise CommandFailedException, "Unexpected output from `#{cmd[1]}`: #{stdout.lines.first.to_s.strip[0, 200]}" + end + + def parse_table(stdout) + keys, *rows = stdout.strip.split("\n").map { _1.split(/\s+/) } + return [] unless keys - keys, *rows = stdout.split("\n").map { _1.split(/\s+/) } rows.map { keys.zip(_1).to_h } end diff --git a/test/litestream/test_commands.rb b/test/litestream/test_commands.rb index 9b96224..49a803b 100644 --- a/test/litestream/test_commands.rb +++ b/test/litestream/test_commands.rb @@ -1,4 +1,6 @@ require "test_helper" +require "fileutils" +require "tmpdir" class TestCommands < ActiveSupport::TestCase def run @@ -204,7 +206,7 @@ def test_restore_with_string_option assert_equal "--config", argv[0] assert_match Regexp.new("dummy/config/litestream.yml"), argv[1] assert_equal "--parallelism", argv[2] - assert_equal 10, argv[3] + assert_equal "10", argv[3] assert_equal "db/test.sqlite3", argv[4] end Litestream::Commands.stub :run, stub do @@ -482,7 +484,7 @@ def test_generations_with_string_option assert_equal "--config", argv[0] assert_match Regexp.new("dummy/config/litestream.yml"), argv[1] assert_equal "--parallelism", argv[2] - assert_equal 10, argv[3] + assert_equal "10", argv[3] assert_equal "db/test.sqlite3", argv[4] end Litestream::Commands.stub :run, stub do @@ -615,7 +617,7 @@ def test_snapshots_with_string_option assert_equal "--config", argv[0] assert_match Regexp.new("dummy/config/litestream.yml"), argv[1] assert_equal "--parallelism", argv[2] - assert_equal 10, argv[3] + assert_equal "10", argv[3] assert_equal "db/test.sqlite3", argv[4] end Litestream::Commands.stub :run, stub do @@ -748,7 +750,7 @@ def test_wal_with_string_option assert_equal "--config", argv[0] assert_match Regexp.new("dummy/config/litestream.yml"), argv[1] assert_equal "--parallelism", argv[2] - assert_equal 10, argv[3] + assert_equal "10", argv[3] assert_equal "db/test.sqlite3", argv[4] end Litestream::Commands.stub :run, stub do @@ -840,6 +842,191 @@ def test_wal_does_not_set_env_var_from_config_when_env_vars_already_set end end + class TestRunner < ActiveSupport::TestCase + def setup + @tmpdir = Dir.mktmpdir + @litestream_install_dir = ENV["LITESTREAM_INSTALL_DIR"] + @executable = File.join(@tmpdir, "litestream") + File.write(@executable, <<~SH) + #!/bin/sh + shift + while [ "$#" -gt 0 ]; do + case "$1" in + --table) + printf 'path replicas\n/tmp/app.sqlite3 s3\n' + exit 0 + ;; + --json-object) + printf '{"txid":"abc123"}\n' + printf 'time=2026-01-01T00:00:00Z level=INFO msg="restore complete"\n' >&2 + exit 0 + ;; + --json-array) + printf '[{"path":"/tmp/app.sqlite3","replicas":["s3"]}]\n' + exit 0 + ;; + --json-empty-list) + printf '[]\n' + exit 0 + ;; + --json-empty) + exit 0 + ;; + --skip-db) + printf 'time=2026-01-01T00:00:00Z level=INFO msg="database already exists, skipping"\n' + exit 0 + ;; + --skip-replica) + printf 'time=2026-01-01T00:00:00Z level=INFO msg="no matching backups found"\n' + exit 0 + ;; + --fail) + printf 'Error: database not found in config\n' >&2 + exit 7 + ;; + --echo) + shift + printf '%s\n' "$1" + exit 0 + ;; + --sleep) + shift + printf '%s\n' "$$" > "$1" + sleep 30 + exit 0 + ;; + --sleep-in-child) + shift + (trap '' TERM; printf '%s\n' "$$" > "$1"; sleep 30) & + exit 0 + ;; + esac + shift + done + SH + File.chmod(0o755, @executable) + Litestream.config_path = File.join(@tmpdir, "litestream.yml") + end + + def teardown + if @litestream_install_dir + ENV["LITESTREAM_INSTALL_DIR"] = @litestream_install_dir + else + ENV.delete("LITESTREAM_INSTALL_DIR") + end + Litestream::Commands.remove_instance_variable(:@litestream_install_dir_noted) if Litestream::Commands.instance_variable_defined?(:@litestream_install_dir_noted) + Litestream.config_path = nil + FileUtils.remove_entry(@tmpdir) + end + + def test_successful_table_output + assert_equal [{"path" => "/tmp/app.sqlite3", "replicas" => "s3"}], run_with_fake { Litestream::Commands.databases("--table" => nil) } + end + + def test_successful_json_object + assert_equal({"txid" => "abc123"}, run_with_fake { Litestream::Commands.restore("db.sqlite3", json: true, "--json-object": nil) }) + end + + def test_successful_json_array + expected = [{"path" => "/tmp/app.sqlite3", "replicas" => ["s3"]}] + + assert_equal expected, run_with_fake { Litestream::Commands.databases(json: true, "--json-array": nil) } + end + + def test_empty_json_list + assert_equal [], run_with_fake { Litestream::Commands.databases("json" => true, "--json-empty-list" => nil) } + end + + def test_empty_json_output + assert_nil run_with_fake { Litestream::Commands.databases("json" => true, "--json-empty" => nil) } + end + + def test_if_db_not_exists_skip_output + expected = {"skipped" => true, "message" => "database already exists, skipping"} + + assert_equal expected, run_with_fake { Litestream::Commands.restore("db.sqlite3", json: true, "--skip-db": nil) } + end + + def test_if_replica_exists_skip_output + expected = {"skipped" => true, "message" => "no matching backups found"} + + assert_equal expected, run_with_fake { Litestream::Commands.restore("db.sqlite3", json: true, "--skip-replica": nil) } + end + + def test_nonzero_exit_raises_with_status_and_stderr + error = assert_raises(Litestream::Commands::CommandFailedException) do + run_with_fake { Litestream::Commands.databases("--fail" => nil) } + end + + assert_includes error.message, "databases" + assert_includes error.message, "exit status 7" + assert_includes error.message, "Error: database not found in config" + end + + def test_argument_with_space_is_passed_intact + output = run_with_fake { Litestream::Commands.restore("db.sqlite3", "--echo" => "argument with space") } + + assert_equal "argument with space\n", output + end + + def test_timeout_kills_and_reaps_child + pid_file = File.join(@tmpdir, "pid") + + assert_raises(Litestream::Commands::CommandTimeoutException) do + run_with_fake { Litestream::Commands.databases(**{"timeout" => 0.1, "--sleep" => pid_file}) } + end + + pid = wait_for_pid(pid_file) + assert_raises(Errno::ECHILD) { Process.wait(pid, Process::WNOHANG) } + end + + def test_timeout_kills_a_descendant_that_outlives_the_direct_child + pid_file = File.join(@tmpdir, "pid") + + assert_raises(Litestream::Commands::CommandTimeoutException) do + run_with_fake { Litestream::Commands.databases(**{"timeout" => 0.1, "--sleep-in-child" => pid_file}) } + end + + pid = wait_for_pid(pid_file) + sleep 0.05 + assert_raises(Errno::ESRCH) { Process.kill(0, pid) } + end + + def test_integer_option_values_are_passed_as_strings + output = run_with_fake { Litestream::Commands.restore("/tmp/app.sqlite3", "--echo" => 10) } + + assert_equal "10\n", output + end + + def test_install_dir_note_is_printed_once + ENV["LITESTREAM_INSTALL_DIR"] = @tmpdir + + _stdout, stderr = capture_io do + Litestream::Commands.executable + Litestream::Commands.executable + end + + assert_equal 1, stderr.scan("NOTE: using LITESTREAM_INSTALL_DIR").size + end + + private + + def run_with_fake(&block) + Litestream::Commands.stub(:executable, @executable, &block) + end + + def wait_for_pid(pid_file) + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + 5 + until File.exist?(pid_file) && !File.read(pid_file).strip.empty? + flunk "fake litestream never wrote its pid" if Process.clock_gettime(Process::CLOCK_MONOTONIC) > deadline + sleep 0.01 + end + pid = File.read(pid_file).to_i + assert_operator pid, :>, 0 + pid + end + end + class TestOutput < ActiveSupport::TestCase def test_output_formatting_generates_table_with_data data = [ From 859616c07a027827f95a99e04bd0106f94a5501d Mon Sep 17 00:00:00 2001 From: Cole Robertson Date: Mon, 14 Sep 2026 18:55:16 +0000 Subject: [PATCH 2/2] Check the exit status of the in-process replicate The async path gets its failures from the caller, but the foreground path ran litestream through IO.popen and never looked at $?. `rails litestream:replicate` therefore exited 0 when litestream failed to start, so a supervisor saw a clean stop rather than a crash to restart, and nothing on the way out said why. A signal is still a normal stop: replicate runs until something signals it, so only a non-zero exit raises. replicate's bare rescue re-wrapped every StandardError into a message built from the whole command line, which swallowed the exit status this adds, so CommandFailedException now passes through it untouched. --- lib/litestream/commands.rb | 10 ++++++++++ test/litestream/test_commands.rb | 30 ++++++++++++++++++++++++++++++ 2 files changed, 40 insertions(+) diff --git a/lib/litestream/commands.rb b/lib/litestream/commands.rb index 35ebddf..780719b 100644 --- a/lib/litestream/commands.rb +++ b/lib/litestream/commands.rb @@ -108,6 +108,8 @@ def executable(exe_path: DEFAULT_DIR) def replicate(async: false, **argv) cmd = prepare("replicate", argv) run_replicate(cmd, async: async) + rescue CommandFailedException + raise rescue raise CommandFailedException, "Failed to execute `#{cmd.join(" ")}`" end @@ -243,6 +245,14 @@ def run_replicate(cmd, async:) IO.popen(cmd, err: [:child, :out]) do |io| io.each_line { |line| puts line } end + status = $? + # `replicate` runs until it is signalled, so a signal is how it is + # meant to end. A non-zero exit is litestream failing to start or + # dying on its own, which is invisible today: the task exits 0 and a + # supervisor sees a clean stop rather than a crash to restart. + if status && !status.success? && !status.signaled? + raise CommandFailedException, "Failed to execute `replicate` (exit status #{status.exitstatus})" + end end end end diff --git a/test/litestream/test_commands.rb b/test/litestream/test_commands.rb index 49a803b..72dd8c1 100644 --- a/test/litestream/test_commands.rb +++ b/test/litestream/test_commands.rb @@ -163,6 +163,36 @@ def test_replicate_does_not_set_env_var_from_config_when_env_vars_already_set assert_equal "original_key", ENV["LITESTREAM_ACCESS_KEY_ID"] assert_equal "original_access", ENV["LITESTREAM_SECRET_ACCESS_KEY"] end + + # These run a real child, and Litestream.replicate_process reads the global + # $? straight after backticks that its own tests stub out, so a failed or + # signalled status left here surfaces over there. Leave a clean one behind. + def teardown + super + system("true") + end + + def test_replicate_in_process_raises_when_litestream_exits_non_zero + error = assert_raises Litestream::Commands::CommandFailedException do + Litestream::Commands.send(:run_replicate, ["/bin/sh", "-c", "echo boom; exit 3"], async: false) + end + + assert_match "exit status 3", error.message + end + + def test_replicate_in_process_does_not_raise_when_litestream_is_signalled + Litestream::Commands.send(:run_replicate, ["/bin/sh", "-c", "kill -TERM $$"], async: false) + end + + def test_replicate_reports_the_exit_status_rather_than_the_command + error = assert_raises Litestream::Commands::CommandFailedException do + Litestream::Commands.stub :prepare, ["/bin/sh", "-c", "exit 3"] do + Litestream::Commands.replicate + end + end + + assert_match "exit status 3", error.message + end end class TestRestoreCommand < TestCommands