File
Blob: tun/server/handler.go
| 1 | package server |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "fmt" |
| 7 | "net" |
| 8 | "net/http" |
| 9 | "net/url" |
| 10 | "sort" |
| 11 | "strconv" |
| 12 | "strings" |
| 13 | "time" |
| 14 | |
| 15 | "go.miragespace.co/specter/spec/protocol" |
| 16 | "go.miragespace.co/specter/spec/tun" |
| 17 | |
| 18 | "github.com/go-chi/chi/v5" |
| 19 | ) |
| 20 | |
| 21 | type connectedClient struct { |
| 22 | ClientID string `json:"clientId"` |
| 23 | Identity string `json:"identity"` |
| 24 | Address string `json:"address"` |
| 25 | Version string `json:"version"` |
| 26 | URL string `json:"url"` |
| 27 | SessionMode string `json:"sessionMode,omitempty"` |
| 28 | Hostname string `json:"hostname,omitempty"` |
| 29 | OwnerIdentity string `json:"ownerIdentity,omitempty"` |
| 30 | OwnerLabel string `json:"ownerLabel,omitempty"` |
| 31 | OwnerURL string `json:"ownerUrl,omitempty"` |
| 32 | } |
| 33 | |
| 34 | type connectedInfo struct { |
| 35 | Node string `json:"node"` |
| 36 | ObservedAt string `json:"observedAt"` |
| 37 | Clients []connectedClient `json:"clients"` |
| 38 | } |
| 39 | |
| 40 | type clientTunnel struct { |
| 41 | Hostname string `json:"hostname"` |
| 42 | Target string `json:"target"` |
| 43 | Configured string `json:"configured"` |
| 44 | Registered string `json:"registered"` |
| 45 | } |
| 46 | |
| 47 | type tunnelsInfo struct { |
| 48 | Identity string `json:"identity"` |
| 49 | Address string `json:"address"` |
| 50 | ObservedAt string `json:"observedAt"` |
| 51 | ConfigurationError string `json:"configurationError,omitempty"` |
| 52 | RegistrationError string `json:"registrationError,omitempty"` |
| 53 | Tunnels []clientTunnel `json:"tunnels"` |
| 54 | } |
| 55 | |
| 56 | // TunnelServerHandler serves JSON observations for connected clients and their tunnels. |
| 57 | func TunnelServerHandler(s *Server) http.Handler { |
| 58 | router := chi.NewRouter() |
| 59 | |
| 60 | router.Get("/", func(w http.ResponseWriter, r *http.Request) { |
| 61 | clients := s.TunnelTransport.ListConnected() |
| 62 | info := connectedInfo{ |
| 63 | Node: s.Identity().GetAddress(), |
| 64 | ObservedAt: time.Now().UTC().Format(time.RFC3339), |
| 65 | Clients: make([]connectedClient, 0, len(clients)), |
| 66 | } |
| 67 | |
| 68 | clientURLs := make(map[string]string, len(clients)) |
| 69 | for _, h := range clients { |
| 70 | mode, hostname, owner := s.sessions.describe(h.Physical) |
| 71 | if hostname != "" && !strings.Contains(hostname, ".") { |
| 72 | hostname += "." + s.Apex |
| 73 | } |
| 74 | clientURL := fmt.Sprintf("/_internal/tun/%d/%s", h.Identity.GetId(), url.PathEscape(h.Identity.GetAddress())) |
| 75 | clientURLs[h.Identity.GetAddress()] = clientURL |
| 76 | info.Clients = append(info.Clients, connectedClient{ |
| 77 | ClientID: strconv.FormatUint(h.Identity.GetId(), 10), |
| 78 | Identity: fmt.Sprintf("%d/%s", h.Identity.GetId(), h.Identity.GetAddress()), |
| 79 | Address: h.Addr.String(), |
| 80 | Version: h.Version, |
| 81 | URL: clientURL, |
| 82 | SessionMode: mode, |
| 83 | Hostname: hostname, |
| 84 | OwnerIdentity: owner, |
| 85 | OwnerLabel: ownerDisplayLabel(owner), |
| 86 | }) |
| 87 | } |
| 88 | for i := range info.Clients { |
| 89 | if owner := info.Clients[i].OwnerIdentity; owner != "" { |
| 90 | info.Clients[i].OwnerURL = clientURLs[owner] |
| 91 | } |
| 92 | } |
| 93 | sort.Slice(info.Clients, func(i, j int) bool { return info.Clients[i].Identity < info.Clients[j].Identity }) |
| 94 | |
| 95 | writeTunnelJSON(w, info) |
| 96 | }) |
| 97 | |
| 98 | router.Get("/{id}/*", func(w http.ResponseWriter, r *http.Request) { |
| 99 | idStr := chi.URLParam(r, "id") |
| 100 | id, err := strconv.ParseUint(idStr, 10, 64) |
| 101 | if err != nil { |
| 102 | http.Error(w, err.Error(), http.StatusBadRequest) |
| 103 | return |
| 104 | } |
| 105 | address := chi.URLParam(r, "*") |
| 106 | // Chi matches RawPath when present, so decode that parameter once. |
| 107 | if r.URL.RawPath != "" { |
| 108 | address, err = url.PathUnescape(address) |
| 109 | if err != nil { |
| 110 | http.Error(w, "invalid client address", http.StatusBadRequest) |
| 111 | return |
| 112 | } |
| 113 | } |
| 114 | |
| 115 | client := &protocol.Node{ |
| 116 | Id: id, |
| 117 | Address: address, |
| 118 | Rendezvous: true, |
| 119 | } |
| 120 | |
| 121 | callCtx, cancel := context.WithTimeout(r.Context(), lookupTimeout) |
| 122 | defer cancel() |
| 123 | |
| 124 | // default to http client pooling |
| 125 | t := http.DefaultTransport.(*http.Transport).Clone() |
| 126 | t.DisableKeepAlives = true |
| 127 | defer t.CloseIdleConnections() |
| 128 | t.DialContext = func(ctx context.Context, network, addr string) (net.Conn, error) { |
| 129 | // Transport detaches the dial context from the request. This transport |
| 130 | // is private to the lookup, so explicitly bound its dial to our deadline. |
| 131 | return s.TunnelTransport.DialStream(callCtx, client, protocol.Stream_RPC) |
| 132 | } |
| 133 | c := &http.Client{ |
| 134 | Transport: t, |
| 135 | } |
| 136 | |
| 137 | rpcClient := protocol.NewClientQueryServiceProtobufClient("http://client", c) |
| 138 | |
| 139 | prefix := tun.ClientHostnamesPrefix(&protocol.ClientToken{ |
| 140 | Token: []byte(client.GetAddress()), |
| 141 | }) |
| 142 | info := tunnelsInfo{ |
| 143 | Identity: fmt.Sprintf("%d/%s", id, address), |
| 144 | Address: address, |
| 145 | ObservedAt: time.Now().UTC().Format(time.RFC3339), |
| 146 | Tunnels: make([]clientTunnel, 0), |
| 147 | } |
| 148 | |
| 149 | // These sources answer different questions. Read both within one deadline |
| 150 | // so an unavailable client or ring cannot hide the other source's data. |
| 151 | type configurationResult struct { |
| 152 | response *protocol.ListTunnelsResponse |
| 153 | err error |
| 154 | } |
| 155 | type registrationResult struct { |
| 156 | hostnames [][]byte |
| 157 | err error |
| 158 | } |
| 159 | configurations := make(chan configurationResult, 1) |
| 160 | registrations := make(chan registrationResult, 1) |
| 161 | go func(results chan<- configurationResult) { |
| 162 | response, err := rpcClient.ListTunnels(callCtx, &protocol.ListTunnelsRequest{}) |
| 163 | results <- configurationResult{response: response, err: err} |
| 164 | }(configurations) |
| 165 | go func(results chan<- registrationResult) { |
| 166 | hostnames, err := s.Chord.PrefixList(callCtx, []byte(prefix)) |
| 167 | results <- registrationResult{hostnames: hostnames, err: err} |
| 168 | }(registrations) |
| 169 | |
| 170 | var configuration configurationResult |
| 171 | var registration registrationResult |
| 172 | for configurations != nil || registrations != nil { |
| 173 | select { |
| 174 | case configuration = <-configurations: |
| 175 | configurations = nil |
| 176 | case registration = <-registrations: |
| 177 | registrations = nil |
| 178 | case <-callCtx.Done(): |
| 179 | // Preserve results already completed when cancellation arrived. |
| 180 | if configurations != nil { |
| 181 | select { |
| 182 | case configuration = <-configurations: |
| 183 | default: |
| 184 | configuration.err = callCtx.Err() |
| 185 | } |
| 186 | } |
| 187 | if registrations != nil { |
| 188 | select { |
| 189 | case registration = <-registrations: |
| 190 | default: |
| 191 | registration.err = callCtx.Err() |
| 192 | } |
| 193 | } |
| 194 | configurations, registrations = nil, nil |
| 195 | } |
| 196 | } |
| 197 | |
| 198 | configured := make(map[string]string) |
| 199 | registered := make(map[string]bool) |
| 200 | if configuration.err != nil { |
| 201 | info.ConfigurationError = configuration.err.Error() |
| 202 | } else { |
| 203 | for _, tunnel := range configuration.response.GetTunnels() { |
| 204 | configured[tunnel.GetHostname()] = tunnel.GetTarget() |
| 205 | } |
| 206 | } |
| 207 | if registration.err != nil { |
| 208 | info.RegistrationError = registration.err.Error() |
| 209 | } else { |
| 210 | for _, hostname := range registration.hostnames { |
| 211 | registered[string(hostname)] = true |
| 212 | } |
| 213 | } |
| 214 | |
| 215 | hostnames := make(map[string]bool, len(configured)+len(registered)) |
| 216 | for hostname := range configured { |
| 217 | hostnames[hostname] = true |
| 218 | } |
| 219 | for hostname := range registered { |
| 220 | hostnames[hostname] = true |
| 221 | } |
| 222 | for hostname := range hostnames { |
| 223 | target, hasConfig := configured[hostname] |
| 224 | info.Tunnels = append(info.Tunnels, clientTunnel{ |
| 225 | Hostname: hostname, |
| 226 | Target: target, |
| 227 | Configured: observedStatus(hasConfig, configuration.err), |
| 228 | Registered: observedStatus(registered[hostname], registration.err), |
| 229 | }) |
| 230 | } |
| 231 | sort.Slice(info.Tunnels, func(i, j int) bool { return info.Tunnels[i].Hostname < info.Tunnels[j].Hostname }) |
| 232 | |
| 233 | writeTunnelJSON(w, info) |
| 234 | }) |
| 235 | |
| 236 | return router |
| 237 | } |
| 238 | |
| 239 | func ownerDisplayLabel(identity string) string { |
| 240 | if remainder, ok := strings.CutPrefix(identity, "v2:"); ok { |
| 241 | if id, _, ok := strings.Cut(remainder, ":"); ok { |
| 242 | if value, err := strconv.ParseUint(id, 10, 64); err == nil { |
| 243 | return strconv.FormatUint(value, 10) |
| 244 | } |
| 245 | } |
| 246 | } |
| 247 | label := []rune(identity) |
| 248 | if len(label) <= 24 { |
| 249 | return identity |
| 250 | } |
| 251 | return string(label[:12]) + "…" + string(label[len(label)-8:]) |
| 252 | } |
| 253 | |
| 254 | func observedStatus(present bool, err error) string { |
| 255 | if err != nil { |
| 256 | return "Unknown" |
| 257 | } |
| 258 | if present { |
| 259 | return "Yes" |
| 260 | } |
| 261 | return "No" |
| 262 | } |
| 263 | |
| 264 | func writeTunnelJSON(w http.ResponseWriter, info any) { |
| 265 | w.Header().Set("Content-Type", "application/json; charset=utf-8") |
| 266 | w.Header().Set("Cache-Control", "no-store") |
| 267 | json.NewEncoder(w).Encode(info) |
| 268 | } |