package main import ( "cmp" "context" "embed" "encoding/binary" "encoding/json" "errors" "flag" "fmt" "html/template" "io" "io/fs" "log" "net" "net/http" "net/url" "os" "os/signal" "strings" "time" "github.com/coder/websocket" "github.com/kami/maven/internal/audio" "github.com/kami/maven/internal/ipc" "github.com/kami/maven/internal/voice" ) // presenceSignals — the only fact keys /api/signal may write. mavweb is a // network-facing surface inside wg; an allowlist keeps a compromised caller // boxed to forging weak presence signals (reachability, multi-source, never // truth) — it can't write arbitrary facts. ponytail: floor auth (wg-only); a // per-signal token belongs here if the tunnel ever hosts untrusted devices. var presenceSignals = map[string]string{ "desk_active": "infer:hyprland", "page_heartbeat": "infer:heartbeat", "wg_handshake": "infer:wg", } //go:embed static/* var staticFiles embed.FS //go:embed dash.html var dashHTML string // dashTmpl — the monitoring read surface, server-rendered from dash.html (no JS, // no client fetch); meta-refresh keeps it live. html/template escapes the user // text in facts/nudges. Read-only: browses the append-only store via CoreAPI, // never writes — the store IS the audit trail, this just shows it. var dashTmpl = template.Must(template.New("dash").Funcs(template.FuncMap{ "ago": func(t time.Time) string { return time.Since(t).Round(time.Second).String() + " ago" }, }).Parse(dashHTML)) func noCache(h http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Cache-Control", "no-cache, no-store, must-revalidate") h.ServeHTTP(w, r) }) } func main() { addr := flag.String("addr", ":9200", "HTTP listen address") voiceAddr := flag.String("voice", "127.0.0.1:9100", "voice server TCP addr (host:port)") // ntfyWS: the ntfy WebSocket subscribe URL the PWA connects to for in-app // nudge delivery, e.g. wss://ntfy.kvmx.ru/maven/ws?auth=. The // client subscribes directly (lowest overhead — mavweb isn't in the path); // we only serve it the URL so the deny-all auth token stays deployment // config, never baked into the static JS. Empty ⇒ /api/ntfy returns 204 and // the PWA skips subscription (voice-only, as before). ntfyWS := flag.String("ntfy", "", "ntfy WebSocket subscribe URL served to the PWA (e.g. wss://host/topic/ws?auth=...)") // coreSock: mavend's IPC socket. When set, /api/signal writes presence // facts through CoreAPI (page heartbeat from the PWA, desk_active from a PC // script). Empty ⇒ /api/signal returns 503 and presence stays unfed. coreSock := flag.String("core", "", "mavend IPC socket path for presence-signal ingest (empty = disabled)") flag.Parse() var core ipc.CoreAPI if *coreSock != "" { c, err := ipc.Dial(*coreSock) if err != nil { log.Fatalf("dial core %s: %v", *coreSock, err) } defer c.Close() core = c } mux := http.NewServeMux() sub, err := fs.Sub(staticFiles, "static") if err != nil { log.Fatalf("static fs: %v", err) } mux.Handle("/", noCache(http.FileServer(http.FS(sub)))) mux.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) { handleWS(w, r, *voiceAddr) }) mux.HandleFunc("/api/ptt", func(w http.ResponseWriter, r *http.Request) { handlePTT(w, r, *voiceAddr) }) mux.HandleFunc("/api/ping", func(w http.ResponseWriter, r *http.Request) { w.Write([]byte("pong")) }) mux.HandleFunc("/api/ntfy", func(w http.ResponseWriter, r *http.Request) { if *ntfyWS == "" { w.WriteHeader(http.StatusNoContent) // not configured → PWA skips return } w.Header().Set("Content-Type", "text/plain") w.Write([]byte(*ntfyWS)) }) mux.HandleFunc("/api/signal", func(w http.ResponseWriter, r *http.Request) { handleSignal(w, r, core) }) mux.HandleFunc("/dash", func(w http.ResponseWriter, r *http.Request) { handleDash(w, r, core) }) // /tools — the authed enable surface. maven proposes acts she can't run; // this page is where a human reviews and enables them (proposed→enabled). // Enabling is the boundary-moving act (maven.md), so it lives ONLY here, // behind wg+nginx+auth — never the voice/chat path. mux.HandleFunc("/tools", func(w http.ResponseWriter, r *http.Request) { handleTools(w, r, core) }) srv := &http.Server{Addr: *addr, Handler: mux} go func() { sig := make(chan os.Signal, 1) signal.Notify(sig, os.Interrupt) <-sig log.Println("shutting down...") srv.Close() }() log.Printf("mavweb listening on %s, voice → %s", *addr, *voiceAddr) if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { log.Fatal(err) } } func handleWS(w http.ResponseWriter, r *http.Request, voiceAddr string) { conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{ OriginPatterns: []string{"*"}, }) if err != nil { log.Printf("ws accept: %v", err) return } defer conn.Close(websocket.StatusNormalClosure, "bye") ctx := r.Context() var d net.Dialer tc, err := d.DialContext(ctx, "tcp", voiceAddr) if err != nil { log.Printf("dial voice: %v", err) writeWSErr(conn, ctx, "voice unavailable") return } defer tc.Close() for { _, msg, err := conn.Read(ctx) if err != nil { log.Printf("ws read: %v", err) return } if len(msg) < 4 { log.Printf("ws msg too short (%d bytes)", len(msg)) continue } log.Printf("ws got %d bytes from client", len(msg)) pcm := audio.Audio{Format: audio.PCM16kMono, Bytes: msg} req := voice.Request{ ID: uint64(time.Now().UnixNano()), Method: voice.MethodPushToTalk, Params: mustMarshal(voice.PushToTalkReq{ Audio: pcm, Lang: "mixed", Surface: voice.SurfacePCClient, }), } if err := writeFrame(tc, &req); err != nil { log.Printf("write voice req: %v", err) return } // Read frames until we get the matching Response (handling any interleaved Pushes) for { resp, push, err := readOneFrame(tc) if err != nil { log.Printf("read voice: %v", err) return } if push != nil { data, _ := json.Marshal(push) conn.Write(ctx, websocket.MessageText, data) continue } if resp.Error != nil { writeWSErr(conn, ctx, resp.Error.Message) break } var pttResp voice.PushToTalkResp if err := json.Unmarshal(resp.Result, &pttResp); err != nil { log.Printf("unmarshal resp: %v", err) break } if pttResp.ReplyText != "" { conn.Write(ctx, websocket.MessageText, []byte(pttResp.ReplyText)) } if len(pttResp.ReplyAudio.Bytes) > 0 { conn.Write(ctx, websocket.MessageBinary, pttResp.ReplyAudio.Bytes) } break } } } // handleSignal ingests one presence signal and writes a fresh fact through // CoreAPI. The fact's timestamp (now) is all the presence scorer reads; value // is a marker. Only allowlisted keys are accepted (see presenceSignals). func handleSignal(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI) { if r.Method != http.MethodPost { http.Error(w, "POST only", http.StatusMethodNotAllowed) return } if core == nil { http.Error(w, "presence ingest disabled (no -core)", http.StatusServiceUnavailable) return } key := r.URL.Query().Get("key") source, ok := presenceSignals[key] if !ok { http.Error(w, "unknown signal key", http.StatusBadRequest) return } // kind=env: an observation about the device/surface, NOT a self-fact — a // passive signal never writes truth about you (spec), it only feeds // presence. confidence 1.0: the reading ("input happened") is certain; // presence applies its own per-signal weight/decay on top. if _, err := core.WriteFact(r.Context(), ipc.WriteFactReq{ Ts: time.Now(), Kind: "env", Key: key, Value: `"active"`, Source: source, Confidence: 1.0, }); err != nil { log.Printf("signal %s: %v", key, err) http.Error(w, "write failed", http.StatusBadGateway) return } w.WriteHeader(http.StatusNoContent) } func handleDash(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI) { if core == nil { http.Error(w, "dash disabled (no -core)", http.StatusServiceUnavailable) return } ctx := r.Context() pres, err1 := core.Presence(ctx) facts, err2 := core.RecentFacts(ctx, 50) nudges, err3 := core.RecentNudges(ctx, 50) notes, err4 := core.RecentNotes(ctx, 50) if err := cmp.Or(err1, err2, err3, err4); err != nil { log.Printf("dash: %v", err) http.Error(w, "core read failed", http.StatusBadGateway) return } w.Header().Set("Content-Type", "text/html; charset=utf-8") if err := dashTmpl.Execute(w, struct { Presence ipc.Presence Facts []ipc.Fact Nudges []ipc.Nudge Notes []ipc.Note }{pres, facts, nudges, notes}); err != nil { log.Printf("dash render: %v", err) } } // toolsTmpl — the enable surface. Server-rendered, no JS: a plain HTML form // POSTs back to /tools to enable a proposal. html/template escapes tool names + // utterances (they came from voice STT — untrusted text). var toolsTmpl = template.Must(template.New("tools").Funcs(template.FuncMap{ "join": strings.Join, }).Parse(toolsHTML)) const toolsHTML = `maven · tools

maven · tools

{{if .Msg}}
{{.Msg}}
{{end}}

proposed ({{len .Proposed}})

{{if .Proposed}}

maven drafted these from acts she couldn't run. Fill the command (argv, space-separated) and enable.

{{range .Proposed}}{{end}}
namefrom utteranceenable as
{{.Name}}{{.Utterance}}
{{else}}

none pending.

{{end}}

enabled ({{len .Enabled}})

{{if .Enabled}} {{range .Enabled}}{{end}}
namecommand
{{.Name}}{{join .Cmd " "}} {{if .Destructive}}destructive{{end}}
{{else}}

none enabled.

{{end}} ` // handleTools serves the enable surface (GET) and applies an enable (POST). // POST fields: name, cmd (space-separated argv), destructive (checkbox). cmd is // whitespace-split — argv with embedded spaces isn't supported (ponytail: no // shell-word parsing; the box owner controls this input, quote a wrapper script // if an arg needs spaces). func handleTools(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI) { if core == nil { http.Error(w, "tools disabled (no -core)", http.StatusServiceUnavailable) return } ctx := r.Context() var msg string if r.Method == http.MethodPost { name := strings.TrimSpace(r.FormValue("name")) cmd := strings.Fields(r.FormValue("cmd")) destructive := r.FormValue("destructive") != "" if name == "" || len(cmd) == 0 { http.Error(w, "name and cmd required", http.StatusBadRequest) return } if err := core.EnableTool(ctx, name, cmd, destructive, time.Now()); err != nil { log.Printf("tools: enable %q: %v", name, err) http.Error(w, "enable failed: "+err.Error(), http.StatusBadGateway) return } msg = "enabled " + name } proposed, err1 := core.ListTools(ctx, "proposed") enabled, err2 := core.ListTools(ctx, "enabled") if err := cmp.Or(err1, err2); err != nil { log.Printf("tools: %v", err) http.Error(w, "core read failed", http.StatusBadGateway) return } w.Header().Set("Content-Type", "text/html; charset=utf-8") if err := toolsTmpl.Execute(w, struct { Msg string Proposed []ipc.Tool Enabled []ipc.Tool }{msg, proposed, enabled}); err != nil { log.Printf("tools render: %v", err) } } func writeWSErr(conn *websocket.Conn, ctx context.Context, msg string) { conn.Write(ctx, websocket.MessageText, []byte(`{"error":"`+msg+`"}`)) } func writeFrame(w io.Writer, v any) error { body, err := json.Marshal(v) if err != nil { return fmt.Errorf("marshal: %w", err) } const maxFrame = 64 << 20 if len(body) > maxFrame { return fmt.Errorf("frame too large: %d", len(body)) } var hdr [4]byte binary.BigEndian.PutUint32(hdr[:], uint32(len(body))) if _, err := w.Write(hdr[:]); err != nil { return err } _, err = w.Write(body) return err } func readFrame(r io.Reader, v any) error { var hdr [4]byte if _, err := io.ReadFull(r, hdr[:]); err != nil { return err } n := binary.BigEndian.Uint32(hdr[:]) const maxFrame = 64 << 20 if n > maxFrame { return fmt.Errorf("frame too large: %d", n) } buf := make([]byte, n) if _, err := io.ReadFull(r, buf); err != nil { return err } return json.Unmarshal(buf, v) } func readOneFrame(r io.Reader) (*voice.Response, *voice.Push, error) { var raw struct { ID uint64 `json:"id"` Result json.RawMessage `json:"r,omitempty"` Error *voice.RpcError `json:"e,omitempty"` Kind voice.PushKind `json:"kind,omitempty"` Params json.RawMessage `json:"p,omitempty"` } if err := readFrame(r, &raw); err != nil { return nil, nil, err } if raw.Kind != "" && raw.ID == 0 { return nil, &voice.Push{Kind: raw.Kind, Params: raw.Params}, nil } return &voice.Response{ID: raw.ID, Result: raw.Result, Error: raw.Error}, nil, nil } func handlePTT(w http.ResponseWriter, r *http.Request, voiceAddr string) { if r.Method != http.MethodPost { http.Error(w, "POST only", 405) return } body, err := io.ReadAll(r.Body) if err != nil { http.Error(w, err.Error(), 400) return } if len(body) < 4 { http.Error(w, "too short", 400) return } log.Printf("ptt got %d bytes from client", len(body)) pcm := audio.Audio{Format: audio.PCM16kMono, Bytes: body} var d net.Dialer tc, err := d.DialContext(r.Context(), "tcp", voiceAddr) if err != nil { log.Printf("ptt dial voice: %v", err) http.Error(w, "voice unavailable", 503) return } defer tc.Close() req := voice.Request{ ID: uint64(time.Now().UnixNano()), Method: voice.MethodPushToTalk, Params: mustMarshal(voice.PushToTalkReq{ Audio: pcm, Lang: "mixed", Surface: voice.SurfacePCClient, }), } if err := writeFrame(tc, &req); err != nil { log.Printf("ptt write: %v", err) http.Error(w, err.Error(), 500) return } for { resp, push, err := readOneFrame(tc) if err != nil { log.Printf("ptt read: %v", err) http.Error(w, err.Error(), 500) return } if push != nil { continue } if resp.Error != nil { http.Error(w, resp.Error.Message, 500) return } var pttResp voice.PushToTalkResp if err := json.Unmarshal(resp.Result, &pttResp); err != nil { http.Error(w, err.Error(), 500) return } w.Header().Set("Content-Type", "audio/l16;rate=16000;channels=1") w.Header().Set("X-Reply-Text", url.QueryEscape(pttResp.ReplyText)) w.Write(pttResp.ReplyAudio.Bytes) return } } func mustMarshal(v any) json.RawMessage { b, err := json.Marshal(v) if err != nil { panic(err) } return b }