Skip to content
Open
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
1 change: 1 addition & 0 deletions Dockerfile
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
FROM python:3.9-slim-bookworm
WORKDIR /app
COPY requirements.txt ./
RUN apt-get update && apt-get install build-essential -y
RUN pip install --upgrade pip
RUN pip install -r requirements.txt
COPY src ./
Expand Down
2 changes: 2 additions & 0 deletions requirements.txt
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ annotated-types==0.6.0
certifi==2024.2.2
charset-normalizer==3.3.2
click==8.1.7
google-cloud-firestore==2.10.0
idna==3.7
markdown-it-py==3.0.0
mdurl==0.1.2
Expand All @@ -18,3 +19,4 @@ typing_extensions==4.11.0
urllib3==2.2.1
watchdog==4.0.0
sdnotify==0.3.2
pyzmq==25.1.2
33 changes: 33 additions & 0 deletions src/communication_module.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
#!/usr/bin/env python3

import logging
import zmq

logger = logging.getLogger(__name__)


class CommunicationModule(object):

def __init__(self, host="127.0.0.1", port=5555):
self.host = host
self.port = port
self.context = zmq.Context()
self.socket = self.context.socket(zmq.REQ)

def start(self):
logger.debug(
f"Starting communication module ('{self.host}:{self.port}')")
self.socket.connect(f"tcp://{self.host}:{self.port}")

def stop(self):
logger.debug("Stopping communication module")
self.socket.close()
self.context.term()

def execute_command(self, command, params=None):
if not command:
logger.error("The command argument is required")
return
self.socket.send_json({"command": command, "params": params})
response = self.socket.recv_json()
return response
167 changes: 167 additions & 0 deletions src/firestore_client_handler.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,167 @@
#!/usr/bin/env python3

from google.cloud.firestore import Client
from google.cloud.firestore_v1 import watch
from google.oauth2.credentials import Credentials
import google.api_core
import json
import logging
import requests
import threading
import time

logger = logging.getLogger(__name__)

watch._should_recover = lambda _: False
watch._should_terminate = lambda _: True


class FirestoreClientHandler:

REFRESH_TOKEN_TIME_MIN = 58
INITIALIZATION_RETRY_TIME_MIN = 1
SERVER_RESPONSE_TIMEOUT_S = 60

def __init__(self, api_key: str, project_id: str, refresh_token: str):
self.api_key = api_key
self.project_id = project_id
self.refresh_token = refresh_token
self.token_expired_timer = None
self.initialization_retry_timer = None
self.server_responde_message_handler = None
self.user_id = None
self.client = None

def initialize_client(self, notify=True):
if self.client:
if notify:
threading.Thread(
target=self.on_client_initialized, daemon=True).start()
return

logger.debug("Initializing Firestore client")
creds = self._get_credentials()

if not creds:
if self.initialization_retry_timer:
return
self.initialization_retry_timer = threading.Timer(
self.INITIALIZATION_RETRY_TIME_MIN*60.0, self.on_initialization_retry)
self.initialization_retry_timer.start()
return

self.client = Client(self.project_id, creds)
logger.debug("Firebase client initialized")
if notify:
threading.Thread(target=self.on_client_initialized,
daemon=True).start()
self.server_responde_message_handler = ServerResponseMessageHandler(
timeout_s=self.SERVER_RESPONSE_TIMEOUT_S, on_server_not_responding=self.on_server_not_responding)
google.api_core.bidi._LOGGER = BidiCustomLogger()
google.api_core.bidi._LOGGER.addHandler(
self.server_responde_message_handler)
self.token_expired_timer = threading.Timer(
self.REFRESH_TOKEN_TIME_MIN*60.0, self.on_token_expired)
self.token_expired_timer.start()

def stop_client(self):
logger.debug("Stopping Firestore client")
self.client = None
if self.server_responde_message_handler:
self.server_responde_message_handler.stop()
self.server_responde_message_handler = None
if self.initialization_retry_timer:
self.initialization_retry_timer.cancel()
self.initialization_retry_timer.join()
self.initialization_retry_timer = None
if self.token_expired_timer:
self.token_expired_timer.cancel()
self.token_expired_timer.join()
self.token_expired_timer = None

def on_client_initialized(self):
logger.warning(
"This function should be overwritten by the child class")

def on_server_not_responding(self):
logger.warning(
"This function should be overwritten by the child class")

def on_token_expired(self):
logger.warning(
"This function should be overwritten by the child class")

def on_initialization_retry(self):
logger.warning("Retrying firebase client initialization...")
self.initialization_retry_timer = None
threading.Thread(target=self.initialize_client).start()

def _get_credentials(self):
user_id, token_id = self._get_ids()
if not user_id or not token_id:
logger.debug(f"Invalid user and token ids ({user_id}, {token_id})")
return None
self.user_id = user_id
creds = Credentials(token_id, self.refresh_token)
return creds

def _get_ids(self):
token_response = self._get_token_response()
user_id = token_response.get("user_id")
token_id = token_response.get("id_token")
return (user_id, token_id)

def _get_token_response(self):
request_ref = f"https://securetoken.googleapis.com/v1/token?key={self.api_key}"
headers = {"content-type": "application/json; charset=UTF-8"}
data = json.dumps({"grantType": "refresh_token",
"refreshToken": self.refresh_token})
try:
response_object = requests.post(
request_ref, headers=headers, data=data)
return response_object.json()
except Exception as e:
return {}


class BidiCustomLogger(logging.Logger):

def __init__(self, name="google.api_core.bidi", level=logging.DEBUG):
super().__init__(name, level)


class ServerResponseMessageHandler(logging.StreamHandler):

def __init__(self, timeout_s=60, on_server_not_responding=lambda _: None, server_response_msg="recved response.", watchdog_timeout_msg="watchdog timeout"):
super().__init__()
self.timeout_s = timeout_s
self.on_server_not_responding = on_server_not_responding
self.server_response_msg = server_response_msg
self.server_response_last_time = time.time()
self.watchdog_timeout_msg = watchdog_timeout_msg
self.not_responding_timer = None
self.setFormatter(logging.Formatter(
'%(asctime)s %(levelname)-8s - %(name)-16s - %(message)s', None, '%'))

def emit(self, record):
msg = self.format(record)
if self.watchdog_timeout_msg in msg:
logger.debug("Server connection timeout detected")
threading.Thread(target=self.on_server_not_responding).start()
if self.server_response_msg in msg:
logger.debug("Server response message caught ({} seconds)".format(
time.time() - self.server_response_last_time))
self.server_response_last_time = time.time()
if self.not_responding_timer:
self.stop()
self.not_responding_timer = threading.Timer(
self.timeout_s, self.on_server_not_responding)
self.not_responding_timer.start()
if record.levelno >= logger.getEffectiveLevel():
super().emit(record)

def stop(self):
if self.not_responding_timer:
self.not_responding_timer.cancel()
self.not_responding_timer.join()
self.not_responding_timer = None
123 changes: 80 additions & 43 deletions src/installed_service_handler.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import logging
import os
import pathlib
from threading import Timer
from threading import Timer, Thread
from typing import List, Optional

from python_on_whales import DockerClient, docker
Expand Down Expand Up @@ -36,66 +36,95 @@ class InstalledServiceHandler(FileSystemEventHandler):
docker: Optional[DockerClient]
"""The docker client of the service compose"""

def __init__(self, service_path: str, wait_seconds: float):
def __init__(self, service_path: str, wait_seconds: float, service_state_change_callback: callable=lambda _: None):
self.service_path = service_path
self.service_name = service_path.split("/")[-1]
self.wait_seconds = wait_seconds
self.service_state_change_callback = service_state_change_callback
self.timer = None
self.observer = Observer()
self.observer = None
self.files = os.listdir(service_path)
self.docker = self._get_docker()

def start(self):
def start(self, mode: str="offline"):
"""Start the handler by starting the observer of the service folder."""
logger.debug(f"Starting '{self.service_name}' Installed Service Handler.")
try:
self.observer.schedule(self, self.service_path, recursive=True)
self.observer.start()
logger.debug(f"{self.service_name} Installed Service Handler started.")
if not self.is_running():
logger.error(f"Service '{self.service_name}' is not running, starting it up")
self.up_compose_services()
if mode == "offline":
self.observer = Observer()
self.observer.schedule(self, self.service_path, recursive=True)
self.observer.start()
except FileNotFoundError:
logger.error(
f"Couldn't start a watcher for the {self.service_name} service, the folder no longer exists."
)
self.down()
self.down_compose_services()
self.clean_compose_services()

def stop(self):
"""Stop the handler by stopping the observer of the service folder."""
logger.debug("Installed Service Handler stopped.")
self.observer.stop()

def up(self):
"""Start the compose of the service by doing "docker compose up"."""
if self.docker:
try:
self.docker.compose.up(detach=True)
logger.info(f"{self.service_name} service compose started.")
except:
logger.error(
f"An error occurred, couldn't start service {self.service_name}."
)

def down(self):
"""Stop the compose of the service.
logger.debug(f"Stopping '{self.service_name}' Installed Service Handler.")
if self.observer:
self.observer.stop()

def up_compose_services(self):
"""Start the services of the compose file by executing "docker compose pull" and "docker compose up"."""
if not self.docker:
return
try:
Thread(target=self.service_state_change_callback, args=[self.service_name, "pulling"]).start()
self.docker.compose.pull()
Thread(target=self.service_state_change_callback, args=[self.service_name, "starting"]).start()
self.docker.compose.up(detach=True, pull="never")
Thread(target=self.service_state_change_callback, args=[self.service_name, "started"]).start()
logger.info(f"'{self.service_name}' compose services started.")
except Exception as e:
logger.error(
f"An error occurred, couldn't start service {self.service_name}: {e}."
)
Thread(target=self.service_state_change_callback, args=[self.service_name, "unknown"]).start()

def down_compose_services(self, remove_images: bool = False):
"""Shut down the services of the compose file by executing "docker compose down"."""
if not self.docker:
return
self.docker.compose.down(remove_images="all" if remove_images else None)
logger.info(f"'{self.service_name}' compose services downed.")

def stop_compose_services(self):
"""Stop the services of the compose file by executing "docker compose stop"."""
if not self.docker:
return
self.docker.compose.stop()
logger.info(f"'{self.service_name}' compose services stopped.")

def clean_compose_services(self):
"""Clean the services of the compose file.

This is done by getting the containers started by the compose file and killing them.
The volumes are pruned to remove any unused ones.

This is done like this because, when the service folder is removed the compose needs to stop.
But, in that case, the compose file no longer exists, so "docker compose down" can't be called.
This is done like this because, when the service folder is removed and, as such, the compose
file no longer exists, "docker compose down" can not be directly called.
"""
if self.docker:
logger.info(f"{self.service_name} service compose stopped.")
containers = self.docker.ps(
filters={
"label": f"com.docker.compose.project.config_files={self.service_path}/{self.compose_file}"
}
)
for container in containers:
docker.container.stop(container)
docker.container.remove(container)

docker.volume.prune()

def reload_service_compose(self):
if not self.docker:
return
containers = self.docker.ps(
filters={
"label": f"com.docker.compose.project.config_files={self.service_path}/{self.compose_file}"
}
)
for container in containers:
docker.container.stop(container)
docker.container.remove(container)

docker.volume.prune()
logger.info(f"'{self.service_name}' compose services cleaned.")

def reload_compose_services(self):
"""Stop the compose, get the docker client again and start the compose.

The docker client is loaded again because, if the compose file is changed, the docker client will not be valid.
Expand All @@ -105,12 +134,20 @@ def reload_service_compose(self):

logger.debug(f"{self.service_name} service modified.")

self.down()
self.down_compose_services()
self.docker = self._get_docker()
self.up()
self.up_compose_services()

self.timer = None

def is_running(self):
"""Check if the service is running"""
running_services = self.docker.compose.ls()
for running_service in running_services:
if running_service.name == self.service_name and running_service.running == 1:
return True
return False

def on_any_event(self, event: FileSystemEvent):
"""Reload the service when the service changes.

Expand All @@ -133,7 +170,7 @@ def on_any_event(self, event: FileSystemEvent):
self.timer.cancel()
self.timer = None

self.timer = Timer(self.wait_seconds, self.reload_service_compose)
self.timer = Timer(self.wait_seconds, self.reload_compose_services)
self.timer.start()

def _get_compose_file_name(self):
Expand Down
Loading