feat(telemetry): retention purge + context field whitelist + size caps (#32)
Telemetry up-gating hardening (code portion of #32): - Context field whitelist: telemetry `context` is filtered to a small set of non-content diagnostic keys (route/retryable/phase/exit_code/duration_ms/ attempt) before persistence. Unknown keys — including potentially identifying ones (email, full file path, prompt, raw IP) — are dropped, so a client regression cannot land arbitrary JSON in the store. Empty/unparseable/no-allowed-key context is dropped to "". - Per-field size cap: stack_top and context are truncated to 8KiB after redaction (backstop against unbounded blobs within batch limits). - Retention: daily master-only task deletes telemetry rows older than HEICODE_TELEMETRY_RETENTION_DAYS (default 30; <=0 disables). HEICODE_TELEMETRY_RETENTION_INTERVAL_HOURS (default 24) sets cadence. model.DeleteTelemetryEventsBefore(cutoff) + controller.StartTelemetryRetentionTask() wired into main.go under IsMasterNode. - GET /api/heicode/config telemetry block now surfaces retention_days for client/admin transparency. Tests: whitelist drop/keep, size cap, redaction-within-allowed-key. go build/vet clean; controller telemetry tests pass. Affects: Manager only (telemetry ingest + retention). No billing/consume-log change (telemetry still never bills). Privacy-doc disclosure + production enable-checklist portions of #32 tracked in heicodeDocs sync (#34) / desktop client API docs (#35). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -24,6 +24,8 @@ func HeicodeConfig(c *gin.Context) {
|
||||
"endpoint": "/api/heicode/telemetry/events",
|
||||
"max_batch": common.GetEnvOrDefault("HEICODE_TELEMETRY_MAX_BATCH", telemetryMaxBatch),
|
||||
"flush_interval_sec": common.GetEnvOrDefault("HEICODE_TELEMETRY_FLUSH_INTERVAL_SEC", 30),
|
||||
// Server retention window (#32): events older than this are purged.
|
||||
"retention_days": telemetryRetentionDays(),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
@@ -25,8 +25,65 @@ const (
|
||||
telemetryMaxBatch = 20
|
||||
telemetryMaxBodySize = 256 * 1024
|
||||
headerDeviceID = "X-Heicode-Device-Id"
|
||||
// Per-field hard cap after redaction (#32): a backstop so a single event
|
||||
// can't park an unbounded blob in the telemetry store even within batch
|
||||
// limits. stack_top / context are truncated past this many bytes.
|
||||
telemetryMaxFieldBytes = 8 * 1024
|
||||
)
|
||||
|
||||
// telemetryContextAllowedKeys whitelists the non-content diagnostic keys the
|
||||
// client may attach to an event's `context` (#32). Anything else is dropped
|
||||
// before persistence, so a client regression can't land arbitrary — possibly
|
||||
// identifying — JSON (prompts, code, tokens, emails, full file paths, raw IPs)
|
||||
// in the telemetry store. Keep in sync with the client telemetry contract
|
||||
// (winos#23) and docs/integration/heicode-desktop-client-api.md. Additions must
|
||||
// be reviewed against the "no identifying content" rule in issue #32.
|
||||
var telemetryContextAllowedKeys = map[string]bool{
|
||||
"route": true, // logical UI route, e.g. "chat" (no params)
|
||||
"retryable": true, // bool
|
||||
"phase": true, // lifecycle phase enum
|
||||
"exit_code": true, // process exit code (int)
|
||||
"duration_ms": true, // numeric timing
|
||||
"attempt": true, // retry attempt count
|
||||
}
|
||||
|
||||
// filterTelemetryContext keeps only whitelisted keys from the client-supplied
|
||||
// context object, then redacts and size-caps the result (#32). Returns "" when
|
||||
// the context is empty, unparseable, or has no allowed keys — telemetry is
|
||||
// best-effort diagnostics, so dropping an unrecognized payload is preferable to
|
||||
// storing arbitrary JSON.
|
||||
func filterTelemetryContext(raw json.RawMessage) string {
|
||||
if len(raw) == 0 {
|
||||
return ""
|
||||
}
|
||||
var obj map[string]json.RawMessage
|
||||
if err := common.Unmarshal(raw, &obj); err != nil {
|
||||
return "" // not an object (or malformed) -> drop
|
||||
}
|
||||
filtered := make(map[string]json.RawMessage, len(obj))
|
||||
for k, v := range obj {
|
||||
if telemetryContextAllowedKeys[k] {
|
||||
filtered[k] = v
|
||||
}
|
||||
}
|
||||
if len(filtered) == 0 {
|
||||
return ""
|
||||
}
|
||||
b, err := common.Marshal(filtered)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
return capTelemetryField(model.RedactText(string(b)))
|
||||
}
|
||||
|
||||
// capTelemetryField truncates an already-redacted field to telemetryMaxFieldBytes.
|
||||
func capTelemetryField(s string) string {
|
||||
if len(s) <= telemetryMaxFieldBytes {
|
||||
return s
|
||||
}
|
||||
return s[:telemetryMaxFieldBytes]
|
||||
}
|
||||
|
||||
type telemetryEventIn struct {
|
||||
ClientId string `json:"client_id"`
|
||||
SchemaVersion int `json:"schema_version"`
|
||||
@@ -73,13 +130,12 @@ func (e telemetryEventIn) toModel(userID int, deviceID string, now int64) model.
|
||||
stackTopJSON := ""
|
||||
if len(e.StackTop) > 0 {
|
||||
if b, err := common.Marshal(e.StackTop); err == nil {
|
||||
stackTopJSON = model.RedactText(string(b))
|
||||
stackTopJSON = capTelemetryField(model.RedactText(string(b)))
|
||||
}
|
||||
}
|
||||
ctxJSON := ""
|
||||
if len(e.Context) > 0 {
|
||||
ctxJSON = model.RedactText(string(e.Context))
|
||||
}
|
||||
// #32: context is restricted to a key whitelist (drop arbitrary/identifying
|
||||
// JSON), then redacted and size-capped.
|
||||
ctxJSON := filterTelemetryContext(e.Context)
|
||||
return model.TelemetryEvent{
|
||||
ReceivedAt: now,
|
||||
UserId: userID,
|
||||
|
||||
@@ -65,12 +65,45 @@ func TestTelemetryToModel_RedactsSecrets(t *testing.T) {
|
||||
ev := telemetryEventIn{
|
||||
ClientId: "dev-1",
|
||||
StackTop: []string{"at boom (auth.ts) key=sk-abcDEF1234567890"},
|
||||
Context: json.RawMessage(`{"hdr":"Authorization: Bearer aZ09tokenVALUE","ok":true}`),
|
||||
// "route" is whitelisted (#32) so it survives the field filter; redaction
|
||||
// must still strip the Bearer token carried inside an allowed key.
|
||||
Context: json.RawMessage(`{"route":"Authorization: Bearer aZ09tokenVALUE","retryable":true}`),
|
||||
}
|
||||
m := ev.toModel(7, "dev-1", 1700000000)
|
||||
|
||||
require.NotContains(t, m.StackTopJSON, "sk-abcDEF1234567890", "sk- secret must be redacted in stack_top")
|
||||
require.Contains(t, m.StackTopJSON, "REDACTED")
|
||||
require.NotContains(t, m.ContextJSON, "aZ09tokenVALUE", "Bearer token must be redacted in context")
|
||||
require.Contains(t, m.ContextJSON, "ok") // non-secret content preserved
|
||||
require.Contains(t, m.ContextJSON, "retryable") // non-secret whitelisted content preserved
|
||||
}
|
||||
|
||||
// #32: context must be restricted to a key whitelist so a client regression
|
||||
// cannot land arbitrary/identifying JSON in the telemetry store.
|
||||
func TestFilterTelemetryContext_Whitelist(t *testing.T) {
|
||||
// allowed keys kept, unknown keys (incl. potentially identifying) dropped
|
||||
out := filterTelemetryContext(json.RawMessage(
|
||||
`{"route":"chat","retryable":true,"email":"a@b.com","file":"C:/Users/x/secret.go","prompt":"hi"}`))
|
||||
require.Contains(t, out, "route")
|
||||
require.Contains(t, out, "retryable")
|
||||
require.NotContains(t, out, "email")
|
||||
require.NotContains(t, out, "a@b.com")
|
||||
require.NotContains(t, out, "secret.go")
|
||||
require.NotContains(t, out, "prompt")
|
||||
|
||||
// no allowed keys -> dropped entirely
|
||||
require.Equal(t, "", filterTelemetryContext(json.RawMessage(`{"email":"a@b.com"}`)))
|
||||
// non-object / malformed -> dropped
|
||||
require.Equal(t, "", filterTelemetryContext(json.RawMessage(`"a string"`)))
|
||||
require.Equal(t, "", filterTelemetryContext(json.RawMessage(`not json`)))
|
||||
require.Equal(t, "", filterTelemetryContext(nil))
|
||||
}
|
||||
|
||||
// #32: per-field size cap is a backstop against unbounded blobs.
|
||||
func TestCapTelemetryField(t *testing.T) {
|
||||
require.Equal(t, "short", capTelemetryField("short"))
|
||||
big := make([]byte, telemetryMaxFieldBytes+100)
|
||||
for i := range big {
|
||||
big[i] = 'a'
|
||||
}
|
||||
require.Len(t, capTelemetryField(string(big)), telemetryMaxFieldBytes)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/heicode/manager/common"
|
||||
"github.com/heicode/manager/model"
|
||||
)
|
||||
|
||||
// Telemetry retention (#32). Account-linkable client error-telemetry must not be
|
||||
// kept indefinitely: a daily task deletes rows older than the retention window.
|
||||
// Disabled when telemetry ingest is off (default) or retention <= 0.
|
||||
//
|
||||
// - HEICODE_TELEMETRY_RETENTION_DAYS (default 30): rows received earlier than
|
||||
// now-RETENTION are deleted. <=0 disables the purge (keep-forever — only for
|
||||
// explicit operator opt-out; not recommended for production).
|
||||
// - HEICODE_TELEMETRY_RETENTION_INTERVAL_HOURS (default 24): sweep cadence.
|
||||
//
|
||||
// Master-only (wired from main.go under IsMasterNode) so multiple nodes don't
|
||||
// all sweep the shared LOG_DB.
|
||||
|
||||
var telemetryRetentionTaskOnce sync.Once
|
||||
|
||||
func telemetryRetentionDays() int {
|
||||
return common.GetEnvOrDefault("HEICODE_TELEMETRY_RETENTION_DAYS", 30)
|
||||
}
|
||||
|
||||
// StartTelemetryRetentionTask launches the daily telemetry retention sweep.
|
||||
func StartTelemetryRetentionTask() {
|
||||
telemetryRetentionTaskOnce.Do(func() {
|
||||
// Only meaningful once ingest is enabled; if the endpoint is off there is
|
||||
// nothing being written, but we still allow the sweep to drain any rows
|
||||
// captured during a prior enabled window. Gate on retention days instead.
|
||||
days := telemetryRetentionDays()
|
||||
if days <= 0 {
|
||||
common.SysLog("telemetry retention task disabled (HEICODE_TELEMETRY_RETENTION_DAYS<=0; rows kept indefinitely)")
|
||||
return
|
||||
}
|
||||
intervalHours := common.GetEnvOrDefault("HEICODE_TELEMETRY_RETENTION_INTERVAL_HOURS", 24)
|
||||
if intervalHours < 1 {
|
||||
intervalHours = 24
|
||||
}
|
||||
go func() {
|
||||
time.Sleep(5 * time.Minute) // avoid startup churn
|
||||
runTelemetryRetentionOnce()
|
||||
ticker := time.NewTicker(time.Duration(intervalHours) * time.Hour)
|
||||
defer ticker.Stop()
|
||||
for range ticker.C {
|
||||
runTelemetryRetentionOnce()
|
||||
}
|
||||
}()
|
||||
common.SysLog(fmt.Sprintf("telemetry retention task started: retention=%dd interval=%dh", days, intervalHours))
|
||||
})
|
||||
}
|
||||
|
||||
func runTelemetryRetentionOnce() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
common.SysLog(fmt.Sprintf("telemetry retention task panic recovered: %v", r))
|
||||
}
|
||||
}()
|
||||
days := telemetryRetentionDays()
|
||||
if days <= 0 {
|
||||
return
|
||||
}
|
||||
cutoff := time.Now().Unix() - int64(days)*86400
|
||||
deleted, err := model.DeleteTelemetryEventsBefore(cutoff)
|
||||
if err != nil {
|
||||
common.SysLog("telemetry retention task: " + err.Error())
|
||||
return
|
||||
}
|
||||
if deleted > 0 {
|
||||
common.SysLog(fmt.Sprintf("telemetry retention task: deleted %d events older than %d days", deleted, days))
|
||||
}
|
||||
}
|
||||
@@ -134,6 +134,9 @@ func main() {
|
||||
// retention window (issue #4). Master-only so multiple nodes don't all purge.
|
||||
if common.IsMasterNode {
|
||||
controller.StartSecretPurgeTask()
|
||||
// Telemetry retention: daily purge of client error-telemetry older than
|
||||
// HEICODE_TELEMETRY_RETENTION_DAYS (#32). Master-only.
|
||||
controller.StartTelemetryRetentionTask()
|
||||
}
|
||||
|
||||
if common.IsMasterNode && constant.UpdateTask {
|
||||
|
||||
@@ -39,3 +39,15 @@ func InsertTelemetryEvents(events []TelemetryEvent) error {
|
||||
}
|
||||
return LOG_DB.Create(&events).Error
|
||||
}
|
||||
|
||||
// DeleteTelemetryEventsBefore removes telemetry rows received before cutoffUnix
|
||||
// (server unix seconds), enforcing the retention window (#32). Returns the
|
||||
// number of rows deleted. Account-linkable device telemetry must not be kept
|
||||
// indefinitely; callers run this from a periodic retention task.
|
||||
func DeleteTelemetryEventsBefore(cutoffUnix int64) (int64, error) {
|
||||
if LOG_DB == nil {
|
||||
return 0, nil
|
||||
}
|
||||
res := LOG_DB.Where("received_at < ?", cutoffUnix).Delete(&TelemetryEvent{})
|
||||
return res.RowsAffected, res.Error
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user