package bridge import ( "bytes" "context" "encoding/json" "errors" "fmt" "io" "net/http" "os" "path/filepath" "strings" "time" ) // Uploader sends one face image to object storage. // // It is an interface because the bridge must work identically when images are // off, when the PC is not enrolled yet, and when the server has no bucket // configured - three states that are normal rather than exceptional. type Uploader interface { // Upload returns the object key the server assigned. Upload(ctx context.Context, path string) (key string, err error) } // SpacesUploader uploads through a URL the server mints. // // The shop PC holds no bucket credentials, only its own agent token. That is // the point: the bucket is shared with other applications and a counter-top PC // is the least trustworthy machine in the estate, so a stolen one gives up a // few minutes of write access to one key rather than a bucket password. type SpacesUploader struct { // BaseURL is the server, e.g. https://mcp.loyaly.ai BaseURL string // Token is the agent's own API credential, issued at enrolment. Separate // from the broker password so rotating either does not break the other. Token string Client *http.Client } // ErrImagesOff means the server stores no images. Distinct from a failure: the // agent should stop trying and carry on sending visits, not retry forever. var ErrImagesOff = errors.New("server does not store images") // maxImageBytes bounds what will be read off disk and sent. The engine writes // ~20 KB crops; anything near this is a bug or a different file that landed in // the outbox, and a shop uplink should not spend minutes discovering that. const maxImageBytes = 2 << 20 func (u *SpacesUploader) httpClient() *http.Client { if u.Client != nil { return u.Client } // Long enough for a slow shop uplink, bounded so a half-open connection // cannot stall the queue behind it. return &http.Client{Timeout: 60 * time.Second} } type uploadTarget struct { Key string `json:"key"` URL string `json:"url"` Headers map[string]string `json:"headers"` ExpiresIn int `json:"expires_in"` } func (u *SpacesUploader) Upload(ctx context.Context, path string) (string, error) { if u.BaseURL == "" || u.Token == "" { // Not claimed yet. The visit still queues; it simply has no photo. return "", ErrImagesOff } info, err := os.Stat(path) if err != nil { return "", err } if info.Size() == 0 { return "", errors.New("image file is empty") } if info.Size() > maxImageBytes { return "", fmt.Errorf("image is %d bytes, over the %d limit", info.Size(), maxImageBytes) } body, err := os.ReadFile(path) if err != nil { return "", err } return u.UploadBytes(ctx, body) } // UploadBytes puts an image already in memory. // // Split out for camera snapshots, which the engine hands over as bytes. The // alternative - writing each frame to a temp file so Upload could read it back // - would put a picture of a shop floor on disk once a minute per camera, on // the one machine in the estate least worth trusting with it. func (u *SpacesUploader) UploadBytes(ctx context.Context, body []byte) (string, error) { if u.BaseURL == "" || u.Token == "" { return "", ErrImagesOff } if len(body) == 0 { return "", errors.New("image is empty") } if int64(len(body)) > maxImageBytes { return "", fmt.Errorf("image is %d bytes, over the %d limit", len(body), maxImageBytes) } target, err := u.target(ctx) if errors.Is(err, ErrImagesOff) { // No object storage on this server. Send the bytes to the API itself, // which holds them for a deployment that has no bucket - the same // fallback camera snapshots already take, and chosen by the SENTINEL // rather than by matching the message, because a prose change would // otherwise silently stop every photo in the estate. // // Only after target() has spoken. The unclaimed case returns the same // sentinel from the guard at the top of this function, and a PC with no // credentials has no server to PUT to either. return u.uploadDirect(ctx, body) } if err != nil { return "", err } req, err := http.NewRequestWithContext(ctx, http.MethodPut, target.URL, bytes.NewReader(body)) if err != nil { return "", err } // Sent exactly as handed back. The ACL is inside the server's signature, so // changing or dropping it does not publish the image - it fails the upload, // which is the safe direction. for k, v := range target.Headers { req.Header.Set(k, v) } req.ContentLength = int64(len(body)) resp, err := u.httpClient().Do(req) if err != nil { return "", fmt.Errorf("upload: %w", err) } defer resp.Body.Close() msg, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<10)) if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated { return "", fmt.Errorf("upload returned %s: %s", resp.Status, strings.TrimSpace(string(msg))) } return target.Key, nil } // uploadDirect posts the image to our own API, for a deployment with no bucket. // // Deliberately the second choice. A presigned PUT never passes a photograph // through the server at all, which is what makes it the right route wherever // object storage exists; this one is what stops "no S3 account" from meaning // "no customer photo, ever" on every local install and every self-hosted site. // // The server decides where it lands and returns the key, exactly as the // presigned route does. That symmetry is the point: the caller cannot tell // which route ran, so the queued visit, the read path and erasure all stay // single implementations. func (u *SpacesUploader) uploadDirect(ctx context.Context, body []byte) (string, error) { req, err := http.NewRequestWithContext(ctx, http.MethodPost, strings.TrimRight(u.BaseURL, "/")+"/api/agent/faces", bytes.NewReader(body)) if err != nil { return "", err } req.Header.Set("Authorization", "Bearer "+u.Token) req.Header.Set("Content-Type", "image/jpeg") req.ContentLength = int64(len(body)) resp, err := u.httpClient().Do(req) if err != nil { return "", fmt.Errorf("upload face: %w", err) } defer resp.Body.Close() if resp.StatusCode == http.StatusNotImplemented { // This server stores no images at all. Stop trying rather than retry // every visitor forever. return "", ErrImagesOff } if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated { msg, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<10)) return "", fmt.Errorf("upload face returned %s: %s", resp.Status, strings.TrimSpace(string(msg))) } var out struct { Key string `json:"key"` } if err := json.NewDecoder(io.LimitReader(resp.Body, 8<<10)).Decode(&out); err != nil { return "", err } if out.Key == "" { return "", errors.New("server stored the face but named no key for it") } return out.Key, nil } func (u *SpacesUploader) target(ctx context.Context) (uploadTarget, error) { var out uploadTarget req, err := http.NewRequestWithContext(ctx, http.MethodPost, strings.TrimRight(u.BaseURL, "/")+"/api/agent/upload-url", nil) if err != nil { return out, err } req.Header.Set("Authorization", "Bearer "+u.Token) resp, err := u.httpClient().Do(req) if err != nil { return out, fmt.Errorf("ask for an upload url: %w", err) } defer resp.Body.Close() if resp.StatusCode == http.StatusNotImplemented { return out, ErrImagesOff } if resp.StatusCode != http.StatusOK { body, _ := io.ReadAll(io.LimitReader(resp.Body, 4<<10)) return out, fmt.Errorf("upload url returned %s: %s", resp.Status, strings.TrimSpace(string(body))) } if err := json.NewDecoder(io.LimitReader(resp.Body, 64<<10)).Decode(&out); err != nil { return out, err } if out.URL == "" || out.Key == "" { return out, errors.New("server returned an incomplete upload target") } return out, nil } // attachImage uploads the engine's face image and returns the object key. // // Every failure is non-fatal and the local file is removed regardless. A visit // without a photo is a real visit and the number the customer pays for; a // visit stuck behind a failed upload is lost footfall. Keeping the file for a // retry would also mean an outbox that grows for as long as the failure lasts, // full of pictures of customers. func (b *Bridge) attachImage(ctx context.Context, visit map[string]any, path string) { if path == "" { return } defer func() { if err := os.Remove(path); err != nil && !os.IsNotExist(err) { b.logf("could not remove %s after upload: %v", filepath.Base(path), err) } }() if b.Uploader == nil { return } // Bounded separately from the caller: an upload that hangs must not hold // up the visit it belongs to. uctx, cancel := context.WithTimeout(ctx, 90*time.Second) defer cancel() key, err := b.Uploader.Upload(uctx, path) if err != nil { if errors.Is(err, ErrImagesOff) { // Normal for a deployment that stores no images, and for a PC that // has not been claimed yet. Not worth a line per visitor. return } b.logf("image upload failed for %s, sending the visit without it: %v", filepath.Base(path), err) return } visit["image_key"] = key }