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
221 lines
5.9 KiB
Go
221 lines
5.9 KiB
Go
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)
|
|
}
|