Protocole WebSocket
Format des messages, cycle de vie de la connexion, heartbeat et catalogue complet des messages echanges entre la plateforme et les workers.
Connexion
Connexion d'un worker
Les workers se connectent a la plateforme a l'adresse suivante :
wss://parsyn.progatis.com/ws/worker?worker_uuid={uuid}L'UUID du worker est transmis dans l'URL pour identification. L'authentification s'effectue dans le premier message :
{
"type": "authenticate",
"api_key": "wkr_abc123...",
"worker_uuid": "abc12345-6789-..."
}Flux d'enrolement (premiere connexion avec une cle d'enrolement)
Si le worker se connecte avec une cle d'enrolement (enroll_...) au lieu d'une cle individuelle, la plateforme repond avec :
{
"type": "enrollment_complete",
"data": {
"worker_id": 42,
"worker_name": "east-worker-3",
"api_key": "wkr_newkey..."
}
}Le worker sauvegarde cette cle individuelle localement et l'utilise pour toutes les connexions suivantes. Apres l'enrolement, la plateforme envoie le message standard de succes d'authentification.
Authentification reussie
{
"type": "authenticated",
"status": "ok",
"worker_id": "abc12345-6789-..."
}Echec d'authentification
{
"type": "auth_failed",
"reason": "Invalid API key"
}Enregistrement (apres l'authentification)
Apres une authentification reussie, le worker envoie les informations materielles et les resultats de benchmark :
{
"type": "worker_register",
"worker_id": "abc12345",
"data": {
"gpu_available": true,
"gpu_name": "NVIDIA A100-SXM4-80GB",
"gpu_memory_gb": 80.0,
"max_batch_size": 32,
"worker_type": "both",
"benchmark": {
"gpu_score": 95.2,
"tflops": 82.6,
"memory_bandwidth_gbps": 1008.0,
"max_batch_size_tested": 256,
"benchmark_version": "1.0.0",
"benchmark_duration_seconds": 45.2,
"is_estimated": false
}
}
}WebSocket du dashboard
Les clients frontend se connectent pour recevoir les mises a jour en temps reel :
wss://parsyn.progatis.com/ws/user/{user_id}?token={websocket_token}Le WebSocket token est un token a courte duree de vie (60 secondes) obtenu via l'API REST. Les utilisateurs s'abonnent aux mises a jour de jobs specifiques :
{ "type": "subscribe_job", "job_id": 123 }
{ "type": "unsubscribe_job", "job_id": 123 }Format des messages
Tous les messages sont au format JSON avec un champ type :
{
"type": "message_type",
"worker_id": "abc12345",
"job_id": 123,
"data": { ... },
"timestamp": "2026-03-10T14:30:00Z"
}Worker vers plateforme
heartbeat
Envoye toutes les 30 secondes.
{
"type": "heartbeat",
"worker_id": "abc12345",
"status": "idle"
}worker_metrics
Metriques systeme envoyees toutes les 5 secondes.
{
"type": "worker_metrics",
"data": {
"gpu_utilization": 85.5,
"gpu_memory_used": 18.2,
"gpu_memory_total": 24.0,
"cpu_utilization": 22.3,
"ram_used": 8.5,
"ram_total": 64.0,
"status": "training",
"status_message": "Epoch 2/3, step 1500/3000"
}
}training_progress
Metriques d'entrainement par step.
{
"type": "training_progress",
"worker_id": "abc12345",
"job_id": 123,
"data": {
"epoch": 2,
"step": 1500,
"total_steps": 3000,
"loss": 0.3421,
"learning_rate": 0.000018,
"grad_norm": 1.24,
"throughput": 42.5,
"perplexity": 1.41,
"token_throughput": 64000,
"loss_ema": 0.35,
"gpu_utilization": 98,
"gpu_memory_used": 71.3,
"cpu_utilization": 15.0,
"ram_used": 8.2,
"eta_seconds": 3600
}
}job_started
Le worker a initialise le job (modele charge, dataset pret).
job_completed
{
"type": "job_completed",
"worker_id": "abc12345",
"job_id": 123,
"data": {
"final_loss": 0.1823,
"total_steps": 3000,
"training_time_seconds": 7200,
"s3_prefix": "models/job-123/final/"
}
}job_failed
{
"type": "job_failed",
"worker_id": "abc12345",
"job_id": 123,
"data": {
"error": "CUDA out of memory. Tried to allocate 2.5 GB...",
"step": 450,
"last_checkpoint": "checkpoints/job-123/step-400/"
}
}checkpoint_saved
{
"type": "checkpoint_saved",
"job_id": 123,
"data": {
"step": 1000,
"epoch": 1,
"path": "checkpoints/job-123/step-1000/",
"metrics": { "loss": 0.4521, "eval_loss": 0.5012 }
}
}inference_tokens
Lots de tokens diffuses en continu pendant l'inference (par lots de 5 tokens ou toutes les 50ms).
{
"type": "inference_tokens",
"inference_id": "infer-456",
"data": { "tokens": ["Gradient", " descent", " is", " an", " optimization"] }
}inference_completed
{
"type": "inference_completed",
"inference_id": "infer-456",
"data": {
"tokens_generated": 256,
"metrics": {
"ttft_ms": 125.3,
"tps": 42.5,
"total_time_ms": 6234.5
}
}
}pong
Reponse au ping de la plateforme.
{
"type": "pong",
"worker_id": "abc12345",
"status": "training",
"busy": true,
"current_job_id": 123
}Plateforme vers worker
start_training
{
"type": "start_training",
"job_id": 123,
"data": {
"dataset": {
"storage_path": "datasets/42/data.jsonl",
"download_url": "https://s3.../presigned-url",
"format": "jsonl",
"id": 42
},
"base_model": {
"huggingface_id": "meta-llama/Llama-3.1-8B"
},
"config": {
"epochs": 3,
"batch_size": 4,
"learning_rate": 0.0002,
"max_length": 1024,
"mixed_precision": "bf16",
"use_peft": true,
"lora_r": 16,
"save_steps": 500,
"logging_steps": 10
},
"output_path": "models/job-123/"
}
}stop_training
Arret progressif. Le worker sauvegarde un checkpoint avant de s'arreter.
{ "type": "stop_training", "job_id": 123 }ping
Verification de disponibilite. Le worker doit repondre avec pong.
start_inference
{
"type": "start_inference",
"inference_id": "infer-456",
"data": {
"model_path": "models/job-123/final/",
"huggingface_id": "meta-llama/Llama-3.1-8B",
"download_urls": [
{ "key": "adapter_config.json", "url": "https://..." },
{ "key": "adapter_model.safetensors", "url": "https://..." }
],
"messages": [
{ "role": "user", "content": "Explain gradient descent." }
],
"system_prompt": "You are a helpful ML tutor.",
"temperature": 0.7,
"max_tokens": 500,
"stream": true
}
}start_evaluation
{
"type": "start_evaluation",
"evaluation_id": 78,
"data": {
"model_url": "models/job-123/final/",
"dataset_url": "datasets/43/eval.jsonl",
"metrics": ["perplexity"]
}
}Messages d'entrainement federe
Demarre une session federee avec un fragment de dataset, les poids globaux et les informations de round.
Indique au worker de charger le dernier modele global apres l'aggregation.
Signale au worker de passer au round suivant d'entrainement local.
Mecanisme de heartbeat
| Parametre | Valeur |
|---|---|
| Intervalle de heartbeat du worker | 30 secondes (configurable : HEARTBEAT_INTERVAL) |
| Timeout de la plateforme | 90 secondes (WORKER_HEARTBEAT_TIMEOUT) |
| Intervalle des metriques systeme | 5 secondes (utilisation GPU, CPU, RAM) |
| Frequence de verification de la plateforme | Toutes les 60 secondes, analyse de tous les workers pour detecter les heartbeats manques |
Si un worker manque 3 heartbeats consecutifs, la plateforme envoie un ping pour verification. Si aucun pong n'est recu, le worker est marque hors ligne et tout job en cours passe a l'etat failed.
Reconnexion
Lorsque la connexion WebSocket est interrompue, le worker se reconnecte automatiquement :
- Attendre
RECONNECT_DELAYsecondes (par defaut : 5). - Tenter la reconnexion.
- En cas d'echec, doubler le delai (plafonne a 60 secondes), avec environ 10 % de jitter pour eviter l'effet de thundering herd.
- En cas de succes, reinitialiser le delai, se reauthentifier et renvoyer les informations materielles.
Le worker retente indefiniment jusqu'a un arret manuel.
Plateforme vers dashboard
La plateforme transmet les mises a jour des workers aux clients dashboard connectes :
job_update
{
"type": "job_update",
"job_id": 123,
"data": {
"status": "running",
"epoch": 2, "step": 1500,
"loss": 0.3421, "learning_rate": 0.000018,
"gpu_utilization": 98
}
}worker_update
{
"type": "worker_update",
"worker_id": "abc12345",
"data": {
"status": "online",
"gpu_utilization": 45,
"gpu_memory_used": 12.5,
"gpu_memory_total": 24.0,
"cpu_utilization": 18.0,
"current_job_id": 123
}
}training_failed
{
"type": "training_failed",
"job_id": 123,
"data": {
"error": "Worker disconnected during training",
"last_checkpoint_step": 1400
}
}