""" AETHER-TaskFlow Custom Algorithms. AETHER — Adaptive task scoring engine with reward-driven weight evolution. RAPTOR — Risk-Aware Priority-Tuned Operational Router: selects action type. AWFRO-X — Adaptive Waste-Free Resource Optimizer: recycles low-value states. """ from __future__ import annotations import math from typing import Any, Dict, List, Optional, Tuple # --------------------------------------------------------------------------- # AETHER – Adaptive Decision Core # --------------------------------------------------------------------------- class AETHER: """ Dynamically weighted task scorer. Weights evolve via a momentum-based gradient update driven by episodic rewards, forcing the agent to generalize rather than memorise. """ def __init__(self) -> None: self.weights: Dict[str, float] = { "priority": 2.0, "deadline_urgency": 1.5, "uncertainty_penalty": -1.2, "value": 1.0, "resource_fit": 0.8, } self._momentum: Dict[str, float] = {k: 0.0 for k in self.weights} self._lr: float = 0.05 self._beta: float = 0.9 # momentum coefficient self._episode_rewards: List[float] = [] self._step: int = 0 # ------------------------------------------------------------------ # Scoring # ------------------------------------------------------------------ def score( self, task: Dict[str, Any], resources: Dict[str, float], step: int, max_steps: int, ) -> float: """ Compute a composite urgency score for a task. Higher = more urgent/valuable to act on now. """ time_left = max(1, max_steps - step) # Deadline urgency: exponential decay — tasks due soon score higher deadline_urgency = math.exp(-task["deadline"] / max(1.0, time_left)) # Resource fit: can we actually execute this task? can_execute = ( resources.get("energy", 0) >= task["required_energy"] and resources.get("budget", 0) >= task["required_budget"] ) resource_fit = 1.0 if can_execute else -0.5 score = ( self.weights["priority"] * task["priority"] + self.weights["deadline_urgency"] * deadline_urgency + self.weights["uncertainty_penalty"] * task["uncertainty"] + self.weights["value"] * (task["value"] / 30.0) # normalise + self.weights["resource_fit"] * resource_fit ) return score def rank_tasks( self, tasks: List[Dict[str, Any]], resources: Dict[str, float], step: int, max_steps: int, ) -> List[Tuple[int, float]]: """Return list of (task_id, score) sorted descending.""" scored = [ (t["task_id"], self.score(t, resources, step, max_steps)) for t in tasks ] scored.sort(key=lambda x: x[1], reverse=True) return scored # ------------------------------------------------------------------ # Online weight update # ------------------------------------------------------------------ def update(self, reward: float) -> None: """Momentum-based weight update after each step.""" self._episode_rewards.append(reward) self._step += 1 # Compute a normalised advantage signal if len(self._episode_rewards) > 1: mean_r = sum(self._episode_rewards) / len(self._episode_rewards) std_r = ( sum((r - mean_r) ** 2 for r in self._episode_rewards) / len(self._episode_rewards) ) ** 0.5 advantage = (reward - mean_r) / max(std_r, 1e-6) else: advantage = reward # Update each weight with momentum for key in self.weights: grad = advantage * self._lr self._momentum[key] = ( self._beta * self._momentum[key] + (1 - self._beta) * grad ) self.weights[key] += self._momentum[key] # Clamp weights to sensible ranges self.weights["priority"] = max(0.5, min(4.0, self.weights["priority"])) self.weights["deadline_urgency"] = max(0.3, min(3.0, self.weights["deadline_urgency"])) self.weights["uncertainty_penalty"] = max(-3.0, min(-0.1, self.weights["uncertainty_penalty"])) self.weights["value"] = max(0.2, min(2.0, self.weights["value"])) self.weights["resource_fit"] = max(0.1, min(2.0, self.weights["resource_fit"])) def reset(self) -> None: self._episode_rewards = [] self._step = 0 # --------------------------------------------------------------------------- # RAPTOR – Execution Strategy Engine # --------------------------------------------------------------------------- class RAPTOR: """ Risk-Aware Priority-Tuned Operational Router. Decides *how* to act on the highest-scored task based on current resource levels, task uncertainty, and deadline pressure. """ def decide( self, task: Dict[str, Any], resources: Dict[str, float], step: int, max_steps: int, ) -> str: """ Return the optimal action type for the given task + resource state. Decision logic (priority order): 1. If resources are critically low → defer 2. If uncertainty is very high → optimize first 3. If deadline is imminent and resources sufficient → execute 4. If task can be delegated cheaply → delegate 5. Default → execute """ time_left = max_steps - step energy = resources.get("energy", 0.0) budget = resources.get("budget", 0.0) uncertainty = task.get("uncertainty", 0.0) deadline = task.get("deadline", 5) req_energy = task.get("required_energy", 1.0) req_budget = task.get("required_budget", 5.0) # Critical resource shortage if energy < req_energy * 0.5 or budget < req_budget * 0.5: if time_left > 2: return "defer" else: return "delegate" # Very high uncertainty – optimize first to reduce risk if uncertainty > 0.75 and time_left > 1: return "optimize" # Imminent deadline – must act now if deadline <= 1 and energy >= req_energy and budget >= req_budget: return "execute" # Low value + sufficient time → delegate to save resources if task.get("value", 10) < 8.0 and time_left > 3: return "delegate" # Sufficient resources – execute if energy >= req_energy and budget >= req_budget: return "execute" # Fallback return "defer" # --------------------------------------------------------------------------- # AWFRO-X – Adaptive Waste-Free Resource Optimizer # --------------------------------------------------------------------------- class AWFROX: """ Converts low-value deferred states into usable outcomes. Filters the task queue to remove tasks that are guaranteed to fail (e.g. deadline passed, insufficient resources with no recovery path) and recycles deferred tasks back into the active queue if conditions improve. """ def filter_viable( self, tasks: List[Dict[str, Any]], resources: Dict[str, float], step: int, max_steps: int, ) -> List[Dict[str, Any]]: """Remove tasks that cannot possibly be completed.""" time_left = max_steps - step viable = [] for task in tasks: # Deadline already passed if task.get("deadline", 0) < 0: continue # No time left if time_left <= 0: continue viable.append(task) return viable def recycle_deferred( self, active: List[Dict[str, Any]], deferred: List[Dict[str, Any]], resources: Dict[str, float], step: int, ) -> Tuple[List[Dict[str, Any]], List[Dict[str, Any]]]: """ Requeue deferred tasks when resources recover. Returns (updated_active, updated_deferred). """ still_deferred = [] for task in deferred: can_execute = ( resources.get("energy", 0) >= task.get("required_energy", 1.0) * 0.8 and resources.get("budget", 0) >= task.get("required_budget", 1.0) * 0.8 and task.get("deadline", 0) > 0 ) if can_execute: task["status"] = "pending" active.append(task) else: still_deferred.append(task) return active, still_deferred