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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
219 changes: 193 additions & 26 deletions nisshi-smoke-test/src/broker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,12 +30,13 @@
use std::{
fmt,
fs::{self, File},
io,
net::TcpListener,
path::{Path, PathBuf},
process::{Child, Command, Stdio},
sync::atomic::{AtomicU16, Ordering},
thread,
time::Duration,
time::{Duration, SystemTime},
};

use crate::{KafkaCli, ScramLogin, StorageUrl, label, settings, timed_command, unique_name};
Expand All @@ -53,6 +54,13 @@ const CONTAINER_START_TIMEOUT: Duration = Duration::from_secs(60);
/// the broker running; `run.sh` stops every broker whose file is still there when the run ends.
const PID_FILE: &str = "broker.pid";

/// The storage URL option that makes a broker run maintenance every 2 seconds instead of every 10
/// minutes, for a test that waits for a maintenance run.
pub const FREQUENT_MAINTENANCE: &str = "maintenance_interval=2s";

/// What `docker cp` prints when the file to copy isn't in the container.
const DOCKER_CP_MISSING_FILE: &str = "Could not find the file";

/// A broker the tests talk to: the shared one `just smoke` started, or one
/// of their own from [`Broker::isolated`].
#[derive(Debug)]
Expand Down Expand Up @@ -132,6 +140,14 @@ impl Broker {
Self::launch(LaunchOptions::new(settings::storage_url_under_test()))
}

/// Launches a broker of this test's own, as [`Broker::isolated`] does, that runs maintenance
/// every 2 seconds instead of every 10 minutes, for a test that waits for a maintenance run.
pub fn isolated_with_frequent_maintenance() -> Self {
Self::launch(LaunchOptions::new(
settings::storage_url_under_test().with_query_option(FREQUENT_MAINTENANCE),
))
}

/// Launches `NISSHI_SMOKE_IMAGE` when set, otherwise `NISSHI_SMOKE_BIN`, and waits until it
/// answers. Panics if it doesn't start, or if something else already listens on the port.
pub fn launch(options: LaunchOptions) -> Self {
Expand Down Expand Up @@ -172,30 +188,54 @@ impl Broker {
/// Stops the broker as [`Broker::stop`] does, and starts it again on the same storage, on a
/// new port. Panics if it fails its checks or doesn't start again.
pub fn restart(self) -> Self {
self.restart_with(|_| {})
self.restart_with(|_| {}, &[] as &[&str])
}

/// Like [`Broker::restart`], but the broker starts again with `--authentication`, so every
/// client must log in. The harness's readiness check logs in as `login`, so the broker starts
/// again only if it still has that user's credentials.
pub fn restart_requiring_login(self, login: &ScramLogin) -> Self {
self.restart_with(|options| options.login = Some(login.clone()))
self.restart_with(
|options| options.login = Some(login.clone()),
&[] as &[&str],
)
}

/// Like [`Broker::restart`], but the broker starts again on `storage`, with the same cluster
/// id and working directory. A relative SQLite path in `storage` therefore names a file next
/// to the broker's own database.
///
/// Each of `files_to_delete`, a path relative to the working directory, is deleted while the
/// broker is stopped, so a test can make sure that the broker can't read them after the
/// restart. A file that doesn't exist is skipped.
pub fn restart_on_storage(
self,
storage: StorageUrl,
files_to_delete: &[impl AsRef<str>],
) -> Self {
self.restart_with(|options| options.storage = storage, files_to_delete)
}

fn restart_with(mut self, change_options: impl FnOnce(&mut LaunchOptions)) -> Self {
fn restart_with(
mut self,
change_options: impl FnOnce(&mut LaunchOptions),
files_to_delete: &[impl AsRef<str>],
) -> Self {
let deployment = self
.deployment
.take()
.expect("only a broker this test launched can restart");

let (stopped, files, mut options) = deployment.stop();
let (stopped, mut files, mut options) = deployment.stop();

if let Err(reason) = stopped {
panic!("{reason}");
}

files.delete(files_to_delete);
options.port = free_port();
change_options(&mut options);
files.create_volume_for(&options.storage);
let cluster_id = options.cluster_id.clone();

Self::start_with_files(options, files).unwrap_or_else(|failed_start| {
Expand Down Expand Up @@ -285,21 +325,33 @@ impl Broker {
&self.bootstrap
}

/// Returns whether `path` exists where the broker runs: on this machine for a broker process,
/// and in its container for a broker container. Panics for the shared broker.
/// Returns whether the file at `path` exists where the broker runs: on this machine for a
/// broker process, and in its container for a broker container. A relative `path` starts at
/// the broker's working directory.
///
/// Panics for the shared broker. `smoke-broker` launches that broker, so this `Broker` doesn't
/// know where the broker runs.
pub fn file_exists(&self, path: &str) -> bool {
let deployment = self
.deployment
.as_ref()
.expect("only a broker this test launched has files the test can check");
self.launched_deployment().file_metadata(path).is_some()
}

match &deployment.host {
Host::Process(_) => Path::new(path).exists(),
/// Returns when the file at `path`, where the broker runs, was last written, or `None` if the
/// file doesn't exist. A relative `path` starts at the broker's working directory.
///
/// Panics for the shared broker. `smoke-broker` launches that broker, so this `Broker` doesn't
/// know where the broker runs.
pub fn file_modification_time(&self, path: &str) -> Option<SystemTime> {
self.launched_deployment()
.file_metadata(path)?
.modified()
.ok()
}

Host::Container { name } => {
docker(&["exec", name, "test", "-e", path]).is_ok_and(|test| test.code == Some(0))
}
}
/// Returns the broker that this `Broker` launched. Panics for the shared broker.
fn launched_deployment(&self) -> &Deployment {
self.deployment
.as_ref()
.expect("only a broker this test launched has files the test can check")
}

/// Stops a broker this test launched, failing if it exited before now, didn't exit with 0 on
Expand Down Expand Up @@ -392,6 +444,49 @@ impl fmt::Display for FailedStart {
}

impl Deployment {
/// Returns the metadata of the file at `path` where the broker runs, or `None` if the file
/// doesn't exist. A relative `path` starts at the broker's working directory.
///
/// Panics if the harness can't read the file for another reason, so that a test never takes
/// such a failure for a missing file.
fn file_metadata(&self, path: &str) -> Option<fs::Metadata> {
match &self.host {
Host::Process(_) => match fs::metadata(self.files.dir.join(path)) {
Ok(metadata) => Some(metadata),
Err(err) if err.kind() == io::ErrorKind::NotFound => None,
Err(err) => panic!("{path}: {err}"),
},

Host::Container { name } => {
// The broker's image holds only the broker, without a shell or `stat`, so the
// harness copies the file out with `docker cp`, which keeps its modification time.
// `docker cp` reads a relative path from the container's root, not from its
// working directory.
let working_dir = inspect(name, "{{.Config.WorkingDir}}")
.unwrap_or_else(|reason| panic!("{reason}"));
let path_in_container = Path::new("/").join(working_dir).join(path);

let copy = self.files.dir.join(unique_name("copied"));
let source = format!("{name}:{}", path_in_container.display());

let copied = docker(&["cp", &source, &copy.display().to_string()])
.unwrap_or_else(|reason| panic!("{reason}"));

let metadata = match copied.code {
Some(0) => Some(
fs::metadata(&copy)
.unwrap_or_else(|err| panic!("{}: {err}", copy.display())),
),
_ if copied.stderr.contains(DOCKER_CP_MISSING_FILE) => None,
_ => panic!("docker cp {source}: {}", copied.stderr.trim()),
};

_ = fs::remove_file(&copy);
metadata
}
}
}

/// Stops the broker and checks it. Returns the files and options too, for a restart.
fn stop(self) -> (Result<(), String>, BrokerFiles, LaunchOptions) {
let mut failures = Vec::new();
Expand Down Expand Up @@ -511,19 +606,90 @@ impl BrokerFiles {
fs::create_dir_all(&dir)
.unwrap_or_else(|err| panic!("broker directory {}: {err}", dir.display()));

let volume = (settings::broker_image().is_some() && storage.is_sqlite()).then(|| {
let mut files = Self {
dir,
volume: None,
starts: 0,
keep_dir: false,
};
files.create_volume_for(storage);

files
}

/// Creates a volume for a SQLite broker in a container, if the files don't have one yet.
///
/// The volume is the container's working directory, where the broker resolves a relative
/// `sqlite://` path. A database outside the volume is removed with the container, at the next
/// stop.
fn create_volume_for(&mut self, storage: &StorageUrl) {
if self.volume.is_none() && settings::broker_image().is_some() && storage.is_sqlite() {
let volume = unique_name("nisshi-smoke-sqlite");

_ = docker(&["volume", "create", &label(), &volume]);

volume
});
self.volume = Some(volume);
}
}

Self {
dir,
volume,
starts: 0,
keep_dir: false,
/// Deletes each of `paths`, relative to the broker's working directory, while the broker is
/// stopped. Skips a path that doesn't exist, and panics if a deletion fails.
///
/// A broker container without a volume keeps its files in the container, which the harness
/// removed when the broker stopped, so there is nothing to delete.
fn delete(&self, paths: &[impl AsRef<str>]) {
if paths.is_empty() {
return;
}

match (&self.volume, settings::broker_image()) {
(_, None) => {
for path in paths {
let path = path.as_ref();
match fs::remove_file(self.dir.join(path)) {
Ok(()) => {}
Err(err) if err.kind() == io::ErrorKind::NotFound => {}
Err(err) => panic!("{path}: {err}"),
}
}
}

(Some(volume), Some(_)) => {
// The broker's image holds only the broker, without `rm`, so a container from the
// Kafka tools' image deletes the files from the volume.
let tools_image = inspect(&settings::kafka_container(), "{{.Config.Image}}")
.unwrap_or_else(|reason| panic!("{reason}"));
let mount = format!("{volume}:/data");
let label = label();
let mut args = vec![
"run",
"--rm",
&label,
// The broker writes the volume as root. The tools' image runs as a user that
// can't delete root's files.
"--user=0",
"--volume",
&mount,
"--entrypoint",
"rm",
&tools_image,
"-f",
];
let paths_in_volume = paths
.iter()
.map(|path| format!("/data/{}", path.as_ref()))
.collect::<Vec<_>>();
args.extend(paths_in_volume.iter().map(String::as_str));

match docker(&args) {
Ok(deleted) if deleted.code == Some(0) => {}
deleted => {
panic!("couldn't delete {paths_in_volume:?} from {volume}: {deleted:?}")
}
}
}

(None, Some(_)) => {}
}
}

Expand Down Expand Up @@ -757,7 +923,8 @@ fn save_logs(name: &str, log: &Path) -> Result<(), String> {
.map_err(|err| format!("couldn't write {}: {err}", log.display()))
}

/// Removes a broker container, but not its volume, which [`BrokerFiles`] removes.
/// Removes a container that the harness started, but not a broker's volume, which [`BrokerFiles`]
/// removes.
fn remove_container(name: &str) {
_ = docker(&["rm", "--force", "--volumes", name]);
}
Expand Down
Loading
Loading