Compare commits

...

8 Commits

Author SHA1 Message Date
brotoskyj 0ac3b54d89 Refactored Progress Bar
Build Library / Build Library (push) Successful in 5m54s
Refactored Progress Bar to remove multiprogress bar and only draw one instance. #36 is still open and not fixed with this push, but I believe this is the way to fix the issue.
Also implemented an Arc Mutex on the progress bar so it can be controlled via different threads.
2025-12-05 17:19:40 -05:00
Zarithas 154a7efcc8 Bug fix for Revoke OTP resolved, no longer crashes when select all is chosen when there are no active sessions, Swapped Revoke OTP and Quiet Hosts locations on menus 2025-12-05 15:30:24 -05:00
Zarithas ab5f00d8e7 Merge branch 'RustImplementation' of https://git.racooncity.org/brotoskyj/Airlocktools into RustImplementation 2025-12-05 15:06:06 -05:00
Zarithas 3ab803c12e Quiet Agent UI improvements 2025-12-05 15:05:46 -05:00
Zarithas 98cb23e5ea Bugfix for Issue 39.
Fixes:
brotoskyj/AirlockTools#39
2025-12-05 14:48:23 -05:00
brotoskyj 0aabbfd36e Style Change
Build Library / Build Library (push) Successful in 6m25s
Cleaned up services.rs file and removed whitespaces
Specified Data Types were needed instead of allowing the compiler to select types
2025-12-05 11:12:29 -05:00
brotoskyj b19eeb6c96 Fixed Progress Bar
Build Library / Build Library (push) Successful in 6m22s
Changed date math so that the last date from the first response from the API is the minimum value for the progress bar. This gives true progress percentages from the first date to the last date on the last request. #36
2025-12-04 11:57:49 -05:00
brotoskyj 76bd3a6087 Feature added - MPSC channel
Build Library / Build Library (push) Successful in 5m44s
closes #38
2025-12-03 17:54:57 -05:00
12 changed files with 403 additions and 270 deletions
+30 -31
View File
@@ -58,9 +58,14 @@ class OTPRevokeWidget(Static):
} }
#button_container { #button_container {
height: auto; height: auto;
width: 100%;
padding: 1; padding: 1;
align: center middle; align: center middle;
} }
#button_container Button {
min-width: 16;
margin: 0 1;
}
#result_container { #result_container {
height: auto; height: auto;
max-height: 10; max-height: 10;
@@ -87,29 +92,10 @@ class OTPRevokeWidget(Static):
# Action buttons # Action buttons
with Horizontal(id="button_container"): with Horizontal(id="button_container"):
self.refresh_button = Button("🔄 Refresh", id="refresh_btn") yield Button("Refresh", id="refresh_btn")
self.refresh_button.styles.width = "15%" yield Button("Select All", id="select_all_btn")
self.refresh_button.styles.margin = (1, 1, 1, 1) yield Button("Clear Selection", id="select_none_btn")
yield self.refresh_button yield Button("Revoke Selected", id="revoke_btn", variant="error")
self.select_all_button = Button("☑️ Select All", id="select_all_btn")
self.select_all_button.styles.width = "15%"
self.select_all_button.styles.margin = (1, 1, 1, 1)
yield self.select_all_button
self.select_none_button = Button(
"❌ Clear Selection", id="select_none_btn"
)
self.select_none_button.styles.width = "20%"
self.select_none_button.styles.margin = (1, 1, 1, 1)
yield self.select_none_button
self.revoke_button = Button(
"🛑 Revoke Selected", id="revoke_btn", variant="error"
)
self.revoke_button.styles.width = "20%"
self.revoke_button.styles.margin = (1, 1, 1, 1)
yield self.revoke_button
# Results display # Results display
with Vertical(id="result_container"): with Vertical(id="result_container"):
@@ -122,7 +108,7 @@ class OTPRevokeWidget(Static):
# Configure sessions table # Configure sessions table
self.sessions_table.clear() self.sessions_table.clear()
self.sessions_table.add_columns( self.sessions_table.add_columns(
"", "OTP ID", "Hostname", "Status", "Purpose", "Granted" "", "OTP ID", "Hostname", "Status", "Purpose", "Granted"
) )
# Enable row selection with checkbox column # Enable row selection with checkbox column
@@ -213,7 +199,11 @@ class OTPRevokeWidget(Static):
elif btn.id == "select_all_btn": elif btn.id == "select_all_btn":
# Select all visible rows # Select all visible rows
if self._filtered_df is not None: if (
self._filtered_df is not None
and not self._filtered_df.empty
and "otpid" in self._filtered_df.columns
):
self._selected_otpids = set(str(x) for x in self._filtered_df["otpid"]) self._selected_otpids = set(str(x) for x in self._filtered_df["otpid"])
await self._refresh_table() await self._refresh_table()
@@ -235,7 +225,12 @@ class OTPRevokeWidget(Static):
# Get the row index from the cursor row # Get the row index from the cursor row
row_index = self.sessions_table.cursor_row row_index = self.sessions_table.cursor_row
if self._filtered_df is not None and row_index < len(self._filtered_df): if (
self._filtered_df is not None
and not self._filtered_df.empty
and "otpid" in self._filtered_df.columns
and row_index < len(self._filtered_df)
):
# Get the OTP ID for this row # Get the OTP ID for this row
otpid = str(self._filtered_df.iloc[row_index]["otpid"]) otpid = str(self._filtered_df.iloc[row_index]["otpid"])
@@ -257,12 +252,12 @@ class OTPRevokeWidget(Static):
async def _revoke_selected(self) -> None: async def _revoke_selected(self) -> None:
"""Revoke the selected OTP sessions.""" """Revoke the selected OTP sessions."""
if not self._selected_otpids: if not self._selected_otpids:
self.results_display.update("No sessions selected for revocation") self.results_display.update("No sessions selected for revocation")
return return
api = getattr(self.app, "api", None) api = getattr(self.app, "api", None)
if not api: if not api:
self.results_display.update("API not available") self.results_display.update("API not available")
return return
# Collect results # Collect results
@@ -303,13 +298,13 @@ class OTPRevokeWidget(Static):
else "No response" else "No response"
) )
results.append( results.append(
f"Failed to revoke OTP {otpid} for {hostname}: {error_msg}" f"Failed to revoke OTP {otpid} for {hostname}: {error_msg}"
) )
logger.error(f"Failed to revoke OTP {otpid}: {error_msg}") logger.error(f"Failed to revoke OTP {otpid}: {error_msg}")
except Exception as e: except Exception as e:
failure_count += 1 failure_count += 1
results.append(f"Error revoking OTP {otpid}: {str(e)}") results.append(f"Error revoking OTP {otpid}: {str(e)}")
logger.exception(f"Exception revoking OTP {otpid}: {e}") logger.exception(f"Exception revoking OTP {otpid}: {e}")
# Update results display # Update results display
@@ -368,7 +363,11 @@ class OTPRevokeScreen(Screen):
async def action_select_all(self) -> None: async def action_select_all(self) -> None:
"""Select all visible sessions.""" """Select all visible sessions."""
if self.widget._filtered_df is not None: if (
self.widget._filtered_df is not None
and not self.widget._filtered_df.empty
and "otpid" in self.widget._filtered_df.columns
):
self.widget._selected_otpids = set( self.widget._selected_otpids = set(
str(x) for x in self.widget._filtered_df["otpid"] str(x) for x in self.widget._filtered_df["otpid"]
) )
+183 -110
View File
@@ -34,7 +34,7 @@ from textual.app import ComposeResult
from textual.containers import Horizontal, Vertical from textual.containers import Horizontal, Vertical
from textual.reactive import reactive from textual.reactive import reactive
from textual.screen import Screen from textual.screen import Screen
from textual.widgets import Button, DataTable, Footer, Header, Static from textual.widgets import Button, DataTable, Footer, Header, Input, Static
from models.policy import Policy from models.policy import Policy
from services.API import AirlockAPIWrapper from services.API import AirlockAPIWrapper
@@ -51,16 +51,17 @@ class QuietAgentWorkflowScreen(Screen):
This screen provides a multi-step workflow: This screen provides a multi-step workflow:
1. Select initial policy to analyze 1. Select initial policy to analyze
2. View categorized agents (enforce ready vs. non-enforce ready) 2. Configure analysis parameters (history period and quiet time period)
3. Select target policies for each category 3. View categorized agents (enforce ready vs. non-enforce ready)
4. Execute agent migrations 4. Select target policies for each category
5. Execute agent migrations
Attributes: Attributes:
api (AirlockAPIWrapper): API wrapper for Airlock operations api (AirlockAPIWrapper): API wrapper for Airlock operations
policies (List[Policy]): List of all available policies policies (List[Policy]): List of all available policies
selected_policy (Optional[Policy]): The initially selected policy to analyze selected_policy (Optional[Policy]): The initially selected policy to analyze
history_days (int): Number of days of history to pull (default: 150) history_days (int): Number of days of history to pull (default: 150, range: 1-365)
quiet_days (int): Number of days without execution to be considered quiet (default: 45) quiet_days (int): Number of days without execution to be considered quiet (default: 45, range: 1-365)
agents_df (Optional[pd.DataFrame]): DataFrame of all agents with analysis results agents_df (Optional[pd.DataFrame]): DataFrame of all agents with analysis results
enforce_ready_df (Optional[pd.DataFrame]): DataFrame of agents ready for enforcement enforce_ready_df (Optional[pd.DataFrame]): DataFrame of agents ready for enforcement
non_enforce_ready_df (Optional[pd.DataFrame]): DataFrame of agents not ready for enforcement non_enforce_ready_df (Optional[pd.DataFrame]): DataFrame of agents not ready for enforcement
@@ -86,7 +87,7 @@ class QuietAgentWorkflowScreen(Screen):
self.api = api self.api = api
self.policies = policies self.policies = policies
self.selected_policy: Optional[Policy] = None self.selected_policy: Optional[Policy] = None
self.history_days = 150 # Fixed as per requirements self.history_days = 150 # Default value, user-selectable
self.quiet_days = 45 # Default value self.quiet_days = 45 # Default value
self.agents_df: Optional[pd.DataFrame] = None self.agents_df: Optional[pd.DataFrame] = None
self.enforce_ready_df: Optional[pd.DataFrame] = None self.enforce_ready_df: Optional[pd.DataFrame] = None
@@ -130,7 +131,7 @@ class QuietAgentWorkflowScreen(Screen):
stage_messages = { stage_messages = {
"select_policy": "Step 1: Select Policy to Analyze", "select_policy": "Step 1: Select Policy to Analyze",
"select_quiet_days": "Step 2: Select Quiet Time Period", "select_history_days": "Step 2: Configure Analysis Parameters",
"analyzing": "Analyzing agent activity...", "analyzing": "Analyzing agent activity...",
"view_results": "Step 3: Review Categorized Agents", "view_results": "Step 3: Review Categorized Agents",
"select_enforce_target": "Step 4: Select Target Policy for Enforce Ready Agents", "select_enforce_target": "Step 4: Select Target Policy for Enforce Ready Agents",
@@ -161,7 +162,7 @@ class QuietAgentWorkflowScreen(Screen):
# Initial policy selection for analysis # Initial policy selection for analysis
self.selected_policy = message.policy self.selected_policy = message.policy
logger.info(f"Selected policy for analysis: {self.selected_policy.name}") logger.info(f"Selected policy for analysis: {self.selected_policy.name}")
self._show_quiet_days_selection() self._show_history_days_selection()
elif self.workflow_stage == "select_enforce_target": elif self.workflow_stage == "select_enforce_target":
# Target policy selection for enforce ready agents # Target policy selection for enforce ready agents
self.enforce_ready_target_policy = message.policy self.enforce_ready_target_policy = message.policy
@@ -177,48 +178,167 @@ class QuietAgentWorkflowScreen(Screen):
) )
self._show_migration_confirmation() self._show_migration_confirmation()
def _show_quiet_days_selection(self) -> None: def _show_history_days_selection(self) -> None:
"""Show the quiet days selection screen.""" """Show the history days and quiet days selection screen."""
self.workflow_stage = "select_quiet_days" self.workflow_stage = "select_history_days"
content = self.query_one("#content_area", Vertical) content = self.query_one("#content_area", Vertical)
content.remove_children() content.remove_children()
# Create info text # Create info text
info_widget = Static( info_widget = Static(
f"Policy Selected: {self.selected_policy.name}\n\n" f"Policy Selected: {self.selected_policy.name}\n\n"
f"History Period: {self.history_days} days\n\n" "Configure Analysis Parameters:",
"Select quiet time period (days without untrusted execution):", id="analysis_params_info",
id="quiet_days_info",
) )
info_widget.styles.margin = (0, 0, 2, 0) info_widget.styles.margin = (0, 0, 2, 0)
content.mount(info_widget) content.mount(info_widget)
# Create button container and mount it first # Create input container
button_container = Vertical(id="quiet_days_buttons") input_container = Vertical(id="analysis_params_input_container")
button_container.styles.height = "auto" input_container.styles.height = "auto"
content.mount(button_container) content.mount(input_container)
# Now add buttons to the mounted container # History days label
for days in [15, 30, 45, 60]: history_label = Static("History Period (days of execution history to pull):")
btn = Button( history_label.styles.margin = (0, 0, 1, 0)
f"{days} days {'(Default)' if days == 45 else ''}", input_container.mount(history_label)
id=f"quiet_days_{days}",
classes="quiet_day_btn", # Add history days input field
history_input = Input(
placeholder="Enter days (1-365, default: 150)",
value="150",
id="history_days_input",
)
history_input.styles.width = "50"
history_input.styles.margin = (0, 0, 2, 0)
input_container.mount(history_input)
# Quiet days label
quiet_label = Static(
"Quiet Time Period (days without execution to be considered quiet):"
)
quiet_label.styles.margin = (0, 0, 1, 0)
input_container.mount(quiet_label)
# Add quiet days input field
quiet_input = Input(
placeholder="Enter days (1-365, default: 45)",
value="45",
id="quiet_days_input",
)
quiet_input.styles.width = "50"
quiet_input.styles.margin = (0, 0, 2, 0)
input_container.mount(quiet_input)
# Add submit button
submit_btn = Button(
"Continue",
id="analysis_params_submit",
variant="primary",
)
submit_btn.styles.width = "50"
submit_btn.styles.margin = (1, 0, 0, 0)
input_container.mount(submit_btn)
# Focus the first input field
history_input.focus()
def _validate_and_submit_history_days(self) -> None:
"""Validate and submit the history days and quiet days inputs."""
try:
history_input = self.query_one("#history_days_input", Input)
quiet_input = self.query_one("#quiet_days_input", Input)
history_value = history_input.value.strip()
quiet_value = quiet_input.value.strip()
# Validate history days
if not history_value:
self.app.notify(
"Please enter a history period value", severity="error", timeout=3
)
history_input.focus()
return
try:
history_days = int(history_value)
except ValueError:
self.app.notify(
"Please enter a valid number for history period",
severity="error",
timeout=3,
)
history_input.focus()
return
if history_days < 1 or history_days > 365:
self.app.notify(
"History period must be between 1 and 365 days",
severity="error",
timeout=3,
)
history_input.focus()
return
# Validate quiet days
if not quiet_value:
self.app.notify(
"Please enter a quiet time period value",
severity="error",
timeout=3,
)
quiet_input.focus()
return
try:
quiet_days = int(quiet_value)
except ValueError:
self.app.notify(
"Please enter a valid number for quiet time period",
severity="error",
timeout=3,
)
quiet_input.focus()
return
if quiet_days < 1 or quiet_days > 365:
self.app.notify(
"Quiet time period must be between 1 and 365 days",
severity="error",
timeout=3,
)
quiet_input.focus()
return
# Check that quiet days doesn't exceed history days
if quiet_days > history_days:
self.app.notify(
"Quiet time period cannot exceed history period",
severity="error",
timeout=3,
)
quiet_input.focus()
return
# All validation passed
self.history_days = history_days
self.quiet_days = quiet_days
logger.info(
f"Selected history days: {history_days}, quiet days: {quiet_days}"
) )
btn.styles.width = "100%" self._start_analysis()
btn.styles.margin = (0, 0, 1, 0)
button_container.mount(btn) except Exception as e:
logger.error(f"Error validating analysis parameters: {e}")
self.app.notify(f"Error: {str(e)}", severity="error", timeout=3)
def on_button_pressed(self, event: Button.Pressed) -> None: def on_button_pressed(self, event: Button.Pressed) -> None:
"""Handle button press events.""" """Handle button press events."""
button_id = event.button.id button_id = event.button.id
# Quiet days selection buttons # Analysis parameters submit button
if button_id and button_id.startswith("quiet_days_"): if button_id == "analysis_params_submit":
days = int(button_id.split("_")[-1]) self._validate_and_submit_history_days()
self.quiet_days = days
logger.info(f"Selected quiet days: {days}")
self._start_analysis()
return return
# Navigation buttons # Navigation buttons
@@ -258,46 +378,44 @@ class QuietAgentWorkflowScreen(Screen):
self._show_policy_selection() self._show_policy_selection()
return return
def on_input_submitted(self, event: Input.Submitted) -> None:
"""Handle input submission (Enter key pressed)."""
if event.input.id in ["history_days_input", "quiet_days_input"]:
self._validate_and_submit_history_days()
def _start_analysis(self) -> None: def _start_analysis(self) -> None:
"""Start the agent activity analysis.""" """Start the agent activity analysis."""
self.workflow_stage = "analyzing" # Show notification that analysis is starting
content = self.query_one("#content_area", Vertical)
content.remove_children()
# Show analyzing message with detailed steps
analyzing_msg = Static(
f"Analyzing Agent Activity\n"
f"{'=' * 50}\n\n"
f"Policy: {self.selected_policy.name}\n"
f"History Period: {self.history_days} days\n"
f"Quiet Threshold: {self.quiet_days} days\n\n"
f"Progress:\n"
f"Step 1/4: Fetching agents from policy...\n"
f"Step 2/4: Pulling execution history (this may take a moment)...\n"
f"Step 3/4: Analyzing activity patterns...\n"
f"Step 4/4: Categorizing agents...\n\n"
f"Please wait - this operation cannot be cancelled.",
id="analyzing_message",
)
analyzing_msg.styles.margin = (2, 1)
content.mount(analyzing_msg)
# Show notification
self.app.notify( self.app.notify(
"Starting analysis - this may take several minutes for large policies", "Starting analysis - this may take several minutes for large policies",
severity="information", severity="information",
timeout=5, timeout=5,
) )
# Perform the analysis asynchronously # Clear the screen to provide a blank canvas for Rust progress output
self.call_later(self._perform_analysis) # (Rust output displays over the TUI, so we clear everything except header/footer)
try:
# Clear title
title_widget = self.query_one("#workflow_title", Static)
title_widget.update("")
def _perform_analysis(self) -> None: # Clear status
status_widget = self.query_one("#workflow_status", Static)
status_widget.update("")
# Clear content area
content = self.query_one("#content_area", Vertical)
content.remove_children()
except Exception as e:
logger.debug(f"Could not clear screen for analysis: {e}")
# Delay the analysis start to ensure UI refresh completes first
# This prevents Rust output from starting before the screen is cleared
self.set_timer(0.5, self._perform_analysis_worker)
def _perform_analysis_worker(self) -> None:
"""Perform the actual agent activity analysis.""" """Perform the actual agent activity analysis."""
try: try:
# Update status: Fetching agents
self._update_analysis_status("Step 1/4: Fetching agents from policy...")
# Get agents in the selected policy # Get agents in the selected policy
agents = self.api.agents_find_by_group(self.selected_policy.groupid) agents = self.api.agents_find_by_group(self.selected_policy.groupid)
@@ -310,32 +428,11 @@ class QuietAgentWorkflowScreen(Screen):
self._show_policy_selection() self._show_policy_selection()
return return
agent_count = len(agents)
self.app.notify(
f"Found {agent_count} agents - fetching execution history...",
severity="information",
timeout=3,
)
# Update status: Pulling execution history
self._update_analysis_status(
f"Step 2/4: Pulling execution history for {agent_count} agents...\n"
f"(This may take several minutes - progress shown in terminal)"
)
# Get execution history (this shows progress bars in terminal via airlock_libs) # Get execution history (this shows progress bars in terminal via airlock_libs)
policy_exec_history = getPolicyInfo( policy_exec_history = getPolicyInfo(
self.api, self.selected_policy, [1, 2, 6, 7], self.history_days self.api, self.selected_policy, [1, 2, 6, 7], self.history_days
) )
# Update status: Analyzing patterns
self._update_analysis_status("Step 3/4: Analyzing activity patterns...")
self.app.notify(
"History retrieved - analyzing patterns...",
severity="information",
timeout=2,
)
if policy_exec_history.empty: if policy_exec_history.empty:
logger.info( logger.info(
"No execution history found for the selected policy and time range." "No execution history found for the selected policy and time range."
@@ -381,9 +478,6 @@ class QuietAgentWorkflowScreen(Screen):
lambda x: True if pd.isna(x) or x > self.quiet_days else False lambda x: True if pd.isna(x) or x > self.quiet_days else False
) )
# Update status: Categorizing
self._update_analysis_status("Step 4/4: Categorizing agents...")
# Sort agents # Sort agents
agents = agents.sort_values( agents = agents.sort_values(
by=["execution_count", "hostname"], ascending=[True, True] by=["execution_count", "hostname"], ascending=[True, True]
@@ -394,7 +488,7 @@ class QuietAgentWorkflowScreen(Screen):
# Categorize agents into DataFrames # Categorize agents into DataFrames
self.enforce_ready_df = agents[agents["enforce_ready"]].copy() self.enforce_ready_df = agents[agents["enforce_ready"]].copy()
self.non_enforce_ready_df = agents[not agents["enforce_ready"]].copy() self.non_enforce_ready_df = agents[~agents["enforce_ready"]].copy()
logger.info( logger.info(
f"Analysis complete: {len(self.enforce_ready_df)} enforce ready, " f"Analysis complete: {len(self.enforce_ready_df)} enforce ready, "
@@ -416,27 +510,6 @@ class QuietAgentWorkflowScreen(Screen):
self.app.notify(f"Analysis failed: {str(e)}", severity="error", timeout=5) self.app.notify(f"Analysis failed: {str(e)}", severity="error", timeout=5)
self._show_policy_selection() self._show_policy_selection()
def _update_analysis_status(self, status_text: str) -> None:
"""Update the analysis status message."""
try:
analyzing_msg = self.query_one("#analyzing_message", Static)
# Build updated message
updated_text = (
f"Analyzing Agent Activity\n"
f"{'=' * 50}\n\n"
f"Policy: {self.selected_policy.name}\n"
f"History Period: {self.history_days} days\n"
f"Quiet Threshold: {self.quiet_days} days\n\n"
f"Progress:\n"
f"{status_text}\n\n"
f"Please wait - this operation cannot be cancelled."
)
analyzing_msg.update(updated_text)
except Exception as e:
logger.debug(f"Could not update analysis status: {e}")
def _show_results(self) -> None: def _show_results(self) -> None:
"""Show the categorized results.""" """Show the categorized results."""
self.workflow_stage = "view_results" self.workflow_stage = "view_results"
@@ -821,7 +894,7 @@ class QuietAgentWorkflowScreen(Screen):
# Depending on stage, go back to previous stage or exit # Depending on stage, go back to previous stage or exit
if self.workflow_stage in ["select_policy", "view_results", "complete"]: if self.workflow_stage in ["select_policy", "view_results", "complete"]:
self.app.pop_screen() self.app.pop_screen()
elif self.workflow_stage == "select_quiet_days": elif self.workflow_stage == "select_history_days":
self._show_policy_selection() self._show_policy_selection()
elif self.workflow_stage == "select_enforce_target": elif self.workflow_stage == "select_enforce_target":
self._show_results() self._show_results()
+4 -4
View File
@@ -92,15 +92,15 @@ class MainMenuScreen(Screen):
BUTTON_DEFS = { BUTTON_DEFS = {
"agent_actions": [ "agent_actions": [
( (
"🖥️ - Find, Move, or Generate OTP for Agents", "🖥️ - Find agent, Move agent, or Generate One Time Pass",
"move_agent_workflow_button", "move_agent_workflow_button",
), ),
("🎫 - Review and appove OTP Activities", "otp_activities_button"), ("🎫 - Review and approve OTP Activities", "otp_activities_button"),
("🔕 - Find and Move Quiet Hosts to Enforcement", "find_quiet_button"), ("🛑 - Revoke Active OTP Session", "otp_revoke_button"),
], ],
"policy": [ "policy": [
("⚖️ - Prepare Policy For Enforcement", "policy_prep_button"), ("⚖️ - Prepare Policy For Enforcement", "policy_prep_button"),
("🛑 - Revoke OTPs", "otp_revoke_button"), ("🔕 - Find and Move Quiet Hosts to Enforcement", "find_quiet_button"),
], ],
} }
+51 -3
View File
@@ -26,10 +26,13 @@ dependencies = [
[[package]] [[package]]
name = "airlock_libs" name = "airlock_libs"
version = "5.0.1" version = "5.2.0"
dependencies = [ dependencies = [
"chrono", "chrono",
"crossbeam",
"flexi_logger",
"indicatif", "indicatif",
"log",
"mongodb", "mongodb",
"opentelemetry 0.18.0", "opentelemetry 0.18.0",
"opentelemetry-otlp", "opentelemetry-otlp",
@@ -504,6 +507,19 @@ version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b"
[[package]]
name = "crossbeam"
version = "0.8.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1137cd7e7fc0fb5d3c5a8678be38ec56e819125d8d7907411fe24ccb943faca8"
dependencies = [
"crossbeam-channel",
"crossbeam-deque",
"crossbeam-epoch",
"crossbeam-queue",
"crossbeam-utils",
]
[[package]] [[package]]
name = "crossbeam-channel" name = "crossbeam-channel"
version = "0.5.15" version = "0.5.15"
@@ -513,6 +529,16 @@ dependencies = [
"crossbeam-utils", "crossbeam-utils",
] ]
[[package]]
name = "crossbeam-deque"
version = "0.8.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9dd111b7b7f7d55b72c0a6ae361660ee5853c9af73f70c3c2ef6858b950e2e51"
dependencies = [
"crossbeam-epoch",
"crossbeam-utils",
]
[[package]] [[package]]
name = "crossbeam-epoch" name = "crossbeam-epoch"
version = "0.9.18" version = "0.9.18"
@@ -522,6 +548,15 @@ dependencies = [
"crossbeam-utils", "crossbeam-utils",
] ]
[[package]]
name = "crossbeam-queue"
version = "0.3.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0f58bbc28f91df819d0aa2a2c00cd19754769c2fad90579b3592b1c9ba7a3115"
dependencies = [
"crossbeam-utils",
]
[[package]] [[package]]
name = "crossbeam-utils" name = "crossbeam-utils"
version = "0.8.21" version = "0.8.21"
@@ -766,6 +801,19 @@ version = "0.4.2"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ce7134b9999ecaf8bcd65542e436736ef32ddca1b3e06094cb6ec5755203b80" checksum = "0ce7134b9999ecaf8bcd65542e436736ef32ddca1b3e06094cb6ec5755203b80"
[[package]]
name = "flexi_logger"
version = "0.31.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "31e5335674a3a259527f97e9176a3767dcc9b220b8e29d643daeb2d6c72caf8b"
dependencies = [
"chrono",
"log",
"nu-ansi-term",
"regex",
"thiserror 2.0.17",
]
[[package]] [[package]]
name = "fnv" name = "fnv"
version = "1.0.7" version = "1.0.7"
@@ -1557,9 +1605,9 @@ dependencies = [
[[package]] [[package]]
name = "log" name = "log"
version = "0.4.28" version = "0.4.29"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "34080505efa8e45a4b816c349525ebe327ceaa8559756f0356cba97ef3bf7432" checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897"
dependencies = [ dependencies = [
"value-bag", "value-bag",
] ]
+4 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "airlock_libs" name = "airlock_libs"
version = "5.0.1" version = "5.2.0"
edition = "2024" edition = "2024"
[lib] [lib]
@@ -25,6 +25,9 @@ tracing = "0.1.41"
tracing-subscriber = "0.3.20" tracing-subscriber = "0.3.20"
tracing-opentelemetry = "0.32.0" tracing-opentelemetry = "0.32.0"
pyo3-async-runtimes = { version = "0.27.0", features = ["async-std", "tokio"] } pyo3-async-runtimes = { version = "0.27.0", features = ["async-std", "tokio"] }
crossbeam = "0.8.4"
log = "0.4.29"
flexi_logger = "0.31.7"
[package.metadata.maturin] [package.metadata.maturin]
generate-abi-stubs = true generate-abi-stubs = true
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "maturin"
[project] [project]
name = "airlock_libs" name = "airlock_libs"
version = "5.0.1" version = "5.2.0"
description = "Airlock Digital API Wrapper" description = "Airlock Digital API Wrapper"
readme = "README.md" readme = "README.md"
license = { text = "AGPL-3.0-only" } license = { text = "AGPL-3.0-only" }
+1 -1
View File
@@ -1,7 +1,7 @@
use pyo3::prelude::*; use pyo3::prelude::*;
pub mod modules; pub mod modules;
pub mod services;
pub mod prelude; pub mod prelude;
pub mod services;
#[pymodule] #[pymodule]
fn airlock_libs(py: Python<'_>, m: &Bound<PyModule>) -> PyResult<()> { fn airlock_libs(py: Python<'_>, m: &Bound<PyModule>) -> PyResult<()> {
m.add_function(wrap_pyfunction!(services::pull_policy_exec_histories, py)?)?; m.add_function(wrap_pyfunction!(services::pull_policy_exec_histories, py)?)?;
+1 -1
View File
@@ -109,4 +109,4 @@ impl SkipBack {
let objectid_hex = format!("{}0000000000000000", hex_timestamp); let objectid_hex = format!("{}0000000000000000", hex_timestamp);
ObjectId::parse_str(&objectid_hex).expect("Invalid ObjectId hex") ObjectId::parse_str(&objectid_hex).expect("Invalid ObjectId hex")
} }
} }
+1 -1
View File
@@ -1 +1 @@
pub mod datatypes; pub mod datatypes;
+1 -1
View File
@@ -24,4 +24,4 @@ pub use std::{
io::{Read, Seek, SeekFrom}, io::{Read, Seek, SeekFrom},
path::PathBuf, path::PathBuf,
str::FromStr, str::FromStr,
}; };
+125 -115
View File
@@ -1,5 +1,8 @@
use crate::prelude::*;
use crate::modules::datatypes::*; use crate::modules::datatypes::*;
use crate::prelude::*;
use crossbeam::channel::unbounded;
use std::sync::{Arc, Mutex};
use std::thread;
#[pyfunction] #[pyfunction]
pub fn pull_policy_exec_histories( pub fn pull_policy_exec_histories(
py: Python<'_>, py: Python<'_>,
@@ -12,12 +15,12 @@ pub fn pull_policy_exec_histories(
ExtractedValues::Headers(h) => h, ExtractedValues::Headers(h) => h,
ExtractedValues::BaseUrl(_) => std::process::abort(), ExtractedValues::BaseUrl(_) => std::process::abort(),
}; };
let base_url = match PyData::convert(py, &py_self, false) { let base_url: String = match PyData::convert(py, &py_self, false) {
ExtractedValues::Headers(_) => std::process::abort(), ExtractedValues::Headers(_) => std::process::abort(),
ExtractedValues::BaseUrl(b) => b, ExtractedValues::BaseUrl(b) => b,
}; };
let handle = std::thread::spawn(move || { let handle: thread::JoinHandle<String> = std::thread::spawn(move || {
let rt = match tokio::runtime::Runtime::new() { let rt: tokio::runtime::Runtime = match tokio::runtime::Runtime::new() {
Ok(rt) => rt, Ok(rt) => rt,
Err(e) => { Err(e) => {
println!("Failed to build Tokio Runtime: {:?}", e); println!("Failed to build Tokio Runtime: {:?}", e);
@@ -27,14 +30,13 @@ pub fn pull_policy_exec_histories(
rt.block_on(async { rt.block_on(async {
let _ = init_tracer(); let _ = init_tracer();
}); });
let tracer = global::tracer("global_tracer"); let tracer: global::BoxedTracer = global::tracer("global_tracer");
let _cx = Context::new(); let _cx: Context = Context::new();
let file_path: PathBuf = format!( let file_path: PathBuf = format!(
"{}\\cache\\chunkinator.json", "{}\\cache\\chunkinator.json",
get_base_directory().display() get_base_directory().display()
) )
.into(); .into();
let writeable_filepath = file_path.clone();
if !&file_path.exists() { if !&file_path.exists() {
if let Some(parent_dir) = &file_path.parent() if let Some(parent_dir) = &file_path.parent()
&& !parent_dir.exists() && !parent_dir.exists()
@@ -55,13 +57,14 @@ pub fn pull_policy_exec_histories(
} }
} }
} }
let data = ApiResponse { let data: ApiResponse = ApiResponse {
error: "Success".to_string(), error: "Success".to_string(),
response: ExecHistories { response: ExecHistories {
exechistories: vec![], exechistories: vec![],
}, },
}; };
let data_write = serde_json::to_string_pretty(&data).expect("Failed to serialize"); let writeable_filepath: PathBuf = file_path.clone();
let data_write: String = serde_json::to_string_pretty(&data).expect("Failed to serialize");
match fs::write(writeable_filepath.clone(), data_write) { match fs::write(writeable_filepath.clone(), data_write) {
Ok(_) => {} Ok(_) => {}
Err(e) => { Err(e) => {
@@ -70,17 +73,18 @@ pub fn pull_policy_exec_histories(
} }
} }
let mut checkpoint_number: String = SkipBack::find_checkpoint(days).to_string(); let mut checkpoint_number: String = SkipBack::find_checkpoint(days).to_string();
let multi_progress = MultiProgress::new(); let progress_bar = Arc::new(Mutex::new(ProgressBar::new(100)));
multi_progress.set_draw_target(ProgressDrawTarget::stdout()); progress_bar
let progress_bar = multi_progress.add(ProgressBar::new(100)); .lock()
progress_bar.set_style( .unwrap()
.set_draw_target(ProgressDrawTarget::stderr());
progress_bar.lock().unwrap().set_style(
ProgressStyle::default_bar() ProgressStyle::default_bar()
.template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len}") .template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len} {message}")
.unwrap(), .unwrap(),
); );
progress_bar.enable_steady_tick(std::time::Duration::from_millis(100)); let client: Client = tracer.in_span("Building HTTP Client", |cx| {
let client = tracer.in_span("Building HTTP Client", |cx| { let client_result: Result<Client, reqwest::Error> = build_client(headers);
let client_result = build_client(headers);
match client_result { match client_result {
Ok(client_result) => { Ok(client_result) => {
cx.span().add_event( cx.span().add_event(
@@ -107,26 +111,84 @@ pub fn pull_policy_exec_histories(
} }
} }
}); });
let cutoff = Local::now().naive_local() - Duration::days(days); let cutoff: chrono::NaiveDateTime = Local::now().naive_local() - Duration::days(days);
let mut f = match File::open(&writeable_filepath) { let (tx, rx) = unbounded::<Vec<Group>>();
Ok(f) => f, let pb_clone = progress_bar.clone();
Err(e) => { thread::spawn(move || {
println!("Failed to Access {:?}: {}", &writeable_filepath, e); let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists()
std::process::abort(); {
let contents: String = fs::read_to_string(&writeable_filepath).unwrap_or_default();
let existing: ApiResponse =
serde_json::from_str(&contents).unwrap_or(ApiResponse {
error: "Success".to_string(),
response: ExecHistories {
exechistories: vec![],
},
});
existing
.response
.exechistories
.into_iter()
.map(|entry: Group| {
(
(
entry.sha256.clone(),
entry.filename.clone(),
entry.hostname.clone(),
),
entry,
)
})
.collect()
} else {
HashMap::new()
};
while let Ok(parsed_responses) = rx.recv() {
for executions in parsed_responses {
if executions.checkpoint.is_empty() || executions.datetime.is_empty() {
continue;
}
let history_date: NaiveDate = match NaiveDate::parse_from_str(
&executions.datetime.replace(" +0000 UTC", ""),
"%Y-%m-%dT%H:%M:%SZ",
) {
Ok(date) => date,
Err(_) => continue,
};
if history_date >= cutoff.into() {
let key: (String, String, String) = (
executions.sha256.clone(),
executions.filename.clone(),
executions.hostname.clone(),
);
seen.entry(key).or_insert(executions.clone());
}
}
let final_response: ApiResponse = ApiResponse {
error: "Success".to_string(),
response: ExecHistories {
exechistories: seen.values().cloned().collect(),
},
};
let data_write: String = serde_json::to_string_pretty(&final_response).unwrap();
match fs::write(&writeable_filepath, data_write) {
Ok(_) => {}
Err(e) => {
println!("Failed to write to: {:?}: {}", &writeable_filepath, e);
}
}
} }
}; });
let mut first_date: Option<NaiveDate> = None;
tracer.in_span("Airlock Data Retreival", |cx| { tracer.in_span("Airlock Data Retreival", |cx| {
let span = cx.span(); pb_clone
.lock()
.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").string(days.to_string()));
span.set_attribute(KeyValue::new("Policy Name", policy_names.clone())); span.set_attribute(KeyValue::new("Policy Name", policy_names.clone()));
loop { loop {
match f.seek(SeekFrom::Start(0)) {
Ok(_) => {}
Err(e) => {
println!("Failed to seek start of {:?}: {}", f, e);
std::process::abort();
}
}
let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| { let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| {
let results: ApiResponse = history_logging( let results: ApiResponse = history_logging(
&base_url, &base_url,
@@ -141,103 +203,51 @@ pub fn pull_policy_exec_histories(
)); ));
results results
}); });
let parsed_responses = execution_histories.response.exechistories; let parsed_responses: Vec<Group> = execution_histories.response.exechistories;
if parsed_responses.is_empty() { if parsed_responses.is_empty() {
break; break;
} }
let mut seen: HashMap<(String, String, String), Group> = tx.send(parsed_responses.clone()).unwrap();
if writeable_filepath.exists() { checkpoint_number = parsed_responses.last().unwrap().checkpoint.clone();
let mut contents = String::new(); if let Some(last_item) = parsed_responses.last()
f.read_to_string(&mut contents).unwrap();
let existing_data: ApiResponse =
serde_json::from_str(&contents).unwrap_or(ApiResponse {
error: "Success".to_string(),
response: ExecHistories {
exechistories: vec![],
},
});
existing_data
.response
.exechistories
.into_iter()
.map(|entry| {
(
(
entry.sha256.clone(),
entry.filename.clone(),
entry.hostname.clone(),
),
entry,
)
})
.collect()
} else {
HashMap::new()
};
for (index, executions) in parsed_responses.iter().enumerate() {
if executions.checkpoint.is_empty() || executions.datetime.is_empty() {
continue;
}
if index == parsed_responses.len() - 1 {
checkpoint_number = executions.checkpoint.clone();
break;
}
let history_date = match NaiveDate::parse_from_str(
&executions.datetime.replace(" +0000 UTC", ""),
"%Y-%m-%dT%H:%M:%SZ",
) {
Ok(date) => date,
Err(_) => continue,
};
if history_date >= cutoff.into() {
let key = (
executions.sha256.clone(),
executions.filename.clone(),
executions.hostname.clone(),
);
seen.entry(key).or_insert(executions.clone());
}
}
let final_response = ApiResponse {
error: "Success".to_string(),
response: ExecHistories {
exechistories: seen.values().cloned().collect(),
},
};
let data_write = serde_json::to_string_pretty(&final_response).unwrap();
match fs::write(&writeable_filepath, data_write) {
Ok(_) => {}
Err(e) => {
println!("Failed to write to: {:?}: {}", &writeable_filepath, e);
}
}
if let Some(last_item) = &final_response.response.exechistories.last()
&& let Ok(last_date) = NaiveDate::parse_from_str( && let Ok(last_date) = NaiveDate::parse_from_str(
&last_item.datetime.replace(" +0000 UTC", ""), &last_item.datetime.replace(" +0000 UTC", ""),
"%Y-%m-%dT%H:%M:%SZ", "%Y-%m-%dT%H:%M:%SZ",
) )
{ {
let date_diff = Local::now().naive_local().date() - last_date; if first_date.is_none() {
let percentage_diff = first_date = Some(last_date);
(days - date_diff.num_days()) as f64 / days as f64 * 100.0; }
progress_bar.set_position(percentage_diff.round() as u64); if let Some(base_date) = first_date {
progress_bar.set_message("Total Percent Complete"); let date_diff: chrono::TimeDelta = last_date - base_date;
let total_span: i64 =
(Local::now().naive_local().date() - base_date).num_days();
let percentage: u64 = ((date_diff.num_days() as f64 / total_span as f64)
* 100.0)
.clamp(0.0, 100.0)
.round() as u64;
pb_clone.lock().unwrap().set_position(percentage);
}
} }
} }
}); });
progress_bar.finish_with_message("All Checkpoints Complete"); progress_bar
let return_data = match fs::read_to_string(&writeable_filepath) { .lock()
.unwrap()
.finish_with_message("All Checkpoints Complete");
let return_data: String = match fs::read_to_string(file_path.clone()) {
Ok(return_data) => return_data, Ok(return_data) => return_data,
Err(e) => { Err(e) => {
println!("Failed to read data from: {:?}: {}", &writeable_filepath, e); println!("Failed to read data from: {:?}: {}", &file_path, e);
std::process::abort(); std::process::abort();
} }
}; };
drop(tx);
shutdown_tracer_provider(); shutdown_tracer_provider();
return_data.to_string() return_data.to_string()
}); });
let gil_value = handle.join().unwrap(); let gil_value: String = handle.join().unwrap();
Python::attach(|py| PyString::new(py, &gil_value).into()) Python::attach(|py: Python<'_>| PyString::new(py, &gil_value).into())
} }
fn build_client(headers: HeaderMap) -> Result<reqwest::Client, reqwest::Error> { fn build_client(headers: HeaderMap) -> Result<reqwest::Client, reqwest::Error> {
@@ -264,7 +274,7 @@ async fn history_logging(
}}"#, }}"#,
exec_types, checkpoint_number, policy_names exec_types, checkpoint_number, policy_names
); );
let res = client let res: Result<reqwest::Response, reqwest::Error> = client
.post(format!("{}/v1/logging/exechistories", base_url)) .post(format!("{}/v1/logging/exechistories", base_url))
.body(payload) .body(payload)
.send() .send()
@@ -305,13 +315,13 @@ pub fn get_base_directory() -> PathBuf {
} }
fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> { fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> {
let cfg = TelemetryConfig::load(); let cfg: TelemetryConfig = TelemetryConfig::load();
if !cfg.TELEMETRY { if !cfg.TELEMETRY {
global::set_tracer_provider(NoopTracerProvider::new()); global::set_tracer_provider(NoopTracerProvider::new());
return Ok(None); return Ok(None);
} }
let endpoint = cfg.TELEM_URL.unwrap_or_default(); let endpoint: String = cfg.TELEM_URL.unwrap_or_default();
let tracer = let tracer: sdktrace::Tracer =
opentelemetry_otlp::new_pipeline() opentelemetry_otlp::new_pipeline()
.tracing() .tracing()
.with_exporter( .with_exporter(
@@ -325,4 +335,4 @@ fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> {
.install_simple() .install_simple()
.unwrap(); .unwrap();
Ok(Some(tracer)) Ok(Some(tracer))
} }
+1 -1
View File
@@ -11,4 +11,4 @@ urllib3==2.5.0
pyperclip==1.11.0 pyperclip==1.11.0
--extra-index-url https://git.racooncity.org/api/packages/brotoskyj/pypi/simple/ --extra-index-url https://git.racooncity.org/api/packages/brotoskyj/pypi/simple/
airlock_libs==5.0.1 airlock_libs==5.2.0