Compare commits

..
19 Commits
Author SHA1 Message Date
Dan Helfman aba45f03d6 Bump version for release. 2026-01-27 13:09:10 -08:00
Dan Helfman f6124528df When the "unsafe_skip_path_validation_before_create" option is enabled, don't log a warning about it (#1244). 2026-01-25 12:17:35 -08:00
Dan Helfman b67dcf829e Fix a regression in which the ntfy monitoring hook failed to send a ping when the "priority" option was set (#1246). 2026-01-24 16:35:55 -08:00
Dan Helfman 71e25756f2 Fix a regression in which the KeePassXC credential hook password prompt was invisible (#1245). 2026-01-24 16:17:32 -08:00
Dan Helfman ff2f9fd5ee Add an additional test and "fix" code coverage (#1242). 2026-01-24 12:28:33 -08:00
Dan Helfman ca4447ffab Fix "spot" check hang (#1242).
Reviewed-on: https://projects.torsion.org/borgmatic-collective/borgmatic/pulls/1247
2026-01-24 20:20:08 +00:00
Dan Helfman 104fe35e39 Add another test to get some additional coverage that's timing dependent (#1242). 2026-01-23 23:04:15 -08:00
Dan Helfman 248fa1db64 Add automated tests for new code (#1242). 2026-01-23 22:50:37 -08:00
Dan Helfman 97f7c65f6c More refactoring and test fixes (#1242). 2026-01-22 23:01:34 -08:00
Dan Helfman 765eba5315 Get existing unit/integration tests passing (#1242). 2026-01-22 17:29:30 -08:00
Dan Helfman bd051beced Structural refactor just to get code out of log_outputs() and into separate utility functions (#1242). 2026-01-22 14:02:48 -08:00
Dan Helfman d5cd4efecd Fix spot check hang (#1242). 2026-01-22 10:27:47 -08:00
Dan Helfman 13fd225a0b Abolish ICE. 2026-01-21 19:29:22 -08:00
Dan Helfman 87c5863218 Some cleanup and also fix delayed logs (#1242). 2026-01-19 22:22:15 -08:00
Dan Helfman 677871aa89 Initial stab at fixing spot check hang by actually draining and consuming buffers after a process exits (#1242). 2026-01-19 16:26:22 -08:00
Dan Helfman d2390581e7 Fix implicit string concatenation instead of trying to paper over it (incidental work included in #1241). 2026-01-18 20:56:29 -08:00
Dan Helfman efd0f0d618 For the "recreate" action, actually pass the "--dry-run" flag through to Borg instead of just skipping the Borg call (#1241). 2026-01-18 18:39:29 -08:00
Dan Helfman 4ff7dccab4 Expand on archive argument to recreate (#1239).
Reviewed-on: https://projects.torsion.org/borgmatic-collective/borgmatic/pulls/1239
2026-01-17 03:15:21 +00:00
Jason Lingohr 76537f6c11 Expand on archive argument
Make a small but specific help expansion on the `archive` option.
2026-01-17 02:45:15 +00:00
19 changed files with 1289 additions and 205 deletions
+10
View File
@@ -1,3 +1,13 @@
2.1.1
* #1241: For the "recreate" action, actually pass the "--dry-run" flag through to Borg instead of
just skipping the Borg call.
* #1242: Fix a regression in which the "spot" check hung while collecting archive contents.
* #1244: When the "unsafe_skip_path_validation_before_create" option is enabled, don't log a
warning about it.
* #1245: Fix a regression in which the KeePassXC credential hook password prompt was invisible.
* #1246: Fix a regression in which the ntfy monitoring hook failed to send a ping when the
"priority" option was set.
2.1.0
* TL;DR: Many logging, memory, and performance improvements. Mind those breaking changes!
* #485: When running commands (database clients, command hooks, etc.), elevate stderr output to
+3 -2
View File
@@ -18,6 +18,7 @@ def run_recreate(
local_borg_version,
recreate_arguments,
global_arguments,
dry_run_label,
local_path,
remote_path,
):
@@ -25,9 +26,9 @@ def run_recreate(
Run the "recreate" action for the given repository.
'''
if recreate_arguments.archive:
logger.answer(f'Recreating archive {recreate_arguments.archive}')
logger.answer(f'Recreating archive {recreate_arguments.archive}{dry_run_label}')
else:
logger.answer('Recreating repository')
logger.answer(f'Recreating repository{dry_run_label}')
# Collect and process patterns.
processed_patterns = borgmatic.actions.pattern.process_patterns(
+1 -1
View File
@@ -270,7 +270,7 @@ def make_base_create_command( # noqa: PLR0912
working_directory = borgmatic.config.paths.get_working_directory(config)
if config.get('unsafe_skip_path_validation_before_create'):
logger.warning(
logger.debug(
'Skipping pre-backup path validation due to "unsafe_skip_path_validation_before_create" option.'
)
+1 -4
View File
@@ -72,6 +72,7 @@ def recreate_archive(
+ (('--chunker-params', chunker_params) if chunker_params else ())
+ (('--recompress', recompress) if recompress else ())
+ exclude_flags
+ (('--dry-run',) if global_arguments.dry_run else ())
+ (tuple(shlex.split(extra_borg_options)) if extra_borg_options else ())
+ (
(
@@ -94,10 +95,6 @@ def recreate_archive(
)
)
if global_arguments.dry_run:
logger.info('Skipping the archive recreation (dry run)')
return
borgmatic.execute.execute_command(
full_command=recreate_command,
output_log_level=logging.INFO,
+1 -1
View File
@@ -1931,7 +1931,7 @@ def make_parsers(schema, unparsed_arguments): # noqa: PLR0915
)
recreate_group.add_argument(
'--archive',
help='Archive name, hash, or series to recreate',
help='Archive name, hash, or series to recreate, defaults to all archives in the repository (if specified), or all archives across all repositories',
)
recreate_group.add_argument(
'--list',
+1
View File
@@ -441,6 +441,7 @@ def run_actions( # noqa: PLR0912, PLR0915
local_borg_version,
action_arguments,
global_arguments,
dry_run_label,
local_path,
remote_path,
)
+1 -1
View File
@@ -39,7 +39,7 @@ def bash_completion():
'check_version() {',
' local this_script="$(cat "$BASH_SOURCE" 2> /dev/null)"',
' local installed_script="$(borgmatic --bash-completion 2> /dev/null)"',
' if [ "$this_script" != "$installed_script" ] && [ "$installed_script" != "" ];'
' if [ "$this_script" != "$installed_script" ] && [ "$installed_script" != "" ];',
f''' then cat << EOF\n{borgmatic.commands.completion.actions.upgrade_message(
'bash',
'sudo sh -c "borgmatic --bash-completion > $BASH_SOURCE"',
+21
View File
@@ -2342,6 +2342,13 @@ properties:
example: Your backups have started.
priority:
type: string
enum:
- max
- urgent
- high
- default
- low
- min
description: |
The priority to set.
example: min
@@ -2366,6 +2373,13 @@ properties:
example: Your backups have finished.
priority:
type: string
enum:
- max
- urgent
- high
- default
- low
- min
description: |
The priority to set.
example: min
@@ -2390,6 +2404,13 @@ properties:
example: Your backups have failed.
priority:
type: string
enum:
- max
- urgent
- high
- default
- low
- min
description: |
The priority to set.
example: max
+241 -125
View File
@@ -3,6 +3,7 @@ import contextlib
import enum
import json
import logging
import os
import select
import subprocess
import textwrap
@@ -41,7 +42,9 @@ def interpret_exit_code(command, exit_code, borg_local_path=None, borg_exit_code
if exit_code == 0:
return Exit_status.SUCCESS
if borg_local_path and command[0] == borg_local_path:
parsed_command = command.split(' ', 1) if isinstance(command, str) else command
if borg_local_path and parsed_command[0] == borg_local_path:
# First try looking for the exit code in the borg_exit_codes configuration.
for entry in borg_exit_codes or ():
if entry.get('code') == exit_code:
@@ -163,7 +166,9 @@ def parse_log_line(line, log_level, elevate_stderr, borg_local_path, command):
came from stderr and the string "warning:" appears at the start of the log line. In that case,
just elevate the log level to a WARN.
'''
if borg_local_path and command[0] == borg_local_path:
parsed_command = command.split(' ', 1) if isinstance(command, str) else command
if borg_local_path and parsed_command[0] == borg_local_path:
log_record = borg_json_log_line_to_record(line, log_level)
if log_record:
@@ -177,18 +182,20 @@ def parse_log_line(line, log_level, elevate_stderr, borg_local_path, command):
return log_line_to_record(line, log_level)
def handle_log_record(log_record, last_lines):
def handle_log_record(log_record, last_lines=None):
'''
Given a log record to be logged and a rolling list of last lines, append the record's message to
the last lines. Then (if the log level is not None), log the record.
the last lines (if given). Then (if the log level is not None), log the record.
Return the log record.
'''
log_message = log_record.getMessage()
last_lines.append(log_message)
if len(last_lines) > ERROR_OUTPUT_MAX_LINE_COUNT:
last_lines.pop(0)
if last_lines is not None:
last_lines.append(log_message)
if len(last_lines) > ERROR_OUTPUT_MAX_LINE_COUNT:
last_lines.pop(0)
if log_record.levelno is not None:
logger.handle(log_record)
@@ -196,7 +203,205 @@ def handle_log_record(log_record, last_lines):
return log_record
def log_outputs( # noqa: PLR0912
READ_CHUNK_SIZE = 4096
def read_lines(buffer, process, line_separator='\n'):
'''
Given a Python buffer (like stdout) ready for reading, its process, and a line separator,
repeatedly yield a tuple of (decoded) lines from the buffer until the process has exited.
It is assumed that this function's generator is used in conjunction with an external select()
call to know when to read more lines. Otherwise, the generator will busywait if it's called in a
tight loop.
'''
data = ''
while True:
chunk = os.read(buffer.fileno(), READ_CHUNK_SIZE).decode()
if not chunk: # EOF
# The process is still running, so we keep running too.
if process.poll() is None: # pragma: no cover
continue
break
data += chunk
lines = []
# Split the data into lines, holding back anything leftover that might
# be a partial line.
while True:
separator_position = data.find(line_separator)
if separator_position == -1:
break
lines.append(data[:separator_position].rstrip())
data = data[separator_position + 1 :]
yield tuple(lines)
# Yield any leftover data from the end of the buffer.
if data:
yield (data.rstrip(),)
Buffer_reader = collections.namedtuple(
'Buffer_reader',
('lines', 'process'),
)
Process_metadata = collections.namedtuple(
'Process_metadata',
('last_lines', 'capture'),
)
def log_buffer_lines(
buffer_readers, process_metadatas, output_log_level, borg_local_path, capture_stderr=False
):
'''
Given a dict from buffer object to Buffer_reader, a dict from subprocess.Popen() instance to
Process_metadata instance, a requested output log level for stdout, Borg's local path, and
whether to capture stderr, read and log any ready output lines from the buffers. Additionally,
if the log level is None for any log record, then yield those log messages for capture.
This function just does one "turn of the crank" of logging buffer output. It is intended to be
called repeatedly to continue to process buffers.
'''
if not buffer_readers:
return
(ready_buffers, _, _) = select.select(buffer_readers.keys(), [], [])
for ready_buffer in ready_buffers:
reader = buffer_readers[ready_buffer]
# The "ready" process has exited, but it might be a pipe destination with other
# processes (pipe sources) waiting to be read from. So as a measure to prevent
# hangs, vent all processes when one exits.
if reader.process and reader.process.poll() is not None:
for other_process in process_metadatas:
if (
other_process.poll() is None
and other_process.stdout
and other_process.stdout not in buffer_readers
):
# Add the process's output to buffer_readers to ensure it'll get read.
buffer_readers[other_process.stdout] = Buffer_reader(
read_lines(other_process.stdout, other_process), other_process
)
try:
lines = next(reader.lines)
except StopIteration:
continue
for line in lines:
if not line or not reader.process:
continue
# Keep the last few lines of output in case the process errors and we need the
# output for the exception below.
log_record = handle_log_record(
parse_log_line(
line=line,
log_level=output_log_level,
elevate_stderr=(ready_buffer == reader.process.stderr and not capture_stderr),
borg_local_path=borg_local_path,
command=reader.process.args,
),
last_lines=process_metadatas[reader.process].last_lines,
)
if log_record.levelno is None and process_metadatas[reader.process].capture:
yield log_record.getMessage()
def raise_for_process_errors(buffer_readers, process_metadatas, borg_local_path, borg_exit_codes):
'''
Given a dict from buffer object to Buffer_reader, a dict from subprocess.Popen() instance to
Process_metadata instance, Borg's local path, a sequence of exit code configuration dicts, check
the given processes for error or warning exit codes. If found, vent or kill any running
processes. In the case of an error exit code, raise. In the case of warning, return
Exit_status.WARNING. Otherwise, return None.
'''
result_status = None
for process in process_metadatas:
exit_code = process.poll() if buffer_readers else process.wait()
if exit_code is None:
continue
exit_status = interpret_exit_code(process.args, exit_code, borg_local_path, borg_exit_codes)
if exit_status not in {Exit_status.ERROR, Exit_status.WARNING}:
continue
# Something has gone wrong. So vent each process' output buffer to prevent it from
# hanging. And then kill the process.
for other_process in process_metadatas:
if other_process.poll() is None:
other_process.stdout.read(0)
other_process.kill()
if exit_status == Exit_status.WARNING:
result_status = Exit_status.WARNING
continue
last_lines = process_metadatas[process].last_lines
# If an error occurs, include its output in the raised exception so that we don't
# inadvertently hide error output.
if len(last_lines) >= ERROR_OUTPUT_MAX_LINE_COUNT:
last_lines.insert(0, '...')
raise subprocess.CalledProcessError(
exit_code,
command_for_process(process),
'\n'.join(last_lines),
)
return result_status
def log_remaining_buffer_lines(
buffer_readers, process_metadatas, output_log_level, borg_local_path, capture_stderr=False
):
'''
Given a dict from buffer object to Buffer_reader, a dict from subprocess.Popen() instance to
Process_metadata instance, a requested output log level for stdout, Borg's local path, and
whether to capture stderr, drain and log any remaining output lines from the buffers until
they're empty. Additionally, if the log level is None for any log record, then yield those log
messages for capture.
'''
for output_buffer, reader in buffer_readers.items():
if not reader.process:
continue
for lines in reader.lines:
for line in lines:
log_record = handle_log_record(
parse_log_line(
line=line.rstrip(),
log_level=output_log_level,
elevate_stderr=(
output_buffer == reader.process.stderr and not capture_stderr
),
borg_local_path=borg_local_path,
command=reader.process.args,
),
)
if log_record.levelno is None and process_metadatas[reader.process].capture:
yield log_record.getMessage()
def log_outputs(
processes,
exclude_stdouts,
output_log_level,
@@ -221,131 +426,42 @@ def log_outputs( # noqa: PLR0912
buffers. Also note that stdout for a process can be None if output is intentionally not
captured, in which case it won't be logged.
'''
# Map from output buffer to sequence of last lines.
process_last_lines = collections.defaultdict(list)
process_for_output_buffer = {
buffer: process
# Map from output buffer to Process_metadata instance. By convention, the last process is the
# process to capture.
process_metadatas = {
process: Process_metadata(last_lines=[], capture=bool(process == processes[-1]))
for process in processes
}
# Map from buffer to Buffer_reader instance.
buffer_readers = {
buffer: Buffer_reader(read_lines(buffer, process), process)
for process in processes
if process.stdout or process.stderr
for buffer in output_buffers_for_process(process, exclude_stdouts)
}
output_buffers = list(process_for_output_buffer.keys())
process_to_capture = processes[-1]
still_running = True
# Log output for each process until they all exit.
while True: # noqa: PLR1702
if output_buffers:
(ready_buffers, _, _) = select.select(output_buffers, [], [])
# Log output lines for each process until they all exit.
while True:
yield from log_buffer_lines(
buffer_readers, process_metadatas, output_log_level, borg_local_path, capture_stderr
)
for ready_buffer in ready_buffers:
ready_process = process_for_output_buffer.get(ready_buffer)
# The "ready" process has exited, but it might be a pipe destination with other
# processes (pipe sources) waiting to be read from. So as a measure to prevent
# hangs, vent all processes when one exits.
if ready_process and ready_process.poll() is not None:
for other_process in processes:
if (
other_process.poll() is None
and other_process.stdout
and other_process.stdout not in output_buffers
):
# Add the process's output to output_buffers to ensure it'll get read.
output_buffers.append(other_process.stdout)
while True:
line = ready_buffer.readline().rstrip().decode()
if not line or not ready_process:
break
command = (
ready_process.args.split(' ')
if isinstance(ready_process.args, str)
else ready_process.args
)
# Keep the last few lines of output in case the process errors and we need the
# output for the exception below.
log_record = handle_log_record(
parse_log_line(
line=line,
log_level=output_log_level,
elevate_stderr=(
ready_buffer == ready_process.stderr and not capture_stderr
),
borg_local_path=borg_local_path,
command=command,
),
last_lines=process_last_lines[ready_process],
)
if log_record.levelno is None and ready_process == process_to_capture:
yield log_record.getMessage()
if not still_running:
if (
raise_for_process_errors(
buffer_readers, process_metadatas, borg_local_path, borg_exit_codes
)
== Exit_status.WARNING
):
break
still_running = False
if all(process.poll() is not None for process in processes):
break
for process in processes:
exit_code = process.poll() if output_buffers else process.wait()
if exit_code is None:
still_running = True
command = process.args.split(' ') if isinstance(process.args, str) else process.args
continue
command = process.args.split(' ') if isinstance(process.args, str) else process.args
exit_status = interpret_exit_code(command, exit_code, borg_local_path, borg_exit_codes)
if exit_status in {Exit_status.ERROR, Exit_status.WARNING}:
last_lines = process_last_lines[process]
# If an error occurs, include its output in the raised exception so that we don't
# inadvertently hide error output.
for output_buffer in output_buffers_for_process(process, exclude_stdouts):
# Collect any straggling output lines that came in since we last gathered output.
while output_buffer: # pragma: no cover
line = output_buffer.readline().rstrip().decode()
if not line:
break
log_record = handle_log_record(
parse_log_line(
line=line,
log_level=output_log_level,
elevate_stderr=(
output_buffer == process.stderr and not capture_stderr
),
borg_local_path=borg_local_path,
command=command,
),
last_lines=last_lines,
)
if log_record.levelno is None and process == process_to_capture:
yield log_record.getMessage()
if len(last_lines) == ERROR_OUTPUT_MAX_LINE_COUNT:
last_lines.insert(0, '...')
# Something has gone wrong. So vent each process' output buffer to prevent it from
# hanging. And then kill the process.
for other_process in processes:
if other_process.poll() is None:
other_process.stdout.read(0)
other_process.kill()
if exit_status == Exit_status.ERROR:
raise subprocess.CalledProcessError(
exit_code,
command_for_process(process),
'\n'.join(last_lines),
)
still_running = False
break
# Now that all processes have exited, drain and consume any last output.
yield from log_remaining_buffer_lines(
buffer_readers, process_metadatas, output_log_level, borg_local_path, capture_stderr
)
SECRET_COMMAND_FLAG_NAMES = {'--password'}
@@ -496,7 +612,7 @@ def execute_command_and_capture_output(
command,
stdin=input_file,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
stderr=subprocess.PIPE if capture_stderr else None,
shell=shell,
env=environment,
cwd=working_directory,
+13 -3
View File
@@ -22,6 +22,16 @@ def initialize_monitor(
'''
PRIORITY_NAME_TO_ID = {
'max': 5,
'urgent': 5,
'high': 4,
'default': 3,
'low': 2,
'min': 1,
}
def ping_monitor(hook_config, config, config_filename, state, monitoring_log_level, dry_run):
'''
Ping the configured Ntfy topic. Use the given configuration filename in any log entries.
@@ -31,13 +41,13 @@ def ping_monitor(hook_config, config, config_filename, state, monitoring_log_lev
if state.name.lower() in run_states:
dry_run_label = ' (dry run; not actually pinging)' if dry_run else ''
default_priority = PRIORITY_NAME_TO_ID['default']
state_config = hook_config.get(
state.name.lower(),
{
'title': f'A borgmatic {state.name} event happened',
'message': f'A borgmatic {state.name} event happened',
'priority': 'default',
'priority': default_priority,
'tags': 'borgmatic',
},
)
@@ -55,7 +65,7 @@ def ping_monitor(hook_config, config, config_filename, state, monitoring_log_lev
'topic': topic,
'title': state_config.get('title'),
'message': state_config.get('message'),
'priority': state_config.get('priority'),
'priority': PRIORITY_NAME_TO_ID.get(state_config.get('priority'), default_priority),
'tags': state_config.get('tags'),
}
+1 -1
View File
@@ -60,7 +60,7 @@ follows:
unsafe_skip_path_validation_before_create: true
```
However, this is indeed unsafe, and could lead to hangs or data being left out
However, this is indeed unsafe and could lead to hangs or data being left out
of backups. Use this option at your own risk.
+3 -1
View File
@@ -129,7 +129,9 @@ Note the lack of "`//`" after `s3:` or `b2:`.
When selecting your cloud hosting provider, be aware that Amazon in particular
has [financially
supported](https://en.wikipedia.org/wiki/White_House_State_Ballroom) the Trump
regime.
regime. Additionally, U.S. Immigration and Customs Enforcement (ICE) is [powered
by
Amazon](https://medium.com/@noazureforapartheid/microsoft-powers-ice-why-doesnt-microsoft-want-to-talk-about-its-contracts-with-immigration-and-bc04fae8d43b).
## Related documentation
+1 -1
View File
@@ -1,6 +1,6 @@
[project]
name = "borgmatic"
version = "2.1.0"
version = "2.1.1"
authors = [
{ name="Dan Helfman", email="witten@torsion.org" },
]
+63 -18
View File
@@ -8,6 +8,55 @@ from flexmock import flexmock
from borgmatic import execute as module
def test_read_lines_yields_single_line():
process = subprocess.Popen(['echo', 'hi'], stdout=subprocess.PIPE)
assert tuple(module.read_lines(process.stdout, process)) == (('hi',),)
def test_read_lines_yields_single_line_longer_than_chunk_size():
process = subprocess.Popen(
['echo', 'this line is longer than the chunk size'], stdout=subprocess.PIPE
)
assert tuple(flexmock(module, READ_CHUNK_SIZE=16).read_lines(process.stdout, process)) == (
(),
(),
('this line is longer than the chunk size',),
)
def test_read_lines_yields_multiple_lines():
process = subprocess.Popen(['echo', 'hi\nthere'], stdout=subprocess.PIPE)
assert tuple(module.read_lines(process.stdout, process)) == (('hi', 'there'),)
def test_read_lines_yields_multiple_lines_plus_partial_line():
process = subprocess.Popen(['echo', '-n', 'hi\nthere\npartial'], stdout=subprocess.PIPE)
assert tuple(module.read_lines(process.stdout, process)) == (('hi', 'there'), ('partial',))
def test_read_lines_with_longer_running_process_yields_many_lines():
process = subprocess.Popen(
[
sys.executable,
'-c',
"import random, string; print('\\n'.join(random.choice(string.ascii_letters) for _ in range(1000)))",
],
stdout=subprocess.PIPE,
)
assert tuple(module.read_lines(process.stdout, process))
def test_read_lines_yields_nothing():
process = subprocess.Popen(['echo', '-n'], stdout=subprocess.PIPE)
assert tuple(module.read_lines(process.stdout, process)) == ()
def test_log_outputs_logs_each_line_separately():
hi_record = flexmock(
msg='hi',
@@ -269,7 +318,7 @@ def test_log_outputs_kills_other_processes_and_raises_when_one_errors():
other_process,
(),
).and_return((other_process.stdout,))
flexmock(other_process).should_receive('kill').once()
flexmock(other_process).should_call('kill').once()
with pytest.raises(subprocess.CalledProcessError) as error:
tuple(
@@ -291,12 +340,6 @@ def test_log_outputs_kills_other_processes_and_returns_when_one_exits_with_warni
flexmock(module).should_receive('command_for_process').and_return('grep')
process = subprocess.Popen(['grep'], stdout=subprocess.PIPE, stderr=subprocess.STDOUT)
flexmock(module).should_receive('interpret_exit_code').with_args(
['grep'],
None,
'borg',
None,
).and_return(module.Exit_status.SUCCESS)
flexmock(module).should_receive('interpret_exit_code').with_args(
['grep'],
2,
@@ -313,7 +356,7 @@ def test_log_outputs_kills_other_processes_and_returns_when_one_exits_with_warni
None,
'borg',
None,
).and_return(module.Exit_status.SUCCESS)
).and_return(module.Exit_status.STILL_RUNNING)
flexmock(module).should_receive('output_buffers_for_process').with_args(process, ()).and_return(
(process.stdout,),
)
@@ -321,7 +364,7 @@ def test_log_outputs_kills_other_processes_and_returns_when_one_exits_with_warni
other_process,
(),
).and_return((other_process.stdout,))
flexmock(other_process).should_receive('kill').once()
flexmock(other_process).should_call('kill').once()
assert (
tuple(
@@ -370,7 +413,15 @@ def test_log_outputs_vents_other_processes_when_one_exits():
other_process,
(process.stdout,),
).and_return((other_process.stdout,))
flexmock(process.stdout).should_call('readline').at_least().once()
flexmock(module.os).should_call('read').with_args(
process.stderr.fileno(), int
).at_least().once()
flexmock(module.os).should_call('read').with_args(
process.stdout.fileno(), int
).at_least().once()
flexmock(module.os).should_call('read').with_args(
other_process.stdout.fileno(), int
).at_least().once()
assert (
tuple(
@@ -433,12 +484,6 @@ def test_log_outputs_truncates_long_error_output():
flexmock(module).should_receive('command_for_process').and_return('grep')
process = subprocess.Popen(['grep'], stdout=subprocess.PIPE, stderr=subprocess.STDOUT)
flexmock(module).should_receive('interpret_exit_code').with_args(
['grep'],
None,
'borg',
None,
).and_return(module.Exit_status.SUCCESS)
flexmock(module).should_receive('interpret_exit_code').with_args(
['grep'],
2,
@@ -487,8 +532,8 @@ def test_log_outputs_with_unfinished_process_re_polls():
flexmock(module.logger).should_receive('log').never()
flexmock(module).should_receive('interpret_exit_code').and_return(module.Exit_status.SUCCESS)
process = subprocess.Popen(['true'], stdout=subprocess.PIPE, stderr=subprocess.STDOUT)
flexmock(process).should_receive('poll').and_return(None).and_return(0).times(3)
process = subprocess.Popen(['sleep', '0.001'], stdout=subprocess.PIPE)
flexmock(process).should_call('poll').at_least().times(3)
flexmock(module).should_receive('output_buffers_for_process').and_return((process.stdout,))
assert (
+7
View File
@@ -22,6 +22,7 @@ def test_run_recreate_does_not_raise():
local_borg_version=None,
recreate_arguments=flexmock(repository=flexmock(), archive=None),
global_arguments=flexmock(),
dry_run_label='',
local_path=None,
remote_path=None,
)
@@ -45,6 +46,7 @@ def test_run_recreate_with_archive_does_not_raise():
local_borg_version=None,
recreate_arguments=flexmock(repository=flexmock(), archive='test-archive'),
global_arguments=flexmock(),
dry_run_label='',
local_path=None,
remote_path=None,
)
@@ -69,6 +71,7 @@ def test_run_recreate_with_leftover_recreate_archive_raises():
local_borg_version=None,
recreate_arguments=flexmock(repository=flexmock(), archive='test-archive.recreate'),
global_arguments=flexmock(),
dry_run_label='',
local_path=None,
remote_path=None,
)
@@ -93,6 +96,7 @@ def test_run_recreate_with_latest_archive_resolving_to_leftover_recreate_archive
local_borg_version=None,
recreate_arguments=flexmock(repository=flexmock(), archive='latest'),
global_arguments=flexmock(),
dry_run_label='',
local_path=None,
remote_path=None,
)
@@ -122,6 +126,7 @@ def test_run_recreate_with_archive_already_exists_error_raises():
local_borg_version=None,
recreate_arguments=flexmock(repository=flexmock(), archive='test-archive', target=None),
global_arguments=flexmock(),
dry_run_label='',
local_path=None,
remote_path=None,
)
@@ -155,6 +160,7 @@ def test_run_recreate_with_target_and_archive_already_exists_error_raises():
target='target-archive',
),
global_arguments=flexmock(),
dry_run_label='',
local_path=None,
remote_path=None,
)
@@ -188,6 +194,7 @@ def test_run_recreate_with_other_called_process_error_passes_it_through():
target='target-archive',
),
global_arguments=flexmock(),
dry_run_label='',
local_path=None,
remote_path=None,
)
-1
View File
@@ -1128,7 +1128,6 @@ def test_make_base_create_command_with_unsafe_skip_path_validation_before_create
(f'repo::{module.flags.get_default_archive_name_format()}',),
)
flexmock(module).should_receive('validate_planned_backup_paths').never()
flexmock(module.logger).should_receive('warning').once()
module.make_base_create_command(
dry_run=False,
+35 -38
View File
@@ -20,44 +20,6 @@ def insert_execute_command_mock(command, working_directory=None, borg_exit_codes
).once()
def test_recreate_archive_dry_run_skips_execution():
flexmock(module.borgmatic.borg.flags).should_receive('make_exclude_flags').and_return(())
flexmock(module.borgmatic.borg.pattern).should_receive('write_patterns_file').and_return(None)
flexmock(module.borgmatic.borg.flags).should_receive('make_list_filter_flags').and_return('')
flexmock(module.borgmatic.borg.flags).should_receive('make_match_archives_flags').and_return(())
flexmock(module.borgmatic.borg.feature).should_receive('available').and_return(True)
flexmock(module.borgmatic.borg.flags).should_receive(
'make_repository_archive_flags',
).and_return(
(
'--repo',
'repo',
),
)
flexmock(module.borgmatic.execute).should_receive('execute_command').never()
recreate_arguments = flexmock(
repository=flexmock(),
list=None,
target=None,
comment=None,
timestamp=None,
match_archives=None,
)
result = module.recreate_archive(
repository='repo',
archive='archive',
config={},
local_borg_version='1.2.3',
recreate_arguments=recreate_arguments,
global_arguments=flexmock(dry_run=True),
local_path='borg',
)
assert result is None
def test_recreate_calls_borg_with_required_flags():
flexmock(module.borgmatic.borg.flags).should_receive('make_exclude_flags').and_return(())
flexmock(module.borgmatic.borg.pattern).should_receive('write_patterns_file').and_return(None)
@@ -93,6 +55,41 @@ def test_recreate_calls_borg_with_required_flags():
)
def test_recreate_with_dry_run_calls_borg_with_dry_run_flag():
flexmock(module.borgmatic.borg.flags).should_receive('make_exclude_flags').and_return(())
flexmock(module.borgmatic.borg.pattern).should_receive('write_patterns_file').and_return(None)
flexmock(module.borgmatic.borg.flags).should_receive('make_list_filter_flags').and_return('')
flexmock(module.borgmatic.borg.flags).should_receive('make_match_archives_flags').and_return(())
flexmock(module.borgmatic.borg.feature).should_receive('available').and_return(True)
flexmock(module.borgmatic.borg.flags).should_receive(
'make_repository_archive_flags',
).and_return(
(
'--repo',
'repo',
),
)
insert_execute_command_mock(('borg', 'recreate', '--log-json', '--dry-run', '--repo', 'repo'))
module.recreate_archive(
repository='repo',
archive='archive',
config={},
local_borg_version='1.2.3',
recreate_arguments=flexmock(
list=None,
target=None,
comment=None,
timestamp=None,
match_archives=None,
),
global_arguments=flexmock(dry_run=True),
local_path='borg',
remote_path=None,
patterns=None,
)
def test_recreate_with_remote_path():
flexmock(module.borgmatic.borg.flags).should_receive('make_exclude_flags').and_return(())
flexmock(module.borgmatic.borg.pattern).should_receive('write_patterns_file').and_return(None)
+2 -2
View File
@@ -24,7 +24,7 @@ CUSTOM_MESSAGE_PAYLOAD = {
'topic': TOPIC,
'title': CUSTOM_MESSAGE_CONFIG['title'],
'message': CUSTOM_MESSAGE_CONFIG['message'],
'priority': CUSTOM_MESSAGE_CONFIG['priority'],
'priority': 1,
'tags': CUSTOM_MESSAGE_CONFIG['tags'],
}
@@ -34,7 +34,7 @@ def default_message_payload(state=Enum):
'topic': TOPIC,
'title': f'A borgmatic {state.name} event happened',
'message': f'A borgmatic {state.name} event happened',
'priority': 'default',
'priority': 3,
'tags': 'borgmatic',
}
+884 -6
View File
@@ -19,12 +19,18 @@ from borgmatic import execute as module
(['borg1'], 1, 'borg1', None, module.Exit_status.WARNING),
(['grep'], 100, None, None, module.Exit_status.ERROR),
(['grep'], 100, 'borg', None, module.Exit_status.ERROR),
('grep', 2, None, None, module.Exit_status.ERROR),
('borg', 2, 'borg', None, module.Exit_status.ERROR),
(['borg'], 100, 'borg', None, module.Exit_status.WARNING),
(['borg1'], 100, 'borg1', None, module.Exit_status.WARNING),
('borg', 100, 'borg', None, module.Exit_status.WARNING),
('borg1', 100, 'borg1', None, module.Exit_status.WARNING),
(['grep'], 0, None, None, module.Exit_status.SUCCESS),
(['grep'], 0, 'borg', None, module.Exit_status.SUCCESS),
(['borg'], 0, 'borg', None, module.Exit_status.SUCCESS),
(['borg1'], 0, 'borg1', None, module.Exit_status.SUCCESS),
('grep', 0, None, None, module.Exit_status.SUCCESS),
('grep', 0, 'borg', None, module.Exit_status.SUCCESS),
# -9 exit code occurs when child process get SIGKILLed.
(['grep'], -9, None, None, module.Exit_status.ERROR),
(['grep'], -9, 'borg', None, module.Exit_status.ERROR),
@@ -174,6 +180,23 @@ def test_parse_log_line_with_borg_command_parses_borg_log_line():
)
def test_parse_log_line_with_borg_command_parses_borg_log_line_with_string_command():
record = flexmock()
flexmock(module).should_receive('borg_json_log_line_to_record').and_return(record).once()
flexmock(module).should_receive('log_line_to_record').never()
assert (
module.parse_log_line(
'All done',
module.logging.INFO,
elevate_stderr=False,
borg_local_path='borg',
command='borg do-stuff',
)
== record
)
def test_parse_log_line_without_borg_command_parses_plain_log_line():
record = flexmock()
flexmock(module).should_receive('borg_json_log_line_to_record').never()
@@ -191,6 +214,23 @@ def test_parse_log_line_without_borg_command_parses_plain_log_line():
)
def test_parse_log_line_without_borg_command_parses_plain_log_line_with_string_command():
record = flexmock()
flexmock(module).should_receive('borg_json_log_line_to_record').never()
flexmock(module).should_receive('log_line_to_record').and_return(record).once()
assert (
module.parse_log_line(
'All done',
module.logging.INFO,
elevate_stderr=False,
borg_local_path='borg',
command='totally-not-borg do-stuff',
)
== record
)
def test_parse_log_line_with_elevate_stderr_makes_error_record():
record = flexmock()
flexmock(module).should_receive('borg_json_log_line_to_record').never()
@@ -262,6 +302,844 @@ def test_handle_log_record_over_max_line_count_trims_and_appends():
assert last_lines == [*original_last_lines[1:], 'line']
def test_handle_log_record_without_last_lines_just_handles():
flexmock(module.logger).should_receive('handle').once()
log_record = flexmock(levelno=module.logging.INFO, getMessage=lambda: 'line')
assert module.handle_log_record(log_record) == log_record
def test_log_buffer_lines_without_buffer_readers_bails():
flexmock(module.select).should_receive('select').never()
assert (
tuple(
module.log_buffer_lines(
buffer_readers={},
process_metadatas={},
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_without_ready_buffers_bails():
buffer_readers = {flexmock(): flexmock()}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return([], [], []).once()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas={},
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_ready_buffer_and_running_process_handles_each_log_line():
process = flexmock(poll=lambda: None, stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(lines=iter((('hi', 'there'),)), process=process)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').and_return(flexmock())
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).twice()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_ready_buffer_and_capture_process_yields_each_line():
process = flexmock(poll=lambda: None, stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(lines=iter((('hi', 'there'),)), process=process)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=True)}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').and_return(flexmock())
flexmock(module).should_receive('handle_log_record').and_return(
flexmock(levelno=None, getMessage=lambda: 'message')
).twice()
assert tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
) == ('message', 'message')
def test_log_buffer_lines_with_ready_buffer_and_log_level_and_capture_process_does_not_yield_each_line():
process = flexmock(poll=lambda: None, stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(lines=iter((('hi', 'there'),)), process=process)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=True)}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').and_return(flexmock())
flexmock(module).should_receive('handle_log_record').and_return(
flexmock(levelno=10, getMessage=lambda: 'message')
).twice()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_ready_buffer_and_finished_process_vents_other_processes():
process_stdout = flexmock()
process = flexmock(poll=lambda: 0, stdout=process_stdout, stderr=flexmock(), args=flexmock())
other_process = flexmock(
poll=lambda: None, stdout=flexmock(), stderr=flexmock(), args=flexmock()
)
buffer_readers = {process_stdout: module.Buffer_reader(lines=iter((('hi',),)), process=process)}
process_metadatas = {
process: module.Process_metadata(last_lines=[], capture=False),
other_process: module.Process_metadata(last_lines=[], capture=False),
}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('read_lines').and_return(iter((('there',),))).once()
flexmock(module).should_receive('parse_log_line').and_return(flexmock())
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).once()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
# log_buffer_lines() vents other processes by adding them to buffer_readers, with the idea that
# subsequent calls will then read from them.
assert len(buffer_readers) == 2
# Assert that the process' buffer has been consumed, indicating that it hasn't been accidentally
# replaced.
assert tuple(buffer_readers[process_stdout].lines) == ()
def test_log_buffer_lines_with_ready_buffer_and_finished_process_does_not_vent_other_finished_processes():
process_stdout = flexmock()
process = flexmock(poll=lambda: 0, stdout=process_stdout, stderr=flexmock(), args=flexmock())
other_process = flexmock(poll=lambda: 0, stdout=flexmock(), stderr=flexmock(), args=flexmock())
buffer_readers = {process_stdout: module.Buffer_reader(lines=iter((('hi',),)), process=process)}
process_metadatas = {
process: module.Process_metadata(last_lines=[], capture=False),
other_process: module.Process_metadata(last_lines=[], capture=False),
}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('read_lines').never()
flexmock(module).should_receive('parse_log_line').and_return(flexmock())
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).once()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
assert len(buffer_readers) == 1
assert tuple(buffer_readers[process_stdout].lines) == ()
def test_log_buffer_lines_with_ready_eof_buffer_and_running_process_skips_it():
process = flexmock(poll=lambda: None, stderr=flexmock(), args=flexmock())
buffer_readers = {flexmock(): module.Buffer_reader(lines=iter(()), process=process)}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').never()
flexmock(module).should_receive('handle_log_record').never()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_ready_buffer_with_empty_line_skips_it():
process = flexmock(poll=lambda: None, stderr=flexmock(), args=flexmock())
buffer_readers = {flexmock(): module.Buffer_reader(lines=iter((('',),)), process=process)}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').never()
flexmock(module).should_receive('handle_log_record').never()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_multiple_ready_buffers_and_running_processes_handles_log_lines_from_each():
process = flexmock(poll=lambda: None, stderr=flexmock(), args=flexmock())
other_process = flexmock(poll=lambda: None, stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(lines=iter((('hi', 'there'),)), process=process),
flexmock(): module.Buffer_reader(lines=iter((('foo', 'bar'),)), process=other_process),
}
process_metadatas = {
process: module.Process_metadata(last_lines=[], capture=False),
other_process: module.Process_metadata(last_lines=[], capture=False),
}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').and_return(flexmock())
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).times(4)
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_multiple_ready_buffers_from_same_running_process_handles_all_log_lines():
process = flexmock(poll=lambda: None, stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(lines=iter((('hi', 'there'),)), process=process),
flexmock(): module.Buffer_reader(lines=iter((('foo', 'bar'),)), process=process),
}
process_metadatas = {
process: module.Process_metadata(last_lines=[], capture=False),
}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').and_return(flexmock())
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).times(4)
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_ready_stderr_buffer_and_running_process_elevates_stderr():
process_stderr = flexmock()
process = flexmock(poll=lambda: None, stderr=process_stderr, args=flexmock())
buffer_readers = {
process_stderr: module.Buffer_reader(lines=iter((('hi', 'there'),)), process=process)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=True, borg_local_path=object, command=object
).and_return(flexmock()).twice()
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).twice()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_ready_stdout_buffer_and_running_process_does_not_elevate_stderr():
process_stdout = flexmock()
process = flexmock(poll=lambda: None, stdout=process_stdout, stderr=flexmock(), args=flexmock())
buffer_readers = {
process_stdout: module.Buffer_reader(lines=iter((('hi', 'there'),)), process=process)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=False, borg_local_path=object, command=object
).and_return(flexmock()).twice()
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).twice()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_buffer_lines_with_ready_stderr_buffer_and_capture_stderr_does_not_elevate_stderr():
process_stderr = flexmock()
process = flexmock(poll=lambda: None, stderr=process_stderr, args=flexmock())
buffer_readers = {
process_stderr: module.Buffer_reader(lines=iter((('hi', 'there'),)), process=process)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module.select).should_receive('select').with_args(
buffer_readers.keys(), [], []
).and_return(list(buffer_readers.keys()), [], [])
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=False, borg_local_path=object, command=object
).and_return(flexmock()).twice()
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).twice()
assert (
tuple(
module.log_buffer_lines(
buffer_readers=buffer_readers,
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
capture_stderr=True,
)
)
== ()
)
def test_raise_for_process_errors_with_no_processes_bails():
process = flexmock()
process.should_receive('poll').never()
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
assert (
module.raise_for_process_errors(
buffer_readers,
process_metadatas={},
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
is None
)
def test_raise_for_process_errors_with_running_process_bails():
process = flexmock(poll=lambda: None)
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
assert (
module.raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
is None
)
def test_raise_for_process_errors_with_running_process_and_no_buffer_readers_waits_and_bails():
process = flexmock()
process.should_receive('poll').never()
process.should_receive('wait').and_return(None).once()
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
assert (
module.raise_for_process_errors(
buffer_readers={},
process_metadatas=process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
is None
)
def test_raise_for_process_errors_with_successful_process_bails():
process = flexmock(poll=lambda: 0, args=flexmock())
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module).should_receive('interpret_exit_code').and_return(module.Exit_status.SUCCESS)
assert (
module.raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
is None
)
def test_raise_for_process_errors_with_warning_process_returns_warning_status():
process = flexmock(poll=lambda: 1, args=flexmock())
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module).should_receive('interpret_exit_code').and_return(module.Exit_status.WARNING)
assert (
module.raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
== module.Exit_status.WARNING
)
def test_raise_for_process_errors_with_error_process_raises():
process = flexmock(poll=lambda: 3, args=flexmock())
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {
process: module.Process_metadata(last_lines=['hi', 'there'], capture=False)
}
flexmock(module).should_receive('interpret_exit_code').and_return(module.Exit_status.ERROR)
command = flexmock()
flexmock(module).should_receive('command_for_process').and_return(command)
with pytest.raises(module.subprocess.CalledProcessError) as error:
module.raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
assert error.value.returncode == 3
assert error.value.cmd == command
assert error.value.output == 'hi\nthere'
def test_raise_for_process_errors_with_success_process_and_warning_process_returns_warning_status():
process = flexmock(poll=lambda: 0, args=flexmock())
other_process = flexmock(poll=lambda: 1, args=flexmock())
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {
process: module.Process_metadata(last_lines=['hi', 'there'], capture=False),
other_process: module.Process_metadata(last_lines=['and', 'stuff'], capture=False),
}
flexmock(module).should_receive('interpret_exit_code').with_args(
object, 0, object, object
).and_return(module.Exit_status.SUCCESS)
flexmock(module).should_receive('interpret_exit_code').with_args(
object, 1, object, object
).and_return(module.Exit_status.WARNING)
flexmock(module).should_receive('command_for_process').never()
assert (
module.raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
== module.Exit_status.WARNING
)
def test_raise_for_process_errors_with_warning_process_and_error_process_raises():
process = flexmock(poll=lambda: 1, args=flexmock())
other_process = flexmock(poll=lambda: 3, args=flexmock())
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {
process: module.Process_metadata(last_lines=['hi', 'there'], capture=False),
other_process: module.Process_metadata(last_lines=['and', 'stuff'], capture=False),
}
flexmock(module).should_receive('interpret_exit_code').with_args(
object, 1, object, object
).and_return(module.Exit_status.WARNING)
flexmock(module).should_receive('interpret_exit_code').with_args(
object, 3, object, object
).and_return(module.Exit_status.ERROR)
command = flexmock()
flexmock(module).should_receive('command_for_process').and_return(command)
with pytest.raises(module.subprocess.CalledProcessError) as error:
module.raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
assert error.value.returncode == 3
assert error.value.cmd == command
assert error.value.output == 'and\nstuff'
def test_raise_for_process_errors_with_warning_process_and_running_process_kills_and_returns_warning_status():
process = flexmock(poll=lambda: 1, args=flexmock())
other_process = flexmock(
poll=lambda: None, stdout=flexmock(read=lambda size: None), args=flexmock()
)
other_process.should_receive('kill').once()
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {
process: module.Process_metadata(last_lines=['hi', 'there'], capture=False),
other_process: module.Process_metadata(last_lines=['and', 'stuff'], capture=False),
}
flexmock(module).should_receive('interpret_exit_code').with_args(
object, 1, object, object
).and_return(module.Exit_status.WARNING)
command = flexmock()
flexmock(module).should_receive('command_for_process').and_return(command)
assert (
module.raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
== module.Exit_status.WARNING
)
def test_raise_for_process_errors_with_error_process_and_running_process_kills_and_raises():
process = flexmock(poll=lambda: 3, args=flexmock())
other_process = flexmock(
poll=lambda: None, stdout=flexmock(read=lambda size: None), args=flexmock()
)
other_process.should_receive('kill').once()
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {
process: module.Process_metadata(last_lines=['hi', 'there'], capture=False),
other_process: module.Process_metadata(last_lines=['and', 'stuff'], capture=False),
}
flexmock(module).should_receive('interpret_exit_code').with_args(
object, 3, object, object
).and_return(module.Exit_status.ERROR)
command = flexmock()
flexmock(module).should_receive('command_for_process').and_return(command)
with pytest.raises(module.subprocess.CalledProcessError) as error:
module.raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
assert error.value.returncode == 3
assert error.value.cmd == command
assert error.value.output == 'hi\nthere'
def test_raise_for_process_errors_with_warning_process_and_long_output_raises_with_truncated_output():
process = flexmock(poll=lambda: 3, args=flexmock())
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=process)}
process_metadatas = {
process: module.Process_metadata(last_lines=['hi', 'there'], capture=False)
}
flexmock(module).should_receive('interpret_exit_code').and_return(module.Exit_status.ERROR)
command = flexmock()
flexmock(module).should_receive('command_for_process').and_return(command)
with pytest.raises(module.subprocess.CalledProcessError) as error:
flexmock(module, ERROR_OUTPUT_MAX_LINE_COUNT=2).raise_for_process_errors(
buffer_readers,
process_metadatas,
borg_local_path=flexmock(),
borg_exit_codes=flexmock(),
)
assert error.value.returncode == 3
assert error.value.cmd == command
assert error.value.output == '...\nhi\nthere'
def test_log_remaining_buffer_lines_without_buffer_readers_bails():
process_metadatas = {flexmock(): module.Process_metadata(last_lines=[], capture=False)}
flexmock(module).should_receive('parse_log_line').never()
assert (
tuple(
module.log_remaining_buffer_lines(
buffer_readers={},
process_metadatas=process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_remaining_buffer_lines_without_reader_process_bails():
buffer_readers = {flexmock(): module.Buffer_reader(lines=flexmock(), process=None)}
process_metadatas = {flexmock(): module.Process_metadata(last_lines=[], capture=False)}
flexmock(module).should_receive('parse_log_line').never()
assert (
tuple(
module.log_remaining_buffer_lines(
buffer_readers,
process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_remaining_buffer_lines_logs_each_line():
process = flexmock(stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(
lines=(('hi', 'there'), ('and', 'stuff')),
process=process,
)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=False, borg_local_path=object, command=object
).and_return(flexmock()).times(4)
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).times(4)
assert (
tuple(
module.log_remaining_buffer_lines(
buffer_readers,
process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_remaining_buffer_lines_with_multiple_buffers_logs_lines_from_each():
process = flexmock(stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(
lines=(('hi', 'there'),),
process=process,
),
flexmock(): module.Buffer_reader(
lines=(('and', 'stuff'),),
process=process,
),
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=False, borg_local_path=object, command=object
).and_return(flexmock()).times(4)
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).times(4)
assert (
tuple(
module.log_remaining_buffer_lines(
buffer_readers,
process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_remaining_buffer_lines_with_stderr_buffer_elevates_stderr():
stderr = flexmock()
process = flexmock(stderr=stderr, args=flexmock())
buffer_readers = {
stderr: module.Buffer_reader(
lines=(('hi', 'there'),),
process=process,
),
flexmock(): module.Buffer_reader(
lines=(('and', 'stuff'),),
process=process,
),
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=True, borg_local_path=object, command=object
).and_return(flexmock()).twice()
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=False, borg_local_path=object, command=object
).and_return(flexmock()).twice()
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).times(4)
assert (
tuple(
module.log_remaining_buffer_lines(
buffer_readers,
process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_log_remaining_buffer_lines_with_stderr_buffer_and_capture_stderr_does_not_elevate_stderr():
stderr = flexmock()
process = flexmock(stderr=stderr, args=flexmock())
buffer_readers = {
stderr: module.Buffer_reader(
lines=(('hi', 'there'),),
process=process,
),
flexmock(): module.Buffer_reader(
lines=(('and', 'stuff'),),
process=process,
),
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=False)}
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=False, borg_local_path=object, command=object
).and_return(flexmock()).times(4)
flexmock(module).should_receive('handle_log_record').and_return(flexmock(levelno=10)).times(4)
assert (
tuple(
module.log_remaining_buffer_lines(
buffer_readers,
process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
capture_stderr=True,
)
)
== ()
)
def test_log_remaining_buffer_lines_with_capture_process_yields_each_line():
process = flexmock(stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(
lines=(('hi', 'there'),),
process=process,
)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=True)}
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=False, borg_local_path=object, command=object
).and_return(flexmock()).twice()
flexmock(module).should_receive('handle_log_record').and_return(
flexmock(levelno=None, getMessage=lambda: 'message')
).twice()
assert tuple(
module.log_remaining_buffer_lines(
buffer_readers,
process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
) == ('message', 'message')
def test_log_remaining_buffer_lines_with_log_level_and_capture_process_does_not_yield_each_line():
process = flexmock(stderr=flexmock(), args=flexmock())
buffer_readers = {
flexmock(): module.Buffer_reader(
lines=(('hi', 'there'),),
process=process,
)
}
process_metadatas = {process: module.Process_metadata(last_lines=[], capture=True)}
flexmock(module).should_receive('parse_log_line').with_args(
line=str, log_level=object, elevate_stderr=False, borg_local_path=object, command=object
).and_return(flexmock()).twice()
flexmock(module).should_receive('handle_log_record').and_return(
flexmock(levelno=10, getMessage=lambda: 'message')
).twice()
assert (
tuple(
module.log_remaining_buffer_lines(
buffer_readers,
process_metadatas,
output_log_level=flexmock(),
borg_local_path=flexmock(),
)
)
== ()
)
def test_mask_command_secrets_masks_password_flag_value():
assert module.mask_command_secrets(('cooldb', '--username', 'bob', '--password', 'pass')) == (
'cooldb',
@@ -507,7 +1385,7 @@ def test_execute_command_and_capture_output_returns_stdout():
full_command,
stdin=None,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
stderr=None,
shell=False,
env=None,
cwd=None,
@@ -552,7 +1430,7 @@ def test_execute_command_and_capture_output_returns_output_when_process_error_is
full_command,
stdin=None,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
stderr=None,
shell=False,
env=None,
cwd=None,
@@ -574,7 +1452,7 @@ def test_execute_command_and_capture_output_raises_when_command_errors():
full_command,
stdin=None,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
stderr=None,
shell=False,
env=None,
cwd=None,
@@ -596,7 +1474,7 @@ def test_execute_command_and_capture_output_with_shell_returns_output():
'foo bar',
stdin=None,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
stderr=None,
shell=True,
env=None,
cwd=None,
@@ -618,7 +1496,7 @@ def test_execute_command_and_capture_output_with_enviroment_returns_output():
full_command,
stdin=None,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
stderr=None,
shell=False,
env={'a': 'b'},
cwd=None,
@@ -646,7 +1524,7 @@ def test_execute_command_and_capture_output_returns_output_with_working_director
full_command,
stdin=None,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
stderr=None,
shell=False,
env=None,
cwd='/working',