from __future__ import annotations

import os
import sys
import time
import logging
import threading
import shutil
import json
from pathlib import Path
from unittest.mock import MagicMock, patch

import pytest

# Add src to sys.path
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src"))

from api_client import ApiClient, ApiError
from config import Config
from command_runner import (
    process_commands,
    load_command_states,
    update_command_state,
    save_command_states,
    recover_interrupted_commands,
    CommandWorkerManager,
    HANDLERS,
)
from update import self_update, UpdateInfo

def _fresh_config(tmp: str) -> Config:
    os.environ["WOLFPANEL_DEV_HOME"] = tmp
    os.environ["WOLFPANEL_API_MOCK"] = "1"
    import importlib
    import config as config_module
    importlib.reload(config_module)
    return config_module.load_config()

def _registered_config(tmp: str) -> Config:
    cfg = _fresh_config(tmp)
    from identity import ensure_dirs, store_credentials
    ensure_dirs(cfg)
    store_credentials(cfg, "server_123", "wp_agent_dev_token")
    return _fresh_config(tmp)

@pytest.fixture(autouse=True)
def clean_env():
    # Store originals
    orig_async = os.environ.get("WOLFPANEL_TEST_ASYNC")
    orig_max = os.environ.get("WOLFPANEL_MAX_CONCURRENT_COMMANDS")
    
    # Pre-test cleanup
    import command_runner
    if command_runner._worker_manager is not None:
        command_runner._worker_manager.shutdown()
        command_runner._worker_manager = None
        
    yield
    
    # Restore
    if orig_async is not None:
        os.environ["WOLFPANEL_TEST_ASYNC"] = orig_async
    else:
        os.environ.pop("WOLFPANEL_TEST_ASYNC", None)
        
    if orig_max is not None:
        os.environ["WOLFPANEL_MAX_CONCURRENT_COMMANDS"] = orig_max
    else:
        os.environ.pop("WOLFPANEL_MAX_CONCURRENT_COMMANDS", None)

    # Post-test cleanup
    if command_runner._worker_manager is not None:
        command_runner._worker_manager.shutdown()
        command_runner._worker_manager = None

def test_queue_capacity_limit(tmp_path):
    cfg = _registered_config(str(tmp_path))
    api_client = MagicMock()
    
    # Set worker manager with async execution
    os.environ["WOLFPANEL_TEST_ASYNC"] = "1"
    os.environ["WOLFPANEL_MAX_CONCURRENT_COMMANDS"] = "1"
    
    manager = CommandWorkerManager(cfg, api_client)
    
    # Mock handlers
    event = threading.Event()
    def slow_handler(payload):
        event.wait(timeout=5.0)
        return {"success": True}
    HANDLERS["test.slow"] = slow_handler
    
    # Enqueue first to start processing and block the worker thread
    manager.enqueue({"id": 1, "command_type": "test.slow", "payload": {}})
    
    # Fill the queue manually up to capacity (maxsize=1000)
    # The first item has been retrieved from the queue, so queue has 1000 empty slots.
    for i in range(2, 1002):
        manager.enqueue({"id": i, "command_type": "test.slow", "payload": {}})
        
    # The queue should now be completely full (1000 items)
    assert manager.queue.full()
    
    # Try to enqueue 1002nd item - should be skipped/rejected and not update state or block
    manager.enqueue({"id": 1002, "command_type": "test.slow", "payload": {}})
    
    # Unblock worker and shut down
    event.set()
    manager.shutdown()
    
    # Verify state for 1002 was not written/updated
    states = load_command_states(cfg.var_dir)
    assert "1002" not in states

def test_race_condition_duplicate_prevention(tmp_path):
    cfg = _registered_config(str(tmp_path))
    api_client = MagicMock()
    
    os.environ["WOLFPANEL_TEST_ASYNC"] = "1"
    os.environ["WOLFPANEL_MAX_CONCURRENT_COMMANDS"] = "1"
    
    event = threading.Event()
    def slow_handler(payload):
        event.wait(timeout=5.0)
        return {"success": True}
    HANDLERS["test.slow"] = slow_handler
    
    # Mock API to return the same command ID twice in subsequent polls
    api_client.get_pending_commands.return_value = [
        {"id": 501, "command_type": "test.slow", "payload": {}}
    ]
    
    with patch("command_runner.load_agent_token", return_value="wp_agent_dev_token"):
        # First poll enqueues and starts running command
        process_commands(cfg, api_client)
        time.sleep(0.1)
        
        # Second poll: since the command is still running locally, it should be skipped/ignored
        process_commands(cfg, api_client)
        
        # Let's check states
        states = load_command_states(cfg.var_dir)
        assert states["501"]["status"] == "running"
        
        # Let's count how many times it was enqueued. 
        # Check that it's only in queue once.
        import command_runner
        assert command_runner._worker_manager.queue.qsize() == 0  # because worker thread retrieved it
        
        event.set()
        time.sleep(0.5)

def test_power_loss_durability(tmp_path):
    cfg = _registered_config(str(tmp_path))
    
    # Call save_command_states to verify durability path does not crash
    states = {"999": {"status": "success", "updated_at": "test-time"}}
    save_command_states(cfg.var_dir, states)
    
    loaded = load_command_states(cfg.var_dir)
    assert loaded["999"]["status"] == "success"
    
    # Check corrupt json loading handling
    state_file = cfg.var_dir / "command_states.json"
    state_file.write_text("{bad json...", encoding="utf-8")
    assert load_command_states(cfg.var_dir) == {}

def test_worker_death_resurrection(tmp_path):
    cfg = _registered_config(str(tmp_path))
    api_client = MagicMock()
    
    os.environ["WOLFPANEL_TEST_ASYNC"] = "1"
    os.environ["WOLFPANEL_MAX_CONCURRENT_COMMANDS"] = "2"
    
    manager = CommandWorkerManager(cfg, api_client)
    manager.start_workers()
    
    assert len(manager.workers) == 2
    # Ensure all are alive
    for t in manager.workers:
        assert t.is_alive()
        
    # Simulate a thread dying by mock
    manager.workers[0].is_alive = MagicMock(return_value=False)
    
    # Trigger start_workers again
    manager.start_workers()
    
    # Verify it resurrected the dead thread to maintain size of 2
    assert len(manager.workers) == 2
    for t in manager.workers:
        assert t.is_alive() or t is not manager.workers[0]
        
    manager.shutdown()

def test_update_waits_for_running_commands(tmp_path):
    cfg = _registered_config(str(tmp_path))
    api_client = MagicMock()
    
    os.environ["WOLFPANEL_TEST_ASYNC"] = "1"
    os.environ["WOLFPANEL_MAX_CONCURRENT_COMMANDS"] = "2"
    
    import command_runner
    command_runner._worker_manager = CommandWorkerManager(cfg, api_client)
    command_runner._worker_manager.start_workers()
    
    # Simulate a running slow command
    event = threading.Event()
    def slow_cmd_handler(payload):
        event.wait(timeout=2.0)
        return {"success": True}
    HANDLERS["test.slow"] = slow_cmd_handler
    
    # Enqueue it
    command_runner._worker_manager.enqueue({"id": 701, "command_type": "test.slow", "payload": {}})
    time.sleep(0.1)
    
    # Verify it is in active commands
    assert "701" in command_runner._worker_manager.active_commands
    
    # Mock update info
    info = UpdateInfo(
        current="0.1.1",
        latest="0.1.2",
        update_available=True,
        url="https://downloads.wolfpanel.net/agent/v0.1.2/agent.tar.gz",
        sha256="mock-sha"
    )
    
    # Mock swap directories atomic and clean_old_versions
    with patch("update.swap_directories_atomic") as mock_swap, \
         patch("update.clean_old_versions"):
         
        # We want to unblock the slow command halfway during wait
        def unblock_soon():
            time.sleep(0.5)
            event.set()
        threading.Thread(target=unblock_soon, daemon=True).start()
        
        # Start update (under pytest it will return without executing sys.exit)
        self_update(cfg, info, {"job_id": "702"})
        
        # The update should have waited for 701 to finish (which takes 0.5s)
        # Verify swap was called after 701 completed
        mock_swap.assert_called_once()
        assert "701" not in command_runner._worker_manager.active_commands
        
    command_runner._worker_manager.shutdown()
    command_runner._worker_manager = None
