Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
125 changes: 122 additions & 3 deletions extensions/business/dauth/dauth_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
from extensions.business.mixins.request_tracking_mixin import _RequestTrackingMixin
from extensions.business.dauth.dauth_mixin import _DauthMixin

__VER__ = '0.2.2'
__VER__ = '0.3.0'

_CONFIG = {
**BasePlugin.CONFIG,
Expand Down Expand Up @@ -77,8 +77,12 @@
"EE_CLOUDFLARE_TOKEN_LIVENESS_API",

"EE_MQTT_HOST_SEED" # this generates dynamically the EE_MQTT_HOST
],

],

"DAUTH_ORACLE_ONLY_SUPERVISOR_KEYS": [
"EE_CLOUDFLARE_TOKEN_DAUTH_MANAGER",
],

'VALIDATION_RULES': {
**BasePlugin.CONFIG['VALIDATION_RULES'],
},
Expand All @@ -99,12 +103,21 @@ class DauthManagerPlugin(

def __init__(self, **kwargs):
super(DauthManagerPlugin, self).__init__(**kwargs)
self._dauth_server_enabled = None
self._dauth_server_enabled_message = None
self._dauth_web_app_initialized = False
self._dauth_pause_teardown_succeeded = True
return



def on_init(self):
self._check_dauth_server_enabled_on_start()
super(DauthManagerPlugin, self).on_init()
self._dauth_web_app_initialized = True
if not self._is_dauth_server_enabled():
self.on_pause()
# endif
my_address = self.bc.address
my_eth_address = self.bc.eth_address
self.P("Started {} plugin on {} / {}\n - Auth keys: {}\n - Predefined keys: {}".format(
Expand All @@ -113,6 +126,106 @@ def on_init(self):
)
self._init_request_tracking()
return

def _check_dauth_server_enabled_on_start(self):
if getattr(self, "_dauth_server_enabled", None) is not None:
return self._dauth_server_enabled
# endif

error = None
try:
enabled = self.bc.is_dauth_oracle() is True
except Exception as e:
enabled = False
error = str(e)
# end try

message = None if enabled else error or "current node is not registered as a dAuth oracle"
self._dauth_server_enabled = enabled
self._dauth_server_enabled_message = message
if enabled:
self.P(f"{self.__class__.__name__} dAuth registry gate is enabled")
else:
self.P(
f"{self.__class__.__name__} dAuth registry gate is disabled. "
f"(cause: {message})",
color='r',
boxed=True
)
# endif
return enabled

def _is_dauth_server_enabled(self):
return getattr(self, "_dauth_server_enabled", None) is True

def should_pause(self):
return not self._is_dauth_server_enabled()

def should_resume(self):
return self._is_dauth_server_enabled()

def on_pause(self):
if not getattr(self, "_dauth_web_app_initialized", False):
return
# endif

self._dauth_pause_teardown_succeeded = False
self._stop_request_monitor.set()
if self._request_monitor_thread is not None:
self._request_monitor_thread.join(timeout=1.0)
# endif

self._maybe_close_start_commands()
running_commands = [
idx for idx, process in enumerate(self.start_commands_processes)
if process is not None and process.poll() is None
]
if running_commands:
raise RuntimeError(f"Failed to stop start commands {running_commands} while pausing")
if self._request_monitor_thread is not None and self._request_monitor_thread.is_alive():
raise RuntimeError("Failed to stop the FastAPI request monitor while pausing")
# endif

with self._incoming_lock:
self._incoming_requests.clear()
# endwith
self.postponed_requests.clear()
while True:
try:
self._server_queue.get(False)
except Exception:
break
# endwhile

self._maybe_read_and_stop_all_log_readers()
self.maybe_stop_tunnel_engine()
tunnel_stop_started = self.time()
while getattr(self, "tunnel_engine_started", False):
if self.time() - tunnel_stop_started >= 1.0:
raise RuntimeError("Failed to stop tunnel engine while pausing")
self.sleep(0.01)
# endwhile
self.reset_tunnel_engine()

nr_commands = len(self.get_start_commands())
self.start_commands_started = [False] * nr_commands
self.start_commands_finished = [False] * nr_commands
self.start_commands_processes = [None] * nr_commands
self.start_commands_start_time = [None] * nr_commands
self._dauth_pause_teardown_succeeded = True
return

def on_resume(self):
if not self._is_dauth_server_enabled():
raise RuntimeError("Cannot resume an ineligible dAuth server")
if not self._dauth_pause_teardown_succeeded:
raise RuntimeError("Cannot resume after an incomplete web app teardown")
# endif

self.failed = False
self._stop_request_monitor.clear()
self._start_request_monitor_thread()
return


def on_request(self, request):
Expand Down Expand Up @@ -212,6 +325,12 @@ def get_auth_data(self, body: dict):
}
}
"""
if not self._is_dauth_server_enabled():
response = self.__get_response({
'error': 'dAuth server is not registered as a dAuth oracle'
})
return response

try:
data = self.process_dauth_request(body)
except Exception as e:
Expand Down
16 changes: 14 additions & 2 deletions extensions/business/dauth/dauth_mixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -298,11 +298,24 @@ def fill_dauth_data(
# end if is_node
# end set node tags

# set the supervisor flag if this is identified as an oracle
# set supervisor secrets if this is a protocol oracle; dAuth-only secrets
# are additionally gated by the dAuth oracle registry.
if is_node and requester_node_address in oracles:
requester_is_dauth_oracle = False
try:
requester_is_dauth_oracle = self.bc.is_dauth_oracle(node_address_eth=sender_eth_address)
except Exception as e:
self.P(
f"dAuth oracle check failed for {sender_eth_address}; omitting supervisor keys: {e}",
color='r'
)
# end try

dauth_data["EE_SUPERVISOR"] = True
for key in self.cfg_supervisor_keys:
if isinstance(key, str) and len(key) > 0:
if key in self.cfg_dauth_oracle_only_supervisor_keys and not requester_is_dauth_oracle:
continue
dauth_data[key] = self.os_environ.get(key)
# end if
# end for
Expand Down Expand Up @@ -556,4 +569,3 @@ def process_dauth_request(self, body):
res = eng.process_dauth_request(request_sdk)
# res = eng.process_dauth_request(request_bad)
l.P(f"Result:\n{json.dumps(res, indent=2)}")

Loading