Files
telemtPanel/apps/agent/main.go
T
Denozordec b5f31c1083
Build and Push Telemt Panel Docker Image / build-and-push (push) Failing after 38s
Build and Push Telemt Panel Docker Image / create-release (push) Skipped
Init
2026-08-04 18:33:48 +07:00

199 lines
5.3 KiB
Go

// Telemt Panel agent (fleet mode)
// Build: go build -o telemt-panel-agent .
// Run: ./telemt-panel-agent -panel-url https://panel -token <enrollment> OR -state /var/lib/...
package main
import (
"bytes"
"encoding/json"
"flag"
"fmt"
"io"
"log"
"net/http"
"os"
"path/filepath"
"time"
)
type enrollResponse struct {
AgentID string `json:"agentId"`
AgentToken string `json:"agentToken"`
PanelURL string `json:"panelUrl"`
}
type job struct {
ID string `json:"id"`
Type string `json:"type"`
Payload map[string]any `json:"payload"`
}
type stateFile struct {
AgentID string `json:"agentId"`
AgentToken string `json:"agentToken"`
PanelURL string `json:"panelUrl"`
}
func main() {
panelURL := flag.String("panel-url", os.Getenv("PANEL_URL"), "panel base URL")
token := flag.String("token", os.Getenv("ENROLLMENT_TOKEN"), "enrollment token")
statePath := flag.String("state", "/var/lib/telemt-panel-agent/state.json", "state file")
telemtURL := flag.String("telemt-url", envOr("TELEMT_API_URL", "http://127.0.0.1:9091"), "local Telemt API")
telemtAuth := flag.String("telemt-auth", os.Getenv("TELEMT_AUTH_HEADER"), "Telemt Authorization header")
flag.Parse()
st, err := loadOrEnroll(*statePath, *panelURL, *token)
if err != nil {
log.Fatal(err)
}
client := &http.Client{Timeout: 30 * time.Second}
log.Printf("agent %s connected to %s", st.AgentID, st.PanelURL)
for {
jobs, err := fetchJobs(client, st)
if err != nil {
log.Printf("jobs poll: %v", err)
time.Sleep(5 * time.Second)
continue
}
for _, j := range jobs {
if j.Type != "telemt.proxy" {
_ = postResult(client, st, j.ID, false, nil, "unsupported job type")
continue
}
result, err := doTelemtProxy(client, *telemtURL, *telemtAuth, j.Payload)
if err != nil {
_ = postResult(client, st, j.ID, false, nil, err.Error())
continue
}
_ = postResult(client, st, j.ID, true, result, "")
}
time.Sleep(2 * time.Second)
}
}
func envOr(k, def string) string {
if v := os.Getenv(k); v != "" {
return v
}
return def
}
func loadOrEnroll(path, panelURL, token string) (*stateFile, error) {
if b, err := os.ReadFile(path); err == nil {
var st stateFile
if json.Unmarshal(b, &st) == nil && st.AgentToken != "" {
return &st, nil
}
}
if panelURL == "" || token == "" {
return nil, fmt.Errorf("need -panel-url and -token for first enroll (or existing state file)")
}
hostname, _ := os.Hostname()
body, _ := json.Marshal(map[string]string{
"token": token,
"hostname": hostname,
"name": hostname,
"agentVersion": "0.1.0",
})
res, err := http.Post(panelURL+"/api/agent/enroll", "application/json", bytes.NewReader(body))
if err != nil {
return nil, err
}
defer res.Body.Close()
raw, _ := io.ReadAll(res.Body)
if res.StatusCode >= 300 {
return nil, fmt.Errorf("enroll %s: %s", res.Status, string(raw))
}
var er enrollResponse
if err := json.Unmarshal(raw, &er); err != nil {
return nil, err
}
st := &stateFile{AgentID: er.AgentID, AgentToken: er.AgentToken, PanelURL: er.PanelURL}
_ = os.MkdirAll(filepath.Dir(path), 0o755)
b, _ := json.MarshalIndent(st, "", " ")
_ = os.WriteFile(path, b, 0o600)
return st, nil
}
func fetchJobs(client *http.Client, st *stateFile) ([]job, error) {
req, _ := http.NewRequest(http.MethodGet, st.PanelURL+"/api/agent/jobs", nil)
req.Header.Set("Authorization", "Bearer "+st.AgentToken)
res, err := client.Do(req)
if err != nil {
return nil, err
}
defer res.Body.Close()
raw, _ := io.ReadAll(res.Body)
if res.StatusCode >= 300 {
return nil, fmt.Errorf("%s: %s", res.Status, string(raw))
}
var jobs []job
if err := json.Unmarshal(raw, &jobs); err != nil {
return nil, err
}
return jobs, nil
}
func postResult(client *http.Client, st *stateFile, id string, ok bool, result any, errMsg string) error {
payload := map[string]any{"ok": ok, "result": result, "error": errMsg}
b, _ := json.Marshal(payload)
req, _ := http.NewRequest(http.MethodPost, st.PanelURL+"/api/agent/jobs/"+id+"/result", bytes.NewReader(b))
req.Header.Set("Authorization", "Bearer "+st.AgentToken)
req.Header.Set("Content-Type", "application/json")
res, err := client.Do(req)
if err != nil {
return err
}
defer res.Body.Close()
return nil
}
func doTelemtProxy(client *http.Client, base, auth string, payload map[string]any) (any, error) {
method, _ := payload["method"].(string)
if method == "" {
method = "GET"
}
path, _ := payload["path"].(string)
if path == "" {
return nil, fmt.Errorf("missing path")
}
url := stringsTrimSlash(base) + path
var body io.Reader
if payload["body"] != nil && method != "GET" && method != "DELETE" {
b, _ := json.Marshal(payload["body"])
body = bytes.NewReader(b)
}
req, err := http.NewRequest(method, url, body)
if err != nil {
return nil, err
}
req.Header.Set("Accept", "application/json")
if auth != "" {
req.Header.Set("Authorization", auth)
}
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
res, err := client.Do(req)
if err != nil {
return nil, err
}
defer res.Body.Close()
raw, _ := io.ReadAll(res.Body)
var out any
if json.Unmarshal(raw, &out) != nil {
out = map[string]any{"ok": false, "raw": string(raw), "status": res.StatusCode}
}
return out, nil
}
func stringsTrimSlash(s string) string {
for len(s) > 0 && s[len(s)-1] == '/' {
s = s[:len(s)-1]
}
return s
}