Skip to content
Open
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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))
Expand Down
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
102 changes: 89 additions & 13 deletions lib/litestream/commands.rb
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
require "json"
require "open3"
require_relative "upstream"

module Litestream
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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)
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

Expand All @@ -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
Expand Down
Loading