diff --git a/Changelog.md b/Changelog.md old mode 100644 new mode 100755 diff --git a/README.md b/README.md old mode 100644 new mode 100755 diff --git a/data_ingestion/__init__.py b/data_ingestion/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/data_ingestion/ingest_data.py b/data_ingestion/ingest_data.py new file mode 100644 index 0000000..826ac7b --- /dev/null +++ b/data_ingestion/ingest_data.py @@ -0,0 +1,88 @@ +import requests +import pandas as pd +from loguru import logger + +BASE_URL = "http://127.0.0.1:8001" + +def get_users(): + """Fetches all users from the API and returns them as a pandas DataFrame.""" + all_users = [] + page = 1 + while True: + response = requests.get(f"{BASE_URL}/users?page={page}&size=100") + if response.status_code == 200: + data = response.json() + users = data["items"] + if not users: + break + all_users.extend(users) + page += 1 + logger.info(f"Fetched page {page} of users, total users fetched: {len(all_users)}") + else: + logger.warning(f"Failed to fetch users: {response.status_code}") + logger.error(f"Error details: {response.text}") + break + return pd.DataFrame(all_users) + + +def get_tracks(): + """Fetches all tracks from the API and returns them as a pandas DataFrame.""" + all_tracks = [] + page = 1 + while True: + response = requests.get(f"{BASE_URL}/tracks?page={page}&size=100") + if response.status_code == 200: + data = response.json() + tracks = data["items"] + if not tracks: + break + all_tracks.extend(tracks) + page += 1 + logger.info(f"Fetched page {page} of tracks, total tracks fetched: {len(all_tracks)}") + else: + logger.warning(f"Failed to fetch tracks: {response.status_code}") + logger.error(f"Error details: {response.text}") + break + return pd.DataFrame(all_tracks) + +def get_listen_history(): + """Fetches all listen history from the API and returns it as a pandas DataFrame.""" + all_listen_history = [] + page = 1 + while True: + response = requests.get(f"{BASE_URL}/listen_history?page={page}&size=100") + if response.status_code == 200: + data = response.json() + listen_history = data["items"] + if not listen_history: + break + all_listen_history.extend(listen_history) + page += 1 + logger.info(f"Fetched page {page} of listen history, total records fetched: {len(all_listen_history)}") + else: + logger.warning(f"Failed to fetch listen history: {response.status_code}") + logger.error(f"Error details: {response.text}") + break + return pd.DataFrame(all_listen_history) + +def main(): + """init log file """ + logger.add("data_ingestion.log", rotation="1 MB", level="INFO") + logger.info("Starting data ingestion process") + logger.info("##########################################################################") + """Main function to ingest data from all endpoints.""" + users_df = get_users() + tracks_df = get_tracks() + listen_history_df = get_listen_history() + + print("Users DataFrame:") + print(users_df.head()) + print("\nTracks DataFrame:") + print(tracks_df.head()) + print("\nListen History DataFrame:") + print(listen_history_df.head()) + logger.info("Data ingestion process completed successfully") + logger.info("##########################################################################") + +if __name__ == "__main__": + main() diff --git a/docs/ANSWERS.md b/docs/ANSWERS.md old mode 100644 new mode 100755 index d038faa..ee77320 --- a/docs/ANSWERS.md +++ b/docs/ANSWERS.md @@ -1,6 +1,24 @@ # Réponses du test ## _Utilisation de la solution (étape 1 à 3)_ +etape 1: +utilisation de uv en installant une version 3.9 puis activation de cette environement : source .env_uv_3919/bin/activate +installation des paquets uv pip install -r requirements.txt + +etape 2: +lancement du serveur: +une erreur dans les methodes getters du main utilise la classe Page qui n'avait pas les methodes get_tracks, get_users et get_listen_history +modification du code en utilisant la classe CustomPage a la place de Page +le port 8000 etant occupe, +j'ai redirige en local le serveur sur le port 8001 avec python -m uvicorn main:app --port 8001 + +Le serveur demarre et l'url http://127.0.0.1:8001/docs affiche l'API + +etape 3 : +je teste trois aspects : +- la reponse 200 du serveur +- la reponse en terme de data +- le schema de la data retourne _Inscrire la documentation technique_ @@ -8,16 +26,133 @@ _Inscrire la documentation technique_ ### Étape 4 -_votre réponse ici_ +nous partons de la data d'origine et du besoin de faire des recommandations a partir des chansons ecoutes. +une chanson seraecoute a une date d par un ou plusieurs users. Un user ecoutera une ou plusieurs chansons. En modele relationnel , nous aurions une table pivot qui correspondrait au coeur de ce qui va faire l'objet des calculs pour les recommandations : la table user_history. +Dans votre cas, la table d'association pourrait être `user_listen_history` (ou `ecoutes`): + +**Tables principales :** + +* **`users`** + * `user_id` (PK) + * `user_name` + * ... + +* **`songs`** + * `song_id` (PK) + * `title` + * `artist` + * ... + +**Table d'association :** + +* **`user_listen_history`** + * `listen_id` (PK, optionnel, mais bonne pratique pour une clé primaire unique) + * `user_id` (FK vers `users.user_id`) + * `song_id` (FK vers `songs.song_id`) + * `listen_timestamp` (Quand l'écoute a eu lieu - très important pour les recommandations basées sur l'historique récent) + * `play_duration_seconds` (Durée d'écoute, si pertinent) + * ... + +Avec cette structure, vous pouvez facilement savoir quelles chansons un utilisateur a écoutées, et quels utilisateurs ont écouté une chanson donnée, ainsi que des détails sur chaque événement d'écoute. + +Ainsi ce modele peut etre porte par un SGDBR comme Postgresql. + +Les donnees massives issue de l'API , elles seraint stocker dans un data lake tel AWS S3. +Quand aux donnees issue des calculs de machine learning seraient stocker dans une base de donnees cle valeur comme MongoDB: +"Calcul des recommandations terminé et servi via MongoDB" +``` + +**Exemple de document dans MongoDB :** + +```json +{ + "_id": "user_12345", // La clé est le user_id pour un lookup O(1) + "recommendations": [ + "song_id_abc", + "song_id_def", + "song_id_ghi", + // ... + ], + "model_version": "v1.2.3", + "calculation_date": "2025-08-04T22:00:00Z" +} +``` +ainsi , on a un workflow hybride Data lake de donnees massives sur un cloud + les donnees maitres sur une base relationnelles + les resultats des calculs de recommandations fait a partir de la table user_history_listen ### Étape 5 -_votre réponse ici_ +le systeme de surveillance de la sante du pipel;ine va reposer sur : +1- un systeme de logging , etape par etape , de chaque composant du pipeline +2- un dashboard via un outil BI pour visualiser les metriques ce qui permet d'avoir une big picture du pipeline +3- des alertes envoyeant aux intervenants cles lies aux actions cles du pipelines dans un canala specifique de communication(email, tchat ...) + +venons en aux metriques: +- les metriques lie au pipeline , aux ressources materielles consommes, aux statutss desou ko des etapes du pipeline. + +- les metriques lie a la data : le nombre de lignes , la preence de null , la coherence du schema , les types de donnees + + ### Étape 6 -_votre réponse ici_ +Pour automatiser le calcul des recommandations de type scoring , je mettrais en place un workflow orchestré qui se déclenche après la réussite du pipeline d'ingestion de données. + +**Architecture et Flux de Travail :** + +``` +[FIN du Pipeline d'Ingestion] -> [DÉCLENCHEMENT du Workflow de Calcul] + | + V +[1. Tâche de Calcul Batch (ex: Spark,..)] + | a. Charge le dernier modèle validé + | b. Charge les nouvelles données d'historique des utilisateurs. + | c. Pour chaque utilisateur, calcule les N recommandations. + | d. Écrit les résultats dans une table "staging". + | + V +[2. Tâche de Validation des Résultats] + | a. Vérifie la qualité des recommandations (ex: nombre d'utilisateurs avec des recommandations, pas de valeurs nulles). + | + V +[3. Tâche de Keep or Replace] + | a. Si la validation est réussie, remplace l'ancienne table de recommandations par la nouvelle (staging). + | b. Archive l'ancienne table pour analyse. + | + V +[4. Notification] -> [Canal Slack/Email: "Calcul des recommandations terminé avec succès"] ### Étape 7 -_votre réponse ici_ +la le reentrainement du modele de recommandation va necessite d'evaluer le modele existant et le modele issue de donnees recentes. Ces evaluations vontse faire sur les metriques de performances des modeles selon un jeu de donnees test. + +[DÉCLENCHEUR (ex: Hebdomadaire, ou baisse de performance)] + | + V +[1. Tâche d'Extraction et Préparation des Données] + | a. Crée un jeu de données d'entraînement (ex: 90 derniers jours) et de test (ex: 7 derniers jours). + | + V +[2. Tâche d'Entraînement Parallèle] + | a. Entraîne un nouveau "modèle candidat" sur le jeu de données d'entraînement. + | + V +[3. Tâche d'Évaluation Comparative] + | a. Charge le "modèle en production" actuel. + | b. Évalue le "modèle candidat" ET le "modèle en production" sur le même jeu de données de test. + | c. Compare leurs métriques de performance. + | + V +[4. Tâche de Validation et Enregistrement (Conditionnelle)] + | a. SI le candidat est meilleur que la production (selon un seuil défini, ex: +5% de précision): + | i. Enregistre le "modèle candidat". + | ii. Attribue au nouveau modèle le tag "staging" ou "validation". + | b. SINON: + | i. Garde le modèle en production et alerte l'équipe (le modèle n'apprend plus). + | + V +[5. Tâche de Déploiement/Promotion (Manuelle ou Automatique)] + | a. Un Data Scientist valide les métriques du modèle en "staging". + | b. Promotion du modèle : le tag passe de "staging" à "production". + | + V +[6. Notification] -> [Rapport de réentraînement envoyé sur Slack/Email] diff --git a/notebooks/01_users_ingestion.ipynb b/notebooks/01_users_ingestion.ipynb new file mode 100644 index 0000000..89ccda4 --- /dev/null +++ b/notebooks/01_users_ingestion.ipynb @@ -0,0 +1,77 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 1, + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + " id first_name last_name email gender \\\n", + "0 8231 Gabrielle Mcguire melinda03@example.org Male \n", + "1 30893 Morgan Hill jshaw@example.com Bigender \n", + "2 13477 Tiffany Wiley crosslisa@example.com Gender questioning \n", + "3 24153 Jonathan Kerr terryjessica@example.com Genderqueer \n", + "4 16236 Robert Hensley dennisoliver@example.net Female \n", + "\n", + " favorite_genres created_at updated_at \n", + "0 Hip Hop 2025-01-05T05:13:07.505433 2025-07-30T15:54:19.063299 \n", + "1 Jazz 2023-10-17T19:55:22.351170 2025-04-28T09:18:47.655815 \n", + "2 Country 2023-12-11T06:51:59.356696 2025-06-28T22:13:55.643900 \n", + "3 Indie 2024-12-04T01:05:42.812122 2024-12-27T23:50:52.591169 \n", + "4 Blues 2024-08-14T02:08:12.199497 2025-01-02T19:13:52.554275 \n" + ] + } + ], + "source": [ + "import requests\n", + "import pandas as pd\n", + "\n", + "BASE_URL = \"http://127.0.0.1:8001\"\n", + "\n", + "def get_users():\n", + " all_users = []\n", + " page = 1\n", + " while True:\n", + " response = requests.get(f\"{BASE_URL}/users?page={page}&size=100\")\n", + " if response.status_code == 200:\n", + " data = response.json()\n", + " users = data[\"items\"]\n", + " if not users:\n", + " break\n", + " all_users.extend(users)\n", + " page += 1\n", + " else:\n", + " print(f\"Failed to fetch users: {response.status_code}\")\n", + " break\n", + " return pd.DataFrame(all_users)\n", + "\n", + "users_df = get_users()\n", + "print(users_df.head())" + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": ".env_uv_3919", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.9.19" + } + }, + "nbformat": 4, + "nbformat_minor": 4 +} diff --git a/notebooks/02_tracks_ingestion.ipynb b/notebooks/02_tracks_ingestion.ipynb new file mode 100644 index 0000000..9fbc2c5 --- /dev/null +++ b/notebooks/02_tracks_ingestion.ipynb @@ -0,0 +1,77 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 1, + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + " id name artist songwriters duration genres \\\n", + "0 53589 because Cody Davis Thomas Mccarthy 41:21 rich \n", + "1 3938 up Marcia Williamson Richard Benson 16:18 response \n", + "2 86307 find Amber Spencer Henry Navarro 26:55 available \n", + "3 96800 sell Brianna Newman DDS Brett Nguyen 24:21 unit \n", + "4 96699 believe Jesse Larsen Angela Travis 55:42 green \n", + "\n", + " album created_at updated_at \n", + "0 movie 2024-05-29T14:58:48.794863 2024-09-21T11:46:29.933336 \n", + "1 talk 2024-09-01T15:28:21.664984 2025-04-27T14:58:42.430469 \n", + "2 audience 2024-12-08T09:55:31.602019 2025-06-11T10:49:13.635377 \n", + "3 hold 2024-09-28T11:15:18.362727 2024-08-24T06:22:06.328800 \n", + "4 evidence 2024-05-22T00:51:45.434858 2025-07-04T05:58:41.597525 \n" + ] + } + ], + "source": [ + "import requests\n", + "import pandas as pd\n", + "\n", + "BASE_URL = \"http://127.0.0.1:8001\"\n", + "\n", + "def get_tracks():\n", + " all_tracks = []\n", + " page = 1\n", + " while True:\n", + " response = requests.get(f\"{BASE_URL}/tracks?page={page}&size=100\")\n", + " if response.status_code == 200:\n", + " data = response.json()\n", + " tracks = data[\"items\"]\n", + " if not tracks:\n", + " break\n", + " all_tracks.extend(tracks)\n", + " page += 1\n", + " else:\n", + " print(f\"Failed to fetch tracks: {response.status_code}\")\n", + " break\n", + " return pd.DataFrame(all_tracks)\n", + "\n", + "tracks_df = get_tracks()\n", + "print(tracks_df.head())" + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": ".env_uv_3919", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.9.19" + } + }, + "nbformat": 4, + "nbformat_minor": 4 +} diff --git a/notebooks/03_listen_history_ingestion.ipynb b/notebooks/03_listen_history_ingestion.ipynb new file mode 100644 index 0000000..6083e40 --- /dev/null +++ b/notebooks/03_listen_history_ingestion.ipynb @@ -0,0 +1,77 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 2, + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + " user_id items created_at \\\n", + "0 8231 [48603, 65779, 39459, 71843, 50528] 2023-10-09T02:44:43.782155 \n", + "1 30893 [39572, 75194, 53485, 78937, 82035] 2023-12-02T00:09:05.153940 \n", + "2 13477 [75096, 36601, 78265, 20994, 31979] 2024-03-15T09:31:51.040543 \n", + "3 24153 [95393, 19603, 5042, 23927, 39632] 2025-03-18T14:09:17.945575 \n", + "4 16236 [14367, 77470, 97452, 13954, 52303] 2024-02-01T17:14:32.075974 \n", + "\n", + " updated_at \n", + "0 2024-09-23T15:39:21.776922 \n", + "1 2025-05-28T13:14:43.812083 \n", + "2 2025-08-04T04:11:56.050109 \n", + "3 2025-03-24T23:21:32.534180 \n", + "4 2024-04-06T03:02:39.391084 \n" + ] + } + ], + "source": [ + "import requests\n", + "import pandas as pd\n", + "\n", + "BASE_URL = \"http://127.0.0.1:8001\"\n", + "\n", + "def get_listen_history():\n", + " all_listen_history = []\n", + " page = 1\n", + " while True:\n", + " response = requests.get(f\"{BASE_URL}/listen_history?page={page}&size=100\")\n", + " if response.status_code == 200:\n", + " data = response.json()\n", + " listen_history = data[\"items\"]\n", + " if not listen_history:\n", + " break\n", + " all_listen_history.extend(listen_history)\n", + " page += 1\n", + " else:\n", + " print(f\"Failed to fetch listen history: {response.status_code}\")\n", + " break\n", + " return pd.DataFrame(all_listen_history)\n", + "\n", + "listen_history_df = get_listen_history()\n", + "print(listen_history_df.head())" + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": ".env_uv_3919", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.9.19" + } + }, + "nbformat": 4, + "nbformat_minor": 4 +} diff --git a/notebooks/get_data_from_api.ipynb b/notebooks/get_data_from_api.ipynb new file mode 100644 index 0000000..8a7137e --- /dev/null +++ b/notebooks/get_data_from_api.ipynb @@ -0,0 +1,136 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 1, + "id": "58c84385", + "metadata": {}, + "outputs": [], + "source": [ + "import pandas as pd \n", + "import requests " + ] + }, + { + "cell_type": "code", + "execution_count": 2, + "id": "aba7a920", + "metadata": {}, + "outputs": [], + "source": [ + "BASE_URL = \"http://127.0.0.1:8001\" \n", + "endpoint_users= \"/users\"\n", + "endpoint_tracks = \"/tracks\"\n", + "endpoint_history = \"/listen_history\"" + ] + }, + { + "cell_type": "code", + "execution_count": 3, + "id": "8bfca5f3", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + " items total page size pages\n", + "0 {'id': 8231, 'first_name': 'Gabrielle', 'last_... 1000 1 100 10\n", + "1 {'id': 30893, 'first_name': 'Morgan', 'last_na... 1000 1 100 10\n", + "2 {'id': 13477, 'first_name': 'Tiffany', 'last_n... 1000 1 100 10\n", + "3 {'id': 24153, 'first_name': 'Jonathan', 'last_... 1000 1 100 10\n", + "4 {'id': 16236, 'first_name': 'Robert', 'last_na... 1000 1 100 10\n" + ] + } + ], + "source": [ + "response = requests.get(f\"{BASE_URL}{endpoint_users}\")\n", + "if response.status_code == 200:\n", + " data = response.json()\n", + " df = pd.DataFrame(data)\n", + " print(df.head())\n", + "else:\n", + " print(f\"Error: {response.status_code} - {response.text}\")" + ] + }, + { + "cell_type": "code", + "execution_count": 4, + "id": "3b70ad47", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + " items total page size pages\n", + "0 {'id': 53589, 'name': 'because', 'artist': 'Co... 1000 1 100 10\n", + "1 {'id': 3938, 'name': 'up', 'artist': 'Marcia W... 1000 1 100 10\n", + "2 {'id': 86307, 'name': 'find', 'artist': 'Amber... 1000 1 100 10\n", + "3 {'id': 96800, 'name': 'sell', 'artist': 'Brian... 1000 1 100 10\n", + "4 {'id': 96699, 'name': 'believe', 'artist': 'Je... 1000 1 100 10\n" + ] + } + ], + "source": [ + "response = requests.get(f\"{BASE_URL}{endpoint_tracks}\")\n", + "if response.status_code == 200:\n", + " data = response.json()\n", + " df = pd.DataFrame(data)\n", + " print(df.head())\n", + "else:\n", + " print(f\"Error: {response.status_code} - {response.text}\")" + ] + }, + { + "cell_type": "code", + "execution_count": 5, + "id": "68a1d7bd", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + " items total page size pages\n", + "0 {'user_id': 8231, 'items': [48603, 65779, 3945... 1000 1 100 10\n", + "1 {'user_id': 30893, 'items': [39572, 75194, 534... 1000 1 100 10\n", + "2 {'user_id': 13477, 'items': [75096, 36601, 782... 1000 1 100 10\n", + "3 {'user_id': 24153, 'items': [95393, 19603, 504... 1000 1 100 10\n", + "4 {'user_id': 16236, 'items': [14367, 77470, 974... 1000 1 100 10\n" + ] + } + ], + "source": [ + "response = requests.get(f\"{BASE_URL}{endpoint_history}\")\n", + "if response.status_code == 200:\n", + " data = response.json()\n", + " df = pd.DataFrame(data)\n", + " print(df.head())\n", + "else:\n", + " print(f\"Error: {response.status_code} - {response.text}\")" + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": ".env_uv_3919", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.9.19" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} diff --git a/requirements.txt b/requirements.txt old mode 100644 new mode 100755 index 55593fd..ae2348e --- a/requirements.txt +++ b/requirements.txt @@ -2,4 +2,6 @@ faker fastapi uvicorn fastapi_pagination -pytest \ No newline at end of file +pytest +requests +pandas \ No newline at end of file diff --git a/requirements_new.txt b/requirements_new.txt new file mode 100644 index 0000000..667c74a --- /dev/null +++ b/requirements_new.txt @@ -0,0 +1,123 @@ +annotated-types==0.7.0 +anyio==4.10.0 +argon2-cffi==25.1.0 +argon2-cffi-bindings==25.1.0 +arrow==1.3.0 +asttokens==3.0.0 +async-lru==2.0.5 +attrs==25.3.0 +babel==2.17.0 +beautifulsoup4==4.13.4 +black==25.1.0 +bleach==6.2.0 +certifi==2025.8.3 +cffi==1.17.1 +charset-normalizer==3.4.2 +click==8.1.8 +comm==0.2.3 +debugpy==1.8.15 +decorator==5.2.1 +defusedxml==0.7.1 +exceptiongroup==1.3.0 +executing==2.2.0 +faker==37.5.3 +fastapi==0.116.1 +fastapi-pagination==0.13.3 +fastjsonschema==2.21.1 +fqdn==1.5.1 +h11==0.16.0 +httpcore==1.0.9 +httpx==0.28.1 +idna==3.10 +importlib-metadata==8.7.0 +iniconfig==2.1.0 +ipykernel==6.30.1 +ipython==8.18.1 +ipywidgets==8.1.7 +isoduration==20.11.0 +jedi==0.19.2 +jinja2==3.1.6 +json5==0.12.0 +jsonpointer==3.0.0 +jsonschema==4.25.0 +jsonschema-specifications==2025.4.1 +jupyter==1.1.1 +jupyter-client==8.6.3 +jupyter-console==6.6.3 +jupyter-core==5.8.1 +jupyter-events==0.12.0 +jupyter-lsp==2.2.6 +jupyter-server==2.16.0 +jupyter-server-terminals==0.5.3 +jupyterlab==4.4.5 +jupyterlab-pygments==0.3.0 +jupyterlab-server==2.27.3 +jupyterlab-widgets==3.0.15 +lark==1.2.2 +markupsafe==3.0.2 +matplotlib-inline==0.1.7 +mistune==3.1.3 +mypy-extensions==1.1.0 +nbclient==0.10.2 +nbconvert==7.16.6 +nbformat==5.10.4 +nest-asyncio==1.6.0 +notebook==7.4.4 +notebook-shim==0.2.4 +numpy==2.0.2 +overrides==7.7.0 +packaging==25.0 +pandas==2.3.1 +pandocfilters==1.5.1 +parso==0.8.4 +pathspec==0.12.1 +pexpect==4.9.0 +platformdirs==4.3.8 +pluggy==1.6.0 +prometheus-client==0.22.1 +prompt-toolkit==3.0.51 +psutil==7.0.0 +ptyprocess==0.7.0 +pure-eval==0.2.3 +pycparser==2.22 +pydantic==2.11.7 +pydantic-core==2.33.2 +pygments==2.19.2 +pytest==8.4.1 +python-dateutil==2.9.0.post0 +python-json-logger==3.3.0 +pytz==2025.2 +pyyaml==6.0.2 +pyzmq==27.0.1 +referencing==0.36.2 +requests==2.32.4 +rfc3339-validator==0.1.4 +rfc3986-validator==0.1.1 +rfc3987-syntax==1.1.0 +rpds-py==0.26.0 +ruff==0.12.7 +send2trash==1.8.3 +setuptools==80.9.0 +six==1.17.0 +sniffio==1.3.1 +soupsieve==2.7 +stack-data==0.6.3 +starlette==0.47.2 +terminado==0.18.1 +tinycss2==1.4.0 +tomli==2.2.1 +tornado==6.5.1 +traitlets==5.14.3 +types-python-dateutil==2.9.0.20250708 +typing-extensions==4.14.1 +typing-inspection==0.4.1 +tzdata==2025.2 +uri-template==1.3.0 +urllib3==2.5.0 +uvicorn==0.35.0 +wcwidth==0.2.13 +webcolors==24.11.1 +webencodings==0.5.1 +websocket-client==1.8.0 +widgetsnbextension==4.0.14 +zipp==3.23.0 diff --git a/src/moovitamix_fastapi/__init__.py b/src/moovitamix_fastapi/__init__.py old mode 100644 new mode 100755 diff --git a/src/moovitamix_fastapi/classes_out.py b/src/moovitamix_fastapi/classes_out.py old mode 100644 new mode 100755 diff --git a/src/moovitamix_fastapi/generate_fake_data.py b/src/moovitamix_fastapi/generate_fake_data.py old mode 100644 new mode 100755 diff --git a/src/moovitamix_fastapi/main.py b/src/moovitamix_fastapi/main.py old mode 100644 new mode 100755 index bd1d01f..638f778 --- a/src/moovitamix_fastapi/main.py +++ b/src/moovitamix_fastapi/main.py @@ -4,10 +4,15 @@ from fastapi.responses import RedirectResponse from fastapi_pagination import Page, add_pagination, paginate from generate_fake_data import FakeDataGenerator +from typing import TypeVar +from fastapi_pagination.customization import CustomizedPage, UseParamsFields -Page = Page.with_custom_options( - size=Query(100, ge=1, le=100), -) +T = TypeVar("T") + +CustomPage = CustomizedPage[ + Page[T], + UseParamsFields(size=Query(100, ge=1, le=100)), +] app = FastAPI( title="MooVitamix", @@ -37,18 +42,18 @@ async def overridden_swagger(): @app.get("/tracks", tags=["HTTP methods"]) -async def get_tracks() -> Page[TracksOut]: +async def get_tracks() -> CustomPage[TracksOut]: return paginate(tracks) @app.get("/users", tags=["HTTP methods"]) -async def get_users() -> Page[UsersOut]: +async def get_users() -> CustomPage[UsersOut]: return paginate(users) @app.get("/listen_history", tags=["HTTP methods"]) -async def get_listen_history() -> Page[ListenHistoryOut]: +async def get_listen_history() -> CustomPage[ListenHistoryOut]: return paginate(listen_history) -add_pagination(app) +add_pagination(app) \ No newline at end of file diff --git a/test/data_ingestion/__init__.py b/test/data_ingestion/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/test/data_ingestion/ingest_data.py b/test/data_ingestion/ingest_data.py new file mode 100644 index 0000000..826ac7b --- /dev/null +++ b/test/data_ingestion/ingest_data.py @@ -0,0 +1,88 @@ +import requests +import pandas as pd +from loguru import logger + +BASE_URL = "http://127.0.0.1:8001" + +def get_users(): + """Fetches all users from the API and returns them as a pandas DataFrame.""" + all_users = [] + page = 1 + while True: + response = requests.get(f"{BASE_URL}/users?page={page}&size=100") + if response.status_code == 200: + data = response.json() + users = data["items"] + if not users: + break + all_users.extend(users) + page += 1 + logger.info(f"Fetched page {page} of users, total users fetched: {len(all_users)}") + else: + logger.warning(f"Failed to fetch users: {response.status_code}") + logger.error(f"Error details: {response.text}") + break + return pd.DataFrame(all_users) + + +def get_tracks(): + """Fetches all tracks from the API and returns them as a pandas DataFrame.""" + all_tracks = [] + page = 1 + while True: + response = requests.get(f"{BASE_URL}/tracks?page={page}&size=100") + if response.status_code == 200: + data = response.json() + tracks = data["items"] + if not tracks: + break + all_tracks.extend(tracks) + page += 1 + logger.info(f"Fetched page {page} of tracks, total tracks fetched: {len(all_tracks)}") + else: + logger.warning(f"Failed to fetch tracks: {response.status_code}") + logger.error(f"Error details: {response.text}") + break + return pd.DataFrame(all_tracks) + +def get_listen_history(): + """Fetches all listen history from the API and returns it as a pandas DataFrame.""" + all_listen_history = [] + page = 1 + while True: + response = requests.get(f"{BASE_URL}/listen_history?page={page}&size=100") + if response.status_code == 200: + data = response.json() + listen_history = data["items"] + if not listen_history: + break + all_listen_history.extend(listen_history) + page += 1 + logger.info(f"Fetched page {page} of listen history, total records fetched: {len(all_listen_history)}") + else: + logger.warning(f"Failed to fetch listen history: {response.status_code}") + logger.error(f"Error details: {response.text}") + break + return pd.DataFrame(all_listen_history) + +def main(): + """init log file """ + logger.add("data_ingestion.log", rotation="1 MB", level="INFO") + logger.info("Starting data ingestion process") + logger.info("##########################################################################") + """Main function to ingest data from all endpoints.""" + users_df = get_users() + tracks_df = get_tracks() + listen_history_df = get_listen_history() + + print("Users DataFrame:") + print(users_df.head()) + print("\nTracks DataFrame:") + print(tracks_df.head()) + print("\nListen History DataFrame:") + print(listen_history_df.head()) + logger.info("Data ingestion process completed successfully") + logger.info("##########################################################################") + +if __name__ == "__main__": + main() diff --git a/test/requirements_new.txt b/test/requirements_new.txt new file mode 100644 index 0000000..d7eb48c --- /dev/null +++ b/test/requirements_new.txt @@ -0,0 +1,124 @@ +annotated-types==0.7.0 +anyio==4.10.0 +argon2-cffi==25.1.0 +argon2-cffi-bindings==25.1.0 +arrow==1.3.0 +asttokens==3.0.0 +async-lru==2.0.5 +attrs==25.3.0 +babel==2.17.0 +beautifulsoup4==4.13.4 +black==25.1.0 +bleach==6.2.0 +certifi==2025.8.3 +cffi==1.17.1 +charset-normalizer==3.4.2 +click==8.1.8 +comm==0.2.3 +debugpy==1.8.15 +decorator==5.2.1 +defusedxml==0.7.1 +exceptiongroup==1.3.0 +executing==2.2.0 +faker==37.5.3 +fastapi==0.116.1 +fastapi-pagination==0.13.3 +fastjsonschema==2.21.1 +fqdn==1.5.1 +h11==0.16.0 +httpcore==1.0.9 +httpx==0.28.1 +idna==3.10 +importlib-metadata==8.7.0 +iniconfig==2.1.0 +ipykernel==6.30.1 +ipython==8.18.1 +ipywidgets==8.1.7 +isoduration==20.11.0 +jedi==0.19.2 +jinja2==3.1.6 +json5==0.12.0 +jsonpointer==3.0.0 +jsonschema==4.25.0 +jsonschema-specifications==2025.4.1 +jupyter==1.1.1 +jupyter-client==8.6.3 +jupyter-console==6.6.3 +jupyter-core==5.8.1 +jupyter-events==0.12.0 +jupyter-lsp==2.2.6 +jupyter-server==2.16.0 +jupyter-server-terminals==0.5.3 +jupyterlab==4.4.5 +jupyterlab-pygments==0.3.0 +jupyterlab-server==2.27.3 +jupyterlab-widgets==3.0.15 +lark==1.2.2 +loguru==0.7.3 +markupsafe==3.0.2 +matplotlib-inline==0.1.7 +mistune==3.1.3 +mypy-extensions==1.1.0 +nbclient==0.10.2 +nbconvert==7.16.6 +nbformat==5.10.4 +nest-asyncio==1.6.0 +notebook==7.4.4 +notebook-shim==0.2.4 +numpy==2.0.2 +overrides==7.7.0 +packaging==25.0 +pandas==2.3.1 +pandocfilters==1.5.1 +parso==0.8.4 +pathspec==0.12.1 +pexpect==4.9.0 +platformdirs==4.3.8 +pluggy==1.6.0 +prometheus-client==0.22.1 +prompt-toolkit==3.0.51 +psutil==7.0.0 +ptyprocess==0.7.0 +pure-eval==0.2.3 +pycparser==2.22 +pydantic==2.11.7 +pydantic-core==2.33.2 +pygments==2.19.2 +pytest==8.4.1 +python-dateutil==2.9.0.post0 +python-json-logger==3.3.0 +pytz==2025.2 +pyyaml==6.0.2 +pyzmq==27.0.1 +referencing==0.36.2 +requests==2.32.4 +rfc3339-validator==0.1.4 +rfc3986-validator==0.1.1 +rfc3987-syntax==1.1.0 +rpds-py==0.26.0 +ruff==0.12.7 +send2trash==1.8.3 +setuptools==80.9.0 +six==1.17.0 +sniffio==1.3.1 +soupsieve==2.7 +stack-data==0.6.3 +starlette==0.47.2 +terminado==0.18.1 +tinycss2==1.4.0 +tomli==2.2.1 +tornado==6.5.1 +traitlets==5.14.3 +types-python-dateutil==2.9.0.20250708 +typing-extensions==4.14.1 +typing-inspection==0.4.1 +tzdata==2025.2 +uri-template==1.3.0 +urllib3==2.5.0 +uvicorn==0.35.0 +wcwidth==0.2.13 +webcolors==24.11.1 +webencodings==0.5.1 +websocket-client==1.8.0 +widgetsnbextension==4.0.14 +zipp==3.23.0 diff --git a/test/test_api.py b/test/test_api.py new file mode 100644 index 0000000..66fbcaa --- /dev/null +++ b/test/test_api.py @@ -0,0 +1,16 @@ +import pytest +import requests + +BASE_URL = "http://127.0.0.1:8001" + +@pytest.mark.parametrize("endpoint", ["/users", "/tracks", "/listen_history"]) +def test_endpoint_returns_200(endpoint): + """ + Tests that the API endpoints return a 200 OK status code. + """ + try: + response = requests.get(f"{BASE_URL}{endpoint}") + assert response.status_code == 200 + except requests.exceptions.ConnectionError as e: + pytest.fail(f"Connection to {BASE_URL} failed: {e}") + diff --git a/test/test_classes_out.py b/test/test_classes_out.py old mode 100644 new mode 100755 diff --git a/test/test_ingest_data.py b/test/test_ingest_data.py new file mode 100755 index 0000000..56862b4 --- /dev/null +++ b/test/test_ingest_data.py @@ -0,0 +1,130 @@ +import pandas as pd +from unittest.mock import patch, MagicMock + +# Import the script to be tested +from data_ingestion.ingest_data import get_listen_history, get_tracks, get_users + +# --- Expected Schemas --- +EXPECTED_USERS_COLUMNS = [ + "id", + "first_name", + "last_name", + "email", + "gender", + "favorite_genres", + "created_at", + "updated_at", +] +EXPECTED_TRACKS_COLUMNS = [ + "id", + "name", + "artist", + "songwriters", + "duration", + "genres", + "album", + "created_at", + "updated_at", +] +EXPECTED_LISTEN_HISTORY_COLUMNS = [ + "id", + "user_id", + "track_id", + "listened_at", + "created_at", + "updated_at", +] + + +# --- Mock Setup --- +def mock_success_response(sample_data, page=1): + """Creates a mock response object for a successful API call.""" + mock_res = MagicMock() + mock_res.status_code = 200 + + def json_func(): + # Simulate pagination: return data for page 1, empty list for subsequent pages + return {"items": sample_data if page == 1 else []} + + mock_res.json = json_func + return mock_res + + +# Sample data that matches the expected schemas +users_sample_data = [{col: "sample" for col in EXPECTED_USERS_COLUMNS}] +tracks_sample_data = [{col: "sample" for col in EXPECTED_TRACKS_COLUMNS}] +listen_history_sample_data = [ + {col: "sample" for col in EXPECTED_LISTEN_HISTORY_COLUMNS} +] + +# --- Original Mock Tests (verifying DataFrame content) --- + + +@patch("data_ingestion.ingest_data.requests.get") +def test_get_users(mock_get): + """Tests the get_users function returns a DataFrame with expected content.""" + mock_get.side_effect = [ + mock_success_response(users_sample_data, page=1), + mock_success_response([], page=2), + ] + users_df = get_users() + assert isinstance(users_df, pd.DataFrame) + assert not users_df.empty + # Check that the content matches the sample data + assert users_df.to_dict("records") == users_sample_data + mock_get.assert_any_call("http://127.0.0.1:8001/users?page=1&size=100") + + +@patch("data_ingestion.ingest_data.requests.get") +def test_get_tracks(mock_get): + """Tests the get_tracks function returns a DataFrame with expected content.""" + mock_get.side_effect = [ + mock_success_response(tracks_sample_data, page=1), + mock_success_response([], page=2), + ] + tracks_df = get_tracks() + assert isinstance(tracks_df, pd.DataFrame) + assert not tracks_df.empty + assert tracks_df.to_dict("records") == tracks_sample_data + mock_get.assert_any_call("http://127.0.0.1:8001/tracks?page=1&size=100") + + +@patch("data_ingestion.ingest_data.requests.get") +def test_get_listen_history(mock_get): + """Tests the get_listen_history function returns a DataFrame with expected content.""" + mock_get.side_effect = [ + mock_success_response(listen_history_sample_data, page=1), + mock_success_response([], page=2), + ] + listen_history_df = get_listen_history() + assert isinstance(listen_history_df, pd.DataFrame) + assert not listen_history_df.empty + assert listen_history_df.to_dict("records") == listen_history_sample_data + mock_get.assert_any_call("http://127.0.0.1:8001/listen_history?page=1&size=100") + + +# --- New Schema Validation Tests --- + + +@patch("data_ingestion.ingest_data.requests.get") +def test_get_users_schema(mock_get): + """Tests that the DataFrame from get_users has the expected columns.""" + mock_get.return_value = mock_success_response(users_sample_data) + users_df = get_users() + assert set(users_df.columns) == set(EXPECTED_USERS_COLUMNS) + + +@patch("data_ingestion.ingest_data.requests.get") +def test_get_tracks_schema(mock_get): + """Tests that the DataFrame from get_tracks has the expected columns.""" + mock_get.return_value = mock_success_response(tracks_sample_data) + tracks_df = get_tracks() + assert set(tracks_df.columns) == set(EXPECTED_TRACKS_COLUMNS) + + +@patch("data_ingestion.ingest_data.requests.get") +def test_get_listen_history_schema(mock_get): + """Tests that the DataFrame from get_listen_history has the expected columns.""" + mock_get.return_value = mock_success_response(listen_history_sample_data) + listen_history_df = get_listen_history() + assert set(listen_history_df.columns) == set(EXPECTED_LISTEN_HISTORY_COLUMNS) \ No newline at end of file