package server import ( "fmt" "runtime" "time" "git.smesh.lol/morly/pkg/pool" "git.smesh.lol/morly/pkg/relay/wire" "git.smesh.lol/morly/pkg/transport" ) // --- Worker lifecycle --- func (s *Server) startProxyWorkers(n int32) { s.proxyIn, s.proxyOut, s.proxyDone = spawnProxyPool(n) s.proxyBusyTime = []int64{:n} s.proxyPool = pool.NewPool(n) } // spawnProxyPool forks the proxy domains in a free function: the channels and // the loop scratch belong to this frame, and only the finished slices land in // the Server's sovereign arena. Forking inside a self-mutating method's loop // left the receiver's slices unusable part-way through the pool. // // The slices are grown with push because an inline []chan T{:n} literal is // miscompiled to an empty slice; see spawnIngestPool. func spawnProxyPool(n int32) (ins []chan wire.ProxyRequest, outs []chan wire.ProxyResponse, dones []chan struct{}) { var in chan wire.ProxyRequest var out chan wire.ProxyResponse var done chan struct{} for i := 0; i < n; i++ { in, out, done = newProxyWorker() ins = push(ins, in) outs = push(outs, out) dones = push(dones, done) } return } // newProxyWorker creates the channel pair and forks one proxy domain. func newProxyWorker() (in chan wire.ProxyRequest, out chan wire.ProxyResponse, done chan struct{}) { // Root lifetime: the server stores these and polls them on every tick. runtime.SovereignSetArena(runtime.RootArena()) in = chan wire.ProxyRequest{} out = chan wire.ProxyResponse{} done = spawn(wire.ProxyWorker, in, out) runtime.SovereignRestoreArena(runtime.RootArena()) return } func (s *Server) startBlossomWorkers(n int32) { s.blossomIn, s.blossomOut, s.blossomDone = spawnBlossomPool(n) s.blossomPool = pool.NewPool(n) } // spawnBlossomPool mirrors spawnProxyPool for the blossom domains. func spawnBlossomPool(n int32) (ins []chan wire.BlossomRequest, outs []chan wire.BlossomResponse, dones []chan struct{}) { var in chan wire.BlossomRequest var out chan wire.BlossomResponse var done chan struct{} for i := 0; i < n; i++ { in, out, done = newBlossomWorker() ins = push(ins, in) outs = push(outs, out) dones = push(dones, done) } return } // newBlossomWorker creates the channel pair and forks one blossom domain. func newBlossomWorker() (in chan wire.BlossomRequest, out chan wire.BlossomResponse, done chan struct{}) { // Root lifetime: the server stores these and polls them on every tick. runtime.SovereignSetArena(runtime.RootArena()) in = chan wire.BlossomRequest{} out = chan wire.BlossomResponse{} done = spawn(wire.BlossomWorker, in, out) runtime.SovereignRestoreArena(runtime.RootArena()) return } func (s *Server) respawnProxyWorker(i int32) { in, out, done := newProxyWorker() s.proxyIn[i] = in s.proxyOut[i] = out s.proxyDone[i] = done s.proxyPool.Busy[i] = false s.proxyBusyTime[i] = 0 fmt.Println("respawned proxy worker", i) } func (s *Server) respawnBlossomWorker(i int32) { in, out, done := newBlossomWorker() s.blossomIn[i] = in s.blossomOut[i] = out s.blossomDone[i] = done s.blossomPool.Busy[i] = false fmt.Println("respawned blossom worker", i) } // --- Proxy dispatch --- func (s *Server) doDispatchProxy(fd int32, path string, headers map[string]string) (ok bool) { target := path[len("/proxy/"):] if target == "" { s.t.SendHTTP(fd, 400, map[string]string{"Content-Type": "text/plain"}, []byte("missing url\n")) return true } if transport.HasPrefix(target, s.selfHost) { direct := target[len(s.selfHost)-1:] s.t.SendHTTP(fd, 302, map[string]string{ "Location": direct, "Access-Control-Allow-Origin": "*", }, nil) return true } targetURL := []byte("https://" | target) s.nextAsyncID++ rid := s.nextAsyncID connClose := headers["connection"] == "close" req := wire.ProxyRequest{ReqID: rid, MaxBytes: 32 * 1024 * 1024, URL: targetURL} s.proxyReapStuck() if i := s.proxyPool.IdleIndex(); i >= 0 { s.asyncPending[rid] = asyncHTTPEntry{connFD: fd, connClose: connClose, createdAt: time.Now().UnixNano()} s.proxyPool.Busy[i] = true s.proxyBusyTime[i] = time.Now().UnixNano() if !s.proxyDispatchFrame(i, req) { s.proxyPool.Busy[i] = false s.proxyBusyTime[i] = 0 delete(s.asyncPending, rid) } else { return true } } if len(s.proxyQueue) < 64 { s.proxyQueue = push(s.proxyQueue, pendingProxy{connFD: fd, connClose: connClose, req: req}) return true } s.t.SendHTTP(fd, 503, map[string]string{ "Content-Type": "text/plain", "Retry-After": "1", }, []byte("proxy workers busy\n")) return true } func (s *Server) drainProxyQueue() { var i int32 var p pendingProxy for len(s.proxyQueue) > 0 { i = s.proxyPool.IdleIndex() if i < 0 { return } p = s.proxyQueue[0] s.proxyQueue = s.proxyQueue[1:] s.asyncPending[p.req.ReqID] = asyncHTTPEntry{connFD: p.connFD, connClose: p.connClose, createdAt: time.Now().UnixNano()} s.proxyPool.Busy[i] = true s.proxyBusyTime[i] = time.Now().UnixNano() if !s.proxyDispatchFrame(i, p.req) { s.proxyPool.Busy[i] = false s.proxyBusyTime[i] = 0 delete(s.asyncPending, p.req.ReqID) s.t.SendHTTP(p.connFD, 502, map[string]string{"Content-Type": "text/plain"}, []byte("proxy dispatch failed\n")) } } } func (s *Server) proxyReapStuck() { now := time.Now().UnixNano() reaped := false for i := 0; i < s.proxyPool.Len(); i++ { if s.proxyPool.Busy[i] && s.proxyBusyTime[i] > 0 { elapsed := now - s.proxyBusyTime[i] if elapsed > 15_000_000_000 { fmt.Println("proxy worker", i, "stuck", elapsed/1_000_000_000, "s, respawning") s.respawnProxyWorker(i) reaped = true } } } if reaped { s.drainProxyQueue() } } func (s *Server) proxyDispatchFrame(i int32, req wire.ProxyRequest) (ok bool) { s.proxyIn[i] <- req return true } // --- Blossom dispatch --- func (s *Server) doDispatchBlossom(fd int32, method, path string, headers map[string]string, body []byte) (ok bool) { bpath := path[len("/blossom"):] i := s.blossomPool.IdleIndex() if i >= 0 { s.nextAsyncID++ rid := s.nextAsyncID req := wire.BlossomRequest{ ReqID: rid, Dir: []byte(s.cfg.BlossomDir), Method: []byte(method), Path: []byte(bpath), ContentType: []byte(headers["content-type"]), Body: body, Upstream: []byte(s.cfg.BlossomUpstream), } s.blossomIn[i] <- req connClose := headers["connection"] == "close" s.asyncPending[rid] = asyncHTTPEntry{connFD: fd, connClose: connClose, createdAt: time.Now().UnixNano()} s.blossomPool.Busy[i] = true return true } s.t.SendHTTP(fd, 503, map[string]string{ "Content-Type": "text/plain", "Retry-After": "1", }, []byte("blossom worker busy\n")) return true } // asyncReapStuck removes asyncPending entries older than 60s. // Does NOT send a response - connFD may have been reused by another // connection by the time 60s have elapsed. func (s *Server) asyncReapStuck() { cutoff := time.Now().UnixNano() - 60_000_000_000 for rid, entry := range s.asyncPending { if entry.createdAt > 0 && entry.createdAt < cutoff { fmt.Println("asyncPending: reaping stale entry rid=", rid, "fd=", entry.connFD) delete(s.asyncPending, rid) } } } // --- Worker response handlers --- func (s *Server) pollProxyWorkers() { var resp wire.ProxyResponse for i := int32(0); i < s.proxyPool.Len(); i++ { if !workerAlive(s.proxyDone[i]) { s.respawnProxyWorker(i) s.drainProxyQueue() continue } select { case resp = <-s.proxyOut[i]: s.proxyPool.Busy[i] = false s.proxyBusyTime[i] = 0 s.drainProxyQueue() s.completeProxyResponse(resp) default: } } } func (s *Server) completeProxyResponse(resp wire.ProxyResponse) { entry, ok := s.asyncPending[resp.ReqID] if !ok { return } delete(s.asyncPending, resp.ReqID) var status int32 var h map[string]string var body []byte switch { case resp.Status < 0: status = 502 h = map[string]string{"Content-Type": "text/plain"} body = []byte("proxy: " | string(resp.Err) | "\n") case resp.Status == 415: status = 415 h = map[string]string{"Content-Type": "text/plain"} body = []byte("content-type not allowed\n") case resp.Status >= 200 && resp.Status < 300: status = 200 h = map[string]string{ "Content-Type": string(resp.ContentType), "Cross-Origin-Resource-Policy": "cross-origin", "Cache-Control": "public, max-age=86400", "Access-Control-Allow-Origin": "*", } body = resp.Body default: status = int32(resp.Status) h = map[string]string{ "Content-Type": "text/plain", "Cache-Control": "public, max-age=3600", } body = []byte(fmt.Sprintf("upstream %d\n", resp.Status)) } s.t.CompleteHTTP(entry.connFD, status, h, body, entry.connClose) resp.Body = nil resp.ContentType = nil resp.Err = nil body = nil } var blossomCORSHeaders map[string]string func (s *Server) pollBlossomWorkers() { var resp wire.BlossomResponse for i := int32(0); i < s.blossomPool.Len(); i++ { if !workerAlive(s.blossomDone[i]) { s.respawnBlossomWorker(i) continue } select { case resp = <-s.blossomOut[i]: s.blossomPool.Busy[i] = false s.completeBlossomResponse(resp) default: } } } func (s *Server) completeBlossomResponse(resp wire.BlossomResponse) { entry, ok := s.asyncPending[resp.ReqID] if !ok { return } delete(s.asyncPending, resp.ReqID) h := map[string]string{} for k, v := range blossomCORSHeaders { h[k] = v } if len(resp.CT) > 0 { h["Content-Type"] = string(resp.CT) } if resp.Size > 0 { h["Content-Length"] = fmt.Sprintf("%d", resp.Size) h["Content-Type"] = "application/octet-stream" } s.t.CompleteHTTP(entry.connFD, int32(resp.Status), h, resp.Body, entry.connClose) resp.Body = nil resp.CT = nil } func init() { blossomCORSHeaders = map[string]string{ "Access-Control-Allow-Origin": "*", "Access-Control-Allow-Methods": "GET, PUT, DELETE, HEAD, OPTIONS", "Access-Control-Allow-Headers": "Authorization, Content-Type", } }