helmd hub: central server på Pi5 + agent-uplink + helmd ctl + web-login
Hub-läge i samma binär (doc/hub-plan.md): agenter ansluter utåt med delad agent-nyckel (hub-sektion i config.json), web-/CLI-klienter loggar in (lokala konton, pbkdf2; Authenticator-interface för LDAP senare) och styr alla agenters sessioner via proxy + merged SSE. Webbklienten autodetekterar hub/local via /api/mode. Nya kommandon: helmd hub [setpass|users|agentkey], helmd ctl (login/agents/sessions/ prompt/screen/answer/events …). Enhetstester + scripts/run-hub-smoke.sh (15-stegs e2e: hub + agent + ctl). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011KikHkfCiC3yELbsMN8fT9
This commit is contained in:
220
helmd/uplink.go
Normal file
220
helmd/uplink.go
Normal file
@@ -0,0 +1,220 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// uplink är agentens utåtgående anslutning till en hub: registrera med
|
||||
// den delade agent-nyckeln, long-polla efter jobb, kör jobben mot den
|
||||
// egna lokala muxen (samma handlers och token-auth som LAN-API:t) och
|
||||
// posta tillbaka svaren. Hubben nere ⇒ backoff + omregistrering;
|
||||
// det lokala API:t fungerar hela tiden oberoende av hubben.
|
||||
type uplink struct {
|
||||
url string // hubbens bas-URL, utan avslutande /
|
||||
name string // agentens namn i hubben
|
||||
host string
|
||||
agentKey string // delad nyckel ur config (hub.key)
|
||||
token string // lokala API-token — injiceras i jobb-requests
|
||||
handler http.Handler
|
||||
client *http.Client
|
||||
stop chan struct{} // stänger ner alla loopar (används av tester)
|
||||
|
||||
mu sync.Mutex
|
||||
pollKey string // per-registrering, utfärdas av hubben
|
||||
}
|
||||
|
||||
func startUplink(cfg *Config, handler http.Handler, mgr *Manager) {
|
||||
name := cfg.Hub.Name
|
||||
if name == "" {
|
||||
name = hostname()
|
||||
}
|
||||
u := &uplink{
|
||||
url: strings.TrimRight(cfg.Hub.URL, "/"), name: name, host: hostname(),
|
||||
agentKey: cfg.Hub.Key, token: cfg.Token, handler: handler,
|
||||
// klient-timeout > hubbens poll-timeout (25 s) så att pollen
|
||||
// hinner svara 204 i lugn och ro
|
||||
client: &http.Client{Timeout: 40 * time.Second},
|
||||
stop: make(chan struct{}),
|
||||
}
|
||||
go u.run()
|
||||
go u.forwardEvents(mgr)
|
||||
log.Printf("hub-uplink aktiv mot %s som %q", u.url, u.name)
|
||||
}
|
||||
|
||||
func (u *uplink) key() string {
|
||||
u.mu.Lock()
|
||||
defer u.mu.Unlock()
|
||||
return u.pollKey
|
||||
}
|
||||
|
||||
func (u *uplink) shutdown() { close(u.stop) }
|
||||
|
||||
func (u *uplink) run() {
|
||||
backoff := time.Second
|
||||
for {
|
||||
select {
|
||||
case <-u.stop:
|
||||
return
|
||||
default:
|
||||
}
|
||||
if err := u.register(); err != nil {
|
||||
log.Printf("hub: registrering mot %s misslyckades: %v (nytt försök om %s)", u.url, err, backoff)
|
||||
time.Sleep(backoff)
|
||||
if backoff < 30*time.Second {
|
||||
backoff *= 2
|
||||
}
|
||||
continue
|
||||
}
|
||||
backoff = time.Second
|
||||
log.Printf("hub: ansluten till %s som %q", u.url, u.name)
|
||||
|
||||
// tre parallella pollers så att flera jobb kan köras samtidigt;
|
||||
// när någon får auth-fel (hubben omstartad) registrerar vi om
|
||||
lost := make(chan struct{})
|
||||
var once sync.Once
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < 3; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
u.pollLoop(lost, &once)
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
}
|
||||
}
|
||||
|
||||
func (u *uplink) register() error {
|
||||
body, _ := json.Marshal(map[string]string{
|
||||
"name": u.name, "host": u.host, "version": version, "key": u.agentKey,
|
||||
})
|
||||
resp, err := u.client.Post(u.url+"/hub/register", "application/json", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
data, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
|
||||
if resp.StatusCode != 200 {
|
||||
return fmt.Errorf("hubben svarade %s: %s", resp.Status, bytes.TrimSpace(data))
|
||||
}
|
||||
var out struct {
|
||||
Key string `json:"key"`
|
||||
}
|
||||
if err := json.Unmarshal(data, &out); err != nil || out.Key == "" {
|
||||
return fmt.Errorf("oväntat registreringssvar")
|
||||
}
|
||||
u.mu.Lock()
|
||||
u.pollKey = out.Key
|
||||
u.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (u *uplink) qs() string {
|
||||
return "?name=" + url.QueryEscape(u.name) + "&key=" + url.QueryEscape(u.key())
|
||||
}
|
||||
|
||||
func (u *uplink) pollLoop(lost chan struct{}, once *sync.Once) {
|
||||
for {
|
||||
select {
|
||||
case <-lost:
|
||||
return
|
||||
case <-u.stop:
|
||||
return
|
||||
default:
|
||||
}
|
||||
resp, err := u.client.Get(u.url + "/hub/work" + u.qs())
|
||||
if err != nil {
|
||||
// hubben nere/onåbar — låt run() börja om från registrering
|
||||
once.Do(func() { close(lost) })
|
||||
return
|
||||
}
|
||||
switch resp.StatusCode {
|
||||
case 200:
|
||||
var job hubJob
|
||||
err := json.NewDecoder(io.LimitReader(resp.Body, 96<<20)).Decode(&job)
|
||||
resp.Body.Close()
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
res := u.execute(&job)
|
||||
u.postResult(res)
|
||||
case 204:
|
||||
resp.Body.Close() // inget jobb — polla igen direkt
|
||||
default:
|
||||
resp.Body.Close() // 401/404: hubben har glömt oss — registrera om
|
||||
once.Do(func() { close(lost) })
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// execute kör ett jobb mot den lokala muxen med agentens egna token.
|
||||
func (u *uplink) execute(j *hubJob) *hubResult {
|
||||
if strings.HasPrefix(j.Path, "/api/events") {
|
||||
// SSE genom jobbkanalen vore en evighetsförfrågan
|
||||
return &hubResult{ID: j.ID, Status: 404, Body: []byte(`{"error":"events proxas inte"}`)}
|
||||
}
|
||||
req, err := http.NewRequest(j.Method, j.Path, bytes.NewReader(j.Body))
|
||||
if err != nil {
|
||||
return &hubResult{ID: j.ID, Status: 400, Body: []byte(`{"error":"ogiltigt jobb"}`)}
|
||||
}
|
||||
req.Header.Set("Authorization", "Bearer "+u.token)
|
||||
if j.ContentType != "" {
|
||||
req.Header.Set("Content-Type", j.ContentType)
|
||||
}
|
||||
rec := &recorder{hdr: http.Header{}, code: 200}
|
||||
u.handler.ServeHTTP(rec, req)
|
||||
return &hubResult{ID: j.ID, Status: rec.code, Headers: rec.hdr, Body: rec.buf.Bytes()}
|
||||
}
|
||||
|
||||
func (u *uplink) postResult(res *hubResult) {
|
||||
body, err := json.Marshal(res)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
resp, err := u.client.Post(u.url+"/hub/result"+u.qs(), "application/json", bytes.NewReader(body))
|
||||
if err == nil {
|
||||
io.Copy(io.Discard, io.LimitReader(resp.Body, 4096))
|
||||
resp.Body.Close()
|
||||
}
|
||||
}
|
||||
|
||||
// forwardEvents speglar agentens lokala SSE-events upp till hubben.
|
||||
func (u *uplink) forwardEvents(mgr *Manager) {
|
||||
ch := mgr.Subscribe()
|
||||
for ev := range ch {
|
||||
if u.key() == "" {
|
||||
continue // inte registrerad än
|
||||
}
|
||||
body, _ := json.Marshal(map[string]interface{}{"events": []Event{ev}})
|
||||
resp, err := u.client.Post(u.url+"/hub/events"+u.qs(), "application/json", bytes.NewReader(body))
|
||||
if err == nil {
|
||||
io.Copy(io.Discard, io.LimitReader(resp.Body, 4096))
|
||||
resp.Body.Close()
|
||||
}
|
||||
// fel ignoreras — pollers sköter återanslutningen
|
||||
}
|
||||
}
|
||||
|
||||
// recorder är en minimal http.ResponseWriter som fångar svaret på ett
|
||||
// jobb (httptest används inte i produktionskod).
|
||||
type recorder struct {
|
||||
hdr http.Header
|
||||
code int
|
||||
buf bytes.Buffer
|
||||
}
|
||||
|
||||
func (r *recorder) Header() http.Header { return r.hdr }
|
||||
func (r *recorder) WriteHeader(code int) { r.code = code }
|
||||
func (r *recorder) Write(b []byte) (int, error) {
|
||||
return r.buf.Write(b)
|
||||
}
|
||||
Reference in New Issue
Block a user