|
| 1 | +// Package billing connects relay to billing-service in realtime. |
| 2 | +// |
| 3 | +// Relay enforces tiers locally (cached, fail-open) and reports metered usage |
| 4 | +// upstream. Env: |
| 5 | +// |
| 6 | +// BILLING_URL e.g. http://billing:8080 (empty = billing disabled) |
| 7 | +// BILLING_API_KEY Bearer key for billing-service /v1 |
| 8 | +// BILLING_ORG_ID org uuid for usage reports + entitlement lookups |
| 9 | +// |
| 10 | +// Cadence: FetchTier caches for 5 minutes; ReportUsage is called by the |
| 11 | +// forwarder every 60s with an idempotency key of org:events:YYYYMMDDHHMM so |
| 12 | +// replays dedupe server-side. A billing outage never blocks ingest — the last |
| 13 | +// known tier stays active and usage is retried on the next tick. |
| 14 | +package billing |
| 15 | + |
| 16 | +import ( |
| 17 | + "bytes" |
| 18 | + "context" |
| 19 | + "encoding/json" |
| 20 | + "fmt" |
| 21 | + "io" |
| 22 | + "net/http" |
| 23 | + "os" |
| 24 | + "strings" |
| 25 | + "sync" |
| 26 | + "time" |
| 27 | +) |
| 28 | + |
| 29 | +// Config for the live billing link. |
| 30 | +type Config struct { |
| 31 | + URL string |
| 32 | + APIKey string |
| 33 | + OrgID string |
| 34 | + Client *http.Client |
| 35 | +} |
| 36 | + |
| 37 | +// ConfigFromEnv loads the link. Enabled() == false means run standalone. |
| 38 | +func ConfigFromEnv() Config { |
| 39 | + return Config{ |
| 40 | + URL: strings.TrimRight(strings.TrimSpace(os.Getenv("BILLING_URL")), "/"), |
| 41 | + APIKey: strings.TrimSpace(os.Getenv("BILLING_API_KEY")), |
| 42 | + OrgID: strings.TrimSpace(os.Getenv("BILLING_ORG_ID")), |
| 43 | + Client: &http.Client{Timeout: 10 * time.Second}, |
| 44 | + } |
| 45 | +} |
| 46 | + |
| 47 | +// Enabled reports whether live calls should be attempted. |
| 48 | +func (c Config) Enabled() bool { return c.URL != "" && c.APIKey != "" && c.OrgID != "" } |
| 49 | + |
| 50 | +// Tier is the cached entitlement snapshot relay enforces. |
| 51 | +type Tier struct { |
| 52 | + Plan string `json:"tier"` |
| 53 | + MaxDevices int64 `json:"max_devices"` |
| 54 | + MaxEventsPerDay int64 `json:"max_events_per_day"` |
| 55 | + FetchedAt time.Time |
| 56 | +} |
| 57 | + |
| 58 | +// Cache holds the last known tier with a 5 minute TTL. |
| 59 | +type Cache struct { |
| 60 | + mu sync.RWMutex |
| 61 | + cfg Config |
| 62 | + tier Tier |
| 63 | +} |
| 64 | + |
| 65 | +func NewCache(cfg Config) *Cache { return &Cache{cfg: cfg, tier: Tier{Plan: "local"}} } |
| 66 | + |
| 67 | +// Get returns the cached tier, refreshing in the background when stale. |
| 68 | +// Fail-open: any fetch error keeps the previous tier. |
| 69 | +func (c *Cache) Get(ctx context.Context) Tier { |
| 70 | + c.mu.RLock() |
| 71 | + tier, stale := c.tier, time.Since(c.tier.FetchedAt) > 5*time.Minute |
| 72 | + c.mu.RUnlock() |
| 73 | + if !stale || !c.cfg.Enabled() { |
| 74 | + return tier |
| 75 | + } |
| 76 | + if fresh, err := c.cfg.FetchTier(ctx); err == nil { |
| 77 | + c.mu.Lock() |
| 78 | + c.tier = fresh |
| 79 | + c.mu.Unlock() |
| 80 | + return fresh |
| 81 | + } |
| 82 | + return tier |
| 83 | +} |
| 84 | + |
| 85 | +// FetchTier GETs /v1/entitlements/{org} live. |
| 86 | +func (c Config) FetchTier(ctx context.Context) (Tier, error) { |
| 87 | + req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.URL+"/v1/entitlements/"+c.OrgID, nil) |
| 88 | + if err != nil { |
| 89 | + return Tier{Plan: "local"}, err |
| 90 | + } |
| 91 | + req.Header.Set("Authorization", "Bearer "+c.APIKey) |
| 92 | + resp, err := c.Client.Do(req) |
| 93 | + if err != nil { |
| 94 | + return Tier{Plan: "local"}, err |
| 95 | + } |
| 96 | + defer resp.Body.Close() |
| 97 | + if resp.StatusCode != http.StatusOK { |
| 98 | + return Tier{Plan: "local"}, fmt.Errorf("billing: tier fetch HTTP %d", resp.StatusCode) |
| 99 | + } |
| 100 | + var out Tier |
| 101 | + if err := json.NewDecoder(resp.Body).Decode(&out); err != nil { |
| 102 | + return Tier{Plan: "local"}, err |
| 103 | + } |
| 104 | + if out.Plan == "" { |
| 105 | + out.Plan = "pilot" |
| 106 | + } |
| 107 | + out.FetchedAt = time.Now() |
| 108 | + return out, nil |
| 109 | +} |
| 110 | + |
| 111 | +// ReportUsage POSTs a meter delta idempotently. |
| 112 | +func (c Config) ReportUsage(ctx context.Context, events int64, key string) error { |
| 113 | + if !c.Enabled() { |
| 114 | + return nil |
| 115 | + } |
| 116 | + body, _ := json.Marshal(map[string]any{ |
| 117 | + "organization_id": c.OrgID, |
| 118 | + "metrics": map[string]int64{"events": events}, |
| 119 | + "idempotency_key": key, |
| 120 | + }) |
| 121 | + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.URL+"/v1/usage", bytes.NewReader(body)) |
| 122 | + if err != nil { |
| 123 | + return err |
| 124 | + } |
| 125 | + req.Header.Set("Content-Type", "application/json") |
| 126 | + req.Header.Set("Authorization", "Bearer "+c.APIKey) |
| 127 | + if key != "" { |
| 128 | + req.Header.Set("Idempotency-Key", key) |
| 129 | + } |
| 130 | + resp, err := c.Client.Do(req) |
| 131 | + if err != nil { |
| 132 | + return err |
| 133 | + } |
| 134 | + _, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 64<<10)) |
| 135 | + resp.Body.Close() |
| 136 | + if resp.StatusCode != http.StatusCreated && resp.StatusCode != http.StatusOK { |
| 137 | + return fmt.Errorf("billing: usage HTTP %d", resp.StatusCode) |
| 138 | + } |
| 139 | + return nil |
| 140 | +} |
0 commit comments