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) }