Repository files navigation

EdgeGrid

EdgeGrid is a decentralized ML training network built on personal computers. Submit a training job via HTTP, and EdgeGrid finds a capable machine in the network, runs the job, streams logs back in real time, and stores the model checkpoint — with no manual setup on the worker machine beyond running the agent.

You Coordinator Worker (friend's PC)
│ │ │
│ POST /jobs │ │
│ {script, requirements, │ │
│ requires_gpu: true} │ │
│─────────────────────────>│ │
│ │ match GPU/RAM/disk reqs │
│ │ assign worker │
│ │─────────────────────────────>│
│ │ │ pip install (cached)
│ │ │ run train.py
│ GET /jobs/{id}/logs │ │
│─────────────────────────>│<── log lines via JetStream ──│
│<── SSE stream ───────────│ │
│ Epoch 1/10 loss=0.84 │ │
│ Epoch 2/10 loss=0.71 │ │
│ ... │ │
│ │<── checkpoint + result ──────│
│ GET /jobs/{id}/artifact │ │
│─────────────────────────>│ │
│<── model.tar.gz ─────────│ │

How it works

EdgeGrid is fully event-driven. Workers have no inbound ports from the coordinator — they pull jobs from NATS JetStream. All state (job lifecycle, worker registry, checkpoints) lives in NATS, not in the coordinator process. The coordinator can crash and restart with zero data loss.

Coordinator — HTTP API for job submission and status. Matches job hardware requirements (GPU, RAM, VRAM, disk) to registered workers. Dispatches jobs directly to the matching worker's personal NATS subject. If no worker is free, the job stays QUEUED and is auto-dispatched when capacity appears.

Worker — Registers its hardware capabilities at startup. Listens for jobs addressed to it. Runs training inside an isolated directory with a cached Python venv. Streams stdout/stderr to NATS as log lines. Pushes output/ as a checkpoint every 5 minutes during training and once on completion.

NATS JetStream — The single source of truth. Carries job messages, log lines, results, heartbeats, and cancel signals. Stores worker state and job state in KV buckets. Stores datasets and checkpoints in object store buckets. Replication is configurable for production clusters.


Features

  • Intelligent routing — jobs matched to workers by GPU, VRAM, RAM, and disk requirements
  • Job queuing — no free worker? job waits and auto-dispatches when one becomes available
  • CAS-safe dispatch — multiple coordinators can run simultaneously without double-assigning workers
  • Log streaming — real-time stdout/stderr via SSE; late-connecting clients get full replay from the start
  • Job cancellationDELETE /jobs/{id} kills the running Python process on the worker
  • Mid-training checkpointingoutput/ uploaded to object store every 5 minutes during training
  • Stale job recovery — if a worker dies, the job is automatically requeued within ~90 seconds
  • Venv caching — SHA256(requirements.txt) keyed venvs; repeated jobs with the same deps skip pip install
  • Single binary — run as coordinator, worker, or both

Getting Started

Prerequisites

  • Go 1.21+
  • NATS Server with JetStream enabled
  • Python 3 (only for training executor — auto-detected on the worker machine)

Build

git clone https://github.com/edgegrid/edgegrid.git
cd edgegrid
go build -o edgegrid ./cmd/edgegrid

Run locally (single node, dev)

# Terminal 1 — NATS
nats-server -js
# Terminal 2 — coordinator + worker (default: both enabled, mock executor)
./edgegrid
# Terminal 3 — submit a training job
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{ "training_script": "import os\nprint(\"training...\")\nopen(os.environ[\"OUTPUT_DIR\"]+\"/model.pt\",\"w\").write(\"weights\")", "dataset_ref": "my-dataset", "requires_gpu": false }'# → {"job_id":"a1b2c3d4","status":"queued"}# Stream logs
curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# Check status
curl http://localhost:8080/jobs/a1b2c3d4
# Cancel
curl -X DELETE http://localhost:8080/jobs/a1b2c3d4

Run as separate coordinator and worker

# Coordinator only
./edgegrid -server -nats nats://localhost:4222 -port 8080
# Worker only (on another machine)
./edgegrid -client -nats nats://coordinator:4222 -executor training -worker-id worker-gpu-01

Configuration

Flags take precedence over environment variables. If neither -server nor -client is passed, both are enabled.

FlagEnv varDefaultDescription
-servertrue*Enable coordinator (HTTP API)
-clienttrue*Enable worker
-natsNATS_URLnats://localhost:4222NATS connection URL
-portPORT8080Coordinator HTTP API port
-worker-idWORKER_IDauto-generatedCustom worker identifier
-executorEXECUTORmockExecutor backend (mock or training)
-replicasNATS_REPLICAS1NATS JetStream replication factor (1=dev, 3=prod)

HTTP API

POST /jobs — Submit a training job

{
"training_script": "print('hello')",
"requirements": "torch==2.0.0\nnumpy==1.24.0",
"dataset_type": "object_store",
"dataset_ref": "my-dataset-key",
"base_model_type": "hf",
"base_model_ref": "bert-base-uncased",
"training_config_json": "{\"epochs\": 10, \"lr\": 0.001}",
"requires_gpu": true,
"min_ram_gb": 16.0,
"min_vram_gb": 8.0,
"min_disk_gb": 20.0
}

Response 202 Accepted:

{"job_id": "a1b2c3d4", "status": "queued"}

Hardware requirement fields are all optional. Omit them and any free worker qualifies.

GET /jobs/{id} — Job status

{
"job_id": "a1b2c3d4",
"state": "COMPLETED",
"worker_id": "worker-gpu-01",
"checkpoint_key": "a1b2c3d4",
"updated_at": "2026-07-01T12:34:56Z"
}

GET /jobs/{id}/logs — Live log streaming (SSE)

curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# data: Epoch 1/10 loss=0.842# data: Epoch 2/10 loss=0.761# ...# event: done# data: COMPLETED

Connects via Server-Sent Events. Late-connecting clients receive all prior log lines from the beginning (JetStream DeliverAll). Stream closes with event: done when the job reaches a terminal state.

DELETE /jobs/{id} — Cancel a job

Cancels a QUEUED or RUNNING job. Returns 202 Accepted. The job state becomes CANCELLED. If running, the training process on the worker is killed within seconds.

POST /jobs/{id}/upload — Upload a dataset

Upload a dataset file for the job before or after submission (referenced by dataset_ref in the job request).

GET /jobs/{id}/artifact — Download checkpoint

Downloads the latest model checkpoint as a .tar.gz archive. Available after training completes or during training (mid-training checkpoints are uploaded every 5 minutes).

GET /health

Returns 200 ok. Used by load balancers and Docker Compose health checks.


Job Lifecycle

QUEUED ──► RUNNING ──► COMPLETED
│ │ └──────────► FAILED
│
└──────────────────► CANCELLED
  • QUEUED — job created, waiting for a capable worker
  • RUNNING — dispatched to a worker, training in progress
  • COMPLETED — training finished, checkpoint available
  • FAILED — training script exited with a non-zero status
  • CANCELLED — cancelled via DELETE /jobs/{id}

If a worker dies mid-job, the stale job recovery process requeues it to QUEUED within ~90 seconds.


Architecture

┌──────────────────────────────────────────────────────────┐
│ NATS JetStream Cluster │
│ │
│ JOBS stream workers_state KV jobs_state KV │
│ jobs.train.* (TTL: 1 min) (TTL: 24h) │
│ jobs.results worker caps, job state, │
│ jobs.logs.* free/busy state RequestProto │
│ jobs.cancel │
│ workers.register datasets store checkpoints store │
│ workers.heartbeat (TTL: 48h) (TTL: 7 days) │
└──────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌────────┴────────┐ ┌────────┴────────┐
│ Coordinator │ │ Worker │
│ │ │ │
│ HTTP API │ │ RegisterWorker │
│ Job routing │ │ StartHeartbeat │
│ TryDispatch │ │ StartJobListen │
│ StaleRecovery │ │ StartCancel │
│ │ │ TrainingExec │
└─────────────────┘ └─────────────────┘

Packages

PackageResponsibility
cmd/edgegridBinary entrypoint
internal/agentBoots coordinator and/or worker from a shared NATS connection
internal/coordinatorHTTP API, job routing, dispatch, stale recovery
internal/coordinator/workermanWorker KV registry, capability matching, CAS assignment
internal/workerJob listener, heartbeat, registration, cancel listener
internal/worker/executorTraining executor (venv cache, script runner) and mock
internal/brokerNATS JetStream, KV, and Object Store helpers
internal/jobstateJob state read/write helpers
internal/proto/workerProtobuf schemas and generated Go code

Documentation

Detailed design docs are in docs/:

DocWhat it covers
intelligent-routing.mdHardware capability matching, why routing is coordinator-owned
training-executor.mdVenv caching, script execution, environment injection
job-queuing.mdFIFO queue, RequestProto persistence, CAS dispatch
log-streaming.mdJetStream + SSE, DeliverAll for late clients
job-cancellation.mdPer-job context, cancel signal broadcast, state protection
reliability.mdStale job recovery, mid-training checkpointing
nats-raft-replicas.mdRaft consensus, replication, stateless coordinator design

About

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Repository files navigation

EdgeGrid

EdgeGrid is a decentralized ML training network built on personal computers. Submit a training job via HTTP, and EdgeGrid finds a capable machine in the network, runs the job, streams logs back in real time, and stores the model checkpoint — with no manual setup on the worker machine beyond running the agent.

You Coordinator Worker (friend's PC)
│ │ │
│ POST /jobs │ │
│ {script, requirements, │ │
│ requires_gpu: true} │ │
│─────────────────────────>│ │
│ │ match GPU/RAM/disk reqs │
│ │ assign worker │
│ │─────────────────────────────>│
│ │ │ pip install (cached)
│ │ │ run train.py
│ GET /jobs/{id}/logs │ │
│─────────────────────────>│<── log lines via JetStream ──│
│<── SSE stream ───────────│ │
│ Epoch 1/10 loss=0.84 │ │
│ Epoch 2/10 loss=0.71 │ │
│ ... │ │
│ │<── checkpoint + result ──────│
│ GET /jobs/{id}/artifact │ │
│─────────────────────────>│ │
│<── model.tar.gz ─────────│ │

How it works

EdgeGrid is fully event-driven. Workers have no inbound ports from the coordinator — they pull jobs from NATS JetStream. All state (job lifecycle, worker registry, checkpoints) lives in NATS, not in the coordinator process. The coordinator can crash and restart with zero data loss.

Coordinator — HTTP API for job submission and status. Matches job hardware requirements (GPU, RAM, VRAM, disk) to registered workers. Dispatches jobs directly to the matching worker's personal NATS subject. If no worker is free, the job stays QUEUED and is auto-dispatched when capacity appears.

Worker — Registers its hardware capabilities at startup. Listens for jobs addressed to it. Runs training inside an isolated directory with a cached Python venv. Streams stdout/stderr to NATS as log lines. Pushes output/ as a checkpoint every 5 minutes during training and once on completion.

NATS JetStream — The single source of truth. Carries job messages, log lines, results, heartbeats, and cancel signals. Stores worker state and job state in KV buckets. Stores datasets and checkpoints in object store buckets. Replication is configurable for production clusters.


Features

  • Intelligent routing — jobs matched to workers by GPU, VRAM, RAM, and disk requirements
  • Job queuing — no free worker? job waits and auto-dispatches when one becomes available
  • CAS-safe dispatch — multiple coordinators can run simultaneously without double-assigning workers
  • Log streaming — real-time stdout/stderr via SSE; late-connecting clients get full replay from the start
  • Job cancellationDELETE /jobs/{id} kills the running Python process on the worker
  • Mid-training checkpointingoutput/ uploaded to object store every 5 minutes during training
  • Stale job recovery — if a worker dies, the job is automatically requeued within ~90 seconds
  • Venv caching — SHA256(requirements.txt) keyed venvs; repeated jobs with the same deps skip pip install
  • Single binary — run as coordinator, worker, or both

Getting Started

Prerequisites

  • Go 1.21+
  • NATS Server with JetStream enabled
  • Python 3 (only for training executor — auto-detected on the worker machine)

Build

git clone https://github.com/edgegrid/edgegrid.git
cd edgegrid
go build -o edgegrid ./cmd/edgegrid

Run locally (single node, dev)

# Terminal 1 — NATS
nats-server -js
# Terminal 2 — coordinator + worker (default: both enabled, mock executor)
./edgegrid
# Terminal 3 — submit a training job
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{ "training_script": "import os\nprint(\"training...\")\nopen(os.environ[\"OUTPUT_DIR\"]+\"/model.pt\",\"w\").write(\"weights\")", "dataset_ref": "my-dataset", "requires_gpu": false }'# → {"job_id":"a1b2c3d4","status":"queued"}# Stream logs
curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# Check status
curl http://localhost:8080/jobs/a1b2c3d4
# Cancel
curl -X DELETE http://localhost:8080/jobs/a1b2c3d4

Run as separate coordinator and worker

# Coordinator only
./edgegrid -server -nats nats://localhost:4222 -port 8080
# Worker only (on another machine)
./edgegrid -client -nats nats://coordinator:4222 -executor training -worker-id worker-gpu-01

Configuration

Flags take precedence over environment variables. If neither -server nor -client is passed, both are enabled.

FlagEnv varDefaultDescription
-servertrue*Enable coordinator (HTTP API)
-clienttrue*Enable worker
-natsNATS_URLnats://localhost:4222NATS connection URL
-portPORT8080Coordinator HTTP API port
-worker-idWORKER_IDauto-generatedCustom worker identifier
-executorEXECUTORmockExecutor backend (mock or training)
-replicasNATS_REPLICAS1NATS JetStream replication factor (1=dev, 3=prod)

HTTP API

POST /jobs — Submit a training job

{
"training_script": "print('hello')",
"requirements": "torch==2.0.0\nnumpy==1.24.0",
"dataset_type": "object_store",
"dataset_ref": "my-dataset-key",
"base_model_type": "hf",
"base_model_ref": "bert-base-uncased",
"training_config_json": "{\"epochs\": 10, \"lr\": 0.001}",
"requires_gpu": true,
"min_ram_gb": 16.0,
"min_vram_gb": 8.0,
"min_disk_gb": 20.0
}

Response 202 Accepted:

{"job_id": "a1b2c3d4", "status": "queued"}

Hardware requirement fields are all optional. Omit them and any free worker qualifies.

GET /jobs/{id} — Job status

{
"job_id": "a1b2c3d4",
"state": "COMPLETED",
"worker_id": "worker-gpu-01",
"checkpoint_key": "a1b2c3d4",
"updated_at": "2026-07-01T12:34:56Z"
}

GET /jobs/{id}/logs — Live log streaming (SSE)

curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# data: Epoch 1/10 loss=0.842# data: Epoch 2/10 loss=0.761# ...# event: done# data: COMPLETED

Connects via Server-Sent Events. Late-connecting clients receive all prior log lines from the beginning (JetStream DeliverAll). Stream closes with event: done when the job reaches a terminal state.

DELETE /jobs/{id} — Cancel a job

Cancels a QUEUED or RUNNING job. Returns 202 Accepted. The job state becomes CANCELLED. If running, the training process on the worker is killed within seconds.

POST /jobs/{id}/upload — Upload a dataset

Upload a dataset file for the job before or after submission (referenced by dataset_ref in the job request).

GET /jobs/{id}/artifact — Download checkpoint

Downloads the latest model checkpoint as a .tar.gz archive. Available after training completes or during training (mid-training checkpoints are uploaded every 5 minutes).

GET /health

Returns 200 ok. Used by load balancers and Docker Compose health checks.


Job Lifecycle

QUEUED ──► RUNNING ──► COMPLETED
│ │ └──────────► FAILED
│
└──────────────────► CANCELLED
  • QUEUED — job created, waiting for a capable worker
  • RUNNING — dispatched to a worker, training in progress
  • COMPLETED — training finished, checkpoint available
  • FAILED — training script exited with a non-zero status
  • CANCELLED — cancelled via DELETE /jobs/{id}

If a worker dies mid-job, the stale job recovery process requeues it to QUEUED within ~90 seconds.


Architecture

┌──────────────────────────────────────────────────────────┐
│ NATS JetStream Cluster │
│ │
│ JOBS stream workers_state KV jobs_state KV │
│ jobs.train.* (TTL: 1 min) (TTL: 24h) │
│ jobs.results worker caps, job state, │
│ jobs.logs.* free/busy state RequestProto │
│ jobs.cancel │
│ workers.register datasets store checkpoints store │
│ workers.heartbeat (TTL: 48h) (TTL: 7 days) │
└──────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌────────┴────────┐ ┌────────┴────────┐
│ Coordinator │ │ Worker │
│ │ │ │
│ HTTP API │ │ RegisterWorker │
│ Job routing │ │ StartHeartbeat │
│ TryDispatch │ │ StartJobListen │
│ StaleRecovery │ │ StartCancel │
│ │ │ TrainingExec │
└─────────────────┘ └─────────────────┘

Packages

PackageResponsibility
cmd/edgegridBinary entrypoint
internal/agentBoots coordinator and/or worker from a shared NATS connection
internal/coordinatorHTTP API, job routing, dispatch, stale recovery
internal/coordinator/workermanWorker KV registry, capability matching, CAS assignment
internal/workerJob listener, heartbeat, registration, cancel listener
internal/worker/executorTraining executor (venv cache, script runner) and mock
internal/brokerNATS JetStream, KV, and Object Store helpers
internal/jobstateJob state read/write helpers
internal/proto/workerProtobuf schemas and generated Go code

Documentation

Detailed design docs are in docs/:

DocWhat it covers
intelligent-routing.mdHardware capability matching, why routing is coordinator-owned
training-executor.mdVenv caching, script execution, environment injection
job-queuing.mdFIFO queue, RequestProto persistence, CAS dispatch
log-streaming.mdJetStream + SSE, DeliverAll for late clients
job-cancellation.mdPer-job context, cancel signal broadcast, state protection
reliability.mdStale job recovery, mid-training checkpointing
nats-raft-replicas.mdRaft consensus, replication, stateless coordinator design

About

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

EdgeGrid

EdgeGrid is a decentralized ML training network built on personal computers. Submit a training job via HTTP, and EdgeGrid finds a capable machine in the network, runs the job, streams logs back in real time, and stores the model checkpoint — with no manual setup on the worker machine beyond running the agent.

You Coordinator Worker (friend's PC)
│ │ │
│ POST /jobs │ │
│ {script, requirements, │ │
│ requires_gpu: true} │ │
│─────────────────────────>│ │
│ │ match GPU/RAM/disk reqs │
│ │ assign worker │
│ │─────────────────────────────>│
│ │ │ pip install (cached)
│ │ │ run train.py
│ GET /jobs/{id}/logs │ │
│─────────────────────────>│<── log lines via JetStream ──│
│<── SSE stream ───────────│ │
│ Epoch 1/10 loss=0.84 │ │
│ Epoch 2/10 loss=0.71 │ │
│ ... │ │
│ │<── checkpoint + result ──────│
│ GET /jobs/{id}/artifact │ │
│─────────────────────────>│ │
│<── model.tar.gz ─────────│ │

How it works

EdgeGrid is fully event-driven. Workers have no inbound ports from the coordinator — they pull jobs from NATS JetStream. All state (job lifecycle, worker registry, checkpoints) lives in NATS, not in the coordinator process. The coordinator can crash and restart with zero data loss.

Coordinator — HTTP API for job submission and status. Matches job hardware requirements (GPU, RAM, VRAM, disk) to registered workers. Dispatches jobs directly to the matching worker's personal NATS subject. If no worker is free, the job stays QUEUED and is auto-dispatched when capacity appears.

Worker — Registers its hardware capabilities at startup. Listens for jobs addressed to it. Runs training inside an isolated directory with a cached Python venv. Streams stdout/stderr to NATS as log lines. Pushes output/ as a checkpoint every 5 minutes during training and once on completion.

NATS JetStream — The single source of truth. Carries job messages, log lines, results, heartbeats, and cancel signals. Stores worker state and job state in KV buckets. Stores datasets and checkpoints in object store buckets. Replication is configurable for production clusters.


Features

  • Intelligent routing — jobs matched to workers by GPU, VRAM, RAM, and disk requirements
  • Job queuing — no free worker? job waits and auto-dispatches when one becomes available
  • CAS-safe dispatch — multiple coordinators can run simultaneously without double-assigning workers
  • Log streaming — real-time stdout/stderr via SSE; late-connecting clients get full replay from the start
  • Job cancellationDELETE /jobs/{id} kills the running Python process on the worker
  • Mid-training checkpointingoutput/ uploaded to object store every 5 minutes during training
  • Stale job recovery — if a worker dies, the job is automatically requeued within ~90 seconds
  • Venv caching — SHA256(requirements.txt) keyed venvs; repeated jobs with the same deps skip pip install
  • Single binary — run as coordinator, worker, or both

Getting Started

Prerequisites

  • Go 1.21+
  • NATS Server with JetStream enabled
  • Python 3 (only for training executor — auto-detected on the worker machine)

Build

git clone https://github.com/edgegrid/edgegrid.git
cd edgegrid
go build -o edgegrid ./cmd/edgegrid

Run locally (single node, dev)

# Terminal 1 — NATS
nats-server -js
# Terminal 2 — coordinator + worker (default: both enabled, mock executor)
./edgegrid
# Terminal 3 — submit a training job
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{ "training_script": "import os\nprint(\"training...\")\nopen(os.environ[\"OUTPUT_DIR\"]+\"/model.pt\",\"w\").write(\"weights\")", "dataset_ref": "my-dataset", "requires_gpu": false }'# → {"job_id":"a1b2c3d4","status":"queued"}# Stream logs
curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# Check status
curl http://localhost:8080/jobs/a1b2c3d4
# Cancel
curl -X DELETE http://localhost:8080/jobs/a1b2c3d4

Run as separate coordinator and worker

# Coordinator only
./edgegrid -server -nats nats://localhost:4222 -port 8080
# Worker only (on another machine)
./edgegrid -client -nats nats://coordinator:4222 -executor training -worker-id worker-gpu-01

Configuration

Flags take precedence over environment variables. If neither -server nor -client is passed, both are enabled.

FlagEnv varDefaultDescription
-servertrue*Enable coordinator (HTTP API)
-clienttrue*Enable worker
-natsNATS_URLnats://localhost:4222NATS connection URL
-portPORT8080Coordinator HTTP API port
-worker-idWORKER_IDauto-generatedCustom worker identifier
-executorEXECUTORmockExecutor backend (mock or training)
-replicasNATS_REPLICAS1NATS JetStream replication factor (1=dev, 3=prod)

HTTP API

POST /jobs — Submit a training job

{
"training_script": "print('hello')",
"requirements": "torch==2.0.0\nnumpy==1.24.0",
"dataset_type": "object_store",
"dataset_ref": "my-dataset-key",
"base_model_type": "hf",
"base_model_ref": "bert-base-uncased",
"training_config_json": "{\"epochs\": 10, \"lr\": 0.001}",
"requires_gpu": true,
"min_ram_gb": 16.0,
"min_vram_gb": 8.0,
"min_disk_gb": 20.0
}

Response 202 Accepted:

{"job_id": "a1b2c3d4", "status": "queued"}

Hardware requirement fields are all optional. Omit them and any free worker qualifies.

GET /jobs/{id} — Job status

{
"job_id": "a1b2c3d4",
"state": "COMPLETED",
"worker_id": "worker-gpu-01",
"checkpoint_key": "a1b2c3d4",
"updated_at": "2026-07-01T12:34:56Z"
}

GET /jobs/{id}/logs — Live log streaming (SSE)

curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# data: Epoch 1/10 loss=0.842# data: Epoch 2/10 loss=0.761# ...# event: done# data: COMPLETED

Connects via Server-Sent Events. Late-connecting clients receive all prior log lines from the beginning (JetStream DeliverAll). Stream closes with event: done when the job reaches a terminal state.

DELETE /jobs/{id} — Cancel a job

Cancels a QUEUED or RUNNING job. Returns 202 Accepted. The job state becomes CANCELLED. If running, the training process on the worker is killed within seconds.

POST /jobs/{id}/upload — Upload a dataset

Upload a dataset file for the job before or after submission (referenced by dataset_ref in the job request).

GET /jobs/{id}/artifact — Download checkpoint

Downloads the latest model checkpoint as a .tar.gz archive. Available after training completes or during training (mid-training checkpoints are uploaded every 5 minutes).

GET /health

Returns 200 ok. Used by load balancers and Docker Compose health checks.


Job Lifecycle

QUEUED ──► RUNNING ──► COMPLETED
│ │ └──────────► FAILED
│
└──────────────────► CANCELLED
  • QUEUED — job created, waiting for a capable worker
  • RUNNING — dispatched to a worker, training in progress
  • COMPLETED — training finished, checkpoint available
  • FAILED — training script exited with a non-zero status
  • CANCELLED — cancelled via DELETE /jobs/{id}

If a worker dies mid-job, the stale job recovery process requeues it to QUEUED within ~90 seconds.


Architecture

┌──────────────────────────────────────────────────────────┐
│ NATS JetStream Cluster │
│ │
│ JOBS stream workers_state KV jobs_state KV │
│ jobs.train.* (TTL: 1 min) (TTL: 24h) │
│ jobs.results worker caps, job state, │
│ jobs.logs.* free/busy state RequestProto │
│ jobs.cancel │
│ workers.register datasets store checkpoints store │
│ workers.heartbeat (TTL: 48h) (TTL: 7 days) │
└──────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌────────┴────────┐ ┌────────┴────────┐
│ Coordinator │ │ Worker │
│ │ │ │
│ HTTP API │ │ RegisterWorker │
│ Job routing │ │ StartHeartbeat │
│ TryDispatch │ │ StartJobListen │
│ StaleRecovery │ │ StartCancel │
│ │ │ TrainingExec │
└─────────────────┘ └─────────────────┘

Packages

PackageResponsibility
cmd/edgegridBinary entrypoint
internal/agentBoots coordinator and/or worker from a shared NATS connection
internal/coordinatorHTTP API, job routing, dispatch, stale recovery
internal/coordinator/workermanWorker KV registry, capability matching, CAS assignment
internal/workerJob listener, heartbeat, registration, cancel listener
internal/worker/executorTraining executor (venv cache, script runner) and mock
internal/brokerNATS JetStream, KV, and Object Store helpers
internal/jobstateJob state read/write helpers
internal/proto/workerProtobuf schemas and generated Go code

Documentation

Detailed design docs are in docs/:

DocWhat it covers
intelligent-routing.mdHardware capability matching, why routing is coordinator-owned
training-executor.mdVenv caching, script execution, environment injection
job-queuing.mdFIFO queue, RequestProto persistence, CAS dispatch
log-streaming.mdJetStream + SSE, DeliverAll for late clients
job-cancellation.mdPer-job context, cancel signal broadcast, state protection
reliability.mdStale job recovery, mid-training checkpointing
nats-raft-replicas.mdRaft consensus, replication, stateless coordinator design

About

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

EdgeGrid

EdgeGrid is a decentralized ML training network built on personal computers. Submit a training job via HTTP, and EdgeGrid finds a capable machine in the network, runs the job, streams logs back in real time, and stores the model checkpoint — with no manual setup on the worker machine beyond running the agent.

You Coordinator Worker (friend's PC)
│ │ │
│ POST /jobs │ │
│ {script, requirements, │ │
│ requires_gpu: true} │ │
│─────────────────────────>│ │
│ │ match GPU/RAM/disk reqs │
│ │ assign worker │
│ │─────────────────────────────>│
│ │ │ pip install (cached)
│ │ │ run train.py
│ GET /jobs/{id}/logs │ │
│─────────────────────────>│<── log lines via JetStream ──│
│<── SSE stream ───────────│ │
│ Epoch 1/10 loss=0.84 │ │
│ Epoch 2/10 loss=0.71 │ │
│ ... │ │
│ │<── checkpoint + result ──────│
│ GET /jobs/{id}/artifact │ │
│─────────────────────────>│ │
│<── model.tar.gz ─────────│ │

How it works

EdgeGrid is fully event-driven. Workers have no inbound ports from the coordinator — they pull jobs from NATS JetStream. All state (job lifecycle, worker registry, checkpoints) lives in NATS, not in the coordinator process. The coordinator can crash and restart with zero data loss.

Coordinator — HTTP API for job submission and status. Matches job hardware requirements (GPU, RAM, VRAM, disk) to registered workers. Dispatches jobs directly to the matching worker's personal NATS subject. If no worker is free, the job stays QUEUED and is auto-dispatched when capacity appears.

Worker — Registers its hardware capabilities at startup. Listens for jobs addressed to it. Runs training inside an isolated directory with a cached Python venv. Streams stdout/stderr to NATS as log lines. Pushes output/ as a checkpoint every 5 minutes during training and once on completion.

NATS JetStream — The single source of truth. Carries job messages, log lines, results, heartbeats, and cancel signals. Stores worker state and job state in KV buckets. Stores datasets and checkpoints in object store buckets. Replication is configurable for production clusters.


Features

  • Intelligent routing — jobs matched to workers by GPU, VRAM, RAM, and disk requirements
  • Job queuing — no free worker? job waits and auto-dispatches when one becomes available
  • CAS-safe dispatch — multiple coordinators can run simultaneously without double-assigning workers
  • Log streaming — real-time stdout/stderr via SSE; late-connecting clients get full replay from the start
  • Job cancellationDELETE /jobs/{id} kills the running Python process on the worker
  • Mid-training checkpointingoutput/ uploaded to object store every 5 minutes during training
  • Stale job recovery — if a worker dies, the job is automatically requeued within ~90 seconds
  • Venv caching — SHA256(requirements.txt) keyed venvs; repeated jobs with the same deps skip pip install
  • Single binary — run as coordinator, worker, or both

Getting Started

Prerequisites

  • Go 1.21+
  • NATS Server with JetStream enabled
  • Python 3 (only for training executor — auto-detected on the worker machine)

Build

git clone https://github.com/edgegrid/edgegrid.git
cd edgegrid
go build -o edgegrid ./cmd/edgegrid

Run locally (single node, dev)

# Terminal 1 — NATS
nats-server -js
# Terminal 2 — coordinator + worker (default: both enabled, mock executor)
./edgegrid
# Terminal 3 — submit a training job
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{ "training_script": "import os\nprint(\"training...\")\nopen(os.environ[\"OUTPUT_DIR\"]+\"/model.pt\",\"w\").write(\"weights\")", "dataset_ref": "my-dataset", "requires_gpu": false }'# → {"job_id":"a1b2c3d4","status":"queued"}# Stream logs
curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# Check status
curl http://localhost:8080/jobs/a1b2c3d4
# Cancel
curl -X DELETE http://localhost:8080/jobs/a1b2c3d4

Run as separate coordinator and worker

# Coordinator only
./edgegrid -server -nats nats://localhost:4222 -port 8080
# Worker only (on another machine)
./edgegrid -client -nats nats://coordinator:4222 -executor training -worker-id worker-gpu-01

Configuration

Flags take precedence over environment variables. If neither -server nor -client is passed, both are enabled.

FlagEnv varDefaultDescription
-servertrue*Enable coordinator (HTTP API)
-clienttrue*Enable worker
-natsNATS_URLnats://localhost:4222NATS connection URL
-portPORT8080Coordinator HTTP API port
-worker-idWORKER_IDauto-generatedCustom worker identifier
-executorEXECUTORmockExecutor backend (mock or training)
-replicasNATS_REPLICAS1NATS JetStream replication factor (1=dev, 3=prod)

HTTP API

POST /jobs — Submit a training job

{
"training_script": "print('hello')",
"requirements": "torch==2.0.0\nnumpy==1.24.0",
"dataset_type": "object_store",
"dataset_ref": "my-dataset-key",
"base_model_type": "hf",
"base_model_ref": "bert-base-uncased",
"training_config_json": "{\"epochs\": 10, \"lr\": 0.001}",
"requires_gpu": true,
"min_ram_gb": 16.0,
"min_vram_gb": 8.0,
"min_disk_gb": 20.0
}

Response 202 Accepted:

{"job_id": "a1b2c3d4", "status": "queued"}

Hardware requirement fields are all optional. Omit them and any free worker qualifies.

GET /jobs/{id} — Job status

{
"job_id": "a1b2c3d4",
"state": "COMPLETED",
"worker_id": "worker-gpu-01",
"checkpoint_key": "a1b2c3d4",
"updated_at": "2026-07-01T12:34:56Z"
}

GET /jobs/{id}/logs — Live log streaming (SSE)

curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# data: Epoch 1/10 loss=0.842# data: Epoch 2/10 loss=0.761# ...# event: done# data: COMPLETED

Connects via Server-Sent Events. Late-connecting clients receive all prior log lines from the beginning (JetStream DeliverAll). Stream closes with event: done when the job reaches a terminal state.

DELETE /jobs/{id} — Cancel a job

Cancels a QUEUED or RUNNING job. Returns 202 Accepted. The job state becomes CANCELLED. If running, the training process on the worker is killed within seconds.

POST /jobs/{id}/upload — Upload a dataset

Upload a dataset file for the job before or after submission (referenced by dataset_ref in the job request).

GET /jobs/{id}/artifact — Download checkpoint

Downloads the latest model checkpoint as a .tar.gz archive. Available after training completes or during training (mid-training checkpoints are uploaded every 5 minutes).

GET /health

Returns 200 ok. Used by load balancers and Docker Compose health checks.


Job Lifecycle

QUEUED ──► RUNNING ──► COMPLETED
│ │ └──────────► FAILED
│
└──────────────────► CANCELLED
  • QUEUED — job created, waiting for a capable worker
  • RUNNING — dispatched to a worker, training in progress
  • COMPLETED — training finished, checkpoint available
  • FAILED — training script exited with a non-zero status
  • CANCELLED — cancelled via DELETE /jobs/{id}

If a worker dies mid-job, the stale job recovery process requeues it to QUEUED within ~90 seconds.


Architecture

┌──────────────────────────────────────────────────────────┐
│ NATS JetStream Cluster │
│ │
│ JOBS stream workers_state KV jobs_state KV │
│ jobs.train.* (TTL: 1 min) (TTL: 24h) │
│ jobs.results worker caps, job state, │
│ jobs.logs.* free/busy state RequestProto │
│ jobs.cancel │
│ workers.register datasets store checkpoints store │
│ workers.heartbeat (TTL: 48h) (TTL: 7 days) │
└──────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌────────┴────────┐ ┌────────┴────────┐
│ Coordinator │ │ Worker │
│ │ │ │
│ HTTP API │ │ RegisterWorker │
│ Job routing │ │ StartHeartbeat │
│ TryDispatch │ │ StartJobListen │
│ StaleRecovery │ │ StartCancel │
│ │ │ TrainingExec │
└─────────────────┘ └─────────────────┘

Packages

PackageResponsibility
cmd/edgegridBinary entrypoint
internal/agentBoots coordinator and/or worker from a shared NATS connection
internal/coordinatorHTTP API, job routing, dispatch, stale recovery
internal/coordinator/workermanWorker KV registry, capability matching, CAS assignment
internal/workerJob listener, heartbeat, registration, cancel listener
internal/worker/executorTraining executor (venv cache, script runner) and mock
internal/brokerNATS JetStream, KV, and Object Store helpers
internal/jobstateJob state read/write helpers
internal/proto/workerProtobuf schemas and generated Go code

Documentation

Detailed design docs are in docs/:

DocWhat it covers
intelligent-routing.mdHardware capability matching, why routing is coordinator-owned
training-executor.mdVenv caching, script execution, environment injection
job-queuing.mdFIFO queue, RequestProto persistence, CAS dispatch
log-streaming.mdJetStream + SSE, DeliverAll for late clients
job-cancellation.mdPer-job context, cancel signal broadcast, state protection
reliability.mdStale job recovery, mid-training checkpointing
nats-raft-replicas.mdRaft consensus, replication, stateless coordinator design

About

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Repository files navigation

EdgeGrid

EdgeGrid is a decentralized ML training network built on personal computers. Submit a training job via HTTP, and EdgeGrid finds a capable machine in the network, runs the job, streams logs back in real time, and stores the model checkpoint — with no manual setup on the worker machine beyond running the agent.

You Coordinator Worker (friend's PC)
│ │ │
│ POST /jobs │ │
│ {script, requirements, │ │
│ requires_gpu: true} │ │
│─────────────────────────>│ │
│ │ match GPU/RAM/disk reqs │
│ │ assign worker │
│ │─────────────────────────────>│
│ │ │ pip install (cached)
│ │ │ run train.py
│ GET /jobs/{id}/logs │ │
│─────────────────────────>│<── log lines via JetStream ──│
│<── SSE stream ───────────│ │
│ Epoch 1/10 loss=0.84 │ │
│ Epoch 2/10 loss=0.71 │ │
│ ... │ │
│ │<── checkpoint + result ──────│
│ GET /jobs/{id}/artifact │ │
│─────────────────────────>│ │
│<── model.tar.gz ─────────│ │

How it works

EdgeGrid is fully event-driven. Workers have no inbound ports from the coordinator — they pull jobs from NATS JetStream. All state (job lifecycle, worker registry, checkpoints) lives in NATS, not in the coordinator process. The coordinator can crash and restart with zero data loss.

Coordinator — HTTP API for job submission and status. Matches job hardware requirements (GPU, RAM, VRAM, disk) to registered workers. Dispatches jobs directly to the matching worker's personal NATS subject. If no worker is free, the job stays QUEUED and is auto-dispatched when capacity appears.

Worker — Registers its hardware capabilities at startup. Listens for jobs addressed to it. Runs training inside an isolated directory with a cached Python venv. Streams stdout/stderr to NATS as log lines. Pushes output/ as a checkpoint every 5 minutes during training and once on completion.

NATS JetStream — The single source of truth. Carries job messages, log lines, results, heartbeats, and cancel signals. Stores worker state and job state in KV buckets. Stores datasets and checkpoints in object store buckets. Replication is configurable for production clusters.


Features

  • Intelligent routing — jobs matched to workers by GPU, VRAM, RAM, and disk requirements
  • Job queuing — no free worker? job waits and auto-dispatches when one becomes available
  • CAS-safe dispatch — multiple coordinators can run simultaneously without double-assigning workers
  • Log streaming — real-time stdout/stderr via SSE; late-connecting clients get full replay from the start
  • Job cancellationDELETE /jobs/{id} kills the running Python process on the worker
  • Mid-training checkpointingoutput/ uploaded to object store every 5 minutes during training
  • Stale job recovery — if a worker dies, the job is automatically requeued within ~90 seconds
  • Venv caching — SHA256(requirements.txt) keyed venvs; repeated jobs with the same deps skip pip install
  • Single binary — run as coordinator, worker, or both

Getting Started

Prerequisites

  • Go 1.21+
  • NATS Server with JetStream enabled
  • Python 3 (only for training executor — auto-detected on the worker machine)

Build

git clone https://github.com/edgegrid/edgegrid.git
cd edgegrid
go build -o edgegrid ./cmd/edgegrid

Run locally (single node, dev)

# Terminal 1 — NATS
nats-server -js
# Terminal 2 — coordinator + worker (default: both enabled, mock executor)
./edgegrid
# Terminal 3 — submit a training job
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{ "training_script": "import os\nprint(\"training...\")\nopen(os.environ[\"OUTPUT_DIR\"]+\"/model.pt\",\"w\").write(\"weights\")", "dataset_ref": "my-dataset", "requires_gpu": false }'# → {"job_id":"a1b2c3d4","status":"queued"}# Stream logs
curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# Check status
curl http://localhost:8080/jobs/a1b2c3d4
# Cancel
curl -X DELETE http://localhost:8080/jobs/a1b2c3d4

Run as separate coordinator and worker

# Coordinator only
./edgegrid -server -nats nats://localhost:4222 -port 8080
# Worker only (on another machine)
./edgegrid -client -nats nats://coordinator:4222 -executor training -worker-id worker-gpu-01

Configuration

Flags take precedence over environment variables. If neither -server nor -client is passed, both are enabled.

FlagEnv varDefaultDescription
-servertrue*Enable coordinator (HTTP API)
-clienttrue*Enable worker
-natsNATS_URLnats://localhost:4222NATS connection URL
-portPORT8080Coordinator HTTP API port
-worker-idWORKER_IDauto-generatedCustom worker identifier
-executorEXECUTORmockExecutor backend (mock or training)
-replicasNATS_REPLICAS1NATS JetStream replication factor (1=dev, 3=prod)

HTTP API

POST /jobs — Submit a training job

{
"training_script": "print('hello')",
"requirements": "torch==2.0.0\nnumpy==1.24.0",
"dataset_type": "object_store",
"dataset_ref": "my-dataset-key",
"base_model_type": "hf",
"base_model_ref": "bert-base-uncased",
"training_config_json": "{\"epochs\": 10, \"lr\": 0.001}",
"requires_gpu": true,
"min_ram_gb": 16.0,
"min_vram_gb": 8.0,
"min_disk_gb": 20.0
}

Response 202 Accepted:

{"job_id": "a1b2c3d4", "status": "queued"}

Hardware requirement fields are all optional. Omit them and any free worker qualifies.

GET /jobs/{id} — Job status

{
"job_id": "a1b2c3d4",
"state": "COMPLETED",
"worker_id": "worker-gpu-01",
"checkpoint_key": "a1b2c3d4",
"updated_at": "2026-07-01T12:34:56Z"
}

GET /jobs/{id}/logs — Live log streaming (SSE)

curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# data: Epoch 1/10 loss=0.842# data: Epoch 2/10 loss=0.761# ...# event: done# data: COMPLETED

Connects via Server-Sent Events. Late-connecting clients receive all prior log lines from the beginning (JetStream DeliverAll). Stream closes with event: done when the job reaches a terminal state.

DELETE /jobs/{id} — Cancel a job

Cancels a QUEUED or RUNNING job. Returns 202 Accepted. The job state becomes CANCELLED. If running, the training process on the worker is killed within seconds.

POST /jobs/{id}/upload — Upload a dataset

Upload a dataset file for the job before or after submission (referenced by dataset_ref in the job request).

GET /jobs/{id}/artifact — Download checkpoint

Downloads the latest model checkpoint as a .tar.gz archive. Available after training completes or during training (mid-training checkpoints are uploaded every 5 minutes).

GET /health

Returns 200 ok. Used by load balancers and Docker Compose health checks.


Job Lifecycle

QUEUED ──► RUNNING ──► COMPLETED
│ │ └──────────► FAILED
│
└──────────────────► CANCELLED
  • QUEUED — job created, waiting for a capable worker
  • RUNNING — dispatched to a worker, training in progress
  • COMPLETED — training finished, checkpoint available
  • FAILED — training script exited with a non-zero status
  • CANCELLED — cancelled via DELETE /jobs/{id}

If a worker dies mid-job, the stale job recovery process requeues it to QUEUED within ~90 seconds.


Architecture

┌──────────────────────────────────────────────────────────┐
│ NATS JetStream Cluster │
│ │
│ JOBS stream workers_state KV jobs_state KV │
│ jobs.train.* (TTL: 1 min) (TTL: 24h) │
│ jobs.results worker caps, job state, │
│ jobs.logs.* free/busy state RequestProto │
│ jobs.cancel │
│ workers.register datasets store checkpoints store │
│ workers.heartbeat (TTL: 48h) (TTL: 7 days) │
└──────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌────────┴────────┐ ┌────────┴────────┐
│ Coordinator │ │ Worker │
│ │ │ │
│ HTTP API │ │ RegisterWorker │
│ Job routing │ │ StartHeartbeat │
│ TryDispatch │ │ StartJobListen │
│ StaleRecovery │ │ StartCancel │
│ │ │ TrainingExec │
└─────────────────┘ └─────────────────┘

Packages

PackageResponsibility
cmd/edgegridBinary entrypoint
internal/agentBoots coordinator and/or worker from a shared NATS connection
internal/coordinatorHTTP API, job routing, dispatch, stale recovery
internal/coordinator/workermanWorker KV registry, capability matching, CAS assignment
internal/workerJob listener, heartbeat, registration, cancel listener
internal/worker/executorTraining executor (venv cache, script runner) and mock
internal/brokerNATS JetStream, KV, and Object Store helpers
internal/jobstateJob state read/write helpers
internal/proto/workerProtobuf schemas and generated Go code

Documentation

Detailed design docs are in docs/:

DocWhat it covers
intelligent-routing.mdHardware capability matching, why routing is coordinator-owned
training-executor.mdVenv caching, script execution, environment injection
job-queuing.mdFIFO queue, RequestProto persistence, CAS dispatch
log-streaming.mdJetStream + SSE, DeliverAll for late clients
job-cancellation.mdPer-job context, cancel signal broadcast, state protection
reliability.mdStale job recovery, mid-training checkpointing
nats-raft-replicas.mdRaft consensus, replication, stateless coordinator design

About

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

EdgeGrid

EdgeGrid is a decentralized ML training network built on personal computers. Submit a training job via HTTP, and EdgeGrid finds a capable machine in the network, runs the job, streams logs back in real time, and stores the model checkpoint — with no manual setup on the worker machine beyond running the agent.

You Coordinator Worker (friend's PC)
│ │ │
│ POST /jobs │ │
│ {script, requirements, │ │
│ requires_gpu: true} │ │
│─────────────────────────>│ │
│ │ match GPU/RAM/disk reqs │
│ │ assign worker │
│ │─────────────────────────────>│
│ │ │ pip install (cached)
│ │ │ run train.py
│ GET /jobs/{id}/logs │ │
│─────────────────────────>│<── log lines via JetStream ──│
│<── SSE stream ───────────│ │
│ Epoch 1/10 loss=0.84 │ │
│ Epoch 2/10 loss=0.71 │ │
│ ... │ │
│ │<── checkpoint + result ──────│
│ GET /jobs/{id}/artifact │ │
│─────────────────────────>│ │
│<── model.tar.gz ─────────│ │

How it works

EdgeGrid is fully event-driven. Workers have no inbound ports from the coordinator — they pull jobs from NATS JetStream. All state (job lifecycle, worker registry, checkpoints) lives in NATS, not in the coordinator process. The coordinator can crash and restart with zero data loss.

Coordinator — HTTP API for job submission and status. Matches job hardware requirements (GPU, RAM, VRAM, disk) to registered workers. Dispatches jobs directly to the matching worker's personal NATS subject. If no worker is free, the job stays QUEUED and is auto-dispatched when capacity appears.

Worker — Registers its hardware capabilities at startup. Listens for jobs addressed to it. Runs training inside an isolated directory with a cached Python venv. Streams stdout/stderr to NATS as log lines. Pushes output/ as a checkpoint every 5 minutes during training and once on completion.

NATS JetStream — The single source of truth. Carries job messages, log lines, results, heartbeats, and cancel signals. Stores worker state and job state in KV buckets. Stores datasets and checkpoints in object store buckets. Replication is configurable for production clusters.


Features

  • Intelligent routing — jobs matched to workers by GPU, VRAM, RAM, and disk requirements
  • Job queuing — no free worker? job waits and auto-dispatches when one becomes available
  • CAS-safe dispatch — multiple coordinators can run simultaneously without double-assigning workers
  • Log streaming — real-time stdout/stderr via SSE; late-connecting clients get full replay from the start
  • Job cancellationDELETE /jobs/{id} kills the running Python process on the worker
  • Mid-training checkpointingoutput/ uploaded to object store every 5 minutes during training
  • Stale job recovery — if a worker dies, the job is automatically requeued within ~90 seconds
  • Venv caching — SHA256(requirements.txt) keyed venvs; repeated jobs with the same deps skip pip install
  • Single binary — run as coordinator, worker, or both

Getting Started

Prerequisites

  • Go 1.21+
  • NATS Server with JetStream enabled
  • Python 3 (only for training executor — auto-detected on the worker machine)

Build

git clone https://github.com/edgegrid/edgegrid.git
cd edgegrid
go build -o edgegrid ./cmd/edgegrid

Run locally (single node, dev)

# Terminal 1 — NATS
nats-server -js
# Terminal 2 — coordinator + worker (default: both enabled, mock executor)
./edgegrid
# Terminal 3 — submit a training job
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{ "training_script": "import os\nprint(\"training...\")\nopen(os.environ[\"OUTPUT_DIR\"]+\"/model.pt\",\"w\").write(\"weights\")", "dataset_ref": "my-dataset", "requires_gpu": false }'# → {"job_id":"a1b2c3d4","status":"queued"}# Stream logs
curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# Check status
curl http://localhost:8080/jobs/a1b2c3d4
# Cancel
curl -X DELETE http://localhost:8080/jobs/a1b2c3d4

Run as separate coordinator and worker

# Coordinator only
./edgegrid -server -nats nats://localhost:4222 -port 8080
# Worker only (on another machine)
./edgegrid -client -nats nats://coordinator:4222 -executor training -worker-id worker-gpu-01

Configuration

Flags take precedence over environment variables. If neither -server nor -client is passed, both are enabled.

FlagEnv varDefaultDescription
-servertrue*Enable coordinator (HTTP API)
-clienttrue*Enable worker
-natsNATS_URLnats://localhost:4222NATS connection URL
-portPORT8080Coordinator HTTP API port
-worker-idWORKER_IDauto-generatedCustom worker identifier
-executorEXECUTORmockExecutor backend (mock or training)
-replicasNATS_REPLICAS1NATS JetStream replication factor (1=dev, 3=prod)

HTTP API

POST /jobs — Submit a training job

{
"training_script": "print('hello')",
"requirements": "torch==2.0.0\nnumpy==1.24.0",
"dataset_type": "object_store",
"dataset_ref": "my-dataset-key",
"base_model_type": "hf",
"base_model_ref": "bert-base-uncased",
"training_config_json": "{\"epochs\": 10, \"lr\": 0.001}",
"requires_gpu": true,
"min_ram_gb": 16.0,
"min_vram_gb": 8.0,
"min_disk_gb": 20.0
}

Response 202 Accepted:

{"job_id": "a1b2c3d4", "status": "queued"}

Hardware requirement fields are all optional. Omit them and any free worker qualifies.

GET /jobs/{id} — Job status

{
"job_id": "a1b2c3d4",
"state": "COMPLETED",
"worker_id": "worker-gpu-01",
"checkpoint_key": "a1b2c3d4",
"updated_at": "2026-07-01T12:34:56Z"
}

GET /jobs/{id}/logs — Live log streaming (SSE)

curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# data: Epoch 1/10 loss=0.842# data: Epoch 2/10 loss=0.761# ...# event: done# data: COMPLETED

Connects via Server-Sent Events. Late-connecting clients receive all prior log lines from the beginning (JetStream DeliverAll). Stream closes with event: done when the job reaches a terminal state.

DELETE /jobs/{id} — Cancel a job

Cancels a QUEUED or RUNNING job. Returns 202 Accepted. The job state becomes CANCELLED. If running, the training process on the worker is killed within seconds.

POST /jobs/{id}/upload — Upload a dataset

Upload a dataset file for the job before or after submission (referenced by dataset_ref in the job request).

GET /jobs/{id}/artifact — Download checkpoint

Downloads the latest model checkpoint as a .tar.gz archive. Available after training completes or during training (mid-training checkpoints are uploaded every 5 minutes).

GET /health

Returns 200 ok. Used by load balancers and Docker Compose health checks.


Job Lifecycle

QUEUED ──► RUNNING ──► COMPLETED
│ │ └──────────► FAILED
│
└──────────────────► CANCELLED
  • QUEUED — job created, waiting for a capable worker
  • RUNNING — dispatched to a worker, training in progress
  • COMPLETED — training finished, checkpoint available
  • FAILED — training script exited with a non-zero status
  • CANCELLED — cancelled via DELETE /jobs/{id}

If a worker dies mid-job, the stale job recovery process requeues it to QUEUED within ~90 seconds.


Architecture

┌──────────────────────────────────────────────────────────┐
│ NATS JetStream Cluster │
│ │
│ JOBS stream workers_state KV jobs_state KV │
│ jobs.train.* (TTL: 1 min) (TTL: 24h) │
│ jobs.results worker caps, job state, │
│ jobs.logs.* free/busy state RequestProto │
│ jobs.cancel │
│ workers.register datasets store checkpoints store │
│ workers.heartbeat (TTL: 48h) (TTL: 7 days) │
└──────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌────────┴────────┐ ┌────────┴────────┐
│ Coordinator │ │ Worker │
│ │ │ │
│ HTTP API │ │ RegisterWorker │
│ Job routing │ │ StartHeartbeat │
│ TryDispatch │ │ StartJobListen │
│ StaleRecovery │ │ StartCancel │
│ │ │ TrainingExec │
└─────────────────┘ └─────────────────┘

Packages

PackageResponsibility
cmd/edgegridBinary entrypoint
internal/agentBoots coordinator and/or worker from a shared NATS connection
internal/coordinatorHTTP API, job routing, dispatch, stale recovery
internal/coordinator/workermanWorker KV registry, capability matching, CAS assignment
internal/workerJob listener, heartbeat, registration, cancel listener
internal/worker/executorTraining executor (venv cache, script runner) and mock
internal/brokerNATS JetStream, KV, and Object Store helpers
internal/jobstateJob state read/write helpers
internal/proto/workerProtobuf schemas and generated Go code

Documentation

Detailed design docs are in docs/:

DocWhat it covers
intelligent-routing.mdHardware capability matching, why routing is coordinator-owned
training-executor.mdVenv caching, script execution, environment injection
job-queuing.mdFIFO queue, RequestProto persistence, CAS dispatch
log-streaming.mdJetStream + SSE, DeliverAll for late clients
job-cancellation.mdPer-job context, cancel signal broadcast, state protection
reliability.mdStale job recovery, mid-training checkpointing
nats-raft-replicas.mdRaft consensus, replication, stateless coordinator design

About

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

EdgeGrid

EdgeGrid is a decentralized ML training network built on personal computers. Submit a training job via HTTP, and EdgeGrid finds a capable machine in the network, runs the job, streams logs back in real time, and stores the model checkpoint — with no manual setup on the worker machine beyond running the agent.

You Coordinator Worker (friend's PC)
│ │ │
│ POST /jobs │ │
│ {script, requirements, │ │
│ requires_gpu: true} │ │
│─────────────────────────>│ │
│ │ match GPU/RAM/disk reqs │
│ │ assign worker │
│ │─────────────────────────────>│
│ │ │ pip install (cached)
│ │ │ run train.py
│ GET /jobs/{id}/logs │ │
│─────────────────────────>│<── log lines via JetStream ──│
│<── SSE stream ───────────│ │
│ Epoch 1/10 loss=0.84 │ │
│ Epoch 2/10 loss=0.71 │ │
│ ... │ │
│ │<── checkpoint + result ──────│
│ GET /jobs/{id}/artifact │ │
│─────────────────────────>│ │
│<── model.tar.gz ─────────│ │

How it works

EdgeGrid is fully event-driven. Workers have no inbound ports from the coordinator — they pull jobs from NATS JetStream. All state (job lifecycle, worker registry, checkpoints) lives in NATS, not in the coordinator process. The coordinator can crash and restart with zero data loss.

Coordinator — HTTP API for job submission and status. Matches job hardware requirements (GPU, RAM, VRAM, disk) to registered workers. Dispatches jobs directly to the matching worker's personal NATS subject. If no worker is free, the job stays QUEUED and is auto-dispatched when capacity appears.

Worker — Registers its hardware capabilities at startup. Listens for jobs addressed to it. Runs training inside an isolated directory with a cached Python venv. Streams stdout/stderr to NATS as log lines. Pushes output/ as a checkpoint every 5 minutes during training and once on completion.

NATS JetStream — The single source of truth. Carries job messages, log lines, results, heartbeats, and cancel signals. Stores worker state and job state in KV buckets. Stores datasets and checkpoints in object store buckets. Replication is configurable for production clusters.


Features

  • Intelligent routing — jobs matched to workers by GPU, VRAM, RAM, and disk requirements
  • Job queuing — no free worker? job waits and auto-dispatches when one becomes available
  • CAS-safe dispatch — multiple coordinators can run simultaneously without double-assigning workers
  • Log streaming — real-time stdout/stderr via SSE; late-connecting clients get full replay from the start
  • Job cancellationDELETE /jobs/{id} kills the running Python process on the worker
  • Mid-training checkpointingoutput/ uploaded to object store every 5 minutes during training
  • Stale job recovery — if a worker dies, the job is automatically requeued within ~90 seconds
  • Venv caching — SHA256(requirements.txt) keyed venvs; repeated jobs with the same deps skip pip install
  • Single binary — run as coordinator, worker, or both

Getting Started

Prerequisites

  • Go 1.21+
  • NATS Server with JetStream enabled
  • Python 3 (only for training executor — auto-detected on the worker machine)

Build

git clone https://github.com/edgegrid/edgegrid.git
cd edgegrid
go build -o edgegrid ./cmd/edgegrid

Run locally (single node, dev)

# Terminal 1 — NATS
nats-server -js
# Terminal 2 — coordinator + worker (default: both enabled, mock executor)
./edgegrid
# Terminal 3 — submit a training job
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{ "training_script": "import os\nprint(\"training...\")\nopen(os.environ[\"OUTPUT_DIR\"]+\"/model.pt\",\"w\").write(\"weights\")", "dataset_ref": "my-dataset", "requires_gpu": false }'# → {"job_id":"a1b2c3d4","status":"queued"}# Stream logs
curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# Check status
curl http://localhost:8080/jobs/a1b2c3d4
# Cancel
curl -X DELETE http://localhost:8080/jobs/a1b2c3d4

Run as separate coordinator and worker

# Coordinator only
./edgegrid -server -nats nats://localhost:4222 -port 8080
# Worker only (on another machine)
./edgegrid -client -nats nats://coordinator:4222 -executor training -worker-id worker-gpu-01

Configuration

Flags take precedence over environment variables. If neither -server nor -client is passed, both are enabled.

FlagEnv varDefaultDescription
-servertrue*Enable coordinator (HTTP API)
-clienttrue*Enable worker
-natsNATS_URLnats://localhost:4222NATS connection URL
-portPORT8080Coordinator HTTP API port
-worker-idWORKER_IDauto-generatedCustom worker identifier
-executorEXECUTORmockExecutor backend (mock or training)
-replicasNATS_REPLICAS1NATS JetStream replication factor (1=dev, 3=prod)

HTTP API

POST /jobs — Submit a training job

{
"training_script": "print('hello')",
"requirements": "torch==2.0.0\nnumpy==1.24.0",
"dataset_type": "object_store",
"dataset_ref": "my-dataset-key",
"base_model_type": "hf",
"base_model_ref": "bert-base-uncased",
"training_config_json": "{\"epochs\": 10, \"lr\": 0.001}",
"requires_gpu": true,
"min_ram_gb": 16.0,
"min_vram_gb": 8.0,
"min_disk_gb": 20.0
}

Response 202 Accepted:

{"job_id": "a1b2c3d4", "status": "queued"}

Hardware requirement fields are all optional. Omit them and any free worker qualifies.

GET /jobs/{id} — Job status

{
"job_id": "a1b2c3d4",
"state": "COMPLETED",
"worker_id": "worker-gpu-01",
"checkpoint_key": "a1b2c3d4",
"updated_at": "2026-07-01T12:34:56Z"
}

GET /jobs/{id}/logs — Live log streaming (SSE)

curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# data: Epoch 1/10 loss=0.842# data: Epoch 2/10 loss=0.761# ...# event: done# data: COMPLETED

Connects via Server-Sent Events. Late-connecting clients receive all prior log lines from the beginning (JetStream DeliverAll). Stream closes with event: done when the job reaches a terminal state.

DELETE /jobs/{id} — Cancel a job

Cancels a QUEUED or RUNNING job. Returns 202 Accepted. The job state becomes CANCELLED. If running, the training process on the worker is killed within seconds.

POST /jobs/{id}/upload — Upload a dataset

Upload a dataset file for the job before or after submission (referenced by dataset_ref in the job request).

GET /jobs/{id}/artifact — Download checkpoint

Downloads the latest model checkpoint as a .tar.gz archive. Available after training completes or during training (mid-training checkpoints are uploaded every 5 minutes).

GET /health

Returns 200 ok. Used by load balancers and Docker Compose health checks.


Job Lifecycle

QUEUED ──► RUNNING ──► COMPLETED
│ │ └──────────► FAILED
│
└──────────────────► CANCELLED
  • QUEUED — job created, waiting for a capable worker
  • RUNNING — dispatched to a worker, training in progress
  • COMPLETED — training finished, checkpoint available
  • FAILED — training script exited with a non-zero status
  • CANCELLED — cancelled via DELETE /jobs/{id}

If a worker dies mid-job, the stale job recovery process requeues it to QUEUED within ~90 seconds.


Architecture

┌──────────────────────────────────────────────────────────┐
│ NATS JetStream Cluster │
│ │
│ JOBS stream workers_state KV jobs_state KV │
│ jobs.train.* (TTL: 1 min) (TTL: 24h) │
│ jobs.results worker caps, job state, │
│ jobs.logs.* free/busy state RequestProto │
│ jobs.cancel │
│ workers.register datasets store checkpoints store │
│ workers.heartbeat (TTL: 48h) (TTL: 7 days) │
└──────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌────────┴────────┐ ┌────────┴────────┐
│ Coordinator │ │ Worker │
│ │ │ │
│ HTTP API │ │ RegisterWorker │
│ Job routing │ │ StartHeartbeat │
│ TryDispatch │ │ StartJobListen │
│ StaleRecovery │ │ StartCancel │
│ │ │ TrainingExec │
└─────────────────┘ └─────────────────┘

Packages

PackageResponsibility
cmd/edgegridBinary entrypoint
internal/agentBoots coordinator and/or worker from a shared NATS connection
internal/coordinatorHTTP API, job routing, dispatch, stale recovery
internal/coordinator/workermanWorker KV registry, capability matching, CAS assignment
internal/workerJob listener, heartbeat, registration, cancel listener
internal/worker/executorTraining executor (venv cache, script runner) and mock
internal/brokerNATS JetStream, KV, and Object Store helpers
internal/jobstateJob state read/write helpers
internal/proto/workerProtobuf schemas and generated Go code

Documentation

Detailed design docs are in docs/:

DocWhat it covers
intelligent-routing.mdHardware capability matching, why routing is coordinator-owned
training-executor.mdVenv caching, script execution, environment injection
job-queuing.mdFIFO queue, RequestProto persistence, CAS dispatch
log-streaming.mdJetStream + SSE, DeliverAll for late clients
job-cancellation.mdPer-job context, cancel signal broadcast, state protection
reliability.mdStale job recovery, mid-training checkpointing
nats-raft-replicas.mdRaft consensus, replication, stateless coordinator design

About

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Repository files navigation

EdgeGrid

EdgeGrid is a decentralized ML training network built on personal computers. Submit a training job via HTTP, and EdgeGrid finds a capable machine in the network, runs the job, streams logs back in real time, and stores the model checkpoint — with no manual setup on the worker machine beyond running the agent.

You Coordinator Worker (friend's PC)
│ │ │
│ POST /jobs │ │
│ {script, requirements, │ │
│ requires_gpu: true} │ │
│─────────────────────────>│ │
│ │ match GPU/RAM/disk reqs │
│ │ assign worker │
│ │─────────────────────────────>│
│ │ │ pip install (cached)
│ │ │ run train.py
│ GET /jobs/{id}/logs │ │
│─────────────────────────>│<── log lines via JetStream ──│
│<── SSE stream ───────────│ │
│ Epoch 1/10 loss=0.84 │ │
│ Epoch 2/10 loss=0.71 │ │
│ ... │ │
│ │<── checkpoint + result ──────│
│ GET /jobs/{id}/artifact │ │
│─────────────────────────>│ │
│<── model.tar.gz ─────────│ │

How it works

EdgeGrid is fully event-driven. Workers have no inbound ports from the coordinator — they pull jobs from NATS JetStream. All state (job lifecycle, worker registry, checkpoints) lives in NATS, not in the coordinator process. The coordinator can crash and restart with zero data loss.

Coordinator — HTTP API for job submission and status. Matches job hardware requirements (GPU, RAM, VRAM, disk) to registered workers. Dispatches jobs directly to the matching worker's personal NATS subject. If no worker is free, the job stays QUEUED and is auto-dispatched when capacity appears.

Worker — Registers its hardware capabilities at startup. Listens for jobs addressed to it. Runs training inside an isolated directory with a cached Python venv. Streams stdout/stderr to NATS as log lines. Pushes output/ as a checkpoint every 5 minutes during training and once on completion.

NATS JetStream — The single source of truth. Carries job messages, log lines, results, heartbeats, and cancel signals. Stores worker state and job state in KV buckets. Stores datasets and checkpoints in object store buckets. Replication is configurable for production clusters.


Features

  • Intelligent routing — jobs matched to workers by GPU, VRAM, RAM, and disk requirements
  • Job queuing — no free worker? job waits and auto-dispatches when one becomes available
  • CAS-safe dispatch — multiple coordinators can run simultaneously without double-assigning workers
  • Log streaming — real-time stdout/stderr via SSE; late-connecting clients get full replay from the start
  • Job cancellationDELETE /jobs/{id} kills the running Python process on the worker
  • Mid-training checkpointingoutput/ uploaded to object store every 5 minutes during training
  • Stale job recovery — if a worker dies, the job is automatically requeued within ~90 seconds
  • Venv caching — SHA256(requirements.txt) keyed venvs; repeated jobs with the same deps skip pip install
  • Single binary — run as coordinator, worker, or both

Getting Started

Prerequisites

  • Go 1.21+
  • NATS Server with JetStream enabled
  • Python 3 (only for training executor — auto-detected on the worker machine)

Build

git clone https://github.com/edgegrid/edgegrid.git
cd edgegrid
go build -o edgegrid ./cmd/edgegrid

Run locally (single node, dev)

# Terminal 1 — NATS
nats-server -js
# Terminal 2 — coordinator + worker (default: both enabled, mock executor)
./edgegrid
# Terminal 3 — submit a training job
curl -X POST http://localhost:8080/jobs \
-H "Content-Type: application/json" \
-d '{ "training_script": "import os\nprint(\"training...\")\nopen(os.environ[\"OUTPUT_DIR\"]+\"/model.pt\",\"w\").write(\"weights\")", "dataset_ref": "my-dataset", "requires_gpu": false }'# → {"job_id":"a1b2c3d4","status":"queued"}# Stream logs
curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# Check status
curl http://localhost:8080/jobs/a1b2c3d4
# Cancel
curl -X DELETE http://localhost:8080/jobs/a1b2c3d4

Run as separate coordinator and worker

# Coordinator only
./edgegrid -server -nats nats://localhost:4222 -port 8080
# Worker only (on another machine)
./edgegrid -client -nats nats://coordinator:4222 -executor training -worker-id worker-gpu-01

Configuration

Flags take precedence over environment variables. If neither -server nor -client is passed, both are enabled.

FlagEnv varDefaultDescription
-servertrue*Enable coordinator (HTTP API)
-clienttrue*Enable worker
-natsNATS_URLnats://localhost:4222NATS connection URL
-portPORT8080Coordinator HTTP API port
-worker-idWORKER_IDauto-generatedCustom worker identifier
-executorEXECUTORmockExecutor backend (mock or training)
-replicasNATS_REPLICAS1NATS JetStream replication factor (1=dev, 3=prod)

HTTP API

POST /jobs — Submit a training job

{
"training_script": "print('hello')",
"requirements": "torch==2.0.0\nnumpy==1.24.0",
"dataset_type": "object_store",
"dataset_ref": "my-dataset-key",
"base_model_type": "hf",
"base_model_ref": "bert-base-uncased",
"training_config_json": "{\"epochs\": 10, \"lr\": 0.001}",
"requires_gpu": true,
"min_ram_gb": 16.0,
"min_vram_gb": 8.0,
"min_disk_gb": 20.0
}

Response 202 Accepted:

{"job_id": "a1b2c3d4", "status": "queued"}

Hardware requirement fields are all optional. Omit them and any free worker qualifies.

GET /jobs/{id} — Job status

{
"job_id": "a1b2c3d4",
"state": "COMPLETED",
"worker_id": "worker-gpu-01",
"checkpoint_key": "a1b2c3d4",
"updated_at": "2026-07-01T12:34:56Z"
}

GET /jobs/{id}/logs — Live log streaming (SSE)

curl -N http://localhost:8080/jobs/a1b2c3d4/logs
# data: Epoch 1/10 loss=0.842# data: Epoch 2/10 loss=0.761# ...# event: done# data: COMPLETED

Connects via Server-Sent Events. Late-connecting clients receive all prior log lines from the beginning (JetStream DeliverAll). Stream closes with event: done when the job reaches a terminal state.

DELETE /jobs/{id} — Cancel a job

Cancels a QUEUED or RUNNING job. Returns 202 Accepted. The job state becomes CANCELLED. If running, the training process on the worker is killed within seconds.

POST /jobs/{id}/upload — Upload a dataset

Upload a dataset file for the job before or after submission (referenced by dataset_ref in the job request).

GET /jobs/{id}/artifact — Download checkpoint

Downloads the latest model checkpoint as a .tar.gz archive. Available after training completes or during training (mid-training checkpoints are uploaded every 5 minutes).

GET /health

Returns 200 ok. Used by load balancers and Docker Compose health checks.


Job Lifecycle

QUEUED ──► RUNNING ──► COMPLETED
│ │ └──────────► FAILED
│
└──────────────────► CANCELLED
  • QUEUED — job created, waiting for a capable worker
  • RUNNING — dispatched to a worker, training in progress
  • COMPLETED — training finished, checkpoint available
  • FAILED — training script exited with a non-zero status
  • CANCELLED — cancelled via DELETE /jobs/{id}

If a worker dies mid-job, the stale job recovery process requeues it to QUEUED within ~90 seconds.


Architecture

┌──────────────────────────────────────────────────────────┐
│ NATS JetStream Cluster │
│ │
│ JOBS stream workers_state KV jobs_state KV │
│ jobs.train.* (TTL: 1 min) (TTL: 24h) │
│ jobs.results worker caps, job state, │
│ jobs.logs.* free/busy state RequestProto │
│ jobs.cancel │
│ workers.register datasets store checkpoints store │
│ workers.heartbeat (TTL: 48h) (TTL: 7 days) │
└──────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌────────┴────────┐ ┌────────┴────────┐
│ Coordinator │ │ Worker │
│ │ │ │
│ HTTP API │ │ RegisterWorker │
│ Job routing │ │ StartHeartbeat │
│ TryDispatch │ │ StartJobListen │
│ StaleRecovery │ │ StartCancel │
│ │ │ TrainingExec │
└─────────────────┘ └─────────────────┘

Packages

PackageResponsibility
cmd/edgegridBinary entrypoint
internal/agentBoots coordinator and/or worker from a shared NATS connection
internal/coordinatorHTTP API, job routing, dispatch, stale recovery
internal/coordinator/workermanWorker KV registry, capability matching, CAS assignment
internal/workerJob listener, heartbeat, registration, cancel listener
internal/worker/executorTraining executor (venv cache, script runner) and mock
internal/brokerNATS JetStream, KV, and Object Store helpers
internal/jobstateJob state read/write helpers
internal/proto/workerProtobuf schemas and generated Go code

Documentation

Detailed design docs are in docs/:

DocWhat it covers
intelligent-routing.mdHardware capability matching, why routing is coordinator-owned
training-executor.mdVenv caching, script execution, environment injection
job-queuing.mdFIFO queue, RequestProto persistence, CAS dispatch
log-streaming.mdJetStream + SSE, DeliverAll for late clients
job-cancellation.mdPer-job context, cancel signal broadcast, state protection
reliability.mdStale job recovery, mid-training checkpointing
nats-raft-replicas.mdRaft consensus, replication, stateless coordinator design

About

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages