Compare commits

..

7 Commits

Author SHA1 Message Date
brotoskyj 53f0b548b0 Bug Fixes
Build Library / Build Library (push) Successful in 4m51s
Changed days to days.to_string() to fix type confusion
closes #47
2025-12-17 11:02:13 -05:00
Zarithas 66eb101c5d Fixed Bitbake 2025-12-16 17:07:35 -05:00
Zarithas 22101c1eba Merge branch 'Zar-Branch' of https://git.racooncity.org/brotoskyj/Airlocktools into Zar-Branch
git commit -m "feat: Multiple UI improvements and new server log functionality

- Add Server Log tab with DataTable display of server activity logs

- Fix keyboard navigation bug in agents tab

- Add execution history viewer for selected agents

- Improve policy tree widget functionality by adding single device operations

- Integrate logging notifications into TUI
  - Add TextualNotificationHandler to setup.py
  - Display ERROR/WARNING/CRITICAL logs as toast notifications
  - Remove terminal output to prevent interference with TUI
  - Logs still written to Loxide.log file

closes #45"
2025-12-16 16:17:25 -05:00
Zarithas fc17c869fc feat: Multiple UI improvements and new server log functionality
- Add Server Log tab with DataTable display of server activity logs

- Fix keyboard navigation bug in agents tab

- Add execution history viewer for selected agents

- Improve policy tree widget functionality by adding single device operations

- Integrate logging notifications into TUI
  - Add TextualNotificationHandler to setup.py
  - Display ERROR/WARNING/CRITICAL logs as toast notifications
  - Remove terminal output to prevent interference with TUI
  - Logs still written to Loxide.log file

closes #45
2025-12-16 16:16:49 -05:00
brotoskyj 1dbbcff5d5 Merge remote-tracking branch 'origin/RustImplementation' into RustImplementation
Build Library / Build Library (push) Successful in 5m18s
2025-12-15 17:04:32 -05:00
brotoskyj f080b0034f Changes:
Features
Moved the init_tracer() function to an implementation in TelemetryConfig for cleaner main file
closes #43

Bug Fixes
Telemetry is now opt-in again. Due to changes in opentelemetry, this required a major overhaul of the telemetryconfig function.
closes #44
2025-12-15 17:04:17 -05:00
Zarithas 57d0f12000 feat: Complete Policy Prep Workflow with UX upgrades, Liftoff API, and TUI merge
- Added intro screen with workflow overview, time estimate, and onboarding controls
- Improved visuals: cleaner checkboxes (/), better loading screen layout
- Enforced mandatory tab reviews for critical steps with warnings and blocked navigation
- Optimized logging: INFO for milestones, DEBUG for internals; cleaner production logs
- Implemented Liftoff API integration: paths, publishers, hashes with granular error handling
- Color-coded completion feedback ( success,  failure,  partial) and detailed summaries
- Consolidated architecture: merged TUI.py into Loxide.py (single entry point, no circular imports)
- Fixed race condition in table creation with concurrency locks
2025-12-15 17:01:58 -05:00
17 changed files with 2310 additions and 605 deletions
+539 -6
View File
@@ -23,23 +23,556 @@
import logging
import os
from typing import Optional
import dotenv
from textual.app import App, ComposeResult
from textual.containers import Vertical
from textual.message import Message
from textual.reactive import reactive
from textual.screen import Screen
from textual.widgets import (
Button,
DirectoryTree,
Footer,
Header,
Static,
Tab,
Tabs,
)
import urllib3
from models.agent import Agent
from models.policy import Policy
from services.API import AirlockAPIWrapper
from services.security import getAPI
from TUI.TUI import run_Loxide
from utils.configmanager import get_system_value
from utils.setup import setup
from utils.utils import irtang
from TUI.Screens.executionhistoryscreen import ExecutionHistoryScreen
from TUI.Screens.moveagentworkflowscreen import MoveAgentWorkflowScreen
from TUI.Screens.otpactivityscreen import OTPActivitiesScreen
from TUI.Screens.otprevokescreen import OTPRevokeScreen
from TUI.Screens.otpworkflowscreen import OTPWorkflowScreen
from TUI.Screens.policyprepworkflowscreen import PolicyPrepWorkflowScreen
from TUI.Screens.quietagentworkflowscreen import QuietAgentWorkflowScreen
from TUI.Themes.theme_amber_terminal import get_amber_terminal_theme
from TUI.Themes.theme_retro_terminal import get_retro_terminal_theme
from TUI.Themes.themeselector import ThemeSelector
from TUI.Widgets.agentmoveoperations import AgentMoveOperations
from TUI.Widgets.multiagentselector import MultiAgentSelector
from TUI.Widgets.policytreewidget import PolicyTreeWidget
from TUI.Widgets.resultsdisplay import ResultsDisplay
from TUI.Widgets.serverlogwidget import ServerLogWidget
from utils.configmanager import (
get_system_value,
get_user_value,
load_env,
save_user_config,
)
from utils.setup import get_base_directory, setup
from utils.utils import irtang, open_directory
dotenv.load_dotenv()
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
# ---------------------------------------------------------------------------
# GLOBAL STASH
# ---------------------------------------------------------------------------
_APP_RESTART_REASON = None
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# helper to persist TEXTUAL_THEME to *user* config and mirror to .env
# ---------------------------------------------------------------------------
def _persist_user_theme(theme_name: str) -> None:
"""
Store the chosen Textual theme in the user's config using the config manager.
No need to touch .env - config manager handles everything.
"""
base_dir = get_base_directory()
config_dir = base_dir / "config"
try:
save_user_config(config_dir, {"TEXTUAL_THEME": theme_name})
logger.debug("Updated user config with TEXTUAL_THEME=%s", theme_name)
except Exception as exc:
logger.error("Failed to save TEXTUAL_THEME: %s", exc)
# ---------------------------------------------------------------------------
# 1) SCREEN
# ---------------------------------------------------------------------------
class MainMenuScreen(Screen):
api: AirlockAPIWrapper
current_tab = reactive("")
BUTTON_DEFS = {
"agent_actions": [
{
"label": "🖥️ - Multi-Agent Operations",
"id": "move_agent_workflow_button",
"description": "Select agents to: Move policies, Generate OTPs, Toggle audit/enforcement, View history, Export data",
},
{
"label": "🎫 - Review and approve OTP Activities",
"id": "otp_activities_button",
},
{
"label": "🛑 - Revoke Active OTP Session",
"id": "otp_revoke_button",
},
],
"policy": [
{
"label": "⚖️ - Prepare Policy For Enforcement",
"id": "policy_prep_button",
},
{
"label": "🔕 - Find and Move Quiet Hosts to Enforcement",
"id": "find_quiet_button",
},
],
}
def __init__(self) -> None:
super().__init__()
self.extras = get_user_value("EXTRAS", str, "NOTTODAY")
wd = load_env("WORKING_DIR") or os.getcwd()
if not os.path.isdir(wd):
wd = os.getcwd()
self.working_dir = wd
def _make_buttons_for(self, tab_id: str) -> Vertical:
defs = self.BUTTON_DEFS.get(tab_id, [])
widgets = []
for item in defs:
# Support both old tuple format and new dict format
if isinstance(item, dict):
label = item["label"]
btn_id = item["id"]
description = item.get("description")
else:
# Old tuple format: (label, id)
label, btn_id = item
description = None
btn = Button(label, id=btn_id)
btn.styles.width = "100%"
widgets.append(btn)
# Add description text if provided
if description:
desc_text = Static(description, classes="button_description")
desc_text.styles.width = "100%"
desc_text.styles.color = "ansi_bright_black"
desc_text.styles.text_align = "center"
desc_text.styles.margin = (0, 0, 1, 0)
widgets.append(desc_text)
return Vertical(*widgets)
def compose(self) -> ComposeResult:
yield Header(show_clock=True, icon="")
tabs = [
Tab("Agents", id="agent_actions"),
Tab("Tree View", id="p_tree"),
Tab("Server Log", id="server_log"),
Tab("Directory", id="dir"),
Tab("Settings", id="settings"),
]
if self.extras == "POLICYPREP":
tabs.insert(2, Tab("Policy Prep", id="policy"))
yield Tabs(*tabs, id="tabs")
yield Vertical(id="content")
yield Footer()
def on_mount(self) -> None:
self.switch_tab("agent_actions")
def on_key(self, event) -> None:
"""Handle up/down arrow keys for button navigation."""
if event.key == "down":
self._focus_nearby_button(1)
event.prevent_default()
event.stop()
elif event.key == "up":
self._focus_nearby_button(-1)
event.prevent_default()
event.stop()
# left/right are handled by Textual's default tab navigation
# focus helpers
def _get_content_buttons(self) -> list[Button]:
content = self.query_one("#content", Vertical)
return list(content.query(Button))
def _focus_first_button(self) -> None:
buttons = self._get_content_buttons()
if buttons:
buttons[0].focus()
def _focus_tabs(self) -> None:
tabs = self.query_one("#tabs", Tabs)
tabs.focus()
def _focus_nearby_button(self, direction: int) -> None:
buttons = self._get_content_buttons()
if not buttons:
return
try:
current = next(i for i, b in enumerate(buttons) if b.has_focus)
except StopIteration:
if direction > 0:
buttons[0].focus()
else:
buttons[-1].focus()
return
if direction < 0 and current == 0:
self._focus_tabs()
return
new_index = current + direction
if 0 <= new_index < len(buttons):
buttons[new_index].focus()
def switch_tab(self, tab_id: str) -> None:
self.current_tab = tab_id
content = self.query_one("#content", Vertical)
content.remove_children()
if tab_id in self.BUTTON_DEFS:
content.mount(self._make_buttons_for(tab_id))
elif tab_id == "server_log":
content.mount(ServerLogWidget(self.app.api))
elif tab_id == "dir":
content.mount(DirectoryTree(self.working_dir, id="dir_tree"))
elif tab_id == "p_tree":
content.mount(PolicyTreeWidget(self.app.policies, self.app.devices))
elif tab_id == "settings":
content.mount(ThemeSelector())
else:
content.mount(Static(f"Unknown tab: {tab_id}"))
def on_tabs_tab_activated(self, event: Tabs.TabActivated) -> None:
self.switch_tab(event.tab.id)
def on_multi_agent_selector_agents_selected(
self, message: MultiAgentSelector.AgentsSelected
) -> None:
"""Handle selected agents from AgentSelector."""
global _APP_RESTART_REASON
selected_agents = message.selected_agents
logger.info("Selected agents: %s", selected_agents)
# TODO: Implement actual handling of selected agents
_APP_RESTART_REASON = ("multi_agent_action", selected_agents)
self.app.exit()
def on_theme_selector_theme_selected(
self, message: ThemeSelector.ThemeSelected
) -> None:
"""Handle theme selection from ThemeSelector."""
global _APP_RESTART_REASON
_persist_user_theme(message.theme_name)
_APP_RESTART_REASON = ("restart",)
self.app.exit()
def on_agent_move_operations_operation_complete(
self, message: AgentMoveOperations.OperationComplete
) -> None:
"""Handle completion of agent move operation - show results."""
logger.info(
"Agent move operation completed: %s, %d successful, %d unsuccessful",
message.operation,
len(message.successful),
len(message.unsuccessful),
)
# Format results for display
successful_text = "\n".join(
[f"{agent.hostname}" for agent, _ in message.successful]
)
unsuccessful_text = "\n".join(
[f"{agent.hostname}: {error}" for agent, error in message.unsuccessful]
)
# Remove the operations widget
try:
ops_widget = self.query_one(AgentMoveOperations)
ops_widget.remove()
except Exception:
pass
# Show results
self.query_one("#content", Vertical).mount(
ResultsDisplay(message.operation, successful_text, unsuccessful_text)
)
def on_results_display_go_back(self, message: ResultsDisplay.GoBack) -> None:
"""Handle back button from results display."""
try:
results_widget = self.query_one(ResultsDisplay)
results_widget.remove()
except Exception:
pass
# Return to main menu
self.app.pop_screen()
def on_policy_tree_widget_view_execution_history(
self, message: PolicyTreeWidget.ViewExecutionHistory
) -> None:
"""Handle request to view execution history for a device from tree view."""
logger.info("Viewing execution history for device: %s", message.device.hostname)
self.app.push_screen(ExecutionHistoryScreen([message.device]))
message.stop()
def on_policy_tree_widget_generate_otp(
self, message: PolicyTreeWidget.GenerateOTP
) -> None:
"""Handle request to generate OTP for a device from tree view."""
logger.info("Generating OTP for device: %s", message.device.hostname)
self.app.push_screen(OTPWorkflowScreen([message.device]))
message.stop()
def on_policy_tree_widget_toggle_enforcement(
self, message: PolicyTreeWidget.ToggleEnforcement
) -> None:
"""Handle request to toggle enforcement for a device from tree view."""
logger.info("Toggling enforcement for device: %s", message.device.hostname)
try:
from services.agenthandler import moveAgentToRelatedPolicy
from utils.configmanager import get_system_json
policy_relationship_map = get_system_json("POLICY_MAP_ENF_AUD", "{}")
# Determine current mode and toggle
if message.device.groupid in policy_relationship_map:
# Currently in enforcement, move to audit
result = moveAgentToRelatedPolicy(self.app.api, message.device, "audit")
mode = "audit"
else:
# Currently in audit, move to enforcement
result = moveAgentToRelatedPolicy(
self.app.api, message.device, "enforcement"
)
mode = "enforcement"
logger.info(f"Successfully toggled {message.device.hostname} to {mode}")
# Refresh data at the app level
self.app.refresh_data()
# Refresh the tree widget with new data
try:
tree_widget = self.query_one(PolicyTreeWidget)
tree_widget.refresh_data(self.app.policies, self.app.devices)
except:
pass
except Exception as e:
logger.error(
f"Failed to toggle enforcement for {message.device.hostname}: {e}"
)
self.app.bell()
message.stop()
def on_directory_tree_file_selected(
self, event: DirectoryTree.FileSelected
) -> None:
path = event.path
logger.debug("Directory file selected: %s", path)
try:
open_directory(str(path))
except Exception as exc:
logger.error("Failed to open %s: %s", path, exc)
self.app.bell()
def on_button_pressed(self, event: Button.Pressed) -> None:
button_id = event.button.id
logger.debug("Button pressed: %s", button_id)
match button_id:
case "move_agent_workflow_button":
self.app.push_screen(MoveAgentWorkflowScreen(self.app.devices))
event.stop()
case "otp_generate_button":
self.app.push_screen(OTPWorkflowScreen(self.app.devices))
event.stop()
case "find_quiet_button":
self.app.push_screen(
QuietAgentWorkflowScreen(self.app.api, self.app.policies)
)
event.stop()
return
case "otp_activities_button":
self.app.push_screen(OTPActivitiesScreen())
event.stop()
return
case "otp_revoke_button":
self.app.push_screen(OTPRevokeScreen())
event.stop()
return
case "policy_prep_button":
# Use the new TUI workflow screen instead of legacy
self.app.push_screen(
PolicyPrepWorkflowScreen(self.app.api, self.app.policies)
)
event.stop()
return
case _:
self.app.bell()
logger.warning("Unknown button pressed: %s", button_id)
return
# ---------------------------------------------------------------------------
# 2) APP
# ---------------------------------------------------------------------------
class Loxide(App[Message]):
api: AirlockAPIWrapper
working_dir: str
policies: Optional[list[Policy]]
devices: Optional[list[Agent]]
CSS = """
#logo {
width: 100%;
content-align: center middle;
text-align: center;
}
"""
BINDINGS = [
("q", "quit", "Quit"),
("f", "open_fe", "Launch Explorer"),
("r", "refresh", "Refresh"),
]
def __init__(self, api: AirlockAPIWrapper):
self._textual_theme = get_user_value("TEXTUAL_THEME", str, "textual-dark")
super().__init__()
self.api = api
wd = load_env("WORKING_DIR") or os.getcwd()
if not os.path.isdir(wd):
wd = os.getcwd()
self.working_dir = wd
# Initial data load
self.refresh_data()
def refresh_data(self) -> None:
"""Public method to refresh policies and devices from the API."""
try:
self.policies = [
Policy(**row.to_dict())
for _, row in self.api.policy_find_all().iterrows()
]
self.devices = [
Agent(**row.to_dict())
for _, row in self.api.agent_find_all().iterrows()
]
if self.policies and self.devices:
for agent in self.devices:
agent.enrich_with_policies(self.policies)
logger.debug(
f"Enriched {len(self.devices)} agents with policy information"
)
except Exception as exc:
logger.error("Failed to load policies/devices: %s", exc)
self.policies = None
self.devices = None
def on_mount(self, api: AirlockAPIWrapper) -> None:
self.register_theme(get_retro_terminal_theme())
self.register_theme(get_amber_terminal_theme())
self.theme = self._textual_theme
self.push_screen(MainMenuScreen())
def action_refresh(self) -> None:
self.refresh_data()
def action_quit(self) -> None:
global _APP_RESTART_REASON
_APP_RESTART_REASON = None
self.exit()
def action_open_fe(self) -> None:
"""Open the working directory in the OS file manager (footer binding)."""
path_to_open = self.working_dir or os.getcwd()
try:
open_directory(path_to_open)
except Exception as exc:
logger.error("Failed to open directory %s: %s", path_to_open, exc)
self.bell() # optional feedback
# ---------------------------------------------------------------------------
# 3) PUBLIC ENTRYPOINT - Updated to accept attach_notification_handler
# ---------------------------------------------------------------------------
def run_Loxide(api: AirlockAPIWrapper, attach_notification_handler=None) -> None:
global _APP_RESTART_REASON
base_dir = get_base_directory()
env_path = base_dir / ".env"
dotenv.load_dotenv(dotenv_path=env_path, override=True)
max_attempts = 5
attempts = 0
while attempts < max_attempts:
attempts += 1
logger.debug("Starting app loop iteration (attempt %d)", attempts)
_APP_RESTART_REASON = None
app = Loxide(api)
# Attach the notification handler if provided
if attach_notification_handler:
attach_notification_handler(app)
try:
app.run()
except SystemExit as exc:
if exc.code != 0:
logger.debug("Caught SystemExit from Textual: %s", exc)
raise
reason = _APP_RESTART_REASON
logger.debug("After app.run(), _APP_RESTART_REASON = %r", reason)
if not reason:
logger.debug("No restart reason, exiting loop")
break
if reason[0] == "restart":
logger.debug("Restarting app loop")
continue
if reason[0] == "multi_agent_action":
logger.info("Multi-agent action with selected agents: %s", reason[1])
continue
logger.error("Unknown restart reason: %r", reason)
break
# ---------------------------------------------------------------------------
# 4) MAIN FUNCTION - Updated to get and pass attach_notification_handler
# ---------------------------------------------------------------------------
def main():
irtang()
# Determine working directory, setup directory, configure logging, sent env, get API and URL if not already stored
setup()
# setup() now returns a function to attach the notification handler
attach_notification_handler = setup()
logger = logging.getLogger(__name__)
try:
@@ -66,7 +599,7 @@ def main():
base_url=str(url),
api_key=api_key,
)
run_Loxide(api)
run_Loxide(api, attach_notification_handler)
if __name__ == "__main__":
+639
View File
@@ -0,0 +1,639 @@
# Copyright (C) 2025 James Brotosky, Brandon Wickline
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU Affero General Public License as published
# by the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Affero General Public License for more details.
#
# You should have received a copy of the GNU Affero General Public License
# along with this program. If not, see <https://www.gnu.org/licenses/>.
from datetime import datetime, timedelta
import logging
import os
from typing import List
import pandas as pd
from textual.app import ComposeResult
from textual.binding import Binding
from textual.containers import Horizontal, Vertical
from textual.screen import Screen
from textual.widgets import (
Button,
DataTable,
Footer,
Header,
Label,
Select,
Static,
)
from models.agent import Agent
from models.execution import ExecutionHistoryRecord
from utils.configmanager import load_env
logger = logging.getLogger(__name__)
class ExecutionHistoryScreen(Screen):
"""
A screen for viewing and exporting execution history for selected agents.
This screen allows users to:
1. Select a start date and end date using dropdown selects
2. Fetch execution history for all selected agents
3. View the results in a DataTable
4. Export the results to CSV using a keybinding
Attributes:
agents (List[Agent]): List of agents to fetch execution history for
execution_data (pd.DataFrame): Combined execution history data
working_dir (str): Directory for CSV exports
"""
DEFAULT_CSS = """
ExecutionHistoryScreen {
align: center top;
}
#main_container {
width: 95%;
height: 1fr;
border: solid $primary;
padding: 1;
}
#title {
text-style: bold;
color: $text;
text-align: center;
margin-bottom: 1;
}
#date_container {
height: auto;
margin-bottom: 1;
}
#start_date_row, #end_date_row {
height: auto;
align-horizontal: left;
margin-bottom: 1;
}
.date_label {
width: 8;
margin-right: 1;
}
.date_selector {
width: 18;
margin: 0 1;
}
#quick_buttons_row {
height: auto;
align-horizontal: center;
margin-bottom: 1;
}
.quick_select_btn {
margin: 0 1;
}
#button_row {
height: auto;
align-horizontal: center;
margin-top: 1;
margin-bottom: 1;
}
Button {
margin: 0 1;
}
#status_label {
text-align: center;
color: $accent;
margin-bottom: 1;
}
#results_container {
height: 1fr;
display: none;
}
#results_button_row {
height: auto;
align-horizontal: center;
margin-bottom: 1;
}
#history_table {
height: 1fr;
border: solid $primary;
}
DataTable > .datatable--header {
text-style: bold;
background: $primary 20%;
}
"""
BINDINGS = [
Binding("escape", "close_screen", "Close"),
Binding("e", "export_csv", "Export CSV"),
Binding("q", "close_screen", "Quit"),
]
def __init__(self, agents: List[Agent]):
"""
Initialize the ExecutionHistoryScreen.
Args:
agents (List[Agent]): List of agents to fetch execution history for
"""
super().__init__()
self.agents = agents
self.execution_data = pd.DataFrame()
self.working_dir = load_env("WORKING_DIR") or os.getcwd()
# Generate dropdown options
today = datetime.now().date()
# Month options - format is (display_text, value)
self.month_options = [
("January", "01"),
("February", "02"),
("March", "03"),
("April", "04"),
("May", "05"),
("June", "06"),
("July", "07"),
("August", "08"),
("September", "09"),
("October", "10"),
("November", "11"),
("December", "12"),
]
# Day options (1-31) - format is (display_text, value)
self.day_options = [(f"{i}", f"{i:02d}") for i in range(1, 32)]
# Year options (current year back 5 years) - format is (display_text, value)
current_year = today.year
self.year_options = [
(str(year), str(year)) for year in range(current_year, current_year - 6, -1)
]
# Default dates: last 30 days
start_date = today - timedelta(days=30)
self.start_month = f"{start_date.month:02d}"
self.start_day = f"{start_date.day:02d}"
self.start_year = str(start_date.year)
self.end_month = f"{today.month:02d}"
self.end_day = f"{today.day:02d}"
self.end_year = str(today.year)
def compose(self) -> ComposeResult:
"""Build the UI layout."""
yield Header(show_clock=True, icon="📊")
with Vertical(id="main_container"):
title_text = f"Execution History - {len(self.agents)} Agent(s)"
yield Static(title_text, id="title")
# Date selection area
with Vertical(id="date_container"):
yield Label("Select Date Range:")
# Start date row
with Horizontal(id="start_date_row"):
yield Label("From:", classes="date_label")
yield Select(
options=self.month_options,
value=self.start_month,
id="start_month_select",
classes="date_selector",
)
yield Select(
options=self.day_options,
value=self.start_day,
id="start_day_select",
classes="date_selector",
)
yield Select(
options=self.year_options,
value=self.start_year,
id="start_year_select",
classes="date_selector",
)
# End date row
with Horizontal(id="end_date_row"):
yield Label("To:", classes="date_label")
yield Select(
options=self.month_options,
value=self.end_month,
id="end_month_select",
classes="date_selector",
)
yield Select(
options=self.day_options,
value=self.end_day,
id="end_day_select",
classes="date_selector",
)
yield Select(
options=self.year_options,
value=self.end_year,
id="end_year_select",
classes="date_selector",
)
# Quick select buttons
with Horizontal(id="quick_buttons_row"):
yield Button(
"1 Day",
id="quick_1day",
classes="quick_select_btn",
variant="default",
)
yield Button(
"1 Week",
id="quick_1week",
classes="quick_select_btn",
variant="default",
)
yield Button(
"30 Days",
id="quick_30days",
classes="quick_select_btn",
variant="default",
)
# Buttons
with Horizontal(id="button_row"):
yield Button("Fetch History", id="fetch_btn", variant="primary")
yield Button("Close", id="close_btn", variant="error")
# Status
yield Static(
"Select date range and click 'Fetch History'", id="status_label"
)
# Results container (hidden initially, shown after fetch)
with Vertical(id="results_container"):
with Horizontal(id="results_button_row"):
yield Button("Export CSV", id="export_btn", variant="success")
yield Button("Back", id="back_btn", variant="default")
yield DataTable(id="history_table")
yield Footer()
def on_mount(self) -> None:
"""Initialize the table when screen is mounted."""
table = self.query_one("#history_table", DataTable)
table.cursor_type = "row"
table.zebra_stripes = True
# Initially empty - will populate after fetch
logger.info(f"ExecutionHistoryScreen mounted with {len(self.agents)} agents")
def on_select_changed(self, event: Select.Changed) -> None:
"""Handle date selection changes."""
select_id = event.select.id
if select_id == "start_month_select":
self.start_month = event.value
logger.debug(f"Start month changed to: {self.start_month}")
elif select_id == "start_day_select":
self.start_day = event.value
logger.debug(f"Start day changed to: {self.start_day}")
elif select_id == "start_year_select":
self.start_year = event.value
logger.debug(f"Start year changed to: {self.start_year}")
elif select_id == "end_month_select":
self.end_month = event.value
logger.debug(f"End month changed to: {self.end_month}")
elif select_id == "end_day_select":
self.end_day = event.value
logger.debug(f"End day changed to: {self.end_day}")
elif select_id == "end_year_select":
self.end_year = event.value
logger.debug(f"End year changed to: {self.end_year}")
def _set_quick_date_range(self, days: int) -> None:
"""Set the date range based on quick select button."""
today = datetime.now().date()
start_date = today - timedelta(days=days)
# Update internal values
self.start_month = f"{start_date.month:02d}"
self.start_day = f"{start_date.day:02d}"
self.start_year = str(start_date.year)
self.end_month = f"{today.month:02d}"
self.end_day = f"{today.day:02d}"
self.end_year = str(today.year)
# Update the Select widgets
try:
self.query_one("#start_month_select", Select).value = self.start_month
self.query_one("#start_day_select", Select).value = self.start_day
self.query_one("#start_year_select", Select).value = self.start_year
self.query_one("#end_month_select", Select).value = self.end_month
self.query_one("#end_day_select", Select).value = self.end_day
self.query_one("#end_year_select", Select).value = self.end_year
self.app.notify(
f"Date range set to last {days} day(s)",
severity="information",
timeout=2,
)
logger.info(f"Quick select: Set date range to last {days} days")
except Exception as e:
logger.error(f"Failed to update date selects: {e}")
def _show_date_selection(self) -> None:
"""Show the date selection view and hide results."""
try:
self.query_one("#date_container").styles.display = "block"
self.query_one("#button_row").styles.display = "block"
self.query_one("#status_label").styles.display = "block"
self.query_one("#results_container").styles.display = "none"
except Exception as e:
logger.error(f"Failed to show date selection: {e}")
def _show_results(self) -> None:
"""Hide date selection view and show results."""
try:
self.query_one("#date_container").styles.display = "none"
self.query_one("#button_row").styles.display = "none"
self.query_one("#status_label").styles.display = "none"
self.query_one("#results_container").styles.display = "block"
except Exception as e:
logger.error(f"Failed to show results: {e}")
def on_button_pressed(self, event: Button.Pressed) -> None:
"""Handle button clicks."""
if event.button.id == "fetch_btn":
self._fetch_execution_history()
elif event.button.id == "export_btn":
self._export_to_csv()
elif event.button.id == "close_btn":
self.app.pop_screen()
elif event.button.id == "back_btn":
self._show_date_selection()
elif event.button.id == "quick_1day":
self._set_quick_date_range(days=1)
elif event.button.id == "quick_1week":
self._set_quick_date_range(days=7)
elif event.button.id == "quick_30days":
self._set_quick_date_range(days=30)
def _fetch_execution_history(self) -> None:
"""Fetch execution history for all selected agents."""
status_label = self.query_one("#status_label", Static)
status_label.update("⏳ Fetching execution history...")
# Disable buttons during fetch
fetch_btn = self.query_one("#fetch_btn", Button)
export_btn = self.query_one("#export_btn", Button)
fetch_btn.disabled = True
export_btn.disabled = True
api = self.app.api
all_history = []
try:
# Construct dates from dropdowns
start_date_str = f"{self.start_year}-{self.start_month}-{self.start_day}"
end_date_str = f"{self.end_year}-{self.end_month}-{self.end_day}"
# Validate dates
try:
start_dt = datetime.strptime(start_date_str, "%Y-%m-%d")
end_dt = datetime.strptime(end_date_str, "%Y-%m-%d")
except ValueError as e:
status_label.update(f"❌ Invalid date: {str(e)}")
fetch_btn.disabled = False
export_btn.disabled = False
self.app.notify(f"Invalid date selected: {str(e)}", severity="error")
return
if start_dt > end_dt:
status_label.update("❌ Error: Start date must be before end date")
fetch_btn.disabled = False
export_btn.disabled = False
return
# Fetch history for each agent
for i, agent in enumerate(self.agents):
try:
status_label.update(
f"⏳ Fetching history for {agent.hostname} ({i+1}/{len(self.agents)})..."
)
# Call API - note the API expects 'dateto' first, then 'datefrom'
history = api.history_execution(
today=end_date_str,
date_selected=start_date_str,
agent_name=agent.hostname,
)
if history:
# Add agent hostname to each record for identification
for record in history:
record["agent_hostname"] = agent.hostname
all_history.extend(history)
logger.info(
f"Fetched {len(history)} records for {agent.hostname}"
)
else:
logger.info(f"No history found for {agent.hostname}")
except Exception as e:
logger.error(f"Failed to fetch history for {agent.hostname}: {e}")
self.app.notify(
f"Warning: Failed to fetch history for {agent.hostname}",
severity="warning",
)
# Convert to DataFrame
if all_history:
status_label.update(
"⏳ Enriching execution data with hash information..."
)
# Normalize field names (handle API typos)
for record in all_history:
if "policver" in record and "policyver" not in record:
record["policyver"] = record.pop("policver")
# Convert dict records to ExecutionHistoryRecord objects
execution_records = []
for record in all_history:
try:
execution_records.append(ExecutionHistoryRecord(**record))
except TypeError as e:
logger.warning(f"Failed to create ExecutionHistoryRecord: {e}")
# If it fails, just keep the dict
continue
# Enrich with hash data if we have ExecutionHistoryRecord objects
if execution_records:
try:
enriched_records = ExecutionHistoryRecord.enrich_with_hashes(
api, execution_records
)
logger.info(
f"Enriched {len(enriched_records)} records with hash data"
)
# Convert back to DataFrame
self.execution_data = pd.DataFrame(
[r.__dict__ for r in enriched_records]
)
# Flatten hash_obj if present
if (
not self.execution_data.empty
and "hash_obj" in self.execution_data.columns
):
hash_df = self.execution_data["hash_obj"].apply(
lambda h: (
h.to_dict() if h and hasattr(h, "to_dict") else {}
)
)
self.execution_data = pd.concat(
[
self.execution_data.drop(columns=["hash_obj"]),
hash_df,
],
axis=1,
)
except Exception as e:
logger.warning(f"Failed to enrich with hashes: {e}")
# Fall back to plain DataFrame
self.execution_data = pd.DataFrame(all_history)
else:
# If we couldn't create any ExecutionHistoryRecord objects, just use raw data
self.execution_data = pd.DataFrame(all_history)
self._populate_table()
self._show_results() # Switch to results view
self.app.notify(
f"Successfully loaded {len(self.execution_data)} records",
severity="information",
)
else:
status_label.update(
"ℹ️ No execution history found for selected agents/dates"
)
self.app.notify("No execution history found", severity="information")
self.execution_data = pd.DataFrame()
except Exception as e:
logger.error(f"Error fetching execution history: {e}")
status_label.update(f"❌ Error: {str(e)}")
self.app.notify(f"Failed to fetch history: {str(e)}", severity="error")
finally:
# Re-enable buttons
fetch_btn.disabled = False
export_btn.disabled = False
def _populate_table(self) -> None:
"""Populate the DataTable with execution history data."""
table = self.query_one("#history_table", DataTable)
table.clear(columns=True)
if self.execution_data.empty:
return
# Define preferred column order (your specified order)
preferred_order = [
"policyname",
"policyver",
"hostname",
"username",
"publisher",
"filename",
"pprocess",
"gprocess",
"sha256",
"commandline",
"agent_hostname", # Our custom field
]
# Get available columns in preferred order, then add any remaining columns
available_cols = []
for col in preferred_order:
if col in self.execution_data.columns:
available_cols.append(col)
# Add any remaining columns not in preferred order
for col in self.execution_data.columns:
if col not in available_cols:
available_cols.append(col)
# Add columns to table
for col in available_cols:
table.add_column(col, key=col)
# Add rows
for idx, row in self.execution_data.iterrows():
row_data = []
for col in available_cols:
value = row[col]
# Convert to string, handle None/NaN
if pd.isna(value):
row_data.append("")
else:
row_data.append(str(value))
table.add_row(*row_data, key=str(idx))
logger.info(f"Populated table with {len(self.execution_data)} rows")
def _export_to_csv(self) -> None:
"""Export the current execution data to CSV."""
if self.execution_data.empty:
self.app.notify("No data to export", severity="warning")
return
try:
# Create filename with timestamp
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
filename = f"execution_history_{timestamp}.csv"
filepath = os.path.join(self.working_dir, filename)
# Export to CSV
self.execution_data.to_csv(filepath, index=False, encoding="utf-8-sig")
self.app.notify(
f"✅ Exported {len(self.execution_data)} records to: {filepath}",
severity="information",
timeout=5,
)
logger.info(f"Exported execution history to: {filepath}")
except Exception as e:
logger.error(f"Failed to export CSV: {e}")
self.app.notify(f"Failed to export CSV: {str(e)}", severity="error")
def action_export_csv(self) -> None:
"""Keybinding action to export CSV."""
self._export_to_csv()
def action_close_screen(self) -> None:
"""Close this screen and return to previous."""
self.app.pop_screen()
+380 -59
View File
@@ -165,6 +165,18 @@ class PolicyPrepWorkflowScreen(Screen):
# Track if we're navigating with keyboard (to prevent selection)
self._keyboard_navigation = False
# Tab review tracking for Step 5 (First Review)
self.approved_tab_reviewed = False
self.needs_review_tab_reviewed = False
# Tab review tracking for Step 6 (Path Review)
self.paths_tab_reviewed = False
self.publishers_tab_reviewed = False
# Lock to prevent concurrent table creation
self._creating_review_table = False
self._creating_path_table = False
def compose(self) -> ComposeResult:
"""Build the UI layout for the workflow screen."""
yield Header(show_clock=True, icon="⚙️")
@@ -191,7 +203,7 @@ class PolicyPrepWorkflowScreen(Screen):
def on_mount(self) -> None:
"""Initialize the screen when mounted."""
self._update_checklist()
self._show_source_policy_selection()
self._show_introduction()
def watch_workflow_stage(self, old_value: str, new_value: str) -> None:
"""React to workflow stage changes."""
@@ -348,6 +360,86 @@ class PolicyPrepWorkflowScreen(Screen):
step8.styles.text_style = "dim"
col2.mount(step8)
def _show_introduction(self) -> None:
"""Show workflow introduction and overview."""
self.workflow_stage = "introduction"
content = self.query_one("#content_area", Vertical)
content.remove_children()
# Title
title = Static("Welcome to Policy Preparation Workflow")
title.styles.margin = (1, 1)
title.styles.text_style = "bold"
title.styles.text_align = "center"
content.mount(title)
# Description
description = Static(
"This workflow will help you:\n"
" • Fetch execution history from selected policies\n"
" • Review and approve safe executions\n"
" • Calculate efficient path exclusions\n"
" • Generate publisher trust rules\n"
" • Apply changes to your destination policy"
)
description.styles.margin = (1, 2)
content.mount(description)
# Process steps
steps_title = Static("The Process:")
steps_title.styles.margin = (1, 2, 0, 2)
steps_title.styles.text_style = "bold"
content.mount(steps_title)
steps = Static(
" 📋 Step 1: Select source policies (data collection)\n"
" 🎯 Step 2: Select destination policy (where changes go)\n"
" 📝 Step 3: Select destination allowlist\n"
" 📊 Step 4: Fetch execution data (may take 1-2 minutes)\n"
" ✅ Step 5: Review approved/needs review executions\n"
" 📁 Step 6: Review path exclusions and publishers\n"
" 🔍 Step 7: Preview changes before applying\n"
" 🚀 Step 8: Liftoff - Apply to production"
)
steps.styles.margin = (0, 2)
content.mount(steps)
# Time estimate
estimate = Static("⏱️ Estimated Time: 15-30 minutes depending on data size")
estimate.styles.margin = (1, 2)
estimate.styles.color = "cyan"
content.mount(estimate)
# Tips
tips_title = Static("💡 Tips:")
tips_title.styles.margin = (1, 2, 0, 2)
tips_title.styles.text_style = "bold"
content.mount(tips_title)
tips = Static(
" • Start with a test policy first\n"
" • Review carefully - changes affect all agents\n"
" • Use path rules when possible (more efficient)\n"
" • Publishers are powerful - use cautiously"
)
tips.styles.margin = (0, 2)
tips.styles.color = "yellow"
content.mount(tips)
# Buttons
button_container = Horizontal()
button_container.styles.margin = (2, 2)
button_container.styles.align = ("center", "middle")
content.mount(button_container)
continue_btn = Button(
"Continue to Policy Selection", id="start_workflow", variant="success"
)
cancel_btn = Button("Cancel", id="cancel_workflow", variant="default")
button_container.mount(continue_btn)
button_container.mount(cancel_btn)
def _show_source_policy_selection(self) -> None:
"""Show the source policy selection screen."""
self.workflow_stage = "select_source"
@@ -366,7 +458,7 @@ class PolicyPrepWorkflowScreen(Screen):
table.zebra_stripes = True
# Add columns - checkbox first, then data columns
table.add_columns("", "Name", "ID", "Parent")
table.add_columns("", "Name", "ID", "Parent")
# Sort policies by name for easier selection
sorted_policies = sorted(self.policies, key=lambda p: p.name.lower())
@@ -376,7 +468,7 @@ class PolicyPrepWorkflowScreen(Screen):
# Skip parent policies
if policy.parent == "global-policy-settings":
continue
checkbox = "" # All start unchecked
checkbox = "" # All start unchecked
table.add_row(
checkbox,
policy.name,
@@ -739,7 +831,7 @@ class PolicyPrepWorkflowScreen(Screen):
def _show_fetch_results(self) -> None:
"""Show the results of data fetching."""
logger.info("=== _show_fetch_results called ===")
logger.debug("=== _show_fetch_results called ===")
self.workflow_stage = "first_review"
content = self.query_one("#content_area", Vertical)
content.remove_children()
@@ -792,15 +884,39 @@ class PolicyPrepWorkflowScreen(Screen):
logger.info("Mounted tab buttons")
# Show approved table by default
logger.info("About to call _show_review_table('approved')")
logger.debug("About to call _show_review_table('approved')")
self._show_review_table("approved")
logger.info("=== _show_fetch_results complete ===")
logger.debug("=== _show_fetch_results complete ===")
def _show_review_table(self, table_type: str) -> None:
"""Show an editable DataTable for reviewing executions."""
logger.info(f"=== _show_review_table called with type: {table_type} ===")
# Prevent concurrent execution
if self._creating_review_table:
logger.warning(
f"Already creating review table, ignoring duplicate call for {table_type}"
)
return
self._creating_review_table = True
try:
self._show_review_table_impl(table_type)
finally:
self._creating_review_table = False
def _show_review_table_impl(self, table_type: str) -> None:
"""Internal implementation of _show_review_table."""
logger.debug(f"_show_review_table called with type: {table_type}")
content = self.query_one("#content_area", Vertical)
# Mark tab as reviewed
if table_type == "approved":
self.approved_tab_reviewed = True
logger.debug("Marked approved tab as reviewed")
else:
self.needs_review_tab_reviewed = True
logger.debug("Marked needs_review tab as reviewed")
# Determine which dataframe and table ID to show
if table_type == "approved":
df = self.approved_df
@@ -841,6 +957,18 @@ class PolicyPrepWorkflowScreen(Screen):
except Exception as e:
logger.debug(f"Error removing existing tables: {e}")
# Remove existing instruction and help text (they accumulate without removal)
# Remove ALL Static widgets - they're just text that needs to be replaced
try:
existing_statics = content.query("Static")
logger.debug(
f"Found {len(existing_statics)} existing Static widgets to remove"
)
for static in existing_statics:
static.remove()
except Exception as e:
logger.debug(f"Error removing Static widgets: {e}")
# Force a refresh to ensure removals are processed
try:
content.refresh()
@@ -850,7 +978,7 @@ class PolicyPrepWorkflowScreen(Screen):
# Note: We no longer remove review_controls or review_continue_container
# They are reused between tabs to avoid DuplicateIds errors
logger.info(
logger.debug(
f"DataFrame for {table_type}: {'empty' if df is None or df.empty else f'{len(df)} rows'}"
)
@@ -861,13 +989,13 @@ class PolicyPrepWorkflowScreen(Screen):
logger.info(f"No data for {table_type}, mounted empty message")
return
# Instructions
# Instructions (no ID needed - we remove all Statics anyway)
instruction = Static(title)
instruction.styles.margin = (1, 1)
instruction.styles.text_style = "bold"
content.mount(instruction)
# Help text
# Help text (no ID needed)
help_text = Static(
"Click to toggle, 'r' for range select (click start, press 'r', click end)\n"
"Space to toggle cursor row, 'd' to delete, 'a' select all, arrows navigate"
@@ -876,6 +1004,19 @@ class PolicyPrepWorkflowScreen(Screen):
help_text.styles.text_style = "dim"
content.mount(help_text)
# CRITICAL: Check if table already exists in content (should not happen after removal above)
try:
existing_check = content.query_one(f"#{table_id}", DataTable)
if existing_check:
logger.error(
f"Table {table_id} STILL EXISTS after removal! This should not happen."
)
# Don't create a new one - just return
return
except Exception:
# Good - table doesn't exist, proceed with creation
pass
# Create the review table
review_table = DataTable(id=table_id)
review_table.styles.height = "40vh" # Increased since we removed button rows
@@ -899,15 +1040,15 @@ class PolicyPrepWorkflowScreen(Screen):
if available_cols:
# Add checkbox column first
review_table.add_columns("", *available_cols)
review_table.add_columns("", *available_cols)
# Add rows with row keys for tracking
for idx, row in df.iterrows():
checkbox = "" # All start unchecked
checkbox = "" # All start unchecked
row_data = [str(row.get(col, "")) for col in available_cols]
review_table.add_row(checkbox, *row_data, key=str(idx))
logger.info(f"About to mount {table_id}")
logger.debug(f"About to mount {table_id}")
# Final safety check - make sure no table with this ID exists before mounting
try:
@@ -923,7 +1064,7 @@ class PolicyPrepWorkflowScreen(Screen):
pass
content.mount(review_table)
logger.info(f"Successfully mounted {table_id} with {len(df)} rows")
logger.debug(f"Successfully mounted {table_id} with {len(df)} rows")
# Row count display only (removed Select All, Clear, Delete buttons)
try:
@@ -1044,13 +1185,19 @@ class PolicyPrepWorkflowScreen(Screen):
content = self.query_one("#content_area", Vertical)
content.remove_children()
# Clear message
# Add spacer to push text to bottom
spacer = Static("")
spacer.styles.height = "1fr"
content.mount(spacer)
# Loading message at bottom (above checklist)
status = Static(
"Building path exclusions and publisher lists...\n\n"
"Building path exclusions and publisher lists...\n"
"This may take a moment for large datasets."
)
status.styles.margin = (2, 1)
status.styles.margin = (1, 1)
status.styles.text_align = "center"
status.styles.color = "cyan"
content.mount(status)
# Force UI refresh to show the loading screen
@@ -1311,8 +1458,8 @@ class PolicyPrepWorkflowScreen(Screen):
logger.info(
f"=== _calculate_paths called with path_exclusion_constant={path_exclusion_constant} ==="
)
logger.info(f"Input DataFrame: {len(df)} rows")
logger.info(f"Columns: {list(df.columns) if not df.empty else 'empty'}")
logger.debug(f"Input DataFrame: {len(df)} rows")
logger.debug(f"Columns: {list(df.columns) if not df.empty else 'empty'}")
if df.empty:
logger.warning("Input DataFrame is empty")
@@ -1327,7 +1474,7 @@ class PolicyPrepWorkflowScreen(Screen):
haslcp = self._split_filepaths_grouped(df, path_exclusion_constant, "filename")
haslcp = haslcp.drop_duplicates()
logger.info(f"After split_filepaths_grouped: {len(haslcp)} rows")
logger.debug(f"After split_filepaths_grouped: {len(haslcp)} rows")
# Filter forbidden paths
badpathparts = get_system_list("BAD_PATH_PARTS")
@@ -1337,7 +1484,7 @@ class PolicyPrepWorkflowScreen(Screen):
forbidden_pattern, case=False, na=False, regex=True
)
logger.info(
logger.debug(
f"Removing forbidden filepaths: {forbidden_lcfp.sum()} paths filtered"
)
lcp_not_forbidden = haslcp[~forbidden_lcfp].copy()
@@ -1347,7 +1494,7 @@ class PolicyPrepWorkflowScreen(Screen):
)
lcp_not_forbidden = haslcp.copy()
logger.info(f"After forbidden filtering: {len(lcp_not_forbidden)} rows")
logger.debug(f"After forbidden filtering: {len(lcp_not_forbidden)} rows")
# Select relevant columns
if "policyname" in lcp_not_forbidden.columns:
@@ -1399,11 +1546,11 @@ class PolicyPrepWorkflowScreen(Screen):
lcp_not_forbidden_review = lcp_not_forbidden_review[
lcp_not_forbidden_review["unique_sha256_count"] >= min_files_for_path
]
logger.info(
logger.debug(
f"After MIN_FILES_FOR_PATH filter ({min_files_for_path}): {len(lcp_not_forbidden_review)} rows (removed {before_filter - len(lcp_not_forbidden_review)})"
)
logger.info(
logger.debug(
f"Final result: {len(lcp_not_forbidden_review)} rows with columns: {list(lcp_not_forbidden_review.columns)}"
)
@@ -1441,7 +1588,7 @@ class PolicyPrepWorkflowScreen(Screen):
def _show_path_results(self) -> None:
"""Show the results of path building."""
logger.info("=== _show_path_results called ===")
logger.debug("=== _show_path_results called ===")
self.workflow_stage = "second_review"
content = self.query_one("#content_area", Vertical)
content.remove_children()
@@ -1490,15 +1637,39 @@ class PolicyPrepWorkflowScreen(Screen):
logger.info("Mounted tab buttons")
# Show paths table by default
logger.info("About to call _show_path_review_table('paths')")
logger.debug("About to call _show_path_review_table('paths')")
self._show_path_review_table("paths")
logger.info("=== _show_path_results complete ===")
logger.debug("=== _show_path_results complete ===")
def _show_path_review_table(self, table_type: str) -> None:
"""Show an editable DataTable for reviewing paths/publishers."""
logger.info(f"=== _show_path_review_table called with type: {table_type} ===")
# Prevent concurrent execution
if self._creating_path_table:
logger.warning(
f"Already creating path table, ignoring duplicate call for {table_type}"
)
return
self._creating_path_table = True
try:
self._show_path_review_table_impl(table_type)
finally:
self._creating_path_table = False
def _show_path_review_table_impl(self, table_type: str) -> None:
"""Internal implementation of _show_path_review_table."""
logger.debug(f"_show_path_review_table called with type: {table_type}")
content = self.query_one("#content_area", Vertical)
# Mark tab as reviewed (only paths and publishers, not remaining)
if table_type == "paths":
self.paths_tab_reviewed = True
logger.debug("Marked paths tab as reviewed")
elif table_type == "publishers":
self.publishers_tab_reviewed = True
logger.debug("Marked publishers tab as reviewed")
# Determine which dataframe to show
if table_type == "paths":
# Combine primary and secondary paths for review
@@ -1608,6 +1779,18 @@ class PolicyPrepWorkflowScreen(Screen):
except Exception as e:
logger.debug(f"Error removing existing controls: {e}")
# Remove existing instruction and help text (they accumulate without removal)
# Remove ALL Static widgets - they're just text that needs to be replaced
try:
existing_statics = content.query("Static")
logger.debug(
f"Found {len(existing_statics)} existing Static widgets to remove"
)
for static in existing_statics:
static.remove()
except Exception as e:
logger.debug(f"Error removing Static widgets: {e}")
# Force refresh to ensure removals complete
try:
content.refresh()
@@ -1620,13 +1803,13 @@ class PolicyPrepWorkflowScreen(Screen):
content.mount(empty_msg)
return
# Instructions
# Instructions (no ID needed)
instruction = Static(title)
instruction.styles.margin = (1, 1)
instruction.styles.text_style = "bold"
content.mount(instruction)
# Help text (different for remaining hashes)
# Help text (different for remaining hashes, no ID needed)
if table_type != "remaining":
help_text = Static(
"Click to toggle, 'r' for range select (click start, press 'r', click end)\n"
@@ -1653,11 +1836,11 @@ class PolicyPrepWorkflowScreen(Screen):
available_cols = [col for col in columns if col in df.columns]
if available_cols:
# Add checkbox column first
review_table.add_columns("", *available_cols)
review_table.add_columns("", *available_cols)
# Add rows with row keys for tracking
for idx, row in df.iterrows():
checkbox = "" # All start unchecked
checkbox = "" # All start unchecked
row_data = []
for col in available_cols:
value = row.get(col, "")
@@ -1669,7 +1852,7 @@ class PolicyPrepWorkflowScreen(Screen):
row_data.append(str(value))
review_table.add_row(checkbox, *row_data, key=str(idx))
logger.info(f"About to mount {table_id}")
logger.debug(f"About to mount {table_id}")
# Final safety check before mounting table
try:
@@ -1682,7 +1865,7 @@ class PolicyPrepWorkflowScreen(Screen):
pass
content.mount(review_table)
logger.info(f"Successfully mounted {table_id}")
logger.debug(f"Successfully mounted {table_id}")
# Row count display only (removed Select All, Clear, Delete buttons)
if table_type != "remaining":
@@ -1999,53 +2182,166 @@ class PolicyPrepWorkflowScreen(Screen):
"""Perform the actual application of changes."""
try:
results = []
errors = []
# Apply path exclusions to policy
if self.destination_policy and self.primary_paths_df is not None:
# This would call the actual API methods
results.append("Applied path exclusions to policy")
# Apply path exclusions to policy (primary + secondary)
if self.destination_policy:
path_rules = []
# Process primary paths
if (
self.primary_paths_df is not None
and not self.primary_paths_df.empty
):
logger.info(
f"Processing {len(self.primary_paths_df)} primary paths"
)
for _, row in self.primary_paths_df.groupby(
["longestcfp", "file_extension"]
):
path = row.iloc[0]["longestcfp"]
ext = row.iloc[0]["file_extension"]
# Format: C:\Path\**.ext
path_rule = f"{path}\\**{ext}"
path_rules.append(path_rule)
# Process secondary paths
if (
self.secondary_paths_df is not None
and not self.secondary_paths_df.empty
):
logger.info(
f"Processing {len(self.secondary_paths_df)} secondary paths"
)
for _, row in self.secondary_paths_df.groupby(
["longestcfp", "file_extension"]
):
path = row.iloc[0]["longestcfp"]
ext = row.iloc[0]["file_extension"]
path_rule = f"{path}\\**{ext}"
path_rules.append(path_rule)
# Apply path rules to policy
if path_rules:
try:
logger.info(
f"Applying {len(path_rules)} path exclusions to policy {self.destination_policy.name}"
)
response = self.api.policy_add_path_exclusions(
str(self.destination_policy.groupid), path_rules
)
results.append(
f"✓ Added {len(path_rules)} path exclusions to policy"
)
logger.info(f"Path exclusions applied successfully: {response}")
except Exception as e:
error_msg = f"✗ Failed to add path exclusions: {str(e)}"
errors.append(error_msg)
logger.error(error_msg, exc_info=True)
# Apply publishers to policy
if self.destination_policy and self.publishers_df is not None:
# This would call the actual API methods
results.append("Applied approved publishers to policy")
if (
self.destination_policy
and self.publishers_df is not None
and not self.publishers_df.empty
):
try:
publishers = self.publishers_df["publisher"].unique().tolist()
logger.info(
f"Applying {len(publishers)} publishers to policy {self.destination_policy.name}"
)
response = self.api.policy_add_publishers(
str(self.destination_policy.groupid), publishers
)
results.append(
f"✓ Added {len(publishers)} trusted publishers to policy"
)
logger.info(f"Publishers applied successfully: {response}")
except Exception as e:
error_msg = f"✗ Failed to add publishers: {str(e)}"
errors.append(error_msg)
logger.error(error_msg, exc_info=True)
# Apply hashes to allowlist
if self.destination_allowlist and self.approved_df is not None:
# This would call the actual API methods
results.append("Applied approved hashes to allowlist")
if (
self.destination_allowlist
and self.approved_df is not None
and not self.approved_df.empty
):
try:
# Get unique hashes
hashes = self.approved_df["sha256"].unique().tolist()
logger.info(
f"Applying {len(hashes)} hashes to allowlist {self.destination_allowlist.name}"
)
response = self.api.hash_add_to_allowlist(
str(self.destination_allowlist.applicationid), hashes
)
results.append(
f"✓ Added {len(hashes):,} approved hashes to allowlist"
)
logger.info(f"Hashes applied successfully: {response}")
except Exception as e:
error_msg = f"✗ Failed to add hashes: {str(e)}"
errors.append(error_msg)
logger.error(error_msg, exc_info=True)
self._show_completion(results)
# Show completion with both results and errors
all_results = results + errors
self._show_completion(all_results, has_errors=len(errors) > 0)
except Exception as e:
logger.error(f"Failed to apply changes: {e}", exc_info=True)
self.app.notify(f"Failed to apply changes: {str(e)}", severity="error")
logger.error(f"Critical failure in _perform_apply: {e}", exc_info=True)
self.app.notify(f"Critical failure: {str(e)}", severity="error")
self._show_test_screen()
def _show_completion(self, results: List[str]) -> None:
def _show_completion(self, results: List[str], has_errors: bool = False) -> None:
"""Show completion screen."""
self.workflow_stage = "complete"
content = self.query_one("#content_area", Vertical)
content.remove_children()
summary = Static(
"Policy Preparation Complete!\n\n"
"The following changes have been applied:"
)
# Title depends on whether there were errors
if has_errors:
title_text = "Policy Preparation Completed with Errors\n\n" "Results:"
title_color = "yellow"
else:
title_text = (
"Policy Preparation Complete!\n\n"
"The following changes have been applied:"
)
title_color = "green"
summary = Static(title_text)
summary.styles.margin = (1, 1)
summary.styles.text_style = "bold"
summary.styles.color = title_color
content.mount(summary)
for result in results:
result_widget = Static(f" {result}")
result_widget.styles.margin = (0, 2)
# Color based on success/failure
if result.startswith(""):
result_widget.styles.color = "green"
elif result.startswith(""):
result_widget.styles.color = "red"
content.mount(result_widget)
# Final message
final = Static(
f"\nPolicy '{self.destination_policy.name}' is now ready for enforcement!"
)
final.styles.margin = (2, 1)
final.styles.color = "green"
if has_errors:
final = Static(
f"\n⚠️ Policy '{self.destination_policy.name}' was partially updated.\n"
"Please review errors above and retry failed operations manually."
)
final.styles.margin = (2, 1)
final.styles.color = "yellow"
else:
final = Static(
f"\n✅ Policy '{self.destination_policy.name}' is now ready for enforcement!"
)
final.styles.margin = (2, 1)
final.styles.color = "green"
content.mount(final)
# Done button
@@ -2078,7 +2374,7 @@ class PolicyPrepWorkflowScreen(Screen):
)
# Determine if this row should be checked
is_selected = row_key_str in selected_keys
checkbox = "☑️" if is_selected else ""
checkbox = "" if is_selected else ""
# Update the checkbox cell (first column, index 0)
try:
@@ -2315,8 +2611,15 @@ class PolicyPrepWorkflowScreen(Screen):
"""Handle button presses."""
button_id = event.button.id
# Introduction screen buttons
if button_id == "start_workflow":
self._show_source_policy_selection()
elif button_id == "cancel_workflow":
self.app.pop_screen()
# Source policy selection buttons
if button_id == "select_none_source":
elif button_id == "select_none_source":
self.selected_source_policy_ids.clear()
# Refresh checkbox display
self._refresh_table_checkboxes(
@@ -2434,6 +2737,15 @@ class PolicyPrepWorkflowScreen(Screen):
)
elif button_id == "continue_from_review":
# Check if both tabs have been reviewed
if not self.approved_tab_reviewed or not self.needs_review_tab_reviewed:
self.app.notify(
"Please review both 'Approved' and 'Needs Review' tabs before continuing.",
severity="warning",
timeout=5,
)
return
# Validate that review is complete
if (self.approved_df is None or self.approved_df.empty) and (
self.needs_review_df is None or self.needs_review_df.empty
@@ -2449,6 +2761,15 @@ class PolicyPrepWorkflowScreen(Screen):
self._show_path_building_screen()
elif button_id == "build_preflight":
# Check if both required tabs have been reviewed
if not self.paths_tab_reviewed or not self.publishers_tab_reviewed:
self.app.notify(
"Please review both 'Paths' and 'Publishers' tabs before continuing.",
severity="warning",
timeout=5,
)
return
# Validate that path review is complete
if (self.primary_paths_df is None or self.primary_paths_df.empty) and (
self.publishers_df is None or self.publishers_df.empty
-444
View File
@@ -1,444 +0,0 @@
# Copyright (C) 2025 James Brotosky, Brandon Wickline
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU Affero General Public License as published
# by the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Affero General Public License for more details.
#
# You should have received a copy of the GNU Affero General Public License
# along with this program. If not, see <https://www.gnu.org/licenses/>.
import logging
import os
from typing import Optional
import dotenv
from textual.app import App, ComposeResult
from textual.containers import Vertical
from textual.message import Message
from textual.reactive import reactive
from textual.screen import Screen
from textual.widgets import (
Button,
DirectoryTree,
Footer,
Header,
Static,
Tab,
Tabs,
)
from models.agent import Agent
from models.policy import Policy
from services.API import AirlockAPIWrapper
from TUI.Screens.moveagentworkflowscreen import MoveAgentWorkflowScreen
from TUI.Screens.otpactivityscreen import OTPActivitiesScreen
from TUI.Screens.otprevokescreen import OTPRevokeScreen
from TUI.Screens.otpworkflowscreen import OTPWorkflowScreen
from TUI.Screens.policyprepworkflowscreen import PolicyPrepWorkflowScreen
from TUI.Screens.quietagentworkflowscreen import QuietAgentWorkflowScreen
from TUI.Themes.theme_amber_terminal import get_amber_terminal_theme
from TUI.Themes.theme_retro_terminal import get_retro_terminal_theme
from TUI.Themes.themeselector import ThemeSelector
from TUI.Widgets.agentmoveoperations import AgentMoveOperations
from TUI.Widgets.multiagentselector import MultiAgentSelector
from TUI.Widgets.policytreewidget import PolicyTreeWidget
from TUI.Widgets.resultsdisplay import ResultsDisplay
from utils.configmanager import get_user_value, load_env, save_user_config
from utils.setup import get_base_directory
from utils.utils import open_directory
dotenv.load_dotenv()
# ---------------------------------------------------------------------------
# GLOBAL STASH
# ---------------------------------------------------------------------------
_APP_RESTART_REASON = None
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# helper to persist TEXTUAL_THEME to *user* config and mirror to .env
# ---------------------------------------------------------------------------
def _persist_user_theme(theme_name: str) -> None:
"""
Store the chosen Textual theme in the user's config using the config manager.
No need to touch .env - config manager handles everything.
"""
base_dir = get_base_directory()
config_dir = base_dir / "config"
try:
save_user_config(config_dir, {"TEXTUAL_THEME": theme_name})
logger.debug("Updated user config with TEXTUAL_THEME=%s", theme_name)
except Exception as exc:
logger.error("Failed to save TEXTUAL_THEME: %s", exc)
# ---------------------------------------------------------------------------
# 1) SCREEN
# ---------------------------------------------------------------------------
class MainMenuScreen(Screen):
api: AirlockAPIWrapper
current_tab = reactive("")
BUTTON_DEFS = {
"agent_actions": [
(
"🖥️ - Find agent, Move agent, or Generate One Time Pass",
"move_agent_workflow_button",
),
("🎫 - Review and approve OTP Activities", "otp_activities_button"),
("🛑 - Revoke Active OTP Session", "otp_revoke_button"),
],
"policy": [
("⚖️ - Prepare Policy For Enforcement", "policy_prep_button"),
("🔕 - Find and Move Quiet Hosts to Enforcement", "find_quiet_button"),
],
}
def __init__(self) -> None:
super().__init__()
self.extras = get_user_value("EXTRAS", str, "NOTTODAY")
wd = load_env("WORKING_DIR") or os.getcwd()
if not os.path.isdir(wd):
wd = os.getcwd()
self.working_dir = wd
def _make_buttons_for(self, tab_id: str) -> Vertical:
defs = self.BUTTON_DEFS.get(tab_id, [])
buttons = []
for label, btn_id in defs:
btn = Button(label, id=btn_id)
btn.styles.width = "100%"
buttons.append(btn)
return Vertical(*buttons)
def compose(self) -> ComposeResult:
yield Header(show_clock=True, icon="")
tabs = [
Tab("Tree View", id="p_tree"),
Tab("Agents", id="agent_actions"),
Tab("Directory", id="dir"),
Tab("Settings", id="settings"),
]
if self.extras == "POLICYPREP":
tabs.insert(2, Tab("Policy Prep", id="policy"))
yield Tabs(*tabs, id="tabs")
yield Vertical(id="content")
yield Footer()
def on_mount(self) -> None:
self.switch_tab("agent_actions")
# focus helpers
def _get_content_buttons(self) -> list[Button]:
content = self.query_one("#content", Vertical)
return list(content.query(Button))
def _focus_first_button(self) -> None:
buttons = self._get_content_buttons()
if buttons:
buttons[0].focus()
def _focus_tabs(self) -> None:
tabs = self.query_one("#tabs", Tabs)
tabs.focus()
def _focus_nearby_button(self, direction: int) -> None:
buttons = self._get_content_buttons()
if not buttons:
return
try:
current = next(i for i, b in enumerate(buttons) if b.has_focus)
except StopIteration:
if direction > 0:
buttons[0].focus()
else:
buttons[-1].focus()
return
if direction < 0 and current == 0:
self._focus_tabs()
return
new_index = current + direction
if 0 <= new_index < len(buttons):
buttons[new_index].focus()
def switch_tab(self, tab_id: str) -> None:
self.current_tab = tab_id
content = self.query_one("#content", Vertical)
content.remove_children()
if tab_id in self.BUTTON_DEFS:
content.mount(self._make_buttons_for(tab_id))
self.call_later(self._focus_first_button)
elif tab_id == "dir":
content.mount(DirectoryTree(self.working_dir, id="dir_tree"))
elif tab_id == "p_tree":
content.mount(PolicyTreeWidget(self.app.policies, self.app.devices))
elif tab_id == "settings":
content.mount(ThemeSelector())
else:
content.mount(Static(f"Unknown tab: {tab_id}"))
def on_tabs_tab_activated(self, event: Tabs.TabActivated) -> None:
self.switch_tab(event.tab.id)
def on_multi_agent_selector_agents_selected(
self, message: MultiAgentSelector.AgentsSelected
) -> None:
"""Handle selected agents from AgentSelector."""
global _APP_RESTART_REASON
selected_agents = message.selected_agents
logger.info("Selected agents: %s", selected_agents)
# TODO: Implement actual handling of selected agents
_APP_RESTART_REASON = ("multi_agent_action", selected_agents)
self.app.exit()
def on_theme_selector_theme_selected(
self, message: ThemeSelector.ThemeSelected
) -> None:
"""Handle theme selection from ThemeSelector."""
global _APP_RESTART_REASON
_persist_user_theme(message.theme_name)
_APP_RESTART_REASON = ("restart",)
self.app.exit()
def on_agent_move_operations_operation_complete(
self, message: AgentMoveOperations.OperationComplete
) -> None:
"""Handle completion of agent move operation - show results."""
logger.info(
"Agent move operation completed: %s, %d successful, %d unsuccessful",
message.operation,
len(message.successful),
len(message.unsuccessful),
)
# Format results for display
successful_text = "\n".join(
[f"{agent.hostname}" for agent, _ in message.successful]
)
unsuccessful_text = "\n".join(
[f"{agent.hostname}: {error}" for agent, error in message.unsuccessful]
)
# Remove the operations widget
try:
ops_widget = self.query_one(AgentMoveOperations)
ops_widget.remove()
except Exception:
pass
# Show results
self.query_one("#content", Vertical).mount(
ResultsDisplay(message.operation, successful_text, unsuccessful_text)
)
def on_results_display_go_back(self, message: ResultsDisplay.GoBack) -> None:
"""Handle back button from results display."""
try:
results_widget = self.query_one(ResultsDisplay)
results_widget.remove()
except Exception:
pass
# Return to main menu
self.app.pop_screen()
def on_directory_tree_file_selected(
self, event: DirectoryTree.FileSelected
) -> None:
path = event.path
logger.debug("Directory file selected: %s", path)
try:
open_directory(str(path))
except Exception as exc:
logger.error("Failed to open %s: %s", path, exc)
self.app.bell()
def on_button_pressed(self, event: Button.Pressed) -> None:
button_id = event.button.id
logger.debug("Button pressed: %s", button_id)
match button_id:
case "move_agent_workflow_button":
self.app.push_screen(MoveAgentWorkflowScreen(self.app.devices))
event.stop()
case "otp_generate_button":
self.app.push_screen(OTPWorkflowScreen(self.app.devices))
event.stop()
case "find_quiet_button":
self.app.push_screen(
QuietAgentWorkflowScreen(self.app.api, self.app.policies)
)
event.stop()
return
case "otp_activities_button":
self.app.push_screen(OTPActivitiesScreen())
event.stop()
return
case "otp_revoke_button":
self.app.push_screen(OTPRevokeScreen())
event.stop()
return
case "policy_prep_button":
# Use the new TUI workflow screen instead of legacy
self.app.push_screen(
PolicyPrepWorkflowScreen(self.app.api, self.app.policies)
)
event.stop()
return
case _:
self.app.bell()
logger.warning("Unknown button pressed: %s", button_id)
return
# ---------------------------------------------------------------------------
# 2) APP
# ---------------------------------------------------------------------------
class Loxide(App[Message]):
api: AirlockAPIWrapper
working_dir: str
policies: Optional[list[Policy]]
devices: Optional[list[Agent]]
CSS = """
#logo {
width: 100%;
content-align: center middle;
text-align: center;
}
"""
BINDINGS = [
("q", "quit", "Quit"),
("f", "open_fe", "Launch Explorer"),
("r", "refresh", "Refresh"),
]
def __init__(self, api: AirlockAPIWrapper):
self._textual_theme = get_user_value("TEXTUAL_THEME", str, "textual-dark")
super().__init__()
self.api = api
wd = load_env("WORKING_DIR") or os.getcwd()
if not os.path.isdir(wd):
wd = os.getcwd()
self.working_dir = wd
# Initial data load
self.refresh_data()
def refresh_data(self) -> None:
"""Public method to refresh policies and devices from the API."""
try:
self.policies = [
Policy(**row.to_dict())
for _, row in self.api.policy_find_all().iterrows()
]
self.devices = [
Agent(**row.to_dict())
for _, row in self.api.agent_find_all().iterrows()
]
if self.policies and self.devices:
for agent in self.devices:
agent.enrich_with_policies(self.policies)
logger.debug(
f"Enriched {len(self.devices)} agents with policy information"
)
except Exception as exc:
logger.error("Failed to load policies/devices: %s", exc)
self.policies = None
self.devices = None
def on_mount(self, api: AirlockAPIWrapper) -> None:
self.register_theme(get_retro_terminal_theme())
self.register_theme(get_amber_terminal_theme())
self.theme = self._textual_theme
self.push_screen(MainMenuScreen())
def action_refresh(self) -> None:
self.refresh_data()
def action_quit(self) -> None:
global _APP_RESTART_REASON
_APP_RESTART_REASON = None
self.exit()
def action_open_fe(self) -> None:
"""Open the working directory in the OS file manager (footer binding)."""
path_to_open = self.working_dir or os.getcwd()
try:
open_directory(path_to_open)
except Exception as exc:
logger.error("Failed to open directory %s: %s", path_to_open, exc)
self.bell() # optional feedback
# ---------------------------------------------------------------------------
# 3) PUBLIC ENTRYPOINT
# ---------------------------------------------------------------------------
def run_Loxide(api: AirlockAPIWrapper) -> None:
global _APP_RESTART_REASON
base_dir = get_base_directory()
env_path = base_dir / ".env"
dotenv.load_dotenv(dotenv_path=env_path, override=True)
max_attempts = 5
attempts = 0
while attempts < max_attempts:
attempts += 1
logger.debug("Starting app loop iteration (attempt %d)", attempts)
_APP_RESTART_REASON = None
app = Loxide(api)
try:
app.run()
except SystemExit as exc:
if exc.code != 0:
logger.debug("Caught SystemExit from Textual: %s", exc)
raise
reason = _APP_RESTART_REASON
logger.debug("After app.run(), _APP_RESTART_REASON = %r", reason)
if not reason:
logger.debug("No restart reason, exiting loop")
break
if reason[0] == "restart":
logger.debug("Restarting app loop")
continue
if reason[0] == "multi_agent_action":
logger.info("Multi-agent action with selected agents: %s", reason[1])
continue
logger.error("Unknown restart reason: %r", reason)
break
# ---------------------------------------------------------------------------
# 4) DEV
# ---------------------------------------------------------------------------
if __name__ == "__main__":
api = AirlockAPIWrapper()
run_Loxide(api)
+46
View File
@@ -29,6 +29,7 @@ from textual.widget import Widget
from textual.widgets import Button, DataTable, Footer, Header, Static, TextArea
from models.agent import Agent
from TUI.Screens.executionhistoryscreen import ExecutionHistoryScreen
from TUI.Screens.otpworkflowscreen import OTPWorkflowScreen
from TUI.Screens.policyselectorscreen import PolicySelectorScreen
from TUI.Widgets.OTP_generate import OTPGenerator
@@ -140,6 +141,7 @@ class AgentMoveOperations(Widget):
toggle_enforcement_btn = self.query_one("#toggle_enforcement_btn", Button)
other_policy_btn = self.query_one("#other_policy_btn", Button)
otp_gen_btn = self.query_one("#otp_gen_btn", Button)
exec_history_btn = self.query_one("#exec_history_btn", Button)
# If operation in progress, disable all
if self.operation_in_progress:
@@ -148,6 +150,7 @@ class AgentMoveOperations(Widget):
local_approval_btn.disabled = True
toggle_enforcement_btn.disabled = True
other_policy_btn.disabled = True
exec_history_btn.disabled = True
else:
# If an operation was selected, disable
if self.selected_operation:
@@ -162,6 +165,9 @@ class AgentMoveOperations(Widget):
other_policy_btn.disabled = (
self.selected_operation == "other_policy"
)
exec_history_btn.disabled = (
self.selected_operation == "exec_history"
)
else:
# Enable all buttons
otp_gen_btn = False
@@ -169,6 +175,7 @@ class AgentMoveOperations(Widget):
local_approval_btn.disabled = False
toggle_enforcement_btn.disabled = False
other_policy_btn.disabled = False
exec_history_btn.disabled = False
except NoMatches:
pass
@@ -316,6 +323,13 @@ class AgentMoveOperations(Widget):
other_policy_btn.styles.margin = (0, 0, 1, 0)
yield other_policy_btn
exec_history_btn = Button(
"📊 View Execution History", id="exec_history_btn"
)
exec_history_btn.styles.width = "100%"
exec_history_btn.styles.margin = (0, 0, 1, 0)
yield exec_history_btn
# Status label
status_label = Static("", id="status_label")
status_label.styles.margin = (2, 0, 0, 0)
@@ -407,6 +421,9 @@ class AgentMoveOperations(Widget):
elif btn_id == "otp_gen_btn":
self._start_OTP_gen_operation()
event.stop()
elif btn_id == "exec_history_btn":
self._start_execution_history_operation()
event.stop()
def _start_local_approval_operation(self) -> None:
"""
@@ -681,6 +698,35 @@ class AgentMoveOperations(Widget):
self.app.push_screen(OTPWorkflowScreen(self.agents))
def _start_execution_history_operation(self) -> None:
"""
Launch the execution history viewer for selected agents.
This operation opens a new screen that allows the user to:
1. Select a date range for execution history
2. Fetch execution logs for all selected agents
3. View the results in a table
4. Export the results to CSV
The screen is pushed onto the screen stack, allowing the user to return
to this screen when done.
"""
status_label = self.query_one("#status_label", Static)
status_label.update("Opening execution history viewer...")
try:
# Push the execution history screen
self.app.push_screen(ExecutionHistoryScreen(self.agents))
logger.info(
f"Opened execution history viewer for {len(self.agents)} agents"
)
except Exception as e:
logger.error(f"Failed to open execution history viewer: {e}")
status_label.update(f"❌ Error: {str(e)}")
self.app.notify(
f"Failed to open execution history: {str(e)}", severity="error"
)
def _execute_move_to_policy(self, target_policy) -> None:
"""
Execute the actual move of agents to the selected policy.
+100 -1
View File
@@ -18,8 +18,9 @@ import logging
from rich.text import Text
from textual.containers import Horizontal, Vertical
from textual.message import Message
from textual.widget import Widget
from textual.widgets import Input, OptionList, Static, Switch, Tree
from textual.widgets import Button, Input, OptionList, Static, Switch, Tree
from textual.widgets.option_list import Option
logger = logging.getLogger(__name__)
@@ -28,6 +29,27 @@ logger = logging.getLogger(__name__)
class PolicyTreeWidget(Widget):
"""Widget for displaying and searching a hierarchical policy tree."""
class ViewExecutionHistory(Message):
"""Message sent when user wants to view execution history for a device."""
def __init__(self, device):
super().__init__()
self.device = device
class GenerateOTP(Message):
"""Message sent when user wants to generate OTP for a device."""
def __init__(self, device):
super().__init__()
self.device = device
class ToggleEnforcement(Message):
"""Message sent when user wants to toggle audit/enforcement for a device."""
def __init__(self, device):
super().__init__()
self.device = device
def __init__(self, policies, devices):
super().__init__()
self.policies = policies
@@ -35,6 +57,7 @@ class PolicyTreeWidget(Widget):
self.last_highlighted_node = None
self.leaf_counts = defaultdict(int)
self.match_type = "Count" # Default to sorting by count
self.selected_device = None # Track currently selected device
def compose(self):
# Create the switch and its label
@@ -56,6 +79,18 @@ class PolicyTreeWidget(Widget):
search_box = Input(
placeholder="Search policies or devices...", id="tree_search"
)
exec_history_button = Button(
"📊 Execution History", id="view_exec_history_button", disabled=True
)
exec_history_button.styles.margin = (0, 1, 0, 0) # Right margin
otp_button = Button("🎫 Generate OTP", id="generate_otp_button", disabled=True)
otp_button.styles.margin = (0, 1, 0, 0) # Right margin
toggle_enforcement_button = Button(
"🔄 Toggle Enforcement/Audit", id="toggle_enforcement_button", disabled=True
)
# No right margin on last button
details_pane = Static("", id="details_pane")
# Layout the UI
@@ -73,6 +108,12 @@ class PolicyTreeWidget(Widget):
# Add the search box and details pane
yield label
yield search_box
# Action buttons in a horizontal row
with Horizontal() as button_row:
button_row.styles.height = "auto"
yield exec_history_button
yield otp_button
yield toggle_enforcement_button
yield details_pane
def on_mount(self) -> None:
@@ -89,6 +130,32 @@ class PolicyTreeWidget(Widget):
# Expand the root node
policy_tree.root.expand()
def refresh_data(self, policies, devices):
"""Refresh the widget with new data and rebuild the tree."""
self.policies = policies
self.devices = devices
self.selected_device = None
# Disable all buttons since selection is lost
try:
self.query_one("#view_exec_history_button", Button).disabled = True
self.query_one("#generate_otp_button", Button).disabled = True
self.query_one("#toggle_enforcement_button", Button).disabled = True
except:
pass
# Rebuild tree with new data
self._precompute_leaf_counts()
total_leaves = sum(
self.leaf_counts.get(policy.groupid, 0)
for policy in self.policies
if policy.parent == "global-policy-settings"
)
policy_tree = self.query_one("#policy_tree", Tree)
policy_tree.root.set_label(f"Agents in Policies: ({total_leaves})")
self._build_tree()
policy_tree.root.expand()
def _precompute_leaf_counts(self):
"""Precompute leaf counts for each policy group."""
device_counts = defaultdict(int)
@@ -184,6 +251,9 @@ class PolicyTreeWidget(Widget):
node = message.node
data = node.data
details_pane = self.query_one("#details_pane", Static)
exec_history_button = self.query_one("#view_exec_history_button", Button)
otp_button = self.query_one("#generate_otp_button", Button)
toggle_enforcement_button = self.query_one("#toggle_enforcement_button", Button)
if self.last_highlighted_node is not None:
original_label = str(self.last_highlighted_node.label).strip()
@@ -198,6 +268,20 @@ class PolicyTreeWidget(Widget):
node.set_label(highlighted_label)
self.last_highlighted_node = node
# Check if selected node is a device (has Agent data)
from models.agent import Agent
if data and isinstance(data, Agent):
self.selected_device = data
exec_history_button.disabled = False
otp_button.disabled = False
toggle_enforcement_button.disabled = False
else:
self.selected_device = None
exec_history_button.disabled = True
otp_button.disabled = True
toggle_enforcement_button.disabled = True
if data:
details = "\n".join(
f"{key}: {value}" for key, value in data.__dict__.items()
@@ -287,3 +371,18 @@ class PolicyTreeWidget(Widget):
option_list.remove()
except:
pass
def on_button_pressed(self, event: Button.Pressed) -> None:
"""Handle button presses."""
if event.button.id == "view_exec_history_button":
if self.selected_device:
self.post_message(self.ViewExecutionHistory(self.selected_device))
event.stop()
elif event.button.id == "generate_otp_button":
if self.selected_device:
self.post_message(self.GenerateOTP(self.selected_device))
event.stop()
elif event.button.id == "toggle_enforcement_button":
if self.selected_device:
self.post_message(self.ToggleEnforcement(self.selected_device))
event.stop()
+159
View File
@@ -0,0 +1,159 @@
# Copyright (C) 2025 James Brotosky, Brandon Wickline
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU Affero General Public License as published
# by the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Affero General Public License for more details.
#
# You should have received a copy of the GNU Affero General Public License
# along with this program. If not, see <https://www.gnu.org/licenses/>.
import datetime
import logging
from bson import ObjectId
from textual.app import ComposeResult
from textual.containers import Container, Vertical
from textual.widgets import Button, DataTable, Static
from services.API import AirlockAPIWrapper
logger = logging.getLogger(__name__)
def skipback(days):
"""
Generate a MongoDB ObjectId for a given number of days ago from today.
"""
adjusted_days = days
date_days_ago = datetime.datetime.now(datetime.UTC) - datetime.timedelta(
days=adjusted_days
)
timestamp = int(date_days_ago.timestamp())
hex_timestamp = format(timestamp, "08x")
objectid_hex = hex_timestamp + "0000000000000000"
return ObjectId(objectid_hex)
class ServerLogWidget(Vertical):
"""Widget for displaying server activity logs in a DataTable."""
DEFAULT_CSS = """
ServerLogWidget {
width: 100%;
height: 100%;
}
ServerLogWidget #status_bar {
width: 100%;
height: auto;
background: $surface;
padding: 1;
margin-bottom: 1;
}
ServerLogWidget DataTable {
height: 1fr;
border: solid $primary;
}
ServerLogWidget #button_container {
width: 100%;
height: auto;
layout: horizontal;
padding: 1;
}
ServerLogWidget Button {
margin-right: 1;
}
"""
def __init__(self, api: AirlockAPIWrapper):
super().__init__()
self.api = api
def compose(self) -> ComposeResult:
yield Static("Loading server logs (last 72 hours)...", id="status_bar")
yield DataTable(id="server_log_table")
with Container(id="button_container"):
yield Button("Refresh", id="refresh_button", variant="primary")
def on_mount(self) -> None:
"""Initialize the DataTable and load server logs."""
self.load_logs()
def load_logs(self) -> None:
"""Load server logs from the API and populate the DataTable."""
table = self.query_one("#server_log_table", DataTable)
status = self.query_one("#status_bar", Static)
try:
status.update("⏳ Loading server logs (last 72 hours)...")
# Create a fake checkpoint for 3 days ago (72 hours)
checkpoint = str(skipback(3))
# Get server logs from API
logs = self.api.server_logs(checkpoint=checkpoint)
if not logs:
status.update("ℹï¸ No server logs found in the last 72 hours.")
table.clear(columns=True)
return
# Clear existing data
table.clear(columns=True)
# Add columns based on the first log entry
if logs:
first_log = logs[0]
columns = [col for col in first_log.keys() if col != "checkpoint"]
for col in columns:
table.add_column(col, key=col)
# Add rows in reverse order so newest entries are at the top
for log_entry in reversed(logs):
row_data = []
for col in columns:
value = log_entry.get(col, "")
# Format datetime column to be more readable
if col == "datetime" and value:
try:
# Parse ISO format and convert to readable format
dt = datetime.datetime.fromisoformat(
str(value).replace("Z", "+00:00")
)
value = dt.strftime("%Y-%m-%d %H:%M:%S")
except Exception:
# If parsing fails, just use the original value
pass
row_data.append(str(value))
table.add_row(*row_data)
status.update(
f"✅ Loaded {len(logs)} log entries from the last 72 hours"
)
logger.info(f"Loaded {len(logs)} server log entries")
else:
status.update("ℹï¸ No log entries found.")
except Exception as exc:
error_msg = f"❌ Error loading server logs: {exc}"
status.update(error_msg)
logger.error(f"Failed to load server logs: {exc}", exc_info=True)
table.clear(columns=True)
def on_button_pressed(self, event: Button.Pressed) -> None:
"""Handle button presses."""
button_id = event.button.id
if button_id == "refresh_button":
self.load_logs()
event.stop()
+308 -28
View File
@@ -24,6 +24,34 @@ dependencies = [
"memchr",
]
[[package]]
name = "airlock_libs"
version = "6.1.1"
dependencies = [
"chrono",
"crossbeam",
"flexi_logger",
"indicatif",
"log",
"mongodb",
"opentelemetry 0.27.1",
"opentelemetry-appender-log",
"opentelemetry-otlp",
"opentelemetry-proto",
"opentelemetry-semantic-conventions",
"opentelemetry_sdk 0.27.1",
"pyo3",
"reqwest",
"serde",
"serde-pyobject",
"serde_json",
"tokio",
"tonic",
"tracing",
"tracing-opentelemetry",
"tracing-subscriber",
]
[[package]]
name = "android_system_properties"
version = "0.1.5"
@@ -39,6 +67,150 @@ version = "1.0.100"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a23eb6b1614318a8071c9b2521f36b424b2c83db5eb3a0fead4a6c0809af6e61"
[[package]]
name = "async-channel"
version = "1.9.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "81953c529336010edd6d8e358f886d9581267795c61b19475b71314bffa46d35"
dependencies = [
"concurrent-queue",
"event-listener 2.5.3",
"futures-core",
]
[[package]]
name = "async-channel"
version = "2.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "924ed96dd52d1b75e9c1a3e6275715fd320f5f9439fb5a4a11fa51f4221158d2"
dependencies = [
"concurrent-queue",
"event-listener-strategy",
"futures-core",
"pin-project-lite",
]
[[package]]
name = "async-executor"
version = "1.13.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "497c00e0fd83a72a79a39fcbd8e3e2f055d6f6c7e025f3b3d91f4f8e76527fb8"
dependencies = [
"async-task",
"concurrent-queue",
"fastrand",
"futures-lite",
"pin-project-lite",
"slab",
]
[[package]]
name = "async-global-executor"
version = "2.4.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "05b1b633a2115cd122d73b955eadd9916c18c8f510ec9cd1686404c60ad1c29c"
dependencies = [
"async-channel 2.5.0",
"async-executor",
"async-io",
"async-lock",
"blocking",
"futures-lite",
"once_cell",
]
[[package]]
name = "async-io"
version = "2.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "456b8a8feb6f42d237746d4b3e9a178494627745c3c56c6ea55d92ba50d026fc"
dependencies = [
"autocfg",
"cfg-if",
"concurrent-queue",
"futures-io",
"futures-lite",
"parking",
"polling",
"rustix",
"slab",
"windows-sys 0.61.2",
]
[[package]]
name = "async-lock"
version = "3.4.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5fd03604047cee9b6ce9de9f70c6cd540a0520c813cbd49bae61f33ab80ed1dc"
dependencies = [
"event-listener 5.4.1",
"event-listener-strategy",
"pin-project-lite",
]
[[package]]
name = "async-process"
version = "2.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fc50921ec0055cdd8a16de48773bfeec5c972598674347252c0399676be7da75"
dependencies = [
"async-channel 2.5.0",
"async-io",
"async-lock",
"async-signal",
"async-task",
"blocking",
"cfg-if",
"event-listener 5.4.1",
"futures-lite",
"rustix",
]
[[package]]
name = "async-signal"
version = "0.2.13"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "43c070bbf59cd3570b6b2dd54cd772527c7c3620fce8be898406dd3ed6adc64c"
dependencies = [
"async-io",
"async-lock",
"atomic-waker",
"cfg-if",
"futures-core",
"futures-io",
"rustix",
"signal-hook-registry",
"slab",
"windows-sys 0.61.2",
]
[[package]]
name = "async-std"
version = "1.13.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2c8e079a4ab67ae52b7403632e4618815d6db36d2a010cfe41b02c1b1578f93b"
dependencies = [
"async-channel 1.9.0",
"async-global-executor",
"async-io",
"async-lock",
"async-process",
"crossbeam-utils",
"futures-channel",
"futures-core",
"futures-io",
"futures-lite",
"gloo-timers",
"kv-log-macro",
"log",
"memchr",
"once_cell",
"pin-project-lite",
"pin-utils",
"slab",
"wasm-bindgen-futures",
]
[[package]]
name = "async-stream"
version = "0.3.6"
@@ -61,6 +233,12 @@ dependencies = [
"syn",
]
[[package]]
name = "async-task"
version = "4.7.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8b75356056920673b02621b35afd0f7dda9306d03c79a30f5c56c44cf256e3de"
[[package]]
name = "async-trait"
version = "0.1.89"
@@ -164,6 +342,19 @@ dependencies = [
"generic-array",
]
[[package]]
name = "blocking"
version = "1.6.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e83f8d02be6967315521be875afa792a316e28d57b5a2d401897e2a7921b7f21"
dependencies = [
"async-channel 2.5.0",
"async-task",
"futures-io",
"futures-lite",
"piper",
]
[[package]]
name = "bson"
version = "2.15.0"
@@ -234,6 +425,15 @@ dependencies = [
"windows-link",
]
[[package]]
name = "concurrent-queue"
version = "2.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4ca0197aee26d1ae37445ee532fefce43251d24cc7c166799f4d46817f1d3973"
dependencies = [
"crossbeam-utils",
]
[[package]]
name = "console"
version = "0.16.1"
@@ -555,6 +755,33 @@ dependencies = [
"windows-sys 0.61.2",
]
[[package]]
name = "event-listener"
version = "2.5.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0206175f82b8d6bf6652ff7d71a1e27fd2e4efde587fd368662814d6ec1d9ce0"
[[package]]
name = "event-listener"
version = "5.4.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e13b66accf52311f30a0db42147dadea9850cb48cd070028831ae5f5d4b856ab"
dependencies = [
"concurrent-queue",
"parking",
"pin-project-lite",
]
[[package]]
name = "event-listener-strategy"
version = "0.5.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8be9f3dfaaffdae2972880079a491a1a8bb7cbed0b8dd7a347f668b4150a3b93"
dependencies = [
"event-listener 5.4.1",
"pin-project-lite",
]
[[package]]
name = "fastrand"
version = "2.3.0"
@@ -649,6 +876,19 @@ version = "0.3.31"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9e5c1b78ca4aae1ac06c48a526a655760685149f0d465d21f37abfe57ce075c6"
[[package]]
name = "futures-lite"
version = "2.6.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f78e10609fe0e0b3f4157ffab1876319b5b0db102a2c60dc4626306dc46b44ad"
dependencies = [
"fastrand",
"futures-core",
"futures-io",
"parking",
"pin-project-lite",
]
[[package]]
name = "futures-macro"
version = "0.3.31"
@@ -732,6 +972,18 @@ version = "0.3.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0cc23270f6e1808e30a928bdc84dea0b9b4136a8bc82338574f23baf47bbd280"
[[package]]
name = "gloo-timers"
version = "0.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bbb143cf96099802033e0d4f4963b19fd2e0b728bcf076cd9cf7f6634f092994"
dependencies = [
"futures-channel",
"futures-core",
"js-sys",
"wasm-bindgen",
]
[[package]]
name = "h2"
version = "0.4.12"
@@ -769,6 +1021,12 @@ version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea"
[[package]]
name = "hermit-abi"
version = "0.5.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fc0fef456e4baa96da950455cd02c081ca953b141298e41db3fc7e36b1da849c"
[[package]]
name = "hex"
version = "0.4.3"
@@ -1198,6 +1456,15 @@ dependencies = [
"wasm-bindgen",
]
[[package]]
name = "kv-log-macro"
version = "1.0.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0de8b303297635ad57c9f5059fd9cee7a47f8e8daa09df0fcd07dd39fb22977f"
dependencies = [
"log",
]
[[package]]
name = "lazy_static"
version = "1.5.0"
@@ -1236,6 +1503,9 @@ name = "log"
version = "0.4.29"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897"
dependencies = [
"value-bag",
]
[[package]]
name = "lru-slab"
@@ -1625,6 +1895,7 @@ version = "0.27.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "231e9d6ceef9b0b2546ddf52335785ce41252bc7474ee8ba05bfad277be13ab8"
dependencies = [
"async-std",
"async-trait",
"futures-channel",
"futures-executor",
@@ -1655,6 +1926,12 @@ dependencies = [
"thiserror 2.0.17",
]
[[package]]
name = "parking"
version = "2.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f38d5652c16fde515bb1ecef450ab0f6a219d619a7274976324d5e377f7dceba"
[[package]]
name = "parking_lot"
version = "0.12.5"
@@ -1725,12 +2002,37 @@ version = "0.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184"
[[package]]
name = "piper"
version = "0.2.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "96c8c490f422ef9a4efd2cb5b42b76c8613d7e7dfc1caf667b8a3350a5acc066"
dependencies = [
"atomic-waker",
"fastrand",
"futures-io",
]
[[package]]
name = "pkg-config"
version = "0.3.32"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c"
[[package]]
name = "polling"
version = "3.11.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5d0e4f59085d47d8241c88ead0f274e8a0cb551f3625263c05eb8dd897c34218"
dependencies = [
"cfg-if",
"concurrent-queue",
"hermit-abi",
"pin-project-lite",
"rustix",
"windows-sys 0.61.2",
]
[[package]]
name = "portable-atomic"
version = "1.11.1"
@@ -2413,34 +2715,6 @@ dependencies = [
"libc",
]
[[package]]
name = "signoz_test"
version = "6.0.0"
dependencies = [
"chrono",
"crossbeam",
"flexi_logger",
"indicatif",
"log",
"mongodb",
"opentelemetry 0.27.1",
"opentelemetry-appender-log",
"opentelemetry-otlp",
"opentelemetry-proto",
"opentelemetry-semantic-conventions",
"opentelemetry_sdk 0.27.1",
"pyo3",
"reqwest",
"serde",
"serde-pyobject",
"serde_json",
"tokio",
"tonic",
"tracing",
"tracing-opentelemetry",
"tracing-subscriber",
]
[[package]]
name = "slab"
version = "0.4.11"
@@ -3083,6 +3357,12 @@ version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65"
[[package]]
name = "value-bag"
version = "1.12.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7ba6f5989077681266825251a52748b8c1d8a4ad098cc37e440103d0ea717fc0"
[[package]]
name = "vcpkg"
version = "0.2.15"
+3 -3
View File
@@ -1,6 +1,6 @@
[package]
name = "signoz_test"
version = "6.0.0"
name = "airlock_libs"
version = "6.1.1"
edition = "2024"
[dependencies]
@@ -25,7 +25,7 @@ crossbeam = "0.8.4"
log = "0.4.29"
flexi_logger = "0.31.7"
opentelemetry-appender-log = "0.27.0"
opentelemetry_sdk = { version = "0.27.0", features = ["rt-tokio", "trace"] }
opentelemetry_sdk = { version = "0.27.0", features = ["rt-tokio", "testing", "trace"] }
[package.metadata.maturin]
generate-abi-stubs = true
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "maturin"
[project]
name = "airlock_libs"
version = "6.0.0"
version = "6.1.1"
description = "Airlock Digital API Wrapper"
readme = "README.md"
license = { text = "AGPL-3.0-only" }
+26 -1
View File
@@ -8,7 +8,32 @@ pub struct TelemetryConfig {
}
impl TelemetryConfig {
pub fn load() -> Self {
pub fn init_tracer() -> opentelemetry_sdk::trace::TracerProvider {
let cfg: TelemetryConfig = TelemetryConfig::load();
if !cfg.TELEMETRY {
return TracerProvider::builder().build();
}
let endpoint = cfg.TELEM_URL.unwrap_or_default();
let channel = Channel::from_shared(endpoint.clone())
.unwrap()
.tls_config(ClientTlsConfig::new().with_native_roots())
.unwrap()
.connect_lazy();
let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_tonic()
.with_endpoint(endpoint.clone())
.with_channel(channel)
.build()
.expect("Failed to build exporter");
opentelemetry_sdk::trace::TracerProvider::builder()
.with_simple_exporter(exporter)
.with_resource(Resource::new(vec![KeyValue::new(
"service.name",
"LoxideLibs",
)]))
.build()
}
fn load() -> Self {
let cfg_path = get_base_directory().join("config\\user_config.json");
if !cfg_path.exists() {
return Self {
View File
+5
View File
@@ -1,4 +1,5 @@
pub use chrono::{Duration, Local, NaiveDate};
pub use crossbeam::channel::unbounded;
pub use indicatif::{MultiProgress, ProgressBar, ProgressDrawTarget, ProgressStyle};
pub use mongodb::bson::oid::ObjectId;
pub use opentelemetry::global::GlobalTracerProvider;
@@ -7,6 +8,7 @@ pub use opentelemetry::trace::{Status, TraceContextExt, Tracer};
pub use opentelemetry::*;
pub use opentelemetry_otlp::ExportConfig;
pub use opentelemetry_otlp::WithExportConfig;
pub use opentelemetry_otlp::WithTonicConfig;
pub use opentelemetry_sdk::Resource;
pub use opentelemetry_sdk::trace::{Config, TracerProvider};
pub use pyo3::{prelude::*, types::PyString};
@@ -16,6 +18,8 @@ pub use reqwest::{
};
pub use serde::{Deserialize, Serialize};
pub use serde_json::Value;
pub use std::sync::{Arc, Mutex};
pub use std::thread;
pub use std::{
collections::HashMap,
env,
@@ -25,3 +29,4 @@ pub use std::{
path::PathBuf,
str::FromStr,
};
pub use tonic::transport::{Channel, ClientTlsConfig};
+3 -34
View File
@@ -1,10 +1,6 @@
use crate::modules::datatypes::*;
use crate::prelude::*;
use crossbeam::channel::unbounded;
use opentelemetry_otlp::WithTonicConfig;
use std::sync::{Arc, Mutex};
use std::thread;
use tonic::transport::{Channel, ClientTlsConfig};
#[pyfunction]
pub fn pull_policy_exec_histories(
py: Python<'_>,
@@ -25,7 +21,7 @@ pub fn pull_policy_exec_histories(
std::process::abort();
}
};
let tracer_provider = rt.block_on(async { init_tracer() });
let tracer_provider = rt.block_on(async { TelemetryConfig::init_tracer() });
global::set_tracer_provider(tracer_provider.clone());
let tracer: global::BoxedTracer = global::tracer("tracer");
let _cx: Context = Context::new();
@@ -184,10 +180,7 @@ pub fn pull_policy_exec_histories(
.unwrap()
.enable_steady_tick(std::time::Duration::from_millis(100));
let span: opentelemetry::trace::SpanRef<'_> = cx.span();
//span.set_attribute(Key::new("Days").string(days.to_string()));
//span.set_attribute(Key::new("Days"));
//span.set_attribute(KeyValue::new("Policy Name", policy_names.clone()));
span.set_attribute(KeyValue::new("Days", days));
span.set_attribute(KeyValue::new("Days", days.to_string()));
span.set_attribute(KeyValue::new("Policy Name", policy_names.clone()));
loop {
let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| {
@@ -320,27 +313,3 @@ pub fn get_base_directory() -> PathBuf {
}
}
}
fn init_tracer() -> opentelemetry_sdk::trace::TracerProvider {
let cfg: TelemetryConfig = TelemetryConfig::load();
let endpoint = cfg.TELEM_URL.unwrap_or_default().clone();
let channel_endpoint = endpoint.clone();
let channel = Channel::from_shared(channel_endpoint.clone())
.unwrap()
.tls_config(ClientTlsConfig::new().with_native_roots())
.unwrap()
.connect_lazy();
let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_tonic()
.with_endpoint(endpoint.clone())
.with_channel(channel)
.build()
.expect("Failed to build exporter");
opentelemetry_sdk::trace::TracerProvider::builder()
.with_simple_exporter(exporter)
.with_resource(Resource::new(vec![KeyValue::new(
"service.name",
"LoxideLibs",
)]))
.build()
}
+18 -6
View File
@@ -1,14 +1,26 @@
# Core TUI dependencies
textual==6.5.0
# API and data handling
Requests==2.32.5
pandas==2.3.3
numpy==2.3.4
# Database
pymongo==4.15.3
# Security and encryption
cryptography==46.0.3
keyring==25.6.0
numpy==2.3.4
pandas==2.3.3
pymongo==4.15.3
# Environment management
python-dotenv==1.2.1
Requests==2.32.5
textual==6.5.0
# Utilities
tqdm==4.67.1
urllib3==2.5.0
pyperclip==1.11.0
# Custom/Private packages
--extra-index-url https://git.racooncity.org/api/packages/brotoskyj/pypi/simple/
airlock_libs==6.0.0
airlock_libs==6.1.1
+8
View File
@@ -351,6 +351,14 @@ class AirlockAPIWrapper:
result = self._post("/v1/getexechistory", payload)
return result["response"]["exechistory"]
def server_logs(self, checkpoint: str | None = None) -> str:
"""Retrieves Server Activity History Logs."""
payload = {}
if checkpoint is not None:
payload["checkpoint"] = checkpoint
result = self._post("/v1/logging/svractivities?checkpoint", payload)
return result["response"]["svractivities"]
"""
from services.API import AirlockAPIWrapper
+75 -22
View File
@@ -29,6 +29,49 @@ from utils.configmanager import (
)
class TextualNotificationHandler(logging.Handler):
"""
Custom logging handler that sends ERROR, WARNING, and CRITICAL logs
to Textual toast notifications.
"""
def __init__(self, app):
super().__init__()
self.app = app
def emit(self, record):
try:
# Only handle ERROR, WARNING, and CRITICAL
if record.levelno >= logging.WARNING:
# Format the message
msg = self.format(record)
# Map log levels to Textual severity
severity_map = {
logging.WARNING: "warning",
logging.ERROR: "error",
logging.CRITICAL: "error",
}
severity = severity_map.get(record.levelno, "information")
# Send to Textual notification
# Use call_from_thread if logging from non-main thread
try:
self.app.notify(msg, severity=severity, timeout=5)
except Exception:
# If we're not on the main thread, schedule it
try:
self.app.call_from_thread(
self.app.notify, msg, severity=severity, timeout=5
)
except Exception:
# Silently fail to avoid breaking the logging system
pass
except Exception:
# Silently fail to avoid breaking the logging system
pass
def get_base_directory() -> Path:
system = platform.system()
home = Path.home()
@@ -44,38 +87,31 @@ def configure_logging(log_dir: Path, log_level: str = "INFO"):
log_file = log_dir / "Loxide.log"
config = {
"version": 1, # Required key for dictConfig format version
"disable_existing_loggers": False, # Keeps existing loggers active
"version": 1,
"disable_existing_loggers": False,
"formatters": {
"detailed": {
"format": "%(asctime)s - %(name)s - %(levelname)s - %(message)s"
# Includes timestamp, logger name, level, and message
},
"simple": {
"format": "%(levelname)s - %(message)s"
# Minimal format for console output
},
"simple": {"format": "%(levelname)s - %(message)s"},
"toast": {"format": "%(name)s: %(message)s"}, # Simpler format for toasts
},
"handlers": {
"file": {
"class": "logging.handlers.TimedRotatingFileHandler",
"filename": str(log_file),
"when": "midnight", # Rotate logs at midnight
"interval": 1, # Every 1 day
"backupCount": 7, # Keep 7 days of logs
"encoding": "utf-8", # Ensure UTF-8 encoding
"level": "DEBUG", # Always log DEBUG and above to file
"formatter": "detailed", # Use detailed format
},
"console": {
"class": "logging.StreamHandler",
"level": log_level.upper(), # System-configured level for console
"formatter": "simple", # Use simple format
"when": "midnight",
"interval": 1,
"backupCount": 7,
"encoding": "utf-8",
"level": "DEBUG",
"formatter": "detailed",
},
# REMOVED console handler - it interferes with Textual TUI
},
"root": {
"level": "DEBUG", # Root logger level
"handlers": ["file", "console"], # Attach both handlers
"level": "DEBUG",
"handlers": ["file"], # Only use file handler, not console
},
}
@@ -96,6 +132,18 @@ def configure_logging(log_dir: Path, log_level: str = "INFO"):
logging.config.dictConfig(config)
logging.getLogger().debug("✅ Logging configured.")
# Return a function to attach the notification handler once the app is created
def attach_notification_handler(app):
"""Attach the Textual notification handler to the root logger."""
handler = TextualNotificationHandler(app)
handler.setLevel(logging.WARNING) # Only WARNING and above
formatter = logging.Formatter("%(name)s: %(message)s")
handler.setFormatter(formatter)
logging.getLogger().addHandler(handler)
logging.getLogger().debug("✅ Textual notification handler attached.")
return attach_notification_handler
def setup():
"""
@@ -105,6 +153,9 @@ def setup():
3. Load user config (mutable, from user_config.json)
4. Configure logging
5. Set up .env with WORKING_DIR only
Returns:
attach_notification_handler: Function to attach notification handler to TUI app
"""
base_dir = get_base_directory()
dirs = {
@@ -122,7 +173,7 @@ def setup():
# Configure logging with system-defined log level
log_level = get_system_value("LOG_LEVEL", str, "INFO")
configure_logging(dirs["logs"], log_level)
attach_handler = configure_logging(dirs["logs"], log_level)
# Load user config (mutable)
load_user_config(dirs["config"])
@@ -145,7 +196,6 @@ def setup():
"Approved": [],
"Needs_Review": ["Review_First", "Review_Second", "HTML"],
"Preflight": ["HTML"],
"Archived": [],
}
for folder_name, subfolders in folders_structure.items():
@@ -158,3 +208,6 @@ def setup():
logging.debug(f"'{subfolder}' subfolder created at: {subfolder_path}")
logging.info("✅ Setup complete")
# Return the attach handler function
return attach_handler