mirror of
https://github.com/basicmachines-co/basic-memory
synced 2026-06-21 13:47:35 +00:00
feat: Implement cloud mount CLI commands for local file access (#306)
Signed-off-by: phernandez <paul@basicmachines.co> Co-authored-by: Claude <noreply@anthropic.com>
This commit is contained in:
@@ -11,78 +11,25 @@ from rich.table import Table
|
||||
|
||||
from basic_memory.cli.app import cloud_app
|
||||
from basic_memory.cli.auth import CLIAuth
|
||||
from basic_memory.config import ConfigManager
|
||||
from basic_memory.cli.commands.cloud.api_client import (
|
||||
CloudAPIError,
|
||||
get_cloud_config,
|
||||
make_api_request,
|
||||
get_authenticated_headers,
|
||||
)
|
||||
from basic_memory.cli.commands.cloud.mount_commands import (
|
||||
mount_cloud_files,
|
||||
setup_cloud_mount,
|
||||
show_mount_status,
|
||||
unmount_cloud_files,
|
||||
)
|
||||
from basic_memory.cli.commands.cloud.rclone_config import MOUNT_PROFILES
|
||||
from basic_memory.ignore_utils import load_gitignore_patterns, should_ignore_path
|
||||
from basic_memory.utils import generate_permalink
|
||||
|
||||
console = Console()
|
||||
|
||||
|
||||
class CloudAPIError(Exception):
|
||||
"""Exception raised for cloud API errors."""
|
||||
|
||||
pass
|
||||
|
||||
|
||||
def get_cloud_config() -> tuple[str, str, str]:
|
||||
"""Get cloud OAuth configuration from config."""
|
||||
config_manager = ConfigManager()
|
||||
config = config_manager.config
|
||||
return config.cloud_client_id, config.cloud_domain, config.cloud_host
|
||||
|
||||
|
||||
async def make_api_request(
|
||||
method: str,
|
||||
url: str,
|
||||
headers: Optional[dict] = None,
|
||||
json_data: Optional[dict] = None,
|
||||
timeout: float = 30.0,
|
||||
) -> httpx.Response:
|
||||
"""Make an API request to the cloud service."""
|
||||
headers = headers or {}
|
||||
auth_headers = await get_authenticated_headers()
|
||||
headers.update(auth_headers)
|
||||
# Add debug headers to help with compression issues
|
||||
headers.setdefault("Accept-Encoding", "identity") # Disable compression for debugging
|
||||
|
||||
async with httpx.AsyncClient(timeout=timeout) as client:
|
||||
try:
|
||||
console.print(f"[dim]Making {method} request to {url}[/dim]")
|
||||
console.print(f"[dim]Headers: {dict(headers)}[/dim]")
|
||||
|
||||
response = await client.request(method=method, url=url, headers=headers, json=json_data)
|
||||
|
||||
console.print(f"[dim]Response status: {response.status_code}[/dim]")
|
||||
console.print(f"[dim]Response headers: {dict(response.headers)}[/dim]")
|
||||
|
||||
response.raise_for_status()
|
||||
return response
|
||||
except httpx.HTTPError as e:
|
||||
console.print(f"[red]HTTP Error details: {e}[/red]")
|
||||
# Check if this is a response error with response details
|
||||
if hasattr(e, "response") and e.response is not None: # pyright: ignore [reportAttributeAccessIssue]
|
||||
response = e.response # type: ignore
|
||||
console.print(f"[red]Response status: {response.status_code}[/red]")
|
||||
console.print(f"[red]Response headers: {dict(response.headers)}[/red]")
|
||||
try:
|
||||
console.print(f"[red]Response text: {response.text}[/red]")
|
||||
except Exception:
|
||||
console.print("[red]Could not read response text[/red]")
|
||||
raise CloudAPIError(f"API request failed: {e}") from e
|
||||
|
||||
|
||||
async def get_authenticated_headers() -> dict[str, str]:
|
||||
"""Get authentication headers with JWT token."""
|
||||
client_id, domain, _ = get_cloud_config()
|
||||
auth = CLIAuth(client_id=client_id, authkit_domain=domain)
|
||||
token = await auth.get_valid_token()
|
||||
if not token:
|
||||
console.print("[red]Not authenticated. Please run 'tenant login' first.[/red]")
|
||||
raise typer.Exit(1)
|
||||
|
||||
return {"Authorization": f"Bearer {token}"}
|
||||
|
||||
|
||||
@cloud_app.command()
|
||||
def login():
|
||||
"""Authenticate with WorkOS using OAuth Device Authorization flow."""
|
||||
@@ -407,3 +354,45 @@ def status() -> None:
|
||||
except Exception as e:
|
||||
console.print(f"[red]Unexpected error: {e}[/red]")
|
||||
raise typer.Exit(1)
|
||||
|
||||
|
||||
# Mount commands
|
||||
|
||||
|
||||
@cloud_app.command("setup")
|
||||
def setup() -> None:
|
||||
"""Set up local file access with automatic rclone installation and configuration."""
|
||||
setup_cloud_mount()
|
||||
|
||||
|
||||
@cloud_app.command("mount")
|
||||
def mount(
|
||||
profile: str = typer.Option(
|
||||
"balanced", help=f"Mount profile: {', '.join(MOUNT_PROFILES.keys())}"
|
||||
),
|
||||
path: Optional[str] = typer.Option(
|
||||
None, help="Custom mount path (default: ~/basic-memory-{tenant-id})"
|
||||
),
|
||||
) -> None:
|
||||
"""Mount cloud files locally for editing."""
|
||||
try:
|
||||
mount_cloud_files(profile_name=profile)
|
||||
except Exception as e:
|
||||
console.print(f"[red]Mount failed: {e}[/red]")
|
||||
raise typer.Exit(1)
|
||||
|
||||
|
||||
@cloud_app.command("unmount")
|
||||
def unmount() -> None:
|
||||
"""Unmount cloud files."""
|
||||
try:
|
||||
unmount_cloud_files()
|
||||
except Exception as e:
|
||||
console.print(f"[red]Unmount failed: {e}[/red]")
|
||||
raise typer.Exit(1)
|
||||
|
||||
|
||||
@cloud_app.command("mount-status")
|
||||
def mount_status() -> None:
|
||||
"""Show current mount status."""
|
||||
show_mount_status()
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
"""Cloud commands package."""
|
||||
@@ -0,0 +1,77 @@
|
||||
"""Cloud API client utilities."""
|
||||
|
||||
from typing import Optional
|
||||
|
||||
import httpx
|
||||
import typer
|
||||
from rich.console import Console
|
||||
|
||||
from basic_memory.cli.auth import CLIAuth
|
||||
from basic_memory.config import ConfigManager
|
||||
|
||||
console = Console()
|
||||
|
||||
|
||||
class CloudAPIError(Exception):
|
||||
"""Exception raised for cloud API errors."""
|
||||
|
||||
pass
|
||||
|
||||
|
||||
def get_cloud_config() -> tuple[str, str, str]:
|
||||
"""Get cloud OAuth configuration from config."""
|
||||
config_manager = ConfigManager()
|
||||
config = config_manager.config
|
||||
return config.cloud_client_id, config.cloud_domain, config.cloud_host
|
||||
|
||||
|
||||
async def get_authenticated_headers() -> dict[str, str]:
|
||||
"""Get authentication headers with JWT token."""
|
||||
client_id, domain, _ = get_cloud_config()
|
||||
auth = CLIAuth(client_id=client_id, authkit_domain=domain)
|
||||
token = await auth.get_valid_token()
|
||||
if not token:
|
||||
console.print("[red]Not authenticated. Please run 'basic-memory cloud login' first.[/red]")
|
||||
raise typer.Exit(1)
|
||||
|
||||
return {"Authorization": f"Bearer {token}"}
|
||||
|
||||
|
||||
async def make_api_request(
|
||||
method: str,
|
||||
url: str,
|
||||
headers: Optional[dict] = None,
|
||||
json_data: Optional[dict] = None,
|
||||
timeout: float = 30.0,
|
||||
) -> httpx.Response:
|
||||
"""Make an API request to the cloud service."""
|
||||
headers = headers or {}
|
||||
auth_headers = await get_authenticated_headers()
|
||||
headers.update(auth_headers)
|
||||
# Add debug headers to help with compression issues
|
||||
headers.setdefault("Accept-Encoding", "identity") # Disable compression for debugging
|
||||
|
||||
async with httpx.AsyncClient(timeout=timeout) as client:
|
||||
try:
|
||||
console.print(f"[dim]Making {method} request to {url}[/dim]")
|
||||
console.print(f"[dim]Headers: {dict(headers)}[/dim]")
|
||||
|
||||
response = await client.request(method=method, url=url, headers=headers, json=json_data)
|
||||
|
||||
console.print(f"[dim]Response status: {response.status_code}[/dim]")
|
||||
console.print(f"[dim]Response headers: {dict(response.headers)}[/dim]")
|
||||
|
||||
response.raise_for_status()
|
||||
return response
|
||||
except httpx.HTTPError as e:
|
||||
console.print(f"[red]HTTP Error details: {e}[/red]")
|
||||
# Check if this is a response error with response details
|
||||
if hasattr(e, "response") and e.response is not None: # pyright: ignore [reportAttributeAccessIssue]
|
||||
response = e.response # type: ignore
|
||||
console.print(f"[red]Response status: {response.status_code}[/red]")
|
||||
console.print(f"[red]Response headers: {dict(response.headers)}[/red]")
|
||||
try:
|
||||
console.print(f"[red]Response text: {response.text}[/red]")
|
||||
except Exception:
|
||||
console.print("[red]Could not read response text[/red]")
|
||||
raise CloudAPIError(f"API request failed: {e}") from e
|
||||
@@ -0,0 +1,295 @@
|
||||
"""Cloud mount commands for Basic Memory CLI."""
|
||||
|
||||
import asyncio
|
||||
import subprocess
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
|
||||
import typer
|
||||
from rich.console import Console
|
||||
from rich.table import Table
|
||||
|
||||
from basic_memory.cli.commands.cloud.api_client import CloudAPIError, make_api_request
|
||||
from basic_memory.cli.commands.cloud.rclone_config import (
|
||||
MOUNT_PROFILES,
|
||||
add_tenant_to_rclone_config,
|
||||
build_mount_command,
|
||||
cleanup_orphaned_rclone_processes,
|
||||
get_default_mount_path,
|
||||
get_rclone_processes,
|
||||
is_path_mounted,
|
||||
unmount_path,
|
||||
)
|
||||
from basic_memory.cli.commands.cloud.rclone_installer import RcloneInstallError, install_rclone
|
||||
from basic_memory.config import ConfigManager
|
||||
|
||||
console = Console()
|
||||
|
||||
|
||||
class MountError(Exception):
|
||||
"""Exception raised for mount-related errors."""
|
||||
|
||||
pass
|
||||
|
||||
|
||||
async def get_tenant_info() -> dict:
|
||||
"""Get current tenant information from cloud API."""
|
||||
try:
|
||||
config_manager = ConfigManager()
|
||||
config = config_manager.config
|
||||
host_url = config.cloud_host.rstrip("/")
|
||||
|
||||
response = await make_api_request(method="GET", url=f"{host_url}/tenant/mount/info")
|
||||
|
||||
return response.json()
|
||||
except Exception as e:
|
||||
raise MountError(f"Failed to get tenant info: {e}") from e
|
||||
|
||||
|
||||
async def generate_mount_credentials(tenant_id: str) -> dict:
|
||||
"""Generate scoped credentials for mounting."""
|
||||
try:
|
||||
config_manager = ConfigManager()
|
||||
config = config_manager.config
|
||||
host_url = config.cloud_host.rstrip("/")
|
||||
|
||||
response = await make_api_request(method="POST", url=f"{host_url}/tenant/mount/credentials")
|
||||
|
||||
return response.json()
|
||||
except Exception as e:
|
||||
raise MountError(f"Failed to generate mount credentials: {e}") from e
|
||||
|
||||
|
||||
def setup_cloud_mount() -> None:
|
||||
"""Set up cloud mount with rclone installation and configuration."""
|
||||
console.print("[bold blue]Basic Memory Cloud Setup[/bold blue]")
|
||||
console.print("Setting up local file access to your cloud tenant...\n")
|
||||
|
||||
try:
|
||||
# Step 1: Install rclone
|
||||
console.print("[blue]Step 1: Installing rclone...[/blue]")
|
||||
install_rclone()
|
||||
|
||||
# Step 2: Get tenant info
|
||||
console.print("\n[blue]Step 2: Getting tenant information...[/blue]")
|
||||
tenant_info = asyncio.run(get_tenant_info())
|
||||
|
||||
tenant_id = tenant_info.get("tenant_id")
|
||||
bucket_name = tenant_info.get("bucket_name")
|
||||
|
||||
if not tenant_id or not bucket_name:
|
||||
raise MountError("Invalid tenant information received from cloud API")
|
||||
|
||||
console.print(f"[green]✓ Found tenant: {tenant_id}[/green]")
|
||||
console.print(f"[green]✓ Bucket: {bucket_name}[/green]")
|
||||
|
||||
# Step 3: Generate mount credentials
|
||||
console.print("\n[blue]Step 3: Generating mount credentials...[/blue]")
|
||||
creds = asyncio.run(generate_mount_credentials(tenant_id))
|
||||
|
||||
access_key = creds.get("access_key")
|
||||
secret_key = creds.get("secret_key")
|
||||
|
||||
if not access_key or not secret_key:
|
||||
raise MountError("Failed to generate mount credentials")
|
||||
|
||||
console.print("[green]✓ Generated secure credentials[/green]")
|
||||
|
||||
# Step 4: Configure rclone
|
||||
console.print("\n[blue]Step 4: Configuring rclone...[/blue]")
|
||||
add_tenant_to_rclone_config(
|
||||
tenant_id=tenant_id,
|
||||
bucket_name=bucket_name,
|
||||
access_key=access_key,
|
||||
secret_key=secret_key,
|
||||
)
|
||||
|
||||
# Step 5: Perform initial mount
|
||||
console.print("\n[blue]Step 5: Mounting cloud files...[/blue]")
|
||||
mount_path = get_default_mount_path(tenant_id)
|
||||
MOUNT_PROFILES["balanced"]
|
||||
|
||||
mount_cloud_files(
|
||||
tenant_id=tenant_id,
|
||||
bucket_name=bucket_name,
|
||||
mount_path=mount_path,
|
||||
profile_name="balanced",
|
||||
)
|
||||
|
||||
console.print("\n[bold green]✓ Cloud setup completed successfully![/bold green]")
|
||||
console.print("\nYour cloud files are now accessible at:")
|
||||
console.print(f" {mount_path}")
|
||||
console.print("\nYou can now edit files locally and they will sync to the cloud!")
|
||||
console.print("\nUseful commands:")
|
||||
console.print(" basic-memory cloud mount-status # Check mount status")
|
||||
console.print(" basic-memory cloud unmount # Unmount files")
|
||||
console.print(" basic-memory cloud mount --profile fast # Remount with faster sync")
|
||||
|
||||
except (RcloneInstallError, MountError, CloudAPIError) as e:
|
||||
console.print(f"\n[red]Setup failed: {e}[/red]")
|
||||
raise typer.Exit(1)
|
||||
except Exception as e:
|
||||
console.print(f"\n[red]Unexpected error during setup: {e}[/red]")
|
||||
raise typer.Exit(1)
|
||||
|
||||
|
||||
def mount_cloud_files(
|
||||
tenant_id: Optional[str] = None,
|
||||
bucket_name: Optional[str] = None,
|
||||
mount_path: Optional[Path] = None,
|
||||
profile_name: str = "balanced",
|
||||
) -> None:
|
||||
"""Mount cloud files with specified profile."""
|
||||
|
||||
try:
|
||||
# Get tenant info if not provided
|
||||
if not tenant_id or not bucket_name:
|
||||
tenant_info = asyncio.run(get_tenant_info())
|
||||
tenant_id = tenant_info.get("tenant_id")
|
||||
bucket_name = tenant_info.get("bucket_name")
|
||||
|
||||
if not tenant_id or not bucket_name:
|
||||
raise MountError("Could not determine tenant information")
|
||||
|
||||
# Set default mount path if not provided
|
||||
if not mount_path:
|
||||
mount_path = get_default_mount_path(tenant_id)
|
||||
|
||||
# Get mount profile
|
||||
if profile_name not in MOUNT_PROFILES:
|
||||
raise MountError(
|
||||
f"Unknown profile: {profile_name}. Available: {list(MOUNT_PROFILES.keys())}"
|
||||
)
|
||||
|
||||
profile = MOUNT_PROFILES[profile_name]
|
||||
|
||||
# Check if already mounted
|
||||
if is_path_mounted(mount_path):
|
||||
console.print(f"[yellow]Path {mount_path} is already mounted[/yellow]")
|
||||
console.print("Use 'basic-memory cloud unmount' first, or mount to a different path")
|
||||
return
|
||||
|
||||
# Create mount directory
|
||||
mount_path.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# Build and execute mount command
|
||||
mount_cmd = build_mount_command(tenant_id, bucket_name, mount_path, profile)
|
||||
|
||||
console.print(
|
||||
f"[blue]Mounting with profile '{profile_name}' ({profile.description})...[/blue]"
|
||||
)
|
||||
console.print(f"[dim]Command: {' '.join(mount_cmd)}[/dim]")
|
||||
|
||||
result = subprocess.run(mount_cmd, capture_output=True, text=True)
|
||||
|
||||
if result.returncode != 0:
|
||||
error_msg = result.stderr or "Unknown error"
|
||||
raise MountError(f"Mount command failed: {error_msg}")
|
||||
|
||||
# Wait a moment for mount to establish
|
||||
time.sleep(2)
|
||||
|
||||
# Verify mount
|
||||
if is_path_mounted(mount_path):
|
||||
console.print(f"[green]✓ Successfully mounted to {mount_path}[/green]")
|
||||
console.print(f"[green]✓ Sync profile: {profile.description}[/green]")
|
||||
else:
|
||||
raise MountError("Mount command succeeded but path is not mounted")
|
||||
|
||||
except MountError:
|
||||
raise
|
||||
except Exception as e:
|
||||
raise MountError(f"Unexpected error during mount: {e}") from e
|
||||
|
||||
|
||||
def unmount_cloud_files(tenant_id: Optional[str] = None) -> None:
|
||||
"""Unmount cloud files."""
|
||||
|
||||
try:
|
||||
# Get tenant info if not provided
|
||||
if not tenant_id:
|
||||
tenant_info = asyncio.run(get_tenant_info())
|
||||
tenant_id = tenant_info.get("tenant_id")
|
||||
|
||||
if not tenant_id:
|
||||
raise MountError("Could not determine tenant ID")
|
||||
|
||||
mount_path = get_default_mount_path(tenant_id)
|
||||
|
||||
if not is_path_mounted(mount_path):
|
||||
console.print(f"[yellow]Path {mount_path} is not mounted[/yellow]")
|
||||
return
|
||||
|
||||
console.print(f"[blue]Unmounting {mount_path}...[/blue]")
|
||||
|
||||
# Unmount the path
|
||||
if unmount_path(mount_path):
|
||||
console.print(f"[green]✓ Successfully unmounted {mount_path}[/green]")
|
||||
|
||||
# Clean up any orphaned rclone processes
|
||||
killed_count = cleanup_orphaned_rclone_processes()
|
||||
if killed_count > 0:
|
||||
console.print(
|
||||
f"[green]✓ Cleaned up {killed_count} orphaned rclone process(es)[/green]"
|
||||
)
|
||||
else:
|
||||
console.print(f"[red]✗ Failed to unmount {mount_path}[/red]")
|
||||
console.print("You may need to manually unmount or restart your system")
|
||||
|
||||
except MountError:
|
||||
raise
|
||||
except Exception as e:
|
||||
raise MountError(f"Unexpected error during unmount: {e}") from e
|
||||
|
||||
|
||||
def show_mount_status() -> None:
|
||||
"""Show current mount status and running processes."""
|
||||
|
||||
try:
|
||||
# Get tenant info
|
||||
tenant_info = asyncio.run(get_tenant_info())
|
||||
tenant_id = tenant_info.get("tenant_id")
|
||||
|
||||
if not tenant_id:
|
||||
console.print("[red]Could not determine tenant ID[/red]")
|
||||
return
|
||||
|
||||
mount_path = get_default_mount_path(tenant_id)
|
||||
|
||||
# Create status table
|
||||
table = Table(title="Cloud Mount Status", show_header=True, header_style="bold blue")
|
||||
table.add_column("Property", style="green", min_width=15)
|
||||
table.add_column("Value", style="dim", min_width=30)
|
||||
|
||||
# Check mount status
|
||||
is_mounted = is_path_mounted(mount_path)
|
||||
mount_status = "[green]✓ Mounted[/green]" if is_mounted else "[red]✗ Not mounted[/red]"
|
||||
|
||||
table.add_row("Tenant ID", tenant_id)
|
||||
table.add_row("Mount Path", str(mount_path))
|
||||
table.add_row("Status", mount_status)
|
||||
|
||||
# Get rclone processes
|
||||
processes = get_rclone_processes()
|
||||
if processes:
|
||||
table.add_row("rclone Processes", f"{len(processes)} running")
|
||||
else:
|
||||
table.add_row("rclone Processes", "None")
|
||||
|
||||
console.print(table)
|
||||
|
||||
# Show running processes details
|
||||
if processes:
|
||||
console.print("\n[bold]Running rclone processes:[/bold]")
|
||||
for proc in processes:
|
||||
console.print(f" PID {proc['pid']}: {proc['command'][:80]}...")
|
||||
|
||||
# Show mount profiles
|
||||
console.print("\n[bold]Available mount profiles:[/bold]")
|
||||
for name, profile in MOUNT_PROFILES.items():
|
||||
console.print(f" {name}: {profile.description}")
|
||||
|
||||
except Exception as e:
|
||||
console.print(f"[red]Error getting mount status: {e}[/red]")
|
||||
raise typer.Exit(1)
|
||||
@@ -0,0 +1,284 @@
|
||||
"""rclone configuration management for Basic Memory Cloud."""
|
||||
|
||||
import configparser
|
||||
import os
|
||||
import shutil
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
from typing import Dict, List, Optional
|
||||
|
||||
from rich.console import Console
|
||||
|
||||
console = Console()
|
||||
|
||||
|
||||
class RcloneConfigError(Exception):
|
||||
"""Exception raised for rclone configuration errors."""
|
||||
|
||||
pass
|
||||
|
||||
|
||||
class RcloneMountProfile:
|
||||
"""Mount profile with optimized settings."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
name: str,
|
||||
cache_time: str,
|
||||
poll_interval: str,
|
||||
attr_timeout: str,
|
||||
write_back: str,
|
||||
description: str,
|
||||
extra_args: Optional[List[str]] = None,
|
||||
):
|
||||
self.name = name
|
||||
self.cache_time = cache_time
|
||||
self.poll_interval = poll_interval
|
||||
self.attr_timeout = attr_timeout
|
||||
self.write_back = write_back
|
||||
self.description = description
|
||||
self.extra_args = extra_args or []
|
||||
|
||||
|
||||
# Mount profiles based on SPEC-7 Phase 4 testing
|
||||
MOUNT_PROFILES = {
|
||||
"fast": RcloneMountProfile(
|
||||
name="fast",
|
||||
cache_time="5s",
|
||||
poll_interval="3s",
|
||||
attr_timeout="3s",
|
||||
write_back="1s",
|
||||
description="Ultra-fast development (5s sync, higher bandwidth)",
|
||||
),
|
||||
"balanced": RcloneMountProfile(
|
||||
name="balanced",
|
||||
cache_time="10s",
|
||||
poll_interval="5s",
|
||||
attr_timeout="5s",
|
||||
write_back="2s",
|
||||
description="Fast development (10-15s sync, recommended)",
|
||||
),
|
||||
"safe": RcloneMountProfile(
|
||||
name="safe",
|
||||
cache_time="15s",
|
||||
poll_interval="10s",
|
||||
attr_timeout="10s",
|
||||
write_back="5s",
|
||||
description="Conflict-aware mount with backup",
|
||||
extra_args=[
|
||||
"--conflict-suffix",
|
||||
".conflict-{DateTimeExt}",
|
||||
"--backup-dir",
|
||||
"~/.basic-memory/conflicts",
|
||||
"--track-renames",
|
||||
],
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def get_rclone_config_path() -> Path:
|
||||
"""Get the path to rclone configuration file."""
|
||||
config_dir = Path.home() / ".config" / "rclone"
|
||||
config_dir.mkdir(parents=True, exist_ok=True)
|
||||
return config_dir / "rclone.conf"
|
||||
|
||||
|
||||
def backup_rclone_config() -> Optional[Path]:
|
||||
"""Create a backup of existing rclone config."""
|
||||
config_path = get_rclone_config_path()
|
||||
if not config_path.exists():
|
||||
return None
|
||||
|
||||
backup_path = config_path.with_suffix(f".conf.backup-{os.getpid()}")
|
||||
shutil.copy2(config_path, backup_path)
|
||||
console.print(f"[dim]Created backup: {backup_path}[/dim]")
|
||||
return backup_path
|
||||
|
||||
|
||||
def load_rclone_config() -> configparser.ConfigParser:
|
||||
"""Load existing rclone configuration."""
|
||||
config = configparser.ConfigParser()
|
||||
config_path = get_rclone_config_path()
|
||||
|
||||
if config_path.exists():
|
||||
config.read(config_path)
|
||||
|
||||
return config
|
||||
|
||||
|
||||
def save_rclone_config(config: configparser.ConfigParser) -> None:
|
||||
"""Save rclone configuration to file."""
|
||||
config_path = get_rclone_config_path()
|
||||
|
||||
with open(config_path, "w") as f:
|
||||
config.write(f)
|
||||
|
||||
console.print(f"[dim]Updated rclone config: {config_path}[/dim]")
|
||||
|
||||
|
||||
def add_tenant_to_rclone_config(
|
||||
tenant_id: str,
|
||||
bucket_name: str,
|
||||
access_key: str,
|
||||
secret_key: str,
|
||||
endpoint: str = "https://fly.storage.tigris.dev",
|
||||
region: str = "auto",
|
||||
) -> str:
|
||||
"""Add tenant configuration to rclone config file."""
|
||||
|
||||
# Backup existing config
|
||||
backup_rclone_config()
|
||||
|
||||
# Load existing config
|
||||
config = load_rclone_config()
|
||||
|
||||
# Create section name
|
||||
section_name = f"basic-memory-{tenant_id}"
|
||||
|
||||
# Add/update the tenant section
|
||||
if not config.has_section(section_name):
|
||||
config.add_section(section_name)
|
||||
|
||||
config.set(section_name, "type", "s3")
|
||||
config.set(section_name, "provider", "Other")
|
||||
config.set(section_name, "access_key_id", access_key)
|
||||
config.set(section_name, "secret_access_key", secret_key)
|
||||
config.set(section_name, "endpoint", endpoint)
|
||||
config.set(section_name, "region", region)
|
||||
|
||||
# Save updated config
|
||||
save_rclone_config(config)
|
||||
|
||||
console.print(f"[green]✓ Added tenant {tenant_id} to rclone config[/green]")
|
||||
return section_name
|
||||
|
||||
|
||||
def remove_tenant_from_rclone_config(tenant_id: str) -> bool:
|
||||
"""Remove tenant configuration from rclone config."""
|
||||
config = load_rclone_config()
|
||||
section_name = f"basic-memory-{tenant_id}"
|
||||
|
||||
if config.has_section(section_name):
|
||||
backup_rclone_config()
|
||||
config.remove_section(section_name)
|
||||
save_rclone_config(config)
|
||||
console.print(f"[green]✓ Removed tenant {tenant_id} from rclone config[/green]")
|
||||
return True
|
||||
|
||||
return False
|
||||
|
||||
|
||||
def get_default_mount_path(tenant_id: str) -> Path:
|
||||
"""Get default mount path for a tenant."""
|
||||
return Path.home() / f"basic-memory-{tenant_id}"
|
||||
|
||||
|
||||
def build_mount_command(
|
||||
tenant_id: str, bucket_name: str, mount_path: Path, profile: RcloneMountProfile
|
||||
) -> List[str]:
|
||||
"""Build rclone mount command with optimized settings."""
|
||||
|
||||
rclone_remote = f"basic-memory-{tenant_id}:{bucket_name}"
|
||||
|
||||
cmd = [
|
||||
"rclone",
|
||||
"nfsmount",
|
||||
rclone_remote,
|
||||
str(mount_path),
|
||||
"--vfs-cache-mode",
|
||||
"writes",
|
||||
"--dir-cache-time",
|
||||
profile.cache_time,
|
||||
"--vfs-cache-poll-interval",
|
||||
profile.poll_interval,
|
||||
"--attr-timeout",
|
||||
profile.attr_timeout,
|
||||
"--vfs-write-back",
|
||||
profile.write_back,
|
||||
"--daemon",
|
||||
]
|
||||
|
||||
# Add profile-specific extra arguments
|
||||
cmd.extend(profile.extra_args)
|
||||
|
||||
return cmd
|
||||
|
||||
|
||||
def is_path_mounted(mount_path: Path) -> bool:
|
||||
"""Check if a path is currently mounted."""
|
||||
if not mount_path.exists():
|
||||
return False
|
||||
|
||||
try:
|
||||
# Check if mount point is actually mounted by looking for mount table entry
|
||||
result = subprocess.run(["mount"], capture_output=True, text=True, check=False)
|
||||
|
||||
if result.returncode == 0:
|
||||
# Look for our mount path in mount output
|
||||
mount_str = str(mount_path.resolve())
|
||||
return mount_str in result.stdout
|
||||
|
||||
return False
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def get_rclone_processes() -> List[Dict[str, str]]:
|
||||
"""Get list of running rclone processes."""
|
||||
try:
|
||||
# Use ps to find rclone processes
|
||||
result = subprocess.run(
|
||||
["ps", "-eo", "pid,args"], capture_output=True, text=True, check=False
|
||||
)
|
||||
|
||||
processes = []
|
||||
if result.returncode == 0:
|
||||
for line in result.stdout.split("\n"):
|
||||
if "rclone" in line and "basic-memory" in line:
|
||||
parts = line.strip().split(None, 1)
|
||||
if len(parts) >= 2:
|
||||
processes.append({"pid": parts[0], "command": parts[1]})
|
||||
|
||||
return processes
|
||||
except Exception:
|
||||
return []
|
||||
|
||||
|
||||
def kill_rclone_process(pid: str) -> bool:
|
||||
"""Kill a specific rclone process."""
|
||||
try:
|
||||
subprocess.run(["kill", pid], check=True)
|
||||
console.print(f"[green]✓ Killed rclone process {pid}[/green]")
|
||||
return True
|
||||
except subprocess.CalledProcessError:
|
||||
console.print(f"[red]✗ Failed to kill rclone process {pid}[/red]")
|
||||
return False
|
||||
|
||||
|
||||
def unmount_path(mount_path: Path) -> bool:
|
||||
"""Unmount a mounted path."""
|
||||
if not is_path_mounted(mount_path):
|
||||
return True
|
||||
|
||||
try:
|
||||
subprocess.run(["umount", str(mount_path)], check=True)
|
||||
console.print(f"[green]✓ Unmounted {mount_path}[/green]")
|
||||
return True
|
||||
except subprocess.CalledProcessError as e:
|
||||
console.print(f"[red]✗ Failed to unmount {mount_path}: {e}[/red]")
|
||||
return False
|
||||
|
||||
|
||||
def cleanup_orphaned_rclone_processes() -> int:
|
||||
"""Clean up orphaned rclone processes for basic-memory."""
|
||||
processes = get_rclone_processes()
|
||||
killed_count = 0
|
||||
|
||||
for proc in processes:
|
||||
console.print(
|
||||
f"[yellow]Found rclone process: {proc['pid']} - {proc['command'][:80]}...[/yellow]"
|
||||
)
|
||||
if kill_rclone_process(proc["pid"]):
|
||||
killed_count += 1
|
||||
|
||||
return killed_count
|
||||
@@ -0,0 +1,198 @@
|
||||
"""Cross-platform rclone installation utilities."""
|
||||
|
||||
import platform
|
||||
import shutil
|
||||
import subprocess
|
||||
from typing import Optional
|
||||
|
||||
from rich.console import Console
|
||||
|
||||
console = Console()
|
||||
|
||||
|
||||
class RcloneInstallError(Exception):
|
||||
"""Exception raised for rclone installation errors."""
|
||||
|
||||
pass
|
||||
|
||||
|
||||
def is_rclone_installed() -> bool:
|
||||
"""Check if rclone is already installed and available in PATH."""
|
||||
return shutil.which("rclone") is not None
|
||||
|
||||
|
||||
def get_platform() -> str:
|
||||
"""Get the current platform identifier."""
|
||||
system = platform.system().lower()
|
||||
if system == "darwin":
|
||||
return "macos"
|
||||
elif system == "linux":
|
||||
return "linux"
|
||||
elif system == "windows":
|
||||
return "windows"
|
||||
else:
|
||||
raise RcloneInstallError(f"Unsupported platform: {system}")
|
||||
|
||||
|
||||
def run_command(command: list[str], check: bool = True) -> subprocess.CompletedProcess:
|
||||
"""Run a command with proper error handling."""
|
||||
try:
|
||||
console.print(f"[dim]Running: {' '.join(command)}[/dim]")
|
||||
result = subprocess.run(command, capture_output=True, text=True, check=check)
|
||||
if result.stdout:
|
||||
console.print(f"[dim]Output: {result.stdout.strip()}[/dim]")
|
||||
return result
|
||||
except subprocess.CalledProcessError as e:
|
||||
console.print(f"[red]Command failed: {e}[/red]")
|
||||
if e.stderr:
|
||||
console.print(f"[red]Error output: {e.stderr}[/red]")
|
||||
raise RcloneInstallError(f"Command failed: {e}") from e
|
||||
except FileNotFoundError as e:
|
||||
raise RcloneInstallError(f"Command not found: {' '.join(command)}") from e
|
||||
|
||||
|
||||
def install_rclone_macos() -> None:
|
||||
"""Install rclone on macOS using Homebrew or official script."""
|
||||
# Try Homebrew first
|
||||
if shutil.which("brew"):
|
||||
try:
|
||||
console.print("[blue]Installing rclone via Homebrew...[/blue]")
|
||||
run_command(["brew", "install", "rclone"])
|
||||
console.print("[green]✓ rclone installed via Homebrew[/green]")
|
||||
return
|
||||
except RcloneInstallError:
|
||||
console.print(
|
||||
"[yellow]Homebrew installation failed, trying official script...[/yellow]"
|
||||
)
|
||||
|
||||
# Fallback to official script
|
||||
console.print("[blue]Installing rclone via official script...[/blue]")
|
||||
try:
|
||||
run_command(["sh", "-c", "curl https://rclone.org/install.sh | sudo bash"])
|
||||
console.print("[green]✓ rclone installed via official script[/green]")
|
||||
except RcloneInstallError:
|
||||
raise RcloneInstallError(
|
||||
"Failed to install rclone. Please install manually: brew install rclone"
|
||||
)
|
||||
|
||||
|
||||
def install_rclone_linux() -> None:
|
||||
"""Install rclone on Linux using package managers or official script."""
|
||||
# Try snap first (most universal)
|
||||
if shutil.which("snap"):
|
||||
try:
|
||||
console.print("[blue]Installing rclone via snap...[/blue]")
|
||||
run_command(["sudo", "snap", "install", "rclone"])
|
||||
console.print("[green]✓ rclone installed via snap[/green]")
|
||||
return
|
||||
except RcloneInstallError:
|
||||
console.print("[yellow]Snap installation failed, trying apt...[/yellow]")
|
||||
|
||||
# Try apt (Debian/Ubuntu)
|
||||
if shutil.which("apt"):
|
||||
try:
|
||||
console.print("[blue]Installing rclone via apt...[/blue]")
|
||||
run_command(["sudo", "apt", "update"])
|
||||
run_command(["sudo", "apt", "install", "-y", "rclone"])
|
||||
console.print("[green]✓ rclone installed via apt[/green]")
|
||||
return
|
||||
except RcloneInstallError:
|
||||
console.print("[yellow]apt installation failed, trying official script...[/yellow]")
|
||||
|
||||
# Fallback to official script
|
||||
console.print("[blue]Installing rclone via official script...[/blue]")
|
||||
try:
|
||||
run_command(["sh", "-c", "curl https://rclone.org/install.sh | sudo bash"])
|
||||
console.print("[green]✓ rclone installed via official script[/green]")
|
||||
except RcloneInstallError:
|
||||
raise RcloneInstallError(
|
||||
"Failed to install rclone. Please install manually: sudo snap install rclone"
|
||||
)
|
||||
|
||||
|
||||
def install_rclone_windows() -> None:
|
||||
"""Install rclone on Windows using package managers."""
|
||||
# Try winget first (built into Windows 10+)
|
||||
if shutil.which("winget"):
|
||||
try:
|
||||
console.print("[blue]Installing rclone via winget...[/blue]")
|
||||
run_command(["winget", "install", "Rclone.Rclone"])
|
||||
console.print("[green]✓ rclone installed via winget[/green]")
|
||||
return
|
||||
except RcloneInstallError:
|
||||
console.print("[yellow]winget installation failed, trying chocolatey...[/yellow]")
|
||||
|
||||
# Try chocolatey
|
||||
if shutil.which("choco"):
|
||||
try:
|
||||
console.print("[blue]Installing rclone via chocolatey...[/blue]")
|
||||
run_command(["choco", "install", "rclone", "-y"])
|
||||
console.print("[green]✓ rclone installed via chocolatey[/green]")
|
||||
return
|
||||
except RcloneInstallError:
|
||||
console.print("[yellow]chocolatey installation failed, trying scoop...[/yellow]")
|
||||
|
||||
# Try scoop
|
||||
if shutil.which("scoop"):
|
||||
try:
|
||||
console.print("[blue]Installing rclone via scoop...[/blue]")
|
||||
run_command(["scoop", "install", "rclone"])
|
||||
console.print("[green]✓ rclone installed via scoop[/green]")
|
||||
return
|
||||
except RcloneInstallError:
|
||||
console.print("[yellow]scoop installation failed[/yellow]")
|
||||
|
||||
# No package manager available
|
||||
raise RcloneInstallError(
|
||||
"Could not install rclone automatically. Please install a package manager "
|
||||
"(winget, chocolatey, or scoop) or install rclone manually from https://rclone.org/downloads/"
|
||||
)
|
||||
|
||||
|
||||
def install_rclone(platform_override: Optional[str] = None) -> None:
|
||||
"""Install rclone for the current platform."""
|
||||
if is_rclone_installed():
|
||||
console.print("[green]rclone is already installed[/green]")
|
||||
return
|
||||
|
||||
platform_name = platform_override or get_platform()
|
||||
console.print(f"[blue]Installing rclone for {platform_name}...[/blue]")
|
||||
|
||||
try:
|
||||
if platform_name == "macos":
|
||||
install_rclone_macos()
|
||||
elif platform_name == "linux":
|
||||
install_rclone_linux()
|
||||
elif platform_name == "windows":
|
||||
install_rclone_windows()
|
||||
else:
|
||||
raise RcloneInstallError(f"Unsupported platform: {platform_name}")
|
||||
|
||||
# Verify installation
|
||||
if not is_rclone_installed():
|
||||
raise RcloneInstallError("rclone installation completed but command not found in PATH")
|
||||
|
||||
console.print("[green]✓ rclone installation completed successfully[/green]")
|
||||
|
||||
except RcloneInstallError:
|
||||
raise
|
||||
except Exception as e:
|
||||
raise RcloneInstallError(f"Unexpected error during installation: {e}") from e
|
||||
|
||||
|
||||
def get_rclone_version() -> Optional[str]:
|
||||
"""Get the installed rclone version."""
|
||||
if not is_rclone_installed():
|
||||
return None
|
||||
|
||||
try:
|
||||
result = run_command(["rclone", "version"], check=False)
|
||||
if result.returncode == 0:
|
||||
# Parse version from output (format: "rclone v1.64.0")
|
||||
lines = result.stdout.strip().split("\n")
|
||||
for line in lines:
|
||||
if line.startswith("rclone v"):
|
||||
return line.split()[1]
|
||||
return "unknown"
|
||||
except Exception:
|
||||
return "unknown"
|
||||
Reference in New Issue
Block a user