diff --git a/cmd/server/main.go b/cmd/server/main.go index c9fe224..1b58f42 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -4,136 +4,135 @@ package main import ( "context" "encoding/json" - "fmt" + "flag" "log" "net/http" + "os" + "time" + "gpio-monitor/internal/adapter" ) -// API обработчики HTTP type API struct { buf *adapter.RingBuffer + startTime time.Time } -func (a *API) HandleLatest(w http.ResponseWriter, r *http.Request) { - if r.Method != http.MethodGet { - http.Error(w, "Method not allowed", http.StatusMethodNotAllowed) - return - } - +func (a *API) HandleHealth(w http.ResponseWriter, r *http.Request) { latest, ok := a.buf.GetLatest() - if !ok { - w.Header().Set("Content-Type", "application/json") - json.NewEncoder(w).Encode(map[string]interface{}{ - "error": "no data yet", - }) - return + idle := time.Since(a.buf.LastWriteTime()) + + status := "ok" + if idle > 15*time.Second { + status = "degraded" } - // Явное преобразование в []int для читаемого JSON - lastN := a.buf.GetLast(10) - historyInts := make([]int, len(lastN)) - for i, v := range lastN { - historyInts[i] = int(v) - } - - w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(map[string]interface{}{ - "latest_byte": latest, - "latest_hex": fmt.Sprintf("0x%02x", latest), - "latest_char": string(rune(latest)), - "history": historyInts, // теперь будет [1,2,3,...] - "stats": a.buf.Stats(), + "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) HandleNibbles(w http.ResponseWriter, r *http.Request) { - // Для 4-битной шины — показываем последние 16 полубайт - data := a.buf.GetLast(16) - - nibbles := make([]map[string]interface{}, len(data)) - for i, b := range data { - nibbles[i] = map[string]interface{}{ - "index": i, - "value": int(b), - "hex": fmt.Sprintf("0x%02x", b), - "bin": fmt.Sprintf("%04b", b), // 4 бита - "bin_full": fmt.Sprintf("%08b", b), // полный байт - } +func (a *API) HandleLatest(w http.ResponseWriter, r *http.Request) { + latest, ok := a.buf.GetLatest() + if !ok { + json.NewEncoder(w).Encode(map[string]string{"error": "no data"}) + return } - w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(map[string]interface{}{ - "nibbles": nibbles, - "count": len(nibbles), + "latest": latest, + "history": a.buf.GetLast(10), }) } func (a *API) HandleHistory(w http.ResponseWriter, r *http.Request) { - // Можно добавить параметр ?n=100 - n := 100 - data := a.buf.GetLast(n) - - w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(map[string]interface{}{ - "bytes": data, - "count": len(data), + "bytes": a.buf.GetLast(100), }) } -// CORS middleware для удобства разработки -func corsMiddleware(next http.HandlerFunc) http.HandlerFunc { +// === FIFO 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 := openPipe(path) + 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, retrying...") + } + }() +} + +func openPipe(path string) (*os.File, error) { + return os.OpenFile(path, os.O_RDONLY, 0) +} + +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(http.StatusOK) - return - } - next(w, r) } } func main() { - // Инициализация - updates := make(chan byte, 256) - buf := adapter.NewRingBuffer(1024 * 10) // 10KB буфер + pipePath := flag.String("pipe", "/tmp/gpio_pipe", "pipe path") + flag.Parse() - // Запускаем чтение stdin в фоне + buf := adapter.NewRingBuffer(1024 * 10) ctx, cancel := context.WithCancel(context.Background()) defer cancel() - reader := adapter.NewStdinReader(updates) - if err := reader.Start(ctx, buf); err != nil { - log.Fatal("Ошибка запуска читателя stdin:", err) + startPipeReader(ctx, buf, *pipePath) + + api := &API{ + buf: buf, + startTime: time.Now(), } + fs := http.FileServer(http.Dir("/home/user/temp/golang/web")) + 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) + } + }))) - // Инициализация API - api := &API{buf: buf} - - // HTTP маршруты (каждый только ОДИН раз!) - http.HandleFunc("/api/latest", corsMiddleware(api.HandleLatest)) - http.HandleFunc("/api/history", corsMiddleware(api.HandleHistory)) - http.HandleFunc("/api/nibbles", corsMiddleware(api.HandleNibbles)) - - // Статические файлы для веб-интерфейса - fs := http.FileServer(http.Dir("web")) http.Handle("/", fs) - // Запуск сервера - log.Println("=== GPIO Monitor Server ===") - log.Println("Сервер запущен на http://localhost:8080") - log.Println("API endpoints:") - log.Println(" GET /api/latest - последний байт + история") - log.Println(" GET /api/history - последние 100 байт") - log.Println(" GET /api/nibbles - 4-битное представление") - log.Println("Ожидание данных из stdin...") - log.Println("") - - if err := http.ListenAndServe(":8080", nil); err != nil { - log.Fatal("Ошибка сервера:", err) - } + log.Println("server :8080") + log.Fatal(http.ListenAndServe(":8080", nil)) } diff --git a/internal/adapter/buffer.go b/internal/adapter/buffer.go index aacba73..0cc3ec9 100644 --- a/internal/adapter/buffer.go +++ b/internal/adapter/buffer.go @@ -6,90 +6,94 @@ import ( "time" ) -// RingBuffer - потокобезопасный кольцевой буфер для байт GPIO type RingBuffer struct { - mu sync.RWMutex - data []byte - writePtr int - count int // сколько байт записано (до заполнения буфера) + mu sync.RWMutex + data []byte + writePtr int + count int + lastWrite time.Time } -// NewRingBuffer создает буфер заданного размера func NewRingBuffer(size int) *RingBuffer { return &RingBuffer{ - data: make([]byte, size), + data: make([]byte, size), lastWrite: time.Now(), } } -// Write добавляет байт в буфер func (rb *RingBuffer) Write(b byte) { rb.mu.Lock() defer rb.mu.Unlock() - + rb.data[rb.writePtr] = b rb.writePtr = (rb.writePtr + 1) % len(rb.data) + if rb.count < len(rb.data) { rb.count++ } + rb.lastWrite = time.Now() } -// GetLatest возвращает последний записанный байт +func (rb *RingBuffer) LastWriteTime() time.Time { + rb.mu.RLock() + defer rb.mu.RUnlock() + return rb.lastWrite +} + func (rb *RingBuffer) GetLatest() (byte, bool) { rb.mu.RLock() defer rb.mu.RUnlock() - + if rb.count == 0 { return 0, false } - + idx := rb.writePtr - 1 if idx < 0 { idx += len(rb.data) } + return rb.data[idx], true } -// GetLast возвращает последние N байт func (rb *RingBuffer) GetLast(n int) []byte { rb.mu.RLock() defer rb.mu.RUnlock() - + if n > rb.count { n = rb.count } - if n == 0 { - return []byte{} - } - - result := make([]byte, n) + + res := make([]byte, n) + start := rb.writePtr - n if start < 0 { start += len(rb.data) - copy(result, rb.data[start:]) - copy(result[len(rb.data)-start:], rb.data[:rb.writePtr]) + copy(res, rb.data[start:]) + copy(res[len(rb.data)-start:], rb.data[:rb.writePtr]) } else { - copy(result, rb.data[start:rb.writePtr]) + copy(res, rb.data[start:rb.writePtr]) } - return result + + return res } -// GetAll возвращает все данные из буфера -func (rb *RingBuffer) GetAll() []byte { - return rb.GetLast(rb.count) -} - -// Stats возвращает статистику буфера func (rb *RingBuffer) Stats() map[string]interface{} { rb.mu.RLock() defer rb.mu.RUnlock() - + return map[string]interface{}{ - "total_bytes": rb.count, - "buffer_size": len(rb.data), - "last_write": rb.lastWrite.Format(time.RFC3339), - "is_full": rb.count == len(rb.data), + "size": len(rb.data), + "filled": rb.count, + "last_write": rb.lastWrite, } } + +func (rb *RingBuffer) IsAlive(timeout time.Duration) bool { + rb.mu.RLock() + defer rb.mu.RUnlock() + + return time.Since(rb.lastWrite) < timeout +} diff --git a/internal/adapter/reader.go b/internal/adapter/reader.go deleted file mode 100644 index f106b8e..0000000 --- a/internal/adapter/reader.go +++ /dev/null @@ -1,60 +0,0 @@ -// internal/adapter/reader.go -package adapter - -import ( - "bufio" - "context" - "io" - "log" - "os" -) - -// StdinReader читает бинарные данные из stdin и пишет в буфер -type StdinReader struct { - updates chan byte // канал для real-time уведомлений (0 = без уведомлений) -} - -// NewStdinReader создает читалку stdin -func NewStdinReader(updates chan byte) *StdinReader { - return &StdinReader{ - updates: updates, - } -} - -// Start начинает чтение stdin в отдельной горутине -func (r *StdinReader) Start(ctx context.Context, buf *RingBuffer) error { - reader := bufio.NewReader(os.Stdin) - - go func() { - for { - select { - case <-ctx.Done(): - log.Println("StdinReader: остановлен") - return - default: - b, err := reader.ReadByte() - if err != nil { - if err == io.EOF { - log.Println("StdinReader: stdin закрыт (EOF)") - return - } - // Другие ошибки игнорируем, продолжаем читать - continue - } - - buf.Write(b) - - // Неблокирующая отправка в канал уведомлений - if r.updates != nil { - select { - case r.updates <- b: - default: - // Канал заполнен, пропускаем уведомление - } - } - } - } - }() - - return nil -} diff --git a/start-server.sh b/start-server.sh new file mode 100755 index 0000000..60f8bd5 --- /dev/null +++ b/start-server.sh @@ -0,0 +1,4 @@ +#!/bin/bash +cd /home/user/temp/golang || exit 1 + +./gpio-monitor-server -pipe /tmp/gpio_pipe diff --git a/web/index.html b/web/index.html index e61f23b..06d21c4 100644 --- a/web/index.html +++ b/web/index.html @@ -1,36 +1,125 @@
- -