P
Parsyn
/Docs
Back to Home

WebSocket Protocol

Message format, connection lifecycle, heartbeat, and the full catalog of messages between the platform and workers.

Connection

Worker connection

Workers connect to the platform at:

wss://parsyn.progatis.com/ws/worker?worker_uuid={uuid}

The worker UUID is passed in the URL for identification. Authentication happens in the first message:

{
  "type": "authenticate",
  "api_key": "wkr_abc123...",
  "worker_uuid": "abc12345-6789-..."
}

Enrollment flow (first connection with enrollment key)

If the worker connects with an enrollment key (enroll_...) instead of an individual key, the platform responds with:

{
  "type": "enrollment_complete",
  "data": {
    "worker_id": 42,
    "worker_name": "east-worker-3",
    "api_key": "wkr_newkey..."
  }
}

The worker saves this individual key locally and uses it for all future connections. After enrollment, the platform sends the standard authentication success message.

Authentication success

{
  "type": "authenticated",
  "status": "ok",
  "worker_id": "abc12345-6789-..."
}

Authentication failure

{
  "type": "auth_failed",
  "reason": "Invalid API key"
}

Registration (after auth)

After successful authentication, the worker sends hardware info and benchmark results:

{
  "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
    }
  }
}

Dashboard WebSocket

Frontend clients connect to receive real-time updates:

wss://parsyn.progatis.com/ws/user/{user_id}?token={websocket_token}

The WebSocket token is a short-lived (60 seconds) token obtained through the REST API. Users subscribe to specific job updates:

{ "type": "subscribe_job", "job_id": 123 }
{ "type": "unsubscribe_job", "job_id": 123 }

Message format

All messages are JSON with a type field:

{
  "type": "message_type",
  "worker_id": "abc12345",
  "job_id": 123,
  "data": { ... },
  "timestamp": "2026-03-10T14:30:00Z"
}

Worker to platform

heartbeat

Sent every 30 seconds.

{
  "type": "heartbeat",
  "worker_id": "abc12345",
  "status": "idle"
}

worker_metrics

System metrics sent every 5 seconds.

{
  "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

Per-step training metrics.

{
  "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

Worker has initialized the job (model loaded, dataset ready).

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

Streamed token batches during inference (batched: 5 tokens or 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

Response to platform ping.

{
  "type": "pong",
  "worker_id": "abc12345",
  "status": "training",
  "busy": true,
  "current_job_id": 123
}

Platform to 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

Graceful stop. Worker saves a checkpoint before terminating.

{ "type": "stop_training", "job_id": 123 }

ping

Health check. Worker must respond with 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"]
  }
}

Federated training messages

start_federated_training

Begins a federated session with dataset shard, global weights, and round info.

load_aggregated_weights

Tells the worker to load the latest global model after aggregation.

continue_round

Signals the worker to proceed to the next round of local training.

Heartbeat mechanism

ParameterValue
Worker heartbeat interval30 seconds (configurable: HEARTBEAT_INTERVAL)
Platform timeout90 seconds (WORKER_HEARTBEAT_TIMEOUT)
System metrics interval5 seconds (GPU, CPU, RAM utilization)
Platform check frequencyEvery 60 seconds, scans all workers for missed heartbeats

If a worker misses 3 consecutive heartbeats, the platform sends a ping to double-check. If no pong is received, the worker is marked offline and any running job transitions to failed.

Reconnection

When the WebSocket connection drops, the worker reconnects automatically:

  1. Wait RECONNECT_DELAY seconds (default: 5).
  2. Attempt reconnection.
  3. On failure, double the delay (capped at 60 seconds), with ~10% jitter to avoid thundering herd.
  4. On success, reset the delay, re-authenticate, re-send hardware info.

The worker retries indefinitely until manually stopped.

Platform to dashboard

The platform forwards worker updates to connected dashboard clients:

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
  }
}