diff --git a/scripts/omega b/scripts/omega index 15954614..a6c11785 100755 --- a/scripts/omega +++ b/scripts/omega @@ -850,36 +850,93 @@ options() { fi } +remove_partly_copied_memory() { + local volume="$1" + local container="$2" + local label="$3" + local leftover + local holder + local holders="" + + if ! docker volume inspect "${volume}" >/dev/null 2>&1; then + return 0 + fi + for leftover in $(docker ps -aq --filter "label=${label}" --filter "volume=${volume}" 2>/dev/null); do + docker rm -f "${leftover}" >/dev/null 2>&1 || true + done + for holder in $(docker ps -a --filter "volume=${volume}" --format '{{.Names}}' 2>/dev/null); do + if [ "${holder}" != "${container}" ]; then + holders="${holders:+${holders} }${holder}" + fi + done + if [ -n "${holders}" ]; then + echo "Could not remove the partly copied volume ${volume}, it is used by containers: ${holders}. Remove them with 'docker rm -f ${holders}' and start again" >&2 + return 1 + fi + if docker container inspect "${container}" >/dev/null 2>&1 \ + && docker rm -f "${container}" >/dev/null; then + echo "Removed container ${container}, it used the partly copied volume ${volume}" + fi + if ! docker volume rm "${volume}" >/dev/null; then + echo "Could not remove the partly copied volume ${volume}" >&2 + return 1 + fi +} + migrate_omegaclaw_memory() { local old_container="omegaclaw" local old_volume="omegaclaw-memory" local new_volume="omega-memory" + local new_container="omega" + local migration_label="omega.memory-migration" local memory_path="/PeTTa/repos/Omega/memory" local marker=".migrated-to-omega" local started_marker=".migration-started" local stopped_old_container=0 local old_state + local new_state - if [ -n "${memory_import_file}" ] \ - || ! docker volume inspect "${old_volume}" >/dev/null 2>&1; then + if ! docker volume inspect "${old_volume}" >/dev/null 2>&1; then return 0 fi - if ! old_state=$(docker run --rm --entrypoint sh --volume "${old_volume}:/from:ro" "${image}" \ + if ! old_state=$(docker run --rm --label "${migration_label}" --entrypoint sh --volume "${old_volume}:/from:ro" "${image}" \ -c "if [ -e /from/${marker} ]; then echo migrated; elif [ -e /from/${started_marker} ]; then echo interrupted; fi"); then echo "Could not read the migration markers in ${old_volume}" >&2 return 1 fi + if [ -n "${memory_import_file}" ]; then + if [ "${old_state}" = "interrupted" ]; then + echo "The last copy from ${old_volume} did not finish, the memory import replaces it" + remove_partly_copied_memory "${new_volume}" "${new_container}" "${migration_label}" || return 1 + if ! docker run --rm --label "${migration_label}" --entrypoint sh --volume "${old_volume}:/from" "${image}" \ + -c "rm -f /from/${started_marker}"; then + echo "Could not clear the migration marker in ${old_volume}" >&2 + return 1 + fi + fi + return 0 + fi case "${old_state}" in migrated) return 0 ;; interrupted) - echo "The last copy from ${old_volume} did not finish, copying again" - if docker volume inspect "${new_volume}" >/dev/null 2>&1 \ - && ! docker volume rm "${new_volume}" >/dev/null; then - echo "Could not remove the partly copied volume ${new_volume}" >&2 - return 1 + if docker volume inspect "${new_volume}" >/dev/null 2>&1; then + if ! new_state=$(docker run --rm --label "${migration_label}" --entrypoint sh \ + --volume "${old_volume}:/from:ro" --volume "${new_volume}:/to:ro" "${image}" \ + -c 'if [ -e /to/history.metta ]; then n=$(wc -c < /to/history.metta); if [ ! -e /from/history.metta ] || [ "$n" -gt "$(wc -c < /from/history.metta)" ] || ! head -c "$n" /from/history.metta | cmp -s - /to/history.metta; then echo kept; fi; fi'); then + echo "Could not compare ${new_volume} with ${old_volume}" >&2 + return 1 + fi + if [ "${new_state}" = "kept" ]; then + echo "${new_volume} holds memory that is not a partial copy of ${old_volume}, for example an imported archive, so it was left as it was" >&2 + echo "To keep it, clear the marker: docker run --rm --entrypoint rm --volume ${old_volume}:/from ${image} /from/${started_marker}" >&2 + echo "To copy ${old_volume} again instead, remove it: docker rm -f ${new_container}; docker volume rm ${new_volume}" >&2 + return 1 + fi fi + echo "The last copy from ${old_volume} did not finish, copying again" + remove_partly_copied_memory "${new_volume}" "${new_container}" "${migration_label}" || return 1 ;; *) if docker volume inspect "${new_volume}" >/dev/null 2>&1; then @@ -898,15 +955,15 @@ migrate_omegaclaw_memory() { echo "Stopped container ${old_container}" fi - if docker run --rm --entrypoint sh --volume "${old_volume}:/from" "${image}" \ + if docker run --rm --label "${migration_label}" --entrypoint sh --volume "${old_volume}:/from" "${image}" \ -c "date -u > /from/${started_marker}" \ && docker volume create "${new_volume}" >/dev/null \ - && docker run --rm --entrypoint sh --volume "${new_volume}:${memory_path}" "${image}" \ + && docker run --rm --label "${migration_label}" --entrypoint sh --volume "${new_volume}:${memory_path}" "${image}" \ -c "test -f ${memory_path}/prompt.txt" \ - && docker run --rm --entrypoint sh \ + && docker run --rm --label "${migration_label}" --entrypoint sh \ --volume "${old_volume}:/from:ro" --volume "${new_volume}:/to" "${image}" \ -c 'list=$(mktemp) && cd /from && find . -mindepth 1 -maxdepth 1 ! -name "prompt*.txt" ! -name tg_prompt.txt ! -name .migration-started > "$list" && while IFS= read -r entry; do rm -rf "/to/$entry" && cp -a "$entry" /to/ || exit 1; done < "$list"' \ - && docker run --rm --entrypoint sh --volume "${old_volume}:/from" "${image}" \ + && docker run --rm --label "${migration_label}" --entrypoint sh --volume "${old_volume}:/from" "${image}" \ -c "date -u > /from/${marker}"; then echo "Memory copied. Carrying memory over is experimental, so keep the volume ${old_volume} until you have checked that the agent remembers what it knew. Then remove the old installation with 'docker rm ${old_container}' and 'docker volume rm ${old_volume}'" return 0 @@ -915,7 +972,7 @@ migrate_omegaclaw_memory() { echo "Could not copy memory from ${old_volume}, the old installation is left as it was" >&2 if docker volume rm "${new_volume}" >/dev/null 2>&1 \ || ! docker volume inspect "${new_volume}" >/dev/null 2>&1; then - docker run --rm --entrypoint sh --volume "${old_volume}:/from" "${image}" \ + docker run --rm --label "${migration_label}" --entrypoint sh --volume "${old_volume}:/from" "${image}" \ -c "rm -f /from/${started_marker}" >/dev/null 2>&1 || true else echo "Could not remove ${new_volume}, the next start will copy the memory again" >&2 diff --git a/tests/test_omega_launcher_migration.py b/tests/test_omega_launcher_migration.py index c67cd641..b02163e3 100644 --- a/tests/test_omega_launcher_migration.py +++ b/tests/test_omega_launcher_migration.py @@ -10,6 +10,7 @@ REPO_ROOT = Path(__file__).resolve().parents[1] LAUNCHER = REPO_ROOT / "scripts" / "omega" IMAGE = "test/omega:new" +MIGRATION_LABEL = "omega.memory-migration" FAKE_DOCKER = r'''#!/usr/bin/env python3 import os @@ -23,6 +24,8 @@ root = Path(os.environ["FAKE_DOCKER_ROOT"]) volumes = root / "volumes" containers = root / "containers" +mounts_dir = root / "mounts" +labels_dir = root / "labels" args = sys.argv[1:] with open(root / "calls.log", "a", encoding="utf-8") as log: @@ -33,6 +36,13 @@ def volume(name): return volumes / name +def containers_using(name): + return sorted( + path.name for path in mounts_dir.iterdir() + if name in path.read_text().split() + ) + + if args[:2] == ["volume", "inspect"]: sys.exit(0 if volume(args[2]).is_dir() else 1) if args[:2] == ["volume", "create"]: @@ -42,8 +52,35 @@ def volume(name): if args[:2] == ["volume", "rm"]: if os.environ.get("FAKE_DOCKER_FAIL_VOLUME_RM") or not volume(args[2]).is_dir(): sys.exit(1) + users = containers_using(args[2]) + if users: + ids = ", ".join(format(abs(hash(user)), "x") for user in users) + print(f"Error response from daemon: remove {args[2]}: volume is in use - [{ids}]", file=sys.stderr) + sys.exit(1) shutil.rmtree(volume(args[2])) sys.exit(0) +if args[:1] == ["rm"]: + name = args[-1] + if not (containers / name).exists(): + print(f"Error response from daemon: No such container: {name}", file=sys.stderr) + sys.exit(0 if "-f" in args else 1) + (containers / name).unlink() + (mounts_dir / name).unlink(missing_ok=True) + (labels_dir / name).unlink(missing_ok=True) + sys.exit(0) +if args[:1] == ["ps"]: + filters = [args[index + 1].split("=", 1) for index, option in enumerate(args) if option == "--filter"] + names = sorted(path.name for path in containers.iterdir()) + for key, value in filters: + if key == "volume": + names = [user for user in containers_using(value) if user in names] + elif key == "label": + names = [name for name in names + if (labels_dir / name).exists() and value in (labels_dir / name).read_text().split()] + print("\n".join(names)) + sys.exit(0) +if args[:2] == ["container", "inspect"]: + sys.exit(0 if (containers / args[-1]).exists() else 1) if args[:1] == ["inspect"]: state = containers / args[-1] if not state.exists(): @@ -59,6 +96,7 @@ def volume(name): if args[:1] == ["run"]: mounts = {} entrypoint = None + container_name = None index = 1 while index < len(args): option = args[index] @@ -66,9 +104,11 @@ def volume(name): source, target = args[index + 1].split(":")[:2] mounts[target] = source index += 2 - elif option in ("--entrypoint", "--user", "--group-add", "-e", "--name", "--security-opt", "--tmpfs"): + elif option in ("--entrypoint", "--user", "--group-add", "-e", "--name", "--security-opt", "--tmpfs", "--label"): if option == "--entrypoint": entrypoint = args[index + 1] + if option == "--name": + container_name = args[index + 1] index += 2 elif option.startswith("-"): index += 1 @@ -86,6 +126,11 @@ def volume(name): volume(source).mkdir(parents=True, exist_ok=True) if target == MEMORY_PATH and not any(volume(source).iterdir()): shutil.copytree(root / "image-memory", volume(source), dirs_exist_ok=True) + if container_name: + (containers / container_name).write_text("running") + (mounts_dir / container_name).write_text( + "\n".join(source for source in mounts.values() if "/" not in source) + ) if entrypoint == "sh" and command[:1] == ["-c"]: script = command[1] for target in sorted(mounts, key=len, reverse=True): @@ -117,6 +162,8 @@ def docker_root(tmp_path): root = tmp_path / "docker" (root / "volumes").mkdir(parents=True) (root / "containers").mkdir() + (root / "mounts").mkdir() + (root / "labels").mkdir() image_memory = root / "image-memory" (image_memory / "chroma_db").mkdir(parents=True) (image_memory / "history.metta").write_text("") @@ -274,8 +321,7 @@ def test_without_old_volume_agent_starts_on_image_memory(docker_root): assert not (docker_root / "volumes" / "omegaclaw-memory").exists() -def test_memory_import_skips_migration(docker_root, tmp_path): - _install_omegaclaw(docker_root) +def _memory_import_arguments(docker_root, tmp_path): transfer = tmp_path / "transfer" transfer.mkdir() transfer.chmod(0o2770) @@ -290,9 +336,7 @@ def test_memory_import_skips_migration(docker_root, tmp_path): encoding="utf-8", ) python3.chmod(0o755) - - result = _launcher( - docker_root, + return ( "start", "-d", IMAGE, @@ -304,6 +348,12 @@ def test_memory_import_skips_migration(docker_root, tmp_path): "memory.tar.gz", ) + +def test_memory_import_skips_migration(docker_root, tmp_path): + _install_omegaclaw(docker_root) + + result = _launcher(docker_root, *_memory_import_arguments(docker_root, tmp_path)) + assert result.returncode == 0, result.stderr assert _read(docker_root / "volumes" / "omega-memory" / "history.metta") == "" assert _container_state(docker_root, "omegaclaw") == "running" @@ -383,6 +433,126 @@ def test_interrupted_copy_is_redone(docker_root): assert _started_agent(docker_root) +def _container_on_volume(root, name, state, volume_names, label=None): + (root / "containers" / name).write_text(state) + (root / "mounts" / name).write_text(volume_names) + if label: + (root / "labels" / name).write_text(label) + + +@pytest.mark.parametrize("state", ["exited", "running"]) +def test_interrupted_copy_is_redone_while_omega_container_holds_new_volume(docker_root, state): + old, new = _interrupted_migration(docker_root, image_files=True) + _container_on_volume(docker_root, "omega", state, "omega-memory") + + result = _launcher(docker_root, "start", "-d", IMAGE) + + assert result.returncode == 0, result.stderr + assert "volume is in use" not in result.stderr + assert "Removed container omega, it used the partly copied volume omega-memory" in result.stdout + assert _read(new / "history.metta") == "(old history)\n" + assert _read(new / "chroma_db" / "chroma.sqlite3") == "old long-term memory" + assert (old / ".migrated-to-omega").exists() + assert _started_agent(docker_root) + + +def test_interrupted_copy_names_other_containers_on_new_volume(docker_root): + old, new = _interrupted_migration(docker_root, image_files=True) + _container_on_volume(docker_root, "omega-leftover", "exited", "omega-memory") + before = _snapshot(new) + + result = _launcher(docker_root, "start", "-d", IMAGE) + + assert result.returncode != 0 + assert "omega-leftover" in result.stderr + assert _snapshot(new) == before + assert (docker_root / "containers" / "omega-leftover").exists() + assert not (old / ".migrated-to-omega").exists() + assert not _started_agent(docker_root) + + +def test_interrupted_copy_keeps_omega_while_another_container_holds_the_volume(docker_root): + old, new = _interrupted_migration(docker_root, image_files=True) + _container_on_volume(docker_root, "omega", "running", "omega-memory") + _container_on_volume(docker_root, "omega-leftover", "exited", "omega-memory") + before = _snapshot(new) + + result = _launcher(docker_root, "start", "-d", IMAGE) + + assert result.returncode != 0 + assert "it is used by containers: omega-leftover." in result.stderr + assert _container_state(docker_root, "omega") == "running" + assert _snapshot(new) == before + assert not _started_agent(docker_root) + + +def test_interrupted_copy_removes_leftover_copy_containers(docker_root): + old, new = _interrupted_migration(docker_root, image_files=True) + _container_on_volume(docker_root, "musing_haslett", "running", "omegaclaw-memory omega-memory", label=MIGRATION_LABEL) + + result = _launcher(docker_root, "start", "-d", IMAGE) + + assert result.returncode == 0, result.stderr + assert _read(new / "history.metta") == "(old history)\n" + assert (old / ".migrated-to-omega").exists() + assert not (docker_root / "containers" / "musing_haslett").exists() + assert _started_agent(docker_root) + + +def test_interrupted_copy_keeps_labelled_containers_on_other_volumes(docker_root): + old, new = _interrupted_migration(docker_root, image_files=True) + _container_on_volume(docker_root, "other_copier", "running", "omegaclaw-memory", label=MIGRATION_LABEL) + + result = _launcher(docker_root, "start", "-d", IMAGE) + + assert result.returncode == 0, result.stderr + assert _read(new / "history.metta") == "(old history)\n" + assert (docker_root / "containers" / "other_copier").exists() + + +def test_migration_containers_carry_the_label(docker_root): + _install_omegaclaw(docker_root) + + result = _launcher(docker_root, "start", "-d", IMAGE) + + assert result.returncode == 0, result.stderr + helper_runs = [call for call in _calls(docker_root) if call.startswith("run --rm")] + assert len(helper_runs) >= 4 + assert all(f"--label {MIGRATION_LABEL}" in call for call in helper_runs) + + +def test_memory_import_replaces_an_interrupted_copy(docker_root, tmp_path): + old, new = _interrupted_migration(docker_root, image_files=True) + (new / "history.metta").write_text("(partly copied)\n") + _container_on_volume(docker_root, "omega", "exited", "omega-memory") + + result = _launcher(docker_root, *_memory_import_arguments(docker_root, tmp_path)) + + assert result.returncode == 0, result.stderr + assert _read(new / "history.metta") == "" + assert not (old / ".migration-started").exists() + assert not (old / ".migrated-to-omega").exists() + assert _started_agent(docker_root) + + (new / "history.metta").write_text("(imported and used)\n") + again = _launcher(docker_root, "start", "-d", IMAGE) + + assert again.returncode == 0, again.stderr + assert _read(new / "history.metta") == "(imported and used)\n" + + +def test_memory_import_keeps_the_interrupted_state_when_the_volume_cannot_be_removed(docker_root, tmp_path): + old, new = _interrupted_migration(docker_root, image_files=True) + (new / "history.metta").write_text("(partly copied)\n") + + result = _launcher(docker_root, *_memory_import_arguments(docker_root, tmp_path), fail_volume_rm=True) + + assert result.returncode != 0 + assert (old / ".migration-started").exists() + assert _read(new / "history.metta") == "(partly copied)\n" + assert not _started_agent(docker_root) + + def test_run_interrupted_before_the_image_files_is_redone(docker_root): old, new = _interrupted_migration(docker_root, image_files=False) @@ -428,3 +598,51 @@ def test_start_stops_when_markers_cannot_be_read(docker_root): assert result.returncode != 0 assert not _started_agent(docker_root) assert _snapshot(new) == before + + +def test_memory_import_stops_when_markers_cannot_be_read(docker_root, tmp_path): + old, new = _interrupted_migration(docker_root, image_files=True) + (new / "history.metta").write_text("(partly copied)\n") + before = _snapshot(new) + + result = _launcher(docker_root, *_memory_import_arguments(docker_root, tmp_path), fail_marker_check=True) + + assert result.returncode != 0 + assert "Could not read the migration markers" in result.stderr + assert (old / ".migration-started").exists() + assert _snapshot(new) == before + assert not _started_agent(docker_root) + + +@pytest.mark.parametrize("history", ["(old hi", None]) +def test_interrupted_copy_with_partly_copied_history_is_redone(docker_root, history): + old, new = _interrupted_migration(docker_root, image_files=True) + if history is None: + (new / "history.metta").unlink() + else: + (new / "history.metta").write_text(history) + + result = _launcher(docker_root, "start", "-d", IMAGE) + + assert result.returncode == 0, result.stderr + assert _read(new / "history.metta") == "(old history)\n" + assert (old / ".migrated-to-omega").exists() + assert _started_agent(docker_root) + + +@pytest.mark.parametrize("history", ["(imported and used)\n", "(old history)\n(used after import)\n"]) +def test_interrupted_copy_keeps_memory_that_is_not_a_partial_copy(docker_root, history): + old, new = _interrupted_migration(docker_root, image_files=True) + (new / "history.metta").write_text(history) + _container_on_volume(docker_root, "omega", "exited", "omega-memory") + before = _snapshot(new) + + result = _launcher(docker_root, "start", "-d", IMAGE) + + assert result.returncode != 0 + assert "docker volume rm omega-memory" in result.stderr + assert ".migration-started" in result.stderr + assert _snapshot(new) == before + assert (old / ".migration-started").exists() + assert _container_state(docker_root, "omega") == "exited" + assert not _started_agent(docker_root)