package event import ( "bytes" "context" "encoding/json" "fmt" "log" "net/http" "time" "tianyan-edge/internal/config" ) const maxBuffer = 500 type Uploader struct { cfg *config.Config events <-chan SuspectedEvent buf []SuspectedEvent } func NewUploader(cfg *config.Config, events <-chan SuspectedEvent) *Uploader { return &Uploader{cfg: cfg, events: events} } func (u *Uploader) Run(ctx context.Context) { client := &http.Client{Timeout: 10 * time.Second} ticker := time.NewTicker(5 * time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): return case ev := <-u.events: if err := u.upload(client, ev); err != nil { if ue, ok := err.(*UploadError); ok && ue.Fatal { log.Printf("uploader: FATAL: %s. Clearing buffer and stopping.", ue.Msg) u.buf = nil // Drop all pending events // Give some time before returning or loop with delay time.Sleep(1 * time.Minute) } else { u.buffer(ev) } } case <-ticker.C: u.flushBuffer(client) } } } func (u *Uploader) buffer(ev SuspectedEvent) { if len(u.buf) < maxBuffer { u.buf = append(u.buf, ev) } else { log.Println("uploader: offline buffer full, dropping event") } } func (u *Uploader) flushBuffer(client *http.Client) { remaining := u.buf[:0] for _, ev := range u.buf { if err := u.upload(client, ev); err != nil { if ue, ok := err.(*UploadError); ok && ue.Fatal { log.Printf("uploader: FATAL flush: %s. Dropping remaining buffer.", ue.Msg) u.buf = nil return } remaining = append(remaining, ev) } } u.buf = remaining } type UploadError struct { Fatal bool Msg string } func (e *UploadError) Error() string { return e.Msg } func (u *Uploader) upload(client *http.Client, ev SuspectedEvent) error { url := fmt.Sprintf("%s/api/v1/edge/events/suspected", u.cfg.CloudURL) body, _ := json.Marshal(ev) req, _ := http.NewRequest(http.MethodPost, url, bytes.NewReader(body)) req.Header.Set("Content-Type", "application/json") req.Header.Set("Authorization", "Bearer "+u.cfg.EdgeToken) resp, err := client.Do(req) if err != nil { log.Printf("uploader: upload error: %v", err) return nil // Network error, retryable } defer resp.Body.Close() if resp.StatusCode == 401 || resp.StatusCode == 403 { return &UploadError{Fatal: true, Msg: fmt.Sprintf("unauthorized (status %d), check token", resp.StatusCode)} } if resp.StatusCode == 200 || resp.StatusCode == 201 { return nil } log.Printf("uploader: upload failed status=%d", resp.StatusCode) return nil // Server error, retryable }