243 lines
5.0 KiB
Go
243 lines
5.0 KiB
Go
package main
|
||
|
||
import (
|
||
"context"
|
||
"embed"
|
||
"encoding/json"
|
||
"flag"
|
||
"fmt"
|
||
"io"
|
||
"io/fs"
|
||
"log"
|
||
"net/http"
|
||
"os"
|
||
"strings"
|
||
"time"
|
||
|
||
"gpio-monitor/internal/adapter"
|
||
)
|
||
|
||
// ===== EMBED WEB =====
|
||
//
|
||
//go:embed web/*
|
||
var webFiles embed.FS
|
||
|
||
type API struct {
|
||
buf *adapter.RingBuffer
|
||
startTime time.Time
|
||
}
|
||
|
||
func writeJSON(w http.ResponseWriter, v any) {
|
||
w.Header().Set("Content-Type", "application/json")
|
||
json.NewEncoder(w).Encode(v)
|
||
}
|
||
|
||
func (a *API) HandleHealth(w http.ResponseWriter, r *http.Request) {
|
||
latest, ok := a.buf.GetLatest()
|
||
idle := time.Since(a.buf.LastWriteTime())
|
||
|
||
status := "На связи"
|
||
if idle > 15*time.Second {
|
||
status = "Нет связи"
|
||
}
|
||
|
||
writeJSON(w, map[string]any{
|
||
"status": status,
|
||
"uptime_sec": time.Since(a.startTime).Seconds(),
|
||
"last_data_ms": idle.Milliseconds(),
|
||
"pipe_alive": a.buf.IsAlive(3 * time.Second),
|
||
"has_data": ok,
|
||
"latest": latest,
|
||
"stats": a.buf.Stats(),
|
||
})
|
||
}
|
||
|
||
func (a *API) HandleLatest(w http.ResponseWriter, r *http.Request) {
|
||
latest, ok := a.buf.GetLatest()
|
||
if !ok {
|
||
writeJSON(w, map[string]string{"error": "no data"})
|
||
return
|
||
}
|
||
|
||
writeJSON(w, map[string]any{
|
||
"latest": latest,
|
||
"history": a.buf.GetLast(10),
|
||
})
|
||
}
|
||
|
||
func (a *API) HandleHistory(w http.ResponseWriter, r *http.Request) {
|
||
raw := a.buf.GetLast(100)
|
||
|
||
out := make([]int, len(raw))
|
||
for i, v := range raw {
|
||
out[i] = int(v)
|
||
}
|
||
|
||
writeJSON(w, map[string]any{
|
||
"bytes": out,
|
||
})
|
||
}
|
||
|
||
func (a *API) HandleStream(w http.ResponseWriter, r *http.Request) {
|
||
// Проверяем доступность MJPEG потока в go2rtc
|
||
resp, err := http.Get("http://localhost:1984/api/streams?src=cam_mjpeg")
|
||
camAvailable := err == nil && resp.StatusCode == 200
|
||
if resp != nil {
|
||
resp.Body.Close()
|
||
}
|
||
|
||
writeJSON(w, map[string]interface{}{
|
||
"cam": "/api/cam",
|
||
"available": camAvailable,
|
||
"source": fmt.Sprintf("http://%s:1984/api/stream.mjpeg?src=cam_mjpeg", strings.Split(r.Host, ":")[0]),
|
||
})
|
||
}
|
||
|
||
// Прокси для MJPEG потока — просто ретранслирует готовый поток из go2rtc
|
||
func handleCamProxy(w http.ResponseWriter, r *http.Request) {
|
||
log.Printf("[cam] proxy request from %s", r.RemoteAddr)
|
||
|
||
resp, err := http.Get("http://127.0.0.1:1984/api/stream.mjpeg?src=cam_mjpeg")
|
||
if err != nil {
|
||
log.Printf("[cam] go2rtc unavailable: %v", err)
|
||
http.Error(w, "camera unavailable", 503)
|
||
return
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
// Прокидываем заголовки
|
||
for k, vv := range resp.Header {
|
||
for _, v := range vv {
|
||
w.Header().Add(k, v)
|
||
}
|
||
}
|
||
|
||
w.Header().Set("Access-Control-Allow-Origin", "*")
|
||
w.WriteHeader(resp.StatusCode)
|
||
|
||
flusher, ok := w.(http.Flusher)
|
||
if !ok {
|
||
http.Error(w, "stream unsupported", 500)
|
||
return
|
||
}
|
||
|
||
buf := make([]byte, 32*1024)
|
||
|
||
for {
|
||
n, err := resp.Body.Read(buf)
|
||
|
||
if n > 0 {
|
||
_, err = w.Write(buf[:n])
|
||
if err != nil {
|
||
log.Printf("[cam] client disconnected")
|
||
return
|
||
}
|
||
flusher.Flush()
|
||
}
|
||
|
||
if err != nil {
|
||
if err != io.EOF {
|
||
log.Printf("[cam] stream ended: %v", err)
|
||
}
|
||
return
|
||
}
|
||
}
|
||
}
|
||
|
||
// ===== PIPE READER =====
|
||
func startPipeReader(ctx context.Context, buf *adapter.RingBuffer, path string) {
|
||
go func() {
|
||
buffer := make([]byte, 64)
|
||
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
default:
|
||
}
|
||
|
||
f, err := os.OpenFile(path, os.O_RDONLY, 0)
|
||
if err != nil {
|
||
time.Sleep(time.Second)
|
||
continue
|
||
}
|
||
|
||
log.Println("pipe connected")
|
||
|
||
for {
|
||
n, err := f.Read(buffer)
|
||
if err != nil {
|
||
f.Close()
|
||
break
|
||
}
|
||
|
||
for i := 0; i < n; i++ {
|
||
buf.Write(buffer[i])
|
||
}
|
||
}
|
||
|
||
log.Println("pipe disconnected, retry...")
|
||
}
|
||
}()
|
||
}
|
||
|
||
func cors(next http.HandlerFunc) http.HandlerFunc {
|
||
return func(w http.ResponseWriter, r *http.Request) {
|
||
w.Header().Set("Access-Control-Allow-Origin", "*")
|
||
w.Header().Set("Access-Control-Allow-Methods", "GET, OPTIONS")
|
||
w.Header().Set("Access-Control-Allow-Headers", "Content-Type")
|
||
|
||
if r.Method == "OPTIONS" {
|
||
w.WriteHeader(200)
|
||
return
|
||
}
|
||
|
||
next(w, r)
|
||
}
|
||
}
|
||
|
||
func main() {
|
||
pipePath := flag.String("pipe", "/tmp/gpio_pipe", "pipe path")
|
||
flag.Parse()
|
||
|
||
buf := adapter.NewRingBuffer(1024 * 10)
|
||
ctx, cancel := context.WithCancel(context.Background())
|
||
defer cancel()
|
||
|
||
startPipeReader(ctx, buf, *pipePath)
|
||
|
||
api := &API{
|
||
buf: buf,
|
||
startTime: time.Now(),
|
||
}
|
||
|
||
// ===== API =====
|
||
http.Handle("/api/", http.StripPrefix("/api", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||
switch r.URL.Path {
|
||
case "/latest":
|
||
cors(api.HandleLatest)(w, r)
|
||
case "/history":
|
||
cors(api.HandleHistory)(w, r)
|
||
case "/health":
|
||
cors(api.HandleHealth)(w, r)
|
||
case "/stream":
|
||
cors(api.HandleStream)(w, r)
|
||
case "/cam":
|
||
handleCamProxy(w, r)
|
||
default:
|
||
http.Error(w, "not found", 404)
|
||
}
|
||
})))
|
||
|
||
// ===== WEB (EMBEDDED) =====
|
||
webFS, err := fs.Sub(webFiles, "web")
|
||
if err != nil {
|
||
log.Fatal(err)
|
||
}
|
||
|
||
http.Handle("/", http.FileServer(http.FS(webFS)))
|
||
|
||
log.Println("Server started on :8080")
|
||
log.Println("Camera proxy available at /api/cam")
|
||
log.Fatal(http.ListenAndServe(":8080", nil))
|
||
} |