Repository navigation
Run Litestream commands without a shell and surface failures #2
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||
|---|---|---|---|---|---|---|---|---|---|---|
| @@ -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 | ||||||||||
|
|
@@ -100,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 | ||||||||||
|
|
@@ -135,14 +145,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 +169,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) | ||||||||||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win Apply one deadline across the three joins. Each join receives the full 🐛 Proposed fix- unless wait_thread.join(timeout) && stdout_reader.join(timeout) && stderr_reader.join(timeout)
+ deadline = timeout && Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout
+ remaining = -> { deadline && [deadline - Process.clock_gettime(Process::CLOCK_MONOTONIC), 0].max }
+ unless [wait_thread, stdout_reader, stderr_reader].all? { |thread| thread.join(remaining.call) }
📝 Committable suggestion
Suggested change
🤖 Prompt for AI Agents |
||||||||||
| 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 | ||||||||||
|
|
||||||||||
|
|
@@ -177,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 | ||||||||||
|
|
||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win
Coerce
timeoutbefore passing it toThread#join.Rake parsing stores
timeout=30as aString.executeremovestimeoutbeforepreparestringifies command arguments, then passes the string directly torun.Thread#joinrequiresnilor aNumerictimeout, so"30"raisesTypeError.🐛 Proposed fix
📝 Committable suggestion
🤖 Prompt for AI Agents