P
Parsyn
/Docs
Retour à l'accueil

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

start_federated_training

Demarre une session federee avec un fragment de dataset, les poids globaux et les informations de round.

load_aggregated_weights

Indique au worker de charger le dernier modele global apres l'aggregation.

continue_round

Signale au worker de passer au round suivant d'entrainement local.

Mecanisme de heartbeat

ParametreValeur
Intervalle de heartbeat du worker30 secondes (configurable : HEARTBEAT_INTERVAL)
Timeout de la plateforme90 secondes (WORKER_HEARTBEAT_TIMEOUT)
Intervalle des metriques systeme5 secondes (utilisation GPU, CPU, RAM)
Frequence de verification de la plateformeToutes 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 :

  1. Attendre RECONNECT_DELAY secondes (par defaut : 5).
  2. Tenter la reconnexion.
  3. En cas d'echec, doubler le delai (plafonne a 60 secondes), avec environ 10 % de jitter pour eviter l'effet de thundering herd.
  4. 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
  }
}